1use 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
24pub const METRIC_TARGET: &str = "turnframe";
26
27pub const MAX_LABEL_VALUE_LEN: usize = 64;
30
31static METADATA: Metadata<'static> =
32 Metadata::new(METRIC_TARGET, Level::INFO, Some(module_path!()));
33
34#[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,
44 Provider,
46 Model,
48 Purpose,
50 Risk,
52 Interaction,
54 ErrorCode,
56 Operation,
59 Effort,
61}
62
63impl LabelKey {
64 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 #[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 #[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#[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#[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 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 Signal::TaskCompleted | Signal::TaskRepaired | Signal::TaskEscalated => TASK_ERROR,
179 Signal::TaskVoteDisagreement | Signal::TaskLatency => TASK,
180 Signal::BudgetExhausted => TURN_ERROR,
182 _ => WORKFLOW,
185 }
186}
187
188#[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#[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#[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 _ => "Turnframe signal.",
271 }
272}
273
274#[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#[derive(Debug, Clone, Copy, PartialEq, Eq)]
289#[non_exhaustive]
290pub struct MetricDoc {
291 pub signal: Signal,
293 pub name: &'static str,
295 pub metric_type: &'static str,
297 pub unit: Unit,
299 pub labels: &'static [LabelKey],
301 pub description: &'static str,
303}
304
305impl MetricDoc {
306 #[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 #[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#[must_use]
335pub fn catalogue() -> Vec<MetricDoc> {
336 Signal::ALL.into_iter().map(MetricDoc::of).collect()
337}
338
339#[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#[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
369pub 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#[derive(Debug, Clone, Copy, Default)]
416pub struct MetricsObserver;
417
418impl MetricsObserver {
419 #[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 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 #[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 struct Fixture {
584 name: &'static str,
585 labels: &'static [&'static str],
586 }
587
588 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 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 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 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}