Skip to content
Closed
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
18 changes: 9 additions & 9 deletions apps/server/src/mcp/OrchestratorMcpService.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -67,7 +67,7 @@ describe("OrchestratorMcpService", () => {
],
} as unknown as OrchestrationV2ThreadProjection;
const childProjection = {
thread: { id: childThreadId },
thread: { id: childThreadId, title: "Child task" },
runs: [
{
id: childRunId,
Expand Down Expand Up @@ -199,7 +199,7 @@ describe("OrchestratorMcpService", () => {
],
} as unknown as OrchestrationV2ThreadProjection;
const childProjection = {
thread: { id: childThreadId },
thread: { id: childThreadId, title: "Child task" },
runs: [{ id: RunId.make("run:mcp-restart-child"), ordinal: 1, status: "cancelled" }],
contextTransfers: [],
messages: [],
Expand Down Expand Up @@ -283,7 +283,7 @@ describe("OrchestratorMcpService", () => {
],
} as unknown as OrchestrationV2ThreadProjection;
const childProjection = {
thread: { id: childThreadId },
thread: { id: childThreadId, title: "Child task" },
runs: [],
contextTransfers: [],
messages: [],
Expand Down Expand Up @@ -359,7 +359,7 @@ describe("OrchestratorMcpService", () => {
],
} as unknown as OrchestrationV2ThreadProjection;
const childProjection = {
thread: { id: childThreadId },
thread: { id: childThreadId, title: "Child task" },
runs: [{ id: childRunId, status: "running" }],
contextTransfers: [],
messages: [],
Expand Down Expand Up @@ -438,7 +438,7 @@ describe("OrchestratorMcpService", () => {
],
} as unknown as OrchestrationV2ThreadProjection;
const childProjection = {
thread: { id: childThreadId },
thread: { id: childThreadId, title: "Child task" },
runs: [{ id: childRunId, status: "running" }],
contextTransfers: [],
messages: [],
Expand Down Expand Up @@ -524,7 +524,7 @@ describe("OrchestratorMcpService", () => {
],
} as unknown as OrchestrationV2ThreadProjection;
const childProjection = {
thread: { id: childThreadId },
thread: { id: childThreadId, title: "Child task" },
runs: [{ id: childRunId, status: "running" }],
contextTransfers: [],
messages: [],
Expand Down Expand Up @@ -624,7 +624,7 @@ describe("OrchestratorMcpService", () => {
[
childThreadId,
{
thread: { id: childThreadId },
thread: { id: childThreadId, title: "Child task" },
runs: [{ id: RunId.make("run:mcp-cancel-grandchild-child"), status: "running" }],
contextTransfers: [],
messages: [],
Expand All @@ -637,7 +637,7 @@ describe("OrchestratorMcpService", () => {
[
grandchildThreadId,
{
thread: { id: grandchildThreadId },
thread: { id: grandchildThreadId, title: "Grandchild task" },
runs: [],
contextTransfers: [],
messages: [],
Expand Down Expand Up @@ -851,7 +851,7 @@ describe("OrchestratorMcpService provider resolution", () => {
}) as unknown as OrchestrationV2ThreadProjection;

const childProjection = {
thread: { id: childThreadId },
thread: { id: childThreadId, title: "Child task" },
runs: [],
contextTransfers: [],
messages: [],
Expand Down
10 changes: 10 additions & 0 deletions apps/server/src/mcp/OrchestratorMcpService.ts
Original file line number Diff line number Diff line change
Expand Up @@ -1325,6 +1325,11 @@ const make = Effect.gen(function* () {
const response = {
taskId: task.id,
childThreadId: task.childThreadId,
childThreadLink: formatThreadLink({
environmentId: scope.environmentId,
threadId: task.childThreadId,
title: childControls.thread.title,
}),
childRunId: childRun?.id ?? null,
childNodeId: task.id,
status,
Expand Down Expand Up @@ -2239,6 +2244,11 @@ const make = Effect.gen(function* () {
);
return {
threadId,
link: formatThreadLink({
environmentId: scope.environmentId,
threadId,
title: projection.thread.title,
}),
runId: run?.id ?? null,
status: run?.status ?? "idle",
title: projection.thread.title,
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -1695,6 +1695,9 @@ describe("orchestrator MCP toolkit", () => {
expect(delegatedSource.messages[0]).toMatchObject({
senderThreadId: parentThreadId,
});
expect(delegated.childThreadLink).toBe(
`[${delegatedSource.thread.title}](t3-thread://v1/environment%3Amcp-orchestrator/${encodeURIComponent(delegated.childThreadId)})`,
);
expect(
delegatedSource.turnItems.find((item) => item.type === "user_message"),
).toMatchObject({
Expand Down Expand Up @@ -2089,6 +2092,9 @@ describe("orchestrator MCP toolkit", () => {
expect(createdSource.messages[0]).toMatchObject({
senderThreadId: parentThreadId,
});
expect(promptedThread.link).toBe(
`[${createdSource.thread.title}](t3-thread://v1/environment%3Amcp-orchestrator/${encodeURIComponent(promptedThread.threadId)})`,
);
expect(
createdSource.turnItems.find((item) => item.type === "user_message"),
).toMatchObject({
Expand Down
79 changes: 79 additions & 0 deletions apps/server/src/mcp/toolkits/core.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -474,6 +474,85 @@ it.effect("a client caller targets any thread within its ceiling and cannot act
),
);

// Serves "source-thread" as the fork source; `forkShell` answers the fork's own shell read.
const forkToolkitLayer = (
forkShell: ThreadManagement.ThreadManagementService["Service"]["getThreadShell"],
) =>
McpHttpServer.layerThreadToolkit.pipe(
Layer.provideMerge(McpServer.McpServer.layer),
Layer.provide(NodeCrypto.layer),
Layer.provide(
Layer.mock(ThreadManagement.ThreadManagementService)({
getThreadShell: (id) =>
id === "source-thread"
? Effect.succeed(McpToolAccessTestkit.liveThreadShell(id, { runtimeMode: "auto" }))
: forkShell(id),
getProjectThreadRecords: () =>
Effect.succeed({
thread: {
id: ThreadId.make("source-thread"),
projectId: "project-a",
title: "Source title",
runtimeMode: "auto",
interactionMode: "default",
deletedAt: null,
},
} as never),
dispatch: () => Effect.succeed({ sequence: 7, storedEvents: [] }),
}),
),
);

const forkSourceThread = (args: Record<string, unknown> = {}) =>
Effect.gen(function* () {
const server = yield* McpServer.McpServer;
const forked = yield* server
.callTool({
name: "t3_thread_fork",
arguments: { threadId: "source-thread", sourcePoint: { type: "latest_stable" }, ...args },
})
.pipe(
Effect.provideService(McpInvocationContext.McpInvocationContext, clientScope("auto")),
Effect.provideService(McpSchema.McpServerClient, client),
);
expect(forked.isError).toBe(false);
return forked.structuredContent as { targetThreadId: string; link: string };
});

it.effect("a fork result links to the new fork, not its source", () =>
Effect.gen(function* () {
const { targetThreadId, link } = yield* forkSourceThread();
expect(targetThreadId).not.toBe("source-thread");
expect(link).toBe(
`[Fork title](t3-thread://v1/mcp-core-environment/${encodeURIComponent(targetThreadId)})`,
);
}).pipe(
Effect.provide(
forkToolkitLayer((id) =>
Effect.succeed({
...McpToolAccessTestkit.liveThreadShell(id, { runtimeMode: "auto" }),
title: "Fork title",
}),
),
),
),
);

it.effect("a committed fork still returns its link when reading its title fails", () =>
Effect.gen(function* () {
const { targetThreadId, link } = yield* forkSourceThread({ title: "Requested title" });
expect(link).toBe(
`[Requested title](t3-thread://v1/mcp-core-environment/${encodeURIComponent(targetThreadId)})`,
);
}).pipe(
Effect.provide(
forkToolkitLayer((threadId) =>
Effect.fail(new OrchestratorProjectionError({ threadId, cause: "read failed" })),
),
),
),
);

it.effect("a read-only client reads threads and is refused every write before it runs", () =>
Effect.gen(function* () {
const server = yield* McpServer.McpServer;
Expand Down
6 changes: 3 additions & 3 deletions apps/server/src/mcp/toolkits/orchestrator/tools.ts
Original file line number Diff line number Diff line change
Expand Up @@ -63,7 +63,7 @@ const OrchestratorCapabilitiesTool = Tool.make("orchestrator_capabilities", {

export const DelegateTaskTool = Tool.make("delegate_task", {
description:
"Needs an agent running inside a T3 thread. Delegate one task to a T3-owned child agent/subagent of THIS thread and run it with only the supplied task prompt, without copying parent conversation history. Choose providers and models from orchestrator_capabilities, which uses the same live catalog as the composer. Prefer native subagent tools for same-provider work only when they support the chosen model. Use this for any model missing from the native tool, including same-provider work, for cross-provider work, or for explicitly T3-owned child tasks. For every T3 delegated review round, call delegate_task again with the original brief, prior findings, responses, and unresolved objections in the task prompt. Track each round by its own taskId and use a distinct clientRequestId per round, stable across retries of that round. The childThreadId is backing storage, not the target for starting another delegated review round through t3_thread_send. Provider, model, model options (see orchestrator_capabilities), runtime mode, and interaction mode inherit unless target overrides them. Prefer mode='async' for long work; mode='wait' blocks until completion or timeout. timeoutMs on mode=wait is only the parent's wait budget and does not cancel the child. waitTimedOut on that wait call means the timeout fired; keep that taskId and read status on later task_status. An async child's completion wakes this thread through a notification, steered into active turns where supported or queued otherwise, so end the turn instead of polling or spawning watchers; use task_status only when the result is needed mid-turn.",
"Needs an agent running inside a T3 thread. Delegate one task to a T3-owned child agent/subagent of THIS thread and run it with only the supplied task prompt, without copying parent conversation history. Choose providers and models from orchestrator_capabilities, which uses the same live catalog as the composer. Prefer native subagent tools for same-provider work only when they support the chosen model. Use this for any model missing from the native tool, including same-provider work, for cross-provider work, or for explicitly T3-owned child tasks. For every T3 delegated review round, call delegate_task again with the original brief, prior findings, responses, and unresolved objections in the task prompt. Track each round by its own taskId and use a distinct clientRequestId per round, stable across retries of that round. The childThreadId is backing storage, not the target for starting another delegated review round through t3_thread_send. Paste childThreadLink when you mention the child thread so the user can open it. Provider, model, model options (see orchestrator_capabilities), runtime mode, and interaction mode inherit unless target overrides them. Prefer mode='async' for long work; mode='wait' blocks until completion or timeout. timeoutMs on mode=wait is only the parent's wait budget and does not cancel the child. waitTimedOut on that wait call means the timeout fired; keep that taskId and read status on later task_status. An async child's completion wakes this thread through a notification, steered into active turns where supported or queued otherwise, so end the turn instead of polling or spawning watchers; use task_status only when the result is needed mid-turn.",
parameters: OrchestratorMcpDelegateTaskInput,
success: OrchestratorMcpDelegateTaskResult,
failure: OrchestratorMcpFailure,
Expand All @@ -76,7 +76,7 @@ export const DelegateTaskTool = Tool.make("delegate_task", {

const TaskStatusTool = Tool.make("task_status", {
description:
"Needs an agent running inside a T3 thread. Read a T3-owned delegated task created by this parent thread. childRunId identifies the original delegated run. workState distinguishes working, waiting_for_children, and result_available; a completed turn with live nested work is not a completed task. summary is the final task result, including provider errors on failure, and remains stable after publication. hasPendingChildRuns reports later queued or executing turns in the backing child thread, even after the task is terminal; it does not reopen the task, and task_cancel stops those turns too. latestTerminal* provides later non-monitor turn results. Reading a terminal result acknowledges its automatic parent delivery.",
"Needs an agent running inside a T3 thread. Read a T3-owned delegated task created by this parent thread. childRunId identifies the original delegated run. workState distinguishes working, waiting_for_children, and result_available; a completed turn with live nested work is not a completed task. summary is the final task result, including provider errors on failure, and remains stable after publication. hasPendingChildRuns reports later queued or executing turns in the backing child thread, even after the task is terminal; it does not reopen the task, and task_cancel stops those turns too. latestTerminal* provides later non-monitor turn results. Reading a terminal result acknowledges its automatic parent delivery. Paste childThreadLink when you mention the child thread so the user can open it.",
parameters: OrchestratorMcpTaskStatusInput,
success: OrchestratorMcpDelegateTaskResult,
failure: OrchestratorMcpFailure,
Expand Down Expand Up @@ -165,7 +165,7 @@ const RequestSecretTool = Tool.make("request_secret", {

export const CreateThreadsTool = Tool.make("create_threads", {
description:
"Needs an agent running inside a T3 thread. Create one or more ORDINARY TOP-LEVEL T3 conversations. This is not delegation and does not create child agents/subagents. For delegated work, choose models from orchestrator_capabilities. Prefer native subagents only when they support the chosen model; otherwise call delegate_task, including for same-provider work. Use create_threads for a batch of separate top-level threads sharing this checkout. Prefer t3_thread_launch for a single thread. Both require the user to request separate/new/top-level threads or conversations. Each entry may override provider, model, options, runtime mode, and interaction mode; omitted settings inherit. Project, branch, and worktree always inherit and cannot be overridden here. For independent implementation or a PR stack in its own worktree, use t3_thread_launch with workspaceStrategy instead of asking the agent to create a worktree in its prompt.",
"Needs an agent running inside a T3 thread. Create one or more ORDINARY TOP-LEVEL T3 conversations. This is not delegation and does not create child agents/subagents. For delegated work, choose models from orchestrator_capabilities. Prefer native subagents only when they support the chosen model; otherwise call delegate_task, including for same-provider work. Use create_threads for a batch of separate top-level threads sharing this checkout. Prefer t3_thread_launch for a single thread. Both require the user to request separate/new/top-level threads or conversations. Each entry may override provider, model, options, runtime mode, and interaction mode; omitted settings inherit. Project, branch, and worktree always inherit and cannot be overridden here. For independent implementation or a PR stack in its own worktree, use t3_thread_launch with workspaceStrategy instead of asking the agent to create a worktree in its prompt. Paste each returned link when you mention a thread.",
parameters: OrchestratorMcpCreateThreadsInput,
success: OrchestratorMcpCreateThreadsResult,
failure: OrchestratorMcpFailure,
Expand Down
17 changes: 15 additions & 2 deletions apps/server/src/mcp/toolkits/thread/handlers.ts
Original file line number Diff line number Diff line change
Expand Up @@ -9,6 +9,7 @@ import {
} from "@t3tools/contracts";
import * as Effect from "effect/Effect";
import { modelSelectionCommandType } from "@t3tools/shared/model";
import { formatThreadLink } from "@t3tools/shared/threadLinks";

import * as McpToolAccess from "../../McpToolAccess.ts";
import {
Expand Down Expand Up @@ -116,7 +117,7 @@ export const layer = McpToolAccess.toLayer(ThreadToolkit, {
),
t3_thread_fork: writesThread((input) =>
Effect.gen(function* () {
const { threads, projection } = yield* readThread(input.threadId);
const { scope, threads, projection } = yield* readThread(input.threadId);
const commandId = yield* newCommandId();
const targetThreadId = ThreadId.make(`${commandId}:fork`);
const result = yield* threads
Expand All @@ -131,7 +132,19 @@ export const layer = McpToolAccess.toLayer(ThreadToolkit, {
creationSource: "mcp",
})
.pipe(Effect.mapError(dispatchFailure));
return { sequence: result.sequence, targetThreadId };
// The fork is committed, so a failed title read must not fail the call and invite a retry.
const fork = yield* threads
.getThreadShell(targetThreadId)
.pipe(Effect.orElseSucceed(() => null));
return {
sequence: result.sequence,
targetThreadId,
link: formatThreadLink({
environmentId: scope.environmentId,
threadId: targetThreadId,
title: fork?.title ?? input.title ?? projection.thread.title,
}),
};
}),
),
t3_thread_merge_back: McpToolAccess.writesThreads(
Expand Down
9 changes: 7 additions & 2 deletions apps/server/src/mcp/toolkits/thread/tools.ts
Original file line number Diff line number Diff line change
Expand Up @@ -196,16 +196,21 @@ const ThreadConfigureTool = Tool.make("t3_thread_configure", {
}).annotate(Tool.Destructive, true);

const transferResult = Schema.Struct({ sequence: NonNegativeInt, targetThreadId: ThreadId });
const forkResult = Schema.Struct({
...transferResult.fields,
/** Paste this whenever you mention the fork, so the user can click to open it. */
link: Schema.String,
});
const ThreadForkTool = Tool.make("t3_thread_fork", {
...commandTool,
description:
"Fork a thread from a stable run or checkpoint using the existing fork command. Omit threadId to fork this thread. The fork inherits the source configuration. Acceptance does not mean a provider turn has completed.",
"Fork a thread from a stable run or checkpoint using the existing fork command. Omit threadId to fork this thread. The fork inherits the source configuration. Acceptance does not mean a provider turn has completed. Paste the returned link when you mention the fork.",
parameters: Schema.Struct({
threadId: Schema.optional(ThreadId),
sourcePoint: OrchestrationV2ThreadForkSourcePoint,
title: Schema.optional(TrimmedNonEmptyString),
}),
success: transferResult,
success: forkResult,
}).annotate(Tool.Destructive, true);
const ThreadMergeBackTool = Tool.make("t3_thread_merge_back", {
...commandTool,
Expand Down
8 changes: 5 additions & 3 deletions docs/orchestration-v2/orchestrator-mcp-server.md
Original file line number Diff line number Diff line change
Expand Up @@ -249,6 +249,7 @@ the timeout expires. A wait timeout does not cancel the child; the result sets
type DelegateTaskResult = {
taskId: string;
childThreadId: string;
childThreadLink: string;
childRunId: string | null;
childNodeId: string;
status: "queued" | "running" | "waiting" | "completed" | "failed" | "cancelled" | "interrupted";
Expand Down Expand Up @@ -352,9 +353,10 @@ and `creationSource: "mcp"`; provider output uses `creationSource: "provider"`.
Actor and ingress are separate so agent-authored user-role messages remain
distinguishable from human-authored messages.

List, read, and launch results include `link`, a Markdown link of the form
`[title](t3-thread://v1/<environmentId>/<threadId>)` that clients open as the
thread. List and read results also report `snoozed` and `snoozedUntil`, and
List, read, launch, `create_threads`, and fork results include `link`, and
`delegate_task` and `task_status` results include `childThreadLink`: a Markdown
Comment thread
coderabbitai[bot] marked this conversation as resolved.
link of the form `[title](t3-thread://v1/<environmentId>/<threadId>)` that
clients open as the thread. List and read results also report `snoozed` and `snoozedUntil`, and
`t3_thread_list` filters on `snoozed`. The server's `isSnoozed` follows the
client's `effectiveSnoozed`, so agents and the sidebar agree: a snoozed thread
wakes early when it has a pending request, fails, or completes after the snooze.
Expand Down
Loading
Loading