Skip to main content

sz_rust_workflow/observability/
event_bus.rs

1// SPDX-License-Identifier: Apache-2.0
2// Copyright (c) 2024-2026 SZ-Rust Team
3//
4use async_trait::async_trait;
5use chrono::{DateTime, Utc};
6use serde::{Deserialize, Serialize};
7use tokio::sync::broadcast;
8
9use crate::error::WorkflowResult;
10
11/// 工作流事件,对齐 design 2.3.2 事件模型。
12#[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/// 事件总线 trait。
121#[async_trait]
122pub trait WorkflowEventBus: Send + Sync + 'static {
123    async fn publish(&self, event: WorkflowEvent) -> WorkflowResult<()>;
124}
125
126/// InMemory 事件总线(broadcast channel)。
127pub 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    /// 订阅事件。
138    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
157/// Noop 事件总线(测试用)。
158pub 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}