File size: 2,156 Bytes
2bbfd2a | 1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 40 41 42 43 44 45 46 47 48 49 50 51 52 53 54 55 56 57 58 59 60 61 62 63 64 65 66 67 68 69 70 71 72 73 74 75 76 77 78 79 80 81 82 | # -*- 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, "<remote>", "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)
|