diff --git a/apps/server/scripts/acp-mock-agent.ts b/apps/server/scripts/acp-mock-agent.ts index 0d89775844d..bc7828dd854 100644 --- a/apps/server/scripts/acp-mock-agent.ts +++ b/apps/server/scripts/acp-mock-agent.ts @@ -19,6 +19,23 @@ const emitInterleavedAssistantToolCalls = const emitGenericToolPlaceholders = process.env.T3_ACP_EMIT_GENERIC_TOOL_PLACEHOLDERS === "1"; const emitAskQuestion = process.env.T3_ACP_EMIT_ASK_QUESTION === "1"; const emitXAiAskUserQuestion = process.env.T3_ACP_EMIT_XAI_ASK_USER_QUESTION === "1"; +const emitXAiPromptCompleteThenHang = process.env.T3_ACP_EMIT_XAI_PROMPT_COMPLETE_THEN_HANG === "1"; +const emitForeignSessionUpdates = process.env.T3_ACP_EMIT_FOREIGN_SESSION_UPDATES === "1"; +const hangPromptForever = process.env.T3_ACP_HANG_PROMPT_FOREVER === "1"; +const hangFirstPromptForever = process.env.T3_ACP_HANG_FIRST_PROMPT_FOREVER === "1"; +const emitLateUpdateAfterCancel = process.env.T3_ACP_EMIT_LATE_UPDATE_AFTER_CANCEL === "1"; +const omitXAiPromptCompleteStopReason = + process.env.T3_ACP_OMIT_XAI_PROMPT_COMPLETE_STOP_REASON === "1"; +const failLoadSession = process.env.T3_ACP_FAIL_LOAD_SESSION === "1"; +const emitLoadReplay = process.env.T3_ACP_EMIT_LOAD_REPLAY === "1"; +const hangLoadSessionAfterReplay = process.env.T3_ACP_HANG_LOAD_SESSION_AFTER_REPLAY === "1"; +const delayLoadSessionAfterReplay = process.env.T3_ACP_DELAY_LOAD_SESSION_AFTER_REPLAY === "1"; +const loadSessionDelayMs = Number(process.env.T3_ACP_LOAD_SESSION_DELAY_MS ?? "5000"); +const emitStaleXAiPromptCompleteBeforeSecondHang = + process.env.T3_ACP_EMIT_STALE_XAI_PROMPT_COMPLETE_BEFORE_SECOND_HANG === "1"; +const emitOverlappingXAiPromptCompleteOutOfOrder = + process.env.T3_ACP_EMIT_OVERLAPPING_XAI_PROMPT_COMPLETE_OUT_OF_ORDER === "1"; +const failPrompt = process.env.T3_ACP_FAIL_PROMPT === "1"; const failSetConfigOption = process.env.T3_ACP_FAIL_SET_CONFIG_OPTION === "1"; const exitOnSetConfigOption = process.env.T3_ACP_EXIT_ON_SET_CONFIG_OPTION === "1"; const promptResponseText = process.env.T3_ACP_PROMPT_RESPONSE_TEXT; @@ -36,8 +53,21 @@ let parameterizedModelPicker = false; let currentReasoning = "medium"; let currentContext = "272k"; let currentFast = false; +let promptCount = 0; +let overlappingFirstPromptId: string | undefined; const cancelledSessions = new Set(); +function promptIdFromRequestMeta( + request: Pick, +): string | undefined { + const meta = request._meta; + if (meta === null || typeof meta !== "object") { + return undefined; + } + const promptId = meta.promptId ?? meta.requestId; + return typeof promptId === "string" && promptId.length > 0 ? promptId : undefined; +} + function logExit(reason: string): void { if (!exitLogPath) { return; @@ -45,6 +75,10 @@ function logExit(reason: string): void { NodeFS.appendFileSync(exitLogPath, `${reason}\n`, "utf8"); } +function writeJsonRpcNotification(method: string, params: unknown): void { + process.stdout.write(`${JSON.stringify({ jsonrpc: "2.0", method, params })}\n`); +} + process.once("SIGTERM", () => { logExit("SIGTERM"); process.exit(0); @@ -284,22 +318,66 @@ const program = Effect.gen(function* () { }), ); + const emitLoadReplayNotifications = (requestedSessionId: string) => { + writeJsonRpcNotification("session/update", { + _meta: { isReplay: true }, + sessionId: requestedSessionId, + update: { + sessionUpdate: "tool_call", + toolCallId: "replay-tool-1", + title: "Replay tool", + kind: "search", + status: "completed", + }, + }); + writeJsonRpcNotification("session/update", { + _meta: { isReplay: true }, + sessionId: requestedSessionId, + update: { + sessionUpdate: "agent_message_chunk", + content: { type: "text", text: "replayed assistant text" }, + }, + }); + }; + yield* agent.handleLoadSession((request) => - agent.client - .sessionUpdate({ - sessionId: String(request.sessionId ?? sessionId), + Effect.gen(function* () { + const requestedSessionId = String(request.sessionId ?? sessionId); + if (failLoadSession) { + return yield* AcpError.AcpRequestError.internalError("Mock load session failure"); + } + if (hangLoadSessionAfterReplay || delayLoadSessionAfterReplay) { + emitLoadReplayNotifications(requestedSessionId); + yield* agent.client.sessionUpdate({ + sessionId: requestedSessionId, + update: { + sessionUpdate: "user_message_chunk", + content: { type: "text", text: "replay-tail" }, + }, + }); + yield* Effect.sleep(loadSessionDelayMs); + return { + modes: modeState(), + models: modelState(), + configOptions: configOptions(), + }; + } + if (emitLoadReplay) { + emitLoadReplayNotifications(requestedSessionId); + } + yield* agent.client.sessionUpdate({ + sessionId: requestedSessionId, update: { sessionUpdate: "user_message_chunk", content: { type: "text", text: "replay" }, }, - }) - .pipe( - Effect.as({ - modes: modeState(), - models: modelState(), - configOptions: configOptions(), - }), - ), + }); + return { + modes: modeState(), + models: modelState(), + configOptions: configOptions(), + }; + }), ); yield* agent.handleSetSessionModel((request) => @@ -356,19 +434,152 @@ const program = Effect.gen(function* () { ); yield* agent.handleCancel(({ sessionId }) => - Effect.sync(() => { - cancelledSessions.add(String(sessionId ?? "mock-session-1")); + Effect.gen(function* () { + const cancelledSessionId = String(sessionId ?? "mock-session-1"); + cancelledSessions.add(cancelledSessionId); + if (emitLateUpdateAfterCancel) { + yield* Effect.sleep("50 millis"); + yield* Effect.sync(() => { + writeJsonRpcNotification("session/update", { + sessionId: cancelledSessionId, + update: { + sessionUpdate: "agent_message_chunk", + content: { type: "text", text: "late after cancel" }, + }, + }); + }); + } }), ); yield* agent.handlePrompt((request) => Effect.gen(function* () { const requestedSessionId = String(request.sessionId ?? sessionId); + promptCount += 1; if (Number.isFinite(promptDelayMs) && promptDelayMs > 0) { yield* Effect.sleep(`${promptDelayMs} millis`); } + if (failPrompt) { + return yield* AcpError.AcpRequestError.internalError("Mock prompt failure"); + } + + if (emitStaleXAiPromptCompleteBeforeSecondHang && promptCount === 1) { + return { + stopReason: "end_turn", + _meta: { + promptId: "mock-stale-xai-prompt-1", + requestId: "mock-stale-xai-prompt-1", + }, + }; + } + + if (emitStaleXAiPromptCompleteBeforeSecondHang && promptCount === 2) { + const currentPromptId = promptIdFromRequestMeta(request) ?? "mock-current-xai-prompt-2"; + writeJsonRpcNotification("_x.ai/session/prompt_complete", { + sessionId: requestedSessionId, + promptId: "mock-stale-xai-prompt-1", + stopReason: "end_turn", + agentResult: null, + }); + + writeJsonRpcNotification("_x.ai/session/prompt_complete", { + sessionId: requestedSessionId, + promptId: currentPromptId, + stopReason: "end_turn", + agentResult: null, + }); + + return yield* Effect.never; + } + + if (emitOverlappingXAiPromptCompleteOutOfOrder && promptCount === 1) { + overlappingFirstPromptId = promptIdFromRequestMeta(request); + return yield* Effect.never; + } + + if (emitOverlappingXAiPromptCompleteOutOfOrder && promptCount === 2) { + const secondPromptId = promptIdFromRequestMeta(request); + if (overlappingFirstPromptId !== undefined && secondPromptId !== undefined) { + writeJsonRpcNotification("_x.ai/session/prompt_complete", { + sessionId: requestedSessionId, + promptId: secondPromptId, + stopReason: "end_turn", + agentResult: null, + }); + writeJsonRpcNotification("_x.ai/session/prompt_complete", { + sessionId: requestedSessionId, + promptId: overlappingFirstPromptId, + stopReason: "end_turn", + agentResult: null, + }); + } + return yield* Effect.never; + } + + if (hangPromptForever || (hangFirstPromptForever && promptCount === 1)) { + return yield* Effect.never; + } + + if (emitXAiPromptCompleteThenHang) { + writeJsonRpcNotification("session/update", { + sessionId: requestedSessionId, + update: { + sessionUpdate: "agent_message_chunk", + content: { type: "text", text: "hello from " }, + }, + }); + + if (emitForeignSessionUpdates) { + writeJsonRpcNotification("session/update", { + sessionId: "mock-child-session-1", + update: { + sessionUpdate: "agent_message_chunk", + content: { type: "text", text: "child before completion" }, + }, + }); + } + + writeJsonRpcNotification("_x.ai/session/prompt_complete", { + sessionId: requestedSessionId, + promptId: promptIdFromRequestMeta(request) ?? "mock-xai-prompt-1", + ...(omitXAiPromptCompleteStopReason ? {} : { stopReason: "end_turn" }), + agentResult: null, + }); + + if (emitForeignSessionUpdates) { + writeJsonRpcNotification("session/update", { + sessionId: "mock-child-session-1", + update: { + sessionUpdate: "tool_call", + toolCallId: "child-tool-call-1", + title: "Child-only tool", + kind: "other", + status: "pending", + rawInput: {}, + }, + }); + writeJsonRpcNotification("session/update", { + sessionId: "mock-child-session-1", + update: { + sessionUpdate: "agent_message_chunk", + content: { type: "text", text: "child after completion" }, + }, + }); + } + + writeJsonRpcNotification("session/update", { + sessionId: requestedSessionId, + update: { + sessionUpdate: "agent_message_chunk", + content: { type: "text", text: "mock" }, + }, + }); + + return yield* Effect.never; + } + if (emitInterleavedAssistantToolCalls) { const toolCallId = "tool-call-1"; @@ -599,6 +810,42 @@ const program = Effect.gen(function* () { return { stopReason: "end_turn" }; } + if (emitForeignSessionUpdates) { + yield* agent.client.sessionUpdate({ + sessionId: requestedSessionId, + update: { + sessionUpdate: "agent_message_chunk", + content: { type: "text", text: "root before child" }, + }, + }); + yield* agent.client.sessionUpdate({ + sessionId: "mock-child-session-1", + update: { + sessionUpdate: "agent_message_chunk", + content: { type: "text", text: "child content" }, + }, + }); + yield* agent.client.sessionUpdate({ + sessionId: "mock-child-session-1", + update: { + sessionUpdate: "tool_call", + toolCallId: "child-tool-call-1", + title: "Child-only tool", + kind: "other", + status: "pending", + rawInput: {}, + }, + }); + yield* agent.client.sessionUpdate({ + sessionId: requestedSessionId, + update: { + sessionUpdate: "agent_message_chunk", + content: { type: "text", text: " root after child" }, + }, + }); + return { stopReason: "end_turn" }; + } + yield* agent.client.sessionUpdate({ sessionId: requestedSessionId, update: { diff --git a/apps/server/src/provider/Layers/CursorAdapter.test.ts b/apps/server/src/provider/Layers/CursorAdapter.test.ts index 9795e5a0680..73dc0967622 100644 --- a/apps/server/src/provider/Layers/CursorAdapter.test.ts +++ b/apps/server/src/provider/Layers/CursorAdapter.test.ts @@ -113,6 +113,23 @@ async function waitForFileContent(filePath: string, attempts = 40) { throw new Error(`Timed out waiting for file content at ${filePath}`); } +function waitForJsonLogMatch( + filePath: string, + predicate: (entry: Record) => boolean, + attempts = 40, +) { + return Effect.gen(function* () { + for (let attempt = 0; attempt < attempts; attempt += 1) { + const requests = yield* Effect.promise(() => readJsonLines(filePath)); + if (requests.some(predicate)) { + return requests; + } + yield* Effect.yieldNow; + } + return yield* Effect.promise(() => readJsonLines(filePath)); + }); +} + // Tests mutate `ServerSettingsService` mid-flight (e.g. setting // `providers.cursor.binaryPath` to a mock ACP wrapper). The adapter // captures `cursorSettings` once at construction, so without a resolver @@ -1004,7 +1021,10 @@ cursorAdapterTestLayer("CursorAdapterLive", (it) => { assert.equal(turnCompleted.payload.stopReason, "cancelled"); } - const requests = yield* Effect.promise(() => readJsonLines(requestLogPath)); + const requests = yield* waitForJsonLogMatch( + requestLogPath, + (entry) => entry.method === "session/cancel", + ); assert.isTrue(requests.some((entry) => entry.method === "session/cancel")); assert.isTrue( requests.some( diff --git a/apps/server/src/provider/Layers/CursorAdapter.ts b/apps/server/src/provider/Layers/CursorAdapter.ts index 9760b2f81fb..59788a2d225 100644 --- a/apps/server/src/provider/Layers/CursorAdapter.ts +++ b/apps/server/src/provider/Layers/CursorAdapter.ts @@ -785,6 +785,9 @@ export function makeCursorAdapter( Stream.mapEffect(acp.getEvents(), (event) => Effect.gen(function* () { switch (event._tag) { + case "EventStreamBarrier": + yield* Deferred.succeed(event.acknowledge, undefined); + return; case "ModeChanged": return; case "AssistantItemStarted": diff --git a/apps/server/src/provider/Layers/GrokAdapter.test.ts b/apps/server/src/provider/Layers/GrokAdapter.test.ts index c871e3c2fc4..7b6f0972ae8 100644 --- a/apps/server/src/provider/Layers/GrokAdapter.test.ts +++ b/apps/server/src/provider/Layers/GrokAdapter.test.ts @@ -10,20 +10,23 @@ import * as Deferred from "effect/Deferred"; import * as Effect from "effect/Effect"; import * as Fiber from "effect/Fiber"; import * as Layer from "effect/Layer"; +import * as Ref from "effect/Ref"; import * as Schema from "effect/Schema"; import * as Stream from "effect/Stream"; +import * as TestClock from "effect/testing/TestClock"; import { ApprovalRequestId, GrokSettings, ProviderDriverKind, - ThreadId, ProviderInstanceId, + ThreadId, + TurnId, type ProviderRuntimeEvent, } from "@t3tools/contracts"; import { ServerConfig } from "../../config.ts"; -import { makeGrokAdapter } from "./GrokAdapter.ts"; +import { grokPromptSettlementBelongsToContext, makeGrokAdapter } from "./GrokAdapter.ts"; const decodeGrokSettings = Schema.decodeSync(GrokSettings); const __dirname = NodePath.dirname(NodeURL.fileURLToPath(import.meta.url)); @@ -45,7 +48,11 @@ exec ${JSON.stringify(mockAgentCommand)} ${JSON.stringify(mockAgentPath)} "$@" return wrapperPath; } -function waitForFileContent(filePath: string, attempts = 40): Effect.Effect { +function waitForFileContent( + filePath: string, + attempts = 40, + expectedContent?: string, +): Effect.Effect { const readAttempt = (remainingAttempts: number): Effect.Effect => Effect.gen(function* () { if (remainingAttempts <= 0) { @@ -54,7 +61,10 @@ function waitForFileContent(filePath: string, attempts = 40): Effect.Effect NodeFSP.readFile(filePath, "utf8")).pipe( Effect.orElseSucceed(() => ""), ); - if (raw.trim().length > 0) { + if ( + raw.trim().length > 0 && + (expectedContent === undefined || raw.includes(expectedContent)) + ) { return raw; } yield* Effect.sleep("25 millis"); @@ -79,6 +89,39 @@ const grokAdapterTestLayer = ServerConfig.layerTest(process.cwd(), { const makeTestAdapter = (binaryPath: string, options?: Parameters[1]) => makeGrokAdapter(decodeGrokSettings({ binaryPath }), options).pipe(Effect.orDie); +it("requires a settlement to match the live Grok turn", () => { + const staleTurnId = TurnId.make("stale-turn"); + const replacementTurnId = TurnId.make("replacement-turn"); + + assert.isFalse( + grokPromptSettlementBelongsToContext({ + liveAcpSessionId: "session-1", + expectedAcpSessionId: "session-1", + liveActiveTurnId: replacementTurnId, + liveSessionActiveTurnId: replacementTurnId, + turnId: staleTurnId, + }), + ); + assert.isFalse( + grokPromptSettlementBelongsToContext({ + liveAcpSessionId: "replacement-session", + expectedAcpSessionId: "stale-session", + liveActiveTurnId: staleTurnId, + liveSessionActiveTurnId: staleTurnId, + turnId: staleTurnId, + }), + ); + assert.isTrue( + grokPromptSettlementBelongsToContext({ + liveAcpSessionId: "session-1", + expectedAcpSessionId: "session-1", + liveActiveTurnId: staleTurnId, + liveSessionActiveTurnId: staleTurnId, + turnId: staleTurnId, + }), + ); +}); + it.layer(grokAdapterTestLayer)("GrokAdapterLive", (it) => { it.effect("starts a session and maps mock ACP prompt flow to runtime events", () => Effect.gen(function* () { @@ -175,6 +218,783 @@ it.layer(grokAdapterTestLayer)("GrokAdapterLive", (it) => { }), ); + it.effect("reports a Grok session running only while the prompt is in flight", () => + Effect.gen(function* () { + const threadId = ThreadId.make("grok-session-ready-after-prompt"); + const wrapperPath = yield* Effect.promise(() => + makeMockGrokWrapper({ + T3_ACP_EMIT_TOOL_CALLS: "1", + }), + ); + const adapter = yield* makeTestAdapter(wrapperPath); + const requestOpened = + yield* Deferred.make>(); + const eventsFiber = yield* Stream.runForEach(adapter.streamEvents, (event) => + event.type === "request.opened" + ? Deferred.succeed(requestOpened, event).pipe(Effect.ignore) + : Effect.void, + ).pipe(Effect.forkChild); + + yield* adapter.startSession({ + threadId, + provider: ProviderDriverKind.make("grok"), + cwd: process.cwd(), + runtimeMode: "approval-required", + modelSelection: { instanceId: ProviderInstanceId.make("grok"), model: "grok-build" }, + }); + + const sendTurnFiber = yield* adapter + .sendTurn({ threadId, input: "check lifecycle", attachments: [] }) + .pipe(Effect.forkChild); + const requestOpenedEvent = yield* Deferred.await(requestOpened); + + const runningSessions = yield* adapter.listSessions(); + const runningSession = runningSessions.find((session) => session.threadId === threadId); + assert.equal(runningSession?.status, "running"); + assert.isDefined(runningSession?.activeTurnId); + + yield* adapter.respondToRequest( + threadId, + ApprovalRequestId.make(String(requestOpenedEvent.requestId)), + "accept", + ); + yield* Fiber.join(sendTurnFiber); + + const readySessions = yield* adapter.listSessions(); + const readySession = readySessions.find((session) => session.threadId === threadId); + assert.equal(readySession?.status, "ready"); + assert.isUndefined(readySession?.activeTurnId); + + yield* Fiber.interrupt(eventsFiber); + yield* adapter.stopSession(threadId); + }), + ); + + it.effect("restores ready without completing an unstarted turn when preparation fails", () => + Effect.gen(function* () { + const threadId = ThreadId.make("grok-preparation-failure-while-connecting"); + const wrapperPath = yield* Effect.promise(() => makeMockGrokWrapper()); + const adapter = yield* makeTestAdapter(wrapperPath); + + const runtimeEvents: ProviderRuntimeEvent[] = []; + const runtimeEventsFiber = yield* Stream.runForEach(adapter.streamEvents, (event) => + Effect.sync(() => { + runtimeEvents.push(event); + }), + ).pipe(Effect.forkChild); + + yield* adapter.startSession({ + threadId, + provider: ProviderDriverKind.make("grok"), + cwd: process.cwd(), + runtimeMode: "full-access", + modelSelection: { instanceId: ProviderInstanceId.make("grok"), model: "grok-build" }, + }); + + const error = yield* Effect.flip( + adapter.sendTurn({ + threadId, + input: "prepare invalid attachment", + attachments: [ + { + type: "image", + id: "missing-image", + name: "missing.png", + mimeType: "image/png", + sizeBytes: 1, + }, + ], + }), + ); + for (let yieldAttempt = 0; yieldAttempt < 4; yieldAttempt += 1) { + yield* Effect.yieldNow; + } + + const turnCompletedEvent = runtimeEvents.find( + (event): event is Extract => + event.type === "turn.completed", + ); + const readySessions = yield* adapter.listSessions(); + const readySession = readySessions.find((session) => session.threadId === threadId); + + assert.equal(error._tag, "ProviderAdapterRequestError"); + assert.isUndefined(turnCompletedEvent); + assert.equal(readySession?.status, "ready"); + assert.isUndefined(readySession?.activeTurnId); + + yield* Fiber.interrupt(runtimeEventsFiber); + yield* adapter.stopSession(threadId); + }), + ); + + it.effect("completes a Grok turn from xAI prompt completion when the prompt RPC hangs", () => + Effect.gen(function* () { + const threadId = ThreadId.make("grok-xai-prompt-complete-fallback"); + const wrapperPath = yield* Effect.promise(() => + makeMockGrokWrapper({ + T3_ACP_EMIT_XAI_PROMPT_COMPLETE_THEN_HANG: "1", + T3_ACP_EMIT_FOREIGN_SESSION_UPDATES: "1", + }), + ); + const adapter = yield* makeTestAdapter(wrapperPath); + + const runtimeEvents: ProviderRuntimeEvent[] = []; + const turnCompleted = yield* Deferred.make(); + const runtimeEventsFiber = yield* Stream.runForEach(adapter.streamEvents, (event) => + Effect.sync(() => { + runtimeEvents.push(event); + }).pipe( + Effect.andThen( + event.type === "turn.completed" + ? Deferred.succeed(turnCompleted, undefined) + : Effect.void, + ), + ), + ).pipe(Effect.forkChild); + + yield* adapter.startSession({ + threadId, + provider: ProviderDriverKind.make("grok"), + cwd: process.cwd(), + runtimeMode: "full-access", + modelSelection: { instanceId: ProviderInstanceId.make("grok"), model: "grok-build" }, + }); + + const sendTurnResult = yield* adapter.sendTurn({ + threadId, + input: "exercise fallback", + attachments: [], + }); + + yield* Deferred.await(turnCompleted); + for (let yieldAttempt = 0; yieldAttempt < 8; yieldAttempt += 1) { + yield* Effect.yieldNow; + } + const readySessions = yield* adapter.listSessions(); + const readySession = readySessions.find((session) => session.threadId === threadId); + const turnCompletedEvent = runtimeEvents.find( + (event): event is Extract => + event.type === "turn.completed", + ); + const eventTypes = runtimeEvents.map((event) => event.type); + const content = runtimeEvents + .filter( + (event): event is Extract => + event.type === "content.delta" && String(event.threadId) === String(threadId), + ) + .map((event) => event.payload.delta) + .join(""); + const terminalIndex = runtimeEvents.findIndex( + (event) => event.type === "turn.completed" && String(event.threadId) === String(threadId), + ); + const turnOutputTypes = new Set([ + "content.delta", + "item.started", + "item.updated", + "item.completed", + "turn.plan.updated", + ]); + const outputAfterTerminal = runtimeEvents + .slice(terminalIndex + 1) + .filter( + (event) => String(event.threadId) === String(threadId) && turnOutputTypes.has(event.type), + ); + const toolTitles = runtimeEvents.flatMap((event) => + event.type === "item.updated" && event.payload.title ? [event.payload.title] : [], + ); + + assert.equal(sendTurnResult.threadId, threadId); + assert.include(eventTypes, "turn.completed"); + assert.equal(content, "hello from mock"); + assert.isAtLeast(terminalIndex, 0); + assert.deepEqual(outputAfterTerminal, []); + assert.notInclude(toolTitles, "Child-only tool"); + assert.equal(turnCompletedEvent?.payload.stopReason, "end_turn"); + assert.equal(readySession?.status, "ready"); + assert.isUndefined(readySession?.activeTurnId); + + yield* Fiber.interrupt(runtimeEventsFiber); + yield* adapter.stopSession(threadId); + }), + ); + + it.effect("retains turn transcript when sendTurn is interrupted after prompt success", () => + Effect.gen(function* () { + const threadId = ThreadId.make("grok-send-turn-interrupt-after-prompt"); + const wrapperPath = yield* Effect.promise(() => + makeMockGrokWrapper({ + T3_ACP_EMIT_XAI_PROMPT_COMPLETE_THEN_HANG: "1", + }), + ); + const adapter = yield* makeTestAdapter(wrapperPath); + const contentDelta = yield* Deferred.make(); + const runtimeEventsFiber = yield* Stream.runForEach(adapter.streamEvents, (event) => + event.type === "content.delta" ? Deferred.succeed(contentDelta, undefined) : Effect.void, + ).pipe(Effect.forkChild); + + yield* adapter.startSession({ + threadId, + provider: ProviderDriverKind.make("grok"), + cwd: process.cwd(), + runtimeMode: "full-access", + modelSelection: { instanceId: ProviderInstanceId.make("grok"), model: "grok-build" }, + }); + + const sendTurnFiber = yield* adapter + .sendTurn({ + threadId, + input: "interrupt after prompt", + attachments: [], + }) + .pipe(Effect.forkChild); + + yield* Deferred.await(contentDelta); + for (let yieldAttempt = 0; yieldAttempt < 6; yieldAttempt += 1) { + yield* Effect.yieldNow; + } + yield* Fiber.interrupt(sendTurnFiber); + for (let yieldAttempt = 0; yieldAttempt < 4; yieldAttempt += 1) { + yield* Effect.yieldNow; + } + + const snapshot = yield* adapter.readThread(threadId); + assert.equal(snapshot.turns.length, 1); + assert.equal(snapshot.turns[0]?.items.length, 1); + + yield* Fiber.interrupt(runtimeEventsFiber); + yield* adapter.stopSession(threadId); + }), + ); + + it.effect("does not report a synthetic stop reason when xAI omits one", () => + Effect.gen(function* () { + const threadId = ThreadId.make("grok-xai-prompt-complete-missing-stop-reason"); + const wrapperPath = yield* Effect.promise(() => + makeMockGrokWrapper({ + T3_ACP_EMIT_XAI_PROMPT_COMPLETE_THEN_HANG: "1", + T3_ACP_OMIT_XAI_PROMPT_COMPLETE_STOP_REASON: "1", + }), + ); + const adapter = yield* makeTestAdapter(wrapperPath); + + const runtimeEvents: ProviderRuntimeEvent[] = []; + const turnCompleted = yield* Deferred.make(); + const runtimeEventsFiber = yield* Stream.runForEach(adapter.streamEvents, (event) => + Effect.sync(() => { + runtimeEvents.push(event); + }).pipe( + Effect.andThen( + event.type === "turn.completed" + ? Deferred.succeed(turnCompleted, undefined) + : Effect.void, + ), + ), + ).pipe(Effect.forkChild); + + yield* adapter.startSession({ + threadId, + provider: ProviderDriverKind.make("grok"), + cwd: process.cwd(), + runtimeMode: "full-access", + modelSelection: { instanceId: ProviderInstanceId.make("grok"), model: "grok-build" }, + }); + + yield* adapter.sendTurn({ + threadId, + input: "exercise missing stop reason", + attachments: [], + }); + + yield* Deferred.await(turnCompleted); + const turnCompletedEvent = runtimeEvents.find( + (event): event is Extract => + event.type === "turn.completed", + ); + + assert.equal(turnCompletedEvent?.payload.state, "completed"); + assert.isNull(turnCompletedEvent?.payload.stopReason); + + yield* Fiber.interrupt(runtimeEventsFiber); + yield* adapter.stopSession(threadId); + }), + ); + + it.effect("lets Stop unblock a fully silent Grok prompt and accept a follow-up turn", () => + Effect.gen(function* () { + const threadId = ThreadId.make("grok-stop-after-full-silence"); + const wrapperPath = yield* Effect.promise(() => + makeMockGrokWrapper({ + T3_ACP_HANG_FIRST_PROMPT_FOREVER: "1", + }), + ); + const adapter = yield* makeTestAdapter(wrapperPath); + + const runtimeEvents: ProviderRuntimeEvent[] = []; + const runtimeEventsFiber = yield* Stream.runForEach(adapter.streamEvents, (event) => + Effect.sync(() => { + runtimeEvents.push(event); + }), + ).pipe(Effect.forkChild); + + yield* adapter.startSession({ + threadId, + provider: ProviderDriverKind.make("grok"), + cwd: process.cwd(), + runtimeMode: "full-access", + modelSelection: { instanceId: ProviderInstanceId.make("grok"), model: "grok-build" }, + }); + + yield* Effect.gen(function* () { + yield* Effect.sleep("500 millis"); + yield* adapter.interruptTurn(threadId); + }).pipe(Effect.forkChild({ startImmediately: true })); + + yield* adapter.sendTurn({ + threadId, + input: "hang forever", + attachments: [], + }); + for (let yieldAttempt = 0; yieldAttempt < 8; yieldAttempt += 1) { + yield* Effect.yieldNow; + } + + const cancelledEvents = runtimeEvents.filter( + (event): event is Extract => + event.type === "turn.completed" && String(event.threadId) === String(threadId), + ); + const readySessions = yield* adapter.listSessions(); + const readySession = readySessions.find((session) => session.threadId === threadId); + + assert.lengthOf(cancelledEvents, 1); + assert.equal(cancelledEvents[0]?.payload.state, "cancelled"); + assert.equal(readySession?.status, "ready"); + assert.isUndefined(readySession?.activeTurnId); + + const followUpEventsBefore = runtimeEvents.length; + yield* adapter.sendTurn({ + threadId, + input: "continue after stop", + attachments: [], + }); + for (let yieldAttempt = 0; yieldAttempt < 8; yieldAttempt += 1) { + yield* Effect.yieldNow; + } + + const followUpCompletedEvents = runtimeEvents + .slice(followUpEventsBefore) + .filter( + (event): event is Extract => + event.type === "turn.completed" && String(event.threadId) === String(threadId), + ); + assert.lengthOf(followUpCompletedEvents, 1); + assert.equal(followUpCompletedEvents[0]?.payload.state, "completed"); + + yield* Fiber.interrupt(runtimeEventsFiber); + yield* adapter.stopSession(threadId); + }).pipe(TestClock.withLive), + ); + + it.effect("does not let a cancelled prompt settlement consume the follow-up prompt slot", () => + Effect.gen(function* () { + const threadId = ThreadId.make("grok-cancelled-settlement-before-follow-up"); + const tempDir = yield* Effect.promise(() => + NodeFSP.mkdtemp(NodePath.join(NodeOS.tmpdir(), "grok-acp-cancel-race-")), + ); + const requestLogPath = NodePath.join(tempDir, "requests.ndjson"); + const wrapperPath = yield* Effect.promise(() => + makeMockGrokWrapper({ + T3_ACP_HANG_FIRST_PROMPT_FOREVER: "1", + T3_ACP_REQUEST_LOG_PATH: requestLogPath, + }), + ); + const adapter = yield* makeTestAdapter(wrapperPath); + + const runtimeEvents: ProviderRuntimeEvent[] = []; + const firstTurnStarted = yield* Deferred.make(); + const twoTurnsCompleted = yield* Deferred.make(); + const completedCountRef = yield* Ref.make(0); + const runtimeEventsFiber = yield* Stream.runForEach(adapter.streamEvents, (event) => + Effect.gen(function* () { + runtimeEvents.push(event); + if (String(event.threadId) !== String(threadId)) { + return; + } + if (event.type === "turn.started" && event.turnId !== undefined) { + yield* Deferred.succeed(firstTurnStarted, event.turnId).pipe(Effect.ignore); + return; + } + if (event.type !== "turn.completed") { + return; + } + const completedCount = yield* Ref.updateAndGet(completedCountRef, (count) => count + 1); + if (completedCount === 2) { + yield* Deferred.succeed(twoTurnsCompleted, undefined); + } + }), + ).pipe(Effect.forkChild); + + yield* adapter.startSession({ + threadId, + provider: ProviderDriverKind.make("grok"), + cwd: process.cwd(), + runtimeMode: "full-access", + }); + + const firstSendTurnFiber = yield* adapter + .sendTurn({ threadId, input: "cancel this prompt", attachments: [] }) + .pipe(Effect.forkChild); + const firstTurnId = yield* Deferred.await(firstTurnStarted).pipe(Effect.timeout("2 seconds")); + yield* waitForFileContent(requestLogPath, 80, '"method":"session/prompt"'); + + yield* adapter.interruptTurn(threadId, firstTurnId).pipe(Effect.timeout("2 seconds")); + const followUp = yield* adapter + .sendTurn({ threadId, input: "complete the follow-up", attachments: [] }) + .pipe(Effect.timeout("2 seconds")); + yield* Fiber.join(firstSendTurnFiber).pipe(Effect.timeout("2 seconds")); + yield* Deferred.await(twoTurnsCompleted).pipe(Effect.timeout("2 seconds")); + + const turnCompletedEvents = runtimeEvents.filter( + (event): event is Extract => + event.type === "turn.completed" && String(event.threadId) === String(threadId), + ); + const readySessions = yield* adapter.listSessions(); + const readySession = readySessions.find((session) => session.threadId === threadId); + + assert.notEqual(String(followUp.turnId), String(firstTurnId)); + assert.deepEqual( + turnCompletedEvents.map((event) => [String(event.turnId), event.payload.state]), + [ + [String(firstTurnId), "cancelled"], + [String(followUp.turnId), "completed"], + ], + ); + assert.equal(readySession?.status, "ready"); + assert.isUndefined(readySession?.activeTurnId); + + yield* Fiber.interrupt(runtimeEventsFiber); + yield* adapter.stopSession(threadId); + }).pipe(TestClock.withLive), + ); + + it.effect("drops late ACP notifications after a turn is cancelled", () => + Effect.gen(function* () { + const threadId = ThreadId.make("grok-drop-late-cancelled-notifications"); + const wrapperPath = yield* Effect.promise(() => + makeMockGrokWrapper({ + T3_ACP_HANG_PROMPT_FOREVER: "1", + T3_ACP_EMIT_LATE_UPDATE_AFTER_CANCEL: "1", + }), + ); + const lateNativeUpdate = yield* Deferred.make(); + const adapter = yield* makeTestAdapter(wrapperPath, { + nativeEventLogger: { + filePath: "memory://grok-cancelled-native-events", + write: (record: unknown) => + JSON.stringify(record).includes("late after cancel") + ? Deferred.succeed(lateNativeUpdate, undefined).pipe(Effect.asVoid) + : Effect.void, + close: () => Effect.void, + }, + }); + + const runtimeEvents: ProviderRuntimeEvent[] = []; + const turnStarted = yield* Deferred.make(); + const runtimeEventsFiber = yield* Stream.runForEach(adapter.streamEvents, (event) => + Effect.sync(() => { + runtimeEvents.push(event); + }).pipe( + Effect.andThen( + event.type === "turn.started" && + event.turnId !== undefined && + String(event.threadId) === String(threadId) + ? Deferred.succeed(turnStarted, event.turnId).pipe(Effect.asVoid) + : Effect.void, + ), + ), + ).pipe(Effect.forkChild); + + yield* adapter.startSession({ + threadId, + provider: ProviderDriverKind.make("grok"), + cwd: process.cwd(), + runtimeMode: "full-access", + }); + + const sendTurnFiber = yield* adapter + .sendTurn({ threadId, input: "cancel before the late update", attachments: [] }) + .pipe(Effect.forkChild); + const turnId = yield* Deferred.await(turnStarted).pipe(Effect.timeout("2 seconds")); + yield* adapter.interruptTurn(threadId, turnId).pipe(Effect.timeout("2 seconds")); + yield* Fiber.join(sendTurnFiber).pipe(Effect.timeout("2 seconds")); + yield* Deferred.await(lateNativeUpdate).pipe(Effect.timeout("2 seconds")); + for (let yieldAttempt = 0; yieldAttempt < 8; yieldAttempt += 1) { + yield* Effect.yieldNow; + } + + const cancelledIndex = runtimeEvents.findIndex( + (event) => + event.type === "turn.completed" && + String(event.threadId) === String(threadId) && + String(event.turnId) === String(turnId) && + event.payload.state === "cancelled", + ); + const turnOutputTypes = new Set([ + "content.delta", + "item.started", + "item.updated", + "item.completed", + "turn.plan.updated", + ]); + const outputAfterCancellation = runtimeEvents + .slice(cancelledIndex + 1) + .filter( + (event) => String(event.threadId) === String(threadId) && turnOutputTypes.has(event.type), + ); + + assert.isAtLeast(cancelledIndex, 0); + assert.deepEqual(outputAfterCancellation, []); + + yield* Fiber.interrupt(runtimeEventsFiber); + yield* adapter.stopSession(threadId); + }).pipe(TestClock.withLive), + ); + + it.effect("lets Stop cancel during the xAI completion drain window", () => + Effect.gen(function* () { + const threadId = ThreadId.make("grok-stop-during-completion-drain"); + const wrapperPath = yield* Effect.promise(() => + makeMockGrokWrapper({ + T3_ACP_EMIT_XAI_PROMPT_COMPLETE_THEN_HANG: "1", + }), + ); + const adapter = yield* makeTestAdapter(wrapperPath); + + const runtimeEvents: ProviderRuntimeEvent[] = []; + const activeTurnIdRef = yield* Ref.make(undefined); + const trailingChunkTurnId = yield* Deferred.make(); + const runtimeEventsFiber = yield* Stream.runForEach(adapter.streamEvents, (event) => + Effect.gen(function* () { + runtimeEvents.push(event); + if (String(event.threadId) !== String(threadId)) { + return; + } + if (event.type === "turn.started") { + yield* Ref.set(activeTurnIdRef, event.turnId); + } + if (event.type !== "content.delta" || event.payload.delta !== "mock") { + return; + } + const turnId = event.turnId ?? (yield* Ref.get(activeTurnIdRef)); + if (turnId === undefined) { + return; + } + yield* Deferred.succeed(trailingChunkTurnId, turnId).pipe(Effect.ignore); + }), + ).pipe(Effect.forkChild); + + yield* adapter.startSession({ + threadId, + provider: ProviderDriverKind.make("grok"), + cwd: process.cwd(), + runtimeMode: "full-access", + modelSelection: { instanceId: ProviderInstanceId.make("grok"), model: "grok-build" }, + }); + + const sendTurnFiber = yield* adapter + .sendTurn({ + threadId, + input: "cancel during completion drain", + attachments: [], + }) + .pipe(Effect.forkChild); + + const turnId = yield* Deferred.await(trailingChunkTurnId).pipe(Effect.timeout("2 seconds")); + yield* adapter.interruptTurn(threadId, turnId).pipe(Effect.timeout("2 seconds")); + yield* Fiber.join(sendTurnFiber).pipe(Effect.timeout("2 seconds")); + + const turnCompletedEvents = runtimeEvents.filter( + (event): event is Extract => + event.type === "turn.completed" && String(event.threadId) === String(threadId), + ); + const readySessions = yield* adapter.listSessions(); + const readySession = readySessions.find((session) => session.threadId === threadId); + + assert.lengthOf(turnCompletedEvents, 1); + assert.equal(turnCompletedEvents[0]?.payload.state, "cancelled"); + assert.equal(readySession?.status, "ready"); + assert.isUndefined(readySession?.activeTurnId); + + yield* Fiber.interrupt(runtimeEventsFiber); + yield* adapter.stopSession(threadId); + }), + ); + + it.effect("settles the in-flight prompt before emitting completion", () => + Effect.gen(function* () { + const threadId = ThreadId.make("grok-completion-before-next-turn"); + const wrapperPath = yield* Effect.promise(() => makeMockGrokWrapper()); + const adapter = yield* makeTestAdapter(wrapperPath); + const completedCountRef = yield* Ref.make(0); + const secondTurnCompleted = yield* Deferred.make(); + + const runtimeEventsFiber = yield* Stream.runForEach(adapter.streamEvents, (event) => { + if (event.type !== "turn.completed" || String(event.threadId) !== String(threadId)) { + return Effect.void; + } + + return Ref.modify(completedCountRef, (count) => { + const nextCount = count + 1; + return [nextCount, nextCount] as const; + }).pipe( + Effect.flatMap((count) => { + if (count === 1) { + return adapter + .sendTurn({ + threadId, + input: "second turn after completion", + attachments: [], + }) + .pipe(Effect.forkChild, Effect.asVoid); + } + if (count === 2) { + return Deferred.succeed(secondTurnCompleted, undefined).pipe(Effect.asVoid); + } + return Effect.void; + }), + ); + }).pipe(Effect.forkChild); + + yield* adapter.startSession({ + threadId, + provider: ProviderDriverKind.make("grok"), + cwd: process.cwd(), + runtimeMode: "full-access", + modelSelection: { instanceId: ProviderInstanceId.make("grok"), model: "grok-build" }, + }); + + yield* adapter.sendTurn({ + threadId, + input: "first turn", + attachments: [], + }); + yield* Deferred.await(secondTurnCompleted); + + const completedCount = yield* Ref.get(completedCountRef); + const readySessions = yield* adapter.listSessions(); + const readySession = readySessions.find((session) => session.threadId === threadId); + + assert.equal(completedCount, 2); + assert.equal(readySession?.status, "ready"); + assert.isUndefined(readySession?.activeTurnId); + + yield* Fiber.interrupt(runtimeEventsFiber); + yield* adapter.stopSession(threadId); + }), + ); + + it.effect("restores a Grok session to ready when the prompt RPC fails", () => + Effect.gen(function* () { + const threadId = ThreadId.make("grok-prompt-failure-ready"); + const wrapperPath = yield* Effect.promise(() => + makeMockGrokWrapper({ + T3_ACP_FAIL_PROMPT: "1", + }), + ); + const adapter = yield* makeTestAdapter(wrapperPath); + const runtimeEvents: ProviderRuntimeEvent[] = []; + const runtimeEventsFiber = yield* Stream.runForEach(adapter.streamEvents, (event) => + Effect.sync(() => { + runtimeEvents.push(event); + }), + ).pipe(Effect.forkChild); + + yield* adapter.startSession({ + threadId, + provider: ProviderDriverKind.make("grok"), + cwd: process.cwd(), + runtimeMode: "full-access", + modelSelection: { instanceId: ProviderInstanceId.make("grok"), model: "grok-build" }, + }); + + const error = yield* Effect.flip( + adapter.sendTurn({ + threadId, + input: "fail prompt", + attachments: [], + }), + ); + const readySessions = yield* adapter.listSessions(); + const readySession = readySessions.find((session) => session.threadId === threadId); + const failedTurnCompleted = runtimeEvents.find( + (event) => event.type === "turn.completed" && event.threadId === threadId, + ); + + assert.equal(error._tag, "ProviderAdapterRequestError"); + assert.equal(readySession?.status, "ready"); + assert.isUndefined(readySession?.activeTurnId); + assert.equal(failedTurnCompleted?.type, "turn.completed"); + if (failedTurnCompleted?.type === "turn.completed") { + assert.equal(failedTurnCompleted.payload.state, "failed"); + assert.isString(failedTurnCompleted.payload.errorMessage); + } + + yield* Fiber.interrupt(runtimeEventsFiber); + yield* adapter.stopSession(threadId); + }), + ); + + it.effect("ignores replayed session/load updates when resuming a Grok session", () => + Effect.gen(function* () { + const threadId = ThreadId.make("grok-load-replay-filter"); + const wrapperPath = yield* Effect.promise(() => + makeMockGrokWrapper({ + T3_ACP_EMIT_LOAD_REPLAY: "1", + }), + ); + const adapter = yield* makeTestAdapter(wrapperPath); + const runtimeEvents: ProviderRuntimeEvent[] = []; + const runtimeEventsFiber = yield* Stream.runForEach(adapter.streamEvents, (event) => + Effect.sync(() => { + runtimeEvents.push(event); + }), + ).pipe(Effect.forkChild); + + const session = yield* adapter.startSession({ + threadId, + provider: ProviderDriverKind.make("grok"), + cwd: process.cwd(), + runtimeMode: "full-access", + modelSelection: { instanceId: ProviderInstanceId.make("grok"), model: "grok-build" }, + resumeCursor: { schemaVersion: 1, sessionId: "mock-session-1" }, + }); + + yield* adapter.sendTurn({ + threadId, + input: "after resume", + attachments: [], + }); + + assert.deepStrictEqual(session.resumeCursor, { + schemaVersion: 1, + sessionId: "mock-session-1", + }); + assert.isFalse( + runtimeEvents.some( + (event) => event.type === "item.completed" && event.payload.title === "Replay tool", + ), + ); + assert.isFalse( + runtimeEvents.some( + (event) => + event.type === "content.delta" && event.payload.delta === "replayed assistant text", + ), + ); + + yield* Fiber.interrupt(runtimeEventsFiber); + yield* adapter.stopSession(threadId); + }), + ); + it.effect("rejects startSession when provider mismatches", () => Effect.gen(function* () { const wrapperPath = yield* Effect.promise(() => makeMockGrokWrapper()); @@ -331,6 +1151,7 @@ it.layer(grokAdapterTestLayer)("GrokAdapterLive", (it) => { assert.deepEqual(resolvedEvent.payload.answers, { "Which scope should Grok use?": "Workspace", }); + assert.equal(String(resolvedEvent.turnId), String(requestedEvent.turnId)); yield* Fiber.join(sendTurnFiber); yield* Fiber.interrupt(eventsFiber); diff --git a/apps/server/src/provider/Layers/GrokAdapter.ts b/apps/server/src/provider/Layers/GrokAdapter.ts index 40f425cbaa1..c22b2180183 100644 --- a/apps/server/src/provider/Layers/GrokAdapter.ts +++ b/apps/server/src/provider/Layers/GrokAdapter.ts @@ -22,6 +22,7 @@ import * as FileSystem from "effect/FileSystem"; import * as Option from "effect/Option"; import * as Path from "effect/Path"; import * as PubSub from "effect/PubSub"; +import * as Ref from "effect/Ref"; import * as Schema from "effect/Schema"; import * as Scope from "effect/Scope"; import * as Semaphore from "effect/Semaphore"; @@ -62,6 +63,7 @@ import { extractXAiAskUserQuestions, makeXAiAskUserQuestionCancelledResponse, makeXAiAskUserQuestionResponse, + promptResponseHasMissingXAiStopReason, XAiAskUserQuestionRequest, } from "../acp/XAiAcpExtension.ts"; import { type GrokAdapterShape } from "../Services/GrokAdapter.ts"; @@ -108,6 +110,8 @@ interface GrokSessionContext { turns: Array<{ id: TurnId; items: Array }>; lastPlanFingerprint: string | undefined; activeTurnId: TurnId | undefined; + /** Turns already interrupted; late prompt RPCs must not resurrect them. */ + interruptedTurnIds: Set; /** Number of sendTurn prompts currently in flight or being prepared. * >0 means a turn is actively running, so a new sendTurn is a steer that * continues it, and only the last remaining prompt settles the turn. */ @@ -136,10 +140,38 @@ function settlePendingUserInputsAsCancelled( ); } +function appendPromptResultToTurn( + ctx: GrokSessionContext, + turnId: TurnId, + promptParts: ReadonlyArray, + result: EffectAcpSchema.PromptResponse, +): void { + const existingTurnRecord = ctx.turns.find((turn) => turn.id === turnId); + ctx.turns = existingTurnRecord + ? ctx.turns.map((turn) => + turn.id === turnId + ? { ...turn, items: [...turn.items, { prompt: promptParts, result }] } + : turn, + ) + : [...ctx.turns, { id: turnId, items: [{ prompt: promptParts, result }] }]; +} + function isRecord(value: unknown): value is Record { return typeof value === "object" && value !== null && !Array.isArray(value); } +const resolveNotificationTurnId = (ctx: GrokSessionContext): TurnId | undefined => ctx.activeTurnId; + +const resolveCallbackTurnId = (ctx: GrokSessionContext): TurnId | undefined => ctx.activeTurnId; + +const resolveSessionCallbackTurnId = ( + sessions: ReadonlyMap, + threadId: ThreadId, +): TurnId | undefined => { + const ctx = sessions.get(threadId); + return ctx ? resolveCallbackTurnId(ctx) : undefined; +}; + function parseGrokResume(raw: unknown): { sessionId: string } | undefined { if (!isRecord(raw)) return undefined; if (raw.schemaVersion !== GROK_RESUME_VERSION) return undefined; @@ -170,6 +202,28 @@ function selectAutoApprovedPermissionOption( ); } +function completedStopReasonFromPromptResponse( + response: EffectAcpSchema.PromptResponse | undefined, +): EffectAcpSchema.StopReason | null { + if (response === undefined || promptResponseHasMissingXAiStopReason(response)) { + return null; + } + return response.stopReason; +} + +export function grokPromptSettlementBelongsToContext(input: { + readonly liveAcpSessionId: string; + readonly expectedAcpSessionId: string; + readonly liveActiveTurnId: TurnId | undefined; + readonly liveSessionActiveTurnId: TurnId | undefined; + readonly turnId: TurnId; +}): boolean { + return ( + input.liveAcpSessionId === input.expectedAcpSessionId && + (input.liveActiveTurnId === input.turnId || input.liveSessionActiveTurnId === input.turnId) + ); +} + export function makeGrokAdapter(grokSettings: GrokSettings, options?: GrokAdapterLiveOptions) { return Effect.gen(function* () { const boundInstanceId = options?.instanceId ?? ProviderInstanceId.make("grok"); @@ -240,6 +294,144 @@ export function makeGrokAdapter(grokSettings: GrokSettings, options?: GrokAdapte const withThreadLock = (threadId: string, effect: Effect.Effect) => Effect.flatMap(getThreadSemaphore(threadId), (semaphore) => semaphore.withPermit(effect)); + const settlePromptInFlight = ( + threadId: ThreadId, + turnId: TurnId, + expectedAcpSessionId: string, + options?: { + readonly errorMessage?: string; + readonly completedStopReason?: EffectAcpSchema.StopReason | null; + readonly emitTurnCompletion?: boolean; + /** Interrupt/cancel: drop every outstanding prompt slot and settle once. */ + readonly settleAllPrompts?: boolean; + }, + ) => + Effect.gen(function* () { + const liveCtx = sessions.get(threadId); + if (!liveCtx) { + return; + } + const settlementBelongsToLiveContext = grokPromptSettlementBelongsToContext({ + liveAcpSessionId: liveCtx.acpSessionId, + expectedAcpSessionId, + liveActiveTurnId: liveCtx.activeTurnId, + liveSessionActiveTurnId: liveCtx.session.activeTurnId, + turnId, + }); + if (!settlementBelongsToLiveContext) { + // interruptTurn already consumed every prompt slot for this turn. A + // late prompt result must neither emit a second terminal event nor + // consume a slot belonging to a newer turn on the same ACP session. + if ( + liveCtx.acpSessionId !== expectedAcpSessionId || + liveCtx.interruptedTurnIds.has(turnId) + ) { + return; + } + if (options?.emitTurnCompletion !== false) { + if (options?.errorMessage !== undefined) { + yield* offerRuntimeEvent({ + type: "turn.completed", + ...(yield* makeEventStamp()), + provider: PROVIDER, + threadId, + turnId, + payload: { + state: "failed", + errorMessage: options.errorMessage, + }, + }); + } else if (options?.completedStopReason !== undefined) { + yield* offerRuntimeEvent({ + type: "turn.completed", + ...(yield* makeEventStamp()), + provider: PROVIDER, + threadId, + turnId, + payload: { + state: options.completedStopReason === "cancelled" ? "cancelled" : "completed", + stopReason: options.completedStopReason ?? null, + }, + }); + } + } + return; + } + let settleTurnId = turnId; + if (options?.settleAllPrompts) { + liveCtx.promptsInFlight = 0; + if (liveCtx.activeTurnId !== turnId && liveCtx.session.activeTurnId !== turnId) { + const fallbackTurnId = liveCtx.activeTurnId ?? liveCtx.session.activeTurnId; + if (!fallbackTurnId) { + if (liveCtx.session.status === "running" || liveCtx.session.status === "connecting") { + const updatedAt = yield* nowIso; + const { activeTurnId: _activeTurnId, ...readySession } = liveCtx.session; + liveCtx.activeTurnId = undefined; + liveCtx.session = { + ...readySession, + status: "ready", + updatedAt, + }; + } + return; + } + settleTurnId = fallbackTurnId; + } + } else { + const remainingPrompts = Math.max(0, liveCtx.promptsInFlight - 1); + if ( + remainingPrompts > 0 || + liveCtx.activeTurnId !== settleTurnId || + liveCtx.session.activeTurnId !== settleTurnId + ) { + liveCtx.promptsInFlight = remainingPrompts; + return; + } + liveCtx.promptsInFlight = remainingPrompts; + } + const updatedAt = yield* nowIso; + const canEmitTurnCompletion = + liveCtx.session.status === "running" || liveCtx.session.status === "connecting"; + const shouldEmitFailedTurn = options?.errorMessage !== undefined && canEmitTurnCompletion; + const shouldEmitCompletedTurn = + options?.completedStopReason !== undefined && canEmitTurnCompletion; + const { activeTurnId: _activeTurnId, ...readySession } = liveCtx.session; + liveCtx.activeTurnId = undefined; + liveCtx.session = { + ...readySession, + status: "ready", + updatedAt, + }; + if (options?.emitTurnCompletion === false) { + return; + } + if (shouldEmitFailedTurn) { + yield* offerRuntimeEvent({ + type: "turn.completed", + ...(yield* makeEventStamp()), + provider: PROVIDER, + threadId, + turnId: settleTurnId, + payload: { + state: "failed", + errorMessage: options.errorMessage, + }, + }); + } else if (shouldEmitCompletedTurn) { + yield* offerRuntimeEvent({ + type: "turn.completed", + ...(yield* makeEventStamp()), + provider: PROVIDER, + threadId, + turnId: settleTurnId, + payload: { + state: options.completedStopReason === "cancelled" ? "cancelled" : "completed", + stopReason: options.completedStopReason ?? null, + }, + }); + } + }); + const logNative = (threadId: ThreadId, method: string, payload: unknown) => Effect.gen(function* () { if (!nativeEventLogger) return; @@ -271,6 +463,8 @@ export function makeGrokAdapter(grokSettings: GrokSettings, options?: GrokAdapte const emitPlanUpdate = ( ctx: GrokSessionContext, + turnId: TurnId | undefined, + stamp: { readonly eventId: EventId; readonly createdAt: string }, payload: { readonly explanation?: string | null; readonly plan: ReadonlyArray<{ @@ -282,17 +476,17 @@ export function makeGrokAdapter(grokSettings: GrokSettings, options?: GrokAdapte method: string, ) => Effect.gen(function* () { - const fingerprint = `${ctx.activeTurnId ?? "no-turn"}:${encodeJsonStringForDiagnostics(payload) ?? "[unserializable payload]"}`; + const fingerprint = `${turnId ?? "no-turn"}:${encodeJsonStringForDiagnostics(payload) ?? "[unserializable payload]"}`; if (ctx.lastPlanFingerprint === fingerprint) { return; } ctx.lastPlanFingerprint = fingerprint; yield* offerRuntimeEvent( makeAcpPlanUpdatedEvent({ - stamp: yield* makeEventStamp(), + stamp, provider: PROVIDER, threadId: ctx.threadId, - turnId: ctx.activeTurnId, + turnId, payload, source: "acp.jsonrpc", method, @@ -424,13 +618,14 @@ export function makeGrokAdapter(grokSettings: GrokSettings, options?: GrokAdapte const requestId = ApprovalRequestId.make(yield* randomUUIDv4); const runtimeRequestId = RuntimeRequestId.make(requestId); const resolution = yield* Deferred.make(); + const turnId = resolveSessionCallbackTurnId(sessions, input.threadId); pendingUserInputs.set(requestId, { resolution }); yield* offerRuntimeEvent({ type: "user-input.requested", ...(yield* makeEventStamp()), provider: PROVIDER, threadId: input.threadId, - turnId: sessions.get(input.threadId)?.activeTurnId, + turnId, requestId: runtimeRequestId, payload: { questions: extractXAiAskUserQuestions(params) }, raw: { @@ -447,7 +642,7 @@ export function makeGrokAdapter(grokSettings: GrokSettings, options?: GrokAdapte ...(yield* makeEventStamp()), provider: PROVIDER, threadId: input.threadId, - turnId: sessions.get(input.threadId)?.activeTurnId, + turnId, requestId: runtimeRequestId, payload: { answers: resolvedAnswers }, raw: { @@ -486,13 +681,14 @@ export function makeGrokAdapter(grokSettings: GrokSettings, options?: GrokAdapte const requestId = ApprovalRequestId.make(yield* randomUUIDv4); const runtimeRequestId = RuntimeRequestId.make(requestId); const decision = yield* Deferred.make(); + const turnId = resolveSessionCallbackTurnId(sessions, input.threadId); pendingApprovals.set(requestId, { decision }); yield* offerRuntimeEvent( makeAcpRequestOpenedEvent({ stamp: yield* makeEventStamp(), provider: PROVIDER, threadId: input.threadId, - turnId: sessions.get(input.threadId)?.activeTurnId, + turnId, requestId: runtimeRequestId, permissionRequest, detail: @@ -512,7 +708,7 @@ export function makeGrokAdapter(grokSettings: GrokSettings, options?: GrokAdapte stamp: yield* makeEventStamp(), provider: PROVIDER, threadId: input.threadId, - turnId: sessions.get(input.threadId)?.activeTurnId, + turnId, requestId: runtimeRequestId, permissionRequest, decision: resolved, @@ -578,6 +774,7 @@ export function makeGrokAdapter(grokSettings: GrokSettings, options?: GrokAdapte turns: [], lastPlanFingerprint: undefined, activeTurnId: undefined, + interruptedTurnIds: new Set(), promptsInFlight: 0, currentModelId: boundModelId, stopped: false, @@ -586,14 +783,39 @@ export function makeGrokAdapter(grokSettings: GrokSettings, options?: GrokAdapte const nf = yield* Stream.runDrain( Stream.mapEffect(acp.getEvents(), (event) => Effect.gen(function* () { + if (event._tag === "EventStreamBarrier") { + yield* Deferred.succeed(event.acknowledge, undefined); + return; + } + if ( + event._tag === "PlanUpdated" || + event._tag === "ToolCallUpdated" || + event._tag === "ContentDelta" + ) { + yield* logNative(ctx.threadId, "session/update", event.rawPayload); + } + + if (event._tag === "ModeChanged") { + return; + } + + const notificationTurnId = resolveNotificationTurnId(ctx); + if ( + notificationTurnId === undefined || + ctx.interruptedTurnIds.has(notificationTurnId) + ) { + return; + } + const stamp = yield* makeEventStamp(); + switch (event._tag) { case "AssistantItemStarted": yield* offerRuntimeEvent( makeAcpAssistantItemEvent({ - stamp: yield* makeEventStamp(), + stamp, provider: PROVIDER, threadId: ctx.threadId, - turnId: ctx.activeTurnId, + turnId: notificationTurnId, itemId: event.itemId, lifecycle: "item.started", }), @@ -602,40 +824,44 @@ export function makeGrokAdapter(grokSettings: GrokSettings, options?: GrokAdapte case "AssistantItemCompleted": yield* offerRuntimeEvent( makeAcpAssistantItemEvent({ - stamp: yield* makeEventStamp(), + stamp, provider: PROVIDER, threadId: ctx.threadId, - turnId: ctx.activeTurnId, + turnId: notificationTurnId, itemId: event.itemId, lifecycle: "item.completed", }), ); return; case "PlanUpdated": - yield* logNative(ctx.threadId, "session/update", event.rawPayload); - yield* emitPlanUpdate(ctx, event.payload, event.rawPayload, "session/update"); + yield* emitPlanUpdate( + ctx, + notificationTurnId, + stamp, + event.payload, + event.rawPayload, + "session/update", + ); return; case "ToolCallUpdated": - yield* logNative(ctx.threadId, "session/update", event.rawPayload); yield* offerRuntimeEvent( makeAcpToolCallEvent({ - stamp: yield* makeEventStamp(), + stamp, provider: PROVIDER, threadId: ctx.threadId, - turnId: ctx.activeTurnId, + turnId: notificationTurnId, toolCall: event.toolCall, rawPayload: event.rawPayload, }), ); return; case "ContentDelta": - yield* logNative(ctx.threadId, "session/update", event.rawPayload); yield* offerRuntimeEvent( makeAcpContentDeltaEvent({ - stamp: yield* makeEventStamp(), + stamp, provider: PROVIDER, threadId: ctx.threadId, - turnId: ctx.activeTurnId, + turnId: notificationTurnId, ...(event.itemId ? { itemId: event.itemId } : {}), text: event.text, rawPayload: event.rawPayload, @@ -697,6 +923,15 @@ export function makeGrokAdapter(grokSettings: GrokSettings, options?: GrokAdapte // resolving from here on does not settle the turn; decremented on // preparation failure here, and after the prompt below otherwise. ctx.promptsInFlight += 1; + // Bind the turn id before cooperative yields so interruptTurn can + // settle this prompt even if stop arrives during preparation. + ctx.activeTurnId = turnId; + ctx.session = { + ...ctx.session, + status: steeringTurnId === undefined ? "connecting" : "running", + activeTurnId: turnId, + updatedAt: yield* nowIso, + }; return yield* Effect.gen(function* () { const turnModelSelection = @@ -765,12 +1000,27 @@ export function makeGrokAdapter(grokSettings: GrokSettings, options?: GrokAdapte const displayModel = currentModelId ? resolveGrokAcpBaseModelId(currentModelId) : undefined; - ctx.activeTurnId = turnId; + for (let yieldAttempt = 0; yieldAttempt < 8; yieldAttempt += 1) { + yield* Effect.yieldNow; + } + if (ctx.interruptedTurnIds.has(turnId)) { + yield* settlePromptInFlight(input.threadId, turnId, ctx.acpSessionId, { + completedStopReason: "cancelled", + emitTurnCompletion: false, + settleAllPrompts: true, + }); + return yield* new ProviderAdapterRequestError({ + provider: PROVIDER, + method: "session/prompt", + detail: "Grok prompt was interrupted during preparation.", + }); + } if (steeringTurnId === undefined) { ctx.lastPlanFingerprint = undefined; } ctx.session = { ...ctx.session, + status: "running", activeTurnId: turnId, updatedAt: yield* nowIso, ...(displayModel ? { model: displayModel } : {}), @@ -796,13 +1046,27 @@ export function makeGrokAdapter(grokSettings: GrokSettings, options?: GrokAdapte }; }).pipe( Effect.tapCause(() => - Effect.sync(() => { - ctx.promptsInFlight = Math.max(0, ctx.promptsInFlight - 1); + Effect.gen(function* () { + const liveCtx = sessions.get(input.threadId); + if (!liveCtx) { + return; + } + yield* settlePromptInFlight(input.threadId, turnId, liveCtx.acpSessionId, { + errorMessage: "Grok prompt preparation failed.", + emitTurnCompletion: false, + }); }), ), ); }), ); + const promptSettled = yield* Ref.make(false); + const promptRpcSucceeded = yield* Ref.make(false); + const promptResultRef = yield* Ref.make( + undefined, + ); + + const promptFailureMessageRef = yield* Ref.make(undefined); return yield* Effect.gen(function* () { const result = yield* prepared.acp @@ -810,6 +1074,18 @@ export function makeGrokAdapter(grokSettings: GrokSettings, options?: GrokAdapte prompt: prepared.promptParts, }) .pipe( + Effect.tap((promptResult) => + Effect.all([ + Ref.set(promptRpcSucceeded, true), + Ref.set(promptResultRef, promptResult), + ]), + ), + Effect.tapError((error) => + Ref.set( + promptFailureMessageRef, + mapAcpToAdapterError(PROVIDER, input.threadId, "session/prompt", error).message, + ).pipe(Effect.andThen(prepared.acp.drainEvents)), + ), Effect.mapError((error) => mapAcpToAdapterError(PROVIDER, input.threadId, "session/prompt", error), ), @@ -820,38 +1096,88 @@ export function makeGrokAdapter(grokSettings: GrokSettings, options?: GrokAdapte Effect.gen(function* () { const ctx = yield* requireSession(input.threadId); if (ctx.acpSessionId !== prepared.acpSessionId) { + yield* settlePromptInFlight( + input.threadId, + prepared.turnId, + prepared.acpSessionId, + { + errorMessage: "Grok session changed before the turn completed.", + settleAllPrompts: true, + }, + ); + yield* Ref.set(promptSettled, true); return yield* new ProviderAdapterRequestError({ provider: PROVIDER, method: "session/prompt", detail: "Grok session changed before the turn completed.", }); } + // Keep prompt settlement atomic with respect to Stop and steering. + // interruptTurn marks its target before waiting for this lock, so + // cancellation can still win while queued ACP events are drained. + for (let yieldAttempt = 0; yieldAttempt < 8; yieldAttempt += 1) { + yield* Effect.yieldNow; + } + yield* prepared.acp.drainEvents; + if (ctx.interruptedTurnIds.has(prepared.turnId)) { + yield* Ref.set(promptSettled, true); + return { + threadId: input.threadId, + turnId: prepared.turnId, + resumeCursor: ctx.session.resumeCursor, + }; + } - const existingTurnRecord = ctx.turns.find((turn) => turn.id === prepared.turnId); - ctx.turns = existingTurnRecord - ? ctx.turns.map((turn) => - turn.id === prepared.turnId - ? { - ...turn, - items: [...turn.items, { prompt: prepared.promptParts, result }], - } - : turn, - ) - : [ - ...ctx.turns, - { id: prepared.turnId, items: [{ prompt: prepared.promptParts, result }] }, - ]; + if ( + ctx.promptsInFlight <= 0 || + ctx.activeTurnId !== prepared.turnId || + ctx.session.activeTurnId !== prepared.turnId + ) { + yield* Ref.set(promptSettled, true); + return { + threadId: input.threadId, + turnId: prepared.turnId, + resumeCursor: ctx.session.resumeCursor, + }; + } + + appendPromptResultToTurn(ctx, prepared.turnId, prepared.promptParts, result); ctx.session = { ...ctx.session, + status: "running", activeTurnId: prepared.turnId, updatedAt: yield* nowIso, ...(prepared.displayModel ? { model: prepared.displayModel } : {}), }; + const remainingPrompts = Math.max(0, ctx.promptsInFlight - 1); + ctx.promptsInFlight = remainingPrompts; - // Only the last remaining prompt settles the turn — a steer- - // superseded prompt resolving (usually cancelled) while another - // is in flight or pending must leave the merged turn running. - if (ctx.promptsInFlight === 1) { + // Only the last remaining prompt settles the turn. A steer- + // superseded prompt resolving while another is in flight or + // pending must leave the merged turn running. + if ( + remainingPrompts === 0 && + ctx.activeTurnId === prepared.turnId && + ctx.session.activeTurnId === prepared.turnId + ) { + if (ctx.interruptedTurnIds.has(prepared.turnId)) { + yield* Ref.set(promptSettled, true); + return { + threadId: input.threadId, + turnId: prepared.turnId, + resumeCursor: ctx.session.resumeCursor, + }; + } + const completedAt = yield* nowIso; + const { activeTurnId: _completedTurnId, ...readySession } = ctx.session; + ctx.activeTurnId = undefined; + ctx.session = { + ...readySession, + status: "ready", + updatedAt: completedAt, + ...(prepared.displayModel ? { model: prepared.displayModel } : {}), + }; + const completedStopReason = completedStopReasonFromPromptResponse(result); yield* offerRuntimeEvent({ type: "turn.completed", ...(yield* makeEventStamp()), @@ -860,9 +1186,13 @@ export function makeGrokAdapter(grokSettings: GrokSettings, options?: GrokAdapte turnId: prepared.turnId, payload: { state: result.stopReason === "cancelled" ? "cancelled" : "completed", - stopReason: result.stopReason ?? null, + stopReason: completedStopReason, }, }); + ctx.interruptedTurnIds.delete(prepared.turnId); + yield* Ref.set(promptSettled, true); + } else if (remainingPrompts > 0) { + yield* Ref.set(promptSettled, true); } return { @@ -874,27 +1204,153 @@ export function makeGrokAdapter(grokSettings: GrokSettings, options?: GrokAdapte ); }).pipe( Effect.ensuring( - Effect.sync(() => { - const liveCtx = sessions.get(input.threadId); - if (liveCtx) { - liveCtx.promptsInFlight = Math.max(0, liveCtx.promptsInFlight - 1); + Effect.gen(function* () { + if (yield* Ref.get(promptSettled)) { + return; } - }), + + if (yield* Ref.get(promptRpcSucceeded)) { + const promptResult = yield* Ref.get(promptResultRef); + if (promptResult === undefined) { + return; + } + yield* withThreadLock( + input.threadId, + Effect.gen(function* () { + const ctx = yield* requireSession(input.threadId); + if (ctx.acpSessionId !== prepared.acpSessionId) { + yield* settlePromptInFlight( + input.threadId, + prepared.turnId, + prepared.acpSessionId, + { + errorMessage: "Grok session changed before the turn completed.", + settleAllPrompts: true, + }, + ); + return; + } + if (ctx.interruptedTurnIds.has(prepared.turnId)) { + return; + } + if ( + ctx.promptsInFlight <= 0 || + ctx.activeTurnId !== prepared.turnId || + ctx.session.activeTurnId !== prepared.turnId + ) { + return; + } + appendPromptResultToTurn( + ctx, + prepared.turnId, + prepared.promptParts, + promptResult, + ); + yield* settlePromptInFlight( + input.threadId, + prepared.turnId, + prepared.acpSessionId, + { + completedStopReason: completedStopReasonFromPromptResponse(promptResult), + }, + ); + }), + ); + return; + } + + const errorMessage = yield* Ref.get(promptFailureMessageRef); + yield* withThreadLock( + input.threadId, + settlePromptInFlight(input.threadId, prepared.turnId, prepared.acpSessionId, { + errorMessage: errorMessage ?? "Grok prompt request failed.", + }), + ); + }).pipe(Effect.catch(() => Effect.void)), ), ); }); - const interruptTurn: GrokAdapterShape["interruptTurn"] = (threadId) => + const interruptTurn: GrokAdapterShape["interruptTurn"] = (threadId, turnId) => Effect.gen(function* () { - const ctx = yield* requireSession(threadId); - yield* settlePendingApprovalsAsCancelled(ctx.pendingApprovals); - yield* settlePendingUserInputsAsCancelled(ctx.pendingUserInputs); - yield* Effect.ignore( - ctx.acp.cancel.pipe( - Effect.mapError((error) => - mapAcpToAdapterError(PROVIDER, threadId, "session/cancel", error), - ), - ), + const observed = yield* Effect.sync(() => { + const ctx = sessions.get(threadId); + if (!ctx || ctx.stopped) { + return { + _tag: "Proceed" as const, + acpSessionId: undefined, + interruptedTurnId: turnId, + }; + } + const activeTurnId = ctx.activeTurnId ?? ctx.session.activeTurnId; + if (turnId !== undefined && activeTurnId !== undefined && activeTurnId !== turnId) { + return { _tag: "Ignore" as const }; + } + const interruptedTurnId = turnId ?? activeTurnId; + if (interruptedTurnId !== undefined) { + ctx.interruptedTurnIds.add(interruptedTurnId); + } + return { + _tag: "Proceed" as const, + acpSessionId: ctx.acpSessionId, + interruptedTurnId, + }; + }); + if (observed._tag === "Ignore") { + return; + } + + yield* withThreadLock( + threadId, + Effect.gen(function* () { + const ctx = yield* requireSession(threadId); + if (observed.acpSessionId !== undefined && ctx.acpSessionId !== observed.acpSessionId) { + return; + } + const activeTurnId = ctx.activeTurnId ?? ctx.session.activeTurnId; + if (turnId !== undefined && activeTurnId !== undefined && activeTurnId !== turnId) { + return; + } + if ( + observed.interruptedTurnId !== undefined && + activeTurnId !== undefined && + activeTurnId !== observed.interruptedTurnId + ) { + return; + } + const interruptedTurnId = + observed.interruptedTurnId ?? turnId ?? activeTurnId ?? ctx.session.activeTurnId; + yield* settlePendingApprovalsAsCancelled(ctx.pendingApprovals); + yield* settlePendingUserInputsAsCancelled(ctx.pendingUserInputs); + yield* Effect.ignore( + ctx.acp.cancel.pipe( + Effect.mapError((error) => + mapAcpToAdapterError(PROVIDER, threadId, "session/cancel", error), + ), + ), + ); + if (interruptedTurnId) { + ctx.interruptedTurnIds.add(interruptedTurnId); + yield* settlePromptInFlight(threadId, interruptedTurnId, ctx.acpSessionId, { + completedStopReason: "cancelled", + settleAllPrompts: true, + }); + } else if ( + ctx.promptsInFlight > 0 || + ctx.session.status === "running" || + ctx.session.status === "connecting" + ) { + const updatedAt = yield* nowIso; + ctx.promptsInFlight = 0; + ctx.activeTurnId = undefined; + const { activeTurnId: _activeTurnId, ...readySession } = ctx.session; + ctx.session = { + ...readySession, + status: "ready", + updatedAt, + }; + } + }), ); }); diff --git a/apps/server/src/provider/acp/AcpJsonRpcConnection.test.ts b/apps/server/src/provider/acp/AcpJsonRpcConnection.test.ts index 5533a04bc83..4e9700dab7d 100644 --- a/apps/server/src/provider/acp/AcpJsonRpcConnection.test.ts +++ b/apps/server/src/provider/acp/AcpJsonRpcConnection.test.ts @@ -7,6 +7,9 @@ import * as NodeFS from "node:fs"; import * as NodeServices from "@effect/platform-node/NodeServices"; import { it } from "@effect/vitest"; import * as Effect from "effect/Effect"; +import * as Fiber from "effect/Fiber"; +import * as Option from "effect/Option"; +import * as TestClock from "effect/testing/TestClock"; import * as Stream from "effect/Stream"; import { describe, expect } from "vite-plus/test"; @@ -113,6 +116,122 @@ describe("AcpSessionRuntime", () => { ), ); + it.effect("drops session updates emitted for a child ACP session", () => + Effect.gen(function* () { + const runtime = yield* AcpSessionRuntime.AcpSessionRuntime; + yield* runtime.start(); + + const promptResult = yield* runtime.prompt({ + prompt: [{ type: "text", text: "hi" }], + }); + expect(promptResult).toMatchObject({ stopReason: "end_turn" }); + + const notes = Array.from(yield* Stream.runCollect(Stream.take(runtime.getEvents(), 4))); + expect(notes.map((note) => note._tag)).toEqual([ + "AssistantItemStarted", + "ContentDelta", + "ContentDelta", + "AssistantItemCompleted", + ]); + expect( + notes + .filter((note) => note._tag === "ContentDelta") + .map((note) => note.text) + .join(""), + ).toBe("root before child root after child"); + expect(notes.some((note) => note._tag === "ToolCallUpdated")).toBe(false); + }).pipe( + Effect.provide( + AcpSessionRuntime.layer({ + spawn: { + command: mockAgentCommand, + args: mockAgentArgs, + env: { + T3_ACP_EMIT_FOREIGN_SESSION_UPDATES: "1", + }, + }, + cwd: process.cwd(), + clientInfo: { name: "t3-test", version: "0.0.0" }, + authMethodId: "test", + }), + ), + Effect.scoped, + Effect.provide(NodeServices.layer), + ), + ); + + it.effect("supports successive standard ACP prompts", () => + Effect.gen(function* () { + const runtime = yield* AcpSessionRuntime.AcpSessionRuntime; + yield* runtime.start(); + + const firstPromptResult = yield* runtime.prompt({ + prompt: [{ type: "text", text: "first" }], + }); + const secondPromptResult = yield* runtime.prompt({ + prompt: [{ type: "text", text: "second" }], + }); + + expect(firstPromptResult).toMatchObject({ stopReason: "end_turn" }); + expect(secondPromptResult).toMatchObject({ stopReason: "end_turn" }); + }).pipe( + Effect.provide( + AcpSessionRuntime.layer({ + spawn: { + command: mockAgentCommand, + args: mockAgentArgs, + }, + cwd: process.cwd(), + clientInfo: { name: "t3-test", version: "0.0.0" }, + authMethodId: "test", + }), + ), + Effect.scoped, + Effect.provide(NodeServices.layer), + ), + ); + + it.effect("releases a fully silent prompt when session/cancel is requested", () => + Effect.gen(function* () { + const runtime = yield* AcpSessionRuntime.AcpSessionRuntime; + yield* runtime.start(); + + const promptFiber = yield* runtime + .prompt({ + prompt: [{ type: "text", text: "hang forever" }], + }) + .pipe(Effect.forkChild({ startImmediately: true })); + + yield* TestClock.adjust("500 millis"); + yield* runtime.cancel; + + const firstPromptResult = yield* Fiber.join(promptFiber); + expect(firstPromptResult).toMatchObject({ stopReason: "cancelled" }); + + const secondPromptResult = yield* runtime.prompt({ + prompt: [{ type: "text", text: "second" }], + }); + expect(secondPromptResult).toMatchObject({ stopReason: "end_turn" }); + }).pipe( + Effect.provide( + AcpSessionRuntime.layer({ + spawn: { + command: mockAgentCommand, + args: mockAgentArgs, + env: { + T3_ACP_HANG_FIRST_PROMPT_FOREVER: "1", + }, + }, + cwd: process.cwd(), + clientInfo: { name: "t3-test", version: "0.0.0" }, + authMethodId: "test", + }), + ), + Effect.scoped, + Effect.provide(NodeServices.layer), + ), + ); + it.effect("segments assistant text around ACP tool calls", () => Effect.gen(function* () { const runtime = yield* AcpSessionRuntime.AcpSessionRuntime; @@ -346,6 +465,109 @@ describe("AcpSessionRuntime", () => { ); }); + it.effect("fails session startup when session/load returns an error", () => + Effect.gen(function* () { + const runtime = yield* AcpSessionRuntime.AcpSessionRuntime; + const error = yield* runtime.start().pipe(Effect.flip); + + expect(error._tag).toBe("AcpRequestError"); + }).pipe( + Effect.provide( + AcpSessionRuntime.layer({ + authMethodId: "test", + spawn: { + command: mockAgentCommand, + args: mockAgentArgs, + env: { + T3_ACP_FAIL_LOAD_SESSION: "1", + }, + }, + cwd: process.cwd(), + resumeSessionId: "stale-session-id", + clientInfo: { name: "t3-test", version: "0.0.0" }, + }), + ), + Effect.scoped, + Effect.provide(NodeServices.layer), + ), + ); + + it.effect("ignores session/update replay notifications during session/load", () => + Effect.gen(function* () { + const runtime = yield* AcpSessionRuntime.AcpSessionRuntime; + yield* runtime.start(); + + yield* runtime.prompt({ + prompt: [{ type: "text", text: "hi" }], + }); + const notes = Array.from(yield* Stream.runCollect(Stream.take(runtime.getEvents(), 4))); + expect(notes.map((note) => note._tag)).toEqual([ + "PlanUpdated", + "AssistantItemStarted", + "ContentDelta", + "AssistantItemCompleted", + ]); + expect(notes.some((note) => note._tag === "ToolCallUpdated")).toBe(false); + }).pipe( + Effect.provide( + AcpSessionRuntime.layer({ + authMethodId: "test", + spawn: { + command: mockAgentCommand, + args: mockAgentArgs, + env: { + T3_ACP_EMIT_LOAD_REPLAY: "1", + }, + }, + cwd: process.cwd(), + resumeSessionId: "mock-session-1", + clientInfo: { name: "t3-test", version: "0.0.0" }, + }), + ), + Effect.scoped, + Effect.provide(NodeServices.layer), + ), + ); + + it.effect("completes session/load after replay becomes idle while its RPC stays pending", () => + Effect.gen(function* () { + const runtime = yield* AcpSessionRuntime.AcpSessionRuntime; + const started = yield* runtime.start().pipe(Effect.timeout("2 seconds")); + + expect(started.sessionId).toBe("mock-session-1"); + expect(started.sessionSetupResult._meta).toMatchObject({ + t3SessionLoadReady: "replay_idle", + }); + + const unexpectedReplayEvent = yield* Stream.runHead(runtime.getEvents()).pipe( + Effect.timeoutOption("100 millis"), + ); + expect(Option.isNone(unexpectedReplayEvent)).toBe(true); + }).pipe( + Effect.provide( + AcpSessionRuntime.layer({ + authMethodId: "test", + spawn: { + command: mockAgentCommand, + args: mockAgentArgs, + env: { + T3_ACP_HANG_LOAD_SESSION_AFTER_REPLAY: "1", + T3_ACP_LOAD_SESSION_DELAY_MS: "10000", + }, + }, + cwd: process.cwd(), + resumeSessionId: "mock-session-1", + sessionLoadReplayIdleGap: "50 millis", + sessionLoadTimeout: "1 second", + clientInfo: { name: "t3-test", version: "0.0.0" }, + }), + ), + Effect.scoped, + Effect.provide(NodeServices.layer), + TestClock.withLive, + ), + ); + it.effect("rejects invalid config option values before sending session/set_config_option", () => { const tempDir = NodeFS.mkdtempSync(NodePath.join(NodeOS.tmpdir(), "acp-runtime-")); const requestLogPath = NodePath.join(tempDir, "requests.ndjson"); diff --git a/apps/server/src/provider/acp/AcpRuntimeModel.test.ts b/apps/server/src/provider/acp/AcpRuntimeModel.test.ts index 54b9e5390e0..7682c5f5f9c 100644 --- a/apps/server/src/provider/acp/AcpRuntimeModel.test.ts +++ b/apps/server/src/provider/acp/AcpRuntimeModel.test.ts @@ -8,6 +8,8 @@ import { parsePermissionRequest, parseSessionModeState, parseSessionUpdateEvent, + sessionUpdateIsReplay, + syntheticLoadSessionResponseFromInitialize, } from "./AcpRuntimeModel.ts"; describe("AcpRuntimeModel", () => { @@ -59,6 +61,95 @@ describe("AcpRuntimeModel", () => { expect(modelConfigId).toBe("model"); }); + it("detects Grok session replay updates from _meta.isReplay", () => { + expect( + sessionUpdateIsReplay({ + _meta: { isReplay: true }, + sessionId: "session-1", + update: { + sessionUpdate: "agent_message_chunk", + content: { type: "text", text: "replayed" }, + }, + } satisfies EffectAcpSchema.SessionNotification), + ).toBe(true); + expect( + sessionUpdateIsReplay({ + sessionId: "session-1", + update: { + sessionUpdate: "agent_message_chunk", + content: { type: "text", text: "live" }, + }, + } satisfies EffectAcpSchema.SessionNotification), + ).toBe(false); + }); + + it("builds a synthetic load response from initialize model state", () => { + const response = syntheticLoadSessionResponseFromInitialize({ + protocolVersion: 1, + _meta: { + modelState: { + currentModelId: "grok-build", + availableModels: [{ modelId: "grok-build", name: "Grok Build" }], + }, + }, + } satisfies EffectAcpSchema.InitializeResponse); + + expect(response.models?.currentModelId).toBe("grok-build"); + expect(response._meta).toMatchObject({ t3SessionLoadReady: "replay_idle" }); + }); + + it("accepts initialize model descriptions with null", () => { + const response = syntheticLoadSessionResponseFromInitialize({ + protocolVersion: 1, + _meta: { + modelState: { + currentModelId: "grok-build", + availableModels: [{ modelId: "grok-build", name: "Grok Build", description: null }], + }, + }, + } satisfies EffectAcpSchema.InitializeResponse); + + expect(response.models?.availableModels[0]?.description).toBeNull(); + }); + + it("ignores malformed initialize model state in synthetic load responses", () => { + const response = syntheticLoadSessionResponseFromInitialize({ + protocolVersion: 1, + _meta: { + modelState: { + currentModelId: "grok-build", + availableModels: [null], + }, + modeState: { + currentModeId: "code", + availableModes: [{ id: "code", name: 12 }], + }, + }, + } as EffectAcpSchema.InitializeResponse); + + expect(response.models).toBeUndefined(); + expect(response.modes).toBeUndefined(); + expect(response._meta).toMatchObject({ t3SessionLoadReady: "replay_idle" }); + }); + + it("builds a synthetic load response with initialize mode state", () => { + const response = syntheticLoadSessionResponseFromInitialize({ + protocolVersion: 1, + _meta: { + modeState: { + currentModeId: "code", + availableModes: [ + { id: "ask", name: "Ask" }, + { id: "code", name: "Code" }, + ], + }, + }, + } satisfies EffectAcpSchema.InitializeResponse); + + expect(response.modes?.currentModeId).toBe("code"); + expect(response.modes?.availableModes).toHaveLength(2); + }); + it("projects typed ACP tool call updates into runtime events", () => { const created = parseSessionUpdateEvent({ sessionId: "session-1", diff --git a/apps/server/src/provider/acp/AcpRuntimeModel.ts b/apps/server/src/provider/acp/AcpRuntimeModel.ts index 3587d703dbd..e6bfc127e6e 100644 --- a/apps/server/src/provider/acp/AcpRuntimeModel.ts +++ b/apps/server/src/provider/acp/AcpRuntimeModel.ts @@ -1,3 +1,8 @@ +import * as Clock from "effect/Clock"; +import * as Duration from "effect/Duration"; +import * as Effect from "effect/Effect"; +import * as Option from "effect/Option"; +import * as Ref from "effect/Ref"; import type * as EffectAcpSchema from "effect-acp/schema"; import { deriveToolActivityPresentation } from "@t3tools/shared/toolActivity"; import type { ToolLifecycleItemType } from "@t3tools/contracts"; @@ -6,6 +11,40 @@ function isRecord(value: unknown): value is Record { return typeof value === "object" && value !== null && !Array.isArray(value); } +function isSessionModelState(value: unknown): value is EffectAcpSchema.SessionModelState { + if (!isRecord(value) || typeof value.currentModelId !== "string") { + return false; + } + if (!Array.isArray(value.availableModels)) { + return false; + } + return value.availableModels.every( + (model) => + isRecord(model) && + typeof model.modelId === "string" && + typeof model.name === "string" && + (model.description === undefined || + model.description === null || + typeof model.description === "string"), + ); +} + +function isSessionModeState(value: unknown): value is EffectAcpSchema.SessionModeState { + if (!isRecord(value) || typeof value.currentModeId !== "string") { + return false; + } + if (!Array.isArray(value.availableModes)) { + return false; + } + return value.availableModes.every( + (mode) => + isRecord(mode) && + typeof mode.id === "string" && + typeof mode.name === "string" && + (mode.description === undefined || typeof mode.description === "string"), + ); +} + export interface AcpSessionMode { readonly id: string; readonly name: string; @@ -414,6 +453,58 @@ export function parsePermissionRequest( }; } +export function sessionUpdateIsReplay(params: EffectAcpSchema.SessionNotification): boolean { + const meta = params._meta; + return isRecord(meta) && meta.isReplay === true; +} + +export interface SessionLoadGate { + readonly active: boolean; + readonly lastActivityAtMillis: number | undefined; + readonly idleGap: Duration.Duration; + readonly initializeResult: EffectAcpSchema.InitializeResponse; +} + +export const waitForSessionLoadReplayIdle = (input: { + readonly gateRef: Ref.Ref>; +}): Effect.Effect => + Effect.gen(function* () { + const pollInterval = Duration.millis(25); + while (true) { + const gate = yield* Ref.get(input.gateRef); + if ( + Option.isSome(gate) && + gate.value.active && + gate.value.lastActivityAtMillis !== undefined + ) { + const idleGapMillis = Duration.toMillis(gate.value.idleGap); + const nowMillis = yield* Clock.currentTimeMillis; + if (nowMillis - gate.value.lastActivityAtMillis >= idleGapMillis) { + return syntheticLoadSessionResponseFromInitialize(gate.value.initializeResult); + } + } + yield* Effect.sleep(pollInterval); + } + }); + +export function syntheticLoadSessionResponseFromInitialize( + initializeResult: EffectAcpSchema.InitializeResponse, +): EffectAcpSchema.LoadSessionResponse { + const meta = initializeResult._meta; + const modelState = isRecord(meta) ? meta.modelState : undefined; + const modeState = isRecord(meta) ? meta.modeState : undefined; + const models = isSessionModelState(modelState) ? modelState : undefined; + const modes = isSessionModeState(modeState) ? modeState : undefined; + + return { + ...(models ? { models } : {}), + ...(modes ? { modes } : {}), + _meta: { + t3SessionLoadReady: "replay_idle", + }, + }; +} + export function parseSessionUpdateEvent(params: EffectAcpSchema.SessionNotification): { readonly modeId?: string; readonly events: ReadonlyArray; diff --git a/apps/server/src/provider/acp/AcpSessionRuntime.ts b/apps/server/src/provider/acp/AcpSessionRuntime.ts index 4fc2c443e11..bc2df3aa8d4 100644 --- a/apps/server/src/provider/acp/AcpSessionRuntime.ts +++ b/apps/server/src/provider/acp/AcpSessionRuntime.ts @@ -1,12 +1,16 @@ import * as Cause from "effect/Cause"; +import * as Clock from "effect/Clock"; import * as Context from "effect/Context"; import * as Deferred from "effect/Deferred"; +import * as Duration from "effect/Duration"; import * as Effect from "effect/Effect"; -import * as Exit from "effect/Exit"; +import * as Fiber from "effect/Fiber"; import * as Layer from "effect/Layer"; +import * as Option from "effect/Option"; import * as Queue from "effect/Queue"; import * as Ref from "effect/Ref"; import * as Scope from "effect/Scope"; +import * as Semaphore from "effect/Semaphore"; import * as Stream from "effect/Stream"; import * as ChildProcess from "effect/unstable/process/ChildProcess"; import * as ChildProcessSpawner from "effect/unstable/process/ChildProcessSpawner"; @@ -23,6 +27,9 @@ import { mergeToolCallState, parseSessionModeState, parseSessionUpdateEvent, + sessionUpdateIsReplay, + waitForSessionLoadReplayIdle, + type SessionLoadGate, type AcpParsedSessionEvent, type AcpSessionModeState, type AcpToolCallState, @@ -32,6 +39,16 @@ function formatConfigOptionValue(value: string | boolean): string { return JSON.stringify(value); } +export interface AcpSessionEventStreamBarrier { + readonly _tag: "EventStreamBarrier"; + readonly acknowledge: Deferred.Deferred; +} + +export type AcpSessionRuntimeEvent = AcpParsedSessionEvent | AcpSessionEventStreamBarrier; + +const defaultSessionLoadTimeout = Duration.seconds(90); +const defaultSessionLoadReplayIdleGap = Duration.seconds(2); + export interface AcpSpawnInput { readonly command: string; readonly args: ReadonlyArray; @@ -43,6 +60,8 @@ export interface AcpSessionRuntimeOptions { readonly spawn: AcpSpawnInput; readonly cwd: string; readonly resumeSessionId?: string; + readonly sessionLoadTimeout?: Duration.Input; + readonly sessionLoadReplayIdleGap?: Duration.Input; readonly clientCapabilities?: EffectAcpSchema.InitializeRequest["clientCapabilities"]; readonly clientInfo: { readonly name: string; @@ -160,7 +179,9 @@ export class AcpSessionRuntime extends Context.Service< */ readonly start: () => Effect.Effect; /** Stream of parsed ACP session events emitted after startup. */ - readonly getEvents: () => Stream.Stream; + readonly getEvents: () => Stream.Stream; + /** Waits until the current event consumer has processed every queued event. */ + readonly drainEvents: Effect.Effect; /** Latest mode state observed from session setup and `session/update` notifications. */ readonly getModeState: Effect.Effect; /** Latest configuration options observed from session setup and configuration writes. */ @@ -254,12 +275,17 @@ export const make = ( Effect.gen(function* () { const spawner = yield* ChildProcessSpawner.ChildProcessSpawner; const runtimeScope = yield* Scope.Scope; - const eventQueue = yield* Queue.unbounded(); + const eventQueue = yield* Queue.unbounded(); const modeStateRef = yield* Ref.make(undefined); const toolCallsRef = yield* Ref.make(new Map()); const assistantSegmentRef = yield* Ref.make({ nextSegmentIndex: 0 }); const configOptionsRef = yield* Ref.make(sessionConfigOptionsFromSetup(undefined)); const startStateRef = yield* Ref.make({ _tag: "NotStarted" }); + const promptSerializationSemaphore = yield* Semaphore.make(1); + const activePromptFiberRef = yield* Ref.make< + Option.Option> + >(Option.none()); + const sessionLoadGateRef = yield* Ref.make>(Option.none()); const logRequest = (event: AcpSessionRequestLogEvent) => options.requestLogger ? options.requestLogger(event) : Effect.void; @@ -331,15 +357,40 @@ export const make = ( const acp = yield* Effect.service(EffectAcpClient.AcpClient).pipe(Effect.provide(acpContext)); yield* acp.handleSessionUpdate((notification) => - handleSessionUpdate({ - queue: eventQueue, - modeStateRef, - toolCallsRef, - assistantSegmentRef, - params: notification, + Effect.gen(function* () { + const gate = yield* Ref.get(sessionLoadGateRef); + if (Option.isSome(gate) && gate.value.active) { + const lastActivityAtMillis = yield* Clock.currentTimeMillis; + yield* Ref.set( + sessionLoadGateRef, + Option.some({ + ...gate.value, + lastActivityAtMillis, + }), + ); + return; + } + if (sessionUpdateIsReplay(notification)) { + return; + } + const startState = yield* Ref.get(startStateRef); + // One runtime projects one root ACP session. Child-session updates need + // explicit lineage routing and must never be flattened into this stream. + if ( + startState._tag !== "Started" || + notification.sessionId !== startState.result.sessionId + ) { + return; + } + yield* handleSessionUpdate({ + queue: eventQueue, + modeStateRef, + toolCallsRef, + assistantSegmentRef, + params: notification, + }); }), ); - const initializeClientCapabilities = { fs: { readTextFile: false, @@ -499,27 +550,74 @@ export const make = ( cwd: options.cwd, mcpServers: options.mcpServers ?? [], } satisfies EffectAcpSchema.LoadSessionRequest; - const resumed = yield* runLoggedRequest( - "session/load", - loadPayload, - acp.agent.loadSession(loadPayload), - ).pipe(Effect.exit); - if (Exit.isSuccess(resumed)) { - sessionId = options.resumeSessionId; - sessionSetupResult = resumed.value; - } else { - const createPayload = { - cwd: options.cwd, - mcpServers: options.mcpServers ?? [], - } satisfies EffectAcpSchema.NewSessionRequest; - const created = yield* runLoggedRequest( - "session/new", - createPayload, - acp.agent.createSession(createPayload), + const sessionLoadTimeout = Duration.fromInputUnsafe( + options.sessionLoadTimeout ?? defaultSessionLoadTimeout, + ); + const sessionLoadReplayIdleGap = Duration.fromInputUnsafe( + options.sessionLoadReplayIdleGap ?? defaultSessionLoadReplayIdleGap, + ); + + yield* Ref.set( + sessionLoadGateRef, + Option.some({ + active: true, + lastActivityAtMillis: undefined, + idleGap: sessionLoadReplayIdleGap, + initializeResult, + }), + ); + + sessionId = options.resumeSessionId; + sessionSetupResult = yield* Effect.gen(function* () { + yield* logRequest({ + method: "session/load", + payload: loadPayload, + status: "started", + }); + + const idleFiber = yield* waitForSessionLoadReplayIdle({ + gateRef: sessionLoadGateRef, + }).pipe(Effect.forkIn(runtimeScope)); + const loaded = yield* Effect.raceFirst( + acp.agent.loadSession(loadPayload), + Fiber.join(idleFiber), + ).pipe( + Effect.ensuring(Fiber.interrupt(idleFiber).pipe(Effect.ignore)), + Effect.timeoutOption(sessionLoadTimeout), + Effect.flatMap((result) => + Option.match(result, { + onNone: () => + Effect.fail( + new EffectAcpErrors.AcpTransportError({ + operation: "call-rpc", + method: "session/load", + detail: "session/load timed out waiting for RPC response or replay idle gap", + cause: undefined, + }), + ), + onSome: Effect.succeed, + }), + ), + Effect.tap((result) => + logRequest({ + method: "session/load", + payload: loadPayload, + status: "succeeded", + result, + }), + ), + Effect.onError((cause) => + logRequest({ + method: "session/load", + payload: loadPayload, + status: "failed", + cause, + }), + ), ); - sessionId = created.sessionId; - sessionSetupResult = created; - } + + return loaded; + }).pipe(Effect.ensuring(Ref.set(sessionLoadGateRef, Option.none()))); } else { const createPayload = { cwd: options.cwd, @@ -596,25 +694,48 @@ export const make = ( handleExtNotification: acp.handleExtNotification, start: () => start, getEvents: () => Stream.fromQueue(eventQueue), + drainEvents: Effect.gen(function* () { + const acknowledge = yield* Deferred.make(); + yield* Queue.offer(eventQueue, { + _tag: "EventStreamBarrier", + acknowledge, + }); + yield* Deferred.await(acknowledge); + }), getModeState: Ref.get(modeStateRef), getConfigOptions: Ref.get(configOptionsRef), prompt: (payload) => - getStartedState.pipe( - Effect.flatMap((started) => { + promptSerializationSemaphore.withPermit( + Effect.gen(function* () { + const started = yield* getStartedState; + yield* closeActiveAssistantSegment({ + queue: eventQueue, + assistantSegmentRef, + }); const requestPayload = { sessionId: started.sessionId, ...payload, } satisfies EffectAcpSchema.PromptRequest; - return closeActiveAssistantSegment({ - queue: eventQueue, - assistantSegmentRef, - }).pipe( - Effect.andThen( - runLoggedRequest( - "session/prompt", - requestPayload, - acp.agent.prompt(requestPayload), - ), + const cancelledResponse = { + stopReason: "cancelled", + } satisfies EffectAcpSchema.PromptResponse; + const promptRpcFiber = yield* runLoggedRequest( + "session/prompt", + requestPayload, + acp.agent.prompt(requestPayload), + ).pipe(Effect.forkIn(runtimeScope)); + yield* Ref.set(activePromptFiberRef, Option.some(promptRpcFiber)); + return yield* Fiber.join(promptRpcFiber).pipe( + Effect.catchCause((cause) => + Cause.hasInterruptsOnly(cause) + ? Effect.succeed(cancelledResponse) + : Effect.failCause(cause), + ), + Effect.ensuring( + Effect.gen(function* () { + yield* Fiber.interrupt(promptRpcFiber).pipe(Effect.ignore); + yield* Ref.set(activePromptFiberRef, Option.none()); + }), ), Effect.tap(() => closeActiveAssistantSegment({ @@ -626,7 +747,17 @@ export const make = ( }), ), cancel: getStartedState.pipe( - Effect.flatMap((started) => acp.agent.cancel({ sessionId: started.sessionId })), + Effect.flatMap((started) => + Effect.gen(function* () { + const activePromptFiber = yield* Ref.get(activePromptFiberRef); + if (Option.isSome(activePromptFiber)) { + yield* Fiber.interrupt(activePromptFiber.value).pipe(Effect.ignore); + } + yield* acp.agent + .cancel({ sessionId: started.sessionId }) + .pipe(Effect.ignore, Effect.forkIn(runtimeScope)); + }), + ), ), setMode: (modeId) => Ref.get(modeStateRef).pipe( @@ -705,7 +836,7 @@ const handleSessionUpdate = ({ assistantSegmentRef, params, }: { - readonly queue: Queue.Queue; + readonly queue: Queue.Queue; readonly modeStateRef: Ref.Ref; readonly toolCallsRef: Ref.Ref>; readonly assistantSegmentRef: Ref.Ref; @@ -801,7 +932,7 @@ const ensureActiveAssistantSegment = ({ assistantSegmentRef, sessionId, }: { - readonly queue: Queue.Queue; + readonly queue: Queue.Queue; readonly assistantSegmentRef: Ref.Ref; readonly sessionId: string; }) => @@ -838,7 +969,7 @@ const closeActiveAssistantSegment = ({ queue, assistantSegmentRef, }: { - readonly queue: Queue.Queue; + readonly queue: Queue.Queue; readonly assistantSegmentRef: Ref.Ref; }) => Ref.modify(assistantSegmentRef, (current) => { diff --git a/apps/server/src/provider/acp/GrokAcpSupport.ts b/apps/server/src/provider/acp/GrokAcpSupport.ts index ee8af1e5266..3e6d63a5393 100644 --- a/apps/server/src/provider/acp/GrokAcpSupport.ts +++ b/apps/server/src/provider/acp/GrokAcpSupport.ts @@ -8,6 +8,7 @@ import type * as EffectAcpSchema from "effect-acp/schema"; import { normalizeModelSlug } from "@t3tools/shared/model"; import * as AcpSessionRuntime from "./AcpSessionRuntime.ts"; +import { makeXAiPromptCompletionRuntime } from "./XAiAcpExtension.ts"; const GROK_API_KEY_ENV = "XAI_API_KEY"; const GROK_OAUTH2_REFERRER_ENV = "GROK_OAUTH2_REFERRER"; @@ -68,9 +69,10 @@ export const makeGrokAcpRuntime = ( ), ), ); - return yield* Effect.service(AcpSessionRuntime.AcpSessionRuntime).pipe( + const runtime = yield* Effect.service(AcpSessionRuntime.AcpSessionRuntime).pipe( Effect.provide(acpContext), ); + return yield* makeXAiPromptCompletionRuntime(runtime); }); export function resolveGrokAcpBaseModelId(model: string | null | undefined): string { diff --git a/apps/server/src/provider/acp/XAiAcpExtension.test.ts b/apps/server/src/provider/acp/XAiAcpExtension.test.ts index a42561d9ae2..c435269fd76 100644 --- a/apps/server/src/provider/acp/XAiAcpExtension.test.ts +++ b/apps/server/src/provider/acp/XAiAcpExtension.test.ts @@ -1,12 +1,39 @@ -import { describe, expect, it } from "vite-plus/test"; +// @effect-diagnostics nodeBuiltinImport:off +import * as NodePath from "node:path"; +import * as NodeURL from "node:url"; + +import * as NodeServices from "@effect/platform-node/NodeServices"; +import { it } from "@effect/vitest"; +import * as Effect from "effect/Effect"; import * as Schema from "effect/Schema"; +import { describe, expect } from "vite-plus/test"; import { extractXAiAskUserQuestions, makeXAiAskUserQuestionCancelledResponse, makeXAiAskUserQuestionResponse, + makeXAiPromptCompletionRuntime, XAiAskUserQuestionRequest, } from "./XAiAcpExtension.ts"; +import * as AcpSessionRuntime from "./AcpSessionRuntime.ts"; + +const __dirname = NodePath.dirname(NodeURL.fileURLToPath(import.meta.url)); +const mockAgentPath = NodePath.join(__dirname, "../../../scripts/acp-mock-agent.ts"); + +const makePromptCompletionRuntime = (env: NodeJS.ProcessEnv) => + Effect.gen(function* () { + const runtime = yield* AcpSessionRuntime.make({ + spawn: { + command: process.execPath, + args: [mockAgentPath], + env, + }, + cwd: process.cwd(), + clientInfo: { name: "t3-test", version: "0.0.0" }, + authMethodId: "test", + }); + return yield* makeXAiPromptCompletionRuntime(runtime); + }); const decodeXAiAskUserQuestionRequest = Schema.decodeUnknownSync(XAiAskUserQuestionRequest); @@ -247,4 +274,59 @@ describe("XAiAcpExtension", () => { }, }); }); + + it.effect("resolves a hung standard prompt from xAI prompt completion", () => + Effect.gen(function* () { + const runtime = yield* makePromptCompletionRuntime({ + T3_ACP_EMIT_XAI_PROMPT_COMPLETE_THEN_HANG: "1", + }); + yield* runtime.start(); + + const promptResult = yield* runtime.prompt({ + prompt: [{ type: "text", text: "hi" }], + }); + const promptId = promptResult._meta?.promptId; + + expect(typeof promptId).toBe("string"); + expect(promptResult).toMatchObject({ + stopReason: "end_turn", + _meta: { + sessionId: "mock-session-1", + promptId, + requestId: promptId, + }, + }); + }).pipe(Effect.scoped, Effect.provide(NodeServices.layer)), + ); + + it.effect("ignores stale xAI completion from an already settled prompt", () => + Effect.gen(function* () { + const runtime = yield* makePromptCompletionRuntime({ + T3_ACP_EMIT_STALE_XAI_PROMPT_COMPLETE_BEFORE_SECOND_HANG: "1", + }); + yield* runtime.start(); + + const firstPromptResult = yield* runtime.prompt({ + prompt: [{ type: "text", text: "first" }], + }); + expect(firstPromptResult).toMatchObject({ + stopReason: "end_turn", + _meta: { promptId: "mock-stale-xai-prompt-1" }, + }); + + const secondPromptResult = yield* runtime.prompt({ + prompt: [{ type: "text", text: "second" }], + }); + const secondPromptId = secondPromptResult._meta?.promptId; + expect(typeof secondPromptId).toBe("string"); + expect(secondPromptId).not.toBe("mock-stale-xai-prompt-1"); + expect(secondPromptResult).toMatchObject({ + stopReason: "end_turn", + _meta: { + promptId: secondPromptId, + requestId: secondPromptId, + }, + }); + }).pipe(Effect.scoped, Effect.provide(NodeServices.layer)), + ); }); diff --git a/apps/server/src/provider/acp/XAiAcpExtension.ts b/apps/server/src/provider/acp/XAiAcpExtension.ts index 6c774c7f8d5..d36a5fcfc89 100644 --- a/apps/server/src/provider/acp/XAiAcpExtension.ts +++ b/apps/server/src/provider/acp/XAiAcpExtension.ts @@ -1,5 +1,29 @@ import type { ProviderUserInputAnswers, UserInputQuestion } from "@t3tools/contracts"; +import * as Deferred from "effect/Deferred"; +import * as Effect from "effect/Effect"; +import * as Ref from "effect/Ref"; import * as Schema from "effect/Schema"; +import type * as EffectAcpSchema from "effect-acp/schema"; + +import type * as AcpSessionRuntime from "./AcpSessionRuntime.ts"; + +const XAiPromptCompleteNotification = Schema.Struct({ + sessionId: Schema.String, + promptId: Schema.optional(Schema.String), + stopReason: Schema.optional(Schema.String), + agentResult: Schema.optional(Schema.NullOr(Schema.Unknown)), +}); + +type XAiPromptCompleteNotification = typeof XAiPromptCompleteNotification.Type; + +interface PendingXAiPromptCompletion { + readonly sessionId: string; + readonly promptId: string; + readonly deferred: Deferred.Deferred; +} + +const completedXAiPromptIdLimit = 128; +const xAiStopReasonMissingMetaKey = "xAiStopReasonMissing"; const XAiAskUserQuestionOption = Schema.Struct({ label: Schema.String, @@ -171,3 +195,238 @@ export function makeXAiAskUserQuestionResponse( export function makeXAiAskUserQuestionCancelledResponse(): XAiAskUserQuestionCancelledResponse { return { outcome: "cancelled" }; } + +/** + * Adds Grok's private prompt-completion fallback around a standards-only ACP runtime. + * The underlying runtime remains unaware of xAI methods and metadata. + */ +export const makeXAiPromptCompletionRuntime = Effect.fn("makeXAiPromptCompletionRuntime")( + function* (runtime: AcpSessionRuntime.AcpSessionRuntime["Service"]) { + const activeSessionIdRef = yield* Ref.make(undefined); + const pendingRef = yield* Ref.make>([]); + const completedPromptIdsRef = yield* Ref.make>([]); + let nextPromptFallbackId = 0; + const allocatePromptFallbackId = Effect.sync(() => { + nextPromptFallbackId += 1; + return `t3-xai-prompt-${nextPromptFallbackId}`; + }); + + yield* runtime.handleExtNotification( + "_x.ai/session/prompt_complete", + XAiPromptCompleteNotification, + (notification) => + resolveXAiPromptCompletionFallback({ + pendingRef, + completedPromptIdsRef, + notification, + }), + ); + + return { + ...runtime, + start: () => + runtime + .start() + .pipe(Effect.tap((started) => Ref.set(activeSessionIdRef, started.sessionId))), + prompt: (payload) => + Effect.gen(function* () { + const sessionId = yield* Ref.get(activeSessionIdRef); + if (sessionId === undefined) { + return yield* runtime.prompt(payload); + } + + const promptId = yield* allocatePromptFallbackId; + const fallback = yield* registerXAiPromptCompletionFallback( + pendingRef, + sessionId, + promptId, + ); + const requestPayload = { + ...payload, + _meta: { + ...payload._meta, + promptId: fallback.promptId, + requestId: fallback.promptId, + }, + } satisfies Omit; + + return yield* Effect.raceFirst( + runtime.prompt(requestPayload), + Deferred.await(fallback.deferred), + ).pipe( + Effect.tap((response) => + rememberCompletedXAiPromptId(completedPromptIdsRef, response, fallback.promptId), + ), + Effect.ensuring(unregisterXAiPromptCompletionFallback(pendingRef, fallback.deferred)), + ); + }), + cancel: Ref.get(activeSessionIdRef).pipe( + Effect.flatMap((sessionId) => + sessionId === undefined + ? runtime.cancel + : abortPendingPromptCompletions(pendingRef, sessionId).pipe( + Effect.andThen(runtime.cancel), + ), + ), + ), + } satisfies AcpSessionRuntime.AcpSessionRuntime["Service"]; + }, +); + +const registerXAiPromptCompletionFallback = ( + pendingRef: Ref.Ref>, + sessionId: string, + promptId: string, +) => + Deferred.make().pipe( + Effect.tap((deferred) => + Ref.update(pendingRef, (pending) => [...pending, { sessionId, promptId, deferred }]), + ), + Effect.map((deferred) => ({ deferred, promptId })), + ); + +const unregisterXAiPromptCompletionFallback = ( + pendingRef: Ref.Ref>, + deferred: Deferred.Deferred, +) => Ref.update(pendingRef, (pending) => pending.filter((entry) => entry.deferred !== deferred)); + +const abortPendingPromptCompletions = ( + pendingRef: Ref.Ref>, + sessionId: string, +) => + Ref.modify(pendingRef, (pending) => { + const [toAbort, remaining] = pending.reduce< + [ReadonlyArray, ReadonlyArray] + >( + ([aborting, kept], entry) => + entry.sessionId === sessionId ? [[...aborting, entry], kept] : [aborting, [...kept, entry]], + [[], []], + ); + if (toAbort.length === 0) { + return [Effect.void, pending] as const; + } + return [ + Effect.forEach( + toAbort, + (entry) => + Deferred.succeed( + entry.deferred, + promptResponseFromXAi({ + sessionId: entry.sessionId, + promptId: entry.promptId, + stopReason: "cancelled", + agentResult: null, + }), + ), + { concurrency: "unbounded" }, + ).pipe(Effect.asVoid), + remaining, + ] as const; + }).pipe(Effect.flatten); + +const resolveXAiPromptCompletionFallback = ({ + pendingRef, + completedPromptIdsRef, + notification, +}: { + readonly pendingRef: Ref.Ref>; + readonly completedPromptIdsRef: Ref.Ref>; + readonly notification: XAiPromptCompleteNotification; +}) => + Ref.get(completedPromptIdsRef).pipe( + Effect.flatMap((completedPromptIds) => { + if ( + notification.promptId !== undefined && + completedPromptIds.includes(notification.promptId) + ) { + return Effect.void; + } + return Ref.modify(pendingRef, (pending) => { + const index = + notification.promptId !== undefined + ? pending.findIndex( + (entry) => + entry.sessionId === notification.sessionId && + entry.promptId === notification.promptId, + ) + : pending.findIndex((entry) => entry.sessionId === notification.sessionId); + if (index < 0) { + return [Effect.void, pending] as const; + } + const entry = pending[index]; + if (!entry) { + return [Effect.void, pending] as const; + } + return [ + Deferred.succeed(entry.deferred, promptResponseFromXAi(notification)).pipe(Effect.asVoid), + [...pending.slice(0, index), ...pending.slice(index + 1)], + ] as const; + }).pipe(Effect.flatten); + }), + ); + +const rememberCompletedXAiPromptId = ( + completedPromptIdsRef: Ref.Ref>, + response: EffectAcpSchema.PromptResponse, + fallbackPromptId: string, +) => { + const promptId = promptIdFromResponse(response) ?? fallbackPromptId; + return Ref.update(completedPromptIdsRef, (completedPromptIds) => { + if (completedPromptIds.includes(promptId)) { + return completedPromptIds; + } + return [...completedPromptIds, promptId].slice(-completedXAiPromptIdLimit); + }); +}; + +function promptIdFromResponse(response: EffectAcpSchema.PromptResponse): string | undefined { + const meta = response._meta; + if (meta === null || typeof meta !== "object") { + return undefined; + } + const promptId = meta.promptId ?? meta.requestId; + return typeof promptId === "string" && promptId.length > 0 ? promptId : undefined; +} + +export function promptResponseHasMissingXAiStopReason( + response: EffectAcpSchema.PromptResponse, +): boolean { + const meta = response._meta; + return meta !== null && typeof meta === "object" && meta[xAiStopReasonMissingMetaKey] === true; +} + +function promptResponseFromXAi( + notification: XAiPromptCompleteNotification, +): EffectAcpSchema.PromptResponse { + const stopReason = normalizeXAiStopReason(notification.stopReason); + const meta: Record = { + sessionId: notification.sessionId, + }; + if (notification.stopReason === undefined) { + meta[xAiStopReasonMissingMetaKey] = true; + } + if (notification.promptId !== undefined) { + meta.promptId = notification.promptId; + meta.requestId = notification.promptId; + } + if (notification.agentResult !== undefined) { + meta.agentResult = notification.agentResult; + } + return { + stopReason, + _meta: meta, + }; +} + +function normalizeXAiStopReason(value: string | undefined): EffectAcpSchema.StopReason { + switch (value) { + case "cancelled": + case "end_turn": + case "max_tokens": + case "max_turn_requests": + case "refusal": + return value; + default: + return "end_turn"; + } +} diff --git a/apps/web/src/components/ChatView.logic.test.ts b/apps/web/src/components/ChatView.logic.test.ts index 43ed895c0db..0a0103df183 100644 --- a/apps/web/src/components/ChatView.logic.test.ts +++ b/apps/web/src/components/ChatView.logic.test.ts @@ -6,6 +6,7 @@ import { MAX_HIDDEN_MOUNTED_PREVIEW_THREADS, MAX_HIDDEN_MOUNTED_TERMINAL_THREADS, buildExpiredTerminalContextToastCopy, + buildThreadTurnInterruptInput, createLocalDispatchSnapshot, deriveComposerSendState, getStartedThreadModelChangeBlockReason, @@ -69,6 +70,30 @@ const readySession = { updatedAt: "2026-03-29T00:00:10.000Z", }; +describe("buildThreadTurnInterruptInput", () => { + it("targets the session's active running turn", () => { + const activeTurnId = TurnId.make("turn-running"); + + expect( + buildThreadTurnInterruptInput( + makeThread({ + session: { + ...readySession, + status: "running", + activeTurnId, + }, + }), + ), + ).toEqual({ threadId, turnId: activeTurnId }); + }); + + it("omits a turn id when the session is not running", () => { + expect(buildThreadTurnInterruptInput(makeThread({ session: readySession }))).toEqual({ + threadId, + }); + }); +}); + describe("deriveComposerSendState", () => { it("treats expired terminal pills as non-sendable content", () => { const state = deriveComposerSendState({ diff --git a/apps/web/src/components/ChatView.logic.ts b/apps/web/src/components/ChatView.logic.ts index 36947caae6f..705793ec77e 100644 --- a/apps/web/src/components/ChatView.logic.ts +++ b/apps/web/src/components/ChatView.logic.ts @@ -74,6 +74,17 @@ export function shouldWriteThreadErrorToCurrentServerThread(input: { ); } +export function buildThreadTurnInterruptInput(thread: Pick): { + threadId: ThreadId; + turnId?: TurnId; +} { + const runningTurnId = thread.session?.status === "running" ? thread.session.activeTurnId : null; + return { + threadId: thread.id, + ...(runningTurnId !== null ? { turnId: runningTurnId } : {}), + }; +} + export function reconcileMountedTerminalThreadIds(input: { currentThreadIds: ReadonlyArray; openThreadIds: ReadonlyArray; diff --git a/apps/web/src/components/ChatView.tsx b/apps/web/src/components/ChatView.tsx index 9fb8d647b4e..81bd44c08cf 100644 --- a/apps/web/src/components/ChatView.tsx +++ b/apps/web/src/components/ChatView.tsx @@ -215,6 +215,7 @@ import { MAX_HIDDEN_MOUNTED_TERMINAL_THREADS, buildExpiredTerminalContextToastCopy, buildLocalDraftThread, + buildThreadTurnInterruptInput, collectUserMessageBlobPreviewUrls, createLocalDispatchSnapshot, deriveComposerSendState, @@ -3942,9 +3943,7 @@ function ChatViewContent(props: ChatViewProps) { if (!activeThread) return; const result = await interruptThreadTurn({ environmentId, - input: { - threadId: activeThread.id, - }, + input: buildThreadTurnInterruptInput(activeThread), }); if (result._tag === "Failure" && !isAtomCommandInterrupted(result)) { const error = squashAtomCommandFailure(result); @@ -4775,6 +4774,11 @@ function ChatViewContent(props: ChatViewProps) { listRef={legendListRef} timelineEntries={timelineEntries} latestTurn={activeLatestTurn} + runningTurnId={ + activeThread.session?.status === "running" + ? activeThread.session.activeTurnId + : null + } turnDiffSummaryByAssistantMessageId={turnDiffSummaryByAssistantMessageId} activeThreadEnvironmentId={activeThread.environmentId} routeThreadKey={routeThreadKey} diff --git a/apps/web/src/components/chat/MessagesTimeline.logic.test.ts b/apps/web/src/components/chat/MessagesTimeline.logic.test.ts index 50ee10b4169..1676d2d7c85 100644 --- a/apps/web/src/components/chat/MessagesTimeline.logic.test.ts +++ b/apps/web/src/components/chat/MessagesTimeline.logic.test.ts @@ -794,6 +794,67 @@ describe("deriveMessagesTimelineRows", () => { ]); }); + it("does not fold the session's running turn when latestTurn regresses", () => { + const rows = deriveMessagesTimelineRows({ + timelineEntries: [ + { + id: "previous-work-entry", + kind: "work", + createdAt: "2026-01-01T00:00:05Z", + entry: { + id: "previous-work", + createdAt: "2026-01-01T00:00:05Z", + turnId: "turn-1" as never, + label: "Read files", + tone: "tool" as const, + }, + }, + { + id: "user-followup-entry", + kind: "message", + createdAt: "2026-01-01T00:01:00Z", + message: { + id: "user-followup" as never, + role: "user", + text: "continue", + turnId: null, + createdAt: "2026-01-01T00:01:00Z", + updatedAt: "2026-01-01T00:01:00Z", + streaming: false, + }, + }, + { + id: "running-work-entry", + kind: "work", + createdAt: "2026-01-01T00:01:05Z", + entry: { + id: "running-work", + createdAt: "2026-01-01T00:01:05Z", + turnId: "turn-2" as never, + label: "Searched files", + tone: "tool" as const, + }, + }, + ], + latestTurn: { + turnId: "turn-1" as never, + state: "completed", + startedAt: "2026-01-01T00:00:00Z", + completedAt: "2026-01-01T00:00:25Z", + }, + runningTurnId: "turn-2" as never, + isWorking: true, + activeTurnStartedAt: "2026-01-01T00:01:00Z", + turnDiffSummaryByAssistantMessageId: new Map(), + revertTurnCountByUserMessageId: new Map(), + }); + + expect(rows.filter((row) => row.kind === "turn-fold").map((row) => row.turnId)).toEqual([ + "turn-1", + ]); + expect(rows.map((row) => row.id)).toContain("running-work-entry"); + }); + it("only shows assistant metadata on the terminal assistant message", () => { const rows = deriveMessagesTimelineRows({ timelineEntries: [ diff --git a/apps/web/src/components/chat/MessagesTimeline.logic.ts b/apps/web/src/components/chat/MessagesTimeline.logic.ts index 1426f1deee2..86212576f0b 100644 --- a/apps/web/src/components/chat/MessagesTimeline.logic.ts +++ b/apps/web/src/components/chat/MessagesTimeline.logic.ts @@ -149,13 +149,20 @@ interface TurnFold { } /** - * The latest turn counts as unsettled while it is still running (or has not - * recorded a completion). This is deliberately keyed on the turn's own - * lifecycle rather than transient working state: right after the user sends - * a message, the previous turn is still the "active" one until the server - * creates the new turn, and folding must not flicker through that window. + * The session's running turn is authoritative when latestTurn briefly lags or + * regresses behind it. Otherwise, the latest turn counts as unsettled while it + * is still running (or has not recorded a completion). This is deliberately + * keyed on turn lifecycle rather than transient working state: right after the + * user sends a message, the previous turn is still the "active" one until the + * server creates the new turn, and folding must not flicker through that window. */ -function deriveUnsettledTurnId(latestTurn: TimelineLatestTurn | null): TurnId | null { +function deriveUnsettledTurnId( + latestTurn: TimelineLatestTurn | null, + runningTurnId: TurnId | null, +): TurnId | null { + if (runningTurnId !== null) { + return runningTurnId; + } if (!latestTurn) { return null; } @@ -291,6 +298,7 @@ function deriveTurnFolds(input: { export function deriveMessagesTimelineRows(input: { timelineEntries: ReadonlyArray; latestTurn?: TimelineLatestTurn | null; + runningTurnId?: TurnId | null; expandedTurnIds?: ReadonlySet; isWorking: boolean; activeTurnStartedAt: string | null; @@ -302,7 +310,10 @@ export function deriveMessagesTimelineRows(input: { input.timelineEntries.flatMap((entry) => (entry.kind === "message" ? [entry.message] : [])), ); const terminalAssistantMessageIds = deriveTerminalAssistantMessageIds(input.timelineEntries); - const unsettledTurnId = deriveUnsettledTurnId(input.latestTurn ?? null); + const unsettledTurnId = deriveUnsettledTurnId( + input.latestTurn ?? null, + input.runningTurnId ?? null, + ); const foldsByAnchorEntryId = deriveTurnFolds({ timelineEntries: input.timelineEntries, terminalAssistantMessageIds, diff --git a/apps/web/src/components/chat/MessagesTimeline.test.tsx b/apps/web/src/components/chat/MessagesTimeline.test.tsx index 3008bd2ba9e..a69ce17c7d2 100644 --- a/apps/web/src/components/chat/MessagesTimeline.test.tsx +++ b/apps/web/src/components/chat/MessagesTimeline.test.tsx @@ -109,6 +109,7 @@ function buildProps() { activeTurnStartedAt: null, listRef: createRef(), latestTurn: null, + runningTurnId: null, turnDiffSummaryByAssistantMessageId: new Map(), routeThreadKey: "environment-local:thread-1", onOpenTurnDiff: () => {}, diff --git a/apps/web/src/components/chat/MessagesTimeline.tsx b/apps/web/src/components/chat/MessagesTimeline.tsx index 55c982c64be..cf8224d51f6 100644 --- a/apps/web/src/components/chat/MessagesTimeline.tsx +++ b/apps/web/src/components/chat/MessagesTimeline.tsx @@ -153,6 +153,7 @@ interface MessagesTimelineProps { listRef: React.RefObject; timelineEntries: ReturnType; latestTurn: TimelineLatestTurn | null; + runningTurnId: TurnId | null; turnDiffSummaryByAssistantMessageId: Map; routeThreadKey: string; onOpenTurnDiff: (turnId: TurnId, filePath?: string) => void; @@ -182,6 +183,7 @@ export const MessagesTimeline = memo(function MessagesTimeline({ listRef, timelineEntries, latestTurn, + runningTurnId, turnDiffSummaryByAssistantMessageId, routeThreadKey, onOpenTurnDiff, @@ -272,6 +274,7 @@ export const MessagesTimeline = memo(function MessagesTimeline({ deriveMessagesTimelineRows({ timelineEntries, latestTurn, + runningTurnId, expandedTurnIds, isWorking, activeTurnStartedAt, @@ -281,6 +284,7 @@ export const MessagesTimeline = memo(function MessagesTimeline({ [ timelineEntries, latestTurn, + runningTurnId, expandedTurnIds, isWorking, activeTurnStartedAt, diff --git a/packages/effect-acp/src/client.test.ts b/packages/effect-acp/src/client.test.ts index c732f80ef35..779a6499e4b 100644 --- a/packages/effect-acp/src/client.test.ts +++ b/packages/effect-acp/src/client.test.ts @@ -17,18 +17,59 @@ import { it, assert } from "@effect/vitest"; import * as AcpClient from "./client.ts"; import * as AcpSchema from "./_generated/schema.gen.ts"; import * as AcpError from "./errors.ts"; -import { encodeJsonl, jsonRpcRequest, jsonRpcResponse } from "./_internal/shared.ts"; +import { + encodeJsonl, + jsonRpcNotification, + jsonRpcRequest, + jsonRpcResponse, +} from "./_internal/shared.ts"; import { makeInMemoryStdio } from "./_internal/stdio.ts"; const InitializeRequest = jsonRpcRequest("initialize", AcpSchema.InitializeRequest); const InitializeResponse = jsonRpcResponse(AcpSchema.InitializeResponse); const ExtRequest = jsonRpcRequest("x/test", Schema.Struct({ hello: Schema.String })); const ExtResponse = jsonRpcResponse(Schema.Struct({ ok: Schema.Boolean })); +const PromptRequest = jsonRpcRequest("session/prompt", AcpSchema.PromptRequest); +const PromptResponse = jsonRpcResponse(AcpSchema.PromptResponse); +const decodePromptRequestLine = Schema.decodeEffect(Schema.fromJsonString(PromptRequest)); +const XAiPromptCompleteNotification = jsonRpcNotification( + "_x.ai/session/prompt_complete", + Schema.Struct({ + sessionId: Schema.String, + promptId: Schema.String, + stopReason: Schema.String, + agentResult: Schema.NullOr(Schema.Unknown), + }), +); +const XAiQueueChangedNotification = jsonRpcNotification( + "_x.ai/queue/changed", + Schema.Struct({ + sessionId: Schema.String, + entries: Schema.Array(Schema.Unknown), + }), +); +const XAiSessionsChangedNotification = jsonRpcNotification( + "_x.ai/sessions/changed", + Schema.Struct({ + upserted: Schema.Array(Schema.Unknown), + removed: Schema.Array(Schema.Unknown), + }), +); const mockPeerPath = Effect.map(Effect.service(Path.Path), (path) => path.join(import.meta.dirname, "../test/fixtures/acp-mock-peer.ts"), ); const mockPeerArgs = (path: string) => [path]; +function concatBytes(chunks: ReadonlyArray): Uint8Array { + const batch = new Uint8Array(chunks.reduce((total, chunk) => total + chunk.length, 0)); + let offset = 0; + for (const chunk of chunks) { + batch.set(chunk, offset); + offset += chunk.length; + } + return batch; +} + it.layer(NodeServices.layer)("effect-acp client", (it) => { const makeHandle = (env?: Record) => Effect.gen(function* () { @@ -446,4 +487,90 @@ it.layer(NodeServices.layer)("effect-acp client", (it) => { yield* Scope.close(scope, Exit.void); }), ); + + it.effect( + "routes a standard prompt response after Grok extension notifications in the same batch", + () => + Effect.gen(function* () { + const { stdio, input, output } = yield* makeInMemoryStdio(); + const scope = yield* Scope.make(); + const acp = yield* AcpClient.make(stdio).pipe(Effect.provideService(Scope.Scope, scope)); + + const promptFiber = yield* acp.agent + .prompt({ + sessionId: "grok-session-1", + prompt: [{ type: "text", text: "run the ls command" }], + }) + .pipe(Effect.forkScoped); + + const outbound = yield* Queue.take(output); + const decodedPrompt = yield* decodePromptRequestLine(outbound); + + const responseBatch = concatBytes( + yield* Effect.all([ + encodeJsonl(XAiQueueChangedNotification, { + jsonrpc: "2.0", + method: "_x.ai/queue/changed", + params: { sessionId: "grok-session-1", entries: [] }, + }), + encodeJsonl(XAiPromptCompleteNotification, { + jsonrpc: "2.0", + method: "_x.ai/session/prompt_complete", + params: { + sessionId: "grok-session-1", + promptId: "prompt-1", + stopReason: "end_turn", + agentResult: null, + }, + }), + encodeJsonl(XAiSessionsChangedNotification, { + jsonrpc: "2.0", + method: "_x.ai/sessions/changed", + params: { + upserted: [ + { + sessionId: "grok-session-1", + title: null, + cwd: process.cwd(), + isWorktree: false, + modelId: "grok-composer-2.5-fast", + yolo: false, + activity: "idle", + resident: true, + lastChangeUnixMs: 1_710_000_000_000, + origin: { kind: "local" }, + }, + ], + removed: [], + }, + }), + encodeJsonl(PromptResponse, { + jsonrpc: "2.0", + id: decodedPrompt.id, + result: { + stopReason: "end_turn", + _meta: { + sessionId: "grok-session-1", + requestId: "prompt-1", + promptId: "prompt-1", + modelId: "grok-composer-2.5-fast", + }, + }, + }), + ]), + ); + yield* Queue.offer(input, responseBatch); + + assert.deepEqual(yield* Fiber.join(promptFiber), { + stopReason: "end_turn", + _meta: { + sessionId: "grok-session-1", + requestId: "prompt-1", + promptId: "prompt-1", + modelId: "grok-composer-2.5-fast", + }, + }); + yield* Scope.close(scope, Exit.void); + }), + ); });