Files
admin c1b5e4d966 chore(工作区): 全量入库 + 补齐 .gitignore(以工作区为准)
- 变更规模:新增 514 / 修改 62 / 重命名 155 / 删除 4(归档重组与文档轮次)
- .gitignore 修:`归档/**/db-cwd归一-备份-*/` —— 原规则写绝对层级(归档/db-cwd归一-…),
  目录搬进 归档/配置与备份/ 后**静默失效**,43 MB 的 DB 备份又变成未跟踪
- .gitignore 补:嵌套 git 内部数据(归档/内嵌git-20261008/、归档/skills-git-旧线-20261007/dotgit-原样移出/)
- .gitignore 补:运行态与部署副本(.workbuddy/collab/、.workbuddy/tools/、.workbuddy/.load-pending、.workbuddy/tmp-*)
- .gitignore 补:备份件(*.bak-*)
- 未跟踪文件从 2190 降到 890(其余为 归档/ 归档件与 .workbuddy/memory/ 知识文件,按口径入库)
2026-10-10 23:13:22 +08:00

2908 lines
173 KiB
Python
Raw Permalink Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
# -*- coding: utf-8 -*-
"""collabd —— **协作守护程序**(通用版 · 一个进程承担"派活之外的全部功能")
配套:`session-mechanism` 技能。配置见同目录 `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/<id>/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/<id>` 成功者得)
② **同线互斥、跨线并行** —— 队列里标 `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/<id>` ⇒ **建不成就是别人取走了** ⇒ ⛔ 换队首/等待",
"2. 取到后写 `claims/<id>/holder`(内容 `<会话名>@<线>@<完整会话id>`)⇒ 「谁在做」的唯一权威;**第 3 段别省**(持有人失活 ⇒ 程序立即出队;⛔ 不写只能等超时兜底)",
"3. **做完 → 删掉 `claims/<id>`**(出队)⇒ 下一个才可能成为队首(⛔ 不删 ⇒ 该线一直被占)",
"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:<port> LISTENING` 候选(1024<p<65535),最多 limit 个。"""
try:
# 🔴 2026-09-29 修(用户:"一会弹出来一会弹出来的,影响我操作电脑"):
# `netstat` 是**控制台程序** —— 不加 CREATE_NO_WINDOW 就**每次都闪一个黑窗**
# (本函数每轮都被调 ⇒ 10~30 秒闪一次)。⇒ 所有子进程一律**不显窗**。
r = subprocess.run(["netstat", "-ano"], stdout=subprocess.PIPE, stderr=subprocess.DEVNULL,
creationflags=0x08000000, # CREATE_NO_WINDOW
timeout=40, errors="replace")
out = r.stdout or ""
except Exception as e:
log("wake: netstat err %s" % e)
return []
ports = []
for ln in out.splitlines():
if "LISTENING" not in ln:
continue
m = re.search(r"127\.0\.0\.1:(\d+)\b", ln)
if not m:
continue
p = int(m.group(1))
if 1024 < p < 65535 and p not in ports:
ports.append(p)
return ports[:limit]
def discover_gateways() -> 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/session-mechanism/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 <goalId>` 显式声明(⚠️ 尚未接线,留待需要时)
⛔ **不命中 ⇒ 不算本项目**。
🔴 **⛔ 刻意不用 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/<id>` 原子取件,做完删掉 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 <id> --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` ② `<COLLABD_WORKSPACE 或 cwd>"
"/.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、而产物根本没刷新,全链路静默失灵)