Give the held-message protocol one home in MessageHold
Test / lint (pull_request) Successful in 23s
Test / test (18) (pull_request) Successful in 34s
Test / test (20) (pull_request) Successful in 31s
Test / test (22) (pull_request) Successful in 39s
Test / test (24) (pull_request) Successful in 32s
Test / test (26) (pull_request) Successful in 39s
Mirror / push (push) Successful in 5s
Test / lint (pull_request) Successful in 23s
Test / test (18) (pull_request) Successful in 34s
Test / test (20) (pull_request) Successful in 31s
Test / test (22) (pull_request) Successful in 39s
Test / test (24) (pull_request) Successful in 32s
Test / test (26) (pull_request) Successful in 39s
Mirror / push (push) Successful in 5s
This commit is contained in:
+41
-2
@@ -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<Held>;
|
||||
@@ -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. */
|
||||
|
||||
Reference in New Issue
Block a user