feat(health): добавить Globalping и мультивыбор источников проб
CD / update-wiki (push) Successful in 8s
quality / commitlint (push) Skipped
quality / changes (push) Successful in 5s
quality / docker-check (push) Skipped
quality / web (push) Successful in 54s
quality / api (push) Successful in 46s
CD / quality (push) Successful in 1m49s
CD / publish (push) Successful in 1m40s

Несколько источников проб сразу и правило агрегации на сервисе вместо XOR Local/Cloudflare.

Co-authored-by: Cursor <[email protected]>
This commit is contained in:
Denozordec
2026-08-19 18:32:18 +07:00
co-authored by Cursor
parent 4c4908558b
commit b9bea44dce
31 changed files with 2671 additions and 271 deletions
@@ -2,6 +2,7 @@ import type { FastifyInstance } from "fastify";
import { AsyncTask, CronJob } from "toad-scheduler";
import {
getAppSettings,
getAppSettingsSecrets,
updateAppSettings,
type HealthEngineFallbacks,
} from "@cfdm/db";
@@ -71,11 +72,17 @@ export function createHealthCheckTask(
successRecoveries: settings.healthSuccessRecoveries,
};
const mailbox = mailboxFromSettings(app.db, app.cf, fallbacks);
const secrets = getAppSettingsSecrets(app.db);
const n = await healthCheckService.runAllChecks(app.db, {
thresholds,
probeGapMs: config.healthProbeGapMs,
mailbox,
staleAfterMs: cronStaleAfterMs(settings.healthCheckCron),
globalping: {
token: secrets.globalpingToken,
locations: secrets.globalpingLocations,
limit: secrets.globalpingLimit,
},
onStatusChange: async (target, prev, next) => {
try {
const label =
+137 -69
View File
@@ -3,11 +3,20 @@ import { resolve4, resolve6 } from "node:dns/promises";
import { Agent, buildConnector, fetch as undiciFetch } from "undici";
import type { Db } from "@cfdm/db";
import { repos } from "@cfdm/db";
import type { HealthCheckTarget, IpHealthState } from "@cfdm/shared";
import type { HealthCheckTarget, IpHealthState, HealthCheckProvider } from "@cfdm/shared";
import {
aggregateHealthOk,
parseHealthAggregate,
targetProviders,
} from "@cfdm/shared";
import { AppError } from "../errors.js";
import { nextHealthState } from "./health/state-machine.js";
import { LocalHealthCheckProvider } from "./health/local.js";
import { workerNotConfiguredResult } from "./health/worker.js";
import {
globalpingNotConfiguredResult,
probeWithGlobalping,
} from "./health/globalping.js";
import {
buildTargetsDoc,
indexResults,
@@ -15,6 +24,7 @@ import {
originProbeKey,
type HealthMailbox,
} from "./health/mailbox.js";
import type { GlobalpingClientOptions } from "../lib/globalping-client.js";
export interface HealthCheckThresholds {
degradedFailures: number;
@@ -275,6 +285,7 @@ export interface RunAllChecksOptions {
mailbox?: HealthMailbox | null;
/** Results older than this are stale (default 10 min). */
staleAfterMs?: number;
globalping?: GlobalpingClientOptions | null;
onStatusChange?: (
target: HealthCheckTarget,
prevState: IpHealthState | null,
@@ -287,20 +298,60 @@ function sleep(ms: number): Promise<void> {
}
/**
* One network hit per key. Group+binding on the same IP share a single TCP/HTTP probe
* so anti-bot / rate-limit on the origin is not tripped by back-to-back checks.
* One network hit per origin+provider. Group+binding on the same IP share a probe.
*/
export function physicalProbeKey(target: HealthCheckTarget): string {
const kind = target.provider === "cloudflare" ? "cloudflare" : "local";
return `${kind}|${originProbeKey(target)}`;
export function physicalProbeKey(
target: HealthCheckTarget,
provider: HealthCheckProvider = target.provider,
): string {
return `${provider}|${originProbeKey(target)}`;
}
function applyProbeResult(
function logSourceResult(
db: Db,
target: HealthCheckTarget,
provider: HealthCheckProvider,
result: ProbeResult,
): void {
repos.insertHealthProbeLog(db, {
scope: target.scope,
refId: target.ref_id,
ip: target.ip,
provider,
status: result.ok ? "up" : "down",
ok: result.ok,
latencyMs: result.latencyMs,
colo: result.colo ?? null,
error: result.error,
});
}
function applyAggregatedStatus(
db: Db,
target: HealthCheckTarget,
sources: Array<{ provider: HealthCheckProvider; result: ProbeResult }>,
options: RunAllChecksOptions,
): void {
const policy = parseHealthAggregate(target.aggregate);
const oks = sources.map((s) => s.result.ok);
const aggregatedOk = aggregateHealthOk(oks, policy);
const latencies = sources.map((s) => s.result.latencyMs);
const latencyMs = latencies.length
? Math.round(latencies.reduce((sum, n) => sum + n, 0) / latencies.length)
: 0;
const colo =
sources.find((s) => s.result.colo)?.result.colo ??
sources[0]?.result.colo ??
null;
const error = aggregatedOk
? null
: sources
.map((s) => s.result.error)
.filter((msg): msg is string => Boolean(msg))
.join("; ") || "health aggregate down";
const statusProvider =
sources.length > 1 ? "aggregate" : (sources[0]?.provider ?? target.provider);
const prev = repos.getIpHealthStatusRow(
db,
target.scope,
@@ -308,8 +359,8 @@ function applyProbeResult(
target.ip,
);
const { state, failures, successes, node } = deriveState(
result.ok,
result.latencyMs,
aggregatedOk,
latencyMs,
prev
? {
consecutive_failures: prev.consecutive_failures,
@@ -322,30 +373,18 @@ function applyProbeResult(
const prevState: IpHealthState | null = prev
? (prev.status as IpHealthState)
: null;
const provider = target.provider === "cloudflare" ? "cloudflare" : "local";
repos.upsertIpHealthStatus(
db,
target.scope,
target.ref_id,
target.ip,
state,
result.latencyMs,
latencyMs,
failures,
result.error,
error,
successes,
{ colo: result.colo ?? null, provider },
{ colo, provider: statusProvider },
);
repos.insertHealthProbeLog(db, {
scope: target.scope,
refId: target.ref_id,
ip: target.ip,
provider,
status: state,
ok: result.ok,
latencyMs: result.latencyMs,
colo: result.colo ?? null,
error: result.error,
});
const matchedNode = repos.findNodeByIp(db, target.ip);
if (matchedNode && matchedNode.enabled) {
repos.updateNode(db, matchedNode.id, {
@@ -353,7 +392,7 @@ function applyProbeResult(
consecutive_failures: failures,
consecutive_successes: successes,
last_check_at: new Date().toISOString().replace("T", " ").slice(0, 19),
last_failure_reason: result.error,
last_failure_reason: error,
});
}
if (prevState !== state) {
@@ -379,42 +418,26 @@ export async function runAllChecks(
const local = new LocalHealthCheckProvider();
const staleAfterMs = options.staleAfterMs ?? 10 * 60_000;
const byPhysical = new Map<string, HealthCheckTarget[]>();
const byOrigin = new Map<string, HealthCheckTarget[]>();
for (const target of targets) {
const key = physicalProbeKey(target);
const list = byPhysical.get(key);
const key = originProbeKey(target);
const list = byOrigin.get(key);
if (list) list.push(target);
else byPhysical.set(key, [target]);
else byOrigin.set(key, [target]);
}
const localGroups: HealthCheckTarget[][] = [];
const cloudflareGroups: HealthCheckTarget[][] = [];
for (const group of byPhysical.values()) {
if (group[0]?.provider === "cloudflare") cloudflareGroups.push(group);
else localGroups.push(group);
}
let probeIndex = 0;
for (const group of localGroups) {
if (probeIndex > 0 && gapMs > 0) {
await sleep(gapMs);
}
probeIndex += 1;
const representative =
group.find((t) => t.scope === "binding") ?? group[0]!;
const result = await local.probe(representative);
for (const target of group) {
applyProbeResult(db, target, result, options);
}
}
if (cloudflareGroups.length > 0) {
const mailbox = options.mailbox ?? null;
const needsCloudflare = targets.some((t) =>
targetProviders(t).includes("cloudflare"),
);
let mailboxResults = new Map<string, { ok: boolean; latencyMs: number; error: string | null }>();
let mailboxColo: string | null = null;
let mailboxStale = true;
const mailbox = options.mailbox ?? null;
if (needsCloudflare) {
const resultsDoc = mailbox ? await mailbox.getResults() : null;
const byKey = indexResults(resultsDoc);
const stale = !mailbox || isResultsStale(resultsDoc, staleAfterMs);
const colo = resultsDoc?.colo ?? null;
mailboxResults = indexResults(resultsDoc);
mailboxStale = !mailbox || isResultsStale(resultsDoc, staleAfterMs);
mailboxColo = resultsDoc?.colo ?? null;
if (mailbox) {
try {
const next = buildTargetsDoc(targets);
@@ -426,28 +449,71 @@ export async function runAllChecks(
// ingest still proceeds
}
}
}
for (const group of cloudflareGroups) {
const representative =
group.find((t) => t.scope === "binding") ?? group[0]!;
const item = byKey.get(originProbeKey(representative));
let result: ProbeResult;
if (!mailbox) {
result = workerNotConfiguredResult();
} else if (stale || !item) {
result = staleWorkerResult(colo);
} else {
const probeCache = new Map<string, ProbeResult>();
let probeIndex = 0;
async function resolveProvider(
provider: HealthCheckProvider,
representative: HealthCheckTarget,
originKey: string,
): Promise<ProbeResult> {
const cacheKey = `${provider}|${originKey}`;
const cached = probeCache.get(cacheKey);
if (cached) return cached;
let result: ProbeResult;
if (provider === "local") {
if (probeIndex > 0 && gapMs > 0) await sleep(gapMs);
probeIndex += 1;
result = await local.probe(representative);
} else if (provider === "cloudflare") {
const item = mailboxResults.get(originKey);
if (!mailbox) result = workerNotConfiguredResult();
else if (mailboxStale || !item) result = staleWorkerResult(mailboxColo);
else {
result = {
ok: item.ok,
latencyMs: item.latencyMs,
error: item.error,
colo,
colo: mailboxColo,
};
}
for (const target of group) {
applyProbeResult(db, target, result, options);
} else {
if (!options.globalping?.token?.trim()) {
result = globalpingNotConfiguredResult();
} else {
if (probeIndex > 0 && gapMs > 0) await sleep(gapMs);
probeIndex += 1;
result = await probeWithGlobalping(representative, options.globalping);
}
}
probeCache.set(cacheKey, result);
return result;
}
for (const [originKey, group] of byOrigin) {
const representative =
group.find((t) => t.scope === "binding") ?? group[0]!;
const needed = new Set<HealthCheckProvider>();
for (const target of group) {
for (const provider of targetProviders(target)) needed.add(provider);
}
for (const provider of needed) {
await resolveProvider(provider, representative, originKey);
}
for (const target of group) {
const providers = targetProviders(target);
const sources = providers.map((provider) => ({
provider,
result: probeCache.get(`${provider}|${originKey}`)!,
}));
for (const source of sources) {
logSourceResult(db, target, source.provider, source.result);
}
applyAggregatedStatus(db, target, sources, options);
}
}
repos.pruneStaleIpHealthStatus(db, targets);
@@ -473,6 +539,8 @@ export async function runDomainMonitors(
timeout_ms: monitor.timeout_ms,
verify_tls: false,
provider: "local",
providers: ["local"],
aggregate: "majority",
};
let result: ProbeResult;
if (monitor.type === "http") {
@@ -0,0 +1,40 @@
import type { HealthCheckTarget } from "@cfdm/shared";
import {
runGlobalpingMeasurement,
type GlobalpingClientOptions,
} from "../../lib/globalping-client.js";
import type { ProbeResult } from "../health-check-service.js";
export function globalpingNotConfiguredResult(): ProbeResult {
return {
ok: false,
latencyMs: 0,
error: "Globalping: токен не задан",
colo: null,
};
}
export async function probeWithGlobalping(
target: HealthCheckTarget,
options: GlobalpingClientOptions,
): Promise<ProbeResult> {
if (!options.token?.trim()) {
return globalpingNotConfiguredResult();
}
try {
const result = await runGlobalpingMeasurement(target, options);
return {
ok: result.ok,
latencyMs: result.latencyMs,
error: result.error,
colo: result.colo,
};
} catch (err) {
return {
ok: false,
latencyMs: 0,
error: err instanceof Error ? err.message : "Globalping: ошибка запроса",
colo: null,
};
}
}
@@ -1,6 +1,6 @@
import { getAppSettings, repos, updateAppSettings, type HealthEngineFallbacks } from "@cfdm/db";
import type { Db } from "@cfdm/db";
import { HEALTH_PROBE_KV_TITLE, HEALTH_PROBE_SCRIPT_NAME } from "@cfdm/shared";
import { HEALTH_PROBE_KV_TITLE, HEALTH_PROBE_SCRIPT_NAME, targetHasProvider } from "@cfdm/shared";
import type { CloudflareClient } from "../../lib/cf-client.js";
import { AppError } from "../../errors.js";
import { loadHealthProbeWorkerSource } from "./health-probe-script.js";
@@ -126,7 +126,7 @@ export async function maybeEnsureHealthWorker(
): Promise<void> {
const hasCloudflare = repos
.listHealthCheckTargets(db)
.some((target) => target.provider === "cloudflare");
.some((target) => targetHasProvider(target, "cloudflare"));
if (!hasCloudflare) return;
const settings = getAppSettings(db, fallbacks);
if (settings.healthWorkerKvNamespaceId.trim() && !settings.healthWorkerError) {
@@ -172,7 +172,7 @@ export function fireEnsureHealthWorker(
if (!cf.isConfigured) return;
const hasCloudflare = repos
.listHealthCheckTargets(db)
.some((target) => target.provider === "cloudflare");
.some((target) => targetHasProvider(target, "cloudflare"));
if (!hasCloudflare) {
void syncCloudflareTargetsToKv(db, cf, fallbacks).catch((err) => {
log?.warn({ err }, "health worker KV sync failed");
+2 -2
View File
@@ -5,7 +5,7 @@ import type {
HealthProbeTargetItem,
HealthProbeTargetsDoc,
} from "@cfdm/shared";
import { HEALTH_KV_RESULTS_KEY, HEALTH_KV_TARGETS_KEY } from "@cfdm/shared";
import { HEALTH_KV_RESULTS_KEY, HEALTH_KV_TARGETS_KEY, targetHasProvider } from "@cfdm/shared";
import type { CloudflareClient } from "../../lib/cf-client.js";
export interface HealthMailbox {
@@ -70,7 +70,7 @@ export function cloudflareMailboxTargets(
): HealthProbeTargetItem[] {
const unique = new Map<string, HealthProbeTargetItem>();
for (const target of targets) {
if (target.provider !== "cloudflare") continue;
if (!targetHasProvider(target, "cloudflare")) continue;
if (target.type !== "tcp" && target.type !== "http") continue;
const key = originProbeKey(target);
if (unique.has(key)) continue;
@@ -2,6 +2,8 @@ import type { Db } from "@cfdm/db";
import { repos } from "@cfdm/db";
import type {
DnsRecord,
HealthCheckAggregate,
HealthCheckProvider,
HealthCheckScope,
HealthCheckType,
IpHealthState,
@@ -51,7 +53,9 @@ export interface ServiceDomainInput {
health_check_interval_sec?: number;
health_check_timeout_ms?: number;
health_check_verify_tls?: boolean;
health_check_provider?: "local" | "cloudflare";
health_check_provider?: HealthCheckProvider;
health_check_providers?: HealthCheckProvider[];
health_check_aggregate?: HealthCheckAggregate;
}
export interface ToggleRequest {
@@ -72,7 +76,9 @@ export interface ServiceGroupBody {
health_check_interval_sec?: number;
health_check_timeout_ms?: number;
health_check_verify_tls?: boolean;
health_check_provider?: "local" | "cloudflare";
health_check_provider?: HealthCheckProvider;
health_check_providers?: HealthCheckProvider[];
health_check_aggregate?: HealthCheckAggregate;
}
export interface UpdateServiceGroupBody {
@@ -89,7 +95,9 @@ export interface UpdateServiceGroupBody {
health_check_interval_sec?: number;
health_check_timeout_ms?: number;
health_check_verify_tls?: boolean;
health_check_provider?: "local" | "cloudflare";
health_check_provider?: HealthCheckProvider;
health_check_providers?: HealthCheckProvider[];
health_check_aggregate?: HealthCheckAggregate;
}
export interface UpdateServiceConfigRequest {
@@ -296,6 +304,10 @@ async function buildView(db: Db, serviceId: number): Promise<ServiceView> {
health_check_timeout_ms: binding.health_check_timeout_ms,
health_check_verify_tls: binding.health_check_verify_tls,
health_check_provider: binding.health_check_provider ?? "local",
health_check_providers: binding.health_check_providers ?? [
binding.health_check_provider ?? "local",
],
health_check_aggregate: binding.health_check_aggregate ?? "majority",
sync_status: aggregateSyncStatus(statuses),
};
});
@@ -1164,7 +1176,9 @@ export async function updateConfig(
input.health_check_interval_sec !== undefined ||
input.health_check_timeout_ms !== undefined ||
input.health_check_verify_tls !== undefined ||
input.health_check_provider !== undefined
input.health_check_provider !== undefined ||
input.health_check_providers !== undefined ||
input.health_check_aggregate !== undefined
) {
repos.updateBindingLbConfig(db, binding.id, {
lb_mode: input.lb_mode,
@@ -1177,6 +1191,8 @@ export async function updateConfig(
health_check_timeout_ms: input.health_check_timeout_ms,
health_check_verify_tls: input.health_check_verify_tls,
health_check_provider: input.health_check_provider,
health_check_providers: input.health_check_providers,
health_check_aggregate: input.health_check_aggregate,
});
}
@@ -1281,6 +1297,8 @@ export async function createGroup(
health_check_timeout_ms: body.health_check_timeout_ms,
health_check_verify_tls: body.health_check_verify_tls,
health_check_provider: body.health_check_provider,
health_check_providers: body.health_check_providers,
health_check_aggregate: body.health_check_aggregate,
},
);
fireEnsureHealthWorker(db, cf, DEFAULT_HEALTH_FALLBACKS);
@@ -1320,6 +1338,8 @@ export async function updateGroup(
health_check_timeout_ms: body.health_check_timeout_ms,
health_check_verify_tls: body.health_check_verify_tls,
health_check_provider: body.health_check_provider,
health_check_providers: body.health_check_providers,
health_check_aggregate: body.health_check_aggregate,
},
);
if (!domain && group.enabled) {