Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
47 changes: 42 additions & 5 deletions docs/Operations/Observability.md
Original file line number Diff line number Diff line change
Expand Up @@ -34,9 +34,29 @@ turn (engine loop + message persistence). It carries:
| `gen_ai.conversation.id` | the `conversation_id` arg | Ties the turn to its conversation. |
| `gen_ai.agent.name` | constant `"smooth-agent-chat"` | The agent/persona driving the turn. |
| `smooai.org_id` | the turn's `org_id` (streaming path) | Set only when an org is resolved. **Matches the monorepo TS chat handler's attribute exactly**, so the observability studio groups Rust + TS turns by org. |
| `gen_ai.usage.input_tokens` | `AgentEvent::Completed.prompt_tokens` | Recorded on completion **only when the engine reported usage** (non-zero). Omitted otherwise — e.g. a mock turn — per the convention's "omit if unknown" rule. |
| `gen_ai.usage.output_tokens` | `AgentEvent::Completed.completion_tokens` | Same gating as input tokens. |
| `gen_ai.usage.input_tokens` | `AgentEvent::Completed.prompt_tokens` | Recorded **only when `prompt_tokens > 0`** — see "Zero prompt tokens means fabricated usage" below. |
| `gen_ai.usage.output_tokens` | `AgentEvent::Completed.completion_tokens` | Same gating as input tokens: both or neither. |
| `gen_ai.usage.cost_usd` | `AgentEvent::Completed.cost_usd` | The turn's cost in USD. **Recorded only when positive** — see "A zero cost is never recorded" below. |
| `smooai.gen_ai.cost_unavailable` | constant `"unpriced"` | Set **instead of** `cost_usd` when no cost could be established. Same attribute name and values as the TypeScript lane, so a consumer never special-cases per engine. |

#### Zero prompt tokens means fabricated usage

`prompt_tokens = 0` does not mean the turn consumed no input — a grounded turn
always does. It is the signature of core's streaming collector **inventing** a
usage struct, which it does whenever no `StreamEvent::Usage` arrives: it
hardcodes `prompt_tokens = 0` and estimates `completion_tokens` as
`content.len() / 4`. LiteLLM at `llm.smoo.ai` drops the usage chunk for
`smooth-*` aliases, so this is the common case, not the edge case.

Both counts are therefore omitted together. The plausible-looking output number
beside a zero input is an estimate, not a measurement, and shipping it next to a
dollar figure makes a guess look authoritative. Absent is honest; `0` is a lie.

The real fix is in `smooth-operator-core` — either stop fabricating, or mark the
struct as estimated. Until then this instrumentation declines to export it.

`telemetry::record_turn_usage` is the single place this policy lives, so the
streaming runner and `KnowledgeChatRuntime` cannot drift apart on it.

#### A zero cost is never recorded

Expand All @@ -54,10 +74,27 @@ fallback prices any model it doesn't recognise at `0`. So a zero always means
render a paid turn as a confident `$0.00`. An absent attribute lets a consumer
say "not measured" instead.

A *positive* value may be either the gateway's authoritative number or the
local `ModelPricing` estimate; the engine sums both onto one `f64`, so the
Cost is judged **independently of the token counts**, on purpose: the gateway
reports cost in an HTTP header and usage in an SSE chunk — two separate channels
— so a turn can legitimately have an authoritative cost and no usage.
Suppressing cost whenever usage was fabricated would throw that away and
recreate the all-zero-rows bug.

A *positive* value may be either the gateway's authoritative number or the local
`ModelPricing` estimate; the engine sums both onto one `f64`, so the
instrumentation cannot tell them apart. Separating them needs a second field on
`AgentEvent::Completed`.
`AgentEvent::Completed`. The residual hazard is a locally-priced model with
fabricated usage: `ModelPricing × 0 input` is an undercount that still looks
authoritative. Production is not exposed today, because its model is not in the
local pricing table — the local path yields exactly `0`, which is dropped.

Independently, `gen_ai_events.response_id` joins to
`LiteLLM_SpendLogs.request_id`, which carries the gateway's own dollars **and**
its real prompt/completion counts. The Rust operator does **not** populate
`gen_ai.response.id` yet: core never captures the `chatcmpl-…` id — neither
`ChatResponse` nor `StreamChunk` deserializes an `id` field, and `LlmResponse`
has nowhere to put one. Adding it is a core change, and it would make
"measured vs estimated" verifiable after the fact rather than a matter of trust.

### `gen_ai.tool` span — one per tool call

Expand Down
22 changes: 9 additions & 13 deletions rust/smooth-operator-server/src/runner.rs
Original file line number Diff line number Diff line change
Expand Up @@ -48,11 +48,11 @@ use smooth_operator::domain::{Citation, Direction, Message as DomainMessage, Mes
use smooth_operator::interaction::{InteractionOutcome, InteractionRegistry, InteractionRequest};
use smooth_operator::rerank::Reranker;
use smooth_operator::telemetry::{
record_cost_usd, redact_tool_arguments, AGENT_NAME, GEN_AI_AGENT_NAME, GEN_AI_CONVERSATION_ID,
GEN_AI_OPERATION_NAME, GEN_AI_REQUEST_MODEL, GEN_AI_SYSTEM, GEN_AI_TOOL_ARGUMENTS,
GEN_AI_TOOL_NAME, GEN_AI_USAGE_COST_USD, GEN_AI_USAGE_INPUT_TOKENS, GEN_AI_USAGE_OUTPUT_TOKENS,
OPERATION_CHAT, OPERATION_TOOL, OTEL_STATUS_CODE, OTEL_STATUS_MESSAGE, SMOOAI_ORG_ID,
SPAN_CHAT, SPAN_TOOL, SYSTEM_NAME,
record_turn_usage, redact_tool_arguments, AGENT_NAME, COST_UNAVAILABLE, GEN_AI_AGENT_NAME,
GEN_AI_CONVERSATION_ID, GEN_AI_OPERATION_NAME, GEN_AI_REQUEST_MODEL, GEN_AI_SYSTEM,
GEN_AI_TOOL_ARGUMENTS, GEN_AI_TOOL_NAME, GEN_AI_USAGE_COST_USD, GEN_AI_USAGE_INPUT_TOKENS,
GEN_AI_USAGE_OUTPUT_TOKENS, OPERATION_CHAT, OPERATION_TOOL, OTEL_STATUS_CODE,
OTEL_STATUS_MESSAGE, SMOOAI_ORG_ID, SPAN_CHAT, SPAN_TOOL, SYSTEM_NAME,
};
use smooth_operator::tool_provider::{ToolProvider, ToolProviderContext};
use smooth_operator::tools::{
Expand Down Expand Up @@ -1023,6 +1023,7 @@ pub async fn run_streaming_turn(
{ GEN_AI_USAGE_INPUT_TOKENS } = tracing::field::Empty,
{ GEN_AI_USAGE_OUTPUT_TOKENS } = tracing::field::Empty,
{ GEN_AI_USAGE_COST_USD } = tracing::field::Empty,
{ COST_UNAVAILABLE } = tracing::field::Empty,
);
if let Some(org) = org_id_for_span.as_deref() {
turn_span.record(SMOOAI_ORG_ID, org);
Expand Down Expand Up @@ -1285,14 +1286,9 @@ pub async fn run_streaming_turn(
// the GenAI conventions); one `gen_ai.tool` child span per tool call with the
// redacted arguments, latency, and an ERROR status on failure.
if let Some(u) = usage.as_ref() {
if u.prompt_tokens > 0 || u.completion_tokens > 0 {
turn_span.record(GEN_AI_USAGE_INPUT_TOKENS, u.prompt_tokens);
turn_span.record(GEN_AI_USAGE_OUTPUT_TOKENS, u.completion_tokens);
}
// Gateway-authoritative cost (LiteLLM's `x-litellm-response-cost`), which
// the engine already accumulates onto `Completed`. Dropped when zero —
// that means "unpriced", not "free". See `record_cost_usd`.
record_cost_usd(&turn_span, u.cost_usd);
// One helper owns the "measured or absent" policy for both turn paths —
// which counts are real, and whether a cost may be trusted beside them.
record_turn_usage(&turn_span, u.prompt_tokens, u.completion_tokens, u.cost_usd);
}
for rec in &tool_records {
// The OTLP ingest merges resource attrs with THIS span's attrs and does
Expand Down
92 changes: 92 additions & 0 deletions rust/smooth-operator-server/tests/telemetry.rs
Original file line number Diff line number Diff line change
Expand Up @@ -393,4 +393,96 @@ async fn unpriced_turn_omits_cost_rather_than_recording_zero() {
missing price must never read as free; fields: {:?}",
chat.fields
);
// …and says WHY, with the same attribute name + value the TS lane emits, so
// a consumer never special-cases per engine. "unpriced" is actionable:
// someone has to price the model.
assert_eq!(
chat.fields
.get("smooai.gen_ai.cost_unavailable")
.map(String::as_str),
Some("unpriced"),
"absence alone makes a consumer infer the reason; fields: {:?}",
chat.fields
);
}

/// The prod signature, four for four across every chat-ws turn ever recorded:
/// `input_tokens = 0` with a plausible-looking output count.
///
/// Its cause is not this crate — core's `collect_stream` fabricates the whole
/// usage struct when the gateway sends no usage chunk (LiteLLM drops it for
/// `smooth-*` aliases), hardcoding `prompt_tokens = 0` and estimating
/// `completion_tokens` as `content.len() / 4`. So the "plausible" output count
/// was never a measurement either. Until core stops fabricating, the honest
/// move is to export neither. Pearl th-126fe6.
#[tokio::test]
async fn fabricated_usage_omits_both_token_counts() {
let sink: SpanSink = Arc::new(Mutex::new(Vec::new()));
let layer = CapturingLayer {
sink: Arc::clone(&sink),
index: Arc::new(Mutex::new(HashMap::new())),
};
let subscriber = tracing_subscriber::registry().with(layer);
let _guard = tracing::subscriber::set_default(subscriber);

// No `StreamEvent::Usage` at all — exactly what the gateway sends today.
// Core will fabricate: prompt 0, completion ≈ len/4.
let mock = MockLlmClient::new();
mock.push_stream(vec![
StreamEvent::Delta {
content: "Items are accepted within 30 days for a full refund.".into(),
},
StreamEvent::Done {
finish_reason: "stop".into(),
},
]);

let (tx, mut rx) = unbounded_channel::<serde_json::Value>();
runner::run_streaming_turn(
TurnRequest {
llm_provider: Some(Arc::new(mock.clone())),
conversation_id: "conv-otel-fabricated",
request_id: "req-otel-fabricated",
..base_turn_request()
},
&tx,
)
.await
.expect("run_streaming_turn");
drop(tx);
while rx.try_recv().is_ok() {}

let spans = sink.lock().expect("sink poisoned").clone();
let chat = spans
.iter()
.find(|s| s.name == "gen_ai.chat")
.unwrap_or_else(|| panic!("expected a `gen_ai.chat` span; got: {spans:#?}"));

assert!(
!chat.fields.contains_key("gen_ai.usage.input_tokens"),
"input_tokens = 0 is impossible on a grounded turn — it must be ABSENT, \
not 0; fields: {:?}",
chat.fields
);
assert!(
!chat.fields.contains_key("gen_ai.usage.output_tokens"),
"the output count beside a fabricated input is core's content.len()/4 \
estimate, not a measurement — shipping it next to a dollar figure \
would look authoritative; fields: {:?}",
chat.fields
);

// Pinning the KNOWN-IMPERFECT half so it can't change unnoticed. Cost is
// judged independently of the counts (gateway sends them over separate
// channels), and `openai/gpt-4o` IS in the local `ModelPricing` table — so
// this turn gets a locally-derived cost computed against a zero input
// count, i.e. an undercount. Only cost provenance on the engine's
// `Completed` event can distinguish that from a gateway figure; prod is not
// exposed because its model isn't in the local table (local path → 0 →
// dropped). If this assertion ever starts failing, provenance landed.
assert!(
chat.fields.contains_key("gen_ai.usage.cost_usd"),
"documents today's residual gap, not desired behaviour; fields: {:?}",
chat.fields
);
}
57 changes: 20 additions & 37 deletions rust/smooth-operator/src/runtime.rs
Original file line number Diff line number Diff line change
Expand Up @@ -22,11 +22,11 @@ use crate::adapter::{MessageQuery, StorageAdapter};
use crate::curation::{CuratedKnowledgeStore, RetrievalFilter};
use crate::domain::{Citation, Direction, Message as DomainMessage, MessageContent};
use crate::telemetry::{
record_cost_usd, redact_tool_arguments, AGENT_NAME, GEN_AI_AGENT_NAME, GEN_AI_CONVERSATION_ID,
GEN_AI_OPERATION_NAME, GEN_AI_REQUEST_MODEL, GEN_AI_SYSTEM, GEN_AI_TOOL_ARGUMENTS,
GEN_AI_TOOL_NAME, GEN_AI_USAGE_COST_USD, GEN_AI_USAGE_INPUT_TOKENS, GEN_AI_USAGE_OUTPUT_TOKENS,
OPERATION_CHAT, OPERATION_TOOL, OTEL_STATUS_CODE, OTEL_STATUS_MESSAGE, SPAN_CHAT, SPAN_TOOL,
SYSTEM_NAME,
record_turn_usage, redact_tool_arguments, AGENT_NAME, COST_UNAVAILABLE, GEN_AI_AGENT_NAME,
GEN_AI_CONVERSATION_ID, GEN_AI_OPERATION_NAME, GEN_AI_REQUEST_MODEL, GEN_AI_SYSTEM,
GEN_AI_TOOL_ARGUMENTS, GEN_AI_TOOL_NAME, GEN_AI_USAGE_COST_USD, GEN_AI_USAGE_INPUT_TOKENS,
GEN_AI_USAGE_OUTPUT_TOKENS, OPERATION_CHAT, OPERATION_TOOL, OTEL_STATUS_CODE,
OTEL_STATUS_MESSAGE, SPAN_CHAT, SPAN_TOOL, SYSTEM_NAME,
};
use crate::tools::{KnowledgeResultSink, KnowledgeSearchTool};
use tracing::Instrument;
Expand Down Expand Up @@ -195,35 +195,22 @@ pub struct TurnOutcome {
pub citations: Vec<Citation>,
}

/// Extract `(input_tokens, output_tokens)` from the engine's terminal
/// [`AgentEvent::Completed`] event, if one is present and carries usage. The
/// engine reports `prompt_tokens` / `completion_tokens` on `Completed`; those
/// map directly onto the GenAI `gen_ai.usage.input_tokens` /
/// `gen_ai.usage.output_tokens` attributes. Returns `None` when there is no
/// `Completed` event (e.g. a mock turn that didn't surface usage), so the
/// caller omits the attributes rather than recording zeros.
fn usage_from_events(events: &[AgentEvent]) -> Option<(u64, u64)> {
/// The turn's token counts and cost from the terminal [`AgentEvent::Completed`],
/// or `None` when the turn produced no `Completed` event (e.g. an offline mock
/// turn) and therefore reported nothing at all.
///
/// Deliberately unfiltered: it hands back whatever the engine said, including a
/// fabricated `prompt_tokens = 0`. Deciding which of those numbers may be
/// exported is [`record_turn_usage`](crate::telemetry::record_turn_usage)'s
/// job, so that policy lives in exactly one place for both turn paths.
fn usage_from_events(events: &[AgentEvent]) -> Option<(u64, u64, f64)> {
events.iter().find_map(|e| match e {
AgentEvent::Completed {
prompt_tokens,
completion_tokens,
cost_usd,
..
} if *prompt_tokens > 0 || *completion_tokens > 0 => {
Some((*prompt_tokens, *completion_tokens))
}
_ => None,
})
}

/// The turn's accumulated cost in USD from the terminal
/// [`AgentEvent::Completed`], or `None` when the turn produced no `Completed`
/// event. The value still needs the zero check in
/// [`record_cost_usd`](crate::telemetry::record_cost_usd) before it can be
/// exported — the engine reports `0.0` both for a free turn and for one it
/// could not price.
fn cost_from_events(events: &[AgentEvent]) -> Option<f64> {
events.iter().find_map(|e| match e {
AgentEvent::Completed { cost_usd, .. } => Some(*cost_usd),
} => Some((*prompt_tokens, *completion_tokens, *cost_usd)),
_ => None,
})
}
Expand Down Expand Up @@ -557,6 +544,7 @@ impl KnowledgeChatRuntime {
{ GEN_AI_USAGE_INPUT_TOKENS } = tracing::field::Empty,
{ GEN_AI_USAGE_OUTPUT_TOKENS } = tracing::field::Empty,
{ GEN_AI_USAGE_COST_USD } = tracing::field::Empty,
{ COST_UNAVAILABLE } = tracing::field::Empty,
);

// Run the turn body inside the span so any engine-internal spans nest
Expand All @@ -568,14 +556,9 @@ impl KnowledgeChatRuntime {

// Record token usage on the turn span if the engine reported it via the
// terminal `Completed` event (omitted otherwise, per the GenAI convs).
if let Some((input, output)) = usage_from_events(&outcome.events) {
turn_span.record(GEN_AI_USAGE_INPUT_TOKENS, input);
turn_span.record(GEN_AI_USAGE_OUTPUT_TOKENS, output);
}
// Gateway-authoritative per-turn cost, dropped when zero (= unpriced,
// not free). See `record_cost_usd`.
if let Some(cost) = cost_from_events(&outcome.events) {
record_cost_usd(&turn_span, cost);
// One helper owns the "measured or absent" policy for both turn paths.
if let Some((input, output, cost)) = usage_from_events(&outcome.events) {
record_turn_usage(&turn_span, input, output, cost);
}

// Emit a child `gen_ai.tool` span per tool call so each invocation is an
Expand Down
Loading
Loading