Name files for what they export, and move the session files into session/
Mirror / push (push) Has been cancelled
Test / lint (pull_request) Successful in 21s
Test / test (18) (pull_request) Successful in 29s
Test / test (20) (pull_request) Successful in 29s
Test / test (22) (pull_request) Successful in 30s
Test / test (24) (pull_request) Successful in 29s
Test / test (26) (pull_request) Successful in 34s
Mirror / push (push) Has been cancelled
Test / lint (pull_request) Successful in 21s
Test / test (18) (pull_request) Successful in 29s
Test / test (20) (pull_request) Successful in 29s
Test / test (22) (pull_request) Successful in 30s
Test / test (24) (pull_request) Successful in 29s
Test / test (26) (pull_request) Successful in 34s
This commit is contained in:
@@ -1,13 +1,13 @@
|
||||
import type { LinkLife } from '../link-life.ts';
|
||||
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 { 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';
|
||||
import { IdleWaiters } from './idle-waiters.ts';
|
||||
import { createSms } from './sms.ts';
|
||||
import { retainedOctets } from '../codec/retained-pdu.ts';
|
||||
|
||||
export type HeldMessagesOptions = {
|
||||
link: LinkLife;
|
||||
|
||||
@@ -1,13 +1,13 @@
|
||||
import type { Concat } from '../protocol/concat.ts';
|
||||
import type { DlrMerger } from '../messages/receipt-merge.ts';
|
||||
import type { ErrorName } from '../codec/statuses.ts';
|
||||
import type { DlrMerger } from '../messages/dlr-merger.ts';
|
||||
import type { ErrorName } from '../codec/errors.ts';
|
||||
import type { HeldMessagesOptions } from './held-messages.ts';
|
||||
import type { LinkLife } from '../link-life.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 { Session } from './session.ts';
|
||||
import type { SmppLog } from '../log.ts';
|
||||
import type { SmsIdFormat } from '../protocol/message-ids.ts';
|
||||
import { HeldMessages } from './held-messages.ts';
|
||||
@@ -15,8 +15,8 @@ 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 { detach } from '../codec/retained-pdu.ts';
|
||||
import { dlrFromPdu } from '../protocol/dlr.ts';
|
||||
import { respIdParams, segmentId } from '../protocol/message-ids.ts';
|
||||
import { respNameFor } from '../codec/commands.ts';
|
||||
|
||||
@@ -0,0 +1,188 @@
|
||||
import type { SmppLog } from '../log.ts';
|
||||
import type { VoidResult } from '../result.ts';
|
||||
|
||||
export type LinkLifeOptions = {
|
||||
log: SmppLog;
|
||||
now?: (() => number) | undefined;
|
||||
/** Whether a dropped link is followed by another one until stop(). */
|
||||
reconnects: boolean;
|
||||
/** How long a request may wait for a link. 0 waits for as long as one may still arrive. */
|
||||
timeout: number;
|
||||
};
|
||||
|
||||
/** `binding`: a socket is attached and its bind is not answered yet, so it carries nothing but that bind. */
|
||||
type Phase = 'binding' | 'down' | 'ended' | 'up';
|
||||
|
||||
type Waiter = (result: VoidResult) => void;
|
||||
|
||||
function aborted(): Error {
|
||||
return new Error('Aborted while waiting for a link');
|
||||
}
|
||||
|
||||
function expired(): Error {
|
||||
return new Error('The link did not come back in time');
|
||||
}
|
||||
|
||||
function over(): Error {
|
||||
return new Error('Session is closed');
|
||||
}
|
||||
|
||||
/** Whether the session's link lives, and where a request with no link to go out on waits for the next one. */
|
||||
export class LinkLife {
|
||||
private readonly log: SmppLog;
|
||||
private readonly now: () => number;
|
||||
private readonly reconnects: boolean;
|
||||
private readonly timeout: number;
|
||||
private readonly waiting = new Set<Waiter>();
|
||||
private drops = 0;
|
||||
private phase: Phase = 'up';
|
||||
private stopped = false;
|
||||
|
||||
constructor(options: LinkLifeOptions) {
|
||||
this.log = options.log;
|
||||
this.now = options.now ?? Date.now;
|
||||
this.reconnects = options.reconnects;
|
||||
this.timeout = options.timeout;
|
||||
}
|
||||
|
||||
/** A socket is on the link, bound or not. */
|
||||
isAttached(): boolean {
|
||||
return this.phase === 'binding' || this.phase === 'up';
|
||||
}
|
||||
|
||||
/** Whether a request can go out right now. */
|
||||
isUp(): boolean {
|
||||
return this.phase === 'up';
|
||||
}
|
||||
|
||||
private isOver(): boolean {
|
||||
return this.phase === 'ended';
|
||||
}
|
||||
|
||||
/** The session is shutting down: nothing new is taken, and no link follows this one. */
|
||||
isStopped(): boolean {
|
||||
return this.stopped;
|
||||
}
|
||||
|
||||
/** Whether a link that drops now is followed by another. */
|
||||
retrying(): boolean {
|
||||
return this.reconnects && !this.stopped;
|
||||
}
|
||||
|
||||
/** Not up and not over, with a link to come. */
|
||||
awaitsNextLink(): boolean {
|
||||
return !this.isUp() && !this.isOver() && this.retrying();
|
||||
}
|
||||
|
||||
/** Changes with every drop, so what was read off one link can tell that link is gone. */
|
||||
generation(): number {
|
||||
return this.drops;
|
||||
}
|
||||
|
||||
/** Why no request will ever be admitted, or undefined while one may still get through. */
|
||||
refusal(): Error | undefined {
|
||||
return this.isUp() || this.awaitsNextLink() ? undefined : over();
|
||||
}
|
||||
|
||||
/** One budget for a request, however many links it waits through. */
|
||||
hold(signal: AbortSignal | undefined): () => Promise<VoidResult> {
|
||||
const deadline = this.timeout > 0 ? this.now() + this.timeout : 0;
|
||||
|
||||
return () => this.wait(deadline, signal);
|
||||
}
|
||||
|
||||
/** A socket from the reconnect loop, not yet bound. An ended session stays ended. */
|
||||
attach(): void {
|
||||
if (this.isOver()) return;
|
||||
|
||||
this.phase = 'binding';
|
||||
}
|
||||
|
||||
/** The link is bound: everything held goes out on it. */
|
||||
open(): void {
|
||||
this.phase = 'up';
|
||||
|
||||
if (this.waiting.size > 0) {
|
||||
this.log.verbose('linkLife - sending what was held for a link', { held: this.waiting.size });
|
||||
}
|
||||
|
||||
this.release({});
|
||||
}
|
||||
|
||||
/** The attached link is gone: the event that says so, or undefined when there was none to lose. */
|
||||
drop(): 'close' | 'disconnected' | undefined {
|
||||
if (!this.isAttached()) return undefined;
|
||||
|
||||
this.phase = 'down';
|
||||
this.drops++;
|
||||
|
||||
return this.retrying() ? 'disconnected' : 'close';
|
||||
}
|
||||
|
||||
stop(): void {
|
||||
this.stopped = true;
|
||||
}
|
||||
|
||||
/** The session is over: nothing held will ever go out. False means it already was. */
|
||||
end(): boolean {
|
||||
if (this.isOver()) return false;
|
||||
|
||||
this.phase = 'ended';
|
||||
this.stopped = true;
|
||||
this.release({ err: over() });
|
||||
|
||||
return true;
|
||||
}
|
||||
|
||||
/** Resolves once a link can carry the request, or with the reason none ever will. */
|
||||
private wait(deadline: number, signal: AbortSignal | undefined): Promise<VoidResult> {
|
||||
if (this.isUp()) return Promise.resolve({});
|
||||
|
||||
const refused = this.refusal();
|
||||
|
||||
if (refused) return Promise.resolve({ err: refused });
|
||||
|
||||
if (signal?.aborted === true) return Promise.resolve({ err: aborted() });
|
||||
|
||||
const left = deadline === 0 ? 0 : deadline - this.now();
|
||||
|
||||
if (deadline !== 0 && left <= 0) return Promise.resolve({ err: expired() });
|
||||
|
||||
return this.waitForLink(left, signal);
|
||||
}
|
||||
|
||||
private waitForLink(left: number, signal: AbortSignal | undefined): Promise<VoidResult> {
|
||||
this.log.verbose('linkLife - holding a request until a link is back', { timeout: left });
|
||||
|
||||
return new Promise<VoidResult>(resolve => {
|
||||
let timer: NodeJS.Timeout | undefined = undefined;
|
||||
const settle = (result: VoidResult): void => {
|
||||
if (timer) clearTimeout(timer);
|
||||
|
||||
signal?.removeEventListener('abort', onAbort);
|
||||
this.waiting.delete(settle);
|
||||
resolve(result);
|
||||
};
|
||||
const giveUp = (): void => {
|
||||
this.log.warn('linkLife - no link came back in time', { timeout: left });
|
||||
settle({ err: expired() });
|
||||
};
|
||||
|
||||
function onAbort(): void {
|
||||
settle({ err: aborted() });
|
||||
}
|
||||
|
||||
// Not unref()'d: a held request is awaited with no other handle, so the process would exit unsettled.
|
||||
if (left > 0) timer = setTimeout(giveUp, left);
|
||||
|
||||
signal?.addEventListener('abort', onAbort, { once: true });
|
||||
this.waiting.add(settle);
|
||||
});
|
||||
}
|
||||
|
||||
private release(result: VoidResult): void {
|
||||
for (const settle of [...this.waiting]) {
|
||||
settle(result);
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -1,6 +1,6 @@
|
||||
import type { LinkLife } from '../link-life.ts';
|
||||
import type { LinkLife } from './link-life.ts';
|
||||
import type { PduObject, PduObjectInput } from '../codec/pdu.ts';
|
||||
import type { PduTransport } from './transport.ts';
|
||||
import type { PduTransport } from './pdu-transport.ts';
|
||||
import type { Result, VoidResult } from '../result.ts';
|
||||
import type { SendOptions } from '../options.ts';
|
||||
import type { SmppLog } from '../log.ts';
|
||||
|
||||
@@ -2,7 +2,7 @@ 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 { PduFramer } from '../codec/pdu-framer.ts';
|
||||
import { PduRefusedError } from '../codec/refusal.ts';
|
||||
import { pduToObj } from '../codec/pdu.ts';
|
||||
|
||||
@@ -0,0 +1,139 @@
|
||||
import type { Result, VoidResult } from '../result.ts';
|
||||
import type { SmppLog } from '../log.ts';
|
||||
import type { Socket } from 'node:net';
|
||||
import { defaults } from '../options.ts';
|
||||
|
||||
export type ReconnectLoopOptions = {
|
||||
connect: () => Promise<Result<{ sock: Socket }>>;
|
||||
log: SmppLog;
|
||||
maxDelay?: number | undefined;
|
||||
minDelay?: number | undefined;
|
||||
now?: (() => number) | undefined;
|
||||
/** Brings the owner back up on a freshly opened socket. An err means try again. */
|
||||
onConnected: (sock: Socket) => Promise<VoidResult>;
|
||||
/** Whether the wait between attempts lets the process exit. Default true. */
|
||||
unref?: boolean | undefined;
|
||||
};
|
||||
|
||||
/** Reopens a dropped connection, backing off between attempts until it is told to stop. */
|
||||
export class ReconnectLoop {
|
||||
private readonly maxDelay: number;
|
||||
private readonly minDelay: number;
|
||||
private readonly now: () => number;
|
||||
private readonly options: ReconnectLoopOptions;
|
||||
private attempting = false;
|
||||
private delay: number;
|
||||
private halted = false;
|
||||
private timer: NodeJS.Timeout | undefined;
|
||||
private upAt: number | undefined;
|
||||
|
||||
constructor(options: ReconnectLoopOptions) {
|
||||
this.maxDelay = options.maxDelay ?? defaults.maxDelay;
|
||||
this.minDelay = options.minDelay ?? defaults.minDelay;
|
||||
this.now = options.now ?? Date.now;
|
||||
this.options = options;
|
||||
this.delay = this.minDelay;
|
||||
}
|
||||
|
||||
/** Read through a method: stop() can land while an attempt is awaiting. */
|
||||
private isStopped(): boolean {
|
||||
return this.halted;
|
||||
}
|
||||
|
||||
schedule(): void {
|
||||
if (this.timer || this.attempting || this.isStopped()) return;
|
||||
|
||||
// Coming up is not proof: a stream we cannot read is only found once the link is bound.
|
||||
if (this.upAt !== undefined && this.now() - this.upAt >= this.maxDelay) {
|
||||
this.delay = this.minDelay;
|
||||
}
|
||||
|
||||
this.upAt = undefined;
|
||||
|
||||
const delay = this.delay;
|
||||
|
||||
// Announced when the wait is over rather than when it starts: a cancelled one never happened.
|
||||
this.timer = setTimeout(() => {
|
||||
this.timer = undefined;
|
||||
this.options.log.info('reconnect - retrying', { delay });
|
||||
void this.run();
|
||||
}, delay);
|
||||
|
||||
if (this.options.unref ?? true) this.timer.unref();
|
||||
|
||||
this.delay = Math.min(delay * 2, this.maxDelay);
|
||||
}
|
||||
|
||||
stop(): void {
|
||||
this.halted = true;
|
||||
|
||||
if (this.timer) clearTimeout(this.timer);
|
||||
|
||||
this.timer = undefined;
|
||||
}
|
||||
|
||||
private async run(): Promise<void> {
|
||||
this.attempting = true;
|
||||
|
||||
// connect() and onConnected() are the application's, so a throw from either lands here.
|
||||
const retry = await this.attempt().catch((thrown: unknown) => {
|
||||
const err = thrown instanceof Error ? thrown : new Error(String(thrown));
|
||||
|
||||
this.options.log.error('reconnect - an attempt threw', { message: err.message });
|
||||
|
||||
return true;
|
||||
});
|
||||
|
||||
this.attempting = false;
|
||||
|
||||
if (retry) this.schedule();
|
||||
}
|
||||
|
||||
/** True means the attempt failed and the loop should try again. */
|
||||
private async attempt(): Promise<boolean> {
|
||||
if (this.isStopped()) return false;
|
||||
|
||||
const opened = await this.options.connect();
|
||||
|
||||
if (opened.err) {
|
||||
this.options.log.warn('reconnect - could not open a socket', {
|
||||
message: opened.err.message,
|
||||
});
|
||||
|
||||
return true;
|
||||
}
|
||||
|
||||
if (this.isStopped()) {
|
||||
opened.sock.destroy();
|
||||
|
||||
return false;
|
||||
}
|
||||
|
||||
const up = await this.bringUp(opened.sock);
|
||||
|
||||
if (up.err) {
|
||||
this.options.log.warn('reconnect - could not come back up', { message: up.err.message });
|
||||
|
||||
return true;
|
||||
}
|
||||
|
||||
this.upAt = this.now();
|
||||
|
||||
return false;
|
||||
}
|
||||
|
||||
/** The loop owns the socket until the owner is up on it, so a failed handover must not leak it. */
|
||||
private async bringUp(sock: Socket): Promise<VoidResult> {
|
||||
try {
|
||||
const up = await this.options.onConnected(sock);
|
||||
|
||||
if (up.err) sock.destroy();
|
||||
|
||||
return up;
|
||||
} catch (thrown: unknown) {
|
||||
sock.destroy();
|
||||
|
||||
return { err: thrown instanceof Error ? thrown : new Error(String(thrown)) };
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -1,6 +1,6 @@
|
||||
import type { SmppLog } from '../log.ts';
|
||||
import type { VoidResult } from '../result.ts';
|
||||
import { IdleWaiters } from './waiting.ts';
|
||||
import { IdleWaiters } from './idle-waiters.ts';
|
||||
|
||||
export type SendWindowOptions = {
|
||||
limit: number;
|
||||
|
||||
@@ -0,0 +1,468 @@
|
||||
import type { Dlr } from '../protocol/dlr.ts';
|
||||
import type { ErrorName } from '../codec/errors.ts';
|
||||
import type { MessageDlr } from '../messages/dlr-merger.ts';
|
||||
import type { ParamValue } from '../codec/types.ts';
|
||||
import type { PduObject, PduObjectInput, TlvInputs } from '../codec/pdu.ts';
|
||||
import type { PduRefusedError } from '../codec/refusal.ts';
|
||||
import type { BindType, LinkEnd, SessionBind } from '../protocol/bind.ts';
|
||||
import type { CloseOptions, ReconnectOptions, SendOptions, SessionOptions } from '../options.ts';
|
||||
import type { Result, VoidResult } from '../result.ts';
|
||||
import type { SendSmsOptions, SendSmsResult } from '../messages/submit.ts';
|
||||
import type { SmppLog } from '../log.ts';
|
||||
import type { Sms } from './sms.ts';
|
||||
import type { Socket } from 'node:net';
|
||||
import { DlrMerger } from '../messages/dlr-merger.ts';
|
||||
import { EventEmitter } from 'node:events';
|
||||
import { IncomingRequests } from './incoming-requests.ts';
|
||||
import { LinkLife } from './link-life.ts';
|
||||
import { LinkTimers } from './link-timers.ts';
|
||||
import { OutgoingRequests } from './outgoing-requests.ts';
|
||||
import { PduTransport } from './pdu-transport.ts';
|
||||
import { ReconnectLoop } from './reconnect-loop.ts';
|
||||
import { leftOf } from './idle-waiters.ts';
|
||||
import { errorFrom } from '../result.ts';
|
||||
import { optionalParamsMinVersion } from '../codec/constants.ts';
|
||||
import { bindCarries, checkedBind } from '../protocol/bind.ts';
|
||||
import { defaults } from '../options.ts';
|
||||
import { isResp, objToPdu, pduReturn } from '../codec/pdu.ts';
|
||||
import { refusalAnswer } from '../codec/refusal.ts';
|
||||
import { guardedLog } from '../log.ts';
|
||||
import { submitSms, unsent } from '../messages/submit.ts';
|
||||
import { ConcatReference } from '../protocol/udh.ts';
|
||||
|
||||
export type {
|
||||
CloseOptions,
|
||||
MessageDlr,
|
||||
ReconnectOptions,
|
||||
SendOptions,
|
||||
SendSmsOptions,
|
||||
SendSmsResult,
|
||||
SessionOptions,
|
||||
};
|
||||
export type { BindType };
|
||||
|
||||
export type SessionEvents = {
|
||||
close: [];
|
||||
data: [Buffer];
|
||||
disconnected: [];
|
||||
dlr: [Dlr, PduObject];
|
||||
incomingPdu: [Buffer];
|
||||
incomingPduObj: [PduObject];
|
||||
messageDlr: [MessageDlr];
|
||||
reconnected: [];
|
||||
sessionError: [Error | PduRefusedError];
|
||||
sms: [Sms];
|
||||
};
|
||||
|
||||
/** A listener may return a promise: an `async` one that rejects is routed like one that throws. */
|
||||
type SessionListener<K extends keyof SessionEvents> = (...args: SessionEvents[K]) => unknown;
|
||||
|
||||
export class Session extends EventEmitter<SessionEvents> {
|
||||
declare addListener: <K extends keyof SessionEvents>(event: K, listener: SessionListener<K>) => this;
|
||||
declare off: <K extends keyof SessionEvents>(event: K, listener: SessionListener<K>) => this;
|
||||
declare on: <K extends keyof SessionEvents>(event: K, listener: SessionListener<K>) => this;
|
||||
declare once: <K extends keyof SessionEvents>(event: K, listener: SessionListener<K>) => this;
|
||||
declare prependListener: <K extends keyof SessionEvents>(event: K, listener: SessionListener<K>) => this;
|
||||
declare prependOnceListener: <K extends keyof SessionEvents>(event: K, listener: SessionListener<K>) => this;
|
||||
declare removeListener: <K extends keyof SessionEvents>(event: K, listener: SessionListener<K>) => this;
|
||||
|
||||
readonly log: SmppLog;
|
||||
|
||||
/** Which end of the link this is. `server()` sets it; a hand-wired SMSC must set it too. */
|
||||
linkEnd: LinkEnd = 'esme';
|
||||
userData: unknown = undefined;
|
||||
|
||||
private bind: SessionBind | undefined = undefined;
|
||||
|
||||
private readonly concatReference = new ConcatReference();
|
||||
private readonly dlrMerger: DlrMerger;
|
||||
private readonly incoming: IncomingRequests;
|
||||
private readonly link: LinkLife;
|
||||
private readonly options: SessionOptions;
|
||||
private readonly outgoing: OutgoingRequests;
|
||||
private readonly reconnectLoop: ReconnectLoop | undefined;
|
||||
private readonly timers: LinkTimers;
|
||||
private readonly transport: PduTransport;
|
||||
|
||||
/** A listener that throws is the application's bug; it must not become ours. Hard rule 1. */
|
||||
override emit<K extends keyof SessionEvents>(
|
||||
event: K,
|
||||
...args: K extends keyof SessionEvents ? SessionEvents[K] : never
|
||||
): boolean {
|
||||
try {
|
||||
return super.emit(event, ...args);
|
||||
} catch (thrown: unknown) {
|
||||
const err = errorFrom(thrown);
|
||||
|
||||
this.log.error('session - a listener threw', { event, message: err.message });
|
||||
|
||||
// Guarded against the listener that throws being the one listening for this.
|
||||
if (event !== 'sessionError') this.emit('sessionError', err);
|
||||
|
||||
return false;
|
||||
}
|
||||
}
|
||||
|
||||
/** The same guard for a listener that rejects rather than throws; captureRejections routes here. */
|
||||
override [EventEmitter.captureRejectionSymbol](
|
||||
reason: unknown,
|
||||
...args: [event: keyof SessionEvents, ...rest: unknown[]]
|
||||
): void {
|
||||
const [event, ...rest] = args;
|
||||
const error = errorFrom(reason);
|
||||
|
||||
this.log.error('session - a listener rejected', { event, message: error.message });
|
||||
|
||||
if (event === 'sms') this.incoming.listenerRejected(rest[0]);
|
||||
|
||||
if (event !== 'sessionError') this.emit('sessionError', error);
|
||||
}
|
||||
|
||||
constructor(options: SessionOptions) {
|
||||
super({ captureRejections: true });
|
||||
|
||||
this.log = guardedLog(options.log);
|
||||
this.options = options;
|
||||
this.dlrMerger = new DlrMerger({ log: this.log, max: defaults.maxDlrMerges, timeout: defaults.dlrMergeTimeout });
|
||||
this.reconnectLoop = this.loopFor(options.reconnect);
|
||||
|
||||
const responseTimeout = options.responseTimeout ?? defaults.responseTimeout;
|
||||
|
||||
this.link = new LinkLife({ log: this.log, reconnects: this.reconnectLoop !== undefined, timeout: responseTimeout });
|
||||
this.timers = new LinkTimers({
|
||||
enquireLinkInterval: options.enquireLinkInterval,
|
||||
idleTimeout: options.idleTimeout,
|
||||
log: this.log,
|
||||
onEnquireLink: () => { void this.send({ cmdName: 'enquire_link' }); },
|
||||
// Not close(): a link that went quiet is a drop, and a drop is what reconnect is for.
|
||||
onIdle: () => { this.teardown(); },
|
||||
});
|
||||
this.transport = this.transportFor(options.sock);
|
||||
this.outgoing = new OutgoingRequests({
|
||||
link: this.link,
|
||||
log: this.log,
|
||||
maxOutstanding: options.maxOutstanding ?? defaults.maxOutstanding,
|
||||
responseTimeout,
|
||||
transport: this.transport,
|
||||
});
|
||||
this.incoming = new IncomingRequests({
|
||||
dlrMerger: this.dlrMerger,
|
||||
link: this.link,
|
||||
log: this.log,
|
||||
maxOctets: options.maxOctets,
|
||||
maxReassembly: options.maxReassembly,
|
||||
onRequest: options.onRequest,
|
||||
reassemblyTimeout: options.reassemblyTimeout,
|
||||
sendPastDrain: input => this.outgoing.requestPastDrain(input, {}),
|
||||
session: this,
|
||||
smsIdFormat: options.smsIdFormat,
|
||||
systemId: options.systemId,
|
||||
});
|
||||
|
||||
this.resetTimers();
|
||||
}
|
||||
|
||||
/** Replaced on reconnect, so hold the session rather than this. */
|
||||
get sock(): Socket {
|
||||
return this.transport.sock;
|
||||
}
|
||||
|
||||
/** The role the ESME bound with, whichever end of the link this is. Undefined before any bind. */
|
||||
get boundAs(): BindType | undefined {
|
||||
return this.bind?.as;
|
||||
}
|
||||
|
||||
/** What the peer declared when binding: 0x00 if it declared none, undefined before any bind. */
|
||||
get peerInterfaceVersion(): number | undefined {
|
||||
return this.bind?.peerVersion;
|
||||
}
|
||||
|
||||
/** Records a bind this link accepted or had accepted, until the next one. */
|
||||
bound(bindType: string, declaredVersion: unknown): VoidResult {
|
||||
const checked = checkedBind(bindType, declaredVersion);
|
||||
|
||||
if (!checked.err) this.bind = checked.bind;
|
||||
|
||||
return checked.err ? { err: checked.err } : {};
|
||||
}
|
||||
|
||||
/** Whether this session's bind direction carries a command. Consulted by the library's senders. */
|
||||
bindAllows(cmdName: string): boolean {
|
||||
return bindCarries(this.boundAs, cmdName, this.linkEnd);
|
||||
}
|
||||
|
||||
/** SMPP 3.4 forbids sending optional parameters to a peer that declared an older version. */
|
||||
acceptsOptionalParams(): boolean {
|
||||
return this.peerInterfaceVersion === undefined || this.peerInterfaceVersion >= optionalParamsMinVersion;
|
||||
}
|
||||
|
||||
/** Sends a request and resolves with the peer's response. */
|
||||
send(input: PduObjectInput, options: SendOptions = {}): Promise<Result<{ pduObj: PduObject }>> {
|
||||
return this.outgoing.request(input, options);
|
||||
}
|
||||
|
||||
/** Answers a request the peer sent us. Responses are never waited on. */
|
||||
sendReturn(
|
||||
pdu: PduObject,
|
||||
status: ErrorName = 'ESME_ROK',
|
||||
params: Record<string, ParamValue> = {},
|
||||
tlvs?: TlvInputs,
|
||||
): Promise<VoidResult> {
|
||||
return Promise.resolve(this.answer(pduReturn(pdu, status, params, tlvs), pdu.cmdName, pdu.seqNr));
|
||||
}
|
||||
|
||||
private answer(built: Result<{ buffer: Buffer }>, cmdName: string, seqNr: number): VoidResult {
|
||||
const sent = built.err ? { err: built.err } : this.transport.write(built.buffer);
|
||||
|
||||
// A peer that unbinds and drops the link takes our response with it; that is not a failure.
|
||||
if (sent.err && this.link.isAttached()) {
|
||||
this.log.warn('session - could not answer a request', {
|
||||
cmdName,
|
||||
message: sent.err.message,
|
||||
seqNr,
|
||||
});
|
||||
this.emit('sessionError', sent.err);
|
||||
}
|
||||
|
||||
return sent;
|
||||
}
|
||||
|
||||
async sendSms(sms: SendSmsOptions, options: SendOptions = {}): Promise<SendSmsResult> {
|
||||
if (!this.bindAllows('submit_sm')) {
|
||||
return unsent(new Error('A receiver-bound session does not carry submit_sm'));
|
||||
}
|
||||
|
||||
const sent = await submitSms({
|
||||
log: this.log,
|
||||
reference: this.concatReference.next(),
|
||||
respIdNotation: this.options.smsIdFormat?.submitResp,
|
||||
send: input => this.send(input, options),
|
||||
}, sms);
|
||||
|
||||
if (!sent.err && sms.dlr === true) this.dlrMerger.expect(sent.smsIds);
|
||||
|
||||
return sent;
|
||||
}
|
||||
|
||||
/**
|
||||
* Drains, unbinds politely, then closes. Many SMSCs drop the link instead of answering the
|
||||
* unbind, which is fine. Reports the unbind's own failure ahead of an unfinished drain.
|
||||
*/
|
||||
async unbind(): Promise<VoidResult> {
|
||||
const drained = await this.drain(undefined);
|
||||
const wasOpen = this.link.isAttached();
|
||||
const sent = wasOpen
|
||||
? await this.outgoing.requestOnCurrentLink({ cmdName: 'unbind' })
|
||||
: { err: new Error('Session is closed') };
|
||||
const closedOnUnbind = wasOpen && !this.link.isAttached();
|
||||
|
||||
this.end();
|
||||
|
||||
return sent.err && !closedOnUnbind ? { err: sent.err } : drained;
|
||||
}
|
||||
|
||||
/**
|
||||
* Closes for good: refuses new sends, waits up to `shutdownTimeout` for the requests already sent
|
||||
* and the messages not yet answered, then tears down whatever is left. A session closed this way never reconnects.
|
||||
*/
|
||||
async close(options: CloseOptions = {}): Promise<VoidResult> {
|
||||
const drained = await this.drain(options.signal);
|
||||
|
||||
this.end();
|
||||
|
||||
return drained;
|
||||
}
|
||||
|
||||
private transportFor(sock: Socket): PduTransport {
|
||||
return new PduTransport({
|
||||
log: this.log,
|
||||
onClose: () => { this.onClose(); },
|
||||
onData: chunk => { this.onData(chunk); },
|
||||
onError: err => { this.emit('sessionError', err); },
|
||||
onFramed: pdu => { this.emit('incomingPdu', pdu); },
|
||||
onPdu: pduObj => { this.dispatch(pduObj); },
|
||||
onRefused: refused => { this.refuse(refused); },
|
||||
onUnreadable: err => {
|
||||
this.emit('sessionError', err);
|
||||
this.teardown();
|
||||
},
|
||||
}, sock);
|
||||
}
|
||||
|
||||
private loopFor(reconnect: ReconnectOptions | undefined): ReconnectLoop | undefined {
|
||||
if (!reconnect) return undefined;
|
||||
|
||||
return new ReconnectLoop({
|
||||
connect: reconnect.connect,
|
||||
log: this.log,
|
||||
maxDelay: reconnect.maxDelay,
|
||||
minDelay: reconnect.minDelay,
|
||||
onConnected: sock => this.comeBackUp(sock, reconnect.onConnected),
|
||||
});
|
||||
}
|
||||
|
||||
private async comeBackUp(
|
||||
sock: Socket,
|
||||
bind: (session: Session) => Promise<VoidResult>,
|
||||
): Promise<VoidResult> {
|
||||
this.attach(sock);
|
||||
|
||||
const bound = await bind(this);
|
||||
|
||||
if (bound.err) {
|
||||
this.teardown();
|
||||
|
||||
return { err: bound.err };
|
||||
}
|
||||
|
||||
// close() can land while the rebind is in flight.
|
||||
if (!this.link.retrying()) {
|
||||
this.teardown();
|
||||
|
||||
return { err: new Error('Session closed while it was coming back up') };
|
||||
}
|
||||
|
||||
this.resetTimers();
|
||||
this.link.open();
|
||||
this.log.info('session - reconnected');
|
||||
this.emit('reconnected');
|
||||
|
||||
return {};
|
||||
}
|
||||
|
||||
private attach(sock: Socket): void {
|
||||
this.transport.attach(sock);
|
||||
this.link.attach();
|
||||
}
|
||||
|
||||
/** Stops new sends and waits out the messages we hold and the requests already issued. */
|
||||
private async drain(signal: AbortSignal | undefined): Promise<VoidResult> {
|
||||
this.stop();
|
||||
|
||||
// No bound link, so nothing is on the wire to wait out.
|
||||
if (!this.outgoing.canCarry()) return {};
|
||||
|
||||
const timeout = this.options.shutdownTimeout ?? defaults.shutdownTimeout;
|
||||
const deadline = timeout > 0 ? Date.now() + timeout : 0;
|
||||
// Answering a message can put a receipt on the wire; nothing on the wire produces a message.
|
||||
const messages = await this.incoming.drain(this.answering(timeout), signal);
|
||||
const requests = await this.outgoing.drain(leftOf(deadline), signal);
|
||||
|
||||
// The link went before the drain finished, so an empty window says nothing about the peer.
|
||||
if (!this.outgoing.canCarry()) {
|
||||
return { err: new Error('The session closed before the drain finished') };
|
||||
}
|
||||
|
||||
if (!messages.err) return requests;
|
||||
|
||||
if (!requests.err) return messages;
|
||||
|
||||
return { err: new Error(`${messages.err.message}; ${requests.err.message}`) };
|
||||
}
|
||||
|
||||
/** The application half's budget, which may never be "forever": nothing else ends that wait. */
|
||||
private answering(timeout: number): number {
|
||||
if (timeout > 0) return timeout;
|
||||
|
||||
const responseTimeout = this.options.responseTimeout ?? defaults.responseTimeout;
|
||||
|
||||
return responseTimeout > 0 ? responseTimeout : defaults.responseTimeout;
|
||||
}
|
||||
|
||||
/** The session is over now, drained or not. Nothing brings it back. */
|
||||
private end(): void {
|
||||
this.stop();
|
||||
this.teardown();
|
||||
this.dlrMerger.clear();
|
||||
this.emitClose();
|
||||
}
|
||||
|
||||
/** No new sends, and no link after this one. */
|
||||
private stop(): void {
|
||||
this.link.stop();
|
||||
this.reconnectLoop?.stop();
|
||||
}
|
||||
|
||||
private emitClose(): void {
|
||||
if (!this.link.end()) return;
|
||||
|
||||
this.outgoing.linkLost();
|
||||
this.emit('close');
|
||||
}
|
||||
|
||||
private teardown(): void {
|
||||
const lost = this.link.drop();
|
||||
|
||||
if (!lost) return;
|
||||
|
||||
this.outgoing.linkLost();
|
||||
this.timers.clear();
|
||||
this.incoming.clear();
|
||||
this.sock.destroy();
|
||||
|
||||
// `lost` is read before clear(): a listener it reaches may close() the session, and the drop still reports as disconnected.
|
||||
if (lost === 'disconnected') this.emit('disconnected');
|
||||
else this.emitClose();
|
||||
}
|
||||
|
||||
private onData(chunk: Buffer): void {
|
||||
this.emit('data', chunk);
|
||||
this.resetTimers();
|
||||
}
|
||||
|
||||
private dispatch(pduObj: PduObject): void {
|
||||
if (isResp(pduObj)) {
|
||||
if (!this.outgoing.deliver(pduObj)) {
|
||||
this.log.debug('session - response with no matching request', { seqNr: pduObj.seqNr });
|
||||
}
|
||||
|
||||
return;
|
||||
}
|
||||
|
||||
this.emit('incomingPduObj', pduObj);
|
||||
// Every application hook and listener reached from an incoming PDU funnels through here.
|
||||
void this.incoming.handle(pduObj).catch((thrown: unknown) => {
|
||||
const err = errorFrom(thrown);
|
||||
|
||||
this.log.error('session - a handler threw', {
|
||||
cmdName: pduObj.cmdName,
|
||||
message: err.message,
|
||||
seqNr: pduObj.seqNr,
|
||||
});
|
||||
this.emit('sessionError', err);
|
||||
});
|
||||
}
|
||||
|
||||
/** A PDU the codec refused. Its header parsed, so the peer gets an answer and the link stays. */
|
||||
private refuse(refused: PduRefusedError): void {
|
||||
const { cmdId, cmdName, seqNr } = refused.header;
|
||||
|
||||
this.emit('sessionError', refused);
|
||||
|
||||
// A response carries a sequence number of ours, so writing one back lands in the peer's space.
|
||||
if (isResp(refused.header)) {
|
||||
this.outgoing.settleRefused(seqNr, refused);
|
||||
|
||||
return;
|
||||
}
|
||||
|
||||
this.answer(objToPdu({ ...refusalAnswer(refused), seqNr }), cmdName ?? String(cmdId), seqNr);
|
||||
}
|
||||
|
||||
private resetTimers(): void {
|
||||
if (!this.link.isAttached()) return;
|
||||
|
||||
this.timers.reset();
|
||||
}
|
||||
|
||||
private onClose(): void {
|
||||
if (this.link.retrying()) {
|
||||
this.teardown();
|
||||
this.reconnectLoop?.schedule();
|
||||
|
||||
return;
|
||||
}
|
||||
|
||||
this.end();
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,243 @@
|
||||
import type { ErrorName } from '../codec/errors.ts';
|
||||
import type { MessageState } from '../codec/constants.ts';
|
||||
import type { PduObject, PduObjectInput, TlvInputs } from '../codec/pdu.ts';
|
||||
import type { Result, VoidResult } from '../result.ts';
|
||||
import type { Session } from './session.ts';
|
||||
import { UnansweredError } from '../unanswered-error.ts';
|
||||
import { consts } from '../codec/constants.ts';
|
||||
import { decodeSegments } from '../messages/reassembly.ts';
|
||||
import { messageClassOf } from '../codec/encodings.ts';
|
||||
import { paramText } from '../codec/types.ts';
|
||||
import { receiptCodes, transientStates } from '../protocol/dlr.ts';
|
||||
import { smppDate } from '../message.ts';
|
||||
import { respIdParams, segmentId } from '../protocol/message-ids.ts';
|
||||
import { uuidv7 } from '../protocol/uuid.ts';
|
||||
|
||||
/** `pduObjs` holds what the peer took, so a partial failure names what is already receipted. */
|
||||
export type SendDlrResult = {
|
||||
err?: Error;
|
||||
pduObjs: PduObject[];
|
||||
/** Segments that went out unanswered. The peer may have taken them, so sending again may duplicate. */
|
||||
unanswered: number;
|
||||
};
|
||||
|
||||
export type SendRespOptions = {
|
||||
/** The id the peer correlates a later delivery receipt by. Defaults to a generated UUID v7. */
|
||||
smsId?: string;
|
||||
status?: ErrorName;
|
||||
};
|
||||
|
||||
/**
|
||||
* A received SMS, and the handle for answering it. Multipart messages arrive as one Sms carrying
|
||||
* every segment's PDU.
|
||||
*/
|
||||
export type Sms = {
|
||||
/**
|
||||
* Whether the peer was answered as the message's segments arrived, which is what a concatenated
|
||||
* message needs and a segment count cannot tell you. `sendResp()` then writes nothing.
|
||||
*/
|
||||
answeredOnArrival: boolean;
|
||||
dlr: boolean;
|
||||
/** GSM 03.38 message class 0: shown on arrival and not stored. */
|
||||
flash: boolean;
|
||||
from: string;
|
||||
message: string;
|
||||
pduObjs: PduObject[];
|
||||
/** Sends a delivery report back to the sender. Defaults to DELIVERED. */
|
||||
sendDlr: (status?: MessageState) => Promise<SendDlrResult>;
|
||||
/**
|
||||
* Answers the message, and says the application is done with it. A concatenated message was
|
||||
* answered segment by segment as it arrived, so there it only releases a shutdown's wait and
|
||||
* refuses an `smsId` or a refusing `status`. Part of the protocol, not optional.
|
||||
*/
|
||||
sendResp: (options?: SendRespOptions) => Promise<VoidResult>;
|
||||
session: Session;
|
||||
/** The id the segments were answered with, the id `sendResp()` was given, or a generated UUID v7. */
|
||||
readonly smsId: string;
|
||||
submitTime: Date;
|
||||
to: string;
|
||||
};
|
||||
|
||||
export type SmsInput = {
|
||||
/** The id base the segments were already answered with; absent leaves the answer to `sendResp()`. */
|
||||
answeredAs?: string | undefined;
|
||||
pduObjs: PduObject[];
|
||||
session: Session;
|
||||
};
|
||||
|
||||
export type SmsHandlers = {
|
||||
answered: () => void;
|
||||
lostLink: () => boolean;
|
||||
send: (input: PduObjectInput) => Promise<Result<{ pduObj: PduObject }>>;
|
||||
};
|
||||
|
||||
/** GSM 03.38 section 4 gives class 0 immediate display; every other class is stored somewhere. */
|
||||
const immediateDisplayClass = 0;
|
||||
|
||||
export function createSms(input: SmsInput, handlers: SmsHandlers): Sms {
|
||||
const first = input.pduObjs[0];
|
||||
const registered = first?.params.registered_delivery;
|
||||
const dataCoding = first?.params.data_coding;
|
||||
const answered = { smsId: input.answeredAs ?? uuidv7() };
|
||||
|
||||
const sms: Sms = {
|
||||
answeredOnArrival: input.answeredAs !== undefined,
|
||||
dlr: typeof registered === 'number' && registered !== 0,
|
||||
flash: typeof dataCoding === 'number' && messageClassOf(dataCoding) === immediateDisplayClass,
|
||||
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
|
||||
? sendResp(sms, input.session, answered, options ?? {}, handlers)
|
||||
: answeredOnArrival(options ?? {}, handlers)),
|
||||
session: input.session,
|
||||
get smsId(): string {
|
||||
return answered.smsId;
|
||||
},
|
||||
submitTime: new Date(),
|
||||
to: paramText(first?.params.destination_addr),
|
||||
};
|
||||
|
||||
return 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, 'answered'>,
|
||||
): Promise<VoidResult> {
|
||||
if (options.smsId !== undefined) {
|
||||
return Promise.resolve({
|
||||
err: new Error('This message\'s id was fixed when its first segment arrived; read sms.smsId'),
|
||||
});
|
||||
}
|
||||
|
||||
if (options.status !== undefined && options.status !== 'ESME_ROK') {
|
||||
return Promise.resolve({
|
||||
err: new Error('Its segments were answered as they arrived, so there is nothing left to refuse; refuse a segment from the onRequest option instead'),
|
||||
});
|
||||
}
|
||||
|
||||
handlers.answered();
|
||||
|
||||
return Promise.resolve({});
|
||||
}
|
||||
|
||||
async function sendResp(
|
||||
sms: Sms,
|
||||
session: Session,
|
||||
answered: { smsId: string },
|
||||
options: SendRespOptions,
|
||||
handlers: Pick<SmsHandlers, 'answered' | 'lostLink'>,
|
||||
): Promise<VoidResult> {
|
||||
const total = sms.pduObjs.length;
|
||||
|
||||
if (total === 0) {
|
||||
return { err: new Error('No PDUs to answer') };
|
||||
}
|
||||
|
||||
if (options.smsId === '') {
|
||||
return { err: new Error('smsId must not be empty') };
|
||||
}
|
||||
|
||||
if (options.smsId !== undefined) answered.smsId = options.smsId;
|
||||
|
||||
// A response carries the sequence number it was asked on, which the next link knows nothing about.
|
||||
if (handlers.lostLink()) {
|
||||
return { err: new Error('The link this message arrived on is gone, so nothing would correlate the response') };
|
||||
}
|
||||
|
||||
const results = await Promise.all(sms.pduObjs.map((pduObj, index) => session.sendReturn(
|
||||
pduObj,
|
||||
options.status ?? 'ESME_ROK',
|
||||
respIdParams(pduObj.cmdName, segmentId(answered.smsId, index, total)),
|
||||
)));
|
||||
|
||||
const failure = results.find(result => result.err);
|
||||
|
||||
if (!failure) handlers.answered();
|
||||
|
||||
return failure ?? {};
|
||||
}
|
||||
|
||||
/** The receipt as text, which is all of it a peer below SMPP 3.4 is allowed to be sent. */
|
||||
function receiptText(sms: Sms, smsId: string, status: MessageState): string {
|
||||
const delivered = status === 'DELIVERED';
|
||||
const failed = !delivered && !transientStates.includes(status);
|
||||
|
||||
return [
|
||||
`id:${smsId}`,
|
||||
'sub:001',
|
||||
`dlvrd:${delivered ? '001' : '000'}`,
|
||||
`submit date:${smppDate(sms.submitTime)}`,
|
||||
`done date:${smppDate(new Date())}`,
|
||||
`stat:${receiptCodes[status]}`,
|
||||
`err:${failed ? '001' : '000'}`,
|
||||
'text:',
|
||||
].join(' ');
|
||||
}
|
||||
|
||||
function receiptTlvs(smsId: string, status: MessageState): TlvInputs {
|
||||
return {
|
||||
message_state: { tagValue: consts.MESSAGE_STATE[status] },
|
||||
receipted_message_id: { tagValue: smsId },
|
||||
};
|
||||
}
|
||||
|
||||
function collectReceipt(sent: Result<{ pduObj: PduObject }>[]): SendDlrResult {
|
||||
const pduObjs: PduObject[] = [];
|
||||
let failure: Error | undefined;
|
||||
let unanswered = 0;
|
||||
|
||||
for (const one of sent) {
|
||||
if (one.err) {
|
||||
if (one.err instanceof UnansweredError) unanswered++;
|
||||
|
||||
failure ??= one.err;
|
||||
} else if (one.pduObj.cmdStatus === 'ESME_ROK') {
|
||||
pduObjs.push(one.pduObj);
|
||||
} else {
|
||||
const refusal = one.pduObj.cmdStatus ?? String(one.pduObj.cmdStatusId);
|
||||
|
||||
failure ??= new Error(`deliver_sm refused by the peer: ${refusal}`);
|
||||
}
|
||||
}
|
||||
|
||||
return failure ? { err: failure, pduObjs, unanswered } : { pduObjs, unanswered };
|
||||
}
|
||||
|
||||
async function sendDlr(
|
||||
sms: Sms,
|
||||
session: Session,
|
||||
handlers: Pick<SmsHandlers, 'send'>,
|
||||
status: MessageState = 'DELIVERED',
|
||||
): Promise<SendDlrResult> {
|
||||
if (!session.bindAllows('deliver_sm')) {
|
||||
return {
|
||||
err: new Error('A transmitter-bound session does not carry deliver_sm'),
|
||||
pduObjs: [],
|
||||
unanswered: 0,
|
||||
};
|
||||
}
|
||||
|
||||
const total = sms.pduObjs.length;
|
||||
// Together, not one after a response: a drain waiting for this message must see the whole receipt.
|
||||
const sent = await Promise.all(sms.pduObjs.map((_segment, index) => {
|
||||
const smsId = segmentId(sms.smsId, index, total);
|
||||
|
||||
return handlers.send({
|
||||
cmdName: 'deliver_sm',
|
||||
params: {
|
||||
destination_addr: sms.from,
|
||||
esm_class: transientStates.includes(status)
|
||||
? consts.ESM_CLASS.INTERMEDIATE_DELIVERY
|
||||
: consts.ESM_CLASS.MC_DELIVERY_RECEIPT,
|
||||
short_message: receiptText(sms, smsId, status),
|
||||
source_addr: sms.to,
|
||||
},
|
||||
...(session.acceptsOptionalParams() ? { tlvs: receiptTlvs(smsId, status) } : {}),
|
||||
});
|
||||
}));
|
||||
return collectReceipt(sent);
|
||||
}
|
||||
Reference in New Issue
Block a user