diff --git a/docs/providers/openai/runtimes.md b/docs/providers/openai/runtimes.md index 3a2d9b974bbb..7fc728600a39 100644 --- a/docs/providers/openai/runtimes.md +++ b/docs/providers/openai/runtimes.md @@ -89,6 +89,9 @@ The selected effort applies to new sessions and later turns in existing sessions. `adaptive` and omitted native efforts use the model default; updating an existing session resets its effort to that default. The single-agent MVP does not support `ultra` delegation. Automatic runtime selection is unchanged. +When reasoning display is enabled, new sessions request native summaries. +Summary generation is fixed at creation; enabling it for an existing session +requires a reset. Native delegation remains disabled. ```json5 { @@ -115,7 +118,8 @@ A restricted API key needs Agents and Responses read/write plus Models read permission so the service can retrieve the selected model when creating a session. Agents API owns the persistent agent session and workspace. OpenClaw stores -the session binding in plugin SQLite state and mirrors text replies into its +the session binding in plugin SQLite state and mirrors assistant commentary, +reasoning summaries, native tool calls and results, and final text into its normal transcript. Follow-up messages reuse the agent session; input during a running turn steers it, and interruption cancels its remote turn. `/new` and `/reset` start a fresh session on the next message. Reset and local session @@ -123,9 +127,44 @@ deletion retire the binding; the Agents API retains the remote history and workspace, which can be managed through its API. If the event stream closes, the harness subscribes again and reconciles saved -turns and input receipts before accepting completion. It does not resubmit the -user's message. Native token usage is best effort; unavailable usage currently -appears as zero in OpenClaw's usage totals. +turns, saved items, and input receipts before accepting completion. It does not +resubmit the user's message. Completion requires a terminal root turn and an +idle native session. Recovered items use stable identities to avoid duplicate +history. If a recovered item is still running, its saved snapshot is displayed +and ambiguous overlapping deltas are suppressed until authoritative completion; +new items continue streaming normally. + +When a retry retains the native conversation, later canonical facts can also +repair missing tool records from earlier terminal turns. These records append +to the existing history without replaying progress or counting earlier work +in the current attempt. Retrieved command invocation facts can be saved while +execution is still running; a durable result requires the item's own terminal +status. MCP arguments and web-search actions are saved after their items finish. + +Commentary remains live while durable commentary and reasoning records wait for +retrieved native history. Missing or unfinished items can prevent exact ordering; +terminal reconciliation retains available completed records instead of dropping +them behind an unresolved item. Host input keeps its existing transcript +placement, and historical repairs append without rewriting earlier messages. + +Native token usage is best effort and is accumulated across all admitted turns, +including work superseded by a steering follow-up. Cached input and reasoning +tokens remain separate usage facts. Billed tokens do not establish active +context occupancy; that value remains unavailable. +Completed, failed, and cancelled turn events can supply usage even when the +saved turn has none; the harness retains that contribution and reports it once. + +Native commands, MCP calls, and web searches use the same activity and output +callbacks as the Codex harness. Assistant commentary preserves its text, so a +progress marker can reach the channel while a command is still running. +Existing channel settings govern output and reasoning visibility. + +Saved native tool items supply canonical history and available command output, +exit codes, duration, arguments, MCP details, and web-search actions. Missing +command exit facts remain unknown even when the assistant claims success. +The native API exposes web-search activity without result bodies or snippets. +Structured plans, diffs, compaction events, native child agents, and pre-execution +approval or hook events are not provided by this harness. New Agents API sessions enable built-in web search in live mode. Sessions created before web search was enabled need `/new` or `/reset` to pick it up. diff --git a/extensions/agentsapi/agentsapi-client.ts b/extensions/agentsapi/agentsapi-client.ts index 6569673549a8..eeb057438943 100644 --- a/extensions/agentsapi/agentsapi-client.ts +++ b/extensions/agentsapi/agentsapi-client.ts @@ -1,15 +1,123 @@ import { randomUUID } from "node:crypto"; import { setTimeout as delay } from "node:timers/promises"; import OpenAI from "openai"; -import type { - AgentReasoningParam, - AgentSessionEvent, - AgentSessionItem, -} from "openai/resources/beta/agents/agents"; +import type { AgentReasoningParam, AgentSessionEvent } from "openai/resources/beta/agents/agents"; import type { Turn } from "openai/resources/beta/agents/sessions/turns"; import { responseWithRelease } from "openclaw/plugin-sdk/fetch-runtime"; import { fetchWithSsrFGuard } from "openclaw/plugin-sdk/ssrf-runtime"; +import { z } from "zod"; +const usageSchema = z.looseObject({ + input_tokens: z.number(), + output_tokens: z.number(), + total_tokens: z.number().optional(), + input_tokens_details: z.looseObject({ cached_tokens: z.number() }).optional(), + output_tokens_details: z.looseObject({ reasoning_tokens: z.number() }).optional(), +}); +const errorSchema = z.looseObject({ + message: z.string(), + code: z.string().nullable().optional(), + type: z.string().optional(), + param: z.string().nullable().optional(), +}); +const functionCallSchema = z.looseObject({ + type: z.literal("function_call"), + turn_id: z.string().min(1), + call_id: z.string().min(1), + name: z.string().min(1), + arguments: z.unknown(), +}); +const sessionSchema = z.looseObject({ + id: z.string(), + status: z.enum(["idle", "in_progress", "requires_action", "failed"]), + error: z.string().nullable(), + usage: usageSchema.nullable().optional(), + environment: z.looseObject({ type: z.literal("openai_hosted"), id: z.string().min(1) }), + required_actions: z.array( + z.union([ + functionCallSchema, + z.looseObject({ type: z.literal("environment_connection"), environment_id: z.string() }), + ]), + ), +}); +const textPartSchema = z.looseObject({ type: z.string(), text: z.string().optional() }); +// Validate native correlation and projection fields while retaining complete payloads. +const itemSchema = z.looseObject({ + id: z.string(), + type: z.string(), + role: z.string().optional(), + phase: z.string().nullable().optional(), + status: z.string().nullable().optional(), + turn_id: z.string().optional(), + content: z.array(textPartSchema).optional(), + summary: z.array(textPartSchema).optional(), + command: z.string().optional(), + cwd: z.string().nullable().optional(), + duration_ms: z.number().nullable().optional(), + exit_code: z.number().nullable().optional(), + name: z.string().optional(), + call_id: z.string().optional(), + server_label: z.string().optional(), + arguments: z.unknown().optional(), + output: z.unknown().optional(), + error: z.unknown().optional(), + action: z + .looseObject({ + type: z.string(), + query: z.string().nullable().optional(), + queries: z.array(z.string()).nullable().optional(), + url: z.string().nullable().optional(), + pattern: z.string().nullable().optional(), + }) + .nullable() + .optional(), +}); +const eventSchema = z.looseObject({ + type: z.string(), + event_id: z.string().optional(), + session_id: z.string().optional(), + turn_id: z.string().nullable().optional(), + item_id: z.string().optional(), + output_index: z.number().nullable().optional(), + content_index: z.number().optional(), + summary_index: z.number().optional(), + status: z.string().nullable().optional(), + delta: z.string().optional(), + text: z.string().optional(), + part: textPartSchema.optional(), + item: itemSchema.optional(), + usage: usageSchema.nullable().optional(), + session: sessionSchema.optional(), + environment: z + .looseObject({ + id: z.string(), + type: z.string(), + status: z.string(), + error: errorSchema.nullable(), + }) + .optional(), + turn: z + .looseObject({ + id: z.string(), + subagent_id: z.string().nullable(), + session_id: z.string().optional(), + status: z.string().optional(), + agent_id: z.string().optional(), + created_at: z.number().optional(), + started_at: z.number().nullable().optional(), + completed_at: z.number().nullable().optional(), + error: errorSchema.nullable().optional(), + usage: usageSchema.nullable().optional(), + }) + .optional(), + error: errorSchema.optional(), +}); +export type AgentsApiEvent = z.infer; +export type AgentsApiItem = z.infer; +export type AgentsApiReasoning = { + effort?: "none" | "minimal" | "low" | "medium" | "high" | "xhigh" | "max" | null; + summary?: "concise" | "detailed" | "auto" | null; +}; /** The SDK owns the wire protocol; OpenClaw retains native session authority. */ export class AgentsApiClient { private readonly sessions: OpenAI["beta"]["agents"]["sessions"]; @@ -54,13 +162,20 @@ export class AgentsApiClient { instructions: string, model: string, reasoningEffort?: AgentReasoningParam["effort"], + extras?: { + reasoning?: AgentsApiReasoning; + }, ): Promise { const session = await this.sessions.create( { agent: { model, instructions, - reasoning: reasoningEffort === undefined ? undefined : { effort: reasoningEffort }, + reasoning: extras?.reasoning + ? { ...extras.reasoning, effort: reasoningEffort } + : reasoningEffort === undefined + ? undefined + : { effort: reasoningEffort }, multi_agent: { enabled: false }, tools: [{ type: "web_search", mode: "live" }], }, @@ -102,7 +217,7 @@ export class AgentsApiClient { stream.controller.abort(); throw error; } - return observeEvents(stream, signal, this.assertCurrent); + return observeEvents(stream, signal, sessionId, this.assertCurrent); } async session(sessionId: string, signal: AbortSignal) { @@ -178,12 +293,20 @@ export class AgentsApiClient { } } - async items(sessionId: string, turnId: string, signal: AbortSignal): Promise { - const items: AgentSessionItem[] = []; + async items( + sessionId: string, + turnId: string | undefined, + signal: AbortSignal, + ): Promise { + const items: AgentsApiItem[] = []; const pages = this.sessions.items.list(sessionId, { order: "asc", limit: 100 }, { signal }); for await (const page of (await pages).iterPages()) { this.assertCurrent(); - items.push(...page.data.filter((item) => item.turn_id === turnId)); + items.push( + ...page.data + .filter((item) => turnId === undefined || item.turn_id === turnId) + .map((item) => itemSchema.parse(item)), + ); if (page.has_more && !page.hasNextPage()) { throw new Error("Agents API items page has no continuation cursor"); } @@ -192,14 +315,48 @@ export class AgentsApiClient { } } +/** Customer-safe native failure facts remain available to host result classification. */ +export class AgentsApiError extends Error { + readonly code: string | null | undefined; + readonly status: number | undefined; + readonly type: string | undefined; + readonly param: string | null | undefined; + + constructor( + message: string, + details: { + code?: string | null; + status?: number; + type?: string; + param?: string | null; + } = {}, + ) { + super(message); + this.name = "AgentsApiError"; + this.code = details.code; + this.status = details.status; + this.type = details.type; + this.param = details.param; + } +} + async function* observeEvents( stream: AsyncIterable, signal: AbortSignal, + sessionId: string, assertCurrent: () => void, -) { - for await (const event of stream) { +): AsyncGenerator { + for await (const rawEvent of stream) { signal.throwIfAborted(); assertCurrent(); + const event = eventSchema.parse(rawEvent); + if ( + (event.session_id && event.session_id !== sessionId) || + (event.session && event.session.id !== sessionId) || + (event.turn?.session_id && event.turn.session_id !== sessionId) + ) { + throw new Error("Agents API returned an event outside the requested session"); + } yield event; } signal.throwIfAborted(); diff --git a/extensions/agentsapi/agentsapi-harness.ts b/extensions/agentsapi/agentsapi-harness.ts index 991d7fb9539f..8cb1619aebe2 100644 --- a/extensions/agentsapi/agentsapi-harness.ts +++ b/extensions/agentsapi/agentsapi-harness.ts @@ -7,6 +7,8 @@ import { createAgentHarnessAttemptLifecycle, emitAgentHarnessAttemptEvent, selectSupportedReasoningEffort, + AgentHarnessProjectionSettlement, + racePromiseWithAbortSignal, type AgentHarnessAttemptTimeout, } from "openclaw/plugin-sdk/agent-harness-attempt-runtime"; import { @@ -39,6 +41,7 @@ import { createAgentsApiBindings } from "./agentsapi-bindings.js"; import { AgentsApiClient } from "./agentsapi-client.js"; import { createAgentsApiMessageProjection } from "./agentsapi-messages.js"; import { createAgentsApiSession } from "./agentsapi-session.js"; +import { recordAgentsApiNativeToolTranscript } from "./agentsapi-transcript.js"; /** Agents API owns native protocol; the host harness runtime owns coordination. */ export function createAgentsApiHarness(runtime: PluginRuntime): AgentHarnessV2 { @@ -184,13 +187,46 @@ async function runAgentsApiSession( assertOwnerCurrent(); controller.signal.throwIfAborted(); }; + let finalizingProjection = false; + let finalizingProjectionSignal: AbortSignal | undefined; + const assertProjectionCurrent = () => { + assertOwnerCurrent(); + if (finalizingProjection) { + finalizingProjectionSignal?.throwIfAborted(); + } else { + controller.signal.throwIfAborted(); + } + }; + let lastToolError: AgentHarnessAttemptResult["lastToolError"]; + let toolTerminalObserved = false; + const observeToolTerminal = params.observeToolTerminal; + const runParams: AgentHarnessAttemptParamsV2 = observeToolTerminal + ? { + ...params, + observeToolTerminal: (observation) => { + assertProjectionCurrent(); + const resolution = observeToolTerminal(observation); + assertProjectionCurrent(); + toolTerminalObserved = true; + lastToolError = resolution.lastToolError; + return resolution; + }, + } + : params; let timeout: AgentHarnessAttemptTimeout | undefined; + let settling = false; + let settlementDeadlineAtMs: number | undefined; const deadlines = createAgentHarnessAttemptDeadlineController({ startedAtMs, timeoutMs: params.timeoutMs, settlementTimeoutMs: 30_000, signal: controller.signal, - onDeadlineChanged: params.onAttemptDeadlineChanged, + onDeadlineChanged: (deadline) => { + if (settling && deadline.kind === "bounded") { + settlementDeadlineAtMs = deadline.deadlineAtMs; + } + params.onAttemptDeadlineChanged?.(deadline); + }, onTimeout: (expired) => { timeout = expired; const error = new Error(`Agents API ${expired.kind} timed out`); @@ -198,6 +234,10 @@ async function runAgentsApiSession( cancellation.abortExplicitly(error); }, }); + const beginSettlement = () => { + settling = true; + deadlines.beginSettlement(Date.now()); + }; const emitEvent = ( event: Parameters>[0], ) => emitAgentHarnessAttemptEvent(params, event, { label: "Agents API", log: embeddedAgentLog }); @@ -214,6 +254,22 @@ async function runAgentsApiSession( let reply: ReturnType["reply"] | undefined; let projection: ReturnType | undefined; let usageRecorded = false; + let projectionClosed = false; + const projectionSettlement = new AgentHarnessProjectionSettlement( + runParams, + () => { + if (projectionClosed || controller.signal.aborted) { + return false; + } + try { + assertOwnerCurrent(); + return true; + } catch { + return false; + } + }, + { label: "Agents API" }, + ); let terminalTurnId: string | undefined; const handle = { kind: "embedded", @@ -277,6 +333,14 @@ async function runAgentsApiSession( .join("\n\n"), params.model.id, reasoningEffort, + { + reasoning: { + effort: reasoningEffort, + ...(params.reasoningLevel && params.reasoningLevel !== "off" + ? { summary: "auto" } + : {}), + }, + }, ); assertCurrent(); await bind({ sessionId: remoteSessionId, authFingerprint: fingerprint }); @@ -284,10 +348,16 @@ async function runAgentsApiSession( await client.setReasoningEffort(remoteSessionId, reasoningEffort, controller.signal); assertCurrent(); } - projection = createAgentsApiMessageProjection(remoteSessionId, (event) => { - void emitEvent(event); - }); - const messageProjection = projection; + projection = createAgentsApiMessageProjection( + projectionSettlement.params, + remoteSessionId, + async (event) => { + assertCurrent(); + await emitEvent(event); + assertCurrent(); + }, + assertProjectionCurrent, + ); reply = projection.reply; native = createAgentsApiSession({ client, @@ -296,11 +366,32 @@ async function runAgentsApiSession( sessionId: remoteSessionId, signal: controller.signal, assertCurrent, - onSettled: () => deadlines.beginSettlement(Date.now()), + onSettled: beginSettlement, + onReconcile: (turn, items) => + projection!.reconcile(turn, items, { presentation: !finalizingProjection }), onUsageError: (error) => embeddedAgentLog.warn("Agents API token accounting unavailable", { error }), - onEvent: (event) => { - messageProjection.observe(event); + onTranscriptOrderingGap: () => projection!.reportTranscriptOrderingGap(), + onReconcileHistory: async (entries) => { + for (const { turn, items } of entries) { + for (const item of items) { + assertProjectionCurrent(); + await recordAgentsApiNativeToolTranscript( + runParams, + remoteSessionId!, + turn.id, + item, + assertProjectionCurrent, + Date.now, + { enclosingStatus: turn.status }, + ); + assertProjectionCurrent(); + } + } + }, + onEvent: async (event) => { + await projection!.observe(event); + assertCurrent(); params.onRunProgress?.({ reason: event.type, provider: "openai", @@ -329,7 +420,7 @@ async function runAgentsApiSession( } else { const items = await client.items(remoteSessionId, result.turn.id, controller.signal); assertCurrent(); - await projection.commit(params, result.turn, items, assertCurrent); + await projection.commit(result.turn, items); assertCurrent(); } } catch (error) { @@ -344,7 +435,20 @@ async function runAgentsApiSession( embeddedAgentLog.warn("Agents API session failed", { error }); } } finally { + beginSettlement(); + // Reuse the owner's absolute settlement boundary. After an upstream abort + // closes that owner, one cleanup budget starts before native retirement. + const cleanupMs = Math.max( + 0, + Math.min(30_000, (settlementDeadlineAtMs ?? Date.now() + 30_000) - Date.now()), + ); + const cleanupSignal = + cleanupMs > 0 + ? AbortSignal.timeout(cleanupMs) + : AbortSignal.abort(new Error("Agents API settlement timed out")); try { + // Retirement retains the native binding lease until admitted POST/cancel + // work settles under its API timeouts; early release could cancel a successor. await native?.close(); } catch (error) { terminal = { kind: "failed", source: "prompt", error }; @@ -359,6 +463,40 @@ async function runAgentsApiSession( } catch (error) { terminal = { kind: "failed", source: "prompt", error }; } + if ((controller.signal.aborted || terminal.kind !== "ok") && native && projection) { + let ownerCurrent = false; + try { + assertOwnerCurrent(); + ownerCurrent = true; + } catch { + // Retired authority cannot publish evidence into a successor session. + } + if (ownerCurrent) { + finalizingProjection = true; + finalizingProjectionSignal = cleanupSignal; + try { + await native.reconcileAfterClose(cleanupSignal); + } catch (error) { + embeddedAgentLog.warn("Agents API terminal history reconciliation failed", { error }); + } finally { + finalizingProjection = false; + finalizingProjectionSignal = undefined; + } + } + } + try { + await racePromiseWithAbortSignal(projectionSettlement.drain(), cleanupSignal); + } catch (error) { + if (!controller.signal.aborted) { + terminal = { kind: "failed", source: "prompt", error }; + } + } + projectionClosed = true; + if (timeout) { + terminal = { kind: "timeout", phase: "prompt", source: "runtime", aborted: true }; + } else if (terminal.kind === "ok" && controller.signal.aborted) { + terminal = { kind: "aborted", source: params.abortSignal?.aborted ? "external" : "runtime" }; + } cancellation.freezeTerminalOutcome(); deadlines.dispose(); cancellation.dispose(); @@ -385,18 +523,24 @@ async function runAgentsApiSession( reply?.lastAssistant && terminalTurnId ? `agentsapi:${remoteSessionId}:${terminalTurnId}` : undefined, - toolMetas: [], + toolMetas: projection?.toolMetas ?? [], + lastToolError: toolTerminalObserved ? lastToolError : projection?.lastToolError, didSendViaMessagingTool: false, messagingToolSentTexts: [], messagingToolSentMediaUrls: [], messagingToolSentTargets: [], cloudCodeAssistFormatError: false, - attemptUsage: reply?.usage, + attemptUsage: projection?.tokenUsage, + agentHarnessResultClassification: projection?.resultClassification, replayMetadata: { hadPotentialSideEffects: native?.wasSubmitted() ?? false, replaySafe: !native?.wasSubmitted(), }, - itemLifecycle: { startedCount: 0, completedCount: 0, activeCount: 0 }, + itemLifecycle: projection?.itemLifecycle ?? { + startedCount: 0, + completedCount: 0, + activeCount: 0, + }, }; assertHarnessCurrent(); const contextWindow = { diff --git a/extensions/agentsapi/agentsapi-messages.test.ts b/extensions/agentsapi/agentsapi-messages.test.ts new file mode 100644 index 000000000000..2f85353baa5c --- /dev/null +++ b/extensions/agentsapi/agentsapi-messages.test.ts @@ -0,0 +1,199 @@ +import type { Turn as SDKTurn } from "openai/resources/beta/agents/sessions/turns"; +import type { AgentHarnessAttemptParamsV2 } from "openclaw/plugin-sdk/agent-harness-runtime"; +import { describe, expect, it } from "vitest"; +import type { AgentsApiEvent, AgentsApiItem } from "./agentsapi-client.js"; +import { createAgentsApiMessageProjection } from "./agentsapi-messages.js"; + +type AgentEvent = Parameters>[0]; + +describe("Agents API commentary projection", () => { + it("completes commentary after an identical final delta without replaying completion", async () => { + const { projection, events } = createProjection(); + const item: AgentsApiItem = { + id: "commentary-fixture", + type: "message", + role: "assistant", + phase: "commentary", + status: "in_progress", + turn_id: "turn-fixture", + content: [{ type: "output_text", text: "" }], + }; + const text = "Checking the command."; + + await projection.observe({ type: "agent.session.turn.item.added", item }); + await projection.observe({ + type: "agent.session.turn.output_text.delta", + item_id: item.id, + turn_id: "turn-fixture", + content_index: 0, + delta: text, + }); + const completed: AgentsApiEvent = { + type: "agent.session.turn.item.done", + item: { + ...item, + status: "completed", + content: [{ type: "output_text", text }], + }, + }; + await projection.observe(completed); + await projection.observe(completed); + + expect(events).toEqual([ + { + stream: "item", + data: { + itemId: "agentsapi:session-fixture:turn-fixture:commentary-fixture", + kind: "preamble", + title: "Preamble", + phase: "update", + progressText: text, + source: "agentsapi", + }, + }, + { + stream: "item", + data: { + itemId: "agentsapi:session-fixture:turn-fixture:commentary-fixture", + kind: "preamble", + title: "Preamble", + phase: "end", + progressText: text, + source: "agentsapi", + }, + }, + ]); + }); +}); + +describe("Agents API final usage accounting", () => { + it("retains completed-turn usage when final REST accounting returns no turns", async () => { + const { projection } = createProjection(); + await projection.observe({ + type: "agent.session.turn.completed", + turn: createTurn("turn-a", observedUsageA), + }); + await projection.observe({ + type: "agent.session.turn.completed", + turn: createTurn("turn-b", observedUsageB), + }); + + projection.recordUsage(usageModel, []); + + expect(projection.tokenUsage).toMatchObject({ + input: 120, + output: 8, + cacheRead: 30, + reasoningTokens: 3, + total: 158, + contextUsage: { state: "unavailable" }, + }); + expect(projection.reply.assistantUsage).toMatchObject({ + input: 120, + output: 8, + cacheRead: 30, + totalTokens: 158, + }); + }); + + it("replaces matching canonical usage while retaining omitted turns without counting them twice", async () => { + const { projection } = createProjection(); + await projection.observe({ + type: "agent.session.turn.completed", + turn: createTurn("turn-a", observedUsageA), + }); + await projection.observe({ + type: "agent.session.turn.completed", + turn: createTurn("turn-b", observedUsageB), + }); + const canonicalTurn = createTurn("turn-a", { + input_tokens: 120, + input_tokens_details: { cached_tokens: 30 }, + output_tokens: 6, + output_tokens_details: { reasoning_tokens: 2 }, + total_tokens: 126, + }); + + projection.recordUsage(usageModel, [canonicalTurn, canonicalTurn]); + + const expected = { + input: 130, + output: 9, + cacheRead: 40, + reasoningTokens: 4, + total: 179, + contextUsage: { state: "unavailable" }, + }; + expect(projection.tokenUsage).toMatchObject(expected); + expect(projection.reply.assistantUsage).toMatchObject({ + input: 130, + output: 9, + cacheRead: 40, + totalTokens: 179, + }); + + projection.recordUsage(usageModel, [canonicalTurn, canonicalTurn]); + expect(projection.tokenUsage).toMatchObject(expected); + expect(projection.reply.assistantUsage).toMatchObject({ totalTokens: 179 }); + }); +}); + +function createProjection() { + const events: AgentEvent[] = []; + // Observation needs no auth or transcript operations; final accounting receives the model. + const params = {} as AgentHarnessAttemptParamsV2; + const projection = createAgentsApiMessageProjection( + params, + "session-fixture", + (event) => { + events.push(event); + }, + () => {}, + ); + return { projection, events }; +} + +function createTurn(id: string, usage: typeof observedUsageA) { + return { + id, + agent_id: "agent-fixture", + session_id: "session-fixture", + object: "agent.session.turn", + created_at: 1, + started_at: 1, + completed_at: 2, + status: "completed", + subagent_id: null, + error: null, + usage, + } satisfies SDKTurn; +} + +const usageModel = { + id: "model-fixture", + name: "Fixture Model", + api: "openai-responses", + provider: "openai", + baseUrl: "https://api.openai.com/v1", + reasoning: true, + input: ["text"], + cost: { input: 0, output: 0, cacheRead: 0, cacheWrite: 0 }, + contextWindow: 1024, + maxTokens: 512, +} satisfies AgentHarnessAttemptParamsV2["model"]; + +const observedUsageA = { + input_tokens: 100, + input_tokens_details: { cached_tokens: 20 }, + output_tokens: 5, + output_tokens_details: { reasoning_tokens: 1 }, + total_tokens: 105, +} satisfies NonNullable; + +const observedUsageB = { + input_tokens: 50, + input_tokens_details: { cached_tokens: 10 }, + output_tokens: 3, + output_tokens_details: { reasoning_tokens: 2 }, + total_tokens: 53, +} satisfies NonNullable; diff --git a/extensions/agentsapi/agentsapi-messages.ts b/extensions/agentsapi/agentsapi-messages.ts index 69d517f05675..6a092655a198 100644 --- a/extensions/agentsapi/agentsapi-messages.ts +++ b/extensions/agentsapi/agentsapi-messages.ts @@ -1,33 +1,559 @@ -import type { AgentSessionEvent, AgentSessionItem } from "openai/resources/beta/agents/agents"; -import type { Turn } from "openai/resources/beta/agents/sessions/turns"; +import type { Turn as SDKTurn } from "openai/resources/beta/agents/sessions/turns"; +import { createAgentHarnessAssistantMessage } from "openclaw/plugin-sdk/agent-harness-attempt-runtime"; import { + classifyAgentHarnessTerminalOutcome, + embeddedAgentLog, normalizeUsage, type AgentHarnessAttemptParamsV2, + type AgentHarnessAttemptResult, + type AgentMessage, type NormalizedUsage, } from "openclaw/plugin-sdk/agent-harness-runtime"; import { calculateCost, type AssistantMessage } from "openclaw/plugin-sdk/llm"; -import { appendSessionTranscriptMessageByIdentityStrict } from "openclaw/plugin-sdk/session-transcript-runtime"; +import type { AgentsApiEvent, AgentsApiItem } from "./agentsapi-client.js"; +import { AgentsApiNativeToolProjection } from "./agentsapi-native-tool-projection.js"; +import { + appendAgentsApiTranscriptMessage, + canRecordAgentsApiTranscriptText, + iterateAgentsApiTranscriptItems, + joinTextParts, + readTextParts, +} from "./agentsapi-transcript.js"; type AgentEvent = Parameters>[0]; +type NativeTurn = SDKTurn | NonNullable; type AgentsApiReply = { lastAssistant?: AssistantMessage; usage?: NormalizedUsage; assistantUsage: AssistantMessage["usage"]; }; +type NativeTextState = { + turnId: string; + item: AgentsApiItem; + terminal: boolean; + completionObserved: boolean; + texts: Map; + summaries: Map; + recoveredPartial: boolean; + lastCommentary?: { phase: "update" | "end"; text: string }; + lastAssistantText?: string; +}; -export function createAgentsApiMessageProjection( - remoteSessionId: string, - emitEvent: (event: AgentEvent) => void | Promise, -) { - const reply: AgentsApiReply = { assistantUsage: emptyUsage() }; - const texts = new Map>(); - const assistantPhases = new Map(); - let visibleAssistantItemId: string | undefined; - const emitAssistantSnapshot = (itemId: string, text: string, delta = "") => { - const replace = visibleAssistantItemId !== itemId; - visibleAssistantItemId = itemId; - // Steering can supersede a completed turn; append-only delivery waits for canonical items. - void emitEvent({ +/** Native identities keep saved-state recovery and live events on the same projection. */ +class AgentsApiMessageProjection { + readonly reply: AgentsApiReply = { assistantUsage: emptyUsage() }; + private readonly items = new Map(); + private readonly turnByItem = new Map(); + private readonly eventIds = new Set(); + private readonly unknownTypes = new Set(); + private readonly usageByTurn = new Map(); + private canonicalUsageRecorded = false; + private readonly nativeTools: AgentsApiNativeToolProjection; + private visibleAssistantItemId: string | undefined; + private reasoningOpen = false; + private lastReasoningText = ""; + private finalTurnId: string | undefined; + private classification: AgentHarnessAttemptResult["agentHarnessResultClassification"]; + private timestamp = Date.now(); + private presentationEnabled = true; + private transcriptOrderingGapReported = false; + private readonly recordedGatewayCallIds = new Set(); + + constructor( + private readonly params: AgentHarnessAttemptParamsV2, + private readonly remoteSessionId: string, + private readonly emitEvent: (event: AgentEvent) => void | Promise, + private readonly assertCurrent: () => void, + ) { + this.nativeTools = new AgentsApiNativeToolProjection( + params, + remoteSessionId, + (event) => this.emit(event), + assertCurrent, + () => this.nextTimestamp(), + () => this.presentationEnabled, + ); + } + + get toolMetas(): AgentHarnessAttemptResult["toolMetas"] { + return this.nativeTools.toolMetas; + } + + get lastToolError(): AgentHarnessAttemptResult["lastToolError"] { + return this.nativeTools.lastToolError; + } + + get itemLifecycle(): AgentHarnessAttemptResult["itemLifecycle"] { + const states = [...this.items.values()]; + const completedCount = states.filter((state) => state.terminal).length; + const tools = this.nativeTools.itemLifecycle; + return { + startedCount: states.length + tools.startedCount, + completedCount: completedCount + tools.completedCount, + activeCount: states.length - completedCount + tools.activeCount, + }; + } + + get hadPotentialSideEffects(): boolean { + return this.nativeTools.hadPotentialSideEffects; + } + + get resultClassification(): AgentHarnessAttemptResult["agentHarnessResultClassification"] { + return this.classification; + } + + recordUsage(model: AgentHarnessAttemptParamsV2["model"], turns: SDKTurn[]): void { + const usage = emptyUsage(); + let observed = false; + let reasoningTokens: number | undefined; + // Canonical usage replaces observed usage by admitted turn identity. A failed + // or partial REST read cannot discard terminal-event usage for omitted turns. + const contributions = new Map(this.usageByTurn); + for (const turn of new Map(turns.map((record) => [record.id, record])).values()) { + const normalized = normalizeUsage(turn.usage); + if (normalized) { + contributions.set(turn.id, normalized); + } + } + this.canonicalUsageRecorded = true; + for (const normalized of contributions.values()) { + observed = true; + usage.input += normalized.input ?? 0; + usage.output += normalized.output ?? 0; + usage.cacheRead += normalized.cacheRead ?? 0; + usage.cacheWrite += normalized.cacheWrite ?? 0; + usage.totalTokens += + normalized.total ?? + (normalized.input ?? 0) + + (normalized.output ?? 0) + + (normalized.cacheRead ?? 0) + + (normalized.cacheWrite ?? 0); + if (normalized.reasoningTokens !== undefined) { + reasoningTokens = (reasoningTokens ?? 0) + normalized.reasoningTokens; + } + } + if (observed) { + calculateCost(model, usage); + this.reply.usage = { + ...normalizeUsage(usage), + ...(reasoningTokens !== undefined ? { reasoningTokens } : {}), + }; + } else { + this.reply.usage = { contextUsage: { state: "unavailable" } }; + } + this.reply.assistantUsage = usage; + } + + get tokenUsage(): NormalizedUsage | undefined { + if (this.canonicalUsageRecorded) { + return this.reply.usage; + } + const usage: NormalizedUsage = { contextUsage: { state: "unavailable" } }; + for (const contribution of this.usageByTurn.values()) { + for (const bucket of [ + "input", + "output", + "cacheRead", + "cacheWrite", + "reasoningTokens", + "total", + ] as const) { + const value = contribution[bucket]; + if (value !== undefined) { + usage[bucket] = (usage[bucket] ?? 0) + value; + } + } + } + return usage; + } + + async observe(event: AgentsApiEvent): Promise { + this.assertCurrent(); + if (this.finalTurnId || (event.event_id && this.eventIds.has(event.event_id))) { + return; + } + if (event.event_id) { + this.eventIds.add(event.event_id); + if (this.eventIds.size > 4096) { + const oldest = this.eventIds.values().next().value; + if (oldest !== undefined) { + this.eventIds.delete(oldest); + } + } + } + if ( + event.turn?.subagent_id === null && + [ + "agent.session.turn.completed", + "agent.session.turn.failed", + "agent.session.turn.cancelled", + ].includes(event.type) + ) { + this.recordTurnUsage({ ...event.turn, usage: event.usage ?? event.turn.usage }); + } + if (event.item) { + const turnId = + event.item.turn_id ?? + event.turn_id ?? + this.turnByItem.get(event.item.id) ?? + this.nativeTools.resolveTurnId(event.item.id); + if (!turnId) { + throw new Error("Agents API output item has no turn identity"); + } + await this.recordItem(turnId, event.item, event.type === "agent.session.turn.item.done"); + return; + } + if (event.type === "agent.output.command_execution_output.delta") { + await this.nativeTools.observeOutput(event); + return; + } + if ( + event.type === "agent.session.turn.output_text.delta" || + event.type === "agent.session.turn.output_text.done" || + ((event.type === "agent.session.turn.content_part.added" || + event.type === "agent.session.turn.content_part.done") && + event.part?.type === "output_text") + ) { + const state = this.eventItem(event); + if (!state || state.terminal) { + return; + } + if (state.recoveredPartial && !event.type.endsWith(".done")) { + return; + } + const index = event.content_index ?? 0; + state.texts.set( + index, + event.type.endsWith(".delta") + ? (state.texts.get(index) ?? "") + (event.delta ?? "") + : (event.text ?? event.part?.text ?? ""), + ); + await this.emitAssistant(state, false, event.delta ?? ""); + return; + } + if (event.type.startsWith("agent.session.turn.reasoning_summary_")) { + const state = this.eventItem(event); + if (!state || state.terminal) { + return; + } + if (state.recoveredPartial && !event.type.endsWith(".done")) { + return; + } + const index = event.summary_index ?? 0; + if (event.type.endsWith("text.delta")) { + state.summaries.set(index, (state.summaries.get(index) ?? "") + (event.delta ?? "")); + } else if (event.type.endsWith("text.done") || event.type.endsWith("part.done")) { + state.summaries.set(index, event.text ?? event.part?.text ?? ""); + } else if (event.part?.text) { + state.summaries.set(index, event.part.text); + } + await this.emitReasoning(); + return; + } + if (event.type.startsWith("agent.session.environment.")) { + await this.emitEnvironment(event); + return; + } + if (!INTERNAL_EVENT_TYPES.has(event.type)) { + this.recordUnknownType(event.type); + } + } + + recordGatewayTranscriptReceipt(turnId: string, callId: string): void { + this.assertCurrent(); + this.recordedGatewayCallIds.add(`${turnId}:${callId}`); + } + + reportTranscriptOrderingGap(): void { + this.assertCurrent(); + if (this.transcriptOrderingGapReported) { + return; + } + this.transcriptOrderingGapReported = true; + embeddedAgentLog.warn( + "Agents API canonical transcript prefix is unavailable; host input and tool receipts retain their existing placement", + ); + } + + async reconcile( + turn: NativeTurn, + items: AgentsApiItem[], + options: { presentation?: boolean } = {}, + ): Promise { + this.assertCurrent(); + if (this.finalTurnId) { + return true; + } + const previousPresentation = this.presentationEnabled; + // Cancellation reconciliation runs only after the stream is retired. It + // records canonical facts under the original owner, without reopening output. + this.presentationEnabled = previousPresentation && options.presentation !== false; + try { + let transcriptReady = true; + const terminalTurn = isTerminalTurn(turn.status); + for (const projected of iterateAgentsApiTranscriptItems( + turn.id, + items, + turn.status, + terminalTurn, + this.recordedGatewayCallIds, + (itemId) => this.items.get(this.identity(turn.id, itemId))?.completionObserved === true, + )) { + const { item, terminal } = projected; + transcriptReady = projected.transcriptReady; + // Settlement preserves available records even when an earlier native + // item is unresolved. Their order is explicitly best effort in that case. + await this.recordItem( + turn.id, + item, + terminal, + turn.status, + true, + transcriptReady || terminalTurn, + ); + } + this.recordTurnUsage(turn); + if (isTerminalTurn(turn.status)) { + for (const state of this.items.values()) { + if (state.turnId !== turn.id || state.terminal) { + continue; + } + state.terminal = true; + } + const remainingReady = await this.nativeTools.reconcileRemaining( + turn.id, + turn.status, + new Set(items.map((item) => item.id)), + ); + transcriptReady = remainingReady && transcriptReady; + if (!transcriptReady) { + this.reportTranscriptOrderingGap(); + } + } + return transcriptReady; + } finally { + this.presentationEnabled = previousPresentation; + } + } + + async commit(turn: NativeTurn, items: AgentsApiItem[]): Promise { + this.assertCurrent(); + if (this.finalTurnId) { + if (this.finalTurnId !== turn.id) { + throw new Error("Agents API reply was already committed for a different terminal turn"); + } + return; + } + if (!(await this.reconcile(turn, items))) { + this.reportTranscriptOrderingGap(); + } + await this.endReasoning(); + const completedMessages = items.filter( + (item) => item.type === "message" && item.role === "assistant" && item.status === "completed", + ); + const finalItems = completedMessages.filter((item) => item.phase === "final_answer"); + const visibleItems = finalItems.length + ? finalItems + : completedMessages.filter((item) => item.phase !== "commentary"); + const text = visibleItems + .map( + (item) => + item.content + ?.filter((part) => part.type === "output_text") + .map((part) => part.text ?? "") + .join("") ?? "", + ) + .join("\n"); + const assistant = createAgentHarnessAssistantMessage(this.attribution(), text, { + tokenUsage: this.tokenUsage, + aborted: turn.status === "cancelled", + promptError: turn.error?.message, + timestamp: this.nextTimestamp(), + }); + if (this.canonicalUsageRecorded) { + assistant.usage = this.reply.assistantUsage; + } else { + if (this.usageByTurn.size > 0) { + calculateCost(this.params.model, assistant.usage); + } + this.reply.usage = normalizeUsage(assistant.usage); + this.reply.assistantUsage = assistant.usage; + } + this.classification = classifyAgentHarnessTerminalOutcome({ + assistantTexts: [text], + reasoningText: this.reasoningText(), + promptError: turn.error, + turnCompleted: isTerminalTurn(turn.status), + }); + if (text) { + this.reply.lastAssistant = await this.append({ + ...assistant, + idempotencyKey: `agentsapi:${this.remoteSessionId}:${turn.id}`, + }); + this.assertCurrent(); + await this.params.onAssistantMessageStart?.(); + this.assertCurrent(); + } + await this.emitAssistantSnapshot(`agentsapi:${this.remoteSessionId}:${turn.id}:reply`, text); + if (text) { + this.assertCurrent(); + await this.params.onPartialReply?.({ text }); + this.assertCurrent(); + } + this.finalTurnId = turn.id; + } + + private async recordItem( + turnId: string, + item: AgentsApiItem, + terminal: boolean, + enclosingStatus?: string, + canonical = false, + recordTranscript = true, + ): Promise { + if ( + item.type === "function_call" || + item.type === "function_call_output" || + item.role === "user" + ) { + return; + } + if ( + await this.nativeTools.recordItem( + turnId, + item, + terminal, + enclosingStatus, + canonical, + recordTranscript, + ) + ) { + return; + } + if (item.type !== "message" && item.type !== "reasoning") { + this.recordUnknownType(`item:${item.type}`); + return; + } + const id = this.identity(turnId, item.id); + let state = this.items.get(id); + if (state?.terminal && !canonical) { + return; + } + if (!state) { + state = { + turnId, + item, + terminal: false, + completionObserved: false, + texts: new Map(), + summaries: new Map(), + recoveredPartial: false, + }; + this.items.set(id, state); + this.turnByItem.set(item.id, turnId); + } + state.item = item; + if (!canonical && terminal) { + state.completionObserved = true; + } + // Saved state has no replay cursor. A recovered partial item stays on + // snapshots until its authoritative completion; new items stream normally. + if (canonical && !terminal) { + state.recoveredPartial = true; + } + if (item.type === "message" && item.role === "assistant") { + if (terminal || canonical || state.texts.size === 0) { + state.texts = readTextParts(item.content, "output_text"); + } + await this.emitAssistant(state, terminal); + if ( + canonical && + recordTranscript && + canRecordAgentsApiTranscriptText(item, enclosingStatus, state.completionObserved) && + item.phase === "commentary" && + joinTextParts(state.texts) + ) { + await this.append({ + ...createAgentHarnessAssistantMessage(this.attribution(), joinTextParts(state.texts), { + aborted: false, + timestamp: this.nextTimestamp(), + }), + openclawStreamFallback: { + replacementText: joinTextParts(state.texts), + source: "segment", + itemId: id, + }, + idempotencyKey: `${id}:commentary`, + }); + } + } else if (item.type === "reasoning") { + if (terminal || canonical || state.summaries.size === 0) { + state.summaries = readTextParts(item.summary, "summary_text"); + } + await this.emitReasoning(); + if (terminal) { + const text = joinTextParts(state.summaries); + if ( + canonical && + recordTranscript && + canRecordAgentsApiTranscriptText(item, enclosingStatus, state.completionObserved) && + text + ) { + await this.append({ + ...createAgentHarnessAssistantMessage(this.attribution(), "", { + aborted: false, + timestamp: this.nextTimestamp(), + content: [{ type: "thinking", thinking: text }], + }), + idempotencyKey: `${id}:reasoning`, + }); + } + await this.endReasoning(); + } + } + state.terminal = terminal; + } + + private async emitAssistant( + state: NativeTextState, + terminal: boolean, + delta = "", + ): Promise { + const id = this.identity(state.turnId, state.item.id); + const text = joinTextParts(state.texts); + if (state.item.phase === "commentary") { + const phase = terminal ? "end" : "update"; + if ( + !text.trim() || + (state.lastCommentary?.phase === phase && state.lastCommentary.text === text) + ) { + return; + } + state.lastCommentary = { phase, text }; + await this.emit({ + stream: "item", + data: { + itemId: id, + kind: "preamble", + title: "Preamble", + phase, + progressText: text, + source: "agentsapi", + }, + }); + return; + } + if (!text || state.lastAssistantText === text) { + return; + } + state.lastAssistantText = text; + await this.emitAssistantSnapshot(id, text, delta); + } + + private async emitAssistantSnapshot(itemId: string, text: string, delta = ""): Promise { + const replace = this.visibleAssistantItemId !== itemId; + this.visibleAssistantItemId = itemId; + await this.emit({ stream: "assistant", data: { itemId, @@ -37,195 +563,134 @@ export function createAgentsApiMessageProjection( ...(replace ? { replace: true } : {}), }, }); - }; - const complete = (turnId: string, text: string): void => { - emitAssistantSnapshot(`agentsapi:${remoteSessionId}:${turnId}:reply`, text); - }; - return { - reply, - recordUsage(model: AgentHarnessAttemptParamsV2["model"], turns: Turn[]): void { - const usage = emptyUsage(); - let observed = false; - let reasoningTokens: number | undefined; - // Canonical turn snapshots replace SSE observations, never add to them. - for (const turn of new Map(turns.map((record) => [record.id, record])).values()) { - if (!turn.usage) { - continue; - } - const normalized = normalizeUsage(turn.usage); - if (!normalized) { - continue; - } - observed = true; - usage.input += normalized.input ?? 0; - usage.output += normalized.output ?? 0; - usage.cacheRead += normalized.cacheRead ?? 0; - usage.totalTokens += normalized.total ?? turn.usage.input_tokens + turn.usage.output_tokens; - if (normalized.reasoningTokens !== undefined) { - reasoningTokens = (reasoningTokens ?? 0) + normalized.reasoningTokens; - } - } - if (observed) { - calculateCost(model, usage); - reply.usage = { - ...normalizeUsage(usage), - ...(reasoningTokens !== undefined ? { reasoningTokens } : {}), - }; - } - reply.assistantUsage = usage; - }, - observe(event: AgentSessionEvent): void { - if ( - (event.type === "agent.session.turn.item.added" || - event.type === "agent.session.turn.item.done") && - event.item.type === "message" && - event.item.role === "assistant" && - event.item.id - ) { - assistantPhases.set(event.item.id, event.item.phase); - if (event.type === "agent.session.turn.item.done") { - const parts = new Map(); - event.item.content?.forEach((part, index) => { - if (part.type === "output_text") { - parts.set(index, part.text ?? ""); - } - }); - texts.set(event.item.id, parts); - } - const parts = texts.get(event.item.id); - if (event.item.phase !== "commentary" && parts) { - emitAssistantSnapshot( - `agentsapi:${remoteSessionId}:${event.item.id}`, - joinTextParts(parts), - ); - } - } - if ( - event.type === "agent.session.turn.output_text.delta" || - event.type === "agent.session.turn.output_text.done" - ) { - if (!event.item_id) { - throw new Error("Agents API text event has no item identity"); - } - const parts = texts.get(event.item_id) ?? new Map(); - const index = event.content_index ?? 0; - parts.set( - index, - event.type === "agent.session.turn.output_text.done" - ? event.text - : (parts.get(index) ?? "") + event.delta, - ); - texts.set(event.item_id, parts); - if ( - assistantPhases.has(event.item_id) && - assistantPhases.get(event.item_id) !== "commentary" - ) { - emitAssistantSnapshot( - `agentsapi:${remoteSessionId}:${event.item_id}`, - joinTextParts(parts), - event.type === "agent.session.turn.output_text.delta" ? event.delta : "", - ); - } - } - }, - complete, - commit( - params: AgentHarnessAttemptParamsV2, - turn: Turn, - items: AgentSessionItem[], - assertCurrent: () => void, - ): Promise { - return commitAgentsApiReply( - params, - remoteSessionId, - turn, - items, - assertCurrent, - reply, - complete, - ); - }, - }; -} - -async function commitAgentsApiReply( - params: AgentHarnessAttemptParamsV2, - remoteSessionId: string, - turn: Turn, - items: AgentSessionItem[], - assertCurrent: () => void, - reply: AgentsApiReply, - emitFinalReply: (turnId: string, text: string) => void | Promise, -): Promise { - assertCurrent(); - const { agentId, sessionId, sessionKey, storePath } = params.sessionTarget ?? {}; - if ( - !agentId || - !sessionId || - !sessionKey || - !storePath || - sessionId !== params.sessionId || - agentId !== params.agentId || - sessionKey !== params.sessionKey - ) { - throw new Error("Agents API requires a matching host-prepared session target"); } - const sessionTarget = { ...params.sessionTarget, agentId, sessionId, sessionKey, storePath }; - const completedMessages = items - .filter((item) => item.type === "message") - .filter((item) => item.role === "assistant" && item.status === "completed"); - const finalItems = completedMessages.filter((item) => item.phase === "final_answer"); - const visibleItems = finalItems.length - ? finalItems - : completedMessages.filter((item) => item.phase !== "commentary"); - const text = visibleItems - .map( - (item) => - item.content - ?.filter((part) => part.type === "output_text") - .map((part) => part.text ?? "") - .join("") ?? "", - ) - .join("\n"); - const usage = reply.assistantUsage; - if (text) { - const assistant: AssistantMessage & { idempotencyKey: string } = { - role: "assistant", - content: [{ type: "text", text }], - api: "openai-responses", - provider: "openai", - model: params.model.id, - usage, - stopReason: "stop", - timestamp: Date.now(), - idempotencyKey: `agentsapi:${remoteSessionId}:${turn.id}`, - }; - assertCurrent(); - const append = await appendSessionTranscriptMessageByIdentityStrict({ - ...sessionTarget, - config: params.config, - message: assistant, - prepareMessageAfterIdempotencyCheck: (message) => { - assertCurrent(); - return message; + + private async emitReasoning(): Promise { + this.assertCurrent(); + if (!this.presentationEnabled) { + return; + } + const text = this.reasoningText(); + if (!text.trim() || text === this.lastReasoningText) { + return; + } + this.lastReasoningText = text; + this.reasoningOpen = true; + this.assertCurrent(); + await this.params.onReasoningStream?.({ text, isReasoningSnapshot: true }); + this.assertCurrent(); + } + + private async endReasoning(): Promise { + if (!this.reasoningOpen) { + return; + } + this.reasoningOpen = false; + this.assertCurrent(); + if (!this.presentationEnabled) { + return; + } + await this.params.onReasoningEnd?.(); + this.assertCurrent(); + } + + private reasoningText(): string { + return [...this.items.values()] + .filter((state) => state.item.type === "reasoning") + .map((state) => joinTextParts(state.summaries)) + .filter((text) => text.trim()) + .join("\n\n"); + } + + private recordTurnUsage(turn: NativeTurn): void { + this.assertCurrent(); + if (this.canonicalUsageRecorded) { + return; + } + const usage = normalizeUsage(turn.usage) ?? this.usageByTurn.get(turn.id); + if (!usage) { + return; + } + this.usageByTurn.set(turn.id, { ...usage, contextUsage: { state: "unavailable" } }); + } + + private async emitEnvironment(event: AgentsApiEvent): Promise { + const status = event.type.slice("agent.session.environment.".length); + if (!ENVIRONMENT_TITLES.has(status)) { + this.recordUnknownType(event.type); + return; + } + await this.emit({ + stream: "item", + data: { + itemId: `agentsapi:${this.remoteSessionId}:environment`, + kind: "environment", + title: ENVIRONMENT_TITLES.get(status), + phase: status === "pending" ? "start" : "end", + ...(status === "disconnected" + ? { summary: "Connection unavailable" } + : { + status: + status === "failed" ? "failed" : status === "pending" ? "running" : "completed", + }), + source: "agentsapi", }, }); - if (append.kind === "result") { - reply.lastAssistant = append.result.message; - } - assertCurrent(); - if (append.kind !== "result") { - throw new Error("Agents API assistant transcript append was refused"); - } - await params.onAssistantMessageStart?.(); } - assertCurrent(); - await emitFinalReply(turn.id, text); - assertCurrent(); - if (text) { - await params.onPartialReply?.({ text }); - assertCurrent(); + + private recordUnknownType(type: string): void { + const safeType = type.replace(/[^a-zA-Z0-9_.:-]/gu, "?").slice(0, 128); + if (this.unknownTypes.size >= 20 || this.unknownTypes.has(safeType)) { + return; + } + this.unknownTypes.add(safeType); + embeddedAgentLog.debug("Agents API projection omitted an unfamiliar native type", { + type: safeType, + }); } + + private eventItem(event: AgentsApiEvent): NativeTextState | undefined { + if (!event.item_id) { + throw new Error("Agents API output event has no item identity"); + } + const turnId = event.turn_id ?? this.turnByItem.get(event.item_id); + return turnId ? this.items.get(this.identity(turnId, event.item_id)) : undefined; + } + + private identity(turnId: string, itemId: string): string { + return `agentsapi:${this.remoteSessionId}:${turnId}:${itemId}`; + } + + private attribution() { + return { api: "openai-responses" as const, provider: "openai", modelId: this.params.model.id }; + } + + private nextTimestamp(): number { + this.timestamp = Math.max(Date.now(), this.timestamp + 1); + return this.timestamp; + } + + private async emit(event: AgentEvent): Promise { + this.assertCurrent(); + if (!this.presentationEnabled) { + return; + } + await this.emitEvent(event); + this.assertCurrent(); + } + + private append(message: TMessage): Promise { + return appendAgentsApiTranscriptMessage(this.params, message, this.assertCurrent); + } +} + +export function createAgentsApiMessageProjection( + params: AgentHarnessAttemptParamsV2, + remoteSessionId: string, + emitEvent: (event: AgentEvent) => void | Promise, + assertCurrent: () => void, +) { + return new AgentsApiMessageProjection(params, remoteSessionId, emitEvent, assertCurrent); } function emptyUsage(): AssistantMessage["usage"] { @@ -241,9 +706,35 @@ function emptyUsage(): AssistantMessage["usage"] { }; } -function joinTextParts(parts: Map): string { - return [...parts.entries()] - .toSorted(([left], [right]) => left - right) - .map(([, text]) => text) - .join(""); +function isTerminalTurn(status?: string): boolean { + return status === "completed" || status === "failed" || status === "cancelled"; } + +const INTERNAL_EVENT_TYPES = new Set([ + "error", + "agent.session.created", + "agent.session.idle", + "agent.session.in_progress", + "agent.session.requires_action", + "agent.session.failed", + "agent.session.error", + "agent.session.turn.created", + "agent.session.turn.in_progress", + "agent.session.turn.completed", + "agent.session.turn.failed", + "agent.session.turn.cancelled", + "agent.session.turn.content_part.added", + "agent.session.turn.content_part.done", + // Native delegation remains disabled until child history and settlement exist. + "agent.session.subagent.active", + "agent.session.subagent.closed", + "agent.session.subagent.created", +]); + +const ENVIRONMENT_TITLES = new Map([ + ["pending", "Agent environment starting"], + ["ready", "Agent environment ready"], + ["connected", "Agent environment connected"], + ["disconnected", "Agent environment disconnected"], + ["failed", "Agent environment failed"], +]); diff --git a/extensions/agentsapi/agentsapi-native-items.ts b/extensions/agentsapi/agentsapi-native-items.ts new file mode 100644 index 000000000000..c45bb6e6294d --- /dev/null +++ b/extensions/agentsapi/agentsapi-native-items.ts @@ -0,0 +1,191 @@ +import { + TOOL_TRANSCRIPT_OUTPUT_MAX_CHARS, + truncateNativeToolTranscriptText, +} from "openclaw/plugin-sdk/agent-harness-attempt-runtime"; +import { + inferToolMetaFromArgs, + sanitizeToolArgs, + sanitizeToolResult, + type AgentHarnessAttemptParamsV2, +} from "openclaw/plugin-sdk/agent-harness-runtime"; +import { asOptionalRecord } from "openclaw/plugin-sdk/string-coerce-runtime"; +import type { AgentsApiItem } from "./agentsapi-client.js"; + +export type AgentsApiNativeTool = { + name: string; + args: Record; + meta?: string; + commandBearing: boolean; +}; + +export type AgentsApiNativeToolOutcome = { + status: "running" | "completed" | "failed" | "cancelled" | "unknown"; + isError: boolean; + outcomeUnknown: boolean; + error?: string; + errorCode?: string; +}; + +export function agentsApiNativeTool( + item: AgentsApiItem, + params: AgentHarnessAttemptParamsV2, +): AgentsApiNativeTool | undefined { + let name: string; + let args: Record; + if (item.type === "command_execution") { + name = "bash"; + args = { + ...(item.command !== undefined ? { command: item.command } : {}), + ...(typeof item.cwd === "string" ? { cwd: item.cwd } : {}), + }; + } else if (item.type === "mcp_call") { + name = item.server_label ? `${item.server_label}.${item.name ?? "mcp"}` : (item.name ?? "mcp"); + args = + asOptionalRecord(item.arguments) ?? + (item.arguments === undefined ? {} : { arguments: item.arguments }); + } else if (item.type === "web_search_call") { + name = "web_search"; + args = item.action ? { ...item.action } : {}; + } else { + return undefined; + } + args = asOptionalRecord(sanitizeToolArgs(args)) ?? {}; + const meta = inferToolMetaFromArgs(name, args, { + detailMode: params.toolProgressDetail === "raw" ? "raw" : "explain", + }); + return { + name, + args, + ...(meta ? { meta } : {}), + commandBearing: item.type === "command_execution", + }; +} + +export function agentsApiNativeToolOutcome( + item: AgentsApiItem, + enclosingStatus?: string, +): AgentsApiNativeToolOutcome { + const nativeError = item.error === null || item.error === undefined ? undefined : item.error; + const errorRecord = asOptionalRecord(nativeError); + const errorCode = typeof errorRecord?.code === "string" ? errorRecord.code : undefined; + const errorText = nativeError === undefined ? undefined : nativeValueText(nativeError); + const failedCommand = + item.type === "command_execution" && typeof item.exit_code === "number" && item.exit_code !== 0; + let status: AgentsApiNativeToolOutcome["status"]; + if (item.status === "failed" || nativeError !== undefined || failedCommand) { + status = "failed"; + } else if (item.status === "completed") { + status = item.type === "command_execution" && item.exit_code == null ? "unknown" : "completed"; + } else if (item.status === "in_progress" && enclosingStatus === undefined) { + status = "running"; + } else if ( + (item.status === "incomplete" || item.status === "in_progress") && + enclosingStatus === "cancelled" + ) { + status = "cancelled"; + } else { + status = "unknown"; + } + const error = + status === "failed" + ? errorText || + (failedCommand + ? `Command exited with code ${item.exit_code}` + : "Agents API native tool failed") + : status === "cancelled" + ? "Agents API native tool was cancelled" + : status === "unknown" + ? "Agents API native tool outcome is unavailable" + : undefined; + return { + status, + isError: status !== "completed" && status !== "running", + outcomeUnknown: status === "unknown", + ...(error ? { error } : {}), + ...(errorCode ? { errorCode } : {}), + }; +} + +export function agentsApiNativeToolDetails( + sessionId: string, + turnId: string, + item: AgentsApiItem, + outcome: AgentsApiNativeToolOutcome, + capturedOutput?: string, +): Record { + const nativeOutput = item.output ?? capturedOutput; + const output = + nativeOutput === undefined || nativeOutput === null + ? undefined + : boundedNativeValue(nativeOutput); + return ( + asOptionalRecord( + sanitizeToolResult({ + status: outcome.status, + native: { + backend: "agentsapi", + sessionId, + turnId, + itemId: item.id, + itemType: item.type, + status: item.status ?? null, + }, + ...(typeof item.exit_code === "number" ? { exitCode: item.exit_code } : {}), + ...(typeof item.duration_ms === "number" ? { durationMs: item.duration_ms } : {}), + ...(typeof item.cwd === "string" ? { cwd: item.cwd } : {}), + ...(item.server_label ? { serverLabel: item.server_label } : {}), + ...(item.type === "web_search_call" + ? { ...(item.action ? { action: item.action } : {}), resultAvailability: "unavailable" } + : { + outputAvailability: + output === undefined + ? "unavailable" + : item.output == null + ? "partial" + : "canonical", + ...(output + ? { output: output.value, ...(output.truncated ? { outputTruncated: true } : {}) } + : {}), + }), + ...(item.error != null ? { nativeError: boundedNativeValue(item.error).value } : {}), + ...(outcome.error ? { error: outcome.error } : {}), + ...(outcome.errorCode ? { errorCode: outcome.errorCode } : {}), + ...(outcome.outcomeUnknown ? { outcomeUnknown: true } : {}), + }), + ) ?? {} + ); +} + +export function agentsApiNativeToolOutput( + item: AgentsApiItem, + streamed?: string, +): string | undefined { + if (item.type === "web_search_call") { + return undefined; + } + return item.output !== undefined && item.output !== null + ? nativeValueText(item.output) + : streamed; +} + +function nativeValueText(value: unknown): string { + const sanitized = sanitizeNativeValue(value); + const text = + typeof sanitized === "string" ? sanitized : (JSON.stringify(sanitized, null, 2) ?? ""); + return truncateNativeToolTranscriptText(text, "Agents API"); +} + +function boundedNativeValue(value: unknown): { value: unknown; truncated: boolean } { + const sanitized = sanitizeNativeValue(value); + const text = + typeof sanitized === "string" ? sanitized : (JSON.stringify(sanitized, null, 2) ?? ""); + return text.length <= TOOL_TRANSCRIPT_OUTPUT_MAX_CHARS + ? { value: sanitized, truncated: false } + : { value: truncateNativeToolTranscriptText(text, "Agents API"), truncated: true }; +} + +function sanitizeNativeValue(value: unknown): unknown { + return Array.isArray(value) + ? asOptionalRecord(sanitizeToolResult({ content: value }))?.content + : sanitizeToolResult(value); +} diff --git a/extensions/agentsapi/agentsapi-native-tool-projection.ts b/extensions/agentsapi/agentsapi-native-tool-projection.ts new file mode 100644 index 000000000000..c17eb430f269 --- /dev/null +++ b/extensions/agentsapi/agentsapi-native-tool-projection.ts @@ -0,0 +1,569 @@ +import { createHash } from "node:crypto"; +import { + formatNativeToolOutput, + formatNativeToolSummary, + MAX_TOOL_OUTPUT_DELTA_MESSAGES_PER_ITEM, + NativeToolOutputAccumulator, + truncateNativeToolTranscriptText, + TOOL_TRANSCRIPT_OUTPUT_MAX_CHARS, +} from "openclaw/plugin-sdk/agent-harness-attempt-runtime"; +import { + formatToolProgressOutput, + projectAgentToolActivity, + sanitizeToolResult, + TOOL_PROGRESS_OUTPUT_MAX_CHARS, + type AgentHarnessAttemptParamsV2, + type AgentHarnessAttemptResult, +} from "openclaw/plugin-sdk/agent-harness-runtime"; +import { truncateUtf16Safe } from "openclaw/plugin-sdk/text-utility-runtime"; +import type { AgentsApiEvent, AgentsApiItem } from "./agentsapi-client.js"; +import { + agentsApiNativeTool, + agentsApiNativeToolDetails, + agentsApiNativeToolOutcome, + agentsApiNativeToolOutput, + type AgentsApiNativeTool, + type AgentsApiNativeToolOutcome, +} from "./agentsapi-native-items.js"; +import { + recordAgentsApiNativeToolInvocation, + recordAgentsApiNativeToolTranscript, +} from "./agentsapi-transcript.js"; + +type AgentEvent = Parameters>[0]; +type NativeToolState = { + turnId: string; + item: AgentsApiItem; + canonicalItem?: AgentsApiItem; + terminal: boolean; + tool?: AgentsApiNativeTool; + startProjected: boolean; + callRecorded: boolean; + resultRecorded: boolean; + provisionalTerminalObserved: boolean; + canonicalTerminalObserved: boolean; + terminalProjectionHash?: string; + recoveredPartial: boolean; + outputProgressChars: number; + outputProgressMessages: number; + commandOutputChars: number; + commandOutputMessages: number; + recoveredOutput?: string; + recoveredOutputTruncated?: boolean; +}; + +/** Native tools own lifecycle, bounded output, and canonical transcript facts. */ +export class AgentsApiNativeToolProjection { + private readonly items = new Map(); + private readonly turnByItem = new Map(); + private readonly output = new NativeToolOutputAccumulator("Agents API"); + private readonly metas = new Map(); + private nativeToolError: AgentHarnessAttemptResult["lastToolError"]; + + constructor( + private readonly params: AgentHarnessAttemptParamsV2, + private readonly remoteSessionId: string, + private readonly emitEvent: (event: AgentEvent) => void | Promise, + private readonly assertCurrent: () => void, + private readonly nextTimestamp: () => number, + private readonly isPresentationEnabled: () => boolean, + ) {} + + get toolMetas(): AgentHarnessAttemptResult["toolMetas"] { + return [...this.metas.values()]; + } + + get lastToolError(): AgentHarnessAttemptResult["lastToolError"] { + return this.nativeToolError; + } + + get itemLifecycle(): NonNullable { + const states = [...this.items.values()]; + const completedCount = states.filter((state) => state.terminal).length; + return { + startedCount: states.length, + completedCount, + activeCount: states.length - completedCount, + }; + } + + get hadPotentialSideEffects(): boolean { + return [...this.items.values()].some( + (state) => state.item.type === "command_execution" || state.item.type === "mcp_call", + ); + } + + resolveTurnId(itemId: string): string | undefined { + return this.turnByItem.get(itemId); + } + + async recordItem( + turnId: string, + item: AgentsApiItem, + terminal: boolean, + enclosingStatus?: string, + canonical = false, + recordTranscript = true, + ): Promise { + this.assertCurrent(); + const tool = agentsApiNativeTool(item, this.params); + if (!tool) { + return false; + } + const id = this.identity(turnId, item.id); + let state = this.items.get(id); + if (state?.terminal && !canonical) { + return true; + } + if (!state) { + state = { + turnId, + item, + tool, + terminal: false, + startProjected: false, + callRecorded: false, + resultRecorded: false, + provisionalTerminalObserved: false, + canonicalTerminalObserved: false, + recoveredPartial: false, + outputProgressChars: 0, + outputProgressMessages: 0, + commandOutputChars: 0, + commandOutputMessages: 0, + }; + this.items.set(id, state); + this.turnByItem.set(item.id, turnId); + } + state.item = item; + state.tool = tool; + if (canonical) { + state.canonicalItem = item; + } + // Saved state has no replay cursor. Recovered partial output waits for an + // authoritative completion; new native tool items continue streaming. + if (canonical && !terminal) { + state.recoveredPartial = true; + if (item.type === "command_execution" && typeof item.output === "string") { + state.recoveredOutput = agentsApiNativeToolOutput(item); + state.recoveredOutputTruncated = item.output.length > TOOL_TRANSCRIPT_OUTPUT_MAX_CHARS; + } + } + await this.startTool(state); + if (canonical && recordTranscript && !state.callRecorded) { + state.callRecorded = await recordAgentsApiNativeToolInvocation( + this.params, + this.remoteSessionId, + turnId, + item, + this.assertCurrent, + this.nextTimestamp, + ); + } + if (terminal) { + await this.finishTool(state, enclosingStatus, canonical, recordTranscript); + } + state.terminal = terminal; + return true; + } + + async reconcileRemaining( + turnId: string, + status: string | undefined, + currentItemIds: ReadonlySet, + ): Promise { + this.assertCurrent(); + let transcriptReady = true; + for (const state of this.items.values()) { + if (state.turnId !== turnId || (state.resultRecorded && state.canonicalTerminalObserved)) { + continue; + } + if (!state.resultRecorded && !currentItemIds.has(state.item.id)) { + // Retain previously retrieved facts, but do not promise their original + // position when the terminal snapshot no longer contains this item. + transcriptReady = false; + } + await this.finishTool(state, status, true); + state.terminal = true; + } + return transcriptReady; + } + + async observeOutput(event: AgentsApiEvent): Promise { + this.assertCurrent(); + const state = this.eventItem(event); + if ( + !state || + state.terminal || + state.recoveredPartial || + state.item.type !== "command_execution" || + !event.delta + ) { + return; + } + const id = this.identity(state.turnId, state.item.id); + const delta = sanitizeToolResult(event.delta); + this.output.append(id, delta); + if ( + state.commandOutputChars < TOOL_PROGRESS_OUTPUT_MAX_CHARS && + state.commandOutputMessages < MAX_TOOL_OUTPUT_DELTA_MESSAGES_PER_ITEM + ) { + const remaining = TOOL_PROGRESS_OUTPUT_MAX_CHARS - state.commandOutputChars; + const text = truncateUtf16Safe(delta, remaining); + state.commandOutputChars += text.length; + state.commandOutputMessages += 1; + await this.emit({ + stream: "command_output", + data: { + itemId: id, + toolCallId: id, + phase: "delta", + name: state.tool?.name ?? "bash", + title: + typeof state.tool?.args.command === "string" + ? state.tool.args.command + : "Command output", + output: delta.length > remaining ? `${text}\n...(truncated)...` : text, + status: "running", + ...(typeof state.tool?.args.cwd === "string" ? { cwd: state.tool.args.cwd } : {}), + }, + }); + } + await this.emitToolOutput(state, event.delta, false); + } + + private async startTool(state: NativeToolState): Promise { + if (!state.tool || state.startProjected) { + return; + } + const id = this.identity(state.turnId, state.item.id); + const { name, args, meta, commandBearing } = state.tool; + state.startProjected = true; + this.metas.set(id, { toolCallId: id, toolName: name, ...(meta ? { meta } : {}) }); + await this.emit({ + stream: "item", + data: projectAgentToolActivity({ toolCallId: id, name, phase: "start", args, meta }), + }); + await this.emit({ + stream: "tool", + data: { + phase: "start", + itemId: id, + toolCallId: id, + name, + args, + ...(meta ? { meta } : {}), + ...(commandBearing ? { commandBearing: true } : {}), + }, + }); + if (this.shouldEmitToolResult()) { + await this.emitToolProgress( + id, + formatNativeToolSummary(name, this.formattedMeta(state.tool)), + ); + } + } + + private async finishTool( + state: NativeToolState, + enclosingStatus?: string, + canonical = false, + recordTranscript = true, + ): Promise { + if (!state.tool) { + return; + } + const id = this.identity(state.turnId, state.item.id); + const item = canonical ? (state.canonicalItem ?? state.item) : state.item; + const tool = + canonical && state.canonicalItem + ? agentsApiNativeTool(state.canonicalItem, this.params) + : state.tool; + if (!tool) { + return; + } + const { name, args, meta, commandBearing } = tool; + const savedItemUnavailable = canonical && !state.canonicalItem; + const outcome: AgentsApiNativeToolOutcome = savedItemUnavailable + ? { + status: enclosingStatus === "cancelled" ? "cancelled" : "unknown", + isError: true, + outcomeUnknown: enclosingStatus !== "cancelled", + error: + enclosingStatus === "cancelled" + ? "Agents API native tool was cancelled" + : "Agents API native tool item is unavailable in saved state", + } + : agentsApiNativeToolOutcome(item, enclosingStatus); + const capturedOutput = state.recoveredOutput ?? this.output.textByItem.get(id); + const output = agentsApiNativeToolOutput(item, capturedOutput); + const details = { + ...agentsApiNativeToolDetails( + this.remoteSessionId, + state.turnId, + item, + outcome, + capturedOutput, + ), + ...(savedItemUnavailable + ? { + savedItemAvailability: "unavailable", + ...(item.type !== "web_search_call" + ? { outputAvailability: output === undefined ? "unavailable" : "partial" } + : {}), + } + : {}), + }; + const captureTruncated = + item.output == null && (state.recoveredOutputTruncated || this.output.isTruncated(id)); + this.metas.set(id, { + toolCallId: id, + toolName: name, + ...(meta ? { meta } : {}), + isError: outcome.isError, + }); + if (canonical && recordTranscript && state.canonicalItem && !state.resultRecorded) { + // Only retrieved items supply durable calls and results. Streamed tool + // names, arguments, and output may still be partial. + state.resultRecorded = await recordAgentsApiNativeToolTranscript( + this.params, + this.remoteSessionId, + state.turnId, + state.canonicalItem, + this.assertCurrent, + this.nextTimestamp, + { enclosingStatus, capturedOutput, captureTruncated }, + ); + } + if ( + canonical && + (state.canonicalItem ? !state.canonicalTerminalObserved : !state.provisionalTerminalObserved) + ) { + this.assertCurrent(); + const resolution = this.params.observeToolTerminal?.({ + toolCallId: id, + toolName: name, + arguments: args, + ...(meta ? { meta } : {}), + executionStarted: true, + outcome: outcome.isError ? "failure" : "success", + ...(outcome.isError + ? { + failure: { + ...(outcome.error ? { error: outcome.error } : {}), + ...(outcome.errorCode ? { errorCode: outcome.errorCode } : {}), + }, + } + : {}), + nativeMutation: { + mutatingAction: item.type !== "web_search_call", + replaySafe: item.type === "web_search_call", + }, + }); + this.assertCurrent(); + if (state.canonicalItem) { + state.canonicalTerminalObserved = true; + } else { + state.provisionalTerminalObserved = true; + } + if (resolution) { + this.nativeToolError = resolution.lastToolError; + } else if (outcome.isError) { + this.nativeToolError = { + toolName: name, + ...(meta ? { meta } : {}), + ...(outcome.error ? { error: outcome.error } : {}), + ...(outcome.errorCode ? { errorCode: outcome.errorCode } : {}), + ...(item.type !== "web_search_call" ? { mutatingAction: true } : {}), + }; + } else if (this.nativeToolError?.mutatingAction !== true) { + this.nativeToolError = undefined; + } + } + this.assertCurrent(); + if (!this.isPresentationEnabled()) { + return; + } + const commandFacts = commandBearing + ? { + title: typeof args.command === "string" ? args.command : "Command output", + ...(!savedItemUnavailable && typeof item.exit_code === "number" + ? { exitCode: item.exit_code } + : {}), + ...(!savedItemUnavailable && typeof item.duration_ms === "number" + ? { durationMs: item.duration_ms } + : {}), + ...(typeof args.cwd === "string" ? { cwd: args.cwd } : {}), + } + : undefined; + const presentationOutput = + output !== undefined && commandBearing ? (formatToolProgressOutput(output) ?? "") : output; + const toolSnapshot = { + phase: "result", + itemId: id, + toolCallId: id, + name, + args, + status: outcome.status, + isError: outcome.isError, + result: details, + ...(presentationOutput !== undefined ? { output: presentationOutput } : {}), + ...(meta ? { meta } : {}), + ...(commandBearing ? { commandBearing: true, ...commandFacts } : {}), + }; + const itemSnapshot = projectAgentToolActivity({ + toolCallId: id, + name, + phase: "result", + args, + meta, + status: outcome.status === "cancelled" ? "unknown" : outcome.status, + result: { details }, + isError: outcome.isError, + }); + const commandSnapshot = commandFacts + ? { + itemId: id, + toolCallId: id, + phase: "end", + name, + ...commandFacts, + ...(presentationOutput !== undefined ? { output: presentationOutput } : {}), + status: outcome.status, + } + : undefined; + // The bounded digest detects changed projected facts without retaining a + // second payload. Canonical enrichment replaces the same terminal item. + const projectionHash = nativeTerminalProjectionHash({ + tool: toolSnapshot, + item: itemSnapshot, + command: commandSnapshot, + }); + if (state.terminalProjectionHash === projectionHash) { + return; + } + const replacement = { + replaceable: true, + ...(state.terminalProjectionHash ? { replace: true } : {}), + }; + await this.emit({ stream: "tool", data: { ...toolSnapshot, ...replacement } }); + await this.emit({ stream: "item", data: { ...itemSnapshot, ...replacement } }); + if (commandSnapshot) { + await this.emit({ stream: "command_output", data: { ...commandSnapshot, ...replacement } }); + } + state.terminalProjectionHash = projectionHash; + if (output !== undefined && state.outputProgressMessages === 0) { + await this.emitToolOutput(state, output, true, outcome.isError); + } + } + + private async emitToolOutput( + state: NativeToolState, + output: string, + terminal: boolean, + isError = false, + ): Promise { + if ( + !state.tool || + !this.shouldEmitToolOutput() || + state.outputProgressChars >= TOOL_PROGRESS_OUTPUT_MAX_CHARS || + state.outputProgressMessages >= MAX_TOOL_OUTPUT_DELTA_MESSAGES_PER_ITEM + ) { + return; + } + const remaining = TOOL_PROGRESS_OUTPUT_MAX_CHARS - state.outputProgressChars; + const sanitized = sanitizeToolResult(output); + const truncated = sanitized.length > remaining; + const text = truncateUtf16Safe(sanitized, remaining); + state.outputProgressChars += text.length; + state.outputProgressMessages += 1; + await this.emitToolProgress( + this.identity(state.turnId, state.item.id), + formatNativeToolOutput( + state.tool.name, + this.formattedMeta(state.tool), + truncated ? `${text}\n...(truncated)...` : text, + ), + terminal && isError, + ); + } + + private async emitToolProgress(itemId: string, text: string, isError = false): Promise { + this.assertCurrent(); + if (!this.isPresentationEnabled()) { + return; + } + await this.params.onToolResult?.({ + text: truncateNativeToolTranscriptText(text, "Agents API"), + ...(this.params.messageChannel || this.params.messageProvider + ? { channelData: { openclawToolProgressId: `tool:${itemId}` } } + : {}), + ...(isError ? { isError: true } : {}), + }); + this.assertCurrent(); + } + + private formattedMeta(tool: AgentsApiNativeTool): string | undefined { + return !tool.commandBearing || + (!(this.params.messageChannel ?? this.params.messageProvider) && this.shouldEmitToolOutput()) + ? tool.meta + : undefined; + } + + private shouldEmitToolResult(): boolean { + this.assertCurrent(); + if (!this.isPresentationEnabled()) { + return false; + } + return ( + this.params.shouldEmitToolResult?.() ?? + (this.params.verboseLevel === "on" || this.params.verboseLevel === "full") + ); + } + + private shouldEmitToolOutput(): boolean { + this.assertCurrent(); + if (!this.isPresentationEnabled()) { + return false; + } + return this.params.shouldEmitToolOutput?.() ?? this.params.verboseLevel === "full"; + } + + private eventItem(event: AgentsApiEvent): NativeToolState | undefined { + if (!event.item_id) { + throw new Error("Agents API output event has no item identity"); + } + const turnId = event.turn_id ?? this.turnByItem.get(event.item_id); + return turnId ? this.items.get(this.identity(turnId, event.item_id)) : undefined; + } + + private identity(turnId: string, itemId: string): string { + return `agentsapi:${this.remoteSessionId}:${turnId}:${itemId}`; + } + + private async emit(event: AgentEvent): Promise { + this.assertCurrent(); + if (!this.isPresentationEnabled()) { + return; + } + await this.emitEvent(event); + this.assertCurrent(); + } +} + +function nativeTerminalProjectionHash(snapshot: Record): string { + return createHash("sha256") + .update( + JSON.stringify(snapshot, (_key, value: unknown) => { + if (value === null || typeof value !== "object" || Array.isArray(value)) { + return value; + } + return Object.fromEntries( + Object.entries(value).toSorted(([left], [right]) => + left < right ? -1 : left > right ? 1 : 0, + ), + ); + }), + ) + .digest("hex"); +} diff --git a/extensions/agentsapi/agentsapi-session.test.ts b/extensions/agentsapi/agentsapi-session.test.ts index 683d080d3a85..38c3a1da0cc0 100644 --- a/extensions/agentsapi/agentsapi-session.test.ts +++ b/extensions/agentsapi/agentsapi-session.test.ts @@ -1,7 +1,8 @@ -import type { AgentSessionEvent } from "openai/resources/beta/agents/agents"; +import type { AgentSession, AgentSessionMessage } from "openai/resources/beta/agents/agents"; +import type { Turn } from "openai/resources/beta/agents/sessions/turns"; import { afterEach, describe, expect, it, vi } from "vitest"; import { z } from "zod"; -import { AgentsApiClient } from "./agentsapi-client.js"; +import { AgentsApiClient, type AgentsApiEvent } from "./agentsapi-client.js"; import { createAgentsApiSession } from "./agentsapi-session.js"; const { fetchWithSsrFGuardMock } = vi.hoisted(() => ({ @@ -21,6 +22,9 @@ describe("Agents API native session receipts", () => { it("waits for session idle and the admitted input receipt after the root turn completes", async () => { const controller = new AbortController(); const stream = createEventStream(); + let savedTurns: Turn[] = []; + let savedItems: AgentSessionMessage[] = []; + let savedSession = createSavedSession("in_progress"); fetchWithSsrFGuardMock.mockImplementation(async (request) => { request.beforeRequest?.(); if (request.init?.method === "POST") { @@ -31,7 +35,7 @@ describe("Agents API native session receipts", () => { } return guardedResponse( request.url, - Response.json({ data: [], has_more: false, last_id: null }), + savedStateResponse(request.url, savedTurns, savedItems, savedSession), ); }); const session = createSession(controller.signal, (event) => stream.observe(event)); @@ -43,16 +47,24 @@ describe("Agents API native session receipts", () => { ); await submitted.promise; + const createdTurn = createSavedTurn("in_progress"); + savedTurns = [createdTurn]; await stream.send({ type: "agent.session.turn.created", - turn: { id: "turn-fixture", subagent_id: null }, + turn: createdTurn, }); + const completedTurn = createSavedTurn("completed"); + savedTurns = [completedTurn]; await stream.send({ type: "agent.session.turn.completed", - turn: { id: "turn-fixture", subagent_id: null }, + turn: completedTurn, }); + // Observing the next event fences the completed turn's saved-state reconciliation. + await stream.send({ type: "agent.session.in_progress" }); expect(session.isSettled()).toBe(false); + savedSession = createSavedSession("idle"); + savedItems = [createSavedMessage("assistant-fixture", "assistant", "Fixture reply")]; await stream.send({ type: "agent.session.idle" }); await stream.send({ type: "agent.session.turn.output_text.done", @@ -62,9 +74,11 @@ describe("Agents API native session receipts", () => { }); expect(session.isSettled()).toBe(false); + const inputReceipt = createSavedMessage("input-fixture", "user", "Fixture prompt"); + savedItems = [inputReceipt, ...savedItems]; await stream.send({ type: "agent.session.turn.item.done", - item: { id: "input-fixture", type: "message", role: "user", turn_id: "turn-fixture" }, + item: inputReceipt, }); expect(session.isSettled()).toBe(false); await stream.send({ type: "agent.session.idle" }); @@ -104,12 +118,15 @@ describe("Agents API native session receipts", () => { if (new Headers(request.init?.headers).get("accept") === "text/event-stream") { return guardedResponse(request.url, stream.response(request.signal)); } - if (!inputTypes.includes("agent.session.input.cancel")) { + if (new URL(request.url).pathname !== "/v1/agents/sessions/session-fixture") { return guardedResponse( request.url, - Response.json({ data: [], has_more: false, last_id: null }), + savedStateResponse(request.url, [], [], createSavedSession("in_progress")), ); } + if (!inputTypes.includes("agent.session.input.cancel")) { + throw new Error("Expected native cancellation before the idle request"); + } idleRequested.resolve(); return guardedResponse(request.url, await idleReceipt.promise); }); @@ -139,7 +156,7 @@ describe("Agents API native session receipts", () => { await idleRequested.promise; expect(inputTypes).toEqual(["agent.session.input.message", "agent.session.input.cancel"]); expect(cancellationSettled).toBe(false); - idleReceipt.resolve(Response.json({ id: "session-fixture", status: "idle", error: null })); + idleReceipt.resolve(Response.json(createSavedSession("idle"))); await cancellation; await expect(run).rejects.toBe(interruption); @@ -160,7 +177,7 @@ describe("Agents API native session receipts", () => { }); }); -function createSession(signal: AbortSignal, onEvent: (event: AgentSessionEvent) => void) { +function createSession(signal: AbortSignal, onEvent: (event: AgentsApiEvent) => void) { return createAgentsApiSession({ client: new AgentsApiClient("fixture-not-a-real-api-key", () => {}), cleanupClient: new AgentsApiClient("fixture-not-a-real-api-key", () => {}), @@ -171,6 +188,92 @@ function createSession(signal: AbortSignal, onEvent: (event: AgentSessionEvent) }); } +function savedStateResponse( + url: string, + turns: Turn[], + items: AgentSessionMessage[], + session: AgentSession, +) { + switch (new URL(url).pathname) { + case "/v1/agents/sessions/session-fixture/turns": + return Response.json({ data: turns, has_more: false }); + case "/v1/agents/sessions/session-fixture/items": + return Response.json({ data: items, has_more: false }); + case "/v1/agents/sessions/session-fixture": + return Response.json(session); + default: + throw new Error(`Unexpected native session request: ${url}`); + } +} + +function createSavedTurn(status: "in_progress" | "completed"): Turn { + return { + id: "turn-fixture", + agent_id: "agent-fixture", + completed_at: status === "completed" ? 2 : null, + created_at: 1, + error: null, + object: "agent.session.turn", + session_id: "session-fixture", + started_at: 1, + status, + subagent_id: null, + usage: null, + }; +} + +function createSavedMessage( + id: string, + role: "user" | "assistant", + text: string, +): AgentSessionMessage { + return { + id, + content: [{ type: role === "user" ? "input_text" : "output_text", text }], + phase: role === "user" ? null : "final_answer", + role, + status: "completed", + turn_id: "turn-fixture", + type: "message", + }; +} + +function createSavedSession(status: AgentSession["status"]): AgentSession { + return { + id: "session-fixture", + agent: { + id: "agent-fixture", + instructions: "Fixture instructions", + model: "fixture-model", + multi_agent: { enabled: false, max_concurrent_subagents: null }, + name: null, + reasoning: { effort: null, summary: null }, + service_tier: "auto", + text: { format: { type: "text" }, verbosity: "medium" }, + tools: [], + }, + created_at: 1, + environment: { + id: "environment-fixture", + capability_directories: [], + files: [], + network: { access: "disabled", allowed_domains: [] }, + packages: { npm: [], python: [], system: [] }, + plugins: [], + skills: [], + type: "openai_hosted", + }, + error: null, + last_active_at: 2, + metadata: {}, + object: "agent.session", + required_actions: [], + status, + usage: null, + vault_ids: [], + }; +} + function guardedResponse(url: string, response: Response) { return { response, finalUrl: url, release: async () => {} }; } @@ -214,7 +317,7 @@ function createEventStream() { streamController.enqueue(encoder.encode(`data: ${JSON.stringify(event)}\n\n`)); return observed.promise; }, - observe(event: AgentSessionEvent) { + observe(event: AgentsApiEvent) { waiters.get(event.type)?.shift()?.(); }, }; diff --git a/extensions/agentsapi/agentsapi-session.ts b/extensions/agentsapi/agentsapi-session.ts index 6208fb6fe754..3ccc0cd9cc98 100644 --- a/extensions/agentsapi/agentsapi-session.ts +++ b/extensions/agentsapi/agentsapi-session.ts @@ -1,8 +1,13 @@ import { setTimeout as delay } from "node:timers/promises"; -import type { AgentSessionEvent } from "openai/resources/beta/agents/agents"; import type { Turn } from "openai/resources/beta/agents/sessions/turns"; import { formatErrorMessage } from "openclaw/plugin-sdk/error-runtime"; -import { AgentsApiClient } from "./agentsapi-client.js"; +import { asOptionalRecord } from "openclaw/plugin-sdk/string-coerce-runtime"; +import { + AgentsApiClient, + AgentsApiError, + type AgentsApiEvent, + type AgentsApiItem, +} from "./agentsapi-client.js"; /** Native input receipts and session idle, together, establish Agents API completion. */ export function createAgentsApiSession(options: { @@ -11,26 +16,34 @@ export function createAgentsApiSession(options: { sessionId: string; signal: AbortSignal; assertCurrent: () => void; - onEvent: (event: AgentSessionEvent) => void; + onEvent: (event: AgentsApiEvent) => void | Promise; + onReconcile?: (turn: Turn, items: AgentsApiItem[]) => Promise; + onReconcileHistory?: (entries: Array<{ turn: Turn; items: AgentsApiItem[] }>) => Promise; onSettled?: () => void; onUsageError?: (error: unknown) => void; + onTranscriptOrderingGap?: () => void; }) { const { client, cleanupClient, sessionId, signal, assertCurrent } = options; let streamController = new AbortController(); let submitted = false; let stopped = false; + let closed = false; let settled = false; - let rootTurn: Turn | undefined; - let turnFailure: string | undefined; + let rootTurn: Turn | AgentsApiEvent["turn"]; + let turnFailure: AgentsApiError | undefined; let cancelled = false; let submission: Promise = Promise.resolve(); let admittedSubmission: Promise = Promise.resolve(); let cancellation: Promise | undefined; let admittedMessageCount = 0; + let baselineTurnId: string | undefined; + let baselineCaptured = false; const observedInputItems = new Set(); const coordinatorTurnIds = new Set(); + const excludedTurnIds = new Set(); + const itemTurnIds = new Map(); + const excludedItemIds = new Set(); let latestInputTurnId: string | undefined; - let baselineTurnId: string | undefined; let usageTurns: Promise | undefined; const isAvailable = () => submitted && !stopped && !settled && !rootTurn && !signal.aborted; @@ -84,13 +97,6 @@ export function createAgentsApiSession(options: { }; signal.addEventListener("abort", onAbort, { once: true }); - const collectInputs = async (turnId: string) => { - for (const item of await client.items(sessionId, turnId, signal)) { - if (item.type === "message" && item.role === "user" && item.id) { - observedInputItems.add(item.id); - } - } - }; const assertSessionUsable = (session: { status: string; error: string | null }) => { if (session.status === "failed") { throw new Error(session.error ?? "Agents API session failed"); @@ -99,6 +105,115 @@ export function createAgentsApiSession(options: { throw new Error("Agents API MVP cannot continue: agent.session.requires_action"); } }; + const readAdmittedTurns = async (readClient: AgentsApiClient, readSignal: AbortSignal) => { + const turns = await readClient.turns(sessionId, readSignal, baselineTurnId); + readSignal.throwIfAborted(); + for (const turn of turns) { + coordinatorTurnIds.add(turn.id); + excludedTurnIds.delete(turn.id); + } + const latest = turns.at(-1); + latestInputTurnId = latest?.id; + rootTurn = latest && isTerminalTurn(latest.status) ? latest : undefined; + turnFailure = + latest?.status === "failed" + ? new AgentsApiError(latest.error?.message ?? "Agents API turn failed", latest.error ?? {}) + : undefined; + cancelled = latest?.status === "cancelled"; + return turns; + }; + const readItemsByTurn = async (readClient: AgentsApiClient, readSignal: AbortSignal) => { + const savedItems = await readClient.items(sessionId, undefined, readSignal); + readSignal.throwIfAborted(); + const itemsByTurn = new Map(); + for (const item of savedItems) { + if (!item.turn_id) { + continue; + } + const items = itemsByTurn.get(item.turn_id) ?? []; + items.push(item); + itemsByTurn.set(item.turn_id, items); + } + return itemsByTurn; + }; + const readSavedState = async (readClient: AgentsApiClient, readSignal: AbortSignal) => { + const turns = await readAdmittedTurns(readClient, readSignal); + const entries: Array<{ turn: Turn; items: AgentsApiItem[] }> = []; + const inputItems = new Set(); + const itemsByTurn = + turns.length || baselineTurnId + ? await readItemsByTurn(readClient, readSignal) + : new Map(); + for (const turn of turns) { + const items = itemsByTurn.get(turn.id) ?? []; + for (const item of items) { + if (item.turn_id && item.turn_id !== turn.id) { + throw new Error("Agents API saved item belongs to a different turn"); + } + rememberItemTurn(item.id, turn.id); + if (item.type === "message" && item.role === "user") { + inputItems.add(item.id); + } + } + entries.push({ turn, items }); + } + observedInputItems.clear(); + for (const id of inputItems) { + observedInputItems.add(id); + } + return { turns, entries, itemsByTurn }; + }; + const projectSavedState = async ( + entries: Array<{ turn: Turn; items: AgentsApiItem[] }>, + readSignal: AbortSignal, + ) => { + let transcriptReady = true; + for (const { turn, items } of entries) { + readSignal.throwIfAborted(); + const ready = await options.onReconcile?.(turn, items); + transcriptReady = ready !== false && transcriptReady; + readSignal.throwIfAborted(); + } + return transcriptReady; + }; + const reconcilePriorHistory = async ( + readClient: AgentsApiClient, + readSignal: AbortSignal, + itemsByTurn: Map, + ) => { + if (!baselineTurnId || !options.onReconcileHistory) { + return; + } + const turns = await readClient.turns(sessionId, readSignal); + readSignal.throwIfAborted(); + const baselineIndex = turns.findIndex((turn) => turn.id === baselineTurnId); + if (baselineIndex < 0) { + throw new Error("Agents API historical reconciliation lost its baseline turn"); + } + const priorTurns = turns + .slice(0, baselineIndex + 1) + .filter((turn) => isTerminalTurn(turn.status)); + if (!priorTurns.length) { + return; + } + // Historical facts repair the retained conversation without entering this + // attempt's admission, live presentation, tool lifecycle, or token accounting. + await options.onReconcileHistory( + priorTurns.map((turn) => ({ + turn, + items: itemsByTurn.get(turn.id) ?? [], + })), + ); + readSignal.throwIfAborted(); + }; + const rememberItemTurn = (itemId: string, turnId: string) => { + const previous = itemTurnIds.get(itemId); + if (previous && previous !== turnId) { + throw new Error("Agents API item belongs to different admitted turns"); + } + itemTurnIds.set(itemId, turnId); + excludedItemIds.delete(itemId); + }; return { isAvailable, @@ -139,101 +254,175 @@ export function createAgentsApiSession(options: { async run(prompt: string, persistInput: () => Promise, onSubmitted: () => void) { signal.throwIfAborted(); baselineTurnId = (await client.turns(sessionId, signal, undefined, true))[0]?.id; + baselineCaptured = true; + if (baselineTurnId) { + excludedTurnIds.add(baselineTurnId); + } let events = await client.subscribe( sessionId, AbortSignal.any([signal, streamController.signal]), ); let nextEvent = events.next(); void nextEvent.catch(() => {}); - let reconciledStream = false; + const settleFromSavedState = async (recover = false): Promise => { + const submissionFence = submission; + await submissionFence; + assertCurrent(); + const admittedCount = admittedMessageCount; + const snapshot = await readSavedState(client, signal); + assertCurrent(); + const session = await client.session(sessionId, signal); + assertCurrent(); + assertSessionUsable(session); + settled = Boolean( + rootTurn && + rootTurn.id === snapshot.turns.at(-1)?.id && + session.status === "idle" && + submissionFence === submission && + admittedCount === admittedMessageCount && + observedInputItems.size === admittedCount, + ); + if (settled) { + options.onSettled?.(); + } + if (settled || recover) { + if (settled) { + await reconcilePriorHistory(client, signal, snapshot.itemsByTurn); + assertCurrent(); + } + await projectSavedState(snapshot.entries, signal); + assertCurrent(); + } + }; + const belongsToAttempt = async (event: AgentsApiEvent) => { + if (event.item && event.item_id && event.item.id !== event.item_id) { + throw new Error("Agents API event contains different item identities"); + } + const itemId = event.item?.id ?? event.item_id; + const turnIds = new Set( + [ + event.turn?.id, + event.turn_id, + event.item?.turn_id, + itemId ? itemTurnIds.get(itemId) : undefined, + ].filter((id): id is string => typeof id === "string" && id.length > 0), + ); + if (turnIds.size > 1) { + throw new Error("Agents API event contains different turn identities"); + } + const turnId = turnIds.values().next().value; + if (!turnId) { + if (itemId) { + if (excludedItemIds.has(itemId)) { + return false; + } + // Command deltas can omit their turn; saved admitted items retain that correlation. + const snapshot = await readSavedState(client, signal); + assertCurrent(); + if (!itemTurnIds.has(itemId)) { + excludedItemIds.add(itemId); + return false; + } + await projectSavedState(snapshot.entries, signal); + assertCurrent(); + return true; + } + if (event.item || event.type === "agent.output.command_execution_output.delta") { + throw new Error("Agents API output event is missing its turn ID"); + } + return true; + } + if (excludedTurnIds.has(turnId)) { + if (itemId) { + excludedItemIds.add(itemId); + } + return false; + } + if (!coordinatorTurnIds.has(turnId)) { + // A delayed same-session event is not proof that this attempt admitted its turn. + await readAdmittedTurns(client, signal); + assertCurrent(); + if (!coordinatorTurnIds.has(turnId)) { + excludedTurnIds.add(turnId); + if (itemId) { + excludedItemIds.add(itemId); + } + return false; + } + } + if (event.item) { + rememberItemTurn(event.item.id, turnId); + } + return true; + }; try { await persistInput(); assertCurrent(); signal.throwIfAborted(); await submit(prompt); onSubmitted(); - while (!settled) { - const chunk = await nextEvent; + while (true) { + if (settled) { + break; + } + let chunk: IteratorResult; + try { + chunk = await nextEvent; + } catch (error) { + signal.throwIfAborted(); + assertCurrent(); + if (!isAgentsApiTransportDisconnect(error)) { + throw error; + } + chunk = { done: true, value: undefined }; + } if (chunk.done) { - reconciledStream = true; streamController.abort(); - await events.return(undefined); - await delay(500, undefined, { signal }); - streamController = new AbortController(); - // Subscribe before reconciliation: Agents API streams do not replay. - events = await client.subscribe( - sessionId, - AbortSignal.any([signal, streamController.signal]), - ); + // A broken reader can reject return() as well as next(). Retire only + // this transport; the admitted native work remains in the session. + await events.return(undefined).catch((error: unknown) => { + if (!isAgentsApiTransportDisconnect(error)) { + throw error; + } + }); + while (true) { + await delay(500, undefined, { signal }); + assertCurrent(); + streamController = new AbortController(); + try { + // Subscribe before reconciliation: Agents API streams do not replay. + events = await client.subscribe( + sessionId, + AbortSignal.any([signal, streamController.signal]), + ); + break; + } catch (error) { + streamController.abort(); + signal.throwIfAborted(); + assertCurrent(); + if (!isAgentsApiTransportDisconnect(error)) { + throw error; + } + } + } nextEvent = events.next(); void nextEvent.catch(() => {}); - await submission; - const admittedCount = admittedMessageCount; - const turns = await client.turns(sessionId, signal, baselineTurnId); - for (const turn of turns) { - coordinatorTurnIds.add(turn.id); - await collectInputs(turn.id); - } - const latestTurn = turns.at(-1); - if (latestTurn) { - latestInputTurnId = latestTurn.id; - rootTurn = ["completed", "failed", "cancelled"].includes(latestTurn.status) - ? latestTurn - : undefined; - turnFailure = - latestTurn.status === "failed" - ? (latestTurn.error?.message ?? "Agents API turn failed") - : undefined; - cancelled = latestTurn.status === "cancelled"; - } - const session = await client.session(sessionId, signal); - assertCurrent(); - assertSessionUsable(session); - settled = Boolean( - rootTurn && - session.status === "idle" && - admittedCount === admittedMessageCount && - observedInputItems.size >= admittedMessageCount, - ); + await settleFromSavedState(true); continue; } const event = chunk.value; nextEvent = events.next(); void nextEvent.catch(() => {}); assertCurrent(); - options.onEvent(event); - if (rootTurn && event.type === "agent.session.idle") { - await submission; - assertCurrent(); - if (observedInputItems.size < admittedMessageCount) { - for (const turnId of coordinatorTurnIds) { - await collectInputs(turnId); - } - } - if ( - observedInputItems.size < admittedMessageCount || - rootTurn.id !== latestInputTurnId - ) { - continue; - } - if (reconciledStream) { - const session = await client.session(sessionId, signal); - assertSessionUsable(session); - if (session.status !== "idle") { - continue; - } - } - settled = true; - break; + if (!(await belongsToAttempt(event))) { + continue; } - if (event.type === "agent.session.turn.created" && event.turn?.subagent_id === null) { - if (!coordinatorTurnIds.has(event.turn.id)) { - coordinatorTurnIds.add(event.turn.id); - latestInputTurnId = event.turn.id; - rootTurn = undefined; - turnFailure = undefined; - cancelled = false; - } + assertCurrent(); + await options.onEvent(event); + assertCurrent(); + if (event.type === "agent.session.idle") { + await settleFromSavedState(); + continue; } if ( (event.type === "agent.session.turn.item.added" || @@ -244,42 +433,48 @@ export function createAgentsApiSession(options: { !observedInputItems.has(event.item.id) ) { observedInputItems.add(event.item.id); - const inputTurnId = event.item.turn_id ?? event.turn_id; + const inputTurnId = + event.item.turn_id ?? event.turn_id ?? itemTurnIds.get(event.item.id); if (!inputTurnId) { throw new Error("Agents API input item is missing its turn ID"); } - if (!latestInputTurnId) { - coordinatorTurnIds.add(inputTurnId); - latestInputTurnId = inputTurnId; - } } if (event.type === "error") { - throw new Error(event.error?.message ?? "Agents API stream error"); + throw new AgentsApiError( + event.error?.message ?? "Agents API stream error", + event.error, + ); } - if ( - [ - "agent.session.failed", - "agent.session.environment.failed", - "agent.session.requires_action", - ].includes(event.type) - ) { - throw new Error(`Agents API MVP cannot continue: ${event.type}`); + if (event.type === "agent.session.requires_action") { + throw new Error("Agents API MVP cannot continue: agent.session.requires_action"); + } + if (["agent.session.failed", "agent.session.environment.failed"].includes(event.type)) { + const nativeError = event.environment?.error ?? event.error; + throw new AgentsApiError( + nativeError?.message ?? + event.session?.error ?? + `Agents API cannot continue: ${event.type}`, + nativeError ?? {}, + ); } if ( (event.type === "agent.session.turn.completed" || event.type === "agent.session.turn.failed" || event.type === "agent.session.turn.cancelled") && - event.turn.subagent_id === null && + event.turn?.subagent_id === null && event.turn.id === latestInputTurnId ) { rootTurn = event.turn; turnFailure = event.type.endsWith(".failed") - ? (event.turn.error?.message ?? "Agents API turn failed") + ? new AgentsApiError( + event.turn.error?.message ?? "Agents API turn failed", + event.turn.error ?? {}, + ) : undefined; cancelled = event.type.endsWith(".cancelled"); + await settleFromSavedState(); } } - options.onSettled?.(); } finally { streamController.abort(); await events.return(undefined); @@ -290,10 +485,28 @@ export function createAgentsApiSession(options: { ); } if (turnFailure) { - throw new Error(turnFailure); + throw turnFailure; } return { turn: rootTurn, cancelled }; }, + async reconcileAfterClose(cleanupSignal: AbortSignal): Promise { + if (!closed) { + throw new Error("Agents API canonical cleanup requires a closed session attempt"); + } + if (!submitted || !baselineCaptured) { + return undefined; + } + cleanupSignal.throwIfAborted(); + const session = await cleanupClient.session(sessionId, cleanupSignal); + cleanupSignal.throwIfAborted(); + if (session.status !== "idle" && session.status !== "failed") { + throw new Error("Agents API canonical cleanup requires native work to be retired"); + } + const snapshot = await readSavedState(cleanupClient, cleanupSignal); + await reconcilePriorHistory(cleanupClient, cleanupSignal, snapshot.itemsByTurn); + await projectSavedState(snapshot.entries, cleanupSignal); + return snapshot.turns.at(-1); + }, async close() { signal.removeEventListener("abort", onAbort); streamController.abort(); @@ -302,6 +515,26 @@ export function createAgentsApiSession(options: { } await cancellation; stopped = true; + closed = true; }, }; } + +function isTerminalTurn(status: string): boolean { + return ["completed", "failed", "cancelled"].includes(status); +} + +function isAgentsApiTransportDisconnect(error: unknown): boolean { + if (!(error instanceof Error) || error instanceof AgentsApiError) { + return false; + } + const code = asOptionalRecord(error)?.code; + if ( + (typeof code === "string" && + ["ECONNRESET", "ECONNREFUSED", "ETIMEDOUT", "UND_ERR_SOCKET"].includes(code)) || + (error instanceof TypeError && ["terminated", "fetch failed"].includes(error.message)) + ) { + return true; + } + return error.cause instanceof Error && isAgentsApiTransportDisconnect(error.cause); +} diff --git a/extensions/agentsapi/agentsapi-transcript.ts b/extensions/agentsapi/agentsapi-transcript.ts new file mode 100644 index 000000000000..2dbeb4ea0dc5 --- /dev/null +++ b/extensions/agentsapi/agentsapi-transcript.ts @@ -0,0 +1,285 @@ +import { + createAgentHarnessToolCallMessage, + createAgentHarnessToolResultMessage, +} from "openclaw/plugin-sdk/agent-harness-attempt-runtime"; +import type { + AgentHarnessAttemptParamsV2, + AgentMessage, +} from "openclaw/plugin-sdk/agent-harness-runtime"; +import { appendSessionTranscriptMessageByIdentityStrict } from "openclaw/plugin-sdk/session-transcript-runtime"; +import type { AgentsApiItem } from "./agentsapi-client.js"; +import { + agentsApiNativeTool, + agentsApiNativeToolDetails, + agentsApiNativeToolOutcome, + agentsApiNativeToolOutput, +} from "./agentsapi-native-items.js"; + +/** Canonical native facts use the same durable identities during live and historical repair. */ +export async function recordAgentsApiNativeToolTranscript( + params: AgentHarnessAttemptParamsV2, + sessionId: string, + turnId: string, + item: AgentsApiItem, + assertCurrent: () => void, + nextTimestamp: () => number, + options: { + enclosingStatus?: string; + capturedOutput?: string; + captureTruncated?: boolean; + } = {}, +): Promise { + assertCurrent(); + const tool = agentsApiNativeTool(item, params); + if (!tool || !["completed", "failed", "incomplete"].includes(item.status ?? "")) { + // A failed parent turn can retire before its command completes. Do not + // freeze a provisional result under the command's durable identity. + return false; + } + const id = `agentsapi:${sessionId}:${turnId}:${item.id}`; + const outcome = agentsApiNativeToolOutcome(item, options.enclosingStatus); + const output = agentsApiNativeToolOutput(item, options.capturedOutput); + const details = agentsApiNativeToolDetails( + sessionId, + turnId, + item, + outcome, + options.capturedOutput, + ); + const text = + output ?? + (item.type === "web_search_call" + ? `Web search ${outcome.status}; native search results are unavailable.` + : (outcome.error ?? `${tool.name} ${outcome.status}`)); + await recordAgentsApiNativeToolInvocation( + params, + sessionId, + turnId, + item, + assertCurrent, + nextTimestamp, + ); + await appendAgentsApiTranscriptMessage( + params, + { + ...createAgentHarnessToolResultMessage( + { id, name: tool.name, text, isError: outcome.isError, details }, + nextTimestamp(), + ), + __openclaw: { + toolOutput: { + source: "execution", + modelInput: "unverified", + ...(outcome.outcomeUnknown ? { outcome: "unknown" } : {}), + ...(options.captureTruncated ? { captureTruncated: true } : {}), + }, + ...(item.type === "web_search_call" ? { resultContentSource: "network" } : {}), + }, + idempotencyKey: `${id}:result`, + }, + assertCurrent, + ); + return true; +} + +/** Canonical invocation facts can precede completion of the native tool. */ +export async function recordAgentsApiNativeToolInvocation( + params: AgentHarnessAttemptParamsV2, + sessionId: string, + turnId: string, + item: AgentsApiItem, + assertCurrent: () => void, + nextTimestamp: () => number, +): Promise { + assertCurrent(); + const tool = agentsApiNativeTool(item, params); + if (!tool || !canRecordAgentsApiNativeToolInvocation(item)) { + return false; + } + const id = `agentsapi:${sessionId}:${turnId}:${item.id}`; + await appendAgentsApiTranscriptMessage( + params, + { + ...createAgentHarnessToolCallMessage( + { api: "openai-responses", provider: "openai", modelId: params.model.id }, + { id, name: tool.name, arguments: tool.args }, + nextTimestamp(), + ), + idempotencyKey: `${id}:call`, + }, + assertCurrent, + ); + return true; +} + +export async function appendAgentsApiTranscriptMessage( + params: AgentHarnessAttemptParamsV2, + message: TMessage, + assertCurrent: () => void, +): Promise { + assertCurrent(); + const { agentId, sessionId, sessionKey, storePath } = params.sessionTarget ?? {}; + if ( + !agentId || + !sessionId || + !sessionKey || + !storePath || + sessionId !== params.sessionId || + agentId !== params.agentId || + sessionKey !== params.sessionKey + ) { + throw new Error("Agents API requires a matching host-prepared session target"); + } + const append = await appendSessionTranscriptMessageByIdentityStrict({ + ...params.sessionTarget, + agentId, + sessionId, + sessionKey, + storePath, + config: params.config, + message, + prepareMessageAfterIdempotencyCheck: (prepared) => { + assertCurrent(); + return prepared; + }, + }); + assertCurrent(); + if (append.kind !== "result") { + throw new Error("Agents API transcript append was refused"); + } + return append.result.message; +} + +/** Walk the retrieved native prefix without another transcript queue or cursor. */ +export function* iterateAgentsApiTranscriptItems( + turnId: string, + items: readonly AgentsApiItem[], + enclosingStatus: string | undefined, + terminalTurn: boolean, + recordedGatewayCallIds: ReadonlySet, + hasObservedCompletion: (itemId: string) => boolean, +): Generator<{ item: AgentsApiItem; terminal: boolean; transcriptReady: boolean }> { + let transcriptReady = true; + let deferredFinalSeen = false; + for (const item of items) { + if (item.turn_id && item.turn_id !== turnId) { + throw new Error("Agents API saved item belongs to a different turn"); + } + const completionObserved = hasObservedCompletion(item.id); + const terminal = + terminalTurn || + ["completed", "failed", "incomplete"].includes(item.status ?? "") || + (item.status == null && completionObserved); + if ( + transcriptReady && + ((deferredFinalSeen && hasAgentsApiTranscriptRecord(item)) || + !canRecordAgentsApiTranscriptItem( + turnId, + item, + enclosingStatus, + recordedGatewayCallIds, + completionObserved, + )) + ) { + transcriptReady = false; + } + // The host's aggregate final is published at settlement. Later transcript + // slots cannot occupy their exact canonical position around that final. + deferredFinalSeen ||= isAgentsApiDeferredFinalText(item); + yield { item, terminal, transcriptReady }; + } +} + +/** Nonterminal calls require retrieved invocation fields, rather than streamed guesses. */ +function canRecordAgentsApiNativeToolInvocation(item: AgentsApiItem): boolean { + if (["completed", "failed", "incomplete"].includes(item.status ?? "")) { + return ["command_execution", "mcp_call", "web_search_call"].includes(item.type); + } + if (item.type === "command_execution") { + return typeof item.command === "string" && (item.cwd === null || typeof item.cwd === "string"); + } + // A retrieved MCP arguments field does not prove that generation is complete. + // Wait for this item's terminal status before freezing its invocation. + return false; +} + +/** Backend publication stops before invocation or text facts that are still incomplete. */ +function canRecordAgentsApiTranscriptItem( + turnId: string, + item: AgentsApiItem, + enclosingStatus: string | undefined, + recordedGatewayCallIds: ReadonlySet, + completionObserved = false, +): boolean { + if (["command_execution", "mcp_call", "web_search_call"].includes(item.type)) { + return canRecordAgentsApiNativeToolInvocation(item); + } + if (item.type === "function_call") { + return ( + typeof item.call_id === "string" && recordedGatewayCallIds.has(`${turnId}:${item.call_id}`) + ); + } + if ( + item.type === "reasoning" || + (item.type === "message" && item.role === "assistant" && item.phase === "commentary") + ) { + return canRecordAgentsApiTranscriptText(item, enclosingStatus, completionObserved); + } + return true; +} + +export function canRecordAgentsApiTranscriptText( + item: AgentsApiItem, + enclosingStatus?: string, + completionObserved = false, +): boolean { + if (["completed", "failed", "incomplete"].includes(item.status ?? "")) { + return true; + } + // A nullable status can use an observed native item completion. Recovery also + // permits the completed coordinator's reasoning snapshot; explicitly running + // items remain provisional even when their current summaries are empty. + return ( + item.status == null && + (completionObserved || (item.type === "reasoning" && enclosingStatus === "completed")) + ); +} + +function isAgentsApiDeferredFinalText(item: AgentsApiItem): boolean { + return ( + item.type === "message" && + item.role === "assistant" && + item.phase !== "commentary" && + item.status === "completed" && + Boolean(item.content?.some((part) => part.type === "output_text" && part.text)) + ); +} + +function hasAgentsApiTranscriptRecord(item: AgentsApiItem): boolean { + return ( + ["command_execution", "mcp_call", "web_search_call", "function_call"].includes(item.type) || + (item.type === "reasoning" && + Boolean(item.summary?.some((part) => part.type === "summary_text" && part.text))) || + (item.type === "message" && + item.role === "assistant" && + item.phase === "commentary" && + Boolean(item.content?.some((part) => part.type === "output_text" && part.text))) + ); +} + +export function readTextParts(parts: AgentsApiItem["content"], type: string): Map { + const texts = new Map(); + parts?.forEach((part, index) => { + if (part.type === type) { + texts.set(index, part.text ?? ""); + } + }); + return texts; +} + +export function joinTextParts(parts: Map): string { + return [...parts.entries()] + .toSorted(([left], [right]) => left - right) + .map(([, text]) => text) + .join(""); +}