diff --git a/src/worker.ts b/src/worker.ts index 92d6a51..5d288bd 100644 --- a/src/worker.ts +++ b/src/worker.ts @@ -10,6 +10,17 @@ export async function startImageWorker(): Promise { await startWorker(concurrency, handle); } +/** Monatliches Budget erreicht? (Summe usage.cost im laufenden Kalendermonat.) */ +async function budgetExceeded(): Promise { + 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 { // Angehaltene/abgebrochene Aufträge nicht verarbeiten. 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 { await query(`UPDATE items SET status='queued' WHERE id=$1 AND status='running'`, [job.itemId]); 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]); try { @@ -52,8 +72,11 @@ async function maybeComplete(jobId: string): Promise { (SELECT count(*) FROM items WHERE job_id=$1 AND status IN ('queued','running'))::int AS open FROM jobs WHERE id=$1`, [jobId]); if (row && row.open === 0) { - await query(`UPDATE jobs SET status='done', finished_at=now() WHERE id=$1 AND status<>'cancelled'`, - [jobId]); + // Idempotent: nur der erste Übergang nach 'done' löst Auslieferung + Rückmeldung aus. + 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. try { await deliverPendingForJob(jobId); } catch (e) { console.error('[worker] Auslieferung:', e); } // Telegram-Rückmeldung, wenn der Auftrag aus einem Chat kam.