# 活动任务锁:同一任务同时只允许一个会话,防止两窗口同时 resume 损坏 jsonl
ACTIVE_TASKS = set()
+# 后台常驻任务:task_id -> entry
+# entry = {proc, master_fd, sid, name, cwd, model, started_at,
+# buffer: list[str], buflen: int, ws: WebSocket|None,
+# done: bool, exit_code: int|None}
+# 关闭浏览器(WS 断开)时进程不杀,保留在 RUNNING 里后台继续跑;
+# 下次重连时 attach 到同一 PTY,回放断开期间的输出。
+RUNNING = {}
+PROC_BUFFER_CAP = 500 * 1024 # 每个任务最多缓存的输出字节数
+REPLAY_LIMIT = 30 * 1024 # 重连时回放给前端的输出字节数上限
+
if not PASSWORD:
PASSWORD = secrets.token_urlsafe(12)
logger.info(f"🔑 自动生成密码: {PASSWORD}")
save_tasks(data)
+def _trim_buffer(entry: dict) -> None:
+ buf = entry["buffer"]
+ total = entry["buflen"]
+ while buf and total > PROC_BUFFER_CAP:
+ total -= len(buf.pop(0))
+ entry["buflen"] = total
+
+
+async def _drain(entry: dict) -> None:
+ """持续把 PTY 输出读入缓冲,并转发给当前连接的 ws(若有)。"""
+ loop = asyncio.get_running_loop()
+ fd = entry["master_fd"]
+ proc = entry["proc"]
+ while True:
+ try:
+ data = await loop.run_in_executor(None, lambda: os.read(fd, 65536))
+ except BlockingIOError:
+ # 非阻塞 fd 在无数据时会抛 BlockingIOError(EAGAIN),
+ # 这是正常空闲状态,不是 EOF —— 只 sleep 重试,绝不能 break。
+ await asyncio.sleep(0.02)
+ continue
+ except OSError:
+ break
+ if not data:
+ if proc.poll() is not None:
+ break
+ await asyncio.sleep(0.02)
+ continue
+ try:
+ text = data.decode("utf-8", errors="replace")
+ except Exception:
+ text = data.decode("latin-1", errors="replace")
+ entry["buffer"].append(text)
+ entry["buflen"] += len(text)
+ _trim_buffer(entry)
+ ws = entry.get("ws")
+ if ws is not None:
+ try:
+ async with entry["send_lock"]:
+ await ws.send_json({"type": "output", "data": text})
+ except Exception:
+ pass
+ entry["exit_code"] = proc.poll()
+ entry["done"] = True
+ try:
+ if entry["master_fd"] is not None:
+ os.close(entry["master_fd"])
+ except OSError:
+ pass
+ entry["master_fd"] = None
+ logger.info(f"任务后台结束: name={entry['name']}, pid={proc.pid}, code={entry['exit_code']}")
+
+
+def last_assistant_summary(sid: str) -> str:
+ """从会话记录里取最后一条 assistant 文本,作为“当前任务”提示。"""
+ matches = glob.glob(str(CODEBUDDY_PROJECTS_DIR / "*" / f"{sid}.jsonl"))
+ if not matches:
+ return ""
+ try:
+ lines = open(matches[0], "r", encoding="utf-8", errors="replace").read().splitlines()
+ except Exception:
+ return ""
+ for raw in reversed(lines):
+ try:
+ o = json.loads(raw)
+ except Exception:
+ continue
+ if o.get("role") != "assistant":
+ continue
+ c = o.get("content")
+ txt = ""
+ if isinstance(c, list):
+ for p in c:
+ if isinstance(p, dict):
+ if p.get("type") == "output_text":
+ txt += p.get("text", "")
+ elif "text" in p:
+ txt += p.get("text", "")
+ elif isinstance(c, str):
+ txt = c
+ txt = " ".join(txt.split())
+ if txt:
+ return txt[:300]
+ return ""
+
+
+def _launch_process(task: dict, sid: str, model: str, cwd: str, codebuddy_path: str) -> dict:
+ """启动 codebuddy 进程,登记到 RUNNING 并开始后台 drain。"""
+ if session_exists(sid):
+ cmd = [codebuddy_path, "--model", model, "--resume", sid]
+ else:
+ cmd = [codebuddy_path, "--model", model, "--session-id", sid]
+ master_fd, slave_fd = pty.openpty()
+ winsize = struct.pack("HHHH", 24, 80, 0, 0)
+ fcntl.ioctl(slave_fd, termios.TIOCSWINSZ, winsize)
+ env = os.environ.copy()
+ env["TERM"] = "xterm-256color"
+ env["COLORTERM"] = "truecolor"
+ env["COLUMNS"] = "80"
+ env["LINES"] = "24"
+ try:
+ proc = subprocess.Popen(
+ cmd, stdin=slave_fd, stdout=slave_fd, stderr=slave_fd,
+ cwd=cwd, env=env, preexec_fn=os.setsid,
+ )
+ finally:
+ try:
+ os.close(slave_fd)
+ except OSError:
+ pass
+ os.set_blocking(master_fd, False)
+ entry = {
+ "proc": proc, "master_fd": master_fd, "sid": sid,
+ "name": task["name"], "cwd": cwd, "model": model,
+ "started_at": time.time(), "buffer": [], "buflen": 0,
+ "ws": None, "done": False, "exit_code": None,
+ "send_lock": asyncio.Lock(),
+ }
+ RUNNING[task["id"]] = entry
+ asyncio.get_running_loop().create_task(_drain(entry))
+ logger.info(f"任务启动(后台常驻): name={task['name']}, cwd={cwd}, mode={'resume' if session_exists(sid) else 'new'}, model={model}, sid={sid}, pid={proc.pid}")
+ return entry
+
+
+def _kill_entry(task_id: str) -> bool:
+ """真正终止任务进程(进程组)并从 RUNNING 移除。"""
+ 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)
+ except Exception:
+ try:
+ proc.kill()
+ except Exception:
+ pass
+ return True
+
+
# ==================== 前端 HTML ====================
# 前后端分离:前端页面单独放在 frontend/index.html,便于维护
FRONTEND_DIR = Path(__file__).resolve().parent.parent / "frontend"
-with open(FRONTEND_DIR / "index.html", "r", encoding="utf-8") as _f:
- HTML_PAGE = _f.read()
+def load_html():
+ # 每次请求时读取,修改前端后刷新浏览器即可生效,无需重启后端
+ with open(FRONTEND_DIR / "index.html", "r", encoding="utf-8") as _f:
+ return _f.read()
# ==================== 后端路由 ====================
@app.get("/")
async def index():
- return HTMLResponse(HTML_PAGE)
+ return HTMLResponse(load_html())
@app.post("/api/login")
tasks = []
for t in data.get("tasks", []):
tt = dict(t)
+ # 四态实时状态:前台已登录 / 后台运行中 / 后台任务待确认 / 空闲
+ if t["id"] in ACTIVE_TASKS:
+ tt["live_status"] = "connected"
+ elif t["id"] in RUNNING:
+ e = RUNNING[t["id"]]
+ if e["proc"].poll() is None:
+ tt["live_status"] = "background"
+ else:
+ # 后台跑完了但当时没人在看 → 提示用户回来确认
+ tt["live_status"] = "done_detached"
+ tt["exit_code"] = e["exit_code"]
+ else:
+ tt["live_status"] = "idle"
tt["active"] = t["id"] in ACTIVE_TASKS
# 兼容旧数据:无 status 字段视为 active
tt.setdefault("status", "active")
async def api_kill(pid: int, request: Request):
if not _auth_ok(request):
return JSONResponse({"ok": False}, status_code=401)
+ # 从 RUNNING 移除对应任务(按 pid 反查)
+ for tid, e in list(RUNNING.items()):
+ if e["proc"].pid == pid:
+ RUNNING.pop(tid, None)
+ break
+ # 杀掉整个进程组(codebuddy 用 setsid 独立成组,子进程也一起杀)
try:
- os.kill(pid, signal.SIGTERM)
+ pgid = os.getpgid(pid)
+ os.killpg(pgid, signal.SIGTERM)
except ProcessLookupError:
pass
+ except Exception:
+ try:
+ os.kill(pid, signal.SIGTERM)
+ except Exception:
+ pass
return JSONResponse({"ok": True})
return
if task_id in ACTIVE_TASKS:
- await websocket.accept()
- await websocket.send_json({"type": "error", "data": "该任务已在另一窗口运行,请先关闭该窗口"})
- await websocket.close()
- return
+ # 可能是上一次连接刚断开、finally 还没来得及 discard 的瞬间竞态;
+ # 若该任务的 ws 已经置空(已 detach),说明旧连接已结束,允许本次重连(会 attach 到同一进程)
+ prev = RUNNING.get(task_id)
+ if prev is not None and prev.get("ws") is None:
+ ACTIVE_TASKS.discard(task_id)
+ else:
+ await websocket.accept()
+ await websocket.send_json({"type": "error", "data": "该任务已在另一窗口运行,请先关闭该窗口"})
+ await websocket.close()
+ return
# 找到 codebuddy 可执行文件
codebuddy_path = None
await websocket.close()
return
- # 决定启动参数:已存在 session 则 resume,否则用固定 session-id 新建
+ # 决定生效模型:任务显式设置 > 默认最新 GLM
sid = task["session_id"]
- # 计算生效模型:任务显式设置 > 默认最新 GLM
with MODELS_CACHE_LOCK:
_default_model = MODELS_CACHE.get("default") or (MODELS_CACHE.get("models") or [""])[0] or "glm-5.2"
model = (task.get("model") or "").strip() or _default_model
- if session_exists(sid):
- cmd = [codebuddy_path, "--model", model, "--resume", sid]
- mode = "resume"
- else:
- cmd = [codebuddy_path, "--model", model, "--session-id", sid]
- mode = "new"
cwd = task["cwd"]
try:
except OSError:
pass # 启动时若目录不可用,交给 codebuddy 报错
+ # attach 到已在后台运行的任务;若断开期间已结束,则稍后汇报并开新会话;否则新启动
+ done_entry = None
+ entry = RUNNING.get(task_id)
+ if entry is not None and entry["done"]:
+ done_entry = entry
+ RUNNING.pop(task_id, None)
+ entry = None
+ if entry is not None and entry["proc"].poll() is None:
+ attached = True
+ else:
+ attached = False
+ entry = _launch_process(task, sid, model, cwd, codebuddy_path)
+
ACTIVE_TASKS.add(task_id)
await websocket.accept()
- logger.info(f"任务启动: name={task['name']}, cwd={cwd}, mode={mode}, model={model}, sid={sid}")
-
- # 创建 PTY
- master_fd, slave_fd = pty.openpty()
- winsize = struct.pack("HHHH", 24, 80, 0, 0)
- fcntl.ioctl(slave_fd, termios.TIOCSWINSZ, winsize)
-
- env = os.environ.copy()
- env["TERM"] = "xterm-256color"
- env["COLORTERM"] = "truecolor"
- env["COLUMNS"] = "80"
- env["LINES"] = "24"
-
- try:
- proc = subprocess.Popen(
- cmd,
- stdin=slave_fd, stdout=slave_fd, stderr=slave_fd,
- cwd=cwd, env=env, preexec_fn=os.setsid,
- )
- except Exception as e:
- os.close(slave_fd)
- await websocket.send_json({"type": "output", "data": f"\r\n\x1b[31m✗ 启动 codebuddy 失败: {e}\x1b[0m\r\n"})
- await websocket.close()
- ACTIVE_TASKS.discard(task_id)
- return
- finally:
+ entry["ws"] = websocket
+ logger.info(f"任务连接: name={task['name']}, mode={'attach' if attached else 'new'}, sid={sid}, pid={entry['proc'].pid}")
+
+ # 汇报状态 + 回放断开期间的输出
+ if done_entry is not None:
+ replay = "".join(done_entry["buffer"])
+ if len(replay) > REPLAY_LIMIT:
+ replay = replay[-REPLAY_LIMIT:]
try:
- os.close(slave_fd)
- except OSError:
+ async with entry["send_lock"]:
+ await websocket.send_json({
+ "type": "reattach", "attached": False, "running": False,
+ "exit_code": done_entry["exit_code"],
+ "elapsed_sec": int(time.time() - done_entry["started_at"]),
+ "replay": replay, "current_task": last_assistant_summary(sid),
+ "pid": entry["proc"].pid,
+ })
+ except Exception:
pass
-
- os.set_blocking(master_fd, False)
- await websocket.send_json({"type": "pid", "pid": proc.pid})
-
- # 更新 last_used
- touch_task(task_id)
-
- async def read_output():
- loop = asyncio.get_event_loop()
- while proc.poll() is None:
- try:
- data = await loop.run_in_executor(None, lambda: os.read(master_fd, 65536))
- if data:
- try:
- text = data.decode("utf-8", errors="replace")
- except Exception:
- text = data.decode("latin-1", errors="replace")
- await websocket.send_json({"type": "output", "data": text})
- else:
- await asyncio.sleep(0.01)
- except BlockingIOError:
- await asyncio.sleep(0.05)
- except OSError:
- break
- except Exception:
- await asyncio.sleep(0.05)
-
- code = proc.poll()
+ elif attached:
+ replay = "".join(entry["buffer"])
+ if len(replay) > REPLAY_LIMIT:
+ replay = replay[-REPLAY_LIMIT:]
try:
- while True:
- try:
- d = os.read(master_fd, 65536)
- if not d:
- break
- await websocket.send_json({"type": "output", "data": d.decode("utf-8", errors="replace")})
- except OSError:
- break
+ async with entry["send_lock"]:
+ await websocket.send_json({
+ "type": "reattach", "attached": True, "running": True,
+ "exit_code": entry["exit_code"],
+ "elapsed_sec": int(time.time() - entry["started_at"]),
+ "replay": replay, "current_task": last_assistant_summary(sid),
+ "pid": entry["proc"].pid,
+ })
except Exception:
pass
- await websocket.send_json({"type": "exit", "code": code})
- logger.info(f"任务结束: name={task['name']}, pid={proc.pid}, code={code}")
+ else:
+ try:
+ async with entry["send_lock"]:
+ await websocket.send_json({"type": "pid", "pid": entry["proc"].pid})
+ except Exception:
+ pass
+
+ touch_task(task_id)
async def write_input():
try:
while True:
msg = await websocket.receive_json()
- if msg.get("type") == "input":
- os.write(master_fd, msg.get("data", "").encode("utf-8"))
+ if msg.get("type") == "kill":
+ _kill_entry(task_id)
+ try:
+ await websocket.close()
+ except Exception:
+ pass
+ return
+ elif msg.get("type") == "input":
+ if entry["master_fd"] is not None:
+ try:
+ os.write(entry["master_fd"], msg.get("data", "").encode("utf-8"))
+ except OSError:
+ pass
elif msg.get("type") == "resize":
cols = msg.get("cols", 80)
rows = msg.get("rows", 24)
- winsize = struct.pack("HHHH", rows, cols, 0, 0)
- try:
- fcntl.ioctl(master_fd, termios.TIOCSWINSZ, winsize)
- except OSError:
- pass
+ if entry["master_fd"] is not None:
+ try:
+ fcntl.ioctl(entry["master_fd"], termios.TIOCSWINSZ,
+ struct.pack("HHHH", rows, cols, 0, 0))
+ except OSError:
+ pass
except Exception:
pass
+ finally:
+ if entry.get("ws") is websocket:
+ entry["ws"] = None
- read_task = asyncio.create_task(read_output())
write_task = asyncio.create_task(write_input())
try:
- await asyncio.wait([read_task, write_task], return_when=asyncio.FIRST_COMPLETED)
+ # 保持连接直到客户端断开(entry["ws"] 被置空)或进程结束
+ while entry.get("ws") is websocket:
+ await asyncio.sleep(1)
+ if entry["done"] and entry["master_fd"] is None:
+ try:
+ async with entry["send_lock"]:
+ await websocket.send_json({"type": "exit", "code": entry["exit_code"]})
+ except Exception:
+ pass
+ break
except Exception:
pass
finally:
- for t in [read_task, write_task]:
- if not t.done():
- t.cancel()
- try:
- os.close(master_fd)
- except Exception:
- pass
- try:
- if proc.poll() is None:
- os.killpg(os.getpgid(proc.pid), signal.SIGTERM)
- try:
- proc.wait(timeout=3)
- except Exception:
- os.killpg(os.getpgid(proc.pid), signal.SIGKILL)
- except Exception:
- pass
- try:
- await websocket.close()
- except Exception:
- pass
+ write_task.cancel()
+ if entry.get("ws") is websocket:
+ entry["ws"] = None
ACTIVE_TASKS.discard(task_id)
- logger.info(f"任务清理完成: name={task['name']}, pid={proc.pid}")
+ # 关键:关闭浏览器(WS 断开)时【不杀进程】,让其后台继续跑(detach)
+ logger.info(f"任务断开(后台保留): name={task['name']}, pid={entry['proc'].pid}")
# ==================== 启动 ====================
.task-card .tc-del { background: rgba(255,107,107,0.12); color: #ff8787; border: 1px solid rgba(255,107,107,0.25); flex: 0 0 56px; }
.task-card .tc-del:hover { background: rgba(255,107,107,0.22); }
.task-card.tc-paused { opacity: 0.6; border-color: rgba(255,230,109,0.3); }
-.task-card .tc-status { font-size: 12px; }
-.task-card.tc-active { border-color: #ffe66d; }
-.task-card.tc-active .tc-status { color: #ffe66d; }
+.task-card .tc-status { font-size: 12px; font-weight: 600; }
+.task-card.tc-connected { border-color: #ffe66d; }
+.task-card.tc-connected .tc-status { color: #ffe66d; }
+.task-card.tc-background { border-color: #4ecca3; }
+.task-card.tc-background .tc-status { color: #4ecca3; }
+.task-card.tc-paused .tc-status { color: #ffd866; }
+.task-card.tc-idle .tc-status { color: #777; }
+.task-card.tc-done { border-color: #ffa94d; }
+.task-card.tc-done .tc-status { color: #ffa94d; }
.task-new {
display: flex; align-items: center; justify-content: center; gap: 8px;
background: rgba(78,204,163,0.06); border: 2px dashed rgba(78,204,163,0.3);
terminalPage.style.display = 'none';
if (page === 'login') { loginPage.style.display = 'flex'; passwordInput.focus(); }
else if (page === 'picker') pickerPage.style.display = 'flex';
- else if (page === 'terminal') terminalPage.style.display = 'flex';
+ else if (page === 'terminal') {
+ terminalPage.style.display = 'flex';
+ // 显示后重新 fit 并聚焦,避免重连后需要先按一次回车才能输入
+ setTimeout(() => { try { fitAddon && fitAddon.fit(); } catch (e) {} try { term && term.focus(); } catch (e) {} }, 30);
+ }
}
function authHeaders() { return { 'Authorization': 'Bearer ' + authToken }; }
filtered.forEach(t => {
const card = document.createElement('div');
const isPaused = (t.status || 'active') === 'paused';
- card.className = 'task-card' + (isPaused ? ' tc-paused' : '') + (t.active ? ' tc-active' : '');
+ // 三态实时状态:前台已登录 / 后台运行中 / 空闲(paused 优先显示已暂停)
+ const ls = t.live_status || (t.active ? 'connected' : 'idle');
+ let statusLine, statusClass;
+ if (ls === 'connected') { statusLine = '● 前台已登录'; statusClass = 'tc-connected'; }
+ else if (ls === 'background') { statusLine = '⚙ 后台运行中'; statusClass = 'tc-background'; }
+ else if (ls === 'done_detached') { statusLine = '⏳ 后台任务待确认' + (t.exit_code != null ? '(退出码 ' + t.exit_code + ')' : ''); statusClass = 'tc-done'; }
+ else if (isPaused) { statusLine = '⏸ 已暂停'; statusClass = 'tc-paused'; }
+ else { statusLine = '○ 空闲'; statusClass = 'tc-idle'; }
+ card.className = 'task-card' + (isPaused ? ' tc-paused' : '') + ' ' + statusClass;
const last = t.last_used ? '上次使用 ' + t.last_used.replace('T', ' ') : '';
// 计算暂停时长(用于判断是否可删除)
let pausedInfo = '';
pausedInfo = '已暂停 ' + elapsedH.toFixed(1) + ' 小时';
canDelete = elapsedH >= 24;
}
- const statusLine = t.active ? '● 运行中' : (isPaused ? '⏸ 已暂停' : '○ 空闲');
+ // statusLine 已按 live_status 在上方计算
card.innerHTML =
'<div class="tc-name">' + escapeHtml(t.name) + '</div>' +
'<div class="tc-cwd">' + escapeHtml(t.cwd) + '</div>' +
term.loadAddon(new WebLinksAddon.WebLinksAddon());
term.open(termContainer);
fitAddon.fit();
+ // 点击/触摸终端区域即聚焦(移动端软键盘需要一次用户手势才会弹出)
+ const focusTerm = () => { try { term.focus(); } catch (e) {} };
+ termContainer.addEventListener('click', focusTerm);
+ termContainer.addEventListener('touchstart', focusTerm, { passive: true });
// 复制支持:有选区时 Ctrl+C 复制到剪贴板(不发中断信号)
// Ctrl+Shift+C 也复制(xterm 默认行为)
if (v && term && term.focus) { try { term.focus(); } catch(e){} }
}
+ function renderReattach(msg) {
+ const sec = msg.elapsed_sec || 0;
+ const m = Math.floor(sec / 60), s = sec % 60;
+ term.writeln('');
+ if (msg.running) {
+ term.writeln('\\x1b[36m🔗 已重新连接 · 任务运行中 · 已运行 ' + m + '分' + s + '秒\\x1b[0m');
+ } else {
+ term.writeln('\\x1b[33mⓘ 任务已结束(退出码 ' + (msg.exit_code == null ? '?' : msg.exit_code) + ')\\x1b[0m');
+ }
+ if (msg.current_task) term.writeln('\\x1b[90m当前/最近任务:' + msg.current_task + '\\x1b[0m');
+ if (msg.replay) {
+ term.writeln('\\x1b[90m──── 以下为断开期间 / 上次的输出 ────\\x1b[0m');
+ term.write(msg.replay);
+ }
+ try { term.focus(); } catch (e) {}
+ }
+
function connectWS() {
if (ws) { ws.close(); ws = null; }
if (!authToken || !TASK_ID) return;
const msg = JSON.parse(event.data);
if (msg.type === 'output' && msg.data) term.write(msg.data);
else if (msg.type === 'pid') pid = msg.pid;
+ else if (msg.type === 'reattach') { pid = msg.pid || pid; renderReattach(msg); }
else if (msg.type === 'error') { term.writeln('\\x1b[31m✗ ' + msg.data + '\\x1b[0m'); }
else if (msg.type === 'exit') { term.writeln('\\x1b[33m⚠ 进程已退出 (code: ' + msg.code + ')\\x1b[0m'); setConnected(false); }
} catch(e) { term.write(event.data); }
window.close();
};
- // 窗口关闭前尽力杀进程
- window.addEventListener('beforeunload', () => {
- if (pid) navigator.sendBeacon && navigator.sendBeacon(API + '/api/kill/' + pid + '?token=' + encodeURIComponent(authToken));
- });
+ // 注意:关闭浏览器不再杀进程——后端会在 WS 断开时保留任务后台继续运行(detach)。
+ // 只有显式点“关闭任务”或“重启”才会通过 /api/kill 杀进程。
// ========== 初始化 ==========
(async () => {