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