Skip to main content

khive_storage/
telemetry.rs

1//! Typed payload structs for the ADR-094 lifecycle telemetry events.
2//!
3//! These are a documentation/test convenience only: the persisted
4//! discriminant on a stored [`crate::Event`] remains `khive_types::EventKind`,
5//! and each payload here is serialized into that event's JSON `payload`
6//! field by the emitting call site. Nothing here changes storage schema.
7
8use serde::{Deserialize, Serialize};
9
10/// Mirrors the eight ADR-094 lifecycle `EventKind` variants plus the three
11/// ADR-103 Stage 1 phase-span variants. Not itself persisted — see the
12/// module docs.
13#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
14#[serde(rename_all = "snake_case")]
15pub enum LifecycleEvent {
16    ChannelPollStarted,
17    ChannelPollSucceeded,
18    ChannelPollFailed,
19    ChannelBackoffArmed,
20    ChannelBackoffReset,
21    ChannelHeartbeatPersistFailed,
22    ConfigLocked,
23    CheckpointOutcomeRecorded,
24    PhaseStarted,
25    PhaseCompleted,
26    PhaseCancelled,
27}
28
29/// Payload for [`khive_types::EventKind::ChannelPollStarted`].
30#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
31pub struct ChannelPollStartedPayload {
32    pub channel_kind: String,
33    pub channel_slug: String,
34    pub since_rfc3339: String,
35}
36
37/// Payload for [`khive_types::EventKind::ChannelPollSucceeded`].
38#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
39pub struct ChannelPollSucceededPayload {
40    pub channel_kind: String,
41    pub channel_slug: String,
42    pub envelope_count: usize,
43    pub previous_backoff_attempt: u32,
44}
45
46/// Payload for [`khive_types::EventKind::ChannelPollFailed`].
47#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
48pub struct ChannelPollFailedPayload {
49    pub channel_kind: String,
50    pub channel_slug: String,
51    pub error_class: String,
52    pub error_message: String,
53}
54
55/// Payload for [`khive_types::EventKind::ChannelBackoffArmed`].
56#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
57pub struct ChannelBackoffArmedPayload {
58    pub channel_kind: String,
59    pub channel_slug: String,
60    pub attempt: u32,
61    pub step_ms: u64,
62    pub delay_ms: u64,
63}
64
65/// Payload for [`khive_types::EventKind::ChannelBackoffReset`].
66#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
67pub struct ChannelBackoffResetPayload {
68    pub channel_kind: String,
69    pub channel_slug: String,
70    pub previous_backoff_attempt: u32,
71}
72
73/// Payload for [`khive_types::EventKind::ChannelHeartbeatPersistFailed`].
74#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
75pub struct ChannelHeartbeatPersistFailedPayload {
76    pub channel_kind: String,
77    pub channel_slug: String,
78    pub error: String,
79}
80
81/// Payload for [`khive_types::EventKind::ConfigLocked`].
82#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
83pub struct ConfigLockedPayload {
84    pub key: String,
85    pub value: String,
86}
87
88/// Payload for [`khive_types::EventKind::CheckpointOutcomeRecorded`].
89#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
90pub struct CheckpointOutcomeRecordedPayload {
91    pub wal_pages: u64,
92    pub warn_pages: u64,
93    pub high_water_pages: u64,
94    pub truncate_high_water_pages: u64,
95    pub above_warn: bool,
96    pub above_high_water: bool,
97    pub above_truncate_high_water: bool,
98    /// Number of elevated observations aggregated into this episode so far.
99    /// Absent on rows written before the #1838 transition-summary contract.
100    #[serde(default, skip_serializing_if = "Option::is_none")]
101    pub episode_elevated_ticks: Option<u64>,
102    /// Highest WAL frame count observed during this episode so far. Absent on
103    /// rows written before the #1838 transition-summary contract.
104    #[serde(default, skip_serializing_if = "Option::is_none")]
105    pub episode_peak_wal_pages: Option<u64>,
106}
107
108/// Payload for [`khive_types::EventKind::PhaseStarted`] (ADR-103 Stage 1).
109///
110/// `work_class` is the closed ADR-103 enum (`interactive` | `warm` |
111/// `maintenance` | `inference`), carried as a plain string here since the
112/// enum itself is defined in a downstream crate; producers are responsible
113/// for using the closed set of values. `corpus_size` is populated only when
114/// it is cheaply known at phase start (e.g. a corpus count already on hand);
115/// `None` when unknown at this point.
116#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
117pub struct PhaseStartedPayload {
118    pub work_class: String,
119    pub phase: String,
120    pub corpus_size: Option<u64>,
121}
122
123/// Payload for [`khive_types::EventKind::PhaseCompleted`] (ADR-103 Stage 1).
124///
125/// `cpu_us` is a process-level `getrusage` delta across the phase (see
126/// `khive_runtime::resource`), not a per-thread measurement — `None` when
127/// the underlying read is unavailable on this platform.
128#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
129pub struct PhaseCompletedPayload {
130    pub work_class: String,
131    pub phase: String,
132    pub wall_us: i64,
133    pub cpu_us: Option<i64>,
134}
135
136/// Payload for [`khive_types::EventKind::PhaseCancelled`] (ADR-103 Stage 1).
137/// Same shape as [`PhaseCompletedPayload`] — the phase ran for `wall_us`
138/// before being cut short rather than returning a result.
139#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
140pub struct PhaseCancelledPayload {
141    pub work_class: String,
142    pub phase: String,
143    pub wall_us: i64,
144    pub cpu_us: Option<i64>,
145}
146
147#[cfg(test)]
148mod tests {
149    use super::*;
150
151    #[test]
152    fn lifecycle_event_roundtrips_through_json() {
153        for kind in [
154            LifecycleEvent::ChannelPollStarted,
155            LifecycleEvent::ChannelPollSucceeded,
156            LifecycleEvent::ChannelPollFailed,
157            LifecycleEvent::ChannelBackoffArmed,
158            LifecycleEvent::ChannelBackoffReset,
159            LifecycleEvent::ChannelHeartbeatPersistFailed,
160            LifecycleEvent::ConfigLocked,
161            LifecycleEvent::CheckpointOutcomeRecorded,
162            LifecycleEvent::PhaseStarted,
163            LifecycleEvent::PhaseCompleted,
164            LifecycleEvent::PhaseCancelled,
165        ] {
166            let json = serde_json::to_string(&kind).expect("serialize");
167            let parsed: LifecycleEvent = serde_json::from_str(&json).expect("deserialize");
168            assert_eq!(parsed, kind);
169        }
170    }
171
172    #[test]
173    fn checkpoint_outcome_recorded_payload_roundtrips() {
174        let payload = CheckpointOutcomeRecordedPayload {
175            wal_pages: 2500,
176            warn_pages: 2000,
177            high_water_pages: 6000,
178            truncate_high_water_pages: 20_000,
179            above_warn: true,
180            above_high_water: false,
181            above_truncate_high_water: false,
182            episode_elevated_ticks: Some(7),
183            episode_peak_wal_pages: Some(3100),
184        };
185        let json = serde_json::to_value(&payload).expect("serialize");
186        let parsed: CheckpointOutcomeRecordedPayload =
187            serde_json::from_value(json).expect("deserialize");
188        assert_eq!(parsed, payload);
189    }
190
191    #[test]
192    fn checkpoint_outcome_recorded_payload_accepts_legacy_rows_without_episode_summary() {
193        let parsed: CheckpointOutcomeRecordedPayload = serde_json::from_value(serde_json::json!({
194            "wal_pages": 2500,
195            "warn_pages": 2000,
196            "high_water_pages": 6000,
197            "truncate_high_water_pages": 20000,
198            "above_warn": true,
199            "above_high_water": false,
200            "above_truncate_high_water": false
201        }))
202        .expect("legacy payload must remain readable");
203
204        assert_eq!(parsed.episode_elevated_ticks, None);
205        assert_eq!(parsed.episode_peak_wal_pages, None);
206    }
207
208    #[test]
209    fn phase_started_payload_roundtrips() {
210        let payload = PhaseStartedPayload {
211            work_class: "warm".into(),
212            phase: "ann_warm".into(),
213            corpus_size: Some(553_000),
214        };
215        let json = serde_json::to_value(&payload).expect("serialize");
216        let parsed: PhaseStartedPayload = serde_json::from_value(json).expect("deserialize");
217        assert_eq!(parsed, payload);
218    }
219
220    #[test]
221    fn phase_completed_payload_roundtrips_with_absent_cpu_us() {
222        let payload = PhaseCompletedPayload {
223            work_class: "warm".into(),
224            phase: "ann_warm".into(),
225            wall_us: 41_000_000,
226            cpu_us: None,
227        };
228        let json = serde_json::to_value(&payload).expect("serialize");
229        let parsed: PhaseCompletedPayload = serde_json::from_value(json).expect("deserialize");
230        assert_eq!(parsed, payload);
231    }
232
233    #[test]
234    fn phase_cancelled_payload_roundtrips() {
235        let payload = PhaseCancelledPayload {
236            work_class: "warm".into(),
237            phase: "ann_warm".into(),
238            wall_us: 12_000,
239            cpu_us: Some(9_500),
240        };
241        let json = serde_json::to_value(&payload).expect("serialize");
242        let parsed: PhaseCancelledPayload = serde_json::from_value(json).expect("deserialize");
243        assert_eq!(parsed, payload);
244    }
245}