import { BoundedTaskQueue, LONG_VIDEO_PROXY_DURATION_SEC, MEDIA_PROBE_OPERATION, MEDIA_PROXY_CACHE_CLEAR_OPERATION, MEDIA_PROXY_CACHE_STATS_OPERATION, MEDIA_PROXY_OPERATION, WorkerClient, mediaProbeToProjectAsset, type BrowserMediaSource, type MediaProbeProgress, type MediaProbeRequest, type MediaProbeResult, type MediaProxyProgress, type MediaProxyRequest, type MediaProxyResult, type PlaybackMediaSource, type ProxyCacheStats, type QueueStats, type ScheduledTaskHandle, type WorkerTransport, } from "@web-video-editor/media-runtime"; import { LogHub, StructuredLogger, type LogSink, } from "@web-video-editor/observability"; import type { PreviewSource } from "@web-video-editor/preview-runtime"; import { createDirectRuntimeSource, createOpfsRuntimeSource, } from "./preview-sources"; import type { BuildOpfsAction, BuildOpfsOptions, ImportAssetAction, ImportAssetOptions, WebVAAction, WebVAAssetEntry, WebVARuntimeSource, WebVAState, } from "./types"; const DEFAULT_PROXY_TIMEOUT_MIN_MS = 5 * 60_000; const DEFAULT_PROXY_TIMEOUT_MAX_MS = 30 * 60_000; const EMPTY_QUEUE_STATS: QueueStats = { active: 0, activePeak: 0, backpressureCount: 0, concurrency: 1, highWatermark: 2, queued: 0, queuedPeak: 0, }; export function createInitialWebVAState(): WebVAState { return { entries: [], importQueue: { ...EMPTY_QUEUE_STATS }, previewSources: {}, proxyQueue: { ...EMPTY_QUEUE_STATS }, }; } export type WebVASdkOptions = { createId?: () => string; createObjectUrl?: (blob: Blob) => string; defaultTimeoutMs?: number; importConcurrency?: number; importHighWatermark?: number; logHub?: LogHub; logger?: StructuredLogger; proxyConcurrency?: number; proxyHighWatermark?: number; revokeObjectUrl?: (url: string) => void; workerFactory?: () => WorkerTransport; workerLogSink?: LogSink; }; export type WebVAActionResult = | WebVAAssetEntry | WebVARuntimeSource | WebVAState | undefined; function errorMessage(error: unknown): string { return error instanceof Error ? error.message : String(error); } function defaultWorkerFactory(): WorkerTransport { if (typeof Worker === "undefined") { throw new Error( "Web Worker is unavailable; pass WebVASdkOptions.workerFactory", ); } return new Worker(new URL("./media.worker.ts", import.meta.url), { name: "web-va-media-worker", type: "module", }); } function proxyTimeoutMs(result: MediaProbeResult): number { return Math.min( DEFAULT_PROXY_TIMEOUT_MAX_MS, Math.max(DEFAULT_PROXY_TIMEOUT_MIN_MS, result.durationSec * 250), ); } export class WebVASdk { private readonly client: WorkerClient; private readonly createId: () => string; private readonly createObjectUrl: (blob: Blob) => string; private readonly importQueue: BoundedTaskQueue; private readonly importTasks = new Map< string, ScheduledTaskHandle >(); private readonly listeners = new Set<() => void>(); private readonly objectUrls = new Set(); private readonly proxyQueue: BoundedTaskQueue; private readonly proxyTasks = new Map< string, ScheduledTaskHandle >(); private readonly revokeObjectUrl: (url: string) => void; private snapshot = createInitialWebVAState(); private disposed = false; constructor(options: WebVASdkOptions = {}) { const logHub = options.logHub; const logger = options.logger ?? (logHub ? new StructuredLogger(logHub, "web-va-sdk") : undefined); this.createId = options.createId ?? (() => crypto.randomUUID()); this.createObjectUrl = options.createObjectUrl ?? ((blob) => URL.createObjectURL(blob)); this.revokeObjectUrl = options.revokeObjectUrl ?? ((url) => URL.revokeObjectURL(url)); this.client = new WorkerClient({ defaultTimeoutMs: options.defaultTimeoutMs ?? 5 * 60_000, logger, workerFactory: options.workerFactory ?? defaultWorkerFactory, workerLogSink: options.workerLogSink ?? logHub, }); this.importQueue = new BoundedTaskQueue(this.client, { concurrency: options.importConcurrency ?? 1, highWatermark: options.importHighWatermark ?? 2, logger, marker: "[IMPORT]", operation: MEDIA_PROBE_OPERATION, }); this.proxyQueue = new BoundedTaskQueue(this.client, { concurrency: options.proxyConcurrency ?? 1, highWatermark: options.proxyHighWatermark ?? 2, logger, marker: "[PROXY]", operation: MEDIA_PROXY_OPERATION, }); this.importQueue.events$.subscribe(() => { queueMicrotask(() => { if (!this.disposed) { this.publish({ importQueue: this.importQueue.stats() }); } }); }); this.proxyQueue.events$.subscribe(() => { queueMicrotask(() => { if (!this.disposed) { this.publish({ proxyQueue: this.proxyQueue.stats() }); } }); }); } getSnapshot = (): WebVAState => this.snapshot; subscribe = (listener: () => void): (() => void) => { this.assertActive(); this.listeners.add(listener); return () => this.listeners.delete(listener); }; getImportQueueStats(): QueueStats { return this.importQueue.stats(); } getProxyQueueStats(): QueueStats { return this.proxyQueue.stats(); } getPreviewSources(): readonly PreviewSource[] { return Object.values(this.snapshot.previewSources).map( (source) => source.preview, ); } getAudioSources(): ReadonlyMap { return new Map( Object.entries(this.snapshot.previewSources).map(([assetId, source]) => [ assetId, source.audio, ]), ); } getOriginalSources(): ReadonlyMap { const sources = new Map(); for (const entry of this.snapshot.entries) { if (entry.asset && !sources.has(entry.asset.id)) { sources.set(entry.asset.id, entry.source); } } return sources; } async dispatch(action: WebVAAction): Promise { switch (action.type) { case "asset.import": return this.importAsset(action.source, action.options); case "opfs.build": return this.buildOpfs(action.entryId, action.options); case "opfs.clear": return this.clearOpfsCache(); case "preview.add": return this.addToPreview(action.entryId); case "preview.remove": this.removeFromPreview(action.assetId); return undefined; } } async importAsset( source: BrowserMediaSource, options: ImportAssetOptions = {}, ): Promise { this.assertActive(); const entryId = options.entryId ?? this.createId(); if (this.findEntry(entryId)) { throw new Error(`Import entry "${entryId}" already exists`); } const entry: WebVAAssetEntry = { id: entryId, name: source.name, proxy: { ratio: 0, status: "idle" }, source, status: "probing", }; this.publish({ entries: [entry, ...this.snapshot.entries] }); let task: ScheduledTaskHandle; try { task = this.importQueue.schedule< MediaProbeRequest, MediaProbeResult, MediaProbeProgress >(0, { source }, { requestId: `probe_${entryId}` }); this.importTasks.set(entryId, task); } catch (error) { this.updateEntry(entryId, { error: errorMessage(error), status: "failed", }); throw error; } task.progress$.subscribe(({ payload }) => { if (payload && this.importTasks.get(entryId) === task) { this.updateEntry(entryId, { progress: payload }); } }); try { const result = await task.result; if (this.importTasks.get(entryId) !== task) { throw new DOMException("Import was cancelled", "AbortError"); } this.importTasks.delete(entryId); const asset = mediaProbeToProjectAsset(result, options.assetId); const direct = createDirectRuntimeSource( asset, result, source, this.createObjectUrl, ); if (direct.objectUrl) { this.objectUrls.add(direct.objectUrl); } this.updateEntry(entryId, { asset, result, runtimeSource: direct.runtimeSource, status: "ready", }); if (options.autoAddToPreview ?? true) { this.addToPreview(entryId); } if (options.autoBuildOpfs ?? true) { void this.buildOpfs(entryId, { parameters: options.proxyParameters, }).catch(() => undefined); } return this.requireEntry(entryId); } catch (error) { if (this.importTasks.get(entryId) === task) { this.importTasks.delete(entryId); this.updateEntry(entryId, { error: errorMessage(error), status: "failed", }); } throw error; } } async probe(source: BrowserMediaSource): Promise { this.assertActive(); return this.importQueue.schedule( 0, { source }, { requestId: `diagnostic-probe-${this.createId()}` }, ).result; } async buildOpfs( entryId: string, options: BuildOpfsOptions = {}, ): Promise { this.assertActive(); const entry = this.requireReadyEntry(entryId); const running = this.proxyTasks.get(entryId); if (running) { await running.result; return this.requireEntry(entryId); } let task: ScheduledTaskHandle; try { task = this.proxyQueue.schedule< MediaProxyRequest, MediaProxyResult, MediaProxyProgress >( 0, { fingerprint: entry.result.fingerprint, parameters: options.parameters ?? (entry.result.durationSec >= LONG_VIDEO_PROXY_DURATION_SEC ? { frameRate: 15, maxThumbnailCount: 120, thumbnailIntervalSec: 30, } : undefined), }, { requestId: `proxy_${entryId}`, timeoutMs: proxyTimeoutMs(entry.result), }, ); this.proxyTasks.set(entryId, task); this.updateEntry(entryId, { proxy: { ratio: 0, status: "queued" }, }); } catch (error) { this.updateEntry(entryId, { proxy: { error: errorMessage(error), ratio: 0, status: "failed", }, }); throw error; } task.progress$.subscribe(({ payload, progress }) => { if (!payload || this.proxyTasks.get(entryId) !== task) { return; } const current = this.requireEntry(entryId); this.updateEntry(entryId, { proxy: { progress: { ...current.proxy.progress, ...payload, elapsedMs: payload.elapsedMs, stage: payload.stage, }, ratio: progress.ratio, status: "running", }, }); }); try { const result = await task.result; if (this.proxyTasks.get(entryId) !== task) { throw new DOMException("OPFS build was cancelled", "AbortError"); } this.proxyTasks.delete(entryId); const current = this.requireReadyEntry(entryId); const runtimeSource = createOpfsRuntimeSource(current.asset, result); this.updateEntry(entryId, { proxy: { progress: current.proxy.progress, ratio: 1, result, status: "ready", }, runtimeSource, }); if (this.snapshot.previewSources[current.asset.id]) { this.publish({ cacheStats: result.cache, previewSources: { ...this.snapshot.previewSources, [current.asset.id]: runtimeSource, }, }); } else { this.publish({ cacheStats: result.cache }); } return this.requireEntry(entryId); } catch (error) { if (this.proxyTasks.get(entryId) === task) { this.proxyTasks.delete(entryId); const current = this.requireEntry(entryId); const cancelled = error instanceof DOMException && error.name === "AbortError"; this.updateEntry(entryId, { proxy: { error: errorMessage(error), progress: current.proxy.progress, ratio: current.proxy.ratio, status: cancelled ? "cancelled" : "failed", }, }); } throw error; } } addToPreview(entryId: string): WebVARuntimeSource { this.assertActive(); const entry = this.requireReadyEntry(entryId); if (!entry.runtimeSource) { throw new Error(`Import entry "${entryId}" has no runtime source`); } this.publish({ previewSources: { ...this.snapshot.previewSources, [entry.asset.id]: entry.runtimeSource, }, }); return entry.runtimeSource; } removeFromPreview(assetId: string): void { this.assertActive(); if (!this.snapshot.previewSources[assetId]) { return; } const previewSources = { ...this.snapshot.previewSources }; delete previewSources[assetId]; this.publish({ previewSources }); } cancelImport(entryId: string): void { this.importTasks.get(entryId)?.cancel("Import cancelled"); } cancelOpfs(entryId: string): void { this.proxyTasks.get(entryId)?.cancel("OPFS build cancelled"); } async getOpfsCacheStats(): Promise { this.assertActive(); const stats = await this.client.request< Record, ProxyCacheStats >(MEDIA_PROXY_CACHE_STATS_OPERATION, 0, {}).result; this.publish({ cacheStats: stats }); return stats; } async clearOpfsCache(): Promise { this.assertActive(); for (const task of this.importTasks.values()) { task.cancel("Import cancelled because the OPFS cache was cleared"); } this.importTasks.clear(); for (const task of this.proxyTasks.values()) { task.cancel("OPFS build cancelled because the cache was cleared"); } this.proxyTasks.clear(); this.revokeAllObjectUrls(); this.publish({ cacheStats: undefined, entries: [], previewSources: {}, }); const stats = await this.client.request< Record, ProxyCacheStats >(MEDIA_PROXY_CACHE_CLEAR_OPERATION, 0, {}).result; this.publish({ cacheStats: stats }); return this.snapshot; } dispose(): void { if (this.disposed) { return; } this.disposed = true; this.listeners.clear(); this.importTasks.clear(); this.proxyTasks.clear(); this.importQueue.dispose(); this.proxyQueue.dispose(); this.client.dispose(); this.revokeAllObjectUrls(); } private assertActive(): void { if (this.disposed) { throw new Error("WebVASdk is disposed"); } } private findEntry(entryId: string): WebVAAssetEntry | undefined { return this.snapshot.entries.find((entry) => entry.id === entryId); } private requireEntry(entryId: string): WebVAAssetEntry { const entry = this.findEntry(entryId); if (!entry) { throw new Error(`Import entry "${entryId}" does not exist`); } return entry; } private requireReadyEntry( entryId: string, ): WebVAAssetEntry & { asset: NonNullable; result: MediaProbeResult } { const entry = this.requireEntry(entryId); if (entry.status !== "ready" || !entry.asset || !entry.result) { throw new Error(`Import entry "${entryId}" is not ready`); } return entry as WebVAAssetEntry & { asset: NonNullable; result: MediaProbeResult; }; } private updateEntry( entryId: string, patch: Partial, ): void { if (this.disposed || !this.findEntry(entryId)) { return; } this.publish({ entries: this.snapshot.entries.map((entry) => entry.id === entryId ? { ...entry, ...patch } : entry, ), }); } private publish(patch: Partial): void { if (this.disposed) { return; } this.snapshot = { ...this.snapshot, ...patch }; for (const listener of this.listeners) { listener(); } } private revokeAllObjectUrls(): void { for (const url of this.objectUrls) { this.revokeObjectUrl(url); } this.objectUrls.clear(); } } export type { BuildOpfsAction, ImportAssetAction };