Stand vor Einführung des Nacht-Agenten
This commit is contained in:
commit
4763548bfb
168 changed files with 12726 additions and 0 deletions
26
apps/worker/package.json
Normal file
26
apps/worker/package.json
Normal file
|
|
@ -0,0 +1,26 @@
|
|||
{
|
||||
"name": "@kc/worker",
|
||||
"private": true,
|
||||
"version": "0.1.0",
|
||||
"type": "module",
|
||||
"scripts": {
|
||||
"dev": "tsx watch src/index.ts",
|
||||
"start": "node dist/index.js",
|
||||
"build": "tsc -p tsconfig.json",
|
||||
"typecheck": "tsc -p tsconfig.json --noEmit",
|
||||
"test": "vitest run --passWithNoTests"
|
||||
},
|
||||
"dependencies": {
|
||||
"@kc/connectors": "workspace:*",
|
||||
"@kc/platform": "workspace:*",
|
||||
"discord.js": "^14.27.0",
|
||||
"mysql2": "^3.24.4",
|
||||
"zod": "^4.6.5"
|
||||
},
|
||||
"devDependencies": {
|
||||
"@types/node": "^26.6.3",
|
||||
"tsx": "^4.23.15",
|
||||
"typescript": "^7.0.2",
|
||||
"vitest": "^5.0.2"
|
||||
}
|
||||
}
|
||||
51
apps/worker/src/discord.ts
Normal file
51
apps/worker/src/discord.ts
Normal file
|
|
@ -0,0 +1,51 @@
|
|||
import { Client, GatewayIntentBits, REST, Routes, SlashCommandBuilder, type ChatInputCommandInteraction } from 'discord.js';
|
||||
import { env } from './env.js';
|
||||
import { pool } from '@kc/platform/db';
|
||||
|
||||
let client: Client | null = null;
|
||||
export const discordEnabled = () => !!env.discord.token;
|
||||
export const discordReady = () => !!client?.isReady();
|
||||
|
||||
const commands = [
|
||||
new SlashCommandBuilder().setName('kc-status').setDescription('Zustand des Kundencenters (Aufträge, Datenbank)'),
|
||||
new SlashCommandBuilder().setName('kc-kunde').setDescription('Kunde suchen (nur Kundennummer, Name, Status)')
|
||||
.addStringOption((o) => o.setName('suche').setDescription('Kundennummer oder Name').setRequired(true).setMaxLength(100)),
|
||||
].map((c) => c.toJSON());
|
||||
|
||||
async function handle(i: ChatInputCommandInteraction): Promise<void> {
|
||||
// Alle Antworten ephemeral: nie vertrauliche Daten in Kanälen sichtbar machen.
|
||||
if (!env.discord.staffUserIds.includes(i.user.id)) { await i.reply({ content: 'Keine Berechtigung.', ephemeral: true }); return; }
|
||||
if (i.commandName === 'kc-status') {
|
||||
const [rows] = await pool.query('SELECT status, COUNT(*) n FROM jobs GROUP BY status') as any;
|
||||
const jobs = rows.length ? rows.map((r: any) => `${r.status}: ${r.n}`).join(', ') : 'keine';
|
||||
await i.reply({ content: `Datenbank: ok\nAufträge: ${jobs}`, ephemeral: true });
|
||||
} else if (i.commandName === 'kc-kunde') {
|
||||
const q = i.options.getString('suche', true);
|
||||
const [rows] = await pool.query("SELECT customer_number, name, status FROM organizations WHERE name LIKE CONCAT('%', ?, '%') OR customer_number = ? LIMIT 10", [q, q]) as any;
|
||||
await i.reply({ content: rows.length ? rows.map((r: any) => `${r.customer_number} · ${r.name} · ${r.status}`).join('\n') : 'Nichts gefunden.', ephemeral: true });
|
||||
}
|
||||
}
|
||||
|
||||
export async function startDiscord(log: (m: string) => void): Promise<void> {
|
||||
if (!env.discord.token) { log('Discord-Bot deaktiviert (DISCORD_BOT_TOKEN nicht gesetzt)'); return; }
|
||||
client = new Client({ intents: [GatewayIntentBits.Guilds] });
|
||||
client.on('interactionCreate', (i) => {
|
||||
if (i.isChatInputCommand()) handle(i).catch((e) => { log(`Befehl fehlgeschlagen: ${(e as Error).message}`); i.replied || i.deferred ? undefined : i.reply({ content: 'Fehler.', ephemeral: true }).catch(() => undefined); });
|
||||
});
|
||||
client.once('ready', async (c) => {
|
||||
log(`Discord verbunden als ${c.user.tag}`);
|
||||
if (env.discord.guildId) await new REST().setToken(env.discord.token!).put(Routes.applicationGuildCommands(c.user.id, env.discord.guildId), { body: commands });
|
||||
else log('DISCORD_GUILD_ID fehlt: Slash-Befehle nicht registriert');
|
||||
});
|
||||
await client.login(env.discord.token);
|
||||
}
|
||||
|
||||
/** Sendet eine Nachricht in den Admin-Kanal. Wirft bei Fehlern, damit der Job wiederholt wird. */
|
||||
export async function notifyAdmin(text: string): Promise<'sent' | 'skipped'> {
|
||||
if (!client || !env.discord.adminChannelId) return 'skipped';
|
||||
const ch = await client.channels.fetch(env.discord.adminChannelId);
|
||||
if (!ch || !ch.isTextBased() || !('send' in ch)) throw new Error('Admin-Kanal nicht gefunden oder kein Textkanal');
|
||||
await ch.send({ content: text, allowedMentions: { parse: [] } });
|
||||
return 'sent';
|
||||
}
|
||||
export async function stopDiscord(): Promise<void> { await client?.destroy(); }
|
||||
13
apps/worker/src/env.ts
Normal file
13
apps/worker/src/env.ts
Normal file
|
|
@ -0,0 +1,13 @@
|
|||
import { config } from '@kc/platform/config';
|
||||
|
||||
/** Worker-spezifische Optionen; DB/Secrets kommen aus @kc/platform. */
|
||||
export const env = {
|
||||
db: config.db,
|
||||
discord: {
|
||||
token: process.env.DISCORD_BOT_TOKEN || null,
|
||||
guildId: process.env.DISCORD_GUILD_ID || null,
|
||||
adminChannelId: process.env.DISCORD_ADMIN_CHANNEL_ID || null,
|
||||
staffUserIds: (process.env.DISCORD_STAFF_USER_IDS ?? '').split(',').map((s) => s.trim()).filter(Boolean),
|
||||
},
|
||||
healthPort: Number(process.env.KC_WORKER_HEALTH_PORT ?? 4102),
|
||||
};
|
||||
48
apps/worker/src/index.ts
Normal file
48
apps/worker/src/index.ts
Normal file
|
|
@ -0,0 +1,48 @@
|
|||
import http from 'node:http';
|
||||
import { env } from './env.js';
|
||||
import { pool } from '@kc/platform/db';
|
||||
import { discordEnabled, discordReady, startDiscord, stopDiscord } from './discord.js';
|
||||
import { recoverStale, runOnce } from './jobs.js';
|
||||
import { enqueue } from '@kc/platform/jobs';
|
||||
import { processContractLifecycle, scheduleDueSyncs } from '@kc/connectors';
|
||||
|
||||
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);
|
||||
startDiscord(log).catch((e) => log(`Discord-Start fehlgeschlagen: ${(e as Error).message}`));
|
||||
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);
|
||||
|
||||
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)); }
|
||||
}
|
||||
66
apps/worker/src/jobs.ts
Normal file
66
apps/worker/src/jobs.ts
Normal file
|
|
@ -0,0 +1,66 @@
|
|||
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<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() : 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<boolean> {
|
||||
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<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, 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`);
|
||||
}
|
||||
1
apps/worker/tsconfig.json
Normal file
1
apps/worker/tsconfig.json
Normal file
|
|
@ -0,0 +1 @@
|
|||
{ "extends": "../../tsconfig.base.json", "compilerOptions": { "rootDir": "src", "outDir": "dist" }, "include": ["src"] }
|
||||
Loading…
Add table
Add a link
Reference in a new issue