1use serde::{Deserialize, Serialize};
9
10#[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#[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#[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#[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#[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#[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#[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#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
83pub struct ConfigLockedPayload {
84 pub key: String,
85 pub value: String,
86}
87
88#[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 #[serde(default, skip_serializing_if = "Option::is_none")]
101 pub episode_elevated_ticks: Option<u64>,
102 #[serde(default, skip_serializing_if = "Option::is_none")]
105 pub episode_peak_wal_pages: Option<u64>,
106}
107
108#[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#[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#[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}