| |
| """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, "<remote>", "exec"), g) |
| q.put(("complete", "ok")) |
| except BaseException: |
| 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) |
|
|