# -*- coding: utf-8 -*- """control_server.py — 远程控制通道 (port 5000). 重构自旧 colab_server 协议. Endpoints: POST /execute_stream {"code": "..."} -> SSE: {type: stdout|stderr|error|complete} GET /health -> {"ok": true} 注意: 这是无鉴权的任意代码执行端点, 仅用于开发调试, URL 为随机子域名不可枚举, 用完建议关掉 Colab 会话. """ import json import queue import sys import threading import traceback from flask import Flask, Response, jsonify, request app = Flask(__name__) class _StreamToQ: def __init__(self, q, tag): self.q, self.tag = q, tag def write(self, s): if s: self.q.put((self.tag, s)) return len(s) def flush(self): pass def writelines(self, lines): for s in lines: self.write(s) def isatty(self): return False @app.route("/health") def health(): return jsonify({"ok": True, "service": "samai-control"}) @app.route("/execute_stream", methods=["POST"]) def execute_stream(): code = (request.get_json(force=True) or {}).get("code", "") q = queue.Queue() def worker(): g = {"__name__": "__main__"} old = sys.stdout, sys.stderr sys.stdout = _StreamToQ(q, "stdout") sys.stderr = _StreamToQ(q, "stderr") try: exec(compile(code, "", "exec"), g) # noqa: S102 q.put(("complete", "ok")) except BaseException: # noqa: BLE001 q.put(("error", traceback.format_exc())) q.put(("complete", "error")) finally: sys.stdout, sys.stderr = old threading.Thread(target=worker, daemon=True).start() def gen(): while True: t, c = q.get() yield f"data: {json.dumps({'type': t, 'content': c})}\n\n" if t == "complete": break return Response(gen(), mimetype="text/event-stream", headers={"Cache-Control": "no-cache", "X-Accel-Buffering": "no"}) if __name__ == "__main__": app.run(host="0.0.0.0", port=5000, threaded=True)