ApplyLobsterCoreTablesMigration.java 8.3 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184
  1. import java.io.IOException;
  2. import java.nio.charset.StandardCharsets;
  3. import java.nio.file.Files;
  4. import java.nio.file.Path;
  5. import java.sql.*;
  6. import java.util.ArrayList;
  7. import java.util.List;
  8. import java.util.regex.Matcher;
  9. import java.util.regex.Pattern;
  10. /**
  11. * Apply all CREATE TABLE IF NOT EXISTS lobster_* blocks from tenant-initTable-migration.sql,
  12. * then patch known column gaps (e.g. lobster_node_execution_log.node_code).
  13. *
  14. * Usage: java -cp mysql-connector-j.jar;. ApplyLobsterCoreTablesMigration [databaseName]
  15. */
  16. public class ApplyLobsterCoreTablesMigration {
  17. private static final String HOST = getenv("DB_HOST", "cq-cdb-8fjmemkb.sql.tencentcdb.com");
  18. private static final int PORT = Integer.parseInt(getenv("DB_PORT", "27220"));
  19. private static final String USER = getenv("DB_USER", "root");
  20. private static final String PASS = getenv("DB_PASSWORD", "Ylrz_1q2w3e4r5t6y");
  21. public static void main(String[] args) throws Exception {
  22. String database = args.length > 0 ? args[0] : "ylrz_saas";
  23. Path migration = Path.of("d:/ylrz_saas_new/java/fs-service/src/main/resources/db/tenant-initTable-migration.sql");
  24. if (!Files.exists(migration)) {
  25. migration = Path.of("../fs-service/src/main/resources/db/tenant-initTable-migration.sql");
  26. }
  27. List<String> statements = extractLobsterCreateStatements(Files.readString(migration, StandardCharsets.UTF_8));
  28. System.out.println("Found " + statements.size() + " lobster CREATE TABLE statements");
  29. String url = "jdbc:mysql://" + HOST + ":" + PORT + "/" + database
  30. + "?useUnicode=true&characterEncoding=utf8&serverTimezone=GMT%2B8&allowMultiQueries=false";
  31. try (Connection conn = DriverManager.getConnection(url, USER, PASS)) {
  32. int ok = 0;
  33. int skip = 0;
  34. for (String sql : statements) {
  35. String table = extractTableName(sql);
  36. try (Statement st = conn.createStatement()) {
  37. st.execute(sql);
  38. System.out.println("[OK] " + table);
  39. ok++;
  40. } catch (SQLException e) {
  41. String msg = e.getMessage() != null ? e.getMessage() : "";
  42. if (msg.contains("already exists")) {
  43. System.out.println("[SKIP exists] " + table);
  44. skip++;
  45. } else {
  46. System.err.println("[FAIL] " + table + ": " + msg);
  47. }
  48. }
  49. }
  50. patchNodeExecutionLog(conn);
  51. ensureToolExecLog(conn);
  52. auditLobsterTables(conn);
  53. System.out.println("Done: created/verified=" + ok + ", skipped=" + skip);
  54. }
  55. }
  56. private static List<String> extractLobsterCreateStatements(String sql) {
  57. Pattern p = Pattern.compile(
  58. "CREATE TABLE IF NOT EXISTS `lobster_[^`]+`\\s*\\([\\s\\S]*?\\)\\s*ENGINE=InnoDB[^;]*;",
  59. Pattern.CASE_INSENSITIVE);
  60. Matcher m = p.matcher(sql);
  61. List<String> out = new ArrayList<>();
  62. while (m.find()) {
  63. out.add(m.group().trim());
  64. }
  65. return out;
  66. }
  67. private static String extractTableName(String createSql) {
  68. int start = createSql.indexOf('`') + 1;
  69. int end = createSql.indexOf('`', start);
  70. return end > start ? createSql.substring(start, end) : "unknown";
  71. }
  72. private static void ensureToolExecLog(Connection conn) throws SQLException {
  73. if (tableExists(conn, "lobster_tool_exec_log")) {
  74. return;
  75. }
  76. try (Statement st = conn.createStatement()) {
  77. st.execute("CREATE TABLE IF NOT EXISTS `lobster_tool_exec_log` ("
  78. + "`id` bigint NOT NULL AUTO_INCREMENT,"
  79. + "`company_id` bigint NOT NULL,"
  80. + "`tool_name` varchar(100) DEFAULT NULL,"
  81. + "`params_json` text,"
  82. + "`result_json` text,"
  83. + "`create_time` datetime DEFAULT CURRENT_TIMESTAMP,"
  84. + "PRIMARY KEY (`id`),"
  85. + "KEY `idx_company_id` (`company_id`)"
  86. + ") ENGINE=InnoDB DEFAULT CHARSET=utf8mb4");
  87. System.out.println("[PATCH] created lobster_tool_exec_log");
  88. }
  89. }
  90. private static void patchNodeExecutionLog(Connection conn) throws SQLException {
  91. if (!tableExists(conn, "lobster_node_execution_log")) {
  92. System.out.println("[PATCH skip] lobster_node_execution_log missing");
  93. return;
  94. }
  95. addColumnIfMissing(conn, "lobster_node_execution_log", "node_code",
  96. "ALTER TABLE lobster_node_execution_log ADD COLUMN node_code varchar(100) DEFAULT NULL AFTER workflow_id");
  97. addColumnIfMissing(conn, "lobster_node_execution_log", "quality_score",
  98. "ALTER TABLE lobster_node_execution_log ADD COLUMN quality_score int DEFAULT NULL AFTER retry_count");
  99. addIndexIfMissing(conn, "lobster_node_execution_log", "idx_node_code",
  100. "ALTER TABLE lobster_node_execution_log ADD KEY idx_node_code (company_id, node_code)");
  101. }
  102. private static void auditLobsterTables(Connection conn) throws SQLException {
  103. try (Statement st = conn.createStatement();
  104. ResultSet rs = st.executeQuery(
  105. "SELECT COUNT(*) FROM information_schema.TABLES "
  106. + "WHERE TABLE_SCHEMA=DATABASE() AND TABLE_NAME LIKE 'lobster_%'")) {
  107. rs.next();
  108. System.out.println("[AUDIT] lobster_* tables: " + rs.getInt(1));
  109. }
  110. String[] required = {
  111. "lobster_execution_config", "lobster_tool_config", "lobster_workflow_instance",
  112. "lobster_node_execution_log", "lobster_message_delivery_log", "lobster_evolution_log"
  113. };
  114. for (String table : required) {
  115. System.out.println(" " + table + ": " + (tableExists(conn, table) ? "OK" : "MISSING"));
  116. }
  117. }
  118. private static void addColumnIfMissing(Connection conn, String table, String column, String ddl) throws SQLException {
  119. if (columnExists(conn, table, column)) {
  120. return;
  121. }
  122. try (Statement st = conn.createStatement()) {
  123. st.execute(ddl);
  124. System.out.println("[PATCH] added " + table + "." + column);
  125. }
  126. }
  127. private static void addIndexIfMissing(Connection conn, String table, String index, String ddl) throws SQLException {
  128. if (indexExists(conn, table, index)) {
  129. return;
  130. }
  131. try (Statement st = conn.createStatement()) {
  132. st.execute(ddl);
  133. System.out.println("[PATCH] added index " + table + "." + index);
  134. }
  135. }
  136. private static boolean tableExists(Connection conn, String table) throws SQLException {
  137. try (PreparedStatement ps = conn.prepareStatement(
  138. "SELECT COUNT(*) FROM information_schema.TABLES WHERE TABLE_SCHEMA=DATABASE() AND TABLE_NAME=?")) {
  139. ps.setString(1, table);
  140. try (ResultSet rs = ps.executeQuery()) {
  141. return rs.next() && rs.getInt(1) > 0;
  142. }
  143. }
  144. }
  145. private static boolean columnExists(Connection conn, String table, String column) throws SQLException {
  146. try (PreparedStatement ps = conn.prepareStatement(
  147. "SELECT COUNT(*) FROM information_schema.COLUMNS WHERE TABLE_SCHEMA=DATABASE() AND TABLE_NAME=? AND COLUMN_NAME=?")) {
  148. ps.setString(1, table);
  149. ps.setString(2, column);
  150. try (ResultSet rs = ps.executeQuery()) {
  151. return rs.next() && rs.getInt(1) > 0;
  152. }
  153. }
  154. }
  155. private static boolean indexExists(Connection conn, String table, String index) throws SQLException {
  156. try (PreparedStatement ps = conn.prepareStatement(
  157. "SELECT COUNT(*) FROM information_schema.STATISTICS WHERE TABLE_SCHEMA=DATABASE() AND TABLE_NAME=? AND INDEX_NAME=?")) {
  158. ps.setString(1, table);
  159. ps.setString(2, index);
  160. try (ResultSet rs = ps.executeQuery()) {
  161. return rs.next() && rs.getInt(1) > 0;
  162. }
  163. }
  164. }
  165. private static String getenv(String key, String defaultValue) {
  166. String v = System.getenv(key);
  167. return v != null && !v.isEmpty() ? v : defaultValue;
  168. }
  169. }