diff --git a/apps/server/integration/providerService.integration.test.ts b/apps/server/integration/providerService.integration.test.ts index e703af4b1f45..73f6caa2cb06 100644 --- a/apps/server/integration/providerService.integration.test.ts +++ b/apps/server/integration/providerService.integration.test.ts @@ -74,7 +74,10 @@ const makeIntegrationFixture = Effect.gen(function* () { Layer.succeed(ProviderEventLoggers, NoOpProviderEventLoggers), ).pipe(Layer.provide(SqlitePersistenceMemory)); - const layer = makeProviderServiceLive().pipe(Layer.provide(shared)); + const layer = makeProviderServiceLive().pipe( + Layer.provide(shared), + Layer.provide(NodeServices.layer), + ); return { cwd, diff --git a/apps/server/src/orchestration/Layers/ProviderCommandReactor.test.ts b/apps/server/src/orchestration/Layers/ProviderCommandReactor.test.ts index ce464565dc5f..0b48dd721d1b 100644 --- a/apps/server/src/orchestration/Layers/ProviderCommandReactor.test.ts +++ b/apps/server/src/orchestration/Layers/ProviderCommandReactor.test.ts @@ -610,6 +610,10 @@ describe("ProviderCommandReactor", () => { it("generates a worktree branch name for the first turn", async () => { const harness = await createHarness(); const now = "2026-01-01T00:00:00.000Z"; + const worktreePath = NodeFS.mkdtempSync( + NodePath.join(NodeOS.tmpdir(), "t3-provider-branch-worktree-"), + ); + createdBaseDirs.add(worktreePath); await Effect.runPromise( harness.engine.dispatch({ @@ -617,7 +621,7 @@ describe("ProviderCommandReactor", () => { commandId: CommandId.make("cmd-thread-branch"), threadId: ThreadId.make("thread-1"), branch: "t3code/1234abcd", - worktreePath: "/tmp/provider-project-worktree", + worktreePath, }), ); @@ -658,7 +662,7 @@ describe("ProviderCommandReactor", () => { expect(harness.generateBranchName.mock.calls[0]?.[0]).toMatchObject({ message: "Add a safer reconnect backoff.", }); - expect(harness.refreshStatus.mock.calls[0]?.[0]).toBe("/tmp/provider-project-worktree"); + expect(harness.refreshStatus.mock.calls[0]?.[0]).toBe(worktreePath); }); it("forwards codex model options through session start and turn send", async () => { @@ -1130,6 +1134,10 @@ describe("ProviderCommandReactor", () => { }, }); const now = "2026-01-01T00:00:00.000Z"; + const worktreePath = NodeFS.mkdtempSync( + NodePath.join(NodeOS.tmpdir(), "t3-provider-session-worktree-"), + ); + createdBaseDirs.add(worktreePath); await Effect.runPromise( harness.engine.dispatch({ @@ -1159,7 +1167,7 @@ describe("ProviderCommandReactor", () => { type: "thread.meta.update", commandId: CommandId.make("cmd-thread-worktree-change"), threadId: ThreadId.make("thread-1"), - worktreePath: "/tmp/provider-project-worktree", + worktreePath, }), ); @@ -1185,7 +1193,7 @@ describe("ProviderCommandReactor", () => { expect(harness.stopSession.mock.calls.length).toBe(0); expect(harness.startSession.mock.calls[1]?.[1]).toMatchObject({ threadId: ThreadId.make("thread-1"), - cwd: "/tmp/provider-project-worktree", + cwd: worktreePath, resumeCursor: { opaque: "resume-1" }, modelSelection: { instanceId: ProviderInstanceId.make("claudeAgent"), @@ -1195,6 +1203,56 @@ describe("ProviderCommandReactor", () => { }); }); + it("falls back to the project checkout when the thread worktree path is gone", async () => { + const harness = await createHarness(); + const now = "2026-01-01T00:00:00.000Z"; + const staleWorktreePath = NodeFS.mkdtempSync( + NodePath.join(NodeOS.tmpdir(), "t3-provider-stale-worktree-"), + ); + + await Effect.runPromise( + harness.engine.dispatch({ + type: "thread.meta.update", + commandId: CommandId.make("cmd-thread-stale-worktree"), + threadId: ThreadId.make("thread-1"), + worktreePath: staleWorktreePath, + }), + ); + NodeFS.rmSync(staleWorktreePath, { recursive: true, force: true }); + + await Effect.runPromise( + harness.engine.dispatch({ + type: "thread.turn.start", + commandId: CommandId.make("cmd-turn-start-stale-worktree"), + threadId: ThreadId.make("thread-1"), + message: { + messageId: asMessageId("user-message-stale-worktree"), + role: "user", + text: "continue after deleted worktree", + attachments: [], + }, + interactionMode: DEFAULT_PROVIDER_INTERACTION_MODE, + runtimeMode: "approval-required", + createdAt: now, + }), + ); + + await waitFor(() => harness.startSession.mock.calls.length === 1); + await waitFor(() => harness.sendTurn.mock.calls.length === 1); + + const readModel = await harness.readModel(); + const thread = readModel.threads.find((entry) => entry.id === ThreadId.make("thread-1")); + expect(thread?.branch).toBeNull(); + expect(thread?.worktreePath).toBeNull(); + expect(harness.startSession.mock.calls[0]?.[1]).toMatchObject({ + threadId: ThreadId.make("thread-1"), + cwd: "/tmp/provider-project", + }); + expect( + thread?.activities.find((activity) => activity.kind === "provider.turn.start.failed"), + ).toBeUndefined(); + }); + it("restarts claude sessions when claude effort changes", async () => { const harness = await createHarness({ threadModelSelection: { @@ -2053,45 +2111,45 @@ describe("ProviderCommandReactor", () => { expect(resolvedActivity).toBeUndefined(); }); - it("reacts to thread.session.stop by stopping provider session and clearing thread session state", async () => { - const harness = await createHarness(); - const now = "2026-01-01T00:00:00.000Z"; + effectIt( + "reacts to thread.session.stop by stopping provider session and clearing thread session state", + () => + Effect.gen(function* () { + const harness = yield* Effect.promise(() => createHarness()); + const now = "2026-01-01T00:00:00.000Z"; - await Effect.runPromise( - harness.engine.dispatch({ - type: "thread.session.set", - commandId: CommandId.make("cmd-session-set-for-stop"), - threadId: ThreadId.make("thread-1"), - session: { + yield* harness.engine.dispatch({ + type: "thread.session.set", + commandId: CommandId.make("cmd-session-set-for-stop"), threadId: ThreadId.make("thread-1"), - status: "ready", - providerName: "codex", - providerInstanceId: ProviderInstanceId.make("codex_work"), - runtimeMode: "approval-required", - activeTurnId: null, - lastError: null, - updatedAt: now, - }, - createdAt: now, - }), - ); + session: { + threadId: ThreadId.make("thread-1"), + status: "ready", + providerName: "codex", + providerInstanceId: ProviderInstanceId.make("codex_work"), + runtimeMode: "approval-required", + activeTurnId: null, + lastError: null, + updatedAt: now, + }, + createdAt: now, + }); - await Effect.runPromise( - harness.engine.dispatch({ - type: "thread.session.stop", - commandId: CommandId.make("cmd-session-stop"), - threadId: ThreadId.make("thread-1"), - createdAt: now, - }), - ); + yield* harness.engine.dispatch({ + type: "thread.session.stop", + commandId: CommandId.make("cmd-session-stop"), + threadId: ThreadId.make("thread-1"), + createdAt: now, + }); - await waitFor(() => harness.stopSession.mock.calls.length === 1); - const readModel = await harness.readModel(); - const thread = readModel.threads.find((entry) => entry.id === ThreadId.make("thread-1")); - expect(thread?.session).not.toBeNull(); - expect(thread?.session?.status).toBe("stopped"); - expect(thread?.session?.threadId).toBe("thread-1"); - expect(thread?.session?.providerInstanceId).toBe(ProviderInstanceId.make("codex_work")); - expect(thread?.session?.activeTurnId).toBeNull(); - }); + yield* Effect.promise(() => waitFor(() => harness.stopSession.mock.calls.length === 1)); + const readModel = yield* Effect.promise(() => harness.readModel()); + const thread = readModel.threads.find((entry) => entry.id === ThreadId.make("thread-1")); + expect(thread?.session).not.toBeNull(); + expect(thread?.session?.status).toBe("stopped"); + expect(thread?.session?.threadId).toBe("thread-1"); + expect(thread?.session?.providerInstanceId).toBe(ProviderInstanceId.make("codex_work")); + expect(thread?.session?.activeTurnId).toBeNull(); + }), + ); }); diff --git a/apps/server/src/orchestration/Layers/ProviderCommandReactor.ts b/apps/server/src/orchestration/Layers/ProviderCommandReactor.ts index 9c7a7c94bb10..8ae1084bd60a 100644 --- a/apps/server/src/orchestration/Layers/ProviderCommandReactor.ts +++ b/apps/server/src/orchestration/Layers/ProviderCommandReactor.ts @@ -19,6 +19,7 @@ import * as Crypto from "effect/Crypto"; import * as Duration from "effect/Duration"; import * as Effect from "effect/Effect"; import * as Equal from "effect/Equal"; +import * as FileSystem from "effect/FileSystem"; import * as Layer from "effect/Layer"; import * as Option from "effect/Option"; import * as Schema from "effect/Schema"; @@ -196,6 +197,7 @@ const make = Effect.gen(function* () { const vcsStatusBroadcaster = yield* VcsStatusBroadcaster; const textGeneration = yield* TextGeneration; const serverSettingsService = yield* ServerSettingsService; + const fileSystem = yield* FileSystem.FileSystem; const serverCommandId = (tag: string) => crypto.randomUUIDv4.pipe(Effect.map((uuid) => CommandId.make(`server:${tag}:${uuid}`))); const serverEventId = () => crypto.randomUUIDv4.pipe(Effect.map(EventId.make)); @@ -349,6 +351,65 @@ const make = Effect.gen(function* () { }); }); + const resolveAvailableThreadWorkspaceCwd = Effect.fnUntraced(function* (input: { + readonly thread: { + readonly id: ThreadId; + readonly projectId: ProjectId; + readonly worktreePath: string | null; + readonly modelSelection: ModelSelection; + readonly session: OrchestrationSession | null; + }; + readonly projects: ReadonlyArray<{ + readonly id: ProjectId; + readonly workspaceRoot: string; + }>; + }) { + const fallbackCwd = resolveThreadWorkspaceCwd({ + thread: { + projectId: input.thread.projectId, + worktreePath: null, + }, + projects: input.projects, + }); + + if (!input.thread.worktreePath) { + return fallbackCwd; + } + + const stat = yield* fileSystem + .stat(input.thread.worktreePath) + .pipe(Effect.orElseSucceed(() => null)); + if (stat?.type === "Directory") { + return input.thread.worktreePath; + } + + if (!fallbackCwd) { + return yield* new ProviderAdapterRequestError({ + provider: providerErrorLabelFromInstanceHint({ + modelSelectionInstanceId: String(input.thread.modelSelection.instanceId), + sessionProvider: input.thread.session?.providerName ?? undefined, + }), + method: "thread.turn.start", + detail: `Thread '${input.thread.id}' worktree path is unavailable and the project checkout could not be resolved: ${input.thread.worktreePath}.`, + }); + } + + yield* Effect.logWarning("provider command reactor falling back from stale worktree path", { + threadId: input.thread.id, + worktreePath: input.thread.worktreePath, + fallbackCwd, + reason: stat === null ? "missing" : `not-${stat.type}`, + }); + yield* orchestrationEngine.dispatch({ + type: "thread.meta.update", + commandId: yield* serverCommandId("stale-worktree-fallback"), + threadId: input.thread.id, + branch: null, + worktreePath: null, + }); + return fallbackCwd; + }); + const ensureSessionForThread = Effect.fn("ensureSessionForThread")(function* ( threadId: ThreadId, createdAt: string, @@ -466,7 +527,7 @@ const make = Effect.gen(function* () { } } const project = yield* resolveProject(thread.projectId); - const effectiveCwd = resolveThreadWorkspaceCwd({ + const effectiveCwd = yield* resolveAvailableThreadWorkspaceCwd({ thread, projects: project ? [project] : [], }); @@ -775,10 +836,12 @@ const make = Effect.gen(function* () { if (isFirstUserMessageTurn) { const project = yield* resolveProject(thread.projectId); const generationCwd = - resolveThreadWorkspaceCwd({ + (yield* resolveAvailableThreadWorkspaceCwd({ thread, projects: project ? [project] : [], - }) ?? process.cwd(); + })) ?? process.cwd(); + const generationWorktreePath = + thread.worktreePath && generationCwd === thread.worktreePath ? thread.worktreePath : null; const generationInput = { messageText: message.text, ...(message.attachments !== undefined ? { attachments: message.attachments } : {}), @@ -787,8 +850,8 @@ const make = Effect.gen(function* () { yield* maybeGenerateAndRenameWorktreeBranchForFirstTurn({ threadId: event.payload.threadId, - branch: thread.branch, - worktreePath: thread.worktreePath, + branch: generationWorktreePath ? thread.branch : null, + worktreePath: generationWorktreePath, ...generationInput, }).pipe(Effect.forkScoped); diff --git a/apps/server/src/provider/Layers/ProviderService.test.ts b/apps/server/src/provider/Layers/ProviderService.test.ts index ccbbce1759f0..f2334a64e56b 100644 --- a/apps/server/src/provider/Layers/ProviderService.test.ts +++ b/apps/server/src/provider/Layers/ProviderService.test.ts @@ -33,6 +33,7 @@ import * as Ref from "effect/Ref"; import * as Scope from "effect/Scope"; import * as Stream from "effect/Stream"; import * as TestClock from "effect/testing/TestClock"; +import * as Cause from "effect/Cause"; import * as SqlClient from "effect/unstable/sql/SqlClient"; import { @@ -299,6 +300,7 @@ function makeProviderServiceLayer() { ProviderEventLoggers.NoOpProviderEventLoggers, ), ), + Layer.provideMerge(NodeServices.layer), ), directoryLayer, @@ -350,6 +352,7 @@ it.effect("ProviderServiceLive catches stopAll failures during shutdown", () => ProviderEventLoggers.NoOpProviderEventLoggers, ), ), + Layer.provideMerge(NodeServices.layer), ), directoryLayer, runtimeRepositoryLayer, @@ -717,6 +720,8 @@ it.effect( const tempDir = NodeFS.mkdtempSync( NodePath.join(NodeOS.tmpdir(), "t3-provider-service-restart-"), ); + const persistedCwd = NodePath.join(tempDir, "project"); + NodeFS.mkdirSync(persistedCwd, { recursive: true }); const dbPath = NodePath.join(tempDir, "orchestration.sqlite"); const persistenceLayer = makeSqlitePersistenceLive(dbPath); const runtimeRepositoryLayer = ProviderSessionRuntime.layer.pipe( @@ -758,7 +763,7 @@ it.effect( const session = yield* provider.startSession(threadId, { provider: ProviderDriverKind.make("codex"), providerInstanceId: codexInstanceId, - cwd: "/tmp/project", + cwd: persistedCwd, runtimeMode: "full-access", threadId, }); @@ -827,7 +832,7 @@ it.effect( threadId?: string; }; assert.equal(startPayload.provider, "codex"); - assert.equal(startPayload.cwd, "/tmp/project"); + assert.equal(startPayload.cwd, persistedCwd); assert.deepEqual(startPayload.resumeCursor, updatedResumeCursor); assert.equal(startPayload.threadId, startedSession.threadId); } @@ -844,12 +849,16 @@ routing.layer("ProviderServiceLive routing", (it) => { it.effect("routes provider operations and rollback conversation", () => Effect.gen(function* () { const provider = yield* ProviderService.ProviderService; + const threadId = asThreadId("thread-routing-ops"); + const persistedCwd = NodeFS.mkdtempSync( + NodePath.join(NodeOS.tmpdir(), "t3-provider-routing-cwd-"), + ); - const session = yield* provider.startSession(asThreadId("thread-1"), { + const session = yield* provider.startSession(threadId, { provider: ProviderDriverKind.make("codex"), providerInstanceId: codexInstanceId, - threadId: asThreadId("thread-1"), - cwd: "/tmp/project", + threadId, + cwd: persistedCwd, runtimeMode: "full-access", }); assert.equal(session.provider, "codex"); @@ -919,23 +928,28 @@ routing.layer("ProviderServiceLive routing", (it) => { threadId?: string; }; assert.equal(startPayload.provider, "codex"); - assert.equal(startPayload.cwd, "/tmp/project"); + assert.equal(startPayload.cwd, persistedCwd); assert.deepEqual(startPayload.resumeCursor, session.resumeCursor); assert.equal(startPayload.threadId, session.threadId); } assert.equal(routing.codex.sendTurn.mock.calls.length, 1); + NodeFS.rmSync(persistedCwd, { recursive: true, force: true }); }), ); it.effect("recovers stale persisted sessions for rollback by resuming thread identity", () => Effect.gen(function* () { const provider = yield* ProviderService.ProviderService; + const threadId = asThreadId("thread-rollback-recovery"); + const persistedCwd = NodeFS.mkdtempSync( + NodePath.join(NodeOS.tmpdir(), "t3-provider-rollback-cwd-"), + ); - const initial = yield* provider.startSession(asThreadId("thread-1"), { + const initial = yield* provider.startSession(threadId, { provider: ProviderDriverKind.make("codex"), providerInstanceId: codexInstanceId, - threadId: asThreadId("thread-1"), - cwd: "/tmp/project", + threadId, + cwd: persistedCwd, runtimeMode: "full-access", }); yield* routing.codex.stopSession(initial.threadId); @@ -958,13 +972,14 @@ routing.layer("ProviderServiceLive routing", (it) => { threadId?: string; }; assert.equal(startPayload.provider, "codex"); - assert.equal(startPayload.cwd, "/tmp/project"); + assert.equal(startPayload.cwd, persistedCwd); assert.deepEqual(startPayload.resumeCursor, initial.resumeCursor); assert.equal(startPayload.threadId, initial.threadId); } assert.equal(routing.codex.rollbackThread.mock.calls.length, 1); const rollbackCall = routing.codex.rollbackThread.mock.calls[0]; assert.equal(rollbackCall?.[1], 1); + NodeFS.rmSync(persistedCwd, { recursive: true, force: true }); }), ); @@ -972,12 +987,15 @@ routing.layer("ProviderServiceLive routing", (it) => { Effect.gen(function* () { const provider = yield* ProviderService.ProviderService; const runtimeRepository = yield* ProviderSessionRuntime.ProviderSessionRuntimeRepository; + const persistedCwd = NodeFS.mkdtempSync( + NodePath.join(NodeOS.tmpdir(), "t3-provider-reap-cwd-"), + ); const initial = yield* provider.startSession(asThreadId("thread-reap-preserve"), { provider: ProviderDriverKind.make("codex"), providerInstanceId: codexInstanceId, threadId: asThreadId("thread-reap-preserve"), - cwd: "/tmp/project-reap-preserve", + cwd: persistedCwd, runtimeMode: "full-access", }); @@ -1012,11 +1030,12 @@ routing.layer("ProviderServiceLive routing", (it) => { threadId?: string; }; assert.equal(startPayload.provider, "codex"); - assert.equal(startPayload.cwd, "/tmp/project-reap-preserve"); + assert.equal(startPayload.cwd, persistedCwd); assert.deepEqual(startPayload.resumeCursor, initial.resumeCursor); assert.equal(startPayload.threadId, initial.threadId); } assert.equal(routing.codex.sendTurn.mock.calls.length, 1); + NodeFS.rmSync(persistedCwd, { recursive: true, force: true }); }), ); @@ -1122,12 +1141,16 @@ routing.layer("ProviderServiceLive routing", (it) => { it.effect("recovers stale sessions for sendTurn using persisted cwd", () => Effect.gen(function* () { const provider = yield* ProviderService.ProviderService; + const persistedCwd = NodeFS.mkdtempSync( + NodePath.join(NodeOS.tmpdir(), "t3-provider-send-turn-cwd-"), + ); + const threadId = asThreadId("thread-send-turn-cwd-recovery"); - const initial = yield* provider.startSession(asThreadId("thread-1"), { + const initial = yield* provider.startSession(threadId, { provider: ProviderDriverKind.make("codex"), providerInstanceId: codexInstanceId, - threadId: asThreadId("thread-1"), - cwd: "/tmp/project-send-turn", + threadId, + cwd: persistedCwd, runtimeMode: "full-access", }); @@ -1152,23 +1175,66 @@ routing.layer("ProviderServiceLive routing", (it) => { threadId?: string; }; assert.equal(startPayload.provider, "codex"); - assert.equal(startPayload.cwd, "/tmp/project-send-turn"); + assert.equal(startPayload.cwd, persistedCwd); assert.deepEqual(startPayload.resumeCursor, initial.resumeCursor); assert.equal(startPayload.threadId, initial.threadId); } assert.equal(routing.codex.sendTurn.mock.calls.length, 1); + NodeFS.rmSync(persistedCwd, { recursive: true, force: true }); + }), + ); + + it.effect("does not reuse a deleted persisted cwd when recovering sendTurn", () => + Effect.gen(function* () { + const provider = yield* ProviderService.ProviderService; + const staleCwd = NodeFS.mkdtempSync(NodePath.join(NodeOS.tmpdir(), "t3-provider-stale-cwd-")); + const threadId = asThreadId("thread-stale-cwd-send-turn"); + + const initial = yield* provider.startSession(threadId, { + provider: ProviderDriverKind.make("codex"), + providerInstanceId: codexInstanceId, + threadId, + cwd: staleCwd, + runtimeMode: "full-access", + }); + + NodeFS.rmSync(staleCwd, { recursive: true, force: true }); + yield* routing.codex.stopAll(); + routing.codex.startSession.mockClear(); + routing.codex.sendTurn.mockClear(); + + const exit = yield* Effect.exit( + provider.sendTurn({ + threadId: initial.threadId, + input: "resume", + attachments: [], + }), + ); + + assert.equal(Exit.isFailure(exit), true); + if (Exit.isFailure(exit)) { + const error = Cause.squash(exit.cause); + assert.instanceOf(error, ProviderValidationError); + assert.include(error.message, staleCwd); + assert.include(error.message, "no longer exists"); + } + assert.equal(routing.codex.startSession.mock.calls.length, 0); + assert.equal(routing.codex.sendTurn.mock.calls.length, 0); }), ); it.effect("recovers stale claudeAgent sessions for sendTurn using persisted cwd", () => Effect.gen(function* () { const provider = yield* ProviderService.ProviderService; + const persistedCwd = NodeFS.mkdtempSync( + NodePath.join(NodeOS.tmpdir(), "t3-provider-claude-send-turn-cwd-"), + ); const initial = yield* provider.startSession(asThreadId("thread-claude-send-turn"), { provider: ProviderDriverKind.make("claudeAgent"), providerInstanceId: claudeAgentInstanceId, threadId: asThreadId("thread-claude-send-turn"), - cwd: "/tmp/project-claude-send-turn", + cwd: persistedCwd, modelSelection: createModelSelection( ProviderInstanceId.make("claudeAgent"), "claude-opus-4-6", @@ -1199,7 +1265,7 @@ routing.layer("ProviderServiceLive routing", (it) => { threadId?: string; }; assert.equal(startPayload.provider, "claudeAgent"); - assert.equal(startPayload.cwd, "/tmp/project-claude-send-turn"); + assert.equal(startPayload.cwd, persistedCwd); assert.deepEqual( startPayload.modelSelection, createModelSelection(ProviderInstanceId.make("claudeAgent"), "claude-opus-4-6", [ @@ -1210,6 +1276,7 @@ routing.layer("ProviderServiceLive routing", (it) => { assert.equal(startPayload.threadId, initial.threadId); } assert.equal(routing.claude.sendTurn.mock.calls.length, 1); + NodeFS.rmSync(persistedCwd, { recursive: true, force: true }); }), ); @@ -1394,6 +1461,8 @@ routing.layer("ProviderServiceLive routing", (it) => { const tempDir = NodeFS.mkdtempSync( NodePath.join(NodeOS.tmpdir(), "t3-provider-service-cwd-"), ); + const persistedCwd = NodePath.join(tempDir, "project-claude-cwd"); + NodeFS.mkdirSync(persistedCwd, { recursive: true }); const dbPath = NodePath.join(tempDir, "orchestration.sqlite"); const persistenceLayer = makeSqlitePersistenceLive(dbPath); const runtimeRepositoryLayer = ProviderSessionRuntime.layer.pipe( @@ -1428,7 +1497,7 @@ routing.layer("ProviderServiceLive routing", (it) => { provider: ProviderDriverKind.make("claudeAgent"), providerInstanceId: claudeAgentInstanceId, threadId: asThreadId("thread-claude-cwd"), - cwd: "/tmp/project-claude-cwd", + cwd: persistedCwd, runtimeMode: "full-access", }); }).pipe(Effect.provide(firstProviderLayer)); @@ -1478,7 +1547,7 @@ routing.layer("ProviderServiceLive routing", (it) => { threadId?: string; }; assert.equal(startPayload.provider, "claudeAgent"); - assert.equal(startPayload.cwd, "/tmp/project-claude-cwd"); + assert.equal(startPayload.cwd, persistedCwd); assert.deepEqual(startPayload.resumeCursor, initial.resumeCursor); assert.equal(startPayload.threadId, initial.threadId); } diff --git a/apps/server/src/provider/Layers/ProviderService.ts b/apps/server/src/provider/Layers/ProviderService.ts index 2eaaeb8ce3c0..3ba6f761e27f 100644 --- a/apps/server/src/provider/Layers/ProviderService.ts +++ b/apps/server/src/provider/Layers/ProviderService.ts @@ -27,6 +27,7 @@ import { import { causeErrorTag } from "@t3tools/shared/observability"; import * as DateTime from "effect/DateTime"; import * as Effect from "effect/Effect"; +import * as FileSystem from "effect/FileSystem"; import * as Layer from "effect/Layer"; import * as Option from "effect/Option"; import * as PubSub from "effect/PubSub"; @@ -212,6 +213,7 @@ const makeProviderService = Effect.fn("makeProviderService")(function* ( const registry = yield* ProviderAdapterRegistry.ProviderAdapterRegistry; const directory = yield* ProviderSessionDirectory.ProviderSessionDirectory; + const fileSystem = yield* FileSystem.FileSystem; const runtimeEventPubSub = yield* PubSub.unbounded(); const nowIso = Effect.map(DateTime.now, DateTime.formatIso); const prepareMcpSession = (threadId: ThreadId, providerInstanceId: ProviderInstanceId) => @@ -297,6 +299,29 @@ const makeProviderService = Effect.fn("makeProviderService")(function* ( ), ); + const readExistingPersistedCwd = Effect.fn("readExistingPersistedCwd")(function* (input: { + readonly operation: string; + readonly threadId: ThreadId; + readonly runtimePayload: ProviderSessionDirectory.ProviderRuntimeBinding["runtimePayload"]; + }) { + const persistedCwd = readPersistedCwd(input.runtimePayload); + if (!persistedCwd) { + return undefined; + } + + const stat = yield* fileSystem.stat(persistedCwd).pipe(Effect.orElseSucceed(() => null)); + if (stat?.type === "Directory") { + return persistedCwd; + } + + return yield* toValidationError( + input.operation, + stat === null + ? `Cannot recover thread '${input.threadId}' because its persisted working directory no longer exists: ${persistedCwd}. Restore the directory or start the thread with a valid checkout.` + : `Cannot recover thread '${input.threadId}' because its persisted working directory is not a directory: ${persistedCwd}. Start the thread with a valid checkout.`, + ); + }); + // `subscribedAdapters` is our source-of-truth for "which instance adapters // are currently wired into the runtime event bus". It both tracks the set // of live subscriptions (so `reconcileInstanceSubscriptions` can diff and @@ -394,7 +419,11 @@ const makeProviderService = Effect.fn("makeProviderService")(function* ( ); } - const persistedCwd = readPersistedCwd(input.binding.runtimePayload); + const persistedCwd = yield* readExistingPersistedCwd({ + operation: input.operation, + threadId: input.binding.threadId, + runtimePayload: input.binding.runtimePayload, + }); const persistedModelSelection = readPersistedModelSelection(input.binding.runtimePayload); yield* prepareMcpSession(input.binding.threadId, bindingInstanceId); @@ -568,7 +597,11 @@ const makeProviderService = Effect.fn("makeProviderService")(function* ( const effectiveCwd = input.cwd ?? (persistedBinding?.providerInstanceId === resolvedInstanceId - ? readPersistedCwd(persistedBinding.runtimePayload) + ? yield* readExistingPersistedCwd({ + operation: "ProviderService.startSession", + threadId, + runtimePayload: persistedBinding.runtimePayload, + }) : undefined); yield* Effect.annotateCurrentSpan({ "provider.kind": resolvedProvider,