Skip to content

Commit 669fe2f

Browse files
fix(notices): serialize unsubscribe with in-flight inbox observations
An unsubscribe that overlaps an observation parked on the durable store now wins: the observation yields before claiming any budget or sending, and the unsubscribe resolves (and the MCP request is acknowledged) only after in-flight observations settle, so a client is never signalled after its unsubscribe succeeded. Subscribe and unsubscribe share the observation queue.
1 parent 72aa40a commit 669fe2f

3 files changed

Lines changed: 133 additions & 18 deletions

File tree

‎packages/agent-bundle/src/mcp-server-runtime.ts‎

Lines changed: 4 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -524,9 +524,11 @@ const installNoticeInboxSubscriptions = (
524524
}
525525
return {};
526526
}) as never);
527-
protocol.setRequestHandler('resources/unsubscribe', ((request: ResourceSubscriptionRequest) => {
527+
protocol.setRequestHandler('resources/unsubscribe', (async (request: ResourceSubscriptionRequest) => {
528528
assertInboxUri(request.params.uri);
529-
notices.unsubscribe();
529+
// Acknowledged only once in-flight observations have settled, so the
530+
// client never receives a signal after its unsubscribe succeeded.
531+
await notices.unsubscribe();
530532
return {};
531533
}) as never);
532534
const send = async (): Promise<void> => {

‎packages/rsc-runtime/src/notices/resource-updated.ts‎

Lines changed: 44 additions & 15 deletions
Original file line numberDiff line numberDiff line change
@@ -73,7 +73,13 @@ export interface AgentNoticeInboxSignaller {
7373
* subscription that could never be honoured.
7474
*/
7575
subscribe(principal: AgentNoticePrincipal): Promise<void>;
76-
unsubscribe(): void;
76+
/**
77+
* Ends the subscription. Resolves only after every observation already in
78+
* flight has settled, so once it resolves no further signal is sent or
79+
* budget spent for the connection — a `resources/unsubscribe` acknowledged
80+
* to the client is honoured even while the store is slow.
81+
*/
82+
unsubscribe(): Promise<void>;
7783
}
7884

7985
interface InboxSubscription {
@@ -115,9 +121,21 @@ export const createNoticeInboxSignaller = (
115121
const now = options.now ?? ((): Date => new Date());
116122
let subscription: InboxSubscription | undefined;
117123
let signalSequence = 0;
118-
// Observations serialize so two renders completing together cannot both
119-
// select the same notice and send two signals for one revision.
124+
// Observations and subscription changes serialize on one queue: two renders
125+
// completing together cannot both select the same notice and send two
126+
// signals for one revision, and an unsubscribe (or re-subscribe) that
127+
// overlaps an observation awaiting the store takes effect only after that
128+
// observation settles, never between its eligibility read and its send.
120129
let queue: Promise<unknown> = Promise.resolve();
130+
// Unsubscribes requested while an observation awaits the store are counted
131+
// synchronously so that observation yields before spending any budget: the
132+
// client asked to stop, so nothing is claimed or sent on its behalf.
133+
let pendingUnsubscribes = 0;
134+
const serialized = <T>(step: () => Promise<T>): Promise<T> => {
135+
const run = queue.then(step);
136+
queue = run.catch(() => undefined);
137+
return run;
138+
};
121139

122140
/**
123141
* Claims the budget for every newly eligible notice as one compare-and-swap
@@ -133,6 +151,7 @@ export const createNoticeInboxSignaller = (
133151
): Promise<
134152
| { readonly at: string; readonly kind: 'claimed'; readonly noticeIds: readonly string[]; readonly revision: number }
135153
| { readonly kind: 'nothing-eligible'; readonly revision: number }
154+
| { readonly kind: 'unsubscribed' }
136155
| { readonly error: unknown; readonly kind: 'failed'; readonly stage: 'read' | 'record' }
137156
> => {
138157
for (let attempt = 0; attempt < MAX_CLAIM_ATTEMPTS; attempt += 1) {
@@ -142,6 +161,7 @@ export const createNoticeInboxSignaller = (
142161
} catch (error) {
143162
return { error, kind: 'failed', stage: 'read' };
144163
}
164+
if (pendingUnsubscribes > 0 || subscription !== current) return { kind: 'unsubscribed' };
145165
const at = now().toISOString();
146166
const nowMs = Date.parse(at);
147167
const noticeIds = Object.freeze(snapshot.notices
@@ -183,12 +203,17 @@ export const createNoticeInboxSignaller = (
183203
} catch (error) {
184204
return Object.freeze({ error, kind: 'failed' as const, stage: 'read' as const });
185205
}
206+
if (pendingUnsubscribes > 0) {
207+
return Object.freeze({ kind: 'idle', reason: 'no-subscription', revision: undefined });
208+
}
186209
const claimed = await claim(ledger, current);
187210
switch (claimed.kind) {
188211
case 'failed':
189212
return Object.freeze({ error: claimed.error, kind: 'failed' as const, stage: claimed.stage });
190213
case 'nothing-eligible':
191214
return Object.freeze({ kind: 'idle', reason: 'nothing-eligible', revision: claimed.revision });
215+
case 'unsubscribed':
216+
return Object.freeze({ kind: 'idle', reason: 'no-subscription', revision: undefined });
192217
case 'claimed':
193218
break;
194219
default: {
@@ -219,21 +244,25 @@ export const createNoticeInboxSignaller = (
219244
return options.store.close();
220245
},
221246
observe(send: () => Promise<void>): Promise<AgentNoticeInboxSignalOutcome> {
222-
const run = queue.then(() => observeOnce(send));
223-
queue = run.catch(() => undefined);
224-
return run;
247+
return serialized(() => observeOnce(send));
225248
},
226-
async subscribe(principal: AgentNoticePrincipal): Promise<void> {
227-
const ledger = await options.store.noticeLedger();
228-
await ledger.read();
229-
subscription = Object.freeze({
230-
id: randomUUID(),
231-
principal,
232-
signalled: new Set<string>(),
249+
subscribe(principal: AgentNoticePrincipal): Promise<void> {
250+
return serialized(async () => {
251+
const ledger = await options.store.noticeLedger();
252+
await ledger.read();
253+
subscription = Object.freeze({
254+
id: randomUUID(),
255+
principal,
256+
signalled: new Set<string>(),
257+
});
233258
});
234259
},
235-
unsubscribe(): void {
236-
subscription = undefined;
260+
unsubscribe(): Promise<void> {
261+
pendingUnsubscribes += 1;
262+
return serialized(async () => {
263+
pendingUnsubscribes -= 1;
264+
subscription = undefined;
265+
});
237266
},
238267
});
239268
};

‎packages/rsc-runtime/tests/notices-resource-updated.test.ts‎

Lines changed: 85 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -155,12 +155,96 @@ describe('notice inbox resources/updated signaller', () => {
155155
await driver.close();
156156
});
157157

158+
it('honours an unsubscribe that overlaps an observation awaiting the store', async () => {
159+
const { driver, ledger } = await openLedger();
160+
await publish(ledger, { sessionId: 's1' });
161+
// A slow ledger: the first read parks until the test releases it, standing
162+
// in for a contended durable store.
163+
let release: (() => void) | undefined;
164+
const gate = new Promise<void>((resolve) => {
165+
release = resolve;
166+
});
167+
let parkNextRead = false;
168+
const slowLedger: AgentNoticeLedger = Object.freeze({
169+
...ledger,
170+
read: async () => {
171+
if (parkNextRead) {
172+
parkNextRead = false;
173+
await gate;
174+
}
175+
return ledger.read();
176+
},
177+
});
178+
const signaller = signallerOver(slowLedger);
179+
const { send, sends } = sender();
180+
await signaller.subscribe(principal('s1'));
181+
182+
const events: string[] = [];
183+
parkNextRead = true;
184+
const observing = signaller.observe(async () => {
185+
events.push('send');
186+
await send();
187+
});
188+
// `resources/unsubscribe` arrives while the observation is parked on the read.
189+
const unsubscribing = signaller.unsubscribe().then(() => {
190+
events.push('unsubscribed');
191+
});
192+
release!();
193+
await expect(observing).resolves.toEqual({ kind: 'idle', reason: 'no-subscription', revision: undefined });
194+
await unsubscribing;
195+
expect(events).toEqual(['unsubscribed']);
196+
expect(sends).toEqual([]);
197+
expect(signaller.subscribed).toBe(false);
198+
// Nothing was claimed on the client's behalf: the durable budget is intact,
199+
// so a later subscriber still receives the signal.
200+
expect((await ledger.read()).notices[0]?.availability).toBeUndefined();
201+
await signaller.subscribe(principal('s1'));
202+
await expect(signaller.observe(send)).resolves.toMatchObject({ kind: 'signalled' });
203+
expect(sends).toHaveLength(1);
204+
await driver.close();
205+
});
206+
207+
it('acknowledges unsubscribe only after an in-flight send has settled', async () => {
208+
const { driver, ledger } = await openLedger();
209+
await publish(ledger, { sessionId: 's1' });
210+
const signaller = signallerOver(ledger);
211+
await signaller.subscribe(principal('s1'));
212+
let releaseSend: (() => void) | undefined;
213+
const sendGate = new Promise<void>((resolve) => {
214+
releaseSend = resolve;
215+
});
216+
const events: string[] = [];
217+
const observing = signaller.observe(async () => {
218+
await sendGate;
219+
events.push('send');
220+
});
221+
// Let the observation pass its eligibility read and claim before the
222+
// unsubscribe arrives mid-send.
223+
await new Promise<void>((resolve) => {
224+
setTimeout(resolve, 0);
225+
});
226+
const unsubscribing = signaller.unsubscribe().then(() => {
227+
events.push('unsubscribed');
228+
});
229+
releaseSend!();
230+
await expect(observing).resolves.toMatchObject({ kind: 'signalled' });
231+
await unsubscribing;
232+
// The wire write that was already committed completes first; the client is
233+
// told it is unsubscribed only afterwards, never the other way round.
234+
expect(events).toEqual(['send', 'unsubscribed']);
235+
await expect(signaller.observe(async () => {
236+
events.push('late');
237+
})).resolves.toMatchObject({ kind: 'idle', reason: 'no-subscription' });
238+
expect(events).toEqual(['send', 'unsubscribed']);
239+
await driver.close();
240+
});
241+
158242
it('stops signalling once unsubscribed and resets tracking on re-subscribe', async () => {
159243
const { driver, ledger } = await openLedger();
160244
const signaller = signallerOver(ledger);
161245
const { send, sends } = sender();
162246
await signaller.subscribe(principal('s1'));
163-
signaller.unsubscribe();
247+
await signaller.unsubscribe();
164248
expect(signaller.subscribed).toBe(false);
165249

166250
await publish(ledger, { retryBudget: 2, sessionId: 's1' });

0 commit comments

Comments
 (0)