feat(platform): expose lifecycle-safe UI observation

This commit is contained in:
2026-09-05 11:50:12 +08:00
parent 2d7655a5d7
commit bd45503d47
2 changed files with 419 additions and 3 deletions
@@ -119,6 +119,8 @@ export class PlatformRuntime {
private fatalError: Error | null = null;
private wireStopped = false;
private startupGeneration = 0;
private readonly stateObservers = new Set<() => void>();
private observerDispatchError: Error | null = null;
constructor(private readonly options: PlatformRuntimeOptions) {
this.clock = options.clock ?? realClock;
@@ -132,6 +134,54 @@ export class PlatformRuntime {
return this.gate.ready;
}
subscribeState(
listener: (state: PlatformState) => void,
onError: (error: unknown) => void,
): () => void {
if (
this.lifecycle === 'kicked'
|| this.lifecycle === 'fatal'
|| this.lifecycle === 'stopped'
) {
throw new Error(`PlatformRuntime cannot subscribe while ${this.lifecycle}`);
}
let active = true;
const unsubscribeStore = this.store.subscribe((state) => {
if (!active) return;
try {
listener(state);
} catch (error) {
off();
try {
onError(error);
} catch (callbackError) {
const failure = this.failureError(
'UI observer',
error,
[callbackError],
);
if (this.isStopped()) {
this.observerDispatchError = failure;
} else {
this.showFatalOnce(failure);
}
}
}
});
const off = (): void => {
if (!active) return;
active = false;
try {
unsubscribeStore();
} finally {
this.stateObservers.delete(off);
}
};
this.stateObservers.add(off);
return off;
}
async start(): Promise<void> {
if (this.lifecycle !== 'idle') {
throw new Error(`PlatformRuntime cannot start while ${this.lifecycle}`);
@@ -260,6 +310,9 @@ export class PlatformRuntime {
this.wireIsOpen = false;
const errors: unknown[] = [];
for (const off of [...this.stateObservers]) {
this.captureCleanupError(errors, off);
}
this.captureCleanupError(errors, () => { this.clearLoginGuard(); });
const unsubscribe = this.unsubscribeWire;
this.unsubscribeWire = null;
@@ -324,6 +377,7 @@ export class PlatformRuntime {
case 'open':
this.wireIsOpen = true;
this.store.setConnectionPhase('connected');
this.throwObserverDispatchError();
this.markReady('socket');
if (this.connectionIntent.type === 'switch') {
const loginEnvelope = this.loginEnvelope;
@@ -341,8 +395,10 @@ export class PlatformRuntime {
this.clearLoginGuard();
}
const stateBeforeDispatch = this.store.getState();
const dispatchResult = this.requireRouter().dispatch(event.message);
this.throwObserverDispatchError();
if (
this.requireRouter().dispatch(event.message) === 'handled-platform'
dispatchResult === 'handled-platform'
&& this.lifecycle === 'running'
&& this.store.getState() !== stateBeforeDispatch
) {
@@ -353,6 +409,8 @@ export class PlatformRuntime {
case 'reconnecting':
this.wireIsOpen = false;
this.store.setConnectionPhase(event.type);
this.throwObserverDispatchError();
if (this.isTerminal()) return;
this.options.scene.showReconnect();
break;
case 'close':
@@ -360,10 +418,11 @@ export class PlatformRuntime {
break;
}
} catch (error) {
const failure = this.withObserverDispatchError(error);
// A lifecycle callback may stop and then throw. Preserve that error without
// replacing the requested terminal stop with a fatal lifecycle transition.
if (this.isStopped()) throw error;
this.showFatalOnce(error);
if (this.isStopped()) throw failure;
this.showFatalOnce(failure);
}
}
@@ -371,6 +430,31 @@ export class PlatformRuntime {
return this.lifecycle === 'stopped';
}
private isTerminal(): boolean {
return this.lifecycle === 'fatal'
|| this.lifecycle === 'kicked'
|| this.lifecycle === 'stopped';
}
private throwObserverDispatchError(): void {
const failure = this.takeObserverDispatchError();
if (failure === null) return;
throw failure;
}
private takeObserverDispatchError(): Error | null {
const failure = this.observerDispatchError;
this.observerDispatchError = null;
return failure;
}
private withObserverDispatchError(error: unknown): unknown {
const observerFailure = this.takeObserverDispatchError();
return observerFailure === null
? error
: this.failureError('UI observer dispatch', observerFailure, [error]);
}
private sendLoginEnvelope(envelope: OutboundEnvelope): void {
const wire = this.requireOpenWire();
const session = this.requireSession();
@@ -441,6 +525,9 @@ export class PlatformRuntime {
this.connectionIntent = { type: 'none' };
this.wireIsOpen = false;
const errors: unknown[] = [];
for (const off of [...this.stateObservers]) {
this.captureCleanupError(errors, off);
}
this.captureCleanupError(errors, () => { this.store.setConnectionPhase('kicked'); });
this.captureCleanupError(errors, () => { this.clearLoginGuard(); });
this.captureCleanupError(errors, () => { this.options.scene.showKicked(data); });
@@ -603,6 +690,9 @@ export class PlatformRuntime {
this.wireIsOpen = false;
const cleanupErrors: unknown[] = [];
for (const off of [...this.stateObservers]) {
this.captureCleanupError(cleanupErrors, off);
}
this.captureCleanupError(cleanupErrors, () => { this.clearLoginGuard(); });
const unsubscribe = this.unsubscribeWire;
@@ -304,6 +304,332 @@ function makeRuntime(config: RuntimeConfig = runtimeConfig(), entry = makeGameEn
return { runtime, wire, scene, clock, factoryConfigs, config };
}
test('UI observes canonical commits and unsubscribe is idempotent', async () => {
const { runtime, wire } = makeRuntime();
const seen: unknown[] = [];
const off = runtime.subscribeState(
state => { seen.push(state); },
error => { throw error; },
);
assert.equal(seen.length, 0);
await runtime.start();
wire.emit({ type: 'open', server: 'ws://agent' });
assert.equal(seen[seen.length - 1], runtime.state);
off();
off();
const count = seen.length;
runtime.login(ACCOUNT);
wire.emit({ type: 'message', message: {
route: 'agent', rpc: 'player_login', data: fixture('player-login-success.json'),
} });
assert.equal(seen.length, count);
runtime.stop();
assert.throws(() => runtime.subscribeState(() => {}, () => {}), /stopped/);
});
test('a failing UI observer reports its original error once while other observers continue', async () => {
const { runtime, wire } = makeRuntime();
const failure = { source: 'render' };
const reported: unknown[] = [];
let failingCalls = 0;
const healthyStates: unknown[] = [];
runtime.subscribeState(
() => {
failingCalls += 1;
throw failure;
},
error => { reported.push(error); },
);
runtime.subscribeState(
state => { healthyStates.push(state); },
error => { throw error; },
);
await runtime.start();
wire.emit({ type: 'open', server: 'ws://agent' });
assert.equal(failingCalls, 1);
assert.deepEqual(reported, [failure]);
assert.equal(reported[0], failure);
assert.equal(healthyStates[0], runtime.state);
runtime.login(ACCOUNT);
wire.emit({ type: 'message', message: {
route: 'agent', rpc: 'player_login', data: fixture('player-login-success.json'),
} });
assert.equal(failingCalls, 1);
assert.equal(healthyStates[healthyStates.length - 1], runtime.state);
runtime.stop();
});
test('an onError failure enters fatal with both UI error identities and cancels observers', async () => {
const { runtime, wire, scene } = makeRuntime();
const listenerFailure = { source: 'listener' };
const onErrorFailure = { source: 'onError' };
let laterObserverCalls = 0;
runtime.subscribeState(
() => { throw listenerFailure; },
() => { throw onErrorFailure; },
);
runtime.subscribeState(
() => { laterObserverCalls += 1; },
error => { throw error; },
);
await runtime.start();
wire.emit({ type: 'open', server: 'ws://agent' });
assert.equal(laterObserverCalls, 0);
assert.equal(wire.stopCalls, 1);
const fatal = scene.calls.find(call => call.name === 'fatal')?.value;
assert.ok(fatal instanceof Error);
const errors = (fatal as Error & { readonly errors: readonly unknown[] }).errors;
assert.equal(errors.length, 2);
assert.equal(errors[0], listenerFailure);
assert.equal(errors[1], onErrorFailure);
assert.throws(() => runtime.subscribeState(() => {}, () => {}), /fatal/);
});
for (const terminal of ['stop', 'fatal', 'kick'] as const) {
test(`${terminal} invalidates a late Store delivery and rejects new UI observers`, async () => {
const originalSubscribe = PlatformStore.prototype.subscribe;
let retainedListener: Parameters<PlatformStore['subscribe']>[0] | null = null;
PlatformStore.prototype.subscribe = function retainFirstListener(listener) {
retainedListener ??= listener;
return originalSubscribe.call(this, listener);
};
try {
const setup = makeRuntime();
let calls = 0;
setup.runtime.subscribeState(
() => { calls += 1; },
error => { throw error; },
);
await setup.runtime.start();
setup.wire.emit({ type: 'open', server: 'ws://agent' });
const callsBeforeTermination = calls;
setup.wire.onAction = (name) => {
if (name !== 'stop') return;
const listener = retainedListener;
assert.ok(listener);
listener(setup.runtime.state, setup.runtime.state);
};
if (terminal === 'stop') {
setup.runtime.stop();
} else if (terminal === 'fatal') {
setup.scene.onCall = (name) => {
if (name === 'reconnect') throw new Error('reconnect render failed');
};
setup.wire.emit({ type: 'slow' });
} else {
setup.runtime.login(ACCOUNT);
setup.wire.emit({ type: 'message', message: {
route: 'agent', rpc: 'kick_server', data: { code: 1, msg: 'terminal' },
} });
}
assert.equal(
calls,
callsBeforeTermination + (terminal === 'fatal' ? 1 : 0),
'observers must be inactive before terminal cleanup callbacks',
);
const callsAtTermination = calls;
const lateListener = retainedListener as Parameters<PlatformStore['subscribe']>[0] | null;
assert.ok(lateListener);
lateListener(setup.runtime.state, setup.runtime.state);
assert.equal(calls, callsAtTermination);
assert.throws(
() => setup.runtime.subscribeState(() => {}, () => {}),
new RegExp(terminal === 'kick' ? 'kicked' : terminal),
);
} finally {
PlatformStore.prototype.subscribe = originalSubscribe;
}
});
}
test('stop during a Store commit cancels later observers and prevents scene revival', async () => {
const setup = makeRuntime();
await setup.runtime.start();
setup.wire.emit({ type: 'open', server: 'ws://agent' });
let laterObserverCalls = 0;
let resubscribeFailure: unknown;
setup.runtime.subscribeState(
() => {
setup.runtime.stop();
resubscribeFailure = thrownValue(
() => setup.runtime.subscribeState(() => {}, () => {}),
);
},
error => { throw error; },
);
setup.runtime.subscribeState(
() => { laterObserverCalls += 1; },
error => { throw error; },
);
const sceneCallsBefore = setup.scene.calls.length;
setup.wire.emit({ type: 'slow' });
assert.equal(laterObserverCalls, 0);
assert.match(String(resubscribeFailure), /stopped/);
assert.equal(setup.scene.calls.length, sceneCallsBefore);
assert.equal(setup.wire.stopCalls, 1);
});
test('listener and onError failures after stopping escape the Store boundary without reviving fatal', async () => {
const setup = makeRuntime();
await setup.runtime.start();
setup.wire.emit({ type: 'open', server: 'ws://agent' });
const listenerFailure = { source: 'listener-after-stop' };
const onErrorFailure = { source: 'onError-after-stop' };
const reported: unknown[] = [];
setup.runtime.subscribeState(
() => {
setup.runtime.stop();
throw listenerFailure;
},
error => {
reported.push(error);
setup.runtime.stop();
throw onErrorFailure;
},
);
const sceneCallsBefore = setup.scene.calls.length;
const thrown = thrownValue(() => { setup.wire.emit({ type: 'slow' }); });
assert.deepEqual(reported, [listenerFailure]);
assert.ok(thrown instanceof Error);
const errors = (thrown as Error & { readonly errors: readonly unknown[] }).errors;
assert.equal(errors.length, 2);
assert.equal(errors[0], listenerFailure);
assert.equal(errors[1], onErrorFailure);
assert.equal(setup.scene.calls.length, sceneCallsBefore);
assert.equal(setup.scene.calls.filter(call => call.name === 'fatal').length, 0);
assert.equal(setup.wire.stopCalls, 1);
assert.throws(() => setup.runtime.subscribeState(() => {}, () => {}), /stopped/);
});
test('stopped observer errors survive a later active-room publish failure', async () => {
const setup = makeRuntime();
await setup.runtime.start();
setup.wire.emit({ type: 'open', server: 'ws://agent' });
setup.runtime.login(ACCOUNT);
const roomLogin = fixture('player-login-room.json');
(roomLogin.players as unknown[])[2] = {
playerid: 430512,
nickname: '同桌玩家',
avatar: 'http://b',
sex: 0,
bean: 10,
isprepare: 0,
onstate: 0,
};
setup.wire.emit({ type: 'message', message: {
route: 'agent', rpc: 'player_login', data: roomLogin,
} });
const listenerFailure = { source: 'room-listener-after-stop' };
const onErrorFailure = { source: 'room-onError-after-stop' };
setup.runtime.subscribeState(
() => {
setup.runtime.stop();
throw listenerFailure;
},
() => { throw onErrorFailure; },
);
const thrown = thrownValue(() => {
setup.wire.emit({ type: 'message', message: {
route: 'room', rpc: 'player_prepare', data: { seat: 2 },
} });
});
assert.ok(thrown instanceof Error);
const errors = (thrown as Error & { readonly errors: readonly unknown[] }).errors;
assert.equal(errors.length, 2);
const observerFailure = errors[0];
assert.ok(observerFailure instanceof Error);
const observerErrors = (
observerFailure as Error & { readonly errors: readonly unknown[] }
).errors;
assert.equal(observerErrors[0], listenerFailure);
assert.equal(observerErrors[1], onErrorFailure);
assert.match(String(errors[1]), /disposed game session/i);
assert.equal(setup.scene.calls.filter(call => call.name === 'fatal').length, 0);
assert.throws(() => setup.runtime.subscribeState(() => {}, () => {}), /stopped/);
});
test('UI observers receive each committed snapshot even when a listener commits reentrantly', async () => {
const setup = makeRuntime();
const delivered: unknown[] = [];
let firstCommit: unknown;
let latestDuringFirstCommit: unknown;
let reentered = false;
setup.runtime.subscribeState(
state => {
if (reentered) return;
reentered = true;
firstCommit = state;
setup.wire.emit({ type: 'slow' });
latestDuringFirstCommit = setup.runtime.state;
},
error => { throw error; },
);
setup.runtime.subscribeState(
state => { delivered.push(state); },
error => { throw error; },
);
await setup.runtime.start();
setup.wire.emit({ type: 'open', server: 'ws://agent' });
assert.notEqual(firstCommit, latestDuringFirstCommit);
assert.equal(delivered[0], firstCommit);
assert.equal(delivered[1], latestDuringFirstCommit);
assert.deepEqual(
delivered.map(state => (state as { app: { phase: string } }).app.phase),
['connected', 'slow'],
);
setup.runtime.stop();
});
test('cancelling and subscribing observers during a commit cannot alter that delivery snapshot', async () => {
const setup = makeRuntime();
let cancelledCalls = 0;
const replacementPhases: string[] = [];
let changedObservers = false;
let offCancelled = (): void => {};
setup.runtime.subscribeState(
() => {
if (changedObservers) return;
changedObservers = true;
offCancelled();
setup.runtime.subscribeState(
state => { replacementPhases.push(state.app.phase); },
error => { throw error; },
);
},
error => { throw error; },
);
offCancelled = setup.runtime.subscribeState(
() => { cancelledCalls += 1; },
error => { throw error; },
);
await setup.runtime.start();
setup.wire.emit({ type: 'open', server: 'ws://agent' });
assert.equal(cancelledCalls, 0);
assert.deepEqual(replacementPhases, []);
setup.wire.emit({ type: 'slow' });
assert.equal(cancelledCalls, 0);
assert.deepEqual(replacementPhases, ['slow']);
setup.runtime.stop();
});
test('a denied join retains the lobby root, reports the raw outcome and permits a valid retry', async () => {
const { runtime, wire, scene } = makeRuntime();
await runtime.start();