fix: lint - 114 errors fixed (unused imports, indentation, blank lines)
CI / test (pull_request) Has been skipped
CI / test (push) Has been skipped
CI / lint (push) Failing after 6s
CI / lint (pull_request) Failing after 6s
CI / notify-on-failure (push) Successful in 1s
CI / notify-on-failure (pull_request) Successful in 3s
CI / test (pull_request) Has been skipped
CI / test (push) Has been skipped
CI / lint (push) Failing after 6s
CI / lint (pull_request) Failing after 6s
CI / notify-on-failure (push) Successful in 1s
CI / notify-on-failure (pull_request) Successful in 3s
This commit is contained in:
@@ -11,8 +11,7 @@ A 类 Skill 由引擎确定性注入全文,不靠 Description 触发。
|
||||
|
||||
import logging
|
||||
import os
|
||||
from pathlib import Path
|
||||
from typing import Any, Dict, List, Optional
|
||||
from typing import Any, List
|
||||
|
||||
logger = logging.getLogger("moziplus-v2.bootstrap")
|
||||
|
||||
|
||||
@@ -73,7 +73,7 @@ class ActiveAgentCounter:
|
||||
cd = seconds if seconds is not None else self._default_cooldown_seconds
|
||||
self._cooldown_until[agent_id] = time.time() + cd
|
||||
logger.info("Cooldown set for %s: %.0fs (until %.0f)",
|
||||
agent_id, cd, self._cooldown_until[agent_id])
|
||||
agent_id, cd, self._cooldown_until[agent_id])
|
||||
|
||||
async def can_acquire(self, agent_id: str, session_id: str = "main") -> bool:
|
||||
"""三层检查:cooldown → global → per agent → per session key"""
|
||||
|
||||
@@ -12,9 +12,9 @@ Dispatcher 负责:
|
||||
from __future__ import annotations
|
||||
|
||||
import json
|
||||
import pathlib
|
||||
import logging
|
||||
import sqlite3
|
||||
from datetime import datetime
|
||||
from enum import Enum
|
||||
from pathlib import Path
|
||||
from typing import Any, Dict, List, Optional
|
||||
@@ -22,7 +22,7 @@ from typing import Any, Dict, List, Optional
|
||||
from src.blackboard.models import Task
|
||||
from src.blackboard.db import get_connection
|
||||
from src.daemon.spawner import AgentBusyError
|
||||
from src.daemon.router import AgentRouter, RouteDecision
|
||||
from src.daemon.router import AgentRouter
|
||||
|
||||
logger = logging.getLogger("moziplus-v2.dispatcher")
|
||||
|
||||
@@ -194,6 +194,7 @@ class Dispatcher:
|
||||
_task_id = task.id
|
||||
_mail_db = db_path
|
||||
_disp = self
|
||||
|
||||
def _mail_on_checks_passed():
|
||||
nonlocal _mail_marked_working
|
||||
if not _disp._mail_auto_working(_task_id, _mail_db):
|
||||
@@ -203,8 +204,8 @@ class Dispatcher:
|
||||
|
||||
# 构建 spawn message
|
||||
message = self._build_spawn_message(task, agent_id, project_config,
|
||||
mode=decision.get("mode", ""),
|
||||
spawn_type=action_type or "executor")
|
||||
mode=decision.get("mode", ""),
|
||||
spawn_type=action_type or "executor")
|
||||
|
||||
# v2.7.2: on_complete 只含业务逻辑,不含 counter.release
|
||||
# counter.release 由 spawn_full_agent 内部的 wrapped_on_complete 保证
|
||||
@@ -269,8 +270,8 @@ class Dispatcher:
|
||||
from src.blackboard.blackboard import Blackboard
|
||||
bb = Blackboard(_task_db)
|
||||
bb.add_comment(_task_id, "daemon",
|
||||
f"@{task_row['assignee']} 审查结论: {verdict_str},请查看详情并决定接受或反驳",
|
||||
comment_type="review")
|
||||
f"@{task_row['assignee']} 审查结论: {verdict_str},请查看详情并决定接受或反驳",
|
||||
comment_type="review")
|
||||
logger.info("Task %s: review verdict=%s, notified assignee=%s",
|
||||
_task_id, verdict_str, task_row["assignee"] if task_row else "?")
|
||||
# 不标 done,保持 review 状态
|
||||
@@ -661,7 +662,7 @@ class Dispatcher:
|
||||
logger.error("Mail %s: failed to revert to pending: %s", task_id, e)
|
||||
|
||||
def _mail_auto_complete(self, task_id: str, agent_id: str,
|
||||
db_path: Path, must_haves: str, outcome=None) -> None:
|
||||
db_path: Path, must_haves: str, outcome=None) -> None:
|
||||
"""Mail 任务:on_complete 后自动标 done/failed(含幻觉门控)"""
|
||||
try:
|
||||
# 解析 performative
|
||||
|
||||
@@ -14,7 +14,7 @@ import logging
|
||||
import re
|
||||
from datetime import datetime
|
||||
from pathlib import Path
|
||||
from typing import Any, Dict, List, Optional, Tuple
|
||||
from typing import Any, Dict, List, Optional
|
||||
|
||||
logger = logging.getLogger("moziplus-v2.experience")
|
||||
|
||||
@@ -68,7 +68,7 @@ class Experience:
|
||||
@classmethod
|
||||
def from_dict(cls, data: Dict[str, Any]) -> Experience:
|
||||
return cls(**{k: v for k, v in data.items() if k != "id"},
|
||||
experience_id=data.get("id"))
|
||||
experience_id=data.get("id"))
|
||||
|
||||
|
||||
class ExperienceStore:
|
||||
@@ -284,7 +284,7 @@ class ExperienceDistiller:
|
||||
all_tags.append(task_type)
|
||||
|
||||
results = self.store.search(tags=all_tags if all_tags else None,
|
||||
query=query, limit=limit)
|
||||
query=query, limit=limit)
|
||||
|
||||
# 按置信度排序
|
||||
results.sort(key=lambda e: e.confidence, reverse=True)
|
||||
|
||||
@@ -4,7 +4,7 @@ from __future__ import annotations
|
||||
|
||||
import logging
|
||||
import re
|
||||
from dataclasses import dataclass, field
|
||||
from dataclasses import dataclass
|
||||
from pathlib import Path
|
||||
from typing import Any, Dict, List, Optional
|
||||
|
||||
|
||||
@@ -9,9 +9,9 @@ from __future__ import annotations
|
||||
import json
|
||||
import logging
|
||||
from pathlib import Path
|
||||
from typing import Any, Dict, Optional
|
||||
from typing import Any, Dict
|
||||
|
||||
from src.blackboard.db import get_connection, init_db
|
||||
from src.blackboard.db import get_connection
|
||||
from src.blackboard.queries import Queries
|
||||
|
||||
logger = logging.getLogger("moziplus-v2.health")
|
||||
@@ -41,7 +41,6 @@ class HealthChecker:
|
||||
{"healthy": bool, "zombie": bool, "stale_ticks": int,
|
||||
"alert_written": bool, "resolved": bool}
|
||||
"""
|
||||
db_key = str(db_path)
|
||||
result: Dict[str, Any] = {
|
||||
"healthy": True,
|
||||
"zombie": False,
|
||||
@@ -58,7 +57,7 @@ class HealthChecker:
|
||||
# 用 event count 变化判断是否有真实变更
|
||||
conn = queries._conn()
|
||||
try:
|
||||
total_events = conn.execute("SELECT COUNT(*) FROM events").fetchone()[0]
|
||||
conn.execute("SELECT COUNT(*) FROM events").fetchone()[0]
|
||||
non_tick_events = conn.execute(
|
||||
"SELECT COUNT(*) FROM events WHERE event_type != 'daemon_tick' "
|
||||
"AND event_type != 'agent_zombie_detected'"
|
||||
|
||||
+2
-3
@@ -15,7 +15,6 @@ from __future__ import annotations
|
||||
import asyncio
|
||||
import json
|
||||
import logging
|
||||
import os
|
||||
from pathlib import Path
|
||||
from typing import Any, Callable, Coroutine, Dict, List, Optional
|
||||
|
||||
@@ -57,7 +56,7 @@ class InboxWatcher:
|
||||
self._running = True
|
||||
self._task = asyncio.create_task(self._loop())
|
||||
logger.info("Inbox watcher started (path=%s, interval=%.1fs)",
|
||||
self.inbox_path, self.watch_interval)
|
||||
self.inbox_path, self.watch_interval)
|
||||
|
||||
async def stop(self) -> None:
|
||||
"""停止监听"""
|
||||
@@ -69,7 +68,7 @@ class InboxWatcher:
|
||||
except asyncio.CancelledError:
|
||||
pass
|
||||
logger.info("Inbox watcher stopped (processed=%d, errors=%d)",
|
||||
self._total_processed, self._total_errors)
|
||||
self._total_processed, self._total_errors)
|
||||
|
||||
@property
|
||||
def is_running(self) -> bool:
|
||||
|
||||
@@ -108,7 +108,7 @@ def notify_mail_failed(db_path: Path, original_mail_id: str,
|
||||
)
|
||||
bb.create_task(notify_task)
|
||||
logger.info("Mail %s: sent failure notification to %s (original_sender=%s, reason=%s, notify_id=%s)",
|
||||
original_mail_id, target_agent, from_agent, reason, notify_id)
|
||||
original_mail_id, target_agent, from_agent, reason, notify_id)
|
||||
|
||||
except Exception as e:
|
||||
logger.warning("notify_mail_failed: failed to send notification for mail %s: %s", original_mail_id, e)
|
||||
|
||||
@@ -8,15 +8,12 @@ from __future__ import annotations
|
||||
|
||||
import json
|
||||
import logging
|
||||
import re
|
||||
from datetime import datetime
|
||||
from enum import Enum
|
||||
from pathlib import Path
|
||||
from typing import Any, Callable, Dict, List, Optional, Tuple
|
||||
from typing import Any, Callable, Dict, List, Optional
|
||||
|
||||
from src.blackboard.models import Task
|
||||
from src.blackboard.operations import Blackboard
|
||||
from src.blackboard.queries import Queries
|
||||
|
||||
logger = logging.getLogger("moziplus-v2.review")
|
||||
|
||||
|
||||
@@ -10,12 +10,11 @@ from __future__ import annotations
|
||||
|
||||
import json
|
||||
import logging
|
||||
import re
|
||||
from dataclasses import dataclass, field
|
||||
from datetime import datetime
|
||||
from enum import Enum
|
||||
from pathlib import Path
|
||||
from typing import Any, Callable, Dict, List, Optional, Tuple
|
||||
from typing import Any, Dict, List, Optional, Tuple
|
||||
|
||||
logger = logging.getLogger("moziplus-v2.skill")
|
||||
|
||||
|
||||
+13
-13
@@ -7,6 +7,7 @@ Subagent: 占位(实际通过 OpenClaw Gateway API sessions_spawn,F17 完善)
|
||||
from __future__ import annotations
|
||||
|
||||
import asyncio
|
||||
import pathlib
|
||||
import json
|
||||
import logging
|
||||
import os
|
||||
@@ -15,7 +16,7 @@ from datetime import datetime
|
||||
from pathlib import Path
|
||||
from typing import Any, Dict, List, Optional
|
||||
|
||||
from src.blackboard.db import get_connection, init_db
|
||||
from src.blackboard.db import get_connection
|
||||
|
||||
logger = logging.getLogger("moziplus-v2.spawner")
|
||||
|
||||
@@ -299,7 +300,7 @@ class AgentSpawner:
|
||||
project_id, agent_id)
|
||||
|
||||
def _build_minimal_fallback(self, task_id, title, description, must_haves,
|
||||
project_id, agent_id):
|
||||
project_id, agent_id):
|
||||
"""最小 fallback:只有任务上下文 + API 指令"""
|
||||
task_section = f"""## 任务
|
||||
{title}
|
||||
@@ -311,7 +312,7 @@ class AgentSpawner:
|
||||
return task_section + "\n\n---\n\n" + api_section
|
||||
|
||||
def _build_api_section(self, project_id: str, task_id: str,
|
||||
agent_id: str) -> str:
|
||||
agent_id: str) -> str:
|
||||
"""构建 API 回写操作指令(BootstrapBuilder 模式下补充)"""
|
||||
# mail 任务直接 done,不走 review
|
||||
success_status = '"done"' if project_id == "_mail" else '"review"'
|
||||
@@ -337,8 +338,8 @@ curl -X POST http://{self.api_host}:{self.api_port}/api/projects/{project_id}/ta
|
||||
"""
|
||||
|
||||
def _build_discussion_prompt(self, task_id: str, title: str,
|
||||
description: str, must_haves: str,
|
||||
project_id: str, agent_id: str) -> str:
|
||||
description: str, must_haves: str,
|
||||
project_id: str, agent_id: str) -> str:
|
||||
"""构建讨论类 spawn prompt(§3.3 框架 + Boids)"""
|
||||
goal_snapshot = description or title
|
||||
constraints = must_haves or "(无特殊约束)"
|
||||
@@ -379,9 +380,8 @@ curl -X POST http://{self.api_host}:{self.api_port}/api/projects/{project_id}/ta
|
||||
return router.agent_profiles.get(agent_id)
|
||||
return None
|
||||
|
||||
|
||||
def _build_mail_prompt(self, task_id: str, title: str, description: str,
|
||||
must_haves: str, agent_id: str) -> str:
|
||||
must_haves: str, agent_id: str) -> str:
|
||||
"""构建 Mail 专用精简模板"""
|
||||
# 解析 must_haves 获取 from 和 performative
|
||||
from_agent = agent_id
|
||||
@@ -575,7 +575,7 @@ curl -X POST http://{self.api_host}:{self.api_port}/api/projects/{project_id}/ta
|
||||
stderr=asyncio.subprocess.PIPE,
|
||||
)
|
||||
self._register_session(session_id, agent_id, task_id, proc.pid,
|
||||
broadcast_task_ids=broadcast_task_ids)
|
||||
broadcast_task_ids=broadcast_task_ids)
|
||||
logger.info("Spawned agent %s (session=%s, pid=%d)",
|
||||
agent_id, session_id, proc.pid)
|
||||
|
||||
@@ -852,7 +852,7 @@ curl -X POST http://{api_host}:{api_port}/api/projects/{project_id}/tasks/{task_
|
||||
self.counter.set_cooldown(agent_id, seconds=60)
|
||||
logger.info("Crash cooldown set for %s: 60s (outcome=%s)", agent_id, outcome)
|
||||
elif outcome in ("compact_failed", "process_crash", "session_stuck",
|
||||
"compact_hanging", "agent_error", "compact_interrupted") and self.counter:
|
||||
"compact_hanging", "agent_error", "compact_interrupted") and self.counter:
|
||||
self.counter.set_cooldown(agent_id, seconds=300) # 5 分钟
|
||||
logger.info("Error cooldown set for %s: 300s (outcome=%s)", agent_id, outcome)
|
||||
# F1: 不可恢复 outcome → 立刻标 failed + 写黑板
|
||||
@@ -881,7 +881,7 @@ curl -X POST http://{api_host}:{api_port}/api/projects/{project_id}/tasks/{task_
|
||||
except Exception:
|
||||
pass
|
||||
|
||||
stderr_text = b"".join(stderr_chunks).decode("utf-8", errors="replace")
|
||||
_ = b"".join(stderr_chunks).decode("utf-8", errors="replace")
|
||||
|
||||
# 检查 session 状态
|
||||
state = self._check_session_state(agent_id)
|
||||
@@ -1428,7 +1428,7 @@ curl -X POST http://{api_host}:{api_port}/api/projects/{project_id}/tasks/{task_
|
||||
return defaults
|
||||
|
||||
def _update_retry_counts(self, db_path: Optional[Path],
|
||||
task_id: Optional[str], counts: dict):
|
||||
task_id: Optional[str], counts: dict):
|
||||
"""将 retry counts 写回最新 task_attempt 的 metadata"""
|
||||
if not db_path or not task_id:
|
||||
return
|
||||
@@ -1488,8 +1488,8 @@ curl -X POST http://{api_host}:{api_port}/api/projects/{project_id}/tasks/{task_
|
||||
from src.blackboard.operations import Blackboard
|
||||
bb = Blackboard(db_path)
|
||||
cid = bb.add_comment(task_id, "daemon",
|
||||
f"@pangtong-fujunshi 任务执行失败: {reason},请评估是否需要介入",
|
||||
comment_type="system")
|
||||
f"@pangtong-fujunshi 任务执行失败: {reason},请评估是否需要介入",
|
||||
comment_type="system")
|
||||
bb.record_mentions(cid, task_id, ["pangtong-fujunshi"])
|
||||
logger.info("Task %s: failure notified pangtong via comment+mention (reason=%s)", task_id, reason)
|
||||
except Exception as e:
|
||||
|
||||
+1
-4
@@ -9,14 +9,11 @@ from __future__ import annotations
|
||||
import asyncio
|
||||
import json
|
||||
import logging
|
||||
import subprocess
|
||||
import uuid
|
||||
from datetime import datetime
|
||||
from enum import Enum
|
||||
from pathlib import Path
|
||||
from typing import Any, Callable, Dict, List, Optional, Set
|
||||
from typing import Any, Dict, List, Optional
|
||||
|
||||
from src.blackboard.models import Event
|
||||
|
||||
logger = logging.getLogger("moziplus-v2.sse")
|
||||
|
||||
|
||||
+23
-24
@@ -21,7 +21,6 @@ from dataclasses import dataclass, field as dc_field
|
||||
|
||||
from src.blackboard.operations import Blackboard
|
||||
from src.blackboard.db import get_connection
|
||||
from src.blackboard.models import Task
|
||||
from src.daemon.spawner import AgentBusyError
|
||||
from src.blackboard.queries import Queries
|
||||
from src.blackboard.registry import ProjectRegistry
|
||||
@@ -35,6 +34,7 @@ class BroadcastRound:
|
||||
responded_agents: set = dc_field(default_factory=set) # 已返回反馈的 Agent(含 NO_REPLY)
|
||||
round_number: int = 0 # 当前第几轮(0=未开始,1=第1轮)
|
||||
|
||||
|
||||
logger = logging.getLogger("moziplus-v2.ticker")
|
||||
|
||||
|
||||
@@ -391,7 +391,7 @@ class Ticker:
|
||||
MAX_ROUNDS = 5 # §4.5 防无限循环
|
||||
|
||||
async def _check_round_complete(self, db_path: Path,
|
||||
project_id: str) -> List[str]:
|
||||
project_id: str) -> List[str]:
|
||||
"""检测 parent task 下所有 sub task 终态 → spawn 庞统 review
|
||||
|
||||
流程(§4.4):
|
||||
@@ -462,7 +462,7 @@ class Ticker:
|
||||
"Round %d review spawned for parent %s (subs: %s)",
|
||||
new_round, parent_id, summary
|
||||
)
|
||||
except Exception as e:
|
||||
except Exception:
|
||||
logger.exception("Round check error for parent %s", parent_id)
|
||||
|
||||
return reviewed
|
||||
@@ -531,9 +531,9 @@ Parent Task ID: {parent_task.id}
|
||||
"""
|
||||
|
||||
async def _spawn_pangtong_review(self, parent_task,
|
||||
review_prompt: str,
|
||||
project_id: str,
|
||||
new_round: int = 0) -> bool:
|
||||
review_prompt: str,
|
||||
project_id: str,
|
||||
new_round: int = 0) -> bool:
|
||||
"""Spawn 庞统进行 review
|
||||
|
||||
流程:
|
||||
@@ -543,7 +543,6 @@ Parent Task ID: {parent_task.id}
|
||||
"""
|
||||
try:
|
||||
agent_id = "pangtong-fujunshi"
|
||||
session_id = f"review-{parent_task.id}-r{new_round}"
|
||||
|
||||
# 构造 on_complete 回调:解析庞统结论,更新 parent 状态
|
||||
async def _on_review_complete(aid: str, outcome: str):
|
||||
@@ -586,7 +585,7 @@ Parent Task ID: {parent_task.id}
|
||||
self._set_parent_reviewing(parent_task.id, project_id)
|
||||
return True
|
||||
return False
|
||||
except Exception as e:
|
||||
except Exception:
|
||||
logger.exception("Failed to spawn pangtong review for %s", parent_task.id)
|
||||
return False
|
||||
|
||||
@@ -603,14 +602,14 @@ Parent Task ID: {parent_task.id}
|
||||
(parent_id,))
|
||||
conn.commit()
|
||||
logger.info("Parent %s → reviewing (round review in progress)",
|
||||
parent_id)
|
||||
parent_id)
|
||||
finally:
|
||||
conn.close()
|
||||
except Exception:
|
||||
logger.exception("Failed to set parent %s to reviewing", parent_id)
|
||||
|
||||
def _handle_review_conclusion(self, parent_id: str, project_id: str,
|
||||
review_text: str, round_num: int):
|
||||
review_text: str, round_num: int):
|
||||
"""解析庞统 review 结论,更新 parent 状态
|
||||
|
||||
review_text 是庞统回复的文本(从 spawner session meta payloads 拼接)。
|
||||
@@ -675,7 +674,7 @@ Parent Task ID: {parent_task.id}
|
||||
MENTION_MAX_RETRIES = 5
|
||||
|
||||
async def _process_mentions(self, db_path: Path,
|
||||
project_id: str) -> List[str]:
|
||||
project_id: str) -> List[str]:
|
||||
"""扫描 pending mentions → spawn 被 @ 的 Agent
|
||||
|
||||
流程(§3.4):
|
||||
@@ -767,8 +766,8 @@ Parent Task ID: {parent_task.id}
|
||||
from src.blackboard.blackboard import Blackboard
|
||||
bb2 = Blackboard(rdb_path)
|
||||
bb2.add_comment(_t_id, "daemon",
|
||||
f"@{t_row['assignee']} 审查结论: {verdict_str},请查看详情并决定接受或反驳",
|
||||
comment_type="review")
|
||||
f"@{t_row['assignee']} 审查结论: {verdict_str},请查看详情并决定接受或反驳",
|
||||
comment_type="review")
|
||||
logger.info("Rebuttal: task %s still %s after rebuttal", _t_id, verdict_str)
|
||||
except Exception:
|
||||
logger.exception("Rebuttal on_complete failed for task %s", _t_id)
|
||||
@@ -805,7 +804,7 @@ Parent Task ID: {parent_task.id}
|
||||
# Agent 忙,不递增 retry_count,等下次 tick 自然重试
|
||||
logger.info("Mention spawn skipped: %s busy, will retry next tick", agent_id)
|
||||
|
||||
except Exception as e:
|
||||
except Exception:
|
||||
logger.exception("Mention processing error for agent %s", agent_id)
|
||||
for item in items:
|
||||
try:
|
||||
@@ -948,7 +947,7 @@ Parent Task ID: {parent_task.id}
|
||||
# ------------------------------------------------------------------
|
||||
|
||||
async def _dispatch_pending(self, db_path: Path,
|
||||
project_id: str) -> List[str]:
|
||||
project_id: str) -> List[str]:
|
||||
"""扫描 pending 任务并调度
|
||||
|
||||
v3.0: 两条路径
|
||||
@@ -1242,7 +1241,7 @@ Parent Task ID: {parent_task.id}
|
||||
return [aid for aid in all_agents if active.get(aid, 0) == 0]
|
||||
|
||||
async def _dispatch_reviews(self, db_path: Path,
|
||||
project_id: str) -> List[str]:
|
||||
project_id: str) -> List[str]:
|
||||
"""扫描 review 状态任务,检查是否有产出,调度审查 Agent"""
|
||||
# mail 任务不走 review 流程,直接跳过
|
||||
if project_id == "_mail":
|
||||
@@ -1344,7 +1343,7 @@ Parent Task ID: {parent_task.id}
|
||||
)
|
||||
reclaimed.append(task.id)
|
||||
logger.warning("Escalated %s: no taker after %d broadcasts",
|
||||
task.id, retry_count)
|
||||
task.id, retry_count)
|
||||
finally:
|
||||
conn.close()
|
||||
else:
|
||||
@@ -1423,7 +1422,7 @@ Parent Task ID: {parent_task.id}
|
||||
if ok:
|
||||
reclaimed.append(task.id)
|
||||
logger.info("Mail %s: ticker recheck found reply, marked done (%.1fm)",
|
||||
task.id, elapsed)
|
||||
task.id, elapsed)
|
||||
finally:
|
||||
conn.close()
|
||||
continue
|
||||
@@ -1440,7 +1439,7 @@ Parent Task ID: {parent_task.id}
|
||||
if ok:
|
||||
reclaimed.append(task.id)
|
||||
logger.warning("Task %s timed out (working %.1fm > %.1fm)",
|
||||
task.id, elapsed, timeout_minutes)
|
||||
task.id, elapsed, timeout_minutes)
|
||||
finally:
|
||||
conn.close()
|
||||
except (ValueError, TypeError):
|
||||
@@ -1501,7 +1500,7 @@ Parent Task ID: {parent_task.id}
|
||||
return True # 保守:查询失败假设有回复
|
||||
|
||||
def _check_recent_routing(self, db_path: Path, task_id: str,
|
||||
action_type: str) -> bool:
|
||||
action_type: str) -> bool:
|
||||
"""检查最近 5 分钟内是否已 dispatch 过指定类型的路由(防重复)"""
|
||||
try:
|
||||
conn = get_connection(db_path)
|
||||
@@ -1579,11 +1578,11 @@ Parent Task ID: {parent_task.id}
|
||||
|
||||
if recovery_report["total_recovered"] > 0:
|
||||
logger.info("Startup recovery: %d tasks recovered across %d projects",
|
||||
recovery_report["total_recovered"],
|
||||
len(recovery_report["projects"]))
|
||||
recovery_report["total_recovered"],
|
||||
len(recovery_report["projects"]))
|
||||
elif recovery_report["total_noop"] > 0:
|
||||
logger.info("Startup recovery: %d tasks kept as-is (no recovery needed)",
|
||||
recovery_report["total_noop"])
|
||||
recovery_report["total_noop"])
|
||||
else:
|
||||
logger.info("Startup recovery: no non-terminal tasks found, clean start")
|
||||
|
||||
@@ -1629,7 +1628,7 @@ Parent Task ID: {parent_task.id}
|
||||
return recovered, noop_count
|
||||
|
||||
def _determine_recovery_action(self, conn, task, status: str,
|
||||
db_path: Path) -> Optional[str]:
|
||||
db_path: Path) -> Optional[str]:
|
||||
"""根据黑板线索决定恢复动作,返回 None 表示不需要干预"""
|
||||
task_id = task["id"]
|
||||
|
||||
|
||||
Reference in New Issue
Block a user