app = FastAPI(title="Codebuddy Web Console")
+@app.on_event("startup")
+async def _startup():
+ # 启动后台回收协程,定期清理已退出/僵死且断开连接的任务 entry
+ asyncio.create_task(_reap_loop())
+ logger.info("已启动后台任务回收协程")
+
+
def create_token() -> str:
expire = datetime.datetime.utcnow() + datetime.timedelta(hours=JWT_EXPIRE_HOURS)
return jwt.encode({"exp": expire, "rand": secrets.token_hex(8)}, JWT_SECRET, algorithm=JWT_ALGORITHM)
def _kill_entry(task_id: str) -> bool:
- """真正终止任务进程(进程组)并从 RUNNING 移除。"""
+ """真正终止任务进程(进程组),先 SIGTERM 再补 SIGKILL 兜底,并从 RUNNING 移除。
+
+ 注意:codebuddy 进程会忽略 SIGTERM,若只发 SIGTERM 旧进程不会死,
+ 重连时 /ws 会 attach 回这个僵死的旧进程,导致终端卡住、无法输入。
+ 因此这里必须在 SIGTERM 之后立刻补 SIGKILL。
+ """
entry = RUNNING.pop(task_id, None)
if not entry:
return False
proc = entry["proc"]
try:
pgid = os.getpgid(proc.pid)
- os.killpg(pgid, signal.SIGTERM)
+ try:
+ os.killpg(pgid, signal.SIGTERM)
+ except ProcessLookupError:
+ return True
+ # 兜底:SIGKILL 确保进程一定退出(codebuddy 忽略 SIGTERM)
+ try:
+ os.killpg(pgid, signal.SIGKILL)
+ except ProcessLookupError:
+ pass
except Exception:
try:
proc.kill()
return True
+def _proc_state(pid: int) -> str:
+ """读取 /proc/<pid>/stat 的进程状态字符(R/S/D/Z/T 等);读不到返回空串。"""
+ try:
+ with open(f"/proc/{pid}/stat") as _f:
+ s = _f.read()
+ # 进程名 comm 可能含空格且被包在括号里,状态位在最后一个 ')' 之后
+ rp = s.rfind(")")
+ return s[rp + 2:rp + 3]
+ except Exception:
+ return ""
+
+
+async def _reap_loop():
+ """后台回收:周期性清理 RUNNING 中已退出/僵死且已断开连接的任务 entry。
+
+ 目的:
+ - 避免“陈旧 entry 仍指向已退出的进程”,导致前端重连时 /ws 误以为进程存活、
+ 从而 attach 回死进程(表现为终端卡死、无法输入)。
+ - 释放已结束且无人观看的任务占用的内存与缓冲。
+
+ 注意:不会触碰“已断开连接但进程仍活着”的 entry——那是“断开后继续在后台跑”
+ 的预期行为,必须保留;也不会触碰“仍连接着”的 entry。
+ """
+ while True:
+ await asyncio.sleep(30)
+ try:
+ for tid, entry in list(RUNNING.items()):
+ proc = entry["proc"]
+ # 仍连着且进程活着:正常会话,跳过
+ if entry.get("ws") is not None and proc.poll() is None:
+ continue
+ dead = proc.poll() is not None
+ if not dead:
+ # 进程可能处于僵尸态(poll 仍返回 None,需 wait 回收)
+ if _proc_state(proc.pid) == "Z":
+ try:
+ proc.wait(timeout=1)
+ except Exception:
+ pass
+ dead = True
+ if dead and entry.get("ws") is None:
+ RUNNING.pop(tid, None)
+ logger.info(f"回收已退出的任务 entry: name={entry.get('name')}, pid={proc.pid}")
+ except Exception as _e:
+ logger.warning(f"回收任务时出错: {_e}")
+
+
# ==================== 前端 HTML ====================
# 前后端分离:前端页面单独放在 frontend/index.html,便于维护
FRONTEND_DIR = Path(__file__).resolve().parent.parent / "frontend"
t = find_task(data, task_id)
if not t:
return JSONResponse({"ok": False, "error": "任务不存在"}, status_code=404)
+ old_model = (t.get("model") or "")
if model:
# 校验是否在已知模型列表中(未知则拒绝,避免无效模型)
with MODELS_CACHE_LOCK:
t["model"] = ""
save_tasks(data)
logger.info(f"设置模型: id={task_id}, name={t.get('name')}, model='{t['model']}'")
+
+ # 模型实际变化且进程在运行:服务端直接终止进程并从 RUNNING 移除 entry,
+ # 前端重连时 /ws 会按新保存的模型重新拉起一个全新会话。
+ # 关键:必须用 SIGKILL 兜底(codebuddy 忽略 SIGTERM),且立刻 pop RUNNING entry,
+ # 否则重连会 attach 回僵死的旧进程导致终端卡死、无法输入。
+ # 不依赖前端 pid——从非挂载窗口(如任务列表页)切换时前端 pid 不可靠,曾导致"切换无效"。
+ if model != old_model:
+ entry = RUNNING.get(task_id)
+ if entry is not None and entry["proc"].poll() is None:
+ old_pid = entry["proc"].pid
+ _kill_entry(task_id)
+ logger.info(f"模型变更,终止旧任务进程等待重启: name={t.get('name')}, old_pid={old_pid}")
return JSONResponse({"ok": True, "task": t})
RUNNING.pop(tid, None)
break
# 杀掉整个进程组(codebuddy 用 setsid 独立成组,子进程也一起杀)
+ # codebuddy 会忽略 SIGTERM,必须再补 SIGKILL 兜底
try:
pgid = os.getpgid(pid)
- os.killpg(pgid, signal.SIGTERM)
- except ProcessLookupError:
- pass
+ try:
+ os.killpg(pgid, signal.SIGTERM)
+ except ProcessLookupError:
+ pass
+ try:
+ os.killpg(pgid, signal.SIGKILL)
+ except ProcessLookupError:
+ pass
except Exception:
try:
- os.kill(pid, signal.SIGTERM)
+ os.kill(pid, signal.SIGKILL)
except Exception:
pass
return JSONResponse({"ok": True})
if task_id in ACTIVE_TASKS:
# 可能是上一次连接刚断开、finally 还没来得及 discard 的瞬间竞态;
- # 若该任务的 ws 已经置空(已 detach),说明旧连接已结束,允许本次重连(会 attach 到同一进程)
+ # - 若该任务的 ws 已置空(已 detach),说明旧连接已结束,允许重连(attach 到同一进程)
+ # - 若 RUNNING 中已无该任务(进程刚被终止,如切换模型/重启会话),同样放行,
+ # 由后续逻辑按新状态启动,避免"已在另一窗口运行"误拒
prev = RUNNING.get(task_id)
- if prev is not None and prev.get("ws") is None:
+ if prev is None or prev.get("ws") is None:
ACTIVE_TASKS.discard(task_id)
else:
await websocket.accept()