696 lines
26 KiB
Python
696 lines
26 KiB
Python
|
|
#!/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()
|