480 lines
13 KiB
TypeScript
480 lines
13 KiB
TypeScript
import { drainOperationEffects } from "../src/features/operations/worker";
|
|
import { drainFurnitureImports } from "../src/lib/services/furni-job-worker";
|
|
import "./load-env";
|
|
import { Cron } from "croner";
|
|
import { lt, sql } from "drizzle-orm";
|
|
import { env } from "../src/env";
|
|
import { db, PasswordReset, WebsiteLoginLogs } from "../src/lib/db";
|
|
import { logger } from "../src/lib/logger";
|
|
import { redis } from "../src/lib/redis";
|
|
import {
|
|
diskPressure,
|
|
emulatorOffline,
|
|
healthDegraded,
|
|
} from "../src/lib/services/alert";
|
|
import { runCatalogExport } from "../src/lib/services/catalog-git-export";
|
|
import {
|
|
DISK_THRESHOLDS,
|
|
diskLevel,
|
|
parseDfOutput,
|
|
} from "../src/lib/services/disk-usage";
|
|
import { publishDueArticles } from "../src/lib/services/news-scheduler";
|
|
import { scheduledAutoCleanFakeNitros } from "../src/lib/services/nitro-cleanup";
|
|
import { rcon } from "../src/lib/services/rcon";
|
|
|
|
function captureWorkerError(err: unknown, context: string): void {
|
|
logger.error(context, {
|
|
module: "jobs",
|
|
err: err instanceof Error ? err.message : String(err),
|
|
});
|
|
}
|
|
|
|
/** In-process cooldown so a flapping probe does not spam Discord/email. */
|
|
const alertCooldownMs = (env.HEALTH_ALERT_COOLDOWN_MIN ?? 15) * 60_000;
|
|
const lastHealthAlertAt = new Map<string, number>();
|
|
|
|
function canAlert(key: string): boolean {
|
|
const now = Date.now();
|
|
const prev = lastHealthAlertAt.get(key) ?? 0;
|
|
if (now - prev < alertCooldownMs) return false;
|
|
lastHealthAlertAt.set(key, now);
|
|
return true;
|
|
}
|
|
|
|
async function probeHealth(): Promise<{
|
|
database: boolean;
|
|
redis: boolean | null;
|
|
emulator: boolean;
|
|
}> {
|
|
const database = await db
|
|
.execute(sql`SELECT 1`)
|
|
.then(() => true)
|
|
.catch(() => false);
|
|
|
|
let redisOk: boolean | null = null;
|
|
if (env.REDIS_URL) {
|
|
if (!redis) {
|
|
redisOk = false;
|
|
} else {
|
|
try {
|
|
redisOk = (await redis.ping()) === "PONG";
|
|
} catch {
|
|
redisOk = false;
|
|
}
|
|
}
|
|
}
|
|
|
|
const emulator = await rcon.send("ping", null).catch(() => false);
|
|
return {
|
|
database,
|
|
redis: redisOk,
|
|
emulator: Boolean(emulator),
|
|
};
|
|
}
|
|
|
|
async function checkOpsHealth(): Promise<void> {
|
|
try {
|
|
const health = await probeHealth();
|
|
const degraded =
|
|
!health.database || health.redis === false || !health.emulator;
|
|
if (!degraded) return;
|
|
|
|
if (!health.emulator && health.database && health.redis !== false) {
|
|
if (canAlert("emulator")) {
|
|
await emulatorOffline("jobs-worker RCON ping failed");
|
|
}
|
|
return;
|
|
}
|
|
|
|
if (canAlert("health")) {
|
|
await healthDegraded(health);
|
|
}
|
|
} catch (err) {
|
|
captureWorkerError(err, "Health probe failed");
|
|
}
|
|
}
|
|
|
|
/**
|
|
* Host-side: probe filesystem fill levels via `df -P -B1` and raise a
|
|
* diskPressure() alert per mount once it crosses 85/90/95% (cooldown-gated per
|
|
* mount+level so a stuck fill level does not spam Discord/email).
|
|
*/
|
|
async function checkDiskUsage(): Promise<void> {
|
|
const { execFile } = await import("node:child_process");
|
|
|
|
let dfOut: string;
|
|
try {
|
|
dfOut = await new Promise<string>((resolve, reject) => {
|
|
execFile(
|
|
"df",
|
|
["-P", "-B1"],
|
|
{
|
|
maxBuffer: 4 * 1024 * 1024,
|
|
timeout: 15_000,
|
|
},
|
|
(err, stdout) => (err ? reject(err) : resolve(stdout)),
|
|
);
|
|
});
|
|
} catch (err) {
|
|
captureWorkerError(err, "Disk probe failed (is df available?)");
|
|
return;
|
|
}
|
|
|
|
const mounts = parseDfOutput(dfOut);
|
|
if (mounts.length === 0) return;
|
|
|
|
for (const usage of mounts) {
|
|
const level = diskLevel(usage.percent);
|
|
if (level === 0) continue;
|
|
if (canAlert(`disk-${usage.mount}-${level}`)) {
|
|
logger.warn("Raising disk usage alert", {
|
|
module: "jobs",
|
|
mount: usage.mount,
|
|
percent: usage.percent,
|
|
});
|
|
await diskPressure(usage);
|
|
}
|
|
|
|
// Self-healing: never let a mount max out. Safe prune from the warning
|
|
// mark, forced prune (drop age windows) from the error mark. Cooldown is
|
|
// keyed separately so a stuck fill level does not re-prune every run
|
|
// while still escalating to the forced path on real pressure.
|
|
if (usage.percent >= DISK_THRESHOLDS.error) {
|
|
if (canAlert(`disk-prune-force-${usage.mount}`)) {
|
|
logger.warn("Disk near full — forcing Docker cache reclaim", {
|
|
module: "jobs",
|
|
mount: usage.mount,
|
|
percent: usage.percent,
|
|
});
|
|
await pruneDockerCache(true);
|
|
}
|
|
} else if (usage.percent >= DISK_THRESHOLDS.warning) {
|
|
if (canAlert(`disk-prune-${usage.mount}`)) {
|
|
logger.info("Disk at high water mark — running safe Docker prune", {
|
|
module: "jobs",
|
|
mount: usage.mount,
|
|
percent: usage.percent,
|
|
});
|
|
await pruneDockerCache(false);
|
|
}
|
|
}
|
|
}
|
|
}
|
|
|
|
async function backupEmulatorJar(): Promise<void> {
|
|
if (!env.EMULATOR_JAR_PATH || !env.EMULATOR_BACKUP_DIR) return;
|
|
|
|
const { copyFileSync, mkdirSync, readdirSync, unlinkSync, existsSync } =
|
|
await import("node:fs");
|
|
const { resolve } = await import("node:path");
|
|
|
|
const timestamp = new Date().toISOString().slice(0, 19).replace(/[T:]/g, "-");
|
|
const backupFile = resolve(
|
|
env.EMULATOR_BACKUP_DIR,
|
|
`emulator-${timestamp}.jar`,
|
|
);
|
|
|
|
if (!existsSync(env.EMULATOR_BACKUP_DIR)) {
|
|
mkdirSync(env.EMULATOR_BACKUP_DIR, { recursive: true });
|
|
}
|
|
|
|
try {
|
|
copyFileSync(env.EMULATOR_JAR_PATH, backupFile);
|
|
logger.info("Backed up emulator JAR", {
|
|
module: "jobs",
|
|
backupFile,
|
|
});
|
|
|
|
const keep = env.EMULATOR_BACKUP_KEEP ?? 7;
|
|
const files = readdirSync(env.EMULATOR_BACKUP_DIR)
|
|
.filter((f) => f.startsWith("emulator-") && f.endsWith(".jar"))
|
|
.sort()
|
|
.reverse();
|
|
|
|
for (let i = keep; i < files.length; i++) {
|
|
unlinkSync(resolve(env.EMULATOR_BACKUP_DIR, files[i]));
|
|
logger.info("Rotated out old backup", {
|
|
module: "jobs",
|
|
file: files[i],
|
|
});
|
|
}
|
|
} catch (err) {
|
|
captureWorkerError(err, "JAR backup failed");
|
|
}
|
|
}
|
|
|
|
/** Optional mysqldump when DB_BACKUP_DIR is set (host must have mysqldump on PATH). */
|
|
async function backupDatabase(): Promise<void> {
|
|
const backupDir = env.DB_BACKUP_DIR;
|
|
if (!backupDir || !env.DATABASE_URL) return;
|
|
|
|
const { mkdirSync, readdirSync, unlinkSync, existsSync, createWriteStream } =
|
|
await import("node:fs");
|
|
const { resolve } = await import("node:path");
|
|
const { spawn } = await import("node:child_process");
|
|
|
|
let parsed: URL;
|
|
try {
|
|
parsed = new URL(env.DATABASE_URL);
|
|
} catch {
|
|
logger.error("Invalid DATABASE_URL for DB backup", { module: "jobs" });
|
|
return;
|
|
}
|
|
|
|
if (!existsSync(backupDir)) {
|
|
mkdirSync(backupDir, { recursive: true });
|
|
}
|
|
|
|
const timestamp = new Date().toISOString().slice(0, 19).replace(/[T:]/g, "-");
|
|
const dbName =
|
|
decodeURIComponent(parsed.pathname.replace(/^\//, "")) || "cms";
|
|
const outFile = resolve(backupDir, `db-${dbName}-${timestamp}.sql`);
|
|
|
|
const args = [
|
|
`-h${parsed.hostname}`,
|
|
`-P${parsed.port || "3306"}`,
|
|
`-u${decodeURIComponent(parsed.username)}`,
|
|
`--single-transaction`,
|
|
`--routines`,
|
|
`--databases`,
|
|
dbName,
|
|
];
|
|
if (parsed.password) {
|
|
args.splice(3, 0, `-p${decodeURIComponent(parsed.password)}`);
|
|
}
|
|
|
|
await new Promise<void>((resolvePromise) => {
|
|
const child = spawn("mysqldump", args, {
|
|
stdio: ["ignore", "pipe", "pipe"],
|
|
});
|
|
const out = createWriteStream(outFile);
|
|
child.stdout.pipe(out);
|
|
let stderr = "";
|
|
child.stderr.on("data", (chunk: Buffer) => {
|
|
stderr += chunk.toString();
|
|
});
|
|
child.on("error", (err) => {
|
|
captureWorkerError(err, "mysqldump spawn failed (is it on PATH?)");
|
|
resolvePromise();
|
|
});
|
|
child.on("close", (code) => {
|
|
out.end();
|
|
if (code !== 0) {
|
|
captureWorkerError(
|
|
new Error(stderr || `mysqldump exit ${code}`),
|
|
"DB backup failed",
|
|
);
|
|
} else {
|
|
logger.info("Backed up database", { module: "jobs", outFile });
|
|
const keep = env.DB_BACKUP_KEEP ?? 7;
|
|
const files = readdirSync(backupDir)
|
|
.filter((f) => f.startsWith("db-") && f.endsWith(".sql"))
|
|
.sort()
|
|
.reverse();
|
|
for (let i = keep; i < files.length; i++) {
|
|
const file = files[i];
|
|
if (file) unlinkSync(resolve(backupDir, file));
|
|
}
|
|
}
|
|
resolvePromise();
|
|
});
|
|
});
|
|
}
|
|
|
|
async function cleanupOldLogs(): Promise<void> {
|
|
try {
|
|
const cutoff = new Date(Date.now() - 30 * 24 * 60 * 60 * 1000);
|
|
await db
|
|
.delete(WebsiteLoginLogs)
|
|
.where(lt(WebsiteLoginLogs.createdAt, cutoff));
|
|
logger.info("Cleaned up login logs older than 30 days", { module: "jobs" });
|
|
} catch (err) {
|
|
captureWorkerError(err, "Log cleanup failed");
|
|
}
|
|
}
|
|
|
|
async function cleanupOldSessions(): Promise<void> {
|
|
try {
|
|
const cutoff = new Date(Date.now() - 7 * 24 * 60 * 60 * 1000);
|
|
await db.delete(PasswordReset).where(lt(PasswordReset.createdAt, cutoff));
|
|
logger.info("Cleaned up expired password reset tokens", {
|
|
module: "jobs",
|
|
});
|
|
} catch (err) {
|
|
captureWorkerError(err, "Session cleanup failed");
|
|
}
|
|
}
|
|
|
|
/** Host-side: reclaim Docker's unused cache (build cache, unreferenced images,
|
|
* stopped containers). Volumes and in-use images are never touched. No-op when
|
|
* docker or the prune script is unavailable. With force=true the age windows
|
|
* are dropped (docker-prune.sh --force) so every unused byte is reclaimed —
|
|
* the emergency path for a nearly-full disk. */
|
|
async function pruneDockerCache(force = false): Promise<void> {
|
|
const { access } = await import("node:fs/promises");
|
|
const { resolve } = await import("node:path");
|
|
const { spawn } = await import("node:child_process");
|
|
const script = resolve(process.cwd(), "scripts", "docker-prune.sh");
|
|
try {
|
|
await access(script);
|
|
} catch {
|
|
return;
|
|
}
|
|
await new Promise<void>((resolvePromise) => {
|
|
const child = spawn("bash", [script, ...(force ? ["--force"] : [])], {
|
|
stdio: "ignore",
|
|
});
|
|
child.on("error", (err) =>
|
|
captureWorkerError(
|
|
err,
|
|
"Docker prune could not start (is bash on PATH?)",
|
|
),
|
|
);
|
|
child.on("close", (code) => {
|
|
if (code !== 0)
|
|
captureWorkerError(
|
|
new Error(`docker-prune.sh exited ${code}`),
|
|
"Docker prune failed",
|
|
);
|
|
resolvePromise();
|
|
});
|
|
});
|
|
}
|
|
|
|
async function publishScheduledArticles(): Promise<void> {
|
|
try {
|
|
const published = await publishDueArticles();
|
|
if (published > 0) {
|
|
logger.info(`Published ${published} scheduled article(s)`, {
|
|
module: "jobs",
|
|
});
|
|
}
|
|
} catch (err) {
|
|
captureWorkerError(err, "Scheduled article publish failed");
|
|
}
|
|
}
|
|
|
|
/** Nightly self-cleaning of old fake .nitro leftovers (skips when a scan is
|
|
* running in the studio). Requires the `items_base` truth list, so a DB outage
|
|
* shows up here as a captured error instead of deleting anything. */
|
|
async function autoCleanNitroFakes(): Promise<void> {
|
|
try {
|
|
const result = await scheduledAutoCleanFakeNitros(30);
|
|
if (result.skipped) {
|
|
logger.info("Skipped nightly nitro auto-clean: scan in progress", {
|
|
module: "jobs",
|
|
});
|
|
} else {
|
|
logger.info("Nightly nitro auto-clean finished", {
|
|
module: "jobs",
|
|
deleted: result.deleted,
|
|
});
|
|
}
|
|
} catch (err) {
|
|
captureWorkerError(err, "Nitro auto-clean failed");
|
|
}
|
|
}
|
|
|
|
async function reportWorkerHeartbeat(): Promise<void> {
|
|
if (!redis) return;
|
|
await redis
|
|
.set("cms:jobs-worker:heartbeat", new Date().toISOString(), "EX", 600)
|
|
.catch((error) => captureWorkerError(error, "WorkerHeartbeat"));
|
|
}
|
|
|
|
async function main() {
|
|
new Cron("* * * * *", () => {
|
|
void drainOperationEffects();
|
|
});
|
|
void drainOperationEffects();
|
|
new Cron("* * * * *", () => {
|
|
void drainFurnitureImports();
|
|
});
|
|
void drainFurnitureImports();
|
|
new Cron("* * * * *", () => {
|
|
runCatalogExport().catch((e) =>
|
|
captureWorkerError(e, "Catalog export failed"),
|
|
);
|
|
});
|
|
logger.info("Worker started", { module: "jobs" });
|
|
|
|
if (env.EMULATOR_JAR_PATH && env.EMULATOR_BACKUP_DIR) {
|
|
new Cron("0 3 * * *", () => {
|
|
backupEmulatorJar().catch((e) => captureWorkerError(e, "Backup error"));
|
|
});
|
|
logger.info("Scheduled: emulator JAR backup (daily 03:00)", {
|
|
module: "jobs",
|
|
});
|
|
}
|
|
|
|
if (env.DB_BACKUP_DIR) {
|
|
new Cron("30 3 * * *", () => {
|
|
backupDatabase().catch((e) => captureWorkerError(e, "DB backup error"));
|
|
});
|
|
logger.info("Scheduled: mysqldump DB backup (daily 03:30)", {
|
|
module: "jobs",
|
|
});
|
|
}
|
|
|
|
new Cron("0 4 * * *", () => {
|
|
Promise.all([cleanupOldLogs(), cleanupOldSessions()]).catch((e) =>
|
|
captureWorkerError(e, "Cleanup error"),
|
|
);
|
|
});
|
|
logger.info("Scheduled: old data cleanup (daily 04:00)", { module: "jobs" });
|
|
|
|
new Cron("0 5 * * *", () => {
|
|
pruneDockerCache().catch((e) =>
|
|
captureWorkerError(e, "Docker prune error"),
|
|
);
|
|
});
|
|
logger.info("Scheduled: Docker cache prune (daily 05:00)", {
|
|
module: "jobs",
|
|
});
|
|
|
|
new Cron("*/5 * * * *", () => {
|
|
checkOpsHealth().catch((e) => captureWorkerError(e, "Health check error"));
|
|
});
|
|
logger.info("Scheduled: ops health probe (every 5 min)", { module: "jobs" });
|
|
|
|
new Cron("*/5 * * * *", () => {
|
|
checkDiskUsage().catch((e) => captureWorkerError(e, "Disk check error"));
|
|
});
|
|
logger.info("Scheduled: disk usage probe (every 5 min)", { module: "jobs" });
|
|
|
|
new Cron("* * * * *", () => {
|
|
publishScheduledArticles().catch((e) =>
|
|
captureWorkerError(e, "ScheduledArticlePublish"),
|
|
);
|
|
});
|
|
logger.info("Scheduled: publish scheduled articles (every minute)", {
|
|
module: "jobs",
|
|
});
|
|
|
|
new Cron("* * * * *", () => {
|
|
void reportWorkerHeartbeat();
|
|
});
|
|
await reportWorkerHeartbeat();
|
|
|
|
new Cron("0 2 * * *", () => {
|
|
autoCleanNitroFakes().catch((e) =>
|
|
captureWorkerError(e, "Nitro auto-clean error"),
|
|
);
|
|
});
|
|
logger.info("Scheduled: nitro auto-clean (daily 02:00)", {
|
|
module: "jobs",
|
|
});
|
|
await Promise.all([
|
|
backupEmulatorJar(),
|
|
cleanupOldLogs(),
|
|
cleanupOldSessions(),
|
|
checkOpsHealth(),
|
|
checkDiskUsage(),
|
|
]);
|
|
}
|
|
|
|
main().catch((err) => {
|
|
captureWorkerError(err, "Fatal");
|
|
process.exit(1);
|
|
});
|