From 6a9cbaf3864ec4ad9709f006fecb091e0f8371d6 Mon Sep 17 00:00:00 2001 From: ViacheslavKlimov Date: Tue, 17 Sep 2024 13:49:23 +0300 Subject: [PATCH 1/3] Rollback upgrade if schema update failed, refactoring --- .../install/SqlDatabaseUpgradeService.java | 260 ++++++++---------- 1 file changed, 118 insertions(+), 142 deletions(-) diff --git a/application/src/main/java/org/thingsboard/server/service/install/SqlDatabaseUpgradeService.java b/application/src/main/java/org/thingsboard/server/service/install/SqlDatabaseUpgradeService.java index 203ab39781..9c611c40f7 100644 --- a/application/src/main/java/org/thingsboard/server/service/install/SqlDatabaseUpgradeService.java +++ b/application/src/main/java/org/thingsboard/server/service/install/SqlDatabaseUpgradeService.java @@ -16,24 +16,22 @@ package org.thingsboard.server.service.install; import lombok.extern.slf4j.Slf4j; -import org.springframework.beans.factory.annotation.Autowired; -import org.springframework.beans.factory.annotation.Value; +import org.intellij.lang.annotations.Language; import org.springframework.context.annotation.Profile; +import org.springframework.jdbc.core.JdbcTemplate; +import org.springframework.jdbc.core.StatementCallback; import org.springframework.stereotype.Service; +import org.springframework.transaction.PlatformTransactionManager; +import org.springframework.transaction.support.TransactionTemplate; import org.thingsboard.server.service.install.update.DefaultDataUpdateService; -import java.nio.charset.Charset; +import java.io.IOException; +import java.io.UncheckedIOException; import java.nio.file.Files; import java.nio.file.Path; import java.nio.file.Paths; -import java.sql.Connection; -import java.sql.DriverManager; -import java.sql.ResultSet; -import java.sql.SQLException; import java.sql.SQLWarning; -import java.sql.Statement; import java.util.concurrent.TimeUnit; -import java.util.function.Consumer; @Service @Profile("install") @@ -42,140 +40,126 @@ public class SqlDatabaseUpgradeService implements DatabaseEntitiesUpgradeService private static final String SCHEMA_UPDATE_SQL = "schema_update.sql"; - @Value("${spring.datasource.url}") - private String dbUrl; + private final InstallScripts installScripts; + private final JdbcTemplate jdbcTemplate; + private final TransactionTemplate transactionTemplate; - @Value("${spring.datasource.username}") - private String dbUserName; - - @Value("${spring.datasource.password}") - private String dbPassword; - - @Autowired - private InstallScripts installScripts; + public SqlDatabaseUpgradeService(InstallScripts installScripts, JdbcTemplate jdbcTemplate, PlatformTransactionManager transactionManager) { + this.installScripts = installScripts; + this.jdbcTemplate = jdbcTemplate; + this.transactionTemplate = new TransactionTemplate(transactionManager); + this.transactionTemplate.setTimeout((int) TimeUnit.MINUTES.toSeconds(120)); + } @Override - public void upgradeDatabase(String fromVersion) throws Exception { + public void upgradeDatabase(String fromVersion) { switch (fromVersion) { - case "3.5.0": - updateSchema("3.5.0", 3005000, "3.5.1", 3005001, null); - break; - case "3.5.1": - updateSchema("3.5.1", 3005001, "3.6.0", 3006000, conn -> { - String[] entityNames = new String[]{"device", "component_descriptor", "customer", "dashboard", "rule_chain", "rule_node", "ota_package", - "asset_profile", "asset", "device_profile", "tb_user", "tenant_profile", "tenant", "widgets_bundle", "entity_view", "edge"}; - for (String entityName : entityNames) { - try { - conn.createStatement().execute("ALTER TABLE " + entityName + " DROP COLUMN search_text CASCADE"); - } catch (Exception e) { - } - } - try { - conn.createStatement().execute("ALTER TABLE component_descriptor ADD COLUMN IF NOT EXISTS configuration_version int DEFAULT 0;"); - } catch (Exception e) { - } - try { - conn.createStatement().execute("ALTER TABLE rule_node ADD COLUMN IF NOT EXISTS configuration_version int DEFAULT 0;"); - } catch (Exception e) { - } - try { - conn.createStatement().execute("CREATE INDEX IF NOT EXISTS idx_rule_node_type_configuration_version ON rule_node(type, configuration_version);"); - } catch (Exception e) { - } - try { - conn.createStatement().execute("UPDATE rule_node SET " + + case "3.5.0" -> updateSchema("3.5.0", 3005000, "3.5.1", 3005001); + case "3.5.1" -> { + updateSchema("3.5.1", 3005001, "3.6.0", 3006000); + + String[] tables = new String[]{"device", "component_descriptor", "customer", "dashboard", "rule_chain", "rule_node", "ota_package", + "asset_profile", "asset", "device_profile", "tb_user", "tenant_profile", "tenant", "widgets_bundle", "entity_view", "edge"}; + for (String table : tables) { + execute("ALTER TABLE " + table + " DROP COLUMN IF EXISTS search_text CASCADE"); + } + execute( + "ALTER TABLE component_descriptor ADD COLUMN IF NOT EXISTS configuration_version int DEFAULT 0;", + "ALTER TABLE rule_node ADD COLUMN IF NOT EXISTS configuration_version int DEFAULT 0;", + "CREATE INDEX IF NOT EXISTS idx_rule_node_type_configuration_version ON rule_node(type, configuration_version);", + "UPDATE rule_node SET " + "configuration = (configuration::jsonb || '{\"updateAttributesOnlyOnValueChange\": \"false\"}'::jsonb)::varchar, " + "configuration_version = 1 " + - "WHERE type = 'org.thingsboard.rule.engine.telemetry.TbMsgAttributesNode' AND configuration_version < 1;"); - } catch (Exception e) { - } - try { - conn.createStatement().execute("CREATE INDEX IF NOT EXISTS idx_notification_recipient_id_unread ON notification(recipient_id) WHERE status <> 'READ';"); - } catch (Exception e) { - } - }); - break; - case "3.6.0": - updateSchema("3.6.0", 3006000, "3.6.1", 3006001, null); - break; - case "3.6.1": - updateSchema("3.6.1", 3006001, "3.6.2", 3006002, connection -> { - try { - Path saveAttributesNodeUpdateFile = Paths.get(installScripts.getDataDir(), "upgrade", "3.6.1", "save_attributes_node_update.sql"); - loadSql(saveAttributesNodeUpdateFile, connection); - } catch (Exception e) { - log.warn("Failed to execute update script for save attributes rule nodes due to: ", e); - } - try { - connection.createStatement().execute("CREATE INDEX IF NOT EXISTS idx_asset_profile_id ON asset(tenant_id, asset_profile_id);"); - } catch (Exception e) { - } - }); - break; - case "3.6.2": - updateSchema("3.6.2", 3006002, "3.6.3", 3006003, null); - break; - case "3.6.3": - updateSchema("3.6.3", 3006003, "3.6.4", 3006004, null); - break; - case "3.6.4": - updateSchema("3.6.4", 3006004, "3.7.0", 3007000, null); - break; - case "3.7.0": - updateSchema("3.7.0", 3007000, "3.7.1", 3007001, connection -> { - try { - connection.createStatement().execute("UPDATE rule_node SET " + - "configuration = CASE " + - " WHEN (configuration::jsonb ->> 'persistAlarmRulesState') = 'false'" + - " THEN (configuration::jsonb || '{\"fetchAlarmRulesStateOnStart\": \"false\"}'::jsonb)::varchar " + - " ELSE configuration " + - "END, " + - "configuration_version = 1 " + - "WHERE type = 'org.thingsboard.rule.engine.profile.TbDeviceProfileNode' " + - "AND configuration_version < 1;"); - } catch (Exception e) { - log.warn("Failed to execute update script for device profile rule nodes due to: ", e); - } - }); - break; - case "3.7.1": - updateSchema("3.7.1", 3007001, "3.7.2", 3007002, null); - break; - default: - throw new RuntimeException("Unable to upgrade SQL database, unsupported fromVersion: " + fromVersion); + "WHERE type = 'org.thingsboard.rule.engine.telemetry.TbMsgAttributesNode' AND configuration_version < 1;", + "CREATE INDEX IF NOT EXISTS idx_notification_recipient_id_unread ON notification(recipient_id) WHERE status <> 'READ';" + ); + } + case "3.6.0" -> updateSchema("3.6.0", 3006000, "3.6.1", 3006001); + case "3.6.1" -> { + updateSchema("3.6.1", 3006001, "3.6.2", 3006002); + + try { + Path saveAttributesNodeUpdateFile = Paths.get(installScripts.getDataDir(), "upgrade", "3.6.1", "save_attributes_node_update.sql"); + loadSql(saveAttributesNodeUpdateFile); + } catch (Exception e) { + log.warn("Failed to execute update script for save attributes rule nodes due to: ", e); + } + execute("CREATE INDEX IF NOT EXISTS idx_asset_profile_id ON asset(tenant_id, asset_profile_id);"); + } + case "3.6.2" -> updateSchema("3.6.2", 3006002, "3.6.3", 3006003); + case "3.6.3" -> updateSchema("3.6.3", 3006003, "3.6.4", 3006004); + case "3.6.4" -> updateSchema("3.6.4", 3006004, "3.7.0", 3007000); + case "3.7.0" -> { + updateSchema("3.7.0", 3007000, "3.7.1", 3007001); + + try { + execute("UPDATE rule_node SET " + + "configuration = CASE " + + " WHEN (configuration::jsonb ->> 'persistAlarmRulesState') = 'false'" + + " THEN (configuration::jsonb || '{\"fetchAlarmRulesStateOnStart\": \"false\"}'::jsonb)::varchar " + + " ELSE configuration " + + "END, " + + "configuration_version = 1 " + + "WHERE type = 'org.thingsboard.rule.engine.profile.TbDeviceProfileNode' " + + "AND configuration_version < 1;", false); + } catch (Exception e) { + log.warn("Failed to execute update script for device profile rule nodes due to: ", e); + } + } + case "3.7.1" -> updateSchema("3.7.1", 3007001, "3.7.2", 3007002); + default -> throw new RuntimeException("Unsupported fromVersion '" + fromVersion + "'"); } } - private void updateSchema(String oldVersionStr, int oldVersion, String newVersionStr, int newVersion, Consumer additionalAction) { - try (Connection conn = DriverManager.getConnection(dbUrl, dbUserName, dbPassword)) { - log.info("Updating schema ..."); - if (isOldSchema(conn, oldVersion)) { - Path schemaUpdateFile = Paths.get(installScripts.getDataDir(), "upgrade", oldVersionStr, SCHEMA_UPDATE_SQL); - loadSql(schemaUpdateFile, conn); - if (additionalAction != null) { - additionalAction.accept(conn); + private void updateSchema(String oldVersionStr, int oldVersion, String newVersionStr, int newVersion) { + try { + transactionTemplate.executeWithoutResult(ts -> { + log.info("Updating schema ..."); + if (isOldSchema(oldVersion)) { + Path schemaUpdateFile = Paths.get(installScripts.getDataDir(), "upgrade", oldVersionStr, SCHEMA_UPDATE_SQL); + loadSql(schemaUpdateFile); + jdbcTemplate.execute("UPDATE tb_schema_settings SET schema_version = " + newVersion); + log.info("Schema updated to version {}", newVersionStr); + } else { + log.info("Skip schema re-update to version {}. Use env flag 'SKIP_SCHEMA_VERSION_CHECK' to force the re-update.", newVersionStr); } - conn.createStatement().execute("UPDATE tb_schema_settings SET schema_version = " + newVersion + ";"); - log.info("Schema updated to version {}", newVersionStr); - } else { - log.info("Skip schema re-update to version {}. Use env flag 'SKIP_SCHEMA_VERSION_CHECK' to force the re-update.", newVersionStr); - } + }); } catch (Exception e) { - log.error("Failed updating schema!!!", e); + throw new RuntimeException("Failed to update schema", e); + } + } + + private void loadSql(Path sqlFile) { + String sql; + try { + sql = Files.readString(sqlFile); + } catch (IOException e) { + throw new UncheckedIOException(e); + } + jdbcTemplate.execute((StatementCallback) stmt -> { + stmt.execute(sql); + printWarnings(stmt.getWarnings()); + return null; + }); + } + + private void execute(@Language("sql") String... statements) { + for (String statement : statements) { + execute(statement, true); } } - private void loadSql(Path sqlFile, Connection conn) throws Exception { - String sql = new String(Files.readAllBytes(sqlFile), Charset.forName("UTF-8")); - Statement st = conn.createStatement(); - st.setQueryTimeout((int) TimeUnit.HOURS.toSeconds(3)); - st.execute(sql);//NOSONAR, ignoring because method used to execute thingsboard database upgrade script - printWarnings(st); - Thread.sleep(5000); + private void execute(@Language("sql") String statement, boolean ignoreErrors) { + try { + jdbcTemplate.execute(statement); + } catch (Exception e) { + if (!ignoreErrors) { + throw e; + } + } } - protected void printWarnings(Statement statement) throws SQLException { - SQLWarning warnings = statement.getWarnings(); + private void printWarnings(SQLWarning warnings) { if (warnings != null) { log.info("{}", warnings.getMessage()); SQLWarning nextWarning = warnings.getNextWarning(); @@ -186,26 +170,18 @@ public class SqlDatabaseUpgradeService implements DatabaseEntitiesUpgradeService } } - protected boolean isOldSchema(Connection conn, long fromVersion) { + private boolean isOldSchema(long fromVersion) { if (DefaultDataUpdateService.getEnv("SKIP_SCHEMA_VERSION_CHECK", false)) { log.info("Skipped DB schema version check due to SKIP_SCHEMA_VERSION_CHECK set to true!"); return true; } + jdbcTemplate.execute("CREATE TABLE IF NOT EXISTS tb_schema_settings (schema_version bigint NOT NULL, CONSTRAINT tb_schema_settings_pkey PRIMARY KEY (schema_version))"); + Long schemaVersion = jdbcTemplate.queryForList("SELECT schema_version FROM tb_schema_settings", Long.class).stream().findFirst().orElse(null); boolean isOldSchema = true; - try { - Statement statement = conn.createStatement(); - statement.execute("CREATE TABLE IF NOT EXISTS tb_schema_settings ( schema_version bigint NOT NULL, CONSTRAINT tb_schema_settings_pkey PRIMARY KEY (schema_version));"); - Thread.sleep(1000); - ResultSet resultSet = statement.executeQuery("SELECT schema_version FROM tb_schema_settings;"); - if (resultSet.next()) { - isOldSchema = resultSet.getLong(1) <= fromVersion; - } else { - resultSet.close(); - statement.execute("INSERT INTO tb_schema_settings (schema_version) VALUES (" + fromVersion + ")"); - } - statement.close(); - } catch (InterruptedException | SQLException e) { - log.info("Failed to check current PostgreSQL schema due to: {}", e.getMessage()); + if (schemaVersion != null) { + isOldSchema = schemaVersion <= fromVersion; + } else { + jdbcTemplate.execute("INSERT INTO tb_schema_settings (schema_version) VALUES (" + fromVersion + ")"); } return isOldSchema; } From 6d367c8ddd8c9ee7444e36e557ec4b9965219bb1 Mon Sep 17 00:00:00 2001 From: ViacheslavKlimov Date: Tue, 17 Sep 2024 14:07:19 +0300 Subject: [PATCH 2/3] Add getSchemaUpdateFile --- .../server/service/install/SqlDatabaseUpgradeService.java | 8 ++++++-- 1 file changed, 6 insertions(+), 2 deletions(-) diff --git a/application/src/main/java/org/thingsboard/server/service/install/SqlDatabaseUpgradeService.java b/application/src/main/java/org/thingsboard/server/service/install/SqlDatabaseUpgradeService.java index 9c611c40f7..a3b7a7b5bc 100644 --- a/application/src/main/java/org/thingsboard/server/service/install/SqlDatabaseUpgradeService.java +++ b/application/src/main/java/org/thingsboard/server/service/install/SqlDatabaseUpgradeService.java @@ -17,6 +17,7 @@ package org.thingsboard.server.service.install; import lombok.extern.slf4j.Slf4j; import org.intellij.lang.annotations.Language; +import org.jetbrains.annotations.NotNull; import org.springframework.context.annotation.Profile; import org.springframework.jdbc.core.JdbcTemplate; import org.springframework.jdbc.core.StatementCallback; @@ -116,8 +117,7 @@ public class SqlDatabaseUpgradeService implements DatabaseEntitiesUpgradeService transactionTemplate.executeWithoutResult(ts -> { log.info("Updating schema ..."); if (isOldSchema(oldVersion)) { - Path schemaUpdateFile = Paths.get(installScripts.getDataDir(), "upgrade", oldVersionStr, SCHEMA_UPDATE_SQL); - loadSql(schemaUpdateFile); + loadSql(getSchemaUpdateFile(oldVersionStr)); jdbcTemplate.execute("UPDATE tb_schema_settings SET schema_version = " + newVersion); log.info("Schema updated to version {}", newVersionStr); } else { @@ -129,6 +129,10 @@ public class SqlDatabaseUpgradeService implements DatabaseEntitiesUpgradeService } } + private Path getSchemaUpdateFile(String version) { + return Paths.get(installScripts.getDataDir(), "upgrade", version, SCHEMA_UPDATE_SQL); + } + private void loadSql(Path sqlFile) { String sql; try { From 1e610d8a5ba07e2fc715f86af698c1071fb38566 Mon Sep 17 00:00:00 2001 From: ViacheslavKlimov Date: Tue, 1 Oct 2024 13:28:07 +0300 Subject: [PATCH 3/3] Fix new upgrade version 3.8.0 --- .../server/service/install/SqlDatabaseUpgradeService.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/application/src/main/java/org/thingsboard/server/service/install/SqlDatabaseUpgradeService.java b/application/src/main/java/org/thingsboard/server/service/install/SqlDatabaseUpgradeService.java index 98c3db4ed2..0bcd9d0014 100644 --- a/application/src/main/java/org/thingsboard/server/service/install/SqlDatabaseUpgradeService.java +++ b/application/src/main/java/org/thingsboard/server/service/install/SqlDatabaseUpgradeService.java @@ -90,7 +90,7 @@ public class SqlDatabaseUpgradeService implements DatabaseEntitiesUpgradeService case "3.6.3" -> updateSchema("3.6.3", 3006003, "3.6.4", 3006004); case "3.6.4" -> updateSchema("3.6.4", 3006004, "3.7.0", 3007000); case "3.7.0" -> { - updateSchema("3.7.0", 3007000, "3.7.1", 3007001); + updateSchema("3.7.0", 3007000, "3.8.0", 3008000); try { execute("UPDATE rule_node SET " +