diff --git a/packages/agent-sessions/src/session-checks.test.ts b/packages/agent-sessions/src/session-checks.test.ts index 3c9490d4de..b04f6097cd 100644 --- a/packages/agent-sessions/src/session-checks.test.ts +++ b/packages/agent-sessions/src/session-checks.test.ts @@ -1,9 +1,10 @@ +import type { AiSessionGenAiValues } from "@maple/domain/http" import { describe, expect, it } from "vitest" import { buildSessionChecks, type SessionCheck } from "./session-checks" import { buildSessionSummary } from "./session-summary" import { buildSessionTurns } from "./session-turns" -import { agentSpan, llmSpan, toolSpan } from "./span-test-support" +import { agentSpan, llmSpan, makeSpan, toolSpan } from "./span-test-support" const SECOND = 1000 const MINUTE = 60 * SECOND @@ -368,6 +369,8 @@ describe("buildSessionChecks", () => { genAi: { conversationId: "t1", // Inclusive of the cache read, as the default convention counts it. + // No write reported: a Claude prompt this long that read nothing + // missed the cache, rather than being too short to cache. usageInputTokens: 10_000, usageCacheReadInputTokens: cacheRead, usageOutputTokens: 100, @@ -397,6 +400,258 @@ describe("buildSessionChecks", () => { expect(byId(checks(firstTurn()), "prompt-cache").status).toBe("skipped") }) + // Short test conversations (DSPy peaked at 870 tokens; Claude Agent SDK on + // Haiku 4.5 sent 1.2K–2.3K against a 4,096-token minimum) cannot be cached, + // so reading them as misses warned "0% ... missed the cache" on every one. + it("skips the prompt cache when no prompt could have been cached", () => { + const call = (spanId: string, startMs: number, genAi: AiSessionGenAiValues) => + llmSpan({ + spanId, + parentSpanId: "a1", + startMs, + durationMs: SECOND, + genAi: { conversationId: "t1", usageOutputTokens: 50, ...genAi }, + }) + const session = (genAi: (i: number) => AiSessionGenAiValues) => + checks([ + agentSpan({ spanId: "a1", startMs: 0, durationMs: MINUTE, genAi: { conversationId: "t1" } }), + ...[0, 1, 2, 3, 4].map((i) => call(`c${i}`, (i + 1) * SECOND, genAi(i))), + ]) + + const short = byId( + session((i) => ({ usageInputTokens: 700 + 40 * i, usageCacheReadInputTokens: 0 })), + "prompt-cache", + ) + expect(short.status).toBe("skipped") + expect(short.headline).toBe( + "Only 0 model calls had a prompt long enough to cache; at least 4 are needed to judge the prompt cache.", + ) + + // Claude calls that neither wrote nor read the cache were below the + // model's minimum (Haiku 4.5: 4,096): Claude Agent SDK, flue through + // OpenRouter (write stamped 0), OpenRouter Broadcast (write key absent). + for (const claude of [ + { providerName: "anthropic", requestModel: "claude-haiku-4-5", usageCacheCreationInputTokens: 0 }, + { + providerName: "openrouter", + requestModel: "anthropic/claude-haiku-4.5", + usageCacheCreationInputTokens: 0, + }, + { providerName: "anthropic", requestModel: "claude-sonnet-4.5" }, + ]) { + const neverWritten = byId( + session((i) => ({ + ...claude, + usageInputTokens: 1_227 + 200 * i, + usageCacheReadInputTokens: 0, + })), + "prompt-cache", + ) + expect(neverWritten.status).toBe("skipped") + } + + // Per call: a session mixing uncached Claude calls with OpenAI misses is + // judged on the OpenAI calls. + const mixed = byId( + session((i) => + i % 2 === 0 + ? { + providerName: "openai", + requestModel: "gpt-4o-mini", + usageInputTokens: 2_000, + usageCacheReadInputTokens: 0, + } + : { + providerName: "openrouter", + requestModel: "anthropic/claude-haiku-4.5", + usageInputTokens: 1_300, + usageCacheReadInputTokens: 0, + usageCacheCreationInputTokens: 0, + }, + ), + "prompt-cache", + ) + expect(mixed.status).toBe("skipped") + expect(mixed.headline).toMatch(/^Only 3 model calls had a prompt long enough to cache/) + + // Long enough on OpenAI: a miss is a miss, whether or not the emitter + // reported a write bucket (most stamp a zero one on every call). + for (const write of [undefined, 0]) { + const missed = byId( + session(() => ({ + providerName: "openai", + usageInputTokens: 2_000, + usageCacheReadInputTokens: 0, + usageCacheCreationInputTokens: write, + })), + "prompt-cache", + ) + expect(missed.status).toBe("warning") + expect(missed.headline).toBe("Cache hit rate 0% over 4 calls; 4 missed the cache") + } + + // A short opening call (a title, a router) is not the one that wrote the + // cache: the first long call is, and it is not a miss. + const shortFirst = byId( + session((i) => + i === 0 + ? { usageInputTokens: 300, usageCacheReadInputTokens: 0 } + : { usageInputTokens: 2_000, usageCacheReadInputTokens: i === 1 ? 0 : 1_000 }, + ), + "prompt-cache", + ) + expect(shortFirst.headline).toBe("Cache hit rate 50% over 3 calls") + }) + + // OpenRouter Broadcast: every request is a `chat` root carrying the usage + // plus a `provider attempt N` and a `generation` child, all op `chat`. The + // headline read "All 48 model calls" for a session of 16. + it("counts model calls in the provider headline once per call, not per chat span", () => { + const request = (n: number) => { + const traceId = `trace-${n}` + const root = `gen-${n}` + const at = n * MINUTE + return [ + llmSpan({ + spanId: root, + traceId, + spanName: "LLM Generation", + startMs: at, + durationMs: 4_842, + genAi: { usageInputTokens: 58_797, usageOutputTokens: 220, responseId: `gen-${n}` }, + }), + llmSpan({ + spanId: `${root}-attempt`, + parentSpanId: root, + traceId, + spanName: "provider attempt 1: OpenAI", + startMs: at + 184, + durationMs: 2_019, + genAi: { responseId: `gen-${n}:attempt-0` }, + }), + llmSpan({ + spanId: `${root}-generation`, + parentSpanId: root, + traceId, + spanName: "generation", + startMs: at + 200, + durationMs: 4_402, + genAi: { responseId: `gen-${n}:generation` }, + }), + ] + } + const report = checks([...request(0), ...request(1)]) + expect(byId(report, "provider").headline).toBe("All 2 model calls were answered first time") + }) + + // Strands TS writes its finish reasons camelCase (`maxTokens`), which the + // lowercased match against `max_tokens` never caught. + it("reads a cut-off reply whatever the finish reason's spelling", () => { + for (const reason of ["maxTokens", "MAX_TOKENS", "max_tokens"]) { + const report = checks([ + agentSpan({ spanId: "a1", startMs: 0, durationMs: 10 * SECOND }), + llmSpan({ + spanId: "l1", + parentSpanId: "a1", + startMs: SECOND, + durationMs: SECOND, + genAi: { requestMaxTokens: 600, responseFinishReasons: [reason] }, + }), + ]) + const replyLength = byId(report, "reply-length") + expect(replyLength.status).toBe("warning") + expect(replyLength.headline).toMatch(/^1 reply hit the output token limit \(max_tokens 600\)/) + } + }) + + // Vercel AI SDK's `ai.response.finishReason` is kebab-case: a + // `content-filter` refusal read as passed. + it("reads a filtered reply whatever the finish reason's spelling", () => { + const report = checks([ + agentSpan({ spanId: "a1", startMs: 0, durationMs: 10 * SECOND }), + llmSpan({ + spanId: "l1", + parentSpanId: "a1", + startMs: SECOND, + durationMs: SECOND, + genAi: { responseFinishReasons: ["content-filter"] }, + }), + ]) + expect(byId(report, "refusals").status).not.toBe("passed") + expect(byId(report, "refusals").headline).toMatch(/^1 reply was refused or filtered/) + }) + + // The app's own span reports no cache fields; its gateway's mirror, in a + // trace of its own with the same response id, does. The call is counted + // once, and judged by the observation that measured the cache. + it("judges the prompt cache off a gateway mirror when the counted span has none", () => { + const report = checks( + [0, 1, 2, 3, 4].flatMap((i) => [ + llmSpan({ + spanId: `app-${i}`, + traceId: `trace-app-${i}`, + startMs: i * MINUTE, + durationMs: 2 * SECOND, + genAi: { + providerName: "openai", + usageInputTokens: 2_000, + usageOutputTokens: 50, + responseId: `gen-${i}`, + }, + }), + llmSpan({ + spanId: `mirror-${i}`, + traceId: `trace-gateway-${i}`, + spanName: "LLM Generation", + startMs: i * MINUTE + 10, + durationMs: 2 * SECOND, + genAi: { + providerName: "openai", + usageInputTokens: 2_000, + usageCacheReadInputTokens: i === 0 ? 0 : 1_500, + usageOutputTokens: 50, + responseId: `gen-${i}`, + }, + }), + ]), + ) + expect(byId(report, "prompt-cache").headline).toBe("Cache hit rate 75% over 4 calls") + }) + + // Google ADK reports each call twice: `call_llm` (no operation, a model) + // over `generate_content`, both with the same usage. The cache check read + // "over 15 calls" for a session of 8. + it("judges the prompt cache once per model call when a wrapper repeats the usage", () => { + const usage = (i: number): AiSessionGenAiValues => ({ + requestModel: "openrouter/openai/gpt-4o-mini", + usageInputTokens: 2_000, + usageCacheReadInputTokens: i === 0 ? 0 : 1_536, + usageOutputTokens: 32, + }) + const report = checks([ + agentSpan({ spanId: "a1", startMs: 0, durationMs: MINUTE, agentName: "assistant" }), + ...[0, 1, 2, 3, 4].flatMap((i) => [ + makeSpan({ + spanId: `call-${i}`, + parentSpanId: "a1", + spanName: "call_llm", + startMs: (i + 1) * SECOND, + durationMs: 2 * SECOND, + genAi: usage(i), + }), + llmSpan({ + spanId: `gen-${i}`, + parentSpanId: `call-${i}`, + spanName: "generate_content", + startMs: (i + 1) * SECOND, + durationMs: 2 * SECOND, + genAi: { operationName: "generate_content", ...usage(i) }, + }), + ]), + ]) + expect(byId(report, "prompt-cache").headline).toBe("Cache hit rate 77% over 4 calls") + }) + // The one rule behind red and amber, pinned per kind: a class that needs a // fix stays red when survived; everything else survived is amber. it("keeps a survived context overflow red and a survived provider error amber", () => { diff --git a/packages/agent-sessions/src/session-checks.ts b/packages/agent-sessions/src/session-checks.ts index ed215ef6c6..c7b08805d1 100644 --- a/packages/agent-sessions/src/session-checks.ts +++ b/packages/agent-sessions/src/session-checks.ts @@ -21,6 +21,7 @@ import { type SessionVerdict, } from "./session-findings" import { + sessionLlmCalls, spanTokenBuckets, type SessionFailureKind, type SessionSummary, @@ -29,6 +30,7 @@ import { import { classifyAiSpan, isLlmCall, + spanModel, spanStartMs, type SessionTurn, type TurnAnchorKind, @@ -39,6 +41,14 @@ import { const CACHE_HIT_MIN_RATE = 0.5 /** Fewer calls than this say nothing about the cache either way. */ const CACHE_MIN_CALLS = 3 +/** The smallest prompt OpenAI, Anthropic or Gemini will cache: a shorter one + * cannot hit the cache, so missing it says nothing about the prefix. */ +const CACHE_MIN_PROMPT_TOKENS = 1024 +/** The largest Claude minimum (Haiku 4.5, Opus 4.5). Anthropic writes as soon + * as a prompt qualifies, so a Claude call that neither wrote nor read the + * cache under it may just have been too short; one above it missed. Only + * Claude's zero says this: other providers stamp a zero write on every call. */ +const CLAUDE_CACHE_MIN_PROMPT_TOKENS = 4096 /** A headline names this many findings before it counts the rest. */ const HEADLINE_MAX_CLAUSES = 3 @@ -135,7 +145,7 @@ export function buildSessionChecks( ...completionCheck(of("incomplete")), contextWindowCheck(of("contextExceeded"), llmCalls), rateLimitCheck(of("rateLimited")), - providerCheck(of("providerError"), of("providerRetry"), llmCalls.length), + providerCheck(of("providerError"), of("providerRetry"), summary.work.llmCalls), refusalCheck(of("refusal")), replyLengthCheck(of("truncation"), llmCalls), structuredOutputCheck(of("invalidOutput")), @@ -149,7 +159,7 @@ export function buildSessionChecks( ...otherErrorsCheck(errors.filter((finding) => finding.tool === undefined)), repetitionCheck(of("repetition"), summary, coverage), stallCheck(of("stall")), - promptCacheCheck(llmCalls), + promptCacheCheck(cacheObservations(spans)), ].sort( (a, b) => STATUS_RANK[a.status] - STATUS_RANK[b.status] || @@ -613,13 +623,35 @@ function stallCheck(findings: readonly SessionFinding[]): SessionCheck { ) } +/** Claude, by provider or through any gateway by model. */ +const isClaude = (span: AiSessionSpan): boolean => + span.genAi.providerName === "anthropic" || /claude/i.test(spanModel(span) ?? "") + +const reportsCache = (span: AiSessionSpan): boolean => + span.genAi.usageCacheReadInputTokens !== undefined || + span.genAi.usageCacheCreationInputTokens !== undefined + +/** + * One span per model call, so a call its framework also rolled up (ADK + * `call_llm` over `generate_content`) is judged once — and, when the span + * counted for it reported no cache usage, another observation of the same + * response that did (an app span beside its gateway's mirror). + */ +function cacheObservations(spans: readonly AiSessionSpan[]): readonly AiSessionSpan[] { + const byResponse = new Map() + for (const span of spans) { + const responseId = span.genAi.responseId + if (responseId !== undefined && responseId !== "" && reportsCache(span)) + byResponse.set(responseId, span) + } + return sessionLlmCalls(spans).map((call) => + reportsCache(call) ? call : (byResponse.get(call.genAi.responseId ?? "") ?? call), + ) +} + function promptCacheCheck(llmCalls: readonly AiSessionSpan[]): SessionCheck { const identity: CheckIdentity = { id: "prompt-cache", name: "Prompt cache", fixArea: "prompt" } - const reporting = llmCalls.filter( - (span) => - span.genAi.usageCacheReadInputTokens !== undefined || - span.genAi.usageCacheCreationInputTokens !== undefined, - ) + const reporting = llmCalls.filter(reportsCache) if (reporting.length === 0) { return check( identity, @@ -627,21 +659,21 @@ function promptCacheCheck(llmCalls: readonly AiSessionSpan[]): SessionCheck { "No model call reported cache usage, so the prompt cache could not be checked.", ) } - // The first call of a session cannot hit a cache nothing has written yet. - const calls = reporting - .slice(1) - .map(spanTokenBuckets) - .filter((buckets) => buckets !== undefined) - .map((buckets) => ({ - read: buckets.cacheRead, - prompt: buckets.input + buckets.cacheRead + buckets.cacheWrite, - })) - .filter((call) => call.prompt > 0) + // The first cacheable call cannot hit a cache nothing has written yet. + const cacheable = reporting.flatMap((span) => { + const buckets = spanTokenBuckets(span) + if (buckets === undefined) return [] + const prompt = buckets.input + buckets.cacheRead + buckets.cacheWrite + const uncachedClaude = isClaude(span) && buckets.cacheRead + buckets.cacheWrite === 0 + const minimum = uncachedClaude ? CLAUDE_CACHE_MIN_PROMPT_TOKENS : CACHE_MIN_PROMPT_TOKENS + return prompt >= minimum ? [{ read: buckets.cacheRead, prompt }] : [] + }) + const calls = cacheable.slice(1) if (calls.length < CACHE_MIN_CALLS) { return check( identity, "skipped", - `Only ${plural(reporting.length, "model call")} reported cache usage; at least ${CACHE_MIN_CALLS + 1} are needed to judge the prompt cache.`, + `Only ${plural(cacheable.length, "model call")} had a prompt long enough to cache; at least ${CACHE_MIN_CALLS + 1} are needed to judge the prompt cache.`, ) } const rate = diff --git a/packages/agent-sessions/src/session-findings.ts b/packages/agent-sessions/src/session-findings.ts index 3478c6c0d5..8c2d44a056 100644 --- a/packages/agent-sessions/src/session-findings.ts +++ b/packages/agent-sessions/src/session-findings.ts @@ -13,6 +13,7 @@ import { canonicalJSON } from "@maple/query-engine" import { clipDetail, failureDetailText } from "./failure-text" import { failureEvents, + finishReasonsIn, findIdleGaps, isProviderAttempt, shadowedAncestorIds, @@ -46,8 +47,9 @@ const IDENTICAL_RUN_MIN_CALLS = 3 */ const MID_TURN_STALL_MIN_MS = 30_000 -/** Finish reasons that mean the reply was cut off at the output token limit. */ -const TRUNCATION_FINISH_REASONS = new Set(["length", "max_tokens", "max_output_tokens"]) +/** Finish reasons that mean the reply was cut off at the output token limit, + * as `finishReasonsIn` keys. */ +const TRUNCATION_FINISH_REASONS = new Set(["length", "maxtokens", "maxoutputtokens"]) /** * Failure kinds that need a fix whether or not the run survived them: the @@ -351,10 +353,7 @@ function truncationFindings( } function truncationSignal(span: AiSessionSpan): string | undefined { - const reasons = (span.genAi.responseFinishReasons ?? []) - .map((reason) => reason.toLowerCase()) - .filter((reason) => TRUNCATION_FINISH_REASONS.has(reason)) - return reasons.length === 0 ? undefined : reasons.join(",") + return finishReasonsIn(span, TRUNCATION_FINISH_REASONS) } function repetitionFindings(turns: readonly SessionTurn[]): SessionFinding[] { diff --git a/packages/agent-sessions/src/session-summary.test.ts b/packages/agent-sessions/src/session-summary.test.ts index 580572b624..a188b861ac 100644 --- a/packages/agent-sessions/src/session-summary.test.ts +++ b/packages/agent-sessions/src/session-summary.test.ts @@ -114,6 +114,38 @@ describe("buildSessionSummary — time", () => { expect(summary.agentTime.segments.map((entry) => entry.kind)).toEqual(["inference", "tool"]) expect(summary.agentTime.totalMs).toBe(7 * SECOND) }) + + // OpenRouter Broadcast nests a `provider attempt` and a `generation` span, + // both op `chat`, under each `LLM Generation`: summing all three read 204.8s + // of agent time in a 125.5s session. + it("charges a model call observed at several levels once, at the outermost", () => { + const summary = summarize([ + llmSpan({ spanId: "root", spanName: "LLM Generation", startMs: 0, durationMs: 4_842 }), + llmSpan({ spanId: "attempt", parentSpanId: "root", startMs: 184, durationMs: 2_019 }), + llmSpan({ + spanId: "generation", + parentSpanId: "root", + startMs: 200, + durationMs: 4_402, + ttftSeconds: 3.6, + }), + ]) + + expect(summary.agentTime.totalMs).toBe(4_842) + // The TTFT only the nested span reported still splits the call. + expect(segment(summary.agentTime.segments, "ttft")).toBe(3_600) + expect(segment(summary.agentTime.segments, "inference")).toBe(1_242) + }) + + it("still charges a model call a tool made inside another call", () => { + const summary = summarize([ + llmSpan({ spanId: "outer", startMs: 0, durationMs: 10 * SECOND }), + toolSpan({ spanId: "tool", parentSpanId: "outer", startMs: SECOND, durationMs: 4 * SECOND }), + llmSpan({ spanId: "inner", parentSpanId: "tool", startMs: 2 * SECOND, durationMs: 2 * SECOND }), + ]) + + expect(segment(summary.agentTime.segments, "inference")).toBe(12 * SECOND) + }) }) describe("buildSessionSummary — failed", () => { diff --git a/packages/agent-sessions/src/session-summary.ts b/packages/agent-sessions/src/session-summary.ts index a7e2e71aa8..9c41affc31 100644 --- a/packages/agent-sessions/src/session-summary.ts +++ b/packages/agent-sessions/src/session-summary.ts @@ -241,7 +241,7 @@ export interface SessionSummary { const RATE_LIMIT_PATTERN = /\b429\b|rate.?limit|too.many.requests|resource.exhausted|overloaded/i const CONTEXT_EXCEEDED_PATTERN = /context.{0,16}(length|window|limit)|maximum.context|prompt is too long|too many tokens/i -const REFUSAL_FINISH_REASONS = new Set(["refusal", "content_filter"]) +const REFUSAL_FINISH_REASONS = new Set(["refusal", "contentfilter"]) export function buildSessionSummary({ spans, @@ -356,29 +356,55 @@ export function findIdleGaps(spans: readonly AiSessionSpan[]): readonly IdleGap[ * A TTFT splits its own span: the wait is not inference, and a session whose * time is mostly first-token latency is a different session from one that is * mostly generation. Agent and non-AI spans contribute nothing — an agent span - * covers its children, and adding it would count the same work twice. + * covers its children, and adding it would count the same work twice. For the + * same reason an inference span inside another, with no agent or tool between + * them, is charged nothing: its time is already inside the outer span's + * (OpenRouter's `LLM Generation` over its `generation`, ADK's `call_llm` over + * `generate_content`). The outermost carries the time, and borrows a TTFT from + * a nested span when it reported none. */ export function computeAgentTime(spans: readonly AiSessionSpan[]): SessionAgentTime { const totals = new Map() const add = (kind: AgentTimeKind, ms: number) => { if (ms > 0) totals.set(kind, (totals.get(kind) ?? 0) + ms) } + const byId = new Map(spans.map((span) => [span.spanId, span])) + const outermostCall = (span: AiSessionSpan): AiSessionSpan => { + let call = span + const seen = new Set([span.spanId]) + let parent = byId.get(span.parentSpanId) + while (parent !== undefined && !seen.has(parent.spanId)) { + seen.add(parent.spanId) + const category = classifyAiSpan(parent) + // An agent or a tool in between ran calls of its own. + if (category === "agent" || category === "tool") break + if (category === "inference") call = parent + parent = byId.get(parent.parentSpanId) + } + return call + } + const calls: AiSessionSpan[] = [] + const nestedTtftMs = new Map() for (const span of spans) { - const spanStart = spanStartMs(span) - const spanEnd = spanEndMs(span) const category = classifyAiSpan(span) - if (category !== "tool" && category !== "inference") continue - if (category === "tool") { - add("tool", spanEnd - spanStart) + if (category === "tool") add("tool", span.durationMs) + if (category !== "inference") continue + const call = outermostCall(span) + if (call === span) { + calls.push(span) continue } const ttftMs = spanTtftMs(span) + if (ttftMs !== undefined && !nestedTtftMs.has(call.spanId)) nestedTtftMs.set(call.spanId, ttftMs) + } + for (const call of calls) { + const ttftMs = spanTtftMs(call) ?? nestedTtftMs.get(call.spanId) // A TTFT longer than the span itself is instrumentation disagreeing with // itself; the span's own duration is the one both classes must fit in. - const ttft = ttftMs === undefined ? 0 : Math.min(ttftMs, spanEnd - spanStart) + const ttft = ttftMs === undefined ? 0 : Math.min(ttftMs, call.durationMs) add("ttft", ttft) - add("inference", spanEnd - spanStart - ttft) + add("inference", call.durationMs - ttft) } const segments = AGENT_TIME_KIND_ORDER.map((kind) => ({ kind, ms: totals.get(kind) ?? 0 })).filter( @@ -582,6 +608,14 @@ function countedLlmCalls( return [...unkeyed, ...byResponse.values()].sort((a, b) => spanStartMs(a) - spanStartMs(b)) } +/** The spans `work.llmCalls` counts, in start order: for a reading that needs + * one span per model call rather than every observation of it. */ +export function sessionLlmCalls(spans: readonly AiSessionSpan[]): readonly AiSessionSpan[] { + const byId = new Map(spans.map((span) => [span.spanId, span])) + const usage = countableUsageSpans(spans, byId) + return countedLlmCalls(spans, byId, usage.bySpan, usage.costs) +} + /** * Each reporter charged to the NEAREST ancestor that also reports, so a * two-level roll-up subtracts each figure once rather than at every level. @@ -851,9 +885,19 @@ function failureSignal(span: AiSessionSpan): string | undefined { } function refusalSignal(span: AiSessionSpan): string | undefined { + return finishReasonsIn(span, REFUSAL_FINISH_REASONS) +} + +/** + * The span's finish reasons that are one of `keys`, lowercased and joined — or + * `undefined`. Matched without case or separators, since vendors spell one + * reason `max_tokens`, `MAX_TOKENS` and `maxTokens` (Strands TS); `keys` are + * written that way too. + */ +export function finishReasonsIn(span: AiSessionSpan, keys: ReadonlySet): string | undefined { const reasons = (span.genAi.responseFinishReasons ?? []) .map((reason) => reason.toLowerCase()) - .filter((reason) => REFUSAL_FINISH_REASONS.has(reason)) + .filter((reason) => keys.has(reason.replace(/[^a-z0-9]/g, ""))) return reasons.length === 0 ? undefined : reasons.join(",") } diff --git a/packages/agent-sessions/src/session-turns.test.ts b/packages/agent-sessions/src/session-turns.test.ts index f73385d983..685e9e2bdb 100644 --- a/packages/agent-sessions/src/session-turns.test.ts +++ b/packages/agent-sessions/src/session-turns.test.ts @@ -3,6 +3,7 @@ import { describe, expect, it } from "vitest" import type { AiSessionSpan } from "@maple/domain/http" import { agentSpan, llmSpan, makeSpan, toolSpan, userMessages } from "./span-test-support" +import { buildSessionSummary } from "./session-summary" import { buildSessionTurns, classifyAiSpan, isLlmCall, spanTtftMs } from "./session-turns" const SECOND = 1000 @@ -170,6 +171,76 @@ describe("buildSessionTurns", () => { expect(Number.isFinite(turns[0]!.startMs)).toBe(true) }) + // Microsoft Agent Framework workflows: `workflow.build` is a one-span trace + // of its own, read as an agent root by its name, ahead of `workflow.run`. + // It became an empty turn 1 with no label, and the session lost its title. + it("folds an anchor that opened no work into the turn after it", () => { + const maf = { + vendorId: "microsoft_agent_framework", + sessionId: "wf-1", + genAi: { conversationId: "wf-1" }, + } + const spans = [ + makeSpan({ + ...maf, + spanId: "build", + traceId: "trace-build", + spanName: "workflow.build", + startMs: 0, + durationMs: 2, + }), + makeSpan({ + ...maf, + spanId: "run", + traceId: "trace-run", + spanName: "workflow.run", + startMs: 50, + durationMs: 10 * SECOND, + }), + agentSpan({ + ...maf, + spanId: "orchestrator", + parentSpanId: "run", + traceId: "trace-run", + startMs: 60, + durationMs: 4 * SECOND, + }), + llmSpan({ + ...maf, + spanId: "chat", + parentSpanId: "orchestrator", + traceId: "trace-run", + startMs: 70, + durationMs: 4 * SECOND, + genAi: { + conversationId: "wf-1", + inputMessages: userMessages("Produce a mini briefing about Amsterdam"), + }, + }), + ] + const turns = buildSessionTurns(spans) + + expect(turns).toHaveLength(1) + expect(turns[0]!.spans.map((span) => span.spanId)).toEqual(["build", "run", "orchestrator", "chat"]) + expect(turns[0]!.label).toBe("Produce a mini briefing about Amsterdam") + expect(buildSessionSummary({ spans, turns }).title).toBe("Produce a mini briefing about Amsterdam") + }) + + // A pause after it says the invocation stood on its own: folding it would + // read the pause as a stall inside the next turn. + it("keeps a workless anchor followed by a pause as its own turn", () => { + const turns = buildSessionTurns([ + agentSpan({ spanId: "quiet", startMs: 0, durationMs: SECOND }), + agentSpan({ spanId: "busy", startMs: 60 * SECOND, durationMs: 10 * SECOND }), + llmSpan({ spanId: "chat", parentSpanId: "busy", startMs: 61 * SECOND, durationMs: SECOND }), + ]) + + expect(turns.map((turn) => turn.spans.map((span) => span.spanId))).toEqual([ + ["quiet"], + ["busy", "chat"], + ]) + }) + it("falls back to root agent invocations when no conversation id exists", () => { const turns = buildSessionTurns([ agentSpan({ spanId: "agent-1", startMs: 0, durationMs: 10 * SECOND }), @@ -298,6 +369,49 @@ describe("buildSessionTurns", () => { expect(turns[0]!.label).toBe("deploy the worker") }) + // LangGraph with a checkpointer, through the OpenInference dual-write: the + // `model` node (a CHAIN span, no operation) starts before its model call and + // carries only the thread's FIRST message, so every turn read turn 1's prompt. + it("labels from a model call before a framework span that started earlier", () => { + const turn = (n: number, prompts: readonly string[]) => { + const at = n * 60 * SECOND + return [ + agentSpan({ + spanId: `assistant-${n}`, + startMs: at, + durationMs: 2 * SECOND, + agentName: "assistant", + }), + makeSpan({ + spanId: `model-${n}`, + parentSpanId: `assistant-${n}`, + spanName: "model", + startMs: at + 10, + durationMs: SECOND, + vendorId: "unknown:openinference", + genAi: { inputMessages: userMessages(prompts[0]!) }, + }), + llmSpan({ + spanId: `chat-${n}`, + parentSpanId: `model-${n}`, + spanName: "ChatOpenAI", + startMs: at + 12, + durationMs: SECOND, + genAi: { inputMessages: userMessages(...prompts) }, + }), + ] + } + const turns = buildSessionTurns([ + ...turn(0, ["Hi! Briefly introduce yourself."]), + ...turn(1, ["Hi! Briefly introduce yourself.", "What's the weather in Berlin?"]), + ]) + + expect(turns.map((t) => t.label)).toEqual([ + "Hi! Briefly introduce yourself.", + "What's the weather in Berlin?", + ]) + }) + it("has no label when message content was not captured", () => { const turns = buildSessionTurns([agentSpan({ spanId: "agent", startMs: 0, durationMs: SECOND })]) @@ -496,6 +610,15 @@ describe("turn labels", () => { expect(long?.endsWith("…")).toBe(true) }) + // smolagents sends every task as "New task:\n", so every turn and the + // session title read "New task:". + it("reads past smolagents' lead-in", () => { + expect(labelFor([{ role: "user", content: "New task:\nWhat is 17 * 23?" }])).toBe("What is 17 * 23?") + expect(labelFor([{ role: "user", content: "New task:" }])).toBe("New task:") + // A user's own line ending in a colon is the prompt itself. + expect(labelFor([{ role: "user", content: "Fix this:\n```ts" }])).toBe("Fix this:") + }) + // Vendors write "User" as readily as "user", and the transcript's own row // builders already read the role case-insensitively. it("reads a capitalised role", () => { diff --git a/packages/agent-sessions/src/session-turns.ts b/packages/agent-sessions/src/session-turns.ts index 1299a5be73..4b85cc9fc3 100644 --- a/packages/agent-sessions/src/session-turns.ts +++ b/packages/agent-sessions/src/session-turns.ts @@ -163,6 +163,11 @@ export interface SessionTurn { readonly traceIds: readonly string[] } +const WORK_CATEGORIES: ReadonlySet = new Set(["inference", "tool"]) +/** How soon before the next turn a workless anchor must end to be its setup: + * the pause the session summary starts calling idle. */ +const SETUP_LEAD_MAX_MS = 5_000 + interface TurnAnchor { readonly span: AiSessionSpan readonly kind: TurnAnchorKind @@ -218,6 +223,37 @@ export function buildSessionTurns(spans: readonly AiSessionSpan[]): readonly Ses buckets[turnOf(span) ?? cursor].push(span) } + // On rules 2 and 3 the boundary is a guess, and an anchor that opened no + // work — no model or tool call, no prompt, nothing failed — right before the + // next turn is that turn's setup: Microsoft Agent Framework's one-span + // `workflow.build` trace ahead of its `workflow.run` became an empty turn 1 + // that also took the session's title. Its spans join the next turn, as spans + // before the first anchor join turn 1; the cursor filled the buckets in start + // order, so they stay in it. One followed by a pause is left alone, and a + // session with no work anywhere keeps its anchors. + const opened = (bucket: readonly AiSessionSpan[]) => + bucket.some( + (span) => + WORK_CATEGORIES.has(classifyAiSpan(span)) || + spanFailed(span) || + lastUserMessageText(span.genAi.inputMessages) !== undefined, + ) + if (anchors[0]?.kind !== "conversation" && buckets.some(opened)) { + for (let i = 0; i < buckets.length - 1; i++) { + const bucket = buckets[i] + const next = buckets[i + 1] + if (opened(bucket)) continue + const endMs = bucket.reduce( + (max, span) => Math.max(max, spanEndMs(span)), + Number.NEGATIVE_INFINITY, + ) + if (next[0] !== undefined && spanStartMs(next[0]) - endMs > SETUP_LEAD_MAX_MS) continue + for (const span of next) bucket.push(span) + buckets[i + 1] = bucket + buckets[i] = [] + } + } + // A turn with no spans has no start, no end and nothing to draw. Rule 1 can no // longer produce one — an anchor always carries its own id, so it lands in its // own bucket even when another anchor shares its millisecond — but two @@ -401,12 +437,15 @@ function findAnchors(ordered: readonly AiSessionSpan[]): readonly TurnAnchor[] { * * The anchor is asked first — on a `chat`-shaped span `gen_ai.input.messages` * is the whole history sent to the model, so a descendant several turns deep - * still carries turn 1's opening prompt. + * still carries turn 1's opening prompt. Model calls come next: the history a + * model was sent ends on the turn's prompt, while a framework's own node span + * may carry only what the thread started with (LangGraph's `model` node under a + * checkpointer holds the thread's first message on every turn). */ function turnLabel(anchor: AiSessionSpan, turnSpans: readonly AiSessionSpan[]): string | undefined { const fromAnchor = lastUserMessageText(anchor.genAi.inputMessages) if (fromAnchor !== undefined) return fromAnchor - for (const span of turnSpans) { + for (const span of [...turnSpans.filter(isLlmCall), ...turnSpans]) { const text = lastUserMessageText(span.genAi.inputMessages) if (text !== undefined) return text } @@ -461,14 +500,23 @@ function messageText(value: unknown): string | undefined { return undefined } -/** The message's first non-empty line, collapsed to one line's worth of text. */ +/** A framework's own heading over the prompt, which names what follows rather + * than saying it: smolagents opens every task with "New task:". Matched + * literally, since a user's "Fix this:" over pasted code is the prompt. */ +const LEAD_IN = /^new task:$/i + +/** The message's first non-empty line, collapsed to one line's worth of text — + * or the line after it, when the first is a lead-in. */ function proseLine(value: string): string | undefined { + const lines: string[] = [] for (const rawLine of value.split("\n")) { const line = rawLine.trim().replace(/\s+/g, " ") - if (line.length === 0) continue - return line.length > MAX_LABEL_LENGTH ? `${line.slice(0, MAX_LABEL_LENGTH - 1)}…` : line + if (line.length > 0 && lines.push(line) === 2) break } - return undefined + const [first, second] = lines + if (first === undefined) return undefined + const line = second !== undefined && LEAD_IN.test(first) ? second : first + return line.length > MAX_LABEL_LENGTH ? `${line.slice(0, MAX_LABEL_LENGTH - 1)}…` : line } function isRecord(value: unknown): value is Record {