fix(queue): isolate process lifecycle per item
Track Streamlink and FFmpeg resources by queue item so parallel jobs can pause, resume, cancel, retry, and clean up independently. Gate queue scheduling during shutdown, await the active queue owner, retain retry artifacts on ordinary failures, remove partial outputs on cancellation, and flush queue state once after active work settles.
This commit is contained in:
@@ -0,0 +1,160 @@
|
||||
import { describe, expect, it, vi } from 'vitest';
|
||||
import { QueueProcessRegistry, QueueRunLifecycle, type QueueProcessResource } from './process-registry';
|
||||
|
||||
function deferred(): { promise: Promise<void>; resolve: () => void } {
|
||||
let resolve!: () => void;
|
||||
const promise = new Promise<void>((done) => {
|
||||
resolve = done;
|
||||
});
|
||||
return { promise, resolve };
|
||||
}
|
||||
|
||||
function createResource(wait: Promise<void> = Promise.resolve()): QueueProcessResource {
|
||||
return {
|
||||
kill: vi.fn(),
|
||||
wait,
|
||||
pause: vi.fn(),
|
||||
resume: vi.fn(),
|
||||
cancel: vi.fn(async () => undefined),
|
||||
cleanup: vi.fn(async () => undefined),
|
||||
};
|
||||
}
|
||||
|
||||
describe('QueueProcessRegistry', () => {
|
||||
it('keeps parallel queue item process groups independent', async () => {
|
||||
const registry = new QueueProcessRegistry();
|
||||
const first = createResource();
|
||||
const second = createResource();
|
||||
|
||||
registry.register('item-a', 'merge', first);
|
||||
registry.register('item-b', 'split', second);
|
||||
|
||||
await registry.cancelItem('item-a');
|
||||
|
||||
expect(first.kill).toHaveBeenCalledOnce();
|
||||
expect(first.cleanup).toHaveBeenCalledOnce();
|
||||
expect(second.kill).not.toHaveBeenCalled();
|
||||
expect(second.cleanup).not.toHaveBeenCalled();
|
||||
expect(registry.activeItemIds()).toEqual(['item-b']);
|
||||
});
|
||||
|
||||
it('pauses and resumes every resource for only the selected item', async () => {
|
||||
const registry = new QueueProcessRegistry();
|
||||
const streamlink = createResource();
|
||||
const postProcessing = createResource();
|
||||
const unrelated = createResource();
|
||||
|
||||
registry.register('item-a', 'streamlink', streamlink);
|
||||
registry.register('item-a', 'post-processing', postProcessing);
|
||||
registry.register('item-b', 'streamlink', unrelated);
|
||||
|
||||
await registry.pauseItem('item-a');
|
||||
let resumed = false;
|
||||
const resumedSignal = registry.whenResumed('item-a').then(() => {
|
||||
resumed = true;
|
||||
});
|
||||
await Promise.resolve();
|
||||
|
||||
expect(registry.isPaused('item-a')).toBe(true);
|
||||
expect(resumed).toBe(false);
|
||||
|
||||
await registry.resumeItem('item-a');
|
||||
await resumedSignal;
|
||||
|
||||
expect(streamlink.pause).toHaveBeenCalledOnce();
|
||||
expect(streamlink.resume).toHaveBeenCalledOnce();
|
||||
expect(postProcessing.pause).toHaveBeenCalledOnce();
|
||||
expect(postProcessing.resume).toHaveBeenCalledOnce();
|
||||
expect(unrelated.pause).not.toHaveBeenCalled();
|
||||
expect(unrelated.resume).not.toHaveBeenCalled();
|
||||
expect(registry.isPaused('item-a')).toBe(false);
|
||||
});
|
||||
|
||||
it.each(['merge', 'split'] as const)('waits for %s termination before removing partial output', async (phase) => {
|
||||
const registry = new QueueProcessRegistry();
|
||||
const closed = deferred();
|
||||
const resource = createResource(closed.promise);
|
||||
|
||||
registry.register('item-a', phase, resource);
|
||||
const cancelling = registry.cancelItem('item-a');
|
||||
|
||||
await Promise.resolve();
|
||||
expect(resource.kill).toHaveBeenCalledOnce();
|
||||
expect(resource.cleanup).not.toHaveBeenCalled();
|
||||
|
||||
closed.resolve();
|
||||
await cancelling;
|
||||
|
||||
expect(resource.cancel).toHaveBeenCalledOnce();
|
||||
expect(resource.cleanup).toHaveBeenCalledOnce();
|
||||
expect(registry.activeItemIds()).toEqual([]);
|
||||
});
|
||||
|
||||
it('allows an explicitly reset item to retry without affecting another item', async () => {
|
||||
const registry = new QueueProcessRegistry();
|
||||
const firstAttempt = createResource();
|
||||
const other = createResource();
|
||||
|
||||
registry.register('item-a', 'streamlink', firstAttempt);
|
||||
registry.register('item-b', 'streamlink', other);
|
||||
await registry.cancelItem('item-a');
|
||||
|
||||
registry.resetItem('item-a');
|
||||
const retry = createResource();
|
||||
const registration = registry.register('item-a', 'streamlink', retry);
|
||||
|
||||
expect(registration.accepted).toBe(true);
|
||||
expect(registry.isCancelled('item-a')).toBe(false);
|
||||
expect(registry.activeItemIds()).toEqual(['item-b', 'item-a']);
|
||||
expect(other.kill).not.toHaveBeenCalled();
|
||||
});
|
||||
|
||||
it('releases completed phase artifacts without deleting retry state', () => {
|
||||
const registry = new QueueProcessRegistry();
|
||||
const downloadedPart = createResource();
|
||||
const mergedIntermediate = createResource();
|
||||
|
||||
registry.register('item-a', 'post-processing', downloadedPart);
|
||||
registry.register('item-a', 'post-processing', mergedIntermediate);
|
||||
registry.releaseItem('item-a');
|
||||
|
||||
expect(downloadedPart.cleanup).not.toHaveBeenCalled();
|
||||
expect(mergedIntermediate.cleanup).not.toHaveBeenCalled();
|
||||
expect(registry.activeItemIds()).toEqual([]);
|
||||
});
|
||||
});
|
||||
|
||||
describe('QueueRunLifecycle', () => {
|
||||
it('latches shutdown, cancels all groups, waits for the queue and persists once', async () => {
|
||||
const registry = new QueueProcessRegistry();
|
||||
const lifecycle = new QueueRunLifecycle(registry);
|
||||
const queueDone = deferred();
|
||||
const processClosed = deferred();
|
||||
const resource = createResource(processClosed.promise);
|
||||
const persist = vi.fn(async () => undefined);
|
||||
const beforeCancel = vi.fn(async () => undefined);
|
||||
|
||||
expect(lifecycle.schedule(async () => queueDone.promise)).toBe(true);
|
||||
registry.register('item-a', 'merge', resource);
|
||||
|
||||
const firstShutdown = lifecycle.shutdown(beforeCancel, persist);
|
||||
const secondShutdown = lifecycle.shutdown(beforeCancel, persist);
|
||||
const late = createResource();
|
||||
const lateRegistration = registry.register('item-late', 'streamlink', late);
|
||||
|
||||
expect(lifecycle.schedule(async () => undefined)).toBe(false);
|
||||
expect(lateRegistration.accepted).toBe(false);
|
||||
|
||||
processClosed.resolve();
|
||||
queueDone.resolve();
|
||||
await Promise.all([firstShutdown, secondShutdown]);
|
||||
|
||||
expect(beforeCancel).toHaveBeenCalledOnce();
|
||||
expect(resource.kill).toHaveBeenCalledOnce();
|
||||
expect(resource.cleanup).toHaveBeenCalledOnce();
|
||||
expect(late.kill).toHaveBeenCalledOnce();
|
||||
expect(late.cleanup).toHaveBeenCalledOnce();
|
||||
expect(persist).toHaveBeenCalledOnce();
|
||||
expect(registry.activeItemIds()).toEqual([]);
|
||||
});
|
||||
});
|
||||
@@ -0,0 +1,209 @@
|
||||
export type QueueProcessPhase = 'streamlink' | 'merge' | 'split' | 'post-processing';
|
||||
|
||||
export interface QueueProcessResource {
|
||||
kill?: () => unknown;
|
||||
wait?: Promise<unknown>;
|
||||
pause?: () => unknown | Promise<unknown>;
|
||||
resume?: () => unknown | Promise<unknown>;
|
||||
cancel?: () => unknown | Promise<unknown>;
|
||||
cleanup?: () => unknown | Promise<unknown>;
|
||||
}
|
||||
|
||||
export interface QueueProcessRegistration {
|
||||
accepted: boolean;
|
||||
release: () => void;
|
||||
}
|
||||
|
||||
interface RegisteredResource {
|
||||
itemId: string;
|
||||
phase: QueueProcessPhase;
|
||||
resource: QueueProcessResource;
|
||||
stopping: Promise<void> | null;
|
||||
}
|
||||
|
||||
export class QueueProcessRegistry {
|
||||
private readonly groups = new Map<string, Set<RegisteredResource>>();
|
||||
private readonly cancelledItems = new Set<string>();
|
||||
private readonly pausedItems = new Set<string>();
|
||||
private readonly cancellationWaiters = new Map<string, Set<() => void>>();
|
||||
private readonly resumeWaiters = new Map<string, Set<() => void>>();
|
||||
private readonly settling = new Set<Promise<void>>();
|
||||
private shuttingDown = false;
|
||||
|
||||
register(itemId: string, phase: QueueProcessPhase, resource: QueueProcessResource): QueueProcessRegistration {
|
||||
const entry: RegisteredResource = { itemId, phase, resource, stopping: null };
|
||||
if (this.shuttingDown || this.cancelledItems.has(itemId)) {
|
||||
this.trackSettlement(this.stopResource(entry));
|
||||
return { accepted: false, release: () => undefined };
|
||||
}
|
||||
|
||||
let group = this.groups.get(itemId);
|
||||
if (!group) {
|
||||
group = new Set();
|
||||
this.groups.set(itemId, group);
|
||||
}
|
||||
group.add(entry);
|
||||
|
||||
return {
|
||||
accepted: true,
|
||||
release: () => this.release(entry),
|
||||
};
|
||||
}
|
||||
|
||||
async pauseItem(itemId: string): Promise<void> {
|
||||
this.pausedItems.add(itemId);
|
||||
await this.invokeItem(itemId, 'pause');
|
||||
}
|
||||
|
||||
async resumeItem(itemId: string): Promise<void> {
|
||||
await this.invokeItem(itemId, 'resume');
|
||||
this.pausedItems.delete(itemId);
|
||||
this.resolveWaiters(this.resumeWaiters, itemId);
|
||||
}
|
||||
|
||||
async cancelItem(itemId: string): Promise<void> {
|
||||
this.cancelledItems.add(itemId);
|
||||
this.pausedItems.delete(itemId);
|
||||
this.resolveWaiters(this.cancellationWaiters, itemId);
|
||||
this.resolveWaiters(this.resumeWaiters, itemId);
|
||||
const entries = [...(this.groups.get(itemId) || [])];
|
||||
await Promise.all(entries.map((entry) => this.stopResource(entry)));
|
||||
if ((this.groups.get(itemId)?.size || 0) === 0) this.groups.delete(itemId);
|
||||
}
|
||||
|
||||
releaseItem(itemId: string): void {
|
||||
this.groups.delete(itemId);
|
||||
this.cancellationWaiters.delete(itemId);
|
||||
this.resumeWaiters.delete(itemId);
|
||||
this.pausedItems.delete(itemId);
|
||||
}
|
||||
|
||||
resetItem(itemId: string): void {
|
||||
if ((this.groups.get(itemId)?.size || 0) > 0) {
|
||||
throw new Error(`Queue item ${itemId} still has active resources`);
|
||||
}
|
||||
this.cancelledItems.delete(itemId);
|
||||
this.pausedItems.delete(itemId);
|
||||
this.cancellationWaiters.delete(itemId);
|
||||
this.resumeWaiters.delete(itemId);
|
||||
}
|
||||
|
||||
isCancelled(itemId: string): boolean {
|
||||
return this.cancelledItems.has(itemId);
|
||||
}
|
||||
|
||||
isPaused(itemId: string): boolean {
|
||||
return this.pausedItems.has(itemId);
|
||||
}
|
||||
|
||||
whenCancelled(itemId: string): Promise<void> {
|
||||
if (this.cancelledItems.has(itemId)) return Promise.resolve();
|
||||
return this.waitFor(this.cancellationWaiters, itemId);
|
||||
}
|
||||
|
||||
whenResumed(itemId: string): Promise<void> {
|
||||
if (!this.pausedItems.has(itemId)) return Promise.resolve();
|
||||
return this.waitFor(this.resumeWaiters, itemId);
|
||||
}
|
||||
|
||||
beginShutdown(): void {
|
||||
this.shuttingDown = true;
|
||||
}
|
||||
|
||||
async cancelAll(): Promise<void> {
|
||||
await Promise.all(this.activeItemIds().map((itemId) => this.cancelItem(itemId)));
|
||||
}
|
||||
|
||||
async waitForIdle(): Promise<void> {
|
||||
while (this.settling.size > 0) {
|
||||
await Promise.allSettled([...this.settling]);
|
||||
}
|
||||
}
|
||||
|
||||
activeItemIds(): string[] {
|
||||
return [...this.groups.entries()]
|
||||
.filter(([, entries]) => entries.size > 0)
|
||||
.map(([itemId]) => itemId);
|
||||
}
|
||||
|
||||
private async invokeItem(itemId: string, operation: 'pause' | 'resume'): Promise<void> {
|
||||
const entries = [...(this.groups.get(itemId) || [])];
|
||||
await Promise.allSettled(entries.map(async ({ resource }) => {
|
||||
await resource[operation]?.();
|
||||
}));
|
||||
}
|
||||
|
||||
private stopResource(entry: RegisteredResource): Promise<void> {
|
||||
if (entry.stopping) return entry.stopping;
|
||||
entry.stopping = (async () => {
|
||||
try { entry.resource.kill?.(); } catch { }
|
||||
try { await entry.resource.cancel?.(); } catch { }
|
||||
try { await entry.resource.wait; } catch { }
|
||||
try { await entry.resource.cleanup?.(); } catch { }
|
||||
this.release(entry);
|
||||
})();
|
||||
return entry.stopping;
|
||||
}
|
||||
|
||||
private trackSettlement(settlement: Promise<void>): void {
|
||||
this.settling.add(settlement);
|
||||
void settlement.finally(() => this.settling.delete(settlement));
|
||||
}
|
||||
|
||||
private release(entry: RegisteredResource): void {
|
||||
const group = this.groups.get(entry.itemId);
|
||||
if (!group) return;
|
||||
group.delete(entry);
|
||||
if (group.size === 0) this.groups.delete(entry.itemId);
|
||||
}
|
||||
|
||||
private waitFor(waiterMap: Map<string, Set<() => void>>, itemId: string): Promise<void> {
|
||||
return new Promise((resolve) => {
|
||||
let waiters = waiterMap.get(itemId);
|
||||
if (!waiters) {
|
||||
waiters = new Set();
|
||||
waiterMap.set(itemId, waiters);
|
||||
}
|
||||
waiters.add(resolve);
|
||||
});
|
||||
}
|
||||
|
||||
private resolveWaiters(waiterMap: Map<string, Set<() => void>>, itemId: string): void {
|
||||
const waiters = waiterMap.get(itemId);
|
||||
if (!waiters) return;
|
||||
waiterMap.delete(itemId);
|
||||
for (const resolve of waiters) resolve();
|
||||
}
|
||||
}
|
||||
|
||||
export class QueueRunLifecycle {
|
||||
private currentRun: Promise<void> | null = null;
|
||||
private shutdownRun: Promise<void> | null = null;
|
||||
|
||||
constructor(private readonly registry: QueueProcessRegistry) { }
|
||||
|
||||
schedule(run: () => Promise<void>, onError?: (error: unknown) => void): boolean {
|
||||
if (this.shutdownRun || this.currentRun) return false;
|
||||
const pending = Promise.resolve()
|
||||
.then(run)
|
||||
.catch((error) => onError?.(error));
|
||||
const tracked = pending.finally(() => {
|
||||
if (this.currentRun === tracked) this.currentRun = null;
|
||||
});
|
||||
this.currentRun = tracked;
|
||||
return true;
|
||||
}
|
||||
|
||||
shutdown(beforeCancel: () => unknown | Promise<unknown>, persist: () => unknown | Promise<unknown>): Promise<void> {
|
||||
if (this.shutdownRun) return this.shutdownRun;
|
||||
this.registry.beginShutdown();
|
||||
this.shutdownRun = (async () => {
|
||||
try { await beforeCancel(); } catch { }
|
||||
await this.registry.cancelAll();
|
||||
if (this.currentRun) await this.currentRun;
|
||||
await this.registry.waitForIdle();
|
||||
await persist();
|
||||
})();
|
||||
return this.shutdownRun;
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user