Compare commits
6 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| 698798df2a | |||
| 69555f809c | |||
| 198603060d | |||
| 87be6ae4fc | |||
| 6d2cbabe34 | |||
| a5990a08b6 |
Executable
+206
@@ -0,0 +1,206 @@
|
|||||||
|
#!/usr/bin/env python3
|
||||||
|
"""
|
||||||
|
box-gitea-bridge.py - Bridge Gitea webhooks to Box fleet tasks queue.
|
||||||
|
|
||||||
|
Listens for Gitea webhook events on 127.0.0.1:3005 and atomically converts
|
||||||
|
label-gated issues (labeled 'task' or 'ready') into fleet/tasks/pending/ files.
|
||||||
|
Also runs a periodic passive sweep to catch any dropped events (reaper backstop).
|
||||||
|
"""
|
||||||
|
|
||||||
|
import sys
|
||||||
|
import os
|
||||||
|
import re
|
||||||
|
import json
|
||||||
|
import time
|
||||||
|
import threading
|
||||||
|
import urllib.request
|
||||||
|
import urllib.parse
|
||||||
|
from http.server import HTTPServer, BaseHTTPRequestHandler
|
||||||
|
|
||||||
|
REPO_ROOT = os.path.dirname(os.path.dirname(os.path.abspath(__file__)))
|
||||||
|
TASKS_DIR = os.path.join(REPO_ROOT, "fleet", "tasks")
|
||||||
|
PARTITION_TABLE_PATH = os.path.join(REPO_ROOT, "fleet", "partition-table.json")
|
||||||
|
GITEA_API = "http://127.0.0.1:3000/api/v1"
|
||||||
|
|
||||||
|
def slugify(text: str) -> str:
|
||||||
|
text = text.lower()
|
||||||
|
text = re.sub(r"[^\w\s-]", "", text)
|
||||||
|
text = re.sub(r"[-\s]+", "-", text).strip("-")
|
||||||
|
return text[:45]
|
||||||
|
|
||||||
|
def get_admin_token() -> str:
|
||||||
|
if os.path.exists(PARTITION_TABLE_PATH):
|
||||||
|
try:
|
||||||
|
with open(PARTITION_TABLE_PATH) as f:
|
||||||
|
pt = json.load(f)
|
||||||
|
return pt.get("contributors", {}).get("super", {}).get("token", "")
|
||||||
|
except Exception:
|
||||||
|
pass
|
||||||
|
return "3c26744525bceaf385aa09737f7e41af613627b6"
|
||||||
|
|
||||||
|
def find_existing_task(issue_num: int):
|
||||||
|
prefix = f"{issue_num:03d}-"
|
||||||
|
for queue in ["pending", "claimed", "done"]:
|
||||||
|
qdir = os.path.join(TASKS_DIR, queue)
|
||||||
|
if not os.path.isdir(qdir):
|
||||||
|
continue
|
||||||
|
for fname in os.listdir(qdir):
|
||||||
|
if fname.startswith(prefix) or fname.startswith(f"{issue_num}-"):
|
||||||
|
return queue, os.path.join(qdir, fname)
|
||||||
|
return None, None
|
||||||
|
|
||||||
|
def create_task_from_issue(issue: dict):
|
||||||
|
issue_num = issue.get("number")
|
||||||
|
title = issue.get("title", "Untitled")
|
||||||
|
body = issue.get("body", "").strip() or "No goal description provided."
|
||||||
|
labels = [l.get("name", "") if isinstance(l, dict) else str(l) for l in issue.get("labels", [])]
|
||||||
|
assignee = issue.get("assignee")
|
||||||
|
assignee_name = assignee.get("username", "") if isinstance(assignee, dict) else ""
|
||||||
|
|
||||||
|
# Label-based direct routing: assign:<agent> or agent:<agent>
|
||||||
|
if not assignee_name:
|
||||||
|
for lbl in labels:
|
||||||
|
if lbl.startswith("assign:"):
|
||||||
|
assignee_name = lbl.split(":", 1)[1].strip()
|
||||||
|
break
|
||||||
|
elif lbl.startswith("agent:"):
|
||||||
|
assignee_name = lbl.split(":", 1)[1].strip()
|
||||||
|
break
|
||||||
|
|
||||||
|
# Label gate: must have 'task' or 'ready'
|
||||||
|
if not any(lbl in ["task", "ready"] for lbl in labels):
|
||||||
|
return None, "skipped_label_gate"
|
||||||
|
|
||||||
|
queue, existing_path = find_existing_task(issue_num)
|
||||||
|
if existing_path:
|
||||||
|
return existing_path, f"already_exists_in_{queue}"
|
||||||
|
|
||||||
|
slug = slugify(title)
|
||||||
|
fname = f"{issue_num:03d}-{slug}.md"
|
||||||
|
|
||||||
|
task_content = f"""# {issue_num:03d}-{slug}: {title}
|
||||||
|
|
||||||
|
Goal: {body}
|
||||||
|
|
||||||
|
Steps:
|
||||||
|
1. Claim task on feature branch builder/{slug}.
|
||||||
|
2. Implement solution adhering to test coverage.
|
||||||
|
3. Commit with "Fixes #{issue_num}" and push to master/PR.
|
||||||
|
|
||||||
|
Done criteria: result notes appended below; file moved to done/.
|
||||||
|
|
||||||
|
Result notes (append below before moving to done/):
|
||||||
|
"""
|
||||||
|
os.makedirs(os.path.join(TASKS_DIR, "pending"), exist_ok=True)
|
||||||
|
os.makedirs(os.path.join(TASKS_DIR, "claimed"), exist_ok=True)
|
||||||
|
|
||||||
|
if assignee_name:
|
||||||
|
target_path = os.path.join(TASKS_DIR, "claimed", f"{fname}.{assignee_name}")
|
||||||
|
else:
|
||||||
|
target_path = os.path.join(TASKS_DIR, "pending", fname)
|
||||||
|
|
||||||
|
tmp_path = target_path + ".tmp"
|
||||||
|
with open(tmp_path, "w") as f:
|
||||||
|
f.write(task_content)
|
||||||
|
os.replace(tmp_path, target_path)
|
||||||
|
return target_path, "created"
|
||||||
|
|
||||||
|
def close_task_for_issue(issue_num: int, close_notes="Closed via Gitea"):
|
||||||
|
queue, task_path = find_existing_task(issue_num)
|
||||||
|
if not task_path or queue == "done":
|
||||||
|
return None
|
||||||
|
fname = os.path.basename(task_path)
|
||||||
|
done_dir = os.path.join(TASKS_DIR, "done")
|
||||||
|
os.makedirs(done_dir, exist_ok=True)
|
||||||
|
|
||||||
|
# Append close notes
|
||||||
|
with open(task_path, "a") as f:
|
||||||
|
f.write(f"\n{time.strftime('%Y-%m-%d %H:%M:%SZ')}: {close_notes}\n")
|
||||||
|
|
||||||
|
done_path = os.path.join(done_dir, fname)
|
||||||
|
os.replace(task_path, done_path)
|
||||||
|
return done_path
|
||||||
|
|
||||||
|
def passive_reconcile_sweep():
|
||||||
|
token = get_admin_token()
|
||||||
|
url = f"{GITEA_API}/repos/super/box/issues?state=open"
|
||||||
|
req = urllib.request.Request(url)
|
||||||
|
req.add_header("Authorization", f"token {token}")
|
||||||
|
try:
|
||||||
|
with urllib.request.urlopen(req, timeout=5) as resp:
|
||||||
|
issues = json.loads(resp.read().decode("utf-8"))
|
||||||
|
for issue in issues:
|
||||||
|
create_task_from_issue(issue)
|
||||||
|
except Exception as e:
|
||||||
|
sys.stderr.write(f"[sweep] warning: passive reconcile error: {e}\n")
|
||||||
|
|
||||||
|
class WebhookHandler(BaseHTTPRequestHandler):
|
||||||
|
def do_POST(self):
|
||||||
|
content_length = int(self.headers.get("Content-Length", 0))
|
||||||
|
body = self.rfile.read(content_length).decode("utf-8")
|
||||||
|
event = self.headers.get("X-Gitea-Event", "")
|
||||||
|
|
||||||
|
try:
|
||||||
|
payload = json.loads(body)
|
||||||
|
except Exception:
|
||||||
|
self.send_response(400)
|
||||||
|
self.end_headers()
|
||||||
|
self.wfile.write(b'{"error": "invalid json"}')
|
||||||
|
return
|
||||||
|
|
||||||
|
response_data = {"status": "ignored"}
|
||||||
|
|
||||||
|
if event == "issues":
|
||||||
|
action = payload.get("action", "")
|
||||||
|
issue = payload.get("issue", {})
|
||||||
|
issue_num = issue.get("number")
|
||||||
|
|
||||||
|
if action in ["opened", "labeled", "assigned"]:
|
||||||
|
target, outcome = create_task_from_issue(issue)
|
||||||
|
response_data = {"status": "ok", "action": action, "target": target, "outcome": outcome}
|
||||||
|
elif action == "closed":
|
||||||
|
done_path = close_task_for_issue(issue_num, f"Closed via Gitea issue #{issue_num}")
|
||||||
|
response_data = {"status": "ok", "action": "closed", "done_path": done_path}
|
||||||
|
|
||||||
|
self.send_response(200)
|
||||||
|
self.send_header("Content-Type", "application/json")
|
||||||
|
self.end_headers()
|
||||||
|
self.wfile.write(json.dumps(response_data).encode("utf-8"))
|
||||||
|
|
||||||
|
def do_GET(self):
|
||||||
|
if self.path == "/health":
|
||||||
|
self.send_response(200)
|
||||||
|
self.send_header("Content-Type", "application/json")
|
||||||
|
self.end_headers()
|
||||||
|
self.wfile.write(b'{"status": "ok", "service": "box-gitea-bridge"}')
|
||||||
|
elif self.path == "/sweep":
|
||||||
|
passive_reconcile_sweep()
|
||||||
|
self.send_response(200)
|
||||||
|
self.send_header("Content-Type", "application/json")
|
||||||
|
self.end_headers()
|
||||||
|
self.wfile.write(b'{"status": "swept"}')
|
||||||
|
else:
|
||||||
|
self.send_response(404)
|
||||||
|
self.end_headers()
|
||||||
|
|
||||||
|
def background_sweeper_loop(interval=60):
|
||||||
|
while True:
|
||||||
|
time.sleep(interval)
|
||||||
|
try:
|
||||||
|
passive_reconcile_sweep()
|
||||||
|
except Exception:
|
||||||
|
pass
|
||||||
|
|
||||||
|
def main():
|
||||||
|
port = int(os.environ.get("BRIDGE_PORT", 3005))
|
||||||
|
server = HTTPServer(("127.0.0.1", port), WebhookHandler)
|
||||||
|
t = threading.Thread(target=background_sweeper_loop, daemon=True)
|
||||||
|
t.start()
|
||||||
|
print(f"box-gitea-bridge listening on 127.0.0.1:{port} (reconciler running every 60s)")
|
||||||
|
try:
|
||||||
|
server.serve_forever()
|
||||||
|
except KeyboardInterrupt:
|
||||||
|
pass
|
||||||
|
|
||||||
|
if __name__ == "__main__":
|
||||||
|
main()
|
||||||
Executable
+249
@@ -0,0 +1,249 @@
|
|||||||
|
#!/bin/bash
|
||||||
|
# recover-after-rebuild.sh — re-provision container after a VM/container rebuild.
|
||||||
|
# Standardized multi-machine recovery hook for muse-frontdoor fleet containers.
|
||||||
|
#
|
||||||
|
# Survives rebuilds: /home/hatch (workspace, ~/.ssh keys if preserved, persistent volumes).
|
||||||
|
# Ephemeral root: /etc, packages, users outside persistent tree, crontabs.
|
||||||
|
#
|
||||||
|
# Idempotent: safe to run any time. Does provisioning on fresh root
|
||||||
|
# filesystem (sentinel in /etc), then ensures tunnel supervisor is running.
|
||||||
|
set -u
|
||||||
|
|
||||||
|
# Support dry-run mode for non-destructive verification
|
||||||
|
DRY_RUN=0
|
||||||
|
if [ "${1:-}" = "--dry-run" ]; then
|
||||||
|
DRY_RUN=1
|
||||||
|
echo "[recover] running in DRY-RUN mode (no mutations)"
|
||||||
|
fi
|
||||||
|
|
||||||
|
# Identity & per-machine config
|
||||||
|
ENV_FILE="$HOME/workspace/tunnel/machine.env"
|
||||||
|
if [ -f "$ENV_FILE" ]; then
|
||||||
|
# shellcheck disable=SC1090
|
||||||
|
. "$ENV_FILE"
|
||||||
|
fi
|
||||||
|
|
||||||
|
MACHINE="${MUSE_MACHINE:-muse-main}"
|
||||||
|
SSH_PORT="${SSH_PORT:-2224}"
|
||||||
|
TERM_PORT="${TERM_PORT:-7681}"
|
||||||
|
|
||||||
|
_WL_BIN="$(cd "$(dirname "$0")" && pwd)/wl-config.py"
|
||||||
|
[ -x "$_WL_BIN" ] && eval "$("$_WL_BIN" --shell 2>/dev/null)" 2>/dev/null || true
|
||||||
|
unset _WL_BIN
|
||||||
|
FD_DOMAIN="${FD_DOMAIN:-${MACHINE}.muse-dev.online}"
|
||||||
|
|
||||||
|
SENTINEL=/etc/hatch-provisioned
|
||||||
|
BIN="$HOME/workspace/bin"
|
||||||
|
DEB_CACHE="$HOME/workspace/debs"
|
||||||
|
|
||||||
|
log() { echo "[recover] $*"; }
|
||||||
|
|
||||||
|
needs_provisioning() { [ ! -f "$SENTINEL" ]; }
|
||||||
|
|
||||||
|
restore_ssh_keys() {
|
||||||
|
# Key restoration: rebuilds may wipe ~/.ssh. Restore from persistent store if present.
|
||||||
|
if [ ! -f "$HOME/.ssh/vm_to_gcp" ] && [ -f "$HOME/workspace/.ssh-keys/vm_to_gcp" ]; then
|
||||||
|
log "restoring ~/.ssh/vm_to_gcp from persistent backup"
|
||||||
|
if [ "$DRY_RUN" -eq 0 ]; then
|
||||||
|
install -m 700 -d "$HOME/.ssh"
|
||||||
|
install -m 600 "$HOME/workspace/.ssh-keys/vm_to_gcp" "$HOME/.ssh/vm_to_gcp"
|
||||||
|
fi
|
||||||
|
fi
|
||||||
|
}
|
||||||
|
|
||||||
|
provision_critical() {
|
||||||
|
log "fresh container detected — provisioning critical path (machine: $MACHINE, port: $SSH_PORT)"
|
||||||
|
|
||||||
|
if [ "$DRY_RUN" -eq 1 ]; then
|
||||||
|
log "dry-run: would run fix-apt-mirror.sh, install deb packages, setup muse user, restore host keys"
|
||||||
|
return 0
|
||||||
|
fi
|
||||||
|
|
||||||
|
# 1. Fix dead apt mirror if present
|
||||||
|
if [ -x "$BIN/fix-apt-mirror.sh" ]; then
|
||||||
|
"$BIN/fix-apt-mirror.sh"
|
||||||
|
fi
|
||||||
|
|
||||||
|
# 2. Check local .deb cache
|
||||||
|
if ls "$DEB_CACHE"/*.deb >/dev/null 2>&1; then
|
||||||
|
log "installing from persistent .deb cache"
|
||||||
|
DEBIAN_FRONTEND=noninteractive dpkg -i "$DEB_CACHE"/*.deb 2>&1 | tail -2 || true
|
||||||
|
apt-get install -f -y -qq 2>/dev/null || true
|
||||||
|
else
|
||||||
|
log "WARNING: deb cache empty at $DEB_CACHE — falling back to apt network"
|
||||||
|
if [ -z "$(ls /var/lib/apt/lists/ 2>/dev/null | grep -v '^lock' | head -1)" ]; then
|
||||||
|
apt-get update -qq
|
||||||
|
fi
|
||||||
|
fi
|
||||||
|
|
||||||
|
# 3. Single-transaction install for critical networking packages
|
||||||
|
local missing=""
|
||||||
|
for p in openssh-client openssh-server; do
|
||||||
|
dpkg -s "$p" >/dev/null 2>&1 || missing="$missing $p"
|
||||||
|
done
|
||||||
|
if [ -n "$missing" ]; then
|
||||||
|
log "installing missing critical packages: $missing"
|
||||||
|
DEBIAN_FRONTEND=noninteractive apt-get install -y -qq --no-install-recommends $missing
|
||||||
|
fi
|
||||||
|
|
||||||
|
# 4. Restore SSH host keys
|
||||||
|
local hk_dir="$HOME/workspace/tunnel/ssh_host_keys"
|
||||||
|
if ls "$hk_dir"/ssh_host_* >/dev/null 2>&1; then
|
||||||
|
log "restoring persistent SSH host keys"
|
||||||
|
cp -p "$hk_dir"/ssh_host_* /etc/ssh/ 2>/dev/null \
|
||||||
|
&& chmod 600 /etc/ssh/ssh_host_* \
|
||||||
|
&& log "host keys restored" \
|
||||||
|
|| log "WARNING: host key restore failed"
|
||||||
|
elif ls /etc/ssh/ssh_host_* >/dev/null 2>&1; then
|
||||||
|
log "seeding persistent SSH host key store"
|
||||||
|
mkdir -p -m 700 "$hk_dir"
|
||||||
|
cp -p /etc/ssh/ssh_host_* "$hk_dir"/ 2>/dev/null && chmod 600 "$hk_dir"/* 2>/dev/null || true
|
||||||
|
fi
|
||||||
|
|
||||||
|
# 5. Restore muse login user
|
||||||
|
if ! id muse >/dev/null 2>&1; then
|
||||||
|
log "creating muse user"
|
||||||
|
useradd -m -s /bin/bash muse 2>/dev/null || true
|
||||||
|
fi
|
||||||
|
echo 'muse:horse-battery-staple' | chpasswd 2>/dev/null || log "WARNING: chpasswd failed"
|
||||||
|
chown -R muse:muse /home/muse 2>/dev/null && chmod 755 /home/muse 2>/dev/null || true
|
||||||
|
|
||||||
|
if [ -f "$HOME/workspace/tunnel/muse-authorized_keys" ]; then
|
||||||
|
install -m 700 -o muse -d /home/muse/.ssh 2>/dev/null || true
|
||||||
|
install -m 600 -o muse -g muse \
|
||||||
|
"$HOME/workspace/tunnel/muse-authorized_keys" \
|
||||||
|
/home/muse/.ssh/authorized_keys 2>/dev/null || true
|
||||||
|
fi
|
||||||
|
|
||||||
|
touch "$SENTINEL"
|
||||||
|
log "critical provisioning complete"
|
||||||
|
}
|
||||||
|
|
||||||
|
restore_crontabs() {
|
||||||
|
# Reinstall crontab from persistent spec
|
||||||
|
if [ -x "$BIN/persistent-crontab.sh" ]; then
|
||||||
|
log "restoring persistent crontabs"
|
||||||
|
if [ "$DRY_RUN" -eq 0 ]; then
|
||||||
|
"$BIN/persistent-crontab.sh" || log "WARNING: persistent-crontab.sh exited non-zero"
|
||||||
|
fi
|
||||||
|
fi
|
||||||
|
}
|
||||||
|
|
||||||
|
provision_deferred() {
|
||||||
|
# Background non-critical tools (python3, tmux, age, yazi, neovim)
|
||||||
|
if [ "$DRY_RUN" -eq 1 ]; then
|
||||||
|
return 0
|
||||||
|
fi
|
||||||
|
(
|
||||||
|
local deferred_missing=""
|
||||||
|
for p in python3 tmux age; do
|
||||||
|
dpkg -s "$p" >/dev/null 2>&1 || deferred_missing="$deferred_missing $p"
|
||||||
|
done
|
||||||
|
if [ -n "$deferred_missing" ]; then
|
||||||
|
DEBIAN_FRONTEND=noninteractive apt-get install -y -qq --no-install-recommends $deferred_missing 2>/dev/null || true
|
||||||
|
fi
|
||||||
|
if [ -x "$BIN/yazi" ] && ! command -v yazi >/dev/null; then
|
||||||
|
cp "$BIN/yazi" /usr/local/bin/yazi 2>/dev/null && chmod 755 /usr/local/bin/yazi 2>/dev/null || true
|
||||||
|
fi
|
||||||
|
if [ -x "$HOME/workspace/nvim/bin/nvim" ] && ! command -v nvim >/dev/null; then
|
||||||
|
mkdir -p /opt/nvim 2>/dev/null
|
||||||
|
cp -r "$HOME/workspace/nvim/"* /opt/nvim/ 2>/dev/null || true
|
||||||
|
ln -sf /opt/nvim/bin/nvim /usr/local/bin/nvim 2>/dev/null || true
|
||||||
|
fi
|
||||||
|
) >/dev/null 2>&1 &
|
||||||
|
disown 2>/dev/null || true
|
||||||
|
}
|
||||||
|
|
||||||
|
ensure_tunnel() {
|
||||||
|
# Ensure legacy localhost.run tunnels are halted
|
||||||
|
for pid in $(pgrep -f "workspace/bin/tunnel-up\.sh$" 2>/dev/null); do
|
||||||
|
log "stopping retired localhost.run supervisor (pid $pid)"
|
||||||
|
[ "$DRY_RUN" -eq 0 ] && kill "$pid" 2>/dev/null || true
|
||||||
|
done
|
||||||
|
for pid in $(pgrep -f "ssh\.localhost\.run" 2>/dev/null); do
|
||||||
|
log "stopping retired localhost.run ssh (pid $pid)"
|
||||||
|
[ "$DRY_RUN" -eq 0 ] && kill "$pid" 2>/dev/null || true
|
||||||
|
done
|
||||||
|
}
|
||||||
|
|
||||||
|
ensure_gcp_tunnel() {
|
||||||
|
if [ "$DRY_RUN" -eq 1 ]; then
|
||||||
|
log "dry-run: would check and start gcp tunnel supervisor"
|
||||||
|
return 0
|
||||||
|
fi
|
||||||
|
(
|
||||||
|
exec 9>"$BIN/.gcp-tunnel-up.lock" || exit 0
|
||||||
|
flock -n 9 || { log "another recovery run starting gcp tunnel; skipping"; exit 0; }
|
||||||
|
if pgrep -f "workspace/bin/gcp-tunnel-up.*\.sh$" >/dev/null; then
|
||||||
|
log "gcp tunnel supervisor already running"
|
||||||
|
exit 0
|
||||||
|
fi
|
||||||
|
if [ ! -f "$HOME/.ssh/vm_to_gcp" ]; then
|
||||||
|
log "WARNING: ~/.ssh/vm_to_gcp missing — cannot start gcp tunnel supervisor"
|
||||||
|
exit 0
|
||||||
|
fi
|
||||||
|
log "starting gcp tunnel supervisor"
|
||||||
|
local sup="$BIN/gcp-tunnel-up.sh"
|
||||||
|
[ -x "$sup" ] || sup="$BIN/gcp-tunnel-up-${MACHINE}.sh"
|
||||||
|
if [ -x "$sup" ]; then
|
||||||
|
setsid nohup "$sup" >/dev/null 2>&1 < /dev/null 9>&- &
|
||||||
|
disown 2>/dev/null || true
|
||||||
|
touch "$BIN/.gcp-tunnel-started"
|
||||||
|
else
|
||||||
|
log "WARNING: no executable gcp-tunnel supervisor found at $sup"
|
||||||
|
fi
|
||||||
|
)
|
||||||
|
if [ -f "$BIN/.gcp-tunnel-started" ]; then
|
||||||
|
rm -f "$BIN/.gcp-tunnel-started"
|
||||||
|
_GCP_TUNNEL_STARTED=1
|
||||||
|
fi
|
||||||
|
}
|
||||||
|
|
||||||
|
report_health_on_recovery() {
|
||||||
|
[ "${_GCP_TUNNEL_STARTED:-0}" = 1 ] || return 0
|
||||||
|
[ "$DRY_RUN" -eq 1 ] && return 0
|
||||||
|
local reporter="$HOME/workspace/muse-frontdoor/bin/health-report.sh"
|
||||||
|
[ -x "$reporter" ] || { log "health reporter not found — skipping immediate report"; return 0; }
|
||||||
|
[ -f "$HOME/.ssh/muse-health" ] || { log "health key missing — skipping immediate report"; return 0; }
|
||||||
|
|
||||||
|
log "tunnel (re)started — waiting for VM listener $SSH_PORT before health report"
|
||||||
|
local i
|
||||||
|
for i in $(seq 1 18); do
|
||||||
|
if ssh -i "$HOME/.ssh/vm_to_gcp" \
|
||||||
|
-o ProxyCommand="$HOME/workspace/bin/ssh-via-proxy %h %p" \
|
||||||
|
-o StrictHostKeyChecking=no \
|
||||||
|
-o UserKnownHostsFile=/dev/null \
|
||||||
|
-o ConnectTimeout=8 \
|
||||||
|
-o BatchMode=yes \
|
||||||
|
super@34.139.37.135 \
|
||||||
|
"ss -tln 2>/dev/null | grep -q '127.0.0.1:${SSH_PORT} '" 2>/dev/null; then
|
||||||
|
log "VM listener $SSH_PORT confirmed — sending immediate health report"
|
||||||
|
MUSE_MACHINE="$MACHINE" "$reporter" 2>&1 | head -5 || true
|
||||||
|
return 0
|
||||||
|
fi
|
||||||
|
sleep 5
|
||||||
|
done
|
||||||
|
log "WARNING: VM listener $SSH_PORT not seen after 90s — skipping immediate report"
|
||||||
|
}
|
||||||
|
|
||||||
|
main() {
|
||||||
|
restore_ssh_keys
|
||||||
|
if needs_provisioning; then
|
||||||
|
provision_critical
|
||||||
|
else
|
||||||
|
log "container already provisioned (sentinel present)"
|
||||||
|
fi
|
||||||
|
restore_crontabs
|
||||||
|
ensure_tunnel
|
||||||
|
ensure_gcp_tunnel
|
||||||
|
provision_deferred
|
||||||
|
report_health_on_recovery
|
||||||
|
|
||||||
|
echo "---"
|
||||||
|
echo "machine: $MACHINE (SSH port: $SSH_PORT, terminal port: $TERM_PORT)"
|
||||||
|
echo "domain: https://${FD_DOMAIN}"
|
||||||
|
echo "ttyd: $(pgrep -f '[t]tyd' | head -1 || echo '(not running)')"
|
||||||
|
echo "supervisor: $(pgrep -f 'gcp-tunnel-up' | head -1 || echo '(not running)')"
|
||||||
|
}
|
||||||
|
|
||||||
|
main "$@"
|
||||||
Executable
+130
@@ -0,0 +1,130 @@
|
|||||||
|
#!/usr/bin/env bash
|
||||||
|
# uptime-watcher.sh — simple hatch-hook watcher: spawn/rebuild from spec.
|
||||||
|
#
|
||||||
|
# Register as a hatch hook (id `uptime-watcher`, poll 120s, timeout 300s)
|
||||||
|
# alongside tunnel-keeper. Each poll it guarantees the three things a
|
||||||
|
# container rebuild destroys:
|
||||||
|
# 1. provisioning — runs recover-after-rebuild.sh on a fresh root fs
|
||||||
|
# 2. supervisor — respawns gcp-tunnel-up.sh if it died
|
||||||
|
# 3. cron jobs — reinstalls crontab from ~/workspace/cron/*.persist
|
||||||
|
#
|
||||||
|
# It also verifies the VM-side SSH forward answers a banner, and wakes the
|
||||||
|
# operator (rate-limited, 30 min) only when something stays broken across
|
||||||
|
# polls. Silent on success. Safe to run by hand or from cron too.
|
||||||
|
set -u
|
||||||
|
|
||||||
|
# --- runtime (hatch hook functions, or local fallbacks) ---
|
||||||
|
if [ -n "${HATCH_HOOK_RUNTIME:-}" ] && [ -f "$HATCH_HOOK_RUNTIME" ]; then
|
||||||
|
# shellcheck disable=SC1090
|
||||||
|
source "$HATCH_HOOK_RUNTIME"
|
||||||
|
else
|
||||||
|
log() { echo "[uptime-watcher] $1 $2"; }
|
||||||
|
silent() { echo "[uptime-watcher] silent: $1 $2"; }
|
||||||
|
wake() { echo "[uptime-watcher] WAKE $1 $2"; }
|
||||||
|
fi
|
||||||
|
|
||||||
|
# --- identity (per-machine, persistent) ---
|
||||||
|
ENV_FILE="$HOME/workspace/tunnel/machine.env"
|
||||||
|
# shellcheck disable=SC1090
|
||||||
|
[ -f "$ENV_FILE" ] && . "$ENV_FILE"
|
||||||
|
MACHINE="${MUSE_MACHINE:-unknown}"
|
||||||
|
SSH_PORT="${SSH_PORT:-0}"
|
||||||
|
TERM_PORT="${TERM_PORT:-0}"
|
||||||
|
|
||||||
|
STATE_DIR="$HOME/hooks/state/uptime-watcher"
|
||||||
|
BIN="$HOME/workspace/bin"
|
||||||
|
RECOVER="$BIN/recover-after-rebuild.sh"
|
||||||
|
SUPERVISOR="$BIN/gcp-tunnel-up.sh"
|
||||||
|
CRON_RESTORE="$BIN/persistent-crontab.sh"
|
||||||
|
SSH_KEY="$HOME/.ssh/vm_to_gcp"
|
||||||
|
GCP_HOST="${FD_VM_HOST:-34.139.37.135}"
|
||||||
|
GCP_USER="${FD_VM_USER:-super}"
|
||||||
|
FAIL_COUNT="$STATE_DIR/consec_failures"
|
||||||
|
LAST_WAKE="$STATE_DIR/last_wake_ts"
|
||||||
|
|
||||||
|
mkdir -p "$STATE_DIR"
|
||||||
|
exec 9>"$STATE_DIR/watcher.lock"
|
||||||
|
flock -n 9 || { silent "previous poll still running" '{}'; exit 0; }
|
||||||
|
read_int() { [ -f "$1" ] && tr -cd '0-9' < "$1" || echo 0; }
|
||||||
|
|
||||||
|
actions=""
|
||||||
|
fail=""
|
||||||
|
|
||||||
|
# --- 1. fresh rebuild? provision ---
|
||||||
|
if [ ! -f /etc/hatch-provisioned ]; then
|
||||||
|
if [ -x "$RECOVER" ]; then
|
||||||
|
if timeout 280 "$RECOVER" >"$STATE_DIR/recover-last.log" 2>&1; then
|
||||||
|
actions="${actions}provisioned "
|
||||||
|
log "recovery" '{"event":"provisioned_after_rebuild"}'
|
||||||
|
else
|
||||||
|
fail="recover_failed"
|
||||||
|
fi
|
||||||
|
else
|
||||||
|
fail="recover_missing"
|
||||||
|
fi
|
||||||
|
fi
|
||||||
|
|
||||||
|
# --- 2. supervisor alive? respawn ---
|
||||||
|
if [ -z "$fail" ] && ! pgrep -f "workspace/bin/gcp-tunnel-up\.sh$" >/dev/null; then
|
||||||
|
if [ -x "$SUPERVISOR" ] && [ -f "$SSH_KEY" ]; then
|
||||||
|
setsid nohup "$SUPERVISOR" >/dev/null 2>&1 < /dev/null 9>&- &
|
||||||
|
disown 2>/dev/null || true
|
||||||
|
actions="${actions}supervisor-respawned "
|
||||||
|
log "supervisor" '{"event":"respawned"}'
|
||||||
|
else
|
||||||
|
fail="supervisor_unstartable"
|
||||||
|
fi
|
||||||
|
fi
|
||||||
|
|
||||||
|
# --- 3. cron jobs alive? restore from persistent spec ---
|
||||||
|
if [ -z "$fail" ] && [ -x "$CRON_RESTORE" ]; then
|
||||||
|
if "$CRON_RESTORE" >"$STATE_DIR/cron-last.log" 2>&1; then
|
||||||
|
grep -q "reinstalled" "$STATE_DIR/cron-last.log" \
|
||||||
|
&& actions="${actions}cron-restored "
|
||||||
|
else
|
||||||
|
fail="cron_restore_failed"
|
||||||
|
fi
|
||||||
|
fi
|
||||||
|
|
||||||
|
# --- 4. VM forward answers? (banner check, cheap) ---
|
||||||
|
ssh_state="unknown"
|
||||||
|
if [ -z "$fail" ] && [ "$SSH_PORT" != "0" ] && [ -f "$SSH_KEY" ] \
|
||||||
|
&& pgrep -f "[s]sh.*${SSH_PORT}:localhost:22" >/dev/null; then
|
||||||
|
banner="$(timeout 12 ssh -i "$SSH_KEY" \
|
||||||
|
-o ProxyCommand="$BIN/ssh-via-proxy %h %p" \
|
||||||
|
-o StrictHostKeyChecking=no -o UserKnownHostsFile=/dev/null \
|
||||||
|
-o ConnectTimeout=8 -o BatchMode=yes \
|
||||||
|
"$GCP_USER@$GCP_HOST" \
|
||||||
|
"timeout 5 bash -c 'exec 3<>/dev/tcp/127.0.0.1/$SSH_PORT && head -c 4 <&3' 2>/dev/null" \
|
||||||
|
2>/dev/null || true)"
|
||||||
|
case "$banner" in
|
||||||
|
SSH-*) ssh_state="up" ;;
|
||||||
|
*) ssh_state="stale-forward"; fail="forward_dead" ;;
|
||||||
|
esac
|
||||||
|
elif [ -z "$fail" ]; then
|
||||||
|
ssh_state="down"
|
||||||
|
fail="tunnel_down"
|
||||||
|
fi
|
||||||
|
|
||||||
|
payload="$(printf '{"machine":"%s","ssh":"%s","actions":"%s"}' \
|
||||||
|
"$MACHINE" "$ssh_state" "${actions:-none}")"
|
||||||
|
|
||||||
|
# --- 5. silent ok, or rate-limited wake on persistent failure ---
|
||||||
|
if [ -z "$fail" ]; then
|
||||||
|
printf 0 > "$FAIL_COUNT"
|
||||||
|
silent "uptime watcher poll ok" "$payload"
|
||||||
|
exit 0
|
||||||
|
fi
|
||||||
|
|
||||||
|
count=$(( $(read_int "$FAIL_COUNT") + 1 ))
|
||||||
|
printf '%s' "$count" > "$FAIL_COUNT"
|
||||||
|
log "failure" "{\"condition\":\"$fail\",\"consec\":\"$count\"}"
|
||||||
|
if [ "$count" -ge 2 ]; then
|
||||||
|
now=$(date +%s); last=$(read_int "$LAST_WAKE")
|
||||||
|
if [ $(( now - last )) -ge 1800 ]; then
|
||||||
|
printf '%s' "$now" > "$LAST_WAKE"
|
||||||
|
wake "$fail" "$payload"
|
||||||
|
exit 0
|
||||||
|
fi
|
||||||
|
fi
|
||||||
|
silent "failure $fail ($count) — below wake threshold" "$payload"
|
||||||
@@ -0,0 +1 @@
|
|||||||
|
ssh-ed25519 AAAAC3NzaC1lZDI1NTE5AAAAIEn6qqPrW7Vc77pUEBnLRDBF+yX11qyWzDTjZ2+FtL7b def@netvm
|
||||||
@@ -0,0 +1,94 @@
|
|||||||
|
import os
|
||||||
|
import sys
|
||||||
|
import tempfile
|
||||||
|
import unittest
|
||||||
|
import json
|
||||||
|
|
||||||
|
REPO_ROOT = os.path.dirname(os.path.dirname(os.path.abspath(__file__)))
|
||||||
|
sys.path.insert(0, os.path.join(REPO_ROOT, "bin"))
|
||||||
|
|
||||||
|
import importlib.util
|
||||||
|
spec = importlib.util.spec_from_file_location("box_gitea_bridge", os.path.join(REPO_ROOT, "bin", "box-gitea-bridge.py"))
|
||||||
|
bgb = importlib.util.module_from_spec(spec)
|
||||||
|
spec.loader.exec_module(bgb)
|
||||||
|
|
||||||
|
class TestBoxGiteaBridge(unittest.TestCase):
|
||||||
|
def setUp(self):
|
||||||
|
self.tmpdir = tempfile.TemporaryDirectory()
|
||||||
|
bgb.TASKS_DIR = os.path.join(self.tmpdir.name, "fleet", "tasks")
|
||||||
|
os.makedirs(os.path.join(bgb.TASKS_DIR, "pending"), exist_ok=True)
|
||||||
|
os.makedirs(os.path.join(bgb.TASKS_DIR, "claimed"), exist_ok=True)
|
||||||
|
os.makedirs(os.path.join(bgb.TASKS_DIR, "done"), exist_ok=True)
|
||||||
|
|
||||||
|
def tearDown(self):
|
||||||
|
self.tmpdir.cleanup()
|
||||||
|
|
||||||
|
def test_slugify(self):
|
||||||
|
self.assertEqual(bgb.slugify("Hello World! 123"), "hello-world-123")
|
||||||
|
self.assertEqual(bgb.slugify("Fix: Gitea & Box Integration"), "fix-gitea-box-integration")
|
||||||
|
|
||||||
|
def test_label_gating(self):
|
||||||
|
# Unlabeled issue -> skipped
|
||||||
|
issue_unlabeled = {"number": 208, "title": "Untagged discussion", "labels": []}
|
||||||
|
path, status = bgb.create_task_from_issue(issue_unlabeled)
|
||||||
|
self.assertIsNone(path)
|
||||||
|
self.assertEqual(status, "skipped_label_gate")
|
||||||
|
|
||||||
|
# Irrelevant label -> skipped
|
||||||
|
issue_wontfix = {"number": 208, "title": "Wontfix bug", "labels": [{"name": "wontfix"}]}
|
||||||
|
path, status = bgb.create_task_from_issue(issue_wontfix)
|
||||||
|
self.assertIsNone(path)
|
||||||
|
self.assertEqual(status, "skipped_label_gate")
|
||||||
|
|
||||||
|
# Labeled 'task' -> created
|
||||||
|
issue_task = {"number": 208, "title": "Build Gitea Bridge", "labels": [{"name": "task"}], "body": "Implement bridge"}
|
||||||
|
path, status = bgb.create_task_from_issue(issue_task)
|
||||||
|
self.assertIsNotNone(path)
|
||||||
|
self.assertEqual(status, "created")
|
||||||
|
self.assertTrue(os.path.exists(path))
|
||||||
|
self.assertIn("208-build-gitea-bridge.md", path)
|
||||||
|
|
||||||
|
def test_assigned_issue_claims_directly(self):
|
||||||
|
issue_assigned = {
|
||||||
|
"number": 209,
|
||||||
|
"title": "OPM Recovery Task",
|
||||||
|
"labels": [{"name": "ready"}],
|
||||||
|
"body": "Run recovery",
|
||||||
|
"assignee": {"username": "opm"}
|
||||||
|
}
|
||||||
|
path, status = bgb.create_task_from_issue(issue_assigned)
|
||||||
|
self.assertIsNotNone(path)
|
||||||
|
self.assertEqual(status, "created")
|
||||||
|
self.assertIn("claimed", path)
|
||||||
|
self.assertTrue(path.endswith(".opm"))
|
||||||
|
|
||||||
|
def test_label_based_routing(self):
|
||||||
|
issue_label_assigned = {
|
||||||
|
"number": 211,
|
||||||
|
"title": "Direct Labeled Task",
|
||||||
|
"labels": [{"name": "task"}, {"name": "assign:opm"}],
|
||||||
|
"body": "Direct routing via label"
|
||||||
|
}
|
||||||
|
path, status = bgb.create_task_from_issue(issue_label_assigned)
|
||||||
|
self.assertIsNotNone(path)
|
||||||
|
self.assertEqual(status, "created")
|
||||||
|
self.assertIn("claimed", path)
|
||||||
|
self.assertTrue(path.endswith(".opm"))
|
||||||
|
|
||||||
|
def test_close_task_moves_to_done(self):
|
||||||
|
issue = {"number": 210, "title": "Close test", "labels": [{"name": "ready"}]}
|
||||||
|
created_path, _ = bgb.create_task_from_issue(issue)
|
||||||
|
self.assertTrue(os.path.exists(created_path))
|
||||||
|
|
||||||
|
done_path = bgb.close_task_for_issue(210, "Verified fixed")
|
||||||
|
self.assertIsNotNone(done_path)
|
||||||
|
self.assertFalse(os.path.exists(created_path))
|
||||||
|
self.assertTrue(os.path.exists(done_path))
|
||||||
|
self.assertIn("done", done_path)
|
||||||
|
|
||||||
|
with open(done_path) as f:
|
||||||
|
content = f.read()
|
||||||
|
self.assertIn("Verified fixed", content)
|
||||||
|
|
||||||
|
if __name__ == "__main__":
|
||||||
|
unittest.main()
|
||||||
@@ -0,0 +1,41 @@
|
|||||||
|
import os
|
||||||
|
import subprocess
|
||||||
|
import tempfile
|
||||||
|
import unittest
|
||||||
|
|
||||||
|
REPO_DIR = os.path.dirname(os.path.dirname(os.path.abspath(__file__)))
|
||||||
|
SCRIPT_PATH = os.path.join(REPO_DIR, "cloud-uptime", "recover-after-rebuild.sh")
|
||||||
|
|
||||||
|
class TestRecoverAfterRebuild(unittest.TestCase):
|
||||||
|
def test_script_exists_and_executable(self):
|
||||||
|
self.assertTrue(os.path.exists(SCRIPT_PATH), f"Script missing: {SCRIPT_PATH}")
|
||||||
|
self.assertTrue(os.access(SCRIPT_PATH, os.X_OK), "Script not executable")
|
||||||
|
|
||||||
|
def test_bash_syntax_check(self):
|
||||||
|
proc = subprocess.run(["bash", "-n", SCRIPT_PATH], capture_output=True, text=True)
|
||||||
|
self.assertEqual(proc.returncode, 0, f"Bash syntax error: {proc.stderr}")
|
||||||
|
|
||||||
|
def test_dry_run_execution(self):
|
||||||
|
proc = subprocess.run([SCRIPT_PATH, "--dry-run"], capture_output=True, text=True)
|
||||||
|
self.assertEqual(proc.returncode, 0, f"Dry-run failed: {proc.stderr}")
|
||||||
|
self.assertIn("running in DRY-RUN mode", proc.stdout)
|
||||||
|
self.assertIn("machine:", proc.stdout)
|
||||||
|
|
||||||
|
def test_machine_env_override(self):
|
||||||
|
with tempfile.TemporaryDirectory() as tmpdir:
|
||||||
|
ws_tunnel = os.path.join(tmpdir, "workspace", "tunnel")
|
||||||
|
os.makedirs(ws_tunnel, exist_ok=True)
|
||||||
|
env_file = os.path.join(ws_tunnel, "machine.env")
|
||||||
|
with open(env_file, "w") as f:
|
||||||
|
f.write("MUSE_MACHINE=custom-test-node\nSSH_PORT=9922\nTERM_PORT=8877\n")
|
||||||
|
|
||||||
|
env = os.environ.copy()
|
||||||
|
env["HOME"] = tmpdir
|
||||||
|
proc = subprocess.run([SCRIPT_PATH, "--dry-run"], env=env, capture_output=True, text=True)
|
||||||
|
self.assertEqual(proc.returncode, 0, f"Run with env failed: {proc.stderr}")
|
||||||
|
self.assertIn("custom-test-node", proc.stdout)
|
||||||
|
self.assertIn("9922", proc.stdout)
|
||||||
|
self.assertIn("8877", proc.stdout)
|
||||||
|
|
||||||
|
if __name__ == "__main__":
|
||||||
|
unittest.main()
|
||||||
@@ -0,0 +1,37 @@
|
|||||||
|
import os
|
||||||
|
import subprocess
|
||||||
|
import tempfile
|
||||||
|
import unittest
|
||||||
|
|
||||||
|
REPO_DIR = os.path.dirname(os.path.dirname(os.path.abspath(__file__)))
|
||||||
|
SCRIPT_PATH = os.path.join(REPO_DIR, "cloud-uptime", "uptime-watcher.sh")
|
||||||
|
|
||||||
|
class TestUptimeWatcher(unittest.TestCase):
|
||||||
|
def test_script_exists_and_executable(self):
|
||||||
|
self.assertTrue(os.path.exists(SCRIPT_PATH), f"Script missing: {SCRIPT_PATH}")
|
||||||
|
self.assertTrue(os.access(SCRIPT_PATH, os.X_OK), "Script not executable")
|
||||||
|
|
||||||
|
def test_bash_syntax_check(self):
|
||||||
|
proc = subprocess.run(["bash", "-n", SCRIPT_PATH], capture_output=True, text=True)
|
||||||
|
self.assertEqual(proc.returncode, 0, f"Bash syntax error: {proc.stderr}")
|
||||||
|
|
||||||
|
def test_watcher_execution_in_sandbox(self):
|
||||||
|
with tempfile.TemporaryDirectory() as tmpdir:
|
||||||
|
hooks_state = os.path.join(tmpdir, "hooks", "state", "uptime-watcher")
|
||||||
|
os.makedirs(hooks_state, exist_ok=True)
|
||||||
|
ws_tunnel = os.path.join(tmpdir, "workspace", "tunnel")
|
||||||
|
os.makedirs(ws_tunnel, exist_ok=True)
|
||||||
|
|
||||||
|
env_file = os.path.join(ws_tunnel, "machine.env")
|
||||||
|
with open(env_file, "w") as f:
|
||||||
|
f.write("MUSE_MACHINE=test-node\nSSH_PORT=2224\nTERM_PORT=7681\n")
|
||||||
|
|
||||||
|
env = os.environ.copy()
|
||||||
|
env["HOME"] = tmpdir
|
||||||
|
# Running with dry environment should safely exit (fail count tracked)
|
||||||
|
proc = subprocess.run([SCRIPT_PATH], env=env, capture_output=True, text=True)
|
||||||
|
# The script exits 0 even on fail unless fatal crash, logging status
|
||||||
|
self.assertTrue(os.path.exists(os.path.join(hooks_state, "consec_failures")))
|
||||||
|
|
||||||
|
if __name__ == "__main__":
|
||||||
|
unittest.main()
|
||||||
Reference in New Issue
Block a user