diff --git a/apps/server/src/mcp/OrchestratorMcpToolkit.integration.test.ts b/apps/server/src/mcp/OrchestratorMcpToolkit.integration.test.ts index be29178896f8..89415684a64e 100644 --- a/apps/server/src/mcp/OrchestratorMcpToolkit.integration.test.ts +++ b/apps/server/src/mcp/OrchestratorMcpToolkit.integration.test.ts @@ -765,6 +765,20 @@ describe("orchestrator MCP toolkit", () => { const invoke = (name: string, args: Record) => 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(); @@ -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* () { diff --git a/apps/server/src/mcp/threadAccess.ts b/apps/server/src/mcp/threadAccess.ts index a965b0823629..e0fb29b7e13a 100644 --- a/apps/server/src/mcp/threadAccess.ts +++ b/apps/server/src/mcp/threadAccess.ts @@ -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"; @@ -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 diff --git a/apps/server/src/mcp/toolkits/core.test.ts b/apps/server/src/mcp/toolkits/core.test.ts index b3930f9c631d..2d907a92c4a2 100644 --- a/apps/server/src/mcp/toolkits/core.test.ts +++ b/apps/server/src/mcp/toolkits/core.test.ts @@ -4,6 +4,7 @@ import { expect, it } from "@effect/vitest"; import { DEFAULT_SERVER_SETTINGS, ChatImageAttachment, + CommandId, EnvironmentId, ProviderInstanceId, RunId, @@ -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"; @@ -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"; @@ -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; diff --git a/apps/server/src/mcp/toolkits/thread/handlers.ts b/apps/server/src/mcp/toolkits/thread/handlers.ts index 0ea058f78130..803402f4621f 100644 --- a/apps/server/src/mcp/toolkits/thread/handlers.ts +++ b/apps/server/src/mcp/toolkits/thread/handlers.ts @@ -11,6 +11,7 @@ import * as Effect from "effect/Effect"; import { modelSelectionCommandType } from "@t3tools/shared/model"; import { + dispatchFailure, newCommandId, readCaller, readFullAccessCaller, @@ -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 }; }); @@ -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) => @@ -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) => @@ -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) => @@ -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) => @@ -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 }; }), });