From d8efeffff74e15b2d0182838cbeac0f7f1cef151 Mon Sep 17 00:00:00 2001 From: Theo Browne Date: Fri, 4 Sep 2026 07:19:43 -0700 Subject: [PATCH 1/4] perf(server): stop retaining unused OpenCode tool history (#9684) (cherry picked from commit f2e3764c257a7e27c8171d7dd1e38d4383074206) --- .../src/provider/Layers/OpenCodeAdapter.ts | 28 ------------------- 1 file changed, 28 deletions(-) diff --git a/apps/server/src/provider/Layers/OpenCodeAdapter.ts b/apps/server/src/provider/Layers/OpenCodeAdapter.ts index 024b00b9a3..3371217141 100644 --- a/apps/server/src/provider/Layers/OpenCodeAdapter.ts +++ b/apps/server/src/provider/Layers/OpenCodeAdapter.ts @@ -342,7 +342,6 @@ interface OpenCodeSessionContext { readonly partById: Map; readonly emittedTextByPartId: Map; readonly completedAssistantPartIds: Set; - readonly turns: Array; activeTurnId: TurnId | undefined; activeAgent: string | undefined; activeVariant: string | undefined; @@ -481,31 +480,6 @@ function mapPermissionDecision(reply: "once" | "always" | "reject"): string { } } -function resolveTurnSnapshot( - context: OpenCodeSessionContext, - turnId: TurnId, -): OpenCodeTurnSnapshot { - const existing = context.turns.find((turn) => turn.id === turnId); - if (existing) { - return existing; - } - - const created: OpenCodeTurnSnapshot = { id: turnId, items: [] }; - context.turns.push(created); - return created; -} - -function appendTurnItem( - context: OpenCodeSessionContext, - turnId: TurnId | undefined, - item: unknown, -): void { - if (!turnId) { - return; - } - resolveTurnSnapshot(context, turnId).items.push(item); -} - const ensureSessionContext = Effect.fn("ensureSessionContext")(function* ( sessions: ReadonlyMap, threadId: ThreadId, @@ -2395,7 +2369,6 @@ export function makeOpenCodeAdapter( : "item.updated", payload, }; - appendTurnItem(context, turnId, part); yield* emit(context.sessionIncarnationId, runtimeEvent); } break; @@ -2888,7 +2861,6 @@ export function makeOpenCodeAdapter( emittedTextByPartId: new Map(), messageRoleById: new Map(), completedAssistantPartIds: new Set(), - turns: [], activeTurnId: undefined, activeAgent: undefined, activeVariant: undefined, From dd339d542f32d3ae96ed7da6ae760ce450043a14 Mon Sep 17 00:00:00 2001 From: Theo Browne Date: Fri, 4 Sep 2026 07:26:51 -0700 Subject: [PATCH 2/4] perf(server): omit repeated OpenCode progress logs (#9689) (cherry picked from commit ec8b2119c377f5c1dbe6235b221ef98eca31a96e) --- .../provider/Layers/EventNdjsonLogger.test.ts | 58 +++++++++++++++++++ .../src/provider/Layers/EventNdjsonLogger.ts | 10 +++- docs/internals/providers.md | 4 ++ 3 files changed, 71 insertions(+), 1 deletion(-) diff --git a/apps/server/src/provider/Layers/EventNdjsonLogger.test.ts b/apps/server/src/provider/Layers/EventNdjsonLogger.test.ts index 88c045c16a..0716b58481 100644 --- a/apps/server/src/provider/Layers/EventNdjsonLogger.test.ts +++ b/apps/server/src/provider/Layers/EventNdjsonLogger.test.ts @@ -341,6 +341,19 @@ describe("EventNdjsonLogger", () => { }, threadId, ); + yield* native.write( + { + event: { + type: "message.part.updated", + payload: { + properties: { + part: { type: "tool", state: { status: "running", output: "progress" } }, + }, + }, + }, + }, + threadId, + ); yield* native.write({ type: "turn.completed", id: "native-final" }, threadId); yield* store.close(); @@ -362,6 +375,51 @@ describe("EventNdjsonLogger", () => { }), ); + it.effect("keeps OpenCode tool input, final output, and errors in native logs", () => + Effect.gen(function* () { + const tempDir = NodeFS.mkdtempSync(NodePath.join(NodeOS.tmpdir(), "t3-provider-log-")); + const basePath = NodePath.join(tempDir, "events.log"); + const threadId = ThreadId.make("thread-tool-lifecycle"); + + try { + const store = yield* makeEventNdjsonLogStore(basePath, { batchWindowMs: 0 }); + const native = store.logger("native"); + const input = { command: "run-build" }; + const states = [ + { status: "pending", input }, + { status: "completed", input, output: "Build complete", time: { start: 1, end: 2 } }, + { status: "error", input, error: "Build failed", time: { start: 3, end: 4 } }, + { status: "unknown", input }, + null, + ]; + const events = states.map((state) => ({ + event: { + type: "message.part.updated", + payload: { properties: { part: { type: "tool", state } } }, + }, + })); + for (const event of events) { + yield* native.write(event, threadId); + } + yield* store.close(); + + const payloads = NodeFS.readFileSync( + ownedLogPath(basePath, "thread-tool-lifecycle"), + "utf8", + ) + .trim() + .split("\n") + .map((line) => parseLogLine(line).payload); + assert.deepEqual( + payloads, + events.map((event) => encodeUnknownJson(event)), + ); + } finally { + NodeFS.rmSync(tempDir, { recursive: true, force: true }); + } + }), + ); + it.effect("contains hostile event accessors inside guarded serialization", () => Effect.gen(function* () { const tempDir = NodeFS.mkdtempSync(NodePath.join(NodeOS.tmpdir(), "t3-provider-log-")); diff --git a/apps/server/src/provider/Layers/EventNdjsonLogger.ts b/apps/server/src/provider/Layers/EventNdjsonLogger.ts index db37ec2b3f..28c762e4c5 100644 --- a/apps/server/src/provider/Layers/EventNdjsonLogger.ts +++ b/apps/server/src/provider/Layers/EventNdjsonLogger.ts @@ -234,7 +234,15 @@ function shouldPersist(stream: EventNdjsonStream, event: unknown): boolean { const part = Reflect.get(properties, "part"); if (typeof part !== "object" || part === null) return true; const partType = Reflect.get(part, "type"); - return partType !== "text" && partType !== "reasoning"; + if (partType === "text" || partType === "reasoning") return false; + if (partType === "tool") { + // Running snapshots repeat growing output. Pending and terminal states stay in the log. + const state = Reflect.get(part, "state"); + if (typeof state === "object" && state !== null) { + return Reflect.get(state, "status") !== "running"; + } + } + return true; } return true; diff --git a/docs/internals/providers.md b/docs/internals/providers.md index 575a26b3a9..33e2f8df33 100644 --- a/docs/internals/providers.md +++ b/docs/internals/providers.md @@ -226,6 +226,10 @@ session's current model" and never sends it in `session/set_model`. ## OpenCode server ownership and catalog +Native OpenCode event logs skip text deltas and running tool snapshots. Pending, completed, +and failed tool states remain, including final output and errors. This filter does not change +events sent to clients or orchestration storage. + Each OpenCode provider instance owns one lazy local server for catalog discovery and text-generation helpers through [`OpenCodeServerOwner.ts`][opencode-server-owner]. Concurrent borrowers share startup. The server closes 30 seconds after the last borrower releases it, or From d99b46ccbdd1bdfbb3733a1bb4b04b382f770014 Mon Sep 17 00:00:00 2001 From: Theo Browne Date: Fri, 4 Sep 2026 11:04:05 -0700 Subject: [PATCH 3/4] perf(server): stop caching unused OpenCode tool parts (#9738) Do not retain OpenCode tool parts after emitting their runtime events. The remaining cache readers need text, reasoning, or step-usage parts, not tool input and output. Keep tool lifecycle events, output bytes, late-role assistant text, and usage handling unchanged. Add focused lifecycle coverage and verify retained memory through the real adapter. Created with GPT-6 Astra (preview) in Codex. (cherry picked from commit c8f77e0d441264efb0acfac312e852c81ae3da83) --- .../provider/Layers/OpenCodeAdapter.test.ts | 128 ++++++++++++++++++ .../src/provider/Layers/OpenCodeAdapter.ts | 5 +- 2 files changed, 132 insertions(+), 1 deletion(-) diff --git a/apps/server/src/provider/Layers/OpenCodeAdapter.test.ts b/apps/server/src/provider/Layers/OpenCodeAdapter.test.ts index 67c6f371c9..87d58de691 100644 --- a/apps/server/src/provider/Layers/OpenCodeAdapter.test.ts +++ b/apps/server/src/provider/Layers/OpenCodeAdapter.test.ts @@ -19,6 +19,7 @@ import type { Event as OpenCodeEvent, PermissionRequest, QuestionRequest, + ToolPart, } from "@opencode-ai/sdk/v2"; import { @@ -5796,6 +5797,133 @@ it.layer(OpenCodeAdapterTestLayer)("OpenCodeAdapterLive", (it) => { }), ); + it.effect("emits tool lifecycle events before late assistant metadata", () => + Effect.gen(function* () { + const adapter = yield* OpenCodeAdapter; + const threadId = asThreadId("thread-tool-lifecycle"); + const sessionID = "http://127.0.0.1:9999/session"; + const messageID = "msg-tools-before-role"; + const start = promiseWithResolvers(); + const input = { command: "pwd" }; + const states = [ + { status: "pending", input, raw: "" }, + { status: "running", input, title: "Working directory", time: { start: 1 } }, + { + status: "completed", + input, + output: "/repo\n", + title: "Working directory", + metadata: {}, + time: { start: 1, end: 2 }, + }, + { status: "error", input, error: "Command failed", time: { start: 3, end: 4 } }, + ] satisfies ReadonlyArray; + runtimeMock.state.subscribedEvents = [ + start.promise, + ...states.map( + (state) => + ({ + id: `evt-tool-${state.status}`, + type: "message.part.updated", + properties: { + sessionID, + time: 4, + part: { + id: state.status === "error" ? "part-failed" : "part-working", + sessionID, + messageID, + type: "tool", + callID: state.status === "error" ? "call-failed" : "call-working", + tool: "bash", + state, + }, + }, + }) satisfies OpenCodeEvent, + ), + { + id: "evt-text-before-role", + type: "message.part.updated", + properties: { + sessionID, + time: 4, + part: { + id: "part-late-text", + sessionID, + messageID, + type: "text", + text: "Tool results received", + time: { start: 4 }, + }, + }, + }, + { + id: "evt-late-assistant-role", + type: "message.updated", + properties: { sessionID, info: { id: messageID, role: "assistant" } }, + }, + { + id: "evt-tool-lifecycle-drained", + type: "session.compacted", + properties: { sessionID }, + }, + ]; + const eventsFiber = yield* adapter.streamEvents.pipe( + Stream.filter((event) => event.threadId === threadId), + Stream.takeUntil((event) => event.type === "thread.state.changed"), + Stream.runCollect, + Effect.forkChild, + ); + yield* adapter.startSession({ + provider: ProviderDriverKind.make("opencode"), + threadId, + runtimeMode: "full-access", + }); + yield* adapter.sendTurn({ + threadId, + input: "Read the working directory", + modelSelection: createModelSelection( + ProviderInstanceId.make("opencode"), + "opencode/kimi-k3", + ), + }); + start.resolve({ + id: "evt-tool-lifecycle-started", + type: "session.status", + properties: { sessionID, status: { type: "busy" } }, + }); + const events = yield* Fiber.join(eventsFiber); + const tools = events.filter( + (event) => + event.type === "item.started" || + event.type === "item.updated" || + event.type === "item.completed", + ); + NodeAssert.deepEqual( + tools.map((event) => [event.type, event.itemId, event.payload.status]), + [ + ["item.started", "call-working", "inProgress"], + ["item.updated", "call-working", "inProgress"], + ["item.completed", "call-working", "completed"], + ["item.completed", "call-failed", "failed"], + ], + ); + NodeAssert.partialDeepStrictEqual(tools[2]?.payload.data, { + command: "pwd", + result: "/repo\n", + }); + NodeAssert.partialDeepStrictEqual(tools[3]?.payload.data, { + command: "pwd", + state: { error: "Command failed" }, + }); + NodeAssert.deepEqual( + events + .filter((event) => event.type === "content.delta") + .map((event) => event.payload.delta), + ["Tool results received"], + ); + }), + ); + it.effect("maps native task progress only while a turn is active", () => Effect.gen(function* () { const adapter = yield* OpenCodeAdapter; diff --git a/apps/server/src/provider/Layers/OpenCodeAdapter.ts b/apps/server/src/provider/Layers/OpenCodeAdapter.ts index 3371217141..29d1bd0991 100644 --- a/apps/server/src/provider/Layers/OpenCodeAdapter.ts +++ b/apps/server/src/provider/Layers/OpenCodeAdapter.ts @@ -2317,7 +2317,10 @@ export function makeOpenCodeAdapter( case "message.part.updated": { const part = event.properties.part; - context.partById.set(part.id, part); + // Tool events use the incoming part and do not need a cached copy. + if (part.type !== "tool") { + context.partById.set(part.id, part); + } const messageRole = messageRoleForPart(context, part); if (messageRole === "assistant") { From 33378ccf9677a8c9ae299ddfd7c1f362f2182941 Mon Sep 17 00:00:00 2001 From: Trevor Walker Date: Sun, 6 Sep 2026 09:16:36 -0600 Subject: [PATCH 4/4] perf(server): index OpenCode text parts by message Adapt upstream 281b92b48810062275798dd4107364384ebaf432 without introducing per-turn token accounting. Preserve Pylon session incarnation ownership and admission, interruption, and rollback semantics. --- .../provider/Layers/OpenCodeAdapter.test.ts | 299 +++++++++++++++++- .../src/provider/Layers/OpenCodeAdapter.ts | 143 +++++---- 2 files changed, 374 insertions(+), 68 deletions(-) diff --git a/apps/server/src/provider/Layers/OpenCodeAdapter.test.ts b/apps/server/src/provider/Layers/OpenCodeAdapter.test.ts index 87d58de691..0ead57c2fb 100644 --- a/apps/server/src/provider/Layers/OpenCodeAdapter.test.ts +++ b/apps/server/src/provider/Layers/OpenCodeAdapter.test.ts @@ -3,6 +3,7 @@ import * as NodeServices from "@effect/platform-node/NodeServices"; import { it } from "@effect/vitest"; import * as Cause from "effect/Cause"; import * as Context from "effect/Context"; +import * as Deferred from "effect/Deferred"; import * as Effect from "effect/Effect"; import * as Exit from "effect/Exit"; import * as Fiber from "effect/Fiber"; @@ -14,7 +15,7 @@ import * as Schema from "effect/Schema"; import * as Scope from "effect/Scope"; import * as Stream from "effect/Stream"; import * as TestClock from "effect/testing/TestClock"; -import { beforeEach } from "vite-plus/test"; +import { beforeEach, vi } from "vite-plus/test"; import type { Event as OpenCodeEvent, PermissionRequest, @@ -568,6 +569,18 @@ function promiseWithResolvers() { return { promise, resolve, reject }; } +function makeOpenCodeEventQueue() { + let pending = promiseWithResolvers(); + const events = [pending.promise]; + runtimeMock.state.subscribedEvents = events; + return (event: unknown) => { + const current = pending; + pending = promiseWithResolvers(); + events.push(pending.promise); + current.resolve(event); + }; +} + const permissionRequest = (id: string, sessionID: string): PermissionRequest => ({ id, sessionID, @@ -5924,6 +5937,290 @@ it.layer(OpenCodeAdapterTestLayer)("OpenCodeAdapterLive", (it) => { }), ); + it.effect("processes late assistant metadata without visiting completed turns", () => + Effect.gen(function* () { + const adapter = yield* OpenCodeAdapter; + const threadId = asThreadId("thread-indexed-opencode-parts"); + const sessionID = "http://127.0.0.1:9999/session"; + const enqueue = makeOpenCodeEventQueue(); + yield* adapter.startSession({ + provider: ProviderDriverKind.make("opencode"), + threadId, + runtimeMode: "full-access", + }); + + for (let index = 0; index < 24; index += 1) { + const completed = yield* adapter.streamEvents.pipe( + Stream.filter((event) => event.threadId === threadId && event.type === "turn.completed"), + Stream.runHead, + Effect.forkChild, + ); + yield* adapter.sendTurn({ + threadId, + input: `Complete turn ${index}`, + modelSelection: createModelSelection( + ProviderInstanceId.make("opencode"), + "opencode/kimi-k3", + ), + }); + enqueue({ + type: "message.updated", + properties: { sessionID, info: { id: `history-message-${index}`, role: "assistant" } }, + }); + enqueue({ + type: "message.part.updated", + properties: { + sessionID, + part: { + id: `history-part-${index}`, + messageID: `history-message-${index}`, + sessionID, + type: "text", + text: `Completed turn ${index}`, + time: { start: 1, end: 2 }, + }, + }, + }); + enqueue({ + type: "session.status", + properties: { sessionID, status: { type: "idle" } }, + }); + yield* Fiber.join(completed); + } + + let visitedHistoryParts = 0; + const values = Map.prototype.values; + yield* Effect.acquireRelease( + Effect.sync(() => + vi + .spyOn(Map.prototype, "values") + .mockImplementation(function (this: Map) { + const iterator = values.call(this); + const next = iterator.next.bind(iterator); + iterator.next = () => { + const result = next(); + const value: unknown = result.value; + if ( + typeof value === "object" && + value !== null && + "id" in value && + typeof value.id === "string" && + value.id.startsWith("history-part-") + ) { + visitedHistoryParts += 1; + } + return result; + }; + return iterator; + }), + ), + (spy) => Effect.sync(() => spy.mockRestore()), + ); + + const stepProcessed = yield* Deferred.make(); + const eventsFiber = yield* adapter.streamEvents.pipe( + Stream.filter((event) => event.threadId === threadId), + Stream.tap((event) => + event.type === "thread.state.changed" + ? Deferred.succeed(stepProcessed, undefined) + : Effect.void, + ), + Stream.takeUntil((event) => event.type === "turn.completed"), + Stream.runCollect, + Effect.forkChild, + ); + yield* adapter.sendTurn({ + threadId, + input: "Process late metadata", + modelSelection: createModelSelection( + ProviderInstanceId.make("opencode"), + "opencode/kimi-k3", + ), + }); + const promptMessageId = (runtimeMock.state.promptCalls.at(-1) as { messageID: string }) + .messageID; + const part = { + id: "current-part", + messageID: "current-message", + sessionID, + type: "text", + text: "Current response", + time: { start: 3, end: 4 }, + }; + const step = { + id: "current-step", + messageID: "current-message", + sessionID, + type: "step-finish", + reason: "stop", + cost: 0, + tokens: { input: 40, output: 10, reasoning: 2, cache: { read: 5, write: 1 } }, + }; + enqueue({ type: "message.part.updated", properties: { sessionID, part } }); + enqueue({ type: "message.part.updated", properties: { sessionID, part: step } }); + enqueue({ type: "session.compacted", properties: { sessionID } }); + yield* Deferred.await(stepProcessed); + yield* adapter.sendTurn({ + threadId, + input: "Steer before metadata arrives", + modelSelection: createModelSelection( + ProviderInstanceId.make("opencode"), + "opencode/kimi-k3", + ), + }); + for (const parentID of ["", promptMessageId]) { + enqueue({ + type: "message.updated", + properties: { sessionID, info: { id: "current-message", role: "assistant", parentID } }, + }); + } + enqueue({ type: "message.part.updated", properties: { sessionID, part } }); + enqueue({ type: "message.part.updated", properties: { sessionID, part: step } }); + enqueue({ + type: "message.part.updated", + properties: { + sessionID, + part: { + id: "history-part-0", + messageID: "history-message-0", + sessionID, + type: "text", + text: "Completed turn zero", + time: { start: 1, end: 2 }, + }, + }, + }); + enqueue({ type: "session.status", properties: { sessionID, status: { type: "idle" } } }); + + const events = yield* Fiber.join(eventsFiber); + NodeAssert.equal(visitedHistoryParts, 0); + NodeAssert.deepEqual( + events + .filter((event) => event.type === "content.delta") + .map((event) => event.payload.delta), + ["Current response", "zero"], + ); + NodeAssert.equal(events.filter((event) => event.type === "item.completed").length, 1); + yield* adapter.stopSession(threadId); + }).pipe(Effect.scoped), + ); + + it.effect("keeps completed text edits and clears removed parts across reconnects", () => + Effect.gen(function* () { + const adapter = yield* OpenCodeAdapter; + const threadId = asThreadId("thread-opencode-text-retention"); + const sessionID = "http://127.0.0.1:9999/session"; + const messageID = "retained-message"; + const metadata = { + type: "message.updated", + properties: { + sessionID, + info: { id: messageID, role: "assistant", time: { created: 1, completed: 2 } }, + }, + }; + const snapshot = ( + text: string, + id = "retained-part", + type: "text" | "reasoning" = "text", + ) => ({ + type: "message.part.updated", + properties: { + sessionID, + part: { id, sessionID, messageID, type, text, time: { start: 1, end: 2 } }, + }, + }); + const delta = (text: string) => ({ + type: "message.part.delta", + properties: { sessionID, messageID, partID: "retained-part", field: "text", delta: text }, + }); + const nonTextReplacement = { + type: "message.part.updated", + properties: { + sessionID, + part: { + id: "retained-part", + sessionID, + messageID, + type: "file", + mime: "text/plain", + url: "file:///repo/result.txt", + }, + }, + }; + runtimeMock.state.subscribedEvents = [ + snapshot("Replaced before metadata"), + nonTextReplacement, + metadata, + snapshot("Thinking", "reasoning-part", "reasoning"), + snapshot("Hello world"), + { type: "server.connected", properties: {} }, + metadata, + snapshot("Thinking", "reasoning-part", "reasoning"), + snapshot("Thinking more", "reasoning-part", "reasoning"), + snapshot("Hello world"), + snapshot("Hello"), + snapshot("Hello there"), + delta(" again"), + snapshot("Hello there again"), + nonTextReplacement, + delta("ignored while file"), + metadata, + snapshot("Hello there again!"), + { + type: "message.part.removed", + properties: { sessionID, messageID, partID: "retained-part" }, + }, + delta("removed part"), + metadata, + snapshot("Fresh"), + snapshot("Second", "second-part"), + { type: "message.removed", properties: { sessionID, messageID } }, + delta("removed message"), + metadata, + snapshot("New thoughts", "reasoning-part", "reasoning"), + snapshot("New"), + { type: "session.compacted", properties: { sessionID } }, + ]; + const eventsFiber = yield* adapter.streamEvents.pipe( + Stream.filter((event) => event.threadId === threadId), + Stream.takeUntil((event) => event.type === "thread.state.changed"), + Stream.runCollect, + Effect.forkChild, + ); + yield* adapter.startSession({ + provider: ProviderDriverKind.make("opencode"), + threadId, + runtimeMode: "full-access", + }); + + const events = yield* Fiber.join(eventsFiber); + NodeAssert.deepEqual( + events + .filter((event) => event.type === "content.delta") + .map((event) => [event.payload.streamKind, event.payload.delta]), + [ + ["reasoning_text", "Thinking"], + ["assistant_text", "Hello world"], + ["reasoning_text", " more"], + ["assistant_text", "there"], + ["assistant_text", " again"], + ["assistant_text", "!"], + ["assistant_text", "Fresh"], + ["assistant_text", "Second"], + ["reasoning_text", "New thoughts"], + ["assistant_text", "New"], + ], + ); + NodeAssert.deepEqual( + events + .filter((event) => event.type === "item.completed") + .map((event) => event.payload.detail), + ["Hello world", "Fresh", "Second", "New"], + ); + yield* adapter.stopSession(threadId); + }), + ); + it.effect("maps native task progress only while a turn is active", () => Effect.gen(function* () { const adapter = yield* OpenCodeAdapter; diff --git a/apps/server/src/provider/Layers/OpenCodeAdapter.ts b/apps/server/src/provider/Layers/OpenCodeAdapter.ts index 29d1bd0991..7dd3a667a4 100644 --- a/apps/server/src/provider/Layers/OpenCodeAdapter.ts +++ b/apps/server/src/provider/Layers/OpenCodeAdapter.ts @@ -323,6 +323,14 @@ function isOpenCodeDefaultTitle(title: string): boolean { return OPENCODE_DEFAULT_TITLE_PATTERN.test(title); } +type OpenCodeTextPart = Extract; + +type OpenCodeTextPartState = Pick & { + text: string | undefined; + emittedText: string | undefined; + completed: boolean; +}; + interface OpenCodeSessionContext { session: ProviderSession; /** Immutable identity captured before the context can be removed or replaced. */ @@ -339,9 +347,8 @@ interface OpenCodeSessionContext { readonly pendingPermissions: Map; readonly pendingQuestions: Map; readonly messageRoleById: Map; - readonly partById: Map; - readonly emittedTextByPartId: Map; - readonly completedAssistantPartIds: Set; + // Completed text can receive later edits; native removal releases retained text. + readonly textPartsByMessageId: Map>; activeTurnId: TurnId | undefined; activeAgent: string | undefined; activeVariant: string | undefined; @@ -513,18 +520,29 @@ function normalizeQuestionRequest(request: QuestionRequest): ReadonlyArray): "assistant_text" | "reasoning_text" { + return part.type === "reasoning" ? "reasoning_text" : "assistant_text"; } -function textFromPart(part: Part): string | undefined { - switch (part.type) { - case "text": - case "reasoning": - return part.text; - default: - return undefined; - } +function retainOpenCodeTextPart( + context: OpenCodeSessionContext, + part: OpenCodeTextPart, +): OpenCodeTextPartState { + const parts = + context.textPartsByMessageId.get(part.messageID) ?? new Map(); + const previous = parts.get(part.id); + const state = { + id: part.id, + messageID: part.messageID, + type: part.type, + text: part.text, + ...(part.time !== undefined ? { time: part.time } : {}), + emittedText: previous?.emittedText, + completed: previous?.completed ?? false, + }; + parts.set(part.id, state); + context.textPartsByMessageId.set(part.messageID, parts); + return state; } function commonPrefixLength(left: string, right: string): number { @@ -1304,6 +1322,7 @@ export function makeOpenCodeAdapter( if (message?.info.id === promptAdmission.messageId && message.info.role === "user") { promptAdmission.messageObserved = true; context.messageRoleById.set(promptAdmission.messageId, "user"); + context.textPartsByMessageId.delete(promptAdmission.messageId); } } @@ -1509,35 +1528,23 @@ export function makeOpenCodeAdapter( /** Emit content.delta and item.completed events for an assistant text part. */ const emitAssistantTextDelta = Effect.fn("emitAssistantTextDelta")(function* ( context: OpenCodeSessionContext, - part: Part, + part: OpenCodeTextPartState, turnId: TurnId | undefined, raw: unknown, ) { - const text = textFromPart(part); - if (text === undefined) { + if (part.text === undefined) { return; } - const previousText = context.emittedTextByPartId.get(part.id); - const { latestText, deltaToEmit } = mergeOpenCodeAssistantText(previousText, text); - context.emittedTextByPartId.set(part.id, latestText); - if (latestText !== text) { - context.partById.set( - part.id, - (part.type === "text" || part.type === "reasoning" - ? { ...part, text: latestText } - : part) satisfies Part, - ); - } + const { latestText, deltaToEmit } = mergeOpenCodeAssistantText(part.emittedText, part.text); + part.emittedText = latestText; + part.text = latestText; if (deltaToEmit.length > 0) { yield* emit(context.sessionIncarnationId, { ...(yield* buildEventBase({ threadId: context.session.threadId, turnId, itemId: part.id, - createdAt: - (part.type === "text" || part.type === "reasoning") && part.time !== undefined - ? isoFromEpochMs(part.time.start) - : undefined, + createdAt: part.time !== undefined ? isoFromEpochMs(part.time.start) : undefined, raw, })), type: "content.delta", @@ -1548,12 +1555,8 @@ export function makeOpenCodeAdapter( }); } - if ( - part.type === "text" && - part.time?.end !== undefined && - !context.completedAssistantPartIds.has(part.id) - ) { - context.completedAssistantPartIds.add(part.id); + if (part.type === "text" && part.time?.end !== undefined && !part.completed) { + part.completed = true; yield* emit(context.sessionIncarnationId, { ...(yield* buildEventBase({ threadId: context.session.threadId, @@ -2250,11 +2253,13 @@ export function makeOpenCodeAdapter( } } context.messageRoleById.set(event.properties.info.id, event.properties.info.role); + if (event.properties.info.role === "user") { + context.textPartsByMessageId.delete(event.properties.info.id); + } if (event.properties.info.role === "assistant") { - for (const part of context.partById.values()) { - if (part.messageID !== event.properties.info.id) { - continue; - } + for (const part of context.textPartsByMessageId + .get(event.properties.info.id) + ?.values() ?? []) { yield* emitAssistantTextDelta(context, part, turnId, event); } } @@ -2263,16 +2268,24 @@ export function makeOpenCodeAdapter( case "message.removed": { context.messageRoleById.delete(event.properties.messageID); + context.textPartsByMessageId.delete(event.properties.messageID); + break; + } + + case "message.part.removed": { + const parts = context.textPartsByMessageId.get(event.properties.messageID); + parts?.delete(event.properties.partID); + if (parts?.size === 0) { + context.textPartsByMessageId.delete(event.properties.messageID); + } break; } case "message.part.delta": { - const existingPart = context.partById.get(event.properties.partID); - if ( - !existingPart || - (existingPart.type !== "text" && existingPart.type !== "reasoning") || - event.properties.field !== "text" - ) { + const existingPart = context.textPartsByMessageId + .get(event.properties.messageID) + ?.get(event.properties.partID); + if (existingPart?.text === undefined || event.properties.field !== "text") { break; } const role = messageRoleForPart(context, existingPart); @@ -2284,21 +2297,13 @@ export function makeOpenCodeAdapter( if (delta.length === 0) { break; } - const previousText = - context.emittedTextByPartId.get(event.properties.partID) ?? - textFromPart(existingPart) ?? - ""; + const previousText = existingPart.emittedText ?? existingPart.text; const { nextText, deltaToEmit } = appendOpenCodeAssistantTextDelta(previousText, delta); if (deltaToEmit.length === 0) { break; } - context.emittedTextByPartId.set(event.properties.partID, nextText); - if (existingPart.type === "text" || existingPart.type === "reasoning") { - context.partById.set(event.properties.partID, { - ...existingPart, - text: nextText, - }); - } + existingPart.emittedText = nextText; + existingPart.text = nextText; yield* emit(context.sessionIncarnationId, { ...(yield* buildEventBase({ threadId: context.session.threadId, @@ -2317,14 +2322,20 @@ export function makeOpenCodeAdapter( case "message.part.updated": { const part = event.properties.part; - // Tool events use the incoming part and do not need a cached copy. - if (part.type !== "tool") { - context.partById.set(part.id, part); - } const messageRole = messageRoleForPart(context, part); - if (messageRole === "assistant") { - yield* emitAssistantTextDelta(context, part, turnId, event); + if ((part.type === "text" || part.type === "reasoning") && messageRole !== "user") { + const state = retainOpenCodeTextPart(context, part); + if (messageRole === "assistant") { + yield* emitAssistantTextDelta(context, state, turnId, event); + } + } else { + const previous = context.textPartsByMessageId.get(part.messageID)?.get(part.id); + if (previous) { + // A non-text PATCH removes the current snapshot. Keep emitted text + // so a later text PATCH still emits only the changed suffix. + previous.text = undefined; + } } if (part.type === "tool") { @@ -2860,10 +2871,8 @@ export function makeOpenCodeAdapter( requestRelationRetries: new Map(), pendingPermissions: new Map(), pendingQuestions: new Map(), - partById: new Map(), - emittedTextByPartId: new Map(), + textPartsByMessageId: new Map(), messageRoleById: new Map(), - completedAssistantPartIds: new Set(), activeTurnId: undefined, activeAgent: undefined, activeVariant: undefined,