Compare commits

1 Commits

Author SHA1 Message Date
lilleman 1b075bf5d2 Give the held-message flow one owner
Test / lint (pull_request) Successful in 23s
Test / test (18) (pull_request) Successful in 30s
Test / test (20) (pull_request) Successful in 30s
Test / test (22) (pull_request) Successful in 31s
Test / test (24) (pull_request) Successful in 30s
Test / test (26) (pull_request) Successful in 30s
Mirror / push (push) Successful in 5s
2026-09-29 18:43:44 +02:00
10 changed files with 144 additions and 140 deletions
+5 -12
View File
@@ -1,8 +1,5 @@
# AGENTS.md # AGENTS.md
Guidance for LLM agents working in this repository. What each file in it is for is under
[Documentation](#documentation).
## What this is ## What this is
A ground-up TypeScript rewrite of `larvitsmpp` 0.4.0, published as `@larvit/smpp` 0.5.0. The branch A ground-up TypeScript rewrite of `larvitsmpp` 0.4.0, published as `@larvit/smpp` 0.5.0. The branch
@@ -13,7 +10,7 @@ not for structure or style.
## Goals ## Goals
The goals, in priority order, live in The goals, in priority order, live in
[README.md](https://gitea.larvit.se/larvit/smpp-js/src/branch/main/README.md#goals). The README states the audience alongside them. [README.md](https://gitea.larvit.se/larvit/smpp-js/src/branch/main/README.md#goals).
## Hard rules ## Hard rules
@@ -50,7 +47,7 @@ src/
dlr-merger.ts DlrMerger: per-segment receipts counted into one MessageDlr 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 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 expiring-groups.ts ExpiringGroups: the capped, weighed, expiring store DlrMerger, HeldMessages and Reassembler share
held-messages.ts HeldMessages: capped, expiring messages the application has not answered, one MessageHold each 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 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 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: whether the link lives, and where a request waits for the next one
@@ -87,8 +84,8 @@ src/
Imports point one way: `defs` knows nothing above it but `result.ts`, `pdu` uses `defs`, `session` 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 uses `pdu`, and `client`/`server` use `session`. The ways back up are the `Session` handed to
`createSms()` and `IncomingRequests`, which call back into it, and to `OnRequest` and `onConnected` `createSms()`, `HeldMessages` and `IncomingRequests`, which call back into it, and to `OnRequest`
in `session-options.ts`, all imported as a type only. and `onConnected` in `session-options.ts`, all imported as a type only.
**Parameter order is wire order.** The key order inside `cmds.*.params` is the order the fields are **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 written to and read from the buffer. Never sort those alphabetically — the alphabetical-ordering
@@ -225,10 +222,6 @@ this is not a changelog.
## Decisions ## Decisions
The decisions themselves live in [docs/decisions.md](docs/decisions.md). Their titles are indexed
here, so a reader sees that a decision exists without carrying its reasoning; the reasoning is in
the file.
### [The public surface](docs/decisions.md#the-public-surface) ### [The public surface](docs/decisions.md#the-public-surface)
- `Session` is publicly constructible, which is what makes `SessionOptions` and `ReconnectOptions` - `Session` is publicly constructible, which is what makes `SessionOptions` and `ReconnectOptions`
@@ -315,7 +308,7 @@ the file.
### [Internals and tests](docs/decisions.md#internals-and-tests) ### [Internals and tests](docs/decisions.md#internals-and-tests)
- #30, #46 and #48 merged under the comprehension floor, and Locality is the next work. - Locality work comes before other work until a scoring run reads 7.0.
- A listener that rejects is routed by Node's `captureRejections`, not by hand-dispatching. - A listener that rejects is routed by Node's `captureRejections`, not by hand-dispatching.
- The four-line abort dance is copied across `LinkLife`, `IdleWaiters`, `PendingRequests` and - The four-line abort dance is copied across `LinkLife`, `IdleWaiters`, `PendingRequests` and
`SendWindow` rather than extracted. `SendWindow` rather than extracted.
+2
View File
@@ -74,6 +74,8 @@
`alert_on_message_delivery` and `broadcast_area_identifier`, the names they read back under. The `alert_on_message_delivery` and `broadcast_area_identifier`, the names they read back under. The
two alternate names are gone from `tlvs` too, which is now typed by `TlvName`: narrow a `string` with `isTlvName()` before indexing it. two alternate names are gone from `tlvs` too, which is now typed by `TlvName`: narrow a `string` with `isTlvName()` before indexing it.
- `cmds.broadcast_sm_resp.tlvMap` is removed; nothing read it. - `cmds.broadcast_sm_resp.tlvMap` is removed; nothing read it.
- The `SmsInput` type is no longer exported; nothing exported took one. Annotate with `Sms`, or a
`Pick<Sms, …>` of the fields you use.
- `session.boundAs` and `session.peerInterfaceVersion` are read-only, and `session.loggedIn` is - `session.boundAs` and `session.peerInterfaceVersion` are read-only, and `session.loggedIn` is
removed: read `session.boundAs !== undefined`. A session you wire yourself records the bind it removed: read `session.boundAs !== undefined`. A session you wire yourself records the bind it
accepted or had accepted with `session.bound(bindType, declaredVersion)`, which returns `err` for a accepted or had accepted with `session.bound(bindType, declaredVersion)`, which returns `err` for a
+3 -4
View File
@@ -713,10 +713,9 @@ one wins. They do not override the [hard rules](https://gitea.larvit.se/larvit/s
alphabet, a receipt format — it gets an option or a hook rather than a fork. A call that passes no alphabet, a receipt format — it gets an option or a hook rather than a fork. A call that passes no
options stays exactly as easy and as safe, and a hook is a seam the library calls, never a way into options stays exactly as easy and as safe, and a hook is a seam the library calls, never a way into
its internals. its internals.
8. **A small, stable public surface over reshapeable internals.** Only what `src/index.ts` exports is 8. **A small, stable public surface over reshapeable internals.** A new option has to beat "the
published. A new option has to beat "the application can do this itself", and has to keep a application can do this itself", and has to keep a promise this library can verify. The low-level
promise this library can verify. The low-level surface is a passthrough: policy binds what the surface is a passthrough: policy binds what the library composes, never what the caller wrote.
library composes, never what the caller wrote.
9. **State wider than one session goes through one store.** A pool of sessions, a limit shared 9. **State wider than one session goes through one store.** A pool of sessions, a limit shared
between processes, and what has to survive a restart — receipts still awaited, a message half between processes, and what has to survive a restart — receipts still awaited, a message half
reassembled — are held through a store interface and never beside it. Without a store the reassembled — are held through a store interface and never beside it. Without a store the
+8 -15
View File
@@ -8,8 +8,7 @@ rule and an index of the titles below.
- **`Session` is publicly constructible, which is what makes `SessionOptions` and `ReconnectOptions` - **`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 public too.** Raised twice as a leak; it is not one. The collaborators `session.ts` delegates to
(`IncomingRequests`, `OutgoingRequests`, `ReconnectLoop`, `LinkTimers`, `DlrMerger`, stay unpublished so they can be reshaped.
`PduTransport`, `submitSms`) stay unpublished so they can be reshaped.
- **`acceptsOptionalParams()` and `bindAllows()` are predicates, not chokepoints.** The library's own - **`acceptsOptionalParams()` and `bindAllows()` are predicates, not chokepoints.** The library's own
senders consult them; `session.send({ tlvs })` is passed through as written, because silently senders consult them; `session.send({ tlvs })` is passed through as written, because silently
@@ -418,8 +417,7 @@ rule and an index of the titles below.
caller would get wrong, where asking the codec is, and publishing it would freeze this library's caller would get wrong, where asking the codec is, and publishing it would freeze this library's
error prose as API for an application whose own refusal should read like itself. Goal 8, from the error prose as API for an application whose own refusal should read like itself. Goal 8, from the
architecture review of [#99](https://github.com/larvit/larvitsmpp/pull/99), 2026-09-09. It is architecture review of [#99](https://github.com/larvit/larvitsmpp/pull/99), 2026-09-09. It is
reached through reached through `encodeBody()` in `message.ts`, which is where the `data_coding`-to-text pair already lives:
`encodeBody()` in `message.ts`, which is where the `data_coding`-to-text pair already lives:
`encodeBody(text, dataCoding)` is `decodeMessage(buffer, dataCoding)`'s mirror and resolves the `encodeBody(text, dataCoding)` is `decodeMessage(buffer, dataCoding)`'s mirror and resolves the
alphabet through the same `encodingByDataCoding()`. `send()` and `sendReturn()` inherit it, alphabet through the same `encodingByDataCoding()`. `send()` and `sendReturn()` inherit it,
since both build through `buildPdu()`; `sendSms()` does not, and keeps its own guard, because since both build through `buildPdu()`; `sendSms()` does not, and keeps its own guard, because
@@ -641,8 +639,7 @@ rule and an index of the titles below.
failure belongs to one session's request, and that channel already carries every failure of one. failure belongs to one session's request, and that channel already carries every failure of one.
The hook is consulted before the bind-direction gate, so it sees a `submit_sm` a receiver-bound The hook is consulted before the bind-direction gate, so it sees a `submit_sm` a receiver-bound
peer may not send; first refusal means first, and one it declines still gets `ESME_RINVBNDSTS`. peer may not send; first refusal means first, and one it declines still gets `ESME_RINVBNDSTS`.
Nothing is held for a request the hook answered: `HeldMessages` is opened by the `sms` event the Nothing is held for a request the hook answered, so the drain waits on none of it. `OnRequest` stays unexported where
hook skipped, so the drain waits on none of it. `OnRequest` stays unexported where
`AuthenticateInput` is exported, because that hook's argument is a shape this library invents and `AuthenticateInput` is exported, because that hook's argument is a shape this library invents and
this one's are two types already published. Rejected: consulting the hook first, which puts this one's are two types already published. Rejected: consulting the hook first, which puts
bind and authentication inside the application's reach for nothing. Rejected: a narrower hook bind and authentication inside the application's reach for nothing. Rejected: a narrower hook
@@ -765,15 +762,11 @@ rule and an index of the titles below.
## Internals and tests ## Internals and tests
- **#30, #46 and #48 merged under the comprehension floor, and Locality is the next work.** Maintainer's call, - **Locality work comes before other work until a scoring run reads 7.0.** Maintainer's call,
2026-09-27. A four-seat scoring run, depth 1, read the project at 6, 6, 7 and 6 (mean 6.25), every 2026-09-27, when #30 merged under the comprehension floor at 6, 6, 7 and 6; #46, #48 and #49
seat capped by Locality in the held-message and shutdown code #30 does not touch, where the floor merged under it on that condition, #49 at 6, 6, 7 and 6 with Locality 5, 5, 6 and 6. Serves goal
is 7.0. The chunks after #30 lift Locality to 7 before any other work. Serves goal 8's 8's reshapeable internals, which a reader has to understand before reshaping. Valid until a
reshapeable internals, which a reader has to understand before reshaping. Valid until a scoring scoring run reads 7.0 or above.
run reads 7.0 or above. #46, maintainer's call 2026-09-28, merged as a step of that work at 6, 6, 7
and 7, Locality 5, 5, 6 and 6, up from 6, 6, 7 and 6 and Locality 5, 5, 6 and 5 the same day.
#48, maintainer's call 2026-09-28, merged at 6, 6, 6 and 6, Locality 5 from every seat, on the
condition that the held-message flow is the next chunk.
- **A listener that rejects is routed by Node's `captureRejections`, not by hand-dispatching.** Both - **A listener that rejects is routed by Node's `captureRejections`, not by hand-dispatching.** Both
emitters construct with `captureRejections: true` and implement emitters construct with `captureRejections: true` and implement
+51 -36
View File
@@ -1,15 +1,23 @@
import type { PduObject } from './pdu.ts'; 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 type { SmppLog } from './log.ts';
import { ExpiringGroups } from './expiring-groups.ts'; import { ExpiringGroups } from './expiring-groups.ts';
import { IdleWaiters } from './idle-waiters.ts'; import { IdleWaiters } from './idle-waiters.ts';
import { createSms } from './sms.ts';
import { retainedOctets } from './retained-pdu.ts'; import { retainedOctets } from './retained-pdu.ts';
export type HeldMessagesOptions = { export type HeldMessagesOptions = {
link: LinkLife;
log: SmppLog; log: SmppLog;
max: number; max: number;
maxOctets: number; maxOctets: number;
/** Injected so expiry can be exercised without a wall clock. */ /** Injected so expiry can be exercised without a wall clock. */
now?: (() => number) | undefined; now?: (() => number) | undefined;
sendPastDrain: SmsHandlers['send'];
session: Session;
timeout: number; timeout: number;
}; };
@@ -20,33 +28,40 @@ function keyOf(pduObjs: PduObject[]): string | undefined {
return first ? String(first.seqNr) : undefined; return first ? String(first.seqNr) : undefined;
} }
type HoldEntry = { type HoldRoute = Pick<HeldMessagesOptions, 'link' | 'sendPastDrain' | 'session'>;
isHeld: () => boolean;
release: () => void;
};
/** /**
* One message offered to the application. A drain waits on it until the first of: `answered()`, * One message offered to the application, and the handlers its `Sms` answers through. A drain
* every listener that took it rejecting, no listener taking it or one throwing, a later message on * waits on it until the first of: `answered()`, every listener that took it rejecting, no listener
* its sequence number, its deadline, or the link going. * taking it or one throwing, a later message on its sequence number, its deadline, or the link going.
*/ */
export class MessageHold { export class MessageHold implements SmsHandlers {
private readonly entry: HoldEntry; private readonly generation: number;
private readonly heldMessages: HeldMessages;
private readonly pduObjs: PduObject[];
private readonly route: HoldRoute;
private working: number; private working: number;
constructor(entry: HoldEntry, listeners: number) { constructor(heldMessages: HeldMessages, route: HoldRoute, pduObjs: PduObject[], listeners: number) {
this.entry = entry; this.generation = route.link.generation();
this.heldMessages = heldMessages;
this.pduObjs = pduObjs;
this.route = route;
this.working = listeners; this.working = listeners;
} }
/** Whether a drain is still waiting for this message to be answered. */ /** Whether a drain is still waiting for this message to be answered. */
isHeld(): boolean { isHeld(): boolean {
return this.entry.isHeld(); return this.heldMessages.holds(this.pduObjs);
} }
/** A turn later, so a listener sending its receipt straight after the response still holds. */ /** A turn later, so a `sendDlr()` called straight after `sendResp()` still goes out past a drain. */
answered(): void { answered(): void {
setImmediate(() => { this.entry.release(); }); 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. */ /** A rejection leaves the other listeners running, so only the last one to fail gives the message up. */
@@ -58,7 +73,12 @@ export class MessageHold {
/** At once, for a message nobody took or a listener threw on: that is not work a shutdown can wait for. */ /** At once, for a message nobody took or a listener threw on: that is not work a shutdown can wait for. */
release(): void { release(): void {
this.entry.release(); 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<Result<{ pduObj: PduObject }>> {
return this.isHeld() ? this.route.sendPastDrain(input) : this.route.session.send(input);
} }
} }
@@ -70,6 +90,7 @@ export class HeldMessages {
private readonly maxOctets: number; private readonly maxOctets: number;
/** A rejecting listener hands the message back as an `unknown`, so its hold is found by identity. */ /** A rejecting listener hands the message back as an `unknown`, so its hold is found by identity. */
private readonly offered = new WeakMap<object, MessageHold>(); private readonly offered = new WeakMap<object, MessageHold>();
private readonly route: HoldRoute;
constructor(options: HeldMessagesOptions) { constructor(options: HeldMessagesOptions) {
this.held = new ExpiringGroups({ this.held = new ExpiringGroups({
@@ -80,6 +101,7 @@ export class HeldMessages {
}); });
this.log = options.log; this.log = options.log;
this.maxOctets = options.maxOctets; this.maxOctets = options.maxOctets;
this.route = { link: options.link, sendPastDrain: options.sendPastDrain, session: options.session };
} }
get octetsHeld(): number { get octetsHeld(): number {
@@ -97,14 +119,8 @@ export class HeldMessages {
return this.held.full || this.held.weight >= this.maxOctets; return this.held.full || this.held.weight >= this.maxOctets;
} }
private hold(pduObjs: PduObject[], listeners: number): MessageHold { private hold(key: string, pduObjs: PduObject[], listeners: number): MessageHold {
const hold = new MessageHold({ const hold = new MessageHold(this, this.route, pduObjs, listeners);
isHeld: () => this.has(pduObjs),
release: () => { this.release(pduObjs); },
}, listeners);
const key = keyOf(pduObjs);
if (key === undefined) return hold;
this.sweep(); this.sweep();
@@ -118,18 +134,17 @@ export class HeldMessages {
return hold; return hold;
} }
offer<T extends object>( offer(pduObjs: PduObject[], answeredAs?: string): MessageHold | undefined {
pduObjs: PduObject[], const key = keyOf(pduObjs);
listeners: number,
build: (hold: MessageHold) => T,
emit: (message: T) => boolean,
): MessageHold {
const hold = this.hold(pduObjs, listeners);
const message = build(hold);
this.offered.set(message, hold); if (key === undefined) return undefined;
if (!emit(message)) hold.release(); const hold = this.hold(key, pduObjs, this.route.session.listenerCount('sms'));
const sms = createSms({ answeredAs, pduObjs, session: this.route.session }, hold);
this.offered.set(sms, hold);
if (!this.route.session.emit('sms', sms)) hold.release();
return hold; return hold;
} }
@@ -141,13 +156,13 @@ export class HeldMessages {
this.offered.get(message)?.listenerGaveUp(); this.offered.get(message)?.listenerGaveUp();
} }
private has(pduObjs: PduObject[]): boolean { holds(pduObjs: PduObject[]): boolean {
const key = keyOf(pduObjs); const key = keyOf(pduObjs);
return key !== undefined && this.held.get(key) === pduObjs; return key !== undefined && this.held.get(key) === pduObjs;
} }
private release(pduObjs: PduObject[]): void { release(pduObjs: PduObject[]): void {
const key = keyOf(pduObjs); const key = keyOf(pduObjs);
// Identity, not the key: a wrapped sequence number must not release someone else's message. // Identity, not the key: a wrapped sequence number must not release someone else's message.
+10 -32
View File
@@ -1,22 +1,21 @@
import type { Concat } from './concat.ts'; import type { Concat } from './concat.ts';
import type { DlrMerger } from './dlr-merger.ts'; import type { DlrMerger } from './dlr-merger.ts';
import type { ErrorName } from './defs/errors.ts'; import type { ErrorName } from './defs/errors.ts';
import type { HeldMessagesOptions } from './held-messages.ts';
import type { LinkLife } from './link-life.ts'; import type { LinkLife } from './link-life.ts';
import type { LostGroup, Refusal } from './reassembly.ts'; import type { LostGroup, Refusal } from './reassembly.ts';
import type { OnRequest } from './session-options.ts'; import type { OnRequest } from './session-options.ts';
import type { PduObject, PduObjectInput } from './pdu.ts'; import type { PduObject } from './pdu.ts';
import type { Result, VoidResult } from './result.ts'; import type { VoidResult } from './result.ts';
import type { Session } from './session.ts'; import type { Session } from './session.ts';
import type { SmppLog } from './log.ts'; import type { SmppLog } from './log.ts';
import type { SmsIdFormat } from './sms-id.ts'; import type { SmsIdFormat } from './sms-id.ts';
import { HeldMessages } from './held-messages.ts'; import { HeldMessages } from './held-messages.ts';
import { Reassembler, decodeSegments } from './reassembly.ts'; import { Reassembler } from './reassembly.ts';
import { bindCommands, defaults, standsInFor } from './session-options.ts'; import { bindCommands, defaults, standsInFor } from './session-options.ts';
import { concatOf } from './concat.ts'; import { concatOf } from './concat.ts';
import { createSms } from './sms.ts';
import { detach } from './retained-pdu.ts'; import { detach } from './retained-pdu.ts';
import { dlrFromPdu } from './dlr.ts'; import { dlrFromPdu } from './dlr.ts';
import { paramText } from './defs/types.ts';
import { respIdParams, segmentId } from './sms-id.ts'; import { respIdParams, segmentId } from './sms-id.ts';
import { respNameFor } from './defs/commands.ts'; import { respNameFor } from './defs/commands.ts';
@@ -52,8 +51,7 @@ export type IncomingRequestsOptions = {
maxReassembly?: number | undefined; maxReassembly?: number | undefined;
onRequest?: OnRequest | undefined; onRequest?: OnRequest | undefined;
reassemblyTimeout?: number | undefined; reassemblyTimeout?: number | undefined;
/** Past a drain's refusal, for a receipt the drain is itself waiting for. */ sendPastDrain: HeldMessagesOptions['sendPastDrain'];
sendPastDrain: (input: PduObjectInput) => Promise<Result<{ pduObj: PduObject }>>;
session: Session; session: Session;
smsIdFormat?: SmsIdFormat | undefined; smsIdFormat?: SmsIdFormat | undefined;
systemId?: string | undefined; systemId?: string | undefined;
@@ -67,7 +65,6 @@ export class IncomingRequests {
private readonly log: SmppLog; private readonly log: SmppLog;
private readonly onRequest: OnRequest | undefined; private readonly onRequest: OnRequest | undefined;
private readonly reassembler: Reassembler; private readonly reassembler: Reassembler;
private readonly sendPastDrain: IncomingRequestsOptions['sendPastDrain'];
private readonly session: Session; private readonly session: Session;
private readonly smsIdFormat: SmsIdFormat; private readonly smsIdFormat: SmsIdFormat;
private readonly systemId: string; private readonly systemId: string;
@@ -76,9 +73,12 @@ export class IncomingRequests {
constructor(options: IncomingRequestsOptions) { constructor(options: IncomingRequestsOptions) {
this.dlrMerger = options.dlrMerger; this.dlrMerger = options.dlrMerger;
this.held = new HeldMessages({ this.held = new HeldMessages({
link: options.link,
log: options.log, log: options.log,
max: defaults.maxHeldMessages, max: defaults.maxHeldMessages,
maxOctets: defaults.maxHeldOctets, maxOctets: defaults.maxHeldOctets,
sendPastDrain: options.sendPastDrain,
session: options.session,
timeout: defaults.heldMessageTimeout, timeout: defaults.heldMessageTimeout,
}); });
this.link = options.link; this.link = options.link;
@@ -91,7 +91,6 @@ export class IncomingRequests {
onLost: lost => { this.reportLost(lost); }, onLost: lost => { this.reportLost(lost); },
timeout: options.reassemblyTimeout ?? defaults.reassemblyTimeout, timeout: options.reassemblyTimeout ?? defaults.reassemblyTimeout,
}); });
this.sendPastDrain = options.sendPastDrain;
this.session = options.session; this.session = options.session;
this.smsIdFormat = options.smsIdFormat ?? {}; this.smsIdFormat = options.smsIdFormat ?? {};
this.systemId = options.systemId ?? defaults.systemId; this.systemId = options.systemId ?? defaults.systemId;
@@ -254,7 +253,7 @@ export class IncomingRequests {
const concat = concatOf(pduObj); const concat = concatOf(pduObj);
if (!concat) { if (!concat) {
this.emitSms([detach(pduObj)]); this.held.offer([detach(pduObj)]);
return; return;
} }
@@ -276,7 +275,7 @@ export class IncomingRequests {
respIdParams(pduObj.cmdName, segmentId(collected.smsId, concat.part - 1, concat.total)), respIdParams(pduObj.cmdName, segmentId(collected.smsId, concat.part - 1, concat.total)),
); );
if (collected.whole) this.emitSms(collected.whole, collected.smsId); if (collected.whole) this.held.offer(collected.whole, collected.smsId);
} }
private reportLost(lost: LostGroup): void { private reportLost(lost: LostGroup): void {
@@ -284,25 +283,4 @@ export class IncomingRequests {
`Gave up ${String(lost.parts)} of ${String(lost.total)} segments of an incomplete concatenated message: ${lostReasons[lost.reason]}`, `Gave up ${String(lost.parts)} of ${String(lost.total)} segments of an incomplete concatenated message: ${lostReasons[lost.reason]}`,
)); ));
} }
private emitSms(pduObjs: PduObject[], answeredAs?: string): void {
const first = pduObjs[0];
if (!first) return;
const generation = this.link.generation();
this.held.offer(pduObjs, this.session.listenerCount('sms'), hold => createSms({
answeredAs,
from: paramText(first.params.source_addr),
message: decodeSegments(pduObjs),
pduObjs,
session: this.session,
to: paramText(first.params.destination_addr),
}, {
lostLink: () => this.link.generation() !== generation,
onAnswered: () => { hold.answered(); },
send: input => (hold.isHeld() ? this.sendPastDrain(input) : this.session.send(input)),
}), sms => this.session.emit('sms', sms));
}
} }
+1 -1
View File
@@ -38,7 +38,7 @@ export { uuidv7 } from './uuid.ts';
export type { BindType, ClientOptions } from './client.ts'; export type { BindType, ClientOptions } from './client.ts';
export type { Dlr, Receipt } from './dlr.ts'; export type { Dlr, Receipt } from './dlr.ts';
export type { SendDlrResult, SendRespOptions, Sms, SmsInput } from './sms.ts'; export type { SendDlrResult, SendRespOptions, Sms } from './sms.ts';
export type { Concat } from './concat.ts'; export type { Concat } from './concat.ts';
export type { ConcatInfo } from './udh.ts'; export type { ConcatInfo } from './udh.ts';
export type { Result, VoidResult } from './result.ts'; export type { Result, VoidResult } from './result.ts';
+10 -12
View File
@@ -5,7 +5,9 @@ import type { Result, VoidResult } from './result.ts';
import type { Session } from './session.ts'; import type { Session } from './session.ts';
import { UnansweredError } from './unanswered-error.ts'; import { UnansweredError } from './unanswered-error.ts';
import { consts } from './defs/constants.ts'; import { consts } from './defs/constants.ts';
import { decodeSegments } from './reassembly.ts';
import { messageClassOf } from './defs/encodings.ts'; import { messageClassOf } from './defs/encodings.ts';
import { paramText } from './defs/types.ts';
import { receiptCodes, transientStates } from './dlr.ts'; import { receiptCodes, transientStates } from './dlr.ts';
import { smppDate } from './message.ts'; import { smppDate } from './message.ts';
import { respIdParams, segmentId } from './sms-id.ts'; import { respIdParams, segmentId } from './sms-id.ts';
@@ -59,17 +61,13 @@ export type Sms = {
export type SmsInput = { export type SmsInput = {
/** The id base the segments were already answered with; absent leaves the answer to `sendResp()`. */ /** The id base the segments were already answered with; absent leaves the answer to `sendResp()`. */
answeredAs?: string | undefined; answeredAs?: string | undefined;
from: string;
message: string;
pduObjs: PduObject[]; pduObjs: PduObject[];
session: Session; session: Session;
to: string;
}; };
/** What the session's incoming side gives a message so it can be answered and accounted for. */
export type SmsHandlers = { export type SmsHandlers = {
answered: () => void;
lostLink: () => boolean; lostLink: () => boolean;
onAnswered: () => void;
send: (input: PduObjectInput) => Promise<Result<{ pduObj: PduObject }>>; send: (input: PduObjectInput) => Promise<Result<{ pduObj: PduObject }>>;
}; };
@@ -86,8 +84,8 @@ export function createSms(input: SmsInput, handlers: SmsHandlers): Sms {
answeredOnArrival: input.answeredAs !== undefined, answeredOnArrival: input.answeredAs !== undefined,
dlr: typeof registered === 'number' && registered !== 0, dlr: typeof registered === 'number' && registered !== 0,
flash: typeof dataCoding === 'number' && messageClassOf(dataCoding) === immediateDisplayClass, flash: typeof dataCoding === 'number' && messageClassOf(dataCoding) === immediateDisplayClass,
from: input.from, from: paramText(first?.params.source_addr),
message: input.message, message: decodeSegments(input.pduObjs),
pduObjs: input.pduObjs, pduObjs: input.pduObjs,
sendDlr: status => sendDlr(sms, input.session, handlers, status), sendDlr: status => sendDlr(sms, input.session, handlers, status),
sendResp: options => (input.answeredAs === undefined sendResp: options => (input.answeredAs === undefined
@@ -98,7 +96,7 @@ export function createSms(input: SmsInput, handlers: SmsHandlers): Sms {
return answered.smsId; return answered.smsId;
}, },
submitTime: new Date(), submitTime: new Date(),
to: input.to, to: paramText(first?.params.destination_addr),
}; };
return sms; return sms;
@@ -107,7 +105,7 @@ export function createSms(input: SmsInput, handlers: SmsHandlers): Sms {
/** Every segment went out answered, so the call is what the shutdown waits for and nothing else. */ /** Every segment went out answered, so the call is what the shutdown waits for and nothing else. */
function answeredOnArrival( function answeredOnArrival(
options: SendRespOptions, options: SendRespOptions,
handlers: Pick<SmsHandlers, 'onAnswered'>, handlers: Pick<SmsHandlers, 'answered'>,
): Promise<VoidResult> { ): Promise<VoidResult> {
if (options.smsId !== undefined) { if (options.smsId !== undefined) {
return Promise.resolve({ return Promise.resolve({
@@ -121,7 +119,7 @@ function answeredOnArrival(
}); });
} }
handlers.onAnswered(); handlers.answered();
return Promise.resolve({}); return Promise.resolve({});
} }
@@ -131,7 +129,7 @@ async function sendResp(
session: Session, session: Session,
answered: { smsId: string }, answered: { smsId: string },
options: SendRespOptions, options: SendRespOptions,
handlers: Pick<SmsHandlers, 'lostLink' | 'onAnswered'>, handlers: Pick<SmsHandlers, 'answered' | 'lostLink'>,
): Promise<VoidResult> { ): Promise<VoidResult> {
const total = sms.pduObjs.length; const total = sms.pduObjs.length;
@@ -158,7 +156,7 @@ async function sendResp(
const failure = results.find(result => result.err); const failure = results.find(result => result.err);
if (!failure) handlers.onAnswered(); if (!failure) handlers.answered();
return failure ?? {}; return failure ?? {};
} }
+35 -18
View File
@@ -5,7 +5,7 @@ import type { Collected, LostGroup } from '../src/reassembly.ts';
import type { Dlr } from '../src/dlr.ts'; import type { Dlr } from '../src/dlr.ts';
import type { ErrorName } from '../src/defs/errors.ts'; import type { ErrorName } from '../src/defs/errors.ts';
import type { IncomingRequestsOptions } from '../src/incoming-requests.ts'; import type { IncomingRequestsOptions } from '../src/incoming-requests.ts';
import type { MessageHold } from '../src/held-messages.ts'; import type { HeldMessagesOptions, MessageHold } from '../src/held-messages.ts';
import type { MessageState } from '../src/defs/constants.ts'; import type { MessageState } from '../src/defs/constants.ts';
import type { MessageDlr } from '../src/session.ts'; import type { MessageDlr } from '../src/session.ts';
import type { PduObject, PduObjectInput } from '../src/pdu.ts'; import type { PduObject, PduObjectInput } from '../src/pdu.ts';
@@ -1543,11 +1543,34 @@ describe('held message bounds', () => {
} }
function offer(held: HeldMessages, seqNr: number): MessageHold { function offer(held: HeldMessages, seqNr: number): MessageHold {
return held.offer(message(seqNr), 1, () => ({}), () => true); const hold = held.offer(message(seqNr));
assert.ok(hold);
return hold;
} }
test('is full at its count, and a re-used sequence number replaces rather than adding', () => { /** Offers to a session with a listener, so an offer is held rather than released as untaken. */
const held = new HeldMessages({ log: silentLog, max: 2, maxOctets: 1_000_000, timeout: 10_000 }); function heldOn(
t: TestContext,
options: Pick<HeldMessagesOptions, 'max' | 'maxOctets' | 'now' | 'timeout'>,
): HeldMessages {
const session = new Session({ sock: new net.Socket() });
closeAfter(t, session);
session.on('sms', () => undefined);
return new HeldMessages({
...options,
link: new LinkLife({ log: silentLog, reconnects: false, timeout: 100 }),
log: silentLog,
sendPastDrain: () => Promise.resolve({ err: new Error('never sent') }),
session,
});
}
test('is full at its count, and a re-used sequence number replaces rather than adding', t => {
const held = heldOn(t, { max: 2, maxOctets: 1_000_000, timeout: 10_000 });
const first = offer(held, 1); const first = offer(held, 1);
const replaced = offer(held, 2); const replaced = offer(held, 2);
@@ -1563,9 +1586,9 @@ describe('held message bounds', () => {
}); });
// submitPdu() holds 1026 octets by the maxOctets charge: its object, and the three text fields. // submitPdu() holds 1026 octets by the maxOctets charge: its object, and the three text fields.
test('is full at its octet cap, until a message leaves by any way out', () => { test('is full at its octet cap, until a message leaves by any way out', t => {
let now = 0; let now = 0;
const held = new HeldMessages({ log: silentLog, max: 10, maxOctets: 2000, now: () => now, timeout: 10_000 }); const held = heldOn(t, { max: 10, maxOctets: 2000, now: () => now, timeout: 10_000 });
const answered = offer(held, 1); const answered = offer(held, 1);
assert.equal(held.full(), false); assert.equal(held.full(), false);
@@ -1658,9 +1681,9 @@ describe('held message bounds', () => {
incoming.clear(); incoming.clear();
}); });
test('gives up on a message the application never answers', () => { test('gives up on a message the application never answers', t => {
let now = 0; let now = 0;
const held = new HeldMessages({ log: silentLog, max: 10, maxOctets: 1_000_000, now: () => now, timeout: 60 }); const held = heldOn(t, { max: 10, maxOctets: 1_000_000, now: () => now, timeout: 60 });
offer(held, 1); offer(held, 1);
now = 61; now = 61;
@@ -1674,9 +1697,9 @@ describe('held message bounds', () => {
}); });
// Without this the drain sits out its whole budget before returning what a sweep already settled. // Without this the drain sits out its whole budget before returning what a sweep already settled.
test('wakes a waiting drain when the last message expires', async () => { test('wakes a waiting drain when the last message expires', async t => {
let now = 0; let now = 0;
const held = new HeldMessages({ log: silentLog, max: 10, maxOctets: 1_000_000, now: () => now, timeout: 60 }); const held = heldOn(t, { max: 10, maxOctets: 1_000_000, now: () => now, timeout: 60 });
offer(held, 1); offer(held, 1);
@@ -1701,14 +1724,11 @@ describe('sendResp()', () => {
session.sendReturn = () => Promise.resolve({ err: new Error('Socket is closed') }); session.sendReturn = () => Promise.resolve({ err: new Error('Socket is closed') });
const sms = createSms({ const sms = createSms({
from: '46701113311',
message: 'never answered',
pduObjs: [submitPdu(1)], pduObjs: [submitPdu(1)],
session, session,
to: '46709771337',
}, { }, {
answered: () => { answered++; },
lostLink: () => false, lostLink: () => false,
onAnswered: () => { answered++; },
send: () => Promise.resolve({ err: new Error('never sent') }), send: () => Promise.resolve({ err: new Error('never sent') }),
}); });
@@ -1726,14 +1746,11 @@ describe('sendDlr()', () => {
let call = 0; let call = 0;
const sms = createSms({ const sms = createSms({
from: '46701113311',
message: 'three segments',
pduObjs: [submitPdu(1), submitPdu(2), submitPdu(3)], pduObjs: [submitPdu(1), submitPdu(2), submitPdu(3)],
session, session,
to: '46709771337',
}, { }, {
answered: () => undefined,
lostLink: () => false, lostLink: () => false,
onAnswered: () => undefined,
send: () => { send: () => {
call++; call++;
+19 -10
View File
@@ -205,17 +205,14 @@ after #46, read 6, 6, 7 and 7, Locality 5, 5, 6 and 6. A fourth, after the link'
owner in #48, read 6, 6, 6 and 6, Locality 5 from every seat: all four still ranked `Session.teardown()` owner in #48, read 6, 6, 6 and 6, Locality 5 from every seat: all four still ranked `Session.teardown()`
hardest, and the held-message flow across `incoming-requests.ts`, `held-messages.ts`, `sms.ts` and hardest, and the held-message flow across `incoming-requests.ts`, `held-messages.ts`, `sms.ts` and
`Session`'s rejection handler second. `Session`'s rejection handler second.
A fifth, after the held-message flow got one owner, read 6, 6, 7 and 6, Locality 5, 5, 6 and 6:
three seats still ranked `MessageHold` hardest — six ways out, a rejection routed from `Session`
through `IncomingRequests` to a `WeakMap`, and the `setImmediate` a receipt relies on — and the
teardown cluster second; the inherited architect scored Shape 5 on the flat `src/` and the names
below.
- [ ] **Lift Locality to 7, and confirm it with a scoring run.** A run reading 7.0 or above also - [ ] **Lift Locality to 7, and confirm it with a scoring run.** A run reading 7.0 or above also
retires the #30, #46 and #48 decision. retires the Locality-first decision.
- [ ] **Give the held-message flow one owner — next, the condition #48 merged under.** Whether a
message is still held, and so whether its receipt may pass the drain, is decided across
`IncomingRequests.emitSms()`, `HeldMessages.offer()`/`MessageHold`, `createSms()`'s handlers in
`sms.ts` and `Session`'s rejection handler, which finds the hold again through a `WeakMap`
keyed on the `Sms`; `MessageHold.answered()` defers its release a `setImmediate` so a
`sendDlr()` straight after `sendResp()` still counts as held. All four seats of the #48 run
ranked it second hardest; the inherited architect put it at about two days.
### Correctness ### Correctness
@@ -325,10 +322,16 @@ hardest, and the held-message flow across `incoming-requests.ts`, `held-messages
- [ ] **Collapse the three objects named `defaults`.** `client.ts`, `server.ts` and - [ ] **Collapse the three objects named `defaults`.** `client.ts`, `server.ts` and
`session-options.ts` each export or hold one; `port: 2775` is written twice and the idle `session-options.ts` each export or hold one; `port: 2775` is written twice and the idle
timeout is derived two ways to the same 40 000. "What is the default for X" has three answers timeout is derived two ways to the same 40 000, and 64 MiB is both `defaultMaxOctets` and
`defaults.maxHeldOctets`. "What is the default for X" has three answers
depending on the entrypoint, and nothing fails when they drift. Named by both architects as the depending on the entrypoint, and nothing fails when they drift. Named by both architects as the
most likely first bug a new contributor ships. most likely first bug a new contributor ships.
- [ ] **Give `hold` one meaning, and rename `IncomingRequests.refusing` for what it does.**
`LinkLife.hold()` is a request's budget waiting for a link, `HeldMessages.hold()` a message the
application owes an answer; `refusing` decides no refusal — `held.full()` does — and only makes
the warn and info lines fire once each way. From the 2026-09-29 scoring run.
- [ ] **Rename `EncodingName`'s `ASCII` to `GSM7`, with `ASCII` a deprecated alias for one minor.** - [ ] **Rename `EncodingName`'s `ASCII` to `GSM7`, with `ASCII` a deprecated alias for one minor.**
It is GSM 03.38, where `$` is 0x02 and `@` is 0x00, and `segmentUnits.ASCII = 153` is a septet It is GSM 03.38, where `$` is 0x02 and `@` is 0x00, and `segmentUnits.ASCII = 153` is a septet
budget under a name that says octets. The 2026-09-09 decision removed `consts.ENCODING.ASCII` budget under a name that says octets. The 2026-09-09 decision removed `consts.ENCODING.ASCII`
@@ -349,6 +352,12 @@ hardest, and the held-message flow across `incoming-requests.ts`, `held-messages
### Self-sufficiency — 6–7 today, and the gate is 7 ### Self-sufficiency — 6–7 today, and the gate is 7
- [ ] **Move the fixture-copy reasoning out of AGENTS.md's Conventions, and split the longest
decision entries.** The fixtures bullet holds four justifications for tolerated copies — decisions,
so they belong in `docs/decisions.md` under Internals and tests; the entries under the wire's
alphabet and body rules run 30–40 lines with their `Rejected:` clauses inline, where a reader
who knows the answer still hunts for it. From the 2026-09-29 prose pass.
- [ ] **Move the one-line facts out of the decision log and back to the code.** Five of nine readers - [ ] **Move the one-line facts out of the decision log and back to the code.** Five of nine readers
independently reported being sent to `docs/decisions.md` for a question they hit while reading, independently reported being sent to `docs/decisions.md` for a question they hit while reading,
with no link from the code; one counted roughly fifty index redirects. The three left to inline with no link from the code; one counted roughly fifty index redirects. The three left to inline