Hand IncomingRequests the session in place of sixteen closures #46

Merged
lilleman merged 1 commits from incoming-deps into main 2026-09-28 18:56:55 +02:00
8 changed files with 123 additions and 170 deletions
Showing only changes of commit fd3b08435a - Show all commits
+5 -4
View File
@@ -51,7 +51,7 @@ src/
dlr.ts Delivery receipts: text and TLV parsing, receipt status codes dlr.ts Delivery receipts: text and TLV parsing, receipt status codes
dlr-merger.ts DlrMerger: per-segment receipts counted into one MessageDlr 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 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 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 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 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 constants.ts consts + constsById, and the SMPP version constants
encodings.ts GSM 03.38, LATIN1, UCS2, detection, data_coding resolution encodings.ts GSM 03.38, LATIN1, UCS2, detection, data_coding resolution
errors.ts errors + errorsById (ESME_*) 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 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 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` 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 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 `createSms()` and `IncomingRequests`, which call back into it, and to `OnRequest` and `onConnected`
only. 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 **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 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) ### [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. - 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 - The four-line abort dance is copied across `LinkGate`, `IdleWaiters`, `PendingRequests` and
`SendWindow` rather than extracted. `SendWindow` rather than extracted.
+1 -1
View File
@@ -1,6 +1,6 @@
# Migrating from larvitsmpp 0.4.0 # 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. shape is the same, connect, send, listen for delivery reports, with callbacks replaced by promises.
## API changes ## API changes
+9 -8
View File
@@ -8,8 +8,8 @@ rule and an index of the titles below.
- **`Session` is publicly constructible, which is what makes `SessionOptions` and `ReconnectOptions` - **`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 public too.** Raised twice as a leak; it is not one. The collaborators `session.ts` delegates to
(`Reassembler`, `PendingRequests`, `SendWindow`, `ReconnectLoop`, `LinkTimers`, `LinkGate`, (`IncomingRequests`, `OutgoingRequests`, `ReconnectLoop`, `LinkTimers`, `DlrMerger`,
`DlrMerger`, `PduTransport`, `submitSms`) stay unpublished so they can be reshaped. `PduTransport`, `submitSms`) stay unpublished so they can be reshaped.
- **`acceptsOptionalParams()` and `bindAllows()` are predicates, not chokepoints.** The library's own - **`acceptsOptionalParams()` and `bindAllows()` are predicates, not chokepoints.** The library's own
senders consult them; `session.send({ tlvs })` is passed through as written, because silently 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 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. 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()` - **A peer that declared no version is pre-3.4, and `undefined` means no bind yet.** `bound()`
records what the ESME declared and the client's `bind()` records the `sc_interface_version` the records what the peer declared, the ESME's `interface_version` or the SMSC's
SMSC answered with; a peer that declared nothing is recorded as `undeclaredInterfaceVersion` (0x00) `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`. 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.** - **`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 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 `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 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. `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, - **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 ## 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 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 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 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 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 - **A listener that rejects is routed by Node's `captureRejections`, not by hand-dispatching.** Both
emitters construct with `captureRejections: true` and implement emitters construct with `captureRejections: true` and implement
+38 -55
View File
@@ -1,19 +1,18 @@
import type { BindType, LinkEnd } from './session-options.ts';
import type { Concat } from './concat.ts'; import type { Concat } from './concat.ts';
import type { Dlr } from './dlr.ts'; import type { DlrMerger } from './dlr-merger.ts';
import type { DlrMerger, MessageDlr } from './dlr-merger.ts';
import type { ErrorName } from './defs/errors.ts'; import type { ErrorName } from './defs/errors.ts';
import type { LostGroup, Refusal } from './reassembly.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 { PduObject, PduObjectInput } from './pdu.ts';
import type { Result, VoidResult } from './result.ts'; import type { Result, VoidResult } from './result.ts';
import type { Session } from './session.ts';
import type { SmppLog } from './log.ts'; import type { SmppLog } from './log.ts';
import type { Sms, SmsHandlers, SmsInput } from './sms.ts';
import type { SmsIdFormat } from './sms-id.ts'; import type { SmsIdFormat } from './sms-id.ts';
import { HeldMessages } from './held-messages.ts'; import { HeldMessages } from './held-messages.ts';
import { Reassembler, decodeSegments } from './reassembly.ts'; import { Reassembler, decodeSegments } from './reassembly.ts';
import { bindCommands, defaults, standsInFor } from './session-options.ts'; import { bindCommands, defaults, standsInFor } from './session-options.ts';
import { concatOf } from './concat.ts'; import { concatOf } from './concat.ts';
import { createSms } from './sms.ts';
import { detach } from './retained-pdu.ts'; import { detach } from './retained-pdu.ts';
import { dlrFromPdu } from './dlr.ts'; import { dlrFromPdu } from './dlr.ts';
import { paramText } from './defs/types.ts'; import { paramText } from './defs/types.ts';
@@ -44,55 +43,35 @@ const lostReasons: Record<LostGroup['reason'], string> = {
linkGone: 'the link they arrived on went', 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 = { export type IncomingRequestsOptions = {
deps: IncomingDeps;
dlrMerger: DlrMerger; dlrMerger: DlrMerger;
log: SmppLog; log: SmppLog;
maxOctets?: number | undefined; maxOctets?: number | undefined;
maxReassembly?: number | undefined; maxReassembly?: number | undefined;
onRequest?: OnRequest | undefined;
reassemblyTimeout?: number | 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 }>>;
session: Session;
smsIdFormat?: SmsIdFormat | undefined; smsIdFormat?: SmsIdFormat | undefined;
systemId?: string | undefined; systemId?: string | undefined;
}; };
/** Everything the peer asks of a session: messages, receipts, links and the answers to them. */ /** Everything the peer asks of a session: messages, receipts, links and the answers to them. */
export class IncomingRequests { export class IncomingRequests {
private readonly deps: IncomingDeps;
private readonly dlrMerger: DlrMerger; private readonly dlrMerger: DlrMerger;
private readonly held: HeldMessages; private readonly held: HeldMessages;
private readonly log: SmppLog; private readonly log: SmppLog;
private readonly onRequest: OnRequest | undefined;
private readonly reassembler: Reassembler; private readonly reassembler: Reassembler;
private readonly sendPastDrain: IncomingRequestsOptions['sendPastDrain'];
private readonly session: Session;
private readonly smsIdFormat: SmsIdFormat; private readonly smsIdFormat: SmsIdFormat;
private readonly systemId: string; private readonly systemId: string;
private linkGeneration = 0; private linkGeneration = 0;
private refusing = false; private refusing = false;
constructor(options: IncomingRequestsOptions) { constructor(options: IncomingRequestsOptions) {
this.deps = options.deps;
this.dlrMerger = options.dlrMerger; this.dlrMerger = options.dlrMerger;
this.held = new HeldMessages({ this.held = new HeldMessages({
log: options.log, log: options.log,
@@ -101,6 +80,7 @@ export class IncomingRequests {
timeout: defaults.heldMessageTimeout, timeout: defaults.heldMessageTimeout,
}); });
this.log = options.log; this.log = options.log;
this.onRequest = options.onRequest;
this.reassembler = new Reassembler({ this.reassembler = new Reassembler({
log: options.log, log: options.log,
max: options.maxReassembly ?? defaults.maxReassembly, max: options.maxReassembly ?? defaults.maxReassembly,
@@ -108,14 +88,18 @@ export class IncomingRequests {
onLost: lost => { this.reportLost(lost); }, onLost: lost => { this.reportLost(lost); },
timeout: options.reassemblyTimeout ?? defaults.reassemblyTimeout, timeout: options.reassemblyTimeout ?? defaults.reassemblyTimeout,
}); });
this.sendPastDrain = options.sendPastDrain;
this.session = options.session;
this.smsIdFormat = options.smsIdFormat ?? {}; this.smsIdFormat = options.smsIdFormat ?? {};
this.systemId = options.systemId ?? defaults.systemId; this.systemId = options.systemId ?? defaults.systemId;
} }
async handle(pduObj: PduObject): Promise<void> { async handle(pduObj: PduObject): Promise<void> {
const generation = this.linkGeneration; 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. // The link it arrived on went while the hook ran, so nothing we answer now correlates.
if (this.linkGeneration !== generation) { if (this.linkGeneration !== generation) {
@@ -124,12 +108,12 @@ export class IncomingRequests {
return; 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', { this.log.info('session - command the peer\'s bind direction does not carry', {
bindType: this.deps.boundAs() ?? '', bindType: this.session.boundAs ?? '',
cmdName: pduObj.cmdName, cmdName: pduObj.cmdName,
}); });
await this.deps.answer(pduObj, 'ESME_RINVBNDSTS'); await this.session.sendReturn(pduObj, 'ESME_RINVBNDSTS');
return; return;
} }
@@ -147,14 +131,15 @@ export class IncomingRequests {
: this.onDelivery(pduObj)); : this.onDelivery(pduObj));
break; break;
case 'enquire_link': case 'enquire_link':
await this.deps.answer(pduObj); await this.session.sendReturn(pduObj);
break; break;
case 'submit_sm': case 'submit_sm':
await this.onMessage(pduObj); await this.onMessage(pduObj);
break; break;
case 'unbind': case 'unbind':
await this.deps.answer(pduObj); await this.session.sendReturn(pduObj);
await this.deps.peerUnbound(); // A peer that has said it is finished will not answer what we still have outstanding.
await this.session.close({ signal: AbortSignal.abort() });
break; break;
default: default:
await this.unhandled(pduObj); await this.unhandled(pduObj);
@@ -187,7 +172,7 @@ export class IncomingRequests {
private async unhandled(pduObj: PduObject): Promise<void> { private async unhandled(pduObj: PduObject): Promise<void> {
if (bindCommands.includes(pduObj.cmdName)) { if (bindCommands.includes(pduObj.cmdName)) {
this.log.info('session - bind on an already bound session', { cmdName: 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; return;
} }
@@ -199,11 +184,11 @@ export class IncomingRequests {
} }
this.log.info('session - no handler for command', { cmdName: pduObj.cmdName }); 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 { 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. */ /** SMPP carries a mobile-originated message and a delivery receipt on the same command. */
@@ -216,13 +201,13 @@ export class IncomingRequests {
return; return;
} }
this.deps.reportDlr(dlr, pduObj); this.session.emit('dlr', dlr, pduObj);
const merged = this.dlrMerger.collect(dlr); 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> { private async refusedAtBound(pduObj: PduObject): Promise<boolean> {
@@ -239,7 +224,7 @@ export class IncomingRequests {
cmdName: pduObj.cmdName, cmdName: pduObj.cmdName,
seqNr: pduObj.seqNr, seqNr: pduObj.seqNr,
}); });
await this.deps.answer(pduObj, throttledStatus(this.carriedAs(pduObj))); await this.session.sendReturn(pduObj, throttledStatus(this.carriedAs(pduObj)));
return true; return true;
} }
@@ -275,7 +260,7 @@ export class IncomingRequests {
const collected = this.reassembler.collect(pduObj, concat); const collected = this.reassembler.collect(pduObj, concat);
if (!collected.kept) { if (!collected.kept) {
await this.deps.answer( await this.session.sendReturn(
pduObj, pduObj,
refusedSegmentStatus(this.carriedAs(pduObj), collected.refusal, concat.spelling), refusedSegmentStatus(this.carriedAs(pduObj), collected.refusal, concat.spelling),
); );
@@ -283,7 +268,7 @@ export class IncomingRequests {
return; return;
} }
await this.deps.answer( await this.session.sendReturn(
pduObj, pduObj,
'ESME_ROK', 'ESME_ROK',
respIdParams(pduObj.cmdName, segmentId(collected.smsId, concat.part - 1, concat.total)), respIdParams(pduObj.cmdName, segmentId(collected.smsId, concat.part - 1, concat.total)),
@@ -293,7 +278,7 @@ export class IncomingRequests {
} }
private reportLost(lost: LostGroup): void { 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]}`, `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; 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, answeredAs,
from: paramText(first.params.source_addr), from: paramText(first.params.source_addr),
message: decodeSegments(pduObjs), message: decodeSegments(pduObjs),
pduObjs, pduObjs,
session: this.session,
to: paramText(first.params.destination_addr), 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, lostLink: () => this.linkGeneration !== generation,
onAnswered: () => { hold.answered(); }, onAnswered: () => { hold.answered(); },
send: input => (hold.isHeld() ? this.deps.sendPastDrain(input) : this.deps.send(input)), send: input => (hold.isHeld() ? this.sendPastDrain(input) : this.session.send(input)),
}), sms => this.deps.offerSms(sms)); }), sms => this.session.emit('sms', sms));
} }
} }
+13 -34
View File
@@ -1,10 +1,9 @@
import type { ErrorName } from './defs/errors.ts'; import type { ErrorName } from './defs/errors.ts';
import type { IncomingDeps } from './incoming-requests.ts';
import type { MessageDlr } from './dlr-merger.ts'; import type { MessageDlr } from './dlr-merger.ts';
import type { ParamValue } from './defs/types.ts'; import type { ParamValue } from './defs/types.ts';
import type { PduObject, PduObjectInput, TlvInputs } from './pdu.ts'; import type { PduObject, PduObjectInput, TlvInputs } from './pdu.ts';
import type { PduRefusedError } from './pdu-refusal.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 { Result, VoidResult } from './result.ts';
import type { SendSmsOptions, SendSmsResult } from './send-sms.ts'; import type { SendSmsOptions, SendSmsResult } from './send-sms.ts';
import type { SmppLog } from './log.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 { refusalAnswer } from './pdu-refusal.ts';
import { guardedLog } from './log.ts'; import { guardedLog } from './log.ts';
import { submitSms, unsent } from './send-sms.ts'; import { submitSms, unsent } from './send-sms.ts';
import { createSms } from './sms.ts';
import { ConcatReference } from './udh.ts'; import { ConcatReference } from './udh.ts';
export type { export type {
@@ -116,16 +114,6 @@ export class Session extends EventEmitter<SessionEvents> {
max: defaults.maxDlrMerges, max: defaults.maxDlrMerges,
timeout: defaults.dlrMergeTimeout, 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.reconnectLoop = this.loopFor(options.reconnect);
this.timers = new LinkTimers({ this.timers = new LinkTimers({
enquireLinkInterval: options.enquireLinkInterval, enquireLinkInterval: options.enquireLinkInterval,
@@ -142,6 +130,18 @@ export class Session extends EventEmitter<SessionEvents> {
responseTimeout: options.responseTimeout ?? defaults.responseTimeout, responseTimeout: options.responseTimeout ?? defaults.responseTimeout,
transport: this.transport, 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(); this.resetTimers();
} }
@@ -257,27 +257,6 @@ export class Session extends EventEmitter<SessionEvents> {
return drained; 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 { private transportFor(sock: Socket): PduTransport {
return new PduTransport({ return new PduTransport({
log: this.log, log: this.log,
+9 -11
View File
@@ -1,6 +1,5 @@
import type { ErrorName } from './defs/errors.ts'; import type { ErrorName } from './defs/errors.ts';
import type { MessageState } from './defs/constants.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 { PduObject, PduObjectInput, TlvInputs } from './pdu.ts';
import type { Result, VoidResult } from './result.ts'; import type { Result, VoidResult } from './result.ts';
import type { Session } from './session.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. */ /** What the session's incoming side gives a message so it can be answered and accounted for. */
export type SmsHandlers = { export type SmsHandlers = {
acceptsOptionalParams: () => boolean;
answer: (pduObj: PduObject, status: ErrorName, params: Record<string, ParamValue>) => Promise<VoidResult>;
bindAllows: (cmdName: string) => boolean;
lostLink: () => boolean; lostLink: () => boolean;
onAnswered: () => void; onAnswered: () => void;
send: (input: PduObjectInput) => Promise<Result<{ pduObj: PduObject }>>; send: (input: PduObjectInput) => Promise<Result<{ pduObj: PduObject }>>;
@@ -93,9 +89,9 @@ export function createSms(input: SmsInput, handlers: SmsHandlers): Sms {
from: input.from, from: input.from,
message: input.message, message: input.message,
pduObjs: input.pduObjs, pduObjs: input.pduObjs,
sendDlr: status => sendDlr(sms, handlers, status), sendDlr: status => sendDlr(sms, input.session, handlers, status),
sendResp: options => (input.answeredAs === undefined sendResp: options => (input.answeredAs === undefined
? sendResp(sms, answered, options ?? {}, handlers) ? sendResp(sms, input.session, answered, options ?? {}, handlers)
: answeredOnArrival(options ?? {}, handlers)), : answeredOnArrival(options ?? {}, handlers)),
session: input.session, session: input.session,
get smsId(): string { get smsId(): string {
@@ -132,9 +128,10 @@ function answeredOnArrival(
async function sendResp( async function sendResp(
sms: Sms, sms: Sms,
session: Session,
answered: { smsId: string }, answered: { smsId: string },
options: SendRespOptions, options: SendRespOptions,
handlers: Pick<SmsHandlers, 'answer' | 'lostLink' | 'onAnswered'>, handlers: Pick<SmsHandlers, 'lostLink' | 'onAnswered'>,
): Promise<VoidResult> { ): Promise<VoidResult> {
const total = sms.pduObjs.length; 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') }; 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, pduObj,
options.status ?? 'ESME_ROK', options.status ?? 'ESME_ROK',
respIdParams(pduObj.cmdName, segmentId(answered.smsId, index, total)), respIdParams(pduObj.cmdName, segmentId(answered.smsId, index, total)),
@@ -214,10 +211,11 @@ function collectReceipt(sent: Result<{ pduObj: PduObject }>[]): SendDlrResult {
async function sendDlr( async function sendDlr(
sms: Sms, sms: Sms,
handlers: Pick<SmsHandlers, 'acceptsOptionalParams' | 'bindAllows' | 'send'>, session: Session,
handlers: Pick<SmsHandlers, 'send'>,
status: MessageState = 'DELIVERED', status: MessageState = 'DELIVERED',
): Promise<SendDlrResult> { ): Promise<SendDlrResult> {
if (!handlers.bindAllows('deliver_sm')) { if (!session.bindAllows('deliver_sm')) {
return { return {
err: new Error('A transmitter-bound session does not carry deliver_sm'), err: new Error('A transmitter-bound session does not carry deliver_sm'),
pduObjs: [], pduObjs: [],
@@ -240,7 +238,7 @@ async function sendDlr(
short_message: receiptText(sms, smsId, status), short_message: receiptText(sms, smsId, status),
source_addr: sms.to, source_addr: sms.to,
}, },
...(handlers.acceptsOptionalParams() ? { tlvs: receiptTlvs(smsId, status) } : {}), ...(session.acceptsOptionalParams() ? { tlvs: receiptTlvs(smsId, status) } : {}),
}); });
})); }));
return collectReceipt(sent); return collectReceipt(sent);
+24 -52
View File
@@ -4,7 +4,7 @@ import test, { describe } from 'node:test';
import type { Collected, LostGroup } from '../src/reassembly.ts'; import type { Collected, LostGroup } from '../src/reassembly.ts';
import type { Dlr } from '../src/dlr.ts'; import type { Dlr } from '../src/dlr.ts';
import type { ErrorName } from '../src/defs/errors.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 { MessageHold } from '../src/held-messages.ts';
import type { MessageState } from '../src/defs/constants.ts'; import type { MessageState } from '../src/defs/constants.ts';
import type { MessageDlr } from '../src/session.ts'; import type { MessageDlr } from '../src/session.ts';
@@ -101,24 +101,14 @@ function abortAfter(
}); });
} }
function stubPort(session: Session, port: Partial<IncomingDeps> = {}): IncomingDeps { function incomingOn(session: Session, options: Partial<IncomingRequestsOptions> = {}): IncomingRequests {
return { return new IncomingRequests({
acceptsOptionalParams: () => true, dlrMerger: new DlrMerger({ log: silentLog, max: 10, timeout: 10_000 }),
answer: (pduObj, status, params) => session.sendReturn(pduObj, status, params), log: silentLog,
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') }), sendPastDrain: () => Promise.resolve({ err: new Error('never sent') }),
smsListeners: () => session.listenerCount('sms'), session,
...port, ...options,
}; });
} }
function submitPdu(seqNr: number, cmdStatus: ErrorName = 'ESME_ROK'): PduObject { function submitPdu(seqNr: number, cmdStatus: ErrorName = 'ESME_ROK'): PduObject {
@@ -758,11 +748,7 @@ describe('reconnect', () => {
closeAfter(t, session); closeAfter(t, session);
const incoming = new IncomingRequests({ const incoming = incomingOn(session, { onRequest: async () => { await delay(10); return false; } });
deps: stubPort(session, { onRequest: async () => { await delay(10); return false; } }),
dlrMerger: new DlrMerger({ log: silentLog, max: 10, timeout: 10_000 }),
log: silentLog,
});
let messages = 0; let messages = 0;
session.on('sms', () => { messages++; }); session.on('sms', () => { messages++; });
@@ -1571,11 +1557,7 @@ describe('held message bounds', () => {
closeAfter(t, session); closeAfter(t, session);
const warnings: string[] = []; const warnings: string[] = [];
const incoming = new IncomingRequests({ const incoming = incomingOn(session, { log: { ...silentLog, warn: message => { warnings.push(message); } } });
deps: stubPort(session),
dlrMerger: new DlrMerger({ log: silentLog, max: 10, timeout: 10_000 }),
log: { ...silentLog, warn: message => { warnings.push(message); } },
});
const answers: (ErrorName | undefined)[] = []; const answers: (ErrorName | undefined)[] = [];
const received: Sms[] = []; const received: Sms[] = [];
@@ -1624,11 +1606,7 @@ describe('held message bounds', () => {
closeAfter(t, session); closeAfter(t, session);
const incoming = new IncomingRequests({ const incoming = incomingOn(session);
deps: stubPort(session),
dlrMerger: new DlrMerger({ log: silentLog, max: 10, timeout: 10_000 }),
log: silentLog,
});
const chunk = Buffer.alloc(64 * 1024); const chunk = Buffer.alloc(64 * 1024);
const carried = submitPdu(1); const carried = submitPdu(1);
let received: Sms | undefined; let received: Sms | undefined;
@@ -1682,6 +1660,9 @@ describe('sendResp()', () => {
closeAfter(t, session); closeAfter(t, session);
let answered = 0; let answered = 0;
session.sendReturn = () => Promise.resolve({ err: new Error('Socket is closed') });
const sms = createSms({ const sms = createSms({
from: '46701113311', from: '46701113311',
message: 'never answered', message: 'never answered',
@@ -1689,9 +1670,6 @@ describe('sendResp()', () => {
session, session,
to: '46709771337', to: '46709771337',
}, { }, {
acceptsOptionalParams: () => true,
answer: () => Promise.resolve({ err: new Error('Socket is closed') }),
bindAllows: () => true,
lostLink: () => false, lostLink: () => false,
onAnswered: () => { answered++; }, onAnswered: () => { answered++; },
send: () => Promise.resolve({ err: new Error('never sent') }), send: () => Promise.resolve({ err: new Error('never sent') }),
@@ -1717,9 +1695,6 @@ describe('sendDlr()', () => {
session, session,
to: '46709771337', to: '46709771337',
}, { }, {
acceptsOptionalParams: () => true,
answer: () => Promise.resolve({}),
bindAllows: () => true,
lostLink: () => false, lostLink: () => false,
onAnswered: () => undefined, onAnswered: () => undefined,
send: () => { send: () => {
@@ -3120,26 +3095,23 @@ describe('graceful shutdown', () => {
closeAfter(t, session); closeAfter(t, session);
const calls: string[] = []; const calls: string[] = [];
const incoming = new IncomingRequests({ const incoming = incomingOn(session);
deps: stubPort(session, { const close = session.close.bind(session);
answer: pduObj => {
session.sendReturn = pduObj => {
calls.push(pduObj.cmdName); calls.push(pduObj.cmdName);
return Promise.resolve({}); return Promise.resolve({});
}, };
peerUnbound: () => { session.close = options => {
calls.push('peerUnbound'); calls.push('close');
return Promise.resolve({}); return close(options);
}, };
}),
dlrMerger: new DlrMerger({ log: silentLog, max: 10, timeout: 10_000 }),
log: silentLog,
});
await incoming.handle({ ...submitPdu(1), cmdId: 0x00000006, cmdName: 'unbind', params: {} }); await incoming.handle({ ...submitPdu(1), cmdId: 0x00000006, cmdName: 'unbind', params: {} });
assert.deepEqual(calls, ['unbind', 'peerUnbound']); assert.deepEqual(calls, ['unbind', 'close']);
}); });
}); });
+22 -3
View File
@@ -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. 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 - [ ] **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. retires the #30 and #46 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 - [ ] **Give the link's liveness one owner.** A third run the same day, after #46, read 6, 6, 7 and
`SmsHandlers`, which makes every inbound call path indirect. Two seats. 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 ### 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 ### 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 - [ ] **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 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 session are up before any bind. The gate decision in `docs/decisions.md` makes the same claim