From 33cb24913a219a46de0e8b806c5b614e7edd54c5 Mon Sep 17 00:00:00 2001 From: Lilleman auf Larv Date: Mon, 28 Sep 2026 11:08:39 +0200 Subject: [PATCH] Hand IncomingRequests the session in place of sixteen closures --- AGENTS.md | 4 +- src/incoming-requests.ts | 92 +++++++++++++++---------------------- src/session.ts | 29 ++---------- src/sms.ts | 14 ++---- test/session-extras.test.ts | 82 +++++++++++---------------------- todo.md | 2 - 6 files changed, 76 insertions(+), 147 deletions(-) diff --git a/AGENTS.md b/AGENTS.md index 593195a..9ac578e 100644 --- a/AGENTS.md +++ b/AGENTS.md @@ -88,8 +88,8 @@ src/ Imports point one way: `defs` knows nothing above it but `result.ts`, `pdu` uses `defs`, `session` uses `pdu`, and `client`/`server` use `session`. The ways back up are the `Session` handed to -`createSms()`, and to `OnRequest` and `onConnected` in `session-options.ts`, all imported as a type -only. +`createSms()` and `IncomingRequests`, and to `OnRequest` and `onConnected` in `session-options.ts`, +all imported as a type only. **Parameter order is wire order.** The key order inside `cmds.*.params` is the order the fields are written to and read from the buffer. Never sort those alphabetically — the alphabetical-ordering diff --git a/src/incoming-requests.ts b/src/incoming-requests.ts index e3a3d18..b4df313 100644 --- a/src/incoming-requests.ts +++ b/src/incoming-requests.ts @@ -1,19 +1,18 @@ -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 { DlrMerger } 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 { OnRequest } from './session-options.ts'; import type { PduObject, PduObjectInput } from './pdu.ts'; import type { Result, VoidResult } from './result.ts'; +import type { Session } from './session.ts'; import type { SmppLog } from './log.ts'; -import type { Sms, SmsHandlers, SmsInput } from './sms.ts'; import type { SmsIdFormat } from './sms-id.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 { createSms } from './sms.ts'; import { detach } from './retained-pdu.ts'; import { dlrFromPdu } from './dlr.ts'; import { paramText } from './defs/types.ts'; @@ -44,55 +43,36 @@ const lostReasons: Record = { 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; + onRequest?: OnRequest | undefined; reassemblyTimeout?: number | undefined; + /** Past a drain's refusal, for a receipt the drain is itself waiting for. */ + sendPastDrain: (input: PduObjectInput) => Promise>; + /** The session these requests arrive on, which answers them and emits what they carry. */ + 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 deps: IncomingDeps; private readonly dlrMerger: DlrMerger; private readonly held: HeldMessages; 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; private linkGeneration = 0; private refusing = false; constructor(options: IncomingRequestsOptions) { - this.deps = options.deps; this.dlrMerger = options.dlrMerger; this.held = new HeldMessages({ log: options.log, @@ -101,6 +81,7 @@ export class IncomingRequests { timeout: defaults.heldMessageTimeout, }); this.log = options.log; + this.onRequest = options.onRequest; this.reassembler = new Reassembler({ log: options.log, max: options.maxReassembly ?? defaults.maxReassembly, @@ -108,6 +89,8 @@ 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; } @@ -115,7 +98,7 @@ export class IncomingRequests { async handle(pduObj: PduObject): Promise { const generation = this.linkGeneration; - if (this.deps.onRequest && await this.deps.onRequest(pduObj)) return; + if (this.onRequest && await this.onRequest(this.session, pduObj)) return; // The link it arrived on went while the hook ran, so nothing we answer now correlates. if (this.linkGeneration !== generation) { @@ -124,12 +107,12 @@ export class IncomingRequests { return; } - if (!this.deps.bindAllows(pduObj.cmdName)) { + if (!this.session.bindAllows(pduObj.cmdName)) { this.log.info('session - command the peer\'s bind direction does not carry', { - bindType: this.deps.boundAs() ?? '', + bindType: this.session.boundAs ?? '', cmdName: pduObj.cmdName, }); - await this.deps.answer(pduObj, 'ESME_RINVBNDSTS'); + await this.session.sendReturn(pduObj, 'ESME_RINVBNDSTS'); return; } @@ -147,14 +130,15 @@ export class IncomingRequests { : this.onDelivery(pduObj)); break; case 'enquire_link': - await this.deps.answer(pduObj); + await this.session.sendReturn(pduObj); break; case 'submit_sm': await this.onMessage(pduObj); break; case 'unbind': - await this.deps.answer(pduObj); - await this.deps.peerUnbound(); + 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); @@ -187,7 +171,7 @@ export class IncomingRequests { 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 }); + await this.session.sendReturn(pduObj, 'ESME_RALYBND', { system_id: this.systemId }); return; } @@ -199,11 +183,11 @@ export class IncomingRequests { } this.log.info('session - no handler for command', { cmdName: pduObj.cmdName }); - await this.deps.answer(pduObj, 'ESME_RINVCMDID'); + await this.session.sendReturn(pduObj, 'ESME_RINVCMDID'); } private carriedAs(pduObj: PduObject): string { - return standsInFor(pduObj.cmdName, this.deps.linkEnd()); + return standsInFor(pduObj.cmdName, this.session.linkEnd); } /** SMPP carries a mobile-originated message and a delivery receipt on the same command. */ @@ -216,13 +200,13 @@ export class IncomingRequests { return; } - this.deps.reportDlr(dlr, pduObj); + this.session.emit('dlr', dlr, pduObj); const merged = this.dlrMerger.collect(dlr); - if (merged) this.deps.reportMessageDlr(merged); + if (merged) this.session.emit('messageDlr', merged); - await this.deps.answer(pduObj); + await this.session.sendReturn(pduObj); } private async refusedAtBound(pduObj: PduObject): Promise { @@ -239,7 +223,7 @@ export class IncomingRequests { cmdName: pduObj.cmdName, seqNr: pduObj.seqNr, }); - await this.deps.answer(pduObj, throttledStatus(this.carriedAs(pduObj))); + await this.session.sendReturn(pduObj, throttledStatus(this.carriedAs(pduObj))); return true; } @@ -275,7 +259,7 @@ export class IncomingRequests { const collected = this.reassembler.collect(pduObj, concat); if (!collected.kept) { - await this.deps.answer( + await this.session.sendReturn( pduObj, refusedSegmentStatus(this.carriedAs(pduObj), collected.refusal, concat.spelling), ); @@ -283,7 +267,7 @@ export class IncomingRequests { return; } - await this.deps.answer( + await this.session.sendReturn( pduObj, 'ESME_ROK', respIdParams(pduObj.cmdName, segmentId(collected.smsId, concat.part - 1, concat.total)), @@ -293,7 +277,7 @@ export class IncomingRequests { } private reportLost(lost: LostGroup): void { - this.deps.reportError(new Error( + this.session.emit('sessionError', new Error( `Gave up ${String(lost.parts)} of ${String(lost.total)} segments of an incomplete concatenated message: ${lostReasons[lost.reason]}`, )); } @@ -305,19 +289,17 @@ export class IncomingRequests { const generation = this.linkGeneration; - this.held.offer(pduObjs, this.deps.smsListeners(), hold => this.deps.createSms({ + 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), }, { - 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)), - }), sms => this.deps.offerSms(sms)); + send: input => (hold.isHeld() ? this.sendPastDrain(input) : this.session.send(input)), + }), sms => this.session.emit('sms', sms)); } } diff --git a/src/session.ts b/src/session.ts index 59ee036..248fafe 100644 --- a/src/session.ts +++ b/src/session.ts @@ -1,10 +1,9 @@ import type { ErrorName } from './defs/errors.ts'; -import type { IncomingDeps } from './incoming-requests.ts'; import type { MessageDlr } from './dlr-merger.ts'; import type { ParamValue } from './defs/types.ts'; import type { PduObject, PduObjectInput, TlvInputs } from './pdu.ts'; import type { PduRefusedError } from './pdu-refusal.ts'; -import type { BindType, CloseOptions, LinkEnd, OnRequest, ReconnectOptions, SendOptions, SessionBind, SessionEvents, SessionOptions } from './session-options.ts'; +import type { BindType, CloseOptions, LinkEnd, ReconnectOptions, SendOptions, SessionBind, SessionEvents, SessionOptions } from './session-options.ts'; import type { Result, VoidResult } from './result.ts'; import type { SendSmsOptions, SendSmsResult } from './send-sms.ts'; import type { SmppLog } from './log.ts'; @@ -24,7 +23,6 @@ import { isResp, objToPdu, pduReturn } from './pdu.ts'; import { refusalAnswer } from './pdu-refusal.ts'; import { guardedLog } from './log.ts'; import { submitSms, unsent } from './send-sms.ts'; -import { createSms } from './sms.ts'; import { ConcatReference } from './udh.ts'; export type { @@ -117,12 +115,14 @@ export class Session extends EventEmitter { timeout: defaults.dlrMergeTimeout, }); this.incoming = new IncomingRequests({ - deps: this.incomingDeps(options.onRequest), dlrMerger: this.dlrMerger, 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, }); @@ -257,27 +257,6 @@ export class Session extends EventEmitter { return drained; } - private incomingDeps(onRequest: OnRequest | undefined): IncomingDeps { - return { - acceptsOptionalParams: () => this.acceptsOptionalParams(), - answer: (pduObj, status, params) => this.sendReturn(pduObj, status, params), - bindAllows: cmdName => this.bindAllows(cmdName), - boundAs: () => this.boundAs, - createSms: (input, handlers) => createSms({ ...input, session: this }, handlers), - linkEnd: () => this.linkEnd, - offerSms: sms => this.emit('sms', sms), - onRequest: onRequest && (pduObj => onRequest(this, pduObj)), - // A peer that has said it is finished will not answer what we still have outstanding. - peerUnbound: () => this.close({ signal: AbortSignal.abort() }), - reportDlr: (dlr, pduObj) => { this.emit('dlr', dlr, pduObj); }, - reportError: err => { this.emit('sessionError', err); }, - reportMessageDlr: merged => { this.emit('messageDlr', merged); }, - send: input => this.send(input), - sendPastDrain: input => this.outgoing.requestPastDrain(input, {}), - smsListeners: () => this.listenerCount('sms'), - }; - } - private transportFor(sock: Socket): PduTransport { return new PduTransport({ log: this.log, diff --git a/src/sms.ts b/src/sms.ts index f8526f3..1175580 100644 --- a/src/sms.ts +++ b/src/sms.ts @@ -1,6 +1,5 @@ import type { ErrorName } from './defs/errors.ts'; import type { MessageState } from './defs/constants.ts'; -import type { ParamValue } from './defs/types.ts'; import type { PduObject, PduObjectInput, TlvInputs } from './pdu.ts'; import type { Result, VoidResult } from './result.ts'; import type { Session } from './session.ts'; @@ -69,9 +68,6 @@ export type SmsInput = { /** What the session's incoming side gives a message so it can be answered and accounted for. */ export type SmsHandlers = { - acceptsOptionalParams: () => boolean; - answer: (pduObj: PduObject, status: ErrorName, params: Record) => Promise; - bindAllows: (cmdName: string) => boolean; lostLink: () => boolean; onAnswered: () => void; send: (input: PduObjectInput) => Promise>; @@ -134,7 +130,7 @@ async function sendResp( sms: Sms, answered: { smsId: string }, options: SendRespOptions, - handlers: Pick, + handlers: Pick, ): Promise { const total = sms.pduObjs.length; @@ -153,7 +149,7 @@ async function sendResp( 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) => handlers.answer( + const results = await Promise.all(sms.pduObjs.map((pduObj, index) => sms.session.sendReturn( pduObj, options.status ?? 'ESME_ROK', respIdParams(pduObj.cmdName, segmentId(answered.smsId, index, total)), @@ -214,10 +210,10 @@ function collectReceipt(sent: Result<{ pduObj: PduObject }>[]): SendDlrResult { async function sendDlr( sms: Sms, - handlers: Pick, + handlers: Pick, status: MessageState = 'DELIVERED', ): Promise { - if (!handlers.bindAllows('deliver_sm')) { + if (!sms.session.bindAllows('deliver_sm')) { return { err: new Error('A transmitter-bound session does not carry deliver_sm'), pduObjs: [], @@ -240,7 +236,7 @@ async function sendDlr( short_message: receiptText(sms, smsId, status), source_addr: sms.to, }, - ...(handlers.acceptsOptionalParams() ? { tlvs: receiptTlvs(smsId, status) } : {}), + ...(sms.session.acceptsOptionalParams() ? { tlvs: receiptTlvs(smsId, status) } : {}), }); })); return collectReceipt(sent); diff --git a/test/session-extras.test.ts b/test/session-extras.test.ts index cce32b1..dd190fa 100644 --- a/test/session-extras.test.ts +++ b/test/session-extras.test.ts @@ -4,7 +4,7 @@ import test, { describe } from 'node:test'; import type { Collected, LostGroup } from '../src/reassembly.ts'; import type { Dlr } from '../src/dlr.ts'; import type { ErrorName } from '../src/defs/errors.ts'; -import type { IncomingDeps } from '../src/incoming-requests.ts'; +import type { IncomingRequestsOptions } from '../src/incoming-requests.ts'; import type { MessageHold } from '../src/held-messages.ts'; import type { MessageState } from '../src/defs/constants.ts'; import type { MessageDlr } from '../src/session.ts'; @@ -101,24 +101,16 @@ function abortAfter( }); } -function stubPort(session: Session, port: Partial = {}): IncomingDeps { - return { - acceptsOptionalParams: () => true, - answer: (pduObj, status, params) => session.sendReturn(pduObj, status, params), - bindAllows: () => true, - boundAs: () => 'transceiver', - createSms: (input, handlers) => createSms({ ...input, session }, handlers), - linkEnd: () => 'smsc', - offerSms: sms => session.emit('sms', sms), - peerUnbound: () => Promise.resolve({}), - reportDlr: () => undefined, - reportError: () => undefined, - reportMessageDlr: () => undefined, - send: () => Promise.resolve({ err: new Error('never sent') }), +function incomingOn(session: Session, options: Partial = {}): IncomingRequests { + session.linkEnd = 'smsc'; + + return new IncomingRequests({ + dlrMerger: new DlrMerger({ log: silentLog, max: 10, timeout: 10_000 }), + log: silentLog, sendPastDrain: () => Promise.resolve({ err: new Error('never sent') }), - smsListeners: () => session.listenerCount('sms'), - ...port, - }; + session, + ...options, + }); } function submitPdu(seqNr: number, cmdStatus: ErrorName = 'ESME_ROK'): PduObject { @@ -758,11 +750,7 @@ describe('reconnect', () => { closeAfter(t, session); - const incoming = new IncomingRequests({ - deps: stubPort(session, { onRequest: async () => { await delay(10); return false; } }), - dlrMerger: new DlrMerger({ log: silentLog, max: 10, timeout: 10_000 }), - log: silentLog, - }); + const incoming = incomingOn(session, { onRequest: async () => { await delay(10); return false; } }); let messages = 0; session.on('sms', () => { messages++; }); @@ -1571,11 +1559,7 @@ describe('held message bounds', () => { closeAfter(t, session); const warnings: string[] = []; - const incoming = new IncomingRequests({ - deps: stubPort(session), - dlrMerger: new DlrMerger({ log: silentLog, max: 10, timeout: 10_000 }), - log: { ...silentLog, warn: message => { warnings.push(message); } }, - }); + const incoming = incomingOn(session, { log: { ...silentLog, warn: message => { warnings.push(message); } } }); const answers: (ErrorName | undefined)[] = []; const received: Sms[] = []; @@ -1624,11 +1608,7 @@ describe('held message bounds', () => { closeAfter(t, session); - const incoming = new IncomingRequests({ - deps: stubPort(session), - dlrMerger: new DlrMerger({ log: silentLog, max: 10, timeout: 10_000 }), - log: silentLog, - }); + const incoming = incomingOn(session); const chunk = Buffer.alloc(64 * 1024); const carried = submitPdu(1); let received: Sms | undefined; @@ -1682,6 +1662,9 @@ describe('sendResp()', () => { closeAfter(t, session); let answered = 0; + + session.sendReturn = () => Promise.resolve({ err: new Error('Socket is closed') }); + const sms = createSms({ from: '46701113311', message: 'never answered', @@ -1689,9 +1672,6 @@ describe('sendResp()', () => { session, to: '46709771337', }, { - acceptsOptionalParams: () => true, - answer: () => Promise.resolve({ err: new Error('Socket is closed') }), - bindAllows: () => true, lostLink: () => false, onAnswered: () => { answered++; }, send: () => Promise.resolve({ err: new Error('never sent') }), @@ -1717,9 +1697,6 @@ describe('sendDlr()', () => { session, to: '46709771337', }, { - acceptsOptionalParams: () => true, - answer: () => Promise.resolve({}), - bindAllows: () => true, lostLink: () => false, onAnswered: () => undefined, send: () => { @@ -3120,26 +3097,23 @@ describe('graceful shutdown', () => { closeAfter(t, session); const calls: string[] = []; - const incoming = new IncomingRequests({ - deps: stubPort(session, { - answer: pduObj => { - calls.push(pduObj.cmdName); + const incoming = incomingOn(session); + const close = session.close.bind(session); - return Promise.resolve({}); - }, - peerUnbound: () => { - calls.push('peerUnbound'); + session.sendReturn = pduObj => { + calls.push(pduObj.cmdName); - return Promise.resolve({}); - }, - }), - dlrMerger: new DlrMerger({ log: silentLog, max: 10, timeout: 10_000 }), - log: silentLog, - }); + return Promise.resolve({}); + }; + session.close = options => { + calls.push('close'); + + return close(options); + }; await incoming.handle({ ...submitPdu(1), cmdId: 0x00000006, cmdName: 'unbind', params: {} }); - assert.deepEqual(calls, ['unbind', 'peerUnbound']); + assert.deepEqual(calls, ['unbind', 'close']); }); }); diff --git a/todo.md b/todo.md index 2ddacdc..3a969f8 100644 --- a/todo.md +++ b/todo.md @@ -204,8 +204,6 @@ and 5. Every seat ranked the session's lifecycle hardest and least wanted to mod - [ ] **Lift Locality to 7, and confirm it with a scoring run.** A run reading 7.0 or above also retires the #30 decision. The sub-items are what the 2026-09-28 run named, most seats first. -- [ ] **Shrink the `IncomingDeps` closure bag.** 16 lambdas, six of them repeated in - `SmsHandlers`, which makes every inbound call path indirect. Two seats. ### Correctness