diff --git a/server/src/__tests__/public-mcp.test.ts b/server/src/__tests__/public-mcp.test.ts index 1577d029e1..9c866b9cb1 100644 --- a/server/src/__tests__/public-mcp.test.ts +++ b/server/src/__tests__/public-mcp.test.ts @@ -406,6 +406,23 @@ describe.skipIf(!support.supported)("public MCP OAuth and tool boundary", () => await f.service.unsubscribe(f.principal, f.input); }); + it("keeps the longest promised lifetime when cached refreshes overlap across replicas", async () => { + const f = await eventFixture(); + const first = await f.service.subscribe(f.principal, { ...f.input, ttlMs: 30_000 }); + const replica = createPublicMcpEvents(db, oauth, f.dispatch, f.options); + const refreshed = await Promise.all([ + f.service.subscribe(f.principal, { ...f.input, ttlMs: 3600_000 }), + replica.subscribe(f.principal, { ...f.input, ttlMs: 30_000 }), + ]); + const [stored] = await db.select().from(mcpEventSubscriptions).where(eq(mcpEventSubscriptions.id, first.id)); + expect(stored!.expiresAt.getTime()).toBe(f.now() + 3600_000); + expect(refreshed.every(r => Date.parse(r.refreshBefore) <= stored!.expiresAt.getTime())).toBe(true); + const shorter = await replica.subscribe(f.principal, { ...f.input, ttlMs: 30_000 }); + expect(Date.parse(shorter.refreshBefore)).toBe(stored!.expiresAt.getTime()); + expect(f.received).toHaveLength(1); + await f.service.unsubscribe(f.principal, f.input); + }); + it("retries a lost delivery with a stable ID and fresh signature, rotates keys and stops on unsubscribe", async () => { const f = await eventFixture(); await f.service.subscribe(f.principal, f.input); diff --git a/server/src/services/public-mcp/events.ts b/server/src/services/public-mcp/events.ts index 21241ad5c4..974fbe142e 100644 --- a/server/src/services/public-mcp/events.ts +++ b/server/src/services/public-mcp/events.ts @@ -131,21 +131,27 @@ export function createPublicMcpEvents(db: Db, oauth: PublicMcpOAuth, api: ApiDis return await db.transaction(async tx => { await tx.execute(sql`select pg_advisory_xact_lock(736721043)`); await tx.execute(sql`select pg_advisory_xact_lock(hashtextextended(${id}, 0))`); + let retainedExpiry = 0; if (lease) { const [held] = await tx.update(admissions).set({ finishedAt: new Date(now()) }).where(and(eq(admissions.id, lease.id), isNull(admissions.finishedAt), gt(admissions.expiresAt, new Date(now())))).returning(); if (!held) throw new McpEventError(-32602, "Subscription verification expired or was stopped. Reconnect the monitor."); if (existing) { const [active] = await tx.select().from(subscriptions).where(and(eq(subscriptions.id, id), isNull(subscriptions.stoppedAt), gt(subscriptions.expiresAt, new Date(now())))); if (!active) throw new McpEventError(-32602, "The monitor expired or was stopped. Subscribe again."); + retainedExpiry = active.expiresAt.getTime(); } } else { // A cached refresh cannot recreate a subscription removed while it was awaiting authority. const [held] = await tx.select().from(subscriptions).where(eq(subscriptions.id, id)); const [pending] = await tx.select().from(admissions).where(and(eq(admissions.subscriptionId, id), isNull(admissions.finishedAt), gt(admissions.expiresAt, new Date(now())))); if (!held || held.stoppedAt || held.expiresAt.getTime() <= now() || pending || canonical(held.deliveryMaterial) !== canonical(existing!.deliveryMaterial)) throw new McpEventError(-32602, "The monitor changed or was stopped. Retry the subscription."); + retainedExpiry = held.expiresAt.getTime(); } if (cloudOrigin && cloud!.expiresAt <= now()) throw new McpEventError(-32602, "Refresh the hosted connection before subscribing."); - const expiresAt = new Date(Math.min(now() + Math.min(Math.max(input.ttlMs ?? lifetime, 30_000), lifetime), cloudOrigin ? Math.min(now() + rotationMs, cloud!.expiresAt) : Infinity)); + // Refreshes extend the same monitor; a shorter overlapping request must + // not revoke a lifetime already promised to another caller. Hosted + // authority still caps the lifetime to its current proof. + const expiresAt = new Date(Math.min(Math.max(retainedExpiry, now() + Math.min(Math.max(input.ttlMs ?? lifetime, 30_000), lifetime)), cloudOrigin ? Math.min(now() + rotationMs, cloud!.expiresAt) : Infinity)); const value = { companyId: principal.grant.companyId, grantId: principal.grant.id, name: input.name, taskId: input.arguments.taskId, arguments: input.arguments, deliveryMaterial: material, expiresAt, stoppedAt: null, verifiedAt: verify ? new Date(now()) : existing!.verifiedAt, startsAt: existing?.startsAt ?? requestedAt, scannedAt: new Date(now()) };