From e938cd0c6293ce4f07afedfb6c86ab9d6f9ede50 Mon Sep 17 00:00:00 2001 From: davidfrigolet Date: Tue, 16 Jun 2026 13:28:20 +0100 Subject: [PATCH 1/4] fix(mongock-importer): skip IGNORED audit entries safely `MongockImporterMongoDB.toAuditEntry` returns `null` for entries whose state is `IGNORED` (via `MongockAuditEntry.shouldBeIgnored()`), but the returned list is iterated downstream without a null check. Booting an application against a Mongock audit history containing an `IGNORED` entry therefore crashed with: NullPointerException: Cannot invoke "AuditEntry.getSystemChange()" because "auditEntryFromOrigin" is null inside the `migration-mongock-to-flamingock-community` system change. Changes: - **MongockImporterMongoDB**: filter `null`s out of the stream produced by `getAuditHistory()` and emit an INFO log when an entry is dropped because of `IGNORED`. The skipped entry's `changeId` is reported so operators know which legacy change was bypassed. - **MongockImportChange**: defensive null-filter at the caller so the same protection applies to any current or future `AuditHistoryReader` implementation (Couchbase, DynamoDB). - **MongoDBImporterTest**: regression test seeding the basic scenario plus one `IGNORED` entry. Asserts the `IGNORED` entry is absent from the Flamingock audit store and the rest of the audit history imports normally. No public API change. Behavior is purely additive (null-safety + log) relative to the previous code path. --- .../mongodb/MongockImporterMongoDB.java | 8 +++ .../mongock/mongodb/MongoDBImporterTest.java | 59 +++++++++++++++++++ .../support/mongock/MongockImportChange.java | 7 ++- 3 files changed, 73 insertions(+), 1 deletion(-) diff --git a/legacy/mongock-importer-mongodb/src/main/java/io/flamingock/importer/mongock/mongodb/MongockImporterMongoDB.java b/legacy/mongock-importer-mongodb/src/main/java/io/flamingock/importer/mongock/mongodb/MongockImporterMongoDB.java index 0ba220ebe..7d4e830f9 100644 --- a/legacy/mongock-importer-mongodb/src/main/java/io/flamingock/importer/mongock/mongodb/MongockImporterMongoDB.java +++ b/legacy/mongock-importer-mongodb/src/main/java/io/flamingock/importer/mongock/mongodb/MongockImporterMongoDB.java @@ -21,6 +21,8 @@ import io.flamingock.internal.common.core.audit.AuditEntry; import io.flamingock.internal.common.core.audit.AuditHistoryReader; import org.bson.Document; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; import java.time.Instant; import java.time.LocalDateTime; @@ -28,10 +30,13 @@ import java.util.ArrayList; import java.util.Date; import java.util.List; +import java.util.Objects; import java.util.stream.Collectors; public class MongockImporterMongoDB implements AuditHistoryReader { + private static final Logger logger = LoggerFactory.getLogger("MongockImporter"); + private final MongoCollection sourceCollection; public MongockImporterMongoDB(MongoDatabase mongoDatabase, String collectionName) { @@ -44,6 +49,7 @@ public List getAuditHistory() { .into(new ArrayList<>()) .stream() .map(MongockImporterMongoDB::toAuditEntry) + .filter(Objects::nonNull) .collect(Collectors.toList()); } @@ -55,6 +61,8 @@ private static AuditEntry toAuditEntry(Document document) { .toLocalDateTime(); if (changeEntry.shouldBeIgnored()) { + logger.info("Skipping Mongock audit entry with changeId[{}]: state=IGNORED (Mongock never executed this change; nothing to import).", + changeEntry.getChangeId()); return null; } return new AuditEntry( diff --git a/legacy/mongock-importer-mongodb/src/test/java/io/flamingock/importer/mongock/mongodb/MongoDBImporterTest.java b/legacy/mongock-importer-mongodb/src/test/java/io/flamingock/importer/mongock/mongodb/MongoDBImporterTest.java index fa4fad124..8ce564540 100644 --- a/legacy/mongock-importer-mongodb/src/test/java/io/flamingock/importer/mongock/mongodb/MongoDBImporterTest.java +++ b/legacy/mongock-importer-mongodb/src/test/java/io/flamingock/importer/mongock/mongodb/MongoDBImporterTest.java @@ -498,6 +498,65 @@ void GIVEN_unknownAuditEntriesAndRelaxedMode_WHEN_migratingToFlamingockCommunity ); } + @Test + @DisplayName("GIVEN Mongock audit history contains an IGNORED entry " + + "WHEN migrating to Flamingock Community " + + "THEN should skip the IGNORED entry without crashing " + + "AND import the rest of the history normally") + void GIVEN_ignoredAuditEntry_WHEN_migratingToFlamingockCommunity_THEN_shouldSkipIgnoredAndImportRest() throws java.text.ParseException { + // Regression test for IGNORED-state null leak in MongockImporterMongoDB.toAuditEntry(). + // Before the fix, the importer returned null for IGNORED entries and the caller + // (MongockImportChange.importHistory) crashed with: + // NullPointerException: Cannot invoke "AuditEntry.getSystemChange()" + // because "auditEntryFromOrigin" is null + // The fix filters nulls at both layers and logs the skipped entry. + + mongockTestHelper.setupBasicScenario(); + mongockTestHelper.write(new MongockChangeEntry( + "ignored-execution-1", + "ignored-change", + "mongock", + MongockTestHelper.DEFAULT_DATE_FORMAT.parse("2025-06-19T05:43:57.200Z"), + MongockChangeState.IGNORED, + io.flamingock.common.test.mongock.MongockChangeType.EXECUTION, + "io.example.IgnoredChangeUnit", + "apply", + null, + 0L, + MongockTestHelper.DEFAULT_HOSTNAME, + null, + false, + null + )); + + MongoDBSyncTargetSystem mongodbTargetSystem = new MongoDBSyncTargetSystem("mongodb-target-system", mongoClient, DATABASE_NAME); + + Runner flamingock = testKit.createBuilder() + .addTargetSystem(mongodbTargetSystem) + .build(); + + flamingock.run(); + + // IGNORED entry must NOT appear in the Flamingock audit store. + assertNull(getAuditEntryByChangeId("ignored-change"), + "IGNORED Mongock entry must not be imported into the Flamingock audit store"); + + // Remaining basic-scenario entries imported as normal, plus native changes executed. + auditHelper.verifyAuditSequenceStrict( + APPLIED("system-change-00001_before"), + APPLIED("system-change-00001"), + APPLIED("mongock-change-1_before"), + APPLIED("mongock-change-1"), + APPLIED("mongock-change-2"), + STARTED("migration-mongock-to-flamingock-community"), + APPLIED("migration-mongock-to-flamingock-community"), + STARTED("create-users-collection-with-index"), + APPLIED("create-users-collection-with-index"), + STARTED("seed-users"), + APPLIED("seed-users") + ); + } + @Test @DisplayName("GIVEN relaxed import flag with invalid value " + "WHEN migrating to Flamingock Community " + diff --git a/legacy/mongock-support/src/main/java/io/flamingock/support/mongock/MongockImportChange.java b/legacy/mongock-support/src/main/java/io/flamingock/support/mongock/MongockImportChange.java index 37a9700ff..563aac132 100644 --- a/legacy/mongock-support/src/main/java/io/flamingock/support/mongock/MongockImportChange.java +++ b/legacy/mongock-support/src/main/java/io/flamingock/support/mongock/MongockImportChange.java @@ -31,7 +31,9 @@ import javax.inject.Named; import java.util.List; +import java.util.Objects; import java.util.Optional; +import java.util.stream.Collectors; import static io.flamingock.internal.common.core.audit.AuditReaderType.MONGOCK; import static io.flamingock.internal.common.core.metadata.Constants.MONGOCK_IMPORT_EMPTY_ORIGIN_ALLOWED_PROPERTY_KEY; @@ -61,7 +63,10 @@ public void importHistory(@Named("change.targetSystem.id") String targetSystemId logger.info("Starting audit log migration from Mongock to Flamingock community audit store"); AuditHistoryReader legacyHistoryReader = getAuditHistoryReader(targetSystemId, targetSystemManager); PipelineHelper pipelineHelper = new PipelineHelper(pipelineDescriptor); - List legacyHistory = legacyHistoryReader.getAuditHistory(); + List legacyHistory = legacyHistoryReader.getAuditHistory() + .stream() + .filter(Objects::nonNull) + .collect(Collectors.toList()); boolean ignoreUnknownEntries = resolveIgnoreUnknownEntries(ignoreUnknownEntriesPropertyValue); validate(legacyHistory, targetSystemId, emptyOriginAllowedPropertyValue); legacyHistory.forEach(auditEntryFromOrigin -> { From eb632ff89c79fe548acfdcf78b8cc661bc61f050 Mon Sep 17 00:00:00 2001 From: davidfrigolet Date: Tue, 16 Jun 2026 13:28:20 +0100 Subject: [PATCH 2/4] fix(mongock-importer): skip IGNORED audit entries safely `MongockImporterMongoDB.toAuditEntry` returns `null` for entries whose state is `IGNORED` (via `MongockAuditEntry.shouldBeIgnored()`), but the returned list is iterated downstream without a null check. Booting an application against a Mongock audit history containing an `IGNORED` entry therefore crashed with: NullPointerException: Cannot invoke "AuditEntry.getSystemChange()" because "auditEntryFromOrigin" is null inside the `migration-mongock-to-flamingock-community` system change. Changes: - **MongockImporterMongoDB**: filter `null`s out of the stream produced by `getAuditHistory()` and emit an INFO log when an entry is dropped because of `IGNORED`. The skipped entry's `changeId` is reported so operators know which legacy change was bypassed. - **MongockImportChange**: defensive null-filter at the caller so the same protection applies to any current or future `AuditHistoryReader` implementation (Couchbase, DynamoDB). - **MongoDBImporterTest**: regression test seeding the basic scenario plus one `IGNORED` entry. Asserts the `IGNORED` entry is absent from the Flamingock audit store and the rest of the audit history imports normally. No public API change. Behavior is purely additive (null-safety + log) relative to the previous code path. --- .../mongodb/MongockImporterMongoDB.java | 8 +++ .../mongock/mongodb/MongoDBImporterTest.java | 59 +++++++++++++++++++ .../support/mongock/MongockImportChange.java | 7 ++- 3 files changed, 73 insertions(+), 1 deletion(-) diff --git a/legacy/mongock-importer-mongodb/src/main/java/io/flamingock/importer/mongock/mongodb/MongockImporterMongoDB.java b/legacy/mongock-importer-mongodb/src/main/java/io/flamingock/importer/mongock/mongodb/MongockImporterMongoDB.java index 0ba220ebe..7d4e830f9 100644 --- a/legacy/mongock-importer-mongodb/src/main/java/io/flamingock/importer/mongock/mongodb/MongockImporterMongoDB.java +++ b/legacy/mongock-importer-mongodb/src/main/java/io/flamingock/importer/mongock/mongodb/MongockImporterMongoDB.java @@ -21,6 +21,8 @@ import io.flamingock.internal.common.core.audit.AuditEntry; import io.flamingock.internal.common.core.audit.AuditHistoryReader; import org.bson.Document; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; import java.time.Instant; import java.time.LocalDateTime; @@ -28,10 +30,13 @@ import java.util.ArrayList; import java.util.Date; import java.util.List; +import java.util.Objects; import java.util.stream.Collectors; public class MongockImporterMongoDB implements AuditHistoryReader { + private static final Logger logger = LoggerFactory.getLogger("MongockImporter"); + private final MongoCollection sourceCollection; public MongockImporterMongoDB(MongoDatabase mongoDatabase, String collectionName) { @@ -44,6 +49,7 @@ public List getAuditHistory() { .into(new ArrayList<>()) .stream() .map(MongockImporterMongoDB::toAuditEntry) + .filter(Objects::nonNull) .collect(Collectors.toList()); } @@ -55,6 +61,8 @@ private static AuditEntry toAuditEntry(Document document) { .toLocalDateTime(); if (changeEntry.shouldBeIgnored()) { + logger.info("Skipping Mongock audit entry with changeId[{}]: state=IGNORED (Mongock never executed this change; nothing to import).", + changeEntry.getChangeId()); return null; } return new AuditEntry( diff --git a/legacy/mongock-importer-mongodb/src/test/java/io/flamingock/importer/mongock/mongodb/MongoDBImporterTest.java b/legacy/mongock-importer-mongodb/src/test/java/io/flamingock/importer/mongock/mongodb/MongoDBImporterTest.java index fa4fad124..8ce564540 100644 --- a/legacy/mongock-importer-mongodb/src/test/java/io/flamingock/importer/mongock/mongodb/MongoDBImporterTest.java +++ b/legacy/mongock-importer-mongodb/src/test/java/io/flamingock/importer/mongock/mongodb/MongoDBImporterTest.java @@ -498,6 +498,65 @@ void GIVEN_unknownAuditEntriesAndRelaxedMode_WHEN_migratingToFlamingockCommunity ); } + @Test + @DisplayName("GIVEN Mongock audit history contains an IGNORED entry " + + "WHEN migrating to Flamingock Community " + + "THEN should skip the IGNORED entry without crashing " + + "AND import the rest of the history normally") + void GIVEN_ignoredAuditEntry_WHEN_migratingToFlamingockCommunity_THEN_shouldSkipIgnoredAndImportRest() throws java.text.ParseException { + // Regression test for IGNORED-state null leak in MongockImporterMongoDB.toAuditEntry(). + // Before the fix, the importer returned null for IGNORED entries and the caller + // (MongockImportChange.importHistory) crashed with: + // NullPointerException: Cannot invoke "AuditEntry.getSystemChange()" + // because "auditEntryFromOrigin" is null + // The fix filters nulls at both layers and logs the skipped entry. + + mongockTestHelper.setupBasicScenario(); + mongockTestHelper.write(new MongockChangeEntry( + "ignored-execution-1", + "ignored-change", + "mongock", + MongockTestHelper.DEFAULT_DATE_FORMAT.parse("2025-06-19T05:43:57.200Z"), + MongockChangeState.IGNORED, + io.flamingock.common.test.mongock.MongockChangeType.EXECUTION, + "io.example.IgnoredChangeUnit", + "apply", + null, + 0L, + MongockTestHelper.DEFAULT_HOSTNAME, + null, + false, + null + )); + + MongoDBSyncTargetSystem mongodbTargetSystem = new MongoDBSyncTargetSystem("mongodb-target-system", mongoClient, DATABASE_NAME); + + Runner flamingock = testKit.createBuilder() + .addTargetSystem(mongodbTargetSystem) + .build(); + + flamingock.run(); + + // IGNORED entry must NOT appear in the Flamingock audit store. + assertNull(getAuditEntryByChangeId("ignored-change"), + "IGNORED Mongock entry must not be imported into the Flamingock audit store"); + + // Remaining basic-scenario entries imported as normal, plus native changes executed. + auditHelper.verifyAuditSequenceStrict( + APPLIED("system-change-00001_before"), + APPLIED("system-change-00001"), + APPLIED("mongock-change-1_before"), + APPLIED("mongock-change-1"), + APPLIED("mongock-change-2"), + STARTED("migration-mongock-to-flamingock-community"), + APPLIED("migration-mongock-to-flamingock-community"), + STARTED("create-users-collection-with-index"), + APPLIED("create-users-collection-with-index"), + STARTED("seed-users"), + APPLIED("seed-users") + ); + } + @Test @DisplayName("GIVEN relaxed import flag with invalid value " + "WHEN migrating to Flamingock Community " + diff --git a/legacy/mongock-support/src/main/java/io/flamingock/support/mongock/MongockImportChange.java b/legacy/mongock-support/src/main/java/io/flamingock/support/mongock/MongockImportChange.java index 37a9700ff..563aac132 100644 --- a/legacy/mongock-support/src/main/java/io/flamingock/support/mongock/MongockImportChange.java +++ b/legacy/mongock-support/src/main/java/io/flamingock/support/mongock/MongockImportChange.java @@ -31,7 +31,9 @@ import javax.inject.Named; import java.util.List; +import java.util.Objects; import java.util.Optional; +import java.util.stream.Collectors; import static io.flamingock.internal.common.core.audit.AuditReaderType.MONGOCK; import static io.flamingock.internal.common.core.metadata.Constants.MONGOCK_IMPORT_EMPTY_ORIGIN_ALLOWED_PROPERTY_KEY; @@ -61,7 +63,10 @@ public void importHistory(@Named("change.targetSystem.id") String targetSystemId logger.info("Starting audit log migration from Mongock to Flamingock community audit store"); AuditHistoryReader legacyHistoryReader = getAuditHistoryReader(targetSystemId, targetSystemManager); PipelineHelper pipelineHelper = new PipelineHelper(pipelineDescriptor); - List legacyHistory = legacyHistoryReader.getAuditHistory(); + List legacyHistory = legacyHistoryReader.getAuditHistory() + .stream() + .filter(Objects::nonNull) + .collect(Collectors.toList()); boolean ignoreUnknownEntries = resolveIgnoreUnknownEntries(ignoreUnknownEntriesPropertyValue); validate(legacyHistory, targetSystemId, emptyOriginAllowedPropertyValue); legacyHistory.forEach(auditEntryFromOrigin -> { From 0ae88a3194f8eac78eb894a69a34779f3f601b54 Mon Sep 17 00:00:00 2001 From: davidfrigolet Date: Wed, 29 Jul 2026 16:03:01 +0100 Subject: [PATCH 3/4] fix(mongock-importer): skip IGNORED audit entries safely --- .../couchbase/CouchbaseChangeEntry.java | 12 +++++ .../couchbase/MongockImporterCouchbase.java | 17 ++++++- .../couchbase/CouchbaseImporterTest.java | 40 +++++++++++++++- .../mongock/dynamodb/MongockAuditEntry.java | 8 ++++ .../dynamodb/MongockImporterDynamoDB.java | 17 ++++++- .../dynamodb/DynamoDBImporterTest.java | 47 +++++++++++++++++++ .../MongockImporterMongoDBReactive.java | 15 ++++-- .../MongockImporterMongoDBReactiveTest.java | 7 ++- .../mongodb/MongockImporterMongoDB.java | 2 +- .../support/mongock/MongockImportChange.java | 7 +-- 10 files changed, 155 insertions(+), 17 deletions(-) diff --git a/legacy/mongock-importer-couchbase/src/main/java/io/flamingock/importer/mongock/couchbase/CouchbaseChangeEntry.java b/legacy/mongock-importer-couchbase/src/main/java/io/flamingock/importer/mongock/couchbase/CouchbaseChangeEntry.java index 4448ed862..bf93ac589 100644 --- a/legacy/mongock-importer-couchbase/src/main/java/io/flamingock/importer/mongock/couchbase/CouchbaseChangeEntry.java +++ b/legacy/mongock-importer-couchbase/src/main/java/io/flamingock/importer/mongock/couchbase/CouchbaseChangeEntry.java @@ -46,6 +46,10 @@ public void setExecutionId(String executionId) { this.executionId = executionId; } + public String getChangeId() { + return changeId; + } + public static CouchbaseChangeEntry fromJson(JsonObject doc) { CouchbaseChangeEntry entry = new CouchbaseChangeEntry(); entry.executionId = doc.getString("executionId"); @@ -71,7 +75,15 @@ private static Long parseLong(Object value) { throw new IllegalArgumentException("Cannot convert value to Long: " + value); } + public boolean shouldBeIgnored() { + return MongockChangeState.valueOf(state) == MongockChangeState.IGNORED; + } + public AuditEntry toAuditEntry() { + if (shouldBeIgnored()) { + return null; + } + LocalDateTime ts = LocalDateTime.ofInstant(Instant.ofEpochMilli(timestamp), ZoneId.systemDefault()); MongockChangeState stateEnum = MongockChangeState.valueOf(state); diff --git a/legacy/mongock-importer-couchbase/src/main/java/io/flamingock/importer/mongock/couchbase/MongockImporterCouchbase.java b/legacy/mongock-importer-couchbase/src/main/java/io/flamingock/importer/mongock/couchbase/MongockImporterCouchbase.java index 9780fa9aa..50a23ae70 100644 --- a/legacy/mongock-importer-couchbase/src/main/java/io/flamingock/importer/mongock/couchbase/MongockImporterCouchbase.java +++ b/legacy/mongock-importer-couchbase/src/main/java/io/flamingock/importer/mongock/couchbase/MongockImporterCouchbase.java @@ -22,12 +22,17 @@ import com.couchbase.client.java.query.QueryScanConsistency; import io.flamingock.internal.common.core.audit.AuditEntry; import io.flamingock.internal.common.core.audit.AuditHistoryReader; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; import java.util.List; +import java.util.Objects; import java.util.stream.Collectors; public class MongockImporterCouchbase implements AuditHistoryReader { + private static final Logger logger = LoggerFactory.getLogger("MongockImporter"); + private static final String MONGOCK_CHANGE_ENTRY_DOCTYPE = "mongockChangeEntry"; private final Cluster cluster; @@ -56,7 +61,17 @@ public List getAuditHistory() { return result.rowsAsObject().stream() .map(CouchbaseChangeEntry::fromJson) - .map(CouchbaseChangeEntry::toAuditEntry) + .map(MongockImporterCouchbase::toAuditEntry) + .filter(Objects::nonNull) .collect(Collectors.toList()); } + + private static AuditEntry toAuditEntry(CouchbaseChangeEntry entry) { + if (entry.shouldBeIgnored()) { + logger.info("Skipping Mongock audit entry with changeId[{}]: state=IGNORED (change was already executed by Mongock; not imported into Flamingock audit history).", + entry.getChangeId()); + return null; + } + return entry.toAuditEntry(); + } } \ No newline at end of file diff --git a/legacy/mongock-importer-couchbase/src/test/java/io/flamingock/importer/mongock/couchbase/CouchbaseImporterTest.java b/legacy/mongock-importer-couchbase/src/test/java/io/flamingock/importer/mongock/couchbase/CouchbaseImporterTest.java index 38d721510..c44f2990d 100644 --- a/legacy/mongock-importer-couchbase/src/test/java/io/flamingock/importer/mongock/couchbase/CouchbaseImporterTest.java +++ b/legacy/mongock-importer-couchbase/src/test/java/io/flamingock/importer/mongock/couchbase/CouchbaseImporterTest.java @@ -204,6 +204,40 @@ void GIVEN_someChangeUnitsAlreadyExecuted_WHEN_migratingToFlamingockCommunity_TH } + @Test + @DisplayName("GIVEN Mongock audit history contains an IGNORED entry " + + "WHEN migrating to Flamingock Community " + + "THEN should skip the IGNORED entry without crashing " + + "AND import the rest of the history normally") + void GIVEN_ignoredAuditEntry_WHEN_migratingToFlamingockCommunity_THEN_shouldSkipIgnoredAndImportRest() { + Collection originCollection = cluster.bucket(MONGOCK_BUCKET_NAME).scope(MONGOCK_SCOPE_NAME).collection(MONGOCK_COLLECTION_NAME); + + originCollection.upsert("mongock-change-1", createAuditObject("mongock-change-1")); + originCollection.upsert("mongock-change-2", createAuditObject("mongock-change-2")); + originCollection.upsert("ignored-change", createAuditObject("ignored-change", true, "io.example.IgnoredChangeUnit", "apply", "IGNORED")); + + Runner flamingock = testKit.createBuilder() + .setAuditStore(auditStore) + .addTargetSystem(targetSystem) + .build(); + + flamingock.run(); + + auditHelper.verifyAuditSequenceStrict( + // Legacy imports from Mongock (APPLIED only - no STARTED for imported changes) + APPLIED("mongock-change-1"), + APPLIED("mongock-change-2"), + + // System stage - actual system importer change + STARTED("migration-mongock-to-flamingock-community"), + APPLIED("migration-mongock-to-flamingock-community"), + + // Application stage - new changes + STARTED("flamingock-change"), + APPLIED("flamingock-change") + ); + } + @Test @DisplayName("GIVEN mongock audit history empty " + "AND no empty origen allowed value provided " + @@ -587,12 +621,16 @@ private static JsonObject createAuditObject(String value) { } private static JsonObject createAuditObject(String value, boolean systemChange, String changeLogClass, String changeSetMethod) { + return createAuditObject(value, systemChange, changeLogClass, changeSetMethod, "EXECUTED"); + } + + private static JsonObject createAuditObject(String value, boolean systemChange, String changeLogClass, String changeSetMethod, String state) { JsonObject doc = JsonObject.create() .put("executionId", "exec-1") .put("changeId", value) .put("author", "author1") .put("timestamp", Instant.now().toEpochMilli()) - .put("state", "EXECUTED") + .put("state", state) .put("type", "EXECUTION") .put("changeLogClass", changeLogClass) .put("changeSetMethod", changeSetMethod) diff --git a/legacy/mongock-importer-dynamodb/src/main/java/io/flamingock/importer/mongock/dynamodb/MongockAuditEntry.java b/legacy/mongock-importer-dynamodb/src/main/java/io/flamingock/importer/mongock/dynamodb/MongockAuditEntry.java index d1e05b55e..6c56abbea 100644 --- a/legacy/mongock-importer-dynamodb/src/main/java/io/flamingock/importer/mongock/dynamodb/MongockAuditEntry.java +++ b/legacy/mongock-importer-dynamodb/src/main/java/io/flamingock/importer/mongock/dynamodb/MongockAuditEntry.java @@ -180,7 +180,15 @@ public void setSystemChange(Boolean systemChange) { this.systemChange = systemChange; } + public boolean shouldBeIgnored() { + return MongockChangeState.valueOf(state) == MongockChangeState.IGNORED; + } + public AuditEntry toAuditEntry() { + if (shouldBeIgnored()) { + return null; + } + long epochMillis; try { epochMillis = Long.parseLong(timestamp); diff --git a/legacy/mongock-importer-dynamodb/src/main/java/io/flamingock/importer/mongock/dynamodb/MongockImporterDynamoDB.java b/legacy/mongock-importer-dynamodb/src/main/java/io/flamingock/importer/mongock/dynamodb/MongockImporterDynamoDB.java index 603ec3276..cdc34db6b 100644 --- a/legacy/mongock-importer-dynamodb/src/main/java/io/flamingock/importer/mongock/dynamodb/MongockImporterDynamoDB.java +++ b/legacy/mongock-importer-dynamodb/src/main/java/io/flamingock/importer/mongock/dynamodb/MongockImporterDynamoDB.java @@ -17,17 +17,22 @@ import io.flamingock.internal.common.core.audit.AuditEntry; import io.flamingock.internal.common.core.audit.AuditHistoryReader; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; import software.amazon.awssdk.enhanced.dynamodb.DynamoDbEnhancedClient; import software.amazon.awssdk.enhanced.dynamodb.DynamoDbTable; import software.amazon.awssdk.enhanced.dynamodb.TableSchema; import software.amazon.awssdk.services.dynamodb.DynamoDbClient; import java.util.List; +import java.util.Objects; import java.util.stream.Collectors; import java.util.stream.StreamSupport; public class MongockImporterDynamoDB implements AuditHistoryReader { + private static final Logger logger = LoggerFactory.getLogger("MongockImporter"); + private final DynamoDbTable sourceTable; public MongockImporterDynamoDB(DynamoDbClient client, String tableName) { @@ -44,7 +49,17 @@ public List getAuditHistory() { .collect(Collectors.toList()); return entries.stream() - .map(MongockAuditEntry::toAuditEntry) + .map(MongockImporterDynamoDB::toAuditEntry) + .filter(Objects::nonNull) .collect(Collectors.toList()); } + + private static AuditEntry toAuditEntry(MongockAuditEntry entry) { + if (entry.shouldBeIgnored()) { + logger.info("Skipping Mongock audit entry with changeId[{}]: state=IGNORED (change was already executed by Mongock; not imported into Flamingock audit history).", + entry.getChangeId()); + return null; + } + return entry.toAuditEntry(); + } } diff --git a/legacy/mongock-importer-dynamodb/src/test/java/io/flamingock/importer/mongock/dynamodb/DynamoDBImporterTest.java b/legacy/mongock-importer-dynamodb/src/test/java/io/flamingock/importer/mongock/dynamodb/DynamoDBImporterTest.java index d4818fbcf..4032417bf 100644 --- a/legacy/mongock-importer-dynamodb/src/test/java/io/flamingock/importer/mongock/dynamodb/DynamoDBImporterTest.java +++ b/legacy/mongock-importer-dynamodb/src/test/java/io/flamingock/importer/mongock/dynamodb/DynamoDBImporterTest.java @@ -17,6 +17,8 @@ import io.flamingock.api.annotations.EnableFlamingock; import io.flamingock.api.annotations.Stage; +import io.flamingock.common.test.mongock.MongockChangeEntry; +import io.flamingock.common.test.mongock.MongockChangeState; import io.flamingock.store.dynamodb.DynamoDBAuditStore; import io.flamingock.core.kit.TestKit; import io.flamingock.core.kit.audit.AuditTestHelper; @@ -420,6 +422,51 @@ void GIVEN_unknownAuditEntriesAndRelaxedMode_WHEN_migratingToFlamingockCommunity ); } + @Test + @DisplayName("GIVEN Mongock audit history contains an IGNORED entry " + + "WHEN migrating to Flamingock Community " + + "THEN should skip the IGNORED entry without crashing " + + "AND import the rest of the history normally") + void GIVEN_ignoredAuditEntry_WHEN_migratingToFlamingockCommunity_THEN_shouldSkipIgnoredAndImportRest() throws java.text.ParseException { + mongockTestHelper.setupBasicScenario(); + mongockTestHelper.write(new MongockChangeEntry( + "ignored-execution-1", + "ignored-change", + "mongock", + io.flamingock.common.test.mongock.MongockTestHelper.DEFAULT_DATE_FORMAT.parse("2025-06-19T05:43:57.200Z"), + MongockChangeState.IGNORED, + io.flamingock.common.test.mongock.MongockChangeType.EXECUTION, + "io.example.IgnoredChangeUnit", + "apply", + null, + 0L, + io.flamingock.common.test.mongock.MongockTestHelper.DEFAULT_HOSTNAME, + null, + false, + null + )); + + DynamoDBTargetSystem dynamodbTargetSystem = new DynamoDBTargetSystem("dynamodb-target-system", client); + + Runner flamingock = testKit.createBuilder() + .addTargetSystem(dynamodbTargetSystem) + .build(); + + flamingock.run(); + + auditHelper.verifyAuditSequenceStrict( + APPLIED("system-change-00001_before"), + APPLIED("system-change-00001"), + APPLIED("mongock-change-1_before"), + APPLIED("mongock-change-1"), + APPLIED("mongock-change-2"), + STARTED("migration-mongock-to-flamingock-community"), + APPLIED("migration-mongock-to-flamingock-community"), + STARTED("create-users-table"), + APPLIED("create-users-table") + ); + } + @Test @DisplayName("GIVEN relaxed import flag with invalid value " + "WHEN migrating to Flamingock Community " + diff --git a/legacy/mongock-importer-mongodb-reactive/src/main/java/io/flamingock/importer/mongock/mongodb/reactive/MongockImporterMongoDBReactive.java b/legacy/mongock-importer-mongodb-reactive/src/main/java/io/flamingock/importer/mongock/mongodb/reactive/MongockImporterMongoDBReactive.java index ed7a666d3..941235e90 100644 --- a/legacy/mongock-importer-mongodb-reactive/src/main/java/io/flamingock/importer/mongock/mongodb/reactive/MongockImporterMongoDBReactive.java +++ b/legacy/mongock-importer-mongodb-reactive/src/main/java/io/flamingock/importer/mongock/mongodb/reactive/MongockImporterMongoDBReactive.java @@ -22,16 +22,21 @@ import io.flamingock.internal.common.core.audit.AuditHistoryReader; import io.flamingock.reactive.util.PublisherSync; import org.bson.Document; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; import java.time.Instant; import java.time.LocalDateTime; import java.time.ZoneId; import java.util.Date; import java.util.List; +import java.util.Objects; import java.util.stream.Collectors; public class MongockImporterMongoDBReactive implements AuditHistoryReader { + private static final Logger logger = LoggerFactory.getLogger("MongockImporter"); + private final MongoCollection sourceCollection; public MongockImporterMongoDBReactive(MongoDatabase mongoDatabase, String collectionName) { @@ -43,19 +48,23 @@ public List getAuditHistory() { return PublisherSync.collect(sourceCollection.find()) .stream() .map(MongockImporterMongoDBReactive::toAuditEntry) + .filter(Objects::nonNull) .collect(Collectors.toList()); } private static AuditEntry toAuditEntry(Document document) { MongockAuditEntry changeEntry = toChangeEntry(document); - LocalDateTime timestamp = Instant.ofEpochMilli(changeEntry.getTimestamp().getTime()) - .atZone(ZoneId.systemDefault()) - .toLocalDateTime(); if (changeEntry.shouldBeIgnored()) { + logger.info("Skipping Mongock audit entry with changeId[{}]: state=IGNORED (change was already executed by Mongock; not imported into Flamingock audit history).", + changeEntry.getChangeId()); return null; } + + LocalDateTime timestamp = Instant.ofEpochMilli(changeEntry.getTimestamp().getTime()) + .atZone(ZoneId.systemDefault()) + .toLocalDateTime(); return new AuditEntry( changeEntry.getExecutionId(), null, diff --git a/legacy/mongock-importer-mongodb-reactive/src/test/java/io/flamingock/importer/mongock/mongodb/reactive/MongockImporterMongoDBReactiveTest.java b/legacy/mongock-importer-mongodb-reactive/src/test/java/io/flamingock/importer/mongock/mongodb/reactive/MongockImporterMongoDBReactiveTest.java index 5823e656d..e2ad7cf16 100644 --- a/legacy/mongock-importer-mongodb-reactive/src/test/java/io/flamingock/importer/mongock/mongodb/reactive/MongockImporterMongoDBReactiveTest.java +++ b/legacy/mongock-importer-mongodb-reactive/src/test/java/io/flamingock/importer/mongock/mongodb/reactive/MongockImporterMongoDBReactiveTest.java @@ -92,17 +92,16 @@ void shouldMapExecutedEntry() { } @Test - @DisplayName("Should map an IGNORED legacy entry to null") - void shouldMapIgnoredEntryToNull() { + @DisplayName("Should skip an IGNORED legacy entry") + void shouldSkipIgnoredEntry() { seed(document("users-initialization", "EXECUTED", "EXECUTION", "pretend-mongock-run")); seed(document("ghost-extra", "IGNORED", "EXECUTION", null)); MongockImporterMongoDBReactive importer = new MongockImporterMongoDBReactive(database, LEGACY_COLLECTION); List history = importer.getAuditHistory(); - Assertions.assertEquals(2, history.size()); + Assertions.assertEquals(1, history.size()); Assertions.assertEquals("users-initialization", history.get(0).getChangeId()); - Assertions.assertNull(history.get(1)); } @Test diff --git a/legacy/mongock-importer-mongodb/src/main/java/io/flamingock/importer/mongock/mongodb/MongockImporterMongoDB.java b/legacy/mongock-importer-mongodb/src/main/java/io/flamingock/importer/mongock/mongodb/MongockImporterMongoDB.java index 7d4e830f9..943ed88c9 100644 --- a/legacy/mongock-importer-mongodb/src/main/java/io/flamingock/importer/mongock/mongodb/MongockImporterMongoDB.java +++ b/legacy/mongock-importer-mongodb/src/main/java/io/flamingock/importer/mongock/mongodb/MongockImporterMongoDB.java @@ -61,7 +61,7 @@ private static AuditEntry toAuditEntry(Document document) { .toLocalDateTime(); if (changeEntry.shouldBeIgnored()) { - logger.info("Skipping Mongock audit entry with changeId[{}]: state=IGNORED (Mongock never executed this change; nothing to import).", + logger.info("Skipping Mongock audit entry with changeId[{}]: state=IGNORED (change was already executed by Mongock; not imported into Flamingock audit history).", changeEntry.getChangeId()); return null; } diff --git a/legacy/mongock-support/src/main/java/io/flamingock/support/mongock/MongockImportChange.java b/legacy/mongock-support/src/main/java/io/flamingock/support/mongock/MongockImportChange.java index 563aac132..37a9700ff 100644 --- a/legacy/mongock-support/src/main/java/io/flamingock/support/mongock/MongockImportChange.java +++ b/legacy/mongock-support/src/main/java/io/flamingock/support/mongock/MongockImportChange.java @@ -31,9 +31,7 @@ import javax.inject.Named; import java.util.List; -import java.util.Objects; import java.util.Optional; -import java.util.stream.Collectors; import static io.flamingock.internal.common.core.audit.AuditReaderType.MONGOCK; import static io.flamingock.internal.common.core.metadata.Constants.MONGOCK_IMPORT_EMPTY_ORIGIN_ALLOWED_PROPERTY_KEY; @@ -63,10 +61,7 @@ public void importHistory(@Named("change.targetSystem.id") String targetSystemId logger.info("Starting audit log migration from Mongock to Flamingock community audit store"); AuditHistoryReader legacyHistoryReader = getAuditHistoryReader(targetSystemId, targetSystemManager); PipelineHelper pipelineHelper = new PipelineHelper(pipelineDescriptor); - List legacyHistory = legacyHistoryReader.getAuditHistory() - .stream() - .filter(Objects::nonNull) - .collect(Collectors.toList()); + List legacyHistory = legacyHistoryReader.getAuditHistory(); boolean ignoreUnknownEntries = resolveIgnoreUnknownEntries(ignoreUnknownEntriesPropertyValue); validate(legacyHistory, targetSystemId, emptyOriginAllowedPropertyValue); legacyHistory.forEach(auditEntryFromOrigin -> { From 76f3111be35c6f22ad9e3bc8dea5b1042f6a0872 Mon Sep 17 00:00:00 2001 From: davidfrigolet Date: Tue, 4 Aug 2026 14:06:43 +0100 Subject: [PATCH 4/4] fix(mongock-importer): skip IGNORED audit entries safely --- .../mongock/couchbase/MongockImporterCouchbase.java | 6 ------ .../importer/mongock/dynamodb/MongockImporterDynamoDB.java | 6 ------ .../mongodb/reactive/MongockImporterMongoDBReactive.java | 6 ------ .../importer/mongock/mongodb/MongockImporterMongoDB.java | 6 ------ 4 files changed, 24 deletions(-) diff --git a/legacy/mongock-importer-couchbase/src/main/java/io/flamingock/importer/mongock/couchbase/MongockImporterCouchbase.java b/legacy/mongock-importer-couchbase/src/main/java/io/flamingock/importer/mongock/couchbase/MongockImporterCouchbase.java index 50a23ae70..e1f19520a 100644 --- a/legacy/mongock-importer-couchbase/src/main/java/io/flamingock/importer/mongock/couchbase/MongockImporterCouchbase.java +++ b/legacy/mongock-importer-couchbase/src/main/java/io/flamingock/importer/mongock/couchbase/MongockImporterCouchbase.java @@ -22,8 +22,6 @@ import com.couchbase.client.java.query.QueryScanConsistency; import io.flamingock.internal.common.core.audit.AuditEntry; import io.flamingock.internal.common.core.audit.AuditHistoryReader; -import org.slf4j.Logger; -import org.slf4j.LoggerFactory; import java.util.List; import java.util.Objects; @@ -31,8 +29,6 @@ public class MongockImporterCouchbase implements AuditHistoryReader { - private static final Logger logger = LoggerFactory.getLogger("MongockImporter"); - private static final String MONGOCK_CHANGE_ENTRY_DOCTYPE = "mongockChangeEntry"; private final Cluster cluster; @@ -68,8 +64,6 @@ public List getAuditHistory() { private static AuditEntry toAuditEntry(CouchbaseChangeEntry entry) { if (entry.shouldBeIgnored()) { - logger.info("Skipping Mongock audit entry with changeId[{}]: state=IGNORED (change was already executed by Mongock; not imported into Flamingock audit history).", - entry.getChangeId()); return null; } return entry.toAuditEntry(); diff --git a/legacy/mongock-importer-dynamodb/src/main/java/io/flamingock/importer/mongock/dynamodb/MongockImporterDynamoDB.java b/legacy/mongock-importer-dynamodb/src/main/java/io/flamingock/importer/mongock/dynamodb/MongockImporterDynamoDB.java index cdc34db6b..eeb0630b0 100644 --- a/legacy/mongock-importer-dynamodb/src/main/java/io/flamingock/importer/mongock/dynamodb/MongockImporterDynamoDB.java +++ b/legacy/mongock-importer-dynamodb/src/main/java/io/flamingock/importer/mongock/dynamodb/MongockImporterDynamoDB.java @@ -17,8 +17,6 @@ import io.flamingock.internal.common.core.audit.AuditEntry; import io.flamingock.internal.common.core.audit.AuditHistoryReader; -import org.slf4j.Logger; -import org.slf4j.LoggerFactory; import software.amazon.awssdk.enhanced.dynamodb.DynamoDbEnhancedClient; import software.amazon.awssdk.enhanced.dynamodb.DynamoDbTable; import software.amazon.awssdk.enhanced.dynamodb.TableSchema; @@ -31,8 +29,6 @@ public class MongockImporterDynamoDB implements AuditHistoryReader { - private static final Logger logger = LoggerFactory.getLogger("MongockImporter"); - private final DynamoDbTable sourceTable; public MongockImporterDynamoDB(DynamoDbClient client, String tableName) { @@ -56,8 +52,6 @@ public List getAuditHistory() { private static AuditEntry toAuditEntry(MongockAuditEntry entry) { if (entry.shouldBeIgnored()) { - logger.info("Skipping Mongock audit entry with changeId[{}]: state=IGNORED (change was already executed by Mongock; not imported into Flamingock audit history).", - entry.getChangeId()); return null; } return entry.toAuditEntry(); diff --git a/legacy/mongock-importer-mongodb-reactive/src/main/java/io/flamingock/importer/mongock/mongodb/reactive/MongockImporterMongoDBReactive.java b/legacy/mongock-importer-mongodb-reactive/src/main/java/io/flamingock/importer/mongock/mongodb/reactive/MongockImporterMongoDBReactive.java index 941235e90..787a07744 100644 --- a/legacy/mongock-importer-mongodb-reactive/src/main/java/io/flamingock/importer/mongock/mongodb/reactive/MongockImporterMongoDBReactive.java +++ b/legacy/mongock-importer-mongodb-reactive/src/main/java/io/flamingock/importer/mongock/mongodb/reactive/MongockImporterMongoDBReactive.java @@ -22,8 +22,6 @@ import io.flamingock.internal.common.core.audit.AuditHistoryReader; import io.flamingock.reactive.util.PublisherSync; import org.bson.Document; -import org.slf4j.Logger; -import org.slf4j.LoggerFactory; import java.time.Instant; import java.time.LocalDateTime; @@ -35,8 +33,6 @@ public class MongockImporterMongoDBReactive implements AuditHistoryReader { - private static final Logger logger = LoggerFactory.getLogger("MongockImporter"); - private final MongoCollection sourceCollection; public MongockImporterMongoDBReactive(MongoDatabase mongoDatabase, String collectionName) { @@ -57,8 +53,6 @@ private static AuditEntry toAuditEntry(Document document) { MongockAuditEntry changeEntry = toChangeEntry(document); if (changeEntry.shouldBeIgnored()) { - logger.info("Skipping Mongock audit entry with changeId[{}]: state=IGNORED (change was already executed by Mongock; not imported into Flamingock audit history).", - changeEntry.getChangeId()); return null; } diff --git a/legacy/mongock-importer-mongodb/src/main/java/io/flamingock/importer/mongock/mongodb/MongockImporterMongoDB.java b/legacy/mongock-importer-mongodb/src/main/java/io/flamingock/importer/mongock/mongodb/MongockImporterMongoDB.java index 943ed88c9..36c3e3e8c 100644 --- a/legacy/mongock-importer-mongodb/src/main/java/io/flamingock/importer/mongock/mongodb/MongockImporterMongoDB.java +++ b/legacy/mongock-importer-mongodb/src/main/java/io/flamingock/importer/mongock/mongodb/MongockImporterMongoDB.java @@ -21,8 +21,6 @@ import io.flamingock.internal.common.core.audit.AuditEntry; import io.flamingock.internal.common.core.audit.AuditHistoryReader; import org.bson.Document; -import org.slf4j.Logger; -import org.slf4j.LoggerFactory; import java.time.Instant; import java.time.LocalDateTime; @@ -35,8 +33,6 @@ public class MongockImporterMongoDB implements AuditHistoryReader { - private static final Logger logger = LoggerFactory.getLogger("MongockImporter"); - private final MongoCollection sourceCollection; public MongockImporterMongoDB(MongoDatabase mongoDatabase, String collectionName) { @@ -61,8 +57,6 @@ private static AuditEntry toAuditEntry(Document document) { .toLocalDateTime(); if (changeEntry.shouldBeIgnored()) { - logger.info("Skipping Mongock audit entry with changeId[{}]: state=IGNORED (change was already executed by Mongock; not imported into Flamingock audit history).", - changeEntry.getChangeId()); return null; } return new AuditEntry(