use std::cell::RefCell;
use std::sync::OnceLock;
use serde::{Deserialize, Serialize};
use serde_json::Value;
use time::OffsetDateTime;
use crate::EVENT_SCHEMA_VERSION;
use crate::ids::{EventId, RunId};
static PROCESS_ORIGIN: OnceLock<Option<String>> = OnceLock::new();
thread_local! {
static THREAD_ORIGIN: RefCell<Option<String>> = const { RefCell::new(None) };
}
pub fn set_process_origin(origin: impl Into<String>) {
let s = origin.into();
let value = if s.trim().is_empty() { None } else { Some(s) };
let _ = PROCESS_ORIGIN.set(value);
}
pub fn process_origin() -> Option<String> {
if let Some(o) = THREAD_ORIGIN.with(|c| c.borrow().clone()) {
return Some(o);
}
PROCESS_ORIGIN.get().cloned().flatten()
}
#[must_use]
pub struct OriginScope {
prev: Option<String>,
}
impl OriginScope {
pub fn new(origin: impl Into<String>) -> Self {
let s = origin.into();
let value = if s.trim().is_empty() { None } else { Some(s) };
let prev = THREAD_ORIGIN.with(|c| c.replace(value));
OriginScope { prev }
}
}
impl Drop for OriginScope {
fn drop(&mut self) {
let prev = self.prev.take();
THREAD_ORIGIN.with(|c| *c.borrow_mut() = prev);
}
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct Event {
pub event_id: EventId,
pub run_id: RunId,
#[serde(with = "time::serde::rfc3339")]
pub ts: OffsetDateTime,
pub parent_event_id: Option<EventId>,
pub kind: String,
pub schema_version: u32,
pub payload: Value,
#[serde(default)]
pub origin: Option<String>,
#[serde(default)]
pub hlc: Option<String>,
}
impl Event {
pub fn new(run_id: RunId, kind: impl Into<String>, payload: Value) -> Self {
Self {
event_id: EventId::new(),
run_id,
ts: OffsetDateTime::now_utc(),
parent_event_id: None,
kind: kind.into(),
schema_version: EVENT_SCHEMA_VERSION,
payload,
origin: process_origin(),
hlc: Some(crate::clock::now().to_canonical()),
}
}
pub fn with_parent(mut self, parent_event_id: EventId) -> Self {
self.parent_event_id = Some(parent_event_id);
self
}
pub fn with_origin(mut self, origin: Option<String>) -> Self {
self.origin = origin;
self
}
pub fn with_hlc(mut self, hlc: Option<String>) -> Self {
self.hlc = hlc;
self
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn origin_scope_overrides_and_restores() {
assert_eq!(process_origin(), None);
{
let _s = OriginScope::new("srv1/user:alice");
assert_eq!(process_origin().as_deref(), Some("srv1/user:alice"));
let e = Event::new(RunId::new(), "memory.cited", serde_json::json!({}));
assert_eq!(e.origin.as_deref(), Some("srv1/user:alice"));
{
let _inner = OriginScope::new("srv1/user:bob");
assert_eq!(process_origin().as_deref(), Some("srv1/user:bob"));
}
assert_eq!(process_origin().as_deref(), Some("srv1/user:alice"));
}
assert_eq!(process_origin(), None);
}
#[test]
fn origin_scope_empty_is_no_override() {
let _s = OriginScope::new("");
assert_eq!(process_origin(), None);
}
}