fix(platform): close startup subscription handoff gaps
This commit is contained in:
@@ -138,20 +138,20 @@ export class PlatformRuntime {
|
||||
const generation = this.startupGeneration + 1;
|
||||
this.startupGeneration = generation;
|
||||
this.lifecycle = 'starting';
|
||||
this.options.scene.showLoading();
|
||||
this.requireActiveStartup(generation);
|
||||
const resourcesReady = this.beginReadyTask(
|
||||
'resources',
|
||||
this.options.loadResources,
|
||||
generation,
|
||||
);
|
||||
const minimumDisplayReady = this.beginReadyTask(
|
||||
'minimum-display',
|
||||
this.options.waitForMinimumDisplay,
|
||||
generation,
|
||||
);
|
||||
|
||||
try {
|
||||
this.options.scene.showLoading();
|
||||
this.requireActiveStartup(generation);
|
||||
const resourcesReady = this.beginReadyTask(
|
||||
'resources',
|
||||
this.options.loadResources,
|
||||
generation,
|
||||
);
|
||||
const minimumDisplayReady = this.beginReadyTask(
|
||||
'minimum-display',
|
||||
this.options.waitForMinimumDisplay,
|
||||
generation,
|
||||
);
|
||||
this.requireActiveStartup(generation);
|
||||
const config = await this.options.resolveRuntimeConfig();
|
||||
this.requireActiveStartup(generation);
|
||||
@@ -204,7 +204,7 @@ export class PlatformRuntime {
|
||||
this.router = router;
|
||||
this.wire = wire;
|
||||
this.requireActiveStartup(generation);
|
||||
this.unsubscribeWire = wire.subscribe((event) => { this.onWireEvent(event); });
|
||||
this.subscribeToWire(wire, generation);
|
||||
this.requireActiveStartup(generation);
|
||||
this.markReady('config');
|
||||
this.lifecycle = 'running';
|
||||
@@ -443,6 +443,49 @@ export class PlatformRuntime {
|
||||
wire.stop();
|
||||
}
|
||||
|
||||
private subscribeToWire(wire: PlatformWireClient, generation: number): void {
|
||||
let actualUnsubscribe: (() => void) | null = null;
|
||||
let releaseRequested = false;
|
||||
let released = false;
|
||||
const temporaryOwner = (): void => {
|
||||
if (released) return;
|
||||
if (actualUnsubscribe === null) {
|
||||
releaseRequested = true;
|
||||
return;
|
||||
}
|
||||
released = true;
|
||||
const unsubscribe = actualUnsubscribe;
|
||||
actualUnsubscribe = null;
|
||||
unsubscribe();
|
||||
};
|
||||
|
||||
this.unsubscribeWire = temporaryOwner;
|
||||
const returnedUnsubscribe: unknown = wire.subscribe(
|
||||
(event) => { this.onWireEvent(event); },
|
||||
);
|
||||
if (typeof returnedUnsubscribe !== 'function') {
|
||||
throw new TypeError(
|
||||
'WireClient.subscribe must return an unsubscribe function',
|
||||
);
|
||||
}
|
||||
actualUnsubscribe = returnedUnsubscribe as () => void;
|
||||
if (!releaseRequested) return;
|
||||
|
||||
const cleanupErrors: unknown[] = [];
|
||||
this.captureCleanupError(cleanupErrors, temporaryOwner);
|
||||
if (cleanupErrors.length === 0) this.requireActiveStartup(generation);
|
||||
|
||||
const primary = this.fatalError
|
||||
?? new Error(`PlatformRuntime startup interrupted while ${this.lifecycle}`);
|
||||
const failure = this.failureError(
|
||||
'subscription handoff',
|
||||
primary,
|
||||
cleanupErrors,
|
||||
);
|
||||
if (this.fatalError !== null) this.fatalError = failure;
|
||||
throw failure;
|
||||
}
|
||||
|
||||
private captureCleanupError(errors: unknown[], operation: () => void): void {
|
||||
try {
|
||||
operation();
|
||||
|
||||
@@ -187,6 +187,68 @@ class RecordingWireClient implements PlatformWireClient {
|
||||
}
|
||||
}
|
||||
|
||||
class SynchronousSubscribeWireClient implements PlatformWireClient {
|
||||
readonly sent: OutboundEnvelope[] = [];
|
||||
readonly switchCalls: string[] = [];
|
||||
startCalls = 0;
|
||||
stopCalls = 0;
|
||||
unsubscribeCalls = 0;
|
||||
private listener: ((event: WireEvent) => void) | null = null;
|
||||
|
||||
constructor(
|
||||
private readonly subscribeEvent: WireEvent,
|
||||
private readonly unsubscribeFailure?: { readonly value: unknown },
|
||||
) {}
|
||||
|
||||
subscribe(listener: (event: WireEvent) => void): () => void {
|
||||
this.listener = listener;
|
||||
listener(this.subscribeEvent);
|
||||
return () => {
|
||||
this.unsubscribeCalls += 1;
|
||||
this.listener = null;
|
||||
if (this.unsubscribeFailure !== undefined) throw this.unsubscribeFailure.value;
|
||||
};
|
||||
}
|
||||
|
||||
start(): void {
|
||||
this.startCalls += 1;
|
||||
}
|
||||
|
||||
send(envelope: OutboundEnvelope): void {
|
||||
this.sent.push(envelope);
|
||||
}
|
||||
|
||||
switchServer(target: string): void {
|
||||
this.switchCalls.push(target);
|
||||
}
|
||||
|
||||
reconnectCurrent(): void {}
|
||||
|
||||
stop(): void {
|
||||
this.stopCalls += 1;
|
||||
}
|
||||
|
||||
emit(event: WireEvent): void {
|
||||
this.listener?.(event);
|
||||
}
|
||||
|
||||
get listenerAttached(): boolean {
|
||||
return this.listener !== null;
|
||||
}
|
||||
}
|
||||
|
||||
class InvalidSubscribeWireClient extends RecordingWireClient {
|
||||
constructor(private readonly invalidUnsubscribe: unknown) {
|
||||
super();
|
||||
}
|
||||
|
||||
override subscribe(_listener: (event: WireEvent) => void): () => void {
|
||||
this.lifecycle.push('subscribe');
|
||||
this.onAction('subscribe');
|
||||
return this.invalidUnsubscribe as () => void;
|
||||
}
|
||||
}
|
||||
|
||||
class RecordingScene implements ScenePort {
|
||||
readonly calls: Array<{ readonly name: string; readonly value?: unknown }> = [];
|
||||
onCall: (name: string, value?: unknown) => void = () => {};
|
||||
@@ -329,6 +391,86 @@ test('start resolves config, subscribes before connect, and never logs in withou
|
||||
assert.deepEqual(setup.scene.calls.map((call) => call.name), ['loading', 'login']);
|
||||
});
|
||||
|
||||
test('showLoading failure enters terminal fatal before any startup dependency runs', async () => {
|
||||
const primary = new Error('loading scene failed');
|
||||
const scene = new RecordingScene();
|
||||
const wire = new RecordingWireClient();
|
||||
let resourcesStarted = 0;
|
||||
let minimumDisplayStarted = 0;
|
||||
let configsResolved = 0;
|
||||
let wireFactories = 0;
|
||||
scene.onCall = (name) => {
|
||||
if (name === 'loading') throw primary;
|
||||
};
|
||||
const runtime = new PlatformRuntime({
|
||||
gameEntry: makeGameEntry(),
|
||||
resolveRuntimeConfig: async () => {
|
||||
configsResolved += 1;
|
||||
return runtimeConfig();
|
||||
},
|
||||
createWireClient: () => {
|
||||
wireFactories += 1;
|
||||
return wire;
|
||||
},
|
||||
scene,
|
||||
loadResources: async () => { resourcesStarted += 1; },
|
||||
waitForMinimumDisplay: async () => { minimumDisplayStarted += 1; },
|
||||
getLoginDeviceSnapshot: () => DEVICE,
|
||||
});
|
||||
|
||||
const thrown = await rejectedValue(runtime.start());
|
||||
|
||||
assert.equal(thrown, primary);
|
||||
assert.deepEqual(scene.calls.map((call) => call.name), ['loading', 'fatal']);
|
||||
assert.equal(scene.calls[1]?.value, primary);
|
||||
assert.equal(resourcesStarted, 0);
|
||||
assert.equal(minimumDisplayStarted, 0);
|
||||
assert.equal(configsResolved, 0);
|
||||
assert.equal(wireFactories, 0);
|
||||
assert.deepEqual(wire.lifecycle, []);
|
||||
assert.equal(runtime.ready, false);
|
||||
|
||||
await assert.rejects(() => runtime.start(), /fatal/i);
|
||||
assert.equal(scene.calls.filter((call) => call.name === 'fatal').length, 1);
|
||||
});
|
||||
|
||||
test('showLoading and showFatal failures preserve both error identities without starting dependencies', async () => {
|
||||
const primary = new Error('loading scene failed');
|
||||
const fatalSceneFailure = { step: 'scene.showFatal' };
|
||||
const scene = new RecordingScene();
|
||||
let dependencyCalls = 0;
|
||||
scene.onCall = (name) => {
|
||||
if (name === 'loading') throw primary;
|
||||
if (name === 'fatal') throw fatalSceneFailure;
|
||||
};
|
||||
const runtime = new PlatformRuntime({
|
||||
gameEntry: makeGameEntry(),
|
||||
resolveRuntimeConfig: async () => {
|
||||
dependencyCalls += 1;
|
||||
return runtimeConfig();
|
||||
},
|
||||
createWireClient: () => {
|
||||
dependencyCalls += 1;
|
||||
return new RecordingWireClient();
|
||||
},
|
||||
scene,
|
||||
loadResources: async () => { dependencyCalls += 1; },
|
||||
waitForMinimumDisplay: async () => { dependencyCalls += 1; },
|
||||
getLoginDeviceSnapshot: () => DEVICE,
|
||||
});
|
||||
|
||||
const thrown = await rejectedValue(runtime.start());
|
||||
assert.ok(thrown instanceof Error);
|
||||
const errors = (thrown as Error & { readonly errors: readonly unknown[] }).errors;
|
||||
|
||||
assert.equal(errors[0], primary);
|
||||
assert.equal(errors[1], fatalSceneFailure);
|
||||
assert.equal(dependencyCalls, 0);
|
||||
assert.deepEqual(scene.calls.map((call) => call.name), ['loading', 'fatal']);
|
||||
await assert.rejects(() => runtime.start(), /fatal/i);
|
||||
assert.equal(scene.calls.filter((call) => call.name === 'fatal').length, 1);
|
||||
});
|
||||
|
||||
test('stop while config is pending prevents every late startup continuation from reviving runtime', async () => {
|
||||
const configReady = deferred<RuntimeConfig>();
|
||||
const resourcesReady = deferred<void>();
|
||||
@@ -375,6 +517,109 @@ test('stop while config is pending prevents every late startup continuation from
|
||||
assert.deepEqual(scene.calls.map((call) => call.name), ['loading']);
|
||||
});
|
||||
|
||||
test('a synchronous subscribe callback that stops runtime hands the returned unsubscribe to teardown exactly once', async () => {
|
||||
const unsubscribeFailure = { step: 'late-unsubscribe' };
|
||||
const wire = new SynchronousSubscribeWireClient(
|
||||
{ type: 'slow' },
|
||||
{ value: unsubscribeFailure },
|
||||
);
|
||||
const scene = new RecordingScene();
|
||||
let runtime!: PlatformRuntime;
|
||||
runtime = new PlatformRuntime({
|
||||
gameEntry: makeGameEntry(),
|
||||
resolveRuntimeConfig: async () => runtimeConfig(),
|
||||
createWireClient: () => wire,
|
||||
scene,
|
||||
loadResources: async () => {},
|
||||
waitForMinimumDisplay: async () => {},
|
||||
getLoginDeviceSnapshot: () => DEVICE,
|
||||
});
|
||||
scene.onCall = (name) => {
|
||||
if (name === 'reconnect') runtime.stop();
|
||||
};
|
||||
|
||||
const thrown = await rejectedValue(runtime.start());
|
||||
assert.ok(thrown instanceof Error);
|
||||
const errors = (thrown as Error & { readonly errors: readonly unknown[] }).errors;
|
||||
|
||||
assert.match(String(errors[0]), /stopped/i);
|
||||
assert.equal(errors[1], unsubscribeFailure);
|
||||
assert.equal(wire.unsubscribeCalls, 1);
|
||||
assert.equal(wire.listenerAttached, false);
|
||||
assert.equal(wire.stopCalls, 1);
|
||||
assert.equal(wire.startCalls, 0);
|
||||
|
||||
const sceneCallCount = scene.calls.length;
|
||||
wire.emit({ type: 'open', server: 'ws://agent' });
|
||||
runtime.stop();
|
||||
assert.equal(scene.calls.length, sceneCallCount);
|
||||
assert.equal(wire.unsubscribeCalls, 1);
|
||||
assert.equal(wire.stopCalls, 1);
|
||||
});
|
||||
|
||||
test('a fatal synchronous subscribe callback still adopts unsubscribe and ignores late events', async () => {
|
||||
const primary = new Error('synchronous reconnect scene failed');
|
||||
const wire = new SynchronousSubscribeWireClient({ type: 'slow' });
|
||||
const scene = new RecordingScene();
|
||||
scene.onCall = (name) => {
|
||||
if (name === 'reconnect') throw primary;
|
||||
};
|
||||
const runtime = new PlatformRuntime({
|
||||
gameEntry: makeGameEntry(),
|
||||
resolveRuntimeConfig: async () => runtimeConfig(),
|
||||
createWireClient: () => wire,
|
||||
scene,
|
||||
loadResources: async () => {},
|
||||
waitForMinimumDisplay: async () => {},
|
||||
getLoginDeviceSnapshot: () => DEVICE,
|
||||
});
|
||||
|
||||
const thrown = await rejectedValue(runtime.start());
|
||||
|
||||
assert.equal(thrown, primary);
|
||||
assert.equal(wire.unsubscribeCalls, 1);
|
||||
assert.equal(wire.listenerAttached, false);
|
||||
assert.equal(wire.stopCalls, 1);
|
||||
assert.equal(wire.startCalls, 0);
|
||||
assert.deepEqual(
|
||||
scene.calls.filter((call) => call.name === 'fatal').map((call) => call.value),
|
||||
[primary],
|
||||
);
|
||||
|
||||
const sceneCallCount = scene.calls.length;
|
||||
wire.emit({ type: 'reconnecting', server: 'ws://agent' });
|
||||
assert.equal(scene.calls.length, sceneCallCount);
|
||||
});
|
||||
|
||||
for (const invalidUnsubscribe of [7, null, undefined] as const) {
|
||||
test(`subscribe returning ${String(invalidUnsubscribe)} is fatal before Wire start`, async () => {
|
||||
const wire = new InvalidSubscribeWireClient(invalidUnsubscribe);
|
||||
const scene = new RecordingScene();
|
||||
const runtime = new PlatformRuntime({
|
||||
gameEntry: makeGameEntry(),
|
||||
resolveRuntimeConfig: async () => runtimeConfig(),
|
||||
createWireClient: () => wire,
|
||||
scene,
|
||||
loadResources: async () => {},
|
||||
waitForMinimumDisplay: async () => {},
|
||||
getLoginDeviceSnapshot: () => DEVICE,
|
||||
});
|
||||
|
||||
const thrown = await rejectedValue(runtime.start());
|
||||
|
||||
assert.ok(thrown instanceof TypeError);
|
||||
assert.match(thrown.message, /subscribe.*unsubscribe.*function/i);
|
||||
assert.deepEqual(wire.lifecycle, ['subscribe', 'stop']);
|
||||
assert.equal(wire.stopCalls, 1);
|
||||
assert.deepEqual(
|
||||
scene.calls.filter((call) => call.name === 'fatal').map((call) => call.value),
|
||||
[thrown],
|
||||
);
|
||||
await assert.rejects(() => runtime.start(), /fatal/i);
|
||||
assert.equal(scene.calls.filter((call) => call.name === 'fatal').length, 1);
|
||||
});
|
||||
}
|
||||
|
||||
test('subscribe startup failure enters fatal once and stops the created wire without starting it', async () => {
|
||||
const setup = makeRuntime();
|
||||
const primary = new Error('subscribe failed');
|
||||
|
||||
Reference in New Issue
Block a user