Files
craftvia/scripts/craftvia-worker.ts
T
msolarczekandClaude Opus 5 b0aedb5d23 L10b Betrieb & Aufräumen: Lotse-Betrieb – Aufbewahrung KI-Protokoll und Token-Kontingent
Aufräumpunkt k (Spec §31):
- Aufbewahrung: services/lotse/retention.ts leert input/output und createdById von
  AiGeneration-Einträgen älter als AI_GENERATION_RETENTION_DAYS (Default 180), Metadaten bleiben,
  Audit je Mandant. Queue/Processor ai-retention, täglicher BullMQ-Job-Scheduler beim Start des
  craftvia-worker.
- Kontingent: services/lotse/budget.ts (Tokens ein+aus je Kalendermonat, TenantSettings-Wert vor
  Env AI_MONTHLY_TOKEN_LIMIT, 0 = unbegrenzt). Lotse-Entwurf und Sprachnotiz-Zusammenfassung
  → blocked budget_exceeded mit Klartext; Import-Extraktion fällt auf manuelle Erfassung zurück
  (Hinweis ai_budget_exceeded). /settings/lotse: Kontingent setzen, Verbrauch anzeigen.
- scripts/test-betrieb-audit.ts: Audit nach Commit/Rollback/verschachtelt, Merge atomar und in
  äußerer Transaktion, Audit „read", Aufbewahrung (Frist, Metadaten, Idempotenz, Mandant B),
  Kontingent (Mandant/Env/Vormonat/unbegrenzt, Rollen, Audit).

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
2026-09-14 18:19:19 +02:00

45 lines
1.7 KiB
TypeScript

import "dotenv/config";
import { Worker } from "bullmq";
import { JOB_QUEUES, workerConnection, closeJobQueues, scheduleRecurringJobs, type JobPayload } from "../src/server/jobs/queues";
import { PROCESSORS } from "../src/server/jobs/processors";
/** Craftvia background worker: `npm run worker:craftvia`. One BullMQ worker per registered queue. */
async function main() {
const connection = workerConnection();
if (!connection) {
console.error("[worker] REDIS_URL is not set — nothing to do.");
process.exit(1);
}
const workers: Worker<JobPayload>[] = [];
for (const name of Object.values(JOB_QUEUES)) {
const load = PROCESSORS[name];
if (!load) {
console.warn(`[worker] no processor registered for ${name} — skipped`);
continue;
}
const processor = await load();
const concurrency = name === JOB_QUEUES.reportPdf ? 2 : 4;
const w = new Worker<JobPayload>(name, async (job) => processor(job.data), { connection, concurrency });
w.on("failed", (job, err) => console.error(`[worker] ${name} job ${job?.id} failed:`, err.message));
workers.push(w);
console.info(`[worker] listening on ${name}`);
}
// L10b: recurring jobs (AI log retention); a scheduling failure must not stop the queue workers
await scheduleRecurringJobs(connection).then(
() => console.info("[worker] recurring jobs scheduled"),
(err) => console.error("[worker] scheduling recurring jobs failed:", (err as Error).message),
);
const shutdown = async () => {
await Promise.all(workers.map((w) => w.close()));
await closeJobQueues();
process.exit(0);
};
process.on("SIGTERM", shutdown);
process.on("SIGINT", shutdown);
}
main().catch((err) => {
console.error(err);
process.exit(1);
});