diff --git a/AGENTS.md b/AGENTS.md index 882bb24..244455b 100644 --- a/AGENTS.md +++ b/AGENTS.md @@ -38,6 +38,7 @@ These are not preferences. Breaking one is a defect. ``` src/ index.ts Public surface. Named exports only, no default export. + bind-direction.ts What a bind declares, and which commands its direction carries at either end client.ts client() -> { err, session } server.ts server() -> { err, server }, server owns the listener + close() session.ts Session: the socket's life, dispatch, events, and the collaborators below @@ -45,17 +46,17 @@ src/ concat.ts How a PDU says it is a segment: its UDH, or the sar_* TLVs dlr.ts Delivery receipts: text and TLV parsing, receipt status codes dlr-merger.ts DlrMerger: per-segment receipts counted into one MessageDlr + drain.ts drain(): the shutdown's two waits and their budgets, and IdleWaiters, the wait on a count error-from.ts An untyped value as error material: errorFrom() an Error, namedValue() a name expiring-groups.ts ExpiringGroups: the capped, weighed, expiring store DlrMerger, HeldMessages and Reassembler share - held-messages.ts HeldMessages: a message from its `sms` event to its answer, capped and expiring, one MessageHold each - idle-waiters.ts IdleWaiters: waiting for a count to fall to zero, and what is left of a budget + held-messages.ts HeldMessages: a message from its `sms` event to its answer, and the six ways a hold ends incoming-requests.ts Every request the peer sends: messages, receipts, links, unknown commands - link-life.ts LinkLife: whether the link lives, and where a request waits for the next one + link-life.ts LinkLife: the link's life as one transition table, and where a request waits for the next link link-timers.ts LinkTimers: the enquire_link heartbeat and the idle timeout log.ts SmppLog, the logger contract, and silentLog — the default message.ts Encoding detection, splitting, bit counting, SMPP date formatting message-body.ts Where an inbound body is: short_message, or the message_payload TLV - outgoing-requests.ts OutgoingRequests: the window, the pending map and the retry + outgoing-requests.ts OutgoingRequests: the window, the pending map and the carry onto the next link pdu.ts pduToObj / objToPdu / pduReturn — synchronous, result-returning pdu-framer.ts PduFramer: a byte stream cut into complete PDUs pdu-refusal.ts A PDU the codec would not read, and the answer SMPP names for it @@ -67,7 +68,7 @@ src/ retained-pdu.ts A PDU copied off the wire so holding it pins nothing else, and what holding it costs send-sms.ts submitSms composition and the submitSmParams builder send-window.ts SendWindow: the maxOutstanding semaphore - session-options.ts SessionOptions, ReconnectOptions, bind direction and the session defaults + session-options.ts SessionOptions, ReconnectOptions, their checks and the session defaults sms-id.ts Message ids: the peer's notation, the - a segment gets, which response carries one udh.ts User data header: its length, the concatenation fields of a long SMS and their reference unanswered-error.ts UnansweredError: it went out and no answer came back diff --git a/DESIGN.md b/DESIGN.md new file mode 100644 index 0000000..b4c948e --- /dev/null +++ b/DESIGN.md @@ -0,0 +1,56 @@ +# Draft B: the session's life as one table, the hold as one file, the drain as one function + +## Structure and who owns what + +| Module | Owns | +| --- | --- | +| `link-life.ts` | The link's phase (`up`, `binding`, `down`, `closed`) and the `stopping` flag. `transition(event)` is the whole lifecycle: one switch, five events (`attached`, `bound`, `lost`, `stopping`, `closed`), returning the effects the session runs in order (`dropLink`, `emitDisconnected`, `scheduleReconnect`, `emitClose`, `linkUp`, `stopReconnect`). Also the waiters for a request with no link, released by the same transitions. | +| `session.ts` | Wiring, and the imperative shell: `apply(event)` runs `link.transition(event)` and then `run(effect)`, a switch that names every side effect the lifecycle has. No `end`/`stop`/`emitClose`/`teardown`/`onClose` — every path in (socket close, unreadable stream, idle timeout, failed rebind, `close()`, `unbind()`) is one `apply()` call. | +| `held-messages.ts` | The whole held-message flow: the store, the `Sms` creation, the `sms` emit, and each of the six exits as a numbered method. Owned by `Session` directly, so a rejected listener is routed `Session → HeldMessages` in one hop. The hysteresis logging (`refusing`) moved here from `IncomingRequests`: `full()` decides and logs. | +| `drain.ts` | `drain()`: the shutdown's two waits, both budgets, the one turn given to the application, and the error text. `IdleWaiters` lives here too, as the primitive the two stores wait with. | +| `incoming-requests.ts` | Routing only: the hook, the bind gate, the dispatch per command, reassembly. It is handed `held` and never constructs, counts or drains it. | +| `outgoing-requests.ts` | `request()` (refused after `stopping`), `requestPastDrain()` (a receipt), and `carry()`: a recursion in place of `for (;;)`, one attempt per link, the same budget throughout. `idle()` reports a count; the drain owns the words. | +| `bind-direction.ts` | `bindCommands`, `BindType`, `LinkEnd`, `standsInFor`, `bindCarries`, `checkedBind`, split out of `session-options.ts`, which now holds option types, their checks and the defaults. | + +`hold` now has one meaning: none. `LinkLife.budget()` is a request's budget, `HeldMessages` keeps messages, `MessageHold` is gone. + +## How each exit reads now + +A hold is `Held = { key, pduObjs, working }` in `HeldMessages`; `offer()` builds the `Sms` with three closures and emits it. + +1. **Answered** — `sendResp()` puts the response on the wire, calls `answered()`, which is `release(held)`. Synchronous, no `setImmediate`. +2. **Every listener rejected** — `Session[captureRejectionSymbol]` calls `held.listenerRejected(sms)`; the `WeakMap` finds the `Held`, `working--`, the last one releases. +3. **No listener, or one threw** — `session.emit('sms')` returned false, so `offer()` releases at once. +4. **Re-used sequence number** — `keep()` replaces the entry; the old `Held` fails `release()`'s identity check and can free nothing. +5. **Deadline** — `sweep()`, before each `keep()` and on the store's timer. +6. **Link gone** — `clear()`, run by the `dropLink` effect. + +The receipt that used to depend on the deferred release is now the drain's business: `sendDlr()` always sends through `requestPastDrain()` (a receipt answers a message the drain waits on, whenever it is sent), and `drain()` waits one `setImmediate` turn after the last message is answered before reading the window, so a receipt sent straight after the answer is in the window by then. One line, in the one place that waits on the application. + +**Shutdown** — `close()`: `apply('stopping')` (effect `stopReconnect`, before the first await) → `drain()`: messages half on `shutdownTimeout` or the `responseTimeout` fallback, one turn, requests half on what is left → `apply('closed')`: `dropLink` if a socket is attached, then `emitClose`. `unbind()` is the same with the unbind request between the drain and `closed`. + +**Reconnect** — socket close, unreadable stream, idle timeout and a failed rebind all `apply('lost')`. While retrying: `down`, effects `dropLink`, `emitDisconnected`, `scheduleReconnect`. Not retrying: `dropLink`, `emitClose`, phase `closed`. The loop's `comeBackUp()` is `attach` → `apply('attached')` → bind → `apply('bound')`, which is `linkUp` (release held sends, timers, `reconnected`), or, if `stopping` landed meanwhile, the `lost` path. Re-entrancy: the phase is updated before any effect runs, so a listener that calls `close()` from `disconnected` sees `down`, and the remaining effects are no-ops. + +## Deleted + +- `MessageHold` class, `SmsHandlers`' `isHeld` branch, the `setImmediate` in `answered()`, `sendPastDrain` plumbing through `IncomingRequests`. +- `Session.end()`, `stop()`, `emitClose()`, `teardown()`, `onClose()`, `resetTimers()`, `attach()`; `LinkLife.drop()`'s returned event name and `end()`'s boolean, `isUp`/`isAttached`/`isOver`/`isStopped`/`retrying` as public predicates (`carries`, `attached`, `isStopping`, `awaitsNextLink`, `refusal` remain, all reading `phase`). +- `IncomingRequests.drain()`, `.listenerRejected()`, `.refusing`; `OutgoingRequests.drain()`; `Session.drain()`'s two-error merge and `answering()`. +- `idle-waiters.ts` (into `drain.ts`); `for (;;)` in `requestPastDrain`. + +## Tests + +`npm test`: lint and typecheck clean; 516 tests, 516 pass, 0 fail (baseline 515/515). + +Changed, all in `test/session-extras.test.ts`, all reaching into internals: + +- `incomingOn()` builds a `HeldMessages` (new `heldMessagesOn()` helper) instead of passing `sendPastDrain`; `standsInFor` imported from `bind-direction.ts`. +- "drops a message whose link went while onRequest was still running": `link.drop()` → `link.transition('lost')`. +- `LinkLife` describe: `drop/attach/open/stop/end` → `transition('lost'|'attached'|'bound'|'stopping'|'closed')`; the "names the event it warrants" test now asserts the effect lists, which is the same claim stated in the new vocabulary; one test added: a link bound after the shutdown began is ended, not brought up. +- Held-message bounds: `offer()` returns the `Sms`, so `first.isHeld()`/`replaced.isHeld()` became "answering the replaced message frees nothing, answering the first frees it" (size before/after `sendResp()`), and `answered.release()` became `await answered.sendResp()`; `heldOn()` stubs `sendReturn` for that. + +No assertion was weakened; the `session.test.ts`, `readme.test.ts` and the shutdown/receipt integration tests run unchanged. + +## Not changed on purpose + +`ExpiringGroups` stays mechanism-only: the three owners' policies genuinely differ (refuse, evict oldest, evict with `spent` memory). `ReconnectLoop.halted` stays: `client()` runs a loop with no session behind it for `fromStart`. The public surface, `index.ts` and README are untouched. diff --git a/docs/decisions.md b/docs/decisions.md index c245097..f98e81e 100644 --- a/docs/decisions.md +++ b/docs/decisions.md @@ -561,7 +561,7 @@ rule and an index of the titles below. - **A stream this library cannot frame is a dead link; one PDU it cannot parse is not.** Maintainer's call, 2026-08-31, narrowed 2026-09-05 via the interop plan: a `command_length` below 16 or above `maxPduLength` leaves nothing that can say where the next PDU starts, so it tears the - link down through `teardown()` and the reconnect loop retries it on a fresh socket with a fresh + link down through the `dropLink` effect and the reconnect loop retries it on a fresh socket with a fresh framer. Every other codec failure honoured `command_length`, so the stream is still in sync and the next PDU starts where it says — tearing the link down there cost one peer half its receipts and its MO to a reconnect loop (`interop-tests/findings/01-smscsim.md`), and left the peer waiting @@ -668,7 +668,7 @@ rule and an index of the titles below. the peer, whose every request is bounded by `responseTimeout` unless the caller set that to 0 as well, and unsafe for the application, which nothing bounds — `close()` is what you reach for when the application is stuck, so it may not block on the application coming unstuck. That half falls - back to `responseTimeout`, the same answer `LinkLife`'s hold already takes — and to that + back to `responseTimeout`, the same answer `LinkLife`'s budget 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, and a message past the cap is @@ -694,10 +694,10 @@ rule and an index of the titles below. 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()` + `onDelivery()` answers each receipt before the group it belongs to is complete, and `dropLink` runs on every path — an idle timeout and a failed rebind, not only `close()` — so clearing the merges there loses receipts no peer has a reason to send again. They are cleared where the session - is over instead. Inbound segments stay in `teardown()`: a concatenation reference is the + is over instead. Inbound segments stay in `dropLink`: a concatenation reference is the peer's own counter, so a half-arrived group kept across a drop would take a later message's segments as readily as the rest of its own, and goal 2 will not hand the application a message assembled that way. What goes there is traffic already answered, which is why each group reaches @@ -705,7 +705,7 @@ rule and an index of the titles below. - **The bind state is the session's, `bound()` alone writes it, and it holds through a reconnect's gap.** Maintainer's call, 2026-09-28. `client()` and `server()` record their bind through - `bound()`, the call a hand-wired session makes, so the state has one writer — goal 8. Clearing the state at `teardown()` was rejected: `bindAllows()` and + `bound()`, the call a hand-wired session makes, so the state has one writer — goal 8. Clearing the state at `dropLink` was rejected: `bindAllows()` and `acceptsOptionalParams()` then answer yes to everything while the link is down, so a receiver-bound client queues a `submit_sm` the peer refuses and a receipt built then carries TLVs a pre-3.4 peer must not get — goal 4. Valid while the reconnect loop binds again with the same bind @@ -757,7 +757,7 @@ rule and an index of the titles below. attached, which admits a send one round trip before the bind is answered, and collaborators that ask the session, which answered the same question two ways at admit and at release. `ReconnectLoop.halted` is the one other flag, because `client()` also runs a loop with no session - behind it for `fromStart`; a session's loop is stopped by `Session.stop()` alone. + behind it for `fromStart`; a session's loop is stopped by the `stopReconnect` effect alone. ## Internals and tests diff --git a/src/bind-direction.ts b/src/bind-direction.ts new file mode 100644 index 0000000..ac251da --- /dev/null +++ b/src/bind-direction.ts @@ -0,0 +1,73 @@ +import type { Result } from './result.ts'; +import { quoted } from './error-from.ts'; + +export const bindCommands: readonly string[] = [ + 'bind_receiver', + 'bind_transceiver', + 'bind_transmitter', +]; + +export type BindType = 'receiver' | 'transceiver' | 'transmitter'; + +/** Which end of the link a session is. Only `server()` is the SMSC; everything else is the ESME. */ +export type LinkEnd = 'esme' | 'smsc'; + +export function bindTypeFromCommand(cmdName: string): BindType | undefined { + if (cmdName === 'bind_receiver') return 'receiver'; + if (cmdName === 'bind_transceiver') return 'transceiver'; + if (cmdName === 'bind_transmitter') return 'transmitter'; + + return undefined; +} + +/** + * Which message-carrying command an inbound one stands in for. Every command but `data_sm` names + * its own direction; that one travels either way, so the end it arrived at is what says. + */ +export function standsInFor(cmdName: string, linkEnd: LinkEnd): string { + if (cmdName !== 'data_sm') return cmdName; + + return linkEnd === 'smsc' ? 'submit_sm' : 'deliver_sm'; +} + +/** + * Whether a bind direction carries a command at all. A receiver-bound ESME submits nothing and a + * transmitter-bound one is delivered nothing, whichever end of the link is looking. A session that + * has not bound carries everything, since nothing has declared a direction yet. + */ +export function bindCarries( + bindType: BindType | undefined, + cmdName: string, + linkEnd: LinkEnd, +): boolean { + const carried = standsInFor(cmdName, linkEnd); + + if (bindType === 'receiver') return carried !== 'submit_sm'; + if (bindType === 'transmitter') return carried !== 'deliver_sm'; + + return true; +} + +/** SMPP 3.4: a peer that declares no version at all is one from before optional parameters. */ +export const undeclaredInterfaceVersion = 0x00; + +export type SessionBind = { as: BindType; peerVersion: number }; + +function isBindType(value: unknown): value is BindType { + return typeof value === 'string' && bindTypeFromCommand(`bind_${value}`) !== undefined; +} + +/** A bind as `Session.bound()` records it: undefined declares no version, which is pre-3.4. */ +export function checkedBind(bindType: unknown, declaredVersion: unknown): Result<{ bind: SessionBind }> { + if (!isBindType(bindType)) { + return { err: new Error(`bindType must be receiver, transceiver or transmitter, the bind command's name without "bind_", got ${quoted(bindType)}`) }; + } + + if (declaredVersion === undefined) return { bind: { as: bindType, peerVersion: undeclaredInterfaceVersion } }; + + if (typeof declaredVersion !== 'number' || !Number.isInteger(declaredVersion) || declaredVersion < 0 || declaredVersion > 0xFF) { + return { err: new Error(`declaredVersion must be an integer 0-255, the interface_version param or the sc_interface_version TLV's tagValue, or undefined where the peer declared none, got ${quoted(declaredVersion)}`) }; + } + + return { bind: { as: bindType, peerVersion: declaredVersion } }; +} diff --git a/src/client.ts b/src/client.ts index 6b7299c..28d9fc9 100644 --- a/src/client.ts +++ b/src/client.ts @@ -1,6 +1,7 @@ import type { ConnectionOptions } from 'node:tls'; import type { Result, VoidResult } from './result.ts'; -import type { BindType, ReconnectOptions } from './session-options.ts'; +import type { BindType } from './bind-direction.ts'; +import type { ReconnectOptions } from './session-options.ts'; import type { SmppLog } from './log.ts'; import type { SmsIdFormat } from './sms-id.ts'; import type { Socket } from 'node:net'; diff --git a/src/drain.ts b/src/drain.ts new file mode 100644 index 0000000..5251ad0 --- /dev/null +++ b/src/drain.ts @@ -0,0 +1,112 @@ +import type { SmppLog } from './log.ts'; +import type { VoidResult } from './result.ts'; +import { defaults } from './session-options.ts'; + +/** Everything waiting for a count to fall to zero, and how such a wait is cut short. */ +export class IdleWaiters { + private readonly waiting: (() => void)[] = []; + + /** Wakes everything waiting, whatever the count reads now. */ + settle(): void { + for (const resolve of this.waiting.splice(0)) { + resolve(); + } + } + + /** + * Resolves 0 once nothing is left, or with what still is when the timeout or the signal cuts the + * wait short. A timeout of 0 waits forever. + */ + wait(remaining: () => number, timeout: number, signal: AbortSignal | undefined): Promise { + if (remaining() === 0) return Promise.resolve(0); + + if (signal?.aborted === true) return Promise.resolve(remaining()); + + return new Promise(resolve => { + let timer: NodeJS.Timeout | undefined = undefined; + const done = (): void => { + const index = this.waiting.indexOf(done); + + if (timer) clearTimeout(timer); + if (index !== -1) this.waiting.splice(index, 1); + + signal?.removeEventListener('abort', done); + resolve(remaining()); + }; + + if (timeout > 0) { + timer = setTimeout(done, timeout); + timer.unref(); + } + + signal?.addEventListener('abort', done, { once: true }); + this.waiting.push(done); + }); + } +} + +/** The two stores a shutdown waits on, and whether the link still carries what they hold. */ +export type Drainable = { + linkCarries: () => boolean; + /** Resolves 0 once the application has answered every message, or with how many it has not. */ + messagesUnanswered: (timeout: number, signal: AbortSignal | undefined) => Promise; + /** Resolves 0 once every request sent is answered, or with how many are not. */ + requestsUnfinished: (timeout: number, signal: AbortSignal | undefined) => Promise; +}; + +export type DrainOptions = { + log: SmppLog; + responseTimeout: number; + shutdownTimeout: number; + signal: AbortSignal | undefined; +}; + +/** What is left of a budget, in the shape a wait takes it: 0 waits forever. */ +function leftOf(deadline: number): number { + return deadline === 0 ? 0 : Math.max(1, deadline - Date.now()); +} + +/** The application half may never wait forever: nothing else ends that wait. */ +function answeringBudget(options: DrainOptions): number { + if (options.shutdownTimeout > 0) return options.shutdownTimeout; + + return options.responseTimeout > 0 ? options.responseTimeout : defaults.responseTimeout; +} + +function report(log: SmppLog, unanswered: number, unfinished: number): VoidResult { + const lost: string[] = []; + + if (unanswered > 0) { + log.warn('drain - shutting down with messages unanswered', { unanswered }); + lost.push(`Shut down with ${String(unanswered)} message(s) unanswered`); + } + + if (unfinished > 0) { + log.warn('drain - shutting down with requests unfinished', { unfinished }); + lost.push(`Shut down with ${String(unfinished)} request(s) unfinished`); + } + + return lost.length > 0 ? { err: new Error(lost.join('; ')) } : {}; +} + +/** + * Waits out the messages the application holds, then the requests on the wire. Answering a message + * can put a receipt on the wire; nothing on the wire produces a message, so that order covers both. + */ +export async function drain(stores: Drainable, options: DrainOptions): Promise { + // No bound link, so nothing is on the wire to wait out. + if (!stores.linkCarries()) return {}; + + const deadline = options.shutdownTimeout > 0 ? Date.now() + options.shutdownTimeout : 0; + const unanswered = await stores.messagesUnanswered(answeringBudget(options), options.signal); + + // One turn, so a receipt sent straight after the last answer is in the window before it is read. + await new Promise(resolve => { setImmediate(resolve); }); + + const unfinished = await stores.requestsUnfinished(leftOf(deadline), options.signal); + + // The link went before the drain finished, so an empty window says nothing about the peer. + if (!stores.linkCarries()) return { err: new Error('The session closed before the drain finished') }; + + return report(options.log, unanswered, unfinished); +} diff --git a/src/error-from.ts b/src/error-from.ts index 5d77236..e622a0a 100644 --- a/src/error-from.ts +++ b/src/error-from.ts @@ -15,3 +15,8 @@ const printable: readonly string[] = ['boolean', 'number', 'string']; export function namedValue(value: unknown): string { return printable.includes(typeof value) ? String(value) : typeof value; } + +/** A string in quotes, so an empty one and one of spaces are visible; anything else as namedValue(). */ +export function quoted(value: unknown): string { + return typeof value === 'string' ? JSON.stringify(value) : namedValue(value); +} diff --git a/src/held-messages.ts b/src/held-messages.ts index b9e740e..b8d004c 100644 --- a/src/held-messages.ts +++ b/src/held-messages.ts @@ -1,26 +1,29 @@ -import type { LinkLife } from './link-life.ts'; -import type { PduObject, PduObjectInput } from './pdu.ts'; -import type { Result } from './result.ts'; +import type { PduObject } from './pdu.ts'; import type { Session } from './session.ts'; -import type { SmsHandlers } from './sms.ts'; +import type { Sms, SmsHandlers } from './sms.ts'; import type { SmppLog } from './log.ts'; import { ExpiringGroups } from './expiring-groups.ts'; -import { IdleWaiters } from './idle-waiters.ts'; +import { IdleWaiters } from './drain.ts'; import { createSms } from './sms.ts'; import { retainedOctets } from './retained-pdu.ts'; export type HeldMessagesOptions = { - link: LinkLife; + /** Changes with every link, so a message can tell the one it arrived on is gone. */ + linkGeneration: () => number; log: SmppLog; max: number; maxOctets: number; /** Injected so expiry can be exercised without a wall clock. */ now?: (() => number) | undefined; - sendPastDrain: SmsHandlers['send']; + /** How a receipt goes out: past a shutdown's refusal, since it answers a message the drain waits for. */ + sendReceipt: SmsHandlers['send']; session: Session; timeout: number; }; +/** A message handed to the application, and how many of its listeners are still working on it. */ +type Held = { key: string; pduObjs: PduObject[]; working: number }; + /** The peer's own sequence number, which is what our answer to this message will carry. */ function keyOf(pduObjs: PduObject[]): string | undefined { const first = pduObjs[0]; @@ -28,69 +31,19 @@ function keyOf(pduObjs: PduObject[]): string | undefined { return first ? String(first.seqNr) : undefined; } -type HoldRoute = Pick; - /** - * One message offered to the application, and the handlers its `Sms` answers through. 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. + * The messages handed to the application that it has not answered yet, which is what a shutdown + * waits on. A hold ends the first of six ways, numbered below: 1 answered, 2 the last listener + * working on it rejecting, 3 no listener taking it or one throwing, 4 a later message on its + * sequence number, 5 its deadline, 6 the link going. */ -export class MessageHold implements SmsHandlers { - private readonly generation: number; - private readonly heldMessages: HeldMessages; - private readonly pduObjs: PduObject[]; - private readonly route: HoldRoute; - private working: number; - - constructor(heldMessages: HeldMessages, route: HoldRoute, pduObjs: PduObject[], listeners: number) { - this.generation = route.link.generation(); - this.heldMessages = heldMessages; - this.pduObjs = pduObjs; - this.route = route; - this.working = listeners; - } - - /** Whether a drain is still waiting for this message to be answered. */ - isHeld(): boolean { - return this.heldMessages.holds(this.pduObjs); - } - - /** A turn later, so a `sendDlr()` called straight after `sendResp()` still goes out past a drain. */ - answered(): void { - setImmediate(() => { this.release(); }); - } - - lostLink(): boolean { - return this.route.link.generation() !== this.generation; - } - - /** A rejection leaves the other listeners running, so only the last one to fail gives the message up. */ - listenerGaveUp(): void { - this.working--; - - if (this.working <= 0) this.answered(); - } - - /** At once, for a message nobody took or a listener threw on: that is not work a shutdown can wait for. */ - release(): void { - this.heldMessages.release(this.pduObjs); - } - - /** A receipt for a message still held is what a drain waits for, so it goes out past the drain. */ - send(input: PduObjectInput): Promise> { - return this.isHeld() ? this.route.sendPastDrain(input) : this.route.session.send(input); - } -} - -/** The messages handed to the application that it has not answered yet, held by their segments. */ export class HeldMessages { - private readonly held: ExpiringGroups; + private readonly held: ExpiringGroups; private readonly idleWaiters = new IdleWaiters(); - private readonly log: SmppLog; - 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(); - private readonly route: HoldRoute; + private readonly offered = new WeakMap(); + private readonly options: HeldMessagesOptions; + private refusing = false; constructor(options: HeldMessagesOptions) { this.held = new ExpiringGroups({ @@ -99,9 +52,7 @@ export class HeldMessages { onSweep: () => { this.sweep(); }, timeout: options.timeout, }); - this.log = options.log; - this.maxOctets = options.maxOctets; - this.route = { link: options.link, sendPastDrain: options.sendPastDrain, session: options.session }; + this.options = options; } get octetsHeld(): number { @@ -112,69 +63,85 @@ export class HeldMessages { return this.held.size; } - /** Whether a message arriving now is past the bound, once the expired are swept. */ + /** Whether a message arriving now is past the bound. Logs once on reaching it, once on coming back to half. */ full(): boolean { this.sweep(); - return this.held.full || this.held.weight >= this.maxOctets; - } - - private hold(key: string, pduObjs: PduObject[], listeners: number): MessageHold { - const hold = new MessageHold(this, this.route, pduObjs, listeners); + const { log, max, maxOctets } = this.options; - this.sweep(); + if (this.held.full || this.held.weight >= maxOctets) { + if (!this.refusing) { + this.refusing = true; + log.warn('heldMessages - unanswered messages at their bound, refusing new ones until the application answers', { + messages: this.size, + octets: this.octetsHeld, + }); + } - if (this.held.get(key)) { - this.log.warn('heldMessages - replacing a message on a re-used sequence number', { seqNr: Number(key) }); + return true; } - this.held.set(key, pduObjs); - this.held.weigh(key, pduObjs.reduce((sum, pduObj) => sum + retainedOctets(pduObj), 0)); + // Half, so a peer keeping its window full does not flip this on every answer. + if (this.refusing && this.size <= max / 2 && this.octetsHeld <= maxOctets / 2) { + this.refusing = false; + log.info('heldMessages - unanswered messages down to half their bound, accepting again', { messages: this.size }); + } - return hold; + return false; } - offer(pduObjs: PduObject[], answeredAs?: string): MessageHold | undefined { + /** Hands the message to the application as an `sms` event and holds it until one of the six ways out. */ + offer(pduObjs: PduObject[], answeredAs?: string): Sms | undefined { const key = keyOf(pduObjs); if (key === undefined) return undefined; - const hold = this.hold(key, pduObjs, this.route.session.listenerCount('sms')); - const sms = createSms({ answeredAs, pduObjs, session: this.route.session }, hold); + const { linkGeneration, sendReceipt, session } = this.options; + const generation = linkGeneration(); + const held = this.keep(key, pduObjs); + const sms = createSms({ answeredAs, pduObjs, session }, { + answered: () => { this.release(held); }, + lostLink: () => linkGeneration() !== generation, + send: sendReceipt, + }); - this.offered.set(sms, hold); + this.offered.set(sms, held); - if (!this.route.session.emit('sms', sms)) hold.release(); + // 3: not work a shutdown can wait for. + if (!session.emit('sms', sms)) this.release(held); - return hold; + return sms; } - /** One listener gave up on a message; the last one to do so is what releases it. */ + /** 2: 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(); - } + const held = this.offered.get(message); - holds(pduObjs: PduObject[]): boolean { - const key = keyOf(pduObjs); + if (!held) return; - return key !== undefined && this.held.get(key) === pduObjs; + held.working--; + + if (held.working <= 0) this.release(held); } - release(pduObjs: PduObject[]): void { - const key = keyOf(pduObjs); + /** 5: drops every message past its deadline. Runs before each keep and on its own timer. */ + sweep(): void { + const expired = this.held.takeExpired(); - // Identity, not the key: a wrapped sequence number must not release someone else's message. - if (key === undefined || this.held.get(key) !== pduObjs) return; + if (expired.length === 0) return; - this.held.delete(key); + this.options.log.warn('heldMessages - messages the application never answered', { + messages: expired.length, + }); this.settle(); } - /** Drops every message: their segments went with the link, so no answer of ours correlates now. */ + /** 6: drops every message. Their segments went with the link, so no answer of ours correlates now. */ clear(): void { this.held.takeAll(); + this.refusing = false; this.idleWaiters.settle(); } @@ -183,15 +150,28 @@ export class HeldMessages { return this.idleWaiters.wait(() => this.held.size, timeout, signal); } - /** Drops every message past its deadline. Runs before each hold and on its own timer. */ - sweep(): void { - const expired = this.held.takeExpired(); + /** 4 is in here: a re-used sequence number replaces the message held on it. */ + private keep(key: string, pduObjs: PduObject[]): Held { + this.sweep(); - if (expired.length === 0) return; + if (this.held.get(key)) { + this.options.log.warn('heldMessages - replacing a message on a re-used sequence number', { seqNr: Number(key) }); + } - this.log.warn('heldMessages - messages the application never answered', { - messages: expired.length, - }); + const held: Held = { key, pduObjs, working: this.options.session.listenerCount('sms') }; + + this.held.set(key, held); + this.held.weigh(key, pduObjs.reduce((sum, pduObj) => sum + retainedOctets(pduObj), 0)); + + return held; + } + + /** 1 comes through here, as do 2 and 3. */ + private release(held: Held): void { + // Identity, not the key: a wrapped sequence number must not release someone else's message. + if (this.held.get(held.key) !== held) return; + + this.held.delete(held.key); this.settle(); } diff --git a/src/idle-waiters.ts b/src/idle-waiters.ts deleted file mode 100644 index dd29a9e..0000000 --- a/src/idle-waiters.ts +++ /dev/null @@ -1,47 +0,0 @@ -/** What is left of a budget, in the shape a wait takes it: 0 waits forever. */ -export function leftOf(deadline: number): number { - return deadline === 0 ? 0 : Math.max(1, deadline - Date.now()); -} - -/** Everything waiting for a count to fall to zero, and how such a wait is cut short. */ -export class IdleWaiters { - private readonly waiting: (() => void)[] = []; - - /** Wakes everything waiting, whatever the count reads now. */ - settle(): void { - for (const resolve of this.waiting.splice(0)) { - resolve(); - } - } - - /** - * Resolves 0 once nothing is left, or with what still is when the timeout or the signal cuts the - * wait short. A timeout of 0 waits forever. - */ - wait(remaining: () => number, timeout: number, signal: AbortSignal | undefined): Promise { - if (remaining() === 0) return Promise.resolve(0); - - if (signal?.aborted === true) return Promise.resolve(remaining()); - - return new Promise(resolve => { - let timer: NodeJS.Timeout | undefined = undefined; - const done = (): void => { - const index = this.waiting.indexOf(done); - - if (timer) clearTimeout(timer); - if (index !== -1) this.waiting.splice(index, 1); - - signal?.removeEventListener('abort', done); - resolve(remaining()); - }; - - if (timeout > 0) { - timer = setTimeout(done, timeout); - timer.unref(); - } - - signal?.addEventListener('abort', done, { once: true }); - this.waiting.push(done); - }); - } -} diff --git a/src/incoming-requests.ts b/src/incoming-requests.ts index 51aeec7..2c13e49 100644 --- a/src/incoming-requests.ts +++ b/src/incoming-requests.ts @@ -1,18 +1,17 @@ import type { Concat } from './concat.ts'; import type { DlrMerger } from './dlr-merger.ts'; import type { ErrorName } from './defs/errors.ts'; -import type { HeldMessagesOptions } from './held-messages.ts'; +import type { HeldMessages } from './held-messages.ts'; import type { LinkLife } from './link-life.ts'; import type { LostGroup, Refusal } from './reassembly.ts'; import type { OnRequest } from './session-options.ts'; import type { PduObject } from './pdu.ts'; -import type { VoidResult } from './result.ts'; import type { Session } from './session.ts'; import type { SmppLog } from './log.ts'; import type { SmsIdFormat } from './sms-id.ts'; -import { HeldMessages } from './held-messages.ts'; import { Reassembler } from './reassembly.ts'; -import { bindCommands, defaults, standsInFor } from './session-options.ts'; +import { bindCommands, standsInFor } from './bind-direction.ts'; +import { defaults } from './session-options.ts'; import { concatOf } from './concat.ts'; import { detach } from './retained-pdu.ts'; import { dlrFromPdu } from './dlr.ts'; @@ -45,13 +44,13 @@ const lostReasons: Record = { export type IncomingRequestsOptions = { dlrMerger: DlrMerger; + held: HeldMessages; link: LinkLife; log: SmppLog; maxOctets?: number | undefined; maxReassembly?: number | undefined; onRequest?: OnRequest | undefined; reassemblyTimeout?: number | undefined; - sendPastDrain: HeldMessagesOptions['sendPastDrain']; session: Session; smsIdFormat?: SmsIdFormat | undefined; systemId?: string | undefined; @@ -68,19 +67,10 @@ export class IncomingRequests { private readonly session: Session; private readonly smsIdFormat: SmsIdFormat; private readonly systemId: string; - private refusing = false; constructor(options: IncomingRequestsOptions) { this.dlrMerger = options.dlrMerger; - this.held = new HeldMessages({ - link: options.link, - log: options.log, - max: defaults.maxHeldMessages, - maxOctets: defaults.maxHeldOctets, - sendPastDrain: options.sendPastDrain, - session: options.session, - timeout: defaults.heldMessageTimeout, - }); + this.held = options.held; this.link = options.link; this.log = options.log; this.onRequest = options.onRequest; @@ -148,28 +138,11 @@ export class IncomingRequests { } } - /** Drops the segments of every message that never became whole, and of every one still held. */ + /** Drops the segments of every message that never became whole. */ clear(): void { - this.refusing = false; - this.held.clear(); this.reassembler.clear(); } - listenerRejected(sms: unknown): void { - this.held.listenerRejected(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); - - if (unanswered === 0) return {}; - - this.log.warn('session - shutting down with messages unanswered', { timeout, unanswered }); - - return { err: new Error(`Shut down with ${String(unanswered)} message(s) unanswered`) }; - } - private async unhandled(pduObj: PduObject): Promise { if (bindCommands.includes(pduObj.cmdName)) { this.log.info('session - bind on an already bound session', { cmdName: pduObj.cmdName }); @@ -212,35 +185,15 @@ export class IncomingRequests { } 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, - }); - } - - this.log.verbose('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))); + if (!this.held.full()) return false; - return true; - } - - // 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 down to half their bound, accepting again', { messages: this.held.size }); - } + this.log.verbose('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 false; + return true; } /** diff --git a/src/link-life.ts b/src/link-life.ts index f44f2ea..d2927b0 100644 --- a/src/link-life.ts +++ b/src/link-life.ts @@ -4,14 +4,33 @@ import type { VoidResult } from './result.ts'; export type LinkLifeOptions = { log: SmppLog; now?: (() => number) | undefined; - /** Whether a dropped link is followed by another one until stop(). */ + /** Whether a lost link is followed by another one, until the session stops. */ reconnects: boolean; /** How long a request may wait for a link. 0 waits for as long as one may still arrive. */ timeout: number; }; -/** `binding`: a socket is attached and its bind is not answered yet, so it carries nothing but that bind. */ -type Phase = 'binding' | 'down' | 'ended' | 'up'; +/** + * `up`: a bound socket carries requests. `binding`: a socket is attached and only its bind may go + * out. `down`: no socket, and the reconnect loop owes one. `closed`: over, and nothing brings it back. + */ +export type LinkPhase = 'binding' | 'closed' | 'down' | 'up'; + +/** + * `attached`: a socket from the reconnect loop. `bound`: its bind was answered. `lost`: the socket + * went, whoever noticed. `stopping`: a shutdown began, so no link follows this one. `closed`: the + * shutdown is over. + */ +export type LinkEvent = 'attached' | 'bound' | 'closed' | 'lost' | 'stopping'; + +/** What the session does after a transition, in the order returned. */ +export type LinkEffect = + | 'dropLink' + | 'emitClose' + | 'emitDisconnected' + | 'linkUp' + | 'scheduleReconnect' + | 'stopReconnect'; type Waiter = (result: VoidResult) => void; @@ -27,7 +46,10 @@ function over(): Error { return new Error('Session is closed'); } -/** Whether the session's link lives, and where a request with no link to go out on waits for the next one. */ +/** + * The link's life as one state machine: `transition()` is the whole table, and every predicate + * below reads the phase it keeps. Also where a request with no link waits for the next one. + */ export class LinkLife { private readonly log: SmppLog; private readonly now: () => number; @@ -35,8 +57,8 @@ export class LinkLife { private readonly timeout: number; private readonly waiting = new Set(); private drops = 0; - private phase: Phase = 'up'; - private stopped = false; + private linkPhase: LinkPhase = 'up'; + private stopping = false; constructor(options: LinkLifeOptions) { this.log = options.log; @@ -45,33 +67,49 @@ export class LinkLife { this.timeout = options.timeout; } - /** A socket is on the link, bound or not. */ - isAttached(): boolean { - return this.phase === 'binding' || this.phase === 'up'; + get phase(): LinkPhase { + return this.linkPhase; } - /** Whether a request can go out right now. */ - isUp(): boolean { - return this.phase === 'up'; + transition(event: LinkEvent): LinkEffect[] { + switch (event) { + case 'attached': + if (this.linkPhase === 'down') this.linkPhase = 'binding'; + + return []; + case 'bound': + if (this.linkPhase !== 'binding') return []; + + return this.stopping ? this.lose() : this.open(); + case 'closed': + return this.end(); + case 'lost': + return this.lose(); + case 'stopping': + this.stopping = true; + + return ['stopReconnect']; + } } - private isOver(): boolean { - return this.phase === 'ended'; + /** A socket is on the link, bound or not. */ + attached(): boolean { + return this.linkPhase === 'binding' || this.linkPhase === 'up'; } - /** The session is shutting down: nothing new is taken, and no link follows this one. */ - isStopped(): boolean { - return this.stopped; + /** Whether a request can go out right now. */ + carries(): boolean { + return this.linkPhase === 'up'; } - /** Whether a link that drops now is followed by another. */ - retrying(): boolean { - return this.reconnects && !this.stopped; + /** A shutdown began: nothing new is taken. */ + isStopping(): boolean { + return this.stopping; } - /** Not up and not over, with a link to come. */ + /** Not carrying, with a link still to come. */ awaitsNextLink(): boolean { - return !this.isUp() && !this.isOver() && this.retrying(); + return (this.linkPhase === 'down' || this.linkPhase === 'binding') && this.retrying(); } /** Changes with every drop, so what was read off one link can tell that link is gone. */ @@ -81,62 +119,60 @@ export class LinkLife { /** Why no request will ever be admitted, or undefined while one may still get through. */ refusal(): Error | undefined { - return this.isUp() || this.awaitsNextLink() ? undefined : over(); + return this.carries() || this.awaitsNextLink() ? undefined : over(); } /** One budget for a request, however many links it waits through. */ - hold(signal: AbortSignal | undefined): () => Promise { + budget(signal: AbortSignal | undefined): () => Promise { const deadline = this.timeout > 0 ? this.now() + this.timeout : 0; return () => this.wait(deadline, signal); } - /** A socket from the reconnect loop, not yet bound. An ended session stays ended. */ - attach(): void { - if (this.isOver()) return; - - this.phase = 'binding'; + private retrying(): boolean { + return this.reconnects && !this.stopping; } - /** The link is bound: everything held goes out on it. */ - open(): void { - this.phase = 'up'; + private open(): LinkEffect[] { + this.linkPhase = 'up'; if (this.waiting.size > 0) { this.log.verbose('linkLife - sending what was held for a link', { held: this.waiting.size }); } this.release({}); + + return ['linkUp']; } - /** The attached link is gone: the event that says so, or undefined when there was none to lose. */ - drop(): 'close' | 'disconnected' | undefined { - if (!this.isAttached()) return undefined; + private lose(): LinkEffect[] { + if (!this.attached()) return []; - this.phase = 'down'; this.drops++; + this.linkPhase = 'down'; - return this.retrying() ? 'disconnected' : 'close'; - } + if (!this.retrying()) return ['dropLink', ...this.end()]; - stop(): void { - this.stopped = true; + return ['dropLink', 'emitDisconnected', 'scheduleReconnect']; } - /** The session is over: nothing held will ever go out. False means it already was. */ - end(): boolean { - if (this.isOver()) return false; + private end(): LinkEffect[] { + if (this.linkPhase === 'closed') return []; + + const dropped = this.attached(); + + if (dropped) this.drops++; - this.phase = 'ended'; - this.stopped = true; + this.linkPhase = 'closed'; + this.stopping = true; this.release({ err: over() }); - return true; + return dropped ? ['dropLink', 'emitClose'] : ['emitClose']; } /** Resolves once a link can carry the request, or with the reason none ever will. */ private wait(deadline: number, signal: AbortSignal | undefined): Promise { - if (this.isUp()) return Promise.resolve({}); + if (this.carries()) return Promise.resolve({}); const refused = this.refusal(); diff --git a/src/outgoing-requests.ts b/src/outgoing-requests.ts index a0adf24..8c5571e 100644 --- a/src/outgoing-requests.ts +++ b/src/outgoing-requests.ts @@ -7,7 +7,7 @@ import type { SmppLog } from './log.ts'; import { PendingRequests } from './pending-requests.ts'; import { SendWindow } from './send-window.ts'; import { UnansweredError } from './unanswered-error.ts'; -import { bindCommands } from './session-options.ts'; +import { bindCommands } from './bind-direction.ts'; import { objToPdu } from './pdu.ts'; export type OutgoingRequestsOptions = { @@ -35,7 +35,6 @@ function misuse(input: PduObjectInput): Error | undefined { /** Everything this end asks of the peer: which link carries it, how many at once, and the answer. */ export class OutgoingRequests { private readonly link: LinkLife; - private readonly log: SmppLog; private readonly pending: PendingRequests; private readonly responseTimeout: number; private readonly transport: PduTransport; @@ -43,7 +42,6 @@ export class OutgoingRequests { constructor(options: OutgoingRequestsOptions) { this.link = options.link; - this.log = options.log; this.pending = new PendingRequests(options.log); this.responseTimeout = options.responseTimeout; this.transport = options.transport; @@ -51,7 +49,7 @@ export class OutgoingRequests { } canCarry(): boolean { - return this.link.isUp() && !this.transport.sock.destroyed; + return this.link.carries() && !this.transport.sock.destroyed; } /** The link is gone, and every answer still owed on it with it. */ @@ -76,44 +74,27 @@ export class OutgoingRequests { if (wrong) return Promise.resolve({ err: wrong }); // With no link, the request is refused as closed further on. - if (this.link.isStopped() && this.canCarry()) { + if (this.link.isStopping() && this.canCarry()) { return Promise.resolve({ err: new Error('Session is shutting down') }); } return this.requestPastDrain(input, options); } - /** request() without the drain's refusal, which a receipt for a held message has to take. */ - async requestPastDrain( - input: PduObjectInput, - options: SendOptions, - ): Promise> { + /** request() without the drain's refusal, which a receipt for a message the drain waits on has to take. */ + requestPastDrain(input: PduObjectInput, options: SendOptions): Promise> { const refused = this.refuse(input, options); - if (refused) return { err: refused }; + if (refused) return Promise.resolve({ err: refused }); // A bind is what makes a link usable, so it cannot wait for one: it takes the link's answer now. if (bindCommands.includes(input.cmdName)) { const shut = this.link.refusal(); - return shut ? { err: shut } : this.requestOnCurrentLink(input, options); + return shut ? Promise.resolve({ err: shut }) : this.requestOnCurrentLink(input, options); } - const waitForLink = this.link.hold(options.signal); - - for (;;) { - const held = await waitForLink(); - - if (held.err) return { err: held.err }; - - const slot = await this.window.acquire(options.signal); - - if (slot.err) return { err: slot.err }; - - const attempt = await this.attempt(input, options).finally(() => { this.window.release(); }); - - if (!this.retriesOnNextLink(attempt)) return attempt.result; - } + return this.carry(input, options, this.link.budget(options.signal)); } /** Straight onto the current link, for what has to go out either way. */ @@ -124,21 +105,31 @@ export class OutgoingRequests { return (await this.attempt(input, options)).result; } - /** Waits out the requests already on the wire, and says how many never finished. */ - async drain(timeout: number, signal: AbortSignal | undefined): Promise { - const unfinished = await this.window.idle(timeout, signal); + /** Resolves 0 once every request sent is answered, or with how many are not: on the wire, or queued behind the window. */ + idle(timeout: number, signal: AbortSignal | undefined): Promise { + return this.window.idle(timeout, signal); + } - if (unfinished === 0) return {}; + /** On the first link the budget admits; one the socket refused untouched waits for the next link and goes again. */ + private async carry( + input: PduObjectInput, + options: SendOptions, + waitForLink: () => Promise, + ): Promise> { + const held = await waitForLink(); - this.log.warn('outgoingRequests - shutting down with requests unfinished', { timeout, unfinished }); + if (held.err) return { err: held.err }; - return { err: new Error(`Shut down with ${String(unfinished)} request(s) unfinished`) }; - } + const slot = await this.window.acquire(options.signal); + + if (slot.err) return { err: slot.err }; + + const attempt = await this.attempt(input, options).finally(() => { this.window.release(); }); + + // Only once the link is dropped: until then the retry lands straight back on the dead socket. + if (attempt.retryOnNextLink && this.link.awaitsNextLink()) return this.carry(input, options, waitForLink); - /** Nothing reached the socket, so the next link carries it. */ - private retriesOnNextLink(attempt: Attempt): boolean { - // Until the link is dropped it admits the retry straight back onto the dead socket, and the loop spins. - return attempt.retryOnNextLink && this.link.awaitsNextLink(); + return attempt.result; } /** Why a request cannot go out at all, as opposed to not yet. */ diff --git a/src/send-window.ts b/src/send-window.ts index e13d67d..50294cb 100644 --- a/src/send-window.ts +++ b/src/send-window.ts @@ -1,6 +1,6 @@ import type { SmppLog } from './log.ts'; import type { VoidResult } from './result.ts'; -import { IdleWaiters } from './idle-waiters.ts'; +import { IdleWaiters } from './drain.ts'; export type SendWindowOptions = { limit: number; diff --git a/src/server.ts b/src/server.ts index 2b9aefe..bcdf532 100644 --- a/src/server.ts +++ b/src/server.ts @@ -1,4 +1,5 @@ -import type { BindType, CloseOptions, OnRequest } from './session-options.ts'; +import type { BindType } from './bind-direction.ts'; +import type { CloseOptions, OnRequest } from './session-options.ts'; import type { PduObject, TlvInputs } from './pdu.ts'; import type { Result, VoidResult } from './result.ts'; import type { Server as NetServer, Socket } from 'node:net'; @@ -6,7 +7,8 @@ import type { Server as TlsServer, TlsOptions } from 'node:tls'; import type { SmppLog } from './log.ts'; import { EventEmitter } from 'node:events'; import { Session, defaultSystemId } from './session.ts'; -import { bindTypeFromCommand, checkSessionOptions } from './session-options.ts'; +import { bindTypeFromCommand } from './bind-direction.ts'; +import { checkSessionOptions } from './session-options.ts'; import { createServer as createNetServer } from 'node:net'; import { createServer as createTlsServer } from 'node:tls'; import { defaultInterfaceVersion } from './defs/constants.ts'; diff --git a/src/session-options.ts b/src/session-options.ts index 0b768c9..93ae835 100644 --- a/src/session-options.ts +++ b/src/session-options.ts @@ -11,7 +11,7 @@ import type { Socket } from 'node:net'; import { backoffDefaults } from './reconnect-loop.ts'; import { defaultMaxOctets } from './reassembly.ts'; import { isSmsIdNotation, smsIdNotations, smsIdPlaces } from './sms-id.ts'; -import { namedValue } from './error-from.ts'; +import { namedValue, quoted } from './error-from.ts'; export type SessionEvents = { close: []; @@ -26,53 +26,6 @@ export type SessionEvents = { sms: [Sms]; }; -export const bindCommands: readonly string[] = [ - 'bind_receiver', - 'bind_transceiver', - 'bind_transmitter', -]; - -export type BindType = 'receiver' | 'transceiver' | 'transmitter'; - -/** Which end of the link a session is. Only `server()` is the SMSC; everything else is the ESME. */ -export type LinkEnd = 'esme' | 'smsc'; - -export function bindTypeFromCommand(cmdName: string): BindType | undefined { - if (cmdName === 'bind_receiver') return 'receiver'; - if (cmdName === 'bind_transceiver') return 'transceiver'; - if (cmdName === 'bind_transmitter') return 'transmitter'; - - return undefined; -} - -/** - * Which message-carrying command an inbound one stands in for. Every command but `data_sm` names - * its own direction; that one travels either way, so the end it arrived at is what says. - */ -export function standsInFor(cmdName: string, linkEnd: LinkEnd): string { - if (cmdName !== 'data_sm') return cmdName; - - return linkEnd === 'smsc' ? 'submit_sm' : 'deliver_sm'; -} - -/** - * Whether a bind direction carries a command at all. A receiver-bound ESME submits nothing and a - * transmitter-bound one is delivered nothing, whichever end of the link is looking. A session that - * has not bound carries everything, since nothing has declared a direction yet. - */ -export function bindCarries( - bindType: BindType | undefined, - cmdName: string, - linkEnd: LinkEnd, -): boolean { - const carried = standsInFor(cmdName, linkEnd); - - if (bindType === 'receiver') return carried !== 'submit_sm'; - if (bindType === 'transmitter') return carried !== 'deliver_sm'; - - return true; -} - export type SendOptions = { signal?: AbortSignal | undefined }; /** An already-aborted signal skips the drain; one that fires during it cuts the wait short. */ @@ -118,34 +71,6 @@ export type SessionOptions = { export const defaultSystemId = ''; -/** SMPP 3.4: a peer that declares no version at all is one from before optional parameters. */ -export const undeclaredInterfaceVersion = 0x00; - -export type SessionBind = { as: BindType; peerVersion: number }; - -function quoted(value: unknown): string { - return typeof value === 'string' ? JSON.stringify(value) : namedValue(value); -} - -function isBindType(value: unknown): value is BindType { - return typeof value === 'string' && bindTypeFromCommand(`bind_${value}`) !== undefined; -} - -/** A bind as `Session.bound()` records it: undefined declares no version, which is pre-3.4. */ -export function checkedBind(bindType: unknown, declaredVersion: unknown): Result<{ bind: SessionBind }> { - if (!isBindType(bindType)) { - return { err: new Error(`bindType must be receiver, transceiver or transmitter, the bind command's name without "bind_", got ${quoted(bindType)}`) }; - } - - if (declaredVersion === undefined) return { bind: { as: bindType, peerVersion: undeclaredInterfaceVersion } }; - - if (typeof declaredVersion !== 'number' || !Number.isInteger(declaredVersion) || declaredVersion < 0 || declaredVersion > 0xFF) { - return { err: new Error(`declaredVersion must be an integer 0-255, the interface_version param or the sc_interface_version TLV's tagValue, or undefined where the peer declared none, got ${quoted(declaredVersion)}`) }; - } - - return { bind: { as: bindType, peerVersion: declaredVersion } }; -} - export const defaults = { /** Receipts of a multipart message can be a working day apart, so the cap does the bounding. */ dlrMergeTimeout: 86_400_000, diff --git a/src/session.ts b/src/session.ts index 1fd9b46..b32419a 100644 --- a/src/session.ts +++ b/src/session.ts @@ -1,25 +1,29 @@ +import type { BindType, LinkEnd, SessionBind } from './bind-direction.ts'; import type { ErrorName } from './defs/errors.ts'; +import type { LinkEffect, LinkEvent } from './link-life.ts'; import type { MessageDlr } from './dlr-merger.ts'; import type { ParamValue } from './defs/types.ts'; import type { PduObject, PduObjectInput, TlvInputs } from './pdu.ts'; import type { PduRefusedError } from './pdu-refusal.ts'; -import type { BindType, CloseOptions, LinkEnd, ReconnectOptions, SendOptions, SessionBind, SessionEvents, SessionOptions } from './session-options.ts'; +import type { CloseOptions, ReconnectOptions, SendOptions, SessionEvents, SessionOptions } from './session-options.ts'; import type { Result, VoidResult } from './result.ts'; import type { SendSmsOptions, SendSmsResult } from './send-sms.ts'; import type { SmppLog } from './log.ts'; import type { Socket } from 'node:net'; import { DlrMerger } from './dlr-merger.ts'; import { EventEmitter } from 'node:events'; +import { HeldMessages } from './held-messages.ts'; import { IncomingRequests } from './incoming-requests.ts'; import { LinkLife } from './link-life.ts'; import { LinkTimers } from './link-timers.ts'; import { OutgoingRequests } from './outgoing-requests.ts'; import { PduTransport } from './pdu-transport.ts'; import { ReconnectLoop } from './reconnect-loop.ts'; -import { leftOf } from './idle-waiters.ts'; +import { bindCarries, bindCommands, checkedBind } from './bind-direction.ts'; +import { drain } from './drain.ts'; import { errorFrom } from './error-from.ts'; import { optionalParamsMinVersion } from './defs/constants.ts'; -import { bindCarries, bindCommands, checkedBind, defaultSystemId, defaults } from './session-options.ts'; +import { defaultSystemId, defaults } from './session-options.ts'; import { isResp, objToPdu, pduReturn } from './pdu.ts'; import { refusalAnswer } from './pdu-refusal.ts'; import { guardedLog } from './log.ts'; @@ -61,6 +65,7 @@ export class Session extends EventEmitter { private readonly concatReference = new ConcatReference(); private readonly dlrMerger: DlrMerger; + private readonly held: HeldMessages; private readonly incoming: IncomingRequests; private readonly link: LinkLife; private readonly options: SessionOptions; @@ -98,7 +103,7 @@ export class Session extends EventEmitter { this.log.error('session - a listener rejected', { event, message: error.message }); - if (event === 'sms') this.incoming.listenerRejected(rest[0]); + if (event === 'sms') this.held.listenerRejected(rest[0]); if (event !== 'sessionError') this.emit('sessionError', error); } @@ -119,8 +124,7 @@ export class Session extends EventEmitter { idleTimeout: options.idleTimeout, log: this.log, onEnquireLink: () => { void this.send({ cmdName: 'enquire_link' }); }, - // Not close(): a link that went quiet is a drop, and a drop is what reconnect is for. - onIdle: () => { this.teardown(); }, + onIdle: () => { this.apply('lost'); }, }); this.transport = this.transportFor(options.sock); this.outgoing = new OutgoingRequests({ @@ -130,21 +134,34 @@ export class Session extends EventEmitter { responseTimeout, transport: this.transport, }); - this.incoming = new IncomingRequests({ + this.held = new HeldMessages({ + linkGeneration: () => this.link.generation(), + log: this.log, + max: defaults.maxHeldMessages, + maxOctets: defaults.maxHeldOctets, + sendReceipt: input => this.outgoing.requestPastDrain(input, {}), + session: this, + timeout: defaults.heldMessageTimeout, + }); + this.incoming = this.incomingFor(options, this.held); + + this.timers.reset(); + } + + private incomingFor(options: SessionOptions, held: HeldMessages): IncomingRequests { + return new IncomingRequests({ dlrMerger: this.dlrMerger, + held, link: this.link, log: this.log, maxOctets: options.maxOctets, maxReassembly: options.maxReassembly, onRequest: options.onRequest, reassemblyTimeout: options.reassemblyTimeout, - sendPastDrain: input => this.outgoing.requestPastDrain(input, {}), session: this, smsIdFormat: options.smsIdFormat, systemId: options.systemId, }); - - this.resetTimers(); } /** Replaced on reconnect, so hold the session rather than this. */ @@ -196,22 +213,6 @@ export class Session extends EventEmitter { return Promise.resolve(this.answer(pduReturn(pdu, status, params, tlvs), pdu.cmdName, pdu.seqNr)); } - private answer(built: Result<{ buffer: Buffer }>, cmdName: string, seqNr: number): VoidResult { - const sent = built.err ? { err: built.err } : this.transport.write(built.buffer); - - // A peer that unbinds and drops the link takes our response with it; that is not a failure. - if (sent.err && this.link.isAttached()) { - this.log.warn('session - could not answer a request', { - cmdName, - message: sent.err.message, - seqNr, - }); - this.emit('sessionError', sent.err); - } - - return sent; - } - async sendSms(sms: SendSmsOptions, options: SendOptions = {}): Promise { if (!this.bindAllows('submit_sm')) { return unsent(new Error('A receiver-bound session does not carry submit_sm')); @@ -235,13 +236,13 @@ export class Session extends EventEmitter { */ async unbind(): Promise { const drained = await this.drain(undefined); - const wasOpen = this.link.isAttached(); + const wasOpen = this.link.attached(); const sent = wasOpen ? await this.outgoing.requestOnCurrentLink({ cmdName: 'unbind' }) : { err: new Error('Session is closed') }; - const closedOnUnbind = wasOpen && !this.link.isAttached(); + const closedOnUnbind = wasOpen && !this.link.attached(); - this.end(); + this.apply('closed'); return sent.err && !closedOnUnbind ? { err: sent.err } : drained; } @@ -253,15 +254,84 @@ export class Session extends EventEmitter { async close(options: CloseOptions = {}): Promise { const drained = await this.drain(options.signal); - this.end(); + this.apply('closed'); return drained; } + /** Runs the link's transition and then, in order, what it asks of the session. */ + private apply(event: LinkEvent): void { + for (const effect of this.link.transition(event)) { + this.run(effect); + } + } + + private run(effect: LinkEffect): void { + switch (effect) { + case 'dropLink': + this.timers.clear(); + this.held.clear(); + this.incoming.clear(); + this.outgoing.linkLost(); + this.sock.destroy(); + break; + case 'emitClose': + this.dlrMerger.clear(); + this.emit('close'); + break; + case 'emitDisconnected': + this.emit('disconnected'); + break; + case 'linkUp': + this.timers.reset(); + this.log.info('session - reconnected'); + this.emit('reconnected'); + break; + case 'scheduleReconnect': + this.reconnectLoop?.schedule(); + break; + case 'stopReconnect': + this.reconnectLoop?.stop(); + break; + } + } + + /** Stops new sends before its first await, then waits out what the session holds. */ + private drain(signal: AbortSignal | undefined): Promise { + this.apply('stopping'); + + return drain({ + linkCarries: () => this.outgoing.canCarry(), + messagesUnanswered: (timeout, cut) => this.held.idle(timeout, cut), + requestsUnfinished: (timeout, cut) => this.outgoing.idle(timeout, cut), + }, { + log: this.log, + responseTimeout: this.options.responseTimeout ?? defaults.responseTimeout, + shutdownTimeout: this.options.shutdownTimeout ?? defaults.shutdownTimeout, + signal, + }); + } + + private answer(built: Result<{ buffer: Buffer }>, cmdName: string, seqNr: number): VoidResult { + const sent = built.err ? { err: built.err } : this.transport.write(built.buffer); + + // A peer that unbinds and drops the link takes our response with it; that is not a failure. + if (sent.err && this.link.attached()) { + this.log.warn('session - could not answer a request', { + cmdName, + message: sent.err.message, + seqNr, + }); + this.emit('sessionError', sent.err); + } + + return sent; + } + private transportFor(sock: Socket): PduTransport { return new PduTransport({ log: this.log, - onClose: () => { this.onClose(); }, + onClose: () => { this.apply('lost'); }, onData: chunk => { this.onData(chunk); }, onError: err => { this.emit('sessionError', err); }, onFramed: pdu => { this.emit('incomingPdu', pdu); }, @@ -269,7 +339,7 @@ export class Session extends EventEmitter { onRefused: refused => { this.refuse(refused); }, onUnreadable: err => { this.emit('sessionError', err); - this.teardown(); + this.apply('lost'); }, }, sock); } @@ -290,109 +360,27 @@ export class Session extends EventEmitter { sock: Socket, bind: (session: Session) => Promise, ): Promise { - this.attach(sock); + this.transport.attach(sock); + this.apply('attached'); const bound = await bind(this); if (bound.err) { - this.teardown(); + this.apply('lost'); return { err: bound.err }; } - // close() can land while the rebind is in flight. - if (!this.link.retrying()) { - this.teardown(); - - return { err: new Error('Session closed while it was coming back up') }; - } - - this.resetTimers(); - this.link.open(); - this.log.info('session - reconnected'); - this.emit('reconnected'); - - return {}; - } - - private attach(sock: Socket): void { - this.transport.attach(sock); - this.link.attach(); - } - - /** Stops new sends and waits out the messages we hold and the requests already issued. */ - private async drain(signal: AbortSignal | undefined): Promise { - this.stop(); - - // No bound link, so nothing is on the wire to wait out. - if (!this.outgoing.canCarry()) return {}; - - const timeout = this.options.shutdownTimeout ?? defaults.shutdownTimeout; - const deadline = timeout > 0 ? Date.now() + timeout : 0; - // Answering a message can put a receipt on the wire; nothing on the wire produces a message. - const messages = await this.incoming.drain(this.answering(timeout), signal); - const requests = await this.outgoing.drain(leftOf(deadline), signal); - - // The link went before the drain finished, so an empty window says nothing about the peer. - if (!this.outgoing.canCarry()) { - return { err: new Error('The session closed before the drain finished') }; - } - - if (!messages.err) return requests; - - if (!requests.err) return messages; - - return { err: new Error(`${messages.err.message}; ${requests.err.message}`) }; - } - - /** The application half's budget, which may never be "forever": nothing else ends that wait. */ - private answering(timeout: number): number { - if (timeout > 0) return timeout; - - const responseTimeout = this.options.responseTimeout ?? defaults.responseTimeout; + this.apply('bound'); - return responseTimeout > 0 ? responseTimeout : defaults.responseTimeout; - } - - /** The session is over now, drained or not. Nothing brings it back. */ - private end(): void { - this.stop(); - this.teardown(); - this.dlrMerger.clear(); - this.emitClose(); - } - - /** No new sends, and no link after this one. */ - private stop(): void { - this.link.stop(); - this.reconnectLoop?.stop(); - } - - private emitClose(): void { - if (!this.link.end()) return; - - this.outgoing.linkLost(); - this.emit('close'); - } - - private teardown(): void { - const lost = this.link.drop(); - - if (!lost) return; - - this.outgoing.linkLost(); - this.timers.clear(); - this.incoming.clear(); - this.sock.destroy(); - - // `lost` is read before clear(): a listener it reaches may close() the session, and the drop still reports as disconnected. - if (lost === 'disconnected') this.emit('disconnected'); - else this.emitClose(); + // close() can land while the rebind is in flight, and 'bound' then ends the link instead. + return this.link.carries() ? {} : { err: new Error('Session closed while it was coming back up') }; } private onData(chunk: Buffer): void { this.emit('data', chunk); - this.resetTimers(); + + if (this.link.attached()) this.timers.reset(); } private dispatch(pduObj: PduObject): void { @@ -433,21 +421,4 @@ export class Session extends EventEmitter { this.answer(objToPdu({ ...refusalAnswer(refused), seqNr }), cmdName ?? String(cmdId), seqNr); } - - private resetTimers(): void { - if (!this.link.isAttached()) return; - - this.timers.reset(); - } - - private onClose(): void { - if (this.link.retrying()) { - this.teardown(); - this.reconnectLoop?.schedule(); - - return; - } - - this.end(); - } } diff --git a/test/session-extras.test.ts b/test/session-extras.test.ts index 6b36279..25d8f74 100644 --- a/test/session-extras.test.ts +++ b/test/session-extras.test.ts @@ -5,7 +5,7 @@ import type { Collected, LostGroup } from '../src/reassembly.ts'; import type { Dlr } from '../src/dlr.ts'; import type { ErrorName } from '../src/defs/errors.ts'; import type { IncomingRequestsOptions } from '../src/incoming-requests.ts'; -import type { HeldMessagesOptions, MessageHold } from '../src/held-messages.ts'; +import type { HeldMessagesOptions } from '../src/held-messages.ts'; import type { MessageState } from '../src/defs/constants.ts'; import type { MessageDlr } from '../src/session.ts'; import type { PduObject, PduObjectInput } from '../src/pdu.ts'; @@ -26,7 +26,8 @@ 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, defaults, standsInFor } from '../src/session-options.ts'; +import { checkSessionOptions, defaults } from '../src/session-options.ts'; +import { standsInFor } from '../src/bind-direction.ts'; import { client } from '../src/client.ts'; import { closeAfter, closeListenerAfter } from './teardown.ts'; import { concatOf } from '../src/concat.ts'; @@ -101,12 +102,30 @@ function abortAfter( }); } +function heldMessagesOn( + session: Session, + options: Partial> = {}, +): HeldMessages { + return new HeldMessages({ + linkGeneration: () => 0, + log: silentLog, + max: defaults.maxHeldMessages, + maxOctets: defaults.maxHeldOctets, + sendReceipt: () => Promise.resolve({ err: new Error('never sent') }), + session, + timeout: defaults.heldMessageTimeout, + ...options, + }); +} + function incomingOn(session: Session, options: Partial = {}): IncomingRequests { + const log = options.log ?? silentLog; + return new IncomingRequests({ - dlrMerger: new DlrMerger({ log: silentLog, max: 10, timeout: 10_000 }), - link: new LinkLife({ log: silentLog, reconnects: false, timeout: 100 }), - log: silentLog, - sendPastDrain: () => Promise.resolve({ err: new Error('never sent') }), + dlrMerger: new DlrMerger({ log, max: 10, timeout: 10_000 }), + held: heldMessagesOn(session, { log }), + link: new LinkLife({ log, reconnects: false, timeout: 100 }), + log, session, ...options, }); @@ -757,7 +776,7 @@ describe('reconnect', () => { const handled = incoming.handle(submitPdu(1)); - link.drop(); + link.transition('lost'); await handled; @@ -1407,9 +1426,9 @@ describe('LinkLife', () => { test('refuses a hold whose deadline has already passed', async () => { let now = 0; const link = new LinkLife({ log: silentLog, now: () => now, reconnects: true, timeout: 100 }); - const waitForLink = link.hold(undefined); + const waitForLink = link.budget(undefined); - link.drop(); + link.transition('lost'); now = 101; const held = await waitForLink(); @@ -1421,14 +1440,15 @@ describe('LinkLife', () => { const link = new LinkLife({ log: silentLog, reconnects: true, timeout: 10_000 }); const timers = (): number => process.getActiveResourcesInfo().filter(name => name === 'Timeout').length; - link.drop(); + link.transition('lost'); const before = timers(); - const held = link.hold(undefined)(); + const held = link.budget(undefined)(); assert.equal(timers(), before + 1, 'an unref\'d timer is not counted here, which is the point'); - link.open(); + link.transition('attached'); + link.transition('bound'); assert.deepEqual(await held, {}); }); @@ -1437,9 +1457,9 @@ describe('LinkLife', () => { test('gives up at once on a signal that was already aborted', async () => { const link = new LinkLife({ log: silentLog, reconnects: true, timeout: 100 }); - link.drop(); + link.transition('lost'); - const held = await link.hold(AbortSignal.abort())(); + const held = await link.budget(AbortSignal.abort())(); assert.match(held.err?.message ?? '', /Aborted while waiting for a link/); }); @@ -1448,47 +1468,58 @@ describe('LinkLife', () => { const link = new LinkLife({ log: silentLog, reconnects: true, timeout: 100 }); assert.equal(link.awaitsNextLink(), false, 'up'); - link.drop(); + link.transition('lost'); assert.equal(link.awaitsNextLink(), true, 'down, returning'); - link.attach(); + link.transition('attached'); assert.equal(link.awaitsNextLink(), true, 'attached, not yet bound'); - link.open(); + link.transition('bound'); assert.equal(link.awaitsNextLink(), false, 'reopened'); - link.drop(); - link.stop(); + link.transition('lost'); + link.transition('stopping'); assert.equal(link.awaitsNextLink(), false, 'down, stopped'); assert.match(link.refusal()?.message ?? '', /closed/, 'stopped while down'); - link.end(); + link.transition('closed'); assert.equal(link.awaitsNextLink(), false, 'ended'); }); - test('drops an attached link once, counts each drop, and names the event it warrants', () => { + test('drops an attached link once, counts each drop, and names what the session does next', () => { const link = new LinkLife({ log: silentLog, reconnects: true, timeout: 100 }); const generation = link.generation(); - assert.equal(link.drop(), 'disconnected'); - assert.equal(link.drop(), undefined, 'already down'); + assert.deepEqual(link.transition('lost'), ['dropLink', 'emitDisconnected', 'scheduleReconnect']); + assert.deepEqual(link.transition('lost'), [], 'already down'); assert.equal(link.generation(), generation + 1); - link.attach(); - link.stop(); - assert.equal(link.drop(), 'close', 'a new link drops again, with none to follow it'); + assert.deepEqual(link.transition('attached'), []); + assert.deepEqual(link.transition('stopping'), ['stopReconnect']); + assert.deepEqual(link.transition('lost'), ['dropLink', 'emitClose'], 'a new link drops again, with none to follow it'); assert.equal(link.generation(), generation + 2); - assert.equal(new LinkLife({ log: silentLog, reconnects: false, timeout: 100 }).drop(), 'close'); + assert.equal(link.phase, 'closed'); + assert.deepEqual(new LinkLife({ log: silentLog, reconnects: false, timeout: 100 }).transition('lost'), ['dropLink', 'emitClose']); + }); + + test('ends a link bound after the shutdown began, rather than bringing it up', () => { + const link = new LinkLife({ log: silentLog, reconnects: true, timeout: 100 }); + + link.transition('lost'); + link.transition('attached'); + link.transition('stopping'); + + assert.deepEqual(link.transition('bound'), ['dropLink', 'emitClose']); + assert.equal(link.carries(), false); }); test('releases a held request with the reason once the link ends', async () => { const link = new LinkLife({ log: silentLog, reconnects: true, timeout: 0 }); - link.drop(); + link.transition('lost'); - const held = link.hold(undefined)(); - - link.end(); + const held = link.budget(undefined)(); + assert.deepEqual(link.transition('closed'), ['emitClose']); assert.match((await held).err?.message ?? '', /Session is closed/); - link.attach(); - assert.equal(link.isAttached(), false, 'ended is final'); - assert.equal(link.end(), false); + link.transition('attached'); + assert.equal(link.attached(), false, 'ended is final'); + assert.deepEqual(link.transition('closed'), []); }); }); @@ -1542,12 +1573,12 @@ describe('held message bounds', () => { return [submitPdu(seqNr)]; } - function offer(held: HeldMessages, seqNr: number): MessageHold { - const hold = held.offer(message(seqNr)); + function offer(held: HeldMessages, seqNr: number): Sms { + const sms = held.offer(message(seqNr)); - assert.ok(hold); + assert.ok(sms); - return hold; + return sms; } /** Offers to a session with a listener, so an offer is held rather than released as untaken. */ @@ -1559,17 +1590,12 @@ describe('held message bounds', () => { closeAfter(t, session); session.on('sms', () => undefined); + session.sendReturn = () => Promise.resolve({}); - return new HeldMessages({ - ...options, - link: new LinkLife({ log: silentLog, reconnects: false, timeout: 100 }), - log: silentLog, - sendPastDrain: () => Promise.resolve({ err: new Error('never sent') }), - session, - }); + return heldMessagesOn(session, options); } - test('is full at its count, and a re-used sequence number replaces rather than adding', t => { + test('is full at its count, and a re-used sequence number replaces rather than adding', async t => { const held = heldOn(t, { max: 2, maxOctets: 1_000_000, timeout: 10_000 }); const first = offer(held, 1); const replaced = offer(held, 2); @@ -1579,14 +1605,16 @@ describe('held message bounds', () => { assert.equal(held.size, 2); assert.equal(held.octetsHeld, 2 * 1026, 'the replaced message leaves its octets with it'); assert.equal(held.full(), true); - assert.equal(first.isHeld(), true); - assert.equal(replaced.isHeld(), false); + await replaced.sendResp(); + assert.equal(held.size, 2, 'answering the replaced message releases nothing'); + await first.sendResp(); + assert.equal(held.size, 1, 'answering the first releases it'); held.clear(); }); // submitPdu() holds 1026 octets by the maxOctets charge: its object, and the three text fields. - test('is full at its octet cap, until a message leaves by any way out', t => { + test('is full at its octet cap, until a message leaves by any way out', async t => { let now = 0; const held = heldOn(t, { max: 10, maxOctets: 2000, now: () => now, timeout: 10_000 }); const answered = offer(held, 1); @@ -1595,8 +1623,8 @@ describe('held message bounds', () => { offer(held, 2); assert.equal(held.full(), true); - answered.release(); - assert.equal(held.full(), false, 'after a release'); + await answered.sendResp(); + assert.equal(held.full(), false, 'after an answer'); offer(held, 3); now = 20_000;