Skip to content

Commit 964c947

Browse files
fix(notices): keep wire-successful receipts owed and renew the hold during pending sends
Codex P2 x2 on #376: - A send that succeeded but whose signalAvailability() commit failed was deduplicated only in memory, so a restarted signaller could resend it once the abandoned hold lapsed. Owed receipts are now retried with the same idempotency key before any later observation spends, and on close(). - A protocol write pending longer than the 30s TTL let another process take the hold and send too. The holder now renews under its key while send() is pending; the reducer refuses a foreign key while a hold is live and refuses a lapsed holder's late renewal once another key has taken over.
1 parent 4665444 commit 964c947

6 files changed

Lines changed: 330 additions & 26 deletions

File tree

‎.changeset/notice-inbox-resource-updated.md‎

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -3,4 +3,4 @@
33
"agent-bundle": patch
44
---
55

6-
Wire the #99 stage-4 `mcp-resource-updated` delivery route into generated stateful MCP servers. `@agent-bundle/runtime/notices` gains `createNoticeInboxSignaller` — one connection's subscription to the reserved inbox resource (`AGENT_NOTICE_INBOX_URI`) that, after each completed render, sends at most one `notifications/resources/updated` for the subscriber's newly eligible pending notices and records it through `signalAvailability()` as an availability receipt (never delivery), honouring `nextAttemptAt` and bounding signals per notice by `retryBudget` across restarts. The ledger gains `reserveAvailability()` / `releaseAvailability()` (and `AgentNotice.availabilityReservation`, `AGENT_NOTICE_AVAILABILITY_RESERVATION_TTL_MS`): the signaller holds a notice's budget slot by compare-and-swap before the wire write so concurrent processes cannot both send, finalizes the hold into the receipt only when the protocol write succeeds, and releases it when the write fails, so a failed send costs no budget and the receipt always means the write succeeded; `unsubscribe()` resolves only after in-flight observations settle; `@agent-bundle/runtime/mount` gains `createGeneratedNoticeRuntime` and `GeneratedRuntimeState.noticeLedger()` so a server process can hold its own handle on the durable store its worker mounts. Generated workspace-durable MCP entries now register `resources/subscribe`/`resources/unsubscribe` for the inbox URI only, advertise `resources.subscribe` exactly when that wiring is active, fail subscriptions closed when the store is unreadable, and the inbox projection exposes the `availability` receipt alongside `exposure`. Volatile lifetimes keep the store in the worker's heap and advertise no subscription capability.
6+
Wire the #99 stage-4 `mcp-resource-updated` delivery route into generated stateful MCP servers. `@agent-bundle/runtime/notices` gains `createNoticeInboxSignaller` — one connection's subscription to the reserved inbox resource (`AGENT_NOTICE_INBOX_URI`) that, after each completed render, sends at most one `notifications/resources/updated` for the subscriber's newly eligible pending notices and records it through `signalAvailability()` as an availability receipt (never delivery), honouring `nextAttemptAt` and bounding signals per notice by `retryBudget` across restarts. The ledger gains `reserveAvailability()` / `releaseAvailability()` (and `AgentNotice.availabilityReservation`, `AGENT_NOTICE_AVAILABILITY_RESERVATION_TTL_MS`): the signaller holds a notice's budget slot by compare-and-swap before the wire write so concurrent processes cannot both send, finalizes the hold into the receipt only when the protocol write succeeds, and releases it when the write fails, so a failed send costs no budget and the receipt always means the write succeeded. The hold is renewed while the write is pending (a different key may take it over only after the TTL, and a lapsed holder cannot steal it back), and a receipt whose commit failed after a successful send is retried idempotently before later observations and on `close()`; `unsubscribe()` resolves only after in-flight observations settle; `@agent-bundle/runtime/mount` gains `createGeneratedNoticeRuntime` and `GeneratedRuntimeState.noticeLedger()` so a server process can hold its own handle on the durable store its worker mounts. Generated workspace-durable MCP entries now register `resources/subscribe`/`resources/unsubscribe` for the inbox URI only, advertise `resources.subscribe` exactly when that wiring is active, fail subscriptions closed when the store is unreadable, and the inbox projection exposes the `availability` receipt alongside `exposure`. Volatile lifetimes keep the store in the worker's heap and advertise no subscription capability.

‎packages/rsc-runtime/README.md‎

Lines changed: 7 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -266,10 +266,15 @@ protocol write failed, so the receipt only ever means the write succeeded and
266266
a failed send costs no budget. Eligibility is recipient-matched against the
267267
subscriber's observed identity, respects `nextAttemptAt`, skips notices whose
268268
slot another signaller currently holds (a hold older than
269-
`AGENT_NOTICE_AVAILABILITY_RESERVATION_TTL_MS` counts as abandoned), and is
269+
`AGENT_NOTICE_AVAILABILITY_RESERVATION_TTL_MS` counts as abandoned; a live
270+
holder renews it under its key while its write is pending, and the reducer
271+
refuses a renewal once another key has legitimately taken over), and is
270272
bounded by `retryBudget` (availability receipts per notice, durable across
271273
restarts); because the slot is held before the wire write, two server processes
272-
over one store can never both signal the same notice. Exposure and availability
274+
over one store can never both signal the same notice. A send that reached the
275+
wire but whose receipt commit failed stays owed: the same idempotent receipt is
276+
retried before any later observation spends (and on `close()`), so a restarted
277+
process cannot resend a notice the wire already carried. Exposure and availability
273278
receipts never re-trigger a signal, so a subscribed client cannot be driven into
274279
a refetch loop. Subscribing fails closed when the store is unreadable,
275280
`unsubscribe()` resolves only after in-flight observations settle, and only the

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

Lines changed: 123 additions & 14 deletions
Original file line numberDiff line numberDiff line change
@@ -28,6 +28,12 @@ export interface AgentNoticeInboxStore {
2828
export interface CreateNoticeInboxSignallerOptions {
2929
/** Clock injection for deterministic tests. */
3030
readonly now?: () => Date;
31+
/**
32+
* How often a hold is renewed while its `send()` is still pending. Defaults
33+
* to a third of `AGENT_NOTICE_AVAILABILITY_RESERVATION_TTL_MS` so a live
34+
* holder never lapses; only a holder whose process is gone does.
35+
*/
36+
readonly reservationRenewalIntervalMs?: number;
3137
readonly store: AgentNoticeInboxStore;
3238
}
3339

@@ -59,11 +65,14 @@ export type AgentNoticeInboxSignalOutcome =
5965
export interface AgentNoticeInboxSignaller {
6066
readonly inboxUri: typeof AGENT_NOTICE_INBOX_URI;
6167
readonly subscribed: boolean;
68+
/** Commits any receipt still owed for a send that reached the wire, then closes the store. */
6269
close(): Promise<void>;
6370
/**
64-
* Runs after one completed render: reads the ledger and, when the
65-
* subscriber has newly eligible pending notices, sends exactly one
66-
* `resources/updated` through `send` and records the availability receipt.
71+
* Runs after one completed render: first commits any receipt still owed from
72+
* an earlier send, then reads the ledger and, when the subscriber has newly
73+
* eligible pending notices, holds their budget slot, sends exactly one
74+
* `resources/updated` through `send` (renewing the hold while the write is
75+
* pending), and records the availability receipt once the write succeeded.
6776
* Never throws; failures are returned so the render path stays unaffected.
6877
*/
6978
observe(send: () => Promise<void>): Promise<AgentNoticeInboxSignalOutcome>;
@@ -89,6 +98,13 @@ interface InboxSubscription {
8998
readonly signalled: Set<string>;
9099
}
91100

101+
/** A send that succeeded on the wire whose availability receipt has not been committed yet. */
102+
interface PendingReceipt {
103+
readonly at: string;
104+
readonly noticeIds: readonly string[];
105+
readonly reservationKey: string;
106+
}
107+
92108
const eligibleForSignal = (
93109
notice: AgentNotice,
94110
principal: AgentNoticePrincipal,
@@ -127,8 +143,16 @@ export const createNoticeInboxSignaller = (
127143
options: CreateNoticeInboxSignallerOptions,
128144
): AgentNoticeInboxSignaller => {
129145
const now = options.now ?? ((): Date => new Date());
146+
const renewalIntervalMs = options.reservationRenewalIntervalMs
147+
?? Math.floor(AGENT_NOTICE_AVAILABILITY_RESERVATION_TTL_MS / 3);
130148
let subscription: InboxSubscription | undefined;
131149
let signalSequence = 0;
150+
// Sends that reached the wire but whose receipt commit failed. The send is
151+
// a fact, so the receipt is retried with the same idempotency key on every
152+
// later observation (and on close) until it lands; the memory dedupe alone
153+
// would otherwise let a restarted signaller spend the same slot again once
154+
// the abandoned hold lapsed.
155+
const pendingReceipts = new Map<string, PendingReceipt>();
132156
// Observations and subscription changes serialize on one queue: two renders
133157
// completing together cannot both select the same notice and send two
134158
// signals for one revision, and an unsubscribe (or re-subscribe) that
@@ -210,9 +234,74 @@ export const createNoticeInboxSignaller = (
210234
};
211235
};
212236

237+
/** Commits the receipt for one wire-successful send; idempotent across retries. */
238+
const commitReceipt = (
239+
ledger: AgentNoticeLedger,
240+
receipt: PendingReceipt,
241+
): Promise<Awaited<ReturnType<AgentNoticeLedger['read']>>> => ledger.signalAvailability({
242+
at: receipt.at,
243+
idempotencyKey: `agent-notices:availability:signal:${receipt.reservationKey}`,
244+
noticeIds: receipt.noticeIds,
245+
reservationKey: receipt.reservationKey,
246+
});
247+
248+
/** Retries every outstanding receipt; the first failure is returned so the caller reports it. */
249+
const drainPendingReceipts = async (ledger: AgentNoticeLedger): Promise<unknown> => {
250+
for (const receipt of [...pendingReceipts.values()]) {
251+
try {
252+
await commitReceipt(ledger, receipt);
253+
pendingReceipts.delete(receipt.reservationKey);
254+
} catch (error) {
255+
return error;
256+
}
257+
}
258+
return undefined;
259+
};
260+
261+
/**
262+
* Keeps a hold alive while its protocol write is pending. A write that
263+
* outlives the TTL would otherwise let another process treat the hold as
264+
* abandoned and send too; renewing under the same key refreshes `at`, and
265+
* the reducer refuses a renewal once a different key has legitimately taken
266+
* over, so a holder that could not renew for a whole TTL never steals back.
267+
* The timer exists only for the duration of one in-flight send.
268+
*/
269+
const renewWhile = async <T>(
270+
ledger: AgentNoticeLedger,
271+
hold: { readonly noticeIds: readonly string[]; readonly reservationKey: string },
272+
pending: Promise<T>,
273+
): Promise<T> => {
274+
let renewals = 0;
275+
let stopped = false;
276+
let timer: ReturnType<typeof setTimeout> | undefined;
277+
const tick = (): void => {
278+
if (stopped) return;
279+
renewals += 1;
280+
const renewal = renewals;
281+
void ledger.reserveAvailability({
282+
at: now().toISOString(),
283+
idempotencyKey: `agent-notices:availability:renew:${hold.reservationKey}:${String(renewal)}`,
284+
noticeIds: hold.noticeIds,
285+
reservationKey: hold.reservationKey,
286+
}).catch(() => undefined).then(() => {
287+
if (stopped) return;
288+
timer = setTimeout(tick, renewalIntervalMs);
289+
timer.unref?.();
290+
});
291+
};
292+
timer = setTimeout(tick, renewalIntervalMs);
293+
timer.unref?.();
294+
try {
295+
return await pending;
296+
} finally {
297+
stopped = true;
298+
if (timer !== undefined) clearTimeout(timer);
299+
}
300+
};
301+
213302
const observeOnce = async (send: () => Promise<void>): Promise<AgentNoticeInboxSignalOutcome> => {
214303
const current = subscription;
215-
if (current === undefined) {
304+
if (current === undefined && pendingReceipts.size === 0) {
216305
return Object.freeze({ kind: 'idle', reason: 'no-subscription', revision: undefined });
217306
}
218307
let ledger: AgentNoticeLedger;
@@ -221,7 +310,13 @@ export const createNoticeInboxSignaller = (
221310
} catch (error) {
222311
return Object.freeze({ error, kind: 'failed' as const, stage: 'read' as const });
223312
}
224-
if (pendingUnsubscribes > 0) {
313+
// Receipts owed from earlier sends come first: they are facts about the
314+
// wire, and a ledger that cannot take them is not one to spend against.
315+
const owed = await drainPendingReceipts(ledger);
316+
if (owed !== undefined) {
317+
return Object.freeze({ error: owed, kind: 'failed' as const, stage: 'record' as const });
318+
}
319+
if (current === undefined || pendingUnsubscribes > 0) {
225320
return Object.freeze({ kind: 'idle', reason: 'no-subscription', revision: undefined });
226321
}
227322
const claimed = await claim(ledger, current);
@@ -245,7 +340,7 @@ export const createNoticeInboxSignaller = (
245340
// send finalizes the hold into the availability receipt. Only the receipt
246341
// means the protocol write succeeded.
247342
try {
248-
await send();
343+
await renewWhile(ledger, claimed, send());
249344
} catch (error) {
250345
try {
251346
await ledger.releaseAvailability({
@@ -259,14 +354,18 @@ export const createNoticeInboxSignaller = (
259354
}
260355
return Object.freeze({ error, kind: 'failed' as const, stage: 'send' as const });
261356
}
357+
// The wire write succeeded: this subscription never sends for these notices
358+
// again, and the receipt is owed until it commits.
262359
for (const id of claimed.noticeIds) current.signalled.add(id);
360+
const receipt: PendingReceipt = Object.freeze({
361+
at: claimed.at,
362+
noticeIds: claimed.noticeIds,
363+
reservationKey: claimed.reservationKey,
364+
});
365+
pendingReceipts.set(receipt.reservationKey, receipt);
263366
try {
264-
const committed = await ledger.signalAvailability({
265-
at: claimed.at,
266-
idempotencyKey: `agent-notices:availability:signal:${claimed.reservationKey}`,
267-
noticeIds: claimed.noticeIds,
268-
reservationKey: claimed.reservationKey,
269-
});
367+
const committed = await commitReceipt(ledger, receipt);
368+
pendingReceipts.delete(receipt.reservationKey);
270369
return Object.freeze({ kind: 'signalled', noticeIds: claimed.noticeIds, revision: committed.revision });
271370
} catch (error) {
272371
return Object.freeze({ error, kind: 'failed' as const, stage: 'record' as const });
@@ -279,8 +378,18 @@ export const createNoticeInboxSignaller = (
279378
return subscription !== undefined;
280379
},
281380
close(): Promise<void> {
282-
subscription = undefined;
283-
return options.store.close();
381+
return serialized(async () => {
382+
subscription = undefined;
383+
if (pendingReceipts.size > 0) {
384+
try {
385+
await drainPendingReceipts(await options.store.noticeLedger());
386+
} catch {
387+
// A receipt still owed at close is lost with the process; the hold
388+
// it left behind lapses after the TTL.
389+
}
390+
}
391+
await options.store.close();
392+
});
284393
},
285394
observe(send: () => Promise<void>): Promise<AgentNoticeInboxSignalOutcome> {
286395
return serialized(() => observeOnce(send));

‎packages/rsc-runtime/src/notices/state.ts‎

Lines changed: 21 additions & 9 deletions
Original file line numberDiff line numberDiff line change
@@ -11,10 +11,11 @@ import type {
1111
AgentStateLifetime,
1212
} from '../state/contract.js';
1313
import { canonicalJson, defineState } from '../state/index.js';
14-
import type {
15-
AgentNotice,
16-
AgentNoticePrincipal,
17-
AgentRecipient,
14+
import {
15+
AGENT_NOTICE_AVAILABILITY_RESERVATION_TTL_MS,
16+
type AgentNotice,
17+
type AgentNoticePrincipal,
18+
type AgentRecipient,
1819
} from './contract.js';
1920

2021
const observed = <T extends z.ZodType>(value: T) => z.discriminatedUnion('state', [
@@ -299,10 +300,12 @@ const transitionAvailability = (
299300
};
300301

301302
/**
302-
* Holds a budget slot without spending it. A live notice may carry one
303-
* reservation; a newer reserver replaces an older one outright because the
304-
* ledger cannot tell abandoned from slow, so the signaller decides staleness
305-
* by `at` against the reservation TTL before it ever reserves.
303+
* Holds a budget slot without spending it. A live notice carries at most one
304+
* reservation: the holder renews it by reserving again under the same key,
305+
* and a different key takes it over only once the current hold is older than
306+
* the reservation TTL at the event's own `at` — a rule the reducer can apply
307+
* deterministically on replay, so a slow holder's renewal can never steal a
308+
* hold back from the process that legitimately took over after it lapsed.
306309
*/
307310
const transitionAvailabilityReservation = (
308311
notice: AgentNotice,
@@ -311,11 +314,20 @@ const transitionAvailabilityReservation = (
311314
if (!input.noticeIds.has(notice.id)) return notice;
312315
switch (notice.state) {
313316
case 'pending':
314-
case 'attempted':
317+
case 'attempted': {
318+
const held = notice.availabilityReservation;
319+
if (
320+
held !== undefined
321+
&& held.key !== input.reservationKey
322+
&& Date.parse(held.at) + AGENT_NOTICE_AVAILABILITY_RESERVATION_TTL_MS > Date.parse(input.at)
323+
) {
324+
return notice;
325+
}
315326
return Object.freeze({
316327
...notice,
317328
availabilityReservation: Object.freeze({ at: input.at, key: input.reservationKey }),
318329
});
330+
}
319331
case 'expired':
320332
case 'unavailable':
321333
case 'withdrawn':

‎packages/rsc-runtime/tests/notices-ledger.test.ts‎

Lines changed: 41 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1,6 +1,7 @@
11
import { describe, expect, it } from '@rstest/core';
22

33
import {
4+
AGENT_NOTICE_AVAILABILITY_RESERVATION_TTL_MS,
45
AGENT_NOTICE_STATES,
56
AgentNoticeError,
67
selectNoticeDeliveryRoutes,
@@ -882,6 +883,46 @@ describe('notice delivery routing receipts (#99 stage 4)', () => {
882883
reservationKey: 'holder-b:1',
883884
})).rejects.toMatchObject({ code: 'revision-conflict' });
884885

886+
// A live hold is not overwritten by a different key even without a
887+
// compare-and-swap; the holder's own key renews it; a lapsed hold is taken.
888+
const contested = await ledger.reserveAvailability({
889+
at: '2026-09-01T19:04:02.000Z',
890+
idempotencyKey: 'reserve:b-uncontended',
891+
noticeIds: [id],
892+
reservationKey: 'holder-b:1',
893+
});
894+
expect(contested.notices[0]?.availabilityReservation).toEqual({ at: '2026-09-01T19:04:00.000Z', key: 'holder-a:1' });
895+
const renewed = await ledger.reserveAvailability({
896+
at: '2026-09-01T19:04:10.000Z',
897+
idempotencyKey: 'renew:a:1',
898+
noticeIds: [id],
899+
reservationKey: 'holder-a:1',
900+
});
901+
expect(renewed.notices[0]?.availabilityReservation).toEqual({ at: '2026-09-01T19:04:10.000Z', key: 'holder-a:1' });
902+
const lapsedAt = new Date(Date.parse('2026-09-01T19:04:10.000Z') + AGENT_NOTICE_AVAILABILITY_RESERVATION_TTL_MS).toISOString();
903+
const takenOver = await ledger.reserveAvailability({
904+
at: lapsedAt,
905+
idempotencyKey: 'reserve:b-after-lapse',
906+
noticeIds: [id],
907+
reservationKey: 'holder-b:1',
908+
});
909+
expect(takenOver.notices[0]?.availabilityReservation).toEqual({ at: lapsedAt, key: 'holder-b:1' });
910+
// The lapsed holder's late renewal cannot steal the hold back.
911+
const lateRenewal = await ledger.reserveAvailability({
912+
at: new Date(Date.parse(lapsedAt) + 1_000).toISOString(),
913+
idempotencyKey: 'renew:a:2',
914+
noticeIds: [id],
915+
reservationKey: 'holder-a:1',
916+
});
917+
expect(lateRenewal.notices[0]?.availabilityReservation).toEqual({ at: lapsedAt, key: 'holder-b:1' });
918+
await ledger.releaseAvailability({ idempotencyKey: 'release:b-takeover', noticeIds: [id], reservationKey: 'holder-b:1' });
919+
await ledger.reserveAvailability({
920+
at: '2026-09-01T19:04:00.000Z',
921+
idempotencyKey: 'reserve:a-again',
922+
noticeIds: [id],
923+
reservationKey: 'holder-a:1',
924+
});
925+
885926
// Releasing with another holder's key leaves the hold intact; the owner's key clears it.
886927
const foreignRelease = await ledger.releaseAvailability({
887928
idempotencyKey: 'release:b',

0 commit comments

Comments
 (0)