taskmanager/scripts/worker.ts

374 lines
11 KiB
TypeScript

import "dotenv/config";
import { prisma as db } from "../src/lib/prisma";
import { loadPrincipal } from "../src/lib/principal";
import { allowed } from "../src/lib/rbac";
import { removeFile, storageRoot, storagePath } from "../src/lib/storage";
import { sendMail } from "../src/lib/mail";
import { createHmac } from "node:crypto";
import { lookup } from "node:dns/promises";
import { request as httpsRequest } from "node:https";
import { readdir, stat, writeFile } from "node:fs/promises";
import { setTimeout as delay } from "node:timers/promises";
import { isIP } from "node:net";
import type { Prisma } from "@prisma/client";
function publicAddress(ip: string) {
if (isIP(ip) === 4) {
const [a, b, c] = ip.split(".").map(Number);
return !(
[0, 10, 127].includes(a) ||
a >= 224 ||
(a === 169 && b === 254) ||
(a === 172 && b >= 16 && b <= 31) ||
(a === 192 && b === 168) ||
(a === 192 && b === 0) ||
(a === 198 && (b === 18 || b === 19 || (b === 51 && c === 100))) ||
(a === 203 && b === 0 && c === 113) ||
(a === 100 && b >= 64 && b <= 127)
);
}
// Restrict IPv6 to global unicast, excluding mapped IPv4 and local ranges.
return (
/^[23][0-9a-f]{3}:/i.test(ip) && !ip.toLowerCase().startsWith("2001:db8:")
);
}
async function deliver(urlString: string, secret: string, event: object) {
const url = new URL(urlString);
if (
url.protocol !== "https:" ||
!(process.env.WEBHOOK_ALLOWED_HOSTS || "")
.split(",")
.includes(url.hostname) ||
url.username ||
url.password ||
(url.port && url.port !== "443")
)
throw new Error("Webhook nicht freigegeben");
const address = await lookup(url.hostname);
if (!publicAddress(address.address))
throw new Error("Webhook muss öffentliche IP verwenden");
const body = JSON.stringify(event);
const timestamp = String(Math.floor(Date.now() / 1000));
const signature = createHmac("sha256", secret)
.update(`${timestamp}.${body}`)
.digest("hex");
await new Promise<void>((resolve, reject) => {
const req = httpsRequest(
url,
{
method: "POST",
timeout: 10000,
lookup: (_host, options, callback) =>
callback(
null,
options.all ? [address] : address.address,
address.family,
),
headers: {
"Content-Type": "application/json",
"Content-Length": Buffer.byteLength(body),
"X-Taskmanager-Timestamp": timestamp,
"X-Taskmanager-Signature": signature,
},
},
(res) => {
res.resume();
if (res.statusCode && res.statusCode >= 200 && res.statusCode < 300)
resolve();
else reject(new Error(`HTTP ${res.statusCode}`));
},
);
req.on("timeout", () => req.destroy(new Error("Timeout")));
req.on("error", reject);
req.end(body);
});
}
async function* scanTasks(where: Prisma.TaskWhereInput) {
let cursor: string | undefined;
for (;;) {
const rows = await db.task.findMany({
where,
orderBy: { id: "asc" },
take: 500,
...(cursor ? { cursor: { id: cursor }, skip: 1 } : {}),
});
if (!rows.length) return;
for (const task of rows) yield task;
cursor = rows[rows.length - 1].id;
}
}
export async function tick() {
const now = new Date();
const day = new Intl.DateTimeFormat("en-CA", {
timeZone: process.env.APP_TIMEZONE || "Europe/Berlin",
year: "numeric",
month: "2-digit",
day: "2-digit",
}).format(now);
const due = scanTasks({
archivedAt: null,
status: { not: "ERLEDIGT" },
bisWann: { lte: new Date(day + "T00:00:00Z") },
});
for await (const task of due) {
const candidates = await db.user.findMany({
where: {
active: true,
OR: [
{ id: task.assigneeId || "" },
{ memberships: { some: { groupId: task.owningGroupId } } },
],
},
select: { id: true },
});
for (const candidate of candidates) {
const user = await loadPrincipal(candidate.id);
if (
!user ||
!allowed(user, "tasks.read", task) ||
(candidate.id !== task.assigneeId &&
!allowed(user, "tasks.update", task))
)
continue;
const overdue = task.bisWann.toISOString().slice(0, 10) < day;
if (candidate.id !== task.assigneeId && !overdue) continue;
const key = `${task.id}:${task.version}:${candidate.id}:${day}`;
await db.notification.upsert({
where: { key },
update: {},
create: {
key,
userId: candidate.id,
taskId: task.id,
text: `${overdue ? "Überfällig" : "Heute fällig"}: ${task.was.slice(0, 150)}`,
},
});
}
}
if (
process.env.SMTP_URL &&
process.env.MAIL_FROM &&
process.env.NEXTAUTH_URL
) {
const pending = await db.notification.findMany({
where: { emailed: false },
take: 100,
orderBy: { createdAt: "asc" },
});
for (const item of pending) {
const user = await loadPrincipal(item.userId);
const task = await db.task.findUnique({ where: { id: item.taskId } });
if (
!user ||
!task ||
!allowed(user, "tasks.read", task) ||
task.archivedAt ||
task.status === "ERLEDIGT"
) {
await db.notification.update({
where: { id: item.id },
data: { emailed: true },
});
continue;
}
try {
await sendMail(
user.email,
"Aufgabenplaner: Erinnerung",
`${item.text}\n${process.env.NEXTAUTH_URL}/dashboard/tasks/${task.id}`,
);
await db.notification.update({
where: { id: item.id },
data: { emailed: true },
});
} catch {
console.error("E-Mail-Zustellung fehlgeschlagen", item.id);
}
}
}
const recurring = scanTasks({
status: "ERLEDIGT",
archivedAt: null,
recurrenceDays: { not: null },
completedAt: { not: null },
});
for await (const source of recurring) {
const actor = await loadPrincipal(source.creatorId);
if (
!actor ||
!allowed(actor, "tasks.create", source) ||
(source.assigneeId && !allowed(actor, "tasks.assign", source))
)
continue;
const group = await db.group.findFirst({
where: { id: source.owningGroupId, active: true },
});
if (!group) continue;
if (
source.assigneeId &&
!(await db.userGroup.findFirst({
where: {
groupId: source.owningGroupId,
userId: source.assigneeId,
user: { active: true },
},
}))
)
continue;
const dueDate = new Date(source.completedAt!);
dueDate.setUTCDate(dueDate.getUTCDate() + source.recurrenceDays!);
dueDate.setUTCHours(0, 0, 0, 0);
await db.task.upsert({
where: { recurrenceSourceId: source.id },
update: {},
create: {
recurrenceSourceId: source.id,
wo: source.wo,
was: source.was,
bisWann: dueDate,
creatorId: source.creatorId,
assigneeId: source.assigneeId,
owningGroupId: source.owningGroupId,
priority: source.priority,
tags: source.tags,
project: source.project,
recurrenceDays: source.recurrenceDays,
checklist: Array.isArray(source.checklist)
? source.checklist.map((c) => ({ ...(c as object), done: false }))
: [],
},
});
}
const cleanup = await db.cleanupJob.findMany({
where: { attempts: { lt: 10 } },
take: 100,
});
for (const job of cleanup)
try {
await removeFile(job.storageKey);
await db.cleanupJob.delete({ where: { id: job.id } });
} catch {
await db.cleanupJob.update({
where: { id: job.id },
data: { attempts: { increment: 1 } },
});
}
// Reconcile files abandoned by a crash between writing bytes and storing metadata.
for (const key of await readdir(storageRoot()).catch(() => [] as string[])) {
try {
const info = await stat(storagePath(key));
if (
now.getTime() - info.mtimeMs > 86400000 &&
!(await db.file.findFirst({ where: { path: key } }))
)
await removeFile(key);
} catch {
/* Unknown filenames are not removed. */
}
}
const deliveries = await db.webhookDelivery.findMany({
where: {
deliveredAt: null,
attempts: { lt: 5 },
nextAttemptAt: { lte: now },
},
take: 50,
});
for (const item of deliveries) {
const hook = await db.webhook.findUnique({ where: { id: item.webhookId } });
if (!hook?.active) {
await db.webhookDelivery.update({
where: { id: item.id },
data: { attempts: 5 },
});
continue;
}
const event = await db.auditEvent.findUnique({
where: { id: item.eventId },
});
if (!event) {
await db.webhookDelivery.update({
where: { id: item.id },
data: { attempts: 5 },
});
continue;
}
try {
await deliver(hook.url, hook.secret, {
id: event.id,
action: event.action,
targetId: event.targetId,
createdAt: event.createdAt,
});
await db.webhookDelivery.update({
where: { id: item.id },
data: { deliveredAt: new Date(), attempts: { increment: 1 } },
});
} catch {
await db.webhookDelivery.update({
where: { id: item.id },
data: {
attempts: { increment: 1 },
nextAttemptAt: new Date(Date.now() + 2 ** item.attempts * 60000),
},
});
}
}
await db.rateLimit.deleteMany({ where: { expiresAt: { lt: now } } });
await db.actionToken.deleteMany({ where: { expiresAt: { lt: now } } });
await db.apiToken.deleteMany({ where: { expiresAt: { lt: now } } });
const retention = Number(process.env.ARCHIVE_RETENTION_DAYS || 0);
if (retention >= 30) {
const old = await db.task.findMany({
where: {
archivedAt: { lt: new Date(Date.now() - retention * 86400000) },
status: "ERLEDIGT",
},
take: 100,
include: { files: true },
});
for (const t of old)
await db.$transaction(async (tx) => {
if (t.files.length)
await tx.cleanupJob.createMany({
data: t.files.map((f) => ({ storageKey: f.path })),
});
await tx.auditEvent.create({
data: {
actorId: "worker",
action: "tasks.retention_delete",
targetId: t.id,
taskId: t.id,
},
});
await tx.task.delete({ where: { id: t.id, version: t.version } });
});
}
}
async function main() {
const shutdown = new AbortController();
const stop = () => shutdown.abort();
process.once("SIGTERM", stop);
process.once("SIGINT", stop);
do {
try {
await tick();
if (process.env.WORKER_HEARTBEAT_FILE)
await writeFile(
process.env.WORKER_HEARTBEAT_FILE,
new Date().toISOString(),
{ mode: 0o600 },
);
} catch (e) {
console.error(
"Worker-Durchlauf fehlgeschlagen",
e instanceof Error ? e.message : e,
);
if (process.argv.includes("--once")) process.exitCode = 1;
}
if (process.argv.includes("--once") || shutdown.signal.aborted) break;
await delay(60000, undefined, { signal: shutdown.signal }).catch(() => {});
} while (!shutdown.signal.aborted);
}
if (!process.env.WORKER_TEST) main().finally(() => db.$disconnect());