Give the held-message protocol one home in MessageHold #35
@@ -52,7 +52,7 @@ src/
|
|||||||
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, expiring store both of those share
|
||||||
held-messages.ts HeldMessages: capped, expiring messages the application has not answered
|
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
|
||||||
link-gate.ts LinkGate: where a request with no link to go out on waits for the next one
|
link-gate.ts LinkGate: where a request with no link to go out on waits for the next one
|
||||||
|
|||||||
+48
-5
@@ -25,6 +25,44 @@ type Held = {
|
|||||||
pduObjs: PduObject[];
|
pduObjs: PduObject[];
|
||||||
};
|
};
|
||||||
|
|
||||||
|
type HoldEntry = {
|
||||||
|
isHeld: () => boolean;
|
||||||
|
release: () => void;
|
||||||
|
};
|
||||||
|
|
||||||
|
/** One message offered to the application, held until it is answered or every listener gives up. */
|
||||||
|
export class MessageHold {
|
||||||
|
private readonly entry: HoldEntry;
|
||||||
|
private working: number;
|
||||||
|
|
||||||
|
constructor(entry: HoldEntry, listeners: number) {
|
||||||
|
this.entry = entry;
|
||||||
|
this.working = listeners;
|
||||||
|
}
|
||||||
|
|
||||||
|
/** Whether a drain is still waiting for this message to be answered. */
|
||||||
|
isHeld(): boolean {
|
||||||
|
return this.entry.isHeld();
|
||||||
|
}
|
||||||
|
|
||||||
|
/** A turn later, so a listener sending its receipt straight after the response still holds. */
|
||||||
|
answered(): void {
|
||||||
|
setImmediate(() => { this.entry.release(); });
|
||||||
|
}
|
||||||
|
|
||||||
|
/** A rejection leaves the other listeners running, so only the last one to fail gives the message up. */
|
||||||
|
listenerGaveUp(): void {
|
||||||
|
this.working--;
|
||||||
|
|
||||||
|
if (this.working <= 0) this.answered();
|
||||||
|
}
|
||||||
|
|
||||||
|
/** At once, for a message nobody took: that is not work a shutdown can wait for. */
|
||||||
|
release(): void {
|
||||||
|
this.entry.release();
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
/** 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<Held>;
|
||||||
@@ -59,10 +97,14 @@ export class HeldMessages {
|
|||||||
return this.held.full || this.octets >= this.maxOctets;
|
return this.held.full || this.octets >= this.maxOctets;
|
||||||
}
|
}
|
||||||
|
|
||||||
hold(pduObjs: PduObject[]): void {
|
hold(pduObjs: PduObject[], listeners: number): MessageHold {
|
||||||
|
const hold = new MessageHold({
|
||||||
|
isHeld: () => this.has(pduObjs),
|
||||||
|
release: () => { this.release(pduObjs); },
|
||||||
|
}, listeners);
|
||||||
const key = keyOf(pduObjs);
|
const key = keyOf(pduObjs);
|
||||||
|
|
||||||
if (key === undefined) return;
|
if (key === undefined) return hold;
|
||||||
|
|
||||||
this.sweep();
|
this.sweep();
|
||||||
|
|
||||||
@@ -77,16 +119,17 @@ export class HeldMessages {
|
|||||||
|
|
||||||
this.held.set(key, { octets, pduObjs });
|
this.held.set(key, { octets, pduObjs });
|
||||||
this.octets += octets;
|
this.octets += octets;
|
||||||
|
|
||||||
|
return hold;
|
||||||
}
|
}
|
||||||
|
|
||||||
/** Whether a drain is still waiting for this message to be answered. */
|
private has(pduObjs: PduObject[]): boolean {
|
||||||
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 === pduObjs;
|
||||||
}
|
}
|
||||||
|
|
||||||
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.
|
||||||
|
|||||||
@@ -10,6 +10,7 @@ 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';
|
||||||
@@ -83,9 +84,9 @@ export type IncomingRequestsOptions = {
|
|||||||
export class IncomingRequests {
|
export class IncomingRequests {
|
||||||
private readonly deps: IncomingDeps;
|
private readonly deps: IncomingDeps;
|
||||||
private readonly dlrMerger: DlrMerger;
|
private readonly dlrMerger: DlrMerger;
|
||||||
/** The rejection handler is handed the Sms back as an `unknown`, so its hold is found by identity. */
|
|
||||||
private readonly emitted = new WeakMap<object, () => void>();
|
|
||||||
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;
|
||||||
@@ -175,7 +176,7 @@ export class IncomingRequests {
|
|||||||
listenerRejected(sms: unknown): void {
|
listenerRejected(sms: unknown): void {
|
||||||
if (typeof sms !== 'object' || sms === null) return;
|
if (typeof sms !== 'object' || sms === null) return;
|
||||||
|
|
||||||
this.emitted.get(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. */
|
||||||
@@ -309,8 +310,7 @@ export class IncomingRequests {
|
|||||||
if (!first) return;
|
if (!first) return;
|
||||||
|
|
||||||
const generation = this.linkGeneration;
|
const generation = this.linkGeneration;
|
||||||
// A turn later, so a listener sending its receipt straight after the response still holds.
|
const hold = this.held.hold(pduObjs, this.deps.smsListeners());
|
||||||
const release = (): void => { setImmediate(() => { this.held.release(pduObjs); }); };
|
|
||||||
|
|
||||||
const sms = this.deps.createSms({
|
const sms = this.deps.createSms({
|
||||||
answeredAs,
|
answeredAs,
|
||||||
@@ -323,22 +323,12 @@ export class IncomingRequests {
|
|||||||
answer: (pduObj, status, params) => this.deps.answer(pduObj, status, params),
|
answer: (pduObj, status, params) => this.deps.answer(pduObj, status, params),
|
||||||
bindAllows: cmdName => this.deps.bindAllows(cmdName),
|
bindAllows: cmdName => this.deps.bindAllows(cmdName),
|
||||||
lostLink: () => this.linkGeneration !== generation,
|
lostLink: () => this.linkGeneration !== generation,
|
||||||
onAnswered: release,
|
onAnswered: () => { hold.answered(); },
|
||||||
// Past the refusal only while a drain is still waiting for this message; an ordinary send after.
|
send: input => (hold.isHeld() ? this.deps.sendPastDrain(input) : this.deps.send(input)),
|
||||||
send: input => (this.held.has(pduObjs) ? this.deps.sendPastDrain(input) : this.deps.send(input)),
|
|
||||||
});
|
});
|
||||||
|
|
||||||
// A rejection leaves the other listeners running, so only the last one to fail gives the message up.
|
this.holds.set(sms, hold);
|
||||||
let working = this.deps.smsListeners();
|
|
||||||
|
|
||||||
this.held.hold(pduObjs);
|
if (!this.deps.offerSms(sms)) hold.release();
|
||||||
this.emitted.set(sms, () => {
|
|
||||||
working--;
|
|
||||||
|
|
||||||
if (working <= 0) release();
|
|
||||||
});
|
|
||||||
|
|
||||||
// A message nobody took is not work a shutdown can wait for.
|
|
||||||
if (!this.deps.offerSms(sms)) this.held.release(pduObjs);
|
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
+13
-15
@@ -1502,15 +1502,14 @@ describe('held message bounds', () => {
|
|||||||
|
|
||||||
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 = message(1);
|
const first = held.hold(message(1), 1);
|
||||||
|
|
||||||
held.hold(first);
|
held.hold(message(2), 1);
|
||||||
held.hold(message(2));
|
held.hold(message(2), 1);
|
||||||
held.hold(message(2));
|
|
||||||
|
|
||||||
assert.equal(held.size, 2);
|
assert.equal(held.size, 2);
|
||||||
assert.equal(held.full(), true);
|
assert.equal(held.full(), true);
|
||||||
assert.equal(held.has(first), true);
|
assert.equal(first.isHeld(), true);
|
||||||
|
|
||||||
held.clear();
|
held.clear();
|
||||||
});
|
});
|
||||||
@@ -1519,23 +1518,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 = message(1);
|
const answered = held.hold(message(1), 1);
|
||||||
|
|
||||||
held.hold(answered);
|
|
||||||
assert.equal(held.full(), false);
|
assert.equal(held.full(), false);
|
||||||
held.hold(message(2));
|
held.hold(message(2), 1);
|
||||||
assert.equal(held.full(), true);
|
assert.equal(held.full(), true);
|
||||||
|
|
||||||
held.release(answered);
|
answered.release();
|
||||||
assert.equal(held.full(), false, 'after a release');
|
assert.equal(held.full(), false, 'after a release');
|
||||||
held.hold(message(3));
|
held.hold(message(3), 1);
|
||||||
|
|
||||||
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));
|
held.hold(message(4), 1);
|
||||||
held.hold(message(5));
|
held.hold(message(5), 1);
|
||||||
|
|
||||||
held.clear();
|
held.clear();
|
||||||
assert.equal(held.full(), false, 'after a clear');
|
assert.equal(held.full(), false, 'after a clear');
|
||||||
@@ -1624,11 +1622,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));
|
held.hold(message(1), 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));
|
held.hold(message(2), 1);
|
||||||
|
|
||||||
assert.equal(held.size, 1);
|
assert.equal(held.size, 1);
|
||||||
|
|
||||||
@@ -1640,7 +1638,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));
|
held.hold(message(1), 1);
|
||||||
|
|
||||||
const waiting = held.idle(1000, undefined);
|
const waiting = held.idle(1000, undefined);
|
||||||
|
|
||||||
|
|||||||
@@ -199,15 +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
|
||||||
|
|
||||||
- [ ] **Give the held-message protocol one name and one home.** `emitSms()` is the unit 8 of 9
|
|
||||||
readers named and 4 would least want to modify, and every one proposed the same fix. It runs
|
|
||||||
five mechanisms in one scope: a hold keyed by array identity, a `working` counter seeded from
|
|
||||||
`listenerCount('sms')`, a `WeakMap` keyed by the `Sms` object, a `setImmediate`-deferred
|
|
||||||
release, and a captured `linkGeneration` — with the counter decremented from `session.ts`'s
|
|
||||||
`captureRejectionSymbol` in another file. A `MessageHold` owning `hold/release/listenerGaveUp`
|
|
||||||
collapses three files into one readable object. Every way of getting it wrong is silent: a hung
|
|
||||||
shutdown, or a receipt refused.
|
|
||||||
|
|
||||||
- [ ] **Derive `Reassembler`'s octet total instead of maintaining it at five sites.** `this.octets`
|
- [ ] **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
|
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
|
`clear`, and `collect()` discovers its own eviction by re-reading the map by identity. Push the
|
||||||
|
|||||||
Reference in New Issue
Block a user