Compare commits
8 Commits
main
...
9d3dbc522c
| Author | SHA1 | Date | |
|---|---|---|---|
| 9d3dbc522c | |||
| 936102fd04 | |||
| 9df7f05669 | |||
| 245d0e2740 | |||
| 1db63e19ed | |||
| b5e6b03ff8 | |||
| d8c2647e75 | |||
| e52baafd29 |
+32
-3
@@ -25,7 +25,11 @@ type HoldEntry = {
|
|||||||
release: () => void;
|
release: () => void;
|
||||||
};
|
};
|
||||||
|
|
||||||
/** One message offered to the application, held until it is answered or every listener gives up. */
|
/**
|
||||||
|
* One message offered to the application. A drain waits on it until the first of: `answered()`,
|
||||||
|
* every listener that took it rejecting, no listener taking it or one throwing, a later message on
|
||||||
|
* its sequence number, its deadline, or the link going.
|
||||||
|
*/
|
||||||
export class MessageHold {
|
export class MessageHold {
|
||||||
private readonly entry: HoldEntry;
|
private readonly entry: HoldEntry;
|
||||||
private working: number;
|
private working: number;
|
||||||
@@ -52,7 +56,7 @@ export class MessageHold {
|
|||||||
if (this.working <= 0) this.answered();
|
if (this.working <= 0) this.answered();
|
||||||
}
|
}
|
||||||
|
|
||||||
/** At once, for a message nobody took: that is not work a shutdown can wait for. */
|
/** At once, for a message nobody took or a listener threw on: that is not work a shutdown can wait for. */
|
||||||
release(): void {
|
release(): void {
|
||||||
this.entry.release();
|
this.entry.release();
|
||||||
}
|
}
|
||||||
@@ -64,6 +68,8 @@ export class HeldMessages {
|
|||||||
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;
|
||||||
|
/** A rejecting listener hands the message back as an `unknown`, so its hold is found by identity. */
|
||||||
|
private readonly offered = new WeakMap<object, MessageHold>();
|
||||||
|
|
||||||
constructor(options: HeldMessagesOptions) {
|
constructor(options: HeldMessagesOptions) {
|
||||||
this.held = new ExpiringGroups({
|
this.held = new ExpiringGroups({
|
||||||
@@ -91,7 +97,7 @@ export class HeldMessages {
|
|||||||
return this.held.full || this.held.weight >= this.maxOctets;
|
return this.held.full || this.held.weight >= this.maxOctets;
|
||||||
}
|
}
|
||||||
|
|
||||||
hold(pduObjs: PduObject[], listeners: number): MessageHold {
|
private hold(pduObjs: PduObject[], listeners: number): MessageHold {
|
||||||
const hold = new MessageHold({
|
const hold = new MessageHold({
|
||||||
isHeld: () => this.has(pduObjs),
|
isHeld: () => this.has(pduObjs),
|
||||||
release: () => { this.release(pduObjs); },
|
release: () => { this.release(pduObjs); },
|
||||||
@@ -112,6 +118,29 @@ export class HeldMessages {
|
|||||||
return hold;
|
return hold;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
offer<T extends object>(
|
||||||
|
pduObjs: PduObject[],
|
||||||
|
listeners: number,
|
||||||
|
build: (hold: MessageHold) => T,
|
||||||
|
emit: (message: T) => boolean,
|
||||||
|
): MessageHold {
|
||||||
|
const hold = this.hold(pduObjs, listeners);
|
||||||
|
const message = build(hold);
|
||||||
|
|
||||||
|
this.offered.set(message, hold);
|
||||||
|
|
||||||
|
if (!emit(message)) hold.release();
|
||||||
|
|
||||||
|
return hold;
|
||||||
|
}
|
||||||
|
|
||||||
|
/** One listener gave up on a message; the last one to do so is what releases it. */
|
||||||
|
listenerRejected(message: unknown): void {
|
||||||
|
if (typeof message !== 'object' || message === null) return;
|
||||||
|
|
||||||
|
this.offered.get(message)?.listenerGaveUp();
|
||||||
|
}
|
||||||
|
|
||||||
private has(pduObjs: PduObject[]): boolean {
|
private has(pduObjs: PduObject[]): boolean {
|
||||||
const key = keyOf(pduObjs);
|
const key = keyOf(pduObjs);
|
||||||
|
|
||||||
|
|||||||
@@ -10,7 +10,6 @@ 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';
|
||||||
@@ -85,8 +84,6 @@ export class IncomingRequests {
|
|||||||
private readonly deps: IncomingDeps;
|
private readonly deps: IncomingDeps;
|
||||||
private readonly dlrMerger: DlrMerger;
|
private readonly dlrMerger: DlrMerger;
|
||||||
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;
|
||||||
@@ -172,11 +169,8 @@ export class IncomingRequests {
|
|||||||
this.reassembler.clear();
|
this.reassembler.clear();
|
||||||
}
|
}
|
||||||
|
|
||||||
/** One `sms` listener gave up on a message; the last one to do so is what releases the hold. */
|
|
||||||
listenerRejected(sms: unknown): void {
|
listenerRejected(sms: unknown): void {
|
||||||
if (typeof sms !== 'object' || sms === null) return;
|
this.held.listenerRejected(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. */
|
||||||
@@ -310,9 +304,8 @@ export class IncomingRequests {
|
|||||||
if (!first) return;
|
if (!first) return;
|
||||||
|
|
||||||
const generation = this.linkGeneration;
|
const generation = this.linkGeneration;
|
||||||
const hold = this.held.hold(pduObjs, this.deps.smsListeners());
|
|
||||||
|
|
||||||
const sms = this.deps.createSms({
|
this.held.offer(pduObjs, this.deps.smsListeners(), hold => this.deps.createSms({
|
||||||
answeredAs,
|
answeredAs,
|
||||||
from: paramText(first.params.source_addr),
|
from: paramText(first.params.source_addr),
|
||||||
message: decodeSegments(pduObjs),
|
message: decodeSegments(pduObjs),
|
||||||
@@ -325,10 +318,6 @@ export class IncomingRequests {
|
|||||||
lostLink: () => this.linkGeneration !== generation,
|
lostLink: () => this.linkGeneration !== generation,
|
||||||
onAnswered: () => { hold.answered(); },
|
onAnswered: () => { hold.answered(); },
|
||||||
send: input => (hold.isHeld() ? this.deps.sendPastDrain(input) : this.deps.send(input)),
|
send: input => (hold.isHeld() ? this.deps.sendPastDrain(input) : this.deps.send(input)),
|
||||||
});
|
}), sms => this.deps.offerSms(sms));
|
||||||
|
|
||||||
this.holds.set(sms, hold);
|
|
||||||
|
|
||||||
if (!this.deps.offerSms(sms)) hold.release();
|
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
+17
-11
@@ -5,6 +5,7 @@ import type { Collected, LostGroup } from '../src/reassembly.ts';
|
|||||||
import type { Dlr } from '../src/dlr.ts';
|
import type { Dlr } from '../src/dlr.ts';
|
||||||
import type { ErrorName } from '../src/defs/errors.ts';
|
import type { ErrorName } from '../src/defs/errors.ts';
|
||||||
import type { IncomingDeps } from '../src/incoming-requests.ts';
|
import type { IncomingDeps } from '../src/incoming-requests.ts';
|
||||||
|
import type { MessageHold } from '../src/held-messages.ts';
|
||||||
import type { MessageState } from '../src/defs/constants.ts';
|
import type { MessageState } from '../src/defs/constants.ts';
|
||||||
import type { MessageDlr } from '../src/session.ts';
|
import type { MessageDlr } from '../src/session.ts';
|
||||||
import type { PduObject, PduObjectInput } from '../src/pdu.ts';
|
import type { PduObject, PduObjectInput } from '../src/pdu.ts';
|
||||||
@@ -1518,17 +1519,22 @@ describe('held message bounds', () => {
|
|||||||
return [submitPdu(seqNr)];
|
return [submitPdu(seqNr)];
|
||||||
}
|
}
|
||||||
|
|
||||||
|
function offer(held: HeldMessages, seqNr: number): MessageHold {
|
||||||
|
return held.offer(message(seqNr), 1, () => ({}), () => true);
|
||||||
|
}
|
||||||
|
|
||||||
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 = held.hold(message(1), 1);
|
const first = offer(held, 1);
|
||||||
|
const replaced = offer(held, 2);
|
||||||
|
|
||||||
held.hold(message(2), 1);
|
offer(held, 2);
|
||||||
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.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);
|
||||||
|
assert.equal(replaced.isHeld(), false);
|
||||||
|
|
||||||
held.clear();
|
held.clear();
|
||||||
});
|
});
|
||||||
@@ -1537,22 +1543,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 = held.hold(message(1), 1);
|
const answered = offer(held, 1);
|
||||||
|
|
||||||
assert.equal(held.full(), false);
|
assert.equal(held.full(), false);
|
||||||
held.hold(message(2), 1);
|
offer(held, 2);
|
||||||
assert.equal(held.full(), true);
|
assert.equal(held.full(), true);
|
||||||
|
|
||||||
answered.release();
|
answered.release();
|
||||||
assert.equal(held.full(), false, 'after a release');
|
assert.equal(held.full(), false, 'after a release');
|
||||||
held.hold(message(3), 1);
|
offer(held, 3);
|
||||||
|
|
||||||
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), 1);
|
offer(held, 4);
|
||||||
held.hold(message(5), 1);
|
offer(held, 5);
|
||||||
|
|
||||||
held.clear();
|
held.clear();
|
||||||
assert.equal(held.full(), false, 'after a clear');
|
assert.equal(held.full(), false, 'after a clear');
|
||||||
@@ -1641,11 +1647,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), 1);
|
offer(held, 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), 1);
|
offer(held, 2);
|
||||||
|
|
||||||
assert.equal(held.size, 1);
|
assert.equal(held.size, 1);
|
||||||
|
|
||||||
@@ -1657,7 +1663,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), 1);
|
offer(held, 1);
|
||||||
|
|
||||||
const waiting = held.idle(1000, undefined);
|
const waiting = held.idle(1000, undefined);
|
||||||
|
|
||||||
|
|||||||
@@ -204,9 +204,6 @@ and 5. Every seat ranked the session's lifecycle hardest and least wanted to mod
|
|||||||
|
|
||||||
- [ ] **Lift Locality to 7, and confirm it with a scoring run.** A run reading 7.0 or above also
|
- [ ] **Lift Locality to 7, and confirm it with a scoring run.** A run reading 7.0 or above also
|
||||||
retires the #30 decision. The sub-items are what the 2026-09-28 run named, most seats first.
|
retires the #30 decision. The sub-items are what the 2026-09-28 run named, most seats first.
|
||||||
- [ ] **Give the held-message flow one place a reader can follow it.** Whether a drain still waits
|
|
||||||
on a message is spread over `emitSms()`, `MessageHold`, the session's rejection route and
|
|
||||||
`Sms.isHeld()`. Four seats.
|
|
||||||
- [ ] **Shrink the `IncomingDeps` closure bag.** 16 lambdas, six of them repeated in
|
- [ ] **Shrink the `IncomingDeps` closure bag.** 16 lambdas, six of them repeated in
|
||||||
`SmsHandlers`, which makes every inbound call path indirect. Two seats.
|
`SmsHandlers`, which makes every inbound call path indirect. Two seats.
|
||||||
|
|
||||||
@@ -268,6 +265,11 @@ and 5. Every seat ranked the session's lifecycle hardest and least wanted to mod
|
|||||||
so `"\x1B("` is detected as GSM, goes out as 0x1B 0x28 and arrives as `{`. From the
|
so `"\x1B("` is detected as GSM, goes out as 0x1B 0x28 and arrives as `{`. From the
|
||||||
2026-09-28 scoring run.
|
2026-09-28 scoring run.
|
||||||
|
|
||||||
|
- [ ] **Count an `sms` listener that throws as one giving up, as a rejection already is.** A
|
||||||
|
synchronous throw makes `Session.emit` return false, and `HeldMessages.offer()` then releases
|
||||||
|
the hold at once, so a drain stops waiting on an async listener still answering beside it.
|
||||||
|
Goal 2. From the stability review of #45.
|
||||||
|
|
||||||
### Throughput — goal 6, and the default window is where we are slowest
|
### Throughput — goal 6, and the default window is where we are slowest
|
||||||
|
|
||||||
- [ ] **Close the gap to jsmpp at `maxOutstanding: 10`.** Measured 2026-09-20 against the same sink,
|
- [ ] **Close the gap to jsmpp at `maxOutstanding: 10`.** Measured 2026-09-20 against the same sink,
|
||||||
|
|||||||
Reference in New Issue
Block a user