Files
lilleman 576b713af4
Test / lint (pull_request) Successful in 23s
Test / test (18) (pull_request) Successful in 30s
Test / test (20) (pull_request) Successful in 31s
Test / test (22) (pull_request) Successful in 31s
Test / test (24) (pull_request) Successful in 30s
Test / test (26) (pull_request) Successful in 30s
Mirror / push (push) Successful in 5s
Plan the comprehension rewrite and file its background
2026-09-30 12:04:56 +02:00

2727 lines
108 KiB
Diff

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 <base>-<n> 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<number>;
+ /** Requests on the wire or queued behind the send window. */
+ requests: (timeout: number, signal: AbortSignal | undefined) => Promise<number>;
+};
+
+/** 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<VoidResult> {
+ 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<HeldMessagesOptions, 'link' | 'sendPastDrain' | 'session'>;
-
/**
- * 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<Result<{ pduObj: PduObject }>> {
- 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<PduObject[]>;
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<object, MessageHold>();
- private readonly route: HoldRoute;
+ private readonly offered = new WeakMap<object, HeldMessage>();
+ 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<LostGroup['reason'], string> = {
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<void> {
- 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<void> {
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<void> {
+ private async route(link: Link, pduObj: PduObject): Promise<void> {
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<VoidResult> {
- 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<void> {
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<void> {
+ private async onDelivery(link: Link, pduObj: PduObject): Promise<void> {
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<boolean> {
- 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<void> {
+ 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<void> {
- 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<Waiter>();
- 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<VoidResult> {
- 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<VoidResult> {
- 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<VoidResult> {
- this.log.verbose('linkLife - holding a request until a link is back', { timeout: left });
-
- return new Promise<VoidResult>(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<HeldMessagesOptions, 'log'>;
+ /** Undefined or 0 keeps a silent peer forever. */
+ idleTimeout?: number | undefined;
+ log: SmppLog;
+ on: LinkEvents;
+ reassembly: Omit<ReassemblerOptions, 'log'>;
+ 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<Attempt> {
+ // 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<Waiter>();
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<Result<{ pduObj: PduObject }>> {
// 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<Result<{ pduObj: PduObject }>> {
- 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<Result<{ pduObj: PduObject }>> {
+ 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<Attempt> {
+ 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<Result<{ pduObj: PduObject }>> {
- 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<VoidResult> {
- 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<number> {
+ 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<Result<{ link: Link }>> {
+ for (;;) {
+ const link = this.links.current();
- private async attempt(input: PduObjectInput, options: SendOptions): Promise<Attempt> {
- // 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<VoidResult> {
+ this.log.verbose('outgoingRequests - holding a request until a link is back', { timeout: left });
+
+ return new Promise<VoidResult>(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<K extends keyof SessionEvents> = (...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<SessionEvents> {
declare addListener: <K extends keyof SessionEvents>(event: K, listener: SessionListener<K>) => this;
declare off: <K extends keyof SessionEvents>(event: K, listener: SessionListener<K>) => this;
@@ -58,16 +61,16 @@ export class Session extends EventEmitter<SessionEvents> {
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<K extends keyof SessionEvents>(
@@ -98,7 +101,7 @@ export class Session extends EventEmitter<SessionEvents> {
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<SessionEvents> {
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<SessionEvents> {
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<SendSmsResult> {
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<SessionEvents> {
*/
async unbind(): Promise<VoidResult> {
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<SessionEvents> {
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<VoidResult> {
+ 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<VoidResult>,
- ): Promise<VoidResult> {
- 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<VoidResult>): Promise<VoidResult> {
+ 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<VoidResult> {
- 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<SessionEvents> {
});
}
- /** 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<IncomingRequestsOptions> = {}): 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<void>; link: Link };
+
+/** The inbound path of one link, fed PDUs by hand. */
+function incomingOn(session: Session, options: Partial<IncomingRequestsOptions> = {}): 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<boolean> => { 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 => {