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}