Class ::nsmcp::FileLogSource (public)

 ::nx::Class ::nsmcp::FileLogSource[i]

Defined in /usr/local/ns/tcl/nsmcp/lib/log-source.tcl

Testcases:
No testcase defined.
Source code:
    :property path:required
    :property store:required
    :property {recordParser ::nsmcp::logRecordParser}
    :property {maxReadBytes:integer 1048576}
    :property {maxEvents:integer 1000}
    :property {maxMessageBytes:integer 65536}
    :property {mutex {}}
    :variable cacheKey

    :method init {} {
        if {${:maxReadBytes} < 1 || ${:maxEvents} < 1 || ${:maxMessageBytes} < 1} {
            error "log source limits must be positive"
        }
        set :path [file normalize ${:path}]
        if {${:mutex} eq ""} {set :mutex [ns_mutex create]}
        set :cacheKey [ns_crypto::md string -digest sha256 --  [list [ns_info server] [self] ${:path} ${:maxEvents} ${:maxMessageBytes} ${:recordParser}]]
    }

    :method emptyState {identity generation prefix} {
        return [dict create identity $identity generation $generation prefix $prefix  cursor 0 fragment {} fragmentOffset 0 fragmentTruncated false  pending {} events {}]
    }

    :method appendText {event text} {
        set message [dict get $event message]
        # Limit retained UTF-8 bytes; a clipped codepoint is replacement-decoded.
        set remaining [expr {max(0, ${:maxMessageBytes} - [string length [encoding convertto utf-8 $message]])}]
        set bytes [encoding convertto utf-8 $text]
        if {[string length $bytes] > $remaining} {
            set bytes [string range $bytes 0 [expr {$remaining-1}]]
            dict set event truncated true
        }
        set text [encoding convertfrom utf-8 $bytes]
        while {[string length [encoding convertto utf-8 $text]] > $remaining} {
            set text [string range $text 0 end-1]
        }
        append message $text
        dict set event message $message
        return $event
    }

    :method joinLine {state line start end clipped} {
        # Same joiner convention as nsstats _ns_stats.log.logfile. Remove ANSI
        # color escapes before checking the first character for ':'.
        if {[string first \x1b $line] >= 0} {set line [${:recordParser} stripAnsi $line]}
        set pending [dict get $state pending]
        if {[string index $line 0] eq ":" && $pending ne ""} {
            set pending [:appendText $pending \n$line]
        } else {
            if {$pending ne ""} {
                dict set pending complete true
                dict lappend state events $pending
                dict set state events [lrange [dict get $state events] end-[expr {${:maxEvents}-1}] end]
            }
            set pending [dict create offset $start endOffset $end message {}  complete false truncated false orphanContinuation [expr {[string index $line 0] eq ":"}]]
            set pending [:appendText $pending $line]
            set pending [dict merge $pending [${:recordParser} metadata $line]]
        }
        dict set pending endOffset $end
        if {$clipped} {dict set pending truncated true}
        dict set state pending $pending
        return $state
    }

    :method consume {state bytes} {
        # Read channels are binary, so offsets and cursor are byte offsets even
        # with multibyte UTF-8. Unfinished physical lines survive between reads.
        set parts [split $bytes \n]
        set last [expr {[llength $parts]-1}]
        set position [dict get $state cursor]
        set fragment [dict get $state fragment]
        set start [dict get $state fragmentOffset]
        set clipped [dict get $state fragmentTruncated]
        for {set i 0} {$i <= $last} {incr i} {
            set part [lindex $parts $i]
            set room [expr {${:maxMessageBytes} - [string length $fragment]}]
            if {[string length $part] > $room} {set clipped true}
            append fragment [string range $part 0 [expr {$room-1}]]
            incr position [string length $part]
            if {$i < $last} {
                incr position
                set line [encoding convertfrom utf-8 [string trimright $fragment \r]]
                set state [:joinLine $state $line $start $position $clipped]
                set fragment {}; set clipped false; set start $position
            }
        }
        dict set state cursor $position
        dict set state fragment $fragment
        dict set state fragmentOffset $start
        dict set state fragmentTruncated $clipped
        return $state
    }

    :public method read {} {
        ns_mutex lock ${:mutex}
        try {
            file stat ${:path} before
            set channel [open ${:path} rb]
            try {
                fconfigure $channel -translation binary -encoding binary
                file stat ${:path} opened
                if {[list $before(dev) $before(ino)] ne [list $opened(dev) $opened(ino)]} {
                    error "log file changed while opening; retry"
                }
                set identity [list $before(dev) $before(ino)]
                set reset false
                set generation 1
                if {[nsv_exists ${:store} ${:cacheKey}] && [nsv_get ${:store} ${:cacheKey} state]} {
                    set generation [dict get $state generation]
                    set prefix [dict get $state prefix]
                    # A fixed prefix also detects common copytruncate-and-regrow
                    # cases that have already grown beyond the saved cursor.
                    set prefixNow [read $channel [string length $prefix]]
                    if {$identity ne [dict get $state identity] ||
                        $before(size) < [dict get $state cursor] || $prefixNow ne $prefix} {
                        set reset true
                        incr generation
                    }
                } else {
                    set reset true
                }
                if {$reset} {
                    seek $channel 0
                    set prefix [read $channel [expr {min(256, $before(size))}]]
                    set state [:emptyState $identity $generation $prefix]
                }
                if {[dict get $state prefix] eq "" && $before(size) > 0} {
                    seek $channel 0
                    dict set state prefix [read $channel [expr {min(256, $before(size))}]]
                }
                seek $channel [dict get $state cursor]
                # Read only bytes visible at the snapshot, within the work budget.
                set count [expr {min(${:maxReadBytes}, max(0, $before(size)-[dict get $state cursor]))}]
                set state [:consume $state [read $channel $count]]
                file stat ${:path} after
                seek $channel 0
                set prefix [dict get $state prefix]
                set prefixAfter [read $channel [string length $prefix]]
                if {[list $after(dev) $after(ino)] ne $identity ||
                    $after(size) < $before(size) || $prefixAfter ne $prefix} {
                    # Do not publish a mixed snapshot if rotation races this read.
                    error "log file changed during read; retry"
                }
                nsv_set ${:store} ${:cacheKey} $state
                set events [dict get $state events]
                if {[dict get $state pending] ne ""} {lappend events [dict get $state pending]}
                set events [lrange $events end-[expr {${:maxEvents}-1}] end]
                return [dict create generation $generation reset $reset  cursor [dict get $state cursor]  backlogBytes [expr {max(0, $before(size)-[dict get $state cursor])}]  partialLineBytes [expr {[dict get $state cursor]-[dict get $state fragmentOffset]}]  events $events]
            } finally {close $channel}
        } finally {ns_mutex unlock ${:mutex}}
    }
XQL Not present:
Generic, PostgreSQL, Oracle
[ hide source ] | [ make this the default ]
Show another procedure: