Skip to content

Commit fd8ee97

Browse files
Upstash Boxclaude
andcommitted
DX-3022: mcp-toolkit: task authorize + principal, subscription limit, review fixes
- Tasks: `task.principal` (the owner) in TaskContext; optional `authorize(args, { principal, auth, request })` on tasks.define, run before the task is stored or queued - Tasks: reject non-positive / non-integer `defaults.ttlMs` and `pollIntervalMs` - Events: `maxSubscriptions` (default 8, Infinity = off) live subscriptions per subscriber, checked before the callback challenge and atomically in the store (Lua over idx:<event> and by:<subscriber>); SubscriptionStore gains count(), put() takes { limit } and returns boolean, delete() takes subscriber - Events: route and authorize deliveries by the payload's input fields parsed with `input`, so transforms apply on both sides; a payload that doesn't fit `input` makes emit throw - Drop QStash deduplicationId for tasks and events (Copilot: the sanitized event id was lossy); the event id still reaches the host as webhook-id - @upstash/redis floor ^1.38.4; aria-pressed on the demo's transport toggle Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01Gtucop3k8CNVLjgvoDTC8n
1 parent e14b0bf commit fd8ee97

18 files changed

Lines changed: 462 additions & 86 deletions

File tree

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

Lines changed: 6 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -11,14 +11,17 @@ lives in a `TaskStore` and the work runs behind a `TaskDispatcher`, and the disp
1111
delivery endpoint (`export const POST = tasks.createExecuteHandler()`) and always verifies QStash's
1212
signature against the URL it published to. Tasks are declared once at module scope with
1313
`tasks.define(...)` and attached to each request's server with `tasks.register(server)`. A required
14-
`principal` scopes tasks to their caller (and subscriptions to theirs).
14+
`principal` scopes tasks to their caller (and subscriptions to theirs); handlers get it as
15+
`task.principal`, and an optional `authorize` on `tasks.define` checks the arguments before a task
16+
is stored or queued.
1517

1618
`@upstash/mcp-toolkit/events` — MCP Events with webhook delivery, as ChatGPT ships it. Typed
1719
`events.define(...)` handles with `emit(payload)`: the payload's values route the event and are what
1820
the required `authorize` is checked against, at subscribe and again before every delivery.
1921
`events/list`, `events/subscribe` and `events/unsubscribe` are registered on your server, with a
20-
signed verification challenge, deterministic subscription ids, expiry and refresh, SSRF checks on
21-
callback URLs, and signing secrets encrypted at rest. `QStashDelivery` signs each attempt with
22+
signed verification challenge, deterministic subscription ids, expiry and refresh, a configurable
23+
limit of live subscriptions per subscriber (8 by default), SSRF checks on callback URLs, and signing
24+
secrets encrypted at rest. `QStashDelivery` signs each attempt with
2225
Standard Webhooks and lets QStash retry failures (`export const POST = events.createDeliveryHandler()`).
2326

2427
The package uses WebCrypto only, so it runs on Node and edge runtimes.

‎CLAUDE.md‎

Lines changed: 12 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -762,11 +762,10 @@ implements TanStack AI's own backend contracts (see its section below) — keep
762762
`dispatcher.cancel(taskId)` when the returned status is `cancelled` — `cancel` must be
763763
idempotent. There is no `dispatchId`: `QStashDispatcher.cancel` is a no-op (a redelivery of a
764764
cancelled task finds it settled and does nothing), Workflow's run id is the task id.
765-
- **Dispatch dedupe on the task id:** QStash `deduplicationId` and Workflow `workflowRunId` are the
766-
task id. That is only safe because ids
767-
are random: if keyed (deterministic) ids ever come back, dedupe must move to a per-record key
768-
(taskId + createdAt), since QStash keeps dedup ids for 10 minutes and Workflow refuses a reused
769-
run id. A failed `dispatch` settles the fresh record `failed` ("Could not be queued") and rethrows.
765+
- **No QStash `deduplicationId`** (tasks or events): the layer dispatches each task once, and an
766+
event's id reaches the host as `webhook-id`, which is where Standard Webhooks dedupes. Workflow's
767+
`workflowRunId` is the task id; that is only safe because ids are random (Workflow refuses a
768+
reused run id). A failed `dispatch` settles the fresh record `failed` ("Could not be queued") and rethrows.
770769
- **Workflow context is inferred** from the dispatcher (`TaskDispatcher<TContext>`), so
771770
`createTaskLayer({ dispatcher: new WorkflowDispatcher(...) })` types `task.run` / `task.sleep` with
772771
no type argument; a test in `workflow.test.ts` keeps it that way.
@@ -873,8 +872,14 @@ implements TanStack AI's own backend contracts (see its section below) — keep
873872
**every IP literal** (v4 after the URL parser's normalization, and any `[…]` v6) plus
874873
localhost/single-label/`.local`/`.internal`; the old private-range/NAT64/6to4 parser is gone.
875874
It does not resolve DNS. Bad params are `-32602` with a `reason`.
876-
**QStash `deduplicationId` cannot contain `:`** (dev server answers 400) — ids are
877-
`${eventId}_${subscriptionId}` sanitized; this was caught only by the e2e smoke.
875+
**Subscription limit:** `maxSubscriptions` (default 8, `Infinity` = off) live subscriptions per
876+
subscriber across all events. The layer checks `store.count` before the challenge (no outbound
877+
POST for a caller at the limit), and `put(sub, { limit })` re-checks atomically (Redis: one Lua
878+
script over `idx:<event>` and `by:<subscriber>`). Refreshes never count. Routing and delivery-time
879+
`authorize` use the payload's input fields *parsed with `input`* (`projection`), so transforms
880+
match on both sides; a payload that doesn't fit `input` makes `emit` throw.
881+
QStash `deduplicationId` cannot contain `:` (dev server answers 400), one more reason not to
882+
derive it from caller-chosen event ids.
878883
- **Local dev needs the QStash dev server** (`npx @upstash/qstash-cli dev`) — it prints deterministic
879884
creds. `APP_URL` must be reachable *from QStash*.
880885

‎examples/mcp-toolkit-demo/app/page.tsx‎

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -161,6 +161,7 @@ export default function Page() {
161161
key={key}
162162
type="button"
163163
className={`driver ${server === key ? "on" : ""}`}
164+
aria-pressed={server === key}
164165
onClick={() => setServer(key)}
165166
>
166167
{SERVERS[key].label}

‎packages/mcp-toolkit/README.md‎

Lines changed: 52 additions & 15 deletions
Original file line numberDiff line numberDiff line change
@@ -121,6 +121,27 @@ progress, then the handler's result) and `task_cancel`. To stop a cancelled task
121121
`await task.isCancelled()` between steps. Cancelling is cooperative, so running code only stops
122122
where it checks.
123123

124+
**Check what the arguments point at.** The model fills in tool arguments, so a `workspaceId` in
125+
them is whatever it was told. Give the task an `authorize`, which runs before the task is stored or
126+
queued, and use `task.principal` (the user who started the task) inside the handler. Never take a
127+
user id from the arguments:
128+
129+
```ts
130+
tasks.define(
131+
"export_workspace",
132+
{
133+
description: "Exports a workspace to CSV.",
134+
inputSchema: z.object({ workspaceId: z.string() }),
135+
// May this user export this workspace? `false` refuses the call.
136+
authorize: ({ workspaceId }, { principal }) => canExport(principal, workspaceId),
137+
},
138+
async ({ workspaceId }, task) => {
139+
const csv = await exportWorkspace(workspaceId, { as: task.principal });
140+
return { content: [{ type: "text", text: csv }] };
141+
},
142+
);
143+
```
144+
124145
<details>
125146
<summary><b>What the model sees</b></summary>
126147

@@ -166,9 +187,13 @@ export const tasks = createTaskLayer({
166187

167188
tasks.define(
168189
"migrate_workspace",
169-
{ description: "Copies a workspace to new storage.", inputSchema: z.object({ workspaceId: z.string() }) },
190+
{
191+
description: "Copies a workspace to new storage.",
192+
inputSchema: z.object({ workspaceId: z.string() }),
193+
authorize: ({ workspaceId }, { principal }) => isOwner(principal, workspaceId),
194+
},
170195
async ({ workspaceId }, task) => {
171-
const batches = await task.run("plan", () => listBatches(workspaceId));
196+
const batches = await task.run("plan", () => listBatches(workspaceId, task.principal));
172197
for (const [i, batch] of batches.entries()) {
173198
if (await task.isCancelled()) return {};
174199
await task.update(`Copying batch ${i + 1}/${batches.length}`);
@@ -278,8 +303,10 @@ the event id stays the same across retries so the host can drop duplicates.
278303
<summary><b>Matching and <code>authorize</code></b></summary>
279304

280305
Every `input` field must also be a `payload` field: the payload is the one place an event's values
281-
come from, and the same values are used to route it and to authorize it. A subscription matches
282-
when each argument it gave equals the payload's value. Emitting
306+
come from, and the same values are used to route it and to authorize it. The payload's values for
307+
the input fields are parsed with the `input` schema, so its transforms apply on both sides and
308+
`authorize` gets the types it declares. A subscription matches when each argument it gave equals
309+
the payload's value. Emitting
283310
`{ documentId: "doc_123", text }` reaches subscribers of `{ documentId: "doc_123" }` and of `{}`.
284311

285312
`authorize(args, caller)` runs twice:
@@ -294,8 +321,8 @@ when each argument it gave equals the payload's value. Emitting
294321

295322
`() => true` lets every authenticated subscriber hear every matching event.
296323

297-
Pass `{ eventId }` as the second argument to `emit` to deduplicate: emitting the same id twice
298-
delivers once.
324+
Pass `{ eventId }` as the second argument to `emit` to give an event a stable id. It is sent as
325+
`webhook-id`, so a host drops a second emit with the same id as a duplicate.
299326

300327
</details>
301328

@@ -318,6 +345,11 @@ delivers once.
318345
random bytes. A refresh with the same secret skips the challenge. If you rotate the key, stored
319346
subscriptions stop receiving events until the host refreshes them.
320347
- **Lifetime.** A subscription lasts 7 days by default and 30 days at most.
348+
- **How many.** A subscriber holds at most 8 live subscriptions across all events (set
349+
`maxSubscriptions`; `Infinity` turns it off). Past that, a new subscribe is refused before the
350+
callback is challenged, with reason `subscription_limit`; refreshing an existing one still works.
351+
Each subscription is a webhook per matching emit, so this bounds what one user can make your
352+
server send.
321353
- **Host answers.** `410` deletes the subscription, `413` and redirects drop the event, and any
322354
other error is retried.
323355
- **Not implemented:** the draft's poll and stream delivery modes (`events/subscribe` refuses
@@ -341,7 +373,9 @@ subscribe yet. The spec draft is
341373
<details>
342374
<summary><b>Who can see what</b></summary>
343375

344-
**Tasks.** The server sets the owner from `principal`, and no tool argument can set it.
376+
**Tasks.** The server sets the owner from `principal`, and no tool argument can set it. Whether a
377+
caller may start a task with given arguments is your `authorize`; the handler gets the owner as
378+
`task.principal`.
345379
`task_status` and `task_cancel` only accept UUID task ids, and compare the stored owner with the
346380
caller. For anyone else, the task looks exactly like an unknown id, so they can't even tell it
347381
exists. Task ids are random UUIDs. No tool lists tasks, and the execute route only accepts signed
@@ -375,8 +409,9 @@ Treat the Redis credentials like any other production secret.
375409
- A Redis client built with `automaticDeserialization: false` is not supported.
376410

377411
**`mcp-events:sub:<id>`**: the subscription (`event`, `args`, `url`, `encryptedSecret`,
378-
`subscriber`, `createdAt`, `expiresAt`), expiring with it. **`mcp-events:idx:<event>`**: a sorted
379-
set of the event's subscription ids, scored by expiry.
412+
`subscriber`, `createdAt`, `expiresAt`), expiring with it. **`mcp-events:idx:<event>`** and
413+
**`mcp-events:by:<subscriber>`**: sorted sets of the event's and the subscriber's subscription ids,
414+
scored by expiry. The second one enforces the limit.
380415

381416
**Not in Redis.** A task message in QStash carries only `{ taskId }`. An event message carries the
382417
full payload, which stays in QStash (and in its DLQ, if every retry fails) until it is delivered.
@@ -389,13 +424,13 @@ The toolkit never stores the caller's token, the request, the plaintext webhook
389424
<summary><b>All options</b></summary>
390425

391426
**`createTaskLayer`**: `store`, `dispatcher` and `principal` are required. Optional:
392-
`defaults.ttlMs` (1 day) and `defaults.pollIntervalMs` (2s).
427+
`defaults.ttlMs` (1 day) and `defaults.pollIntervalMs` (2s), in positive whole milliseconds.
393428

394429
**`tasks.define(name, config, handler)`**: `description` and `inputSchema` are required. Optional:
395-
`title`, `completedMessage`.
430+
`title`, `completedMessage`, `authorize(args, { principal, auth, request })`.
396431

397432
**`createEventLayer`**: `store`, `delivery` and `principal` are required, and so is `secretKey`
398-
unless `MCP_EVENTS_SECRET_KEY` is set. Optional: `allowInsecureCallbacks`.
433+
unless `MCP_EVENTS_SECRET_KEY` is set. Optional: `maxSubscriptions` (8), `allowInsecureCallbacks`.
399434

400435
**`events.define(name, config)`**: `description`, `payload` and `authorize` are required.
401436
Optional: `title`, `input`.
@@ -436,15 +471,17 @@ interface TaskStore {
436471
}
437472

438473
interface TaskDispatcher<TContext = unknown> {
439-
dispatch(task: Task): Promise<void>; // idempotent per task id
474+
dispatch(task: Task): Promise<void>; // called once per task
440475
cancel(taskId: string): Promise<void>;
441476
createExecuteHandler(endpoints: TaskEndpoints<TContext>): (request: Request) => Promise<Response>;
442477
}
443478

444479
interface SubscriptionStore {
445-
put(subscription: Subscription): Promise<void>;
480+
// false, storing nothing, when a new one would put its subscriber over `limit` (atomically)
481+
put(subscription: Subscription, options: { limit: number }): Promise<boolean>;
446482
get(id: string): Promise<Subscription | null>;
447-
delete(subscription: { id: string; event: string }): Promise<void>;
483+
count(subscriber: string): Promise<number>; // live subscriptions, across all events
484+
delete(subscription: { id: string; event: string; subscriber: string }): Promise<void>;
448485
find(event: string): Promise<Subscription[]>; // every live subscription to the event
449486
}
450487

‎packages/mcp-toolkit/package.json‎

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -53,7 +53,7 @@
5353
],
5454
"dependencies": {
5555
"@upstash/qstash": "^2.11.3",
56-
"@upstash/redis": "^1.38.0",
56+
"@upstash/redis": "^1.38.4",
5757
"@upstash/workflow": "^1.3.3",
5858
"zod": "^4.2.0"
5959
},

‎packages/mcp-toolkit/src/events/backends/qstash.test.ts‎

Lines changed: 43 additions & 10 deletions
Original file line numberDiff line numberDiff line change
@@ -35,14 +35,14 @@ describe.skipIf(!hasRedisCreds)("RedisSubscriptionStore (real Redis)", () => {
3535

3636
it("round-trips a subscription, keeping a numeric-looking subscriber a string", async () => {
3737
const sub = makeSub();
38-
await store.put(sub);
38+
await store.put(sub, { limit: Infinity });
3939
expect(await store.get(sub.id)).toEqual(sub);
4040
expect(await store.get("sub_missing")).toBeNull();
4141
});
4242

4343
it("expires the record with the subscription", async () => {
4444
const sub = makeSub({ expiresAt: Date.now() + 30_000 });
45-
await store.put(sub);
45+
await store.put(sub, { limit: Infinity });
4646
const pttl = await redis.pttl(subKey(sub.id));
4747
expect(pttl).toBeGreaterThan(0);
4848
expect(pttl).toBeLessThanOrEqual(30_000);
@@ -53,7 +53,7 @@ describe.skipIf(!hasRedisCreds)("RedisSubscriptionStore (real Redis)", () => {
5353
const one = makeSub({ event: "find", args: { repo: "a" } });
5454
const elsewhere = makeSub({ event: "other", args: {} });
5555
const expired = makeSub({ event: "find" });
56-
for (const sub of [all, one, elsewhere]) await store.put(sub);
56+
for (const sub of [all, one, elsewhere]) await store.put(sub, { limit: Infinity });
5757
// Written live, then left in the index past its expiry.
5858
await redis.zadd(indexKey("find"), { score: Date.now() - 1, member: expired.id });
5959

@@ -68,7 +68,7 @@ describe.skipIf(!hasRedisCreds)("RedisSubscriptionStore (real Redis)", () => {
6868
makeSub({ event: "popular", args: {}, subscriber: `member-${i}` }),
6969
);
7070
for (let i = 0; i < subs.length; i += 100) {
71-
await Promise.all(subs.slice(i, i + 100).map((sub) => store.put(sub)));
71+
await Promise.all(subs.slice(i, i + 100).map((sub) => store.put(sub, { limit: Infinity })));
7272
}
7373
const found = await store.find("popular");
7474
expect(found).toHaveLength(1100);
@@ -77,21 +77,54 @@ describe.skipIf(!hasRedisCreds)("RedisSubscriptionStore (real Redis)", () => {
7777

7878
it("refreshing replaces the record and keeps one index entry", async () => {
7979
const sub = makeSub({ event: "refresh", expiresAt: Date.now() + 10_000 });
80-
await store.put(sub);
81-
await store.put({ ...sub, expiresAt: Date.now() + 50_000 });
80+
await store.put(sub, { limit: Infinity });
81+
await store.put({ ...sub, expiresAt: Date.now() + 50_000 }, { limit: Infinity });
8282
expect((await store.get(sub.id))?.expiresAt).toBeGreaterThan(Date.now() + 40_000);
8383
expect(await store.find("refresh")).toHaveLength(1);
8484
});
8585

8686
it("prunes expired index entries on the next write", async () => {
8787
await redis.zadd(indexKey("prune"), { score: Date.now() - 1000, member: "sub_old" });
88-
await store.put(makeSub({ event: "prune" }));
88+
await store.put(makeSub({ event: "prune" }), { limit: Infinity });
8989
expect(await redis.zscore(indexKey("prune"), "sub_old")).toBeNull();
9090
});
9191

92+
it("refuses a new subscription past the subscriber's limit, atomically, but allows a refresh", async () => {
93+
const subscriber = `limited-${Date.now()}`;
94+
const [a, b, c] = [1, 2, 3].map((i) => makeSub({ event: `limit-${i}`, subscriber }));
95+
expect(await store.put(a!, { limit: 2 })).toBe(true);
96+
expect(await store.put(b!, { limit: 2 })).toBe(true);
97+
expect(await store.count(subscriber)).toBe(2);
98+
expect(await store.put(c!, { limit: 2 })).toBe(false);
99+
expect(await store.get(c!.id)).toBeNull();
100+
expect(await store.find("limit-3")).toEqual([]);
101+
// Replacing one of theirs is not a new subscription.
102+
expect(await store.put({ ...a!, expiresAt: Date.now() + 90_000 }, { limit: 2 })).toBe(true);
103+
// Concurrent puts can't overshoot: the check and the write are one script.
104+
const racers = Array.from({ length: 5 }, (_, i) =>
105+
makeSub({ event: `race-${i}`, subscriber: `${subscriber}-race` }),
106+
);
107+
const results = await Promise.all(racers.map((sub) => store.put(sub, { limit: 2 })));
108+
expect(results.filter(Boolean)).toHaveLength(2);
109+
expect(await store.count(`${subscriber}-race`)).toBe(2);
110+
});
111+
112+
it("frees a slot on delete and when a subscription expires", async () => {
113+
const subscriber = `freed-${Date.now()}`;
114+
const a = makeSub({ event: "freed", subscriber });
115+
await store.put(a, { limit: 1 });
116+
await store.delete(a);
117+
expect(await store.count(subscriber)).toBe(0);
118+
const b = makeSub({ event: "freed", subscriber, expiresAt: Date.now() + 1_500 });
119+
expect(await store.put(b, { limit: 1 })).toBe(true);
120+
await new Promise((resolve) => setTimeout(resolve, 1_700));
121+
expect(await store.count(subscriber)).toBe(0);
122+
expect(await store.put(makeSub({ event: "freed", subscriber }), { limit: 1 })).toBe(true);
123+
});
124+
92125
it("deletes the record and its index entry", async () => {
93126
const sub = makeSub({ event: "delete" });
94-
await store.put(sub);
127+
await store.put(sub, { limit: Infinity });
95128
await store.delete(sub);
96129
expect(await store.get(sub.id)).toBeNull();
97130
expect(await store.find("delete")).toEqual([]);
@@ -171,7 +204,7 @@ describe("QStashDelivery", () => {
171204
expect(sent).toEqual([]);
172205
});
173206

174-
it("publishes one deduplicated message per subscription, in batches of 100", async () => {
207+
it("publishes one message per subscription, in batches of 100", async () => {
175208
const batches: Record<string, unknown>[][] = [];
176209
const qstash = {
177210
batchJSON: async (messages: Record<string, unknown>[]) => {
@@ -187,7 +220,7 @@ describe("QStashDelivery", () => {
187220
url: URL,
188221
body: jobs[0],
189222
retries: 3,
190-
deduplicationId: "evt_1_sub_0",
191223
});
224+
expect(batches[0]?.[0]).not.toHaveProperty("deduplicationId");
192225
});
193226
});

0 commit comments

Comments
 (0)