From cc18569a058caaf711115d2fcc796a9b65d11ce7 Mon Sep 17 00:00:00 2001 From: Lilleman auf Larv Date: Tue, 29 Sep 2026 11:03:02 +0200 Subject: [PATCH] Give the held-message flow one owner --- AGENTS.md | 4 +- src/held-messages.ts | 87 +++++++++++++++++++++++++------------ src/incoming-requests.ts | 34 +++------------ src/sms.ts | 10 ++--- test/session-extras.test.ts | 47 +++++++++++++++----- todo.md | 8 ---- 6 files changed, 107 insertions(+), 83 deletions(-) diff --git a/AGENTS.md b/AGENTS.md index 8ed8839..a238b2f 100644 --- a/AGENTS.md +++ b/AGENTS.md @@ -50,7 +50,7 @@ src/ 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 DlrMerger, HeldMessages and Reassembler share - held-messages.ts HeldMessages: capped, expiring messages the application has not answered, one MessageHold each + held-messages.ts HeldMessages: a message from its `sms` event to its answer, capped and expiring, 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 link-life.ts LinkLife: whether the link lives, and where a request waits for the next one @@ -87,7 +87,7 @@ 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 `IncomingRequests`, which call back into it, and to `OnRequest` and `onConnected` +`createSms()`, `HeldMessages` 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 diff --git a/src/held-messages.ts b/src/held-messages.ts index 4ae9e09..a67f82d 100644 --- a/src/held-messages.ts +++ b/src/held-messages.ts @@ -1,15 +1,28 @@ -import type { PduObject } from './pdu.ts'; +import type { LinkLife } from './link-life.ts'; +import type { PduObject, PduObjectInput } from './pdu.ts'; +import type { Result } from './result.ts'; +import type { Session } from './session.ts'; +import type { SmsHandlers } from './sms.ts'; import type { SmppLog } from './log.ts'; import { ExpiringGroups } from './expiring-groups.ts'; import { IdleWaiters } from './idle-waiters.ts'; +import { createSms } from './sms.ts'; +import { decodeSegments } from './reassembly.ts'; +import { paramText } from './defs/types.ts'; import { retainedOctets } from './retained-pdu.ts'; +type Send = (input: PduObjectInput) => Promise>; + export type HeldMessagesOptions = { + link: LinkLife; log: SmppLog; max: number; maxOctets: number; /** Injected so expiry can be exercised without a wall clock. */ now?: (() => number) | undefined; + /** Past a drain's refusal, for a receipt the drain is itself waiting for. */ + sendPastDrain: Send; + session: Session; timeout: number; }; @@ -20,33 +33,40 @@ function keyOf(pduObjs: PduObject[]): string | undefined { return first ? String(first.seqNr) : undefined; } -type HoldEntry = { - isHeld: () => boolean; - release: () => void; -}; +type HoldRoute = Pick; /** * One message offered to the application. A drain waits on it until the first of: `answered()`, * every listener that took it rejecting, no listener taking it or one throwing, a later message on * its sequence number, its deadline, or the link going. */ -export class MessageHold { - private readonly entry: HoldEntry; +export class MessageHold implements SmsHandlers { + private readonly generation: number; + private readonly pduObjs: PduObject[]; + private readonly route: HoldRoute; + private readonly store: HeldMessages; private working: number; - constructor(entry: HoldEntry, listeners: number) { - this.entry = entry; + constructor(store: HeldMessages, route: HoldRoute, pduObjs: PduObject[], listeners: number) { + this.generation = route.link.generation(); + this.pduObjs = pduObjs; + this.route = route; + this.store = store; this.working = listeners; } /** Whether a drain is still waiting for this message to be answered. */ isHeld(): boolean { - return this.entry.isHeld(); + return this.store.holds(this.pduObjs); } /** A turn later, so a listener sending its receipt straight after the response still holds. */ answered(): void { - setImmediate(() => { this.entry.release(); }); + setImmediate(() => { this.release(); }); + } + + lostLink(): boolean { + return this.route.link.generation() !== this.generation; } /** A rejection leaves the other listeners running, so only the last one to fail gives the message up. */ @@ -58,7 +78,12 @@ export class MessageHold { /** At once, for a message nobody took or a listener threw on: that is not work a shutdown can wait for. */ release(): void { - this.entry.release(); + this.store.release(this.pduObjs); + } + + /** A receipt for a message still held is what a drain waits for, so it goes out past the drain. */ + send(input: PduObjectInput): Promise> { + return this.isHeld() ? this.route.sendPastDrain(input) : this.route.session.send(input); } } @@ -70,6 +95,7 @@ export class HeldMessages { private readonly maxOctets: number; /** A rejecting listener hands the message back as an `unknown`, so its hold is found by identity. */ private readonly offered = new WeakMap(); + private readonly route: HoldRoute; constructor(options: HeldMessagesOptions) { this.held = new ExpiringGroups({ @@ -80,6 +106,7 @@ export class HeldMessages { }); this.log = options.log; this.maxOctets = options.maxOctets; + this.route = { link: options.link, sendPastDrain: options.sendPastDrain, session: options.session }; } get octetsHeld(): number { @@ -98,10 +125,7 @@ export class HeldMessages { } private hold(pduObjs: PduObject[], listeners: number): MessageHold { - const hold = new MessageHold({ - isHeld: () => this.has(pduObjs), - release: () => { this.release(pduObjs); }, - }, listeners); + const hold = new MessageHold(this, this.route, pduObjs, listeners); const key = keyOf(pduObjs); if (key === undefined) return hold; @@ -118,18 +142,25 @@ export class HeldMessages { return hold; } - offer( - pduObjs: PduObject[], - listeners: number, - build: (hold: MessageHold) => T, - emit: (message: T) => boolean, - ): MessageHold { - const hold = this.hold(pduObjs, listeners); - const message = build(hold); + /** Hands a whole message to the application as an `sms` event, held until it is answered. */ + offer(pduObjs: PduObject[], answeredAs?: string): MessageHold | undefined { + const first = pduObjs[0]; - this.offered.set(message, hold); + if (!first) return undefined; - if (!emit(message)) hold.release(); + const hold = this.hold(pduObjs, this.route.session.listenerCount('sms')); + const sms = createSms({ + answeredAs, + from: paramText(first.params.source_addr), + message: decodeSegments(pduObjs), + pduObjs, + session: this.route.session, + to: paramText(first.params.destination_addr), + }, hold); + + this.offered.set(sms, hold); + + if (!this.route.session.emit('sms', sms)) hold.release(); return hold; } @@ -141,13 +172,13 @@ export class HeldMessages { this.offered.get(message)?.listenerGaveUp(); } - private has(pduObjs: PduObject[]): boolean { + holds(pduObjs: PduObject[]): boolean { const key = keyOf(pduObjs); return key !== undefined && this.held.get(key) === pduObjs; } - private release(pduObjs: PduObject[]): void { + release(pduObjs: PduObject[]): void { const key = keyOf(pduObjs); // Identity, not the key: a wrapped sequence number must not release someone else's message. diff --git a/src/incoming-requests.ts b/src/incoming-requests.ts index f2a39f3..cf61499 100644 --- a/src/incoming-requests.ts +++ b/src/incoming-requests.ts @@ -10,13 +10,11 @@ import type { Session } from './session.ts'; import type { SmppLog } from './log.ts'; import type { SmsIdFormat } from './sms-id.ts'; import { HeldMessages } from './held-messages.ts'; -import { Reassembler, decodeSegments } from './reassembly.ts'; +import { Reassembler } 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'; import { respIdParams, segmentId } from './sms-id.ts'; import { respNameFor } from './defs/commands.ts'; @@ -67,7 +65,6 @@ export class IncomingRequests { 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; @@ -76,9 +73,12 @@ export class IncomingRequests { constructor(options: IncomingRequestsOptions) { this.dlrMerger = options.dlrMerger; this.held = new HeldMessages({ + link: options.link, log: options.log, max: defaults.maxHeldMessages, maxOctets: defaults.maxHeldOctets, + sendPastDrain: options.sendPastDrain, + session: options.session, timeout: defaults.heldMessageTimeout, }); this.link = options.link; @@ -91,7 +91,6 @@ export class IncomingRequests { onLost: lost => { this.reportLost(lost); }, timeout: options.reassemblyTimeout ?? defaults.reassemblyTimeout, }); - this.sendPastDrain = options.sendPastDrain; this.session = options.session; this.smsIdFormat = options.smsIdFormat ?? {}; this.systemId = options.systemId ?? defaults.systemId; @@ -254,7 +253,7 @@ export class IncomingRequests { const concat = concatOf(pduObj); if (!concat) { - this.emitSms([detach(pduObj)]); + this.held.offer([detach(pduObj)]); return; } @@ -276,7 +275,7 @@ export class IncomingRequests { respIdParams(pduObj.cmdName, segmentId(collected.smsId, concat.part - 1, concat.total)), ); - if (collected.whole) this.emitSms(collected.whole, collected.smsId); + if (collected.whole) this.held.offer(collected.whole, collected.smsId); } private reportLost(lost: LostGroup): void { @@ -284,25 +283,4 @@ export class IncomingRequests { `Gave up ${String(lost.parts)} of ${String(lost.total)} segments of an incomplete concatenated message: ${lostReasons[lost.reason]}`, )); } - - private emitSms(pduObjs: PduObject[], answeredAs?: string): void { - const first = pduObjs[0]; - - if (!first) return; - - const generation = this.link.generation(); - - 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), - }, { - lostLink: () => this.link.generation() !== generation, - onAnswered: () => { hold.answered(); }, - send: input => (hold.isHeld() ? this.sendPastDrain(input) : this.session.send(input)), - }), sms => this.session.emit('sms', sms)); - } } diff --git a/src/sms.ts b/src/sms.ts index 93f75c7..52eb357 100644 --- a/src/sms.ts +++ b/src/sms.ts @@ -68,8 +68,8 @@ export type SmsInput = { /** What the session's incoming side gives a message so it can be answered and accounted for. */ export type SmsHandlers = { + answered: () => void; lostLink: () => boolean; - onAnswered: () => void; send: (input: PduObjectInput) => Promise>; }; @@ -107,7 +107,7 @@ export function createSms(input: SmsInput, handlers: SmsHandlers): Sms { /** Every segment went out answered, so the call is what the shutdown waits for and nothing else. */ function answeredOnArrival( options: SendRespOptions, - handlers: Pick, + handlers: Pick, ): Promise { if (options.smsId !== undefined) { return Promise.resolve({ @@ -121,7 +121,7 @@ function answeredOnArrival( }); } - handlers.onAnswered(); + handlers.answered(); return Promise.resolve({}); } @@ -131,7 +131,7 @@ async function sendResp( session: Session, answered: { smsId: string }, options: SendRespOptions, - handlers: Pick, + handlers: Pick, ): Promise { const total = sms.pduObjs.length; @@ -158,7 +158,7 @@ async function sendResp( const failure = results.find(result => result.err); - if (!failure) handlers.onAnswered(); + if (!failure) handlers.answered(); return failure ?? {}; } diff --git a/test/session-extras.test.ts b/test/session-extras.test.ts index a6b75a1..b2104ba 100644 --- a/test/session-extras.test.ts +++ b/test/session-extras.test.ts @@ -5,7 +5,7 @@ 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 { IncomingRequestsOptions } from '../src/incoming-requests.ts'; -import type { MessageHold } from '../src/held-messages.ts'; +import type { HeldMessagesOptions, MessageHold } from '../src/held-messages.ts'; import type { MessageState } from '../src/defs/constants.ts'; import type { MessageDlr } from '../src/session.ts'; import type { PduObject, PduObjectInput } from '../src/pdu.ts'; @@ -1543,11 +1543,34 @@ describe('held message bounds', () => { } function offer(held: HeldMessages, seqNr: number): MessageHold { - return held.offer(message(seqNr), 1, () => ({}), () => true); + const hold = held.offer(message(seqNr)); + + assert.ok(hold); + + return hold; } - test('is full at its count, and a re-used sequence number replaces rather than adding', () => { - const held = new HeldMessages({ log: silentLog, max: 2, maxOctets: 1_000_000, timeout: 10_000 }); + /** Offers to a session with a listener, so an offer is held rather than released as untaken. */ + function heldOn( + t: TestContext, + options: Pick, + ): HeldMessages { + const session = new Session({ sock: new net.Socket() }); + + closeAfter(t, session); + session.on('sms', () => undefined); + + return new HeldMessages({ + ...options, + link: new LinkLife({ log: silentLog, reconnects: false, timeout: 100 }), + log: silentLog, + sendPastDrain: () => Promise.resolve({ err: new Error('never sent') }), + session, + }); + } + + test('is full at its count, and a re-used sequence number replaces rather than adding', t => { + const held = heldOn(t, { max: 2, maxOctets: 1_000_000, timeout: 10_000 }); const first = offer(held, 1); const replaced = offer(held, 2); @@ -1563,9 +1586,9 @@ describe('held message bounds', () => { }); // submitPdu() holds 1026 octets by the maxOctets charge: its object, and the three text fields. - test('is full at its octet cap, until a message leaves by any way out', () => { + test('is full at its octet cap, until a message leaves by any way out', t => { let now = 0; - const held = new HeldMessages({ log: silentLog, max: 10, maxOctets: 2000, now: () => now, timeout: 10_000 }); + const held = heldOn(t, { max: 10, maxOctets: 2000, now: () => now, timeout: 10_000 }); const answered = offer(held, 1); assert.equal(held.full(), false); @@ -1658,9 +1681,9 @@ describe('held message bounds', () => { incoming.clear(); }); - test('gives up on a message the application never answers', () => { + test('gives up on a message the application never answers', t => { let now = 0; - const held = new HeldMessages({ log: silentLog, max: 10, maxOctets: 1_000_000, now: () => now, timeout: 60 }); + const held = heldOn(t, { max: 10, maxOctets: 1_000_000, now: () => now, timeout: 60 }); offer(held, 1); now = 61; @@ -1674,9 +1697,9 @@ describe('held message bounds', () => { }); // Without this the drain sits out its whole budget before returning what a sweep already settled. - test('wakes a waiting drain when the last message expires', async () => { + test('wakes a waiting drain when the last message expires', async t => { let now = 0; - const held = new HeldMessages({ log: silentLog, max: 10, maxOctets: 1_000_000, now: () => now, timeout: 60 }); + const held = heldOn(t, { max: 10, maxOctets: 1_000_000, now: () => now, timeout: 60 }); offer(held, 1); @@ -1708,7 +1731,7 @@ describe('sendResp()', () => { to: '46709771337', }, { lostLink: () => false, - onAnswered: () => { answered++; }, + answered: () => { answered++; }, send: () => Promise.resolve({ err: new Error('never sent') }), }); @@ -1733,7 +1756,7 @@ describe('sendDlr()', () => { to: '46709771337', }, { lostLink: () => false, - onAnswered: () => undefined, + answered: () => undefined, send: () => { call++; diff --git a/todo.md b/todo.md index 6406b80..4960d3d 100644 --- a/todo.md +++ b/todo.md @@ -209,14 +209,6 @@ hardest, and the held-message flow across `incoming-requests.ts`, `held-messages - [ ] **Lift Locality to 7, and confirm it with a scoring run.** A run reading 7.0 or above also retires the #30, #46 and #48 decision. -- [ ] **Give the held-message flow one owner — next, the condition #48 merged under.** Whether a - message is still held, and so whether its receipt may pass the drain, is decided across - `IncomingRequests.emitSms()`, `HeldMessages.offer()`/`MessageHold`, `createSms()`'s handlers in - `sms.ts` and `Session`'s rejection handler, which finds the hold again through a `WeakMap` - keyed on the `Sms`; `MessageHold.answered()` defers its release a `setImmediate` so a - `sendDlr()` straight after `sendResp()` still counts as held. All four seats of the #48 run - ranked it second hardest; the inherited architect put it at about two days. - ### Correctness - [ ] **Refuse to open a link that dropped while its rebind was answered.** A peer sending