import type { BindType, LinkEnd } from './session-options.ts'; import type { Concat } from './concat.ts'; import type { Dlr } from './dlr.ts'; import type { DlrMerger, MessageDlr } from './dlr-merger.ts'; import type { ErrorName } from './defs/errors.ts'; import type { LostGroup, Refusal } from './reassembly.ts'; import type { ParamValue } from './defs/types.ts'; import type { PduObject, PduObjectInput } from './pdu.ts'; import type { Result, VoidResult } from './result.ts'; import type { SmppLog } from './log.ts'; import type { Sms, SmsHandlers, SmsInput } from './sms.ts'; import type { SmsIdFormat } from './sms-id.ts'; import type { MessageHold } from './held-messages.ts'; import { HeldMessages } from './held-messages.ts'; import { Reassembler, decodeSegments } from './reassembly.ts'; import { bindCommands, defaults, standsInFor } from './session-options.ts'; import { concatOf } from './concat.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'; /** 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 = { evicted: 'the reassembly buffer filled', expired: 'no further segment arrived in time', linkGone: 'the link they arrived on went', }; type Send = (input: PduObjectInput) => Promise>; /** What the incoming side asks of the session it serves; the session decides how. */ export type IncomingDeps = { acceptsOptionalParams: () => boolean; answer: (pduObj: PduObject, status?: ErrorName, params?: Record) => Promise; bindAllows: (cmdName: string) => boolean; boundAs: () => BindType | undefined; createSms: (input: Omit, handlers: SmsHandlers) => Sms; linkEnd: () => LinkEnd; /** Answers whether any listener took the message. */ offerSms: (sms: Sms) => boolean; /** The application's first refusal, answering true where it took the request itself. */ onRequest?: ((pduObj: PduObject) => Promise | boolean) | undefined; peerUnbound: () => Promise; reportDlr: (dlr: Dlr, pduObj: PduObject) => void; reportError: (err: Error) => void; reportMessageDlr: (merged: MessageDlr) => void; send: Send; /** Past a drain's refusal, for a receipt the drain is itself waiting for. */ sendPastDrain: Send; smsListeners: () => number; }; export type IncomingRequestsOptions = { deps: IncomingDeps; dlrMerger: DlrMerger; log: SmppLog; maxOctets?: number | undefined; maxReassembly?: number | undefined; reassemblyTimeout?: number | undefined; 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 deps: IncomingDeps; private readonly dlrMerger: DlrMerger; private readonly held: HeldMessages; /** The rejection handler is handed the Sms back as an `unknown`, so its hold is found by identity. */ private readonly holds = new WeakMap(); private readonly log: SmppLog; private readonly reassembler: Reassembler; private readonly smsIdFormat: SmsIdFormat; private readonly systemId: string; private linkGeneration = 0; private refusing = false; constructor(options: IncomingRequestsOptions) { this.deps = options.deps; this.dlrMerger = options.dlrMerger; this.held = new HeldMessages({ log: options.log, max: defaults.maxHeldMessages, maxOctets: defaults.maxHeldOctets, timeout: defaults.heldMessageTimeout, }); this.log = options.log; 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.smsIdFormat = options.smsIdFormat ?? {}; this.systemId = options.systemId ?? defaults.systemId; } async handle(pduObj: PduObject): Promise { const generation = this.linkGeneration; if (this.deps.onRequest && await this.deps.onRequest(pduObj)) return; // The link it arrived on went while the hook ran, so nothing we answer now correlates. if (this.linkGeneration !== generation) { this.log.info('session - dropping a request whose link went', { cmdName: pduObj.cmdName }); return; } if (!this.deps.bindAllows(pduObj.cmdName)) { this.log.info('session - command the peer\'s bind direction does not carry', { bindType: this.deps.boundAs() ?? '', cmdName: pduObj.cmdName, }); await this.deps.answer(pduObj, 'ESME_RINVBNDSTS'); return; } await this.route(pduObj); } private async route(pduObj: PduObject): Promise { 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.deps.answer(pduObj); break; case 'submit_sm': await this.onMessage(pduObj); break; case 'unbind': await this.deps.answer(pduObj); await this.deps.peerUnbound(); 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.linkGeneration++; this.refusing = false; this.held.clear(); this.reassembler.clear(); } /** One `sms` listener gave up on a message; the last one to do so is what releases the hold. */ listenerRejected(sms: unknown): void { if (typeof sms !== 'object' || sms === null) return; this.holds.get(sms)?.listenerGaveUp(); } /** Waits out the messages the application still holds, and says how many it never answered. */ async drain(timeout: number, signal: AbortSignal | undefined): Promise { 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 { if (bindCommands.includes(pduObj.cmdName)) { this.log.info('session - bind on an already bound session', { cmdName: pduObj.cmdName }); await this.deps.answer(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.deps.answer(pduObj, 'ESME_RINVCMDID'); } private carriedAs(pduObj: PduObject): string { return standsInFor(pduObj.cmdName, this.deps.linkEnd()); } /** SMPP carries a mobile-originated message and a delivery receipt on the same command. */ private async onDelivery(pduObj: PduObject): Promise { const dlr = dlrFromPdu(pduObj, this.smsIdFormat); if (!dlr) { await this.onMessage(pduObj); return; } this.deps.reportDlr(dlr, pduObj); const merged = this.dlrMerger.collect(dlr); if (merged) this.deps.reportMessageDlr(merged); await this.deps.answer(pduObj); } private async refusedAtBound(pduObj: PduObject): Promise { 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.deps.answer(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 { if (await this.refusedAtBound(pduObj)) return; const concat = concatOf(pduObj); if (!concat) { this.emitSms([detach(pduObj)]); return; } const collected = this.reassembler.collect(pduObj, concat); if (!collected.kept) { await this.deps.answer( pduObj, refusedSegmentStatus(this.carriedAs(pduObj), collected.refusal, concat.spelling), ); return; } await this.deps.answer( pduObj, 'ESME_ROK', respIdParams(pduObj.cmdName, segmentId(collected.smsId, concat.part - 1, concat.total)), ); if (collected.whole) this.emitSms(collected.whole, collected.smsId); } private reportLost(lost: LostGroup): void { this.deps.reportError(new Error( `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.linkGeneration; const hold = this.held.hold(pduObjs, this.deps.smsListeners()); const sms = this.deps.createSms({ answeredAs, from: paramText(first.params.source_addr), message: decodeSegments(pduObjs), pduObjs, to: paramText(first.params.destination_addr), }, { acceptsOptionalParams: () => this.deps.acceptsOptionalParams(), answer: (pduObj, status, params) => this.deps.answer(pduObj, status, params), bindAllows: cmdName => this.deps.bindAllows(cmdName), lostLink: () => this.linkGeneration !== generation, onAnswered: () => { hold.answered(); }, send: input => (hold.isHeld() ? this.deps.sendPastDrain(input) : this.deps.send(input)), }); this.holds.set(sms, hold); if (!this.deps.offerSms(sms)) hold.release(); } }