sz_rust_workflow/observability/
event_bus.rs1use async_trait::async_trait;
5use chrono::{DateTime, Utc};
6use serde::{Deserialize, Serialize};
7use tokio::sync::broadcast;
8
9use crate::error::WorkflowResult;
10
11#[derive(Debug, Clone, Serialize, Deserialize)]
13#[serde(tag = "type", rename_all = "snake_case")]
14pub enum WorkflowEvent {
15 InstanceStarted {
16 instance_id: String,
17 flow_key: String,
18 initiator: String,
19 timestamp: DateTime<Utc>,
20 },
21 InstanceSuspended {
22 instance_id: String,
23 actor: String,
24 timestamp: DateTime<Utc>,
25 },
26 InstanceResumed {
27 instance_id: String,
28 actor: String,
29 timestamp: DateTime<Utc>,
30 },
31 InstanceTerminated {
32 instance_id: String,
33 actor: String,
34 timestamp: DateTime<Utc>,
35 },
36 InstanceCompleted {
37 instance_id: String,
38 timestamp: DateTime<Utc>,
39 },
40 InstanceWithdrawn {
41 instance_id: String,
42 actor: String,
43 timestamp: DateTime<Utc>,
44 },
45 TransitionFired {
46 instance_id: String,
47 from: String,
48 to: String,
49 event: String,
50 timestamp: DateTime<Utc>,
51 },
52 NodeEntered {
53 instance_id: String,
54 node_id: String,
55 timestamp: DateTime<Utc>,
56 },
57 NodeLeft {
58 instance_id: String,
59 node_id: String,
60 timestamp: DateTime<Utc>,
61 },
62 TaskCreated {
63 instance_id: String,
64 task_id: String,
65 node_id: String,
66 candidates: Vec<String>,
67 timestamp: DateTime<Utc>,
68 },
69 TaskHandled {
70 instance_id: String,
71 task_id: String,
72 actor: String,
73 action: String,
74 timestamp: DateTime<Utc>,
75 },
76 PluginNodeCompleted {
77 instance_id: String,
78 node_id: String,
79 capability_name: String,
80 timestamp: DateTime<Utc>,
81 },
82}
83
84impl WorkflowEvent {
85 pub fn instance_id(&self) -> &str {
86 match self {
87 Self::InstanceStarted { instance_id, .. }
88 | Self::InstanceSuspended { instance_id, .. }
89 | Self::InstanceResumed { instance_id, .. }
90 | Self::InstanceTerminated { instance_id, .. }
91 | Self::InstanceCompleted { instance_id, .. }
92 | Self::InstanceWithdrawn { instance_id, .. }
93 | Self::TransitionFired { instance_id, .. }
94 | Self::NodeEntered { instance_id, .. }
95 | Self::NodeLeft { instance_id, .. }
96 | Self::TaskCreated { instance_id, .. }
97 | Self::TaskHandled { instance_id, .. }
98 | Self::PluginNodeCompleted { instance_id, .. } => instance_id,
99 }
100 }
101
102 pub fn timestamp(&self) -> DateTime<Utc> {
103 match self {
104 Self::InstanceStarted { timestamp, .. }
105 | Self::InstanceSuspended { timestamp, .. }
106 | Self::InstanceResumed { timestamp, .. }
107 | Self::InstanceTerminated { timestamp, .. }
108 | Self::InstanceCompleted { timestamp, .. }
109 | Self::InstanceWithdrawn { timestamp, .. }
110 | Self::TransitionFired { timestamp, .. }
111 | Self::NodeEntered { timestamp, .. }
112 | Self::NodeLeft { timestamp, .. }
113 | Self::TaskCreated { timestamp, .. }
114 | Self::TaskHandled { timestamp, .. }
115 | Self::PluginNodeCompleted { timestamp, .. } => *timestamp,
116 }
117 }
118}
119
120#[async_trait]
122pub trait WorkflowEventBus: Send + Sync + 'static {
123 async fn publish(&self, event: WorkflowEvent) -> WorkflowResult<()>;
124}
125
126pub struct InMemoryEventBus {
128 sender: broadcast::Sender<WorkflowEvent>,
129}
130
131impl InMemoryEventBus {
132 pub fn new(capacity: usize) -> Self {
133 let (sender, _) = broadcast::channel(capacity.max(1));
134 Self { sender }
135 }
136
137 pub fn subscribe(&self) -> broadcast::Receiver<WorkflowEvent> {
139 self.sender.subscribe()
140 }
141}
142
143impl Default for InMemoryEventBus {
144 fn default() -> Self {
145 Self::new(256)
146 }
147}
148
149#[async_trait]
150impl WorkflowEventBus for InMemoryEventBus {
151 async fn publish(&self, event: WorkflowEvent) -> WorkflowResult<()> {
152 let _ = self.sender.send(event);
153 Ok(())
154 }
155}
156
157pub struct NoopEventBus;
159
160#[async_trait]
161impl WorkflowEventBus for NoopEventBus {
162 async fn publish(&self, _event: WorkflowEvent) -> WorkflowResult<()> {
163 Ok(())
164 }
165}
166
167#[cfg(test)]
168mod tests {
169 use super::*;
170
171 #[tokio::test]
172 async fn in_memory_event_bus_publish_subscribe() {
173 let bus = InMemoryEventBus::new(16);
174 let mut rx = bus.subscribe();
175 let event = WorkflowEvent::InstanceStarted {
176 instance_id: "i1".into(),
177 flow_key: "test".into(),
178 initiator: "u1".into(),
179 timestamp: Utc::now(),
180 };
181 bus.publish(event.clone()).await.unwrap();
182 let received = rx.recv().await.unwrap();
183 assert_eq!(received.instance_id(), "i1");
184 }
185
186 #[tokio::test]
187 async fn noop_event_bus() {
188 let bus = NoopEventBus;
189 let event = WorkflowEvent::InstanceStarted {
190 instance_id: "i1".into(),
191 flow_key: "test".into(),
192 initiator: "u1".into(),
193 timestamp: Utc::now(),
194 };
195 bus.publish(event).await.unwrap();
196 }
197
198 #[test]
199 fn event_instance_id_and_timestamp() {
200 let ts = Utc::now();
201 let event = WorkflowEvent::TransitionFired {
202 instance_id: "i1".into(),
203 from: "draft".into(),
204 to: "review".into(),
205 event: "submit".into(),
206 timestamp: ts,
207 };
208 assert_eq!(event.instance_id(), "i1");
209 assert_eq!(event.timestamp(), ts);
210 }
211}