diff --git a/AGENTS.md b/AGENTS.md index b93c9c6..7f13a4a 100644 --- a/AGENTS.md +++ b/AGENTS.md @@ -88,7 +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 one way back up is the `Session` handed to -`IncomingRequests` and `createSms()`, imported as a type only. +`createSms()`, imported as a type only; `IncomingRequests` reaches its session through the +`IncomingDeps` port. **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 13e84fe..ae49491 100644 --- a/src/incoming-requests.ts +++ b/src/incoming-requests.ts @@ -1,18 +1,19 @@ +import type { BindType, LinkEnd } from './session-options.ts'; import type { Concat } from './concat.ts'; -import type { DlrMerger } from './dlr-merger.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 { OnRequest } from './session-options.ts'; +import type { ParamValue } from './defs/types.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'; @@ -43,37 +44,56 @@ 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 = { + 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>; - 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; /** The rejection handler is handed the Sms back as an `unknown`, so its hold is found by identity. */ private readonly emitted = new WeakMap void>(); 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, @@ -82,7 +102,6 @@ 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, @@ -90,8 +109,6 @@ 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; } @@ -99,7 +116,7 @@ export class IncomingRequests { async handle(pduObj: PduObject): Promise { const generation = this.linkGeneration; - if (this.onRequest && await this.onRequest(this.session, pduObj)) return; + 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) { @@ -108,12 +125,12 @@ export class IncomingRequests { return; } - if (!this.session.bindAllows(pduObj.cmdName)) { + if (!this.deps.bindAllows(pduObj.cmdName)) { this.log.info('session - command the peer\'s bind direction does not carry', { - bindType: this.session.boundAs ?? '', + bindType: this.deps.boundAs() ?? '', cmdName: pduObj.cmdName, }); - await this.session.sendReturn(pduObj, 'ESME_RINVBNDSTS'); + await this.deps.answer(pduObj, 'ESME_RINVBNDSTS'); return; } @@ -131,15 +148,14 @@ export class IncomingRequests { : this.onDelivery(pduObj)); break; case 'enquire_link': - await this.session.sendReturn(pduObj); + await this.deps.answer(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() }); + await this.deps.answer(pduObj); + await this.deps.peerUnbound(); break; default: await this.unhandled(pduObj); @@ -175,7 +191,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.session.sendReturn(pduObj, 'ESME_RALYBND', { system_id: this.systemId }); + await this.deps.answer(pduObj, 'ESME_RALYBND', { system_id: this.systemId }); return; } @@ -187,11 +203,11 @@ export class IncomingRequests { } this.log.info('session - no handler for command', { cmdName: pduObj.cmdName }); - await this.session.sendReturn(pduObj, 'ESME_RINVCMDID'); + await this.deps.answer(pduObj, 'ESME_RINVCMDID'); } private carriedAs(pduObj: PduObject): string { - return standsInFor(pduObj.cmdName, this.session.linkEnd); + return standsInFor(pduObj.cmdName, this.deps.linkEnd()); } /** SMPP carries a mobile-originated message and a delivery receipt on the same command. */ @@ -204,13 +220,13 @@ export class IncomingRequests { return; } - this.session.emit('dlr', dlr, pduObj); + this.deps.reportDlr(dlr, pduObj); const merged = this.dlrMerger.collect(dlr); - if (merged) this.session.emit('messageDlr', merged); + if (merged) this.deps.reportMessageDlr(merged); - await this.session.sendReturn(pduObj); + await this.deps.answer(pduObj); } private async refusedAtBound(pduObj: PduObject): Promise { @@ -227,7 +243,7 @@ export class IncomingRequests { cmdName: pduObj.cmdName, seqNr: pduObj.seqNr, }); - await this.session.sendReturn(pduObj, throttledStatus(this.carriedAs(pduObj))); + await this.deps.answer(pduObj, throttledStatus(this.carriedAs(pduObj))); return true; } @@ -263,7 +279,7 @@ export class IncomingRequests { const collected = this.reassembler.collect(pduObj, concat); if (!collected.kept) { - await this.session.sendReturn( + await this.deps.answer( pduObj, refusedSegmentStatus(this.carriedAs(pduObj), collected.refusal, concat.spelling), ); @@ -271,7 +287,7 @@ export class IncomingRequests { return; } - await this.session.sendReturn( + await this.deps.answer( pduObj, 'ESME_ROK', respIdParams(pduObj.cmdName, segmentId(collected.smsId, concat.part - 1, concat.total)), @@ -281,7 +297,7 @@ export class IncomingRequests { } private reportLost(lost: LostGroup): void { - this.session.emit('sessionError', new Error( + this.deps.reportError(new Error( `Gave up ${String(lost.parts)} of ${String(lost.total)} segments of an incomplete concatenated message: ${lostReasons[lost.reason]}`, )); } @@ -295,22 +311,21 @@ export class IncomingRequests { // A turn later, so a listener sending its receipt straight after the response still holds. const release = (): void => { setImmediate(() => { this.held.release(pduObjs); }); }; - const sms = createSms({ + const sms = this.deps.createSms({ answeredAs, from: paramText(first.params.source_addr), message: decodeSegments(pduObjs), pduObjs, - session: this.session, to: paramText(first.params.destination_addr), }, { lostLink: () => this.linkGeneration !== generation, onAnswered: release, // Past the refusal only while a drain is still waiting for this message; an ordinary send after. - send: input => (this.held.has(pduObjs) ? this.sendPastDrain(input) : this.session.send(input)), + send: input => (this.held.has(pduObjs) ? this.deps.sendPastDrain(input) : this.deps.send(input)), }); // A rejection leaves the other listeners running, so only the last one to fail gives the message up. - let working = this.session.listenerCount('sms'); + let working = this.deps.smsListeners(); this.held.hold(pduObjs); this.emitted.set(sms, () => { @@ -320,6 +335,6 @@ export class IncomingRequests { }); // A message nobody took is not work a shutdown can wait for. - if (!this.session.emit('sms', sms)) this.held.release(pduObjs); + if (!this.deps.offerSms(sms)) this.held.release(pduObjs); } } diff --git a/src/session.ts b/src/session.ts index 71370c9..2b878b6 100644 --- a/src/session.ts +++ b/src/session.ts @@ -1,9 +1,10 @@ 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, ReconnectOptions, SendOptions, SessionEvents, SessionOptions } from './session-options.ts'; +import type { BindType, CloseOptions, LinkEnd, OnRequest, ReconnectOptions, SendOptions, 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'; @@ -23,6 +24,7 @@ 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 { @@ -118,14 +120,12 @@ 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.pastDrain(input, {}), - session: this, smsIdFormat: options.smsIdFormat, systemId: options.systemId, }); @@ -245,6 +245,26 @@ export class Session extends EventEmitter { return drained; } + private incomingDeps(onRequest: OnRequest | undefined): IncomingDeps { + return { + 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.pastDrain(input, {}), + smsListeners: () => this.listenerCount('sms'), + }; + } + private transportFor(sock: Socket): PduTransport { return new PduTransport({ log: this.log, diff --git a/test/session-extras.test.ts b/test/session-extras.test.ts index 43c29ae..2f533c0 100644 --- a/test/session-extras.test.ts +++ b/test/session-extras.test.ts @@ -4,6 +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 { MessageState } from '../src/defs/constants.ts'; import type { MessageDlr } from '../src/session.ts'; import type { PduObject, PduObjectInput } from '../src/pdu.ts'; @@ -99,6 +100,26 @@ function abortAfter( }); } +/** A port with no link behind it; the session is only what an Sms carries and answers through. */ +function stubPort(session: Session, port: Partial = {}): IncomingDeps { + return { + 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') }), + sendPastDrain: () => Promise.resolve({ err: new Error('never sent') }), + smsListeners: () => session.listenerCount('sms'), + ...port, + }; +} + function submitPdu(seqNr: number, cmdStatus: ErrorName = 'ESME_ROK'): PduObject { return { cmdId: 0x00000004, @@ -729,14 +750,11 @@ describe('reconnect', () => { const session = new Session({ sock: new net.Socket() }); closeAfter(t, session); - session.boundAs = 'transceiver'; 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, - onRequest: async () => { await delay(10); return false; }, - sendPastDrain: () => Promise.resolve({ err: new Error('never sent') }), - session, }); let messages = 0; @@ -755,6 +773,36 @@ describe('reconnect', () => { assert.equal(messages, 1, 'the harness delivers a message whose link stayed'); }); + test('answers an unbind before asking the session to end, and answers a command outside the bind', async t => { + const session = new Session({ sock: new net.Socket() }); + + closeAfter(t, session); + + const calls: string[] = []; + const incoming = new IncomingRequests({ + deps: stubPort(session, { + answer: (pduObj, status) => { + calls.push(`${pduObj.cmdName} ${status ?? 'ESME_ROK'}`); + + return Promise.resolve({}); + }, + bindAllows: cmdName => cmdName !== 'submit_sm', + peerUnbound: () => { + calls.push('peerUnbound'); + + return Promise.resolve(); + }, + }), + dlrMerger: new DlrMerger({ log: silentLog, max: 10, timeout: 10_000 }), + log: silentLog, + }); + + await incoming.handle(submitPdu(1)); + await incoming.handle({ ...submitPdu(2), cmdId: 0x00000006, cmdName: 'unbind', params: {} }); + + assert.deepEqual(calls, ['submit_sm ESME_RINVBNDSTS', 'unbind ESME_ROK', 'peerUnbound']); + }); + test('does not reconnect after an explicit close', async t => { const smpp = await startServer(t); const { session } = await connect(t, smpp, { reconnect: { maxDelay: 50, minDelay: 10 } }); @@ -1528,14 +1576,12 @@ describe('held message bounds', () => { const session = new Session({ sock: new net.Socket() }); closeAfter(t, session); - session.boundAs = 'transceiver'; 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); } }, - sendPastDrain: () => Promise.resolve({ err: new Error('never sent') }), - session, }); const answers: (ErrorName | undefined)[] = []; const received: Sms[] = []; @@ -1584,13 +1630,11 @@ describe('held message bounds', () => { const session = new Session({ sock: new net.Socket() }); closeAfter(t, session); - session.boundAs = 'transceiver'; const incoming = new IncomingRequests({ + deps: stubPort(session), dlrMerger: new DlrMerger({ log: silentLog, max: 10, timeout: 10_000 }), log: silentLog, - sendPastDrain: () => Promise.resolve({ err: new Error('never sent') }), - session, }); const chunk = Buffer.alloc(64 * 1024); const carried = submitPdu(1); diff --git a/todo.md b/todo.md index 69b4c5a..ceb9900 100644 --- a/todo.md +++ b/todo.md @@ -199,22 +199,10 @@ next work ([decision](docs/decisions.md#internals-and-tests)). ### Locality — next, ahead of everything below; 5–6 today, and the gate is 7 -- [ ] **Give `IncomingRequests` a port instead of the `Session` it drives.** It holds its owner and - calls eight members of it 18 times, including `this.session.close()` on an inbound `unbind` — - a collaborator ending its owner's life. `OutgoingRequests` is the mirror half of the same - boundary and takes no session at all. AGENTS.md names this as the one way back up; the port - removes that exception, and `docs/decisions.md` already states the rule under The session's life: "a collaborator - that has to ask does not own its decision". It is also the missing test seam — inbound routing, - reassembly dispatch, `onRequest` ordering and bind-direction refusal have no unit test because - the class cannot be built without a live socket. Carry the eight members as `IncomingDeps`, - exactly as `sendPastDrain` is carried now. No public surface changes. **Do this before the - store (goal 9), or the back-edge is baked into the store's published interface.** - - [ ] **Route `sms.ts` through its handlers, all of it.** `createSms()` already injects `handlers.send`, and then reaches `sms.session.sendReturn()`, `sms.session.bindAllows()` and `sms.session.acceptsOptionalParams()` anyway — two channels to one collaborator. `Sms.session` - stays public as data the application reads. The cheaper half of the item above, and the one - that shows the shape. + stays public as data the application reads. - [ ] **Give the held-message protocol one name and one home.** `emitSms()` is the unit 8 of 9 readers named and 4 would least want to modify, and every one proposed the same fix. It runs