// The overwritten-list wake pass: each entry is woken via `_wake` with // epoch-CAS'd backoff meta, parking repeat failures until the ttl drops them. import type { Resolution, World } from "./placement.js"; import type { createOutlierDetector } from "./outlier.js"; import type { createPinboard } from "./server.js"; import { attestationHeaders } from "./attestation.js"; import { EPOCH_HEADER, type Transport } from "./proxy.js"; import { runScript } from "./placement.js"; import { overwrittenCasScript } from "./scripts/placement.js"; import { validateNamespace } from "./router.js"; import { logger, serializeError } from "./logger.js"; import { parseOverwrittenField, type Keys } from "./keys.js"; import { discardResponseBody, fetchHeld } from "./http.js"; type OverwrittenMeta = { /** Placement epoch stamped at the latest re-list; fences every delete. */ epoch: string; failures: number; delayUntil: number | null; } & ({ parked: false; parkedAt: null } | { parked: true; parkedAt: number }); // Fresh markers from the resolve script are `{epoch}`: absent backoff fields // default, present ones must be well-typed or the marker counts as corrupt. const parseMeta = (raw: string): OverwrittenMeta => { const bad = (): never => { throw new Error(`malformed overwritten marker: ${raw}`); }; let meta: Record; try { const parsed: unknown = JSON.parse(raw); meta = typeof parsed === "object" && parsed !== null ? (parsed as Record) : bad(); } catch { return bad(); } const { epoch, failures = 0, delayUntil = null, parked = false, parkedAt = null, } = meta; if ( typeof epoch !== "string" || epoch === "" || !Number.isFinite(Number(epoch)) ) { return bad(); } if (typeof failures !== "number") return bad(); if (delayUntil !== null && typeof delayUntil !== "number") return bad(); if (parked === true) { return typeof parkedAt === "number" ? { epoch, failures, delayUntil, parked: true, parkedAt } : bad(); } if (parked !== false) return bad(); return { epoch, failures, delayUntil, parked: false, parkedAt: null }; }; export namespace createWakePass { export type Deps = { redis: createPinboard.RedisKV; keys: Keys; resolve(world: World, ns: string, id: string): Promise; fetch: Transport.FetchLike; outlier?: createOutlierDetector.Instance | undefined; now(): number; random(): number; bindingQuarantined( bucket: number, hostname: string, ns: string, ): Promise; /** Without a DurabilityStore there is nothing to wake: lists are deleted. */ hasStore: boolean; wakeTimeoutMs: number; wakeLimit: number; wakeConcurrency: number; parkedTtlMs: number; baseDelayMs: number; maxDelayMs: number; parkAfterFailures: number; }; } export const createWakePass = ( deps: createWakePass.Deps, ): (() => Promise) => { const { redis, keys, now, random } = deps; const buckets = Array.from({ length: keys.shards }, (_, s) => s); // Background resolves never initiate wind-downs: takeover "deny". A 202 // means the worker is hydrating the instance: neither alive nor failed. const wake = async ( env: string, hostname: string, ns: string, id: string, ): Promise<{ outcome: "ok" | "activating"; epoch: string }> => { const resolution = await deps.resolve( { env, hostname, takeover: "deny" }, ns, id, ); if (resolution.kind !== "forward") { throw new Error(`placement refused: ${resolution.kind}`); } const res = await fetchHeld( deps.fetch, new Request(`${resolution.worker.advertiseUrl}${ns}/${id}/_wake`, { method: "POST", headers: { ...(await attestationHeaders(resolution.worker, {}, now())), [EPOCH_HEADER]: resolution.epoch, }, redirect: "manual", signal: AbortSignal.timeout(deps.wakeTimeoutMs), }), ); await discardResponseBody(res).catch(() => undefined); if (res.status === 202) { return { outcome: "activating", epoch: resolution.epoch }; } await deps.outlier?.report(resolution.worker.sessionId, res.ok); if (!res.ok) throw new Error(`wake answered ${res.status}`); return { outcome: "ok", epoch: resolution.epoch }; }; type DueEntry = { key: string; field: string; entry: NonNullable>; meta: OverwrittenMeta; }; // Backoff meta re-writes CAS on the epoch read this pass, so a concurrent // re-list's fresher marker is never clobbered. const casMeta = (key: string, field: string, meta: OverwrittenMeta) => runScript( redis, overwrittenCasScript, [key], [field, meta.epoch, JSON.stringify(meta)], ); const wakeOne = async ({ key, field, entry, meta }: DueEntry) => { try { const { outcome, epoch } = await wake( entry.env, entry.hostname, validateNamespace(entry.ns), entry.id, ); if (outcome === "ok") { await runScript(redis, overwrittenCasScript, [key], [field, epoch, ""]); return; } await casMeta(key, field, { ...meta, delayUntil: now() + deps.wakeTimeoutMs, }); } catch (err) { logger.warn( { operation: "durability", namespace: entry.ns, id: entry.id, err: serializeError(err), }, "durable wake failed", ); const failures = meta.failures + 1; const delay = Math.min( deps.baseDelayMs * 2 ** (failures - 1), deps.maxDelayMs, ); await casMeta(key, field, { epoch: meta.epoch, failures, delayUntil: now() + delay / 2 + (random() * delay) / 2, ...(failures >= deps.parkAfterFailures ? { parked: true, parkedAt: now() } : { parked: false, parkedAt: null }), }); } }; return async (): Promise => { const due: DueEntry[] = []; for (const bucket of buckets) { const key = keys.overwritten(bucket); if (!deps.hasStore) { await redis.del(key); continue; } for (const [field, metaRaw] of Object.entries(await redis.hgetall(key))) { const entry = parseOverwrittenField(field); if (entry === null) { logger.warn( { operation: "durability", field }, "dropping malformed overwritten entry", ); await redis.hdel(key, field); continue; } // Entries freeze env at eviction; once the hostname's binding moved // to another env the entry is unservable — waking would strand a // placement row in the old env's hash that no sweep scans. The store // still holds the marked row, so the id resurrects if the env returns. const boundEnv = await redis.hget( keys.bindings(bucket), `h:${entry.hostname}${entry.ns}`, ); if (boundEnv !== null && boundEnv !== entry.env) { logger.warn( { operation: "durability", namespace: entry.ns, id: entry.id, env: entry.env, boundEnv, }, "dropping overwritten entry displaced to another env", ); await redis.hdel(key, field); continue; } let meta: OverwrittenMeta; try { meta = parseMeta(metaRaw); } catch (err) { logger.warn( { operation: "durability", namespace: entry.ns, id: entry.id, err: serializeError(err), }, "dropping malformed overwritten marker", ); await redis.hdel(key, field); continue; } if (meta.parked) { if (meta.parkedAt + deps.parkedTtlMs < now()) { logger.warn( { operation: "durability", namespace: entry.ns, id: entry.id }, "dropping parked overwritten entry past ttl", ); await redis.hdel(key, field); } continue; } if ( (meta.delayUntil !== null && meta.delayUntil > now()) || (await deps.bindingQuarantined(bucket, entry.hostname, entry.ns)) ) { continue; } due.push({ key, field, entry, meta }); } } // Stable oldest-first: failed entries always leave with a delayUntil, so // repeat visitors sort behind fresh ones and cannot monopolize the pass. due.sort( (a, b) => (a.meta.delayUntil ?? 0) - (b.meta.delayUntil ?? 0) || a.field.localeCompare(b.field), ); const batch = due.slice(0, deps.wakeLimit); for (let i = 0; i < batch.length; i += deps.wakeConcurrency) { await Promise.all(batch.slice(i, i + deps.wakeConcurrency).map(wakeOne)); } }; };