D2: Widen retry to upload-source jobs

- 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.
This commit is contained in:
Giancarmine Salucci
2026-07-08 01:25:49 +02:00
parent d633f74689
commit f03b757174
7 changed files with 217 additions and 17 deletions
+25 -5
View File
@@ -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<string> {
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<void> {
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<Downl
return { type: 'audio', audioPath, title };
}
/** Save an uploaded file to tmp dir and return its path. */
/** Save an uploaded file to persistent storage (survives cleanupJobTmp). */
export async function saveUploadedFile(
buffer: Buffer,
filename: string,
jobId: string
): Promise<string> {
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;
+17 -3
View File
@@ -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<void> {
const job = getJob(jobId);
if (!job) throw new Error('Job not found');
resetJob(jobId);
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);
}
@@ -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 });
}
+2 -1
View File
@@ -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) {
+56 -1
View File
@@ -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([
+109
View File
@@ -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> = {}): 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();
});
});
+5 -1
View File
@@ -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',