Merge branch 'fleet/story-d3-foundation'

This commit is contained in:
Giancarmine Salucci
2026-07-08 00:55:44 +02:00
7 changed files with 322 additions and 3 deletions
+28 -1
View File
@@ -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<string, unknown>): 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<Job> & { 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
});
}
+58
View File
@@ -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<T>(
fn: () => Promise<T>,
options: RetryOptions = {}
): Promise<T> {
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<void> {
return new Promise((resolve) => setTimeout(resolve, ms));
}
+23 -1
View File
@@ -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<boolean> {
try {
const { default: fetch } = await import('node-fetch');
@@ -236,3 +246,15 @@ export async function checkHealth(): Promise<boolean> {
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<WhisperHealth> {
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<WhisperHealth>;
}
+2
View File
@@ -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;
}
+23
View File
@@ -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 ──────────────────────────────────────────────────────────────
+129
View File
@@ -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);
});
});
+59 -1
View File
@@ -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);
});
});