Give the held-message flow one owner
Test / lint (pull_request) Successful in 23s
Test / test (18) (pull_request) Successful in 30s
Test / test (20) (pull_request) Successful in 30s
Test / test (22) (pull_request) Successful in 31s
Test / test (24) (pull_request) Successful in 30s
Test / test (26) (pull_request) Successful in 30s
Mirror / push (push) Successful in 5s

This commit was merged in pull request #49.
This commit is contained in:
2026-09-29 18:43:44 +02:00
parent 57f53d56f1
commit 1b075bf5d2
10 changed files with 144 additions and 140 deletions
+51 -36
View File
@@ -1,15 +1,23 @@
import type { PduObject } from './pdu.ts';
import type { LinkLife } from './link-life.ts';
import type { PduObject, PduObjectInput } from './pdu.ts';
import type { Result } from './result.ts';
import type { Session } from './session.ts';
import type { SmsHandlers } from './sms.ts';
import type { SmppLog } from './log.ts';
import { ExpiringGroups } from './expiring-groups.ts';
import { IdleWaiters } from './idle-waiters.ts';
import { createSms } from './sms.ts';
import { retainedOctets } from './retained-pdu.ts';
export type HeldMessagesOptions = {
link: LinkLife;
log: SmppLog;
max: number;
maxOctets: number;
/** Injected so expiry can be exercised without a wall clock. */
now?: (() => number) | undefined;
sendPastDrain: SmsHandlers['send'];
session: Session;
timeout: number;
};
@@ -20,33 +28,40 @@ function keyOf(pduObjs: PduObject[]): string | undefined {
return first ? String(first.seqNr) : undefined;
}
type HoldEntry = {
isHeld: () => boolean;
release: () => void;
};
type HoldRoute = Pick<HeldMessagesOptions, 'link' | 'sendPastDrain' | 'session'>;
/**
* 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.
* One message offered to the application, and the handlers its `Sms` answers through. 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;
export class MessageHold implements SmsHandlers {
private readonly generation: number;
private readonly heldMessages: HeldMessages;
private readonly pduObjs: PduObject[];
private readonly route: HoldRoute;
private working: number;
constructor(entry: HoldEntry, listeners: number) {
this.entry = entry;
constructor(heldMessages: HeldMessages, route: HoldRoute, pduObjs: PduObject[], listeners: number) {
this.generation = route.link.generation();
this.heldMessages = heldMessages;
this.pduObjs = pduObjs;
this.route = route;
this.working = listeners;
}
/** Whether a drain is still waiting for this message to be answered. */
isHeld(): boolean {
return this.entry.isHeld();
return this.heldMessages.holds(this.pduObjs);
}
/** A turn later, so a listener sending its receipt straight after the response still holds. */
/** A turn later, so a `sendDlr()` called straight after `sendResp()` still goes out past a drain. */
answered(): void {
setImmediate(() => { this.entry.release(); });
setImmediate(() => { this.release(); });
}
lostLink(): boolean {
return this.route.link.generation() !== this.generation;
}
/** A rejection leaves the other listeners running, so only the last one to fail gives the message up. */
@@ -58,7 +73,12 @@ export class MessageHold {
/** 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();
this.heldMessages.release(this.pduObjs);
}
/** A receipt for a message still held is what a drain waits for, so it goes out past the drain. */
send(input: PduObjectInput): Promise<Result<{ pduObj: PduObject }>> {
return this.isHeld() ? this.route.sendPastDrain(input) : this.route.session.send(input);
}
}
@@ -70,6 +90,7 @@ export class HeldMessages {
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>();
private readonly route: HoldRoute;
constructor(options: HeldMessagesOptions) {
this.held = new ExpiringGroups({
@@ -80,6 +101,7 @@ export class HeldMessages {
});
this.log = options.log;
this.maxOctets = options.maxOctets;
this.route = { link: options.link, sendPastDrain: options.sendPastDrain, session: options.session };
}
get octetsHeld(): number {
@@ -97,14 +119,8 @@ export class HeldMessages {
return this.held.full || this.held.weight >= this.maxOctets;
}
private 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;
private hold(key: string, pduObjs: PduObject[], listeners: number): MessageHold {
const hold = new MessageHold(this, this.route, pduObjs, listeners);
this.sweep();
@@ -118,18 +134,17 @@ 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);
offer(pduObjs: PduObject[], answeredAs?: string): MessageHold | undefined {
const key = keyOf(pduObjs);
this.offered.set(message, hold);
if (key === undefined) return undefined;
if (!emit(message)) hold.release();
const hold = this.hold(key, pduObjs, this.route.session.listenerCount('sms'));
const sms = createSms({ answeredAs, pduObjs, session: this.route.session }, hold);
this.offered.set(sms, hold);
if (!this.route.session.emit('sms', sms)) hold.release();
return hold;
}
@@ -141,13 +156,13 @@ export class HeldMessages {
this.offered.get(message)?.listenerGaveUp();
}
private has(pduObjs: PduObject[]): boolean {
holds(pduObjs: PduObject[]): boolean {
const key = keyOf(pduObjs);
return key !== undefined && this.held.get(key) === pduObjs;
}
private release(pduObjs: PduObject[]): void {
release(pduObjs: PduObject[]): void {
const key = keyOf(pduObjs);
// Identity, not the key: a wrapped sequence number must not release someone else's message.
+10 -32
View File
@@ -1,22 +1,21 @@
import type { Concat } from './concat.ts';
import type { DlrMerger } from './dlr-merger.ts';
import type { ErrorName } from './defs/errors.ts';
import type { HeldMessagesOptions } from './held-messages.ts';
import type { LinkLife } from './link-life.ts';
import type { LostGroup, Refusal } from './reassembly.ts';
import type { OnRequest } from './session-options.ts';
import type { PduObject, PduObjectInput } from './pdu.ts';
import type { Result, VoidResult } from './result.ts';
import type { PduObject } from './pdu.ts';
import type { VoidResult } from './result.ts';
import type { Session } from './session.ts';
import type { SmppLog } from './log.ts';
import type { SmsIdFormat } from './sms-id.ts';
import { HeldMessages } from './held-messages.ts';
import { Reassembler, decodeSegments } from './reassembly.ts';
import { Reassembler } from './reassembly.ts';
import { bindCommands, defaults, standsInFor } from './session-options.ts';
import { concatOf } from './concat.ts';
import { createSms } from './sms.ts';
import { detach } from './retained-pdu.ts';
import { dlrFromPdu } from './dlr.ts';
import { paramText } from './defs/types.ts';
import { respIdParams, segmentId } from './sms-id.ts';
import { respNameFor } from './defs/commands.ts';
@@ -52,8 +51,7 @@ export type IncomingRequestsOptions = {
maxReassembly?: number | undefined;
onRequest?: OnRequest | undefined;
reassemblyTimeout?: number | undefined;
/** Past a drain's refusal, for a receipt the drain is itself waiting for. */
sendPastDrain: (input: PduObjectInput) => Promise<Result<{ pduObj: PduObject }>>;
sendPastDrain: HeldMessagesOptions['sendPastDrain'];
session: Session;
smsIdFormat?: SmsIdFormat | undefined;
systemId?: string | undefined;
@@ -67,7 +65,6 @@ export class IncomingRequests {
private readonly log: SmppLog;
private readonly onRequest: OnRequest | undefined;
private readonly reassembler: Reassembler;
private readonly sendPastDrain: IncomingRequestsOptions['sendPastDrain'];
private readonly session: Session;
private readonly smsIdFormat: SmsIdFormat;
private readonly systemId: string;
@@ -76,9 +73,12 @@ export class IncomingRequests {
constructor(options: IncomingRequestsOptions) {
this.dlrMerger = options.dlrMerger;
this.held = new HeldMessages({
link: options.link,
log: options.log,
max: defaults.maxHeldMessages,
maxOctets: defaults.maxHeldOctets,
sendPastDrain: options.sendPastDrain,
session: options.session,
timeout: defaults.heldMessageTimeout,
});
this.link = options.link;
@@ -91,7 +91,6 @@ export class IncomingRequests {
onLost: lost => { this.reportLost(lost); },
timeout: options.reassemblyTimeout ?? defaults.reassemblyTimeout,
});
this.sendPastDrain = options.sendPastDrain;
this.session = options.session;
this.smsIdFormat = options.smsIdFormat ?? {};
this.systemId = options.systemId ?? defaults.systemId;
@@ -254,7 +253,7 @@ export class IncomingRequests {
const concat = concatOf(pduObj);
if (!concat) {
this.emitSms([detach(pduObj)]);
this.held.offer([detach(pduObj)]);
return;
}
@@ -276,7 +275,7 @@ export class IncomingRequests {
respIdParams(pduObj.cmdName, segmentId(collected.smsId, concat.part - 1, concat.total)),
);
if (collected.whole) this.emitSms(collected.whole, collected.smsId);
if (collected.whole) this.held.offer(collected.whole, collected.smsId);
}
private reportLost(lost: LostGroup): void {
@@ -284,25 +283,4 @@ export class IncomingRequests {
`Gave up ${String(lost.parts)} of ${String(lost.total)} segments of an incomplete concatenated message: ${lostReasons[lost.reason]}`,
));
}
private emitSms(pduObjs: PduObject[], answeredAs?: string): void {
const first = pduObjs[0];
if (!first) return;
const generation = this.link.generation();
this.held.offer(pduObjs, this.session.listenerCount('sms'), hold => createSms({
answeredAs,
from: paramText(first.params.source_addr),
message: decodeSegments(pduObjs),
pduObjs,
session: this.session,
to: paramText(first.params.destination_addr),
}, {
lostLink: () => this.link.generation() !== generation,
onAnswered: () => { hold.answered(); },
send: input => (hold.isHeld() ? this.sendPastDrain(input) : this.session.send(input)),
}), sms => this.session.emit('sms', sms));
}
}
+1 -1
View File
@@ -38,7 +38,7 @@ export { uuidv7 } from './uuid.ts';
export type { BindType, ClientOptions } from './client.ts';
export type { Dlr, Receipt } from './dlr.ts';
export type { SendDlrResult, SendRespOptions, Sms, SmsInput } from './sms.ts';
export type { SendDlrResult, SendRespOptions, Sms } from './sms.ts';
export type { Concat } from './concat.ts';
export type { ConcatInfo } from './udh.ts';
export type { Result, VoidResult } from './result.ts';
+10 -12
View File
@@ -5,7 +5,9 @@ import type { Result, VoidResult } from './result.ts';
import type { Session } from './session.ts';
import { UnansweredError } from './unanswered-error.ts';
import { consts } from './defs/constants.ts';
import { decodeSegments } from './reassembly.ts';
import { messageClassOf } from './defs/encodings.ts';
import { paramText } from './defs/types.ts';
import { receiptCodes, transientStates } from './dlr.ts';
import { smppDate } from './message.ts';
import { respIdParams, segmentId } from './sms-id.ts';
@@ -59,17 +61,13 @@ export type Sms = {
export type SmsInput = {
/** The id base the segments were already answered with; absent leaves the answer to `sendResp()`. */
answeredAs?: string | undefined;
from: string;
message: string;
pduObjs: PduObject[];
session: Session;
to: string;
};
/** What the session's incoming side gives a message so it can be answered and accounted for. */
export type SmsHandlers = {
answered: () => void;
lostLink: () => boolean;
onAnswered: () => void;
send: (input: PduObjectInput) => Promise<Result<{ pduObj: PduObject }>>;
};
@@ -86,8 +84,8 @@ export function createSms(input: SmsInput, handlers: SmsHandlers): Sms {
answeredOnArrival: input.answeredAs !== undefined,
dlr: typeof registered === 'number' && registered !== 0,
flash: typeof dataCoding === 'number' && messageClassOf(dataCoding) === immediateDisplayClass,
from: input.from,
message: input.message,
from: paramText(first?.params.source_addr),
message: decodeSegments(input.pduObjs),
pduObjs: input.pduObjs,
sendDlr: status => sendDlr(sms, input.session, handlers, status),
sendResp: options => (input.answeredAs === undefined
@@ -98,7 +96,7 @@ export function createSms(input: SmsInput, handlers: SmsHandlers): Sms {
return answered.smsId;
},
submitTime: new Date(),
to: input.to,
to: paramText(first?.params.destination_addr),
};
return sms;
@@ -107,7 +105,7 @@ export function createSms(input: SmsInput, handlers: SmsHandlers): Sms {
/** Every segment went out answered, so the call is what the shutdown waits for and nothing else. */
function answeredOnArrival(
options: SendRespOptions,
handlers: Pick<SmsHandlers, 'onAnswered'>,
handlers: Pick<SmsHandlers, 'answered'>,
): Promise<VoidResult> {
if (options.smsId !== undefined) {
return Promise.resolve({
@@ -121,7 +119,7 @@ function answeredOnArrival(
});
}
handlers.onAnswered();
handlers.answered();
return Promise.resolve({});
}
@@ -131,7 +129,7 @@ async function sendResp(
session: Session,
answered: { smsId: string },
options: SendRespOptions,
handlers: Pick<SmsHandlers, 'lostLink' | 'onAnswered'>,
handlers: Pick<SmsHandlers, 'answered' | 'lostLink'>,
): Promise<VoidResult> {
const total = sms.pduObjs.length;
@@ -158,7 +156,7 @@ async function sendResp(
const failure = results.find(result => result.err);
if (!failure) handlers.onAnswered();
if (!failure) handlers.answered();
return failure ?? {};
}