Skip to main content

pulse_system_types/
monitoring.rs

1//! Shared monitoring types and traits for the pulse-null plugin ecosystem.
2//!
3//! This module defines the contract between pulse-null core and its monitoring
4//! plugins (praxis-echo, vigil-echo, pulse-echo). Core calls through trait
5//! objects; plugins implement the traits.
6
7use std::path::Path;
8
9use serde::{Deserialize, Serialize};
10
11// ---------------------------------------------------------------------------
12// Pipeline types (praxis-echo domain)
13// ---------------------------------------------------------------------------
14
15/// Threshold configuration for pipeline documents.
16#[derive(Debug, Clone, Serialize, Deserialize)]
17pub struct PipelineThresholds {
18    pub learning_soft: usize,
19    pub learning_hard: usize,
20    pub thoughts_soft: usize,
21    pub thoughts_hard: usize,
22    pub curiosity_soft: usize,
23    pub curiosity_hard: usize,
24    pub reflections_soft: usize,
25    pub reflections_hard: usize,
26    pub praxis_soft: usize,
27    pub praxis_hard: usize,
28}
29
30impl Default for PipelineThresholds {
31    fn default() -> Self {
32        Self {
33            learning_soft: 5,
34            learning_hard: 8,
35            thoughts_soft: 5,
36            thoughts_hard: 10,
37            curiosity_soft: 3,
38            curiosity_hard: 7,
39            reflections_soft: 15,
40            reflections_hard: 20,
41            praxis_soft: 5,
42            praxis_hard: 10,
43        }
44    }
45}
46
47/// Status of a document relative to its thresholds.
48#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
49pub enum ThresholdStatus {
50    Green,
51    Yellow,
52    Red,
53}
54
55impl std::fmt::Display for ThresholdStatus {
56    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
57        match self {
58            Self::Green => write!(f, "green"),
59            Self::Yellow => write!(f, "yellow"),
60            Self::Red => write!(f, "red"),
61        }
62    }
63}
64
65/// Health report for a single pipeline document.
66#[derive(Debug, Clone, Serialize, Deserialize)]
67pub struct DocumentHealth {
68    pub count: usize,
69    pub soft: usize,
70    pub hard: usize,
71    pub status: ThresholdStatus,
72}
73
74/// Full pipeline health report across all documents.
75#[derive(Debug, Clone)]
76pub struct PipelineHealth {
77    pub learning: DocumentHealth,
78    pub thoughts: DocumentHealth,
79    pub curiosity: DocumentHealth,
80    pub reflections: DocumentHealth,
81    pub praxis: DocumentHealth,
82    pub warnings: Vec<String>,
83}
84
85/// Document entry counts for freeze detection.
86#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize, Default)]
87pub struct DocumentCounts {
88    pub learning: usize,
89    pub thoughts: usize,
90    pub curiosity: usize,
91    pub reflections: usize,
92    pub praxis: usize,
93}
94
95/// Persistent pipeline state tracked across sessions.
96#[derive(Debug, Clone, Serialize, Deserialize, Default)]
97pub struct PipelineState {
98    pub last_updated: Option<String>,
99    pub session_count: u32,
100    pub sessions_without_movement: u32,
101    pub last_counts: DocumentCounts,
102}
103
104impl PipelineState {
105    /// Update counts and detect pipeline freezes.
106    ///
107    /// `now_iso` should be an ISO 8601 timestamp string (e.g. from `chrono::Utc::now().to_rfc3339()`).
108    pub fn update_counts(&mut self, new_counts: &DocumentCounts, now_iso: &str) {
109        if *new_counts == self.last_counts {
110            self.sessions_without_movement += 1;
111        } else {
112            self.sessions_without_movement = 0;
113        }
114        self.last_counts = new_counts.clone();
115        self.session_count += 1;
116        self.last_updated = Some(now_iso.to_string());
117    }
118}
119
120// ---------------------------------------------------------------------------
121// Cognitive types (vigil-echo domain)
122// ---------------------------------------------------------------------------
123
124/// Cognitive health status levels.
125#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
126pub enum CognitiveStatus {
127    Healthy,
128    Watch,
129    Concern,
130    Alert,
131}
132
133impl std::fmt::Display for CognitiveStatus {
134    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
135        match self {
136            Self::Healthy => write!(f, "HEALTHY"),
137            Self::Watch => write!(f, "WATCH"),
138            Self::Concern => write!(f, "CONCERN"),
139            Self::Alert => write!(f, "ALERT"),
140        }
141    }
142}
143
144/// Signal trend direction.
145#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
146pub enum Trend {
147    Improving,
148    Stable,
149    Declining,
150}
151
152impl std::fmt::Display for Trend {
153    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
154        match self {
155            Self::Improving => write!(f, "improving"),
156            Self::Stable => write!(f, "stable"),
157            Self::Declining => write!(f, "declining"),
158        }
159    }
160}
161
162/// Full cognitive health assessment.
163#[derive(Debug, Clone)]
164pub struct CognitiveHealth {
165    pub status: CognitiveStatus,
166    pub vocabulary_trend: Trend,
167    pub question_trend: Trend,
168    pub evidence_trend: Trend,
169    pub progress_trend: Trend,
170    pub suggestions: Vec<String>,
171    pub sufficient_data: bool,
172}
173
174/// A single frame of cognitive signals extracted from LLM output.
175#[derive(Debug, Clone, Serialize, Deserialize)]
176pub struct SignalFrame {
177    pub timestamp: String,
178    pub task_id: String,
179    pub vocabulary_diversity: f64,
180    pub question_count: usize,
181    pub evidence_references: usize,
182    pub thought_progress: bool,
183}
184
185// ---------------------------------------------------------------------------
186// Outcome types (pulse-echo domain)
187// ---------------------------------------------------------------------------
188
189/// A record of what happened during task execution.
190#[derive(Debug, Clone, Serialize, Deserialize)]
191pub struct OutcomeRecord {
192    pub task_id: String,
193    pub timestamp: String,
194    pub domain: String,
195    pub task_type: String,
196    pub description: String,
197    pub outcome: String,
198    pub tokens_used: u32,
199    pub tool_rounds: u32,
200}
201
202// ---------------------------------------------------------------------------
203// Calibration types (praxis-echo domain)
204// ---------------------------------------------------------------------------
205
206/// Confidence level for a threshold recommendation.
207#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
208pub enum Confidence {
209    Low,
210    Medium,
211    High,
212}
213
214/// A recommendation to adjust a document's thresholds.
215#[derive(Debug, Clone, Serialize, Deserialize)]
216pub struct ThresholdRecommendation {
217    pub document: String,
218    pub current_soft: usize,
219    pub current_hard: usize,
220    pub recommended_soft: Option<usize>,
221    pub recommended_hard: Option<usize>,
222    pub reason: String,
223    pub confidence: Confidence,
224    pub evidence_count: usize,
225}
226
227/// Summary of outcome data used in calibration.
228#[derive(Debug, Clone, Serialize, Deserialize)]
229pub struct OutcomeSummary {
230    pub total: usize,
231    pub success_rate: f64,
232    /// (domain, count, success_rate)
233    pub domains: Vec<(String, usize, f64)>,
234}
235
236/// Result of a calibration analysis run.
237#[derive(Debug, Clone, Serialize, Deserialize)]
238pub struct CalibrationReport {
239    pub generated_at: String,
240    pub recommendations: Vec<ThresholdRecommendation>,
241    pub sample_size: usize,
242    pub outcome_summary: OutcomeSummary,
243}
244
245/// A point-in-time snapshot of pipeline document counts.
246#[derive(Debug, Clone, Serialize, Deserialize)]
247pub struct PipelineSnapshot {
248    pub timestamp: String,
249    pub learning: usize,
250    pub thoughts: usize,
251    pub curiosity: usize,
252    pub reflections: usize,
253    pub praxis: usize,
254}
255
256// ---------------------------------------------------------------------------
257// Traits
258// ---------------------------------------------------------------------------
259
260/// Pipeline document monitoring and archiving.
261///
262/// Implemented by praxis-echo. Used by pulse-null core for:
263/// - Banner rendering (startup health bars)
264/// - System prompt injection (pipeline state for LLM context)
265/// - Post-execution hooks (update state, auto-archive)
266/// - CLI commands (pipeline health, archive)
267/// - Dashboard endpoint (JSON health)
268pub trait PipelineMonitor: Send + Sync {
269    /// Calculate pipeline health from document files on disk.
270    fn calculate(&self, root_dir: &Path, thresholds: &PipelineThresholds) -> PipelineHealth;
271
272    /// Render pipeline health as text for prompt injection.
273    fn render_for_prompt(
274        &self,
275        health: &PipelineHealth,
276        sessions_frozen: u32,
277        freeze_threshold: u32,
278    ) -> String;
279
280    /// Extract document counts from a health report (for freeze detection).
281    fn counts_from_health(&self, health: &PipelineHealth) -> DocumentCounts;
282
283    /// Load persistent pipeline state from disk.
284    fn load_state(&self, root_dir: &Path) -> PipelineState;
285
286    /// Save persistent pipeline state to disk.
287    fn save_state(
288        &self,
289        root_dir: &Path,
290        state: &PipelineState,
291    ) -> Result<(), Box<dyn std::error::Error>>;
292
293    /// Check all documents against hard limits and auto-archive overflow.
294    /// Returns list of document names that were archived.
295    fn check_and_archive(
296        &self,
297        root_dir: &Path,
298        thresholds: &PipelineThresholds,
299        health: &PipelineHealth,
300    ) -> Vec<String>;
301
302    /// List archived files, optionally filtered by document name.
303    fn list_archives(
304        &self,
305        root_dir: &Path,
306        document: Option<&str>,
307    ) -> Result<Vec<String>, Box<dyn std::error::Error>>;
308
309    /// Manually archive a specific document by name.
310    fn archive_by_name(
311        &self,
312        root_dir: &Path,
313        document: &str,
314    ) -> Result<String, Box<dyn std::error::Error>>;
315}
316
317/// Metacognitive signal tracking and health assessment.
318///
319/// Implemented by vigil-echo. Used by pulse-null core for:
320/// - Banner rendering (cognitive health status + trend arrows)
321/// - System prompt injection (assessment text for LLM context)
322/// - Post-execution hooks (extract and record signals, detect changes)
323/// - Dashboard endpoint (JSON signals)
324pub trait CognitiveMonitor: Send + Sync {
325    /// Perform a cognitive health assessment from signal history.
326    fn assess(&self, root_dir: &Path, window_size: usize, min_samples: usize) -> CognitiveHealth;
327
328    /// Render cognitive health as text for prompt injection.
329    fn render_for_prompt(&self, health: &CognitiveHealth) -> String;
330
331    /// Extract cognitive signals from LLM output text.
332    fn extract(&self, content: &str, task_id: &str) -> SignalFrame;
333
334    /// Append a signal frame to the rolling window on disk.
335    fn record(
336        &self,
337        root_dir: &Path,
338        frame: SignalFrame,
339        window_size: usize,
340    ) -> Result<(), Box<dyn std::error::Error>>;
341}
342
343/// Task execution outcome recording.
344///
345/// Implemented by pulse-echo. Used by pulse-null core for:
346/// - Post-execution hooks (build and record outcomes after tasks/intents)
347pub trait OutcomeTracker: Send + Sync {
348    /// Build an outcome record from task execution results.
349    fn build_outcome(
350        &self,
351        task_id: &str,
352        task_name: &str,
353        response_text: &str,
354        tool_rounds: u32,
355        input_tokens: u32,
356        output_tokens: u32,
357    ) -> OutcomeRecord;
358
359    /// Record an outcome to persistent storage.
360    fn record_outcome(
361        &self,
362        docs_dir: &Path,
363        outcome: OutcomeRecord,
364        max_outcomes: usize,
365    ) -> Result<(), Box<dyn std::error::Error>>;
366}
367
368// ---------------------------------------------------------------------------
369// Tests
370// ---------------------------------------------------------------------------
371
372#[cfg(test)]
373mod tests {
374    use super::*;
375
376    #[test]
377    fn threshold_status_display() {
378        assert_eq!(ThresholdStatus::Green.to_string(), "green");
379        assert_eq!(ThresholdStatus::Yellow.to_string(), "yellow");
380        assert_eq!(ThresholdStatus::Red.to_string(), "red");
381    }
382
383    #[test]
384    fn cognitive_status_display() {
385        assert_eq!(CognitiveStatus::Healthy.to_string(), "HEALTHY");
386        assert_eq!(CognitiveStatus::Watch.to_string(), "WATCH");
387        assert_eq!(CognitiveStatus::Concern.to_string(), "CONCERN");
388        assert_eq!(CognitiveStatus::Alert.to_string(), "ALERT");
389    }
390
391    #[test]
392    fn trend_display() {
393        assert_eq!(Trend::Improving.to_string(), "improving");
394        assert_eq!(Trend::Stable.to_string(), "stable");
395        assert_eq!(Trend::Declining.to_string(), "declining");
396    }
397
398    #[test]
399    fn pipeline_thresholds_default() {
400        let t = PipelineThresholds::default();
401        assert_eq!(t.learning_soft, 5);
402        assert_eq!(t.learning_hard, 8);
403        assert_eq!(t.curiosity_hard, 7);
404    }
405
406    #[test]
407    fn pipeline_state_update_counts_detects_freeze() {
408        let mut state = PipelineState::default();
409        let counts = DocumentCounts {
410            learning: 3,
411            thoughts: 2,
412            curiosity: 1,
413            reflections: 5,
414            praxis: 2,
415        };
416
417        state.update_counts(&counts, "2026-03-05T12:00:00Z");
418        assert_eq!(state.sessions_without_movement, 0);
419        assert_eq!(state.session_count, 1);
420
421        // Same counts — should increment freeze counter
422        state.update_counts(&counts, "2026-03-05T13:00:00Z");
423        assert_eq!(state.sessions_without_movement, 1);
424        assert_eq!(state.session_count, 2);
425
426        // Different counts — should reset freeze counter
427        let new_counts = DocumentCounts {
428            learning: 4,
429            ..counts
430        };
431        state.update_counts(&new_counts, "2026-03-05T14:00:00Z");
432        assert_eq!(state.sessions_without_movement, 0);
433        assert_eq!(state.session_count, 3);
434    }
435
436    #[test]
437    fn document_counts_default() {
438        let counts = DocumentCounts::default();
439        assert_eq!(counts.learning, 0);
440        assert_eq!(counts.thoughts, 0);
441    }
442
443    #[test]
444    fn threshold_status_equality() {
445        assert_eq!(ThresholdStatus::Red, ThresholdStatus::Red);
446        assert_ne!(ThresholdStatus::Red, ThresholdStatus::Green);
447    }
448
449    #[test]
450    fn signal_frame_serializes() {
451        let frame = SignalFrame {
452            timestamp: "2026-03-05T12:00:00Z".to_string(),
453            task_id: "test".to_string(),
454            vocabulary_diversity: 0.72,
455            question_count: 3,
456            evidence_references: 5,
457            thought_progress: true,
458        };
459        let json = serde_json::to_string(&frame).unwrap();
460        let back: SignalFrame = serde_json::from_str(&json).unwrap();
461        assert_eq!(back.task_id, "test");
462        assert!((back.vocabulary_diversity - 0.72).abs() < f64::EPSILON);
463    }
464
465    #[test]
466    fn outcome_record_serializes() {
467        let record = OutcomeRecord {
468            task_id: "task-1".to_string(),
469            timestamp: "2026-03-05T12:00:00Z".to_string(),
470            domain: "research".to_string(),
471            task_type: "research".to_string(),
472            description: "Deep dive".to_string(),
473            outcome: "success".to_string(),
474            tokens_used: 1500,
475            tool_rounds: 3,
476        };
477        let json = serde_json::to_string(&record).unwrap();
478        let back: OutcomeRecord = serde_json::from_str(&json).unwrap();
479        assert_eq!(back.task_id, "task-1");
480        assert_eq!(back.tokens_used, 1500);
481    }
482}