refactor(api): add staleness detection for inflight worker operation states

Integrate isWorkerOperationStateStale checks into document parse API endpoints to ensure inflight worker operation states are not reused if stale. Introduce helper for staleness detection and corresponding unit tests. Enhance SSE event streaming with keepalive intervals and improve progress acknowledgment error handling.
This commit is contained in:
Richard R 2026-06-03 01:32:43 -06:00
parent 0c4a71ad47
commit b139120e4e
5 changed files with 204 additions and 103 deletions

View file

@ -69,6 +69,7 @@ const COMPUTE_STATE_BUCKET = 'compute_state';
const COMPUTE_STATE_TTL_MS = 24 * 60 * 60 * 1000; const COMPUTE_STATE_TTL_MS = 24 * 60 * 60 * 1000;
const LOOP_ERROR_BACKOFF_MS = 500; const LOOP_ERROR_BACKOFF_MS = 500;
const RUNNING_HEARTBEAT_MS = 5000; const RUNNING_HEARTBEAT_MS = 5000;
const OP_EVENTS_KEEPALIVE_MS = 15_000;
const DOCUMENT_ID_REGEX = /^[a-f0-9]{64}$/i; const DOCUMENT_ID_REGEX = /^[a-f0-9]{64}$/i;
const SAFE_NAMESPACE_REGEX = /^[a-zA-Z0-9._-]{1,128}$/; const SAFE_NAMESPACE_REGEX = /^[a-zA-Z0-9._-]{1,128}$/;
const WHISPER_MAX_DELIVER = 1; const WHISPER_MAX_DELIVER = 1;
@ -863,6 +864,7 @@ export async function createComputeWorkerApp(options: CreateComputeWorkerAppOpti
let closed = false; let closed = false;
let unsubscribe: (() => void) | null = null; let unsubscribe: (() => void) | null = null;
let keepalive: NodeJS.Timeout | null = null;
const writeSnapshot = (snapshot: StreamedOperationState, eventId: number): void => { const writeSnapshot = (snapshot: StreamedOperationState, eventId: number): void => {
if (closed || reply.raw.writableEnded) return; if (closed || reply.raw.writableEnded) return;
@ -884,6 +886,10 @@ export async function createComputeWorkerApp(options: CreateComputeWorkerAppOpti
unsubscribe(); unsubscribe();
unsubscribe = null; unsubscribe = null;
} }
if (keepalive) {
clearInterval(keepalive);
keepalive = null;
}
activeSse = Math.max(0, activeSse - 1); activeSse = Math.max(0, activeSse - 1);
markActivity(); markActivity();
if (!reply.raw.writableEnded) { if (!reply.raw.writableEnded) {
@ -903,6 +909,11 @@ export async function createComputeWorkerApp(options: CreateComputeWorkerAppOpti
return reply; return reply;
} }
keepalive = setInterval(() => {
if (closed || reply.raw.writableEnded) return;
reply.raw.write(': keepalive\n\n');
}, OP_EVENTS_KEEPALIVE_MS);
unsubscribe = await operationEventStream.subscribe({ unsubscribe = await operationEventStream.subscribe({
opId: params.data.opId, opId: params.data.opId,
sinceEventId, sinceEventId,
@ -1247,6 +1258,17 @@ export async function createComputeWorkerApp(options: CreateComputeWorkerAppOpti
const result = await input.run(decoded.payload, context.queueWaitTiming?.queueWaitMs ?? 0, { const result = await input.run(decoded.payload, context.queueWaitTiming?.queueWaitMs ?? 0, {
onProgress: async (progress) => { 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()); await markProgress(context!, progress, Date.now());
}, },
}); });

View file

@ -5,7 +5,7 @@ import { documents } from '@/db/schema';
import { requireAuthContext } from '@/lib/server/auth/auth'; import { requireAuthContext } from '@/lib/server/auth/auth';
import { getWorkerClientConfigFromEnv } from '@/lib/server/compute/worker'; import { getWorkerClientConfigFromEnv } from '@/lib/server/compute/worker';
import { isAbortLikeError } from '@/lib/server/compute/abort-like-error'; 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 { fetchWorkerOperationState } from '@/lib/server/compute/worker-op-state';
import { isValidDocumentId } from '@/lib/server/documents/blobstore'; import { isValidDocumentId } from '@/lib/server/documents/blobstore';
import { normalizeParseStatus, parseDocumentParseState } from '@/lib/server/documents/parse-state'; 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 { logDegraded, logServerError } from '@/lib/server/errors/logging';
import type { PdfParseProgress, PdfParseStatus } from '@/types/parsed-pdf'; import type { PdfParseProgress, PdfParseStatus } from '@/types/parsed-pdf';
import { parseSseEventId, parseSsePayload } from '@openreader/compute-core'; import { parseSseEventId, parseSsePayload } from '@openreader/compute-core';
import { getComputeOpStaleMs } from '@openreader/compute-core';
import type { PdfLayoutJobResult, WorkerOperationEvent, WorkerOperationState } from '@openreader/compute-core/api-contracts'; import type { PdfLayoutJobResult, WorkerOperationEvent, WorkerOperationState } from '@openreader/compute-core/api-contracts';
export const dynamic = 'force-dynamic'; export const dynamic = 'force-dynamic';
@ -55,6 +56,7 @@ function sleep(ms: number): Promise<void> {
} }
async function toSnapshotState(row: ParseRow, preferredOpId?: string | null): Promise<SnapshotState> { async function toSnapshotState(row: ParseRow, preferredOpId?: string | null): Promise<SnapshotState> {
const opStaleMs = getComputeOpStaleMs();
const state = await healStaleDocumentParseState({ const state = await healStaleDocumentParseState({
documentId: row.id, documentId: row.id,
userId: row.userId, 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. // the per-user document row currently says "ready" or has a different opId.
if (opId && (requestedOpId !== null || parseStatus !== 'ready')) { if (opId && (requestedOpId !== null || parseStatus !== 'ready')) {
const workerState = await fetchWorkerOperationState<PdfLayoutJobResult>(opId); const workerState = await fetchWorkerOperationState<PdfLayoutJobResult>(opId);
if (workerState && workerState.opId === opId) { if (
workerState
&& workerState.opId === opId
&& !isWorkerOperationStateStale(workerState, opStaleMs)
) {
return { return {
snapshot: { snapshot: {
...snapshotFromWorkerState(workerState), ...snapshotFromWorkerState(workerState),
@ -325,119 +331,136 @@ export async function GET(req: NextRequest, ctx: { params: Promise<{ id: string
continue; continue;
} }
workerAbort = new AbortController(); try {
const query = lastEventId && lastEventId > 0 workerAbort = new AbortController();
? `?sinceEventId=${encodeURIComponent(String(lastEventId))}` const query = lastEventId && lastEventId > 0
: ''; ? `?sinceEventId=${encodeURIComponent(String(lastEventId))}`
const response = await fetch( : '';
`${workerCfg.baseUrl}/ops/${encodeURIComponent(currentOpId)}/events${query}`, const response = await fetch(
{ `${workerCfg.baseUrl}/ops/${encodeURIComponent(currentOpId)}/events${query}`,
method: 'GET', {
headers: { method: 'GET',
Authorization: `Bearer ${workerCfg.token}`, headers: {
Accept: 'text/event-stream', Authorization: `Bearer ${workerCfg.token}`,
...(lastEventId && lastEventId > 0 ? { 'Last-Event-ID': String(lastEventId) } : {}), 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) { if (!response.ok) {
const upstreamResponseBody = await response.text().catch(() => ''); const upstreamResponseBody = await response.text().catch(() => '');
logger.warn({ logger.warn({
event: 'documents.parsed.events.worker_stream_request_failed', event: 'documents.parsed.events.worker_stream_request_failed',
degraded: true, degraded: true,
step: 'worker_stream_request', step: 'worker_stream_request',
documentId: id, documentId: id,
opId: currentOpId, opId: currentOpId,
status: response.status, status: response.status,
upstreamResponseBody, upstreamResponseBody,
error: { error: {
name: 'WorkerStreamRequestFailed', name: 'WorkerStreamRequestFailed',
message: `Worker stream request failed with status ${response.status}`, message: `Worker stream request failed with status ${response.status}`,
}, },
}, 'Worker stream request failed'); }, 'Worker stream request failed');
await sleep(500); await sleep(500);
continue; continue;
} }
if (!response.body) { if (!response.body) {
logger.warn({ logger.warn({
event: 'documents.parsed.events.worker_stream_missing_body', event: 'documents.parsed.events.worker_stream_missing_body',
degraded: true, degraded: true,
step: 'worker_stream_body', step: 'worker_stream_body',
documentId: id, documentId: id,
opId: currentOpId, opId: currentOpId,
error: { error: {
name: 'WorkerStreamMissingBody', name: 'WorkerStreamMissingBody',
message: 'Worker stream response missing body', message: 'Worker stream response missing body',
}, },
}, 'Worker stream response missing body'); }, 'Worker stream response missing body');
await sleep(500); await sleep(500);
continue; 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;
} }
buffer += decoder.decode(read.value, { stream: true }); const reader = response.body.getReader();
const decoder = new TextDecoder();
let buffer = '';
let streamEnded = false;
while (true) { while (!closed && !streamEnded) {
const frameEnd = buffer.indexOf('\n\n'); const read = await reader.read();
if (frameEnd < 0) break; if (read.done) {
const frame = buffer.slice(0, frameEnd); streamEnded = true;
buffer = buffer.slice(frameEnd + 2); break;
const eventId = parseSseEventId(frame);
if (eventId && eventId > 0) {
lastEventId = eventId;
} }
const payload = parseSsePayload(frame); buffer += decoder.decode(read.value, { stream: true });
if (!payload) continue;
let parsed: WorkerOperationEvent<PdfLayoutJobResult> | WorkerOperationState<PdfLayoutJobResult>; while (true) {
try { const frameEnd = buffer.indexOf('\n\n');
parsed = JSON.parse(payload) as WorkerOperationEvent<PdfLayoutJobResult> | WorkerOperationState<PdfLayoutJobResult>; if (frameEnd < 0) break;
} catch { const frame = buffer.slice(0, frameEnd);
continue; buffer = buffer.slice(frameEnd + 2);
}
const workerSnapshot: WorkerOperationState<PdfLayoutJobResult> = ( const eventId = parseSseEventId(frame);
parsed && typeof parsed === 'object' && 'snapshot' in parsed if (eventId && eventId > 0) {
? parsed.snapshot lastEventId = eventId;
: parsed as WorkerOperationState<PdfLayoutJobResult> }
);
if (!workerSnapshot || workerSnapshot.opId !== currentOpId) continue;
const nextSnapshot: ParsedSnapshot = { const payload = parseSsePayload(frame);
...snapshotFromWorkerState(workerSnapshot), if (!payload) continue;
opId: workerSnapshot.opId,
};
const nextSignature = JSON.stringify(nextSnapshot);
if (nextSignature !== signature) {
current = nextSnapshot;
signature = nextSignature;
currentFromWorker = true;
writeSnapshot(current);
}
if (shouldCloseForTerminalSnapshot(current, currentFromWorker)) { let parsed: WorkerOperationEvent<PdfLayoutJobResult> | WorkerOperationState<PdfLayoutJobResult>;
closeStream(); try {
return; parsed = JSON.parse(payload) as WorkerOperationEvent<PdfLayoutJobResult> | WorkerOperationState<PdfLayoutJobResult>;
} catch {
continue;
}
const workerSnapshot: WorkerOperationState<PdfLayoutJobResult> = (
parsed && typeof parsed === 'object' && 'snapshot' in parsed
? parsed.snapshot
: parsed as WorkerOperationState<PdfLayoutJobResult>
);
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; if (closed) return;

View file

@ -7,6 +7,7 @@ import { requireAuthContext } from '@/lib/server/auth/auth';
import { createOrReusePdfWorkerOperation } from '@/lib/server/compute/worker-op-create'; import { createOrReusePdfWorkerOperation } from '@/lib/server/compute/worker-op-create';
import { import {
documentParseStateFromWorkerState, documentParseStateFromWorkerState,
isWorkerOperationStateStale,
snapshotFromWorkerState, snapshotFromWorkerState,
} from '@/lib/server/compute/worker-parse-state'; } from '@/lib/server/compute/worker-parse-state';
import { fetchWorkerOperationState } from '@/lib/server/compute/worker-op-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 { errorResponse } from '@/lib/server/errors/next-response';
import { logDegraded } from '@/lib/server/errors/logging'; import { logDegraded } from '@/lib/server/errors/logging';
import type { ParsedPdfDocument } from '@/types/parsed-pdf'; import type { ParsedPdfDocument } from '@/types/parsed-pdf';
import { getComputeOpStaleMs } from '@openreader/compute-core';
import type { PdfLayoutJobResult, WorkerOperationState } from '@openreader/compute-core/api-contracts'; import type { PdfLayoutJobResult, WorkerOperationState } from '@openreader/compute-core/api-contracts';
export const dynamic = 'force-dynamic'; export const dynamic = 'force-dynamic';
@ -192,6 +194,7 @@ export async function GET(req: NextRequest, ctx: { params: Promise<{ id: string
request: req, request: req,
}); });
try { try {
const opStaleMs = getComputeOpStaleMs();
if (!isS3Configured()) return s3NotConfiguredResponse(); if (!isS3Configured()) return s3NotConfiguredResponse();
const authCtxOrRes = await requireAuthContext(req); const authCtxOrRes = await requireAuthContext(req);
@ -255,7 +258,11 @@ export async function GET(req: NextRequest, ctx: { params: Promise<{ id: string
if (effectiveOpId && effectiveStatus !== 'ready') { if (effectiveOpId && effectiveStatus !== 'ready') {
const workerState = await fetchWorkerOperationState<PdfLayoutJobResult>(effectiveOpId); const workerState = await fetchWorkerOperationState<PdfLayoutJobResult>(effectiveOpId);
if (workerState && workerState.opId === effectiveOpId) { if (
workerState
&& workerState.opId === effectiveOpId
&& !isWorkerOperationStateStale(workerState, opStaleMs)
) {
return finalizeFromWorkerState({ return finalizeFromWorkerState({
workerState, workerState,
row, row,
@ -348,6 +355,7 @@ export async function POST(req: NextRequest, ctx: { params: Promise<{ id: string
request: req, request: req,
}); });
try { try {
const opStaleMs = getComputeOpStaleMs();
if (!isS3Configured()) return s3NotConfiguredResponse(); if (!isS3Configured()) return s3NotConfiguredResponse();
const authCtxOrRes = await requireAuthContext(req); 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); const existingOpId = normalizeOpId(state.opId);
if (existingOpId) { if (existingOpId) {
const existing = await fetchWorkerOperationState<PdfLayoutJobResult>(existingOpId); const existing = await fetchWorkerOperationState<PdfLayoutJobResult>(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); const snapshot = snapshotFromWorkerState(existing);
return NextResponse.json({ return NextResponse.json({
error: 'Parse operation already in progress', error: 'Parse operation already in progress',

View file

@ -2,6 +2,10 @@ import type { PdfLayoutJobResult, WorkerOperationState } from '@openreader/compu
import type { PdfParseProgress, PdfParseStatus } from '@/types/parsed-pdf'; import type { PdfParseProgress, PdfParseStatus } from '@/types/parsed-pdf';
import type { DocumentParseState } from '@/lib/server/documents/parse-state'; 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 { export function mapWorkerStatusToParseStatus(status: WorkerOperationState['status']): PdfParseStatus {
switch (status) { switch (status) {
case 'queued': case 'queued':
@ -44,6 +48,18 @@ export function documentParseStateFromWorkerState(
}; };
} }
export function isWorkerOperationStateStale(
state: WorkerOperationState<PdfLayoutJobResult>,
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: { export function mergeNonReadyParseSnapshot(input: {
parseStatus: PdfParseStatus; parseStatus: PdfParseStatus;
parseProgress: PdfParseProgress | null; parseProgress: PdfParseProgress | null;

View file

@ -2,6 +2,7 @@ import { describe, expect, test } from 'vitest';
import type { PdfLayoutJobResult, WorkerOperationState } from '@openreader/compute-core/api-contracts'; import type { PdfLayoutJobResult, WorkerOperationState } from '@openreader/compute-core/api-contracts';
import { import {
documentParseStateFromWorkerState, documentParseStateFromWorkerState,
isWorkerOperationStateStale,
snapshotFromWorkerState, snapshotFromWorkerState,
} from '../../src/lib/server/compute/worker-parse-state'; } from '../../src/lib/server/compute/worker-parse-state';
@ -88,4 +89,30 @@ describe('worker parse state mapping', () => {
error: 'layout model crashed', 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);
});
}); });