You can not select more than 25 topics
Topics must start with a letter or number, can include dashes ('-') and can be up to 35 characters long.
174 lines
5.0 KiB
174 lines
5.0 KiB
import { createEngineTimers } from '$timr';
|
|
import { createConnection } from './connection.ts';
|
|
import {
|
|
CONNECTION_METHOD_CLOSE_ALL,
|
|
CONNECTION_METHOD_CLOSE_CONNECTION,
|
|
CONNECTION_METHOD_CONNECTION,
|
|
CONNECTION_METHOD_CREATE_CONNECTION,
|
|
CONNECTION_METHOD_OPEN_ALL,
|
|
CONNECTION_METHOD_OPEN_CONNECTION,
|
|
CONNECTION_METHOD_RECONNECT_ALL,
|
|
CONNECTION_METHOD_RECONNECT_CONNECTION,
|
|
alreadyExistsErrorMessage,
|
|
disposedErrorMessage,
|
|
invalidNameErrorMessage,
|
|
notFoundErrorMessage
|
|
} from './consts.ts';
|
|
import {
|
|
ConnConnectionAlreadyExistsError,
|
|
ConnConnectionNotFoundError,
|
|
ConnDisposedError,
|
|
ConnInvalidConnectionNameError
|
|
} from './errors.ts';
|
|
import type {
|
|
Connection,
|
|
ConnectionChannelMap,
|
|
ConnectionMap,
|
|
ConnectionOptions,
|
|
EngineConnections,
|
|
EngineConnectionsOptions
|
|
} from './types.ts';
|
|
|
|
function validateConnectionName(name: unknown): asserts name is string {
|
|
if (typeof name !== 'string' || name.trim().length === 0) {
|
|
throw new ConnInvalidConnectionNameError(invalidNameErrorMessage(name));
|
|
}
|
|
}
|
|
|
|
function mergeObject<T extends object>(
|
|
defaults: T | undefined,
|
|
next: T | undefined
|
|
): T | undefined {
|
|
if (defaults === undefined) return next;
|
|
if (next === undefined) return defaults;
|
|
return { ...defaults, ...next };
|
|
}
|
|
|
|
function mergeConnectionOptions<TChannels extends ConnectionChannelMap>(
|
|
defaults: Partial<ConnectionOptions> | undefined,
|
|
options: ConnectionOptions<TChannels>
|
|
): ConnectionOptions<TChannels> {
|
|
const reconnect =
|
|
options.reconnect === false
|
|
? false
|
|
: defaults?.reconnect === false
|
|
? (options.reconnect ?? false)
|
|
: mergeObject(defaults?.reconnect, options.reconnect);
|
|
const heartbeat =
|
|
options.heartbeat === false
|
|
? false
|
|
: defaults?.heartbeat === false
|
|
? (options.heartbeat ?? false)
|
|
: mergeObject(defaults?.heartbeat, options.heartbeat);
|
|
return {
|
|
...defaults,
|
|
...options,
|
|
reconnect,
|
|
heartbeat,
|
|
buffer: mergeObject(defaults?.buffer, options.buffer),
|
|
channels: mergeObject(
|
|
defaults?.channels,
|
|
options.channels
|
|
) as ConnectionOptions<TChannels>['channels']
|
|
} as ConnectionOptions<TChannels>;
|
|
}
|
|
|
|
export function createEngineConnections<TConnections extends ConnectionMap = ConnectionMap>(
|
|
options: EngineConnectionsOptions = {}
|
|
): EngineConnections<TConnections> {
|
|
const connections = new Map<string, Connection>();
|
|
const timers =
|
|
options.timers ??
|
|
createEngineTimers({
|
|
logger: options.logger
|
|
});
|
|
const ownsTimers = options.timers === undefined;
|
|
let disposed = false;
|
|
|
|
function ensureLive(method: string): void {
|
|
if (disposed) throw new ConnDisposedError(disposedErrorMessage(method));
|
|
}
|
|
|
|
function getConnection(name: string, method: string): Connection {
|
|
ensureLive(method);
|
|
validateConnectionName(name);
|
|
const connection = connections.get(name);
|
|
if (connection === undefined) {
|
|
throw new ConnConnectionNotFoundError(notFoundErrorMessage(name), name);
|
|
}
|
|
return connection;
|
|
}
|
|
|
|
const engine: EngineConnections<TConnections> = {
|
|
createConnection<TChannels extends ConnectionChannelMap = ConnectionChannelMap>(
|
|
name: string,
|
|
connectionOptions: ConnectionOptions<TChannels>
|
|
): Connection<TChannels> {
|
|
ensureLive(CONNECTION_METHOD_CREATE_CONNECTION);
|
|
validateConnectionName(name);
|
|
if (connections.has(name)) {
|
|
throw new ConnConnectionAlreadyExistsError(alreadyExistsErrorMessage(name), name);
|
|
}
|
|
const connection = createConnection<TChannels>(
|
|
name,
|
|
mergeConnectionOptions<TChannels>(options.defaults, connectionOptions),
|
|
{
|
|
logger: options.logger,
|
|
timers,
|
|
session: options.session
|
|
}
|
|
);
|
|
connections.set(name, connection as Connection);
|
|
return connection;
|
|
},
|
|
connection<TChannels extends ConnectionChannelMap = ConnectionChannelMap>(
|
|
name: string
|
|
): Connection<TChannels> {
|
|
return getConnection(name, CONNECTION_METHOD_CONNECTION) as Connection<TChannels>;
|
|
},
|
|
has(name) {
|
|
if (disposed) return false;
|
|
return connections.has(name);
|
|
},
|
|
names() {
|
|
if (disposed) return [];
|
|
return [...connections.keys()];
|
|
},
|
|
openConnection(name) {
|
|
return getConnection(name, CONNECTION_METHOD_OPEN_CONNECTION).connect();
|
|
},
|
|
closeConnection(name, reason) {
|
|
getConnection(name, CONNECTION_METHOD_CLOSE_CONNECTION).disconnect(reason);
|
|
},
|
|
reconnectConnection(name, reason) {
|
|
return getConnection(name, CONNECTION_METHOD_RECONNECT_CONNECTION).reconnect(reason);
|
|
},
|
|
openAll() {
|
|
ensureLive(CONNECTION_METHOD_OPEN_ALL);
|
|
return Promise.all([...connections.values()].map((connection) => connection.connect()));
|
|
},
|
|
closeAll(reason) {
|
|
ensureLive(CONNECTION_METHOD_CLOSE_ALL);
|
|
for (const connection of connections.values()) connection.disconnect(reason);
|
|
},
|
|
reconnectAll(reason) {
|
|
ensureLive(CONNECTION_METHOD_RECONNECT_ALL);
|
|
return Promise.all(
|
|
[...connections.values()].map((connection) => connection.reconnect(reason))
|
|
);
|
|
},
|
|
close(name, reason) {
|
|
this.closeConnection(name, reason);
|
|
},
|
|
dispose() {
|
|
if (disposed) return;
|
|
disposed = true;
|
|
for (const connection of connections.values()) connection.dispose();
|
|
connections.clear();
|
|
if (ownsTimers) timers.dispose();
|
|
}
|
|
};
|
|
|
|
return engine;
|
|
}
|