Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -263,6 +263,7 @@ private JournalEventSequencer newSequencer() {
private void occupyStreamPosition(long sequence) {
JournalEvent<AuditEntry> squatter = new JournalEvent<>(
"pre-existing-event",
"key-pre-existing-event",
JournalEventType.CHANGE_STATE,
JournalEvent.DEFAULT_VERSION,
STREAM_ID,
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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());

Expand Down Expand Up @@ -350,6 +351,7 @@ void mapperOmitsBothPendingKeysWhenAcknowledged() {
JournalEvent<AuditEntry> source = journalEvent("stage-1", 1L, "event-1");
JournalEvent<AuditEntry> acknowledged = new JournalEvent<>(
source.getEventId(),
source.getIdempotencyKey(),
source.getEventType(),
source.getEventVersion(),
source.getStreamId(),
Expand Down Expand Up @@ -535,7 +537,7 @@ private JournalEvent<AuditEntry> 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<JournalEvent<AuditEntry>> awaitUnacknowledgedCount(int expected) throws InterruptedException {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -316,7 +316,7 @@ private MongoDBReactiveAuditPersistence persistenceWithoutTransactions(

private void occupyStreamPosition(long streamSequence) {
JournalEvent<AuditEntry> 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)));
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -331,6 +331,7 @@ private static JournalEvent<AuditEntry> fixedEvent(String eventId,
boolean acknowledged) {
return new JournalEvent<>(
eventId,
"key-" + eventId,
JournalEventType.CHANGE_STATE,
3,
streamId,
Expand All @@ -353,6 +354,7 @@ private static JournalEvent<AuditEntry> event(String eventId,
boolean acknowledged) {
return new JournalEvent<>(
eventId,
"key-" + eventId,
JournalEventType.CHANGE_STATE,
JournalEvent.DEFAULT_VERSION,
streamId,
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -344,7 +344,7 @@ private MongoDBSyncAuditPersistence persistenceWithoutTransactions(
*/
private void occupyStreamPosition(long streamSequence) {
JournalEvent<AuditEntry> 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));
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -275,6 +275,7 @@ private static List<String> ids(List<JournalEvent<AuditEntry>> events) {
private static JournalEvent<AuditEntry> event(String eventId, String streamId, long sequence, boolean acknowledged) {
return new JournalEvent<>(
eventId,
"key-" + eventId,
JournalEventType.CHANGE_STATE,
JournalEvent.DEFAULT_VERSION,
streamId,
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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";
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -67,6 +67,7 @@ List<String> getIndexNames(String tableName) {
List<ColumnDefinition> 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),
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -42,13 +42,14 @@ final class SqlJournalEventMapper {
void bind(PreparedStatement statement, JournalEvent<AuditEntry> 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<AuditEntry> fromResultSet(ResultSet resultSet) throws SQLException {
Expand All @@ -67,6 +68,7 @@ JournalEvent<AuditEntry> 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),
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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());
Expand All @@ -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());
Expand Down Expand Up @@ -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();
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -662,6 +662,7 @@ private static JournalEvent<AuditEntry> event(String streamId,
String changeId) {
return new JournalEvent<>(
eventId,
"key-" + eventId,
JournalEventType.CHANGE_STATE,
JournalEvent.DEFAULT_VERSION,
streamId,
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -45,6 +45,7 @@ void generatesTypedSchemaAndIndexesForEveryDialect(SqlDialect dialect) {
List<String> 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"));
Expand All @@ -65,7 +66,7 @@ void generatesTypedSchemaAndIndexesForEveryDialect(SqlDialect dialect) {
.allMatch(name -> name.length() <= helper.getMaximumIndexNameLength()));

List<String> 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)));
}
Expand Down Expand Up @@ -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"));
Expand Down Expand Up @@ -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<String> columnNames(List<SqlJournalDialectHelper.ColumnDefinition> definitions) {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -51,7 +51,7 @@ void roundTripsTypedChangeStateEvent() throws Exception {
AuditEntry auditEntry = auditEntry();
Instant occurredAt = Instant.parse("2026-08-11T10:20:30.123456Z");
JournalEvent<AuditEntry> 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")) {
Expand Down Expand Up @@ -80,6 +80,7 @@ void roundTripsTypedChangeStateEvent() throws Exception {
JournalEvent<AuditEntry> 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());
Expand All @@ -96,7 +97,7 @@ void roundTripsTypedChangeStateEvent() throws Exception {
@DisplayName("rejects event types whose payload mapping is not implemented")
void rejectsUnsupportedEventType() throws Exception {
JournalEvent<AuditEntry> 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")) {
Expand All @@ -114,7 +115,7 @@ void roundTripsAcknowledgedEventAndNullableFields() throws Exception {
null, null, null, 0L, null, null, false, null, null,
null, null, null, null);
JournalEvent<AuditEntry> 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);

Expand Down Expand Up @@ -148,7 +149,7 @@ void roundTripsAcknowledgedEventAndNullableFields() throws Exception {
void keepsEnvelopeAndPayloadTimesDistinct() throws Exception {
AuditEntry auditEntry = auditEntry();
JournalEvent<AuditEntry> 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")) {
Expand All @@ -175,11 +176,11 @@ void keepsEnvelopeAndPayloadTimesDistinct() throws Exception {
@DisplayName("enforces stream position uniqueness without requiring globally unique event IDs")
void enforcesCompositeStreamPositionAndAllowsDuplicateEventIds() throws Exception {
JournalEvent<AuditEntry> 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<AuditEntry> 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<AuditEntry> 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);
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -495,6 +495,7 @@ private static JournalEvent<AuditEntry> event(String streamId,
boolean acknowledged) {
return new JournalEvent<>(
eventId,
"key-" + eventId,
JournalEventType.CHANGE_STATE,
JournalEvent.DEFAULT_VERSION,
streamId,
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -32,6 +32,7 @@ public final class JournalEvent<T> {
public static final int DEFAULT_VERSION = 1;

private final String eventId;
private final String idempotencyKey;
private final JournalEventType eventType;
private final int eventVersion;

Expand All @@ -44,15 +45,17 @@ public final class JournalEvent<T> {
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,
Expand All @@ -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");

Expand All @@ -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;
}
Expand Down
Loading
Loading