Give the link's liveness one owner
Test / lint (pull_request) Successful in 24s
Test / test (18) (pull_request) Successful in 32s
Test / test (20) (pull_request) Successful in 32s
Test / test (22) (pull_request) Successful in 32s
Test / test (24) (pull_request) Successful in 32s
Mirror / push (push) Has been cancelled
Test / test (26) (pull_request) Successful in 32s
Test / lint (pull_request) Successful in 24s
Test / test (18) (pull_request) Successful in 32s
Test / test (20) (pull_request) Successful in 32s
Test / test (22) (pull_request) Successful in 32s
Test / test (24) (pull_request) Successful in 32s
Mirror / push (push) Has been cancelled
Test / test (26) (pull_request) Successful in 32s
This commit is contained in:
+33
-35
@@ -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' })
|
||||
: { 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<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,13 +317,12 @@ 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.
|
||||
if (!this.outgoing.canCarry()) return {};
|
||||
@@ -355,29 +355,32 @@ 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;
|
||||
// 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<SessionEvents> {
|
||||
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<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