Derive the retained octets from ExpiringGroups' weight #36
@@ -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
|
||||
|
||||
+47
-4
@@ -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<T> = {
|
||||
deadline: number;
|
||||
group: T;
|
||||
weight: number;
|
||||
};
|
||||
|
||||
/**
|
||||
@@ -19,13 +22,16 @@ type Entry<T> = {
|
||||
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 +45,19 @@ 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. */
|
||||
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<T> {
|
||||
}
|
||||
|
||||
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<T> {
|
||||
}
|
||||
|
||||
this.entries.clear();
|
||||
this.total = 0;
|
||||
this.idle();
|
||||
|
||||
return taken;
|
||||
@@ -81,7 +115,7 @@ export class ExpiringGroups<T> {
|
||||
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<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;
|
||||
|
||||
|
||||
+8
-32
@@ -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,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();
|
||||
}
|
||||
|
||||
+21
-33
@@ -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);
|
||||
}
|
||||
|
||||
@@ -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();
|
||||
});
|
||||
|
||||
@@ -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
|
||||
|
||||
Reference in New Issue
Block a user