diff --git a/packages/server/src/server/agent/providers/opencode-agent.test.ts b/packages/server/src/server/agent/providers/opencode-agent.test.ts index e5ef7f681..26f3b89cb 100644 --- a/packages/server/src/server/agent/providers/opencode-agent.test.ts +++ b/packages/server/src/server/agent/providers/opencode-agent.test.ts @@ -91,6 +91,40 @@ function providerAssistantMessages(events: AgentStreamEvent[], text: string): Ag ); } +type TurnEventSignature = [type: AgentStreamEvent["type"], turnId: string | undefined]; + +function turnEventSignatures(events: AgentStreamEvent[]): TurnEventSignature[] { + return events.map((event) => [event.type, "turnId" in event ? event.turnId : undefined]); +} + +function userMessageEvents(params: { + sessionId: string; + messageId: string; + text: string; +}): unknown[] { + return [ + { + type: "message.updated", + properties: { + info: { id: params.messageId, sessionID: params.sessionId, role: "user" }, + }, + }, + { + type: "message.part.updated", + properties: { + part: { + id: `part_${params.messageId}`, + sessionID: params.sessionId, + messageID: params.messageId, + type: "text", + text: params.text, + time: { start: 1, end: 2 }, + }, + }, + }, + ]; +} + function assistantTurnEvents({ sessionId = "session-1", text = "Hello from OpenCode", @@ -123,6 +157,29 @@ function assistantTurnEvents({ ]; } +async function createParentSession( + sessionId: string, + configure?: (openCode: TestOpenCodeClient) => void, +): Promise<{ + readonly parent: Awaited>; + readonly openCode: TestOpenCodeClient; +}> { + const runtime = new TestOpenCodeHarness(); + const openCode = new TestOpenCodeClient(); + openCode.sessionCreateResponse = { data: { id: sessionId } }; + configure?.(openCode); + runtime.enqueueClient(openCode); + const client = new OpenCodeAgentClient(createTestLogger(), undefined, { + serverManager: runtime, + createClient: runtime.createClient, + }); + const parent = await client.createSession({ + provider: "opencode", + cwd: "/workspace/repo", + }); + return { parent, openCode }; +} + function manualCompactEvents({ sessionId = "session-1", summaryText = "## Goal\n- Preserve context while continuing the task.", @@ -1844,56 +1901,91 @@ describe("OpenCode adapter startTurn error handling", () => { } }); - test("delays the next prompt until a slow interrupt abort settles", async () => { - vi.useFakeTimers(); - const abortDeferred = createTestDeferred<{ data: boolean; error: undefined }>(); - const promptAsync = vi.fn().mockResolvedValue({ data: {}, error: undefined }); - const abort = vi - .fn() - .mockReturnValueOnce(abortDeferred.promise) - .mockResolvedValue({ data: true, error: undefined }); - const fakeClient = { - global: { - event: vi.fn().mockImplementation( - async (options: { - signal: AbortSignal; - }): Promise<{ stream: AsyncIterable }> => ({ - stream: abortableOpenCodeStream(options.signal), - }), - ), - }, - session: { - promptAsync, - abort, - }, - } as never; + test("retries abort and waits for provider idle before starting the next prompt", async () => { + const { parent: session, openCode } = await createParentSession("ses_unit_test"); + openCode.sessionPromptAsyncEvents = []; + openCode.sessionAbortImplementation = async () => { + if (openCode.calls.sessionAbort.length === 1) { + throw new Error("abort transport failed"); + } + openCode.emitEvent({ + type: "session.idle", + properties: { sessionID: "ses_unit_test" }, + }); + return { data: true }; + }; + try { + await session.startTurn("first"); + await session.interrupt(); + expect(openCode.calls.sessionPromptAsync).toHaveLength(1); - const session = new __openCodeInternals.OpenCodeAgentSession( - { provider: "opencode", cwd: "/tmp/test" }, - fakeClient, - "ses_unit_test", - createTestLogger(), - ); - - await session.startTurn("first"); - expect(promptAsync).toHaveBeenCalledTimes(1); - - const interruptPromise = session.interrupt(); - await vi.advanceTimersByTimeAsync(2_000); - await interruptPromise; - expect(abort).toHaveBeenCalledTimes(1); - - const secondTurnPromise = session.startTurn("second"); - await vi.advanceTimersByTimeAsync(1_000); - expect(promptAsync).toHaveBeenCalledTimes(1); - - abortDeferred.resolve({ data: true, error: undefined }); - await secondTurnPromise; - expect(promptAsync).toHaveBeenCalledTimes(2); - - await session.interrupt(); - vi.useRealTimers(); + await session.startTurn("second"); + expect(openCode.calls.sessionAbort).toHaveLength(2); + expect(openCode.calls.sessionPromptAsync).toHaveLength(2); + } finally { + await session.close(); + } }); + + test("does not send a prompt while the previous provider turn is still stopping", async () => { + vi.useFakeTimers(); + const { parent: session, openCode } = await createParentSession("ses_unit_test"); + openCode.sessionPromptAsyncEvents = []; + + try { + await session.startTurn("first"); + await session.interrupt(); + + const secondTurn = expect(session.startTurn("second")).rejects.toThrow( + "OpenCode previous turn to stop", + ); + await vi.advanceTimersByTimeAsync(10_000); + await secondTurn; + expect(openCode.calls.sessionPromptAsync).toHaveLength(1); + } finally { + vi.useRealTimers(); + await session.close(); + } + }); + + test("reconnects and confirms provider idle when the event stream ends while stopping", async () => { + const firstStreamEnd = createTestDeferred(); + let subscriptionCount = 0; + const { parent: session, openCode } = await createParentSession( + "ses_stream_reconnect", + (client) => { + client.globalEventImplementation = async (options) => { + subscriptionCount += 1; + const signal = (options as { signal: AbortSignal }).signal; + const end = subscriptionCount === 1 ? firstStreamEnd.promise : waitForAbort(signal); + return { + stream: { + async *[Symbol.asyncIterator]() { + yield { type: "server.connected", properties: {} }; + await end; + }, + }, + }; + }; + }, + ); + openCode.sessionPromptAsyncEvents = []; + + try { + await session.startTurn("first"); + await session.interrupt(); + + const secondTurn = session.startTurn("second"); + firstStreamEnd.resolve(); + + await expect(secondTurn).resolves.toEqual({ turnId: "opencode-turn-1" }); + expect(openCode.calls.globalEvent).toHaveLength(2); + expect(openCode.calls.sessionStatus).toEqual([{ directory: "/workspace/repo" }]); + expect(openCode.calls.sessionPromptAsync).toHaveLength(2); + } finally { + await session.close(); + } + }, 15_000); }); describe("OpenCodeAgentClient env", () => { @@ -2464,7 +2556,7 @@ describe("OpenCode provider subagent contract", () => { ]); }); - test("synthesizes a turn for externally driven adopted child timeline events", async () => { + test("synthesizes an autonomous turn for adopted child timeline events", async () => { const { child, childClient, parent } = await createAdoptedChildSession(); const completed = createTestDeferred(); const events: AgentStreamEvent[] = []; @@ -2515,7 +2607,184 @@ describe("OpenCode provider subagent contract", () => { ); }); - test("synthesizes a turn for externally driven adopted child permissions", async () => { + test("surfaces an autonomous OpenCode wake as a new parent turn", async () => { + const { parent, openCode } = await createParentSession("ses_parent_plugin_wake"); + const events: AgentStreamEvent[] = []; + parent.subscribe((event) => events.push(event)); + + try { + for (const event of userMessageEvents({ + sessionId: "ses_parent_plugin_wake", + messageId: "msg_plugin_wake", + text: "All background tasks are complete.", + })) { + openCode.emitEvent(event); + } + openCode.emitEvent({ + type: "session.status", + properties: { sessionID: "ses_parent_plugin_wake", status: { type: "busy" } }, + }); + for (const event of assistantTurnEvents({ + sessionId: "ses_parent_plugin_wake", + text: "Parent consumed the background result.", + })) { + openCode.emitEvent(event); + } + + await vi.waitFor(() => { + expect(turnEventSignatures(events)).toEqual([ + ["turn_started", "opencode-turn-0"], + ["timeline", "opencode-turn-0"], + ["timeline", "opencode-turn-0"], + ["turn_completed", "opencode-turn-0"], + ]); + }); + expect(events).toContainEqual( + expect.objectContaining({ + type: "timeline", + item: expect.objectContaining({ + type: "assistant_message", + text: "Parent consumed the background result.", + }), + }), + ); + } finally { + await parent.close(); + } + }); + + test("does not mistake OpenCode session metadata for an autonomous turn", async () => { + const { parent, openCode } = await createParentSession("ses_parent_metadata"); + const events: AgentStreamEvent[] = []; + const streamDrained = createTestDeferred(); + parent.subscribe((event) => { + events.push(event); + if (event.type === "provider_subagent") { + streamDrained.resolve(); + } + }); + + try { + openCode.emitEvent({ + type: "session.created", + properties: { info: { id: "ses_parent_metadata" } }, + }); + openCode.emitEvent({ + type: "session.idle", + properties: { sessionID: "ses_parent_metadata" }, + }); + for (const event of assistantTurnEvents({ + sessionId: "ses_parent_metadata", + text: "Autonomous response.", + })) { + openCode.emitEvent(event); + } + openCode.emitEvent({ + type: "session.created", + properties: { + info: { + id: "ses_child_after_parent_metadata", + parentID: "ses_parent_metadata", + title: "Stream drain marker", + directory: "/workspace/repo", + }, + }, + }); + + await streamDrained.promise; + expect( + turnEventSignatures(events.filter((event) => event.type !== "provider_subagent")), + ).toEqual([]); + } finally { + await parent.close(); + } + }); + + test("keeps an autonomous turn active until OpenCode reaches its terminal event", async () => { + const { parent, openCode } = await createParentSession("ses_parent_active_wake"); + const events: AgentStreamEvent[] = []; + parent.subscribe((event) => events.push(event)); + + try { + openCode.emitEvent( + userMessageEvents({ + sessionId: "ses_parent_active_wake", + messageId: "msg_autonomous_wake", + text: "Autonomous wake", + })[0], + ); + await vi.waitFor(() => { + expect(events).toEqual([ + { type: "turn_started", provider: "opencode", turnId: "opencode-turn-0" }, + ]); + }); + + await expect(parent.startTurn("Continue from Paseo")).rejects.toThrow( + "A foreground turn is already active", + ); + expect(openCode.calls.sessionAbort).toEqual([]); + expect(openCode.calls.sessionPromptAsync).toEqual([]); + } finally { + await parent.close(); + } + }); + test("does not adopt late output from an interrupted Paseo turn", async () => { + const { parent, openCode } = await createParentSession("ses_parent_interrupted"); + openCode.sessionPromptAsyncEvents = []; + const events: AgentStreamEvent[] = []; + const streamDrained = createTestDeferred(); + parent.subscribe((event) => { + events.push(event); + if (event.type === "provider_subagent") { + streamDrained.resolve(); + } + }); + + try { + await parent.startTurn("Start from Paseo"); + await parent.interrupt(); + + for (const event of userMessageEvents({ + sessionId: "ses_parent_interrupted", + messageId: "msg_delayed_interrupted_user", + text: "Delayed interrupted prompt", + })) { + openCode.emitEvent(event); + } + for (const event of assistantTurnEvents({ + sessionId: "ses_parent_interrupted", + text: "Late interrupted response.", + })) { + openCode.emitEvent(event); + } + openCode.emitEvent({ + type: "session.created", + properties: { + info: { + id: "ses_child_after_interrupted_output", + parentID: "ses_parent_interrupted", + title: "Stream drain marker", + directory: "/workspace/repo", + }, + }, + }); + + await streamDrained.promise; + expect(events.filter((event) => event.type !== "provider_subagent")).toEqual([ + { type: "turn_started", provider: "opencode", turnId: "opencode-turn-0" }, + { + type: "turn_canceled", + provider: "opencode", + reason: "interrupted", + turnId: "opencode-turn-0", + }, + ]); + } finally { + await parent.close(); + } + }); + + test("synthesizes an autonomous turn for adopted child permissions", async () => { const { child, childClient, parent } = await createAdoptedChildSession(); const completed = createTestDeferred(); const events: AgentStreamEvent[] = []; @@ -2617,6 +2886,13 @@ describe("OpenCode provider subagent contract", () => { }), ); }); + expect(events.some((event) => event.type === "turn_started")).toBe(false); + expect( + events.find( + (event) => + event.type === "permission_requested" && event.request.id === "perm_provider_child", + ), + ).not.toHaveProperty("turnId"); expect( events.filter( (event) => @@ -3446,31 +3722,3 @@ function waitForAbort(signal: AbortSignal): Promise { signal.addEventListener("abort", () => resolve(), { once: true }); }); } - -function abortableOpenCodeStream(signal: AbortSignal): AsyncIterable { - return { - [Symbol.asyncIterator]: () => { - let emittedConnected = false; - return { - next: () => { - if (!emittedConnected) { - emittedConnected = true; - return Promise.resolve({ - done: false, - value: { type: "server.connected", properties: {} } as OpenCodeEvent, - }); - } - return new Promise>((resolve) => { - if (signal.aborted) { - resolve({ done: true, value: undefined }); - return; - } - signal.addEventListener("abort", () => resolve({ done: true, value: undefined }), { - once: true, - }); - }); - }, - }; - }, - }; -} diff --git a/packages/server/src/server/agent/providers/opencode-agent.ts b/packages/server/src/server/agent/providers/opencode-agent.ts index 0eeca0e7f..954896b6a 100644 --- a/packages/server/src/server/agent/providers/opencode-agent.ts +++ b/packages/server/src/server/agent/providers/opencode-agent.ts @@ -2756,6 +2756,11 @@ function createDeferred(): Deferred { return { promise, resolve, reject }; } +type OpenCodeTurnState = + | { status: "idle" } + | { status: "running"; turnId: string } + | { status: "stopping"; idle: Deferred }; + function unwrapOpenCodeGlobalEvent(event: unknown): OpenCodeEvent | null { const record = readOpenCodeRecord(event); if (!record) { @@ -2791,14 +2796,6 @@ function getOpenCodeEventSessionId(event: OpenCodeEvent): string | null { ); } -function isOpenCodeUserMessageEvent(event: OpenCodeEvent, sessionId: string): boolean { - return ( - event.type === "message.updated" && - event.properties.info.sessionID === sessionId && - event.properties.info.role === "user" - ); -} - function isOpenCodeTerminalEvent(event: OpenCodeEvent, sessionId: string): boolean { if (event.type === "session.idle" || event.type === "session.error") { return event.properties.sessionID === sessionId; @@ -2893,7 +2890,6 @@ class OpenCodeAgentSession implements AgentSession { private autoAcceptEnabled = false; private pendingPermissions = new Map(); private abortController: AbortController | null = null; - private pendingAbortPromise: Promise | null = null; private accumulatedUsage: AgentUsage = {}; private sessionTotalCostUsd: number | undefined; private mcpConfigured = false; @@ -2915,8 +2911,7 @@ class OpenCodeAgentSession implements AgentSession { private availableModesCache: AgentMode[] | null = null; private readonly subscribers = new Set<(event: AgentStreamEvent) => void>(); private nextTurnOrdinal = 0; - private activeForegroundTurnId: string | null = null; - private activeForegroundTurnSource: "paseo" | "external" | null = null; + private turnState: OpenCodeTurnState = { status: "idle" }; private readonly runningToolCalls = new Map(); private subAgentsByCallId = new Map(); private subAgentCallIdByChildSessionId = new Map(); @@ -2932,7 +2927,6 @@ class OpenCodeAgentSession implements AgentSession { private eventStreamAbortController: AbortController | null = null; private eventStreamReady: Deferred | null = null; private eventStreamTask: Promise | null = null; - private suppressTerminalUntilNextUserMessage = false; private closed = false; private readonly persistSession: boolean; private deletedFromProvider = false; @@ -2967,6 +2961,10 @@ class OpenCodeAgentSession implements AgentSession { return this.sessionId; } + private get activeForegroundTurnId(): string | null { + return this.turnState.status === "running" ? this.turnState.turnId : null; + } + get features(): AgentFeature[] { return [buildOpenCodeAutoAcceptFeature(this.config)]; } @@ -3008,29 +3006,24 @@ class OpenCodeAgentSession implements AgentSession { } async interrupt(): Promise { - let turnId = this.activeForegroundTurnId; - const turnAbortController = this.abortController; - turnAbortController?.abort(); + const turnId = this.activeForegroundTurnId; + this.abortController?.abort(); + if (turnId) { + this.beginStoppingTurn(turnId); + } // COMPAT(opencodeSlowAbort): OpenCode 1.14.42+ blocks session.abort until // the running tool actually stops, which can be tens of seconds for // long-running tools. Cap the wait so the user-visible cancel lands // quickly while still giving OpenCode a chance to confirm the abort // cleanly. Drop the timeout once upstream returns abort acknowledgement // before tool teardown. - const abortPromise = this.beginSessionAbort(turnId, "interrupt"); + const abortPromise = this.abortSession(turnId, "interrupt"); await withTimeout(abortPromise, 2_000, "OpenCode session.abort").catch((error) => { this.logger.warn( { err: error, sessionId: this.sessionId, turnId }, "OpenCode session.abort exceeded the cancel cap; proceeding with local cancel", ); }); - if (turnId) { - this.suppressTerminalUntilNextUserMessage = true; - this.finishForegroundTurn( - { type: "turn_canceled", provider: "opencode", reason: "interrupted" }, - turnId, - ); - } } async revertBoth(input: { messageId: string }): Promise { @@ -3042,8 +3035,8 @@ class OpenCodeAgentSession implements AgentSession { }); } - private beginSessionAbort(turnId: string | null, reason: string): Promise { - const abortPromise = this.client.session + private abortSession(turnId: string | null, reason: string): Promise { + return this.client.session .abort({ sessionID: this.sessionId, directory: this.config.cwd, @@ -3055,49 +3048,77 @@ class OpenCodeAgentSession implements AgentSession { "OpenCode session.abort rejected", ); }); - const trackedAbortPromise = abortPromise.finally(() => { - if (this.pendingAbortPromise === trackedAbortPromise) { - this.pendingAbortPromise = null; - } - }); - this.pendingAbortPromise = trackedAbortPromise; - return trackedAbortPromise; } - private async awaitPendingAbortBeforeStartingTurn(): Promise { - const pendingAbortPromise = this.pendingAbortPromise; - if (!pendingAbortPromise) { + private async waitUntilProviderIdle(): Promise { + if (this.turnState.status !== "stopping") { return; } + const stopping = this.turnState; + // OpenCode creates prompts before joining its single session runner. Sending + // while that runner is stopping can strand the new prompt behind the old + // run, so only provider-confirmed idle makes the session reusable. await withTimeout( - pendingAbortPromise, + this.observeProviderStopBoundary(stopping), OPENCODE_PENDING_ABORT_START_TIMEOUT_MS, - "OpenCode pending session.abort", - ).catch((error) => { - this.logger.warn( - { err: error, sessionId: this.sessionId }, - "OpenCode session.abort was still pending before starting the next turn", - ); - }); + "OpenCode previous turn to stop", + ); + } + + private async observeProviderStopBoundary( + stopping: Extract, + ): Promise { + while (this.turnState === stopping) { + void this.abortSession(null, "turn_start_boundary"); + const eventStreamTask = this.eventStreamTask; + const boundary = await (eventStreamTask + ? Promise.race([ + stopping.idle.promise.then(() => "idle" as const), + eventStreamTask.then( + () => "stream_ended" as const, + () => "stream_ended" as const, + ), + ]) + : Promise.resolve("stream_ended" as const)); + if (boundary === "idle") { + return; + } + + await this.ensureEventStreamReady(); + const response = await this.client.session.status({ directory: this.config.cwd }); + if (response.error) { + throw new Error( + `Failed to confirm OpenCode session status: ${toDiagnosticErrorMessage(response.error)}`, + ); + } + const statuses = readOpenCodeRecord(response.data); + if (!statuses) { + throw new Error("OpenCode returned an invalid session status response"); + } + const status = readOpenCodeRecord(statuses[this.sessionId]); + const statusType = readNonEmptyString(status?.type); + if (!status || statusType === "idle") { + this.finishStoppingTurn(); + return; + } + if (statusType !== "busy" && statusType !== "retry") { + throw new Error(`OpenCode returned an unknown session status '${statusType ?? "missing"}'`); + } + } } async startTurn( prompt: AgentPromptInput, options?: AgentRunOptions, ): Promise<{ turnId: string }> { - if (this.activeForegroundTurnId) { - if (this.activeForegroundTurnSource === "external") { - // A direct Paseo prompt owns the foreground; close the adopted child run first. - this.finishForegroundTurn( - { type: "turn_completed", provider: "opencode", usage: undefined }, - this.activeForegroundTurnId, - ); - } else { - throw new Error("A foreground turn is already active"); - } + if (this.turnState.status === "running") { + throw new Error("A foreground turn is already active"); + } + await this.waitUntilProviderIdle(); + if (this.turnState.status !== "idle") { + throw new Error("OpenCode is still stopping the previous turn"); } - await this.awaitPendingAbortBeforeStartingTurn(); this.runningToolCalls.clear(); this.subAgentsByCallId.clear(); @@ -3127,8 +3148,7 @@ class OpenCodeAgentSession implements AgentSession { } const turnId = this.createTurnId(); - this.activeForegroundTurnId = turnId; - this.activeForegroundTurnSource = "paseo"; + this.turnState = { status: "running", turnId }; this.notifySubscribers({ type: "turn_started", provider: "opencode" }, turnId); const slashCommand = await this.resolveSlashCommandInvocation(prompt); @@ -3555,6 +3575,20 @@ class OpenCodeAgentSession implements AgentSession { if (!event) { return; } + if ( + this.turnState.status === "stopping" && + getOpenCodeEventSessionId(event) === this.sessionId + ) { + if (isOpenCodeTerminalEvent(event, this.sessionId)) { + this.finishStoppingTurn(); + } + this.traceOpenCode("provider.opencode.event.skip", { + n: eventCount, + reason: "turn_stopping", + type: event.type, + }); + return; + } const translated = await this.translateEvent(event); const foregroundEvents: AgentStreamEvent[] = []; for (const translatedEvent of translated) { @@ -3564,8 +3598,8 @@ class OpenCodeAgentSession implements AgentSession { foregroundEvents.push(translatedEvent); } } - if (!turnId && this.shouldStartExternalDrivenTurn(event, foregroundEvents)) { - turnId = this.startExternalDrivenTurn(); + if (!turnId && this.shouldStartAutonomousTurn(event, foregroundEvents)) { + turnId = this.startAutonomousTurn(); } if (!turnId) { this.emitBackgroundPermissionRequests(foregroundEvents); @@ -3576,18 +3610,6 @@ class OpenCodeAgentSession implements AgentSession { }); return; } - if (this.suppressTerminalUntilNextUserMessage) { - if (isOpenCodeUserMessageEvent(event, this.sessionId)) { - this.suppressTerminalUntilNextUserMessage = false; - } else if (isOpenCodeTerminalEvent(event, this.sessionId)) { - this.traceOpenCode("provider.opencode.event.skip", { - n: eventCount, - reason: "stale_interrupt_terminal", - type: event.type, - }); - return; - } - } this.traceOpenCode("provider.opencode.parsed_event", { turnId, n: eventCount, @@ -3625,30 +3647,39 @@ class OpenCodeAgentSession implements AgentSession { } } - private shouldStartExternalDrivenTurn( + private shouldStartAutonomousTurn( event: OpenCodeEvent, foregroundEvents: readonly AgentStreamEvent[], ): boolean { + if (this.turnState.status !== "idle") { + return false; + } + if (getOpenCodeEventSessionId(event) !== this.sessionId) { + return false; + } + // OpenCode publishes the persisted user message before it marks the runner + // busy. That message is the earliest unambiguous boundary for a plugin- + // initiated parent turn; session metadata and assistant echoes are not. + if (event.type === "message.updated" && event.properties.info.role === "user") { + return true; + } if (!this.externallyDriven) { return false; } - if (this.activeForegroundTurnId) { - return false; - } - if (foregroundEvents.some((foregroundEvent) => !toTerminalTurnEvent(foregroundEvent))) { + if ( + foregroundEvents.some( + (foregroundEvent) => + foregroundEvent.type !== "thread_started" && !toTerminalTurnEvent(foregroundEvent), + ) + ) { return true; } - return ( - event.type === "session.status" && - event.properties.sessionID === this.sessionId && - event.properties.status.type === "busy" - ); + return event.type === "session.status" && event.properties.status.type === "busy"; } - private startExternalDrivenTurn(): string { + private startAutonomousTurn(): string { const turnId = this.createTurnId(); - this.activeForegroundTurnId = turnId; - this.activeForegroundTurnSource = "external"; + this.turnState = { status: "running", turnId }; this.runningToolCalls.clear(); this.subAgentsByCallId.clear(); this.subAgentCallIdByChildSessionId.clear(); @@ -3680,12 +3711,36 @@ class OpenCodeAgentSession implements AgentSession { } this.pendingUserMessageText = null; this.pendingClientMessageId = null; - this.activeForegroundTurnId = null; - this.activeForegroundTurnSource = null; + this.turnState = { status: "idle" }; this.abortController = null; this.notifySubscribers(event, turnId); } + private beginStoppingTurn(turnId: string): void { + if (this.turnState.status !== "running" || this.turnState.turnId !== turnId) { + return; + } + this.synthesizeInterruptedToolCalls(turnId); + this.pendingUserMessageText = null; + this.pendingClientMessageId = null; + this.abortController = null; + this.turnState = { status: "stopping", idle: createDeferred() }; + this.notifySubscribers( + { type: "turn_canceled", provider: "opencode", reason: "interrupted" }, + turnId, + ); + } + + private finishStoppingTurn(): void { + if (this.turnState.status !== "stopping") { + return; + } + const stopping = this.turnState; + resetOpenCodeTurnTrackingState(this.createTranslationState()); + this.turnState = { status: "idle" }; + stopping.idle.resolve(); + } + private trackToolCall(item: ToolCallTimelineItem): void { if (item.status === "running") { this.runningToolCalls.set(item.callId, item); @@ -3932,7 +3987,10 @@ class OpenCodeAgentSession implements AgentSession { logger: this.logger, }); await this.deleteProviderSessionIfEphemeral(); - this.activeForegroundTurnId = null; + if (this.turnState.status === "stopping") { + this.turnState.idle.resolve(); + } + this.turnState = { status: "idle" }; } finally { await this.releaseServer?.(); this.releaseServer = null; diff --git a/packages/server/src/server/agent/providers/opencode/test-utils/test-opencode-harness.ts b/packages/server/src/server/agent/providers/opencode/test-utils/test-opencode-harness.ts index 3c3ac4dc7..b59399279 100644 --- a/packages/server/src/server/agent/providers/opencode/test-utils/test-opencode-harness.ts +++ b/packages/server/src/server/agent/providers/opencode/test-utils/test-opencode-harness.ts @@ -89,6 +89,7 @@ export class TestOpenCodeClient { sessionGet: [] as unknown[], sessionMessages: [] as unknown[], sessionPromptAsync: [] as unknown[], + sessionStatus: [] as unknown[], sessionSummarize: [] as unknown[], sessionUpdate: [] as unknown[], }; @@ -102,9 +103,13 @@ export class TestOpenCodeClient { permissionReplyResponse: OpenCodeResponse = {}; providerListResponse: OpenCodeResponse = { data: { connected: [], all: [] } }; providerListImplementation: (() => Promise) | null = null; + globalEventImplementation: + | ((options: unknown) => Promise<{ stream: AsyncIterable }>) + | null = null; questionRejectResponse: OpenCodeResponse = {}; questionReplyResponse: OpenCodeResponse = {}; sessionAbortResponse: OpenCodeResponse = {}; + sessionAbortImplementation: ((parameters: unknown) => Promise) | null = null; sessionCommandError: unknown = null; sessionCommandEvents: unknown[] = [idleEvent()]; sessionCommandResponse: OpenCodeResponse = {}; @@ -118,6 +123,7 @@ export class TestOpenCodeClient { sessionMessagesResponse: OpenCodeResponse = { data: [] }; sessionPromptAsyncEvents: unknown[] = [idleEvent()]; sessionPromptAsyncResponse: OpenCodeResponse = {}; + sessionStatusResponse: OpenCodeResponse = { data: {} }; sessionSummarizeEvents: unknown[] = [idleEvent()]; sessionSummarizeResponse: OpenCodeResponse = { data: {} }; sessionUpdateResponse: OpenCodeResponse = {}; @@ -162,6 +168,9 @@ export class TestOpenCodeClient { global: { event: async (options: unknown) => { this.calls.globalEvent.push(options); + if (this.globalEventImplementation) { + return await this.globalEventImplementation(options); + } const signal = (options as { signal?: AbortSignal }).signal; return { stream: signal ? stopEventStreamOnAbort(this.eventStream, signal) : this.eventStream, @@ -205,7 +214,9 @@ export class TestOpenCodeClient { session: { abort: async (parameters: unknown) => { this.calls.sessionAbort.push(parameters); - return this.sessionAbortResponse; + return this.sessionAbortImplementation + ? await this.sessionAbortImplementation(parameters) + : this.sessionAbortResponse; }, command: async (parameters: unknown) => { this.calls.sessionCommand.push(parameters); @@ -247,6 +258,10 @@ export class TestOpenCodeClient { } return this.sessionPromptAsyncResponse; }, + status: async (parameters: unknown) => { + this.calls.sessionStatus.push(parameters); + return this.sessionStatusResponse; + }, summarize: async (parameters: unknown) => { this.calls.sessionSummarize.push(parameters); for (const event of this.sessionSummarizeEvents) {