From 047560310308c8d708208f4d69b2770f44b3dcc6 Mon Sep 17 00:00:00 2001 From: Nikol Georgieva Date: Mon, 5 Oct 2026 10:42:02 +0300 Subject: [PATCH 1/2] migrations: a versioned .migration artefact applies once per tenant schema after table sync, recorded in a DIRIGIBLE_MIGRATIONS ledger (#7636) Structure evolves on publish and CSVIM seeds reference data, but there was no artefact for a data migration - a backfill, a split field, a recomputed roll-up - so such changes were hand-run SQL per instance and per tenant schema, or a job inventing its own "already ran" marker. A new components/data/data-migrations module adds the .migration artefact: /.../__.migration, SQL, optionally headed by "-- tenant: each|system" (default each) and "-- idempotent: true|false" (default false). MigrationsSynchronizer runs at SynchronizersOrder.MIGRATION (260): after schema/table/view/entity, before BPMN and CSVIM. - Each migration runs in one transaction together with the row recording it in DIRIGIBLE_MIGRATIONS (project, version, location, checksum, tenant, applied at, duration). The ledger lives in the database the migration changes - every tenant schema, or the system database for tenant: system - and is keyed by project/version, so a concurrent second node's insert fails and rolls its run back. - The ledger decides, not the artefact lifecycle, so every re-run is safe: the post-provisioning re-trigger migrates the new tenant and finds the others done; a fresh system database re-parses without re-applying. - An applied file whose checksum changed is a FAILED artefact (not a re-run) unless it declares idempotent: true. Checksums ignore line-ending conversion. - A project's migrations apply in version order (numeric, segment-wise): a later version waits for an earlier one still applying in the pass, and fails if the earlier one failed. - The synchronizer iterates the tenants itself and registers one state for all of them; the per-tenant completion of BaseSynchronizer keeps the state the LAST tenant registered, so one tenant's success would hide another's failure. - A failed migration is a failed artefact, so the census already counts it and DIRIGIBLE_READINESS_REQUIRE_CLEAN_BOOT withholds readiness. The new "migrations" health component reports total / pending / failed, and GET /services/core/migrations lists the caller's tenant ledger (plus the system ledger for the default tenant). Verified: data-migrations unit tests (22, H2-backed executor tests), the core-base and core-liquibase suites, DataMigrationIT on H2 and on PostgreSQL 16, formatter:validate with the cache wiped, and the release javadoc build on the touched modules. Co-Authored-By: Claude Opus 5.5 (1M context) --- .claude/docs/synchronizer-model.md | 2 +- .../base/synchronizer/SynchronizersOrder.java | 7 + .../db/changelog/dirigible-system.json | 212 +++++++++++++ components/data/data-migrations/pom.xml | 53 ++++ .../components/data/migrations/Migration.java | 143 +++++++++ .../data/migrations/MigrationExecutor.java | 191 ++++++++++++ .../data/migrations/MigrationLedger.java | 247 ++++++++++++++++ .../data/migrations/MigrationRepository.java | 38 +++ .../data/migrations/MigrationScript.java | 179 +++++++++++ .../data/migrations/MigrationService.java | 40 +++ .../data/migrations/MigrationsEndpoint.java | 69 +++++ .../migrations/MigrationsHealthIndicator.java | 55 ++++ .../migrations/MigrationsSynchronizer.java | 279 ++++++++++++++++++ .../migrations/MigrationExecutorTest.java | 185 ++++++++++++ .../data/migrations/MigrationScriptTest.java | 115 ++++++++ .../MigrationsSynchronizerTest.java | 229 ++++++++++++++ components/group/group-database/pom.xml | 4 + components/pom.xml | 6 + .../tests/api/DataMigrationIT.java | 255 ++++++++++++++++ 19 files changed, 2308 insertions(+), 1 deletion(-) create mode 100644 components/data/data-migrations/pom.xml create mode 100644 components/data/data-migrations/src/main/java/org/eclipse/dirigible/components/data/migrations/Migration.java create mode 100644 components/data/data-migrations/src/main/java/org/eclipse/dirigible/components/data/migrations/MigrationExecutor.java create mode 100644 components/data/data-migrations/src/main/java/org/eclipse/dirigible/components/data/migrations/MigrationLedger.java create mode 100644 components/data/data-migrations/src/main/java/org/eclipse/dirigible/components/data/migrations/MigrationRepository.java create mode 100644 components/data/data-migrations/src/main/java/org/eclipse/dirigible/components/data/migrations/MigrationScript.java create mode 100644 components/data/data-migrations/src/main/java/org/eclipse/dirigible/components/data/migrations/MigrationService.java create mode 100644 components/data/data-migrations/src/main/java/org/eclipse/dirigible/components/data/migrations/MigrationsEndpoint.java create mode 100644 components/data/data-migrations/src/main/java/org/eclipse/dirigible/components/data/migrations/MigrationsHealthIndicator.java create mode 100644 components/data/data-migrations/src/main/java/org/eclipse/dirigible/components/data/migrations/MigrationsSynchronizer.java create mode 100644 components/data/data-migrations/src/test/java/org/eclipse/dirigible/components/data/migrations/MigrationExecutorTest.java create mode 100644 components/data/data-migrations/src/test/java/org/eclipse/dirigible/components/data/migrations/MigrationScriptTest.java create mode 100644 components/data/data-migrations/src/test/java/org/eclipse/dirigible/components/data/migrations/MigrationsSynchronizerTest.java create mode 100644 tests/tests-integrations/src/main/java/org/eclipse/dirigible/integration/tests/api/DataMigrationIT.java 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: + *
    + *
  • recorded with the same checksum - already applied, nothing runs;
  • + *
  • recorded with another checksum - the file was edited after it applied: a failure, unless the + * migration declares itself idempotent, in which case it runs again and the row is updated;
  • + *
  • not recorded, but an earlier version of the same project is not recorded either - it waits, + * so a project's migrations apply in version order;
  • + *
  • otherwise - it runs.
  • + *
+ */ +@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..2e32fba4be7 --- /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 using sql [{}]", sql); + } 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; + } +} From f966d2f472dfbff3f26aac47c5788314ef11a7cf Mon Sep 17 00:00:00 2001 From: Nikol Georgieva Date: Mon, 5 Oct 2026 14:26:03 +0300 Subject: [PATCH 2/2] migrations: the ledger creation logs the table, not the SQL (#7636) The CREATE TABLE is built from constants, but CodeQL reads it as derived from user input (every datasource leaves DataSourcesManager through one static map) and flagged the log line as log injection. Logging the table name says what happened and carries nothing tainted. Verified: data-migrations unit tests (22) green; formatter:validate on the module with the cache wiped. Co-Authored-By: Claude Opus 5.5 (1M context) --- .../dirigible/components/data/migrations/MigrationLedger.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) 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 index 2e32fba4be7..c17596d629b 100644 --- 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 @@ -103,7 +103,7 @@ void prepare(DataSource dataSource) throws SQLException { .build(); try (PreparedStatement statement = connection.prepareStatement(sql)) { statement.executeUpdate(); - LOGGER.info("Created the migrations ledger using sql [{}]", sql); + 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)