Give the gate, the window and the pending map one owner

This commit is contained in:
2026-09-01 16:52:15 +02:00
parent 3bcea108e3
commit 8becbc0ab6
12 changed files with 276 additions and 168 deletions
+5
View File
@@ -10,6 +10,11 @@ export class HeldMessages {
this.held.add(pduObjs);
}
/** Whether a drain is still waiting for this message to be answered. */
has(pduObjs: PduObject[]): boolean {
return this.held.has(pduObjs);
}
release(pduObjs: PduObject[]): void {
if (!this.held.delete(pduObjs)) return;
+6 -5
View File
@@ -21,8 +21,8 @@ export type IncomingRequestsOptions = {
maxReassembly?: number | undefined;
onRequest?: OnRequest | undefined;
reassemblyTimeout?: number | undefined;
/** Past the drain gate: a receipt answering a message the shutdown is still waiting for. */
sendHeld: (input: PduObjectInput) => Promise<Result<{ pduObj: PduObject }>>;
/** Past a drain's refusal, for a receipt the drain is itself waiting for. */
sendPastDrain: (input: PduObjectInput) => Promise<Result<{ pduObj: PduObject }>>;
session: Session;
smsIdFormat?: SmsIdFormat | undefined;
systemId?: string | undefined;
@@ -35,7 +35,7 @@ export class IncomingRequests {
private readonly log: SmppLog;
private readonly onRequest: OnRequest | undefined;
private readonly reassembler: Reassembler;
private readonly sendHeld: IncomingRequestsOptions['sendHeld'];
private readonly sendPastDrain: IncomingRequestsOptions['sendPastDrain'];
private readonly session: Session;
private readonly smsIdFormat: SmsIdFormat;
private readonly systemId: string;
@@ -50,7 +50,7 @@ export class IncomingRequests {
maxOctets: options.maxOctets,
timeout: options.reassemblyTimeout ?? defaults.reassemblyTimeout,
});
this.sendHeld = options.sendHeld;
this.sendPastDrain = options.sendPastDrain;
this.session = options.session;
this.smsIdFormat = options.smsIdFormat ?? {};
this.systemId = options.systemId ?? defaults.systemId;
@@ -167,7 +167,8 @@ export class IncomingRequests {
}, {
// A turn later, so a listener sending its receipt straight after the response still holds.
onAnswered: () => { setImmediate(() => { this.held.release(pduObjs); }); },
send: this.sendHeld,
// Past the refusal only while a drain is still waiting for this message; an ordinary send after.
send: input => (this.held.has(pduObjs) ? this.sendPastDrain(input) : this.session.send(input)),
});
this.held.hold(pduObjs);
+8 -6
View File
@@ -47,9 +47,11 @@ export class LinkGate {
return this.up || this.returning ? undefined : over();
}
/** When a hold starting now has to give up. 0 never does. */
deadline(): number {
return this.timeout > 0 ? this.now() + this.timeout : 0;
/** One budget for a request, however many links it waits through. 0 never gives up. */
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. */
@@ -73,7 +75,7 @@ export class LinkGate {
}
/** Resolves once a link can carry the request, or with the reason none ever will. */
wait(deadline: number, signal: AbortSignal | undefined): Promise<VoidResult> {
private wait(deadline: number, signal: AbortSignal | undefined): Promise<VoidResult> {
if (this.up) return Promise.resolve({});
const refused = this.refusal();
@@ -86,10 +88,10 @@ export class LinkGate {
if (deadline !== 0 && left <= 0) return Promise.resolve({ err: expired() });
return this.hold(left, signal);
return this.waitForLink(left, signal);
}
private hold(left: number, signal: AbortSignal | undefined): Promise<VoidResult> {
private waitForLink(left: number, signal: AbortSignal | undefined): Promise<VoidResult> {
this.log.verbose('linkGate - holding a request until a link is back', { timeout: left });
return new Promise<VoidResult>(resolve => {
+168
View File
@@ -0,0 +1,168 @@
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';
import { bindCommands } from './session-options.ts';
import { objToPdu } from './pdu.ts';
export type OutgoingRequestsOptions = {
log: SmppLog;
maxOutstanding: number;
responseTimeout: number;
transport: PduTransport;
};
/** `retryOnNextLink`: the write failed, so nothing reached the socket and another link may carry it. */
type Attempt = { result: Result<{ pduObj: PduObject }>; retryOnNextLink: boolean };
function abortedBeforeSend(): Error {
return new Error('Aborted before the request was sent');
}
/** 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 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.log = options.log;
this.pending = new PendingRequests(options.log);
this.responseTimeout = options.responseTimeout;
this.transport = options.transport;
this.window = new SendWindow(options.maxOutstanding);
}
/** Read through a method: a drop can land while a request is awaiting. */
linkDown(): boolean {
return !this.gate.isUp() || this.transport.sock.destroyed;
}
/** 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);
this.pending.settleAll(new Error('Session closed before a response arrived'));
}
/** Hands a response to the request waiting for it. False means nothing was. */
deliver(pduObj: PduObject): boolean {
return this.pending.deliver(pduObj);
}
/** Sends a request and resolves with the peer's response. */
request(input: PduObjectInput, options: SendOptions): Promise<Result<{ pduObj: PduObject }>> {
// A drain on a live link. A link that is down is the gate's answer, which says closed instead.
if (this.draining && !this.linkDown()) {
return Promise.resolve({ err: new Error('Session is shutting down') });
}
return this.pastDrain(input, options);
}
/** The same path without that refusal, which a receipt for a held message has to take. */
async pastDrain(
input: PduObjectInput,
options: SendOptions,
): Promise<Result<{ pduObj: PduObject }>> {
const refused = this.refuse(input, options);
if (refused) return { err: refused };
// A bind is what makes a link usable, so it cannot wait for one.
if (bindCommands.includes(input.cmdName)) return this.now(input, options);
const waitForLink = this.gate.hold(options.signal);
for (;;) {
const held = await waitForLink();
if (held.err) return { err: held.err };
await this.window.acquire();
const attempt = await this.attempt(input, options).finally(() => { this.window.release(); });
// Nothing reached the socket, so the next link carries it instead of the caller resending.
if (!attempt.retryOnNextLink || this.gate.isUp() || this.gate.refusal()) return attempt.result;
}
}
/** Past the gate, the window and a drain, for what has to go out either way. */
async now(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);
if (unfinished === 0) return {};
this.log.warn('session - shutting down with requests unfinished', { timeout, unfinished });
return { err: new Error(`Shut down with ${String(unfinished)} request(s) unfinished`) };
}
/** Why a request cannot go out at all, as opposed to not yet. */
private refuse(input: PduObjectInput, options: SendOptions): Error | undefined {
if (input.cmdName.endsWith('_resp')) {
return new Error(`Use sendReturn() for responses, not send(): ${input.cmdName}`);
}
// Before the gate and the window, or an aborted call waits for what it will never use.
if (options.signal?.aborted === true) return abortedBeforeSend();
// A bind skips the gate below, so the answer it would have given is given here instead.
return bindCommands.includes(input.cmdName) ? this.gate.refusal() : undefined;
}
private async attempt(input: PduObjectInput, options: SendOptions): Promise<Attempt> {
// pending.wait() alone settles the caller while the request still goes out to the peer.
if (options.signal?.aborted === true) {
return { result: { err: abortedBeforeSend() }, retryOnNextLink: false };
}
const seqNr = this.pending.nextSeqNr();
const built = objToPdu({ ...input, seqNr });
if (built.err) return { result: { err: built.err }, retryOnNextLink: false };
const response = this.pending.wait(seqNr, {
signal: options.signal,
timeout: this.responseTimeout,
});
const written = this.transport.write(built.buffer);
if (written.err) {
this.pending.settle(seqNr, { err: written.err });
return { result: { err: written.err }, retryOnNextLink: true };
}
const answered = await response;
// It went out, so a failure now means the peer may have taken it and the answer was the loss.
return { result: answered.err ? { err: new UnansweredError(answered.err) } : answered, retryOnNextLink: false };
}
}
-8
View File
@@ -12,14 +12,6 @@ type Pending = {
settle: (result: Result<{ pduObj: PduObject }>) => void;
};
/** The request went out and no answer came back: the peer may have accepted it. */
export class UnansweredError extends Error {
constructor(cause: Error) {
super(`No answer came back, so the peer may have accepted it: ${cause.message}`, { cause });
this.name = 'UnansweredError';
}
}
/** Hands out sequence numbers and matches responses to the requests waiting for them. */
export class PendingRequests {
private readonly log: SmppLog;
+1 -1
View File
@@ -4,7 +4,7 @@ import type { PduObject, PduObjectInput } from './pdu.ts';
import type { Result } from './result.ts';
import type { SmppLog } from './log.ts';
import type { SmsIdNotation } from './sms-id.ts';
import { UnansweredError } from './pending-requests.ts';
import { UnansweredError } from './unanswered-error.ts';
import { consts } from './defs/constants.ts';
import { detect } from './defs/encodings.ts';
import { normaliseSmsId } from './sms-id.ts';
+27 -117
View File
@@ -10,17 +10,15 @@ import type { Socket } from 'node:net';
import { DlrMerger } from './dlr-merger.ts';
import { EventEmitter } from 'node:events';
import { IncomingRequests } from './incoming-requests.ts';
import { LinkGate } from './link-gate.ts';
import { LinkTimers } from './link-timers.ts';
import { OutgoingRequests } from './outgoing-requests.ts';
import { PduTransport } from './pdu-transport.ts';
import { PendingRequests, UnansweredError } from './pending-requests.ts';
import { ReconnectLoop } from './reconnect-loop.ts';
import { SendWindow } from './send-window.ts';
import { leftOf } from './idle-waiters.ts';
import { errorFrom } from './error-from.ts';
import { optionalParamsMinVersion } from './defs/constants.ts';
import { bindCarries, bindCommands, defaultSystemId, defaults } from './session-options.ts';
import { isResp, objToPdu, pduReturn } from './pdu.ts';
import { isResp, pduReturn } from './pdu.ts';
import { silentLog } from './log.ts';
import { submitSms, unsent } from './send-sms.ts';
import { ConcatReference } from './udh.ts';
@@ -41,13 +39,6 @@ export { bindCommands, defaultSystemId };
/** A listener may return a promise: an `async` one that rejects is routed like one that throws. */
type SessionListener<K extends keyof SessionEvents> = (...args: SessionEvents[K]) => unknown;
function abortedBeforeSend(): Error {
return new Error('Aborted before the request was sent');
}
/** `retryOnNextLink`: the write failed, so nothing reached the socket and another link may carry it. */
type Attempt = { result: Result<{ pduObj: PduObject }>; retryOnNextLink: boolean };
export class Session extends EventEmitter<SessionEvents> {
declare addListener: <K extends keyof SessionEvents>(event: K, listener: SessionListener<K>) => this;
declare off: <K extends keyof SessionEvents>(event: K, listener: SessionListener<K>) => this;
@@ -68,17 +59,14 @@ export class Session extends EventEmitter<SessionEvents> {
private readonly concatReference = new ConcatReference();
private readonly dlrMerger: DlrMerger;
private readonly gate: LinkGate;
private readonly incoming: IncomingRequests;
private readonly options: SessionOptions;
private readonly pending: PendingRequests;
private readonly outgoing: OutgoingRequests;
private readonly reconnectLoop: ReconnectLoop | undefined;
private readonly timers: LinkTimers;
private readonly transport: PduTransport;
private readonly window: SendWindow;
private closed = false;
private draining = false;
private ended = false;
/** A listener that throws is the application's bug; it must not become ours. Hard rule 1. */
@@ -123,7 +111,6 @@ export class Session extends EventEmitter<SessionEvents> {
max: defaults.maxDlrMerges,
timeout: defaults.dlrMergeTimeout,
});
this.gate = new LinkGate({ log: this.log, timeout: options.responseTimeout ?? defaults.responseTimeout });
this.incoming = new IncomingRequests({
dlrMerger: this.dlrMerger,
log: this.log,
@@ -131,12 +118,11 @@ export class Session extends EventEmitter<SessionEvents> {
maxReassembly: options.maxReassembly,
onRequest: options.onRequest,
reassemblyTimeout: options.reassemblyTimeout,
sendHeld: input => this.sendThrough(input, {}),
sendPastDrain: input => this.outgoing.pastDrain(input, {}),
session: this,
smsIdFormat: options.smsIdFormat,
systemId: options.systemId,
});
this.pending = new PendingRequests(this.log);
this.reconnectLoop = this.loopFor(options.reconnect);
this.timers = new LinkTimers({
enquireLinkInterval: options.enquireLinkInterval,
@@ -147,7 +133,12 @@ export class Session extends EventEmitter<SessionEvents> {
onIdle: () => { this.teardown(); },
});
this.transport = this.transportFor(options.sock);
this.window = new SendWindow(options.maxOutstanding ?? defaults.maxOutstanding);
this.outgoing = new OutgoingRequests({
log: this.log,
maxOutstanding: options.maxOutstanding ?? defaults.maxOutstanding,
responseTimeout: options.responseTimeout ?? defaults.responseTimeout,
transport: this.transport,
});
this.resetTimers();
}
@@ -170,56 +161,7 @@ export class Session extends EventEmitter<SessionEvents> {
/** Sends a request and resolves with the peer's response. */
send(input: PduObjectInput, options: SendOptions = {}): Promise<Result<{ pduObj: PduObject }>> {
// A drain on a live link. A link that is down is the gate's answer, which says closed instead.
if (this.draining && !this.linkDown()) return Promise.resolve({ err: new Error('Session is shutting down') });
return this.sendThrough(input, options);
}
/** The same path without that refusal, which a receipt for a held message has to take. */
private async sendThrough(
input: PduObjectInput,
options: SendOptions = {},
): Promise<Result<{ pduObj: PduObject }>> {
const refused = this.refuseSend(input, options);
if (refused) return { err: refused };
// A bind is what makes a link usable, so it cannot wait for one.
if (bindCommands.includes(input.cmdName)) return (await this.attempt(input, options)).result;
const deadline = this.gate.deadline();
for (;;) {
const held = await this.gate.wait(deadline, options.signal);
if (held.err) return { err: held.err };
await this.window.acquire();
const attempt = await this.attempt(input, options).finally(() => { this.window.release(); });
// Nothing reached the socket, so the next link carries it instead of the caller resending.
if (!attempt.retryOnNextLink || this.gate.isUp() || !this.retrying()) return attempt.result;
}
}
/** Why a request cannot go out at all, as opposed to not yet. */
private refuseSend(input: PduObjectInput, options: SendOptions): Error | undefined {
if (input.cmdName.endsWith('_resp')) {
return new Error(`Use sendReturn() for responses, not send(): ${input.cmdName}`);
}
// Before the gate and the window, or an aborted call waits for what it will never use.
if (options.signal?.aborted === true) return abortedBeforeSend();
// A bind skips the gate below, so the answer it would have given is given here instead.
return bindCommands.includes(input.cmdName) ? this.gate.refusal() : undefined;
}
/** Read through a method: a drop can land while a send is awaiting. */
private linkDown(): boolean {
return !this.gate.isUp() || this.sock.destroyed;
return this.outgoing.request(input, options);
}
/** Answers a request the peer sent us. Responses are never waited on. */
@@ -269,9 +211,9 @@ export class Session extends EventEmitter<SessionEvents> {
async unbind(): Promise<VoidResult> {
const drained = await this.drain(undefined);
const wasOpen = !this.closed;
// attempt(), not send(): the drain gate refuses a send, and the unbind goes out either way.
// now(), not send(): a drain refuses a send, and the unbind goes out either way.
const sent = wasOpen
? (await this.attempt({ cmdName: 'unbind' }, {})).result
? await this.outgoing.now({ cmdName: 'unbind' })
: { err: new Error('Session is closed') };
const closedOnUnbind = wasOpen && this.closed;
@@ -341,7 +283,7 @@ export class Session extends EventEmitter<SessionEvents> {
}
this.resetTimers();
this.gate.open();
this.outgoing.linkUp();
this.log.info('session - reconnected');
this.emit('reconnected');
@@ -353,58 +295,27 @@ export class Session extends EventEmitter<SessionEvents> {
this.closed = false;
}
private async attempt(input: PduObjectInput, options: SendOptions): Promise<Attempt> {
// pending.wait() alone settles the caller while the request still goes out to the peer.
if (options.signal?.aborted === true) {
return { result: { err: abortedBeforeSend() }, retryOnNextLink: false };
}
const seqNr = this.pending.nextSeqNr();
const built = objToPdu({ ...input, seqNr });
if (built.err) return { result: { err: built.err }, retryOnNextLink: false };
const response = this.pending.wait(seqNr, {
signal: options.signal,
timeout: this.options.responseTimeout ?? defaults.responseTimeout,
});
const written = this.transport.write(built.buffer);
if (written.err) {
this.pending.settle(seqNr, { err: written.err });
return { result: { err: written.err }, retryOnNextLink: true };
}
const answered = await response;
// It went out, so a failure now means the peer may have taken it and the answer was the loss.
return { result: answered.err ? { err: new UnansweredError(answered.err) } : answered, retryOnNextLink: false };
}
/** 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.draining = true;
this.outgoing.stopAccepting();
if (this.linkDown()) return {};
if (this.outgoing.linkDown()) return {};
const timeout = this.options.shutdownTimeout ?? defaults.shutdownTimeout;
const deadline = timeout > 0 ? Date.now() + timeout : 0;
// Only the application answers a held message, so that half falls back rather than wait forever.
const answering = timeout > 0 ? timeout : (this.options.responseTimeout ?? defaults.responseTimeout);
// Answering a message can put a receipt on the wire; nothing on the wire produces a message.
const messages = await this.incoming.drain(timeout, signal);
const unfinished = await this.window.idle(leftOf(deadline), signal);
const messages = await this.incoming.drain(answering, signal);
const requests = await this.outgoing.drain(leftOf(deadline), signal);
// The window empties on a teardown too, which settles everything the link was carrying.
if (this.linkDown()) return { err: new Error('The session closed before the drain finished') };
if (this.outgoing.linkDown()) {
return { err: new Error('The session closed before the drain finished') };
}
if (messages.err) return messages;
if (unfinished === 0) return {};
this.log.warn('session - shutting down with requests unfinished', { timeout, unfinished });
return { err: new Error(`Shut down with ${String(unfinished)} request(s) unfinished`) };
return messages.err ? messages : requests;
}
/** The session is over now, drained or not. Nothing brings it back. */
@@ -419,7 +330,7 @@ export class Session extends EventEmitter<SessionEvents> {
if (this.ended) return;
this.ended = true;
this.gate.shut(false);
this.outgoing.linkLost(false);
this.emit('close');
}
@@ -427,9 +338,8 @@ export class Session extends EventEmitter<SessionEvents> {
if (this.closed) return;
this.closed = true;
this.gate.shut(this.retrying());
this.outgoing.linkLost(this.retrying());
this.timers.clear();
this.pending.settleAll(new Error('Session closed before a response arrived'));
this.incoming.clear();
this.sock.destroy();
@@ -448,7 +358,7 @@ export class Session extends EventEmitter<SessionEvents> {
private dispatch(pduObj: PduObject): void {
if (isResp(pduObj)) {
if (!this.pending.deliver(pduObj)) {
if (!this.outgoing.deliver(pduObj)) {
this.log.debug('session - response with no matching request', { seqNr: pduObj.seqNr });
}
+7
View File
@@ -0,0 +1,7 @@
/** The request went out and no answer came back: the peer may have accepted it. */
export class UnansweredError extends Error {
constructor(cause: Error) {
super(`No answer came back, so the peer may have accepted it: ${cause.message}`, { cause });
this.name = 'UnansweredError';
}
}