From 855ccbe75f7cb4195c215784120d865071abd611 Mon Sep 17 00:00:00 2001 From: Lilleman auf Larv Date: Sat, 26 Sep 2026 13:18:56 +0200 Subject: [PATCH 1/5] Refuse a message past the held-message bound with ESME_RTHROTTLED, and answer a full reassembly buffer the same way --- AGENTS.md | 5 +- CHANGELOG.md | 9 ++-- README.md | 11 ++-- docs/decisions.md | 27 +++++++--- src/held-messages.ts | 33 +++--------- src/incoming-requests.ts | 18 ++++++- test/session-extras.test.ts | 101 +++++++++++++++++++----------------- test/session.test.ts | 2 +- todo.md | 11 ---- 9 files changed, 113 insertions(+), 104 deletions(-) diff --git a/AGENTS.md b/AGENTS.md index ce2ac3b..cc83610 100644 --- a/AGENTS.md +++ b/AGENTS.md @@ -302,7 +302,10 @@ the file. - The drain waits on the messages the application holds, and `sendResp()` is what says it is done with one. - The drain's wait on the application ignores `shutdownTimeout: 0`. -- What the application holds unanswered is capped on constants. +- What the application holds unanswered is capped on constants, and a message past the cap is + refused. +- A store at its bound answers `ESME_RTHROTTLED` to a `submit_sm` and `ESME_RX_T_APPN` to a + `deliver_sm`. - A reconnect keeps the delivery-receipt merges; everything else the link held is dropped. - A message id base is merged at most once. - A send that never reached the socket waits for the next link; one that did is counted, not resent. diff --git a/CHANGELOG.md b/CHANGELOG.md index 193be3f..2cb73c8 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -40,9 +40,12 @@ count as next to nothing, so a peer could hold far more than the cap. **Raise a `maxOctets` you tuned low**: it now holds several times fewer segments, and an incomplete message evicted over the cap is lost, since its segments were already answered. -- Unanswered `sms` messages are capped at 64 MiB per session by the `maxOctets` charge, beside the - 1000-message cap: the oldest is dropped with a warning, as over the count. The library used to - hold up to 1000 messages of any size for an application that answered none of them. +- A message arriving while the application holds 1000 unanswered, or 64 MiB of them by the + `maxOctets` charge, is refused with `ESME_RTHROTTLED` (`ESME_RX_T_APPN` on a `deliver_sm`), so the + peer keeps it and retries. The oldest used to be dropped to make room, which freed nothing while + the application still held it and let `close()` stop waiting for a message the peer was owed. +- A segment the reassembly buffer has no room for is refused with `ESME_RTHROTTLED`, where it was + `ESME_RMSGQFUL`. - `server()` refuses a `maxOctets` below 1 or not a whole number, `Infinity` included, like its other limits. `server({ maxOctets: 0 })` used to start and then refuse every multipart message. - `callback_num`, `callback_num_atag`, `callback_num_pres_ind`, `broadcast_area_identifier` and diff --git a/README.md b/README.md index e6f1987..2b125cd 100644 --- a/README.md +++ b/README.md @@ -383,9 +383,10 @@ holds for `session.send()`. message has failed. Answering through `sendReturn()` instead leaves the wait running. 3. Tear down what is left, resolving to an `err` that says what was lost. -At most 1000 unanswered messages, and 64 MiB of them by the `maxOctets` charge, are held for five -minutes each; what falls out of a bound is dropped with a warning on the log and waited for no -longer. None of the bounds is an option. +At most 1000 unanswered messages, and 64 MiB of them by the `maxOctets` charge, are held. Past +either bound a new message is refused with `ESME_RTHROTTLED` (`ESME_RX_T_APPN` on a `deliver_sm`), +so the peer retries it. One held five minutes is dropped with a warning on the log and waited for +no longer. None of the bounds is an option. `close({ signal })` cuts the wait short. `unbind()` takes no signal, and waits a further `responseTimeout` for its own response. @@ -488,8 +489,8 @@ even where the SMSC took some of its segments; their receipts still arrive as `d **Multipart is answered on arrival.** Each segment is answered as it lands, because a relaying SMSC will not send the next until the last is answered. The answer is `ESME_ROK`, unless the segment -numbers itself into no message this session can join, which refuses it, or the reassembly buffer is -full, which asks the SMSC to keep it and try again. `sms.answeredOnArrival` says whether the message +numbers itself into no message this session can join, which refuses it, or the reassembly buffer or +the unanswered messages are at their bound, which asks the SMSC to keep it and try again. `sms.answeredOnArrival` says whether the message you hold was answered that way; a segment count cannot, since a peer may number a message one part of one. diff --git a/docs/decisions.md b/docs/decisions.md index 72acb16..4190bd6 100644 --- a/docs/decisions.md +++ b/docs/decisions.md @@ -603,7 +603,7 @@ rule and an index of the titles below. untouched, and is where a caller-chosen id and a refusal live; `onRequest` is the escape hatch for an application that must refuse a PDU the `sms` event could not have shown it yet. `collect()` answers every segment it will not carry rather than leaving it unanswered, which is the same stall - in miniature: the field that numbered it where the segment belongs to no group, `ESME_RMSGQFUL` + in miniature: the field that numbered it where the segment belongs to no group, the retry status where the segment's own arrival overran the octet cap, since a peer told that still holds it. Rejected: answering every segment but the one that completes the group, which leaves the peer holding some segments accepted and one refused with nothing in SMPP to retract the rest, and still cannot honour @@ -677,12 +677,25 @@ rule and an index of the titles below. back to `responseTimeout`, the same answer the link gate's hold already takes — and to that option's default where it is 0 as well, since neither option is an answer about the application. -- **What the application holds unanswered is capped on constants.** A bound the application cannot - raise is the point: an application that answers nothing would otherwise grow it for the life of - the link, which goal 4 forbids. Reassembly's `maxOctets` is an option because it bounds what the - peer sends; this bounds what the application leaves unanswered. A message that falls out of a bound - is one the drain stops waiting for, so `close()` can report fewer unanswered than there were — - accepted, because the alternative is holding what nothing will answer. +- **What the application holds unanswered is capped on constants, and a message past the cap is + refused.** A bound the application cannot raise is the point: an application that answers nothing + would otherwise grow it for the life of the link, which goal 4 forbids. Reassembly's `maxOctets` + is an option because it bounds what the peer sends; this bounds what the application leaves + unanswered. Maintainer's call, 2026-09-26: a message arriving past the cap is refused with the + retry status, so the peer keeps it (goal 2). Rejected: dropping the oldest to make room, which + frees nothing while the application still holds its `Sms`, and stops the drain waiting for a + message the peer is owed. Rejected: pausing the socket, which also stalls every answer and + `enquire_link` on the link. A message held past its timeout is still dropped, so `close()` can + report fewer unanswered than there were — accepted, because the alternative is holding what + nothing will answer. + +- **A store at its bound answers `ESME_RTHROTTLED` to a `submit_sm` and `ESME_RX_T_APPN` to a + `deliver_sm`.** Maintainer's call, 2026-09-26, for reassembly and held messages alike, so "keep it + and retry" has one spelling per direction. `ESME_RTHROTTLED` asks the sender to slow down, which + is what the peer outrunning us needs, and operators send it (Vonage, LINK Mobility, Route Mobile, + Jasmin), so clients built against them meet it (goal 1). Rejected: `ESME_RMSGQFUL`, which names an + exhausted queue and no rate. `ESME_RTHROTTLED` is the SMSC's to send, so an ESME answers with + SMPP 3.4's temporary receiver error, the one an SMSC retries on (goal 3). - **A reconnect keeps the delivery-receipt merges; everything else the link held is dropped.** `onDelivery()` answers each receipt before the group it belongs to is complete, and `teardown()` diff --git a/src/held-messages.ts b/src/held-messages.ts index dcf3dca..adab745 100644 --- a/src/held-messages.ts +++ b/src/held-messages.ts @@ -30,7 +30,6 @@ export class HeldMessages { private readonly held: ExpiringGroups; private readonly idleWaiters = new IdleWaiters(); private readonly log: SmppLog; - private readonly max: number; private readonly maxOctets: number; private octets = 0; @@ -42,7 +41,6 @@ export class HeldMessages { timeout: options.timeout, }); this.log = options.log; - this.max = options.max; this.maxOctets = options.maxOctets; } @@ -50,7 +48,13 @@ export class HeldMessages { return this.held.size; } - /** An application that answers no message at all may not grow this without end. */ + /** Whether a message arriving now is past the bound, once the expired are swept. */ + full(): boolean { + this.sweep(); + + return this.held.full || this.octets >= this.maxOctets; + } + hold(pduObjs: PduObject[]): void { const key = keyOf(pduObjs); @@ -63,17 +67,10 @@ export class HeldMessages { if (replaced) { this.log.warn('heldMessages - replacing a message on a re-used sequence number', { seqNr: Number(key) }); this.delete(key, replaced); - } else if (this.held.full) { - this.dropOldest(); } const octets = pduObjs.reduce((sum, pduObj) => sum + retainedOctets(pduObj), 0); - // The message just held stays even alone past the cap: the peer is still owed its answer. - while (this.held.size > 0 && this.octets + octets > this.maxOctets) { - this.dropOldest(); - } - this.held.set(key, { octets, pduObjs }); this.octets += octets; } @@ -109,22 +106,6 @@ export class HeldMessages { return this.idleWaiters.wait(() => this.held.size, timeout, signal); } - private dropOldest(): void { - const oldest = this.held.takeOldest(); - - if (!oldest) return; - - const [seqNr, held] = oldest; - - this.octets -= held.octets; - this.log.warn('heldMessages - buffer full, dropping the oldest message', { - max: this.max, - maxOctets: this.maxOctets, - octets: this.octets, - seqNr: Number(seqNr), - }); - } - /** Drops every message past its deadline. Runs before each hold and on its own timer. */ sweep(): void { const expired = this.held.takeExpired(); diff --git a/src/incoming-requests.ts b/src/incoming-requests.ts index 244f897..98e6f81 100644 --- a/src/incoming-requests.ts +++ b/src/incoming-requests.ts @@ -19,7 +19,11 @@ import { paramText } from './defs/types.ts'; import { respIdParams, segmentId } from './sms-id.ts'; import { respNameFor } from './defs/commands.ts'; -/** SMPP 3.4 lists ESME_RMSGQFUL under submit_sm_resp only; 4.6.2's retryable code is another. */ +/** Asks the peer to keep the message and retry: RTHROTTLED is the SMSC's to send, so an ESME's is another. */ +function throttledStatus(carriedAs: string): ErrorName { + return carriedAs === 'submit_sm' ? 'ESME_RTHROTTLED' : 'ESME_RX_T_APPN'; +} + export function refusedSegmentStatus( carriedAs: string, refusal: Refusal, @@ -30,7 +34,7 @@ export function refusedSegmentStatus( return spelling === 'sar' ? 'ESME_RINVTLVVAL' : 'ESME_RINVESMCLASS'; } - return carriedAs === 'submit_sm' ? 'ESME_RMSGQFUL' : 'ESME_RX_T_APPN'; + return throttledStatus(carriedAs); } const lostReasons: Record = { @@ -212,6 +216,16 @@ export class IncomingRequests { * one request at a time never sends the second segment until the first has been answered. */ private async onMessage(pduObj: PduObject): Promise { + if (this.held.full()) { + this.log.info('session - unanswered messages at their bound, asking the peer to retry', { + cmdName: pduObj.cmdName, + seqNr: pduObj.seqNr, + }); + await this.session.sendReturn(pduObj, throttledStatus(this.carriedAs(pduObj))); + + return; + } + const concat = concatOf(pduObj); if (!concat) { diff --git a/test/session-extras.test.ts b/test/session-extras.test.ts index f210725..91027c9 100644 --- a/test/session-extras.test.ts +++ b/test/session-extras.test.ts @@ -1482,76 +1482,82 @@ describe('held message bounds', () => { return [submitPdu(seqNr)]; } - test('drops the message held longest rather than holding every one', () => { + 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 oldest = message(1); + const first = message(1); - held.hold(oldest); + held.hold(first); held.hold(message(2)); held.hold(message(2)); - assert.equal(held.size, 2, 'a re-used sequence number replaces rather than evicting'); - assert.equal(held.has(oldest), true); - - held.hold(message(3)); - assert.equal(held.size, 2); - assert.equal(held.has(oldest), false); + assert.equal(held.full(), true); + assert.equal(held.has(first), true); held.clear(); }); // submitPdu() holds 1026 octets by the maxOctets charge: its object, and the three text fields. - test('drops the message held longest once the octets held pass the cap', () => { + test('is full at its octet cap, until a message leaves by any way out', () => { let now = 0; - const held = new HeldMessages({ log: silentLog, max: 10, maxOctets: 2100, now: () => now, timeout: 10_000 }); - const oldest = message(1); - - held.hold(oldest); - held.hold(message(2)); - - assert.equal(held.size, 2); - - held.hold(message(3)); - - assert.equal(held.size, 2); - assert.equal(held.has(oldest), false); - - // A message that leaves any other way gives its octets back, so two still fit afterwards. - const answered = message(4); + const held = new HeldMessages({ log: silentLog, max: 10, maxOctets: 2000, now: () => now, timeout: 10_000 }); + const answered = message(1); held.hold(answered); + assert.equal(held.full(), false); + held.hold(message(2)); + assert.equal(held.full(), true); + held.release(answered); - held.hold(message(5)); - assert.equal(held.size, 2, 'after a release'); + assert.equal(held.full(), false, 'after a release'); + held.hold(message(3)); now = 20_000; held.sweep(); now = 0; - held.hold(message(6)); - held.hold(message(7)); - assert.equal(held.size, 2, 'after a sweep'); - - held.clear(); - held.hold(message(8)); - held.hold(message(9)); - assert.equal(held.size, 2, 'after a clear'); + assert.equal(held.full(), false, 'after a sweep'); + held.hold(message(4)); + held.hold(message(5)); held.clear(); + assert.equal(held.full(), false, 'after a clear'); }); - // Dropping it would leave the drain blind to a message the peer is still owed an answer for. - test('keeps a message larger than the cap on its own', () => { - const held = new HeldMessages({ log: silentLog, max: 10, maxOctets: 1000, timeout: 10_000 }); - const large = message(2); + // Dropping one the application still holds frees nothing, and the drain stops waiting for it. + test('refuses what arrives past the bound with a status that asks the peer to retry', async t => { + const session = new Session({ sock: new net.Socket() }); - held.hold(message(1)); - held.hold(large); + closeAfter(t, session); + session.boundAs = 'transceiver'; - assert.equal(held.size, 1); - assert.equal(held.has(large), true); + const incoming = new IncomingRequests({ + dlrMerger: new DlrMerger({ log: silentLog, max: 10, timeout: 10_000 }), + log: silentLog, + sendPastDrain: () => Promise.resolve({ err: new Error('never sent') }), + session, + }); + const answers: (ErrorName | undefined)[] = []; + let messages = 0; - held.clear(); + session.sendReturn = (_pdu, status) => { + answers.push(status); + + return Promise.resolve({}); + }; + session.on('sms', () => { messages++; }); + + for (let seqNr = 1; seqNr <= 1000; seqNr++) { + await incoming.handle(submitPdu(seqNr)); + } + + assert.equal(messages, 1000); + + await incoming.handle(submitPdu(1001)); + await incoming.handle(segment(7, 1, 2)); + + assert.equal(messages, 1000); + assert.deepEqual(answers, ['ESME_RTHROTTLED', 'ESME_RTHROTTLED']); + incoming.clear(); }); test('holds a message detached from the chunk it was read from', async t => { @@ -2213,16 +2219,15 @@ describe('reassembly bounds', () => { // Jasmin dispatches one request per connector at a time: holding a group unanswered until it was // whole deadlocked every multi-segment message against it (interop-tests/findings/03-jasmin.md). describe('the status a refused segment is answered with', () => { - // SMPP 3.4 lists ESME_RMSGQFUL under submit_sm_resp only; 4.6.2's retryable code is another. test('names one the command the segment arrived on defines', () => { - assert.equal(refusedSegmentStatus('submit_sm', 'full', 'udh'), 'ESME_RMSGQFUL'); + assert.equal(refusedSegmentStatus('submit_sm', 'full', 'udh'), 'ESME_RTHROTTLED'); assert.equal(refusedSegmentStatus('deliver_sm', 'full', 'udh'), 'ESME_RX_T_APPN'); assert.equal(refusedSegmentStatus('submit_sm', 'unplaceable', 'udh'), 'ESME_RINVESMCLASS'); assert.equal(refusedSegmentStatus('deliver_sm', 'unplaceable', 'udh'), 'ESME_RINVESMCLASS'); // esm_class is 0x00 on a sar_* segment and entirely valid: the TLV values are what cannot be honoured. assert.equal(refusedSegmentStatus('submit_sm', 'unplaceable', 'sar'), 'ESME_RINVTLVVAL'); assert.equal(refusedSegmentStatus('deliver_sm', 'unplaceable', 'sar'), 'ESME_RINVTLVVAL'); - assert.equal(refusedSegmentStatus('submit_sm', 'full', 'sar'), 'ESME_RMSGQFUL'); + assert.equal(refusedSegmentStatus('submit_sm', 'full', 'sar'), 'ESME_RTHROTTLED'); }); // Which command that is, for the one that travels both ways, is what the end it arrived at says. @@ -2232,7 +2237,7 @@ describe('the status a refused segment is answered with', () => { assert.equal(standsInFor('deliver_sm', 'esme'), 'deliver_sm'); assert.equal(standsInFor('submit_sm', 'smsc'), 'submit_sm'); assert.equal(standsInFor('enquire_link', 'smsc'), 'enquire_link'); - assert.equal(refusedSegmentStatus(standsInFor('data_sm', 'smsc'), 'full', 'udh'), 'ESME_RMSGQFUL'); + assert.equal(refusedSegmentStatus(standsInFor('data_sm', 'smsc'), 'full', 'udh'), 'ESME_RTHROTTLED'); assert.equal(refusedSegmentStatus(standsInFor('data_sm', 'esme'), 'full', 'udh'), 'ESME_RX_T_APPN'); }); }); diff --git a/test/session.test.ts b/test/session.test.ts index 9aab5e6..ce6af2c 100644 --- a/test/session.test.ts +++ b/test/session.test.ts @@ -932,7 +932,7 @@ describe('receiving', () => { assert.ok(refused.pduObj); assert.equal(refused.pduObj.cmdName, 'data_sm_resp'); - assert.equal(refused.pduObj.cmdStatus, 'ESME_RMSGQFUL'); + assert.equal(refused.pduObj.cmdStatus, 'ESME_RTHROTTLED'); }); test('reassembles a concatenated message whose segments arrived in message_payload', async t => { diff --git a/todo.md b/todo.md index 1cac818..843a501 100644 --- a/todo.md +++ b/todo.md @@ -6,17 +6,6 @@ hard rules first — they constrain every item below. This is a working file that sets its own rules. The documentation conventions in AGENTS.md do not govern it, and nothing here is a source anything else may cite. -## Security - -- [ ] **Bound what a peer can make the application hold, not only the library.** An `sms` the - application is still answering pins its PDUs through `Sms.pduObjs` and its `sendResp`/`sendDlr` - closures, so the held-message caps free nothing while it works, and nothing slows the peer: a - peer faster than an application answering asynchronously (a DB write per message) grows the - heap without bound. Each eviction by the caps also stops the drain waiting for a message still - being answered, so `close()` can cut it off and the peer re-sends it. Flow control is a wire - change under goal 4 and the maintainer's call: answer `ESME_RTHROTTLED`, or stop reading the - socket, once the held count or octets are at the cap. From the stability review of #28. - ## Status The rewrite is **feature complete and green**: the suite, lint and typecheck are clean on Node 18 -- 2.52.0 From fb6b125aab88546399254454b1952f7fd7b51105 Mon Sep 17 00:00:00 2001 From: Lilleman auf Larv Date: Sat, 26 Sep 2026 14:37:20 +0200 Subject: [PATCH 2/5] Give the dumbclient peers a namespace owner that outlives them, count window 2000's throttled messages, and correct what the findings blamed on the peer --- CHANGELOG.md | 4 +- README.md | 14 ++--- interop-tests/compose.dumbclient.yaml | 37 ++++++++----- interop-tests/dumbclient.test.ts | 45 ++++++++-------- interop-tests/findings/07-load.md | 76 ++++++++++----------------- src/incoming-requests.ts | 2 +- test/session-extras.test.ts | 21 +++++--- 7 files changed, 98 insertions(+), 101 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index 2cb73c8..27ad67a 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -44,8 +44,8 @@ `maxOctets` charge, is refused with `ESME_RTHROTTLED` (`ESME_RX_T_APPN` on a `deliver_sm`), so the peer keeps it and retries. The oldest used to be dropped to make room, which freed nothing while the application still held it and let `close()` stop waiting for a message the peer was owed. -- A segment the reassembly buffer has no room for is refused with `ESME_RTHROTTLED`, where it was - `ESME_RMSGQFUL`. +- A `submit_sm` segment the reassembly buffer has no room for is refused with `ESME_RTHROTTLED`, + where it was `ESME_RMSGQFUL`. - `server()` refuses a `maxOctets` below 1 or not a whole number, `Infinity` included, like its other limits. `server({ maxOctets: 0 })` used to start and then refuse every multipart message. - `callback_num`, `callback_num_atag`, `callback_num_pres_ind`, `broadcast_area_identifier` and diff --git a/README.md b/README.md index 2b125cd..c9f69b0 100644 --- a/README.md +++ b/README.md @@ -383,10 +383,10 @@ holds for `session.send()`. message has failed. Answering through `sendReturn()` instead leaves the wait running. 3. Tear down what is left, resolving to an `err` that says what was lost. -At most 1000 unanswered messages, and 64 MiB of them by the `maxOctets` charge, are held. Past -either bound a new message is refused with `ESME_RTHROTTLED` (`ESME_RX_T_APPN` on a `deliver_sm`), -so the peer retries it. One held five minutes is dropped with a warning on the log and waited for -no longer. None of the bounds is an option. +A message arriving while 1000 unanswered messages, or 64 MiB of them by the `maxOctets` charge, are +held is refused with `ESME_RTHROTTLED` (`ESME_RX_T_APPN` on a `deliver_sm`), so the peer retries it. +One held five minutes is dropped with a warning on the log and waited for no longer. None of the +bounds is an option. `close({ signal })` cuts the wait short. `unbind()` takes no signal, and waits a further `responseTimeout` for its own response. @@ -490,9 +490,9 @@ even where the SMSC took some of its segments; their receipts still arrive as `d **Multipart is answered on arrival.** Each segment is answered as it lands, because a relaying SMSC will not send the next until the last is answered. The answer is `ESME_ROK`, unless the segment numbers itself into no message this session can join, which refuses it, or the reassembly buffer or -the unanswered messages are at their bound, which asks the SMSC to keep it and try again. `sms.answeredOnArrival` says whether the message -you hold was answered that way; a segment count cannot, since a peer may number a message one part -of one. +the unanswered messages are at their bound, which asks the SMSC to keep it and try again. +`sms.answeredOnArrival` says whether the message you hold was answered that way; a segment count +cannot, since a peer may number a message one part of one. - The id was fixed with the first segment, so `sendResp()` there only says you are done, and returns `err` for an `smsId` or a refusing `status`. diff --git a/interop-tests/compose.dumbclient.yaml b/interop-tests/compose.dumbclient.yaml index a04a6a7..b29e726 100644 --- a/interop-tests/compose.dumbclient.yaml +++ b/interop-tests/compose.dumbclient.yaml @@ -12,27 +12,36 @@ x-dumbclient-healthcheck: &dumbclient-healthcheck timeout: 2s services: - # Network owner for every dumbclient-* service and the capture sidecar below (see the comment on - # `capture`): all four are pure outbound TCP clients with nothing of their own listening, so - # sharing one netns is only ever a source-IP detail, never a port collision. + # Network owner for every dumbclient-* service and the capture sidecar below: all four are pure + # outbound TCP clients, so sharing one netns is only a source-IP detail. It never exits, because a + # client that finishes early would take the namespace, and every other conversation, with it. + dumbclient-netns: + image: nicolaka/netshoot:v0.16 + <<: *log-limits + command: ["sleep", "infinity"] + dumbclient-w2000: build: ./interop-tests/peers/dumbclient image: interop-dumbclient:de0334b <<: *log-limits + network_mode: "service:dumbclient-netns" command: ["conf/window2000.yml"] healthcheck: *dumbclient-healthcheck + depends_on: + dumbclient-netns: + condition: service_started # S9's comparison run: window below maxHeldMessages (1000, session-options.ts), where nothing - # should ever be evicted - see findings/07-load.md. + # should ever be throttled - see findings/07-load.md. dumbclient-w500: build: ./interop-tests/peers/dumbclient image: interop-dumbclient:de0334b <<: *log-limits - network_mode: "service:dumbclient-w2000" + network_mode: "service:dumbclient-netns" command: ["conf/window500.yml"] healthcheck: *dumbclient-healthcheck depends_on: - dumbclient-w2000: + dumbclient-netns: condition: service_started # S6: sends one message, then never speaks again - the no-ping binary (see the Dockerfile) sends @@ -41,13 +50,13 @@ services: build: ./interop-tests/peers/dumbclient image: interop-dumbclient:de0334b <<: *log-limits - network_mode: "service:dumbclient-w2000" + network_mode: "service:dumbclient-netns" environment: DUMBCLIENT_BIN: /app/smpp-dumb-client-noping command: ["conf/idle.yml"] healthcheck: *dumbclient-healthcheck depends_on: - dumbclient-w2000: + dumbclient-netns: condition: service_started # The long soak: the longest run the time-box allows, fast handler, watched for anything that @@ -56,25 +65,25 @@ services: build: ./interop-tests/peers/dumbclient image: interop-dumbclient:de0334b <<: *log-limits - network_mode: "service:dumbclient-w2000" + network_mode: "service:dumbclient-netns" command: ["conf/soak.yml"] healthcheck: *dumbclient-healthcheck depends_on: - dumbclient-w2000: + dumbclient-netns: condition: service_started - # Every dumbclient-* service shares dumbclient-w2000's netns (see above), so this one sidecar + # Every dumbclient-* service shares dumbclient-netns's namespace (see above), so this one sidecar # sees all four conversations with node:2775 - the same pattern compose.kannel.yaml uses for its # four bearerbox variants, one namespace deeper. capture: image: nicolaka/netshoot:v0.16 - network_mode: "service:dumbclient-w2000" + network_mode: "service:dumbclient-netns" cap_add: - NET_ADMIN - NET_RAW depends_on: - dumbclient-w2000: - condition: service_healthy + dumbclient-netns: + condition: service_started command: ["dumpcap", "-i", "any", "-f", "tcp port 2775", "-w", "/captures/dumbclient.pcapng"] volumes: - ./interop-tests/captures:/captures diff --git a/interop-tests/dumbclient.test.ts b/interop-tests/dumbclient.test.ts index 9a292e7..81bc9ff 100644 --- a/interop-tests/dumbclient.test.ts +++ b/interop-tests/dumbclient.test.ts @@ -206,17 +206,26 @@ after(async () => { // S9 (target 11) and the backpressure-at-server scenario: window 2000 at a high rate against a // handler slowed enough to build a real backlog. window500 is the same shape with a window below -// maxHeldMessages (1000, session-options.ts defaults.maxHeldMessages) - see findings/07-load.md for -// what that constant, rather than maxOutstanding, turns out to be the one that interacts with a -// peer's window. -describe('S9 - bounded window against a slowed handler', () => { - for (const [name, expectedCount] of [['dumb-w500', 20_000], ['dumb-w2000', 20_000]] as const) { - test(`${name}: every message answered exactly once, ordering holds`, async () => { - const done = await waitFor(() => (statsFor(name).answered >= expectedCount ? true : undefined), 180_000); +// maxHeldMessages (1000, session-options.ts defaults.maxHeldMessages), the bound past which a +// peer's window is answered ESME_RTHROTTLED. smpp-dumb-client counts a throttled message as sent +// and never resends it, so window 2000 accounts for 20,000 as answered plus throttled. +const throttleMessage = 'session - unanswered messages at their bound, asking the peer to retry'; - assert.ok(done, `${name} did not answer ${String(expectedCount)} messages within budget`); +// window500's peak (<=500) and the soak's never reach the 1000 default, so every refusal is +// necessarily from the w2000 session - the runs share one server and one log. +function throttled(name: string): number { + return name === 'dumb-w2000' ? logEntries.filter(entry => entry.message === throttleMessage).length : 0; +} + +describe('S9 - bounded window against a slowed handler', () => { + for (const name of ['dumb-w500', 'dumb-w2000'] as const) { + test(`${name}: every message answered or throttled exactly once, ordering holds`, async () => { + const done = await waitFor(() => (statsFor(name).answered + throttled(name) >= 20_000 ? true : undefined), 180_000); + + assert.ok(done, `${name} did not account for 20000 messages within budget`); const s = statsFor(name); + const expectedCount = 20_000 - throttled(name); assert.equal(s.arrived, expectedCount); assert.equal(s.answered, expectedCount); @@ -227,17 +236,9 @@ describe('S9 - bounded window against a slowed handler', () => { }); } - test('window 2000 pressed past maxHeldMessages (1000): the internal held-message cap evicts, window500 never does', async () => { - await waitFor(() => (statsFor('dumb-w2000').answered >= 20_000 ? true : undefined), 180_000); - - const evictions = logEntries.filter(entry => entry.message === 'heldMessages - buffer full, dropping the oldest message'); - - // window500's peak (<=500) never reaches the 1000 default, so any eviction observed is - // necessarily from the w2000 session - the two runs share one server and one log. - assert.ok(evictions.length > 0, 'expected at least one held-message eviction under window 2000'); - // The peer's own window, respected exactly both runs (peakOutstanding read 500 and 2000 on - // the nose) - the lower bound is what distinguishes this from window500's own eviction-free run. - assert.ok(statsFor('dumb-w2000').peakOutstanding > 1000 && statsFor('dumb-w2000').peakOutstanding <= 2000); + test('window 2000 pressed past maxHeldMessages (1000): the peer is throttled, window500 never is', () => { + assert.ok(throttled('dumb-w2000') > 0, 'expected at least one ESME_RTHROTTLED under window 2000'); + assert.ok(statsFor('dumb-w2000').peakOutstanding <= 1000); assert.equal(statsFor('dumb-w500').peakOutstanding <= 500, true); }); @@ -284,10 +285,8 @@ describe('S6 - idle peer, no enquire_link at all', () => { }); // The long soak: the longest run the time-box allows, fast handler, watched for anything that -// grows without bound (held messages, listeners, memory). Bounded by wall-clock rather than a -// target count: smpp-dumb-client's own TX-tracking window bookkeeping stalls under sustained load -// (findings/07-load.md, Peer quirks) well short of the configured count, on the client's side only -// - our own arrived/answered stay in lockstep throughout, which is what this asserts. +// grows without bound (held messages, listeners, memory). Bounded by wall-clock, and asserting that +// arrived/answered stay in lockstep. describe('Long soak', () => { const SOAK_DURATION_MS = 300_000; diff --git a/interop-tests/findings/07-load.md b/interop-tests/findings/07-load.md index 93e3b5e..73cddd4 100644 --- a/interop-tests/findings/07-load.md +++ b/interop-tests/findings/07-load.md @@ -65,7 +65,7 @@ keeps a live reproducer asserting what our server does when it receives it (refu unframeable - see Scenarios) rather than removing the peer. `smpp-dumb-client` covers S9, and substitutes for S6 and (partially) S8 - see below. -### smpp-dumb-client: builds and interoperates cleanly; its own window bookkeeping stalls under sustained load +### smpp-dumb-client: builds and interoperates cleanly No build friction. Two binaries from the same pinned source: `smpp-dumb-client` (unmodified) and `smpp-dumb-client-noping` (its two `enquireSender()` call sites in `smpp.go` commented out at build @@ -78,25 +78,18 @@ One integration snag, not a build one: `smpp.remote` in `config.yml` is fed stra there directly. Fixed in the entrypoint: every `conf/*.yml` carries a `NODE_HOST` placeholder, resolved with `getent hosts` and substituted into a writable copy before the real binary starts. -Four one-shot scenarios share `dumbclient-w2000`'s network namespace (`network_mode: -"service:dumbclient-w2000"`) - they are pure outbound clients with nothing of their own listening, -so the only shared cost is a source IP, and one capture sidecar sees all four conversations with -`node:2775` the same way `compose.kannel.yaml`'s does for its four bearerbox variants. +The four scenarios share one network namespace, owned by `dumbclient-netns`, a container that +never exits - they are pure outbound clients with nothing of their own listening, so the only shared +cost is a source IP, and one capture sidecar sees all four conversations with `node:2775` the same +way `compose.kannel.yaml`'s does for its four bearerbox variants. -The long soak (below) surfaced a peer-side limit worth designing around rather than fighting: with -a fast, immediate-response handler and a window of 100 - nothing our server should ever have -trouble draining - the peer's own reported in-flight count (`GetTrackQueueSize`, read from -`len(TrackTX)`) gets stuck pinned at the window within the first minute, and its log fills with -`Expired TX packet` lines (`libsmpp`'s hardcoded, non-configurable 7000ms `TX_MAX_TIMEOUT_MS`) - -throughput drops from ~500/s to a trickle of tens per second, gated by how many tracked entries -individually cross that 7s mark each second rather than by real responses being matched. Our own -server-side counters (`arrived`/`answered`/`peakOutstanding`, tracked independently in -`dumbclient.test.ts`) stay in lockstep throughout with a low peak - see Scenarios - which places the -stall entirely on the peer's own window bookkeeping, not on anything our server did or failed to -do. The soak test was redesigned around this: bounded by wall-clock (5 minutes) rather than a -target count, asserting the invariants that matter regardless of how much the peer's own bug lets -through (every arrival answered, nothing duplicated, memory shape), and reporting whatever -throughput was actually reached rather than requiring a specific one. +Runs 1 and 2 had `dumbclient-w2000` own the namespace. It exits once it has sent its 20,000, which +took every other client's network with it: the soak's responses stopped arriving, and its log filled +with `Expired TX packet` lines (`libsmpp`'s 7000ms `TX_MAX_TIMEOUT_MS`). Those runs read that as the +peer's own window bookkeeping stalling; run 3 (2026-09-26), with the namespace owned by a container +that outlives them all, reached 173,820 soak messages in 300s where run 2 reached 22,440. The soak +stays bounded by wall-clock (5 minutes), asserting every arrival answered, nothing duplicated, and +the memory shape. One test-harness bug found and fixed between the two runs below, not a library defect: the S6 test's first version attached its `session.on('close', ...)` listener lazily inside the test body, after @@ -105,9 +98,9 @@ the S6 test ran, the idle session had already closed, and an `EventEmitter` neve event to a listener added after it fired. Fixed by attaching every session's `close` listener at `session`-creation time, recording it in the same per-scenario stats every other assertion reads. -Two runs of `./interop-tests/run.py dumbclient`. Run 1 (the original 300,000-count soak) surfaced -both the peer's TX-tracking stall and the S6 harness bug above; run 2, after both fixes, is the one -reported below. `smppload.test.ts` passed on every run it was given (three, across the investigation +Three runs of `./interop-tests/run.py dumbclient`. Run 1 (the original 300,000-count soak) surfaced +the S6 harness bug above; run 2 fixed it; run 3, with the namespace owner above and the held-message +throttle in `src/`, is the one Scenarios reports. The capture figures below are run 2's. `smppload.test.ts` passed on every run it was given (three, across the investigation above); its one scenario needs no repeat - a second run reproduces the identical corrupted PDU, adding nothing. @@ -127,29 +120,24 @@ every session's own `arrived` exactly, and every session's own `answered` matche ## Throughput and memory -`dumb-w500` and `dumb-w2000` (S9) both ran to their full 20,000-message count in ~44s each, -concurrently, against a handler serialised to answer roughly one message every 2ms -(`SLOW_HANDLER_DELAY_MS`) - `peakOutstanding` read exactly 500 and exactly 2000, the two configured -windows, confirming the peer never let more than its own window ride at once. +Run 3. `dumb-w500` ran to its full 20,000 in ~44s against a handler serialised to answer roughly +one message every 2ms (`SLOW_HANDLER_DELAY_MS`), `peakOutstanding` exactly 500. `dumb-w2000`, run +concurrently, held exactly 1000 and was throttled for the rest. -The soak (fast, immediate-response handler; window 100) reached 22,440 `submit_sm` over its fixed -300s observation window - about 75/s, well under the peer's own configured `rate: 500` and under -what our server can sustain (see Setup: `smpp-dumb-client`'s own TX-tracking bookkeeping is the -ceiling here, not our server - `peakOutstanding` stayed at 25 throughout). Sampled every 5s across -the whole run (69 samples over 340s, all four scenarios combined): rss first=170MiB, min=124MiB, -max=306MiB (during the two window runs' backlog), last=125MiB, heapUsed at the last sample 15MiB - -back below its own starting point once the backlog drained, not merely flat. No monotonic trend in -either direction. +The soak (fast, immediate-response handler; window 100) reached 173,820 `submit_sm` over its fixed +300s, about 580/s, `peakOutstanding` 15. Sampled every 5s across the whole run (69 samples over +340s, all four scenarios combined, the harness's own per-message bookkeeping included): rss +first=165MiB, min=165MiB, max=298MiB, last=298MiB, heapUsed at the last sample 81MiB. ## Scenarios (PLAN.md) | Id | Result | Evidence | | --- | --- | --- | | S6 (idleTimeout, no peer ever pings) | pass | `dumbclient.test.ts` "S6 - idle peer..." - dropped at idleTimeout, `linkTimers - closing an idle peer` logged, no response past the one owed | -| S8 (throughput, long messages, receipts) | blocked (smppload) / partial substitute | smppload's own scenario is blocked - see Setup. The soak below gives a genuine submit_sm/s figure without long messages or receipts, which `smpp-dumb-client` does not support (`research/esme-clients-and-validators.md` section B) - and is itself capped well below what our server can sustain by the peer's own TX-tracking stall, also see Setup | -| S9 (bounded window) | pass | `dumbclient.test.ts` "S9 - bounded window..." - 20,000/20,000 answered on both window 500 and window 2000, in arrival order, no duplicate ids, `peakOutstanding` exactly 500 and exactly 2000 | -| Backpressure at the server | pass | Same run: `peakOutstanding` 2000 exceeds `maxHeldMessages` (1000, session-options.ts) and the eviction warning fires; window 500 (`peakOutstanding` 500) never does; memory sampled before/after the window runs (170MiB before, 306MiB after, 125MiB once the soak's own run had also settled) | -| Long soak | pass (run once at the redesigned, wall-clock-bounded shape - see Setup) | `dumbclient.test.ts` "Long soak" - 22,440 arrived, 22,440 answered, 0 duplicates, 0 unanswered errors, `close()` drains with no error | +| S8 (throughput, long messages, receipts) | blocked (smppload) / partial substitute | smppload's own scenario is blocked - see Setup. The soak below gives a genuine submit_sm/s figure without long messages or receipts, which `smpp-dumb-client` does not support (`research/esme-clients-and-validators.md` section B) | +| S9 (bounded window) | pass | Run 3: `dumbclient.test.ts` "S9 - bounded window..." - window 500 20,000/20,000 answered in ~44s, `peakOutstanding` exactly 500; window 2000 5,639 answered and 14,361 throttled, in arrival order, no duplicate ids | +| Backpressure at the server | pass | Run 3: window 2000 holds exactly 1000 (`maxHeldMessages`, session-options.ts) and the rest is answered `ESME_RTHROTTLED`; smpp-dumb-client counts a throttled message as sent and never resends it; window 500 is never throttled | +| Long soak | pass | Run 3: `dumbclient.test.ts` "Long soak" - 173,820 arrived, 173,820 answered, 0 duplicates, 0 unanswered errors, `close()` drains with no error; rss 165MiB first, 298MiB max and last, heapUsed 81MiB last | | smppload bind corruption (not in PLAN.md - found this phase) | blocked | `smppload.test.ts` - our server refuses the unreadable stream instead of hanging | ## Defects in @larvit/smpp @@ -157,16 +145,14 @@ either direction. None found. `smppload.test.ts`'s own scenario is smppload's defect, not ours: our server's reaction (refusing the stream as unframeable, per the decision in the root `AGENTS.md`, "A stream this library cannot frame...") is the documented behaviour working exactly as designed against a peer -that never gets as far as a readable PDU. The soak's throughput ceiling is the peer's own defect -(see Setup) - our own `arrived`/`answered`/`peakOutstanding` counters stayed clean throughout every -run. +that never gets as far as a readable PDU. ## Peer quirks - **smppload's `bind_transceiver` is corrupted on the wire** - see Setup. Not chased past `oserl`'s `pack/2` (which is correct on inspection) given the time-box. -- **`smpp-dumb-client`'s window bookkeeping stalls under sustained load, throttling its own - throughput far below what a promptly-answering server can sustain** - see Setup. Its `enquire_link` +- **`smpp-dumb-client` treats `ESME_RTHROTTLED` as final** - a throttled message counts as sent + and is never resubmitted. Its `enquire_link` interval (10s once bound as an ESME) is also hardcoded (`smpp.go`, `enquireSender(10)`), not exposed through `config.yml` at all - the no-ping binary built for S6 patches the call site out rather than configuring it. @@ -177,7 +163,3 @@ run. - Whether smppload's bind corruption is in `oserl`'s `gen_esme_session`/`smpp_session` send path (not reached, given the time-box) or something specific to this build's dependency versions. -- Whether `smpp-dumb-client`'s stall is a sequence-number correlation bug (a response failing to - match its `TrackTX` entry, falling back to the 7s expiry) or something else in its own window - accounting - not chased past the observation in Setup, given the time-box and that the fault is - clearly on the peer's side (our own counters stayed clean throughout). diff --git a/src/incoming-requests.ts b/src/incoming-requests.ts index 98e6f81..54ffb3d 100644 --- a/src/incoming-requests.ts +++ b/src/incoming-requests.ts @@ -217,7 +217,7 @@ export class IncomingRequests { */ private async onMessage(pduObj: PduObject): Promise { if (this.held.full()) { - this.log.info('session - unanswered messages at their bound, asking the peer to retry', { + this.log.verbose('session - unanswered messages at their bound, asking the peer to retry', { cmdName: pduObj.cmdName, seqNr: pduObj.seqNr, }); diff --git a/test/session-extras.test.ts b/test/session-extras.test.ts index 91027c9..e541a93 100644 --- a/test/session-extras.test.ts +++ b/test/session-extras.test.ts @@ -24,7 +24,7 @@ import { Session } from '../src/session.ts'; import { DlrMerger } from '../src/dlr-merger.ts'; import { PduRefusedError } from '../src/pdu-refusal.ts'; import { objToPdu } from '../src/pdu.ts'; -import { checkSessionOptions, standsInFor } from '../src/session-options.ts'; +import { checkSessionOptions, defaults, standsInFor } from '../src/session-options.ts'; import { client } from '../src/client.ts'; import { closeAfter, closeListenerAfter } from './teardown.ts'; import { concatOf } from '../src/concat.ts'; @@ -1537,26 +1537,33 @@ describe('held message bounds', () => { session, }); const answers: (ErrorName | undefined)[] = []; - let messages = 0; + const received: Sms[] = []; session.sendReturn = (_pdu, status) => { answers.push(status); return Promise.resolve({}); }; - session.on('sms', () => { messages++; }); + session.on('sms', sms => { received.push(sms); }); - for (let seqNr = 1; seqNr <= 1000; seqNr++) { + for (let seqNr = 1; seqNr <= defaults.maxHeldMessages; seqNr++) { await incoming.handle(submitPdu(seqNr)); } - assert.equal(messages, 1000); + assert.equal(received.length, defaults.maxHeldMessages); - await incoming.handle(submitPdu(1001)); + await incoming.handle(submitPdu(defaults.maxHeldMessages + 1)); await incoming.handle(segment(7, 1, 2)); - assert.equal(messages, 1000); + assert.equal(received.length, defaults.maxHeldMessages); assert.deepEqual(answers, ['ESME_RTHROTTLED', 'ESME_RTHROTTLED']); + + // The refused first segment joined no group, so the second one completes nothing. + await received[0]?.sendResp(); + await delay(0); + await incoming.handle(segment(7, 2, 2)); + + assert.equal(received.length, defaults.maxHeldMessages); incoming.clear(); }); -- 2.52.0 From b9167073d44cc087b5ab6e3420af8d2ec9ff8eda Mon Sep 17 00:00:00 2001 From: Lilleman auf Larv Date: Sat, 26 Sep 2026 14:41:10 +0200 Subject: [PATCH 3/5] Warn once at the held-message bound, tell receivers sendResp() is what releases a message, and pin that a refused segment joins no group --- CHANGELOG.md | 9 +++++---- README.md | 18 +++++++++++------- docs/decisions.md | 4 +++- interop-tests/compose.dumbclient.yaml | 1 + interop-tests/findings/07-load.md | 11 +++++------ src/incoming-requests.ts | 14 ++++++++++++++ test/session-extras.test.ts | 9 ++++++--- todo.md | 5 +++++ 8 files changed, 50 insertions(+), 21 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index 27ad67a..fc2bcb3 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -40,10 +40,11 @@ count as next to nothing, so a peer could hold far more than the cap. **Raise a `maxOctets` you tuned low**: it now holds several times fewer segments, and an incomplete message evicted over the cap is lost, since its segments were already answered. -- A message arriving while the application holds 1000 unanswered, or 64 MiB of them by the - `maxOctets` charge, is refused with `ESME_RTHROTTLED` (`ESME_RX_T_APPN` on a `deliver_sm`), so the - peer keeps it and retries. The oldest used to be dropped to make room, which freed nothing while - the application still held it and let `close()` stop waiting for a message the peer was owed. +- A message arriving while the application holds 1000 unanswered, or 64 MiB of them counted the way + `maxOctets` counts segments, is refused with `ESME_RTHROTTLED` (`ESME_RX_T_APPN` on a + `deliver_sm`), so the peer keeps it and retries. **Call `sendResp()` on every `sms`, multipart + included**: 1000 left unanswered now stop inbound traffic for up to five minutes, where the oldest + used to be dropped with a warning. - A `submit_sm` segment the reassembly buffer has no room for is refused with `ESME_RTHROTTLED`, where it was `ESME_RMSGQFUL`. - `server()` refuses a `maxOctets` below 1 or not a whole number, `Infinity` included, like its diff --git a/README.md b/README.md index c9f69b0..988e2e6 100644 --- a/README.md +++ b/README.md @@ -101,9 +101,10 @@ session.on('sms', async sms => { }); ``` -Call `sendResp()` for every message; it is part of the protocol. Delivery receipts reach you as -`dlr` events, not here. A multipart message arrives reassembled and already answered segment by -segment, so `sendResp()` there only says you are done with it: [Receiving in depth](#receiving-in-depth). +Call `sendResp()` for every message, multipart included: until you do, it counts toward the bound +past which the peer's messages are refused. Delivery receipts reach you as `dlr` events, not here. A +multipart message arrives reassembled and already answered segment by segment, so `sendResp()` there +puts nothing on the wire and releases it: [Receiving in depth](#receiving-in-depth). ## Run an SMPP server @@ -383,10 +384,7 @@ holds for `session.send()`. message has failed. Answering through `sendReturn()` instead leaves the wait running. 3. Tear down what is left, resolving to an `err` that says what was lost. -A message arriving while 1000 unanswered messages, or 64 MiB of them by the `maxOctets` charge, are -held is refused with `ESME_RTHROTTLED` (`ESME_RX_T_APPN` on a `deliver_sm`), so the peer retries it. -One held five minutes is dropped with a warning on the log and waited for no longer. None of the -bounds is an option. +A message left unanswered for five minutes is no longer waited for. `close({ signal })` cuts the wait short. `unbind()` takes no signal, and waits a further `responseTimeout` for its own response. @@ -434,6 +432,12 @@ const { err, pduObj } = await session.send({ each is two messages. - **Answered on arrival.** Each segment was answered as it landed, before you see the message: [Server in depth](#server-in-depth). +- **Unanswered messages.** While 1000 messages you have not called `sendResp()` on, or 64 MiB of + them counted the way `maxOctets` counts segments, are held, every new message is refused with + `ESME_RTHROTTLED` (`ESME_RX_T_APPN` on a `deliver_sm`) so the peer retries it, and no `sms` fires. + Reaching the bound logs one `warn`, and falling back below it one `info`. A message left five + minutes is dropped from the count with a `warn`; a later `sendResp()` still answers it. Neither + bound is an option. - **Where the body is.** A body in the `message_payload` TLV, SMPP's way of carrying up to 64 KB and the only place a `data_sm` has, reads exactly like one in `short_message`, concatenated messages and receipts included. A PDU filling both is read from `short_message`. diff --git a/docs/decisions.md b/docs/decisions.md index 4190bd6..ea9cad1 100644 --- a/docs/decisions.md +++ b/docs/decisions.md @@ -685,7 +685,9 @@ rule and an index of the titles below. retry status, so the peer keeps it (goal 2). Rejected: dropping the oldest to make room, which frees nothing while the application still holds its `Sms`, and stops the drain waiting for a message the peer is owed. Rejected: pausing the socket, which also stalls every answer and - `enquire_link` on the link. A message held past its timeout is still dropped, so `close()` can + `enquire_link` on the link. Reaching the bound is a log line and nothing more (goal 8): an event + or a public count would be surface for what the application already knows, since it is the one + not answering. A message held past its timeout is still dropped, so `close()` can report fewer unanswered than there were — accepted, because the alternative is holding what nothing will answer. diff --git a/interop-tests/compose.dumbclient.yaml b/interop-tests/compose.dumbclient.yaml index b29e726..e6c4278 100644 --- a/interop-tests/compose.dumbclient.yaml +++ b/interop-tests/compose.dumbclient.yaml @@ -19,6 +19,7 @@ services: image: nicolaka/netshoot:v0.16 <<: *log-limits command: ["sleep", "infinity"] + init: true dumbclient-w2000: build: ./interop-tests/peers/dumbclient diff --git a/interop-tests/findings/07-load.md b/interop-tests/findings/07-load.md index 73cddd4..3b6a73e 100644 --- a/interop-tests/findings/07-load.md +++ b/interop-tests/findings/07-load.md @@ -100,8 +100,8 @@ event to a listener added after it fired. Fixed by attaching every session's `cl Three runs of `./interop-tests/run.py dumbclient`. Run 1 (the original 300,000-count soak) surfaced the S6 harness bug above; run 2 fixed it; run 3, with the namespace owner above and the held-message -throttle in `src/`, is the one Scenarios reports. The capture figures below are run 2's. `smppload.test.ts` passed on every run it was given (three, across the investigation -above); its one scenario needs no repeat - a second run reproduces the identical corrupted PDU, +throttle in `src/`, is the one Scenarios reports. The capture figures below are run 2's. +`smppload.test.ts` passed on every run it was given (three, across the investigation above); its one scenario needs no repeat - a second run reproduces the identical corrupted PDU, adding nothing. ``` @@ -152,10 +152,9 @@ that never gets as far as a readable PDU. - **smppload's `bind_transceiver` is corrupted on the wire** - see Setup. Not chased past `oserl`'s `pack/2` (which is correct on inspection) given the time-box. - **`smpp-dumb-client` treats `ESME_RTHROTTLED` as final** - a throttled message counts as sent - and is never resubmitted. Its `enquire_link` - interval (10s once bound as an ESME) is also hardcoded (`smpp.go`, `enquireSender(10)`), not - exposed through `config.yml` at all - the no-ping binary built for S6 patches the call site out - rather than configuring it. + and is never resubmitted. Its `enquire_link` interval (10s once bound as an ESME) is hardcoded + (`smpp.go`, `enquireSender(10)`), not exposed through `config.yml` at all - the no-ping binary + built for S6 patches the call site out rather than configuring it. - **`smpp.remote` takes a literal IP, never a hostname** (`net.ParseIP`, no DNS resolution) - see Setup. diff --git a/src/incoming-requests.ts b/src/incoming-requests.ts index 54ffb3d..817227e 100644 --- a/src/incoming-requests.ts +++ b/src/incoming-requests.ts @@ -71,6 +71,7 @@ export class IncomingRequests { private readonly smsIdFormat: SmsIdFormat; private readonly systemId: string; private linkGeneration = 0; + private refusing = false; constructor(options: IncomingRequestsOptions) { this.dlrMerger = options.dlrMerger; @@ -148,6 +149,7 @@ export class IncomingRequests { /** Drops the segments of every message that never became whole, and of every one still held. */ clear(): void { this.linkGeneration++; + this.refusing = false; this.held.clear(); this.reassembler.clear(); } @@ -217,6 +219,13 @@ export class IncomingRequests { */ private async onMessage(pduObj: PduObject): Promise { if (this.held.full()) { + if (!this.refusing) { + this.refusing = true; + this.log.warn('session - unanswered messages at their bound, refusing new ones until the application answers', { + messages: this.held.size, + }); + } + this.log.verbose('session - unanswered messages at their bound, asking the peer to retry', { cmdName: pduObj.cmdName, seqNr: pduObj.seqNr, @@ -226,6 +235,11 @@ export class IncomingRequests { return; } + if (this.refusing) { + this.refusing = false; + this.log.info('session - unanswered messages below their bound, accepting again', { messages: this.held.size }); + } + const concat = concatOf(pduObj); if (!concat) { diff --git a/test/session-extras.test.ts b/test/session-extras.test.ts index e541a93..98a540d 100644 --- a/test/session-extras.test.ts +++ b/test/session-extras.test.ts @@ -1530,9 +1530,10 @@ describe('held message bounds', () => { closeAfter(t, session); session.boundAs = 'transceiver'; + const warnings: string[] = []; const incoming = new IncomingRequests({ dlrMerger: new DlrMerger({ log: silentLog, max: 10, timeout: 10_000 }), - log: silentLog, + log: { ...silentLog, warn: message => { warnings.push(message); } }, sendPastDrain: () => Promise.resolve({ err: new Error('never sent') }), session, }); @@ -1557,12 +1558,14 @@ describe('held message bounds', () => { assert.equal(received.length, defaults.maxHeldMessages); assert.deepEqual(answers, ['ESME_RTHROTTLED', 'ESME_RTHROTTLED']); + assert.equal(warnings.length, 1, 'reaching the bound warns once, not per refusal'); - // The refused first segment joined no group, so the second one completes nothing. + // The refused first segment joined no group, so the second one is taken and completes nothing. await received[0]?.sendResp(); - await delay(0); + await new Promise(resolve => { setImmediate(resolve); }); await incoming.handle(segment(7, 2, 2)); + assert.equal(answers.at(-1), 'ESME_ROK'); assert.equal(received.length, defaults.maxHeldMessages); incoming.clear(); }); diff --git a/todo.md b/todo.md index 843a501..ee74247 100644 --- a/todo.md +++ b/todo.md @@ -405,6 +405,11 @@ and is also what the panel ranked hardest — two methods, one answer. ## Worth doing, not blocking +- [ ] **Make the dumbclient soak's memory sample evidence of no library leak again.** Its rss ends at + its maximum (298 MiB, heapUsed 81 MiB after 173,820 messages), which the harness's own per-id + `Set` and `answerOrder` explain but cannot separate from a leak in `src/`: sample the heap + after the bookkeeping is cleared. From the stability review of #29. + - [ ] **Decide whether `alert_notification` reaches the application as more than `incomingPduObj`.** It is the SMSC saying a handset it could not reach is reachable again (`esme_addr`, `ms_availability_status`); a client has no `onRequest`, so the raw PDU event is the only way in. -- 2.52.0 From f2977638f23e6ee6c0f84ded90eb1f6773994ab5 Mon Sep 17 00:00:00 2001 From: Lilleman auf Larv Date: Sat, 26 Sep 2026 14:44:00 +0200 Subject: [PATCH 4/5] Stop refusing only once the held messages are down to half, so a full window warns once, and name the octets in the warning --- README.md | 7 ++++--- interop-tests/findings/07-load.md | 4 ++-- src/held-messages.ts | 4 ++++ src/incoming-requests.ts | 28 ++++++++++++++++++++-------- test/session-extras.test.ts | 10 ++++++++++ 5 files changed, 40 insertions(+), 13 deletions(-) diff --git a/README.md b/README.md index 988e2e6..4f5f349 100644 --- a/README.md +++ b/README.md @@ -435,7 +435,8 @@ const { err, pduObj } = await session.send({ - **Unanswered messages.** While 1000 messages you have not called `sendResp()` on, or 64 MiB of them counted the way `maxOctets` counts segments, are held, every new message is refused with `ESME_RTHROTTLED` (`ESME_RX_T_APPN` on a `deliver_sm`) so the peer retries it, and no `sms` fires. - Reaching the bound logs one `warn`, and falling back below it one `info`. A message left five + Reaching the bound logs one `warn`, and the first message accepted once both are down to half one + `info`. A message left five minutes is dropped from the count with a `warn`; a later `sendResp()` still answers it. Neither bound is an option. - **Where the body is.** A body in the `message_payload` TLV, SMPP's way of carrying up to 64 KB and @@ -498,8 +499,8 @@ the unanswered messages are at their bound, which asks the SMSC to keep it and t `sms.answeredOnArrival` says whether the message you hold was answered that way; a segment count cannot, since a peer may number a message one part of one. -- The id was fixed with the first segment, so `sendResp()` there only says you are done, and - returns `err` for an `smsId` or a refusing `status`. +- The id was fixed with the first segment, so `sendResp()` there puts nothing on the wire and + releases the message, and returns `err` for an `smsId` or a refusing `status`. - `sms.smsId` is the base. `sendDlr()` names `-1`, `-2` and so on: the ids the `submit_sm` responses carried. - A `deliver_sm` is answered with no id at all, since SMPP marks that field unused, so an inbound diff --git a/interop-tests/findings/07-load.md b/interop-tests/findings/07-load.md index 3b6a73e..51ba1fd 100644 --- a/interop-tests/findings/07-load.md +++ b/interop-tests/findings/07-load.md @@ -101,8 +101,8 @@ event to a listener added after it fired. Fixed by attaching every session's `cl Three runs of `./interop-tests/run.py dumbclient`. Run 1 (the original 300,000-count soak) surfaced the S6 harness bug above; run 2 fixed it; run 3, with the namespace owner above and the held-message throttle in `src/`, is the one Scenarios reports. The capture figures below are run 2's. -`smppload.test.ts` passed on every run it was given (three, across the investigation above); its one scenario needs no repeat - a second run reproduces the identical corrupted PDU, -adding nothing. +`smppload.test.ts` passed on every run it was given (three, across the investigation above); its one +scenario needs no repeat - a second run reproduces the identical corrupted PDU, adding nothing. ``` dumbclient run 2: frames 111300, bind_transceiver 4/4, enquire_link 12 (enquire_link_resp 9 - the diff --git a/src/held-messages.ts b/src/held-messages.ts index adab745..9585382 100644 --- a/src/held-messages.ts +++ b/src/held-messages.ts @@ -44,6 +44,10 @@ export class HeldMessages { this.maxOctets = options.maxOctets; } + get octetsHeld(): number { + return this.octets; + } + get size(): number { return this.held.size; } diff --git a/src/incoming-requests.ts b/src/incoming-requests.ts index 817227e..64867c2 100644 --- a/src/incoming-requests.ts +++ b/src/incoming-requests.ts @@ -213,16 +213,13 @@ export class IncomingRequests { await this.session.sendReturn(pduObj); } - /** - * A concatenated message is answered segment by segment as it arrives: a peer that dispatches - * one request at a time never sends the second segment until the first has been answered. - */ - private async onMessage(pduObj: PduObject): Promise { + private async refusedAtBound(pduObj: PduObject): Promise { if (this.held.full()) { if (!this.refusing) { this.refusing = true; this.log.warn('session - unanswered messages at their bound, refusing new ones until the application answers', { messages: this.held.size, + octets: this.held.octetsHeld, }); } @@ -232,14 +229,29 @@ export class IncomingRequests { }); await this.session.sendReturn(pduObj, throttledStatus(this.carriedAs(pduObj))); - return; + return true; } - if (this.refusing) { + // Half, so a peer keeping its window full does not flip this on every answer. + if ( + this.refusing + && this.held.size <= defaults.maxHeldMessages / 2 + && this.held.octetsHeld <= defaults.maxHeldOctets / 2 + ) { this.refusing = false; - this.log.info('session - unanswered messages below their bound, accepting again', { messages: this.held.size }); + this.log.info('session - unanswered messages down to half their bound, accepting again', { messages: this.held.size }); } + return false; + } + + /** + * A concatenated message is answered segment by segment as it arrives: a peer that dispatches + * one request at a time never sends the second segment until the first has been answered. + */ + private async onMessage(pduObj: PduObject): Promise { + if (await this.refusedAtBound(pduObj)) return; + const concat = concatOf(pduObj); if (!concat) { diff --git a/test/session-extras.test.ts b/test/session-extras.test.ts index 98a540d..fa52f1a 100644 --- a/test/session-extras.test.ts +++ b/test/session-extras.test.ts @@ -1567,6 +1567,16 @@ describe('held message bounds', () => { assert.equal(answers.at(-1), 'ESME_ROK'); assert.equal(received.length, defaults.maxHeldMessages); + + // A peer keeping its window full crosses the bound on every answer, and that is still one warning. + await received[1]?.sendResp(); + await new Promise(resolve => { setImmediate(resolve); }); + await incoming.handle(submitPdu(defaults.maxHeldMessages + 2)); + await incoming.handle(submitPdu(defaults.maxHeldMessages + 3)); + await incoming.handle(submitPdu(defaults.maxHeldMessages + 4)); + + assert.equal(answers.at(-1), 'ESME_RTHROTTLED'); + assert.equal(warnings.length, 1); incoming.clear(); }); -- 2.52.0 From 42b1239f3d7df53dfac6d915cb1acc42f74b4413 Mon Sep 17 00:00:00 2001 From: Lilleman auf Larv Date: Sat, 26 Sep 2026 14:53:09 +0200 Subject: [PATCH 5/5] Name the retry status by direction, and correct what the held-message decision claims --- AGENTS.md | 4 ++-- CHANGELOG.md | 2 +- README.md | 13 ++++++------- docs/decisions.md | 36 ++++++++++++++++++------------------ src/incoming-requests.ts | 2 +- 5 files changed, 28 insertions(+), 29 deletions(-) diff --git a/AGENTS.md b/AGENTS.md index cc83610..b05736e 100644 --- a/AGENTS.md +++ b/AGENTS.md @@ -304,8 +304,8 @@ the file. - The drain's wait on the application ignores `shutdownTimeout: 0`. - What the application holds unanswered is capped on constants, and a message past the cap is refused. -- A store at its bound answers `ESME_RTHROTTLED` to a `submit_sm` and `ESME_RX_T_APPN` to a - `deliver_sm`. +- A store at its bound answers `ESME_RTHROTTLED` to a submission and `ESME_RX_T_APPN` to a delivery, + a `data_sm` by whichever it stands in for. - A reconnect keeps the delivery-receipt merges; everything else the link held is dropped. - A message id base is merged at most once. - A send that never reached the socket waits for the next link; one that did is counted, not resent. diff --git a/CHANGELOG.md b/CHANGELOG.md index fc2bcb3..3a0dbda 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -42,7 +42,7 @@ the cap is lost, since its segments were already answered. - A message arriving while the application holds 1000 unanswered, or 64 MiB of them counted the way `maxOctets` counts segments, is refused with `ESME_RTHROTTLED` (`ESME_RX_T_APPN` on a - `deliver_sm`), so the peer keeps it and retries. **Call `sendResp()` on every `sms`, multipart + delivery), so the peer keeps it and retries. **Call `sendResp()` on every `sms`, multipart included**: 1000 left unanswered now stop inbound traffic for up to five minutes, where the oldest used to be dropped with a warning. - A `submit_sm` segment the reassembly buffer has no room for is refused with `ESME_RTHROTTLED`, diff --git a/README.md b/README.md index 4f5f349..aee9db5 100644 --- a/README.md +++ b/README.md @@ -432,13 +432,12 @@ const { err, pduObj } = await session.send({ each is two messages. - **Answered on arrival.** Each segment was answered as it landed, before you see the message: [Server in depth](#server-in-depth). -- **Unanswered messages.** While 1000 messages you have not called `sendResp()` on, or 64 MiB of - them counted the way `maxOctets` counts segments, are held, every new message is refused with - `ESME_RTHROTTLED` (`ESME_RX_T_APPN` on a `deliver_sm`) so the peer retries it, and no `sms` fires. - Reaching the bound logs one `warn`, and the first message accepted once both are down to half one - `info`. A message left five - minutes is dropped from the count with a `warn`; a later `sendResp()` still answers it. Neither - bound is an option. +- **Unanswered messages.** While a session holds 1000 messages you have not called `sendResp()` on, + or 64 MiB of them counted the way `maxOctets` counts segments, every new message and segment is + refused so the peer retries it: `ESME_RTHROTTLED` on a submission, `ESME_RX_T_APPN` on a delivery. + No `sms` fires. Reaching the bound logs one `warn`, and the first message accepted once both are + down to half one `info`. A message left five minutes is dropped from the count with a `warn`; a + later `sendResp()` still answers it. None of the three is an option. - **Where the body is.** A body in the `message_payload` TLV, SMPP's way of carrying up to 64 KB and the only place a `data_sm` has, reads exactly like one in `short_message`, concatenated messages and receipts included. A PDU filling both is read from `short_message`. diff --git a/docs/decisions.md b/docs/decisions.md index ea9cad1..b77ffaa 100644 --- a/docs/decisions.md +++ b/docs/decisions.md @@ -679,25 +679,25 @@ rule and an index of the titles below. - **What the application holds unanswered is capped on constants, and a message past the cap is refused.** A bound the application cannot raise is the point: an application that answers nothing - would otherwise grow it for the life of the link, which goal 4 forbids. Reassembly's `maxOctets` - is an option because it bounds what the peer sends; this bounds what the application leaves - unanswered. Maintainer's call, 2026-09-26: a message arriving past the cap is refused with the - retry status, so the peer keeps it (goal 2). Rejected: dropping the oldest to make room, which - frees nothing while the application still holds its `Sms`, and stops the drain waiting for a - message the peer is owed. Rejected: pausing the socket, which also stalls every answer and - `enquire_link` on the link. Reaching the bound is a log line and nothing more (goal 8): an event - or a public count would be surface for what the application already knows, since it is the one - not answering. A message held past its timeout is still dropped, so `close()` can - report fewer unanswered than there were — accepted, because the alternative is holding what - nothing will answer. + would otherwise hold ever more messages, which goal 4 forbids. Reassembly's `maxOctets` is an + option because it bounds what the peer sends; this bounds what the application leaves unanswered. + Maintainer's call, 2026-09-26. Refusing leaves the message with the peer, which will send it again + (goal 2). Rejected: dropping the oldest to make room, which frees nothing while the application + still holds its `Sms`, and stops the drain waiting for a message the peer is owed. Rejected: + pausing the socket, which also stalls every answer and `enquire_link` on the link. Reaching the + bound shows only in the log (goal 8): an event or a public count would be surface for what the + application already knows, since it is the one not answering. A message held past its timeout is + still dropped, so `close()` can report fewer unanswered than there were — accepted, because the + alternative is holding what nothing will answer. -- **A store at its bound answers `ESME_RTHROTTLED` to a `submit_sm` and `ESME_RX_T_APPN` to a - `deliver_sm`.** Maintainer's call, 2026-09-26, for reassembly and held messages alike, so "keep it - and retry" has one spelling per direction. `ESME_RTHROTTLED` asks the sender to slow down, which - is what the peer outrunning us needs, and operators send it (Vonage, LINK Mobility, Route Mobile, - Jasmin), so clients built against them meet it (goal 1). Rejected: `ESME_RMSGQFUL`, which names an - exhausted queue and no rate. `ESME_RTHROTTLED` is the SMSC's to send, so an ESME answers with - SMPP 3.4's temporary receiver error, the one an SMSC retries on (goal 3). +- **A store at its bound answers `ESME_RTHROTTLED` to a submission and `ESME_RX_T_APPN` to a + delivery, a `data_sm` by whichever it stands in for.** Maintainer's call, 2026-09-26, for + reassembly and held messages alike, so "keep it and retry" has one spelling per direction. + `ESME_RTHROTTLED` asks the sender to slow down, which is what the peer outrunning us needs, and + operators send it (Vonage, LINK Mobility, Route Mobile, Jasmin), so clients built against them + meet it (goal 1). Rejected: `ESME_RMSGQFUL`, which names an exhausted queue and no rate. + `ESME_RTHROTTLED` is the SMSC's to send, so an ESME answers with SMPP 3.4's temporary receiver + error, the one an SMSC retries on (goal 3). - **A reconnect keeps the delivery-receipt merges; everything else the link held is dropped.** `onDelivery()` answers each receipt before the group it belongs to is complete, and `teardown()` diff --git a/src/incoming-requests.ts b/src/incoming-requests.ts index 64867c2..13e84fe 100644 --- a/src/incoming-requests.ts +++ b/src/incoming-requests.ts @@ -19,7 +19,7 @@ import { paramText } from './defs/types.ts'; import { respIdParams, segmentId } from './sms-id.ts'; import { respNameFor } from './defs/commands.ts'; -/** Asks the peer to keep the message and retry: RTHROTTLED is the SMSC's to send, so an ESME's is another. */ +/** Asks the peer to keep the message and retry. */ function throttledStatus(carriedAs: string): ErrorName { return carriedAs === 'submit_sm' ? 'ESME_RTHROTTLED' : 'ESME_RX_T_APPN'; } -- 2.52.0