From cce1d7c02ed4ecc2a8c13ca0766c880190630d13 Mon Sep 17 00:00:00 2001 From: Lilleman auf Larv Date: Mon, 28 Sep 2026 21:24:34 +0200 Subject: [PATCH] Give the link's liveness one owner --- AGENTS.md | 8 +-- docs/decisions.md | 22 ++++--- src/incoming-requests.ts | 14 ++-- src/{link-gate.ts => link-life.ts} | 101 ++++++++++++++++++++++------- src/outgoing-requests.ts | 42 +++++------- src/reconnect-loop.ts | 2 +- src/session.ts | 68 ++++++++++--------- test/session-extras.test.ts | 79 +++++++++++++++------- todo.md | 21 +++--- 9 files changed, 214 insertions(+), 143 deletions(-) rename src/{link-gate.ts => link-life.ts} (52%) diff --git a/AGENTS.md b/AGENTS.md index 86e66fb..1d9b98e 100644 --- a/AGENTS.md +++ b/AGENTS.md @@ -55,12 +55,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 @@ -314,13 +314,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. +- `LinkLife` 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. - 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 diff --git a/docs/decisions.md b/docs/decisions.md index c7d939f..4988b42 100644 --- a/docs/decisions.md +++ b/docs/decisions.md @@ -508,9 +508,9 @@ 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 + `close` ends up holding two binds on one account. `teardown()` picks the event by + `LinkLife.retrying()`, and `end()` stops the session before tearing down, so every deliberate + shutdown emits `close`. A retry that opens a socket and then loses it is attached again through `attach()`, which is why a second drop emits again. - **An answer belongs to the link the message arrived on; a receipt does not.** Maintainer's call, @@ -542,7 +542,7 @@ rule and an index of the titles below. 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 + `LinkLife`'s hold is not — it is awaited with no other handle, so a process whose only work is `client()` would exit unbound. - **`connectTimeout` defaults to 10 s, bounds the whole connect including the TLS handshake, and @@ -760,15 +760,17 @@ rule and an index of the titles below. an abort at the gate 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.** +- **`LinkLife` 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 + that window go out unbound and come back `ESME_RINVBNDSTS`. `LinkLife` 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. + release. Every other collaborator reads whether the link lives from it and keeps no copy: five + copies held in step by statement order were what the 2026-09-28 comprehension runs ranked hardest. + The retry in `requestPastDrain()` asks `link.awaitsNextLink()` rather than `canCarry()`, which also + reads the socket: a loop condition the link does not gate on spins against a link that admits it + straight back. ## Internals and tests @@ -790,7 +792,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 diff --git a/src/incoming-requests.ts b/src/incoming-requests.ts index cb8bed0..f2a39f3 100644 --- a/src/incoming-requests.ts +++ b/src/incoming-requests.ts @@ -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 = { 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 { - 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)); diff --git a/src/link-gate.ts b/src/link-life.ts similarity index 52% rename from src/link-gate.ts rename to src/link-life.ts index e3fe884..0408b0e 100644 --- a/src/link-gate.ts +++ b/src/link-life.ts @@ -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 no request. */ +type Phase = 'binding' | 'down' | 'ended' | 'up'; + type Waiter = (result: VoidResult) => void; function aborted(): Error { @@ -22,34 +27,61 @@ 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(); - 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'; + } + + /** False once the session is shutting down: nothing new is taken, and no link follows this one. */ + isAccepting(): boolean { + return !this.stopped; + } + + /** Whether a link that drops now is followed by another. */ + retrying(): boolean { + return this.reconnects && !this.stopped; + } + + /** Down, with another link on its way. */ 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. */ @@ -59,29 +91,50 @@ export class LinkGate { 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. */ + attach(): void { + 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. False means there was none to lose. */ + drop(): boolean { + if (!this.isAttached()) return false; - if (!returning) this.release({ err: over() }); + this.phase = 'down'; + this.drops++; + + return true; + } + + 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 { - if (this.up) return Promise.resolve({}); + if (this.isUp()) return Promise.resolve({}); const refused = this.refusal(); @@ -97,7 +150,7 @@ export class LinkGate { } private waitForLink(left: number, signal: AbortSignal | undefined): Promise { - 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(resolve => { let timer: NodeJS.Timeout | undefined = undefined; @@ -109,7 +162,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() }); }; diff --git a/src/outgoing-requests.ts b/src/outgoing-requests.ts index f2c92f9..233eced 100644 --- a/src/outgoing-requests.ts +++ b/src/outgoing-requests.ts @@ -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,16 @@ 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(); + return !this.link.isAccepting() && !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')); } @@ -89,7 +82,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.isAccepting() && this.canCarry()) { return Promise.resolve({ err: new Error('Session is shutting down') }); } @@ -105,14 +98,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); } - const waitForLink = this.gate.hold(options.signal); + const waitForLink = this.link.hold(options.signal); for (;;) { const held = await waitForLink(); @@ -137,11 +130,6 @@ export class OutgoingRequests { 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 { const unfinished = await this.window.idle(timeout, signal); @@ -155,13 +143,13 @@ export class OutgoingRequests { /** Nothing reached the socket, so the next link carries it instead of the caller resending. */ 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); } diff --git a/src/reconnect-loop.ts b/src/reconnect-loop.ts index e59fb4b..af8044d 100644 --- a/src/reconnect-loop.ts +++ b/src/reconnect-loop.ts @@ -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; } diff --git a/src/session.ts b/src/session.ts index d7e6d85..76d0856 100644 --- a/src/session.ts +++ b/src/session.ts @@ -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 { 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( event: K, @@ -109,12 +108,12 @@ export class Session extends EventEmitter { 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 { }); 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 { 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 { */ async unbind(): Promise { 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' }) : { err: new Error('Session is closed') }; - const closedOnUnbind = wasOpen && this.lifecycle !== 'attached'; + const closedOnUnbind = wasOpen && !this.link.isAttached(); this.end(); @@ -300,14 +301,14 @@ export class Session extends EventEmitter { } // 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,13 +317,12 @@ export class Session extends EventEmitter { 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 { - this.reconnectLoop?.stop(); - this.outgoing.stopAccepting(); + this.stop(); // No link, so nothing is on the wire to wait out. if (!this.outgoing.canCarry()) return {}; @@ -355,29 +355,32 @@ export class Session extends EventEmitter { /** 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; + // Read once: clear() reports lost segments, and a listener could stop the session between reads. + const retrying = this.link.retrying(); - this.lifecycle = 'torn-down'; + if (!this.link.drop()) 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(); @@ -386,11 +389,6 @@ export class Session extends EventEmitter { 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 +434,13 @@ export class Session extends EventEmitter { } 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(); diff --git a/test/session-extras.test.ts b/test/session-extras.test.ts index 2cdb5e5..c130236 100644 --- a/test/session-extras.test.ts +++ b/test/session-extras.test.ts @@ -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 = {}): 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,73 @@ 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'); + assert.equal(new LinkLife({ log: silentLog, reconnects: false, timeout: 100 }).drop(), true); + }); + + test('drops an attached link once, and counts each drop', () => { + const link = new LinkLife({ log: silentLog, reconnects: true, timeout: 100 }); + const generation = link.generation(); + + assert.equal(link.drop(), true); + assert.equal(link.drop(), false, 'already down'); + assert.equal(link.generation(), generation + 1); + link.attach(); + assert.equal(link.drop(), true, 'a new link drops again'); + assert.equal(link.generation(), generation + 2); + }); + + 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/); }); }); diff --git a/todo.md b/todo.md index fee07a4..a4f01c6 100644 --- a/todo.md +++ b/todo.md @@ -200,16 +200,11 @@ 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; the link's liveness now has one owner. - [ ] **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 and #46 decision. ### Correctness @@ -369,9 +364,9 @@ and 5. Every seat ranked the session's lifecycle hardest and least wanted to mod cross-check in MIGRATION.md; the planned work in `interop-tests/AGENTS.md` (an expected malformed count per peer) and `benchmarks/README.md`. -- [ ] **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 +- [ ] **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 +480,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.