Class ::nsmcp::AccessLogIngestor (public)

 ::nx::Class ::nsmcp::AccessLogIngestor[i]

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