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};