From 849951be1d1a7ee9f9302006ccb187bf5b4e36f3 Mon Sep 17 00:00:00 2001 From: Paul Buetow Date: Wed, 22 Jul 2026 23:51:18 +0300 Subject: =?UTF-8?q?feat:=20DTail=20fork=20=E2=80=94=20server/client=20feat?= =?UTF-8?q?ure=20development?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Squashed development of the snonux/dtail fork's product code (internal/, cmd/) since diverging from mimecast/dtail. Major areas: - Read/output path: the former "turbo" channel-less path is now the single, default server-side read/output path for cat/grep/tail and MapReduce; the old channel-based path and its config/env toggles were removed. - MapReduce: single aggregate implementation (server + serverless) fed directly by a processor pipeline, with input-exhausted finalization via the shutdown coordinator; high-concurrency and data-race fixes. - Journal source reads (journal:unit.service) via journalctl, Linux-gated behind a journal-v1 capability. - Auth-key fast reconnect: in-memory per-user public-key cache with TTL/max-keys, registered over an authenticated session (AUTHKEY), checked before authorized_keys. - Interactive query reload (--interactive-query) with SESSION START/UPDATE generation boundaries and capability negotiation. - Client-side deadlines: --timeout / --shutdownAfter as context deadlines; follow shutdown handling. - Client logging: diagnostics-only daily log by default, opt-in payload tee via --log-payload. - Numerous correctness fixes (buffer-pool double-recycle races, EOF-sentinel leaks, glob-expansion cap, TOCTOU in CSV parsing) with accompanying unit tests. Co-Authored-By: Claude Opus 4.8 --- internal/mapr/client/session_state.go | 95 +++++++++++++++++++++++++++++++++++ 1 file changed, 95 insertions(+) create mode 100644 internal/mapr/client/session_state.go (limited to 'internal/mapr/client/session_state.go') diff --git a/internal/mapr/client/session_state.go b/internal/mapr/client/session_state.go new file mode 100644 index 0000000..1983644 --- /dev/null +++ b/internal/mapr/client/session_state.go @@ -0,0 +1,95 @@ +package client + +import ( + "fmt" + "sync" + + "github.com/mimecast/dtail/internal/mapr" +) + +// SessionSnapshot captures the current client-side mapreduce session state. +type SessionSnapshot struct { + Generation uint64 + Query *mapr.Query + GlobalGroup *mapr.GlobalGroupSet + LastResult string +} + +// SessionState keeps the mutable mapreduce query state shared by the client +// reporter and per-server handlers. +type SessionState struct { + mu sync.RWMutex + generation uint64 + query *mapr.Query + global *mapr.GlobalGroupSet + lastResult string + changedCh chan struct{} +} + +// NewSessionState returns a new shared mapreduce session state. +func NewSessionState(query *mapr.Query) *SessionState { + return &SessionState{ + query: query, + global: mapr.NewGlobalGroupSet(), + changedCh: make(chan struct{}, 1), + } +} + +// Snapshot returns a point-in-time copy of the shared mapreduce state. +func (s *SessionState) Snapshot() SessionSnapshot { + s.mu.RLock() + defer s.mu.RUnlock() + + return SessionSnapshot{ + Generation: s.generation, + Query: s.query, + GlobalGroup: s.global, + LastResult: s.lastResult, + } +} + +// Changes returns a channel that is signaled whenever a new generation is committed. +func (s *SessionState) Changes() <-chan struct{} { + return s.changedCh +} + +// CommitQuery resets the shared aggregation state for a newly accepted query generation. +func (s *SessionState) CommitQuery(rawQuery string, generation uint64) (*mapr.Query, error) { + query, err := mapr.NewQuery(rawQuery) + if err != nil { + return nil, fmt.Errorf("parse session query: %w", err) + } + + s.mu.Lock() + s.generation = generation + s.query = query + s.global = mapr.NewGlobalGroupSet() + s.lastResult = "" + s.mu.Unlock() + + s.notifyChange() + return query, nil +} + +// CommitRenderedResult stores the last rendered result for the active generation. +func (s *SessionState) CommitRenderedResult(generation uint64, result string) (changed bool, ok bool) { + s.mu.Lock() + defer s.mu.Unlock() + + if s.generation != generation { + return false, false + } + if s.lastResult == result { + return false, true + } + + s.lastResult = result + return true, true +} + +func (s *SessionState) notifyChange() { + select { + case s.changedCh <- struct{}{}: + default: + } +} -- cgit v1.2.3