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-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
+48 -12
View File
@@ -1,8 +1,10 @@
export type ExpiringGroupsOptions = {
max: number;
/** 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;
};
@@ -10,22 +12,23 @@ export type ExpiringGroupsOptions = {
type Entry<T> = {
deadline: number;
group: T;
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<T> {
private readonly entries = new Map<string, Entry<T>>();
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 +42,18 @@ export class ExpiringGroups<T> {
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. */
/** Replacing a key restarts its deadline and zeroes its weight; weigh() it again. */
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;
@@ -54,11 +62,31 @@ export class ExpiringGroups<T> {
}
delete(key: string): void {
this.entries.delete(key);
this.remove(key);
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][] {
const taken: [string, T][] = [];
@@ -67,12 +95,12 @@ export class ExpiringGroups<T> {
}
this.entries.clear();
this.total = 0;
this.idle();
return taken;
}
/** Removes every group past its deadline and hands them over. */
takeExpired(): [string, T][] {
const now = this.now();
const taken: [string, T][] = [];
@@ -81,7 +109,7 @@ export class ExpiringGroups<T> {
if (entry.deadline > now) continue;
taken.push([key, entry.group]);
this.entries.delete(key);
this.remove(key);
}
this.idle();
@@ -89,7 +117,6 @@ export class ExpiringGroups<T> {
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();
@@ -102,6 +129,15 @@ export class ExpiringGroups<T> {
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;
+9 -32
View File
@@ -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<Held>;
private readonly held: ExpiringGroups<PduObject[]>;
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,12 @@ 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);
this.held.weigh(key, pduObjs.reduce((sum, pduObj) => sum + retainedOctets(pduObj), 0));
return hold;
}
@@ -126,25 +115,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 +145,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();
}
+21 -33
View File
@@ -45,7 +45,6 @@ export type Collected =
export const defaultMaxOctets = 64 * 1024 * 1024;
type Group = {
octets: number;
parts: Map<number, PduObject>;
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<Group>({
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);
}
+3
View File
@@ -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);
@@ -2191,6 +2192,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();
});
-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
- [ ] **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