diff --git a/AGENTS.md b/AGENTS.md
index 882bb24..06cc397 100644
--- a/AGENTS.md
+++ b/AGENTS.md
@@ -40,14 +40,16 @@ src/
index.ts Public surface. Named exports only, no default export.
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
- sms.ts The live handle emitted as the 'sms' event (sendResp/sendDlr)
+ session.ts Session: the link's life (linkLost, finish, comeBackUp), the drain, dispatch and events
+ sms.ts The live handle handed to onSms (sendResp/sendDlr)
+ bind-direction.ts Bind types and commands, which end of the link this is, what a bind direction carries
concat.ts How a PDU says it is a segment: its UDH, or the sar_* TLVs
+ defaults.ts Every default the session layer runs on, in one object
dlr.ts Delivery receipts: text and TLV parsing, receipt status codes
dlr-merger.ts DlrMerger: per-segment receipts counted into one MessageDlr
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
+ expiring-groups.ts ExpiringGroups: the capped, weighed, expiring store DlrMerger, HandledMessages and Reassembler share
+ handled-messages.ts HandledMessages: the messages whose onSms handler is running, capped, expiring, waited on by a drain
idle-waiters.ts IdleWaiters: waiting for a count to fall to zero, and what is left of a budget
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
@@ -55,7 +57,7 @@ src/
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: request() through the link wait and the window, requestOnLink() for a bind or unbind
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 +69,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, OnSms, OnRequest, and the option checks
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
@@ -84,8 +86,10 @@ src/
Imports point one way: `defs` knows nothing above it but `result.ts`, `pdu` uses `defs`, `session`
uses `pdu`, and `client`/`server` use `session`. The ways back up are the `Session` handed to
-`createSms()`, `HeldMessages` and `IncomingRequests`, which call back into it, and to `OnRequest`
-and `onConnected` in `session-options.ts`, all imported as a type only.
+`IncomingRequests` (the hook's argument, the bind predicates, `emit` and `close`) and to
+`createSms()` (the public `sms.session`), and to `OnRequest`, `OnSms` and `onConnected` in
+`session-options.ts`, all imported as a type only. Everything else a collaborator does to the
+session is a named closure in its options: `answer`, `sendReceipt`, `report`.
**Parameter order is wire order.** The key order inside `cmds.*.params` is the order the fields are
written to and read from the buffer. Never sort those alphabetically — the alphabetical-ordering
@@ -224,6 +228,7 @@ this is not a changelog.
### [The public surface](docs/decisions.md#the-public-surface)
+- A request the application must answer is a handler option; a fact it may watch is an event.
- `Session` is publicly constructible, which is what makes `SessionOptions` and `ReconnectOptions`
public too.
- `acceptsOptionalParams()` and `bindAllows()` are predicates, not chokepoints.
@@ -288,15 +293,15 @@ this is not a changelog.
- A stream this library cannot frame is a dead link; one PDU it cannot parse is not.
- A deliberate shutdown drains; an unusable link and an abort do not.
- `sendSms()` puts every segment of a message on the wire together.
-- Every segment of a concatenated message is answered as it arrives, so `sendResp()` on one is the
- application's own signal rather than the peer's answer.
+- Every segment of a concatenated message is answered as it arrives, so `sendResp()` on one writes
+ nothing.
- `server()` composes the application's `onRequest` after its own bind handling, and offers it every
request that handling did not answer.
-- The drain waits on the messages the application holds, and `sendResp()` is what says it is done
- with one.
-- The drain's wait on the application ignores `shutdownTimeout: 0`.
-- What the application holds unanswered is capped on constants, and a message past the cap is
- refused.
+- The drain waits on the `onSms` handlers still running, and a handler's promise is what says it
+ is done with a message.
+- A handler that fails before answering has the message refused for it; one that fails after has
+ its answer stand.
+- What the running handlers hold is capped on constants, and a message past the cap is refused.
- A store at its bound answers `ESME_RTHROTTLED` to a submission and `ESME_RX_T_APPN` to a delivery,
a `data_sm` by whichever it stands in for.
- A reconnect keeps the delivery-receipt merges; everything else the link held is dropped.
diff --git a/CHANGELOG.md b/CHANGELOG.md
index 7cb5ff7..ab02284 100644
--- a/CHANGELOG.md
+++ b/CHANGELOG.md
@@ -2,6 +2,25 @@
## 0.6.0 (unreleased)
+- **Inbound messages go to an `onSms` handler option, and the `sms` event is gone.** Give it to
+ `client()`, `server()` or `Session`; `sms.session` says which session a server's message came in
+ on. The message is held — counted toward the bound, and waited for by `close()` — until the
+ promise the handler returns settles. A handler that throws or rejects before answering has the
+ message refused for it (`ESME_RTHROTTLED` on a submission, `ESME_RX_T_APPN` on a delivery), so the
+ peer retries; with no `onSms` at all, every message is refused that way and reported on
+ `sessionError`. The release used to come from `sendResp()`, one turn later, or from every
+ listener rejecting.
+- `close()` and `unbind()` refuse new `sendSms()` and `send()` calls, and let `sendResp()` and
+ `sendDlr()` through for as long as the session lives, so a receipt sent after the answer goes out
+ wherever in the handler it is sent. `shutdownTimeout: 0` now waits for a handler as long as it
+ waits for a request; it used to fall back to `responseTimeout` for the messages.
+- `sendResp()` refuses a second call on the same message, and a call that failed leaves `sms.smsId`
+ as it was.
+- `encoding: 'GSM'` names the GSM 03.38 alphabet, in `sendSms()`, `EncodingName`, `encodings`,
+ `dataCodingByEncoding`, `detect()` and the rest. It was `'ASCII'`, which named the one thing the
+ alphabet is not.
+- A client whose rebind was answered on a link the peer had already dropped no longer comes up on
+ the dead socket; the reconnect loop retries instead.
- `client()` now bounds each connect attempt at 10 seconds, the TLS handshake included, and reports
one that expires as an ordinary connect failure, so `reconnect` retries it on its usual backoff.
A connect previously waited the operating system out, around 130 s on Linux against a host that
@@ -42,9 +61,9 @@
the cap is lost, since its segments were already answered.
- A message arriving while the application holds 1000 unanswered, or 64 MiB of them counted the way
`maxOctets` counts segments, is refused with `ESME_RTHROTTLED` (`ESME_RX_T_APPN` on a
- delivery), so the peer keeps it and retries. **Call `sendResp()` on every `sms`, multipart
- included**: 1000 left unanswered now stop inbound traffic for up to five minutes, where the oldest
- used to be dropped with a warning.
+ delivery), so the peer keeps it and retries. **Return from `onSms` on every message, multipart
+ included**: 1000 handlers still running now stop inbound traffic for up to five minutes, where the
+ oldest message used to be dropped with a warning.
- A `submit_sm` segment the reassembly buffer has no room for is refused with `ESME_RTHROTTLED`,
where it was `ESME_RMSGQFUL`.
- `server()` refuses a `maxOctets` below 1 or not a whole number, `Infinity` included, like its
diff --git a/DESIGN.md b/DESIGN.md
new file mode 100644
index 0000000..66a8b5d
--- /dev/null
+++ b/DESIGN.md
@@ -0,0 +1,83 @@
+# Draft D: the message handler is the contract
+
+Round one showed the internals were as simple as the `sms` event allowed: a hold with six exits,
+three of them about how many listeners there were and how they failed. This draft changes the
+contract so the hold has one exit, and the internals follow.
+
+## The contract, as an implementer sees it
+
+- **Receiving.** `onSms: sms => Promise | void`, an option on `client()`, `server()` and
+ `Session`. One handler per session; a server's handler sees every session, `sms.session` says which.
+- **Answering.** `await sms.sendResp()` once, `sendResp({ smsId, status })` to name or refuse. A
+ multipart message was answered per segment as it arrived (`sms.answeredOnArrival`), so there it
+ writes nothing. A second call is refused. `sms.sendDlr()` sends the receipt, before or after the
+ handler returns.
+- **What happens if you don't.** Return without answering: the peer waits on you, the shutdown does
+ not. Throw or reject before answering: the message is refused with the retry status and reported
+ on `sessionError`. Throw after answering: the answer stands. No `onSms`: every message is refused
+ the same way. Never return: after five minutes the message stops counting.
+- **Bound.** 1000 running handlers, or 64 MiB of messages under them, refuse new messages until half.
+- **Sending.** `sendSms()`/`send()` wait for a bound link and a window slot; what never reached the
+ socket waits for the next link, what did is `unanswered`. Unchanged.
+- **Shutdown.** `close()`/`unbind()`: refuse new `sendSms()`/`send()`; `sendResp()`/`sendDlr()`
+ still go out. Wait up to `shutdownTimeout` for running handlers, then for requests on the wire.
+ `0` waits forever for both. Tear down, report what did not finish.
+- **Events** are facts nothing waits on: `dlr`, `messageDlr`, `close`, `disconnected`,
+ `reconnected`, `sessionError`, `data`, `incomingPdu`, `incomingPduObj`. Rule in README: a request
+ you must answer is a handler; a fact you may watch is an event.
+- `encoding: 'GSM'` names GSM 03.38.
+
+## Internal structure and ownership
+
+| Module | Owns |
+| --- | --- |
+| `session.ts` | The link's life in three methods: `linkLost()` (socket close, unreadable stream, idle, failed rebind — one entry, decides retry or end), `finish()` (once; drop, clear merges, `close`), `comeBackUp()`. `drain()` is the shutdown in eleven lines. Collaborators get named closures (`answer`, `sendReceipt`, `report`), the session itself only where the public contract needs it (`sms.session`, the hook's argument). |
+| `handled-messages.ts` | `HandledMessages`: the whole hold. `offer()` creates the `Sms`, runs the handler, releases on settle; answers a failed handler's message; `refuses()` with its hysteresis; `idle()` for the drain; `sweep()`/`clear()`. |
+| `incoming-requests.ts` | Routing only, and synchronous but for the hook and `close()`: hook, link-gone check, bind gate, one switch, reassembly, offer. |
+| `outgoing-requests.ts` | `request()`: one path (refusal, link wait, slot, attempt, retry where nothing reached the socket), with `pastDrain` for a receipt. `requestOnLink()`: a bind or unbind, straight onto the attached socket, outside the window. Binds are routed there by `request()` so `send()` still carries a hand-wired bind. |
+| `sms.ts` | `sendResp()` marks the message answered synchronously before writing, so a second call and a failed handler both find it; a failed write leaves the id alone. `SmsHandlers` is `answer`/`answered`/`lostLink`/`send`, all explicit. |
+| `bind-direction.ts` | Bind types, commands (derived from one list), `LinkEnd`, `standsInFor`, `bindCarries`, `checkedBind`; out of `session-options.ts`, which keeps options and their checks. |
+| `defaults.ts` | Every default, one object; `client.ts`, `server.ts`, `reassembly.ts`, `reconnect-loop.ts` read it. |
+
+The other modules are unchanged in shape; `DlrMerger.close()` is `spend()`.
+
+## Deleted
+
+`held-messages.ts` (`HeldMessages`, `MessageHold`, the `WeakMap`, the `setImmediate`, the listener
+count, the re-used-sequence-number exit); `Session[captureRejectionSymbol]`'s routing into the hold;
+`IncomingRequests.listenerRejected()` and its `refusing` flag; `OutgoingRequests.requestPastDrain()`
+and `requestOnCurrentLink()` (folded into `request()` + `requestOnLink()`); `Session.end()`,
+`stop()`, `emitClose()`, `teardown()`, `onClose()`, `attach()`, `resetTimers()`, `answering()` and
+the `shutdownTimeout: 0` fallback; three `defaults` objects, `backoffDefaults`, `defaultMaxOctets`,
+`defaultSystemId`'s re-export; `checkHooks()` in `server.ts` (now in `checkSessionOptions()`); the
+`sms` key of `SessionEvents`; `EncodingName` `'ASCII'`. Also fixed on the way: a rebind answered on
+a link the peer had dropped no longer opens the dead socket (todo item, now in `comeBackUp()`).
+
+## Tests
+
+`docker compose run --rm node npm test`: lint and typecheck clean, **520 tests, 520 pass, 0 fail**
+(baseline 515). Every `session.on('sms', …)` in the suite became an `onSms` option (a test helper
+`inbox()` hands the first message to the test and holds nothing; `holding()` and `stuck` hold). None
+weakened on the wire or on goal 2. Changed beyond that mechanical move:
+
+- `graceful shutdown`: `submitInFlight()` holds through the handler and returns `release()`. Replaced
+ `gives up on a message the application never answers` → `…a handler that never returns`;
+ `falls back to responseTimeout …` → `waits out a handler for as long as it takes at shutdownTimeout 0`
+ (the fallback is gone); `a listener that rejected before answering …` → `a handler that rejected before
+ answering has the message refused for it, and holds nothing` (asserts the peer's `ESME_RTHROTTLED`);
+ `waits for the listener still working when another one rejected` (no second listener exists) → three
+ tests: `close() keeps waiting for a handler that answered and is still working`, `close() stops waiting
+ once the handler returns, answered or not`, `a handler that rejected after answering leaves that answer alone`;
+ `a message no listener took …` → `… no handler takes is refused at once …`; `leaves a message the library
+ refused to answer unanswered` → `keeps waiting on a handler whose answer the library refused`;
+ tests that relied on "no listener" to keep a request in flight use a `stuck` handler. Messages read
+ `still being handled` where they read `unanswered`.
+- `held message bounds` → `handled message bounds`: same bounds, exits are now handler return, sweep,
+ clear; added `answers for a handler that failed before answering, and reports it` and `answers for a
+ message no handler takes`.
+- `sendResp()`: added `answers once, and keeps the id it was given only once that answer is on the wire`.
+- `refuses to answer a message whose link went, held or already answered`: the already-answered
+ message is now refused as `already answered` (the held one still as link gone).
+- `turns a throwing sms listener …` and siblings, and `session-error.test.ts`: listener → handler.
+- `encodings`/`message`/`unsendable`/`message-class`/`declared-alphabet`: `'ASCII'` → `'GSM'`.
+- `readme.test.ts` follows the README examples; `incomingOn()` passes `answer` and `sendReceipt`.
diff --git a/MIGRATION-NOTES.md b/MIGRATION-NOTES.md
new file mode 100644
index 0000000..3c2968a
--- /dev/null
+++ b/MIGRATION-NOTES.md
@@ -0,0 +1,29 @@
+# Breaking changes for a 0.5.0 user
+
+Each line names the change and its one-line replacement.
+
+- **The `sms` event is gone; inbound messages go to an `onSms` handler option.**
+ `session.on('sms', async sms => { … })` → `client({ onSms: async sms => { … } })`.
+- **A server's handler is a server option, not a per-session listener.**
+ `smpp.on('session', s => s.on('sms', h))` → `server({ onSms: h })`; `sms.session` is the session.
+- **A hand-wired `Session` takes the handler the same way.**
+ `session.on('sms', h)` → `new Session({ onSms: h, sock })`.
+- **The hold on a message is the handler's promise, not `sendResp()`.** `close()` waits until your
+ handler returns; do the answer and the receipt inside it. A handler that returns before answering
+ releases the hold; the message then waits on you, not on the shutdown.
+- **A handler that throws or rejects before answering has the message refused for it**
+ (`ESME_RTHROTTLED` on a submission, `ESME_RX_T_APPN` on a delivery), so the peer retries. Catch
+ what you want answered otherwise. Rejected listeners used to leave the message unanswered.
+- **No `onSms` means every inbound message is refused that way and reported on `sessionError`.** A
+ transceiver that only sends sees `dlr` events as before; give it `onSms` only if it receives.
+- **`sendDlr()` no longer needs to follow `sendResp()` in the same turn to pass a shutdown.** Send it
+ wherever in the handler you like; nothing to replace.
+- **`shutdownTimeout: 0` waits for a handler as long as it takes**, where it used to give the
+ messages `responseTimeout`. Set `shutdownTimeout` if you want a bound.
+- **`sendResp()` twice on one message is refused** with `This message was already answered`.
+ Answer once.
+- **`sendResp()` that fails leaves `sms.smsId` unchanged.** Read `sms.smsId` after a successful
+ answer, not before.
+- **`encoding: 'ASCII'` is `encoding: 'GSM'`**, and so are `encodings.GSM`, `dataCodingByEncoding.GSM`
+ and what `detect()` and `encodingByDataCoding()` return. `'ASCII'` is refused by name.
+- **`SessionEvents` has no `sms` key**; a listener typed against it moves to `OnSms`.
diff --git a/MIGRATION.md b/MIGRATION.md
index 0296ede..29f9d76 100644
--- a/MIGRATION.md
+++ b/MIGRATION.md
@@ -11,6 +11,8 @@ shape is the same, connect, send, listen for delivery reports, with callbacks re
`close()`, or the socket outlives the call.
- **`server()` resolves once, when it is listening**, with a handle carrying `close()`, `port` and
a `session` event. It no longer calls back once per connection.
+- **The `sms` event is the `onSms` option**, on `client()` and `server()` alike:
+ `onSms: async sms => { await sms.sendResp(); }`. The message is held until the handler returns.
- **The id a message is answered with goes to `sendResp({ smsId })`.** `sms.smsId` is read-only: the
id the segments were answered with, the id `sendResp()` was given, or the generated UUID v7.
Assigning to it throws a `TypeError` in strict-mode code (every ES module, and any file under
@@ -85,7 +87,7 @@ have for these:
directions. 0.4.0 kept only the last one it read.
- A body carried in the `message_payload` TLV was ignored, so the message arrived empty, and a
`data_sm` was answered `ESME_RINVCMDID`, so a receipt thrown on one was lost silently. Both reach
- the application now: a receipt as `dlr`, answered for you, and a message as `sms` for you to answer.
+ the application now: a receipt as `dlr`, answered for you, and a message to `onSms` for you to answer.
- A long message segmented by the `sar_msg_ref_num`, `sar_total_segments` and `sar_segment_seqnum`
TLVs rather than a user data header was never reassembled, so each segment arrived as its own
message. Both spellings reassemble now.
diff --git a/README.md b/README.md
index 9ff2b7d..2cbbe73 100644
--- a/README.md
+++ b/README.md
@@ -10,7 +10,8 @@ window, long messages and delivery receipts. TypeScript, ESM, no dependencies.
- **Send window.** 10 requests in flight; further sends queue instead of overrunning the SMSC.
- **Long messages.** Split on send, reassembled on receive, in both the UDH and `sar_*` spellings.
- **Delivery receipts.** Read from TLVs or from receipt text, matched to the ids you were given.
-- **Graceful shutdown.** `close()` waits for what is in flight, so neither end has to guess.
+- **Graceful shutdown.** `close()` waits for your message handlers and for what is in flight, so
+ neither end has to guess.
- **Never throws.** Every fallible call resolves to `{ err?, … }`.
- **Interoperable.** Tested as a client against Jasmin and SMPPSim, and as a server against Kannel,
jsmpp, Cloudhopper, python-smpplib and php-smpp:
@@ -92,34 +93,37 @@ receipts into one, and SMSCs that write ids in two notations: [Delivery receipts
## Receive SMS
-A `receiver` or `transceiver` client gets mobile-originated messages as `sms` events:
+A `receiver` or `transceiver` client hands mobile-originated messages to its `onSms` handler:
```javascript
-session.on('sms', async sms => {
- // sms.from, sms.to, sms.message
- await sms.sendResp();
+const { err, session } = await client({
+ onSms: async sms => {
+ // sms.from, sms.to, sms.message
+ await sms.sendResp();
+ },
});
```
-Call `sendResp()` for every message, multipart included: until you do, it counts toward the bound
-past which the peer's messages are refused. Delivery receipts reach you as `dlr` events, not here. A
-multipart message arrives reassembled and already answered segment by segment, so `sendResp()` there
-puts nothing on the wire and releases it: [Receiving in depth](#receiving-in-depth).
+Call `sendResp()` once for every message, multipart included; a message you never answer is one
+the peer waits on. The message is held for as long as your handler runs: it counts toward the bound
+past which the peer's messages are refused, and `close()` waits for the handler to return. A handler
+that throws or rejects before answering has the message refused for it, so the peer retries; without
+an `onSms` at all, every message is refused that way. Delivery receipts reach you as `dlr` events,
+not here. A multipart message arrives reassembled and already answered segment by segment, so
+`sendResp()` there puts nothing on the wire: [Receiving in depth](#receiving-in-depth).
## Run an SMPP server
```javascript
import { server } from '@larvit/smpp';
-const { err, server: smpp } = await server();
-if (err) throw err;
-
-smpp.on('session', session => {
- session.on('sms', async sms => {
- // sms.from, sms.to, sms.message, sms.dlr
+const { err, server: smpp } = await server({
+ onSms: async sms => {
+ // sms.from, sms.to, sms.message, sms.dlr, sms.session
await sms.sendResp();
- });
+ },
});
+if (err) throw err;
```
With authentication and delivery reports:
@@ -134,13 +138,10 @@ const { err, server: smpp } = await server({
return { userData: { userId: 123 } };
},
-});
-if (err) throw err;
-
-smpp.on('session', session => {
- session.on('sms', async sms => {
+ onSms: async sms => {
+ // sms.session.userData is what authenticate returned for this peer
if (sms.answeredOnArrival) {
- await sms.sendResp(); // multipart: already answered per segment; this only releases the shutdown drain
+ await sms.sendResp(); // multipart: already answered per segment, so this writes nothing
} else {
// no args: ESME_ROK + generated id; or sendResp({ smsId, status: 'ESME_RMSGQFUL' })
await sms.sendResp();
@@ -149,19 +150,22 @@ smpp.on('session', session => {
if (sms.dlr) {
await sms.sendDlr(); // same as sms.sendDlr('DELIVERED')
}
- });
+ },
});
+if (err) throw err;
console.log(smpp.port); // the port actually bound, useful when 0 was requested
await smpp.close(); // stop listening, then drain and close every live session
```
- `sendResp()` answers `ESME_ROK` with a generated UUID v7 as the message id.
- `sendResp({ smsId, status })` names the id or refuses the message.
+ `sendResp({ smsId, status })` names the id or refuses the message. A second call is refused.
- `sms.dlr` is true where the sender asked for a receipt. `sendDlr()` reports `DELIVERED`,
`sendDlr('UNDELIVERABLE')` any other state: [Server in depth](#server-in-depth).
- A message that arrived in several segments was answered as they arrived, so `sendResp()` there
takes no `smsId` or refusing `status`. `sms.answeredOnArrival` says which case you are in.
+- `smpp.close()` waits for every running `onSms` before it closes a session, so answer and send the
+ receipt inside the handler, as above, and nothing is cut off.
## Errors
@@ -183,7 +187,7 @@ Neither is named `error`, because Node throws on an unhandled `error` event.
| --- | --- | --- |
| A PDU the peer sent that the codec could not read. The link stays up; only that PDU is lost. | `PduRefusedError` | Count it. |
| A concatenated message given up on before it was whole. Its segments were answered, so the peer will not resend, and no `sms` fired for it. | `Error` | Count it as lost traffic. |
-| The session or socket failing, or a hook or listener that threw or rejected. | `Error` | Alert. |
+| The session or socket failing, or a hook, handler or listener that threw or rejected. | `Error` | Alert. |
The last two are told apart by message text only, so this alerts on both:
@@ -231,8 +235,9 @@ All optional. Timeouts and delays are milliseconds.
| `enquireLinkInterval` | `20000` | Interval between `enquire_link` on a quiet link. |
| `idleTimeout` | `2 × enquireLinkInterval` | Give up on a link the peer has stopped answering, and re-bind unless `reconnect` is `false`. |
| `responseTimeout` | `30000` | How long to wait for a response, and how long a send with no link waits for the next one. `0` waits forever. |
-| `shutdownTimeout` | `5000` | How long `close()` and `unbind()` wait for requests already sent and messages not yet answered. `0` waits forever for the requests, which end when the peer answers or `responseTimeout` expires, so both at `0` never ends. The messages then fall back to `responseTimeout`, or to its default where that is `0` too. |
+| `shutdownTimeout` | `5000` | How long `close()` and `unbind()` wait for `onSms` handlers still running and for requests already sent. `0` waits forever: a request ends when the peer answers or `responseTimeout` expires, so both at `0` never ends, and a handler ends when it returns or after the five minutes past which one is no longer counted. |
| `maxOutstanding` | `10` | Requests on the wire at once; further sends queue. |
+| `onSms` | refuse | `(sms) => Promise \| void`. Takes every inbound message: [Receive SMS](#receive-sms). |
| `smsIdFormat` | — | The notation the SMSC writes message ids in, per place: `{ receipt: 'decimal', submitResp: 'hex' }`. Only where the two disagree: [Delivery receipts](#delivery-receipts). |
| `reconnect` | on | `{ minDelay, maxDelay }` retunes the backoff; `false` turns it off, so a drop ends the session; `{ fromStart: true }` retries the first connect too. |
| `log` | silent | Any object with `debug`, `error`, `info`, `verbose` and `warn` methods: [Logging](#logging). |
@@ -257,6 +262,7 @@ All optional. Timeouts are milliseconds.
| `host`, `port` | all interfaces, `2775` | Where to listen. `port: 0` takes any free port; `smpp.port` says which. |
| `authenticate` | accept everything | `({ password, session, systemId, systemType }) => false \| { userData }`, sync or async. |
| `onRequest` | none | `(session, pduObj) => true \| false`, sync or async. First refusal on every request a bound peer sends: [Server in depth](#server-in-depth). |
+| `onSms` | refuse | `(sms) => Promise \| void`. Takes every message a bound peer submits, on every session; `sms.session` says which: [Run an SMPP server](#run-an-smpp-server). |
| `systemId` | `''` | The SMSC identity returned in the bind response. |
| `interfaceVersion` | `0x34` | The SMPP version advertised in the bind response. Optional parameters are sent to a peer from `0x34` up, whatever this is set to. |
| `tls` | `false` | A `tls.TlsOptions` object with your certificate and key. A bare `true` is refused. |
@@ -298,11 +304,11 @@ refuse a non-ASCII sender of its own accord, which reaches you as a refusal such
| `encoding` | Alphabet | Characters per SMS | Per segment of a long message |
| --- | --- | --- | --- |
-| `ASCII` | GSM 03.38 7-bit | 160 | 153 |
+| `GSM` | GSM 03.38 7-bit | 160 | 153 |
| `LATIN1` | ISO 8859-1 | 140 | 134 |
| `UCS2` | UCS-2 | 70 | 67 |
-- Omitted: `ASCII` where the message fits GSM 7-bit, otherwise `UCS2`. `LATIN1` only when named.
+- Omitted: `GSM` where the message fits GSM 7-bit, otherwise `UCS2`. `LATIN1` only when named.
Any other name is refused.
- GSM extension characters (`{}[]\~^|€` and form feed) count as two, as does a character outside
the basic multilingual plane in `UCS2`.
@@ -357,9 +363,11 @@ holds for `session.send()`.
### Events
+A request you must answer is a handler option: `onSms`, and `onRequest` on a server. A fact you may
+watch is an event, and nothing waits on its listeners.
+
| Event | Fires when |
| --- | --- |
-| `sms` | An SMS arrives, reassembled if it was multipart. Carries `sendResp()`, `sendDlr()` and `smsId`. |
| `dlr` | A delivery report arrives, one per segment, with its PDU as the second argument: [Delivery receipts](#delivery-receipts). |
| `messageDlr` | Every segment of a long message sent with `dlr: true` has a final report: [Delivery receipts](#delivery-receipts). |
| `close` | The session is over and nothing will bring the link back. Fires once, whether you closed it or the link failed for good. |
@@ -376,17 +384,14 @@ holds for `session.send()`.
**Shutdown.** `close()` and `unbind()` both:
-1. Refuse further sends. `sendDlr()` is the one send let past, when issued straight after
- `sendResp()`; await anything in between and it races the shutdown like any other send.
-2. Wait up to `shutdownTimeout` for the requests already sent, and for every `sms` the application
- has not answered. That wait ends when `sendResp()` puts the response on the wire (or, for a
- message answered on arrival, when it is called at all), or when every listener that took the
- message has failed. Answering through `sendReturn()` instead leaves the wait running.
+1. Refuse further `sendSms()` and `send()`. `sendResp()` and `sendDlr()` still go out, because they
+ finish a message the shutdown is waiting on.
+2. Wait up to `shutdownTimeout` for every `onSms` handler still running, then for the requests
+ already sent.
3. Tear down what is left, resolving to an `err` that says what was lost.
-A message left unanswered for five minutes is no longer waited for.
-`close({ signal })` cuts the wait short. `unbind()` takes no signal, and waits a further
-`responseTimeout` for its own response.
+A handler still running after five minutes is no longer waited for. `close({ signal })` cuts the wait
+short. `unbind()` takes no signal, and waits a further `responseTimeout` for its own response.
**Sends and the link.**
@@ -442,12 +447,15 @@ const { err, pduObj } = await session.send({
each is two messages.
- **Answered on arrival.** Each segment was answered as it landed, before you see the message:
[Server in depth](#server-in-depth).
-- **Unanswered messages.** While a session holds 1000 messages you have not called `sendResp()` on,
- or 64 MiB of them counted the way `maxOctets` counts segments, every new message and segment is
- refused so the peer retries it: `ESME_RTHROTTLED` on a submission, `ESME_RX_T_APPN` on a delivery.
- No `sms` fires. Reaching the bound logs one `warn`, and the first message accepted once both are
- down to half one `info`. A message left five minutes is dropped from the count with a `warn`; a
- later `sendResp()` still answers it. None of the three is an option.
+- **Messages being handled.** While 1000 `onSms` handlers are running on a session, or the messages
+ they hold weigh 64 MiB counted the way `maxOctets` counts segments, every new message and segment
+ is refused so the peer retries it: `ESME_RTHROTTLED` on a submission, `ESME_RX_T_APPN` on a
+ delivery. The handler is not called. Reaching the bound logs one `warn`, and the first message
+ accepted once both are down to half one `info`. A handler running five minutes is dropped from the
+ count with a `warn`; its `sendResp()` still answers. None of the three is an option.
+- **A handler that fails.** One that throws or rejects reaches `sessionError`. Where it had not
+ answered, the message is refused with the same retry status, so the peer sends it again; where it
+ had, that answer stands. The same happens to every message where no `onSms` was given.
- **Where the body is.** A body in the `message_payload` TLV, SMPP's way of carrying up to 64 KB and
the only place a `data_sm` has, reads exactly like one in `short_message`, concatenated messages
and receipts included. A PDU filling both is read from `short_message`.
@@ -517,7 +525,7 @@ cannot, since a peer may number a message one part of one.
**Refusing a request** for a reason in the request rather than the message (a full queue, an unknown
recipient, an unauthorised sender) has to land before a segment is answered. `onRequest` runs on
-every request a bound peer sends, before reassembly and before the `sms` event:
+every request a bound peer sends, before reassembly and before `onSms`:
```javascript
import { isCommand, server } from '@larvit/smpp';
@@ -668,7 +676,7 @@ if (isCommand(pduObj, 'submit_sm')) {
| Receipts | `dlrFromPdu`, `parseReceipt`, `receiptCodes` |
| Time and ids | `smppDate`, `smppTime`, `uuidv7` |
| Spec tables | `cmds`, `consts`, `encodings`, `errors`, `tlvs`, `types`, the `cmdsById`, `constsById`, `errorsById` and `tlvsById` maps, and all of them grouped as `defs`. `isCommandName`, `isErrorName`, `isEncodingName`, `isTlvName`, `commandNameById` and `errorNameById` narrow a value into them. |
-| Types | Every option, result, event payload and table entry has a named type: `ClientOptions`, `ServerOptions`, `SendSmsOptions`, `SendSmsResult`, `Sms`, `Dlr`, `MessageDlr`, `Receipt`, `PduObject`, `PduHeader`, `SmppLog`, `Result` and the rest in `dist/index.d.ts`. |
+| Types | Every option, hook, result, event payload and table entry has a named type: `ClientOptions`, `ServerOptions`, `OnSms`, `SendSmsOptions`, `SendSmsResult`, `Sms`, `Dlr`, `MessageDlr`, `Receipt`, `PduObject`, `PduHeader`, `SmppLog`, `Result` and the rest in `dist/index.d.ts`. |
## What changed per release
diff --git a/benchmarks/smsc-sink.ts b/benchmarks/smsc-sink.ts
index 324c11e..5779ba0 100644
--- a/benchmarks/smsc-sink.ts
+++ b/benchmarks/smsc-sink.ts
@@ -5,22 +5,20 @@ import { server } from '../src/server.ts';
* library's own ceiling rather than an SMSC's storage. Prints the bound port on stdout, then waits.
*/
const port = Number(process.env.PORT ?? 0);
-const { err, server: smpp } = await server({ port });
+let answered = 0;
+const { err, server: smpp } = await server({
+ onSms: async sms => {
+ answered++;
+ await sms.sendResp();
+ },
+ port,
+});
if (err) {
process.stderr.write(`sink failed to listen: ${err.message}\n`);
process.exit(1);
}
-let answered = 0;
-
-smpp.on('session', session => {
- session.on('sms', async sms => {
- answered++;
- await sms.sendResp();
- });
-});
-
smpp.on('serverError', reason => {
process.stderr.write(`sink serverError: ${reason.message}\n`);
});
diff --git a/docs/decisions.md b/docs/decisions.md
index c245097..c34cb05 100644
--- a/docs/decisions.md
+++ b/docs/decisions.md
@@ -6,6 +6,24 @@ rule and an index of the titles below.
## The public surface
+- **A request the application must answer is a handler option; a fact it may watch is an event.**
+ Maintainer's call, 2026-09-29, from the comprehension panel's round two: every seat ranked the
+ `sms` event's hold hardest — released one turn after `sendResp()`, or when every listener had
+ rejected against a count taken at emit, or when no listener took it, or on a re-used sequence
+ number — and the internals could not be made simpler than the contract demanded. `onSms` on
+ `client()`, `server()` and `Session` takes each message, and the returned promise is the hold:
+ the message counts toward the bound and `close()` waits, until the handler settles. Goal 2 settles
+ the failure case: a handler that throws before answering has decided nothing, so the message is
+ refused with the retry status and the peer sends it again; one that threw after answering keeps
+ its answer, since a second response is goal 1's wire violation. No `onSms` at all takes the same
+ path, per message, so a receiver bound with nothing to receive into is visible rather than silent.
+ `dlr`, `messageDlr` and the life events stay events, since nothing waits on their listeners.
+ Rejected: keeping the event and waiting on the listener's own promise, which needs `listeners()`
+ re-declared and a count of them at emit — the contract that made the hold six exits. Rejected: a
+ handler that returns the answer for the library to write, which cannot express answering first and
+ working after, nor a receipt that must follow the answer inside the same hold. Rejected: a
+ settable `session.onSms`, a second spelling of the option.
+
- **`Session` is publicly constructible, which is what makes `SessionOptions` and `ReconnectOptions`
public too.** Raised twice as a leak; it is not one. The collaborators `session.ts` delegates to
stay unpublished so they can be reshaped.
@@ -17,17 +35,16 @@ rule and an index of the titles below.
three the library dispatches by it, of which it sends the first two.
- **Both emitters re-declare their listener methods to accept a promise.** Maintainer's call,
- 2026-08-27: `EventEmitter` types every listener as void-returning, so the
- `session.on('sms', async sms => …)` README documents reads as a misused promise in any strict
- consumer. `declare on: …` and its six siblings re-type the inherited methods to return `unknown`,
+ 2026-08-27: `EventEmitter` types every listener as void-returning, so an
+ `session.on('dlr', async dlr => …)` reads as a misused promise in any strict consumer. `declare on: …` and its six siblings re-type the inherited methods to return `unknown`,
which emits nothing and needs no cast; overriding them as real methods cannot work, because the
`super.on()` call needs one. The cost is that a subclass can no longer reach those seven through
`super` — re-declaring them the same way is its way out. `unknown` rather than
`void | Promise` because a listener may return anything: `session.on('close', () =>
- set.delete(session))` returns a boolean. This also settles what the drain can wait on: a listener's
- own promise would be the better completion signal, and reaching it needs `listeners()`, which
- cannot be re-declared the same way — Node types it invariantly enough that widening `void` to
- `unknown` is `TS2416`. Re-probed 2026-09-01; `sendResp()` stays the signal.
+ set.delete(session))` returns a boolean. A listener's own promise cannot be waited on, since
+ reaching it needs `listeners()`, which cannot be re-declared the same way — Node types it
+ invariantly enough that widening `void` to `unknown` is `TS2416` — which is why a message goes to
+ a handler option rather than an event.
- **`PduRefusedError` is exported, and `sessionError` names it in the event's type.** Maintainer's
call, 2026-09-05, from a product review: one event carries both a PDU the peer malformed and the
@@ -225,8 +242,8 @@ rule and an index of the titles below.
carried where bit 4 says so in every group below 0x80 and always in the 0xF0 group, and
`encodingByDataCoding()` reads the alphabet off that same test rather than repeating the group
masks beside it. It is exported for the reason `concatOf()` is — an application that needs a class
- other than 0 would otherwise rewrite the read this fixed. Rejected: a `messageClass` field on the
- `sms` event, which pays goal 8 for three classes nothing here acts on, where the boolean the
+ other than 0 would otherwise rewrite the read this fixed. Rejected: a `messageClass` field on
+ `Sms`, which pays goal 8 for three classes nothing here acts on, where the boolean the
application already had covers the one it does. Compressed text is out of scope and stays out —
nothing here implements 3GPP TS 23.042, so a compressed body reaches the application as whatever
its declared alphabet makes of it — but bit 5 does not move the class bits, so 0x30 is read as
@@ -583,8 +600,8 @@ rule and an index of the titles below.
one round trip rather than one per segment. Rejected: sending each segment once the last is
answered, which a receiver waiting for the whole message before answering would deadlock.
-- **Every segment of a concatenated message is answered as it arrives, so `sendResp()` on one is the
- application's own signal rather than the peer's answer.** Maintainer's call, 2026-09-06, from the
+- **Every segment of a concatenated message is answered as it arrives, so `sendResp()` on one writes
+ nothing.** Maintainer's call, 2026-09-06, from the
Jasmin interoperability phase: Jasmin dispatches one `submit_sm` per connector at a time and will
not send segment 2 until segment 1 is answered, so holding a group unanswered until it was whole
deadlocked every multi-segment message against a production gateway
@@ -596,7 +613,7 @@ rule and an index of the titles below.
no-op. `answeredOnArrival` is on `Sms` because nothing the application can compute says it, and the
discriminant a reader would reach for instead is wrong. A message `sendResp()` still answers itself is
untouched, and is where a caller-chosen id and a refusal live; `onRequest` is the escape hatch for
- an application that must refuse a PDU the `sms` event could not have shown it yet. `collect()`
+ an application that must refuse a PDU `onSms` could not have shown it yet. `collect()`
answers every segment it will not carry rather than leaving it unanswered, which is the same stall
in miniature: the field that numbered it where the segment belongs to no group, the retry status
where the segment's own arrival overran the octet cap, since a peer told that still holds it. Rejected:
@@ -613,7 +630,7 @@ rule and an index of the titles below.
distinguishable. Accepted: a completing segment whose own answer the socket would not carry still
reaches the application, because the message is whole and correct and the failed answer is on
`sessionError` — a peer that re-sends after the drop is the smaller risk than dropping a message
- in hand. The answer goes out before the `sms` event either way, so a listener's own receipt can
+ in hand. The answer goes out before `onSms` is called either way, so a handler's own receipt can
never precede the acceptance of the message it reports on.
- **`server()` composes the application's `onRequest` after its own bind handling, and offers it
@@ -653,36 +670,39 @@ rule and an index of the titles below.
fall-through could be gated on it, which buys a fail-open path with state and an internal contract
no other collaborator needs.
-- **The drain waits on the messages the application holds, and `sendResp()` is what says it is done
- with one.** Maintainer's call, 2026-09-01: waiting on the send window alone tore a server session
- down while the application was still answering a `submit_sm`, so the peer timed out and re-sent —
- the duplicate goal 2 forbids, in the direction the window already covers. No completion signal was
- added to the `sms` event: `sendResp()` is what an application already calls when it is done with a
- message, so it is the one the drain waits for. Counting every inbound request until `sendReturn()`
- answered it was rejected: an `onRequest` that deliberately answers nothing would then cost a full
- `shutdownTimeout` on every close. The response reaching the wire ends the wait, so a `sendResp()`
- the library refused or the socket would not carry leaves `close()` still reporting the message the
- peer is owed.
-
-- **The drain's wait on the application ignores `shutdownTimeout: 0`.** Waiting forever is safe for
- 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
- 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
- refused.** A bound the application cannot raise is the point: an application that answers nothing
- would otherwise hold ever more messages, which goal 4 forbids. Reassembly's `maxOctets` is an
- option because it bounds what the peer sends; this bounds what the application leaves unanswered.
+- **The drain waits on the `onSms` handlers still running, and a handler's promise is what says it
+ is done with a message.** Maintainer's call, 2026-09-01, re-settled 2026-09-29 with the handler
+ option: waiting on the send window alone tore a server session down while the application was
+ still answering a `submit_sm`, so the peer timed out and re-sent — the duplicate goal 2 forbids.
+ The handler returning is the one signal, so a receipt sent inside it goes out before the drain
+ ends and `shutdownTimeout: 0` waits for it as it waits for a request, bounded by the five-minute
+ deadline past which a handler is no longer counted. Rejected: `sendResp()` as the signal, which
+ had to let a receipt sent one turn later past the drain and could not see a handler still working.
+ Rejected: counting every inbound request until `sendReturn()` answered it — an `onRequest` that
+ deliberately answers nothing would then cost a full `shutdownTimeout` on every close. Rejected: a
+ `responseTimeout` fallback for a `shutdownTimeout` of 0, a second bound where the deadline already
+ is one.
+
+- **A handler that fails before answering has the message refused for it; one that fails after has
+ its answer stand.** Goal 2: a handler that threw decided nothing, so the peer keeps the message and
+ retries on `ESME_RTHROTTLED` or `ESME_RX_T_APPN`, the same answer the bound gives; after an answer
+ nothing more is written, since a second response is goal 1's wire violation. `sendResp()` marks
+ the message answered before its write and refuses a second call, which is what makes the two
+ cases separable. Rejected: leaving the message unanswered, which is what the `sms` event did and
+ what stalls a peer that dispatches one request at a time.
+
+- **What the running handlers hold is capped on constants, and a message past the cap is
+ refused.** A bound the application cannot raise is the point: a handler that never returns would
+ otherwise hold ever more messages, which goal 4 forbids. Reassembly's `maxOctets` is an
+ option because it bounds what the peer sends; this bounds what the application is still handling.
Maintainer's call, 2026-09-26. Refusing leaves the message with the peer, which will send it again
- (goal 2). Rejected: dropping the oldest to make room, which frees nothing while the application
- still holds its `Sms`, and stops the drain waiting for a message the peer is owed. Rejected:
+ (goal 2). Rejected: dropping the oldest to make room, which frees nothing while the handler still
+ holds its `Sms`, and stops the drain waiting for a message the peer is owed. Rejected:
pausing the socket, which also stalls every answer and `enquire_link` on the link. Reaching the
bound shows only in the log (goal 8): an event or a public count would be surface for what the
- application already knows, since it is the one not answering. A message held past its timeout is
- still dropped, so `close()` can report fewer unanswered than there were — accepted, because the
- alternative is holding what nothing will answer.
+ application already knows, since it is the one not returning. A handler past its timeout is still
+ dropped from the count, so `close()` can report fewer than there were — accepted, because the
+ alternative is holding what nothing will finish.
- **A store at its bound answers `ESME_RTHROTTLED` to a submission and `ESME_RX_T_APPN` to a
delivery, a `data_sm` by whichever it stands in for.** Maintainer's call, 2026-09-26, for
diff --git a/interop-tests/cloudhopper.test.ts b/interop-tests/cloudhopper.test.ts
index dca5b5b..2ddda53 100644
--- a/interop-tests/cloudhopper.test.ts
+++ b/interop-tests/cloudhopper.test.ts
@@ -41,23 +41,20 @@ async function driver(path: string, params: Record = {}): Promis
const manualTexts = new Set();
const allSms: { session: Session; sms: Sms }[] = [];
-function attach(session: Session): void {
- session.on('sms', sms => {
- allSms.push({ session, sms });
+function onSms(sms: Sms): void {
+ allSms.push({ session: sms.session, sms });
- if (manualTexts.has(sms.message)) return;
+ if (manualTexts.has(sms.message)) return;
- // The slow server this phase's window scenarios need: every ordinary submit is held for
- // SLOW_DELAY_MS before being answered, so a burst genuinely presses on a small window.
- void delay(SLOW_DELAY_MS).then(() => sms.sendResp());
- });
+ // The slow server this phase's window scenarios need: every ordinary submit is held for
+ // SLOW_DELAY_MS before being answered, so a burst genuinely presses on a small window.
+ void delay(SLOW_DELAY_MS).then(() => sms.sendResp());
}
-const { err: serverErr, server: smpp } = await server({ authenticate: () => true, idleTimeout: 40_000, port: SMPP_PORT });
+const { err: serverErr, server: smpp } = await server({ authenticate: () => true, idleTimeout: 40_000, onSms, port: SMPP_PORT });
assert.equal(serverErr, undefined);
assert.ok(smpp);
-smpp.on('session', attach);
const key = readFileSync('/shared-certs/server.key');
const cert = readFileSync('/shared-certs/server.crt');
@@ -67,13 +64,13 @@ const cert = readFileSync('/shared-certs/server.crt');
const { err: tlsServerErr, server: tlsSmpp } = await server({
authenticate: () => true,
idleTimeout: 40_000,
+ onSms,
port: TLS_PORT,
tls: { cert, key, maxVersion: 'TLSv1.2' },
});
assert.equal(tlsServerErr, undefined);
assert.ok(tlsSmpp);
-tlsSmpp.on('session', attach);
after(async () => {
await smpp.close();
diff --git a/interop-tests/dumbclient.test.ts b/interop-tests/dumbclient.test.ts
index 81bc9ff..61b5d80 100644
--- a/interop-tests/dumbclient.test.ts
+++ b/interop-tests/dumbclient.test.ts
@@ -105,6 +105,23 @@ const { err, server: smpp } = await server({
authenticate: ({ systemId }) => ({ userData: { systemId } satisfies ScenarioUserData }),
idleTimeout: 40_000,
log,
+ onSms: sms => {
+ const { session } = sms;
+ const name = scenarioOf(session);
+
+ sessionByScenario.set(name, session);
+
+ const s = statsFor(name);
+ const arrivalIndex = s.arrived;
+
+ s.arrived++;
+ if (s.ids.has(sms.smsId)) s.duplicateIds++;
+ else s.ids.add(sms.smsId);
+ s.peakOutstanding = Math.max(s.peakOutstanding, s.arrived - s.answered);
+
+ if (name === 'dumb-w500' || name === 'dumb-w2000') slowRespond(session, sms, arrivalIndex);
+ else fastRespond(session, sms, arrivalIndex);
+ },
port: SMPP_PORT,
});
@@ -141,23 +158,6 @@ function fastRespond(session: Session, sms: Sms, arrivalIndex: number): void {
smpp.on('session', session => {
// Attached now, not lazily in a test body - see the comment on ScenarioStats.closed.
session.on('close', () => { statsFor(scenarioOf(session)).closed = true; });
-
- session.on('sms', sms => {
- const name = scenarioOf(session);
-
- sessionByScenario.set(name, session);
-
- const s = statsFor(name);
- const arrivalIndex = s.arrived;
-
- s.arrived++;
- if (s.ids.has(sms.smsId)) s.duplicateIds++;
- else s.ids.add(sms.smsId);
- s.peakOutstanding = Math.max(s.peakOutstanding, s.arrived - s.answered);
-
- if (name === 'dumb-w500' || name === 'dumb-w2000') slowRespond(session, sms, arrivalIndex);
- else fastRespond(session, sms, arrivalIndex);
- });
});
function memShape(): string {
@@ -206,7 +206,7 @@ after(async () => {
// S9 (target 11) and the backpressure-at-server scenario: window 2000 at a high rate against a
// handler slowed enough to build a real backlog. window500 is the same shape with a window below
-// maxHeldMessages (1000, session-options.ts defaults.maxHeldMessages), the bound past which a
+// maxHandledMessages (1000, defaults.ts), the bound past which a
// peer's window is answered ESME_RTHROTTLED. smpp-dumb-client counts a throttled message as sent
// and never resends it, so window 2000 accounts for 20,000 as answered plus throttled.
const throttleMessage = 'session - unanswered messages at their bound, asking the peer to retry';
@@ -236,7 +236,7 @@ describe('S9 - bounded window against a slowed handler', () => {
});
}
- test('window 2000 pressed past maxHeldMessages (1000): the peer is throttled, window500 never is', () => {
+ test('window 2000 pressed past maxHandledMessages (1000): the peer is throttled, window500 never is', () => {
assert.ok(throttled('dumb-w2000') > 0, 'expected at least one ESME_RTHROTTLED under window 2000');
assert.ok(statsFor('dumb-w2000').peakOutstanding <= 1000);
assert.equal(statsFor('dumb-w500').peakOutstanding <= 500, true);
diff --git a/interop-tests/jasmin.test.ts b/interop-tests/jasmin.test.ts
index c6c9b92..f2fa41d 100644
--- a/interop-tests/jasmin.test.ts
+++ b/interop-tests/jasmin.test.ts
@@ -86,6 +86,20 @@ const { err: upstreamErr, server: upstream } = await server({
// was already in flight for Jasmin's own requeue_delay (120s default) before it retries - far past
// any per-test wait budget here - so this is generous specifically to never be the trigger.
idleTimeout: 300_000,
+ onSms: sms => {
+ const variant = (sms.session.userData as { variant?: UpstreamVariant } | undefined)?.variant;
+
+ if (variant) upstreamSms.push({ sms, variant });
+
+ void (async () => {
+ await sms.sendResp();
+
+ if (sms.dlr) {
+ await delay(150);
+ await sms.sendDlr('DELIVERED');
+ }
+ })();
+ },
port: UPSTREAM_PORT,
});
@@ -97,7 +111,7 @@ const upstreamServer = upstream;
upstreamServer.on('session', session => {
// `session` fires on raw connect, before authenticate() has run - session.userData is not set
// yet, so the map is populated off the bind PDU itself (like kannel.test.ts's bindPdus), not off
- // userData; userData is only read later, from 'sms', where authenticate() has long since run.
+ // userData; userData is only read later, in onSms, where authenticate() has long since run.
session.on('incomingPduObj', pduObj => {
if (!pduObj.cmdName.startsWith('bind_')) return;
@@ -105,21 +119,6 @@ upstreamServer.on('session', session => {
if (variant) upstreamSessions.set(variant, session);
});
-
- session.on('sms', sms => {
- const variant = (session.userData as { variant?: UpstreamVariant } | undefined)?.variant;
-
- if (variant) upstreamSms.push({ sms, variant });
-
- void (async () => {
- await sms.sendResp();
-
- if (sms.dlr) {
- await delay(150);
- await sms.sendDlr('DELIVERED');
- }
- })();
- });
});
async function waitForUpstreamSession(variant: UpstreamVariant, budget = 20_000): Promise {
@@ -227,7 +226,7 @@ async function sendUdhMo(session: Session, opts: { from: string; message: string
const multipart = segments.length > 1;
for (const segment of segments) {
- const params = submitSmParams({ from: opts.from, message: opts.message, to: opts.to }, segment, { encoding: 'ASCII', multipart });
+ const params = submitSmParams({ from: opts.from, message: opts.message, to: opts.to }, segment, { encoding: 'GSM', multipart });
const sent = await session.send({ cmdName: 'deliver_sm', params });
assert.equal(sent.err, undefined);
@@ -247,7 +246,7 @@ async function sendMessagePayloadMo(session: Session, opts: { from: string; mess
},
tlvs: {
// The body is octets under the PDU's own data_coding wherever it is carried, and 0 is GSM.
- message_payload: { tagValue: encodeMessage(opts.message, 'ASCII').buffer },
+ message_payload: { tagValue: encodeMessage(opts.message, 'GSM').buffer },
},
});
}
@@ -397,15 +396,12 @@ describe('S7 - Jasmin as the ESME against our server (HTTP send API, DLR callbac
test('a long GSM message from our server reassembles at Jasmin (or is recorded as fragments)', async () => {
const upstreamSession = await waitForUpstreamSession('main');
- const { err, session } = await bind(USERNAME, PASSWORD, { bindType: 'receiver' });
+ const sms: Sms[] = [];
+ const { err, session } = await bind(USERNAME, PASSWORD, { bindType: 'receiver', onSms: s => { sms.push(s); } });
assert.equal(err, undefined);
assert.ok(session);
- const sms: Sms[] = [];
-
- session.on('sms', s => { sms.push(s); });
-
const text = `s7-long-${'p'.repeat(300)}`;
await sendUdhMo(upstreamSession, { from: TO, message: text, to: FROM });
@@ -418,15 +414,12 @@ describe('S7 - Jasmin as the ESME against our server (HTTP send API, DLR callbac
test('a UCS-2 message with 一 and an emoji from our server (or is recorded as fragments)', async () => {
const upstreamSession = await waitForUpstreamSession('main');
- const { err, session } = await bind(USERNAME, PASSWORD, { bindType: 'receiver' });
+ const sms: Sms[] = [];
+ const { err, session } = await bind(USERNAME, PASSWORD, { bindType: 'receiver', onSms: s => { sms.push(s); } });
assert.equal(err, undefined);
assert.ok(session);
- const sms: Sms[] = [];
-
- session.on('sms', s => { sms.push(s); });
-
const text = `一😀${'q'.repeat(60)}`;
await sendSarMo(upstreamSession, { from: TO, message: text, to: FROM });
@@ -502,15 +495,12 @@ describe('C3+C7 - long MT through the fake upstream, receipts and id consistency
describe('C8 (target 3) - long MO from an upstream SMSC, SAR vs UDH segmentation', () => {
test('SAR-segmented deliver_sm from the fake upstream', async () => {
const upstreamSession = await waitForUpstreamSession('main');
- const { err, session } = await bind(USERNAME, PASSWORD, { bindType: 'receiver' });
+ const sms: Sms[] = [];
+ const { err, session } = await bind(USERNAME, PASSWORD, { bindType: 'receiver', onSms: s => { sms.push(s); } });
assert.equal(err, undefined);
assert.ok(session);
- const sms: Sms[] = [];
-
- session.on('sms', s => { sms.push(s); });
-
const text = `sar-mo-${'m'.repeat(300)}`;
await sendSarMo(upstreamSession, { from: TO, message: text, to: FROM });
@@ -526,15 +516,12 @@ describe('C8 (target 3) - long MO from an upstream SMSC, SAR vs UDH segmentation
test('UDH-segmented deliver_sm from the fake upstream', async () => {
const upstreamSession = await waitForUpstreamSession('main');
- const { err, session } = await bind(USERNAME, PASSWORD, { bindType: 'receiver' });
+ const sms: Sms[] = [];
+ const { err, session } = await bind(USERNAME, PASSWORD, { bindType: 'receiver', onSms: s => { sms.push(s); } });
assert.equal(err, undefined);
assert.ok(session);
- const sms: Sms[] = [];
-
- session.on('sms', s => { sms.push(s); });
-
const text = `udh-mo-${'n'.repeat(300)}`;
await sendUdhMo(upstreamSession, { from: TO, message: text, to: FROM });
@@ -551,15 +538,12 @@ describe('C8 (target 3) - long MO from an upstream SMSC, SAR vs UDH segmentation
describe('C8 (target 2) - message_payload with sm_length 0', () => {
test('a deliver_sm carrying message_payload instead of short_message', async () => {
const upstreamSession = await waitForUpstreamSession('main');
- const { err, session } = await bind(USERNAME, PASSWORD, { bindType: 'receiver' });
+ const sms: Sms[] = [];
+ const { err, session } = await bind(USERNAME, PASSWORD, { bindType: 'receiver', onSms: s => { sms.push(s); } });
assert.equal(err, undefined);
assert.ok(session);
- const sms: Sms[] = [];
-
- session.on('sms', s => { sms.push(s); });
-
const text = 'message-payload only, sm_length 0';
const pushed = await sendMessagePayloadMo(upstreamSession, { from: TO, message: text, to: FROM });
diff --git a/interop-tests/jsmpp.test.ts b/interop-tests/jsmpp.test.ts
index db96513..21d2c45 100644
--- a/interop-tests/jsmpp.test.ts
+++ b/interop-tests/jsmpp.test.ts
@@ -46,6 +46,10 @@ const manualTexts = new Set();
const { err: serverErr, server: smpp } = await server({
authenticate: () => true,
idleTimeout: 40_000,
+ onSms: sms => {
+ allSms.push({ session: sms.session, sms });
+ if (!manualTexts.has(sms.message)) void sms.sendResp();
+ },
port: SMPP_PORT,
});
@@ -58,10 +62,6 @@ smppServer.on('session', session => {
session.on('incomingPduObj', pduObj => {
if (pduObj.cmdName.startsWith('bind_')) bindPdus.push(pduObj.params);
});
- session.on('sms', sms => {
- allSms.push({ session, sms });
- if (!manualTexts.has(sms.message)) void sms.sendResp();
- });
session.on('sessionError', err => { allSessionErrors.push({ err, session }); });
});
diff --git a/interop-tests/kannel.test.ts b/interop-tests/kannel.test.ts
index 0d911d0..7324305 100644
--- a/interop-tests/kannel.test.ts
+++ b/interop-tests/kannel.test.ts
@@ -154,6 +154,15 @@ const { err: serverErr, server: smpp } = await server({
return variant ? { userData: { variant } } : false;
},
idleTimeout: 40_000,
+ onSms: sms => {
+ const variant = (sms.session.userData as { variant?: Variant } | undefined)?.variant;
+
+ if (variant) allSms.push({ sms, variant });
+
+ // max-pending-submits=1 means bearerbox holds the link to one outstanding submit_sm at a
+ // time - answer each as it lands, or the whole burst stalls behind the first message.
+ if (variant === 'maxp1') void sms.sendResp();
+ },
port: SMPP_PORT,
});
@@ -171,12 +180,6 @@ smppServer.on('session', session => {
if (variant) bindPdus.push({ params: pduObj.params, variant });
});
- session.on('sms', sms => {
- const variant = (session.userData as { variant?: Variant } | undefined)?.variant;
-
- if (variant) allSms.push({ sms, variant });
- });
-
session.on('dlr', dlr => {
const variant = (session.userData as { variant?: Variant } | undefined)?.variant;
@@ -517,10 +520,6 @@ describe('maxp1 variant - max-pending-submits 1', () => {
assert.ok(session);
- // max-pending-submits=1 means bearerbox holds the link to one outstanding submit_sm at a
- // time - answer each as it lands, or the whole burst stalls behind the first message.
- session.on('sms', sms => { void sms.sendResp(); });
-
const texts = Array.from({ length: 20 }, (_, i) => `burst-${String(i).padStart(2, '0')}`);
// Sequential, not Promise.all: concurrent fetch()es reach smsbox's HTTP listener in whatever
diff --git a/interop-tests/php.test.ts b/interop-tests/php.test.ts
index 0415ce9..c273f93 100644
--- a/interop-tests/php.test.ts
+++ b/interop-tests/php.test.ts
@@ -48,6 +48,17 @@ const { err: serverErr, server: smpp } = await server({
return true;
},
+ onSms: sms => {
+ const systemId = systemIdBySession.get(sms.session);
+
+ if (systemId) allSms.push({ sms, systemId });
+
+ // php-smpp's submit_sm() blocks synchronously reading the response on the same connection
+ // that sent it, so answering here (rather than after this event's own test observes the
+ // sms) is the only way that read ever completes - unlike python-smpplib's driver, this one
+ // has no separate reader thread to poll afterwards.
+ void sms.sendResp();
+ },
port: SMPP_PORT,
});
@@ -62,18 +73,6 @@ smppServer.on('session', session => {
systemIdBySession.set(session, paramText(pduObj.params.system_id));
});
-
- session.on('sms', sms => {
- const systemId = systemIdBySession.get(session);
-
- if (systemId) allSms.push({ sms, systemId });
-
- // php-smpp's submit_sm() blocks synchronously reading the response on the same connection
- // that sent it, so answering here (rather than after this event's own test observes the
- // sms) is the only way that read ever completes - unlike python-smpplib's driver, this one
- // has no separate reader thread to poll afterwards.
- void sms.sendResp();
- });
});
after(async () => {
diff --git a/interop-tests/python.test.ts b/interop-tests/python.test.ts
index 84f1c66..97631c0 100644
--- a/interop-tests/python.test.ts
+++ b/interop-tests/python.test.ts
@@ -68,6 +68,11 @@ const { err: serverErr, server: smpp } = await server({
return true;
},
+ onSms: sms => {
+ const systemId = systemIdBySession.get(sms.session);
+
+ if (systemId) allSms.push({ sms, systemId });
+ },
port: SMPP_PORT,
});
@@ -82,12 +87,6 @@ smppServer.on('session', session => {
systemIdBySession.set(session, paramText(pduObj.params.system_id));
});
-
- session.on('sms', sms => {
- const systemId = systemIdBySession.get(session);
-
- if (systemId) allSms.push({ sms, systemId });
- });
});
after(async () => {
@@ -149,7 +148,7 @@ async function waitForAck(name: string, sequence: number, budget = 8000): Promis
}
async function echoBack(session: Session, sms: Sms, dataCoding: number, text: string): Promise {
- const encName = dataCoding === 3 ? 'LATIN1' : dataCoding === 8 ? 'UCS2' : 'ASCII';
+ const encName = dataCoding === 3 ? 'LATIN1' : dataCoding === 8 ? 'UCS2' : 'GSM';
const buf = encodings[encName].encode(text);
const sent = await session.send({
cmdName: 'deliver_sm',
diff --git a/interop-tests/smppsim.test.ts b/interop-tests/smppsim.test.ts
index 8a08817..b5e775d 100644
--- a/interop-tests/smppsim.test.ts
+++ b/interop-tests/smppsim.test.ts
@@ -72,12 +72,11 @@ function collectDlrs(session: Session): Received[] {
return received;
}
-function collectSms(session: Session): Sms[] {
+/** Every message a client's handler was handed; `onSms` goes to the bind. */
+function smsInbox(): { collected: Sms[]; onSms: (sms: Sms) => void } {
const collected: Sms[] = [];
- session.on('sms', sms => { collected.push(sms); });
-
- return collected;
+ return { collected, onSms: sms => { collected.push(sms); } };
}
const DLR_RETRY_BUDGET_MS = 3000;
@@ -202,14 +201,15 @@ describe('smppsim - C3+C7 long MT, receipts and loopback reassembly', () => {
for (const testCase of cases) {
test(testCase.label, async t => {
- const { err, session } = await bind(PEER_HOST);
+ const inbox = smsInbox();
+ const { err, session } = await bind(PEER_HOST, { onSms: inbox.onSms });
assert.equal(err, undefined);
assert.ok(session);
closeAfter(t, session);
const dlrs = collectDlrs(session);
- const sms = collectSms(session);
+ const sms = inbox.collected;
const { reassembled, smsIds } = await sendUntilComplete(
session,
@@ -588,13 +588,14 @@ describe('smppsim - C15 bind version negotiation', () => {
describe('smppsim - C17 encodings round trip over loopback', () => {
test('Latin-1 (å ä ö)', async t => {
- const { err, session } = await bind(PEER_HOST);
+ const inbox = smsInbox();
+ const { err, session } = await bind(PEER_HOST, { onSms: inbox.onSms });
assert.equal(err, undefined);
assert.ok(session);
closeAfter(t, session);
- const sms = collectSms(session);
+ const sms = inbox.collected;
await session.sendSms({ encoding: 'LATIN1', from: FROM, message: 'å ä ö', to: TO });
@@ -605,13 +606,14 @@ describe('smppsim - C17 encodings round trip over loopback', () => {
});
test('UCS-2', async t => {
- const { err, session } = await bind(PEER_HOST);
+ const inbox = smsInbox();
+ const { err, session } = await bind(PEER_HOST, { onSms: inbox.onSms });
assert.equal(err, undefined);
assert.ok(session);
closeAfter(t, session);
- const sms = collectSms(session);
+ const sms = inbox.collected;
await session.sendSms({ encoding: 'UCS2', from: FROM, message: 'ucs2 round trip', to: TO });
@@ -622,13 +624,14 @@ describe('smppsim - C17 encodings round trip over loopback', () => {
});
test('flash (data_coding records the message-class group)', async t => {
- const { err, session } = await bind(PEER_HOST);
+ const inbox = smsInbox();
+ const { err, session } = await bind(PEER_HOST, { onSms: inbox.onSms });
assert.equal(err, undefined);
assert.ok(session);
closeAfter(t, session);
- const sms = collectSms(session);
+ const sms = inbox.collected;
await session.sendSms({ flash: true, from: FROM, message: 'flash test', to: TO });
@@ -640,13 +643,14 @@ describe('smppsim - C17 encodings round trip over loopback', () => {
});
test('a raw submit_sm with data_coding 0xF0 is read as flash', async t => {
- const { err, session } = await bind(PEER_HOST);
+ const inbox = smsInbox();
+ const { err, session } = await bind(PEER_HOST, { onSms: inbox.onSms });
assert.equal(err, undefined);
assert.ok(session);
closeAfter(t, session);
- const sms = collectSms(session);
+ const sms = inbox.collected;
const body = 'message class test';
const sent = await session.send({
@@ -668,13 +672,14 @@ describe('smppsim - C17 encodings round trip over loopback', () => {
});
test('a raw submit_sm with 8-bit binary and a UDH (esm_class 0x40)', async t => {
- const { err, session } = await bind(PEER_HOST);
+ const inbox = smsInbox();
+ const { err, session } = await bind(PEER_HOST, { onSms: inbox.onSms });
assert.equal(err, undefined);
assert.ok(session);
closeAfter(t, session);
- const sms = collectSms(session);
+ const sms = inbox.collected;
// A UDH carrying no recognised concatenation IE (0x00/0x08): one element in GSM 03.40's
// reserved-for-future-use range (0x70), so Wireshark's gsm_sms_ud dissector - which
// validates the *typed* IEs' own lengths (0x01 "Special SMS Message Indication" must be
diff --git a/interop-tests/smscsim.test.ts b/interop-tests/smscsim.test.ts
index ba23042..0c1785c 100644
--- a/interop-tests/smscsim.test.ts
+++ b/interop-tests/smscsim.test.ts
@@ -163,9 +163,11 @@ describe('smscsim - multipart segments', () => {
describe('smscsim - MO injection through the web UI', () => {
test('a message posted to the web page arrives as an sms event', async t => {
+ const incoming: Sms[] = [];
const { err, session } = await client({
bindType: 'transceiver',
host: PEER_HOST,
+ onSms: sms => { incoming.push(sms); },
port: PEER_PORT,
username: 'mo-inject',
});
@@ -174,10 +176,6 @@ describe('smscsim - MO injection through the web UI', () => {
assert.ok(session);
closeAfter(t, session);
- const incoming: Sms[] = [];
-
- session.on('sms', sms => { incoming.push(sms); });
-
const response = await fetch(`http://${PEER_HOST}:${String(PEER_WEB_PORT)}/`, {
body: new URLSearchParams({
message: 'hello from the web UI',
diff --git a/src/bind-direction.ts b/src/bind-direction.ts
new file mode 100644
index 0000000..b810cf6
--- /dev/null
+++ b/src/bind-direction.ts
@@ -0,0 +1,71 @@
+import type { Result } from './result.ts';
+import { namedValue } from './error-from.ts';
+
+export type BindType = 'receiver' | 'transceiver' | 'transmitter';
+
+const bindTypes: readonly BindType[] = ['receiver', 'transceiver', 'transmitter'];
+
+export const bindCommands: readonly string[] = bindTypes.map(bindType => `bind_${bindType}`);
+
+export function bindTypeFromCommand(cmdName: string): BindType | undefined {
+ return bindTypes.find(bindType => `bind_${bindType}` === cmdName);
+}
+
+function isBindType(value: unknown): value is BindType {
+ return bindTypes.some(bindType => bindType === value);
+}
+
+/** Which end of the link a session is. Only `server()` is the SMSC; everything else is the ESME. */
+export type LinkEnd = 'esme' | 'smsc';
+
+/**
+ * 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 quoted(value: unknown): string {
+ return typeof value === 'string' ? JSON.stringify(value) : namedValue(value);
+}
+
+/** 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..3ddd716 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 { OnSms, ReconnectOptions } from './session-options.ts';
import type { SmppLog } from './log.ts';
import type { SmsIdFormat } from './sms-id.ts';
import type { Socket } from 'node:net';
@@ -11,7 +12,7 @@ import { Session } from './session.ts';
import { checkSessionOptions } from './session-options.ts';
import { connect as netConnect } from 'node:net';
import { connect as tlsConnect } from 'node:tls';
-import { defaultInterfaceVersion } from './defs/constants.ts';
+import { defaults } from './defaults.ts';
import { guardedLog } from './log.ts';
/** `fromStart` puts the very first connect and bind through the same backoff loop as a drop. */
@@ -29,6 +30,7 @@ export type ClientOptions = {
interfaceVersion?: number;
log?: SmppLog;
maxOutstanding?: number;
+ onSms?: OnSms;
password?: string;
port?: number;
reconnect?: ReconnectTuning | false;
@@ -41,19 +43,6 @@ export type ClientOptions = {
username?: string;
};
-const defaults = {
- bindType: 'transceiver',
- connectTimeout: 10_000,
- enquireLinkInterval: 20_000,
- host: 'localhost',
- /** The idle timeout is what notices a dead link, so it has to outlast one silent probe. */
- idleTimeoutFactor: 2,
- interfaceVersion: defaultInterfaceVersion,
- password: 'pass',
- port: 2775,
- username: 'user',
-} as const;
-
function armConnectTimeout(
sock: Socket,
connectTimeout: number | false,
@@ -201,9 +190,11 @@ function createSession(options: ClientOptions, log: SmppLog, sock: Socket): Sess
return new Session({
enquireLinkInterval,
- idleTimeout: options.idleTimeout ?? enquireLinkInterval * defaults.idleTimeoutFactor,
+ // Two silent probes: one lost enquire_link must not drop a live link.
+ idleTimeout: options.idleTimeout ?? enquireLinkInterval * 2,
log,
maxOutstanding: options.maxOutstanding,
+ onSms: options.onSms,
reconnect: reconnectFor(options, log),
responseTimeout: options.responseTimeout,
shutdownTimeout: options.shutdownTimeout,
diff --git a/src/defaults.ts b/src/defaults.ts
new file mode 100644
index 0000000..dcaa84b
--- /dev/null
+++ b/src/defaults.ts
@@ -0,0 +1,31 @@
+import { defaultInterfaceVersion } from './defs/constants.ts';
+
+/** Every default the session layer runs on. README's option tables quote these. */
+export const defaults = {
+ bindType: 'transceiver',
+ connectTimeout: 10_000,
+ /** Receipts of a multipart message can be a working day apart, so the cap does the bounding. */
+ dlrMergeTimeout: 86_400_000,
+ enquireLinkInterval: 20_000,
+ /** A handler still running this long is no longer counted or waited for. */
+ handlerTimeout: 300_000,
+ host: 'localhost',
+ /** Two silent enquire_link intervals, so one lost probe does not drop a live link. */
+ idleTimeout: 40_000,
+ interfaceVersion: defaultInterfaceVersion,
+ maxDlrMerges: 1000,
+ maxHandledMessages: 1000,
+ maxHandledOctets: 64 * 1024 * 1024,
+ maxOutstanding: 10,
+ maxReassembly: 1000,
+ maxReassemblyOctets: 64 * 1024 * 1024,
+ password: 'pass',
+ port: 2775,
+ reassemblyTimeout: 300_000,
+ reconnectMaxDelay: 30_000,
+ reconnectMinDelay: 1000,
+ responseTimeout: 30_000,
+ shutdownTimeout: 5000,
+ systemId: '',
+ username: 'user',
+} as const;
diff --git a/src/defs/encodings.ts b/src/defs/encodings.ts
index c454b02..a1c6216 100644
--- a/src/defs/encodings.ts
+++ b/src/defs/encodings.ts
@@ -1,4 +1,4 @@
-export type EncodingName = 'ASCII' | 'LATIN1' | 'UCS2';
+export type EncodingName = 'GSM' | 'LATIN1' | 'UCS2';
export type Encoding = {
decode: (buffer: Uint8Array) => string;
@@ -103,7 +103,7 @@ const ucs2: Encoding = {
};
export const encodings: Record = {
- ASCII: ascii,
+ GSM: ascii,
LATIN1: latin1,
UCS2: ucs2,
};
@@ -115,7 +115,7 @@ export function isEncodingName(value: unknown): value is EncodingName {
}
export function detect(value: string): EncodingName {
- if (encodings.ASCII.match(value)) return 'ASCII';
+ if (encodings.GSM.match(value)) return 'GSM';
if (encodings.LATIN1.match(value)) return 'LATIN1';
return 'UCS2';
@@ -163,21 +163,21 @@ function messageClassEncoding(dataCoding: number): EncodingName | undefined {
if (messageClassOf(dataCoding) === undefined) return undefined;
if ((dataCoding & 0xF0) === 0xF0) {
- return (dataCoding & 0x04) === 0x04 ? 'LATIN1' : 'ASCII';
+ return (dataCoding & 0x04) === 0x04 ? 'LATIN1' : 'GSM';
}
const alphabet = (dataCoding >> 2) & 0x03;
if (alphabet === 0x01) return 'LATIN1';
- return alphabet === 0x02 ? 'UCS2' : 'ASCII';
+ return alphabet === 0x02 ? 'UCS2' : 'GSM';
}
/**
* SMPP data_coding is a flat table for 0x00-0x0E, and the message class ranges are how a flash UCS2
* message arrives as 0x18. The 8-bit binary codings resolve to LATIN1, the one codec here that maps
* every octet to a code point and back unchanged, so a binary payload survives; alphabets with no
- * codec fall back to ASCII.
+ * codec fall back to GSM.
*/
export function encodingByDataCoding(dataCoding: number): EncodingName {
const messageClass = messageClassEncoding(dataCoding);
@@ -186,7 +186,7 @@ export function encodingByDataCoding(dataCoding: number): EncodingName {
if (dataCoding === 0x08) return 'UCS2';
// 0x02 and 0x04 are 8-bit binary, 0x03 is Latin-1.
- return dataCoding >= 0x02 && dataCoding <= 0x04 ? 'LATIN1' : 'ASCII';
+ return dataCoding >= 0x02 && dataCoding <= 0x04 ? 'LATIN1' : 'GSM';
}
/**
@@ -194,7 +194,7 @@ export function encodingByDataCoding(dataCoding: number): EncodingName {
* takes 0x00, the SMSC default alphabet, rather than SMPP 3.4 5.2.19's 0x01, which is IA5.
*/
export const dataCodingByEncoding: Readonly> = {
- ASCII: 0x00,
+ GSM: 0x00,
LATIN1: 0x03,
UCS2: 0x08,
};
diff --git a/src/dlr-merger.ts b/src/dlr-merger.ts
index 0c20cae..b61529a 100644
--- a/src/dlr-merger.ts
+++ b/src/dlr-merger.ts
@@ -126,7 +126,7 @@ export class DlrMerger {
if (group.parts.size < group.expected.size) return undefined;
- this.close(base);
+ this.spend(base);
const segments = [...group.parts.entries()].sort(([a], [b]) => a - b).map(([, one]) => one);
const worst = segments.reduce((carry, one) => (severity[one.statusMsg] > severity[carry.statusMsg] ? one : carry));
@@ -142,7 +142,7 @@ export class DlrMerger {
/** Drops every group past its deadline. Runs before each collect and on its own timer. */
sweep(): void {
for (const [base, group] of this.groups.takeExpired()) {
- this.close(base);
+ this.spend(base);
this.log.info('dlrMerger - incomplete receipts expired', { base, expected: group.expected.size });
}
}
@@ -151,7 +151,7 @@ export class DlrMerger {
this.spent.takeExpired();
if (this.groups.get(base) !== undefined || this.spent.get(base) === true) {
- this.close(base);
+ this.spend(base);
this.log.info('dlrMerger - message id handed out again, leaving its receipts unmerged', { base });
return;
@@ -162,7 +162,8 @@ export class DlrMerger {
this.groups.set(base, { expected, parts: new Map() });
}
- private close(base: string): void {
+ /** The base is finished with, merged or not, and never opens again. */
+ private spend(base: string): void {
this.groups.delete(base);
this.spent.delete(base);
@@ -178,7 +179,7 @@ export class DlrMerger {
const [base] = oldest;
- this.close(base);
+ this.spend(base);
this.log.warn('dlrMerger - buffer full, dropping the oldest message', { base, max: this.max });
}
}
diff --git a/src/handled-messages.ts b/src/handled-messages.ts
new file mode 100644
index 0000000..47776b3
--- /dev/null
+++ b/src/handled-messages.ts
@@ -0,0 +1,157 @@
+import type { ErrorName } from './defs/errors.ts';
+import type { OnSms } from './session-options.ts';
+import type { Sms, SmsHandlers, SmsInput } from './sms.ts';
+import type { SmppLog } from './log.ts';
+import { ExpiringGroups } from './expiring-groups.ts';
+import { IdleWaiters } from './idle-waiters.ts';
+import { createSms } from './sms.ts';
+import { errorFrom } from './error-from.ts';
+import { retainedOctets } from './retained-pdu.ts';
+
+export type HandledMessagesOptions = {
+ log: SmppLog;
+ max: number;
+ maxOctets: number;
+ /** Injected so expiry can be exercised without a wall clock. */
+ now?: (() => number) | undefined;
+ onSms: OnSms | undefined;
+ /** Where a handler's failure is reported. */
+ report: (err: Error) => void;
+ timeout: number;
+};
+
+/** What a message needs from the session, less the notice this class takes for itself. */
+export type SmsRoute = Omit;
+
+type Handled = { answered: boolean; sms: Sms };
+
+/**
+ * The messages whose handler is running. Each counts toward the bound and holds a shutdown until
+ * the handler settles, its deadline passes, or the link goes. A message a handler failed on, or
+ * that no handler takes, is answered here where it was not, asking the peer to retry.
+ */
+export class HandledMessages {
+ private readonly idleWaiters = new IdleWaiters();
+ private readonly log: SmppLog;
+ private readonly max: number;
+ private readonly maxOctets: number;
+ private readonly onSms: OnSms | undefined;
+ private readonly report: (err: Error) => void;
+ private readonly running: ExpiringGroups;
+ private atBound = false;
+ private keys = 0;
+
+ constructor(options: HandledMessagesOptions) {
+ this.log = options.log;
+ this.max = options.max;
+ this.maxOctets = options.maxOctets;
+ this.onSms = options.onSms;
+ this.report = options.report;
+ this.running = new ExpiringGroups({
+ max: options.max,
+ now: options.now,
+ onSweep: () => { this.sweep(); },
+ timeout: options.timeout,
+ });
+ }
+
+ get octets(): number {
+ return this.running.weight;
+ }
+
+ get size(): number {
+ return this.running.size;
+ }
+
+ /** Whether a message arriving now is refused: at the bound, and until the store is half empty again. */
+ refuses(): boolean {
+ this.sweep();
+
+ if (this.running.full || this.running.weight >= this.maxOctets) {
+ if (!this.atBound) {
+ this.atBound = true;
+ this.log.warn('handledMessages - messages at their bound, refusing new ones until handlers return', {
+ messages: this.size,
+ octets: this.octets,
+ });
+ }
+
+ return true;
+ }
+
+ // Half, so a peer keeping its window full does not flip this on every answer.
+ if (this.atBound && this.size <= this.max / 2 && this.octets <= this.maxOctets / 2) {
+ this.atBound = false;
+ this.log.info('handledMessages - messages down to half their bound, accepting again', { messages: this.size });
+ }
+
+ return false;
+ }
+
+ /** Hands the message to the handler. `retryStatus` answers it where the handler fails first. */
+ offer(input: SmsInput, route: SmsRoute, retryStatus: ErrorName): Sms {
+ const key = String(this.keys++);
+ const handled: Handled = { answered: false, sms: createSms(input, { ...route, answered: () => { handled.answered = true; } }) };
+
+ this.sweep();
+ this.running.set(key, handled);
+ this.running.weigh(key, input.pduObjs.reduce((sum, pduObj) => sum + retainedOctets(pduObj), 0));
+ void this.run(key, handled, retryStatus);
+
+ return handled.sms;
+ }
+
+ /** Drops every message: their segments went with the link, so no answer of ours correlates now. */
+ clear(): void {
+ this.running.takeAll();
+ this.idleWaiters.settle();
+ }
+
+ /** Resolves 0 once every handler has returned, or with how many have not. */
+ idle(timeout: number, signal: AbortSignal | undefined): Promise {
+ return this.idleWaiters.wait(() => this.running.size, timeout, signal);
+ }
+
+ /** Drops every message past its deadline. Runs before each offer and on its own timer. */
+ sweep(): void {
+ const expired = this.running.takeExpired();
+
+ if (expired.length === 0) return;
+
+ this.log.warn('handledMessages - handlers still running past their deadline', { messages: expired.length });
+ this.settle();
+ }
+
+ private async run(key: string, handled: Handled, retryStatus: ErrorName): Promise {
+ const failure = await this.handle(handled.sms);
+
+ if (failure) {
+ this.log.error('handledMessages - a handler failed', { message: failure.message });
+ this.report(failure);
+
+ // A handler that failed before answering has decided nothing, so the peer keeps the message.
+ if (!handled.answered) await handled.sms.sendResp({ status: retryStatus });
+ }
+
+ if (this.running.get(key) !== handled) return;
+
+ this.running.delete(key);
+ this.settle();
+ }
+
+ private async handle(sms: Sms): Promise {
+ if (!this.onSms) return new Error('No onSms handler takes inbound messages');
+
+ try {
+ await this.onSms(sms);
+
+ return undefined;
+ } catch (thrown: unknown) {
+ return errorFrom(thrown);
+ }
+ }
+
+ private settle(): void {
+ if (this.running.size === 0) this.idleWaiters.settle();
+ }
+}
diff --git a/src/held-messages.ts b/src/held-messages.ts
deleted file mode 100644
index b9e740e..0000000
--- a/src/held-messages.ts
+++ /dev/null
@@ -1,201 +0,0 @@
-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';
-import type { SmsHandlers } from './sms.ts';
-import type { SmppLog } from './log.ts';
-import { ExpiringGroups } from './expiring-groups.ts';
-import { IdleWaiters } from './idle-waiters.ts';
-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;
- sendPastDrain: SmsHandlers['send'];
- session: Session;
- timeout: 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];
-
- 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.
- */
-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 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