refactor(worker,tts): unify PDF and Whisper retry config, add alignment timeouts

Update compute worker to use separate environment variable for PDF job attempts
(`COMPUTE_PDF_JOB_ATTEMPTS`) and set Whisper job max deliveries to 1, clarifying
and separating retry logic between job types. Adjust `.env.example` and deployment
docs to comment out defaults and document new/renamed variables for advanced tuning.

In TTS segment ensure route, introduce explicit timeouts for Whisper alignment
operations, with configurable durations for local and worker modes, improving
robustness and error handling for long-running alignments.
This commit is contained in:
Richard R 2026-05-21 12:37:21 -06:00
parent fe9104e8fc
commit a2c714c2a5
4 changed files with 87 additions and 34 deletions

View file

@ -2,9 +2,9 @@
# Platform note: # Platform note:
# - Local/manual: keep PORT=8081 # - Local/manual: keep PORT=8081
# - Railway/Render/Fly/etc: platform injects PORT # - Railway/Render/Fly/etc: platform injects PORT
COMPUTE_WORKER_HOST=0.0.0.0 # COMPUTE_WORKER_HOST=0.0.0.0
PORT=8081 # PORT=8081
COMPUTE_LOG_FORMAT=pretty # COMPUTE_LOG_FORMAT=pretty
# COMPUTE_LOG_LEVEL=info # COMPUTE_LOG_LEVEL=info
# App <-> worker auth # App <-> worker auth
@ -23,20 +23,20 @@ S3_BUCKET=openreader-documents
S3_REGION=us-east-1 S3_REGION=us-east-1
S3_ACCESS_KEY_ID=devkey S3_ACCESS_KEY_ID=devkey
S3_SECRET_ACCESS_KEY=devsecret S3_SECRET_ACCESS_KEY=devsecret
S3_PREFIX=openreader # S3_PREFIX=openreader
# Optional for non-AWS S3-compatible endpoints: # Optional for non-AWS S3-compatible endpoints:
S3_ENDPOINT=http://host.docker.internal:8333 S3_ENDPOINT=http://host.docker.internal:8333
S3_FORCE_PATH_STYLE=true S3_FORCE_PATH_STYLE=true
# Queue + execution tuning # Queue + execution tuning
COMPUTE_PREWARM_MODELS=true # COMPUTE_PREWARM_MODELS=true
COMPUTE_WHISPER_CONCURRENCY=1 # COMPUTE_WHISPER_CONCURRENCY=1
COMPUTE_PDF_CONCURRENCY=2 # COMPUTE_PDF_CONCURRENCY=2
COMPUTE_WHISPER_TIMEOUT_MS=30000 # COMPUTE_WHISPER_TIMEOUT_MS=30000
COMPUTE_PDF_TIMEOUT_MS=90000 # COMPUTE_PDF_TIMEOUT_MS=90000
COMPUTE_JOB_ATTEMPTS=2 # COMPUTE_PDF_JOB_ATTEMPTS=2
COMPUTE_JOBS_STREAM_MAX_BYTES=268435456 # COMPUTE_JOBS_STREAM_MAX_BYTES=268435456
COMPUTE_JOB_STATES_MAX_BYTES=67108864 # COMPUTE_JOB_STATES_MAX_BYTES=67108864
# Optional stale window for reusing in-flight opKey entries before forcing a new attempt # Optional stale window for reusing in-flight opKey entries before forcing a new attempt
# Default is max(30m, 4x max compute timeout); running jobs also refresh heartbeat every 5s # Default is max(30m, 4x max compute timeout); running jobs also refresh heartbeat every 5s
# COMPUTE_OP_STALE_MS=1800000 # COMPUTE_OP_STALE_MS=1800000

View file

@ -54,6 +54,7 @@ const SSE_POLL_INTERVAL_MS = 400;
const RUNNING_HEARTBEAT_MS = 5000; const RUNNING_HEARTBEAT_MS = 5000;
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;
interface QueuedJob<TPayload> { interface QueuedJob<TPayload> {
jobId: string; jobId: string;
@ -311,7 +312,7 @@ async function ensureJetStreamResources(
jsm: JetStreamManager, jsm: JetStreamManager,
whisperTimeoutMs: number, whisperTimeoutMs: number,
pdfTimeoutMs: number, pdfTimeoutMs: number,
attempts: number, pdfAttempts: number,
maxBytes: number, maxBytes: number,
): Promise<void> { ): Promise<void> {
const streamConfig = { const streamConfig = {
@ -332,7 +333,12 @@ async function ensureJetStreamResources(
}); });
} }
const ensureConsumer = async (name: string, subject: string, ackWaitMs: number): Promise<void> => { const ensureConsumer = async (
name: string,
subject: string,
ackWaitMs: number,
maxDeliver: number,
): Promise<void> => {
const config = { const config = {
durable_name: name, durable_name: name,
ack_policy: AckPolicy.Explicit, ack_policy: AckPolicy.Explicit,
@ -340,7 +346,7 @@ async function ensureJetStreamResources(
replay_policy: ReplayPolicy.Instant, replay_policy: ReplayPolicy.Instant,
filter_subject: subject, filter_subject: subject,
ack_wait: nanos(Math.max(ackWaitMs, 1_000)), ack_wait: nanos(Math.max(ackWaitMs, 1_000)),
max_deliver: attempts, max_deliver: maxDeliver,
}; };
try { try {
@ -350,14 +356,14 @@ async function ensureJetStreamResources(
await jsm.consumers.update(JOBS_STREAM_NAME, name, { await jsm.consumers.update(JOBS_STREAM_NAME, name, {
filter_subject: subject, filter_subject: subject,
ack_wait: nanos(Math.max(ackWaitMs, 1_000)), ack_wait: nanos(Math.max(ackWaitMs, 1_000)),
max_deliver: attempts, max_deliver: maxDeliver,
}); });
} }
}; };
await Promise.all([ await Promise.all([
ensureConsumer(WHISPER_CONSUMER_NAME, WHISPER_JOBS_SUBJECT, whisperTimeoutMs + 15_000), ensureConsumer(WHISPER_CONSUMER_NAME, WHISPER_JOBS_SUBJECT, whisperTimeoutMs + 15_000, WHISPER_MAX_DELIVER),
ensureConsumer(LAYOUT_CONSUMER_NAME, LAYOUT_JOBS_SUBJECT, pdfTimeoutMs + 15_000), ensureConsumer(LAYOUT_CONSUMER_NAME, LAYOUT_JOBS_SUBJECT, pdfTimeoutMs + 15_000, pdfAttempts),
]); ]);
} }
@ -371,7 +377,7 @@ async function main(): Promise<void> {
const pdfConcurrency = readIntEnv('COMPUTE_PDF_CONCURRENCY', 2); const pdfConcurrency = readIntEnv('COMPUTE_PDF_CONCURRENCY', 2);
const whisperTimeoutMs = readIntEnv('COMPUTE_WHISPER_TIMEOUT_MS', 30_000); const whisperTimeoutMs = readIntEnv('COMPUTE_WHISPER_TIMEOUT_MS', 30_000);
const pdfTimeoutMs = readIntEnv('COMPUTE_PDF_TIMEOUT_MS', 90_000); const pdfTimeoutMs = readIntEnv('COMPUTE_PDF_TIMEOUT_MS', 90_000);
const attempts = readIntEnv('COMPUTE_JOB_ATTEMPTS', 2); const pdfAttempts = readIntEnv('COMPUTE_PDF_JOB_ATTEMPTS', 2);
const prewarmModels = parseBoolEnv('COMPUTE_PREWARM_MODELS', true); const prewarmModels = parseBoolEnv('COMPUTE_PREWARM_MODELS', true);
const jobsStreamMaxBytes = readIntEnv('COMPUTE_JOBS_STREAM_MAX_BYTES', 256 * 1024 * 1024); const jobsStreamMaxBytes = readIntEnv('COMPUTE_JOBS_STREAM_MAX_BYTES', 256 * 1024 * 1024);
const jobStatesMaxBytes = readIntEnv('COMPUTE_JOB_STATES_MAX_BYTES', 64 * 1024 * 1024); const jobStatesMaxBytes = readIntEnv('COMPUTE_JOB_STATES_MAX_BYTES', 64 * 1024 * 1024);
@ -398,7 +404,7 @@ async function main(): Promise<void> {
const js: JetStreamClient = jetstream(nc); const js: JetStreamClient = jetstream(nc);
const jsm: JetStreamManager = await jetstreamManager(nc); const jsm: JetStreamManager = await jetstreamManager(nc);
await ensureJetStreamResources(jsm, whisperTimeoutMs, pdfTimeoutMs, attempts, jobsStreamMaxBytes); await ensureJetStreamResources(jsm, whisperTimeoutMs, pdfTimeoutMs, pdfAttempts, jobsStreamMaxBytes);
const kv = await new Kvm(js).create(COMPUTE_STATE_BUCKET, { const kv = await new Kvm(js).create(COMPUTE_STATE_BUCKET, {
history: 1, history: 1,
@ -962,7 +968,9 @@ async function main(): Promise<void> {
} catch (error) { } catch (error) {
const message = toErrorMessage(error); const message = toErrorMessage(error);
const deliveryCount = input.msg.info.deliveryCount; const deliveryCount = input.msg.info.deliveryCount;
const hasRetriesLeft = deliveryCount < attempts; const isWhisperAlign = decoded?.kind === 'whisper_align';
const maxAttempts = isWhisperAlign ? WHISPER_MAX_DELIVER : pdfAttempts;
const hasRetriesLeft = !isWhisperAlign && deliveryCount < maxAttempts;
if (decoded) { if (decoded) {
const now = Date.now(); const now = Date.now();
@ -1020,7 +1028,7 @@ async function main(): Promise<void> {
jobId: decoded?.jobId, jobId: decoded?.jobId,
error: message, error: message,
deliveryCount, deliveryCount,
maxAttempts: attempts, maxAttempts,
}, 'job failed, nacked for retry'); }, 'job failed, nacked for retry');
} else { } else {
input.msg.term(message); input.msg.term(message);
@ -1030,7 +1038,8 @@ async function main(): Promise<void> {
jobId: decoded?.jobId, jobId: decoded?.jobId,
error: message, error: message,
deliveryCount, deliveryCount,
maxAttempts: attempts, maxAttempts,
retrySuppressed: isWhisperAlign ? 'whisper_align' : undefined,
}, 'job failed, max attempts reached'); }, 'job failed, max attempts reached');
} }
} finally { } finally {

View file

@ -44,9 +44,18 @@ Common optional:
- `COMPUTE_WORKER_HOST=0.0.0.0` - `COMPUTE_WORKER_HOST=0.0.0.0`
- `PORT=8081` (local/manual; on Railway platform injects this) - `PORT=8081` (local/manual; on Railway platform injects this)
- `COMPUTE_LOG_FORMAT=pretty` (default) or `json` - `COMPUTE_LOG_FORMAT=pretty` (default) or `json`
Advanced tuning (usually leave unset unless you need overrides):
- `COMPUTE_PREWARM_MODELS=true` - `COMPUTE_PREWARM_MODELS=true`
- `COMPUTE_WHISPER_CONCURRENCY=1`
- `COMPUTE_PDF_CONCURRENCY=2`
- `COMPUTE_WHISPER_TIMEOUT_MS=30000`
- `COMPUTE_PDF_TIMEOUT_MS=90000`
- `COMPUTE_PDF_JOB_ATTEMPTS=2` (PDF layout retry attempts)
- `COMPUTE_JOBS_STREAM_MAX_BYTES=268435456` (256MB JetStream jobs stream cap) - `COMPUTE_JOBS_STREAM_MAX_BYTES=268435456` (256MB JetStream jobs stream cap)
- `COMPUTE_JOB_STATES_MAX_BYTES=67108864` (64MB JetStream KV bucket cap) - `COMPUTE_JOB_STATES_MAX_BYTES=67108864` (64MB JetStream KV bucket cap)
- `COMPUTE_OP_STALE_MS=1800000` (stale op replacement window)
## App server environment variables (worker mode) ## App server environment variables (worker mode)
@ -110,8 +119,15 @@ COMPUTE_WORKER_HOST=0.0.0.0
# PORT=8081 # PORT=8081
# Railway: rely on injected PORT # Railway: rely on injected PORT
COMPUTE_WORKER_TOKEN=<long-random-shared-token> COMPUTE_WORKER_TOKEN=<long-random-shared-token>
COMPUTE_JOBS_STREAM_MAX_BYTES=268435456 # Optional advanced tuning overrides (defaults shown):
COMPUTE_JOB_STATES_MAX_BYTES=67108864 # COMPUTE_PREWARM_MODELS=true
# COMPUTE_WHISPER_CONCURRENCY=1
# COMPUTE_PDF_CONCURRENCY=2
# COMPUTE_WHISPER_TIMEOUT_MS=30000
# COMPUTE_PDF_TIMEOUT_MS=90000
# COMPUTE_PDF_JOB_ATTEMPTS=2
# COMPUTE_JOBS_STREAM_MAX_BYTES=268435456
# COMPUTE_JOB_STATES_MAX_BYTES=67108864
NATS_URL=tls://connect.ngs.global:4222 NATS_URL=tls://connect.ngs.global:4222
NATS_CREDS="-----BEGIN NATS USER JWT----- NATS_CREDS="-----BEGIN NATS USER JWT-----

View file

@ -44,7 +44,9 @@ import type {
export const runtime = 'nodejs'; export const runtime = 'nodejs';
export const dynamic = 'force-dynamic'; export const dynamic = 'force-dynamic';
const GENERATING_STALE_MS = 45_000; const GENERATING_STALE_MS = 360_000;
const LOCAL_ALIGNMENT_TIMEOUT_MS = 30_000;
const WORKER_ALIGNMENT_TIMEOUT_MS = 60_000;
function attachDeviceIdCookie(response: NextResponse, deviceId: string | null, didCreate: boolean) { function attachDeviceIdCookie(response: NextResponse, deviceId: string | null, didCreate: boolean) {
if (didCreate && deviceId) { if (didCreate && deviceId) {
@ -129,6 +131,24 @@ function isAbortLikeError(error: unknown): boolean {
return /abort/i.test(message); return /abort/i.test(message);
} }
async function withTimeout<T>(promise: Promise<T>, timeoutMs: number, label: string): Promise<T> {
let timer: NodeJS.Timeout | null = null;
try {
return await Promise.race([
promise,
new Promise<T>((_, reject) => {
timer = setTimeout(() => reject(new Error(`${label} timed out after ${timeoutMs}ms`)), timeoutMs);
}),
]);
} finally {
if (timer) clearTimeout(timer);
}
}
function getAlignmentTimeoutMs(mode: 'local' | 'worker'): number {
return mode === 'worker' ? WORKER_ALIGNMENT_TIMEOUT_MS : LOCAL_ALIGNMENT_TIMEOUT_MS;
}
async function deleteEntryIfUnused(userId: string, segmentEntryId: string): Promise<void> { async function deleteEntryIfUnused(userId: string, segmentEntryId: string): Promise<void> {
const stillReferenced = await db const stillReferenced = await db
.select({ segmentId: ttsSegmentVariants.segmentId }) .select({ segmentId: ttsSegmentVariants.segmentId })
@ -386,10 +406,14 @@ export async function POST(request: NextRequest) {
try { try {
const alignStartedAt = Date.now(); const alignStartedAt = Date.now();
const computeBackend = await getComputeBackend(); const computeBackend = await getComputeBackend();
const aligned = (await computeBackend.alignWords({ const aligned = (await withTimeout(
audioObjectKey: existing.audioKey, computeBackend.alignWords({
text: segment.text, audioObjectKey: existing.audioKey,
})).alignments; text: segment.text,
}),
getAlignmentTimeoutMs(computeBackend.mode),
`Whisper alignment (${computeBackend.mode})`,
)).alignments;
stageTimings.selfHealAlignMs = Date.now() - alignStartedAt; stageTimings.selfHealAlignMs = Date.now() - alignStartedAt;
alignment = aligned[0] ? { ...aligned[0], sentenceIndex: segment.original.segmentIndex } : null; alignment = aligned[0] ? { ...aligned[0], sentenceIndex: segment.original.segmentIndex } : null;
@ -636,10 +660,14 @@ export async function POST(request: NextRequest) {
failedStage = 'whisper.align'; failedStage = 'whisper.align';
const alignStartedAt = Date.now(); const alignStartedAt = Date.now();
const computeBackend = await getComputeBackend(); const computeBackend = await getComputeBackend();
const aligned = (await computeBackend.alignWords({ const aligned = (await withTimeout(
audioObjectKey: audioKey, computeBackend.alignWords({
text: segment.text, audioObjectKey: audioKey,
})).alignments; text: segment.text,
}),
getAlignmentTimeoutMs(computeBackend.mode),
`Whisper alignment (${computeBackend.mode})`,
)).alignments;
stageTimings.whisperAlignMs = Date.now() - alignStartedAt; stageTimings.whisperAlignMs = Date.now() - alignStartedAt;
alignment = aligned[0] ? { ...aligned[0], sentenceIndex: segment.original.segmentIndex } : null; alignment = aligned[0] ? { ...aligned[0], sentenceIndex: segment.original.segmentIndex } : null;
} catch (alignError) { } catch (alignError) {