Class ::nsmcp::CachedLogSearch (public)

 ::nx::Class ::nsmcp::CachedLogSearch[i]

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

Testcases:
No testcase defined.
Source code:
    :property cacheDirectory:required
    :property {maxSearches:integer 16}
    :property {maxMatchesPerCall:integer 1000}
    :property {maxCacheBytes:integer 67108864}
    :method init {} {
        next
        if {${:maxSearches} < 1 || ${:maxMatchesPerCall} < 1 || ${:maxCacheBytes} < 1} {
            error "search cache limits must be positive"
        }
        set :cacheDirectory [file normalize ${:cacheDirectory}]
        if {![file exists ${:cacheDirectory}]} {
            file mkdir ${:cacheDirectory}
            file attributes ${:cacheDirectory} -permissions 0700
        }
    }
    :method removeQueries {queries} {
        dict for {key query} $queries {
            file delete -force [dict get $query resultFile] [dict get $query pendingFile]
        }
    }
    :public method flush {} {
        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}}
    }
    # Cursors are byte offsets in the completed JSONL cache, never source paths.
    :public method results {handle {cursor 0} {limit 20}} {
        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}}
    }
    :method temporaryFile {} {
        set channel [file tempfile name [file join ${:cacheDirectory} nsmcp-]]
        close $channel
        return $name
    }
    :method lineStart {channel offset} {
        set position $offset
        while {$position > 0} {
            set start [expr {max(0, $position-32768)}]
            seek $channel $start
            set block [read $channel [expr {$position-$start}]]
            set newline [string last \n $block]
            if {$newline >= 0} {return [expr {$start+$newline+1}]}
            set position $start
        }
        return 0
    }
    :method physicalLine {channel limit} {
        set start [tell $channel]; set sample {}; set clipped false; set terminated false
        while {[tell $channel] < $limit} {
            set position [tell $channel]
            set block [read $channel [expr {min(32768,$limit-$position)}]]
            if {$block eq ""} break
            set newline [string first \n $block]
            if {$newline >= 0} {
                set part [string range $block 0 [expr {$newline-1}]]
                seek $channel [expr {$position+$newline+1}]
                set terminated true
            } else {set part $block}
            set room [expr {${:maxMessageBytes}-[string length $sample]}]
            if {[string length $part] > $room} {set clipped true}
            append sample [string range $part 0 [expr {$room-1}]]
            if {$terminated} break
        }
        set text [encoding convertfrom utf-8 [string trimright $sample \r]]
        regsub -all {\x1b\[[0-9;]*m} $text {} text
        return [dict create start $start end [tell $channel] text $text  truncated $clipped terminated $terminated]
    }
    :method recordStart {channel offset limit} {
        set start [:lineStart $channel $offset]
        while {1} {
            seek $channel $start
            set line [:physicalLine $channel $limit]
            if {$start == 0 || [string index [dict get $line text] 0] ne ":"} {return $start}
            set start [:lineStart $channel [expr {$start-1}]]
        }
    }
    :method record {channel start limit} {
        seek $channel $start
        set event [dict create offset $start endOffset $start message {}  complete false truncated false orphanContinuation false]
        set first true
        while {[tell $channel] < $limit} {
            set line [:physicalLine $channel $limit]
            set text [dict get $line text]
            if {!$first && [string index $text 0] ne ":"} {
                dict set event complete true
                seek $channel [dict get $line start]
                return $event
            }
            if {$first} {
                dict set event orphanContinuation [expr {[string index $text 0] eq ":"}]
            } else {set text \n$text}
            set event [:appendText $event $text]
            if {[dict get $line truncated]} {dict set event truncated true}
            dict set event endOffset [dict get $line end]
            set first false
        }
        return $event
    }
    :method encodeEvent {event} {
        set event [dict merge $event [${:recordParser} metadata [dict get $event message]]]
        set triples {}
        foreach key {timestamp processThread logThread severity requestIdentifier peerAddress method url protocol} {
            lappend triples $key string [dict get $event $key]
        }
        foreach key {status responseBytes} {
            set value [dict get $event $key]
            lappend triples $key [expr {$value eq "" ? "string" : "number"}] $value
        }
        foreach key {offset endOffset} {lappend triples $key number [dict get $event $key]}
        lappend triples message string [dict get $event message]
        foreach key {complete truncated orphanContinuation} {
            lappend triples $key boolean [dict get $event $key]
        }
        return [ns_json value -type object -- $triples]
    }
    :public method search {literal {severity {}} {correlationField {}}} {
        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}}
    }
XQL Not present:
Generic, PostgreSQL, Oracle
[ hide source ] | [ make this the default ]
Show another procedure: