nsmcp::AccessLogIngestor method ingestLocked (protected)
<instance of nsmcp::AccessLogIngestor> ingestLocked from \ allowHistory
Defined in /usr/local/ns/tcl/nsmcp/lib/access-analytics.tcl
- Parameters:
- from (required)
- allowHistory (required)
- Testcases:
- No testcase defined.
Source code: 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]]XQL Not present: Generic, PostgreSQL, Oracle
![[i]](/resources/acs-subsite/ZoomIn16.gif)