Skip to content
Merged
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
20 changes: 10 additions & 10 deletions platform/src/lib/youtube-worker-extraction.ts
Original file line number Diff line number Diff line change
Expand Up @@ -94,16 +94,16 @@ async function boundedBody(response: Response, signal: AbortSignal): Promise<Uin
export interface WorkerExtractionDependencies {
execute: (operation: WorkerYouTubeOperation, fetchImpl: typeof fetch) => Promise<WorkerYouTubeResult>;
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<T extends WorkerYouTubeOperation>(env: Env, operation: T, onDiagnostic?: ExtractionDiagnosticSink, signal?: AbortSignal): Promise<YouTubeOperationResult<T>> {
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));
Expand All @@ -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();
Expand All @@ -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;
Expand All @@ -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;
Expand Down Expand Up @@ -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 });
30 changes: 14 additions & 16 deletions platform/test/worker-extraction-live/run.mjs
Original file line number Diff line number Diff line change
Expand Up @@ -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']) {
Expand Down
2 changes: 0 additions & 2 deletions platform/test/worker-extraction-live/worker.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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();
Expand Down
81 changes: 42 additions & 39 deletions platform/test/youtube-worker-extraction.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -13,52 +13,54 @@ const env = (overrides: Record<string, string> = {}) => ({
}) as unknown as Env;

function harness(execute: WorkerExtractionDependencies['execute']) {
const directFetch = vi.fn<typeof fetch>(async () => Response.json({}));
const close = vi.fn(async () => {});
const proxyFetch = vi.fn<typeof fetch>(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<WorkerExtractionDependencies['execute']>().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<WorkerExtractionDependencies['execute']>().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 () => {
Expand All @@ -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 () => {
Expand All @@ -87,61 +89,61 @@ 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 () => {
vi.useFakeTimers();
const execute = vi.fn<WorkerExtractionDependencies['execute']>(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<WorkerExtractionDependencies['execute']>(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);
});

test('Retry-After waits are bounded by the total operation deadline', async () => {
vi.useFakeTimers();
const execute = vi.fn<WorkerExtractionDependencies['execute']>(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');
Expand All @@ -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 });
Expand All @@ -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();
});
Expand All @@ -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<WorkerExtractionDependencies['execute']>().mockRejectedValueOnce(unavailable()).mockResolvedValue(transcript);
const execute = vi.fn<WorkerExtractionDependencies['execute']>().mockResolvedValue(transcript);
const { run, close } = harness(execute);
close.mockImplementation(() => new Promise(() => {}));
const pending = run(env(), operation);
Expand Down
Loading
Loading