Class ::nsmcp::AccessAnalytics (public)
::nx::Class ::nsmcp::AccessAnalytics
Defined in /usr/local/ns/tcl/nsmcp/lib/access-analytics.tcl
- Testcases:
- No testcase defined.
Source code: :property database:required :property {maxRows:integer 2000000} :property {evictionLimit:integer 64} :method init {} { if {${:maxRows}<1 || ${:evictionLimit}<1 || ${:evictionLimit}>100} {error "positive maxRows and evictionLimit between 1 and 100 required"} if {![nsf::object::exists ${:database}] || ![${:database} info has type ::nsmcp::SQLiteDatabase]} {error "AccessAnalytics requires a SQLiteDatabase object"} next set db [${:database} connect] try { # Migrate the first SQLite draft without retaining duplicate tables. set old [$db scalar {SELECT count(*) FROM sqlite_master WHERE type='table' AND name='metadata'}] set current [$db scalar {SELECT count(*) FROM sqlite_master WHERE type='table' AND name='nsmcp_access_metadata'}] if {$old && !$current && [$db scalar {SELECT count(*) FROM sqlite_master WHERE type='table' AND name IN ('metadata','requests','evictions','diagnostics')}] == 4 && [$db scalar {SELECT value FROM metadata WHERE key='schemaVersion'}] eq "1"} { set source [$db scalar {SELECT value FROM metadata WHERE key='source'}] if {$source ne [list [ns_info server] ${:path}]} {error "draft SQLite cache belongs to another server or log source"} $db transaction { $db execute { ALTER TABLE metadata RENAME TO nsmcp_access_metadata; ALTER TABLE requests RENAME TO nsmcp_access_requests; ALTER TABLE evictions RENAME TO nsmcp_access_evictions; DROP INDEX IF EXISTS requests_time; DROP INDEX IF EXISTS requests_path_time; DROP TABLE IF EXISTS diagnostics; } } } $db execute { CREATE TABLE IF NOT EXISTS nsmcp_access_metadata (key TEXT PRIMARY KEY, value TEXT NOT NULL); CREATE TABLE IF NOT EXISTS nsmcp_access_requests ( generation INTEGER NOT NULL, offset INTEGER NOT NULL, start REAL NOT NULL, completed TEXT NOT NULL, request_id TEXT NOT NULL, path TEXT NOT NULL, pool TEXT NOT NULL, method TEXT NOT NULL, status INTEGER NOT NULL, bytes INTEGER NOT NULL, fallback INTEGER NOT NULL, excluded INTEGER NOT NULL, accept REAL, queue REAL, filter REAL, run REAL, PRIMARY KEY(generation,offset)); CREATE INDEX IF NOT EXISTS nsmcp_access_requests_time ON nsmcp_access_requests(start); CREATE INDEX IF NOT EXISTS nsmcp_access_requests_path_time ON nsmcp_access_requests(path,start); CREATE TABLE IF NOT EXISTS nsmcp_access_evictions (id INTEGER PRIMARY KEY, dropped_at INTEGER NOT NULL, first_time REAL NOT NULL, last_time REAL NOT NULL, rows INTEGER NOT NULL, reason TEXT NOT NULL); } set version [$db scalar {SELECT value FROM nsmcp_access_metadata WHERE key='schemaVersion'}] if {$version ne "" && $version ni {1 2 3}} {error "unsupported SQLite cache schema"} :peerSchema $db set identity [list [ns_info server] ${:path}] set stored [$db scalar {SELECT value FROM nsmcp_access_metadata WHERE key='source'}] if {$stored ne "" && $stored ne $identity} {error "SQLite cache belongs to another server or log source"} $db execute {INSERT OR IGNORE INTO nsmcp_access_metadata VALUES ('source',$identity)} } finally {$db close} } :method peerSchema {db} { # Additive migration also works after definition reload on a live object. set found 0 $db execute {SELECT name FROM pragma_table_info('nsmcp_access_requests')} row { if {$row(name) eq "peer"} {set found 1} } if {!$found} { $db execute {ALTER TABLE nsmcp_access_requests ADD COLUMN peer TEXT NOT NULL DEFAULT ''} } $db execute { CREATE INDEX IF NOT EXISTS nsmcp_access_requests_peer_time ON nsmcp_access_requests(peer,start); INSERT OR REPLACE INTO nsmcp_access_metadata VALUES('schemaVersion','3'); } } :method state {} { set value [${:db} scalar {SELECT value FROM nsmcp_access_metadata WHERE key='checkpoint'}] if {$value ne "" && [dict exists $value anchorHex]} { dict set value anchor [binary decode hex [dict get $value anchorHex]] dict unset value anchorHex } if {$value eq ""} {return [dict merge [:recentState] [:initial]]} if {[dict get $value version]!=1} {error "unsupported SQLite checkpoint version"} return [dict merge [dict create latestCompletionObserved {} lastMalformedAt {}] [:recentState] $value] } :method ingestLocked {from allowHistory} { set :db [${:database} connect] set :batchRows 0 set :loggedRejects 0 try { ${:db} execute {BEGIN IMMEDIATE} :peerSchema ${:db} set :rowCount [${:db} scalar {SELECT value FROM nsmcp_access_metadata WHERE key='rowCount'}] if {${:rowCount} eq ""} {set :rowCount [${:db} scalar {SELECT count(*) FROM nsmcp_access_requests}]} :maintain set result [next] if {${:loggedRejects}>5} {ns_log warning "nsmcp: suppressed [expr {${:loggedRejects}-5}] additional malformed access-record samples from ${:path}"} ${:db} execute {COMMIT} return $result } on error {message options} { catch {${:db} execute {ROLLBACK}} return -options $options $message } finally {${:db} close; unset -nocomplain :db :batchRows :rowCount :loggedRejects} } :method rejected {line offset state} { incr :loggedRejects if {${:loggedRejects}>5} {return} if {![nsf::object::exists ::nsmcp::accessParseRedactor]} {::nsmcp::LogRedactor create ::nsmcp::accessParseRedactor} set sample [string range [encoding convertfrom utf-8 $line] 0 1023] set sample [::nsmcp::accessParseRedactor message $sample] set sample [string map [list "\n" {\n} "\r" {\r} "\x00" {\0}] $sample] set reason [${:parser} failureReason [encoding convertfrom utf-8 $line]] ns_log warning "nsmcp: malformed access record source=${:path} generation=[dict get $state generation] offset=$offset reason=$reason sample=[list $sample]" } :method accumulate {entry offset changedVar stateVar cutoff} { upvar 1 $stateVar state set start [dict get $entry startTime] if {$start < $cutoff*60 || $start > [clock seconds]+60} {return} set db ${:db}; set generation [dict get $state generation] set path [dict get $entry path]; set excluded [expr {[string length $path]>512}] if {$excluded} {set path {}} set peer [dict get $entry peerAddress] if {[string length $peer]>128 || ![ns_ip valid $peer]} {set peer {}} set completed [dict get $entry timestamp] set request_id [string range [dict get $entry requestIdentifier] 0 511] set pool [string range [dict get $entry pool] 0 127] set method [string range [dict get $entry method] 0 127] set status [dict get $entry status]; set bytes [dict get $entry responseBytes] set fallback [expr {[dict get $entry timingSource] eq "completion"}] unset -nocomplain accept queue filter run if {!$fallback} { foreach {column field} {accept acceptTime queue queueTime filter filterTime run runTime} {set $column [dict get $entry $field]} } # Proactive eviction reuses pages before reaching SQLite's hard limit. incr :batchRows set batchLimit [expr {min(1000,max(1,int([${:database} cget -maxBytes]*0.045/2048)))}] if {${:batchRows}%$batchLimit==0} {:maintain} $db execute {INSERT OR IGNORE INTO nsmcp_access_requests (generation,offset,start,completed,request_id,path,pool,method,status,bytes,fallback,excluded,accept,queue,filter,run,peer) VALUES($generation,$offset,$start,$completed,$request_id,$path,$pool,$method,$status,$bytes,$fallback,$excluded,$accept,$queue,$filter,$run,$peer)} incr :rowCount [$db changes] } :method drop {reason cutoff {limit 0}} { set db ${:db}; set now [clock seconds] if {$reason eq "age"} { set where {start < $cutoff} } else { set where {rowid IN (SELECT rowid FROM nsmcp_access_requests ORDER BY start LIMIT $limit)} } set count 0 $db execute "SELECT count(*) AS n,coalesce(min(start),0) AS first,coalesce(max(start),0) AS last FROM nsmcp_access_requests WHERE $where" row { set count $row(n); set first $row(first); set last $row(last) } if {$count==0} {return} $db execute "DELETE FROM nsmcp_access_requests WHERE $where" incr :rowCount -$count $db execute {INSERT INTO nsmcp_access_evictions(dropped_at,first_time,last_time,rows,reason) VALUES($now,$first,$last,$count,$reason)} set old [$db scalar {SELECT value FROM nsmcp_access_metadata WHERE key='droppedThrough'}] set through [expr {$old eq "" ? $last : max(double($old),$last)}] $db execute {INSERT OR REPLACE INTO nsmcp_access_metadata VALUES('droppedThrough',$through); INSERT OR REPLACE INTO nsmcp_access_metadata VALUES('lastEvictionAt',$now)} set limit ${:evictionLimit} $db execute {DELETE FROM nsmcp_access_evictions WHERE id NOT IN (SELECT id FROM nsmcp_access_evictions ORDER BY id DESC LIMIT $limit)} } :method maintain {} { set db ${:db}; set cutoff [expr {([clock seconds]/60-${:retentionMinutes}+1)*60}] :drop age $cutoff set limit ${:maxRows} set count ${:rowCount} if {$count>$limit} {:drop rows 0 [expr {$count-$limit}]} set pageSize [$db scalar {PRAGMA page_size}] set cap [expr {int([${:database} cget -maxBytes]*0.45)}] # Leave room for the next batch, indexes and eviction history. while {([$db scalar {PRAGMA page_count}]-[$db scalar {PRAGMA freelist_count}])*$pageSize > $cap*0.8} { if {${:rowCount}==0} {error "SQLite cache exceeds budget and no analytics rows remain to evict"} :drop size 0 5000 } } :method publish {state changed cutoff} { :maintain dict set state anchorHex [binary encode hex [dict get $state anchor]] dict unset state anchor set count ${:rowCount} ${:db} execute {INSERT OR REPLACE INTO nsmcp_access_metadata VALUES('checkpoint',$state); INSERT OR REPLACE INTO nsmcp_access_metadata VALUES('rowCount',$count)} } :public method validateWindow {options} { set extra {} foreach key {path status pool peer windowSeconds includePeerAddresses} { if {![dict exists $options $key]} continue set value [dict get $options $key] if {$key eq "windowSeconds"} { if {![string is entier -strict $value] || $value<1 || $value>3600} {error "windowSeconds must be between 1 and 3600"} } elseif {$key eq "includePeerAddresses"} { if {![string is boolean -strict $value]} {error "includePeerAddresses must be boolean"} } elseif {$key eq "peer"} { if {![ns_ip valid $value]} {error "peer must be a valid IP address"} } elseif {$key eq "status"} { if {![string is entier -strict $value] || $value<100 || $value>599} {error "status must be between 100 and 599"} } elseif {$value eq "" || [string length $value]>512 || [string first \x00 $value]>=0} {error "invalid exact path or pool filter"} dict set extra $key $value; dict unset options $key } return [dict merge [next [list $options]] $extra] } :public method summary {options} { # Use the scanner coverage and a single snapshot for all aggregates. # Hold a single SQLite snapshot for coverage and all aggregate queries. ns_mutex lock ${:mutex} try { set :db [${:database} connect] set :batchRows 0 ${:db} execute {BEGIN} set options [:validateWindow $options] :peerSchema ${:db} set plain $options foreach key {path status pool peer windowSeconds includePeerAddresses} {dict unset plain $key} # Coverage and counts share the same read transaction. set result [:summaryLocked $plain] set filters {}; foreach key {path status pool peer} {if {[dict exists $options $key]} {dict set filters $key [dict get $options $key]}} dict set result filters $filters set db ${:db}; set from [dict get $result from]; set to [dict get $result to] set where {start >= $from AND start < $to} foreach key {path status pool peer} { if {[dict exists $options $key]} {set $key [dict get $options $key]; append where " AND $key = \$$key"} } $db execute "SELECT count(*) AS requestCount,coalesce(sum(bytes),0) AS responseBytes, coalesce(sum(fallback),0) AS fallbackCount,coalesce(sum(excluded),0) AS excludedPaths, count(run) AS timingCount,coalesce(sum(accept),0) AS acceptTotal, coalesce(sum(queue),0) AS queueTotal,coalesce(sum(filter),0) AS filterTotal, coalesce(sum(run),0) AS runTotal,coalesce(max(run),0) AS runMax FROM nsmcp_access_requests WHERE $where" row { foreach key $row(*) {dict set result $key $row($key)} } foreach {dimension column} {pools pool methods method statuses status paths path} { set counts {} # Bounded output, but count all matching rows in totals above. $db execute "SELECT $column AS label,count(*) AS n FROM nsmcp_access_requests WHERE $where GROUP BY $column ORDER BY n DESC LIMIT 1000" row { if {$dimension eq "paths" && $row(label) eq ""} continue dict set counts $row(label) $row(n) } dict set result $dimension $counts } set tracked 0 dict for {label count} [dict get $result paths] {incr tracked $count} dict set result excludedPaths [expr {[dict get $result requestCount]-$tracked}] set statusTimings {} $db execute "SELECT status,count(*) AS n,count(run) AS timed,coalesce(sum(filter),0) AS filters, coalesce(sum(run),0) AS runs,coalesce(max(run),0) AS maximum FROM nsmcp_access_requests WHERE $where GROUP BY status" row { lappend statusTimings [dict create status $row(status) requestCount $row(n) timingCount $row(timed) filterTotal $row(filters) runTotal $row(runs) runMax $row(maximum)] } dict set result statusTimings $statusTimings set unknownPeers [$db scalar "SELECT count(*) FROM nsmcp_access_requests WHERE $where AND peer = ''"] dict set result coverage unknownPeerCount $unknownPeers dict set result coverage peerCountsComplete [expr {$unknownPeers==0 && [dict get $result coverage windowComplete]}] if {[dict exists $options windowSeconds]} { # RANGE counts all equal timestamps together and uses an inclusive # [end-windowSeconds,end] range clipped to the selected period. # SQLite performs the rolling calculation rather than Tcl fetching rows. set seconds [dict get $options windowSeconds] set burst [dict create windowSeconds $seconds requestCount 0 from {} to {}] $db execute "SELECT peer,start,n FROM ( SELECT peer,start,count(*) OVER ( PARTITION BY peer ORDER BY start RANGE BETWEEN \$seconds PRECEDING AND CURRENT ROW) AS n FROM nsmcp_access_requests WHERE $where AND peer != '' ) ORDER BY n DESC,start,peer LIMIT 1" row { dict set burst requestCount $row(n) dict set burst from [expr {max(double($from),$row(start)-$seconds)}] dict set burst to $row(start) if {[dict exists $options includePeerAddresses] && [dict get $options includePeerAddresses]} { dict set burst peerAddress $row(peer) } } dict set result peerBurst $burst } set size [dict get $result bucketSeconds]; set counts [dict get $result buckets] $db execute "SELECT cast(start/\$size AS INTEGER)*\$size AS slot,count(*) AS n FROM nsmcp_access_requests WHERE $where GROUP BY slot" row { dict set counts $row(slot) $row(n) } dict set result buckets $counts dict set result pathCountsComplete [expr {[dict get $result excludedPaths]==0}] set retained [$db scalar {SELECT coalesce(min(start),'') FROM nsmcp_access_requests}] set through [$db scalar {SELECT value FROM nsmcp_access_metadata WHERE key='droppedThrough'}] set evictionAt [$db scalar {SELECT value FROM nsmcp_access_metadata WHERE key='lastEvictionAt'}] dict set result coverage source ${:path} dict set result coverage retainedRequestFrom $retained dict set result coverage droppedThrough $through dict set result coverage lastEvictionAt $evictionAt set evicted [expr {$through ne "" && $from <= $through}] dict set result coverage evictionOverlap $evicted if {$evicted} {dict set result coverage windowComplete false; set reasons [dict get $result coverage incompleteReasons]; lappend reasons eviction; dict set result coverage incompleteReasons $reasons} dict set result coverage peerCountsComplete [expr {$unknownPeers==0 && [dict get $result coverage windowComplete]}] dict set result coverage databaseBytes [file size [${:database} path]] dict set result coverage storageBudgetBytes [${:database} cget -maxBytes] set events {} $db execute {SELECT * FROM nsmcp_access_evictions ORDER BY id DESC} row { lappend events [dict create droppedAt $row(dropped_at) from $row(first_time) to $row(last_time) rowCount $row(rows) reason $row(reason)] } dict set result evictions $events ${:db} execute {COMMIT} return $result } finally {if {[info exists :db]} {${:db} close; unset :db}; unset -nocomplain :batchRows; ns_mutex unlock ${:mutex}} }XQL Not present: Generic, PostgreSQL, Oracle
![[i]](/resources/acs-subsite/ZoomIn16.gif)