diff --git a/.claude/docs/synchronizer-model.md b/.claude/docs/synchronizer-model.md index b53bb9d3a3a..2f2c06c06d3 100644 --- a/.claude/docs/synchronizer-model.md +++ b/.claude/docs/synchronizer-model.md @@ -8,7 +8,7 @@ Each artefact type has three collaborating pieces: 2. A **Synchronizer** (`extends BaseSynchronizer` or `MultitenantBaseSynchronizer`) — registered as a Spring bean, scans the repository for a specific file extension/pattern, parses it, upserts the artefact, and reacts to lifecycle phases (`CREATE`, `UPDATE`, `DELETE`, `START`, `STOP`). Ordering across synchronizer types is fixed in `SynchronizersOrder`. 3. An optional **engine/service/endpoint** that consumes the live artefact (e.g. Quartz scheduler for jobs, Flowable for `.bpmn`, Camel for routes, Spring MVC endpoints for `expose` declarations). -Existing synchronizer implementations (grep `extends BaseSynchronizer` / `extends MultitenantBaseSynchronizer`) give a complete inventory of supported artefact types: `Job`, `Bpmn`, `Camel`, `Listener`, `Csvim`, `DataSource`, `Table`, `View`, `Access` (security), `Expose` (URL routing), `ExtensionPoint` / `Extension`, `Markdown`, `Proxy`, `PrintTemplate` (`.print` → CMS seeding), etc. Adding a new artefact type means producing a new synchronizer + entity + (usually) a service, then registering it in the relevant `group-*` aggregator. +Existing synchronizer implementations (grep `extends BaseSynchronizer` / `extends MultitenantBaseSynchronizer`) give a complete inventory of supported artefact types: `Job`, `Bpmn`, `Camel`, `Listener`, `Csvim`, `DataSource`, `Table`, `View`, `Access` (security), `Expose` (URL routing), `ExtensionPoint` / `Extension`, `Markdown`, `Proxy`, `PrintTemplate` (`.print` → CMS seeding), `Migration` (`.migration` → a data migration applied once per tenant schema, recorded in that schema's `DIRIGIBLE_MIGRATIONS`), etc. Adding a new artefact type means producing a new synchronizer + entity + (usually) a service, then registering it in the relevant `group-*` aggregator. JS/TS user code is **not** synchronized — it is loaded on demand by `engine-javascript` (`JavascriptEndpoint` at `/services/js/...`, `/public/js/...`) via `DirigibleJavascriptCodeRunner` backed by Graalium/GraalVM polyglot. The `api-*` Java modules under `components/api/` register the JS-callable APIs (`@dirigible/db`, `@dirigible/http`, etc.) into the GraalJS context — pre-built TS/JS bundles for those APIs live in `components/api/api-modules-javascript/src/main/resources/META-INF/dirigible/modules/`. diff --git a/components/core/core-base/src/main/java/org/eclipse/dirigible/components/base/synchronizer/SynchronizersOrder.java b/components/core/core-base/src/main/java/org/eclipse/dirigible/components/base/synchronizer/SynchronizersOrder.java index df70325c025..c4ae35d5e13 100644 --- a/components/core/core-base/src/main/java/org/eclipse/dirigible/components/base/synchronizer/SynchronizersOrder.java +++ b/components/core/core-base/src/main/java/org/eclipse/dirigible/components/base/synchronizer/SynchronizersOrder.java @@ -66,6 +66,13 @@ public interface SynchronizersOrder { /** The component. */ int COMPONENT = 250; + /** + * The data migration ({@code .migration}). After every synchronizer that evolves structure (schema, + * table, view, entity), so a migration finds the columns it backfills, and before the CSVIM seed, + * so a seed lands on migrated data. + */ + int MIGRATION = 260; + /** The bpmn. */ int BPMN = 300; diff --git a/components/core/core-liquibase/src/main/resources/db/changelog/dirigible-system.json b/components/core/core-liquibase/src/main/resources/db/changelog/dirigible-system.json index 7a7416143dc..92af0b8cbcb 100644 --- a/components/core/core-liquibase/src/main/resources/db/changelog/dirigible-system.json +++ b/components/core/core-liquibase/src/main/resources/db/changelog/dirigible-system.json @@ -11283,6 +11283,218 @@ } ] } + }, + { + "changeSet": { + "id": "create-DIRIGIBLE_MIGRATION_FILES", + "author": "dirigible", + "preConditions": [ + { + "onFail": "MARK_RAN" + }, + { + "not": [ + { + "tableExists": { + "tableName": "DIRIGIBLE_MIGRATION_FILES" + } + } + ] + } + ], + "changes": [ + { + "createTable": { + "columns": [ + { + "column": { + "autoIncrement": true, + "constraints": { + "nullable": false, + "primaryKey": true, + "primaryKeyName": "PK_DIRIGIBLE_MIGRATION_FILES" + }, + "name": "MIGRATION_ID", + "type": "BIGINT" + } + }, + { + "column": { + "name": "CREATED_AT", + "type": "TIMESTAMP" + } + }, + { + "column": { + "name": "CREATED_BY", + "type": "CLOB" + } + }, + { + "column": { + "name": "UPDATED_AT", + "type": "TIMESTAMP" + } + }, + { + "column": { + "name": "UPDATED_BY", + "type": "CLOB" + } + }, + { + "column": { + "name": "ARTEFACT_DEPENDENCIES", + "type": "CLOB" + } + }, + { + "column": { + "name": "ARTEFACT_DESCRIPTION", + "type": "CLOB" + } + }, + { + "column": { + "name": "ARTEFACT_ERROR", + "type": "CLOB" + } + }, + { + "column": { + "constraints": { + "nullable": false + }, + "name": "ARTEFACT_KEY", + "type": "VARCHAR(2000)" + } + }, + { + "column": { + "name": "ARTEFACT_STATUS", + "type": "CLOB" + } + }, + { + "column": { + "constraints": { + "nullable": false + }, + "name": "ARTEFACT_LOCATION", + "type": "CLOB" + } + }, + { + "column": { + "constraints": { + "nullable": false + }, + "name": "ARTEFACT_NAME", + "type": "CLOB" + } + }, + { + "column": { + "name": "ARTEFACT_PHASE", + "type": "CLOB" + } + }, + { + "column": { + "name": "ARTEFACT_RUNNING", + "type": "BOOLEAN" + } + }, + { + "column": { + "constraints": { + "nullable": false + }, + "name": "ARTEFACT_TYPE", + "type": "CLOB" + } + }, + { + "column": { + "constraints": { + "nullable": false + }, + "name": "MIGRATION_PROJECT", + "type": "VARCHAR(255)" + } + }, + { + "column": { + "constraints": { + "nullable": false + }, + "name": "MIGRATION_VERSION", + "type": "VARCHAR(64)" + } + }, + { + "column": { + "constraints": { + "nullable": false + }, + "name": "MIGRATION_SCOPE", + "type": "VARCHAR(16)" + } + }, + { + "column": { + "constraints": { + "nullable": false + }, + "name": "MIGRATION_IDEMPOTENT", + "type": "BOOLEAN" + } + }, + { + "column": { + "constraints": { + "nullable": false + }, + "name": "MIGRATION_CHECKSUM", + "type": "VARCHAR(64)" + } + } + ], + "tableName": "DIRIGIBLE_MIGRATION_FILES" + } + } + ] + } + }, + { + "changeSet": { + "id": "UK_DIRIGIBLE_MIGRATION_FILES_ARTEFACT_KEY", + "author": "dirigible", + "preConditions": [ + { + "onFail": "MARK_RAN" + }, + { + "not": [ + { + "uniqueConstraintExists": { + "tableName": "DIRIGIBLE_MIGRATION_FILES", + "constraintName": "UK_DIRIGIBLE_MIGRATION_FILES_ARTEFACT_KEY" + } + } + ] + } + ], + "changes": [ + { + "addUniqueConstraint": { + "columnNames": "ARTEFACT_KEY", + "constraintName": "UK_DIRIGIBLE_MIGRATION_FILES_ARTEFACT_KEY", + "tableName": "DIRIGIBLE_MIGRATION_FILES" + } + } + ] + } } ] } diff --git a/components/data/data-migrations/pom.xml b/components/data/data-migrations/pom.xml new file mode 100644 index 00000000000..65ddc710b3b --- /dev/null +++ b/components/data/data-migrations/pom.xml @@ -0,0 +1,53 @@ + + + 4.0.0 + + + org.eclipse.dirigible + dirigible-components-parent + 15.0.0-SNAPSHOT + ../../pom.xml + + + Components - Data - Migrations + dirigible-components-data-migrations + jar + + + + org.eclipse.dirigible + dirigible-components-core-base + + + org.eclipse.dirigible + dirigible-components-core-repository + + + + org.eclipse.dirigible + dirigible-components-data-sources + + + + org.springframework + spring-jdbc + + + org.springframework.boot + spring-boot-health + + + + + org.eclipse.dirigible + dirigible-database-sql-h2 + test + + + + + ../../../licensing-header.txt + ../../../ + + + diff --git a/components/data/data-migrations/src/main/java/org/eclipse/dirigible/components/data/migrations/Migration.java b/components/data/data-migrations/src/main/java/org/eclipse/dirigible/components/data/migrations/Migration.java new file mode 100644 index 00000000000..45755a1f9e2 --- /dev/null +++ b/components/data/data-migrations/src/main/java/org/eclipse/dirigible/components/data/migrations/Migration.java @@ -0,0 +1,143 @@ +/* + * Copyright (c) 2010-2026 Eclipse Dirigible contributors + * + * All rights reserved. This program and the accompanying materials are made available under the + * terms of the Eclipse Public License v2.0 which accompanies this distribution, and is available at + * http://www.eclipse.org/legal/epl-v20.html + * + * SPDX-FileCopyrightText: Eclipse Dirigible contributors SPDX-License-Identifier: EPL-2.0 + */ +package org.eclipse.dirigible.components.data.migrations; + +import org.eclipse.dirigible.components.base.artefact.Artefact; + +import com.google.gson.annotations.Expose; + +import jakarta.persistence.Column; +import jakarta.persistence.Entity; +import jakarta.persistence.GeneratedValue; +import jakarta.persistence.GenerationType; +import jakarta.persistence.Id; +import jakarta.persistence.Table; + +/** + * The artefact of one {@code .migration} file. It records what the file declares; whether the + * migration has been applied is a fact of each target database, kept in that database's + * {@code DIRIGIBLE_MIGRATIONS} ledger, never here. + */ +@Entity +@Table(name = "DIRIGIBLE_MIGRATION_FILES") +public class Migration extends Artefact { + + /** The artefact type. */ + public static final String ARTEFACT_TYPE = "migration"; + + @Id + @GeneratedValue(strategy = GenerationType.IDENTITY) + @Column(name = "MIGRATION_ID", nullable = false) + private Long id; + + @Column(name = "MIGRATION_PROJECT", columnDefinition = "VARCHAR", nullable = false, length = 255) + @Expose + private String project; + + @Column(name = "MIGRATION_VERSION", columnDefinition = "VARCHAR", nullable = false, length = 64) + @Expose + private String version; + + @Column(name = "MIGRATION_SCOPE", columnDefinition = "VARCHAR", nullable = false, length = 16) + @Expose + private String scope; + + @Column(name = "MIGRATION_IDEMPOTENT", columnDefinition = "BOOLEAN", nullable = false) + @Expose + private boolean idempotent; + + @Column(name = "MIGRATION_CHECKSUM", columnDefinition = "VARCHAR", nullable = false, length = 64) + @Expose + private String checksum; + + Migration(String location, String name, MigrationScript script) { + super(location, name, ARTEFACT_TYPE, null, null); + this.project = script.project(); + this.version = script.version(); + this.scope = script.scope() + .name(); + this.idempotent = script.idempotent(); + this.checksum = script.checksum(); + } + + /** For JPA. */ + public Migration() { + super(); + } + + /** + * Gets the id. + * + * @return the id + */ + public Long getId() { + return id; + } + + /** + * Sets the id. + * + * @param id the id + */ + public void setId(Long id) { + this.id = id; + } + + /** + * Gets the project the migration belongs to. + * + * @return the project + */ + public String getProject() { + return project; + } + + /** + * Gets the version. + * + * @return the version + */ + public String getVersion() { + return version; + } + + /** + * Gets where the migration applies: {@code EACH} tenant schema or the {@code SYSTEM} database. + * + * @return the scope + */ + public String getScope() { + return scope; + } + + /** + * Whether an edit of the applied migration re-applies it instead of failing it. + * + * @return true when the migration declares itself idempotent + */ + public boolean isIdempotent() { + return idempotent; + } + + /** + * Gets the checksum of the content the file had when it was parsed. + * + * @return the checksum + */ + public String getChecksum() { + return checksum; + } + + @Override + public String toString() { + return "Migration {id=" + id + ", location='" + location + '\'' + ", project='" + project + '\'' + ", version='" + version + '\'' + + ", scope='" + scope + '\'' + ", lifecycle=" + lifecycle + '}'; + } +} diff --git a/components/data/data-migrations/src/main/java/org/eclipse/dirigible/components/data/migrations/MigrationExecutor.java b/components/data/data-migrations/src/main/java/org/eclipse/dirigible/components/data/migrations/MigrationExecutor.java new file mode 100644 index 00000000000..ca5b14f4394 --- /dev/null +++ b/components/data/data-migrations/src/main/java/org/eclipse/dirigible/components/data/migrations/MigrationExecutor.java @@ -0,0 +1,191 @@ +/* + * Copyright (c) 2010-2026 Eclipse Dirigible contributors + * + * All rights reserved. This program and the accompanying materials are made available under the + * terms of the Eclipse Public License v2.0 which accompanies this distribution, and is available at + * http://www.eclipse.org/legal/epl-v20.html + * + * SPDX-FileCopyrightText: Eclipse Dirigible contributors SPDX-License-Identifier: EPL-2.0 + */ +package org.eclipse.dirigible.components.data.migrations; + +import java.nio.charset.StandardCharsets; +import java.sql.Connection; +import java.sql.SQLException; +import java.time.Instant; +import java.util.Comparator; +import java.util.List; +import java.util.Map; +import java.util.concurrent.TimeUnit; + +import javax.sql.DataSource; + +import org.eclipse.dirigible.components.base.artefact.ArtefactLifecycle; +import org.eclipse.dirigible.components.data.migrations.MigrationLedger.Entry; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; +import org.springframework.core.io.ByteArrayResource; +import org.springframework.core.io.support.EncodedResource; +import org.springframework.jdbc.datasource.init.ScriptException; +import org.springframework.jdbc.datasource.init.ScriptUtils; +import org.springframework.stereotype.Component; + +/** + * Applies one migration to one database: the statements and the ledger row that records them in one + * transaction, so a migration is either applied and recorded or neither. What it decides, per + * database: + * + */ +@Component +class MigrationExecutor { + + private static final Logger LOGGER = LoggerFactory.getLogger(MigrationExecutor.class); + + /** What applying a migration to one database came to. */ + enum Status { + /** It ran now. */ + APPLIED, + /** The ledger already records it with the same content. */ + ALREADY_APPLIED, + /** An earlier version of the project has not applied yet, but still may in this pass. */ + WAITING, + /** It did not apply and will not by itself: a failed statement, an edit, a failed predecessor. */ + FAILED + } + + /** + * The outcome for one database. + * + * @param tenant the tenant whose database it is, {@code system} for the system database + * @param status what happened + * @param message why, for anything but a success + * @param cause the exception behind a failure, if one was thrown + */ + record Outcome(String tenant, Status status, String message, Throwable cause) { + + Outcome(String tenant, Status status, String message) { + this(tenant, status, message, null); + } + } + + private final MigrationLedger ledger; + + MigrationExecutor(MigrationLedger ledger) { + this.ledger = ledger; + } + + /** + * Applies a migration to a database unless its ledger says otherwise. Never throws: every failure + * is an outcome, so one tenant's failure cannot keep the next tenant from being migrated. + * + * @param script the migration + * @param location the registry-relative location of its file + * @param tenant the tenant whose database it is, {@code system} for the system database + * @param dataSource the database + * @param predecessors the project's other migrations of the same scope with a lower version + * @return the outcome + */ + Outcome apply(MigrationScript script, String location, String tenant, DataSource dataSource, List predecessors) { + try { + ledger.prepare(dataSource); + try (Connection connection = dataSource.getConnection()) { + Map applied = ledger.findByProject(connection, script.project()); + Entry recorded = applied.get(script.version()); + if (recorded == null) { + Outcome waiting = waitForPredecessors(applied, script, tenant, predecessors); + return waiting != null ? waiting : run(connection, script, location, tenant, false); + } + if (recorded.checksum() + .equals(script.checksum())) { + return new Outcome(tenant, Status.ALREADY_APPLIED, null); + } + if (!script.idempotent()) { + return new Outcome(tenant, Status.FAILED, "Migration [" + location + "] was applied to tenant [" + tenant + "] on [" + + recorded.appliedAt() + "] with checksum [" + recorded.checksum() + "], but the file now has checksum [" + + script.checksum() + + "]. An applied migration is never re-run: restore the file and ship the change as a new version, or declare the migration [idempotent: true] if running it again is safe."); + } + return run(connection, script, location, tenant, true); + } + } catch (SQLException | RuntimeException ex) { + // DEBUG: the synchronizer registers the failure with this cause, which logs it at ERROR the + // first time and quietly when the retry of the FAILED artefact meets it again (#7248). + LOGGER.debug("Failed to apply migration [{}] to tenant [{}]", location, tenant, ex); + return new Outcome(tenant, Status.FAILED, "Migration [" + location + "] failed for tenant [" + tenant + "]: " + describe(ex), + ex); + } + } + + private static Outcome waitForPredecessors(Map applied, MigrationScript script, String tenant, + List predecessors) { + return predecessors.stream() + .filter(predecessor -> !applied.containsKey(predecessor.getVersion())) + .min(Comparator.comparing(Migration::getVersion, MigrationScript::compareVersions)) + .map(pending -> isFailed(pending) + ? new Outcome(tenant, Status.FAILED, + "Migration version [" + script.version() + "] of project [" + script.project() + + "] waits for version [" + pending.getVersion() + "] (" + pending.getLocation() + + "), which failed for tenant [" + tenant + "]") + : new Outcome(tenant, Status.WAITING, + "Migration version [" + script.version() + "] of project [" + script.project() + + "] waits for version [" + pending.getVersion() + "] (" + pending.getLocation() + + ") to apply to tenant [" + tenant + "] first")) + .orElse(null); + } + + private static boolean isFailed(Migration migration) { + return ArtefactLifecycle.FAILED.equals(migration.getLifecycle()) || ArtefactLifecycle.FATAL.equals(migration.getLifecycle()); + } + + private Outcome run(Connection connection, MigrationScript script, String location, String tenant, boolean reapply) + throws SQLException { + boolean autoCommit = connection.getAutoCommit(); + connection.setAutoCommit(false); + try { + long started = System.nanoTime(); + ScriptUtils.executeSqlScript(connection, new EncodedResource(new ByteArrayResource(script.sql() + .getBytes(StandardCharsets.UTF_8), + location), StandardCharsets.UTF_8)); + long durationMillis = TimeUnit.NANOSECONDS.toMillis(System.nanoTime() - started); + Entry entry = new Entry(script.project(), script.version(), location, script.checksum(), tenant, Instant.now(), durationMillis); + if (reapply) { + ledger.update(connection, entry); + } else { + ledger.insert(connection, entry); + } + connection.commit(); + LOGGER.info("{} migration [{}] to tenant [{}] in [{}] ms", reapply ? "Re-applied" : "Applied", location, tenant, + durationMillis); + return new Outcome(tenant, Status.APPLIED, null); + } catch (SQLException | ScriptException ex) { + connection.rollback(); + // Another node may have applied it in the meantime - its ledger row is what made this + // insert fail, and it is the proof that the migration is applied. + Entry recorded = ledger.findByProject(connection, script.project()) + .get(script.version()); + if (recorded != null && recorded.checksum() + .equals(script.checksum())) { + LOGGER.info("Migration [{}] was applied to tenant [{}] by another node", location, tenant, ex); + return new Outcome(tenant, Status.ALREADY_APPLIED, null); + } + throw ex; + } finally { + connection.setAutoCommit(autoCommit); + } + } + + private static String describe(Throwable ex) { + Throwable root = ex; + while (root.getCause() != null && root.getCause() != root) { + root = root.getCause(); + } + return root == ex ? String.valueOf(ex.getMessage()) : ex.getMessage() + " - " + root.getMessage(); + } +} diff --git a/components/data/data-migrations/src/main/java/org/eclipse/dirigible/components/data/migrations/MigrationLedger.java b/components/data/data-migrations/src/main/java/org/eclipse/dirigible/components/data/migrations/MigrationLedger.java new file mode 100644 index 00000000000..c17596d629b --- /dev/null +++ b/components/data/data-migrations/src/main/java/org/eclipse/dirigible/components/data/migrations/MigrationLedger.java @@ -0,0 +1,247 @@ +/* + * Copyright (c) 2010-2026 Eclipse Dirigible contributors + * + * All rights reserved. This program and the accompanying materials are made available under the + * terms of the Eclipse Public License v2.0 which accompanies this distribution, and is available at + * http://www.eclipse.org/legal/epl-v20.html + * + * SPDX-FileCopyrightText: Eclipse Dirigible contributors SPDX-License-Identifier: EPL-2.0 + */ +package org.eclipse.dirigible.components.data.migrations; + +import java.sql.Connection; +import java.sql.PreparedStatement; +import java.sql.ResultSet; +import java.sql.SQLException; +import java.sql.Timestamp; +import java.time.Instant; +import java.util.ArrayList; +import java.util.HashMap; +import java.util.List; +import java.util.Map; + +import javax.sql.DataSource; + +import org.eclipse.dirigible.database.sql.DataType; +import org.eclipse.dirigible.database.sql.SqlFactory; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; +import org.springframework.stereotype.Component; + +/** + * The {@code DIRIGIBLE_MIGRATIONS} ledger: one row per migration applied to the database it lives + * in. It lives in the database the migrations change - every tenant schema has its own, the system + * database has one for the {@code tenant: system} migrations - so a migration and the row recording + * it commit in one transaction, and a restored backup carries the ledger that matches its data. The + * table is created on first use. + * + *

+ * The primary key is {@code /}: two nodes applying the same migration at once + * cannot both record it, and the loser's insert fails its transaction, rolling its run back with + * it. + */ +@Component +class MigrationLedger { + + private static final Logger LOGGER = LoggerFactory.getLogger(MigrationLedger.class); + + /** Unquoted table name - used for metadata existence checks and by the DML builders. */ + static final String TABLE_NAME = "DIRIGIBLE_MIGRATIONS"; + + private static final String COLUMN_KEY = "MIGRATION_KEY"; + private static final String COLUMN_PROJECT = "MIGRATION_PROJECT"; + private static final String COLUMN_VERSION = "MIGRATION_VERSION"; + private static final String COLUMN_LOCATION = "MIGRATION_LOCATION"; + private static final String COLUMN_CHECKSUM = "MIGRATION_CHECKSUM"; + private static final String COLUMN_TENANT = "MIGRATION_TENANT"; + private static final String COLUMN_APPLIED_AT = "MIGRATION_APPLIED_AT"; + private static final String COLUMN_DURATION = "MIGRATION_DURATION_MILLIS"; + + /** + * One applied migration. + * + * @param project the project the migration belongs to + * @param version the migration's version + * @param location the registry-relative location of the file that was applied + * @param checksum the checksum of the content that was applied + * @param tenant the tenant whose database the migration was applied to, {@code system} for the + * system database + * @param appliedAt when the migration was (last) applied + * @param durationMillis how long the run took + */ + record Entry(String project, String version, String location, String checksum, String tenant, Instant appliedAt, long durationMillis) { + + String key() { + return project + "/" + version; + } + } + + /** + * Makes sure the ledger exists in the given database, on a connection of its own so that no DDL + * runs inside a migration's transaction. + * + * @param dataSource the database + * @throws SQLException if the table is missing and cannot be created + */ + void prepare(DataSource dataSource) throws SQLException { + try (Connection connection = dataSource.getConnection()) { + if (SqlFactory.getNative(connection) + .existsTable(connection, TABLE_NAME)) { + return; + } + String sql = SqlFactory.getNative(connection) + .create() + .table(quoted(TABLE_NAME)) + .column(quoted(COLUMN_KEY), DataType.VARCHAR, true, false, false, "(512)") + .column(quoted(COLUMN_PROJECT), DataType.VARCHAR, false, false, false, "(255)") + .column(quoted(COLUMN_VERSION), DataType.VARCHAR, false, false, false, "(64)") + .column(quoted(COLUMN_LOCATION), DataType.VARCHAR, false, false, false, "(1024)") + .column(quoted(COLUMN_CHECKSUM), DataType.VARCHAR, false, false, false, "(64)") + .column(quoted(COLUMN_TENANT), DataType.VARCHAR, false, false, false, "(255)") + .column(quoted(COLUMN_APPLIED_AT), DataType.TIMESTAMP, false, false, false) + .column(quoted(COLUMN_DURATION), DataType.BIGINT, false, false, false) + .build(); + try (PreparedStatement statement = connection.prepareStatement(sql)) { + statement.executeUpdate(); + LOGGER.info("Created the migrations ledger [{}]", TABLE_NAME); + } catch (SQLException ex) { + // Another node or tenant pass may have created it in the meantime; tolerate it. + if (SqlFactory.getNative(connection) + .existsTable(connection, TABLE_NAME)) { + LOGGER.debug("The migrations ledger already exists after a concurrent creation", ex); + return; + } + throw ex; + } + } + } + + /** + * Reads the applied migrations of a project, keyed by version. + * + * @param connection a connection to the database the ledger lives in + * @param project the project + * @return the project's applied migrations + * @throws SQLException if the read fails + */ + Map findByProject(Connection connection, String project) throws SQLException { + String sql = SqlFactory.getNative(connection) + .select() + .column("*") + .from(TABLE_NAME) + .where(COLUMN_PROJECT + " = ?") + .build(); + Map applied = new HashMap<>(); + try (PreparedStatement statement = connection.prepareStatement(sql)) { + statement.setString(1, project); + try (ResultSet resultSet = statement.executeQuery()) { + while (resultSet.next()) { + Entry entry = read(resultSet); + applied.put(entry.version(), entry); + } + } + } + return applied; + } + + /** + * Records a migration applied for the first time. + * + * @param connection the connection of the transaction that ran the migration + * @param entry the applied migration + * @throws SQLException if the insert fails - also when another node recorded it first + */ + void insert(Connection connection, Entry entry) throws SQLException { + String sql = SqlFactory.getNative(connection) + .insert() + .into(TABLE_NAME) + .column(COLUMN_KEY) + .column(COLUMN_PROJECT) + .column(COLUMN_VERSION) + .column(COLUMN_LOCATION) + .column(COLUMN_CHECKSUM) + .column(COLUMN_TENANT) + .column(COLUMN_APPLIED_AT) + .column(COLUMN_DURATION) + .build(); + try (PreparedStatement statement = connection.prepareStatement(sql)) { + statement.setString(1, entry.key()); + statement.setString(2, entry.project()); + statement.setString(3, entry.version()); + statement.setString(4, entry.location()); + statement.setString(5, entry.checksum()); + statement.setString(6, entry.tenant()); + statement.setTimestamp(7, Timestamp.from(entry.appliedAt())); + statement.setLong(8, entry.durationMillis()); + statement.executeUpdate(); + } + } + + /** + * Records the re-application of an idempotent migration whose content changed. + * + * @param connection the connection of the transaction that re-ran the migration + * @param entry the re-applied migration + * @throws SQLException if the update fails + */ + void update(Connection connection, Entry entry) throws SQLException { + String sql = SqlFactory.getNative(connection) + .update() + .table(TABLE_NAME) + .set(COLUMN_LOCATION, "?") + .set(COLUMN_CHECKSUM, "?") + .set(COLUMN_APPLIED_AT, "?") + .set(COLUMN_DURATION, "?") + .where(COLUMN_KEY + " = ?") + .build(); + try (PreparedStatement statement = connection.prepareStatement(sql)) { + statement.setString(1, entry.location()); + statement.setString(2, entry.checksum()); + statement.setTimestamp(3, Timestamp.from(entry.appliedAt())); + statement.setLong(4, entry.durationMillis()); + statement.setString(5, entry.key()); + statement.executeUpdate(); + } + } + + /** + * Reads every applied migration of a database, by project and version. + * + * @param dataSource the database + * @return the applied migrations, empty when the database has no ledger yet + * @throws SQLException if the read fails + */ + List findAll(DataSource dataSource) throws SQLException { + try (Connection connection = dataSource.getConnection()) { + if (!SqlFactory.getNative(connection) + .existsTable(connection, TABLE_NAME)) { + return List.of(); + } + String sql = SqlFactory.getNative(connection) + .select() + .column("*") + .from(TABLE_NAME) + .order(COLUMN_PROJECT) + .order(COLUMN_APPLIED_AT) + .build(); + List entries = new ArrayList<>(); + try (PreparedStatement statement = connection.prepareStatement(sql); ResultSet resultSet = statement.executeQuery()) { + while (resultSet.next()) { + entries.add(read(resultSet)); + } + } + return entries; + } + } + + private static Entry read(ResultSet resultSet) throws SQLException { + return new Entry(resultSet.getString(COLUMN_PROJECT), resultSet.getString(COLUMN_VERSION), resultSet.getString(COLUMN_LOCATION), + resultSet.getString(COLUMN_CHECKSUM), resultSet.getString(COLUMN_TENANT), resultSet.getTimestamp(COLUMN_APPLIED_AT) + .toInstant(), + resultSet.getLong(COLUMN_DURATION)); + } + + private static String quoted(String identifier) { + return "\"" + identifier + "\""; + } +} diff --git a/components/data/data-migrations/src/main/java/org/eclipse/dirigible/components/data/migrations/MigrationRepository.java b/components/data/data-migrations/src/main/java/org/eclipse/dirigible/components/data/migrations/MigrationRepository.java new file mode 100644 index 00000000000..61a0b177e0b --- /dev/null +++ b/components/data/data-migrations/src/main/java/org/eclipse/dirigible/components/data/migrations/MigrationRepository.java @@ -0,0 +1,38 @@ +/* + * Copyright (c) 2010-2026 Eclipse Dirigible contributors + * + * All rights reserved. This program and the accompanying materials are made available under the + * terms of the Eclipse Public License v2.0 which accompanies this distribution, and is available at + * http://www.eclipse.org/legal/epl-v20.html + * + * SPDX-FileCopyrightText: Eclipse Dirigible contributors SPDX-License-Identifier: EPL-2.0 + */ +package org.eclipse.dirigible.components.data.migrations; + +import java.util.List; + +import org.eclipse.dirigible.components.base.artefact.ArtefactRepository; +import org.springframework.data.jpa.repository.Modifying; +import org.springframework.data.jpa.repository.Query; +import org.springframework.data.repository.query.Param; +import org.springframework.stereotype.Repository; +import org.springframework.transaction.annotation.Transactional; + +/** The repository of the migration artefacts. */ +@Repository("migrationRepository") +public interface MigrationRepository extends ArtefactRepository { + + @Override + @Modifying + @Transactional + @Query(value = "UPDATE Migration SET running = :running") + void setRunningToAll(@Param("running") boolean running); + + /** + * Finds the migrations of a project. + * + * @param project the project + * @return the project's migrations + */ + List findAllByProject(String project); +} diff --git a/components/data/data-migrations/src/main/java/org/eclipse/dirigible/components/data/migrations/MigrationScript.java b/components/data/data-migrations/src/main/java/org/eclipse/dirigible/components/data/migrations/MigrationScript.java new file mode 100644 index 00000000000..9141d22f882 --- /dev/null +++ b/components/data/data-migrations/src/main/java/org/eclipse/dirigible/components/data/migrations/MigrationScript.java @@ -0,0 +1,179 @@ +/* + * Copyright (c) 2010-2026 Eclipse Dirigible contributors + * + * All rights reserved. This program and the accompanying materials are made available under the + * terms of the Eclipse Public License v2.0 which accompanies this distribution, and is available at + * http://www.eclipse.org/legal/epl-v20.html + * + * SPDX-FileCopyrightText: Eclipse Dirigible contributors SPDX-License-Identifier: EPL-2.0 + */ +package org.eclipse.dirigible.components.data.migrations; + +import java.nio.charset.StandardCharsets; +import java.security.MessageDigest; +import java.security.NoSuchAlgorithmException; +import java.text.ParseException; +import java.util.HexFormat; +import java.util.regex.Matcher; +import java.util.regex.Pattern; + +/** + * One parsed {@code .migration} file: {@code /.../__.migration}, + * SQL, optionally headed by comment lines declaring how it applies: + * + *

+ * -- tenant: each        (default; once per tenant schema) | system (once, on the system database)
+ * -- idempotent: false   (default; an edit after it applied FAILS it) | true (an edit re-applies it)
+ * UPDATE ORDERS SET STATUS = 'OPEN' WHERE STATUS IS NULL;
+ * 
+ * + * The version is one or more dot-separated numbers, optionally prefixed with {@code V}; versions + * compare numerically segment by segment, so {@code 2} precedes {@code 10}. The checksum is taken + * over the content with line endings normalized, so a checkout that converts them does not read as + * an edit. + * + * @param project the project the file belongs to - the first segment of its location + * @param version the version, without the optional {@code V} prefix + * @param description the part of the file name after {@code __} + * @param scope where the migration applies + * @param idempotent whether the migration may run again when its content changes + * @param checksum the SHA-256 of the content, hex-encoded + * @param sql the content, as it is executed + */ +record MigrationScript(String project, String version, String description, Scope scope, boolean idempotent, String checksum, String sql) { + + /** Where a migration applies. */ + enum Scope { + /** Once in every tenant's schema, through the tenant-routed default datasource. */ + EACH, + /** Once, on the system database. */ + SYSTEM + } + + static final String FILE_EXTENSION = ".migration"; + + private static final Pattern FILE_NAME = Pattern.compile("[Vv]?(\\d+(?:\\.\\d+)*)__([A-Za-z0-9][A-Za-z0-9_.\\-]*)\\.migration"); + + private static final Pattern HEADER = Pattern.compile("--\\s*(tenant|idempotent)\\s*:\\s*(\\S*)\\s*", Pattern.CASE_INSENSITIVE); + + /** + * Parses a migration file. + * + * @param location the registry-relative location, {@code //.../.migration} + * @param content the file content + * @return the parsed migration + * @throws ParseException when the file name, a header or the content is not a valid migration + */ + static MigrationScript parse(String location, byte[] content) throws ParseException { + String[] segments = location.split("/"); + if (segments.length < 3 || segments[1].isEmpty()) { + throw new ParseException("The migration [" + location + "] must live inside a project", 0); + } + String fileName = segments[segments.length - 1]; + Matcher name = FILE_NAME.matcher(fileName); + if (!name.matches()) { + throw new ParseException("The migration file name [" + fileName + "] in [" + location + + "] must be __.migration, the version being dot-separated numbers (e.g. 001__backfill_status.migration or V1.2__split_name.migration)", + 0); + } + String sql = new String(content, StandardCharsets.UTF_8).replace("\r\n", "\n"); + + Scope scope = Scope.EACH; + boolean idempotent = false; + boolean hasStatement = false; + for (String line : sql.split("\n")) { + String trimmed = line.trim(); + if (trimmed.isEmpty()) { + continue; + } + if (!trimmed.startsWith("--")) { + hasStatement = true; + break; + } + Matcher header = HEADER.matcher(trimmed); + if (!header.matches()) { + continue; + } + String value = header.group(2) + .toLowerCase(); + if ("tenant".equalsIgnoreCase(header.group(1))) { + scope = parseScope(location, value); + } else { + idempotent = parseIdempotent(location, value); + } + } + if (!hasStatement) { + throw new ParseException("The migration [" + location + "] contains no SQL statement", 0); + } + return new MigrationScript(segments[1], name.group(1), name.group(2), scope, idempotent, checksum(sql), sql); + } + + /** + * Compares two versions numerically, segment by segment; a version that is a prefix of the other + * precedes it ({@code 1} before {@code 1.0}). + * + * @param left a version + * @param right a version + * @return negative, zero or positive as {@code left} precedes, equals or follows {@code right} + */ + static int compareVersions(String left, String right) { + String[] leftSegments = left.split("\\."); + String[] rightSegments = right.split("\\."); + for (int i = 0; i < Math.min(leftSegments.length, rightSegments.length); i++) { + int compared = compareNumbers(leftSegments[i], rightSegments[i]); + if (compared != 0) { + return compared; + } + } + return Integer.compare(leftSegments.length, rightSegments.length); + } + + /** + * Compares two unbounded decimal numbers without parsing them, so a timestamp version cannot + * overflow. + */ + private static int compareNumbers(String left, String right) { + String leftDigits = stripLeadingZeros(left); + String rightDigits = stripLeadingZeros(right); + if (leftDigits.length() != rightDigits.length()) { + return Integer.compare(leftDigits.length(), rightDigits.length()); + } + return leftDigits.compareTo(rightDigits); + } + + private static String stripLeadingZeros(String digits) { + int start = 0; + while (start < digits.length() - 1 && digits.charAt(start) == '0') { + start++; + } + return digits.substring(start); + } + + private static Scope parseScope(String location, String value) throws ParseException { + return switch (value) { + case "each" -> Scope.EACH; + case "system" -> Scope.SYSTEM; + default -> throw new ParseException("The migration [" + location + "] declares [tenant: " + value + + "] - expected [each] (once per tenant schema) or [system] (once, on the system database)", 0); + }; + } + + private static boolean parseIdempotent(String location, String value) throws ParseException { + return switch (value) { + case "true" -> true; + case "false" -> false; + default -> throw new ParseException( + "The migration [" + location + "] declares [idempotent: " + value + "] - expected [true] or [false]", 0); + }; + } + + private static String checksum(String normalizedContent) { + try { + MessageDigest digest = MessageDigest.getInstance("SHA-256"); + return HexFormat.of() + .formatHex(digest.digest(normalizedContent.getBytes(StandardCharsets.UTF_8))); + } catch (NoSuchAlgorithmException ex) { + throw new IllegalStateException("SHA-256 is a mandatory JDK algorithm", ex); + } + } +} diff --git a/components/data/data-migrations/src/main/java/org/eclipse/dirigible/components/data/migrations/MigrationService.java b/components/data/data-migrations/src/main/java/org/eclipse/dirigible/components/data/migrations/MigrationService.java new file mode 100644 index 00000000000..6ce8241ad6a --- /dev/null +++ b/components/data/data-migrations/src/main/java/org/eclipse/dirigible/components/data/migrations/MigrationService.java @@ -0,0 +1,40 @@ +/* + * Copyright (c) 2010-2026 Eclipse Dirigible contributors + * + * All rights reserved. This program and the accompanying materials are made available under the + * terms of the Eclipse Public License v2.0 which accompanies this distribution, and is available at + * http://www.eclipse.org/legal/epl-v20.html + * + * SPDX-FileCopyrightText: Eclipse Dirigible contributors SPDX-License-Identifier: EPL-2.0 + */ +package org.eclipse.dirigible.components.data.migrations; + +import java.util.List; + +import org.eclipse.dirigible.components.base.artefact.BaseArtefactService; +import org.springframework.stereotype.Service; +import org.springframework.transaction.annotation.Transactional; + +/** The service of the migration artefacts. */ +@Service +@Transactional +public class MigrationService extends BaseArtefactService { + + private final MigrationRepository repository; + + MigrationService(MigrationRepository repository) { + super(repository); + this.repository = repository; + } + + /** + * Finds the migrations of a project. + * + * @param project the project + * @return the project's migrations + */ + @Transactional(readOnly = true) + public List findAllByProject(String project) { + return repository.findAllByProject(project); + } +} diff --git a/components/data/data-migrations/src/main/java/org/eclipse/dirigible/components/data/migrations/MigrationsEndpoint.java b/components/data/data-migrations/src/main/java/org/eclipse/dirigible/components/data/migrations/MigrationsEndpoint.java new file mode 100644 index 00000000000..a9d51a7a226 --- /dev/null +++ b/components/data/data-migrations/src/main/java/org/eclipse/dirigible/components/data/migrations/MigrationsEndpoint.java @@ -0,0 +1,69 @@ +/* + * Copyright (c) 2010-2026 Eclipse Dirigible contributors + * + * All rights reserved. This program and the accompanying materials are made available under the + * terms of the Eclipse Public License v2.0 which accompanies this distribution, and is available at + * http://www.eclipse.org/legal/epl-v20.html + * + * SPDX-FileCopyrightText: Eclipse Dirigible contributors SPDX-License-Identifier: EPL-2.0 + */ +package org.eclipse.dirigible.components.data.migrations; + +import java.sql.SQLException; +import java.util.ArrayList; +import java.util.List; + +import org.eclipse.dirigible.components.base.endpoint.BaseEndpoint; +import org.eclipse.dirigible.components.base.tenant.TenantContext; +import org.eclipse.dirigible.components.data.migrations.MigrationLedger.Entry; +import org.eclipse.dirigible.components.data.sources.manager.DataSourcesManager; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; +import org.springframework.http.HttpStatus; +import org.springframework.web.bind.annotation.GetMapping; +import org.springframework.web.bind.annotation.RequestMapping; +import org.springframework.web.bind.annotation.RestController; +import org.springframework.web.server.ResponseStatusException; + +import jakarta.annotation.security.RolesAllowed; + +/** + * The migrations ledger, read-only: what a deployment's smoke test asserts after a release + * ("version N shipped three migrations and each is applied once"). It lists the ledger of the + * caller's tenant; a caller in the default tenant - the operator's - also sees the + * {@code tenant: system} migrations, which no other tenant's administrator has any business + * reading. + */ +@RestController +@RequestMapping(BaseEndpoint.PREFIX_ENDPOINT_CORE + "migrations") +@RolesAllowed({"ADMINISTRATOR", "OPERATOR"}) +class MigrationsEndpoint extends BaseEndpoint { + + private static final Logger LOGGER = LoggerFactory.getLogger(MigrationsEndpoint.class); + + private final MigrationLedger ledger; + private final DataSourcesManager dataSourcesManager; + private final TenantContext tenantContext; + + MigrationsEndpoint(MigrationLedger ledger, DataSourcesManager dataSourcesManager, TenantContext tenantContext) { + this.ledger = ledger; + this.dataSourcesManager = dataSourcesManager; + this.tenantContext = tenantContext; + } + + @GetMapping + List list() { + try { + List entries = new ArrayList<>(ledger.findAll(dataSourcesManager.getDefaultDataSource())); + if (tenantContext.getCurrentTenant() + .isDefault()) { + entries.addAll(ledger.findAll(dataSourcesManager.getSystemDataSource())); + } + return entries; + } catch (SQLException ex) { + LOGGER.error("Failed to read the migrations ledger", ex); + throw new ResponseStatusException(HttpStatus.INTERNAL_SERVER_ERROR, "Failed to read the migrations ledger: " + ex.getMessage(), + ex); + } + } +} diff --git a/components/data/data-migrations/src/main/java/org/eclipse/dirigible/components/data/migrations/MigrationsHealthIndicator.java b/components/data/data-migrations/src/main/java/org/eclipse/dirigible/components/data/migrations/MigrationsHealthIndicator.java new file mode 100644 index 00000000000..53b6ca23551 --- /dev/null +++ b/components/data/data-migrations/src/main/java/org/eclipse/dirigible/components/data/migrations/MigrationsHealthIndicator.java @@ -0,0 +1,55 @@ +/* + * Copyright (c) 2010-2026 Eclipse Dirigible contributors + * + * All rights reserved. This program and the accompanying materials are made available under the + * terms of the Eclipse Public License v2.0 which accompanies this distribution, and is available at + * http://www.eclipse.org/legal/epl-v20.html + * + * SPDX-FileCopyrightText: Eclipse Dirigible contributors SPDX-License-Identifier: EPL-2.0 + */ +package org.eclipse.dirigible.components.data.migrations; + +import java.util.List; + +import org.eclipse.dirigible.components.base.artefact.ArtefactLifecycle; +import org.eclipse.dirigible.components.base.readiness.PlatformReadiness; +import org.springframework.boot.health.contributor.Health; +import org.springframework.boot.health.contributor.HealthIndicator; +import org.springframework.stereotype.Component; + +/** + * The {@code migrations} health component, next to the {@code artefacts} census (#7533): how many + * migrations are registered, how many are still pending (parsed, not yet applied everywhere) and + * how many failed in some database. Like the census it is a quality signal and never turns DOWN; a + * FAILED migration is a failed artefact, so {@code DIRIGIBLE_READINESS_REQUIRE_CLEAN_BOOT} already + * withholds readiness for it, and the boot latch withholds it while the first pass is still + * applying them. UNKNOWN until the first pass has depleted. + */ +@Component +class MigrationsHealthIndicator implements HealthIndicator { + + private final MigrationService migrationService; + + MigrationsHealthIndicator(MigrationService migrationService) { + this.migrationService = migrationService; + } + + @Override + public Health health() { + List migrations = migrationService.getAll(); + long pending = migrations.stream() + .filter(migration -> ArtefactLifecycle.NEW.equals(migration.getLifecycle()) + || ArtefactLifecycle.MODIFIED.equals(migration.getLifecycle())) + .count(); + long failed = migrations.stream() + .filter(migration -> ArtefactLifecycle.FAILED.equals(migration.getLifecycle()) + || ArtefactLifecycle.FATAL.equals(migration.getLifecycle())) + .count(); + Health.Builder builder = PlatformReadiness.getInstance() + .isBootCompleted() ? Health.up() : Health.unknown(); + return builder.withDetail("total", migrations.size()) + .withDetail("pending", pending) + .withDetail("failed", failed) + .build(); + } +} diff --git a/components/data/data-migrations/src/main/java/org/eclipse/dirigible/components/data/migrations/MigrationsSynchronizer.java b/components/data/data-migrations/src/main/java/org/eclipse/dirigible/components/data/migrations/MigrationsSynchronizer.java new file mode 100644 index 00000000000..a6e5e33d49e --- /dev/null +++ b/components/data/data-migrations/src/main/java/org/eclipse/dirigible/components/data/migrations/MigrationsSynchronizer.java @@ -0,0 +1,279 @@ +/* + * Copyright (c) 2010-2026 Eclipse Dirigible contributors + * + * All rights reserved. This program and the accompanying materials are made available under the + * terms of the Eclipse Public License v2.0 which accompanies this distribution, and is available at + * http://www.eclipse.org/legal/epl-v20.html + * + * SPDX-FileCopyrightText: Eclipse Dirigible contributors SPDX-License-Identifier: EPL-2.0 + */ +package org.eclipse.dirigible.components.data.migrations; + +import java.text.ParseException; +import java.util.List; +import java.util.Objects; +import java.util.stream.Collectors; + +import org.eclipse.dirigible.components.base.artefact.ArtefactLifecycle; +import org.eclipse.dirigible.components.base.artefact.ArtefactPhase; +import org.eclipse.dirigible.components.base.artefact.ArtefactService; +import org.eclipse.dirigible.components.base.artefact.topology.TopologyWrapper; +import org.eclipse.dirigible.components.base.synchronizer.BaseSynchronizer; +import org.eclipse.dirigible.components.base.synchronizer.SynchronizerCallback; +import org.eclipse.dirigible.components.base.synchronizer.SynchronizersOrder; +import org.eclipse.dirigible.components.base.tenant.TenantContext; +import org.eclipse.dirigible.components.base.tenant.TenantResult; +import org.eclipse.dirigible.components.data.migrations.MigrationExecutor.Outcome; +import org.eclipse.dirigible.components.data.migrations.MigrationExecutor.Status; +import org.eclipse.dirigible.components.data.migrations.MigrationScript.Scope; +import org.eclipse.dirigible.components.data.sources.manager.DataSourcesManager; +import org.eclipse.dirigible.repository.api.IRepository; +import org.eclipse.dirigible.repository.api.IRepositoryStructure; +import org.eclipse.dirigible.repository.api.IResource; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; +import org.springframework.core.annotation.Order; +import org.springframework.stereotype.Component; + +/** + * Synchronizes {@code .migration} artefacts (see {@link MigrationScript} for the file shape): each + * one is applied exactly once to every tenant's schema, or once to the system database, and + * recorded in that database's {@code DIRIGIBLE_MIGRATIONS} ledger in the same transaction. + * + *

+ * Applying is idempotent by construction - the ledger, not the artefact's lifecycle, says whether a + * database has a migration - so every path that runs it again is safe: the re-parse that follows a + * new tenant's provisioning (it migrates that tenant and finds the others done), a boot against a + * fresh system database, the retry of a FAILED artefact. The artefact is FAILED when any database + * did not take the migration: a failed statement (rolled back with its ledger row), an applied file + * edited afterwards (an applied migration is not re-run unless it declares itself + * {@code idempotent: true}), or an earlier version of its project that failed. + * + *

+ * The tenants are iterated here rather than by {@code BaseSynchronizer}, which completes a + * multitenant artefact once per tenant and keeps the state the LAST tenant registered - one + * tenant's success would hide another's failure. This synchronizer still reports + * {@link #multitenantExecution()} so the post-provisioning re-trigger includes it. + * + *

+ * DELETE and cleanup remove the artefact only. A migration that ran is history: its data change and + * its ledger row stay. + */ +@Component +@Order(SynchronizersOrder.MIGRATION) +class MigrationsSynchronizer extends BaseSynchronizer { + + /** The tenant column of a ledger row written by a {@code tenant: system} migration. */ + static final String SYSTEM_TENANT = "system"; + + private static final Logger LOGGER = LoggerFactory.getLogger(MigrationsSynchronizer.class); + + private final MigrationService migrationService; + private final MigrationExecutor executor; + private final DataSourcesManager dataSourcesManager; + private final TenantContext tenantContext; + private final IRepository repository; + + private SynchronizerCallback callback; + + MigrationsSynchronizer(MigrationService migrationService, MigrationExecutor executor, DataSourcesManager dataSourcesManager, + TenantContext tenantContext, IRepository repository) { + this.migrationService = migrationService; + this.executor = executor; + this.dataSourcesManager = dataSourcesManager; + this.tenantContext = tenantContext; + this.repository = repository; + } + + @Override + public boolean isAccepted(String type) { + return Migration.ARTEFACT_TYPE.equals(type); + } + + @Override + protected List parseImpl(String location, byte[] content) throws ParseException { + MigrationScript script = MigrationScript.parse(location, content); + rejectDuplicateVersion(location, script); + Migration migration = new Migration(location, location.substring(location.lastIndexOf('/') + 1), script); + migration.updateKey(); + try { + Migration existing = migrationService.findByKey(migration.getKey()); + if (existing != null) { + migration.setId(existing.getId()); + } + return List.of(migrationService.save(migration)); + } catch (RuntimeException ex) { + LOGGER.error("Failed to save migration [{}]", location, ex); + throw new ParseException(ex.getMessage(), 0); + } + } + + /** + * Two files may not claim the same version of a project: the ledger keys a migration by project and + * version. A claim by a file that is no longer in the registry does not count - that is a rename, + * whose old artefact the cleanup of this very pass removes. + */ + private void rejectDuplicateVersion(String location, MigrationScript script) throws ParseException { + for (Migration other : migrationService.findAllByProject(script.project())) { + if (!other.getLocation() + .equals(location) + && other.getVersion() + .equals(script.version()) + && repository.getResource(IRepositoryStructure.PATH_REGISTRY_PUBLIC + other.getLocation()) + .exists()) { + throw new ParseException("Version [" + script.version() + "] of project [" + script.project() + "] is claimed by both [" + + location + "] and [" + other.getLocation() + "] - every migration of a project needs its own version", 0); + } + } + } + + @Override + public ArtefactService getService() { + return migrationService; + } + + @Override + public List retrieve(String location) { + return migrationService.findByLocation(location); + } + + @Override + public void setStatus(Migration artefact, ArtefactLifecycle lifecycle, String error) { + artefact.setLifecycle(lifecycle); + artefact.setError(error); + migrationService.save(artefact); + } + + @Override + protected boolean completeImpl(TopologyWrapper wrapper, ArtefactPhase flow) { + Migration migration = wrapper.getArtefact(); + ArtefactLifecycle lifecycle = migration.getLifecycle(); + return switch (flow) { + case CREATE -> !ArtefactLifecycle.NEW.equals(lifecycle) || apply(wrapper, ArtefactLifecycle.CREATED); + case UPDATE -> !ArtefactLifecycle.MODIFIED.equals(lifecycle) || apply(wrapper, ArtefactLifecycle.UPDATED); + // The retry of a FAILED migration - in this pass, and on an idle instance (#7248). + case START -> !ArtefactLifecycle.FAILED.equals(lifecycle) || apply(wrapper, ArtefactLifecycle.CREATED); + case DELETE -> { + migrationService.delete(migration); + callback.registerState(this, wrapper, ArtefactLifecycle.DELETED); + yield true; + } + case PREPARE, STOP -> true; + }; + } + + /** + * Applies the migration to every database its scope names and registers the one state that sums + * them up. + * + * @return false only while the migration waits for an earlier version still being applied in this + * pass - the depleter then offers it again once that one has completed + */ + private boolean apply(TopologyWrapper wrapper, ArtefactLifecycle success) { + Migration migration = wrapper.getArtefact(); + String location = migration.getLocation(); + MigrationScript script; + try { + script = read(location); + } catch (ParseException ex) { + fail(wrapper, ex.getMessage(), ex); + return true; + } + + List predecessors = predecessors(script, location); + List outcomes = script.scope() == Scope.SYSTEM + ? List.of(executor.apply(script, location, SYSTEM_TENANT, dataSourcesManager.getSystemDataSource(), predecessors)) + : tenantContext.executeForEachTenant(() -> executor.apply(script, location, tenantContext.getCurrentTenant() + .getId(), + dataSourcesManager.getDefaultDataSource(), predecessors)) + .stream() + .map(TenantResult::getResult) + .toList(); + + String failures = messages(outcomes, Status.FAILED); + if (!failures.isEmpty()) { + Throwable cause = outcomes.stream() + .map(Outcome::cause) + .filter(Objects::nonNull) + .findFirst() + .orElse(null); + fail(wrapper, failures, cause); + return true; + } + String waiting = messages(outcomes, Status.WAITING); + if (!waiting.isEmpty()) { + // No state is registered: the lifecycle the phase matched on must survive for the next + // round. The message reaches the "undepleted" report should the wait never end. + migration.setError(waiting); + return false; + } + callback.registerState(this, wrapper, success); + return true; + } + + /** The file as it is in the registry now - the artefact keeps what it declares, not the SQL. */ + private MigrationScript read(String location) throws ParseException { + IResource resource = repository.getResource(IRepositoryStructure.PATH_REGISTRY_PUBLIC + location); + if (!resource.exists()) { + throw new ParseException("The migration [" + location + "] is no longer in the registry", 0); + } + return MigrationScript.parse(location, resource.getContent()); + } + + private List predecessors(MigrationScript script, String location) { + return migrationService.findAllByProject(script.project()) + .stream() + .filter(other -> !other.getLocation() + .equals(location)) + .filter(other -> Objects.equals(other.getScope(), script.scope() + .name())) + .filter(other -> MigrationScript.compareVersions(other.getVersion(), script.version()) < 0) + .toList(); + } + + private static String messages(List outcomes, Status status) { + return outcomes.stream() + .filter(outcome -> outcome.status() == status) + .map(Outcome::message) + .collect(Collectors.joining("; ")); + } + + private void fail(TopologyWrapper wrapper, String message, Throwable cause) { + callback.addError(message); + callback.registerState(this, wrapper, ArtefactLifecycle.FAILED, message, cause); + } + + @Override + public void cleanupImpl(Migration migration) { + try { + migrationService.delete(migration); + } catch (RuntimeException ex) { + callback.addError(ex.getMessage()); + callback.registerState(this, migration, ArtefactLifecycle.DELETED, ex); + } + } + + /** + * Included in the post-provisioning re-trigger, so a new tenant is migrated. The tenants are then + * iterated in {@link #apply}, not per artefact by the base class - see the class comment. + */ + @Override + public boolean multitenantExecution() { + return true; + } + + @Override + public void setCallback(SynchronizerCallback callback) { + this.callback = callback; + } + + @Override + public String getFileExtension() { + return MigrationScript.FILE_EXTENSION; + } + + @Override + public String getArtefactType() { + return Migration.ARTEFACT_TYPE; + } +} diff --git a/components/data/data-migrations/src/test/java/org/eclipse/dirigible/components/data/migrations/MigrationExecutorTest.java b/components/data/data-migrations/src/test/java/org/eclipse/dirigible/components/data/migrations/MigrationExecutorTest.java new file mode 100644 index 00000000000..97b45e9e004 --- /dev/null +++ b/components/data/data-migrations/src/test/java/org/eclipse/dirigible/components/data/migrations/MigrationExecutorTest.java @@ -0,0 +1,185 @@ +/* + * Copyright (c) 2010-2026 Eclipse Dirigible contributors + * + * All rights reserved. This program and the accompanying materials are made available under the + * terms of the Eclipse Public License v2.0 which accompanies this distribution, and is available at + * http://www.eclipse.org/legal/epl-v20.html + * + * SPDX-FileCopyrightText: Eclipse Dirigible contributors SPDX-License-Identifier: EPL-2.0 + */ +package org.eclipse.dirigible.components.data.migrations; + +import static org.assertj.core.api.Assertions.assertThat; + +import java.nio.charset.StandardCharsets; +import java.sql.Connection; +import java.sql.ResultSet; +import java.sql.SQLException; +import java.sql.Statement; +import java.text.ParseException; +import java.util.List; + +import org.eclipse.dirigible.components.base.artefact.ArtefactLifecycle; +import org.eclipse.dirigible.components.data.migrations.MigrationExecutor.Outcome; +import org.eclipse.dirigible.components.data.migrations.MigrationExecutor.Status; +import org.eclipse.dirigible.components.data.migrations.MigrationLedger.Entry; +import org.h2.jdbcx.JdbcDataSource; +import org.junit.jupiter.api.AfterEach; +import org.junit.jupiter.api.BeforeEach; +import org.junit.jupiter.api.Test; + +/** + * The executor on a real in-memory database: the transaction and the ledger are what it is about. + */ +class MigrationExecutorTest { + + private static final String TENANT = "tenant-a"; + + private static final String V1 = "/orders/migrations/001__backfill_status.migration"; + private static final String V2 = "/orders/migrations/002__close_old.migration"; + + private final MigrationLedger ledger = new MigrationLedger(); + private final MigrationExecutor executor = new MigrationExecutor(ledger); + + private JdbcDataSource dataSource; + + @BeforeEach + void createTheTable() throws SQLException { + dataSource = new JdbcDataSource(); + dataSource.setURL("jdbc:h2:mem:migration_executor_test;DB_CLOSE_DELAY=-1"); + execute("CREATE TABLE ORDERS (ID INT PRIMARY KEY, STATUS VARCHAR(20))"); + execute("INSERT INTO ORDERS VALUES (1, NULL), (2, 'CLOSED')"); + } + + @AfterEach + void dropEverything() throws SQLException { + execute("DROP ALL OBJECTS"); + } + + @Test + void aMigrationRunsAndIsRecordedOnce() throws Exception { + MigrationScript script = script(V1, "UPDATE ORDERS SET STATUS = 'OPEN' WHERE STATUS IS NULL;"); + + assertThat(executor.apply(script, V1, TENANT, dataSource, List.of()) + .status()).isEqualTo(Status.APPLIED); + assertThat(status(1)).isEqualTo("OPEN"); + + // What a second boot, a new tenant's re-trigger or a retry does: nothing. + execute("UPDATE ORDERS SET STATUS = NULL WHERE ID = 1"); + assertThat(executor.apply(script, V1, TENANT, dataSource, List.of()) + .status()).isEqualTo(Status.ALREADY_APPLIED); + assertThat(status(1)).isNull(); + + List entries = ledger.findAll(dataSource); + assertThat(entries).singleElement() + .satisfies(entry -> { + assertThat(entry.project()).isEqualTo("orders"); + assertThat(entry.version()).isEqualTo("001"); + assertThat(entry.location()).isEqualTo(V1); + assertThat(entry.checksum()).isEqualTo(script.checksum()); + assertThat(entry.tenant()).isEqualTo(TENANT); + assertThat(entry.appliedAt()).isNotNull(); + }); + } + + @Test + void anAppliedMigrationEditedAfterwardsFailsAndDoesNotRun() throws Exception { + executor.apply(script(V1, "UPDATE ORDERS SET STATUS = 'OPEN' WHERE STATUS IS NULL;"), V1, TENANT, dataSource, List.of()); + execute("UPDATE ORDERS SET STATUS = NULL WHERE ID = 1"); + + Outcome outcome = + executor.apply(script(V1, "UPDATE ORDERS SET STATUS = 'NEW' WHERE STATUS IS NULL;"), V1, TENANT, dataSource, List.of()); + + assertThat(outcome.status()).isEqualTo(Status.FAILED); + assertThat(outcome.message()).contains(V1, "never re-run", "idempotent: true"); + assertThat(status(1)).isNull(); + } + + @Test + void anIdempotentMigrationEditedAfterwardsRunsAgainAndIsRecordedWithItsNewContent() throws Exception { + executor.apply(script(V1, "-- idempotent: true\nUPDATE ORDERS SET STATUS = 'OPEN' WHERE STATUS IS NULL;"), V1, TENANT, dataSource, + List.of()); + execute("UPDATE ORDERS SET STATUS = NULL WHERE ID = 1"); + MigrationScript edited = script(V1, "-- idempotent: true\nUPDATE ORDERS SET STATUS = 'NEW' WHERE STATUS IS NULL;"); + + assertThat(executor.apply(edited, V1, TENANT, dataSource, List.of()) + .status()).isEqualTo(Status.APPLIED); + + assertThat(status(1)).isEqualTo("NEW"); + assertThat(ledger.findAll(dataSource)).singleElement() + .extracting(Entry::checksum) + .isEqualTo(edited.checksum()); + } + + @Test + void aFailingStatementRollsBackTheWholeMigrationAndRecordsNothing() throws Exception { + MigrationScript script = script(V1, """ + UPDATE ORDERS SET STATUS = 'OPEN' WHERE STATUS IS NULL; + UPDATE NO_SUCH_TABLE SET STATUS = 'OPEN'; + """); + + Outcome outcome = executor.apply(script, V1, TENANT, dataSource, List.of()); + + assertThat(outcome.status()).isEqualTo(Status.FAILED); + assertThat(outcome.message()).contains(V1, TENANT, "NO_SUCH_TABLE"); + assertThat(outcome.cause()).isNotNull(); + assertThat(status(1)).as("the first statement is rolled back with the second") + .isNull(); + assertThat(ledger.findAll(dataSource)).isEmpty(); + } + + @Test + void aLaterVersionWaitsForAnEarlierOneThatHasNotAppliedYet() throws Exception { + MigrationScript first = script(V1, "UPDATE ORDERS SET STATUS = 'OPEN' WHERE STATUS IS NULL;"); + MigrationScript second = script(V2, "UPDATE ORDERS SET STATUS = 'ARCHIVED' WHERE STATUS = 'OPEN';"); + Migration pendingFirst = artefact(V1, first, ArtefactLifecycle.NEW); + + Outcome waiting = executor.apply(second, V2, TENANT, dataSource, List.of(pendingFirst)); + + assertThat(waiting.status()).isEqualTo(Status.WAITING); + assertThat(waiting.message()).contains("waits for version [001]"); + assertThat(status(1)).isNull(); + + executor.apply(first, V1, TENANT, dataSource, List.of()); + assertThat(executor.apply(second, V2, TENANT, dataSource, List.of(pendingFirst)) + .status()).isEqualTo(Status.APPLIED); + assertThat(status(1)).isEqualTo("ARCHIVED"); + } + + @Test + void aLaterVersionFailsWhenTheEarlierOneFailed() throws Exception { + MigrationScript first = script(V1, "UPDATE NO_SUCH_TABLE SET STATUS = 'OPEN';"); + MigrationScript second = script(V2, "UPDATE ORDERS SET STATUS = 'ARCHIVED';"); + + Outcome outcome = executor.apply(second, V2, TENANT, dataSource, List.of(artefact(V1, first, ArtefactLifecycle.FAILED))); + + assertThat(outcome.status()).isEqualTo(Status.FAILED); + assertThat(outcome.message()).contains("waits for version [001]", "which failed"); + assertThat(status(2)).isEqualTo("CLOSED"); + } + + private static MigrationScript script(String location, String sql) throws ParseException { + return MigrationScript.parse(location, sql.getBytes(StandardCharsets.UTF_8)); + } + + private static Migration artefact(String location, MigrationScript script, ArtefactLifecycle lifecycle) { + Migration migration = new Migration(location, location.substring(location.lastIndexOf('/') + 1), script); + migration.setLifecycle(lifecycle); + return migration; + } + + private String status(int id) throws SQLException { + try (Connection connection = dataSource.getConnection(); + Statement statement = connection.createStatement(); + ResultSet resultSet = statement.executeQuery("SELECT STATUS FROM ORDERS WHERE ID = " + id)) { + assertThat(resultSet.next()).isTrue(); + return resultSet.getString(1); + } + } + + private void execute(String sql) throws SQLException { + try (Connection connection = dataSource.getConnection(); Statement statement = connection.createStatement()) { + statement.execute(sql); + } + } +} diff --git a/components/data/data-migrations/src/test/java/org/eclipse/dirigible/components/data/migrations/MigrationScriptTest.java b/components/data/data-migrations/src/test/java/org/eclipse/dirigible/components/data/migrations/MigrationScriptTest.java new file mode 100644 index 00000000000..d564e3021b1 --- /dev/null +++ b/components/data/data-migrations/src/test/java/org/eclipse/dirigible/components/data/migrations/MigrationScriptTest.java @@ -0,0 +1,115 @@ +/* + * Copyright (c) 2010-2026 Eclipse Dirigible contributors + * + * All rights reserved. This program and the accompanying materials are made available under the + * terms of the Eclipse Public License v2.0 which accompanies this distribution, and is available at + * http://www.eclipse.org/legal/epl-v20.html + * + * SPDX-FileCopyrightText: Eclipse Dirigible contributors SPDX-License-Identifier: EPL-2.0 + */ +package org.eclipse.dirigible.components.data.migrations; + +import static org.assertj.core.api.Assertions.assertThat; +import static org.assertj.core.api.Assertions.assertThatThrownBy; + +import java.nio.charset.StandardCharsets; +import java.text.ParseException; + +import org.eclipse.dirigible.components.data.migrations.MigrationScript.Scope; +import org.junit.jupiter.api.Test; + +class MigrationScriptTest { + + private static MigrationScript parse(String location, String content) throws ParseException { + return MigrationScript.parse(location, content.getBytes(StandardCharsets.UTF_8)); + } + + @Test + void theFileNameCarriesProjectVersionAndDescription() throws ParseException { + MigrationScript script = parse("/orders/migrations/V1.2__backfill_status.migration", "UPDATE ORDERS SET STATUS = 'OPEN';"); + + assertThat(script.project()).isEqualTo("orders"); + assertThat(script.version()).isEqualTo("1.2"); + assertThat(script.description()).isEqualTo("backfill_status"); + assertThat(script.sql()).isEqualTo("UPDATE ORDERS SET STATUS = 'OPEN';"); + } + + @Test + void withoutHeadersAMigrationAppliesToEachTenantAndIsNotIdempotent() throws ParseException { + MigrationScript script = parse("/orders/001__backfill.migration", "-- recomputes the totals\nUPDATE ORDERS SET TOTAL = 0;"); + + assertThat(script.scope()).isEqualTo(Scope.EACH); + assertThat(script.idempotent()).isFalse(); + } + + @Test + void theHeadersDeclareScopeAndIdempotence() throws ParseException { + MigrationScript script = parse("/orders/001__backfill.migration", """ + -- Moves the legacy rows. + -- tenant: system + -- IDEMPOTENT: True + + UPDATE ORDERS SET TOTAL = 0; + -- tenant: each (a comment after the first statement is SQL, not a header) + """); + + assertThat(script.scope()).isEqualTo(Scope.SYSTEM); + assertThat(script.idempotent()).isTrue(); + } + + @Test + void anUnknownHeaderValueIsRejected() { + assertThatThrownBy( + () -> parse("/orders/001__backfill.migration", "-- tenant: some\nUPDATE ORDERS SET TOTAL = 0;")).isInstanceOf( + ParseException.class) + .hasMessageContaining( + "[tenant: some]"); + assertThatThrownBy(() -> parse("/orders/001__backfill.migration", + "-- idempotent: yes\nUPDATE ORDERS SET TOTAL = 0;")).isInstanceOf(ParseException.class) + .hasMessageContaining("[idempotent: yes]"); + } + + @Test + void aFileNameWithoutAVersionIsRejected() { + assertThatThrownBy(() -> parse("/orders/backfill.migration", "UPDATE ORDERS SET TOTAL = 0;")).isInstanceOf(ParseException.class) + .hasMessageContaining( + "__.migration"); + assertThatThrownBy(() -> parse("/orders/1a__backfill.migration", "UPDATE ORDERS SET TOTAL = 0;")).isInstanceOf( + ParseException.class); + } + + @Test + void aMigrationOutsideAProjectIsRejected() { + assertThatThrownBy(() -> parse("/001__backfill.migration", "UPDATE ORDERS SET TOTAL = 0;")).isInstanceOf(ParseException.class) + .hasMessageContaining( + "inside a project"); + } + + @Test + void aMigrationWithoutAStatementIsRejected() { + assertThatThrownBy( + () -> parse("/orders/001__backfill.migration", "-- tenant: each\n\n-- nothing yet\n")).isInstanceOf(ParseException.class) + .hasMessageContaining( + "no SQL statement"); + } + + @Test + void convertedLineEndingsAreNotAnEdit() throws ParseException { + MigrationScript unix = parse("/orders/001__backfill.migration", "UPDATE ORDERS\nSET TOTAL = 0;\n"); + MigrationScript windows = parse("/orders/001__backfill.migration", "UPDATE ORDERS\r\nSET TOTAL = 0;\r\n"); + MigrationScript edited = parse("/orders/001__backfill.migration", "UPDATE ORDERS\nSET TOTAL = 1;\n"); + + assertThat(windows.checksum()).isEqualTo(unix.checksum()); + assertThat(edited.checksum()).isNotEqualTo(unix.checksum()); + } + + @Test + void versionsCompareNumericallySegmentBySegment() { + assertThat(MigrationScript.compareVersions("2", "10")).isNegative(); + assertThat(MigrationScript.compareVersions("1.2", "1.10")).isNegative(); + assertThat(MigrationScript.compareVersions("001", "1")).isZero(); + assertThat(MigrationScript.compareVersions("1", "1.0")).isNegative(); + assertThat(MigrationScript.compareVersions("20261005120000", "20261005115959")).isPositive(); + assertThat(MigrationScript.compareVersions("99999999999999999999", "100000000000000000000")).isNegative(); + } +} diff --git a/components/data/data-migrations/src/test/java/org/eclipse/dirigible/components/data/migrations/MigrationsSynchronizerTest.java b/components/data/data-migrations/src/test/java/org/eclipse/dirigible/components/data/migrations/MigrationsSynchronizerTest.java new file mode 100644 index 00000000000..54cb3ec83ad --- /dev/null +++ b/components/data/data-migrations/src/test/java/org/eclipse/dirigible/components/data/migrations/MigrationsSynchronizerTest.java @@ -0,0 +1,229 @@ +/* + * Copyright (c) 2010-2026 Eclipse Dirigible contributors + * + * All rights reserved. This program and the accompanying materials are made available under the + * terms of the Eclipse Public License v2.0 which accompanies this distribution, and is available at + * http://www.eclipse.org/legal/epl-v20.html + * + * SPDX-FileCopyrightText: Eclipse Dirigible contributors SPDX-License-Identifier: EPL-2.0 + */ +package org.eclipse.dirigible.components.data.migrations; + +import static org.assertj.core.api.Assertions.assertThat; +import static org.mockito.ArgumentMatchers.any; +import static org.mockito.ArgumentMatchers.anyList; +import static org.mockito.ArgumentMatchers.anyString; +import static org.mockito.ArgumentMatchers.contains; +import static org.mockito.ArgumentMatchers.eq; +import static org.mockito.Mockito.mock; +import static org.mockito.Mockito.never; +import static org.mockito.Mockito.verify; +import static org.mockito.Mockito.verifyNoInteractions; +import static org.mockito.Mockito.when; + +import java.nio.charset.StandardCharsets; +import java.text.ParseException; +import java.util.ArrayList; +import java.util.HashMap; +import java.util.List; +import java.util.Map; +import java.util.concurrent.atomic.AtomicReference; + +import org.eclipse.dirigible.components.base.artefact.ArtefactLifecycle; +import org.eclipse.dirigible.components.base.artefact.ArtefactPhase; +import org.eclipse.dirigible.components.base.artefact.topology.TopologyWrapper; +import org.eclipse.dirigible.components.base.callable.CallableResultAndException; +import org.eclipse.dirigible.components.base.synchronizer.SynchronizerCallback; +import org.eclipse.dirigible.components.base.tenant.Tenant; +import org.eclipse.dirigible.components.base.tenant.TenantContext; +import org.eclipse.dirigible.components.base.tenant.TenantResult; +import org.eclipse.dirigible.components.data.migrations.MigrationExecutor.Outcome; +import org.eclipse.dirigible.components.data.migrations.MigrationExecutor.Status; +import org.eclipse.dirigible.components.data.sources.manager.DataSourcesManager; +import org.eclipse.dirigible.components.database.DirigibleDataSource; +import org.eclipse.dirigible.repository.api.IRepository; +import org.eclipse.dirigible.repository.api.IRepositoryStructure; +import org.eclipse.dirigible.repository.api.IResource; +import org.junit.jupiter.api.BeforeEach; +import org.junit.jupiter.api.Test; +import org.mockito.ArgumentCaptor; + +/** How the per-database outcomes become the ONE state of the artefact. */ +class MigrationsSynchronizerTest { + + private static final String V1 = "/orders/migrations/001__first.migration"; + private static final String V2 = "/orders/migrations/002__second.migration"; + private static final String V3 = "/orders/migrations/003__third.migration"; + + private final MigrationService migrationService = mock(MigrationService.class); + private final MigrationExecutor executor = mock(MigrationExecutor.class); + private final DataSourcesManager dataSourcesManager = mock(DataSourcesManager.class); + private final TenantContext tenantContext = mock(TenantContext.class); + private final IRepository repository = mock(IRepository.class); + private final SynchronizerCallback callback = mock(SynchronizerCallback.class); + + private final DirigibleDataSource tenantDataSource = mock(DirigibleDataSource.class); + private final DirigibleDataSource systemDataSource = mock(DirigibleDataSource.class); + + private final Map registry = new HashMap<>(); + + private MigrationsSynchronizer synchronizer; + + @BeforeEach + void setUp() throws Exception { + synchronizer = new MigrationsSynchronizer(migrationService, executor, dataSourcesManager, tenantContext, repository); + synchronizer.setCallback(callback); + when(dataSourcesManager.getDefaultDataSource()).thenReturn(tenantDataSource); + when(dataSourcesManager.getSystemDataSource()).thenReturn(systemDataSource); + when(repository.getResource(anyString())).thenAnswer(invocation -> resource(invocation.getArgument(0))); + givenTenants("tenant-a", "tenant-b"); + } + + @Test + void oneTenantsFailureIsNotHiddenByAnotherTenantsSuccess() throws Exception { + Migration migration = published(V1, "UPDATE ORDERS SET STATUS = 'OPEN';", ArtefactLifecycle.NEW); + givenOutcome("tenant-a", new Outcome("tenant-a", Status.APPLIED, null)); + givenOutcome("tenant-b", new Outcome("tenant-b", Status.FAILED, "boom in tenant-b")); + + boolean depleted = synchronizer.completeImpl(wrap(migration), ArtefactPhase.CREATE); + + assertThat(depleted).isTrue(); + verify(callback).registerState(eq(synchronizer), any(TopologyWrapper.class), eq(ArtefactLifecycle.FAILED), + contains("boom in tenant-b"), any()); + verify(callback, never()).registerState(any(), any(TopologyWrapper.class), eq(ArtefactLifecycle.CREATED)); + } + + @Test + void aMigrationEveryTenantTookIsCreated() throws Exception { + Migration migration = published(V1, "UPDATE ORDERS SET STATUS = 'OPEN';", ArtefactLifecycle.NEW); + givenOutcome("tenant-a", new Outcome("tenant-a", Status.APPLIED, null)); + givenOutcome("tenant-b", new Outcome("tenant-b", Status.ALREADY_APPLIED, null)); + + assertThat(synchronizer.completeImpl(wrap(migration), ArtefactPhase.CREATE)).isTrue(); + + verify(callback).registerState(eq(synchronizer), any(TopologyWrapper.class), eq(ArtefactLifecycle.CREATED)); + } + + @Test + void aMigrationWaitingForAnEarlierVersionStaysUndepletedWithItsLifecycle() throws Exception { + Migration migration = published(V2, "UPDATE ORDERS SET STATUS = 'OPEN';", ArtefactLifecycle.NEW); + givenOutcome("tenant-a", new Outcome("tenant-a", Status.APPLIED, null)); + givenOutcome("tenant-b", new Outcome("tenant-b", Status.WAITING, "waits for version [001]")); + + boolean depleted = synchronizer.completeImpl(wrap(migration), ArtefactPhase.CREATE); + + assertThat(depleted).isFalse(); + assertThat(migration.getLifecycle()).isEqualTo(ArtefactLifecycle.NEW); + assertThat(migration.getError()).contains("waits for version [001]"); + verifyNoInteractions(callback); + } + + @Test + void aSystemMigrationRunsOnceOnTheSystemDatabase() throws Exception { + Migration migration = published(V1, "-- tenant: system\nUPDATE ORDERS SET STATUS = 'OPEN';", ArtefactLifecycle.NEW); + when(executor.apply(any(), eq(V1), eq(MigrationsSynchronizer.SYSTEM_TENANT), eq(systemDataSource), anyList())).thenReturn( + new Outcome(MigrationsSynchronizer.SYSTEM_TENANT, Status.APPLIED, null)); + + assertThat(synchronizer.completeImpl(wrap(migration), ArtefactPhase.CREATE)).isTrue(); + + verify(tenantContext, never()).executeForEachTenant(any()); + verify(callback).registerState(eq(synchronizer), any(TopologyWrapper.class), eq(ArtefactLifecycle.CREATED)); + } + + @Test + void thePredecessorsAreTheLowerVersionsOfTheSameProjectAndScope() throws Exception { + Migration first = published(V1, "UPDATE ORDERS SET STATUS = 'A';", ArtefactLifecycle.CREATED); + Migration systemZero = + published("/orders/migrations/000__system.migration", "-- tenant: system\nUPDATE X SET Y = 1;", ArtefactLifecycle.CREATED); + Migration third = published(V3, "UPDATE ORDERS SET STATUS = 'C';", ArtefactLifecycle.NEW); + Migration second = published(V2, "UPDATE ORDERS SET STATUS = 'B';", ArtefactLifecycle.NEW); + when(migrationService.findAllByProject("orders")).thenReturn(List.of(first, systemZero, third, second)); + when(executor.apply(any(), any(), any(), any(), anyList())).thenAnswer( + invocation -> new Outcome(invocation.getArgument(2), Status.APPLIED, null)); + + synchronizer.completeImpl(wrap(second), ArtefactPhase.CREATE); + + @SuppressWarnings("unchecked") + ArgumentCaptor> predecessors = ArgumentCaptor.forClass(List.class); + verify(executor, org.mockito.Mockito.atLeastOnce()).apply(any(), eq(V2), any(), any(), predecessors.capture()); + assertThat(predecessors.getValue()).containsExactly(first); + } + + @Test + void aMigrationOnlyItsPhaseMatchesIsTouched() throws Exception { + Migration applied = published(V1, "UPDATE ORDERS SET STATUS = 'OPEN';", ArtefactLifecycle.CREATED); + + assertThat(synchronizer.completeImpl(wrap(applied), ArtefactPhase.CREATE)).isTrue(); + assertThat(synchronizer.completeImpl(wrap(applied), ArtefactPhase.UPDATE)).isTrue(); + assertThat(synchronizer.completeImpl(wrap(applied), ArtefactPhase.START)).isTrue(); + + verifyNoInteractions(executor); + } + + @Test + void anApplyAfterTheFileWasRemovedFailsTheArtefact() throws Exception { + Migration migration = published(V1, "UPDATE ORDERS SET STATUS = 'OPEN';", ArtefactLifecycle.FAILED); + registry.remove(V1); + + assertThat(synchronizer.completeImpl(wrap(migration), ArtefactPhase.START)).isTrue(); + + verify(callback).registerState(eq(synchronizer), any(TopologyWrapper.class), eq(ArtefactLifecycle.FAILED), + contains("no longer in the registry"), any(ParseException.class)); + verifyNoInteractions(executor); + } + + private Migration published(String location, String sql, ArtefactLifecycle lifecycle) throws ParseException { + registry.put(location, sql); + Migration migration = new Migration(location, location.substring(location.lastIndexOf('/') + 1), + MigrationScript.parse(location, sql.getBytes(StandardCharsets.UTF_8))); + migration.setLifecycle(lifecycle); + return migration; + } + + private TopologyWrapper wrap(Migration migration) { + return new TopologyWrapper<>(migration, new HashMap<>(), synchronizer); + } + + private IResource resource(String path) { + String location = path.substring(IRepositoryStructure.PATH_REGISTRY_PUBLIC.length()); + IResource resource = mock(IResource.class); + String content = registry.get(location); + when(resource.exists()).thenReturn(content != null); + when(resource.getContent()).thenReturn(content == null ? null : content.getBytes(StandardCharsets.UTF_8)); + return resource; + } + + private void givenOutcome(String tenant, Outcome outcome) { + when(executor.apply(any(), any(), eq(tenant), eq(tenantDataSource), anyList())).thenReturn(outcome); + } + + /** Runs the callable once per tenant with that tenant current, as the platform's context does. */ + @SuppressWarnings("unchecked") + private void givenTenants(String... ids) throws Exception { + AtomicReference current = new AtomicReference<>(); + when(tenantContext.getCurrentTenant()).thenAnswer(invocation -> current.get()); + when(tenantContext.executeForEachTenant(any())).thenAnswer(invocation -> { + CallableResultAndException callable = invocation.getArgument(0); + List> results = new ArrayList<>(); + for (String id : ids) { + Tenant tenant = mock(Tenant.class); + when(tenant.getId()).thenReturn(id); + current.set(tenant); + Object result = callable.call(); + results.add(new TenantResult<>() { + @Override + public Tenant getTenant() { + return tenant; + } + + @Override + public Object getResult() { + return result; + } + }); + } + current.set(null); + return results; + }); + } +} diff --git a/components/group/group-database/pom.xml b/components/group/group-database/pom.xml index 9b9f99919ef..6eeb9841a3d 100644 --- a/components/group/group-database/pom.xml +++ b/components/group/group-database/pom.xml @@ -39,6 +39,10 @@ org.eclipse.dirigible dirigible-components-data-csvim + + org.eclipse.dirigible + dirigible-components-data-migrations + org.eclipse.dirigible dirigible-components-data-export diff --git a/components/pom.xml b/components/pom.xml index ddd479df2f8..20825b16e1f 100644 --- a/components/pom.xml +++ b/components/pom.xml @@ -53,6 +53,7 @@ data/data-store data/data-store-java data/data-csvim + data/data-migrations data/data-export data/data-import data/data-anonymize @@ -646,6 +647,11 @@ dirigible-components-data-csvim ${project.version} + + org.eclipse.dirigible + dirigible-components-data-migrations + ${project.version} + org.eclipse.dirigible dirigible-components-data-transfer diff --git a/tests/tests-integrations/src/main/java/org/eclipse/dirigible/integration/tests/api/DataMigrationIT.java b/tests/tests-integrations/src/main/java/org/eclipse/dirigible/integration/tests/api/DataMigrationIT.java new file mode 100644 index 00000000000..43b165ff019 --- /dev/null +++ b/tests/tests-integrations/src/main/java/org/eclipse/dirigible/integration/tests/api/DataMigrationIT.java @@ -0,0 +1,255 @@ +/* + * Copyright (c) 2010-2026 Eclipse Dirigible contributors + * + * All rights reserved. This program and the accompanying materials are made available under the + * terms of the Eclipse Public License v2.0 which accompanies this distribution, and is available at + * http://www.eclipse.org/legal/epl-v20.html + * + * SPDX-FileCopyrightText: Eclipse Dirigible contributors SPDX-License-Identifier: EPL-2.0 + */ +package org.eclipse.dirigible.integration.tests.api; + +import static io.restassured.RestAssured.given; +import static org.assertj.core.api.Assertions.assertThat; +import static org.hamcrest.Matchers.containsInAnyOrder; +import static org.hamcrest.Matchers.equalTo; +import static org.hamcrest.Matchers.hasKey; +import static org.hamcrest.Matchers.not; + +import java.nio.charset.StandardCharsets; +import java.sql.Connection; +import java.sql.PreparedStatement; +import java.sql.ResultSet; +import java.sql.SQLException; +import java.sql.Statement; +import java.util.ArrayList; +import java.util.List; +import java.util.concurrent.TimeUnit; + +import org.awaitility.Awaitility; +import org.eclipse.dirigible.components.base.readiness.PlatformReadiness; +import org.eclipse.dirigible.components.base.tenant.TenantContext; +import org.eclipse.dirigible.components.data.sources.manager.DataSourcesManager; +import org.eclipse.dirigible.components.initializers.synchronizer.SynchronizationProcessor; +import org.eclipse.dirigible.repository.api.IRepository; +import org.eclipse.dirigible.repository.api.IRepositoryStructure; +import org.eclipse.dirigible.tests.base.IntegrationTest; +import org.eclipse.dirigible.tests.framework.restassured.RestAssuredExecutor; +import org.eclipse.dirigible.tests.framework.tenant.DirigibleTestTenant; +import org.junit.jupiter.api.Test; +import org.springframework.beans.factory.annotation.Autowired; + +/** + * A {@code .migration} applies exactly once per tenant schema, after the table it changes is + * synchronized, recorded in that schema's {@code DIRIGIBLE_MIGRATIONS} (#7636): + *

    + *
  1. a project's table and two migrations land in the default tenant and a provisioned tenant, the + * second migration (a backfill) after the first (the rows it backfills);
  2. + *
  3. every multitenant artefact is synchronized again - what provisioning a new tenant does, and + * the same re-parse a boot against a fresh system database runs - and nothing applies twice, while + * the new tenant is migrated;
  4. + *
  5. an applied file edited afterwards is a FAILED artefact - counted by the health endpoint and + * keeping the boot from being clean - and restoring it heals it;
  6. + *
  7. the ledger endpoint lists what applied.
  8. + *
+ * Pure HTTP / JDBC - no Selenide, no IDE. + */ +class DataMigrationIT extends IntegrationTest { + + private static final String PROJECT = "data-migration-it"; + + private static final String TABLE_PATH = registryPath("/tables/order.table"); + private static final String SEED_PATH = registryPath("/migrations/V1__seed_orders.migration"); + private static final String BACKFILL_LOCATION = "/" + PROJECT + "/migrations/V2__backfill_status.migration"; + private static final String BACKFILL_PATH = IRepositoryStructure.PATH_REGISTRY_PUBLIC + BACKFILL_LOCATION; + + private static final String TABLE_SOURCE = """ + { + "name": "DATA_MIGRATION_IT_ORDER", + "type": "TABLE", + "columns": [ + { + "type": "INTEGER", + "primaryKey": true, + "nullable": false, + "name": "ORDER_ID" + }, + { + "type": "VARCHAR", + "length": 20, + "nullable": true, + "name": "ORDER_STATUS" + } + ] + } + """; + + private static final String SEED_SOURCE = """ + -- The rows a later migration backfills. + INSERT INTO "DATA_MIGRATION_IT_ORDER" ("ORDER_ID", "ORDER_STATUS") VALUES (1, NULL); + INSERT INTO "DATA_MIGRATION_IT_ORDER" ("ORDER_ID", "ORDER_STATUS") VALUES (2, 'CLOSED'); + INSERT INTO "DATA_MIGRATION_IT_ORDER" ("ORDER_ID", "ORDER_STATUS") VALUES (3, NULL); + """; + + private static final String BACKFILL_SOURCE = """ + -- tenant: each + UPDATE "DATA_MIGRATION_IT_ORDER" SET "ORDER_STATUS" = 'OPEN' WHERE "ORDER_STATUS" IS NULL; + """; + + private static final String BACKFILL_EDITED = """ + -- tenant: each + UPDATE "DATA_MIGRATION_IT_ORDER" SET "ORDER_STATUS" = 'NEW' WHERE "ORDER_STATUS" IS NULL; + """; + + private static final long TIMEOUT_SECONDS = 120; + + @Autowired + private IRepository repository; + + @Autowired + private SynchronizationProcessor synchronizationProcessor; + + @Autowired + private DataSourcesManager dataSourcesManager; + + @Autowired + private TenantContext tenantContext; + + @Autowired + private RestAssuredExecutor restAssuredExecutor; + + @Test + void aMigrationAppliesOncePerTenantAndAnEditAfterwardsFailsIt() throws Exception { + DirigibleTestTenant tenant = new DirigibleTestTenant(PROJECT); + createTenants(tenant); + waitForTenantProvisioning(tenant); + + write(TABLE_PATH, TABLE_SOURCE); + write(SEED_PATH, SEED_SOURCE); + write(BACKFILL_PATH, BACKFILL_SOURCE); + synchronizationProcessor.forceProcessSynchronizers(); + + assertThat(defaultTenant(this::ledgerVersions)).containsExactlyInAnyOrder("1", "2"); + assertThat(defaultTenant(this::statuses)).containsExactly("OPEN", "CLOSED", "OPEN"); + assertThat(inTenant(tenant, this::ledgerVersions)).containsExactlyInAnyOrder("1", "2"); + assertThat(inTenant(tenant, this::statuses)).containsExactly("OPEN", "CLOSED", "OPEN"); + + // A row the backfill would rewrite if it ran again. + defaultTenant(() -> execute("UPDATE \"DATA_MIGRATION_IT_ORDER\" SET \"ORDER_STATUS\" = NULL WHERE \"ORDER_ID\" = 3")); + inTenant(tenant, () -> execute("UPDATE \"DATA_MIGRATION_IT_ORDER\" SET \"ORDER_STATUS\" = NULL WHERE \"ORDER_ID\" = 3")); + + // Provisioning a tenant re-synchronizes every multitenant artefact in every tenant. + DirigibleTestTenant newTenant = new DirigibleTestTenant(PROJECT + "-late"); + createTenants(newTenant); + waitForTenantProvisioning(newTenant); + Awaitility.await() + .pollInterval(2, TimeUnit.SECONDS) + .atMost(TIMEOUT_SECONDS, TimeUnit.SECONDS) + .until(() -> inTenant(newTenant, this::ledgerVersions).size() == 2); + synchronizationProcessor.forceProcessSynchronizers(); + + assertThat(inTenant(newTenant, this::statuses)).containsExactly("OPEN", "CLOSED", "OPEN"); + assertThat(defaultTenant(this::ledgerVersions)).as("one ledger row per migration, however often it was synchronized") + .containsExactlyInAnyOrder("1", "2"); + assertThat(inTenant(tenant, this::ledgerVersions)).containsExactlyInAnyOrder("1", "2"); + assertThat(defaultTenant(this::statuses)).as("an applied migration never runs again") + .containsExactly("OPEN", "CLOSED", null); + assertThat(inTenant(tenant, this::statuses)).containsExactly("OPEN", "CLOSED", null); + + restAssuredExecutor.execute(() -> given().when() + .get("/services/core/migrations") + .then() + .statusCode(200) + .body("findAll { it.project == '" + PROJECT + "' }.version", containsInAnyOrder("1", "2")) + .body("find { it.version == '2' && it.project == '" + PROJECT + "' }.location", + equalTo(BACKFILL_LOCATION)), + TIMEOUT_SECONDS); + + // An applied migration edited afterwards: FAILED, and nothing runs. + write(BACKFILL_PATH, BACKFILL_EDITED); + synchronizationProcessor.forceProcessSynchronizers(); + + restAssuredExecutor.execute(() -> given().when() + .get("/actuator/health") + .then() + .statusCode(200) + .body("components.migrations.status", equalTo("UP")) + .body("components.migrations.details.failed", equalTo(1)) + .body("components.migrations.details.pending", equalTo(0)) + .body("components.artefacts.details.failedByType.migration", equalTo(1)), + TIMEOUT_SECONDS); + assertThat(PlatformReadiness.getInstance() + .isCleanBoot()).as("DIRIGIBLE_READINESS_REQUIRE_CLEAN_BOOT withholds readiness on a failed migration") + .isFalse(); + assertThat(defaultTenant(this::statuses)).containsExactly("OPEN", "CLOSED", null); + + // Restoring the file heals it - the ledger already has it. + write(BACKFILL_PATH, BACKFILL_SOURCE); + synchronizationProcessor.forceProcessSynchronizers(); + + restAssuredExecutor.execute(() -> given().when() + .get("/actuator/health") + .then() + .statusCode(200) + .body("components.migrations.details.failed", equalTo(0)) + .body("components.artefacts.details.failedByType", not(hasKey("migration"))), + TIMEOUT_SECONDS); + assertThat(defaultTenant(this::statuses)).containsExactly("OPEN", "CLOSED", null); + } + + private List ledgerVersions() throws SQLException { + return column("SELECT \"MIGRATION_VERSION\" FROM \"DIRIGIBLE_MIGRATIONS\" WHERE \"MIGRATION_PROJECT\" = ?", PROJECT); + } + + private List statuses() throws SQLException { + return column("SELECT \"ORDER_STATUS\" FROM \"DATA_MIGRATION_IT_ORDER\" ORDER BY \"ORDER_ID\""); + } + + private List column(String sql, String... parameters) throws SQLException { + try (Connection connection = dataSourcesManager.getDefaultDataSource() + .getConnection(); + PreparedStatement statement = connection.prepareStatement(sql)) { + for (int i = 0; i < parameters.length; i++) { + statement.setString(i + 1, parameters[i]); + } + List values = new ArrayList<>(); + try (ResultSet resultSet = statement.executeQuery()) { + while (resultSet.next()) { + values.add(resultSet.getString(1)); + } + } + return values; + } + } + + private Void execute(String sql) throws SQLException { + try (Connection connection = dataSourcesManager.getDefaultDataSource() + .getConnection(); + Statement statement = connection.createStatement()) { + statement.executeUpdate(sql); + } + return null; + } + + /** Outside a tenant scope the default datasource is the default tenant's database. */ + private static T defaultTenant(SqlCall call) throws SQLException { + return call.call(); + } + + private T inTenant(DirigibleTestTenant tenant, SqlCall call) throws SQLException { + return tenantContext.execute(tenant.getId(), call::call); + } + + private void write(String path, String source) { + repository.createResource(path, source.getBytes(StandardCharsets.UTF_8), false, "text/plain", true); + } + + private static String registryPath(String path) { + return IRepositoryStructure.PATH_REGISTRY_PUBLIC + "/" + PROJECT + path; + } + + @FunctionalInterface + private interface SqlCall { + T call() throws SQLException; + } +}