143 lines
4.6 KiB
JavaScript
143 lines
4.6 KiB
JavaScript
import {DatabaseSync} from "node:sqlite";
|
|
|
|
const sourcePath = process.env.KUMA_DATABASE || "/migration/kuma.db";
|
|
const targetPath = process.env.REGISTRY_DATABASE || "/data/registry.db";
|
|
const rawRetentionDays = Number.parseInt(process.env.RAW_RETENTION_DAYS || "30", 10);
|
|
const aggregateRetentionDays = Number.parseInt(process.env.AGGREGATE_RETENTION_DAYS || "365", 10);
|
|
|
|
const source = new DatabaseSync(sourcePath, {readOnly: true});
|
|
const target = new DatabaseSync(targetPath);
|
|
target.exec("PRAGMA foreign_keys = ON; PRAGMA busy_timeout = 5000;");
|
|
|
|
const components = target.prepare(`
|
|
SELECT id, monitor_json FROM registry_component WHERE monitor_json IS NOT NULL
|
|
`).all();
|
|
const sourceMonitors = source.prepare(`
|
|
SELECT m.id, m.name, m.url, COUNT(h.id) AS checks
|
|
FROM monitor m
|
|
LEFT JOIN heartbeat h ON h.monitor_id = m.id
|
|
GROUP BY m.id
|
|
ORDER BY checks DESC, m.id ASC
|
|
`).all();
|
|
const heartbeats = source.prepare(`
|
|
SELECT status, time, ping, msg
|
|
FROM heartbeat
|
|
WHERE monitor_id = ? AND time >= ?
|
|
ORDER BY time ASC
|
|
`);
|
|
const dailyStats = source.prepare(`
|
|
SELECT timestamp, ping, up, down
|
|
FROM stat_daily
|
|
WHERE monitor_id = ? AND timestamp >= ?
|
|
ORDER BY timestamp ASC
|
|
`);
|
|
|
|
const insertCheck = target.prepare(`
|
|
INSERT OR IGNORE INTO monitor_check (
|
|
component_id, checked_at, ok, status_code, latency_ms, error
|
|
) VALUES (?, ?, ?, ?, ?, ?)
|
|
`);
|
|
const upsertHourly = target.prepare(`
|
|
INSERT INTO monitor_hourly (component_id, bucket, total, successful, latency_sum)
|
|
VALUES (?, ?, 1, ?, ?)
|
|
ON CONFLICT(component_id, bucket) DO UPDATE SET
|
|
total = total + 1,
|
|
successful = successful + excluded.successful,
|
|
latency_sum = latency_sum + excluded.latency_sum
|
|
`);
|
|
const upsertDaily = target.prepare(`
|
|
INSERT INTO monitor_daily (component_id, bucket, total, successful, latency_sum)
|
|
VALUES (?, ?, ?, ?, ?)
|
|
ON CONFLICT(component_id, bucket) DO NOTHING
|
|
`);
|
|
|
|
function isoFromKuma(value) {
|
|
return new Date(value.endsWith("Z") ? value : `${value}Z`).toISOString();
|
|
}
|
|
|
|
function statusCode(message) {
|
|
const match = typeof message === "string" ? message.match(/\b(\d{3})\b/) : null;
|
|
return match ? Number.parseInt(match[1], 10) : null;
|
|
}
|
|
|
|
function safeError(message, ok) {
|
|
if (ok || typeof message !== "string") {
|
|
return null;
|
|
}
|
|
return message.replace(/[\r\n\t]+/g, " ").slice(0, 300);
|
|
}
|
|
|
|
function normalizedUrl(value) {
|
|
try {
|
|
const url = new URL(value);
|
|
url.pathname = url.pathname === "/" ? "" : url.pathname.replace(/\/$/, "");
|
|
return url.href.replace(/\/$/, "");
|
|
} catch {
|
|
return value;
|
|
}
|
|
}
|
|
|
|
const rawCutoff = new Date(Date.now() - rawRetentionDays * 86400000);
|
|
const kumaRawCutoff = rawCutoff.toISOString().replace("T", " ").replace("Z", "");
|
|
const aggregateCutoff = Math.floor((Date.now() - aggregateRetentionDays * 86400000) / 1000);
|
|
const imported = [];
|
|
|
|
target.exec("BEGIN IMMEDIATE");
|
|
try {
|
|
for (const component of components) {
|
|
const monitor = JSON.parse(component.monitor_json);
|
|
const candidates = sourceMonitors.filter((candidate) => normalizedUrl(candidate.url) === normalizedUrl(monitor.url));
|
|
if (!candidates.length) {
|
|
continue;
|
|
}
|
|
|
|
const canonical = candidates[0];
|
|
let rawChecks = 0;
|
|
for (const heartbeat of heartbeats.all(canonical.id, kumaRawCutoff)) {
|
|
const checkedAt = isoFromKuma(heartbeat.time);
|
|
const ok = heartbeat.status === 1;
|
|
const result = insertCheck.run(
|
|
component.id,
|
|
checkedAt,
|
|
ok ? 1 : 0,
|
|
statusCode(heartbeat.msg),
|
|
Number.isFinite(heartbeat.ping) ? Math.max(0, Math.round(heartbeat.ping)) : null,
|
|
safeError(heartbeat.msg, ok)
|
|
);
|
|
if (result.changes > 0) {
|
|
const latency = Number.isFinite(heartbeat.ping) ? Math.max(0, Math.round(heartbeat.ping)) : 0;
|
|
upsertHourly.run(component.id, `${checkedAt.slice(0, 13)}:00:00.000Z`, ok ? 1 : 0, latency);
|
|
rawChecks += 1;
|
|
}
|
|
}
|
|
|
|
let dailyRows = 0;
|
|
for (const daily of dailyStats.all(canonical.id, aggregateCutoff)) {
|
|
const total = daily.up + daily.down;
|
|
const bucket = new Date(daily.timestamp * 1000).toISOString().slice(0, 10) + "T00:00:00.000Z";
|
|
const result = upsertDaily.run(
|
|
component.id,
|
|
bucket,
|
|
total,
|
|
daily.up,
|
|
Math.max(0, Math.round((daily.ping || 0) * total))
|
|
);
|
|
dailyRows += result.changes;
|
|
}
|
|
|
|
imported.push({
|
|
component: component.id,
|
|
monitorId: canonical.id,
|
|
monitorName: canonical.name,
|
|
rawChecks,
|
|
dailyRows
|
|
});
|
|
}
|
|
target.exec("COMMIT");
|
|
} catch (error) {
|
|
target.exec("ROLLBACK");
|
|
throw error;
|
|
}
|
|
|
|
console.log(JSON.stringify({imported}, null, 2));
|