# -*- coding: utf-8 -*- """ Lobster Phase 2: migrate + verify replay archive / hot buffer / GEPA API. Usage (PowerShell): cd d:\\ylrz_saas_new\\java\\scripts $env:LOBSTER_DB_HOST="your-host" $env:LOBSTER_DB_PORT="27220" $env:LOBSTER_DB_USER="root" $env:LOBSTER_DB_PASSWORD="your-password" $env:LOBSTER_TENANT_DB="ylrz_saas_tenant_1" $env:LOBSTER_COMPANY_ID="1" python verify_lobster_phase2.py --migrate python verify_lobster_phase2.py python verify_lobster_phase2.py --api-base http://127.0.0.1:8006 --token "Bearer xxx" Or reuse same DB vars as verify_test.py (edit host/password there and pass --use-verify-test-db). """ from __future__ import annotations import argparse import json import os import sys import urllib.error import urllib.request try: import pymysql except ImportError: print("FAIL: pip install pymysql") sys.exit(1) SQL_MIGRATION = """ CREATE TABLE IF NOT EXISTS `lobster_learning_replay_archive` ( `id` bigint NOT NULL AUTO_INCREMENT, `company_id` bigint NOT NULL, `instance_id` bigint DEFAULT NULL, `node_code` varchar(100) DEFAULT NULL, `customer_message` text, `ai_reply` text, `quality_score` double DEFAULT NULL, `create_time` datetime DEFAULT NULL, `archived_time` datetime DEFAULT CURRENT_TIMESTAMP, PRIMARY KEY (`id`), KEY `idx_company_archived` (`company_id`, `archived_time`), KEY `idx_company_quality` (`company_id`, `quality_score`) ) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COLLATE=utf8mb4_0900_ai_ci; """ INDEX_NAME = "idx_company_quality_time" def load_db_config(args): if args.use_verify_test_db: return { "host": "cq-cdb-8fjmemkb.sql.tencentcdb.com", "port": 27220, "user": "root", "password": "Ylrz_1q2w3e4r5t6y", "database": args.tenant_db or "ylrz_saas_tenant_1", } return { "host": os.environ.get("LOBSTER_DB_HOST", "127.0.0.1"), "port": int(os.environ.get("LOBSTER_DB_PORT", "3306")), "user": os.environ.get("LOBSTER_DB_USER", "root"), "password": os.environ.get("LOBSTER_DB_PASSWORD", ""), "database": args.tenant_db or os.environ.get("LOBSTER_TENANT_DB", "ylrz_saas_tenant_1"), } def connect(cfg): return pymysql.connect( host=cfg["host"], port=cfg["port"], user=cfg["user"], password=cfg["password"], database=cfg["database"], charset="utf8mb4", connect_timeout=15, read_timeout=30, ) def run_migration(conn): cur = conn.cursor() cur.execute(SQL_MIGRATION) conn.commit() cur.execute( "SELECT COUNT(*) FROM information_schema.statistics " "WHERE table_schema = DATABASE() AND table_name = 'lobster_learning_replay_buffer' " "AND index_name = %s", (INDEX_NAME,), ) if cur.fetchone()[0] == 0: try: cur.execute( f"ALTER TABLE lobster_learning_replay_buffer " f"ADD INDEX `{INDEX_NAME}` (`company_id`, `quality_score`, `create_time`)" ) conn.commit() print("OK added index", INDEX_NAME) except pymysql.err.OperationalError as e: if e.args[0] == 1061: print("SKIP index already exists (1061)") else: raise else: print("OK index exists:", INDEX_NAME) print("OK migration done on", conn.db.decode() if isinstance(conn.db, bytes) else conn.db) def verify_schema(conn): cur = conn.cursor() checks = [] for table in ("lobster_learning_replay_buffer", "lobster_learning_replay_archive"): cur.execute( "SELECT COUNT(*) FROM information_schema.tables " "WHERE table_schema = DATABASE() AND table_name = %s", (table,), ) ok = cur.fetchone()[0] == 1 checks.append((f"table {table}", ok)) cur.execute( "SELECT COUNT(*) FROM information_schema.statistics " "WHERE table_schema = DATABASE() AND table_name = 'lobster_learning_replay_buffer' " "AND index_name = %s", (INDEX_NAME,), ) checks.append((f"index {INDEX_NAME}", cur.fetchone()[0] == 1)) return checks def verify_data(conn, company_id): cur = conn.cursor() cur.execute( "SELECT COUNT(*), COALESCE(MAX(quality_score),0), COALESCE(MIN(create_time),'') " "FROM lobster_learning_replay_buffer WHERE company_id = %s", (company_id,), ) hot_count, hot_max_q, hot_min_t = cur.fetchone() cur.execute( "SELECT COUNT(*) FROM lobster_learning_replay_archive WHERE company_id = %s", (company_id,), ) archive_count = cur.fetchone()[0] cur.execute( "SELECT COUNT(*) FROM lobster_learned_pattern " "WHERE company_id = %s AND pattern_type IN ('skill_staged','skill')", (company_id,), ) skill_count = cur.fetchone()[0] cur.execute( "SELECT COUNT(*) FROM lobster_learning_event_log WHERE company_id = %s", (company_id,), ) event_count = cur.fetchone()[0] cur.execute( "SELECT COUNT(*) FROM lobster_learned_pattern " "WHERE company_id = %s AND source LIKE %s", (company_id, "%GepaEvolution%"), ) gepa_staged = cur.fetchone()[0] return { "hot_replay_rows": hot_count, "hot_max_quality": hot_max_q, "hot_oldest_time": str(hot_min_t), "archive_rows": archive_count, "skill_patterns": skill_count, "event_log_rows": event_count, "gepa_staged_patterns": gepa_staged, } def call_gepa_api(base_url, token, company_id): url = f"{base_url.rstrip('/')}/api/lobster/admin/learning/gepa/{company_id}" req = urllib.request.Request(url, method="POST") req.add_header("Authorization", token if token.startswith("Bearer") else f"Bearer {token}") req.add_header("Content-Type", "application/json") try: with urllib.request.urlopen(req, timeout=60) as resp: body = resp.read().decode("utf-8") return resp.status, json.loads(body) if body else {} except urllib.error.HTTPError as e: return e.code, e.read().decode("utf-8", errors="replace") except urllib.error.URLError as e: return 0, str(e.reason) def main(): parser = argparse.ArgumentParser(description="Lobster Phase2 migrate & verify") parser.add_argument("--migrate", action="store_true", help="Run SQL migration on tenant DB") parser.add_argument("--tenant-db", default=None, help="Tenant database name") parser.add_argument("--company-id", type=int, default=int(os.environ.get("LOBSTER_COMPANY_ID", "1"))) parser.add_argument("--use-verify-test-db", action="store_true", help="Use same remote DB as verify_test.py") parser.add_argument("--api-base", default=os.environ.get("LOBSTER_API_BASE", "http://127.0.0.1:8006")) parser.add_argument("--token", default=os.environ.get("LOBSTER_API_TOKEN", "")) parser.add_argument("--gepa", action="store_true", help="POST GEPA trigger API after DB checks") args = parser.parse_args() cfg = load_db_config(args) print("=== Lobster Phase2 verify ===") print("DB:", cfg["host"], cfg["port"], cfg["database"], "company_id=", args.company_id) conn = connect(cfg) try: if args.migrate: run_migration(conn) print("\n--- Schema ---") all_ok = True for name, ok in verify_schema(conn): print(("OK " if ok else "FAIL") + " " + name) all_ok = all_ok and ok print("\n--- Data ---") stats = verify_data(conn, args.company_id) for k, v in stats.items(): print(f" {k}: {v}") hot_limit = 500 if stats["hot_replay_rows"] > hot_limit: print(f"WARN hot buffer {stats['hot_replay_rows']} > {hot_limit} (trim runs on insert; check async writer)") finally: conn.close() if args.gepa: if not args.token: print("\nSKIP GEPA API: set --token or LOBSTER_API_TOKEN") else: print("\n--- GEPA API ---") status, body = call_gepa_api(args.api_base, args.token, args.company_id) print("HTTP", status, body) if status == 200 and isinstance(body, dict): print("skillsEvolved:", body.get("skillsEvolved", body)) print("\n=== Done ===") if not all_ok: sys.exit(2) sys.exit(0) if __name__ == "__main__": main()