1use std::collections::VecDeque;
5
6use tokio::sync::watch;
7use zeph_common::SecurityEventCategory;
8
9pub use zeph_llm::{ClassifierMetricsSnapshot, TaskMetricsSnapshot};
10pub use zeph_memory::{CategoryScore, ProbeCategory, ProbeVerdict};
11
12#[derive(Debug, Clone)]
14pub struct SecurityEvent {
15 pub timestamp: u64,
17 pub category: SecurityEventCategory,
18 pub source: String,
20 pub detail: String,
22}
23
24impl SecurityEvent {
25 #[must_use]
26 pub fn new(
27 category: SecurityEventCategory,
28 source: impl Into<String>,
29 detail: impl Into<String>,
30 ) -> Self {
31 let source: String = source
33 .into()
34 .chars()
35 .filter(|c| !c.is_ascii_control())
36 .take(64)
37 .collect();
38 let detail = detail.into();
40 let detail = if detail.len() > 128 {
41 let end = detail.floor_char_boundary(127);
42 format!("{}…", &detail[..end])
43 } else {
44 detail
45 };
46 Self {
47 timestamp: std::time::SystemTime::now()
48 .duration_since(std::time::UNIX_EPOCH)
49 .unwrap_or_default()
50 .as_secs(),
51 category,
52 source,
53 detail,
54 }
55 }
56}
57
58pub const SECURITY_EVENT_CAP: usize = 100;
60
61#[derive(Debug, Clone)]
65pub struct TaskSnapshotRow {
66 pub id: u32,
67 pub title: String,
68 pub status: String,
70 pub agent: Option<String>,
71 pub duration_ms: u64,
72 pub error: Option<String>,
74}
75
76#[derive(Debug, Clone, Default)]
78pub struct TaskGraphSnapshot {
79 pub graph_id: String,
80 pub goal: String,
81 pub status: String,
83 pub tasks: Vec<TaskSnapshotRow>,
84 pub completed_at: Option<std::time::Instant>,
85}
86
87impl TaskGraphSnapshot {
88 #[must_use]
91 pub fn is_stale(&self) -> bool {
92 self.completed_at
93 .is_some_and(|t| t.elapsed().as_secs() > 30)
94 }
95}
96
97#[derive(Debug, Clone, Default)]
101pub struct OrchestrationMetrics {
102 pub plans_total: u64,
103 pub tasks_total: u64,
104 pub tasks_completed: u64,
105 pub tasks_failed: u64,
106 pub tasks_skipped: u64,
107}
108
109#[non_exhaustive]
110#[derive(Debug, Clone, PartialEq, Eq)]
112pub enum McpServerConnectionStatus {
113 Connected,
114 Failed,
115}
116
117#[derive(Debug, Clone)]
119pub struct McpServerStatus {
120 pub id: String,
121 pub status: McpServerConnectionStatus,
122 pub tool_count: usize,
124 pub error: String,
126}
127
128#[derive(Debug, Clone, Default)]
130pub struct SkillConfidence {
131 pub name: String,
132 pub posterior: f64,
133 pub total_uses: u32,
134}
135
136#[derive(Debug, Clone, Default)]
138pub struct SubAgentMetrics {
139 pub id: String,
140 pub name: String,
141 pub state: String,
143 pub turns_used: u32,
144 pub max_turns: u32,
145 pub background: bool,
146 pub elapsed_secs: u64,
147 pub permission_mode: String,
150 pub transcript_dir: Option<String>,
153}
154
155#[derive(Debug, Clone, Default)]
160pub struct TurnTimings {
161 pub prepare_context_ms: u64,
162 pub llm_chat_ms: u64,
163 pub tool_exec_ms: u64,
164 pub persist_message_ms: u64,
165}
166
167#[derive(Debug, Clone, Default)]
174#[allow(clippy::struct_excessive_bools)] pub struct MetricsSnapshot {
176 pub prompt_tokens: u64,
177 pub completion_tokens: u64,
178 pub total_tokens: u64,
179 pub reasoning_tokens: u64,
183 pub context_tokens: u64,
184 pub api_calls: u64,
185 pub active_skills: Vec<String>,
186 pub total_skills: usize,
187 pub mcp_server_count: usize,
189 pub mcp_tool_count: usize,
190 pub mcp_connected_count: usize,
192 pub mcp_servers: Vec<McpServerStatus>,
194 pub active_mcp_tools: Vec<String>,
195 pub sqlite_message_count: u64,
196 pub sqlite_conversation_id: Option<zeph_memory::ConversationId>,
197 pub qdrant_available: bool,
198 pub vector_backend: String,
199 pub embeddings_generated: u64,
200 pub last_llm_latency_ms: u64,
201 pub uptime_seconds: u64,
202 pub provider_name: String,
203 pub model_name: String,
204 pub summaries_count: u64,
205 pub context_compactions: u64,
206 pub compaction_hard_count: u64,
209 pub compaction_turns_after_hard: Vec<u64>,
213 pub compression_events: u64,
214 pub compression_tokens_saved: u64,
215 pub acon_results_compressed: u64,
217 pub acon_tokens_saved: u64,
219 pub tool_output_prunes: u64,
220 pub compaction_probe_passes: u64,
222 pub compaction_probe_soft_failures: u64,
224 pub compaction_probe_failures: u64,
226 pub compaction_probe_errors: u64,
228 pub last_probe_verdict: Option<zeph_memory::ProbeVerdict>,
230 pub last_probe_score: Option<f32>,
233 pub last_probe_category_scores: Option<Vec<zeph_memory::CategoryScore>>,
235 pub compaction_probe_threshold: f32,
237 pub compaction_probe_hard_fail_threshold: f32,
239 pub cache_read_tokens: u64,
240 pub cache_creation_tokens: u64,
241 pub cost_spent_cents: f64,
242 pub cost_cps_cents: Option<f64>,
244 pub cost_successful_tasks: u64,
246 pub provider_cost_breakdown: Vec<(String, crate::cost::ProviderUsage)>,
248 pub filter_raw_tokens: u64,
249 pub filter_saved_tokens: u64,
250 pub filter_applications: u64,
251 pub filter_total_commands: u64,
252 pub filter_filtered_commands: u64,
253 pub filter_confidence_full: u64,
254 pub filter_confidence_partial: u64,
255 pub filter_confidence_fallback: u64,
256 pub cancellations: u64,
257 pub server_compaction_events: u64,
258 pub sanitizer_runs: u64,
259 pub sanitizer_injection_flags: u64,
260 pub sanitizer_injection_fp_local: u64,
266 pub sanitizer_truncations: u64,
267 pub quarantine_invocations: u64,
268 pub quarantine_failures: u64,
269 pub classifier_tool_blocks: u64,
271 pub classifier_tool_suspicious: u64,
273 pub causal_ipi_flags: u64,
275 pub vigil_flags_total: u64,
277 pub vigil_blocks_total: u64,
279 pub exfiltration_images_blocked: u64,
280 pub exfiltration_tool_urls_flagged: u64,
281 pub exfiltration_memory_guards: u64,
282 pub pii_scrub_count: u64,
283 pub pii_ner_timeouts: u64,
285 pub pii_ner_circuit_breaker_trips: u64,
287 pub memory_validation_failures: u64,
288 pub rate_limit_trips: u64,
289 pub pre_execution_blocks: u64,
290 pub pre_execution_warnings: u64,
291 pub guardrail_enabled: bool,
293 pub guardrail_warn_mode: bool,
295 pub nli_enabled: bool,
297 pub nli_checks: u64,
299 pub nli_flags: u64,
301 pub secret_masking_enabled: bool,
303 pub secret_mask_registrations: u64,
305 pub secret_mask_applied: u64,
308 pub secret_unmask_misses: u64,
313 pub sub_agents: Vec<SubAgentMetrics>,
314 pub skill_confidence: Vec<SkillConfidence>,
315 pub scheduled_tasks: Vec<[String; 4]>,
317 pub router_thompson_stats: Vec<(String, f64, f64)>,
319 pub security_events: VecDeque<SecurityEvent>,
321 pub orchestration: OrchestrationMetrics,
322 pub orchestration_graph: Option<TaskGraphSnapshot>,
324 pub graph_community_detection_failures: u64,
325 pub graph_entities_total: u64,
326 pub graph_edges_total: u64,
327 pub graph_communities_total: u64,
328 pub graph_extraction_count: u64,
329 pub graph_extraction_failures: u64,
330 pub extended_context: bool,
333 pub guidelines_version: u32,
335 pub guidelines_updated_at: String,
337 pub tool_cache_hits: u64,
338 pub tool_cache_misses: u64,
339 pub tool_cache_entries: usize,
340 pub semantic_fact_count: u64,
342 pub stt_model: Option<String>,
344 pub compaction_model: Option<String>,
346 pub provider_temperature: Option<f32>,
348 pub provider_top_p: Option<f32>,
350 pub embedding_model: String,
352 pub token_budget: Option<u64>,
354 pub compaction_threshold: Option<u32>,
356 pub vault_backend: String,
358 pub active_channel: String,
360 pub bg_inflight: u64,
362 pub bg_dropped: u64,
364 pub bg_completed: u64,
366 pub bg_enrichment_inflight: u64,
368 pub bg_telemetry_inflight: u64,
370 pub shell_background_runs: Vec<ShellBackgroundRunRow>,
372 pub self_learning_enabled: bool,
374 pub semantic_cache_enabled: bool,
376 pub cache_enabled: bool,
378 pub autosave_enabled: bool,
380 pub classifier: ClassifierMetricsSnapshot,
382 pub last_turn_timings: TurnTimings,
384 pub avg_turn_timings: TurnTimings,
386 pub max_turn_timings: TurnTimings,
390 pub timing_sample_count: u64,
392 pub egress_requests_total: u64,
394 pub egress_dropped_total: u64,
396 pub egress_blocked_total: u64,
398 pub context_max_tokens: u64,
404 pub compaction_last_before: u64,
406 pub compaction_last_after: u64,
408 pub compaction_last_at_ms: u64,
410 pub active_goal: Option<crate::goal::GoalSnapshot>,
412 pub cocoon_connected: Option<bool>,
415 pub cocoon_worker_count: u32,
417 pub cocoon_model_count: usize,
419 pub cocoon_ton_balance: Option<f64>,
421}
422
423#[derive(Debug, Clone, Default, serde::Serialize)]
429pub struct ShellBackgroundRunRow {
430 pub run_id: String,
432 pub command: String,
434 pub elapsed_secs: u64,
436}
437
438#[derive(Debug, Default)]
456pub struct StaticMetricsInit {
457 pub stt_model: Option<String>,
459 pub compaction_model: Option<String>,
461 pub semantic_cache_enabled: bool,
466 pub embedding_model: String,
468 pub self_learning_enabled: bool,
470 pub active_channel: String,
472 pub token_budget: Option<u64>,
474 pub compaction_threshold: Option<u32>,
476 pub vault_backend: String,
478 pub autosave_enabled: bool,
480 pub model_name_override: Option<String>,
484}
485
486fn strip_ctrl(s: &str) -> String {
492 let mut out = String::with_capacity(s.len());
493 let mut chars = s.chars().peekable();
494 while let Some(c) = chars.next() {
495 if c == '\x1b' {
496 if chars.peek() == Some(&'[') {
498 chars.next(); for inner in chars.by_ref() {
500 if ('\x40'..='\x7e').contains(&inner) {
501 break;
502 }
503 }
504 }
505 } else if c.is_control() && c != '\t' && c != '\n' && c != '\r' {
507 } else {
509 out.push(c);
510 }
511 }
512 out
513}
514
515impl From<&zeph_orchestration::TaskGraph> for TaskGraphSnapshot {
517 fn from(graph: &zeph_orchestration::TaskGraph) -> Self {
518 let tasks = graph
519 .tasks
520 .iter()
521 .map(|t| {
522 let error = t
523 .result
524 .as_ref()
525 .filter(|_| t.status == zeph_orchestration::TaskStatus::Failed)
526 .and_then(|r| {
527 if r.output.is_empty() {
528 None
529 } else {
530 let s = strip_ctrl(&r.output);
532 if s.len() > 80 {
533 let end = s.floor_char_boundary(79);
534 Some(format!("{}…", &s[..end]))
535 } else {
536 Some(s)
537 }
538 }
539 });
540 let duration_ms = t.result.as_ref().map_or(0, |r| r.duration_ms);
541 TaskSnapshotRow {
542 id: t.id.as_u32(),
543 title: strip_ctrl(&t.title),
544 status: t.status.to_string(),
545 agent: t.assigned_agent.as_deref().map(strip_ctrl),
546 duration_ms,
547 error,
548 }
549 })
550 .collect();
551 Self {
552 graph_id: graph.id.to_string(),
553 goal: strip_ctrl(&graph.goal),
554 status: graph.status.to_string(),
555 tasks,
556 completed_at: None,
557 }
558 }
559}
560
561pub struct MetricsCollector {
562 tx: watch::Sender<MetricsSnapshot>,
563}
564
565impl MetricsCollector {
566 #[must_use]
567 pub fn new() -> (Self, watch::Receiver<MetricsSnapshot>) {
568 let (tx, rx) = watch::channel(MetricsSnapshot::default());
569 (Self { tx }, rx)
570 }
571
572 pub fn update(&self, f: impl FnOnce(&mut MetricsSnapshot)) {
573 self.tx.send_modify(f);
574 }
575
576 pub fn set_context_max_tokens(&self, max_tokens: u64) {
591 self.tx.send_modify(|m| m.context_max_tokens = max_tokens);
592 }
593
594 pub fn record_compaction(&self, before: u64, after: u64, at_ms: u64) {
612 self.tx.send_modify(|m| {
613 m.compaction_last_before = before;
614 m.compaction_last_after = after;
615 m.compaction_last_at_ms = at_ms;
616 });
617 }
618
619 #[must_use]
625 pub fn sender(&self) -> watch::Sender<MetricsSnapshot> {
626 self.tx.clone()
627 }
628}
629
630pub trait HistogramRecorder: Send + Sync {
668 fn observe_llm_latency(&self, duration: std::time::Duration);
670
671 fn observe_turn_duration(&self, duration: std::time::Duration);
673
674 fn observe_tool_execution(&self, duration: std::time::Duration);
676
677 fn observe_bg_task(&self, class_label: &str, duration: std::time::Duration);
681}
682
683#[cfg(test)]
684mod tests {
685 #![allow(clippy::field_reassign_with_default)]
686
687 use super::*;
688
689 #[test]
690 fn default_metrics_snapshot() {
691 let m = MetricsSnapshot::default();
692 assert_eq!(m.total_tokens, 0);
693 assert_eq!(m.api_calls, 0);
694 assert!(m.active_skills.is_empty());
695 assert!(m.active_mcp_tools.is_empty());
696 assert_eq!(m.mcp_tool_count, 0);
697 assert_eq!(m.mcp_server_count, 0);
698 assert!(m.provider_name.is_empty());
699 assert_eq!(m.summaries_count, 0);
700 assert!(m.stt_model.is_none());
702 assert!(m.compaction_model.is_none());
703 assert!(m.provider_temperature.is_none());
704 assert!(m.provider_top_p.is_none());
705 assert!(m.active_channel.is_empty());
706 assert!(m.embedding_model.is_empty());
707 assert!(m.token_budget.is_none());
708 assert!(!m.self_learning_enabled);
709 assert!(!m.semantic_cache_enabled);
710 }
711
712 #[test]
713 fn metrics_collector_update_phase2_fields() {
714 let (collector, rx) = MetricsCollector::new();
715 collector.update(|m| {
716 m.stt_model = Some("whisper-1".into());
717 m.compaction_model = Some("haiku".into());
718 m.provider_temperature = Some(0.7);
719 m.provider_top_p = Some(0.95);
720 m.active_channel = "tui".into();
721 m.embedding_model = "nomic-embed-text".into();
722 m.token_budget = Some(200_000);
723 m.self_learning_enabled = true;
724 m.semantic_cache_enabled = true;
725 });
726 let s = rx.borrow();
727 assert_eq!(s.stt_model.as_deref(), Some("whisper-1"));
728 assert_eq!(s.compaction_model.as_deref(), Some("haiku"));
729 assert_eq!(s.provider_temperature, Some(0.7));
730 assert_eq!(s.provider_top_p, Some(0.95));
731 assert_eq!(s.active_channel, "tui");
732 assert_eq!(s.embedding_model, "nomic-embed-text");
733 assert_eq!(s.token_budget, Some(200_000));
734 assert!(s.self_learning_enabled);
735 assert!(s.semantic_cache_enabled);
736 }
737
738 #[test]
739 fn metrics_collector_update() {
740 let (collector, rx) = MetricsCollector::new();
741 collector.update(|m| {
742 m.api_calls = 5;
743 m.total_tokens = 1000;
744 });
745 let snapshot = rx.borrow().clone();
746 assert_eq!(snapshot.api_calls, 5);
747 assert_eq!(snapshot.total_tokens, 1000);
748 }
749
750 #[test]
751 fn metrics_collector_multiple_updates() {
752 let (collector, rx) = MetricsCollector::new();
753 collector.update(|m| m.api_calls = 1);
754 collector.update(|m| m.api_calls += 1);
755 assert_eq!(rx.borrow().api_calls, 2);
756 }
757
758 #[test]
759 fn metrics_snapshot_clone() {
760 let mut m = MetricsSnapshot::default();
761 m.provider_name = "ollama".into();
762 let cloned = m.clone();
763 assert_eq!(cloned.provider_name, "ollama");
764 }
765
766 #[test]
767 fn filter_metrics_tracking() {
768 let (collector, rx) = MetricsCollector::new();
769 collector.update(|m| {
770 m.filter_raw_tokens += 250;
771 m.filter_saved_tokens += 200;
772 m.filter_applications += 1;
773 });
774 collector.update(|m| {
775 m.filter_raw_tokens += 100;
776 m.filter_saved_tokens += 80;
777 m.filter_applications += 1;
778 });
779 let s = rx.borrow();
780 assert_eq!(s.filter_raw_tokens, 350);
781 assert_eq!(s.filter_saved_tokens, 280);
782 assert_eq!(s.filter_applications, 2);
783 }
784
785 #[test]
786 fn filter_confidence_and_command_metrics() {
787 let (collector, rx) = MetricsCollector::new();
788 collector.update(|m| {
789 m.filter_total_commands += 1;
790 m.filter_filtered_commands += 1;
791 m.filter_confidence_full += 1;
792 });
793 collector.update(|m| {
794 m.filter_total_commands += 1;
795 m.filter_confidence_partial += 1;
796 });
797 let s = rx.borrow();
798 assert_eq!(s.filter_total_commands, 2);
799 assert_eq!(s.filter_filtered_commands, 1);
800 assert_eq!(s.filter_confidence_full, 1);
801 assert_eq!(s.filter_confidence_partial, 1);
802 assert_eq!(s.filter_confidence_fallback, 0);
803 }
804
805 #[test]
806 fn summaries_count_tracks_summarizations() {
807 let (collector, rx) = MetricsCollector::new();
808 collector.update(|m| m.summaries_count += 1);
809 collector.update(|m| m.summaries_count += 1);
810 assert_eq!(rx.borrow().summaries_count, 2);
811 }
812
813 #[test]
814 fn cancellations_counter_increments() {
815 let (collector, rx) = MetricsCollector::new();
816 assert_eq!(rx.borrow().cancellations, 0);
817 collector.update(|m| m.cancellations += 1);
818 collector.update(|m| m.cancellations += 1);
819 assert_eq!(rx.borrow().cancellations, 2);
820 }
821
822 #[test]
823 fn security_event_detail_exact_128_not_truncated() {
824 let s = "a".repeat(128);
825 let ev = SecurityEvent::new(SecurityEventCategory::InjectionFlag, "src", s.clone());
826 assert_eq!(ev.detail, s, "128-char string must not be truncated");
827 }
828
829 #[test]
830 fn security_event_detail_129_is_truncated() {
831 let s = "a".repeat(129);
832 let ev = SecurityEvent::new(SecurityEventCategory::InjectionFlag, "src", s);
833 assert!(
834 ev.detail.ends_with('…'),
835 "129-char string must end with ellipsis"
836 );
837 assert!(
838 ev.detail.len() <= 130,
839 "truncated detail must be at most 130 bytes"
840 );
841 }
842
843 #[test]
844 fn security_event_detail_multibyte_utf8_no_panic() {
845 let s = "中".repeat(43);
847 let ev = SecurityEvent::new(SecurityEventCategory::InjectionFlag, "src", s);
848 assert!(ev.detail.ends_with('…'));
849 }
850
851 #[test]
852 fn security_event_source_capped_at_64_chars() {
853 let long_source = "x".repeat(200);
854 let ev = SecurityEvent::new(SecurityEventCategory::InjectionFlag, long_source, "detail");
855 assert_eq!(ev.source.len(), 64);
856 }
857
858 #[test]
859 fn security_event_source_strips_control_chars() {
860 let source = "tool\x00name\x1b[31m";
861 let ev = SecurityEvent::new(SecurityEventCategory::InjectionFlag, source, "detail");
862 assert!(!ev.source.contains('\x00'));
863 assert!(!ev.source.contains('\x1b'));
864 }
865
866 #[test]
867 fn security_event_category_as_str() {
868 assert_eq!(SecurityEventCategory::InjectionFlag.as_str(), "injection");
869 assert_eq!(SecurityEventCategory::ExfiltrationBlock.as_str(), "exfil");
870 assert_eq!(SecurityEventCategory::Quarantine.as_str(), "quarantine");
871 assert_eq!(SecurityEventCategory::Truncation.as_str(), "truncation");
872 assert_eq!(
873 SecurityEventCategory::CrossBoundaryMcpToAcp.as_str(),
874 "cross_boundary_mcp_to_acp"
875 );
876 }
877
878 #[test]
879 fn ring_buffer_respects_cap_via_update() {
880 let (collector, rx) = MetricsCollector::new();
881 for i in 0..110u64 {
882 let event = SecurityEvent::new(
883 SecurityEventCategory::InjectionFlag,
884 "src",
885 format!("event {i}"),
886 );
887 collector.update(|m| {
888 if m.security_events.len() >= SECURITY_EVENT_CAP {
889 m.security_events.pop_front();
890 }
891 m.security_events.push_back(event);
892 });
893 }
894 let snap = rx.borrow();
895 assert_eq!(snap.security_events.len(), SECURITY_EVENT_CAP);
896 assert!(snap.security_events.back().unwrap().detail.contains("109"));
898 }
899
900 #[test]
901 fn security_events_empty_by_default() {
902 let m = MetricsSnapshot::default();
903 assert!(m.security_events.is_empty());
904 }
905
906 #[test]
907 fn orchestration_metrics_default_zero() {
908 let m = OrchestrationMetrics::default();
909 assert_eq!(m.plans_total, 0);
910 assert_eq!(m.tasks_total, 0);
911 assert_eq!(m.tasks_completed, 0);
912 assert_eq!(m.tasks_failed, 0);
913 assert_eq!(m.tasks_skipped, 0);
914 }
915
916 #[test]
917 fn metrics_snapshot_includes_orchestration_default_zero() {
918 let m = MetricsSnapshot::default();
919 assert_eq!(m.orchestration.plans_total, 0);
920 assert_eq!(m.orchestration.tasks_total, 0);
921 assert_eq!(m.orchestration.tasks_completed, 0);
922 }
923
924 #[test]
925 fn orchestration_metrics_update_via_collector() {
926 let (collector, rx) = MetricsCollector::new();
927 collector.update(|m| {
928 m.orchestration.plans_total += 1;
929 m.orchestration.tasks_total += 5;
930 m.orchestration.tasks_completed += 3;
931 m.orchestration.tasks_failed += 1;
932 m.orchestration.tasks_skipped += 1;
933 });
934 let s = rx.borrow();
935 assert_eq!(s.orchestration.plans_total, 1);
936 assert_eq!(s.orchestration.tasks_total, 5);
937 assert_eq!(s.orchestration.tasks_completed, 3);
938 assert_eq!(s.orchestration.tasks_failed, 1);
939 assert_eq!(s.orchestration.tasks_skipped, 1);
940 }
941
942 #[test]
943 fn strip_ctrl_removes_escape_sequences() {
944 let input = "hello\x1b[31mworld\x00end";
945 let result = strip_ctrl(input);
946 assert_eq!(result, "helloworldend");
947 }
948
949 #[test]
950 fn strip_ctrl_allows_tab_lf_cr() {
951 let input = "a\tb\nc\rd";
952 let result = strip_ctrl(input);
953 assert_eq!(result, "a\tb\nc\rd");
954 }
955
956 #[test]
957 fn task_graph_snapshot_is_stale_after_30s() {
958 let mut snap = TaskGraphSnapshot::default();
959 assert!(!snap.is_stale());
961 snap.completed_at = Some(std::time::Instant::now());
963 assert!(!snap.is_stale());
964 snap.completed_at = Some(
966 std::time::Instant::now()
967 .checked_sub(std::time::Duration::from_secs(31))
968 .unwrap(),
969 );
970 assert!(snap.is_stale());
971 }
972
973 #[test]
975 fn task_graph_snapshot_from_task_graph_maps_fields() {
976 use zeph_orchestration::{GraphStatus, TaskGraph, TaskNode, TaskResult, TaskStatus};
977
978 let mut graph = TaskGraph::new("My goal");
979 let mut task = TaskNode::new(0, "Do work", "description");
980 task.status = TaskStatus::Failed;
981 task.assigned_agent = Some("agent-1".into());
982 task.result = Some(TaskResult {
983 output: "error occurred here".into(),
984 artifacts: vec![],
985 duration_ms: 1234,
986 agent_id: None,
987 agent_def: None,
988 });
989 graph.tasks.push(task);
990 graph.status = GraphStatus::Failed;
991
992 let snap = TaskGraphSnapshot::from(&graph);
993 assert_eq!(snap.goal, "My goal");
994 assert_eq!(snap.status, "failed");
995 assert_eq!(snap.tasks.len(), 1);
996 let row = &snap.tasks[0];
997 assert_eq!(row.title, "Do work");
998 assert_eq!(row.status, "failed");
999 assert_eq!(row.agent.as_deref(), Some("agent-1"));
1000 assert_eq!(row.duration_ms, 1234);
1001 assert!(row.error.as_deref().unwrap().contains("error occurred"));
1002 }
1003
1004 #[test]
1006 fn task_graph_snapshot_from_compiles_with_feature() {
1007 use zeph_orchestration::TaskGraph;
1008 let graph = TaskGraph::new("feature flag test");
1009 let snap = TaskGraphSnapshot::from(&graph);
1010 assert_eq!(snap.goal, "feature flag test");
1011 assert!(snap.tasks.is_empty());
1012 assert!(!snap.is_stale());
1013 }
1014
1015 #[test]
1017 fn task_graph_snapshot_error_truncated_at_80_chars() {
1018 use zeph_orchestration::{TaskGraph, TaskNode, TaskResult, TaskStatus};
1019
1020 let mut graph = TaskGraph::new("goal");
1021 let mut task = TaskNode::new(0, "t", "d");
1022 task.status = TaskStatus::Failed;
1023 task.result = Some(TaskResult {
1024 output: "e".repeat(100),
1025 artifacts: vec![],
1026 duration_ms: 0,
1027 agent_id: None,
1028 agent_def: None,
1029 });
1030 graph.tasks.push(task);
1031
1032 let snap = TaskGraphSnapshot::from(&graph);
1033 let err = snap.tasks[0].error.as_ref().unwrap();
1034 assert!(err.ends_with('…'), "truncated error must end with ellipsis");
1035 assert!(
1036 err.len() <= 83,
1037 "truncated error must not exceed 80 chars + ellipsis"
1038 );
1039 }
1040
1041 #[test]
1043 fn task_graph_snapshot_strips_control_chars_from_title() {
1044 use zeph_orchestration::{TaskGraph, TaskNode};
1045
1046 let mut graph = TaskGraph::new("goal\x1b[31m");
1047 let task = TaskNode::new(0, "title\x00injected", "d");
1048 graph.tasks.push(task);
1049
1050 let snap = TaskGraphSnapshot::from(&graph);
1051 assert!(!snap.goal.contains('\x1b'), "goal must not contain escape");
1052 assert!(
1053 !snap.tasks[0].title.contains('\x00'),
1054 "title must not contain null byte"
1055 );
1056 }
1057
1058 #[test]
1059 fn graph_metrics_default_zero() {
1060 let m = MetricsSnapshot::default();
1061 assert_eq!(m.graph_entities_total, 0);
1062 assert_eq!(m.graph_edges_total, 0);
1063 assert_eq!(m.graph_communities_total, 0);
1064 assert_eq!(m.graph_extraction_count, 0);
1065 assert_eq!(m.graph_extraction_failures, 0);
1066 }
1067
1068 #[test]
1069 fn graph_metrics_update_via_collector() {
1070 let (collector, rx) = MetricsCollector::new();
1071 collector.update(|m| {
1072 m.graph_entities_total = 5;
1073 m.graph_edges_total = 10;
1074 m.graph_communities_total = 2;
1075 m.graph_extraction_count = 7;
1076 m.graph_extraction_failures = 1;
1077 });
1078 let snapshot = rx.borrow().clone();
1079 assert_eq!(snapshot.graph_entities_total, 5);
1080 assert_eq!(snapshot.graph_edges_total, 10);
1081 assert_eq!(snapshot.graph_communities_total, 2);
1082 assert_eq!(snapshot.graph_extraction_count, 7);
1083 assert_eq!(snapshot.graph_extraction_failures, 1);
1084 }
1085
1086 #[test]
1087 fn histogram_recorder_trait_is_object_safe() {
1088 use std::sync::Arc;
1089 use std::time::Duration;
1090
1091 struct NoOpRecorder;
1092 impl HistogramRecorder for NoOpRecorder {
1093 fn observe_llm_latency(&self, _: Duration) {}
1094 fn observe_turn_duration(&self, _: Duration) {}
1095 fn observe_tool_execution(&self, _: Duration) {}
1096 fn observe_bg_task(&self, _: &str, _: Duration) {}
1097 }
1098
1099 let recorder: Arc<dyn HistogramRecorder> = Arc::new(NoOpRecorder);
1101 recorder.observe_llm_latency(Duration::from_millis(500));
1102 recorder.observe_turn_duration(Duration::from_secs(3));
1103 recorder.observe_tool_execution(Duration::from_millis(100));
1104 }
1105}