Give the link's liveness one owner
Test / lint (pull_request) Successful in 24s
Test / test (18) (pull_request) Successful in 34s
Test / test (20) (pull_request) Successful in 32s
Test / test (22) (pull_request) Successful in 32s
Test / test (24) (pull_request) Successful in 38s
Test / test (26) (pull_request) Successful in 32s
Mirror / push (push) Successful in 7s
Test / lint (pull_request) Successful in 24s
Test / test (18) (pull_request) Successful in 34s
Test / test (20) (pull_request) Successful in 32s
Test / test (22) (pull_request) Successful in 32s
Test / test (24) (pull_request) Successful in 38s
Test / test (26) (pull_request) Successful in 32s
Mirror / push (push) Successful in 7s
This commit was merged in pull request #48.
This commit is contained in:
@@ -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();
|
||||
|
||||
|
||||
Reference in New Issue
Block a user