Skip to main content

af_workflow/
event.rs

1//! Events flowing through a branch.
2//!
3//! Port of `platform/types.py::Event`. An `Event` is immutable: a node never
4//! mutates one in place; it derives a new event via [`Event::with_payload`] /
5//! [`Event::with_metadata`] (same id). One ingress emission = one event = one
6//! workflow run.
7
8use std::sync::atomic::{AtomicU64, Ordering};
9
10use serde_json::{Map, Value};
11
12static EVENT_SEQ: AtomicU64 = AtomicU64::new(1);
13
14/// An immutable unit of work moving through a branch's step chain.
15#[derive(Debug, Clone)]
16pub struct Event {
17    /// Stable id, minted fresh per ingress emission.
18    pub id: String,
19    /// Trigger data (the ingress payload: cron tick, bus message, feed item, …).
20    pub payload: Map<String, Value>,
21    /// Platform/trace metadata. Nodes add via [`Event::with_metadata`].
22    pub metadata: Map<String, Value>,
23}
24
25impl Event {
26    /// Mint a new event with a fresh id from a payload object.
27    pub fn new(payload: Map<String, Value>) -> Self {
28        let seq = EVENT_SEQ.fetch_add(1, Ordering::Relaxed);
29        Self {
30            id: format!("evt_{seq:016x}"),
31            payload,
32            metadata: Map::new(),
33        }
34    }
35
36    /// Convenience constructor from any JSON value; non-object payloads are
37    /// wrapped as `{ "value": <v> }` so the payload is always an object.
38    pub fn from_json(value: Value) -> Self {
39        match value {
40            Value::Object(map) => Self::new(map),
41            other => {
42                let mut m = Map::new();
43                m.insert("value".to_string(), other);
44                Self::new(m)
45            }
46        }
47    }
48
49    /// Return a copy with `patch` merged into the payload. Same id.
50    pub fn with_payload(&self, patch: Map<String, Value>) -> Event {
51        let mut payload = self.payload.clone();
52        for (k, v) in patch {
53            payload.insert(k, v);
54        }
55        Event {
56            id: self.id.clone(),
57            payload,
58            metadata: self.metadata.clone(),
59        }
60    }
61
62    /// Return a copy with `patch` merged into the metadata. Same id.
63    pub fn with_metadata(&self, patch: Map<String, Value>) -> Event {
64        let mut metadata = self.metadata.clone();
65        for (k, v) in patch {
66            metadata.insert(k, v);
67        }
68        Event {
69            id: self.id.clone(),
70            payload: self.payload.clone(),
71            metadata,
72        }
73    }
74
75    /// Read a dotted path out of the payload (`a.b.c`), or `None` if absent.
76    pub fn payload_path(&self, path: &str) -> Option<&Value> {
77        let mut segs = path.split('.');
78        let mut cur = self.payload.get(segs.next()?)?;
79        for seg in segs {
80            cur = cur.get(seg)?;
81        }
82        Some(cur)
83    }
84}
85
86/// Read a dotted path out of a JSON value.
87pub(crate) fn read_path<'a>(value: &'a Value, path: &str) -> Option<&'a Value> {
88    let mut cur = value;
89    for seg in path.split('.') {
90        cur = cur.get(seg)?;
91    }
92    Some(cur)
93}