rig-core 0.44.0

An opinionated library for building LLM powered applications.
Documentation
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
375
376
377
378
379
380
381
382
383
384
385
386
387
388
389
390
391
392
393
394
395
396
397
398
399
400
401
402
403
404
405
406
407
408
409
410
411
412
413
414
415
416
417
418
419
420
421
422
423
424
425
426
427
428
429
430
431
432
433
434
435
436
437
438
439
440
441
442
443
444
445
446
447
448
449
450
451
452
453
454
455
456
457
458
459
460
461
462
463
464
465
466
467
468
469
470
471
472
473
474
475
476
477
478
479
480
481
482
483
484
485
486
487
488
489
490
491
492
493
494
495
496
497
498
499
500
501
502
503
504
505
506
507
508
509
510
511
512
513
514
515
516
517
518
519
520
521
522
523
524
525
526
527
528
529
530
531
532
533
534
535
536
537
538
539
540
541
542
543
544
545
546
547
548
549
550
551
552
553
554
555
556
557
558
559
560
561
562
563
564
565
566
567
568
569
570
571
572
573
//! Typed runtime and policy observations, separate from effect replay records.
//! [`Witness`] receives facts and [`ObservationLog`] retains a bounded trace.
//! Optional timestamps come from a host-supplied [`Clock`], not replay identity.
//!
//! ```
//! use rig_core::observe::ObservationLog;
//!
//! let log = ObservationLog::with_capacity(128);
//! log.finalize();
//! assert!(log.trace().finalized);
//! ```

use std::{
    sync::{Arc, Mutex, PoisonError},
    time::Duration,
};

use serde::{Deserialize, Serialize};

use crate::{
    effect::{EffectFamily, EffectId, HandlerKey, Outcome},
    error::ErrorReport,
    wasm_compat::{WasmCompatSend, WasmCompatSync},
};

mod adapter;
pub(crate) mod sse_tail;
pub use adapter::{
    AdapterAnalysis, AdapterContext, AdapterEnding, AdapterErrorBoundary, AdapterErrorEnvelope,
    AdapterEvent, AdapterObservation, AdapterUsage, AdapterVerdict, ObservationSink,
    diagnostic_url_secrets, scrub_diagnostic,
};
pub(crate) use adapter::{AdapterSlot, ObservedError};

#[cfg(test)]
mod tests;

/// One observed fact.
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct Observation {
    /// Sink-assigned sequence position in observation order.
    pub seq: u64,
    /// What the fact is about.
    pub subject: Subject,
    /// Where in the pipeline it was seen.
    pub stage: Stage,
    /// Who owns the decision or the observation.
    pub emitter: Emitter,
    /// The fact.
    pub action: Action,
    /// When, as a host-owned monotonic elapsed duration; a measurement,
    /// never a semantic field. `None` when the sink has no [`Clock`].
    #[serde(default, skip_serializing_if = "Option::is_none")]
    pub at: Option<Duration>,
}

impl Observation {
    /// A fact with no sequence yet; the sink assigns one.
    pub fn new(subject: Subject, stage: Stage, emitter: Emitter, action: Action) -> Self {
        Self {
            seq: 0,
            subject,
            stage,
            emitter,
            action,
            at: None,
        }
    }
}

/// What an observation is about. Every field is optional because facts
/// exist before an effect has an id (a gate decides a pending intent), and
/// some have no effect at all (a run ending). Correlation never depends on
/// a runtime handle: `scope` is the program's serde id, `order` the
/// driver's dispatch order, `effect`/`parent` the record's ids.
#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize)]
pub struct Subject {
    /// The program scope (the run or agent), as the record names it.
    #[serde(default, skip_serializing_if = "Option::is_none")]
    pub scope: Option<String>,
    /// The driver's dispatch order for the effect, when it has one; stable
    /// before an id is issued, so a pre-dispatch decision correlates.
    #[serde(default, skip_serializing_if = "Option::is_none")]
    pub order: Option<u64>,
    /// The effect's id once issued.
    #[serde(default, skip_serializing_if = "Option::is_none")]
    pub effect: Option<EffectId>,
    /// The dispatch this one was made from, when a handler made it.
    #[serde(default, skip_serializing_if = "Option::is_none")]
    pub parent: Option<EffectId>,
    /// The key the effect is routed to.
    #[serde(default, skip_serializing_if = "Option::is_none")]
    pub key: Option<HandlerKey>,
    /// The family of the effect, when known.
    #[serde(default, skip_serializing_if = "Option::is_none")]
    pub family: Option<EffectFamily>,
}

impl Subject {
    /// A subject with nothing but a scope: a run-level fact.
    pub fn scoped(scope: impl Into<String>) -> Self {
        Self {
            scope: Some(scope.into()),
            ..Self::default()
        }
    }
}

/// Where in the pipeline a fact was seen. The bus's four sets, the agent
/// runtime's sets folded into one, the handler side, and the host.
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum Stage {
    /// Before dispatch: a policy held, released or denied an intent.
    Gate,
    /// The driver took or refused an intent.
    Dispatch,
    /// Handler-side layer decision on entry or exit.
    Handler,
    /// The driver landed what a handler produced.
    Collect,
    /// After the record: a policy replaced an answer.
    Judge,
    /// The agent runtime: run endings.
    Runtime,
    /// The application's own policies and state.
    Host,
}

/// Stable emitter name and optional version. [`Self::unknown`] explicitly
/// represents unavailable attribution; it is not inferred from the outcome.
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct Emitter {
    /// The emitter's stable name (`rig-ecs/bus`, a layer's name, a host
    /// system's name).
    pub name: String,
    /// Its declared version, when it has one.
    #[serde(default, skip_serializing_if = "Option::is_none")]
    pub version: Option<String>,
}

impl Emitter {
    /// A named, unversioned emitter.
    pub fn named(name: impl Into<String>) -> Self {
        Self {
            name: name.into(),
            version: None,
        }
    }

    /// A named, versioned emitter.
    pub fn versioned(name: impl Into<String>, version: impl Into<String>) -> Self {
        Self {
            name: name.into(),
            version: Some(version.into()),
        }
    }

    /// The fact landed; the runtime does not know which policy made it.
    /// The name `unknown` is reserved for this: a host emitter must name
    /// itself otherwise.
    pub fn unknown() -> Self {
        Self::named("unknown")
    }

    /// Whether this is the explicit unknown emitter (by its reserved name).
    pub fn is_unknown(&self) -> bool {
        self.name == "unknown"
    }
}

/// A structured reason: a stable code and an optional free-text detail.
/// The code is what a comparison keys on; the detail is for a reader.
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct Reason {
    /// A stable, machine-readable code: an [`crate::error::ErrorKind::code`]
    /// for a fact carrying a report, or an emitter's own (`intake_bound`,
    /// `serial_key_busy`, `reentrant`, `ids_exhausted`, `layer_discarded`,
    /// `despawned_before_dispatch`, `never_served`, `settled`, `max_turns`,
    /// …).
    pub code: String,
    /// What a reader wants to know.
    #[serde(default, skip_serializing_if = "Option::is_none")]
    pub detail: Option<String>,
}

impl Reason {
    /// A reason with a code and no detail.
    pub fn code(code: impl Into<String>) -> Self {
        Self {
            code: code.into(),
            detail: None,
        }
    }

    /// A reason with a code and a detail.
    pub fn with_detail(code: impl Into<String>, detail: impl Into<String>) -> Self {
        Self {
            code: code.into(),
            detail: Some(detail.into()),
        }
    }

    /// The reason an error report carries: its kind's stable code
    /// ([`crate::error::ErrorKind::code`]), its message as the detail.
    pub fn from_report(report: &ErrorReport) -> Self {
        Self::with_detail(report.kind.code(), report.message.clone())
    }

    /// The explicit unknown reason.
    pub fn unknown() -> Self {
        Self::code("unknown")
    }
}

/// The compact form of an outcome an observation carries: enough to
/// classify without duplicating the exchange record that holds the value.
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
#[serde(tag = "outcome", rename_all = "snake_case")]
pub enum OutcomeSummary {
    /// The handler answered with an outcome of this family.
    Ok {
        /// The answer's family.
        family: EffectFamily,
    },
    /// The handler (or a decision) answered with this report.
    Err {
        /// The report's kind and message.
        reason: Reason,
        /// Whether the report says a retry may succeed.
        retryable: bool,
    },
}

impl OutcomeSummary {
    /// The summary of an outcome.
    pub fn of(outcome: &Result<Outcome, ErrorReport>) -> Self {
        match outcome {
            Ok(outcome) => Self::Ok {
                family: outcome.family(),
            },
            Err(report) => Self::Err {
                reason: Reason::from_report(report),
                retryable: report.retryable,
            },
        }
    }
}

/// The fact. Every variant carries its own before/after data or reason;
/// a host adds its own through [`Action::Host`].
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
#[serde(tag = "action", rename_all = "snake_case")]
pub enum Action {
    /// A fact emitted by the provider request boundary.
    Adapter {
        /// Correlation and typed boundary metadata.
        observation: AdapterObservation,
    },
    /// A pending intent was held before dispatch.
    Held {
        /// Why, when the holder said.
        reason: Reason,
    },
    /// A held intent was released to dispatch.
    Released,
    /// An intent was denied before any handler served it.
    Denied {
        /// The report the consumer receives.
        reason: Reason,
    },
    /// The driver took an intent: the id it issued.
    Issued,
    /// The driver refused an intent before any handler: no record.
    Refused {
        /// Why: `handler_unavailable`, `reentrant` or `ids_exhausted`, with
        /// the report's message as the detail.
        reason: Reason,
    },
    /// An outcome landed and the record closed.
    Landed {
        /// The outcome, in brief.
        outcome: OutcomeSummary,
    },
    /// A stream ended before its terminal record.
    StreamTruncated {
        /// Items the consumer had received.
        delivered: usize,
        /// What the last items seen were ([`StreamEvent::name`], or
        /// `"Unknown"`), bounded, for after-the-fact classification.
        ///
        /// [`StreamEvent::name`]: crate::streaming::StreamEvent::name
        tail: Vec<String>,
        /// Error items seen in the stream, if any.
        errors: Vec<Reason>,
    },
    /// An answer was replaced after the record closed.
    Replaced {
        /// What the record holds.
        recorded: OutcomeSummary,
        /// What the consumer received.
        consumed: OutcomeSummary,
    },
    /// An in-flight dispatch was cancelled.
    Cancelled {
        /// Why.
        reason: Reason,
    },
    /// A program ended.
    Ended {
        /// How (`settled`, `max_turns`, `provider`, `cancelled`, …).
        ending: Reason,
    },
    /// A host policy's own fact: named by its kind, carried as its serde
    /// payload. Build one through [`HostAction`].
    Host {
        /// The host action's declared kind.
        kind: String,
        /// Its payload.
        payload: serde_json::Value,
    },
}

/// Maximum serialized truncation-tail size enforced by [`Action::stream_truncated`].
pub const LARGEST_PAYLOAD_BYTES: usize = 64 * 1024;

impl Action {
    /// A truncation observation whose `tail` is cut from the front until it
    /// fits [`LARGEST_PAYLOAD_BYTES`]; `delivered` and `errors` are kept.
    pub fn stream_truncated(delivered: usize, mut tail: Vec<String>, errors: Vec<Reason>) -> Self {
        while !tail.is_empty()
            && serde_json::to_vec(&tail).map_or(usize::MAX, |bytes| bytes.len())
                > LARGEST_PAYLOAD_BYTES
        {
            tail.remove(0);
        }
        Self::StreamTruncated {
            delivered,
            tail,
            errors,
        }
    }
}

/// A host-defined action: a named, serde-typed fact a host policy emits
/// through [`Action::Host`]. The kind is declared once per type, so a
/// trace names every host fact and a consumer deserializes it back.
pub trait HostAction: Serialize + serde::de::DeserializeOwned {
    /// The stable kind name (`rigcoder/approval`).
    const KIND: &'static str;

    /// This fact as an [`Action::Host`]; an unserializable fact is an error,
    /// never a silent omission.
    fn action(&self) -> Result<Action, serde_json::Error> {
        Ok(Action::Host {
            kind: Self::KIND.to_owned(),
            payload: serde_json::to_value(self)?,
        })
    }

    /// The fact back out of an [`Action::Host`] of this kind.
    fn from_action(action: &Action) -> Option<Result<Self, serde_json::Error>> {
        match action {
            Action::Host { kind, payload } if kind == Self::KIND => {
                Some(serde_json::from_value(payload.clone()))
            }
            _ => None,
        }
    }
}

/// A host-owned monotonic clock: what a sink stamps [`Observation::at`]
/// from. rig-core reads no clock itself; a test supplies a counter.
pub trait Clock: WasmCompatSend + WasmCompatSync {
    /// Elapsed time since the clock's origin.
    fn elapsed(&self) -> Duration;
}

/// Where observations go: the seam a driver and a host emit through. A
/// witness is shared, so it takes `&self`; it must never block the caller.
pub trait Witness: WasmCompatSend + WasmCompatSync + 'static {
    /// One fact. The sink assigns its sequence.
    fn observe(&self, observation: Observation);
}

/// The serializable trace a sink produces: the analysis artifact.
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct ObservationTrace {
    /// The sink's declared session, when it has one (a run's scope, a
    /// trial id). Lineage, not semantics.
    #[serde(default, skip_serializing_if = "Option::is_none")]
    pub session: Option<String>,
    /// The facts, in sequence.
    pub observations: Vec<Observation>,
    /// Facts discarded because sink capacity was reached. Nonzero means the
    /// trace is incomplete.
    #[serde(default)]
    pub dropped: u64,
    /// Whether the sink was told the session finished normally.
    #[serde(default)]
    pub finalized: bool,
}

impl ObservationTrace {
    /// Whether every fact the session produced is here.
    pub fn is_complete(&self) -> bool {
        self.dropped == 0
    }
}

/// The bounded in-memory sink, in one of two shapes. A capture
/// ([`Self::with_capacity`], the default) keeps the first `capacity` facts
/// and counts the rest as dropped: the trace is a complete prefix or says
/// how much of the session it missed, which is what a comparison wants. A
/// ring ([`Self::ring`]) keeps the *last* `capacity` facts and counts what
/// it let go: a long-running host always sees its recent past, which is
/// what a dashboard wants. Either way [`Self::drain`] takes what is kept
/// and starts over, so a host that exports periodically never fills up.
pub struct ObservationLog {
    inner: Mutex<LogState>,
    capacity: usize,
    ring: bool,
    clock: Option<Arc<dyn Clock + Send + Sync>>,
}

struct LogState {
    session: Option<String>,
    observations: std::collections::VecDeque<Observation>,
    next: u64,
    dropped: u64,
    finalized: bool,
}

impl std::fmt::Debug for ObservationLog {
    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
        let state = self.lock();
        f.debug_struct("ObservationLog")
            .field("observations", &state.observations.len())
            .field("dropped", &state.dropped)
            .field("capacity", &self.capacity)
            .finish_non_exhaustive()
    }
}

/// The default capacity of an [`ObservationLog`].
pub const DEFAULT_CAPACITY: usize = 65_536;

impl Default for ObservationLog {
    fn default() -> Self {
        Self::with_capacity(DEFAULT_CAPACITY)
    }
}

impl ObservationLog {
    /// A capture keeping the first `capacity` facts; later ones are counted
    /// as dropped.
    pub fn with_capacity(capacity: usize) -> Self {
        Self {
            inner: Mutex::new(LogState {
                session: None,
                observations: std::collections::VecDeque::new(),
                next: 0,
                dropped: 0,
                finalized: false,
            }),
            capacity,
            ring: false,
            clock: None,
        }
    }

    /// Keeps the last `capacity` facts and counts evictions as dropped.
    /// Capacity zero is raised to one.
    pub fn ring(capacity: usize) -> Self {
        Self {
            ring: true,
            // A ring of nothing would keep every fact and count it dropped.
            ..Self::with_capacity(capacity.max(1))
        }
    }

    /// Take the kept facts and start over: the session name and the
    /// sequence continue, the kept facts and the dropped count reset. A
    /// host that exports its trace periodically drains rather than
    /// letting a capture fill and stop.
    pub fn drain(&self) -> ObservationTrace {
        let mut state = self.lock();
        let trace = ObservationTrace {
            session: state.session.clone(),
            observations: state.observations.drain(..).collect(),
            dropped: state.dropped,
            finalized: state.finalized,
        };
        state.dropped = 0;
        trace
    }

    /// Stamp every fact with `clock`'s elapsed time.
    pub fn with_clock(mut self, clock: Arc<dyn Clock + Send + Sync>) -> Self {
        self.clock = Some(clock);
        self
    }

    /// Name the session the trace belongs to.
    pub fn with_session(self, session: impl Into<String>) -> Self {
        self.lock().session = Some(session.into());
        self
    }

    fn lock(&self) -> std::sync::MutexGuard<'_, LogState> {
        self.inner.lock().unwrap_or_else(PoisonError::into_inner)
    }

    /// The session finished normally: later readers know the trace is not
    /// a partial artifact of a killed process. A later fact reopens the capture.
    /// Hosts must drain their producers before exporting a finalized snapshot.
    pub fn finalize(&self) {
        self.lock().finalized = true;
    }

    /// The facts so far.
    pub fn trace(&self) -> ObservationTrace {
        let state = self.lock();
        ObservationTrace {
            session: state.session.clone(),
            observations: state.observations.iter().cloned().collect(),
            dropped: state.dropped,
            finalized: state.finalized,
        }
    }

    /// How many facts are kept.
    pub fn len(&self) -> usize {
        self.lock().observations.len()
    }

    /// Whether no observations are currently retained.
    pub fn is_empty(&self) -> bool {
        self.lock().observations.is_empty()
    }
}

impl Witness for ObservationLog {
    fn observe(&self, mut observation: Observation) {
        let at = self.clock.as_ref().map(|clock| clock.elapsed());
        let mut state = self.lock();
        state.finalized = false;
        observation.seq = state.next;
        state.next += 1;
        if state.observations.len() >= self.capacity {
            state.dropped += 1;
            if !self.ring {
                return;
            }
            state.observations.pop_front();
        }
        observation.at = at;
        state.observations.push_back(observation);
    }
}

impl<W: Witness + ?Sized> Witness for Arc<W> {
    fn observe(&self, observation: Observation) {
        (**self).observe(observation);
    }
}

// A trace serializes and crosses threads on every target.
const _: fn() = || {
    fn assert_wire<T: Clone + Send + Sync + 'static + Serialize + serde::de::DeserializeOwned>() {}
    assert_wire::<Observation>();
    assert_wire::<ObservationTrace>();
};