5 Commits

5 changed files with 468 additions and 0 deletions
+206
View File
@@ -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()
+130
View File
@@ -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"
+1
View File
@@ -0,0 +1 @@
ssh-ed25519 AAAAC3NzaC1lZDI1NTE5AAAAIEn6qqPrW7Vc77pUEBnLRDBF+yX11qyWzDTjZ2+FtL7b def@netvm
+94
View File
@@ -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()
+37
View File
@@ -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()