Skip to main content

rig_core/observe/
mod.rs

1//! Typed runtime and policy observations, separate from effect replay records.
2//! [`Witness`] receives facts and [`ObservationLog`] retains a bounded trace.
3//! Optional timestamps come from a host-supplied [`Clock`], not replay identity.
4//!
5//! ```
6//! use rig_core::observe::ObservationLog;
7//!
8//! let log = ObservationLog::with_capacity(128);
9//! log.finalize();
10//! assert!(log.trace().finalized);
11//! ```
12
13use std::{
14    sync::{Arc, Mutex, PoisonError},
15    time::Duration,
16};
17
18use serde::{Deserialize, Serialize};
19
20use crate::{
21    effect::{EffectFamily, EffectId, HandlerKey, Outcome},
22    error::ErrorReport,
23    wasm_compat::{WasmCompatSend, WasmCompatSync},
24};
25
26mod adapter;
27pub(crate) mod sse_tail;
28pub use adapter::{
29    AdapterAnalysis, AdapterContext, AdapterEnding, AdapterErrorBoundary, AdapterErrorEnvelope,
30    AdapterEvent, AdapterObservation, AdapterUsage, AdapterVerdict, ObservationSink,
31    diagnostic_url_secrets, scrub_diagnostic,
32};
33pub(crate) use adapter::{AdapterSlot, ObservedError, lenient_count};
34
35#[cfg(test)]
36mod tests;
37
38/// One observed fact.
39#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
40pub struct Observation {
41    /// Sink-assigned sequence position in observation order.
42    pub seq: u64,
43    /// What the fact is about.
44    pub subject: Subject,
45    /// Where in the pipeline it was seen.
46    pub stage: Stage,
47    /// Who owns the decision or the observation.
48    pub emitter: Emitter,
49    /// The fact.
50    pub action: Action,
51    /// When, as a host-owned monotonic elapsed duration; a measurement,
52    /// never a semantic field. `None` when the sink has no [`Clock`].
53    #[serde(default, skip_serializing_if = "Option::is_none")]
54    pub at: Option<Duration>,
55}
56
57impl Observation {
58    /// A fact with no sequence yet; the sink assigns one.
59    pub fn new(subject: Subject, stage: Stage, emitter: Emitter, action: Action) -> Self {
60        Self {
61            seq: 0,
62            subject,
63            stage,
64            emitter,
65            action,
66            at: None,
67        }
68    }
69}
70
71/// What an observation is about. Every field is optional because facts
72/// exist before an effect has an id (a gate decides a pending intent), and
73/// some have no effect at all (a run ending). Correlation never depends on
74/// a runtime handle: `scope` is the program's serde id, `order` the
75/// driver's dispatch order, `effect`/`parent` the record's ids.
76#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize)]
77pub struct Subject {
78    /// The program scope (the run or agent), as the record names it.
79    #[serde(default, skip_serializing_if = "Option::is_none")]
80    pub scope: Option<String>,
81    /// The driver's dispatch order for the effect, when it has one; stable
82    /// before an id is issued, so a pre-dispatch decision correlates.
83    #[serde(default, skip_serializing_if = "Option::is_none")]
84    pub order: Option<u64>,
85    /// The effect's id once issued.
86    #[serde(default, skip_serializing_if = "Option::is_none")]
87    pub effect: Option<EffectId>,
88    /// The dispatch this one was made from, when a handler made it.
89    #[serde(default, skip_serializing_if = "Option::is_none")]
90    pub parent: Option<EffectId>,
91    /// The key the effect is routed to.
92    #[serde(default, skip_serializing_if = "Option::is_none")]
93    pub key: Option<HandlerKey>,
94    /// The family of the effect, when known.
95    #[serde(default, skip_serializing_if = "Option::is_none")]
96    pub family: Option<EffectFamily>,
97}
98
99impl Subject {
100    /// A subject with nothing but a scope: a run-level fact.
101    pub fn scoped(scope: impl Into<String>) -> Self {
102        Self {
103            scope: Some(scope.into()),
104            ..Self::default()
105        }
106    }
107}
108
109/// Where in the pipeline a fact was seen. The bus's four sets, the agent
110/// runtime's sets folded into one, the handler side, and the host.
111#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
112#[serde(rename_all = "snake_case")]
113pub enum Stage {
114    /// Before dispatch: a policy held, released or denied an intent.
115    Gate,
116    /// The driver took or refused an intent.
117    Dispatch,
118    /// Handler-side layer decision on entry or exit.
119    Handler,
120    /// The driver landed what a handler produced.
121    Collect,
122    /// After the record: a policy replaced an answer.
123    Judge,
124    /// The agent runtime: run endings.
125    Runtime,
126    /// The application's own policies and state.
127    Host,
128}
129
130/// Stable emitter name and optional version. [`Self::unknown`] explicitly
131/// represents unavailable attribution; it is not inferred from the outcome.
132#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
133pub struct Emitter {
134    /// The emitter's stable name (`rig-ecs/bus`, a layer's name, a host
135    /// system's name).
136    pub name: String,
137    /// Its declared version, when it has one.
138    #[serde(default, skip_serializing_if = "Option::is_none")]
139    pub version: Option<String>,
140}
141
142impl Emitter {
143    /// A named, unversioned emitter.
144    pub fn named(name: impl Into<String>) -> Self {
145        Self {
146            name: name.into(),
147            version: None,
148        }
149    }
150
151    /// A named, versioned emitter.
152    pub fn versioned(name: impl Into<String>, version: impl Into<String>) -> Self {
153        Self {
154            name: name.into(),
155            version: Some(version.into()),
156        }
157    }
158
159    /// The fact landed; the runtime does not know which policy made it.
160    /// The name `unknown` is reserved for this: a host emitter must name
161    /// itself otherwise.
162    pub fn unknown() -> Self {
163        Self::named("unknown")
164    }
165
166    /// Whether this is the explicit unknown emitter (by its reserved name).
167    pub fn is_unknown(&self) -> bool {
168        self.name == "unknown"
169    }
170}
171
172/// A structured reason: a stable code and an optional free-text detail.
173/// The code is what a comparison keys on; the detail is for a reader.
174#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
175pub struct Reason {
176    /// A stable, machine-readable code: an [`crate::error::ErrorKind::code`]
177    /// for a fact carrying a report, or an emitter's own (`intake_bound`,
178    /// `serial_key_busy`, `reentrant`, `ids_exhausted`, `layer_discarded`,
179    /// `despawned_before_dispatch`, `never_served`, `settled`, `max_turns`,
180    /// …).
181    pub code: String,
182    /// What a reader wants to know.
183    #[serde(default, skip_serializing_if = "Option::is_none")]
184    pub detail: Option<String>,
185}
186
187impl Reason {
188    /// A reason with a code and no detail.
189    pub fn code(code: impl Into<String>) -> Self {
190        Self {
191            code: code.into(),
192            detail: None,
193        }
194    }
195
196    /// A reason with a code and a detail.
197    pub fn with_detail(code: impl Into<String>, detail: impl Into<String>) -> Self {
198        Self {
199            code: code.into(),
200            detail: Some(detail.into()),
201        }
202    }
203
204    /// The reason an error report carries: its kind's stable code
205    /// ([`crate::error::ErrorKind::code`]), its message as the detail.
206    pub fn from_report(report: &ErrorReport) -> Self {
207        Self::with_detail(report.kind.code(), report.message.clone())
208    }
209
210    /// The explicit unknown reason.
211    pub fn unknown() -> Self {
212        Self::code("unknown")
213    }
214}
215
216/// The compact form of an outcome an observation carries: enough to
217/// classify without duplicating the exchange record that holds the value.
218#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
219#[serde(tag = "outcome", rename_all = "snake_case")]
220pub enum OutcomeSummary {
221    /// The handler answered with an outcome of this family.
222    Ok {
223        /// The answer's family.
224        family: EffectFamily,
225    },
226    /// The handler (or a decision) answered with this report.
227    Err {
228        /// The report's kind and message.
229        reason: Reason,
230        /// Whether the report says a retry may succeed.
231        retryable: bool,
232    },
233}
234
235impl OutcomeSummary {
236    /// The summary of an outcome.
237    pub fn of(outcome: &Result<Outcome, ErrorReport>) -> Self {
238        match outcome {
239            Ok(outcome) => Self::Ok {
240                family: outcome.family(),
241            },
242            Err(report) => Self::Err {
243                reason: Reason::from_report(report),
244                retryable: report.retryable,
245            },
246        }
247    }
248}
249
250/// The fact. Every variant carries its own before/after data or reason;
251/// a host adds its own through [`Action::Host`].
252#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
253#[serde(tag = "action", rename_all = "snake_case")]
254pub enum Action {
255    /// A fact emitted by the provider request boundary.
256    Adapter {
257        /// Correlation and typed boundary metadata.
258        observation: AdapterObservation,
259    },
260    /// A pending intent was held before dispatch.
261    Held {
262        /// Why, when the holder said.
263        reason: Reason,
264    },
265    /// A held intent was released to dispatch.
266    Released,
267    /// An intent was denied before any handler served it.
268    Denied {
269        /// The report the consumer receives.
270        reason: Reason,
271    },
272    /// The driver took an intent: the id it issued.
273    Issued,
274    /// The driver refused an intent before any handler: no record.
275    Refused {
276        /// Why: `handler_unavailable`, `reentrant` or `ids_exhausted`, with
277        /// the report's message as the detail.
278        reason: Reason,
279    },
280    /// An outcome landed and the record closed.
281    Landed {
282        /// The outcome, in brief.
283        outcome: OutcomeSummary,
284    },
285    /// A stream ended before its terminal record.
286    StreamTruncated {
287        /// Items the consumer had received.
288        delivered: usize,
289        /// What the last items seen were ([`StreamEvent::name`], or
290        /// `"Unknown"`), bounded, for after-the-fact classification.
291        ///
292        /// [`StreamEvent::name`]: crate::streaming::StreamEvent::name
293        tail: Vec<String>,
294        /// Error items seen in the stream, if any.
295        errors: Vec<Reason>,
296    },
297    /// An answer was replaced after the record closed.
298    Replaced {
299        /// What the record holds.
300        recorded: OutcomeSummary,
301        /// What the consumer received.
302        consumed: OutcomeSummary,
303    },
304    /// An in-flight dispatch was cancelled.
305    Cancelled {
306        /// Why.
307        reason: Reason,
308    },
309    /// A program ended.
310    Ended {
311        /// How (`settled`, `max_turns`, `provider`, `cancelled`, …).
312        ending: Reason,
313    },
314    /// A host policy's own fact: named by its kind, carried as its serde
315    /// payload. Build one through [`HostAction`].
316    Host {
317        /// The host action's declared kind.
318        kind: String,
319        /// Its payload.
320        payload: serde_json::Value,
321    },
322}
323
324/// Maximum serialized truncation-tail size enforced by [`Action::stream_truncated`].
325pub const LARGEST_PAYLOAD_BYTES: usize = 64 * 1024;
326
327impl Action {
328    /// A truncation observation whose `tail` is cut from the front until it
329    /// fits [`LARGEST_PAYLOAD_BYTES`]; `delivered` and `errors` are kept.
330    pub fn stream_truncated(delivered: usize, mut tail: Vec<String>, errors: Vec<Reason>) -> Self {
331        while !tail.is_empty()
332            && serde_json::to_vec(&tail).map_or(usize::MAX, |bytes| bytes.len())
333                > LARGEST_PAYLOAD_BYTES
334        {
335            tail.remove(0);
336        }
337        Self::StreamTruncated {
338            delivered,
339            tail,
340            errors,
341        }
342    }
343}
344
345/// A host-defined action: a named, serde-typed fact a host policy emits
346/// through [`Action::Host`]. The kind is declared once per type, so a
347/// trace names every host fact and a consumer deserializes it back.
348pub trait HostAction: Serialize + serde::de::DeserializeOwned {
349    /// The stable kind name (`rigcoder/approval`).
350    const KIND: &'static str;
351
352    /// This fact as an [`Action::Host`]; an unserializable fact is an error,
353    /// never a silent omission.
354    fn action(&self) -> Result<Action, serde_json::Error> {
355        Ok(Action::Host {
356            kind: Self::KIND.to_owned(),
357            payload: serde_json::to_value(self)?,
358        })
359    }
360
361    /// The fact back out of an [`Action::Host`] of this kind.
362    fn from_action(action: &Action) -> Option<Result<Self, serde_json::Error>> {
363        match action {
364            Action::Host { kind, payload } if kind == Self::KIND => {
365                Some(serde_json::from_value(payload.clone()))
366            }
367            _ => None,
368        }
369    }
370}
371
372/// A host-owned monotonic clock: what a sink stamps [`Observation::at`]
373/// from. rig-core reads no clock itself; a test supplies a counter.
374pub trait Clock: WasmCompatSend + WasmCompatSync {
375    /// Elapsed time since the clock's origin.
376    fn elapsed(&self) -> Duration;
377}
378
379/// Where observations go: the seam a driver and a host emit through. A
380/// witness is shared, so it takes `&self`; it must never block the caller.
381pub trait Witness: WasmCompatSend + WasmCompatSync + 'static {
382    /// One fact. The sink assigns its sequence.
383    fn observe(&self, observation: Observation);
384}
385
386/// The serializable trace a sink produces: the analysis artifact.
387#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
388pub struct ObservationTrace {
389    /// The sink's declared session, when it has one (a run's scope, a
390    /// trial id). Lineage, not semantics.
391    #[serde(default, skip_serializing_if = "Option::is_none")]
392    pub session: Option<String>,
393    /// The facts, in sequence.
394    pub observations: Vec<Observation>,
395    /// Facts discarded because sink capacity was reached. Nonzero means the
396    /// trace is incomplete.
397    #[serde(default)]
398    pub dropped: u64,
399    /// Whether the sink was told the session finished normally.
400    #[serde(default)]
401    pub finalized: bool,
402}
403
404impl ObservationTrace {
405    /// Whether every fact the session produced is here.
406    pub fn is_complete(&self) -> bool {
407        self.dropped == 0
408    }
409}
410
411/// The bounded in-memory sink, in one of two shapes. A capture
412/// ([`Self::with_capacity`], the default) keeps the first `capacity` facts
413/// and counts the rest as dropped: the trace is a complete prefix or says
414/// how much of the session it missed, which is what a comparison wants. A
415/// ring ([`Self::ring`]) keeps the *last* `capacity` facts and counts what
416/// it let go: a long-running host always sees its recent past, which is
417/// what a dashboard wants. Either way [`Self::drain`] takes what is kept
418/// and starts over, so a host that exports periodically never fills up.
419pub struct ObservationLog {
420    inner: Mutex<LogState>,
421    capacity: usize,
422    ring: bool,
423    clock: Option<Arc<dyn Clock + Send + Sync>>,
424}
425
426struct LogState {
427    session: Option<String>,
428    observations: std::collections::VecDeque<Observation>,
429    next: u64,
430    dropped: u64,
431    finalized: bool,
432}
433
434impl std::fmt::Debug for ObservationLog {
435    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
436        let state = self.lock();
437        f.debug_struct("ObservationLog")
438            .field("observations", &state.observations.len())
439            .field("dropped", &state.dropped)
440            .field("capacity", &self.capacity)
441            .finish_non_exhaustive()
442    }
443}
444
445/// The default capacity of an [`ObservationLog`].
446pub const DEFAULT_CAPACITY: usize = 65_536;
447
448impl Default for ObservationLog {
449    fn default() -> Self {
450        Self::with_capacity(DEFAULT_CAPACITY)
451    }
452}
453
454impl ObservationLog {
455    /// A capture keeping the first `capacity` facts; later ones are counted
456    /// as dropped.
457    pub fn with_capacity(capacity: usize) -> Self {
458        Self {
459            inner: Mutex::new(LogState {
460                session: None,
461                observations: std::collections::VecDeque::new(),
462                next: 0,
463                dropped: 0,
464                finalized: false,
465            }),
466            capacity,
467            ring: false,
468            clock: None,
469        }
470    }
471
472    /// Keeps the last `capacity` facts and counts evictions as dropped.
473    /// Capacity zero is raised to one.
474    pub fn ring(capacity: usize) -> Self {
475        Self {
476            ring: true,
477            // A ring of nothing would keep every fact and count it dropped.
478            ..Self::with_capacity(capacity.max(1))
479        }
480    }
481
482    /// Take the kept facts and start over: the session name and the
483    /// sequence continue, the kept facts and the dropped count reset. A
484    /// host that exports its trace periodically drains rather than
485    /// letting a capture fill and stop.
486    pub fn drain(&self) -> ObservationTrace {
487        let mut state = self.lock();
488        let trace = ObservationTrace {
489            session: state.session.clone(),
490            observations: state.observations.drain(..).collect(),
491            dropped: state.dropped,
492            finalized: state.finalized,
493        };
494        state.dropped = 0;
495        trace
496    }
497
498    /// Stamp every fact with `clock`'s elapsed time.
499    pub fn with_clock(mut self, clock: Arc<dyn Clock + Send + Sync>) -> Self {
500        self.clock = Some(clock);
501        self
502    }
503
504    /// Name the session the trace belongs to.
505    pub fn with_session(self, session: impl Into<String>) -> Self {
506        self.lock().session = Some(session.into());
507        self
508    }
509
510    fn lock(&self) -> std::sync::MutexGuard<'_, LogState> {
511        self.inner.lock().unwrap_or_else(PoisonError::into_inner)
512    }
513
514    /// The session finished normally: later readers know the trace is not
515    /// a partial artifact of a killed process. A later fact reopens the capture.
516    /// Hosts must drain their producers before exporting a finalized snapshot.
517    pub fn finalize(&self) {
518        self.lock().finalized = true;
519    }
520
521    /// The facts so far.
522    pub fn trace(&self) -> ObservationTrace {
523        let state = self.lock();
524        ObservationTrace {
525            session: state.session.clone(),
526            observations: state.observations.iter().cloned().collect(),
527            dropped: state.dropped,
528            finalized: state.finalized,
529        }
530    }
531
532    /// How many facts are kept.
533    pub fn len(&self) -> usize {
534        self.lock().observations.len()
535    }
536
537    /// Whether no observations are currently retained.
538    pub fn is_empty(&self) -> bool {
539        self.lock().observations.is_empty()
540    }
541}
542
543impl Witness for ObservationLog {
544    fn observe(&self, mut observation: Observation) {
545        let at = self.clock.as_ref().map(|clock| clock.elapsed());
546        let mut state = self.lock();
547        state.finalized = false;
548        observation.seq = state.next;
549        state.next += 1;
550        if state.observations.len() >= self.capacity {
551            state.dropped += 1;
552            if !self.ring {
553                return;
554            }
555            state.observations.pop_front();
556        }
557        observation.at = at;
558        state.observations.push_back(observation);
559    }
560}
561
562impl<W: Witness + ?Sized> Witness for Arc<W> {
563    fn observe(&self, observation: Observation) {
564        (**self).observe(observation);
565    }
566}
567
568// A trace serializes and crosses threads on every target.
569const _: fn() = || {
570    fn assert_wire<T: Clone + Send + Sync + 'static + Serialize + serde::de::DeserializeOwned>() {}
571    assert_wire::<Observation>();
572    assert_wire::<ObservationTrace>();
573};