feat: add strict local login websocket transport
This commit is contained in:
+117
@@ -0,0 +1,117 @@
|
|||||||
|
import {
|
||||||
|
assertLocalConnectionTarget,
|
||||||
|
type LocalStartupContext,
|
||||||
|
} from '../../framework/config/local-startup.ts';
|
||||||
|
import type { Transport } from '../../framework/net/transport.ts';
|
||||||
|
|
||||||
|
export type LocalLoginEvidence = { kind: string; data: unknown };
|
||||||
|
export type EvidenceSink = (event: LocalLoginEvidence) => void;
|
||||||
|
|
||||||
|
interface LocalLoginTransportOptions {
|
||||||
|
context: LocalStartupContext;
|
||||||
|
createSocket: (url: string) => WebSocket;
|
||||||
|
onFault: (error: unknown) => void;
|
||||||
|
record: EvidenceSink;
|
||||||
|
}
|
||||||
|
|
||||||
|
/** Local-only socket boundary; raw evidence precedes protocol decoding. */
|
||||||
|
export class LocalLoginTransport implements Transport {
|
||||||
|
private socket: WebSocket | null = null;
|
||||||
|
private generation = 0;
|
||||||
|
private openCallback?: () => void;
|
||||||
|
private messageCallback?: (frame: string) => void;
|
||||||
|
private closeCallback?: () => void;
|
||||||
|
|
||||||
|
constructor(private readonly options: LocalLoginTransportOptions) {}
|
||||||
|
|
||||||
|
connect(url: string): void {
|
||||||
|
try {
|
||||||
|
assertLocalConnectionTarget(this.options.context, url);
|
||||||
|
} catch (error) {
|
||||||
|
this.options.onFault(error);
|
||||||
|
throw error;
|
||||||
|
}
|
||||||
|
this.close();
|
||||||
|
let socket: WebSocket;
|
||||||
|
try {
|
||||||
|
socket = this.options.createSocket(url);
|
||||||
|
} catch (error) {
|
||||||
|
this.options.onFault(error);
|
||||||
|
throw error;
|
||||||
|
}
|
||||||
|
this.socket = socket;
|
||||||
|
const generation = this.generation;
|
||||||
|
socket.onopen = () => {
|
||||||
|
if (this.socket === socket) this.openCallback?.();
|
||||||
|
};
|
||||||
|
socket.onmessage = (event) => {
|
||||||
|
if (this.socket !== socket) return;
|
||||||
|
if (typeof event.data !== 'string') {
|
||||||
|
this.stopWithFault(new TypeError('Local login WebSocket requires text frames'));
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
try {
|
||||||
|
this.options.record({ kind: 'transport-receive', data: event.data });
|
||||||
|
} catch (error) {
|
||||||
|
this.stopWithFault(error);
|
||||||
|
throw error;
|
||||||
|
}
|
||||||
|
if (this.socket === socket) this.messageCallback?.(event.data);
|
||||||
|
};
|
||||||
|
socket.onclose = () => {
|
||||||
|
if (this.socket !== socket) return;
|
||||||
|
this.socket = null;
|
||||||
|
this.closeCallback?.();
|
||||||
|
};
|
||||||
|
socket.onerror = (event) => {
|
||||||
|
if (this.socket !== socket) return;
|
||||||
|
this.stopWithFault(event);
|
||||||
|
if (this.generation === generation) this.closeCallback?.();
|
||||||
|
};
|
||||||
|
}
|
||||||
|
|
||||||
|
send(frame: string): void {
|
||||||
|
const socket = this.socket;
|
||||||
|
if (!socket || socket.readyState !== 1) {
|
||||||
|
throw new Error('Local login WebSocket must be OPEN before send');
|
||||||
|
}
|
||||||
|
try {
|
||||||
|
socket.send(frame);
|
||||||
|
this.options.record({ kind: 'transport-send', data: frame });
|
||||||
|
} catch (error) {
|
||||||
|
this.stopWithFault(error);
|
||||||
|
throw error;
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
close(): void {
|
||||||
|
this.generation++;
|
||||||
|
const socket = this.socket;
|
||||||
|
this.socket = null;
|
||||||
|
if (!socket) return;
|
||||||
|
try {
|
||||||
|
socket.close();
|
||||||
|
} catch (error) {
|
||||||
|
this.options.onFault(error);
|
||||||
|
throw error;
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
onOpen(callback: () => void): void { this.openCallback = callback; }
|
||||||
|
onMessage(callback: (frame: string) => void): void { this.messageCallback = callback; }
|
||||||
|
onClose(callback: () => void): void { this.closeCallback = callback; }
|
||||||
|
|
||||||
|
private stopWithFault(error: unknown): void {
|
||||||
|
const socket = this.socket;
|
||||||
|
this.socket = null;
|
||||||
|
const cleanupErrors: unknown[] = [];
|
||||||
|
try {
|
||||||
|
socket?.close();
|
||||||
|
} catch (cleanupError) {
|
||||||
|
cleanupErrors.push(cleanupError);
|
||||||
|
}
|
||||||
|
this.options.onFault(error);
|
||||||
|
// Preserve both independent failures, including non-Error thrown values.
|
||||||
|
for (const cleanupError of cleanupErrors) this.options.onFault(cleanupError);
|
||||||
|
}
|
||||||
|
}
|
||||||
@@ -0,0 +1,230 @@
|
|||||||
|
import assert from 'node:assert/strict';
|
||||||
|
import { test } from 'node:test';
|
||||||
|
import { resolveLocalStartup } from '../../YouleNexus/assets/framework/config/local-startup.ts';
|
||||||
|
import { WireClient } from '../../YouleNexus/assets/framework/net/wire-client.ts';
|
||||||
|
import type { Clock } from '../../YouleNexus/assets/framework/net/heartbeat.ts';
|
||||||
|
|
||||||
|
class FakeSocket {
|
||||||
|
readyState = 0;
|
||||||
|
sent: string[] = [];
|
||||||
|
closes = 0;
|
||||||
|
sendError?: Error;
|
||||||
|
closeError?: Error;
|
||||||
|
onopen: ((event: Event) => void) | null = null;
|
||||||
|
onmessage: ((event: MessageEvent) => void) | null = null;
|
||||||
|
onclose: ((event: CloseEvent) => void) | null = null;
|
||||||
|
onerror: ((event: Event) => void) | null = null;
|
||||||
|
send(frame: string): void {
|
||||||
|
if (this.sendError) throw this.sendError;
|
||||||
|
this.sent.push(frame);
|
||||||
|
}
|
||||||
|
close(): void {
|
||||||
|
this.closes++;
|
||||||
|
if (this.closeError) throw this.closeError;
|
||||||
|
this.readyState = 3;
|
||||||
|
this.onclose?.({} as CloseEvent);
|
||||||
|
}
|
||||||
|
open(): void { this.readyState = 1; this.onopen?.({} as Event); }
|
||||||
|
message(data: unknown): void { this.onmessage?.({ data } as MessageEvent); }
|
||||||
|
}
|
||||||
|
|
||||||
|
async function setup(overrides: {
|
||||||
|
createSocket?: (url: string) => WebSocket;
|
||||||
|
record?: (event: { kind: string; data: unknown }) => void;
|
||||||
|
onFault?: (error: unknown) => void;
|
||||||
|
} = {}) {
|
||||||
|
const module = await import('../../YouleNexus/assets/scripts/local-platform-login/local-login-transport.ts')
|
||||||
|
.catch((error: unknown) => assert.fail(`LocalLoginTransport implementation must exist: ${String(error)}`));
|
||||||
|
const context = await resolveLocalStartup({
|
||||||
|
mode: 'debug', hostKind: 'h5', win: {}, search: '?profile=local',
|
||||||
|
fetcher: { async fetch(): Promise<never> { throw new Error('unexpected fetch'); } },
|
||||||
|
});
|
||||||
|
const sockets: FakeSocket[] = [];
|
||||||
|
const urls: string[] = [];
|
||||||
|
const faults: unknown[] = [];
|
||||||
|
const evidence: { kind: string; data: unknown }[] = [];
|
||||||
|
const transport = new module.LocalLoginTransport({
|
||||||
|
context,
|
||||||
|
createSocket: overrides.createSocket ?? ((url) => {
|
||||||
|
urls.push(url);
|
||||||
|
const socket = new FakeSocket();
|
||||||
|
sockets.push(socket);
|
||||||
|
return socket as unknown as WebSocket;
|
||||||
|
}),
|
||||||
|
record: overrides.record ?? ((event) => evidence.push(event)),
|
||||||
|
onFault: (error) => { faults.push(error); overrides.onFault?.(error); },
|
||||||
|
});
|
||||||
|
return { transport, context, sockets, urls, faults, evidence, url: context.config.servers[0]! };
|
||||||
|
}
|
||||||
|
|
||||||
|
test('requires OPEN and records only frames successfully sent', async () => {
|
||||||
|
const h = await setup();
|
||||||
|
assert.throws(() => h.transport.send('{}'), /open/i);
|
||||||
|
h.transport.connect(h.url);
|
||||||
|
assert.throws(() => h.transport.send('{}'), /open/i);
|
||||||
|
const socket = h.sockets[0]!;
|
||||||
|
socket.open();
|
||||||
|
h.transport.send('frame');
|
||||||
|
assert.deepEqual(socket.sent, ['frame']);
|
||||||
|
assert.deepEqual(h.evidence, [{ kind: 'transport-send', data: 'frame' }]);
|
||||||
|
const error = new Error('send failed');
|
||||||
|
socket.sendError = error;
|
||||||
|
assert.throws(() => h.transport.send('unsent'), (caught) => caught === error);
|
||||||
|
assert.equal(h.evidence.length, 1);
|
||||||
|
assert.equal(h.faults[h.faults.length - 1], error);
|
||||||
|
});
|
||||||
|
|
||||||
|
test('validates every connection target before constructing and reports the thrown error', async () => {
|
||||||
|
const h = await setup();
|
||||||
|
h.transport.connect(h.url);
|
||||||
|
assert.throws(() => h.transport.connect('wss://remote.example/'), (error) => {
|
||||||
|
assert.equal(h.faults[h.faults.length - 1], error);
|
||||||
|
return /loopback/i.test(String(error));
|
||||||
|
});
|
||||||
|
assert.deepEqual(h.urls, [h.url]);
|
||||||
|
});
|
||||||
|
|
||||||
|
test('socket constructor failure is reported before rethrowing the same object and cause', async () => {
|
||||||
|
const cause = new Error('root');
|
||||||
|
const error = Object.assign(new Error('constructor failed'), { cause });
|
||||||
|
const h = await setup({ createSocket: () => { throw error; } });
|
||||||
|
assert.throws(() => h.transport.connect(h.url), (caught) => {
|
||||||
|
assert.equal(h.faults[0], error);
|
||||||
|
return caught === error && error.cause === cause;
|
||||||
|
});
|
||||||
|
});
|
||||||
|
|
||||||
|
test('records raw handshake before message delivery', async () => {
|
||||||
|
const order: unknown[] = [];
|
||||||
|
const h = await setup({ record: (event) => order.push(event) });
|
||||||
|
h.transport.onMessage((frame) => order.push(frame));
|
||||||
|
h.transport.connect(h.url);
|
||||||
|
h.sockets[0]!.open();
|
||||||
|
h.sockets[0]!.message('@toconcon raw handshake');
|
||||||
|
assert.deepEqual(order, [
|
||||||
|
{ kind: 'transport-receive', data: '@toconcon raw handshake' }, '@toconcon raw handshake',
|
||||||
|
]);
|
||||||
|
});
|
||||||
|
|
||||||
|
test('non-text frames stop, report TypeError and never reach consumers or reconnect callback', async () => {
|
||||||
|
const h = await setup();
|
||||||
|
let delivered = 0;
|
||||||
|
let closed = 0;
|
||||||
|
h.transport.onMessage(() => delivered++);
|
||||||
|
h.transport.onClose(() => closed++);
|
||||||
|
h.transport.connect(h.url);
|
||||||
|
const socket = h.sockets[0]!;
|
||||||
|
socket.open();
|
||||||
|
socket.message(new ArrayBuffer(2));
|
||||||
|
socket.message('late frame');
|
||||||
|
assert.equal(delivered, 0);
|
||||||
|
assert.equal(closed, 0);
|
||||||
|
assert.equal(socket.closes, 1);
|
||||||
|
assert.equal(h.faults.length, 1);
|
||||||
|
assert.ok(h.faults[0] instanceof TypeError);
|
||||||
|
assert.deepEqual(h.evidence, []);
|
||||||
|
assert.throws(() => h.transport.send('late'), /open/i);
|
||||||
|
});
|
||||||
|
|
||||||
|
for (const direction of ['receive', 'send']) {
|
||||||
|
test(`${direction} evidence failure preserves original error and stops without reconnect`, async () => {
|
||||||
|
const error = new Error('evidence storage failed');
|
||||||
|
const h = await setup({ record: () => { throw error; } });
|
||||||
|
let delivered = 0;
|
||||||
|
let closed = 0;
|
||||||
|
h.transport.onMessage(() => delivered++);
|
||||||
|
h.transport.onClose(() => closed++);
|
||||||
|
h.transport.connect(h.url);
|
||||||
|
const socket = h.sockets[0]!;
|
||||||
|
socket.open();
|
||||||
|
const run = () => direction === 'receive' ? socket.message('frame') : h.transport.send('frame');
|
||||||
|
assert.throws(run, (caught) => caught === error);
|
||||||
|
assert.equal(h.faults[0], error);
|
||||||
|
assert.equal(socket.closes, 1);
|
||||||
|
assert.equal(delivered, 0);
|
||||||
|
assert.equal(closed, 0);
|
||||||
|
assert.throws(() => h.transport.send('late'), /open/i);
|
||||||
|
});
|
||||||
|
}
|
||||||
|
|
||||||
|
test('error followed by close notifies once and retains original error event', async () => {
|
||||||
|
const h = await setup();
|
||||||
|
let closed = 0;
|
||||||
|
h.transport.onClose(() => closed++);
|
||||||
|
h.transport.connect(h.url);
|
||||||
|
const socket = h.sockets[0]!;
|
||||||
|
const error = {} as Event;
|
||||||
|
socket.onerror?.(error);
|
||||||
|
socket.onclose?.({} as CloseEvent);
|
||||||
|
assert.equal(closed, 1);
|
||||||
|
assert.equal(h.faults[0], error);
|
||||||
|
});
|
||||||
|
|
||||||
|
test('explicit close and superseded sockets invalidate every old callback', async () => {
|
||||||
|
const h = await setup();
|
||||||
|
const calls: string[] = [];
|
||||||
|
h.transport.onOpen(() => calls.push('open'));
|
||||||
|
h.transport.onMessage(() => calls.push('message'));
|
||||||
|
h.transport.onClose(() => calls.push('close'));
|
||||||
|
h.transport.connect(h.url);
|
||||||
|
const old = h.sockets[0]!;
|
||||||
|
h.transport.connect(h.url);
|
||||||
|
old.open(); old.message('old'); old.onerror?.({} as Event); old.onclose?.({} as CloseEvent);
|
||||||
|
assert.equal(old.closes, 1);
|
||||||
|
const current = h.sockets[1]!;
|
||||||
|
current.open();
|
||||||
|
h.transport.close();
|
||||||
|
current.open(); current.message('late'); current.onerror?.({} as Event);
|
||||||
|
assert.deepEqual(calls, ['open']);
|
||||||
|
assert.deepEqual(h.faults, []);
|
||||||
|
assert.deepEqual(h.evidence, []);
|
||||||
|
});
|
||||||
|
|
||||||
|
test('close failure is reported and thrown after invalidation', async () => {
|
||||||
|
const h = await setup();
|
||||||
|
h.transport.connect(h.url);
|
||||||
|
const socket = h.sockets[0]!;
|
||||||
|
const error = new Error('close failed');
|
||||||
|
socket.closeError = error;
|
||||||
|
assert.throws(() => h.transport.close(), (caught) => caught === error);
|
||||||
|
assert.equal(h.faults[0], error);
|
||||||
|
assert.throws(() => h.transport.send('late'), /open/i);
|
||||||
|
});
|
||||||
|
|
||||||
|
test('explicit stop from fault observer suppresses the pending error close callback', async () => {
|
||||||
|
const h = await setup({ onFault: () => h.transport.close() });
|
||||||
|
let closed = 0;
|
||||||
|
h.transport.onClose(() => closed++);
|
||||||
|
h.transport.connect(h.url);
|
||||||
|
h.sockets[0]!.onerror?.({} as Event);
|
||||||
|
assert.equal(closed, 0);
|
||||||
|
});
|
||||||
|
|
||||||
|
test('a replacement connection from fault observer cannot receive the old close callback', async () => {
|
||||||
|
const h = await setup({ onFault: () => h.transport.connect(h.url) });
|
||||||
|
let closed = 0;
|
||||||
|
h.transport.onClose(() => closed++);
|
||||||
|
h.transport.connect(h.url);
|
||||||
|
h.sockets[0]!.onerror?.({} as Event);
|
||||||
|
assert.equal(h.sockets.length, 2);
|
||||||
|
assert.equal(closed, 0);
|
||||||
|
});
|
||||||
|
|
||||||
|
test('natural close lets real WireClient reconnect while explicit stop cancels it', async () => {
|
||||||
|
const h = await setup();
|
||||||
|
const timers = new Map<unknown, () => void>();
|
||||||
|
const clock: Clock = {
|
||||||
|
setTimeout: (fn) => { const key = {}; timers.set(key, fn); return key; },
|
||||||
|
clearTimeout: (key) => { timers.delete(key); },
|
||||||
|
};
|
||||||
|
const wire = new WireClient({ servers: h.context.config.servers, transportFactory: () => h.transport, clock });
|
||||||
|
wire.start();
|
||||||
|
h.sockets[0]!.open();
|
||||||
|
h.sockets[0]!.close();
|
||||||
|
assert.equal(timers.size, 1);
|
||||||
|
const retry = [...timers.values()][0]!;
|
||||||
|
timers.clear(); retry();
|
||||||
|
assert.equal(h.sockets.length, 2);
|
||||||
|
wire.stop();
|
||||||
|
assert.equal(timers.size, 0);
|
||||||
|
});
|
||||||
Reference in New Issue
Block a user