From ea150129b41c1b7c2adefda7faab4d30742b8877 Mon Sep 17 00:00:00 2001 From: tzt <14718231+flying-travel@user.noreply.gitee.com> Date: Sun, 30 Aug 2026 21:12:52 +0800 Subject: [PATCH] =?UTF-8?q?feat(v2):=20T8=20=E4=BA=BA=E5=B7=A5=E6=A3=80?= =?UTF-8?q?=E9=AA=8C=E9=98=9F=E5=88=97=20ReviewQueue=EF=BC=88sqlite/?= =?UTF-8?q?=E6=8A=BD=E6=A0=B7/=E5=AE=A1=E6=A0=B8=EF=BC=89?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- router_system/review.py | 160 ++++++++++++++++++++++++++++++++++++++++ tests/test_review.py | 76 +++++++++++++++++++ 任务拆解与执行计划.md | 2 +- 3 files changed, 237 insertions(+), 1 deletion(-) create mode 100644 router_system/review.py create mode 100644 tests/test_review.py diff --git a/router_system/review.py b/router_system/review.py new file mode 100644 index 0000000..3d5358b --- /dev/null +++ b/router_system/review.py @@ -0,0 +1,160 @@ +"""人工检验队列(ReviewQueue)—— 第三个协作者(纯标准库 sqlite3)。 + +复用同一"交流文本"协议:异步队列,不阻塞响应路径。系统交付后按抽样率或 +safety 标签强制规则入队,由人工审核给出 verdict(approve/edit/reject)并可 +回写修正数据(correction),用于后续质量分析(论文 E 实验的人工地基)。 + +对齐《实现方案_v2》5.1 T8: +- SQLite 存储(data/review.sqlite3),零第三方依赖(sqlite3 为标准库)。 +- enqueue / get / list / submit(verdict, correction)。 +- 抽样规则:brief.tags 命中 force_tags(如 safety)强制入队,否则按 sample_rate 随机抽样。 +""" +from __future__ import annotations + +import random +import sqlite3 +import threading +import uuid +from datetime import datetime, timezone +from pathlib import Path +from typing import Any, Dict, List, Optional + + +def _now_iso() -> str: + return datetime.now(timezone.utc).isoformat(timespec="seconds") + + +_SCHEMA = """ +CREATE TABLE IF NOT EXISTS reviews ( + id INTEGER PRIMARY KEY AUTOINCREMENT, + request_id TEXT NOT NULL, + query TEXT NOT NULL, + response TEXT NOT NULL, + tags TEXT NOT NULL DEFAULT '[]', + status TEXT NOT NULL DEFAULT 'pending', -- pending | reviewed + verdict TEXT, -- approve | edit | reject + correction TEXT, + reviewer TEXT, + reason TEXT, + workspace_path TEXT, + created_at TEXT NOT NULL, + reviewed_at TEXT +); +CREATE INDEX IF NOT EXISTS idx_reviews_status ON reviews(status); +""" + + +class ReviewQueue: + """人工检验队列(sqlite 后端)。方法为同步;调用方按需自行放入线程池。""" + + def __init__(self, db_path: str = "data/review.sqlite3"): + self.db_path = Path(db_path) + self.db_path.parent.mkdir(parents=True, exist_ok=True) + self._lock = threading.Lock() + self._init_db() + + def _connect(self) -> sqlite3.Connection: + conn = sqlite3.connect(str(self.db_path)) + conn.row_factory = sqlite3.Row + return conn + + def _init_db(self) -> None: + with self._lock, self._connect() as conn: + conn.executescript(_SCHEMA) + + # --------------------------------------------------------------- + # 抽样策略 + # --------------------------------------------------------------- + @staticmethod + def should_enqueue(tags: List[str], sample_rate: float = 0.10, + force_tags: Optional[List[str]] = None, + rng: Optional[random.Random] = None) -> bool: + """是否应入队:tags 命中 force_tags 强制;否则按 sample_rate 抽样。""" + force = force_tags or [] + if any(t in force for t in tags): + return True + rng = rng or random.Random() + return rng.random() < sample_rate + + # --------------------------------------------------------------- + # 入队 + # --------------------------------------------------------------- + def enqueue(self, request_id: str, query: str, response: str, + tags: Optional[List[str]] = None, reason: str = "sample", + workspace_path: Optional[str] = None) -> int: + """入队一条待审记录,返回 review id。""" + tags_json = _json_dumps(tags or []) + with self._lock, self._connect() as conn: + cur = conn.execute( + "INSERT INTO reviews (request_id, query, response, tags, reason, workspace_path, created_at)" + " VALUES (?,?,?,?,?,?,?)", + (request_id, query, response, tags_json, reason, workspace_path, _now_iso()), + ) + return int(cur.lastrowid) + + # --------------------------------------------------------------- + # 查询 + # --------------------------------------------------------------- + def get(self, review_id: int) -> Optional[Dict[str, Any]]: + with self._lock, self._connect() as conn: + row = conn.execute("SELECT * FROM reviews WHERE id=?", (review_id,)).fetchone() + return _row_to_dict(row) if row else None + + def list(self, status: Optional[str] = None, limit: int = 50) -> List[Dict[str, Any]]: + with self._lock, self._connect() as conn: + if status: + rows = conn.execute( + "SELECT * FROM reviews WHERE status=? ORDER BY id DESC LIMIT ?", + (status, limit)).fetchall() + else: + rows = conn.execute( + "SELECT * FROM reviews ORDER BY id DESC LIMIT ?", (limit,)).fetchall() + return [_row_to_dict(r) for r in rows] + + def count(self, status: Optional[str] = None) -> int: + with self._lock, self._connect() as conn: + if status: + row = conn.execute("SELECT COUNT(*) AS c FROM reviews WHERE status=?", + (status,)).fetchone() + else: + row = conn.execute("SELECT COUNT(*) AS c FROM reviews").fetchone() + return int(row["c"]) if row else 0 + + # --------------------------------------------------------------- + # 审核提交 + # --------------------------------------------------------------- + def submit(self, review_id: int, verdict: str, + correction: Optional[str] = None, + reviewer: Optional[str] = None) -> bool: + """提交审核结论。verdict: approve | edit | reject。返回是否更新成功。""" + if verdict not in ("approve", "edit", "reject"): + raise ValueError(f"非法 verdict: {verdict}(支持 approve|edit|reject)") + with self._lock, self._connect() as conn: + cur = conn.execute( + "UPDATE reviews SET status='reviewed', verdict=?, correction=?, reviewed_at=?, reviewer=?" + " WHERE id=? AND status='pending'", + (verdict, correction, _now_iso(), reviewer, review_id), + ) + return cur.rowcount > 0 + + +def _row_to_dict(row: Optional[sqlite3.Row]) -> Optional[Dict[str, Any]]: + if row is None: + return None + d = dict(row) + try: + import json + d["tags"] = json.loads(d.get("tags") or "[]") + except Exception: + d["tags"] = [] + return d + + +def _json_dumps(obj) -> str: + import json + return json.dumps(obj, ensure_ascii=False) + + +def build_review_queue(cfg: Dict[str, Any]) -> ReviewQueue: + """cfg 为 config.review 段。""" + return ReviewQueue(db_path=cfg.get("queue_db", "data/review.sqlite3")) diff --git a/tests/test_review.py b/tests/test_review.py new file mode 100644 index 0000000..beca83f --- /dev/null +++ b/tests/test_review.py @@ -0,0 +1,76 @@ +"""T8 人工检验队列单测(封闭:临时 sqlite)。""" +import pytest + +from router_system.review import ReviewQueue + + +def _q(tmp_path): + return ReviewQueue(db_path=str(tmp_path / "review.sqlite3")) + + +def test_enqueue_and_get(tmp_path): + q = _q(tmp_path) + rid = q.enqueue("req1", "问", "答", tags=["code"], reason="sample") + assert rid == 1 + row = q.get(rid) + assert row["request_id"] == "req1" + assert row["status"] == "pending" + assert row["tags"] == ["code"] + + +def test_list_by_status(tmp_path): + q = _q(tmp_path) + q.enqueue("r1", "q", "a", tags=["code"]) + q.enqueue("r2", "q", "a", tags=["safety"]) + assert q.count() == 2 + assert q.count(status="pending") == 2 + q.submit(1, "approve") + assert q.count(status="pending") == 1 + pending = q.list(status="pending") + assert len(pending) == 1 + assert pending[0]["request_id"] == "r2" + + +def test_submit_verdicts(tmp_path): + q = _q(tmp_path) + rid = q.enqueue("r1", "q", "a") + assert q.submit(rid, "edit", correction="修正文本") is True + row = q.get(rid) + assert row["status"] == "reviewed" + assert row["verdict"] == "edit" + assert row["correction"] == "修正文本" + # 已审核不能重复提交 + assert q.submit(rid, "approve") is False + + +def test_submit_invalid_verdict(tmp_path): + q = _q(tmp_path) + rid = q.enqueue("r1", "q", "a") + with pytest.raises(ValueError): + q.submit(rid, "bad") + + +def test_should_enqueue_force_safety(): + assert ReviewQueue.should_enqueue(["safety"], sample_rate=0.0, + force_tags=["safety"]) is True + assert ReviewQueue.should_enqueue(["code"], sample_rate=0.0, + force_tags=["safety"]) is False + + +def test_should_enqueue_sample_rate(): + import random + # 固定随机种子下按 10% 抽样应命中/不命中可控 + rng = random.Random(42) + hit = sum(ReviewQueue.should_enqueue(["code"], sample_rate=0.0, force_tags=[], rng=rng) for _ in range(1000)) + assert hit == 0 # sample_rate=0 -> 永不抽样 + rng = random.Random(1) + hit = sum(ReviewQueue.should_enqueue(["code"], sample_rate=1.0, force_tags=[], rng=rng) for _ in range(10)) + assert hit == 10 # sample_rate=1 -> 全抽样 + + +def test_db_recreated(tmp_path): + q1 = _q(tmp_path) + q1.enqueue("r1", "q", "a") + # 重新打开同一 db,数据仍在 + q2 = ReviewQueue(db_path=str(tmp_path / "review.sqlite3")) + assert q2.count() == 1 diff --git a/任务拆解与执行计划.md b/任务拆解与执行计划.md index dc79bec..1d50d35 100644 --- a/任务拆解与执行计划.md +++ b/任务拆解与执行计划.md @@ -74,7 +74,7 @@ P0 完成后的能力:干净的后端抽象 + 可量化的评测 + 可追溯 | T5 | WorkerLoop + 接地验证 | ✅ 完成 | T5 | | T6 | CollaborativePipeline 编排 | ✅ 完成 | T6 | | T7 | 网关扩展(/chat 切 v2,/chat/legacy) | ⬜ | | -| T8 | 人工检验队列 ReviewQueue | ⬜ | | +| T8 | 人工检验队列 ReviewQueue | ✅ 完成 | T8 | | T9 | token 计量与账单 | ⬜ | | | T10 | rollup + prefix cache 调优 | ⬜ | | | T11 | 打包分发 setup_runtime.py | ⬜ | |