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
202 changes: 195 additions & 7 deletions apps/ingest/src/ai_session.rs
Original file line number Diff line number Diff line change
Expand Up @@ -246,7 +246,8 @@ fn resource_facts(attrs: &[KeyValue]) -> ResourceFacts<'_> {

/// Scope-level facts, computed once per `ScopeSpans`. `any` is true when the
/// scope alone can decide a vendor; the flags that only narrow span evidence
/// (`spring_boot`, `vercel_ai`, `matches_service_name`, the crewai refusal)
/// (`genkit`, `spring_boot`, `vercel_ai`, `matches_service_name`, the crewai
/// refusal)
/// deliberately don't set it, so they never force predicate evaluation on
/// evidence-free spans.
#[expect(
Expand All @@ -260,8 +261,10 @@ struct ScopeFacts {
dspy: bool,
eve: bool,
flue: bool,
genkit: bool,
google_adk: bool,
haystack: bool,
haystack_openinference: bool,
langchain: bool,
litellm: bool,
llamaindex: bool,
Expand Down Expand Up @@ -292,6 +295,7 @@ const SCOPE_NAMES: &[&str] = &[
"openinference.instrumentation.",
"eve",
"@flue/opentelemetry",
"genkit-tracer",
"gcp.vertex.agent",
"haystack",
"langsmith",
Expand All @@ -301,6 +305,7 @@ const SCOPE_NAMES: &[&str] = &[
"agent_framework",
"Experimental.Microsoft.Agents.AI",
"@arizeai/openinference-instrumentation-openai-agents",
"@arizeai/openinference-instrumentation-langchain",
"agent_runtime ",
"openrouter",
"crewai.telemetry",
Expand All @@ -310,6 +315,7 @@ const SCOPE_NAMES: &[&str] = &[
"ai",
"gen_ai",
"semantic_kernel.",
"Microsoft.SemanticKernel.Diagnostics",
];

static SCOPE_SCREEN: [u64; 256] = build_screen(&[SCOPE_NAMES]);
Expand All @@ -334,9 +340,16 @@ fn scope_facts(scope_name: &str, resource: &ResourceFacts) -> ScopeFacts {
"openinference.instrumentation.dspy" => facts.dspy = true,
"eve" => facts.eve = true,
"@flue/opentelemetry" => facts.flue = true,
"gcp.vertex.agent" => facts.google_adk = true,
"genkit-tracer" => facts.genkit = true,
"gcp.vertex.agent" | "openinference.instrumentation.google_adk" => {
facts.google_adk = true;
}
"haystack" => facts.haystack = true,
"langsmith" | "openinference.instrumentation.langchain" => facts.langchain = true,
"openinference.instrumentation.haystack" => facts.haystack_openinference = true,
// The `@arizeai/` name is the TypeScript instrumentor.
"langsmith"
| "openinference.instrumentation.langchain"
| "@arizeai/openinference-instrumentation-langchain" => facts.langchain = true,
"litellm" => facts.litellm = true,
"llamaindex.opentelemetry.tracer" | "openinference.instrumentation.llama_index" => {
facts.llamaindex = true;
Expand All @@ -354,7 +367,9 @@ fn scope_facts(scope_name: &str, resource: &ResourceFacts) -> ScopeFacts {
"openrouter" => facts.openrouter = true,
"openinference.instrumentation.crewai" | "crewai.telemetry" => facts.crewai = true,
"openinference.instrumentation.smolagents" => facts.smolagents = true,
"pydantic-ai" => facts.pydantic = true,
"pydantic-ai" | "openinference.instrumentation.pydantic_ai" => facts.pydantic = true,
// The .NET build's ActivitySource; Python's are `semantic_kernel.*`.
"Microsoft.SemanticKernel.Diagnostics" => facts.semantic_kernel = true,
"strands.telemetry.tracer" => facts.strands = true,
"org.springframework.boot" => facts.spring_boot = true,
"ai" | "gen_ai" => facts.vercel_ai = true,
Expand All @@ -377,6 +392,7 @@ fn scope_facts(scope_name: &str, resource: &ResourceFacts) -> ScopeFacts {
|| facts.flue
|| facts.google_adk
|| facts.haystack
|| facts.haystack_openinference
|| facts.langchain
|| facts.litellm
|| facts.llamaindex
Expand Down Expand Up @@ -879,10 +895,27 @@ static VENDORS: &[Vendor] = &[
detect: detect_flue,
session_keys: &["gen_ai.conversation.id"],
},
Vendor {
id: "genkit",
detect: detect_genkit,
session_keys: CONVERSATION_ID_ONLY,
},
Vendor {
id: "google_adk",
detect: detect_google_adk,
session_keys: &["gen_ai.conversation.id", "gcp.vertex.agent.session_id"],
// `session.id` is OpenInference's.
session_keys: &[
"gen_ai.conversation.id",
"gcp.vertex.agent.session_id",
"session.id",
],
},
// One vendor, two dialects with their own session keys: `session.id` is
// the OpenInference instrumentor's, and means nothing on a native span.
Vendor {
id: "haystack",
detect: detect_haystack_openinference,
session_keys: &["session.id", "gen_ai.conversation.id"],

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

🟡 Native Haystack sessions grouped by unrelated ID

When a native Haystack span carries both session.id and gen_ai.conversation.id, it now selects session.id. That key belongs to the OpenInference instrumentor, not the native tracer, so the span joins the wrong AI session.

Learn more

The gateway selects the first nonempty key from each vendor's session list in run_predicates. Haystack's native tracer and its OpenInference instrumentor share one vendor ID, but only the latter uses session.id as its session key. On native spans, an application-supplied session.id can therefore override the previous canonical conversation ID. This ID flows into the session grouping in sessionKey.

Example: A native Haystack span has gen_ai.conversation.id=agent-42 and a generic session.id=browser-9. The gateway now stamps browser-9; previously it stamped agent-42, keeping that agent's traces together.

Recommended fix: Choose session.id for Haystack only when the scope identifies its OpenInference instrumentor, and retain gen_ai.conversation.id as the native scope's session key. Cover a span containing both keys in each scope.

Devin Review


Was this helpful? React with 👍 or 👎 to provide feedback.

Copy link
Copy Markdown
Collaborator Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Valid, fixed in b88400a. Haystack now has two entries under the same vendor id: the openinference.instrumentation.haystack scope reads session.id first, and native haystack spans read only gen_ai.conversation.id again. google_adk and pydantic_ai were already fine because session.id sits after their own keys on both scopes. New test session_key_order_follows_the_emitting_dialect covers a span with both keys on each scope for all three vendors.

},
Vendor {
id: "haystack",
Expand Down Expand Up @@ -948,7 +981,8 @@ static VENDORS: &[Vendor] = &[
Vendor {
id: "pydantic_ai",
detect: detect_pydantic_ai,
session_keys: &["gen_ai.conversation.id"],
// `session.id` is OpenInference's.
session_keys: &["gen_ai.conversation.id", "session.id"],
},
Vendor {
id: "semantic_kernel",
Expand Down Expand Up @@ -1099,10 +1133,21 @@ fn detect_flue(c: &Ctx) -> bool {
c.scope.flue || c.ev.flue
}

fn detect_genkit(c: &Ctx) -> bool {
// Genkit's raw spans carry only its own `genkit:*` dialect, which nothing
// here decodes; the spans an app's GenAI processor has restated name their
// operation.
c.scope.genkit && c.ev.has_gen_ai_operation_name
}

fn detect_google_adk(c: &Ctx) -> bool {
c.scope.google_adk || c.ev.gcp_vertex_agent || c.ev.gen_ai_system == "gcp.vertex.agent"
}

fn detect_haystack_openinference(c: &Ctx) -> bool {
c.scope.haystack_openinference
}

fn detect_haystack(c: &Ctx) -> bool {
c.scope.haystack
|| c.ev.haystack
Expand Down Expand Up @@ -1256,7 +1301,10 @@ fn detect_vercel_ai_sdk(c: &Ctx) -> bool {
}

fn detect_unknown_genai(c: &Ctx) -> bool {
c.ev.has_gen_ai_operation_name
// An OpenInference span that also dual-writes the GenAI operation (every
// provider-client instrumentor does) belongs to the bucket that decodes
// its dialect.
c.ev.has_gen_ai_operation_name && !c.ev.openinference_span_kind
}

fn detect_unknown_openinference(c: &Ctx) -> bool {
Expand Down Expand Up @@ -2047,6 +2095,69 @@ mod tests {
"openai_agents_sdk",
"oa-1",
),
(
"@arizeai/openinference-instrumentation-langchain",
"ChatOpenAI",
&[
("openinference.span.kind", "LLM"),
("gen_ai.operation.name", "chat"),
("session.id", "lc-3"),
],
"langchain",
"lc-3",
),
(
"openinference.instrumentation.haystack",
"OpenAIChatGenerator.run",
&[("openinference.span.kind", "LLM"), ("session.id", "h-1")],
"haystack",
"h-1",
),
(
"openinference.instrumentation.google_adk",
"invoke_agent weather_agent",
&[("openinference.span.kind", "AGENT"), ("session.id", "g-2")],
"google_adk",
"g-2",
),
(
"openinference.instrumentation.pydantic_ai",
"agent run",
&[("openinference.span.kind", "AGENT"), ("session.id", "p-2")],
"pydantic_ai",
"p-2",
),
(
"Microsoft.SemanticKernel.Diagnostics",
"chat.completions gpt-4o",
&[
("gen_ai.operation.name", "chat"),
("gen_ai.conversation.id", "sk-1"),
],
"semantic_kernel",
"sk-1",
),
(
"Microsoft.SemanticKernel.Diagnostics",
"invoke_agent WeatherAgent",
&[
("gen_ai.operation.name", "invoke_agent"),
("gen_ai.conversation.id", "sk-2"),
],
"semantic_kernel",
"sk-2",
),
(
"genkit-tracer",
"generate",
&[
("genkit:type", "action"),
("gen_ai.operation.name", "chat"),
("gen_ai.conversation.id", "gk-1"),
],
"genkit",
"gk-1",
),
];
for (scope_name, span_name, span_attrs, vendor, session_id) in cases {
classified(
Expand All @@ -2060,6 +2171,83 @@ mod tests {
}
}

#[test]
fn session_key_order_follows_the_emitting_dialect() {
// Every span carries both keys.
let both = &[
("gen_ai.conversation.id", "conv-1"),
("session.id", "app-1"),
];
for (scope, span_name, vendor, session_id) in [
// OpenInference's own session key leads on its instrumentor's spans.
(
"openinference.instrumentation.haystack",
"OpenAIChatGenerator.run",
"haystack",
"app-1",
),
// A native span's `session.id` is not the vendor's: the
// conversation id wins.
("haystack", "haystack.pipeline.run", "haystack", "conv-1"),
// google_adk and pydantic_ai rank `session.id` after their own keys
// on both scopes.
(
"openinference.instrumentation.google_adk",
"invoke_agent weather_agent",
"google_adk",
"conv-1",
),
("gcp.vertex.agent", "invoke_agent", "google_adk", "conv-1"),
(
"openinference.instrumentation.pydantic_ai",
"agent run",
"pydantic_ai",
"conv-1",
),
("pydantic-ai", "agent run", "pydantic_ai", "conv-1"),
] {
classified(scope, span_name, both, &[], vendor, Some(session_id));
}
}

#[test]
fn genkit_spans_without_a_gen_ai_operation_are_not_ai() {
assert_eq!(
classify(
"genkit-tracer",
"generate",
&[("genkit:type", "action"), ("genkit:name", "generate")],
&[],
),
None,
);
}

#[test]
fn openinference_provider_clients_land_in_the_openinference_bucket() {
// The provider-client instrumentors dual-write the GenAI operation, which
// used to file them under `unknown:genai`, whose read path skips the
// OpenInference decoding.
for scope_name in [
"openinference.instrumentation.anthropic",
"openinference.instrumentation.google_genai",
"@arizeai/openinference-instrumentation-anthropic",
] {
classified(
scope_name,
"Messages",
&[
("openinference.span.kind", "LLM"),
("gen_ai.operation.name", "chat"),
("gen_ai.conversation.id", "conv-4"),
],
&[],
"unknown:openinference",
Some("conv-4"),
);
}
}

#[test]
fn strands_typescript_is_detected_under_a_custom_service_name() {
// docs_strands_ts: tracer and provider are both named after the
Expand Down
4 changes: 2 additions & 2 deletions apps/ingest/src/ai_session/usage.rs
Original file line number Diff line number Diff line change
Expand Up @@ -551,7 +551,7 @@ mod tests {
}

/// (c) PROD `blind-ts-langchain-demo-001`, OpenInference LangChain JS with
/// the GenAI mirror (`unknown:genai`): (prompt, completion, reasoning).
/// the GenAI mirror (`langchain`): (prompt, completion, reasoning).
/// Three calls report more reasoning than completion; the total is still
/// the sum of `llm.token_count.total`, 17763. The list used to read output
/// 1206 / reasoning 3098 / total 17827, the page 4240 / 0 / 17763.
Expand Down Expand Up @@ -606,7 +606,7 @@ mod tests {
"@arizeai/openinference-instrumentation-langchain",
&slices(&spans),
);
assert!(stamped.iter().all(|span| span.vendor == "unknown:genai"));
assert!(stamped.iter().all(|span| span.vendor == "langchain"));
// 73 completion, 95 reasoning: all of the completion was reasoning.
assert_eq!(stamped[7].buckets, Some([732, 0, 0, 0, 73]));
let total = sum(&stamped);
Expand Down
3 changes: 3 additions & 0 deletions packages/query-engine-integrations/src/ai/ai-vendors.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -128,9 +128,12 @@ describe("openinference", () => {
"agno",
"crewai",
"dspy",
"google_adk",
"haystack",
"langchain",
"llamaindex",
"openai_agents_sdk",
"pydantic_ai",
"smolagents",
] as const) {
expect(AI_VENDOR_INTEGRATIONS[vendorId].id).toBe("openinference")
Expand Down
9 changes: 6 additions & 3 deletions packages/query-engine-integrations/src/ai/ai-vendors.ts
Original file line number Diff line number Diff line change
Expand Up @@ -90,9 +90,9 @@ const OPENINFERENCE_SPAN_KIND_OPERATIONS = new Map([
* `openinference-openai` (the gateway's id for the OpenAI instrumentor),
* `unknown:openinference` (its generic bucket for any other OpenInference
* scope) and the framework ids an `openinference.instrumentation.<framework>`
* scope is stamped with (agno, crewai, dspy, openai_agents_sdk, smolagents;
* langchain and llamaindex once the gateway fingerprints those scopes),
* because the dialect is identical; only the detection path differs. A
* scope is stamped with (agno, crewai, dspy, google_adk, haystack, langchain,
* llamaindex, openai_agents_sdk, pydantic_ai, smolagents), because the dialect
* is identical; only the detection path differs. A
* framework's native spans carry none of these keys, so the entry costs them
* nothing.
*
Expand Down Expand Up @@ -219,9 +219,12 @@ export const AI_VENDOR_INTEGRATIONS = {
agno: openInferenceIntegration,
crewai: openInferenceIntegration,
dspy: openInferenceIntegration,
google_adk: openInferenceIntegration,
haystack: openInferenceIntegration,
langchain: openInferenceIntegration,
llamaindex: openInferenceIntegration,
openai_agents_sdk: openInferenceIntegration,
pydantic_ai: openInferenceIntegration,
smolagents: openInferenceIntegration,
eve: eveIntegration,
maple: mapleIntegration,
Expand Down
Loading