Give the held-message flow one place a reader can follow it
Test / lint (pull_request) Successful in 24s
Test / test (18) (pull_request) Successful in 32s
Test / test (20) (pull_request) Successful in 31s
Test / test (22) (pull_request) Successful in 32s
Test / test (24) (pull_request) Successful in 32s
Test / test (26) (pull_request) Successful in 33s
Mirror / push (push) Successful in 7s

This commit was merged in pull request #45.
This commit is contained in:
2026-09-28 10:54:55 +02:00
parent 1aadcbe3c9
commit 0bea36c069
4 changed files with 57 additions and 31 deletions
+32 -3
View File
@@ -25,7 +25,11 @@ type HoldEntry = {
release: () => void;
};
/** One message offered to the application, held until it is answered or every listener gives up. */
/**
* One message offered to the application. A drain waits on it until the first of: `answered()`,
* every listener that took it rejecting, no listener taking it or one throwing, a later message on
* its sequence number, its deadline, or the link going.
*/
export class MessageHold {
private readonly entry: HoldEntry;
private working: number;
@@ -52,7 +56,7 @@ export class MessageHold {
if (this.working <= 0) this.answered();
}
/** At once, for a message nobody took: that is not work a shutdown can wait for. */
/** At once, for a message nobody took or a listener threw on: that is not work a shutdown can wait for. */
release(): void {
this.entry.release();
}
@@ -64,6 +68,8 @@ export class HeldMessages {
private readonly idleWaiters = new IdleWaiters();
private readonly log: SmppLog;
private readonly maxOctets: number;
/** A rejecting listener hands the message back as an `unknown`, so its hold is found by identity. */
private readonly offered = new WeakMap<object, MessageHold>();
constructor(options: HeldMessagesOptions) {
this.held = new ExpiringGroups({
@@ -91,7 +97,7 @@ export class HeldMessages {
return this.held.full || this.held.weight >= this.maxOctets;
}
hold(pduObjs: PduObject[], listeners: number): MessageHold {
private hold(pduObjs: PduObject[], listeners: number): MessageHold {
const hold = new MessageHold({
isHeld: () => this.has(pduObjs),
release: () => { this.release(pduObjs); },
@@ -112,6 +118,29 @@ export class HeldMessages {
return hold;
}
offer<T extends object>(
pduObjs: PduObject[],
listeners: number,
build: (hold: MessageHold) => T,
emit: (message: T) => boolean,
): MessageHold {
const hold = this.hold(pduObjs, listeners);
const message = build(hold);
this.offered.set(message, hold);
if (!emit(message)) hold.release();
return hold;
}
/** One listener gave up on a message; the last one to do so is what releases it. */
listenerRejected(message: unknown): void {
if (typeof message !== 'object' || message === null) return;
this.offered.get(message)?.listenerGaveUp();
}
private has(pduObjs: PduObject[]): boolean {
const key = keyOf(pduObjs);
+3 -14
View File
@@ -10,7 +10,6 @@ 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';
@@ -85,8 +84,6 @@ export class IncomingRequests {
private readonly deps: IncomingDeps;
private readonly dlrMerger: DlrMerger;
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<object, MessageHold>();
private readonly log: SmppLog;
private readonly reassembler: Reassembler;
private readonly smsIdFormat: SmsIdFormat;
@@ -172,11 +169,8 @@ export class IncomingRequests {
this.reassembler.clear();
}
/** One `sms` listener gave up on a message; the last one to do so is what releases the hold. */
listenerRejected(sms: unknown): void {
if (typeof sms !== 'object' || sms === null) return;
this.holds.get(sms)?.listenerGaveUp();
this.held.listenerRejected(sms);
}
/** Waits out the messages the application still holds, and says how many it never answered. */
@@ -310,9 +304,8 @@ export class IncomingRequests {
if (!first) return;
const generation = this.linkGeneration;
const hold = this.held.hold(pduObjs, this.deps.smsListeners());
const sms = this.deps.createSms({
this.held.offer(pduObjs, this.deps.smsListeners(), hold => this.deps.createSms({
answeredAs,
from: paramText(first.params.source_addr),
message: decodeSegments(pduObjs),
@@ -325,10 +318,6 @@ export class IncomingRequests {
lostLink: () => this.linkGeneration !== generation,
onAnswered: () => { hold.answered(); },
send: input => (hold.isHeld() ? this.deps.sendPastDrain(input) : this.deps.send(input)),
});
this.holds.set(sms, hold);
if (!this.deps.offerSms(sms)) hold.release();
}), sms => this.deps.offerSms(sms));
}
}