Skip to content
Merged
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
28 changes: 28 additions & 0 deletions apps/server/src/mcp/OrchestratorMcpToolkit.integration.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -765,6 +765,20 @@ describe("orchestrator MCP toolkit", () => {
const invoke = (name: string, args: Record<string, unknown>) =>
invokeAs(invocation, name, args);

const refusedSettle = yield* invoke("t3_thread_organize", { action: "settle" });
expect(refusedSettle.isError).toBe(true);
expect(refusedSettle.structuredContent).toBeUndefined();
expect(declaredFailure(refusedSettle)).toEqual({
_tag: "OrchestratorMcpFailure",
code: "orchestration_error",
message: `Thread ${parentThreadId} has active or blocked work and cannot be settled.`,
});
const afterRefusedSettle = yield* orchestrator.getThreadProjection(parentThreadId);
expect(afterRefusedSettle.thread.settledOverride).not.toBe("settled");
expect(afterRefusedSettle.runs.find((run) => run.id === parentRun?.id)?.status).toBe(
"running",
);

const pinned = yield* invoke("t3_thread_organize", { action: "pin" });
expect(pinned.structuredContent).toHaveProperty("sequence");
expect((yield* orchestrator.getThreadShell(parentThreadId))?.pinnedAt).not.toBeNull();
Expand All @@ -774,6 +788,20 @@ describe("orchestrator MCP toolkit", () => {
if (parentRun === undefined || parentRun.rootNodeId === null) {
return yield* Effect.die(new Error("Parent run missing."));
}
for (const name of ["t3_queue_edit", "t3_queue_cancel"]) {
const refusedQueueMutation = yield* invoke(name, {
queuedRunId: parentRun.id,
...(name === "t3_queue_edit" ? { text: "Keep the active turn." } : {}),
});
expect(refusedQueueMutation.isError).toBe(true);
expect(refusedQueueMutation.structuredContent).toBeUndefined();
expect(declaredFailure(refusedQueueMutation)).toEqual({
_tag: "OrchestratorMcpFailure",
code: "orchestration_error",
message: `Run ${parentRun.id} is not queued.`,
});
}

let parentRootNodeId = parentRun.rootNodeId;
const queueAutomaticCompletion = (suffix: string, taskText: string) =>
Effect.gen(function* () {
Expand Down
13 changes: 13 additions & 0 deletions apps/server/src/mcp/threadAccess.ts
Original file line number Diff line number Diff line change
Expand Up @@ -11,6 +11,7 @@ import {
import * as Crypto from "effect/Crypto";
import * as Effect from "effect/Effect";

import type { OrchestratorV2Error } from "../orchestration-v2/Orchestrator.ts";
import * as ThreadManagement from "../orchestration-v2/ThreadManagementService.ts";
import * as OrchestrationMcp from "./OrchestratorMcpService.ts";
import * as McpInvocationContext from "./McpInvocationContext.ts";
Expand All @@ -21,6 +22,18 @@ export const unavailable = () =>
message: "The operation could not be completed.",
});

/** Decider string rejections are public; wrapped storage and hydration causes are not. */
export const dispatchFailure = (error: OrchestratorV2Error) =>
(error._tag === "OrchestratorDispatchError" ||
error._tag === "OrchestratorCommandRejectedError") &&
typeof error.cause === "string" &&
error.cause.length > 0
? new OrchestratorMcpFailure({
code: "orchestration_error",
message: Array.from(error.cause).slice(0, 1000).join(""),
})
: unavailable();

/**
* The most a caller may hand to the threads it targets. A thread caller is
* capped by its own thread's modes; an OAuth client by the ceiling chosen
Expand Down
45 changes: 44 additions & 1 deletion apps/server/src/mcp/toolkits/core.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -4,6 +4,7 @@ import { expect, it } from "@effect/vitest";
import {
DEFAULT_SERVER_SETTINGS,
ChatImageAttachment,
CommandId,
EnvironmentId,
ProviderInstanceId,
RunId,
Expand All @@ -17,8 +18,13 @@ import { McpAttachmentInput } from "./attachment/input.ts";
import { McpSchema, McpServer, Tool } from "effect/ai";
import { FetchHttpClient } from "effect/http";

import {
OrchestratorCommandRejectedError,
OrchestratorDispatchError,
OrchestratorProjectionError,
} from "../../orchestration-v2/Orchestrator.ts";

import * as ServerConfig from "../../config.ts";
import { OrchestratorProjectionError } from "../../orchestration-v2/Orchestrator.ts";
import * as ProviderAdapterRegistry from "../../orchestration-v2/ProviderAdapterRegistry.ts";
import * as ThreadManagement from "../../orchestration-v2/ThreadManagementService.ts";
import * as PreviewBrowser from "../../preview/PreviewBrowser.ts";
Expand All @@ -28,6 +34,7 @@ import * as SecretRequests from "../../secrets/SecretRequests.ts";
import * as ScheduledTaskService from "../../scheduledTasks/ScheduledTaskService.ts";
import * as McpHttpServer from "../McpHttpServer.ts";
import * as McpInvocationContext from "../McpInvocationContext.ts";
import { dispatchFailure } from "../threadAccess.ts";
import { OrchestratorToolkit } from "./orchestrator/tools.ts";
import { PreviewToolkit } from "./preview/tools.ts";
import { PreviewControlsToolkit } from "./previewControls/tools.ts";
Expand Down Expand Up @@ -182,6 +189,42 @@ it.effect("returns a bounded public failure without serializing storage causes",
),
);

it("bounds public command rejections and redacts internal dispatch causes", () => {
const command = { commandId: CommandId.make("mcp-core-command"), commandType: "thread.settle" };
expect(
dispatchFailure(new OrchestratorDispatchError({ ...command, cause: "🙂".repeat(1001) }))
.message,
).toBe("🙂".repeat(1000));
expect(
dispatchFailure(
new OrchestratorCommandRejectedError({ ...command, cause: "Run is not queued." }),
).message,
).toBe("Run is not queued.");
for (const cause of [
undefined,
"",
new Error("private-storage-path"),
{ message: "private-storage-path" },
]) {
expect(dispatchFailure(new OrchestratorDispatchError({ ...command, cause }))).toMatchObject({
code: "orchestration_error",
message: "The operation could not be completed.",
});
expect(
dispatchFailure(new OrchestratorCommandRejectedError({ ...command, cause })),
).toMatchObject({
code: "orchestration_error",
message: "The operation could not be completed.",
});
}
expect(
dispatchFailure(new OrchestratorProjectionError({ threadId, cause: "private-storage-path" })),
).toMatchObject({
code: "orchestration_error",
message: "The operation could not be completed.",
});
});

it.effect("returns an HTML render reference that Codex and Claude tool rows both carry", () =>
Effect.gen(function* () {
const server = yield* McpServer.McpServer;
Expand Down
13 changes: 7 additions & 6 deletions apps/server/src/mcp/toolkits/thread/handlers.ts
Original file line number Diff line number Diff line change
Expand Up @@ -11,6 +11,7 @@ import * as Effect from "effect/Effect";
import { modelSelectionCommandType } from "@t3tools/shared/model";

import {
dispatchFailure,
newCommandId,
readCaller,
readFullAccessCaller,
Expand Down Expand Up @@ -45,7 +46,7 @@ const dispatch = Effect.fn("mcp.dispatchThreadCommand")(function* (
const { threads, projection } = yield* readWritableThread(threadId);
const result = yield* threads
.dispatch(command({ commandId: yield* newCommandId(), threadId: projection.thread.id }))
.pipe(Effect.mapError(unavailable));
.pipe(Effect.mapError(dispatchFailure));
return { sequence: result.sequence };
});

Expand Down Expand Up @@ -131,7 +132,7 @@ export const layer = ThreadToolkit.toLayer({
createdBy: "agent",
creationSource: "mcp",
})
.pipe(Effect.mapError(unavailable));
.pipe(Effect.mapError(dispatchFailure));
return { sequence: result.sequence, targetThreadId };
}),
t3_thread_merge_back: (input) =>
Expand All @@ -148,7 +149,7 @@ export const layer = ThreadToolkit.toLayer({
createdBy: "agent",
creationSource: "mcp",
})
.pipe(Effect.mapError(unavailable));
.pipe(Effect.mapError(dispatchFailure));
return { sequence: result.sequence, targetThreadId: input.targetThreadId };
}),
t3_thread_transfers: (input) =>
Expand Down Expand Up @@ -191,7 +192,7 @@ export const layer = ThreadToolkit.toLayer({
commandId: yield* newCommandId(),
modelSelection: input.modelSelection,
})
.pipe(Effect.mapError(unavailable));
.pipe(Effect.mapError(dispatchFailure));
return { sequence: result.sequence };
}),
t3_pending_request_list: (input) =>
Expand Down Expand Up @@ -219,7 +220,7 @@ export const layer = ThreadToolkit.toLayer({
requestId: input.requestId,
answers: input.answers,
})
.pipe(Effect.mapError(unavailable));
.pipe(Effect.mapError(dispatchFailure));
return { sequence: result.sequence };
}),
t3_queue_list: (input) =>
Expand Down Expand Up @@ -300,7 +301,7 @@ export const layer = ThreadToolkit.toLayer({
default:
command = { ...common, type: `thread.${input.action}` };
}
const result = yield* threads.dispatch(command).pipe(Effect.mapError(unavailable));
const result = yield* threads.dispatch(command).pipe(Effect.mapError(dispatchFailure));
return { sequence: result.sequence };
}),
});
Loading