import java.io.IOException; import java.nio.charset.StandardCharsets; import java.nio.file.Files; import java.nio.file.Path; import java.sql.*; import java.util.ArrayList; import java.util.List; import java.util.regex.Matcher; import java.util.regex.Pattern; /** * Apply all CREATE TABLE IF NOT EXISTS lobster_* blocks from tenant-initTable-migration.sql, * then patch known column gaps (e.g. lobster_node_execution_log.node_code). * * Usage: java -cp mysql-connector-j.jar;. ApplyLobsterCoreTablesMigration [databaseName] */ public class ApplyLobsterCoreTablesMigration { private static final String HOST = getenv("DB_HOST", "cq-cdb-8fjmemkb.sql.tencentcdb.com"); private static final int PORT = Integer.parseInt(getenv("DB_PORT", "27220")); private static final String USER = getenv("DB_USER", "root"); private static final String PASS = getenv("DB_PASSWORD", "Ylrz_1q2w3e4r5t6y"); public static void main(String[] args) throws Exception { String database = args.length > 0 ? args[0] : "ylrz_saas"; Path migration = Path.of("d:/ylrz_saas_new/java/fs-service/src/main/resources/db/tenant-initTable-migration.sql"); if (!Files.exists(migration)) { migration = Path.of("../fs-service/src/main/resources/db/tenant-initTable-migration.sql"); } List statements = extractLobsterCreateStatements(Files.readString(migration, StandardCharsets.UTF_8)); System.out.println("Found " + statements.size() + " lobster CREATE TABLE statements"); String url = "jdbc:mysql://" + HOST + ":" + PORT + "/" + database + "?useUnicode=true&characterEncoding=utf8&serverTimezone=GMT%2B8&allowMultiQueries=false"; try (Connection conn = DriverManager.getConnection(url, USER, PASS)) { int ok = 0; int skip = 0; for (String sql : statements) { String table = extractTableName(sql); try (Statement st = conn.createStatement()) { st.execute(sql); System.out.println("[OK] " + table); ok++; } catch (SQLException e) { String msg = e.getMessage() != null ? e.getMessage() : ""; if (msg.contains("already exists")) { System.out.println("[SKIP exists] " + table); skip++; } else { System.err.println("[FAIL] " + table + ": " + msg); } } } patchNodeExecutionLog(conn); ensureToolExecLog(conn); auditLobsterTables(conn); System.out.println("Done: created/verified=" + ok + ", skipped=" + skip); } } private static List extractLobsterCreateStatements(String sql) { Pattern p = Pattern.compile( "CREATE TABLE IF NOT EXISTS `lobster_[^`]+`\\s*\\([\\s\\S]*?\\)\\s*ENGINE=InnoDB[^;]*;", Pattern.CASE_INSENSITIVE); Matcher m = p.matcher(sql); List out = new ArrayList<>(); while (m.find()) { out.add(m.group().trim()); } return out; } private static String extractTableName(String createSql) { int start = createSql.indexOf('`') + 1; int end = createSql.indexOf('`', start); return end > start ? createSql.substring(start, end) : "unknown"; } private static void ensureToolExecLog(Connection conn) throws SQLException { if (tableExists(conn, "lobster_tool_exec_log")) { return; } try (Statement st = conn.createStatement()) { st.execute("CREATE TABLE IF NOT EXISTS `lobster_tool_exec_log` (" + "`id` bigint NOT NULL AUTO_INCREMENT," + "`company_id` bigint NOT NULL," + "`tool_name` varchar(100) DEFAULT NULL," + "`params_json` text," + "`result_json` text," + "`create_time` datetime DEFAULT CURRENT_TIMESTAMP," + "PRIMARY KEY (`id`)," + "KEY `idx_company_id` (`company_id`)" + ") ENGINE=InnoDB DEFAULT CHARSET=utf8mb4"); System.out.println("[PATCH] created lobster_tool_exec_log"); } } private static void patchNodeExecutionLog(Connection conn) throws SQLException { if (!tableExists(conn, "lobster_node_execution_log")) { System.out.println("[PATCH skip] lobster_node_execution_log missing"); return; } addColumnIfMissing(conn, "lobster_node_execution_log", "node_code", "ALTER TABLE lobster_node_execution_log ADD COLUMN node_code varchar(100) DEFAULT NULL AFTER workflow_id"); addColumnIfMissing(conn, "lobster_node_execution_log", "quality_score", "ALTER TABLE lobster_node_execution_log ADD COLUMN quality_score int DEFAULT NULL AFTER retry_count"); addIndexIfMissing(conn, "lobster_node_execution_log", "idx_node_code", "ALTER TABLE lobster_node_execution_log ADD KEY idx_node_code (company_id, node_code)"); } private static void auditLobsterTables(Connection conn) throws SQLException { try (Statement st = conn.createStatement(); ResultSet rs = st.executeQuery( "SELECT COUNT(*) FROM information_schema.TABLES " + "WHERE TABLE_SCHEMA=DATABASE() AND TABLE_NAME LIKE 'lobster_%'")) { rs.next(); System.out.println("[AUDIT] lobster_* tables: " + rs.getInt(1)); } String[] required = { "lobster_execution_config", "lobster_tool_config", "lobster_workflow_instance", "lobster_node_execution_log", "lobster_message_delivery_log", "lobster_evolution_log" }; for (String table : required) { System.out.println(" " + table + ": " + (tableExists(conn, table) ? "OK" : "MISSING")); } } private static void addColumnIfMissing(Connection conn, String table, String column, String ddl) throws SQLException { if (columnExists(conn, table, column)) { return; } try (Statement st = conn.createStatement()) { st.execute(ddl); System.out.println("[PATCH] added " + table + "." + column); } } private static void addIndexIfMissing(Connection conn, String table, String index, String ddl) throws SQLException { if (indexExists(conn, table, index)) { return; } try (Statement st = conn.createStatement()) { st.execute(ddl); System.out.println("[PATCH] added index " + table + "." + index); } } private static boolean tableExists(Connection conn, String table) throws SQLException { try (PreparedStatement ps = conn.prepareStatement( "SELECT COUNT(*) FROM information_schema.TABLES WHERE TABLE_SCHEMA=DATABASE() AND TABLE_NAME=?")) { ps.setString(1, table); try (ResultSet rs = ps.executeQuery()) { return rs.next() && rs.getInt(1) > 0; } } } private static boolean columnExists(Connection conn, String table, String column) throws SQLException { try (PreparedStatement ps = conn.prepareStatement( "SELECT COUNT(*) FROM information_schema.COLUMNS WHERE TABLE_SCHEMA=DATABASE() AND TABLE_NAME=? AND COLUMN_NAME=?")) { ps.setString(1, table); ps.setString(2, column); try (ResultSet rs = ps.executeQuery()) { return rs.next() && rs.getInt(1) > 0; } } } private static boolean indexExists(Connection conn, String table, String index) throws SQLException { try (PreparedStatement ps = conn.prepareStatement( "SELECT COUNT(*) FROM information_schema.STATISTICS WHERE TABLE_SCHEMA=DATABASE() AND TABLE_NAME=? AND INDEX_NAME=?")) { ps.setString(1, table); ps.setString(2, index); try (ResultSet rs = ps.executeQuery()) { return rs.next() && rs.getInt(1) > 0; } } } private static String getenv(String key, String defaultValue) { String v = System.getenv(key); return v != null && !v.isEmpty() ? v : defaultValue; } }