From 0b9b210fb27cde7956fb655f13a55867e0462ea3 Mon Sep 17 00:00:00 2001 From: Lilleman auf Larv Date: Mon, 28 Sep 2026 01:03:15 +0200 Subject: [PATCH] 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(); } }