Derive the retained octets from ExpiringGroups' weight #36

Merged
lilleman merged 4 commits from reassembly-octets into main 2026-09-28 01:20:22 +02:00
6 changed files with 82 additions and 84 deletions
+1 -1
View File
@@ -51,7 +51,7 @@ src/
dlr.ts Delivery receipts: text and TLV parsing, receipt status codes dlr.ts Delivery receipts: text and TLV parsing, receipt status codes
dlr-merger.ts DlrMerger: per-segment receipts counted into one MessageDlr 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 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 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 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 incoming-requests.ts Every request the peer sends: messages, receipts, links, unknown commands
+48 -12
View File
@@ -1,8 +1,10 @@
export type ExpiringGroupsOptions = { export type ExpiringGroupsOptions = {
max: number; max: number;
/** Enforced by weigh() alone; set() never evicts. */
maxWeight?: number | undefined;
/** Injected so expiry can be exercised without a wall clock. */ /** Injected so expiry can be exercised without a wall clock. */
now?: (() => number) | undefined; 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; onSweep: () => void;
timeout: number; timeout: number;
}; };
@@ -10,22 +12,23 @@ export type ExpiringGroupsOptions = {
type Entry<T> = { type Entry<T> = {
deadline: number; deadline: number;
group: T; group: T;
weight: number;
}; };
/** /** Enforces neither max nor timeout itself: owners check full and call takeExpired(); only weigh() evicts. */
* 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.
*/
export class ExpiringGroups<T> { export class ExpiringGroups<T> {
private readonly entries = new Map<string, Entry<T>>(); private readonly entries = new Map<string, Entry<T>>();
private readonly max: number; private readonly max: number;
private readonly maxWeight: number;
private readonly now: () => number; private readonly now: () => number;
private readonly onSweep: () => void; private readonly onSweep: () => void;
private readonly timeout: number; private readonly timeout: number;
private sweeper: NodeJS.Timeout | undefined; private sweeper: NodeJS.Timeout | undefined;
private total = 0;
constructor(options: ExpiringGroupsOptions) { constructor(options: ExpiringGroupsOptions) {
this.max = options.max; this.max = options.max;
this.maxWeight = options.maxWeight ?? Infinity;
this.now = options.now ?? Date.now; this.now = options.now ?? Date.now;
this.onSweep = options.onSweep; this.onSweep = options.onSweep;
this.timeout = options.timeout; this.timeout = options.timeout;
@@ -39,13 +42,18 @@ export class ExpiringGroups<T> {
return this.entries.size; return this.entries.size;
} }
get weight(): number {
return this.total;
}
get(key: string): T | undefined { get(key: string): T | undefined {
return this.entries.get(key)?.group; 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 { set(key: string, group: T): void {
this.entries.set(key, { deadline: this.now() + this.timeout, group }); this.remove(key);
this.entries.set(key, { deadline: this.now() + this.timeout, group, weight: 0 });
if (this.sweeper) return; if (this.sweeper) return;
@@ -54,11 +62,31 @@ export class ExpiringGroups<T> {
} }
delete(key: string): void { delete(key: string): void {
this.entries.delete(key); this.remove(key);
this.idle(); this.idle();
} }
/** Removes every group and hands them over, so an owner that must account for them can. */ /** 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][] = [];
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;
}
takeAll(): [string, T][] { takeAll(): [string, T][] {
const taken: [string, T][] = []; const taken: [string, T][] = [];
@@ -67,12 +95,12 @@ export class ExpiringGroups<T> {
} }
this.entries.clear(); this.entries.clear();
this.total = 0;
this.idle(); this.idle();
return taken; return taken;
} }
/** Removes every group past its deadline and hands them over. */
takeExpired(): [string, T][] { takeExpired(): [string, T][] {
const now = this.now(); const now = this.now();
const taken: [string, T][] = []; const taken: [string, T][] = [];
@@ -81,7 +109,7 @@ export class ExpiringGroups<T> {
if (entry.deadline > now) continue; if (entry.deadline > now) continue;
taken.push([key, entry.group]); taken.push([key, entry.group]);
this.entries.delete(key); this.remove(key);
} }
this.idle(); this.idle();
@@ -89,7 +117,6 @@ export class ExpiringGroups<T> {
return taken; return taken;
} }
/** Removes the group held longest and hands it over. Undefined means there was none. */
takeOldest(): [string, T] | undefined { takeOldest(): [string, T] | undefined {
const oldest = this.entries.entries().next(); const oldest = this.entries.entries().next();
@@ -102,6 +129,15 @@ export class ExpiringGroups<T> {
return [key, entry.group]; 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 { private idle(): void {
if (!this.sweeper || this.entries.size > 0) return; if (!this.sweeper || this.entries.size > 0) return;
+9 -32
View File
@@ -20,11 +20,6 @@ function keyOf(pduObjs: PduObject[]): string | undefined {
return first ? String(first.seqNr) : undefined; return first ? String(first.seqNr) : undefined;
} }
type Held = {
octets: number;
pduObjs: PduObject[];
};
type HoldEntry = { type HoldEntry = {
isHeld: () => boolean; isHeld: () => boolean;
release: () => void; 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. */ /** The messages handed to the application that it has not answered yet, held by their segments. */
export class HeldMessages { export class HeldMessages {
private readonly held: ExpiringGroups<Held>; private readonly held: ExpiringGroups<PduObject[]>;
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;
private octets = 0;
constructor(options: HeldMessagesOptions) { constructor(options: HeldMessagesOptions) {
this.held = new ExpiringGroups({ this.held = new ExpiringGroups({
@@ -83,7 +77,7 @@ export class HeldMessages {
} }
get octetsHeld(): number { get octetsHeld(): number {
return this.octets; return this.held.weight;
} }
get size(): number { get size(): number {
@@ -94,7 +88,7 @@ export class HeldMessages {
full(): boolean { full(): boolean {
this.sweep(); 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 { hold(pduObjs: PduObject[], listeners: number): MessageHold {
@@ -108,17 +102,12 @@ export class HeldMessages {
this.sweep(); this.sweep();
const replaced = this.held.get(key); if (this.held.get(key)) {
if (replaced) {
this.log.warn('heldMessages - replacing a message on a re-used sequence number', { seqNr: Number(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, pduObjs);
this.held.weigh(key, pduObjs.reduce((sum, pduObj) => sum + retainedOctets(pduObj), 0));
this.held.set(key, { octets, pduObjs });
this.octets += octets;
return hold; return hold;
} }
@@ -126,25 +115,22 @@ export class HeldMessages {
private has(pduObjs: PduObject[]): boolean { private has(pduObjs: PduObject[]): boolean {
const key = keyOf(pduObjs); 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 { private release(pduObjs: PduObject[]): void {
const key = keyOf(pduObjs); const key = keyOf(pduObjs);
// Identity, not the key: a wrapped sequence number must not release someone else's message. // 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.held.delete(key);
this.delete(key, held);
this.settle(); this.settle();
} }
/** Drops every message: their segments went with the link, so no answer of ours correlates now. */ /** Drops every message: their segments went with the link, so no answer of ours correlates now. */
clear(): void { clear(): void {
this.held.takeAll(); this.held.takeAll();
this.octets = 0;
this.idleWaiters.settle(); this.idleWaiters.settle();
} }
@@ -159,21 +145,12 @@ export class HeldMessages {
if (expired.length === 0) return; if (expired.length === 0) return;
for (const [, held] of expired) {
this.octets -= held.octets;
}
this.log.warn('heldMessages - messages the application never answered', { this.log.warn('heldMessages - messages the application never answered', {
messages: expired.length, messages: expired.length,
}); });
this.settle(); this.settle();
} }
private delete(key: string, held: Held): void {
this.held.delete(key);
this.octets -= held.octets;
}
private settle(): void { private settle(): void {
if (this.held.size === 0) this.idleWaiters.settle(); if (this.held.size === 0) this.idleWaiters.settle();
} }
+21 -33
View File
@@ -45,7 +45,6 @@ export type Collected =
export const defaultMaxOctets = 64 * 1024 * 1024; export const defaultMaxOctets = 64 * 1024 * 1024;
type Group = { type Group = {
octets: number;
parts: Map<number, PduObject>; parts: Map<number, PduObject>;
smsId: string; smsId: string;
total: number; total: number;
@@ -88,18 +87,18 @@ export class Reassembler {
private readonly maxOctets: number; private readonly maxOctets: number;
private readonly newId: () => string; private readonly newId: () => string;
private readonly onLost: (lost: LostGroup) => void; private readonly onLost: (lost: LostGroup) => void;
private octets = 0;
constructor(options: ReassemblerOptions) { constructor(options: ReassemblerOptions) {
this.maxOctets = options.maxOctets ?? defaultMaxOctets;
this.groups = new ExpiringGroups<Group>({ this.groups = new ExpiringGroups<Group>({
max: options.max, max: options.max,
maxWeight: this.maxOctets,
now: options.now, now: options.now,
onSweep: () => { this.sweep(); }, onSweep: () => { this.sweep(); },
timeout: options.timeout, timeout: options.timeout,
}); });
this.log = options.log; this.log = options.log;
this.max = options.max; this.max = options.max;
this.maxOctets = options.maxOctets ?? defaultMaxOctets;
this.newId = options.newId ?? uuidv7; this.newId = options.newId ?? uuidv7;
this.onLost = options.onLost; this.onLost = options.onLost;
} }
@@ -118,25 +117,17 @@ export class Reassembler {
if (!this.placeable(concat, existing)) return { kept: false, refusal: 'unplaceable' }; if (!this.placeable(concat, existing)) return { kept: false, refusal: 'unplaceable' };
const group = existing ?? this.open(key, concat.total); 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.parts.set(concat.part, detach(pduObj));
group.octets += delta;
this.octets += delta;
if (group.parts.size < group.total) { 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. // 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 }; return { kept: true, smsId: group.smsId };
} }
this.groups.delete(key); this.groups.delete(key);
this.octets -= group.octets;
return { return {
kept: true, kept: true,
@@ -149,14 +140,11 @@ export class Reassembler {
for (const [, group] of this.groups.takeAll()) { for (const [, group] of this.groups.takeAll()) {
this.lost(group, 'linkGone'); this.lost(group, 'linkGone');
} }
this.octets = 0;
} }
/** Drops every group past its deadline. Runs before each collect and on its own timer. */ /** Drops every group past its deadline. Runs before each collect and on its own timer. */
sweep(): void { sweep(): void {
for (const [, group] of this.groups.takeExpired()) { for (const [, group] of this.groups.takeExpired()) {
this.octets -= group.octets;
this.lost(group, 'expired'); this.lost(group, 'expired');
} }
} }
@@ -189,37 +177,37 @@ export class Reassembler {
private open(key: string, total: number): Group { private open(key: string, total: number): Group {
if (this.groups.full) this.dropOldest(); 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); this.groups.set(key, group);
return group; return group;
} }
/** Drops the oldest groups until the retained payload is back under the octet cap. */ /** Drops the oldest groups until the retained payload is back under the octet cap. False if the current one went. */
private trim(current: string): void { private trim(current: string, group: Group): boolean {
while (this.octets > this.maxOctets) { let octets = 0;
const oldest = this.takeOldest();
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. // 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 { return survived;
const oldest = this.groups.takeOldest();
if (oldest) this.octets -= oldest[1].octets;
return oldest;
} }
private dropOldest(): void { private dropOldest(): void {
const oldest = this.takeOldest(); const oldest = this.groups.takeOldest();
if (oldest) this.lost(oldest[1], 'evicted'); if (oldest) this.lost(oldest[1], 'evicted');
} }
@@ -232,7 +220,7 @@ export class Reassembler {
...lost, ...lost,
max: this.max, max: this.max,
maxOctets: this.maxOctets, maxOctets: this.maxOctets,
octets: this.octets, octets: this.groups.weight,
}); });
this.onLost(lost); this.onLost(lost);
} }
+3
View File
@@ -1508,6 +1508,7 @@ describe('held message bounds', () => {
held.hold(message(2), 1); 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.full(), true); assert.equal(held.full(), true);
assert.equal(first.isHeld(), true); assert.equal(first.isHeld(), true);
@@ -2191,6 +2192,8 @@ describe('reassembly bounds', () => {
assert.ok(reopened.kept); assert.ok(reopened.kept);
assert.equal(reopened.whole, undefined); 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(); reassembler.clear();
}); });
-6
View File
@@ -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 ### 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()`.** - [ ] **Let the two address arrays size a C-Octet String through `cstring.size()`.**
`sizeDestAddresses()` and `sizeUnsuccessSmes()` spell "len + 1" themselves, and each `offset +=` `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 after a write spells it a third time, so `dest_address_array` and `unsuccess_sme_array` each