- 变更规模:新增 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/ 知识文件,按口径入库)
2908 lines
173 KiB
Python
2908 lines
173 KiB
Python
# -*- 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、而产物根本没刷新,全链路静默失灵)
|