lumbda/factory/bend-supervisor.sh
russell@unturf.com 1665893321
factory + quantum + sweep-doctrine: AGPLv3 share-back from foxhop ecdsa
29 new files publish factory infra (V2 autoscaler with live VRAM
sampling + EWMA peak tracking, HUGE solo-dispatch, two-tier DLQ/rDLQ
classifier + retry), general quantum circuit primitives (Cuccaro
ripple-carry adder, Clifford gate library, Clifford tableau simulator,
mod-arith family, dialog GCD reversible inverse, Karatsuba multiplier,
Solinas fast reduction), and a TCRAUDT reducer harness. Originally
developed in ~/git/www.foxhop.net/ecdsa/ for secp256k1 attack-surface
research; published upstream as obligated by AGPLv3.

Parametrization contract at factory/CONTRACT.md. Consumers export
LUMBDA_REPO_DIR + LUMBDA_QUEUE_DIR + LUMBDA_BACKEND_CMD + LUMBDA_EMITTER_CMD
then exec factory scripts. No fork-and-modify; single source of truth
upstream.

Integration tests gate 7 V2 defect classes that wedged a live factory
on 2026-06-12 (skewed-demand starve, zero-floor reservation,
multi-tier greedy, +-25%% damping, cold-start ramp, DLQ surge halve,
post-damp CPU ceiling) + 28 DLQ classifier cases (auto-retry vs
escalate partition) + bash -n syntax lint across every script.

GPU backend stays in consumer trees; rationale in
factory/GPU-BACKEND-NOTE.md. Bend wire protocol + gpu-worker.lsp
already upstream at examples/cuda-fanout/.

make factory-lint                bash -n on every factory/*.sh
make test-integration            V2 reducer + DLQ classifier + syntax gate
make sweep-doctrine              TCRAUDT reducer gate (serial)
make sweep-doctrine-parallel     xargs -P fan-out

Verified on neoblanka: factory-lint 12 scripts PASS; test-integration
14 V2 cases + 28 DLQ classifier cases + 12 syntax cases all PASS.
2026-06-14 10:37:35 -04:00

353 lines
14 KiB
Bash
Executable file

#!/usr/bin/env bash
# bend-supervisor.sh — per-host local watchdog for a bend factory.
#
# Runs ON each backend host. Owns local liveness for that host's components:
# 1. backend (gpu-worker) on $BEND_PORT_LISTEN
# 2. bend-emit-pool.sh (consumes .lsp from $LUMBDA_QUEUE_DIR)
# 3. bend-dispatcher.sh (consumes .ready, writes results.tsv)
# 4. bend-supervisor-dlq-runner.sh (DLQ classifier + auto-retry)
# 5. bend-autoscaler.sh (V2 live-VRAM controller; optional)
#
# Reads heartbeats produced by lib-heartbeat.sh:
# $LUMBDA_QUEUE_DIR/pool.heartbeat (pool)
# $LUMBDA_QUEUE_DIR/dispatcher-<wid>.heartbeat (dispatcher per worker)
# backend liveness: TCP LISTEN socket on $BEND_PORT_LISTEN (not an hb)
#
# Restart policy:
# - backend: restart when port not listening
# - pool: restart when heartbeat stale AND any .lsp lacks marker
# - dispatcher: restart when ALL workers' heartbeats stale AND .ready
# or .bin exists (work waiting).
#
# Persistent loop. Exit only on:
# - $LUMBDA_QUEUE_DIR/SUPERVISOR_STOP marker
# - SIGINT / SIGTERM
# No auto-exit on idle.
#
# Env vars consumed: see $LUMBDA_FACTORY_DIR/CONTRACT.md (LUMBDA_QUEUE_DIR,
# LUMBDA_FACTORY_DIR, LUMBDA_DOMAIN_DIR, LUMBDA_REPO_DIR, LUMBDA_BACKEND_CMD,
# LUMBDA_EMITTER_CMD, BEND_PORT_LISTEN, BEND_ENDPOINTS, POOL_MAX,
# POOL_MIN_BIN, POOL_BATCHES, POLL_SECONDS, POOL_STALE_SEC, DISP_STALE_SEC,
# SUPERVISOR_LOG).
#
# Usage:
# LUMBDA_QUEUE_DIR=/tmp/lumbda-queue BEND_PORT_LISTEN=8320 \
# bend-supervisor.sh
set -u
QUEUE_DIR="${LUMBDA_QUEUE_DIR:-${QUEUE_DIR:-/tmp/lumbda-queue}}"
LUMBDA_DOMAIN_DIR="${LUMBDA_DOMAIN_DIR:-}"
LUMBDA_REPO_DIR="${LUMBDA_REPO_DIR:-$HOME/git/lumbda}"
LUMBDA_FACTORY_DIR="${LUMBDA_FACTORY_DIR:-$(dirname "$(readlink -f "$0")")}"
# Backend launch knobs — set by consumer.
LUMBDA_BACKEND_CMD="${LUMBDA_BACKEND_CMD:-bend-cuda}"
# Optional cuda-fanout dir for lumbda-worker bootstrap script.
CUDA_FANOUT_DIR="${CUDA_FANOUT_DIR:-$LUMBDA_REPO_DIR/examples/cuda-fanout}"
# Lumbda C binary — needed only for the gpu-worker.lsp bootstrap path.
# Consumer can override; if missing we skip the worker launch helper and
# the operator handles backend lifecycle out-of-band.
LUMBDA_C_BIN="${LUMBDA_C_BIN:-$LUMBDA_REPO_DIR/c/lumbda}"
BEND_PORT_LISTEN="${BEND_PORT_LISTEN:-8320}"
BEND_ENDPOINTS="${BEND_ENDPOINTS:-127.0.0.1:${BEND_PORT_LISTEN}}"
POOL_MAX="${POOL_MAX:-12}"
POOL_MIN_BIN="${POOL_MIN_BIN:-268435456}"
POOL_BATCHES="${POOL_BATCHES:-141}"
POLL_SECONDS="${POLL_SECONDS:-30}"
POOL_STALE_SEC="${POOL_STALE_SEC:-180}"
DISP_STALE_SEC="${DISP_STALE_SEC:-300}"
SUPERVISOR_LOG="${SUPERVISOR_LOG:-$QUEUE_DIR/supervisor.log}"
LIB_HEARTBEAT="${LIB_HEARTBEAT:-$LUMBDA_FACTORY_DIR/lib-heartbeat.sh}"
LIB_TIER="${LIB_TIER:-$LUMBDA_FACTORY_DIR/lib-tier.sh}"
# shellcheck source=/dev/null
. "$LIB_HEARTBEAT"
# Source tier policy so reload_config's tier_fingerprint sees defaults at
# first boot. supervisor.config overrides flow through `set -a; . "$CONFIG_FILE";
# set +a` below + propagate to child env.
# shellcheck source=/dev/null
. "$LIB_TIER"
mkdir -p "$QUEUE_DIR"
PID_FILE="$QUEUE_DIR/supervisor.pid"
STOP_FILE="$QUEUE_DIR/SUPERVISOR_STOP"
# Singleton via PID file. If a prior supervisor's pid stays alive, refuse
# to start. If pidfile stays stale (pid dead), claim it.
if [ -f "$PID_FILE" ]; then
old_pid=$(cat "$PID_FILE" 2>/dev/null)
if [ -n "$old_pid" ] && kill -0 "$old_pid" 2>/dev/null; then
echo "[$(date +%H:%M:%S)] bend-supervisor: another instance alive (pid=$old_pid); refusing to start" >&2
exit 1
fi
fi
echo $$ > "$PID_FILE"
SUP_EXIT_REASON=""
on_exit() {
local r="${SUP_EXIT_REASON:-unknown-exit}"
heartbeat_exit_cause "$QUEUE_DIR" "supervisor" "$r"
rm -f "$PID_FILE"
echo "[$(date +%H:%M:%S)] bend-supervisor exit reason=$r" >> "$SUPERVISOR_LOG"
}
trap on_exit EXIT
trap 'SUP_EXIT_REASON=signal-int; exit 130' INT
trap 'SUP_EXIT_REASON=signal-term; exit 143' TERM
log() {
local ts msg
ts=$(date +%H:%M:%S)
msg="[$ts] $*"
echo "$msg" | tee -a "$SUPERVISOR_LOG"
}
log "bend-supervisor start queue=$QUEUE_DIR port=$BEND_PORT_LISTEN"
log " pool-stale=${POOL_STALE_SEC}s disp-stale=${DISP_STALE_SEC}s poll=${POLL_SECONDS}s"
# ── component-1: backend (gpu-worker) ───────────────────────────────
# Listening on TCP $BEND_PORT_LISTEN means alive. No heartbeat file —
# backend = lumbda + cuda, not a script we own.
bend_port_listening() {
ss -tnlp 2>/dev/null | grep -q "0.0.0.0:${BEND_PORT_LISTEN} "
}
restart_bend() {
log "restart backend gpu-worker on :${BEND_PORT_LISTEN}"
local launch_lsp="/tmp/launch-worker-${BEND_PORT_LISTEN}.lsp"
local bend_log="/tmp/bend-worker-${BEND_PORT_LISTEN}.log"
if [ ! -x "$LUMBDA_C_BIN" ] || [ ! -d "$CUDA_FANOUT_DIR" ]; then
log "skip backend restart — LUMBDA_C_BIN ($LUMBDA_C_BIN) or CUDA_FANOUT_DIR ($CUDA_FANOUT_DIR) missing; consumer owns backend lifecycle"
return 1
fi
printf '%s\n' \
"(define *argv* (quote (\"--port\" \"${BEND_PORT_LISTEN}\")))" \
"(load \"wire.lsp\")" \
"(load \"gpu-worker.lsp\")" \
"(main)" \
> "$launch_lsp"
(
cd "$CUDA_FANOUT_DIR" || exit 1
PYTHONUNBUFFERED=1 nohup "$LUMBDA_C_BIN" "$launch_lsp" \
>> "$bend_log" 2>&1 &
)
sleep 5
if bend_port_listening; then
log "backend back up on :${BEND_PORT_LISTEN}"
return 0
fi
log "backend FAILED to come back up — will retry next poll"
return 1
}
# ── component-2: bend-emit-pool ─────────────────────────────────────
# Heartbeat at $QUEUE_DIR/pool.heartbeat. Stale check delegated to lib.
# Restart predicate: stale AND there's actually un-emitted work.
pool_alive() {
heartbeat_is_alive "$QUEUE_DIR" "pool" "$POOL_STALE_SEC"
}
count_unworked_lsp() {
local n=0
for f in "$QUEUE_DIR"/*.lsp; do
[ -f "$f" ] || continue
local tag
tag=$(basename "$f" .lsp)
[ -f "$QUEUE_DIR/${tag}.ready" ] && continue
[ -f "$QUEUE_DIR/${tag}.done" ] && continue
[ -f "$QUEUE_DIR/${tag}.emitting" ] && continue
n=$((n + 1))
done
echo "$n"
}
restart_pool() {
# Factory-protection: bash -n the pool script BEFORE spawning. A
# one-character apostrophe in a comment can close the outer
# exec -a NAME bash -c quote and crash-loop a pool every 30s.
# bash -n catches that statically. If a script breaks, refuse to
# spawn + log loudly. Last-known-good in-memory pool keeps running
# until a human fixes master.
local pool_script="$LUMBDA_FACTORY_DIR/bend-emit-pool.sh"
if ! bash -n "$pool_script" 2>/tmp/pool-syntax.err; then
log "REFUSING to spawn pool — bend-emit-pool.sh failed bash -n:"
sed 's/^/ /' /tmp/pool-syntax.err | while read line; do log "$line"; done
log "fix master + push; supervisor will retry on next loop"
rm -f /tmp/pool-syntax.err
return 1
fi
rm -f /tmp/pool-syntax.err
log "restart bend-emit-pool"
# Clean stale .emitting + STOP markers before relaunch — the old pool
# may have left these around.
rm -f "$QUEUE_DIR"/*.emitting "$QUEUE_DIR/STOP" 2>/dev/null
LUMBDA_DOMAIN_DIR="$LUMBDA_DOMAIN_DIR" \
LUMBDA_REPO_DIR="$LUMBDA_REPO_DIR" \
LUMBDA_FACTORY_DIR="$LUMBDA_FACTORY_DIR" \
LUMBDA_EMITTER_CMD="${LUMBDA_EMITTER_CMD:-}" \
setsid nohup "$pool_script" \
"$POOL_MAX" "$QUEUE_DIR" "$POOL_MIN_BIN" \
>> "$QUEUE_DIR/pool.log" 2>&1 < /dev/null &
}
# ── component-3: bend-dispatcher ────────────────────────────────────
# Multi-worker dispatcher writes per-worker heartbeats. Considered alive
# if ANY worker's heartbeat stays fresh — single-worker stalls absorbed
# by others. Considered dead only if NO worker heartbeat stays fresh
# AND .ready or .bin work waits.
#
# Discovery: glob $QUEUE_DIR/dispatcher-*.heartbeat. Empty glob → no
# dispatcher ever ran → restart if work waits.
dispatcher_any_alive() {
local hb
for hb in "$QUEUE_DIR"/dispatcher-*.heartbeat; do
[ -f "$hb" ] || continue
local svc
svc=$(basename "$hb" .heartbeat)
if heartbeat_is_alive "$QUEUE_DIR" "$svc" "$DISP_STALE_SEC"; then
return 0
fi
done
# Fall through: singleton dispatcher (this script) writes
# dispatcher.heartbeat directly when present.
heartbeat_is_alive "$QUEUE_DIR" "dispatcher" "$DISP_STALE_SEC"
}
count_dispatchable() {
local n
n=$(ls "$QUEUE_DIR"/*.ready "$QUEUE_DIR"/*.bin 2>/dev/null | wc -l)
echo "$n"
}
restart_dispatcher() {
log "restart bend-dispatcher BEND_ENDPOINTS=$BEND_ENDPOINTS"
# Stale heartbeats / inflight markers will be reclaimed by a new
# dispatcher's startup orphan-reclaim block. Clean only the lock so
# new instance can flock.
rm -f "$QUEUE_DIR/dispatcher.lock" "$QUEUE_DIR"/dispatcher-*.heartbeat 2>/dev/null
BEND_ENDPOINTS="$BEND_ENDPOINTS" \
LUMBDA_DOMAIN_DIR="$LUMBDA_DOMAIN_DIR" \
LUMBDA_BACKEND_CMD="$LUMBDA_BACKEND_CMD" \
setsid nohup "$LUMBDA_FACTORY_DIR/bend-dispatcher.sh" \
"$QUEUE_DIR" "$QUEUE_DIR/results.tsv" "$POOL_BATCHES" \
>> "$QUEUE_DIR/dispatcher.log" 2>&1 < /dev/null &
}
# DLQ runner — auto-resolves transient DLQ classes (missing-bin,
# bisect-pool-race, transient cuda-error, no-portal) and escalates the
# persistent classes to rdlq/.
DLQ_RUNNER_STALE_SEC="${DLQ_RUNNER_STALE_SEC:-180}"
dlq_runner_alive() {
heartbeat_is_alive "$QUEUE_DIR" "dlq-runner" "$DLQ_RUNNER_STALE_SEC"
}
restart_dlq_runner() {
log "restart bend-supervisor-dlq-runner"
LUMBDA_QUEUE_DIR="$QUEUE_DIR" \
LUMBDA_FACTORY_DIR="$LUMBDA_FACTORY_DIR" \
setsid nohup "$LUMBDA_FACTORY_DIR/bend-supervisor-dlq-runner.sh" \
>> "$QUEUE_DIR/dlq-runner.log" 2>&1 < /dev/null &
}
# ── hot-reload tunable settings ─────────────────────────────────────
# Lets operators bump POOL_MAX, BEND_ENDPOINTS, POOL_BATCHES live without
# restarting the supervisor. Drop a KEY=VALUE file at $QUEUE_DIR/supervisor.config;
# next poll iter reads it, kills affected child (pool / dispatcher),
# normal restart loop respawns with new value.
CONFIG_FILE="$QUEUE_DIR/supervisor.config"
# Fingerprint TIER_* keys so hot-reload knows when any per-tier cap or
# threshold changed. Compares concatenated env values, not file mtime, so
# unrelated edits (POOL_MAX bump alone) don't churn the dispatcher.
tier_fingerprint() {
printf '%s' \
"${TIER_MICRO_MIN:-}|${TIER_MICRO_MAX:-}|${TIER_SMALL_MAX:-}|${TIER_MEDIUM_MAX:-}|${TIER_LARGE_MAX:-}|" \
"${TIER_MICRO_POOL_MAX:-}|${TIER_MICRO_DISPATCH:-}|" \
"${TIER_SMALL_POOL_MAX:-}|${TIER_SMALL_DISPATCH:-}|" \
"${TIER_MEDIUM_POOL_MAX:-}|${TIER_MEDIUM_DISPATCH:-}|" \
"${TIER_LARGE_POOL_MAX:-}|${TIER_LARGE_DISPATCH:-}"
}
reload_config() {
[ -f "$CONFIG_FILE" ] || return 0
local prev_pool_max="$POOL_MAX"
local prev_endpoints="$BEND_ENDPOINTS"
local prev_batches="$POOL_BATCHES"
local prev_pool_min_bin="$POOL_MIN_BIN"
local prev_tier_fp
prev_tier_fp=$(tier_fingerprint)
# `set -a` exports every var assigned during the source so children
# (pool, dispatcher) inherit them. Without this, TIER_* + POOL_* set
# in supervisor.config stay local to this shell + dispatcher respawn
# would pick up lib-tier defaults instead of operator overrides.
set -a
# shellcheck source=/dev/null
. "$CONFIG_FILE"
set +a
local new_tier_fp
new_tier_fp=$(tier_fingerprint)
if [ "$POOL_MAX" != "$prev_pool_max" ] || [ "$POOL_MIN_BIN" != "$prev_pool_min_bin" ]; then
log "hot-reload POOL_MAX=$POOL_MAX POOL_MIN_BIN=$POOL_MIN_BIN (was MAX=$prev_pool_max MIN=$prev_pool_min_bin); killing pool, restart loop respawns"
pkill -f "bend-emit-pool.sh" 2>/dev/null
fi
# Dispatcher kill ONLY on endpoint / batches changes — tier-only
# changes flow live into workers via tier_reload_caps_from_file in
# lib-tier.sh.
if [ "$BEND_ENDPOINTS" != "$prev_endpoints" ] || [ "$POOL_BATCHES" != "$prev_batches" ]; then
log "hot-reload dispatcher config changed (killing dispatcher, restart loop respawns)"
pkill -f "bend-dispatcher.sh" 2>/dev/null
elif [ "$new_tier_fp" != "$prev_tier_fp" ]; then
log "hot-reload tier policy changed (live-reload — no dispatcher kill, workers read tier caps live)"
fi
}
# ── main loop ───────────────────────────────────────────────────────
while true; do
# Heartbeat first thing — supervisor must prove its own liveness too.
# factory-status reads supervisor.heartbeat to know whether a host
# has its own caretaker.
heartbeat_touch "$QUEUE_DIR" "supervisor"
reload_config
if [ -f "$STOP_FILE" ]; then
log "STOP marker seen — exiting"
rm -f "$STOP_FILE"
SUP_EXIT_REASON="stop-marker"
break
fi
# backend first — pool + dispatcher both need a live backend.
if ! bend_port_listening; then
log "backend DOWN on :${BEND_PORT_LISTEN}"
restart_bend
fi
# pool
if ! pool_alive; then
n_unworked=$(count_unworked_lsp)
if [ "$n_unworked" -gt 0 ]; then
log "pool DEAD with $n_unworked un-emitted .lsp"
restart_pool
else
log "pool DEAD but queue clean — nothing to restart for"
fi
fi
# dispatcher
if ! dispatcher_any_alive; then
n_dispatchable=$(count_dispatchable)
if [ "$n_dispatchable" -gt 0 ]; then
log "dispatcher DEAD with $n_dispatchable .ready+.bin waiting"
restart_dispatcher
else
log "dispatcher DEAD but no .ready / .bin — nothing to restart for"
fi
fi
# DLQ runner — auto-resolve or escalate. Always restart when stale;
# cheap to run idle (scans an empty dlq/ in milliseconds).
if ! dlq_runner_alive; then
log "dlq-runner DEAD (no fresh heartbeat) — restarting"
restart_dlq_runner
fi
sleep "$POLL_SECONDS"
done