diff --git a/apps/server/scripts/acp-mock-agent.ts b/apps/server/scripts/acp-mock-agent.ts index 9fbdeee03c86..621a518b0ae9 100644 --- a/apps/server/scripts/acp-mock-agent.ts +++ b/apps/server/scripts/acp-mock-agent.ts @@ -619,6 +619,12 @@ const program = Effect.gen(function* () { Effect.gen(function* () { const requestedSessionId = String(request.sessionId ?? sessionId); promptCount += 1; + if ( + process.env.T3_ACP_CRASH_PROMPT === "1" && + request.prompt.some((part) => part.type === "text" && part.text === "crash now") + ) { + return yield* Effect.sync(() => process.exit(23)); + } if (completeFirstPromptOnCancel && promptCount === 1) { yield* agent.client.sessionUpdate({ diff --git a/apps/server/src/provider/Layers/GrokAdapter.test.ts b/apps/server/src/provider/Layers/GrokAdapter.test.ts index ecd73af72dbe..fb17d6839fda 100644 --- a/apps/server/src/provider/Layers/GrokAdapter.test.ts +++ b/apps/server/src/provider/Layers/GrokAdapter.test.ts @@ -363,6 +363,66 @@ it.layer(grokAdapterTestLayer)("GrokAdapterLive", (it) => { }), ); + it.effect("retires a crashed process so a deliberate retry can resume", () => + Effect.gen(function* () { + const threadId = ThreadId.make("grok-crash-recovery"); + const tempDir = yield* Effect.promise(() => + NodeFSP.mkdtemp(NodePath.join(NodeOS.tmpdir(), "grok-crash-recovery-")), + ); + const requestLogPath = NodePath.join(tempDir, "requests.ndjson"); + const wrapper = yield* Effect.promise(() => + makeMockGrokWrapper({ + T3_ACP_CRASH_PROMPT: "1", + T3_ACP_REQUEST_LOG_PATH: requestLogPath, + }), + ); + const adapter = yield* makeTestAdapter(wrapper); + const exited = + yield* Deferred.make>(); + const events = yield* Stream.runForEach(adapter.streamEvents, (event) => + event.type === "session.exited" + ? Deferred.succeed(exited, event).pipe(Effect.asVoid) + : Effect.void, + ).pipe(Effect.forkChild); + const input = { + threadId, + provider: ProviderDriverKind.make("grok"), + cwd: process.cwd(), + runtimeMode: "full-access" as const, + }; + const session = yield* adapter.startSession(input); + const failure = yield* Effect.flip( + adapter.sendTurn({ threadId, input: "crash now", attachments: [] }), + ); + assert.isDefined(failure); + assert.isFalse(yield* adapter.hasSession(threadId)); + const retryDuringTeardown = yield* Effect.flip( + adapter.sendTurn({ threadId, input: "retry during teardown", attachments: [] }), + ); + assert.equal(retryDuringTeardown._tag, "ProviderAdapterSessionNotFoundError"); + assert.equal((yield* Deferred.await(exited)).payload.exitKind, "error"); + assert.deepStrictEqual(yield* adapter.listSessions(), []); + yield* Fiber.interrupt(events); + yield* adapter.startSession({ ...input, resumeCursor: session.resumeCursor }); + const turn = yield* adapter.sendTurn({ threadId, input: "retry now", attachments: [] }); + assert.equal(turn.threadId, threadId); + yield* adapter.stopSession(threadId); + const requests = yield* Effect.promise(() => readJsonLines(requestLogPath)); + assert.equal(requests.filter((request) => request.method === "session/new").length, 1); + assert.deepStrictEqual(session.resumeCursor, { + schemaVersion: 1, + sessionId: "mock-session-1", + }); + const resumes = requests.filter((request) => request.method === "session/load"); + assert.equal(resumes.length, 1); + assert.deepStrictEqual(resumes[0]?.params, { + sessionId: "mock-session-1", + cwd: process.cwd(), + mcpServers: [], + }); + }), + ); + it.effect("starts a session and maps mock ACP prompt flow to runtime events", () => Effect.gen(function* () { const threadId = ThreadId.make("grok-mock-thread"); diff --git a/apps/server/src/provider/Layers/GrokAdapter.ts b/apps/server/src/provider/Layers/GrokAdapter.ts index bec0fc39faa0..576d7a322977 100644 --- a/apps/server/src/provider/Layers/GrokAdapter.ts +++ b/apps/server/src/provider/Layers/GrokAdapter.ts @@ -175,6 +175,7 @@ interface GrokSessionContext { currentModelId: string | undefined; currentReasoningEffort: string | undefined; stopped: boolean; + terminated: boolean; /** Live monitor/shell identities and their originating turns. */ readonly backgroundTasks: Map; } @@ -345,6 +346,7 @@ export function grokPromptSettlementBelongsToContext(input: { export function makeGrokAdapter(grokSettings: GrokSettings, options?: GrokAdapterLiveOptions) { return Effect.gen(function* () { + const ownerScope = yield* Effect.scope; const boundInstanceId = options?.instanceId ?? ProviderInstanceId.make("grok"); const fileSystem = yield* FileSystem.FileSystem; const path = yield* Path.Path; @@ -923,7 +925,7 @@ export function makeGrokAdapter(grokSettings: GrokSettings, options?: GrokAdapte threadId: ThreadId, ): Effect.Effect => { const ctx = sessions.get(threadId); - if (!ctx || ctx.stopped) { + if (!ctx || ctx.stopped || ctx.terminated) { return Effect.fail( new ProviderAdapterSessionNotFoundError({ provider: PROVIDER, threadId }), ); @@ -947,7 +949,7 @@ export function makeGrokAdapter(grokSettings: GrokSettings, options?: GrokAdapte ...(yield* makeEventStamp()), provider: PROVIDER, threadId: ctx.threadId, - payload: { exitKind: "graceful" }, + payload: { exitKind: ctx.terminated ? "error" : "graceful" }, }); }); @@ -1317,12 +1319,35 @@ export function makeGrokAdapter(grokSettings: GrokSettings, options?: GrokAdapte ? normalizeGrokReasoningEffort(requestedStartReasoningEffort) : currentStartReasoningEffort, stopped: false, + terminated: false, backgroundTasks: new Map(), }; const nf = yield* Stream.runDrain( Stream.mapEffect(acp.getEvents(), (event) => Effect.gen(function* () { + if (event._tag === "ConnectionTerminated") { + ctx.terminated = true; + yield* withThreadLock( + ctx.threadId, + Effect.gen(function* () { + if (sessions.get(ctx.threadId) !== ctx) return; + if (ctx.activeTurnId) { + yield* settlePromptInFlight( + ctx.threadId, + ctx.activeTurnId, + ctx.acpSessionId, + { + errorMessage: "Grok connection terminated.", + settleAllPrompts: true, + }, + ); + } + yield* stopSessionInternal(ctx); + }), + ).pipe(Effect.forkIn(ownerScope)); + return; + } if (event._tag === "EventStreamBarrier") { yield* Deferred.succeed(event.acknowledge, undefined); return; @@ -2166,12 +2191,16 @@ export function makeGrokAdapter(grokSettings: GrokSettings, options?: GrokAdapte ); const listSessions: GrokAdapterShape["listSessions"] = () => - Effect.sync(() => Array.from(sessions.values(), (c) => ({ ...c.session }))); + Effect.sync(() => + Array.from(sessions.values()) + .filter((c) => !c.terminated) + .map((c) => ({ ...c.session })), + ); const hasSession: GrokAdapterShape["hasSession"] = (threadId) => Effect.sync(() => { const c = sessions.get(threadId); - return c !== undefined && !c.stopped; + return c !== undefined && !c.stopped && !c.terminated; }); const stopAll: GrokAdapterShape["stopAll"] = () =>