Skip to main content

turnframe_core/
observe.rs

1//! Observability signals (spec §26.2, §28).
2//!
3//! [`Signal`] mirrors every metric of the spec under the `turnframe.` prefix,
4//! plus the separately measured latencies of spec §28. An [`Observer`] receives
5//! signals with optional structured [`SignalLabels`]; duration signals travel
6//! through [`Observer::observe_duration`].
7//!
8//! **Never put user text, case text, quotes or free-form payloads in labels.**
9//! Labels are identifiers and enumerations only (workflow key, provider key,
10//! model key, request purpose, risk class, interaction kind, stable error
11//! code). Case ids are deliberately absent: they belong in traces, not in
12//! metric dimensions.
13
14use std::time::Duration;
15
16use crate::command::RiskClass;
17use crate::ids::{ModelKey, OperationKey, ProviderKey, WorkflowKey};
18use crate::interaction::InteractionKind;
19
20/// A metric event. Names are stable and match spec §26.2 with the
21/// `turnframe.` prefix; the duration variants follow spec §28.
22#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
23#[non_exhaustive]
24pub enum Signal {
25    /// A turn was received.
26    TurnReceived,
27    /// A turn completed.
28    TurnCompleted,
29    /// A turn failed.
30    TurnFailed,
31    /// A target was ambiguous.
32    TargetAmbiguous,
33    /// A target was missing.
34    TargetMissing,
35    /// An act named a target that produced no resolution at all: an act on the active
36    /// card when none is open, a target shape the operation does not accept.
37    ///
38    /// Distinct from [`Self::TargetAmbiguous`] and [`Self::TargetMissing`], which are
39    /// resolutions. This act leaves no resolution, no policy decision and no command,
40    /// so without this nothing records why the turn did less; the label carries the
41    /// rejection code.
42    TargetUnresolved,
43    /// The case directory refused to authorize a candidate for this actor.
44    ///
45    /// A candidate reaching the directory and being refused is ordinary in an
46    /// application with a scope narrower than the account: a card outlives the
47    /// scope it was written in and stops being addressable. A rate that climbs
48    /// is not ordinary, and without this an operator has only a log line, so
49    /// the one containment path below the account has no alert.
50    CaseNotAuthorized,
51    /// A command needs confirmation.
52    CommandConfirmationRequired,
53    /// A command executed.
54    CommandExecuted,
55    /// A command was rejected.
56    CommandRejected,
57    /// A command hit an idempotency replay.
58    CommandIdempotencyReplay,
59    /// A command hit a revision conflict.
60    CommandRevisionConflict,
61    /// An interaction was created.
62    InteractionCreated,
63    /// An interaction was resolved.
64    InteractionResolved,
65    /// An interaction response was stale.
66    InteractionStale,
67    /// An interaction failed.
68    InteractionFailed,
69    /// A receipt was emitted.
70    ClaimReceiptEmitted,
71    /// An external outcome is unknown.
72    ExternalOutcomeUnknown,
73    /// An external outcome was reconciled.
74    ExternalReconciled,
75    /// Provider fallback happened.
76    ProviderFallback,
77    /// A provider lacked a required capability.
78    ProviderCapabilityMismatch,
79    /// A workflow invariant was violated.
80    WorkflowInvariantViolation,
81    /// A question was answered.
82    QuestionAnswered,
83    /// A question could not be answered.
84    QuestionUnanswered,
85    /// An act the domain refused during reduction.
86    ///
87    /// Distinct from [`Self::CommandRejected`], which counts a command refused
88    /// at execution: this one never became a command. Both are a turn doing
89    /// less than its plan said, and until this existed only the second was
90    /// countable.
91    ActRefused,
92    /// An act was dropped because a correction or a cancel in the same message replaced
93    /// it. Usually right, and still a turn doing less than its acts said, so it is
94    /// counted and labelled with the operation dropped.
95    ActSuperseded,
96    /// Wall-clock duration of a whole turn, from accepted input to persisted
97    /// response, in milliseconds (spec §28).
98    TurnDuration,
99    /// Pure projection time of one case, in microseconds (spec §28).
100    ProjectionDuration,
101    /// Whole-turn reduction time, in microseconds (spec §28).
102    ReductionDuration,
103    /// Time spent in persistence for one turn, in milliseconds (spec §28).
104    PersistenceDuration,
105    /// Latency of one provider call, in milliseconds (spec §28).
106    ProviderLatency,
107    /// Latency of one external command dispatch, in milliseconds (spec §28).
108    ExternalLatency,
109    /// Latency of one call that writes or reviews the reply, in milliseconds (spec §28).
110    NarrationLatency,
111    /// A model task finished; labelled with its purpose and verdict.
112    TaskCompleted,
113    /// A model task's answer was sent back for a repair.
114    TaskRepaired,
115    /// A model task was escalated to a stronger model.
116    TaskEscalated,
117    /// The votes of a model task found no majority.
118    TaskVoteDisagreement,
119    /// A turn's model calls reached a bound; labelled with the bound.
120    BudgetExhausted,
121    /// Latency of one model task call, in milliseconds.
122    TaskLatency,
123}
124
125/// How a [`Signal`] is measured.
126#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
127#[non_exhaustive]
128pub enum SignalKind {
129    /// A monotonically increasing count of occurrences.
130    Counter,
131    /// A distribution of durations recorded in milliseconds.
132    DurationMillis,
133    /// A distribution of durations recorded in microseconds.
134    DurationMicros,
135}
136
137impl Signal {
138    /// Every signal, for exhaustive registration.
139    pub const ALL: [Self; 39] = [
140        Self::TurnReceived,
141        Self::TurnCompleted,
142        Self::TurnFailed,
143        Self::TargetAmbiguous,
144        Self::TargetUnresolved,
145        Self::TargetMissing,
146        Self::CaseNotAuthorized,
147        Self::CommandConfirmationRequired,
148        Self::CommandExecuted,
149        Self::CommandRejected,
150        Self::CommandIdempotencyReplay,
151        Self::CommandRevisionConflict,
152        Self::InteractionCreated,
153        Self::InteractionResolved,
154        Self::InteractionStale,
155        Self::InteractionFailed,
156        Self::ClaimReceiptEmitted,
157        Self::ExternalOutcomeUnknown,
158        Self::ExternalReconciled,
159        Self::ProviderFallback,
160        Self::ProviderCapabilityMismatch,
161        Self::WorkflowInvariantViolation,
162        Self::QuestionAnswered,
163        Self::QuestionUnanswered,
164        Self::ActSuperseded,
165        Self::ActRefused,
166        Self::TurnDuration,
167        Self::ProjectionDuration,
168        Self::ReductionDuration,
169        Self::PersistenceDuration,
170        Self::ProviderLatency,
171        Self::ExternalLatency,
172        Self::NarrationLatency,
173        Self::TaskCompleted,
174        Self::TaskRepaired,
175        Self::TaskEscalated,
176        Self::TaskVoteDisagreement,
177        Self::BudgetExhausted,
178        Self::TaskLatency,
179    ];
180
181    /// The metric name.
182    #[must_use]
183    pub const fn name(&self) -> &'static str {
184        match self {
185            Self::TurnReceived => "turnframe.turn.received",
186            Self::TurnCompleted => "turnframe.turn.completed",
187            Self::TurnFailed => "turnframe.turn.failed",
188            Self::TargetAmbiguous => "turnframe.target.ambiguous",
189            Self::TargetUnresolved => "turnframe.target.unresolved",
190            Self::TargetMissing => "turnframe.target.missing",
191            Self::CaseNotAuthorized => "turnframe.case.not_authorized",
192            Self::CommandConfirmationRequired => "turnframe.command.confirmation_required",
193            Self::CommandExecuted => "turnframe.command.executed",
194            Self::CommandRejected => "turnframe.command.rejected",
195            Self::CommandIdempotencyReplay => "turnframe.command.idempotency_replay",
196            Self::CommandRevisionConflict => "turnframe.command.revision_conflict",
197            Self::InteractionCreated => "turnframe.interaction.created",
198            Self::InteractionResolved => "turnframe.interaction.resolved",
199            Self::InteractionStale => "turnframe.interaction.stale",
200            Self::InteractionFailed => "turnframe.interaction.failed",
201            Self::ClaimReceiptEmitted => "turnframe.claim.receipt_emitted",
202            Self::ExternalOutcomeUnknown => "turnframe.external.outcome_unknown",
203            Self::ExternalReconciled => "turnframe.external.reconciled",
204            Self::ProviderFallback => "turnframe.provider.fallback",
205            Self::ProviderCapabilityMismatch => "turnframe.provider.capability_mismatch",
206            Self::WorkflowInvariantViolation => "turnframe.workflow.invariant_violation",
207            Self::QuestionAnswered => "turnframe.question.answered",
208            Self::QuestionUnanswered => "turnframe.question.unanswered",
209            Self::ActSuperseded => "turnframe.act.superseded",
210            Self::ActRefused => "turnframe.act.refused",
211            Self::TurnDuration => "turnframe.turn.duration_ms",
212            Self::ProjectionDuration => "turnframe.projection.duration_us",
213            Self::ReductionDuration => "turnframe.reduction.duration_us",
214            Self::PersistenceDuration => "turnframe.persistence.duration_ms",
215            Self::ProviderLatency => "turnframe.provider.latency_ms",
216            Self::ExternalLatency => "turnframe.external.latency_ms",
217            Self::NarrationLatency => "turnframe.narration.latency_ms",
218            Self::TaskCompleted => "turnframe.task.completed",
219            Self::TaskRepaired => "turnframe.task.repaired",
220            Self::TaskEscalated => "turnframe.task.escalated",
221            Self::TaskVoteDisagreement => "turnframe.task.vote_disagreement",
222            Self::BudgetExhausted => "turnframe.budget.exhausted",
223            Self::TaskLatency => "turnframe.task.latency_ms",
224        }
225    }
226
227    /// How the signal is measured: an occurrence counter or a duration
228    /// distribution in the unit the name ends with (`_ms`, `_us`).
229    #[must_use]
230    pub const fn kind(&self) -> SignalKind {
231        match self {
232            Self::TurnDuration
233            | Self::PersistenceDuration
234            | Self::ProviderLatency
235            | Self::ExternalLatency
236            | Self::NarrationLatency
237            | Self::TaskLatency => SignalKind::DurationMillis,
238            Self::ProjectionDuration | Self::ReductionDuration => SignalKind::DurationMicros,
239            Self::TurnReceived
240            | Self::TurnCompleted
241            | Self::TurnFailed
242            | Self::TargetAmbiguous
243            | Self::TargetUnresolved
244            | Self::ActSuperseded
245            | Self::ActRefused
246            | Self::TargetMissing
247            | Self::CaseNotAuthorized
248            | Self::CommandConfirmationRequired
249            | Self::CommandExecuted
250            | Self::CommandRejected
251            | Self::CommandIdempotencyReplay
252            | Self::CommandRevisionConflict
253            | Self::InteractionCreated
254            | Self::InteractionResolved
255            | Self::InteractionStale
256            | Self::InteractionFailed
257            | Self::ClaimReceiptEmitted
258            | Self::ExternalOutcomeUnknown
259            | Self::ExternalReconciled
260            | Self::ProviderFallback
261            | Self::ProviderCapabilityMismatch
262            | Self::WorkflowInvariantViolation
263            | Self::QuestionAnswered
264            | Self::QuestionUnanswered
265            | Self::TaskCompleted
266            | Self::TaskRepaired
267            | Self::TaskEscalated
268            | Self::TaskVoteDisagreement
269            | Self::BudgetExhausted => SignalKind::Counter,
270        }
271    }
272
273    /// Returns `true` for signals that indicate a safety-integrity problem
274    /// rather than a language-quality problem (spec §26.3).
275    #[must_use]
276    pub const fn is_safety_signal(&self) -> bool {
277        matches!(
278            self,
279            Self::CommandRevisionConflict
280                | Self::InteractionStale
281                | Self::ExternalOutcomeUnknown
282                | Self::WorkflowInvariantViolation
283                | Self::ProviderCapabilityMismatch
284        )
285    }
286}
287
288/// Structured, text-free labels attached to a signal.
289///
290/// Every field is an identifier or an enumeration with a bounded value set.
291/// Observers may drop labels a signal does not document; they must never add
292/// free text.
293#[derive(Debug, Clone, Default, PartialEq, Eq)]
294#[non_exhaustive]
295pub struct SignalLabels {
296    /// Workflow the signal concerns.
297    pub workflow: Option<WorkflowKey>,
298    /// Provider key, for provider signals.
299    pub provider: Option<ProviderKey>,
300    /// Risk class, for command signals.
301    pub risk: Option<RiskClass>,
302    /// Model key, for provider signals. Configured model identifiers only.
303    pub model: Option<ModelKey>,
304    /// Normalized request purpose (e.g. `"extract"`), for provider
305    /// signals (spec §20.2).
306    pub purpose: Option<String>,
307    /// Interaction kind, for interaction signals.
308    pub interaction: Option<InteractionKind>,
309    /// Stable machine-readable code that qualifies a failure or rejection
310    /// (a [`RejectionCode`](crate::error::RejectionCode), a provider failure
311    /// code, an invariant kind...). Never a message.
312    pub error_code: Option<String>,
313    /// Operation an act names, for signals about one act rather than one turn.
314    ///
315    /// Bounded like [`Self::workflow`] is, because the set of operations is the
316    /// catalog and the catalog is finite. Never an argument or a value.
317    pub operation: Option<OperationKey>,
318    /// The effort of the turn the signal belongs to.
319    pub effort: Option<crate::effort::Effort>,
320}
321
322impl SignalLabels {
323    /// No labels.
324    #[must_use]
325    pub fn none() -> Self {
326        Self::default()
327    }
328
329    /// Labels with a workflow.
330    #[must_use]
331    pub fn workflow(workflow: WorkflowKey) -> Self {
332        Self {
333            workflow: Some(workflow),
334            ..Self::default()
335        }
336    }
337
338    /// Adds a provider key.
339    #[must_use]
340    pub fn with_provider(mut self, provider: impl Into<ProviderKey>) -> Self {
341        self.provider = Some(provider.into());
342        self
343    }
344
345    /// Adds a risk class.
346    #[must_use]
347    pub fn with_risk(mut self, risk: RiskClass) -> Self {
348        self.risk = Some(risk);
349        self
350    }
351
352    /// Adds a model key.
353    #[must_use]
354    pub fn with_model(mut self, model: impl Into<ModelKey>) -> Self {
355        self.model = Some(model.into());
356        self
357    }
358
359    /// Adds a normalized request purpose.
360    #[must_use]
361    pub fn with_purpose(mut self, purpose: impl Into<String>) -> Self {
362        self.purpose = Some(purpose.into());
363        self
364    }
365
366    /// Labels the signal with the turn's effort.
367    #[must_use]
368    pub fn with_effort(mut self, effort: crate::effort::Effort) -> Self {
369        self.effort = Some(effort);
370        self
371    }
372
373    /// Adds an interaction kind.
374    #[must_use]
375    pub fn with_interaction(mut self, kind: InteractionKind) -> Self {
376        self.interaction = Some(kind);
377        self
378    }
379
380    /// Adds a stable error or rejection code.
381    #[must_use]
382    pub fn with_operation(mut self, operation: impl Into<OperationKey>) -> Self {
383        self.operation = Some(operation.into());
384        self
385    }
386
387    /// Returns a copy carrying a failure or rejection code.
388    #[must_use]
389    pub fn with_error_code(mut self, code: impl Into<String>) -> Self {
390        self.error_code = Some(code.into());
391        self
392    }
393}
394
395/// Receives signals. Implementations must not block.
396pub trait Observer: Send + Sync {
397    /// Observes a signal without labels.
398    fn observe(&self, signal: &Signal);
399
400    /// Observes a signal with labels. The default drops the labels.
401    fn observe_labeled(&self, signal: &Signal, labels: &SignalLabels) {
402        let _ = labels;
403        self.observe(signal);
404    }
405
406    /// Observes a duration signal ([`SignalKind::DurationMillis`] or
407    /// [`SignalKind::DurationMicros`]) with its measured value. The default
408    /// drops the value and forwards to [`Observer::observe_labeled`].
409    fn observe_duration(&self, signal: &Signal, duration: Duration, labels: &SignalLabels) {
410        let _ = duration;
411        self.observe_labeled(signal, labels);
412    }
413}
414
415/// An observer that ignores everything.
416#[derive(Debug, Clone, Copy, Default)]
417pub struct NoopObserver;
418
419impl Observer for NoopObserver {
420    fn observe(&self, _signal: &Signal) {}
421}
422
423#[cfg(test)]
424mod tests {
425    use super::*;
426    use std::collections::{BTreeSet, HashSet};
427
428    #[test]
429    fn names_are_unique_and_prefixed() {
430        let names: BTreeSet<&str> = Signal::ALL.iter().map(Signal::name).collect();
431        assert_eq!(names.len(), Signal::ALL.len());
432        assert!(names.iter().all(|n| n.starts_with("turnframe.")));
433    }
434
435    /// Fails to compile when a variant is added without touching this match,
436    /// and fails at runtime when a variant is missing from `ALL`.
437    #[test]
438    fn all_lists_every_variant() {
439        for signal in Signal::ALL {
440            let listed = match signal {
441                Signal::TurnReceived
442                | Signal::TurnCompleted
443                | Signal::TurnFailed
444                | Signal::TargetAmbiguous
445                | Signal::TargetUnresolved
446                | Signal::TargetMissing
447                | Signal::CaseNotAuthorized
448                | Signal::CommandConfirmationRequired
449                | Signal::CommandExecuted
450                | Signal::CommandRejected
451                | Signal::CommandIdempotencyReplay
452                | Signal::CommandRevisionConflict
453                | Signal::InteractionCreated
454                | Signal::InteractionResolved
455                | Signal::InteractionStale
456                | Signal::InteractionFailed
457                | Signal::ClaimReceiptEmitted
458                | Signal::ExternalOutcomeUnknown
459                | Signal::ExternalReconciled
460                | Signal::ProviderFallback
461                | Signal::ProviderCapabilityMismatch
462                | Signal::WorkflowInvariantViolation
463                | Signal::QuestionAnswered
464                | Signal::QuestionUnanswered
465                | Signal::ActSuperseded
466                | Signal::ActRefused
467                | Signal::TurnDuration
468                | Signal::ProjectionDuration
469                | Signal::ReductionDuration
470                | Signal::PersistenceDuration
471                | Signal::ProviderLatency
472                | Signal::ExternalLatency
473                | Signal::NarrationLatency
474                | Signal::TaskCompleted
475                | Signal::TaskRepaired
476                | Signal::TaskEscalated
477                | Signal::TaskVoteDisagreement
478                | Signal::BudgetExhausted
479                | Signal::TaskLatency => true,
480            };
481            assert!(listed);
482        }
483        let distinct: HashSet<Signal> = Signal::ALL.into_iter().collect();
484        assert_eq!(distinct.len(), Signal::ALL.len());
485    }
486
487    #[test]
488    fn kind_matches_name_suffix() {
489        for signal in Signal::ALL {
490            let name = signal.name();
491            match signal.kind() {
492                SignalKind::DurationMillis => assert!(name.ends_with("_ms"), "{name}"),
493                SignalKind::DurationMicros => assert!(name.ends_with("_us"), "{name}"),
494                SignalKind::Counter => {
495                    assert!(!name.ends_with("_ms") && !name.ends_with("_us"), "{name}");
496                }
497            }
498        }
499    }
500
501    #[test]
502    fn labels_builders_set_every_field() {
503        let labels = SignalLabels::workflow(WorkflowKey::from("trip"))
504            .with_provider("openai")
505            .with_model("gpt-x")
506            .with_purpose("extract")
507            .with_risk(RiskClass::Destructive)
508            .with_interaction(InteractionKind::ConfirmCommand)
509            .with_error_code("rate_limited");
510        assert_eq!(
511            labels.workflow.as_ref().map(WorkflowKey::as_str),
512            Some("trip")
513        );
514        assert_eq!(
515            labels.provider.as_ref().map(ProviderKey::as_str),
516            Some("openai")
517        );
518        assert_eq!(labels.model.as_ref().map(ModelKey::as_str), Some("gpt-x"));
519        assert_eq!(labels.purpose.as_deref(), Some("extract"));
520        assert_eq!(labels.risk, Some(RiskClass::Destructive));
521        assert_eq!(labels.interaction, Some(InteractionKind::ConfirmCommand));
522        assert_eq!(labels.error_code.as_deref(), Some("rate_limited"));
523    }
524
525    #[test]
526    fn noop_observer_accepts_everything() {
527        let observer = NoopObserver;
528        for signal in Signal::ALL {
529            observer.observe(&signal);
530            observer.observe_labeled(&signal, &SignalLabels::workflow(WorkflowKey::from("w")));
531            observer.observe_duration(&signal, Duration::from_millis(3), &SignalLabels::none());
532        }
533    }
534}