1use std::fmt;
19use std::sync::Arc;
20use std::time::Duration;
21
22use tracing::field::Empty;
23use turnframe_core::hash::{Digest, digest_hex};
24use turnframe_core::ids::{
25 AccountId, BlockId, CaseId, CaseRevision, CommandId, ConversationId, EventId, InteractionId,
26 OutboxId, TurnId, WorkflowKey, WorkflowVersion,
27};
28use turnframe_core::observe::{Observer, Signal, SignalLabels};
29use turnframe_core::replay::{ProviderAttemptRecord, ReplayRecord};
30
31use crate::{attrs, opt_str};
32
33pub const TRACE_TARGET: &str = "turnframe";
35
36const ACCOUNT_HASH_DOMAIN: &str = "turnframe.account";
39
40pub const ACCOUNT_HASH_LEN: usize = 16;
44
45#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, PartialOrd, Ord)]
50#[non_exhaustive]
51pub enum PipelineStage {
52 Projection,
54 Understanding,
56 TargetResolution,
58 Reduction,
60 Policy,
62 Execution,
64 ExternalDispatch,
66 Composition,
68 Narration,
70 Persistence,
72 Reconciliation,
74}
75
76impl PipelineStage {
77 pub const ALL: [Self; 11] = [
79 Self::Projection,
80 Self::Understanding,
81 Self::TargetResolution,
82 Self::Reduction,
83 Self::Policy,
84 Self::Execution,
85 Self::ExternalDispatch,
86 Self::Composition,
87 Self::Narration,
88 Self::Persistence,
89 Self::Reconciliation,
90 ];
91
92 #[must_use]
94 pub const fn as_str(self) -> &'static str {
95 match self {
96 Self::Projection => "projection",
97 Self::Understanding => "understanding",
98 Self::TargetResolution => "target_resolution",
99 Self::Reduction => "reduction",
100 Self::Policy => "policy",
101 Self::Execution => "execution",
102 Self::ExternalDispatch => "external_dispatch",
103 Self::Composition => "composition",
104 Self::Narration => "narration",
105 Self::Persistence => "persistence",
106 Self::Reconciliation => "reconciliation",
107 }
108 }
109}
110
111impl std::fmt::Display for PipelineStage {
112 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
113 f.write_str(self.as_str())
114 }
115}
116
117#[must_use]
123pub fn account_hash(account_id: &AccountId) -> String {
124 let mut material =
125 String::with_capacity(ACCOUNT_HASH_DOMAIN.len() + 1 + account_id.as_str().len());
126 material.push_str(ACCOUNT_HASH_DOMAIN);
127 material.push('\0');
128 material.push_str(account_id.as_str());
129 let mut digest = digest_hex(material.as_bytes());
130 digest.truncate(ACCOUNT_HASH_LEN);
131 digest
132}
133
134#[derive(Debug, Clone, Default, PartialEq, Eq)]
141#[non_exhaustive]
142pub struct TurnIdentifiers {
143 pub turn_id: Option<TurnId>,
145 pub conversation_id: Option<ConversationId>,
147 pub account_hash: Option<String>,
149 pub workflow: Option<WorkflowKey>,
151 pub workflow_version: Option<WorkflowVersion>,
153 pub case_id: Option<CaseId>,
155 pub case_revision: Option<CaseRevision>,
157 pub interaction_id: Option<InteractionId>,
159 pub provider: Option<String>,
161 pub model: Option<String>,
163 pub attempt: Option<String>,
165 pub plan_hash: Option<Digest>,
167 pub command_id: Option<CommandId>,
169 pub event_ids: Vec<EventId>,
171 pub block_ids: Vec<BlockId>,
173 pub outbox_id: Option<OutboxId>,
175}
176
177impl TurnIdentifiers {
178 #[must_use]
180 pub fn of_turn(
181 turn_id: TurnId,
182 conversation_id: ConversationId,
183 account_id: &AccountId,
184 ) -> Self {
185 Self {
186 turn_id: Some(turn_id),
187 conversation_id: Some(conversation_id),
188 account_hash: Some(account_hash(account_id)),
189 ..Self::default()
190 }
191 }
192
193 #[must_use]
199 pub fn from_replay(record: &ReplayRecord) -> Self {
200 let workflow_version = record.workflow_versions.first();
201 let case = record.loaded_cases.first();
202 let attempt = record.provider_attempts.last();
203 let command = record
204 .command_outcomes
205 .first()
206 .map(|outcome| outcome.command_ref.command_id);
207
208 Self {
209 turn_id: Some(record.turn_id),
210 conversation_id: Some(record.conversation_id),
211 account_hash: Some(account_hash(&record.account_id)),
212 workflow: workflow_version.map(|entry| entry.key.clone()),
213 workflow_version: workflow_version.map(|entry| entry.version.clone()),
214 case_id: case.map(|case| case.case_id.clone()),
215 case_revision: case.map(|case| case.expected_revision),
216 interaction_id: record.interactions_created.first().copied(),
217 provider: attempt.map(|attempt| attempt.provider_key.to_string()),
218 model: attempt.map(|attempt| attempt.model_key.to_string()),
219 attempt: attempt.map(|attempt| attempt.attempt.to_string()),
220 plan_hash: record.plan_hash.clone(),
221 command_id: command,
222 event_ids: record.event_ids.clone(),
223 block_ids: record.response_block_ids.clone(),
224 outbox_id: record.outbox_ids.first().copied(),
225 }
226 }
227
228 #[must_use]
234 pub fn fields(&self) -> Vec<(&'static str, String)> {
235 let mut fields: Vec<(&'static str, String)> = Vec::new();
236 let mut push = |name: &'static str, value: Option<String>| {
237 if let Some(value) = value {
238 fields.push((name, value));
239 }
240 };
241
242 push(field::TURN_ID, self.turn_id.map(|id| id.to_string()));
243 push(
244 field::CONVERSATION_ID,
245 self.conversation_id.map(|id| id.to_string()),
246 );
247 push(field::ACCOUNT_HASH, self.account_hash.clone());
248 push(field::WORKFLOW, opt_str(&self.workflow).map(str::to_owned));
249 push(
250 field::WORKFLOW_VERSION,
251 opt_str(&self.workflow_version).map(str::to_owned),
252 );
253 push(field::CASE_ID, opt_str(&self.case_id).map(str::to_owned));
254 push(
255 field::CASE_REVISION,
256 self.case_revision.map(|revision| revision.to_string()),
257 );
258 push(
259 field::INTERACTION_ID,
260 self.interaction_id.map(|id| id.to_string()),
261 );
262 push(field::PROVIDER, opt_str(&self.provider).map(str::to_owned));
263 push(field::MODEL, opt_str(&self.model).map(str::to_owned));
264 push(field::ATTEMPT, opt_str(&self.attempt).map(str::to_owned));
265 push(
266 field::PLAN_HASH,
267 self.plan_hash.as_ref().map(|hash| hash.as_str().to_owned()),
268 );
269 push(field::COMMAND_ID, self.command_id.map(|id| id.to_string()));
270 push(field::EVENT_IDS, join_ids(&self.event_ids));
271 push(field::BLOCK_IDS, join_ids(&self.block_ids));
272 push(field::OUTBOX_ID, self.outbox_id.map(|id| id.to_string()));
273 fields
274 }
275}
276
277fn join_ids<T: ToString>(ids: &[T]) -> Option<String> {
280 if ids.is_empty() {
281 return None;
282 }
283 Some(
284 ids.iter()
285 .map(ToString::to_string)
286 .collect::<Vec<_>>()
287 .join(","),
288 )
289}
290
291pub mod field {
293 pub const TURN_ID: &str = "turn_id";
295 pub const CONVERSATION_ID: &str = "conversation_id";
297 pub const ACCOUNT_HASH: &str = "account_hash";
299 pub const WORKFLOW: &str = "workflow";
301 pub const WORKFLOW_VERSION: &str = "workflow_version";
303 pub const CASE_ID: &str = "case_id";
305 pub const CASE_REVISION: &str = "case_revision";
307 pub const INTERACTION_ID: &str = "interaction_id";
309 pub const PROVIDER: &str = "provider";
311 pub const MODEL: &str = "model";
313 pub const ATTEMPT: &str = "attempt";
315 pub const PLAN_HASH: &str = "plan_hash";
317 pub const COMMAND_ID: &str = "command_id";
319 pub const EVENT_IDS: &str = "event_ids";
321 pub const BLOCK_IDS: &str = "block_ids";
323 pub const OUTBOX_ID: &str = "outbox_id";
325 pub const SIGNAL: &str = "signal";
327 pub const STAGE: &str = "stage";
329 pub const DURATION_MS: &str = "duration_ms";
331 pub const RISK: &str = "risk";
333 pub const INTERACTION: &str = "interaction";
335 pub const PURPOSE: &str = "purpose";
337 pub const EFFORT: &str = "effort";
339 pub const ERROR_CODE: &str = "error_code";
341}
342
343#[must_use]
356pub fn turn_span(
357 turn_id: TurnId,
358 conversation_id: ConversationId,
359 account_id: &AccountId,
360) -> ::tracing::Span {
361 ::tracing::span!(
362 target: TRACE_TARGET,
363 ::tracing::Level::INFO,
364 "turnframe.turn",
365 turn_id = %turn_id,
366 conversation_id = %conversation_id,
367 account_hash = %account_hash(account_id),
368 "session.id" = Empty,
369 "user.id" = Empty,
370 tags = Empty,
371 "deployment.environment.name" = Empty,
372 "service.version" = Empty,
373 )
374}
375
376#[must_use]
386pub fn stage_span(stage: PipelineStage, turn_id: TurnId) -> ::tracing::Span {
387 ::tracing::span!(
388 target: TRACE_TARGET,
389 ::tracing::Level::DEBUG,
390 "turnframe.stage",
391 stage = stage.as_str(),
392 turn_id = %turn_id,
393 "session.id" = Empty,
394 "user.id" = Empty,
395 tags = Empty,
396 "deployment.environment.name" = Empty,
397 "service.version" = Empty,
398 )
399}
400
401#[derive(Debug, Clone, Default, PartialEq)]
422#[non_exhaustive]
423pub struct ProviderCall {
424 pub system: Option<String>,
426 pub operation: Option<String>,
428 pub request_model: Option<String>,
430 pub temperature: Option<f64>,
432 pub response_id: Option<String>,
434 pub response_model: Option<String>,
436 pub finish_reasons: Vec<String>,
438 pub input_tokens: Option<u64>,
440 pub output_tokens: Option<u64>,
442 pub input_cached_tokens: Option<u64>,
444}
445
446impl ProviderCall {
447 #[must_use]
449 pub fn new(
450 system: impl Into<String>,
451 operation: impl Into<String>,
452 request_model: impl Into<String>,
453 ) -> Self {
454 Self {
455 system: Some(system.into()),
456 operation: Some(operation.into()),
457 request_model: Some(request_model.into()),
458 ..Self::default()
459 }
460 }
461
462 #[must_use]
464 pub fn from_attempt(attempt: &ProviderAttemptRecord) -> Self {
465 Self {
466 system: Some(attempt.provider_key.to_string()),
467 operation: Some(attempt.purpose.to_string()),
468 request_model: Some(attempt.model_key.to_string()),
469 temperature: attempt.temperature.map(f64::from),
470 response_id: Some(attempt.request_id.to_string()),
471 response_model: Some(attempt.model_key.to_string()),
472 finish_reasons: attempt.finish_reasons.clone(),
473 input_tokens: attempt.input_tokens,
474 output_tokens: attempt.output_tokens,
475 input_cached_tokens: None,
476 }
477 }
478
479 #[must_use]
481 pub fn with_temperature(mut self, temperature: f64) -> Self {
482 self.temperature = Some(temperature);
483 self
484 }
485
486 #[must_use]
488 pub fn with_response(
489 mut self,
490 response_id: impl Into<String>,
491 response_model: impl Into<String>,
492 ) -> Self {
493 self.response_id = Some(response_id.into());
494 self.response_model = Some(response_model.into());
495 self
496 }
497
498 #[must_use]
500 pub fn with_finish_reason(mut self, reason: impl Into<String>) -> Self {
501 self.finish_reasons.push(reason.into());
502 self
503 }
504
505 #[must_use]
507 pub fn with_usage(mut self, input_tokens: u64, output_tokens: u64) -> Self {
508 self.input_tokens = Some(input_tokens);
509 self.output_tokens = Some(output_tokens);
510 self
511 }
512
513 #[must_use]
528 pub fn with_cached_usage(
529 mut self,
530 total_input_tokens: u64,
531 cached_tokens: u64,
532 output_tokens: u64,
533 ) -> Self {
534 self.input_tokens = Some(total_input_tokens.saturating_sub(cached_tokens));
535 self.input_cached_tokens = Some(cached_tokens);
536 self.output_tokens = Some(output_tokens);
537 self
538 }
539
540 #[must_use]
542 pub fn total_input_tokens(&self) -> Option<u64> {
543 match (self.input_tokens, self.input_cached_tokens) {
544 (None, None) => None,
545 (input, cached) => Some(
546 input
547 .unwrap_or_default()
548 .saturating_add(cached.unwrap_or_default()),
549 ),
550 }
551 }
552
553 #[must_use]
560 pub fn attributes(&self) -> Vec<(&'static str, String)> {
561 let mut out: Vec<(&'static str, String)> = Vec::new();
562 let mut push = |key: &'static str, value: Option<String>| {
563 if let Some(value) = value {
564 out.push((key, value));
565 }
566 };
567 push(attrs::GEN_AI_SYSTEM, self.system.clone());
568 push(attrs::GEN_AI_OPERATION_NAME, self.operation.clone());
569 push(attrs::GEN_AI_REQUEST_MODEL, self.request_model.clone());
570 push(
571 attrs::GEN_AI_REQUEST_TEMPERATURE,
572 self.temperature.map(|value| value.to_string()),
573 );
574 push(attrs::GEN_AI_RESPONSE_ID, self.response_id.clone());
575 push(attrs::GEN_AI_RESPONSE_MODEL, self.response_model.clone());
576 push(
577 attrs::GEN_AI_RESPONSE_FINISH_REASONS,
578 if self.finish_reasons.is_empty() {
579 None
580 } else {
581 Some(self.finish_reasons.join(","))
582 },
583 );
584 push(
585 attrs::GEN_AI_USAGE_INPUT_TOKENS,
586 self.input_tokens.map(|value| value.to_string()),
587 );
588 push(
589 attrs::GEN_AI_USAGE_OUTPUT_TOKENS,
590 self.output_tokens.map(|value| value.to_string()),
591 );
592 push(
593 attrs::GEN_AI_USAGE_INPUT_CACHED_TOKENS,
594 self.input_cached_tokens.map(|value| value.to_string()),
595 );
596 out
597 }
598}
599
600#[must_use]
616pub fn provider_call_span(call: &ProviderCall) -> ::tracing::Span {
617 ::tracing::span!(
618 target: TRACE_TARGET,
619 ::tracing::Level::INFO,
620 "gen_ai.client.operation",
621 "gen_ai.system" = call.system.as_deref(),
622 "gen_ai.operation.name" = call.operation.as_deref(),
623 "gen_ai.request.model" = call.request_model.as_deref(),
624 "gen_ai.request.temperature" = call.temperature,
625 "gen_ai.response.id" = call.response_id.as_deref(),
626 "gen_ai.response.model" = call.response_model.as_deref(),
627 "gen_ai.response.finish_reasons" = (!call.finish_reasons.is_empty())
628 .then(|| call.finish_reasons.join(",")),
629 "gen_ai.usage.input_tokens" = call.input_tokens,
630 "gen_ai.usage.output_tokens" = call.output_tokens,
631 "gen_ai.usage.input_cached_tokens" = call.input_cached_tokens,
632 "gen_ai.input.messages" = Empty,
633 "gen_ai.output.messages" = Empty,
634 "session.id" = Empty,
635 "user.id" = Empty,
636 tags = Empty,
637 "deployment.environment.name" = Empty,
638 "service.version" = Empty,
639 )
640}
641
642#[derive(Debug, Clone, Default, PartialEq, Eq)]
670#[non_exhaustive]
671pub struct TraceGrouping {
672 pub session_id: Option<String>,
674 pub end_user_hash: Option<String>,
676 pub tags: Vec<String>,
678 pub environment: Option<String>,
680 pub release: Option<String>,
682}
683
684impl TraceGrouping {
685 #[must_use]
687 pub fn for_conversation(conversation_id: ConversationId) -> Self {
688 Self {
689 session_id: Some(conversation_id.to_string()),
690 ..Self::default()
691 }
692 }
693
694 #[must_use]
697 pub fn with_account(mut self, account_id: &AccountId) -> Self {
698 self.end_user_hash = Some(account_hash(account_id));
699 self
700 }
701
702 #[must_use]
704 pub fn with_end_user_hash(mut self, hash: impl Into<String>) -> Self {
705 self.end_user_hash = Some(hash.into());
706 self
707 }
708
709 #[must_use]
711 pub fn with_tag(mut self, tag: impl Into<String>) -> Self {
712 self.tags.push(tag.into());
713 self
714 }
715
716 #[must_use]
718 pub fn with_environment(mut self, environment: impl Into<String>) -> Self {
719 self.environment = Some(environment.into());
720 self
721 }
722
723 #[must_use]
725 pub fn with_release(mut self, release: impl Into<String>) -> Self {
726 self.release = Some(release.into());
727 self
728 }
729
730 #[must_use]
734 pub fn fields(&self) -> Vec<(&'static str, String)> {
735 let mut out: Vec<(&'static str, String)> = Vec::new();
736 let mut push = |key: &'static str, value: Option<String>| {
737 if let Some(value) = value {
738 out.push((key, value));
739 }
740 };
741 push(attrs::SESSION_ID, self.session_id.clone());
742 push(attrs::USER_ID, self.end_user_hash.clone());
743 push(
744 attrs::TAGS,
745 if self.tags.is_empty() {
746 None
747 } else {
748 Some(self.tags.join(","))
749 },
750 );
751 push(attrs::DEPLOYMENT_ENVIRONMENT, self.environment.clone());
752 push(attrs::SERVICE_VERSION, self.release.clone());
753 out
754 }
755
756 pub fn stamp(&self, span: &::tracing::Span) {
762 for (key, value) in self.fields() {
763 span.record(key, value.as_str());
764 }
765 }
766
767 pub fn stamp_current(&self) {
769 self.stamp(&::tracing::Span::current());
770 }
771
772 #[must_use]
789 pub fn scope_span(&self) -> ::tracing::Span {
790 let span = ::tracing::span!(
791 target: TRACE_TARGET,
792 ::tracing::Level::INFO,
793 "turnframe.grouping",
794 "session.id" = Empty,
795 "user.id" = Empty,
796 tags = Empty,
797 "deployment.environment.name" = Empty,
798 "service.version" = Empty,
799 );
800 self.stamp(&span);
801 span
802 }
803}
804
805#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, PartialOrd, Ord)]
807#[non_exhaustive]
808pub enum ContentRole {
809 Input,
811 Output,
813}
814
815impl ContentRole {
816 #[must_use]
818 pub const fn attribute(self) -> &'static str {
819 match self {
820 Self::Input => attrs::GEN_AI_INPUT_MESSAGES,
821 Self::Output => attrs::GEN_AI_OUTPUT_MESSAGES,
822 }
823 }
824}
825
826pub trait ContentRedactor: Send + Sync + fmt::Debug {
832 fn redact(&self, role: ContentRole, text: &str) -> Option<String>;
834}
835
836#[derive(Debug, Clone, Copy, Default)]
839pub struct DropAllContent;
840
841impl ContentRedactor for DropAllContent {
842 fn redact(&self, _role: ContentRole, _text: &str) -> Option<String> {
843 None
844 }
845}
846
847#[derive(Clone)]
872pub struct ContentRecorder {
873 enabled: bool,
874 redactor: Arc<dyn ContentRedactor>,
875}
876
877impl ContentRecorder {
878 #[must_use]
881 pub fn disabled() -> Self {
882 Self {
883 enabled: false,
884 redactor: Arc::new(DropAllContent),
885 }
886 }
887
888 #[must_use]
890 pub fn enabled(redactor: Arc<dyn ContentRedactor>) -> Self {
891 Self {
892 enabled: true,
893 redactor,
894 }
895 }
896
897 #[must_use]
899 pub fn is_enabled(&self) -> bool {
900 self.enabled
901 }
902
903 #[must_use]
907 pub fn rendered(&self, role: ContentRole, text: &str) -> Option<String> {
908 if !self.enabled {
909 return None;
910 }
911 self.redactor.redact(role, text)
912 }
913
914 pub fn record(&self, span: &::tracing::Span, role: ContentRole, text: &str) {
918 if let Some(rendered) = self.rendered(role, text) {
919 span.record(role.attribute(), rendered.as_str());
920 }
921 }
922}
923
924impl Default for ContentRecorder {
925 fn default() -> Self {
926 Self::disabled()
927 }
928}
929
930impl fmt::Debug for ContentRecorder {
931 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
932 f.debug_struct("ContentRecorder")
933 .field("enabled", &self.enabled)
934 .field("redactor", &self.redactor)
935 .finish()
936 }
937}
938
939pub fn record_turn(ids: &TurnIdentifiers) {
944 ::tracing::event!(
945 target: TRACE_TARGET,
946 ::tracing::Level::INFO,
947 turn_id = ids.turn_id.map(|id| id.to_string()),
948 conversation_id = ids.conversation_id.map(|id| id.to_string()),
949 account_hash = ids.account_hash.as_deref(),
950 workflow = opt_str(&ids.workflow),
951 workflow_version = opt_str(&ids.workflow_version),
952 case_id = opt_str(&ids.case_id),
953 case_revision = ids.case_revision.map(|revision| revision.value()),
954 interaction_id = ids.interaction_id.map(|id| id.to_string()),
955 provider = opt_str(&ids.provider),
956 model = opt_str(&ids.model),
957 attempt = opt_str(&ids.attempt),
958 plan_hash = ids.plan_hash.as_ref().map(Digest::as_str),
959 command_id = ids.command_id.map(|id| id.to_string()),
960 event_ids = join_ids(&ids.event_ids),
961 block_ids = join_ids(&ids.block_ids),
962 outbox_id = ids.outbox_id.map(|id| id.to_string()),
963 "turn recorded",
964 );
965}
966
967#[must_use]
973pub fn signal_fields(labels: &SignalLabels) -> Vec<(&'static str, String)> {
974 let mut fields: Vec<(&'static str, String)> = Vec::new();
975 let mut push = |name: &'static str, value: Option<String>| {
976 if let Some(value) = value {
977 fields.push((name, value));
978 }
979 };
980 push(
981 field::WORKFLOW,
982 opt_str(&labels.workflow).map(str::to_owned),
983 );
984 push(
985 field::PROVIDER,
986 opt_str(&labels.provider).map(str::to_owned),
987 );
988 push(field::MODEL, opt_str(&labels.model).map(str::to_owned));
989 push(field::PURPOSE, opt_str(&labels.purpose).map(str::to_owned));
990 push(
991 field::RISK,
992 labels.risk.as_ref().and_then(crate::enum_label),
993 );
994 push(
995 field::INTERACTION,
996 labels.interaction.as_ref().and_then(crate::enum_label),
997 );
998 push(
999 field::ERROR_CODE,
1000 opt_str(&labels.error_code).map(str::to_owned),
1001 );
1002 push(
1003 field::EFFORT,
1004 labels.effort.map(|effort| effort.as_str().to_owned()),
1005 );
1006 fields
1007}
1008
1009#[derive(Debug, Clone, Copy, Default)]
1031pub struct TracingObserver;
1032
1033impl TracingObserver {
1034 #[must_use]
1037 pub const fn new() -> Self {
1038 Self
1039 }
1040
1041 fn emit(signal: Signal, labels: &SignalLabels, duration: Option<Duration>) {
1042 let duration_ms = duration.map(|value| value.as_secs_f64() * 1_000.0);
1043 if signal.is_safety_signal() {
1044 ::tracing::event!(
1045 target: TRACE_TARGET,
1046 ::tracing::Level::WARN,
1047 signal = signal.name(),
1048 workflow = opt_str(&labels.workflow),
1049 provider = opt_str(&labels.provider),
1050 model = opt_str(&labels.model),
1051 purpose = opt_str(&labels.purpose),
1052 risk = labels.risk.as_ref().and_then(crate::enum_label),
1053 interaction = labels.interaction.as_ref().and_then(crate::enum_label),
1054 error_code = opt_str(&labels.error_code),
1055 effort = labels.effort.map(|effort| effort.as_str()),
1056 duration_ms = duration_ms,
1057 "turnframe safety signal",
1058 );
1059 } else {
1060 ::tracing::event!(
1061 target: TRACE_TARGET,
1062 ::tracing::Level::DEBUG,
1063 signal = signal.name(),
1064 workflow = opt_str(&labels.workflow),
1065 provider = opt_str(&labels.provider),
1066 model = opt_str(&labels.model),
1067 purpose = opt_str(&labels.purpose),
1068 risk = labels.risk.as_ref().and_then(crate::enum_label),
1069 interaction = labels.interaction.as_ref().and_then(crate::enum_label),
1070 error_code = opt_str(&labels.error_code),
1071 effort = labels.effort.map(|effort| effort.as_str()),
1072 duration_ms = duration_ms,
1073 "turnframe signal",
1074 );
1075 }
1076 }
1077}
1078
1079impl Observer for TracingObserver {
1080 fn observe(&self, signal: &Signal) {
1081 Self::emit(*signal, &SignalLabels::none(), None);
1082 }
1083
1084 fn observe_labeled(&self, signal: &Signal, labels: &SignalLabels) {
1085 Self::emit(*signal, labels, None);
1086 }
1087
1088 fn observe_duration(&self, signal: &Signal, duration: Duration, labels: &SignalLabels) {
1089 Self::emit(*signal, labels, Some(duration));
1090 }
1091}
1092
1093#[cfg(test)]
1094mod tests {
1095 use chrono::DateTime;
1096 use turnframe_core::case::CaseRef;
1097 use turnframe_core::command::RiskClass;
1098 use turnframe_core::ids::{ModelKey, ProviderKey};
1099 use turnframe_core::interaction::InteractionKind;
1100 use turnframe_core::replay::{ProviderAttemptOutcome, WorkflowVersionRecord};
1101
1102 use super::*;
1103
1104 fn replay() -> ReplayRecord {
1105 let now = DateTime::from_timestamp(1_700_000_000, 0).expect("valid timestamp");
1106 let mut record = ReplayRecord::received(
1107 TurnId::nil(),
1108 ConversationId::nil(),
1109 AccountId::from("acct-1"),
1110 now,
1111 );
1112 record.workflow_versions.push(WorkflowVersionRecord {
1113 key: WorkflowKey::from("trip"),
1114 version: WorkflowVersion::from("3"),
1115 });
1116 record.loaded_cases.push(CaseRef::new(
1117 WorkflowKey::from("trip"),
1118 CaseId::from("trip-7"),
1119 CaseRevision(4),
1120 ));
1121 record.interactions_created.push(InteractionId::nil());
1122 record.plan_hash = Some(Digest(String::from("abc123")));
1123 record.event_ids.push(EventId::nil());
1124 record.response_block_ids.push(BlockId::from("b1"));
1125 record.response_block_ids.push(BlockId::from("b2"));
1126 record
1127 }
1128
1129 #[test]
1130 fn account_hash_is_stable_short_and_hides_the_raw_id() {
1131 let account = AccountId::from("acct-1");
1132 let hash = account_hash(&account);
1133 assert_eq!(hash.len(), ACCOUNT_HASH_LEN);
1134 assert!(hash.chars().all(|c| c.is_ascii_hexdigit()));
1135 assert_eq!(hash, account_hash(&account));
1136 assert!(!hash.contains("acct"));
1137 assert_ne!(hash, account_hash(&AccountId::from("acct-2")));
1138 }
1139
1140 #[test]
1141 fn account_hash_is_domain_separated() {
1142 let plain = digest_hex(b"acct-1")[..ACCOUNT_HASH_LEN].to_owned();
1144 assert_ne!(account_hash(&AccountId::from("acct-1")), plain);
1145 }
1146
1147 #[test]
1148 fn identifiers_from_replay_carry_the_stable_ids_of_26_1() {
1149 let ids = TurnIdentifiers::from_replay(&replay());
1150 assert_eq!(ids.turn_id, Some(TurnId::nil()));
1151 assert_eq!(ids.conversation_id, Some(ConversationId::nil()));
1152 assert_eq!(ids.workflow.as_ref().map(WorkflowKey::as_str), Some("trip"));
1153 assert_eq!(
1154 ids.workflow_version.as_ref().map(WorkflowVersion::as_str),
1155 Some("3")
1156 );
1157 assert_eq!(ids.case_id.as_ref().map(CaseId::as_str), Some("trip-7"));
1158 assert_eq!(ids.case_revision, Some(CaseRevision(4)));
1159 assert_eq!(ids.interaction_id, Some(InteractionId::nil()));
1160 assert_eq!(ids.plan_hash.as_ref().map(Digest::as_str), Some("abc123"));
1161 assert_eq!(ids.event_ids.len(), 1);
1162 assert_eq!(ids.block_ids.len(), 2);
1163 assert_eq!(
1164 ids.account_hash,
1165 Some(account_hash(&AccountId::from("acct-1")))
1166 );
1167 }
1168
1169 #[test]
1170 fn fields_are_ordered_and_omit_what_is_unknown() {
1171 let ids = TurnIdentifiers::from_replay(&replay());
1172 let fields = ids.fields();
1173 let names: Vec<&str> = fields.iter().map(|(name, _)| *name).collect();
1174 assert_eq!(
1175 names,
1176 vec![
1177 field::TURN_ID,
1178 field::CONVERSATION_ID,
1179 field::ACCOUNT_HASH,
1180 field::WORKFLOW,
1181 field::WORKFLOW_VERSION,
1182 field::CASE_ID,
1183 field::CASE_REVISION,
1184 field::INTERACTION_ID,
1185 field::PLAN_HASH,
1186 field::EVENT_IDS,
1187 field::BLOCK_IDS,
1188 ]
1189 );
1190 let by_name = |name: &str| {
1191 fields
1192 .iter()
1193 .find(|(field, _)| *field == name)
1194 .map(|(_, value)| value.clone())
1195 };
1196 assert_eq!(by_name(field::CASE_ID).as_deref(), Some("trip-7"));
1197 assert_eq!(by_name(field::CASE_REVISION).as_deref(), Some("4"));
1198 assert_eq!(by_name(field::BLOCK_IDS).as_deref(), Some("b1,b2"));
1199 }
1200
1201 #[test]
1202 fn fields_of_an_empty_identifier_set_are_empty() {
1203 assert!(TurnIdentifiers::default().fields().is_empty());
1204 }
1205
1206 #[test]
1207 fn of_turn_hashes_the_account() {
1208 let ids = TurnIdentifiers::of_turn(
1209 TurnId::nil(),
1210 ConversationId::nil(),
1211 &AccountId::from("acct-9"),
1212 );
1213 let account = ids.account_hash.clone().expect("hashed");
1214 assert_eq!(account, account_hash(&AccountId::from("acct-9")));
1215 let rendered = ids.fields();
1216 assert!(rendered.iter().all(|(_, value)| value != "acct-9"));
1217 }
1218
1219 #[test]
1220 fn signal_fields_render_only_typed_labels() {
1221 let labels = SignalLabels::workflow(WorkflowKey::from("trip"))
1222 .with_provider("openai")
1223 .with_model("gpt-x")
1224 .with_purpose("extract")
1225 .with_risk(RiskClass::Destructive)
1226 .with_interaction(InteractionKind::ConfirmCommand)
1227 .with_error_code("rate_limited");
1228 assert_eq!(
1229 signal_fields(&labels),
1230 vec![
1231 (field::WORKFLOW, String::from("trip")),
1232 (field::PROVIDER, String::from("openai")),
1233 (field::MODEL, String::from("gpt-x")),
1234 (field::PURPOSE, String::from("extract")),
1235 (field::RISK, String::from("destructive")),
1236 (field::INTERACTION, String::from("confirm_command")),
1237 (field::ERROR_CODE, String::from("rate_limited")),
1238 ]
1239 );
1240 assert!(signal_fields(&SignalLabels::none()).is_empty());
1241 }
1242
1243 #[test]
1244 fn signal_fields_carry_the_turns_effort() {
1245 let labels = SignalLabels::none().with_effort(turnframe_core::effort::Effort::Low);
1246 assert_eq!(
1247 signal_fields(&labels),
1248 vec![(field::EFFORT, String::from("low"))]
1249 );
1250 }
1251
1252 #[test]
1253 fn join_ids_is_empty_for_no_ids() {
1254 assert_eq!(join_ids::<BlockId>(&[]), None);
1255 assert_eq!(
1256 join_ids(&[BlockId::from("a"), BlockId::from("b")]).as_deref(),
1257 Some("a,b")
1258 );
1259 }
1260
1261 #[test]
1262 fn stage_names_are_distinct() {
1263 let mut names: Vec<&str> = PipelineStage::ALL.iter().map(|s| s.as_str()).collect();
1264 names.sort_unstable();
1265 let total = names.len();
1266 names.dedup();
1267 assert_eq!(names.len(), total);
1268 assert_eq!(PipelineStage::Understanding.to_string(), "understanding");
1269 }
1270
1271 #[test]
1272 fn provider_call_attributes_use_the_genai_keys_in_order() {
1273 let call = ProviderCall::new("openai", "extract", "gpt-x")
1274 .with_temperature(0.2)
1275 .with_response("resp-1", "gpt-x-2026-05")
1276 .with_finish_reason("stop")
1277 .with_finish_reason("length")
1278 .with_cached_usage(1_000, 800, 120);
1279
1280 assert_eq!(
1281 call.attributes(),
1282 vec![
1283 (attrs::GEN_AI_SYSTEM, String::from("openai")),
1284 (attrs::GEN_AI_OPERATION_NAME, String::from("extract")),
1285 (attrs::GEN_AI_REQUEST_MODEL, String::from("gpt-x")),
1286 (attrs::GEN_AI_REQUEST_TEMPERATURE, String::from("0.2")),
1287 (attrs::GEN_AI_RESPONSE_ID, String::from("resp-1")),
1288 (attrs::GEN_AI_RESPONSE_MODEL, String::from("gpt-x-2026-05")),
1289 (
1290 attrs::GEN_AI_RESPONSE_FINISH_REASONS,
1291 String::from("stop,length")
1292 ),
1293 (attrs::GEN_AI_USAGE_INPUT_TOKENS, String::from("200")),
1294 (attrs::GEN_AI_USAGE_OUTPUT_TOKENS, String::from("120")),
1295 (attrs::GEN_AI_USAGE_INPUT_CACHED_TOKENS, String::from("800")),
1296 ]
1297 );
1298 }
1299
1300 #[test]
1301 fn input_tokens_are_net_of_cached_tokens() {
1302 let call = ProviderCall::default().with_cached_usage(1_000, 800, 10);
1303 assert_eq!(call.input_tokens, Some(200));
1304 assert_eq!(call.input_cached_tokens, Some(800));
1305 assert_eq!(call.total_input_tokens(), Some(1_000));
1307
1308 let odd = ProviderCall::default().with_cached_usage(10, 40, 1);
1310 assert_eq!(odd.input_tokens, Some(0));
1311
1312 let plain = ProviderCall::default().with_usage(300, 20);
1314 assert_eq!(plain.input_tokens, Some(300));
1315 assert_eq!(plain.input_cached_tokens, None);
1316 assert_eq!(plain.total_input_tokens(), Some(300));
1317 assert_eq!(ProviderCall::default().total_input_tokens(), None);
1318 }
1319
1320 #[test]
1321 fn a_provider_call_carries_no_content() {
1322 let call = ProviderCall::new("openai", "extract", "gpt-x");
1323 let keys: Vec<&str> = call.attributes().into_iter().map(|(key, _)| key).collect();
1324 assert!(!keys.contains(&attrs::GEN_AI_INPUT_MESSAGES));
1325 assert!(!keys.contains(&attrs::GEN_AI_OUTPUT_MESSAGES));
1326 }
1327
1328 #[test]
1329 fn a_provider_call_can_be_built_from_a_persisted_attempt() {
1330 let mut record = replay();
1331 record.provider_attempts.push(ProviderAttemptRecord {
1332 attempt: 2,
1333 purpose: String::from("extract"),
1334 provider_key: ProviderKey::from("openai"),
1335 model_key: ModelKey::from("gpt-x"),
1336 request_id: String::from("req-1"),
1337 prompt_version: None,
1338 prompt_ref: None,
1339 outcome: ProviderAttemptOutcome::Succeeded,
1340 latency_ms: Some(120),
1341 input_tokens: Some(300),
1342 output_tokens: Some(40),
1343 temperature: Some(0.2),
1344 finish_reasons: vec![String::from("stop")],
1345 });
1346 let attempt = record.provider_attempts.last().expect("attempt");
1347 let call = ProviderCall::from_attempt(attempt);
1348 assert_eq!(call.system.as_deref(), Some("openai"));
1349 assert_eq!(call.operation.as_deref(), Some("extract"));
1350 assert_eq!(call.request_model.as_deref(), Some("gpt-x"));
1351 assert_eq!(call.response_id.as_deref(), Some("req-1"));
1352 assert_eq!(call.input_tokens, Some(300));
1353 assert_eq!(call.output_tokens, Some(40));
1354 assert!(call.temperature.is_some_and(|t| (t - 0.2).abs() < 1e-6));
1357 assert_eq!(call.finish_reasons, vec![String::from("stop")]);
1358 let attributes = call.attributes();
1359 assert!(
1360 attributes
1361 .iter()
1362 .any(|(key, _)| *key == crate::attrs::GEN_AI_REQUEST_TEMPERATURE)
1363 );
1364 assert!(
1365 attributes
1366 .iter()
1367 .any(|(key, _)| *key == crate::attrs::GEN_AI_RESPONSE_FINISH_REASONS)
1368 );
1369
1370 let ids = TurnIdentifiers::from_replay(&record);
1371 assert_eq!(ids.attempt.as_deref(), Some("2"));
1372 assert_eq!(ids.provider.as_deref(), Some("openai"));
1373 }
1374
1375 #[test]
1376 fn the_outbox_identifier_comes_from_the_replay_record() {
1377 let mut record = replay();
1378 assert_eq!(
1379 TurnIdentifiers::from_replay(&record).outbox_id,
1380 None,
1381 "a turn with no external effect names no outbox row"
1382 );
1383
1384 let enqueued = OutboxId::new();
1385 record.outbox_ids.push(enqueued);
1386 let ids = TurnIdentifiers::from_replay(&record);
1387 assert_eq!(ids.outbox_id, Some(enqueued));
1388 assert!(
1389 ids.fields()
1390 .iter()
1391 .any(|(key, value)| *key == field::OUTBOX_ID && value == &enqueued.to_string())
1392 );
1393 }
1394
1395 #[test]
1396 fn trace_grouping_fields_are_vendor_neutral_and_hashed() {
1397 let grouping = TraceGrouping::for_conversation(ConversationId::nil())
1398 .with_account(&AccountId::from("acct-1"))
1399 .with_environment("production")
1400 .with_release("v0.1.0")
1401 .with_tag("trip")
1402 .with_tag("beta");
1403
1404 assert_eq!(
1405 grouping.fields(),
1406 vec![
1407 (attrs::SESSION_ID, ConversationId::nil().to_string()),
1408 (attrs::USER_ID, account_hash(&AccountId::from("acct-1"))),
1409 (attrs::TAGS, String::from("trip,beta")),
1410 (attrs::DEPLOYMENT_ENVIRONMENT, String::from("production")),
1411 (attrs::SERVICE_VERSION, String::from("v0.1.0")),
1412 ]
1413 );
1414 assert!(grouping.fields().iter().all(|(_, value)| value != "acct-1"));
1416 assert_eq!(
1417 grouping.end_user_hash.as_deref(),
1418 Some(account_hash(&AccountId::from("acct-1")).as_str())
1419 );
1420 }
1421
1422 #[test]
1423 fn an_empty_grouping_stamps_nothing() {
1424 assert!(TraceGrouping::default().fields().is_empty());
1425 }
1426
1427 #[test]
1428 fn a_grouping_can_take_an_already_hashed_end_user() {
1429 let grouping = TraceGrouping::default().with_end_user_hash("deadbeefdeadbeef");
1430 assert_eq!(
1431 grouping.fields(),
1432 vec![(attrs::USER_ID, String::from("deadbeefdeadbeef"))]
1433 );
1434 }
1435
1436 #[test]
1437 fn stamping_a_grouping_is_harmless_without_a_subscriber() {
1438 let grouping = TraceGrouping::for_conversation(ConversationId::nil()).with_tag("trip");
1439 let span = grouping.scope_span();
1440 let _entered = span.enter();
1441 grouping.stamp_current();
1442 grouping.stamp(&turn_span(
1443 TurnId::nil(),
1444 ConversationId::nil(),
1445 &AccountId::from("acct-1"),
1446 ));
1447 grouping.stamp(&stage_span(PipelineStage::Understanding, TurnId::nil()));
1448 grouping.stamp(&provider_call_span(&ProviderCall::new(
1449 "openai", "extract", "gpt-x",
1450 )));
1451 }
1452
1453 #[test]
1454 fn content_recording_is_off_by_default() {
1455 let recorder = ContentRecorder::default();
1456 assert!(!recorder.is_enabled());
1457 assert_eq!(
1458 recorder.rendered(ContentRole::Input, "withdraw trip 17"),
1459 None
1460 );
1461 assert_eq!(recorder.rendered(ContentRole::Output, "done"), None);
1462 assert!(!ContentRecorder::disabled().is_enabled());
1463 assert!(format!("{recorder:?}").contains("enabled: false"));
1464 }
1465
1466 #[derive(Debug)]
1468 struct LengthOnly;
1469
1470 impl ContentRedactor for LengthOnly {
1471 fn redact(&self, role: ContentRole, text: &str) -> Option<String> {
1472 match role {
1473 ContentRole::Input => Some(format!("{} chars", text.chars().count())),
1474 _ => None,
1475 }
1476 }
1477 }
1478
1479 #[test]
1480 fn enabled_content_still_goes_through_the_redactor() {
1481 let recorder = ContentRecorder::enabled(Arc::new(LengthOnly));
1482 assert!(recorder.is_enabled());
1483 assert_eq!(
1484 recorder
1485 .rendered(ContentRole::Input, "withdraw trip 17")
1486 .as_deref(),
1487 Some("16 chars")
1488 );
1489 assert_eq!(recorder.rendered(ContentRole::Output, "done"), None);
1491
1492 let strict = ContentRecorder::enabled(Arc::new(DropAllContent));
1494 assert!(strict.is_enabled());
1495 assert_eq!(
1496 strict.rendered(ContentRole::Input, "withdraw trip 17"),
1497 None
1498 );
1499
1500 let span = provider_call_span(&ProviderCall::new("openai", "extract", "gpt-x"));
1501 recorder.record(&span, ContentRole::Input, "withdraw trip 17");
1502 recorder.record(&span, ContentRole::Output, "done");
1503 }
1504
1505 #[test]
1506 fn content_roles_map_to_the_genai_keys() {
1507 assert_eq!(ContentRole::Input.attribute(), attrs::GEN_AI_INPUT_MESSAGES);
1508 assert_eq!(
1509 ContentRole::Output.attribute(),
1510 attrs::GEN_AI_OUTPUT_MESSAGES
1511 );
1512 }
1513
1514 #[test]
1515 fn observers_accept_every_signal_without_a_subscriber() {
1516 let observer = TracingObserver::new();
1517 for signal in Signal::ALL {
1518 observer.observe(&signal);
1519 observer.observe_labeled(&signal, &SignalLabels::none());
1520 observer.observe_duration(&signal, Duration::from_millis(1), &SignalLabels::none());
1521 }
1522 record_turn(&TurnIdentifiers::from_replay(&replay()));
1523 }
1524}