import http from 'node:http'; import { env } from './env.js'; import { pool } from '@kc/platform/db'; import { discordEnabled, discordReady, syncDiscord, stopDiscord } from './discord.js'; import { recoverStale, runOnce } from './jobs.js'; import { enqueue } from '@kc/platform/jobs'; import { processContractLifecycle, scheduleDueSyncs } from '@kc/connectors'; import { checkMailbox } from './mailbox.js'; const log = (m: string) => console.log(JSON.stringify({ t: new Date().toISOString(), svc: 'worker', msg: m })); let lastLoop = Date.now(); let stopping = false; // Health: Liveness (Schleife lebt) + Readiness (DB), Discord-Status nur informativ http.createServer(async (req, res) => { const alive = Date.now() - lastLoop < 60_000; let db = false; try { await pool.query('SELECT 1'); db = true; } catch { /* db down */ } const ok = alive && db; res.writeHead(req.url === '/health' ? (alive ? 200 : 503) : (ok ? 200 : 503), { 'content-type': 'application/json' }); res.end(JSON.stringify({ status: ok ? 'ok' : 'degraded', loop: alive, db, discord: discordEnabled() ? (discordReady() ? 'connected' : 'connecting') : 'disabled' })); }).listen(env.healthPort, '127.0.0.1'); await recoverStale(log); void syncDiscord(log); setInterval(() => void syncDiscord(log), 60_000); // erkennt geänderte Einstellungen (Einstellungen > Discord) und verbindet bei Bedarf neu setInterval(() => recoverStale(log).catch(() => undefined), 60_000); // Regelmäßiger Abgleich: fällige Connector-Instanzen als Aufträge einplanen (idempotent pro Zeitfenster) const schedule = () => scheduleDueSyncs((t, p, o) => enqueue(t, p, o)).catch((e) => log(`Planung fehlgeschlagen: ${(e as Error).message}`)); setInterval(schedule, 30_000); void schedule(); // Verträge: Kündigungen wirksam machen, verlängern, auslaufen lassen (idempotent) const lifecycle = () => processContractLifecycle(new Date(), (t, p, o) => enqueue(t, p, o)).then((r) => { if (r.ended || r.renewed) log(`Verträge: ${r.ended} beendet, ${r.renewed} verlängert`); }).catch((e) => log(`Vertragslauf fehlgeschlagen: ${(e as Error).message}`)); setInterval(lifecycle, 60_000); void lifecycle(); // Backup-Frische: fehlt ein erfolgreiches Backup seit >26 h, wird gemeldet (höchstens alle 12 h, ohne Nutzdaten) import { readFile } from 'node:fs/promises'; const checkBackup = async () => { try { const st = JSON.parse(await readFile(process.env.BACKUP_STATUS_FILE ?? '/var/lib/kundencenter/backup-status.json', 'utf8')); const age = st.lastRun?.at ? (Date.now() - new Date(st.lastRun.at).getTime()) / 3600000 : Infinity; if (age > 26) await enqueue('discord.notify', { event: 'backup.stale', detail: Number.isFinite(age) ? `letzter Lauf vor ${Math.round(age)} Stunden` : 'noch nie gelaufen' }, { idempotencyKey: `backup.stale:${Math.floor(Date.now() / (12 * 3600000))}` }); } catch { /* Backup nicht eingerichtet: keine Meldung */ } }; setInterval(() => void checkBackup(), 3600_000); // Ticket-Posteingang: eingehende Mails abrufen und zuordnen (siehe mailbox.ts) setInterval(() => void checkMailbox(log).catch((e) => log(`IMAP-Lauf fehlgeschlagen: ${(e as Error).message}`)), 120_000); void checkMailbox(log).catch((e) => log(`IMAP-Lauf fehlgeschlagen: ${(e as Error).message}`)); for (const sig of ['SIGTERM', 'SIGINT'] as const) process.on(sig, async () => { stopping = true; await stopDiscord(); await pool.end(); process.exit(0); }); log('Worker gestartet'); while (!stopping) { try { lastLoop = Date.now(); const worked = await runOnce(log); if (!worked) await new Promise((r) => setTimeout(r, 2000)); } catch (e) { log(`Schleifenfehler: ${(e as Error).message}`); await new Promise((r) => setTimeout(r, 5000)); } }