diff --git a/AGENTS.md b/AGENTS.md index 882bb24..55d0934 100644 --- a/AGENTS.md +++ b/AGENTS.md @@ -38,28 +38,28 @@ These are not preferences. Breaking one is a defect. ``` src/ index.ts Public surface. Named exports only, no default export. + bind-direction.ts The three bind types, which end of the link is which, what a bind carries, and how one is recorded 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 + session.ts Session: its life (open, closing, closed), the current Link, dispatch and events sms.ts The live handle emitted as the 'sms' event (sendResp/sendDlr) 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(): a shutdown's wait for the messages, then the requests, on one budget 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: the messages one link handed to the application, one HeldMessage each, and its six exits + idle-waiters.ts IdleWaiters: waiting for a count to fall to zero 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-timers.ts LinkTimers: the enquire_link heartbeat and the idle timeout + link.ts Link: one socket — the PDUs read off it, its timers, the requests waiting on it, what arrived on it — closed once 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 wait for a link and the retry 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 - pdu-transport.ts PduTransport: the socket a session reads complete PDUs off pending-requests.ts PendingRequests: sequence numbers, correlation, timeout, abort reassembly.ts Reassembler: capped, expiring multipart groups reconnect-loop.ts ReconnectLoop: backoff, retry timer, stopped-ness @@ -67,7 +67,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 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 @@ -310,7 +310,7 @@ this is not a changelog. - Locality work comes before other work until a scoring run reads 7.0. - A listener that rejects is routed by Node's `captureRejections`, not by hand-dispatching. -- The four-line abort dance is copied across `LinkLife`, `IdleWaiters`, `PendingRequests` and +- The four-line abort dance is copied across `OutgoingRequests`, `IdleWaiters`, `PendingRequests` and `SendWindow` rather than extracted. - `SmppLog` is a five-method contract this library declares, not a dependency. - The TLS tests build their own self-signed certificate in DER diff --git a/DESIGN.md b/DESIGN.md new file mode 100644 index 0000000..d819b39 --- /dev/null +++ b/DESIGN.md @@ -0,0 +1,65 @@ +# Draft A: a `Link` per socket, a three-word session life, one drain + +The confusion came from *one* socket's state being spread over a mutable phase (`LinkLife`), a +generation counter, a separate stopped flag, and five files. The redesign gives every socket an +object of its own and lets identity do what the counter and the phase did. + +## Structure and ownership + +| Module | Owns | +| --- | --- | +| `link.ts` — `Link` | One socket: framer, enquire/idle timers, `pending` (responses owed on it), `held` (messages handed out on it), `reassembler` (its half-arrived groups). `canCarry()` = bound, not closed, socket alive. `close()` runs once and ends all of it. | +| `session.ts` — `Session` | `life: 'open' \| 'closing' \| 'closed'` and `link: Link` (the latest, closed or not). The four transitions sit together: `openLink()`, `comeBackUp()`, `linkLost()`, `end()`. Nothing else writes `life` or `link`; every life event is emitted from one of them. | +| `outgoing-requests.ts` | The send window, the wait for a carrying link, the retry. Reads the session through `LinkView` (`closing`, `current`, `nextExpected`) and keeps no copy. | +| `incoming-requests.ts` | A per-session router, `handle(link, pduObj)`: what a request touches is the link it arrived on. | +| `held-messages.ts` | `HeldMessages` (per link) and `HeldMessage`, with the six exits listed once on the class. Owns the at-bound log hysteresis. | +| `drain.ts` | `drain(waits, budget, signal)`: messages first, then requests, one deadline, the never-forever rule for the application half, the combined error. | +| `bind-direction.ts` | What `session-options.ts` held that was not an option: bind types, `LinkEnd`, `standsInFor`, `bindCarries`, `checkedBind`. | + +## How each flow reads now + +**Held message, six exits** (all in `held-messages.ts`): 1 `sendResp()` reached the wire → +`HeldMessage.answered()` → release a turn later. 2 last listener rejects → `Session`'s +`captureRejectionSymbol` → `this.link.held.rejected(sms)` → `listenerGaveUp()`. 3 no listener or one +threw → `offer()` sees `emit()` false → `release()`. 4 re-used sequence number → `offer()` replaces +the entry. 5 deadline → `sweep()`. 6 link gone → `Link.close()` → `held.clear()`, which also makes +`lostLink()` true for every `sendResp()` after it. The `WeakMap` stays (a rejecting listener hands the +`Sms` back as `unknown`), but the route is two hops, not four, and there is no generation counter. + +**Shutdown** (`close()`/`unbind()`): `drain()` sets `life = 'closing'`, stops the loop, and if the +link can carry, waits through `drain.ts`; `end()` then sets `'closed'`, closes the link, clears the +merges, releases the link waiters with "closed", emits `close`. Both are idempotent by state, so a +listener re-entering `close()` mid-teardown changes nothing. + +**Link loss** (`linkLost(link)`): ignored unless `link` is the current one and open. Emits the error +it came with, decides `disconnected` vs `close` from `life === 'open' && loop` *before* `link.close()` +(a listener may close the session inside it), then either emits `disconnected` and schedules the +loop, or calls `end()`. + +**Reconnect** (`comeBackUp()`): a new `Link` with `bound: false` becomes `this.link`; the bind goes +out through `send()` because a bind is let onto the current link. A failed bind is `linkLost(link)`. +A `close()` that landed meanwhile shows as `link.isClosed()`. Otherwise `markBound()`, +`outgoing.linkBound()` (releases held sends), `reconnected`. + +**A send** (`requestDuringDrain`): `carrier()` returns the current link once `canCarry()`, waits +while `nextExpected()`, else refuses as closed; a write that reached no socket loops for the next +link. A destroyed-but-not-yet-closed socket cannot spin: `canCarry()` is false for it and the wait +ends only on `linkBound()`/`over()`. + +## Deleted + +`link-life.ts` (phase + stopped + generation + waiters), `link-timers.ts` and `pdu-transport.ts` +(both folded into `Link`), `IncomingRequests.listenerRejected/drain/clear`, `OutgoingRequests. +linkLost/deliver/settleRefused/canCarry/drain`, `Session.stop/emitClose/teardown/onClose/attach/ +resetTimers/transportFor`, `leftOf` (private to `drain.ts`), the `refusing` flag (now +`HeldMessages.atBound`), and `MessageHold` (now `HeldMessage`). Three modules out, three in +(`link.ts`, `drain.ts`, `bind-direction.ts`); `src/` stays at 35 top-level files. + +## Tests + +`docker compose run --rm node npm test`: lint and typecheck clean, **515 tests, 515 pass, 0 fail** +(same count as before). Changed, all in `test/session-extras.test.ts`, none weakened: + +- Helpers `incomingOn()`/`heldOn()`: build a `Link` (new `linkOn()`) instead of a `LinkLife`; `handle(pdu)` and `close()` go through that link. Assertions untouched. +- "drops a message whose link went while onRequest was still running": `link.drop()` → `link.close()`; the second message goes through a fresh link, since a closed link never delivers (the old generation counter let the same object be re-used). +- The six `LinkLife` unit tests, whose subject no longer exists, became five `OutgoingRequests waiting for a link` tests and one `Link` test asserting the same behaviour on the new objects: budget already spent when a link is released, the un-`unref`'d timer, abort while waiting, waits only while a link is on its way and is told "closed" otherwise, the drain refusal only while a link could carry; and `Link.close()` once, settling what waited on it and making `canCarry()` final. diff --git a/docs/decisions.md b/docs/decisions.md index c245097..802e34a 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 `Link.close()` and the reconnect loop retries it on a fresh `Link` 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 a send held for a link 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 `Link.close()` 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 go with the link: 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 when a link closes 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 @@ -752,10 +752,10 @@ rule and an index of the titles below. - **One owner decides whether a link can carry a request, and a bind is what makes it one.** Maintainer's call, 2026-09-01, extended 2026-09-28; goal 1, since a send on a link not yet bound - comes back `ESME_RINVBNDSTS`. `LinkLife` is told what happened and never reads back into the - session; every other collaborator reads it and keeps no copy. Rejected: gating on the socket being - 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. + comes back `ESME_RINVBNDSTS`. `Session` owns its life and the current `Link`, and `Link.canCarry()` + is the one answer; the senders read it through `LinkView` and keep no copy. Rejected: gating on the + socket being 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. @@ -777,7 +777,7 @@ rule and an index of the titles below. handlers normalise through `errorFrom()` rather than inline — a route out of the handler would land on a bare `process.nextTick` with nothing to catch it. -- **The four-line abort dance is copied across `LinkLife`, `IdleWaiters`, `PendingRequests` and +- **The four-line abort dance is copied across `OutgoingRequests`, `IdleWaiters`, `PendingRequests` and `SendWindow` rather than extracted.** Architecture review, 2026-09-06: pre-check `aborted`, attach `{ once: true }`, detach on settle, leave the registry. What differs at each site is the registry and what settling means — a FIFO handing over a slot, a set released together, a map keyed by @@ -797,7 +797,7 @@ rule and an index of the titles below. - **`src/` stays flat until a module has to move for another reason.** Architecture review, 2026-09-06: the grouping the [file map](../AGENTS.md#architecture) already implies — `wire/` for `pdu*` and `defs`, - `link/` for `link-*`, `reconnect-*`, `pdu-transport` and `send-window`, `messages/` for `sms*`, + `link/` for `link`, `reconnect-loop`, `outgoing-requests` and `send-window`, `messages/` for `sms*`, `dlr*`, `message*`, `reassembly` and `udh` — rewrites every import for no change to `dist/index.js`, the one published entry. Valid while that map is what a reader navigates by. diff --git a/src/bind-direction.ts b/src/bind-direction.ts new file mode 100644 index 0000000..9f62ffc --- /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'; + +/** 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 }; + +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; +} + +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..9bd407a --- /dev/null +++ b/src/drain.ts @@ -0,0 +1,57 @@ +import type { SmppLog } from './log.ts'; +import type { VoidResult } from './result.ts'; +import { defaults } from './session-options.ts'; + +export type DrainBudget = { + responseTimeout: number; + /** 0 waits forever for the requests; the messages then fall back to `responseTimeout`. */ + shutdownTimeout: number; +}; + +/** Each resolves 0 once nothing is left, or with what still is when the timeout or the signal cuts it short. */ +export type DrainWaits = { + /** Messages handed to the application and not yet answered. */ + messages: (timeout: number, signal: AbortSignal | undefined) => Promise; + /** Requests on the wire or queued behind the send window. */ + requests: (timeout: number, signal: AbortSignal | undefined) => Promise; +}; + +/** 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's half may never be "forever": nothing else ends that wait. */ +function messagesBudget(budget: DrainBudget): number { + if (budget.shutdownTimeout > 0) return budget.shutdownTimeout; + + return budget.responseTimeout > 0 ? budget.responseTimeout : defaults.responseTimeout; +} + +/** Waits out what a shutdown owes the peer: the messages first, then the requests, on one budget. */ +export async function drain( + waits: DrainWaits, + budget: DrainBudget, + signal: AbortSignal | undefined, + log: SmppLog, +): Promise { + const deadline = budget.shutdownTimeout > 0 ? Date.now() + budget.shutdownTimeout : 0; + const problems: string[] = []; + + // Messages first: answering one can put a receipt on the wire; nothing on the wire produces a message. + const unanswered = await waits.messages(messagesBudget(budget), signal); + + if (unanswered > 0) { + log.warn('drain - shutting down with messages unanswered', { unanswered }); + problems.push(`Shut down with ${String(unanswered)} message(s) unanswered`); + } + + const unfinished = await waits.requests(leftOf(deadline), signal); + + if (unfinished > 0) { + log.warn('drain - shutting down with requests unfinished', { unfinished }); + problems.push(`Shut down with ${String(unfinished)} request(s) unfinished`); + } + + return problems.length === 0 ? {} : { err: new Error(problems.join('; ')) }; +} diff --git a/src/error-from.ts b/src/error-from.ts index 5d77236..f5b3ff9 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; } + +/** namedValue() with a string in quotes, so an empty one and a wrong one both show. */ +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..0d89d8b 100644 --- a/src/held-messages.ts +++ b/src/held-messages.ts @@ -1,4 +1,3 @@ -import type { LinkLife } from './link-life.ts'; import type { PduObject, PduObjectInput } from './pdu.ts'; import type { Result } from './result.ts'; import type { Session } from './session.ts'; @@ -10,12 +9,12 @@ import { createSms } from './sms.ts'; import { retainedOctets } from './retained-pdu.ts'; export type HeldMessagesOptions = { - link: LinkLife; log: SmppLog; max: number; maxOctets: number; /** Injected so expiry can be exercised without a wall clock. */ now?: (() => number) | undefined; + /** A send the shutdown drain lets through, for the receipt of a message it is waiting on. */ sendPastDrain: SmsHandlers['send']; session: Session; timeout: number; @@ -28,31 +27,30 @@ 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. + * One message offered to the application, and what its `Sms` answers through. It is held until + * the first of six exits, each a method here or on `HeldMessages`: + * 1. `answered()`: `sendResp()` put the response on the wire, a turn later. + * 2. `listenerGaveUp()` from the last listener that took it and rejected. + * 3. `release()` at once, from `offer()`: no listener took it, or one threw. + * 4. `offer()` of a later message on the same sequence number replaces it. + * 5. `sweep()`: it passed its deadline. + * 6. `clear()`: the link it arrived on went, so nothing correlates its answer now. */ -export class MessageHold implements SmsHandlers { - private readonly generation: number; - private readonly heldMessages: HeldMessages; +export class HeldMessage implements SmsHandlers { + private readonly held: 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; + constructor(held: HeldMessages, pduObjs: PduObject[], listeners: number) { + this.held = held; 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); + return this.held.holds(this.pduObjs); } /** A turn later, so a `sendDlr()` called straight after `sendResp()` still goes out past a drain. */ @@ -61,7 +59,7 @@ export class MessageHold implements SmsHandlers { } lostLink(): boolean { - return this.route.link.generation() !== this.generation; + return this.held.isGone(); } /** A rejection leaves the other listeners running, so only the last one to fail gives the message up. */ @@ -71,26 +69,30 @@ export class MessageHold implements SmsHandlers { 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); + this.held.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); + return this.isHeld() ? this.held.sendPastDrain(input) : this.held.session.send(input); } } -/** The messages handed to the application that it has not answered yet, held by their segments. */ +/** The messages one link handed to the application that it has not answered yet. */ export class HeldMessages { + readonly sendPastDrain: SmsHandlers['send']; + readonly session: Session; + private readonly held: ExpiringGroups; private readonly idleWaiters = new IdleWaiters(); private readonly log: SmppLog; + private readonly max: number; private readonly maxOctets: number; /** A rejecting listener hands the message back as an `unknown`, so its hold is found by identity. */ - private readonly offered = new WeakMap(); - private readonly route: HoldRoute; + private readonly offered = new WeakMap(); + private atBound = false; + private gone = false; constructor(options: HeldMessagesOptions) { this.held = new ExpiringGroups({ @@ -100,8 +102,10 @@ export class HeldMessages { timeout: options.timeout, }); this.log = options.log; + this.max = options.max; this.maxOctets = options.maxOctets; - this.route = { link: options.link, sendPastDrain: options.sendPastDrain, session: options.session }; + this.sendPastDrain = options.sendPastDrain; + this.session = options.session; } get octetsHeld(): number { @@ -112,48 +116,64 @@ export class HeldMessages { return this.held.size; } - /** Whether a message arriving now is past the bound, once the expired are swept. */ - full(): boolean { - this.sweep(); - - return this.held.full || this.held.weight >= this.maxOctets; + /** The link these messages arrived on is gone. */ + isGone(): boolean { + return this.gone; } - private hold(key: string, pduObjs: PduObject[], listeners: number): MessageHold { - const hold = new MessageHold(this, this.route, pduObjs, listeners); - + /** Whether a message arriving now is past the bound. Logs the crossing once, each way. */ + full(): boolean { this.sweep(); - if (this.held.get(key)) { - this.log.warn('heldMessages - replacing a message on a re-used sequence number', { seqNr: Number(key) }); + const full = this.held.full || this.held.weight >= this.maxOctets; + + if (full && !this.atBound) { + this.atBound = true; + this.log.warn('heldMessages - unanswered messages at their bound, refusing new ones until the application answers', { + messages: this.held.size, + octets: this.held.weight, + }); } - 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 (!full && this.atBound && this.held.size <= this.max / 2 && this.held.weight <= this.maxOctets / 2) { + this.atBound = false; + this.log.info('heldMessages - unanswered messages down to half their bound, accepting again', { messages: this.held.size }); + } - return hold; + return full; } - offer(pduObjs: PduObject[], answeredAs?: string): MessageHold | undefined { + /** Hands the message to the application. `answeredAs` is the id base its segments were already answered with. */ + offer(pduObjs: PduObject[], answeredAs?: string): HeldMessage | 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 hold = new HeldMessage(this, pduObjs, this.session.listenerCount('sms')); + const sms = createSms({ answeredAs, pduObjs, session: this.session }, hold); + this.sweep(); + + if (this.held.get(key)) { + this.log.warn('heldMessages - replacing a message on a re-used sequence number', { seqNr: Number(key) }); + } + + this.held.set(key, pduObjs); + this.held.weigh(key, pduObjs.reduce((sum, pduObj) => sum + retainedOctets(pduObj), 0)); this.offered.set(sms, hold); - if (!this.route.session.emit('sms', sms)) hold.release(); + // False: no listener, or one threw. That is not work a shutdown can wait for. + if (!this.session.emit('sms', sms)) hold.release(); return hold; } /** One listener gave up on a message; the last one to do so is what releases it. */ - listenerRejected(message: unknown): void { - if (typeof message !== 'object' || message === null) return; + rejected(sms: unknown): void { + if (typeof sms !== 'object' || sms === null) return; - this.offered.get(message)?.listenerGaveUp(); + this.offered.get(sms)?.listenerGaveUp(); } holds(pduObjs: PduObject[]): boolean { @@ -172,8 +192,10 @@ export class HeldMessages { this.settle(); } - /** Drops every message: their segments went with the link, so no answer of ours correlates now. */ + /** The link went: every message goes with it, since no answer of ours correlates now. */ clear(): void { + this.gone = true; + this.atBound = false; this.held.takeAll(); this.idleWaiters.settle(); } @@ -183,7 +205,7 @@ 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. */ + /** Drops every message past its deadline. Runs before each offer and on its own timer. */ sweep(): void { const expired = this.held.takeExpired(); diff --git a/src/idle-waiters.ts b/src/idle-waiters.ts index dd29a9e..1240687 100644 --- a/src/idle-waiters.ts +++ b/src/idle-waiters.ts @@ -1,8 +1,3 @@ -/** 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)[] = []; diff --git a/src/incoming-requests.ts b/src/incoming-requests.ts index 51aeec7..04cb6d6 100644 --- a/src/incoming-requests.ts +++ b/src/incoming-requests.ts @@ -1,19 +1,16 @@ 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 { LinkLife } from './link-life.ts'; +import type { Link } from './link.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 { concatOf } from './concat.ts'; +import { defaults } from './session-options.ts'; import { detach } from './retained-pdu.ts'; import { dlrFromPdu } from './dlr.ts'; import { respIdParams, segmentId } from './sms-id.ts'; @@ -43,15 +40,17 @@ const lostReasons: Record = { linkGone: 'the link they arrived on went', }; +/** A group given up on is traffic the peer will not send again, so it is reported as lost. */ +export function lostGroupError(lost: LostGroup): Error { + return new Error( + `Gave up ${String(lost.parts)} of ${String(lost.total)} segments of an incomplete concatenated message: ${lostReasons[lost.reason]}`, + ); +} + export type IncomingRequestsOptions = { dlrMerger: DlrMerger; - 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; @@ -60,51 +59,30 @@ export type IncomingRequestsOptions = { /** Everything the peer asks of a session: messages, receipts, links and the answers to them. */ export class IncomingRequests { private readonly dlrMerger: DlrMerger; - private readonly held: HeldMessages; - private readonly link: LinkLife; private readonly log: SmppLog; private readonly onRequest: OnRequest | undefined; - private readonly reassembler: Reassembler; 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.link = options.link; this.log = options.log; this.onRequest = options.onRequest; - this.reassembler = new Reassembler({ - log: options.log, - max: options.maxReassembly ?? defaults.maxReassembly, - maxOctets: options.maxOctets, - onLost: lost => { this.reportLost(lost); }, - timeout: options.reassemblyTimeout ?? defaults.reassemblyTimeout, - }); this.session = options.session; this.smsIdFormat = options.smsIdFormat ?? {}; this.systemId = options.systemId ?? defaults.systemId; } - async handle(pduObj: PduObject): Promise { - const generation = this.link.generation(); + /** `link` is the one the request arrived on: what it holds is answered there, or not at all. */ + async handle(link: Link, pduObj: PduObject): Promise { const { onRequest } = this; // Called unbound, so the application's hook never sees this class as its `this`. if (onRequest && await onRequest(this.session, pduObj)) return; // The link it arrived on went while the hook ran, so nothing we answer now correlates. - if (this.link.generation() !== generation) { + if (link.isClosed()) { this.log.info('session - dropping a request whose link went', { cmdName: pduObj.cmdName }); return; @@ -120,23 +98,23 @@ export class IncomingRequests { return; } - await this.route(pduObj); + await this.route(link, pduObj); } - private async route(pduObj: PduObject): Promise { + private async route(link: Link, pduObj: PduObject): Promise { switch (pduObj.cmdName) { case 'data_sm': case 'deliver_sm': // A data_sm at the SMSC end is a submission, and a submission is never a report. await (this.carriedAs(pduObj) === 'submit_sm' - ? this.onMessage(pduObj) - : this.onDelivery(pduObj)); + ? this.onMessage(link, pduObj) + : this.onDelivery(link, pduObj)); break; case 'enquire_link': await this.session.sendReturn(pduObj); break; case 'submit_sm': - await this.onMessage(pduObj); + await this.onMessage(link, pduObj); break; case 'unbind': await this.session.sendReturn(pduObj); @@ -148,28 +126,6 @@ export class IncomingRequests { } } - /** Drops the segments of every message that never became whole, and of every one still held. */ - 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 }); @@ -193,11 +149,11 @@ export class IncomingRequests { } /** SMPP carries a mobile-originated message and a delivery receipt on the same command. */ - private async onDelivery(pduObj: PduObject): Promise { + private async onDelivery(link: Link, pduObj: PduObject): Promise { const dlr = dlrFromPdu(pduObj, this.smsIdFormat); if (!dlr) { - await this.onMessage(pduObj); + await this.onMessage(link, pduObj); return; } @@ -211,54 +167,30 @@ export class IncomingRequests { await this.session.sendReturn(pduObj); } - 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, - }); - } - + /** + * 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(link: Link, pduObj: PduObject): Promise { + if (link.held.full()) { 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 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 }); + return; } - 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) { - this.held.offer([detach(pduObj)]); + link.held.offer([detach(pduObj)]); return; } - const collected = this.reassembler.collect(pduObj, concat); + const collected = link.reassembler.collect(pduObj, concat); if (!collected.kept) { await this.session.sendReturn( @@ -275,12 +207,6 @@ export class IncomingRequests { respIdParams(pduObj.cmdName, segmentId(collected.smsId, concat.part - 1, concat.total)), ); - if (collected.whole) this.held.offer(collected.whole, collected.smsId); - } - - private reportLost(lost: LostGroup): void { - this.session.emit('sessionError', new Error( - `Gave up ${String(lost.parts)} of ${String(lost.total)} segments of an incomplete concatenated message: ${lostReasons[lost.reason]}`, - )); + if (collected.whole) link.held.offer(collected.whole, collected.smsId); } } diff --git a/src/link-life.ts b/src/link-life.ts deleted file mode 100644 index f44f2ea..0000000 --- a/src/link-life.ts +++ /dev/null @@ -1,188 +0,0 @@ -import type { SmppLog } from './log.ts'; -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(). */ - 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'; - -type Waiter = (result: VoidResult) => void; - -function aborted(): Error { - return new Error('Aborted while waiting for a link'); -} - -function expired(): Error { - return new Error('The link did not come back in time'); -} - -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. */ -export class LinkLife { - private readonly log: SmppLog; - private readonly now: () => number; - private readonly reconnects: boolean; - private readonly timeout: number; - private readonly waiting = new Set(); - private drops = 0; - private phase: Phase = 'up'; - private stopped = false; - - constructor(options: LinkLifeOptions) { - this.log = options.log; - this.now = options.now ?? Date.now; - this.reconnects = options.reconnects; - this.timeout = options.timeout; - } - - /** A socket is on the link, bound or not. */ - isAttached(): boolean { - return this.phase === 'binding' || this.phase === 'up'; - } - - /** Whether a request can go out right now. */ - isUp(): boolean { - return this.phase === 'up'; - } - - private isOver(): boolean { - return this.phase === 'ended'; - } - - /** The session is shutting down: nothing new is taken, and no link follows this one. */ - isStopped(): boolean { - return this.stopped; - } - - /** Whether a link that drops now is followed by another. */ - retrying(): boolean { - return this.reconnects && !this.stopped; - } - - /** Not up and not over, with a link to come. */ - awaitsNextLink(): boolean { - return !this.isUp() && !this.isOver() && this.retrying(); - } - - /** Changes with every drop, so what was read off one link can tell that link is gone. */ - generation(): number { - return this.drops; - } - - /** 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(); - } - - /** One budget for a request, however many links it waits through. */ - hold(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'; - } - - /** The link is bound: everything held goes out on it. */ - open(): void { - this.phase = 'up'; - - if (this.waiting.size > 0) { - this.log.verbose('linkLife - sending what was held for a link', { held: this.waiting.size }); - } - - this.release({}); - } - - /** 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; - - this.phase = 'down'; - this.drops++; - - return this.retrying() ? 'disconnected' : 'close'; - } - - stop(): void { - this.stopped = true; - } - - /** The session is over: nothing held will ever go out. False means it already was. */ - end(): boolean { - if (this.isOver()) return false; - - this.phase = 'ended'; - this.stopped = true; - this.release({ err: over() }); - - return true; - } - - /** 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({}); - - const refused = this.refusal(); - - if (refused) return Promise.resolve({ err: refused }); - - if (signal?.aborted === true) return Promise.resolve({ err: aborted() }); - - const left = deadline === 0 ? 0 : deadline - this.now(); - - if (deadline !== 0 && left <= 0) return Promise.resolve({ err: expired() }); - - return this.waitForLink(left, signal); - } - - private waitForLink(left: number, signal: AbortSignal | undefined): Promise { - this.log.verbose('linkLife - holding a request until a link is back', { timeout: left }); - - return new Promise(resolve => { - let timer: NodeJS.Timeout | undefined = undefined; - const settle = (result: VoidResult): void => { - if (timer) clearTimeout(timer); - - signal?.removeEventListener('abort', onAbort); - this.waiting.delete(settle); - resolve(result); - }; - const giveUp = (): void => { - this.log.warn('linkLife - no link came back in time', { timeout: left }); - settle({ err: expired() }); - }; - - function onAbort(): void { - settle({ err: aborted() }); - } - - // Not unref()'d: a held request is awaited with no other handle, so the process would exit unsettled. - if (left > 0) timer = setTimeout(giveUp, left); - - signal?.addEventListener('abort', onAbort, { once: true }); - this.waiting.add(settle); - }); - } - - private release(result: VoidResult): void { - for (const settle of [...this.waiting]) { - settle(result); - } - } -} diff --git a/src/link-timers.ts b/src/link-timers.ts deleted file mode 100644 index cd6710a..0000000 --- a/src/link-timers.ts +++ /dev/null @@ -1,50 +0,0 @@ -import type { SmppLog } from './log.ts'; - -export type LinkTimersOptions = { - /** How long between enquire_link probes. Undefined or 0 never probes. */ - enquireLinkInterval?: number | undefined; - /** How long a silent peer is kept. Undefined or 0 keeps it forever. */ - idleTimeout?: number | undefined; - log: SmppLog; - onEnquireLink: () => void; - onIdle: () => void; -}; - -/** Keeps a quiet connection honest: probes the peer, and gives up on one that stays silent. */ -export class LinkTimers { - private readonly options: LinkTimersOptions; - private enquireLink: NodeJS.Timeout | undefined; - private idle: NodeJS.Timeout | undefined; - - constructor(options: LinkTimersOptions) { - this.options = options; - } - - /** Starts both timers over, which every sign of life from the peer should do. */ - reset(): void { - const { enquireLinkInterval, idleTimeout, log, onEnquireLink, onIdle } = this.options; - - this.clear(); - - if (enquireLinkInterval !== undefined && enquireLinkInterval > 0) { - this.enquireLink = setTimeout(onEnquireLink, enquireLinkInterval); - this.enquireLink.unref(); - } - - if (idleTimeout !== undefined && idleTimeout > 0) { - this.idle = setTimeout(() => { - log.info('linkTimers - closing an idle peer', { idleTimeout }); - onIdle(); - }, idleTimeout); - this.idle.unref(); - } - } - - clear(): void { - if (this.enquireLink) clearTimeout(this.enquireLink); - if (this.idle) clearTimeout(this.idle); - - this.enquireLink = undefined; - this.idle = undefined; - } -} diff --git a/src/link.ts b/src/link.ts new file mode 100644 index 0000000..efd05b1 --- /dev/null +++ b/src/link.ts @@ -0,0 +1,243 @@ +import type { HeldMessagesOptions } from './held-messages.ts'; +import type { PduObject, PduObjectInput } from './pdu.ts'; +import type { ReassemblerOptions } from './reassembly.ts'; +import type { Result, VoidResult } from './result.ts'; +import type { SendOptions } from './session-options.ts'; +import type { SmppLog } from './log.ts'; +import type { Socket } from 'node:net'; +import { HeldMessages } from './held-messages.ts'; +import { PduFramer } from './pdu-framer.ts'; +import { PduRefusedError } from './pdu-refusal.ts'; +import { PendingRequests } from './pending-requests.ts'; +import { Reassembler } from './reassembly.ts'; +import { UnansweredError } from './unanswered-error.ts'; +import { isResp, objToPdu, pduToObj } from './pdu.ts'; + +export type LinkEvents = { + /** Raw bytes, before framing. */ + data: (chunk: Buffer) => void; + enquireLink: () => void; + /** A complete PDU, before it is parsed. */ + framed: (pdu: Buffer) => void; + /** Nothing further will be read off the socket: it closed, errored, fell silent or lost sync. */ + lost: (link: Link, err?: Error) => void; + /** A framed PDU the codec could not read. The stream is still in sync, so the link is not lost. */ + refused: (refused: PduRefusedError) => void; + /** A request the peer sent. Responses never reach here: they settle the request waiting on this link. */ + request: (link: Link, pduObj: PduObject) => void; +}; + +export type LinkOptions = { + /** Whether requests may go out from the start. A reconnect's socket carries only its bind until `markBound()`. */ + bound: boolean; + /** Undefined or 0 never probes. */ + enquireLinkInterval?: number | undefined; + held: Omit; + /** Undefined or 0 keeps a silent peer forever. */ + idleTimeout?: number | undefined; + log: SmppLog; + on: LinkEvents; + reassembly: Omit; + responseTimeout: number; + sock: Socket; +}; + +/** `retryOnNextLink`: the write failed, so nothing reached the socket and another link may carry it. */ +export type Attempt = { result: Result<{ pduObj: PduObject }>; retryOnNextLink: boolean }; + +function abortedBeforeSend(): Error { + return new Error('Aborted before the request was sent'); +} + +/** + * One socket: the PDUs read off it, the timers that keep it honest, the requests waiting on it and + * the messages that arrived on it. A reconnect opens a new one; `close()` ends this one for good. + */ +export class Link { + readonly held: HeldMessages; + readonly pending: PendingRequests; + readonly reassembler: Reassembler; + readonly sock: Socket; + + private readonly framer = new PduFramer(); + private readonly log: SmppLog; + private readonly options: LinkOptions; + private bound: boolean; + private closed = false; + private enquireLinkTimer: NodeJS.Timeout | undefined; + private idleTimer: NodeJS.Timeout | undefined; + + constructor(options: LinkOptions) { + this.bound = options.bound; + this.held = new HeldMessages({ ...options.held, log: options.log }); + this.log = options.log; + this.options = options; + this.pending = new PendingRequests(options.log); + this.reassembler = new Reassembler({ ...options.reassembly, log: options.log }); + this.sock = options.sock; + + options.sock.on('data', chunk => { this.read(chunk); }); + options.sock.on('close', () => { options.on.lost(this); }); + options.sock.on('error', err => { + this.log.warn('link - socket error', { message: err.message }); + options.on.lost(this, err); + }); + + if (this.bound) this.resetTimers(); + } + + /** Whether a request can go out on it right now. */ + canCarry(): boolean { + return this.bound && !this.closed && !this.sock.destroyed; + } + + isClosed(): boolean { + return this.closed; + } + + /** The bind was answered, so requests may go out on it. */ + markBound(): void { + this.bound = true; + this.resetTimers(); + } + + write(pdu: Buffer): VoidResult { + if (this.sock.destroyed) return { err: new Error('Socket is closed') }; + + this.sock.write(pdu); + + return {}; + } + + /** Puts a request on the wire and resolves with the peer's response. */ + async send(input: PduObjectInput, options: SendOptions): Promise { + // pending.wait() alone settles the caller while the request still goes out to the peer. + if (options.signal?.aborted === true) { + return { result: { err: abortedBeforeSend() }, retryOnNextLink: false }; + } + + const seqNr = this.pending.nextSeqNr(); + const built = objToPdu({ ...input, seqNr }); + + if (built.err) return { result: { err: built.err }, retryOnNextLink: false }; + + const response = this.pending.wait(seqNr, { + signal: options.signal, + timeout: this.options.responseTimeout, + }); + const written = this.write(built.buffer); + + if (written.err) { + this.pending.settle(seqNr, { err: written.err }); + + return { result: { err: written.err }, retryOnNextLink: true }; + } + + const answered = await response; + + // It went out, so a failure now means the peer may have taken it and the answer was the loss. + return { result: answered.err ? { err: new UnansweredError(answered.err) } : answered, retryOnNextLink: false }; + } + + /** + * Once: fails every request waiting on it, drops every message that arrived on it, and destroys + * the socket. Nothing is read off it after this, so the socket's own close event changes nothing. + */ + close(): void { + if (this.closed) return; + + this.closed = true; + this.pending.settleAll(new Error('Session closed before a response arrived')); + this.clearTimers(); + this.held.clear(); + this.reassembler.clear(); + this.sock.destroy(); + } + + private read(chunk: Buffer): void { + this.options.on.data(chunk); + this.resetTimers(); + this.framer.push(chunk); + + const framed = this.framer.next(); + + if (framed.err) { + this.log.warn('link - unusable stream', { message: framed.err.message }); + this.options.on.lost(this, framed.err); + + return; + } + + for (const pdu of framed.pdus) { + this.options.on.framed(pdu); + + const parsed = pduToObj(pdu); + + if (parsed.err) { + this.refuse(parsed.err); + + continue; + } + + this.dispatch(parsed.pduObj); + } + } + + private dispatch(pduObj: PduObject): void { + if (!isResp(pduObj)) { + this.options.on.request(this, pduObj); + + return; + } + + if (!this.pending.deliver(pduObj)) { + this.log.debug('link - response with no matching request', { seqNr: pduObj.seqNr }); + } + } + + private refuse(err: Error): void { + if (!(err instanceof PduRefusedError)) { + // The framer applies framingRefusal() first, so only a caller that skips it lands here. + this.log.warn('link - could not parse an incoming PDU', { message: err.message }); + this.options.on.lost(this, err); + + return; + } + + this.log.warn('link - refusing a PDU it could not read', { message: err.message, reason: err.reason }); + this.options.on.refused(err); + + // A response carries a sequence number of ours, so it settles the request instead of being answered. + if (isResp(err.header)) this.pending.settle(err.header.seqNr, { err }); + } + + /** Starts both timers over, which every sign of life from the peer does. */ + private resetTimers(): void { + const { enquireLinkInterval, idleTimeout, on } = this.options; + + this.clearTimers(); + + if (this.closed) return; + + if (enquireLinkInterval !== undefined && enquireLinkInterval > 0) { + this.enquireLinkTimer = setTimeout(on.enquireLink, enquireLinkInterval); + this.enquireLinkTimer.unref(); + } + + if (idleTimeout !== undefined && idleTimeout > 0) { + this.idleTimer = setTimeout(() => { + this.log.info('link - closing an idle peer', { idleTimeout }); + on.lost(this); + }, idleTimeout); + this.idleTimer.unref(); + } + } + + private clearTimers(): void { + if (this.enquireLinkTimer) clearTimeout(this.enquireLinkTimer); + if (this.idleTimer) clearTimeout(this.idleTimer); + + this.enquireLinkTimer = undefined; + this.idleTimer = undefined; + } +} diff --git a/src/outgoing-requests.ts b/src/outgoing-requests.ts index a0adf24..664b53b 100644 --- a/src/outgoing-requests.ts +++ b/src/outgoing-requests.ts @@ -1,30 +1,48 @@ -import type { LinkLife } from './link-life.ts'; +import type { Attempt, Link } from './link.ts'; import type { PduObject, PduObjectInput } from './pdu.ts'; -import type { PduTransport } from './pdu-transport.ts'; import type { Result, VoidResult } from './result.ts'; import type { SendOptions } from './session-options.ts'; 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 { objToPdu } from './pdu.ts'; +import { bindCommands } from './bind-direction.ts'; + +/** What the senders read of the session's life. The session owns it and answers each from its state. */ +export type LinkView = { + /** A shutdown has begun: no new request is taken, and no link follows the current one. */ + closing: () => boolean; + /** The latest link, closed or not. */ + current: () => Link; + /** Whether a link that is gone is followed by another. */ + nextExpected: () => boolean; +}; export type OutgoingRequestsOptions = { - link: LinkLife; + links: LinkView; log: SmppLog; maxOutstanding: number; + now?: (() => number) | undefined; + /** How long a request may wait for a link. 0 waits for as long as one may still arrive. */ responseTimeout: number; - transport: PduTransport; }; -/** `retryOnNextLink`: the write failed, so nothing reached the socket and another link may carry it. */ -type Attempt = { result: Result<{ pduObj: PduObject }>; retryOnNextLink: boolean }; +type Waiter = (result: VoidResult) => void; function abortedBeforeSend(): Error { return new Error('Aborted before the request was sent'); } +function abortedWaiting(): Error { + return new Error('Aborted while waiting for a link'); +} + +function expired(): Error { + return new Error('The link did not come back in time'); +} + +function over(): Error { + return new Error('Session is closed'); +} + /** A response carries the request's sequence number, which only sendReturn() has. */ function misuse(input: PduObjectInput): Error | undefined { return input.cmdName.endsWith('_resp') @@ -32,43 +50,28 @@ function misuse(input: PduObjectInput): Error | undefined { : undefined; } +/** Why a request cannot go out at all. Before the link and the window, or an aborted call waits for what it will never use. */ +function refusal(input: PduObjectInput, options: SendOptions): Error | undefined { + return misuse(input) ?? (options.signal?.aborted === true ? abortedBeforeSend() : 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 links: LinkView; private readonly log: SmppLog; - private readonly pending: PendingRequests; + private readonly now: () => number; private readonly responseTimeout: number; - private readonly transport: PduTransport; + private readonly waiting = new Set(); private readonly window: SendWindow; constructor(options: OutgoingRequestsOptions) { - this.link = options.link; + this.links = options.links; this.log = options.log; - this.pending = new PendingRequests(options.log); + this.now = options.now ?? Date.now; this.responseTimeout = options.responseTimeout; - this.transport = options.transport; this.window = new SendWindow({ limit: options.maxOutstanding, log: options.log }); } - canCarry(): boolean { - return this.link.isUp() && !this.transport.sock.destroyed; - } - - /** The link is gone, and every answer still owed on it with it. */ - linkLost(): void { - this.pending.settleAll(new Error('Session closed before a response arrived')); - } - - /** Hands a response to the request waiting for it. False means nothing was. */ - deliver(pduObj: PduObject): boolean { - return this.pending.deliver(pduObj); - } - - /** A response the codec refused settles its request instead of leaving it to time out. */ - settleRefused(seqNr: number, err: Error): void { - this.pending.settle(seqNr, { err }); - } - request(input: PduObjectInput, options: SendOptions): Promise> { // Ahead of the drain, so a misuse is named as one rather than blamed on the shutdown. const wrong = misuse(input); @@ -76,44 +79,54 @@ 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.links.closing() && this.links.current().canCarry()) { return Promise.resolve({ err: new Error('Session is shutting down') }); } - return this.requestPastDrain(input, options); + return this.requestDuringDrain(input, options); } /** request() without the drain's refusal, which a receipt for a held message has to take. */ - async requestPastDrain( + async requestDuringDrain( input: PduObjectInput, options: SendOptions, ): Promise> { - const refused = this.refuse(input, options); + const refused = refusal(input, options); if (refused) return { 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(); + if (bindCommands.includes(input.cmdName)) return this.bindOnCurrentLink(input, options); - return shut ? { err: shut } : this.requestOnCurrentLink(input, options); - } - - const waitForLink = this.link.hold(options.signal); + const deadline = this.responseTimeout > 0 ? this.now() + this.responseTimeout : 0; for (;;) { - const held = await waitForLink(); + const carrier = await this.carrier(deadline, options.signal); - if (held.err) return { err: held.err }; + if (carrier.err) return { err: carrier.err }; - const slot = await this.window.acquire(options.signal); + const attempt = await this.attemptOn(carrier.link, input, options); - if (slot.err) return { err: slot.err }; + // Nothing reached the socket, so the next link carries it; with none to come, this is the answer. + if (!attempt.retryOnNextLink || !this.links.nextExpected()) return attempt.result; + } + } - const attempt = await this.attempt(input, options).finally(() => { this.window.release(); }); + private async bindOnCurrentLink(input: PduObjectInput, options: SendOptions): Promise> { + const link = this.links.current(); - if (!this.retriesOnNextLink(attempt)) return attempt.result; - } + if (!link.canCarry() && !this.links.nextExpected()) return { err: over() }; + + return (await link.send(input, options)).result; + } + + /** One try on one link, under a send-window slot. */ + private async attemptOn(link: Link, input: PduObjectInput, options: SendOptions): Promise { + const slot = await this.window.acquire(options.signal); + + if (slot.err) return { result: { err: slot.err }, retryOnNextLink: false }; + + return link.send(input, options).finally(() => { this.window.release(); }); } /** Straight onto the current link, for what has to go out either way. */ @@ -121,58 +134,81 @@ export class OutgoingRequests { input: PduObjectInput, options: SendOptions = {}, ): Promise> { - return (await this.attempt(input, options)).result; + return (await this.links.current().send(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); - - if (unfinished === 0) return {}; + /** Resolves 0 once nothing is on the wire or queued for it, or with what still is. */ + idle(timeout: number, signal: AbortSignal | undefined): Promise { + return this.window.idle(timeout, signal); + } - this.log.warn('outgoingRequests - shutting down with requests unfinished', { timeout, unfinished }); + /** A link is bound: everything held for one goes out on it. */ + linkBound(): void { + if (this.waiting.size > 0) { + this.log.verbose('outgoingRequests - sending what was held for a link', { held: this.waiting.size }); + } - return { err: new Error(`Shut down with ${String(unfinished)} request(s) unfinished`) }; + this.release({}); } - /** 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(); + /** No link will follow, so nothing held for one will ever go out. */ + over(): void { + this.release({ err: over() }); } - /** Why a request cannot go out at all, as opposed to not yet. */ - private refuse(input: PduObjectInput, options: SendOptions): Error | undefined { - // Before the link and the window, or an aborted call waits for what it will never use. - return misuse(input) ?? (options.signal?.aborted === true ? abortedBeforeSend() : undefined); - } + /** Resolves with a link that can carry the request, or with the reason none ever will. */ + private async carrier(deadline: number, signal: AbortSignal | undefined): Promise> { + for (;;) { + const link = this.links.current(); - private async attempt(input: PduObjectInput, options: SendOptions): Promise { - // pending.wait() alone settles the caller while the request still goes out to the peer. - if (options.signal?.aborted === true) { - return { result: { err: abortedBeforeSend() }, retryOnNextLink: false }; - } + if (link.canCarry()) return { link }; - const seqNr = this.pending.nextSeqNr(); - const built = objToPdu({ ...input, seqNr }); + if (!this.links.nextExpected()) return { err: over() }; - if (built.err) return { result: { err: built.err }, retryOnNextLink: false }; + if (signal?.aborted === true) return { err: abortedWaiting() }; - const response = this.pending.wait(seqNr, { - signal: options.signal, - timeout: this.responseTimeout, - }); - const written = this.transport.write(built.buffer); + const left = deadline === 0 ? 0 : deadline - this.now(); + + if (deadline !== 0 && left <= 0) return { err: expired() }; - if (written.err) { - this.pending.settle(seqNr, { err: written.err }); + const waited = await this.waitForLink(left, signal); - return { result: { err: written.err }, retryOnNextLink: true }; + if (waited.err) return { err: waited.err }; } + } + + private waitForLink(left: number, signal: AbortSignal | undefined): Promise { + this.log.verbose('outgoingRequests - holding a request until a link is back', { timeout: left }); + + return new Promise(resolve => { + let timer: NodeJS.Timeout | undefined = undefined; + const settle = (result: VoidResult): void => { + if (timer) clearTimeout(timer); + + signal?.removeEventListener('abort', onAbort); + this.waiting.delete(settle); + resolve(result); + }; + const giveUp = (): void => { + this.log.warn('outgoingRequests - no link came back in time', { timeout: left }); + settle({ err: expired() }); + }; - const answered = await response; + function onAbort(): void { + settle({ err: abortedWaiting() }); + } - // It went out, so a failure now means the peer may have taken it and the answer was the loss. - return { result: answered.err ? { err: new UnansweredError(answered.err) } : answered, retryOnNextLink: false }; + // Not unref()'d: a held request is awaited with no other handle, so the process would exit unsettled. + if (left > 0) timer = setTimeout(giveUp, left); + + signal?.addEventListener('abort', onAbort, { once: true }); + this.waiting.add(settle); + }); + } + + private release(result: VoidResult): void { + for (const settle of [...this.waiting]) { + settle(result); + } } } diff --git a/src/pdu-transport.ts b/src/pdu-transport.ts deleted file mode 100644 index 966cdee..0000000 --- a/src/pdu-transport.ts +++ /dev/null @@ -1,108 +0,0 @@ -import type { PduObject } from './pdu.ts'; -import type { SmppLog } from './log.ts'; -import type { Socket } from 'node:net'; -import type { VoidResult } from './result.ts'; -import { PduFramer } from './pdu-framer.ts'; -import { PduRefusedError } from './pdu-refusal.ts'; -import { pduToObj } from './pdu.ts'; - -export type PduTransportOptions = { - log: SmppLog; - onClose: () => void; - /** Raw bytes, before framing. */ - onData: (chunk: Buffer) => void; - onError: (err: Error) => void; - /** A complete PDU, before it is parsed. */ - onFramed: (pdu: Buffer) => void; - onPdu: (pduObj: PduObject) => void; - /** A framed PDU the codec could not read. The stream is still in sync, so the link is not lost. */ - onRefused: (refused: PduRefusedError) => void; - /** Nothing further can be read off this stream, whatever the socket does next. */ - onUnreadable: (err: Error) => void; -}; - -/** A socket read as a stream of complete PDUs. A reconnect attaches a new socket in its place. */ -export class PduTransport { - private readonly options: PduTransportOptions; - private framer = new PduFramer(); - private socket: Socket; - - constructor(options: PduTransportOptions, sock: Socket) { - this.options = options; - this.socket = sock; - this.wire(sock); - } - - get sock(): Socket { - return this.socket; - } - - /** Takes over a freshly opened socket. Half a PDU left on the old one must not prefix this one. */ - attach(sock: Socket): void { - // The socket being replaced is already dead, and its three handlers still point here. - this.socket.removeAllListeners(); - this.socket = sock; - this.framer = new PduFramer(); - this.wire(sock); - } - - private wire(sock: Socket): void { - sock.on('data', chunk => { this.read(chunk); }); - sock.on('close', () => { this.options.onClose(); }); - sock.on('error', err => { - this.options.log.warn('transport - socket error', { message: err.message }); - this.options.onError(err); - this.options.onClose(); - }); - } - - write(pdu: Buffer): VoidResult { - if (this.socket.destroyed) return { err: new Error('Socket is closed') }; - - this.socket.write(pdu); - - return {}; - } - - private read(chunk: Buffer): void { - this.options.onData(chunk); - this.framer.push(chunk); - - const framed = this.framer.next(); - - if (framed.err) { - this.options.log.warn('transport - unusable stream', { message: framed.err.message }); - this.options.onUnreadable(framed.err); - - return; - } - - for (const pdu of framed.pdus) { - this.options.onFramed(pdu); - - const parsed = pduToObj(pdu); - - if (parsed.err instanceof PduRefusedError) { - this.options.log.warn('transport - refusing a PDU it could not read', { - message: parsed.err.message, - reason: parsed.err.reason, - }); - this.options.onRefused(parsed.err); - - continue; - } - - // The framer applies framingRefusal() first, so only a caller that skips it lands here. - if (parsed.err) { - this.options.log.warn('transport - could not parse an incoming PDU', { - message: parsed.err.message, - }); - this.options.onUnreadable(parsed.err); - - return; - } - - this.options.onPdu(parsed.pduObj); - } - } -} 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..86e303b 100644 --- a/src/session.ts +++ b/src/session.ts @@ -1,25 +1,25 @@ +import type { BindType, LinkEnd, SessionBind } from './bind-direction.ts'; import type { ErrorName } from './defs/errors.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 { IncomingRequests } from './incoming-requests.ts'; -import { LinkLife } from './link-life.ts'; -import { LinkTimers } from './link-timers.ts'; +import { IncomingRequests, lostGroupError } from './incoming-requests.ts'; +import { Link } from './link.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'; @@ -42,6 +42,9 @@ export { bindCommands, defaultSystemId }; /** A listener may return a promise: an `async` one that rejects is routed like one that throws. */ type SessionListener = (...args: SessionEvents[K]) => unknown; +/** `closing`: a shutdown began; new sends are refused and no link follows the current one. */ +type Life = 'closed' | 'closing' | 'open'; + export class Session extends EventEmitter { declare addListener: (event: K, listener: SessionListener) => this; declare off: (event: K, listener: SessionListener) => this; @@ -58,16 +61,16 @@ export class Session extends EventEmitter { userData: unknown = undefined; private bind: SessionBind | undefined = undefined; + private life: Life = 'open'; + /** The latest socket. Between links it is the closed one, so `sock` still answers. */ + private link: Link; private readonly concatReference = new ConcatReference(); private readonly dlrMerger: DlrMerger; private readonly incoming: IncomingRequests; - private readonly link: LinkLife; private readonly options: SessionOptions; private readonly outgoing: OutgoingRequests; private readonly reconnectLoop: ReconnectLoop | undefined; - private readonly timers: LinkTimers; - private readonly transport: PduTransport; /** A listener that throws is the application's bug; it must not become ours. Hard rule 1. */ override emit( @@ -98,7 +101,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.link.held.rejected(rest[0]); if (event !== 'sessionError') this.emit('sessionError', error); } @@ -110,46 +113,31 @@ export class Session extends EventEmitter { this.options = options; this.dlrMerger = new DlrMerger({ log: this.log, max: defaults.maxDlrMerges, timeout: defaults.dlrMergeTimeout }); this.reconnectLoop = this.loopFor(options.reconnect); - - const responseTimeout = options.responseTimeout ?? defaults.responseTimeout; - - this.link = new LinkLife({ log: this.log, reconnects: this.reconnectLoop !== undefined, timeout: responseTimeout }); - this.timers = new LinkTimers({ - enquireLinkInterval: options.enquireLinkInterval, - 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(); }, - }); - this.transport = this.transportFor(options.sock); this.outgoing = new OutgoingRequests({ - link: this.link, + links: { + closing: () => this.life !== 'open', + current: () => this.link, + nextExpected: () => this.nextLinkExpected(), + }, log: this.log, maxOutstanding: options.maxOutstanding ?? defaults.maxOutstanding, - responseTimeout, - transport: this.transport, + responseTimeout: this.responseTimeout(), }); this.incoming = new IncomingRequests({ dlrMerger: this.dlrMerger, - 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(); + // The first socket carries from the start: its bind goes out through send(). + this.link = this.openLink(options.sock, true); } /** Replaced on reconnect, so hold the session rather than this. */ get sock(): Socket { - return this.transport.sock; + return this.link.sock; } /** The role the ESME bound with, whichever end of the link this is. Undefined before any bind. */ @@ -196,22 +184,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,11 +207,12 @@ export class Session extends EventEmitter { */ async unbind(): Promise { const drained = await this.drain(undefined); - const wasOpen = this.link.isAttached(); + const link = this.link; + const wasOpen = !link.isClosed(); const sent = wasOpen ? await this.outgoing.requestOnCurrentLink({ cmdName: 'unbind' }) : { err: new Error('Session is closed') }; - const closedOnUnbind = wasOpen && !this.link.isAttached(); + const closedOnUnbind = wasOpen && link.isClosed(); this.end(); @@ -258,155 +231,165 @@ export class Session extends EventEmitter { return drained; } - private transportFor(sock: Socket): PduTransport { - return new PduTransport({ - log: this.log, - onClose: () => { this.onClose(); }, - onData: chunk => { this.onData(chunk); }, - onError: err => { this.emit('sessionError', err); }, - onFramed: pdu => { this.emit('incomingPdu', pdu); }, - onPdu: pduObj => { this.dispatch(pduObj); }, - onRefused: refused => { this.refuse(refused); }, - onUnreadable: err => { - this.emit('sessionError', err); - this.teardown(); - }, - }, sock); + private answer(built: Result<{ buffer: Buffer }>, cmdName: string, seqNr: number): VoidResult { + const sent = built.err ? { err: built.err } : this.link.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.isClosed()) { + this.log.warn('session - could not answer a request', { + cmdName, + message: sent.err.message, + seqNr, + }); + this.emit('sessionError', sent.err); + } + + return sent; } - private loopFor(reconnect: ReconnectOptions | undefined): ReconnectLoop | undefined { - if (!reconnect) return undefined; + /** Stops new sends and the reconnect loop, then waits out what the peer is owed on the link. */ + private async drain(signal: AbortSignal | undefined): Promise { + if (this.life === 'open') this.life = 'closing'; - return new ReconnectLoop({ - connect: reconnect.connect, + this.reconnectLoop?.stop(); + + const link = this.link; + + // No bound link, so nothing is on the wire to wait out. + if (!link.canCarry()) return {}; + + const drained = await drain({ + messages: (timeout, cut) => link.held.idle(timeout, cut), + requests: (timeout, cut) => this.outgoing.idle(timeout, cut), + }, { + responseTimeout: this.responseTimeout(), + shutdownTimeout: this.options.shutdownTimeout ?? defaults.shutdownTimeout, + }, signal, this.log); + + // The link went before the drain finished, so an empty window says nothing about the peer. + if (!link.canCarry()) return { err: new Error('The session closed before the drain finished') }; + + return drained; + } + + // The session's life. A link opens, binds and is lost; the session ends once. Every event about + // the life is emitted from one of these four, and nothing else changes `life` or `link`. + + private openLink(sock: Socket, bound: boolean): Link { + return new Link({ + bound, + enquireLinkInterval: this.options.enquireLinkInterval, + held: { + max: defaults.maxHeldMessages, + maxOctets: defaults.maxHeldOctets, + sendPastDrain: input => this.outgoing.requestDuringDrain(input, {}), + session: this, + timeout: defaults.heldMessageTimeout, + }, + idleTimeout: this.options.idleTimeout, log: this.log, - maxDelay: reconnect.maxDelay, - minDelay: reconnect.minDelay, - onConnected: sock => this.comeBackUp(sock, reconnect.onConnected), + on: { + data: chunk => { this.emit('data', chunk); }, + enquireLink: () => { void this.send({ cmdName: 'enquire_link' }); }, + framed: pdu => { this.emit('incomingPdu', pdu); }, + lost: (link, err) => { this.linkLost(link, err); }, + refused: refused => { this.refuse(refused); }, + request: (link, pduObj) => { this.dispatch(link, pduObj); }, + }, + reassembly: { + max: this.options.maxReassembly ?? defaults.maxReassembly, + maxOctets: this.options.maxOctets, + onLost: lost => { this.emit('sessionError', lostGroupError(lost)); }, + timeout: this.options.reassemblyTimeout ?? defaults.reassemblyTimeout, + }, + responseTimeout: this.responseTimeout(), + sock, }); } - private async comeBackUp( - sock: Socket, - bind: (session: Session) => Promise, - ): Promise { - this.attach(sock); + /** Brings the session up on the loop's fresh socket. An err means the loop tries again. */ + private async comeBackUp(sock: Socket, bind: (session: Session) => Promise): Promise { + const link = this.openLink(sock, false); + + this.link = link; const bound = await bind(this); if (bound.err) { - this.teardown(); + this.linkLost(link); 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') }; - } + if (link.isClosed()) return { err: new Error('Session closed while it was coming back up') }; - this.resetTimers(); - this.link.open(); + link.markBound(); + this.outgoing.linkBound(); 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(); + /** The one exit for a socket: it closed, errored, fell silent, lost sync or failed to rebind. */ + private linkLost(link: Link, err?: Error): void { + if (link !== this.link || link.isClosed()) return; - // No bound link, so nothing is on the wire to wait out. - if (!this.outgoing.canCarry()) return {}; + if (err) this.emit('sessionError', err); - 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') }; - } + // Read before close(): a listener it reaches may close() the session, and the drop still reports as disconnected. + const disconnected = this.nextLinkExpected(); - if (!messages.err) return requests; + link.close(); - if (!requests.err) return messages; + if (!disconnected) { + this.end(); - 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; + return; + } - return responseTimeout > 0 ? responseTimeout : defaults.responseTimeout; + this.emit('disconnected'); + this.reconnectLoop?.schedule(); } /** The session is over now, drained or not. Nothing brings it back. */ private end(): void { - this.stop(); - this.teardown(); - this.dlrMerger.clear(); - this.emitClose(); - } + if (this.life === 'closed') return; - /** No new sends, and no link after this one. */ - private stop(): void { - this.link.stop(); + this.life = 'closed'; this.reconnectLoop?.stop(); - } - - private emitClose(): void { - if (!this.link.end()) return; - - this.outgoing.linkLost(); + this.link.close(); + this.dlrMerger.clear(); + this.outgoing.over(); this.emit('close'); } - private teardown(): void { - const lost = this.link.drop(); - - if (!lost) return; + /** Whether a link that is gone is followed by another. */ + private nextLinkExpected(): boolean { + return this.life === 'open' && this.reconnectLoop !== undefined; + } - this.outgoing.linkLost(); - this.timers.clear(); - this.incoming.clear(); - this.sock.destroy(); + private loopFor(reconnect: ReconnectOptions | undefined): ReconnectLoop | undefined { + if (!reconnect) return undefined; - // `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(); + return new ReconnectLoop({ + connect: reconnect.connect, + log: this.log, + maxDelay: reconnect.maxDelay, + minDelay: reconnect.minDelay, + onConnected: sock => this.comeBackUp(sock, reconnect.onConnected), + }); } - private onData(chunk: Buffer): void { - this.emit('data', chunk); - this.resetTimers(); + private responseTimeout(): number { + return this.options.responseTimeout ?? defaults.responseTimeout; } - private dispatch(pduObj: PduObject): void { - if (isResp(pduObj)) { - if (!this.outgoing.deliver(pduObj)) { - this.log.debug('session - response with no matching request', { seqNr: pduObj.seqNr }); - } - - return; - } - + private dispatch(link: Link, pduObj: PduObject): void { this.emit('incomingPduObj', pduObj); // Every application hook and listener reached from an incoming PDU funnels through here. - void this.incoming.handle(pduObj).catch((thrown: unknown) => { + void this.incoming.handle(link, pduObj).catch((thrown: unknown) => { const err = errorFrom(thrown); this.log.error('session - a handler threw', { @@ -418,36 +401,15 @@ export class Session extends EventEmitter { }); } - /** A PDU the codec refused. Its header parsed, so the peer gets an answer and the link stays. */ + /** A PDU the codec refused. Its header parsed, so a request gets an answer and the link stays. */ private refuse(refused: PduRefusedError): void { const { cmdId, cmdName, seqNr } = refused.header; this.emit('sessionError', refused); // A response carries a sequence number of ours, so writing one back lands in the peer's space. - if (isResp(refused.header)) { - this.outgoing.settleRefused(seqNr, refused); - - return; - } + if (isResp(refused.header)) return; 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..92b7524 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 { HeldMessage, 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'; @@ -17,16 +17,18 @@ import type { SmppServer } from '../src/server.ts'; import type { TestContext } from 'node:test'; import { HeldMessages } from '../src/held-messages.ts'; import { IncomingRequests, refusedSegmentStatus } from '../src/incoming-requests.ts'; +import { OutgoingRequests } from '../src/outgoing-requests.ts'; import { UnansweredError } from '../src/unanswered-error.ts'; import { createSms } from '../src/sms.ts'; -import { LinkLife } from '../src/link-life.ts'; +import { Link } from '../src/link.ts'; import { SendWindow } from '../src/send-window.ts'; import { Reassembler, decodeSegments } from '../src/reassembly.ts'; 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,15 +103,41 @@ function abortAfter( }); } -function incomingOn(session: Session, options: Partial = {}): IncomingRequests { - 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') }), +const noop = (): void => undefined; + +/** A link over a socket nobody reads or writes, so what arrives on it is what the test hands it. */ +function linkOn(session: Session, options: { bound?: boolean; log?: SmppLog } = {}): Link { + return new Link({ + bound: options.bound ?? true, + held: { + max: defaults.maxHeldMessages, + maxOctets: defaults.maxHeldOctets, + sendPastDrain: () => Promise.resolve({ err: new Error('never sent') }), + session, + timeout: defaults.heldMessageTimeout, + }, + log: options.log ?? silentLog, + on: { data: noop, enquireLink: noop, framed: noop, lost: noop, refused: noop, request: noop }, + reassembly: { max: 10, onLost: noop, timeout: 10_000 }, + responseTimeout: 100, + sock: new net.Socket(), + }); +} + +type Inbound = { close: () => void; handle: (pduObj: PduObject) => Promise; link: Link }; + +/** The inbound path of one link, fed PDUs by hand. */ +function incomingOn(session: Session, options: Partial = {}): Inbound { + const log = options.log ?? silentLog; + const link = linkOn(session, { log }); + const incoming = new IncomingRequests({ + dlrMerger: new DlrMerger({ log, max: 10, timeout: 10_000 }), + log, session, ...options, }); + + return { close: () => { link.close(); }, handle: pduObj => incoming.handle(link, pduObj), link }; } function submitPdu(seqNr: number, cmdStatus: ErrorName = 'ESME_ROK'): PduObject { @@ -749,21 +777,21 @@ describe('reconnect', () => { closeAfter(t, session); - const link = new LinkLife({ log: silentLog, reconnects: true, timeout: 100 }); - const incoming = incomingOn(session, { link, onRequest: async () => { await delay(10); return false; } }); + const onRequest = async (): Promise => { await delay(10); return false; }; + const incoming = incomingOn(session, { onRequest }); let messages = 0; session.on('sms', () => { messages++; }); const handled = incoming.handle(submitPdu(1)); - link.drop(); + incoming.link.close(); await handled; assert.equal(messages, 0); - await incoming.handle(submitPdu(2)); + await incomingOn(session, { onRequest }).handle(submitPdu(2)); assert.equal(messages, 1, 'the harness delivers a message whose link stayed'); }); @@ -1403,92 +1431,130 @@ describe('sends across a reconnect', () => { }); }); -describe('LinkLife', () => { - test('refuses a hold whose deadline has already passed', async () => { +describe('OutgoingRequests waiting for a link', () => { + type View = { closing: boolean; link: Link; nextExpected: boolean }; + + /** Requests against a link that is down, with the session's answers under the test's control. */ + function outgoingOn(t: TestContext, view: View, options: { now?: () => number; responseTimeout: number }): OutgoingRequests { + closeAfter(t, view.link.held.session); + + return new OutgoingRequests({ + links: { + closing: () => view.closing, + current: () => view.link, + nextExpected: () => view.nextExpected, + }, + log: silentLog, + maxOutstanding: 10, + now: options.now, + responseTimeout: options.responseTimeout, + }); + } + + function down(): View { + return { closing: false, link: linkOn(new Session({ sock: new net.Socket() }), { bound: false }), nextExpected: true }; + } + + // A link released to it that still carries nothing sends it back to wait, and the budget is one. + test('refuses a request whose budget ran out while it waited for a link', async t => { let now = 0; - const link = new LinkLife({ log: silentLog, now: () => now, reconnects: true, timeout: 100 }); - const waitForLink = link.hold(undefined); + const view = down(); + const outgoing = outgoingOn(t, view, { now: () => now, responseTimeout: 100 }); + const held = outgoing.request({ cmdName: 'enquire_link' }, {}); - link.drop(); now = 101; + outgoing.linkBound(); - const held = await waitForLink(); - - assert.match(held.err?.message ?? '', /did not come back in time/); + assert.match((await held).err?.message ?? '', /did not come back in time/); }); test('holds on a timer that keeps the process alive', async () => { - const link = new LinkLife({ log: silentLog, reconnects: true, timeout: 10_000 }); + const view = down(); + const outgoing = new OutgoingRequests({ + links: { closing: () => false, current: () => view.link, nextExpected: () => true }, + log: silentLog, + maxOutstanding: 10, + responseTimeout: 10_000, + }); const timers = (): number => process.getActiveResourcesInfo().filter(name => name === 'Timeout').length; - - link.drop(); - const before = timers(); - const held = link.hold(undefined)(); + const held = outgoing.request({ cmdName: 'enquire_link' }, {}); assert.equal(timers(), before + 1, 'an unref\'d timer is not counted here, which is the point'); - link.open(); + outgoing.over(); - assert.deepEqual(await held, {}); + assert.match((await held).err?.message ?? '', /Session is closed/); }); - // addEventListener never fires for a signal that already aborted, so it would wait out the timeout. - 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(); + test('gives up at once on a signal that fires while it waits', async t => { + const view = down(); + const outgoing = outgoingOn(t, view, { responseTimeout: 100 }); + const controller = new AbortController(); + const held = outgoing.request({ cmdName: 'enquire_link' }, { signal: controller.signal }); - const held = await link.hold(AbortSignal.abort())(); + controller.abort(); - assert.match(held.err?.message ?? '', /Aborted while waiting for a link/); + assert.match((await held).err?.message ?? '', /Aborted while waiting for a link/); }); - test('awaits the next link only while down with one on its way', () => { - const link = new LinkLife({ log: silentLog, reconnects: true, timeout: 100 }); + test('waits only while a link is on its way, and is refused as closed otherwise', async t => { + const view = down(); + const outgoing = outgoingOn(t, view, { responseTimeout: 0 }); + const waiting = outgoing.request({ cmdName: 'enquire_link' }, {}); + + assert.equal(await within(30, waiting), undefined, 'down with a link to come: it waits'); + + view.nextExpected = false; + + const refused = await outgoing.request({ cmdName: 'submit_sm' }, {}); + + assert.match(refused.err?.message ?? '', /closed/, 'down with none to come'); + + outgoing.over(); - assert.equal(link.awaitsNextLink(), false, 'up'); - link.drop(); - assert.equal(link.awaitsNextLink(), true, 'down, returning'); - link.attach(); - assert.equal(link.awaitsNextLink(), true, 'attached, not yet bound'); - link.open(); - assert.equal(link.awaitsNextLink(), false, 'reopened'); - link.drop(); - link.stop(); - assert.equal(link.awaitsNextLink(), false, 'down, stopped'); - assert.match(link.refusal()?.message ?? '', /closed/, 'stopped while down'); - link.end(); - assert.equal(link.awaitsNextLink(), false, 'ended'); + assert.match((await waiting).err?.message ?? '', /Session is closed/, 'and what waited is told the same'); }); - test('drops an attached link once, counts each drop, and names the event it warrants', () => { - const link = new LinkLife({ log: silentLog, reconnects: true, timeout: 100 }); - const generation = link.generation(); + test('refuses a new request during a drain only while a link could carry it', async t => { + const view = down(); + const outgoing = outgoingOn(t, view, { responseTimeout: 100 }); - assert.equal(link.drop(), 'disconnected'); - assert.equal(link.drop(), undefined, '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.equal(link.generation(), generation + 2); - assert.equal(new LinkLife({ log: silentLog, reconnects: false, timeout: 100 }).drop(), 'close'); + view.closing = true; + view.nextExpected = false; + + assert.match((await outgoing.request({ cmdName: 'enquire_link' }, {})).err?.message ?? '', /Session is closed/); + + view.link.markBound(); + + assert.match((await outgoing.request({ cmdName: 'enquire_link' }, {})).err?.message ?? '', /shutting down/); }); +}); + +describe('Link', () => { + test('closes once, taking the requests waiting on it and its bound-ness with it', async t => { + const session = new Session({ sock: new net.Socket() }); - test('releases a held request with the reason once the link ends', async () => { - const link = new LinkLife({ log: silentLog, reconnects: true, timeout: 0 }); + closeAfter(t, session); - link.drop(); + const link = linkOn(session, { bound: false }); - const held = link.hold(undefined)(); + assert.equal(link.canCarry(), false, 'a reconnect\'s link carries nothing until its bind is answered'); + link.markBound(); + assert.equal(link.canCarry(), true); - link.end(); + const waiting = link.pending.wait(1, { timeout: 0 }); - 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.close(); + link.close(); + + assert.equal(link.isClosed(), true); + assert.equal(link.canCarry(), false); + assert.equal(link.sock.destroyed, true); + assert.match((await waiting).err?.message ?? '', /Session closed before a response arrived/); + + link.markBound(); + assert.equal(link.canCarry(), false, 'closed is final'); }); }); @@ -1542,7 +1608,7 @@ describe('held message bounds', () => { return [submitPdu(seqNr)]; } - function offer(held: HeldMessages, seqNr: number): MessageHold { + function offer(held: HeldMessages, seqNr: number): HeldMessage { const hold = held.offer(message(seqNr)); assert.ok(hold); @@ -1562,7 +1628,6 @@ describe('held message bounds', () => { return new HeldMessages({ ...options, - link: new LinkLife({ log: silentLog, reconnects: false, timeout: 100 }), log: silentLog, sendPastDrain: () => Promise.resolve({ err: new Error('never sent') }), session, @@ -1658,7 +1723,7 @@ describe('held message bounds', () => { assert.equal(answers.at(-1), 'ESME_RTHROTTLED'); assert.equal(warnings.length, 1); - incoming.clear(); + incoming.close(); }); test('holds a message detached from the chunk it was read from', async t => { @@ -1678,7 +1743,7 @@ describe('held message bounds', () => { assert.ok(Buffer.isBuffer(retained)); assert.notEqual(retained.buffer, chunk.buffer); - incoming.clear(); + incoming.close(); }); test('gives up on a message the application never answers', t => {