feat(health): деплоить probe-Worker из CFDM и опрашивать цели с edge
quality / changes (push) Successful in 9s
quality / commitlint (push) Skipped
quality / docker-check (push) Skipped
CD / update-wiki (push) Successful in 6s
quality / web (push) Successful in 1m4s
quality / api (push) Successful in 54s
CD / quality (push) Successful in 2m17s
CD / publish (push) Successful in 2m21s

Worker сам ходит на origin по Cron Trigger; CFDM кладёт цели в KV и забирает результаты без POST /probe.

Co-authored-by: Cursor <[email protected]>
This commit is contained in:
Denozordec
2026-08-19 17:22:07 +07:00
co-authored by Cursor
parent 2c92e78b24
commit 4c4908558b
38 changed files with 1941 additions and 555 deletions
+10
View File
@@ -35,8 +35,10 @@ import { auditRoutes } from "./routes/audit.js";
import * as certificateService from "./services/certificate-service.js";
import {
createHealthCheckTask,
healthEngineFallbacksFromConfig,
scheduleHealthCheckJob,
} from "./services/health-check-scheduler.js";
import { fireEnsureHealthWorker } from "./services/health/health-worker-deploy.js";
import { AsyncTask, CronJob } from "toad-scheduler";
export interface BuildAppOptions {
@@ -131,6 +133,14 @@ export async function buildApp(opts: BuildAppOptions = {}) {
app.decorate("reloadHealthCheckJob", () => {
scheduleHealthCheckJob(app, config, healthTask);
});
if (config.cloudflareApiToken) {
fireEnsureHealthWorker(
app.db,
app.cf,
healthEngineFallbacksFromConfig(config),
app.log,
);
}
}
return app;
+53
View File
@@ -7,6 +7,8 @@ import type {
} from "@cfdm/shared";
import { createDnsAdapter } from "./cloudflare/dns-service.js";
import { createHealthCheckAdapter, type CfHealthCheckPayload } from "./cloudflare/healthcheck-service.js";
import { createKvAdapter } from "./cloudflare/kv-service.js";
import { createWorkersAdapter } from "./cloudflare/workers-service.js";
import { createZoneAdapter } from "./cloudflare/zone-service.js";
export type { CfHealthCheckPayload };
@@ -15,11 +17,21 @@ export class CloudflareClient {
private readonly zones;
private readonly dns;
private readonly healthchecks;
private readonly kv;
private readonly workers;
private readonly token;
constructor(token: string) {
this.token = token.trim();
this.zones = createZoneAdapter(token);
this.dns = createDnsAdapter(token);
this.healthchecks = createHealthCheckAdapter(token);
this.kv = createKvAdapter(token);
this.workers = createWorkersAdapter(token);
}
get isConfigured(): boolean {
return this.token.length > 0;
}
listZones(): Promise<CfZone[]> {
@@ -81,4 +93,45 @@ export class CloudflareClient {
deleteHealthCheck(zoneId: string, id: string): Promise<void> {
return this.healthchecks.deleteHealthCheck(zoneId, id);
}
listAccounts() {
return this.workers.listAccounts();
}
listKvNamespaces(accountId: string) {
return this.kv.listNamespaces(accountId);
}
createKvNamespace(accountId: string, title: string) {
return this.kv.createNamespace(accountId, title);
}
kvGet(accountId: string, namespaceId: string, key: string) {
return this.kv.getValue(accountId, namespaceId, key);
}
kvPut(accountId: string, namespaceId: string, key: string, value: string) {
return this.kv.putValue(accountId, namespaceId, key, value);
}
putWorkerScript(opts: {
accountId: string;
scriptName: string;
source: string;
kvNamespaceId: string;
}) {
return this.workers.putScript(opts);
}
putWorkerSchedules(accountId: string, scriptName: string, crons: string[]) {
return this.workers.putSchedules(accountId, scriptName, crons);
}
enableWorkersDev(accountId: string, scriptName: string) {
return this.workers.enableWorkersDev(accountId, scriptName);
}
getWorkersSubdomain(accountId: string) {
return this.workers.getWorkersSubdomain(accountId);
}
}
+43
View File
@@ -17,6 +17,15 @@ export function mapCloudflareFailure(
): AppError {
const lower = message.toLowerCase();
if (status === 401 || status === 403 || lower.includes("authentication")) {
if (
operation.includes("workers") ||
operation.includes("kv_") ||
operation.includes("accounts")
) {
return AppError.cloudflareAuthFailed(
"Токену нужны права Account: Workers Scripts Write и Workers KV Storage Write. Zone DNS недостаточно.",
);
}
return AppError.cloudflareAuthFailed(
"Cloudflare отклонил токен. Проверьте CLOUDFLARE_API_TOKEN.",
);
@@ -67,6 +76,40 @@ export async function handleCfResponse<T>(
return body.result;
}
/** KV PUT / schedules often return `{ success: true }` without `result`. */
export async function handleCfSuccess(
response: Response,
operation: string,
): Promise<void> {
if (response.status === 429) {
const wait = parseRetryAfter(response.headers) ?? 5000;
throw AppError.rateLimited(
`Cloudflare временно ограничил запросы. Повторите через ${Math.ceil(wait / 1000)} с.`,
);
}
const text = await response.text();
if (!text) {
if (!response.ok) {
throw mapCloudflareFailure(operation, response.status, String(response.status));
}
return;
}
let body: CfResponse<unknown>;
try {
body = JSON.parse(text) as CfResponse<unknown>;
} catch {
if (!response.ok) {
throw mapCloudflareFailure(operation, response.status, text.slice(0, 180));
}
return;
}
if (!body.success) {
const msg =
body.errors?.map((e) => e.message).join("; ") ?? "unknown cloudflare error";
throw mapCloudflareFailure(operation, response.status, msg);
}
}
export async function cfRequest<T>(
token: string,
path: string,
+100
View File
@@ -0,0 +1,100 @@
import { CF_API_BASE, handleCfResponse, handleCfSuccess, mapCloudflareFailure } from "./http.js";
export interface CfKvNamespace {
id: string;
title: string;
}
export function createKvAdapter(token: string) {
return {
async listNamespaces(accountId: string): Promise<CfKvNamespace[]> {
const all: CfKvNamespace[] = [];
let page = 1;
while (true) {
const url = new URL(
`${CF_API_BASE}/accounts/${accountId}/storage/kv/namespaces`,
);
url.searchParams.set("per_page", "100");
url.searchParams.set("page", String(page));
const response = await fetch(url.toString(), {
headers: { Authorization: `Bearer ${token}` },
signal: AbortSignal.timeout(30_000),
});
if (response.status >= 500 || response.status === 429) {
throw mapCloudflareFailure("kv_list", response.status, String(response.status));
}
const batch = await handleCfResponse<CfKvNamespace[]>(response, "kv_list");
all.push(...batch);
if (batch.length < 100) break;
page += 1;
}
return all;
},
async createNamespace(accountId: string, title: string): Promise<CfKvNamespace> {
const response = await fetch(
`${CF_API_BASE}/accounts/${accountId}/storage/kv/namespaces`,
{
method: "POST",
headers: {
Authorization: `Bearer ${token}`,
"Content-Type": "application/json",
},
body: JSON.stringify({ title }),
signal: AbortSignal.timeout(30_000),
},
);
if (response.status >= 500 || response.status === 429) {
throw mapCloudflareFailure("kv_create", response.status, String(response.status));
}
return handleCfResponse<CfKvNamespace>(response, "kv_create");
},
async getValue(
accountId: string,
namespaceId: string,
key: string,
): Promise<string | null> {
const response = await fetch(
`${CF_API_BASE}/accounts/${accountId}/storage/kv/namespaces/${namespaceId}/values/${encodeURIComponent(key)}`,
{
headers: { Authorization: `Bearer ${token}` },
signal: AbortSignal.timeout(30_000),
},
);
if (response.status === 404) return null;
if (response.status >= 500 || response.status === 429) {
throw mapCloudflareFailure("kv_get", response.status, String(response.status));
}
if (!response.ok) {
const text = await response.text().catch(() => "");
throw mapCloudflareFailure("kv_get", response.status, text.slice(0, 180));
}
return response.text();
},
async putValue(
accountId: string,
namespaceId: string,
key: string,
value: string,
): Promise<void> {
const response = await fetch(
`${CF_API_BASE}/accounts/${accountId}/storage/kv/namespaces/${namespaceId}/values/${encodeURIComponent(key)}`,
{
method: "PUT",
headers: {
Authorization: `Bearer ${token}`,
"Content-Type": "text/plain",
},
body: value,
signal: AbortSignal.timeout(30_000),
},
);
if (response.status >= 500 || response.status === 429) {
throw mapCloudflareFailure("kv_put", response.status, String(response.status));
}
await handleCfSuccess(response, "kv_put");
},
};
}
@@ -0,0 +1,142 @@
import {
CF_API_BASE,
handleCfResponse,
handleCfSuccess,
mapCloudflareFailure,
} from "./http.js";
export interface CfAccount {
id: string;
name?: string;
}
export interface CfWorkersSubdomain {
subdomain?: string;
enabled?: boolean;
}
export function createWorkersAdapter(token: string) {
return {
async listAccounts(): Promise<CfAccount[]> {
const response = await fetch(`${CF_API_BASE}/accounts?per_page=50`, {
headers: { Authorization: `Bearer ${token}` },
signal: AbortSignal.timeout(30_000),
});
if (response.status >= 500 || response.status === 429) {
throw mapCloudflareFailure("list_accounts", response.status, String(response.status));
}
return handleCfResponse<CfAccount[]>(response, "list_accounts");
},
async putScript(opts: {
accountId: string;
scriptName: string;
source: string;
kvNamespaceId: string;
filename?: string;
}): Promise<void> {
const filename = opts.filename ?? "index.mjs";
const metadata = {
main_module: filename,
compatibility_date: "2025-04-01",
bindings: [
{
type: "kv_namespace",
name: "HEALTH_KV",
namespace_id: opts.kvNamespaceId,
},
],
};
const form = new FormData();
form.append(
"metadata",
new Blob([JSON.stringify(metadata)], { type: "application/json" }),
);
form.append(
filename,
new Blob([opts.source], { type: "application/javascript+module" }),
filename,
);
const response = await fetch(
`${CF_API_BASE}/accounts/${opts.accountId}/workers/scripts/${opts.scriptName}`,
{
method: "PUT",
headers: { Authorization: `Bearer ${token}` },
body: form,
signal: AbortSignal.timeout(60_000),
},
);
if (response.status >= 500 || response.status === 429) {
throw mapCloudflareFailure("workers_put_script", response.status, String(response.status));
}
await handleCfSuccess(response, "workers_put_script");
},
async putSchedules(
accountId: string,
scriptName: string,
crons: string[],
): Promise<void> {
const response = await fetch(
`${CF_API_BASE}/accounts/${accountId}/workers/scripts/${scriptName}/schedules`,
{
method: "PUT",
headers: {
Authorization: `Bearer ${token}`,
"Content-Type": "application/json",
},
body: JSON.stringify(crons.map((cron) => ({ cron }))),
signal: AbortSignal.timeout(30_000),
},
);
if (response.status >= 500 || response.status === 429) {
throw mapCloudflareFailure("workers_put_schedules", response.status, String(response.status));
}
await handleCfSuccess(response, "workers_put_schedules");
},
async enableWorkersDev(
accountId: string,
scriptName: string,
): Promise<void> {
const response = await fetch(
`${CF_API_BASE}/accounts/${accountId}/workers/scripts/${scriptName}/subdomain`,
{
method: "POST",
headers: {
Authorization: `Bearer ${token}`,
"Content-Type": "application/json",
},
body: JSON.stringify({ enabled: true }),
signal: AbortSignal.timeout(30_000),
},
);
if (response.status === 409) return;
if (response.status >= 500 || response.status === 429) {
throw mapCloudflareFailure("workers_subdomain", response.status, String(response.status));
}
if (!response.ok && response.status !== 200 && response.status !== 201) {
await handleCfSuccess(response, "workers_subdomain");
}
},
async getWorkersSubdomain(accountId: string): Promise<string | null> {
const response = await fetch(
`${CF_API_BASE}/accounts/${accountId}/workers/subdomain`,
{
headers: { Authorization: `Bearer ${token}` },
signal: AbortSignal.timeout(30_000),
},
);
if (response.status === 404) return null;
if (response.status >= 500 || response.status === 429) {
throw mapCloudflareFailure("workers_get_subdomain", response.status, String(response.status));
}
const result = await handleCfResponse<CfWorkersSubdomain>(
response,
"workers_get_subdomain",
);
return result.subdomain?.trim() || null;
},
};
}
+6 -3
View File
@@ -5,8 +5,9 @@ import * as healthCheckService from "../services/health-check-service.js";
import * as serviceConfigService from "../services/service-config-service.js";
import {
healthEngineFallbacksFromConfig,
resolveWorkerProbeConfig,
} from "../services/health-check-scheduler.js";
import { mailboxFromSettings } from "../services/health/health-worker-deploy.js";
import { cronStaleAfterMs } from "../services/health/mailbox.js";
export async function healthCheckRoutes(app: FastifyInstance) {
app.get("/health-status", async (request) => {
@@ -20,9 +21,10 @@ export async function healthCheckRoutes(app: FastifyInstance) {
app.post("/health-check/run", async (request) => {
const config = request.server.config;
const fallbacks = healthEngineFallbacksFromConfig(config);
const settings = getAppSettings(
request.server.db,
healthEngineFallbacksFromConfig(config),
fallbacks,
);
const thresholds = {
degradedFailures: settings.healthDegradedFailures,
@@ -33,7 +35,8 @@ export async function healthCheckRoutes(app: FastifyInstance) {
const checked = await healthCheckService.runAllChecks(request.server.db, {
thresholds,
probeGapMs: config.healthProbeGapMs,
worker: resolveWorkerProbeConfig(request.server.db, config),
mailbox: mailboxFromSettings(request.server.db, request.server.cf, fallbacks),
staleAfterMs: cronStaleAfterMs(settings.healthCheckCron),
onStatusChange: async (target, prev, next) => {
try {
const label =
+21 -2
View File
@@ -10,6 +10,7 @@ import {
assertValidHealthCron,
healthEngineFallbacksFromConfig,
} from "../services/health-check-scheduler.js";
import { ensureHealthWorker } from "../services/health/health-worker-deploy.js";
export async function settingsRoutes(app: FastifyInstance) {
app.get("/settings", async (request) => {
@@ -40,11 +41,29 @@ export async function settingsRoutes(app: FastifyInstance) {
"ошибок до down не меньше, чем до degraded",
);
}
const next = updateAppSettings(request.server.db, body, fallbacks);
updateAppSettings(request.server.db, body, fallbacks);
if (body.healthCheckCron !== undefined) {
request.server.reloadHealthCheckJob?.();
const after = getAppSettings(request.server.db, fallbacks);
if (after.healthWorkerKvNamespaceId) {
try {
await ensureHealthWorker(
request.server.db,
request.server.cf,
fallbacks,
);
} catch {
// error stored in settings
}
}
}
return next;
return getAppSettings(request.server.db, fallbacks);
});
app.post("/settings/health/worker/ensure", async (request) => {
const fallbacks = healthEngineFallbacksFromConfig(request.server.config);
await ensureHealthWorker(request.server.db, request.server.cf, fallbacks);
return getAppSettings(request.server.db, fallbacks);
});
app.post("/settings/vps-tracker/test", async (request) => {
+15 -14
View File
@@ -2,7 +2,7 @@ import type { FastifyInstance } from "fastify";
import { AsyncTask, CronJob } from "toad-scheduler";
import {
getAppSettings,
getAppSettingsSecrets,
updateAppSettings,
type HealthEngineFallbacks,
} from "@cfdm/db";
import { repos } from "@cfdm/db";
@@ -10,6 +10,10 @@ import type { AppConfig } from "../config.js";
import { AppError } from "../errors.js";
import * as healthCheckService from "./health-check-service.js";
import * as serviceConfigService from "./service-config-service.js";
import {
mailboxFromSettings,
} from "./health/health-worker-deploy.js";
import { cronStaleAfterMs } from "./health/mailbox.js";
declare module "fastify" {
interface FastifyInstance {
@@ -33,18 +37,6 @@ export function healthEngineFallbacksFromConfig(
};
}
export function resolveWorkerProbeConfig(
db: import("@cfdm/db").Db,
config: AppConfig,
): { url: string; token: string } | null {
const settings = getAppSettings(db, healthEngineFallbacksFromConfig(config));
const secrets = getAppSettingsSecrets(db);
const url = settings.healthWorkerUrl.trim();
const token = (secrets.healthWorkerToken || config.healthWorkerToken).trim();
if (!url || !token) return null;
return { url, token };
}
export function assertValidHealthCron(expr: string): void {
const cronExpression = expr.trim();
const parts = cronExpression.split(/\s+/).filter(Boolean);
@@ -78,10 +70,12 @@ export function createHealthCheckTask(
latencyWarnMs: settings.healthLatencyWarnMs,
successRecoveries: settings.healthSuccessRecoveries,
};
const mailbox = mailboxFromSettings(app.db, app.cf, fallbacks);
const n = await healthCheckService.runAllChecks(app.db, {
thresholds,
probeGapMs: config.healthProbeGapMs,
worker: resolveWorkerProbeConfig(app.db, config),
mailbox,
staleAfterMs: cronStaleAfterMs(settings.healthCheckCron),
onStatusChange: async (target, prev, next) => {
try {
const label =
@@ -114,6 +108,13 @@ export function createHealthCheckTask(
}
},
});
if (mailbox) {
updateAppSettings(
app.db,
{ healthWorkerLastIngestAt: new Date().toISOString() },
fallbacks,
);
}
const monitors = await healthCheckService.runDomainMonitors(
app.db,
thresholds,
+135 -94
View File
@@ -7,10 +7,14 @@ import type { HealthCheckTarget, IpHealthState } 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 {
CloudflareWorkerHealthCheckProvider,
workerNotConfiguredResult,
} from "./health/worker.js";
buildTargetsDoc,
indexResults,
isResultsStale,
originProbeKey,
type HealthMailbox,
} from "./health/mailbox.js";
export interface HealthCheckThresholds {
degradedFailures: number;
@@ -267,8 +271,10 @@ export interface RunAllChecksOptions {
thresholds: HealthCheckThresholds;
/** Pause between unique physical probes (default 2000). Same IP is only probed once. */
probeGapMs?: number;
/** Cloudflare Worker URL+token. Missing → cloudflare targets fail, never Local fallback. */
worker?: { url: string; token: string } | null;
/** KV mailbox with Worker results. Missing → cloudflare targets fail, never Local fallback. */
mailbox?: HealthMailbox | null;
/** Results older than this are stale (default 10 min). */
staleAfterMs?: number;
onStatusChange?: (
target: HealthCheckTarget,
prevState: IpHealthState | null,
@@ -286,33 +292,82 @@ function sleep(ms: number): Promise<void> {
*/
export function physicalProbeKey(target: HealthCheckTarget): string {
const kind = target.provider === "cloudflare" ? "cloudflare" : "local";
const port = target.port ?? (target.type === "http" ? 80 : 80);
const ip = String(target.ip || "").trim().toLowerCase();
if (target.type === "http") {
const path = (target.path?.trim() || "/") || "/";
const expected = target.expected_status ?? "";
return `${kind}|http|${ip}|${port}|${path}|${expected}`;
}
if (target.type === "tcp") return `${kind}|tcp|${ip}|${port}`;
if (target.type === "ping") {
return `${kind}|ping|${String(target.hostname || target.ip || "").trim().toLowerCase()}`;
}
if (target.type === "dns") {
return `${kind}|dns|${String(target.hostname || target.ip || "").trim().toLowerCase()}`;
}
return `${kind}|${target.type}|${ip}|${port}`;
return `${kind}|${originProbeKey(target)}`;
}
async function executeProbe(
function applyProbeResult(
db: Db,
target: HealthCheckTarget,
local: LocalHealthCheckProvider,
worker: CloudflareWorkerHealthCheckProvider | null,
): Promise<ProbeResult> {
if (target.provider === "cloudflare") {
if (!worker) return workerNotConfiguredResult();
return worker.probe(target);
result: ProbeResult,
options: RunAllChecksOptions,
): void {
const prev = repos.getIpHealthStatusRow(
db,
target.scope,
target.ref_id,
target.ip,
);
const { state, failures, successes, node } = deriveState(
result.ok,
result.latencyMs,
prev
? {
consecutive_failures: prev.consecutive_failures,
consecutive_successes: prev.consecutive_successes,
status: prev.status,
}
: null,
options.thresholds,
);
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,
failures,
result.error,
successes,
{ colo: result.colo ?? null, provider },
);
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, {
health_status: node,
consecutive_failures: failures,
consecutive_successes: successes,
last_check_at: new Date().toISOString().replace("T", " ").slice(0, 19),
last_failure_reason: result.error,
});
}
return local.probe(target);
if (prevState !== state) {
options.onStatusChange?.(target, prevState, state);
}
}
function staleWorkerResult(colo: string | null): ProbeResult {
return {
ok: false,
latencyMs: 0,
error: "Cloudflare Worker: результаты устарели или KV пуст",
colo,
};
}
export async function runAllChecks(
@@ -322,10 +377,7 @@ export async function runAllChecks(
const targets = repos.listHealthCheckTargets(db);
const gapMs = Math.max(0, options.probeGapMs ?? 2000);
const local = new LocalHealthCheckProvider();
const worker =
options.worker?.url && options.worker.token
? new CloudflareWorkerHealthCheckProvider(options.worker)
: null;
const staleAfterMs = options.staleAfterMs ?? 10 * 60_000;
const byPhysical = new Map<string, HealthCheckTarget[]>();
for (const target of targets) {
@@ -335,80 +387,69 @@ export async function runAllChecks(
else byPhysical.set(key, [target]);
}
let probeIndex = 0;
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;
// Prefer binding hostname for SNI when several scopes share one IP.
const representative =
group.find((t) => t.scope === "binding") ?? group[0]!;
const result = await executeProbe(representative, local, worker);
const result = await local.probe(representative);
for (const target of group) {
const prev = repos.getIpHealthStatusRow(
db,
target.scope,
target.ref_id,
target.ip,
);
const { state, failures, successes, node } = deriveState(
result.ok,
result.latencyMs,
prev
? {
consecutive_failures: prev.consecutive_failures,
consecutive_successes: prev.consecutive_successes,
status: prev.status,
}
: null,
options.thresholds,
);
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,
failures,
result.error,
successes,
{ colo: result.colo ?? null, provider },
);
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, {
health_status: node,
consecutive_failures: failures,
consecutive_successes: successes,
last_check_at: new Date().toISOString().replace("T", " ").slice(0, 19),
last_failure_reason: result.error,
});
applyProbeResult(db, target, result, options);
}
}
if (cloudflareGroups.length > 0) {
const mailbox = options.mailbox ?? null;
const resultsDoc = mailbox ? await mailbox.getResults() : null;
const byKey = indexResults(resultsDoc);
const stale = !mailbox || isResultsStale(resultsDoc, staleAfterMs);
const colo = resultsDoc?.colo ?? null;
if (mailbox) {
try {
const next = buildTargetsDoc(targets);
const current = await mailbox.getTargets();
if (current?.fingerprint !== next.fingerprint) {
await mailbox.putTargets(next);
}
} catch {
// ingest still proceeds
}
if (prevState !== state) {
options.onStatusChange?.(target, prevState, state);
}
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 {
result = {
ok: item.ok,
latencyMs: item.latencyMs,
error: item.error,
colo,
};
}
for (const target of group) {
applyProbeResult(db, target, result, options);
}
}
}
// Orphan rows (old IPs / hostname keys) still feed MAX latency on group badge.
repos.pruneStaleIpHealthStatus(db, targets);
return targets.length;
}
@@ -0,0 +1,21 @@
import { existsSync, readFileSync } from "node:fs";
import { dirname, join } from "node:path";
import { fileURLToPath } from "node:url";
export function loadHealthProbeWorkerSource(): string {
const dir = dirname(fileURLToPath(import.meta.url));
const candidates = [
join(dir, "health-probe-worker.mjs"),
join(process.cwd(), "dist/health-probe-worker.mjs"),
join(process.cwd(), "health-probe-worker.mjs"),
join(dir, "../../../../../workers/health-probe/src/index.mjs"),
join(process.cwd(), "../../workers/health-probe/src/index.mjs"),
join(process.cwd(), "workers/health-probe/src/index.mjs"),
];
for (const path of candidates) {
if (existsSync(path)) {
return readFileSync(path, "utf8");
}
}
throw new Error("не найден исходник Worker health-probe");
}
@@ -0,0 +1,185 @@
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 type { CloudflareClient } from "../../lib/cf-client.js";
import { AppError } from "../../errors.js";
import { loadHealthProbeWorkerSource } from "./health-probe-script.js";
import {
buildTargetsDoc,
createCloudflareKvMailbox,
toCloudflareCron,
type HealthMailbox,
} from "./mailbox.js";
export const DEFAULT_HEALTH_FALLBACKS: HealthEngineFallbacks = {
healthCheckCron: "0 */2 * * * *",
healthDegradedFailures: 1,
healthDownFailures: 2,
healthLatencyWarnMs: 1000,
healthSuccessRecoveries: 2,
healthWorkerUrl: "",
healthWorkerTokenSet: false,
};
export async function resolveAccountId(
cf: CloudflareClient,
db: Db,
cached?: string | null,
): Promise<string> {
const trimmed = cached?.trim();
if (trimmed) return trimmed;
const domains = repos.listDomains(db);
for (const domain of domains) {
if (!domain.cf_zone_id) continue;
try {
const zone = await cf.getZone(domain.cf_zone_id);
const id = zone.account?.id?.trim();
if (id) return id;
} catch {
// try next zone / accounts list
}
}
const accounts = await cf.listAccounts();
const id = accounts[0]?.id?.trim();
if (!id) {
throw AppError.cloudflare(
"Не удалось определить Cloudflare account_id. Добавьте зону или расширьте права токена (Account Settings Read).",
);
}
return id;
}
export async function ensureKvNamespace(
cf: CloudflareClient,
accountId: string,
existingId?: string | null,
): Promise<string> {
if (existingId?.trim()) return existingId.trim();
const listed = await cf.listKvNamespaces(accountId);
const found = listed.find((ns) => ns.title === HEALTH_PROBE_KV_TITLE);
if (found?.id) return found.id;
const created = await cf.createKvNamespace(accountId, HEALTH_PROBE_KV_TITLE);
if (!created.id) {
throw AppError.cloudflare("Cloudflare не вернул id KV namespace");
}
return created.id;
}
export async function ensureHealthWorker(
db: Db,
cf: CloudflareClient,
fallbacks: HealthEngineFallbacks,
): Promise<{ url: string; kvNamespaceId: string; accountId: string }> {
const settings = getAppSettings(db, fallbacks);
try {
const accountId = await resolveAccountId(cf, db, settings.healthWorkerAccountId);
const kvNamespaceId = await ensureKvNamespace(
cf,
accountId,
settings.healthWorkerKvNamespaceId,
);
const source = loadHealthProbeWorkerSource();
await cf.putWorkerScript({
accountId,
scriptName: HEALTH_PROBE_SCRIPT_NAME,
source,
kvNamespaceId,
});
await cf.putWorkerSchedules(accountId, HEALTH_PROBE_SCRIPT_NAME, [
toCloudflareCron(settings.healthCheckCron),
]);
try {
await cf.enableWorkersDev(accountId, HEALTH_PROBE_SCRIPT_NAME);
} catch {
// workers.dev may already be on
}
const subdomain = await cf.getWorkersSubdomain(accountId);
const url = subdomain
? `https://${HEALTH_PROBE_SCRIPT_NAME}.${subdomain}.workers.dev`
: settings.healthWorkerUrl || `https://${HEALTH_PROBE_SCRIPT_NAME}.workers.dev`;
updateAppSettings(
db,
{
healthWorkerAccountId: accountId,
healthWorkerKvNamespaceId: kvNamespaceId,
healthWorkerUrl: url,
healthWorkerError: null,
healthWorkerDeployedAt: new Date().toISOString(),
},
fallbacks,
);
await syncCloudflareTargetsToKv(db, cf, fallbacks);
return { url, kvNamespaceId, accountId };
} catch (err) {
const message = err instanceof Error ? err.message : String(err);
updateAppSettings(db, { healthWorkerError: message }, fallbacks);
throw err;
}
}
export async function maybeEnsureHealthWorker(
db: Db,
cf: CloudflareClient,
fallbacks: HealthEngineFallbacks,
): Promise<void> {
const hasCloudflare = repos
.listHealthCheckTargets(db)
.some((target) => target.provider === "cloudflare");
if (!hasCloudflare) return;
const settings = getAppSettings(db, fallbacks);
if (settings.healthWorkerKvNamespaceId.trim() && !settings.healthWorkerError) {
await syncCloudflareTargetsToKv(db, cf, fallbacks);
return;
}
await ensureHealthWorker(db, cf, fallbacks);
}
export function mailboxFromSettings(
db: Db,
cf: CloudflareClient,
fallbacks: HealthEngineFallbacks,
): HealthMailbox | null {
const settings = getAppSettings(db, fallbacks);
const accountId = settings.healthWorkerAccountId.trim();
const ns = settings.healthWorkerKvNamespaceId.trim();
if (!accountId || !ns) return null;
return createCloudflareKvMailbox(cf, accountId, ns);
}
export async function syncCloudflareTargetsToKv(
db: Db,
cf: CloudflareClient,
fallbacks: HealthEngineFallbacks,
mailbox?: HealthMailbox | null,
): Promise<void> {
const box = mailbox ?? mailboxFromSettings(db, cf, fallbacks);
if (!box) return;
const next = buildTargetsDoc(repos.listHealthCheckTargets(db));
const current = await box.getTargets();
if (current?.fingerprint === next.fingerprint) return;
await box.putTargets(next);
}
export function fireEnsureHealthWorker(
db: Db,
cf: CloudflareClient,
fallbacks: HealthEngineFallbacks,
log?: { warn: (obj: unknown, msg: string) => void },
): void {
if (process.env.VITEST) return;
if (!cf.isConfigured) return;
const hasCloudflare = repos
.listHealthCheckTargets(db)
.some((target) => target.provider === "cloudflare");
if (!hasCloudflare) {
void syncCloudflareTargetsToKv(db, cf, fallbacks).catch((err) => {
log?.warn({ err }, "health worker KV sync failed");
});
return;
}
void maybeEnsureHealthWorker(db, cf, fallbacks).catch((err) => {
log?.warn({ err }, "health worker ensure failed");
});
}
+145
View File
@@ -0,0 +1,145 @@
import type {
HealthCheckTarget,
HealthProbeResultItem,
HealthProbeResultsDoc,
HealthProbeTargetItem,
HealthProbeTargetsDoc,
} from "@cfdm/shared";
import { HEALTH_KV_RESULTS_KEY, HEALTH_KV_TARGETS_KEY } from "@cfdm/shared";
import type { CloudflareClient } from "../../lib/cf-client.js";
export interface HealthMailbox {
getTargets(): Promise<HealthProbeTargetsDoc | null>;
putTargets(doc: HealthProbeTargetsDoc): Promise<void>;
getResults(): Promise<HealthProbeResultsDoc | null>;
}
export function createCloudflareKvMailbox(
cf: CloudflareClient,
accountId: string,
namespaceId: string,
): HealthMailbox {
return {
async getTargets() {
return readJson<HealthProbeTargetsDoc>(cf, accountId, namespaceId, HEALTH_KV_TARGETS_KEY);
},
async putTargets(doc) {
await cf.kvPut(accountId, namespaceId, HEALTH_KV_TARGETS_KEY, JSON.stringify(doc));
},
async getResults() {
return readJson<HealthProbeResultsDoc>(cf, accountId, namespaceId, HEALTH_KV_RESULTS_KEY);
},
};
}
async function readJson<T>(
cf: CloudflareClient,
accountId: string,
namespaceId: string,
key: string,
): Promise<T | null> {
const raw = await cf.kvGet(accountId, namespaceId, key);
if (!raw) return null;
try {
return JSON.parse(raw) as T;
} catch {
return null;
}
}
export function originProbeKey(target: HealthCheckTarget): string {
const port = target.port ?? (target.type === "http" ? 80 : 80);
const ip = String(target.ip || "").trim().toLowerCase();
if (target.type === "http") {
const path = (target.path?.trim() || "/") || "/";
const expected = target.expected_status ?? "";
return `http|${ip}|${port}|${path}|${expected}`;
}
if (target.type === "tcp") return `tcp|${ip}|${port}`;
if (target.type === "ping") {
return `ping|${String(target.hostname || target.ip || "").trim().toLowerCase()}`;
}
if (target.type === "dns") {
return `dns|${String(target.hostname || target.ip || "").trim().toLowerCase()}`;
}
return `${target.type}|${ip}|${port}`;
}
export function cloudflareMailboxTargets(
targets: HealthCheckTarget[],
): HealthProbeTargetItem[] {
const unique = new Map<string, HealthProbeTargetItem>();
for (const target of targets) {
if (target.provider !== "cloudflare") continue;
if (target.type !== "tcp" && target.type !== "http") continue;
const key = originProbeKey(target);
if (unique.has(key)) continue;
unique.set(key, {
key,
ip: target.ip,
hostname: target.hostname || target.ip,
type: target.type,
port: target.port ?? (target.type === "http" ? 80 : 80),
path: target.path ?? "/",
expectedStatus: target.expected_status,
timeoutMs: target.timeout_ms ?? 3000,
verifyTls: Boolean(target.verify_tls),
});
}
return [...unique.values()].sort((a, b) => a.key.localeCompare(b.key));
}
export function fingerprintTargets(items: HealthProbeTargetItem[]): string {
return items
.map(
(item) =>
`${item.key}|${item.hostname}|${item.timeoutMs ?? ""}|${item.verifyTls ? "1" : "0"}`,
)
.join(";");
}
export function buildTargetsDoc(targets: HealthCheckTarget[]): HealthProbeTargetsDoc {
const items = cloudflareMailboxTargets(targets);
return {
fingerprint: fingerprintTargets(items),
updatedAt: new Date().toISOString(),
items,
};
}
export function indexResults(
doc: HealthProbeResultsDoc | null,
): Map<string, HealthProbeResultItem> {
const map = new Map<string, HealthProbeResultItem>();
if (!doc?.items) return map;
for (const item of doc.items) {
map.set(item.key, item);
}
return map;
}
export function isResultsStale(doc: HealthProbeResultsDoc | null, staleAfterMs: number): boolean {
if (!doc?.probedAt) return true;
const ts = Date.parse(doc.probedAt);
if (!Number.isFinite(ts)) return true;
return Date.now() - ts > staleAfterMs;
}
/** Drop seconds from toad 6-field cron for Cloudflare Workers (5-field). */
export function toCloudflareCron(expr: string): string {
const parts = expr.trim().split(/\s+/).filter(Boolean);
if (parts.length === 6) return parts.slice(1).join(" ");
if (parts.length === 5) return parts.join(" ");
throw new Error("некорректное cron-выражение");
}
export function cronStaleAfterMs(expr: string): number {
const cf = toCloudflareCron(expr);
const minute = cf.split(/\s+/)[0] ?? "*";
if (minute.startsWith("*/")) {
const n = Number(minute.slice(2));
if (Number.isFinite(n) && n > 0) return Math.max(n * 2, 5) * 60_000;
}
if (minute === "*") return 10 * 60_000;
return 10 * 60_000;
}
+7 -87
View File
@@ -1,93 +1,13 @@
import type { HealthCheckTarget } from "@cfdm/shared";
import type { ProbeResult } from "../health-check-service.js";
import type { HealthCheckProvider } from "./provider.js";
export interface WorkerProbeConfig {
url: string;
token: string;
}
const WORKER_NOT_CONFIGURED = "Cloudflare Worker не настроен (URL и токен)";
export class CloudflareWorkerHealthCheckProvider implements HealthCheckProvider {
readonly kind = "cloudflare" as const;
constructor(private readonly config: WorkerProbeConfig) {}
async probe(target: HealthCheckTarget): Promise<ProbeResult> {
const base = this.config.url.replace(/\/$/, "");
if (!base || !this.config.token) {
return {
ok: false,
latencyMs: 0,
error: WORKER_NOT_CONFIGURED,
colo: null,
};
}
const timeoutMs = Math.max(100, target.timeout_ms ?? 3000);
const controller = new AbortController();
const timer = setTimeout(() => controller.abort(), timeoutMs + 2500);
try {
const res = await fetch(`${base}/probe`, {
method: "POST",
headers: {
Authorization: `Bearer ${this.config.token}`,
"Content-Type": "application/json",
},
body: JSON.stringify({
type: target.type === "http" ? "http" : "tcp",
ip: target.ip,
hostname: target.hostname,
port: target.port ?? (target.type === "http" ? 80 : 80),
path: target.path ?? "/",
expected_status: target.expected_status,
timeout_ms: timeoutMs,
verify_tls: Boolean(target.verify_tls),
method: "GET",
}),
signal: controller.signal,
});
if (!res.ok) {
const text = await res.text().catch(() => "");
return {
ok: false,
latencyMs: 0,
error: `Worker HTTP ${res.status}${text ? `: ${text.slice(0, 180)}` : ""}`,
colo: null,
};
}
const body = (await res.json()) as {
ok?: boolean;
latencyMs?: number;
error?: string | null;
colo?: string | null;
};
const ok = Boolean(body.ok);
return {
ok,
latencyMs: typeof body.latencyMs === "number" ? body.latencyMs : 0,
error: ok ? null : (body.error ?? "probe failed"),
colo: body.colo ?? null,
};
} catch (err) {
const message =
err instanceof Error
? err.name === "AbortError"
? "Worker timeout"
: err.message
: "Worker probe failed";
return { ok: false, latencyMs: 0, error: message, colo: null };
} finally {
clearTimeout(timer);
}
}
}
export function workerNotConfiguredResult(): ProbeResult {
export function workerNotConfiguredResult(): {
ok: false;
latencyMs: number;
error: string;
colo: null;
} {
return {
ok: false,
latencyMs: 0,
error: WORKER_NOT_CONFIGURED,
error: "Cloudflare Worker не настроен (нет KV mailbox)",
colo: null,
};
}
@@ -24,6 +24,7 @@ import { isValidIpv4 } from "../lib/validators.js";
import * as dnsService from "./dns-service.js";
import * as domainService from "./domain-service.js";
import { syncServiceToVpsTracker } from "./vps-tracker-sync.js";
import { fireEnsureHealthWorker, DEFAULT_HEALTH_FALLBACKS } from "./health/health-worker-deploy.js";
import {
selectActiveIpsByMode,
withBindingLock,
@@ -1250,6 +1251,8 @@ export async function updateConfig(
void syncServiceToVpsTracker(db, id, removedBindingIds);
fireEnsureHealthWorker(db, cf, DEFAULT_HEALTH_FALLBACKS);
const [view] = attachServiceHealth(db, [await buildView(db, id)]);
return view!;
}
@@ -1261,7 +1264,7 @@ export async function createGroup(
): Promise<ServiceGroup> {
const groupType = body.type?.trim() || "custom";
const domain = await normalizeGroupDomain(db, cf, body.domain);
return repos.createServiceGroup(
const group = repos.createServiceGroup(
db,
body.name,
groupType,
@@ -1280,6 +1283,8 @@ export async function createGroup(
health_check_provider: body.health_check_provider,
},
);
fireEnsureHealthWorker(db, cf, DEFAULT_HEALTH_FALLBACKS);
return group;
}
export async function updateGroup(
@@ -1322,6 +1327,7 @@ export async function updateGroup(
group = repos.getServiceGroup(db, id);
}
await syncEnabledServicesInGroup(db, cf, id);
fireEnsureHealthWorker(db, cf, DEFAULT_HEALTH_FALLBACKS);
return group;
}