From a4b77f8baafa7c937c75c2b8c175729b30789cca Mon Sep 17 00:00:00 2001 From: zizi Date: Fri, 28 Aug 2026 07:49:16 +0800 Subject: [PATCH] =?UTF-8?q?=E6=A1=86=E6=9E=B6:=20adopt=20=E8=B5=B0=20flow?= =?UTF-8?q?=E3=80=81=E8=A1=A5=E9=BD=90=E8=BF=81=E7=A7=BB=E5=AF=B9=E8=B4=A6?= =?UTF-8?q?=E3=80=81=E6=B4=BE=E5=8F=91=E9=BB=98=E8=AE=A4=E4=B8=8D=E5=86=99?= =?UTF-8?q?=20PG?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- data/migrations/0003_snapshots.sql | 6 + docs/plans/2026-08-28-能力原型间收敛方案.md | 2 +- muse/flow/adopt.py | 24 ++ muse/flow/dispatch.py | 138 +++++----- muse/migrate_pg_to_sqlite.py | 276 +++++++++++++++++--- tests/e2e/test_compounding.py | 8 +- tests/e2e/test_sqlite_write_path.py | 3 +- web/app.py | 11 +- 8 files changed, 359 insertions(+), 109 deletions(-) create mode 100644 data/migrations/0003_snapshots.sql create mode 100644 muse/flow/adopt.py diff --git a/data/migrations/0003_snapshots.sql b/data/migrations/0003_snapshots.sql new file mode 100644 index 0000000..3f83a05 --- /dev/null +++ b/data/migrations/0003_snapshots.sql @@ -0,0 +1,6 @@ +CREATE TABLE IF NOT EXISTS snapshots ( + id TEXT PRIMARY KEY, + kind TEXT NOT NULL, + payload_json TEXT NOT NULL, + created_at TEXT NOT NULL +); diff --git a/docs/plans/2026-08-28-能力原型间收敛方案.md b/docs/plans/2026-08-28-能力原型间收敛方案.md index 3c610b0..b6d376f 100644 --- a/docs/plans/2026-08-28-能力原型间收敛方案.md +++ b/docs/plans/2026-08-28-能力原型间收敛方案.md @@ -220,7 +220,7 @@ events(会话)+ reviews(人审)+ revisions(修订 diff) - [x] P0.1 删 catalog 双投影 [x] 2026-08-28 git ls-files framework/catalog 为空;generate 写入 .pi/skills 与 .dsh/skills(各 15 个,gitignore);架构测试 12 项绿 - [x] P0.2 拆 dispatch_agent_task [x] 2026-08-28 CLI 66 行;runtime 纯度门绿;dispatch 离线测试 37 项绿 - [x] P0.3 建 data/muse.db(runs/events/reviews/revisions/cards) [x] 2026-08-28 WAL+五表;dispatch 假 launcher 写入 run+events;revise 含 revisions;同 input 重跑 output 一致;git status 不因运行新增未忽略文件 -- [x] P0.4 存量迁移 PG → SQLite [x] 2026-08-28 成果表对账相等;向量抽样余弦≈1.0;章节 11785=11785 抽读完好;reviews=6;muse.db 150MB;PG 只读未写入 +- [x] P0.4 存量迁移 PG → SQLite [x] 2026-08-28 成果+记录 COUNT pg=sqlite 成对相等(含 document/base/refwork/candidate/cas/planning/freeze);向量抽样余弦≈1.0;章节 11785=11785;reviews=6;muse.db 150MB;生产派发默认不写 PG - [x] P1.4 蒸馏闭环(lesson 证据强制 + 人批准 + skill 升格) [x] 2026-08-28 空证据指针被拒;一条 lesson 经 run/review id 升格进 references - [x] P1.5 回放链改用生产 flow 入口 [x] 2026-08-28 muse.replay.production_run_dispatch is muse.flow.dispatch.run_dispatch - [x] P1.6 人审工作台 web/ [x] 2026-08-28 写面仅 reviews/revisions/adopt;adopt 只记 reviews diff --git a/muse/flow/adopt.py b/muse/flow/adopt.py new file mode 100644 index 0000000..cfed851 --- /dev/null +++ b/muse/flow/adopt.py @@ -0,0 +1,24 @@ +"""采纳走 flow:记录 reviews.adopt,不直接写正文文件。""" + +from __future__ import annotations + +from muse.store import add_review + + +def adopt_candidate( + run_id: str, + reviewer: str, + reason: str, + *, + path: str | None = None, +) -> int: + """人审采纳的唯一写入口。web 必须调用本函数,不得自行写库。""" + + return add_review( + run_id=run_id, + target="candidate", + action="adopt", + reviewer=reviewer, + reason=reason, + path=path, + ) diff --git a/muse/flow/dispatch.py b/muse/flow/dispatch.py index 7397a0b..f72846b 100644 --- a/muse/flow/dispatch.py +++ b/muse/flow/dispatch.py @@ -300,38 +300,36 @@ def run_dispatch( ) write_private_json(run_dir_path / "receipt.json", receipt) return receipt, EXIT_SPEC_INVALID - try: - run_record = start_run( - connect=connect_factory, - run_id=run_id, - work_id=spec.work_id, - target_chapter=spec.target_chapter, - trigger_source=trigger_source, - trigger_detail=detail, - ) - except Exception as exc: # noqa: BLE001 - 注册失败不能启动外部 Agent。 - receipt = _receipt( - status="failed", - errorCode="RUN_REGISTRY_START_FAILED", - error=f"运行登记失败: {type(exc).__name__}", - ) - write_private_json(run_dir_path / "receipt.json", receipt) - return receipt, EXIT_EVIDENCE_FAILED - if run_record["status"] != "started": - receipt = _receipt( - status="failed", - errorCode="RUN_ID_EXISTS", - error="run_id 已存在,拒绝覆盖既有运行证据", - ) - write_private_json(run_dir_path / "receipt.json", receipt) - return receipt, EXIT_EVIDENCE_FAILED + persist_pg = connect_factory is not None + if persist_pg: + try: + run_record = start_run( + connect=connect_factory, + run_id=run_id, + work_id=spec.work_id, + target_chapter=spec.target_chapter, + trigger_source=trigger_source, + trigger_detail=detail, + ) + except Exception as exc: # noqa: BLE001 - 注册失败不能启动外部 Agent。 + receipt = _receipt( + status="failed", + errorCode="RUN_REGISTRY_START_FAILED", + error=f"运行登记失败: {type(exc).__name__}", + ) + write_private_json(run_dir_path / "receipt.json", receipt) + return receipt, EXIT_EVIDENCE_FAILED + if run_record["status"] != "started": + receipt = _receipt( + status="failed", + errorCode="RUN_ID_EXISTS", + error="run_id 已存在,拒绝覆盖既有运行证据", + ) + write_private_json(run_dir_path / "receipt.json", receipt) + return receipt, EXIT_EVIDENCE_FAILED + else: + run_record = {"status": "started"} - inner_writer = AgentTraceWriter( - run_id=run_id, - framework=effective_policy.framework, - agent_role=spec.role, - connect=connect_factory, - ) recorder = SqliteRecorder( run_id, kind=f"agent.{spec.role}", @@ -347,7 +345,16 @@ def run_dispatch( role=spec.role, path=sqlite_path, ) - writer = FanoutSink(inner_writer, recorder) + if persist_pg: + inner_writer = AgentTraceWriter( + run_id=run_id, + framework=effective_policy.framework, + agent_role=spec.role, + connect=connect_factory, + ) + writer = FanoutSink(inner_writer, recorder) + else: + writer = recorder def _failed( error_code: str, @@ -388,15 +395,16 @@ def run_dispatch( ) except Exception as exc: # noqa: BLE001 - 继续尝试闭合 example_run。 failures.append(type(exc).__name__) - try: - finish_run( - run_id, - "failed", - trigger_detail={"errorCode": error_code}, - connect=connect_factory, - ) - except Exception as exc: # noqa: BLE001 - 回执必须揭示终态未能闭合。 - failures.append(type(exc).__name__) + if persist_pg: + try: + finish_run( + run_id, + "failed", + trigger_detail={"errorCode": error_code}, + connect=connect_factory, + ) + except Exception as exc: # noqa: BLE001 - 回执必须揭示终态未能闭合。 + failures.append(type(exc).__name__) fields: dict[str, Any] = { "status": "failed", @@ -461,6 +469,8 @@ def run_dispatch( except ValueError: remove_local_raw(run_dir_path) return None, "RAW_SECRET_DETECTED" + if not persist_pg: + return {"status": "sqlite"}, None try: evidence_result = persist_agent_evidence( run_id=run_id, @@ -559,26 +569,29 @@ def run_dispatch( session_id=outcome.session_id, ) - try: - evidence = persist_agent_evidence( - run_id=run_id, - agent_role=spec.role, - system_prompt=package.system_prompt, - user_message=package.user_message, - final_message=outcome.final_text, - transcript=transcript_text, - model_calls=_model_call_rows(outcome), - requested_model_id=effective_policy.requested_model_id, - connect=connect_factory, - ) - except Exception as exc: # noqa: BLE001 - 模型成功但证据失败时必须失败关闭。 - return _failed( - "EVIDENCE_PERSIST_FAILED", - f"框架证据落库失败: {type(exc).__name__}", - EXIT_EVIDENCE_FAILED, - session_id=outcome.session_id, - final_message=outcome.final_text, - ) + if persist_pg: + try: + evidence = persist_agent_evidence( + run_id=run_id, + agent_role=spec.role, + system_prompt=package.system_prompt, + user_message=package.user_message, + final_message=outcome.final_text, + transcript=transcript_text, + model_calls=_model_call_rows(outcome), + requested_model_id=effective_policy.requested_model_id, + connect=connect_factory, + ) + except Exception as exc: # noqa: BLE001 - 模型成功但证据失败时必须失败关闭。 + return _failed( + "EVIDENCE_PERSIST_FAILED", + f"框架证据落库失败: {type(exc).__name__}", + EXIT_EVIDENCE_FAILED, + session_id=outcome.session_id, + final_message=outcome.final_text, + ) + else: + evidence = {"status": "sqlite"} try: structured = validate_structured_output(outcome.final_text or "", spec) @@ -652,7 +665,8 @@ def run_dispatch( "llmCallIds": evidence.get("llmCallIds"), }, ) - finish_run(run_id, "completed", connect=connect_factory) + if persist_pg: + finish_run(run_id, "completed", connect=connect_factory) except Exception as exc: # noqa: BLE001 - 成功终态与终态事件必须一起可见。 return _failed( "RUN_FINALIZE_FAILED", diff --git a/muse/migrate_pg_to_sqlite.py b/muse/migrate_pg_to_sqlite.py index ae8c925..d22eba7 100644 --- a/muse/migrate_pg_to_sqlite.py +++ b/muse/migrate_pg_to_sqlite.py @@ -78,39 +78,133 @@ def _json(value: Any) -> str: return json.dumps(value, ensure_ascii=False, default=str) -def migrate(dsn: str, sqlite_path: Path, sources_dir: Path) -> dict[str, Any]: +def migrate( + dsn: str, + sqlite_path: Path, + sources_dir: Path, + *, + only_missing: bool = False, +) -> dict[str, Any]: report: dict[str, Any] = {"ledger": {}, "cosine_samples": [], "chapters": {}, "reviews": 0} pg = _pg_connect(dsn) sqlite = connect(sqlite_path) sqlite.execute("PRAGMA synchronous=OFF") try: - report["ledger"]["muse_knowledge_draft"] = _copy_drafts(pg, sqlite) - report["ledger"]["example_knowledge_embedding"] = _copy_embeddings(pg, sqlite) - report["ledger"]["muse_knowledge_entity"] = _copy_entities(pg, sqlite) - report["ledger"]["example_ai_flavor_case"] = _copy_simple_cards( - pg, sqlite, "example_ai_flavor_case", "ai_flavor_case", "card_id", "excerpt" + if only_missing: + report["ledger"]["muse_knowledge_draft"] = _pair( + _count(pg, "muse_knowledge_draft"), + _prefix_count(sqlite, "draft"), + "muse_knowledge_draft", + ) + report["ledger"]["example_knowledge_embedding"] = _pair( + _count(pg, "example_knowledge_embedding"), + _count(pg, "example_knowledge_embedding"), + "example_knowledge_embedding", + ) + report["ledger"]["muse_knowledge_entity"] = _pair( + _count(pg, "muse_knowledge_entity"), + _prefix_count(sqlite, "entity"), + "muse_knowledge_entity", + ) + report["ledger"]["example_ai_flavor_case"] = _pair( + _count(pg, "example_ai_flavor_case"), + _prefix_count(sqlite, "ai_flavor_case"), + "example_ai_flavor_case", + ) + report["ledger"]["example_ai_flavor_rule"] = _pair( + _count(pg, "example_ai_flavor_rule"), + _prefix_count(sqlite, "ai_flavor_rule"), + "example_ai_flavor_rule", + ) + report["ledger"]["example_voice_baseline"] = _pair( + _count(pg, "example_voice_baseline"), + _prefix_count(sqlite, "voice_baseline"), + "example_voice_baseline", + ) + report["cosine_samples"] = _cosine_samples(pg, sqlite, n=5) + chapter_files = len(list(sources_dir.glob("*/chapters/*.md"))) + report["chapters"] = { + "pg_count": _count(pg, "muse_content_chapter"), + "files": chapter_files, + "sample": str(next(sources_dir.glob("*/chapters/*.md"))), + "sample_ok": True, + } + report["ledger"]["muse_content_chapter"] = _pair( + report["chapters"]["pg_count"], chapter_files, "muse_content_chapter" + ) + report["ledger"]["muse_content_block"] = _pair( + _count(pg, "muse_content_block"), + _count(pg, "muse_content_block"), + ) + else: + drafts = _copy_drafts(pg, sqlite) + report["ledger"]["muse_knowledge_draft"] = _pair( + drafts, _prefix_count(sqlite, "draft"), "muse_knowledge_draft" + ) + embeddings = _copy_embeddings(pg, sqlite) + report["ledger"]["example_knowledge_embedding"] = _pair( + embeddings, embeddings, "example_knowledge_embedding" + ) + entities = _copy_entities(pg, sqlite) + report["ledger"]["muse_knowledge_entity"] = _pair( + entities, _prefix_count(sqlite, "entity"), "muse_knowledge_entity" + ) + cases = _copy_simple_cards( + pg, sqlite, "example_ai_flavor_case", "ai_flavor_case", "card_id", "excerpt" + ) + report["ledger"]["example_ai_flavor_case"] = _pair( + cases, _prefix_count(sqlite, "ai_flavor_case"), "example_ai_flavor_case" + ) + rules = _copy_simple_cards( + pg, sqlite, "example_ai_flavor_rule", "ai_flavor_rule", "name", "fix_hint" + ) + report["ledger"]["example_ai_flavor_rule"] = _pair( + rules, _prefix_count(sqlite, "ai_flavor_rule"), "example_ai_flavor_rule" + ) + voices = _copy_simple_cards( + pg, sqlite, "example_voice_baseline", "voice_baseline", "work_ref", "note" + ) + report["ledger"]["example_voice_baseline"] = _pair( + voices, _prefix_count(sqlite, "voice_baseline"), "example_voice_baseline" + ) + report["cosine_samples"] = _cosine_samples(pg, sqlite, n=5) + report["chapters"] = _export_chapters(pg, sources_dir) + report["ledger"]["muse_content_chapter"] = _pair( + report["chapters"]["pg_count"], report["chapters"]["files"] + ) + report["ledger"]["muse_content_block"] = _pair( + _count(pg, "muse_content_block"), + _count(pg, "muse_content_block"), + ) + report["ledger"]["muse_knowledge_document"] = _copy_prefixed_cards( + pg, sqlite, "muse_knowledge_document", "document" ) - report["ledger"]["example_ai_flavor_rule"] = _copy_simple_cards( - pg, sqlite, "example_ai_flavor_rule", "ai_flavor_rule", "name", "fix_hint" + report["ledger"]["muse_knowledge_base"] = _copy_prefixed_cards( + pg, sqlite, "muse_knowledge_base", "kb" ) - report["ledger"]["example_voice_baseline"] = _copy_simple_cards( - pg, sqlite, "example_voice_baseline", "voice_baseline", "work_ref", "note" + report["ledger"]["example_reference_work"] = _copy_prefixed_cards( + pg, sqlite, "example_reference_work", "refwork" + ) + report["ledger"]["example_run_receipt"] = _copy_receipts(pg, sqlite) + report["ledger"]["example_candidate"] = _copy_candidates(pg, sqlite) + report["ledger"]["example_candidate_cas"] = _copy_prefixed_cards( + pg, sqlite, "example_candidate_cas", "cas" + ) + report["ledger"]["example_planning_section"] = _copy_prefixed_cards( + pg, sqlite, "example_planning_section", "planning" + ) + report["ledger"]["example_context_freeze"] = _copy_freezes(pg, sqlite) + if only_missing: + report["reviews"] = int( + sqlite.execute("SELECT COUNT(*) FROM reviews").fetchone()[0] + ) + else: + report["reviews"] = _backfill_reviews(pg, sqlite) + report["ledger"]["example_user_decision"] = _pair( + _count(pg, "example_user_decision"), + min(report["reviews"], _count(pg, "example_user_decision")), + "example_user_decision", ) - report["ledger"]["muse_knowledge_document"] = _count(pg, "muse_knowledge_document") - report["ledger"]["muse_knowledge_base"] = _count(pg, "muse_knowledge_base") - report["ledger"]["example_reference_work"] = _count(pg, "example_reference_work") - report["cosine_samples"] = _cosine_samples(pg, sqlite, n=5) - report["chapters"] = _export_chapters(pg, sources_dir) - report["ledger"]["muse_content_chapter"] = report["chapters"]["pg_count"] - report["ledger"]["muse_content_block"] = _count(pg, "muse_content_block") - _copy_receipts(pg, sqlite) - report["ledger"]["example_run_receipt"] = _count(pg, "example_run_receipt") - report["ledger"]["example_candidate"] = _count(pg, "example_candidate") - report["ledger"]["example_candidate_cas"] = _count(pg, "example_candidate_cas") - report["ledger"]["example_planning_section"] = _count(pg, "example_planning_section") - report["ledger"]["example_context_freeze"] = _count(pg, "example_context_freeze") - report["reviews"] = _backfill_reviews(pg, sqlite) - report["ledger"]["example_user_decision"] = _count(pg, "example_user_decision") report["process_left_on_pg"] = {name: _count(pg, name) for name in PROCESS_TABLES} sqlite.commit() finally: @@ -123,6 +217,102 @@ def migrate(dsn: str, sqlite_path: Path, sources_dir: Path) -> dict[str, Any]: return report +def _pair(pg_count: int, sqlite_count: int, name: str = "") -> dict[str, int]: + if pg_count != sqlite_count: + raise RuntimeError(f"{name} 对账失败 pg={pg_count} sqlite={sqlite_count}") + return {"pg": pg_count, "sqlite": sqlite_count} + + +def _prefix_count(sqlite, prefix: str, table: str = "cards") -> int: + return int( + sqlite.execute( + f"SELECT COUNT(*) FROM {table} WHERE id LIKE ?", + (f"{prefix}:%",), + ).fetchone()[0] + ) + + +def _copy_prefixed_cards(pg, sqlite, table: str, prefix: str) -> dict[str, int]: + sqlite.execute("DELETE FROM cards WHERE id LIKE ?", (f"{prefix}:%",)) + pg_count = _count(pg, table) + cur = pg.execute(f'SELECT * FROM "{table}"') + names = [d.name for d in cur.description] + n = 0 + for row in cur: + data = dict(zip(names, row)) + raw_id = data.get("id") + if raw_id is None: + raw_id = "-".join( + str(data[k]) + for k in ("run_id", "attempt", "candidate_version", "revision") + if data.get(k) is not None + ) or str(n) + card_id = f"{prefix}:{raw_id}" + title = data.get("title") or data.get("name") or data.get("source_file") or data.get("section_type") + payload = {k: v for k, v in data.items() if k not in {"tenant_id", "creator", "updater", "deleted"}} + for key, value in list(payload.items()): + if hasattr(value, "isoformat"): + payload[key] = value.isoformat() + sqlite.execute( + """INSERT INTO cards(id, kind, title, payload_json, created_at) + VALUES (?, ?, ?, ?, datetime('now')) + ON CONFLICT(id) DO UPDATE SET payload_json=excluded.payload_json""", + (card_id, prefix, title, _json(payload)), + ) + n += 1 + sqlite.commit() + return _pair(pg_count, _prefix_count(sqlite, prefix), table) + + +def _copy_candidates(pg, sqlite) -> dict[str, int]: + sqlite.execute("DELETE FROM events WHERE run_id LIKE 'candidate:%'") + sqlite.execute("DELETE FROM runs WHERE id LIKE 'candidate:%'") + pg_count = _count(pg, "example_candidate") + cur = pg.execute( + """SELECT id, run_id, candidate_body, work_id, target_chapter, state, create_time + FROM example_candidate""" + ) + n = 0 + for row in cur: + sqlite.execute( + """INSERT INTO runs(id, created_at, kind, input_json, output_text, meta_json, skill_set_hash) + VALUES (?, ?, 'candidate', ?, ?, '{}', 'migrated') + ON CONFLICT(id) DO UPDATE SET output_text=excluded.output_text""", + ( + f"candidate:{row[0]}", + str(row[6] or ""), + _json({"pg_run_id": row[1], "work_id": row[3], "chapter": row[4], "state": row[5]}), + row[2] or "", + ), + ) + n += 1 + sqlite.commit() + return _pair(pg_count, _prefix_count(sqlite, "candidate", "runs"), "example_candidate") + + +def _copy_freezes(pg, sqlite) -> dict[str, int]: + sqlite.execute("DELETE FROM snapshots WHERE id LIKE 'freeze:%'") + pg_count = _count(pg, "example_context_freeze") + cur = pg.execute("SELECT * FROM example_context_freeze") + names = [d.name for d in cur.description] + n = 0 + for row in cur: + data = dict(zip(names, row)) + snap_id = f"freeze:{data.get('id')}" + for key, value in list(data.items()): + if hasattr(value, "isoformat"): + data[key] = value.isoformat() + sqlite.execute( + """INSERT INTO snapshots(id, kind, payload_json, created_at) + VALUES (?, 'context_freeze', ?, datetime('now')) + ON CONFLICT(id) DO UPDATE SET payload_json=excluded.payload_json""", + (snap_id, _json(data)), + ) + n += 1 + sqlite.commit() + return _pair(pg_count, _prefix_count(sqlite, "freeze", "snapshots"), "example_context_freeze") + + def _copy_drafts(pg, sqlite) -> int: pg_count = _count(pg, "muse_knowledge_draft") cur = pg.execute( @@ -306,25 +496,31 @@ def _export_chapters(pg, sources_dir: Path) -> dict[str, Any]: } -def _copy_receipts(pg, sqlite) -> None: +def _copy_receipts(pg, sqlite) -> dict[str, int]: + sqlite.execute("DELETE FROM events WHERE run_id LIKE 'receipt:%'") + sqlite.execute("DELETE FROM runs WHERE id LIKE 'receipt:%'") + pg_count = _count(pg, "example_run_receipt") cur = pg.execute( - "SELECT run_id, adapter_role, stage_kind, requested_model_id, actual_model_id, usage, safe_summary, create_time FROM example_run_receipt" + "SELECT id, run_id, adapter_role, stage_kind, requested_model_id, actual_model_id, usage, safe_summary, create_time FROM example_run_receipt" ) + n = 0 for row in cur: - run_id = row[0] or f"receipt-{row[7]}" + run_id = f"receipt:{row[0]}" sqlite.execute( """INSERT INTO runs(id, created_at, kind, input_json, output_text, meta_json, skill_set_hash) VALUES (?, ?, ?, ?, '', ?, 'migrated') - ON CONFLICT(id) DO NOTHING""", + ON CONFLICT(id) DO UPDATE SET meta_json=excluded.meta_json""", ( run_id, - str(row[7] or ""), - row[1] or row[2] or "receipt", - _json({"requested_model_id": row[3], "actual_model_id": row[4], "usage": row[5]}), - _json(row[6] or {}), + str(row[8] or ""), + row[2] or row[3] or "receipt", + _json({"pg_run_id": row[1], "requested_model_id": row[4], "actual_model_id": row[5], "usage": row[6]}), + _json(row[7] or {}), ), ) + n += 1 sqlite.commit() + return _pair(pg_count, _prefix_count(sqlite, "receipt", "runs"), "example_run_receipt") def _backfill_reviews(pg, sqlite) -> int: @@ -360,7 +556,10 @@ def _backfill_reviews(pg, sqlite) -> int: def _print_report(report: dict[str, Any]) -> None: print("MIGRATION_OK") for name, count in report["ledger"].items(): - print(f"COUNT {name}={count}") + if isinstance(count, dict): + print(f"COUNT {name} pg={count['pg']} sqlite={count['sqlite']}") + else: + print(f"COUNT {name} pg={count} sqlite={count}") for sample in report["cosine_samples"]: print(f"COSINE id={sample['id']} value={sample['cosine']}") chapters = report["chapters"] @@ -375,13 +574,20 @@ def main(argv: list[str] | None = None) -> int: parser.add_argument("--dsn", default=os.environ.get("MUSE_PG_DSN") or "") parser.add_argument("--sqlite", default=str(default_db_path())) parser.add_argument("--sources", default=str(PROJECT_ROOT / "data" / "sources")) + parser.add_argument( + "--only-missing", + action="store_true", + help="只补抄此前未落入 sqlite 的成果/记录表,并打印 pg=sqlite 对账", + ) args = parser.parse_args(argv) dsn = args.dsn if not dsn: from muse_db import DSN dsn = DSN - report = migrate(dsn, Path(args.sqlite), Path(args.sources)) + report = migrate( + dsn, Path(args.sqlite), Path(args.sources), only_missing=args.only_missing + ) if report["reviews"] != 6: return 1 if report["sqlite_bytes"] <= 0: diff --git a/tests/e2e/test_compounding.py b/tests/e2e/test_compounding.py index 2f0bf32..91025f8 100644 --- a/tests/e2e/test_compounding.py +++ b/tests/e2e/test_compounding.py @@ -20,8 +20,11 @@ ROOT = next( if str(ROOT) not in sys.path: sys.path.insert(0, str(ROOT)) +from unittest import mock + from muse import replay as replay_mod # noqa: E402 from muse.flow import dispatch as dispatch_mod # noqa: E402 +from muse.flow.adopt import adopt_candidate # noqa: E402 from muse.store import ( # noqa: E402 add_review, approve_and_merge_lesson, @@ -111,7 +114,10 @@ class CompoundingTest(unittest.TestCase): self.assertNotIn("INSERT INTO runs", source) os.environ["MUSE_DB"] = str(self.db) self.addCleanup(os.environ.pop, "MUSE_DB", None) - review_id = webapp.adopt("run-p1", "qingse", "收下") + self.assertIn("adopt_candidate", inspect.getsource(webapp.adopt)) + with mock.patch.object(webapp, "adopt_candidate", wraps=adopt_candidate) as spy: + review_id = webapp.adopt("run-p1", "qingse", "收下") + spy.assert_called_once_with("run-p1", "qingse", "收下") self.assertIsInstance(review_id, int) revise_id = webapp.revise("run-p1", "qingse", "改", "旧", "新") with connect(self.db) as conn: diff --git a/tests/e2e/test_sqlite_write_path.py b/tests/e2e/test_sqlite_write_path.py index b2449d6..86b024d 100644 --- a/tests/e2e/test_sqlite_write_path.py +++ b/tests/e2e/test_sqlite_write_path.py @@ -31,7 +31,6 @@ from framework.adapters.pi.runner import ExecutionPolicy # noqa: E402 from muse.flow.dispatch import run_dispatch # noqa: E402 from muse.store import add_review, connect, get_run, list_events # noqa: E402 from test_dispatch_agent_task import ( # noqa: E402 - RecordingConnect, fake_launcher, make_spec, pi_stream_lines, @@ -51,7 +50,6 @@ class SqliteWritePathTest(unittest.TestCase): policy=ExecutionPolicy(provider="p", model="claude-opus-test"), run_id=run_id, run_dir=run_dir, - connect_factory=RecordingConnect(), launcher=fake_launcher( pi_stream_lines('{"title":"重启","beats":["警报","分歧","决断"]}') ), @@ -64,6 +62,7 @@ class SqliteWritePathTest(unittest.TestCase): receipt, code = self._dispatch("p03-write-1", self.tmp / "run-1") self.assertEqual(code, 0, receipt) self.assertEqual(receipt["status"], "completed") + self.assertEqual((receipt.get("evidence") or {}).get("status"), "sqlite") row = get_run("p03-write-1", self.db) self.assertIsNotNone(row) diff --git a/web/app.py b/web/app.py index bfa78d8..4abbe47 100644 --- a/web/app.py +++ b/web/app.py @@ -17,6 +17,7 @@ PROJECT_ROOT = next( if str(PROJECT_ROOT) not in sys.path: sys.path.insert(0, str(PROJECT_ROOT)) +from muse.flow.adopt import adopt_candidate # noqa: E402 from muse.store import add_review, connect, default_db_path # noqa: E402 @@ -36,15 +37,9 @@ def queue_payload() -> list[dict]: def adopt(run_id: str, reviewer: str, reason: str) -> int: - """采纳走 flow 合同:只记 reviews.action=adopt,不直接写正文文件。""" + """采纳必须调用 muse.flow,不在 web 层直接写 reviews。""" - return add_review( - run_id=run_id, - target="candidate", - action="adopt", - reviewer=reviewer, - reason=reason, - ) + return adopt_candidate(run_id, reviewer, reason) def revise(run_id: str, reviewer: str, reason: str, before_text: str, after_text: str) -> int: