diff --git a/AGENTS.md b/AGENTS.md
index 882bb24..244455b 100644
--- a/AGENTS.md
+++ b/AGENTS.md
@@ -38,6 +38,7 @@ These are not preferences. Breaking one is a defect.
```
src/
index.ts Public surface. Named exports only, no default export.
+ bind-direction.ts What a bind declares, and which commands its direction carries at either end
client.ts client() -> { err, session }
server.ts server() -> { err, server }, server owns the listener + close()
session.ts Session: the socket's life, dispatch, events, and the collaborators below
@@ -45,17 +46,17 @@ src/
concat.ts How a PDU says it is a segment: its UDH, or the sar_* TLVs
dlr.ts Delivery receipts: text and TLV parsing, receipt status codes
dlr-merger.ts DlrMerger: per-segment receipts counted into one MessageDlr
+ drain.ts drain(): the shutdown's two waits and their budgets, and IdleWaiters, the wait on a count
error-from.ts An untyped value as error material: errorFrom() an Error, namedValue() a name
expiring-groups.ts ExpiringGroups: the capped, weighed, expiring store DlrMerger, HeldMessages and Reassembler share
- held-messages.ts HeldMessages: a message from its `sms` event to its answer, capped and expiring, one MessageHold each
- idle-waiters.ts IdleWaiters: waiting for a count to fall to zero, and what is left of a budget
+ held-messages.ts HeldMessages: a message from its `sms` event to its answer, and the six ways a hold ends
incoming-requests.ts Every request the peer sends: messages, receipts, links, unknown commands
- link-life.ts LinkLife: whether the link lives, and where a request waits for the next one
+ link-life.ts LinkLife: the link's life as one transition table, and where a request waits for the next link
link-timers.ts LinkTimers: the enquire_link heartbeat and the idle timeout
log.ts SmppLog, the logger contract, and silentLog — the default
message.ts Encoding detection, splitting, bit counting, SMPP date formatting
message-body.ts Where an inbound body is: short_message, or the message_payload TLV
- outgoing-requests.ts OutgoingRequests: the window, the pending map and the retry
+ outgoing-requests.ts OutgoingRequests: the window, the pending map and the carry onto the next link
pdu.ts pduToObj / objToPdu / pduReturn — synchronous, result-returning
pdu-framer.ts PduFramer: a byte stream cut into complete PDUs
pdu-refusal.ts A PDU the codec would not read, and the answer SMPP names for it
@@ -67,7 +68,7 @@ src/
retained-pdu.ts A PDU copied off the wire so holding it pins nothing else, and what holding it costs
send-sms.ts submitSms composition and the submitSmParams builder
send-window.ts SendWindow: the maxOutstanding semaphore
- session-options.ts SessionOptions, ReconnectOptions, bind direction and the session defaults
+ session-options.ts SessionOptions, ReconnectOptions, their checks and the session defaults
sms-id.ts Message ids: the peer's notation, the - a segment gets, which response carries one
udh.ts User data header: its length, the concatenation fields of a long SMS and their reference
unanswered-error.ts UnansweredError: it went out and no answer came back
diff --git a/DESIGN.md b/DESIGN.md
new file mode 100644
index 0000000..b4c948e
--- /dev/null
+++ b/DESIGN.md
@@ -0,0 +1,56 @@
+# Draft B: the session's life as one table, the hold as one file, the drain as one function
+
+## Structure and who owns what
+
+| Module | Owns |
+| --- | --- |
+| `link-life.ts` | The link's phase (`up`, `binding`, `down`, `closed`) and the `stopping` flag. `transition(event)` is the whole lifecycle: one switch, five events (`attached`, `bound`, `lost`, `stopping`, `closed`), returning the effects the session runs in order (`dropLink`, `emitDisconnected`, `scheduleReconnect`, `emitClose`, `linkUp`, `stopReconnect`). Also the waiters for a request with no link, released by the same transitions. |
+| `session.ts` | Wiring, and the imperative shell: `apply(event)` runs `link.transition(event)` and then `run(effect)`, a switch that names every side effect the lifecycle has. No `end`/`stop`/`emitClose`/`teardown`/`onClose` — every path in (socket close, unreadable stream, idle timeout, failed rebind, `close()`, `unbind()`) is one `apply()` call. |
+| `held-messages.ts` | The whole held-message flow: the store, the `Sms` creation, the `sms` emit, and each of the six exits as a numbered method. Owned by `Session` directly, so a rejected listener is routed `Session → HeldMessages` in one hop. The hysteresis logging (`refusing`) moved here from `IncomingRequests`: `full()` decides and logs. |
+| `drain.ts` | `drain()`: the shutdown's two waits, both budgets, the one turn given to the application, and the error text. `IdleWaiters` lives here too, as the primitive the two stores wait with. |
+| `incoming-requests.ts` | Routing only: the hook, the bind gate, the dispatch per command, reassembly. It is handed `held` and never constructs, counts or drains it. |
+| `outgoing-requests.ts` | `request()` (refused after `stopping`), `requestPastDrain()` (a receipt), and `carry()`: a recursion in place of `for (;;)`, one attempt per link, the same budget throughout. `idle()` reports a count; the drain owns the words. |
+| `bind-direction.ts` | `bindCommands`, `BindType`, `LinkEnd`, `standsInFor`, `bindCarries`, `checkedBind`, split out of `session-options.ts`, which now holds option types, their checks and the defaults. |
+
+`hold` now has one meaning: none. `LinkLife.budget()` is a request's budget, `HeldMessages` keeps messages, `MessageHold` is gone.
+
+## How each exit reads now
+
+A hold is `Held = { key, pduObjs, working }` in `HeldMessages`; `offer()` builds the `Sms` with three closures and emits it.
+
+1. **Answered** — `sendResp()` puts the response on the wire, calls `answered()`, which is `release(held)`. Synchronous, no `setImmediate`.
+2. **Every listener rejected** — `Session[captureRejectionSymbol]` calls `held.listenerRejected(sms)`; the `WeakMap` finds the `Held`, `working--`, the last one releases.
+3. **No listener, or one threw** — `session.emit('sms')` returned false, so `offer()` releases at once.
+4. **Re-used sequence number** — `keep()` replaces the entry; the old `Held` fails `release()`'s identity check and can free nothing.
+5. **Deadline** — `sweep()`, before each `keep()` and on the store's timer.
+6. **Link gone** — `clear()`, run by the `dropLink` effect.
+
+The receipt that used to depend on the deferred release is now the drain's business: `sendDlr()` always sends through `requestPastDrain()` (a receipt answers a message the drain waits on, whenever it is sent), and `drain()` waits one `setImmediate` turn after the last message is answered before reading the window, so a receipt sent straight after the answer is in the window by then. One line, in the one place that waits on the application.
+
+**Shutdown** — `close()`: `apply('stopping')` (effect `stopReconnect`, before the first await) → `drain()`: messages half on `shutdownTimeout` or the `responseTimeout` fallback, one turn, requests half on what is left → `apply('closed')`: `dropLink` if a socket is attached, then `emitClose`. `unbind()` is the same with the unbind request between the drain and `closed`.
+
+**Reconnect** — socket close, unreadable stream, idle timeout and a failed rebind all `apply('lost')`. While retrying: `down`, effects `dropLink`, `emitDisconnected`, `scheduleReconnect`. Not retrying: `dropLink`, `emitClose`, phase `closed`. The loop's `comeBackUp()` is `attach` → `apply('attached')` → bind → `apply('bound')`, which is `linkUp` (release held sends, timers, `reconnected`), or, if `stopping` landed meanwhile, the `lost` path. Re-entrancy: the phase is updated before any effect runs, so a listener that calls `close()` from `disconnected` sees `down`, and the remaining effects are no-ops.
+
+## Deleted
+
+- `MessageHold` class, `SmsHandlers`' `isHeld` branch, the `setImmediate` in `answered()`, `sendPastDrain` plumbing through `IncomingRequests`.
+- `Session.end()`, `stop()`, `emitClose()`, `teardown()`, `onClose()`, `resetTimers()`, `attach()`; `LinkLife.drop()`'s returned event name and `end()`'s boolean, `isUp`/`isAttached`/`isOver`/`isStopped`/`retrying` as public predicates (`carries`, `attached`, `isStopping`, `awaitsNextLink`, `refusal` remain, all reading `phase`).
+- `IncomingRequests.drain()`, `.listenerRejected()`, `.refusing`; `OutgoingRequests.drain()`; `Session.drain()`'s two-error merge and `answering()`.
+- `idle-waiters.ts` (into `drain.ts`); `for (;;)` in `requestPastDrain`.
+
+## Tests
+
+`npm test`: lint and typecheck clean; 516 tests, 516 pass, 0 fail (baseline 515/515).
+
+Changed, all in `test/session-extras.test.ts`, all reaching into internals:
+
+- `incomingOn()` builds a `HeldMessages` (new `heldMessagesOn()` helper) instead of passing `sendPastDrain`; `standsInFor` imported from `bind-direction.ts`.
+- "drops a message whose link went while onRequest was still running": `link.drop()` → `link.transition('lost')`.
+- `LinkLife` describe: `drop/attach/open/stop/end` → `transition('lost'|'attached'|'bound'|'stopping'|'closed')`; the "names the event it warrants" test now asserts the effect lists, which is the same claim stated in the new vocabulary; one test added: a link bound after the shutdown began is ended, not brought up.
+- Held-message bounds: `offer()` returns the `Sms`, so `first.isHeld()`/`replaced.isHeld()` became "answering the replaced message frees nothing, answering the first frees it" (size before/after `sendResp()`), and `answered.release()` became `await answered.sendResp()`; `heldOn()` stubs `sendReturn` for that.
+
+No assertion was weakened; the `session.test.ts`, `readme.test.ts` and the shutdown/receipt integration tests run unchanged.
+
+## Not changed on purpose
+
+`ExpiringGroups` stays mechanism-only: the three owners' policies genuinely differ (refuse, evict oldest, evict with `spent` memory). `ReconnectLoop.halted` stays: `client()` runs a loop with no session behind it for `fromStart`. The public surface, `index.ts` and README are untouched.
diff --git a/docs/decisions.md b/docs/decisions.md
index c245097..f98e81e 100644
--- a/docs/decisions.md
+++ b/docs/decisions.md
@@ -561,7 +561,7 @@ rule and an index of the titles below.
- **A stream this library cannot frame is a dead link; one PDU it cannot parse is not.**
Maintainer's call, 2026-08-31, narrowed 2026-09-05 via the interop plan: a `command_length` below
16 or above `maxPduLength` leaves nothing that can say where the next PDU starts, so it tears the
- link down through `teardown()` and the reconnect loop retries it on a fresh socket with a fresh
+ link down through the `dropLink` effect and the reconnect loop retries it on a fresh socket with a fresh
framer. Every other codec failure honoured `command_length`, so the stream is still in sync and
the next PDU starts where it says — tearing the link down there cost one peer half its receipts
and its MO to a reconnect loop (`interop-tests/findings/01-smscsim.md`), and left the peer waiting
@@ -668,7 +668,7 @@ rule and an index of the titles below.
the peer, whose every request is bounded by `responseTimeout` unless the caller set that to 0 as
well, and unsafe for the application, which nothing bounds — `close()` is what you reach for when
the application is stuck, so it may not block on the application coming unstuck. That half falls
- back to `responseTimeout`, the same answer `LinkLife`'s hold already takes — and to that
+ back to `responseTimeout`, the same answer `LinkLife`'s budget already takes — and to that
option's default where it is 0 as well, since neither option is an answer about the application.
- **What the application holds unanswered is capped on constants, and a message past the cap is
@@ -694,10 +694,10 @@ rule and an index of the titles below.
error, the one an SMSC retries on (goal 3).
- **A reconnect keeps the delivery-receipt merges; everything else the link held is dropped.**
- `onDelivery()` answers each receipt before the group it belongs to is complete, and `teardown()`
+ `onDelivery()` answers each receipt before the group it belongs to is complete, and `dropLink`
runs on every path — an idle timeout and a failed rebind, not only `close()` — so clearing the
merges there loses receipts no peer has a reason to send again. They are cleared where the session
- is over instead. Inbound segments stay in `teardown()`: a concatenation reference is the
+ is over instead. Inbound segments stay in `dropLink`: a concatenation reference is the
peer's own counter, so a half-arrived group kept across a drop would take a later message's
segments as readily as the rest of its own, and goal 2 will not hand the application a message
assembled that way. What goes there is traffic already answered, which is why each group reaches
@@ -705,7 +705,7 @@ rule and an index of the titles below.
- **The bind state is the session's, `bound()` alone writes it, and it holds through a reconnect's
gap.** Maintainer's call, 2026-09-28. `client()` and `server()` record their bind through
- `bound()`, the call a hand-wired session makes, so the state has one writer — goal 8. Clearing the state at `teardown()` was rejected: `bindAllows()` and
+ `bound()`, the call a hand-wired session makes, so the state has one writer — goal 8. Clearing the state at `dropLink` was rejected: `bindAllows()` and
`acceptsOptionalParams()` then answer yes to everything while the link is down, so a
receiver-bound client queues a `submit_sm` the peer refuses and a receipt built then carries TLVs a
pre-3.4 peer must not get — goal 4. Valid while the reconnect loop binds again with the same bind
@@ -757,7 +757,7 @@ rule and an index of the titles below.
attached, which admits a send one round trip before the bind is answered, and collaborators that
ask the session, which answered the same question two ways at admit and at release.
`ReconnectLoop.halted` is the one other flag, because `client()` also runs a loop with no session
- behind it for `fromStart`; a session's loop is stopped by `Session.stop()` alone.
+ behind it for `fromStart`; a session's loop is stopped by the `stopReconnect` effect alone.
## Internals and tests
diff --git a/src/bind-direction.ts b/src/bind-direction.ts
new file mode 100644
index 0000000..ac251da
--- /dev/null
+++ b/src/bind-direction.ts
@@ -0,0 +1,73 @@
+import type { Result } from './result.ts';
+import { quoted } from './error-from.ts';
+
+export const bindCommands: readonly string[] = [
+ 'bind_receiver',
+ 'bind_transceiver',
+ 'bind_transmitter',
+];
+
+export type BindType = 'receiver' | 'transceiver' | 'transmitter';
+
+/** Which end of the link a session is. Only `server()` is the SMSC; everything else is the ESME. */
+export type LinkEnd = 'esme' | 'smsc';
+
+export function bindTypeFromCommand(cmdName: string): BindType | undefined {
+ if (cmdName === 'bind_receiver') return 'receiver';
+ if (cmdName === 'bind_transceiver') return 'transceiver';
+ if (cmdName === 'bind_transmitter') return 'transmitter';
+
+ return undefined;
+}
+
+/**
+ * Which message-carrying command an inbound one stands in for. Every command but `data_sm` names
+ * its own direction; that one travels either way, so the end it arrived at is what says.
+ */
+export function standsInFor(cmdName: string, linkEnd: LinkEnd): string {
+ if (cmdName !== 'data_sm') return cmdName;
+
+ return linkEnd === 'smsc' ? 'submit_sm' : 'deliver_sm';
+}
+
+/**
+ * Whether a bind direction carries a command at all. A receiver-bound ESME submits nothing and a
+ * transmitter-bound one is delivered nothing, whichever end of the link is looking. A session that
+ * has not bound carries everything, since nothing has declared a direction yet.
+ */
+export function bindCarries(
+ bindType: BindType | undefined,
+ cmdName: string,
+ linkEnd: LinkEnd,
+): boolean {
+ const carried = standsInFor(cmdName, linkEnd);
+
+ if (bindType === 'receiver') return carried !== 'submit_sm';
+ if (bindType === 'transmitter') return carried !== 'deliver_sm';
+
+ return true;
+}
+
+/** SMPP 3.4: a peer that declares no version at all is one from before optional parameters. */
+export const undeclaredInterfaceVersion = 0x00;
+
+export type SessionBind = { as: BindType; peerVersion: number };
+
+function isBindType(value: unknown): value is BindType {
+ return typeof value === 'string' && bindTypeFromCommand(`bind_${value}`) !== undefined;
+}
+
+/** A bind as `Session.bound()` records it: undefined declares no version, which is pre-3.4. */
+export function checkedBind(bindType: unknown, declaredVersion: unknown): Result<{ bind: SessionBind }> {
+ if (!isBindType(bindType)) {
+ return { err: new Error(`bindType must be receiver, transceiver or transmitter, the bind command's name without "bind_", got ${quoted(bindType)}`) };
+ }
+
+ if (declaredVersion === undefined) return { bind: { as: bindType, peerVersion: undeclaredInterfaceVersion } };
+
+ if (typeof declaredVersion !== 'number' || !Number.isInteger(declaredVersion) || declaredVersion < 0 || declaredVersion > 0xFF) {
+ return { err: new Error(`declaredVersion must be an integer 0-255, the interface_version param or the sc_interface_version TLV's tagValue, or undefined where the peer declared none, got ${quoted(declaredVersion)}`) };
+ }
+
+ return { bind: { as: bindType, peerVersion: declaredVersion } };
+}
diff --git a/src/client.ts b/src/client.ts
index 6b7299c..28d9fc9 100644
--- a/src/client.ts
+++ b/src/client.ts
@@ -1,6 +1,7 @@
import type { ConnectionOptions } from 'node:tls';
import type { Result, VoidResult } from './result.ts';
-import type { BindType, ReconnectOptions } from './session-options.ts';
+import type { BindType } from './bind-direction.ts';
+import type { ReconnectOptions } from './session-options.ts';
import type { SmppLog } from './log.ts';
import type { SmsIdFormat } from './sms-id.ts';
import type { Socket } from 'node:net';
diff --git a/src/drain.ts b/src/drain.ts
new file mode 100644
index 0000000..5251ad0
--- /dev/null
+++ b/src/drain.ts
@@ -0,0 +1,112 @@
+import type { SmppLog } from './log.ts';
+import type { VoidResult } from './result.ts';
+import { defaults } from './session-options.ts';
+
+/** Everything waiting for a count to fall to zero, and how such a wait is cut short. */
+export class IdleWaiters {
+ private readonly waiting: (() => void)[] = [];
+
+ /** Wakes everything waiting, whatever the count reads now. */
+ settle(): void {
+ for (const resolve of this.waiting.splice(0)) {
+ resolve();
+ }
+ }
+
+ /**
+ * Resolves 0 once nothing is left, or with what still is when the timeout or the signal cuts the
+ * wait short. A timeout of 0 waits forever.
+ */
+ wait(remaining: () => number, timeout: number, signal: AbortSignal | undefined): Promise {
+ if (remaining() === 0) return Promise.resolve(0);
+
+ if (signal?.aborted === true) return Promise.resolve(remaining());
+
+ return new Promise(resolve => {
+ let timer: NodeJS.Timeout | undefined = undefined;
+ const done = (): void => {
+ const index = this.waiting.indexOf(done);
+
+ if (timer) clearTimeout(timer);
+ if (index !== -1) this.waiting.splice(index, 1);
+
+ signal?.removeEventListener('abort', done);
+ resolve(remaining());
+ };
+
+ if (timeout > 0) {
+ timer = setTimeout(done, timeout);
+ timer.unref();
+ }
+
+ signal?.addEventListener('abort', done, { once: true });
+ this.waiting.push(done);
+ });
+ }
+}
+
+/** The two stores a shutdown waits on, and whether the link still carries what they hold. */
+export type Drainable = {
+ linkCarries: () => boolean;
+ /** Resolves 0 once the application has answered every message, or with how many it has not. */
+ messagesUnanswered: (timeout: number, signal: AbortSignal | undefined) => Promise;
+ /** Resolves 0 once every request sent is answered, or with how many are not. */
+ requestsUnfinished: (timeout: number, signal: AbortSignal | undefined) => Promise;
+};
+
+export type DrainOptions = {
+ log: SmppLog;
+ responseTimeout: number;
+ shutdownTimeout: number;
+ signal: AbortSignal | undefined;
+};
+
+/** What is left of a budget, in the shape a wait takes it: 0 waits forever. */
+function leftOf(deadline: number): number {
+ return deadline === 0 ? 0 : Math.max(1, deadline - Date.now());
+}
+
+/** The application half may never wait forever: nothing else ends that wait. */
+function answeringBudget(options: DrainOptions): number {
+ if (options.shutdownTimeout > 0) return options.shutdownTimeout;
+
+ return options.responseTimeout > 0 ? options.responseTimeout : defaults.responseTimeout;
+}
+
+function report(log: SmppLog, unanswered: number, unfinished: number): VoidResult {
+ const lost: string[] = [];
+
+ if (unanswered > 0) {
+ log.warn('drain - shutting down with messages unanswered', { unanswered });
+ lost.push(`Shut down with ${String(unanswered)} message(s) unanswered`);
+ }
+
+ if (unfinished > 0) {
+ log.warn('drain - shutting down with requests unfinished', { unfinished });
+ lost.push(`Shut down with ${String(unfinished)} request(s) unfinished`);
+ }
+
+ return lost.length > 0 ? { err: new Error(lost.join('; ')) } : {};
+}
+
+/**
+ * Waits out the messages the application holds, then the requests on the wire. Answering a message
+ * can put a receipt on the wire; nothing on the wire produces a message, so that order covers both.
+ */
+export async function drain(stores: Drainable, options: DrainOptions): Promise {
+ // No bound link, so nothing is on the wire to wait out.
+ if (!stores.linkCarries()) return {};
+
+ const deadline = options.shutdownTimeout > 0 ? Date.now() + options.shutdownTimeout : 0;
+ const unanswered = await stores.messagesUnanswered(answeringBudget(options), options.signal);
+
+ // One turn, so a receipt sent straight after the last answer is in the window before it is read.
+ await new Promise(resolve => { setImmediate(resolve); });
+
+ const unfinished = await stores.requestsUnfinished(leftOf(deadline), options.signal);
+
+ // The link went before the drain finished, so an empty window says nothing about the peer.
+ if (!stores.linkCarries()) return { err: new Error('The session closed before the drain finished') };
+
+ return report(options.log, unanswered, unfinished);
+}
diff --git a/src/error-from.ts b/src/error-from.ts
index 5d77236..e622a0a 100644
--- a/src/error-from.ts
+++ b/src/error-from.ts
@@ -15,3 +15,8 @@ const printable: readonly string[] = ['boolean', 'number', 'string'];
export function namedValue(value: unknown): string {
return printable.includes(typeof value) ? String(value) : typeof value;
}
+
+/** A string in quotes, so an empty one and one of spaces are visible; anything else as namedValue(). */
+export function quoted(value: unknown): string {
+ return typeof value === 'string' ? JSON.stringify(value) : namedValue(value);
+}
diff --git a/src/held-messages.ts b/src/held-messages.ts
index b9e740e..b8d004c 100644
--- a/src/held-messages.ts
+++ b/src/held-messages.ts
@@ -1,26 +1,29 @@
-import type { LinkLife } from './link-life.ts';
-import type { PduObject, PduObjectInput } from './pdu.ts';
-import type { Result } from './result.ts';
+import type { PduObject } from './pdu.ts';
import type { Session } from './session.ts';
-import type { SmsHandlers } from './sms.ts';
+import type { Sms, SmsHandlers } from './sms.ts';
import type { SmppLog } from './log.ts';
import { ExpiringGroups } from './expiring-groups.ts';
-import { IdleWaiters } from './idle-waiters.ts';
+import { IdleWaiters } from './drain.ts';
import { createSms } from './sms.ts';
import { retainedOctets } from './retained-pdu.ts';
export type HeldMessagesOptions = {
- link: LinkLife;
+ /** Changes with every link, so a message can tell the one it arrived on is gone. */
+ linkGeneration: () => number;
log: SmppLog;
max: number;
maxOctets: number;
/** Injected so expiry can be exercised without a wall clock. */
now?: (() => number) | undefined;
- sendPastDrain: SmsHandlers['send'];
+ /** How a receipt goes out: past a shutdown's refusal, since it answers a message the drain waits for. */
+ sendReceipt: SmsHandlers['send'];
session: Session;
timeout: number;
};
+/** A message handed to the application, and how many of its listeners are still working on it. */
+type Held = { key: string; pduObjs: PduObject[]; working: number };
+
/** The peer's own sequence number, which is what our answer to this message will carry. */
function keyOf(pduObjs: PduObject[]): string | undefined {
const first = pduObjs[0];
@@ -28,69 +31,19 @@ function keyOf(pduObjs: PduObject[]): string | undefined {
return first ? String(first.seqNr) : undefined;
}
-type HoldRoute = Pick;
-
/**
- * One message offered to the application, and the handlers its `Sms` answers through. A drain
- * waits on it until the first of: `answered()`, every listener that took it rejecting, no listener
- * taking it or one throwing, a later message on its sequence number, its deadline, or the link going.
+ * The messages handed to the application that it has not answered yet, which is what a shutdown
+ * waits on. A hold ends the first of six ways, numbered below: 1 answered, 2 the last listener
+ * working on it rejecting, 3 no listener taking it or one throwing, 4 a later message on its
+ * sequence number, 5 its deadline, 6 the link going.
*/
-export class MessageHold implements SmsHandlers {
- private readonly generation: number;
- private readonly heldMessages: HeldMessages;
- private readonly pduObjs: PduObject[];
- private readonly route: HoldRoute;
- private working: number;
-
- constructor(heldMessages: HeldMessages, route: HoldRoute, pduObjs: PduObject[], listeners: number) {
- this.generation = route.link.generation();
- this.heldMessages = heldMessages;
- this.pduObjs = pduObjs;
- this.route = route;
- this.working = listeners;
- }
-
- /** Whether a drain is still waiting for this message to be answered. */
- isHeld(): boolean {
- return this.heldMessages.holds(this.pduObjs);
- }
-
- /** A turn later, so a `sendDlr()` called straight after `sendResp()` still goes out past a drain. */
- answered(): void {
- setImmediate(() => { this.release(); });
- }
-
- lostLink(): boolean {
- return this.route.link.generation() !== this.generation;
- }
-
- /** A rejection leaves the other listeners running, so only the last one to fail gives the message up. */
- listenerGaveUp(): void {
- this.working--;
-
- if (this.working <= 0) this.answered();
- }
-
- /** At once, for a message nobody took or a listener threw on: that is not work a shutdown can wait for. */
- release(): void {
- this.heldMessages.release(this.pduObjs);
- }
-
- /** A receipt for a message still held is what a drain waits for, so it goes out past the drain. */
- send(input: PduObjectInput): Promise> {
- return this.isHeld() ? this.route.sendPastDrain(input) : this.route.session.send(input);
- }
-}
-
-/** The messages handed to the application that it has not answered yet, held by their segments. */
export class HeldMessages {
- private readonly held: ExpiringGroups;
+ private readonly held: ExpiringGroups;
private readonly idleWaiters = new IdleWaiters();
- private readonly log: SmppLog;
- private readonly maxOctets: number;
/** A rejecting listener hands the message back as an `unknown`, so its hold is found by identity. */
- private readonly offered = new WeakMap