Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
4 changes: 4 additions & 0 deletions apps/server/src/auth/RpcAuthorization.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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";

Expand Down Expand Up @@ -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)({
Expand Down Expand Up @@ -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(
Expand Down
23 changes: 16 additions & 7 deletions apps/server/src/mcp/linkOrigin.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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" }),
),
Expand Down
5 changes: 5 additions & 0 deletions apps/server/src/mcp/toolkits/core.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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";
Expand Down Expand Up @@ -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", () => {
Expand Down Expand Up @@ -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)({})),
Expand Down Expand Up @@ -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)({})),
),
),
Expand Down Expand Up @@ -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)({})),
Expand Down
45 changes: 45 additions & 0 deletions apps/server/src/mcp/toolkits/orchestrator/handlers.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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";
Expand Down Expand Up @@ -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;
Expand Down
17 changes: 17 additions & 0 deletions apps/server/src/mcp/toolkits/orchestrator/tools.ts
Original file line number Diff line number Diff line change
Expand Up @@ -6,6 +6,8 @@ import {
OrchestratorMcpEnvironmentLinksResult,
OrchestratorMcpEnvironmentUnlinkInput,
OrchestratorMcpEnvironmentUnlinkResult,
OrchestratorMcpThreadHandoffInput,
OrchestratorMcpThreadHandoffResult,
OrchestratorMcpCreateThreadsInput,
OrchestratorMcpCreateThreadsResult,
OrchestratorMcpDelegateTaskInput,
Expand Down Expand Up @@ -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";
Expand Down Expand Up @@ -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.",
Expand Down Expand Up @@ -333,4 +349,5 @@ export const OrchestratorToolkit = Toolkit.make(
EnvironmentLinksTool,
EnvironmentLinkTool,
EnvironmentUnlinkTool,
ThreadHandoffTool,
);
43 changes: 33 additions & 10 deletions apps/server/src/mcp/toolkits/project/handlers.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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";
Expand Down Expand Up @@ -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,
}),
),
);
}),
Expand Down
2 changes: 2 additions & 0 deletions apps/server/src/mcp/toolkits/project/tools.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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";
Expand Down Expand Up @@ -172,6 +173,7 @@ const ThreadImportTool = Tool.make("t3_thread_import", {
dependencies: [
...shared.dependencies,
ThreadImportService.ThreadImportService,
HandoffImport.HandoffImport,
GitVcsDriver.GitVcsDriver,
FileSystem.FileSystem,
],
Expand Down
4 changes: 4 additions & 0 deletions apps/server/src/mcp/toolkits/worktree/registration.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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)({}),
Expand All @@ -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)({}),
Expand Down
3 changes: 3 additions & 0 deletions apps/server/src/observability/RpcInstrumentation.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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",
Expand Down
Loading
Loading