1use std::sync::atomic::{AtomicU64, Ordering};
2use std::time::{SystemTime, UNIX_EPOCH};
3
4use serde::{Deserialize, Serialize};
5use serde_json::Value;
6
7use crate::{EffectId, EventId, RunError, RunId, Usage};
8
9#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
11pub struct EventMeta {
12 pub event_id: EventId,
14 pub sequence: u64,
16 pub run_id: RunId,
18 pub parent_run_id: Option<RunId>,
20 pub caused_by: Option<EventId>,
22 pub timestamp_ms: u64,
24}
25
26#[derive(Clone, Debug, Deserialize, PartialEq, Serialize)]
28pub struct RunEvent {
29 pub meta: EventMeta,
31 pub kind: RunEventKind,
33}
34
35#[derive(Clone, Debug, Deserialize, PartialEq, Serialize)]
37#[non_exhaustive]
38pub enum RunEventKind {
39 Lifecycle(LifecycleEvent),
41 Effect(EffectEvent),
43 Child(ChildEvent),
45 Budget(BudgetEvent),
47 Domain(DomainEvent),
49}
50
51#[derive(Clone, Debug, Deserialize, PartialEq, Serialize)]
53#[non_exhaustive]
54pub enum LifecycleEvent {
55 Started,
57 Completed {
59 output: Value,
61 },
62 Failed {
64 error: RunError,
66 },
67 Cancelled,
69}
70
71#[derive(Clone, Debug, Deserialize, PartialEq, Serialize)]
73#[non_exhaustive]
74pub enum EffectEvent {
75 Requested {
77 effect_id: EffectId,
79 },
80 Started {
82 effect_id: EffectId,
84 },
85 Completed {
87 effect_id: EffectId,
89 output: Value,
91 },
92 Failed {
94 effect_id: EffectId,
96 error: RunError,
98 },
99}
100
101#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
103#[non_exhaustive]
104pub enum ChildEvent {
105 Started {
107 child_run_id: RunId,
109 },
110 Completed {
112 child_run_id: RunId,
114 },
115 Failed {
117 child_run_id: RunId,
119 },
120 Cancelled {
122 child_run_id: RunId,
124 },
125}
126
127#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
129#[non_exhaustive]
130pub enum BudgetEvent {
131 Updated {
133 usage: Usage,
135 },
136}
137
138#[derive(Clone, Debug, Deserialize, PartialEq, Serialize)]
140pub struct DomainEvent {
141 pub namespace: String,
143 pub name: String,
145 pub payload: Value,
147}
148
149#[derive(Debug)]
151pub struct EventFactory {
152 run_id: RunId,
153 parent_run_id: Option<RunId>,
154 sequence: AtomicU64,
155}
156
157impl EventFactory {
158 pub const fn new(run_id: RunId, parent_run_id: Option<RunId>) -> Self {
160 Self {
161 run_id,
162 parent_run_id,
163 sequence: AtomicU64::new(0),
164 }
165 }
166
167 pub fn emit(&self, kind: RunEventKind, caused_by: Option<EventId>) -> RunEvent {
169 let sequence = self.sequence.fetch_add(1, Ordering::Relaxed);
170 let timestamp_ms = SystemTime::now()
171 .duration_since(UNIX_EPOCH)
172 .map_or(0, |duration| {
173 u64::try_from(duration.as_millis()).unwrap_or(u64::MAX)
174 });
175
176 RunEvent {
177 meta: EventMeta {
178 event_id: EventId::new(),
179 sequence,
180 run_id: self.run_id,
181 parent_run_id: self.parent_run_id,
182 caused_by,
183 timestamp_ms,
184 },
185 kind,
186 }
187 }
188}
189
190#[cfg(test)]
191mod tests {
192 use super::{EventFactory, LifecycleEvent, RunEventKind};
193 use crate::RunId;
194
195 #[test]
196 fn event_sequences_are_monotonic_and_causal() {
197 let factory = EventFactory::new(RunId::new(), None);
198 let started = factory.emit(RunEventKind::Lifecycle(LifecycleEvent::Started), None);
199 let completed = factory.emit(
200 RunEventKind::Lifecycle(LifecycleEvent::Completed {
201 output: serde_json::json!("done"),
202 }),
203 Some(started.meta.event_id),
204 );
205
206 assert_eq!(started.meta.sequence, 0);
207 assert_eq!(completed.meta.sequence, 1);
208 assert_eq!(completed.meta.caused_by, Some(started.meta.event_id));
209 }
210}