feat: D3 foundation — retry schema, shared types, health client, backoff helper
- Add retry_count/next_retry_at to jobs table via guarded idempotent ALTER TABLE (PRAGMA table_info existence check, try/catch for safety on restart) - Add retryCount/nextRetryAt to Job type, rowToJob(), updateJob(), resetJob() - Add getHealth() returning structured WhisperHealth (status, gpu_name, vram_total_mb, model, queue_depth, model_state) distinct from existing checkHealth() boolean helper - Create retryWithBackoff() pure helper in src/lib/server/retry.ts with exponential backoff + jitter and onAttempt callback — no persistence owned - Tests: retry column read/write/reset, getHealth() smoke tests, retryWithBackoff() edge cases (success, failure, backoff timing, onAttempt)
This commit is contained in:
+28
-1
@@ -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
|
||||
});
|
||||
}
|
||||
|
||||
|
||||
@@ -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));
|
||||
}
|
||||
@@ -196,7 +196,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');
|
||||
@@ -206,3 +216,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>;
|
||||
}
|
||||
|
||||
@@ -41,6 +41,8 @@ export interface Job {
|
||||
outputDir: string | null;
|
||||
segmentsJson: string | null;
|
||||
error: string | null;
|
||||
retryCount: number | null;
|
||||
nextRetryAt: string | null;
|
||||
createdAt: string;
|
||||
updatedAt: string;
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user