import { pool } from '@kc/platform/db'; import { notifyAdmin } from './discord.js'; import { JobFailure } from '@kc/platform/jobs'; import { executeAction, executeChild, provisionOrder, syncInstance, type ChildPayload, type ExecutePayload } from '@kc/connectors'; interface JobMeta { id: string; correlationId: string; attempt: number; maxAttempts: number } type Handler = (payload: any, job: JobMeta) => Promise; /** Job-Handler nach Typ. Neue Module registrieren hier ihre Aufträge. */ const handlers: Record = { 'discord.notify': async (p) => { const text = p.event === 'customer.created' ? `Neuer Kunde angelegt: ${p.customerNumber}` : p.event === 'backup.failed' ? `⚠ Backup fehlgeschlagen: ${p.detail ?? ''}`.trim() : p.event === 'backup.stale' ? `⚠ Kein aktuelles Backup: ${p.detail ?? ''}`.trim() : p.event === 'restoretest.failed' ? `⚠ Wiederherstellungstest fehlgeschlagen: ${p.detail ?? ''}`.trim() : p.event === 'ticket.created' ? `🎫 Neues Ticket ${p.number}: ${p.subject ?? ''}`.trim() : p.event === 'ticket.message' ? `🎫 Neue Nachricht zu Ticket ${p.number}` : p.event === 'invoice.issued' ? `🧾 Rechnung ${p.number} ausgestellt (${(Number(p.gross ?? 0) / 100).toLocaleString('de-DE', { style: 'currency', currency: 'EUR' })})` : null; // keine personenbezogenen Daten if (!text) throw new Error(`Unbekanntes Ereignis: ${p.event}`); const r = await notifyAdmin(text); if (r === 'skipped') return 'übersprungen: Discord nicht konfiguriert'; }, }; handlers['connector.sync'] = async (p, j) => { const r = await syncInstance(p.instanceId, j.correlationId); return r.ok ? `${r.count} Ressourcen abgeglichen` : `Abgleich fehlgeschlagen: ${r.error}`; // Ausfall ist ein Zustand der Instanz, kein Jobfehler (nächster Lauf folgt planmäßig) }; handlers['order.provision'] = (p: { orderId: string }, j) => provisionOrder(p.orderId, j); handlers['connector.child'] = (p: ChildPayload, j) => executeChild(p, j.id, j.correlationId, () => j.attempt >= j.maxAttempts); handlers['connector.execute'] = (p: ExecutePayload, j) => executeAction(p, j.id, j.correlationId); const backoff = (attempt: number) => Math.min(3600, 30 * 2 ** attempt); // Sekunden, exponentiell /** Holt einen fälligen Auftrag (SKIP LOCKED) und führt ihn aus. Liefert false, wenn nichts zu tun war. */ export async function runOnce(log: (m: string) => void): Promise { const c = await pool.getConnection(); let job: any; try { await c.beginTransaction(); const [rows] = await c.query("SELECT * FROM jobs WHERE status IN ('scheduled','retrying') AND run_at <= UTC_TIMESTAMP(3) ORDER BY run_at LIMIT 1 FOR UPDATE SKIP LOCKED") as any; job = rows[0]; if (!job) { await c.rollback(); return false; } await c.query("UPDATE jobs SET status='running', locked_at=UTC_TIMESTAMP(3), attempts=attempts+1 WHERE id=?", [job.id]); await c.commit(); } catch (e) { await c.rollback(); throw e; } finally { c.release(); } const attempt = job.attempts + 1; try { const h = handlers[job.type]; if (!h) throw new Error(`Kein Handler für ${job.type}`); const note = await h(typeof job.payload === 'string' ? JSON.parse(job.payload) : job.payload, { id: job.id, correlationId: job.correlation_id ?? job.id, attempt, maxAttempts: job.max_attempts }); await pool.query("UPDATE jobs SET status='succeeded', locked_at=NULL, last_error=? WHERE id=?", [note ?? null, job.id]); } catch (e) { const msg = String((e as Error).message).slice(0, 500); const jf = e instanceof JobFailure ? e : null; const noRetry = jf ? !jf.retry : false; const dead = noRetry || attempt >= job.max_attempts; const final = jf?.finalStatus ?? 'needs_review'; await pool.query("UPDATE jobs SET status=?, locked_at=NULL, last_error=?, run_at=DATE_ADD(UTC_TIMESTAMP(3), INTERVAL ? SECOND) WHERE id=?", [dead ? final : 'retrying', msg, jf?.delaySec ?? backoff(attempt), job.id]); log(`Job ${job.id} (${job.type}) Versuch ${attempt} fehlgeschlagen: ${msg}${dead ? ` → ${final}` : ' → wird wiederholt'}`); } return true; } /** Nach Absturz hängengebliebene Aufträge wieder freigeben. Nicht-idempotente Typen kämen hier in 'needs_review'. */ export async function recoverStale(log: (m: string) => void): Promise { // Destruktive Aufträge werden nach einem Absturz NICHT automatisch wiederholt, sondern zur Prüfung markiert. await pool.query("UPDATE jobs SET status='needs_review', locked_at=NULL, last_error='Unterbrochen während der Ausführung; Ergebnis beim Provider prüfen' WHERE status='running' AND (JSON_VALUE(payload, '$.destructive') IN ('1','true') OR JSON_VALUE(payload, '$.action') = 'terminate') AND locked_at < DATE_SUB(UTC_TIMESTAMP(3), INTERVAL 5 MINUTE)"); const [r] = await pool.query("UPDATE jobs SET status='retrying', locked_at=NULL, last_error='Worker-Neustart während der Ausführung' WHERE status='running' AND locked_at < DATE_SUB(UTC_TIMESTAMP(3), INTERVAL 5 MINUTE)") as any; if (r.affectedRows) log(`${r.affectedRows} hängende Aufträge wieder eingeplant`); }