Skip to content

Commit c7aeab8

Browse files
CahidArdaclaude
andauthored
DX-3022: mcp-toolkit: isError fails the task, keep failure details from the model (#62)
* DX-3022: mcp-toolkit: isError fails the task, keep failure details from the model; README matches the docs - A handler returning { isError: true } now settles the task failed (no retry); task_status returns its content with isError, as the synchronous tool would - structuredContent carries only error.code and error.message: error.data (DLQ id, the failing response, Workflow's failResponse) stays in the store - A dispatch failure is logged and answered with a generic tool error instead of rethrowing the transport's message - README: lib/auth.ts with verifyToken + principal, mcp-handler withMcpAuth route, QSTASH_URL; changeset mentions the failure behavior Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01Gtucop3k8CNVLjgvoDTC8n * DX-3022: mcp-toolkit: task_cancel cancels the Workflow run; never store what the handler threw - Client.trigger turns workflowRunId into wfr_<id> but Client.cancel takes the id as is, so cancel sent the bare task id and matched no run: the run kept executing. Cancel now sends runIdOf(taskId). - Tested against the real @upstash/workflow client with fetch stubbed (fails on the old code), not a stubbed client. - failureFunction logs failResponse (the thrown error's message) instead of storing it in error.data, so it reaches neither Redis nor, via task_status, the task's owner. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01Gtucop3k8CNVLjgvoDTC8n * DX-3022: mcp-toolkit: fix lint (no RequestInit global in tests) Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01Gtucop3k8CNVLjgvoDTC8n --------- Co-authored-by: Claude Opus 5.5 <noreply@anthropic.com>
1 parent 5644daa commit c7aeab8

7 files changed

Lines changed: 173 additions & 40 deletions

File tree

‎.changeset/mcp-toolkit-initial.md‎

Lines changed: 4 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -24,6 +24,10 @@ limit of live subscriptions per subscriber (8 by default), SSRF checks on callba
2424
secrets encrypted at rest. `QStashDelivery` signs each attempt with
2525
Standard Webhooks and lets QStash retry failures (`export const POST = events.createDeliveryHandler()`).
2626

27+
A handler that returns a tool error (`isError: true`) fails its task without a retry, and
28+
`task_status` shows that error the way the synchronous tool would. The model only ever sees a
29+
failure's `code` and `message`: transport details stay in the store, and dispatch errors are logged.
30+
2731
The package uses WebCrypto only, so it runs on Node and edge runtimes.
2832

2933
`@upstash/mcp-toolkit/upstash` holds the Upstash backends for both: `RedisTaskStore`,

‎packages/mcp-toolkit/README.md‎

Lines changed: 36 additions & 25 deletions
Original file line numberDiff line numberDiff line change
@@ -12,14 +12,15 @@ Durable building blocks for MCP servers built on the official TypeScript SDK
1212
`WorkflowDispatcher`, `RedisSubscriptionStore`, `QStashDelivery`.
1313

1414
```bash
15-
npm install @upstash/mcp-toolkit @modelcontextprotocol/server zod
15+
npm install @upstash/mcp-toolkit @modelcontextprotocol/server mcp-handler zod
1616
```
1717

1818
Environment variables:
1919

2020
```bash
2121
UPSTASH_REDIS_REST_URL=...
2222
UPSTASH_REDIS_REST_TOKEN=...
23+
QSTASH_URL=... # your QStash region's URL, from the console
2324
QSTASH_TOKEN=...
2425
QSTASH_CURRENT_SIGNING_KEY=... # required: delivery routes refuse to run without them
2526
QSTASH_NEXT_SIGNING_KEY=...
@@ -30,14 +31,24 @@ APP_URL=https://your-app.com # used in the snippets below. QStash has to be
3031
## Who is calling
3132

3233
Both layers need `principal`: a function that returns the id of the user making the call. Tasks
33-
belong to the user who started them, and subscriptions to the user who subscribed. Read it from the
34-
token your MCP route has verified, and throw when there is none:
34+
belong to the user who started them, and subscriptions to the user who subscribed. Keep it next to
35+
the function that verifies the token, so the whole path is in one file:
3536

3637
```ts
3738
// lib/auth.ts
3839
import type { AuthInfo } from "@modelcontextprotocol/server";
3940

40-
// `auth` is the AuthInfo your MCP route verified and passed to the SDK (see below).
41+
// 1. Runs on every MCP request (wired up with withMcpAuth below). Returning undefined answers 401.
42+
export async function verifyToken(
43+
req: Request,
44+
bearerToken?: string,
45+
): Promise<AuthInfo | undefined> {
46+
if (!bearerToken) return undefined;
47+
const { sub, client_id } = await verifyJwt(bearerToken); // Clerk, WorkOS, Auth0, your own
48+
return { token: bearerToken, clientId: client_id, scopes: [], extra: { userId: sub } };
49+
}
50+
51+
// 2. The toolkit calls this with the AuthInfo above. Return the user id, or throw.
4152
export function principal({ auth }: { auth?: AuthInfo }): string {
4253
const userId = auth?.extra?.userId;
4354
if (typeof userId !== "string") throw new Error("Not authenticated");
@@ -52,8 +63,9 @@ export function principal({ auth }: { auth?: AuthInfo }): string {
5263
returns anything but a non-empty string, the call is refused as not authenticated.
5364
- **Use the user id, not `auth.clientId`.** The client id identifies the OAuth app, and every
5465
ChatGPT user shares the same one.
55-
- `auth` is only what your route passed to `handler.fetch(request, { authInfo })` (see below). The
56-
SDK never fills it from headers.
66+
- `auth` is exactly what `verifyToken` returned: `withMcpAuth` runs it and the toolkit hands the
67+
result to `principal`. Nothing is read from headers on its own. Without `mcp-handler`, pass it
68+
yourself with the SDK's `createMcpHandler`: `handler.fetch(request, { authInfo })`.
5769
- `principal` also receives `request`, for cookie or session apps: type its argument as `Caller`
5870
(`{ auth, request }`, from `@upstash/mcp-toolkit/tasks` or `/events`). `request` is unverified,
5971
so check the session yourself, and never trust a header like `x-user-id`.
@@ -65,7 +77,6 @@ export function principal({ auth }: { auth?: AuthInfo }): string {
6577

6678
```ts
6779
// lib/tasks.ts
68-
import { McpServer } from "@modelcontextprotocol/server";
6980
import { createTaskLayer } from "@upstash/mcp-toolkit/tasks";
7081
import { QStashDispatcher, RedisTaskStore } from "@upstash/mcp-toolkit/upstash";
7182
import * as z from "zod";
@@ -86,32 +97,26 @@ tasks.define(
8697
return { content: [{ type: "text", text: await writeReport(topic) }] };
8798
},
8899
);
89-
90-
export function createServer() {
91-
const server = new McpServer({ name: "reports", version: "1.0.0" });
92-
tasks.register(server); // adds generate_report, task_status and task_cancel
93-
return server;
94-
}
95100
```
96101

97102
```ts
98103
// app/api/mcp/route.ts
99-
import { createMcpHandler, type AuthInfo } from "@modelcontextprotocol/server";
100-
import { createServer } from "../../lib/tasks";
104+
import { createMcpHandler, withMcpAuth } from "mcp-handler";
105+
import { verifyToken } from "@/lib/auth";
106+
import { tasks } from "@/lib/tasks";
101107

102-
const handler = createMcpHandler(() => createServer());
108+
const handler = createMcpHandler((server) => {
109+
tasks.register(server); // adds generate_report, task_status and task_cancel
110+
});
103111

104-
export async function POST(request: Request) {
105-
// Clerk, WorkOS, Auth0, your own. `principal` reads the user id from `authInfo.extra`.
106-
const authInfo: AuthInfo | undefined = await verifyToken(request);
107-
if (!authInfo) return new Response("Unauthorized", { status: 401 });
108-
return handler.fetch(request, { authInfo });
109-
}
112+
// Runs verifyToken before any tool sees the request; no AuthInfo means 401.
113+
const authHandler = withMcpAuth(handler, verifyToken, { required: true });
114+
export { authHandler as GET, authHandler as POST };
110115
```
111116

112117
```ts
113118
// app/api/execute/route.ts: QStash delivers the work here
114-
import { tasks } from "../../lib/tasks";
119+
import { tasks } from "@/lib/tasks";
115120

116121
export const POST = tasks.createExecuteHandler();
117122
```
@@ -238,6 +243,12 @@ can take longer.
238243
Without signing keys the route throws instead of running anything unverified.
239244
- When your handler throws, the task is not marked failed. Only the dispatcher marks it `failed`,
240245
and only after QStash has stopped retrying. The failed message stays in the QStash DLQ.
246+
- When your handler _returns_ a tool error (`isError: true`), the task is `failed` at once, with no
247+
retry, and `task_status` returns your content with `isError`, as the synchronous tool would.
248+
- What a failure looks like to the model is `error.code` and `error.message`. The transport's
249+
details (`error.data`: the QStash DLQ id, the Workflow run id) stay in Redis for you. What your
250+
handler threw (Workflow's failure response) and a dispatch error are logged, never stored or
251+
returned.
241252
- By default QStash tries 5 times with backoff `min(pow(3, retried) * 1000, 300000)`, about two
242253
minutes in total, so a task survives a server restart. The free tier and the local dev server
243254
allow at most 5 retries.
@@ -279,12 +290,12 @@ export const commentCreated = events.define("comment.created", {
279290
});
280291
```
281292

282-
Call `events.register(server)` in `createServer()`, next to `tasks.register(server)`, and add the
293+
Call `events.register(server)` in the `createMcpHandler` callback, next to `tasks.register(server)`, and add the
283294
route QStash delivers to:
284295

285296
```ts
286297
// app/api/events/route.ts
287-
import { events } from "../../lib/events";
298+
import { events } from "@/lib/events";
288299

289300
export const POST = events.createDeliveryHandler();
290301
```

‎packages/mcp-toolkit/src/tasks/backends/workflow.test.ts‎

Lines changed: 37 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -6,8 +6,8 @@
66
* `task.run(...)` becomes a journaled step under Workflow and a plain call without it.
77
*/
88
import { Client as QStashClient } from "@upstash/qstash";
9-
import { WorkflowContext } from "@upstash/workflow";
10-
import { describe, expect, it } from "vitest";
9+
import { Client as WorkflowClient, WorkflowContext } from "@upstash/workflow";
10+
import { describe, expect, it, vi } from "vitest";
1111
import * as z from "zod";
1212
import { insideStep, WorkflowDispatcher, workflowRoute } from "./workflow.js";
1313
import { createTaskLayer } from "../core.js";
@@ -95,7 +95,41 @@ describe("WorkflowDispatcher", () => {
9595
await dispatcher.cancel("task-1");
9696

9797
// Unlike a queue, a workflow run can be stopped mid-flight rather than only un-queued.
98-
expect(cancelled).toEqual(["task-1"]);
98+
expect(cancelled).toEqual(["wfr_task-1"]);
99+
});
100+
101+
it("cancels the run it triggered, with the real Workflow client", async () => {
102+
// The real client, only its HTTP stubbed: a stubbed client can't catch how the real one
103+
// names runs (`trigger` prefixes `wfr_`, `cancel` doesn't).
104+
const requests: { method: string; url: string; body: string; headers: string }[] = [];
105+
vi.stubGlobal(
106+
"fetch",
107+
async (input: string | URL | Request, init?: Parameters<typeof fetch>[1]) => {
108+
const url = String(input instanceof Request ? input.url : input);
109+
const method = init?.method ?? "GET";
110+
const headers = JSON.stringify(Object.fromEntries(new Headers(init?.headers).entries()));
111+
requests.push({ method, url, body: String(init?.body ?? ""), headers });
112+
return method === "DELETE"
113+
? Response.json({ cancelled: 1 })
114+
: Response.json([{ messageId: "msg_1" }]);
115+
},
116+
);
117+
try {
118+
const client = new WorkflowClient({ token: "test-token", baseUrl: "https://qstash.test" });
119+
const dispatcher = new WorkflowDispatcher({ url: URL, client });
120+
const taskId = "6f1c2c5e-0000-4000-8000-000000000000";
121+
122+
await dispatcher.dispatch(task({ taskId }));
123+
await dispatcher.cancel(taskId);
124+
125+
const runId = `wfr_${taskId}`;
126+
const trigger = requests.find((r) => r.method !== "DELETE");
127+
const cancel = requests.find((r) => r.method === "DELETE");
128+
expect(`${trigger?.body}${trigger?.headers}`).toContain(runId);
129+
expect(new globalThis.URL(cancel!.url).searchParams.get("workflowRunIds")).toBe(runId);
130+
} finally {
131+
vi.unstubAllGlobals();
132+
}
99133
});
100134
});
101135

‎packages/mcp-toolkit/src/tasks/backends/workflow.ts‎

Lines changed: 16 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -63,14 +63,14 @@ export class WorkflowDispatcher implements TaskDispatcher<WorkflowContext<Workfl
6363
url: this.config.url,
6464
body: { taskId: task.taskId } satisfies WorkflowPayload,
6565
retries: this.config.retries,
66-
// One run per task: Workflow refuses a run id that was already used.
66+
// One run per task: Workflow refuses a run id that was already used. It prefixes this with
67+
// `wfr_` itself, so the run is `runIdOf(taskId)`.
6768
workflowRunId: task.taskId,
6869
});
6970
}
7071

7172
async cancel(taskId: string): Promise<void> {
72-
// The run id is the task id.
73-
await this.client().cancel(taskId);
73+
await this.client().cancel(runIdOf(taskId));
7474
}
7575

7676
/** The workflow endpoint. Its `failureFunction` settles the task `failed` once retries run out. */
@@ -97,10 +97,16 @@ export class WorkflowDispatcher implements TaskDispatcher<WorkflowContext<Workfl
9797
failureFunction: async ({ context, failStatus, failResponse }) => {
9898
const taskId = (context.requestPayload as WorkflowPayload | undefined)?.taskId;
9999
if (!taskId) return;
100+
// `failResponse` is the thrown error's message: database errors, internal hosts, upstream
101+
// responses. It goes to your logs, never into the task, which its owner can read.
102+
console.error(
103+
`[mcp-toolkit] task ${taskId} failed (workflow run ${context.workflowRunId}):`,
104+
failResponse,
105+
);
100106
await endpoints.fail(taskId, {
101107
code: INTERNAL_ERROR,
102108
message: `Workflow run failed${failStatus ? ` (status ${failStatus})` : ""}`,
103-
data: { response: failResponse, workflowRunId: context.workflowRunId },
109+
data: { workflowRunId: context.workflowRunId },
104110
});
105111
},
106112
// Workflow verifies only body and signature; binding the URL refuses a signature that
@@ -114,6 +120,12 @@ export class WorkflowDispatcher implements TaskDispatcher<WorkflowContext<Workfl
114120
}
115121
}
116122

123+
/**
124+
* The Workflow run that serves a task. `Client.trigger` turns the `workflowRunId` it is given into
125+
* `wfr_<id>`, while `Client.cancel` takes the run id as is, so a cancel must add the prefix.
126+
*/
127+
export const runIdOf = (taskId: string): string => `wfr_${taskId}`;
128+
117129
/**
118130
* The workflow's route function. Its first act is always a step: Workflow authorizes every request,
119131
* the failure callback included, by running the route function until its first step, and refuses

‎packages/mcp-toolkit/src/tasks/core.test.ts‎

Lines changed: 45 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -543,6 +543,48 @@ describe("createTaskLayer over MCP", () => {
543543
expect(text(after)).toMatch(/permanent/);
544544
});
545545

546+
it("fails a task whose handler returns a tool error, and shows that error like a tool would", async () => {
547+
live = await harness(async () => ({
548+
isError: true,
549+
content: [{ type: "text", text: "Workspace is read-only" }],
550+
}));
551+
const taskId = await live.start("x");
552+
await live.dispatcher.drain();
553+
const after = await live.status(taskId);
554+
expect(after.isError).toBe(true);
555+
expect(after.structuredContent).toMatchObject({
556+
status: "failed",
557+
// Settled by the layer, not by a retried throw (whose message would be the thrown one).
558+
error: { message: "The task returned an error" },
559+
});
560+
expect(text(after)).toMatch(/Workspace is read-only/);
561+
});
562+
563+
it("keeps a failure's details in the store but out of what the model sees", async () => {
564+
live = await harness(steppedHandler(1, 1));
565+
const taskId = await live.start("x");
566+
await live.dispatcher.drain();
567+
const stored = (await live.store.get(taskId))!;
568+
await live.store.create({
569+
...stored,
570+
status: "failed",
571+
error: {
572+
code: -32603,
573+
message: "Workflow run failed",
574+
data: { response: "pg: password auth failed" },
575+
},
576+
});
577+
const after = await live.status(taskId);
578+
expect(after.structuredContent?.error).toEqual({
579+
code: -32603,
580+
message: "Workflow run failed",
581+
});
582+
expect(JSON.stringify(after)).not.toMatch(/password auth/);
583+
expect((await live.store.get(taskId))?.error?.data).toEqual({
584+
response: "pg: password auth failed",
585+
});
586+
});
587+
546588
it("fails the task when it cannot be dispatched, instead of leaving it working", async () => {
547589
const manual = new ManualDispatcher();
548590
manual.failNext = new Error("queue unavailable");
@@ -555,6 +597,9 @@ describe("createTaskLayer over MCP", () => {
555597
};
556598
const result = await live.call("generate_report", { topic: "x" });
557599
expect(result.isError).toBe(true);
600+
// The transport's error goes to the logs, not to the model.
601+
expect(text(result)).not.toMatch(/queue unavailable/);
602+
expect(text(result)).toMatch(/could not be queued/);
558603
expect(await live.store.get(created!)).toMatchObject({
559604
status: "failed",
560605
statusMessage: "Could not be queued",

‎packages/mcp-toolkit/src/tasks/core.ts‎

Lines changed: 34 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -250,7 +250,9 @@ export function createTaskLayer<TContext = unknown>(
250250
error: { code: INTERNAL_ERROR, message: "The task could not be dispatched" },
251251
})
252252
.catch(() => undefined);
253-
throw error;
253+
// The transport's error (a QStash response, a URL) is for your logs, not the model.
254+
console.error(`[mcp-toolkit] could not dispatch task ${task.taskId}:`, error);
255+
return errorResult("The task could not be queued. Try again later.");
254256
}
255257
return {
256258
content: [
@@ -294,6 +296,17 @@ export function createTaskLayer<TContext = unknown>(
294296

295297
const result = await definition.handler(task.args, mergeContext(taskContext, context));
296298
// A cancel that landed meanwhile wins: settle refuses the second terminal write.
299+
if (result?.isError === true) {
300+
// The handler answered with a tool error, the way a synchronous tool would: the task failed,
301+
// and its own content explains why. Not retried: it returned, it didn't throw.
302+
await store.settle(taskId, {
303+
status: "failed",
304+
statusMessage: "Failed",
305+
result,
306+
error: { code: INTERNAL_ERROR, message: "The task returned an error" },
307+
});
308+
return;
309+
}
297310
await store.settle(taskId, {
298311
status: "completed",
299312
statusMessage: definition.config.completedMessage ?? "Completed",
@@ -336,18 +349,28 @@ function errorResult(text: string): Record<string, unknown> {
336349
return { isError: true, content: [{ type: "text", text }] };
337350
}
338351

339-
/** A status line, followed by the task's own result content once it completed. */
352+
/**
353+
* A status line, followed by the task's own result content once it completed, or once it failed by
354+
* returning a tool error. That last case also carries `isError`, as the synchronous tool would have.
355+
*/
340356
function statusResult(task: Task): Record<string, unknown> {
341357
const line = `Task ${task.taskId} is ${task.status}${task.statusMessage ? `: ${task.statusMessage}` : "."}`;
358+
const handlerError = task.status === "failed" && task.result?.isError === true;
342359
const extra =
343-
task.status === "completed" && Array.isArray(task.result?.content) ? task.result.content : [];
360+
(task.status === "completed" || handlerError) && Array.isArray(task.result?.content)
361+
? task.result.content
362+
: [];
344363
const text =
345364
task.status === "failed"
346365
? `${line} Error: ${task.error?.message ?? "unknown"}`
347366
: task.status === "working"
348367
? `${line} Check again in about ${seconds(task.pollIntervalMs)}.`
349368
: line;
350-
return { content: [{ type: "text", text }, ...extra], structuredContent: toWire(task) };
369+
return {
370+
content: [{ type: "text", text }, ...extra],
371+
structuredContent: toWire(task),
372+
...(handlerError ? { isError: true } : {}),
373+
};
351374
}
352375

353376
const seconds = (ms: number) => `${Math.max(1, Math.round(ms / 1000))}s`;
@@ -364,8 +387,12 @@ function mergeContext<TContext>(
364387
return Object.assign(supplied as object, taskContext) as TaskContext & TContext;
365388
}
366389

367-
/** Strips the server-only fields, leaving what the model sees in `structuredContent`. */
390+
/**
391+
* Strips the server-only fields, leaving what the model sees in `structuredContent`. That includes
392+
* `error.data`: transports put the failing response there (a stack trace, a database error), which
393+
* stays in the store for your logs.
394+
*/
368395
function toWire(task: Task): WireTask {
369-
const { name: _name, args: _args, owner: _owner, ...wire } = task;
370-
return wire;
396+
const { name: _name, args: _args, owner: _owner, error, ...wire } = task;
397+
return error ? { ...wire, error: { code: error.code, message: error.message } } : wire;
371398
}

‎packages/mcp-toolkit/src/tasks/types.ts‎

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -27,7 +27,7 @@ export type WireTask = {
2727
/** Retention window in milliseconds, from creation. */
2828
ttlMs: number;
2929
pollIntervalMs: number;
30-
/** The tool result, once `completed`. */
30+
/** The tool result, once `completed`, or once `failed` because the handler returned `isError`. */
3131
result?: Record<string, unknown>;
3232
/** Once `failed`. */
3333
error?: TaskError;

0 commit comments

Comments
 (0)