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