import type { Db } from "../db.js"; import { nowIso, tx } from "../db.js"; import { newId, ID_PREFIX } from "./ids.js"; import type { JobStatus } from "./types.js"; export interface JobRow { id: string; type: string; status: JobStatus; payload_json: string; provider: string | null; provider_request_id: string | null; idempotency_key: string | null; attempts: number; max_attempts: number; run_after: string | null; locked_at: string | null; result_json: string | null; error: string | null; created_at: string; updated_at: string; } /** * Zapisuje zadanie PRZED pracą zewnętrzną. Idempotency key zapobiega * duplikatom przy retry klienta; provider_request_id pozwala wznowić * polling zamiast zlecać kolejną generację. */ export function enqueueJob( db: Db, input: { type: string; payload: Record; provider?: string | null; idempotencyKey?: string | null; maxAttempts?: number; } ): JobRow { const now = nowIso(); if (input.idempotencyKey) { const existing = db .prepare("SELECT * FROM jobs WHERE idempotency_key = ?") .get(input.idempotencyKey) as unknown as JobRow | undefined; if (existing) return existing; } const id = newId(ID_PREFIX.job); db.prepare( `INSERT INTO jobs (id, type, status, payload_json, provider, idempotency_key, max_attempts, created_at, updated_at) VALUES (?, ?, 'queued', ?, ?, ?, ?, ?, ?)` ).run( id, input.type, JSON.stringify(input.payload), input.provider ?? null, input.idempotencyKey ?? null, input.maxAttempts ?? 3, now, now ); return getJob(db, id)!; } export function getJob(db: Db, id: string): JobRow | null { const row = db.prepare("SELECT * FROM jobs WHERE id = ?").get(id) as unknown as | JobRow | undefined; return row ?? null; } export function listJobs(db: Db, limit = 50): JobRow[] { return db .prepare("SELECT * FROM jobs ORDER BY created_at DESC LIMIT ?") .all(limit) as unknown as JobRow[]; } /** * Atomowe pobranie kolejnego zadania. Status 'waiting' = czeka na wynik * providera (polling bez zużywania prób); 'queued' = do wykonania * (claim zużywa jedną próbę z max_attempts). */ export function claimNextJob(db: Db): JobRow | null { const now = nowIso(); return tx(db, () => { const row = db .prepare( `SELECT * FROM jobs WHERE (status = 'queued' OR status = 'waiting') AND (run_after IS NULL OR run_after <= ?) ORDER BY created_at LIMIT 1` ) .get(now) as unknown as JobRow | undefined; if (!row) return null; if (row.status === "queued") { db.prepare( "UPDATE jobs SET status = 'running', locked_at = ?, attempts = attempts + 1, updated_at = ? WHERE id = ?" ).run(now, now, row.id); return { ...row, status: "running", attempts: row.attempts + 1 }; } db.prepare( "UPDATE jobs SET status = 'running', locked_at = ?, updated_at = ? WHERE id = ?" ).run(now, now, row.id); return { ...row, status: "running" }; }); } /** Oznacza zadanie jako oczekujące na providera — polling bez zużycia prób. */ export function deferJob(db: Db, jobId: string, waitSeconds: number): void { const runAfter = new Date(Date.now() + waitSeconds * 1000).toISOString(); db.prepare( "UPDATE jobs SET status = 'waiting', run_after = ?, updated_at = ? WHERE id = ?" ).run(runAfter, nowIso(), jobId); } export function setJobRequestId(db: Db, jobId: string, providerRequestId: string): void { db.prepare( "UPDATE jobs SET provider_request_id = ?, updated_at = ? WHERE id = ?" ).run(providerRequestId, nowIso(), jobId); } export function finishJob( db: Db, jobId: string, outcome: { ok: true; result?: Record } | { ok: false; error: string; retryable?: boolean } ): void { const now = nowIso(); const job = getJob(db, jobId); if (!job) return; if (outcome.ok) { db.prepare( "UPDATE jobs SET status = 'done', result_json = ?, error = NULL, updated_at = ? WHERE id = ?" ).run(JSON.stringify(outcome.result ?? {}), now, jobId); return; } const retryable = outcome.retryable ?? false; if (retryable && job.attempts < job.max_attempts) { const backoffSec = Math.min(60, 5 * job.attempts); const runAfter = new Date(Date.now() + backoffSec * 1000).toISOString(); db.prepare( "UPDATE jobs SET status = 'queued', run_after = ?, error = ?, updated_at = ? WHERE id = ?" ).run(runAfter, outcome.error, now, jobId); } else { db.prepare( "UPDATE jobs SET status = 'error', error = ?, updated_at = ? WHERE id = ?" ).run(outcome.error, now, jobId); } }