diff --git a/apps/worker/src/jobs.ts b/apps/worker/src/jobs.ts index 2d0d8ab..c542d9a 100644 --- a/apps/worker/src/jobs.ts +++ b/apps/worker/src/jobs.ts @@ -1,9 +1,15 @@ +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; /** Job-Handler nach Typ. Neue Module registrieren hier ihre Aufträge. */ @@ -39,41 +45,53 @@ handlers['connector.child'] = (p: ChildPayload, j) => executeChild(p, j.id, j.co 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. */ +/** 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 { const c = await pool.getConnection(); - let job: any; + 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; } - await c.query("UPDATE jobs SET status='running', locked_at=UTC_TIMESTAMP(3), attempts=attempts+1 WHERE id=?", [job.id]); + 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 }); - await pool.query("UPDATE jobs SET status='succeeded', locked_at=NULL, last_error=? WHERE id=?", [note ?? null, job.id]); + 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'; - 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'}`); - } + 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'. */ +/** 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 { // 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; + 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`); } diff --git a/migrations/032_jobs_lease.sql b/migrations/032_jobs_lease.sql new file mode 100644 index 0000000..0b2aa49 --- /dev/null +++ b/migrations/032_jobs_lease.sql @@ -0,0 +1,6 @@ +-- Nightbot-Befund #99/#16: recoverStale gab "running"-Jobs nach starren 5 Minuten wieder frei, ohne zu prüfen, +-- ob der ursprüngliche Worker noch aktiv daran arbeitet. Ein legitim länger laufender Job konnte dadurch von +-- einem zweiten Worker parallel erneut gestartet werden (Doppelausführung). +ALTER TABLE jobs + ADD COLUMN lease_id CHAR(36) NULL, -- pro Übernahme neu vergeben; nur der aktuelle Lease darf das Ergebnis schreiben + ADD COLUMN locked_by VARCHAR(100) NULL; -- Worker-Kennung (Host:PID:Zufall) für Diagnose, kein Sicherheitsmerkmal