Preserve promised monitor lifetime across concurrent refreshes

Co-Authored-By: Paperclip <noreply@paperclip.ing>
This commit is contained in:
DottaandPaperclip committed 2026-10-03 14:43:30 -05:00
1 parent 84b137420a
commit dc9dda92ac
2 files changed
+24 -1

No files matched your search

+17
View File
@@ -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);
+7 -1
View File
@@ -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()) };