From f581be4b3042104238f8c671bc61b8e593fc84a0 Mon Sep 17 00:00:00 2001 From: Julius Marminge <51714798+juliusmarminge@users.noreply.github.com> Date: Tue, 6 Oct 2026 23:49:38 -0700 Subject: [PATCH] feat(server): hand a thread off to a linked environment A thread can now move to a linked environment with its conversation and its git work, with exactly one copy live. ThreadHandoff marks it departing, which refuses new turns and keeps queued runs from starting; packs the branch and working tree with HandoffGit; uploads the bundle through the other side's signed attachment route; and calls t3_thread_import there, which now takes the bundle and applies it in a new worktree before creating the thread. Only then is the thread departed here and read-only. A failed move marks it failed, which takes turns again, and nothing is left there. An agent can move its own thread with t3_thread_handoff. The move then waits as pending until its turn ends, starts the thread there with the agent's continuationPrompt, and is cancelled by a user message sent before the turn ends. A startup sweep runs moves that a restart cut short; import ids derive from the handoff, so a retry finds the thread already there. threadHandoff.options / start / cancel serve the clients. Co-Authored-By: Claude Opus 5.5 (1M context) --- apps/server/src/auth/RpcAuthorization.ts | 4 + ...OrchestratorMcpToolkit.integration.test.ts | 3 + apps/server/src/mcp/linkOrigin.test.ts | 23 +- apps/server/src/mcp/toolkits/core.test.ts | 5 + .../src/mcp/toolkits/orchestrator/handlers.ts | 45 ++ .../src/mcp/toolkits/orchestrator/tools.ts | 17 + .../src/mcp/toolkits/project/handlers.ts | 43 +- apps/server/src/mcp/toolkits/project/tools.ts | 2 + .../toolkits/worktree/registration.test.ts | 4 + .../src/observability/RpcInstrumentation.ts | 3 + .../src/orchestration-v2/Orchestrator.ts | 88 +++ .../src/orchestration-v2/ProjectionStore.ts | 2 + .../orchestration-v2/ThreadImportService.ts | 23 +- .../orchestration-v2/ThreadMessageIntake.ts | 1 + .../testkit/CapturingCodexAdapter.ts | 155 +++-- apps/server/src/peer/PeerForwarding.test.ts | 27 +- apps/server/src/peer/PeerLinkRequests.test.ts | 2 + apps/server/src/peer/PeerLinks.testkit.ts | 21 +- apps/server/src/peer/RemoteDelegation.test.ts | 5 + apps/server/src/peer/handoff/HandoffImport.ts | 131 ++++ .../src/peer/handoff/ThreadHandoff.test.ts | 619 ++++++++++++++++++ apps/server/src/peer/handoff/ThreadHandoff.ts | 577 ++++++++++++++++ apps/server/src/server.ts | 15 + apps/server/src/ws.ts | 18 + docs/internals/remote.md | 10 + docs/user/remote-access.md | 15 + packages/contracts/src/orchestrationV2.ts | 53 ++ packages/contracts/src/orchestratorMcp.ts | 42 ++ packages/contracts/src/peerLink.ts | 13 + packages/contracts/src/rpc.ts | 56 +- packages/shared/src/t3McpToolPresentation.ts | 4 + 31 files changed, 1929 insertions(+), 97 deletions(-) create mode 100644 apps/server/src/peer/handoff/HandoffImport.ts create mode 100644 apps/server/src/peer/handoff/ThreadHandoff.test.ts create mode 100644 apps/server/src/peer/handoff/ThreadHandoff.ts diff --git a/apps/server/src/auth/RpcAuthorization.ts b/apps/server/src/auth/RpcAuthorization.ts index 2285c580ee7f..17cb603ffc4c 100644 --- a/apps/server/src/auth/RpcAuthorization.ts +++ b/apps/server/src/auth/RpcAuthorization.ts @@ -83,6 +83,10 @@ export const RPC_REQUIRED_SCOPES = { [WS_METHODS.serverRemoveKeybinding]: AuthSettingsWriteScope, [WS_METHODS.serverGetSettings]: AuthOrchestrationReadScope, [WS_METHODS.peerLinksList]: AuthAccessReadScope, + // Moving a thread drives the linked environment through this one's link. + [WS_METHODS.threadHandoffOptions]: AuthOrchestrationReadScope, + [WS_METHODS.threadHandoffStart]: AuthOrchestrationOperateScope, + [WS_METHODS.threadHandoffCancel]: AuthOrchestrationOperateScope, [WS_METHODS.serverUpdateSettings]: AuthSettingsWriteScope, [WS_METHODS.serverSearchAcpRegistry]: AuthOrchestrationReadScope, [WS_METHODS.serverPrepareAcpRegistryAgent]: AuthProvidersManageScope, diff --git a/apps/server/src/mcp/OrchestratorMcpToolkit.integration.test.ts b/apps/server/src/mcp/OrchestratorMcpToolkit.integration.test.ts index 0fbb30778fcf..4cefa3d91058 100644 --- a/apps/server/src/mcp/OrchestratorMcpToolkit.integration.test.ts +++ b/apps/server/src/mcp/OrchestratorMcpToolkit.integration.test.ts @@ -78,6 +78,7 @@ import * as McpInvocationContext from "./McpInvocationContext.ts"; import * as PeerForwarding from "../peer/PeerForwarding.ts"; import * as PeerLinkRequests from "../peer/PeerLinkRequests.ts"; import * as PeerLinks from "../peer/PeerLinks.ts"; +import * as ThreadHandoff from "../peer/handoff/ThreadHandoff.ts"; import { delegatedTaskRun, hasPendingChildRuns } from "./OrchestratorMcpService.ts"; import * as ProviderAdapter from "@t3tools/provider-core/server/ProviderAdapter"; @@ -694,6 +695,7 @@ describe("orchestrator MCP toolkit", () => { Layer.provide(Layer.mock(PeerForwarding.PeerForwarding)({})), Layer.provide(Layer.mock(PeerLinkRequests.PeerLinkRequests)({})), Layer.provide(Layer.mock(PeerLinks.PeerLinks)({})), + Layer.provide(Layer.mock(ThreadHandoff.ThreadHandoff)({})), Layer.provide(Layer.mock(ThreadSearch.ThreadSearch)({})), Layer.provide( Layer.mock(ProjectService.ProjectService)({ @@ -3850,6 +3852,7 @@ describe("orchestrator MCP toolkit", () => { Layer.provide(Layer.mock(PeerForwarding.PeerForwarding)({})), Layer.provide(Layer.mock(PeerLinkRequests.PeerLinkRequests)({})), Layer.provide(Layer.mock(PeerLinks.PeerLinks)({})), + Layer.provide(Layer.mock(ThreadHandoff.ThreadHandoff)({})), Layer.provide(Layer.mock(ProjectService.ProjectService)({})), Layer.provideMerge( SecretRequests.layer.pipe( diff --git a/apps/server/src/mcp/linkOrigin.test.ts b/apps/server/src/mcp/linkOrigin.test.ts index 289451306051..5299f46dec75 100644 --- a/apps/server/src/mcp/linkOrigin.test.ts +++ b/apps/server/src/mcp/linkOrigin.test.ts @@ -44,6 +44,8 @@ import * as ThreadSearch from "../orchestration-v2/ThreadSearch.ts"; import * as PeerLinkRequests from "../peer/PeerLinkRequests.ts"; import * as PeerLinks from "../peer/PeerLinks.ts"; import * as ThreadImportService from "../orchestration-v2/ThreadImportService.ts"; +import * as ThreadHandoff from "../peer/handoff/ThreadHandoff.ts"; +import * as HandoffImport from "../peer/handoff/HandoffImport.ts"; // A linked environment's session drives this one's real orchestrator through // its real T3 tools. What the link starts carries its origin, and the link @@ -130,13 +132,20 @@ const layerTools = Layer.mergeAll( update: () => Effect.die("A linked caller must never reach a project update."), }), ), - Layer.provide(Layer.mock(SecretRequests.SecretRequests)({})), - Layer.provide(Layer.mock(PeerForwarding.PeerForwarding)({})), - Layer.provide(Layer.mock(PeerLinkRequests.PeerLinkRequests)({})), - Layer.provide(Layer.mock(PeerLinks.PeerLinks)({})), - Layer.provide(Layer.mock(ThreadImportService.ThreadImportService)({})), - Layer.provide(Layer.mock(ThreadSearch.ThreadSearch)({})), - Layer.provide(Layer.mock(RemoteDelegation.RemoteDelegation)({})), + // Services the toolkits declare that these calls never reach. + Layer.provide( + Layer.mergeAll( + Layer.mock(SecretRequests.SecretRequests)({}), + Layer.mock(PeerForwarding.PeerForwarding)({}), + Layer.mock(PeerLinkRequests.PeerLinkRequests)({}), + Layer.mock(PeerLinks.PeerLinks)({}), + Layer.mock(ThreadImportService.ThreadImportService)({}), + Layer.mock(ThreadHandoff.ThreadHandoff)({}), + Layer.mock(HandoffImport.HandoffImport)({}), + Layer.mock(ThreadSearch.ThreadSearch)({}), + Layer.mock(RemoteDelegation.RemoteDelegation)({}), + ), + ), Layer.provide( Layer.mock(ManagedProjectFolders.ManagedProjectFolders)({ namedProjectsRoot: "/p" }), ), diff --git a/apps/server/src/mcp/toolkits/core.test.ts b/apps/server/src/mcp/toolkits/core.test.ts index 063e4f352728..3f2cba8cc638 100644 --- a/apps/server/src/mcp/toolkits/core.test.ts +++ b/apps/server/src/mcp/toolkits/core.test.ts @@ -53,6 +53,7 @@ import * as AttachmentHandlers from "./attachment/handlers.ts"; import * as PeerForwarding from "../../peer/PeerForwarding.ts"; import * as PeerLinkRequests from "../../peer/PeerLinkRequests.ts"; import * as PeerLinks from "../../peer/PeerLinks.ts"; +import * as ThreadHandoff from "../../peer/handoff/ThreadHandoff.ts"; import { ThreadToolkit } from "./thread/tools.ts"; import { WorktreeToolkit } from "./worktree/tools.ts"; import { DeviceToolkit } from "./device/tools.ts"; @@ -81,6 +82,7 @@ const layerThreadToolkit = McpHttpServer.layerThreadToolkit.pipe( Layer.provide(Layer.mock(PeerForwarding.PeerForwarding)({})), Layer.provide(Layer.mock(PeerLinkRequests.PeerLinkRequests)({})), Layer.provide(Layer.mock(PeerLinks.PeerLinks)({})), + Layer.provide(Layer.mock(ThreadHandoff.ThreadHandoff)({})), ); it("publishes unique tool names with reference-free object-root inputs", () => { @@ -588,6 +590,7 @@ it.effect("refuses act-as-caller tools to a client caller", () => Layer.provide(Layer.mock(PeerForwarding.PeerForwarding)({})), Layer.provide(Layer.mock(PeerLinkRequests.PeerLinkRequests)({})), Layer.provide(Layer.mock(PeerLinks.PeerLinks)({})), + Layer.provide(Layer.mock(ThreadHandoff.ThreadHandoff)({})), Layer.provide(Layer.mock(ProjectService.ProjectService)({})), Layer.provide(Layer.mock(SecretRequests.SecretRequests)({})), Layer.provide(Layer.mock(RemoteDelegation.RemoteDelegation)({})), @@ -646,6 +649,7 @@ it.effect("a caller cannot rewrite a scheduled task that runs above its own mode Layer.provide(Layer.mock(PeerForwarding.PeerForwarding)({})), Layer.provide(Layer.mock(PeerLinkRequests.PeerLinkRequests)({})), Layer.provide(Layer.mock(PeerLinks.PeerLinks)({})), + Layer.provide(Layer.mock(ThreadHandoff.ThreadHandoff)({})), Layer.provide(Layer.mock(RemoteDelegation.RemoteDelegation)({})), ), ), @@ -717,6 +721,7 @@ it.effect("a caller cannot interrupt a thread that runs above its own modes", () Layer.provide(Layer.mock(PeerForwarding.PeerForwarding)({})), Layer.provide(Layer.mock(PeerLinkRequests.PeerLinkRequests)({})), Layer.provide(Layer.mock(PeerLinks.PeerLinks)({})), + Layer.provide(Layer.mock(ThreadHandoff.ThreadHandoff)({})), Layer.provide(Layer.mock(ProjectService.ProjectService)({})), Layer.provide(Layer.mock(SecretRequests.SecretRequests)({})), Layer.provide(Layer.mock(RemoteDelegation.RemoteDelegation)({})), diff --git a/apps/server/src/mcp/toolkits/orchestrator/handlers.ts b/apps/server/src/mcp/toolkits/orchestrator/handlers.ts index 42ec46b46ae1..0e4e9b60d101 100644 --- a/apps/server/src/mcp/toolkits/orchestrator/handlers.ts +++ b/apps/server/src/mcp/toolkits/orchestrator/handlers.ts @@ -7,6 +7,7 @@ import { OrchestratorToolkit } from "./tools.ts"; import * as PeerForwarding from "../../../peer/PeerForwarding.ts"; import * as PeerLinkRequests from "../../../peer/PeerLinkRequests.ts"; import * as PeerLinks from "../../../peer/PeerLinks.ts"; +import * as ThreadHandoff from "../../../peer/handoff/ThreadHandoff.ts"; import * as McpInvocationContext from "../../McpInvocationContext.ts"; import * as McpToolAccess from "../../McpToolAccess.ts"; import * as OrchestratorMcpService from "../../OrchestratorMcpService.ts"; @@ -173,6 +174,50 @@ const handlers = { return yield* service.waitForThread(scope, input); }), ), + t3_thread_handoff: McpToolAccess.writesThreads( + (input) => [input.threadId], + (input) => + Effect.gen(function* () { + const scope = yield* McpInvocationContext.McpInvocationContext; + const handoff = yield* ThreadHandoff.ThreadHandoff; + const state = yield* handoff + .start({ + threadId: input.threadId, + callerThreadId: scope.thread?.threadId, + environmentId: input.environmentId, + projectId: input.projectId, + continuationPrompt: input.continuationPrompt, + }) + .pipe( + Effect.mapError( + (error) => + new OrchestratorMcpFailure({ + code: + error.reason === "not_found" + ? "thread_not_found" + : error.reason === "refused" + ? "invalid_request" + : "orchestration_error", + message: error.message, + }), + ), + ); + return { + state: state.state, + environmentId: state.environmentId, + label: state.label, + threadId: state.state === "departed" ? state.threadId : null, + message: + state.state === "pending" + ? `This thread moves to ${state.label} when this turn ends. End the turn now; a message to this thread before then cancels the move.` + : state.state === "departed" + ? `The thread moved to ${state.label}.` + : state.state === "failed" + ? `The move failed: ${state.lastError}` + : `The thread is moving to ${state.label}.`, + }; + }), + ), t3_environment_links: McpToolAccess.reads(() => Effect.gen(function* () { const scope = yield* McpInvocationContext.McpInvocationContext; diff --git a/apps/server/src/mcp/toolkits/orchestrator/tools.ts b/apps/server/src/mcp/toolkits/orchestrator/tools.ts index 2504c72f6026..6faab0c64bf9 100644 --- a/apps/server/src/mcp/toolkits/orchestrator/tools.ts +++ b/apps/server/src/mcp/toolkits/orchestrator/tools.ts @@ -6,6 +6,8 @@ import { OrchestratorMcpEnvironmentLinksResult, OrchestratorMcpEnvironmentUnlinkInput, OrchestratorMcpEnvironmentUnlinkResult, + OrchestratorMcpThreadHandoffInput, + OrchestratorMcpThreadHandoffResult, OrchestratorMcpCreateThreadsInput, OrchestratorMcpCreateThreadsResult, OrchestratorMcpDelegateTaskInput, @@ -42,6 +44,7 @@ import * as ThreadManagementService from "../../../orchestration-v2/ThreadManage import * as PeerForwarding from "../../../peer/PeerForwarding.ts"; import * as PeerLinkRequests from "../../../peer/PeerLinkRequests.ts"; import * as PeerLinks from "../../../peer/PeerLinks.ts"; +import * as ThreadHandoff from "../../../peer/handoff/ThreadHandoff.ts"; import * as McpInvocationContext from "../../McpInvocationContext.ts"; import * as OrchestratorMcpService from "../../OrchestratorMcpService.ts"; import * as ThreadMetadataMcpService from "../../ThreadMetadataMcpService.ts"; @@ -274,6 +277,19 @@ const ThreadInterruptTool = Tool.make("t3_thread_interrupt", { .annotate(Tool.Title, "Interrupt a T3 thread") .annotate(Tool.Destructive, true); +const ThreadHandoffTool = Tool.make("t3_thread_handoff", { + description: + "Move a thread to a linked environment (t3_environment_links), such as the user's VPS, with its conversation and its git work: unpushed commits, uncommitted edits and new files. It continues there in a new worktree of the project with the same repository, and the copy here becomes read-only. Omit threadId to move the calling thread: the move then happens as soon as this turn ends, and continuationPrompt becomes its first message there, so end the turn after saying what you will do there. A message to the thread before then cancels the move. Another thread must be idle to move.", + parameters: OrchestratorMcpThreadHandoffInput, + success: OrchestratorMcpThreadHandoffResult, + failure: OrchestratorMcpFailure, + failureMode: "return", + dependencies: [...dependencies, ThreadHandoff.ThreadHandoff], +}) + .annotate(Tool.Title, "Move a thread to a linked environment") + .annotate(Tool.Destructive, true) + .annotate(Tool.OpenWorld, true); + const EnvironmentLinksTool = Tool.make("t3_environment_links", { description: "List the other T3 Code environments this one is linked to, such as the user's other machines, and whether each answers now. A machine the user names that is not listed is not linked yet: call t3_environment_link for it rather than looking elsewhere. Pass one's environmentId to t3_project_list, t3_thread_launch, t3_thread_list, t3_thread_read, t3_thread_send, t3_thread_wait, t3_thread_interrupt, or orchestrator_capabilities to act there. A linked environment lets this one's agents do at most its access there, and never more than the calling agent may do here. An expired link has to be linked again from a new pairing code.", @@ -333,4 +349,5 @@ export const OrchestratorToolkit = Toolkit.make( EnvironmentLinksTool, EnvironmentLinkTool, EnvironmentUnlinkTool, + ThreadHandoffTool, ); diff --git a/apps/server/src/mcp/toolkits/project/handlers.ts b/apps/server/src/mcp/toolkits/project/handlers.ts index 1c515f1d5161..5688c4355eb3 100644 --- a/apps/server/src/mcp/toolkits/project/handlers.ts +++ b/apps/server/src/mcp/toolkits/project/handlers.ts @@ -10,6 +10,7 @@ import * as Effect from "effect/Effect"; import * as FileSystem from "effect/FileSystem"; import * as Option from "effect/Option"; import * as ThreadImportService from "../../../orchestration-v2/ThreadImportService.ts"; +import * as HandoffImport from "../../../peer/handoff/HandoffImport.ts"; import * as ThreadMessageIntake from "../../../orchestration-v2/ThreadMessageIntake.ts"; import * as Claims from "../../../orchestration-v2/AttachmentClaims.ts"; import * as Project from "../../../project/ProjectService.ts"; @@ -210,21 +211,43 @@ export const layer = McpToolAccess.toLayer(ProjectToolkit, { code: "invalid_request", message: "The project was not found.", }); + if (input.worktreePath !== undefined && input.bundle !== undefined) + return yield* new OrchestratorMcpFailure({ + code: "invalid_request", + message: "A bundle lands in a new worktree; omit worktreePath.", + }); if (input.worktreePath !== undefined) yield* assertProjectWorktree(project.workspaceRoot, input.worktreePath); const imports = yield* ThreadImportService.ThreadImportService; + const handoffs = yield* HandoffImport.HandoffImport; + const bundle = input.bundle; + const workspace = + bundle === undefined + ? undefined + : handoffs.applyBundle({ + repoRoot: project.workspaceRoot, + handoffId: input.source.handoffId, + bundle, + }); return yield* imports - .importThread({ ...input, runtimeMode, interactionMode, linkOrigin }) + .importThread({ + ...input, + runtimeMode, + interactionMode, + linkOrigin, + ...(workspace === undefined ? {} : { workspace }), + }) .pipe( - Effect.mapError( - (error) => - new OrchestratorMcpFailure({ - code: - error._tag === "ThreadImportProjectMismatchError" - ? "invalid_request" - : "orchestration_error", - message: error.message, - }), + Effect.mapError((error) => + error._tag === "OrchestratorMcpFailure" + ? error + : new OrchestratorMcpFailure({ + code: + error._tag === "ThreadImportProjectMismatchError" + ? "invalid_request" + : "orchestration_error", + message: error.message, + }), ), ); }), diff --git a/apps/server/src/mcp/toolkits/project/tools.ts b/apps/server/src/mcp/toolkits/project/tools.ts index a18a0d08ac38..9a45e0e409a4 100644 --- a/apps/server/src/mcp/toolkits/project/tools.ts +++ b/apps/server/src/mcp/toolkits/project/tools.ts @@ -25,6 +25,7 @@ import { import * as FileSystem from "effect/FileSystem"; import * as ServerConfig from "../../../config.ts"; import * as ThreadImportService from "../../../orchestration-v2/ThreadImportService.ts"; +import * as HandoffImport from "../../../peer/handoff/HandoffImport.ts"; import * as ThreadLaunchService from "../../../orchestration-v2/ThreadLaunchService.ts"; import * as Crypto from "effect/Crypto"; import * as Schema from "effect/Schema"; @@ -172,6 +173,7 @@ const ThreadImportTool = Tool.make("t3_thread_import", { dependencies: [ ...shared.dependencies, ThreadImportService.ThreadImportService, + HandoffImport.HandoffImport, GitVcsDriver.GitVcsDriver, FileSystem.FileSystem, ], diff --git a/apps/server/src/mcp/toolkits/worktree/registration.test.ts b/apps/server/src/mcp/toolkits/worktree/registration.test.ts index a3ac46b90fb4..c58fc478c0a6 100644 --- a/apps/server/src/mcp/toolkits/worktree/registration.test.ts +++ b/apps/server/src/mcp/toolkits/worktree/registration.test.ts @@ -38,6 +38,8 @@ import * as PeerForwarding from "../../../peer/PeerForwarding.ts"; import * as PeerLinkRequests from "../../../peer/PeerLinkRequests.ts"; import * as PeerLinks from "../../../peer/PeerLinks.ts"; import * as ThreadImportService from "../../../orchestration-v2/ThreadImportService.ts"; +import * as ThreadHandoff from "../../../peer/handoff/ThreadHandoff.ts"; +import * as HandoffImport from "../../../peer/handoff/HandoffImport.ts"; const layerStubServices = Layer.mergeAll( Layer.mock(Orchestrator.OrchestratorV2)({}), @@ -51,6 +53,8 @@ const layerStubServices = Layer.mergeAll( Layer.mock(PeerLinkRequests.PeerLinkRequests)({}), Layer.mock(PeerLinks.PeerLinks)({}), Layer.mock(ThreadImportService.ThreadImportService)({}), + Layer.mock(ThreadHandoff.ThreadHandoff)({}), + Layer.mock(HandoffImport.HandoffImport)({}), Layer.mock(SecretRequests.SecretRequests)({}), Layer.mock(RemoteDelegation.RemoteDelegation)({}), Layer.mock(ProjectService.ProjectService)({}), diff --git a/apps/server/src/observability/RpcInstrumentation.ts b/apps/server/src/observability/RpcInstrumentation.ts index 83c1157fc233..cd323515c3a4 100644 --- a/apps/server/src/observability/RpcInstrumentation.ts +++ b/apps/server/src/observability/RpcInstrumentation.ts @@ -59,6 +59,9 @@ const RPC_AGGREGATES = { [WS_METHODS.peerLinksLink]: "peerLinks", [WS_METHODS.peerLinksUnlink]: "peerLinks", [WS_METHODS.peerLinksAnswerRequest]: "peerLinks", + [WS_METHODS.threadHandoffOptions]: "threadHandoff", + [WS_METHODS.threadHandoffStart]: "threadHandoff", + [WS_METHODS.threadHandoffCancel]: "threadHandoff", [WS_METHODS.serverUpdateSettings]: "server", [WS_METHODS.serverSearchAcpRegistry]: "server", [WS_METHODS.serverPrepareAcpRegistryAgent]: "server", diff --git a/apps/server/src/orchestration-v2/Orchestrator.ts b/apps/server/src/orchestration-v2/Orchestrator.ts index 8519397203f0..4860df99d199 100644 --- a/apps/server/src/orchestration-v2/Orchestrator.ts +++ b/apps/server/src/orchestration-v2/Orchestrator.ts @@ -196,6 +196,23 @@ export class OrchestratorSubagentThreadReadOnlyError extends Schema.TaggedError< } } +/** The thread is moving, or has moved, to a linked environment and takes no new turns here. */ +export class OrchestratorThreadMovedError extends Schema.TaggedError()( + "OrchestratorThreadMovedError", + { + commandId: CommandId, + threadId: ThreadId, + label: Schema.String, + state: Schema.Literals(["departing", "departed"]), + }, +) { + override get message(): string { + return this.state === "departed" + ? `This thread moved to ${this.label}. Continue it there.` + : `This thread is moving to ${this.label}.`; + } +} + /** The command's thread runs above the modes its sender may touch (see `DispatchModeLimit`). */ export class OrchestratorThreadAboveModeLimitError extends Schema.TaggedError()( "OrchestratorThreadAboveModeLimitError", @@ -261,6 +278,7 @@ export const OrchestratorV2Error = Schema.Union([ OrchestratorCommandPreviouslyRejectedError, OrchestratorCommandIdConflictError, OrchestratorSubagentThreadReadOnlyError, + OrchestratorThreadMovedError, OrchestratorThreadAboveModeLimitError, ]); export type OrchestratorV2Error = typeof OrchestratorV2Error.Type; @@ -472,6 +490,7 @@ function commandThreadId(command: OrchestrationV2ServerCommand): ThreadId { case "checkpoint.rollback.fail": case "thread.background-work.settle": case "thread.stop": + case "thread.handoff.update": case "provider.switch": return command.threadId; case "delegated_task.request": @@ -1275,6 +1294,8 @@ const makeOrchestrator = Effect.fn("orchestrationV2.Orchestrator.layer")(functio if ( projection.thread.archivedAt !== null || projection.thread.deletedAt !== null || + // A thread waiting to move, moving or moved starts nothing more here. + (projection.thread.handoff !== undefined && projection.thread.handoff.state !== "failed") || projection.runs.some(isBlockingRun) || projection.runs.some((run) => run.status === "queued" && run.queueHeld === true) ) { @@ -6747,6 +6768,61 @@ const makeOrchestrator = Effect.fn("orchestrationV2.Orchestrator.layer")(functio }, ); + /** Moves a thread through a handoff to a linked environment; see the command. */ + const dispatchThreadHandoffUpdate = Effect.fn("orchestrationV2.dispatch.threadHandoffUpdate")( + function* ( + command: Extract, + events: Ref.Ref>, + ) { + const projection = yield* projectionStore + .getThreadRecords(command.threadId, ["runs"]) + .pipe( + Effect.mapError( + (cause) => new OrchestratorProjectionError({ threadId: command.threadId, cause }), + ), + ); + const thread = projection.thread; + const reject = (cause: string) => + new OrchestratorDispatchError({ + commandId: command.commandId, + commandType: command.type, + cause, + }); + if (thread.deletedAt !== null) return yield* reject(`Thread ${thread.id} is deleted.`); + const current = thread.handoff?.state ?? null; + if (current !== command.expected) { + return yield* reject( + current === null + ? "This thread is not being moved." + : `This thread's move is ${current}, not ${command.expected ?? "absent"}.`, + ); + } + const next = command.handoff?.state ?? null; + if ( + next === "departing" && + projection.runs.some((run) => isBlockingRun(run) || run.status === "queued") + ) { + return yield* reject("A thread moves only while no turn runs or waits on it."); + } + const now = yield* DateTime.now; + const { handoff: _previous, ...rest } = thread; + yield* emit( + events, + command, + )({ + type: "thread.metadata-updated", + threadId: command.threadId, + providerInstanceId: thread.providerInstanceId, + occurredAt: now, + payload: { + ...rest, + ...(command.handoff === null ? {} : { handoff: command.handoff }), + updatedAt: now, + }, + }); + }, + ); + /** * Records a task that runs in a linked environment on its parent: the * task, its node and its turn item, as a local delegation does, but no child @@ -10739,6 +10815,15 @@ const makeOrchestrator = Effect.fn("orchestrationV2.Orchestrator.layer")(functio threadId: command.threadId, }); } + // A thread moving away, or gone, takes no new turns here. + if (thread.handoff?.state === "departing" || thread.handoff?.state === "departed") { + return yield* new OrchestratorThreadMovedError({ + commandId: command.commandId, + threadId: command.threadId, + label: thread.handoff.label, + state: thread.handoff.state, + }); + } yield* dispatchMessage(command, events, effects); break; } @@ -10874,6 +10959,9 @@ const makeOrchestrator = Effect.fn("orchestrationV2.Orchestrator.layer")(functio case "delegated_task.remote.request": yield* dispatchRemoteDelegatedTaskRequest(command, events); break; + case "thread.handoff.update": + yield* dispatchThreadHandoffUpdate(command, events); + break; case "delegated_task.remote.complete": yield* dispatchRemoteDelegatedTaskComplete(command, events); break; diff --git a/apps/server/src/orchestration-v2/ProjectionStore.ts b/apps/server/src/orchestration-v2/ProjectionStore.ts index 7995810867df..ceebdb0ba58b 100644 --- a/apps/server/src/orchestration-v2/ProjectionStore.ts +++ b/apps/server/src/orchestration-v2/ProjectionStore.ts @@ -1476,6 +1476,7 @@ export function threadShellFromProjection( ...(projection.thread.delegatedFrom === undefined ? {} : { delegatedFrom: projection.thread.delegatedFrom }), + ...(projection.thread.handoff === undefined ? {} : { handoff: projection.thread.handoff }), latestRunId: latestRun?.id ?? null, latestRunRequestedAt: latestRun?.requestedAt ?? null, latestRunStartedAt: latestRun?.startedAt ?? null, @@ -1750,6 +1751,7 @@ function shellFromState(input: { ...(input.state.thread.delegatedFrom === undefined ? {} : { delegatedFrom: input.state.thread.delegatedFrom }), + ...(input.state.thread.handoff === undefined ? {} : { handoff: input.state.thread.handoff }), latestRunId: input.state.latestRunId, latestRunRequestedAt: input.state.latestRunRequestedAt, latestRunStartedAt: input.state.latestRunStartedAt, diff --git a/apps/server/src/orchestration-v2/ThreadImportService.ts b/apps/server/src/orchestration-v2/ThreadImportService.ts index fdb048edbc83..a07814df73a1 100644 --- a/apps/server/src/orchestration-v2/ThreadImportService.ts +++ b/apps/server/src/orchestration-v2/ThreadImportService.ts @@ -68,6 +68,13 @@ export type ThreadImportError = | ThreadImportProjectMismatchError | ThreadImportContinuationError; +/** A checkout an import prepared for its thread. */ +export interface ImportedWorkspace { + readonly worktreePath: string; + readonly branch: string | null; + readonly undo: Effect.Effect; +} + /** * Creates a thread here seeded with a conversation from another environment, * for a thread moving here. The history is runless and marked as imported, @@ -78,13 +85,18 @@ export type ThreadImportError = export class ThreadImportService extends Context.Service< ThreadImportService, { - readonly importThread: ( + readonly importThread: ( input: OrchestratorMcpThreadImportInput & { readonly runtimeMode: RuntimeMode; readonly interactionMode: ProviderInteractionMode; readonly linkOrigin: OrchestrationV2LinkOrigin | undefined; + /** + * Prepares the checkout the thread works in, run only when this import + * creates the thread; `undo` runs if the thread then cannot be written. + */ + readonly workspace?: Effect.Effect; }, - ) => Effect.Effect; + ) => Effect.Effect; } >()("t3/orchestration-v2/ThreadImportService") {} @@ -190,6 +202,7 @@ const make = Effect.gen(function* () { .pipe(Effect.orElseSucceed(() => null)); let created = false; if (existing === null) { + const workspace = input.workspace === undefined ? undefined : yield* input.workspace; const now = yield* DateTime.now; const thread: OrchestrationV2AppThread = { createdBy: "agent", @@ -201,8 +214,9 @@ const make = Effect.gen(function* () { modelSelection: input.modelSelection, runtimeMode: input.runtimeMode, interactionMode: input.interactionMode, - branch: input.branch ?? null, - worktreePath: input.worktreePath ?? null, + branch: workspace === undefined ? (input.branch ?? null) : workspace.branch, + worktreePath: + workspace === undefined ? (input.worktreePath ?? null) : workspace.worktreePath, activeProviderThreadId: null, historyOrigin: "v1_import", ...(input.linkOrigin === undefined ? {} : { linkOrigin: input.linkOrigin }), @@ -240,6 +254,7 @@ const make = Effect.gen(function* () { Effect.flatMap((raced) => (raced === null ? Effect.failCause(cause) : Effect.void)), ), ), + Effect.tapError(() => workspace?.undo ?? Effect.void), Effect.mapError((cause) => new ThreadImportWriteError({ threadId, cause })), ); created = true; diff --git a/apps/server/src/orchestration-v2/ThreadMessageIntake.ts b/apps/server/src/orchestration-v2/ThreadMessageIntake.ts index c8f65fd0340a..215a812a456c 100644 --- a/apps/server/src/orchestration-v2/ThreadMessageIntake.ts +++ b/apps/server/src/orchestration-v2/ThreadMessageIntake.ts @@ -23,6 +23,7 @@ function dispatchWasNotAccepted( case "OrchestratorCommandPreviouslyRejectedError": case "OrchestratorCommandIdConflictError": case "OrchestratorSubagentThreadReadOnlyError": + case "OrchestratorThreadMovedError": case "OrchestratorThreadAboveModeLimitError": return true; default: diff --git a/apps/server/src/orchestration-v2/testkit/CapturingCodexAdapter.ts b/apps/server/src/orchestration-v2/testkit/CapturingCodexAdapter.ts index 645c2e7c3f9c..8c4b1ee6ac40 100644 --- a/apps/server/src/orchestration-v2/testkit/CapturingCodexAdapter.ts +++ b/apps/server/src/orchestration-v2/testkit/CapturingCodexAdapter.ts @@ -38,7 +38,12 @@ const unimplemented = (detail: string) => */ export const makeCapturingCodexAdapter = ( capturedTurns: Ref.Ref>, - options: { readonly response: string; readonly modelSelection: ModelSelection }, + options: { + readonly response: string; + readonly modelSelection: ModelSelection; + /** Holds each turn open until this completes, so a test can act mid-turn. */ + readonly holdTurn?: Effect.Effect; + }, ) => ({ instanceId, @@ -61,6 +66,78 @@ export const makeCapturingCodexAdapter = ( updatedAt: now, lastError: null, }; + const providerTurnUpdate = ( + turnInput: ProviderAdapter.ProviderAdapterV2TurnInput, + status: "running" | "completed", + startedAt: DateTime.Utc, + completedAt: DateTime.Utc | null, + ): ProviderAdapter.ProviderAdapterV2Event => ({ + type: "provider_turn.updated", + driver, + providerTurn: { + id: ProviderTurnId.make(`provider-turn:${turnInput.threadId}:${turnInput.runOrdinal}`), + providerThreadId: turnInput.providerThread.id, + nodeId: turnInput.rootNodeId, + runAttemptId: turnInput.attemptId, + nativeTurnRef: { + driver, + nativeId: `native-turn:${turnInput.threadId}:${turnInput.runOrdinal}`, + strength: "strong", + }, + ordinal: turnInput.runOrdinal, + status, + startedAt, + completedAt, + }, + }); + const finishTurn = (turnInput: ProviderAdapter.ProviderAdapterV2TurnInput) => + Effect.gen(function* () { + const eventTime = yield* DateTime.now; + const providerTurnId = ProviderTurnId.make( + `provider-turn:${turnInput.threadId}:${turnInput.runOrdinal}`, + ); + yield* PubSub.publishAll(events, [ + providerTurnUpdate(turnInput, "completed", eventTime, eventTime), + { + type: "turn_item.updated", + driver, + turnItem: { + id: TurnItemId.make( + `turn-item:${turnInput.threadId}:${turnInput.runOrdinal}:assistant`, + ), + threadId: turnInput.threadId, + runId: turnInput.runId, + nodeId: turnInput.rootNodeId, + providerThreadId: turnInput.providerThread.id, + providerTurnId, + nativeItemRef: null, + parentItemId: null, + ordinal: turnInput.runOrdinal * 100 + 1, + status: "completed", + title: null, + startedAt: eventTime, + completedAt: eventTime, + updatedAt: eventTime, + type: "assistant_message", + messageId: MessageId.make( + `message:${turnInput.threadId}:${turnInput.runOrdinal}:assistant`, + ), + text: options.response, + streaming: false, + }, + }, + { + type: "turn.terminal", + driver, + providerThreadId: turnInput.providerThread.id, + providerTurnId, + runOrdinal: turnInput.runOrdinal, + status: "completed", + failure: null, + threadDisposition: "reusable", + }, + ] satisfies ReadonlyArray); + }); return { instanceId, driver, @@ -104,69 +181,19 @@ export const makeCapturingCodexAdapter = ( text: turnInput.message.text, }, ]); - const eventTime = yield* DateTime.now; - const providerTurnId = ProviderTurnId.make( - `provider-turn:${turnInput.threadId}:${turnInput.runOrdinal}`, - ); - yield* PubSub.publishAll(events, [ - { - type: "provider_turn.updated", - driver, - providerTurn: { - id: providerTurnId, - providerThreadId: turnInput.providerThread.id, - nodeId: turnInput.rootNodeId, - runAttemptId: turnInput.attemptId, - nativeTurnRef: { - driver, - nativeId: `native-turn:${turnInput.threadId}:${turnInput.runOrdinal}`, - strength: "strong", - }, - ordinal: turnInput.runOrdinal, - status: "completed", - startedAt: eventTime, - completedAt: eventTime, - }, - }, - { - type: "turn_item.updated", - driver, - turnItem: { - id: TurnItemId.make( - `turn-item:${turnInput.threadId}:${turnInput.runOrdinal}:assistant`, - ), - threadId: turnInput.threadId, - runId: turnInput.runId, - nodeId: turnInput.rootNodeId, - providerThreadId: turnInput.providerThread.id, - providerTurnId, - nativeItemRef: null, - parentItemId: null, - ordinal: turnInput.runOrdinal * 100 + 1, - status: "completed", - title: null, - startedAt: eventTime, - completedAt: eventTime, - updatedAt: eventTime, - type: "assistant_message", - messageId: MessageId.make( - `message:${turnInput.threadId}:${turnInput.runOrdinal}:assistant`, - ), - text: options.response, - streaming: false, - }, - }, - { - type: "turn.terminal", - driver, - providerThreadId: turnInput.providerThread.id, - providerTurnId, - runOrdinal: turnInput.runOrdinal, - status: "completed", - failure: null, - threadDisposition: "reusable", - }, - ] satisfies ReadonlyArray); + if (options.holdTurn !== undefined) { + // Ends the turn later, as a live provider does; startTurn returns now. + // Until then the turn is running, so Stop and steering can target it. + yield* PubSub.publish( + events, + providerTurnUpdate(turnInput, "running", yield* DateTime.now, null), + ); + yield* Effect.forkDetach( + options.holdTurn.pipe(Effect.andThen(finishTurn(turnInput))), + ); + return; + } + yield* finishTurn(turnInput); }), steerTurn: () => Effect.void, interruptTurn: () => Effect.void, diff --git a/apps/server/src/peer/PeerForwarding.test.ts b/apps/server/src/peer/PeerForwarding.test.ts index 77b983606a08..0bc6d0760b6a 100644 --- a/apps/server/src/peer/PeerForwarding.test.ts +++ b/apps/server/src/peer/PeerForwarding.test.ts @@ -35,6 +35,8 @@ import * as PeerForwarding from "./PeerForwarding.ts"; import * as PeerLinks from "./PeerLinks.ts"; import * as PeerLinkRequests from "./PeerLinkRequests.ts"; import * as ThreadImportService from "../orchestration-v2/ThreadImportService.ts"; +import * as ThreadHandoff from "./handoff/ThreadHandoff.ts"; +import * as HandoffImport from "./handoff/HandoffImport.ts"; import { descriptorOf, layerLinkingEnvironment, linkTo, servePeer } from "./PeerLinks.testkit.ts"; const laptop = descriptorOf("environment-laptop", "Laptop"); @@ -177,6 +179,8 @@ const serveBox = (seen: Ref.Ref, linkSession: Ref.Ref) => Layer.provide(Layer.mock(PeerLinkRequests.PeerLinkRequests)({})), Layer.provide(Layer.mock(PeerLinks.PeerLinks)({})), Layer.provide(Layer.mock(ThreadImportService.ThreadImportService)({})), + Layer.provide(Layer.mock(ThreadHandoff.ThreadHandoff)({})), + Layer.provide(Layer.mock(HandoffImport.HandoffImport)({})), ), ); @@ -238,14 +242,21 @@ const makeLaptop = Effect.gen(function* () { ), }), ), - Layer.provide(Layer.mock(ProviderRegistry.ProviderRegistry)({})), - Layer.provide(Layer.mock(ProviderAdapterRegistry.ProviderAdapterRegistryV2)({})), - Layer.provide(Layer.mock(ScheduledTaskService.ScheduledTaskService)({})), - Layer.provide(Layer.mock(ProjectService.ProjectService)({})), - Layer.provide(Layer.mock(SecretRequests.SecretRequests)({})), - Layer.provide(Layer.mock(RemoteDelegation.RemoteDelegation)({})), - Layer.provide(Layer.mock(PeerLinkRequests.PeerLinkRequests)({})), - Layer.provide(Layer.mock(ThreadImportService.ThreadImportService)({})), + // Services the toolkits declare that these calls never reach. + Layer.provide( + Layer.mergeAll( + Layer.mock(ProviderRegistry.ProviderRegistry)({}), + Layer.mock(ProviderAdapterRegistry.ProviderAdapterRegistryV2)({}), + Layer.mock(ScheduledTaskService.ScheduledTaskService)({}), + Layer.mock(ProjectService.ProjectService)({}), + Layer.mock(SecretRequests.SecretRequests)({}), + Layer.mock(RemoteDelegation.RemoteDelegation)({}), + Layer.mock(PeerLinkRequests.PeerLinkRequests)({}), + Layer.mock(ThreadImportService.ThreadImportService)({}), + Layer.mock(ThreadHandoff.ThreadHandoff)({}), + Layer.mock(HandoffImport.HandoffImport)({}), + ), + ), Layer.fresh, ); const here = yield* Layer.build(layerHere); diff --git a/apps/server/src/peer/PeerLinkRequests.test.ts b/apps/server/src/peer/PeerLinkRequests.test.ts index 8d021263f8c6..dd60a726c1ba 100644 --- a/apps/server/src/peer/PeerLinkRequests.test.ts +++ b/apps/server/src/peer/PeerLinkRequests.test.ts @@ -53,6 +53,7 @@ import { type ServedPeer, } from "./PeerLinks.testkit.ts"; import * as RemoteDelegation from "./RemoteDelegation.ts"; +import * as ThreadHandoff from "./handoff/ThreadHandoff.ts"; // The laptop's agent asks the user to link the box. The laptop is its real // orchestrator and toolkit; the box is its real descriptor and MCP OAuth on a @@ -141,6 +142,7 @@ const makeLaptop = Effect.gen(function* () { Layer.provideMerge(PeerLinkRequests.layer), Layer.provide(Layer.mock(PeerForwarding.PeerForwarding)({})), Layer.provide(Layer.mock(RemoteDelegation.RemoteDelegation)({})), + Layer.provide(Layer.mock(ThreadHandoff.ThreadHandoff)({})), Layer.provide(Layer.succeedContext(linking)), Layer.provide(NodeCrypto.layer), Layer.provideMerge(layerThreads), diff --git a/apps/server/src/peer/PeerLinks.testkit.ts b/apps/server/src/peer/PeerLinks.testkit.ts index d4a35c251677..942a4e8a5bc9 100644 --- a/apps/server/src/peer/PeerLinks.testkit.ts +++ b/apps/server/src/peer/PeerLinks.testkit.ts @@ -57,9 +57,16 @@ export const descriptorOf = ( * socket. The descriptor can be swapped, as when its address comes to answer * as another environment, and every bearer it receives is recorded. */ -export const servePeer = ( +export const servePeer = ( initial: ExecutionEnvironmentDescriptor, tools: Layer.Layer, + /** More routes the peer serves, such as attachment uploads. */ + routes: Layer.Layer = Layer.empty as Layer.Layer, + /** + * The peer's own config and secret store, when its tools need the same ones + * its auth uses, as signed uploads do. A fresh pair otherwise. + */ + base?: Context.Context, ) => Effect.gen(function* () { const descriptor = yield* Ref.make(initial); @@ -69,11 +76,18 @@ export const servePeer = ( getEnvironmentId: Ref.get(descriptor).pipe(Effect.map((current) => current.environmentId)), getDescriptor: Ref.get(descriptor), }); + const layerBase = + base === undefined + ? ServerSecretStore.layer.pipe( + Layer.provideMerge( + ServerConfig.layerTest(process.cwd(), { prefix: "t3-peer-link-peer-" }), + ), + ) + : Layer.succeedContext(base); const authContext = yield* EnvironmentAuth.layer.pipe( Layer.provide(Sqlite.layerMemory), - Layer.provideMerge(ServerSecretStore.layer), Layer.provideMerge(ServerEnvironment.layerIdentity), - Layer.provideMerge(ServerConfig.layerTest(process.cwd(), { prefix: "t3-peer-link-peer-" })), + Layer.provideMerge(layerBase), Layer.provideMerge(NodeServices.layer), Layer.fresh, Layer.build, @@ -108,6 +122,7 @@ export const servePeer = ( ), Layer.provide(McpOAuth.layerMcpClientAuthenticator), ), + routes, ).pipe( Layer.provide(layerRecordBearers), Layer.provide(layerEnvironment), diff --git a/apps/server/src/peer/RemoteDelegation.test.ts b/apps/server/src/peer/RemoteDelegation.test.ts index 584e39ef0349..accbe2468fad 100644 --- a/apps/server/src/peer/RemoteDelegation.test.ts +++ b/apps/server/src/peer/RemoteDelegation.test.ts @@ -61,6 +61,8 @@ import { descriptorOf, layerLinkingEnvironment, linkTo, servePeer } from "./Peer import * as RemoteDelegation from "./RemoteDelegation.ts"; import * as PeerLinkRequests from "./PeerLinkRequests.ts"; import * as ThreadImportService from "../orchestration-v2/ThreadImportService.ts"; +import * as ThreadHandoff from "./handoff/ThreadHandoff.ts"; +import * as HandoffImport from "./handoff/HandoffImport.ts"; // The laptop's agent delegates a task to the box through a link. The box is // its real /mcp behind real OAuth, with one thread the test finishes; the @@ -370,6 +372,8 @@ const boxToolkitLayer = (thread: BoxThread) => { Layer.provide(Layer.mock(PeerLinkRequests.PeerLinkRequests)({})), Layer.provide(Layer.mock(PeerLinks.PeerLinks)({})), Layer.provide(Layer.mock(ThreadImportService.ThreadImportService)({})), + Layer.provide(Layer.mock(ThreadHandoff.ThreadHandoff)({})), + Layer.provide(Layer.mock(HandoffImport.HandoffImport)({})), Layer.provide( Layer.mock(ManagedProjectFolders.ManagedProjectFolders)({ namedProjectsRoot: "/p" }), ), @@ -480,6 +484,7 @@ const makeLaptop = ( ), ), Layer.provide(Layer.succeedContext(linking)), + Layer.provide(Layer.mock(ThreadHandoff.ThreadHandoff)({})), Layer.provide(Layer.mock(PeerLinkRequests.PeerLinkRequests)({})), Layer.provide(NodeCrypto.layer), Layer.provideMerge(layerThreads), diff --git a/apps/server/src/peer/handoff/HandoffImport.ts b/apps/server/src/peer/handoff/HandoffImport.ts new file mode 100644 index 000000000000..a0f0093074de --- /dev/null +++ b/apps/server/src/peer/handoff/HandoffImport.ts @@ -0,0 +1,131 @@ +import { OrchestratorMcpFailure } from "@t3tools/contracts"; +import * as Context from "effect/Context"; +import * as Effect from "effect/Effect"; +import * as FileSystem from "effect/FileSystem"; +import * as Layer from "effect/Layer"; +import * as Path from "effect/Path"; +import * as HostProcess from "@t3tools/shared/HostProcess"; + +import { deletePendingAttachment } from "../../assets/AttachmentUpload.ts"; +import { + parseThreadSegmentFromAttachmentId, + resolveAttachmentPathById, +} from "../../attachmentStore.ts"; +import * as ServerConfig from "../../config.ts"; +import type * as ThreadImportService from "../../orchestration-v2/ThreadImportService.ts"; +import * as ServerSettings from "../../serverSettings.ts"; +import * as VcsProcess from "../../vcs/VcsProcess.ts"; +import { resolveWorktreesDirectory } from "../../worktreesDirectory.ts"; +import * as HandoffGit from "./HandoffGit.ts"; + +/** + * The checkout a handed-off thread works in here: its bundle, uploaded as a + * pending attachment, applied in a new worktree of the project. The upload is + * spent either way; the worktree and branch are undone if the import fails. + */ +export class HandoffImport extends Context.Service< + HandoffImport, + { + readonly applyBundle: (input: { + readonly repoRoot: string; + readonly handoffId: string; + readonly bundle: { + readonly attachmentId: string; + readonly branch: string | null; + readonly tip: string; + readonly snapshot: string; + }; + }) => Effect.Effect; + } +>()("t3/peer/handoff/HandoffImport") {} + +const make = Effect.gen(function* () { + const config = yield* ServerConfig.ServerConfig; + const settings = yield* ServerSettings.ServerSettingsService; + const path = yield* Path.Path; + const fileSystem = yield* FileSystem.FileSystem; + const processes = yield* VcsProcess.VcsProcess; + const git = yield* HandoffGit.HandoffGit; + + const applyBundle: HandoffImport["Service"]["applyBundle"] = (input) => + Effect.gen(function* () { + if (parseThreadSegmentFromAttachmentId(input.bundle.attachmentId) !== "pending") { + return yield* new OrchestratorMcpFailure({ + code: "invalid_request", + message: "The bundle must be a pending upload from t3_attachment_prepare_upload.", + }); + } + const bundlePath = resolveAttachmentPathById({ + attachmentsDir: config.attachmentsDir, + attachmentId: input.bundle.attachmentId, + }); + if (bundlePath === null) { + return yield* new OrchestratorMcpFailure({ + code: "invalid_request", + message: "The bundle upload was not found. Upload it again.", + }); + } + const worktreesSetting = yield* settings.getSettings.pipe( + Effect.map((current) => current.worktreesDirectory), + Effect.orElseSucceed(() => ""), + ); + const parent = resolveWorktreesDirectory( + worktreesSetting, + config.worktreesDir, + path, + yield* HostProcess.HomeDirectory, + ); + if (parent === null) { + return yield* new OrchestratorMcpFailure({ + code: "invalid_request", + message: + "This environment's worktree location is not usable. Fix it in Settings → Storage.", + }); + } + const name = (input.bundle.branch ?? `handoff-${input.handoffId}`).replace(/\//g, "-"); + const worktreePath = path.join(parent, path.basename(input.repoRoot), name); + const applied = yield* git + .apply({ + repoRoot: input.repoRoot, + worktreePath, + handoffId: input.handoffId, + bundlePath, + branch: input.bundle.branch, + tip: input.bundle.tip, + snapshot: input.bundle.snapshot, + }) + .pipe( + Effect.mapError( + (error) => + new OrchestratorMcpFailure({ + code: + error._tag === "HandoffGitStepError" ? "orchestration_error" : "invalid_request", + message: error.message, + }), + ), + Effect.ensuring( + deletePendingAttachment(input.bundle.attachmentId).pipe( + Effect.provideService(ServerConfig.ServerConfig, config), + Effect.provideService(FileSystem.FileSystem, fileSystem), + Effect.ignore, + ), + ), + ); + return { + worktreePath, + branch: applied.branch, + undo: processes + .run({ + operation: "HandoffImport.undo", + command: "git", + cwd: input.repoRoot, + args: ["worktree", "remove", "--force", worktreePath], + }) + .pipe(Effect.ignore), + } satisfies ThreadImportService.ImportedWorkspace; + }); + + return HandoffImport.of({ applyBundle }); +}); + +export const layer = Layer.effect(HandoffImport, make); diff --git a/apps/server/src/peer/handoff/ThreadHandoff.test.ts b/apps/server/src/peer/handoff/ThreadHandoff.test.ts new file mode 100644 index 000000000000..80ff221ce4a9 --- /dev/null +++ b/apps/server/src/peer/handoff/ThreadHandoff.test.ts @@ -0,0 +1,619 @@ +import { NodeHttpServer } from "@effect/platform-node"; +import * as NodeCrypto from "@effect/platform-node/NodeCrypto"; +import * as NodeServices from "@effect/platform-node/NodeServices"; +import { + CommandId, + MessageId, + type OrchestrationV2ThreadProjection, + type Project, + ProjectId, + ProviderInstanceId, + ThreadId, +} from "@t3tools/contracts"; +import { expect, it } from "@effect/vitest"; +import * as Context from "effect/Context"; +import * as Queue from "effect/Queue"; +import * as Effect from "effect/Effect"; +import * as FileSystem from "effect/FileSystem"; +import * as Layer from "effect/Layer"; +import * as Option from "effect/Option"; +import * as Path from "effect/Path"; +import * as Ref from "effect/Ref"; +import { FetchHttpClient } from "effect/http"; + +import * as ServerSecretStore from "../../auth/ServerSecretStore.ts"; +import * as ServerConfig from "../../config.ts"; +import * as ServerEnvironment from "../../environment/ServerEnvironment.ts"; +import * as ServerHttp from "../../http.ts"; +import * as McpHttpServer from "../../mcp/McpHttpServer.ts"; +import * as EventSink from "../../orchestration-v2/EventSink.ts"; +import * as EventStore from "../../orchestration-v2/EventStore.ts"; +import * as Orchestrator from "../../orchestration-v2/Orchestrator.ts"; +import * as ProjectionStore from "../../orchestration-v2/ProjectionStore.ts"; +import * as ProviderAdapterRegistry from "../../orchestration-v2/ProviderAdapterRegistry.ts"; +import * as ThreadImportService from "../../orchestration-v2/ThreadImportService.ts"; +import * as ThreadLaunch from "../../orchestration-v2/ThreadLaunchService.ts"; +import * as ThreadManagement from "../../orchestration-v2/ThreadManagementService.ts"; +import { + type CapturedTurn, + makeCapturingCodexAdapter, +} from "../../orchestration-v2/testkit/CapturingCodexAdapter.ts"; +import * as ProviderReplayHarness from "../../orchestration-v2/testkit/ProviderReplayHarness.ts"; +import * as SqlitePersistence from "../../persistence/Sqlite.ts"; +import * as ManagedProjectFolders from "../../project/ManagedProjectFolders.ts"; +import * as ProjectService from "../../project/ProjectService.ts"; +import * as RepositoryIdentityResolver from "../../project/RepositoryIdentityResolver.ts"; +import * as ProviderRegistry from "../../provider/ProviderRegistry.ts"; +import * as ScheduledTaskService from "../../scheduledTasks/ScheduledTaskService.ts"; +import * as SecretRequests from "../../secrets/SecretRequests.ts"; +import * as ServerSettings from "../../serverSettings.ts"; +import * as SourceControlRepositoryService from "../../sourceControl/SourceControlRepositoryService.ts"; +import * as GitVcsDriver from "../../vcs/GitVcsDriver.ts"; +import * as VcsProcess from "../../vcs/VcsProcess.ts"; +import * as PeerForwarding from "../PeerForwarding.ts"; +import * as PeerLinks from "../PeerLinks.ts"; +import { descriptorOf, layerLinkingEnvironment, linkTo, servePeer } from "../PeerLinks.testkit.ts"; +import * as RemoteDelegation from "../RemoteDelegation.ts"; +import * as HandoffGit from "./HandoffGit.ts"; +import * as HandoffImport from "./HandoffImport.ts"; +import * as ThreadHandoff from "./ThreadHandoff.ts"; +import * as PeerLinkRequests from "../PeerLinkRequests.ts"; + +// A thread on the laptop moves to the box: both are real orchestrators, the +// box behind its real /mcp, OAuth and upload route, the laptop linked to it, +// and both with a real clone of the same repository. The thread arrives with +// its conversation and its git work, and only one copy stays live. + +const laptop = descriptorOf("environment-laptop", "Laptop"); +const box = descriptorOf("environment-box", "Box"); +const modelSelection = { instanceId: ProviderInstanceId.make("codex"), model: "gpt-5.4" }; + +const git = (cwd: string, args: ReadonlyArray) => + VcsProcess.VcsProcess.pipe( + Effect.flatMap((processes) => + processes.run({ + operation: "ThreadHandoff.test", + command: "git", + cwd, + args, + env: { + ...process.env, + GIT_AUTHOR_NAME: "Test", + GIT_AUTHOR_EMAIL: "test@example.com", + GIT_COMMITTER_NAME: "Test", + GIT_COMMITTER_EMAIL: "test@example.com", + }, + }), + ), + Effect.map((result) => result.stdout), + Effect.orDie, + ); + +/** An origin and the laptop's and the box's clones of it. */ +const makeRepos = Effect.gen(function* () { + const fileSystem = yield* FileSystem.FileSystem; + const path = yield* Path.Path; + const root = yield* fileSystem.makeTempDirectoryScoped({ prefix: "t3-thread-handoff-" }); + const origin = path.join(root, "origin.git"); + const seed = path.join(root, "seed"); + yield* fileSystem.makeDirectory(seed); + yield* git(root, ["init", "--quiet", "--bare", "--initial-branch=main", origin]); + yield* git(seed, ["init", "--quiet", "--initial-branch=main"]); + yield* fileSystem.writeFileString(path.join(seed, "app.ts"), "export const version = 1;\n"); + yield* git(seed, ["add", "."]); + yield* git(seed, ["commit", "--quiet", "-m", "start"]); + yield* git(seed, ["remote", "add", "origin", origin]); + yield* git(seed, ["push", "--quiet", "origin", "main"]); + const laptopRepo = path.join(root, "laptop"); + const boxRepo = path.join(root, "box"); + yield* git(root, ["clone", "--quiet", origin, laptopRepo]); + yield* git(root, ["clone", "--quiet", origin, boxRepo]); + return { root, laptopRepo, boxRepo }; +}); + +const projectAt = (id: string, workspaceRoot: string): Project => ({ + id: ProjectId.make(id), + title: "app", + workspaceRoot, + // Blank, as a cold enrichment cache leaves it: both sides resolve the + // clones' shared origin themselves. + repositoryIdentity: null, + defaultModelSelection: null, + scripts: [], + createdAt: "2026-10-01T00:00:00.000Z", + updatedAt: "2026-10-01T00:00:00.000Z", + deletedAt: null, +}); + +/** One project, whose repository identity is read from its real checkout, as the service does. */ +const layerProjects = (project: Project) => + Layer.effect( + ProjectService.ProjectService, + Effect.gen(function* () { + const identities = yield* RepositoryIdentityResolver.RepositoryIdentityResolver; + return ProjectService.ProjectService.of({ + getById: (id: ProjectId) => + Effect.succeed(id === project.id ? Option.some(project) : Option.none()), + snapshot: Effect.succeed({ projects: [project], updatedAt: project.updatedAt }), + listPage: () => + identities.resolve(project.workspaceRoot).pipe( + Effect.map((repositoryIdentity) => ({ + projects: [{ ...project, repositoryIdentity }], + nextCursor: null, + })), + ), + } as unknown as ProjectService.ProjectService["Service"]); + }), + ); + +/** What each test runs on: a test HTTP server, git, and Node. */ +const layerTestHost = Layer.mergeAll(NodeHttpServer.layerTest, VcsProcess.layer).pipe( + Layer.provideMerge(NodeServices.layer), +); + +/** A real orchestrator on an in-memory database, with a capturing provider. */ +const layerOrchestration = ( + name: string, + workspace: string, + captured: Ref.Ref>, + holdTurn?: Effect.Effect, +) => { + const layerDatabase = SqlitePersistence.layerMemory; + const layerStores = Layer.mergeAll( + layerDatabase, + EventStore.layer.pipe(Layer.provideMerge(layerDatabase)), + ProjectionStore.layer.pipe(Layer.provideMerge(layerDatabase)), + ); + const layerOrchestrator = ProviderReplayHarness.layerWithRegistry( + { + name, + runtimePolicyOverride: { + cwd: workspace, + approvalPolicy: "never", + sandboxPolicy: { type: "readOnly", access: { type: "fullAccess" }, networkAccess: false }, + }, + }, + ProviderAdapterRegistry.layerSingle( + makeCapturingCodexAdapter(captured, { + response: `${name} replied.`, + modelSelection, + ...(holdTurn === undefined ? {} : { holdTurn }), + }), + ), + { databaseLayer: layerDatabase }, + ); + return ThreadManagement.layer.pipe( + Layer.provideMerge(layerOrchestrator), + Layer.provideMerge(EventSink.layer.pipe(Layer.provide(layerStores))), + Layer.provideMerge(layerStores), + ); +}; + +const unusedServices = Layer.mergeAll( + Layer.mock(ProviderRegistry.ProviderRegistry)({ getProviders: Effect.succeed([]) }), + Layer.mock(ProviderAdapterRegistry.ProviderAdapterRegistryV2)({ list: () => Effect.succeed([]) }), + Layer.mock(ScheduledTaskService.ScheduledTaskService)({}), + Layer.mock(SecretRequests.SecretRequests)({}), + Layer.mock(ThreadLaunch.ThreadLaunchService)({}), + Layer.mock(ManagedProjectFolders.ManagedProjectFolders)({ namedProjectsRoot: "/p" }), + Layer.mock(SourceControlRepositoryService.SourceControlRepositoryService)({}), + Layer.mock(RemoteDelegation.RemoteDelegation)({}), + Layer.mock(PeerForwarding.PeerForwarding)({}), + Layer.mock(PeerLinkRequests.PeerLinkRequests)({}), + Layer.mock(PeerLinks.PeerLinks)({}), + // The box takes moves; it never starts one. + Layer.mock(ThreadHandoff.ThreadHandoff)({}), +); + +/** The box: its real tools on a real /mcp, with uploads, on its own repo. */ +const serveBox = ( + boxRepo: string, + stateDir: string, + captured: Ref.Ref>, +) => + Effect.gen(function* () { + const core = yield* Layer.mergeAll( + ThreadImportService.layer.pipe( + Layer.provideMerge(layerOrchestration("box", boxRepo, captured)), + ), + VcsProcess.layer, + ServerSettings.layerTest(), + ).pipe( + Layer.provide(NodeCrypto.layer), + Layer.provideMerge( + ServerSecretStore.layer.pipe( + Layer.provideMerge(ServerConfig.layerTest(process.cwd(), stateDir)), + ), + ), + Layer.provideMerge(NodeServices.layer), + Layer.build, + ); + const layerCore = Layer.succeedContext(core); + const tools = Layer.mergeAll( + McpHttpServer.layerOrchestratorToolkit, + McpHttpServer.layerProjectRegistration, + McpHttpServer.layerAttachmentToolkit, + ).pipe( + Layer.provide(HandoffImport.layer.pipe(Layer.provide(HandoffGit.layer))), + Layer.provide( + layerProjects(projectAt("project:box", boxRepo)).pipe( + Layer.provideMerge(RepositoryIdentityResolver.layer), + ), + ), + Layer.provide( + Layer.mock(GitVcsDriver.GitVcsDriver)({ + listWorktreePaths: () => Effect.succeed([]), + }), + ), + Layer.provide(unusedServices), + Layer.provide(layerCore), + Layer.provide(NodeCrypto.layer), + ); + // Upload URLs are signed and checked with the box's own secret store. + const served = yield* servePeer(box, tools, ServerHttp.layerAttachmentUploadRoute, core).pipe( + Effect.provide(layerCore), + ); + return { + ...served, + orchestrator: Context.get(core, Orchestrator.OrchestratorV2), + }; + }); + +/** The laptop: a real orchestrator with a thread mid-conversation, linked to the box. */ +const makeLaptop = ( + laptopRepo: string, + captured: Ref.Ref>, + holdTurn?: Effect.Effect, +) => + Effect.gen(function* () { + const linking = yield* layerLinkingEnvironment(laptop).pipe(Layer.build); + const core = yield* layerOrchestration("laptop", laptopRepo, captured, holdTurn).pipe( + Layer.provideMerge( + ServerSecretStore.layer.pipe( + Layer.provideMerge( + ServerConfig.layerTest(process.cwd(), { prefix: "t3-handoff-laptop-" }), + ), + ), + ), + Layer.provideMerge(NodeServices.layer), + Layer.build, + ); + const layerHere = ThreadHandoff.layer.pipe( + Layer.provideMerge(PeerForwarding.layer), + Layer.provide(HandoffGit.layer), + Layer.provide( + layerProjects(projectAt("project:laptop", laptopRepo)).pipe( + Layer.provideMerge(RepositoryIdentityResolver.layer), + ), + ), + Layer.provide( + Layer.succeed(ServerEnvironment.ServerEnvironment, { + getEnvironmentId: Effect.succeed(laptop.environmentId), + getDescriptor: Effect.succeed(laptop), + }), + ), + Layer.provide(Layer.succeedContext(linking)), + Layer.provide(Layer.succeedContext(core)), + Layer.provide(VcsProcess.layer), + Layer.provide(NodeCrypto.layer), + // Uploads go to the box's own address, not the test server's. + Layer.provide(FetchHttpClient.layer), + Layer.provideMerge(NodeServices.layer), + Layer.fresh, + ); + const here = yield* Layer.build(layerHere); + return { + links: Context.get(linking, PeerLinks.PeerLinks), + orchestrator: Context.get(core, Orchestrator.OrchestratorV2), + threads: Context.get(core, ThreadManagement.ThreadManagementService), + handoff: Context.get(here, ThreadHandoff.ThreadHandoff), + }; + }); + +const threadId = ThreadId.make("thread:laptop-login"); + +/** The laptop's thread: one finished turn about the bug, on the laptop's checkout. */ +const createThread = ( + laptopEnv: Effect.Success>, + laptopRepo: string, +) => + laptopEnv.orchestrator.dispatch({ + type: "thread.create", + createdBy: "user", + creationSource: "web", + commandId: CommandId.make("command:create:laptop-login"), + threadId, + projectId: ProjectId.make("project:laptop"), + title: "Flaky login test", + modelSelection, + runtimeMode: "full-access", + interactionMode: "default", + branch: null, + worktreePath: laptopRepo, + }); + +const seedThread = (laptopEnv: Effect.Success>, laptopRepo: string) => + Effect.gen(function* () { + yield* laptopEnv.orchestrator.dispatch({ + type: "thread.create", + createdBy: "user", + creationSource: "web", + commandId: CommandId.make("command:create:laptop-login"), + threadId, + projectId: ProjectId.make("project:laptop"), + title: "Flaky login test", + modelSelection, + runtimeMode: "full-access", + interactionMode: "default", + branch: null, + worktreePath: laptopRepo, + }); + yield* laptopEnv.orchestrator.dispatch({ + type: "message.dispatch", + createdBy: "user", + creationSource: "web", + commandId: CommandId.make("command:send:laptop-login"), + threadId, + messageId: MessageId.make("message:laptop-login:1"), + text: "LAPTOP_QUESTION: why does the login test flake?", + attachments: [], + dispatchMode: { type: "start_immediately" }, + }); + yield* idle(laptopEnv.orchestrator, threadId); + }); + +/** Polls `check` until it holds; for states only another fiber reaches. */ +const waitFor = (check: Effect.Effect) => + Effect.gen(function* () { + for (let attempt = 0; attempt < 2_000; attempt += 1) { + if (yield* check.pipe(Effect.orElseSucceed(() => false))) return; + yield* Effect.sleep("5 millis"); + } + return yield* Effect.die(new Error("The awaited state never came.")); + }); + +const idle = (orchestrator: Orchestrator.OrchestratorV2["Service"], id: ThreadId) => + Effect.gen(function* () { + for (let attempt = 0; attempt < 2_000; attempt += 1) { + const projection: OrchestrationV2ThreadProjection = + yield* orchestrator.getThreadProjection(id); + if ( + projection.runs.length > 0 && + projection.runs.every( + (run) => !["queued", "starting", "running", "waiting", "preparing"].includes(run.status), + ) + ) { + return projection; + } + yield* Effect.sleep("5 millis"); + } + return yield* Effect.die(new Error(`${id} never went idle.`)); + }); + +it.live("a thread moves to the box with its conversation and its git work", () => + Effect.scoped( + Effect.gen(function* () { + const fileSystem = yield* FileSystem.FileSystem; + const path = yield* Path.Path; + const { root, laptopRepo, boxRepo } = yield* makeRepos; + const boxTurns = yield* Ref.make>([]); + const laptopTurns = yield* Ref.make>([]); + const b = yield* serveBox(boxRepo, path.join(root, "box-state"), boxTurns); + const a = yield* makeLaptop(laptopRepo, laptopTurns); + yield* linkTo(a.links, b, "full-access"); + yield* seedThread(a, laptopRepo); + + // The thread's work on the laptop: an unpushed commit, then uncommitted edits. + yield* git(laptopRepo, ["checkout", "--quiet", "-b", "fix/login"]); + yield* fileSystem.writeFileString( + path.join(laptopRepo, "app.ts"), + "export const version = 2;\n", + ); + yield* git(laptopRepo, ["commit", "--quiet", "-am", "unpushed fix"]); + yield* fileSystem.writeFileString( + path.join(laptopRepo, "app.ts"), + "export const version = 3;\n", + ); + yield* fileSystem.writeFileString(path.join(laptopRepo, "notes.md"), "# findings\n"); + + expect(yield* a.handoff.options(threadId)).toEqual([ + { environmentId: box.environmentId, label: "Box", projectId: "project:box", reason: null }, + ]); + const moved = yield* a.handoff.start({ threadId, environmentId: box.environmentId }); + expect(moved.state).toBe("departed"); + if (moved.state !== "departed") return; + + // On the box: the thread, its history, its first run with that context. + const there = yield* idle(b.orchestrator, moved.threadId); + expect(there.thread.historyOrigin).toBe("v1_import"); + expect(there.messages.some((message) => message.text.includes("LAPTOP_QUESTION"))).toBe(true); + const [firstTurn] = yield* Ref.get(boxTurns); + expect(firstTurn?.text).toContain("LAPTOP_QUESTION"); + expect(firstTurn?.text).toContain("Continue where it left off."); + // ...in a new worktree on the branch, with the commit and the edits. + const worktree = there.thread.worktreePath!; + expect(there.thread.branch).toBe("fix/login"); + expect((yield* git(worktree, ["log", "-1", "--format=%s"])).trim()).toBe("unpushed fix"); + expect(yield* fileSystem.readFileString(path.join(worktree, "app.ts"))).toBe( + "export const version = 3;\n", + ); + expect(yield* fileSystem.readFileString(path.join(worktree, "notes.md"))).toBe( + "# findings\n", + ); + + // On the laptop: read-only, pointing at the box. + const here = yield* a.orchestrator.getThreadProjection(threadId); + expect(here.thread.handoff).toMatchObject({ state: "departed", threadId: moved.threadId }); + const refused = yield* a.orchestrator + .dispatch({ + type: "message.dispatch", + createdBy: "user", + creationSource: "web", + commandId: CommandId.make("command:send:after-move"), + threadId, + messageId: MessageId.make("message:after-move"), + text: "Still here?", + attachments: [], + dispatchMode: { type: "start_immediately" }, + }) + .pipe(Effect.flip); + expect(refused.message).toBe("This thread moved to Box. Continue it there."); + }), + ).pipe(Effect.provide(layerTestHost)), +); + +it.live("a move the box refuses releases the thread here, untouched", () => + Effect.scoped( + Effect.gen(function* () { + const fileSystem = yield* FileSystem.FileSystem; + const path = yield* Path.Path; + const { root, laptopRepo, boxRepo } = yield* makeRepos; + const turns = yield* Ref.make>([]); + const b = yield* serveBox(boxRepo, path.join(root, "box-state"), turns); + const a = yield* makeLaptop(laptopRepo, turns); + yield* linkTo(a.links, b, "full-access"); + yield* seedThread(a, laptopRepo); + // The box has its own commit on the branch the laptop's thread is on. + yield* git(laptopRepo, ["checkout", "--quiet", "-b", "fix/login"]); + yield* fileSystem.writeFileString( + path.join(laptopRepo, "app.ts"), + "export const version = 2;\n", + ); + yield* git(laptopRepo, ["commit", "--quiet", "-am", "laptop fix"]); + yield* git(boxRepo, ["checkout", "--quiet", "-b", "fix/login"]); + yield* fileSystem.writeFileString(path.join(boxRepo, "box.ts"), "export const b = 1;\n"); + yield* git(boxRepo, ["add", "."]); + yield* git(boxRepo, ["commit", "--quiet", "-m", "box fix"]); + yield* git(boxRepo, ["checkout", "--quiet", "main"]); + + const failed = yield* a.handoff + .start({ threadId, environmentId: box.environmentId }) + .pipe(Effect.flip); + expect(failed.message).toContain("commits the handed-off work lacks"); + const here = yield* a.orchestrator.getThreadProjection(threadId); + expect(here.thread.handoff).toMatchObject({ state: "failed" }); + // Released: the thread takes turns here again. + yield* a.orchestrator.dispatch({ + type: "message.dispatch", + createdBy: "user", + creationSource: "web", + commandId: CommandId.make("command:send:after-failure"), + threadId, + messageId: MessageId.make("message:after-failure"), + text: "Keep going here then.", + attachments: [], + dispatchMode: { type: "start_immediately" }, + }); + // Nothing was created on the box. + expect( + (yield* b.orchestrator.getShellSnapshot()).threads.filter((shell) => + shell.title.includes("Flaky"), + ), + ).toEqual([]); + }), + ).pipe(Effect.provide(layerTestHost)), +); + +it.live("an agent's own move waits for its turn to end, and a message before then cancels it", () => + Effect.scoped( + Effect.gen(function* () { + const path = yield* Path.Path; + const { root, laptopRepo, boxRepo } = yield* makeRepos; + const turns = yield* Ref.make>([]); + const b = yield* serveBox(boxRepo, path.join(root, "box-state"), turns); + // The laptop's turns stay open until the test lets one finish; each + // release lets exactly one waiting turn end. + const releases = yield* Queue.unbounded(); + const hold = Queue.take(releases); + const releaseTurn = Queue.offer(releases, undefined); + const a = yield* makeLaptop(laptopRepo, turns, hold); + yield* linkTo(a.links, b, "full-access"); + yield* a.handoff.start_(); + yield* createThread(a, laptopRepo); + const send = (id: string, text: string) => + a.orchestrator.dispatch({ + type: "message.dispatch", + createdBy: "user", + creationSource: "web", + commandId: CommandId.make(`command:send:${id}`), + threadId, + messageId: MessageId.make(`message:${id}`), + text, + attachments: [], + dispatchMode: { type: "start_immediately" }, + }); + const handoffState = a.orchestrator + .getThreadProjection(threadId) + .pipe(Effect.map((projection) => projection.thread.handoff?.state ?? null)); + // The agent asks from inside the thread, so the move waits for its turn. + const askToMove = a.handoff.start({ + callerThreadId: threadId, + environmentId: box.environmentId, + continuationPrompt: "SELF_HANDOFF_CONTINUE: run the e2e suite there.", + }); + + // Mid-turn, the agent asks to move: nothing moves yet. + yield* send("turn-1", "Wrap up and move this to the box."); + yield* waitFor( + a.orchestrator + .getThreadProjection(threadId) + .pipe( + Effect.map((projection) => projection.runs.some((run) => run.status === "running")), + ), + ); + expect((yield* askToMove).state).toBe("pending"); + // An agent the user works through steers the same turn before it ends: + // new instructions win. + yield* waitFor( + a.orchestrator + .getThreadProjection(threadId) + .pipe( + Effect.map((projection) => + projection.providerTurns.some((turn) => turn.status === "running"), + ), + ), + ); + const running = (yield* a.orchestrator.getThreadProjection(threadId)).runs.find( + (run) => run.status === "running", + ); + if (running === undefined) return yield* Effect.die("no running turn"); + yield* Effect.sleep("2 millis"); + yield* a.orchestrator.dispatch({ + type: "message.dispatch", + createdBy: "agent", + creationSource: "web", + commandId: CommandId.make("command:send:steer"), + threadId, + messageId: MessageId.make("message:steer"), + text: "Actually, stay here.", + attachments: [], + dispatchMode: { type: "steer_active", targetRunId: running.id }, + }); + yield* releaseTurn; + yield* waitFor(handoffState.pipe(Effect.map((state) => state === null))); + yield* idle(a.orchestrator, threadId); + expect(yield* handoffState).toBe(null); + + // Asked again, with nothing after it: it moves once the turn ends. + yield* send("turn-3", "Now move it."); + yield* waitFor( + a.orchestrator + .getThreadProjection(threadId) + .pipe( + Effect.map((projection) => projection.runs.some((run) => run.status === "running")), + ), + ); + expect((yield* askToMove).state).toBe("pending"); + // Settling it while the turn still runs, as a restart's sweep does, leaves it waiting. + yield* a.handoff.settle(threadId); + expect(yield* handoffState).toBe("pending"); + yield* releaseTurn; + yield* waitFor(handoffState.pipe(Effect.map((state) => state === "departed"))); + const moved = (yield* a.orchestrator.getThreadProjection(threadId)).thread.handoff; + if (moved?.state !== "departed") return yield* Effect.die("not departed"); + yield* idle(b.orchestrator, moved.threadId); + // The box's first turn is the agent's own continuation. + expect((yield* Ref.get(turns)).at(-1)?.text).toContain("SELF_HANDOFF_CONTINUE"); + }), + ).pipe(Effect.provide(layerTestHost)), +); diff --git a/apps/server/src/peer/handoff/ThreadHandoff.ts b/apps/server/src/peer/handoff/ThreadHandoff.ts new file mode 100644 index 000000000000..c26e54e151d0 --- /dev/null +++ b/apps/server/src/peer/handoff/ThreadHandoff.ts @@ -0,0 +1,577 @@ +import { + CommandId, + type EnvironmentId, + type OrchestratorMcpFailure, + ThreadHandoffError, + type OrchestrationV2ThreadHandoff, + type OrchestratorMcpImportedMessage, + type ProjectId, + type ThreadId, +} from "@t3tools/contracts"; +import * as Context from "effect/Context"; +import * as Crypto from "effect/Crypto"; +import * as DateTime from "effect/DateTime"; +import * as Effect from "effect/Effect"; +import * as FileSystem from "effect/FileSystem"; +import * as Layer from "effect/Layer"; +import * as Option from "effect/Option"; +import type * as Scope from "effect/Scope"; +import * as Stream from "effect/Stream"; +import * as HttpBody from "effect/http/HttpBody"; +import * as HttpClient from "effect/http/HttpClient"; +import * as HttpClientRequest from "effect/http/HttpClientRequest"; + +import type * as McpInvocationContext from "../../mcp/McpInvocationContext.ts"; +import { AttachmentToolkit } from "../../mcp/toolkits/attachment/tools.ts"; +import { ProjectToolkit } from "../../mcp/toolkits/project/tools.ts"; +import * as ProjectionStore from "../../orchestration-v2/ProjectionStore.ts"; +import * as ThreadManagement from "../../orchestration-v2/ThreadManagementService.ts"; +import * as ServerEnvironment from "../../environment/ServerEnvironment.ts"; +import * as ProjectService from "../../project/ProjectService.ts"; +import * as RepositoryIdentityResolver from "../../project/RepositoryIdentityResolver.ts"; +import { forkParked } from "../../serverActivation.ts"; +import * as PeerForwarding from "../PeerForwarding.ts"; +import * as PeerLinks from "../PeerLinks.ts"; +import * as HandoffGit from "./HandoffGit.ts"; + +/** What the thread there starts with, unless the move names its own. */ +const DEFAULT_CONTINUATION = + "This thread was just moved here from another machine, with its conversation and its uncommitted work. Continue where it left off."; + +type HandoffState = OrchestrationV2ThreadHandoff["state"]; + +/** + * Moves a thread to a linked environment: its conversation, and its git work + * (unpushed commits, uncommitted edits, new files). Exactly one copy stays + * live: the thread here takes no new turns once it starts moving, and reads + * only once it has moved. A failed move releases it with the reason. + * + * An agent can ask to move its own thread mid-turn; the move then waits for + * that turn to end (`pending`), and a message to the thread before then cancels it. + */ +export class ThreadHandoff extends Context.Service< + ThreadHandoff, + { + /** Where the thread can go: linked environments with a project of its repository. */ + readonly options: (threadId: ThreadId) => Effect.Effect< + ReadonlyArray<{ + readonly environmentId: EnvironmentId; + readonly label: string; + readonly projectId: ProjectId | null; + readonly reason: string | null; + }>, + ThreadHandoffError + >; + /** + * Moves `threadId` to `environmentId` now. An agent moving the thread it + * runs in (`callerThreadId`, the default target) is mid-turn, so its + * thread moves once that turn ends. Returns the state the move is in. + */ + readonly start: (input: { + readonly threadId?: ThreadId | undefined; + readonly callerThreadId?: ThreadId | undefined; + readonly environmentId: EnvironmentId; + readonly projectId?: ProjectId | undefined; + readonly continuationPrompt?: string | undefined; + }) => Effect.Effect; + /** Cancels a move that is still waiting for the turn to end. */ + readonly cancel: (threadId: ThreadId) => Effect.Effect; + /** Runs moves whose turn ended, and settles any a restart cut short. */ + readonly start_: () => Effect.Effect; + /** Settles one thread's move now: what the follower does when its turn ends. */ + readonly settle: (threadId: ThreadId) => Effect.Effect; + } +>()("t3/peer/handoff/ThreadHandoff") {} + +/** A failure the other environment reported, carried as the move's own. */ +const fromRemote = (cause: OrchestratorMcpFailure) => + new ThreadHandoffError({ reason: "failed", message: cause.message, cause }); + +const make = Effect.gen(function* () { + const threads = yield* ThreadManagement.ThreadManagementService; + const projections = yield* ProjectionStore.ProjectionStoreV2; + const projects = yield* ProjectService.ProjectService; + const repositoryIdentities = yield* RepositoryIdentityResolver.RepositoryIdentityResolver; + const forwarding = yield* PeerForwarding.PeerForwarding; + const links = yield* PeerLinks.PeerLinks; + const git = yield* HandoffGit.HandoffGit; + const httpClient = yield* HttpClient.HttpClient; + const fileSystem = yield* FileSystem.FileSystem; + const crypto = yield* Crypto.Crypto; + const hereId = yield* ServerEnvironment.ServerEnvironment.pipe( + Effect.flatMap((environment) => environment.getEnvironmentId), + ); + + /** The thread's agent, acting on its own behalf there through the link. */ + const scopeFor = (thread: { + readonly id: ThreadId; + readonly providerInstanceId: string; + }): McpInvocationContext.McpThreadInvocationScope => ({ + environmentId: hereId, + requestNamespace: `handoff:${thread.id}`, + thread: { + threadId: thread.id, + providerSessionId: `handoff:${thread.id}`, + providerInstanceId: thread.providerInstanceId as never, + }, + client: undefined, + capabilities: new Set(["orchestration"]), + issuedAt: 0, + }); + + const setHandoff = ( + threadId: ThreadId, + expected: HandoffState | null, + handoff: OrchestrationV2ThreadHandoff | null, + step: string, + ) => + threads + .dispatch({ + type: "thread.handoff.update", + commandId: CommandId.make( + `command:handoff:${threadId}:${handoff?.handoffId ?? "clear"}:${step}`, + ), + threadId, + expected, + handoff, + }) + .pipe( + Effect.asVoid, + Effect.mapError( + (error) => new ThreadHandoffError({ reason: "refused", message: error.message }), + ), + ); + + /** The project there with this thread's repository, or the one named. */ + const remoteProject = ( + scope: McpInvocationContext.McpThreadInvocationScope, + environmentId: EnvironmentId, + projectId: ProjectId, + requested: ProjectId | undefined, + ) => + Effect.gen(function* () { + if (requested !== undefined) return requested; + const here = yield* projects.getById(projectId).pipe( + Effect.map(Option.getOrUndefined), + Effect.orElseSucceed(() => undefined), + ); + // The project's own identity reads blank while its cache is cold. + const key = + here === undefined + ? undefined + : (yield* repositoryIdentities.resolve(here.workspaceRoot))?.canonicalKey; + if (key === undefined) { + return yield* new ThreadHandoffError({ + reason: "refused", + message: + "This thread's project has no repository to match there. Pick the project there.", + }); + } + const listed = yield* forwarding + .call(scope, ProjectToolkit.tools.t3_project_list, environmentId, { limit: 100 }) + .pipe(Effect.mapError(fromRemote)); + const matches = listed.projects.filter( + (project) => project.repositoryIdentity?.canonicalKey === key, + ); + if (matches.length !== 1) { + return yield* new ThreadHandoffError({ + reason: "refused", + message: + matches.length === 0 + ? `No project there has this repository (${key}). Add it there first.` + : `Several projects there have this repository; pick one: ${matches.map((p) => p.title).join(", ")}.`, + }); + } + return matches[0]!.id; + }); + + /** The conversation as the thread there will see it. */ + const conversation = (threadId: ThreadId) => + threads.getThreadRecords(threadId, ["messages"], { messageRoles: ["user", "assistant"] }).pipe( + Effect.map((records) => + records.messages + .filter( + (message) => + !message.streaming && + message.notification === undefined && + message.delegatedCompletion === undefined && + message.text.trim() !== "", + ) + .slice(-2_000) + .map((message): OrchestratorMcpImportedMessage => ({ + role: message.role === "user" ? "user" : "assistant", + text: message.text.slice(0, 200_000), + createdAt: DateTime.formatIso(message.createdAt), + })), + ), + Effect.mapError( + () => + new ThreadHandoffError({ + reason: "failed", + message: "The conversation could not be read.", + }), + ), + ); + + /** Uploads the bundle there through its signed upload route. */ + const upload = ( + scope: McpInvocationContext.McpThreadInvocationScope, + environmentId: EnvironmentId, + packed: HandoffGit.HandoffPackage, + ) => + Effect.gen(function* () { + const prepared = yield* forwarding + .call(scope, AttachmentToolkit.tools.t3_attachment_prepare_upload, environmentId, { + upload: { + type: "file", + name: `${packed.snapshot.slice(0, 12)}.bundle`, + mimeType: "application/x-git-bundle", + sizeBytes: packed.bundleBytes, + }, + }) + .pipe(Effect.mapError(fromRemote)); + const resolved = yield* links + .resolve(environmentId) + .pipe( + Effect.mapError( + (error) => new ThreadHandoffError({ reason: "failed", message: error.message }), + ), + ); + const bytes = yield* fileSystem.readFile(packed.bundlePath).pipe( + Effect.mapError( + () => + new ThreadHandoffError({ + reason: "failed", + message: "The bundle could not be read.", + }), + ), + ); + const response = yield* httpClient + .execute( + HttpClientRequest.post(`${resolved.url}${prepared.relativeUrl}`).pipe( + HttpClientRequest.setBody(HttpBody.uint8Array(bytes, "application/x-git-bundle")), + ), + ) + .pipe( + Effect.mapError( + () => + new ThreadHandoffError({ + reason: "failed", + message: "The linked environment stopped answering mid-upload.", + }), + ), + ); + if (response.status !== 204) { + return yield* new ThreadHandoffError({ + reason: "failed", + message: `The linked environment refused the upload (HTTP ${response.status}).`, + }); + } + return prepared.attachmentId; + }); + + /** The move itself, from `departing`: package, send, import, then mark departed. */ + const depart = (threadId: ThreadId, projectHint?: ProjectId | undefined) => + Effect.gen(function* () { + const records = yield* threads + .getThreadRecords(threadId, []) + .pipe( + Effect.mapError( + () => + new ThreadHandoffError({ reason: "not_found", message: "The thread was not found." }), + ), + ); + const thread = records.thread; + const handoff = thread.handoff; + if (handoff?.state !== "departing") return handoff; + const scope = scopeFor(thread); + const result = yield* Effect.gen(function* () { + const projectId = yield* remoteProject( + scope, + handoff.environmentId, + thread.projectId, + projectHint, + ); + const project = yield* projects.getById(thread.projectId).pipe( + Effect.map(Option.getOrUndefined), + Effect.orElseSucceed(() => undefined), + ); + const cwd = thread.worktreePath ?? project?.workspaceRoot; + const tempDir = yield* fileSystem.makeTempDirectoryScoped({ prefix: "t3-handoff-" }).pipe( + Effect.mapError( + () => + new ThreadHandoffError({ + reason: "failed", + message: "No temp space for the move.", + }), + ), + ); + const packed = + cwd === undefined + ? undefined + : yield* git.pack({ cwd, handoffId: handoff.handoffId, outDir: tempDir }).pipe( + Effect.catchTags({ + // A thread outside git moves with its conversation only. + HandoffNotARepositoryError: () => Effect.succeed(undefined), + }), + Effect.mapError( + (error) => new ThreadHandoffError({ reason: "refused", message: error.message }), + ), + ); + const attachmentId = + packed === undefined ? undefined : yield* upload(scope, handoff.environmentId, packed); + return yield* forwarding + .call(scope, ProjectToolkit.tools.t3_thread_import, handoff.environmentId, { + source: { + environmentId: hereId, + threadId, + handoffId: handoff.handoffId, + }, + projectId, + title: thread.title, + modelSelection: thread.modelSelection, + runtimeMode: thread.runtimeMode, + interactionMode: thread.interactionMode, + messages: yield* conversation(threadId), + continuationPrompt: handoff.continuationPrompt ?? DEFAULT_CONTINUATION, + ...(packed === undefined || attachmentId === undefined + ? {} + : { + bundle: { + attachmentId, + branch: packed.branch, + tip: packed.tip, + snapshot: packed.snapshot, + }, + }), + }) + .pipe(Effect.mapError(fromRemote)); + }).pipe(Effect.scoped, Effect.result); + if (result._tag === "Failure") { + yield* setHandoff( + threadId, + "departing", + { + state: "failed", + handoffId: handoff.handoffId, + environmentId: handoff.environmentId, + label: handoff.label, + lastError: result.failure.message, + }, + "failed", + ).pipe(Effect.ignore); + return yield* Effect.fail(result.failure); + } + const departed: OrchestrationV2ThreadHandoff = { + state: "departed", + handoffId: handoff.handoffId, + environmentId: handoff.environmentId, + label: handoff.label, + threadId: result.success.threadId, + }; + yield* setHandoff(threadId, "departing", departed, "departed"); + return departed; + }); + + const options: ThreadHandoff["Service"]["options"] = (threadId) => + Effect.gen(function* () { + const records = yield* threads + .getThreadRecords(threadId, []) + .pipe( + Effect.mapError( + () => + new ThreadHandoffError({ reason: "not_found", message: "The thread was not found." }), + ), + ); + const listed = yield* links.list.pipe(Effect.orElseSucceed(() => [])); + return yield* Effect.forEach( + listed, + (link) => + link.status !== "reachable" + ? Effect.succeed({ + environmentId: link.environmentId, + label: link.label, + projectId: null, + reason: + link.status === "expired" ? "The link expired." : "Not answering right now.", + }) + : remoteProject( + scopeFor(records.thread), + link.environmentId, + records.thread.projectId, + undefined, + ).pipe( + Effect.map((projectId) => ({ + environmentId: link.environmentId, + label: link.label, + projectId, + reason: null, + })), + Effect.catch((error) => + Effect.succeed({ + environmentId: link.environmentId, + label: link.label, + projectId: null, + reason: error.message, + }), + ), + ), + { concurrency: 4 }, + ); + }); + + const start: ThreadHandoff["Service"]["start"] = (request) => + Effect.gen(function* () { + const threadId = request.threadId ?? request.callerThreadId; + if (threadId === undefined) { + return yield* new ThreadHandoffError({ + reason: "refused", + message: "Pass threadId: this MCP client is not running inside a T3 thread.", + }); + } + const input = { + ...request, + threadId, + whenTurnEnds: threadId === request.callerThreadId, + }; + const records = yield* threads + .getThreadRecords(input.threadId, ["runs"]) + .pipe( + Effect.mapError( + () => + new ThreadHandoffError({ reason: "not_found", message: "The thread was not found." }), + ), + ); + const thread = records.thread; + if (thread.archivedAt !== null) { + return yield* new ThreadHandoffError({ + reason: "refused", + message: "Unarchive the thread to move it.", + }); + } + const current = thread.handoff?.state ?? null; + if (current === "pending" || current === "departing" || current === "departed") { + return yield* new ThreadHandoffError({ + reason: "refused", + message: `This thread's move is already ${current}.`, + }); + } + const link = (yield* links.list.pipe(Effect.orElseSucceed(() => []))).find( + (candidate) => candidate.environmentId === input.environmentId, + ); + if (link === undefined) { + return yield* new ThreadHandoffError({ + reason: "refused", + message: "That environment is not linked.", + }); + } + // Fails now, while the asking agent can still say so, rather than after its turn. + yield* remoteProject( + scopeFor(thread), + input.environmentId, + thread.projectId, + input.projectId, + ); + const handoffId = (yield* crypto.randomUUIDv4.pipe(Effect.orDie)).replaceAll("-", ""); + const busy = records.runs.some( + (run) => + run.status === "running" || + run.status === "starting" || + run.status === "waiting" || + run.status === "queued" || + run.status === "preparing", + ); + const base = { + handoffId, + environmentId: input.environmentId, + label: link.label, + ...(input.continuationPrompt === undefined + ? {} + : { continuationPrompt: input.continuationPrompt }), + }; + if (busy) { + if (input.whenTurnEnds !== true) { + return yield* new ThreadHandoffError({ + reason: "refused", + message: "The thread is running a turn. Stop it, or let it move when the turn ends.", + }); + } + const pending: OrchestrationV2ThreadHandoff = { state: "pending", ...base }; + yield* setHandoff(input.threadId, current, pending, "pending"); + return pending; + } + yield* setHandoff(input.threadId, current, { state: "departing", ...base }, "departing"); + return (yield* depart(input.threadId, input.projectId)) ?? { state: "departing", ...base }; + }); + + const cancel: ThreadHandoff["Service"]["cancel"] = (threadId) => + setHandoff(threadId, "pending", null, "cancel"); + + /** A pending move whose turn ended departs; one a restart cut short is settled. */ + const settle = (threadId: ThreadId) => + Effect.gen(function* () { + const records = yield* threads.getThreadRecords(threadId, ["runs", "messages"], { + messageRoles: ["user"], + }); + const handoff = records.thread.handoff; + if (handoff === undefined) return; + if (handoff.state === "pending") { + const live = records.runs.some( + (run) => + run.status === "running" || + run.status === "starting" || + run.status === "waiting" || + run.status === "preparing", + ); + if (live) return; + // Someone wrote during the turn the agent asked in, steered into it or + // queued after it: new instructions win. That is the user, or an agent + // the user works through, such as an orchestrator. + const asked = records.runs.at(-1)?.requestedAt; + const writtenSince = records.messages.some( + (message) => asked !== undefined && DateTime.isGreaterThan(message.createdAt, asked), + ); + if (writtenSince || records.runs.some((run) => run.status === "queued")) { + yield* setHandoff(threadId, "pending", null, "cancel-by-user"); + return; + } + yield* setHandoff(threadId, "pending", { ...handoff, state: "departing" }, "departing"); + } + const now = yield* threads.getThreadRecords(threadId, []); + if (now.thread.handoff?.state === "departing") { + yield* depart(threadId).pipe(Effect.ignoreCause({ log: true })); + } + }).pipe(Effect.ignoreCause({ log: true })); + + const start_: ThreadHandoff["Service"]["start_"] = () => + forkParked( + Effect.gen(function* () { + // Moves a restart cut short, or that were waiting for a turn that ended meanwhile. + const shells = yield* projections.getShellSnapshot().pipe( + Effect.map((snapshot) => snapshot.threads), + Effect.orElseSucceed(() => []), + ); + yield* Effect.forEach( + shells.filter( + (shell) => shell.handoff?.state === "pending" || shell.handoff?.state === "departing", + ), + (shell) => settle(shell.id), + { discard: true }, + ); + yield* threads.streamStoredEvents.pipe( + Stream.filter( + (stored) => + stored.event.type === "run.updated" && + ["completed", "failed", "cancelled", "interrupted", "rolled_back"].includes( + stored.event.payload.status, + ), + ), + Stream.runForEach((stored) => settle(stored.event.threadId)), + ); + }).pipe(Effect.ignoreCause({ log: true })), + ); + + return ThreadHandoff.of({ options, start, cancel, start_, settle }); +}); + +export const layer = Layer.effect(ThreadHandoff, make); diff --git a/apps/server/src/server.ts b/apps/server/src/server.ts index 98f930a001ec..824272c1fe3b 100644 --- a/apps/server/src/server.ts +++ b/apps/server/src/server.ts @@ -60,6 +60,9 @@ import * as McpHttpServer from "./mcp/McpHttpServer.ts"; import * as PeerForwarding from "./peer/PeerForwarding.ts"; import * as RemoteDelegation from "./peer/RemoteDelegation.ts"; import * as PeerLinkRequests from "./peer/PeerLinkRequests.ts"; +import * as HandoffGit from "./peer/handoff/HandoffGit.ts"; +import * as HandoffImport from "./peer/handoff/HandoffImport.ts"; +import * as ThreadHandoff from "./peer/handoff/ThreadHandoff.ts"; import * as PeerLinks from "./peer/PeerLinks.ts"; import * as PeerMcpClient from "./peer/PeerMcpClient.ts"; import * as McpSessionRegistry from "./mcp/McpSessionRegistry.ts"; @@ -725,11 +728,23 @@ const layerMakeRoutes = Layer.mergeAll( Layer.effectDiscard( RemoteDelegation.RemoteDelegation.pipe(Effect.flatMap((remote) => remote.start())), ), + // Runs moves that waited for a turn to end, and settles any a restart cut short. + Layer.effectDiscard( + ThreadHandoff.ThreadHandoff.pipe(Effect.flatMap((handoff) => handoff.start_())), + ), ).pipe( // delegate_task to a linked environment, from /mcp and the follower above. Layer.provide(RemoteDelegation.layer.pipe(Layer.provide(ProjectionStoreV2.layer))), // Links agents ask for: shown in a thread from /mcp, answered over WebSocket. Layer.provide(PeerLinkRequests.layer), + // Moving a thread to a linked environment, from Settings, menus and agents, + // and applying one moved here (t3_thread_import). + Layer.provide( + Layer.merge( + ThreadHandoff.layer.pipe(Layer.provide(ProjectionStoreV2.layer)), + HandoffImport.layer, + ).pipe(Layer.provideMerge(HandoffGit.layer)), + ), // Both transports consume the same service instance, so caches single-flight across clients // and mutations observed on WebSocket invalidate patches subsequently read over HTTP. Layer.provide(layerPullRequestService), diff --git a/apps/server/src/ws.ts b/apps/server/src/ws.ts index 2754327fee13..4d7c98b72352 100644 --- a/apps/server/src/ws.ts +++ b/apps/server/src/ws.ts @@ -3,6 +3,7 @@ import * as Crypto from "effect/Crypto"; import * as Orchestrator from "./orchestration-v2/Orchestrator.ts"; import * as PeerLinkRequests from "./peer/PeerLinkRequests.ts"; import * as PeerLinks from "./peer/PeerLinks.ts"; +import * as ThreadHandoff from "./peer/handoff/ThreadHandoff.ts"; import * as DateTime from "effect/DateTime"; import * as Duration from "effect/Duration"; @@ -1283,6 +1284,7 @@ const layerWsRpc = ( const serverSettings = yield* ServerSettings.ServerSettingsService; const peerLinks = yield* PeerLinks.PeerLinks; const peerLinkRequests = yield* PeerLinkRequests.PeerLinkRequests; + const threadHandoff = yield* ThreadHandoff.ThreadHandoff; const startup = yield* ServerRuntimeStartup.ServerRuntimeStartup; const workspaceEntries = yield* WorkspaceEntries.WorkspaceEntries; const workspaceFileSystem = yield* WorkspaceFileSystem.WorkspaceFileSystem; @@ -2395,6 +2397,20 @@ const layerWsRpc = ( Effect.annotateCurrentSpan({ "orchestration_v2.thread_id": input.threadId }).pipe( Effect.andThen(peerLinkRequests.answer(input)), ), + [WS_METHODS.threadHandoffOptions]: ({ threadId }) => + Effect.annotateCurrentSpan({ "orchestration_v2.thread_id": threadId }).pipe( + Effect.andThen(threadHandoff.options(threadId)), + Effect.map((options) => ({ options })), + ), + [WS_METHODS.threadHandoffStart]: (input) => + Effect.annotateCurrentSpan({ "orchestration_v2.thread_id": input.threadId }).pipe( + Effect.andThen(threadHandoff.start(input)), + ), + [WS_METHODS.threadHandoffCancel]: ({ threadId }) => + Effect.annotateCurrentSpan({ "orchestration_v2.thread_id": threadId }).pipe( + Effect.andThen(threadHandoff.cancel(threadId)), + Effect.as({}), + ), [WS_METHODS.serverGetSettings]: (_input) => serverSettings.getSettings.pipe(Effect.map(ServerSettings.redactServerSettingsForClient)), [WS_METHODS.serverUpdateSettings]: ({ patch, providerInstanceMutation }) => @@ -3143,6 +3159,7 @@ export const layer = Layer.unwrap( const pullRequests = yield* PullRequestService.PullRequestService; const peerLinks = yield* PeerLinks.PeerLinks; const peerLinkRequests = yield* PeerLinkRequests.PeerLinkRequests; + const threadHandoff = yield* ThreadHandoff.ThreadHandoff; const sql = yield* SqlClient.SqlClient; return HttpRouter.add( "GET", @@ -3211,6 +3228,7 @@ export const layer = Layer.unwrap( Layer.provide(Layer.succeed(PullRequestService.PullRequestService, pullRequests)), Layer.provide(Layer.succeed(PeerLinks.PeerLinks, peerLinks)), Layer.provide(Layer.succeed(PeerLinkRequests.PeerLinkRequests, peerLinkRequests)), + Layer.provide(Layer.succeed(ThreadHandoff.ThreadHandoff, threadHandoff)), Layer.provide( SourceControlDiscovery.layer.pipe( Layer.provide( diff --git a/docs/internals/remote.md b/docs/internals/remote.md index 1fbe9dcafba5..29eeeabf11a6 100644 --- a/docs/internals/remote.md +++ b/docs/internals/remote.md @@ -108,6 +108,16 @@ internal command reuses the parent half of a local finalize, so the parent wakes the same way. Open remote tasks are followed again at startup, which is also what keeps restart recovery from treating them as abandoned provider work. +**A thread moves with exactly one live copy.** [`ThreadHandoff`](../../apps/server/src/peer/handoff/ThreadHandoff.ts) +marks the thread `departing`, which the orchestrator refuses new turns for, +packs its git work with [`HandoffGit`](../../apps/server/src/peer/handoff/HandoffGit.ts), +uploads the bundle through the other side's signed attachment route, and calls +its `t3_thread_import`. Only then is the thread `departed` here. Any failure +marks it `failed`, which takes turns again. Import ids derive from the handoff, +so a move that a restart cut short is retried by the startup sweep without +creating a second thread there. An agent's move of its own thread waits as +`pending` until its run ends, and any message to the thread in between cancels it. + The link is routing, not isolation: an agent the link starts runs as the receiving environment's user, inside the limits above. diff --git a/docs/user/remote-access.md b/docs/user/remote-access.md index 96eb83da0680..186ac067169b 100644 --- a/docs/user/remote-access.md +++ b/docs/user/remote-access.md @@ -262,6 +262,21 @@ that session; otherwise paste a pairing code from that machine's removes it here. - Links last 30 days. Link again with a new pairing code to renew one. +### Continue a thread on another machine + +A thread can move to a linked machine with its conversation and its code: +commits you have not pushed, uncommitted edits and new files. It continues +there in a new worktree of the project with the same repository, and the copy +here becomes read-only with a link to it. Ignored files such as `.env` stay +behind, and a move larger than 50 MB asks you to push the branch first. + +You can also tell the agent to move it, for example "I need to wrap up, move +this to my VPS". It finishes its reply first, then the thread moves and picks +up where it said it would. Sending a message before then cancels the move. + +If the other machine's branch has commits this one lacks, the move is refused +and nothing changes on either side. + An agent can also hand a task to a linked machine with `delegate_task`: it runs there as an ordinary thread, in the project with the same repository (or the one `target.projectId` names), and this thread wakes with its result when it diff --git a/packages/contracts/src/orchestrationV2.ts b/packages/contracts/src/orchestrationV2.ts index 41ca5007e21c..f835f32ba70e 100644 --- a/packages/contracts/src/orchestrationV2.ts +++ b/packages/contracts/src/orchestrationV2.ts @@ -118,6 +118,44 @@ export const OrchestrationV2DelegatedFrom = Schema.Struct({ }); export type OrchestrationV2DelegatedFrom = typeof OrchestrationV2DelegatedFrom.Type; +/** + * A thread moving to a linked environment. `pending` waits for the thread's + * own turn to end (the agent asked to move); `departing` is moving now and + * takes no new runs; `departed` lives on there as `threadId`, and this copy + * only reads. Cleared when a move fails, with the reason in `lastError`. + */ +export const OrchestrationV2ThreadHandoff = Schema.Union([ + Schema.Struct({ + state: Schema.Literal("pending"), + handoffId: TrimmedNonEmptyString, + environmentId: EnvironmentId, + label: Schema.String, + continuationPrompt: Schema.optional(Schema.String), + }), + Schema.Struct({ + state: Schema.Literal("departing"), + handoffId: TrimmedNonEmptyString, + environmentId: EnvironmentId, + label: Schema.String, + continuationPrompt: Schema.optional(Schema.String), + }), + Schema.Struct({ + state: Schema.Literal("departed"), + handoffId: TrimmedNonEmptyString, + environmentId: EnvironmentId, + label: Schema.String, + threadId: ThreadId, + }), + Schema.Struct({ + state: Schema.Literal("failed"), + handoffId: TrimmedNonEmptyString, + environmentId: EnvironmentId, + label: Schema.String, + lastError: Schema.String, + }), +]); +export type OrchestrationV2ThreadHandoff = typeof OrchestrationV2ThreadHandoff.Type; + export const OrchestrationV2NativeRefStrength = Schema.Literals(["strong", "weak", "none"]); export type OrchestrationV2NativeRefStrength = typeof OrchestrationV2NativeRefStrength.Type; @@ -407,6 +445,7 @@ export const OrchestrationV2AppThread = Schema.Struct({ historyOrigin: Schema.optional(OrchestrationV2ThreadHistoryOrigin), linkOrigin: Schema.optional(OrchestrationV2LinkOrigin), delegatedFrom: Schema.optional(OrchestrationV2DelegatedFrom), + handoff: Schema.optional(OrchestrationV2ThreadHandoff), lineage: OrchestrationV2AppThreadLineage, forkedFrom: Schema.NullOr( Schema.Union([ @@ -1948,6 +1987,7 @@ export const OrchestrationV2ThreadShell = Schema.Struct({ historyOrigin: Schema.optional(OrchestrationV2ThreadHistoryOrigin), linkOrigin: Schema.optional(OrchestrationV2LinkOrigin), delegatedFrom: Schema.optional(OrchestrationV2DelegatedFrom), + handoff: Schema.optional(OrchestrationV2ThreadHandoff), latestRunId: Schema.NullOr(RunId), latestRunRequestedAt: Schema.optional(Schema.NullOr(Schema.DateTimeUtc)), latestRunStartedAt: Schema.optional(Schema.NullOr(Schema.DateTimeUtc)), @@ -3186,6 +3226,19 @@ const OrchestrationV2InternalCommand = Schema.Union([ }), ), }), + /** + * Sets where a thread is in a move to a linked environment, or clears it + * (`handoff: null`). Moving to `pending` or `departing` requires none in + * progress, and `departing` requires no live or queued run, so a turn and + * a move never overlap. `expected` is the state this transition is from. + */ + Schema.Struct({ + type: Schema.Literal("thread.handoff.update"), + commandId: CommandId, + threadId: ThreadId, + expected: Schema.NullOr(Schema.Literals(["pending", "departing", "departed", "failed"])), + handoff: Schema.NullOr(OrchestrationV2ThreadHandoff), + }), /** * Records a delegated task that runs in a linked environment: the task, its * node and its turn item on the parent, as `delegated_task.request` does, diff --git a/packages/contracts/src/orchestratorMcp.ts b/packages/contracts/src/orchestratorMcp.ts index dbf8942ac0c7..6cf4126ee228 100644 --- a/packages/contracts/src/orchestratorMcp.ts +++ b/packages/contracts/src/orchestratorMcp.ts @@ -557,6 +557,19 @@ export const OrchestratorMcpThreadImportInput = Schema.Struct({ /** The checkout the thread works in here, already prepared. Omit for the project root. */ worktreePath: Schema.optional(TrimmedNonEmptyString), branch: Schema.optional(TrimmedNonEmptyString), + /** + * The thread's git work, uploaded with t3_attachment_prepare_upload: a + * bundle of its branch and working tree. It lands in a new worktree here, + * which the thread then works in; worktreePath must then be omitted. + */ + bundle: Schema.optional( + Schema.Struct({ + attachmentId: TrimmedNonEmptyString.check(Schema.isMaxLength(256)), + branch: Schema.NullOr(TrimmedNonEmptyString), + tip: TrimmedNonEmptyString.check(Schema.isMaxLength(64)), + snapshot: TrimmedNonEmptyString.check(Schema.isMaxLength(64)), + }), + ), messages: Schema.Array(OrchestratorMcpImportedMessage).check(Schema.isMaxLength(2_000)), /** Sent as the thread's first message here, so its agent picks up where it left off. */ continuationPrompt: Schema.optional(OrchestratorMcpPrompt), @@ -572,6 +585,35 @@ export const OrchestratorMcpThreadImportResult = Schema.Struct({ }); export type OrchestratorMcpThreadImportResult = typeof OrchestratorMcpThreadImportResult.Type; +export const OrchestratorMcpThreadHandoffInput = Schema.Struct({ + /** The thread to move; omit for the calling thread. */ + threadId: Schema.optional(ThreadId), + environmentId: EnvironmentId.annotate({ + description: "The linked environment to move to; t3_environment_links lists them.", + }), + projectId: Schema.optional( + ProjectId.annotate({ + description: "That environment's project; omit for the one with this thread's repository.", + }), + ), + continuationPrompt: Schema.optional( + OrchestratorMcpPrompt.annotate({ + description: "The first message there, so the agent picks up where it said it would.", + }), + ), +}); +export type OrchestratorMcpThreadHandoffInput = typeof OrchestratorMcpThreadHandoffInput.Type; + +export const OrchestratorMcpThreadHandoffResult = Schema.Struct({ + state: Schema.Literals(["pending", "departing", "departed", "failed"]), + environmentId: EnvironmentId, + label: Schema.String, + /** The thread there, once it has moved. */ + threadId: Schema.NullOr(ThreadId), + message: Schema.String, +}); +export type OrchestratorMcpThreadHandoffResult = typeof OrchestratorMcpThreadHandoffResult.Type; + export const OrchestratorMcpCapabilitiesInput = Schema.Struct({ environmentId: OrchestratorMcpEnvironmentTarget, }); diff --git a/packages/contracts/src/peerLink.ts b/packages/contracts/src/peerLink.ts index 7882e8d2c74c..3d5718bc685b 100644 --- a/packages/contracts/src/peerLink.ts +++ b/packages/contracts/src/peerLink.ts @@ -128,3 +128,16 @@ export class PeerLinkRequestError extends Schema.TaggedError()( + "ThreadHandoffError", + { + reason: Schema.Literals(["not_found", "refused", "failed"]), + message: Schema.String, + cause: Schema.optionalKey(Schema.Defect()), + }, +) {} diff --git a/packages/contracts/src/rpc.ts b/packages/contracts/src/rpc.ts index 584b4c3c3d6e..293161a9d96c 100644 --- a/packages/contracts/src/rpc.ts +++ b/packages/contracts/src/rpc.ts @@ -22,6 +22,7 @@ import { PeerLinkRemoveResult, PeerLinkRequestAnswerInput, PeerLinkRequestError, + ThreadHandoffError, PeerLinkSummary, } from "./peerLink.ts"; import { @@ -35,7 +36,13 @@ import * as Schema from "effect/Schema"; import * as Rpc from "effect/rpc/Rpc"; import * as RpcGroup from "effect/rpc/RpcGroup"; import * as RpcMiddleware from "effect/rpc/RpcMiddleware"; -import { NonNegativeInt, TrimmedNonEmptyString } from "./baseSchemas.ts"; +import { + EnvironmentId, + NonNegativeInt, + ProjectId, + ThreadId, + TrimmedNonEmptyString, +} from "./baseSchemas.ts"; import { CodexAuthCallbackInput, CodexAuthCallbackState, @@ -219,6 +226,7 @@ import { OrchestrationV2GetShellSnapshotError, OrchestrationV2GetThreadProjectionError, OrchestrationV2RpcSchemas, + OrchestrationV2ThreadHandoff, OrchestrationV2ThreadLaunchError, } from "./orchestrationV2.ts"; import { @@ -479,6 +487,9 @@ export const WS_METHODS = { serverUpsertKeybinding: "server.upsertKeybinding", serverRemoveKeybinding: "server.removeKeybinding", peerLinksList: "peerLinks.list", + threadHandoffOptions: "threadHandoff.options", + threadHandoffStart: "threadHandoff.start", + threadHandoffCancel: "threadHandoff.cancel", peerLinksLink: "peerLinks.link", peerLinksUnlink: "peerLinks.unlink", peerLinksAnswerRequest: "peerLinks.answerRequest", @@ -744,6 +755,46 @@ const WsServerCommitDesktopUpdateRpc = Rpc.make(WS_METHODS.serverCommitDesktopUp error: Schema.Union([ServerSelfUpdateError, EnvironmentAuthorizationError]), }); +const WsThreadHandoffOptionsRpc = Rpc.make(WS_METHODS.threadHandoffOptions, { + payload: Schema.Struct({ threadId: ThreadId }), + success: Schema.Struct({ + options: Schema.Array( + Schema.Struct({ + environmentId: EnvironmentId, + label: Schema.String, + projectId: Schema.NullOr(ProjectId), + reason: Schema.NullOr(Schema.String), + }), + ), + }), + error: Schema.Union([ + ThreadHandoffError, + EnvironmentAuthorizationError, + ]), +}); + +const WsThreadHandoffStartRpc = Rpc.make(WS_METHODS.threadHandoffStart, { + payload: Schema.Struct({ + threadId: ThreadId, + environmentId: EnvironmentId, + projectId: Schema.optional(ProjectId), + }), + success: OrchestrationV2ThreadHandoff, + error: Schema.Union([ + ThreadHandoffError, + EnvironmentAuthorizationError, + ]), +}); + +const WsThreadHandoffCancelRpc = Rpc.make(WS_METHODS.threadHandoffCancel, { + payload: Schema.Struct({ threadId: ThreadId }), + success: Schema.Struct({}), + error: Schema.Union([ + ThreadHandoffError, + EnvironmentAuthorizationError, + ]), +}); + const WsPeerLinksListRpc = Rpc.make(WS_METHODS.peerLinksList, { payload: Schema.Struct({}), success: Schema.Struct({ links: Schema.Array(PeerLinkSummary) }), @@ -1873,6 +1924,9 @@ export const WsRpcGroup = RpcGroup.make( WsServerGetSettingsRpc, WsServerUpdateSettingsRpc, WsPeerLinksListRpc, + WsThreadHandoffOptionsRpc, + WsThreadHandoffStartRpc, + WsThreadHandoffCancelRpc, WsPeerLinksLinkRpc, WsPeerLinksUnlinkRpc, WsPeerLinksAnswerRequestRpc, diff --git a/packages/shared/src/t3McpToolPresentation.ts b/packages/shared/src/t3McpToolPresentation.ts index e6984a6f48b6..1fc1cc41b515 100644 --- a/packages/shared/src/t3McpToolPresentation.ts +++ b/packages/shared/src/t3McpToolPresentation.ts @@ -296,6 +296,10 @@ const T3_MCP_TOOLS: Readonly> = { ["Read", "Reading", "Read", "environment preferences"], "environment-read", ), + t3_thread_handoff: tool( + ["Move", "Moving", "Moved", "a thread to another environment"], + "thread-update", + ), t3_environment_links: tool( ["List", "Listing", "Listed", "linked environments"], "environment-links",