diff --git a/platform/src/lib/youtube-worker-extraction.ts b/platform/src/lib/youtube-worker-extraction.ts index 3be90f6..92f990e 100644 --- a/platform/src/lib/youtube-worker-extraction.ts +++ b/platform/src/lib/youtube-worker-extraction.ts @@ -94,16 +94,16 @@ async function boundedBody(response: Response, signal: AbortSignal): Promise Promise; proxyTransport: (url: string) => YouTubeFetchTransport; - directFetch: typeof fetch; } -/** One direct attempt, then bounded whole-operation retries across the proxy pool. */ +/** Bounded whole-operation retries across the required proxy pool. */ export function createWorkerExtractionRunner(deps: WorkerExtractionDependencies) { return async function run(env: Env, operation: T, onDiagnostic?: ExtractionDiagnosticSink, signal?: AbortSignal): Promise> { const urls = workerProxyUrls(env); - const order = processorSlotOrder(Math.max(1, urls.length), randomProcessorSlot(Math.max(1, urls.length))); - const proxyAttempts = urls.length ? bounded(env.YOUTUBE_PROXY_MAX_ATTEMPTS, 4, 1, 4) : 0; - const routes: Array<{ egress: 'direct' | 'proxy'; slot: number; url?: string }> = [{ egress: 'direct', slot: 0 }]; + if (!urls.length) throw new YouTubeProcessorError('PROCESSOR_UNAVAILABLE', 'A YouTube proxy must be configured for Worker extraction.', 503); + const order = processorSlotOrder(urls.length, randomProcessorSlot(urls.length)); + const proxyAttempts = bounded(env.YOUTUBE_PROXY_MAX_ATTEMPTS, 4, 1, 4); + const routes: Array<{ egress: 'proxy'; slot: number; url: string }> = []; for (let i = 0; i < proxyAttempts; i++) { const slot = order[i % order.length]!; routes.push({ egress: 'proxy', slot, url: urls[slot]! }); } const total = new AbortController(); const totalTimer = setTimeout(() => total.abort(new DOMException('Extraction deadline', 'TimeoutError')), bounded(env.YOUTUBE_EXTRACTION_TIMEOUT_MS, 120_000, 1_000, 300_000)); @@ -115,7 +115,7 @@ export function createWorkerExtractionRunner(deps: WorkerExtractionDependencies) for (const [index, route] of routes.entries()) { deadline.throwIfAborted(); const attempt = new AbortController(); - const timeout = route.egress === 'direct' ? bounded(env.YOUTUBE_DIRECT_TIMEOUT_MS, 8_000, 100, 25_000) : bounded(env.YOUTUBE_PROXY_TIMEOUT_MS, 25_000, 100, 60_000); + const timeout = bounded(env.YOUTUBE_PROXY_TIMEOUT_MS, 25_000, 100, 60_000); const timer = setTimeout(() => attempt.abort(new DOMException('Attempt deadline', 'TimeoutError')), timeout); const attemptSignal = AbortSignal.any([deadline, attempt.signal]); const started = Date.now(); @@ -130,8 +130,8 @@ export function createWorkerExtractionRunner(deps: WorkerExtractionDependencies) let droppedEvents = 0; const record = (event: ExtractionAttempt['events'][number]) => { if (events.length < 64) events.push(event); else droppedEvents++; }; try { - transport = route.url ? deps.proxyTransport(route.url) : undefined; - const fetchImpl = transport?.fetch ?? deps.directFetch; + transport = deps.proxyTransport(route.url); + const fetchImpl = transport.fetch; const trackedFetch: typeof fetch = async (input, init = {}) => { const requestSignal = init.signal ?? (input instanceof Request ? input.signal : undefined); const activeSignal = requestSignal ? AbortSignal.any([attemptSignal, requestSignal]) : attemptSignal; @@ -156,7 +156,7 @@ export function createWorkerExtractionRunner(deps: WorkerExtractionDependencies) attemptSignal.throwIfAborted(); if (operation.kind === 'video' && isVideoMetadataBotChallenge(value)) throw new YouTubeProcessorError('UNAVAILABLE', 'YouTube blocked this connection.', 503, true); // Partial catalogs probe every distinct route once, without repeated pool passes. - if (shouldFallbackResult(operation, value) && index + 1 < Math.min(routes.length, urls.length + 1)) { + if (shouldFallbackResult(operation, value) && index + 1 < Math.min(routes.length, urls.length)) { outcome = 'fallback'; retry = true; } else { outcome = 'success'; status = 200; @@ -199,4 +199,4 @@ export function createWorkerExtractionRunner(deps: WorkerExtractionDependencies) }; } -export const runWorkerYouTubeOperation = createWorkerExtractionRunner({ execute: executeWorkerYouTubeOperation, proxyTransport: createWorkerProxyTransport, directFetch: (input, init) => fetch(input, init) }); +export const runWorkerYouTubeOperation = createWorkerExtractionRunner({ execute: executeWorkerYouTubeOperation, proxyTransport: createWorkerProxyTransport }); diff --git a/platform/test/worker-extraction-live/run.mjs b/platform/test/worker-extraction-live/run.mjs index 31a42f8..dc49718 100644 --- a/platform/test/worker-extraction-live/run.mjs +++ b/platform/test/worker-extraction-live/run.mjs @@ -5,22 +5,20 @@ const token = process.env.TEST_TOKEN; if (!base || !token) throw new Error('Supply the temporary Worker URL and TEST_TOKEN'); const cases = ['transcript', 'translated', 'words', 'transcriptSecond', 'transcriptThird', 'tracks', 'video', 'signals', 'search', 'browse', 'channel', 'channelVideos', 'channelPlaylists', 'playlist', 'comments', 'allComments', 'endscreen']; let failures = 0; -for (const proxy of [false, true]) { - for (const test of cases) { - try { - const response = await fetch(`${base}/${test}${proxy ? '?proxy=1' : ''}`, { headers: { authorization: `Bearer ${token}` }, signal: AbortSignal.timeout(80000) }); - const result = await response.json(); - const transcript = ['transcript', 'translated', 'words', 'transcriptSecond', 'transcriptThird'].includes(test); - const passed = response.ok && result.ok && (!proxy || result.proxyConnections > 0) - && (!transcript || result.summary?.textLength > 0 && result.summary?.segmentsLength > 0) - && (test !== 'translated' || result.summary?.translatedTo?.languageCode === 'fr') - && (test !== 'allComments' || result.summary?.pagesFetched === 2); - if (!passed) failures++; - console.log(JSON.stringify({ test, proxy, passed, ...result })); - } catch (error) { - failures++; - console.log(JSON.stringify({ test, proxy, passed: false, errorType: error.name })); - } +for (const test of cases) { + try { + const response = await fetch(`${base}/${test}`, { headers: { authorization: `Bearer ${token}` }, signal: AbortSignal.timeout(80000) }); + const result = await response.json(); + const transcript = ['transcript', 'translated', 'words', 'transcriptSecond', 'transcriptThird'].includes(test); + const passed = response.ok && result.ok && (result.proxyConnections > 0 && result.attempts?.every(attempt => attempt.egress === 'proxy')) + && (!transcript || result.summary?.textLength > 0 && result.summary?.segmentsLength > 0) + && (test !== 'translated' || result.summary?.translatedTo?.languageCode === 'fr') + && (test !== 'allComments' || result.summary?.pagesFetched === 2); + if (!passed) failures++; + console.log(JSON.stringify({ test, passed, ...result })); + } catch (error) { + failures++; + console.log(JSON.stringify({ test, passed: false, errorType: error.name })); } } for (const test of ['cert-expired', 'cert-hostname']) { diff --git a/platform/test/worker-extraction-live/worker.ts b/platform/test/worker-extraction-live/worker.ts index 4a028f4..e07d2b9 100644 --- a/platform/test/worker-extraction-live/worker.ts +++ b/platform/test/worker-extraction-live/worker.ts @@ -43,11 +43,9 @@ export default { } const operation = cases[url.pathname.slice(1)]; if (!operation) return new Response('Unknown test', { status: 404 }); - const forced = url.searchParams.has('proxy'); const attempts: ExtractionAttempt[] = []; let proxyConnections = 0; const run = createWorkerExtractionRunner({ execute: executeWorkerYouTubeOperation, - directFetch: forced ? async () => new Response('', { status: 429 }) : fetch, proxyTransport: proxy => { proxyConnections++; return createWorkerProxyTransport(proxy); }, }); const started = Date.now(); diff --git a/platform/test/youtube-worker-extraction.test.ts b/platform/test/youtube-worker-extraction.test.ts index b10dabb..7f25976 100644 --- a/platform/test/youtube-worker-extraction.test.ts +++ b/platform/test/youtube-worker-extraction.test.ts @@ -13,52 +13,54 @@ const env = (overrides: Record = {}) => ({ }) as unknown as Env; function harness(execute: WorkerExtractionDependencies['execute']) { - const directFetch = vi.fn(async () => Response.json({})); const close = vi.fn(async () => {}); const proxyFetch = vi.fn(async () => Response.json({})); const proxyTransport = vi.fn((_url: string) => ({ fetch: proxyFetch, close })); - const run = createWorkerExtractionRunner({ execute, directFetch, proxyTransport }); - return { run, directFetch, proxyFetch, proxyTransport, close }; + const run = createWorkerExtractionRunner({ execute, proxyTransport }); + return { run, proxyFetch, proxyTransport, close }; } beforeEach(() => { vi.spyOn(console, 'info').mockImplementation(() => {}); }); afterEach(() => { vi.restoreAllMocks(); vi.useRealTimers(); }); -test('direct success never creates a proxy and returns the same result with safe diagnostics', async () => { +test('first attempt uses a proxy and returns the same result with safe diagnostics', async () => { const execute = vi.fn(async () => transcript); const { run, proxyTransport } = harness(execute); const diagnostic = vi.fn(); expect(await run(env(), operation, diagnostic)).toBe(transcript); - expect(proxyTransport).not.toHaveBeenCalled(); - expect(diagnostic).toHaveBeenCalledWith(expect.objectContaining({ backend: 'worker', egress: 'direct', attempt: 1, outcome: 'success' })); + expect(proxyTransport).toHaveBeenCalledTimes(1); + expect(diagnostic).toHaveBeenCalledWith(expect.objectContaining({ backend: 'worker', egress: 'proxy', attempt: 1, outcome: 'success' })); }); -test('whole-operation fallback visits each proxy before repeating, including a fifth attempt', async () => { - const execute = vi.fn().mockRejectedValueOnce(unavailable()).mockRejectedValueOnce(unavailable()).mockRejectedValueOnce(unavailable()).mockRejectedValueOnce(unavailable()).mockResolvedValue(transcript); +test('whole-operation fallback visits each proxy before repeating, including a fourth attempt', async () => { + const execute = vi.fn().mockRejectedValueOnce(unavailable()).mockRejectedValueOnce(unavailable()).mockRejectedValueOnce(unavailable()).mockResolvedValue(transcript); const { run, proxyTransport, close } = harness(execute); const diagnostics: ExtractionAttempt[] = []; await expect(run(env(), operation, event => diagnostics.push(event))).resolves.toBe(transcript); - expect(execute).toHaveBeenCalledTimes(5); + expect(execute).toHaveBeenCalledTimes(4); for (const [op] of execute.mock.calls) expect(op).toBe(operation); const urls = proxyTransport.mock.calls.map(call => call[0]); expect(urls[0]).not.toEqual(urls[1]); expect(urls[2]).toEqual(urls[0]); expect(urls[3]).toEqual(urls[1]); expect(close).toHaveBeenCalledTimes(4); - expect(diagnostics).toHaveLength(5); - expect(diagnostics.at(-1)).toMatchObject({ attempt: 5, outcome: 'success', egress: 'proxy' }); + expect(diagnostics).toHaveLength(4); + expect(diagnostics.at(-1)).toMatchObject({ attempt: 4, outcome: 'success', egress: 'proxy' }); }); test.each(['INVALID_INPUT', 'AUTH_REQUIRED'] as const)('%s is terminal even when marked retryable', async code => { const { run, proxyTransport } = harness(async () => { throw new YouTubeProcessorError(code, 'private secret', 400, true); }); await expect(run(env(), operation)).rejects.toMatchObject({ code }); - expect(proxyTransport).not.toHaveBeenCalled(); + expect(proxyTransport).toHaveBeenCalledTimes(1); }); -test('no configured proxy means exactly one attempt', async () => { - const execute = vi.fn(async () => { throw unavailable(); }); - const { run } = harness(execute); - await expect(run(env({ OUTBOUND_PROXY_URLS: '' }), operation)).rejects.toMatchObject({ code: 'UNAVAILABLE' }); - expect(execute).toHaveBeenCalledTimes(1); +test('missing proxy configuration fails before extraction or network access', async () => { + const execute = vi.fn(async () => transcript); + const { run, proxyTransport } = harness(execute); + const native = vi.spyOn(globalThis, 'fetch'); + await expect(run(env({ OUTBOUND_PROXY_URLS: '' }), operation)).rejects.toMatchObject({ code: 'PROCESSOR_UNAVAILABLE' }); + expect(execute).not.toHaveBeenCalled(); + expect(proxyTransport).not.toHaveBeenCalled(); + expect(native).not.toHaveBeenCalled(); }); test('an earlier upstream error is not overwritten by a later transcript NOT_FOUND', async () => { @@ -74,7 +76,7 @@ test('partial track catalogs visit each distinct route only once', async () => { const execute = vi.fn(async () => partial); const { run } = harness(execute); await expect(run(env(), { kind: 'caption-tracks', id: operation.id })).resolves.toBe(partial); - expect(execute).toHaveBeenCalledTimes(3); + expect(execute).toHaveBeenCalledTimes(2); }); test('bot-challenged metadata falls back but legitimate restricted metadata does not', async () => { @@ -87,7 +89,7 @@ test('bot-challenged metadata falls back but legitimate restricted metadata does const restricted = { availability: { status: 'LOGIN_REQUIRED', reason: 'This video is private' } } as WorkerYouTubeResult; const restrictedHarness = harness(async () => restricted); await expect(restrictedHarness.run(env(), { kind: 'video', id: operation.id })).resolves.toBe(restricted); - expect(restrictedHarness.proxyTransport).not.toHaveBeenCalled(); + expect(restrictedHarness.proxyTransport).toHaveBeenCalledTimes(1); }); test('timeouts cover body reads and use a fresh proxy attempt', async () => { @@ -95,20 +97,20 @@ test('timeouts cover body reads and use a fresh proxy attempt', async () => { const execute = vi.fn(async (_op, fetchImpl) => { await fetchImpl('https://www.youtube.com/api/timedtext'); return transcript; }); - const { run, directFetch, proxyTransport } = harness(execute); - directFetch.mockImplementationOnce(async () => new Response(new ReadableStream({ start(controller) { controller.enqueue(new Uint8Array([1])); } }))); - const result = run(env({ YOUTUBE_DIRECT_TIMEOUT_MS: '100' }), operation); + const { run, proxyFetch, proxyTransport } = harness(execute); + proxyFetch.mockImplementationOnce(async () => new Response(new ReadableStream({ start(controller) { controller.enqueue(new Uint8Array([1])); } }))); + const result = run(env({ YOUTUBE_PROXY_TIMEOUT_MS: '100' }), operation); await vi.advanceTimersByTimeAsync(105); await expect(result).resolves.toBe(transcript); - expect(proxyTransport).toHaveBeenCalledTimes(1); + expect(proxyTransport).toHaveBeenCalledTimes(2); }); -test('ignoring fetch cancellation cannot hold the direct attempt forever', async () => { +test('ignoring fetch cancellation cannot hold a proxy attempt forever', async () => { vi.useFakeTimers(); const execute = vi.fn(async (_op, fetchImpl) => { await fetchImpl('https://www.youtube.com/api/timedtext'); return transcript; }); - const { run, directFetch } = harness(execute); - directFetch.mockImplementationOnce(() => new Promise(() => {})); - const pending = run(env({ YOUTUBE_DIRECT_TIMEOUT_MS: '100' }), operation); + const { run, proxyFetch } = harness(execute); + proxyFetch.mockImplementationOnce(() => new Promise(() => {})); + const pending = run(env({ YOUTUBE_PROXY_TIMEOUT_MS: '100' }), operation); await vi.advanceTimersByTimeAsync(105); await expect(pending).resolves.toBe(transcript); }); @@ -116,32 +118,32 @@ test('ignoring fetch cancellation cannot hold the direct attempt forever', async test('Retry-After waits are bounded by the total operation deadline', async () => { vi.useFakeTimers(); const execute = vi.fn(async (_op, fetchImpl) => { await fetchImpl('https://www.youtube.com/watch'); throw unavailable(); }); - const { run, directFetch, proxyTransport } = harness(execute); - directFetch.mockResolvedValue(new Response('', { status: 429, headers: { 'retry-after': '60' } })); + const { run, proxyFetch, proxyTransport } = harness(execute); + proxyFetch.mockResolvedValue(new Response('', { status: 429, headers: { 'retry-after': '60' } })); const pending = run(env({ YOUTUBE_EXTRACTION_TIMEOUT_MS: '1000' }), operation); const rejected = expect(pending).rejects.toMatchObject({ code: 'UNAVAILABLE' }); await vi.advanceTimersByTimeAsync(1001); await rejected; - expect(proxyTransport).not.toHaveBeenCalled(); + expect(proxyTransport).toHaveBeenCalledTimes(1); }); test('canceling an operation prevents later routes', async () => { const controller = new AbortController(); const { run, proxyTransport } = harness(async () => { controller.abort(); throw unavailable(); }); await expect(run(env(), operation, undefined, controller.signal)).rejects.toMatchObject({ code: 'UNAVAILABLE' }); - expect(proxyTransport).not.toHaveBeenCalled(); + expect(proxyTransport).toHaveBeenCalledTimes(1); }); test('oversized upstream bodies fail before reaching the parser', async () => { - const { run, directFetch } = harness(async (_op, fetchImpl) => { await fetchImpl('https://www.youtube.com/watch'); return transcript; }); - directFetch.mockImplementation(async () => new Response(new Uint8Array(8 * 1024 * 1024 + 1))); - await expect(run(env({ OUTBOUND_PROXY_URLS: '' }), operation)).rejects.toMatchObject({ code: 'INVALID_RESPONSE' }); + const { run, proxyFetch } = harness(async (_op, fetchImpl) => { await fetchImpl('https://www.youtube.com/watch'); return transcript; }); + proxyFetch.mockImplementation(async () => new Response(new Uint8Array(8 * 1024 * 1024 + 1))); + await expect(run(env({ YOUTUBE_PROXY_MAX_ATTEMPTS: '1' }), operation)).rejects.toMatchObject({ code: 'INVALID_RESPONSE' }); }); test('raw transport errors cannot disclose URLs, credentials or nested causes', async () => { const { run } = harness(async () => { throw new Error('http://user:secret@proxy-a.example signed?token=secret', { cause: new Error('secret') }); }); const sink = vi.fn(); - try { await run(env({ OUTBOUND_PROXY_URLS: '' }), operation, sink); throw new Error('expected failure'); } + try { await run(env({ YOUTUBE_PROXY_MAX_ATTEMPTS: '1' }), operation, sink); throw new Error('expected failure'); } catch (error) { expect(String(error)).not.toContain('secret'); } expect(JSON.stringify(sink.mock.calls)).not.toContain('secret'); expect(JSON.stringify(vi.mocked(console.info).mock.calls)).not.toContain('secret'); @@ -162,9 +164,9 @@ test('single URL fallback and malformed pool handling do not expose credentials' expect(() => workerProxyUrls({ OUTBOUND_PROXY_URLS: '["http://a.test", "http://a.test/"]' })).toThrow(); }); -test('actual library keeps source metadata and caption download on the same fallback transport', async () => { +test('actual library keeps source metadata and caption download on the same proxy transport', async () => { + const nativeFetch = vi.spyOn(globalThis, 'fetch'); const captionRequests: string[] = []; - const directFetch: typeof fetch = async () => new Response('', { status: 429 }); const proxyFetch: typeof fetch = async input => { const url = String(input); if (url.includes('/watch?')) return new Response('', { status: 404 }); @@ -175,8 +177,9 @@ test('actual library keeps source metadata and caption download on the same fall return Response.json({ events: [{ tStartMs: 0, dDurationMs: 1000, segs: [{ utf8: 'Bonjour' }] }] }); }; const close = vi.fn(async () => {}); - const run = createWorkerExtractionRunner({ execute: executeWorkerYouTubeOperation, directFetch, proxyTransport: () => ({ fetch: proxyFetch, close }) }); + const run = createWorkerExtractionRunner({ execute: executeWorkerYouTubeOperation, proxyTransport: () => ({ fetch: proxyFetch, close }) }); await expect(run(env(), operation)).resolves.toMatchObject({ text: 'Bonjour', translatedTo: { languageCode: 'fr' } }); + expect(nativeFetch).not.toHaveBeenCalled(); expect(captionRequests).toEqual(['https://captions.test/en?fmt=json3&tlang=fr']); expect(close).toHaveBeenCalledOnce(); }); @@ -191,7 +194,7 @@ test('Worker switch leaves storyboard on the container path', async () => { test('a stalled transport close cannot hold the operation indefinitely', async () => { vi.useFakeTimers(); - const execute = vi.fn().mockRejectedValueOnce(unavailable()).mockResolvedValue(transcript); + const execute = vi.fn().mockResolvedValue(transcript); const { run, close } = harness(execute); close.mockImplementation(() => new Promise(() => {})); const pending = run(env(), operation); diff --git a/platform/worker-configuration.d.ts b/platform/worker-configuration.d.ts index fbbcd58..c836771 100644 --- a/platform/worker-configuration.d.ts +++ b/platform/worker-configuration.d.ts @@ -1,5 +1,5 @@ /* eslint-disable */ -// Generated by Wrangler by running `wrangler types --config=platform/wrangler.jsonc platform/worker-configuration.d.ts` (hash: 14780a3f802759f16f37aa3dd856cd79) +// Generated by Wrangler by running `wrangler types` (hash: 76dbd247b3d900e5c45b5d4e59a00c15) // Runtime types generated with workerd@1.20260801.1 2026-08-08 nodejs_compat interface __BaseEnv_Env { YOUTUBE_CACHE: KVNamespace; @@ -35,7 +35,6 @@ interface __BaseEnv_Env { LANDING_DEMO_RATE_LIMIT_MODE: "enforced"; YOUTUBE_EXTRACTION_BACKEND: "worker"; YOUTUBE_EXTRACTION_TIMEOUT_MS: "120000"; - YOUTUBE_DIRECT_TIMEOUT_MS: "8000"; YOUTUBE_PROXY_TIMEOUT_MS: "25000"; YOUTUBE_PROXY_MAX_ATTEMPTS: "4"; YOUTUBE_EXTRACTION_RETRY_BASE_MS: "250"; @@ -75,7 +74,7 @@ type StringifyValues> = { [Binding in keyof EnvType]: EnvType[Binding] extends string ? EnvType[Binding] : string; }; declare namespace NodeJS { - interface ProcessEnv extends StringifyValues> {} + interface ProcessEnv extends StringifyValues> {} } // Begin runtime types diff --git a/platform/wrangler.jsonc b/platform/wrangler.jsonc index 61834e2..3a6e3f6 100644 --- a/platform/wrangler.jsonc +++ b/platform/wrangler.jsonc @@ -30,7 +30,6 @@ // Core extraction runs in Workers. Set to container to roll back. "YOUTUBE_EXTRACTION_BACKEND": "worker", "YOUTUBE_EXTRACTION_TIMEOUT_MS": "120000", - "YOUTUBE_DIRECT_TIMEOUT_MS": "8000", "YOUTUBE_PROXY_TIMEOUT_MS": "25000", "YOUTUBE_PROXY_MAX_ATTEMPTS": "4", "YOUTUBE_EXTRACTION_RETRY_BASE_MS": "250", diff --git a/reference/agents/platform-internals.md b/reference/agents/platform-internals.md index 501f06e..eea2cce 100644 --- a/reference/agents/platform-internals.md +++ b/reference/agents/platform-internals.md @@ -9,7 +9,7 @@ Read the root `README.md`, `docs/open-source/local-development.mdx`, and `refere ## Layer boundaries - `platform/` owns authentication, authorization, credit metering, cache policy, and the public HTTP contract. -- `platform/youtube-processor/` owns the rollback outbound provider backend and storyboard calls. By default, with `YOUTUBE_EXTRACTION_BACKEND=worker`, core provider operations run through `platform/src/lib/youtube-worker-extraction.ts` using repository library source, native direct fetch and a Worker proxy transport. Read `reference/engineering/WORKER_EXTRACTION.md` before changing that path. `platform/youtube-frames/` owns agent-only frame extraction and its media traffic. Frame extraction has no public data API route. Route media requests through these private containers. The Worker owns authentication, metering, and response validation. +- `platform/youtube-processor/` owns the rollback outbound provider backend and storyboard calls. By default, with `YOUTUBE_EXTRACTION_BACKEND=worker`, core provider operations run through `platform/src/lib/youtube-worker-extraction.ts` using repository library source, a required Worker proxy transport. Read `reference/engineering/WORKER_EXTRACTION.md` before changing that path. `platform/youtube-frames/` owns agent-only frame extraction and its media traffic. Frame extraction has no public data API route. Route media requests through these private containers. The Worker owns authentication, metering, and response validation. - For individual-frame extraction, read `reference/engineering/FRAME_EXTRACTION.md`. The new container bundles the shared watch implementation and pins the published library independently. - `packages/all-things-youtube/` is the extraction library. The processor installs an exact published npm version using its lockfile for general provider calls. Publish library changes before updating that dependency. The storyboard path is the narrow exception: `platform/youtube-processor/storyboard-extractor.mjs` is a committed bundle of the shared storyboard source, checked in CI and deployed with the processor. Regenerate it with `npm --prefix platform/youtube-processor run bundle`; do not edit the generated file. The processor-directory Docker context uses an allowlist that excludes credentials, tests, and local artifacts. - `packages/video2ctx-cli/` is the independently published hosted-service CLI. Keep authentication and transport behavior compatible with both `video2ctx-platform` branches, and verify the npm tarball before releasing it. diff --git a/reference/engineering/WORKER_EXTRACTION.md b/reference/engineering/WORKER_EXTRACTION.md index fca9291..82d2d27 100644 --- a/reference/engineering/WORKER_EXTRACTION.md +++ b/reference/engineering/WORKER_EXTRACTION.md @@ -1,6 +1,6 @@ # Worker YouTube extraction -The Worker can execute core YouTube operations using the shared extraction library. It tries YouTube directly, then retries eligible failures through the configured Decodo gateways. Decodo selects the exit IP; the Worker still runs the YouTube client and verifies YouTube's TLS certificate. +The Worker can execute core YouTube operations using the shared extraction library. Every attempt uses the configured proxy gateways, including the first attempt. Eligible failures retry through the proxy pool. Decodo selects the exit IP; the Worker still runs the YouTube client and verifies YouTube's TLS certificate. The rollout switch defaults to `worker`. `YOUTUBE_EXTRACTION_BACKEND=worker` moves search, browse, video metadata and signals, channels, playlists, comments, caption catalogs, transcripts and end screens into the Worker. Storyboards remain in the processor container, including its image conversion. Exact frames remain in the FFmpeg container. Caching, coalescing, authentication, billing and public result shapes stay at their existing boundaries. @@ -14,12 +14,9 @@ sequenceDiagram participant Decodo participant Media as Media containers Caller->>Worker: Core data request - Worker->>YouTube: Direct extraction on cache miss - alt Eligible extraction failure - Worker->>Decodo: CONNECT tunnel - Decodo->>YouTube: Forward encrypted connection - YouTube-->>Worker: Transcript or other data through tunnel - end + Worker->>Decodo: CONNECT tunnel on cache miss + Decodo->>YouTube: Forward encrypted connection + YouTube-->>Worker: Transcript or other data through tunnel Worker-->>Caller: Existing result shape opt Storyboard or exact frame request Worker->>Media: Existing image or FFmpeg operation @@ -33,33 +30,32 @@ sequenceDiagram The library adds the supported `all-things-youtube/client` export for client operations such as browse and video signals. That export will be available to external consumers in the next package release. No package publication is required for this platform PR. -Direct access uses native Worker `fetch`. Proxy access uses pinned `tunnelfetch@1.13.0` with `cloudflare:sockets`, HTTP CONNECT and its JavaScript TLS implementation. System-root certificate verification remains enabled. A native CONNECT plus `startTls()` prototype failed after CONNECT succeeded in the deployed runtime; the alternative transport completed transcript downloads. There is no certificate-verification bypass. +All core Worker extraction uses pinned `tunnelfetch@1.13.0` with `cloudflare:sockets`, HTTP CONNECT and its JavaScript TLS implementation. System-root certificate verification remains enabled. A native CONNECT plus `startTls()` prototype failed after CONNECT succeeded in the deployed runtime; the alternative transport completed transcript downloads. There is no certificate-verification bypass. Each operation attempt owns a fresh transport and closes it before the next route. No sockets are shared across Worker invocations. The library's request retry count is one, so the operation runner controls retry count and fresh metadata discovery. ## Configuration -Existing `video2ctx` secrets are reused. `OUTBOUND_PROXY_URLS` is a JSON array of one to four distinct HTTP(S) URLs and takes precedence over the legacy single `OUTBOUND_PROXY_URL`. Do not put credentials in Wrangler vars or commit `.dev.vars`. Invalid configuration fails without echoing the URL. Without a proxy setting, the runner performs only the direct attempt. +Existing `video2ctx` secrets are reused. `OUTBOUND_PROXY_URLS` is a JSON array of one to four distinct HTTP(S) URLs and takes precedence over the legacy single `OUTBOUND_PROXY_URL`. Do not put credentials in Wrangler vars or commit `.dev.vars`. Invalid configuration fails without echoing the URL. A proxy setting is required. Missing configuration fails with `PROCESSOR_UNAVAILABLE` before any YouTube request; there is no direct fallback. | Variable | Default | Meaning | | --- | --- | --- | | `YOUTUBE_EXTRACTION_BACKEND` | `worker` | Set to `container` to roll back core extraction | | `YOUTUBE_EXTRACTION_TIMEOUT_MS` | `120000` | Entire operation including retry waits | -| `YOUTUBE_DIRECT_TIMEOUT_MS` | `8000` | Direct attempt budget | | `YOUTUBE_PROXY_TIMEOUT_MS` | `25000` | Budget for each proxy attempt | -| `YOUTUBE_PROXY_MAX_ATTEMPTS` | `4` | Proxy attempts after direct access | +| `YOUTUBE_PROXY_MAX_ATTEMPTS` | `4` | Maximum proxy attempts, including the first | | `YOUTUBE_EXTRACTION_RETRY_BASE_MS` | `250` | Initial jittered exponential delay | The pool starts at a random slot and visits every configured slot before repeating. A single gateway can receive several attempts; a changed exit IP depends on the Decodo session configuration. The runner honors `Retry-After`, bounded by the total deadline. Invalid input and confirmed authorization restrictions are terminal. Transcript failures otherwise retain the existing fallback policy because missing-caption labels can result from blocked upstream requests. Partial empty caption catalogs probe distinct routes once. Bot-challenged video metadata triggers fallback; ordinary private-video metadata does not. Limits are 8 MiB per response and 32 MiB across an attempt. Timeouts cover response reads as well as connection setup. The proxy library has additional per-request timeouts, including a 25-second total; increasing the operation setting does not raise that transport ceiling. Cleanup is attempted even after cancellation, with at most one second spent waiting for it. -Safe attempt logs contain route, slot, outcome, duration, byte count and status, without proxy credentials or signed YouTube URLs. Transcript diagnostics add optional `backend` and `egress` fields and allow five attempts, preserving older records. An earlier upstream transcript failure is retained when a later route reports `NOT_FOUND`. +Safe attempt logs contain route, slot, outcome, duration, byte count and status, without proxy credentials or signed YouTube URLs. Transcript diagnostics add optional `backend` and `egress` fields and allow five attempts for historical records. New operations have at most four proxy attempts. An earlier upstream transcript failure is retained when a later route reports `NOT_FOUND`. ## Verification and rollout -`platform/test/youtube-worker-extraction.test.ts` covers fallback ordering, terminal failures, time budgets, cleanup, body limits, redaction, translation through the real library and retained storyboard routing. `platform/test/worker-extraction-live/worker.ts` provides token-protected, fixed test cases for an isolated deployment with no production resource bindings. It must receive its own `TEST_TOKEN` and a copied proxy secret, and be deleted after use. +`platform/test/youtube-worker-extraction.test.ts` covers proxy-only routing, missing configuration, retry ordering, terminal failures, time budgets, cleanup, body limits, redaction, translation through the real library and retained storyboard routing. `platform/test/worker-extraction-live/worker.ts` provides token-protected, fixed test cases for an isolated deployment with no production resource bindings. It must receive its own `TEST_TOKEN` and a copied proxy secret, and be deleted after use. Local checks include the full platform and container suites, library packed-package tests, auth integration, documentation generation and Worker startup profiling. See `WORKER_EXTRACTION_RESULTS.md` for the deployed checks. -Deploying the merged configuration enables Worker extraction by default. No additional toggle is required. Watch success rate, latency, proxy usage and CPU for uncached operations after deployment. Sustained load testing and independent review of the new TLS dependency remain follow-up work. Set the switch back to `container` and deploy to roll back core extraction; retain processor bindings and image configuration throughout this rollout. There is no automatic container fallback in Worker mode. +Deploying the merged configuration enables Worker extraction by default. No additional toggle is required. Existing production proxy secrets are reused; local Worker extraction also requires a proxy setting. Watch success rate, latency, proxy usage and CPU for uncached operations after deployment. Sustained load testing and independent review of the new TLS dependency remain follow-up work. Set the switch back to `container` and deploy to roll back core extraction; retain processor bindings and image configuration throughout this rollout. There is no automatic container fallback in Worker mode. diff --git a/reference/engineering/WORKER_EXTRACTION_RESULTS.md b/reference/engineering/WORKER_EXTRACTION_RESULTS.md index 51f5865..8e37676 100644 --- a/reference/engineering/WORKER_EXTRACTION_RESULTS.md +++ b/reference/engineering/WORKER_EXTRACTION_RESULTS.md @@ -1,5 +1,27 @@ # Worker extraction verification +## Proxy-only follow-up after PR #99 + +Tested on 2026-09-27 using temporary Worker `youtube-proxy-only-check-927`, with the production extraction runner, shared library, unchanged timeout settings and the single Decodo gateway from local configuration. The Worker had no production resource bindings. These are standalone Worker extraction timings, not full agent-tool durations or a comparison against the same requests through containers. + +All 17 extraction cases passed. Every case created a proxy transport; all recorded transcript attempts used proxy egress. Both invalid-certificate checks also passed. The remaining cases covered caption tracks, metadata, signals, search, browse, channels, playlists, comments and end screens. + +| Transcript case | Duration | Proxy attempts | +| --- | --- | --- | +| Segment transcript, first video | 2,196 ms | 1 | +| French translation | 8,052 ms | 3 | +| Word transcript | 2,422 ms | 1 | +| Segment transcript, second video | 2,998 ms | 1 | +| Segment transcript, third video | 2,033 ms | 1 | + +Other extraction cases took 898–3,086 ms. Translation required retries, so removing direct access does not eliminate upstream variability. This small live sample does not establish production latency percentiles or validate the production multi-gateway pool. + +Local verification: 817 platform unit tests passed, 11 skipped. The platform build, regenerated Worker types, generated API documentation check and diff whitespace check passed. Tests cover the required proxy configuration, proxy-only first attempt, pool rotation, timeouts, cleanup, translation without native fetch and retained storyboard routing. + +The temporary Worker and local copied secrets were deleted after verification. Production was not deployed by this test. The follow-up requires the existing proxy secrets after normal deployment and does not include PR #100's timeout changes. + +## Original direct-first rollout + Tested on 2026-09-27 in temporary Cloudflare Worker `youtube-extraction-check-da29c5e2`, with no production resource bindings. The test used the repository's extraction runner and shared All Things YouTube library, plus the configured Decodo gateway. No production Worker deployment was performed. ## Deployed checks