From 26512a57a50b084b3b0bdcf33479514e31a0a925 Mon Sep 17 00:00:00 2001 From: Shahbaz !! Date: Fri, 21 Aug 2026 07:09:43 +0000 Subject: [PATCH] fix: prevent duplicate stream terminal events --- src/typings/workers/worker.types.ts | 1 + src/workers/main.ts | 36 +++++++++++++++---------- src/workers/streamFinalizer.ts | 42 +++++++++++++++++++++++++++++ 3 files changed, 65 insertions(+), 14 deletions(-) create mode 100644 src/workers/streamFinalizer.ts diff --git a/src/typings/workers/worker.types.ts b/src/typings/workers/worker.types.ts index 6b2370e0..bb1a1cd9 100644 --- a/src/typings/workers/worker.types.ts +++ b/src/typings/workers/worker.types.ts @@ -104,6 +104,7 @@ export interface ActiveStreamEntry { pcmStream: PCMStream fetched: TrackStreamResult & { type?: string } cancelled: boolean + cleaned: boolean } /** diff --git a/src/workers/main.ts b/src/workers/main.ts index 049003b4..404953bb 100644 --- a/src/workers/main.ts +++ b/src/workers/main.ts @@ -64,6 +64,7 @@ import { enqueueHeadQueue, getHeadQueueLength } from './headQueue.ts' +import { createStreamFinalizer } from './streamFinalizer.ts' type WorkerPlayerClass = typeof import('../playback/player.ts').Player type CreatePCMStreamFn = @@ -1890,7 +1891,9 @@ function cleanupActiveStream( entry?: ActiveStreamEntry ): void { const current = entry || activeStreams.get(streamId) - if (!current) return + if (!current || current.cleaned) return + + current.cleaned = true if (current.pcmStream && !current.pcmStream.destroyed) { current.pcmStream.destroy() @@ -1953,27 +1956,31 @@ async function startLoadStream( payload?.filters || {} ) as unknown as PCMStream - const entry: ActiveStreamEntry = { pcmStream, fetched, cancelled: false } + const entry: ActiveStreamEntry = { + pcmStream, + fetched, + cancelled: false, + cleaned: false + } activeStreams.set(streamId, entry) streamLifecycle.created++ - const finish = (err?: unknown) => { - if (entry.cancelled) { + const finish = createStreamFinalizer(entry, { + onCancelled: () => { streamLifecycle.cancelled++ - cleanupActiveStream(streamId, entry) - return - } - - if (err) { + }, + onError: (error) => { streamLifecycle.errored++ - sendStreamError(streamId, getErrorMessage(err)) - } else { + sendStreamError(streamId, getErrorMessage(error)) + }, + onEnd: () => { streamLifecycle.ended++ sendStreamEnd(streamId) + }, + onCleanup: () => { + cleanupActiveStream(streamId, entry) } - - cleanupActiveStream(streamId, entry) - } + }) pcmStream.on('data', (chunk) => { if (!entry.cancelled) sendStreamChunk(streamId, chunk) @@ -1988,6 +1995,7 @@ function cancelStream(streamId: string): boolean { const entry = activeStreams.get(streamId) if (!entry) return false entry.cancelled = true + streamLifecycle.cancelled++ cleanupActiveStream(streamId, entry) return true } diff --git a/src/workers/streamFinalizer.ts b/src/workers/streamFinalizer.ts new file mode 100644 index 00000000..47a350de --- /dev/null +++ b/src/workers/streamFinalizer.ts @@ -0,0 +1,42 @@ +export interface StreamFinalizerEntry { + cancelled: boolean + cleaned: boolean +} + +export interface StreamFinalizerHandlers { + onCancelled: () => void + onError: (error: unknown) => void + onEnd: () => void + onCleanup: () => void +} + +/** + * Creates an idempotent terminal handler for a PCM stream. + * + * Streams commonly emit `end` followed by `close`, and cancellation can + * destroy a stream before its terminal event arrives. Both cases must produce + * one terminal notification and one cleanup. + */ +export function createStreamFinalizer( + entry: StreamFinalizerEntry, + handlers: StreamFinalizerHandlers +): (error?: unknown) => void { + let finished = false + + return (error?: unknown): void => { + if (finished || entry.cleaned) return + finished = true + + try { + if (entry.cancelled) { + handlers.onCancelled() + } else if (error) { + handlers.onError(error) + } else { + handlers.onEnd() + } + } finally { + handlers.onCleanup() + } + } +}