Class ::nsmcp::AccessLogIngestor (public)
::nx::Class ::nsmcp::AccessLogIngestor
Defined in /usr/local/ns/tcl/nsmcp/lib/access-analytics.tcl
- Testcases:
- No testcase defined.
Source code: :property path:required :property {parser ::nsmcp::accessRecordParser} :property {maxReadBytes:integer 16777216} :property {maxLineBytes:integer 65536} :property {retentionSeconds:double 86400} :property {maxSeconds:integer 2} :property {mutex {}} :method retentionCutoff {} { # Retain the complete minute containing the cutoff, including fractions. return [expr {wide(floor(([clock seconds] - ${:retentionSeconds}) / 60.0)) * 60}] } :method init {} { foreach value [list ${:maxReadBytes} ${:maxLineBytes} ${:retentionSeconds} ${:maxSeconds}] { if {$value<1} {error "positive analytics limits required"} } if {${:retentionSeconds}<60} {error "retentionSeconds must be at least 60"} if {${:maxReadBytes}<=${:maxLineBytes}} {error "maxReadBytes must exceed maxLineBytes"} set :path [file normalize ${:path}] if {${:mutex} eq ""} {set :mutex [ns_mutex create]} } :method initial {} { return [dict create version 1 identity {} generation 0 cursor 0 anchor {} dropping 0 parsed 0 malformed 0 oversized 0 rotationGaps 0 lastRotationAt 0 observedSize 0 updatedAt 0 pending 0 firstObserved {} latestObserved {} latestCompletionObserved {} lastMalformedAt {}] } :method recentState {} { return [dict create recentInitialized false recentFrom {} initialOffset 0 historyCursor 0 historyEnd 0 historyTurn false historyDropping 0 seek {}] } :method sample {f position length} { seek $f $position set bytes [read $f $length] set start 0 if {$position > 0} { set newline [string first \n $bytes] if {$newline < 0} {return {}} set start [expr {$newline+1}] } # Skip malformed/oversized physical rows without unbounded probing. for {set i 0} {$i < 64} {incr i} { set end [string first \n $bytes $start] if {$end < 0} break if {$end-$start <= ${:maxLineBytes}} { set line [encoding convertfrom utf-8 [string range $bytes $start [expr {$end-1}]]] set time [${:parser} completionTime $line] if {$time ne ""} {return [list [expr {$position+$start}] $time]} } set start [expr {$end+1}] } return {} } :method locateRecent {f stateVar size from deadline budgetVar} { upvar 1 $stateVar state $budgetVar budget set origin [dict get $state cursor] set block [expr {min(1048576,${:maxReadBytes})}] set plan [dict get $state seek] if {$plan eq ""} { set plan [dict create phase expand high $size low $origin distance $block end $size origin $origin from $from target [expr {($from/60)*60-300}]] } set found -1 while {$budget >= $block && [clock milliseconds] < $deadline} { if {[dict get $plan phase] eq "expand"} { set position [expr {max($origin,[dict get $plan end]-[dict get $plan distance])}] } else { if {[dict get $plan high]-[dict get $plan low] <= $block} {set found [dict get $plan low]; break} set position [expr {([dict get $plan low]+[dict get $plan high])/2}] } set sample [:sample $f $position $block] incr budget -$block if {$sample eq ""} { # Do not skip unknown log formats or oversized regions. set found $origin break } lassign $sample offset time if {$position == $origin && [dict get $plan high]-$origin <= $block} { set found $origin break } if {$time < [dict get $plan target]} { dict set plan low $offset dict set plan phase bisect } elseif {$position == $origin} { set found $origin break } else { dict set plan high $position if {[dict get $plan phase] eq "expand"} {dict set plan distance [expr {2*[dict get $plan distance]}]} } } if {$found < 0} {dict set state seek $plan; return false} dict set state seek {} dict set state recentInitialized true dict set state recentFrom [expr {([dict get $plan from]/60)*60}] dict set state initialOffset $found dict set state historyCursor $origin dict set state historyEnd $found dict set state cursor $found seek $f [expr {max(0,$found-64)}] dict set state anchor [read $f [expr {min(64,$found)}]] return true } :public method ingest {{from {}} {allowHistory false}} { if {$from eq ""} {set from [expr {[clock seconds]-3600}]} if {![string is entier -strict $from] || $from < 0} {error "expected an epoch-second window start"} ns_mutex lock ${:mutex} try {return [:ingestLocked $from $allowHistory]} finally {ns_mutex unlock ${:mutex}} } :method ingestLocked {from allowHistory} { set state [:state] file stat ${:path} before set identity [list $before(dev) $before(ino)] set deadline [expr {[clock milliseconds]+1000*${:maxSeconds}}] set locating false set historical false set f [open ${:path} rb] try { # Compare a checkpoint anchor as well as inode/size (copytruncate # may have grown past the old offset before the next observation). set cursor [dict get $state cursor] seek $f [expr {max(0,$cursor-64)}] set anchor [read $f [expr {min(64,$cursor)}]] if {$identity ne [dict get $state identity] || $before(size) < $cursor || $anchor ne [dict get $state anchor]} { if {[dict get $state generation] > 0} { dict incr state rotationGaps dict set state lastRotationAt [clock seconds] } dict incr state generation dict set state identity $identity dict set state cursor 0 dict set state dropping 0 set state [dict merge $state [:recentState]] set cursor 0 } set readBudget ${:maxReadBytes} if {![dict get $state recentInitialized]} { # Existing checkpoints already near the requested interval can # keep their forward cursor; lagging old caches migrate in place. set latest [dict get $state latestObserved] if {$latest ne "" && $latest >= ($from/60)*60-300} { dict set state recentInitialized true dict set state recentFrom 0 } elseif {![:locateRecent $f state $before(size) $from $deadline readBudget]} { set locating true } set cursor [dict get $state cursor] } set liveCursor $cursor set historyRemaining [expr {[dict get $state historyEnd]-[dict get $state historyCursor]}] if {!$locating && $historyRemaining > 0 && (($from < [dict get $state recentFrom]) || ($allowHistory && [dict get $state historyTurn] && $before(size)-$cursor < ${:maxReadBytes}))} { set historical true set cursor [dict get $state historyCursor] set liveDropping [dict get $state dropping] dict set state dropping [dict get $state historyDropping] } if {$allowHistory && !$locating} {dict set state historyTurn [expr {!$historical}]} set limit $readBudget if {$locating} {set limit 0} if {$historical} {set limit [expr {min($limit,$historyRemaining)}]} seek $f $cursor set bytes [read $f $limit] set changed {}; set consumed 0; set lastCompletionStamp {}; set completion {}; set batchCompletion {} set pending 0 set cutoff [expr {[:retentionCutoff] / 60}] while {$consumed < [string length $bytes]} { if {[clock milliseconds] >= $deadline} break set end [string first \n $bytes $consumed] if {$end < 0} { if {[string length $bytes]-$consumed > ${:maxLineBytes} || [dict get $state dropping]} { if {![dict get $state dropping]} {dict incr state oversized} dict set state dropping 1 set consumed [string length $bytes] } set pending [expr {$cursor+[string length $bytes] >= $before(size)}] break } set line [string range $bytes $consumed [expr {$end-1}]] set consumed [expr {$end+1}] if {[dict get $state dropping]} {dict set state dropping 0; continue} if {[string length $line] > ${:maxLineBytes}} {dict incr state oversized; continue} set entry [${:parser} parse [encoding convertfrom utf-8 $line]] if {$entry eq ""} {dict incr state malformed; dict set state lastMalformedAt [clock seconds]; :rejected $line [expr {$cursor+$consumed-[string length $line]-1}] $state; continue} set stamp [dict get $entry timestamp] if {$stamp ne $lastCompletionStamp} { set completion [${:parser} completionTime [encoding convertfrom utf-8 $line]] set lastCompletionStamp $stamp } if {$completion ne "" && ($batchCompletion eq "" || $completion>$batchCompletion)} {set batchCompletion $completion} dict incr state parsed set time [dict get $entry startTime] if {[dict get $state firstObserved] eq "" || $time < [dict get $state firstObserved]} {dict set state firstObserved $time} if {[dict get $state latestObserved] eq "" || $time > [dict get $state latestObserved]} {dict set state latestObserved $time} :accumulate $entry [expr {$cursor+$consumed-[string length $line]-1}] changed state $cutoff } # Reuse completion timestamps while the logged second is unchanged. if {$batchCompletion ne ""} { set previous [dict get $state latestCompletionObserved] if {$previous eq "" || $batchCompletion>$previous} {dict set state latestCompletionObserved $batchCompletion} } if {$historical} { dict set state historyCursor [expr {$cursor+$consumed}] dict set state historyDropping [dict get $state dropping] dict set state dropping $liveDropping set cursor $liveCursor } else { dict set state cursor [expr {$cursor+$consumed}] set cursor [dict get $state cursor] seek $f [expr {max(0,$cursor-64)}] dict set state anchor [read $f [expr {min(64,$cursor)}]] } } finally {close $f} file stat ${:path} after if {[list $after(dev) $after(ino)] ne $identity || $after(size) < $cursor} { error "access log rotated during ingestion; retry" } if {!$historical && !$locating} { dict set state pending [expr {$pending && $cursor+[string length $bytes]-$consumed >= $after(size)}] } dict set state observedSize $after(size) dict set state updatedAt [clock seconds] :publish $state $changed $cutoff return [dict create sourceBytes $after(size) cursor $cursor scanIncomplete [expr {$locating || ($cursor < $after(size) && ![dict get $state pending])}] historyUnreadBytes [expr {[dict get $state historyEnd]-[dict get $state historyCursor]}] generation [dict get $state generation]] } :public method validateWindow {options} { foreach key [dict keys $options] {if {$key ni {from to bucketSeconds}} {error "unsupported summary option"}} set now [clock seconds] set options [dict merge [dict create from [expr {$now-3600}] to $now bucketSeconds 60] $options] foreach key {from to bucketSeconds} { if {![string is entier -strict [dict get $options $key]]} {error "expected integer epoch times and bucketSeconds"} } set from [dict get $options from]; set to [dict get $options to]; set size [dict get $options bucketSeconds] if {$from >= $to || $to-$from > 86400 || $from < 0 || $to > $now || $size < 60 || $size > 3600 || $size%60 != 0} { error "expected a window of at most 24 hours ending no later than now and bucketSeconds a multiple of 60 between 60 and 3600" } return $options } :method summaryLocked {options} { set requestedFrom [dict get $options from]; set requestedTo [dict get $options to] dict set options from [expr {([dict get $options from]/60)*60}] dict set options to [expr {(([dict get $options to]+59)/60)*60}] set state [:state] set result [dict merge $options [dict create server [ns_info server] requestCount 0 responseBytes 0 fallbackCount 0 pools {} methods {} statuses {} paths {} excludedPaths 0 timingCount 0 acceptTotal 0.0 queueTotal 0.0 filterTotal 0.0 runTotal 0.0 runMax 0.0 buckets {}]] set counts [dict get $result buckets] for {set slot [dict get $options from]} {$slot < [dict get $options to]} {incr slot [dict get $options bucketSeconds]} { set key [expr {($slot/[dict get $options bucketSeconds])*[dict get $options bucketSeconds]}] if {![dict exists $counts $key]} {dict set counts $key 0} } dict set result buckets $counts set unread [expr {max(0,[dict get $state observedSize]-[dict get $state cursor])}] set historyUnread [expr {[dict get $state historyEnd]-[dict get $state historyCursor]}] set historyIncomplete [expr {$historyUnread>0 && [dict get $options from]<[dict get $state recentFrom]}] set seekPending [expr {[dict get $state seek] ne ""}] dict set result coverage [dict create updatedAt [dict get $state updatedAt] firstObserved [dict get $state firstObserved] latestObserved [dict get $state latestObserved] sourceBytes [dict get $state observedSize] unreadBytes $unread scanIncomplete [expr {$seekPending || $historyIncomplete || ($unread>0 && ![dict get $state pending])}] historyUnreadBytes $historyUnread totalUnreadBytes [expr {$unread+$historyUnread}] processedBytes [expr {max(0,[dict get $state cursor]-$historyUnread)}] historyIncomplete $historyIncomplete seekPending $seekPending recentFrom [dict get $state recentFrom] initialOffset [dict get $state initialOffset] seekMarginSeconds 300 completionOrderAssumed [expr {$historyUnread>0}] pending [dict get $state pending] generation [dict get $state generation] rotationGaps [dict get $state rotationGaps] lastRotationAt [dict get $state lastRotationAt] malformed [dict get $state malformed] oversized [dict get $state oversized] retentionSeconds ${:retentionSeconds} minuteAligned true] set retainedFrom [:retentionCutoff] dict set result coverage retainedFrom $retainedFrom set first [dict get $state firstObserved] dict set result coverage windowComplete [expr {$first ne "" && $first <= [dict get $options from] && [dict get $options from] >= $retainedFrom && $unread == 0 && !$historyIncomplete && !$seekPending && ([dict get $state rotationGaps] == 0 || [dict get $options from] >= (([dict get $state lastRotationAt]+59)/60)*60) && [dict get $state malformed] == 0 && [dict get $state oversized] == 0}] set reasons {} if {$first eq ""} {lappend reasons no-observations} elseif {$first>[dict get $options from]} {lappend reasons unobserved-start} if {[dict get $options from]<$retainedFrom} {lappend reasons retention} if {$seekPending} {lappend reasons seeking} if {$historyIncomplete} {lappend reasons historical-backfill} if {$unread>0} {lappend reasons [expr {[dict get $state pending] ? "pending-line" : "live-tail"}]} if {[dict get $state rotationGaps]>0 && [dict get $options from]<(([dict get $state lastRotationAt]+59)/60)*60} {lappend reasons rotation} if {[dict get $state malformed]>0} {lappend reasons cumulative-parse-errors} if {[dict get $state oversized]>0} {lappend reasons cumulative-oversized-records} set completion [dict get $state latestCompletionObserved] set age {} if {$completion ne ""} {set age [expr {max(0,[dict get $state updatedAt]-$completion)}]} dict set result coverage requestedFrom $requestedFrom dict set result coverage requestedTo $requestedTo dict set result coverage roundedTailSeconds [expr {[dict get $options to]-$requestedTo}] dict set result coverage latestCompletionObserved $completion dict set result coverage completionObservationAgeSeconds $age dict set result coverage sourceCaughtUp [expr {$unread==0 || [dict get $state pending]}] dict set result coverage liveTailIncomplete [expr {$unread>0 && ![dict get $state pending]}] dict set result coverage historicalBackfillPending [expr {$historyUnread>0}] dict set result coverage parseErrorScopeUnknown [expr {[dict get $state malformed]>0}] dict set result coverage oversizedScopeUnknown [expr {[dict get $state oversized]>0}] dict set result coverage lastMalformedAt [dict get $state lastMalformedAt] dict set result coverage incompleteReasons $reasons dict set result pathCountsComplete [expr {[dict get $result excludedPaths]==0}] return $result } :public method refresh {} { # Scheduler callback performs bounded work in a separate thread. if {[catch {:ingest {} true} message]} {ns_log warning "nsmcp: access ingestion failed: $message"} }XQL Not present: Generic, PostgreSQL, Oracle
![[i]](/resources/acs-subsite/ZoomIn16.gif)