import type { createPinboard } from "./server.js"; import { runScript } from "./placement.js"; import { membershipScript } from "./scripts/lease.js"; import { logger, serializeError } from "./logger.js"; type RedisKV = createPinboard.RedisKV; export namespace createMembership { export type Options = { key: string; nodeId: string; /** Canonical config JSON; live members must match it byte-for-byte. */ config: string; ttlMs: number; operation: string; now?: (() => number) | undefined; /** Fires once per transition into incompatibility, with the other member's config. */ onMismatch?: ((other: string) => void) | undefined; }; export type Instance = { ok(): Promise; refresh(): Promise; start(): void; close(): Promise; }; } /** * Cluster-membership fence on the `key` hash: each node upserts itself with a * liveness horizon every tick and removes itself on close; a joiner is * refused while any live member carries a different config — an incompatible * fleet is never evicted, only waited out. The verdict is cached for one * tick, mirroring createLease. */ export const createMembership = ( redis: RedisKV, options: createMembership.Options, ): createMembership.Instance => { if (!Number.isFinite(options.ttlMs) || options.ttlMs <= 0) { throw new Error("ttlMs must be a positive finite number"); } const now = options.now ?? (() => Date.now()); const intervalMs = Math.floor(options.ttlMs / 3); let verdict: { ok: boolean; expires: number } | null = null; let inflight: Promise | null = null; const run = (op: "join" | "leave"): Promise => runScript( redis, membershipScript, [options.key], [ options.nodeId, options.config, String(now()), String(options.ttlMs), op, ], ); const refresh = (): Promise => (inflight ??= (async () => { const result = await run("join"); const status = Array.isArray(result) ? result[0] : result; if (status !== "ok" && status !== "incompatible") { throw new Error(`unknown membership status: ${JSON.stringify(result)}`); } const ok = status === "ok"; if (!ok && verdict?.ok !== false) { options.onMismatch?.(String((result as unknown[])[1])); } verdict = { ok, expires: now() + intervalMs }; return ok; })().finally(() => { inflight = null; })); let timer: NodeJS.Timeout | null = null; let closed = false; const schedule = (): void => { if (closed) return; timer = setTimeout(() => { refresh() .catch((err) => { logger.error( { operation: options.operation, err: serializeError(err) }, "membership refresh failed", ); }) .finally(schedule); }, intervalMs); timer.unref?.(); }; return { ok: () => verdict !== null && verdict.expires > now() ? Promise.resolve(verdict.ok) : refresh(), refresh, start: schedule, close: async () => { closed = true; if (timer !== null) clearTimeout(timer); await inflight?.catch(() => {}); await run("leave").catch(() => {}); }, }; };