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    /// Self-prediction before execution ("I predict this session will...")
201    #[serde(default, skip_serializing_if = "Option::is_none")]
202    pub prediction: Option<String>,
203    /// Affective tag after execution (positive / negative / neutral / surprising)
204    #[serde(default, skip_serializing_if = "Option::is_none")]
205    pub valence: Option<String>,
206}
207
208// ---------------------------------------------------------------------------
209// Calibration types (praxis-echo domain)
210// ---------------------------------------------------------------------------
211
212/// Confidence level for a threshold recommendation.
213#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
214pub enum Confidence {
215    Low,
216    Medium,
217    High,
218}
219
220/// A recommendation to adjust a document's thresholds.
221#[derive(Debug, Clone, Serialize, Deserialize)]
222pub struct ThresholdRecommendation {
223    pub document: String,
224    pub current_soft: usize,
225    pub current_hard: usize,
226    pub recommended_soft: Option<usize>,
227    pub recommended_hard: Option<usize>,
228    pub reason: String,
229    pub confidence: Confidence,
230    pub evidence_count: usize,
231}
232
233/// Summary of outcome data used in calibration.
234#[derive(Debug, Clone, Serialize, Deserialize)]
235pub struct OutcomeSummary {
236    pub total: usize,
237    pub success_rate: f64,
238    /// (domain, count, success_rate)
239    pub domains: Vec<(String, usize, f64)>,
240}
241
242/// Result of a calibration analysis run.
243#[derive(Debug, Clone, Serialize, Deserialize)]
244pub struct CalibrationReport {
245    pub generated_at: String,
246    pub recommendations: Vec<ThresholdRecommendation>,
247    pub sample_size: usize,
248    pub outcome_summary: OutcomeSummary,
249}
250
251/// A point-in-time snapshot of pipeline document counts.
252#[derive(Debug, Clone, Serialize, Deserialize)]
253pub struct PipelineSnapshot {
254    pub timestamp: String,
255    pub learning: usize,
256    pub thoughts: usize,
257    pub curiosity: usize,
258    pub reflections: usize,
259    pub praxis: usize,
260}
261
262// ---------------------------------------------------------------------------
263// Traits
264// ---------------------------------------------------------------------------
265
266/// Pipeline document monitoring and archiving.
267///
268/// Implemented by praxis-echo. Used by pulse-null core for:
269/// - Banner rendering (startup health bars)
270/// - System prompt injection (pipeline state for LLM context)
271/// - Post-execution hooks (update state, auto-archive)
272/// - CLI commands (pipeline health, archive)
273/// - Dashboard endpoint (JSON health)
274pub trait PipelineMonitor: Send + Sync {
275    /// Calculate pipeline health from document files on disk.
276    fn calculate(&self, root_dir: &Path, thresholds: &PipelineThresholds) -> PipelineHealth;
277
278    /// Render pipeline health as text for prompt injection.
279    fn render_for_prompt(
280        &self,
281        health: &PipelineHealth,
282        sessions_frozen: u32,
283        freeze_threshold: u32,
284    ) -> String;
285
286    /// Extract document counts from a health report (for freeze detection).
287    fn counts_from_health(&self, health: &PipelineHealth) -> DocumentCounts;
288
289    /// Load persistent pipeline state from disk.
290    fn load_state(&self, root_dir: &Path) -> PipelineState;
291
292    /// Save persistent pipeline state to disk.
293    fn save_state(
294        &self,
295        root_dir: &Path,
296        state: &PipelineState,
297    ) -> Result<(), Box<dyn std::error::Error>>;
298
299    /// Check all documents against hard limits and auto-archive overflow.
300    /// Returns list of document names that were archived.
301    fn check_and_archive(
302        &self,
303        root_dir: &Path,
304        thresholds: &PipelineThresholds,
305        health: &PipelineHealth,
306    ) -> Vec<String>;
307
308    /// List archived files, optionally filtered by document name.
309    fn list_archives(
310        &self,
311        root_dir: &Path,
312        document: Option<&str>,
313    ) -> Result<Vec<String>, Box<dyn std::error::Error>>;
314
315    /// Manually archive a specific document by name.
316    fn archive_by_name(
317        &self,
318        root_dir: &Path,
319        document: &str,
320    ) -> Result<String, Box<dyn std::error::Error>>;
321}
322
323/// Metacognitive signal tracking and health assessment.
324///
325/// Implemented by vigil-echo. Used by pulse-null core for:
326/// - Banner rendering (cognitive health status + trend arrows)
327/// - System prompt injection (assessment text for LLM context)
328/// - Post-execution hooks (extract and record signals, detect changes)
329/// - Dashboard endpoint (JSON signals)
330pub trait CognitiveMonitor: Send + Sync {
331    /// Perform a cognitive health assessment from signal history.
332    fn assess(&self, root_dir: &Path, window_size: usize, min_samples: usize) -> CognitiveHealth;
333
334    /// Render cognitive health as text for prompt injection.
335    fn render_for_prompt(&self, health: &CognitiveHealth) -> String;
336
337    /// Extract cognitive signals from LLM output text.
338    fn extract(&self, content: &str, task_id: &str) -> SignalFrame;
339
340    /// Append a signal frame to the rolling window on disk.
341    fn record(
342        &self,
343        root_dir: &Path,
344        frame: SignalFrame,
345        window_size: usize,
346    ) -> Result<(), Box<dyn std::error::Error>>;
347}
348
349/// Task execution outcome recording.
350///
351/// Implemented by pulse-echo. Used by pulse-null core for:
352/// - Post-execution hooks (build and record outcomes after tasks/intents)
353pub trait OutcomeTracker: Send + Sync {
354    /// Build an outcome record from task execution results.
355    fn build_outcome(
356        &self,
357        task_id: &str,
358        task_name: &str,
359        response_text: &str,
360        tool_rounds: u32,
361        input_tokens: u32,
362        output_tokens: u32,
363    ) -> OutcomeRecord;
364
365    /// Record an outcome to persistent storage.
366    fn record_outcome(
367        &self,
368        docs_dir: &Path,
369        outcome: OutcomeRecord,
370        max_outcomes: usize,
371    ) -> Result<(), Box<dyn std::error::Error>>;
372}
373
374// ---------------------------------------------------------------------------
375// Tests
376// ---------------------------------------------------------------------------
377
378#[cfg(test)]
379mod tests {
380    use super::*;
381
382    #[test]
383    fn threshold_status_display() {
384        assert_eq!(ThresholdStatus::Green.to_string(), "green");
385        assert_eq!(ThresholdStatus::Yellow.to_string(), "yellow");
386        assert_eq!(ThresholdStatus::Red.to_string(), "red");
387    }
388
389    #[test]
390    fn cognitive_status_display() {
391        assert_eq!(CognitiveStatus::Healthy.to_string(), "HEALTHY");
392        assert_eq!(CognitiveStatus::Watch.to_string(), "WATCH");
393        assert_eq!(CognitiveStatus::Concern.to_string(), "CONCERN");
394        assert_eq!(CognitiveStatus::Alert.to_string(), "ALERT");
395    }
396
397    #[test]
398    fn trend_display() {
399        assert_eq!(Trend::Improving.to_string(), "improving");
400        assert_eq!(Trend::Stable.to_string(), "stable");
401        assert_eq!(Trend::Declining.to_string(), "declining");
402    }
403
404    #[test]
405    fn pipeline_thresholds_default() {
406        let t = PipelineThresholds::default();
407        assert_eq!(t.learning_soft, 5);
408        assert_eq!(t.learning_hard, 8);
409        assert_eq!(t.curiosity_hard, 7);
410    }
411
412    #[test]
413    fn pipeline_state_update_counts_detects_freeze() {
414        let mut state = PipelineState::default();
415        let counts = DocumentCounts {
416            learning: 3,
417            thoughts: 2,
418            curiosity: 1,
419            reflections: 5,
420            praxis: 2,
421        };
422
423        state.update_counts(&counts, "2026-03-05T12:00:00Z");
424        assert_eq!(state.sessions_without_movement, 0);
425        assert_eq!(state.session_count, 1);
426
427        // Same counts — should increment freeze counter
428        state.update_counts(&counts, "2026-03-05T13:00:00Z");
429        assert_eq!(state.sessions_without_movement, 1);
430        assert_eq!(state.session_count, 2);
431
432        // Different counts — should reset freeze counter
433        let new_counts = DocumentCounts {
434            learning: 4,
435            ..counts
436        };
437        state.update_counts(&new_counts, "2026-03-05T14:00:00Z");
438        assert_eq!(state.sessions_without_movement, 0);
439        assert_eq!(state.session_count, 3);
440    }
441
442    #[test]
443    fn document_counts_default() {
444        let counts = DocumentCounts::default();
445        assert_eq!(counts.learning, 0);
446        assert_eq!(counts.thoughts, 0);
447    }
448
449    #[test]
450    fn threshold_status_equality() {
451        assert_eq!(ThresholdStatus::Red, ThresholdStatus::Red);
452        assert_ne!(ThresholdStatus::Red, ThresholdStatus::Green);
453    }
454
455    #[test]
456    fn signal_frame_serializes() {
457        let frame = SignalFrame {
458            timestamp: "2026-03-05T12:00:00Z".to_string(),
459            task_id: "test".to_string(),
460            vocabulary_diversity: 0.72,
461            question_count: 3,
462            evidence_references: 5,
463            thought_progress: true,
464        };
465        let json = serde_json::to_string(&frame).unwrap();
466        let back: SignalFrame = serde_json::from_str(&json).unwrap();
467        assert_eq!(back.task_id, "test");
468        assert!((back.vocabulary_diversity - 0.72).abs() < f64::EPSILON);
469    }
470
471    #[test]
472    fn outcome_record_serializes() {
473        let record = OutcomeRecord {
474            task_id: "task-1".to_string(),
475            timestamp: "2026-03-05T12:00:00Z".to_string(),
476            domain: "research".to_string(),
477            task_type: "research".to_string(),
478            description: "Deep dive".to_string(),
479            outcome: "success".to_string(),
480            tokens_used: 1500,
481            tool_rounds: 3,
482            prediction: None,
483            valence: None,
484        };
485        let json = serde_json::to_string(&record).unwrap();
486        let back: OutcomeRecord = serde_json::from_str(&json).unwrap();
487        assert_eq!(back.task_id, "task-1");
488        assert_eq!(back.tokens_used, 1500);
489    }
490}