From 0bea36c069d411259366aad495333bfffcf16135 Mon Sep 17 00:00:00 2001 From: Lilleman auf Larv Date: Mon, 28 Sep 2026 10:54:55 +0200 Subject: [PATCH] Give the held-message flow one place a reader can follow it --- src/held-messages.ts | 35 ++++++++++++++++++++++++++++++++--- src/incoming-requests.ts | 17 +++-------------- test/session-extras.test.ts | 28 +++++++++++++++++----------- todo.md | 8 +++++--- 4 files changed, 57 insertions(+), 31 deletions(-) diff --git a/src/held-messages.ts b/src/held-messages.ts index cfacbd9..4ae9e09 100644 --- a/src/held-messages.ts +++ b/src/held-messages.ts @@ -25,7 +25,11 @@ type HoldEntry = { release: () => void; }; -/** One message offered to the application, held until it is answered or every listener gives up. */ +/** + * 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; private working: number; @@ -52,7 +56,7 @@ export class MessageHold { if (this.working <= 0) this.answered(); } - /** At once, for a message nobody took: that is not work a shutdown can wait for. */ + /** 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(); } @@ -64,6 +68,8 @@ export class HeldMessages { private readonly idleWaiters = new IdleWaiters(); private readonly log: SmppLog; 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(); constructor(options: HeldMessagesOptions) { this.held = new ExpiringGroups({ @@ -91,7 +97,7 @@ export class HeldMessages { return this.held.full || this.held.weight >= this.maxOctets; } - hold(pduObjs: PduObject[], listeners: number): MessageHold { + private hold(pduObjs: PduObject[], listeners: number): MessageHold { const hold = new MessageHold({ isHeld: () => this.has(pduObjs), release: () => { this.release(pduObjs); }, @@ -112,6 +118,29 @@ 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); + + this.offered.set(message, hold); + + if (!emit(message)) hold.release(); + + return hold; + } + + /** One listener gave up on a message; the last one to do so is what releases it. */ + listenerRejected(message: unknown): void { + if (typeof message !== 'object' || message === null) return; + + this.offered.get(message)?.listenerGaveUp(); + } + private has(pduObjs: PduObject[]): boolean { const key = keyOf(pduObjs); diff --git a/src/incoming-requests.ts b/src/incoming-requests.ts index 4f926de..e3a3d18 100644 --- a/src/incoming-requests.ts +++ b/src/incoming-requests.ts @@ -10,7 +10,6 @@ import type { Result, VoidResult } from './result.ts'; import type { SmppLog } from './log.ts'; import type { Sms, SmsHandlers, SmsInput } from './sms.ts'; import type { SmsIdFormat } from './sms-id.ts'; -import type { MessageHold } from './held-messages.ts'; import { HeldMessages } from './held-messages.ts'; import { Reassembler, decodeSegments } from './reassembly.ts'; import { bindCommands, defaults, standsInFor } from './session-options.ts'; @@ -85,8 +84,6 @@ export class IncomingRequests { private readonly deps: IncomingDeps; private readonly dlrMerger: DlrMerger; private readonly held: HeldMessages; - /** The rejection handler is handed the Sms back as an `unknown`, so its hold is found by identity. */ - private readonly holds = new WeakMap(); private readonly log: SmppLog; private readonly reassembler: Reassembler; private readonly smsIdFormat: SmsIdFormat; @@ -172,11 +169,8 @@ export class IncomingRequests { this.reassembler.clear(); } - /** One `sms` listener gave up on a message; the last one to do so is what releases the hold. */ listenerRejected(sms: unknown): void { - if (typeof sms !== 'object' || sms === null) return; - - this.holds.get(sms)?.listenerGaveUp(); + this.held.listenerRejected(sms); } /** Waits out the messages the application still holds, and says how many it never answered. */ @@ -310,9 +304,8 @@ export class IncomingRequests { if (!first) return; const generation = this.linkGeneration; - const hold = this.held.hold(pduObjs, this.deps.smsListeners()); - const sms = this.deps.createSms({ + this.held.offer(pduObjs, this.deps.smsListeners(), hold => this.deps.createSms({ answeredAs, from: paramText(first.params.source_addr), message: decodeSegments(pduObjs), @@ -325,10 +318,6 @@ export class IncomingRequests { lostLink: () => this.linkGeneration !== generation, onAnswered: () => { hold.answered(); }, send: input => (hold.isHeld() ? this.deps.sendPastDrain(input) : this.deps.send(input)), - }); - - this.holds.set(sms, hold); - - if (!this.deps.offerSms(sms)) hold.release(); + }), sms => this.deps.offerSms(sms)); } } diff --git a/test/session-extras.test.ts b/test/session-extras.test.ts index 209dfa3..cce32b1 100644 --- a/test/session-extras.test.ts +++ b/test/session-extras.test.ts @@ -5,6 +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 { IncomingDeps } 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'; import type { PduObject, PduObjectInput } from '../src/pdu.ts'; @@ -1518,17 +1519,22 @@ describe('held message bounds', () => { return [submitPdu(seqNr)]; } + function offer(held: HeldMessages, seqNr: number): MessageHold { + return held.offer(message(seqNr), 1, () => ({}), () => true); + } + 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 }); - const first = held.hold(message(1), 1); + const first = offer(held, 1); + const replaced = offer(held, 2); - held.hold(message(2), 1); - held.hold(message(2), 1); + offer(held, 2); assert.equal(held.size, 2); assert.equal(held.octetsHeld, 2 * 1026, 'the replaced message leaves its octets with it'); assert.equal(held.full(), true); assert.equal(first.isHeld(), true); + assert.equal(replaced.isHeld(), false); held.clear(); }); @@ -1537,22 +1543,22 @@ describe('held message bounds', () => { test('is full at its octet cap, until a message leaves by any way out', () => { let now = 0; const held = new HeldMessages({ log: silentLog, max: 10, maxOctets: 2000, now: () => now, timeout: 10_000 }); - const answered = held.hold(message(1), 1); + const answered = offer(held, 1); assert.equal(held.full(), false); - held.hold(message(2), 1); + offer(held, 2); assert.equal(held.full(), true); answered.release(); assert.equal(held.full(), false, 'after a release'); - held.hold(message(3), 1); + offer(held, 3); now = 20_000; held.sweep(); now = 0; assert.equal(held.full(), false, 'after a sweep'); - held.hold(message(4), 1); - held.hold(message(5), 1); + offer(held, 4); + offer(held, 5); held.clear(); assert.equal(held.full(), false, 'after a clear'); @@ -1641,11 +1647,11 @@ describe('held message bounds', () => { let now = 0; const held = new HeldMessages({ log: silentLog, max: 10, maxOctets: 1_000_000, now: () => now, timeout: 60 }); - held.hold(message(1), 1); + offer(held, 1); now = 61; // The next message sweeps the one that expired, so only the new one is still waited for. - held.hold(message(2), 1); + offer(held, 2); assert.equal(held.size, 1); @@ -1657,7 +1663,7 @@ describe('held message bounds', () => { let now = 0; const held = new HeldMessages({ log: silentLog, max: 10, maxOctets: 1_000_000, now: () => now, timeout: 60 }); - held.hold(message(1), 1); + offer(held, 1); const waiting = held.idle(1000, undefined); diff --git a/todo.md b/todo.md index 4602348..2ddacdc 100644 --- a/todo.md +++ b/todo.md @@ -204,9 +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. -- [ ] **Give the held-message flow one place a reader can follow it.** Whether a drain still waits - on a message is spread over `emitSms()`, `MessageHold`, the session's rejection route and - `Sms.isHeld()`. Four seats. - [ ] **Shrink the `IncomingDeps` closure bag.** 16 lambdas, six of them repeated in `SmsHandlers`, which makes every inbound call path indirect. Two seats. @@ -268,6 +265,11 @@ and 5. Every seat ranked the session's lifecycle hardest and least wanted to mod so `"\x1B("` is detected as GSM, goes out as 0x1B 0x28 and arrives as `{`. From the 2026-09-28 scoring run. +- [ ] **Count an `sms` listener that throws as one giving up, as a rejection already is.** A + synchronous throw makes `Session.emit` return false, and `HeldMessages.offer()` then releases + the hold at once, so a drain stops waiting on an async listener still answering beside it. + Goal 2. From the stability review of #45. + ### Throughput — goal 6, and the default window is where we are slowest - [ ] **Close the gap to jsmpp at `maxOutstanding: 10`.** Measured 2026-09-20 against the same sink,