Skip to main content

ocre/
events.rs

1//! Structured events (Rails 8.1's `Rails.event`): named facts about what
2//! the app did, with a payload, for analytics, audit trails or a data
3//! warehouse.
4//!
5//! Every [`Ctx`](crate::Ctx) has an [`Events`] notifier,
6//! [`Ctx::events`](crate::Ctx::events). An event is logged at once (an
7//! `info` line `event <name> <payload>` carrying the request's log fields,
8//! so Workers Logs can query it) and handed to each [`Subscriber`] the app
9//! registered with [`subscribe`].
10//!
11//! ```no_run
12//! use axum::extract::State;
13//! use ocre::{Ctx, Result};
14//! use serde_json::json;
15//!
16//! async fn checkout(State(ctx): State<Ctx>) -> Result<&'static str> {
17//!     ctx.events().set_context("tenant", "acme"); // on every event of this request
18//!     ctx.events().notify("order.placed", json!({ "order_id": 42, "total": "19.99" }));
19//!     // Tags group events of one part of the code (Rails' `Rails.event.tagged`).
20//!     ctx.events().tagged("step", "payment").notify("payment.captured", json!({ "order_id": 42 }));
21//!     Ok("OK")
22//! }
23//! # let _ = checkout;
24//! ```
25//!
26//! # Subscribers
27//!
28//! Register them once per Worker instance, in the `start` event of
29//! `src/lib.rs`. A subscriber turns an [`Event`] into the HTTP request to
30//! send ([`Delivery`]); Ocre sends it with `fetch` after the handler, like
31//! error reports.
32//!
33//! # Free plan
34//!
35//! Logging costs nothing. Each delivery is one subrequest (50 per request
36//! on the free plan): batch on the receiving side, or keep events in the
37//! logs only.
38
39use std::sync::{Arc, Mutex, PoisonError};
40
41use serde::Serialize;
42use serde_json::{Map, Value};
43
44pub use crate::errors::Delivery;
45use crate::log::Logger;
46
47/// A structured event, from [`Events::notify`].
48///
49/// # Examples
50///
51/// ```
52/// let events = ocre::events::Events::default();
53/// events.notify("user.signed_up", serde_json::json!({ "user_id": 7 }));
54/// let event = &events.take()[0];
55/// assert_eq!((event.name.as_str(), &event.payload["user_id"]), ("user.signed_up", &serde_json::json!(7)));
56/// ```
57#[derive(Debug, Clone, PartialEq)]
58pub struct Event {
59    /// What happened, e.g. `order.placed`.
60    pub name: String,
61    /// Its data; a payload that is not a JSON object is kept under `value`.
62    pub payload: Map<String, Value>,
63    /// Tags of the [`Events`] that notified it ([`Events::tagged`]).
64    pub tags: Map<String, Value>,
65    /// The request's or job's context ([`Events::set_context`]).
66    pub context: Map<String, Value>,
67    /// When it happened, in Unix seconds.
68    pub timestamp: i64,
69}
70
71/// A destination for events (Rails' event subscribers): turns an [`Event`] into the request to send.
72///
73/// `vars` looks up a Worker variable or secret by name. Return `None` to
74/// send nothing, e.g. for events it ignores.
75///
76/// # Examples
77///
78/// ```
79/// use ocre::events::{Delivery, Event, Subscriber};
80///
81/// /// Posts every `order.*` event to an analytics endpoint.
82/// struct Analytics;
83///
84/// impl Subscriber for Analytics {
85///     fn name(&self) -> &'static str {
86///         "analytics"
87///     }
88///
89///     fn emit(&self, event: &Event, vars: &dyn Fn(&str) -> Option<String>) -> Option<Delivery> {
90///         if !event.name.starts_with("order.") {
91///             return None;
92///         }
93///         let body = serde_json::json!({ "event": event.name, "properties": event.payload }).to_string();
94///         Some(Delivery { url: vars("ANALYTICS_URL")?, headers: vec![("content-type".into(), "application/json".into())], body })
95///     }
96/// }
97///
98/// ocre::events::subscribe(Analytics);
99/// ```
100pub trait Subscriber: Send + Sync {
101    /// A short name, for failure logs: `analytics`.
102    fn name(&self) -> &'static str;
103
104    /// The request to send for `event`, if any.
105    fn emit(&self, event: &Event, vars: &dyn Fn(&str) -> Option<String>) -> Option<Delivery>;
106}
107
108static SUBSCRIBERS: Mutex<Vec<Arc<dyn Subscriber>>> = Mutex::new(Vec::new());
109
110/// Registers a subscriber for every event of this Worker instance (Rails' `Rails.event.subscribe`).
111///
112/// Call it from the `start` event, which runs once per instance; a second
113/// subscriber with the same [`name`](Subscriber::name) replaces the first.
114///
115/// # Examples
116///
117/// ```
118/// # struct Audit;
119/// # impl ocre::events::Subscriber for Audit {
120/// #     fn name(&self) -> &'static str { "audit" }
121/// #     fn emit(&self, _: &ocre::events::Event, _: &dyn Fn(&str) -> Option<String>) -> Option<ocre::events::Delivery> { None }
122/// # }
123/// ocre::events::subscribe(Audit);
124/// ```
125pub fn subscribe(subscriber: impl Subscriber + 'static) {
126    let mut subscribers = SUBSCRIBERS.lock().unwrap_or_else(PoisonError::into_inner);
127    subscribers.retain(|existing| existing.name() != subscriber.name());
128    subscribers.push(Arc::new(subscriber));
129}
130
131/// The deliveries of `events` for every registered subscriber, with the subscriber's name.
132pub(crate) fn deliveries(events: &[Event], vars: &dyn Fn(&str) -> Option<String>) -> Vec<(&'static str, Delivery)> {
133    let subscribers = SUBSCRIBERS.lock().unwrap_or_else(PoisonError::into_inner).clone();
134    let mut out = Vec::new();
135    for event in events {
136        for subscriber in &subscribers {
137            if let Some(delivery) = subscriber.emit(event, vars) {
138                out.push((subscriber.name(), delivery));
139            }
140        }
141    }
142    out
143}
144
145#[derive(Debug, Default)]
146struct State {
147    context: Map<String, Value>,
148    pending: Vec<Event>,
149}
150
151/// The event notifier of one request, job batch or cron run: [`Ctx::events`](crate::Ctx::events).
152///
153/// Cheap to clone; clones share their context and pending events.
154/// `Events::default()` builds one for unit tests, whose events are read back
155/// with [`take`](Self::take).
156///
157/// # Examples
158///
159/// ```
160/// let events = ocre::events::Events::default();
161/// events.notify("cache.miss", "posts/12");
162/// assert_eq!(events.take()[0].payload["value"], "posts/12");
163/// ```
164#[derive(Debug, Clone, Default)]
165pub struct Events {
166    state: Arc<Mutex<State>>,
167    tags: Map<String, Value>,
168    log: Logger,
169}
170
171impl Events {
172    /// A notifier whose log lines carry the fields of `log` (the request id...).
173    pub(crate) fn new(log: Logger) -> Self {
174        Self { state: Arc::default(), tags: Map::new(), log }
175    }
176
177    fn lock(&self) -> std::sync::MutexGuard<'_, State> {
178        self.state.lock().unwrap_or_else(PoisonError::into_inner)
179    }
180
181    /// Records an event (Rails' `Rails.event.notify`): logs it and queues it for the subscribers.
182    ///
183    /// # Examples
184    ///
185    /// ```
186    /// let events = ocre::events::Events::default();
187    /// events.notify("post.published", serde_json::json!({ "post_id": 3 }));
188    /// assert_eq!(events.take()[0].payload["post_id"], 3);
189    /// ```
190    pub fn notify(&self, name: &str, payload: impl Serialize) {
191        self.notify_value(name, serde_json::to_value(payload).unwrap_or(Value::Null));
192    }
193
194    fn notify_value(&self, name: &str, payload: Value) {
195        let payload = match payload {
196            Value::Object(map) => map,
197            value => Map::from_iter([("value".to_owned(), value)]),
198        };
199        let mut log = self.log.clone();
200        for (key, value) in &self.tags {
201            log = log.with(key, value);
202        }
203        log.info(format_args!("event {name} {}", Value::Object(payload.clone())));
204        let mut state = self.lock();
205        let context = state.context.clone();
206        let event = Event { name: name.to_owned(), payload, tags: self.tags.clone(), context, timestamp: crate::now() };
207        state.pending.push(event);
208    }
209
210    /// This notifier with a tag added to its events (Rails' `Rails.event.tagged`).
211    ///
212    /// # Examples
213    ///
214    /// ```
215    /// let events = ocre::events::Events::default();
216    /// events.tagged("importer", "csv").notify("row.skipped", serde_json::json!({ "line": 12 }));
217    /// assert_eq!(events.take()[0].tags["importer"], "csv");
218    /// ```
219    #[must_use]
220    pub fn tagged(&self, key: &str, value: impl Serialize) -> Self {
221        let mut tagged = self.clone();
222        tagged.tags.insert(key.to_owned(), serde_json::to_value(value).unwrap_or(Value::Null));
223        tagged
224    }
225
226    /// Adds context to every later event of this request or job (Rails' `Rails.event.set_context`).
227    ///
228    /// # Examples
229    ///
230    /// ```
231    /// let events = ocre::events::Events::default();
232    /// events.set_context("user_id", 7);
233    /// events.notify("search", serde_json::json!({ "query": "rust" }));
234    /// assert_eq!(events.take()[0].context["user_id"], 7);
235    /// ```
236    pub fn set_context(&self, key: &str, value: impl Serialize) {
237        self.lock().context.insert(key.to_owned(), serde_json::to_value(value).unwrap_or(Value::Null));
238    }
239
240    /// Takes the events not yet sent to the subscribers, oldest first.
241    ///
242    /// # Examples
243    ///
244    /// ```
245    /// let events = ocre::events::Events::default();
246    /// events.notify("a", ());
247    /// assert_eq!(events.take().len(), 1);
248    /// assert!(events.take().is_empty());
249    /// ```
250    pub fn take(&self) -> Vec<Event> {
251        std::mem::take(&mut self.lock().pending)
252    }
253}
254
255#[cfg(test)]
256#[path = "../tests/events.rs"]
257mod tests;