#!/usr/bin/env python3 """Rumeng bridge — 监听 8795,桥接 App 与 tmux/CC session。""" import json import os import re import secrets import subprocess import sys import datetime import time import threading from http.server import BaseHTTPRequestHandler, HTTPServer from socketserver import ThreadingMixIn from urllib.parse import urlparse, parse_qs import httpx import jwt from group_chat import ( ROSTER_BY_ID, REPLY_AGENT_IDS, append_group_message, read_since, roster, agent_status, context_lines, ) AUTH_TOKEN = os.environ.get("BRIDGE_AUTH_TOKEN", "") if not AUTH_TOKEN: print("ERROR: BRIDGE_AUTH_TOKEN not set", file=sys.stderr) sys.exit(1) BASE_DIR = os.path.dirname(os.path.abspath(__file__)) DATA_DIR = os.path.join(BASE_DIR, "data") os.makedirs(DATA_DIR, exist_ok=True) KNOCK_DIR = os.path.expanduser("~/.shenyan-knock") os.makedirs(KNOCK_DIR, exist_ok=True) CHAT_FILE = os.path.join(DATA_DIR, "chat.jsonl") CODEX_FILE = os.path.join(DATA_DIR, "codex.jsonl") FRIDGE_FILE = os.path.join(DATA_DIR, "fridge.jsonl") MOOD_FILE = os.path.join(DATA_DIR, "mood.jsonl") STATUS_FILE = os.path.join(DATA_DIR, "status.json") APNS_TOKEN_FILE = os.path.join(KNOCK_DIR, "apns_token") APNS_KEY_FILE = os.path.join(KNOCK_DIR, "AuthKey_3VJ5X6V9Q8.p8") APNS_TEAM_ID = "ZZMLW32PHH" APNS_KEY_ID = "3VJ5X6V9Q8" APNS_BUNDLE_ID = "rumeng-v1.0.Rumeng" APNS_ENDPOINT = "https://api.sandbox.push.apple.com" CODEX_TERMINAL_FILE = "/tmp/codex_rumeng_terminal.log" CODEX_TERMINAL_MAX_LINES = 120 CODEX_TMUX_TARGET_FILE = os.path.join(KNOCK_DIR, "codex_tmux_target") CODEX_DEFAULT_TMUX_TARGET = "codex" CONTROL_KEYS_PATTERN = re.compile(r"^([A-Z][A-Za-z]*|C-.|M-.|S-.|Space|Tab|Enter|Escape|BSpace)$") APP_ECHO_PATTERN = re.compile(r"^\[App\]\[\d{2}:\d{2}:\d{2}\]\s+") def iso_now(): return datetime.datetime.now(datetime.timezone.utc).isoformat(timespec="milliseconds").replace("+00:00", "Z") def random_id(prefix): return f"{prefix}_{secrets.token_hex(4)}" def set_app_online(): try: with open(STATUS_FILE, "w", encoding="utf-8") as f: f.write(iso_now()) except Exception: pass def get_app_status(): now = datetime.datetime.now(datetime.timezone.utc) result = {"cc_online": tmux_session_alive()} if not os.path.exists(STATUS_FILE): result.update({"online": False, "seconds_ago": None}) return result try: with open(STATUS_FILE, "r", encoding="utf-8") as f: last_ts_str = f.read().strip() last_ts = datetime.datetime.fromisoformat(last_ts_str.replace("Z", "+00:00")) delta = now - last_ts result.update({"online": delta.total_seconds() <= 30, "seconds_ago": round(delta.total_seconds())}) return result except Exception: result.update({"online": False, "seconds_ago": None}) return result def tmux_session_alive(session="shenyan"): try: subprocess.run(["tmux", "has-session", "-t", session], capture_output=True, check=True, timeout=5) return True except Exception: return False def codex_tmux_target(): target = os.environ.get("CODEX_TMUX_TARGET", "").strip() if target and tmux_target_exists(target): return target try: with open(CODEX_TMUX_TARGET_FILE, "r", encoding="utf-8") as f: target = f.read().strip() if target and tmux_target_exists(target): return target except FileNotFoundError: pass except Exception as e: print(f"[Codex] target read failed: {e}", file=sys.stderr) detected = detect_codex_tmux_target() if detected: return detected return CODEX_DEFAULT_TMUX_TARGET def tmux_target_exists(target): try: subprocess.run(["tmux", "display-message", "-t", target, "-p", "#{pane_id}"], capture_output=True, text=True, check=True, timeout=5) return True except Exception: return False def detect_codex_tmux_target(): try: result = subprocess.run( [ "tmux", "list-panes", "-a", "-F", "#{pane_id}\t#{session_name}:#{window_index}.#{pane_index}\t#{pane_current_command}\t#{pane_title}\t#{pane_start_command}", ], capture_output=True, text=True, check=True, timeout=5, ) except Exception: return "" for line in result.stdout.splitlines(): parts = line.split("\t") if len(parts) < 5: continue pane_id, target, command, title, start_command = parts[:5] haystack = " ".join([command, title, start_command]).lower() if "codex" in haystack and "codex_rumeng_terminal.log" not in haystack: return pane_id or target return "" def tmux_send_text(target, text, enter=True): subprocess.run(["tmux", "send-keys", "-t", target, "-l", text], check=True, timeout=10) if enter: subprocess.run(["tmux", "send-keys", "-t", target, "Enter"], check=True, timeout=10) def read_jsonl(filepath): records = [] if not os.path.exists(filepath): return records with open(filepath, "r", encoding="utf-8") as f: for line in f: line = line.strip() if line: records.append(json.loads(line)) return records def append_jsonl(filepath, record): with open(filepath, "a", encoding="utf-8") as f: f.write(json.dumps(record, ensure_ascii=False) + "\n") def rewrite_jsonl(filepath, records): with open(filepath, "w", encoding="utf-8") as f: for record in records: f.write(json.dumps(record, ensure_ascii=False) + "\n") def read_tail_text(filepath, max_lines): try: with open(filepath, "r", encoding="utf-8") as f: lines = f.read().splitlines() except FileNotFoundError: return "" except Exception as e: return f"[codex-sync log read error: {e}]" return "\n".join(lines[-max_lines:]).strip() def append_codex_terminal_overlay(content): codex_log = read_tail_text(CODEX_TERMINAL_FILE, CODEX_TERMINAL_MAX_LINES) if not codex_log: return content base = content.rstrip() if base: return f"{base}\n\n--- Codex sync ---\n{codex_log}\n" return f"--- Codex sync ---\n{codex_log}\n" def send_apns_alert(text): try: with open(APNS_TOKEN_FILE, "r", encoding="utf-8") as f: device_token = f.read().strip() with open(APNS_KEY_FILE, "r", encoding="utf-8") as f: private_key = f.read() except FileNotFoundError as e: result = {"ok": False, "error": f"missing file: {e}"} print(f"[APNs] skipped: {result['error']}", file=sys.stderr) return result if not device_token: result = {"ok": False, "error": "empty device token"} print(f"[APNs] skipped: {result['error']}", file=sys.stderr) return result provider_token = jwt.encode( {"iss": APNS_TEAM_ID, "iat": int(time.time())}, private_key, algorithm="ES256", headers={"kid": APNS_KEY_ID}, ) payload = { "aps": { "alert": { "title": "", "body": text[:100], }, "sound": "default", } } headers = { "authorization": f"bearer {provider_token}", "apns-topic": APNS_BUNDLE_ID, "apns-push-type": "alert", "apns-priority": "10", } try: with httpx.Client(http2=True, timeout=10) as client: resp = client.post(f"{APNS_ENDPOINT}/3/device/{device_token}", json=payload, headers=headers) if 200 <= resp.status_code < 300: result = {"ok": True, "status": resp.status_code, "apns_id": resp.headers.get("apns-id", "")} print(f"[APNs] sent: HTTP {resp.status_code} apns-id={result['apns_id']}", file=sys.stderr) return result result = {"ok": False, "status": resp.status_code, "response": resp.text} print(f"[APNs] failed: HTTP {resp.status_code} {resp.text}", file=sys.stderr) return result except Exception as e: result = {"ok": False, "error": str(e)} print(f"[APNs] failed: {e}", file=sys.stderr) return result class ThreadingHTTPServer(ThreadingMixIn, HTTPServer): daemon_threads = True class Handler(BaseHTTPRequestHandler): def log_message(self, format, *args): sys.stderr.write(f"[{self.log_date_time_string()}] {' '.join(str(a) for a in args)}\n") def send_json(self, data, status=200): body = json.dumps(data, ensure_ascii=False).encode("utf-8") self.send_response(status) self.send_header("Content-Type", "application/json; charset=utf-8") self.send_header("Content-Length", str(len(body))) self.end_headers() self.wfile.write(body) def read_json_body(self): length = int(self.headers.get("Content-Length", 0)) if length == 0: return {} return json.loads(self.rfile.read(length).decode("utf-8")) def check_auth(self): return self.headers.get("X-Auth-Token") == AUTH_TOKEN def do_GET(self): if not self.check_auth(): return self.send_json({"error": "unauthorized"}, 401) set_app_online() parsed = urlparse(self.path) path, params = parsed.path, parse_qs(parsed.query, keep_blank_values=True) if path == "/chat/status": self.send_json(get_app_status()) elif path == "/chat/history": self.handle_chat_history(params) elif path == "/codex/history": self.handle_codex_history(params) elif path == "/codex/terminal": self.handle_codex_terminal(params) elif path == "/fridge/list": self.handle_fridge_list() elif path == "/mood/list": self.handle_mood_list() elif path == "/group/poll": self.handle_group_poll(params) elif path == "/group/roster": self.send_json({"ok": True, "roster": roster(), "status": agent_status(tmux_session_alive)}) elif path == "/tmux/capture": self.handle_tmux_capture(params) elif path == "/apns/token": self.handle_apns_token() elif path == "/health": self.send_json({"ok": True}) else: self.send_json({"error": "not found"}, 404) def do_POST(self): if not self.check_auth(): return self.send_json({"error": "unauthorized"}, 401) set_app_online() path = urlparse(self.path).path try: body = self.read_json_body() except Exception as e: return self.send_json({"error": f"invalid JSON: {e}"}, 400) try: if path == "/chat/send": self.handle_chat_send(body) elif path == "/chat/append": self.handle_chat_append(body) elif path == "/codex/append": self.handle_codex_append(body) elif path == "/codex/send": self.handle_codex_send(body) elif path == "/codex/terminal/send": self.handle_codex_terminal_send(body) elif path == "/codex/terminal/key": self.handle_codex_terminal_key(body) elif path == "/fridge/add": self.handle_fridge_add(body) elif path == "/fridge/reply": self.handle_fridge_reply(body) elif path == "/fridge/delete": self.handle_fridge_delete(body) elif path == "/mood/upsert": self.handle_mood_upsert(body) elif path == "/group/send": self.handle_group_send(body) elif path == "/tmux/send": self.handle_tmux_send(body) elif path == "/apns/register": self.handle_apns_register(body) elif path == "/apns/test": self.handle_apns_test(body) else: self.send_json({"error": "not found"}, 404) except Exception as e: self.send_json({"error": str(e)}, 500) def handle_chat_send(self, body): text = body.get("text", "") if not text: return self.send_json({"error": "missing text"}, 400) record = {"ts": iso_now(), "role": "眠眠", "text": text} append_jsonl(CHAT_FILE, record) local_hms = datetime.datetime.now().strftime("%H:%M:%S") injected = f"[App][{local_hms}] {text}" try: subprocess.run(["tmux", "send-keys", "-t", "shenyan", "-l", injected], check=True, timeout=10) subprocess.run(["tmux", "send-keys", "-t", "shenyan", "Enter"], check=True, timeout=10) except Exception as e: self.send_json({"ok": True, "record": record, "tmux_error": str(e)}) return self.send_json({"ok": True, "record": record}) def handle_chat_append(self, body): record = { "ts": body.get("ts", iso_now()), "role": body.get("role", ""), "text": body.get("text", ""), } if "thinking" in body: record["thinking"] = body["thinking"] if "source" in body: record["source"] = body["source"] append_jsonl(CHAT_FILE, record) if record["text"] and record["role"] != "眠眠": threading.Thread(target=send_apns_alert, args=(record["text"],), daemon=True).start() self.send_json({"ok": True}) def handle_chat_history(self, params): since = params.get("since", [None])[0] try: limit = int(params.get("limit", ["50"])[0]) except ValueError: limit = 50 records = read_jsonl(CHAT_FILE) if since: records = [r for r in records if r.get("ts", "") > since] records = records[-limit:] self.send_json({"records": records}) def handle_codex_append(self, body): text = str(body.get("text", "")) source = str(body.get("source", "")) if source == "codex-user" and APP_ECHO_PATTERN.match(text): self.send_json({"ok": True, "skipped": "app_echo"}) return record = { "ts": body.get("ts", iso_now()), "role": body.get("role", ""), "text": text, } if "phase" in body: record["phase"] = body["phase"] if "source" in body: record["source"] = body["source"] append_jsonl(CODEX_FILE, record) self.send_json({"ok": True}) def handle_codex_send(self, body): text = str(body.get("text", "")).strip() if not text: return self.send_json({"error": "missing text"}, 400) record = { "ts": iso_now(), "role": "眠眠", "text": text, "source": "codex-app-user", } append_jsonl(CODEX_FILE, record) self.append_codex_terminal_line(record["role"], text, record["ts"], "app") local_hms = datetime.datetime.now().strftime("%H:%M:%S") injected = f"[App][{local_hms}] {text}" target = codex_tmux_target() try: tmux_send_text(target, injected) except subprocess.CalledProcessError as e: err = e.stderr.strip() if e.stderr else str(e) self.append_codex_terminal_line("bridge", f"tmux send failed target={target}: {err}", iso_now(), "error") self.send_json({"ok": True, "record": record, "tmux_target": target, "tmux_error": err}) return except subprocess.TimeoutExpired: self.append_codex_terminal_line("bridge", f"tmux send timeout target={target}", iso_now(), "error") self.send_json({"ok": True, "record": record, "tmux_target": target, "tmux_error": "tmux send timeout"}) return self.append_codex_terminal_line("bridge", f"sent target={target}", iso_now(), "tmux") self.send_json({"ok": True, "record": record, "tmux_target": target}) def handle_codex_terminal_send(self, body): text = str(body.get("text", "")).strip() if not text: return self.send_json({"error": "missing text"}, 400) ts = iso_now() self.append_codex_terminal_line("$", text, ts, "app-terminal") target = codex_tmux_target() try: tmux_send_text(target, text) except subprocess.CalledProcessError as e: err = e.stderr.strip() if e.stderr else str(e) self.append_codex_terminal_line("bridge", f"tmux send failed target={target}: {err}", iso_now(), "error") self.send_json({"ok": False, "ts": ts, "tmux_target": target, "tmux_error": err}, 500) return except subprocess.TimeoutExpired: self.append_codex_terminal_line("bridge", f"tmux send timeout target={target}", iso_now(), "error") self.send_json({"ok": False, "ts": ts, "tmux_target": target, "tmux_error": "tmux send timeout"}, 500) return self.append_codex_terminal_line("bridge", f"terminal sent target={target}", iso_now(), "tmux") self.send_json({"ok": True, "ts": ts, "tmux_target": target}) def handle_codex_terminal_key(self, body): key = str(body.get("key", "")).strip() if not key: return self.send_json({"error": "missing key"}, 400) ts = iso_now() self.append_codex_terminal_line("key", key, ts, "app-terminal") target = codex_tmux_target() try: if CONTROL_KEYS_PATTERN.match(key): subprocess.run(["tmux", "send-keys", "-t", target, key], check=True, timeout=10) else: tmux_send_text(target, key, enter=False) except subprocess.CalledProcessError as e: err = e.stderr.strip() if e.stderr else str(e) self.append_codex_terminal_line("bridge", f"tmux key failed target={target}: {err}", iso_now(), "error") self.send_json({"ok": False, "ts": ts, "tmux_target": target, "tmux_error": err}, 500) return except subprocess.TimeoutExpired: self.append_codex_terminal_line("bridge", f"tmux key timeout target={target}", iso_now(), "error") self.send_json({"ok": False, "ts": ts, "tmux_target": target, "tmux_error": "tmux key timeout"}, 500) return self.append_codex_terminal_line("bridge", f"key sent target={target} key={key}", iso_now(), "tmux") self.send_json({"ok": True, "ts": ts, "tmux_target": target}) def append_codex_terminal_line(self, role, text, ts, source): hms = ts[11:19] if len(ts) >= 19 else datetime.datetime.now().strftime("%H:%M:%S") prefix = f"[Codex/{source}][{hms}] {role}: " line = prefix + text.replace("\n", "\n" + " " * len(prefix)) try: with open(CODEX_TERMINAL_FILE, "a", encoding="utf-8") as f: f.write(line + "\n") except Exception as e: print(f"[Codex] terminal log write failed: {e}", file=sys.stderr) def handle_codex_history(self, params): since = params.get("since", [None])[0] try: limit = int(params.get("limit", ["80"])[0]) except ValueError: limit = 80 records = read_jsonl(CODEX_FILE) if since: records = [r for r in records if r.get("ts", "") > since] self.send_json({"records": records[-limit:]}) def handle_codex_terminal(self, params): try: limit = int(params.get("lines", ["160"])[0]) except ValueError: limit = 160 self.send_json({"content": read_tail_text(CODEX_TERMINAL_FILE, limit)}) def handle_fridge_list(self): self.send_json({"notes": read_jsonl(FRIDGE_FILE)}) def handle_fridge_add(self, body): note = { "id": random_id("fridge"), "text": body.get("text", ""), "role": body.get("role", "眠眠"), "created_at": iso_now(), "replies": [], } append_jsonl(FRIDGE_FILE, note) self.send_json({"ok": True, "id": note["id"]}) def handle_fridge_reply(self, body): note_id = body.get("id") if not note_id: return self.send_json({"error": "missing id"}, 400) records = read_jsonl(FRIDGE_FILE) target = next((r for r in records if r.get("id") == note_id), None) if not target: return self.send_json({"error": "note not found"}, 404) reply = { "id": random_id("reply"), "text": body.get("text", ""), "role": body.get("role", "眠眠"), "created_at": iso_now(), } target.setdefault("replies", []).append(reply) rewrite_jsonl(FRIDGE_FILE, records) self.send_json({"ok": True, "id": reply["id"]}) def handle_fridge_delete(self, body): note_id = body.get("id") if not note_id: return self.send_json({"error": "missing id"}, 400) records = read_jsonl(FRIDGE_FILE) new_records = [r for r in records if r.get("id") != note_id] if len(new_records) == len(records): return self.send_json({"error": "note not found"}, 404) rewrite_jsonl(FRIDGE_FILE, new_records) self.send_json({"ok": True}) # ── group endpoints ────────────────────────────────────── def handle_mood_list(self): self.send_json(read_jsonl(MOOD_FILE)) def handle_mood_upsert(self, body): date = body.get("date", "") role = body.get("role", "") if not date or not role: return self.send_json({"error": "missing date or role"}, 400) records = read_jsonl(MOOD_FILE) existing = next((r for r in records if r.get("date") == date and r.get("role") == role), None) if existing: existing["tag"] = body.get("tag", "") existing["color"] = body.get("color", "") existing["updated_at"] = iso_now() rewrite_jsonl(MOOD_FILE, records) else: record = { "id": random_id("mood"), "date": date, "role": role, "tag": body.get("tag", ""), "color": body.get("color", ""), "created_at": iso_now(), } append_jsonl(MOOD_FILE, record) if role == "眠眠": notify = { "ts": iso_now(), "role": "system", "text": f"[mood] 猫猫今天:{body.get('tag', '')}", "source": "mood", } append_jsonl(CHAT_FILE, notify) self.send_json({"ok": True}) def handle_group_poll(self, params): since = params.get("since", [None])[0] if since: since = since.replace(" ", "+") # parse_qs decodes + as space try: limit = int(params.get("limit", ["100"])[0]) except ValueError: limit = 100 records = read_since(since, limit) self.send_json({ "ok": True, "records": records, "count": len(records), "last_ts": records[-1]["ts"] if records else since, "roster": roster(), "status": agent_status(tmux_session_alive), }) def handle_group_send(self, body): sender_id = body.get("sender_id", "amian") text = body.get("text", "") if not text: return self.send_json({"error": "missing text"}, 400) try: record = append_group_message( sender_id=sender_id, text=text, mentions=body.get("mentions"), reply_to=body.get("reply_to"), ) except ValueError as e: return self.send_json({"error": str(e)}, 400) self.send_json({"ok": True, "record": record}) def handle_tmux_capture(self, params): session = params.get("session", ["shenyan"])[0] try: lines = int(params.get("lines", ["200"])[0]) except ValueError: lines = 200 try: result = subprocess.run( ["tmux", "capture-pane", "-t", session, "-p", "-S", f"-{lines}"], capture_output=True, text=True, check=True, timeout=10, ) self.send_json({"content": result.stdout}) except subprocess.CalledProcessError as e: self.send_json({"error": f"tmux: {e.stderr.strip()}"}, 404) except subprocess.TimeoutExpired: self.send_json({"error": "tmux capture timeout"}, 500) def handle_tmux_send(self, body): keys = body.get("keys", "") session = body.get("session", "shenyan") enter = body.get("enter", True) try: if keys: if CONTROL_KEYS_PATTERN.match(keys): subprocess.run(["tmux", "send-keys", "-t", session, keys], check=True, timeout=10) else: subprocess.run(["tmux", "send-keys", "-t", session, "-l", keys], check=True, timeout=10) if enter: subprocess.run(["tmux", "send-keys", "-t", session, "Enter"], check=True, timeout=10) self.send_json({"ok": True}) except subprocess.CalledProcessError as e: self.send_json({"error": f"tmux: {e.stderr.strip() if e.stderr else e}"}, 500) except subprocess.TimeoutExpired: self.send_json({"error": "tmux send timeout"}, 500) def handle_apns_register(self, body): token = str(body.get("token", "")).strip() if not token: return self.send_json({"error": "missing token"}, 400) with open(APNS_TOKEN_FILE, "w", encoding="utf-8") as f: f.write(token + "\n") self.send_json({"ok": True}) def handle_apns_token(self): token = "" try: with open(APNS_TOKEN_FILE, "r", encoding="utf-8") as f: token = f.read().strip() except FileNotFoundError: pass self.send_json({"token": token}) def handle_apns_test(self, body): text = str(body.get("text", "测试横幅。")) self.send_json(send_apns_alert(text)) def main(): host = "0.0.0.0" port = int(os.environ.get("BRIDGE_PORT", "8795")) server = ThreadingHTTPServer((host, port), Handler) print(f"rumeng-bridge listening on {host}:{port}", file=sys.stderr) try: server.serve_forever() except KeyboardInterrupt: print("\nshutting down...", file=sys.stderr) server.shutdown() if __name__ == "__main__": main()