Skip to main content

remem/db/query/observability/
types.rs

1use std::collections::BTreeMap;
2
3use serde::Serialize;
4
5pub const OBSERVABILITY_SCHEMA_VERSION: u32 = 1;
6pub const CURRENT_MEMORY_CONTRACT_SPEC_PATH: &str = "docs/specs/current-memory-contracts/TECH.md";
7
8#[derive(Debug, Clone, PartialEq, Serialize)]
9pub struct ObservabilityReport {
10    pub schema_version: u32,
11    pub generated_at_epoch: i64,
12    pub spec_path: &'static str,
13    pub checks: Vec<ObservabilityCheck>,
14    pub metrics: ObservabilityMetrics,
15}
16
17#[derive(Debug, Clone, PartialEq, Eq, Serialize)]
18pub struct ObservabilityCheck {
19    pub code: String,
20    pub severity: &'static str,
21    pub scope: &'static str,
22    pub message: String,
23    pub metrics: BTreeMap<String, i64>,
24    pub actions: Vec<String>,
25}
26
27#[derive(Debug, Clone, Default, PartialEq, Serialize)]
28pub struct ObservabilityMetrics {
29    pub capture: CaptureObservabilityMetrics,
30    pub promotion: PromotionObservabilityMetrics,
31    pub context_injection: ContextInjectionObservabilityMetrics,
32    pub usage_feedback: UsageFeedbackObservabilityMetrics,
33    pub temporal_facts: TemporalFactObservabilityMetrics,
34    pub staleness: StalenessObservabilityMetrics,
35    pub queue: QueueObservabilityMetrics,
36    pub worker: WorkerObservabilityMetrics,
37}
38
39#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize)]
40pub struct CaptureObservabilityMetrics {
41    pub captured_events: i64,
42    pub capture_drop_events: i64,
43    pub actionable_capture_drops: i64,
44    pub unrecovered_capture_spills: i64,
45    pub pending_extraction_tasks: i64,
46    pub processing_extraction_tasks: i64,
47    pub expired_processing_extraction_tasks: i64,
48    pub failed_extraction_tasks: i64,
49}
50
51#[derive(Debug, Clone, Default, PartialEq, Serialize)]
52pub struct PromotionObservabilityMetrics {
53    pub observations: i64,
54    pub candidates: i64,
55    pub promoted: i64,
56    pub pending_review: i64,
57    pub candidate_rate_percent: f64,
58    pub promoted_rate_percent: f64,
59}
60
61#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize)]
62pub struct ContextInjectionObservabilityMetrics {
63    pub output_table_exists: bool,
64    pub item_table_exists: bool,
65    pub output_rows: i64,
66    pub output_emit_count: i64,
67    pub output_suppress_count: i64,
68    pub output_modes: Vec<CountBucket>,
69    pub item_rows: i64,
70    pub item_statuses: Vec<CountBucket>,
71    pub item_channels: Vec<CountBucket>,
72    pub item_drop_reasons: Vec<CountBucket>,
73    pub item_staleness_source_anchors: Vec<CountBucket>,
74    pub item_staleness_ages: Vec<CountBucket>,
75}
76
77#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize)]
78pub struct UsageFeedbackObservabilityMetrics {
79    pub citation_table_exists: bool,
80    pub usage_table_exists: bool,
81    pub citation_events: i64,
82    pub citation_line_present_events: i64,
83    pub matched_events: i64,
84    pub inserted_events: i64,
85    pub no_citation_events: i64,
86    pub unmatched_events: i64,
87    pub usage_events: i64,
88}
89
90#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize)]
91pub struct TemporalFactObservabilityMetrics {
92    pub table_exists: bool,
93    pub total_rows: i64,
94    pub retrieval_eligible_rows: i64,
95    pub invalidated_rows: i64,
96    pub expired_rows: i64,
97    pub orphan_source_memory_rows: i64,
98    pub unlinked_source_rows: i64,
99}
100
101#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize)]
102pub struct StalenessObservabilityMetrics {
103    pub memory_table_exists: bool,
104    pub total_memories: i64,
105    pub source_anchors: Vec<CountBucket>,
106    pub ages: Vec<CountBucket>,
107    pub error_count: i64,
108}
109
110#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize)]
111pub struct QueueObservabilityMetrics {
112    pub pending_observations: i64,
113    pub ready_pending_observations: i64,
114    pub delayed_pending_observations: i64,
115    pub processing_pending_observations: i64,
116    pub expired_processing_pending_observations: i64,
117    pub failed_pending_observations: i64,
118    pub pending_jobs: i64,
119    pub processing_jobs: i64,
120    pub failed_jobs: i64,
121    pub stuck_jobs: i64,
122    pub retryable_extraction_replay_ranges: i64,
123}
124
125#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize)]
126pub struct WorkerObservabilityMetrics {
127    pub daemon_healthy: bool,
128    pub heartbeat_age_secs: Option<i64>,
129    pub heartbeat_owner_present: bool,
130}
131
132#[derive(Debug, Clone, PartialEq, Eq, Serialize)]
133pub struct CountBucket {
134    pub value: String,
135    pub count: i64,
136}
137
138impl ObservabilityReport {
139    pub fn unavailable(generated_at_epoch: i64, reason: impl Into<String>) -> Self {
140        Self {
141            schema_version: OBSERVABILITY_SCHEMA_VERSION,
142            generated_at_epoch,
143            spec_path: CURRENT_MEMORY_CONTRACT_SPEC_PATH,
144            checks: vec![ObservabilityCheck::new(
145                "observability_database_unavailable",
146                "warn",
147                "database",
148                reason,
149            )
150            .action("open the remem database before requesting runtime observability")],
151            metrics: ObservabilityMetrics::default(),
152        }
153    }
154}
155
156impl ObservabilityCheck {
157    pub fn new(
158        code: impl Into<String>,
159        severity: &'static str,
160        scope: &'static str,
161        message: impl Into<String>,
162    ) -> Self {
163        Self {
164            code: code.into(),
165            severity,
166            scope,
167            message: message.into(),
168            metrics: BTreeMap::new(),
169            actions: Vec::new(),
170        }
171    }
172
173    pub fn metric(mut self, key: &'static str, value: i64) -> Self {
174        self.metrics.insert(key.to_string(), value);
175        self
176    }
177
178    pub fn action(mut self, action: impl Into<String>) -> Self {
179        self.actions.push(action.into());
180        self
181    }
182}