Compare commits
1 Commits
de572803a0
...
57f53d56f1
| Author | SHA1 | Date | |
|---|---|---|---|
| 57f53d56f1 |
@@ -13,9 +13,7 @@ not for structure or style.
|
||||
## Goals
|
||||
|
||||
The goals, in priority order, live in
|
||||
[README.md](https://gitea.larvit.se/larvit/smpp-js/src/branch/main/README.md#goals) — they say where this library is heading, which an outside
|
||||
reader judges it by. The README states the audience alongside them.
|
||||
|
||||
[README.md](https://gitea.larvit.se/larvit/smpp-js/src/branch/main/README.md#goals). The README states the audience alongside them.
|
||||
|
||||
## Hard rules
|
||||
|
||||
@@ -55,12 +53,12 @@ src/
|
||||
held-messages.ts HeldMessages: capped, expiring messages the application has not answered, one MessageHold each
|
||||
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-gate.ts LinkGate: where a request with no link to go out on waits for the next one
|
||||
link-life.ts LinkLife: whether the link lives, and where a request waits for the next one
|
||||
link-timers.ts LinkTimers: the enquire_link heartbeat and the idle timeout
|
||||
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 gate, the window, the pending map and the retry
|
||||
outgoing-requests.ts OutgoingRequests: the window, the pending map and the retry
|
||||
pdu.ts pduToObj / objToPdu / pduReturn — synchronous, result-returning
|
||||
pdu-framer.ts PduFramer: a byte stream cut into complete PDUs
|
||||
pdu-refusal.ts A PDU the codec would not read, and the answer SMPP names for it
|
||||
@@ -180,7 +178,6 @@ decision under [The wire](docs/decisions.md#the-wire).
|
||||
collaborator the type system already keeps in step: `recordingDeps()` in `messaging-mode.test.ts`,
|
||||
`message-class.test.ts` and `unsendable.test.ts` is one `SendSmsDeps.send` that answers nothing,
|
||||
and a field added to that type fails to compile in every copy at once.
|
||||
- `message_id` values the library generates are UUID v7.
|
||||
- A socket a test opens and never reads must be `resume()`d, and a `data` listener counts. An unread
|
||||
socket never processes the peer's FIN, so `server.close()` hangs forever — that is a test bug, not
|
||||
a library one.
|
||||
@@ -314,13 +311,13 @@ the file.
|
||||
- A message id base is merged at most once.
|
||||
- A send that never reached the socket waits for the next link; one that did is counted, not resent.
|
||||
- A send queued for a send-window slot is bounded by the caller's `signal`, and by nothing else.
|
||||
- The gate decides whether a link can carry a request, and a bind is what makes it one.
|
||||
- One owner decides whether a link can carry a request, and a bind is what makes it one.
|
||||
|
||||
### [Internals and tests](docs/decisions.md#internals-and-tests)
|
||||
|
||||
- #30 and #46 merged under the comprehension floor, and Locality is the next work.
|
||||
- #30, #46 and #48 merged under the comprehension floor, and Locality is the next work.
|
||||
- A listener that rejects is routed by Node's `captureRejections`, not by hand-dispatching.
|
||||
- The four-line abort dance is copied across `LinkGate`, `IdleWaiters`, `PendingRequests` and
|
||||
- The four-line abort dance is copied across `LinkLife`, `IdleWaiters`, `PendingRequests` and
|
||||
`SendWindow` rather than extracted.
|
||||
- `SmppLog` is a five-method contract this library declares, not a dependency.
|
||||
- The TLS tests build their own self-signed certificate in DER
|
||||
|
||||
@@ -746,10 +746,7 @@ Who depends on this library, and what they may rely on.
|
||||
spooling, scheduling, retry policy and billing belong to whatever this is the edge of. State shared
|
||||
between instances is for goal 9's store, which has not shipped: today every session keeps its own,
|
||||
in memory.
|
||||
- **Pre-1.0, so the minor is the breaking unit** and a patch never breaks. What a 0.4.0 consumer has
|
||||
to change is in [MIGRATION.md](https://gitea.larvit.se/larvit/smpp-js/src/branch/main/MIGRATION.md);
|
||||
what each later minor changes is in
|
||||
[CHANGELOG.md](https://gitea.larvit.se/larvit/smpp-js/src/branch/main/CHANGELOG.md).
|
||||
- **Pre-1.0, so the minor is the breaking unit** and a patch never breaks.
|
||||
|
||||
Personas this README serves, in order:
|
||||
|
||||
|
||||
+18
-24
@@ -508,10 +508,7 @@ rule and an index of the titles below.
|
||||
|
||||
- **`close` means the session is over, and a drop the loop will retry is `disconnected`.**
|
||||
Maintainer's call, 2026-08-31: without the split, an application that opens a replacement client on
|
||||
`close` ends up holding two binds on one account. `teardown()` picks the event by whether the
|
||||
reconnect loop is still live, and `end()` stops that loop before tearing down, so every deliberate
|
||||
shutdown emits `close`. A retry that opens a socket and then loses it resets `lifecycle` through
|
||||
`attach()`, which is why a second drop emits again.
|
||||
`close` ends up holding two binds on one account, which goal 4 forbids.
|
||||
|
||||
- **An answer belongs to the link the message arrived on; a receipt does not.** Maintainer's call,
|
||||
2026-09-01. Rejected: answering on the new link, which succeeds and reports `{}` for a response
|
||||
@@ -541,9 +538,7 @@ rule and an index of the titles below.
|
||||
gives up on the operator whose provisioning lands a minute later; the backoff is what bounds the
|
||||
rate goal 4 cares about. The attempts before the first link report nothing, because the session
|
||||
running one has not reached the application: `disconnected` would have no listener and `close`
|
||||
would be a lie. Its wait is the one retry timer that is not `unref()`'d, for the reason
|
||||
`LinkGate`'s hold is not — it is awaited with no other handle, so a process whose only work is
|
||||
`client()` would exit unbound.
|
||||
would be a lie.
|
||||
|
||||
- **`connectTimeout` defaults to 10 s, bounds the whole connect including the TLS handshake, and
|
||||
`false` is the one way to turn it off.** Maintainer's call, 2026-09-20, serving goal 5: a connect
|
||||
@@ -676,7 +671,7 @@ rule and an index of the titles below.
|
||||
the peer, whose every request is bounded by `responseTimeout` unless the caller set that to 0 as
|
||||
well, and unsafe for the application, which nothing bounds — `close()` is what you reach for when
|
||||
the application is stuck, so it may not block on the application coming unstuck. That half falls
|
||||
back to `responseTimeout`, the same answer the link gate's hold already takes — and to that
|
||||
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
|
||||
@@ -738,10 +733,8 @@ rule and an index of the titles below.
|
||||
optional so every construction site answers. `UnansweredError` stays unexported: `unanswered` is
|
||||
the one spelling on the public surface. The hold is bounded by `responseTimeout` rather than an
|
||||
option of its own — that is already the answer to how long one request may wait — and its clock
|
||||
starts when the send is issued rather than when it first finds the gate shut, so one budget covers
|
||||
every hold a single call makes. That timer is the one here that is not `unref()`'d: a held request
|
||||
is awaited with the socket already destroyed, so an unref'd one lets a process whose only remaining
|
||||
work is that send exit without settling it.
|
||||
starts when the send is issued rather than when it first finds the link down, so one budget covers
|
||||
every hold a single call makes.
|
||||
|
||||
- **A send queued for a send-window slot is bounded by the caller's `signal`, and by nothing else.**
|
||||
Maintainer's call, 2026-09-06, from a review of PR #71: the hold above observes the signal and the
|
||||
@@ -757,29 +750,30 @@ rule and an index of the titles below.
|
||||
window is this end's own concurrency draining as the peer answers rather than a link going nowhere,
|
||||
and that bound would fail a message with more segments than `maxOutstanding` partway through
|
||||
against a slow peer. The failure is a plain `Error` rather than `UnansweredError`, the same answer
|
||||
an abort at the gate already gives. The drain half needs nothing: `close({ signal })` already hands the signal to
|
||||
an abort while held for a link already gives. The drain half needs nothing: `close({ signal })` already hands the signal to
|
||||
`window.idle()`, and `unbind()` taking none is the shape README states.
|
||||
|
||||
- **The gate decides whether a link can carry a request, and a bind is what makes it one.**
|
||||
Maintainer's call, 2026-09-01: `attach()` marks the session attached the moment a socket is
|
||||
handed over, one round trip before the bind is answered, so gating on that let a send arriving in
|
||||
that window go out unbound and come back `ESME_RINVBNDSTS`. The gate is told what happened and
|
||||
never reads back into the session: a collaborator that has to ask does not own its decision, which
|
||||
is how the first cut ended up answering the same question two different ways at admit and at
|
||||
release. The retry in `requestPastDrain()` asks `gate.awaitsNextLink()` rather than `canCarry()`,
|
||||
which also reads the socket: a loop condition the gate does not gate on spins against a gate that
|
||||
admits it straight back.
|
||||
- **One owner decides whether a link can carry a request, and a bind is what makes it one.**
|
||||
Maintainer's call, 2026-09-01, extended 2026-09-28; goal 1, since a send on a link not yet bound
|
||||
comes back `ESME_RINVBNDSTS`. `LinkLife` is told what happened and never reads back into the
|
||||
session; every other collaborator reads it and keeps no copy. Rejected: gating on the socket being
|
||||
attached, which admits a send one round trip before the bind is answered, and collaborators that
|
||||
ask the session, which answered the same question two ways at admit and at release.
|
||||
`ReconnectLoop.halted` is the one other flag, because `client()` also runs a loop with no session
|
||||
behind it for `fromStart`; a session's loop is stopped by `Session.stop()` alone.
|
||||
|
||||
|
||||
## Internals and tests
|
||||
|
||||
- **#30 and #46 merged under the comprehension floor, and Locality is the next work.** Maintainer's call,
|
||||
- **#30, #46 and #48 merged under the comprehension floor, and Locality is the next work.** 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
|
||||
seat capped by Locality in the held-message and shutdown code #30 does not touch, where the floor
|
||||
is 7.0. The chunks after #30 lift Locality to 7 before any other work. Serves goal 8's
|
||||
reshapeable internals, which a reader has to understand before reshaping. Valid until a scoring
|
||||
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
|
||||
emitters construct with `captureRejections: true` and implement
|
||||
@@ -790,7 +784,7 @@ rule and an index of the titles below.
|
||||
handlers normalise through `errorFrom()` rather than inline — a route out of the handler would land
|
||||
on a bare `process.nextTick` with nothing to catch it.
|
||||
|
||||
- **The four-line abort dance is copied across `LinkGate`, `IdleWaiters`, `PendingRequests` and
|
||||
- **The four-line abort dance is copied across `LinkLife`, `IdleWaiters`, `PendingRequests` and
|
||||
`SendWindow` rather than extracted.** Architecture review, 2026-09-06: pre-check `aborted`, attach
|
||||
`{ once: true }`, detach on settle, leave the registry. What differs at each site is the registry
|
||||
and what settling means — a FIFO handing over a slot, a set released together, a map keyed by
|
||||
|
||||
@@ -1,6 +1,7 @@
|
||||
import type { Concat } from './concat.ts';
|
||||
import type { DlrMerger } from './dlr-merger.ts';
|
||||
import type { ErrorName } from './defs/errors.ts';
|
||||
import type { LinkLife } from './link-life.ts';
|
||||
import type { LostGroup, Refusal } from './reassembly.ts';
|
||||
import type { OnRequest } from './session-options.ts';
|
||||
import type { PduObject, PduObjectInput } from './pdu.ts';
|
||||
@@ -45,6 +46,7 @@ const lostReasons: Record<LostGroup['reason'], string> = {
|
||||
|
||||
export type IncomingRequestsOptions = {
|
||||
dlrMerger: DlrMerger;
|
||||
link: LinkLife;
|
||||
log: SmppLog;
|
||||
maxOctets?: number | undefined;
|
||||
maxReassembly?: number | undefined;
|
||||
@@ -61,6 +63,7 @@ export type IncomingRequestsOptions = {
|
||||
export class IncomingRequests {
|
||||
private readonly dlrMerger: DlrMerger;
|
||||
private readonly held: HeldMessages;
|
||||
private readonly link: LinkLife;
|
||||
private readonly log: SmppLog;
|
||||
private readonly onRequest: OnRequest | undefined;
|
||||
private readonly reassembler: Reassembler;
|
||||
@@ -68,7 +71,6 @@ export class IncomingRequests {
|
||||
private readonly session: Session;
|
||||
private readonly smsIdFormat: SmsIdFormat;
|
||||
private readonly systemId: string;
|
||||
private linkGeneration = 0;
|
||||
private refusing = false;
|
||||
|
||||
constructor(options: IncomingRequestsOptions) {
|
||||
@@ -79,6 +81,7 @@ export class IncomingRequests {
|
||||
maxOctets: defaults.maxHeldOctets,
|
||||
timeout: defaults.heldMessageTimeout,
|
||||
});
|
||||
this.link = options.link;
|
||||
this.log = options.log;
|
||||
this.onRequest = options.onRequest;
|
||||
this.reassembler = new Reassembler({
|
||||
@@ -95,14 +98,14 @@ export class IncomingRequests {
|
||||
}
|
||||
|
||||
async handle(pduObj: PduObject): Promise<void> {
|
||||
const generation = this.linkGeneration;
|
||||
const generation = this.link.generation();
|
||||
const { onRequest } = this;
|
||||
|
||||
// Called unbound, so the application's hook never sees this class as its `this`.
|
||||
if (onRequest && await onRequest(this.session, pduObj)) return;
|
||||
|
||||
// The link it arrived on went while the hook ran, so nothing we answer now correlates.
|
||||
if (this.linkGeneration !== generation) {
|
||||
if (this.link.generation() !== generation) {
|
||||
this.log.info('session - dropping a request whose link went', { cmdName: pduObj.cmdName });
|
||||
|
||||
return;
|
||||
@@ -148,7 +151,6 @@ export class IncomingRequests {
|
||||
|
||||
/** Drops the segments of every message that never became whole, and of every one still held. */
|
||||
clear(): void {
|
||||
this.linkGeneration++;
|
||||
this.refusing = false;
|
||||
this.held.clear();
|
||||
this.reassembler.clear();
|
||||
@@ -288,7 +290,7 @@ export class IncomingRequests {
|
||||
|
||||
if (!first) return;
|
||||
|
||||
const generation = this.linkGeneration;
|
||||
const generation = this.link.generation();
|
||||
|
||||
this.held.offer(pduObjs, this.session.listenerCount('sms'), hold => createSms({
|
||||
answeredAs,
|
||||
@@ -298,7 +300,7 @@ export class IncomingRequests {
|
||||
session: this.session,
|
||||
to: paramText(first.params.destination_addr),
|
||||
}, {
|
||||
lostLink: () => this.linkGeneration !== generation,
|
||||
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,13 +1,18 @@
|
||||
import type { SmppLog } from './log.ts';
|
||||
import type { VoidResult } from './result.ts';
|
||||
|
||||
export type LinkGateOptions = {
|
||||
export type LinkLifeOptions = {
|
||||
log: SmppLog;
|
||||
now?: (() => number) | undefined;
|
||||
/** Whether a dropped link is followed by another one until stop(). */
|
||||
reconnects: boolean;
|
||||
/** How long a request may wait for a link. 0 waits for as long as one may still arrive. */
|
||||
timeout: number;
|
||||
};
|
||||
|
||||
/** `binding`: a socket is attached and its bind is not answered yet, so it carries nothing but that bind. */
|
||||
type Phase = 'binding' | 'down' | 'ended' | 'up';
|
||||
|
||||
type Waiter = (result: VoidResult) => void;
|
||||
|
||||
function aborted(): Error {
|
||||
@@ -22,66 +27,116 @@ function over(): Error {
|
||||
return new Error('Session is closed');
|
||||
}
|
||||
|
||||
/** Where a request with no link to go out on waits for the next one. */
|
||||
export class LinkGate {
|
||||
/** Whether the session's link lives, and where a request with no link to go out on waits for the next one. */
|
||||
export class LinkLife {
|
||||
private readonly log: SmppLog;
|
||||
private readonly now: () => number;
|
||||
private readonly reconnects: boolean;
|
||||
private readonly timeout: number;
|
||||
private readonly waiting = new Set<Waiter>();
|
||||
private returning = false;
|
||||
private up = true;
|
||||
private drops = 0;
|
||||
private phase: Phase = 'up';
|
||||
private stopped = false;
|
||||
|
||||
constructor(options: LinkGateOptions) {
|
||||
constructor(options: LinkLifeOptions) {
|
||||
this.log = options.log;
|
||||
this.now = options.now ?? Date.now;
|
||||
this.reconnects = options.reconnects;
|
||||
this.timeout = options.timeout;
|
||||
}
|
||||
|
||||
/** Whether a request can go out right now. A link that is attached but not yet bound cannot. */
|
||||
/** A socket is on the link, bound or not. */
|
||||
isAttached(): boolean {
|
||||
return this.phase === 'binding' || this.phase === 'up';
|
||||
}
|
||||
|
||||
/** Whether a request can go out right now. */
|
||||
isUp(): boolean {
|
||||
return this.up;
|
||||
return this.phase === 'up';
|
||||
}
|
||||
|
||||
/** Shut, with another link on its way to reopen it. */
|
||||
private isOver(): boolean {
|
||||
return this.phase === 'ended';
|
||||
}
|
||||
|
||||
/** The session is shutting down: nothing new is taken, and no link follows this one. */
|
||||
isStopped(): boolean {
|
||||
return this.stopped;
|
||||
}
|
||||
|
||||
/** Whether a link that drops now is followed by another. */
|
||||
retrying(): boolean {
|
||||
return this.reconnects && !this.stopped;
|
||||
}
|
||||
|
||||
/** Not up and not over, with a link to come. */
|
||||
awaitsNextLink(): boolean {
|
||||
return !this.up && this.returning;
|
||||
return !this.isUp() && !this.isOver() && this.retrying();
|
||||
}
|
||||
|
||||
/** Why the gate will never admit a request, or undefined while one may still get through. */
|
||||
/** Changes with every drop, so what was read off one link can tell that link is gone. */
|
||||
generation(): number {
|
||||
return this.drops;
|
||||
}
|
||||
|
||||
/** Why no request will ever be admitted, or undefined while one may still get through. */
|
||||
refusal(): Error | undefined {
|
||||
return this.up || this.returning ? undefined : over();
|
||||
return this.isUp() || this.awaitsNextLink() ? undefined : over();
|
||||
}
|
||||
|
||||
/** One budget for a request, however many links it waits through. 0 never gives up. */
|
||||
/** One budget for a request, however many links it waits through. */
|
||||
hold(signal: AbortSignal | undefined): () => Promise<VoidResult> {
|
||||
const deadline = this.timeout > 0 ? this.now() + this.timeout : 0;
|
||||
|
||||
return () => this.wait(deadline, signal);
|
||||
}
|
||||
|
||||
/** A link is up and bound: everything held goes out on it. */
|
||||
/** A socket from the reconnect loop, not yet bound. An ended session stays ended. */
|
||||
attach(): void {
|
||||
if (this.isOver()) return;
|
||||
|
||||
this.phase = 'binding';
|
||||
}
|
||||
|
||||
/** The link is bound: everything held goes out on it. */
|
||||
open(): void {
|
||||
this.up = true;
|
||||
this.returning = false;
|
||||
this.phase = 'up';
|
||||
|
||||
if (this.waiting.size > 0) {
|
||||
this.log.verbose('linkGate - sending what was held for a link', { held: this.waiting.size });
|
||||
this.log.verbose('linkLife - sending what was held for a link', { held: this.waiting.size });
|
||||
}
|
||||
|
||||
this.release({});
|
||||
}
|
||||
|
||||
/** The link is gone. `returning` says whether another one is on its way. */
|
||||
shut(returning: boolean): void {
|
||||
this.up = false;
|
||||
this.returning = returning;
|
||||
/** The attached link is gone: the event that says so, or undefined when there was none to lose. */
|
||||
drop(): 'close' | 'disconnected' | undefined {
|
||||
if (!this.isAttached()) return undefined;
|
||||
|
||||
if (!returning) this.release({ err: over() });
|
||||
this.phase = 'down';
|
||||
this.drops++;
|
||||
|
||||
return this.retrying() ? 'disconnected' : 'close';
|
||||
}
|
||||
|
||||
stop(): void {
|
||||
this.stopped = true;
|
||||
}
|
||||
|
||||
/** The session is over: nothing held will ever go out. False means it already was. */
|
||||
end(): boolean {
|
||||
if (this.isOver()) return false;
|
||||
|
||||
this.phase = 'ended';
|
||||
this.stopped = true;
|
||||
this.release({ err: over() });
|
||||
|
||||
return true;
|
||||
}
|
||||
|
||||
/** Resolves once a link can carry the request, or with the reason none ever will. */
|
||||
private wait(deadline: number, signal: AbortSignal | undefined): Promise<VoidResult> {
|
||||
if (this.up) return Promise.resolve({});
|
||||
if (this.isUp()) return Promise.resolve({});
|
||||
|
||||
const refused = this.refusal();
|
||||
|
||||
@@ -97,7 +152,7 @@ export class LinkGate {
|
||||
}
|
||||
|
||||
private waitForLink(left: number, signal: AbortSignal | undefined): Promise<VoidResult> {
|
||||
this.log.verbose('linkGate - holding a request until a link is back', { timeout: left });
|
||||
this.log.verbose('linkLife - holding a request until a link is back', { timeout: left });
|
||||
|
||||
return new Promise<VoidResult>(resolve => {
|
||||
let timer: NodeJS.Timeout | undefined = undefined;
|
||||
@@ -109,7 +164,7 @@ export class LinkGate {
|
||||
resolve(result);
|
||||
};
|
||||
const giveUp = (): void => {
|
||||
this.log.warn('linkGate - no link came back in time', { timeout: left });
|
||||
this.log.warn('linkLife - no link came back in time', { timeout: left });
|
||||
settle({ err: expired() });
|
||||
};
|
||||
|
||||
+17
-35
@@ -1,9 +1,9 @@
|
||||
import type { LinkLife } from './link-life.ts';
|
||||
import type { PduObject, PduObjectInput } from './pdu.ts';
|
||||
import type { PduTransport } from './pdu-transport.ts';
|
||||
import type { Result, VoidResult } from './result.ts';
|
||||
import type { SendOptions } from './session-options.ts';
|
||||
import type { SmppLog } from './log.ts';
|
||||
import { LinkGate } from './link-gate.ts';
|
||||
import { PendingRequests } from './pending-requests.ts';
|
||||
import { SendWindow } from './send-window.ts';
|
||||
import { UnansweredError } from './unanswered-error.ts';
|
||||
@@ -11,6 +11,7 @@ import { bindCommands } from './session-options.ts';
|
||||
import { objToPdu } from './pdu.ts';
|
||||
|
||||
export type OutgoingRequestsOptions = {
|
||||
link: LinkLife;
|
||||
log: SmppLog;
|
||||
maxOutstanding: number;
|
||||
responseTimeout: number;
|
||||
@@ -33,17 +34,15 @@ function misuse(input: PduObjectInput): Error | undefined {
|
||||
|
||||
/** Everything this end asks of the peer: which link carries it, how many at once, and the answer. */
|
||||
export class OutgoingRequests {
|
||||
private readonly gate: LinkGate;
|
||||
private readonly link: LinkLife;
|
||||
private readonly log: SmppLog;
|
||||
private readonly pending: PendingRequests;
|
||||
private readonly responseTimeout: number;
|
||||
private readonly transport: PduTransport;
|
||||
private readonly window: SendWindow;
|
||||
|
||||
private draining = false;
|
||||
|
||||
constructor(options: OutgoingRequestsOptions) {
|
||||
this.gate = new LinkGate({ log: options.log, timeout: options.responseTimeout });
|
||||
this.link = options.link;
|
||||
this.log = options.log;
|
||||
this.pending = new PendingRequests(options.log);
|
||||
this.responseTimeout = options.responseTimeout;
|
||||
@@ -52,22 +51,11 @@ export class OutgoingRequests {
|
||||
}
|
||||
|
||||
canCarry(): boolean {
|
||||
return this.gate.isUp() && !this.transport.sock.destroyed;
|
||||
return this.link.isUp() && !this.transport.sock.destroyed;
|
||||
}
|
||||
|
||||
/** The link went before the drain finished, so an empty window says nothing about the peer. */
|
||||
droppedWhileDraining(): boolean {
|
||||
return this.draining && !this.canCarry();
|
||||
}
|
||||
|
||||
/** A link is up and bound, so everything held for one goes out on it. */
|
||||
linkUp(): void {
|
||||
this.gate.open();
|
||||
}
|
||||
|
||||
/** The link is gone; `returning` says whether another one is on its way. */
|
||||
linkLost(returning: boolean): void {
|
||||
this.gate.shut(returning);
|
||||
/** The link is gone, and every answer still owed on it with it. */
|
||||
linkLost(): void {
|
||||
this.pending.settleAll(new Error('Session closed before a response arrived'));
|
||||
}
|
||||
|
||||
@@ -81,7 +69,6 @@ export class OutgoingRequests {
|
||||
this.pending.settle(seqNr, { err });
|
||||
}
|
||||
|
||||
/** Sends a request and resolves with the peer's response. */
|
||||
request(input: PduObjectInput, options: SendOptions): Promise<Result<{ pduObj: PduObject }>> {
|
||||
// Ahead of the drain, so a misuse is named as one rather than blamed on the shutdown.
|
||||
const wrong = misuse(input);
|
||||
@@ -89,7 +76,7 @@ export class OutgoingRequests {
|
||||
if (wrong) return Promise.resolve({ err: wrong });
|
||||
|
||||
// With no link, the request is refused as closed further on.
|
||||
if (this.draining && this.canCarry()) {
|
||||
if (this.link.isStopped() && this.canCarry()) {
|
||||
return Promise.resolve({ err: new Error('Session is shutting down') });
|
||||
}
|
||||
|
||||
@@ -105,14 +92,14 @@ export class OutgoingRequests {
|
||||
|
||||
if (refused) return { err: refused };
|
||||
|
||||
// A bind is what makes a link usable, so it cannot wait for one: it takes the gate's answer now.
|
||||
// A bind is what makes a link usable, so it cannot wait for one: it takes the link's answer now.
|
||||
if (bindCommands.includes(input.cmdName)) {
|
||||
const shut = this.gate.refusal();
|
||||
const shut = this.link.refusal();
|
||||
|
||||
return shut ? { err: shut } : this.requestPastDrainGateAndWindow(input, options);
|
||||
return shut ? { err: shut } : this.requestOnCurrentLink(input, options);
|
||||
}
|
||||
|
||||
const waitForLink = this.gate.hold(options.signal);
|
||||
const waitForLink = this.link.hold(options.signal);
|
||||
|
||||
for (;;) {
|
||||
const held = await waitForLink();
|
||||
@@ -130,18 +117,13 @@ export class OutgoingRequests {
|
||||
}
|
||||
|
||||
/** Straight onto the current link, for what has to go out either way. */
|
||||
async requestPastDrainGateAndWindow(
|
||||
async requestOnCurrentLink(
|
||||
input: PduObjectInput,
|
||||
options: SendOptions = {},
|
||||
): Promise<Result<{ pduObj: PduObject }>> {
|
||||
return (await this.attempt(input, options)).result;
|
||||
}
|
||||
|
||||
/** Refuses every request from here on, on a link that is already down as much as a live one. */
|
||||
stopAccepting(): void {
|
||||
this.draining = true;
|
||||
}
|
||||
|
||||
/** Waits out the requests already on the wire, and says how many never finished. */
|
||||
async drain(timeout: number, signal: AbortSignal | undefined): Promise<VoidResult> {
|
||||
const unfinished = await this.window.idle(timeout, signal);
|
||||
@@ -153,15 +135,15 @@ export class OutgoingRequests {
|
||||
return { err: new Error(`Shut down with ${String(unfinished)} request(s) unfinished`) };
|
||||
}
|
||||
|
||||
/** Nothing reached the socket, so the next link carries it instead of the caller resending. */
|
||||
/** Nothing reached the socket, so the next link carries it. */
|
||||
private retriesOnNextLink(attempt: Attempt): boolean {
|
||||
// Until the gate is shut it admits the retry straight back onto the dead socket, and the loop spins.
|
||||
return attempt.retryOnNextLink && this.gate.awaitsNextLink();
|
||||
// Until the link is dropped it admits the retry straight back onto the dead socket, and the loop spins.
|
||||
return attempt.retryOnNextLink && this.link.awaitsNextLink();
|
||||
}
|
||||
|
||||
/** Why a request cannot go out at all, as opposed to not yet. */
|
||||
private refuse(input: PduObjectInput, options: SendOptions): Error | undefined {
|
||||
// Before the gate and the window, or an aborted call waits for what it will never use.
|
||||
// Before the link and the window, or an aborted call waits for what it will never use.
|
||||
return misuse(input) ?? (options.signal?.aborted === true ? abortedBeforeSend() : undefined);
|
||||
}
|
||||
|
||||
|
||||
@@ -40,7 +40,7 @@ export class ReconnectLoop {
|
||||
}
|
||||
|
||||
/** Read through a method: stop() can land while an attempt is awaiting. */
|
||||
isStopped(): boolean {
|
||||
private isStopped(): boolean {
|
||||
return this.halted;
|
||||
}
|
||||
|
||||
|
||||
+40
-41
@@ -11,6 +11,7 @@ import type { Socket } from 'node:net';
|
||||
import { DlrMerger } from './dlr-merger.ts';
|
||||
import { EventEmitter } from 'node:events';
|
||||
import { IncomingRequests } from './incoming-requests.ts';
|
||||
import { LinkLife } from './link-life.ts';
|
||||
import { LinkTimers } from './link-timers.ts';
|
||||
import { OutgoingRequests } from './outgoing-requests.ts';
|
||||
import { PduTransport } from './pdu-transport.ts';
|
||||
@@ -61,15 +62,13 @@ export class Session extends EventEmitter<SessionEvents> {
|
||||
private readonly concatReference = new ConcatReference();
|
||||
private readonly dlrMerger: DlrMerger;
|
||||
private readonly incoming: IncomingRequests;
|
||||
private readonly link: LinkLife;
|
||||
private readonly options: SessionOptions;
|
||||
private readonly outgoing: OutgoingRequests;
|
||||
private readonly reconnectLoop: ReconnectLoop | undefined;
|
||||
private readonly timers: LinkTimers;
|
||||
private readonly transport: PduTransport;
|
||||
|
||||
/** `ended` is final: end() stops the reconnect loop before any attach() can run. */
|
||||
private lifecycle: 'attached' | 'ended' | 'torn-down' = 'attached';
|
||||
|
||||
/** A listener that throws is the application's bug; it must not become ours. Hard rule 1. */
|
||||
override emit<K extends keyof SessionEvents>(
|
||||
event: K,
|
||||
@@ -109,12 +108,12 @@ export class Session extends EventEmitter<SessionEvents> {
|
||||
|
||||
this.log = guardedLog(options.log);
|
||||
this.options = options;
|
||||
this.dlrMerger = new DlrMerger({
|
||||
log: this.log,
|
||||
max: defaults.maxDlrMerges,
|
||||
timeout: defaults.dlrMergeTimeout,
|
||||
});
|
||||
this.dlrMerger = new DlrMerger({ log: this.log, max: defaults.maxDlrMerges, timeout: defaults.dlrMergeTimeout });
|
||||
this.reconnectLoop = this.loopFor(options.reconnect);
|
||||
|
||||
const responseTimeout = options.responseTimeout ?? defaults.responseTimeout;
|
||||
|
||||
this.link = new LinkLife({ log: this.log, reconnects: this.reconnectLoop !== undefined, timeout: responseTimeout });
|
||||
this.timers = new LinkTimers({
|
||||
enquireLinkInterval: options.enquireLinkInterval,
|
||||
idleTimeout: options.idleTimeout,
|
||||
@@ -125,13 +124,15 @@ export class Session extends EventEmitter<SessionEvents> {
|
||||
});
|
||||
this.transport = this.transportFor(options.sock);
|
||||
this.outgoing = new OutgoingRequests({
|
||||
link: this.link,
|
||||
log: this.log,
|
||||
maxOutstanding: options.maxOutstanding ?? defaults.maxOutstanding,
|
||||
responseTimeout: options.responseTimeout ?? defaults.responseTimeout,
|
||||
responseTimeout,
|
||||
transport: this.transport,
|
||||
});
|
||||
this.incoming = new IncomingRequests({
|
||||
dlrMerger: this.dlrMerger,
|
||||
link: this.link,
|
||||
log: this.log,
|
||||
maxOctets: options.maxOctets,
|
||||
maxReassembly: options.maxReassembly,
|
||||
@@ -199,7 +200,7 @@ export class Session extends EventEmitter<SessionEvents> {
|
||||
const sent = built.err ? { err: built.err } : this.transport.write(built.buffer);
|
||||
|
||||
// A peer that unbinds and drops the link takes our response with it; that is not a failure.
|
||||
if (sent.err && this.lifecycle === 'attached') {
|
||||
if (sent.err && this.link.isAttached()) {
|
||||
this.log.warn('session - could not answer a request', {
|
||||
cmdName,
|
||||
message: sent.err.message,
|
||||
@@ -234,11 +235,11 @@ export class Session extends EventEmitter<SessionEvents> {
|
||||
*/
|
||||
async unbind(): Promise<VoidResult> {
|
||||
const drained = await this.drain(undefined);
|
||||
const wasOpen = this.lifecycle === 'attached';
|
||||
const wasOpen = this.link.isAttached();
|
||||
const sent = wasOpen
|
||||
? await this.outgoing.requestPastDrainGateAndWindow({ cmdName: 'unbind' })
|
||||
? await this.outgoing.requestOnCurrentLink({ cmdName: 'unbind' })
|
||||
: { err: new Error('Session is closed') };
|
||||
const closedOnUnbind = wasOpen && this.lifecycle !== 'attached';
|
||||
const closedOnUnbind = wasOpen && !this.link.isAttached();
|
||||
|
||||
this.end();
|
||||
|
||||
@@ -246,8 +247,8 @@ export class Session extends EventEmitter<SessionEvents> {
|
||||
}
|
||||
|
||||
/**
|
||||
* Closes for good: refuses new sends, waits out the requests already on the wire up to
|
||||
* `shutdownTimeout`, then tears down whatever is left. A session closed this way never reconnects.
|
||||
* Closes for good: refuses new sends, waits up to `shutdownTimeout` for the requests already sent
|
||||
* and the messages not yet answered, then tears down whatever is left. A session closed this way never reconnects.
|
||||
*/
|
||||
async close(options: CloseOptions = {}): Promise<VoidResult> {
|
||||
const drained = await this.drain(options.signal);
|
||||
@@ -300,14 +301,14 @@ export class Session extends EventEmitter<SessionEvents> {
|
||||
}
|
||||
|
||||
// close() can land while the rebind is in flight.
|
||||
if (this.reconnectLoop?.isStopped() === true) {
|
||||
if (!this.link.retrying()) {
|
||||
this.teardown();
|
||||
|
||||
return { err: new Error('Session closed while it was coming back up') };
|
||||
}
|
||||
|
||||
this.resetTimers();
|
||||
this.outgoing.linkUp();
|
||||
this.link.open();
|
||||
this.log.info('session - reconnected');
|
||||
this.emit('reconnected');
|
||||
|
||||
@@ -316,15 +317,14 @@ export class Session extends EventEmitter<SessionEvents> {
|
||||
|
||||
private attach(sock: Socket): void {
|
||||
this.transport.attach(sock);
|
||||
this.lifecycle = 'attached';
|
||||
this.link.attach();
|
||||
}
|
||||
|
||||
/** Stops new sends and waits out the messages we hold and the requests already issued. */
|
||||
private async drain(signal: AbortSignal | undefined): Promise<VoidResult> {
|
||||
this.reconnectLoop?.stop();
|
||||
this.outgoing.stopAccepting();
|
||||
this.stop();
|
||||
|
||||
// No link, so nothing is on the wire to wait out.
|
||||
// No bound link, so nothing is on the wire to wait out.
|
||||
if (!this.outgoing.canCarry()) return {};
|
||||
|
||||
const timeout = this.options.shutdownTimeout ?? defaults.shutdownTimeout;
|
||||
@@ -333,7 +333,8 @@ export class Session extends EventEmitter<SessionEvents> {
|
||||
const messages = await this.incoming.drain(this.answering(timeout), signal);
|
||||
const requests = await this.outgoing.drain(leftOf(deadline), signal);
|
||||
|
||||
if (this.outgoing.droppedWhileDraining()) {
|
||||
// The link went before the drain finished, so an empty window says nothing about the peer.
|
||||
if (!this.outgoing.canCarry()) {
|
||||
return { err: new Error('The session closed before the drain finished') };
|
||||
}
|
||||
|
||||
@@ -355,42 +356,40 @@ export class Session extends EventEmitter<SessionEvents> {
|
||||
|
||||
/** The session is over now, drained or not. Nothing brings it back. */
|
||||
private end(): void {
|
||||
this.reconnectLoop?.stop();
|
||||
this.stop();
|
||||
this.teardown();
|
||||
this.dlrMerger.clear();
|
||||
this.emitClose();
|
||||
}
|
||||
|
||||
private emitClose(): void {
|
||||
if (this.lifecycle === 'ended') return;
|
||||
/** No new sends, and no link after this one. */
|
||||
private stop(): void {
|
||||
this.link.stop();
|
||||
this.reconnectLoop?.stop();
|
||||
}
|
||||
|
||||
this.lifecycle = 'ended';
|
||||
this.outgoing.linkLost(false);
|
||||
private emitClose(): void {
|
||||
if (!this.link.end()) return;
|
||||
|
||||
this.outgoing.linkLost();
|
||||
this.emit('close');
|
||||
}
|
||||
|
||||
private teardown(): void {
|
||||
if (this.lifecycle !== 'attached') return;
|
||||
const lost = this.link.drop();
|
||||
|
||||
this.lifecycle = 'torn-down';
|
||||
if (!lost) return;
|
||||
|
||||
// Read once: clear() reports lost segments, and a listener could stop the loop between reads.
|
||||
const retrying = this.retrying();
|
||||
|
||||
this.outgoing.linkLost(retrying);
|
||||
this.outgoing.linkLost();
|
||||
this.timers.clear();
|
||||
this.incoming.clear();
|
||||
this.sock.destroy();
|
||||
|
||||
if (retrying) this.emit('disconnected');
|
||||
// `lost` is read before clear(): a listener it reaches may close() the session, and the drop still reports as disconnected.
|
||||
if (lost === 'disconnected') this.emit('disconnected');
|
||||
else this.emitClose();
|
||||
}
|
||||
|
||||
// Copied into the gate at teardown, so stopping the loop anywhere but drain() and end() has to shut the gate too.
|
||||
private retrying(): boolean {
|
||||
return this.reconnectLoop !== undefined && !this.reconnectLoop.isStopped();
|
||||
}
|
||||
|
||||
private onData(chunk: Buffer): void {
|
||||
this.emit('data', chunk);
|
||||
this.resetTimers();
|
||||
@@ -436,13 +435,13 @@ export class Session extends EventEmitter<SessionEvents> {
|
||||
}
|
||||
|
||||
private resetTimers(): void {
|
||||
if (this.lifecycle !== 'attached') return;
|
||||
if (!this.link.isAttached()) return;
|
||||
|
||||
this.timers.reset();
|
||||
}
|
||||
|
||||
private onClose(): void {
|
||||
if (this.retrying()) {
|
||||
if (this.link.retrying()) {
|
||||
this.teardown();
|
||||
this.reconnectLoop?.schedule();
|
||||
|
||||
|
||||
+93
-23
@@ -19,7 +19,7 @@ import { HeldMessages } from '../src/held-messages.ts';
|
||||
import { IncomingRequests, refusedSegmentStatus } from '../src/incoming-requests.ts';
|
||||
import { UnansweredError } from '../src/unanswered-error.ts';
|
||||
import { createSms } from '../src/sms.ts';
|
||||
import { LinkGate } from '../src/link-gate.ts';
|
||||
import { LinkLife } from '../src/link-life.ts';
|
||||
import { SendWindow } from '../src/send-window.ts';
|
||||
import { Reassembler, decodeSegments } from '../src/reassembly.ts';
|
||||
import { Session } from '../src/session.ts';
|
||||
@@ -104,6 +104,7 @@ function abortAfter(
|
||||
function incomingOn(session: Session, options: Partial<IncomingRequestsOptions> = {}): IncomingRequests {
|
||||
return new IncomingRequests({
|
||||
dlrMerger: new DlrMerger({ log: silentLog, max: 10, timeout: 10_000 }),
|
||||
link: new LinkLife({ log: silentLog, reconnects: false, timeout: 100 }),
|
||||
log: silentLog,
|
||||
sendPastDrain: () => Promise.resolve({ err: new Error('never sent') }),
|
||||
session,
|
||||
@@ -748,14 +749,15 @@ describe('reconnect', () => {
|
||||
|
||||
closeAfter(t, session);
|
||||
|
||||
const incoming = incomingOn(session, { onRequest: async () => { await delay(10); return false; } });
|
||||
const link = new LinkLife({ log: silentLog, reconnects: true, timeout: 100 });
|
||||
const incoming = incomingOn(session, { link, onRequest: async () => { await delay(10); return false; } });
|
||||
let messages = 0;
|
||||
|
||||
session.on('sms', () => { messages++; });
|
||||
|
||||
const handled = incoming.handle(submitPdu(1));
|
||||
|
||||
incoming.clear();
|
||||
link.drop();
|
||||
|
||||
await handled;
|
||||
|
||||
@@ -1401,13 +1403,13 @@ describe('sends across a reconnect', () => {
|
||||
});
|
||||
});
|
||||
|
||||
describe('LinkGate', () => {
|
||||
describe('LinkLife', () => {
|
||||
test('refuses a hold whose deadline has already passed', async () => {
|
||||
let now = 0;
|
||||
const gate = new LinkGate({ log: silentLog, now: () => now, timeout: 100 });
|
||||
const waitForLink = gate.hold(undefined);
|
||||
const link = new LinkLife({ log: silentLog, now: () => now, reconnects: true, timeout: 100 });
|
||||
const waitForLink = link.hold(undefined);
|
||||
|
||||
gate.shut(true);
|
||||
link.drop();
|
||||
now = 101;
|
||||
|
||||
const held = await waitForLink();
|
||||
@@ -1416,42 +1418,77 @@ describe('LinkGate', () => {
|
||||
});
|
||||
|
||||
test('holds on a timer that keeps the process alive', async () => {
|
||||
const gate = new LinkGate({ log: silentLog, timeout: 10_000 });
|
||||
const link = new LinkLife({ log: silentLog, reconnects: true, timeout: 10_000 });
|
||||
const timers = (): number => process.getActiveResourcesInfo().filter(name => name === 'Timeout').length;
|
||||
|
||||
gate.shut(true);
|
||||
link.drop();
|
||||
|
||||
const before = timers();
|
||||
const held = gate.hold(undefined)();
|
||||
const held = link.hold(undefined)();
|
||||
|
||||
assert.equal(timers(), before + 1, 'an unref\'d timer is not counted here, which is the point');
|
||||
|
||||
gate.open();
|
||||
link.open();
|
||||
|
||||
assert.deepEqual(await held, {});
|
||||
});
|
||||
|
||||
// addEventListener never fires for a signal that already aborted, so it would wait out the timeout.
|
||||
test('gives up at once on a signal that was already aborted', async () => {
|
||||
const gate = new LinkGate({ log: silentLog, timeout: 100 });
|
||||
const link = new LinkLife({ log: silentLog, reconnects: true, timeout: 100 });
|
||||
|
||||
gate.shut(true);
|
||||
link.drop();
|
||||
|
||||
const held = await gate.hold(AbortSignal.abort())();
|
||||
const held = await link.hold(AbortSignal.abort())();
|
||||
|
||||
assert.match(held.err?.message ?? '', /Aborted while waiting for a link/);
|
||||
});
|
||||
|
||||
test('awaits the next link only while shut with one on its way', () => {
|
||||
const gate = new LinkGate({ log: silentLog, timeout: 100 });
|
||||
test('awaits the next link only while down with one on its way', () => {
|
||||
const link = new LinkLife({ log: silentLog, reconnects: true, timeout: 100 });
|
||||
|
||||
assert.equal(gate.awaitsNextLink(), false, 'up');
|
||||
gate.shut(true);
|
||||
assert.equal(gate.awaitsNextLink(), true, 'shut, returning');
|
||||
gate.open();
|
||||
assert.equal(gate.awaitsNextLink(), false, 'reopened');
|
||||
gate.shut(false);
|
||||
assert.equal(gate.awaitsNextLink(), false, 'shut for good');
|
||||
assert.equal(link.awaitsNextLink(), false, 'up');
|
||||
link.drop();
|
||||
assert.equal(link.awaitsNextLink(), true, 'down, returning');
|
||||
link.attach();
|
||||
assert.equal(link.awaitsNextLink(), true, 'attached, not yet bound');
|
||||
link.open();
|
||||
assert.equal(link.awaitsNextLink(), false, 'reopened');
|
||||
link.drop();
|
||||
link.stop();
|
||||
assert.equal(link.awaitsNextLink(), false, 'down, stopped');
|
||||
assert.match(link.refusal()?.message ?? '', /closed/, 'stopped while down');
|
||||
link.end();
|
||||
assert.equal(link.awaitsNextLink(), false, 'ended');
|
||||
});
|
||||
|
||||
test('drops an attached link once, counts each drop, and names the event it warrants', () => {
|
||||
const link = new LinkLife({ log: silentLog, reconnects: true, timeout: 100 });
|
||||
const generation = link.generation();
|
||||
|
||||
assert.equal(link.drop(), 'disconnected');
|
||||
assert.equal(link.drop(), undefined, 'already down');
|
||||
assert.equal(link.generation(), generation + 1);
|
||||
link.attach();
|
||||
link.stop();
|
||||
assert.equal(link.drop(), 'close', 'a new link drops again, with none to follow it');
|
||||
assert.equal(link.generation(), generation + 2);
|
||||
assert.equal(new LinkLife({ log: silentLog, reconnects: false, timeout: 100 }).drop(), 'close');
|
||||
});
|
||||
|
||||
test('releases a held request with the reason once the link ends', async () => {
|
||||
const link = new LinkLife({ log: silentLog, reconnects: true, timeout: 0 });
|
||||
|
||||
link.drop();
|
||||
|
||||
const held = link.hold(undefined)();
|
||||
|
||||
link.end();
|
||||
|
||||
assert.match((await held).err?.message ?? '', /Session is closed/);
|
||||
link.attach();
|
||||
assert.equal(link.isAttached(), false, 'ended is final');
|
||||
assert.equal(link.end(), false);
|
||||
});
|
||||
});
|
||||
|
||||
@@ -2424,6 +2461,39 @@ describe('a peer that sends the next segment only once the last one is answered'
|
||||
assert.deepEqual(await peerOf(smpp).close(), {}, 'a group nothing completed is not held');
|
||||
});
|
||||
|
||||
test('reports a drop once as disconnected when a listener closes the session over the segments it lost', async t => {
|
||||
const smpp = await startServer(t);
|
||||
const { session } = await connect(t, smpp, { reconnect: { maxDelay: 100, minDelay: 20 } });
|
||||
|
||||
assert.ok(session);
|
||||
|
||||
const events: string[] = [];
|
||||
const closed = once<true>(resolve => { session.on('close', () => { resolve(true); }); });
|
||||
|
||||
session.on('close', () => { events.push('close'); });
|
||||
session.on('disconnected', () => { events.push('disconnected'); });
|
||||
session.on('sessionError', () => { void session.close(); });
|
||||
|
||||
const [first] = segmentsOf(0x2D);
|
||||
|
||||
assert.ok(first);
|
||||
|
||||
const delivered = await peerOf(smpp).send({
|
||||
cmdName: 'deliver_sm',
|
||||
params: submitSmParams(
|
||||
{ from: '46701113311', message: text, to: '46709771337' },
|
||||
first,
|
||||
{ encoding: 'ASCII', multipart: true },
|
||||
),
|
||||
});
|
||||
|
||||
assert.equal(delivered.err, undefined);
|
||||
await peerOf(smpp).close();
|
||||
await closed;
|
||||
|
||||
assert.deepEqual(events, ['disconnected', 'close']);
|
||||
});
|
||||
|
||||
test('close() still waits for a concatenated message the application has not answered', async t => {
|
||||
const smpp = await startServer(t, { shutdownTimeout: 50 });
|
||||
const incoming = once<Sms>(resolve => {
|
||||
|
||||
@@ -200,19 +200,38 @@ next work ([decision](docs/decisions.md#internals-and-tests)).
|
||||
### Locality — next, ahead of everything below; 5–6 today, and the gate is 7
|
||||
|
||||
A second four-seat run on 2026-09-28, after #35–#40, read 6, 6, 7 and 6 again, Locality 5, 5, 6
|
||||
and 5. Every seat ranked the session's lifecycle hardest and least wanted to modify it.
|
||||
and 5. Every seat ranked the session's lifecycle hardest and least wanted to modify it. A third,
|
||||
after #46, read 6, 6, 7 and 7, Locality 5, 5, 6 and 6. A fourth, after the link's liveness got one
|
||||
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
|
||||
`Session`'s rejection handler second.
|
||||
|
||||
- [ ] **Lift Locality to 7, and confirm it with a scoring run.** A run reading 7.0 or above also
|
||||
retires the #30 and #46 decision. The sub-items are what the 2026-09-28 run named, most seats first.
|
||||
- [ ] **Give the link's liveness one owner.** A third run the same day, after #46, read 6, 6, 7 and
|
||||
7, Locality 5, 5, 6 and 6; all four seats ranked `drain()`/`end()`/`teardown()`/`retrying()`
|
||||
in `session.ts` hardest, because whether the link lives is kept in `Session.lifecycle`,
|
||||
`ReconnectLoop.halted`, `LinkGate.up`/`returning`, `OutgoingRequests.draining` and
|
||||
`IncomingRequests.linkGeneration`, held in step by statement order and the comment above
|
||||
`retrying()`. Four seats.
|
||||
retires the #30, #46 and #48 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
|
||||
|
||||
- [ ] **Refuse to open a link that dropped while its rebind was answered.** A peer sending
|
||||
`bind_resp` and FIN together can tear the link down before `comeBackUp()` resumes; it then
|
||||
calls `link.open()` on a `down` link, `resetTimers()` skips, and `attempt()` reports success,
|
||||
so the session is `up` on a destroyed socket with nothing to reconnect it until `close()`.
|
||||
`open()` accepting only `binding`, and `comeBackUp()` returning an err when the link is no
|
||||
longer attached, lets the loop retry. Unreproduced; from the stability review of #48.
|
||||
|
||||
- [ ] **Register a multipart send's receipt merge before its segments go out.** `Session.sendSms()`
|
||||
calls `dlrMerger.expect()` only once `submitSms()` resolves, after the last segment's response,
|
||||
so a receipt for an early segment that arrives first is logged at `debug` as naming no merge,
|
||||
and the group then waits out `dlrMergeTimeout` with no `messageDlr`. Likeliest with a fast SMSC
|
||||
or more segments than `maxOutstanding`. Goal 2. From the 2026-09-28 scoring run on #48.
|
||||
|
||||
- [ ] **Settle what a repeated tag not marked `multiple` reads as, and pin it in a test.** A vendor
|
||||
tag or a known single-value tag a peer sends twice keeps the last occurrence and drops the
|
||||
rest silently, which goal 3 argues against; listing it would change every such tag's shape.
|
||||
@@ -367,11 +386,20 @@ and 5. Every seat ranked the session's lifecycle hardest and least wanted to mod
|
||||
`error`-event reason (hard rule 3 owns it) and the Audience bullets restating goal 8 and
|
||||
Install; the `'use strict'` clause in both MIGRATION.md and CHANGELOG.md; the node-smpp
|
||||
cross-check in MIGRATION.md; the planned work in `interop-tests/AGENTS.md` (an expected
|
||||
malformed count per peer) and `benchmarks/README.md`.
|
||||
malformed count per peer) and `benchmarks/README.md`. The prose sweep of #48 adds: the smppload
|
||||
note in both `benchmarks/README.md` and `interop-tests/README.md`; the summary after the
|
||||
`AGENTS.md` link in `interop-tests/README.md`; the `run.py` foreground rule tacked onto rule 5 in
|
||||
`interop-tests/AGENTS.md`, which wants its own number.
|
||||
|
||||
- [ ] **Make `LinkGate.isUp()`'s doc true or its state match it.** It says a link attached but not
|
||||
yet bound cannot carry a request, while `up` starts `true`, so the first link and a server
|
||||
session are up before any bind. The gate decision in `docs/decisions.md` makes the same claim
|
||||
- [ ] **Give this library one figure at window 50 in `benchmarks/README.md`.** Its "same sink, same
|
||||
host" table reads 37,125/s where the table below it reads 38,675/s; re-measure or cite one run.
|
||||
|
||||
- [ ] **Name the goal and the premise of every `docs/decisions.md` entry.** The prose sweep of #48
|
||||
counted 30 of 58 entries naming no goal and 52 with no "valid while" premise.
|
||||
|
||||
- [ ] **Make `LinkLife` start unbound, or its decision's title true.** A link attached but not
|
||||
yet bound cannot carry a request, while `phase` starts `up`, so the first link and a server
|
||||
session are up before any bind. The `LinkLife` decision in `docs/decisions.md` makes the same claim
|
||||
in its title, and carries the same fix. From the comprehension panel of #25.
|
||||
|
||||
- [ ] **Move `checkSessionOptions()`'s doc comment to what it describes.** It explains why a count
|
||||
@@ -485,9 +513,9 @@ and 5. Every seat ranked the session's lifecycle hardest and least wanted to mod
|
||||
support Gitea, so mirror each Gitea pull request to GitHub for it to review there.
|
||||
Maintainer's ask, 2026-09-14; not started until asked.
|
||||
|
||||
- [ ] **Count what is left of a budget one way in `leftOf()` and the link gate.** Today they are one
|
||||
- [ ] **Count what is left of a budget one way in `leftOf()` and `LinkLife`.** Today they are one
|
||||
concept counted twice. `idle-waiters.ts` reads what is left of a budget as `Math.max(1,
|
||||
deadline - now)`, because 0 means "forever" there; `link-gate.ts` runs the same subtraction
|
||||
deadline - now)`, because 0 means "forever" there; `link-life.ts` runs the same subtraction
|
||||
and calls `<= 0` expired. Neither is reachable from the other, so nothing can disagree today,
|
||||
but a reader who learns one and applies it to the other is wrong. A budget type both take
|
||||
would close it. Raised by review, 2026-09-01.
|
||||
|
||||
Reference in New Issue
Block a user