From bd1284b661f91b227a17d3bff8c0f1be8e6eee19 Mon Sep 17 00:00:00 2001 From: Lilleman auf Larv Date: Mon, 28 Sep 2026 00:59:38 +0200 Subject: [PATCH 1/3] Give the held-message protocol one home in MessageHold --- AGENTS.md | 2 +- src/held-messages.ts | 43 ++++++++++++++++++++++++++++++++++++++-- src/incoming-requests.ts | 27 +++++++++---------------- todo.md | 9 --------- 4 files changed, 51 insertions(+), 30 deletions(-) diff --git a/AGENTS.md b/AGENTS.md index c3858e6..92a902d 100644 --- a/AGENTS.md +++ b/AGENTS.md @@ -52,7 +52,7 @@ src/ dlr-merger.ts DlrMerger: per-segment receipts counted into one MessageDlr error-from.ts An untyped value as error material: errorFrom() an Error, namedValue() a name expiring-groups.ts ExpiringGroups: the capped, expiring store both of those share - held-messages.ts HeldMessages: capped, expiring messages the application has not answered + held-messages.ts HeldMessages: capped, expiring messages the application has not answered, one MessageHold each idle-waiters.ts IdleWaiters: waiting for a count to fall to zero, and what is left of a budget incoming-requests.ts Every request the peer sends: messages, receipts, links, unknown commands link-gate.ts LinkGate: where a request with no link to go out on waits for the next one diff --git a/src/held-messages.ts b/src/held-messages.ts index 9585382..20f4cbd 100644 --- a/src/held-messages.ts +++ b/src/held-messages.ts @@ -25,6 +25,41 @@ type Held = { pduObjs: PduObject[]; }; +/** One message offered to the application, held until it is answered or every listener gives up. */ +export class MessageHold { + private readonly held: HeldMessages; + private readonly pduObjs: PduObject[]; + private working: number; + + constructor(held: HeldMessages, pduObjs: PduObject[], listeners: number) { + this.held = held; + this.pduObjs = pduObjs; + this.working = listeners; + } + + /** Whether a drain is still waiting for this message to be answered. */ + isHeld(): boolean { + return this.held.has(this.pduObjs); + } + + /** A turn later, so a listener sending its receipt straight after the response still holds. */ + answered(): void { + setImmediate(() => { this.held.release(this.pduObjs); }); + } + + /** A rejection leaves the other listeners running, so only the last one to fail gives the message up. */ + listenerGaveUp(): void { + this.working--; + + if (this.working <= 0) this.answered(); + } + + /** A message nobody took is not work a shutdown can wait for. */ + untaken(): void { + this.held.release(this.pduObjs); + } +} + /** The messages handed to the application that it has not answered yet, held by their segments. */ export class HeldMessages { private readonly held: ExpiringGroups; @@ -59,10 +94,12 @@ export class HeldMessages { return this.held.full || this.octets >= this.maxOctets; } - hold(pduObjs: PduObject[]): void { + /** Holds a message about to be offered to `listeners` listeners. */ + hold(pduObjs: PduObject[], listeners = 1): MessageHold { + const hold = new MessageHold(this, pduObjs, listeners); const key = keyOf(pduObjs); - if (key === undefined) return; + if (key === undefined) return hold; this.sweep(); @@ -77,6 +114,8 @@ export class HeldMessages { this.held.set(key, { octets, pduObjs }); this.octets += octets; + + return hold; } /** Whether a drain is still waiting for this message to be answered. */ diff --git a/src/incoming-requests.ts b/src/incoming-requests.ts index e018b55..8985e5c 100644 --- a/src/incoming-requests.ts +++ b/src/incoming-requests.ts @@ -10,6 +10,7 @@ import type { Result, VoidResult } from './result.ts'; import type { SmppLog } from './log.ts'; import type { Sms, SmsHandlers, SmsInput } from './sms.ts'; import type { SmsIdFormat } from './sms-id.ts'; +import type { MessageHold } from './held-messages.ts'; import { HeldMessages } from './held-messages.ts'; import { Reassembler, decodeSegments } from './reassembly.ts'; import { bindCommands, defaults, standsInFor } from './session-options.ts'; @@ -83,9 +84,9 @@ export type IncomingRequestsOptions = { export class IncomingRequests { private readonly deps: IncomingDeps; private readonly dlrMerger: DlrMerger; - /** The rejection handler is handed the Sms back as an `unknown`, so its hold is found by identity. */ - private readonly emitted = new WeakMap void>(); private readonly held: HeldMessages; + /** The rejection handler is handed the Sms back as an `unknown`, so its hold is found by identity. */ + private readonly holds = new WeakMap(); private readonly log: SmppLog; private readonly reassembler: Reassembler; private readonly smsIdFormat: SmsIdFormat; @@ -175,7 +176,7 @@ export class IncomingRequests { listenerRejected(sms: unknown): void { if (typeof sms !== 'object' || sms === null) return; - this.emitted.get(sms)?.(); + this.holds.get(sms)?.listenerGaveUp(); } /** Waits out the messages the application still holds, and says how many it never answered. */ @@ -309,8 +310,7 @@ export class IncomingRequests { if (!first) return; const generation = this.linkGeneration; - // A turn later, so a listener sending its receipt straight after the response still holds. - const release = (): void => { setImmediate(() => { this.held.release(pduObjs); }); }; + const hold = this.held.hold(pduObjs, this.deps.smsListeners()); const sms = this.deps.createSms({ answeredAs, @@ -323,22 +323,13 @@ export class IncomingRequests { answer: (pduObj, status, params) => this.deps.answer(pduObj, status, params), bindAllows: cmdName => this.deps.bindAllows(cmdName), lostLink: () => this.linkGeneration !== generation, - onAnswered: release, + onAnswered: () => { hold.answered(); }, // Past the refusal only while a drain is still waiting for this message; an ordinary send after. - send: input => (this.held.has(pduObjs) ? this.deps.sendPastDrain(input) : this.deps.send(input)), + send: input => (hold.isHeld() ? this.deps.sendPastDrain(input) : this.deps.send(input)), }); - // A rejection leaves the other listeners running, so only the last one to fail gives the message up. - let working = this.deps.smsListeners(); + this.holds.set(sms, hold); - this.held.hold(pduObjs); - this.emitted.set(sms, () => { - working--; - - if (working <= 0) release(); - }); - - // A message nobody took is not work a shutdown can wait for. - if (!this.deps.offerSms(sms)) this.held.release(pduObjs); + if (!this.deps.offerSms(sms)) hold.untaken(); } } diff --git a/todo.md b/todo.md index 0551ce7..3a4009a 100644 --- a/todo.md +++ b/todo.md @@ -199,15 +199,6 @@ next work ([decision](docs/decisions.md#internals-and-tests)). ### Locality — next, ahead of everything below; 5–6 today, and the gate is 7 -- [ ] **Give the held-message protocol one name and one home.** `emitSms()` is the unit 8 of 9 - readers named and 4 would least want to modify, and every one proposed the same fix. It runs - five mechanisms in one scope: a hold keyed by array identity, a `working` counter seeded from - `listenerCount('sms')`, a `WeakMap` keyed by the `Sms` object, a `setImmediate`-deferred - release, and a captured `linkGeneration` — with the counter decremented from `session.ts`'s - `captureRejectionSymbol` in another file. A `MessageHold` owning `hold/release/listenerGaveUp` - collapses three files into one readable object. Every way of getting it wrong is silent: a hung - shutdown, or a receipt refused. - - [ ] **Derive `Reassembler`'s octet total instead of maintaining it at five sites.** `this.octets` and each `group.octets` must agree, adjusted in `collect`, `trim`, `takeOldest`, `sweep` and `clear`, and `collect()` discovers its own eviction by re-reading the map by identity. Push the -- 2.52.0 From 9dda1a6917637cafd34446071e26072c22592bc6 Mon Sep 17 00:00:00 2001 From: Lilleman auf Larv Date: Mon, 28 Sep 2026 01:02:43 +0200 Subject: [PATCH 2/3] Release a held message through the hold it was given --- test/session-extras.test.ts | 28 +++++++++++++--------------- 1 file changed, 13 insertions(+), 15 deletions(-) diff --git a/test/session-extras.test.ts b/test/session-extras.test.ts index 8e23ed3..9bf67ed 100644 --- a/test/session-extras.test.ts +++ b/test/session-extras.test.ts @@ -1502,15 +1502,14 @@ describe('held message bounds', () => { test('is full at its count, and a re-used sequence number replaces rather than adding', () => { const held = new HeldMessages({ log: silentLog, max: 2, maxOctets: 1_000_000, timeout: 10_000 }); - const first = message(1); + const first = held.hold(message(1), 1); - held.hold(first); - held.hold(message(2)); - held.hold(message(2)); + held.hold(message(2), 1); + held.hold(message(2), 1); assert.equal(held.size, 2); assert.equal(held.full(), true); - assert.equal(held.has(first), true); + assert.equal(first.isHeld(), true); held.clear(); }); @@ -1519,23 +1518,22 @@ describe('held message bounds', () => { test('is full at its octet cap, until a message leaves by any way out', () => { let now = 0; const held = new HeldMessages({ log: silentLog, max: 10, maxOctets: 2000, now: () => now, timeout: 10_000 }); - const answered = message(1); + const answered = held.hold(message(1), 1); - held.hold(answered); assert.equal(held.full(), false); - held.hold(message(2)); + held.hold(message(2), 1); assert.equal(held.full(), true); - held.release(answered); + answered.release(); assert.equal(held.full(), false, 'after a release'); - held.hold(message(3)); + held.hold(message(3), 1); now = 20_000; held.sweep(); now = 0; assert.equal(held.full(), false, 'after a sweep'); - held.hold(message(4)); - held.hold(message(5)); + held.hold(message(4), 1); + held.hold(message(5), 1); held.clear(); assert.equal(held.full(), false, 'after a clear'); @@ -1624,11 +1622,11 @@ describe('held message bounds', () => { let now = 0; const held = new HeldMessages({ log: silentLog, max: 10, maxOctets: 1_000_000, now: () => now, timeout: 60 }); - held.hold(message(1)); + held.hold(message(1), 1); now = 61; // The next message sweeps the one that expired, so only the new one is still waited for. - held.hold(message(2)); + held.hold(message(2), 1); assert.equal(held.size, 1); @@ -1640,7 +1638,7 @@ describe('held message bounds', () => { let now = 0; const held = new HeldMessages({ log: silentLog, max: 10, maxOctets: 1_000_000, now: () => now, timeout: 60 }); - held.hold(message(1)); + held.hold(message(1), 1); const waiting = held.idle(1000, undefined); -- 2.52.0 From 0b9b210fb27cde7956fb655f13a55867e0462ea3 Mon Sep 17 00:00:00 2001 From: Lilleman auf Larv Date: Mon, 28 Sep 2026 01:03:15 +0200 Subject: [PATCH 3/3] Keep a hold's release behind MessageHold, and require its listener count --- src/held-messages.ts | 36 ++++++++++++++++++++---------------- src/incoming-requests.ts | 3 +-- 2 files changed, 21 insertions(+), 18 deletions(-) diff --git a/src/held-messages.ts b/src/held-messages.ts index 20f4cbd..07ac2fa 100644 --- a/src/held-messages.ts +++ b/src/held-messages.ts @@ -25,26 +25,29 @@ type Held = { pduObjs: PduObject[]; }; +type HoldEntry = { + isHeld: () => boolean; + release: () => void; +}; + /** One message offered to the application, held until it is answered or every listener gives up. */ export class MessageHold { - private readonly held: HeldMessages; - private readonly pduObjs: PduObject[]; + private readonly entry: HoldEntry; private working: number; - constructor(held: HeldMessages, pduObjs: PduObject[], listeners: number) { - this.held = held; - this.pduObjs = pduObjs; + constructor(entry: HoldEntry, listeners: number) { + this.entry = entry; this.working = listeners; } /** Whether a drain is still waiting for this message to be answered. */ isHeld(): boolean { - return this.held.has(this.pduObjs); + return this.entry.isHeld(); } /** A turn later, so a listener sending its receipt straight after the response still holds. */ answered(): void { - setImmediate(() => { this.held.release(this.pduObjs); }); + setImmediate(() => { this.entry.release(); }); } /** A rejection leaves the other listeners running, so only the last one to fail gives the message up. */ @@ -54,9 +57,9 @@ export class MessageHold { if (this.working <= 0) this.answered(); } - /** A message nobody took is not work a shutdown can wait for. */ - untaken(): void { - this.held.release(this.pduObjs); + /** At once, for a message nobody took: that is not work a shutdown can wait for. */ + release(): void { + this.entry.release(); } } @@ -94,9 +97,11 @@ export class HeldMessages { return this.held.full || this.octets >= this.maxOctets; } - /** Holds a message about to be offered to `listeners` listeners. */ - hold(pduObjs: PduObject[], listeners = 1): MessageHold { - const hold = new MessageHold(this, pduObjs, listeners); + hold(pduObjs: PduObject[], listeners: number): MessageHold { + const hold = new MessageHold({ + isHeld: () => this.has(pduObjs), + release: () => { this.release(pduObjs); }, + }, listeners); const key = keyOf(pduObjs); if (key === undefined) return hold; @@ -118,14 +123,13 @@ export class HeldMessages { return hold; } - /** Whether a drain is still waiting for this message to be answered. */ - has(pduObjs: PduObject[]): boolean { + private has(pduObjs: PduObject[]): boolean { const key = keyOf(pduObjs); return key !== undefined && this.held.get(key)?.pduObjs === pduObjs; } - release(pduObjs: PduObject[]): void { + private release(pduObjs: PduObject[]): void { const key = keyOf(pduObjs); // Identity, not the key: a wrapped sequence number must not release someone else's message. diff --git a/src/incoming-requests.ts b/src/incoming-requests.ts index 8985e5c..4f926de 100644 --- a/src/incoming-requests.ts +++ b/src/incoming-requests.ts @@ -324,12 +324,11 @@ export class IncomingRequests { bindAllows: cmdName => this.deps.bindAllows(cmdName), lostLink: () => this.linkGeneration !== generation, onAnswered: () => { hold.answered(); }, - // Past the refusal only while a drain is still waiting for this message; an ordinary send after. send: input => (hold.isHeld() ? this.deps.sendPastDrain(input) : this.deps.send(input)), }); this.holds.set(sms, hold); - if (!this.deps.offerSms(sms)) hold.untaken(); + if (!this.deps.offerSms(sms)) hold.release(); } } -- 2.52.0