Compare commits

8 Commits

Author SHA1 Message Date
lilleman 9d3dbc522c Drop the doc offer() restated
Mirror / push (push) Has been cancelled
Test / lint (pull_request) Successful in 24s
Test / test (18) (pull_request) Successful in 32s
Test / test (22) (pull_request) Successful in 31s
Test / test (20) (pull_request) Successful in 31s
Test / test (24) (pull_request) Successful in 34s
Test / test (26) (pull_request) Successful in 31s
2026-09-28 10:54:50 +02:00
lilleman 936102fd04 Assert a re-used sequence number ends the hold it replaced
Mirror / push (push) Has been cancelled
Test / test (20) (pull_request) Successful in 32s
Test / test (22) (pull_request) Successful in 33s
Test / test (26) (pull_request) Successful in 32s
Test / test (18) (pull_request) Successful in 32s
Test / test (24) (pull_request) Successful in 32s
Test / lint (pull_request) Failing after 14m40s
2026-09-28 10:54:31 +02:00
lilleman 9df7f05669 Make offer() the only way to hold a message
Mirror / push (push) Successful in 7s
Test / lint (pull_request) Successful in 22s
Test / test (18) (pull_request) Successful in 32s
Test / test (20) (pull_request) Successful in 31s
Test / test (22) (pull_request) Successful in 31s
Test / test (24) (pull_request) Successful in 33s
Test / test (26) (pull_request) Successful in 32s
2026-09-28 10:52:36 +02:00
lilleman 245d0e2740 Hold the bounds tests through offer() 2026-09-28 10:51:39 +02:00
lilleman 1db63e19ed Return the hold offer() took 2026-09-28 10:51:07 +02:00
lilleman b5e6b03ff8 File a throwing sms listener releasing its hold 2026-09-28 10:51:07 +02:00
lilleman d8c2647e75 Name every way a hold ends 2026-09-28 10:50:54 +02:00
lilleman e52baafd29 Give the held-message flow one place a reader can follow it
Mirror / push (push) Successful in 6s
Test / lint (pull_request) Successful in 24s
Test / test (18) (pull_request) Successful in 32s
Test / test (20) (pull_request) Successful in 31s
Test / test (24) (pull_request) Successful in 32s
Test / test (22) (pull_request) Successful in 32s
Test / test (26) (pull_request) Successful in 31s
2026-09-28 10:48:38 +02:00
4 changed files with 57 additions and 31 deletions
+32 -3
View File
@@ -25,7 +25,11 @@ type HoldEntry = {
release: () => void; 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 { export class MessageHold {
private readonly entry: HoldEntry; private readonly entry: HoldEntry;
private working: number; private working: number;
@@ -52,7 +56,7 @@ export class MessageHold {
if (this.working <= 0) this.answered(); 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 { release(): void {
this.entry.release(); this.entry.release();
} }
@@ -64,6 +68,8 @@ export class HeldMessages {
private readonly idleWaiters = new IdleWaiters(); private readonly idleWaiters = new IdleWaiters();
private readonly log: SmppLog; private readonly log: SmppLog;
private readonly maxOctets: number; 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<object, MessageHold>();
constructor(options: HeldMessagesOptions) { constructor(options: HeldMessagesOptions) {
this.held = new ExpiringGroups({ this.held = new ExpiringGroups({
@@ -91,7 +97,7 @@ export class HeldMessages {
return this.held.full || this.held.weight >= this.maxOctets; 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({ const hold = new MessageHold({
isHeld: () => this.has(pduObjs), isHeld: () => this.has(pduObjs),
release: () => { this.release(pduObjs); }, release: () => { this.release(pduObjs); },
@@ -112,6 +118,29 @@ export class HeldMessages {
return hold; return hold;
} }
offer<T extends object>(
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 { private has(pduObjs: PduObject[]): boolean {
const key = keyOf(pduObjs); const key = keyOf(pduObjs);
+3 -14
View File
@@ -10,7 +10,6 @@ import type { Result, VoidResult } from './result.ts';
import type { SmppLog } from './log.ts'; import type { SmppLog } from './log.ts';
import type { Sms, SmsHandlers, SmsInput } from './sms.ts'; import type { Sms, SmsHandlers, SmsInput } from './sms.ts';
import type { SmsIdFormat } from './sms-id.ts'; import type { SmsIdFormat } from './sms-id.ts';
import type { MessageHold } from './held-messages.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';
@@ -85,8 +84,6 @@ export class IncomingRequests {
private readonly deps: IncomingDeps; private readonly deps: IncomingDeps;
private readonly dlrMerger: DlrMerger; private readonly dlrMerger: DlrMerger;
private readonly held: HeldMessages; 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<object, MessageHold>();
private readonly log: SmppLog; private readonly log: SmppLog;
private readonly reassembler: Reassembler; private readonly reassembler: Reassembler;
private readonly smsIdFormat: SmsIdFormat; private readonly smsIdFormat: SmsIdFormat;
@@ -172,11 +169,8 @@ export class IncomingRequests {
this.reassembler.clear(); 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 { listenerRejected(sms: unknown): void {
if (typeof sms !== 'object' || sms === null) return; this.held.listenerRejected(sms);
this.holds.get(sms)?.listenerGaveUp();
} }
/** Waits out the messages the application still holds, and says how many it never answered. */ /** 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; if (!first) return;
const generation = this.linkGeneration; 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, answeredAs,
from: paramText(first.params.source_addr), from: paramText(first.params.source_addr),
message: decodeSegments(pduObjs), message: decodeSegments(pduObjs),
@@ -325,10 +318,6 @@ export class IncomingRequests {
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.deps.sendPastDrain(input) : this.deps.send(input)),
}); }), sms => this.deps.offerSms(sms));
this.holds.set(sms, hold);
if (!this.deps.offerSms(sms)) hold.release();
} }
} }
+17 -11
View File
@@ -5,6 +5,7 @@ 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 { IncomingDeps } from '../src/incoming-requests.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';
import type { PduObject, PduObjectInput } from '../src/pdu.ts'; import type { PduObject, PduObjectInput } from '../src/pdu.ts';
@@ -1518,17 +1519,22 @@ describe('held message bounds', () => {
return [submitPdu(seqNr)]; 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', () => { 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 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); offer(held, 2);
held.hold(message(2), 1);
assert.equal(held.size, 2); assert.equal(held.size, 2);
assert.equal(held.octetsHeld, 2 * 1026, 'the replaced message leaves its octets with it'); assert.equal(held.octetsHeld, 2 * 1026, 'the replaced message leaves its octets with it');
assert.equal(held.full(), true); assert.equal(held.full(), true);
assert.equal(first.isHeld(), true); assert.equal(first.isHeld(), true);
assert.equal(replaced.isHeld(), false);
held.clear(); 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', () => { test('is full at its octet cap, until a message leaves by any way out', () => {
let now = 0; let now = 0;
const held = new HeldMessages({ log: silentLog, max: 10, maxOctets: 2000, now: () => now, timeout: 10_000 }); 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); assert.equal(held.full(), false);
held.hold(message(2), 1); offer(held, 2);
assert.equal(held.full(), true); assert.equal(held.full(), true);
answered.release(); answered.release();
assert.equal(held.full(), false, 'after a release'); assert.equal(held.full(), false, 'after a release');
held.hold(message(3), 1); offer(held, 3);
now = 20_000; now = 20_000;
held.sweep(); held.sweep();
now = 0; now = 0;
assert.equal(held.full(), false, 'after a sweep'); assert.equal(held.full(), false, 'after a sweep');
held.hold(message(4), 1); offer(held, 4);
held.hold(message(5), 1); offer(held, 5);
held.clear(); held.clear();
assert.equal(held.full(), false, 'after a clear'); assert.equal(held.full(), false, 'after a clear');
@@ -1641,11 +1647,11 @@ describe('held message bounds', () => {
let now = 0; let now = 0;
const held = new HeldMessages({ log: silentLog, max: 10, maxOctets: 1_000_000, now: () => now, timeout: 60 }); 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; now = 61;
// The next message sweeps the one that expired, so only the new one is still waited for. // 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); assert.equal(held.size, 1);
@@ -1657,7 +1663,7 @@ describe('held message bounds', () => {
let now = 0; let now = 0;
const held = new HeldMessages({ log: silentLog, max: 10, maxOctets: 1_000_000, now: () => now, timeout: 60 }); 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); const waiting = held.idle(1000, undefined);
+5 -3
View File
@@ -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 - [ ] **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 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 - [ ] **Shrink the `IncomingDeps` closure bag.** 16 lambdas, six of them repeated in
`SmsHandlers`, which makes every inbound call path indirect. Two seats. `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 so `"\x1B("` is detected as GSM, goes out as 0x1B 0x28 and arrives as `{`. From the
2026-09-28 scoring run. 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 ### 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, - [ ] **Close the gap to jsmpp at `maxOutstanding: 10`.** Measured 2026-09-20 against the same sink,