diff --git a/bin/box-gitea-bridge.py b/bin/box-gitea-bridge.py new file mode 100755 index 0000000..4dd30d0 --- /dev/null +++ b/bin/box-gitea-bridge.py @@ -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: or 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() diff --git a/tests/test_box_gitea_bridge.py b/tests/test_box_gitea_bridge.py new file mode 100644 index 0000000..1ebc2ff --- /dev/null +++ b/tests/test_box_gitea_bridge.py @@ -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()