diff --git a/apps/server/src/keybindings.ts b/apps/server/src/keybindings.ts index 7a1bcd7f356a..8c8090ab09bc 100644 --- a/apps/server/src/keybindings.ts +++ b/apps/server/src/keybindings.ts @@ -40,6 +40,7 @@ import * as Context from "effect/Context"; import * as Scope from "effect/Scope"; import * as Stream from "effect/Stream"; import * as Semaphore from "effect/Semaphore"; +import { keepWatching } from "./restartingWatch.ts"; import * as ServerConfig from "./config.ts"; import { writeFileStringAtomically } from "./atomicWrite.ts"; import { fromJsonStringPretty, fromLenientJson } from "@t3tools/shared/schemaJson"; @@ -577,11 +578,11 @@ const make = Effect.gen(function* () { Stream.debounce(Duration.millis(100)), ); - yield* Stream.runForEach(debouncedKeybindingsEvents, () => revalidateAndEmitSafely).pipe( - Effect.ignoreCause({ log: true }), - Effect.forkIn(watcherScope), - Effect.asVoid, - ); + yield* keepWatching({ + label: "Keybindings", + watch: Stream.runForEach(debouncedKeybindingsEvents, () => revalidateAndEmitSafely), + revalidate: revalidateAndEmitSafely, + }).pipe(Effect.forkIn(watcherScope), Effect.asVoid); }); const start = Effect.gen(function* () { diff --git a/apps/server/src/mcp/OrchestratorMcpService.test.ts b/apps/server/src/mcp/OrchestratorMcpService.test.ts index 3f927c283ef8..34332f102579 100644 --- a/apps/server/src/mcp/OrchestratorMcpService.test.ts +++ b/apps/server/src/mcp/OrchestratorMcpService.test.ts @@ -2,13 +2,18 @@ import * as NodeServices from "@effect/platform-node/NodeServices"; import { assert, describe, it } from "@effect/vitest"; import { EnvironmentId, + IsoDateTime, NodeId, ProjectId, ProviderDriverKind, ProviderInstanceId, RunId, + ScheduledTaskId, ThreadId, type OrchestrationV2ThreadProjection, + type OrchestrationV2ThreadShell, + type RuntimeMode, + type ScheduledTask, type ServerProvider, } from "@t3tools/contracts"; import * as Effect from "effect/Effect"; @@ -1037,3 +1042,124 @@ describe("OrchestratorMcpService provider resolution", () => { }), ); }); + +describe("OrchestratorMcpService scheduled tasks", () => { + it.effect("lets a restricted agent edit or delete only tasks with no more access", () => + Effect.gen(function* () { + const threadId = ThreadId.make("thread:mcp-scheduled-caller"); + const projectId = ProjectId.make("project:mcp-scheduled"); + const timestamp = IsoDateTime.make("2026-07-01T09:00:00.000Z"); + const fullAccessThreadId = ThreadId.make("thread:mcp-scheduled-full-access"); + const makeTask = ( + id: string, + runtimeMode: RuntimeMode, + threadId: ThreadId | null = null, + ): ScheduledTask => ({ + id: ScheduledTaskId.make(id), + title: id, + prompt: "check the build", + enabled: false, + schedule: { type: "interval", everyMs: 60_000 }, + projectId, + threadId, + workspaceStrategy: { type: "worktree", baseRef: "main", startFromOrigin: true }, + modelSelection: { instanceId: ProviderInstanceId.make("codex"), model: "gpt-5.6-terra" }, + runtimeMode, + interactionMode: "default", + createdBy: "user", + creationSource: "web", + createdAt: timestamp, + updatedAt: timestamp, + nextRunAt: null, + lastRunAt: null, + lastRunStatus: "never", + lastRunError: null, + runCount: 0, + }); + const fullAccessTask = makeTask("scheduled-task:full-access", "full-access"); + const restrictedTask = makeTask("scheduled-task:restricted", "approval-required"); + // Saved as restricted, but it runs in a full-access thread. + const boundTask = makeTask("scheduled-task:bound", "approval-required", fullAccessThreadId); + const writes = yield* Ref.make>([]); + const caller = { + thread: { + id: threadId, + projectId, + archivedAt: null, + runtimeMode: "approval-required", + interactionMode: "default", + }, + } as unknown as OrchestrationV2ThreadProjection; + const dependencies = Layer.mergeAll( + NodeServices.layer, + Layer.mock(ThreadManagementService.ThreadManagementService)({ + getThreadRecords: () => Effect.succeed(caller), + getThreadShell: (threadId) => + Effect.succeed( + threadId === fullAccessThreadId + ? ({ + runtimeMode: "full-access", + interactionMode: "default", + } as OrchestrationV2ThreadShell) + : null, + ), + }), + Layer.mock(ProviderRegistry.ProviderRegistry)({ getProviders: Effect.succeed([]) }), + Layer.mock(ProviderAdapterRegistry.ProviderAdapterRegistryV2)({ + list: () => Effect.succeed([]), + }), + Layer.mock(ScheduledTaskService.ScheduledTaskService)({ + list: () => Effect.succeed({ tasks: [fullAccessTask, restrictedTask, boundTask] }), + upsert: (input) => + Ref.update(writes, (all) => [...all, `upsert:${input.id}`]).pipe( + Effect.as({ task: { ...restrictedTask, enabled: input.enabled } }), + ), + delete: (input) => + Ref.update(writes, (all) => [...all, `delete:${input.id}`]).pipe( + Effect.as({ id: input.id }), + ), + }), + ); + const scope: McpInvocationScope = { + environmentId: EnvironmentId.make("environment:mcp-scheduled"), + threadId, + providerSessionId: "provider-session:mcp-scheduled", + providerInstanceId: ProviderInstanceId.make("codex"), + capabilities: new Set(["orchestration"]), + issuedAt: 1, + }; + + yield* Effect.gen(function* () { + const service = yield* OrchestratorMcpService.OrchestratorMcpService; + const updateError = yield* service + .updateScheduledTask(scope, { + scheduledTaskId: fullAccessTask.id, + prompt: "push to main", + enabled: true, + }) + .pipe(Effect.flip); + assert.equal(updateError.code, "capability_denied"); + const deleteError = yield* service + .deleteScheduledTask(scope, { scheduledTaskId: fullAccessTask.id }) + .pipe(Effect.flip); + assert.equal(deleteError.code, "capability_denied"); + const boundError = yield* service + .updateScheduledTask(scope, { scheduledTaskId: boundTask.id, prompt: "push to main" }) + .pipe(Effect.flip); + assert.equal(boundError.code, "capability_denied"); + assert.deepStrictEqual(yield* Ref.get(writes), []); + + const updated = yield* service.updateScheduledTask(scope, { + scheduledTaskId: restrictedTask.id, + enabled: true, + }); + assert.isTrue(updated.enabled); + yield* service.deleteScheduledTask(scope, { scheduledTaskId: restrictedTask.id }); + assert.deepStrictEqual(yield* Ref.get(writes), [ + `upsert:${restrictedTask.id}`, + `delete:${restrictedTask.id}`, + ]); + }).pipe(Effect.provide(OrchestratorMcpService.layer.pipe(Layer.provide(dependencies)))); + }), + ); +}); diff --git a/apps/server/src/mcp/OrchestratorMcpService.ts b/apps/server/src/mcp/OrchestratorMcpService.ts index 5fdc7eb81c80..56c5be078dfb 100644 --- a/apps/server/src/mcp/OrchestratorMcpService.ts +++ b/apps/server/src/mcp/OrchestratorMcpService.ts @@ -1184,6 +1184,44 @@ const make = Effect.gen(function* () { return task; }); + // Agents may edit or delete only tasks that run with no more access than + // their own live thread. Otherwise a restricted agent could rewrite a task + // that the scheduler runs with full access. An unbound task runs with its + // saved modes; a bound one runs in its thread with that thread's modes. + const requireScheduledTaskAuthority = ( + caller: OrchestrationV2ThreadProjection["thread"], + task: ScheduledTask, + ) => + Effect.gen(function* () { + const boundThread = + task.threadId === null || task.threadId === caller.id + ? null + : yield* threadManagement + .getThreadShell(task.threadId) + .pipe( + Effect.mapError((error) => + failure( + "orchestration_error", + `Unable to read thread ${task.threadId}: ${errorMessage(error)}`, + ), + ), + ); + const allowed = + caller.archivedAt === null && + [task, ...(boundThread === null ? [] : [boundThread])].every( + (required) => + runtimeModeRank(caller.runtimeMode) >= runtimeModeRank(required.runtimeMode) && + interactionModeRank(caller.interactionMode) >= + interactionModeRank(required.interactionMode), + ); + if (!allowed) { + return yield* failure( + "capability_denied", + `Scheduled task ${task.id} needs a live calling thread with at least the modes it runs with.`, + ); + } + }); + return OrchestratorMcpService.of({ scheduleTask: (scope, input) => Effect.gen(function* () { @@ -1253,6 +1291,7 @@ const make = Effect.gen(function* () { parent.thread.projectId, input.scheduledTaskId, ); + yield* requireScheduledTaskAuthority(parent.thread, existing); const threadId = input.bindToCurrentThread === undefined ? existing.threadId @@ -1298,6 +1337,7 @@ const make = Effect.gen(function* () { parent.thread.projectId, input.scheduledTaskId, ); + yield* requireScheduledTaskAuthority(parent.thread, existing); yield* scheduledTasks .delete({ id: existing.id }) .pipe( diff --git a/apps/server/src/mcp/toolkits/orchestrator/tools.ts b/apps/server/src/mcp/toolkits/orchestrator/tools.ts index 833e88bd2ae8..204aead3a628 100644 --- a/apps/server/src/mcp/toolkits/orchestrator/tools.ts +++ b/apps/server/src/mcp/toolkits/orchestrator/tools.ts @@ -122,7 +122,7 @@ const ListScheduledTasksTool = Tool.make("list_scheduled_tasks", { const UpdateScheduledTaskTool = Tool.make("update_scheduled_task", { description: - "Update an existing scheduled task by scheduledTaskId (from list_scheduled_tasks). Only the provided fields change; omit a field to leave it as-is. Use enabled=false to pause a task without deleting it. Set bindToCurrentThread to move the task between posting into this thread and launching a fresh thread per run.", + "Update an existing scheduled task by scheduledTaskId (from list_scheduled_tasks). Only the provided fields change; omit a field to leave it as-is. Use enabled=false to pause a task without deleting it. Set bindToCurrentThread to move the task between posting into this thread and launching a fresh thread per run. Requires a live caller with at least the task's runtime and interaction modes.", parameters: OrchestratorMcpUpdateScheduledTaskInput, success: OrchestratorMcpScheduleTaskResult, failure: OrchestratorMcpFailure, @@ -134,7 +134,7 @@ const UpdateScheduledTaskTool = Tool.make("update_scheduled_task", { const DeleteScheduledTaskTool = Tool.make("delete_scheduled_task", { description: - "Permanently delete a scheduled task by scheduledTaskId (from list_scheduled_tasks). The task stops running immediately. To keep it but stop runs, use update_scheduled_task with enabled=false instead.", + "Permanently delete a scheduled task by scheduledTaskId (from list_scheduled_tasks). The task stops running immediately. To keep it but stop runs, use update_scheduled_task with enabled=false instead. Requires a live caller with at least the task's runtime and interaction modes.", parameters: OrchestratorMcpDeleteScheduledTaskInput, success: OrchestratorMcpDeleteScheduledTaskResult, failure: OrchestratorMcpFailure, diff --git a/apps/server/src/orchestration-v2/Adapters/AcpAdapterV2.test.ts b/apps/server/src/orchestration-v2/Adapters/AcpAdapterV2.test.ts index 146924586763..31dcdc69764b 100644 --- a/apps/server/src/orchestration-v2/Adapters/AcpAdapterV2.test.ts +++ b/apps/server/src/orchestration-v2/Adapters/AcpAdapterV2.test.ts @@ -3212,6 +3212,113 @@ describe("AcpAdapterV2", () => { }).pipe(Effect.provide(testLayer), Effect.scoped), ); + it.effect("activates a saved session with session/resume when the flavor prefers it", () => + Effect.gen(function* () { + const childProcessSpawner = yield* ChildProcessSpawner.ChildProcessSpawner; + const path = yield* Path.Path; + const mockAgentPath = yield* path.fromFileUrl( + new URL("../../../scripts/acp-mock-agent.ts", import.meta.url), + ); + const protocolEvents = yield* Queue.bounded(256); + const instanceId = ProviderInstanceId.make("acp-test"); + const adapter = makeAcpAdapterV2({ + crypto: yield* Crypto.Crypto, + instanceId, + flavor: { + driver: ACP_TEST_DRIVER, + capabilities: AcpProviderCapabilitiesV2, + preferResumeSession: true, + makeRuntime: makeMockRuntime({ childProcessSpawner, mockAgentPath, protocolEvents }), + }, + fileSystem: yield* FileSystem.FileSystem, + idAllocator: yield* IdAllocator.IdAllocatorV2, + serverConfig: yield* ServerConfig.ServerConfig, + selfInvocation: yield* resolveSelfInvocation(), + }); + const firstThreadId = ThreadId.make("thread-acp-prefer-resume:first"); + const secondThreadId = ThreadId.make("thread-acp-prefer-resume:second"); + const runtimePolicy = ProviderAdapterV2RuntimePolicy.make({ + runtimeMode: "full-access", + interactionMode: "default", + cwd: process.cwd(), + }); + const modelSelection = { instanceId, model: "default" } satisfies ModelSelection; + const runtime = yield* adapter.openSession({ + threadId: firstThreadId, + providerSessionId: ProviderSessionId.make("provider-session-acp-prefer-resume"), + modelSelection, + runtimePolicy, + }); + const firstProviderThread = yield* runtime.ensureThread({ + threadId: firstThreadId, + modelSelection, + runtimePolicy, + }); + const now = yield* DateTime.now; + yield* runtime.startTurn( + makeTurnInput({ + threadId: firstThreadId, + providerThread: firstProviderThread, + instanceId, + runtimePolicy, + modelSelection, + now, + }), + ); + yield* runtime.events.pipe( + Stream.filter((event) => event.type === "turn.terminal"), + Stream.runHead, + ); + + const secondProviderThread: OrchestrationV2ProviderThread = { + ...firstProviderThread, + id: ProviderThreadId.make("provider-thread-acp-prefer-resume:second"), + appThreadId: secondThreadId, + nativeThreadRef: { + driver: ACP_TEST_DRIVER, + nativeId: "mock-session-2", + strength: "strong", + }, + status: "idle", + }; + yield* runtime.resumeThread({ + providerThread: secondProviderThread, + modelSelection, + runtimePolicy, + }); + yield* runtime.startTurn( + makeTurnInput({ + threadId: secondThreadId, + providerThread: secondProviderThread, + instanceId, + runtimePolicy, + modelSelection, + now, + ordinal: 2, + }), + ); + yield* runtime.events.pipe( + Stream.filter((event) => event.type === "turn.terminal"), + Stream.runHead, + ); + + // ACP v2 sends a load as `session/resume` with `replayFrom`; a plain + // resume skips the slow history replay. + const activations = Array.from(yield* Queue.takeAll(protocolEvents)) + .filter( + (event) => + event.direction === "outgoing" && + (rawProtocolMethod(event) === "session/load" || + rawProtocolMethod(event) === "session/resume"), + ) + .map((event) => ({ + method: rawProtocolMethod(event), + replayFrom: rawProtocolRequestParam(event, "replayFrom"), + })); + assert.deepStrictEqual(activations, [{ method: "session/resume", replayFrom: undefined }]); + }).pipe(Effect.provide(testLayer), Effect.scoped), + ); + it.live("terminalizes an empty successful foreground Bash tool when the turn completes", () => Effect.gen(function* () { const childProcessSpawner = yield* ChildProcessSpawner.ChildProcessSpawner; diff --git a/apps/server/src/orchestration-v2/Adapters/AcpAdapterV2.ts b/apps/server/src/orchestration-v2/Adapters/AcpAdapterV2.ts index bf01d9cb2d7d..9cb984b9ac48 100644 --- a/apps/server/src/orchestration-v2/Adapters/AcpAdapterV2.ts +++ b/apps/server/src/orchestration-v2/Adapters/AcpAdapterV2.ts @@ -6094,14 +6094,15 @@ export function makeAcpAdapterV2( } const activationOptions = acpMcpActivation(threadId, self); prepareTerminalEnvironment(threadId, sessionId); - const activated = canLoadSession - ? yield* runtime.loadSession(sessionId, activationOptions) - : canResumeSession + const activated = + canResumeSession && (flavor.preferResumeSession === true || !canLoadSession) ? yield* runtime.resumeSession(sessionId, activationOptions) - : yield* new ProviderAdapter.ProviderAdapterProtocolError({ - driver, - detail: `ACP driver cannot load or resume session ${sessionId}`, - }); + : canLoadSession + ? yield* runtime.loadSession(sessionId, activationOptions) + : yield* new ProviderAdapter.ProviderAdapterProtocolError({ + driver, + detail: `ACP driver cannot load or resume session ${sessionId}`, + }); rememberTerminalEnvironment(activated.sessionId, threadId); return activated; }); diff --git a/apps/server/src/orchestration-v2/Adapters/ClaudeAdapterV2.test.ts b/apps/server/src/orchestration-v2/Adapters/ClaudeAdapterV2.test.ts index a4164c936fdf..40ec3dc9f272 100644 --- a/apps/server/src/orchestration-v2/Adapters/ClaudeAdapterV2.test.ts +++ b/apps/server/src/orchestration-v2/Adapters/ClaudeAdapterV2.test.ts @@ -7,6 +7,7 @@ import type { SDKUserMessage, } from "@anthropic-ai/claude-agent-sdk"; import type { AskUserQuestionInput } from "@anthropic-ai/claude-agent-sdk/sdk-tools"; +import * as NodeCrypto from "@effect/platform-node/NodeCrypto"; import * as NodeServices from "@effect/platform-node/NodeServices"; import { ChatAttachmentId, @@ -45,6 +46,8 @@ import * as Scope from "effect/Scope"; import * as Stream from "effect/Stream"; import { Tool } from "effect/unstable/ai"; import { formatClaudeResumeCompactionQuestion } from "@t3tools/shared/claudeCompaction"; +import { HostProcessPlatform } from "@t3tools/shared/hostProcess"; +import { SpawnExecutableResolution } from "@t3tools/shared/shell"; import { attachmentRelativePath } from "../../attachmentStore.ts"; import * as ServerConfig from "../../config.ts"; @@ -55,7 +58,9 @@ import { ProjectToolkit } from "../../mcp/toolkits/project/tools.ts"; import { WorktreeToolkit } from "../../mcp/toolkits/worktree/tools.ts"; import { ThreadToolkit } from "../../mcp/toolkits/thread/tools.ts"; import { OrchestratorToolkit } from "../../mcp/toolkits/orchestrator/tools.ts"; +import { ClaudeExecutableFileCheck } from "../../provider/Drivers/ClaudeExecutable.ts"; import type { EventNdjsonLogger } from "../../provider/Layers/EventNdjsonLogger.ts"; +import * as ProviderEventLoggers from "../../provider/Layers/ProviderEventLoggers.ts"; import { ProviderAdapterV2RuntimePolicy, type ProviderAdapterV2Event, @@ -861,6 +866,22 @@ describe("ClaudeAdapterV2 context usage", () => { updatedAt: "2026-08-29T00:00:00.000Z", }); }); + + it("reads the context window from the model catalog", () => { + const maxTokens = (model: string, options: ModelSelection["options"] = []) => + ClaudeAdapterV2.claudeProviderTurnTokenUsage( + { input_tokens: 1, output_tokens: 1 }, + { instanceId: CLAUDE_TEST_MODEL_SELECTION.instanceId, model, options }, + "2026-08-29T00:00:00.000Z", + ).maxTokens; + + // Opus 4.8 has a fixed 1M window in the catalog. + assert.equal(maxTokens("claude-opus-4-8"), 1_000_000); + // Opus 4.6 follows its 200k/1m selection. + assert.equal(maxTokens("claude-opus-4-6", [{ id: "contextWindow", value: "200k" }]), 200_000); + assert.equal(maxTokens("claude-opus-4-6", [{ id: "contextWindow", value: "1m" }]), 1_000_000); + assert.equal(maxTokens("custom-claude-model"), 200_000); + }); }); describe("ClaudeAdapterV2 session permissions", () => { @@ -1025,6 +1046,75 @@ describe("ClaudeAdapterV2 Auto-accept edits", () => { ); }); +describe("ClaudeAgentSdkQueryRunner forkSession", () => { + it.effect("forks a session stored under the instance CLAUDE_CONFIG_DIR", () => + Effect.scoped( + Effect.gen(function* () { + const fileSystem = yield* FileSystem.FileSystem; + const path = yield* Path.Path; + const configDir = yield* fileSystem.makeTempDirectoryScoped({ + prefix: "t3-claude-fork-config-", + }); + const projectDir = path.join(configDir, "projects", "-repo"); + const sessionId = "11111111-1111-4111-8111-111111111111"; + const entry = (uuid: string, parentUuid: string | null, message: unknown) => ({ + type: parentUuid === null ? "user" : "assistant", + uuid, + parentUuid, + sessionId, + cwd: "/repo", + isSidechain: false, + timestamp: "2026-10-01T00:00:00.000Z", + message, + }); + yield* fileSystem.makeDirectory(projectDir, { recursive: true }); + yield* fileSystem.writeFileString( + path.join(projectDir, `${sessionId}.jsonl`), + [ + entry("aaaaaaaa-0000-4000-8000-000000000001", null, { + role: "user", + content: "hello", + }), + entry("aaaaaaaa-0000-4000-8000-000000000002", "aaaaaaaa-0000-4000-8000-000000000001", { + role: "assistant", + content: [{ type: "text", text: "hi" }], + }), + ] + .map((line) => JSON.stringify(line)) + .join("\n"), + ); + + const runner = yield* ClaudeAdapterV2.ClaudeAgentSdkQueryRunner; + const forked = yield* runner.forkSession({ + sessionId, + options: {}, + environment: { ...process.env, CLAUDE_CONFIG_DIR: configDir }, + threadId: ThreadId.make("thread-claude-fork-config-dir"), + providerSessionId: ProviderSessionId.make("provider-session-claude-fork-config-dir"), + }); + + assert.notEqual(forked.sessionId, sessionId); + assert.isTrue(yield* fileSystem.exists(path.join(projectDir, `${forked.sessionId}.jsonl`))); + }), + ).pipe( + Effect.provide( + ClaudeAdapterV2.claudeAgentSdkQueryRunnerLiveLayer.pipe( + Layer.provideMerge( + Layer.mergeAll( + NodeServices.layer, + NodeCrypto.layer, + Layer.succeed( + ProviderEventLoggers.ProviderEventLoggers, + ProviderEventLoggers.NoOpProviderEventLoggers, + ), + ), + ), + ), + ), + ), + ); +}); + describe("ClaudeAdapterV2 approval cancellation", () => { it.effect("observes an approval signal that was already aborted", () => Effect.gen(function* () { @@ -1065,10 +1155,10 @@ describe("ClaudeAdapterV2 approval cancellation", () => { }); describe("ClaudeAdapterV2 executable path", () => { - it.effect("expands ~ in the configured binary path for the SDK", () => + // Opens one turn and returns the executable path the SDK query received. + const sdkExecutablePath = (binaryPath: string) => Effect.scoped( Effect.gen(function* () { - const path = yield* Path.Path; const executablePaths: Array = []; const adapter = yield* ClaudeAdapterV2.createClaudeAdapterV2( { @@ -1076,7 +1166,7 @@ describe("ClaudeAdapterV2 executable path", () => { displayName: undefined, environment: [], enabled: true, - config: { ...DEFAULT_CLAUDE_SETTINGS, binaryPath: "~/bin/claude" }, + config: { ...DEFAULT_CLAUDE_SETTINGS, binaryPath }, }, {}, ).pipe( @@ -1125,10 +1215,29 @@ describe("ClaudeAdapterV2 executable path", () => { attachments: [], }), ); - - assert.deepEqual(executablePaths, [path.join(NodeOS.homedir(), "bin", "claude")]); + return executablePaths; }), - ).pipe(Effect.provide(Layer.merge(IdAllocator.layer, NodeServices.layer))), + ).pipe(Effect.provide(Layer.merge(IdAllocator.layer, NodeServices.layer))); + + it.effect("expands ~ in the configured binary path for the SDK", () => + Effect.gen(function* () { + const path = yield* Path.Path; + assert.deepEqual(yield* sdkExecutablePath("~/bin/claude"), [ + path.join(NodeOS.homedir(), "bin", "claude"), + ]); + }).pipe(Effect.provide(NodeServices.layer)), + ); + + it.effect("follows a Windows npm launcher shim to the real executable", () => + Effect.gen(function* () { + const executable = "C:\\npm\\node_modules\\@anthropic-ai\\claude-code\\bin\\claude.exe"; + const executablePaths = yield* sdkExecutablePath("claude").pipe( + Effect.provideService(HostProcessPlatform, "win32"), + Effect.provideService(SpawnExecutableResolution, () => "C:\\npm\\claude.cmd"), + Effect.provideService(ClaudeExecutableFileCheck, (filePath) => filePath === executable), + ); + assert.deepEqual(executablePaths, [executable]); + }), ); }); @@ -1584,6 +1693,7 @@ describe("ClaudeAdapterV2 native fork", () => { const forkCalls: Array<{ readonly sessionId: string; readonly options: unknown; + readonly environment: NodeJS.ProcessEnv; readonly threadId: ThreadId; readonly providerSessionId: ProviderSessionId; }> = []; @@ -1677,6 +1787,7 @@ describe("ClaudeAdapterV2 native fork", () => { dir: "/workspace", upToMessageId: "assistant-message-cursor", }, + environment: {}, threadId: targetThreadId, providerSessionId, }, diff --git a/apps/server/src/orchestration-v2/Adapters/ClaudeAdapterV2.ts b/apps/server/src/orchestration-v2/Adapters/ClaudeAdapterV2.ts index 2cf91c619594..92222febda68 100644 --- a/apps/server/src/orchestration-v2/Adapters/ClaudeAdapterV2.ts +++ b/apps/server/src/orchestration-v2/Adapters/ClaudeAdapterV2.ts @@ -33,7 +33,7 @@ import type { WebSearchOutput, } from "@anthropic-ai/claude-agent-sdk/sdk-tools"; import { parseCliArgs } from "@t3tools/shared/cliArgs"; -import { HostProcessEnvironment } from "@t3tools/shared/hostProcess"; +import { HostProcessEnvironment, HostProcessIsExecutable } from "@t3tools/shared/hostProcess"; import { applyClaudePromptEffortPrefix } from "@t3tools/shared/model"; import { CLAUDE_RESUME_COMPACTION_NEVER_ANSWER, @@ -81,6 +81,7 @@ import * as FileSystem from "effect/FileSystem"; import * as Layer from "effect/Layer"; import * as Option from "effect/Option"; import * as Path from "effect/Path"; +import { ChildProcess, ChildProcessSpawner } from "effect/unstable/process"; import * as Queue from "effect/Queue"; import * as Ref from "effect/Ref"; import * as Schema from "effect/Schema"; @@ -92,13 +93,13 @@ import { discoverClaudeSkills } from "../../provider/Drivers/ClaudeSkills.ts"; import { compileClaudeModelSelection } from "../../claudeModelOptions.ts"; import * as ServerConfig from "../../config.ts"; import { expandHomePath } from "../../pathExpansion.ts"; +import { resolveClaudeSdkExecutablePath } from "../../provider/Drivers/ClaudeExecutable.ts"; import { claudeSignedOutMessage, makeClaudeEnvironment, } from "../../provider/Drivers/ClaudeHome.ts"; import { BUNDLED_CLAUDE_MODEL_CATALOG, - resolveClaudeCatalogContextWindow, resolveClaudeCatalogContextWindowTokens, } from "../../provider/ClaudeModelCatalog.ts"; import { @@ -107,6 +108,7 @@ import { shouldPersistProviderEvent, } from "../../provider/Layers/EventNdjsonLogger.ts"; import * as ProviderEventLoggers from "../../provider/Layers/ProviderEventLoggers.ts"; +import { spawnAndCollect } from "../../provider/providerSnapshot.ts"; import { claudeRateLimitEventToUpdate, type ClaudeScopedLimitNames, @@ -137,13 +139,12 @@ import { export const CLAUDE_PROVIDER = ProviderDriverKind.make("claudeAgent"); export const CLAUDE_AGENT_SDK_QUERY_PROTOCOL = "claude-agent-sdk.query" as const; -function claudeContextWindow(modelSelection: ModelSelection): number | null { - if (modelSelection.model === "claude-opus-4-6" || modelSelection.model === "claude-opus-4-7") { - return 1_000_000; - } - return resolveClaudeCatalogContextWindow(BUNDLED_CLAUDE_MODEL_CATALOG, modelSelection) === "1m" - ? 1_000_000 - : 200_000; +// The catalog knows fixed windows (Opus 4.7, 4.8) and the selected option for +// models with a 200k/1m choice. Unknown and custom models fall back to 200k. +function claudeContextWindow(modelSelection: ModelSelection): number { + return ( + resolveClaudeCatalogContextWindowTokens(BUNDLED_CLAUDE_MODEL_CATALOG, modelSelection) ?? 200_000 + ); } export function claudeProviderTurnTokenUsage( @@ -376,6 +377,8 @@ export class ClaudeAgentSdkQueryRunner extends Context.Service< export interface ClaudeAgentSdkSessionForkInput { readonly sessionId: string; readonly options: ForkSessionOptions; + /** The instance environment; its CLAUDE_CONFIG_DIR locates the session files. */ + readonly environment: NodeJS.ProcessEnv; readonly threadId: ThreadId; readonly providerSessionId: OrchestrationV2ProviderSession["id"]; } @@ -579,15 +582,74 @@ export function makeClaudeAgentSdkProtocolLogger(input: { }; } +const encodeClaudeHistoryOptions = Schema.encodeSync(Schema.fromJsonString(Schema.Unknown)); +const decodeClaudeHistoryFork = Schema.decodeUnknownEffect( + Schema.fromJsonString(Schema.Struct({ sessionId: Schema.String })), +); + +/** + * SDK history helpers read session files through `process.env`. A fork for an + * instance with its own CLAUDE_CONFIG_DIR runs in the history worker process + * with that instance's environment, so the server's environment never changes + * while other providers run. + */ +const forkClaudeSessionInWorker = Effect.fn("forkClaudeSessionInWorker")(function* ( + input: ClaudeAgentSdkSessionForkInput, +) { + // The single executable has no Node to run a sibling script, so it hosts + // the worker as a hidden subcommand of itself. + const workerArguments = (yield* HostProcessIsExecutable) + ? ["__claude-history"] + : [ + yield* Path.Path.pipe( + Effect.flatMap((path) => + path.fromFileUrl( + new URL( + import.meta.url.endsWith(".ts") + ? "../../claude-history-worker.ts" + : "./claude-history-worker.mjs", + import.meta.url, + ), + ), + ), + ), + ]; + const result = yield* spawnAndCollect( + process.execPath, + ChildProcess.make( + process.execPath, + [ + ...workerArguments, + "forkSession", + input.sessionId, + encodeClaudeHistoryOptions(input.options), + ], + { env: { ...input.environment, ELECTRON_RUN_AS_NODE: "1" } }, + ), + ).pipe(Effect.timeout("30 seconds")); + if (result.code !== 0) { + return yield* new ClaudeAgentSdkQueryRunnerError({ + method: "forkSession", + cause: result.stderr || "Claude history worker failed.", + }); + } + return yield* decodeClaudeHistoryFork(result.stdout); +}); + export const claudeAgentSdkQueryRunnerLiveLayer: Layer.Layer< ClaudeAgentSdkQueryRunner, never, - Crypto.Crypto | ProviderEventLoggers.ProviderEventLoggers + | ChildProcessSpawner.ChildProcessSpawner + | Crypto.Crypto + | Path.Path + | ProviderEventLoggers.ProviderEventLoggers > = Layer.effect( ClaudeAgentSdkQueryRunner, Effect.gen(function* () { const crypto = yield* Crypto.Crypto; const { native: nativeEventLogger } = yield* ProviderEventLoggers.ProviderEventLoggers; + const spawner = yield* ChildProcessSpawner.ChildProcessSpawner; + const path = yield* Path.Path; return ClaudeAgentSdkQueryRunner.of({ allocateSessionId: crypto.randomUUIDv4.pipe( @@ -716,10 +778,16 @@ export const claudeAgentSdkQueryRunnerLiveLayer: Layer.Layer< options: input.options, }, }); - const result = yield* Effect.tryPromise({ - try: () => forkClaudeSession(input.sessionId, input.options), - catch: (cause) => queryRunnerError(cause, "forkSession"), - }); + const fork: Effect.Effect = + input.environment.CLAUDE_CONFIG_DIR === process.env.CLAUDE_CONFIG_DIR + ? Effect.tryPromise(() => forkClaudeSession(input.sessionId, input.options)) + : forkClaudeSessionInWorker(input).pipe( + Effect.provideService(ChildProcessSpawner.ChildProcessSpawner, spawner), + Effect.provideService(Path.Path, path), + ); + const result = yield* fork.pipe( + Effect.mapError((cause) => queryRunnerError(cause, "forkSession")), + ); yield* logProtocolEvent({ direction: "incoming", stage: "decoded", @@ -7588,6 +7656,7 @@ export function makeClaudeAdapterV2( const forked = yield* queryRunner.forkSession({ sessionId: sourceNativeThreadId, options: forkOptions, + environment: adapterOptions.environment, threadId: forkInput.targetThreadId, providerSessionId: input.providerSessionId, }); @@ -7666,9 +7735,15 @@ export const createClaudeAdapterV2 = Effect.fn("ClaudeAdapterV2Driver.create")( const baseEnvironment = mergeProviderInstanceEnvironment(environment, hostEnvironment); const claudeEnvironment = yield* makeClaudeEnvironment(config, baseEnvironment); const path = yield* Path.Path; + // The SDK spawns this path directly, so Windows `claude` and `claude.cmd` + // must resolve to the real executable first. + const binaryPath = yield* resolveClaudeSdkExecutablePath( + expandHomePath(config.binaryPath), + claudeEnvironment, + ); return makeClaudeAdapterV2({ instanceId, - settings: { ...config, enabled, binaryPath: expandHomePath(config.binaryPath) }, + settings: { ...config, enabled, binaryPath }, environment: claudeEnvironment, attachmentsDir: serverConfig.attachmentsDir, fileSystem, diff --git a/apps/server/src/orchestration-v2/Orchestrator.ts b/apps/server/src/orchestration-v2/Orchestrator.ts index c52087778d63..69daad2b2a47 100644 --- a/apps/server/src/orchestration-v2/Orchestrator.ts +++ b/apps/server/src/orchestration-v2/Orchestrator.ts @@ -2247,6 +2247,17 @@ const makeOrchestrator = Effect.fn("orchestrationV2.Orchestrator.layer")(functio cause: `Thread ${command.threadId} worktree changed before the metadata update could be applied.`, }); } + if ( + command.type === "thread.metadata.update" && + command.expectedTitle !== undefined && + command.expectedTitle !== thread.title + ) { + return yield* new OrchestratorDispatchError({ + commandId: command.commandId, + commandType: command.type, + cause: `Thread ${command.threadId} title changed before the metadata update could be applied.`, + }); + } if (command.type === "thread.metadata.update" && command.expectedEmpty === true) { const records = yield* projectionStore .getThreadRecords(command.threadId, ["runs"]) diff --git a/apps/server/src/orchestration-v2/ProjectionStore.test.ts b/apps/server/src/orchestration-v2/ProjectionStore.test.ts index c64fece9db07..a2ef52c1226a 100644 --- a/apps/server/src/orchestration-v2/ProjectionStore.test.ts +++ b/apps/server/src/orchestration-v2/ProjectionStore.test.ts @@ -9,6 +9,7 @@ import { MessageId, type ModelSelection, NodeId, + PlanId, ProjectId, ProviderDriverKind, ProviderInstanceId, @@ -488,6 +489,53 @@ it.layer(TestLayer)("ProjectionStoreV2", (it) => { }), ); + it.effect("reports an unreadable plan row as a read error, not a defect", () => + Effect.gen(function* () { + const projectionStore = yield* ProjectionStore.ProjectionStoreV2; + const sql = yield* SqlClient.SqlClient; + const now = yield* DateTime.now; + const threadId = ThreadId.make("thread:bad-plan-row"); + const planId = PlanId.make("plan:bad-plan-row"); + yield* projectionStore.apply({ + id: EventId.make("event:bad-plan-row:thread"), + type: "thread.created", + threadId, + occurredAt: now, + payload: { + createdBy: "user", + creationSource: "web", + id: threadId, + projectId: ProjectId.make("project:bad-plan-row"), + title: "Bad plan row", + providerInstanceId, + modelSelection, + runtimeMode: "full-access", + interactionMode: "default", + branch: null, + worktreePath: null, + activeProviderThreadId: null, + lineage: { parentThreadId: null, relationshipToParent: null, rootThreadId: threadId }, + forkedFrom: null, + createdAt: now, + updatedAt: now, + archivedAt: null, + settledOverride: null, + settledAt: null, + lastVisitedAt: null, + deletedAt: null, + }, + }); + yield* sql` + INSERT INTO orchestration_v2_projection_plans + (plan_id, thread_id, run_id, node_id, kind, status, payload_json) + VALUES (${planId}, ${threadId}, NULL, 'node:bad-plan-row', 'proposed', 'active', '{not json') + `; + + const error = yield* projectionStore.getPlan(threadId, planId).pipe(Effect.flip); + assert.equal(error._tag, "ProjectionStoreReadError"); + }), + ); + it.effect("pages complete user turns through SQL regardless of tool count or payload size", () => Effect.gen(function* () { const projectionStore = yield* ProjectionStore.ProjectionStoreV2; diff --git a/apps/server/src/orchestration-v2/ProjectionStore.ts b/apps/server/src/orchestration-v2/ProjectionStore.ts index ac0624b05806..d569c1f7908e 100644 --- a/apps/server/src/orchestration-v2/ProjectionStore.ts +++ b/apps/server/src/orchestration-v2/ProjectionStore.ts @@ -1020,8 +1020,9 @@ const decodeRuntimeRequestPayload = Schema.decodeUnknownEffect( const decodeMessagePayload = Schema.decodeUnknownEffect( Schema.fromJsonString(OrchestrationV2ConversationMessageJsonSchema), ); -const decodePlanArtifact = Schema.decodeUnknownEffect(OrchestrationV2PlanArtifactSchema); -const decodePlanPayload = (json: string) => decodePlanArtifact(parseEncodedPayload(json)); +const decodePlanPayload = Schema.decodeUnknownEffect( + Schema.fromJsonString(OrchestrationV2PlanArtifactSchema), +); const decodeTurnItemPayload = Schema.decodeUnknownEffect( Schema.fromJsonString(OrchestrationV2TurnItemJsonSchema), ); @@ -3548,7 +3549,7 @@ export const layer: Layer.Layer = id: "message_id", decode: (payload) => decodeMessagePayload(payload), }, - { name: "plans", id: "plan_id", decode: decodePlanPayload }, + { name: "plans", id: "plan_id", decode: (payload) => decodePlanPayload(payload) }, { name: "turn_items", id: "turn_item_id", diff --git a/apps/server/src/orchestration-v2/ProviderTurnStartService.test.ts b/apps/server/src/orchestration-v2/ProviderTurnStartService.test.ts index 746fbe235d99..0c01fb435886 100644 --- a/apps/server/src/orchestration-v2/ProviderTurnStartService.test.ts +++ b/apps/server/src/orchestration-v2/ProviderTurnStartService.test.ts @@ -18,6 +18,7 @@ import { } from "@t3tools/contracts"; import * as Cause from "effect/Cause"; import * as DateTime from "effect/DateTime"; +import * as Deferred from "effect/Deferred"; import * as Effect from "effect/Effect"; import * as FileSystem from "effect/FileSystem"; import * as Layer from "effect/Layer"; @@ -27,6 +28,7 @@ import * as Schema from "effect/Schema"; import * as GitWorkflow from "../git/GitWorkflowService.ts"; import * as ProjectService from "../project/ProjectService.ts"; import * as ProviderAuthService from "../provider/Services/ProviderAuthService.ts"; +import * as ServerSettings from "../serverSettings.ts"; import * as ContextHandoffService from "./ContextHandoffService.ts"; import * as EventSink from "./EventSink.ts"; import * as IdAllocator from "./IdAllocator.ts"; @@ -36,6 +38,7 @@ import * as ProviderSessionManager from "./ProviderSessionManager.ts"; import * as ProviderTurnStart from "./ProviderTurnStartService.ts"; import * as RunExecutionService from "./RunExecutionService.ts"; import * as RuntimePolicy from "./RuntimePolicy.ts"; +import * as TemporaryBranchRename from "./TemporaryBranchRename.ts"; const isDomainEvent = Schema.is(OrchestrationV2DomainEvent); @@ -130,6 +133,8 @@ it("does not commit running state when inherited background routing cannot be re }), Layer.mock(RunExecutionService.RunExecutionServiceV2)({ startRootRun }), Layer.mock(RuntimePolicy.RuntimePolicyV2)({}), + ServerSettings.ServerSettingsService.layerTest({ worktreeSubmodules: "none" }), + Layer.mock(TemporaryBranchRename.TemporaryBranchRename)({}), ), ), ); @@ -142,16 +147,106 @@ it("does not commit running state when inherited background routing cannot be re expect(error._tag).toBe("ProviderTurnStartError"); expect(projectionReadCount).toBe(2); expect(pruneWorktrees).toHaveBeenCalledWith({ cwd: "/tmp/provider-turn-start-project" }); - expect(createWorktree).toHaveBeenCalledWith({ - cwd: "/tmp/provider-turn-start-project", - refName: "feature/restore", - path: "/tmp/missing-provider-turn-start-worktree", - }); + expect(createWorktree).toHaveBeenCalledWith( + { + cwd: "/tmp/provider-turn-start-project", + refName: "feature/restore", + path: "/tmp/missing-provider-turn-start-worktree", + }, + { submodules: "none" }, + ); expect(writeIfRunCurrent).not.toHaveBeenCalled(); expect(startRootRun).not.toHaveBeenCalled(); }).pipe(Effect.provide(layer), Effect.runPromise); }); +// The thread's first run was a /compact, so this is the first prompt to reach the provider. +it("names a temporary worktree branch from the first run that reaches the provider", async () => { + const threadId = ThreadId.make("thread_provider_turn_start_branch_name"); + const runId = RunId.make("run_provider_turn_start_branch_name"); + const messageId = MessageId.make("message_provider_turn_start_branch_name"); + const rootNodeId = NodeId.make("node_provider_turn_start_branch_name"); + const attemptId = RunAttemptId.make("attempt_provider_turn_start_branch_name"); + const providerThreadId = ProviderThreadId.make("provider_thread_provider_turn_start_branch_name"); + const checkpointScopeId = CheckpointScopeId.make("checkpoint_scope_provider_turn_start_branch"); + const projection = { + thread: { + id: threadId, + projectId: ProjectId.make("project_provider_turn_start_branch_name"), + branch: "t3code/1a2b3c4d", + worktreePath: "/tmp/provider-turn-start-branch-name", + }, + runs: [ + { + id: runId, + status: "starting", + rootNodeId, + activeAttemptId: attemptId, + providerThreadId, + userMessageId: messageId, + ordinal: 2, + }, + ], + nodes: [{ id: rootNodeId, checkpointScopeId }], + attempts: [{ id: attemptId }], + providerThreads: [ + { + id: providerThreadId, + providerSessionId: ProviderSessionId.make("provider_session_branch_name"), + }, + ], + messages: [{ id: messageId, text: "Fix the login bug", attachments: [] }], + checkpointScopes: [{ id: checkpointScopeId }], + contextHandoffs: [], + contextTransfers: [], + turnItems: [], + } as unknown as OrchestrationV2ThreadProjection; + await Effect.gen(function* () { + const renamed = yield* Deferred.make(); + const layer = ProviderTurnStart.layer.pipe( + Layer.provide( + Layer.mergeAll( + Layer.mock(ContextHandoffService.ContextHandoffServiceV2)({}), + Layer.mock(EventSink.EventSinkV2)({}), + IdAllocator.layer, + Layer.succeed(FileSystem.FileSystem, { exists: () => Effect.succeed(true) } as never), + Layer.mock(GitWorkflow.GitWorkflowService)({}), + Layer.mock(ProjectService.ProjectService)({}), + Layer.mock(ProjectionStore.ProjectionStoreV2)({ + getTurnStartContext: () => Effect.succeed({ ...projection, hasConversation: false }), + // Stops the start right after the rename is forked. + getRuntimeRecoveryProjection: () => + Effect.fail( + new ProjectionStore.ProjectionStoreReadError({ threadId, cause: "stop" }), + ), + }), + Layer.mock(ProviderSessionManager.ProviderSessionManagerV2)({}), + Layer.mock(ProviderAuthService.ProviderAuthService)({ + tryHandlePromptCommand: () => Effect.succeed(false), + }), + Layer.mock(RunExecutionService.RunExecutionServiceV2)({}), + Layer.mock(RuntimePolicy.RuntimePolicyV2)({}), + ServerSettings.ServerSettingsService.layerTest(), + Layer.mock(TemporaryBranchRename.TemporaryBranchRename)({ + rename: (input) => Deferred.succeed(renamed, input).pipe(Effect.asVoid), + }), + ), + ), + ); + yield* Effect.gen(function* () { + yield* (yield* ProviderTurnStart.ProviderTurnStartServiceV2) + .start({ threadId, runId }) + .pipe(Effect.ignore); + expect(yield* Deferred.await(renamed)).toMatchObject({ + threadId, + branch: "t3code/1a2b3c4d", + worktreePath: "/tmp/provider-turn-start-branch-name", + message: { text: "Fix the login bug" }, + }); + }).pipe(Effect.provide(layer)); + }).pipe(Effect.runPromise); +}); + function makeLocalCommandHarness(input: { readonly text: string; readonly previousNativeSession?: boolean; @@ -493,6 +588,8 @@ function makeLocalCommandHarness(input: { Layer.mock(RuntimePolicy.RuntimePolicyV2)({ resolve: () => Effect.succeed({} as never), }), + ServerSettings.ServerSettingsService.layerTest(), + Layer.mock(TemporaryBranchRename.TemporaryBranchRename)({}), ), ), ); diff --git a/apps/server/src/orchestration-v2/ProviderTurnStartService.ts b/apps/server/src/orchestration-v2/ProviderTurnStartService.ts index 781491ac7fd6..ce28082f1e4f 100644 --- a/apps/server/src/orchestration-v2/ProviderTurnStartService.ts +++ b/apps/server/src/orchestration-v2/ProviderTurnStartService.ts @@ -1,5 +1,6 @@ import { modelSelectionsEqual } from "@t3tools/shared/model"; import { projectComposerContextForProvider } from "@t3tools/shared/composerContextReferences"; +import { resolveProjectSettings } from "@t3tools/shared/projectSettings"; import { CommandId, type OrchestrationV2DomainEvent, @@ -15,14 +16,17 @@ import * as Context from "effect/Context"; import * as Cause from "effect/Cause"; import * as DateTime from "effect/DateTime"; import * as Effect from "effect/Effect"; +import * as Exit from "effect/Exit"; import * as FileSystem from "effect/FileSystem"; import * as Layer from "effect/Layer"; import * as Option from "effect/Option"; +import * as Scope from "effect/Scope"; import * as Schema from "effect/Schema"; import * as GitWorkflowService from "../git/GitWorkflowService.ts"; import * as ProjectService from "../project/ProjectService.ts"; import * as ProviderAuthService from "../provider/Services/ProviderAuthService.ts"; +import * as ServerSettings from "../serverSettings.ts"; import * as EventSink from "./EventSink.ts"; import * as ContextHandoffService from "./ContextHandoffService.ts"; import { @@ -47,6 +51,7 @@ import * as ProviderSessionManager from "./ProviderSessionManager.ts"; import { makeProviderFailure } from "./ProviderFailure.ts"; import * as RunExecutionService from "./RunExecutionService.ts"; import * as RuntimePolicy from "./RuntimePolicy.ts"; +import * as TemporaryBranchRename from "./TemporaryBranchRename.ts"; import { isRestartNoteContinuation, pendingRestartCancelledBackgroundWork, @@ -95,6 +100,8 @@ export const layer: Layer.Layer< | ProviderSessionManager.ProviderSessionManagerV2 | RunExecutionService.RunExecutionServiceV2 | RuntimePolicy.RuntimePolicyV2 + | ServerSettings.ServerSettingsService + | TemporaryBranchRename.TemporaryBranchRename > = Layer.effect( ProviderTurnStartServiceV2, Effect.gen(function* () { @@ -109,6 +116,10 @@ export const layer: Layer.Layer< const providerSessions = yield* ProviderSessionManager.ProviderSessionManagerV2; const runExecution = yield* RunExecutionService.RunExecutionServiceV2; const runtimePolicy = yield* RuntimePolicy.RuntimePolicyV2; + const serverSettings = yield* ServerSettings.ServerSettingsService; + const branchRename = yield* TemporaryBranchRename.TemporaryBranchRename; + const backgroundScope = yield* Scope.make("sequential"); + yield* Effect.addFinalizer(() => Scope.close(backgroundScope, Exit.void)); // These callbacks outlive startup while a run drains background work. Build // them outside start's scope so they cannot retain its full thread history. @@ -470,13 +481,22 @@ export const layer: Layer.Layer< worktreePath, branch, }); + // Best effort like the rest of this recovery: a settings read + // failure falls back to the checkout's t3.json. + const submodules = yield* serverSettings.getSettings.pipe( + Effect.map( + (settings) => + resolveProjectSettings(settings, projection.thread.projectId).settings + .worktreeSubmodules, + ), + Effect.orElseSucceed(() => null), + ); yield* gitWorkflow.pruneWorktrees({ cwd: project.workspaceRoot }).pipe( Effect.andThen( - gitWorkflow.createWorktree({ - cwd: project.workspaceRoot, - refName: branch, - path: worktreePath, - }), + gitWorkflow.createWorktree( + { cwd: project.workspaceRoot, refName: branch, path: worktreePath }, + { submodules }, + ), ), Effect.catchCause((cause) => Cause.hasInterruptsOnly(cause) @@ -491,6 +511,20 @@ export const layer: Layer.Layer< } } } + // A run that reaches the provider names a still-temporary worktree branch + // from its message, for threads launched without one. Commands handled + // above, such as /compact, return first and never name it. The service attempts + // each thread once, in the background so naming never delays the turn. + yield* branchRename + .rename({ + threadId: projection.thread.id, + projectId: projection.thread.projectId, + commandId: CommandId.make(run.id), + branch, + worktreePath, + message, + }) + .pipe(Effect.forkIn(backgroundScope)); const selectInheritedBackgroundItems = ( current: ProjectionStore.ProjectionRuntimeRecoveryState, ): ReturnType => diff --git a/apps/server/src/orchestration-v2/TemporaryBranchRename.ts b/apps/server/src/orchestration-v2/TemporaryBranchRename.ts new file mode 100644 index 000000000000..9e7765ed62e3 --- /dev/null +++ b/apps/server/src/orchestration-v2/TemporaryBranchRename.ts @@ -0,0 +1,139 @@ +import { + CommandId, + type ChatAttachment, + type OrchestrationMessageContext, + type ProjectId, + type ThreadId, +} from "@t3tools/contracts"; +import { isTemporaryWorktreeBranch } from "@t3tools/shared/git"; +import { resolveProjectSettings } from "@t3tools/shared/projectSettings"; +import * as Context from "effect/Context"; +import * as Effect from "effect/Effect"; +import * as Layer from "effect/Layer"; +import * as Option from "effect/Option"; + +import * as GitWorkflow from "../git/GitWorkflowService.ts"; +import * as ProviderRegistry from "../provider/Services/ProviderRegistry.ts"; +import * as ServerSettings from "../serverSettings.ts"; +import * as TextGeneration from "../textGeneration/TextGeneration.ts"; +import * as ThreadManagement from "./ThreadManagementService.ts"; + +export interface TemporaryBranchRenameInput { + readonly threadId: ThreadId; + readonly projectId: ProjectId; + readonly commandId: CommandId; + readonly branch: string | null; + readonly worktreePath: string | null; + readonly message: { + readonly text: string; + readonly attachments: ReadonlyArray; + readonly context?: OrchestrationMessageContext | undefined; + }; +} + +/** + * Names a worktree thread's temporary `t3code/` branch from its first + * message. Thread launch calls it when it has that message, and the thread's + * first run start calls it too, which covers launches without one. Each thread + * is attempted once per server, so the two callers never rename twice. A + * failure is logged and the temporary name stays. + */ +export class TemporaryBranchRename extends Context.Service< + TemporaryBranchRename, + { + readonly rename: (input: TemporaryBranchRenameInput) => Effect.Effect; + } +>()("t3/orchestration-v2/TemporaryBranchRename") {} + +const make = Effect.gen(function* () { + const git = yield* GitWorkflow.GitWorkflowService; + // Optional: without it the writer model is used without an availability check. + const providerRegistry = yield* Effect.serviceOption(ProviderRegistry.ProviderRegistry); + const serverSettings = yield* ServerSettings.ServerSettingsService; + const textGeneration = yield* TextGeneration.TextGeneration; + const threads = yield* ThreadManagement.ThreadManagementService; + const attempted = new Set(); + + const rename = Effect.fn("TemporaryBranchRename.rename")(function* ( + input: TemporaryBranchRenameInput, + ) { + const { branch: oldBranch, worktreePath: cwd } = input; + if (cwd === null || oldBranch === null || !isTemporaryWorktreeBranch(oldBranch)) return; + const first = yield* Effect.sync(() => { + if (attempted.has(input.threadId)) return false; + attempted.add(input.threadId); + return true; + }); + if (!first) return; + yield* Effect.gen(function* () { + const settings = resolveProjectSettings( + yield* serverSettings.getSettings, + input.projectId, + ).settings; + const modelSelection = + settings.sourceControlWriterModelSelection === null + ? settings.textGenerationModelSelection + : ServerSettings.resolveSourceControlWriterModelSelection( + settings, + Option.isSome(providerRegistry) + ? yield* providerRegistry.value.getProviders + : undefined, + ); + const generated = yield* textGeneration.generateBranchName({ + naming: { + mode: settings.branchNamingMode, + prefix: settings.branchNamePrefix, + instructions: settings.branchNameInstructions, + }, + cwd, + message: input.message.text, + attachments: input.message.attachments, + ...(input.message.context ? { context: input.message.context } : {}), + modelSelection, + }); + const renamed = yield* git.renameBranch({ + cwd, + oldBranch, + newBranch: generated.branch, + ...(settings.branchNamingMode === "custom" ? { exactName: true } : {}), + }); + // The update is rejected if the thread moved to another worktree during + // generation. Any failed update puts the old name back, so the thread + // never points at a branch that no longer exists. + yield* threads + .dispatch({ + type: "thread.metadata.update", + commandId: CommandId.make(`${input.commandId}:branch-rename`), + threadId: input.threadId, + branch: renamed.branch, + worktreePath: cwd, + expectedWorktreePath: cwd, + }) + .pipe( + Effect.tapError(() => + git + .renameBranch({ + cwd, + oldBranch: renamed.branch, + newBranch: oldBranch, + exactName: true, + }) + .pipe(Effect.ignore), + ), + ); + }).pipe( + Effect.catchCause((cause) => + Effect.logWarning("Thread worktree branch rename failed", { + commandId: input.commandId, + threadId: input.threadId, + oldBranch, + cause, + }), + ), + ); + }); + + return TemporaryBranchRename.of({ rename }); +}); + +export const layer = Layer.effect(TemporaryBranchRename, make); diff --git a/apps/server/src/orchestration-v2/ThreadLaunchService.test.ts b/apps/server/src/orchestration-v2/ThreadLaunchService.test.ts index 10eb4ac9485e..14997735b25e 100644 --- a/apps/server/src/orchestration-v2/ThreadLaunchService.test.ts +++ b/apps/server/src/orchestration-v2/ThreadLaunchService.test.ts @@ -55,6 +55,7 @@ import * as IdAllocator from "./IdAllocator.ts"; import type { ProviderAdapterV2Shape } from "./ProviderAdapter.ts"; import * as ProviderAdapterRegistry from "./ProviderAdapterRegistry.ts"; import * as ThreadLaunch from "./ThreadLaunchService.ts"; +import * as TemporaryBranchRename from "./TemporaryBranchRename.ts"; import * as ThreadManagement from "./ThreadManagementService.ts"; import * as ThreadTitleRegeneration from "./ThreadTitleRegenerationService.ts"; import { makeOrchestratorV2ReplayLayerWithRegistry } from "./testkit/ProviderReplayHarness.ts"; @@ -183,8 +184,13 @@ function makeHarness(options: HarnessOptions = {}) { folderForThread: () => Effect.succeed(Option.none()), }), ); + const branchRename = TemporaryBranchRename.layer.pipe( + Layer.provide(Layer.merge(externalServices, threadManagement)), + ); const launch = ThreadLaunch.layer.pipe( - Layer.provide(Layer.mergeAll(externalServices, threadManagement, receipts, IdAllocator.layer)), + Layer.provide( + Layer.mergeAll(externalServices, threadManagement, receipts, IdAllocator.layer, branchRename), + ), ); const projectedProjects = Layer.mock(ProjectStore.ProjectStoreV2)({ get: (requestedProjectId) => @@ -1118,6 +1124,7 @@ it.effect("renames a temporary t3code/ branch off the provisioning critica Effect.andThen(Deferred.await(allowBranchName)), Effect.as({ branch: "generated-branch" }), ), + serverSettings: { worktreeSubmodules: "top-level" }, }); yield* Effect.gen(function* () { const launches = yield* ThreadLaunch.ThreadLaunchService; @@ -1132,6 +1139,7 @@ it.effect("renames a temporary t3code/ branch off the provisioning critica ); yield* Deferred.await(branchNameStarted); assert.equal(harness.createWorktree.mock.calls[0]?.[0]?.newRefName, "t3code/abcd1234"); + assert.equal(harness.createWorktree.mock.calls[0]?.[1]?.submodules, "top-level"); yield* waitUntil(() => threads .getThreadProjection(launched.threadId) diff --git a/apps/server/src/orchestration-v2/ThreadLaunchService.ts b/apps/server/src/orchestration-v2/ThreadLaunchService.ts index faf613376920..af1ed1ce1d0c 100644 --- a/apps/server/src/orchestration-v2/ThreadLaunchService.ts +++ b/apps/server/src/orchestration-v2/ThreadLaunchService.ts @@ -28,19 +28,18 @@ import * as Option from "effect/Option"; import * as Ref from "effect/Ref"; import * as Schema from "effect/Schema"; import * as Scope from "effect/Scope"; -import { buildTemporaryWorktreeBranchName, isTemporaryWorktreeBranch } from "@t3tools/shared/git"; +import { buildTemporaryWorktreeBranchName } from "@t3tools/shared/git"; import * as GitWorkflow from "../git/GitWorkflowService.ts"; import * as ProjectService from "../project/ProjectService.ts"; import * as ProjectSetupScriptRunner from "../project/ProjectSetupScriptRunner.ts"; import * as ManagedProjectFolders from "../project/ManagedProjectFolders.ts"; -import * as ProviderRegistry from "../provider/Services/ProviderRegistry.ts"; import * as ServerSettings from "../serverSettings.ts"; -import * as TextGeneration from "../textGeneration/TextGeneration.ts"; import * as CommandReceiptStore from "./CommandReceiptStore.ts"; import * as IdAllocator from "./IdAllocator.ts"; import { makeProviderFailure } from "./ProviderFailure.ts"; import { randomUuidV4 } from "./RandomUuid.ts"; +import * as TemporaryBranchRename from "./TemporaryBranchRename.ts"; import * as ThreadManagement from "./ThreadManagementService.ts"; export type ThreadLaunchWorkspaceStrategy = @@ -149,13 +148,12 @@ const make = Effect.gen(function* () { const terminals = yield* TerminalManager.TerminalManager; const git = yield* GitWorkflow.GitWorkflowService; const setupScripts = yield* ProjectSetupScriptRunner.ProjectSetupScriptRunner; - const providerRegistry = yield* ProviderRegistry.ProviderRegistry; const serverSettings = yield* ServerSettings.ServerSettingsService; - const textGeneration = yield* TextGeneration.TextGeneration; const receipts = yield* CommandReceiptStore.CommandReceiptStoreV2; const ids = yield* IdAllocator.IdAllocatorV2; const threads = yield* ThreadManagement.ThreadManagementService; const managedFolders = yield* ManagedProjectFolders.ManagedProjectFolders; + const branchRename = yield* TemporaryBranchRename.TemporaryBranchRename; const preparationScope = yield* Scope.make("sequential"); const scheduledLaunches = yield* Ref.make>(new Set()); yield* Effect.addFinalizer(() => Scope.close(preparationScope, Exit.void)); @@ -229,41 +227,6 @@ const make = Effect.gen(function* () { }); } yield* Effect.gen(function* () { - const initialMessage = input.initialMessage; - const generateBranchNameFor = (cwd: string, message: ThreadLaunchInitialMessage) => - Effect.gen(function* () { - const settings = resolveProjectSettings( - yield* serverSettings.getSettings, - input.projectId, - ).settings; - const modelSelection = - settings.sourceControlWriterModelSelection === null - ? settings.textGenerationModelSelection - : ServerSettings.resolveSourceControlWriterModelSelection( - settings, - yield* providerRegistry.getProviders, - ); - return yield* textGeneration - .generateBranchName({ - naming: { - mode: settings.branchNamingMode, - prefix: settings.branchNamePrefix, - instructions: settings.branchNameInstructions, - }, - cwd, - message: message.text, - attachments: message.attachments, - ...(message.context ? { context: message.context } : {}), - modelSelection, - }) - .pipe( - Effect.map((result) => ({ - branch: result.branch, - exactName: settings.branchNamingMode === "custom", - })), - ); - }); - // The server owns worktree naming: without an explicit branch, provision // under a temporary `t3code/` name so the worktree never waits on // name generation, then rename in the background below. @@ -330,6 +293,14 @@ const make = Effect.gen(function* () { } if (startFromOrigin) yield* setupTracker.stageStatus(threadId, "fetch", "done"); yield* setupTracker.stageStatus(threadId, "checkout", "running"); + // Best effort: a settings read failure falls back to the checkout's t3.json. + const submodules = yield* serverSettings.getSettings.pipe( + Effect.map( + (settings) => + resolveProjectSettings(settings, input.projectId).settings.worktreeSubmodules, + ), + Effect.orElseSucceed(() => null), + ); const worktree = yield* git .createWorktree( { @@ -340,6 +311,7 @@ const make = Effect.gen(function* () { path: null, }, { + submodules, progress: { onWorktreeClaimed: (path) => Effect.sync(() => { @@ -370,44 +342,19 @@ const make = Effect.gen(function* () { // Rename temporary branches (server-invented above, or sent by clients // that name worktrees themselves) in the background so generation latency - // never delays provisioning or the provider turn. The temporary name - // simply sticks if generation or the rename fails. - if ( - worktreePath !== null && - branch !== null && - initialMessage !== undefined && - isTemporaryWorktreeBranch(branch) - ) { - const oldBranch = branch; - const worktreeCwd = worktreePath; - yield* generateBranchNameFor(worktreeCwd, initialMessage).pipe( - Effect.flatMap(({ branch: newBranch, exactName }) => - git.renameBranch({ - cwd: worktreeCwd, - oldBranch, - newBranch, - ...(exactName ? { exactName: true } : {}), - }), - ), - Effect.flatMap((renamed) => - threads.dispatch({ - type: "thread.metadata.update", - commandId: CommandId.make(`${input.commandId}:branch-rename`), - threadId, - branch: renamed.branch, - worktreePath: worktreeCwd, - }), - ), - Effect.catchCause((cause) => - Effect.logWarning("Thread worktree branch rename failed", { - commandId: input.commandId, - threadId, - oldBranch, - cause, - }), - ), - Effect.forkIn(preparationScope), - ); + // never delays provisioning or the provider turn. Without a first message + // here, the thread's first run start renames it instead. + if (input.initialMessage !== undefined) { + yield* branchRename + .rename({ + threadId, + projectId: input.projectId, + commandId: input.commandId, + branch, + worktreePath, + message: input.initialMessage, + }) + .pipe(Effect.forkIn(preparationScope)); } const cwd = worktreePath ?? project.workspaceRoot; diff --git a/apps/server/src/orchestration-v2/ThreadTitleRegenerationService.test.ts b/apps/server/src/orchestration-v2/ThreadTitleRegenerationService.test.ts index 21a22a60dd94..565b8ccd669e 100644 --- a/apps/server/src/orchestration-v2/ThreadTitleRegenerationService.test.ts +++ b/apps/server/src/orchestration-v2/ThreadTitleRegenerationService.test.ts @@ -2,22 +2,27 @@ import { assert, describe, it, vi } from "@effect/vitest"; import { type ChatAttachment, CommandId, + EventId, MessageId, + type OrchestrationV2AppThread, ProjectId, ProviderDriverKind, ProviderInstanceId, ThreadId, TextGenerationError, } from "@t3tools/contracts"; +import * as DateTime from "effect/DateTime"; import * as Effect from "effect/Effect"; import * as Deferred from "effect/Deferred"; import * as Fiber from "effect/Fiber"; import * as TestClock from "effect/testing/TestClock"; import * as Layer from "effect/Layer"; import * as Option from "effect/Option"; +import * as Stream from "effect/Stream"; import { SqlitePersistenceMemory } from "../persistence/Layers/Sqlite.ts"; import * as ProjectStore from "./ProjectStore.ts"; +import * as EventSink from "./EventSink.ts"; import * as ServerSettings from "../serverSettings.ts"; import * as TextGeneration from "../textGeneration/TextGeneration.ts"; import { CodexProviderCapabilitiesV2 } from "./Adapters/CodexAdapterV2.ts"; @@ -59,6 +64,7 @@ function makeHarness( const generateThreadTitle = vi.fn( options.generateTitle ?? (() => Effect.succeed({ title: "Generated title" })), ); + const projectedProjects = Layer.mock(ProjectStore.ProjectStoreV2)({ get: (requestedProjectId) => Effect.succeed( @@ -91,11 +97,27 @@ function makeHarness( ), ); return { - layer: Layer.mergeAll(threadManagement, titleRegeneration, outbox, database), + layer: Layer.mergeAll(orchestrator, threadManagement, titleRegeneration, outbox, database), generateThreadTitle, }; } +// Waits on the stored event stream, so a change that already landed counts. +function awaitThreadUpdate( + threadId: ThreadId, + predicate: (thread: OrchestrationV2AppThread) => boolean, +) { + return Effect.gen(function* () { + const threads = yield* ThreadManagement.ThreadManagementService; + return yield* threads.streamStoredEventsFrom({ threadId, afterSequence: 0 }).pipe( + Stream.filter( + ({ event }) => event.type === "thread.metadata-updated" && predicate(event.payload), + ), + Stream.runHead, + ); + }); +} + function createThread(input: { readonly command: string; readonly thread: string }) { return Effect.gen(function* () { const threads = yield* ThreadManagement.ThreadManagementService; @@ -438,6 +460,74 @@ describe("ThreadTitleRegenerationService", () => { ); }); +describe("ThreadTitleRegenerationService refinement", () => { + it.effect("does not arm a refinement after the title changed", () => + Effect.gen(function* () { + const harness = makeHarness(); + yield* Effect.gen(function* () { + const threads = yield* ThreadManagement.ThreadManagementService; + const threadId = yield* createThread({ + command: "create:renamed", + thread: "thread:renamed", + }); + const error = yield* threads + .dispatch({ + type: "thread.metadata.update", + commandId: CommandId.make("refine:renamed"), + threadId, + expectedTitle: "Fix this", + regenerateTitle: true, + }) + .pipe(Effect.flip); + assert.include(String(error.cause), "title changed"); + assert.isNotOk((yield* threads.getThreadProjection(threadId)).thread.titleRegeneration); + }).pipe(Effect.provide(harness.layer)); + }), + ); + + it.effect("refines a generic first title once the first run completes", () => + Effect.gen(function* () { + const harness = makeHarness({ + generateTitle: () => Effect.succeed({ title: "Fix this", needsRefinement: true }), + }); + yield* Effect.gen(function* () { + const threads = yield* ThreadManagement.ThreadManagementService; + const eventSink = yield* EventSink.EventSinkV2; + const service = yield* ThreadTitleRegeneration.ThreadTitleRegenerationService; + const threadId = yield* createThread({ command: "create:refine", thread: "thread:refine" }); + yield* dispatchUserMessage({ command: "message:refine", threadId, text: "fix this" }); + const requestId = yield* armRegeneration({ command: "title:refine", threadId }); + yield* service.execute({ + threadId, + requestId, + kind: { type: "initial", messageId: MessageId.make("message:refine:message") }, + }); + assert.equal((yield* threads.getThreadProjection(threadId)).thread.title, "Fix this"); + + const { runs } = yield* threads.getThreadProjection(threadId); + const run = runs[0]!; + yield* eventSink.write({ + events: [ + { + id: EventId.make("event:refine:run-completed"), + type: "run.updated", + threadId, + occurredAt: yield* DateTime.now, + payload: { ...run, status: "completed" }, + }, + ], + }); + + const refinement = yield* awaitThreadUpdate( + threadId, + (thread) => thread.titleRegeneration?.requestId === `${requestId}:title-refine`, + ); + assert.isTrue(Option.isSome(refinement)); + }).pipe(Effect.provide(harness.layer)); + }), + ); +}); + it.effect.each(["success", "exhausted", "stale", "interrupted"] as const)( "initial title retry: %s", (outcome) => diff --git a/apps/server/src/orchestration-v2/ThreadTitleRegenerationService.ts b/apps/server/src/orchestration-v2/ThreadTitleRegenerationService.ts index 4d3961019314..809868dd7e39 100644 --- a/apps/server/src/orchestration-v2/ThreadTitleRegenerationService.ts +++ b/apps/server/src/orchestration-v2/ThreadTitleRegenerationService.ts @@ -9,9 +9,12 @@ import { import * as Cause from "effect/Cause"; import * as Context from "effect/Context"; import * as Effect from "effect/Effect"; +import * as Exit from "effect/Exit"; import * as Layer from "effect/Layer"; import * as Option from "effect/Option"; import * as Schedule from "effect/Schedule"; +import * as Scope from "effect/Scope"; +import * as Stream from "effect/Stream"; import * as ServerSettings from "../serverSettings.ts"; import * as TextGeneration from "../textGeneration/TextGeneration.ts"; @@ -43,6 +46,53 @@ const make = Effect.gen(function* () { const projects = yield* ProjectStore.ProjectStoreV2; const serverSettings = yield* ServerSettings.ServerSettingsService; const textGeneration = yield* TextGeneration.TextGeneration; + const backgroundScope = yield* Scope.make("sequential"); + yield* Effect.addFinalizer(() => Scope.close(backgroundScope, Exit.void)); + + // A generic first title ("Fix this") is refined once, after the first run + // completes, from the whole conversation. Replaying the thread's events from + // the start also sees a run that already ended. Any title change in the + // meantime, such as a user rename, cancels the refinement. + const refineAfterFirstRun = (input: { + readonly threadId: ThreadId; + readonly requestId: CommandId; + readonly title: string; + }) => + threads.streamStoredEventsFrom({ threadId: input.threadId, afterSequence: 0 }).pipe( + Stream.filter( + ({ event }) => + event.type === "run.updated" && + event.payload.ordinal === 1 && + ThreadManagementService.isTerminalRunStatus(event.payload.status), + ), + Stream.runHead, + Effect.flatMap((ended) => + Option.isNone(ended) || + ended.value.event.type !== "run.updated" || + ended.value.event.payload.status !== "completed" + ? Effect.void + : threads.getThreadShell(input.threadId).pipe( + Effect.flatMap((thread) => + thread === null || thread.title !== input.title || thread.titleRegeneration + ? Effect.void + : threads.dispatch({ + type: "thread.metadata.update", + commandId: CommandId.make(`${input.requestId}:title-refine`), + threadId: input.threadId, + // A rename that lands after the read above still wins. + expectedTitle: input.title, + regenerateTitle: true, + }), + ), + ), + ), + Effect.catchCause((cause) => + Effect.logWarning("Thread title refinement failed", { + threadId: input.threadId, + cause, + }), + ), + ); const complete = (input: { readonly threadId: ThreadId; @@ -64,7 +114,12 @@ const make = Effect.gen(function* () { )(function* (input) { const outcome: | { readonly type: "stale" } - | { readonly type: "complete"; readonly title?: string } = yield* Effect.gen(function* () { + | { + readonly type: "complete"; + readonly title?: string; + /** Set when a generic initial title should be refined after the first run. */ + readonly refineTitle?: string; + } = yield* Effect.gen(function* () { const projection = yield* threads.getThreadRecords( input.threadId, ["messages"], @@ -115,10 +170,19 @@ const make = Effect.gen(function* () { modelSelection: settings.textGenerationModelSelection, }); const generatedTitle = result.title.trim(); - return generatedTitle === "New thread" || + const title = + generatedTitle === "New thread" || (input.kind.type === "regenerate" && generatedTitle === projection.thread.title.trim()) - ? { type: "complete" as const } - : { type: "complete" as const, title: result.title }; + ? undefined + : result.title; + const refine = + input.kind.type === "initial" && + (result.needsRefinement === true || generatedTitle === "New thread"); + return { + type: "complete" as const, + ...(title === undefined ? {} : { title }), + ...(refine ? { refineTitle: title ?? projection.thread.title } : {}), + }; }).pipe( Effect.retry({ times: input.kind.type === "initial" ? 2 : 0, @@ -142,6 +206,13 @@ const make = Effect.gen(function* () { ...input, ...(outcome.title === undefined ? {} : { title: outcome.title }), }); + if (outcome.refineTitle !== undefined) { + yield* refineAfterFirstRun({ + threadId: input.threadId, + requestId: input.requestId, + title: outcome.refineTitle, + }).pipe(Effect.forkIn(backgroundScope)); + } }); return ThreadTitleRegenerationService.of({ execute }); diff --git a/apps/server/src/orchestration-v2/runtimeLayer.ts b/apps/server/src/orchestration-v2/runtimeLayer.ts index 8573c6570795..3bbb97d2242f 100644 --- a/apps/server/src/orchestration-v2/runtimeLayer.ts +++ b/apps/server/src/orchestration-v2/runtimeLayer.ts @@ -47,6 +47,7 @@ import * as RuntimePolicy from "./RuntimePolicy.ts"; import { layer as runtimeRequestServiceLayer } from "./RuntimeRequestService.ts"; import { layerWithLegacyImporter as threadManagementServiceLayer } from "./ThreadManagementService.ts"; import { layer as threadLaunchServiceLayer } from "./ThreadLaunchService.ts"; +import * as TemporaryBranchRename from "./TemporaryBranchRename.ts"; import { layer as threadLifecycleServiceLayer } from "./ThreadLifecycleService.ts"; import { layer as threadForkServiceLayer } from "./ThreadForkService.ts"; import { layer as turnItemPositionStoreLayer } from "./TurnItemPositionStore.ts"; @@ -146,21 +147,6 @@ const runExecutionServiceProvided = runExecutionServiceLayer.pipe( ), ); -const providerTurnStartServiceProvided = providerTurnStartServiceLayer.pipe( - Layer.provide( - Layer.mergeAll( - contextHandoffServiceProvided, - eventSinkProvided, - idAllocatorLayer, - projectionStoreLayer, - providerSessionManagerProvided, - providerAuthServiceProvided, - runExecutionServiceProvided, - runtimePolicyProvided, - ), - ), -); - const providerTurnControlServiceProvided = providerTurnControlServiceLayer.pipe( Layer.provide(Layer.merge(projectionStoreLayer, providerSessionManagerProvided)), ); @@ -235,6 +221,25 @@ const agentSessionImporterProvided = agentSessionImporterLayer.pipe( const threadManagementProvided = threadManagementServiceLayer.pipe( Layer.provide(Layer.merge(orchestratorProvided, legacyV1ThreadImporterProvided)), ); +// One instance, so launch and the first run start share its once-per-thread guard. +const temporaryBranchRenameProvided = TemporaryBranchRename.layer.pipe( + Layer.provide(Layer.merge(threadManagementProvided, TextGeneration.layer)), +); +const providerTurnStartServiceProvided = providerTurnStartServiceLayer.pipe( + Layer.provide( + Layer.mergeAll( + contextHandoffServiceProvided, + eventSinkProvided, + idAllocatorLayer, + projectionStoreLayer, + providerSessionManagerProvided, + providerAuthServiceProvided, + runExecutionServiceProvided, + runtimePolicyProvided, + temporaryBranchRenameProvided, + ), + ), +); export const ProjectSetupScriptRunnerLayerLive = projectSetupScriptRunnerLayer.pipe( Layer.provide(ProjectServiceLayerLive), ); @@ -250,6 +255,7 @@ const threadLaunchProvided = threadLaunchServiceLayer.pipe( threadManagementProvided, commandReceiptStoreProvided, idAllocatorLayer, + temporaryBranchRenameProvided, ), ), ); diff --git a/apps/server/src/orchestration-v2/testkit/ProviderReplayHarness.ts b/apps/server/src/orchestration-v2/testkit/ProviderReplayHarness.ts index f4a1201f1f3e..316336560d1f 100644 --- a/apps/server/src/orchestration-v2/testkit/ProviderReplayHarness.ts +++ b/apps/server/src/orchestration-v2/testkit/ProviderReplayHarness.ts @@ -15,6 +15,7 @@ import * as CheckpointStore from "../../checkpointing/CheckpointStore.ts"; import * as ServerConfig from "../../config.ts"; import { SqlitePersistenceMemory } from "../../persistence/Layers/Sqlite.ts"; import * as ServerSettings from "../../serverSettings.ts"; +import * as TemporaryBranchRename from "../TemporaryBranchRename.ts"; import * as ThreadManagementService from "../ThreadManagementService.ts"; import * as McpSessionRegistryTestkit from "../../mcp/McpSessionRegistry.testkit.ts"; import * as VcsDriverRegistry from "../../vcs/VcsDriverRegistry.ts"; @@ -374,6 +375,9 @@ export function makeOrchestratorV2ReplayLayerWithRegistry( }), runExecutionServiceProvided, runtimeLayer, + serverSettingsLayer, + // Replay threads have no temporary worktree branch to name. + Layer.mock(TemporaryBranchRename.TemporaryBranchRename)({ rename: () => Effect.void }), ), ), ); diff --git a/apps/server/src/orchestration-v2/testkit/ProviderTurnStartMemory.fixture.mjs b/apps/server/src/orchestration-v2/testkit/ProviderTurnStartMemory.fixture.mjs index 36d5ce1adb1b..5e870abafd9d 100644 --- a/apps/server/src/orchestration-v2/testkit/ProviderTurnStartMemory.fixture.mjs +++ b/apps/server/src/orchestration-v2/testkit/ProviderTurnStartMemory.fixture.mjs @@ -13,20 +13,35 @@ const [Effect, Layer, FileSystem] = await Promise.all([ load("FileSystem"), ]); const app = (file) => import(NodeURL.pathToFileURL(root + "/apps/server/src/" + file + ".ts")); -const [Start, Projection, Run, Sessions, Policy, Id, Sink, Handoff, Git, Project, Auth] = - await Promise.all([ - app("orchestration-v2/ProviderTurnStartService"), - app("orchestration-v2/ProjectionStore"), - app("orchestration-v2/RunExecutionService"), - app("orchestration-v2/ProviderSessionManager"), - app("orchestration-v2/RuntimePolicy"), - app("orchestration-v2/IdAllocator"), - app("orchestration-v2/EventSink"), - app("orchestration-v2/ContextHandoffService"), - app("git/GitWorkflowService"), - app("project/ProjectService"), - app("provider/Services/ProviderAuthService"), - ]); +const [ + Start, + Projection, + Run, + Sessions, + Policy, + Id, + Sink, + Handoff, + Git, + Project, + Auth, + Settings, + BranchRename, +] = await Promise.all([ + app("orchestration-v2/ProviderTurnStartService"), + app("orchestration-v2/ProjectionStore"), + app("orchestration-v2/RunExecutionService"), + app("orchestration-v2/ProviderSessionManager"), + app("orchestration-v2/RuntimePolicy"), + app("orchestration-v2/IdAllocator"), + app("orchestration-v2/EventSink"), + app("orchestration-v2/ContextHandoffService"), + app("git/GitWorkflowService"), + app("project/ProjectService"), + app("provider/Services/ProviderAuthService"), + app("serverSettings"), + app("orchestration-v2/TemporaryBranchRename"), +]); let current; let fullReads = 0; const liveRuns = []; @@ -49,6 +64,8 @@ const dependencies = Layer.mergeAll( Layer.mock(Git.GitWorkflowService)({}), Layer.mock(Project.ProjectService)({}), Layer.mock(Auth.ProviderAuthService)({}), + Layer.mock(Settings.ServerSettingsService)({}), + Layer.mock(BranchRename.TemporaryBranchRename)({ rename: () => Effect.void }), Layer.mock(Projection.ProjectionStoreV2)({ getThreadProjection: () => Effect.sync(() => { diff --git a/apps/server/src/orchestration-v2/workflowScriptQuery.test.ts b/apps/server/src/orchestration-v2/workflowScriptQuery.test.ts index 311b89427697..8a40ee8fc130 100644 --- a/apps/server/src/orchestration-v2/workflowScriptQuery.test.ts +++ b/apps/server/src/orchestration-v2/workflowScriptQuery.test.ts @@ -2,23 +2,30 @@ import * as NodeFS from "node:fs"; import * as NodeOS from "node:os"; import * as NodePath from "node:path"; +import * as NodeServices from "@effect/platform-node/NodeServices"; import { it as effectIt } from "@effect/vitest"; import * as Effect from "effect/Effect"; -import { afterAll, assert, describe } from "vite-plus/test"; +import * as Layer from "effect/Layer"; +import { afterAll, assert } from "vite-plus/test"; +import { ProviderDriverKind, ProviderInstanceId } from "@t3tools/contracts"; import { symlinksSupported } from "@t3tools/shared/testing/symlinks"; +import * as ServerSettings from "../serverSettings.ts"; import { readWorkflowScript } from "./workflowScriptQuery.ts"; -const root = NodePath.join(NodeOS.homedir(), ".claude", "projects", "__wf_script_test__"); +// A Claude instance with its own home: scripts live under /projects, +// never the real ~/.claude/projects. +const sandbox = NodeFS.mkdtempSync(NodePath.join(NodeOS.tmpdir(), "wf-script-test-")); +const claudeHome = NodePath.join(sandbox, "claude-home"); +const root = NodePath.join(claudeHome, "projects", "-repo"); NodeFS.mkdirSync(root, { recursive: true }); const scriptPath = NodePath.join(root, "run.js"); NodeFS.writeFileSync(scriptPath, "export const meta = {};\n"); -const outside = NodePath.join(NodeOS.tmpdir(), "wf-outside.js"); +const outside = NodePath.join(sandbox, "wf-outside.js"); NodeFS.writeFileSync(outside, "evil\n"); const link = NodePath.join(root, "sneaky.js"); // Planted only where the host allows it; the escape test is skipped // otherwise rather than passing vacuously on "not-found". if (symlinksSupported) { - NodeFS.rmSync(link, { force: true }); NodeFS.symlinkSync(outside, link); if (!NodeFS.lstatSync(link).isSymbolicLink()) { throw new Error("test setup: sneaky.js must be a symlink"); @@ -26,12 +33,23 @@ if (symlinksSupported) { } afterAll(() => { - NodeFS.rmSync(root, { recursive: true, force: true }); - NodeFS.rmSync(outside, { force: true }); + NodeFS.rmSync(sandbox, { recursive: true, force: true }); }); -describe("readWorkflowScript containment", () => { - effectIt.effect("serves a real script under the projects root", () => +const testLayer = Layer.mergeAll( + NodeServices.layer, + ServerSettings.ServerSettingsService.layerTest({ + providerInstances: { + [ProviderInstanceId.make("claudeAgent")]: { + driver: ProviderDriverKind.make("claudeAgent"), + config: { homePath: claudeHome }, + }, + }, + }), +); + +effectIt.layer(testLayer)("readWorkflowScript containment", (it) => { + it.effect("serves a real script under the projects root", () => Effect.gen(function* () { const result = yield* readWorkflowScript({ scriptPath }); assert.include(result.contents, "export const meta"); @@ -39,7 +57,7 @@ describe("readWorkflowScript containment", () => { }), ); - effectIt.effect("rejects relative and non-js paths", () => + it.effect("rejects relative and non-js paths", () => Effect.gen(function* () { const relative = yield* Effect.exit(readWorkflowScript({ scriptPath: "run.js" })); assert.equal(relative._tag, "Failure"); @@ -50,25 +68,23 @@ describe("readWorkflowScript containment", () => { }), ); - effectIt.effect.skipIf(!symlinksSupported)( - "rejects paths outside the root and symlink escapes", - () => - Effect.gen(function* () { - const escaped = yield* Effect.exit(readWorkflowScript({ scriptPath: outside })); - assert.equal(escaped._tag, "Failure"); - // A symlink INSIDE the root pointing outside must fail specifically on - // realpath re-containment — a "not-found" would mean the link was - // never exercised and the assertion proves nothing. - const sneaky = yield* Effect.exit( - readWorkflowScript({ scriptPath: link }).pipe( - Effect.flip, - Effect.map((error) => error.reason), - ), - ); - assert.equal(sneaky._tag, "Success"); - if (sneaky._tag === "Success") { - assert.equal(sneaky.value, "outside-root"); - } - }), + it.effect.skipIf(!symlinksSupported)("rejects paths outside the root and symlink escapes", () => + Effect.gen(function* () { + const escaped = yield* Effect.exit(readWorkflowScript({ scriptPath: outside })); + assert.equal(escaped._tag, "Failure"); + // A symlink INSIDE the root pointing outside must fail specifically on + // realpath re-containment — a "not-found" would mean the link was + // never exercised and the assertion proves nothing. + const sneaky = yield* Effect.exit( + readWorkflowScript({ scriptPath: link }).pipe( + Effect.flip, + Effect.map((error) => error.reason), + ), + ); + assert.equal(sneaky._tag, "Success"); + if (sneaky._tag === "Success") { + assert.equal(sneaky.value, "outside-root"); + } + }), ); }); diff --git a/apps/server/src/orchestration-v2/workflowScriptQuery.ts b/apps/server/src/orchestration-v2/workflowScriptQuery.ts index 184abc6147ce..214821c399b7 100644 --- a/apps/server/src/orchestration-v2/workflowScriptQuery.ts +++ b/apps/server/src/orchestration-v2/workflowScriptQuery.ts @@ -4,9 +4,10 @@ * "{} script" affordance. * * Containment rules (lifted from the reviewed #3650 inspection service): - * - the resolved realpath must live under ~/.claude/projects (where the - * Claude harness persists workflow scripts) — realpath re-containment - * defeats symlink escapes, including a symlinked leaf file; + * - the resolved realpath must live under the `projects` dir of a configured + * Claude instance (where the Claude harness persists workflow scripts) — + * realpath re-containment defeats symlink escapes, including a symlinked + * leaf file; * - only .js leaf files are served; * - reads are size-capped rather than failed, with a truncation marker. * @@ -14,17 +15,42 @@ * never trusted beyond these checks. */ import * as NodeFSP from "node:fs/promises"; -import * as NodeOS from "node:os"; import * as NodePath from "node:path"; -import { OrchestrationGetWorkflowScriptError } from "@t3tools/contracts"; +import { ClaudeSettings, OrchestrationGetWorkflowScriptError } from "@t3tools/contracts"; +import { HostProcessEnvironment } from "@t3tools/shared/hostProcess"; import * as Effect from "effect/Effect"; +import * as Option from "effect/Option"; +import * as Schema from "effect/Schema"; + +import { resolveClaudeHomePath } from "../provider/Drivers/ClaudeHome.ts"; +import { deriveProviderInstanceConfigMap } from "../provider/Layers/ProviderInstanceRegistryHydration.ts"; +import { mergeProviderInstanceEnvironment } from "../provider/ProviderInstanceEnvironment.ts"; +import * as ServerSettings from "../serverSettings.ts"; const SCRIPT_BYTE_CAP = 256 * 1024; +const decodeClaudeSettings = Schema.decodeUnknownOption(ClaudeSettings); -function scriptsRoot(): string { - return NodePath.join(NodeOS.homedir(), ".claude", "projects"); -} +// Each Claude instance can set its own home or CLAUDE_CONFIG_DIR, and a +// thread's workflow may come from any of them (subagents included). +const claudeProjectsRoots = Effect.fn("orchestration.claudeProjectsRoots")(function* () { + const settings = yield* ServerSettings.ServerSettingsService.pipe( + Effect.flatMap((service) => service.getSettings), + ); + const hostEnvironment = yield* HostProcessEnvironment; + const roots = new Set(); + for (const instance of Object.values(deriveProviderInstanceConfigMap(settings))) { + if (instance.driver !== "claudeAgent") continue; + const config = decodeClaudeSettings(instance.config ?? {}); + if (Option.isNone(config)) continue; + const home = yield* resolveClaudeHomePath( + config.value, + mergeProviderInstanceEnvironment(instance.environment, hostEnvironment), + ); + roots.add(NodePath.join(home, "projects")); + } + return [...roots]; +}); export const readWorkflowScript = Effect.fn("orchestration.readWorkflowScript")(function* (input: { readonly scriptPath: string; @@ -38,15 +64,28 @@ export const readWorkflowScript = Effect.fn("orchestration.readWorkflowScript")( }); } - const root = yield* Effect.tryPromise({ - try: () => NodeFSP.realpath(scriptsRoot()), - catch: (cause) => - new OrchestrationGetWorkflowScriptError({ - reason: "root-unavailable", - scriptPath: requested, - cause, - }), - }); + const roots = yield* claudeProjectsRoots().pipe( + Effect.flatMap((candidates) => + Effect.promise(() => + Promise.all(candidates.map((root) => NodeFSP.realpath(root).catch(() => null))), + ), + ), + Effect.map((resolvedRoots) => resolvedRoots.filter((root) => root !== null)), + Effect.mapError( + (cause) => + new OrchestrationGetWorkflowScriptError({ + reason: "root-unavailable", + scriptPath: requested, + cause, + }), + ), + ); + if (roots.length === 0) { + return yield* new OrchestrationGetWorkflowScriptError({ + reason: "root-unavailable", + scriptPath: requested, + }); + } // Realpath the FILE itself (not just its directory): a symlink named // like a script inside a contained directory must not escape. @@ -60,7 +99,7 @@ export const readWorkflowScript = Effect.fn("orchestration.readWorkflowScript")( }), }); - if (resolved !== root && !resolved.startsWith(`${root}${NodePath.sep}`)) { + if (!roots.some((root) => resolved === root || resolved.startsWith(`${root}${NodePath.sep}`))) { return yield* new OrchestrationGetWorkflowScriptError({ reason: "outside-root", scriptPath: resolved, diff --git a/apps/server/src/provider/ClaudeModelCatalog.ts b/apps/server/src/provider/ClaudeModelCatalog.ts index 1fb80cb4a57b..e601e9dc9392 100644 --- a/apps/server/src/provider/ClaudeModelCatalog.ts +++ b/apps/server/src/provider/ClaudeModelCatalog.ts @@ -215,7 +215,7 @@ export function isClaudeCatalogUltracodeEffort(effort: string | null | undefined return effort === "ultracode"; } -export function resolveClaudeCatalogContextWindow( +function resolveClaudeCatalogContextWindow( catalog: ClaudeModelCatalog, modelSelection: ModelSelection | undefined, ): string | undefined { diff --git a/apps/server/src/restartingWatch.test.ts b/apps/server/src/restartingWatch.test.ts new file mode 100644 index 000000000000..b3bbe71ecf66 --- /dev/null +++ b/apps/server/src/restartingWatch.test.ts @@ -0,0 +1,38 @@ +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 Ref from "effect/Ref"; +import * as TestClock from "effect/testing/TestClock"; + +import { keepWatching } from "./restartingWatch.ts"; + +it.effect("restarts a failed watcher and revalidates before it watches again", () => + Effect.gen(function* () { + const events = yield* Ref.make>([]); + const record = (event: string) => Ref.update(events, (all) => [...all, event]); + const secondWatchStarted = yield* Deferred.make(); + const watchStarts = yield* Ref.make(0); + const watch = Ref.updateAndGet(watchStarts, (count) => count + 1).pipe( + Effect.flatMap((count) => + count === 1 + ? record("watch").pipe(Effect.andThen(Effect.fail("watch error"))) + : record("watch").pipe( + Effect.andThen(Deferred.succeed(secondWatchStarted, undefined)), + Effect.andThen(Effect.never), + ), + ), + ); + + const fiber = yield* keepWatching({ + label: "Test", + watch, + revalidate: record("revalidate"), + }).pipe(Effect.forkChild); + yield* TestClock.adjust("30 seconds"); + yield* Deferred.await(secondWatchStarted); + + assert.deepStrictEqual(yield* Ref.get(events), ["watch", "revalidate", "watch"]); + yield* Fiber.interrupt(fiber); + }), +); diff --git a/apps/server/src/restartingWatch.ts b/apps/server/src/restartingWatch.ts new file mode 100644 index 000000000000..e32b9aea3df3 --- /dev/null +++ b/apps/server/src/restartingWatch.ts @@ -0,0 +1,37 @@ +import * as Cause from "effect/Cause"; +import * as Duration from "effect/Duration"; +import * as Effect from "effect/Effect"; +import * as Ref from "effect/Ref"; +import * as Schedule from "effect/Schedule"; + +// A watcher that keeps failing retries every 30 seconds at most. +const restartSchedule = Schedule.exponential("500 millis").pipe( + Schedule.modifyDelay(({ duration }) => + Effect.succeed(Duration.min(duration, Duration.seconds(30))), + ), +); + +/** + * Keeps a file watcher running for the life of the calling fiber. One watch + * error or a closed watch stream would otherwise stop reloads until the server + * restarts, so the watcher restarts with capped backoff. Each restart runs + * `revalidate` first, because changes made while the watcher was down were + * missed. + */ +export const keepWatching = (input: { + readonly label: string; + readonly watch: Effect.Effect; + readonly revalidate: Effect.Effect; +}) => + Effect.gen(function* () { + const restarted = yield* Ref.make(false); + const runOnce = Ref.getAndSet(restarted, true).pipe( + Effect.flatMap((isRestart) => (isRestart ? input.revalidate : Effect.void)), + Effect.andThen(input.watch), + Effect.catchCauseIf( + (cause) => !Cause.hasInterrupts(cause), + (cause) => Effect.logWarning(`${input.label} watcher failed; restarting`, { cause }), + ), + ); + return yield* Effect.repeat(runOnce, restartSchedule); + }); diff --git a/apps/server/src/serverSettings.ts b/apps/server/src/serverSettings.ts index 372506e0e508..9c611a24ac06 100644 --- a/apps/server/src/serverSettings.ts +++ b/apps/server/src/serverSettings.ts @@ -51,6 +51,7 @@ import * as Semaphore from "effect/Semaphore"; import * as Scope from "effect/Scope"; import * as Stream from "effect/Stream"; import * as SqlClient from "effect/unstable/sql/SqlClient"; +import { keepWatching } from "./restartingWatch.ts"; import { writeFileStringAtomically } from "./atomicWrite.ts"; import * as ServerConfig from "./config.ts"; import { type DeepPartial, deepMerge } from "@t3tools/shared/Struct"; @@ -1213,11 +1214,11 @@ const make = Effect.gen(function* () { Stream.debounce(Duration.millis(100)), ); - yield* Stream.runForEach(debouncedSettingsEvents, () => revalidateAndEmitSafely).pipe( - Effect.ignoreCause({ log: true }), - Effect.forkIn(watcherScope), - Effect.asVoid, - ); + yield* keepWatching({ + label: "Settings", + watch: Stream.runForEach(debouncedSettingsEvents, () => revalidateAndEmitSafely), + revalidate: revalidateAndEmitSafely, + }).pipe(Effect.forkIn(watcherScope), Effect.asVoid); }); const start = Effect.gen(function* () { diff --git a/apps/server/src/ws.ts b/apps/server/src/ws.ts index 2a0006ddb24b..5cb1c64cbc1e 100644 --- a/apps/server/src/ws.ts +++ b/apps/server/src/ws.ts @@ -2388,10 +2388,22 @@ const makeWsRpcLayer = ( providerAuth.complete(input, currentSessionId), { "rpc.aggregate": "provider" }, ), - [WS_METHODS.chatGptReconnectProfile]: (input) => providerAuth.reconnectProfile(input), - [WS_METHODS.chatGptImportProfile]: (input) => providerAuth.importProfile(input), + [WS_METHODS.chatGptReconnectProfile]: (input) => + observeRpcEffect( + WS_METHODS.chatGptReconnectProfile, + providerAuth.reconnectProfile(input), + { "rpc.aggregate": "provider" }, + ), + [WS_METHODS.chatGptImportProfile]: (input) => + observeRpcEffect(WS_METHODS.chatGptImportProfile, providerAuth.importProfile(input), { + "rpc.aggregate": "provider", + }), [WS_METHODS.chatGptHandoffSubscribe]: (input) => - subscribeChatGptHandoff(input, currentSessionId), + observeRpcStream( + WS_METHODS.chatGptHandoffSubscribe, + subscribeChatGptHandoff(input, currentSessionId), + { "rpc.aggregate": "provider" }, + ), [WS_METHODS.codexAuthCallbackSubscribe]: (input) => observeRpcStream( WS_METHODS.codexAuthCallbackSubscribe, @@ -3718,12 +3730,39 @@ const makeWsRpcLayer = ( }), ); +// Built once per server, not per connection: every client shares the provider +// update guard, the agent session scan, and the source-control discovery state. +const sharedWsRpcServicesLayer = Layer.mergeAll( + AgentSessionScanner.layer, + ProviderMaintenanceRunner.layer, + SourceControlDiscovery.layer.pipe( + Layer.provide( + SourceControlProviderRegistry.layer.pipe( + Layer.provide( + Layer.mergeAll( + AzureDevOpsCli.layer, + BitbucketApi.layer, + GitHubCli.layer, + GitLabCli.layer, + ForgejoCli.layer, + ), + ), + Layer.provideMerge(GitVcsDriver.layer), + Layer.provide(VcsDriverRegistry.layer.pipe(Layer.provide(VcsProjectConfig.layer))), + ), + ), + ), +); + export const websocketRpcRouteLayer = Layer.unwrap( Effect.gen(function* () { const previewAutomationBroker = yield* PreviewAutomationBroker.PreviewAutomationBroker; const serverSelfUpdate = yield* ServerSelfUpdate.ServerSelfUpdate; const pullRequests = yield* PullRequestService.PullRequestService; const sql = yield* SqlClient.SqlClient; + const agentSessionScanner = yield* AgentSessionScanner.AgentSessionScanner; + const providerMaintenance = yield* ProviderMaintenanceRunner.ProviderMaintenanceRunner; + const sourceControlDiscovery = yield* SourceControlDiscovery.SourceControlDiscovery; return HttpRouter.add( "GET", "/ws", @@ -3776,31 +3815,23 @@ export const websocketRpcRouteLayer = Layer.unwrap( ).pipe( Layer.provideMerge(RpcSerialization.layerJson), Layer.provide(Layer.succeed(SqlClient.SqlClient, sql)), - Layer.provide(AgentSessionScanner.layer), - Layer.provide(ProviderMaintenanceRunner.layer), + Layer.provide( + Layer.succeed(AgentSessionScanner.AgentSessionScanner, agentSessionScanner), + ), + Layer.provide( + Layer.succeed( + ProviderMaintenanceRunner.ProviderMaintenanceRunner, + providerMaintenance, + ), + ), Layer.provide(Layer.succeed(ServerSelfUpdate.ServerSelfUpdate, serverSelfUpdate)), // One server-lifetime service means clients share the same PR caches, and a WS // mutation invalidates the HTTP diff cache that every client reads from. Layer.provide(Layer.succeed(PullRequestService.PullRequestService, pullRequests)), Layer.provide( - SourceControlDiscovery.layer.pipe( - Layer.provide( - SourceControlProviderRegistry.layer.pipe( - Layer.provide( - Layer.mergeAll( - AzureDevOpsCli.layer, - BitbucketApi.layer, - GitHubCli.layer, - GitLabCli.layer, - ForgejoCli.layer, - ), - ), - Layer.provideMerge(GitVcsDriver.layer), - Layer.provide( - VcsDriverRegistry.layer.pipe(Layer.provide(VcsProjectConfig.layer)), - ), - ), - ), + Layer.succeed( + SourceControlDiscovery.SourceControlDiscovery, + sourceControlDiscovery, ), ), ), @@ -3819,4 +3850,4 @@ export const websocketRpcRouteLayer = Layer.unwrap( ), ); }), -); +).pipe(Layer.provide(sharedWsRpcServicesLayer)); diff --git a/packages/contracts/src/orchestrationV2.ts b/packages/contracts/src/orchestrationV2.ts index 1820008b77dd..1c8d41bfc0db 100644 --- a/packages/contracts/src/orchestrationV2.ts +++ b/packages/contracts/src/orchestrationV2.ts @@ -2572,6 +2572,8 @@ export const OrchestrationV2Command = Schema.Union([ branch: Schema.optional(Schema.NullOr(TrimmedNonEmptyString)), worktreePath: Schema.optional(Schema.NullOr(TrimmedNonEmptyString)), expectedWorktreePath: Schema.optional(Schema.NullOr(TrimmedNonEmptyString)), + /** Reject unless the thread title still equals this, e.g. after a user rename. */ + expectedTitle: Schema.optional(Schema.String), /** Reject unless no message or run has landed on this thread. */ expectedEmpty: Schema.optional(Schema.Boolean), limitRecovery: Schema.optional(Schema.NullOr(OrchestrationV2LimitRecoveryUpdate)),