forked from cerc-io/stack-orchestrator
feat: graceful shutdown, ZFS upgrade, storage migration, sync-tools build
- entrypoint.py: Python stays PID 1, traps SIGTERM, requests graceful exit via admin RPC (agave-validator exit --force) before falling back to signals - snapshot_download.py: fix break-on-failure bug in incremental download loop (continue + re-probe instead of giving up) - biscayne-upgrade-zfs.yml: upgrade ZFS 2.2.2 → 2.2.9 via arter97/zfs-lts PPA to fix io_uring deadlock at kernel module level - biscayne-migrate-storage.yml: one-time migration from zvol/XFS to ZFS dataset (zvol workaround no longer needed with graceful shutdown + ZFS fix) - biscayne-stop.yml: patch terminationGracePeriodSeconds to 300 before scaling to 0, updated docs for admin RPC shutdown - biscayne-sync-tools.yml: fix SSH agent forwarding (vars: ansible_become), add --tags build-container support, add set -e to shell blocks - biscayne-recover.yml: updated for graceful shutdown awareness - check-status.py: add --pane flag for tmux, clean redraw in watch mode - CLAUDE.md: update docs for ZFS dataset storage, graceful shutdown Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>
This commit is contained in:
co-authored by
Claude Opus 4.6
parent
173b807451
commit
b88af2be70
@@ -2,12 +2,17 @@
|
||||
"""Agave validator entrypoint — snapshot management, arg construction, liveness probe.
|
||||
|
||||
Two subcommands:
|
||||
entrypoint.py serve (default) — snapshot freshness check + exec agave-validator
|
||||
entrypoint.py serve (default) — snapshot freshness check + run agave-validator
|
||||
entrypoint.py probe — liveness probe (slot lag check, exits 0/1)
|
||||
|
||||
Replaces the bash entrypoint.sh / start-rpc.sh / start-validator.sh with a single
|
||||
Python module. Test mode still dispatches to start-test.sh.
|
||||
|
||||
Python stays as PID 1 and traps SIGTERM. On SIGTERM, it runs
|
||||
``agave-validator exit --force --ledger /data/ledger`` which connects to the
|
||||
admin RPC Unix socket and tells the validator to flush I/O and exit cleanly.
|
||||
This avoids the io_uring/ZFS deadlock that occurs when the process is killed.
|
||||
|
||||
All configuration comes from environment variables — same vars as the original
|
||||
bash scripts. See compose files for defaults.
|
||||
"""
|
||||
@@ -18,8 +23,10 @@ import json
|
||||
import logging
|
||||
import os
|
||||
import re
|
||||
import signal
|
||||
import subprocess
|
||||
import sys
|
||||
import threading
|
||||
import time
|
||||
import urllib.error
|
||||
import urllib.request
|
||||
@@ -365,11 +372,77 @@ def append_extra_args(args: list[str]) -> list[str]:
|
||||
return args
|
||||
|
||||
|
||||
# -- Graceful shutdown --------------------------------------------------------
|
||||
|
||||
# Timeout for graceful exit via admin RPC. Leave 30s margin for k8s
|
||||
# terminationGracePeriodSeconds (300s).
|
||||
GRACEFUL_EXIT_TIMEOUT = 270
|
||||
|
||||
|
||||
def graceful_exit(child: subprocess.Popen[bytes]) -> None:
|
||||
"""Request graceful shutdown via the admin RPC Unix socket.
|
||||
|
||||
Runs ``agave-validator exit --force --ledger /data/ledger`` which connects
|
||||
to the admin RPC socket at ``/data/ledger/admin.rpc`` and sets the
|
||||
validator's exit flag. The validator flushes all I/O and exits cleanly,
|
||||
avoiding the io_uring/ZFS deadlock.
|
||||
|
||||
If the admin RPC exit fails or the child doesn't exit within the timeout,
|
||||
falls back to SIGTERM then SIGKILL.
|
||||
"""
|
||||
log.info("SIGTERM received — requesting graceful exit via admin RPC")
|
||||
try:
|
||||
result = subprocess.run(
|
||||
["agave-validator", "exit", "--force", "--ledger", LEDGER_DIR],
|
||||
capture_output=True, text=True, timeout=30,
|
||||
)
|
||||
if result.returncode == 0:
|
||||
log.info("Admin RPC exit requested successfully")
|
||||
else:
|
||||
log.warning(
|
||||
"Admin RPC exit returned %d: %s",
|
||||
result.returncode, result.stderr.strip(),
|
||||
)
|
||||
except subprocess.TimeoutExpired:
|
||||
log.warning("Admin RPC exit command timed out after 30s")
|
||||
except FileNotFoundError:
|
||||
log.warning("agave-validator binary not found for exit command")
|
||||
|
||||
# Wait for child to exit
|
||||
try:
|
||||
child.wait(timeout=GRACEFUL_EXIT_TIMEOUT)
|
||||
log.info("Validator exited cleanly with code %d", child.returncode)
|
||||
return
|
||||
except subprocess.TimeoutExpired:
|
||||
log.warning(
|
||||
"Validator did not exit within %ds — sending SIGTERM",
|
||||
GRACEFUL_EXIT_TIMEOUT,
|
||||
)
|
||||
|
||||
# Fallback: SIGTERM
|
||||
child.terminate()
|
||||
try:
|
||||
child.wait(timeout=15)
|
||||
log.info("Validator exited after SIGTERM with code %d", child.returncode)
|
||||
return
|
||||
except subprocess.TimeoutExpired:
|
||||
log.warning("Validator did not exit after SIGTERM — sending SIGKILL")
|
||||
|
||||
# Last resort: SIGKILL
|
||||
child.kill()
|
||||
child.wait()
|
||||
log.info("Validator killed with SIGKILL, code %d", child.returncode)
|
||||
|
||||
|
||||
# -- Serve subcommand ---------------------------------------------------------
|
||||
|
||||
|
||||
def cmd_serve() -> None:
|
||||
"""Main serve flow: snapshot check, setup, exec agave-validator."""
|
||||
"""Main serve flow: snapshot check, setup, run agave-validator as child.
|
||||
|
||||
Python stays as PID 1 and traps SIGTERM to perform graceful shutdown
|
||||
via the admin RPC Unix socket.
|
||||
"""
|
||||
mode = env("AGAVE_MODE", "test")
|
||||
log.info("AGAVE_MODE=%s", mode)
|
||||
|
||||
@@ -407,7 +480,21 @@ def cmd_serve() -> None:
|
||||
Path("/tmp/entrypoint-start").write_text(str(time.time()))
|
||||
|
||||
log.info("Starting agave-validator with %d arguments", len(args))
|
||||
os.execvp("agave-validator", ["agave-validator"] + args)
|
||||
child = subprocess.Popen(["agave-validator"] + args)
|
||||
|
||||
# Forward SIGUSR1 to child (log rotation)
|
||||
signal.signal(signal.SIGUSR1, lambda _sig, _frame: child.send_signal(signal.SIGUSR1))
|
||||
|
||||
# Trap SIGTERM — run graceful_exit in a thread so the signal handler returns
|
||||
# immediately and child.wait() in the main thread can observe the exit.
|
||||
def _on_sigterm(_sig: int, _frame: object) -> None:
|
||||
threading.Thread(target=graceful_exit, args=(child,), daemon=True).start()
|
||||
|
||||
signal.signal(signal.SIGTERM, _on_sigterm)
|
||||
|
||||
# Wait for child — if it exits on its own (crash, normal exit), propagate code
|
||||
child.wait()
|
||||
sys.exit(child.returncode)
|
||||
|
||||
|
||||
# -- Probe subcommand ---------------------------------------------------------
|
||||
|
||||
@@ -655,8 +655,9 @@ def download_best_snapshot(
|
||||
log.info("Downloading incremental %s (%d mirrors, slot %d, gap %d slots)",
|
||||
inc_fn, len(inc_mirrors), inc_slot, gap)
|
||||
if not download_aria2c(inc_mirrors, output_dir, inc_fn, connections):
|
||||
log.error("Failed to download incremental %s", inc_fn)
|
||||
break
|
||||
log.warning("Failed to download incremental %s — re-probing in 10s", inc_fn)
|
||||
time.sleep(10)
|
||||
continue
|
||||
|
||||
prev_inc_filename = inc_fn
|
||||
|
||||
|
||||
+22
-3
@@ -18,6 +18,7 @@ from __future__ import annotations
|
||||
|
||||
import argparse
|
||||
import json
|
||||
import os
|
||||
import subprocess
|
||||
import sys
|
||||
import time
|
||||
@@ -206,9 +207,11 @@ def display(iteration: int = 0) -> None:
|
||||
snapshots = check_snapshots()
|
||||
ramdisk = check_ramdisk()
|
||||
|
||||
print(f"\n{'=' * 60}")
|
||||
print(f" Biscayne Agave Status — {ts}")
|
||||
print(f"{'=' * 60}")
|
||||
# Clear screen and home cursor for clean redraw in watch mode
|
||||
if iteration > 0:
|
||||
print("\033[2J\033[H", end="")
|
||||
|
||||
print(f"\n Biscayne Agave Status — {ts}\n")
|
||||
|
||||
# Pod
|
||||
print(f"\n Pod: {pod['phase']}")
|
||||
@@ -275,14 +278,30 @@ def display(iteration: int = 0) -> None:
|
||||
# -- Main ---------------------------------------------------------------------
|
||||
|
||||
|
||||
def spawn_tmux_pane(interval: int) -> None:
|
||||
"""Launch this script with --watch in a new tmux pane."""
|
||||
script = os.path.abspath(__file__)
|
||||
cmd = f"python3 {script} --watch -i {interval}"
|
||||
subprocess.run(
|
||||
["tmux", "split-window", "-h", "-d", cmd],
|
||||
check=True,
|
||||
)
|
||||
|
||||
|
||||
def main() -> int:
|
||||
p = argparse.ArgumentParser(description=__doc__,
|
||||
formatter_class=argparse.RawDescriptionHelpFormatter)
|
||||
p.add_argument("--watch", action="store_true", help="Repeat every interval")
|
||||
p.add_argument("--pane", action="store_true",
|
||||
help="Launch --watch in a new tmux pane")
|
||||
p.add_argument("-i", "--interval", type=int, default=30,
|
||||
help="Watch interval in seconds (default: 30)")
|
||||
args = p.parse_args()
|
||||
|
||||
if args.pane:
|
||||
spawn_tmux_pane(args.interval)
|
||||
return 0
|
||||
|
||||
discover()
|
||||
|
||||
try:
|
||||
|
||||
Reference in New Issue
Block a user