Skip to main content

turnframe_telemetry/
metrics.rs

1//! Metric emission through the [`metrics`](https://docs.rs/metrics) facade
2//! (spec §26.2, §28).
3//!
4//! Every [`Signal`] becomes one metric named exactly by
5//! [`Signal::name`]: a counter for occurrence signals, a histogram for the
6//! latency signals of spec §28, recorded in the unit its name ends with.
7//!
8//! The dimensions of those metrics are deliberately narrow. A signal declares
9//! the labels it documents in [`documented_labels`]; anything else the caller
10//! passed is dropped, and every surviving value must satisfy
11//! [`is_safe_label_value`]. Neither the user's words nor a case identifier can
12//! reach a label: [`SignalLabels`] has no free-text field, and the safety check
13//! rejects UUIDs, whitespace and anything longer than
14//! [`MAX_LABEL_VALUE_LEN`].
15
16use std::time::Duration;
17
18use metrics::{Key, Label, Level, Metadata, Unit};
19use serde::{Deserialize, Serialize};
20use turnframe_core::observe::{Observer, Signal, SignalKind, SignalLabels};
21
22use crate::{enum_label, opt_str};
23
24/// Target every Turnframe metric is registered under, for exporter filtering.
25pub const METRIC_TARGET: &str = "turnframe";
26
27/// Longest label value that is emitted. Longer values are dropped rather than
28/// truncated: a truncated identifier is still an identifier.
29pub const MAX_LABEL_VALUE_LEN: usize = 64;
30
31static METADATA: Metadata<'static> =
32    Metadata::new(METRIC_TARGET, Level::INFO, Some(module_path!()));
33
34/// Name of a metric dimension.
35///
36/// The set is closed on purpose: it is exactly the typed, bounded fields of
37/// [`SignalLabels`]. There is no variant for anything an end user typed.
38#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, PartialOrd, Ord, Serialize, Deserialize)]
39#[serde(rename_all = "snake_case")]
40#[non_exhaustive]
41pub enum LabelKey {
42    /// Workflow key, e.g. `trip`.
43    Workflow,
44    /// Configured provider key, e.g. `openai`.
45    Provider,
46    /// Configured model key, e.g. `gpt-x`.
47    Model,
48    /// Normalized request purpose, e.g. `extract` (spec §20.2).
49    Purpose,
50    /// Risk class of a command (spec §14.3).
51    Risk,
52    /// Kind of an interaction (spec §15).
53    Interaction,
54    /// Stable machine-readable failure or rejection code. Never a message.
55    ErrorCode,
56    /// Operation an act names, e.g. `trip.set_name`. Bounded by the
57    /// catalog, like [`Self::Workflow`] is bounded by the registry.
58    Operation,
59    /// The effort of the turn, `low`, `medium` or `high`.
60    Effort,
61}
62
63impl LabelKey {
64    /// Every label name, in the order they are emitted.
65    pub const ALL: [Self; 9] = [
66        Self::Workflow,
67        Self::Provider,
68        Self::Model,
69        Self::Purpose,
70        Self::Risk,
71        Self::Interaction,
72        Self::ErrorCode,
73        Self::Operation,
74        Self::Effort,
75    ];
76
77    /// The label name as it appears on the metric.
78    #[must_use]
79    pub const fn as_str(self) -> &'static str {
80        match self {
81            Self::Workflow => "workflow",
82            Self::Provider => "provider",
83            Self::Model => "model",
84            Self::Purpose => "purpose",
85            Self::Risk => "risk",
86            Self::Interaction => "interaction",
87            Self::ErrorCode => "error_code",
88            Self::Operation => "operation",
89            Self::Effort => "effort",
90        }
91    }
92
93    /// Reads this dimension out of a label set, if the caller supplied it.
94    #[must_use]
95    pub fn value_of(self, labels: &SignalLabels) -> Option<String> {
96        match self {
97            Self::Workflow => opt_str(&labels.workflow).map(ToOwned::to_owned),
98            Self::Provider => opt_str(&labels.provider).map(ToOwned::to_owned),
99            Self::Model => opt_str(&labels.model).map(ToOwned::to_owned),
100            Self::Purpose => opt_str(&labels.purpose).map(ToOwned::to_owned),
101            Self::Risk => labels.risk.as_ref().and_then(enum_label),
102            Self::Interaction => labels.interaction.as_ref().and_then(enum_label),
103            Self::ErrorCode => opt_str(&labels.error_code).map(ToOwned::to_owned),
104            Self::Operation => labels.operation.as_ref().map(|key| key.as_str().to_owned()),
105            Self::Effort => labels.effort.map(|effort| effort.as_str().to_owned()),
106        }
107    }
108}
109
110/// Returns `true` when `value` is safe to use as a metric dimension.
111///
112/// A value is rejected when it is empty, longer than [`MAX_LABEL_VALUE_LEN`]
113/// bytes, contains whitespace (a hallmark of prose rather than a code) or
114/// parses as a UUID (the shape of every server-generated record identifier).
115/// Configured keys, enumeration names and stable error codes all pass.
116#[must_use]
117pub fn is_safe_label_value(value: &str) -> bool {
118    !value.is_empty()
119        && value.len() <= MAX_LABEL_VALUE_LEN
120        && !value.chars().any(char::is_whitespace)
121        && uuid::Uuid::parse_str(value).is_err()
122}
123
124/// The dimensions a signal documents.
125///
126/// Labels outside this set are dropped, which is what keeps one metric's
127/// cardinality bounded no matter what a caller attaches. Signals added to
128/// [`Signal`] after this crate was compiled fall back to the workflow
129/// dimension alone.
130#[must_use]
131pub fn documented_labels(signal: Signal) -> &'static [LabelKey] {
132    use LabelKey::{
133        Effort, ErrorCode, Interaction, Model, Operation, Provider, Purpose, Risk, Workflow,
134    };
135
136    const WORKFLOW: &[LabelKey] = &[Workflow];
137    const WORKFLOW_ERROR: &[LabelKey] = &[Workflow, ErrorCode];
138    const COMMAND: &[LabelKey] = &[Workflow, Risk];
139    const COMMAND_ERROR: &[LabelKey] = &[Workflow, Risk, ErrorCode];
140    const INTERACTION: &[LabelKey] = &[Workflow, Interaction];
141    const INTERACTION_ERROR: &[LabelKey] = &[Workflow, Interaction, ErrorCode];
142    const PROVIDER_ERROR: &[LabelKey] = &[Provider, Model, Purpose, ErrorCode];
143    const PROVIDER: &[LabelKey] = &[Provider, Model, Purpose];
144    const OPERATION: &[LabelKey] = &[Workflow, Operation];
145    const TURN: &[LabelKey] = &[Workflow, Effort];
146    const TURN_ERROR: &[LabelKey] = &[Workflow, ErrorCode, Effort];
147    const TASK: &[LabelKey] = &[Provider, Model, Purpose, Effort];
148    const TASK_ERROR: &[LabelKey] = &[Provider, Model, Purpose, ErrorCode, Effort];
149
150    match signal {
151        Signal::TurnReceived | Signal::TurnCompleted | Signal::TurnDuration => TURN,
152        Signal::TurnFailed => TURN_ERROR,
153        Signal::TargetAmbiguous | Signal::TargetMissing | Signal::CaseNotAuthorized => WORKFLOW,
154        // The reason is the whole value here: "an act went nowhere" is a number,
155        // "an act named a card that was not open" is a diagnosis.
156        Signal::TargetUnresolved => WORKFLOW_ERROR,
157        Signal::ActSuperseded => OPERATION,
158        Signal::ActRefused => WORKFLOW,
159        Signal::CommandConfirmationRequired | Signal::CommandExecuted => COMMAND,
160        Signal::CommandRejected => COMMAND_ERROR,
161        Signal::CommandIdempotencyReplay | Signal::CommandRevisionConflict => COMMAND,
162        Signal::InteractionCreated | Signal::InteractionResolved | Signal::InteractionStale => {
163            INTERACTION
164        }
165        Signal::InteractionFailed => INTERACTION_ERROR,
166        Signal::ClaimReceiptEmitted => WORKFLOW,
167        Signal::ExternalOutcomeUnknown => WORKFLOW_ERROR,
168        Signal::ExternalReconciled => WORKFLOW,
169        Signal::ProviderFallback | Signal::ProviderCapabilityMismatch => PROVIDER_ERROR,
170        Signal::WorkflowInvariantViolation => WORKFLOW_ERROR,
171        Signal::QuestionAnswered | Signal::QuestionUnanswered => WORKFLOW,
172        Signal::ProjectionDuration
173        | Signal::ReductionDuration
174        | Signal::PersistenceDuration
175        | Signal::ExternalLatency => WORKFLOW,
176        Signal::ProviderLatency | Signal::NarrationLatency => PROVIDER,
177        // A task's verdict and a repair's cause travel as the error code.
178        Signal::TaskCompleted | Signal::TaskRepaired | Signal::TaskEscalated => TASK_ERROR,
179        Signal::TaskVoteDisagreement | Signal::TaskLatency => TASK,
180        // The bound that stopped the turn is the error code.
181        Signal::BudgetExhausted => TURN_ERROR,
182        // `Signal` is `#[non_exhaustive]`; an unknown signal keeps the safest
183        // dimension rather than none, so it is still attributable.
184        _ => WORKFLOW,
185    }
186}
187
188/// The metric type a signal is exported as: `counter` or `histogram`.
189#[must_use]
190pub fn metric_type(signal: Signal) -> &'static str {
191    match signal.kind() {
192        SignalKind::DurationMillis | SignalKind::DurationMicros => "histogram",
193        _ => "counter",
194    }
195}
196
197/// The unit a signal is recorded in.
198#[must_use]
199pub fn unit(signal: Signal) -> Unit {
200    match signal.kind() {
201        SignalKind::DurationMillis => Unit::Milliseconds,
202        SignalKind::DurationMicros => Unit::Microseconds,
203        _ => Unit::Count,
204    }
205}
206
207/// One-line meaning of a signal, registered as the metric description and
208/// published in the crate's metric catalogue.
209#[must_use]
210pub fn description(signal: Signal) -> &'static str {
211    match signal {
212        Signal::TurnReceived => "User turns accepted for processing.",
213        Signal::TurnCompleted => "Turns that produced a persisted response.",
214        Signal::TurnFailed => "Turns that ended in an orchestrator error.",
215        Signal::TargetAmbiguous => "Acts whose target matched more than one case.",
216        Signal::TargetMissing => "Acts whose target matched no case.",
217        Signal::TargetUnresolved => {
218            "Acts whose target produced no resolution at all, leaving no target record, no policy \
219             decision and no command."
220        }
221        Signal::CaseNotAuthorized => {
222            "Candidates the case directory refused to authorize for the actor."
223        }
224        Signal::CommandConfirmationRequired => {
225            "Commands held back until a confirmation interaction is answered."
226        }
227        Signal::CommandExecuted => "Commands that committed.",
228        Signal::CommandRejected => "Commands the domain refused.",
229        Signal::CommandIdempotencyReplay => {
230            "Commands whose idempotency key was already in the journal."
231        }
232        Signal::CommandRevisionConflict => "Commands planned against a stale case revision.",
233        Signal::InteractionCreated => "Interaction cards persisted for the user.",
234        Signal::InteractionResolved => "Interaction cards answered and resolved.",
235        Signal::InteractionStale => "Answers that arrived after the case had moved on.",
236        Signal::InteractionFailed => "Interaction responses that could not be accepted.",
237        Signal::ClaimReceiptEmitted => "Operational receipts derived from committed events.",
238        Signal::ExternalOutcomeUnknown => {
239            "External side effects whose outcome is unknown and needs reconciliation."
240        }
241        Signal::ExternalReconciled => "External side effects whose outcome was later established.",
242        Signal::ProviderFallback => "Provider calls that fell back to another candidate.",
243        Signal::ProviderCapabilityMismatch => {
244            "Provider calls refused because the model lacked a required capability."
245        }
246        Signal::WorkflowInvariantViolation => "Projected views that broke a Flow Map invariant.",
247        Signal::QuestionAnswered => {
248            "Questions answered from the facts of their records or a knowledge source."
249        }
250        Signal::QuestionUnanswered => "Questions the turn could not answer.",
251        Signal::ActRefused => {
252            "Acts the domain refused during reduction, which never became commands."
253        }
254        Signal::ActSuperseded => "Acts a correction or a cancel in the same message replaced.",
255        Signal::TurnDuration => "Wall-clock time from accepted input to persisted response.",
256        Signal::ProjectionDuration => "Pure projection time of one case.",
257        Signal::ReductionDuration => "Whole-turn reduction time.",
258        Signal::PersistenceDuration => "Time spent in persistence for one turn.",
259        Signal::ProviderLatency => "Latency of one provider call.",
260        Signal::ExternalLatency => "Latency of one external command dispatch.",
261        Signal::NarrationLatency => "Latency of one call that writes or reviews the reply.",
262        Signal::TaskCompleted => "Model tasks finished, by purpose and verdict.",
263        Signal::TaskRepaired => "Model task answers sent back for a repair.",
264        Signal::TaskEscalated => "Model tasks re-run on a stronger model.",
265        Signal::TaskVoteDisagreement => "Model task votes that found no majority.",
266        Signal::BudgetExhausted => "Turns whose model calls reached a bound.",
267        Signal::TaskLatency => "Latency of one model task call.",
268        // `Signal` is `#[non_exhaustive]`: a signal this crate predates still
269        // gets a metric, just a generic description.
270        _ => "Turnframe signal.",
271    }
272}
273
274/// Builds the labels actually emitted for a signal: the dimensions the signal
275/// documents, restricted to the values the caller supplied and to those that
276/// pass [`is_safe_label_value`].
277#[must_use]
278pub fn metric_labels(signal: Signal, labels: &SignalLabels) -> Vec<Label> {
279    documented_labels(signal)
280        .iter()
281        .filter_map(|key| key.value_of(labels).map(|value| (*key, value)))
282        .filter(|(_, value)| is_safe_label_value(value))
283        .map(|(key, value)| Label::new(key.as_str(), value))
284        .collect()
285}
286
287/// One row of the metric catalogue: everything the README table publishes.
288#[derive(Debug, Clone, Copy, PartialEq, Eq)]
289#[non_exhaustive]
290pub struct MetricDoc {
291    /// The signal this metric reports.
292    pub signal: Signal,
293    /// The metric name, identical to [`Signal::name`].
294    pub name: &'static str,
295    /// `counter` or `histogram`.
296    pub metric_type: &'static str,
297    /// The recorded unit.
298    pub unit: Unit,
299    /// The dimensions the metric documents.
300    pub labels: &'static [LabelKey],
301    /// What the metric means.
302    pub description: &'static str,
303}
304
305impl MetricDoc {
306    /// The catalogue row for one signal.
307    #[must_use]
308    pub fn of(signal: Signal) -> Self {
309        Self {
310            signal,
311            name: signal.name(),
312            metric_type: metric_type(signal),
313            unit: unit(signal),
314            labels: documented_labels(signal),
315            description: description(signal),
316        }
317    }
318
319    /// The label names, comma separated, or `—` when the metric has none.
320    #[must_use]
321    pub fn label_list(&self) -> String {
322        if self.labels.is_empty() {
323            return String::from("—");
324        }
325        self.labels
326            .iter()
327            .map(|key| format!("`{}`", key.as_str()))
328            .collect::<Vec<_>>()
329            .join(", ")
330    }
331}
332
333/// The whole metric catalogue, in the order of [`Signal::ALL`].
334#[must_use]
335pub fn catalogue() -> Vec<MetricDoc> {
336    Signal::ALL.into_iter().map(MetricDoc::of).collect()
337}
338
339/// Renders the metric catalogue as the markdown table published in the crate's
340/// README. A test compares the README against this function so the two cannot
341/// drift.
342#[must_use]
343pub fn catalogue_markdown() -> String {
344    let mut out = String::from("| Metric | Type | Unit | Labels | Meaning |\n");
345    out.push_str("| --- | --- | --- | --- | --- |\n");
346    for doc in catalogue() {
347        out.push_str(&format!(
348            "| `{}` | {} | {} | {} | {} |\n",
349            doc.name,
350            doc.metric_type,
351            unit_label(doc.unit),
352            doc.label_list(),
353            doc.description,
354        ));
355    }
356    out
357}
358
359/// The unit as it is written in the catalogue table.
360#[must_use]
361pub fn unit_label(unit: Unit) -> &'static str {
362    match unit {
363        Unit::Milliseconds => "ms",
364        Unit::Microseconds => "µs",
365        _ => "count",
366    }
367}
368
369/// Registers the description and unit of every metric with the installed
370/// `metrics` recorder.
371///
372/// Call it once at start-up, after the exporter is installed and before the
373/// first turn, so an exporter that publishes metadata (Prometheus `HELP` and
374/// `TYPE`, OTLP descriptions) has it. Calling it with no recorder installed is
375/// a no-op.
376pub fn describe_all() {
377    metrics::with_recorder(|recorder| {
378        for signal in Signal::ALL {
379            let name = metrics::KeyName::from(signal.name());
380            let metric_unit = unit(signal);
381            let text = metrics::SharedString::from(description(signal));
382            match signal.kind() {
383                SignalKind::DurationMillis | SignalKind::DurationMicros => {
384                    recorder.describe_histogram(name, Some(metric_unit), text);
385                }
386                _ => recorder.describe_counter(name, Some(metric_unit), text),
387            }
388        }
389    });
390}
391
392/// An [`Observer`] that turns signals into `metrics` counters and histograms.
393///
394/// Counters are incremented by one per occurrence. Duration signals are
395/// recorded as histograms in the unit their name ends with — milliseconds for
396/// `*_ms`, microseconds for `*_us` — with sub-unit precision preserved.
397///
398/// A duration signal reported through
399/// [`observe_labeled`](Observer::observe_labeled) carries no measurement and is
400/// dropped; a counter signal reported through
401/// [`observe_duration`](Observer::observe_duration) is counted and the duration
402/// ignored.
403///
404/// ```rust
405/// use turnframe_core::ids::WorkflowKey;
406/// use turnframe_core::observe::{Observer, Signal, SignalLabels};
407/// use turnframe_telemetry::MetricsObserver;
408///
409/// let observer = MetricsObserver::new();
410/// observer.observe_labeled(
411///     &Signal::TurnCompleted,
412///     &SignalLabels::workflow(WorkflowKey::from("trip")),
413/// );
414/// ```
415#[derive(Debug, Clone, Copy, Default)]
416pub struct MetricsObserver;
417
418impl MetricsObserver {
419    /// Builds the observer. Emission goes to whichever recorder the
420    /// application installed, so there is nothing to configure here.
421    #[must_use]
422    pub const fn new() -> Self {
423        Self
424    }
425
426    fn record(signal: Signal, labels: &SignalLabels, duration: Option<Duration>) {
427        let key = Key::from_parts(signal.name(), metric_labels(signal, labels));
428        metrics::with_recorder(|recorder| match signal.kind() {
429            SignalKind::DurationMillis => {
430                if let Some(duration) = duration {
431                    recorder
432                        .register_histogram(&key, &METADATA)
433                        .record(duration.as_secs_f64() * 1_000.0);
434                }
435            }
436            SignalKind::DurationMicros => {
437                if let Some(duration) = duration {
438                    recorder
439                        .register_histogram(&key, &METADATA)
440                        .record(duration.as_secs_f64() * 1_000_000.0);
441                }
442            }
443            _ => recorder.register_counter(&key, &METADATA).increment(1),
444        });
445    }
446}
447
448impl Observer for MetricsObserver {
449    fn observe(&self, signal: &Signal) {
450        Self::record(*signal, &SignalLabels::none(), None);
451    }
452
453    fn observe_labeled(&self, signal: &Signal, labels: &SignalLabels) {
454        Self::record(*signal, labels, None);
455    }
456
457    fn observe_duration(&self, signal: &Signal, duration: Duration, labels: &SignalLabels) {
458        Self::record(*signal, labels, Some(duration));
459    }
460}
461
462#[cfg(test)]
463mod tests {
464    use std::sync::{Arc, Mutex, PoisonError};
465
466    use metrics::{
467        Counter, CounterFn, Gauge, Histogram, HistogramFn, KeyName, Recorder, SharedString,
468    };
469    use turnframe_core::command::RiskClass;
470    use turnframe_core::ids::WorkflowKey;
471    use turnframe_core::interaction::InteractionKind;
472
473    use super::*;
474
475    /// `Signal` is `#[non_exhaustive]`, so a downstream crate cannot write a
476    /// match over it without a wildcard arm. This const is the compile-time
477    /// gate instead: `Signal::ALL` is a fixed-size array, so adding a variant
478    /// in core changes its length and breaks this assertion until the fixture
479    /// below is extended.
480    const _: () = assert!(
481        Signal::ALL.len() == 39,
482        "a signal was added to turnframe-core: extend the fixture in this module"
483    );
484
485    #[derive(Debug, Default)]
486    struct Captured {
487        counters: Vec<(Key, u64)>,
488        histograms: Vec<(Key, f64)>,
489        described: Vec<(String, &'static str, Option<Unit>, String)>,
490    }
491
492    type Shared = Arc<Mutex<Captured>>;
493
494    fn lock(shared: &Shared) -> std::sync::MutexGuard<'_, Captured> {
495        shared.lock().unwrap_or_else(PoisonError::into_inner)
496    }
497
498    #[derive(Debug)]
499    struct CapturedCounter {
500        key: Key,
501        shared: Shared,
502    }
503
504    impl CounterFn for CapturedCounter {
505        fn increment(&self, value: u64) {
506            lock(&self.shared).counters.push((self.key.clone(), value));
507        }
508
509        fn absolute(&self, value: u64) {
510            lock(&self.shared).counters.push((self.key.clone(), value));
511        }
512    }
513
514    #[derive(Debug)]
515    struct CapturedHistogram {
516        key: Key,
517        shared: Shared,
518    }
519
520    impl HistogramFn for CapturedHistogram {
521        fn record(&self, value: f64) {
522            lock(&self.shared)
523                .histograms
524                .push((self.key.clone(), value));
525        }
526    }
527
528    /// A recorder that keeps everything it was handed, so a test can assert on
529    /// metric names, label sets and recorded values without an exporter.
530    #[derive(Debug, Default)]
531    struct TestRecorder {
532        shared: Shared,
533    }
534
535    impl Recorder for TestRecorder {
536        fn describe_counter(&self, key: KeyName, unit: Option<Unit>, description: SharedString) {
537            lock(&self.shared).described.push((
538                key.as_str().to_owned(),
539                "counter",
540                unit,
541                description.to_string(),
542            ));
543        }
544
545        fn describe_gauge(&self, key: KeyName, unit: Option<Unit>, description: SharedString) {
546            lock(&self.shared).described.push((
547                key.as_str().to_owned(),
548                "gauge",
549                unit,
550                description.to_string(),
551            ));
552        }
553
554        fn describe_histogram(&self, key: KeyName, unit: Option<Unit>, description: SharedString) {
555            lock(&self.shared).described.push((
556                key.as_str().to_owned(),
557                "histogram",
558                unit,
559                description.to_string(),
560            ));
561        }
562
563        fn register_counter(&self, key: &Key, _metadata: &Metadata<'_>) -> Counter {
564            Counter::from_arc(Arc::new(CapturedCounter {
565                key: key.clone(),
566                shared: Arc::clone(&self.shared),
567            }))
568        }
569
570        fn register_gauge(&self, _key: &Key, _metadata: &Metadata<'_>) -> Gauge {
571            Gauge::noop()
572        }
573
574        fn register_histogram(&self, key: &Key, _metadata: &Metadata<'_>) -> Histogram {
575            Histogram::from_arc(Arc::new(CapturedHistogram {
576                key: key.clone(),
577                shared: Arc::clone(&self.shared),
578            }))
579        }
580    }
581
582    /// What one signal is expected to emit.
583    struct Fixture {
584        name: &'static str,
585        labels: &'static [&'static str],
586    }
587
588    /// The expected wire shape of every signal. Kept as a match rather than a
589    /// table so the fixture reads next to the signal it describes.
590    fn fixture(signal: Signal) -> Fixture {
591        const NONE: &[&str] = &[];
592        const W: &[&str] = &["workflow"];
593        const WE: &[&str] = &["workflow", "error_code"];
594        const WR: &[&str] = &["workflow", "risk"];
595        const WRE: &[&str] = &["workflow", "risk", "error_code"];
596        const WI: &[&str] = &["workflow", "interaction"];
597        const WIE: &[&str] = &["workflow", "interaction", "error_code"];
598        const PMPE: &[&str] = &["provider", "model", "purpose", "error_code"];
599        const PMP: &[&str] = &["provider", "model", "purpose"];
600        const WO: &[&str] = &["workflow", "operation"];
601
602        let (name, labels) = match signal {
603            Signal::TurnReceived => ("turnframe.turn.received", W),
604            Signal::TurnCompleted => ("turnframe.turn.completed", W),
605            Signal::TurnFailed => ("turnframe.turn.failed", WE),
606            Signal::TargetAmbiguous => ("turnframe.target.ambiguous", W),
607            Signal::TargetMissing => ("turnframe.target.missing", W),
608            Signal::TargetUnresolved => ("turnframe.target.unresolved", WE),
609            Signal::CommandConfirmationRequired => ("turnframe.command.confirmation_required", WR),
610            Signal::CommandExecuted => ("turnframe.command.executed", WR),
611            Signal::CommandRejected => ("turnframe.command.rejected", WRE),
612            Signal::CaseNotAuthorized => ("turnframe.case.not_authorized", W),
613            Signal::CommandIdempotencyReplay => ("turnframe.command.idempotency_replay", WR),
614            Signal::CommandRevisionConflict => ("turnframe.command.revision_conflict", WR),
615            Signal::InteractionCreated => ("turnframe.interaction.created", WI),
616            Signal::InteractionResolved => ("turnframe.interaction.resolved", WI),
617            Signal::InteractionStale => ("turnframe.interaction.stale", WI),
618            Signal::InteractionFailed => ("turnframe.interaction.failed", WIE),
619            Signal::ClaimReceiptEmitted => ("turnframe.claim.receipt_emitted", W),
620            Signal::ExternalOutcomeUnknown => ("turnframe.external.outcome_unknown", WE),
621            Signal::ExternalReconciled => ("turnframe.external.reconciled", W),
622            Signal::ProviderFallback => ("turnframe.provider.fallback", PMPE),
623            Signal::ProviderCapabilityMismatch => ("turnframe.provider.capability_mismatch", PMPE),
624            Signal::WorkflowInvariantViolation => ("turnframe.workflow.invariant_violation", WE),
625            Signal::QuestionAnswered => ("turnframe.question.answered", W),
626            Signal::QuestionUnanswered => ("turnframe.question.unanswered", W),
627            Signal::ActSuperseded => ("turnframe.act.superseded", WO),
628            Signal::ActRefused => ("turnframe.act.refused", W),
629            Signal::TurnDuration => ("turnframe.turn.duration_ms", W),
630            Signal::ProjectionDuration => ("turnframe.projection.duration_us", W),
631            Signal::ReductionDuration => ("turnframe.reduction.duration_us", W),
632            Signal::PersistenceDuration => ("turnframe.persistence.duration_ms", W),
633            Signal::ProviderLatency => ("turnframe.provider.latency_ms", PMP),
634            Signal::ExternalLatency => ("turnframe.external.latency_ms", W),
635            Signal::NarrationLatency => ("turnframe.narration.latency_ms", PMP),
636            Signal::TaskCompleted => ("turnframe.task.completed", PMPE),
637            Signal::TaskRepaired => ("turnframe.task.repaired", PMPE),
638            Signal::TaskEscalated => ("turnframe.task.escalated", PMPE),
639            Signal::TaskVoteDisagreement => ("turnframe.task.vote_disagreement", PMP),
640            Signal::BudgetExhausted => ("turnframe.budget.exhausted", WE),
641            Signal::TaskLatency => ("turnframe.task.latency_ms", PMP),
642            _ => ("", NONE),
643        };
644        assert!(!name.is_empty(), "no fixture for {signal:?}");
645        Fixture { name, labels }
646    }
647
648    /// A label set that fills every dimension with a plausible value.
649    fn full_labels() -> SignalLabels {
650        SignalLabels::workflow(WorkflowKey::from("trip"))
651            .with_provider("openai")
652            .with_model("gpt-x")
653            .with_purpose("extract")
654            .with_risk(RiskClass::Destructive)
655            .with_interaction(InteractionKind::ConfirmCommand)
656            .with_error_code("rate_limited")
657            .with_operation("trip.set_name")
658    }
659
660    fn label_names(key: &Key) -> Vec<String> {
661        key.labels().map(|l| l.key().to_owned()).collect()
662    }
663
664    fn label_values(key: &Key) -> Vec<String> {
665        key.labels().map(|l| l.value().to_owned()).collect()
666    }
667
668    #[test]
669    fn every_signal_emits_its_name_and_documented_label_set() {
670        for signal in Signal::ALL {
671            let expected = fixture(signal);
672            let recorder = TestRecorder::default();
673            let shared = Arc::clone(&recorder.shared);
674            metrics::with_local_recorder(&recorder, || {
675                let observer = MetricsObserver::new();
676                observer.observe_duration(&signal, Duration::from_millis(7), &full_labels());
677            });
678
679            let captured = lock(&shared);
680            let key = match signal.kind() {
681                SignalKind::DurationMillis | SignalKind::DurationMicros => {
682                    assert_eq!(captured.counters.len(), 0, "{signal:?} is not a counter");
683                    assert_eq!(captured.histograms.len(), 1, "{signal:?}");
684                    captured.histograms[0].0.clone()
685                }
686                _ => {
687                    assert_eq!(
688                        captured.histograms.len(),
689                        0,
690                        "{signal:?} is not a histogram"
691                    );
692                    assert_eq!(captured.counters.len(), 1, "{signal:?}");
693                    assert_eq!(captured.counters[0].1, 1, "{signal:?}");
694                    captured.counters[0].0.clone()
695                }
696            };
697
698            assert_eq!(key.name(), expected.name, "{signal:?}");
699            assert_eq!(label_names(&key), expected.labels, "{signal:?}");
700        }
701    }
702
703    #[test]
704    fn durations_are_recorded_in_the_unit_of_the_name() {
705        let cases = [
706            (Signal::TurnDuration, 1_500.0_f64),
707            (Signal::PersistenceDuration, 1_500.0),
708            (Signal::ProviderLatency, 1_500.0),
709            (Signal::ExternalLatency, 1_500.0),
710            (Signal::NarrationLatency, 1_500.0),
711            (Signal::TaskLatency, 1_500.0),
712            (Signal::ProjectionDuration, 1_500_000.0),
713            (Signal::ReductionDuration, 1_500_000.0),
714        ];
715        for (signal, expected) in cases {
716            let recorder = TestRecorder::default();
717            let shared = Arc::clone(&recorder.shared);
718            metrics::with_local_recorder(&recorder, || {
719                MetricsObserver::new().observe_duration(
720                    &signal,
721                    Duration::from_millis(1_500),
722                    &SignalLabels::none(),
723                );
724            });
725            let captured = lock(&shared);
726            assert_eq!(captured.histograms.len(), 1, "{signal:?}");
727            assert!(
728                (captured.histograms[0].1 - expected).abs() < f64::EPSILON,
729                "{signal:?}: {} != {expected}",
730                captured.histograms[0].1
731            );
732        }
733    }
734
735    #[test]
736    fn a_duration_signal_without_a_measurement_emits_nothing() {
737        let recorder = TestRecorder::default();
738        let shared = Arc::clone(&recorder.shared);
739        metrics::with_local_recorder(&recorder, || {
740            MetricsObserver::new().observe_labeled(&Signal::TurnDuration, &full_labels());
741        });
742        let captured = lock(&shared);
743        assert!(captured.counters.is_empty());
744        assert!(captured.histograms.is_empty());
745    }
746
747    #[test]
748    fn no_label_value_is_a_uuid_or_longer_than_the_limit() {
749        let adversarial =
750            SignalLabels::workflow(WorkflowKey::from("0191f0f6-7d5b-7c2e-9a1e-2b3c4d5e6f70"))
751                .with_provider("018f4e7c-3d2a-7b19-8f6e-112233445566")
752                .with_model("m".repeat(MAX_LABEL_VALUE_LEN + 1))
753                .with_purpose("the user asked to withdraw the trip for Marta Bianchi")
754                .with_risk(RiskClass::Irreversible)
755                .with_interaction(InteractionKind::Freeform)
756                .with_error_code(String::new());
757
758        for signal in Signal::ALL {
759            let labels = metric_labels(signal, &adversarial);
760            for label in &labels {
761                let value = label.value();
762                assert!(
763                    uuid::Uuid::parse_str(value).is_err(),
764                    "{signal:?}: {value} is a uuid"
765                );
766                assert!(
767                    value.chars().count() <= MAX_LABEL_VALUE_LEN,
768                    "{signal:?}: {value} is too long"
769                );
770                assert!(
771                    !value.chars().any(char::is_whitespace),
772                    "{signal:?}: {value} looks like prose"
773                );
774                assert!(!value.is_empty(), "{signal:?}: empty label value");
775            }
776            // Only the bounded enumerations survive an adversarial label set.
777            let names: Vec<&str> = labels.iter().map(metrics::Label::key).collect();
778            assert!(
779                names.iter().all(|n| *n == "risk" || *n == "interaction"),
780                "{signal:?}: {names:?}"
781            );
782        }
783    }
784
785    #[test]
786    fn safe_label_values_accept_keys_and_codes() {
787        assert!(is_safe_label_value("trip"));
788        assert!(is_safe_label_value("openai"));
789        assert!(is_safe_label_value("gpt-5.4-mini"));
790        assert!(is_safe_label_value("trip.traveler_missing"));
791        assert!(is_safe_label_value("external_regulated"));
792        assert!(!is_safe_label_value(""));
793        assert!(!is_safe_label_value("a b"));
794        assert!(!is_safe_label_value(&"x".repeat(MAX_LABEL_VALUE_LEN + 1)));
795        assert!(is_safe_label_value(&"x".repeat(MAX_LABEL_VALUE_LEN)));
796        assert!(!is_safe_label_value("0191f0f6-7d5b-7c2e-9a1e-2b3c4d5e6f70"));
797        assert!(!is_safe_label_value("0191f0f67d5b7c2e9a1e2b3c4d5e6f70"));
798    }
799
800    #[test]
801    fn labels_not_documented_by_a_signal_are_dropped() {
802        // `turn.received` documents only the workflow, so the provider and the
803        // error code a caller attached never reach the metric.
804        let labels = metric_labels(Signal::TurnReceived, &full_labels());
805        assert_eq!(
806            labels.iter().map(metrics::Label::key).collect::<Vec<_>>(),
807            vec!["workflow"]
808        );
809    }
810
811    #[test]
812    fn describe_all_registers_every_metric_with_its_unit() {
813        let recorder = TestRecorder::default();
814        let shared = Arc::clone(&recorder.shared);
815        metrics::with_local_recorder(&recorder, describe_all);
816
817        let captured = lock(&shared);
818        assert_eq!(captured.described.len(), Signal::ALL.len());
819        for signal in Signal::ALL {
820            let row = captured
821                .described
822                .iter()
823                .find(|(name, ..)| name == signal.name())
824                .unwrap_or_else(|| panic!("{signal:?} not described"));
825            assert_eq!(row.1, metric_type(signal), "{signal:?}");
826            assert_eq!(row.2, Some(unit(signal)), "{signal:?}");
827            assert_eq!(row.3, description(signal), "{signal:?}");
828            assert!(!row.3.is_empty());
829        }
830    }
831
832    #[test]
833    fn observe_without_labels_still_emits_the_metric() {
834        let recorder = TestRecorder::default();
835        let shared = Arc::clone(&recorder.shared);
836        metrics::with_local_recorder(&recorder, || {
837            MetricsObserver::new().observe(&Signal::TurnReceived);
838        });
839        let captured = lock(&shared);
840        assert_eq!(captured.counters.len(), 1);
841        assert_eq!(captured.counters[0].0.name(), "turnframe.turn.received");
842        assert!(label_values(&captured.counters[0].0).is_empty());
843    }
844
845    #[test]
846    fn a_turn_and_its_tasks_are_counted_by_effort() {
847        let labels = SignalLabels::none().with_effort(turnframe_core::effort::Effort::High);
848        for signal in [
849            Signal::TurnCompleted,
850            Signal::TurnDuration,
851            Signal::TaskCompleted,
852            Signal::TaskLatency,
853        ] {
854            assert!(
855                documented_labels(signal).contains(&LabelKey::Effort),
856                "{signal:?}"
857            );
858        }
859        assert_eq!(LabelKey::Effort.value_of(&labels).as_deref(), Some("high"));
860    }
861
862    #[test]
863    fn catalogue_covers_every_signal_and_names_match() {
864        let catalogue = catalogue();
865        assert_eq!(catalogue.len(), Signal::ALL.len());
866        for doc in &catalogue {
867            assert_eq!(doc.name, doc.signal.name());
868            assert_eq!(doc.labels, documented_labels(doc.signal));
869            assert!(!doc.description.is_empty());
870        }
871        let markdown = catalogue_markdown();
872        for signal in Signal::ALL {
873            assert!(markdown.contains(signal.name()), "{signal:?}");
874        }
875    }
876
877    #[test]
878    fn label_keys_have_distinct_names() {
879        let mut names: Vec<&str> = LabelKey::ALL.iter().map(|k| k.as_str()).collect();
880        names.sort_unstable();
881        let unique = names.len();
882        names.dedup();
883        assert_eq!(names.len(), unique);
884    }
885}