nsmcp::AccessLogIngestor method ingestLocked (protected)

 <instance of nsmcp::AccessLogIngestor[i]> 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
[ hide source ] | [ make this the default ]
Show another procedure: