From 3aadb64df8c81776896b1a30a7b3ba63af67e75e Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Mikael=20=27Lilleman=27=20G=C3=B6ransson?= Date: Wed, 2 Sep 2026 08:15:26 +0200 Subject: [PATCH] Release a message's hold when the response goes out, or when its listener failed --- AGENTS.md | 6 +++++- README.md | 2 +- src/incoming-requests.ts | 10 ++++++++++ src/session.ts | 4 +++- src/sms.ts | 10 ++++++---- test/session-extras.test.ts | 39 +++++++++++++++++++++++++++++++++++++ todo.md | 13 +------------ 7 files changed, 65 insertions(+), 19 deletions(-) diff --git a/AGENTS.md b/AGENTS.md index 8e083bc..bb7369c 100644 --- a/AGENTS.md +++ b/AGENTS.md @@ -351,7 +351,11 @@ Grouped by what each one constrains. drain waits for. Counting every inbound request until `sendReturn()` answered it was rejected — an `onRequest` that deliberately answers nothing would then cost a full `shutdownTimeout` on every close — and a message no listener took is released at once, since nothing is going to answer it. - `teardown()` drops what is still held for the same reason it drops inbound segments. The release + A listener that threw or rejected before answering releases it the same way, for the same reason, + with `sessionError` carrying the failure. What ends the wait is the response reaching the wire, not + the call: a `sendResp()` the library refused leaves the message held, so `close()` still reports the + one the peer is owed. `teardown()` drops what is still held for the same reason it drops inbound + segments. The release is one turn late, so a listener that sends its receipt straight after the response is still holding when the drain looks; `sendDlr()` is the one send that goes out past the drain's refusal, and only while the message is still held — past that it is an ordinary send, because the drain it diff --git a/README.md b/README.md index 48108aa..fb3d8de 100644 --- a/README.md +++ b/README.md @@ -332,7 +332,7 @@ TypeScript users can import `SmppLog` to have the compiler check one. `sendSms()`, `send()`, `sendReturn()`, `unbind()` and `close()`. Both `close()` and `unbind()` refuse further sends, wait out the requests this end already sent for up to `shutdownTimeout`, and then tear down whatever is left, resolving to an `err` that says what was lost. They also wait for -every `sms` the application has not called `sendResp()` on, so a peer whose `submit_sm` is still +every `sms` that `sendResp()` has not answered, so a peer whose `submit_sm` is still being handled is answered rather than left to re-send it — answering its PDUs through `sendReturn()` instead leaves that wait running until it gives up. `sendDlr()` is the one send the refusal lets past, and it catches the wait when issued straight after `sendResp()`; await anything in between and diff --git a/src/incoming-requests.ts b/src/incoming-requests.ts index 547015a..fca8fa7 100644 --- a/src/incoming-requests.ts +++ b/src/incoming-requests.ts @@ -39,6 +39,8 @@ export class IncomingRequests { private readonly session: Session; private readonly smsIdFormat: SmsIdFormat; private readonly systemId: string; + /** Identity, not a type guard: the rejected event carries the Sms back as an `unknown`. */ + private readonly emitted = new WeakMap void>(); private linkGeneration = 0; constructor(options: IncomingRequestsOptions) { @@ -111,6 +113,13 @@ export class IncomingRequests { this.reassembler.clear(); } + /** A listener that rejected answered nothing and never will, so a shutdown may not wait for it. */ + listenerRejected(sms: unknown): void { + if (typeof sms !== 'object' || sms === null) return; + + this.emitted.get(sms)?.(); + } + /** Waits out the messages the application still holds, and says how many it never answered. */ async drain(timeout: number, signal: AbortSignal | undefined): Promise { const unanswered = await this.held.idle(timeout, signal); @@ -191,6 +200,7 @@ export class IncomingRequests { }); this.held.hold(pduObjs); + this.emitted.set(sms, () => { this.held.release(pduObjs); }); // A message nobody took is not work a shutdown can wait for. if (!this.session.emit('sms', sms)) this.held.release(pduObjs); diff --git a/src/session.ts b/src/session.ts index 34868e9..9097030 100644 --- a/src/session.ts +++ b/src/session.ts @@ -93,11 +93,13 @@ export class Session extends EventEmitter { reason: unknown, ...args: [event: keyof SessionEvents, ...rest: unknown[]] ): void { - const [event] = args; + const [event, ...rest] = args; const error = errorFrom(reason); this.log.error('session - a listener rejected', { event, message: error.message }); + if (event === 'sms') this.incoming.listenerRejected(rest[0]); + if (event !== 'sessionError') this.emit('sessionError', error); } diff --git a/src/sms.ts b/src/sms.ts index 8b74b1d..750ea02 100644 --- a/src/sms.ts +++ b/src/sms.ts @@ -77,8 +77,7 @@ export function createSms(input: SmsInput, handlers: SmsHandlers): Sms { message: input.message, pduObjs: input.pduObjs, sendDlr: status => sendDlr(sms, handlers.send, status), - sendResp: options => sendResp(sms, answered, options ?? {}, handlers.lostLink) - .finally(handlers.onAnswered), + sendResp: options => sendResp(sms, answered, options ?? {}, handlers), session: input.session, get smsId(): string { return answered.smsId; @@ -94,7 +93,7 @@ async function sendResp( sms: Sms, answered: { smsId: string }, options: SendRespOptions, - lostLink: () => boolean, + handlers: SmsHandlers, ): Promise { const total = sms.pduObjs.length; @@ -109,7 +108,7 @@ async function sendResp( if (options.smsId !== undefined) answered.smsId = options.smsId; // A response carries the sequence number it was asked on, which the next link knows nothing about. - if (lostLink()) { + if (handlers.lostLink()) { return { err: new Error('The link this message arrived on is gone, so nothing would correlate the response') }; } @@ -119,6 +118,9 @@ async function sendResp( { message_id: segmentId(answered.smsId, index, total) }, ))); + // Only here: a refusal above put nothing on the wire, so the peer is still owed its response. + handlers.onAnswered(); + return results.find(result => result.err) ?? {}; } diff --git a/test/session-extras.test.ts b/test/session-extras.test.ts index 1608de2..f6b9212 100644 --- a/test/session-extras.test.ts +++ b/test/session-extras.test.ts @@ -1401,6 +1401,45 @@ describe('graceful shutdown', () => { assert.ok((await sent).err instanceof Error); }); + // emit() releases the hold of a listener that throws; one that rejects may cost no more than that. + test('a listener that rejected before answering does not hold the shutdown up', async t => { + const smpp = await startServer(t, { shutdownTimeout: 30_000 }); + const failed = once(resolve => { + smpp.on('session', bound => { + bound.on('sessionError', resolve); + bound.on('sms', () => Promise.reject(new Error('the listener gave up'))); + }); + }); + const { session } = await connect(t, smpp); + + assert.ok(session); + + const sent = session.sendSms({ + from: '46701113311', + message: 'the listener rejects', + to: '46709771337', + }); + + assert.equal((await failed).message, 'the listener gave up'); + + const started = Date.now(); + + assert.deepEqual(await peerOf(smpp).close(), {}); + assert.ok(Date.now() - started < 1000); + assert.ok((await sent).err instanceof Error); + }); + + // Nothing reached the peer, so a drain counting this answered would report an outcome that never was. + test('leaves a message the library refused to answer unanswered', async t => { + const { sms, smpp } = await submitInFlight(t, {}, { shutdownTimeout: 50 }); + const refused = await sms.sendResp({ smsId: '' }); + const closed = await peerOf(smpp).close(); + + assert.match(refused.err?.message ?? '', /smsId must not be empty/); + assert.ok(closed.err instanceof Error); + assert.match(closed.err.message, /1 message\(s\) unanswered/); + }); + test('gives up on a request that outlasts shutdownTimeout', async t => { const { sent, session } = await submitInFlight(t, { shutdownTimeout: 50 }); const closed = await session.close(); diff --git a/todo.md b/todo.md index 7762c37..b583cce 100644 --- a/todo.md +++ b/todo.md @@ -61,6 +61,7 @@ Rules the API follows: | Held messages capped and expiring, so an application that answers nothing cannot grow them | `test/session-extras.test.ts` | | A send with no link held for the next one, and one the link dropped under counted as `unanswered` | `test/session-extras.test.ts` | | A message whose link dropped refused an answer, with its receipt still allowed out | `test/session-extras.test.ts` | +| The hold released exactly when the peer was answered: a refused `sendResp()` keeps it, a listener that rejected drops it | `test/session-extras.test.ts` | | Every runnable README example | `test/readme.test.ts` | | Receipt-versus-message classification by `esm_class` | `test/dlr.test.ts`, `test/session.test.ts` | | A listener that throws, or rejects, reaching `sessionError`/`serverError` rather than the process | `test/session.test.ts`, `test/error-from.test.ts` | @@ -134,18 +135,6 @@ session message is a change to every call site. and applies it to the other is wrong. A budget type both take would close it. Raised by review, 2026-09-01. -- [ ] **An `sms` listener that rejects before answering costs a whole `shutdownTimeout`.** - One that *throws* is fine: `emit()` catches it, returns false, and `emitSms()` releases the - hold. A rejecting `async` one reaches `sessionError` through `captureRejections`, which hands - the handler an `unknown[]` the `Sms` cannot be read out of without a cast, so nothing releases - until the message expires. Same cost as the `onRequest`-answers-nothing case that was declined, - but reached by a bug rather than a policy. Raised by review, 2026-09-01. - -- [ ] **A refused `sendResp()` releases the hold anyway.** Every early return runs - `.finally(onAnswered)`, so `sendResp({ smsId: '' })` refuses and stops the drain waiting for a - message the peer was never answered, which `close()` then reports as answered. Raised by - review, 2026-09-01. - - [ ] **Does an intermediate delivery notification deserve to be a `dlr`?** `esm_class` message type `INTERMEDIATE_DELIVERY` (0x20) is classified as a message today, so a peer that reports non-final states with it hands the application a raw `id:… stat:ENROUTE` text as an inbound