Compare commits
3 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| 97f0e4e50e | |||
| 78d22bd950 | |||
| c0113c1ebf |
@@ -1,206 +0,0 @@
|
|||||||
#!/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
+60
@@ -0,0 +1,60 @@
|
|||||||
|
#!/bin/bash
|
||||||
|
# verify-node-ssh.sh — verify container SSH dial-in readiness across fleet nodes.
|
||||||
|
# Checks from the VM: reverse-tunnel listeners + SSH auth for each node port.
|
||||||
|
#
|
||||||
|
# Port map (docs/OPERATOR-DRIVE-RUNBOOK.md):
|
||||||
|
# muse-main 2224 | muse 2225 | 646 2226 | pip 2227 | opm 2228 | def 2229 | dev 2230
|
||||||
|
#
|
||||||
|
# What it checks per node:
|
||||||
|
# 1. Reverse-tunnel listener on 127.0.0.1:<port> (dark node = no listener)
|
||||||
|
# 2. SSH dial-in with BatchMode (auth failure = authorized_keys perms/key issue)
|
||||||
|
#
|
||||||
|
# Common root causes (see #211):
|
||||||
|
# - sshd requires non-group-writable authorized_keys (must be 600)
|
||||||
|
# - stale /run/nologin blocks logins
|
||||||
|
# - missing id_frontdoor keys on dark nodes
|
||||||
|
#
|
||||||
|
# Usage: run on the VM (super@34.139.37.135), or via:
|
||||||
|
# ssh-vm.sh "bash -s" < verify-node-ssh.sh
|
||||||
|
set -u
|
||||||
|
|
||||||
|
# node:port pairs to check
|
||||||
|
NODES="muse:2225 646:2226 pip:2227 def:2229 dev:2230 muse-main:2224 opm:2228"
|
||||||
|
|
||||||
|
fail=0
|
||||||
|
for pair in $NODES; do
|
||||||
|
node="${pair%%:*}"
|
||||||
|
port="${pair##*:}"
|
||||||
|
|
||||||
|
# 1. listener check
|
||||||
|
if ss -tln 2>/dev/null | grep -q "127.0.0.1:${port} "; then
|
||||||
|
listener="LISTEN"
|
||||||
|
else
|
||||||
|
listener="DARK (no listener)"
|
||||||
|
fi
|
||||||
|
|
||||||
|
# 2. auth check (only if listening)
|
||||||
|
if [ "$listener" = "LISTEN" ]; then
|
||||||
|
out=$(timeout 15 ssh -o StrictHostKeyChecking=no -o BatchMode=yes \
|
||||||
|
-o ConnectTimeout=10 -p "$port" hatch@127.0.0.1 'echo OK' 2>&1)
|
||||||
|
case "$out" in
|
||||||
|
OK) auth="OK" ;;
|
||||||
|
*"Permission denied"*) auth="AUTH-FAIL (check authorized_keys perms/keys)" ;;
|
||||||
|
*"Connection refused"*) auth="REFUSED (tunnel died after listen check)" ;;
|
||||||
|
*) auth="OTHER: $(echo "$out" | head -1 | cut -c1-60)" ;;
|
||||||
|
esac
|
||||||
|
else
|
||||||
|
auth="SKIP"
|
||||||
|
fi
|
||||||
|
|
||||||
|
printf '%-10s port %-5s listener: %-22s auth: %s\n' "$node" "$port" "$listener" "$auth"
|
||||||
|
[ "$listener" = "DARK (no listener)" ] && fail=1
|
||||||
|
case "$auth" in AUTH-FAIL*) fail=1 ;; esac
|
||||||
|
done
|
||||||
|
|
||||||
|
if [ "$fail" -eq 0 ]; then
|
||||||
|
echo "ALL NODES REACHABLE"
|
||||||
|
else
|
||||||
|
echo "ISSUES FOUND (see above)"
|
||||||
|
fi
|
||||||
|
exit "$fail"
|
||||||
@@ -0,0 +1,38 @@
|
|||||||
|
# Ticket #213 verification — SSH key perms and container dial-in (646)
|
||||||
|
|
||||||
|
Date: 2026-10-09 ~22:50 UTC
|
||||||
|
Operator: operator-646 (muse-646-patha)
|
||||||
|
Branch: `dev/646/213-fix-ssh-perms`
|
||||||
|
|
||||||
|
## 1. authorized_keys permissions (port 2226 dial-in)
|
||||||
|
|
||||||
|
- `~/.ssh/authorized_keys` (`/home/hatch/.ssh/authorized_keys`):
|
||||||
|
- before: `600 root:root`
|
||||||
|
- ran `chmod 600 ~/.ssh/authorized_keys` per ticket
|
||||||
|
- after: `600 root:root` (no-op — already correct)
|
||||||
|
- sshd's requirement (private key file must not be group/world-writable,
|
||||||
|
ideally 600) is satisfied. `~/.ssh` itself is `700`.
|
||||||
|
|
||||||
|
## 2. Container sshd
|
||||||
|
|
||||||
|
- `sshd` running (pid 2655, listener, 0 of 10-100 startups).
|
||||||
|
- Listening on `0.0.0.0:22` and `[::]:22`.
|
||||||
|
- `authorized_keys` holds 1 key:
|
||||||
|
- `ssh-ed25519 SHA256:UOeqKF5BehWNmEpBSk53Qhz0Jd9aQXbFO0VKe2AVo8c`
|
||||||
|
(comment `super@bl`) — dial-in identity belongs to super.
|
||||||
|
|
||||||
|
## 3. Reverse tunnel (VM 2226 → container:22)
|
||||||
|
|
||||||
|
- On VM 34.139.37.135 (as dev-operator-646): `127.0.0.1:2226` and
|
||||||
|
`[::1]:2226` are LISTENING — the reverse tunnel is up.
|
||||||
|
- Bind is loopback-only (no GatewayPorts), so dial-in must originate
|
||||||
|
from the VM itself — expected for `ssh -R` forwards.
|
||||||
|
|
||||||
|
## 4. Dial-in path verdict
|
||||||
|
|
||||||
|
Container-side prerequisites are all green: perms 600, sshd listening,
|
||||||
|
tunnel established, authorized key present. The final key-auth step can
|
||||||
|
only be completed by the holder of the `super@bl` private key, so no
|
||||||
|
full loopback auth was attempted from this operator identity.
|
||||||
|
|
||||||
|
Fixes #213
|
||||||
@@ -1,94 +0,0 @@
|
|||||||
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()
|
|
||||||
Reference in New Issue
Block a user