use std::sync::atomic::{AtomicU64, Ordering};
use serde_json::{Map, Value};
static EVENT_SEQ: AtomicU64 = AtomicU64::new(1);
#[derive(Debug, Clone)]
pub struct Event {
pub id: String,
pub payload: Map<String, Value>,
pub metadata: Map<String, Value>,
}
impl Event {
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(),
}
}
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)
}
}
}
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(),
}
}
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,
}
}
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)
}
}
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)
}