diff --git a/compute/worker/src/runtime.ts b/compute/worker/src/runtime.ts index 2a43774..61b18ca 100644 --- a/compute/worker/src/runtime.ts +++ b/compute/worker/src/runtime.ts @@ -69,6 +69,7 @@ const COMPUTE_STATE_BUCKET = 'compute_state'; const COMPUTE_STATE_TTL_MS = 24 * 60 * 60 * 1000; const LOOP_ERROR_BACKOFF_MS = 500; const RUNNING_HEARTBEAT_MS = 5000; +const OP_EVENTS_KEEPALIVE_MS = 15_000; const DOCUMENT_ID_REGEX = /^[a-f0-9]{64}$/i; const SAFE_NAMESPACE_REGEX = /^[a-zA-Z0-9._-]{1,128}$/; const WHISPER_MAX_DELIVER = 1; @@ -863,6 +864,7 @@ export async function createComputeWorkerApp(options: CreateComputeWorkerAppOpti let closed = false; let unsubscribe: (() => void) | null = null; + let keepalive: NodeJS.Timeout | null = null; const writeSnapshot = (snapshot: StreamedOperationState, eventId: number): void => { if (closed || reply.raw.writableEnded) return; @@ -884,6 +886,10 @@ export async function createComputeWorkerApp(options: CreateComputeWorkerAppOpti unsubscribe(); unsubscribe = null; } + if (keepalive) { + clearInterval(keepalive); + keepalive = null; + } activeSse = Math.max(0, activeSse - 1); markActivity(); if (!reply.raw.writableEnded) { @@ -903,6 +909,11 @@ export async function createComputeWorkerApp(options: CreateComputeWorkerAppOpti return reply; } + keepalive = setInterval(() => { + if (closed || reply.raw.writableEnded) return; + reply.raw.write(': keepalive\n\n'); + }, OP_EVENTS_KEEPALIVE_MS); + unsubscribe = await operationEventStream.subscribe({ opId: params.data.opId, sinceEventId, @@ -1247,6 +1258,17 @@ export async function createComputeWorkerApp(options: CreateComputeWorkerAppOpti const result = await input.run(decoded.payload, context.queueWaitTiming?.queueWaitMs ?? 0, { onProgress: async (progress) => { + try { + input.msg.working(); + } catch (ackError) { + app.log.warn({ + worker: input.workerLabel, + kind: context?.decoded.kind, + opId: context?.decoded.opId, + jobId: context?.decoded.jobId, + error: toErrorMessage(ackError), + }, 'failed to extend JetStream ack wait on progress'); + } await markProgress(context!, progress, Date.now()); }, }); diff --git a/src/app/api/documents/[id]/parsed/events/route.ts b/src/app/api/documents/[id]/parsed/events/route.ts index 31af7fe..2649ae5 100644 --- a/src/app/api/documents/[id]/parsed/events/route.ts +++ b/src/app/api/documents/[id]/parsed/events/route.ts @@ -5,7 +5,7 @@ import { documents } from '@/db/schema'; import { requireAuthContext } from '@/lib/server/auth/auth'; import { getWorkerClientConfigFromEnv } from '@/lib/server/compute/worker'; import { isAbortLikeError } from '@/lib/server/compute/abort-like-error'; -import { snapshotFromWorkerState } from '@/lib/server/compute/worker-parse-state'; +import { isWorkerOperationStateStale, snapshotFromWorkerState } from '@/lib/server/compute/worker-parse-state'; import { fetchWorkerOperationState } from '@/lib/server/compute/worker-op-state'; import { isValidDocumentId } from '@/lib/server/documents/blobstore'; import { normalizeParseStatus, parseDocumentParseState } from '@/lib/server/documents/parse-state'; @@ -16,6 +16,7 @@ import { errorResponse } from '@/lib/server/errors/next-response'; import { logDegraded, logServerError } from '@/lib/server/errors/logging'; import type { PdfParseProgress, PdfParseStatus } from '@/types/parsed-pdf'; import { parseSseEventId, parseSsePayload } from '@openreader/compute-core'; +import { getComputeOpStaleMs } from '@openreader/compute-core'; import type { PdfLayoutJobResult, WorkerOperationEvent, WorkerOperationState } from '@openreader/compute-core/api-contracts'; export const dynamic = 'force-dynamic'; @@ -55,6 +56,7 @@ function sleep(ms: number): Promise { } async function toSnapshotState(row: ParseRow, preferredOpId?: string | null): Promise { + const opStaleMs = getComputeOpStaleMs(); const state = await healStaleDocumentParseState({ documentId: row.id, userId: row.userId, @@ -68,7 +70,11 @@ async function toSnapshotState(row: ParseRow, preferredOpId?: string | null): Pr // the per-user document row currently says "ready" or has a different opId. if (opId && (requestedOpId !== null || parseStatus !== 'ready')) { const workerState = await fetchWorkerOperationState(opId); - if (workerState && workerState.opId === opId) { + if ( + workerState + && workerState.opId === opId + && !isWorkerOperationStateStale(workerState, opStaleMs) + ) { return { snapshot: { ...snapshotFromWorkerState(workerState), @@ -325,119 +331,136 @@ export async function GET(req: NextRequest, ctx: { params: Promise<{ id: string continue; } - workerAbort = new AbortController(); - const query = lastEventId && lastEventId > 0 - ? `?sinceEventId=${encodeURIComponent(String(lastEventId))}` - : ''; - const response = await fetch( - `${workerCfg.baseUrl}/ops/${encodeURIComponent(currentOpId)}/events${query}`, - { - method: 'GET', - headers: { - Authorization: `Bearer ${workerCfg.token}`, - Accept: 'text/event-stream', - ...(lastEventId && lastEventId > 0 ? { 'Last-Event-ID': String(lastEventId) } : {}), + try { + workerAbort = new AbortController(); + const query = lastEventId && lastEventId > 0 + ? `?sinceEventId=${encodeURIComponent(String(lastEventId))}` + : ''; + const response = await fetch( + `${workerCfg.baseUrl}/ops/${encodeURIComponent(currentOpId)}/events${query}`, + { + method: 'GET', + headers: { + Authorization: `Bearer ${workerCfg.token}`, + Accept: 'text/event-stream', + ...(lastEventId && lastEventId > 0 ? { 'Last-Event-ID': String(lastEventId) } : {}), + }, + cache: 'no-store', + signal: workerAbort.signal, }, - cache: 'no-store', - signal: workerAbort.signal, - }, - ); + ); - if (closed) return; + if (closed) return; - if (!response.ok) { - const upstreamResponseBody = await response.text().catch(() => ''); - logger.warn({ - event: 'documents.parsed.events.worker_stream_request_failed', - degraded: true, - step: 'worker_stream_request', - documentId: id, - opId: currentOpId, - status: response.status, - upstreamResponseBody, - error: { - name: 'WorkerStreamRequestFailed', - message: `Worker stream request failed with status ${response.status}`, - }, - }, 'Worker stream request failed'); - await sleep(500); - continue; - } - if (!response.body) { - logger.warn({ - event: 'documents.parsed.events.worker_stream_missing_body', - degraded: true, - step: 'worker_stream_body', - documentId: id, - opId: currentOpId, - error: { - name: 'WorkerStreamMissingBody', - message: 'Worker stream response missing body', - }, - }, 'Worker stream response missing body'); - await sleep(500); - continue; - } - - const reader = response.body.getReader(); - const decoder = new TextDecoder(); - let buffer = ''; - let streamEnded = false; - - while (!closed && !streamEnded) { - const read = await reader.read(); - if (read.done) { - streamEnded = true; - break; + if (!response.ok) { + const upstreamResponseBody = await response.text().catch(() => ''); + logger.warn({ + event: 'documents.parsed.events.worker_stream_request_failed', + degraded: true, + step: 'worker_stream_request', + documentId: id, + opId: currentOpId, + status: response.status, + upstreamResponseBody, + error: { + name: 'WorkerStreamRequestFailed', + message: `Worker stream request failed with status ${response.status}`, + }, + }, 'Worker stream request failed'); + await sleep(500); + continue; + } + if (!response.body) { + logger.warn({ + event: 'documents.parsed.events.worker_stream_missing_body', + degraded: true, + step: 'worker_stream_body', + documentId: id, + opId: currentOpId, + error: { + name: 'WorkerStreamMissingBody', + message: 'Worker stream response missing body', + }, + }, 'Worker stream response missing body'); + await sleep(500); + continue; } - buffer += decoder.decode(read.value, { stream: true }); + const reader = response.body.getReader(); + const decoder = new TextDecoder(); + let buffer = ''; + let streamEnded = false; - while (true) { - const frameEnd = buffer.indexOf('\n\n'); - if (frameEnd < 0) break; - const frame = buffer.slice(0, frameEnd); - buffer = buffer.slice(frameEnd + 2); - - const eventId = parseSseEventId(frame); - if (eventId && eventId > 0) { - lastEventId = eventId; + while (!closed && !streamEnded) { + const read = await reader.read(); + if (read.done) { + streamEnded = true; + break; } - const payload = parseSsePayload(frame); - if (!payload) continue; + buffer += decoder.decode(read.value, { stream: true }); - let parsed: WorkerOperationEvent | WorkerOperationState; - try { - parsed = JSON.parse(payload) as WorkerOperationEvent | WorkerOperationState; - } catch { - continue; - } + while (true) { + const frameEnd = buffer.indexOf('\n\n'); + if (frameEnd < 0) break; + const frame = buffer.slice(0, frameEnd); + buffer = buffer.slice(frameEnd + 2); - const workerSnapshot: WorkerOperationState = ( - parsed && typeof parsed === 'object' && 'snapshot' in parsed - ? parsed.snapshot - : parsed as WorkerOperationState - ); - if (!workerSnapshot || workerSnapshot.opId !== currentOpId) continue; + const eventId = parseSseEventId(frame); + if (eventId && eventId > 0) { + lastEventId = eventId; + } - const nextSnapshot: ParsedSnapshot = { - ...snapshotFromWorkerState(workerSnapshot), - opId: workerSnapshot.opId, - }; - const nextSignature = JSON.stringify(nextSnapshot); - if (nextSignature !== signature) { - current = nextSnapshot; - signature = nextSignature; - currentFromWorker = true; - writeSnapshot(current); - } + const payload = parseSsePayload(frame); + if (!payload) continue; - if (shouldCloseForTerminalSnapshot(current, currentFromWorker)) { - closeStream(); - return; + let parsed: WorkerOperationEvent | WorkerOperationState; + try { + parsed = JSON.parse(payload) as WorkerOperationEvent | WorkerOperationState; + } catch { + continue; + } + + const workerSnapshot: WorkerOperationState = ( + parsed && typeof parsed === 'object' && 'snapshot' in parsed + ? parsed.snapshot + : parsed as WorkerOperationState + ); + if (!workerSnapshot || workerSnapshot.opId !== currentOpId) continue; + + const nextSnapshot: ParsedSnapshot = { + ...snapshotFromWorkerState(workerSnapshot), + opId: workerSnapshot.opId, + }; + const nextSignature = JSON.stringify(nextSnapshot); + if (nextSignature !== signature) { + current = nextSnapshot; + signature = nextSignature; + currentFromWorker = true; + writeSnapshot(current); + } + + if (shouldCloseForTerminalSnapshot(current, currentFromWorker)) { + closeStream(); + return; + } } } + } catch (error) { + if (closed || isAbortLikeError(error)) return; + logDegraded(logger, { + event: 'documents.parsed.events.worker_stream_read_failed', + msg: 'Worker stream read failed; reconnecting', + step: 'worker_stream_read', + context: { + documentId: id, + opId: currentOpId, + requestId, + }, + error, + }); + await sleep(500); + continue; } if (closed) return; diff --git a/src/app/api/documents/[id]/parsed/route.ts b/src/app/api/documents/[id]/parsed/route.ts index 9738b90..0a783f2 100644 --- a/src/app/api/documents/[id]/parsed/route.ts +++ b/src/app/api/documents/[id]/parsed/route.ts @@ -7,6 +7,7 @@ import { requireAuthContext } from '@/lib/server/auth/auth'; import { createOrReusePdfWorkerOperation } from '@/lib/server/compute/worker-op-create'; import { documentParseStateFromWorkerState, + isWorkerOperationStateStale, snapshotFromWorkerState, } from '@/lib/server/compute/worker-parse-state'; import { fetchWorkerOperationState } from '@/lib/server/compute/worker-op-state'; @@ -33,6 +34,7 @@ import { createRequestLogger, hashForLog, type ServerLogger } from '@/lib/server import { errorResponse } from '@/lib/server/errors/next-response'; import { logDegraded } from '@/lib/server/errors/logging'; import type { ParsedPdfDocument } from '@/types/parsed-pdf'; +import { getComputeOpStaleMs } from '@openreader/compute-core'; import type { PdfLayoutJobResult, WorkerOperationState } from '@openreader/compute-core/api-contracts'; export const dynamic = 'force-dynamic'; @@ -192,6 +194,7 @@ export async function GET(req: NextRequest, ctx: { params: Promise<{ id: string request: req, }); try { + const opStaleMs = getComputeOpStaleMs(); if (!isS3Configured()) return s3NotConfiguredResponse(); const authCtxOrRes = await requireAuthContext(req); @@ -255,7 +258,11 @@ export async function GET(req: NextRequest, ctx: { params: Promise<{ id: string if (effectiveOpId && effectiveStatus !== 'ready') { const workerState = await fetchWorkerOperationState(effectiveOpId); - if (workerState && workerState.opId === effectiveOpId) { + if ( + workerState + && workerState.opId === effectiveOpId + && !isWorkerOperationStateStale(workerState, opStaleMs) + ) { return finalizeFromWorkerState({ workerState, row, @@ -348,6 +355,7 @@ export async function POST(req: NextRequest, ctx: { params: Promise<{ id: string request: req, }); try { + const opStaleMs = getComputeOpStaleMs(); if (!isS3Configured()) return s3NotConfiguredResponse(); const authCtxOrRes = await requireAuthContext(req); @@ -388,7 +396,12 @@ export async function POST(req: NextRequest, ctx: { params: Promise<{ id: string const existingOpId = normalizeOpId(state.opId); if (existingOpId) { const existing = await fetchWorkerOperationState(existingOpId); - if (existing && (existing.status === 'queued' || existing.status === 'running') && !replace) { + if ( + existing + && !isWorkerOperationStateStale(existing, opStaleMs) + && (existing.status === 'queued' || existing.status === 'running') + && !replace + ) { const snapshot = snapshotFromWorkerState(existing); return NextResponse.json({ error: 'Parse operation already in progress', diff --git a/src/lib/server/compute/worker-parse-state.ts b/src/lib/server/compute/worker-parse-state.ts index 0e2b861..cce75a4 100644 --- a/src/lib/server/compute/worker-parse-state.ts +++ b/src/lib/server/compute/worker-parse-state.ts @@ -2,6 +2,10 @@ import type { PdfLayoutJobResult, WorkerOperationState } from '@openreader/compu import type { PdfParseProgress, PdfParseStatus } from '@/types/parsed-pdf'; import type { DocumentParseState } from '@/lib/server/documents/parse-state'; +function isInflightWorkerStatus(status: WorkerOperationState['status']): boolean { + return status === 'queued' || status === 'running'; +} + export function mapWorkerStatusToParseStatus(status: WorkerOperationState['status']): PdfParseStatus { switch (status) { case 'queued': @@ -44,6 +48,18 @@ export function documentParseStateFromWorkerState( }; } +export function isWorkerOperationStateStale( + state: WorkerOperationState, + staleMs: number, + nowMs = Date.now(), +): boolean { + if (!isInflightWorkerStatus(state.status)) return false; + if (!Number.isFinite(staleMs) || staleMs <= 0) return false; + const updatedAt = Number(state.updatedAt ?? 0); + if (!Number.isFinite(updatedAt) || updatedAt <= 0) return false; + return (nowMs - updatedAt) > staleMs; +} + export function mergeNonReadyParseSnapshot(input: { parseStatus: PdfParseStatus; parseProgress: PdfParseProgress | null; diff --git a/tests/unit/worker-parse-state.vitest.spec.ts b/tests/unit/worker-parse-state.vitest.spec.ts index 062cf87..dfdd37e 100644 --- a/tests/unit/worker-parse-state.vitest.spec.ts +++ b/tests/unit/worker-parse-state.vitest.spec.ts @@ -2,6 +2,7 @@ import { describe, expect, test } from 'vitest'; import type { PdfLayoutJobResult, WorkerOperationState } from '@openreader/compute-core/api-contracts'; import { documentParseStateFromWorkerState, + isWorkerOperationStateStale, snapshotFromWorkerState, } from '../../src/lib/server/compute/worker-parse-state'; @@ -88,4 +89,30 @@ describe('worker parse state mapping', () => { error: 'layout model crashed', }); }); + + test('treats old inflight worker states as stale', () => { + const workerState = makeWorkerState({ + status: 'running', + updatedAt: 1_000, + progress: { + totalPages: 500, + pagesParsed: 250, + currentPage: 251, + phase: 'infer', + }, + }); + + expect(isWorkerOperationStateStale(workerState, 5_000, 6_001)).toBe(true); + expect(isWorkerOperationStateStale(workerState, 5_000, 5_999)).toBe(false); + }); + + test('never treats terminal worker states as stale', () => { + const failedState = makeWorkerState({ + status: 'failed', + updatedAt: 1_000, + error: { code: 'PDF_PARSE_FAILED', message: 'crashed' }, + }); + + expect(isWorkerOperationStateStale(failedState, 5_000, 99_999)).toBe(false); + }); });