Move src into codec, protocol, messages, session, client and server
Mirror / push (push) Has been cancelled
Test / lint (pull_request) Successful in 22s
Test / test (18) (pull_request) Successful in 32s
Test / test (20) (pull_request) Successful in 29s
Test / test (22) (pull_request) Successful in 32s
Test / test (24) (pull_request) Successful in 29s
Test / test (26) (pull_request) Successful in 36s

This commit is contained in:
2026-09-30 20:15:04 +02:00
parent 108eb1e5b2
commit 980e9d4569
79 changed files with 573 additions and 610 deletions
+201
View File
@@ -0,0 +1,201 @@
import type { LinkLife } from '../link-life.ts';
import type { PduObject, PduObjectInput } from '../codec/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 '../messages/expiring-groups.ts';
import { IdleWaiters } from './waiting.ts';
import { createSms } from '../sms.ts';
import { retainedOctets } from '../codec/retained.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;
};
/** 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 HoldRoute = Pick<HeldMessagesOptions, 'link' | 'sendPastDrain' | 'session'>;
/**
* 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 implements SmsHandlers {
private readonly generation: number;
private readonly heldMessages: HeldMessages;
private readonly pduObjs: PduObject[];
private readonly route: HoldRoute;
private working: number;
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.heldMessages.holds(this.pduObjs);
}
/** A turn later, so a `sendDlr()` called straight after `sendResp()` still goes out past a drain. */
answered(): void {
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. */
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.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);
}
}
/** 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>();
private readonly route: HoldRoute;
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;
this.route = { link: options.link, sendPastDrain: options.sendPastDrain, session: options.session };
}
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;
}
private hold(key: string, pduObjs: PduObject[], listeners: number): MessageHold {
const hold = new MessageHold(this, this.route, pduObjs, listeners);
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;
}
offer(pduObjs: PduObject[], answeredAs?: string): MessageHold | undefined {
const key = keyOf(pduObjs);
if (key === undefined) return undefined;
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;
}
/** 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();
}
holds(pduObjs: PduObject[]): boolean {
const key = keyOf(pduObjs);
return key !== undefined && this.held.get(key) === pduObjs;
}
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();
}
}
+50
View File
@@ -0,0 +1,50 @@
import type { SmppLog } from '../log.ts';
export type LinkTimersOptions = {
/** How long between enquire_link probes. Undefined or 0 never probes. */
enquireLinkInterval?: number | undefined;
/** How long a silent peer is kept. Undefined or 0 keeps it forever. */
idleTimeout?: number | undefined;
log: SmppLog;
onEnquireLink: () => void;
onIdle: () => void;
};
/** Keeps a quiet connection honest: probes the peer, and gives up on one that stays silent. */
export class LinkTimers {
private readonly options: LinkTimersOptions;
private enquireLink: NodeJS.Timeout | undefined;
private idle: NodeJS.Timeout | undefined;
constructor(options: LinkTimersOptions) {
this.options = options;
}
/** Starts both timers over, which every sign of life from the peer should do. */
reset(): void {
const { enquireLinkInterval, idleTimeout, log, onEnquireLink, onIdle } = this.options;
this.clear();
if (enquireLinkInterval !== undefined && enquireLinkInterval > 0) {
this.enquireLink = setTimeout(onEnquireLink, enquireLinkInterval);
this.enquireLink.unref();
}
if (idleTimeout !== undefined && idleTimeout > 0) {
this.idle = setTimeout(() => {
log.info('linkTimers - closing an idle peer', { idleTimeout });
onIdle();
}, idleTimeout);
this.idle.unref();
}
}
clear(): void {
if (this.enquireLink) clearTimeout(this.enquireLink);
if (this.idle) clearTimeout(this.idle);
this.enquireLink = undefined;
this.idle = undefined;
}
}
+178
View File
@@ -0,0 +1,178 @@
import type { LinkLife } from '../link-life.ts';
import type { PduObject, PduObjectInput } from '../codec/pdu.ts';
import type { PduTransport } from './transport.ts';
import type { Result, VoidResult } from '../result.ts';
import type { SendOptions } from '../options.ts';
import type { SmppLog } from '../log.ts';
import { PendingRequests } from './pending-requests.ts';
import { SendWindow } from './send-window.ts';
import { UnansweredError } from '../unanswered-error.ts';
import { bindCommands } from '../protocol/bind.ts';
import { objToPdu } from '../codec/pdu.ts';
export type OutgoingRequestsOptions = {
link: LinkLife;
log: SmppLog;
maxOutstanding: number;
responseTimeout: number;
transport: PduTransport;
};
/** `retryOnNextLink`: the write failed, so nothing reached the socket and another link may carry it. */
type Attempt = { result: Result<{ pduObj: PduObject }>; retryOnNextLink: boolean };
function abortedBeforeSend(): Error {
return new Error('Aborted before the request was sent');
}
/** A response carries the request's sequence number, which only sendReturn() has. */
function misuse(input: PduObjectInput): Error | undefined {
return input.cmdName.endsWith('_resp')
? new Error(`Use sendReturn() for responses, not send(): ${input.cmdName}`)
: undefined;
}
/** Everything this end asks of the peer: which link carries it, how many at once, and the answer. */
export class OutgoingRequests {
private readonly link: LinkLife;
private readonly log: SmppLog;
private readonly pending: PendingRequests;
private readonly responseTimeout: number;
private readonly transport: PduTransport;
private readonly window: SendWindow;
constructor(options: OutgoingRequestsOptions) {
this.link = options.link;
this.log = options.log;
this.pending = new PendingRequests(options.log);
this.responseTimeout = options.responseTimeout;
this.transport = options.transport;
this.window = new SendWindow({ limit: options.maxOutstanding, log: options.log });
}
canCarry(): boolean {
return this.link.isUp() && !this.transport.sock.destroyed;
}
/** The link is gone, and every answer still owed on it with it. */
linkLost(): void {
this.pending.settleAll(new Error('Session closed before a response arrived'));
}
/** Hands a response to the request waiting for it. False means nothing was. */
deliver(pduObj: PduObject): boolean {
return this.pending.deliver(pduObj);
}
/** A response the codec refused settles its request instead of leaving it to time out. */
settleRefused(seqNr: number, err: Error): void {
this.pending.settle(seqNr, { err });
}
request(input: PduObjectInput, options: SendOptions): Promise<Result<{ pduObj: PduObject }>> {
// Ahead of the drain, so a misuse is named as one rather than blamed on the shutdown.
const wrong = misuse(input);
if (wrong) return Promise.resolve({ err: wrong });
// With no link, the request is refused as closed further on.
if (this.link.isStopped() && this.canCarry()) {
return Promise.resolve({ err: new Error('Session is shutting down') });
}
return this.requestPastDrain(input, options);
}
/** request() without the drain's refusal, which a receipt for a held message has to take. */
async requestPastDrain(
input: PduObjectInput,
options: SendOptions,
): Promise<Result<{ pduObj: PduObject }>> {
const refused = this.refuse(input, options);
if (refused) return { err: refused };
// A bind is what makes a link usable, so it cannot wait for one: it takes the link's answer now.
if (bindCommands.includes(input.cmdName)) {
const shut = this.link.refusal();
return shut ? { err: shut } : this.requestOnCurrentLink(input, options);
}
const waitForLink = this.link.hold(options.signal);
for (;;) {
const held = await waitForLink();
if (held.err) return { err: held.err };
const slot = await this.window.acquire(options.signal);
if (slot.err) return { err: slot.err };
const attempt = await this.attempt(input, options).finally(() => { this.window.release(); });
if (!this.retriesOnNextLink(attempt)) return attempt.result;
}
}
/** Straight onto the current link, for what has to go out either way. */
async requestOnCurrentLink(
input: PduObjectInput,
options: SendOptions = {},
): Promise<Result<{ pduObj: PduObject }>> {
return (await this.attempt(input, options)).result;
}
/** Waits out the requests already on the wire, and says how many never finished. */
async drain(timeout: number, signal: AbortSignal | undefined): Promise<VoidResult> {
const unfinished = await this.window.idle(timeout, signal);
if (unfinished === 0) return {};
this.log.warn('outgoingRequests - shutting down with requests unfinished', { timeout, unfinished });
return { err: new Error(`Shut down with ${String(unfinished)} request(s) unfinished`) };
}
/** Nothing reached the socket, so the next link carries it. */
private retriesOnNextLink(attempt: Attempt): boolean {
// Until the link is dropped it admits the retry straight back onto the dead socket, and the loop spins.
return attempt.retryOnNextLink && this.link.awaitsNextLink();
}
/** Why a request cannot go out at all, as opposed to not yet. */
private refuse(input: PduObjectInput, options: SendOptions): Error | undefined {
// Before the link and the window, or an aborted call waits for what it will never use.
return misuse(input) ?? (options.signal?.aborted === true ? abortedBeforeSend() : undefined);
}
private async attempt(input: PduObjectInput, options: SendOptions): Promise<Attempt> {
// pending.wait() alone settles the caller while the request still goes out to the peer.
if (options.signal?.aborted === true) {
return { result: { err: abortedBeforeSend() }, retryOnNextLink: false };
}
const seqNr = this.pending.nextSeqNr();
const built = objToPdu({ ...input, seqNr });
if (built.err) return { result: { err: built.err }, retryOnNextLink: false };
const response = this.pending.wait(seqNr, {
signal: options.signal,
timeout: this.responseTimeout,
});
const written = this.transport.write(built.buffer);
if (written.err) {
this.pending.settle(seqNr, { err: written.err });
return { result: { err: written.err }, retryOnNextLink: true };
}
const answered = await response;
// It went out, so a failure now means the peer may have taken it and the answer was the loss.
return { result: answered.err ? { err: new UnansweredError(answered.err) } : answered, retryOnNextLink: false };
}
}
+88
View File
@@ -0,0 +1,88 @@
import type { PduObject } from '../codec/pdu.ts';
import type { Result } from '../result.ts';
import type { SmppLog } from '../log.ts';
import { maxSeqNr } from '../codec/pdu.ts';
export type WaitOptions = {
signal?: AbortSignal | undefined;
timeout: number;
};
type Pending = {
settle: (result: Result<{ pduObj: PduObject }>) => void;
};
/** Hands out sequence numbers and matches responses to the requests waiting for them. */
export class PendingRequests {
private readonly log: SmppLog;
private readonly pending = new Map<number, Pending>();
private ourSeqNr = 1;
constructor(log: SmppLog) {
this.log = log;
}
nextSeqNr(): number {
const seqNr = this.ourSeqNr;
this.ourSeqNr = this.ourSeqNr >= maxSeqNr ? 1 : this.ourSeqNr + 1;
return seqNr;
}
wait(seqNr: number, options: WaitOptions): Promise<Result<{ pduObj: PduObject }>> {
const { signal, timeout } = options;
return new Promise(resolve => {
const abort = (): void => {
this.settle(seqNr, { err: new Error('Aborted before a response arrived') });
};
const timer = timeout > 0 ? this.expire(seqNr, timeout) : undefined;
this.pending.set(seqNr, {
settle: result => {
if (timer) clearTimeout(timer);
signal?.removeEventListener('abort', abort);
this.pending.delete(seqNr);
resolve(result);
},
});
if (signal?.aborted === true) abort();
else signal?.addEventListener('abort', abort, { once: true });
});
}
/** Hands a response to whoever is waiting for it. False means nobody was. */
deliver(pduObj: PduObject): boolean {
const pending = this.pending.get(pduObj.seqNr);
if (!pending) return false;
pending.settle({ pduObj });
return true;
}
settle(seqNr: number, result: Result<{ pduObj: PduObject }>): void {
this.pending.get(seqNr)?.settle(result);
}
settleAll(err: Error): void {
for (const [seqNr] of this.pending) {
this.settle(seqNr, { err });
}
}
private expire(seqNr: number, timeout: number): NodeJS.Timeout {
const timer = setTimeout(() => {
this.log.warn('pendingRequests - no response before the timeout', { seqNr, timeout });
this.settle(seqNr, { err: new Error(`No response to seqNr ${String(seqNr)}`) });
}, timeout);
timer.unref();
return timer;
}
}
+287
View File
@@ -0,0 +1,287 @@
import type { Concat } from '../protocol/concat.ts';
import type { DlrMerger } from '../messages/receipt-merge.ts';
import type { ErrorName } from '../codec/statuses.ts';
import type { HeldMessagesOptions } from './held-messages.ts';
import type { LinkLife } from '../link-life.ts';
import type { LostGroup, Refusal } from '../messages/reassembly.ts';
import type { OnRequest } from '../options.ts';
import type { PduObject } from '../codec/pdu.ts';
import type { VoidResult } from '../result.ts';
import type { Session } from '../session.ts';
import type { SmppLog } from '../log.ts';
import type { SmsIdFormat } from '../protocol/message-ids.ts';
import { HeldMessages } from './held-messages.ts';
import { Reassembler } from '../messages/reassembly.ts';
import { bindCommands, standsInFor } from '../protocol/bind.ts';
import { defaults } from '../options.ts';
import { concatOf } from '../protocol/concat.ts';
import { detach } from '../codec/retained.ts';
import { dlrFromPdu } from '../protocol/receipt.ts';
import { respIdParams, segmentId } from '../protocol/message-ids.ts';
import { respNameFor } from '../codec/commands.ts';
/** Asks the peer to keep the message and retry. */
function throttledStatus(carriedAs: string): ErrorName {
return carriedAs === 'submit_sm' ? 'ESME_RTHROTTLED' : 'ESME_RX_T_APPN';
}
export function refusedSegmentStatus(
carriedAs: string,
refusal: Refusal,
spelling: Concat['spelling'],
): ErrorName {
// A sar_* segment's esm_class is 0x00 and correct: naming it would name the part the peer got right.
if (refusal === 'unplaceable') {
return spelling === 'sar' ? 'ESME_RINVTLVVAL' : 'ESME_RINVESMCLASS';
}
return throttledStatus(carriedAs);
}
const lostReasons: Record<LostGroup['reason'], string> = {
evicted: 'the reassembly buffer filled',
expired: 'no further segment arrived in time',
linkGone: 'the link they arrived on went',
};
export type IncomingRequestsOptions = {
dlrMerger: DlrMerger;
link: LinkLife;
log: SmppLog;
maxOctets?: number | undefined;
maxReassembly?: number | undefined;
onRequest?: OnRequest | undefined;
reassemblyTimeout?: number | undefined;
sendPastDrain: HeldMessagesOptions['sendPastDrain'];
session: Session;
smsIdFormat?: SmsIdFormat | undefined;
systemId?: string | undefined;
};
/** Everything the peer asks of a session: messages, receipts, links and the answers to them. */
export class IncomingRequests {
private readonly dlrMerger: DlrMerger;
private readonly held: HeldMessages;
private readonly link: LinkLife;
private readonly log: SmppLog;
private readonly onRequest: OnRequest | undefined;
private readonly reassembler: Reassembler;
private readonly session: Session;
private readonly smsIdFormat: SmsIdFormat;
private readonly systemId: string;
private refusing = false;
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;
this.log = options.log;
this.onRequest = options.onRequest;
this.reassembler = new Reassembler({
log: options.log,
max: options.maxReassembly ?? defaults.maxReassembly,
maxOctets: options.maxOctets,
onLost: lost => { this.reportLost(lost); },
timeout: options.reassemblyTimeout ?? defaults.reassemblyTimeout,
});
this.session = options.session;
this.smsIdFormat = options.smsIdFormat ?? {};
this.systemId = options.systemId ?? defaults.systemId;
}
async handle(pduObj: PduObject): Promise<void> {
const generation = this.link.generation();
const { onRequest } = this;
// Called unbound, so the application's hook never sees this class as its `this`.
if (onRequest && await onRequest(this.session, pduObj)) return;
// The link it arrived on went while the hook ran, so nothing we answer now correlates.
if (this.link.generation() !== generation) {
this.log.info('session - dropping a request whose link went', { cmdName: pduObj.cmdName });
return;
}
if (!this.session.bindAllows(pduObj.cmdName)) {
this.log.info('session - command the peer\'s bind direction does not carry', {
bindType: this.session.boundAs ?? '',
cmdName: pduObj.cmdName,
});
await this.session.sendReturn(pduObj, 'ESME_RINVBNDSTS');
return;
}
await this.route(pduObj);
}
private async route(pduObj: PduObject): Promise<void> {
switch (pduObj.cmdName) {
case 'data_sm':
case 'deliver_sm':
// A data_sm at the SMSC end is a submission, and a submission is never a report.
await (this.carriedAs(pduObj) === 'submit_sm'
? this.onMessage(pduObj)
: this.onDelivery(pduObj));
break;
case 'enquire_link':
await this.session.sendReturn(pduObj);
break;
case 'submit_sm':
await this.onMessage(pduObj);
break;
case 'unbind':
await this.session.sendReturn(pduObj);
// A peer that has said it is finished will not answer what we still have outstanding.
await this.session.close({ signal: AbortSignal.abort() });
break;
default:
await this.unhandled(pduObj);
}
}
/** Drops the segments of every message that never became whole, and of every one still held. */
clear(): void {
this.refusing = false;
this.held.clear();
this.reassembler.clear();
}
listenerRejected(sms: unknown): void {
this.held.listenerRejected(sms);
}
/** Waits out the messages the application still holds, and says how many it never answered. */
async drain(timeout: number, signal: AbortSignal | undefined): Promise<VoidResult> {
const unanswered = await this.held.idle(timeout, signal);
if (unanswered === 0) return {};
this.log.warn('session - shutting down with messages unanswered', { timeout, unanswered });
return { err: new Error(`Shut down with ${String(unanswered)} message(s) unanswered`) };
}
private async unhandled(pduObj: PduObject): Promise<void> {
if (bindCommands.includes(pduObj.cmdName)) {
this.log.info('session - bind on an already bound session', { cmdName: pduObj.cmdName });
await this.session.sendReturn(pduObj, 'ESME_RALYBND', { system_id: this.systemId });
return;
}
if (!respNameFor(pduObj.cmdName)) {
this.log.verbose('session - ignoring a command SMPP gives no response', { cmdName: pduObj.cmdName });
return;
}
this.log.info('session - no handler for command', { cmdName: pduObj.cmdName });
await this.session.sendReturn(pduObj, 'ESME_RINVCMDID');
}
private carriedAs(pduObj: PduObject): string {
return standsInFor(pduObj.cmdName, this.session.linkEnd);
}
/** SMPP carries a mobile-originated message and a delivery receipt on the same command. */
private async onDelivery(pduObj: PduObject): Promise<void> {
const dlr = dlrFromPdu(pduObj, this.smsIdFormat);
if (!dlr) {
await this.onMessage(pduObj);
return;
}
this.session.emit('dlr', dlr, pduObj);
const merged = this.dlrMerger.collect(dlr);
if (merged) this.session.emit('messageDlr', merged);
await this.session.sendReturn(pduObj);
}
private async refusedAtBound(pduObj: PduObject): Promise<boolean> {
if (this.held.full()) {
if (!this.refusing) {
this.refusing = true;
this.log.warn('session - unanswered messages at their bound, refusing new ones until the application answers', {
messages: this.held.size,
octets: this.held.octetsHeld,
});
}
this.log.verbose('session - unanswered messages at their bound, asking the peer to retry', {
cmdName: pduObj.cmdName,
seqNr: pduObj.seqNr,
});
await this.session.sendReturn(pduObj, throttledStatus(this.carriedAs(pduObj)));
return true;
}
// Half, so a peer keeping its window full does not flip this on every answer.
if (
this.refusing
&& this.held.size <= defaults.maxHeldMessages / 2
&& this.held.octetsHeld <= defaults.maxHeldOctets / 2
) {
this.refusing = false;
this.log.info('session - unanswered messages down to half their bound, accepting again', { messages: this.held.size });
}
return false;
}
/**
* A concatenated message is answered segment by segment as it arrives: a peer that dispatches
* one request at a time never sends the second segment until the first has been answered.
*/
private async onMessage(pduObj: PduObject): Promise<void> {
if (await this.refusedAtBound(pduObj)) return;
const concat = concatOf(pduObj);
if (!concat) {
this.held.offer([detach(pduObj)]);
return;
}
const collected = this.reassembler.collect(pduObj, concat);
if (!collected.kept) {
await this.session.sendReturn(
pduObj,
refusedSegmentStatus(this.carriedAs(pduObj), collected.refusal, concat.spelling),
);
return;
}
await this.session.sendReturn(
pduObj,
'ESME_ROK',
respIdParams(pduObj.cmdName, segmentId(collected.smsId, concat.part - 1, concat.total)),
);
if (collected.whole) this.held.offer(collected.whole, collected.smsId);
}
private reportLost(lost: LostGroup): void {
this.session.emit('sessionError', new Error(
`Gave up ${String(lost.parts)} of ${String(lost.total)} segments of an incomplete concatenated message: ${lostReasons[lost.reason]}`,
));
}
}
+91
View File
@@ -0,0 +1,91 @@
import type { SmppLog } from '../log.ts';
import type { VoidResult } from '../result.ts';
import { IdleWaiters } from './waiting.ts';
export type SendWindowOptions = {
limit: number;
log: SmppLog;
};
type Waiter = (result: VoidResult) => void;
function aborted(): Error {
return new Error('Aborted while waiting for a send window slot');
}
/** Caps how many requests are on the wire at once; anything past the limit waits its turn. */
export class SendWindow {
private readonly idleWaiters = new IdleWaiters();
private readonly limit: number;
private readonly log: SmppLog;
private readonly waiting = new Set<Waiter>();
private inFlight = 0;
constructor(options: SendWindowOptions) {
this.limit = options.limit;
this.log = options.log;
}
/** Resolves once a slot is the caller's, or with the reason it stopped waiting for one. */
acquire(signal: AbortSignal | undefined): Promise<VoidResult> {
if (this.inFlight < this.limit) {
this.inFlight++;
return Promise.resolve({});
}
if (signal?.aborted === true) return Promise.resolve({ err: aborted() });
return this.queue(signal);
}
release(): void {
const next = this.waiting.values().next().value;
if (next) {
this.waiting.delete(next);
next({});
return;
}
this.inFlight--;
if (this.inFlight > 0) return;
this.idleWaiters.settle();
}
/** Everything the caller is still owed: on the wire, plus queued behind a full window. */
unfinished(): number {
return this.inFlight + this.waiting.size;
}
/** Resolves 0 once nothing is left on the wire, or with what still is. */
idle(timeout: number, signal: AbortSignal | undefined): Promise<number> {
return this.idleWaiters.wait(() => this.unfinished(), timeout, signal);
}
/** A waiter leaves the queue as it settles, so release() can only hand a slot to one still in it. */
private queue(signal: AbortSignal | undefined): Promise<VoidResult> {
this.log.verbose('sendWindow - queueing a request behind a full window', {
limit: this.limit,
queued: this.waiting.size + 1,
});
return new Promise<VoidResult>(resolve => {
const settle = (result: VoidResult): void => {
this.waiting.delete(settle);
signal?.removeEventListener('abort', onAbort);
resolve(result);
};
function onAbort(): void {
settle({ err: aborted() });
}
signal?.addEventListener('abort', onAbort, { once: true });
this.waiting.add(settle);
});
}
}
+108
View File
@@ -0,0 +1,108 @@
import type { PduObject } from '../codec/pdu.ts';
import type { SmppLog } from '../log.ts';
import type { Socket } from 'node:net';
import type { VoidResult } from '../result.ts';
import { PduFramer } from '../codec/framer.ts';
import { PduRefusedError } from '../codec/refusal.ts';
import { pduToObj } from '../codec/pdu.ts';
export type PduTransportOptions = {
log: SmppLog;
onClose: () => void;
/** Raw bytes, before framing. */
onData: (chunk: Buffer) => void;
onError: (err: Error) => void;
/** A complete PDU, before it is parsed. */
onFramed: (pdu: Buffer) => void;
onPdu: (pduObj: PduObject) => void;
/** A framed PDU the codec could not read. The stream is still in sync, so the link is not lost. */
onRefused: (refused: PduRefusedError) => void;
/** Nothing further can be read off this stream, whatever the socket does next. */
onUnreadable: (err: Error) => void;
};
/** A socket read as a stream of complete PDUs. A reconnect attaches a new socket in its place. */
export class PduTransport {
private readonly options: PduTransportOptions;
private framer = new PduFramer();
private socket: Socket;
constructor(options: PduTransportOptions, sock: Socket) {
this.options = options;
this.socket = sock;
this.wire(sock);
}
get sock(): Socket {
return this.socket;
}
/** Takes over a freshly opened socket. Half a PDU left on the old one must not prefix this one. */
attach(sock: Socket): void {
// The socket being replaced is already dead, and its three handlers still point here.
this.socket.removeAllListeners();
this.socket = sock;
this.framer = new PduFramer();
this.wire(sock);
}
private wire(sock: Socket): void {
sock.on('data', chunk => { this.read(chunk); });
sock.on('close', () => { this.options.onClose(); });
sock.on('error', err => {
this.options.log.warn('transport - socket error', { message: err.message });
this.options.onError(err);
this.options.onClose();
});
}
write(pdu: Buffer): VoidResult {
if (this.socket.destroyed) return { err: new Error('Socket is closed') };
this.socket.write(pdu);
return {};
}
private read(chunk: Buffer): void {
this.options.onData(chunk);
this.framer.push(chunk);
const framed = this.framer.next();
if (framed.err) {
this.options.log.warn('transport - unusable stream', { message: framed.err.message });
this.options.onUnreadable(framed.err);
return;
}
for (const pdu of framed.pdus) {
this.options.onFramed(pdu);
const parsed = pduToObj(pdu);
if (parsed.err instanceof PduRefusedError) {
this.options.log.warn('transport - refusing a PDU it could not read', {
message: parsed.err.message,
reason: parsed.err.reason,
});
this.options.onRefused(parsed.err);
continue;
}
// The framer applies framingRefusal() first, so only a caller that skips it lands here.
if (parsed.err) {
this.options.log.warn('transport - could not parse an incoming PDU', {
message: parsed.err.message,
});
this.options.onUnreadable(parsed.err);
return;
}
this.options.onPdu(parsed.pduObj);
}
}
}
+47
View File
@@ -0,0 +1,47 @@
/** What is left of a budget, in the shape a wait takes it: 0 waits forever. */
export function leftOf(deadline: number): number {
return deadline === 0 ? 0 : Math.max(1, deadline - Date.now());
}
/** Everything waiting for a count to fall to zero, and how such a wait is cut short. */
export class IdleWaiters {
private readonly waiting: (() => void)[] = [];
/** Wakes everything waiting, whatever the count reads now. */
settle(): void {
for (const resolve of this.waiting.splice(0)) {
resolve();
}
}
/**
* Resolves 0 once nothing is left, or with what still is when the timeout or the signal cuts the
* wait short. A timeout of 0 waits forever.
*/
wait(remaining: () => number, timeout: number, signal: AbortSignal | undefined): Promise<number> {
if (remaining() === 0) return Promise.resolve(0);
if (signal?.aborted === true) return Promise.resolve(remaining());
return new Promise<number>(resolve => {
let timer: NodeJS.Timeout | undefined = undefined;
const done = (): void => {
const index = this.waiting.indexOf(done);
if (timer) clearTimeout(timer);
if (index !== -1) this.waiting.splice(index, 1);
signal?.removeEventListener('abort', done);
resolve(remaining());
};
if (timeout > 0) {
timer = setTimeout(done, timeout);
timer.unref();
}
signal?.addEventListener('abort', done, { once: true });
this.waiting.push(done);
});
}
}