feat: add executor routing for research agents + research CLI subcommand
在 _make_agent() 中添加 rtp-researcher / net-researcher / kernel-analyzer 三条路由, 新增 `python scripts/run.py research` 子命令可一键将三项基础技术调研任务入队。 Co-Authored-By: Claude Sonnet 4.6 <noreply@anthropic.com>
This commit is contained in:
@@ -29,6 +29,7 @@ def main() -> None:
|
||||
p_sf.add_argument("value", help="事实值")
|
||||
sub.add_parser("projects", help="查看自动发现的所有项目(含类型/模式)")
|
||||
sub.add_parser("collect-env", help="收集本机+SSH设备环境信息(模型路径/磁盘/包)")
|
||||
sub.add_parser("research", help="触发三项基础技术调研任务入队(RTP/网络/内核)")
|
||||
args = parser.parse_args()
|
||||
|
||||
if args.cmd == "scan":
|
||||
@@ -194,6 +195,34 @@ def main() -> None:
|
||||
ts = facts.get("collected_at", "")
|
||||
print(f"\n已保存到 ProjectMemory (__env__) [{ts}]")
|
||||
|
||||
elif args.cmd == "research":
|
||||
from rockchip_agents.config import load_config
|
||||
from rockchip_agents.core.queue import TaskQueue, Task
|
||||
cfg = load_config()
|
||||
q = TaskQueue()
|
||||
research_tasks = [
|
||||
Task(project="research", type="research",
|
||||
title="RTP 实时传输优化调研",
|
||||
context="调研 RK3588 平台 RTP 协议栈延迟瓶颈,输出 [调研报告] facts 和 [优化任务]",
|
||||
priority=3, mode="report", agent_role="rtp-researcher", initiator="user"),
|
||||
Task(project="research", type="research",
|
||||
title="网络栈优化调研",
|
||||
context="调研 RK3588 TCP/UDP 调优参数,输出 [调研报告] facts 和 [优化任务]",
|
||||
priority=3, mode="report", agent_role="net-researcher", initiator="user"),
|
||||
Task(project="research", type="research",
|
||||
title="内核资源调度优化调研",
|
||||
context="分析 RK3588 多核调度与 NPU 资源竞争,输出 [调研报告] facts 和 [优化任务]",
|
||||
priority=3, mode="report", agent_role="kernel-analyzer", initiator="user"),
|
||||
]
|
||||
for t in research_tasks:
|
||||
tid = q.enqueue(t)
|
||||
if tid > 0:
|
||||
print(f" 入队: [{t.agent_role}] {t.title} → #{tid}")
|
||||
else:
|
||||
print(f" 跳过(已在队列): {t.title}")
|
||||
print(f"\n已入队 {len(research_tasks)} 个基础技术调研任务")
|
||||
print("执行调研:python scripts/run.py run-next(需在普通终端,非 Claude Code 内)")
|
||||
|
||||
else:
|
||||
parser.print_help()
|
||||
|
||||
|
||||
@@ -0,0 +1,286 @@
|
||||
from __future__ import annotations
|
||||
|
||||
import asyncio
|
||||
import json
|
||||
import logging
|
||||
import re
|
||||
import urllib.request
|
||||
from typing import Any, Optional
|
||||
|
||||
from rockchip_agents.agents.developer import DeveloperAgent
|
||||
from rockchip_agents.config import AgentsConfig
|
||||
from rockchip_agents.core.queue import TaskQueue, Task
|
||||
from rockchip_agents.tools.feishu import FeishuNotifier
|
||||
from rockchip_agents.tools.project_memory import ProjectMemory
|
||||
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
_TAG_PATTERN = re.compile(
|
||||
r"\[(?P<type>SPAWN|AWAIT|APPROVE|REJECT|VETO|PROPOSAL|SELECT)"
|
||||
r"(?::(?P<role>[a-zA-Z0-9_-]+))?\]"
|
||||
r"(?P<desc>[^\[]*)",
|
||||
re.MULTILINE,
|
||||
)
|
||||
|
||||
|
||||
def _parse_agent_tags(text: str) -> list[dict]:
|
||||
"""从 agent 输出中解析结构化控制标签"""
|
||||
results = []
|
||||
for m in _TAG_PATTERN.finditer(text):
|
||||
tag_type = m.group("type")
|
||||
role = (m.group("role") or "").strip()
|
||||
desc = (m.group("desc") or "").strip()
|
||||
entry: dict = {"type": tag_type, "role": role}
|
||||
if tag_type in ("VETO", "REJECT"):
|
||||
entry["reason"] = desc
|
||||
elif tag_type in ("SPAWN", "PROPOSAL", "AWAIT"):
|
||||
entry["desc"] = desc
|
||||
results.append(entry)
|
||||
return results
|
||||
|
||||
|
||||
def _make_agent(config: AgentsConfig, task: Task, memory: ProjectMemory) -> Any:
|
||||
"""根据 agent_role 路由到对应 Agent,传入共享 memory 实例。"""
|
||||
role = task.agent_role
|
||||
if role == "tester":
|
||||
from rockchip_agents.agents.tester import TesterAgent
|
||||
return TesterAgent(config)
|
||||
if role == "productizer":
|
||||
from rockchip_agents.agents.productizer import ProductizerAgent
|
||||
return ProductizerAgent(config, memory=memory)
|
||||
if role == "architect":
|
||||
from rockchip_agents.agents.architect import ArchitectAgent
|
||||
return ArchitectAgent(config, memory=memory)
|
||||
if role == "planner":
|
||||
from rockchip_agents.agents.planner import PlannerAgent
|
||||
from rockchip_agents.core.queue import TaskQueue as _TQ
|
||||
return PlannerAgent(config, queue=_TQ(), memory=memory)
|
||||
if role == "dev-kernel":
|
||||
from rockchip_agents.agents.kernel_dev import KernelDevAgent
|
||||
return KernelDevAgent(config, memory=memory)
|
||||
if role == "dev-lowlevel":
|
||||
from rockchip_agents.agents.lowlevel_dev import LowlevelDevAgent
|
||||
return LowlevelDevAgent(config, memory=memory)
|
||||
if role == "hw-engineer":
|
||||
from rockchip_agents.agents.hw_engineer import HwEngineerAgent
|
||||
return HwEngineerAgent(config, memory=memory)
|
||||
if role == "os-engineer":
|
||||
from rockchip_agents.agents.os_engineer import OsEngineerAgent
|
||||
return OsEngineerAgent(config, memory=memory)
|
||||
if role == "algo-researcher":
|
||||
from rockchip_agents.agents.algo_researcher import AlgoResearcherAgent
|
||||
return AlgoResearcherAgent(config, memory=memory)
|
||||
if role == "algo-antishake":
|
||||
from rockchip_agents.agents.algo_antishake import AlgoAntishakeAgent
|
||||
return AlgoAntishakeAgent(config, memory=memory)
|
||||
if role == "algo-position":
|
||||
from rockchip_agents.agents.algo_position import AlgoPositionAgent
|
||||
return AlgoPositionAgent(config, memory=memory)
|
||||
if role == "algo-nav":
|
||||
from rockchip_agents.agents.algo_nav import AlgoNavAgent
|
||||
return AlgoNavAgent(config, memory=memory)
|
||||
if role == "base-validator":
|
||||
from rockchip_agents.agents.base_validator import BaseValidatorAgent
|
||||
return BaseValidatorAgent(config, memory=memory)
|
||||
if role == "base-architect":
|
||||
from rockchip_agents.agents.base_architect import BaseArchitectAgent
|
||||
from rockchip_agents.core.queue import TaskQueue as _TQ2
|
||||
return BaseArchitectAgent(config, queue=_TQ2(), memory=memory)
|
||||
if role == "system-tester":
|
||||
from rockchip_agents.agents.system_tester import SystemTesterAgent
|
||||
return SystemTesterAgent(config, memory=memory)
|
||||
if role == "vision-analyst":
|
||||
from rockchip_agents.agents.vision_analyst import VisionAnalystAgent
|
||||
from rockchip_agents.core.queue import TaskQueue as _TQ3
|
||||
return VisionAnalystAgent(config, queue=_TQ3(), memory=memory)
|
||||
if role == "market-pm":
|
||||
from rockchip_agents.agents.market_pm import MarketPmAgent
|
||||
return MarketPmAgent(config, memory=memory)
|
||||
if role == "media-producer":
|
||||
from rockchip_agents.agents.media_producer import MediaProducerAgent
|
||||
return MediaProducerAgent(config, memory=memory)
|
||||
if role == "rtp-researcher":
|
||||
from rockchip_agents.agents.rtp_researcher import RtpResearcherAgent
|
||||
return RtpResearcherAgent(config, memory=memory)
|
||||
if role == "net-researcher":
|
||||
from rockchip_agents.agents.net_researcher import NetResearcherAgent
|
||||
return NetResearcherAgent(config, memory=memory)
|
||||
if role == "kernel-analyzer":
|
||||
from rockchip_agents.agents.kernel_analyzer import KernelAnalyzerAgent
|
||||
return KernelAnalyzerAgent(config, memory=memory)
|
||||
return DeveloperAgent(config, memory=memory)
|
||||
|
||||
|
||||
class Executor:
|
||||
def __init__(self, config: AgentsConfig,
|
||||
queue: Optional[TaskQueue] = None,
|
||||
memory: Optional[ProjectMemory] = None) -> None:
|
||||
self._config = config
|
||||
self._queue = queue or TaskQueue()
|
||||
self._memory = memory or ProjectMemory()
|
||||
self._feishu = FeishuNotifier(
|
||||
webhook_url=config.feishu.webhook_url,
|
||||
bot_token=config.feishu.bot_token,
|
||||
)
|
||||
self._semaphore = asyncio.Semaphore(config.scheduler.max_concurrent)
|
||||
|
||||
async def run_next(self) -> bool:
|
||||
"""取出并执行下一个任务,返回是否有任务执行。"""
|
||||
task = self._queue.dequeue()
|
||||
if not task:
|
||||
logger.debug("任务队列为空")
|
||||
return False
|
||||
async with self._semaphore:
|
||||
await self._run_task(task)
|
||||
return True
|
||||
|
||||
async def run_all_pending(
|
||||
self,
|
||||
skip_projects: frozenset[str] | None = None,
|
||||
role_filter: frozenset[str] | None = None,
|
||||
) -> int:
|
||||
"""并发执行所有 pending 任务(受 semaphore 限流),返回执行数量。
|
||||
skip_projects: 跳过这些项目的任务(用于项目级暂停)。
|
||||
role_filter: 仅执行指定 agent_role 的任务(用于多队列隔离)。
|
||||
"""
|
||||
tasks_to_run: list[Task] = []
|
||||
while True:
|
||||
task = self._queue.dequeue(
|
||||
exclude_projects=skip_projects,
|
||||
role_filter=role_filter,
|
||||
)
|
||||
if not task:
|
||||
break
|
||||
tasks_to_run.append(task)
|
||||
|
||||
async def _run_with_sem(t: Task) -> None:
|
||||
async with self._semaphore:
|
||||
await self._run_task(t)
|
||||
|
||||
if tasks_to_run:
|
||||
await asyncio.gather(*[_run_with_sem(t) for t in tasks_to_run],
|
||||
return_exceptions=True)
|
||||
return len(tasks_to_run)
|
||||
|
||||
def _notify_wx(self, path: str, data: dict) -> None:
|
||||
"""非阻塞通知 claude-wx 服务(失败静默降级)"""
|
||||
if not self._config.claude_wx_url:
|
||||
return
|
||||
try:
|
||||
url = self._config.claude_wx_url.rstrip("/") + path
|
||||
body = json.dumps(data).encode()
|
||||
req = urllib.request.Request(
|
||||
url, data=body,
|
||||
headers={"Content-Type": "application/json"},
|
||||
)
|
||||
urllib.request.urlopen(req, timeout=3)
|
||||
except Exception as e:
|
||||
logger.debug("claude-wx 通知失败(非关键): %s", e)
|
||||
|
||||
async def _run_task(self, task: Task) -> None:
|
||||
logger.info("开始执行任务 [%s] %s (role=%s, mode=%s)",
|
||||
task.project, task.title, task.agent_role, task.mode)
|
||||
try:
|
||||
agent = _make_agent(self._config, task, self._memory)
|
||||
self._notify_wx("/agent/task-start", {
|
||||
"task_id": task.id,
|
||||
"project": task.project,
|
||||
"title": task.title,
|
||||
"mode": task.mode,
|
||||
})
|
||||
result = await asyncio.get_running_loop().run_in_executor(
|
||||
None, agent.run, task
|
||||
)
|
||||
# 解析结构化输出标签
|
||||
tags = _parse_agent_tags(result.summary)
|
||||
handled = await self._handle_tags(task, tags)
|
||||
if handled:
|
||||
return # 标签处理接管状态,不走普通 mark_done
|
||||
self._queue.mark_done(task.id, result.summary)
|
||||
self._notify_wx("/agent/task-end", {
|
||||
"task_id": task.id,
|
||||
"project": task.project,
|
||||
"status": result.status,
|
||||
"summary": result.summary[:200],
|
||||
})
|
||||
self._feishu.send_task_result(
|
||||
project=task.project, task_title=task.title,
|
||||
status=result.status, summary=result.summary,
|
||||
)
|
||||
logger.info("任务完成 [%s] %s", task.project, task.title)
|
||||
try:
|
||||
self._memory.add_history(
|
||||
project=task.project,
|
||||
task_title=task.title,
|
||||
task_type=task.type,
|
||||
status=result.status,
|
||||
summary=result.summary[:300],
|
||||
)
|
||||
except Exception as hist_err:
|
||||
logger.warning("写入任务历史失败: %s", hist_err)
|
||||
except Exception as e:
|
||||
logger.error("任务失败 [%s] %s: %s", task.project, task.title, e)
|
||||
self._queue.mark_failed(task.id, str(e))
|
||||
self._notify_wx("/agent/task-end", {
|
||||
"task_id": task.id,
|
||||
"project": task.project,
|
||||
"status": "failed",
|
||||
"summary": str(e)[:200],
|
||||
})
|
||||
try:
|
||||
self._memory.add_history(
|
||||
project=task.project,
|
||||
task_title=task.title,
|
||||
task_type=task.type,
|
||||
status="failed",
|
||||
summary=str(e)[:300],
|
||||
)
|
||||
except Exception as hist_err:
|
||||
logger.warning("写入任务历史失败: %s", hist_err)
|
||||
|
||||
async def _handle_tags(self, task: Task, tags: list[dict]) -> bool:
|
||||
"""处理结构化输出标签,返回 True 表示状态已由标签接管"""
|
||||
if not tags:
|
||||
return False
|
||||
handled = False
|
||||
for tag in tags:
|
||||
t = tag["type"]
|
||||
if t == "SPAWN":
|
||||
child = Task(
|
||||
project=task.project,
|
||||
type="code_review",
|
||||
title=f"[审查] {task.title}",
|
||||
context=tag.get("desc", ""),
|
||||
priority=task.priority,
|
||||
mode=task.mode,
|
||||
agent_role=tag["role"],
|
||||
parent_task_id=task.id,
|
||||
)
|
||||
self._queue.enqueue(child)
|
||||
logger.info("SPAWN 子任务: role=%s, title=%s", tag["role"], child.title)
|
||||
elif t == "AWAIT":
|
||||
self._queue.mark_waiting_approval(task.id, awaiting_role=tag["role"])
|
||||
review = Task(
|
||||
project=task.project,
|
||||
type="code_review",
|
||||
title=f"[审批] {task.title}",
|
||||
context=f"请审批任务 #{task.id}: {task.title}",
|
||||
priority=1,
|
||||
mode=task.mode,
|
||||
agent_role=tag["role"],
|
||||
parent_task_id=task.id,
|
||||
)
|
||||
self._queue.enqueue(review)
|
||||
logger.info("AWAIT: 任务 #%d 进入 waiting_approval, 等待 %s", task.id, tag["role"])
|
||||
handled = True
|
||||
elif t == "VETO" and task.parent_task_id:
|
||||
self._queue.mark_done(task.id, tag.get("reason", ""))
|
||||
self._queue.mark_vetoed(task.parent_task_id, reason=tag.get("reason", ""))
|
||||
logger.info("VETO: 父任务 #%d 被否决,重入队", task.parent_task_id)
|
||||
handled = True
|
||||
elif t == "APPROVE" and task.parent_task_id:
|
||||
self._queue.mark_done(task.id, "approved")
|
||||
self._queue.approve_waiting(task.parent_task_id)
|
||||
logger.info("APPROVE: 父任务 #%d 审批通过", task.parent_task_id)
|
||||
handled = True
|
||||
return handled
|
||||
@@ -81,3 +81,49 @@ def test_kernel_analyzer_extract_research_facts(tmp_path):
|
||||
agent._extract_research_facts(report)
|
||||
facts = mem.get_facts("__research__")
|
||||
assert facts.get("kernel.summary") == "CFS 调度延迟在 NPU 高负载时达 20ms"
|
||||
|
||||
|
||||
def test_executor_routes_rtp_researcher(tmp_path):
|
||||
"""executor._make_agent 应将 rtp-researcher 路由到 RtpResearcherAgent。"""
|
||||
from rockchip_agents.core.executor import _make_agent
|
||||
from rockchip_agents.core.queue import Task
|
||||
from rockchip_agents.tools.project_memory import ProjectMemory
|
||||
from rockchip_agents.agents.rtp_researcher import RtpResearcherAgent
|
||||
from rockchip_agents.config import AgentsConfig, ClaudeConfig, SchedulerConfig, FeishuConfig
|
||||
|
||||
cfg = AgentsConfig(projects={}, scheduler=SchedulerConfig(),
|
||||
claude=ClaudeConfig(api_key="fake"), feishu=FeishuConfig(), devices={})
|
||||
task = Task(project="research", type="research", title="RTP调研",
|
||||
context="", priority=3, mode="report", agent_role="rtp-researcher")
|
||||
agent = _make_agent(cfg, task, ProjectMemory(db_path=tmp_path / "m.db"))
|
||||
assert isinstance(agent, RtpResearcherAgent)
|
||||
|
||||
|
||||
def test_executor_routes_net_researcher(tmp_path):
|
||||
from rockchip_agents.core.executor import _make_agent
|
||||
from rockchip_agents.core.queue import Task
|
||||
from rockchip_agents.tools.project_memory import ProjectMemory
|
||||
from rockchip_agents.agents.net_researcher import NetResearcherAgent
|
||||
from rockchip_agents.config import AgentsConfig, ClaudeConfig, SchedulerConfig, FeishuConfig
|
||||
|
||||
cfg = AgentsConfig(projects={}, scheduler=SchedulerConfig(),
|
||||
claude=ClaudeConfig(api_key="fake"), feishu=FeishuConfig(), devices={})
|
||||
task = Task(project="research", type="research", title="网络调研",
|
||||
context="", priority=3, mode="report", agent_role="net-researcher")
|
||||
agent = _make_agent(cfg, task, ProjectMemory(db_path=tmp_path / "m.db"))
|
||||
assert isinstance(agent, NetResearcherAgent)
|
||||
|
||||
|
||||
def test_executor_routes_kernel_analyzer(tmp_path):
|
||||
from rockchip_agents.core.executor import _make_agent
|
||||
from rockchip_agents.core.queue import Task
|
||||
from rockchip_agents.tools.project_memory import ProjectMemory
|
||||
from rockchip_agents.agents.kernel_analyzer import KernelAnalyzerAgent
|
||||
from rockchip_agents.config import AgentsConfig, ClaudeConfig, SchedulerConfig, FeishuConfig
|
||||
|
||||
cfg = AgentsConfig(projects={}, scheduler=SchedulerConfig(),
|
||||
claude=ClaudeConfig(api_key="fake"), feishu=FeishuConfig(), devices={})
|
||||
task = Task(project="research", type="research", title="内核调研",
|
||||
context="", priority=3, mode="report", agent_role="kernel-analyzer")
|
||||
agent = _make_agent(cfg, task, ProjectMemory(db_path=tmp_path / "m.db"))
|
||||
assert isinstance(agent, KernelAnalyzerAgent)
|
||||
|
||||
Reference in New Issue
Block a user