Files
box/bin/box-gitea-bridge.py

207 lines
7.3 KiB
Python
Executable File

#!/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()