Files
tonemark/src/lib/server/pipeline.ts
T
Giancarmine Salucci 0a418e773f D3: wire health preflight into pipeline before whisper submission
Add health preflight check in runJob() after ensureWhisperRunning()
and before submitJob(). Uses existing getHealth() from whisper.ts
with retryWithBackoff() for transient failures.

- Healthy whisper (status 'ok') -> job proceeds to submission
- Unreachable whisper -> job fails early with 'whisper health check failed: ...'
- Non-ok status (e.g. 'starting') -> job fails early with actual status
- getHealth() 5000ms timeout preserved, no new HTTP client

Tests: health-preflight.test.ts covers error message contracts
and pipeline integration (YouTube + upload source)
2026-07-09 03:59:52 +02:00

224 lines
8.8 KiB
TypeScript

import { createJob, updateJob, setJobStatus, getJob, resetJob } from './db.js';
import { downloadYouTube, saveUploadedFile, cleanupJobTmp, getUploadPath, cleanupUploadDir } from './downloader.js';
import { prepareAudio, cleanup as cleanupFiles } from './audio.js';
import { submitJob, streamJob, getHealth } from './whisper.js';
import { ensureWhisperRunning } from './docker.js';
import { retryWithBackoff } from './retry.js';
import type { AudioMode, Segment } from '$lib/types.js';
const WEBHOOK_BASE_URL = process.env.WEBHOOK_BASE_URL ?? 'http://localhost:3000';
/** Progress listeners: jobId → set of callbacks */
const progressListeners = new Map<string, Set<(data: string) => void>>();
export function subscribeProgress(jobId: string, cb: (data: string) => void): () => void {
if (!progressListeners.has(jobId)) progressListeners.set(jobId, new Set());
progressListeners.get(jobId)!.add(cb);
return () => progressListeners.get(jobId)?.delete(cb);
}
export function emitProgress(jobId: string, payload: object) {
const listeners = progressListeners.get(jobId);
if (!listeners) return;
const data = JSON.stringify(payload);
for (const cb of listeners) cb(data);
}
/** Start a transcription job for a YouTube URL. Runs async — returns immediately. */
export async function startYouTubeJob(
url: string,
audioMode: AudioMode = 'auto',
language?: string
): Promise<string> {
const job = createJob(url, 'Downloading…', audioMode);
runJob(job.id, { type: 'youtube', url }, audioMode, language).catch((err) => {
console.error(`[pipeline] job ${job.id} failed:`, err);
});
return job.id;
}
/** Start a transcription job for an uploaded file. Runs async — returns immediately. */
export async function startUploadJob(
buffer: Buffer,
filename: string,
audioMode: AudioMode = 'auto',
language?: string
): Promise<string> {
const job = createJob(filename, filename, audioMode);
runJob(job.id, { type: 'upload', buffer, filename }, audioMode, language).catch((err) => {
console.error(`[pipeline] job ${job.id} failed:`, err);
});
return job.id;
}
/** Retry a failed/cancelled job — YouTube or upload-source.
*
* Upload retry reads the original file asynchronously (fire-and-forget IIFE)
* so a missing file produces an async job failure instead of a synchronous
* 500 crash. The outer function returns immediately. */
export async function retryJob(jobId: string): Promise<void> {
const job = getJob(jobId);
if (!job) throw new Error('Job not found');
resetJob(jobId);
const isUpload = !job.source.startsWith('http');
if (isUpload) {
// Fire-and-forget: missing file -> async job failure, not 500
(async () => {
try {
const { readFile } = await import('fs/promises');
const uploadPath = getUploadPath(jobId, job.source);
const buffer = await readFile(uploadPath);
runJob(jobId, { type: 'upload', buffer, filename: job.source }, job.audioMode as AudioMode).catch((err) => {
console.error(`[pipeline] retry upload job ${jobId} failed:`, err);
});
} catch (err: unknown) {
const message = err instanceof Error ? err.message : String(err);
console.error(`[pipeline] retry upload job ${jobId} failed to read file:`, message);
updateJob({ id: jobId, status: 'failed', error: `retry failed: ${message}` });
emitProgress(jobId, { type: 'error', message: `retry failed: ${message}` });
}
})();
} else {
runJob(jobId, { type: 'youtube', url: job.source }, job.audioMode as AudioMode).catch((err) => {
console.error(`[pipeline] retry job ${jobId} failed:`, err);
});
}
}
async function runJob(
jobId: string,
input: { type: 'youtube'; url: string } | { type: 'upload'; buffer: Buffer; filename: string },
audioMode: AudioMode,
language?: string
) {
let rawAudioPath: string | null = null;
let wavPath: string | null = null;
try {
// ── 1. Download / save input ──────────────────────────────────────────
setJobStatus(jobId, 'downloading', 0);
emitProgress(jobId, { type: 'status', status: 'downloading' });
let title = 'Untitled';
let captionSegments: Segment[] | null = null;
if (input.type === 'youtube') {
const result = await downloadYouTube(input.url, jobId);
if (result.type === 'captions') {
// Fast path — use captions directly
title = result.title;
captionSegments = result.segments;
} else {
rawAudioPath = result.audioPath;
title = result.title;
}
} else {
rawAudioPath = await saveUploadedFile(input.buffer, input.filename, jobId);
title = input.filename.replace(/\.[^.]+$/, '');
// Remux webm files (MediaRecorder output) to fix Cues/duration
const { remuxUpload } = await import('./remux.js');
rawAudioPath = await remuxUpload(rawAudioPath, jobId);
}
updateJob({ id: jobId, title });
if (captionSegments) {
// Caption fast path — skip whisper
const { writeOutputs } = await import('./formatter.js');
const paths = await writeOutputs(captionSegments, title, jobId);
updateJob({
id: jobId,
status: 'done',
progress: 100,
segmentsJson: JSON.stringify(captionSegments),
outputDir: paths.srt.replace(/\/[^/]+$/, '')
});
emitProgress(jobId, { type: 'done' });
const { sendNotification } = await import('./push.js');
await sendNotification(jobId, '✅ Transcript ready', title);
await cleanupJobTmp(jobId);
await cleanupUploadDir(jobId).catch(() => {});
return;
}
// ── 2. Prepare audio ─────────────────────────────────────────────────
setJobStatus(jobId, 'preparing', 5);
emitProgress(jobId, { type: 'status', status: 'preparing' });
const { wavPath: wp, analysis } = await prepareAudio(rawAudioPath!, jobId, audioMode);
wavPath = wp;
updateJob({ id: jobId, meanVolume: analysis.meanVolume });
// ── 3. Ensure whisper is running ──────────────────────────────────────
await ensureWhisperRunning();
// ── 3b. Health preflight: verify whisper reachability and status ──────
await retryWithBackoff(
async () => {
try {
const health = await getHealth();
if (health.status !== 'ok') {
throw new Error(`whisper health check failed: status is "${health.status}"`);
}
} catch (err) {
if (err instanceof Error && err.message.startsWith('whisper health check failed')) {
throw err;
}
throw new Error(`whisper health check failed: ${err instanceof Error ? err.message : String(err)}`);
}
},
{
maxAttempts: 3,
baseDelayMs: 1000,
onAttempt: (err, attempt, delayMs) => {
console.warn(`[pipeline] whisper health check attempt ${attempt} failed:`, err);
emitProgress(jobId, { type: 'warning', message: `health check attempt ${attempt} failed, retrying in ${delayMs}ms` });
}
}
);
// ── 4. Submit to whisper with webhook ────────────────────────────────
setJobStatus(jobId, 'transcribing', 10);
emitProgress(jobId, { type: 'status', status: 'transcribing', progress: 10 });
const webhookUrl = `${WEBHOOK_BASE_URL}/api/webhook/${jobId}`;
const whisperJobId = await submitJob(wavPath, webhookUrl, language, (state, retryAfterSecs) => {
setJobStatus(jobId, 'warming_model', 10);
emitProgress(jobId, { type: 'model_warming', status: 'warming_model', state, retryAfterSecs, progress: 10 });
});
updateJob({ id: jobId, whisperJobId });
setJobStatus(jobId, 'transcribing', 10);
emitProgress(jobId, { type: 'status', status: 'transcribing', progress: 10 });
// ── 5. Open SSE for live progress (non-blocking relay) ───────────────
streamJob(
whisperJobId,
(percent, chunk, total) => {
const progress = 10 + Math.round(percent * 0.8);
setJobStatus(jobId, 'transcribing', progress);
emitProgress(jobId, { type: 'progress', percent, chunk, total, progress });
},
() => { /* webhook will handle completion */ },
(msg) => {
setJobStatus(jobId, 'failed', 0);
updateJob({ id: jobId, error: msg });
emitProgress(jobId, { type: 'error', message: msg });
}
).catch((err) => console.warn('[pipeline] SSE relay error:', err));
// Clean up wav after submitting (webhook handles the rest)
await cleanupFiles(wavPath);
wavPath = null;
} catch (err: unknown) {
const message = err instanceof Error ? err.message : String(err);
updateJob({ id: jobId, status: 'failed', error: message });
emitProgress(jobId, { type: 'error', message });
// rawAudioPath NOT cleaned here — for upload jobs the file lives in
// DATA_DIR/uploads/{jobId}/ so it survives for retry. For YouTube jobs,
// cleanupJobTmp below handles TMP_DIR/{jobId} which contains the audio.
if (wavPath) await cleanupFiles(wavPath).catch(() => {});
await cleanupJobTmp(jobId);
}
}