Skip to main content

lattice_agent/log/
ai_log.rs

1// `linkme`'s distributed-slice expansion (behind
2// `register_event!`) uses `#[link_section]` declarations, which
3// the workspace's `unsafe_code = "deny"` lint flags. Same shape
4// `lattice-lsp::events` / `lattice-config::core_options` use.
5#![allow(unsafe_code)]
6
7//! AI agent logging facade -- the producer side of the
8//! `*ai:<provider>:<index>*` buffer views (AI-1b).
9//!
10//! ## Inspiration
11//!
12//! A direct port of `lattice-lsp`'s logging facade
13//! (`lattice_lsp::logging`), adapted per-process for AI agent sessions. LSP puts
14//! per-server stderr / protocol events in `*lsp:<server-id>*`
15//! buffers keyed by `(server_id, workspace)`; AI agents are the
16//! same shape one level simpler -- a `SessionKey { provider,
17//! index }` distinguishes a second `opencode` session (its own
18//! ring/buffer) from the first, exactly as two `rust-analyzer`
19//! instances against different workspaces stay distinct.
20//!
21//! ## Producer / consumer split
22//!
23//! This module is the **producer**: the ACP connection and
24//! session machinery emit [`AiLogRecord`]s through
25//! [`AiLogger::log`]. Records land in bounded [`LogRing`]s held
26//! inside the logger.
27//!
28//! The **consumer** -- buffer-backed log views in
29//! `lattice-ui-tui` (AI-1b, later slice) -- snapshots a ring on
30//! demand via [`AiLogger::snapshot_global`] /
31//! [`AiLogger::snapshot_session`] and renders the records.
32//! Auto-scroll-to-tail and live tail-follow are buffer-side
33//! concerns; the producer just appends.
34//!
35//! ## Tracing crate fan-out
36//!
37//! Every [`AiLogger::log`] call also emits a `tracing::*` event
38//! at the matching level (target `"ai"`) so users who prefer
39//! `RUST_LOG=ai=debug ./lattice` still see everything. The two
40//! paths are independent: the in-memory rings survive even when
41//! no `tracing` subscriber is installed, and `tracing` users see
42//! the same events whether or not the buffer views are open.
43//!
44//! ## Deviation from the LSP template
45//!
46//! `lattice_lsp::logging::LspLogger` also gates a per-instance
47//! JSON-RPC wire trace (`instance_trace` / `enable_trace` /
48//! `is_tracing` / `toggle_trace` / `disable_trace`). AI-1b has no
49//! per-message wire trace to toggle -- streamed agent text is a
50//! first-class `AiLogSource::AgentText`, not an opt-in trace --
51//! so that machinery is intentionally omitted here. Everything
52//! else mirrors the LSP logger verbatim.
53
54use std::collections::{HashMap, VecDeque};
55use std::sync::Arc;
56use std::time::SystemTime;
57
58use std::sync::Mutex;
59
60/// Per-process key for the per-session log ring -- the
61/// `(provider, index)` pair that distinguishes concurrent agent
62/// processes of the same provider. A second `opencode` session
63/// (`index = 2`) is a distinct ring/buffer from the first,
64/// analogous to LSP's `InstanceKey(server_id, workspace)`. Cheap
65/// to clone (one `Arc<str>` + a `Copy` integer).
66#[derive(Debug, Clone, PartialEq, Eq, Hash)]
67pub struct SessionKey {
68    pub provider: Arc<str>,
69    pub index: u32,
70}
71
72impl SessionKey {
73    /// Construct a session key from an owned/borrowed provider
74    /// name + process index. Accepts anything that converts into
75    /// `Arc<str>` so callers don't need to pre-Arc.
76    pub fn new(provider: impl Into<Arc<str>>, index: u32) -> Self {
77        Self {
78            provider: provider.into(),
79            index,
80        }
81    }
82}
83
84/// Format one log record line for append to a synthetic AI log
85/// buffer. Shape: `HH:MM:SS.mmm [<provider>:<index>] <level>
86/// <source>: <message>`. Trailing newline is the caller's
87/// responsibility (the drain batches many records into one
88/// buffer-append).
89///
90/// `timestamp` is the record's own emission time, not the format
91/// time: `:ai-log` seeds a buffer from history, and stamping at
92/// format time would render every historical line with the moment
93/// the buffer happened to be opened.
94///
95/// The clock is UTC (`(secs / 3600) % 24`), matching
96/// `lattice_lsp::logging::format_log_event_line`; rendering local
97/// time would need a timezone dependency neither logger carries.
98pub fn format_ai_log_line(
99    timestamp: SystemTime,
100    session: Option<&SessionKey>,
101    level: &str,
102    source: &str,
103    message: &str,
104) -> String {
105    let elapsed = timestamp.duration_since(std::time::UNIX_EPOCH).ok();
106    let secs = elapsed.map(|d| d.as_secs()).unwrap_or(0);
107    let ms = elapsed.map(|d| d.subsec_millis()).unwrap_or(0);
108    let hh = (secs / 3600) % 24;
109    let mm = (secs / 60) % 60;
110    let ss = secs % 60;
111    let prefix = session
112        .map(|key| format!("[{}:{}] ", key.provider, key.index))
113        .unwrap_or_default();
114    let msg = one_line(message);
115    format!("{hh:02}:{mm:02}:{ss:02}.{ms:03} {prefix}{level} {source:>6}: {msg}")
116}
117
118/// Collapse newlines / carriage returns / tabs into spaces so the
119/// formatted record fits on one buffer line.
120fn one_line(s: &str) -> String {
121    s.replace(['\n', '\r', '\t'], " ")
122}
123
124/// Closure invoked on every successful append. Wired by the App
125/// (or test harness) to publish [`AiLogPushed`] onto the runtime
126/// event bus, which lets log buffers refresh live as records
127/// arrive. Optional: `AiLogger::with_defaults()` starts with no
128/// publisher; `set_event_publisher` installs one.
129pub type AiLogEventPublisher = Arc<dyn Fn(AiLogPushed) + Send + Sync>;
130
131/// Compact severity tag for [`AiLogPushed`]. Mirrors
132/// [`AiLogSource::tag`]'s shape for the level discriminator.
133pub fn level_tag(l: AiLogLevel) -> &'static str {
134    match l {
135        AiLogLevel::Trace => "trace",
136        AiLogLevel::Debug => "debug",
137        AiLogLevel::Info => "info",
138        AiLogLevel::Warn => "warn",
139        AiLogLevel::Error => "error",
140    }
141}
142
143/// Internal helper: lock with consistent expect-message. Logger
144/// mutexes are never held across `.await` and never panic in the
145/// critical section, so `PoisonError` is unreachable in safe
146/// usage.
147fn lock<'a, T>(m: &'a Mutex<T>) -> std::sync::MutexGuard<'a, T> {
148    m.lock().expect("AiLogger mutex poisoned")
149}
150
151/// Severity levels. Ordered low-to-high so the derived
152/// `PartialOrd` matches "Error > Warn > ... > Trace". Per-session
153/// min level filters records BELOW it (i.e. `level < min` drops;
154/// `level >= min` keeps).
155#[derive(Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord, Hash)]
156pub enum AiLogLevel {
157    /// Finest-grained detail (raw wire payload fragments, mailbox
158    /// commands).
159    Trace,
160    /// Per-message, non-trace detail (tool-call argument dumps,
161    /// permission-request bookkeeping).
162    Debug,
163    /// Lifecycle milestones (session started, agent attached,
164    /// turn completed, ...).
165    Info,
166    /// Recoverable problems / unexpected protocol behaviour.
167    Warn,
168    /// Unrecoverable failures.
169    Error,
170}
171
172impl AiLogLevel {
173    /// Parse a string like `"info"` / `"debug"` (case-insensitive).
174    /// Used by config loaders.
175    pub fn parse(s: &str) -> Option<Self> {
176        match s.trim().to_ascii_lowercase().as_str() {
177            "error" => Some(AiLogLevel::Error),
178            "warn" | "warning" => Some(AiLogLevel::Warn),
179            "info" => Some(AiLogLevel::Info),
180            "debug" => Some(AiLogLevel::Debug),
181            "trace" => Some(AiLogLevel::Trace),
182            _ => None,
183        }
184    }
185
186    /// Iterate over all levels, low to high.
187    pub fn all() -> &'static [AiLogLevel] {
188        &[
189            AiLogLevel::Trace,
190            AiLogLevel::Debug,
191            AiLogLevel::Info,
192            AiLogLevel::Warn,
193            AiLogLevel::Error,
194        ]
195    }
196
197    /// Single-letter compact form for log rendering.
198    pub fn short(self) -> char {
199        match self {
200            AiLogLevel::Error => 'E',
201            AiLogLevel::Warn => 'W',
202            AiLogLevel::Info => 'I',
203            AiLogLevel::Debug => 'D',
204            AiLogLevel::Trace => 'T',
205        }
206    }
207}
208
209/// Where a log record originated.
210#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
211pub enum AiLogSource {
212    /// Lattice-side (session spawn, handshake, permission
213    /// requests, shutdown sequence, decode failures).
214    Client,
215    /// Streamed assistant text from `session/update`.
216    AgentText,
217    /// Streamed assistant reasoning / thinking content.
218    Reasoning,
219    /// Tool-call request / result content.
220    ToolCall,
221    /// Lifecycle milestones (session started, agent attached,
222    /// turn completed, process exited).
223    Lifecycle,
224}
225
226impl AiLogSource {
227    /// Compact tag used in log rendering.
228    pub fn tag(self) -> &'static str {
229        match self {
230            AiLogSource::Client => "client",
231            AiLogSource::AgentText => "agent",
232            AiLogSource::Reasoning => "reason",
233            AiLogSource::ToolCall => "tool",
234            AiLogSource::Lifecycle => "life",
235        }
236    }
237}
238
239/// One log entry. Cheap to clone (`Arc<str>` for the session's
240/// provider name; the message is a plain `String`).
241#[derive(Debug, Clone)]
242pub struct AiLogRecord {
243    /// Wall-clock time of emission. Used for buffer rendering;
244    /// not load-bearing for ordering (records preserve insertion
245    /// order in their ring).
246    pub timestamp: SystemTime,
247    /// `None` for subsystem-wide records (supervisor events).
248    pub session: Option<SessionKey>,
249    pub level: AiLogLevel,
250    pub source: AiLogSource,
251    pub message: String,
252}
253
254/// Bounded ring of log records. Append-only; oldest is evicted
255/// when capacity is reached.
256#[derive(Debug)]
257pub struct LogRing {
258    buf: VecDeque<AiLogRecord>,
259    capacity: usize,
260}
261
262impl LogRing {
263    /// Construct with capacity. `capacity` of 0 produces a ring
264    /// that drops everything (useful for "logging disabled").
265    pub fn new(capacity: usize) -> Self {
266        Self {
267            buf: VecDeque::with_capacity(capacity.min(1024)),
268            capacity,
269        }
270    }
271
272    /// Append a record, evicting the oldest if at capacity.
273    pub fn push(&mut self, record: AiLogRecord) {
274        if self.capacity == 0 {
275            return;
276        }
277        while self.buf.len() >= self.capacity {
278            self.buf.pop_front();
279        }
280        self.buf.push_back(record);
281    }
282
283    /// Number of records currently stored.
284    pub fn len(&self) -> usize {
285        self.buf.len()
286    }
287
288    pub fn is_empty(&self) -> bool {
289        self.buf.is_empty()
290    }
291
292    /// Snapshot every record. Cheap clone (each `AiLogRecord`'s
293    /// heavy field is the message `String`; we're optimising for
294    /// correctness here, not for log-view paint cost).
295    pub fn snapshot(&self) -> Vec<AiLogRecord> {
296        self.buf.iter().cloned().collect()
297    }
298
299    /// Drop everything. Used by `:ai-log clear`.
300    pub fn clear(&mut self) {
301        self.buf.clear();
302    }
303
304    /// Capacity setter. Used when the user adjusts
305    /// `ai.log_capacity` at runtime.
306    pub fn set_capacity(&mut self, capacity: usize) {
307        self.capacity = capacity;
308        while self.buf.len() > self.capacity {
309            self.buf.pop_front();
310        }
311    }
312}
313
314/// The typed event fired on every successful append. Analogue of
315/// `lattice_lsp::events::LspLogPushed`.
316#[derive(Debug, Clone)]
317pub struct AiLogPushed {
318    /// Wall-clock time of emission, carried so a live-tailing buffer
319    /// renders the same timestamp the record will show once it is
320    /// replayed from history by a later `:ai-log`.
321    pub timestamp: SystemTime,
322    /// `None` for subsystem-wide records; `Some(key)` per-session.
323    pub session: Option<SessionKey>,
324    /// Severity tag (`"trace"`, `"debug"`, `"info"`, `"warn"`,
325    /// `"error"`).
326    pub level: String,
327    /// Source tag (`"client"`, `"agent"`, `"reason"`, `"tool"`,
328    /// `"life"`).
329    pub source: String,
330    /// The record's message text.
331    pub message: String,
332}
333
334lattice_protocol::register_event!(
335    AiLogPushed,
336    "ai.log-pushed",
337    "Fired after an AI agent log record is appended to its ring; \
338     drives live refresh of the per-process *ai:<provider>:<index>* buffers.",
339    "lattice-ai",
340);
341
342/// Producer-side log facade. One per AI subsystem; each agent
343/// connection gets a clone (cheap -- internal state is
344/// `Arc<Mutex<...>>`).
345///
346/// Cheap-by-design: append is O(1) amortized, level gating is one
347/// HashMap lookup. Logging is **not** a hot path -- per-chunk cost
348/// is the `session/update` decode + apply, not the log emission.
349#[derive(Clone)]
350pub struct AiLogger {
351    state: Arc<LoggerState>,
352}
353
354struct LoggerState {
355    /// Subsystem-wide ring (records where `session` is `None`).
356    global: Mutex<LogRing>,
357    /// Per-session rings, keyed by `(provider, index)`. A second
358    /// `opencode` session is a different ring than the first.
359    per_session: Mutex<HashMap<SessionKey, LogRing>>,
360    /// Optional publisher fired on every successful append. Wired
361    /// by the App at boot to feed the runtime event bus. `None` ->
362    /// no events emitted (test paths, or pre-wire).
363    event_publisher: Mutex<Option<AiLogEventPublisher>>,
364    /// Default capacity for new per-session rings.
365    default_capacity: Mutex<usize>,
366    /// Default min level (records below are dropped).
367    default_min_level: Mutex<AiLogLevel>,
368    /// Per-session overrides for min level. Falls back to the
369    /// default when absent.
370    session_levels: Mutex<HashMap<SessionKey, AiLogLevel>>,
371}
372
373impl AiLogger {
374    /// Construct a logger with the given default min level and
375    /// default per-session ring capacity (e.g. `AiLogLevel::Info`,
376    /// 10_000). Conservative defaults: Info filters out debug/trace
377    /// noise; 10k records ≈ a few MB of memory at most.
378    pub fn new(default_min_level: AiLogLevel, default_capacity: usize) -> Self {
379        Self {
380            state: Arc::new(LoggerState {
381                global: Mutex::new(LogRing::new(default_capacity)),
382                per_session: Mutex::new(HashMap::new()),
383                default_capacity: Mutex::new(default_capacity),
384                default_min_level: Mutex::new(default_min_level),
385                session_levels: Mutex::new(HashMap::new()),
386                event_publisher: Mutex::new(None),
387            }),
388        }
389    }
390
391    /// Install / replace the event publisher. Subsequent `log`
392    /// calls fire the closure with [`AiLogPushed`] after the
393    /// record lands in its ring. The App wires this at boot so the
394    /// runtime event bus sees every append; subscribers (live log
395    /// views) drain the bus on tick.
396    pub fn set_event_publisher(&self, publisher: AiLogEventPublisher) {
397        *lock(&self.state.event_publisher) = Some(publisher);
398    }
399
400    /// Sensible defaults: Info level, 10k records / ring.
401    pub fn with_defaults() -> Self {
402        Self::new(AiLogLevel::Info, 10_000)
403    }
404
405    /// Append a record. Gated by per-session min level (or the
406    /// default). Records with `session = None` route to the
407    /// subsystem-wide ring (`*ai*`) and use the default min level.
408    pub fn log(
409        &self,
410        session: Option<&SessionKey>,
411        level: AiLogLevel,
412        source: AiLogSource,
413        message: impl Into<String>,
414    ) {
415        let min = self.effective_min_level(session);
416        if level < min {
417            return;
418        }
419
420        let message = message.into();
421
422        // tracing fan-out -- always fires, regardless of buffer
423        // views being open. RUST_LOG users see the same events.
424        let provider_disp = session.map(|key| key.provider.to_string());
425        let index_disp = session.map(|key| key.index);
426        match level {
427            AiLogLevel::Error => tracing::error!(
428                target: "ai",
429                provider = provider_disp.as_deref(),
430                index = index_disp,
431                source = source.tag(),
432                "{}",
433                message
434            ),
435            AiLogLevel::Warn => tracing::warn!(
436                target: "ai",
437                provider = provider_disp.as_deref(),
438                index = index_disp,
439                source = source.tag(),
440                "{}",
441                message
442            ),
443            AiLogLevel::Info => tracing::info!(
444                target: "ai",
445                provider = provider_disp.as_deref(),
446                index = index_disp,
447                source = source.tag(),
448                "{}",
449                message
450            ),
451            AiLogLevel::Debug => tracing::debug!(
452                target: "ai",
453                provider = provider_disp.as_deref(),
454                index = index_disp,
455                source = source.tag(),
456                "{}",
457                message
458            ),
459            AiLogLevel::Trace => tracing::trace!(
460                target: "ai",
461                provider = provider_disp.as_deref(),
462                index = index_disp,
463                source = source.tag(),
464                "{}",
465                message
466            ),
467        }
468
469        let record = AiLogRecord {
470            timestamp: SystemTime::now(),
471            session: session.cloned(),
472            level,
473            source,
474            message,
475        };
476
477        // Fan out to the runtime event bus before / after the ring
478        // push. Payload is a typed `AiLogPushed`; subscribers use
479        // `session` to route to the correct per-process buffer.
480        let publish_payload = AiLogPushed {
481            timestamp: record.timestamp,
482            session: record.session.clone(),
483            level: level_tag(record.level).to_string(),
484            source: record.source.tag().to_string(),
485            message: record.message.clone(),
486        };
487
488        match session {
489            None => {
490                lock(&self.state.global).push(record);
491            }
492            Some(key) => {
493                let cap = *lock(&self.state.default_capacity);
494                let mut per = lock(&self.state.per_session);
495                per.entry(key.clone())
496                    .or_insert_with(|| LogRing::new(cap))
497                    .push(record);
498            }
499        }
500
501        // Snapshot the publisher under the mutex, drop the lock,
502        // then call -- the bus's internal mutex is independent of
503        // ours and we never hold both at once.
504        let publisher = lock(&self.state.event_publisher).clone();
505        if let Some(p) = publisher {
506            p(publish_payload);
507        }
508    }
509
510    /// Resolve the min level for a session (or subsystem-wide).
511    fn effective_min_level(&self, session: Option<&SessionKey>) -> AiLogLevel {
512        if let Some(key) = session
513            && let Some(level) = lock(&self.state.session_levels).get(key).copied()
514        {
515            return level;
516        }
517        *lock(&self.state.default_min_level)
518    }
519
520    /// Set per-session min level. `None` removes the override and
521    /// reverts to the default.
522    pub fn set_session_level(&self, session: SessionKey, level: Option<AiLogLevel>) {
523        let mut guard = lock(&self.state.session_levels);
524        match level {
525            Some(l) => {
526                guard.insert(session, l);
527            }
528            None => {
529                guard.remove(&session);
530            }
531        }
532    }
533
534    /// Set the default min level (applies to subsystem-wide
535    /// records and to sessions without an override).
536    pub fn set_default_level(&self, level: AiLogLevel) {
537        *lock(&self.state.default_min_level) = level;
538    }
539
540    /// Set the default ring capacity. Existing rings are resized;
541    /// future per-session rings inherit the new value.
542    pub fn set_default_capacity(&self, capacity: usize) {
543        *lock(&self.state.default_capacity) = capacity;
544        lock(&self.state.global).set_capacity(capacity);
545        for (_, ring) in lock(&self.state.per_session).iter_mut() {
546            ring.set_capacity(capacity);
547        }
548    }
549
550    /// Snapshot the subsystem-wide ring. Used by the `*ai*` buffer
551    /// view to populate its body.
552    pub fn snapshot_global(&self) -> Vec<AiLogRecord> {
553        lock(&self.state.global).snapshot()
554    }
555
556    /// Snapshot a session's ring (empty if the session has never
557    /// logged anything).
558    pub fn snapshot_session(&self, session: &SessionKey) -> Vec<AiLogRecord> {
559        lock(&self.state.per_session)
560            .get(session)
561            .map(LogRing::snapshot)
562            .unwrap_or_default()
563    }
564
565    /// List every session with a per-session ring. Useful for
566    /// `:ai-status` and the `*ai*` buffer's "sessions" header.
567    pub fn known_sessions(&self) -> Vec<SessionKey> {
568        lock(&self.state.per_session).keys().cloned().collect()
569    }
570
571    /// Drop the subsystem-wide ring's contents.
572    pub fn clear_global(&self) {
573        lock(&self.state.global).clear();
574    }
575
576    /// Drop a session's ring contents.
577    pub fn clear_session(&self, session: &SessionKey) {
578        if let Some(ring) = lock(&self.state.per_session).get_mut(session) {
579            ring.clear();
580        }
581    }
582}
583
584impl Default for AiLogger {
585    fn default() -> Self {
586        Self::with_defaults()
587    }
588}
589
590impl std::fmt::Debug for AiLogger {
591    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
592        let global_len = lock(&self.state.global).len();
593        let n_sessions = lock(&self.state.per_session).len();
594        f.debug_struct("AiLogger")
595            .field("global_records", &global_len)
596            .field("session_count", &n_sessions)
597            .finish_non_exhaustive()
598    }
599}
600
601#[cfg(test)]
602mod tests {
603    use std::time::Duration;
604
605    use super::*;
606
607    fn key(provider: &str, index: u32) -> SessionKey {
608        SessionKey::new(Arc::<str>::from(provider), index)
609    }
610
611    #[test]
612    fn ai_log_level_parse_round_trips() {
613        assert_eq!(AiLogLevel::parse("error"), Some(AiLogLevel::Error));
614        assert_eq!(AiLogLevel::parse("WARN"), Some(AiLogLevel::Warn));
615        assert_eq!(AiLogLevel::parse("warning"), Some(AiLogLevel::Warn));
616        assert_eq!(AiLogLevel::parse("Info"), Some(AiLogLevel::Info));
617        assert_eq!(AiLogLevel::parse("debug"), Some(AiLogLevel::Debug));
618        assert_eq!(AiLogLevel::parse("trace"), Some(AiLogLevel::Trace));
619        assert_eq!(AiLogLevel::parse("nope"), None);
620    }
621
622    #[test]
623    fn level_ordering_matches_severity() {
624        // Error is the highest -- a min-level filter of Error means
625        // "only show errors".
626        assert!(AiLogLevel::Error > AiLogLevel::Warn);
627        assert!(AiLogLevel::Warn > AiLogLevel::Info);
628        assert!(AiLogLevel::Info > AiLogLevel::Debug);
629        assert!(AiLogLevel::Debug > AiLogLevel::Trace);
630    }
631
632    #[test]
633    fn ring_evicts_oldest_at_capacity() {
634        let mut ring = LogRing::new(3);
635        for i in 0..5 {
636            ring.push(AiLogRecord {
637                timestamp: SystemTime::now(),
638                session: None,
639                level: AiLogLevel::Info,
640                source: AiLogSource::Client,
641                message: format!("msg {i}"),
642            });
643        }
644        assert_eq!(ring.len(), 3);
645        let snap = ring.snapshot();
646        assert_eq!(snap[0].message, "msg 2");
647        assert_eq!(snap[1].message, "msg 3");
648        assert_eq!(snap[2].message, "msg 4");
649    }
650
651    #[test]
652    fn ring_zero_capacity_drops_everything() {
653        let mut ring = LogRing::new(0);
654        ring.push(AiLogRecord {
655            timestamp: SystemTime::now(),
656            session: None,
657            level: AiLogLevel::Info,
658            source: AiLogSource::Client,
659            message: "lost".into(),
660        });
661        assert_eq!(ring.len(), 0);
662    }
663
664    #[test]
665    fn log_routes_to_correct_ring() {
666        let logger = AiLogger::with_defaults();
667        let oc1 = key("opencode", 1);
668        let oc2 = key("claude-code", 1);
669
670        logger.log(None, AiLogLevel::Info, AiLogSource::Client, "subsys event");
671        logger.log(
672            Some(&oc1),
673            AiLogLevel::Info,
674            AiLogSource::AgentText,
675            "opencode evt",
676        );
677        logger.log(
678            Some(&oc2),
679            AiLogLevel::Warn,
680            AiLogSource::Lifecycle,
681            "claude-code evt",
682        );
683
684        let g = logger.snapshot_global();
685        assert_eq!(g.len(), 1);
686        assert_eq!(g[0].message, "subsys event");
687
688        let r = logger.snapshot_session(&oc1);
689        assert_eq!(r.len(), 1);
690        assert_eq!(r[0].message, "opencode evt");
691
692        let p = logger.snapshot_session(&oc2);
693        assert_eq!(p.len(), 1);
694        assert_eq!(p[0].message, "claude-code evt");
695
696        // Snapshotting an unknown session returns empty, not a
697        // panic.
698        let unknown = key("zzz", 1);
699        assert!(logger.snapshot_session(&unknown).is_empty());
700    }
701
702    #[test]
703    fn different_sessions_stay_distinct() {
704        // The per-process requirement: two `opencode` sessions
705        // (index 1 vs index 2) must NOT share a ring -- otherwise
706        // `*ai:opencode:1*` and `*ai:opencode:2*` would surface
707        // each other's records.
708        let logger = AiLogger::with_defaults();
709        let oc1 = key("opencode", 1);
710        let oc2 = key("opencode", 2);
711        logger.log(
712            Some(&oc1),
713            AiLogLevel::Info,
714            AiLogSource::Client,
715            "msg from session 1",
716        );
717        logger.log(
718            Some(&oc2),
719            AiLogLevel::Info,
720            AiLogSource::Client,
721            "msg from session 2",
722        );
723        let a = logger.snapshot_session(&oc1);
724        let b = logger.snapshot_session(&oc2);
725        assert_eq!(a.len(), 1, "session 1 has its own record");
726        assert_eq!(a[0].message, "msg from session 1");
727        assert_eq!(b.len(), 1, "session 2 has its own record");
728        assert_eq!(b[0].message, "msg from session 2");
729        // The record carries the session so renderers can verify
730        // which process produced it.
731        assert_eq!(a[0].session.as_ref(), Some(&oc1));
732        assert_eq!(b[0].session.as_ref(), Some(&oc2));
733    }
734
735    #[test]
736    fn log_below_min_level_is_dropped() {
737        let logger = AiLogger::new(AiLogLevel::Warn, 100);
738        let id = key("opencode", 1);
739        logger.log(Some(&id), AiLogLevel::Info, AiLogSource::Client, "below");
740        logger.log(Some(&id), AiLogLevel::Warn, AiLogSource::Client, "at");
741        logger.log(Some(&id), AiLogLevel::Error, AiLogSource::Client, "above");
742        let snap = logger.snapshot_session(&id);
743        assert_eq!(snap.len(), 2);
744        assert_eq!(snap[0].message, "at");
745        assert_eq!(snap[1].message, "above");
746    }
747
748    #[test]
749    fn per_session_level_overrides_default() {
750        let logger = AiLogger::new(AiLogLevel::Info, 100);
751        let oc = key("opencode", 1);
752        let cc = key("claude-code", 1);
753        logger.set_session_level(oc.clone(), Some(AiLogLevel::Debug));
754
755        logger.log(Some(&oc), AiLogLevel::Debug, AiLogSource::Client, "oc dbg");
756        logger.log(Some(&cc), AiLogLevel::Debug, AiLogSource::Client, "cc dbg");
757
758        // opencode override: Debug accepted.
759        assert_eq!(logger.snapshot_session(&oc).len(), 1);
760        // claude-code default: Debug below Info, dropped.
761        assert_eq!(logger.snapshot_session(&cc).len(), 0);
762    }
763
764    #[test]
765    fn known_sessions_lists_only_seen() {
766        let logger = AiLogger::with_defaults();
767        let oc = key("opencode", 1);
768        let cc = key("claude-code", 1);
769        logger.log(Some(&oc), AiLogLevel::Info, AiLogSource::Client, "x");
770        logger.log(Some(&cc), AiLogLevel::Info, AiLogSource::Client, "y");
771        let mut known: Vec<String> = logger
772            .known_sessions()
773            .into_iter()
774            .map(|k| k.provider.to_string())
775            .collect();
776        known.sort();
777        assert_eq!(
778            known,
779            vec!["claude-code".to_string(), "opencode".to_string()]
780        );
781    }
782
783    #[test]
784    fn clearing_drops_records() {
785        let logger = AiLogger::with_defaults();
786        let id = key("opencode", 1);
787        logger.log(None, AiLogLevel::Info, AiLogSource::Client, "a");
788        logger.log(Some(&id), AiLogLevel::Info, AiLogSource::Client, "b");
789        logger.clear_global();
790        assert!(logger.snapshot_global().is_empty());
791        // Per-session ring still has the entry until we clear it.
792        assert_eq!(logger.snapshot_session(&id).len(), 1);
793        logger.clear_session(&id);
794        assert!(logger.snapshot_session(&id).is_empty());
795    }
796
797    #[test]
798    fn set_default_capacity_resizes_existing_rings() {
799        let logger = AiLogger::new(AiLogLevel::Info, 100);
800        let id = key("opencode", 1);
801        for i in 0..50 {
802            logger.log(
803                Some(&id),
804                AiLogLevel::Info,
805                AiLogSource::Client,
806                format!("r{i}"),
807            );
808        }
809        assert_eq!(logger.snapshot_session(&id).len(), 50);
810        // Shrink to 10 -- the most recent 10 survive.
811        logger.set_default_capacity(10);
812        let snap = logger.snapshot_session(&id);
813        assert_eq!(snap.len(), 10);
814        assert_eq!(snap[0].message, "r40");
815        assert_eq!(snap[9].message, "r49");
816    }
817
818    #[test]
819    fn cheap_clone_shares_state() {
820        let a = AiLogger::with_defaults();
821        let b = a.clone();
822        a.log(None, AiLogLevel::Info, AiLogSource::Client, "from a");
823        // Clone sees the same ring.
824        assert_eq!(b.snapshot_global().len(), 1);
825    }
826
827    #[test]
828    fn format_ai_log_line_shape() {
829        let session = key("opencode", 2);
830        let line = format_ai_log_line(
831            SystemTime::now(),
832            Some(&session),
833            "info",
834            "agent",
835            "hello\nworld",
836        );
837        // HH:MM:SS.mmm [opencode:2] info  agent: hello world
838        assert!(
839            line.contains("[opencode:2]"),
840            "line should carry the session tag: {line}"
841        );
842        assert!(line.contains("info"), "line should carry the level: {line}");
843        assert!(
844            line.contains("agent"),
845            "line should carry the source: {line}"
846        );
847        assert!(
848            line.contains("hello world"),
849            "newline should collapse to a space: {line}"
850        );
851        assert!(!line.contains('\n'), "no embedded newline: {line}");
852
853        let global_line =
854            format_ai_log_line(SystemTime::now(), None, "warn", "client", "no session");
855        assert!(
856            !global_line.contains('['),
857            "subsystem-wide line has no session prefix: {global_line}"
858        );
859    }
860
861    /// The rendered clock must come from the record's own emission time,
862    /// not from the moment the line happens to be formatted. `:ai-log`
863    /// seeds a buffer by replaying history, so a format-time clock would
864    /// stamp every historical line with the buffer's open time.
865    #[test]
866    fn format_ai_log_line_renders_the_records_own_timestamp() {
867        // 01:01:01.500 UTC on the epoch day.
868        let stamp = std::time::UNIX_EPOCH + Duration::from_millis(3_661_500);
869        let line = format_ai_log_line(stamp, None, "info", "agent", "pong");
870        assert!(
871            line.starts_with("01:01:01.500 "),
872            "line must render the record's timestamp, got {line}"
873        );
874
875        // A different record time renders a different clock, even though
876        // both lines are formatted at the same instant.
877        let later = std::time::UNIX_EPOCH + Duration::from_millis(7_322_250);
878        let later_line = format_ai_log_line(later, None, "info", "agent", "pong");
879        assert!(later_line.starts_with("02:02:02.250 "), "got {later_line}");
880    }
881
882    /// The bus event carries the record's timestamp, so a live-tailing
883    /// buffer and a later history replay render the same clock for the
884    /// same record.
885    #[test]
886    fn published_event_carries_the_records_timestamp() {
887        let logger = AiLogger::with_defaults();
888        let seen: Arc<Mutex<Vec<AiLogPushed>>> = Arc::new(Mutex::new(Vec::new()));
889        let sink = seen.clone();
890        logger.set_event_publisher(Arc::new(move |event: AiLogPushed| {
891            sink.lock().expect("sink lock").push(event);
892        }));
893
894        let session = key("opencode", 1);
895        logger.log(
896            Some(&session),
897            AiLogLevel::Info,
898            AiLogSource::AgentText,
899            "hi",
900        );
901
902        let events = seen.lock().expect("sink lock");
903        let published = events.first().expect("one published event");
904        let record = logger
905            .snapshot_session(&session)
906            .first()
907            .cloned()
908            .expect("one ring record");
909        assert_eq!(
910            published.timestamp, record.timestamp,
911            "the published event and the ring record must agree on emission time"
912        );
913    }
914}