verify_lobster_phase2.py 8.3 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245
  1. # -*- coding: utf-8 -*-
  2. """
  3. Lobster Phase 2: migrate + verify replay archive / hot buffer / GEPA API.
  4. Usage (PowerShell):
  5. cd d:\\ylrz_saas_new\\java\\scripts
  6. $env:LOBSTER_DB_HOST="your-host"
  7. $env:LOBSTER_DB_PORT="27220"
  8. $env:LOBSTER_DB_USER="root"
  9. $env:LOBSTER_DB_PASSWORD="your-password"
  10. $env:LOBSTER_TENANT_DB="ylrz_saas_tenant_1"
  11. $env:LOBSTER_COMPANY_ID="1"
  12. python verify_lobster_phase2.py --migrate
  13. python verify_lobster_phase2.py
  14. python verify_lobster_phase2.py --api-base http://127.0.0.1:8006 --token "Bearer xxx"
  15. Or reuse same DB vars as verify_test.py (edit host/password there and pass --use-verify-test-db).
  16. """
  17. from __future__ import annotations
  18. import argparse
  19. import json
  20. import os
  21. import sys
  22. import urllib.error
  23. import urllib.request
  24. try:
  25. import pymysql
  26. except ImportError:
  27. print("FAIL: pip install pymysql")
  28. sys.exit(1)
  29. SQL_MIGRATION = """
  30. CREATE TABLE IF NOT EXISTS `lobster_learning_replay_archive` (
  31. `id` bigint NOT NULL AUTO_INCREMENT,
  32. `company_id` bigint NOT NULL,
  33. `instance_id` bigint DEFAULT NULL,
  34. `node_code` varchar(100) DEFAULT NULL,
  35. `customer_message` text,
  36. `ai_reply` text,
  37. `quality_score` double DEFAULT NULL,
  38. `create_time` datetime DEFAULT NULL,
  39. `archived_time` datetime DEFAULT CURRENT_TIMESTAMP,
  40. PRIMARY KEY (`id`),
  41. KEY `idx_company_archived` (`company_id`, `archived_time`),
  42. KEY `idx_company_quality` (`company_id`, `quality_score`)
  43. ) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COLLATE=utf8mb4_0900_ai_ci;
  44. """
  45. INDEX_NAME = "idx_company_quality_time"
  46. def load_db_config(args):
  47. if args.use_verify_test_db:
  48. return {
  49. "host": "cq-cdb-8fjmemkb.sql.tencentcdb.com",
  50. "port": 27220,
  51. "user": "root",
  52. "password": "Ylrz_1q2w3e4r5t6y",
  53. "database": args.tenant_db or "ylrz_saas_tenant_1",
  54. }
  55. return {
  56. "host": os.environ.get("LOBSTER_DB_HOST", "127.0.0.1"),
  57. "port": int(os.environ.get("LOBSTER_DB_PORT", "3306")),
  58. "user": os.environ.get("LOBSTER_DB_USER", "root"),
  59. "password": os.environ.get("LOBSTER_DB_PASSWORD", ""),
  60. "database": args.tenant_db or os.environ.get("LOBSTER_TENANT_DB", "ylrz_saas_tenant_1"),
  61. }
  62. def connect(cfg):
  63. return pymysql.connect(
  64. host=cfg["host"],
  65. port=cfg["port"],
  66. user=cfg["user"],
  67. password=cfg["password"],
  68. database=cfg["database"],
  69. charset="utf8mb4",
  70. connect_timeout=15,
  71. read_timeout=30,
  72. )
  73. def run_migration(conn):
  74. cur = conn.cursor()
  75. cur.execute(SQL_MIGRATION)
  76. conn.commit()
  77. cur.execute(
  78. "SELECT COUNT(*) FROM information_schema.statistics "
  79. "WHERE table_schema = DATABASE() AND table_name = 'lobster_learning_replay_buffer' "
  80. "AND index_name = %s",
  81. (INDEX_NAME,),
  82. )
  83. if cur.fetchone()[0] == 0:
  84. try:
  85. cur.execute(
  86. f"ALTER TABLE lobster_learning_replay_buffer "
  87. f"ADD INDEX `{INDEX_NAME}` (`company_id`, `quality_score`, `create_time`)"
  88. )
  89. conn.commit()
  90. print("OK added index", INDEX_NAME)
  91. except pymysql.err.OperationalError as e:
  92. if e.args[0] == 1061:
  93. print("SKIP index already exists (1061)")
  94. else:
  95. raise
  96. else:
  97. print("OK index exists:", INDEX_NAME)
  98. print("OK migration done on", conn.db.decode() if isinstance(conn.db, bytes) else conn.db)
  99. def verify_schema(conn):
  100. cur = conn.cursor()
  101. checks = []
  102. for table in ("lobster_learning_replay_buffer", "lobster_learning_replay_archive"):
  103. cur.execute(
  104. "SELECT COUNT(*) FROM information_schema.tables "
  105. "WHERE table_schema = DATABASE() AND table_name = %s",
  106. (table,),
  107. )
  108. ok = cur.fetchone()[0] == 1
  109. checks.append((f"table {table}", ok))
  110. cur.execute(
  111. "SELECT COUNT(*) FROM information_schema.statistics "
  112. "WHERE table_schema = DATABASE() AND table_name = 'lobster_learning_replay_buffer' "
  113. "AND index_name = %s",
  114. (INDEX_NAME,),
  115. )
  116. checks.append((f"index {INDEX_NAME}", cur.fetchone()[0] == 1))
  117. return checks
  118. def verify_data(conn, company_id):
  119. cur = conn.cursor()
  120. cur.execute(
  121. "SELECT COUNT(*), COALESCE(MAX(quality_score),0), COALESCE(MIN(create_time),'') "
  122. "FROM lobster_learning_replay_buffer WHERE company_id = %s",
  123. (company_id,),
  124. )
  125. hot_count, hot_max_q, hot_min_t = cur.fetchone()
  126. cur.execute(
  127. "SELECT COUNT(*) FROM lobster_learning_replay_archive WHERE company_id = %s",
  128. (company_id,),
  129. )
  130. archive_count = cur.fetchone()[0]
  131. cur.execute(
  132. "SELECT COUNT(*) FROM lobster_learned_pattern "
  133. "WHERE company_id = %s AND pattern_type IN ('skill_staged','skill')",
  134. (company_id,),
  135. )
  136. skill_count = cur.fetchone()[0]
  137. cur.execute(
  138. "SELECT COUNT(*) FROM lobster_learning_event_log WHERE company_id = %s",
  139. (company_id,),
  140. )
  141. event_count = cur.fetchone()[0]
  142. cur.execute(
  143. "SELECT COUNT(*) FROM lobster_learned_pattern "
  144. "WHERE company_id = %s AND source LIKE %s",
  145. (company_id, "%GepaEvolution%"),
  146. )
  147. gepa_staged = cur.fetchone()[0]
  148. return {
  149. "hot_replay_rows": hot_count,
  150. "hot_max_quality": hot_max_q,
  151. "hot_oldest_time": str(hot_min_t),
  152. "archive_rows": archive_count,
  153. "skill_patterns": skill_count,
  154. "event_log_rows": event_count,
  155. "gepa_staged_patterns": gepa_staged,
  156. }
  157. def call_gepa_api(base_url, token, company_id):
  158. url = f"{base_url.rstrip('/')}/api/lobster/admin/learning/gepa/{company_id}"
  159. req = urllib.request.Request(url, method="POST")
  160. req.add_header("Authorization", token if token.startswith("Bearer") else f"Bearer {token}")
  161. req.add_header("Content-Type", "application/json")
  162. try:
  163. with urllib.request.urlopen(req, timeout=60) as resp:
  164. body = resp.read().decode("utf-8")
  165. return resp.status, json.loads(body) if body else {}
  166. except urllib.error.HTTPError as e:
  167. return e.code, e.read().decode("utf-8", errors="replace")
  168. except urllib.error.URLError as e:
  169. return 0, str(e.reason)
  170. def main():
  171. parser = argparse.ArgumentParser(description="Lobster Phase2 migrate & verify")
  172. parser.add_argument("--migrate", action="store_true", help="Run SQL migration on tenant DB")
  173. parser.add_argument("--tenant-db", default=None, help="Tenant database name")
  174. parser.add_argument("--company-id", type=int, default=int(os.environ.get("LOBSTER_COMPANY_ID", "1")))
  175. parser.add_argument("--use-verify-test-db", action="store_true", help="Use same remote DB as verify_test.py")
  176. parser.add_argument("--api-base", default=os.environ.get("LOBSTER_API_BASE", "http://127.0.0.1:8006"))
  177. parser.add_argument("--token", default=os.environ.get("LOBSTER_API_TOKEN", ""))
  178. parser.add_argument("--gepa", action="store_true", help="POST GEPA trigger API after DB checks")
  179. args = parser.parse_args()
  180. cfg = load_db_config(args)
  181. print("=== Lobster Phase2 verify ===")
  182. print("DB:", cfg["host"], cfg["port"], cfg["database"], "company_id=", args.company_id)
  183. conn = connect(cfg)
  184. try:
  185. if args.migrate:
  186. run_migration(conn)
  187. print("\n--- Schema ---")
  188. all_ok = True
  189. for name, ok in verify_schema(conn):
  190. print(("OK " if ok else "FAIL") + " " + name)
  191. all_ok = all_ok and ok
  192. print("\n--- Data ---")
  193. stats = verify_data(conn, args.company_id)
  194. for k, v in stats.items():
  195. print(f" {k}: {v}")
  196. hot_limit = 500
  197. if stats["hot_replay_rows"] > hot_limit:
  198. print(f"WARN hot buffer {stats['hot_replay_rows']} > {hot_limit} (trim runs on insert; check async writer)")
  199. finally:
  200. conn.close()
  201. if args.gepa:
  202. if not args.token:
  203. print("\nSKIP GEPA API: set --token or LOBSTER_API_TOKEN")
  204. else:
  205. print("\n--- GEPA API ---")
  206. status, body = call_gepa_api(args.api_base, args.token, args.company_id)
  207. print("HTTP", status, body)
  208. if status == 200 and isinstance(body, dict):
  209. print("skillsEvolved:", body.get("skillsEvolved", body))
  210. print("\n=== Done ===")
  211. if not all_ok:
  212. sys.exit(2)
  213. sys.exit(0)
  214. if __name__ == "__main__":
  215. main()