From fd3b08435aa380f345ce7b8ac9b7a1bdf01243fb Mon Sep 17 00:00:00 2001 From: Lilleman auf Larv Date: Mon, 28 Sep 2026 18:53:20 +0200 Subject: [PATCH] Hand IncomingRequests the session in place of sixteen closures --- AGENTS.md | 9 ++-- MIGRATION.md | 2 +- docs/decisions.md | 17 +++---- src/incoming-requests.ts | 93 +++++++++++++++---------------------- src/session.ts | 47 ++++++------------- src/sms.ts | 20 ++++---- test/session-extras.test.ts | 80 +++++++++++-------------------- todo.md | 25 ++++++++-- 8 files changed, 123 insertions(+), 170 deletions(-) diff --git a/AGENTS.md b/AGENTS.md index 593195a..86e66fb 100644 --- a/AGENTS.md +++ b/AGENTS.md @@ -51,7 +51,7 @@ src/ dlr.ts Delivery receipts: text and TLV parsing, receipt status codes dlr-merger.ts DlrMerger: per-segment receipts counted into one MessageDlr error-from.ts An untyped value as error material: errorFrom() an Error, namedValue() a name - expiring-groups.ts ExpiringGroups: the capped, weighed, expiring store both of those share + expiring-groups.ts ExpiringGroups: the capped, weighed, expiring store DlrMerger, HeldMessages and Reassembler share held-messages.ts HeldMessages: capped, expiring messages the application has not answered, one MessageHold each idle-waiters.ts IdleWaiters: waiting for a count to fall to zero, and what is left of a budget incoming-requests.ts Every request the peer sends: messages, receipts, links, unknown commands @@ -82,14 +82,15 @@ src/ constants.ts consts + constsById, and the SMPP version constants encodings.ts GSM 03.38, LATIN1, UCS2, detection, data_coding resolution errors.ts errors + errorsById (ESME_*) + index.ts defs: every table as one group tlvs.ts TLV definitions, tlvsById, the typed read and input shapes, and reading and writing a TLV stream types.ts Wire types: int8/int16/int32/string/cstring/buffer/arrays ``` 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`, which call back into it, 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 @@ -317,7 +318,7 @@ the file. ### [Internals and tests](docs/decisions.md#internals-and-tests) -- #30 merged under the comprehension floor, and Locality is the next work. +- #30 and #46 merged under the comprehension floor, and Locality is the next work. - A listener that rejects is routed by Node's `captureRejections`, not by hand-dispatching. - The four-line abort dance is copied across `LinkGate`, `IdleWaiters`, `PendingRequests` and `SendWindow` rather than extracted. diff --git a/MIGRATION.md b/MIGRATION.md index 8f47396..0296ede 100644 --- a/MIGRATION.md +++ b/MIGRATION.md @@ -1,6 +1,6 @@ # Migrating from larvitsmpp 0.4.0 -`@larvit/smpp` 0.5.0 succeeds [larvitsmpp](https://www.npmjs.com/package/larvitsmpp) 0.4.0. The +`@larvit/smpp` succeeds [larvitsmpp](https://www.npmjs.com/package/larvitsmpp) 0.4.0. The shape is the same, connect, send, listen for delivery reports, with callbacks replaced by promises. ## API changes diff --git a/docs/decisions.md b/docs/decisions.md index 6c2247e..c7d939f 100644 --- a/docs/decisions.md +++ b/docs/decisions.md @@ -8,8 +8,8 @@ rule and an index of the titles below. - **`Session` is publicly constructible, which is what makes `SessionOptions` and `ReconnectOptions` public too.** Raised twice as a leak; it is not one. The collaborators `session.ts` delegates to - (`Reassembler`, `PendingRequests`, `SendWindow`, `ReconnectLoop`, `LinkTimers`, `LinkGate`, - `DlrMerger`, `PduTransport`, `submitSms`) stay unpublished so they can be reshaped. + (`IncomingRequests`, `OutgoingRequests`, `ReconnectLoop`, `LinkTimers`, `DlrMerger`, + `PduTransport`, `submitSms`) stay unpublished so they can be reshaped. - **`acceptsOptionalParams()` and `bindAllows()` are predicates, not chokepoints.** The library's own senders consult them; `session.send({ tlvs })` is passed through as written, because silently @@ -98,9 +98,9 @@ rule and an index of the titles below. optional-parameter threshold.** That threshold is fixed at 0x34 by the spec, so an implementation that must declare 5.0 throughout can, without moving it. -- **A peer that declared no version is pre-3.4, and `undefined` means no bind yet.** `acceptBind()` - records what the ESME declared and the client's `bind()` records the `sc_interface_version` the - SMSC answered with; a peer that declared nothing is recorded as `undeclaredInterfaceVersion` (0x00) +- **A peer that declared no version is pre-3.4, and `undefined` means no bind yet.** `bound()` + records what the peer declared, the ESME's `interface_version` or the SMSC's + `sc_interface_version`; a peer that declared nothing is recorded as `undeclaredInterfaceVersion` (0x00) and sent no optional parameters, which is how the spec reads an absent `sc_interface_version`. - **`esm_class` decides what a `deliver_sm` is, and the body is read only when it names nothing.** @@ -510,7 +510,7 @@ rule and an index of the titles below. Maintainer's call, 2026-08-31: without the split, an application that opens a replacement client on `close` ends up holding two binds on one account. `teardown()` picks the event by whether the reconnect loop is still live, and `end()` stops that loop before tearing down, so every deliberate - shutdown emits `close`. A retry that opens a socket and then loses it clears `closed` through + shutdown emits `close`. A retry that opens a socket and then loses it resets `lifecycle` through `attach()`, which is why a second drop emits again. - **An answer belongs to the link the message arrived on; a receipt does not.** Maintainer's call, @@ -773,12 +773,13 @@ rule and an index of the titles below. ## Internals and tests -- **#30 merged under the comprehension floor, and Locality is the next work.** Maintainer's call, +- **#30 and #46 merged under the comprehension floor, and Locality is the next work.** Maintainer's call, 2026-09-27. A four-seat scoring run, depth 1, read the project at 6, 6, 7 and 6 (mean 6.25), every seat capped by Locality in the held-message and shutdown code #30 does not touch, where the floor is 7.0. The chunks after #30 lift Locality to 7 before any other work. Serves goal 8's reshapeable internals, which a reader has to understand before reshaping. Valid until a scoring - run reads 7.0 or above. + run reads 7.0 or above. #46, maintainer's call 2026-09-28, merged as a step of that work at 6, 6, 7 + and 7, Locality 5, 5, 6 and 6, up from 6, 6, 7 and 6 and Locality 5, 5, 6 and 5 the same day. - **A listener that rejects is routed by Node's `captureRejections`, not by hand-dispatching.** Both emitters construct with `captureRejections: true` and implement diff --git a/src/incoming-requests.ts b/src/incoming-requests.ts index e3a3d18..cb8bed0 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,35 @@ 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>; + 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 +80,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,14 +88,18 @@ 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; } async handle(pduObj: PduObject): Promise { const generation = this.linkGeneration; + const { onRequest } = this; - if (this.deps.onRequest && await this.deps.onRequest(pduObj)) return; + // Called unbound, so the application's hook never sees this class as its `this`. + if (onRequest && await onRequest(this.session, pduObj)) return; // The link it arrived on went while the hook ran, so nothing we answer now correlates. if (this.linkGeneration !== generation) { @@ -124,12 +108,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 +131,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 +172,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 +184,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 +201,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 +224,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 +260,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 +268,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 +278,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 +290,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..d7e6d85 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 { @@ -116,16 +114,6 @@ export class Session extends EventEmitter { max: defaults.maxDlrMerges, timeout: defaults.dlrMergeTimeout, }); - this.incoming = new IncomingRequests({ - deps: this.incomingDeps(options.onRequest), - dlrMerger: this.dlrMerger, - log: this.log, - maxOctets: options.maxOctets, - maxReassembly: options.maxReassembly, - reassemblyTimeout: options.reassemblyTimeout, - smsIdFormat: options.smsIdFormat, - systemId: options.systemId, - }); this.reconnectLoop = this.loopFor(options.reconnect); this.timers = new LinkTimers({ enquireLinkInterval: options.enquireLinkInterval, @@ -142,6 +130,18 @@ export class Session extends EventEmitter { responseTimeout: options.responseTimeout ?? defaults.responseTimeout, transport: this.transport, }); + this.incoming = new IncomingRequests({ + 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, + }); this.resetTimers(); } @@ -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..93f75c7 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>; @@ -93,9 +89,9 @@ export function createSms(input: SmsInput, handlers: SmsHandlers): Sms { from: input.from, message: input.message, pduObjs: input.pduObjs, - sendDlr: status => sendDlr(sms, handlers, status), + sendDlr: status => sendDlr(sms, input.session, handlers, status), sendResp: options => (input.answeredAs === undefined - ? sendResp(sms, answered, options ?? {}, handlers) + ? sendResp(sms, input.session, answered, options ?? {}, handlers) : answeredOnArrival(options ?? {}, handlers)), session: input.session, get smsId(): string { @@ -132,9 +128,10 @@ function answeredOnArrival( async function sendResp( sms: Sms, + session: Session, answered: { smsId: string }, options: SendRespOptions, - handlers: Pick, + handlers: Pick, ): Promise { const total = sms.pduObjs.length; @@ -153,7 +150,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) => session.sendReturn( pduObj, options.status ?? 'ESME_ROK', respIdParams(pduObj.cmdName, segmentId(answered.smsId, index, total)), @@ -214,10 +211,11 @@ function collectReceipt(sent: Result<{ pduObj: PduObject }>[]): SendDlrResult { async function sendDlr( sms: Sms, - handlers: Pick, + session: Session, + handlers: Pick, status: MessageState = 'DELIVERED', ): Promise { - if (!handlers.bindAllows('deliver_sm')) { + if (!session.bindAllows('deliver_sm')) { return { err: new Error('A transmitter-bound session does not carry deliver_sm'), pduObjs: [], @@ -240,7 +238,7 @@ async function sendDlr( short_message: receiptText(sms, smsId, status), source_addr: sms.to, }, - ...(handlers.acceptsOptionalParams() ? { tlvs: receiptTlvs(smsId, status) } : {}), + ...(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..2cdb5e5 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,14 @@ 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 { + 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 +748,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 +1557,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 +1606,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 +1660,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 +1670,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 +1695,6 @@ describe('sendDlr()', () => { session, to: '46709771337', }, { - acceptsOptionalParams: () => true, - answer: () => Promise.resolve({}), - bindAllows: () => true, lostLink: () => false, onAnswered: () => undefined, send: () => { @@ -3120,26 +3095,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..aeeecc2 100644 --- a/todo.md +++ b/todo.md @@ -203,9 +203,13 @@ A second four-seat run on 2026-09-28, after #35–#40, read 6, 6, 7 and 6 again, and 5. Every seat ranked the session's lifecycle hardest and least wanted to modify it. - [ ] **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. + retires the #30 and #46 decision. The sub-items are what the 2026-09-28 run named, most seats first. +- [ ] **Give the link's liveness one owner.** A third run the same day, after #46, read 6, 6, 7 and + 7, Locality 5, 5, 6 and 6; all four seats ranked `drain()`/`end()`/`teardown()`/`retrying()` + in `session.ts` hardest, because whether the link lives is kept in `Session.lifecycle`, + `ReconnectLoop.halted`, `LinkGate.up`/`returning`, `OutgoingRequests.draining` and + `IncomingRequests.linkGeneration`, held in step by statement order and the comment above + `retrying()`. Four seats. ### Correctness @@ -351,6 +355,21 @@ and 5. Every seat ranked the session's lifecycle hardest and least wanted to mod ### Doc claims this review falsified +- [ ] **Log-cap every interop peer and probe Kannel by protocol, as `interop-tests/AGENTS.md` says + every peer is.** `compose.kannel.yaml` and `compose.smscsim.yaml` carry no `logging:` block, + and Kannel's healthcheck is the bare TCP probe that file warns against. From the prose sweep + of #46. + +- [ ] **Cut what the prose sweep of #46 found restated or misplaced.** `docs/decisions.md` entries + of 30–45 lines carrying pre-fix history (133–160, 217–255, 362–402, 404–440, 442–472, + 593–624, 626–662); AGENTS.md's 14-line shared-fixtures bullet, whose tolerated-copies + reasoning is a decision; the defect table rows MIGRATION.md already carries; README's + `error`-event reason (hard rule 3 owns it) and the Audience bullets restating goal 8 and + Install; the `'use strict'` clause in both MIGRATION.md and CHANGELOG.md; the node-smpp + cross-check in MIGRATION.md; the planned work in `interop-tests/AGENTS.md` (an expected + malformed count per peer) and `benchmarks/README.md`. README persona 3 names a store nothing + ships yet — the maintainer's call, since it is the audience. + - [ ] **Make `LinkGate.isUp()`'s doc true or its state match it.** It says a link attached but not yet bound cannot carry a request, while `up` starts `true`, so the first link and a server session are up before any bind. The gate decision in `docs/decisions.md` makes the same claim