#!/usr/bin/env node import fs from 'node:fs'; import path from 'node:path'; import { Agent, fetch as undiciFetch } from 'undici'; import { createClient } from 'redis'; function loadEnvFile(filePath) { if (!fs.existsSync(filePath)) return; for (const line of fs.readFileSync(filePath, 'utf8').split('\n')) { const trimmed = line.trim(); if (!trimmed || trimmed.startsWith('#')) continue; const idx = trimmed.indexOf('='); if (idx < 0) continue; const key = trimmed.slice(0, idx).trim(); const value = trimmed.slice(idx + 1).trim(); if (!process.env[key]) process.env[key] = value; } } function positiveInteger(value, fallback) { const n = Number(value); if (!Number.isFinite(n) || n <= 0) return fallback; return Math.floor(n); } function targetWorkers() { const targets = (process.env.TKMIND_API_TARGETS || process.env.TKMIND_API_TARGET || '') .split(',') .map((value) => value.trim()) .filter(Boolean); return targets.map((target, index) => ({ id: `goosed-${index + 1}`, target, })); } function workerKey(namespace, id, field) { return [namespace, 'worker', id, field].join(':'); } async function probeTarget(target, timeoutMs) { const startedAt = Date.now(); const controller = new AbortController(); const timeout = setTimeout(() => controller.abort(), timeoutMs); const dispatcher = new Agent({ connect: { rejectUnauthorized: false } }); try { const url = new URL('/status', target.endsWith('/') ? target : `${target}/`).toString(); const res = await undiciFetch(url, { signal: controller.signal, dispatcher: url.startsWith('https://127.0.0.1') ? dispatcher : undefined, headers: process.env.TKMIND_SERVER__SECRET_KEY ? { 'X-Secret-Key': process.env.TKMIND_SERVER__SECRET_KEY } : undefined, }); const body = (await res.text()).trim().slice(0, 80); return { ok: res.ok && body === 'ok', statusCode: res.status, body, latencyMs: Date.now() - startedAt, error: null, }; } catch (err) { return { ok: false, statusCode: 0, body: '', latencyMs: Date.now() - startedAt, error: err instanceof Error ? err.message : String(err), }; } finally { clearTimeout(timeout); dispatcher.close(); } } async function writeHeartbeats(client, namespace, workers, options) { const now = Date.now(); const results = []; for (const worker of workers) { const probe = await probeTarget(worker.target, options.timeoutMs); const ttlSeconds = Math.max(5, Math.ceil(options.ttlMs / 1000)); const multi = client .multi() .set(workerKey(namespace, worker.id, 'heartbeat'), String(now), { EX: ttlSeconds }) .set(workerKey(namespace, worker.id, 'heartbeat_at'), String(now), { EX: ttlSeconds }) .set(workerKey(namespace, worker.id, 'heartbeat_source'), 'sidecar', { EX: ttlSeconds }) .set(workerKey(namespace, worker.id, 'heartbeat_target'), worker.target, { EX: ttlSeconds }) .set(workerKey(namespace, worker.id, 'heartbeat_ok'), probe.ok ? '1' : '0', { EX: ttlSeconds }) .set(workerKey(namespace, worker.id, 'heartbeat_status_code'), String(probe.statusCode), { EX: ttlSeconds }) .set(workerKey(namespace, worker.id, 'heartbeat_latency_ms'), String(probe.latencyMs), { EX: ttlSeconds }); if (probe.error) { multi.set(workerKey(namespace, worker.id, 'heartbeat_error'), probe.error, { EX: ttlSeconds }); } else { multi.del(workerKey(namespace, worker.id, 'heartbeat_error')); } await multi.exec(); results.push({ ...worker, ...probe, heartbeatAt: now }); } return results; } loadEnvFile(process.env.MEMIND_ENV_FILE || path.join(process.cwd(), '.env')); const redisUrl = process.env.MEMIND_RUNTIME_REDIS_URL || 'redis://127.0.0.1:6379/0'; const namespace = process.env.MEMIND_RUNTIME_REDIS_NAMESPACE || 'memind:runtime'; const action = process.argv[2] || 'once'; const intervalMs = positiveInteger(process.env.MEMIND_RUNTIME_HEARTBEAT_INTERVAL_MS, 15_000); const timeoutMs = positiveInteger(process.env.MEMIND_RUNTIME_HEARTBEAT_TIMEOUT_MS, 5_000); const ttlMs = positiveInteger(process.env.MEMIND_RUNTIME_HEARTBEAT_TTL_MS, Math.max(45_000, intervalMs * 3)); if (!['once', 'serve'].includes(action)) { console.error('Usage: node scripts/runtime-worker-heartbeat.mjs '); process.exit(2); } const workers = targetWorkers(); if (workers.length === 0) { console.error('No TKMIND_API_TARGETS configured'); process.exit(2); } const client = createClient({ url: redisUrl }); client.on('error', (err) => { console.error(`Redis error: ${err instanceof Error ? err.message : err}`); }); await client.connect(); let stopping = false; const stop = async () => { stopping = true; await client.quit().catch(() => {}); }; process.on('SIGTERM', () => void stop().finally(() => process.exit(0))); process.on('SIGINT', () => void stop().finally(() => process.exit(0))); async function tick() { const result = await writeHeartbeats(client, namespace, workers, { timeoutMs, ttlMs }); const payload = { ok: result.every((worker) => worker.ok), action, namespace, intervalMs: action === 'serve' ? intervalMs : undefined, timeoutMs, ttlMs, workers: result, }; console.log(JSON.stringify(payload)); return payload; } if (action === 'once') { const payload = await tick(); await client.quit(); process.exit(payload.ok ? 0 : 1); } await tick(); while (!stopping) { await new Promise((resolve) => setTimeout(resolve, intervalMs)); if (!stopping) await tick().catch((err) => { console.error(err instanceof Error ? err.message : String(err)); }); }