1use std::sync::atomic::{AtomicU64, Ordering};
9
10use serde_json::{Map, Value};
11
12static EVENT_SEQ: AtomicU64 = AtomicU64::new(1);
13
14#[derive(Debug, Clone)]
16pub struct Event {
17 pub id: String,
19 pub payload: Map<String, Value>,
21 pub metadata: Map<String, Value>,
23}
24
25impl Event {
26 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 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 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 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 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
86pub(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}