Skip to main content

turnframe_telemetry/
tracing.rs

1//! Structured tracing: the end-to-end trace of one turn (spec §26.1).
2//!
3//! A metric says how often something happened; a trace says which turn it
4//! happened to. This module carries the stable identifiers spec §26.1 asks for
5//! — turn, conversation, workflow key and version, case id and revision,
6//! interaction id, provider, model, attempt, plan hash, command id, event ids,
7//! response block ids and outbox id — as structured `tracing` fields.
8//!
9//! Two rules hold everywhere here:
10//!
11//! * **Identifiers and codes only.** Nothing the user wrote, nothing a case
12//!   contains and no provider payload is ever a field. [`TurnIdentifiers`] has
13//!   no free-text member, so there is no place to put one.
14//! * **The account id is never logged raw.** [`account_hash`] logs a truncated
15//!   digest instead, which correlates a tenant's turns without naming the
16//!   tenant (spec §25.5).
17
18use 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
33/// Tracing target of every event and span this crate emits.
34pub const TRACE_TARGET: &str = "turnframe";
35
36/// Domain separator for the account digest, so the same account id hashed for
37/// another purpose does not produce the same value.
38const ACCOUNT_HASH_DOMAIN: &str = "turnframe.account";
39
40/// Number of hexadecimal characters kept from the account digest. Sixteen
41/// characters (64 bits) keep collisions negligible at any realistic tenant
42/// count while staying short enough to read in a log line.
43pub const ACCOUNT_HASH_LEN: usize = 16;
44
45/// One stage of the turn pipeline, used as the name of a child span.
46///
47/// The set is the deterministic path of the architecture, so a trace of a turn
48/// reads as the pipeline it went through.
49#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, PartialOrd, Ord)]
50#[non_exhaustive]
51pub enum PipelineStage {
52    /// Projecting persisted state into workflow views.
53    Projection,
54    /// The small model tasks that understand the message.
55    Understanding,
56    /// Server-side target resolution.
57    TargetResolution,
58    /// Whole-turn reduction.
59    Reduction,
60    /// Policy evaluation over the planned commands.
61    Policy,
62    /// Command execution and commit.
63    Execution,
64    /// Dispatch of an external side effect.
65    ExternalDispatch,
66    /// Composition of the response blocks.
67    Composition,
68    /// The tasks that write and review the reply.
69    Narration,
70    /// Persistence of the turn and its replay record.
71    Persistence,
72    /// Reconciliation of an unknown external outcome.
73    Reconciliation,
74}
75
76impl PipelineStage {
77    /// Every stage, in pipeline order.
78    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    /// The stage name as it appears on the span.
93    #[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/// Returns a short, domain-separated digest of an account id.
118///
119/// The raw tenant identifier never reaches a log line; the digest is stable for
120/// the life of the account, so every turn of one tenant still groups together
121/// during an investigation.
122#[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/// The stable identifiers of one turn (spec §26.1).
135///
136/// Every member is an identifier, a version label or a digest. Build one with
137/// [`TurnIdentifiers::of_turn`] and fill in what a stage knows, or derive a
138/// whole one from a persisted [`ReplayRecord`] with
139/// [`TurnIdentifiers::from_replay`].
140#[derive(Debug, Clone, Default, PartialEq, Eq)]
141#[non_exhaustive]
142pub struct TurnIdentifiers {
143    /// The turn.
144    pub turn_id: Option<TurnId>,
145    /// The conversation the turn belongs to.
146    pub conversation_id: Option<ConversationId>,
147    /// Digest of the tenant, from [`account_hash`]. Never the raw account id.
148    pub account_hash: Option<String>,
149    /// Workflow in force.
150    pub workflow: Option<WorkflowKey>,
151    /// Version of that workflow.
152    pub workflow_version: Option<WorkflowVersion>,
153    /// Case the turn acted on.
154    pub case_id: Option<CaseId>,
155    /// Revision that case was loaded at.
156    pub case_revision: Option<CaseRevision>,
157    /// Interaction created or answered.
158    pub interaction_id: Option<InteractionId>,
159    /// Provider key of the model call.
160    pub provider: Option<String>,
161    /// Model key of the model call.
162    pub model: Option<String>,
163    /// Attempt within the stage.
164    pub attempt: Option<String>,
165    /// Hash of the accepted, normalized plan.
166    pub plan_hash: Option<Digest>,
167    /// Command the event concerns.
168    pub command_id: Option<CommandId>,
169    /// Events committed by the turn.
170    pub event_ids: Vec<EventId>,
171    /// Response block ids, in order.
172    pub block_ids: Vec<BlockId>,
173    /// Outbox row of an external side effect.
174    pub outbox_id: Option<OutboxId>,
175}
176
177impl TurnIdentifiers {
178    /// The identifiers every stage of a turn knows.
179    #[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    /// Everything spec §26.1 asks for that a persisted replay record holds.
194    ///
195    /// The workflow, case and interaction taken are the first of each list and
196    /// the provider attempt the last, which is the one a failure concerns; the
197    /// full lists stay in the replay record itself.
198    #[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    /// Renders the identifiers as `(field name, value)` pairs in a stable
229    /// order, omitting what is not known.
230    ///
231    /// This is the exact set a subscriber sees; it is a plain function so a
232    /// test can assert on it without standing up a subscriber.
233    #[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
277/// Joins a list of identifiers into one comma-separated field value, or `None`
278/// when the list is empty.
279fn 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
291/// Names of the structured fields this module emits (spec §26.1).
292pub mod field {
293    /// The turn.
294    pub const TURN_ID: &str = "turn_id";
295    /// The conversation.
296    pub const CONVERSATION_ID: &str = "conversation_id";
297    /// Digest of the tenant.
298    pub const ACCOUNT_HASH: &str = "account_hash";
299    /// Workflow key.
300    pub const WORKFLOW: &str = "workflow";
301    /// Workflow version.
302    pub const WORKFLOW_VERSION: &str = "workflow_version";
303    /// Case identifier.
304    pub const CASE_ID: &str = "case_id";
305    /// Case revision.
306    pub const CASE_REVISION: &str = "case_revision";
307    /// Interaction identifier.
308    pub const INTERACTION_ID: &str = "interaction_id";
309    /// Provider key.
310    pub const PROVIDER: &str = "provider";
311    /// Model key.
312    pub const MODEL: &str = "model";
313    /// Provider attempt.
314    pub const ATTEMPT: &str = "attempt";
315    /// Hash of the accepted plan.
316    pub const PLAN_HASH: &str = "plan_hash";
317    /// Command identifier.
318    pub const COMMAND_ID: &str = "command_id";
319    /// Committed event identifiers.
320    pub const EVENT_IDS: &str = "event_ids";
321    /// Response block identifiers.
322    pub const BLOCK_IDS: &str = "block_ids";
323    /// Outbox row identifier.
324    pub const OUTBOX_ID: &str = "outbox_id";
325    /// The signal a telemetry event reports.
326    pub const SIGNAL: &str = "signal";
327    /// The pipeline stage a span covers.
328    pub const STAGE: &str = "stage";
329    /// Measured duration of a latency signal, in milliseconds.
330    pub const DURATION_MS: &str = "duration_ms";
331    /// Risk class of a command.
332    pub const RISK: &str = "risk";
333    /// Kind of an interaction.
334    pub const INTERACTION: &str = "interaction";
335    /// Normalized request purpose.
336    pub const PURPOSE: &str = "purpose";
337    /// The effort of the turn.
338    pub const EFFORT: &str = "effort";
339    /// Stable failure or rejection code.
340    pub const ERROR_CODE: &str = "error_code";
341}
342
343/// Opens the span that covers one whole turn.
344///
345/// The account id is hashed rather than logged, so a tenant is correlatable but
346/// not identifiable from the logs alone.
347///
348/// ```rust
349/// use turnframe_core::ids::{AccountId, ConversationId, TurnId};
350/// use turnframe_telemetry::tracing::turn_span;
351///
352/// let span = turn_span(TurnId::nil(), ConversationId::nil(), &AccountId::from("acct-1"));
353/// let _entered = span.enter();
354/// ```
355#[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/// Opens a child span for one stage of the pipeline.
377///
378/// ```rust
379/// use turnframe_core::ids::TurnId;
380/// use turnframe_telemetry::tracing::{PipelineStage, stage_span};
381///
382/// let span = stage_span(PipelineStage::Understanding, TurnId::nil());
383/// let _entered = span.enter();
384/// ```
385#[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/// One provider call, described with the OpenTelemetry GenAI semantic
402/// conventions (`gen_ai.*`).
403///
404/// These are the attributes every LLM observability backend already reads —
405/// Langfuse, Datadog LLM Observability, Phoenix, Braintrust — so a Turnframe
406/// application shows up in them as a model call without any vendor adapter.
407/// The keys live in [`crate::attrs`] so a provider crate and an application
408/// spell them identically.
409///
410/// Token accounting is the part that is easy to get wrong, so it is explicit
411/// here: [`ProviderCall::input_tokens`] is always recorded **net of cached
412/// tokens**. A consumer that reads only that field sees the uncached prompt
413/// cost; one that adds [`ProviderCall::input_cached_tokens`] sees the whole
414/// prompt. Neither double counts. Use
415/// [`ProviderCall::with_cached_usage`] when the provider reports a total and a
416/// cached figure, and [`ProviderCall::with_usage`] when it already reports the
417/// net one.
418///
419/// No member holds prompt or completion text: content is user data and travels
420/// only through [`ContentRecorder`], which is off by default.
421#[derive(Debug, Clone, Default, PartialEq)]
422#[non_exhaustive]
423pub struct ProviderCall {
424    /// Provider key, recorded as `gen_ai.system`.
425    pub system: Option<String>,
426    /// Normalized request purpose, recorded as `gen_ai.operation.name`.
427    pub operation: Option<String>,
428    /// Model asked for, recorded as `gen_ai.request.model`.
429    pub request_model: Option<String>,
430    /// Sampling temperature, recorded as `gen_ai.request.temperature`.
431    pub temperature: Option<f64>,
432    /// The provider's identifier for the response.
433    pub response_id: Option<String>,
434    /// Model that actually answered, which is not always the one asked for.
435    pub response_model: Option<String>,
436    /// Why generation stopped.
437    pub finish_reasons: Vec<String>,
438    /// Input tokens billed, **net of `input_cached_tokens`**.
439    pub input_tokens: Option<u64>,
440    /// Output tokens generated.
441    pub output_tokens: Option<u64>,
442    /// Input tokens served from the provider's prompt cache.
443    pub input_cached_tokens: Option<u64>,
444}
445
446impl ProviderCall {
447    /// A call to `system` for `operation` with `model`.
448    #[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    /// Everything a persisted provider attempt already knows (spec §20.7).
463    #[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    /// Sets the requested sampling temperature.
480    #[must_use]
481    pub fn with_temperature(mut self, temperature: f64) -> Self {
482        self.temperature = Some(temperature);
483        self
484    }
485
486    /// Sets the response identifier and the model that answered.
487    #[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    /// Adds a finish reason.
499    #[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    /// Records usage the provider already reports net of its prompt cache.
506    #[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    /// Records usage a provider reports as a **total** prompt size plus a
514    /// cached figure, and stores the input tokens net of the cache.
515    ///
516    /// ```rust
517    /// use turnframe_telemetry::tracing::ProviderCall;
518    ///
519    /// let call = ProviderCall::new("openai", "extract", "gpt-x")
520    ///     .with_cached_usage(1_000, 800, 120);
521    ///
522    /// // 200 uncached + 800 cached = the 1_000 the provider reported.
523    /// assert_eq!(call.input_tokens, Some(200));
524    /// assert_eq!(call.input_cached_tokens, Some(800));
525    /// assert_eq!(call.output_tokens, Some(120));
526    /// ```
527    #[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    /// The total prompt size the provider saw: net input plus cached.
541    #[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    /// The call as `(attribute key, value)` pairs in a stable order, omitting
554    /// what is unknown.
555    ///
556    /// This is the exact set the span carries; it is a plain function so a test
557    /// can assert on it without standing up a subscriber, and an adopter can
558    /// rename the keys for a backend that wants its own spelling.
559    #[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/// Opens the span of one provider call, carrying the GenAI semantic
601/// conventions of [`crate::attrs`].
602///
603/// The span also declares the trace-grouping fields and the two content fields
604/// empty, so [`TraceGrouping::stamp`] and [`ContentRecorder`] can fill them in
605/// on the same span.
606///
607/// ```rust
608/// use turnframe_telemetry::tracing::{ProviderCall, provider_call_span};
609///
610/// let call = ProviderCall::new("openai", "extract", "gpt-x")
611///     .with_cached_usage(1_000, 800, 120);
612/// let span = provider_call_span(&call);
613/// let _entered = span.enter();
614/// ```
615#[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/// Vendor-neutral grouping of spans: which session, which end user, which
643/// tags, which environment, which release.
644///
645/// LLM observability backends filter at the level of the individual span, not
646/// only at the trace root, so the grouping has to reach every span rather than
647/// sit on the first one. Two ways to make that happen:
648///
649/// * Without the `otel` feature, [`TraceGrouping::stamp`] fills the grouping
650///   fields — which every span this crate opens declares empty — on whichever
651///   span you hand it, and [`TraceGrouping::scope_span`] opens a parent span
652///   that carries them for a subscriber that flattens ancestors.
653/// * With the `otel` feature, [`crate::otel::attach_grouping`] puts the same
654///   values in OpenTelemetry baggage on the current context, where a
655///   baggage-copying span processor in the application's SDK setup stamps them
656///   onto every span that starts underneath. The processor itself lives in the
657///   adopter's code because it needs `opentelemetry_sdk`, which this crate does
658///   not depend on; [`crate::otel::grouping_from_baggage`] gives it the
659///   key/values to copy.
660///
661/// The end-user reference is always a digest — build it with
662/// [`TraceGrouping::with_account`], which hashes the account id, or supply your
663/// own digest with [`TraceGrouping::with_end_user_hash`]. There is no
664/// constructor that takes a raw account or user identifier.
665///
666/// A backend that insists on its own attribute names is a renaming function
667/// over [`TraceGrouping::fields`] in the adopter's code, not something this
668/// crate hardcodes.
669#[derive(Debug, Clone, Default, PartialEq, Eq)]
670#[non_exhaustive]
671pub struct TraceGrouping {
672    /// The session, recorded as `session.id`. Turnframe uses the conversation.
673    pub session_id: Option<String>,
674    /// Digest of the end user, recorded as `user.id`. Never a raw identifier.
675    pub end_user_hash: Option<String>,
676    /// Free-form grouping labels the application chose.
677    pub tags: Vec<String>,
678    /// Deployment environment, e.g. `production`.
679    pub environment: Option<String>,
680    /// Release or build of the running service.
681    pub release: Option<String>,
682}
683
684impl TraceGrouping {
685    /// Groups spans by conversation, which is the session a backend shows.
686    #[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    /// Adds the end user as the digest of an account id. The raw id is hashed
695    /// by [`account_hash`] and never stored on the grouping.
696    #[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    /// Adds an end-user reference the caller has already hashed.
703    #[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    /// Adds a grouping label.
710    #[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    /// Sets the deployment environment.
717    #[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    /// Sets the release or build.
724    #[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    /// The grouping as `(attribute key, value)` pairs in a stable order,
731    /// omitting what is unset. Testable without a subscriber, and the set an
732    /// adopter renames for a backend with its own spelling.
733    #[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    /// Records the grouping on `span`.
757    ///
758    /// Every span this crate opens declares the grouping fields empty, so this
759    /// fills them in. Recording on a span that did not declare them is a no-op
760    /// rather than an error.
761    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    /// Records the grouping on the span that is currently entered.
768    pub fn stamp_current(&self) {
769        self.stamp(&::tracing::Span::current());
770    }
771
772    /// Opens a span that carries the grouping and becomes the parent of every
773    /// span opened while it is entered.
774    ///
775    /// ```rust
776    /// use turnframe_core::ids::{AccountId, ConversationId};
777    /// use turnframe_telemetry::tracing::TraceGrouping;
778    ///
779    /// let grouping = TraceGrouping::for_conversation(ConversationId::nil())
780    ///     .with_account(&AccountId::from("acct-1"))
781    ///     .with_environment("production")
782    ///     .with_release("v0.1.0")
783    ///     .with_tag("trip");
784    ///
785    /// let span = grouping.scope_span();
786    /// let _entered = span.enter();
787    /// ```
788    #[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/// Which half of a conversation a piece of content is.
806#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, PartialOrd, Ord)]
807#[non_exhaustive]
808pub enum ContentRole {
809    /// The prompt sent to the model.
810    Input,
811    /// The completion the model returned.
812    Output,
813}
814
815impl ContentRole {
816    /// The attribute key this role is recorded under.
817    #[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
826/// Decides what, if anything, of a prompt or completion may be recorded.
827///
828/// The hook is mandatory: content recording cannot be switched on without one,
829/// so there is no configuration in which raw user text reaches the backend
830/// unexamined. Returning `None` drops the content entirely.
831pub trait ContentRedactor: Send + Sync + fmt::Debug {
832    /// Returns the text that may be recorded, or `None` to record nothing.
833    fn redact(&self, role: ContentRole, text: &str) -> Option<String>;
834}
835
836/// A redactor that records nothing. Useful as an explicit "not yet" while the
837/// real one is being written.
838#[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/// Records prompt and completion text on a provider-call span.
848///
849/// **Disabled by default, and for a reason.** Prompts and completions are user
850/// data: enabling this sends what the user typed, and what the model said back,
851/// to whatever tracing backend is configured. That is a decision about data
852/// residency, retention and consent, not a debugging convenience, so it takes
853/// an explicit [`ContentRecorder::enabled`] call and a
854/// [`ContentRedactor`] to switch on (spec §25.5).
855///
856/// ```rust
857/// use std::sync::Arc;
858///
859/// use turnframe_telemetry::tracing::{ContentRecorder, ContentRole, DropAllContent};
860///
861/// // The default records nothing at all.
862/// let off = ContentRecorder::disabled();
863/// assert!(!off.is_enabled());
864/// assert_eq!(off.rendered(ContentRole::Input, "withdraw trip 17"), None);
865///
866/// // Even switched on, everything goes through the redaction hook.
867/// let on = ContentRecorder::enabled(Arc::new(DropAllContent));
868/// assert!(on.is_enabled());
869/// assert_eq!(on.rendered(ContentRole::Input, "withdraw trip 17"), None);
870/// ```
871#[derive(Clone)]
872pub struct ContentRecorder {
873    enabled: bool,
874    redactor: Arc<dyn ContentRedactor>,
875}
876
877impl ContentRecorder {
878    /// Records nothing. The default, and what an application gets unless it
879    /// deliberately asks for something else.
880    #[must_use]
881    pub fn disabled() -> Self {
882        Self {
883            enabled: false,
884            redactor: Arc::new(DropAllContent),
885        }
886    }
887
888    /// Records content, every piece of it through `redactor` first.
889    #[must_use]
890    pub fn enabled(redactor: Arc<dyn ContentRedactor>) -> Self {
891        Self {
892            enabled: true,
893            redactor,
894        }
895    }
896
897    /// Whether content recording is switched on.
898    #[must_use]
899    pub fn is_enabled(&self) -> bool {
900        self.enabled
901    }
902
903    /// What would be recorded for this text: `None` when recording is off or
904    /// the redactor dropped it. Exposed so an application can test its own
905    /// redaction hook without a subscriber.
906    #[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    /// Records content on `span`, if recording is on and the redactor allows
915    /// it. The span must be one this crate opened, which declares the two
916    /// content fields empty.
917    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
939/// Emits an event carrying every known identifier of a turn (spec §26.1).
940///
941/// Use it at the end of a turn, or when a replay record is written, so one log
942/// line ties the turn to its plan, its commands, its events and its blocks.
943pub 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/// The label fields of a signal, as `(field name, value)` pairs.
968///
969/// Only the typed members of [`SignalLabels`] appear; there is no free-text
970/// member to leak. Exposed as a plain function so the extraction can be tested
971/// without a subscriber.
972#[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/// An [`Observer`] that emits one structured `tracing` event per signal.
1010///
1011/// Signals that report a safety-integrity problem
1012/// ([`Signal::is_safety_signal`]) are emitted at `WARN`, everything else at
1013/// `DEBUG`: a claim violation or a stale interaction should surface without a
1014/// filter change, while the ordinary flow of a busy system should not.
1015///
1016/// The event carries the signal name, the typed labels and, for a latency
1017/// signal, the measured duration in milliseconds. It never carries text.
1018///
1019/// ```rust
1020/// use turnframe_core::ids::WorkflowKey;
1021/// use turnframe_core::observe::{Observer, Signal, SignalLabels};
1022/// use turnframe_telemetry::TracingObserver;
1023///
1024/// let observer = TracingObserver::new();
1025/// observer.observe_labeled(
1026///     &Signal::WorkflowInvariantViolation,
1027///     &SignalLabels::workflow(WorkflowKey::from("trip")).with_error_code("two_open_drafts"),
1028/// );
1029/// ```
1030#[derive(Debug, Clone, Copy, Default)]
1031pub struct TracingObserver;
1032
1033impl TracingObserver {
1034    /// Builds the observer. Events go to whichever subscriber the application
1035    /// installed, so there is nothing to configure here.
1036    #[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        // The digest is not the plain digest of the identifier.
1143        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        // Net plus cached is the total the provider reported: no double count.
1306        assert_eq!(call.total_input_tokens(), Some(1_000));
1307
1308        // A cache figure larger than the total cannot underflow.
1309        let odd = ProviderCall::default().with_cached_usage(10, 40, 1);
1310        assert_eq!(odd.input_tokens, Some(0));
1311
1312        // Usage reported already net leaves the cache unset.
1313        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        // The persisted attempt now carries what the request asked for and what
1355        // the provider answered, so the span needs no second source.
1356        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        // The raw tenant identifier is nowhere in the grouping.
1415        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    /// A redactor that keeps only the length, to prove the hook is consulted.
1467    #[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        // The redactor dropped the completion entirely.
1490        assert_eq!(recorder.rendered(ContentRole::Output, "done"), None);
1491
1492        // And a redactor that drops everything records nothing even when on.
1493        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}