diff --git a/apps/desktop/src/main/__tests__/runtime-host-session-observer.test.ts b/apps/desktop/src/main/__tests__/runtime-host-session-observer.test.ts index b022782d46..18cabd68ee 100644 --- a/apps/desktop/src/main/__tests__/runtime-host-session-observer.test.ts +++ b/apps/desktop/src/main/__tests__/runtime-host-session-observer.test.ts @@ -1134,7 +1134,7 @@ test('does not let one backpressured transcript consumer block another', async ( await observer.close(); }); -test('releases an idle session when transcript delivery fails', async () => { +test('keeps a transcript consumer available after a delivery fails', async () => { const events = new AsyncFrameQueue(); let closeCount = 0; const observer = new RuntimeHostSessionObserver({ @@ -1176,12 +1176,18 @@ test('releases an idle session when transcript delivery fails', async () => { }, emitSessionsChanged() {}, }); - let opened = false; + let failDelivery = false; + let failedDeliveries = 0; + let successfulDeliveries = 0; const consumerId = 'consumer-failing'; - await observer.openTranscript('session-1', consumerId, { + const opened = await observer.openTranscript('session-1', consumerId, { id: 25, send(_channel, batch) { - if (opened) throw new Error('renderer unavailable'); + if (failDelivery) { + failedDeliveries += 1; + throw new Error('renderer unavailable'); + } + successfulDeliveries += 1; queueMicrotask(() => observer.acknowledgeTranscript( consumerId, @@ -1194,7 +1200,7 @@ test('releases an idle session when transcript delivery fails', async () => { once() {}, off() {}, }); - opened = true; + failDelivery = true; events.push({ kind: 'subscription.transcript_advanced', hostEpoch: 'host-1', @@ -1204,6 +1210,22 @@ test('releases an idle session when transcript delivery fails', async () => { throughSequence: 0, }); + await waitFor(() => failedDeliveries === 1); + failDelivery = false; + await assert.doesNotReject( + observer.loadTranscriptAround( + { + consumerId, + generation: opened.generation, + anchorSequence: 0, + maxBytes: DESKTOP_TRANSCRIPT_FRAGMENT_MAX_BYTES, + }, + 25, + ), + ); + assert.ok(successfulDeliveries > 1); + assert.equal(closeCount, 0); + await observer.closeTranscript(consumerId, 25); await waitFor(() => closeCount === 1); await observer.close(); }); diff --git a/apps/desktop/src/main/runtime-host-session-observer.ts b/apps/desktop/src/main/runtime-host-session-observer.ts index 815b960a5a..390ee5f34e 100644 --- a/apps/desktop/src/main/runtime-host-session-observer.ts +++ b/apps/desktop/src/main/runtime-host-session-observer.ts @@ -124,7 +124,6 @@ interface ObservedSessionState { interface TranscriptConsumer { readonly consumerId: string; readonly target: RuntimeHostTranscriptTarget; - readonly destroyedListener: () => void; generation: string; deliverySequence: number; deliveryBytes: number; @@ -281,13 +280,9 @@ export class RuntimeHostSessionObserver { void this.#closeIfIdle(state); } } - const destroyedListener = () => { - void this.closeTranscript(consumerId); - }; const consumer: TranscriptConsumer = { consumerId, target, - destroyedListener, generation: replica.generation, deliverySequence: 0, deliveryBytes: 0, @@ -296,7 +291,6 @@ export class RuntimeHostSessionObserver { }; state.transcriptConsumers.set(consumerId, consumer); this.#transcriptConsumers.set(consumerId, state); - target.once('destroyed', destroyedListener); try { consumer.resetRequested = true; await this.#scheduleTranscriptDelivery(state, consumer); @@ -1081,8 +1075,7 @@ export class RuntimeHostSessionObserver { if (consumer.generation !== replica.generation) { this.#requestTranscriptReset(state, consumer); } else if (!this.#mergeTranscriptChange(consumer, change)) { - this.#detachTranscriptConsumer(state, consumer); - void this.#closeIfIdle(state); + this.#requestTranscriptReset(state, consumer); } else { void this.#scheduleTranscriptDelivery(state, consumer).catch(() => undefined); } @@ -1156,8 +1149,9 @@ export class RuntimeHostSessionObserver { } } } catch (error) { - this.#detachTranscriptConsumer(state, consumer); - void this.#closeIfIdle(state); + if (state.transcriptConsumers.get(consumer.consumerId) === consumer) { + consumer.resetRequested = true; + } throw error; } })().finally(() => { @@ -1362,7 +1356,6 @@ export class RuntimeHostSessionObserver { pending.reject(new Error('Desktop transcript consumer was closed')); } consumer.pendingDeliveries.clear(); - consumer.target.off('destroyed', consumer.destroyedListener); } #touchReplica(state: ObservedSessionState, protectedState?: ObservedSessionState): boolean {