::nsmcp::CachedLogSearch ::nsmcp::CachedLogSearch encodeEvent flush init lineStart physicalLine record recordStart removeQueries results search temporaryFile ::nsmcp::FileLogSource ::nsmcp::FileLogSource appendText consume emptyState init joinLine read ::nsmcp::CachedLogSearch->::nsmcp::FileLogSource ::nsmcp::LogSource ::nsmcp::LogSource read ::nsmcp::FileLogSource->::nsmcp::LogSource

Class ::nsmcp::CachedLogSearch

::nsmcp::CachedLogSearch[i] create ... \
           [ -cacheDirectory:required cacheDirectory:required ] \
           [ -cacheKey cacheKey ] \
           [ -maxCacheBytes:integer (default "67108864") ] \
           [ -maxEvents:integer (default "1000") ] \
           [ -maxMatchesPerCall:integer (default "1000") ] \
           [ -maxMessageBytes:integer (default "65536") ] \
           [ -maxReadBytes:integer (default "1048576") ] \
           [ -maxSearches:integer (default "16") ] \
           [ -mutex (default "") ] \
           [ -path:required path:required ] \
           [ -recordParser (default "::nsmcp::logRecordParser") ] \
           [ -store:required store:required ]

Defined in /usr/local/ns/tcl/nsmcp/lib/log-search.tcl

Class Relations

  • class: ::nx::Class[i]
  • superclass: ::nsmcp::FileLogSource[i]
::nx::Class create ::nsmcp::CachedLogSearch \
     -superclass ::nsmcp::FileLogSource

Methods (to be applied on instances)

  • flush (scripted, public)

     <instance of nsmcp::CachedLogSearch[i]> flush

    Testcases:
    No testcase defined.
    ns_mutex lock ${:mutex}
    try {
        set key ${:cacheKey}:search
        if {[nsv_exists ${:store} $key]} {
            set meta [nsv_get ${:store} $key]
            :removeQueries [dict get $meta queries]
            dict set meta queries {}
            dict incr meta generation
            nsv_set ${:store} $key $meta
        }
    } finally {ns_mutex unlock ${:mutex}}
  • results (scripted, public)

     <instance of nsmcp::CachedLogSearch[i]> results handle [ cursor ] \
        [ limit ]
    Parameters:
    handle (required)
    cursor (optional, defaults to "0")
    limit (optional, defaults to "20")

    Testcases:
    No testcase defined.
    if {![regexp {^[0-9a-f]{64}$} $handle]} {error "invalid result handle"}
    if {![string is entier -strict $cursor] || $cursor < 0 ||
        ![string is integer -strict $limit] || $limit < 1 || $limit > 100} {
        error "expected a nonnegative cursor and a limit between 1 and 100"
    }
    ns_mutex lock ${:mutex}
    try {
        set key ${:cacheKey}:search
        set query {}
        if {[nsv_exists ${:store} $key]} {
            set meta [nsv_get ${:store} $key]
            dict for {q candidate} [dict get $meta queries] {
                if {[dict exists $candidate resultHandle] &&
                    [dict get $candidate resultHandle] eq $handle} {
                    set query $candidate; break
                }
            }
        }
        if {$query eq ""} {
            return -code error -errorcode {NSMCP RESULT UNKNOWN} "unknown or expired result handle"
        }
        file stat ${:path} current
        set probe [open ${:path} rb]
        try {set prefix [read $probe [string length [dict get $meta prefix]]]} finally {close $probe}
        if {[list $current(dev) $current(ino)] ne [dict get $meta identity] ||
            $current(size) < [dict get $meta observedSize] || $prefix ne [dict get $meta prefix]} {
            error "result handle expired after log rotation; repeat the search"
        }
        set end [dict get $query fileBytes]
        if {$cursor > $end} {error "cursor is beyond the published results"}
        set f [open [dict get $query resultFile] rb]
        try {
            if {$cursor > 0} {
                seek $f [expr {$cursor-1}]
                if {[read $f 1] ne "\n"} {error "cursor must be a record boundary"}
            }
            seek $f $cursor
            set records {}; set bytes 0
            while {[tell $f] < $end && [llength $records] < $limit} {
                set start [tell $f]
                set line {}
                while {1} {
                    set position [tell $f]
                    set chunk [read $f [expr {min(65536,$end-$position,2097153-[string length $line])}]]
                    set newline [string first \n $chunk]
                    if {$newline >= 0} {
                        append line [string range $chunk 0 [expr {$newline-1}]]
                        seek $f [expr {$position+$newline+1}]
                        break
                    }
                    append line $chunk
                    if {[string length $line] >= 2097153} {error "cached record exceeds the page byte limit"}
                    if {$chunk eq "" || [tell $f] >= $end} {error "incomplete result cache"}
                }
                set length [expr {[tell $f]-$start}]
                if {$length > 2097152} {error "cached record exceeds the page byte limit"}
                if {$bytes+$length > 2097152} {seek $f $start; break}
                incr bytes $length
                lappend records [ns_json parse [encoding convertfrom utf-8 $line]]
            }
            set next [tell $f]
        } finally {close $f}
        set more [expr {$next < $end}]
        return [dict create resultHandle $handle records $records  nextCursor [expr {$more ? $next : ""}] more $more  resultCount [dict get $query resultCount]  scanIncomplete [dict get $query scanIncomplete] pending [dict get $query pending]  cacheFull [dict get $query cacheFull] generation [dict get $meta generation]]
    } finally {ns_mutex unlock ${:mutex}}
  • search (scripted, public)

     <instance of nsmcp::CachedLogSearch[i]> search literal [ severity ] \
        [ correlationField ]
    Parameters:
    literal (required)
    severity (optional)
    correlationField (optional)

    Testcases:
    No testcase defined.
    if {$correlationField ni {{} system access}} {
        error "expected an empty, system or access correlation field"
    }
    set identifierCheck [expr {$correlationField eq "access" ? "validConnectionIdentifier" : "validRequestIdentifier"}]
    if {$correlationField ne "" && ![${:recordParser} $identifierCheck $literal]} {
        error "correlation requires a complete request identifier"
    }
    if {$severity ne "" && ![regexp {^[A-Za-z][A-Za-z0-9_()-]*$} $severity]} {
        error "expected a log severity name"
    }
    if {$literal eq "" || [string length [encoding convertto utf-8 $literal]] > 4096 ||
        [string first \n $literal] >= 0 || [string first \r $literal] >= 0 ||
        [string first \x00 $literal] >= 0} {
        error "expected a nonempty, single-line literal search key of at most 4096 bytes"
    }
    ns_mutex lock ${:mutex}
    try {
        set metaKey ${:cacheKey}:search
        set queryKey [ns_crypto::md string -digest sha256 -- [list access-metadata-v5 $literal $severity $correlationField]]
        file stat ${:path} before
        set source [open ${:path} rb]
        set scanner [open ${:path} rb]
        try {
            fconfigure $source -translation binary -encoding binary
            fconfigure $scanner -translation binary -encoding binary
            file stat ${:path} opened
            set identity [list $before(dev) $before(ino)]
            if {$identity ne [list $opened(dev) $opened(ino)]} {error "log rotated while opening; retry"}
            set reset false; set generation 1
            if {[nsv_exists ${:store} $metaKey]} {
                set meta [nsv_get ${:store} $metaKey]
                set generation [dict get $meta generation]
                set prefix [dict get $meta prefix]
                if {$identity ne [dict get $meta identity] ||
                    $before(size) < [dict get $meta observedSize] ||
                    [read $source [string length $prefix]] ne $prefix} {
                    :removeQueries [dict get $meta queries]
                    set reset true; incr generation
                }
            } else {set reset true}
            if {$reset} {
                seek $source 0
                set meta [dict create identity $identity generation $generation  observedSize $before(size) queries {}  prefix [read $source [expr {min(256,$before(size))}]]]
            }
            if {[dict get $meta prefix] eq "" && $before(size) > 0} {
                seek $source 0
                dict set meta prefix [read $source [expr {min(256,$before(size))}]]
            }
            set queries [dict get $meta queries]
            if {![dict exists $queries $queryKey]} {
                if {[dict size $queries] >= ${:maxSearches}} {
                    set oldest {}; set oldestTime {}
                    dict for {key query} $queries {
                        if {$oldest eq "" || [dict get $query used] < $oldestTime} {
                            set oldest $key; set oldestTime [dict get $query used]
                        }
                    }
                    :removeQueries [dict create $oldest [dict get $queries $oldest]]
                    dict unset queries $oldest
                }
                dict set queries $queryKey [dict create resultFile [:temporaryFile]  pendingFile [:temporaryFile] resultHandle [ns_crypto::randombytes -encoding hex 32] cursor 0 resultCount 0 fileBytes 0  pending false used [clock clicks -milliseconds] cacheFull false  sourceBytes -1 scanIncomplete true]
            }
            set query [dict get $queries $queryKey]
            # Publish ownership before processing, so an interrupted scan's
            # files remain discoverable for invalidation/eviction.
            dict set meta queries $queries
            nsv_set ${:store} $metaKey $meta
            if {[dict get $query sourceBytes] == $before(size) &&
                ![dict get $query scanIncomplete]} {
                dict set query used [clock clicks -milliseconds]
                dict set meta queries $queryKey $query
                dict set meta observedSize $before(size)
                nsv_set ${:store} $metaKey $meta
                return [dict merge $query [dict create generation $generation reset $reset]]
            }
            set output [open [dict get $query resultFile] r+]
            set pending [open [dict get $query pendingFile] w]
            try {
                fconfigure $output -encoding utf-8 -translation lf
                fconfigure $pending -encoding utf-8 -translation lf
                # Discard writes beyond the last published checkpoint after
                # a failed scan, rather than appending duplicate results.
                ns_ftruncate $output [dict get $query fileBytes]
                seek $output 0 end
                dict set query pending false
                set cursor [dict get $query cursor]; set matched 0
                while {$cursor < $before(size) && $matched < ${:maxMatchesPerCall} &&
                       ![dict get $query cacheFull]} {
                    # Keep native searching separate from buffered Tcl record reads.
                    seek $scanner $cursor
                    set found [ns_fseekchars $scanner $literal]
                    if {$found < 0 || $found >= $before(size)} {
                        set cursor [:recordStart $source [expr {$before(size)-1}] $before(size)]
                        break
                    }
                    set start [:recordStart $source $found $before(size)]
                    set event [:record $source $start $before(size)]
                    incr matched
                    set metadata [${:recordParser} metadata [dict get $event message]]
                    set accepted [expr {$severity eq "" || [dict get $metadata severity] eq $severity}]
                    if {$correlationField eq "system"} {
                        set accepted [expr {$accepted && [dict get $metadata requestIdentifier] eq $literal}]
                    } elseif {$correlationField eq "access"} {
                        set accepted [expr {$accepted && [${:recordParser} accessRequestIdentifier [dict get $event message]] eq $literal}]
                    }
                    if {!$accepted} {
                        if {![dict get $event complete]} {set cursor $start; break}
                        set cursor [dict get $event endOffset]
                        continue
                    }
                    set json [:encodeEvent $event]
                    if {![dict get $event complete]} {
                        puts $pending $json
                        dict set query pending true
                        set cursor $start
                        break
                    }
                    set bytes [string length [encoding convertto utf-8 $json\n]]
                    if {[tell $output]+$bytes > ${:maxCacheBytes}} {
                        dict set query cacheFull true
                        set cursor $start
                        break
                    }
                    puts $output $json
                    dict incr query resultCount
                    set cursor [dict get $event endOffset]
                }
                flush $output
                dict set query fileBytes [tell $output]
                dict set query cursor $cursor
                dict set query used [clock clicks -milliseconds]
            } finally {close $output; close $pending}
            file stat ${:path} after
            seek $source 0
            set prefix [dict get $meta prefix]
            if {$identity ne [list $after(dev) $after(ino)] ||
                $after(size) < $before(size) || [read $source [string length $prefix]] ne $prefix} {
                error "log changed during scan; retry"
            }
            # cursor may point to the pending final message. Distinguish that
            # deliberate recheck from a match-budget backlog.
            dict set query scanIncomplete [expr {$matched >= ${:maxMatchesPerCall} || [dict get $query cacheFull]}]
            dict set query sourceBytes $before(size)
            dict set meta observedSize $before(size)
            dict set meta queries $queryKey $query
            nsv_set ${:store} $metaKey $meta
            return [dict merge $query [dict create generation $generation reset $reset]]
        } finally {close $source; close $scanner}
    } finally {ns_mutex unlock ${:mutex}}