feat: 版本流水线 + 算法组分组重组

tools/version_pipeline.py(新增):
- 并行开发流水线:v1测试时v2已在开发,不等待
- v1测试问题 → v3 backlog(跳一版本,不阻塞v2)
- [版本]/[测试通过]/[测试问题] 标签驱动流水线状态
- pipeline_summary() 注入 agent context,让 agent 知道当前版本状态

algo_vision.py(新增):
- AlgoVisionAgent:视觉算法组基类,rk3566主验证+rk3588可扩展
- _update_version_pipeline():解析版本标签,自动触发下版本开发任务
- _build_context() 注入版本流水线状态

算法组分组:
- 视觉组:algo-antishake → AlgoVisionAgent(Rockchip平台)
- MCU组:algo-position/nav → AlgoMcuAgent(暂缓)

scheduler:_enqueue_vision_algo_cycle() 只调度视觉组,MCU组注释保留

Co-Authored-By: Claude Sonnet 4.6 <noreply@anthropic.com>
This commit is contained in:
2026-03-09 21:45:51 +08:00
co-authored by Claude Sonnet 4.6
parent bb6e611dab
commit 28ed413037
2 changed files with 301 additions and 11 deletions
+66 -11
View File
@@ -8,6 +8,7 @@ from nmfs_agents.agents.developer import AgentResult, DeveloperAgent
from nmfs_agents.config import AgentsConfig
from nmfs_agents.core.queue import Task, TaskQueue
from nmfs_agents.tools.project_memory import ProjectMemory
from nmfs_agents.tools.version_pipeline import VersionPipeline
logger = logging.getLogger(__name__)
@@ -19,20 +20,30 @@ ALGO_VISION_SYSTEM = """你是视觉算法组工程师,专注瑞芯微(Rockc
- 平台能力:RKNN Lite NPUrk3566 ~1 TOPS/ RKNN NPUrk3588 ~6 TOPS
- 语言栈:Python 3.10 + numpy + opencv-python + rknn-lite
## 工作模式(每轮必须产出代码改进
1. **读记忆** - ProjectMemory __hardware__.algo.{domain}.* 获取上轮实测 benchmark
2. **看代码** - 读 src/ 现有实现,找本轮最值得优化的瓶颈
3. **改代码** - 针对性改进(不推倒重来),优先 numpy 向量化/NPU 路径
4. **rk3566 验证** - SSH 到设备运行 tests/,取实测数字
5. **写记忆** - [记忆] 更新 benchmark 数据
6. **定方向** - [后续] 一句话指明下轮重点
## 版本流水线规则(重要
- **并行开发**:v1 进入测试时立即开始 v2 开发,不等 v1 测试结果
- **问题跨版本**:v1 测试发现的问题 → 进入 v3 backlog(不 hotfix,不阻塞 v2
- **每轮输出版本标记**:`[版本] vN: 描述本轮改动` 让系统记录版本号
- **测试结果标记**:`[测试通过] vN` 或 `[测试问题] vN: 问题描述`
## 算法优先级原则
实时性(延迟 ms > 精度 > FPS > 内存占用
## 工作模式(每轮必须产出代码改进)
1. **读流水线** - 看上方版本状态,确认本轮做哪个版本
2. **读记忆** - ProjectMemory __hardware__.algo.{domain}.* 获取上轮实测数据
3. **读 backlog** - 本轮应解决哪些历史遗留问题
4. **改代码** - 针对性改进,优先 numpy 向量化/NPU 路径
5. **rk3566 验证** - SSH 运行 tests/ 获取实测数字
6. **标记版本** - [版本] + [测试通过]/[测试问题]
7. **写记忆** - [记忆] 更新 benchmark
8. **定方向** - [后续] 下轮改进点
## 算法优先级
实时性(延迟 ms > 精度 > FPS > 内存
## 输出格式(必须包含)
- [版本] vN: <本轮改动摘要>
- [测试通过] vN 或 [测试问题] vN: <问题描述>
- [记忆] __hardware__: algo.{domain}.{metric}=值
- [后续] <下一轮具体改进点>
- [后续] <下一轮聚焦点>
"""
@@ -44,14 +55,17 @@ class AlgoVisionAgent(DeveloperAgent):
config: AgentsConfig,
queue: Optional[TaskQueue] = None,
memory: Optional[ProjectMemory] = None,
version_pipeline: Optional[VersionPipeline] = None,
) -> None:
super().__init__(config, memory=memory)
self._queue = queue or TaskQueue()
self._vp = version_pipeline or VersionPipeline()
def run(self, task: Task) -> AgentResult:
result = super().run(task)
if result.status == "done":
self._extract_vision_facts(result.summary, task)
self._update_version_pipeline(result.summary, task)
self._dispatch_followup(result.summary, task)
return result
@@ -87,7 +101,9 @@ class AlgoVisionAgent(DeveloperAgent):
if algo_facts:
lines = "\n".join(f" {k}: {v}" for k, v in algo_facts.items())
hw_section = f"\n【上轮 Benchmark(来自记忆)】\n{lines}\n"
return base + hw_section + self._build_device_context() + "\n" + ALGO_VISION_SYSTEM
# 注入版本流水线状态
pipeline_section = "\n" + self._vp.pipeline_summary(task.project) + "\n"
return base + hw_section + pipeline_section + self._build_device_context() + "\n" + ALGO_VISION_SYSTEM
def _extract_vision_facts(self, report_text: str, task: Task) -> None:
"""解析 [记忆] 标签,写入 ProjectMemory __hardware__ facts。"""
@@ -101,6 +117,45 @@ class AlgoVisionAgent(DeveloperAgent):
except Exception as e:
logger.warning("视觉记忆写入失败 %s.%s: %s", project_key, sub_key, e)
def _update_version_pipeline(self, report_text: str, task: Task) -> None:
"""解析版本标签,更新流水线状态并在需要时触发下版本开发任务。"""
project = task.project
# [版本] vN: 描述 → 登记新版本
for m in re.finditer(r"\[版本\]\s*(v\d+)\s*[:]\s*(.+)", report_text):
version, desc = m.group(1), m.group(2).strip()[:200]
self._vp.start_version(project, version, desc)
# 版本登记后立即进入测试态
self._vp.start_testing(project, version)
# [测试通过] vN → 标记完成,触发下一版本开发
for m in re.finditer(r"\[测试通过\]\s*(v\d+)", report_text):
version = m.group(1)
self._vp.mark_done(project, version)
# 自动启动下下个版本(跳过已在开发中的)
next_v = self._vp.get_next_version_name(project)
backlog = self._vp.get_next_backlog(project)
backlog_ctx = "\n".join(f"- {b}" for b in backlog[:5]) if backlog else "无历史遗留问题"
self._queue.enqueue(Task(
project=project, type="code_improve",
title=f"视觉算法 {next_v} 开发",
context=(
f"【并行流水线:{version} 已完成,开始 {next_v} 开发】\n"
f"本轮 backlog(来自历史测试问题):\n{backlog_ctx}\n"
f"在此基础上继续迭代改进。"
),
agent_role=task.agent_role,
priority=task.priority, mode=task.mode,
parent_task_id=task.id, initiator=task.agent_role,
))
logger.info("[%s] %s 完成,已触发 %s 开发", project, version, next_v)
# [测试问题] vN: 问题描述 → 记录到 backlog(进入 v+2
for m in re.finditer(r"\[测试问题\]\s*(v\d+)\s*[:]\s*(.+)", report_text):
version, issue = m.group(1), m.group(2).strip()[:200]
self._vp.add_test_issue(project, version, issue)
logger.info("[%s] %s 问题记录: %s", project, version, issue)
def _dispatch_followup(self, report_text: str, task: Task) -> None:
"""解析 [后续] 标签,入队下轮迭代任务。"""
for m in re.finditer(r"\[后续\]\s*(.+)", report_text):
+235
View File
@@ -0,0 +1,235 @@
from __future__ import annotations
"""version_pipeline.py — 并行版本开发流水线管理。
流水线模式:
v1 开发完毕 → v1 测试(入队 tester)
v2 立即开始开发(不等 v1 测试完)
v1 测试发现问题 → 问题进入 v3 backlog(不 hotfix v1
v2 测试发现问题 → 进入 v4 backlog
...
每个项目独立维护自己的版本流水线状态(SQLite)。
"""
import json
import logging
import sqlite3
from dataclasses import dataclass, asdict
from datetime import datetime
from pathlib import Path
from typing import Optional
logger = logging.getLogger(__name__)
_DEFAULT_DB = Path("data/version_pipeline.db")
@dataclass
class VersionRecord:
id: int
project: str
version: str # v1, v2, v3, ...
description: str
status: str # developing | testing | done | abandoned
created_at: str
tested_at: Optional[str]
done_at: Optional[str]
issues: str # JSON list of issue descriptions found during testing
snapshot_path: str # optional: path to code snapshot
class VersionPipeline:
"""项目并行版本流水线管理器。
Usage (in agent context):
vp = VersionPipeline()
# 完成 v1 开发后
vid = vp.start_testing("antishake", "v1", "EIS 基础光流实现")
# 立即开始 v2 开发
v2 = vp.start_next_dev("antishake", "v2", "添加卡尔曼平滑")
# v1 测试发现问题 → 进入 v3 backlog(不是 v2
vp.add_test_issue("antishake", "v1", "光流在暗场景下漂移严重")
# v1 测试通过
vp.mark_done("antishake", "v1")
# 获取 v3 需要处理的问题列表
issues = vp.get_next_backlog("antishake")
"""
def __init__(self, db_path: Path = _DEFAULT_DB) -> None:
self._db_path = db_path
self._db_path.parent.mkdir(parents=True, exist_ok=True)
self._init_db()
def _init_db(self) -> None:
with self._connect() as conn:
conn.execute("""
CREATE TABLE IF NOT EXISTS versions (
id INTEGER PRIMARY KEY AUTOINCREMENT,
project TEXT NOT NULL,
version TEXT NOT NULL,
description TEXT NOT NULL DEFAULT '',
status TEXT NOT NULL DEFAULT 'developing',
created_at TEXT NOT NULL,
tested_at TEXT,
done_at TEXT,
issues TEXT NOT NULL DEFAULT '[]',
snapshot_path TEXT NOT NULL DEFAULT '',
UNIQUE(project, version)
)
""")
def _connect(self) -> sqlite3.Connection:
conn = sqlite3.connect(str(self._db_path))
conn.row_factory = sqlite3.Row
return conn
def _now(self) -> str:
return datetime.now().strftime("%Y-%m-%d %H:%M")
# ─── 版本生命周期 ──────────────────────────────────────────────
def start_version(self, project: str, version: str, description: str) -> int:
"""登记新版本开始开发,返回 id。"""
with self._connect() as conn:
try:
cur = conn.execute(
"INSERT INTO versions(project,version,description,status,created_at) "
"VALUES(?,?,?,'developing',?)",
(project, version, description, self._now()),
)
vid = cur.lastrowid
logger.info("[%s] %s 开始开发: %s", project, version, description)
return vid
except sqlite3.IntegrityError:
logger.warning("[%s] %s 已存在", project, version)
return -1
def start_testing(self, project: str, version: str) -> bool:
"""将版本状态置为 testing,同时应立即启动下一版本开发。"""
with self._connect() as conn:
cur = conn.execute(
"UPDATE versions SET status='testing', tested_at=? "
"WHERE project=? AND version=? AND status='developing'",
(self._now(), project, version),
)
ok = cur.rowcount > 0
if ok:
logger.info("[%s] %s 进入测试阶段", project, version)
return ok
def add_test_issue(self, project: str, version: str, issue: str) -> None:
"""记录测试中发现的问题(将进入 version+2 的 backlog)。"""
with self._connect() as conn:
row = conn.execute(
"SELECT issues FROM versions WHERE project=? AND version=?",
(project, version),
).fetchone()
if not row:
return
issues = json.loads(row["issues"])
issues.append({"issue": issue, "found_at": self._now()})
conn.execute(
"UPDATE versions SET issues=? WHERE project=? AND version=?",
(json.dumps(issues, ensure_ascii=False), project, version),
)
logger.info("[%s] %s 发现问题: %s", project, version, issue)
def mark_done(self, project: str, version: str) -> None:
"""版本测试通过,标记为 done。"""
with self._connect() as conn:
conn.execute(
"UPDATE versions SET status='done', done_at=? WHERE project=? AND version=?",
(self._now(), project, version),
)
logger.info("[%s] %s 测试通过,已完成", project, version)
# ─── 查询 ──────────────────────────────────────────────────────
def get_status(self, project: str) -> list[dict]:
"""返回项目所有版本状态。"""
with self._connect() as conn:
rows = conn.execute(
"SELECT version,description,status,created_at,tested_at,done_at,issues "
"FROM versions WHERE project=? ORDER BY id",
(project,),
).fetchall()
return [dict(r) for r in rows]
def get_current_dev_version(self, project: str) -> Optional[str]:
"""返回当前正在开发的版本号(最新的 developing 状态)。"""
with self._connect() as conn:
row = conn.execute(
"SELECT version FROM versions WHERE project=? AND status='developing' "
"ORDER BY id DESC LIMIT 1",
(project,),
).fetchone()
return row["version"] if row else None
def get_testing_version(self, project: str) -> Optional[str]:
"""返回当前正在测试的版本号。"""
with self._connect() as conn:
row = conn.execute(
"SELECT version FROM versions WHERE project=? AND status='testing' "
"ORDER BY id DESC LIMIT 1",
(project,),
).fetchone()
return row["version"] if row else None
def get_next_backlog(self, project: str) -> list[str]:
"""收集所有已完成版本的测试问题,作为下一轮开发的 backlog。
规则:v1 发现的问题 → v3 backlog(跳一个版本),
这样 v2 已在开发中的工作不被打断。
"""
with self._connect() as conn:
rows = conn.execute(
"SELECT version, issues FROM versions "
"WHERE project=? AND (status='done' OR status='testing') AND issues != '[]' "
"ORDER BY id",
(project,),
).fetchall()
backlog: list[str] = []
for row in rows:
issues = json.loads(row["issues"])
for item in issues:
backlog.append(f"[来自{row['version']}测试] {item['issue']}")
return backlog
def get_next_version_name(self, project: str) -> str:
"""自动推算下一个版本号(v1→v2→v3...)。"""
with self._connect() as conn:
row = conn.execute(
"SELECT version FROM versions WHERE project=? ORDER BY id DESC LIMIT 1",
(project,),
).fetchone()
if not row:
return "v1"
last = row["version"] # e.g. "v3"
try:
num = int(last.lstrip("v")) + 1
except ValueError:
num = 1
return f"v{num}"
def pipeline_summary(self, project: str) -> str:
"""返回流水线状态的人类可读摘要,用于注入 agent context。"""
status = self.get_status(project)
if not status:
return f"[{project}] 版本流水线:尚无版本记录,请从 v1 开始。"
lines = [f"[{project}] 版本流水线状态:"]
for v in status:
issues = json.loads(v["issues"])
issue_str = f"{len(issues)}个问题待处理" if issues else ""
lines.append(f" {v['version']}: {v['status']}{issue_str}{v['description']}")
backlog = self.get_next_backlog(project)
if backlog:
lines.append(f"\n下一版本 backlog{len(backlog)}条):")
for b in backlog[:5]:
lines.append(f" - {b}")
next_v = self.get_next_version_name(project)
lines.append(f"\n下一个版本号:{next_v}")
return "\n".join(lines)