diff --git a/cocoscreator_projects/YouleNexus/assets/framework/platform/runtime.ts b/cocoscreator_projects/YouleNexus/assets/framework/platform/runtime.ts index 169fd7a..a73995c 100644 --- a/cocoscreator_projects/YouleNexus/assets/framework/platform/runtime.ts +++ b/cocoscreator_projects/YouleNexus/assets/framework/platform/runtime.ts @@ -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 { 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; diff --git a/cocoscreator_projects/framework-tests/platform/runtime.test.ts b/cocoscreator_projects/framework-tests/platform/runtime.test.ts index 9b2fb99..cfc56dc 100644 --- a/cocoscreator_projects/framework-tests/platform/runtime.test.ts +++ b/cocoscreator_projects/framework-tests/platform/runtime.test.ts @@ -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[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[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();