diff --git a/packages/cli/src/acp/event.ts b/packages/cli/src/acp/event.ts index 11bf4c11c5..277f107a11 100644 --- a/packages/cli/src/acp/event.ts +++ b/packages/cli/src/acp/event.ts @@ -198,8 +198,8 @@ export async function streamTurn(input: { toolCallId: event.data.callID, toolName: current.name, input: current.input, - structured: current.structured, - content: current.content, + structured: event.data.structured ?? current.structured, + content: event.data.content ?? current.content, error: event.data.error.message, cwd: input.cwd, }), diff --git a/packages/cli/src/run/noninteractive.ts b/packages/cli/src/run/noninteractive.ts index 102294f64b..59fda40e3b 100644 --- a/packages/cli/src/run/noninteractive.ts +++ b/packages/cli/src/run/noninteractive.ts @@ -401,8 +401,8 @@ export async function runNonInteractivePrompt(input: Input) { state: { status: "error", input: current.input, - structured: current.structured, - content: current.content, + structured: event.data.structured ?? current.structured, + content: event.data.content ?? current.content, error: event.data.error, result: event.data.result, }, diff --git a/packages/cli/test/acp/event-behavior.test.ts b/packages/cli/test/acp/event-behavior.test.ts index 9de17cf4e1..3c63f063a4 100644 --- a/packages/cli/test/acp/event-behavior.test.ts +++ b/packages/cli/test/acp/event-behavior.test.ts @@ -214,7 +214,7 @@ describe("acp event behavior", () => { }), ) send( - durableEvent("session.tool.progress", { + ephemeralEvent("session.tool.progress", { sessionID: "ses_tools", assistantMessageID: "msg_tools", callID: "call_ok", @@ -251,7 +251,7 @@ describe("acp event behavior", () => { }), ) send( - durableEvent("session.tool.progress", { + ephemeralEvent("session.tool.progress", { sessionID: "ses_tools", assistantMessageID: "msg_tools", callID: "call_fail", @@ -265,6 +265,8 @@ describe("acp event behavior", () => { assistantMessageID: "msg_tools", callID: "call_fail", error: { type: "tool.error", message: "not found" }, + structured: { bytes: 0 }, + content: [{ type: "text", text: "opening" }], executed: true, }), ) diff --git a/packages/cli/test/run/noninteractive.test.ts b/packages/cli/test/run/noninteractive.test.ts index f27869a991..dc0559acfe 100644 --- a/packages/cli/test/run/noninteractive.test.ts +++ b/packages/cli/test/run/noninteractive.test.ts @@ -130,7 +130,6 @@ function failedTool(inputID: string): V2Event[] { id: "evt_failed_tool_progress", created: 3, type: "session.tool.progress", - durable: { aggregateID: "ses_1", seq: 3, version: 1 }, data: { sessionID: "ses_1", assistantMessageID: "msg_failed_tool", @@ -149,6 +148,8 @@ function failedTool(inputID: string): V2Event[] { assistantMessageID: "msg_failed_tool", callID: "call_failed_tool", error: { type: "unknown", message: "tool failed" }, + structured: { checkpoint: 1 }, + content: [{ type: "text", text: "partial output" }], executed: true, }, }, diff --git a/packages/client/src/promise/generated/types.ts b/packages/client/src/promise/generated/types.ts index 6664c1424c..efbb4e8c44 100644 --- a/packages/client/src/promise/generated/types.ts +++ b/packages/client/src/promise/generated/types.ts @@ -1385,24 +1385,6 @@ export type SessionToolCalled = { } } -export type SessionToolFailed = { - id: string - created: number - metadata?: { [x: string]: any } - type: "session.tool.failed" - durable: { aggregateID: string; seq: number; version: 1 } - location?: LocationRef - data: { - sessionID: string - assistantMessageID: string - callID: string - error: SessionStructuredError - result?: any - executed: boolean - resultState?: SessionMessageProviderState7 - } -} - export type ModelCost = { tier?: { type: "context"; size: number } input: MoneyUSDPerMillionTokens @@ -1847,22 +1829,6 @@ export type SessionMessageToolStateError = { result?: JsonValue } -export type SessionToolProgress = { - id: string - created: number - metadata?: { [x: string]: any } - type: "session.tool.progress" - durable: { aggregateID: string; seq: number; version: 1 } - location?: LocationRef - data: { - sessionID: string - assistantMessageID: string - callID: string - structured: { [x: string]: any } - content: Array - } -} - export type SessionToolSuccess = { id: string created: number @@ -1882,6 +1848,41 @@ export type SessionToolSuccess = { } } +export type SessionToolFailed = { + id: string + created: number + metadata?: { [x: string]: any } + type: "session.tool.failed" + durable: { aggregateID: string; seq: number; version: 1 } + location?: LocationRef + data: { + sessionID: string + assistantMessageID: string + callID: string + error: SessionStructuredError + structured?: { [x: string]: any } + content?: Array + result?: any + executed: boolean + resultState?: SessionMessageProviderState7 + } +} + +export type SessionToolProgress = { + id: string + created: number + metadata?: { [x: string]: any } + type: "session.tool.progress" + location?: LocationRef + data: { + sessionID: string + assistantMessageID: string + callID: string + structured: { [x: string]: any } + content: Array + } +} + export type SessionMessageCompaction = | SessionMessageCompactionRunning | SessionMessageCompactionCompleted @@ -2219,7 +2220,6 @@ export type SessionEventDurable = | SessionToolInputStarted | SessionToolInputEnded | SessionToolCalled - | SessionToolProgress | SessionToolSuccess | SessionToolFailed | SessionRetryScheduled diff --git a/packages/core/src/database/migration.gen.ts b/packages/core/src/database/migration.gen.ts index c4067c4765..dd6214b095 100644 --- a/packages/core/src/database/migration.gen.ts +++ b/packages/core/src/database/migration.gen.ts @@ -55,5 +55,6 @@ export const migrations = ( import("./migration/20260709190621_session_pending_table"), import("./migration/20260710025429_instruction_sync"), import("./migration/20260716020354_kv"), + import("./migration/20260722011141_delete_tool_progress_events"), ]) ).map((module) => module.default) satisfies DatabaseMigration.Migration[] diff --git a/packages/core/src/database/migration/20260722011141_delete_tool_progress_events.ts b/packages/core/src/database/migration/20260722011141_delete_tool_progress_events.ts new file mode 100644 index 0000000000..afb8f44334 --- /dev/null +++ b/packages/core/src/database/migration/20260722011141_delete_tool_progress_events.ts @@ -0,0 +1,11 @@ +import { Effect } from "effect" +import type { DatabaseMigration } from "../migration" + +export default { + id: "20260722011141_delete_tool_progress_events", + up(tx) { + return Effect.gen(function* () { + yield* tx.run(`DELETE FROM \`event\` WHERE \`type\` = 'session.tool.progress.1';`) + }) + }, +} satisfies DatabaseMigration.Migration diff --git a/packages/core/src/session/message-updater.ts b/packages/core/src/session/message-updater.ts index 00b0927cf4..fa7b86a773 100644 --- a/packages/core/src/session/message-updater.ts +++ b/packages/core/src/session/message-updater.ts @@ -402,8 +402,8 @@ export function update(adapter: Adapter, event: SessionEvent.Event) { status: "error", error: event.data.error, input: typeof match.state.input === "string" ? {} : match.state.input, - structured: match.state.status === "running" ? match.state.structured : {}, - content: match.state.status === "running" ? match.state.content : [], + structured: event.data.structured ?? (match.state.status === "running" ? match.state.structured : {}), + content: event.data.content ?? (match.state.status === "running" ? match.state.content : []), result: event.data.result, }), ) diff --git a/packages/core/src/session/projector.ts b/packages/core/src/session/projector.ts index 97044675d1..f759a3757f 100644 --- a/packages/core/src/session/projector.ts +++ b/packages/core/src/session/projector.ts @@ -697,7 +697,6 @@ const layer = Layer.effectDiscard( yield* events.project(SessionEvent.Tool.Input.Started, (event) => run(db, event)) yield* events.project(SessionEvent.Tool.Input.Ended, (event) => run(db, event)) yield* events.project(SessionEvent.Tool.Called, (event) => run(db, event)) - yield* events.project(SessionEvent.Tool.Progress, (event) => run(db, event)) yield* events.project(SessionEvent.Tool.Success, (event) => run(db, event)) yield* events.project(SessionEvent.Tool.Failed, (event) => run(db, event)) yield* events.project(SessionEvent.Reasoning.Started, (event) => run(db, event)) diff --git a/packages/core/src/session/runner/llm.ts b/packages/core/src/session/runner/llm.ts index 2a8aee08cc..8c38de47bb 100644 --- a/packages/core/src/session/runner/llm.ts +++ b/packages/core/src/session/runner/llm.ts @@ -62,6 +62,8 @@ const layer = Layer.effect( assistantMessageID: message.id, callID: tool.id, error: { type: "aborted", message: `Tool execution interrupted: ${tool.name}` }, + structured: tool.state.status === "running" ? tool.state.structured : {}, + content: tool.state.status === "running" ? tool.state.content : [], executed: tool.executed === true, }) } @@ -163,16 +165,7 @@ const layer = Layer.effect( agent: agent.id, messageID: assistantMessageID, call: event, - progress: (update) => - serialized( - events.publish(SessionEvent.Tool.Progress, { - sessionID: session.id, - assistantMessageID, - callID: event.id, - structured: { ...update.structured }, - content: [...update.content], - }), - ), + progress: (update) => serialized(publisher.progress(event.id, update)), }), ).pipe( Effect.flatMap((settlement) => diff --git a/packages/core/src/session/runner/publish-llm-event.ts b/packages/core/src/session/runner/publish-llm-event.ts index 1344d2627a..0f9e0586dc 100644 --- a/packages/core/src/session/runner/publish-llm-event.ts +++ b/packages/core/src/session/runner/publish-llm-event.ts @@ -11,6 +11,7 @@ import { AgentV2 } from "../../agent" import { Snapshot } from "../../snapshot" import { RelativePath } from "../../schema" import { SessionUsage } from "../usage" +import type { ToolRegistry } from "../../tool/registry" type Input = { readonly sessionID: SessionSchema.ID @@ -54,8 +55,13 @@ export const createLLMEventPublisher = (events: Pick() + const failureState = (tool: { readonly progress?: ToolRegistry.Progress }) => ({ + structured: tool.progress?.structured ?? {}, + content: tool.progress?.content ?? [], + }) let assistantMessageID = input.assistantMessageID let stepStarted = false let stepFailed = false @@ -232,6 +238,7 @@ export const createLLMEventPublisher = (events: Pick - context.progress({ - structured: { truncated: capture.truncated }, - content: [{ type: "text", text: capture.output }], + Effect.gen(function* () { + if ( + previousProgress?.output === capture.output && + previousProgress.truncated === capture.truncated + ) + return + previousProgress = capture + yield* context.progress({ + structured: { truncated: capture.truncated }, + content: [{ type: "text", text: capture.output }], + }) }), ), ), diff --git a/packages/core/test/database-migration.test.ts b/packages/core/test/database-migration.test.ts index 2c68d2e186..58673f2039 100644 --- a/packages/core/test/database-migration.test.ts +++ b/packages/core/test/database-migration.test.ts @@ -24,6 +24,7 @@ import renameInstructionsMigration from "@opencode-ai/core/database/migration/20 import addSessionForkMigration from "@opencode-ai/core/database/migration/20260706223930_add-session-fork" import timeSuspendedMigration from "@opencode-ai/core/database/migration/20260709163752_time_suspended" import instructionSyncMigration from "@opencode-ai/core/database/migration/20260710025429_instruction_sync" +import deleteToolProgressEventsMigration from "@opencode-ai/core/database/migration/20260722011141_delete_tool_progress_events" import { AppNodeBuilder } from "@opencode-ai/core/effect/app-node-builder" import { LayerNode } from "@opencode-ai/util/effect/layer-node" import { EventV2 } from "@opencode-ai/core/event" @@ -557,6 +558,31 @@ describe("DatabaseMigration", () => { ) }) + test("deletes durable tool progress without changing aggregate sequence watermarks", async () => { + await run( + Effect.gen(function* () { + const db = yield* makeDb + yield* db.run(sql`CREATE TABLE event_sequence (aggregate_id text PRIMARY KEY, seq integer NOT NULL)`) + yield* db.run( + sql`CREATE TABLE event (id text PRIMARY KEY, aggregate_id text NOT NULL, seq integer NOT NULL, type text NOT NULL, data text NOT NULL)`, + ) + yield* db.run(sql`INSERT INTO event_sequence VALUES ('ses_test', 5)`) + yield* db.run(sql`INSERT INTO event VALUES ('evt_success', 'ses_test', 4, 'session.tool.success.1', '{}')`) + yield* db.run(sql`INSERT INTO event VALUES ('evt_progress', 'ses_test', 5, 'session.tool.progress.1', '{}')`) + + yield* DatabaseMigration.applyOnly(db, [deleteToolProgressEventsMigration]) + + expect(yield* db.all(sql`SELECT id, seq, type, data FROM event ORDER BY seq`)).toEqual([ + { id: "evt_success", seq: 4, type: "session.tool.success.1", data: "{}" }, + ]) + expect(yield* db.get(sql`SELECT aggregate_id, seq FROM event_sequence`)).toEqual({ + aggregate_id: "ses_test", + seq: 5, + }) + }), + ) + }) + test("records the authoritative parent sequence on existing forks", async () => { await run( Effect.gen(function* () { diff --git a/packages/core/test/session-runner-tool-events.test.ts b/packages/core/test/session-runner-tool-events.test.ts index b6f44dccdb..1b3ecaaac8 100644 --- a/packages/core/test/session-runner-tool-events.test.ts +++ b/packages/core/test/session-runner-tool-events.test.ts @@ -1,5 +1,5 @@ import { expect, test } from "bun:test" -import { Effect, Schema } from "effect" +import { Cause, Effect, Exit, Schema } from "effect" import { LLMEvent } from "@opencode-ai/ai" import { Money } from "@opencode-ai/schema/money" import { EventV2 } from "@opencode-ai/core/event" @@ -16,11 +16,11 @@ import { createLLMEventPublisher } from "@opencode-ai/core/session/runner/publis const sessionID = SessionV2.ID.make("ses_tool_event_test") const base64 = "iVBORw0KGgoAAAANSUhEUgAAAAEAAAAB" -const capture = (providerMetadataKey = "anthropic") => { +const capture = (providerMetadataKey = "anthropic", options?: { readonly interruptProgress?: boolean }) => { const published: Array<{ readonly type: string; readonly data: unknown }> = [] const events: Pick = { - publish: (definition, data) => - Effect.sync(() => { + publish: (definition, data) => { + const publish = Effect.sync(() => { const event = { id: EventV2.ID.create(), type: definition.type, data } as EventV2.Payload published.push({ type: definition.durable @@ -29,7 +29,11 @@ const capture = (providerMetadataKey = "anthropic") => { data, }) return event - }), + }) + return definition.type === SessionEvent.Tool.Progress.type && options?.interruptProgress + ? publish.pipe(Effect.andThen(Effect.interrupt)) + : publish + }, } return { published, @@ -92,6 +96,24 @@ test("provider-executed success retains its raw provider result", async () => { expect(success?.data).toHaveProperty("result") }) +test("interrupted progress publication remains in the terminal failure snapshot", async () => { + const { published, publisher } = capture("anthropic", { interruptProgress: true }) + await Effect.runPromise(publisher.publish(call)) + const exit = await Effect.runPromiseExit( + publisher.progress(call.id, { + structured: { phase: "visible" }, + content: [{ type: "text", text: "visible" }], + }), + ) + expect(Exit.isFailure(exit) && Cause.hasInterruptsOnly(exit.cause)).toBe(true) + await Effect.runPromise(publisher.failUnsettledTools({ type: "aborted", message: "interrupted" })) + + expect(published.find((event) => event.type === "session.tool.failed.1")?.data).toMatchObject({ + structured: { phase: "visible" }, + content: [{ type: "text", text: "visible" }], + }) +}) + test("provider metadata is flattened using the route key", async () => { const { published, publisher } = capture() await Effect.runPromise( diff --git a/packages/core/test/session-runner.test.ts b/packages/core/test/session-runner.test.ts index dc2cb2d9df..59bcd8a34b 100644 --- a/packages/core/test/session-runner.test.ts +++ b/packages/core/test/session-runner.test.ts @@ -892,6 +892,52 @@ describe("SessionRunnerLLM", () => { }), ) + it.effect("persists the final progress snapshot when a tool fails", () => + Effect.gen(function* () { + const session = yield* setup + const registry = yield* ToolRegistry.Service + yield* registry.register({ + failing_progress: Tool.make({ + description: "Report progress and fail", + input: Schema.Struct({}), + output: Schema.Struct({}), + execute: (_, context) => + Effect.gen(function* () { + yield* context.progress({ + structured: { phase: "running" }, + content: [{ type: "text", text: "before failure" }], + }) + return yield* new ToolFailure({ message: "failed after progress" }) + }), + }), + }, { codemode: false }) + yield* admit(session, "Run failing progress") + responses = [reply.tool("call-failing-progress", "failing_progress", {}), reply.stop()] + + yield* session.resume(sessionID) + + expect(yield* session.context(sessionID)).toMatchObject([ + { type: "user", text: "Run failing progress" }, + { + type: "assistant", + content: [ + { + type: "tool", + id: "call-failing-progress", + state: { + status: "error", + structured: { phase: "running" }, + content: [{ type: "text", text: "before failure" }], + error: { message: "failed after progress" }, + }, + }, + ], + }, + { type: "assistant", finish: "stop" }, + ]) + }), + ) + it.effect("executes the tool advertised before a registry reload", () => Effect.gen(function* () { const session = yield* setup diff --git a/packages/core/test/session-tool-progress.test.ts b/packages/core/test/session-tool-progress.test.ts index 0935ee3633..9e11a7aef3 100644 --- a/packages/core/test/session-tool-progress.test.ts +++ b/packages/core/test/session-tool-progress.test.ts @@ -25,7 +25,7 @@ const model = { id: ModelV2.ID.make("model"), providerID: ProviderV2.ID.make("pr const content = (text: string) => [{ type: "text" as const, text }] describe("Tool.Progress", () => { - it.effect("projects durable progress and keeps final settlements durable", () => + it.effect("keeps progress live-only and terminal settlements replayable", () => Effect.gen(function* () { const { db } = yield* Database.Service const service = yield* EventV2.Service @@ -87,7 +87,7 @@ describe("Tool.Progress", () => { state: { status: "running", structured: {}, content: [] }, }) - yield* service.publish(SessionEvent.Tool.Progress, { + const progress = yield* service.publish(SessionEvent.Tool.Progress, { sessionID, assistantMessageID, callID: "call-success", @@ -95,7 +95,7 @@ describe("Tool.Progress", () => { content: content("saved"), }) expect((yield* readAssistant).content[0]).toMatchObject({ - state: { status: "running", structured: { phase: "checkpoint" }, content: content("saved") }, + state: { status: "running", structured: {}, content: [] }, }) const success = yield* service.publish(SessionEvent.Tool.Success, { @@ -123,6 +123,8 @@ describe("Tool.Progress", () => { assistantMessageID, callID: "call-failed", error: { type: "unknown", message: "boom" }, + structured: { phase: "checkpoint" }, + content: content("before failure"), executed: false, }) expect((yield* readAssistant).content[1]).toMatchObject({ @@ -133,6 +135,7 @@ describe("Tool.Progress", () => { error: { type: "unknown", message: "boom" }, }, }) + expect(Schema.is(SessionEvent.Durable)(progress)).toBe(false) expect(Schema.is(SessionEvent.Durable)(success)).toBe(true) expect(Schema.is(SessionEvent.Durable)(failed)).toBe(true) @@ -143,9 +146,10 @@ describe("Tool.Progress", () => { .orderBy(asc(EventTable.seq)) .all() .pipe(Effect.orDie) - expect(rows.map((row) => row.type)).toContain(EventV2.versionedType(SessionEvent.Tool.Progress.type, 1)) + expect(rows.map((row) => row.type)).not.toContain(EventV2.versionedType(SessionEvent.Tool.Progress.type, 1)) expect(rows.map((row) => row.type)).toContain(EventV2.versionedType(SessionEvent.Tool.Success.type, 1)) expect(rows.map((row) => row.type)).toContain(EventV2.versionedType(SessionEvent.Tool.Failed.type, 1)) + }), ) }) diff --git a/packages/core/test/tool-shell.test.ts b/packages/core/test/tool-shell.test.ts index 2b66b2ec8d..8d79cde7a9 100644 --- a/packages/core/test/tool-shell.test.ts +++ b/packages/core/test/tool-shell.test.ts @@ -159,6 +159,9 @@ const mixedOutputCommand = isWindows ? "[Console]::Out.Write('stdout'); Start-Sleep -Milliseconds 50; [Console]::Error.Write('stderr'); Start-Sleep -Milliseconds 100" : "printf stdout; sleep 0.05; printf stderr >&2" const idleCommand = isWindows ? "Start-Sleep -Seconds 60" : "sleep 60" +const steadyProgressCommand = isWindows + ? "[Console]::Out.Write('steady'); Start-Sleep -Milliseconds 3400" + : "printf steady; sleep 3.4" const bodyExitCommand = isWindows ? "[Console]::Out.Write('body'); Start-Sleep -Milliseconds 100; exit 7" : "printf body && exit 7" @@ -462,6 +465,34 @@ describe("ShellTool", () => { { timeout: 15_000 }, ) + it.live( + "does not repeat unchanged shell progress", + () => + Effect.acquireUseRelease( + Effect.promise(() => tmpdir()), + (tmp) => { + reset() + return withSession(tmp.path, (registry) => + Effect.gen(function* () { + const updates: ToolRegistry.Progress[] = [] + yield* settleTool(registry, { + ...call({ command: steadyProgressCommand }, "call-steady-progress"), + progress: (update) => Effect.sync(() => updates.push(update)), + }) + expect(updates).toEqual([ + { + structured: { truncated: false }, + content: [{ type: "text", text: "steady" }], + }, + ]) + }), + ) + }, + (tmp) => Effect.promise(() => tmp[Symbol.asyncDispose]().then(() => undefined)), + ), + { timeout: 10_000 }, + ) + it.live("returns a useful timeout settlement", () => Effect.acquireUseRelease( Effect.promise(() => tmpdir()), diff --git a/packages/schema/src/session-event.ts b/packages/schema/src/session-event.ts index 730b0b1ffb..5c9c40a871 100644 --- a/packages/schema/src/session-event.ts +++ b/packages/schema/src/session-event.ts @@ -3,7 +3,6 @@ export * as SessionEvent from "./session-event.js" import { Schema } from "effect" import { optional } from "./schema.js" import { Event } from "./event.js" -import { ToolContent } from "./llm.js" import { FinishReason } from "./llm.js" import { Model } from "./model.js" import { NonNegativeInt, PositiveInt, RelativePath } from "./schema.js" @@ -410,18 +409,15 @@ export namespace Tool { }) export type Called = typeof Called.Type - /** - * Replayable bounded running-tool state. Tools should checkpoint semantic - * transitions or at a bounded cadence, not persist every stdout/stderr chunk. - */ - export const Progress = Event.durable({ + const ToolOutputFields = { + structured: SessionMessage.ToolStateRunning.fields.structured, + content: SessionMessage.ToolStateRunning.fields.content, + } + + /** Live replacement snapshot for a running tool. Terminal events own durable output. */ + export const Progress = Event.ephemeral({ type: "session.tool.progress", - ...options, - schema: { - ...ToolBase, - structured: Schema.Record(Schema.String, Schema.Unknown), - content: Schema.Array(ToolContent), - }, + schema: { ...ToolBase, ...ToolOutputFields }, }) export type Progress = typeof Progress.Type @@ -430,8 +426,7 @@ export namespace Tool { ...options, schema: { ...ToolBase, - structured: Schema.Record(Schema.String, Schema.Unknown), - content: Schema.Array(ToolContent), + ...ToolOutputFields, result: Schema.Unknown.pipe(optional), executed: Schema.Boolean, resultState: SessionMessage.ProviderState.pipe(optional), @@ -445,6 +440,9 @@ export namespace Tool { schema: { ...ToolBase, error: SessionError.Error, + // Optional only for compatibility with existing v1 failure rows. + structured: ToolOutputFields.structured.pipe(optional), + content: ToolOutputFields.content.pipe(optional), result: Schema.Unknown.pipe(optional), executed: Schema.Boolean, resultState: SessionMessage.ProviderState.pipe(optional), diff --git a/packages/schema/test/event-manifest.test.ts b/packages/schema/test/event-manifest.test.ts index 512528f556..89766d3d26 100644 --- a/packages/schema/test/event-manifest.test.ts +++ b/packages/schema/test/event-manifest.test.ts @@ -125,7 +125,6 @@ describe("public event manifest", () => { "session.tool.input.started.1", "session.tool.input.ended.1", "session.tool.called.1", - "session.tool.progress.1", "session.tool.success.1", "session.tool.failed.1", "session.reasoning.started.1", @@ -152,6 +151,8 @@ describe("public event manifest", () => { expect(EventManifest.Latest.has("session.usage.recorded")).toBe(false) expect(SessionEvent.UsageUpdated.durability).toBe("ephemeral") expect(SessionEvent.Compaction.Delta.durability).toBe("ephemeral") + expect(SessionEvent.Tool.Progress.durability).toBe("ephemeral") + expect(EventManifest.Server.get("session.tool.progress")).toBe(SessionEvent.Tool.Progress) expect(EventManifest.Durable.has("session.compaction.delta.1")).toBe(false) expect(EventManifest.ServerDefinitions).toContain(SessionEvent.UsageUpdated) expect(EventManifest.Definitions.every((definition) => definition.durability !== undefined)).toBe(true) diff --git a/packages/tui/src/context/data.tsx b/packages/tui/src/context/data.tsx index 41e9b68236..8ea43a921a 100644 --- a/packages/tui/src/context/data.tsx +++ b/packages/tui/src/context/data.tsx @@ -643,8 +643,8 @@ export const { use: useData, provider: DataProvider } = createSimpleContext({ status: "error", error: event.data.error, input: typeof match.state.input === "string" ? {} : match.state.input, - structured: match.state.status === "running" ? match.state.structured : {}, - content: match.state.status === "running" ? match.state.content : [], + structured: event.data.structured ?? (match.state.status === "running" ? match.state.structured : {}), + content: event.data.content ?? (match.state.status === "running" ? match.state.content : []), result: event.data.result, } match.executed = event.data.executed || match.executed === true diff --git a/packages/tui/src/mini/stream-v2.subagent.ts b/packages/tui/src/mini/stream-v2.subagent.ts index 164aa9a179..0a8ff2bdde 100644 --- a/packages/tui/src/mini/stream-v2.subagent.ts +++ b/packages/tui/src/mini/stream-v2.subagent.ts @@ -838,8 +838,9 @@ export function createSubagentTracker(input: SubagentTrackerInput): SubagentTrac ? { status: "error", input: part && part.state.status !== "streaming" ? part.state.input : {}, - structured: part && part.state.status !== "streaming" ? part.state.structured : {}, - content: part && part.state.status !== "streaming" ? part.state.content : [], + structured: + event.data.structured ?? (part && part.state.status !== "streaming" ? part.state.structured : {}), + content: event.data.content ?? (part && part.state.status !== "streaming" ? part.state.content : []), error: event.data.error, result: event.data.result, } @@ -958,14 +959,15 @@ export function createSubagentTracker(input: SubagentTrackerInput): SubagentTrac if (pendingCalls.has(key)) pendingCalls.set(key, event.data.input) return } - if (event.type === "session.tool.failed") { - pendingCalls.delete(sourceKey(event.data.assistantMessageID, event.data.callID)) + if ( + event.type !== "session.tool.progress" && + event.type !== "session.tool.success" && + event.type !== "session.tool.failed" + ) return - } - if (event.type !== "session.tool.progress" && event.type !== "session.tool.success") return const key = sourceKey(event.data.assistantMessageID, event.data.callID) const pending = pendingCalls.get(key) - if (event.type === "session.tool.success") pendingCalls.delete(key) + if (event.type !== "session.tool.progress") pendingCalls.delete(key) const found = childSessionID(record(event.data.structured)) if (!found) return const child = admitChild(found.sessionID) diff --git a/packages/tui/src/mini/stream-v2.transport.ts b/packages/tui/src/mini/stream-v2.transport.ts index b90f5e682d..383bfa1e7f 100644 --- a/packages/tui/src/mini/stream-v2.transport.ts +++ b/packages/tui/src/mini/stream-v2.transport.ts @@ -1067,8 +1067,9 @@ export async function createSessionTransport(input: StreamInput): Promise { id: "evt_progress_1", created: 0, type: "session.tool.progress", - durable: durable("session-1", 5), data: { sessionID: "session-1", assistantMessageID: "msg_explicit_assistant_9", diff --git a/packages/tui/test/mini/stream-v2.transport.test.ts b/packages/tui/test/mini/stream-v2.transport.test.ts index 371611972e..7fca422dc2 100644 --- a/packages/tui/test/mini/stream-v2.transport.test.ts +++ b/packages/tui/test/mini/stream-v2.transport.test.ts @@ -2015,7 +2015,6 @@ describe("V2 mini transport", () => { id: "evt_progress", created: 3, type: "session.tool.progress", - durable: durable("ses_1", 2), data: { sessionID: "ses_1", assistantMessageID: "msg_progress", @@ -2034,6 +2033,8 @@ describe("V2 mini transport", () => { assistantMessageID: "msg_progress", callID: "call_progress", error: { type: "unknown", message: "boom" }, + structured: { checkpoint: 1 }, + content: [{ type: "text", text: "partial" }], executed: true, }, }) @@ -2806,6 +2807,71 @@ describe("V2 mini transport", () => { await transport.close() }) + test("discovers a subagent from its terminal failure snapshot", async () => { + const events = feed() + events.push(connected()) + const client = sdk({ streams: [events], messages: { ses_child_failed: [] } }) + const ui = footer() + const transport = await createSessionTransport({ + sdk: client, + sessionID: "ses_1", + thinking: false, + footer: ui.api, + }) + const states = () => ui.events.flatMap((event) => (event.type === "stream.subagent" ? [event.state] : [])) + events.push({ + id: "evt_failed_subagent_input", + created: 1, + type: "session.tool.input.started", + durable: durable("ses_1"), + data: { + sessionID: "ses_1", + assistantMessageID: "msg_failed_subagent", + callID: "call_failed_subagent", + name: "subagent", + }, + }) + events.push({ + id: "evt_failed_subagent_called", + created: 2, + type: "session.tool.called", + durable: durable("ses_1", 1), + data: { + sessionID: "ses_1", + assistantMessageID: "msg_failed_subagent", + callID: "call_failed_subagent", + input: { agent: "explore", description: "Inspect failure", prompt: "inspect" }, + executed: true, + }, + }) + events.push({ + id: "evt_failed_subagent", + created: 3, + type: "session.tool.failed", + durable: durable("ses_1", 2), + data: { + sessionID: "ses_1", + assistantMessageID: "msg_failed_subagent", + callID: "call_failed_subagent", + error: { type: "unknown", message: "subagent failed" }, + structured: { sessionID: "ses_child_failed", status: "running" }, + content: [], + executed: true, + }, + }) + + while (!states().some((state) => state.tabs.some((tab) => tab.sessionID === "ses_child_failed"))) + await Bun.sleep(0) + expect(states().at(-1)?.tabs).toMatchObject([ + { + sessionID: "ses_child_failed", + label: "Explore", + description: "Inspect failure", + }, + ]) + await transport.close() + }) + test("discovers current subagents from progress and reduces descendant tool state", async () => { const events = feed() events.push(connected()) @@ -2847,7 +2913,6 @@ describe("V2 mini transport", () => { id: "evt_subagent_progress", created: 3, type: "session.tool.progress", - durable: durable("ses_1", 2), data: { sessionID: "ses_1", assistantMessageID: "msg_subagent", @@ -2899,7 +2964,6 @@ describe("V2 mini transport", () => { id: "evt_child_tool_progress", created: 6, type: "session.tool.progress", - durable: durable("ses_child_progress", 2), data: { sessionID: "ses_child_progress", assistantMessageID: "msg_child_tool", @@ -2930,6 +2994,8 @@ describe("V2 mini transport", () => { assistantMessageID: "msg_child_tool", callID: "call_child_shell", error: { type: "unknown", message: "child boom" }, + structured: { checkpoint: "child" }, + content: [{ type: "text", text: "child partial" }], executed: true, }, })