Hand IncomingRequests the session in place of sixteen closures
Mirror / push (push) Successful in 6s
Test / lint (pull_request) Successful in 24s
Test / test (18) (pull_request) Successful in 32s
Test / test (20) (pull_request) Successful in 32s
Test / test (22) (pull_request) Successful in 32s
Test / test (24) (pull_request) Successful in 32s
Test / test (26) (pull_request) Successful in 31s

This commit is contained in:
2026-09-28 11:08:39 +02:00
parent 0bea36c069
commit 33cb24913a
6 changed files with 76 additions and 147 deletions
+2 -2
View File
@@ -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
+37 -55
View File
@@ -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<LostGroup['reason'], string> = {
linkGone: 'the link they arrived on went',
};
type Send = (input: PduObjectInput) => Promise<Result<{ pduObj: PduObject }>>;
/** 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<string, ParamValue>) => Promise<VoidResult>;
bindAllows: (cmdName: string) => boolean;
boundAs: () => BindType | undefined;
createSms: (input: Omit<SmsInput, 'session'>, 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> | boolean) | undefined;
peerUnbound: () => Promise<VoidResult>;
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<Result<{ pduObj: PduObject }>>;
/** 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<void> {
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<void> {
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<boolean> {
@@ -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));
}
}
+4 -25
View File
@@ -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<SessionEvents> {
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<SessionEvents> {
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,
+5 -9
View File
@@ -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<string, ParamValue>) => Promise<VoidResult>;
bindAllows: (cmdName: string) => boolean;
lostLink: () => boolean;
onAnswered: () => void;
send: (input: PduObjectInput) => Promise<Result<{ pduObj: PduObject }>>;
@@ -134,7 +130,7 @@ async function sendResp(
sms: Sms,
answered: { smsId: string },
options: SendRespOptions,
handlers: Pick<SmsHandlers, 'answer' | 'lostLink' | 'onAnswered'>,
handlers: Pick<SmsHandlers, 'lostLink' | 'onAnswered'>,
): Promise<VoidResult> {
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<SmsHandlers, 'acceptsOptionalParams' | 'bindAllows' | 'send'>,
handlers: Pick<SmsHandlers, 'send'>,
status: MessageState = 'DELIVERED',
): Promise<SendDlrResult> {
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);
+28 -54
View File
@@ -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> = {}): 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<IncomingRequestsOptions> = {}): 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']);
});
});
-2
View File
@@ -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