diff --git a/src/lib/server/db.ts b/src/lib/server/db.ts index d4df810..1d01220 100644 --- a/src/lib/server/db.ts +++ b/src/lib/server/db.ts @@ -4,6 +4,18 @@ import { existsSync, mkdirSync } from 'fs'; import { join } from 'path'; import type { Job, JobStatus, AudioMode, PushSubscription } from '$lib/types.js'; +/** + * Add a column to a table if it does not already exist. + * Uses PRAGMA table_info for an existence check, consistent with the + * existing CREATE TABLE IF NOT EXISTS idiom — no migration tool introduced. + */ +function addColumnIfMissing(table: string, column: string, typeDef: string): void { + const columns = db.pragma(`table_info(${table})`) as { name: string }[]; + if (!columns.some((c) => c.name === column)) { + db.exec(`ALTER TABLE ${table} ADD COLUMN ${column} ${typeDef}`); + } +} + const DATA_DIR = process.env.DATA_DIR ?? join(process.env.HOME ?? '/tmp', '.whisper-pwa'); if (!existsSync(DATA_DIR)) mkdirSync(DATA_DIR, { recursive: true }); @@ -37,6 +49,14 @@ db.exec(` ); `); +// Add retry columns idempotently (may already exist on restart) +try { + addColumnIfMissing('jobs', 'retry_count', 'INTEGER DEFAULT 0'); + addColumnIfMissing('jobs', 'next_retry_at', 'TEXT'); +} catch (e) { + console.warn('[db] failed to add retry columns:', e); +} + const stmts = { insertJob: db.prepare(` INSERT INTO jobs (id, status, title, source, audio_mode) @@ -52,6 +72,8 @@ const stmts = { output_dir = @outputDir, segments_json = @segmentsJson, error = @error, + retry_count = @retryCount, + next_retry_at = @nextRetryAt, updated_at = datetime('now') WHERE id = @id `), @@ -64,6 +86,7 @@ const stmts = { status = 'pending', progress = 0, error = NULL, mean_volume = NULL, whisper_job_id = NULL, output_dir = NULL, segments_json = NULL, + retry_count = 0, next_retry_at = NULL, updated_at = datetime('now') WHERE id = @id `), @@ -92,6 +115,8 @@ function rowToJob(row: Record): Job { outputDir: row.output_dir as string | null, segmentsJson: row.segments_json as string | null, error: row.error as string | null, + retryCount: row.retry_count as number | null, + nextRetryAt: row.next_retry_at as string | null, createdAt: row.created_at as string, updatedAt: row.updated_at as string }; @@ -125,7 +150,9 @@ export function updateJob(job: Partial & { id: string }): void { progress: merged.progress, outputDir: merged.outputDir, segmentsJson: merged.segmentsJson, - error: merged.error + error: merged.error, + retryCount: merged.retryCount ?? null, + nextRetryAt: merged.nextRetryAt ?? null }); } diff --git a/src/lib/server/retry.ts b/src/lib/server/retry.ts new file mode 100644 index 0000000..2c7a948 --- /dev/null +++ b/src/lib/server/retry.ts @@ -0,0 +1,58 @@ +/** + * Shared bounded retry helper — no persistence owned here. + * + * Retries an async function with exponential backoff + jitter. + * The caller owns the retry decision; this helper only drives the + * timing and calls `onAttempt` after each failure for observability + * (e.g. emitting SSE error/status events — same event types, not new ones). + */ + +export interface RetryOptions { + /** Maximum number of attempts before rejecting (default: 3). */ + maxAttempts?: number; + /** Base delay in ms before the first retry (doubled each attempt, default: 1000). */ + baseDelayMs?: number; + /** + * Called after each failed attempt with the error, attempt number (1-indexed), + * and delay before the next retry. Use to emit SSE events or log warnings. + */ + onAttempt?: (error: unknown, attempt: number, delayMs: number) => void; +} + +/** + * Retry `fn` with exponential backoff + jitter. + * + * - Resolves with the value of `fn()` on first success. + * - Rejects with the last error after `maxAttempts` total tries. + * - `onAttempt` is called after each failure (never on success). + * - Does NOT persist anything itself — the caller manages queue state. + */ +export async function retryWithBackoff( + fn: () => Promise, + options: RetryOptions = {} +): Promise { + const { maxAttempts = 3, baseDelayMs = 1000, onAttempt } = options; + let lastError: unknown; + + for (let attempt = 1; attempt <= maxAttempts; attempt++) { + try { + return await fn(); + } catch (err) { + lastError = err; + if (attempt < maxAttempts) { + // Exponential backoff: delay = base * 2^(attempt-1) + random jitter up to 25% + const delay = baseDelayMs * Math.pow(2, attempt - 1); + const jitter = Math.random() * delay * 0.25; + const totalDelay = Math.round(delay + jitter); + onAttempt?.(err, attempt, totalDelay); + await sleep(totalDelay); + } + } + } + + throw lastError; +} + +function sleep(ms: number): Promise { + return new Promise((resolve) => setTimeout(resolve, ms)); +} diff --git a/src/lib/server/whisper.ts b/src/lib/server/whisper.ts index 32f442b..a7edb52 100644 --- a/src/lib/server/whisper.ts +++ b/src/lib/server/whisper.ts @@ -226,7 +226,17 @@ export async function streamJob( } } -/** Check if the whisper server is healthy. */ +/** Structured health info returned by the whisper /health endpoint. */ +export interface WhisperHealth { + status: string; + gpu_name?: string; + vram_total_mb?: number; + model?: string; + queue_depth?: number; + model_state?: string; +} + +/** Check if the whisper server is healthy (boolean only). */ export async function checkHealth(): Promise { try { const { default: fetch } = await import('node-fetch'); @@ -236,3 +246,15 @@ export async function checkHealth(): Promise { return false; } } + +/** + * Fetch structured health info from whisper-rtx2080. + * Returns the full JSON body: status, gpu_name, vram_total_mb, model, queue_depth, model_state. + * Throws on network error or non-ok response. + */ +export async function getHealth(): Promise { + const { default: fetch } = await import('node-fetch'); + const res = await fetch(`${whisperUrl()}/health`, { signal: AbortSignal.timeout(5000) }); + if (!res.ok) throw new Error(`/health returned ${res.status}`); + return res.json() as Promise; +} diff --git a/src/lib/types.ts b/src/lib/types.ts index 532e4c8..455522e 100644 --- a/src/lib/types.ts +++ b/src/lib/types.ts @@ -50,6 +50,8 @@ export interface Job { outputDir: string | null; segmentsJson: string | null; error: string | null; + retryCount: number | null; + nextRetryAt: string | null; createdAt: string; updatedAt: string; } diff --git a/src/tests/db.test.ts b/src/tests/db.test.ts index 4d3b17d..5b69661 100644 --- a/src/tests/db.test.ts +++ b/src/tests/db.test.ts @@ -14,6 +14,7 @@ import { listJobs, updateJob, setJobStatus, + resetJob, savePushSubscription, getAllSubscriptions, deletePushSubscription @@ -107,6 +108,28 @@ describe('updateJob', () => { it('does nothing for an unknown id', () => { expect(() => updateJob({ id: 'no-such-id', title: 'Ghost' })).not.toThrow(); }); + + it('stores and reads retryCount and nextRetryAt', () => { + const job = createJob('src', 'retry test', 'auto'); + // Default retry_count is 0 (DEFAULT 0 in schema) + expect(getJob(job.id)!.retryCount).toBe(0); + expect(getJob(job.id)!.nextRetryAt).toBeNull(); + + const nextRetry = new Date(Date.now() + 60000).toISOString(); + updateJob({ id: job.id, retryCount: 2, nextRetryAt: nextRetry }); + const updated = getJob(job.id)!; + expect(updated.retryCount).toBe(2); + expect(updated.nextRetryAt).toBe(nextRetry); + }); + + it('resetJob clears retry fields', () => { + const job = createJob('src', 'reset retry', 'auto'); + updateJob({ id: job.id, retryCount: 3, nextRetryAt: '2026-07-08T12:00:00.000Z' }); + resetJob(job.id); + const after = getJob(job.id)!; + expect(after.retryCount).toBe(0); + expect(after.nextRetryAt).toBeNull(); + }); }); // ── setJobStatus ────────────────────────────────────────────────────────────── diff --git a/src/tests/retry.test.ts b/src/tests/retry.test.ts new file mode 100644 index 0000000..5a3d461 --- /dev/null +++ b/src/tests/retry.test.ts @@ -0,0 +1,129 @@ +import { describe, it, expect, vi, beforeEach, afterEach } from 'vitest'; +import { retryWithBackoff } from '$lib/server/retry.js'; + +describe('retryWithBackoff', () => { + beforeEach(() => { + vi.useFakeTimers(); + }); + + afterEach(() => { + vi.useRealTimers(); + }); + + it('resolves immediately when fn succeeds on first attempt', async () => { + const fn = vi.fn().mockResolvedValue('ok'); + const result = await retryWithBackoff(fn, { maxAttempts: 3 }); + expect(result).toBe('ok'); + expect(fn).toHaveBeenCalledTimes(1); + }); + + it('retries when fn fails and resolves on subsequent attempt', async () => { + const fn = vi.fn() + .mockRejectedValueOnce(new Error('fail 1')) + .mockRejectedValueOnce(new Error('fail 2')) + .mockResolvedValue('recovered'); + + const promise = retryWithBackoff(fn, { maxAttempts: 5, baseDelayMs: 10 }); + + // Advance through timers for each retry + await vi.runAllTimersAsync(); + const result = await promise; + expect(result).toBe('recovered'); + expect(fn).toHaveBeenCalledTimes(3); + }); + + it('rejects after maxAttempts exhausted without success', async () => { + const error = new Error('persistent failure'); + const fn = vi.fn().mockRejectedValue(error); + + const promise = retryWithBackoff(fn, { maxAttempts: 3, baseDelayMs: 10 }); + promise.catch(() => {}); + await vi.runAllTimersAsync(); + await expect(promise).rejects.toThrow('persistent failure'); + expect(fn).toHaveBeenCalledTimes(3); + }); + + it('calls onAttempt with error, attempt number, and delay after each failure', async () => { + const fn = vi.fn() + .mockRejectedValueOnce(new Error('fail 1')) + .mockRejectedValueOnce(new Error('fail 2')) + .mockResolvedValue('done'); + const onAttempt = vi.fn(); + + const promise = retryWithBackoff(fn, { + maxAttempts: 5, + baseDelayMs: 100, + onAttempt + }); + + await vi.runAllTimersAsync(); + await promise; + + // First failure: attempt=1, delay=baseDelayMs*2^0 + jitter ≈ 100-125ms + // Second failure: attempt=2, delay=baseDelayMs*2^1 + jitter ≈ 200-250ms + expect(onAttempt).toHaveBeenCalledTimes(2); + expect(onAttempt).toHaveBeenNthCalledWith( + 1, + expect.objectContaining({ message: 'fail 1' }), + 1, + expect.any(Number) + ); + expect(onAttempt).toHaveBeenNthCalledWith( + 2, + expect.objectContaining({ message: 'fail 2' }), + 2, + expect.any(Number) + ); + + // Verify delays are increasing (exponential backoff) + const delay1 = onAttempt.mock.calls[0][2]; + const delay2 = onAttempt.mock.calls[1][2]; + expect(delay2).toBeGreaterThan(delay1); + }); + + it('uses default options when none provided', async () => { + const fn = vi.fn().mockResolvedValue('defaults'); + const result = await retryWithBackoff(fn); + expect(result).toBe('defaults'); + }); + + it('does not call onAttempt on success path', async () => { + const fn = vi.fn().mockResolvedValue('ok'); + const onAttempt = vi.fn(); + await retryWithBackoff(fn, { maxAttempts: 3, onAttempt }); + expect(onAttempt).not.toHaveBeenCalled(); + }); + + it('uses increasing delays with exponential backoff', async () => { + const fn = vi.fn().mockRejectedValue(new Error('always fail')); + const onAttempt = vi.fn(); + + const promise = retryWithBackoff(fn, { + maxAttempts: 4, + baseDelayMs: 100, + onAttempt + }); + promise.catch(() => {}); + await vi.runAllTimersAsync(); + await expect(promise).rejects.toThrow('always fail'); + expect(fn).toHaveBeenCalledTimes(4); + + const delays = onAttempt.mock.calls.map((c: unknown[]) => c[2] as number); + // delay=100 * 2^(n-1) + jitter: base is 100, 200, 400 + expect(delays[0]).toBeGreaterThanOrEqual(100); + expect(delays[0]).toBeLessThan(130); + expect(delays[1]).toBeGreaterThanOrEqual(200); + expect(delays[1]).toBeLessThan(260); + expect(delays[2]).toBeGreaterThanOrEqual(400); + expect(delays[2]).toBeLessThan(520); + }); + + it('defaults to 3 maxAttempts', async () => { + const fn = vi.fn().mockRejectedValue(new Error('fail')); + const promise = retryWithBackoff(fn, { baseDelayMs: 10 }); + promise.catch(() => {}); + await vi.runAllTimersAsync(); + await expect(promise).rejects.toThrow('fail'); + expect(fn).toHaveBeenCalledTimes(3); + }); +}); diff --git a/src/tests/whisper.test.ts b/src/tests/whisper.test.ts index e420870..a7aa7ae 100644 --- a/src/tests/whisper.test.ts +++ b/src/tests/whisper.test.ts @@ -22,7 +22,7 @@ vi.mock('form-data', () => ({ vi.mock('fs', () => ({ createReadStream: mocks.createReadStream })); -import { submitJob, streamJob, getModelStatus, cancelJob, unloadModel } from '$lib/server/whisper.js'; +import { submitJob, streamJob, getModelStatus, cancelJob, unloadModel, checkHealth, getHealth } from '$lib/server/whisper.js'; afterEach(() => vi.clearAllMocks()); @@ -669,3 +669,61 @@ describe('streamJob — SSE event parsing', () => { expect(onProgress).toHaveBeenCalledWith(60, 0, 0); }); }); + +// ── getHealth ───────────────────────────────────────────────────────────────── + +describe('getHealth', () => { + it('returns structured health when whisper responds', async () => { + const healthData = { + status: 'ok', + gpu_name: 'NVIDIA GeForce RTX 2080 Ti', + vram_total_mb: 11264, + model: 'ggml-base-q5_1', + queue_depth: 0, + model_state: 'loaded' + }; + mocks.fetch.mockResolvedValue({ + ok: true, + json: () => Promise.resolve(healthData) + }); + + const result = await getHealth(); + expect(result.status).toBe('ok'); + expect(result.gpu_name).toBe('NVIDIA GeForce RTX 2080 Ti'); + expect(result.vram_total_mb).toBe(11264); + expect(result.model).toBe('ggml-base-q5_1'); + expect(result.queue_depth).toBe(0); + expect(result.model_state).toBe('loaded'); + }); + + it('calls /health on the configured WHISPER_URL', async () => { + vi.stubEnv('WHISPER_URL', 'http://whisper-rtx2080:8080'); + mocks.fetch.mockResolvedValue({ + ok: true, + json: () => Promise.resolve({ status: 'ok' }) + }); + await getHealth(); + expect(mocks.fetch).toHaveBeenCalledWith( + 'http://whisper-rtx2080:8080/health', + expect.objectContaining({ signal: expect.anything() }) + ); + vi.unstubAllEnvs(); + }); + + it('throws on non-ok response', async () => { + mocks.fetch.mockResolvedValue({ ok: false, status: 503 }); + await expect(getHealth()).rejects.toThrow('/health'); + }); + + it('distinct from checkHealth — checkHealth returns boolean', async () => { + mocks.fetch.mockResolvedValue({ ok: true, json: () => Promise.resolve({ status: 'ok' }) }); + const healthResult = await getHealth(); + expect(typeof healthResult).toBe('object'); + expect(healthResult.status).toBe('ok'); + + mocks.fetch.mockResolvedValue({ ok: true }); + const checkResult = await checkHealth(); + expect(typeof checkResult).toBe('boolean'); + expect(checkResult).toBe(true); + }); +});