Keep a hold's release behind MessageHold, and require its listener count
Test / lint (pull_request) Successful in 23s
Test / test (18) (pull_request) Successful in 31s
Test / test (20) (pull_request) Successful in 31s
Test / test (22) (pull_request) Successful in 31s
Test / test (24) (pull_request) Successful in 31s
Test / test (26) (pull_request) Successful in 32s
Mirror / push (push) Successful in 5s

This commit was merged in pull request #35.
This commit is contained in:
2026-09-28 01:03:15 +02:00
parent 9dda1a6917
commit 0b9b210fb2
2 changed files with 21 additions and 18 deletions
+20 -16
View File
@@ -25,26 +25,29 @@ type Held = {
pduObjs: PduObject[]; pduObjs: PduObject[];
}; };
type HoldEntry = {
isHeld: () => boolean;
release: () => void;
};
/** One message offered to the application, held until it is answered or every listener gives up. */ /** One message offered to the application, held until it is answered or every listener gives up. */
export class MessageHold { export class MessageHold {
private readonly held: HeldMessages; private readonly entry: HoldEntry;
private readonly pduObjs: PduObject[];
private working: number; private working: number;
constructor(held: HeldMessages, pduObjs: PduObject[], listeners: number) { constructor(entry: HoldEntry, listeners: number) {
this.held = held; this.entry = entry;
this.pduObjs = pduObjs;
this.working = listeners; this.working = listeners;
} }
/** Whether a drain is still waiting for this message to be answered. */ /** Whether a drain is still waiting for this message to be answered. */
isHeld(): boolean { 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. */ /** A turn later, so a listener sending its receipt straight after the response still holds. */
answered(): void { 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. */ /** 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(); if (this.working <= 0) this.answered();
} }
/** A message nobody took is not work a shutdown can wait for. */ /** At once, for a message nobody took: that is not work a shutdown can wait for. */
untaken(): void { release(): void {
this.held.release(this.pduObjs); this.entry.release();
} }
} }
@@ -94,9 +97,11 @@ export class HeldMessages {
return this.held.full || this.octets >= this.maxOctets; return this.held.full || this.octets >= this.maxOctets;
} }
/** Holds a message about to be offered to `listeners` listeners. */ hold(pduObjs: PduObject[], listeners: number): MessageHold {
hold(pduObjs: PduObject[], listeners = 1): MessageHold { const hold = new MessageHold({
const hold = new MessageHold(this, pduObjs, listeners); isHeld: () => this.has(pduObjs),
release: () => { this.release(pduObjs); },
}, listeners);
const key = keyOf(pduObjs); const key = keyOf(pduObjs);
if (key === undefined) return hold; if (key === undefined) return hold;
@@ -118,14 +123,13 @@ export class HeldMessages {
return hold; return hold;
} }
/** Whether a drain is still waiting for this message to be answered. */ private has(pduObjs: PduObject[]): boolean {
has(pduObjs: PduObject[]): boolean {
const key = keyOf(pduObjs); const key = keyOf(pduObjs);
return key !== undefined && this.held.get(key)?.pduObjs === pduObjs; return key !== undefined && this.held.get(key)?.pduObjs === pduObjs;
} }
release(pduObjs: PduObject[]): void { private release(pduObjs: PduObject[]): void {
const key = keyOf(pduObjs); const key = keyOf(pduObjs);
// Identity, not the key: a wrapped sequence number must not release someone else's message. // Identity, not the key: a wrapped sequence number must not release someone else's message.
+1 -2
View File
@@ -324,12 +324,11 @@ export class IncomingRequests {
bindAllows: cmdName => this.deps.bindAllows(cmdName), bindAllows: cmdName => this.deps.bindAllows(cmdName),
lostLink: () => this.linkGeneration !== generation, lostLink: () => this.linkGeneration !== generation,
onAnswered: () => { hold.answered(); }, 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)), send: input => (hold.isHeld() ? this.deps.sendPastDrain(input) : this.deps.send(input)),
}); });
this.holds.set(sms, hold); this.holds.set(sms, hold);
if (!this.deps.offerSms(sms)) hold.untaken(); if (!this.deps.offerSms(sms)) hold.release();
} }
} }