assay_workflow/
signals.rs1use anyhow::Result;
4
5use crate::ctx::{WorkflowCtx, timestamp_now};
6use crate::events::WorkflowBusEvent;
7use crate::store::WorkflowStore;
8use crate::types::*;
9
10impl<S: WorkflowStore> WorkflowCtx<S> {
11 pub async fn send_signal(
16 &self,
17 workflow_id: &str,
18 name: &str,
19 payload: Option<&str>,
20 ) -> Result<()> {
21 let now = timestamp_now();
22 let payload_value: serde_json::Value = payload
27 .and_then(|s| serde_json::from_str(s).ok())
28 .unwrap_or(serde_json::Value::Null);
29 let event_payload =
30 serde_json::json!({ "signal": name, "payload": payload_value }).to_string();
31
32 self.store
33 .deliver_signal(
34 &WorkflowSignal {
35 id: None,
36 workflow_id: workflow_id.to_string(),
37 name: name.to_string(),
38 payload: payload.map(String::from),
39 consumed: false,
40 received_at: now,
41 },
42 &event_payload,
43 )
44 .await?;
45
46 self.emit_needs_dispatch(workflow_id).await;
48
49 let ns = self
52 .store
53 .get_workflow(workflow_id)
54 .await?
55 .map(|w| w.namespace)
56 .unwrap_or_default();
57 self.emit(
58 &ns,
59 WorkflowBusEvent::SignalReceived {
60 workflow_id: workflow_id.to_string(),
61 signal_name: name.to_string(),
62 },
63 )
64 .await;
65
66 Ok(())
67 }
68
69 pub async fn get_events(&self, workflow_id: &str) -> Result<Vec<WorkflowEvent>> {
70 self.store.list_events(workflow_id).await
71 }
72
73 pub async fn get_events_page(
74 &self,
75 workflow_id: &str,
76 cursor: Option<i32>,
77 limit: i64,
78 descending: bool,
79 ) -> Result<Vec<WorkflowEvent>> {
80 self.store
81 .list_events_page(workflow_id, cursor, limit, descending)
82 .await
83 }
84}