#!/usr/bin/env python3 """humanization 规则/样例种子同步:Git YAML(迁移种子)→ muse-example 运行时权威表。 合同: - 种子前先过文件侧激活门(load.load_rules(samples=...)):样例不齐的规则拒绝进库。 - 单事务:先比对 content_sha256 生成计划,再执行 upsert;任何错误整批回滚。 - 幂等:内容哈希一致的行跳过;内容变化的行 UPDATE,并记一条 synced 生命周期事件。 - 数据库独有行不自动删除:默认报告清单;--strict 时失败关闭,交人工裁决。 - --dry-run 只输出目标清单,不连接数据库。 """ from __future__ import annotations import argparse import json import sys from pathlib import Path TOOL_DIR = Path(__file__).resolve().parent AGENT_ROOT = TOOL_DIR.parent.parent for _p in ( AGENT_ROOT / "humanization" / "src", ): if str(_p) not in sys.path: sys.path.insert(0, str(_p)) from deai import load, load_db # noqa: E402 class SeedError(ValueError): pass _RULE_INSERT = ( "INSERT INTO example_ai_flavor_rule " "(rule_id, name, layer, carrier_scope, default_disposition, status, version, " "trigger_json, carve_out, function_check, sample_refs, case_card_ids, " "fix_hint, evidence, payload, content_sha256, creator, tenant_id) " "VALUES (%s, %s, %s, %s, %s, %s, %s, %s::jsonb, %s::jsonb, %s::jsonb, %s::jsonb, %s::jsonb, " "%s, %s, %s::jsonb, %s, %s, %s)" ) _RULE_UPDATE = ( "UPDATE example_ai_flavor_rule SET name=%s, layer=%s, carrier_scope=%s, " "default_disposition=%s, status=%s, version=%s, trigger_json=%s::jsonb, carve_out=%s::jsonb, " "function_check=%s::jsonb, sample_refs=%s::jsonb, case_card_ids=%s::jsonb, fix_hint=%s, " "evidence=%s, payload=%s::jsonb, content_sha256=%s, updater=%s " "WHERE tenant_id=%s AND rule_id=%s" ) _SAMPLE_INSERT = ( "INSERT INTO example_ai_flavor_sample " "(sample_id, sample_type, carrier, source, source_license, rules, text, note, " "case_card_id, source_ref, payload, content_sha256, creator, tenant_id) " "VALUES (%s, %s, %s, %s, %s, %s::jsonb, %s, %s, %s, %s, %s::jsonb, %s, %s, %s)" ) _SAMPLE_UPDATE = ( "UPDATE example_ai_flavor_sample SET sample_type=%s, carrier=%s, source=%s, " "source_license=%s, rules=%s::jsonb, text=%s, note=%s, case_card_id=%s, source_ref=%s, " "payload=%s::jsonb, content_sha256=%s, updater=%s " "WHERE tenant_id=%s AND sample_id=%s" ) _RULE_EVENT_INSERT = ( "INSERT INTO example_ai_flavor_rule_event " "(rule_id, event, rule_version, status, approver, note, content_sha256, creator, tenant_id) " "VALUES (%s, 'synced', %s, %s, '', %s, %s, %s, %s)" ) def _dumps(value) -> str: return json.dumps(value, ensure_ascii=False) def _rule_insert_params(row: dict, creator: str, tenant_id: int) -> tuple: return ( row["rule_id"], row["name"], row["layer"], row["carrier_scope"], row["default_disposition"], row["status"], row["version"], _dumps(row["trigger_json"]), _dumps(row["carve_out"]), _dumps(row["function_check"]), _dumps(row["sample_refs"]), _dumps(row["case_card_ids"]), row["fix_hint"], row["evidence"], _dumps(row["payload"]), row["content_sha256"], creator, tenant_id, ) def _rule_update_params(row: dict, creator: str, tenant_id: int) -> tuple: return ( row["name"], row["layer"], row["carrier_scope"], row["default_disposition"], row["status"], row["version"], _dumps(row["trigger_json"]), _dumps(row["carve_out"]), _dumps(row["function_check"]), _dumps(row["sample_refs"]), _dumps(row["case_card_ids"]), row["fix_hint"], row["evidence"], _dumps(row["payload"]), row["content_sha256"], creator, tenant_id, row["rule_id"], ) def _sample_insert_params(row: dict, creator: str, tenant_id: int) -> tuple: return ( row["sample_id"], row["sample_type"], row["carrier"], row["source"], row["source_license"], _dumps(row["rules"]), row["text"], row["note"], row["case_card_id"], row["source_ref"], _dumps(row["payload"]), row["content_sha256"], creator, tenant_id, ) def _sample_update_params(row: dict, creator: str, tenant_id: int) -> tuple: return ( row["sample_type"], row["carrier"], row["source"], row["source_license"], _dumps(row["rules"]), row["text"], row["note"], row["case_card_id"], row["source_ref"], _dumps(row["payload"]), row["content_sha256"], creator, tenant_id, row["sample_id"], ) def seed(conn, *, rules: dict, samples: dict, creator: str = "1", tenant_id: int = load_db.TENANT_ID, strict: bool = False) -> dict: """单事务同步规则与样例;调用方负责事务边界(真实路径用 conn.transaction())。""" summary = { "rules": {"inserted": 0, "updated": 0, "unchanged": 0}, "samples": {"inserted": 0, "updated": 0, "unchanged": 0}, "db_only_rules": [], "db_only_samples": [], "events": 0, } existing_rules = { row[0]: {"content_sha256": row[1]} for row in conn.execute( "SELECT rule_id, content_sha256 FROM example_ai_flavor_rule " "WHERE tenant_id = %s AND deleted = FALSE", (tenant_id,) ).fetchall() } existing_samples = { row[0]: row[1] for row in conn.execute( "SELECT sample_id, content_sha256 FROM example_ai_flavor_sample " "WHERE tenant_id = %s AND deleted = FALSE", (tenant_id,) ).fetchall() } for rule in sorted(rules.values(), key=lambda item: item["id"]): row = load_db.rule_row(rule) current = existing_rules.get(row["rule_id"]) if current is None: conn.execute(_RULE_INSERT, _rule_insert_params(row, creator, tenant_id)) conn.execute(_RULE_EVENT_INSERT, ( row["rule_id"], row["version"], row["status"], "seed insert", row["content_sha256"], creator, tenant_id, )) summary["rules"]["inserted"] += 1 summary["events"] += 1 elif current["content_sha256"] != row["content_sha256"]: conn.execute(_RULE_UPDATE, _rule_update_params(row, creator, tenant_id)) conn.execute(_RULE_EVENT_INSERT, ( row["rule_id"], row["version"], row["status"], "seed update", row["content_sha256"], creator, tenant_id, )) summary["rules"]["updated"] += 1 summary["events"] += 1 else: summary["rules"]["unchanged"] += 1 for sample in sorted(samples.values(), key=lambda item: item["id"]): row = load_db.sample_row(sample) current_sha = existing_samples.get(row["sample_id"]) if current_sha is None: conn.execute(_SAMPLE_INSERT, _sample_insert_params(row, creator, tenant_id)) summary["samples"]["inserted"] += 1 elif current_sha != row["content_sha256"]: conn.execute(_SAMPLE_UPDATE, _sample_update_params(row, creator, tenant_id)) summary["samples"]["updated"] += 1 else: summary["samples"]["unchanged"] += 1 summary["db_only_rules"] = sorted(set(existing_rules) - set(rules)) summary["db_only_samples"] = sorted(set(existing_samples) - set(samples)) if strict and (summary["db_only_rules"] or summary["db_only_samples"]): raise SeedError( "数据库存在 YAML 种子之外的行,失败关闭: rules=" + ",".join(summary["db_only_rules"]) + "; samples=" + ",".join(summary["db_only_samples"]) ) return summary def main(argv: list[str] | None = None) -> int: parser = argparse.ArgumentParser(description="humanization 规则/样例种子同步(YAML → muse-example)") parser.add_argument("--dry-run", action="store_true", help="只输出目标清单,不连接数据库") parser.add_argument("--strict", action="store_true", help="数据库独有行存在时失败关闭") parser.add_argument("--tenant", type=int, default=load_db.TENANT_ID) args = parser.parse_args(argv) try: # 种子前先过激活门:样例不齐的 active 规则在文件侧就被拒绝 samples = load.load_samples() rules = load.load_rules(samples=samples) if args.dry_run: plan = { "status": "dry_run", "rules": len(rules), "samples": len(samples), "rule_ids": sorted(rules), "library_version": load.rule_library_version(rules), } print(json.dumps(plan, ensure_ascii=False)) return 0 from muse_db import connect # noqa: E402 with connect() as conn: with conn.transaction(): summary = seed(conn, rules=rules, samples=samples, tenant_id=args.tenant, strict=args.strict) print(json.dumps({"status": "seeded", **summary}, ensure_ascii=False)) return 0 except (SeedError, load.LoadError, ValueError, OSError) as exc: print(f"SEED_RULES_FAILED: {exc}") return 2 if __name__ == "__main__": raise SystemExit(main())