Give IncomingRequests a port instead of the Session it drives
Test / lint (pull_request) Successful in 23s
Test / test (18) (pull_request) Successful in 31s
Test / test (20) (pull_request) Successful in 37s
Test / test (22) (pull_request) Successful in 32s
Test / test (24) (pull_request) Successful in 34s
Test / test (26) (pull_request) Successful in 37s
Mirror / push (push) Successful in 5s

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