auto-sync: 2026-05-31 23:35:05

This commit is contained in:
cfdaily
2026-05-31 23:35:05 +08:00
parent a2cf49ee99
commit 3e1d4d066b
3 changed files with 181 additions and 9 deletions
+54 -1
View File
@@ -229,6 +229,12 @@ class Dispatcher:
_dispatcher = self
_is_review = action_type == "review"
# v2.8.1 Fix-2: 明确需要回退 current_agent 的 outcome
ROLLBACK_CURRENT_AGENT_OUTCOMES = frozenset({
"crashed", "compact_failed", "process_crash",
"session_stuck", "compact_hanging",
})
def _task_on_complete(aid, outcome):
try:
if _is_review:
@@ -239,7 +245,26 @@ class Dispatcher:
_dispatcher._mark_task_status(_task_db, _task_id, "done")
logger.info("Task %s: review complete (%s), marking done", _task_id, outcome)
else:
logger.warning("Task %s: review agent %s, NOT marking done", _task_id, outcome)
logger.warning("Task %s: review agent %s (%s), NOT marking done", _task_id, aid, outcome)
# v2.8.1 Fix-2: crash 后回退 current_agent,避免 exclude_current 卡死
if outcome in ROLLBACK_CURRENT_AGENT_OUTCOMES:
try:
conn = get_connection(_task_db)
try:
conn.execute(
"UPDATE tasks SET current_agent = "
"(SELECT assignee FROM tasks WHERE id=?) "
"WHERE id=? AND current_agent=?",
(_task_id, _task_id, aid)
)
conn.commit()
finally:
conn.close()
logger.info("Task %s: rolled back current_agent from %s to assignee",
_task_id, aid)
except Exception as e:
logger.warning("Task %s: failed to rollback current_agent: %s",
_task_id, e)
else:
_dispatcher._task_auto_complete(_task_id, _task_db)
except Exception as e:
@@ -759,3 +784,31 @@ class Dispatcher:
conn.close()
except Exception as e:
logger.error("Task %s: mark status error: %s", task_id, e)
@staticmethod
def _check_crash_limit(task_id: str, db_path: pathlib.Path, limit: int = 3,
window_minutes: int = 30) -> bool:
"""v2.8.1 Fix-3c: 检查 task 最近 window_minutes 内的 crash 次数是否超限。
基于 task_attempts 表(持久化),PM2 重启不丢失。
Returns: True = 已超限,应 escalate。
"""
try:
conn = get_connection(db_path)
try:
row = conn.execute(
"SELECT COUNT(*) as cnt FROM task_attempts "
"WHERE task_id=? AND outcome='crashed' "
"AND started_at > datetime('now', ?)",
(task_id, f'-{window_minutes} minutes')
).fetchone()
count = row["cnt"] if row else 0
if count >= limit:
logger.warning("Task %s: crash limit reached (%d/%d in %dm)",
task_id, count, limit, window_minutes)
return True
return False
finally:
conn.close()
except Exception:
return False