kundencenter/apps/worker/src/jobs.ts
Kundencenter a9893b6d32 fix(worker): Lease/Heartbeat statt starrer 5-Minuten-Grenze bei Jobwiederaufnahme
Nightbot-Befund #99/#16 (P1): recoverStale gab "running"-Jobs nach starren
fünf Minuten wieder frei, ohne zu prüfen, ob der ursprüngliche Worker noch
aktiv daran arbeitet. Ein legitim länger laufender Job (z. B. ein
langsamer Connector-Abgleich) konnte dadurch von einem zweiten Worker
parallel erneut gestartet werden – Doppelausführung, z. B. doppelte
Provisionierung oder doppelter Mailversand.

- Migration 032: jobs.lease_id (pro Übernahme neu vergeben) und
  jobs.locked_by (Worker-Kennung, nur Diagnose).
- runOnce vergibt beim Übernehmen einen frischen Lease und hält locked_at
  per Heartbeat alle 60 s aktuell, solange der Handler läuft. Jedes
  Abschluss-UPDATE (Erfolg wie Fehler) ist an genau diesen Lease gebunden
  (WHERE lease_id = ?); hat recoverStale die Zeile inzwischen doch
  freigegeben, verpufft ein verspätetes Ergebnis wirkungslos statt den
  neuen Versuch zu überschreiben.
- recoverStale reißt jetzt nur noch Jobs an sich, deren Heartbeat
  tatsächlich ausgeblieben ist (locked_at älter als 5 Minuten), nicht
  mehr solche, die einfach nur lange laufen. Löscht lease_id beim
  Freigeben mit.

Verifiziert: (a) SQL-Ebene direkt durchexerziert – frischer Heartbeat
verhindert das Anreißen, ausgebliebener Heartbeat löst es aus, ein
verspätetes Abschluss-UPDATE mit altem Lease betrifft 0 Zeilen; (b) echter
Code über runOnce/recoverStale gegen die Testdatenbank – Job durchläuft
Übernahme, Fehlschlag, lease-gebundenes Abschluss-UPDATE korrekt. Worker
neu gestartet, läuft fehlerfrei, 1437 bestehende Jobs weiterhin 'succeeded'.

Co-Authored-By: Claude Sonnet 5 <noreply@anthropic.com>
2026-09-29 11:41:43 +02:00

97 lines
7.9 KiB
TypeScript
Raw Blame History

This file contains ambiguous Unicode characters

This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.

import { randomUUID } from 'node:crypto';
import { hostname } from 'node:os';
import { pool } from '@kc/platform/db';
import { notifyAdmin, postTicketMessage } from './discord.js';
import { JobFailure } from '@kc/platform/jobs';
import { sendTemplateMail } from '@kc/platform/mail';
import { executeAction, executeChild, provisionOrder, syncInstance, type ChildPayload, type ExecutePayload } from '@kc/connectors';
// Kennung dieses Worker-Prozesses (nur zur Diagnose in locked_by, kein Sicherheitsmerkmal).
const WORKER_ID = `${hostname()}:${process.pid}:${randomUUID().slice(0, 8)}`;
const HEARTBEAT_MS = 60_000; // deutlich unter der 5-Minuten-Schwelle in recoverStale
interface JobMeta { id: string; correlationId: string; attempt: number; maxAttempts: number }
type Handler = (payload: any, job: JobMeta) => Promise<string | void>;
/** Job-Handler nach Typ. Neue Module registrieren hier ihre Aufträge. */
const handlers: Record<string, Handler> = {
'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';
},
'mail.template': async (p: { to: string; key: string; vars: Record<string, string> }) => {
const status = await sendTemplateMail(p.to, p.key, p.vars);
if (status === 'failed') throw new Error('Mailversand fehlgeschlagen (siehe Protokoll unter Einstellungen > E-Mail)');
if (status === 'no_template') throw new Error(`Unbekannte Mailvorlage: ${p.key}`);
return status === 'not_configured' ? 'übersprungen: SMTP nicht konfiguriert' : status === 'disabled' ? 'übersprungen: Vorlage deaktiviert' : 'gesendet';
},
'discord.ticket_sync': async (p: { ticketId: string; body: string; authorLabel: string }) => {
await postTicketMessage(p.ticketId, p.body, p.authorLabel);
},
};
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.
* Bindet Ausführung an einen frischen Lease: solange der Handler läuft, wird locked_at per Heartbeat aktuell
* gehalten (recoverStale reißt einen noch aktiv laufenden Job dadurch nicht mehr an sich), und jedes
* Abschluss-UPDATE ist an genau diesen Lease gebunden – hat recoverStale den Job inzwischen doch freigegeben
* (z. B. weil dieser Prozess abgestürzt und der Heartbeat ausgeblieben ist), verpufft ein verspätetes Ergebnis
* dieses Prozesses wirkungslos, statt den neuen Versuch eines anderen Workers zu überschreiben. */
export async function runOnce(log: (m: string) => void): Promise<boolean> {
const c = await pool.getConnection();
let job: any; let lease: string;
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; }
lease = randomUUID();
await c.query("UPDATE jobs SET status='running', locked_at=UTC_TIMESTAMP(3), lease_id=?, locked_by=?, attempts=attempts+1 WHERE id=?", [lease, WORKER_ID, job.id]);
await c.commit();
} catch (e) { await c.rollback(); throw e; } finally { c.release(); }
const attempt = job.attempts + 1;
const heartbeat = setInterval(() => { pool.query("UPDATE jobs SET locked_at=UTC_TIMESTAMP(3) WHERE id=? AND lease_id=?", [job.id, lease]).catch(() => undefined); }, HEARTBEAT_MS);
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 });
const [r] = await pool.query("UPDATE jobs SET status='succeeded', locked_at=NULL, lease_id=NULL, locked_by=NULL, last_error=? WHERE id=? AND lease_id=?", [note ?? null, job.id, lease]) as any;
if (!r.affectedRows) log(`Job ${job.id} (${job.type}): Lease inzwischen einem anderen Worker zugewiesen – Ergebnis (erfolgreich) verworfen, um dessen Versuch nicht zu überschreiben`);
} 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';
const [r] = await pool.query("UPDATE jobs SET status=?, locked_at=NULL, lease_id=NULL, locked_by=NULL, last_error=?, run_at=DATE_ADD(UTC_TIMESTAMP(3), INTERVAL ? SECOND) WHERE id=? AND lease_id=?", [dead ? final : 'retrying', msg, jf?.delaySec ?? backoff(attempt), job.id, lease]) as any;
if (r.affectedRows) log(`Job ${job.id} (${job.type}) Versuch ${attempt} fehlgeschlagen: ${msg}${dead ? ` → ${final}` : ' → wird wiederholt'}`);
else log(`Job ${job.id} (${job.type}): Lease inzwischen einem anderen Worker zugewiesen – Ergebnis (Fehler) verworfen, um dessen Versuch nicht zu überschreiben`);
} finally { clearInterval(heartbeat); }
return true;
}
/** Nach Absturz hängengebliebene Aufträge wieder freigeben. Nicht-idempotente Typen kämen hier in 'needs_review'.
* locked_at ist der letzte Heartbeat, nicht der Übernahmezeitpunkt – ein noch aktiv laufender Job (Heartbeat alle
* 60 s) wird dadurch nicht angerissen, egal wie lange er tatsächlich läuft. lease_id wird beim Freigeben gelöscht,
* damit ein verspätetes Abschluss-UPDATE des ursprünglichen (mutmaßlich abgestürzten) Workers wirkungslos bleibt. */
export async function recoverStale(log: (m: string) => void): Promise<void> {
// 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, lease_id=NULL, locked_by=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, lease_id=NULL, locked_by=NULL, last_error='Worker-Neustart oder ausgebliebener Heartbeat 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`);
}