From 9dc94b8d484e26474b3d146f34046011c39b5c74 Mon Sep 17 00:00:00 2001 From: bercianor Date: Wed, 23 Sep 2026 18:44:20 +0200 Subject: [PATCH] feat(journal): persist idempotency keys --- .../DynamoDBAuditPersistenceJournalTest.java | 1 + .../DynamoDBJournalEventStoreTest.java | 4 +- ...DBReactiveAuditPersistenceJournalTest.java | 2 +- ...ngoDBReactiveJournalEventStoreE2ETest.java | 2 + ...ongoDBSyncAuditPersistenceJournalTest.java | 2 +- .../MongoDBSyncJournalEventStoreE2ETest.java | 1 + .../sql/internal/JournalEventConstants.java | 1 + .../sql/internal/SqlJournalDialectHelper.java | 1 + .../sql/internal/SqlJournalEventMapper.java | 16 +++-- .../store/sql/SqlAuditStoreTest.java | 15 ++++ .../SqlAuditPersistenceJournalTest.java | 1 + .../internal/SqlJournalDialectHelperTest.java | 16 +++-- .../internal/SqlJournalEventMapperTest.java | 15 ++-- .../SqlJournalEventStoreJdbcTest.java | 1 + .../common/core/journal/JournalEvent.java | 13 +++- .../common/core/journal/JournalEventTest.java | 39 +++++++++++ .../core/journal/JournalEventSequencer.java | 40 +++++++++++ .../journal/JournalEventSequencerTest.java | 68 +++++++++++++++++++ .../journal/DynamoDBJournalEventMapper.java | 2 + .../entities/journal/JournalEventEntity.java | 9 +++ .../journal/JournalEventFieldConstants.java | 3 +- .../mongodb/MongoDBJournalEventMapper.java | 3 + .../journal/JournalEventFieldConstants.java | 1 + .../MongoDBJournalEventMapperTest.java | 5 ++ 24 files changed, 235 insertions(+), 26 deletions(-) create mode 100644 core/flamingock-core-commons/src/test/java/io/flamingock/internal/common/core/journal/JournalEventTest.java create mode 100644 core/flamingock-core/src/test/java/io/flamingock/internal/core/journal/JournalEventSequencerTest.java diff --git a/community/flamingock-dynamodb-auditstore/src/test/java/io/flamingock/store/dynamodb/internal/DynamoDBAuditPersistenceJournalTest.java b/community/flamingock-dynamodb-auditstore/src/test/java/io/flamingock/store/dynamodb/internal/DynamoDBAuditPersistenceJournalTest.java index 93b8d12b5..08d0bd189 100644 --- a/community/flamingock-dynamodb-auditstore/src/test/java/io/flamingock/store/dynamodb/internal/DynamoDBAuditPersistenceJournalTest.java +++ b/community/flamingock-dynamodb-auditstore/src/test/java/io/flamingock/store/dynamodb/internal/DynamoDBAuditPersistenceJournalTest.java @@ -263,6 +263,7 @@ private JournalEventSequencer newSequencer() { private void occupyStreamPosition(long sequence) { JournalEvent squatter = new JournalEvent<>( "pre-existing-event", + "key-pre-existing-event", JournalEventType.CHANGE_STATE, JournalEvent.DEFAULT_VERSION, STREAM_ID, diff --git a/community/flamingock-dynamodb-auditstore/src/test/java/io/flamingock/store/dynamodb/internal/DynamoDBJournalEventStoreTest.java b/community/flamingock-dynamodb-auditstore/src/test/java/io/flamingock/store/dynamodb/internal/DynamoDBJournalEventStoreTest.java index 0cee48f3c..d3a1818cb 100644 --- a/community/flamingock-dynamodb-auditstore/src/test/java/io/flamingock/store/dynamodb/internal/DynamoDBJournalEventStoreTest.java +++ b/community/flamingock-dynamodb-auditstore/src/test/java/io/flamingock/store/dynamodb/internal/DynamoDBJournalEventStoreTest.java @@ -291,6 +291,7 @@ void getLastEventByStreamReturnsHighestSequence() throws InterruptedException { assertTrue(last.isPresent()); assertEquals(3L, last.get().getStreamSequence()); assertEquals("event-3", last.get().getEventId()); + assertEquals("key-event-3", last.get().getIdempotencyKey()); assertEquals(JournalEventType.CHANGE_STATE, last.get().getEventType()); assertEquals("change-3", last.get().getData().getChangeId()); @@ -350,6 +351,7 @@ void mapperOmitsBothPendingKeysWhenAcknowledged() { JournalEvent source = journalEvent("stage-1", 1L, "event-1"); JournalEvent acknowledged = new JournalEvent<>( source.getEventId(), + source.getIdempotencyKey(), source.getEventType(), source.getEventVersion(), source.getStreamId(), @@ -535,7 +537,7 @@ private JournalEvent journalEvent(String streamId, long sequence, St "1", RecoveryStrategy.MANUAL_INTERVENTION, true); - return new JournalEvent<>(eventId, JournalEventType.CHANGE_STATE, streamId, sequence, Instant.now(), auditEntry); + return new JournalEvent<>(eventId, "key-" + eventId, JournalEventType.CHANGE_STATE, streamId, sequence, Instant.now(), auditEntry); } private List> awaitUnacknowledgedCount(int expected) throws InterruptedException { diff --git a/community/flamingock-mongodb-reactive-auditstore/src/test/java/io/flamingock/store/mongodb/reactive/internal/MongoDBReactiveAuditPersistenceJournalTest.java b/community/flamingock-mongodb-reactive-auditstore/src/test/java/io/flamingock/store/mongodb/reactive/internal/MongoDBReactiveAuditPersistenceJournalTest.java index b9d1008ff..94aa63ccf 100644 --- a/community/flamingock-mongodb-reactive-auditstore/src/test/java/io/flamingock/store/mongodb/reactive/internal/MongoDBReactiveAuditPersistenceJournalTest.java +++ b/community/flamingock-mongodb-reactive-auditstore/src/test/java/io/flamingock/store/mongodb/reactive/internal/MongoDBReactiveAuditPersistenceJournalTest.java @@ -316,7 +316,7 @@ private MongoDBReactiveAuditPersistence persistenceWithoutTransactions( private void occupyStreamPosition(long streamSequence) { JournalEvent event = new JournalEvent<>( - "pre-existing-event", JournalEventType.CHANGE_STATE, JournalEvent.DEFAULT_VERSION, + "pre-existing-event", "key-pre-existing-event", JournalEventType.CHANGE_STATE, JournalEvent.DEFAULT_VERSION, STREAM_ID, streamSequence, Instant.now(), auditEntry("pre-existing-change", AuditEntry.Status.APPLIED), false); PublisherSync.first(database.getCollection(JOURNAL_COLLECTION).insertOne(mapper.toDocument(event))); } diff --git a/community/flamingock-mongodb-reactive-auditstore/src/test/java/io/flamingock/store/mongodb/reactive/internal/MongoDBReactiveJournalEventStoreE2ETest.java b/community/flamingock-mongodb-reactive-auditstore/src/test/java/io/flamingock/store/mongodb/reactive/internal/MongoDBReactiveJournalEventStoreE2ETest.java index dcdd98c7b..edd3f6898 100644 --- a/community/flamingock-mongodb-reactive-auditstore/src/test/java/io/flamingock/store/mongodb/reactive/internal/MongoDBReactiveJournalEventStoreE2ETest.java +++ b/community/flamingock-mongodb-reactive-auditstore/src/test/java/io/flamingock/store/mongodb/reactive/internal/MongoDBReactiveJournalEventStoreE2ETest.java @@ -331,6 +331,7 @@ private static JournalEvent fixedEvent(String eventId, boolean acknowledged) { return new JournalEvent<>( eventId, + "key-" + eventId, JournalEventType.CHANGE_STATE, 3, streamId, @@ -353,6 +354,7 @@ private static JournalEvent event(String eventId, boolean acknowledged) { return new JournalEvent<>( eventId, + "key-" + eventId, JournalEventType.CHANGE_STATE, JournalEvent.DEFAULT_VERSION, streamId, diff --git a/community/flamingock-mongodb-sync-auditstore/src/test/java/io/flamingock/store/mongodb/sync/internal/MongoDBSyncAuditPersistenceJournalTest.java b/community/flamingock-mongodb-sync-auditstore/src/test/java/io/flamingock/store/mongodb/sync/internal/MongoDBSyncAuditPersistenceJournalTest.java index 31b2d356b..4740d39d5 100644 --- a/community/flamingock-mongodb-sync-auditstore/src/test/java/io/flamingock/store/mongodb/sync/internal/MongoDBSyncAuditPersistenceJournalTest.java +++ b/community/flamingock-mongodb-sync-auditstore/src/test/java/io/flamingock/store/mongodb/sync/internal/MongoDBSyncAuditPersistenceJournalTest.java @@ -344,7 +344,7 @@ private MongoDBSyncAuditPersistence persistenceWithoutTransactions( */ private void occupyStreamPosition(long streamSequence) { JournalEvent squatter = new JournalEvent<>( - "pre-existing-event", JournalEventType.CHANGE_STATE, JournalEvent.DEFAULT_VERSION, + "pre-existing-event", "key-pre-existing-event", JournalEventType.CHANGE_STATE, JournalEvent.DEFAULT_VERSION, STREAM_ID, streamSequence, Instant.now(), auditEntry("pre-existing-change"), false); database.getCollection(JOURNAL_COLLECTION).insertOne(mapper.toDocument(squatter)); } diff --git a/community/flamingock-mongodb-sync-auditstore/src/test/java/io/flamingock/store/mongodb/sync/internal/MongoDBSyncJournalEventStoreE2ETest.java b/community/flamingock-mongodb-sync-auditstore/src/test/java/io/flamingock/store/mongodb/sync/internal/MongoDBSyncJournalEventStoreE2ETest.java index ef5005135..92e6093d2 100644 --- a/community/flamingock-mongodb-sync-auditstore/src/test/java/io/flamingock/store/mongodb/sync/internal/MongoDBSyncJournalEventStoreE2ETest.java +++ b/community/flamingock-mongodb-sync-auditstore/src/test/java/io/flamingock/store/mongodb/sync/internal/MongoDBSyncJournalEventStoreE2ETest.java @@ -275,6 +275,7 @@ private static List ids(List> events) { private static JournalEvent event(String eventId, String streamId, long sequence, boolean acknowledged) { return new JournalEvent<>( eventId, + "key-" + eventId, JournalEventType.CHANGE_STATE, JournalEvent.DEFAULT_VERSION, streamId, diff --git a/community/flamingock-sql-auditstore/src/main/java/io/flamingock/store/sql/internal/JournalEventConstants.java b/community/flamingock-sql-auditstore/src/main/java/io/flamingock/store/sql/internal/JournalEventConstants.java index 1743df382..d6933d265 100644 --- a/community/flamingock-sql-auditstore/src/main/java/io/flamingock/store/sql/internal/JournalEventConstants.java +++ b/community/flamingock-sql-auditstore/src/main/java/io/flamingock/store/sql/internal/JournalEventConstants.java @@ -21,6 +21,7 @@ final class JournalEventConstants { static final String EVENT_ID = "event_id"; + static final String IDEMPOTENCY_KEY = "idempotency_key"; static final String EVENT_TYPE = "event_type"; static final String EVENT_VERSION = "event_version"; static final String STREAM_ID = "stream_id"; diff --git a/community/flamingock-sql-auditstore/src/main/java/io/flamingock/store/sql/internal/SqlJournalDialectHelper.java b/community/flamingock-sql-auditstore/src/main/java/io/flamingock/store/sql/internal/SqlJournalDialectHelper.java index 9e93d7df5..9e531fc91 100644 --- a/community/flamingock-sql-auditstore/src/main/java/io/flamingock/store/sql/internal/SqlJournalDialectHelper.java +++ b/community/flamingock-sql-auditstore/src/main/java/io/flamingock/store/sql/internal/SqlJournalDialectHelper.java @@ -67,6 +67,7 @@ List getIndexNames(String tableName) { List getColumnDefinitions() { return Collections.unmodifiableList(Arrays.asList( new ColumnDefinition(JournalEventConstants.EVENT_ID, ColumnType.VARCHAR, 255, false), + new ColumnDefinition(JournalEventConstants.IDEMPOTENCY_KEY, ColumnType.VARCHAR, 64, false), new ColumnDefinition(JournalEventConstants.EVENT_TYPE, ColumnType.VARCHAR, 32, false), new ColumnDefinition(JournalEventConstants.EVENT_VERSION, ColumnType.INTEGER, 0, false), new ColumnDefinition(JournalEventConstants.STREAM_ID, ColumnType.VARCHAR, 255, false), diff --git a/community/flamingock-sql-auditstore/src/main/java/io/flamingock/store/sql/internal/SqlJournalEventMapper.java b/community/flamingock-sql-auditstore/src/main/java/io/flamingock/store/sql/internal/SqlJournalEventMapper.java index cc638f319..c08b2402f 100644 --- a/community/flamingock-sql-auditstore/src/main/java/io/flamingock/store/sql/internal/SqlJournalEventMapper.java +++ b/community/flamingock-sql-auditstore/src/main/java/io/flamingock/store/sql/internal/SqlJournalEventMapper.java @@ -42,13 +42,14 @@ final class SqlJournalEventMapper { void bind(PreparedStatement statement, JournalEvent event) throws SQLException { requireSupportedEvent(event); statement.setString(1, event.getEventId()); - statement.setString(2, event.getEventType().name()); - statement.setInt(3, event.getEventVersion()); - statement.setString(4, event.getStreamId()); - statement.setLong(5, event.getStreamSequence()); - statement.setTimestamp(6, Timestamp.from(event.getOccurredAt())); - statement.setBoolean(7, event.isAcknowledged()); - statement.setString(8, SqlJournalPayloadCodec.serialize(event.getData())); + statement.setString(2, event.getIdempotencyKey()); + statement.setString(3, event.getEventType().name()); + statement.setInt(4, event.getEventVersion()); + statement.setString(5, event.getStreamId()); + statement.setLong(6, event.getStreamSequence()); + statement.setTimestamp(7, Timestamp.from(event.getOccurredAt())); + statement.setBoolean(8, event.isAcknowledged()); + statement.setString(9, SqlJournalPayloadCodec.serialize(event.getData())); } JournalEvent fromResultSet(ResultSet resultSet) throws SQLException { @@ -67,6 +68,7 @@ JournalEvent fromResultSet(ResultSet resultSet) throws SQLException return new JournalEvent<>( resultSet.getString(JournalEventConstants.EVENT_ID), + resultSet.getString(JournalEventConstants.IDEMPOTENCY_KEY), eventType, resultSet.getInt(JournalEventConstants.EVENT_VERSION), resultSet.getString(JournalEventConstants.STREAM_ID), diff --git a/community/flamingock-sql-auditstore/src/test/java/io/flamingock/store/sql/SqlAuditStoreTest.java b/community/flamingock-sql-auditstore/src/test/java/io/flamingock/store/sql/SqlAuditStoreTest.java index 89a323f20..75b6fb997 100644 --- a/community/flamingock-sql-auditstore/src/test/java/io/flamingock/store/sql/SqlAuditStoreTest.java +++ b/community/flamingock-sql-auditstore/src/test/java/io/flamingock/store/sql/SqlAuditStoreTest.java @@ -465,6 +465,7 @@ void journalEnabledRoundTripsAcrossRuntimeDialects(SqlDialect sqlDialect, String persistence.writeEntry(auditEntry("matrix-change", AuditEntry.Status.APPLIED)); assertEquals(1, auditStore.getAuditReader().getAuditHistory().size()); + String expectedIdempotencyKey = persistedJournalIdempotencyKey("matrix-stage", 2L); SqlJournalEventStore journalEventStore = new SqlJournalEventStore( context.dataSource, "flamingockJournalEvents", targetSystem.getTxWrapper()); @@ -474,6 +475,7 @@ void journalEnabledRoundTripsAcrossRuntimeDialects(SqlDialect sqlDialect, String assertEquals("matrix-stage", journalEvent.getStreamId()); assertEquals(2L, journalEvent.getStreamSequence()); assertEquals(JournalEventType.CHANGE_STATE, journalEvent.getEventType()); + assertEquals(expectedIdempotencyKey, journalEvent.getIdempotencyKey()); assertEquals("matrix-change", journalEvent.getData().getChangeId()); assertEquals(AuditEntry.Status.APPLIED, journalEvent.getData().getState()); assertEquals(AuditTxType.NON_TX, journalEvent.getData().getTxType()); @@ -635,6 +637,19 @@ void stageFactoryDoesNotReinitializeAuditReadiness(SqlDialect sqlDialect, String assertTrue(tableExists("flamingockJournalEvents")); } + private String persistedJournalIdempotencyKey(String streamId, long streamSequence) throws SQLException { + try (Connection connection = context.dataSource.getConnection(); + PreparedStatement statement = connection.prepareStatement( + "SELECT idempotency_key FROM flamingockJournalEvents WHERE stream_id = ? AND stream_sequence = ?")) { + statement.setString(1, streamId); + statement.setLong(2, streamSequence); + try (ResultSet resultSet = statement.executeQuery()) { + assertTrue(resultSet.next(), "Expected persisted journal event"); + return resultSet.getString("idempotency_key"); + } + } + } + private int countRows(String tableName) throws SQLException { try (Connection connection = context.dataSource.getConnection(); Statement statement = connection.createStatement(); diff --git a/community/flamingock-sql-auditstore/src/test/java/io/flamingock/store/sql/internal/SqlAuditPersistenceJournalTest.java b/community/flamingock-sql-auditstore/src/test/java/io/flamingock/store/sql/internal/SqlAuditPersistenceJournalTest.java index 36f793c2d..7c7645c41 100644 --- a/community/flamingock-sql-auditstore/src/test/java/io/flamingock/store/sql/internal/SqlAuditPersistenceJournalTest.java +++ b/community/flamingock-sql-auditstore/src/test/java/io/flamingock/store/sql/internal/SqlAuditPersistenceJournalTest.java @@ -662,6 +662,7 @@ private static JournalEvent event(String streamId, String changeId) { return new JournalEvent<>( eventId, + "key-" + eventId, JournalEventType.CHANGE_STATE, JournalEvent.DEFAULT_VERSION, streamId, diff --git a/community/flamingock-sql-auditstore/src/test/java/io/flamingock/store/sql/internal/SqlJournalDialectHelperTest.java b/community/flamingock-sql-auditstore/src/test/java/io/flamingock/store/sql/internal/SqlJournalDialectHelperTest.java index 3890f3fc1..878c2a9b5 100644 --- a/community/flamingock-sql-auditstore/src/test/java/io/flamingock/store/sql/internal/SqlJournalDialectHelperTest.java +++ b/community/flamingock-sql-auditstore/src/test/java/io/flamingock/store/sql/internal/SqlJournalDialectHelperTest.java @@ -45,6 +45,7 @@ void generatesTypedSchemaAndIndexesForEveryDialect(SqlDialect dialect) { List indexSql = helper.getCreateIndexSqlStrings(TABLE_NAME); assertTrue(ddl.contains("EVENT_ID")); + assertTrue(ddl.contains("IDEMPOTENCY_KEY")); assertTrue(ddl.contains("EVENT_TYPE")); assertTrue(ddl.contains("EVENT_VERSION")); assertTrue(ddl.contains("STREAM_ID")); @@ -65,7 +66,7 @@ void generatesTypedSchemaAndIndexesForEveryDialect(SqlDialect dialect) { .allMatch(name -> name.length() <= helper.getMaximumIndexNameLength())); List definitionNames = columnNames(helper.getColumnDefinitions()); - assertTrue(Arrays.asList("event_id", "stream_id", "stream_sequence", "occurred_at", "acknowledged") + assertTrue(Arrays.asList("event_id", "idempotency_key", "stream_id", "stream_sequence", "occurred_at", "acknowledged") .stream().allMatch(definitionNames::contains)); assertEquals(definitionNames, insertColumnNames(helper.getInsertSqlString(TABLE_NAME))); } @@ -103,6 +104,7 @@ void usesExactPortableTypes(SqlDialect dialect) { String ddl = helper.getCreateTableSqlString(TABLE_NAME).toUpperCase(Locale.ROOT); assertTrue(ddl.contains("EVENT_ID " + varcharType(dialect, 255) + " NOT NULL")); + assertTrue(ddl.contains("IDEMPOTENCY_KEY " + varcharType(dialect, 64) + " NOT NULL")); assertTrue(ddl.contains("EVENT_TYPE " + varcharType(dialect, 32) + " NOT NULL")); assertTrue(ddl.contains("EVENT_VERSION INTEGER NOT NULL")); assertTrue(ddl.contains("STREAM_ID " + varcharType(dialect, 255) + " NOT NULL")); @@ -130,13 +132,13 @@ void keepsTypedColumnDefinitionsAndTextCapacity() { SqlJournalDialectHelper helper = new SqlJournalDialectHelper(SqlDialect.H2); assertEquals(Arrays.asList( - "event_id", "event_type", "event_version", "stream_id", "stream_sequence", "occurred_at", - "acknowledged", "payload"), + "event_id", "idempotency_key", "event_type", "event_version", "stream_id", "stream_sequence", + "occurred_at", "acknowledged", "payload"), columnNames(helper.getColumnDefinitions())); - assertEquals(8, helper.getColumnDefinitions().size()); - assertEquals(SqlJournalDialectHelper.ColumnType.TEXT, helper.getColumnDefinitions().get(7).type); - assertEquals(2048, helper.getColumnDefinitions().get(7).size); - assertFalse(helper.getColumnDefinitions().get(7).nullable); + assertEquals(9, helper.getColumnDefinitions().size()); + assertEquals(SqlJournalDialectHelper.ColumnType.TEXT, helper.getColumnDefinitions().get(8).type); + assertEquals(2048, helper.getColumnDefinitions().get(8).size); + assertFalse(helper.getColumnDefinitions().get(8).nullable); } private static List columnNames(List definitions) { diff --git a/community/flamingock-sql-auditstore/src/test/java/io/flamingock/store/sql/internal/SqlJournalEventMapperTest.java b/community/flamingock-sql-auditstore/src/test/java/io/flamingock/store/sql/internal/SqlJournalEventMapperTest.java index 46b537e01..c72425950 100644 --- a/community/flamingock-sql-auditstore/src/test/java/io/flamingock/store/sql/internal/SqlJournalEventMapperTest.java +++ b/community/flamingock-sql-auditstore/src/test/java/io/flamingock/store/sql/internal/SqlJournalEventMapperTest.java @@ -51,7 +51,7 @@ void roundTripsTypedChangeStateEvent() throws Exception { AuditEntry auditEntry = auditEntry(); Instant occurredAt = Instant.parse("2026-08-11T10:20:30.123456Z"); JournalEvent source = new JournalEvent<>( - "event-1", JournalEventType.CHANGE_STATE, JournalEvent.DEFAULT_VERSION, + "event-1", "key-event-1", JournalEventType.CHANGE_STATE, JournalEvent.DEFAULT_VERSION, "stage-1", 7L, occurredAt, auditEntry, false); try (Connection connection = DriverManager.getConnection("jdbc:h2:mem:journal_mapper;DB_CLOSE_DELAY=-1")) { @@ -80,6 +80,7 @@ void roundTripsTypedChangeStateEvent() throws Exception { JournalEvent actual = new SqlJournalEventMapper().fromResultSet(resultSet); assertEquals(source.getEventId(), actual.getEventId()); + assertEquals(source.getIdempotencyKey(), actual.getIdempotencyKey()); assertEquals(source.getEventType(), actual.getEventType()); assertEquals(source.getEventVersion(), actual.getEventVersion()); assertEquals(source.getStreamId(), actual.getStreamId()); @@ -96,7 +97,7 @@ void roundTripsTypedChangeStateEvent() throws Exception { @DisplayName("rejects event types whose payload mapping is not implemented") void rejectsUnsupportedEventType() throws Exception { JournalEvent unsupported = new JournalEvent<>( - "event-unsupported", JournalEventType.EXECUTION_STATE, "stage-1", 1L, Instant.now(), auditEntry()); + "event-unsupported", "key-event-unsupported", JournalEventType.EXECUTION_STATE, "stage-1", 1L, Instant.now(), auditEntry()); try (Connection connection = DriverManager.getConnection("jdbc:h2:mem:journal_mapper_unsupported;DB_CLOSE_DELAY=-1"); PreparedStatement statement = connection.prepareStatement("SELECT 1")) { @@ -114,7 +115,7 @@ void roundTripsAcknowledgedEventAndNullableFields() throws Exception { null, null, null, 0L, null, null, false, null, null, null, null, null, null); JournalEvent source = new JournalEvent<>( - "event-acknowledged", JournalEventType.CHANGE_STATE, JournalEvent.DEFAULT_VERSION, + "event-acknowledged", "key-event-acknowledged", JournalEventType.CHANGE_STATE, JournalEvent.DEFAULT_VERSION, "stage-nullable", 2L, Instant.parse("2026-08-11T10:20:30Z"), auditEntry, true); @@ -148,7 +149,7 @@ void roundTripsAcknowledgedEventAndNullableFields() throws Exception { void keepsEnvelopeAndPayloadTimesDistinct() throws Exception { AuditEntry auditEntry = auditEntry(); JournalEvent source = new JournalEvent<>( - "event-time", JournalEventType.CHANGE_STATE, "stage-time", 1L, + "event-time", "key-event-time", JournalEventType.CHANGE_STATE, "stage-time", 1L, Instant.parse("2026-08-11T12:00:00Z"), auditEntry); try (Connection connection = DriverManager.getConnection("jdbc:h2:mem:journal_mapper_times;DB_CLOSE_DELAY=-1")) { @@ -175,11 +176,11 @@ void keepsEnvelopeAndPayloadTimesDistinct() throws Exception { @DisplayName("enforces stream position uniqueness without requiring globally unique event IDs") void enforcesCompositeStreamPositionAndAllowsDuplicateEventIds() throws Exception { JournalEvent first = new JournalEvent<>( - "event-shared", JournalEventType.CHANGE_STATE, "stage-1", 1L, Instant.now(), auditEntry()); + "event-shared", "key-event-shared-1", JournalEventType.CHANGE_STATE, "stage-1", 1L, Instant.now(), auditEntry()); JournalEvent otherStream = new JournalEvent<>( - "event-shared", JournalEventType.CHANGE_STATE, "stage-2", 1L, Instant.now(), auditEntry()); + "event-shared", "key-event-shared-2", JournalEventType.CHANGE_STATE, "stage-2", 1L, Instant.now(), auditEntry()); JournalEvent collidingPosition = new JournalEvent<>( - "event-other", JournalEventType.CHANGE_STATE, "stage-1", 1L, Instant.now(), auditEntry()); + "event-other", "key-event-other", JournalEventType.CHANGE_STATE, "stage-1", 1L, Instant.now(), auditEntry()); try (Connection connection = DriverManager.getConnection("jdbc:h2:mem:journal_mapper_identity;DB_CLOSE_DELAY=-1")) { SqlJournalDialectHelper dialectHelper = new SqlJournalDialectHelper(SqlDialect.H2); diff --git a/community/flamingock-sql-auditstore/src/test/java/io/flamingock/store/sql/internal/SqlJournalEventStoreJdbcTest.java b/community/flamingock-sql-auditstore/src/test/java/io/flamingock/store/sql/internal/SqlJournalEventStoreJdbcTest.java index 1ae7f69a5..5a3f97e68 100644 --- a/community/flamingock-sql-auditstore/src/test/java/io/flamingock/store/sql/internal/SqlJournalEventStoreJdbcTest.java +++ b/community/flamingock-sql-auditstore/src/test/java/io/flamingock/store/sql/internal/SqlJournalEventStoreJdbcTest.java @@ -495,6 +495,7 @@ private static JournalEvent event(String streamId, boolean acknowledged) { return new JournalEvent<>( eventId, + "key-" + eventId, JournalEventType.CHANGE_STATE, JournalEvent.DEFAULT_VERSION, streamId, diff --git a/core/flamingock-core-commons/src/main/java/io/flamingock/internal/common/core/journal/JournalEvent.java b/core/flamingock-core-commons/src/main/java/io/flamingock/internal/common/core/journal/JournalEvent.java index fbef90afa..9d9066d92 100644 --- a/core/flamingock-core-commons/src/main/java/io/flamingock/internal/common/core/journal/JournalEvent.java +++ b/core/flamingock-core-commons/src/main/java/io/flamingock/internal/common/core/journal/JournalEvent.java @@ -32,6 +32,7 @@ public final class JournalEvent { public static final int DEFAULT_VERSION = 1; private final String eventId; + private final String idempotencyKey; private final JournalEventType eventType; private final int eventVersion; @@ -44,15 +45,17 @@ public final class JournalEvent { private boolean acknowledged; public JournalEvent(String eventId, + String idempotencyKey, JournalEventType eventType, String streamId, long streamSequence, Instant occurredAt, T data) { - this(eventId, eventType, DEFAULT_VERSION, streamId, streamSequence, occurredAt, data, false); + this(eventId, idempotencyKey, eventType, DEFAULT_VERSION, streamId, streamSequence, occurredAt, data, false); } public JournalEvent(String eventId, + String idempotencyKey, JournalEventType eventType, int eventVersion, String streamId, @@ -62,6 +65,7 @@ public JournalEvent(String eventId, boolean acknowledged) { this.acknowledged = acknowledged; this.eventId = requireNotBlank(eventId, "eventId"); + this.idempotencyKey = requireNotBlank(idempotencyKey, "idempotencyKey"); this.eventType = Objects.requireNonNull(eventType, "eventType"); this.streamId = requireNotBlank(streamId, "streamId"); @@ -79,6 +83,13 @@ public String getEventId() { return eventId; } + /** + * Returns the stable logical identity used by downstream receivers to deduplicate this event. + */ + public String getIdempotencyKey() { + return idempotencyKey; + } + public JournalEventType getEventType() { return eventType; } diff --git a/core/flamingock-core-commons/src/test/java/io/flamingock/internal/common/core/journal/JournalEventTest.java b/core/flamingock-core-commons/src/test/java/io/flamingock/internal/common/core/journal/JournalEventTest.java new file mode 100644 index 000000000..1ed545cdc --- /dev/null +++ b/core/flamingock-core-commons/src/test/java/io/flamingock/internal/common/core/journal/JournalEventTest.java @@ -0,0 +1,39 @@ +/* + * Copyright 2026 Flamingock (https://www.flamingock.io) + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package io.flamingock.internal.common.core.journal; + +import org.junit.jupiter.api.DisplayName; +import org.junit.jupiter.api.Test; + +import java.time.Instant; + +import static org.junit.jupiter.api.Assertions.assertThrows; + +class JournalEventTest { + + @Test + @DisplayName("rejects null and blank idempotency keys") + void rejectsNullAndBlankIdempotencyKeys() { + assertThrows(IllegalArgumentException.class, () -> event(null)); + assertThrows(IllegalArgumentException.class, () -> event(" ")); + } + + private static JournalEvent event(String idempotencyKey) { + return new JournalEvent<>( + "event-1", idempotencyKey, JournalEventType.CHANGE_STATE, "stream-1", 1L, + Instant.parse("2026-09-18T00:00:00Z"), "payload"); + } +} diff --git a/core/flamingock-core/src/main/java/io/flamingock/internal/core/journal/JournalEventSequencer.java b/core/flamingock-core/src/main/java/io/flamingock/internal/core/journal/JournalEventSequencer.java index 1aeb52f02..3d9dee8da 100644 --- a/core/flamingock-core/src/main/java/io/flamingock/internal/core/journal/JournalEventSequencer.java +++ b/core/flamingock-core/src/main/java/io/flamingock/internal/core/journal/JournalEventSequencer.java @@ -20,6 +20,9 @@ import io.flamingock.internal.common.core.journal.JournalEventType; import org.jetbrains.annotations.NotNull; +import java.nio.charset.StandardCharsets; +import java.security.MessageDigest; +import java.security.NoSuchAlgorithmException; import java.time.Instant; import java.util.UUID; @@ -69,6 +72,7 @@ private JournalEvent getAuditEntryJournalEvent(AuditEntry payload, J pendingConfirmation = true; return new JournalEvent<>( UUID.randomUUID().toString(), // eventId + deriveIdempotencyKey(type, payload), type, streamId, nextSequence, // spent only on confirm(), so a failed write leaves no gap @@ -76,4 +80,40 @@ private JournalEvent getAuditEntryJournalEvent(AuditEntry payload, J payload); } + private String deriveIdempotencyKey(JournalEventType eventType, AuditEntry payload) { + if (eventType != JournalEventType.CHANGE_STATE) { + throw new UnsupportedOperationException("No idempotency-key derivation is defined for " + eventType); + } + try { + MessageDigest digest = MessageDigest.getInstance("SHA-256"); + updateCanonicalField(digest, streamId); + updateCanonicalField(digest, payload.getExecutionId()); + updateCanonicalField(digest, payload.getChangeId()); + updateCanonicalField(digest, payload.getState() == null ? null : payload.getState().name()); + return toHex(digest.digest()); + } catch (NoSuchAlgorithmException exception) { + throw new IllegalStateException("SHA-256 is unavailable", exception); + } + } + + private static void updateCanonicalField(MessageDigest digest, String value) { + if (value == null) { + digest.update((byte) 0); + return; + } + byte[] bytes = value.getBytes(StandardCharsets.UTF_8); + digest.update((byte) 1); + digest.update(Integer.toString(bytes.length).getBytes(StandardCharsets.US_ASCII)); + digest.update((byte) ':'); + digest.update(bytes); + } + + private static String toHex(byte[] bytes) { + StringBuilder result = new StringBuilder(bytes.length * 2); + for (byte value : bytes) { + result.append(String.format("%02x", value & 0xff)); + } + return result.toString(); + } + } diff --git a/core/flamingock-core/src/test/java/io/flamingock/internal/core/journal/JournalEventSequencerTest.java b/core/flamingock-core/src/test/java/io/flamingock/internal/core/journal/JournalEventSequencerTest.java new file mode 100644 index 000000000..66e0548e5 --- /dev/null +++ b/core/flamingock-core/src/test/java/io/flamingock/internal/core/journal/JournalEventSequencerTest.java @@ -0,0 +1,68 @@ +/* + * Copyright 2026 Flamingock (https://www.flamingock.io) + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package io.flamingock.internal.core.journal; + +import io.flamingock.api.RecoveryStrategy; +import io.flamingock.internal.common.core.audit.AuditEntry; +import io.flamingock.internal.common.core.audit.AuditTxType; +import io.flamingock.internal.common.core.journal.JournalEvent; +import org.junit.jupiter.api.DisplayName; +import org.junit.jupiter.api.Test; + +import java.time.LocalDateTime; + +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertNotEquals; +import static org.junit.jupiter.api.Assertions.assertNotNull; + +class JournalEventSequencerTest { + + @Test + @DisplayName("derives a stable idempotency key from the CHANGE_STATE identity") + void derivesStableIdempotencyKeyFromChangeStateIdentity() { + AuditEntry source = auditEntry("execution-1", "change-1", AuditEntry.Status.APPLIED); + AuditEntry differentNonIdentityFields = auditEntry( + "execution-1", "change-1", AuditEntry.Status.APPLIED, "author-2", LocalDateTime.of(2026, 2, 1, 0, 0)); + JournalEvent first = new JournalEventSequencer("stage-1", 1L).newEvent(source); + JournalEvent second = new JournalEventSequencer("stage-1", 1L).newEvent(differentNonIdentityFields); + JournalEvent differentStream = new JournalEventSequencer("stage-2", 1L).newEvent(source); + JournalEvent differentExecution = new JournalEventSequencer("stage-1", 1L) + .newEvent(auditEntry("execution-2", "change-1", AuditEntry.Status.APPLIED)); + JournalEvent differentChange = new JournalEventSequencer("stage-1", 1L) + .newEvent(auditEntry("execution-1", "change-2", AuditEntry.Status.APPLIED)); + JournalEvent differentState = new JournalEventSequencer("stage-1", 1L) + .newEvent(auditEntry("execution-1", "change-1", AuditEntry.Status.FAILED)); + + assertNotNull(first.getIdempotencyKey()); + assertEquals(first.getIdempotencyKey(), second.getIdempotencyKey()); + assertNotEquals(first.getIdempotencyKey(), differentStream.getIdempotencyKey()); + assertNotEquals(first.getIdempotencyKey(), differentExecution.getIdempotencyKey()); + assertNotEquals(first.getIdempotencyKey(), differentChange.getIdempotencyKey()); + assertNotEquals(first.getIdempotencyKey(), differentState.getIdempotencyKey()); + } + + private static AuditEntry auditEntry(String executionId, String changeId, AuditEntry.Status state) { + return auditEntry(executionId, changeId, state, "author-1", LocalDateTime.of(2026, 1, 1, 0, 0)); + } + + private static AuditEntry auditEntry( + String executionId, String changeId, AuditEntry.Status state, String author, LocalDateTime createdAt) { + return new AuditEntry( + executionId, "stage-1", changeId, author, createdAt, state, + AuditEntry.ChangeType.STANDARD_CODE, "example.Change", "apply", "Change.java", 1L, "host-1", + null, false, null, AuditTxType.NON_TX, "mongodb", "001", RecoveryStrategy.MANUAL_INTERVENTION, true); + } +} diff --git a/utils/dynamodb-util/src/main/java/io/flamingock/internal/util/dynamodb/entities/journal/DynamoDBJournalEventMapper.java b/utils/dynamodb-util/src/main/java/io/flamingock/internal/util/dynamodb/entities/journal/DynamoDBJournalEventMapper.java index d51c0e51e..a9aa882a5 100644 --- a/utils/dynamodb-util/src/main/java/io/flamingock/internal/util/dynamodb/entities/journal/DynamoDBJournalEventMapper.java +++ b/utils/dynamodb-util/src/main/java/io/flamingock/internal/util/dynamodb/entities/journal/DynamoDBJournalEventMapper.java @@ -50,6 +50,7 @@ public static JournalEventEntity toEntity(JournalEvent event) { requireSupportedType(event.getEventType()); JournalEventEntity entity = new JournalEventEntity(); entity.setEventId(event.getEventId()); + entity.setIdempotencyKey(event.getIdempotencyKey()); entity.setEventType(event.getEventType().name()); entity.setEventVersion(event.getEventVersion()); entity.setStreamId(event.getStreamId()); @@ -79,6 +80,7 @@ public static JournalEvent fromEntity(JournalEventEntity entity) { boolean acknowledged = partitionKeyMissing; return new JournalEvent<>( entity.getEventId(), + entity.getIdempotencyKey(), eventType, entity.getEventVersion() != null ? entity.getEventVersion() : JournalEvent.DEFAULT_VERSION, entity.getStreamId(), diff --git a/utils/dynamodb-util/src/main/java/io/flamingock/internal/util/dynamodb/entities/journal/JournalEventEntity.java b/utils/dynamodb-util/src/main/java/io/flamingock/internal/util/dynamodb/entities/journal/JournalEventEntity.java index 721ef4e80..bf660a7da 100644 --- a/utils/dynamodb-util/src/main/java/io/flamingock/internal/util/dynamodb/entities/journal/JournalEventEntity.java +++ b/utils/dynamodb-util/src/main/java/io/flamingock/internal/util/dynamodb/entities/journal/JournalEventEntity.java @@ -38,6 +38,7 @@ public class JournalEventEntity { private String pendingPartitionKey; private String pendingOrderKey; private String eventId; + private String idempotencyKey; private String eventType; private String occurredAt; private Integer eventVersion; @@ -91,6 +92,14 @@ public void setEventId(String eventId) { this.eventId = eventId; } + public String getIdempotencyKey() { + return idempotencyKey; + } + + public void setIdempotencyKey(String idempotencyKey) { + this.idempotencyKey = idempotencyKey; + } + public String getEventType() { return eventType; } diff --git a/utils/dynamodb-util/src/main/java/io/flamingock/internal/util/dynamodb/entities/journal/JournalEventFieldConstants.java b/utils/dynamodb-util/src/main/java/io/flamingock/internal/util/dynamodb/entities/journal/JournalEventFieldConstants.java index 6b55747cb..0aa5a2308 100644 --- a/utils/dynamodb-util/src/main/java/io/flamingock/internal/util/dynamodb/entities/journal/JournalEventFieldConstants.java +++ b/utils/dynamodb-util/src/main/java/io/flamingock/internal/util/dynamodb/entities/journal/JournalEventFieldConstants.java @@ -23,7 +23,7 @@ * ({@link #PENDING_EVENTS_INDEX}) carries only items with a {@code pendingPartitionKey} and * {@code pendingOrderKey}, * and the non-unique eventId GSI ({@link #EVENT_ID_INDEX}) serves acknowledgement - * lookups only. Event identity is enforced transactionally by a reserved item in this table. + * lookups only. Stream-position uniqueness is provided by the table's composite base key. */ public final class JournalEventFieldConstants { @@ -34,6 +34,7 @@ public final class JournalEventFieldConstants { public static final String KEY_PENDING_PARTITION_KEY = "pendingPartitionKey"; public static final String KEY_PENDING_ORDER_KEY = "pendingOrderKey"; public static final String KEY_EVENT_ID = "eventId"; + public static final String KEY_IDEMPOTENCY_KEY = "idempotencyKey"; public static final String PENDING_PARTITION_VALUE = "pending"; diff --git a/utils/mongodb-util/src/main/java/io/flamingock/internal/common/mongodb/MongoDBJournalEventMapper.java b/utils/mongodb-util/src/main/java/io/flamingock/internal/common/mongodb/MongoDBJournalEventMapper.java index 9513c3207..c29c034ae 100644 --- a/utils/mongodb-util/src/main/java/io/flamingock/internal/common/mongodb/MongoDBJournalEventMapper.java +++ b/utils/mongodb-util/src/main/java/io/flamingock/internal/common/mongodb/MongoDBJournalEventMapper.java @@ -26,6 +26,7 @@ import static io.flamingock.internal.common.mongodb.journal.JournalEventFieldConstants.KEY_ACKNOWLEDGED; import static io.flamingock.internal.common.mongodb.journal.JournalEventFieldConstants.KEY_DATA; import static io.flamingock.internal.common.mongodb.journal.JournalEventFieldConstants.KEY_EVENT_ID; +import static io.flamingock.internal.common.mongodb.journal.JournalEventFieldConstants.KEY_IDEMPOTENCY_KEY; import static io.flamingock.internal.common.mongodb.journal.JournalEventFieldConstants.KEY_EVENT_TYPE; import static io.flamingock.internal.common.mongodb.journal.JournalEventFieldConstants.KEY_EVENT_VERSION; import static io.flamingock.internal.common.mongodb.journal.JournalEventFieldConstants.KEY_OCCURRED_AT; @@ -55,6 +56,7 @@ public Document toDocument(JournalEvent event) { requireSupportedType(event.getEventType()); Document document = new Document(); document.append(KEY_EVENT_ID, event.getEventId()); + document.append(KEY_IDEMPOTENCY_KEY, event.getIdempotencyKey()); document.append(KEY_EVENT_TYPE, event.getEventType().name()); document.append(KEY_EVENT_VERSION, event.getEventVersion()); document.append(KEY_STREAM_ID, event.getStreamId()); @@ -72,6 +74,7 @@ public JournalEvent fromDocument(Document document) { Instant occurredAt = ((Date) document.get(KEY_OCCURRED_AT)).toInstant(); return new JournalEvent<>( document.getString(KEY_EVENT_ID), + document.getString(KEY_IDEMPOTENCY_KEY), eventType, ((Number) document.get(KEY_EVENT_VERSION)).intValue(), document.getString(KEY_STREAM_ID), diff --git a/utils/mongodb-util/src/main/java/io/flamingock/internal/common/mongodb/journal/JournalEventFieldConstants.java b/utils/mongodb-util/src/main/java/io/flamingock/internal/common/mongodb/journal/JournalEventFieldConstants.java index e74c286ae..756043c33 100644 --- a/utils/mongodb-util/src/main/java/io/flamingock/internal/common/mongodb/journal/JournalEventFieldConstants.java +++ b/utils/mongodb-util/src/main/java/io/flamingock/internal/common/mongodb/journal/JournalEventFieldConstants.java @@ -24,6 +24,7 @@ public final class JournalEventFieldConstants { public static final String KEY_EVENT_ID = "eventId"; + public static final String KEY_IDEMPOTENCY_KEY = "idempotencyKey"; public static final String KEY_EVENT_TYPE = "eventType"; public static final String KEY_EVENT_VERSION = "eventVersion"; public static final String KEY_STREAM_ID = "streamId"; diff --git a/utils/mongodb-util/src/test/java/io/flamingock/internal/common/mongodb/MongoDBJournalEventMapperTest.java b/utils/mongodb-util/src/test/java/io/flamingock/internal/common/mongodb/MongoDBJournalEventMapperTest.java index 671d99083..2e8d5ae47 100644 --- a/utils/mongodb-util/src/test/java/io/flamingock/internal/common/mongodb/MongoDBJournalEventMapperTest.java +++ b/utils/mongodb-util/src/test/java/io/flamingock/internal/common/mongodb/MongoDBJournalEventMapperTest.java @@ -38,6 +38,7 @@ void roundTripsChangeStateEventWithAuditEntryPayload() { Instant occurredAt = Instant.parse("2026-07-21T10:15:30Z"); JournalEvent event = new JournalEvent<>( "evt-1", + "key-evt-1", JournalEventType.CHANGE_STATE, JournalEvent.DEFAULT_VERSION, "stageA", @@ -49,6 +50,7 @@ void roundTripsChangeStateEventWithAuditEntryPayload() { JournalEvent restored = mapper.fromDocument(mapper.toDocument(event)); assertEquals("evt-1", restored.getEventId()); + assertEquals("key-evt-1", restored.getIdempotencyKey()); assertEquals(JournalEventType.CHANGE_STATE, restored.getEventType()); assertEquals(JournalEvent.DEFAULT_VERSION, restored.getEventVersion()); assertEquals("stageA", restored.getStreamId()); @@ -63,6 +65,7 @@ void roundTripsChangeStateEventWithAuditEntryPayload() { void storesEventFieldsUnderTheAgreedBsonNames() { JournalEvent event = new JournalEvent<>( "evt-4", + "key-evt-4", JournalEventType.CHANGE_STATE, "stageA", 3L, @@ -72,6 +75,7 @@ void storesEventFieldsUnderTheAgreedBsonNames() { Document document = mapper.toDocument(event); assertEquals("evt-4", document.getString("eventId")); + assertEquals("key-evt-4", document.getString("idempotencyKey")); assertEquals("stageA", document.getString("streamId")); assertEquals(3L, document.get("streamSequence")); // acknowledged must be persisted as a real boolean false: the partial index filters on it. @@ -83,6 +87,7 @@ void storesEventFieldsUnderTheAgreedBsonNames() { void toDocumentRejectsUnsupportedEventType() { JournalEvent executionEvent = new JournalEvent<>( "evt-2", + "key-evt-2", JournalEventType.EXECUTION_STATE, JournalEvent.DEFAULT_VERSION, "stageA",