1#![allow(unsafe_code)]
6
7use std::collections::{HashMap, VecDeque};
55use std::sync::Arc;
56use std::time::SystemTime;
57
58use std::sync::Mutex;
59
60#[derive(Debug, Clone, PartialEq, Eq, Hash)]
67pub struct SessionKey {
68 pub provider: Arc<str>,
69 pub index: u32,
70}
71
72impl SessionKey {
73 pub fn new(provider: impl Into<Arc<str>>, index: u32) -> Self {
77 Self {
78 provider: provider.into(),
79 index,
80 }
81 }
82}
83
84pub 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
118fn one_line(s: &str) -> String {
121 s.replace(['\n', '\r', '\t'], " ")
122}
123
124pub type AiLogEventPublisher = Arc<dyn Fn(AiLogPushed) + Send + Sync>;
130
131pub 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
143fn lock<'a, T>(m: &'a Mutex<T>) -> std::sync::MutexGuard<'a, T> {
148 m.lock().expect("AiLogger mutex poisoned")
149}
150
151#[derive(Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord, Hash)]
156pub enum AiLogLevel {
157 Trace,
160 Debug,
163 Info,
166 Warn,
168 Error,
170}
171
172impl AiLogLevel {
173 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 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 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#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
211pub enum AiLogSource {
212 Client,
215 AgentText,
217 Reasoning,
219 ToolCall,
221 Lifecycle,
224}
225
226impl AiLogSource {
227 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#[derive(Debug, Clone)]
242pub struct AiLogRecord {
243 pub timestamp: SystemTime,
247 pub session: Option<SessionKey>,
249 pub level: AiLogLevel,
250 pub source: AiLogSource,
251 pub message: String,
252}
253
254#[derive(Debug)]
257pub struct LogRing {
258 buf: VecDeque<AiLogRecord>,
259 capacity: usize,
260}
261
262impl LogRing {
263 pub fn new(capacity: usize) -> Self {
266 Self {
267 buf: VecDeque::with_capacity(capacity.min(1024)),
268 capacity,
269 }
270 }
271
272 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 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 pub fn snapshot(&self) -> Vec<AiLogRecord> {
296 self.buf.iter().cloned().collect()
297 }
298
299 pub fn clear(&mut self) {
301 self.buf.clear();
302 }
303
304 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#[derive(Debug, Clone)]
317pub struct AiLogPushed {
318 pub timestamp: SystemTime,
322 pub session: Option<SessionKey>,
324 pub level: String,
327 pub source: String,
330 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#[derive(Clone)]
350pub struct AiLogger {
351 state: Arc<LoggerState>,
352}
353
354struct LoggerState {
355 global: Mutex<LogRing>,
357 per_session: Mutex<HashMap<SessionKey, LogRing>>,
360 event_publisher: Mutex<Option<AiLogEventPublisher>>,
364 default_capacity: Mutex<usize>,
366 default_min_level: Mutex<AiLogLevel>,
368 session_levels: Mutex<HashMap<SessionKey, AiLogLevel>>,
371}
372
373impl AiLogger {
374 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 pub fn set_event_publisher(&self, publisher: AiLogEventPublisher) {
397 *lock(&self.state.event_publisher) = Some(publisher);
398 }
399
400 pub fn with_defaults() -> Self {
402 Self::new(AiLogLevel::Info, 10_000)
403 }
404
405 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 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 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 let publisher = lock(&self.state.event_publisher).clone();
505 if let Some(p) = publisher {
506 p(publish_payload);
507 }
508 }
509
510 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 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 pub fn set_default_level(&self, level: AiLogLevel) {
537 *lock(&self.state.default_min_level) = level;
538 }
539
540 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 pub fn snapshot_global(&self) -> Vec<AiLogRecord> {
553 lock(&self.state.global).snapshot()
554 }
555
556 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 pub fn known_sessions(&self) -> Vec<SessionKey> {
568 lock(&self.state.per_session).keys().cloned().collect()
569 }
570
571 pub fn clear_global(&self) {
573 lock(&self.state.global).clear();
574 }
575
576 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 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 let unknown = key("zzz", 1);
699 assert!(logger.snapshot_session(&unknown).is_empty());
700 }
701
702 #[test]
703 fn different_sessions_stay_distinct() {
704 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 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 assert_eq!(logger.snapshot_session(&oc).len(), 1);
760 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 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 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 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 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 #[test]
866 fn format_ai_log_line_renders_the_records_own_timestamp() {
867 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 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 #[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}