1use std::time::Duration;
15
16use crate::command::RiskClass;
17use crate::ids::{ModelKey, OperationKey, ProviderKey, WorkflowKey};
18use crate::interaction::InteractionKind;
19
20#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
23#[non_exhaustive]
24pub enum Signal {
25 TurnReceived,
27 TurnCompleted,
29 TurnFailed,
31 TargetAmbiguous,
33 TargetMissing,
35 TargetUnresolved,
43 CaseNotAuthorized,
51 CommandConfirmationRequired,
53 CommandExecuted,
55 CommandRejected,
57 CommandIdempotencyReplay,
59 CommandRevisionConflict,
61 InteractionCreated,
63 InteractionResolved,
65 InteractionStale,
67 InteractionFailed,
69 ClaimReceiptEmitted,
71 ExternalOutcomeUnknown,
73 ExternalReconciled,
75 ProviderFallback,
77 ProviderCapabilityMismatch,
79 WorkflowInvariantViolation,
81 QuestionAnswered,
83 QuestionUnanswered,
85 ActRefused,
92 ActSuperseded,
96 TurnDuration,
99 ProjectionDuration,
101 ReductionDuration,
103 PersistenceDuration,
105 ProviderLatency,
107 ExternalLatency,
109 NarrationLatency,
111 TaskCompleted,
113 TaskRepaired,
115 TaskEscalated,
117 TaskVoteDisagreement,
119 BudgetExhausted,
121 TaskLatency,
123}
124
125#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
127#[non_exhaustive]
128pub enum SignalKind {
129 Counter,
131 DurationMillis,
133 DurationMicros,
135}
136
137impl Signal {
138 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 #[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 #[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 #[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#[derive(Debug, Clone, Default, PartialEq, Eq)]
294#[non_exhaustive]
295pub struct SignalLabels {
296 pub workflow: Option<WorkflowKey>,
298 pub provider: Option<ProviderKey>,
300 pub risk: Option<RiskClass>,
302 pub model: Option<ModelKey>,
304 pub purpose: Option<String>,
307 pub interaction: Option<InteractionKind>,
309 pub error_code: Option<String>,
313 pub operation: Option<OperationKey>,
318 pub effort: Option<crate::effort::Effort>,
320}
321
322impl SignalLabels {
323 #[must_use]
325 pub fn none() -> Self {
326 Self::default()
327 }
328
329 #[must_use]
331 pub fn workflow(workflow: WorkflowKey) -> Self {
332 Self {
333 workflow: Some(workflow),
334 ..Self::default()
335 }
336 }
337
338 #[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 #[must_use]
347 pub fn with_risk(mut self, risk: RiskClass) -> Self {
348 self.risk = Some(risk);
349 self
350 }
351
352 #[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 #[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 #[must_use]
368 pub fn with_effort(mut self, effort: crate::effort::Effort) -> Self {
369 self.effort = Some(effort);
370 self
371 }
372
373 #[must_use]
375 pub fn with_interaction(mut self, kind: InteractionKind) -> Self {
376 self.interaction = Some(kind);
377 self
378 }
379
380 #[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 #[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
395pub trait Observer: Send + Sync {
397 fn observe(&self, signal: &Signal);
399
400 fn observe_labeled(&self, signal: &Signal, labels: &SignalLabels) {
402 let _ = labels;
403 self.observe(signal);
404 }
405
406 fn observe_duration(&self, signal: &Signal, duration: Duration, labels: &SignalLabels) {
410 let _ = duration;
411 self.observe_labeled(signal, labels);
412 }
413}
414
415#[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 #[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}