Files
craftvia/scripts/craftvia-worker.ts
T
msolarczekandClaude Opus 5 213bcca3a1 L13 Planung: Geocoding der Objekte über OpenStreetMap Nominatim
Provider-Interface (nominatim|none), strukturierte Suche mit User-Agent und Accept-Language, Drosselung 1/s je Prozess, Job geocode-site (Worker: Concurrency 1 + Limiter), Cache am Objekt, manuelle Koordinaten bleiben, Auslöser Anlage/Adressänderung/Import-Bestätigung nur per Queue, Backfill-Skript, Env-Beispiele. Tests mit Fake-Provider, Nominatim nie aufgerufen.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
2026-09-15 10:13:26 +02:00

46 lines
1.9 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 geocode = name === JOB_QUEUES.geocodeSite; // L13: OSM Nominatim policy — max. 1 request/s
const concurrency = name === JOB_QUEUES.reportPdf ? 2 : geocode ? 1 : 4;
const w = new Worker<JobPayload>(name, async (job) => processor(job.data), { connection, concurrency, ...(geocode ? { limiter: { max: 1, duration: 1_000 } } : {}) });
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);
});