fix: idempotent job completion + monthly budget cap enforcement
Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01XNQ8ghPfzAfsyVYd6HgFb6
This commit is contained in:
+25
-2
@@ -10,6 +10,17 @@ export async function startImageWorker(): Promise<void> {
|
|||||||
await startWorker(concurrency, handle);
|
await startWorker(concurrency, handle);
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/** Monatliches Budget erreicht? (Summe usage.cost im laufenden Kalendermonat.) */
|
||||||
|
async function budgetExceeded(): Promise<boolean> {
|
||||||
|
const s = await one<{ monthly_budget: number | null }>('SELECT monthly_budget FROM settings WHERE id=1');
|
||||||
|
const cap = s?.monthly_budget ? Number(s.monthly_budget) : 0;
|
||||||
|
if (!cap || cap <= 0) return false;
|
||||||
|
const row = await one<{ spent: number }>(
|
||||||
|
`SELECT COALESCE(sum(cost),0)::float AS spent FROM items
|
||||||
|
WHERE cost IS NOT NULL AND created_at >= date_trunc('month', now())`);
|
||||||
|
return (row?.spent || 0) >= cap;
|
||||||
|
}
|
||||||
|
|
||||||
async function handle(job: GenerateJob): Promise<void> {
|
async function handle(job: GenerateJob): Promise<void> {
|
||||||
// Angehaltene/abgebrochene Aufträge nicht verarbeiten.
|
// Angehaltene/abgebrochene Aufträge nicht verarbeiten.
|
||||||
const j = await one<{ status: string }>('SELECT status FROM jobs WHERE id=$1', [job.jobId]);
|
const j = await one<{ status: string }>('SELECT status FROM jobs WHERE id=$1', [job.jobId]);
|
||||||
@@ -17,6 +28,15 @@ async function handle(job: GenerateJob): Promise<void> {
|
|||||||
await query(`UPDATE items SET status='queued' WHERE id=$1 AND status='running'`, [job.itemId]);
|
await query(`UPDATE items SET status='queued' WHERE id=$1 AND status='running'`, [job.itemId]);
|
||||||
return;
|
return;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// Kostendeckel: bei Erreichen Auftrag anhalten, Position zurück in die Schlange.
|
||||||
|
if (await budgetExceeded()) {
|
||||||
|
await query(`UPDATE jobs SET status='paused' WHERE id=$1`, [job.jobId]);
|
||||||
|
await query(`UPDATE items SET status='queued', error_message='Monatsbudget erreicht' WHERE id=$1`, [job.itemId]);
|
||||||
|
console.error('[worker] Monatsbudget erreicht — Auftrag angehalten', job.jobId);
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
|
||||||
await query(`UPDATE jobs SET status='running' WHERE id=$1 AND status='queued'`, [job.jobId]);
|
await query(`UPDATE jobs SET status='running' WHERE id=$1 AND status='queued'`, [job.jobId]);
|
||||||
|
|
||||||
try {
|
try {
|
||||||
@@ -52,8 +72,11 @@ async function maybeComplete(jobId: string): Promise<void> {
|
|||||||
(SELECT count(*) FROM items WHERE job_id=$1 AND status IN ('queued','running'))::int AS open
|
(SELECT count(*) FROM items WHERE job_id=$1 AND status IN ('queued','running'))::int AS open
|
||||||
FROM jobs WHERE id=$1`, [jobId]);
|
FROM jobs WHERE id=$1`, [jobId]);
|
||||||
if (row && row.open === 0) {
|
if (row && row.open === 0) {
|
||||||
await query(`UPDATE jobs SET status='done', finished_at=now() WHERE id=$1 AND status<>'cancelled'`,
|
// Idempotent: nur der erste Übergang nach 'done' löst Auslieferung + Rückmeldung aus.
|
||||||
[jobId]);
|
const done = await query(
|
||||||
|
`UPDATE jobs SET status='done', finished_at=now()
|
||||||
|
WHERE id=$1 AND status NOT IN ('done','cancelled') RETURNING id`, [jobId]);
|
||||||
|
if (done.length === 0) return;
|
||||||
// Picdrop-Auslieferung für alle offenen Positionen anstoßen.
|
// Picdrop-Auslieferung für alle offenen Positionen anstoßen.
|
||||||
try { await deliverPendingForJob(jobId); } catch (e) { console.error('[worker] Auslieferung:', e); }
|
try { await deliverPendingForJob(jobId); } catch (e) { console.error('[worker] Auslieferung:', e); }
|
||||||
// Telegram-Rückmeldung, wenn der Auftrag aus einem Chat kam.
|
// Telegram-Rückmeldung, wenn der Auftrag aus einem Chat kam.
|
||||||
|
|||||||
Reference in New Issue
Block a user