From 91c155175fc21de1ff1b9ae495384fdc24b65c79 Mon Sep 17 00:00:00 2001 From: Lilleman auf Larv Date: Mon, 28 Sep 2026 01:09:12 +0200 Subject: [PATCH 1/4] Derive the retained octets from ExpiringGroups' weight --- AGENTS.md | 2 +- src/expiring-groups.ts | 51 ++++++++++++++++++++++++++++++++--- src/held-messages.ts | 40 ++++++--------------------- src/reassembly.ts | 54 +++++++++++++++---------------------- test/session-extras.test.ts | 2 ++ todo.md | 6 ----- 6 files changed, 79 insertions(+), 76 deletions(-) diff --git a/AGENTS.md b/AGENTS.md index 92a902d..a197a6a 100644 --- a/AGENTS.md +++ b/AGENTS.md @@ -51,7 +51,7 @@ src/ dlr.ts Delivery receipts: text and TLV parsing, receipt status codes 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, expiring store both of those share + expiring-groups.ts ExpiringGroups: the capped, weighed, expiring store both of those share 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 incoming-requests.ts Every request the peer sends: messages, receipts, links, unknown commands diff --git a/src/expiring-groups.ts b/src/expiring-groups.ts index 9c154cb..80440b6 100644 --- a/src/expiring-groups.ts +++ b/src/expiring-groups.ts @@ -1,5 +1,7 @@ export type ExpiringGroupsOptions = { max: number; + /** What the groups may weigh together before weigh() evicts the oldest. */ + maxWeight?: number | undefined; /** Injected so expiry can be exercised without a wall clock. */ now?: (() => number) | undefined; /** Runs on the sweeper's own timer; the owner reports whatever it takes out. */ @@ -10,6 +12,7 @@ export type ExpiringGroupsOptions = { type Entry = { deadline: number; group: T; + weight: number; }; /** @@ -19,13 +22,16 @@ type Entry = { export class ExpiringGroups { private readonly entries = new Map>(); private readonly max: number; + private readonly maxWeight: number; private readonly now: () => number; private readonly onSweep: () => void; private readonly timeout: number; private sweeper: NodeJS.Timeout | undefined; + private total = 0; constructor(options: ExpiringGroupsOptions) { this.max = options.max; + this.maxWeight = options.maxWeight ?? Infinity; this.now = options.now ?? Date.now; this.onSweep = options.onSweep; this.timeout = options.timeout; @@ -39,13 +45,19 @@ export class ExpiringGroups { return this.entries.size; } + get weight(): number { + return this.total; + } + get(key: string): T | undefined { return this.entries.get(key)?.group; } /** Starts the group's deadline, and the sweeper if this is the only group held. */ - set(key: string, group: T): void { - this.entries.set(key, { deadline: this.now() + this.timeout, group }); + set(key: string, group: T, weight = 0): void { + this.remove(key); + this.entries.set(key, { deadline: this.now() + this.timeout, group, weight }); + this.total += weight; if (this.sweeper) return; @@ -54,10 +66,31 @@ export class ExpiringGroups { } delete(key: string): void { - this.entries.delete(key); + this.remove(key); this.idle(); } + /** Records what a group weighs now, and hands over the oldest groups evicted to bring the total under maxWeight. */ + weigh(key: string, weight: number): [string, T][] { + const entry = this.entries.get(key); + const taken: [string, T][] = []; + + if (entry) { + this.total += weight - entry.weight; + entry.weight = weight; + } + + while (this.total > this.maxWeight) { + const oldest = this.takeOldest(); + + if (!oldest) break; + + taken.push(oldest); + } + + return taken; + } + /** Removes every group and hands them over, so an owner that must account for them can. */ takeAll(): [string, T][] { const taken: [string, T][] = []; @@ -67,6 +100,7 @@ export class ExpiringGroups { } this.entries.clear(); + this.total = 0; this.idle(); return taken; @@ -81,7 +115,7 @@ export class ExpiringGroups { if (entry.deadline > now) continue; taken.push([key, entry.group]); - this.entries.delete(key); + this.remove(key); } this.idle(); @@ -102,6 +136,15 @@ export class ExpiringGroups { return [key, entry.group]; } + private remove(key: string): void { + const entry = this.entries.get(key); + + if (!entry) return; + + this.entries.delete(key); + this.total -= entry.weight; + } + private idle(): void { if (!this.sweeper || this.entries.size > 0) return; diff --git a/src/held-messages.ts b/src/held-messages.ts index 07ac2fa..a519d1a 100644 --- a/src/held-messages.ts +++ b/src/held-messages.ts @@ -20,11 +20,6 @@ function keyOf(pduObjs: PduObject[]): string | undefined { return first ? String(first.seqNr) : undefined; } -type Held = { - octets: number; - pduObjs: PduObject[]; -}; - type HoldEntry = { isHeld: () => boolean; release: () => void; @@ -65,11 +60,10 @@ export class MessageHold { /** The messages handed to the application that it has not answered yet, held by their segments. */ export class HeldMessages { - private readonly held: ExpiringGroups; + private readonly held: ExpiringGroups; private readonly idleWaiters = new IdleWaiters(); private readonly log: SmppLog; private readonly maxOctets: number; - private octets = 0; constructor(options: HeldMessagesOptions) { this.held = new ExpiringGroups({ @@ -83,7 +77,7 @@ export class HeldMessages { } get octetsHeld(): number { - return this.octets; + return this.held.weight; } get size(): number { @@ -94,7 +88,7 @@ export class HeldMessages { full(): boolean { this.sweep(); - return this.held.full || this.octets >= this.maxOctets; + return this.held.full || this.held.weight >= this.maxOctets; } hold(pduObjs: PduObject[], listeners: number): MessageHold { @@ -108,17 +102,11 @@ export class HeldMessages { this.sweep(); - const replaced = this.held.get(key); - - if (replaced) { + if (this.held.get(key)) { this.log.warn('heldMessages - replacing a message on a re-used sequence number', { seqNr: Number(key) }); - this.delete(key, replaced); } - const octets = pduObjs.reduce((sum, pduObj) => sum + retainedOctets(pduObj), 0); - - this.held.set(key, { octets, pduObjs }); - this.octets += octets; + this.held.set(key, pduObjs, pduObjs.reduce((sum, pduObj) => sum + retainedOctets(pduObj), 0)); return hold; } @@ -126,25 +114,22 @@ export class HeldMessages { private has(pduObjs: PduObject[]): boolean { const key = keyOf(pduObjs); - return key !== undefined && this.held.get(key)?.pduObjs === pduObjs; + return key !== undefined && this.held.get(key) === pduObjs; } private release(pduObjs: PduObject[]): void { const key = keyOf(pduObjs); // Identity, not the key: a wrapped sequence number must not release someone else's message. - const held = key === undefined ? undefined : this.held.get(key); + if (key === undefined || this.held.get(key) !== pduObjs) return; - if (key === undefined || held?.pduObjs !== pduObjs) return; - - this.delete(key, held); + this.held.delete(key); this.settle(); } /** Drops every message: their segments went with the link, so no answer of ours correlates now. */ clear(): void { this.held.takeAll(); - this.octets = 0; this.idleWaiters.settle(); } @@ -159,21 +144,12 @@ export class HeldMessages { if (expired.length === 0) return; - for (const [, held] of expired) { - this.octets -= held.octets; - } - this.log.warn('heldMessages - messages the application never answered', { messages: expired.length, }); this.settle(); } - private delete(key: string, held: Held): void { - this.held.delete(key); - this.octets -= held.octets; - } - private settle(): void { if (this.held.size === 0) this.idleWaiters.settle(); } diff --git a/src/reassembly.ts b/src/reassembly.ts index 1c2d6bd..4f3d2c5 100644 --- a/src/reassembly.ts +++ b/src/reassembly.ts @@ -45,7 +45,6 @@ export type Collected = export const defaultMaxOctets = 64 * 1024 * 1024; type Group = { - octets: number; parts: Map; smsId: string; total: number; @@ -88,18 +87,18 @@ export class Reassembler { private readonly maxOctets: number; private readonly newId: () => string; private readonly onLost: (lost: LostGroup) => void; - private octets = 0; constructor(options: ReassemblerOptions) { + this.maxOctets = options.maxOctets ?? defaultMaxOctets; this.groups = new ExpiringGroups({ max: options.max, + maxWeight: this.maxOctets, now: options.now, onSweep: () => { this.sweep(); }, timeout: options.timeout, }); this.log = options.log; this.max = options.max; - this.maxOctets = options.maxOctets ?? defaultMaxOctets; this.newId = options.newId ?? uuidv7; this.onLost = options.onLost; } @@ -118,25 +117,17 @@ export class Reassembler { if (!this.placeable(concat, existing)) return { kept: false, refusal: 'unplaceable' }; const group = existing ?? this.open(key, concat.total); - const replaced = group.parts.get(concat.part); - const segment = detach(pduObj); - const delta = retainedOctets(segment) - (replaced === undefined ? 0 : retainedOctets(replaced)); - group.parts.set(concat.part, segment); - group.octets += delta; - this.octets += delta; + group.parts.set(concat.part, detach(pduObj)); if (group.parts.size < group.total) { - this.trim(key); - // Its own arrival overran the octet cap, so the peer keeps it rather than being told we did. - if (this.groups.get(key) !== group) return { kept: false, refusal: 'full' }; + if (!this.trim(key, group)) return { kept: false, refusal: 'full' }; return { kept: true, smsId: group.smsId }; } this.groups.delete(key); - this.octets -= group.octets; return { kept: true, @@ -149,14 +140,11 @@ export class Reassembler { for (const [, group] of this.groups.takeAll()) { this.lost(group, 'linkGone'); } - - this.octets = 0; } /** Drops every group past its deadline. Runs before each collect and on its own timer. */ sweep(): void { for (const [, group] of this.groups.takeExpired()) { - this.octets -= group.octets; this.lost(group, 'expired'); } } @@ -189,37 +177,37 @@ export class Reassembler { private open(key: string, total: number): Group { if (this.groups.full) this.dropOldest(); - const group: Group = { octets: 0, parts: new Map(), smsId: this.newId(), total }; + const group: Group = { parts: new Map(), smsId: this.newId(), total }; this.groups.set(key, group); return group; } - /** Drops the oldest groups until the retained payload is back under the octet cap. */ - private trim(current: string): void { - while (this.octets > this.maxOctets) { - const oldest = this.takeOldest(); + /** Drops the oldest groups until the retained payload is back under the octet cap. False if the current one went. */ + private trim(current: string, group: Group): boolean { + let octets = 0; - if (!oldest) return; + for (const part of group.parts.values()) { + octets += retainedOctets(part); + } + + let survived = true; + + for (const [key, oldest] of this.groups.weigh(current, octets)) { + if (key === current) survived = false; // The refused segment is in the group but stays with the peer, so it is none of the loss. - const answered = oldest[0] === current ? oldest[1].parts.size - 1 : oldest[1].parts.size; + const answered = key === current ? oldest.parts.size - 1 : oldest.parts.size; - if (answered > 0) this.lost(oldest[1], 'evicted', answered); + if (answered > 0) this.lost(oldest, 'evicted', answered); } - } - private takeOldest(): [string, Group] | undefined { - const oldest = this.groups.takeOldest(); - - if (oldest) this.octets -= oldest[1].octets; - - return oldest; + return survived; } private dropOldest(): void { - const oldest = this.takeOldest(); + const oldest = this.groups.takeOldest(); if (oldest) this.lost(oldest[1], 'evicted'); } @@ -232,7 +220,7 @@ export class Reassembler { ...lost, max: this.max, maxOctets: this.maxOctets, - octets: this.octets, + octets: this.groups.weight, }); this.onLost(lost); } diff --git a/test/session-extras.test.ts b/test/session-extras.test.ts index 9bf67ed..c8b3dc2 100644 --- a/test/session-extras.test.ts +++ b/test/session-extras.test.ts @@ -2191,6 +2191,8 @@ describe('reassembly bounds', () => { assert.ok(reopened.kept); assert.equal(reopened.whole, undefined); + assert.ok(collect(reassembler, 3, 1, 2).kept); + assert.equal(reassembler.size, 2, 'a segment sent again replaces its octets rather than adding them'); reassembler.clear(); }); diff --git a/todo.md b/todo.md index 3a4009a..364c196 100644 --- a/todo.md +++ b/todo.md @@ -199,12 +199,6 @@ next work ([decision](docs/decisions.md#internals-and-tests)). ### Locality — next, ahead of everything below; 5–6 today, and the gate is 7 -- [ ] **Derive `Reassembler`'s octet total instead of maintaining it at five sites.** `this.octets` - and each `group.octets` must agree, adjusted in `collect`, `trim`, `takeOldest`, `sweep` and - `clear`, and `collect()` discovers its own eviction by re-reading the map by identity. Push the - budget into `ExpiringGroups` as a weighed capacity, and have `trim()` report whether the - current group survived. Named by 6 of 9 readers. - - [ ] **Let the two address arrays size a C-Octet String through `cstring.size()`.** `sizeDestAddresses()` and `sizeUnsuccessSmes()` spell "len + 1" themselves, and each `offset +=` after a write spells it a third time, so `dest_address_array` and `unsuccess_sme_array` each -- 2.52.0 From 7117c91eaba9d808728de3728501f28c87b3a83c Mon Sep 17 00:00:00 2001 From: Lilleman auf Larv Date: Mon, 28 Sep 2026 01:13:51 +0200 Subject: [PATCH 2/4] Assert a replaced held message leaves its octets with it --- test/session-extras.test.ts | 1 + 1 file changed, 1 insertion(+) diff --git a/test/session-extras.test.ts b/test/session-extras.test.ts index c8b3dc2..fb11cad 100644 --- a/test/session-extras.test.ts +++ b/test/session-extras.test.ts @@ -1508,6 +1508,7 @@ describe('held message bounds', () => { held.hold(message(2), 1); 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); -- 2.52.0 From b1a4201a7dcf85c8f2a43278fa1f71c8039c2593 Mon Sep 17 00:00:00 2001 From: Lilleman auf Larv Date: Mon, 28 Sep 2026 01:13:51 +0200 Subject: [PATCH 3/4] Record a group's weight through weigh() alone --- src/expiring-groups.ts | 5 ++--- src/held-messages.ts | 3 ++- 2 files changed, 4 insertions(+), 4 deletions(-) diff --git a/src/expiring-groups.ts b/src/expiring-groups.ts index 80440b6..ee40936 100644 --- a/src/expiring-groups.ts +++ b/src/expiring-groups.ts @@ -54,10 +54,9 @@ export class ExpiringGroups { } /** Starts the group's deadline, and the sweeper if this is the only group held. */ - set(key: string, group: T, weight = 0): void { + set(key: string, group: T): void { this.remove(key); - this.entries.set(key, { deadline: this.now() + this.timeout, group, weight }); - this.total += weight; + this.entries.set(key, { deadline: this.now() + this.timeout, group, weight: 0 }); if (this.sweeper) return; diff --git a/src/held-messages.ts b/src/held-messages.ts index a519d1a..cfacbd9 100644 --- a/src/held-messages.ts +++ b/src/held-messages.ts @@ -106,7 +106,8 @@ export class HeldMessages { this.log.warn('heldMessages - replacing a message on a re-used sequence number', { seqNr: Number(key) }); } - this.held.set(key, pduObjs, pduObjs.reduce((sum, pduObj) => sum + retainedOctets(pduObj), 0)); + this.held.set(key, pduObjs); + this.held.weigh(key, pduObjs.reduce((sum, pduObj) => sum + retainedOctets(pduObj), 0)); return hold; } -- 2.52.0 From acad60c88973f20dfae182da316a7f0ca2a4c321 Mon Sep 17 00:00:00 2001 From: Lilleman auf Larv Date: Mon, 28 Sep 2026 01:16:06 +0200 Subject: [PATCH 4/4] Say what ExpiringGroups leaves to its owners --- src/expiring-groups.ts | 16 +++++----------- 1 file changed, 5 insertions(+), 11 deletions(-) diff --git a/src/expiring-groups.ts b/src/expiring-groups.ts index ee40936..3e276aa 100644 --- a/src/expiring-groups.ts +++ b/src/expiring-groups.ts @@ -1,10 +1,10 @@ export type ExpiringGroupsOptions = { max: number; - /** What the groups may weigh together before weigh() evicts the oldest. */ + /** Enforced by weigh() alone; set() never evicts. */ maxWeight?: number | undefined; /** Injected so expiry can be exercised without a wall clock. */ now?: (() => number) | undefined; - /** Runs on the sweeper's own timer; the owner reports whatever it takes out. */ + /** Must call takeExpired(): the timer itself removes nothing. */ onSweep: () => void; timeout: number; }; @@ -15,10 +15,7 @@ type Entry = { weight: number; }; -/** - * A capped store of groups that expire. Nothing is dropped silently: the owner takes the expired - * and the evicted out itself, so the accounting and the log line stay where the group is understood. - */ +/** Enforces neither max nor timeout itself: owners check full and call takeExpired(); only weigh() evicts. */ export class ExpiringGroups { private readonly entries = new Map>(); private readonly max: number; @@ -53,7 +50,7 @@ export class ExpiringGroups { return this.entries.get(key)?.group; } - /** Starts the group's deadline, and the sweeper if this is the only group held. */ + /** Replacing a key restarts its deadline and zeroes its weight; weigh() it again. */ set(key: string, group: T): void { this.remove(key); this.entries.set(key, { deadline: this.now() + this.timeout, group, weight: 0 }); @@ -69,7 +66,7 @@ export class ExpiringGroups { this.idle(); } - /** Records what a group weighs now, and hands over the oldest groups evicted to bring the total under maxWeight. */ + /** The returned groups are already removed, and may include key itself. */ weigh(key: string, weight: number): [string, T][] { const entry = this.entries.get(key); const taken: [string, T][] = []; @@ -90,7 +87,6 @@ export class ExpiringGroups { return taken; } - /** Removes every group and hands them over, so an owner that must account for them can. */ takeAll(): [string, T][] { const taken: [string, T][] = []; @@ -105,7 +101,6 @@ export class ExpiringGroups { return taken; } - /** Removes every group past its deadline and hands them over. */ takeExpired(): [string, T][] { const now = this.now(); const taken: [string, T][] = []; @@ -122,7 +117,6 @@ export class ExpiringGroups { return taken; } - /** Removes the group held longest and hands it over. Undefined means there was none. */ takeOldest(): [string, T] | undefined { const oldest = this.entries.entries().next(); -- 2.52.0