From 72b17a6a314b359f059c80b557a1f0f9edaffcb5 Mon Sep 17 00:00:00 2001 From: Waishnav Date: Sat, 8 Aug 2026 23:44:07 +0530 Subject: [PATCH 1/2] feat(agents): normalize provider observations --- src/local-agent-provider-observations.test.ts | 80 ++++++ src/local-agent-provider-observations.ts | 257 ++++++++++++++++++ 2 files changed, 337 insertions(+) create mode 100644 src/local-agent-provider-observations.test.ts create mode 100644 src/local-agent-provider-observations.ts diff --git a/src/local-agent-provider-observations.test.ts b/src/local-agent-provider-observations.test.ts new file mode 100644 index 00000000..82cdc8c0 --- /dev/null +++ b/src/local-agent-provider-observations.test.ts @@ -0,0 +1,80 @@ +import assert from "node:assert/strict"; +import { + extractAcpObservations, + extractAcpUsage, + extractClaudeObservations, + extractClaudeUsage, + extractCodexObservations, + extractCodexUsage, + extractOpenCodeObservations, + extractOpenCodeUsage, + extractPiObservations, + extractPiUsage, +} from "./local-agent-provider-observations.js"; + +{ + const item = { + type: "command_execution", + id: "cmd-1", + command: "npm test", + status: "completed", + usage: { input_tokens: 12, output_tokens: 8 }, + }; + const codexActivity = extractCodexObservations([item]).find((entry) => entry.kind === "activity"); + assert.equal(codexActivity?.activityId, "cmd-1"); + assert.equal(codexActivity?.toolStatus, "completed"); + assert.deepEqual(extractCodexUsage({ usage: { input_tokens: 12, output_tokens: 8 } }), { + inputTokens: 12, + outputTokens: 8, + totalTokens: 20, + }); +} + +{ + const messages = [ + { type: "assistant", content: [{ type: "tool_use", id: "tool-1", name: "bash", input: { command: "pwd" } }] }, + { type: "user", content: [{ type: "tool_result", tool_use_id: "tool-1", content: "ok" }] }, + { type: "result", usage: { input_tokens: 20, output_tokens: 4 } }, + ]; + assert.deepEqual(extractClaudeObservations(messages).filter((entry) => entry.kind === "activity").map((entry) => entry.toolStatus), ["started", "completed"]); + assert.deepEqual(extractClaudeUsage(messages), { inputTokens: 20, outputTokens: 4, totalTokens: 24 }); +} + +{ + const messages = [{ + info: { tokens: { input: 30, output: 5 } }, + parts: [{ type: "tool", callID: "call-1", tool: "grep", state: { status: "completed", output: "match" } }], + }]; + const openCodeActivity = extractOpenCodeObservations(messages).find((entry) => entry.kind === "activity"); + assert.equal(openCodeActivity?.toolName, "grep"); + assert.equal(openCodeActivity?.toolStatus, "completed"); + assert.deepEqual(extractOpenCodeUsage(messages), { inputTokens: 30, outputTokens: 5, totalTokens: 35 }); +} + +{ + const event = { type: "tool_result", toolCallId: "pi-1", toolName: "read", status: "success", usage: { input: 4, output: 3 } }; + const piActivity = extractPiObservations([event]).find((entry) => entry.kind === "activity"); + assert.equal(piActivity?.activityId, "pi-1"); + assert.equal(piActivity?.toolStatus, "completed"); + assert.deepEqual(extractPiUsage(event), { inputTokens: 4, outputTokens: 3, totalTokens: 7 }); +} + +for (const provider of ["cursor", "copilot"] as const) { + const update = { + sessionUpdate: "tool_call_update", + toolCallId: `${provider}-1`, + toolName: "edit", + status: "completed", + usage_update: { input_tokens: 9, output_tokens: 2 }, + }; + const acpActivity = extractAcpObservations(update).find((entry) => entry.kind === "activity"); + assert.equal(acpActivity?.toolName, "edit"); + assert.equal(acpActivity?.toolStatus, "completed"); + assert.deepEqual(extractAcpUsage(update), { + inputTokens: 9, + outputTokens: 2, + totalTokens: 11, + }); +} + +console.log("local-agent-provider-observations.test.ts: ok"); diff --git a/src/local-agent-provider-observations.ts b/src/local-agent-provider-observations.ts new file mode 100644 index 00000000..d22e272e --- /dev/null +++ b/src/local-agent-provider-observations.ts @@ -0,0 +1,257 @@ +import type { + LocalAgentActivityObservation, + LocalAgentObservation, + LocalAgentTokenUsage, + LocalAgentToolStatus, +} from "./local-agent-observations.js"; + +/** Structural provider parsers kept independent from each SDK's type surface. */ + +export function extractCodexObservations(value: unknown): LocalAgentObservation[] { + return extractArray(value).flatMap((item) => { + const record = asRecord(item); + if (!record) return []; + const type = stringValue(record, ["type", "item_type"]) ?? ""; + if (!/(tool|command|mcp|function)/i.test(type) && !hasAny(record, ["tool_name", "toolName", "command"])) { + return []; + } + return [activityFromRecord(record, type || "tool")]; + }); +} + +export function extractCodexUsage(value: unknown): LocalAgentTokenUsage | undefined { + return findUsage(value, ["usage", "token_usage", "tokens"]); +} + +export function extractClaudeObservations(value: unknown): LocalAgentObservation[] { + return extractArray(value).flatMap((item) => { + const record = asRecord(item); + if (!record) return []; + const type = stringValue(record, ["type"]) ?? ""; + const content = Array.isArray(record.content) ? record.content : []; + const results: LocalAgentObservation[] = []; + for (const part of content) { + const partRecord = asRecord(part); + if (!partRecord) continue; + const partType = stringValue(partRecord, ["type"]) ?? ""; + if (partType === "tool_use") { + results.push({ + kind: "activity", + activityId: stringValue(partRecord, ["id", "tool_use_id"]), + toolName: stringValue(partRecord, ["name", "tool_name"]) ?? "tool", + toolStatus: "started", + detail: stringifyDetail(partRecord.input), + }); + } else if (partType === "tool_result") { + results.push({ + kind: "activity", + activityId: stringValue(partRecord, ["tool_use_id", "id"]), + toolName: "tool", + toolStatus: partRecord.is_error === true ? "failed" : "completed", + detail: stringifyDetail(partRecord.content), + }); + } + } + if (type === "tool_progress" || type === "tool_use") { + results.push(activityFromRecord(record, type)); + } + return results; + }); +} + +export function extractClaudeUsage(value: unknown): LocalAgentTokenUsage | undefined { + const direct = findUsage(value, ["usage", "token_usage", "tokens"]); + if (direct) return direct; + const record = asRecord(value); + const modelUsage = asRecord(record?.modelUsage ?? record?.model_usage); + if (!modelUsage) return undefined; + const totals: LocalAgentTokenUsage = {}; + for (const model of Object.values(modelUsage)) { + const usage = normalizeUsage(model); + if (!usage) continue; + addUsage(totals, usage); + } + return Object.keys(totals).length ? totals : undefined; +} + +export function extractOpenCodeObservations(value: unknown): LocalAgentObservation[] { + const root = unwrap(value); + const messages = Array.isArray(root) ? root : arrayValue(asRecord(root), "messages"); + const records = messages ?? (asRecord(root) ? [root] : []); + const output: LocalAgentObservation[] = []; + for (const message of records) { + const messageRecord = asRecord(message); + if (!messageRecord) continue; + const parts = arrayValue(messageRecord, "parts") ?? arrayValue(messageRecord, "content") ?? []; + for (const part of parts) { + const partRecord = asRecord(part); + if (!partRecord || stringValue(partRecord, ["type"]) !== "tool") continue; + const state = asRecord(partRecord.state) ?? partRecord; + output.push({ + kind: "activity", + activityId: stringValue(partRecord, ["callID", "callId", "id"]) ?? stringValue(state, ["callID", "callId", "id"]), + toolName: stringValue(partRecord, ["tool", "name"]) ?? "tool", + toolStatus: mapToolStatus(stringValue(state, ["status"]) ?? "updated"), + detail: stringifyDetail(state.output ?? state.error ?? state.title), + }); + } + } + return output; +} + +export function extractOpenCodeUsage(value: unknown): LocalAgentTokenUsage | undefined { + return findUsage(value, ["usage", "tokens"]); +} + +export function extractPiObservations(value: unknown): LocalAgentObservation[] { + return extractArray(value).flatMap((item) => { + const record = asRecord(item); + if (!record) return []; + const type = stringValue(record, ["type", "event"]) ?? ""; + if (!/(tool|function)/i.test(type)) return []; + return [activityFromRecord(record, type)]; + }); +} + +export function extractPiUsage(value: unknown): LocalAgentTokenUsage | undefined { + return findUsage(value, ["usage", "tokens"]); +} + +export function extractAcpObservations(value: unknown): LocalAgentObservation[] { + const record = asRecord(value); + if (!record) return []; + const update = asRecord(record.update) ?? record; + const sessionUpdate = stringValue(update, ["sessionUpdate", "session_update"]) ?? ""; + if (!/(tool|function)/i.test(sessionUpdate) && !hasAny(update, ["toolCallId", "tool_call_id", "toolName"])) { + return []; + } + return [activityFromRecord(update, sessionUpdate || "tool")]; +} + +export function extractAcpUsage(value: unknown): LocalAgentTokenUsage | undefined { + return findUsage(value, ["usage", "usage_update", "tokens"]); +} + +function activityFromRecord(record: Record, fallbackName: string): LocalAgentActivityObservation { + const type = stringValue(record, ["type", "event", "sessionUpdate", "session_update"]) ?? fallbackName; + return { + kind: "activity", + activityId: stringValue(record, ["id", "callId", "call_id", "toolCallId", "tool_call_id"]), + toolName: stringValue(record, ["toolName", "tool_name", "tool", "name", "command"]) ?? fallbackName, + toolStatus: mapToolStatus(stringValue(record, ["status", "state", "event"]) ?? type), + message: stringValue(record, ["message", "title", "text"]), + detail: stringifyDetail(record.output ?? record.result ?? record.error), + }; +} + +function mapToolStatus(value: string): LocalAgentToolStatus { + const normalized = value.toLowerCase(); + if (/(fail|error|abort|cancel)/.test(normalized)) return "failed"; + if (/(complete|done|success|finish|result)/.test(normalized)) return "completed"; + if (/(start|begin|call|pending)/.test(normalized)) return "started"; + return "updated"; +} + +function findUsage(value: unknown, keys: string[]): LocalAgentTokenUsage | undefined { + const candidates = [value, ...extractArray(value)]; + for (const record of candidates.reverse()) { + for (const key of keys) { + const usage = normalizeUsage(asRecord(record)?.[key]); + if (usage) return usage; + } + for (const nestedKey of ["info", "message", "data", "result"]) { + const nested = asRecord(asRecord(record)?.[nestedKey]); + if (!nested) continue; + for (const key of keys) { + const usage = normalizeUsage(nested[key]); + if (usage) return usage; + } + const usage = normalizeUsage(nested); + if (usage) return usage; + } + const usage = normalizeUsage(record); + if (usage) return usage; + } + return undefined; +} + +function normalizeUsage(value: unknown): LocalAgentTokenUsage | undefined { + if (!value || typeof value !== "object" || Array.isArray(value)) return undefined; + const record = value as Record; + const usage: LocalAgentTokenUsage = {}; + const input = numberValue(record, ["input", "inputTokens", "input_tokens", "promptTokens", "prompt_tokens"]); + const output = numberValue(record, ["output", "outputTokens", "output_tokens", "completionTokens", "completion_tokens"]); + const total = numberValue(record, ["total", "totalTokens", "total_tokens"]); + const cache = asRecord(record.cache); + const cacheRead = numberValue(record, ["cacheRead", "cache_read", "cacheReadTokens", "cache_read_tokens"]) ?? numberValue(cache, ["read", "readTokens", "read_tokens"]); + const cacheWrite = numberValue(record, ["cacheWrite", "cache_write", "cacheWriteTokens", "cache_write_tokens"]) ?? numberValue(cache, ["write", "writeTokens", "write_tokens"]); + if (input !== undefined) usage.inputTokens = input; + if (output !== undefined) usage.outputTokens = output; + if (total !== undefined) usage.totalTokens = total; + if (cacheRead !== undefined) usage.cacheReadTokens = cacheRead; + if (cacheWrite !== undefined) usage.cacheWriteTokens = cacheWrite; + if (usage.totalTokens === undefined && usage.inputTokens !== undefined && usage.outputTokens !== undefined) { + usage.totalTokens = usage.inputTokens + usage.outputTokens; + } + return Object.keys(usage).length ? usage : undefined; +} + +function addUsage(target: LocalAgentTokenUsage, next: LocalAgentTokenUsage): void { + for (const key of ["inputTokens", "outputTokens", "totalTokens", "cacheReadTokens", "cacheWriteTokens"] as const) { + const value = next[key]; + if (value !== undefined) target[key] = (target[key] ?? 0) + value; + } +} + +function extractArray(value: unknown): unknown[] { + if (Array.isArray(value)) return value; + const record = asRecord(value); + if (!record) return []; + for (const key of ["items", "messages", "events", "updates", "data", "result"]) { + if (Array.isArray(record[key])) return record[key] as unknown[]; + } + return [value]; +} + +function unwrap(value: unknown): unknown { + const record = asRecord(value); + return record?.data ?? record?.result ?? value; +} + +function arrayValue(record: Record | undefined, key: string): unknown[] | undefined { + return record && Array.isArray(record[key]) ? record[key] as unknown[] : undefined; +} + +function asRecord(value: unknown): Record | undefined { + if (!value || typeof value !== "object" || Array.isArray(value)) return undefined; + return value as Record; +} + +function hasAny(record: Record, keys: string[]): boolean { + return keys.some((key) => record[key] !== undefined); +} + +function stringValue(record: Record | undefined, keys: string[]): string | undefined { + if (!record) return undefined; + for (const key of keys) if (typeof record[key] === "string" && record[key]) return record[key] as string; + return undefined; +} + +function numberValue(record: Record | undefined, keys: string[]): number | undefined { + if (!record) return undefined; + for (const key of keys) { + const value = record[key]; + if (typeof value === "number" && Number.isFinite(value) && value >= 0) return Math.floor(value); + } + return undefined; +} + +function stringifyDetail(value: unknown): string | undefined { + if (value === undefined || value === null) return undefined; + if (typeof value === "string") return value; + try { + return JSON.stringify(value); + } catch { + return undefined; + } +} From d5144bbbca310fe2c21e6b2bf14f79e015650861 Mon Sep 17 00:00:00 2001 From: Waishnav Date: Sat, 8 Aug 2026 23:44:12 +0530 Subject: [PATCH 2/2] feat(agents): wire provider observation streams --- package.json | 2 +- src/cli-output.test.ts | 3 +++ src/cli-output.ts | 13 ++++++++++++ src/local-agent-adapters.ts | 42 ++++++++++++++++++++++++++++++++++++- src/local-agent-runtime.ts | 26 +++++++++++++++++++++++ 5 files changed, 84 insertions(+), 2 deletions(-) diff --git a/package.json b/package.json index a6f01226..7182bb67 100644 --- a/package.json +++ b/package.json @@ -28,7 +28,7 @@ "dev": "node scripts/dev-server.mjs", "postinstall": "node scripts/fix-node-pty-permissions.mjs", "start": "node dist/cli.js serve", - "test": "tsx src/config.test.ts && tsx src/cli-workspace.test.ts && tsx src/cli-output.test.ts && tsx src/open-workspace-capabilities.test.ts && tsx src/ui/card-types.test.ts && tsx src/ui/patch-display.test.ts && tsx src/ui/tool-display.test.ts && tsx src/apply-patch.test.ts && tsx src/process-platform.test.ts && tsx src/process-sessions.test.ts && tsx src/mcp-sessions.test.ts && tsx src/server-shutdown.test.ts && tsx src/local-agent-runtime.test.ts && tsx src/local-agent-adapters.test.ts && tsx src/local-agent-availability.test.ts && tsx src/local-agent-capabilities.test.ts && tsx src/local-agent-catalog.test.ts && tsx src/local-agent-profiles.test.ts && tsx src/local-agent-resolution.test.ts && tsx src/local-agent-targets.test.ts && tsx src/local-agent-store.test.ts && tsx src/roots.test.ts && tsx src/skill-install.test.ts && tsx src/skills.test.ts && tsx src/workspaces.test.ts && tsx src/review-checkpoints.test.ts && tsx src/oauth-store.test.ts && tsx src/cli.test.ts && tsx src/workflow-contracts.test.ts && tsx src/workflow-errors.test.ts && tsx src/workflow-types.test.ts && tsx src/workflow-store.test.ts && tsx src/workflow-lifecycle.test.ts && tsx src/workflow-view.test.ts && tsx src/workflow-summary.test.ts && tsx src/workflow-tui.test.ts && tsx src/workflow-script.test.ts && tsx src/workflow-sandbox.test.ts && tsx src/workflow-engine.test.ts && tsx src/workflow-files.test.ts && tsx src/workflow-launch.test.ts && tsx src/workflow-replay.test.ts && tsx src/workflow-schema.test.ts", + "test": "tsx src/config.test.ts && tsx src/cli-workspace.test.ts && tsx src/cli-output.test.ts && tsx src/open-workspace-capabilities.test.ts && tsx src/ui/card-types.test.ts && tsx src/ui/patch-display.test.ts && tsx src/ui/tool-display.test.ts && tsx src/apply-patch.test.ts && tsx src/process-platform.test.ts && tsx src/process-sessions.test.ts && tsx src/mcp-sessions.test.ts && tsx src/server-shutdown.test.ts && tsx src/local-agent-runtime.test.ts && tsx src/local-agent-adapters.test.ts && tsx src/local-agent-provider-observations.test.ts && tsx src/local-agent-availability.test.ts && tsx src/local-agent-capabilities.test.ts && tsx src/local-agent-catalog.test.ts && tsx src/local-agent-profiles.test.ts && tsx src/local-agent-resolution.test.ts && tsx src/local-agent-targets.test.ts && tsx src/local-agent-store.test.ts && tsx src/roots.test.ts && tsx src/skill-install.test.ts && tsx src/skills.test.ts && tsx src/workspaces.test.ts && tsx src/review-checkpoints.test.ts && tsx src/oauth-store.test.ts && tsx src/cli.test.ts && tsx src/workflow-contracts.test.ts && tsx src/workflow-errors.test.ts && tsx src/workflow-types.test.ts && tsx src/workflow-store.test.ts && tsx src/workflow-lifecycle.test.ts && tsx src/workflow-view.test.ts && tsx src/workflow-summary.test.ts && tsx src/workflow-tui.test.ts && tsx src/workflow-script.test.ts && tsx src/workflow-sandbox.test.ts && tsx src/workflow-engine.test.ts && tsx src/workflow-files.test.ts && tsx src/workflow-launch.test.ts && tsx src/workflow-replay.test.ts && tsx src/workflow-schema.test.ts", "typecheck": "tsc -p tsconfig.json --noEmit" }, "keywords": [], diff --git a/src/cli-output.test.ts b/src/cli-output.test.ts index 721afb51..b04a2a55 100644 --- a/src/cli-output.test.ts +++ b/src/cli-output.test.ts @@ -87,6 +87,7 @@ const call: WorkflowAgentCallRecord = { startedAt: now, completedAt: now, updatedAt: now, + finalUsage: { inputTokens: 8, outputTokens: 2, totalTokens: 10 }, }; const runJson = workflowRunOutput(run, [call]); assert.deepEqual(runJson.result, { ok: true }); @@ -98,11 +99,13 @@ assert.deepEqual(runJson.calls, { total: 1, }); assert.equal("scriptHash" in runJson, false); +assert.deepEqual(runJson.usage, { inputTokens: 8, outputTokens: 2, totalTokens: 10 }); const callJson = workflowCallOutput(call, { detailed: true }); assert.deepEqual(callJson.structured, { bugs: [] }); assert.equal("cacheKey" in callJson, false); assert.equal("providerSessionId" in callJson, false); assert.equal("profileFingerprint" in callJson, false); +assert.deepEqual(callJson.usage, { inputTokens: 8, outputTokens: 2, totalTokens: 10 }); console.log("cli-output.test.ts: ok"); diff --git a/src/cli-output.ts b/src/cli-output.ts index 667771b7..d9b3686a 100644 --- a/src/cli-output.ts +++ b/src/cli-output.ts @@ -56,6 +56,7 @@ export function workflowRunOutput( resumedFromRunId: run.resumedFromRunId, cancelRequested: run.cancelRequested, calls: calls ? workflowCallCounts(calls) : undefined, + usage: calls ? sumUsage(calls.map((call) => call.finalUsage ?? call.usage)) : undefined, result: parseStoredJson(run.resultJson), error: run.error ? { kind: run.errorKind, message: parseStoredJson(run.error) } @@ -82,6 +83,7 @@ export function workflowCallOutput( effort: call.effort, cached: call.fromCache, durationMs: workflowCallDurationMs(call), + usage: call.finalUsage ?? call.usage, isolation: call.isolation, worktree: call.worktreePath ? { path: call.worktreePath, dirty: call.dirty } @@ -130,6 +132,17 @@ function workflowCallDurationMs(call: WorkflowAgentCallRecord): number | undefin return Math.max(0, Date.parse(call.completedAt) - Date.parse(call.startedAt)); } +function sumUsage(calls: Array): WorkflowAgentCallRecord["usage"] | undefined { + const result: NonNullable = {}; + for (const key of ["inputTokens", "outputTokens", "totalTokens", "cacheReadTokens", "cacheWriteTokens"] as const) { + const values = calls + .map((usage) => usage?.[key]) + .filter((value): value is number => value !== undefined); + if (values.length > 0) result[key] = values.reduce((total, value) => total + value, 0); + } + return Object.keys(result).length > 0 ? result : undefined; +} + function parseStoredJson(value: string | undefined): unknown { if (value === undefined) return undefined; try { diff --git a/src/local-agent-adapters.ts b/src/local-agent-adapters.ts index ce91c90f..45627fbb 100644 --- a/src/local-agent-adapters.ts +++ b/src/local-agent-adapters.ts @@ -15,11 +15,24 @@ import type { LocalAgentProvider } from "./local-agent-profiles.js"; import { removeDevspaceNodeModulesBinFromPath } from "./local-agent-path.js"; import { createCodexSdkLocalAgentRuntime, + createLocalAgentObservationEmitter, isNativeSchemaUnsupportedFailure, ProviderSchemaUnsupportedError, type LocalAgentRunInput, type LocalAgentRunResult, } from "./local-agent-runtime.js"; +import { + extractAcpObservations, + extractAcpUsage, + extractClaudeObservations, + extractClaudeUsage, + extractCodexObservations, + extractCodexUsage, + extractOpenCodeObservations, + extractOpenCodeUsage, + extractPiObservations, + extractPiUsage, +} from "./local-agent-provider-observations.js"; export interface LocalAgentAdapter { readonly provider: LocalAgentProvider; @@ -103,9 +116,13 @@ class ClaudeLocalAgentAdapter implements LocalAgentAdapter { let providerSessionId = input.providerSessionId ?? null; let finalResponse = ""; let structured: unknown | undefined; + let usage: LocalAgentRunResult["usage"]; + const emitObservation = createLocalAgentObservationEmitter(input); const items: unknown[] = []; for await (const message of messages) { items.push(message); + for (const observation of extractClaudeObservations(message)) emitObservation(observation); + usage = extractClaudeUsage(message) ?? usage; const record = message as Record; if (typeof record.session_id === "string") providerSessionId = record.session_id; if (record.type !== "result") continue; @@ -124,6 +141,7 @@ class ClaudeLocalAgentAdapter implements LocalAgentAdapter { providerSessionId, finalResponse, items, + usage, ...(structured !== undefined ? { structured } : {}), }; } catch (error) { @@ -226,8 +244,13 @@ class OpencodeLocalAgentAdapter implements LocalAgentAdapter { try { const sessionId = input.providerSessionId ?? await createOpencodeSession(client, input); const promptResult = await promptOpencodeSession(client, sessionId, input); + const emitObservation = createLocalAgentObservationEmitter(input); + for (const observation of extractOpenCodeObservations(promptResult)) emitObservation(observation); await waitForOpencodeSession(client, sessionId); const messages = await readOpencodeMessages(client, sessionId); + for (const observation of extractOpenCodeObservations(messages)) emitObservation(observation); + const usage = extractOpenCodeUsage(messages) ?? extractOpenCodeUsage(promptResult); + if (usage) emitObservation({ kind: "usage", usage }); const finalResponse = requireFinalResponse( "OpenCode", extractOpenCodeFinalResponse(messages) || extractOpenCodeFinalResponse(promptResult), @@ -237,6 +260,7 @@ class OpencodeLocalAgentAdapter implements LocalAgentAdapter { providerSessionId: sessionId, finalResponse, items: [promptResult, messages], + usage, }; } finally { server.close(); @@ -263,6 +287,8 @@ class AcpLocalAgentAdapter implements LocalAgentAdapter { }); assertPipedChild(child); let stderr = ""; + const emitObservation = createLocalAgentObservationEmitter(input); + let usage: LocalAgentRunResult["usage"]; child.stderr.on("data", (chunk: Buffer) => { stderr += chunk.toString("utf8"); }); @@ -302,6 +328,8 @@ class AcpLocalAgentAdapter implements LocalAgentAdapter { } const update = message.update; + for (const observation of extractAcpObservations(update)) emitObservation(observation); + usage = extractAcpUsage(update) ?? usage; if (update.sessionUpdate !== "agent_message_chunk") continue; const content = update.content; if (content.type === "text") textParts.push(content.text); @@ -315,6 +343,7 @@ class AcpLocalAgentAdapter implements LocalAgentAdapter { providerSessionId, finalResponse: finalResponse.trim(), items: [], + usage, }; } catch (error) { throw new Error(`${this.provider} ACP run failed: ${errorMessage(error)}${stderr ? `\n${stderr.trim()}` : ""}`); @@ -429,7 +458,13 @@ class PiRpcLocalAgentAdapter implements LocalAgentAdapter { assertPipedChild(child); const rpc = new JsonLineRpc(child); const events: unknown[] = []; - rpc.onEvent((event) => events.push(event)); + const emitObservation = createLocalAgentObservationEmitter(input); + let usage: LocalAgentRunResult["usage"]; + rpc.onEvent((event) => { + events.push(event); + for (const observation of extractPiObservations(event)) emitObservation(observation); + usage = extractPiUsage(event) ?? usage; + }); try { const state = await rpc.request({ type: "get_state" }); const providerSessionId = readNestedString(state, ["sessionId"]) ?? input.providerSessionId ?? null; @@ -437,6 +472,10 @@ class PiRpcLocalAgentAdapter implements LocalAgentAdapter { await rpc.request({ type: "prompt", message: input.prompt }); const agentEnd = await done; const sessionMessages = await rpc.request({ type: "get_messages" }); + for (const observation of extractPiObservations(agentEnd)) emitObservation(observation); + for (const observation of extractPiObservations(sessionMessages)) emitObservation(observation); + usage = extractPiUsage(agentEnd) ?? extractPiUsage(sessionMessages) ?? usage; + if (usage) emitObservation({ kind: "usage", usage }); const finalResponse = extractPiFinalResponse(agentEnd) || extractPiFinalResponse(sessionMessages) || @@ -454,6 +493,7 @@ class PiRpcLocalAgentAdapter implements LocalAgentAdapter { providerSessionId, finalResponse, items: [...events, sessionMessages], + usage, }; } finally { child.kill(); diff --git a/src/local-agent-runtime.ts b/src/local-agent-runtime.ts index 3641762d..8032fb28 100644 --- a/src/local-agent-runtime.ts +++ b/src/local-agent-runtime.ts @@ -13,6 +13,10 @@ import type { LocalAgentObservation, LocalAgentTokenUsage, } from "./local-agent-observations.js"; +import { + extractCodexObservations, + extractCodexUsage, +} from "./local-agent-provider-observations.js"; import { isNativeSchemaUnsupportedFailure, ProviderSchemaUnsupportedError, @@ -67,6 +71,20 @@ export function notifyLocalAgentObservation( } } +export function createLocalAgentObservationEmitter( + input: LocalAgentRunInput, +): (observation: LocalAgentObservation) => void { + const emitted = new Set(); + return (observation) => { + if (observation.kind === "activity" && observation.activityId) { + const key = `${observation.activityId}:${observation.toolStatus ?? "updated"}:${observation.message ?? ""}:${observation.detail ?? ""}`; + if (emitted.has(key)) return; + emitted.add(key); + } + notifyLocalAgentObservation(input, observation); + }; +} + interface CodexThreadLike { readonly id: string | null; run(prompt: string, turnOptions?: TurnOptions): Promise; @@ -125,11 +143,19 @@ export class CodexSdkLocalAgentRuntime implements LocalAgentRuntime { throw error; } + const emitObservation = createLocalAgentObservationEmitter(input); + for (const item of turn.items) { + for (const observation of extractCodexObservations(item)) emitObservation(observation); + } + const usage = extractCodexUsage(turn); + if (usage) emitObservation({ kind: "usage", usage }); + return { provider: this.provider, providerSessionId: thread.id, finalResponse: turn.finalResponse, items: turn.items, + usage, ...(input.schema ? { structured: tryParseJson(turn.finalResponse) } : {}), }; }