From f03b757174a0de1d176de5c12d1bf17f1ac3d7a6 Mon Sep 17 00:00:00 2001 From: Giancarmine Salucci Date: Wed, 8 Jul 2026 01:25:49 +0200 Subject: [PATCH] D2: Widen retry to upload-source jobs MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - Remove http-only guard in retry endpoint (src/routes/api/jobs/[id]/retry/+server.ts) - Extend pipeline.ts retryJob() to handle upload-source jobs: read file from persistent storage, pass buffer to runJob - Change saveUploadedFile() to save to DATA_DIR/uploads/{jobId}/ (survives cleanupJobTmp, available for retry on failure) - Add getUploadPath(), cleanupUploadDir() to downloader.ts - Remove cleanupFiles(rawAudioPath) from catch block in runJob() — cleanupJobTmp handles TMP_DIR for YouTube, uploads survive - Add cleanupUploadDir() on webhook success path - Add retry-endpoint.test.ts with 9 tests covering: - 404 for unknown job - 409 for non-retryable status (done, pending) - 200 for cancelled (retryable) - 200 for failed YouTube job (existing behavior) - 200 for failed upload-source job (new behavior, .webm and .mp3) - Negative: retryJob not called on 404/409 - Add downloader.test.ts tests for saveUploadedFile persistent path, getUploadPath, cleanupUploadDir, isolation, and non-existent cleanup - Fix webhook.test.ts mock to include cleanupUploadDir All 189 tests pass across 12 test files. --- src/lib/server/downloader.ts | 30 +++++- src/lib/server/pipeline.ts | 26 ++++-- src/routes/api/jobs/[id]/retry/+server.ts | 3 - src/routes/api/webhook/[jobId]/+server.ts | 3 +- src/tests/downloader.test.ts | 57 ++++++++++- src/tests/retry-endpoint.test.ts | 109 ++++++++++++++++++++++ src/tests/webhook.test.ts | 6 +- 7 files changed, 217 insertions(+), 17 deletions(-) create mode 100644 src/tests/retry-endpoint.test.ts diff --git a/src/lib/server/downloader.ts b/src/lib/server/downloader.ts index a641a59..60241e1 100644 --- a/src/lib/server/downloader.ts +++ b/src/lib/server/downloader.ts @@ -6,12 +6,34 @@ import { join } from 'path'; import { fetchTranscript, type TranscriptResponse } from 'youtube-transcript'; const execFileAsync = promisify(execFile); -const TMP_DIR = join(process.env.DATA_DIR ?? '/tmp/.whisper-pwa', 'downloads'); +const DATA_DIR = process.env.DATA_DIR ?? '/tmp/.whisper-pwa'; +const TMP_DIR = join(DATA_DIR, 'downloads'); +const UPLOADS_DIR = join(DATA_DIR, 'uploads'); export async function ensureTmpDir() { if (!existsSync(TMP_DIR)) await mkdir(TMP_DIR, { recursive: true }); } +async function ensureUploadsDir(jobId: string): Promise { + const dir = join(UPLOADS_DIR, jobId); + if (!existsSync(dir)) await mkdir(dir, { recursive: true }); + return dir; +} + +/** Path to a job's persistent upload file. Survives cleanupJobTmp. */ +export function getUploadPath(jobId: string, filename: string): string { + return join(UPLOADS_DIR, jobId, filename); +} + +/** Delete a job's persistent upload dir (call on success). */ +export async function cleanupUploadDir(jobId: string): Promise { + const dir = join(UPLOADS_DIR, jobId); + try { + const { rm } = await import('fs/promises'); + await rm(dir, { recursive: true, force: true }); + } catch { /* ignore */ } +} + export interface CaptionResult { type: 'captions'; segments: Array<{ index: number; start: number; end: number; text: string; words: [] }>; @@ -93,15 +115,13 @@ export async function downloadYouTube(url: string, jobId: string): Promise { - await ensureTmpDir(); - const outDir = join(TMP_DIR, jobId); - await mkdir(outDir, { recursive: true }); + const outDir = await ensureUploadsDir(jobId); const dest = join(outDir, filename); await writeFile(dest, buffer); return dest; diff --git a/src/lib/server/pipeline.ts b/src/lib/server/pipeline.ts index b3c89ff..828b3bb 100644 --- a/src/lib/server/pipeline.ts +++ b/src/lib/server/pipeline.ts @@ -1,5 +1,5 @@ import { createJob, updateJob, setJobStatus, getJob, resetJob } from './db.js'; -import { downloadYouTube, saveUploadedFile, cleanupJobTmp } from './downloader.js'; +import { downloadYouTube, saveUploadedFile, cleanupJobTmp, getUploadPath, cleanupUploadDir } from './downloader.js'; import { prepareAudio, cleanup as cleanupFiles } from './audio.js'; import { submitJob, streamJob } from './whisper.js'; import { ensureWhisperRunning } from './docker.js'; @@ -50,14 +50,25 @@ export async function startUploadJob( return job.id; } -/** Retry a failed/cancelled YouTube job by resetting and re-running the pipeline. */ +/** Retry a failed/cancelled job — YouTube or upload-source. */ export async function retryJob(jobId: string): Promise { const job = getJob(jobId); if (!job) throw new Error('Job not found'); resetJob(jobId); - runJob(jobId, { type: 'youtube', url: job.source }, job.audioMode as AudioMode).catch((err) => { - console.error(`[pipeline] retry job ${jobId} failed:`, err); - }); + + const isUpload = !job.source.startsWith('http'); + if (isUpload) { + const { readFileSync } = await import('fs'); + const uploadPath = getUploadPath(jobId, job.source); + const buffer = readFileSync(uploadPath); + runJob(jobId, { type: 'upload', buffer, filename: job.source }, job.audioMode as AudioMode).catch((err) => { + console.error(`[pipeline] retry upload job ${jobId} failed:`, err); + }); + } 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( @@ -112,6 +123,7 @@ async function runJob( const { sendNotification } = await import('./push.js'); await sendNotification(jobId, '✅ Transcript ready', title); await cleanupJobTmp(jobId); + await cleanupUploadDir(jobId).catch(() => {}); return; } @@ -162,7 +174,9 @@ async function runJob( const message = err instanceof Error ? err.message : String(err); updateJob({ id: jobId, status: 'failed', error: message }); emitProgress(jobId, { type: 'error', message }); - if (rawAudioPath) await cleanupFiles(rawAudioPath).catch(() => {}); + // 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); } diff --git a/src/routes/api/jobs/[id]/retry/+server.ts b/src/routes/api/jobs/[id]/retry/+server.ts index 80810d2..bb286b9 100644 --- a/src/routes/api/jobs/[id]/retry/+server.ts +++ b/src/routes/api/jobs/[id]/retry/+server.ts @@ -8,9 +8,6 @@ export async function POST({ params }) { if (!['failed', 'cancelled'].includes(job.status)) { throw error(409, 'Only failed or cancelled jobs can be retried'); } - if (!job.source.startsWith('http')) { - throw error(422, 'Cannot retry a file upload — please re-upload the file'); - } await retryJob(params.id); return json({ ok: true }); } diff --git a/src/routes/api/webhook/[jobId]/+server.ts b/src/routes/api/webhook/[jobId]/+server.ts index 320fd00..e7102a4 100644 --- a/src/routes/api/webhook/[jobId]/+server.ts +++ b/src/routes/api/webhook/[jobId]/+server.ts @@ -2,7 +2,7 @@ import { json, error } from '@sveltejs/kit'; import { getJob, updateJob, setJobStatus } from '$lib/server/db.js'; import { writeOutputs } from '$lib/server/formatter.js'; import { sendNotification } from '$lib/server/push.js'; -import { cleanupJobTmp } from '$lib/server/downloader.js'; +import { cleanupJobTmp, cleanupUploadDir } from '$lib/server/downloader.js'; import { emitProgress } from '$lib/server/pipeline.js'; import type { Segment, WhisperJob } from '$lib/types.js'; @@ -79,6 +79,7 @@ export async function POST({ params, request }) { await sendNotification(jobId, '✅ Transcript ready', job.title); await cleanupJobTmp(jobId); + await cleanupUploadDir(jobId).catch(() => {}); return json({ ok: true }); } catch (err: unknown) { diff --git a/src/tests/downloader.test.ts b/src/tests/downloader.test.ts index c04fac2..5b56977 100644 --- a/src/tests/downloader.test.ts +++ b/src/tests/downloader.test.ts @@ -18,7 +18,7 @@ vi.mock('youtube-transcript', () => ({ fetchTranscript: mockFetchTranscript })); -import { downloadYouTube, transcriptEntriesToSegments } from '$lib/server/downloader.js'; +import { downloadYouTube, transcriptEntriesToSegments, saveUploadedFile, getUploadPath, cleanupUploadDir } from '$lib/server/downloader.js'; beforeEach(() => { vi.clearAllMocks(); @@ -57,6 +57,61 @@ describe('transcriptEntriesToSegments', () => { }); }); +describe('saveUploadedFile', () => { + it('saves file to persistent DATA_DIR/uploads/{jobId}/ location', async () => { + const buffer = Buffer.from('fake audio content'); + const filePath = await saveUploadedFile(buffer, 'test-recording.webm', 'upload-test-1'); + + // Must be under DATA_DIR/uploads/, NOT under TMP_DIR/downloads/ + expect(filePath).toContain('/uploads/upload-test-1/test-recording.webm'); + expect(filePath).not.toContain('/downloads/'); + + // File must exist on disk + const { accessSync } = await import('fs'); + expect(() => accessSync(filePath)).not.toThrow(); + + // Content must match + const { readFileSync } = await import('fs'); + expect(readFileSync(filePath)).toEqual(buffer); + + // cleanup + await cleanupUploadDir('upload-test-1'); + }); + + it('getUploadPath returns correct path without creating file', async () => { + const p = getUploadPath('some-job', 'audio.mp3'); + expect(p).toContain('/uploads/some-job/audio.mp3'); + expect(p).not.toContain('/downloads/'); + }); + + it('cleanupUploadDir removes the job upload directory', async () => { + const buffer = Buffer.from('data'); + await saveUploadedFile(buffer, 'temp-file.txt', 'cleanup-test'); + const p = getUploadPath('cleanup-test', 'temp-file.txt'); + + const { accessSync } = await import('fs'); + expect(() => accessSync(p)).not.toThrow(); + + await cleanupUploadDir('cleanup-test'); + expect(() => accessSync(p)).toThrow(); + }); + + it('cleanupUploadDir does not throw for non-existent dir', async () => { + await expect(cleanupUploadDir('non-existent-job')).resolves.toBeUndefined(); + }); + + it('each job gets an isolated upload directory', async () => { + const p1 = await saveUploadedFile(Buffer.from('a'), 'a.wav', 'job-iso-1'); + const p2 = await saveUploadedFile(Buffer.from('b'), 'b.wav', 'job-iso-2'); + + expect(p1).toContain('/uploads/job-iso-1/a.wav'); + expect(p2).toContain('/uploads/job-iso-2/b.wav'); + + await cleanupUploadDir('job-iso-1'); + await cleanupUploadDir('job-iso-2'); + }); +}); + describe('downloadYouTube', () => { it('uses fetched transcript entries directly for caption jobs', async () => { mockFetchTranscript.mockResolvedValue([ diff --git a/src/tests/retry-endpoint.test.ts b/src/tests/retry-endpoint.test.ts new file mode 100644 index 0000000..07b553b --- /dev/null +++ b/src/tests/retry-endpoint.test.ts @@ -0,0 +1,109 @@ +import { describe, it, expect, vi, beforeEach } from 'vitest'; +import type { Job } from '$lib/types.js'; + +const { mockGetJob, mockRetryJob } = vi.hoisted(() => ({ + mockGetJob: vi.fn(), + mockRetryJob: vi.fn() +})); + +vi.mock('$lib/server/db.js', () => ({ + getJob: mockGetJob +})); + +vi.mock('$lib/server/pipeline.js', () => ({ + retryJob: mockRetryJob +})); + +import { POST } from '$lib/../routes/api/jobs/[id]/retry/+server.js'; + +function makeEvent(jobId: string) { + return { params: { id: jobId } }; +} + +function makeJob(overrides: Partial = {}): Job { + return { + id: 'test-job', + status: 'failed', + title: 'Test', + source: 'https://example.com/video', + audioMode: 'auto', + meanVolume: null, + whisperJobId: null, + progress: 0, + outputDir: null, + segmentsJson: null, + error: 'Something failed', + retryCount: 0, + nextRetryAt: null, + createdAt: new Date().toISOString(), + updatedAt: new Date().toISOString(), + ...overrides + }; +} + +describe('POST /api/jobs/:id/retry', () => { + beforeEach(() => { + vi.clearAllMocks(); + mockRetryJob.mockResolvedValue(undefined); + }); + + it('returns 404 when job does not exist', async () => { + mockGetJob.mockReturnValue(null); + await expect(POST(makeEvent('ghost') as any)).rejects.toMatchObject({ status: 404 }); + }); + + it('returns 409 when job is done (non-retryable status)', async () => { + mockGetJob.mockReturnValue(makeJob({ status: 'done' })); + await expect(POST(makeEvent('job-done') as any)).rejects.toMatchObject({ status: 409 }); + }); + + it('returns 409 for pending status', async () => { + mockGetJob.mockReturnValue(makeJob({ status: 'pending' })); + await expect(POST(makeEvent('job-pending') as any)).rejects.toMatchObject({ status: 409 }); + }); + + it('allows retry for cancelled status', async () => { + mockGetJob.mockReturnValue(makeJob({ status: 'cancelled', source: 'https://youtube.com/watch?v=abc' })); + const res = await POST(makeEvent('job-cancelled') as any); + expect(res.status).toBe(200); + expect(await res.json()).toEqual({ ok: true }); + }); + + it('returns 200 and calls retryJob for failed YouTube job', async () => { + const jobId = 'yt-job-1'; + mockGetJob.mockReturnValue(makeJob({ id: jobId, source: 'https://youtube.com/watch?v=abc' })); + const res = await POST(makeEvent(jobId) as any); + expect(res.status).toBe(200); + expect(await res.json()).toEqual({ ok: true }); + expect(mockRetryJob).toHaveBeenCalledWith(jobId); + }); + + it('returns 200 and calls retryJob for failed upload-source job (guard removed)', async () => { + const jobId = 'upload-job-1'; + mockGetJob.mockReturnValue(makeJob({ id: jobId, source: 'my-recording.webm' })); + const res = await POST(makeEvent(jobId) as any); + expect(res.status).toBe(200); + expect(await res.json()).toEqual({ ok: true }); + expect(mockRetryJob).toHaveBeenCalledWith(jobId); + }); + + it('returns 200 for upload with arbitrary filename (not just .webm)', async () => { + const jobId = 'upload-job-2'; + mockGetJob.mockReturnValue(makeJob({ id: jobId, source: 'podcast.mp3' })); + const res = await POST(makeEvent(jobId) as any); + expect(res.status).toBe(200); + expect(mockRetryJob).toHaveBeenCalledWith(jobId); + }); + + it('does not call retryJob when job not found', async () => { + mockGetJob.mockReturnValue(null); + await expect(POST(makeEvent('ghost') as any)).rejects.toThrow(); + expect(mockRetryJob).not.toHaveBeenCalled(); + }); + + it('does not call retryJob for non-retryable status', async () => { + mockGetJob.mockReturnValue(makeJob({ status: 'done' })); + await expect(POST(makeEvent('job-done') as any)).rejects.toThrow(); + expect(mockRetryJob).not.toHaveBeenCalled(); + }); +}); diff --git a/src/tests/webhook.test.ts b/src/tests/webhook.test.ts index 282d4c3..f32440f 100644 --- a/src/tests/webhook.test.ts +++ b/src/tests/webhook.test.ts @@ -10,6 +10,7 @@ const { mockWriteOutputs, mockSendNotification, mockCleanupJobTmp, + mockCleanupUploadDir, mockEmitProgress } = vi.hoisted(() => ({ mockGetJob: vi.fn(), @@ -18,6 +19,7 @@ const { mockWriteOutputs: vi.fn(), mockSendNotification: vi.fn(), mockCleanupJobTmp: vi.fn(), + mockCleanupUploadDir: vi.fn(), mockEmitProgress: vi.fn() })); @@ -36,7 +38,8 @@ vi.mock('$lib/server/push.js', () => ({ })); vi.mock('$lib/server/downloader.js', () => ({ - cleanupJobTmp: mockCleanupJobTmp + cleanupJobTmp: mockCleanupJobTmp, + cleanupUploadDir: mockCleanupUploadDir })); vi.mock('$lib/server/pipeline.js', () => ({ @@ -85,6 +88,7 @@ function makeSeg(index: number, text: string): Segment { beforeEach(() => { vi.clearAllMocks(); + mockCleanupUploadDir.mockResolvedValue(undefined); mockWriteOutputs.mockResolvedValue({ srt: '/out/dir/title.srt', txt: '/out/dir/title.txt',