diff --git a/apps/server/scripts/acp-mock-agent.ts b/apps/server/scripts/acp-mock-agent.ts index 541d8762b8..e2f35a9a3b 100644 --- a/apps/server/scripts/acp-mock-agent.ts +++ b/apps/server/scripts/acp-mock-agent.ts @@ -662,6 +662,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 e38de31156..666b8fdae0 100644 --- a/apps/server/src/provider/Layers/GrokAdapter.test.ts +++ b/apps/server/src/provider/Layers/GrokAdapter.test.ts @@ -9,6 +9,7 @@ import { assert, it } from "@effect/vitest"; import * as Deferred from "effect/Deferred"; import * as Effect from "effect/Effect"; import * as Fiber from "effect/Fiber"; +import * as FileSystem from "effect/FileSystem"; import * as Layer from "effect/Layer"; import * as Ref from "effect/Ref"; import * as Schema from "effect/Schema"; @@ -373,6 +374,84 @@ 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 incarnation = RuntimeSessionId.make("grok-crashed-incarnation"); + const fs = yield* FileSystem.FileSystem; + const tempDir = yield* fs.makeTempDirectoryScoped({ prefix: "grok-crash-recovery-" }); + const requestLogPath = NodePath.join(tempDir, "requests.ndjson"); + const wrapper = yield* Effect.sync(() => + writeFakeCli({ + directory: tempDir, + name: "fake-grok", + env: { T3_ACP_CRASH_PROMPT: "1", T3_ACP_REQUEST_LOG_PATH: requestLogPath }, + source: execScriptSource({ scriptPath: mockAgentPath }), + }), + ); + const adapter = yield* makeTestAdapter(wrapper); + const exited = + yield* Deferred.make>(); + const crashEvents: Array = []; + const events = yield* Stream.runForEach(adapter.streamEvents, (event) => + Effect.gen(function* () { + crashEvents.push(event); + if (event.type === "session.exited") yield* Deferred.succeed(exited, event); + }), + ).pipe(Effect.forkChild); + const input = { + threadId, + provider: ProviderDriverKind.make("grok"), + sessionIncarnationId: incarnation, + 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"); + const exitEvent = yield* Deferred.await(exited); + assert.equal(exitEvent.payload.exitKind, "error"); + assert.equal(exitEvent.sessionIncarnationId, incarnation); + const failedTurns = crashEvents.filter((event) => event.type === "turn.completed"); + assert.lengthOf(failedTurns, 1); + assert.equal(failedTurns[0]?.sessionIncarnationId, incarnation); + assert.equal(failedTurns[0]?.payload.state, "failed"); + assert.deepStrictEqual(yield* adapter.listSessions(), []); + yield* Fiber.interrupt(events); + const replacementIncarnation = RuntimeSessionId.make("grok-resumed-incarnation"); + yield* adapter.startSession({ + ...input, + sessionIncarnationId: replacementIncarnation, + resumeCursor: session.resumeCursor, + }); + const liveSessions = yield* adapter.listSessions(); + assert.equal(liveSessions[0]?.sessionIncarnationId, replacementIncarnation); + 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 2a92dba3d3..328e166c82 100644 --- a/apps/server/src/provider/Layers/GrokAdapter.ts +++ b/apps/server/src/provider/Layers/GrokAdapter.ts @@ -178,6 +178,7 @@ interface GrokSessionContext { currentModelId: string | undefined; currentReasoningEffort: string | undefined; stopped: boolean; + terminated: boolean; /** Live monitor and shell identities, with their originating turns. */ readonly backgroundTasks: Map; readonly completedBackgroundTaskIds: Map; @@ -352,6 +353,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; @@ -942,7 +944,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 }), ); @@ -969,7 +971,7 @@ export function makeGrokAdapter(grokSettings: GrokSettings, options?: GrokAdapte ...(ctx.sessionIncarnationId !== undefined ? { sessionIncarnationId: ctx.sessionIncarnationId } : {}), - payload: { exitKind: "graceful" }, + payload: { exitKind: ctx.terminated ? "error" : "graceful" }, }); }); @@ -1345,6 +1347,7 @@ export function makeGrokAdapter(grokSettings: GrokSettings, options?: GrokAdapte ? normalizeGrokReasoningEffort(requestedStartReasoningEffort) : currentStartReasoningEffort, stopped: false, + terminated: false, backgroundTasks: new Map(), completedBackgroundTaskIds: new Map(), ambiguousBackgroundTaskIds: new Set(), @@ -1353,6 +1356,29 @@ export function makeGrokAdapter(grokSettings: GrokSettings, options?: GrokAdapte 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, + ctx.sessionIncarnationId, + { + 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; @@ -2219,12 +2245,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"] = () =>