Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -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({
Expand Down Expand Up @@ -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,
Expand All @@ -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',
Expand All @@ -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();
});
Expand Down
15 changes: 4 additions & 11 deletions apps/desktop/src/main/runtime-host-session-observer.ts
Original file line number Diff line number Diff line change
Expand Up @@ -124,7 +124,6 @@ interface ObservedSessionState {
interface TranscriptConsumer {
readonly consumerId: string;
readonly target: RuntimeHostTranscriptTarget;
readonly destroyedListener: () => void;
generation: string;
deliverySequence: number;
deliveryBytes: number;
Expand Down Expand Up @@ -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,
Expand All @@ -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);
Expand Down Expand Up @@ -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);
}
Expand Down Expand Up @@ -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(() => {
Expand Down Expand Up @@ -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 {
Expand Down