188 lines
5.2 KiB
TypeScript
188 lines
5.2 KiB
TypeScript
import type { PduObject } from './pdu.ts';
|
|
import type { SmppLog } from './log.ts';
|
|
import { ExpiringGroups } from './expiring-groups.ts';
|
|
import { IdleWaiters } from './idle-waiters.ts';
|
|
import { retainedOctets } from './retained-pdu.ts';
|
|
|
|
export type HeldMessagesOptions = {
|
|
log: SmppLog;
|
|
max: number;
|
|
maxOctets: number;
|
|
/** Injected so expiry can be exercised without a wall clock. */
|
|
now?: (() => number) | undefined;
|
|
timeout: number;
|
|
};
|
|
|
|
/** The peer's own sequence number, which is what our answer to this message will carry. */
|
|
function keyOf(pduObjs: PduObject[]): string | undefined {
|
|
const first = pduObjs[0];
|
|
|
|
return first ? String(first.seqNr) : undefined;
|
|
}
|
|
|
|
type HoldEntry = {
|
|
isHeld: () => boolean;
|
|
release: () => void;
|
|
};
|
|
|
|
/**
|
|
* 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;
|
|
|
|
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.entry.isHeld();
|
|
}
|
|
|
|
/** A turn later, so a listener sending its receipt straight after the response still holds. */
|
|
answered(): void {
|
|
setImmediate(() => { this.entry.release(); });
|
|
}
|
|
|
|
/** 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();
|
|
}
|
|
|
|
/** 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();
|
|
}
|
|
}
|
|
|
|
/** The messages handed to the application that it has not answered yet, held by their segments. */
|
|
export class HeldMessages {
|
|
private readonly held: ExpiringGroups<PduObject[]>;
|
|
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({
|
|
max: options.max,
|
|
now: options.now,
|
|
onSweep: () => { this.sweep(); },
|
|
timeout: options.timeout,
|
|
});
|
|
this.log = options.log;
|
|
this.maxOctets = options.maxOctets;
|
|
}
|
|
|
|
get octetsHeld(): number {
|
|
return this.held.weight;
|
|
}
|
|
|
|
get size(): number {
|
|
return this.held.size;
|
|
}
|
|
|
|
/** Whether a message arriving now is past the bound, once the expired are swept. */
|
|
full(): boolean {
|
|
this.sweep();
|
|
|
|
return this.held.full || this.held.weight >= this.maxOctets;
|
|
}
|
|
|
|
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;
|
|
|
|
this.sweep();
|
|
|
|
if (this.held.get(key)) {
|
|
this.log.warn('heldMessages - replacing a message on a re-used sequence number', { seqNr: Number(key) });
|
|
}
|
|
|
|
this.held.set(key, pduObjs);
|
|
this.held.weigh(key, pduObjs.reduce((sum, pduObj) => sum + retainedOctets(pduObj), 0));
|
|
|
|
return hold;
|
|
}
|
|
|
|
/** Holds a message, builds what the application is handed, and emits it. */
|
|
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);
|
|
|
|
return key !== undefined && this.held.get(key) === pduObjs;
|
|
}
|
|
|
|
private release(pduObjs: PduObject[]): void {
|
|
const key = keyOf(pduObjs);
|
|
|
|
// Identity, not the key: a wrapped sequence number must not release someone else's message.
|
|
if (key === undefined || this.held.get(key) !== pduObjs) return;
|
|
|
|
this.held.delete(key);
|
|
this.settle();
|
|
}
|
|
|
|
/** Drops every message: their segments went with the link, so no answer of ours correlates now. */
|
|
clear(): void {
|
|
this.held.takeAll();
|
|
this.idleWaiters.settle();
|
|
}
|
|
|
|
/** Resolves 0 once every message has been answered, or with how many have not. */
|
|
idle(timeout: number, signal: AbortSignal | undefined): Promise<number> {
|
|
return this.idleWaiters.wait(() => this.held.size, timeout, signal);
|
|
}
|
|
|
|
/** Drops every message past its deadline. Runs before each hold and on its own timer. */
|
|
sweep(): void {
|
|
const expired = this.held.takeExpired();
|
|
|
|
if (expired.length === 0) return;
|
|
|
|
this.log.warn('heldMessages - messages the application never answered', {
|
|
messages: expired.length,
|
|
});
|
|
this.settle();
|
|
}
|
|
|
|
private settle(): void {
|
|
if (this.held.size === 0) this.idleWaiters.settle();
|
|
}
|
|
}
|