validate_db.py 2.8 KB

12345678910111213141516171819202122232425262728293031323334353637383940414243444546474849505152535455565758596061626364656667686970717273747576777879
  1. """M1 校验:连库、确认 schema/表/列齐全、做一次可逆写探测。
  2. 用法(在能连到 RDS 的机器上,如海外开发机):
  3. python scripts/validate_db.py [path/to/.env]
  4. """
  5. from __future__ import annotations
  6. import sys
  7. from creation_knowledge.config import PgConfig
  8. from creation_knowledge.integrations.db import CkStore, _connect
  9. from creation_knowledge.models import Post
  10. EXPECTED = {
  11. "ck_post": {"id", "platform", "url", "raw", "extracted", "screening",
  12. "stage", "created_at", "updated_at"},
  13. "ck_knowledge_item": {"id", "post_id", "item", "deconstruction",
  14. "ingest_payload", "ingest_status", "knowledge_id",
  15. "created_at", "updated_at"},
  16. }
  17. PROBE_ID = "_validate_probe"
  18. def main() -> int:
  19. env_file = sys.argv[1] if len(sys.argv) > 1 else ".env"
  20. cfg = PgConfig.from_env(env_file)
  21. print(f"[cfg] host={cfg.host} db={cfg.database} schema={cfg.schema} user={cfg.user}")
  22. # 1) 连通 + 版本
  23. with _connect(cfg) as conn:
  24. with conn.cursor() as cur:
  25. cur.execute("SELECT version()")
  26. print("[conn] OK:", cur.fetchone()[0][:70])
  27. store = CkStore(cfg)
  28. # 2) 表/列齐全
  29. ok = True
  30. for table, expected_cols in EXPECTED.items():
  31. cols = set(store.table_columns(table))
  32. if not cols:
  33. print(f"[schema] MISSING table: {cfg.schema}.{table}")
  34. ok = False
  35. continue
  36. missing = expected_cols - cols
  37. if missing:
  38. print(f"[schema] {table} 缺列: {sorted(missing)}")
  39. ok = False
  40. else:
  41. print(f"[schema] {table} OK ({len(cols)} cols)")
  42. if not ok:
  43. print("RESULT: FAIL(表结构不符)")
  44. return 1
  45. # 3) 可逆写探测:插入一帖+一片段,读回,删除
  46. store.delete_post(PROBE_ID) # 清理上次残留
  47. try:
  48. store.upsert_post(Post(id=PROBE_ID, url="probe://x", content_id="x",
  49. raw={"probe": True}))
  50. store.set_extracted(PROBE_ID, {"text": "probe", "is_empty": False})
  51. item_id = store.save_item(PROBE_ID, {"title": "probe item",
  52. "knowledge_types": ["how"]},
  53. {"stages": ["脚本"]}, {"title": "probe"})
  54. post = store.read_post(PROBE_ID)
  55. items = store.read_items(PROBE_ID)
  56. assert post and post["stage"] == "extracted", post
  57. assert items and items[0]["id"] == item_id, items
  58. print(f"[write] OK: post.stage={post['stage']}, item_id={item_id}, "
  59. f"ingest_status={items[0]['ingest_status']}")
  60. finally:
  61. store.delete_post(PROBE_ID)
  62. print("[write] 探测数据已清理")
  63. print("RESULT: PASS")
  64. return 0
  65. if __name__ == "__main__":
  66. raise SystemExit(main())