Skip to content
Closed
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
105 changes: 105 additions & 0 deletions docs/api-reference/openapi.json
Original file line number Diff line number Diff line change
Expand Up @@ -7959,6 +7959,41 @@
"ffmpeg_success"
]
},
"transportCode": {
"type": "string",
"enum": [
"TIMEOUT_CONNECT",
"TIMEOUT_HANDSHAKE",
"TIMEOUT_HEADERS",
"TIMEOUT_IDLE",
"TIMEOUT_TOTAL"
]
},
"timeoutPhase": {
"type": "string",
"enum": [
"connect",
"handshake",
"headers",
"idle",
"total",
"request",
"attempt",
"extraction"
]
},
"requestPhase": {
"type": "string",
"enum": [
"headers",
"body"
]
},
"requestElapsedMs": {
"type": "number",
"minimum": 0,
"maximum": 9007199254740991
},
"profile": {
"type": "string",
"enum": [
Expand Down Expand Up @@ -9335,6 +9370,41 @@
"ffmpeg_success"
]
},
"transportCode": {
"type": "string",
"enum": [
"TIMEOUT_CONNECT",
"TIMEOUT_HANDSHAKE",
"TIMEOUT_HEADERS",
"TIMEOUT_IDLE",
"TIMEOUT_TOTAL"
]
},
"timeoutPhase": {
"type": "string",
"enum": [
"connect",
"handshake",
"headers",
"idle",
"total",
"request",
"attempt",
"extraction"
]
},
"requestPhase": {
"type": "string",
"enum": [
"headers",
"body"
]
},
"requestElapsedMs": {
"type": "number",
"minimum": 0,
"maximum": 9007199254740991
},
"profile": {
"type": "string",
"enum": [
Expand Down Expand Up @@ -10164,6 +10234,41 @@
"ffmpeg_success"
]
},
"transportCode": {
"type": "string",
"enum": [
"TIMEOUT_CONNECT",
"TIMEOUT_HANDSHAKE",
"TIMEOUT_HEADERS",
"TIMEOUT_IDLE",
"TIMEOUT_TOTAL"
]
},
"timeoutPhase": {
"type": "string",
"enum": [
"connect",
"handshake",
"headers",
"idle",
"total",
"request",
"attempt",
"extraction"
]
},
"requestPhase": {
"type": "string",
"enum": [
"headers",
"body"
]
},
"requestElapsedMs": {
"type": "number",
"minimum": 0,
"maximum": 9007199254740991
},
"profile": {
"type": "string",
"enum": [
Expand Down
4 changes: 4 additions & 0 deletions platform/src/lib/extraction-diagnostics.ts
Original file line number Diff line number Diff line change
Expand Up @@ -6,6 +6,10 @@ const metric = z.number().nonnegative().max(Number.MAX_SAFE_INTEGER);
export const extractionEventSchema = z.object({
stage: z.enum(['player', 'player_response', 'download', 'complete', 'request', 'image_normalized',
'caption_metadata', 'caption_retry', 'media_candidates', 'media_http', 'media_transfer', 'media_retry', 'media_retry_skipped', 'ffmpeg', 'ffmpeg_success']),
transportCode: z.enum(['TIMEOUT_CONNECT', 'TIMEOUT_HANDSHAKE', 'TIMEOUT_HEADERS', 'TIMEOUT_IDLE', 'TIMEOUT_TOTAL']).optional(),
timeoutPhase: z.enum(['connect', 'handshake', 'headers', 'idle', 'total', 'request', 'attempt', 'extraction']).optional(),
requestPhase: z.enum(['headers', 'body']).optional(),
requestElapsedMs: metric.optional(),
profile: z.enum(['IOS', 'ANDROID_VR', 'MWEB', 'WEB', 'ios', 'android', 'android_vr', 'mweb', 'web']).optional(),
outcome: z.enum(['selected', 'skipped', 'error', 'success']).optional(),
playabilityStatus: z.enum(['OK', 'LOGIN_REQUIRED', 'UNPLAYABLE', 'ERROR', 'LIVE_STREAM_OFFLINE', 'CONTENT_CHECK_REQUIRED', 'AGE_CHECK_REQUIRED', 'UNKNOWN']).optional(),
Expand Down
2 changes: 1 addition & 1 deletion platform/src/lib/youtube-cache-coordinator.ts
Original file line number Diff line number Diff line change
Expand Up @@ -116,7 +116,7 @@ export class YouTubeCacheCoordinatorCore {

const diagnostics: ExtractionAttempt[] = [];
const onDiagnostic: ExtractionDiagnosticSink = event => {
if (['transcript','storyboard','frames'].includes(request.operation.kind) && diagnostics.length < 4) {
if (['transcript','storyboard','frames'].includes(request.operation.kind) && diagnostics.length < 5) {
emitExtractionDiagnostic(item => { diagnostics.push(item); }, event);
}
};
Expand Down
53 changes: 37 additions & 16 deletions platform/src/lib/youtube-worker-extraction.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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 = route.egress === 'direct' ? bounded(env.YOUTUBE_DIRECT_TIMEOUT_MS, 5_000, 100, 25_000) : bounded(env.YOUTUBE_PROXY_TIMEOUT_MS, 20_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 @@ -136,21 +136,39 @@ export function createWorkerExtractionRunner(deps: WorkerExtractionDependencies)
const requestSignal = init.signal ?? (input instanceof Request ? input.signal : undefined);
const activeSignal = requestSignal ? AbortSignal.any([attemptSignal, requestSignal]) : attemptSignal;
activeSignal.throwIfAborted();
const response = await abortable(activeSignal, () => fetchImpl(input, { ...init, signal: activeSignal }));
if (response.status === 429 || response.status >= 500) {
const raw = response.headers.get('retry-after');
const seconds = raw === null ? NaN : Number(raw);
const delay = Number.isFinite(seconds) ? seconds * 1000 : Date.parse(raw ?? '') - Date.now();
if (Number.isFinite(delay)) retryAfter = Math.max(retryAfter, delay, 0);
}
const bytes = await boundedBody(response, activeSignal);
bytesRead += bytes.length;
if (bytesRead > MAX_ATTEMPT_BYTES) throw new YouTubeProcessorError('INVALID_RESPONSE', 'YouTube extraction exceeded the byte budget.', 502, true);
const requestStarted = Date.now();
const path = new URL(input instanceof Request ? input.url : String(input)).pathname;
record({ stage: path === '/watch' || path.endsWith('/player') ? 'caption_metadata' : 'download', outcome: response.ok ? 'success' : 'error', status: response.status, elapsedMs: Date.now() - started });
const headers = new Headers(response.headers);
headers.delete('content-length'); headers.delete('content-encoding');
return new Response([204, 205, 304].includes(response.status) ? null : bytes, { status: response.status, statusText: response.statusText, headers });
const stage = path === '/watch' || path.endsWith('/player') ? 'caption_metadata' : 'download';
let requestPhase: 'headers' | 'body' = 'headers';
try {
const response = await abortable(activeSignal, () => fetchImpl(input, { ...init, signal: activeSignal }));
if (response.status === 429 || response.status >= 500) {
const raw = response.headers.get('retry-after');
const seconds = raw === null ? NaN : Number(raw);
const delay = Number.isFinite(seconds) ? seconds * 1000 : Date.parse(raw ?? '') - Date.now();
if (Number.isFinite(delay)) retryAfter = Math.max(retryAfter, delay, 0);
}
requestPhase = 'body';
const bytes = await boundedBody(response, activeSignal);
bytesRead += bytes.length;
if (bytesRead > MAX_ATTEMPT_BYTES) throw new YouTubeProcessorError('INVALID_RESPONSE', 'YouTube extraction exceeded the byte budget.', 502, true);
record({ stage, outcome: response.ok ? 'success' : 'error', status: response.status, elapsedMs: Date.now() - started });
const headers = new Headers(response.headers);
headers.delete('content-length'); headers.delete('content-encoding');
return new Response([204, 205, 304].includes(response.status) ? null : bytes, { status: response.status, statusText: response.statusText, headers });
} catch (error) {
// Capture before the library wraps transport errors as UPSTREAM_ERROR.
// Never copy error messages, detail objects, URLs, or nested causes.
const phases = { TIMEOUT_CONNECT: 'connect', TIMEOUT_HANDSHAKE: 'handshake', TIMEOUT_HEADERS: 'headers', TIMEOUT_IDLE: 'idle', TIMEOUT_TOTAL: 'total' } as const;
const transportCode = Object.keys(phases).find(code => code === (error as { code?: unknown } | null)?.code) as keyof typeof phases | undefined;
const timeoutPhase = deadline.aborted && deadline.reason?.name === 'TimeoutError' ? 'extraction'
: attempt.signal.aborted && attempt.signal.reason?.name === 'TimeoutError' ? 'attempt'
: requestSignal?.aborted && requestSignal.reason?.name === 'TimeoutError' ? 'request'
: transportCode ? phases[transportCode] : undefined;
record({ stage, outcome: 'error', transportCode, timeoutPhase, requestPhase,
requestElapsedMs: Date.now() - requestStarted, elapsedMs: Date.now() - started });
throw error;
}
};
const value = await abortable(attemptSignal, () => deps.execute(operation, trackedFetch));
attemptSignal.throwIfAborted();
Expand All @@ -171,7 +189,10 @@ export function createWorkerExtractionRunner(deps: WorkerExtractionDependencies)
failureKind = extractionFailureKind(error, attemptSignal);
retry = !deadline.aborted && index + 1 < routes.length && shouldFallbackError(operation, failure);
outcome = retry ? 'fallback' : 'failed';
record({ stage: 'request', outcome: 'error', code: SAFE_CODES.find(code => code === failure.code) ?? 'UNKNOWN', elapsedMs: Date.now() - started });
record({ stage: 'request', outcome: 'error', code: SAFE_CODES.find(code => code === failure.code) ?? 'UNKNOWN',
timeoutPhase: deadline.aborted && deadline.reason?.name === 'TimeoutError' ? 'extraction'
: attempt.signal.aborted && attempt.signal.reason?.name === 'TimeoutError' ? 'attempt' : undefined,
elapsedMs: Date.now() - started });
if (!retry) throw operation.kind === 'transcript' && failure.code === 'NOT_FOUND' && upstreamFailure ? upstreamFailure : failure;
} finally {
clearTimeout(timer);
Expand Down
4 changes: 3 additions & 1 deletion platform/src/lib/youtube-worker-runtime.ts
Original file line number Diff line number Diff line change
Expand Up @@ -11,7 +11,9 @@ export type WorkerYouTubeResult = YouTubeOperationResult<WorkerYouTubeOperation>
export async function executeWorkerYouTubeOperation(operation: WorkerYouTubeOperation, fetchImpl: typeof fetch): Promise<WorkerYouTubeResult> {
// Avoid multiplying library retries by operation retries. Fresh metadata is
// retrieved on every operation retry, including malformed-caption recovery.
const options = { fetch: fetchImpl, retry: { policy: { maxAttempts: 1 } } };
// Match the proxy transport ceiling instead of inheriting the library's
// 10-second request deadline. The runner still enforces the shorter direct budget.
const options = { fetch: fetchImpl, retry: { policy: { maxAttempts: 1, attemptTimeoutMs: 20_000 } } };
const client = createYouTubeClient(options);
switch (operation.kind) {
case 'search': return client.search(operation.query, operation.filters ?? {});
Expand Down
2 changes: 1 addition & 1 deletion platform/src/lib/youtube-worker-transport.ts
Original file line number Diff line number Diff line change
Expand Up @@ -15,7 +15,7 @@ export function createWorkerProxyTransport(proxy: string): YouTubeFetchTransport
trust: { mode: 'system' },
maxBodyBytes: 8 * 1024 * 1024,
maxRedirects: 3,
timeouts: { connectMs: 10_000, handshakeMs: 15_000, headersMs: 15_000, idleMs: 10_000, totalMs: 25_000 },
timeouts: { connectMs: 5_000, handshakeMs: 8_000, headersMs: 12_000, idleMs: 8_000, totalMs: 20_000 },
});
return {
fetch: (input, init) => client.fetch(input, init),
Expand Down
2 changes: 1 addition & 1 deletion platform/src/lib/youtube.ts
Original file line number Diff line number Diff line change
Expand Up @@ -144,7 +144,7 @@ async function cached<T extends VideoResourceOperation>(
);
}
if (Array.isArray(response.diagnostics)) {
for (const event of response.diagnostics.slice(0, 4)) emitExtractionDiagnostic(onDiagnostic, event);
for (const event of response.diagnostics.slice(0, 5)) emitExtractionDiagnostic(onDiagnostic, event);
}
if (!response.ok && response.error) {
if (response.error.apiStatus) throw new ApiError(response.error.apiStatus,response.error.code,response.error.message);
Expand Down
14 changes: 14 additions & 0 deletions platform/test/youtube-cache-coordinator.test.ts
Original file line number Diff line number Diff line change
@@ -1,3 +1,4 @@
import { extractionFixture } from './fixtures/extraction-diagnostic';
import { readYouTubeCacheEntry, YouTubeCacheCoordinatorCore } from '../src/lib/youtube-cache-coordinator';
import type { YouTubeOperation } from '../src/lib/youtube-processor-client';

Expand Down Expand Up @@ -110,3 +111,16 @@ describe('YouTube cache coordinator', () => {
expect(loader).not.toHaveBeenCalled();
});
});


test('forwards all five Worker attempts including precise timeout events', async () => {
const cache = { get: vi.fn(async () => null), put: vi.fn(async () => {}) };
const coordinator = new YouTubeCacheCoordinatorCore(environment(cache), async (_env, _op, onDiagnostic) => {
for (let attempt = 1; attempt <= 5; attempt++) onDiagnostic?.({ ...extractionFixture, attempt,
events: [{ stage: 'download', outcome: 'error', transportCode: 'TIMEOUT_IDLE', timeoutPhase: 'idle', requestPhase: 'body', requestElapsedMs: 8000 }] });
return { text: 'Recovered', segments: [] };
});
const response = await coordinator.getOrLoad({ ...request, operation: { kind: 'transcript', id: 'abcdefghijk', granularity: 'word' }, resourceType: 'transcript' });
expect(response.diagnostics).toHaveLength(5);
expect(response.diagnostics?.[4]).toMatchObject({ attempt: 5, events: [expect.objectContaining({ transportCode: 'TIMEOUT_IDLE', timeoutPhase: 'idle' })] });
});
Loading
Loading