|
|
|
|
@ -38,7 +38,7 @@ import { wireConnectionSession } from './session-wiring.ts';
|
|
|
|
|
import { createConnectionSender } from './sender.ts';
|
|
|
|
|
import { createFrame } from './serializer.ts';
|
|
|
|
|
import { jsonConnectionSerializer } from './serializers/json.ts';
|
|
|
|
|
import { attachConnectionTransport } from './transport-wiring.ts';
|
|
|
|
|
import { createConnectionTransportRuntime } from './connection-transport-runtime.ts';
|
|
|
|
|
import type { Logger } from '$libs/logr';
|
|
|
|
|
import type {
|
|
|
|
|
Connection,
|
|
|
|
|
@ -49,8 +49,7 @@ import type {
|
|
|
|
|
ConnectionConnectResult,
|
|
|
|
|
ConnectionOptions,
|
|
|
|
|
ConnectionSessionSource,
|
|
|
|
|
ConnectionState,
|
|
|
|
|
ConnectionTransport
|
|
|
|
|
ConnectionState
|
|
|
|
|
} from './types.ts';
|
|
|
|
|
|
|
|
|
|
interface ConnectionRuntime {
|
|
|
|
|
@ -79,11 +78,17 @@ export function createConnection<TChannels extends ConnectionChannelMap = Connec
|
|
|
|
|
emitState: (change) => events.emitState(change)
|
|
|
|
|
});
|
|
|
|
|
const channelRegistry = createConnectionChannelRegistry<TChannels>(reportChannelListenerError);
|
|
|
|
|
const transportRuntime = createConnectionTransportRuntime(options.transport, {
|
|
|
|
|
onOpen: handleTransportOpen,
|
|
|
|
|
onMessage: handleTransportMessage,
|
|
|
|
|
onClose: handleTransportClose,
|
|
|
|
|
onError: handleTransportError
|
|
|
|
|
});
|
|
|
|
|
const sender = createConnectionSender({
|
|
|
|
|
serializer,
|
|
|
|
|
frameBuffer,
|
|
|
|
|
isConnected,
|
|
|
|
|
getTransport: () => transport,
|
|
|
|
|
getTransport: () => transportRuntime.current,
|
|
|
|
|
setError: (err) => {
|
|
|
|
|
stateTracker.setError(err);
|
|
|
|
|
},
|
|
|
|
|
@ -97,31 +102,10 @@ export function createConnection<TChannels extends ConnectionChannelMap = Connec
|
|
|
|
|
sender
|
|
|
|
|
});
|
|
|
|
|
|
|
|
|
|
let transport: ConnectionTransport | null = null;
|
|
|
|
|
let detachTransportListeners = (): void => {};
|
|
|
|
|
let disposed = false;
|
|
|
|
|
let intentionalClose = false;
|
|
|
|
|
let connectPromise: Promise<ConnectionConnectResult> | null = null;
|
|
|
|
|
|
|
|
|
|
function resolveTransport(): ConnectionTransport {
|
|
|
|
|
return typeof options.transport === 'function' ? options.transport() : options.transport;
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
function detachTransport(): void {
|
|
|
|
|
detachTransportListeners();
|
|
|
|
|
detachTransportListeners = (): void => {};
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
function attachTransport(next: ConnectionTransport): void {
|
|
|
|
|
detachTransport();
|
|
|
|
|
detachTransportListeners = attachConnectionTransport(next, {
|
|
|
|
|
onOpen: handleTransportOpen,
|
|
|
|
|
onMessage: handleTransportMessage,
|
|
|
|
|
onClose: handleTransportClose,
|
|
|
|
|
onError: handleTransportError
|
|
|
|
|
});
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
function markOpen(): void {
|
|
|
|
|
if (stateTracker.markOpen()) heartbeat.start();
|
|
|
|
|
}
|
|
|
|
|
@ -181,7 +165,7 @@ export function createConnection<TChannels extends ConnectionChannelMap = Connec
|
|
|
|
|
scheduleTimer: timers.schedule,
|
|
|
|
|
cancelTimer: timers.cancel,
|
|
|
|
|
closeTransport: (reason) => {
|
|
|
|
|
transport?.close(undefined, reason);
|
|
|
|
|
transportRuntime.close(reason);
|
|
|
|
|
},
|
|
|
|
|
diagnostics
|
|
|
|
|
});
|
|
|
|
|
@ -206,7 +190,7 @@ export function createConnection<TChannels extends ConnectionChannelMap = Connec
|
|
|
|
|
function isConnected(): boolean {
|
|
|
|
|
return (
|
|
|
|
|
stateTracker.state === CONNECTION_STATE_OPEN &&
|
|
|
|
|
transport?.state === CONNECTION_TRANSPORT_STATE_OPEN
|
|
|
|
|
transportRuntime.current?.state === CONNECTION_TRANSPORT_STATE_OPEN
|
|
|
|
|
);
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
@ -233,12 +217,9 @@ export function createConnection<TChannels extends ConnectionChannelMap = Connec
|
|
|
|
|
return await runConnectionConnect({
|
|
|
|
|
reconnect,
|
|
|
|
|
timers,
|
|
|
|
|
resolveTransport,
|
|
|
|
|
setTransport: (nextTransport) => {
|
|
|
|
|
transport = nextTransport;
|
|
|
|
|
},
|
|
|
|
|
attachTransport,
|
|
|
|
|
detachTransport,
|
|
|
|
|
resolveTransport: transportRuntime.resolve,
|
|
|
|
|
attachTransport: transportRuntime.attach,
|
|
|
|
|
detachTransport: transportRuntime.detach,
|
|
|
|
|
markOpen,
|
|
|
|
|
markClosedFailed: (error) => {
|
|
|
|
|
markClosed(CONNECTION_STATE_FAILED, error);
|
|
|
|
|
@ -263,7 +244,7 @@ export function createConnection<TChannels extends ConnectionChannelMap = Connec
|
|
|
|
|
heartbeat.stop();
|
|
|
|
|
acks.resolveAll({ ok: false, reason: CONNECTION_ACK_REASON_CLOSED });
|
|
|
|
|
stateTracker.transition(CONNECTION_STATE_CLOSING);
|
|
|
|
|
transport?.close(undefined, reason);
|
|
|
|
|
transportRuntime.close(reason);
|
|
|
|
|
if (stateTracker.state !== CONNECTION_STATE_CLOSED) markClosed(CONNECTION_STATE_CLOSED);
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
@ -372,7 +353,7 @@ export function createConnection<TChannels extends ConnectionChannelMap = Connec
|
|
|
|
|
closeTransport(CONNECTION_CLOSE_REASON_DISPOSE);
|
|
|
|
|
detachSession();
|
|
|
|
|
detachBrowserReconnect();
|
|
|
|
|
detachTransport();
|
|
|
|
|
transportRuntime.detach();
|
|
|
|
|
channelRegistry.dispose();
|
|
|
|
|
events.clear();
|
|
|
|
|
frameBuffer.clear();
|
|
|
|
|
|