# -*- coding: utf-8 -*- """collabd —— **协作守护程序**(通用版 · 一个进程承担"派活之外的全部功能") 配套:`multi-session-collab` 技能。配置见同目录 `collabd.config.json`(模板 `collabd.config.example.json`)。 **它承担**:① 通路/服务自愈(连续 N 次不通才动手·带冷却 ⇒ 防误杀)② 监控+**证据分级**+推进判定 ③ 机械靶点判定 ④ **真空判定**(需求未完成 ∧ 无人执行) ⑤ 断链/中断告警 ⑥ **机制体检**(仅空闲) ⑦ **任务图**(可派/等待/关键路径/**可派未派=浪费**) ⑧ 视图覆写(看板/摘要) ⑨ **单例**守护 ⑩ **唤醒回路**(2026-09-29 立 · 用户拍板):`NEXT.md` 出现且闸门开 ⇒ **把「有活」投给活会话** —— 走网关官方 `POST /api/v1/sessions/{id}/reply`(**投递不夺 ACP writer** · 桌面不受干扰)。 **⛔ 它不做**:**派活 / 开新会话** —— 那必须由会话做(且"新建自动化"受白名单+确认制约束)。 边界:⛔ 读取宿主库只读 ·⛔ 不写宿主状态 ·⛔ 不碰别人的文件 · 全静默 · fail-open ⚠️ **唯一写操作 = ⑩ 的那一次 reply 投递**(口令**只进程内用**:⛔ 不落盘 ⛔ 不进日志 ⛔ 不回显); 四道闸门全满足才投(有 `NEXT.md` ∧ `gate=free` ∧ **内容哈希≠上次** ∧ 距上次 ≥`wake_min_gap` 秒), 每次投递**追加登记** `inbox/wakeups.jsonl`(`{ts, port, sessionId, hash, http}` ⇒ 可列出/可删)。 🔴 认口口径:同 cwd 可能有多个口,**取「第一个带活会话的」**(按 uptime 降序里第一个 `sessionId` 非空者)—— ⛔ 不能只取 uptime 最大那个(实测它可能没有活会话 ⇒ 白等)。 用法(🔴 2026-09-30 定型 → **同日按用户改口修正:允许自动任务**;⛔ 仍不靠常驻): 🔴 **判据:投递这个动作只有一个出口(=本程序)。** 谁"拨这一下"有两种 —— ① **宿主钩子**(有会话在动时,事件驱动,见 `--tick` 处注释) ② **自动任务**("谁都没动"时的**时钟缺口**,经 `goalctl.py wake` 调 `--tick`) 用户原话:「自动任务(符合条件时)是通过 协作程序去唤醒 主会话,这样流程统一」。 ⛔ 差别只在"谁拨",判断与投递**仍然只在本程序一处** ⇒ 流程不分裂。 `--once` 协作程序:读库/判定/信号/看板(**⛔ 不投递**) `--tick` 监督程序:读队列 ⇒ **投给主会话**(单条+握手/三条件心跳)—— 由**宿主钩子**或**自动任务**唤起 ⚠️ 判读 `deliver=` 字段:投成(http 2xx)/`target-busy`(主会话**在跑** ⇒ ⛔ 别降级)/ `main-not-live`·`no-main-session`·`no-live-session`(**没有可唤醒的对象** ⇒ 才降级为"自己干") `--where` 路径自证 (无参数 = 常驻模式,**旧形态**,保留但不再是投递的前置) """ from __future__ import annotations import hashlib import json import os import re import socket import sqlite3 import subprocess import sys import time import urllib.error import urllib.request from pathlib import Path HERE = Path(__file__).resolve().parent DEFAULTS = { "workspace": "", "inbox": "tmp/supervise-inbox", "live": "", # 视图(看板类产物);空 ⇒ inbox/realtime.md "taskgraph": "", # 任务图 json;空 ⇒ inbox/taskgraph.json "lines": {}, # {"<目录名>": "<友好名>"} 用于靶点判定 "goal_docs": [], # [["显示名", "相对 workspace 的路径"], ...] "targets": {}, # {"<目录名>": [["显示名","glob 相对该目录"], ...]} "host_db": "", # 宿主库路径;空 ⇒ $CODEBUDDY_CONFIG_DIR/workbuddy.db "shim_port": 0, # 服务端口探针(0 ⇒ 关闭) "shim_streak": 3, "client_entry": "", # 端口不通时拉起的启动器(.mjs/.py/.exe) "client_runner": "", # 用什么跑启动器(如 node 绝对路径);空 ⇒ 直接执行 "singleton_port": 20099, "interval": 10, "idle_min": 12, "stuck_min": 30, "vacuum_min": 5, "health_every": 300, "client_cooldown": 600, # ⑩ 唤醒回路(2026-09-29 立)。⚠️ wake_text 内**一律用「」**,⛔ 别写 ASCII 双引号(P11 引号坑) "wake_enable": True, "wake_min_gap": 180, # 秒:同一件两次投递的最小间隔 "wake_max_ports": 60, # 最多探多少个 loopback 监听口(认网关口用) "wake_text": "读 tmp/supervise-inbox/NEXT.md 处理那一条;收尾时若 NEXT.md 还在才接续(一次只做这一条)", } # 🔴 部署配置**属于使用方**(用户 2026-09-30 定则:「技能就是技能 程序就是程序, # 谁用产生的文件 放在他自己那里」)⇒ **技能目录里⛔ 不放生产配置**。 # 查找顺序(前面的赢): # ① 环境变量 `COLLABD_CONFIG` —— 钩子/启动器显式指定(首选) # ② `<工作区>/.workbuddy/collab/collabd.config.json` —— 使用方的**标准落点** # ③ 都没有 ⇒ 用 DEFAULTS + `COLLABD_WORKSPACE`/cwd,并**明确记一条"未找到配置"**(⛔ 不静默) CFG_USED = "" CFG_MISSING = True def _cfg_candidates() -> list: out = [] env = os.environ.get("COLLABD_CONFIG") if env: out.append(Path(env)) ws = os.environ.get("COLLABD_WORKSPACE") or os.getcwd() out.append(Path(ws) / ".workbuddy" / "collab" / "collabd.config.json") return out def load_cfg() -> dict: global CFG_USED, CFG_MISSING cfg = dict(DEFAULTS) for p in _cfg_candidates(): if p.is_file(): try: cfg.update(json.loads(p.read_text(encoding="utf-8"))) CFG_USED, CFG_MISSING = str(p), False break except Exception: pass if not cfg["workspace"]: cfg["workspace"] = os.environ.get("COLLABD_WORKSPACE") or os.getcwd() return cfg C = load_cfg() WS = Path(C["workspace"]) INBOX = WS / C["inbox"] def _rp(rel: str, default: Path) -> Path: """🔴 把配置里的**相对路径解析到 workspace 下**(⛔ 别按脚本 CWD 解析 —— 踩过:视图被写进技能目录)""" p = Path(rel) if rel else None if p is None: return default return p if p.is_absolute() else (WS / p) LIVE = _rp(C["live"], INBOX / "realtime.md") TG = _rp(C["taskgraph"], INBOX / "taskgraph.json") STALL, VACUUM, READY = INBOX / "STALL.md", INBOX / "VACUUM.md", INBOX / "READY.md" HEALTH = INBOX / "体检报告.md" STATE = INBOX / "collabd-state.json" # ⛔ 日志**不进技能目录**(技能是只读的能力件)⇒ 默认落在使用方的工作区内 LOG = _rp(C.get("log") or "", INBOX / "_collabd.log") def log(m: str) -> None: try: LOG.parent.mkdir(parents=True, exist_ok=True) # ⚠️ 日志落点在**使用方**,目录可能还没建 with open(LOG, "a", encoding="utf-8") as f: f.write("[%s] %s\n" % (time.strftime("%Y-%m-%d %H:%M:%S"), m)) except Exception: pass # 🔴 2026-10-01 体检修加:**同一句结论每轮都写一遍 ⇒ 日志被噪声淹没**(实测:连续 40+ 行同一句 # "acceptance_state 没有有效判据",真正的新事件反而看不见)⇒ 同一 key 在 gap 秒内只写一次。 _LOG_LAST: dict = {} def log_throttled(key: str, gap: float, m: str) -> None: """同一 `key` 在 `gap` 秒内**只写一次**(⛔ 不改语义,只是降噪)。""" _now = time.time() if _now - float(_LOG_LAST.get(key) or 0) < float(gap): return _LOG_LAST[key] = _now log(m) if CFG_MISSING: # ⛔ **找不到配置时不许落任何文件** —— 那一刻我们根本不知道"使用方的地方"在哪; # 旧写法 `log(...)` 会按默认落点写进 **cwd**,而 cwd 常常就是技能目录 # ⇒ 恰好又造出它要避免的污染(2026-09-30 实测:一跑就往技能里长出 `tmp/supervise-inbox/_collabd.log`)。 # ⇒ 只打 **stderr**(调用方若 devnull 掉了就当静默;人手动跑时看得见)。 sys.stderr.write("⚠️ collabd:未找到部署配置(COLLABD_CONFIG 未设,且 %s 不存在)\n" " ⇒ 已拒跑,⛔ 不会往任何目录写文件。\n" % (Path(os.environ.get("COLLABD_WORKSPACE") or os.getcwd()) / ".workbuddy" / "collab" / "collabd.config.json")) def _mt(p: Path) -> float: try: return p.stat().st_mtime except Exception: return 0.0 def _db(): cfg_dir = os.environ.get("CODEBUDDY_CONFIG_DIR") or r"E:\ProgramData\.workbuddy" db = Path(C["host_db"]) if C["host_db"] else Path(cfg_dir, "workbuddy.db") return sqlite3.connect("file:%s?mode=ro" % str(db).replace("\\", "/"), uri=True, timeout=4) # ── 只读取数(宿主库) ─────────────────────────────────────── def fetch() -> dict: d = {"runs": [], "wm": 0, "pending": [], "upcoming": [], "soon": 0, "working": [], "self_working": False, "sessions": [], "locks": [], "due": [], "have": set()} now_ms = int(time.time() * 1000) try: con = _db() rows = list(con.execute("select rowid,status,thread_title,automation_id from automation_runs " "order by rowid desc limit 8")) d["runs"] = [{"rowid": r[0], "status": r[1], "title": (r[2] or "").strip(), "aid": (r[3] or "")[:8]} for r in rows] d["wm"] = max([r["rowid"] for r in d["runs"]], default=0) for r in con.execute("select id,name,cwds,schedule_type,next_run_at,scheduled_at from automations " "where status='ACTIVE' and (deleted_at is null or deleted_at='')"): # 🔴 2026-09-29 修(P1·实测):原实现把「ACTIVE ∧ next_run_at>now」全算 `pending`, # 而两条**周期**自动化(日报 / 体检)恒满足 ⇒ `pending` 恒非空 ⇒「没有会话在跑」 # 那条分支**永远不可达** ⇒ `STALL.md` 从未生成(实测两日零生成)。⇒ 拆开: # · `pending` = **一次性**排活(刚建好、还没跑的下一棒 ⇒ 用于判「只在排活、未见成果」) # · `upcoming` = **周期**自动化(⛔ 只展示,**不代表有人在跑**) if (r[4] or 0) > now_ms: if (r[3] or "").lower() == "recurring": d["upcoming"].append(r) else: d["pending"].append(r) if (r[4] or 0) <= now_ms + 60 * 60 * 1000: # 🔴 「**未来 1 小时内**会不会有人自己动起来」—— 这才是判"确定性静默"的口径。 # ⛔ 不是"全库有没有排期":日报 / 体检那两条明天才跑,与**本线**毫无关系 # (实测 16:34→20:07:全库有 2 条排期,本线 0 条,于是静默 4 小时)。 d["soon"] += 1 if r[5]: d["due"].append((str(r[0]), r[1] or "", r[5])) # 🔴 2026-09-29 修(P1·实测):**唯一可得的「有人在跑」信号 = `sessions.status='working'`**。 # ⛔ 不能用 `automation_runs.status`(实测 54 行**全 ACCEPTED**、零 IN_PROGRESS) # ⛔ 不能用一次性棒的 `next_run_at`(实测**全是 0/None**,跑完即归零) for r in con.execute("select id, coalesce(nullif(custom_title,''),title,'(未命名)'), " " coalesce(last_activity_at, updated_at), coalesce(cwd,'') " "from sessions where status='working'"): # 🔴 2026-09-29:**排除观察者自己**。否则「主会话自己在跑」也会被算成「有人在推进」 # ⇒ 停滞 / 真空永远发现不了(用户问「是不是又发呆了」答不出来的根因之一)。 wid = str(r[0] or "") if SELF_SID and wid.lower().startswith(SELF_SID[:8].lower()): d["self_working"] = True continue d["working"].append({"id": wid, "name": str(r[1] or "")[:34], "at": r[2], "cwd": str(r[3] or "")}) d["have"] = {str(r[0]) for r in con.execute("select distinct automation_id from automation_runs") if r[0]} # 🔴 2026-09-29(回架构·去自造状态源):**「多久没成果」的权威在宿主库里** —— # `automation_runs.updated_at`(实测有该列,毫秒 epoch,形如 '1790667377928')。 # ⛔ 不要让程序自己维护 `last_progress_at` 这类时钟(那就是"第二状态源":会漂、要清理、 # 丢了就永久失真 —— 今天整天的毛病都长在这类自造状态上)。 try: row = con.execute("select updated_at from automation_runs " "where status='ACCEPTED' and coalesce(thread_title,'')<>'' " "order by updated_at desc limit 1").fetchone() if row and row[0]: v = float(row[0]) d["last_result_at"] = v / (1000.0 if v > 1e11 else 1.0) except Exception as e2: log("last_result_at 取失败 %s" % e2) # ⚠️ sessions 表没有 name 列(正确列:title/custom_title/last_activity_at/unread)——踩过一次 for r in con.execute( "select id, coalesce(nullif(custom_title,''), title, '(未命名)'), " " coalesce(last_activity_at, updated_at), status, unread from sessions " "order by coalesce(last_activity_at, updated_at) desc limit 10"): ts, t = r[2], 0.0 if isinstance(ts, (int, float)): t = float(ts) / (1000.0 if float(ts) > 1e11 else 1.0) elif isinstance(ts, str) and ts: try: t = time.mktime(time.strptime(ts[:19], "%Y-%m-%dT%H:%M:%S")) except Exception: t = 0.0 d["sessions"].append({"id": str(r[0] or "")[:8], "name": str(r[1] or "")[:36], "age": (time.time() - t) / 60 if t else 1e9, "tag": ("未读 %s" % r[4]) if r[4] else (r[3] or "")}) con.close() except Exception as e: log("fetch err %s" % e) return d def probe_port() -> bool: if not C["shim_port"]: return True try: s = socket.create_connection(("127.0.0.1", int(C["shim_port"])), timeout=1.5) s.close() return True except Exception: return False def targets() -> tuple[bool, list[str]]: rows, ok_all = [], True for line, items in (C["targets"] or {}).items(): for label, pat in items: base = Path(WS).parent / line if not Path(line).is_absolute() else Path(line) hits = sorted(base.glob(pat), key=_mt) if pat else [] ok = bool(hits) ok_all = ok_all and ok rows.append("- %s **%s** %s" % ("✅" if ok else "⛔", line, label)) return ok_all, rows # ── 任务图(防干等) ──────────────────────────────────────── def taskgraph(cur: list, busy: bool) -> dict: try: g = json.loads(TG.read_text(encoding="utf-8")) except Exception: return {} raw = g.get("nodes", []) # 🔴🔴 2026-10-01 加(治「上报完成 ⇒ 队首不消失」)—— **台账优先**。 # 病根:`tasks.json` 的注释自称「**唯一权威**(队列本体)」,会话「上报完成」也只写它; # 但**本仓没有任何代码回写 `交付物/任务图.json`**(全仓 grep:TG 只被读、从不被写) # ⇒ 任务图是**纯手维护**的分解件 ⇒ 两边必然漂。 # 实测(M5):台账 done、上报单 done、产物已核对(wakeups.jsonl 第 82 行 ok:true), # 任务图仍留 `running` ⇒ ① 它永远是队首 ② `READY.md`/`STALL.md` 每小时喊同一件白跑 # ③ 控制台判「目标未完成」⇒ 整条链被一个**早已做完的件**卡死。 # ⇒ 判定改用**有效状态**:`图 done ∪ 台账 done` 都算 done(`deps` 解析同理)。 # ⛔ 不回写任务图(它仍是人可手改的分解件);只让「算不算做完」以台账为准。 try: _led = _load_tasks() except Exception: _led = {} def _eff(n): if (n.get("status") or "") == "done": return "done" return "done" if str(((_led.get(n.get("id")) or {}).get("state")) or "") == "done" \ else (n.get("status") or "") eff = {n.get("id"): _eff(n) for n in raw} ns = {n.get("id"): n for n in raw} out = {"ready": [], "waiting": [], "critical": list(g.get("critical_path") or []), "waste": False, "nodes_raw": raw} for n in raw: if eff.get(n.get("id")) == "done": continue undone = [x for x in (n.get("deps") or []) if eff.get(x) != "done"] (out["waiting"] if undone else out["ready"]).append( (n.get("id"), (n.get("title") or "")[:40], n.get("line"), undone) if undone else (n.get("id"), (n.get("title") or "")[:40], n.get("line"), n.get("status"))) # 🔴 2026-09-29 加:**受阻件不算「可派」** —— 否则 `READY.md` 与机器摘要会把受损件当 # 「可派未派(浪费)」反复喊,与 `queue.json` 的 blocked **自相矛盾**(实测:N9 被标受阻后, # READY.md 仍在喊「N9 可派」)。 try: blk = json.loads(BLOCKED.read_text(encoding="utf-8")) if BLOCKED.exists() else {} if not isinstance(blk, dict): blk = {} except Exception: blk = {} if blk: out["blocked"] = [(r[0], blk.get(r[0], "")) for r in out["ready"] if r[0] in blk] out["ready"] = [r for r in out["ready"] if r[0] not in blk] out["waste"] = bool(out["ready"]) and not cur and not busy return out # ── 判定 ──────────────────────────────────────────────────── def verdicts(d: dict, up: bool, st: dict) -> dict: now = time.time() fresh = [r for r in d["runs"] if r["rowid"] > int(st.get("rowid") or 0)] done = [r for r in fresh if r["status"] == "ACCEPTED" and r["title"]] started = [r for r in fresh if r not in done] # 🔴 2026-09-29 修:原写成 `... if st.get("aids") else []` ⇒ `aids` 一旦为空(实测刚发生过) # 就**永远判不出「新增排活」** ⇒「只在排活、未见成果」这条警戒失效。改成对空集比较。 newp = [r for r in d["pending"] if str(r[0]) not in set(st.get("aids") or [])] # 🔴 2026-09-29 修(P3·实测):探针**翻转不是推进** —— 原实现任何翻转(含 通→不通)都重置 # `last_progress_at` ⇒ 探针每次抖动都把「无成果时长」清零 ⇒ 卡住 / 脱节告警**永不触发**。 # 🔴 2026-09-29 再修(P3b · 20:07 实测):**「恢复」也不得重置** —— 16:31→20:07 整段停滞里 # 程序只跑过一轮(20:07),恰好赶上探针 不通→通 ⇒ idle 被清成 **1 分钟**,摘要写「1 分钟无 # 成果」而不是「4.5 小时」⇒ **连喊一声都没有**。⇒ `idle` **只认真成果**;探针状态变化降级为 # 后置说明(服务层,⛔ 不算成果、⛔ 不清零)。 rec = (st.get("up") is False and up is True) lost = (st.get("up") is True and up is False) # 🔴 权威顺序:本轮有新成果 ⇒ now;否则 **宿主库最后一次有结论的 run**(`last_result_at`); # ⛔ 只有读库失败时才退回 `state.last_progress_at`(那是旧的自维护时钟,已降级为兜底, # 且一取到库值就被忽略)。 src_pa = float(d.get("last_result_at") or 0) if done: pa = now elif src_pa: pa = src_pa else: pa = float(st.get("last_progress_at") or now) idle = (now - pa) / 60 def _dur(m): return ("%.0f 分钟" % m) if m < 90 else ("%.1f 小时" % (m / 60.0)) if done: v, ic = "🟢 **真成果** —— %d 条有结论的完成" % len(done), "ok" elif newp: v, ic = "🟡 **只在排活,未见新成果**(新增 %d 条接续任务·排期≠成果)" % len(newp), "wait" elif d["working"]: v, ic = "🟡 有会话在跑,暂无新成果(在跑 %d 个;%s 无成果)" % ( len(d["working"]), _dur(idle)), "wait" else: v, ic = "🔴 **中断/脱节** —— 没有任何会话在跑,且 %s 无成果" % _dur(idle), "stall" if rec or lost: v += ";探针 %s(服务层变化,⛔ 不算成果、不清零)" % ("不通→通" if rec else "通→不通") vs = float(st.get("vacuum_since") or 0) vac = (not up) and not busy_(d) vs = (vs or now) if vac else 0.0 # 🔴 2026-09-29(回架构):**「零排期」是一等事实,不是背景噪音**。 # 地基(顶层设计 §1/§2):宿主不提供常驻 ⇒ 一切"等待"只能靠**宿主排期** ⇒ # **排期数 = 0 ⇒ 没人发消息时"确定性静默"**(实测 16:31→20:07 四小时零运行,根因即此)。 # ⇒ 单独成一档告警:不等 idle 攒够,直接喊。 zero_sched = int(d.get("soon") or 0) == 0 if zero_sched and not d["working"] and not done: v += ";⛔ **未来 1 小时内零排期** ⇒ 不等人发消息就是**确定性静默**" ic = "stall" alert = None if zero_sched and not d["working"] and not done: alert = "确定性静默(未来 1 小时内**零排期** + 无人接活)" elif ic == "stall" and idle > float(C["idle_min"]): alert = "脱节(无活无人)" if not d["working"] else "中断(有活没人接)" elif ic == "wait" and idle > float(C["stuck_min"]): alert = "卡住(有会话在跑但 %.0f 分钟无成果%s)" % (idle, ",且探针不通" if (lost or not up) else "") return {"fresh": fresh, "done": done, "started": started, "newp": newp, "up": up, "busy": busy_(d), "verdict": v, "icon": ic, "idle": idle, "alert": alert, "vacuum": vac, "vdur": (time.time() - vs) / 60 if vs else 0.0, "vs": vs, # ★ 2026-10-01 加:`digest_text()` 要按「段落」重讲这两条事实 —— 原实现只把它们 # 拼进 `verdict` 的 `;`/`⇒` 长串里 ⇒ 必须把**原始事实**带出来,⛔ 不靠字符串反解。 "zero_sched": zero_sched, "probe": ("rec" if rec else ("lost" if lost else "")), "pa": pa, "wm": d["wm"], "cur": [str(r[0]) for r in d["pending"]]} def busy_(d: dict) -> bool: """🔴 2026-09-29 修(实测):原实现两个信号**恒为假** —— `d["locks"]` 从未被填充; `automation_runs.status` 实测 54 行**全 ACCEPTED**(零 IN_PROGRESS)⇒ `busy` 恒 False ⇒ ① 体检 `health()` 从不跳过 ② `真空` / `可派未派` 判定失真。 改用实测唯一可得的「有人在跑」信号:`sessions.status='working'`(含本会话 ⇒ 语义正确: 「有会话在执行」就不算真空)。""" return bool(d["locks"]) or bool(d.get("working")) def health(d: dict, V: dict) -> dict: out = {"skipped": V["busy"], "issues": []} if V["busy"]: return out cut = time.strftime("%Y-%m-%dT%H:%M", time.localtime(time.time() - 15 * 60)) for aid, nm, sat in d["due"]: s = (sat or "")[:16] if s and s < cut and aid not in d["have"]: out["issues"].append("**哑火**:`%s`(%s)到点未触发" % (aid[:8], nm[:30])) f = sum(1 for r in d["runs"] if "抢锁失败" in (r["title"] or "")) if f >= 3: out["issues"].append("近 %d 条里 %d 次抢锁失败(摩擦,非故障)" % (len(d["runs"]), f)) if len([r for r in d["runs"] if r["status"] == "IN_PROGRESS"]) > 1: out["issues"].append("同时多个 IN_PROGRESS(疑似并发)") for nm, rel in (C["goal_docs"] or []): if not (WS / rel).exists(): out["issues"].append("**关键件缺失**:%s" % nm) return out # ── 自愈(带去抖) ─────────────────────────────────────────── QUEUE, QUEUE_MD = INBOX / "queue.json", INBOX / "queue.md" CLAIMS, STALE = INBOX / "claims", INBOX / "claims-stale" DOING_TTL = 20 * 60 # 秒:doing 超过多久没人续 ⇒ 判卡死并回退(⛔ 判活优先,见 _reap_dead) BLOCKED = INBOX / "blocked.json" # {"<节点id>": "<受阻原因·谁在等>"} —— 受阻件 ⛔ 不当队首 _SID_RE = re.compile(r"[0-9a-fA-F]{8}-[0-9a-fA-F]{4}-[0-9a-fA-F]{4}-[0-9a-fA-F]{4}-[0-9a-fA-F]{12}") SID_DEAD = {"completed", "error", "archived"} # 实测取自 sessions.status 取值域 SELF_SID = os.environ.get("CODEBUDDY_SESSION_ID", "") or "" # 观察者自己 —— 判「有没有人在推进」时 ⛔ 必须排除 # 🔴 主会话的**显式前缀**(用户 2026-09-30 定名:「主会话前缀 可以叫 主控」)。 # ⚠️ 实测现役主会话写成 `主控 · 协作机制 · …` —— **中点分隔、无方括号** # ⇒ `parse_session_name()` 必须把这种形态也认成 main(否则主会话被判"角色未知")。 MAIN_PREFIX = "主控" def _reap_dead(doing: dict) -> dict: """认领持有人**是否还在**(🔴 P5 治本:claims 与「人」绑定)。 契约:`claims//holder` = `<会话名>@<线>@<完整会话id>`(第 3 段可省)。 · holder 里有会话 id ⇒ 查宿主库;状态是**终结态**(completed/error/archived)⇒ 立即把该 claim 移到 `claims-stale/`(⛔ 不删,可追溯)。实测病根:棒已终结、claim 还在 ⇒ 队首被僵尸占住; 等 20 分钟 TTL 回退后队首又**复活** ⇒ 同一件被反复派(N17、N9 各复现一次)。 · ⛔ 无会话 id / 读库失败 / 状态未知 ⇒ **一律不猜**,交给 TTL 兜底(行为=修前)。 """ if not doing: return doing sids = {} for k, v in doing.items(): m = _SID_RE.search(v or "") if m: sids[k] = m.group(0).lower() if not sids: return doing try: con = _db() rows = list(con.execute("select id, status from sessions")) con.close() except Exception as e: log("reap 跳过(宿主库读不到 ⇒ 回退 TTL):%s" % e) return doing smap = {str(r[0] or "").lower(): str(r[1] or "").lower() for r in rows} for k, s in list(sids.items()): if smap.get(s) not in SID_DEAD: continue # 在跑 / 状态未知 ⇒ 不动 try: STALE.mkdir(parents=True, exist_ok=True) p = CLAIMS / k if p.is_dir(): p.rename(STALE / ("%s-持有会话已终结-%s" % (k, time.strftime("%H%M%S")))) doing.pop(k, None) log("claim 出队 %s(持有会话 %s = %s)" % (k, s[:8], smap.get(s))) except Exception as e: log("reap 移出失败 %s:%s" % (k, e)) return doing def queue_view(T: dict, d: dict | None = None) -> dict: """🔴 **严格队列**(2026-09-29 用户要求:"严格用队列的方式处理,避免打架,处理好一个再处理下一个")。 语义(与"域锁"分工不同:**队列管"谁做下一件";域锁管"能不能动这个资源"**): ① **调度串行** —— 一次只呈现**一个队首**;取件必须**原子**(`mkdir claims/` 成功者得) ② **同线互斥、跨线并行** —— 队列里标 `line`:同线已有 doing ⇒ 该线其他件**不呈现为队首** ③ **卡死回退** —— claim 目录超 `DOING_TTL` 未被续 ⇒ 移到 `claims-stale/`(⛔ 不删,可追溯)并记日志 ④ ⛔ 本程序**不派活**:只"排队 + 呈现队首 + 兜卡死";**取件与派活由会话做**。 """ try: CLAIMS.mkdir(parents=True, exist_ok=True) except Exception: pass now = time.time() # ① 卡死回退(只看 claim 目录的 mtime) for c in list(CLAIMS.iterdir()) if CLAIMS.exists() else []: try: if c.is_dir() and now - c.stat().st_mtime > DOING_TTL: STALE.mkdir(parents=True, exist_ok=True) c.rename(STALE / c.name) log("claim stale -> %s(超 %d 分钟未续,已回退)" % (c.name, DOING_TTL // 60)) except Exception: pass doing = {} for c in (CLAIMS.iterdir() if CLAIMS.exists() else []): if c.is_dir(): h = c / "holder" doing[c.name] = (h.read_text(encoding="utf-8", errors="replace").strip() if h.exists() else "?") # 🔴 2026-09-29 修(P5):**持有人失活 ⇒ 立即出队**(⛔ 不再干等 20 分钟 TTL)—— # 病根:claim 只记「有人取了」,不记「那个人还在不在」⇒ 队首被僵尸占住 / TTL 回退后复活。 doing = _reap_dead(doing) # 🔴 2026-09-29 修(P5·`holder` 契约):契约是 `<会话名>@<线>[@<会话id>]`(见 queue.md 模板)。 # 原实现取 `split("@")[0]` = **会话名**,却拿去比 **线名** ⇒ 恒不相等 ⇒ # **「同线互斥」从未生效** ⇒ 同一条线可被同时派多件(打架)。取第 2 段才是线; # 缺第 2 段 ⇒ 回退按节点 id 反查它是哪条线。 line_of = {n.get("id"): n.get("line") for n in (T.get("nodes_raw") or [])} # 🔴 2026-09-29(回架构 · 去自造状态源):「哪条线被占」**由宿主库推导** —— # `sessions.status='working'` 的 `cwd` 末段就是线名(与任务图 `line` 同构)。 # ⛔ 不再依赖解析自己写的 `holder` 文本(那正是"第二状态源":会漂、要清理、今天整天的 # 缺陷都长在这类自造状态上)。`claims` 降级为**仅给程序看的投影提示**,丢了不影响判定。 busy_lines = set() for w in (d or {}).get("working") or []: seg = str(w.get("cwd") or "").replace("\\", "/").rstrip("/").split("/")[-1] if seg: busy_lines.add(seg) for cid, v in doing.items(): # 兜底:宿主库取不到时才看自造件 parts = [x.strip() for x in v.split("@")] busy_lines.add((parts[1] if len(parts) > 1 and parts[1] else "") or line_of.get(cid) or "") busy_lines.discard("") crit = list(T.get("critical") or []) # 🔴 优先级:**"在关键路径的上游闭包内"⇒ 0**(做它能解锁关键路径);否则 1 # ⛔ 不能只看"是不是关键路径节点本身"(那样会漏掉它的前置,反而去干无关的活) crit_up = set(crit) changed = True while changed: # 迭代求上游闭包 changed = False for n in (T.get("nodes_raw") or []): if n.get("id") in crit_up: continue if any(d in crit_up for d in (n.get("deps") or [])): crit_up.add(n.get("id")) changed = True pend, head = [], None for nid, title, line, status in (T.get("ready") or []): if nid in doing: continue num = int("".join(ch for ch in nid if ch.isdigit()) or 9999) # N6 < N7 < N11 pend.append({"id": nid, "title": title, "line": line, "status": status, "prio": 0 if nid in crit_up else 1, "num": num}) pend.sort(key=lambda x: (x["prio"], x["num"])) # 🔴 2026-09-29 加(P4·队列生命周期):**受阻件不得当队首**。否则队首永远是它、`NEXT.md` # 永远是它 ⇒ 每次唤醒都白跑一遍同一个受阻件(实测:N9 被泄漏锁挡住后反复复活)。 # 受阻清单由**会话**维护:`tmp/supervise-inbox/blocked.json` = {"<节点id>": "<原因·谁在等>"} # (⛔ 不改任务图:这是"排队状态",不是"任务状态") try: blocked = json.loads(BLOCKED.read_text(encoding="utf-8")) if BLOCKED.exists() else {} if not isinstance(blocked, dict): blocked = {} except Exception: blocked = {} if blocked: pend = [p for p in pend if p["id"] not in blocked] for p in pend: # ② 同线互斥 ⇒ 只挑一个未被占线的作队首 if p["line"] not in busy_lines: head = p break # 🔴 闸门(gate):用**钩子事件**驱动"继续读下一条",⛔ 不靠轮询 # 规则:有 claim(有人在做)且 `gate-done.stamp` 不新于该 claim ⇒ gate=busy(处理期间不放行); # 收尾钩子会更新 stamp ⇒ gate=free ⇒ 放行下一条(写 NEXT.md,一次只写一条)。 GATE_DONE = INBOX / "gate-done.stamp" NEXT_MD = INBOX / "NEXT.md" gate = "free" try: claim_mt = max([c.stat().st_mtime for c in CLAIMS.iterdir() if c.is_dir()] or [0]) done_mt = GATE_DONE.stat().st_mtime if GATE_DONE.exists() else 0 if doing and done_mt <= claim_mt: gate = "busy" except Exception: pass out = {"head": head, "doing": doing, "pending": len(pend), "blocked_lines": sorted(busy_lines), "blocked": blocked, "gate": gate, "conflict": bool(len(doing) > 1)} try: if gate == "free" and head: NEXT_MD.write_text( "# 📌 下一条(**只此一条** · 处理期间闸门关闭)\n\n- 条目:**%s** %s\n- 线:`%s`\n- 状态:%s\n\n" "## 主会话该做的\n" "1. 抢锁(按此件涉及的文件声明域)\n" "2. `mkdir \"…/tmp/supervise-inbox/claims/%s\"` **原子取件**;建不成 ⇒ 别人取了 ⇒ ⛔ 停\n" "3. 立刻写 `claims/%s/holder` = `<会话名>@<线>@<本会话 id>`(**第 3 段别省**)—— 这是\n" " 「谁在做」的唯一权威;持有人失活程序会**立即出队**,⛔ 不写就只能等 20 分钟超时兜底\n" "4. 派活(白名单「派活」⇒ 免确认)**或**自己做完\n" "5. 更新任务图(**有可核对产物才标 done**)⇒ **删 `claims/%s` 出队**\n" "6. 收尾(你的 SessionEnd 钩子会自动通知程序放行下一条)\n" "7. 🔴 **收尾=自己判本线缺口并接上**(顶层设计 §4「收尾自判」· 定稿 §4.2 必含条款):\n" " ⚠️ **派活命名规范**:新棒的自动化名必须用**两级前缀** `[协作]-[<类别>]-<具体>`\n" " —— **第 2 级要带方括号**,值取 `goal.json` 的 **`topics`(任务类别)**,⛔ **不是 `short`**。\n" " ⛔ 写成 `[协作]-<某个昵称>-…`(漏方括号、或用简称)⇒ `parse_session_name()` 与\n" " `_in_project()` **都认不出**它 ⇒ 它会被**静默漏管**(不进协作会话层、`--ready-next` 也不算它在跑)。\n" " 做完就判「我这条线还有没有缺口」⇒ 有 ⇒ **在收尾那一刻**写下一行一次性 `automations`\n" " (排期 = 现在 + 3~4 分钟)—— ⛔ 不叫「监管会话」(监管棒这个**角色**已退役)。\n" " ⚠️ 棒中途死掉时这步不会发生 ⇒ 链条会静默断掉 ⇒ 靠**每小时兜底心跳**接(定稿 §3②)。\n" "\n> ⛔ 一次只做这一条:**闸门 gate=%s**;遇阻 ⇒ 写 `NEED-USER.md` 并**明确喊「需用户介入」**,\n" "> 同时把该件写进 `blocked.json`({id: 原因})—— 否则它会一直当队首、每次唤醒都白跑。\n" % (head["id"], head["title"], head["line"], head["status"], head["id"], head["id"], head["id"], gate), encoding="utf-8") elif gate == "free" and blocked: # 🔴 2026-09-29 补(治「受阻即静默」):**受阻不是没活** —— 只是没人能开工。 # 原实现选不出队首就删 NEXT.md ⇒ 唤醒回路的触发条件不再满足 ⇒ 变成 # 「有活 ∧ 谁都动不了 ∧ 程序一声不响」(实测 16:31→20:07 零投递)。 # ⛔ 内容**不得含时间戳**:唤醒回路按**内容哈希**去重,带时间戳会退化成每轮都投一遍。 NEXT_MD.write_text( "# 🛑 队列受阻(**有活,但没人能开工** · 需人处理)\n\n" + "\n".join("- **%s** —— %s" % (k, v) for k, v in blocked.items()) + "\n\n## 主会话该做的\n" "1. 读 `tmp/supervise-inbox/NEED-USER.md`(受阻原因与解除条件都在里面)\n" "2. **明确报告用户**:要哪一句授权、要哪个决定\n" "3. 用户放行后办掉,并**摘掉 `blocked.json` 里对应项** ⇒ 队首自动恢复\n" "4. ⛔ 不要因为受阻就硬派别的棒(那会变成「空转链条」)\n" "\n> ⛔ 一次只处理这一件;**受阻 ≠ 停滞** —— ⛔ 不抢锁、不接管。\n", encoding="utf-8") else: try: NEXT_MD.unlink() except Exception: pass except Exception: pass try: QUEUE.write_text(json.dumps(out, ensure_ascii=False, indent=1), encoding="utf-8") QUEUE_MD.write_text("\n".join([ "# 严格队列(一次一件 · 原子取件 · 同线互斥)", "", "- 生成于:%s" % time.strftime("%Y-%m-%d %H:%M:%S"), "- **队首(下一个该做)**:%s" % ("**%s** %s(线:%s,%s)" % (head["id"], head["title"], head["line"], "在关键路径上游·做它可解锁关键路径" if head["prio"] == 0 else "非关键路径") if head else "(无可派:都在做或等前置)"), "- **正在做**:%s" % ("、".join("%s→%s" % (k, v) for k, v in doing.items()) or "(无)"), "- **受阻(⛔ 不当队首)**:%s" % ("、".join("**%s**(%s)" % (k, v) for k, v in blocked.items()) or "(无)"), "- **待办**:%d 件%s" % (len(pend), " ⚠️ 同时在做多件(可能打架)" if out["conflict"] else ""), "", "## 取件规矩(**严格队列**)", "1. **只取队首那一件**;取件用**原子动作**:`mkdir tmp/supervise-inbox/claims/` ⇒ **建不成就是别人取走了** ⇒ ⛔ 换队首/等待", "2. 取到后写 `claims//holder`(内容 `<会话名>@<线>@<完整会话id>`)⇒ 「谁在做」的唯一权威;**第 3 段别省**(持有人失活 ⇒ 程序立即出队;⛔ 不写只能等超时兜底)", "3. **做完 → 删掉 `claims/`**(出队)⇒ 下一个才可能成为队首(⛔ 不删 ⇒ 该线一直被占)", "4. **同线互斥、跨线并行**:同一条线同时只允许一件;不同线可并行(队列按 `line` 判)", "5. **卡死**:claim 超 %d 分钟未续 ⇒ 程序自动移到 `claims-stale/`(可追溯)⇒ 队列自动放行" % (DOING_TTL // 60), "", ]), encoding="utf-8") except Exception: pass return out # ── ⑩ 唤醒回路(程序把「有活」投给**活会话**) ─────────────────── # 🔴 2026-09-29 立(用户拍板)。事实依据(源码实读 + 本机实测,⛔ 别再试错): # · `POST /api/v1/sessions/{id}/reply` = **官方投递语义:不夺 ACP writer**(桌面不被切换/打断); # ⚠️ 只对**该网关的当前活会话**成立,其余会话会 409 ⇒ 故先 `GET /api/v1/sessions/live` 取 id。 # · 鉴权头**只有** `x-access-token: <口令>` 通(`?password=` 在受保护路径上恒失效); # 口令从 `os.environ[CODEBUDDY_GATEWAY_PASSWORD]` 取(桥跑在 WorkBuddy 进程树内即自动可得)。 # · 本机**并存多个实例**(各带一个网关)⇒ 认口必须加 **cwd == 本工作区** 这一条, # 否则会把「有活」投进**别的工作区**的会话里。同 cwd 多口时取 **uptime 最大**(=最稳的常驻实例)。 NEXT_MD = INBOX / "NEXT.md" WAKEUPS = INBOX / "wakeups.jsonl" NEED_USER = INBOX / "NEED-USER.md" GW_TITLE_KEYS = ("CodeBuddy Gateway", "CodeBuddy Remote Control") GW_HDR = "x-access-token" GW_ENV = "CODEBUDDY_GATEWAY_PASSWORD" # 🔴 2026-09-30 加:本进程**是不是由宿主钩子唤起**(`--tick`)—— # 用来区分两种"没口令",两者的处置完全相反: # · 钩子唤起(进程在**宿主进程树**内)却仍没口令 ⇒ **真异常** ⇒ 必须落 `NEED-USER.md` 让人看见; # · 常驻进程(从「启动」文件夹/命令行起,**天然在宿主进程树外**)没口令 ⇒ **设计使然**(⛔ 不是故障) # ⇒ 只记日志,⛔ 不要每小时写一条 `NEED-USER.md` 把真问题淹掉。 FROM_HOOK = False def _norm_p(p) -> str: return str(p or "").replace("\\", "/").rstrip("/").lower() def _gw_http(port: int, method: str, path: str, obj, token: str, timeout: float = 6.0): """直调本机网关。⛔ 口令只作请求头使用:不写日志、不落盘、不回显、不进异常消息。""" h = {"Accept": "application/json"} data = None if obj is not None: h["Content-Type"] = "application/json" data = json.dumps(obj, ensure_ascii=False).encode("utf-8") if token: h[GW_HDR] = token req = urllib.request.Request("http://127.0.0.1:%d%s" % (int(port), path), data=data, headers=h, method=method) try: with urllib.request.urlopen(req, timeout=timeout) as r: return r.status, r.read().decode("utf-8", "replace") except urllib.error.HTTPError as e: try: return e.code, e.read().decode("utf-8", "replace") except Exception: return e.code, "" except Exception as e: return 0, str(e) def _listen_ports(limit: int) -> list: """扫 `netstat -ano` 取 `127.0.0.1: LISTENING` 候选(1024 list: """认网关口:① `GET /` 200 且正文含 CodeBuddy Gateway/Remote Control ② `/api/v1/info.cwd` **等于本工作区**(多实例区分的关键)。返回按 uptime 降序 ⇒ `[0]` 即最稳的那个。""" tok = os.environ.get(GW_ENV) or "" hits = [] for p in _listen_ports(int(C.get("wake_max_ports") or 60)): st, body = _gw_http(p, "GET", "/", None, "", timeout=3.0) if st != 200 or not any(k in body for k in GW_TITLE_KEYS): continue cwd, up = "", 0.0 st2, t2 = _gw_http(p, "GET", "/api/v1/info", None, tok, timeout=5.0) try: dd = json.loads(t2).get("data") or {} cwd, up = dd.get("cwd") or "", float(dd.get("uptime") or 0.0) except Exception: pass if _norm_p(cwd) != _norm_p(WS): continue sid = "" st3, t3 = _gw_http(p, "GET", "/api/v1/sessions/live", None, tok, timeout=5.0) try: sid = (json.loads(t3).get("data") or {}).get("sessionId") or "" except Exception: pass hits.append({"port": p, "cwd": cwd, "uptime": up, "sessionId": sid, "live_http": st3}) hits.sort(key=lambda x: -x["uptime"]) return hits def _wakeup_log(rec: dict) -> None: try: with open(WAKEUPS, "a", encoding="utf-8") as f: f.write(json.dumps(rec, ensure_ascii=False) + "\n") except Exception as e: log("wake: append wakeups err %s" % e) def need_user(reason: str) -> None: """遇阻 ⇒ 落 `NEED-USER.md`(报告里必须**明确喊「需用户介入」**)。""" try: NEED_USER.parent.mkdir(parents=True, exist_ok=True) NEED_USER.write_text( "# 🔴 需用户介入\n\n- 时间:%s\n- 原因:%s\n\n## 喊话\n**需用户介入**\n" % (time.strftime("%Y-%m-%d %H:%M:%S"), reason), encoding="utf-8") except Exception: pass def wake_round(st: dict, T: dict, V: dict | None = None) -> dict: """⛔ **已停用(2026-09-29)**:投递已**唯一化到监督程序**(`supervise()`),本函数不再被任何路径调用; 它依赖的 `NEXT.md` 旧队列一并退役。⛔ 本轮不删(删除要按"函数边界+夹层常量"的规矩单独做, 避免重演误删常量那次),列入 ④ 收敛清单。""" """闸门全满足才投**一次**(一份内容只投一次)。返回本轮判读(写进 state,供视图/排查用)。 触发面(🔴 2026-09-29 扩):① `NEXT.md` 存在且 `gate=free`(有活派) ② **`NEXT.md` 不存在但出现「停滞」或「可派未派」** —— ⛔ 原实现只认 ① ⇒「没人发消息 ⇒ 没人跑程序 ⇒ 永远发现不了停滞」的死循环。""" info = {"skipped": ""} if not C.get("wake_enable", True): info["skipped"] = "disabled" return info nodes = T.get("nodes_raw") or [] if nodes and all((n.get("status") == "done") for n in nodes): info["skipped"] = "all-done" # 终止:任务图全 done ⇒ 不再唤醒 return info synth = "" if NEXT_MD.exists(): try: txt = NEXT_MD.read_text(encoding="utf-8", errors="replace") except Exception: info["skipped"] = "next-unreadable" return info try: gate = (json.loads(QUEUE.read_text(encoding="utf-8")) or {}).get("gate") except Exception: gate = None if gate != "free": info["skipped"] = "gate=%s" % gate return info else: # 🔴 2026-09-29 补(治「没人发消息 ⇒ 彻底静默」):NEXT.md 只是「有活派」的载体, # **不是唯一该把人叫起来的理由**。**停滞**(无人接活 / 有活可派未派)同样必须投一次; # 否则「有活 ∧ 谁都动不了 ∧ 程序一声不响」—— 实测 16:31→20:07 整 4.5 小时零投递, # 用户看到的就是「一直在发呆」。 # ⚠️ 哈希必须用**粗粒度键**(种类 + 空闲按 30 分钟取整):若拿带时间戳的摘要正文去算, # 会退化成每 180 秒投一次(刷屏)。 idle_v = float((V or {}).get("idle") or 0) alert = str((V or {}).get("alert") or "") if alert.startswith("确定性静默"): # ⛔ 这一档是「**要用户拍板**」的状态 ⇒ **只投一次**(用常量键、不带时间桶): # 它不会自己好,每 30 分钟重复喊一遍只会变成噪音;状态一变(排期恢复 / 有人跑起来 / # 新结论落库)键自然变 ⇒ 会再投一次。 synth = "确定性静默" elif alert: synth = "僵住|%d|%s" % (int(idle_v // 30), alert) elif T.get("waste"): synth = "可派未派|%d" % int(idle_v // 30) if not synth: info["skipped"] = "no-next" return info txt = synth info["trigger"] = synth h = hashlib.sha1(txt.encode("utf-8", "replace")).hexdigest()[:16] wtext = str(C.get("wake_text") or "") if synth: wtext = ("协作程序报警:%s。请读 tmp/supervise-inbox/digest.md 与 NEED-USER.md,判「该谁动」" "——**有活就派、被卡就明确报给用户**;本轮只做这一件,做完即停。" % ((V or {}).get("alert") or "有活可派但无人接")) W = st.get("wake") or {} if W.get("hash") == h: info["skipped"] = "same-item" return info if time.time() - float(W.get("ts") or 0.0) < float(C.get("wake_min_gap") or 180): info["skipped"] = "too-soon" return info tok = os.environ.get(GW_ENV) or "" if not tok: info["skipped"] = "no-token" # ⛔ 口令缺失只跳过,绝不落盘/回显 return info gws = discover_gateways() if not gws: info["skipped"] = "no-gateway" return info # 🔴 2026-09-29 补(实测驱动):**必须挑「有活会话」的那个口**,⛔ 不能只取 `gws[0]`。 # 实测(本机 13:24):同 cwd 两个口 —— `:62213` uptime 更大但 `/sessions/live` 为空, # 有活会话的 `:56944` 排在后面 ⇒ 原逻辑取 `gws[0]` ⇒ 直接落 `no-live-session` 跳过, # **「有活会话也投不出去」**。改法最小:沿用 uptime 降序,取**第一个带 sessionId 的**。 # ⛔ 同上:**不做盲选**(2026-09-30)。本路径虽已停用,但**留着盲选=留着一颗雷** —— # 将来谁把它复活,就会重新开始"随手投给任意一条会话"。 _rm2 = resolve_main(st, [x.get("sessionId") for x in gws] if gws else []) _w2 = str(_rm2.get("sid") or "") g = next((x for x in gws if x["sessionId"] and _w2 and str(x["sessionId"]) == _w2), None) if gws else None if g is None: info["skipped"] = ("no-main-session" if _w2 == "" else "main-not-live") return info if _session_status(str(g.get("sessionId") or "")) == "working": # 同上:忙时不投 log("延后投递(旧路径):目标会话正在执行") info["skipped"] = "target-busy" return info stx, body = _gw_http(g["port"], "POST", "/api/v1/sessions/%s/reply" % g["sessionId"], {"text": wtext}, tok, timeout=20.0) ok = stx in (200, 201, 202) rec = {"ts": time.strftime("%Y-%m-%dT%H:%M:%S"), "epoch": round(time.time(), 1), "port": g["port"], "sessionId": g["sessionId"], "hash": h, "http": stx, "ok": ok} _wakeup_log(rec) st["wake"] = {"ts": time.time(), "hash": h, "port": g["port"], "sessionId": g["sessionId"], "http": stx} log("wake reply -> %s@%s http=%s hash=%s" % (str(g["sessionId"])[:8], g["port"], stx, h)) if not ok: need_user("唤醒投递失败:网关口 %s 回 %s(%s)" % (g["port"], stx, str(body)[:140])) info["skipped"] = "deliver-failed" else: info["delivered"] = rec return info # ── ⑪ 任务队列(🔴 架构定案 2026-09-29 用户):**状态由协作会话上报,⛔ 不靠猜** ────────── # ① 协作会话**开始执行** ⇒ 告诉协作程序 ⇒ 条目置 `running`(记录执行者+开始时刻) # ② 协作会话**处理完毕** ⇒ 告诉协作程序 ⇒ 条目置 `done`(记录产物+结束时刻) # ③ **监督程序**逐条读队列 ⇒ 有新变化 ⇒ 告诉主会话跟进; # **一段时间没有「执行中/执行完毕」的队列**(=静默)⇒ **发心跳**让主会话检查状态。 # ⛔ 本程序只"读队列 + 写通知",⛔ 不派活、⛔ 不开会话。 TASKS = INBOX / "tasks.json" # 🔴 **唯一权威**(队列本体) TASK_EVENTS = INBOX / "tasks-events.jsonl" # 只作审计(append-only,⛔ 不参与判定) TO_MAIN = INBOX / "TO-MAIN.md" # 监督程序给主会话的通知(投影 · 供人/AI 直接读) TASK_STATES = ("pending", "running", "done", "blocked") # 需求台账四态:待执行/执行中/已完成/有阻碍 QUEUE_IDLE_MIN = 5 # 心跳**限流桶**(分钟):同一状态下最多每 **5** 分钟发一次(2026-09-30 用户拍板:30→10→**5**);⛔ 触发条件是 supervise() 的**四条件**(不是计时器) HEARTBEAT_MIN_IDLE = 15 # 🔴 2026-09-30 加:**「真停滞」门槛**(分钟)—— # 距**上次任何进展** ≥ 15 分钟才允许发心跳;**有进展 ⇒ 一条都不发**。 # 为什么:心跳原只要求"三条件成立",而"目标未完成"在长任务期长期成立 # ⇒ 变成定期噪音,且**每次投递都会唤起主会话跑一轮 ⇒ 用户发消息撞 busy ⇒ 体感"卡"**。 def _load_tasks() -> dict: try: t = json.loads(TASKS.read_text(encoding="utf-8")) return t if isinstance(t, dict) else {} except Exception: return {} def _save_tasks(t: dict) -> None: try: TASKS.write_text(json.dumps(t, ensure_ascii=False, indent=1), encoding="utf-8") except Exception as e: log("tasks 写失败 %s" % e) def task_report(tid: str, state: str, by: str = "", artifact: str = "", reason: str = "", line: str = "") -> int: """协作会话**上报**状态转移(需求台账四态)。用法: `python collabd.py --report <需求id> --state pending|running|done|blocked [--by 会话名] [--artifact 产物] [--reason 阻碍原因]` ⚠️ `--state blocked` **必须**带 `--reason`(⛔ 不许只标"卡了"不说卡在哪、谁在等)。 """ tid = (tid or "").strip() state = (state or "").strip().lower() if not tid or state not in TASK_STATES: print("用法:--report <需求id> --state pending|running|done|blocked [--by 会话名] [--artifact 产物] [--reason 原因]") return 2 t = _load_tasks() rec = dict(t.get(tid) or {}) prev = str(rec.get("state") or "") rec.update({"state": state, "by": by or rec.get("by", ""), "t": time.time()}) if line: rec["line"] = line # 线归属**写进条目**(⛔ 不再从会话 cwd 反推) if artifact: rec["artifact"] = artifact if state == "blocked": rec["block_reason"] = reason or rec.get("block_reason", "") else: rec.pop("block_reason", None) # 解除阻碍 ⇒ 原因一并清掉(⛔ 不留过期原因误导后续判断) if state == "running" and not rec.get("t_start"): rec["t_start"] = time.time() if state == "done": rec["t_end"] = time.time() t[tid] = rec # 🔴🔴 2026-09-30 改:**乐观重试**(治"多个主会话同时调用本 skill ⇒ 台账丢更新") # ⛔ **为什么不用文件锁**:实测在**隔离目录**里 OS 锁三项全绿(同进程反复 lock/unlock 0.01s; # 两进程真互斥 0.00s;释放后立刻拿到 0.00s),但**接进生产后** `--once` 与 6 个并发 `--report` # **全部 rc=124 卡住不退出**(⚠️ 数据其实写成功了)⇒ **本机这个方案不能用** ⇒ 换**零锁**方案。 # ✅ 乐观重试:**写 → 回读核对 → 被覆盖 ⇒ 以磁盘为准重新合并自己这条**(最多 3 次)。 # 适用性:本场景是"**低冲突、写小文件**" ⇒ 足够,且**没有任何卡死风险**。 # ⚠️ 实测登记:**8 并发 `--report` 在"无任何保护"的旧实现下也是 8/8 全成功**(没压出丢更新) # ⇒ 说明这是**低概率事件**;本补丁是"真冲突时能自愈"的保险,⛔ 不是复现过的故障修复。 for _try in range(3): _save_tasks(t) try: _back = _load_tasks() except Exception: _back = {} if str((_back.get(tid) or {}).get("state") or "") == state: break # ✅ 自己这条已写进去 if isinstance(_back, dict) and _back: t = dict(_back) # ⚠️ 被覆盖过 ⇒ 以磁盘为准,重新合并自己这条 t[tid] = rec time.sleep(0.05) try: with open(TASK_EVENTS, "a", encoding="utf-8") as f: f.write(json.dumps({"ts": time.strftime("%Y-%m-%dT%H:%M:%S"), "id": tid, "from": prev, "to": state, "by": by, "artifact": artifact}, ensure_ascii=False) + "\n") except Exception: pass log("task report %s: %s -> %s (%s)" % (tid, prev or "-", state, by or "?")) print("OK 已上报:%s %s -> %s" % (tid, prev or "(新)", state)) return 0 def _deliver_str(text: str, key: str, st: dict, topic: str = "") -> dict: """把一段文字投给**活会话**(网关官方 reply,不夺 ACP writer)。⛔ 口令只进程内用。 🔴 `topic`(2026-09-30 加)=**这条内容属于哪个任务类别**(=台账条目的 `line`)—— 用户要求「**同一个工作区** 多会话协作(通过协作会话名称前缀区分具体任务会话)」 ⇒ 投递**按类别选主会话**(见 `main_for_topic()`)。 · `topic` 省略 ⇒ 取 `default`(⛔ **默认类别**,与旧版行为一致 —— 心跳这类"全局"内容就该走它) · `topic` 是已登记类别但**没解析出该类别的主会话** ⇒ **降级报告**(⛔ 不投给别类别) """ info = {"skipped": ""} if not C.get("wake_enable", True): info["skipped"] = "disabled" return info # 🔴 2026-09-29 用户定案:**投递前先判目标会话是否在跑 —— 只有"没在跑"才投**。 # (正在 `working` 的会话,投进去会插进它当前的轮次里 ⇒ 等它空闲再投,与 §2.1 握手同源) # 🔴 2026-09-30 改:目标从"唯一主会话"改成"**该任务类别的主会话**"。 # ⚠️ 此处**不做网关发现**(那一步在后面)⇒ 用**不含 live 信息**的解析做粗判; # 精确判定在拿锁后**用活网关列表重算一遍**(下面 `_rms = resolve_mains(...)`)。 _tgt = str(main_for_topic(resolve_mains(st, []), topic).get("sid") or "") if _tgt and _session_status(_tgt) == "working": log("延后投递:主会话 %s 正在执行 ⇒ 等它空闲" % _tgt[:8]) info["skipped"] = "target-busy" return info # 🔴 2026-10-01 加:目标会话**已哑**(诊断日志撞 ~10 MiB ⇒ EPERM ⇒ 界面不再刷新) # ⇒ 投了也看不见,记 ok:true 是假绿 ⇒ ⛔ 不投,改喊用户换会话。 if _tgt and _tgt in _deaf_sids(): need_user("🔴 主会话 `%s` 的诊断日志已撞 ~10 MiB 上限(宿主拒写 EPERM)⇒ 它的界面再也刷不出内容," "投给它的通知全部落空。**请把主会话换到一条新会话**(旧的那条只能弃用)。" % _tgt[:8]) log("停投:主会话 %s 已哑(日志撞上限)⇒ 改喊用户" % _tgt[:8]) info["skipped"] = "target-deaf" return info # 🔴🔴 2026-09-30 加:**快速否决(⛔ 不取锁、⛔ 不删锁)** —— 治"常驻被护栏杀掉"(实测真因): # `_deliver_str` 每轮尝试都会在 finally 里 `unlink(wake.lock)`;宿主有 **SafeDelete 批量删除护栏** # (按"本轮删除次数"计数,达阈值即要求确认并**拒绝**)⇒ 常驻监督程序跑约 50 次后**被系统终止**。 # 实测:后台任务 `Syz5DD` 跑 **48m43s** 后 `failed`,stdout 原文 # `[safe-delete][SAFE_DELETE_BULK_CONFIRM_REQUIRED] {"count":50,"threshold":50,…targets:[…\wake.lock"]}` # ⇒ **心跳的时钟就是这么断的**(用户当晚的关切正是"别回来还在发呆")。 # 修法:把 hash / 最小间隔的**预判提到取锁之前** —— 高频的 `same-item` / `too-soon` 路径 # (实测占绝大多数:日志里成片 "投递未成(too-soon)")**根本不碰锁文件** ⇒ 删除次数降到"只在可能真投时"。 # ⚠️ 取舍:锁内仍会**重读 STATE 再判一次**(那是权威判据),所以竞态窗口由 `wake_min_gap` 兜底 ⇒ ⛔ 不会双投。 _h0 = hashlib.sha1((key + "|" + text).encode("utf-8", "replace")).hexdigest()[:16] _W0 = st.get("wake") or {} if _W0.get("hash") == _h0: info["skipped"] = "same-item" return info if time.time() - float(_W0.get("ts") or 0.0) < float(C.get("wake_min_gap") or 300): info["skipped"] = "too-soon" return info # 🔴 2026-09-29 23:12 修("卡 working / 发消息没反应"复盘): # **两个常驻进程(协作程序 + 监督程序)各自投递、共用状态文件但无互斥** ⇒ 都读到旧哈希 # ⇒ **同一内容成对重复投递**(实测同一 hash 秒级出现两次:22:24:04/30、22:40:23/26、 # 22:48:32/55、22:51:01×2)⇒ 每次都往主会话的会话队列里**压消息**(22:53 起该会话进入 # parkInQueue「只进不出」⇒ 用户随后发的消息排在后面 ⇒ 表现为"卡死、没反应")。 # ⇒ 三招:**原子取锁(跨进程互斥)+ 拿锁后重读状态 + 投完立刻落盘**。 _lk = INBOX / "wake.lock" _got = False try: _fd = os.open(str(_lk), os.O_CREAT | os.O_EXCL | os.O_WRONLY) os.write(_fd, str(os.getpid()).encode()) os.close(_fd) _got = True except FileExistsError: try: if time.time() - _lk.stat().st_mtime > 120: # 残锁 >120 秒 ⇒ 抢占 _lk.unlink() _fd = os.open(str(_lk), os.O_CREAT | os.O_EXCL | os.O_WRONLY) os.write(_fd, str(os.getpid()).encode()) os.close(_fd) _got = True except Exception: _got = False except Exception: _got = True # 取锁本身出错 ⇒ ⛔ 不因它挡死投递 if not _got: info["skipped"] = "locked" return info try: try: # 拿锁后**重读**状态(⛔ 不用陈旧内存哈希) _st2 = json.loads(STATE.read_text(encoding="utf-8")) or {} if isinstance(_st2, dict) and _st2.get("wake"): st["wake"] = _st2["wake"] except Exception: pass h = hashlib.sha1((key + "|" + text).encode("utf-8", "replace")).hexdigest()[:16] W = st.get("wake") or {} if W.get("hash") == h: info["skipped"] = "same-item" return info if time.time() - float(W.get("ts") or 0.0) < float(C.get("wake_min_gap") or 300): info["skipped"] = "too-soon" return info finally: try: _lk.unlink() except Exception: pass tok = os.environ.get(GW_ENV) or "" if not tok: # 🔴 投递不到 ≠ 可以沉默 —— 但**两种"没口令"要分开处置**(见 `FROM_HOOK` 处说明): if FROM_HOOK: need_user("有内容要投给主会话,但**连宿主钩子进程里都拿不到网关口令**(%s)" "—— 属异常(宿主版本/注册面变了?),⛔ 请勿当成常驻进程的正常现象" % key) else: log("投递跳过:本进程不在宿主进程树内(无 %s)—— **设计使然**(投递由宿主钩子唤起)" % GW_ENV) info["skipped"] = "no-token" return info gws = discover_gateways() # 🔴 2026-09-29 修(用户截图暴露):**桌面上会有很多会话窗口** ⇒ 不能只取"第一个带会话 id 的口", # 否则通知会**投错窗口**。⇒ 目标必须是**声明为主会话**的那个(`--declare --role main`); # 没声明才回落到旧行为。 # 🔴 2026-09-30 重写:**主会话改成"解析出来的"**(见 `resolve_main` 的说明)。 # ⛔ 同时**删掉了原来的盲选回落** `next((x for x in gws if x["sessionId"]), None)` —— # 目标不在活会话里时它会**随手投给任意一条会话**("投错窗口",上面的注释自己就警告过)。 _rms = resolve_mains(st, [x.get("sessionId") for x in gws] if gws else []) rm = main_for_topic(_rms, topic) # 🔴 按**任务类别**取主会话(⛔ 不是"唯一那一条") if rm.get("switched_from"): _follow_main(st, rm) # 换了 ⇒ 跟随 + 显式告警(⛔ 不静默) _want = str(rm.get("sid") or "") g = next((x for x in gws if x["sessionId"] and _want and str(x["sessionId"]) == _want), None) if gws else None if g is None: # 🔴 **⛔ 拒绝盲投**:解析不出就明确报,绝不投给任意一条会话 if _want == "": if topic and topic in (_rms.get("topics") or []): need_user("**找不到「%s」这条任务类别的主会话**:本工作区里既没有标题带 `%s` 的会话," "也没有显式标 `主控` 的那条 ⇒ 本程序拒绝盲投" "(⛔ **不投给别的任务类别**的会话)。请确认该类别的主会话是哪条。" % (topic, topic)) else: need_user("**找不到主会话**:登记 %s 已不在活会话里,按**工作区**(%s)也解析不出来 ⇒ " "本程序拒绝盲投(⛔ 不投给任意会话)。请确认哪条是主会话。" % (str(_main_sid(st))[:8] or "(空)", str(WS))) info["skipped"] = "no-main-session" else: need_user("主会话 %s 此刻不在活会话里(可能已关或已换)⇒ 本程序拒绝盲投。请确认。" % _want[:8]) info["skipped"] = "main-not-live" return info if g is not None and _session_status(str(g.get("sessionId") or "")) == "working": log("延后投递:目标会话 %s 正在执行 ⇒ 等它空闲" % str(g.get("sessionId"))[:8]) info["skipped"] = "target-busy" return info # 🔴 2026-10-01 加:目标**已哑**(诊断日志撞 ~10 MiB ⇒ EPERM ⇒ 界面不再刷新)⇒ 投了也看不见。 # 实测 f8a792ab:completed 却仍在网关口 live ⇒ 登记 main 命中 ⇒ 绕过 _pick_live ⇒ 必须在这里再拦一道。 if g is not None and str(g.get("sessionId") or "") in _deaf_sids(): need_user("🔴 主会话 `%s` 的诊断日志已撞 ~10 MiB 上限(宿主拒写 EPERM)⇒ 它的界面再也刷不出内容," "投给它的通知全部落空(记 ok:true 是假绿)。**请把主会话换到一条新会话**,旧的那条只能弃用。" % str(g.get("sessionId"))[:8]) log("停投:目标 %s 已哑(日志撞上限)⇒ 改喊用户" % str(g.get("sessionId"))[:8]) info["skipped"] = "target-deaf" return info if g is None: need_user("有内容要投给主会话,但**没有活会话**可投(%s)—— 请打开任意工作区会话" % key) info["skipped"] = "no-live-session" return info stx, _body = _gw_http(g["port"], "POST", "/api/v1/sessions/%s/reply" % g["sessionId"], {"text": text}, tok, timeout=20.0) ok = stx in (200, 201, 202) _wakeup_log({"ts": time.strftime("%Y-%m-%dT%H:%M:%S"), "epoch": round(time.time(), 1), "port": g["port"], "sessionId": g["sessionId"], "hash": h, "http": stx, "ok": ok, "kind": key}) # 🔴 2026-09-30 加:投递成功**不等于**被消费 ⇒ 记下"期望"(下一轮回查该会话有没有真的动)。 st["wake"] = {"ts": time.time(), "hash": h, "port": g["port"], "sessionId": g["sessionId"], "http": stx, "expect": {"sid": g["sessionId"], "at": time.time()}} try: STATE.write_text(json.dumps(st, ensure_ascii=False), encoding="utf-8") # 立刻落盘 ⇒ 另一进程读到新哈希 except Exception: pass log("notify(%s) -> %s@%s http=%s hash=%s" % (key, str(g["sessionId"])[:8], g["port"], stx, h)) info["delivered"] = {"http": stx, "kind": key} return info def _session_status(sid: str) -> str: """只读宿主库拿某会话的状态(`working` = 正在执行)。取不到 ⇒ 空串(⛔ 不猜)。""" if not sid: return "" try: con = _db() row = con.execute("select status from sessions where id=?", (sid,)).fetchone() con.close() return str(((row or [""])[0]) or "").lower() except Exception as e: log("session_status 读失败 %s" % e) return "" # ── 🔴 2026-09-30 加:**宿主侧「消息卡住」探针**(用户给出窗口 22:55–23:20,实测定型) ──── # 指纹:工作区日志里 `PromptIterator … route=parkInQueue … hasWaiter=false` # ⇒ **消息进了队列、但没有消费者** ⇒ 界面一直转、发消息没反应。 # 真因=**宿主客户端连接状态丢失**(同窗口 `sendToClient: No state found for connectionId` # 303 次,其他小时 0 次);实测队列 `queueLen` 0→1→2 逐条堆积。 # 🔴 **AI 侧修不了这个根因** —— 唯一解法是**让客户端重新挂上该会话**(切走再切回/重开窗口)。 # ⇒ 我们能做的只有一件:**一发生就发现,并把"该点哪一下"直接写进 `NEED-USER.md`**。 # 详见 `references/pitfalls.md P0-5`。 PARK_TAIL = 400 * 1024 # 只读日志尾部(⛔ 不许整文件扫:工作区日志可达几十 MB) PARK_FRESH = 30 * 60 # 秒:指纹超过这个时长就不算"正在卡" PARK_STAMP = INBOX / "_park.stamp" # 告警节流(⛔ 不每轮刷同一条) PARK_GAP = 180 # 秒 CONSUME_GRACE = 240 # 秒:投递后多久没见目标会话活动 ⇒ 判「没被消费」 def _log_dirs() -> list: """候选工作区日志路径:`<配置根>/logs/<今天|昨天>/<本工作区名>__*.log`。 ⚠️ 日志**跨日不换名**(实测:主会话 09-29 那份一直写到 23:59)⇒ 两个日期都试。 """ root = os.path.join(os.environ.get("CODEBUDDY_CONFIG_DIR") or r"E:\ProgramData\.workbuddy", "logs") name = os.path.basename(str(WS).replace("\\", "/").rstrip("/")) or "" out = [] for d in (time.strftime("%Y-%m-%d"), time.strftime("%Y-%m-%d", time.localtime(time.time() - 86400))): try: for fn in os.listdir(os.path.join(root, d)): if name and fn.startswith(name + "__") and fn.endswith(".log"): out.append(os.path.join(root, d, fn)) except Exception: continue return out def _is_park_line(ln: str) -> bool: """判据:**只认真正的 AcpView 记录行**。 🔴 **⛔ 不能只搜字符串**(2026-09-30 当场踩到):取证命令(`grep parkInQueue …`)的输出 会**被写进同一份工作区日志**(`SandboxShell ProcessOutput … content=`)⇒ 探针命中自己的回声 ⇒ **假阳性**。 """ return ("parkInQueue" in ln and "hasWaiter=false" in ln and "[AcpView][PromptIterator] received prompt" in ln and "ProcessOutput" not in ln and "content=" not in ln) def probe_host_park() -> dict: """扫日志尾部找「parkInQueue + hasWaiter=false」指纹。⛔ 只读、⛔ 不写任何东西。""" hit = {"n": 0, "last": "", "file": ""} for p in _log_dirs(): try: sz = os.path.getsize(p) with open(p, "rb") as f: if sz > PARK_TAIL: f.seek(sz - PARK_TAIL) buf = f.read().decode("utf-8", "replace") except Exception: continue for ln in buf.splitlines(): if _is_park_line(ln): # ⚠️ 判据见 `_is_park_line`(⛔ 不搜字符串,防"回声假阳性") hit["n"] += 1 hit["file"] = os.path.basename(p) m = re.search(r"\[(\d{4}/\d{1,2}/\d{1,2} (\d\d:\d\d:\d\d))", ln) if m: hit["last"] = m.group(2) return hit def _sess_activity(sid: str) -> float: """某会话最后活动时刻(epoch 秒)。⛔ 取不到 ⇒ 0(不猜)。""" if not sid: return 0.0 try: con = _db() row = con.execute("select updated_at from sessions where id=?", (sid,)).fetchone() con.close() return float(((row or [0])[0]) or 0) / 1000.0 except Exception: return 0.0 def check_delivery_consumed(st: dict) -> dict: """🔴 **「投出去」≠「它跑起来了」** —— 投递后必须**回查目标会话有没有真的动**。 判据:目标会话 `updated_at` 是否**晚于**投递时刻。⛔ 不靠"网关回了 200"(那只证明对方收下了)。 超 `CONSUME_GRACE` 秒仍未动 ⇒ 落 `NEED-USER.md`(含**用户唯一该做的那一下**),⛔ 不静默。 """ w = st.get("wake") or {} exp = w.get("expect") or None if not exp: return {} sid, at = str(exp.get("sid") or ""), float(exp.get("at") or 0) if not sid or not at: return {} act = _sess_activity(sid) if act and act > at: w.pop("expect", None) st["wake"] = w return {"consumed": True, "sid": sid[:8], "lag": round(act - at, 1)} if time.time() - at < CONSUME_GRACE: return {"waiting": True, "sid": sid[:8]} w.pop("expect", None) st["wake"] = w need_user( "投给会话 `%s` 的通知**没被消费**(投出后 %.1f 分钟该会话零活动)。\n" "典型的**宿主侧卡住**:`parkInQueue` + `hasWaiter=false` —— 消息进了队列、**没有消费者**。\n" "🔴 **这一步 AI 侧无法自救**。请你**让客户端重新挂上该会话**:把窗口**切走再切回**," "或**关掉再打开该会话窗口** —— 队列随即会被排空。⛔ 别反复发消息试探(只会在队尾再堆一条)。" % (sid[:8], (time.time() - at) / 60.0)) return {"unconsumed": True, "sid": sid[:8]} # ── 🔴 2026-09-30 加:**前置探针(设备侧 worker + 垫片)** ──────────────────────── # 为什么必须有:这两条腿**会周期性掉**,且**掉的时候没有任何声音** —— # · 实测 ①:垫片启动器跑 **1h09m** 后死,死因是 `reg.exe` 被内置程序黑名单拦(启动器某步碰注册表); # · 实测 ②:常驻监督程序跑 **48m43s** 后死,死因是 SafeDelete 批量删除护栏(见 P0-6)。 # ⚠️ 二者一掉 ⇒ ④闸 `503 device-unreachable` ⇒ 整条手机接入链断,而**用户只会看到"发呆"**。 # ⛔ **本函数只报警、⛔ 不自动重起** —— 本机"起一个能长期活的进程"只有**宿主后台任务机制**能做到, # 脚本自身起的任何进程都活不过工具调用边界(实测)⇒ 自动重起在本机做不到。 # ⇒ 报警必须写明"照 `references/deploy.md §5b` 重起"。节流 `FRONT_GAP` 秒,⛔ 不刷屏。 FRONT_STAMP = INBOX / "_front.stamp" FRONT_GAP = 180 FRONT_SHIM_PORT = 20090 FRONT_WORKER_LOG_GLOB = r"E:/dsh-worker-dev/logs/overlay-bg-*.out.log" def check_frontline() -> dict: """探测前置两条腿。返回 `{"shim": bool, "worker": bool, "note": str}`。⛔ 只读、⛔ 不启停任何进程。""" out = {"shim": False, "worker": False, "note": ""} # ① 垫片:127.0.0.1:20090 有没有 LISTENING(⛔ 不用 curl —— 本机有代理会伪造 502) try: r = subprocess.run(["netstat", "-ano"], stdout=subprocess.PIPE, stderr=subprocess.DEVNULL, creationflags=0x08000000, timeout=40, errors="replace") for ln in (r.stdout or "").splitlines(): if "LISTENING" in ln and ("127.0.0.1:%d" % FRONT_SHIM_PORT) in ln: out["shim"] = True break except Exception as e: out["note"] = "netstat 失败:%s" % e # ② worker:最新一份 out 日志的最后一行是否 `state=up` try: import glob as _glob fs = sorted(_glob.glob(FRONT_WORKER_LOG_GLOB), key=lambda p: os.path.getmtime(p)) if fs: with open(fs[-1], "rb") as f: f.seek(max(0, os.path.getsize(fs[-1]) - 4096)) tail = f.read().decode("utf-8", "replace").strip().splitlines() if tail and "state=up" in tail[-1]: out["worker"] = True except Exception as e: out["note"] = (out["note"] + " worker 读日志失败:%s" % e).strip() return out def alert_frontline(fr: dict) -> None: """前置掉了 ⇒ 落 `NEED-USER.md`(节流)。⛔ 两条腿都要报,并写明"哪一条 + 怎么起"。""" if fr.get("shim") and fr.get("worker"): return try: if FRONT_STAMP.exists() and (time.time() - FRONT_STAMP.stat().st_mtime) < FRONT_GAP: return except Exception: pass miss = [] if not fr.get("shim"): miss.append("**垫片**(`127.0.0.1:20090` 无 LISTENING)—— 掉了 ⇒ 入口④闸会 503") if not fr.get("worker"): miss.append("**设备侧 worker**(relay 日志无 `state=up`)—— 掉了 ⇒ relay 无本机通道") need_user( "🔴 **手机接入链路的前置掉了**(这一条**必须人工重起**,AI 侧做不到):\n- " + "\n- ".join(miss) + "\n\n⚠️ 本机要起「能长期活的进程」,只有**宿主后台任务机制**一条路;脚本自己起的进程活不过工具调用边界。" "\n✅ **怎么起**:见 `E:/ProgramData/.workbuddy/skills/multi-session-collab/references/deploy.md` **§5b**" "(两条命令 + 成功判据 + 三条死路)。\n" "⚠️ 已知复发周期:垫片约 **1 小时**(启动器某步碰 `reg.exe` ⇒ 被内置黑名单拦死);worker 尚未见固定周期。") try: FRONT_STAMP.write_text(time.strftime("%Y-%m-%d %H:%M:%S"), encoding="utf-8") except Exception: pass def _acc_short() -> list: """读 `goal.json.acceptance_state` ⇒ 返回**非 pass** 的验收项(供心跳正文直接列出"还差什么")。 ⚠️ 读不到 ⇒ 返回提示行(⛔ 不因它判"全过");`unknown` 视为**未过**。 🔴 2026-09-30 修(**"没判据"被说成了"全部 pass"** —— 与 `goal_state()` 同一个 bug 的另一半): 旧实现在**没有任何有效项**时走 `or [...]` 兜底 ⇒ 打印「(全部 pass —— 目标已达成,可收口)」, 而这恰恰是**最不该说"完成"**的场合(连判据都没有)。⇒ 现在**分开说**。 """ try: _as = (json.loads((INBOX / "goal.json").read_text(encoding="utf-8")).get("acceptance_state") or {}) _eff = [_k for _k in _as if not str(_k).startswith("_")] rows = ["- **%s**:%s" % (k, _as[k]) for k in sorted(_eff) if str(_as[k] or "").strip().lower() != "pass"] if rows: return rows if not _eff: return ["- ⚠️ **`acceptance_state` 里一条有效判据都没有**(只有说明行)" " ⇒ **判不出来**,⛔ 不因此判完成;请补判据,或明确写「本目标不需要验收判据」。"] return ["- (全部 pass —— 目标已达成,可收口)"] except Exception: return ["- ⚠️ 读不到 `goal.json.acceptance_state`(⛔ 不因此判完成)"] def _goal_short() -> str: """本需求目标的 `short`(两级命名前缀里的第 2 级)。""" try: return str(json.loads((INBOX / "goal.json").read_text(encoding="utf-8")).get("short") or "") except Exception: return "" def _goal_topics() -> list: """🔴 **本需求下的"任务类别"清单**(`goal.json.topics`)。 2026-09-30 用户要求:「要能支持**同一个工作区** 多会话协作(主会话根据任务**自动梳理任务类别**: **通过协作会话名称前缀的方式区分具体任务会话**,所有主会话,协作会话,自动唤醒任务,都在**一个工作区**)」 ⇒ 归属锚点从"**跨工作区**"(cwd)改成"**一个工作区 + 任务类别前缀**": **任务类别 ≡ 会话标题二级前缀 `[主题]` ≡ 台账条目的 `line` ≡ 分工板的一行**(同一个词,四处同义)。 ⛔ `topics` 缺省 ⇒ 回落 `[goal.short]` ⇒ **单类别时的行为与旧版逐字一致**(向后兼容)。 ⛔ 读不到 goal.json ⇒ 返回 `[]`(调用方按"无类别"处理,⚠️ 此时**只剩主会话判据**能命中)。 """ try: g = json.loads((INBOX / "goal.json").read_text(encoding="utf-8")) or {} except Exception: return [] out: list = [] for _t in (g.get("topics") or []): _t = str(_t or "").strip() if _t and _t not in out: out.append(_t) if not out: _s = str(g.get("short") or "").strip() if _s: out.append(_s) return out def _as_list(v) -> list: """把 `str | list | tuple | set` 归一成 `list[str]`(⛔ 空值 ⇒ `[]`)。 用途:`_in_project()` 的 `main_sid` / `topics` 两个参数都要能吃**单值或集合** —— 单类别(旧版)传 str、多类别传 list,调用点⛔ 不必两套写法。 """ if v is None: return [] if isinstance(v, str): return [v] if v.strip() else [] try: return [str(x) for x in v if str(x or "").strip()] except Exception: return [str(v)] def _in_project(sid: str, title: str, main_sid, short) -> bool: """🔴 **判定一个会话是否属于本需求目标项目**(2026-09-30 用户明示:「要明确哪些会话是属于某个需求目标项目的」)。 **三级判据(任一命中即算)—— ⛔ 与 cwd 无关** ① **主会话**:`sid == main_sid`(= 解析出来的主会话;🔴 **多任务类别时 main_sid 可以是多条**, 传 str 或 list/set 都行 —— 见 `resolve_mains()`) ② **协作会话**:**标题含 `[<任一任务类别>]`** —— 即两级命名前缀的**二级** (🔴 2026-09-30 泛化:第 4 个参数 `short` 实际收到的是**任务类别集合** `goal.topics`, 可传 str 或 list;⛔ 传 str 时行为与旧版**逐字一致**) ③ 其它会话:`--declare --goal ` 显式声明(⚠️ 尚未接线,留待需要时) ⛔ **不命中 ⇒ 不算本项目**。 🔴 **⛔ 刻意不用 cwd 推断** —— 架构 §2.3 明令「**绝不回落到 cwd 推断**」(我 2026-09-30 一度用 cwd 判,属违规,已改回)。 🔴 **为什么"多类别"必须在这一层就支持**:同一工作区里跑**多个任务类别**时, 若只认**一个** `goal.short`,则除那一个之外的类别会**全部被判"不算本项目"** ⇒ 看板不列它们、`--ready-next` 也不把它们的在跑会话算进来 ⇒ **静默漏管**。 """ if sid and sid in _as_list(main_sid): return True ti = str(title or "") for _t in _as_list(short): if ("[%s]" % _t) in ti: return True return False def ready_next(st: dict) -> dict: """🔴 **收尾确认 —— 用它替代"盲等 5~8 分钟"**(用户 2026-09-30 说明其由来): 用户原话:「5-8 应该是之前出现 **前面还没执行完 就开始下一棒了**, **要是能避免这个问题 可以压缩时间**」 ⇒ 5~8 分钟 = **用固定延迟赌"上一棒已收尾"**;既然能**真的确认**,就该压到"确认即起"。 **三条判据(全满足 ⇒ 可立刻排下一棒,排期 = 现在 + 30~60 秒)** ① 除**主会话**外,**没有别的会话处于 `working`** —— 读**宿主状态**,⛔ 不是"它自己说完了" ② 台账里**没有 `running` 的条目** —— ⛔ 不认"它说做完了",要终态(done/blocked) ③ 上一棒**没有留下待消费的投递**(`wake.expect` 已清)—— ⛔ 不等一个还没被取走的通知 ⚠️ 任一不满足 ⇒ **⛔ 不排**,等下一次心跳/反馈再判(⛔ 不盲等、也⛔ 不硬上)。 """ out = {"ok": False, "why": [], "working_others": [], "running_tasks": []} # ① 除主会话外还有谁在跑? # 🔴 2026-09-30 改:**多任务类别** ⇒ "主会话"可能不止一条 ⇒ # 判据从"排除**那一条**"改成"排除**解析出来的全部主会话**"(⛔ 否则别类别的主会话会被算成"在跑"⇒ 永远拦住) # ⚠️ 传 `live_sids=[]` ⇒ **不做网关探测**(本命令只需工作区扫描,⛔ 别为它加网络开销) _rms0 = resolve_mains(st, []) mains = list(_rms0.get("all_sids") or []) topics = _goal_topics() try: con = _db() for row in con.execute("select id, coalesce(nullif(custom_title,''),title,'') " "from sessions where status='working'"): sid, ti = str(row[0] or ""), str(row[1] or "") if sid in mains: continue # ⛔ 主会话自己不算(否则永远拦住) if not _in_project(sid, ti, mains, topics): # ⛔ **别的线的会话不算本项目在跑**(2026-09-30 修:曾把「AI变现日报」算进来 ⇒ 误拦派棒) out.setdefault("ignored_others", []).append("%s(%s)" % (sid[:8], ti[:24])) continue out["working_others"].append("%s(%s)" % (sid[:8], ti[:28])) con.close() except Exception as e: out["why"].append("读宿主库失败:%s" % e) if out["working_others"]: out["why"].append("还有会话在跑:%s" % "、".join(out["working_others"])) # ② 台账里还有 running? for k, v in (_load_tasks() or {}).items(): if str((v or {}).get("state") or "") == "running": out["running_tasks"].append(k) if out["running_tasks"]: out["why"].append("台账仍有执行中:%s" % "、".join(out["running_tasks"])) # ③ 上一棒的**单条反馈**是否还等着主会话处理?(握手没走完 ⇒ 本轮别急着重排) # 🔴 2026-09-30 修(**本命令上线 3 分钟就踩到**):原实现查的是 `wake.expect`, # 而**心跳投递同样会留下 `expect`** ⇒ `--ready-next` 会被"心跳尚未被消费"拦死 ⇒ # 实测输出「⛔ 先别排:上一棒投递尚未被消费」—— 而那条"未消费的投递"正是**刚发的心跳**, # 主会话正在处理它 ⇒ **假拦**(若不改,提速全部作废)。 # ⇒ 判据只认**握手等待**(`notify_awaiting`:反馈了一条、等主会话处理),⛔ 不看心跳的 expect。 _aw = st.get("notify_awaiting") or None if _aw: out["why"].append("上一棒的单条反馈仍在等主会话处理(%s,阶段 %s)" % (_aw.get("item"), _aw.get("phase"))) out["ok"] = not out["why"] return out def _guard_says_stop() -> bool: """🔴 **被守护拉起时**(`DSH_GUARDED=1`):守护停下 ⇒ 本程序**优雅退出**。 判据只有两个(⛔ 不猜):① 停止标志 `guard.stop` ② 守护心跳过期(>90 秒 ⇒ 守护已死)。 ⚠️ 手动跑(`--once`/钩子)没有该环境变量 ⇒ ⛔ 不受影响。 """ if os.environ.get("DSH_GUARDED") != "1": return False try: if (INBOX / "guard.stop").exists(): return True except Exception: pass try: ts = float(json.loads((INBOX / "guard.json").read_text(encoding="utf-8")).get("ts") or 0) return (time.time() - ts) > 90 except Exception: return False GOAL_F = INBOX / "goal.json" def load_goal() -> dict: """🔴 **任务目标**(整套机制的运行中心)。缺失/无 title ⇒ 返回空 dict —— 调用方**必须显式处理"未声明目标"**,⛔ 不得假装有目标。""" try: g = json.loads(GOAL_F.read_text(encoding="utf-8")) return g if isinstance(g, dict) and str(g.get("title") or "").strip() else {} except Exception: return {} def goal_line() -> str: """摘要/通知里显示的一行目标。 🔴 2026-09-30 改:显示**任务类别清单**(`goal.topics`)而不是单个 `short` —— 用户要求「主会话根据任务**自动梳理任务类别**,通过**协作会话名称前缀**区分具体任务会话」 ⇒ 这一行是**每轮钩子注入**给主会话的,必须把"本工作区里有哪几个类别"直接摆出来, 主会话才好据此给新棒命名(`[协作]-<类别>-<具体>`)并解析"这活属于哪一类"。 """ g = load_goal() if not g: return "⚠️ **未声明任务目标**(机制不知道该围绕什么跑)" tps = _goal_topics() if tps: return "🎯 围绕目标:%s(任务类别前缀:%s)" % ( g["title"], "、".join("`[%s]`" % _t for _t in tps)) return "🎯 围绕目标:%s" % g["title"] def pw_fingerprint() -> str: """网关口令的**指纹**(sha256 前 12 位)—— **只用于判"变没变"**。 ⛔ 绝不落口令本身;指纹不可逆、**不能用于鉴权**(所以不算"口令落盘")。""" tok = os.environ.get(GW_ENV) or "" return hashlib.sha256(tok.encode("utf-8")).hexdigest()[:12] if tok else "" def goal_state() -> str: """🔴 **目标三态**(2026-09-30 加):`"open"`(有未过项)|`"pass"`(有判据且全过)| `"undeclared"`(**一条有效判据都没有**)。 ⛔ **为什么要第三态**:旧实现把"**没有非 pass 项**"直接当成"全过" —— 可 `acceptance_state` 里**全是 `_` 开头的说明行**时(很常见:只写了 `_说明`/`_更新`), 遍历循环**空转** ⇒ **这一路什么也没说**,前三路又都 `done` ⇒ 整体判「已完成」。 ⇒ 这正是红线里那句「**一路读不到 ⇒ ⛔ 不许当成已完成**」被违反的形状: **不是崩溃,是少说一句话**(实测:`goalctl status` 因此显示"三路全过")。 🔴 `"undeclared"` **不算 pass**(调用方一律按"未完成/判不出"处理,⛔ 不臆断为过)。 """ open_, eff = False, 0 try: for _k, v in (_load_tasks() or {}).items(): if str((v or {}).get("state") or "") != "done": open_ = True except Exception as e: log("goal_state 读台账失败 %s" % e) try: g = json.loads(TG.read_text(encoding="utf-8")) for n in (g.get("nodes") or []): if str(n.get("status") or "") != "done": open_ = True except Exception as e: log("goal_state 读任务图失败 %s" % e) # 🔴🔴 2026-09-30 加(用户点破「目标都还未达成,为什么没触发目标状态评估」): # **前两路测的是"活干完了吗",⛔ 不是"目标达成了吗"** —— 铁证:任务图里 **N14 节点是 `done`, # 而 N14 判的正是「V1 ❌ 未过」**。⇒ 必须**再读一路"验收判据的真实状态"**。 # ⛔ 为什么不直接把前两路删掉:它们防的是「**漏待办**」(已规划但还没派棒 ⇒ 不在台账); # 本路防的是「**误判完成**」;二者**互补**,取并集。 # ⚠️ `unknown` 视为**未过**(⛔ 不臆断为过)⇒ 宁可多提醒,⛔ 不漏。 try: _gj = json.loads((INBOX / "goal.json").read_text(encoding="utf-8")) _as = _gj.get("acceptance_state") or {} eff = len([_k for _k in _as if not str(_k).startswith("_")]) if eff == 0: # ⚠️ 节流 600 s:这行**每轮都一样**(实测 40+ 行连刷),留一行就够,⛔ 别淹没真事件。 log_throttled("goal_state-no-kpi", 600.0, "goal_state:acceptance_state **没有有效判据**(只有说明行)⇒ 判「未声明」,⛔ 不算过" "(本行每 10 分钟最多一条)") for _k, _v in _as.items(): if str(_k).startswith("_"): continue if str(_v or "").strip().lower() != "pass": open_ = True except Exception as e: log("goal_state 读 acceptance_state 失败 %s" % e) return "undeclared" # 🔴 读失败 ⇒ ⛔ 不判完成(旧实现这里只记日志 ⇒ 静默放行) if open_: return "open" return "pass" if eff else "undeclared" def goals_open() -> bool: """**需求是否仍未完成** —— 🔴 `goal_state() != "pass"`(⛔ 包括"未声明判据"那种判不出来的)。 真值表:`open`(有未过项)⇒ True|`undeclared`(一条有效判据都没有)⇒ **True**(判不出来 ⇒ ⛔ 不判完成)| `pass`(有判据且全过)⇒ False。 """ return goal_state() != "pass" def goal_paused() -> bool: """🔴 **需求总闸**(2026-10-01 体检修):`goal.json.run != "active"` ⇒ **已停**。 为什么要这道闸(2026-10-01 实测):`goal.json.run=paused` 时,本程序**零处**引用 `run` ⇒ `--once` 照跑、照写 `NEXT.md`/`STALL.md`/`NEED-USER.md`/`digest.md`,而 `digest.md` 的内容**经钩子注入每个会话的每一轮** ⇒ 一个"人为停掉"的需求被读成"机制坏了"。 ⛔ 读不到 `goal.json` ⇒ 默认**在跑**(⛔ 不因读失败就把活机制判停)。 """ try: _g = json.loads((INBOX / "goal.json").read_text(encoding="utf-8")) except Exception: return False return str(_g.get("run") or "active").strip().lower() != "active" # 已停时**必须清掉**的告警投影(⛔ 留着就会被钩子继续注入 ⇒ "已停还在喊")。 # 依据:这些都是**投影**(可删可重建,见技能 §10),⛔ 不是人工输入(人工输入只有 `blocked.json`)。 PAUSED_SIGNALS = ("NEXT.md", "STALL.md", "NEED-USER.md", "queue.md") def paused_round(st: dict, why: str = "") -> dict: """**静默一轮**:只留一行说明,⛔ 不排队、⛔ 不产告警、⛔ 不投递。 两种入参,都属「**没有可推进的活**」这一大类: · `why` 为空 ⇒ **人为已停**(`goal.json.run != active`)—— ⛔ 不是故障。 · `why` 非空 ⇒ **目标已完成**(三路全过)—— 🔴 2026-10-01 加。 病根(实测 M5 解卡后立刻复现):原实现只认 `run`,目标三路全过而 `run` 仍 `active` ⇒ 照产告警 ⇒ 终态下每轮喊「卡住:有会话在跑但 N 分钟无成果」、`VACUUM`/`READY` 也复活 ⇒ **把「做完了」读成「坏了」**,还会经钩子注入每个会话。 """ _ts = time.strftime("%Y-%m-%d %H:%M:%S") if why: _line = ("✅ **终态**(%s @ %s)⇒ 目标三路全过、没有可推进的活,本轮不排队、不产告警、不投递。" "要开新活 ⇒ 重新声明目标(`goalctl.py declare …`)。" % (why, _ts)) else: _line = ("⏸ **已停**(`goal.json.run != active` @ %s)⇒ 本轮不排队、不产告警、不投递。" "恢复用 `goalctl.py start --yes`。⚠️ 这是人为停掉,⛔ 不是故障。" % _ts) try: (INBOX / "digest.md").write_text( "# 协作机制摘要\n\n- %s\n" % _line, encoding="utf-8", newline="\n") (INBOX / "advance.md").write_text( "# 机械推进快照(与 digest.md 同源)\n\n- %s\n" % _line, encoding="utf-8", newline="\n") # 🔴 2026-10-01 加:**看板也要跟着收口** —— 否则已停/终态之后看板永久停在最后一张 # "还在跑"的快照(实测:目标三路全过后看板仍显示「队首 M5 · 确定性静默」)⇒ 读的人被误导。 LIVE.parent.mkdir(parents=True, exist_ok=True) LIVE.write_text("# 协作实时状态\n\n%s\n" % _line, encoding="utf-8", newline="\n") except Exception: pass for _n in PAUSED_SIGNALS: try: (INBOX / _n).unlink() except Exception: pass try: (INBOX / "collabd-once.stamp").write_text(_ts, encoding="utf-8") except Exception: pass st["paused"] = True try: STATE.write_text(json.dumps(st, ensure_ascii=False), encoding="utf-8") except Exception: pass return st def is_continuation(name: str) -> bool: """🔴 标题是不是一条**接续会话**(=由**会话自己**建出来的"下一棒")。 **为什么必须认得出它**(2026-10-01 实测坐实的自指死结): 接续会话也是**自动化起的新会话**,标题由**建它的那条会话**写 ⇒ 常写成 `[唤醒机制] 接续 · 日志事前叫停钩子落地(第 2 棒)` / `接续棒:日志增长治理(任务 0 → 任务 A)` —— **没有角色方括号**。而旧判据只排 `[协作]` ⇒ 这类棒会被 `_scan_mains()` 收进 **主会话候选**,又因为它标题里有 `[<类别>]` ⇒ 直接被解析成"**该类别的主会话**" ⇒ 投递把通知**投给这条棒自己**(自己叫自己、白判一次)。 ⛔ 判据只用标题里的字面标记,⛔ **不回落 cwd、⛔ 不猜**(守架构 §2.3)。 """ return "接续" in str(name or "") def parse_session_name(name: str) -> dict: """解析会话命名 —— **四种形态**(返回 `{"role","topic","ok","form"}`): | form | 形态 | role | ok | |---|---|---|---| | `prefix` | `[角色]-[类别]-<具体>`(**合规**) | main/worker/waker | ✅ True | | `prefix` | 只有一级 `[协作]N9-…`(缺 `[类别]`) | worker | ❌ False | | `continuation` | 🆕 **接续会话**:`[<类别>] 接续 · …` / `接续棒:…` | **worker** | ❌ False | | `main-prefix` | 🆕 **主控前缀**:`主控 · <类别> · …`(中点分隔、无方括号) | **main** | ❌ False | 🔴 `ok=False` 恒表示"**没按两级前缀约定命名**"(⛔ 语义不变,老调用方照旧); 🆕 新加 `role` 的**兜底判定**:形态不合规但角色明确时**照样给出 role** —— 否则接续会话/口语命名的主控会被判"角色未知"⇒ 前者冒充主会话、后者解析不出(两条都实测踩过)。 ⛔ 解析不出就返回空(⛔ 不猜);⛔ 绝不回落到 `cwd` 推断(架构 §2.3)。 """ nm = str(name or "").strip() if not nm.startswith("["): # ── 无方括号的两类 ───────────────────────────────────────── head = nm.split("·")[0].split(":")[0].split(":")[0].strip() if head.startswith(MAIN_PREFIX): # 🆕 ④ 主控前缀(实测现役主会话标题就是 `主控 · 协作机制 · …`) return {"role": "main", "topic": "", "ok": False, "form": "main-prefix"} if is_continuation(nm): # 🆕 ③ 无方括号的接续会话(`接续棒:…` / `接续 · 线 · 具体`)⇒ 协作棒 return {"role": "worker", "topic": "", "ok": False, "form": "continuation"} return {"role": "", "topic": "", "ok": False, "form": ""} # 🔴 方括号解析改成"**只剥第一组**"(旧写法 `parts[0].strip("[]")` 会把 `[协作]N9-2156` # 剥成 `协作]N9`、角色判空 ⇒ 靠 `role_of` 里的兜底才没出事;现直接在这里判准)。 _m0 = re.match(r"^\[([^\]]*)\]\s*(.*)$", nm) r = _m0.group(1).strip() if _m0 else "" _rest = _m0.group(2) if _m0 else "" # 🔴 2026-10-01 用户加第三类:**唤醒会话**(判「都空闲」时要把它排除,见 `_anyone_busy`)。 role = {"主": "main", "协作": "worker", "唤醒": "waker"}.get(r, "") if not role: # 🆕 ③ `[<类别>] 接续 · …` —— 第一对方括号里装的是**任务类别**(⛔ 不是角色): # 实测库里就有 `[唤醒机制] 接续 · 日志事前叫停钩子落地(第 2 棒)`。 if is_continuation(nm): return {"role": "worker", "topic": r, "ok": False, "form": "continuation"} return {"role": "", "topic": "", "ok": False, "form": ""} _m1 = re.match(r"^[-\s]*\[([^\]]*)\]", _rest) topic = _m1.group(1).strip() if _m1 else "" return {"role": role, "topic": topic, "ok": bool(role and topic), "form": "prefix"} def topic_of(st: dict, sid: str = "") -> str: """会话的**二级前缀="在协作什么"**:① 命名里的 `[主题]` ② 声明(`--declare --topic`)③ 空。""" if sid: try: con = _db() row = con.execute("select coalesce(nullif(custom_title,''), title, '') from sessions where id=?", (sid,)).fetchone() con.close() tp = parse_session_name(str((row or [""])[0] or "")).get("topic") or "" if tp: return tp except Exception as e: log("topic_of 读标题失败 %s" % e) return str((st.get("topics") or {}).get(sid) or "") def role_of(st: dict, sid: str = "") -> str: """会话角色 —— 🔴 **以「会话命名前缀」为准**(2026-09-29 用户定案): **同工作区、跨工作区都适用**,因为前缀与 `cwd` 无关。 优先级:**① 命名前缀(`[主]`/`[协作]`)② 显式声明(`--declare`)③ 未声明** · 命名前缀:**派活时给自动化命名加前缀** ⇒ 宿主把它带成会话标题(实测:会话标题=自动化名) · 显式声明:会话自己跑 `collabd.py --declare --role main|worker`(sid 自动从环境取)—— 留给"主会话自己改名不方便/标题没前缀"的场合 · ⛔ 都取不到 ⇒ 返回空串(=**未声明**)——调用方必须按"未声明"处理, ⛔ **绝不回落到 cwd 推断**(那正是历史故障的来源) """ if sid: try: con = _db() row = con.execute("select coalesce(nullif(custom_title,''), title, '') from sessions where id=?", (sid,)).fetchone() con.close() t = str((row or [""])[0] or "") pr = parse_session_name(t) if pr.get("role"): # 🔴 2026-10-01 改:判据从 `pr["ok"]`(**两级前缀都齐**)放宽到 `pr["role"]`(**角色明确**)—— # 形态不合规但角色明确的有三支:① 一级前缀 `[协作]N9-…` ② **接续会话** ③ **`主控 · …` 前缀**。 # 旧判据把它们一律判"角色未知" ⇒ 接续会话冒充主会话、口语命名的主控解析不出(都实测踩过)。 if pr.get("form") == "continuation": log("⚠️ 会话 %s 是**接续会话**但标题没按 `[协作]-[<类别>]-<具体>` 命名" " ⇒ 已按 worker 处理(建议改名,否则看板归不了类别)" % str(sid)[:8]) elif not pr["ok"]: log("⚠️ 会话 %s 命名缺 `[<类别>]`(只有一级前缀)⇒ 建议改成 [角色]-[类别]-<具体>" % str(sid)[:8]) return str(pr["role"]) except Exception as e: log("role_of 读标题失败 %s" % e) roles = dict(st.get("roles") or {}) if sid and roles.get(sid): return str(roles[sid]) if not sid: return "" return "" def _main_sid(st: dict) -> str: """主会话的 session id:① 首选**声明为 main** 的那个 ② 退回最近一次投递到的会话 ③ 都取不到 ⇒ 空串(⛔ 不猜)。""" for _sid, _r in (st.get("roles") or {}).items(): if str(_r) == "main" and str(_sid).startswith(("0", "1", "2", "3", "4", "5", "6", "7", "8", "9", "a", "b", "c", "d", "e", "f")): return str(_sid) return str(((st.get("wake") or {}).get("sessionId")) or "") def _same_ws(cwd: str) -> bool: """cwd 是否就是**本工作区**。⚠️ 归一化必须**小写 + 统一斜杠** —— 宿主的分组去重键是 `path.trim().toLowerCase()`(**只小写、不统一斜杠**)⇒ 这里若只做小写,`E:\\x` 与 `E:/x` 会被判成两个工作区。""" a = str(cwd or "").strip().lower().replace("\\", "/").rstrip("/") b = str(WS).strip().lower().replace("\\", "/").rstrip("/") return bool(a) and a == b def _deaf_sids() -> set: """🔴 会话诊断日志已撞 ~10 MiB 上限 ⇒ 宿主拒写(EPERM)⇒ 该会话**界面不再刷新**=哑。 2026-10-01 实测:主会话 f8a792ab 停在 10,485,715 B,累计 45,422 行被丢 (daemon.log `[conversations] diagnostic log write failed` code=EPERM), 而投递仍记 `ok:true · http:200` = **假绿** —— 通知全落空,机制毫无察觉。 🔴 判据(**2026-10-01 修订:只认最新一天**):对**每个会话**,取它**最新一天**的那份日志 (今天有这份就用今天、没有才退到昨天)⇒ ≥ 9.5 MiB 才判哑。 🔴 为什么必须收窄(修「**永久判哑**」):宿主会话日志**按天分目录**(`logs/<日期>/sdk/conversations/`), 昨天那份**永不再增长** ⇒ 旧写法「今天+昨天**任一**命中即判哑」会把 **昨天撞过上限、今天已正常轮转**的会话**永久判哑**(实证 `a202550c`:昨 10.0 MB、 今 4.2 MB 且宿主**仍在写它**)⇒ 本该投出去的唤醒被**静默拦下**(方向相反的新假绿)。 ⚠️ 仍**保留昨天的兜底**:某会话今天**没有**日志(宿主今天没写过它)⇒ 退到昨天那份判 —— 比"直接放行"保守(宁可少投,不可朝已哑界面白投)。 (轮转正常时 .log 会归零重写;到 9.5 说明又快哑了,同样排除 —— 宁可保守。) ⛔ 只读;扫不到 ⇒ 空集(⛔ 不猜)。 """ root = os.path.join(os.environ.get("CODEBUDDY_CONFIG_DIR") or r"E:\ProgramData\.workbuddy", "logs") size = {} # sid -> 它**最新一天**那份日志的字节数 for d in (time.strftime("%Y-%m-%d", time.localtime(time.time() - 86400)), # ⚠️ 先旧 time.strftime("%Y-%m-%d")): # 后新 ⇒ 新的覆盖旧的 cdir = os.path.join(root, d, "sdk", "conversations") try: for fn in os.listdir(cdir): if not fn.endswith(".log"): continue try: size[fn[:-4]] = os.path.getsize(os.path.join(cdir, fn)) except OSError: continue except Exception: continue return {sid for sid, n in size.items() if n >= 9_500_000} def _live_sids() -> set: """当前**活会话** id 集合(各网关口的 `/sessions/live` 取并集)。取不到 ⇒ 空集(⛔ 不猜)。""" out = set() try: for x in discover_gateways(): sid = str(x.get("sessionId") or "") if sid: out.add(sid) except Exception: pass return out def _topic_in_title(title: str, topics=None) -> str: """标题里出现的**任务类别**(🔴 **最长优先** —— 防止"唤醒机制"被"机制"这种短名抢走)。 ⚠️ 判据是**子串**(不要求 `[...]` 包裹)—— 因为用户给主会话的命名是口语式的 (「主控 · 唤醒机制」没有方括号),而协作棒是严格的 `[协作]-[类别]-<具体>`。 两种写法都要能认出来 ⇒ 统一用"标题里出现该类别名"。 ⛔ 都不命中 ⇒ 空串(="这条会话没标类别" ⇒ 调用方按**默认类别**处理,⛔ 不猜)。 """ ti = str(title or "") tps = _as_list(topics) if topics is not None else _goal_topics() hit = "" for _t in tps: if _t and _t in ti and len(_t) > len(hit): hit = _t return hit def _pick_live(cands, live) -> str: """🔴 从候选里挑**活着的那条**(2026-10-01 体检修加)。 为什么要它(实测断点):解析出的主会话 `a202550c`(标题带「主控」)**根本不在活会话里**, 而唯一活着的会话被它压在后面 ⇒ 投递一路报 `main-not-live`(**程序在跑,通知永远发不出去**)。 ⛔ `live` 为空(探测失败)⇒ 回退第一条(**不改旧行为**:探测不到时不许假装知道谁活着)。 🔴 候选**都不活** ⇒ 返回**空**(让调用方落到下一级:主控 ⇒ 工作区最近活动), ⛔ 不许硬认第一条 —— 那正是"解析出来却投不进去(main-not-live)"的来源。 """ c = [str(x) for x in (cands or []) if str(x)] if not c: return "" if not live: return c[0] _deaf = _deaf_sids() # 2026-10-01:哑会话(日志撞上限)⛔ 不当主会话 for x in c: if x in live and x not in _deaf: return x return "" def _scan_mains(st: dict, live=None) -> dict: """🔴 扫一遍宿主库,挑出**本工作区里所有"可能是主会话"的会话**(按最近活动倒序)。 🔴 2026-10-01 加参数 `live`(各网关口的活会话集合):**活着的「主控」排在不活的「主控」前面** (⛔ 传 None ⇒ 行为与旧版逐字一致,`selftest` 的离线用例不受影响)。 ⛔ **排除全部"协作侧"命名**(🔴 2026-10-01 扩):`[协作]-…`(协作棒)、`[唤醒]-…`(唤醒会话)、 以及**接续会话**(`[<类别>] 接续 · …` / `接续棒:…`)—— 三者都是**干活的棒**,⛔ 不是主会话候选。 ⚠️ 旧版只排 `[协作]` ⇒ **接续会话会被收成主会话候选**(实测:`[唤醒机制] 接续 · …` 又被 `_topic_in_title()` 认出类别 ⇒ 直接变成"该类别的主会话")⇒ 通知投给它自己、真主会话被架空。 ⚠️ **不能加 `status='working'`**:主会话在两轮之间是空闲(不是 working)⇒ 加了就**永远漏掉它**。 ⚠️ **也不能按 `is_background_automation` 排**(2026-09-30 实测踩到):**接续会话本身就是自动化起的** —— 连当时那位主会话 `fe146dd9` 都是 `is_background_automation=1` ⇒ 那么排会把**真主会话也排掉**。 @returns `{"cand": [sid…], "named": [sid…], "by_topic": {类别: {sid, source, explicit}}, "topics": [类别…], "err": ""}` · `cand` / `named` 与旧版 `resolve_main()` 的两个列表**逐字同义**(`resolve_main` 就靠它) · `by_topic` = 新增:**按任务类别各挑一条**(用户 2026-09-30「同一个工作区 多会话协作」) """ out = {"cand": [], "named": [], "by_topic": {}, "topics": _goal_topics(), "err": ""} try: con = _db() for sid, cwd, title in con.execute( "select id, cwd, coalesce(custom_title, title, '') from sessions " "where deleted_at is null " "order by last_activity_at desc limit 50"): t = str(title or "") # 🔴 2026-10-01 改:从"只排 `[协作]`"改成"**排掉全部"协作侧"命名**"—— # **接续会话**(`[<类别>] 接续 · …` / `接续棒:…`)与**唤醒会话**(`[唤醒]-…`) # 都是**干活的棒**,⛔ 不是主会话候选。旧判据漏掉它们 ⇒ 接续棒被解析成"该类别 # 的主会话" ⇒ 投递把通知**投给它自己**(**自指死结**,2026-10-01 实测坐实)。 # ⚠️ 判据**只看标题**:⛔ 不看 cwd 归属、⛔ 不看 status(下方两条 ⚠️ 各有一条实测理由)。 _role = str(parse_session_name(t).get("role") or "") if _same_ws(str(cwd or "")) and _role not in ("worker", "waker"): out["cand"].append(str(sid)) # 🔴 「主控」=主会话的**显式前缀**(用户 2026-09-30 定名:「主会话前缀 可以叫 主控」) _explicit = t.strip().startswith("主控") _alive = True if not live else (str(sid) in live) if _explicit: # 活着的「主控」插到队首(⛔ 不活的排后面 —— 那是"解析出来却投不进去"的元凶) (out["named"].insert(0, str(sid)) if _alive else out["named"].append(str(sid))) _tp = _topic_in_title(t, out["topics"]) if _tp: _cur = out["by_topic"].get(_tp) # 同一类别有多条候选 ⇒ ① 显式「主控」优先 ② 否则取**最近活动**(首见即最近) # 🔴 2026-10-01 加第 ③ 条:**现任不活而新候选活 ⇒ 换**(否则同上那个断点) _better = (_cur is None or (_explicit and not _cur.get("explicit")) or (_cur.get("explicit") and not _cur.get("alive") and _alive)) if _better: out["by_topic"][_tp] = { "sid": str(sid), "source": ("prefix:主控" if _explicit else "topic"), "explicit": _explicit, "alive": _alive} con.close() except Exception as e: out["err"] = str(e) log("resolve_main 按工作区解析失败 %s" % e) return out def resolve_mains(st: dict, live_sids=None) -> dict: """🔴 **按任务类别解析主会话**(= `resolve_main()` 的多类别版)。 2026-09-30 用户要求:「要能支持**同一个工作区** 多会话协作(主会话根据任务**自动梳理任务类别**: **通过协作会话名称前缀的方式区分具体任务会话**)」 ⇒ 同一工作区里会**并存多个任务类别**,**每个类别有自己的主会话** ⇒ 投递必须**按件所属的类别**选主会话,⛔ 否则会**投错窗口**(投到另一条类别的会话里)。 🔑 **规则(简单、可判,不猜)**: ① 每个类别各挑一条:**显式 `主控` 优先 → 否则最近活动**(见 `_scan_mains`) ② 某类别**没扫到** 且 它就是**默认类别**(`topics[0]`)⇒ 回落 `default` (⛔ 这样单类别时的行为与旧版**完全一致**,向后兼容) ③ 某类别**没扫到** 且 不是默认类别 ⇒ **sid 留空** ⇒ 调用方**明确报"找不到该类别的会话"** (⛔ **绝不回落到别类别的主会话** —— 那正是"投错窗口") `default` 的取法与 `resolve_main()` **同源同判据**: 登记有效 → 显式「主控」→ 工作区最近活动 → 空 @returns `{"topics": [...], "mains": {类别: {sid, source, switched_from}}, "default": {sid, source, switched_from}, "all_sids": [...], "err": ""}` """ reg = _main_sid(st) live = _live_sids() if live_sids is None else set(str(s) for s in live_sids if s) sc = _scan_mains(st, live) tps = list(sc["topics"]) named, cand = sc["named"], sc["cand"] # `default`:与 resolve_main 逐条相同(登记有效 ⇒ 认;否则 主控 ⇒ 最近活动), # 🔴 且两条**都只认活着的**(2026-10-01:`_pick_live`)。 _nm = _pick_live(named, live) _cd = _pick_live(cand, live) if reg and (not live or reg in live): dft = {"sid": reg, "source": "roles:main", "switched_from": ""} elif _nm: dft = {"sid": _nm, "source": "prefix:主控", "switched_from": reg if (reg and reg != _nm) else ""} elif _cd: dft = {"sid": _cd, "source": "workspace", "switched_from": reg if (reg and reg != _cd) else ""} else: dft = {"sid": "", "source": ("stale-roles" if reg else "no-register"), "switched_from": ""} mains: dict = {} for _tp in tps: _got = sc["by_topic"].get(_tp) or {} _sid = str(_got.get("sid") or "") if not _sid and tps and _tp == tps[0]: # ② 默认类别兜底(⛔ 只为**向后兼容**单类别部署:那时 topics 只有一项,这里=dft) _sid = str(dft.get("sid") or "") _src2 = str(dft.get("source") or "") or "default" # ⚠️ 只有当"默认那条**确实属于本类别**(或本类别就是唯一类别)"时才认,否则会串类别 if _sid and len(tps) > 1 and _topic_in_title(_title_of(_sid), tps) not in ("", _tp): _sid, _src2 = "", "" else: _src2 = str(_got.get("source") or "") # 🔴 2026-10-01:**该类别选出的那条不活** ⇒ 不能拿它投递(`main-not-live`)。 # · 单类别(≤1)⇒ 用 `dft`(已按活挑过);⛔ 多类别 ⇒ **置空**(宁可报"找不到",⛔ 不串类别) if _sid and live and _sid not in live: _sid, _src2 = (str(dft.get("sid") or ""), "live-fallback") if len(tps) <= 1 else ("", "") if _sid: _src2 = "live-fallback" mains[_tp] = {"sid": _sid, "source": _src2 or ("none" if _tp else ""), "switched_from": (reg if (reg and _sid and _sid != reg) else "")} _all = [] for _m in (list(mains.values()) + [dft]): _s = str(_m.get("sid") or "") if _s and _s not in _all: _all.append(_s) return {"topics": tps, "mains": mains, "default": dft, "all_sids": _all, "err": sc["err"]} def main_for_topic(rms: dict, topic: str = "") -> dict: """从 `resolve_mains()` 的结果里取**某个任务类别的主会话**。 🔑 两分支(⛔ 不猜): · `topic` **是已登记的任务类别** ⇒ **严格**取该类别的解析结果(可能 sid 为空 ⇒ 调用方**降级报告**) · `topic` **不是已登记类别**(空 / 旧数据里的工作区名等)⇒ 回落 `default` (⛔ 这是**向后兼容** `resolve_main()` 的旧调用点,⛔ 不是"盲投" —— default 仍是解析出来的) """ rms = rms or {} tp = str(topic or "").strip() if tp and tp in (rms.get("topics") or []): m = (rms.get("mains") or {}).get(tp) or {} return {"sid": str(m.get("sid") or ""), "source": str(m.get("source") or ""), "switched_from": str(m.get("switched_from") or ""), "topic": tp} d = dict(rms.get("default") or {}) d["topic"] = "" return d def _title_of(sid: str) -> str: """按 sid 读会话标题(拿不到 ⇒ 空串,⛔ 不猜)。""" if not sid: return "" try: con = _db() row = con.execute("select coalesce(nullif(custom_title,''), title, '') from sessions where id=?", (sid,)).fetchone() con.close() return str((row or [""])[0] or "") except Exception: return "" def resolve_main(st: dict, live_sids=None) -> dict: """🔴 解析**当前**主会话 —— 不只是读登记(用户 2026-09-30:「协作机制能发现 主会话 换了吗」)。 旧实现的答案是**不能**:`_main_sid()` 只有"① 登记为 main ② 退回上次投给谁"两条路, 两条都**静态/路径依赖**;而且看板与投递**共用同一份登记**(`board.py::_main_sid` 与这里逐字同款) ⇒ 换主会话后**一起继续指向旧会话**,连交叉校验都没有 ⇒ **静默断链**。 更糟的是投递侧原有一处**盲选回落**(目标不在活会话里就取"第一个带会话的口")⇒ 可能**投错窗口**。 本函数把"主会话"从**登记出来的**改成**解析出来的**: ① 登记为 main 且**此刻确实活着** ⇒ 认(登记有效) ② 登记失效(不在活会话里)⇒ 按**工作区**解析:`cwd == 本工作区` 且**不是协作棒**(标题不带 `[协作]`) ⇒ 取最近活动的那条(用户 2026-09-30:「以**一个工作区**为主会话的工作区」) ③ 还是解析不出 ⇒ `sid=""` ⇒ **调用方⛔ 不许盲投**,必须明确报"找不到主会话" @returns `{sid, source, switched_from}`;`switched_from` 非空=**换了**(调用方须跟随+告警) """ reg = _main_sid(st) live = _live_sids() if live_sids is None else set(str(s) for s in live_sids if s) if reg and (not live or reg in live): return {"sid": reg, "source": "roles:main", "switched_from": ""} # ② 按**工作区**认(用户 2026-09-30:「以 **一个工作区** 为主会话的工作区」+ # 「是**一个**工作区,下面**应该都可能是主会话**,一般不会切其他工作区」)—— # 所以判据只用**工作区**这一个锚点:本工作区里那一条,就是它。⛔ 不要求标题匹配目标名 # (实测主会话标题叫「接续 · 机制线(…)」,**根本不含目标短名** ⇒ 按标题认必然落空)。 # ⚠️ **不能加 `status='working'`**:主会话在两轮之间是空闲(不是 working)⇒ 加了就**永远漏掉它**。 # ⚠️ **也不能按 `is_background_automation` 排**(2026-09-30 实测踩到):**接续会话本身就是自动化起的** # —— 连当前这位主会话 `fe146dd9` 都是 `is_background_automation=1` ⇒ 那么排会把**真主会话也排掉**。 # 只按**显式标记**排除"棒":标题带 `[协作]`(那是本机制给协作棒定下的前缀,见 SKILL.md 命名)。 sc = _scan_mains(st, live) # 🔴 与 `resolve_mains()` **共用同一遍扫描**(⛔ 不各扫各的=防两边漂移) named, cand = sc["named"], sc["cand"] # 优先认**显式标了 `主控`** 的那条 —— 它比"最近活动"可信得多(用户自己标的,不是猜的) # 🔴 2026-10-01:但**必须活着**(`_pick_live`)—— 不活的"显式标记"会把投递卡死在 `main-not-live`。 _nm = _pick_live(named, live) if _nm: return {"sid": _nm, "source": "prefix:主控", "switched_from": reg if (reg and reg != _nm) else ""} _cd = _pick_live(cand, live) if _cd: return {"sid": _cd, "source": "workspace", "switched_from": reg if (reg and reg != _cd) else ""} return {"sid": "", "source": ("stale-roles" if reg else "no-register"), "switched_from": ""} def _follow_main(st: dict, rm: dict) -> None: """主会话**换了** ⇒ 跟随 + **显式告警**(⛔ 不静默)。""" old, new = str(rm.get("switched_from") or ""), str(rm.get("sid") or "") if not new: return try: roles = st.setdefault("roles", {}) if old and str(roles.get(old)) == "main": roles[old] = "worker" roles[new] = "main" except Exception: pass log("主会话变更:%s -> %s(依据 %s)⇒ 已跟随" % (old[:8] or "?", new[:8], rm.get("source"))) need_user("主会话从 %s 换到了 %s(依据:%s)—— 本程序**已自动跟随**并把后续通知投给新主会话;" "若判断有误,请纠正登记(`roles` / `--declare`)。" % (old[:8] or "(无)", new[:8], rm.get("source"))) def _main_state(st: dict) -> str: """主会话**当前状态**(`working` ⇒ 正在执行)。取不到 ⇒ 空串(⛔ 不猜)。""" sid = _main_sid(st) try: con = _db() if sid: row = con.execute("select status from sessions where id=?", (sid,)).fetchone() else: row = con.execute("select status from sessions where status='working' limit 1").fetchone() con.close() return str(((row or [""])[0]) or "").lower() except Exception as e: log("main_state 取失败 %s" % e) return "" WAKER_PREFIX = "[唤醒]" """🔴 2026-10-01 用户定:**唤醒会话**用**专属名字前缀** `[唤醒]` 认。 判「都空闲」时必须把它排除 —— 它自己就是被叫醒的那条会话, 不排除 ⇒ 它用自己的存在否掉「都空闲」⇒ **永远不叫**(自指死结)。""" def _busy_sessions() -> list: """当前 `status='working'` 的会话 `[(id, title), …]`;读不到 ⇒ `[]`(⛔ 不猜)。""" try: con = _db() rows = con.execute( "select id, title from sessions where status='working'").fetchall() con.close() return [(str(r[0] or ""), str(r[1] or "")) for r in rows] except Exception as e: log("busy_sessions 取失败 %s" % e) return [] def _anyone_busy(rows: list) -> bool: """🔴 **除唤醒会话(和判断者自己)外,还有没有会话在跑**。用户 2026-10-01 原话: 「**把自己排除不就行了,通过会话名称前缀区分**,还有**也不用他自己判断, 可以通过代码判断**」 两道排除(互补,都要): · `SELF_SID` —— **精确**:程序跑在某个会话里(钩子/会话内)时,那就是「自己」; · `WAKER_PREFIX` —— **按名字**:程序在会话外跑(常驻监督程序)时拿不到 `SELF_SID`, 只能认名字(唤醒会话前缀 `[唤醒]`)。 ⚠️ 判据由**代码**做,⛔ 不让会话自己去查 —— 会话自查必然把自己算进去 ⇒ 条件①永不成立。`rows` 由调用方一次性取好(⛔ 别在循环里反复查库)。 """ for sid, title in rows: if SELF_SID and sid == SELF_SID: continue # 判断者自己(精确) if str(title or "").strip().startswith(WAKER_PREFIX): continue # 唤醒会话(按名字) return True return False def reconcile(st: dict) -> dict: """启动对账(收编)—— 治「守护启动前已有协作会话在跑」:那些棒没进过台账 ⇒ 监督程序看不见 ⇒ 主会话以为没人在跑 ⇒ 可能重复派活。 收编:宿主库 working 会话(排除观察者自己)若台账没有 ⇒ 建条目 sess:<前8位>,state=running, _source=启动对账收编;线留空并标「待确认」(不从标题猜线)。 僵尸:台账里 running 但除自己外无人跑、也无 IN_PROGRESS ⇒ 标「有阻碍 ⇒ 待主会话重派」。 不投递、不改别人文件;幂等。 """ out = {"adopted": [], "zombie": [], "skip": ""} tk = _load_tasks() try: con = _db() allrun = [(str(r[0] or ""), str(r[1] or "")) for r in con.execute("select id, coalesce(nullif(custom_title,''),title,'') " "from sessions where status='working'")] inprog = int(con.execute("select count(*) from automation_runs where status='IN_PROGRESS'").fetchone()[0]) con.close() except Exception as e: out["skip"] = "读宿主库失败:%s" % e return out me = (os.environ.get("CODEBUDDY_SESSION_ID") or "")[:8] others = [(s, ti) for s, ti in allrun if not (me and s[:8] == me)] for sid, title in others: key = "sess:" + sid[:8] if key in tk: continue tk[key] = {"state": "running", "by": (title[:40] or sid[:8]), "t": time.time(), "title": title[:80], "line": "", "_source": "启动对账收编(队列开启前就已在跑)", "_line": "待确认(不从标题猜线)"} out["adopted"].append(key) if not others and not inprog: for k, v in list(tk.items()): if str((v or {}).get("state")) == "running": nv = dict(v or {}) nv["state"] = "blocked" nv["block_reason"] = "执行者已消失(除观察者外无在跑会话、无 IN_PROGRESS 运行)⇒ 待主会话重派" nv["t"] = time.time() tk[k] = nv out["zombie"].append(k) if out["adopted"] or out["zombie"]: _save_tasks(tk) log("reconcile adopted=%s zombie=%s" % (out["adopted"], out["zombie"])) return out def supervise(st: dict, deliver: bool = True, mutate: bool = True) -> dict: """**监督程序职责** —— 🔴 **反馈协议 = 单条 + 握手**(2026-09-29 用户定案): ① **一次只反馈一条**任务状态 ⇒ 退回等主会话处理; ② 反馈成功后进入**等待状态**:**定期监督主会话是否在执行**(`sessions.status='working'`); ③ **主会话执行结束**(working 消失)⇒ 才反馈**下一条**; ④ 另:**一段时间没有「执行中/执行完毕」的队列** ⇒ 发心跳让主会话检查状态。 ⚠️ 等待期间**不发心跳**(⛔ 不叠加消息,一条没处理完就不发第二条)。 ⛔ 只读队列 + 写通知;⛔ 不派活、⛔ 不开会话。 ⚠️ 通知正文**不含秒级时间**(逐字稳定,否则内容哈希去重失效 ⇒ 刷屏)。 ⛔ 兜底:反馈后 **5 分钟**(2026-09-30 由 20 分钟提速)仍未见主会话执行 ⇒ 放行下一条并记「需用户介入」,⛔ 不无限卡死。 🔴 `mutate`(2026-09-30 加)—— **只有「投递方」才配推进队列**: · `deliver=True, mutate=True`(`--tick`,**宿主钩子唤起**)= **唯一的投递方 ⇒ 唯一的推进方**; · `deliver=False, mutate=False`(`--once`/协作程序)= **纯投影**:只算、只写 `TO_MAIN.md`, ⛔ **绝不**改 `notify_pending`/`notify_awaiting`/`notified`/`hb_since`/`q_sig`/`pw_fp`。 ⚠️ 为什么必须分开(今天实测到的一型**静默丢件**):钩子在 `UserPromptSubmit` 上**先**跑 `--once` (节流 3 分钟,谁发话都会跑)—— 若它也推进队列,就会**把待反馈项"消费"掉却不投递** ⇒ 紧接着的 `--tick` 看到队列已空 ⇒ **通知永远发不出去**(表现为"程序在跑,主会话什么也没收到")。 """ t = _load_tasks() now = time.time() # 🔴 **口令指纹探针**(2026-09-29 用户问"口令会变化吗" ⇒ 不猜,装探针): # 指纹变了 ⇒ 记日志 + 落 `NEED-USER.md`(常驻进程需重启一次才能继续投递)。 _fp = pw_fingerprint() _prev = str(st.get("pw_fp") or "") if mutate and _fp and _prev and _fp != _prev: log("⚠️ 网关口令指纹变化:%s -> %s" % (_prev, _fp)) need_user("网关口令**已变化**(指纹 %s -> %s)⇒ 进程需重启一次才能继续投递" % (_prev, _fp)) if _fp and mutate: st["pw_fp"] = _fp sig = "|".join("%s=%s" % (k, (v or {}).get("state")) for k, v in sorted(t.items())) if mutate and sig != str(st.get("q_sig") or ""): st["q_sig"] = sig st["q_transition_at"] = now notified = dict(st.get("notified") or {}) info = {"n": len(t), "awaiting": "", "sent": "", "phase": "", "pw": {"have": bool(_fp), "fp": _fp}} # ① 待反馈序列:**按发生顺序**排队(⛔ 不一次倾倒) pend = list(st.get("notify_pending") or []) for k, v in sorted(t.items()): stt = str((v or {}).get("state") or "") if stt in ("running", "done") and notified.get(k) != stt and ("%s=%s" % (k, stt)) not in pend: pend.append("%s=%s" % (k, stt)) if mutate: st["notify_pending"] = pend aw = st.get("notify_awaiting") or None if aw: # ②③ 等待状态:定期监督主会话是否在执行 ⇒ 执行结束才放行下一条 ms = _main_state(st) phase = str(aw.get("phase") or "wait-start") if phase == "wait-start": if ms == "working": # 主会话**已在执行** ⇒ 进第二阶段 if mutate: aw["phase"], aw["t_working"] = "wait-done", now log("握手:主会话已开始执行(%s)" % aw.get("item")) elif mutate and now - float(aw.get("t") or now) > 5 * 60: log("握手:等主会话 5 分钟未见执行(%s)⇒ 放行并记 NEED" % aw.get("item")) need_user("上报后 5 分钟未见主会话执行:%s" % aw.get("item")) st.pop("notify_awaiting", None) aw = None if aw and mutate and str(aw.get("phase")) == "wait-done": if ms != "working": # **执行结束** ⇒ 放行下一条 log("握手:主会话执行结束(%s,执行 %.0f 秒)⇒ 放行下一条" % (aw.get("item"), now - float(aw.get("t_working") or now))) st.pop("notify_awaiting", None) aw = None if aw: info["awaiting"] = str(aw.get("item")) info["phase"] = str(aw.get("phase")) # ①′ 没有「等待中」才发**下一条**(单条 + 握手) if not aw and pend: item = pend[0] kid, _, stt = item.partition("=") v = t.get(kid) or {} head = ["# 队列变化(上报 -> 主会话)", "", "本轮**只反馈这一条** —— 你处理完(我一看到你执行结束)才会发下一条:", "", "- **%s** -> %s%s%s" % (kid, stt, ("(%s)" % v.get("by")) if v.get("by") else "", (" 产物:%s" % v.get("artifact")) if v.get("artifact") else "")] if stt == "blocked": head += ["", "> 🔴 本条状态 = **有阻碍**:原因 —— %s" % (str(v.get("block_reason") or "(未填原因,⛔ 属不合格上报)")), "> ⇒ 除下面三步外,**必须明确向用户喊「需用户介入」**(谁在等、要哪一句话)。"] txt = "\n".join(head + [ "", "## 主会话该做的", "1. 按产物核对是否真做完(不认自述)", "2. 需要就写一行排期(派活)", "3. ⛔ 不重派已在跑/已完成的件", ]) + "\n" try: TO_MAIN.write_text(txt, encoding="utf-8") except Exception: pass info["sent"] = item if not (deliver and mutate): # 纯投影(`--once`/协作程序):只写通知文件,⛔ 不消费队列 info["deliver"] = {"skipped": "read-only"} return info # 🔴 2026-09-30 加 `topic=`:反馈**按件所属的任务类别**投给它自己的主会话 # (用户「同一个工作区 多会话协作」⇒ 一个工作区里多个类别各有各的主会话)。 # ⛔ 件没标 `line`(旧数据)⇒ 传空 ⇒ 走 `default`(与旧版一致)。 _dv = _deliver_str(txt, "上报·单条", st, topic=str((v or {}).get("line") or "")) info["deliver"] = _dv # 🔴 2026-09-30 修(**静默丢件**):⛔ **只有真投出去了才允许消费队列**。 # 旧实现**无条件** `notified[kid] = stt` ⇒ 被 `target-busy`/`too-soon`/`locked`/ # `no-token` 挡下时,这一条**从此消失**(既没投到主会话、也不再重试) # = 又一型「程序在跑、主会话什么也没收到」。⇒ 未投出 ⇒ **队列原样保留**,下一轮再试。 # ⚠️ `same-item` 例外:内容哈希与上次**逐字相同** ⇒ 说明主会话本来就收到了 ⇒ 允许消费。 if not (_dv.get("delivered") or _dv.get("skipped") == "same-item"): log("投递未成(%s)⇒ 保留 %s 在队首,下一轮重试" % (_dv.get("skipped") or "?", item)) return info st["notify_awaiting"] = {"item": item, "t": now, "phase": "wait-start"} st["notify_pending"] = pend[1:] notified[kid] = stt st["notified"] = notified return info # ④ **心跳**(🔴 2026-09-29 用户定案 —— **三条件合取,⛔ 不是计时器**): # (a) 定期探测到主会话**未在处理**(不在 `working`) # ∧ (b) 队列中**没有待反馈的任务**(无 awaiting ∧ 无 pending) # ∧ (c) **需求仍未完成**(任务图里还有非 done 的节点) # ⇒ 触发心跳,让主会话**核对需求推进状态**。 # ⚠️ 同一状态下**最多每 QUEUE_IDLE_MIN 分钟一次**:正文逐字稳定(⛔ 不带秒级时间)+ # **持续时长按 30 分钟桶**写进正文 ⇒ 靠内容哈希天然限流,⛔ 不会每次探测都刷屏。 # 🔴 2026-10-01 用户收窄 (a):「**是主会话 和 协作会话都空闲**」—— # 原实现只看主会话 ⇒ 协作会话还在跑时也会叫主会话(多余打扰); # ⚠️ 且必须**排除「唤醒会话」自己**(前缀 `[唤醒]`):它正在跑,不排除 # ⇒ 它用自己的存在否掉「都空闲」⇒ **永远不叫**(用户点名的自指死结)。 main_busy = (_main_state(st) == "working") or _anyone_busy(_busy_sessions()) # ⚠️ 用**刚算出来的** `pend`,⛔ 不用 `st["notify_pending"]`:投影轮(mutate=False)不写 st, # 读 st 会拿到**陈旧值**(可能已空)⇒ 误判"无待反馈"。 no_fb = (not aw) and (not pend) goal_open = goals_open() # 🔴 2026-10-01 用户:「**需求有阻碍(唤醒任务暂停)**」「需求已完成(唤醒任务暂停)」 # ⇒ 只剩受阻项、没有别的活 ⇒ **不叫**(阻碍已在上报里通报过,重复叫没意义)。 # ⚠️ 若还有 `pending`/`running` 的项 ⇒ 需求仍算在推进 ⇒ 照常叫。 _blk = [k for k, v in (t or {}).items() if str((v or {}).get("state")) == "blocked"] _mov = [k for k, v in (t or {}).items() if str((v or {}).get("state")) in ("pending", "running")] goal_stuck = bool(_blk) and not _mov # 🔴🔴 2026-09-30 加第 4 个条件「**真停滞**」(**用户报"发消息卡"查出来的**): # 现象:用户"发消息发不出去卡住,看上去发出去了、实际没有"。 # 取证:① 05:20 后**全部会话** 14 条 prompt **全部 `resolveWaiter`、`park=0`** ⇒ **没丢消息**; # ② 主会话 05:40 后 **`busy=true` 占 324/342(95%)**。 # 因果:**每次"程序投递"都会唤起主会话跑一轮 ⇒ 计入 busy**;而心跳从 30 分钟改到 **5 分钟(×6)** # ⇒ 主会话被唤起次数 ×6 ⇒ **用户发消息时更容易撞上 busy ⇒ 排队 ⇒ 体感"发不出去"**。 # 根因(语义错):心跳的触发只要求"三条件成立",而**"目标未完成"在长任务期长期成立** # ⇒ 心跳变成"**目标没完成就定期提醒**"(噪音),⛔ 而它本该表达"**异常停滞**"。 # ⇒ **加一条**:距**上次任何进展**(`last_progress_at`)≥ `HEARTBEAT_MIN_IDLE` 分钟才发。 # **有进展 ⇒ 一条都不发**;**真停滞 ⇒ 才发**(且保留 5 分钟桶)。⇒ 既快发现停滞,又不无事刷。 _prog = float(st.get("last_progress_at") or 0) prog_age_min = ((now - _prog) / 60.0) if _prog else 1e9 info["probe"] = {"main_busy": main_busy, "no_feedback": no_fb, "goal_open": goal_open, "goal_stuck": goal_stuck, "prog_age_min": round(prog_age_min, 1)} if ((not main_busy) and no_fb and goal_open and (not goal_stuck) and prog_age_min >= HEARTBEAT_MIN_IDLE): hb_since = float(st.get("hb_since") or 0) or now if mutate: st["hb_since"] = hb_since held = (now - hb_since) / 60.0 bucket = int(held // float(QUEUE_IDLE_MIN)) * int(QUEUE_IDLE_MIN) open_ids = [k for k, v in sorted(t.items()) if (v or {}).get("state") != "done"] info["heartbeat"] = {"held_min": round(held, 1), "bucket": bucket} txt = "\n".join( ["# 唤醒:请核对需求完成状态(上报 -> 主会话)", "", goal_line(), "", "**触发条件同时成立**:主会话与协作会话都空闲(⛔ 不含唤醒会话自己)· 队列无待反馈 · **需求仍未完成**" "(已持续 %d 分钟)。" % bucket, "", "## 请你做的", "1. 对着任务图核对:哪些还没做、哪些卡住了、下一步该派谁", "2. 需要用户拍板 -> 明确喊「需用户介入」", "", # 🔴 2026-09-30 加(用户「加快效率」):**把"还差什么"直接写进心跳正文**, # 省掉主会话每轮自己去翻目标文件 —— 心跳一到手就知道该干什么。 "## 目标验收:还差这些(非 pass 的项)"] + _acc_short() + [ "", "## 队列现状(未完结 %d)" % len(open_ids)] + (["- **%s**:%s" % (k, (t[k] or {}).get("state")) for k in open_ids] or ["- (队列为空 —— 待办在任务图/旧队列里,尚未并入)"]) ) + "\n" try: TO_MAIN.write_text(txt, encoding="utf-8") except Exception: pass info["notice"] = "心跳" # 🔴 2026-09-29 修(实测:18 秒内两次同 hash 投递 ⇒ 去重失效): # 根因=**两个进程并发读写同一个 state 文件、互相覆盖**(协作程序 one_round 与监督程序 # 的 `--supervise` 循环都会投)。⇒ **投递只由专用监督进程做**(架构上"投递"本就是它的职责), # 另一处调用传 `deliver=False`(仍写通知文件,⛔ 不投)。 if deliver and mutate: info["deliver"] = _deliver_str(txt, "上报·唤醒", st) else: info["deliver"] = {"skipped": "read-only" if not deliver else "not-deliverer"} else: if mutate: st.pop("hb_since", None) if not aw: # ⚠️ 正在等主会话处理(aw 在身)时**别删**刚投出去的通知文件 try: TO_MAIN.unlink() except Exception: pass return info # ── 视图 ──────────────────────────────────────────────────── def render(d: dict, V: dict, H: dict, T: dict, tk: list) -> None: Q_ = queue_view(T, d) # 🔴 严格队列(一次一件 · 原子取件) sh = [] for s in d["sessions"]: a = s["age"] sh.append("- `%s` %-36s %s · %s · %s" % ( s["id"], s["name"], "🟢 活跃" if a < 5 else ("🟡 可能停了" if a < 30 else "⚪ 闲"), ("%.0f 分钟前" % a) if a < 1e8 else "—", s["tag"])) # ★ 2026-10-01 改(用户:看板与注入正文同族,排版必须一致): # 本块原为 `- ` **碎片流**(一行一个半句、`·`/`;`/`⇒` 硬拼),用户原话「**非常不利于阅读**」。 # ⇒ 统一改「段落排版」:`##` 小标题 + **每个无序段落以圆点开头、一句话说清一个事实** # (复杂信息放句末括号内),⛔ 不再把多个事实用 `;` 串成一行。 md = "\n".join([ "# 协作实时状态", "", "每 %ss 覆写一次,由协作程序生成,不需要任何会话在跑。" % C["interval"], "", "## 概览", "- 快照生成于 %s。" % time.strftime("%Y-%m-%d %H:%M:%S"), "- 本轮反馈新增 %d 条,水位到 %s(其中真成果 %d 条、新排期 %d 条、在跑 %d 个)。" % (len(V["fresh"]), V["wm"], len(V["done"]), len(V["newp"]), len(V["cur"])), "- 推进判定:%s。" % V["verdict"], "- 真空:%s。" % ("超过 %.0f 分钟需求未完成且无人执行,需主会话接管" % V["vdur"] if (V["vacuum"] and V["vdur"] >= float(C["vacuum_min"])) else ("真空中已 %.0f 分钟" % V["vdur"] if V["vacuum"] else "无")), "- 体检:%s。" % ("跳过(忙)" if H["skipped"] else ("通过" if not H["issues"] else "%d 项待看" % len(H["issues"]))), "- 告警:%s。" % (V["alert"] or "无"), "- 进程:pid %d,单例 :%s,零令牌、只读、不派活。" % (os.getpid(), C["singleton_port"]), "", "## 🎯 服务探针", "- 前置端口 `%s`:%s。" % (C["shim_port"] or "(未配置)", "通" if V["up"] else "不通"), "", "## 🧭 任务图", ("- 可派未派(浪费):%s,有能干的活却没人接。" % "、".join("%s(%s)" % (r[0], r[2]) for r in T["ready"])) if T.get("waste") else ("- 可派:%s。" % "、".join("%s(%s)" % (r[0], r[2]) for r in T["ready"]) if T.get("ready") else "- 可派:无(防干等)。"), *["- 等待:%s(等前置 %s)。" % (w[0], "、".join(w[3] or [])) for w in (T.get("waiting") or [])[:6]], ("- 关键路径:%s。" % " → ".join(T["critical"])) if T.get("critical") else "- 关键路径:未配置。", "", "## 📋 严格队列", "- 队首(下一个该做):%s。" % ("%s %s(线:%s,%s)" % ( Q_["head"]["id"], Q_["head"]["title"], Q_["head"]["line"], "在关键路径上游" if Q_["head"]["prio"] == 0 else "非关键路径") if Q_.get("head") else "无可派"), "- 正在做:%s。" % ("、".join("%s→%s" % (k, v) for k, v in (Q_.get("doing") or {}).items()) or "无人"), "- 待办:%d 件%s。" % (Q_.get("pending", 0), "(同时在做多件,可能打架)" if Q_.get("conflict") else ""), "- 取件规矩:只取队首,用 `mkdir claims/` 原子取件,做完删掉 claim 出队,卡死 %d 分钟自动回退。" % (DOING_TTL // 60), "", "## ✅ 真成果", *(["- `[%s]` %s · %s" % (r["status"], r["aid"], r["title"][:140] or "(空)") for r in V["done"][:5]] or ["- 无。"]), "", "## 🟡 接续任务(只是计划,不是成果)", *(["- %s · %s" % (str(r[0])[:8], (r[1] or "")[:46]) for r in V["newp"][:6]] or ["- 无。"]), "", "## 🖥 会话状态", ("- 持锁:%s。" % "、".join(h[:30] for h in d["locks"])) if d["locks"] else "- 无锁。", "", *(sh or ["- 无会话。"]), "", "## 🧭 靶点", *(tk or ["- 未配置。"]), "", "## 🫀 在跑的棒(口径 = `sessions.status='working'`,不含本会话)", *(["- `%s` %s" % (w["id"][:8], w["name"]) for w in d["working"]] or ["- 无(没有任何会话在做正事%s)。" % (",只有本会话这个观察者在跑" if d.get("self_working") else "")]), "", "> 本程序承担「派活之外的全部功能」,不派活(派活须由会话做)。", "", ]) LIVE.parent.mkdir(parents=True, exist_ok=True) LIVE.write_text(md, encoding="utf-8") # ★ 2026-10-01 改:摘要改「段落排版」(见 `digest_text` docstring)—— # 原实现把 `verdict` 与 `goal_line()` **原样拼成 `- ` 碎片**,用户反馈「非常不利于阅读」。 (INBOX / "digest.md").write_text(digest_text(d, V, T), encoding="utf-8") def _dur_txt(m: float) -> str: """分钟数 → 人读的时长文本(供摘要段落使用)。""" return ("%.0f 分钟" % m) if m < 90 else ("%.1f 小时" % (m / 60.0)) def _plain_goal() -> str: """目标行去掉装饰前缀(`🎯 围绕目标:`),供「标签:段落」排版复用。""" g = goal_line() for pre in ("🎯 围绕目标:", "🎯 "): if g.startswith(pre): return g[len(pre):] return g def digest_text(d: dict, V: dict, T: dict) -> str: """会话侧的「机械摘要」—— **段落排版**(2026-10-01 改)。 🔴 用户原话:「会话的反馈信息排版 **非常不利于阅读**,改为**段落排版**」,并给出目标形态: **一行 header + 其后每行「标签:一段话」**。 ⛔ 因此这里**不用 `- ` 列表符号**、⛔ **不用 `;`/`⇒` 串联短句** —— 那是给程序看的形态,不是给人看的。 🔴 为什么值得改:本文件是**每轮钩子注入给每个会话**的正文 ⇒ 排版差=每个会话每轮都要 额外花注意力去拆句,是**全平台共用的阅读成本**(同 §二十九「无成本反馈」那类隐形成本)。 ⚠️ 事实来源:全部取自 `V` / `T` 的**结构化字段**,⛔ 不对 `V["verdict"]` 做字符串反解 (反解会随措辞改动而静默失效)。`verdict` 原串仍保留,供 `realtime.md` 与通知使用。 """ if V["done"]: st = "本轮有 %d 条任务真正做出了结论。" % len(V["done"]) elif V["newp"]: st = "本轮只排了 %d 条接续任务,还没有出现新成果,排期不等于成果。" % len(V["newp"]) elif d["working"]: st = "有 %d 个会话正在跑,暂时没有新成果,已经 %s没有进展。" % ( len(d["working"]), _dur_txt(V["idle"])) else: st = "没有任何会话在跑,而且已经 %s没有出新成果。" % _dur_txt(V["idle"]) if V.get("zero_sched"): st += "未来 1 小时内也没有任何排期,只要没人主动开口,这里就会一直静默下去。" if V.get("probe") == "rec": st += "服务探针刚从不通恢复成通,但那是服务层的变化,不算成果、也不会清零停滞计时。" elif V.get("probe") == "lost": st += "服务探针从通变成了不通,服务层有变化。" rd = T.get("ready") or [] if T.get("waste"): rdy = "有能干的活却没有任何棒在跑,现在就该派:" + "、".join( "%s(%s)" % (r[0], r[2]) for r in rd) + "。" elif rd: rdy = "可以派下一棒:" + "、".join("%s(%s)" % (r[0], r[2]) for r in rd) + "。" else: rdy = "没有可派的任务。" wt = T.get("waiting") or [] if wt: wtx = "".join("%s 在等 %s。" % (w[0], "、".join(w[3] or []) or "前置条件") for w in wt[:6]) else: wtx = "没有任务在等前置条件。" cp = T.get("critical") or [] gl = _plain_goal() return "\n".join([ "# 机械摘要", "", "状态:%s" % st, "目标:%s" % (gl if gl.endswith(("。", ")", "`")) else gl + "。"), "可派:%s" % rdy, "等待:%s" % wtx, "关键路径:%s" % ("%s。" % " → ".join(cp) if cp else "尚未配置。"), "", ]) def _sig_head(alert: str) -> str: """信号文件的 header 行(原先写在标题里,且带 `**` 与 `+` 等排版噪声)。""" return str(alert or "").replace("**", "").strip() def signals(V: dict, T: dict) -> None: def w(p, txt): try: p.write_text(txt, encoding="utf-8") except Exception: pass def rm(p): try: p.unlink() except Exception: pass # ★ 2026-10-01 改(用户:会话的反馈信息「非常不利于阅读」): # 三个信号文件与 `digest.md` **同族、同被注入给会话** ⇒ 一律改「段落排版」: # **一行 header + 其后每行「标签:一段话」**;⛔ 不用 `- ` 列表、⛔ 不用 `1. 2. 3.` 步骤碎片。 if V["alert"]: w(STALL, "# ⚠️ %s\n\n" "发生时间:%s\n" "判定依据:程序按「有没有新成果、有没有棒在跑、未来有没有排期」三条自动判出上面这条告警。\n" "该谁做:本程序不派活。请主会话抢一把细粒度域锁,读完摘要与任务图后按缺口派下一棒(属白名单「派活」)。\n" % (_sig_head(V["alert"]), time.strftime("%Y-%m-%d %H:%M:%S"))) else: rm(STALL) if V["vacuum"] and V["vdur"] >= float(C["vacuum_min"]): w(VACUUM, "# 🔴 真空:需求还没做完,却没有任何会话在执行\n\n" "已持续:%.0f 分钟\n" "该谁做:主会话抢域锁,读摘要与任务图判断缺口,然后派下一棒;或者自己做最靠前的那一步。" "做完写记忆、释锁、停。\n" % V["vdur"]) else: rm(VACUUM) if T.get("waste"): w(READY, "# 🔴 可派未派(浪费):有能干的活,却没有任何棒在跑\n\n" "现在就能派:%s\n" "该谁做:按任务图优先关键路径派活,抢细粒度域锁,禁用整工作区粗域。\n" % ";".join("%s %s(线:%s)" % (r[0], r[1], r[2]) for r in T["ready"])) else: rm(READY) def one_round(st: dict) -> dict: d = fetch() up = probe_port() V = verdicts(d, up, st) T = taskgraph(V["cur"], V["busy"]) H = health(d, V) if not H["skipped"] and time.time() - float(st.get("health_at") or 0) > float(C["health_every"]): st["health_at"] = time.time() try: HEALTH.write_text("# 机制体检(空闲时执行)\n\n- %s\n\n## 结果:%s\n%s\n" % (time.strftime("%Y-%m-%d %H:%M:%S"), "✅ 通过" if not H["issues"] else "⚠️ %d 项" % len(H["issues"]), "\n".join("- %s" % i for i in H["issues"]) or "- 无异常"), encoding="utf-8") except Exception: pass tk_ok, tk = targets() render(d, V, H, T, tk) # 🔧 **兼容别名**:旧引用(`CODEBUDDY.md §1.5 D⑤`/SOP)读 `advance.md`。由**同一份内容**产出, # ⛔ 不是第二真相(同源同内容,仅供旧读者过渡)。 try: (INBOX / "advance.md").write_text( "# 机械推进快照(与 digest.md 同源)\n\n" + (INBOX / "digest.md").read_text(encoding="utf-8"), encoding="utf-8") except Exception: pass signals(V, T) # ⑪ 读队列/写通知 —— 🔴 **纯投影**(`mutate=False`):⛔ 绝不推进队列状态。 # **投递 + 推进队列 = 「监督程序」的唯一职责**(`--tick`,由宿主钩子唤起)。 # ⚠️ 2026-09-30 修:旧版这里 `deliver=False` 但**仍然消费** `notify_pending`/`notified` # ⇒ 钩子先跑 `--once`(节流 3 分钟,谁发话都会跑)**把待反馈项吃掉却不投递**, # 紧接着的 `--tick` 便看到空队列 ⇒ **通知永远发不出去**(="程序在跑、主会话没收到")。 st["queue_info"] = supervise(st, deliver=False, mutate=False) # 🔴 2026-09-29 职责纠正(用户定案:"协作程序怎么会投递呢,应该是只维护协作队列,让监督程序来读"): # **投递(通知主会话)=「监督程序」的职责**,⛔ 不是协作程序的。⇒ 这里**不再调用 `wake_round()`** # (它与 `--supervise` 里的 `_deliver_str` 各判各的重 ⇒ 正是"同一内容成对重复投递"的真身)。 # ⇒ 协作程序只做:读库/判定/信号文件/看板/台账投影。 # ⚠️ 副作用提醒:**只起协作程序、不起监督程序 ⇒ 没人投递**(所以守护必须把两个都拉起来)。 wi = st.get("wake_info") or {} # 观测项保留旧值(投递台账现在只由监督程序写) st.update({"rowid": V["wm"], "aids": V["cur"], "last_progress_at": V["pa"], "vacuum_since": V["vs"], "up": up, "wake_info": wi}) try: STATE.write_text(json.dumps(st, ensure_ascii=False), encoding="utf-8") except Exception: pass return st def main() -> int: # 🔴 未知参数 ⇒ **拒绝并退出**(⛔ 不许落进常驻模式)—— # 实测踩过:误传一个未处理的参数(`--reqs`)会**静默变成常驻进程**,很难发现。 _KNOWN = {"--where", "--once", "--supervise", "--tick", "--ready-next", "--report", "--reqs", "--declare", "--state", "--by", "--artifact", "--reason", "--line", "--role", "--name", "--topic", "--reconcile"} _bad = [a for a in sys.argv[1:] if a.startswith("--") and a not in _KNOWN] if _bad: print("未知参数:%s\n用法:collabd.py [--once | --tick | --supervise | --reqs | " "--report --state pending|running|done|blocked [--by ...] [--artifact ...] [--reason ...] | --where]" "\n⚠️ ⛔ 不带参数 = 常驻模式(应由守护程序拉起)" % " ".join(_bad)) return 2 if "--where" in sys.argv: print("collabd.py =", Path(__file__).resolve()) print("workspace =", WS, "\nLIVE =", LIVE, "\nTG =", TG, "\nINBOX =", INBOX) print("配置来源 =", CFG_USED or "⚠️ 未找到(在用 DEFAULTS)") return 0 # 🔴🔴 **找不到使用方的部署配置 ⇒ 拒跑(fail-closed)** # 为什么必须拒:不配配置时 `workspace` 会回落到 **cwd**,而 cwd 常常就是**技能目录** # ⇒ 会在技能里长出 `tmp/supervise-inbox/`、`_collabd.log` 等**使用方的产物** # (2026-09-30 实测:跑一次 `--where` 就在技能里生成了 `tmp/supervise-inbox/_collabd.log`)。 # ⇒ 宁可**什么都不做并说明原因**,⛔ 也不把技能目录当使用方用。 if CFG_MISSING: print("⛔ 未找到部署配置,**拒跑**(避免把技能目录当工作区用、在里面长出运行产物)。\n" " 查找顺序:① 环境变量 `COLLABD_CONFIG` ② `" "/.workbuddy/collab/collabd.config.json`。\n" " 使用方请把配置放在 ② 那个位置(或由钩子/启动器传 ①)。\n" " 只想看路径 ⇒ 用 `--where`。") return 2 INBOX.mkdir(parents=True, exist_ok=True) st = {} try: st = json.loads(STATE.read_text(encoding="utf-8")) except Exception: pass if "--report" in sys.argv: # 协作会话**上报**状态(执行中/执行完毕) def _arg(k, dv=""): return sys.argv[sys.argv.index(k) + 1] if (k in sys.argv and sys.argv.index(k) + 1 < len(sys.argv)) else dv rc = task_report(_arg("--report"), _arg("--state", "running"), _arg("--by"), _arg("--artifact"), _arg("--reason"), _arg("--line")) if rc == 0: # 上报**只写台账**(⛔ 不投递)—— # 🔴 2026-09-29 职责纠正:投递归监督程序;协作程序上报后**只是把状态落盘**, # 由**常驻的监督程序**在它自己的轮次里读到变化并投递(这样才有"单条+握手"的顺序)。 try: STATE.write_text(json.dumps(st, ensure_ascii=False), encoding="utf-8") except Exception: pass return rc if "--declare" in sys.argv: # 会话**声明自己的角色**(sid 自动从环境取,⛔ 不用手填) def _d(k, dv=""): return sys.argv[sys.argv.index(k) + 1] if (k in sys.argv and sys.argv.index(k) + 1 < len(sys.argv)) else dv _role = _d("--role", "worker").strip().lower() if _role not in ("main", "worker"): print("--role 只能是 main / worker") return 2 _sid = os.environ.get("CODEBUDDY_SESSION_ID") or "" _key = _sid or ("name:" + _d("--name", "?")) st["roles"] = dict(st.get("roles") or {}) st["roles"][_key] = _role try: STATE.write_text(json.dumps(st, ensure_ascii=False), encoding="utf-8") except Exception: pass _topic = _d("--topic", "").strip() if _topic: st["topics"] = dict(st.get("topics") or {}) st["topics"][_key] = _topic print("OK 已声明角色:%s = %s%s%s" % (_key[:8], _role, (",主题 = %s" % _topic) if _topic else "(⚠️ 未给 --topic ⇒ 二级前缀缺失)", "" if _sid else "(⛔ 环境里取不到会话 id ⇒ 用名字作键)")) return 0 if "--ready-next" in sys.argv: # 🔴 收尾确认(替代"盲等 5~8 分钟") _r = ready_next(st) if _r["ok"]: print("✅ 可以排下一棒(上一棒已收尾)⇒ 排期建议 = 现在 + 30~60 秒") else: print("⛔ 先别排,原因:") for w in _r["why"]: print(" - %s" % w) return 0 if "--reconcile" in sys.argv: _r = reconcile(st) print("OK 启动对账:收编 %d 条(%s);僵尸 %d 条(%s)" % (len(_r["adopted"]), "、".join(_r["adopted"]) or "-", len(_r["zombie"]), "、".join(_r["zombie"]) or "-")) if _r.get("skip"): print(" WARN: " + _r["skip"]) return 0 if "--reqs" in sys.argv: # 打印**需求台账**(四态;协作程序持有) _t = _load_tasks() _n = {"pending": "待执行", "running": "执行中", "done": "已完成", "blocked": "有阻碍"} if not _t: print("(需求台账为空 —— 待办可能还在任务图里,尚未并入)") for _k, _v in sorted(_t.items()): _s = str((_v or {}).get("state") or "?") print("- %-16s %-6s %s%s%s" % (_k, _n.get(_s, _s), ("线 %s|" % _v.get("line")) if _v.get("line") else "", ("执行者 %s " % _v.get("by")) if _v.get("by") else "", ("|阻碍:%s" % _v.get("block_reason")) if _v.get("block_reason") else "")) return 0 if "--supervise" in sys.argv: # 监督程序**常驻**模式(由守护程序看护;⛔ 不派活) log("supervise loop start pid=%d" % os.getpid()) while True: if _guard_says_stop(): # 守护停了 ⇒ 优雅退出 log("守护已停 ⇒ 监督程序优雅退出") return 0 try: st["queue_info"] = supervise(st) STATE.write_text(json.dumps(st, ensure_ascii=False), encoding="utf-8") # 🔴 前置探针(设备侧 worker + 垫片):掉了必须被立刻看见,⛔ 不静默(节流 10 min) _fr = check_frontline() if not (_fr.get("shim") and _fr.get("worker")): log("前置探针:shim=%s worker=%s" % (_fr.get("shim"), _fr.get("worker"))) alert_frontline(_fr) except Exception as e: log("supervise loop err %s" % e) time.sleep(float(C.get("supervise_interval") or 30)) if "--tick" in sys.argv: # 监督程序 · **宿主钩子唤起的一次性投递轮** # 🔴 2026-09-30 立(用户口径:「又给我整到自动任务去了」)—— # **投递 ⛔ 不靠自动任务排期、⛔ 不靠常驻进程**:由**宿主钩子**唤起本程序跑一轮。 # 依据(源码 + 本机实测,⛔ 非推断): # · 钩子是**宿主起的子进程** ⇒ ① 继承到网关口令 `CODEBUDDY_GATEWAY_PASSWORD` # (2026-09-30 实测:本机会话内进程 `len=43` 有值)② ⛔ 不占任何会话 ③ 零 token; # · `UserPromptSubmit` 两日 26 次 spawn 实测**会被投递**(见 pitfalls §钩子)。 # · ⇒ **监督程序 = 被钩子唤起的一次性进程**(`--tick`),⛔ 不需要常驻、⛔ 不需要排期。 # **覆盖度(为什么这样就够)**:需要投递的时刻**全都是由某个会话在动产生的** —— # 协作会话上报=一次 Bash 调用、主会话处理完=一次工具调用 ⇒ 那一刻钩子必然响 ⇒ **自洽**。 # ⚠️ 唯一缺口=「谁都没动」的**停滞心跳**(需要时钟)⇒ 由会话外守护**只落标记** # (`STALL.md`),等下一次任意钩子触发时**补投**。⇒ ⛔ 常驻进程不再是投递的前置。 # 🔴🔴 **2026-09-30 用户改口(取代本块里"⛔ 不靠自动任务排期"那句)**: # 「自动任务(符合条件时)是通过 协作程序去唤醒 主会话,这样流程统一」 # ⇒ **允许自动任务**,但它**只准当闹钟**:跑 `goalctl.py wake` ⇒ 调本 `--tick`; # ⛔ 不自己判断、⛔ 不自己派活、⛔ 不自己起会话。**"唤醒"仍只有本程序一个出口**,区别只在"谁拨这一下": # 有会话在动 ⇒ 钩子;**谁都没动 ⇒ 自动任务那一下(这正是上面那个唯一缺口的补丁)**。 # ⇒ 判读三分(`goalctl wake` 已内建):投成/主会话在忙/刚投过 ⇒ **拨钟方结束**; # 没有可唤醒的对象(`main-not-live`/`no-main-session`/`no-live-session`)⇒ **降级**:拨钟方自己就是那个会话。 # ⚠️ 本路径**必须**在宿主树内被调起:否则拿不到口令 ⇒ `_deliver_str` 会落 `NEED-USER.md` # (设计如此,⛔ 不静默失败)。手工在会话外跑 `--tick` 会看到 `no-token`,属预期。 if goal_paused(): # 🔴 已停 ⇒ ⛔ 不投递(停掉的目标不该被叫醒) print("tick: paused (goal.run != active) => no delivery") return 0 if not goals_open(): # 🔴 2026-10-01 加:**目标三路全过 ⇒ 终态,⛔ 不叫醒任何人** print("tick: goal complete (三路全过) => no delivery") return 0 global FROM_HOOK FROM_HOOK = True # 本进程=宿主钩子唤起 ⇒ 没口令就属**异常**(见 `_deliver_str`) try: st["queue_info"] = supervise(st, deliver=True) # ① 回查"上次投递有没有被消费"(投出去 ≠ 它跑起来了) _cc = check_delivery_consumed(st) # ② 宿主侧「消息卡住」指纹探针(⛔ AI 侧修不了,只能发现 + 告诉用户点哪一下;节流 10 分钟) _pk = {"n": 0, "last": ""} try: _pk = probe_host_park() _fresh = False if _pk["n"] and _pk["last"]: try: hh, mm, ss = [int(x) for x in _pk["last"].split(":")] _fresh = ((time.time() - (hh * 3600 + mm * 60 + ss)) % 86400) < PARK_FRESH except Exception: _fresh = False _due = True try: if PARK_STAMP.exists() and (time.time() - PARK_STAMP.stat().st_mtime) < PARK_GAP: _due = False except Exception: pass if _fresh and _due: need_user( "检测到**宿主侧「消息卡住」指纹**(%s,最近一次 %s):`parkInQueue` + `hasWaiter=false`" " ⇒ 消息进了队列、**没有消费者**。\n🔴 **AI 侧无法自救**。请你**把那个会话窗口切走再切回**" "(或关掉重开)—— 队列随即排空。⛔ 别反复发消息试探。" % (_pk["file"], _pk["last"])) PARK_STAMP.write_text(time.strftime("%Y-%m-%d %H:%M:%S"), encoding="utf-8") except Exception as e: log("park 探针失败(已忽略)%s" % e) # ③ 前置探针(垫片/worker 掉了 ⇒ 立刻报警,⛔ 不静默;节流 10 min) try: alert_frontline(check_frontline()) except Exception as e: log("前置探针失败(已忽略)%s" % e) STATE.write_text(json.dumps(st, ensure_ascii=False), encoding="utf-8") _qi = st.get("queue_info") or {} _de = _qi.get("deliver") or {} _dd = _de.get("delivered") or {} print("tick: item=%s awaiting=%s phase=%s deliver=%s consume=%s park=%s" % (_qi.get("sent") or _qi.get("notice") or "-", _qi.get("awaiting") or "-", _qi.get("phase") or "-", ("http=%s" % _dd.get("http")) if _dd else (_de.get("skipped") or "-"), ("%s@%s" % (list(_cc.keys())[0], _cc.get("sid"))) if _cc else "-", ("%d 次 最近 %s(%s)" % (_pk["n"], _pk["last"], _pk["file"])) if _pk["n"] else "-")) except Exception as e: log("tick err %s" % e) return 1 return 0 if "--once" in sys.argv: if goal_paused(): # 🔴 已停 ⇒ 只写一行"已停"(⛔ 不产告警) paused_round(st) print("once: paused (goal.run != active) => only one line, no alerts") return 0 if not goals_open(): # 🔴 2026-10-01 加:**目标三路全过 ⇒ 终态** paused_round(st, why="目标三路全过") print("once: goal complete (三路全过) => only one line, no alerts") return 0 st.pop("paused", None) one_round(st) return 0 try: # 单例 g = socket.socket() g.setsockopt(socket.SOL_SOCKET, socket.SO_REUSEADDR, 0) g.bind(("127.0.0.1", int(C["singleton_port"]))) g.listen(1) except Exception as e: log("singleton bind failed (%s) => another instance runs" % e) return 0 log("collabd start pid=%d ws=%s" % (os.getpid(), WS)) n = 0 while True: if _guard_says_stop(): # 守护停了 ⇒ 优雅退出(⛔ 不硬杀、不成孤儿) log("守护已停 ⇒ 协作程序优雅退出") break n += 1 try: st = one_round(st) # ⛔ **业务自愈已从协作机制移出**(2026-09-29 用户定案:「协作机制就是协作机制, # 手机控制 workbuddy 是另一回事」)⇒ 本程序只做 队列/判定/通知,⛔ 不探业务端口、 # ⛔ 不拉业务客户端。若手机接入线需要自愈 ⇒ 由**该线自己的件**做。 except Exception as e: log("round %d err %s" % (n, e)) if n % 30 == 0: log("alive rounds=%d" % n) time.sleep(float(C["interval"])) if __name__ == "__main__": try: sys.exit(main()) except KeyboardInterrupt: sys.exit(0) except Exception as e: log("fatal %s" % e) sys.exit(1) # 🔴 不再 exit 0:**rc=0 必须是真成功**(今天栽过一次"假绿"—— # 异常被吞、rc=0、而产物根本没刷新,全链路静默失灵)