Skip to content

Commit 9876cdd

Browse files
authored
refactor: journal event foundation (#942)
1 parent 67b06bf commit 9876cdd

42 files changed

Lines changed: 1966 additions & 772 deletions

File tree

Some content is hidden

Large Commits have some content hidden by default. Use the searchbox below for content that may be hidden.

cloud/flamingock-cloud/src/main/java/io/flamingock/cloud/CloudAuditPersistenceImpl.java

Lines changed: 3 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -16,6 +16,7 @@
1616
package io.flamingock.cloud;
1717

1818
import io.flamingock.internal.common.core.audit.AuditEntry;
19+
import io.flamingock.internal.common.core.audit.AuditWriter;
1920
import io.flamingock.internal.common.core.audit.issue.AuditEntryIssue;
2021
import io.flamingock.internal.common.core.context.ContextContributor;
2122
import io.flamingock.internal.util.Result;
@@ -24,7 +25,6 @@
2425
import io.flamingock.internal.util.id.ServiceId;
2526
import io.flamingock.internal.core.external.store.audit.cloud.CloudAuditPersistence;
2627
import io.flamingock.internal.common.core.context.ContextInjectable;
27-
import io.flamingock.internal.core.external.store.audit.LifecycleAuditWriter;
2828
import io.flamingock.internal.core.plan.ExecutionPlanner;
2929

3030
import java.util.List;
@@ -36,15 +36,15 @@ public final class CloudAuditPersistenceImpl implements CloudAuditPersistence, C
3636

3737
private final ServiceId serviceId;
3838

39-
private final LifecycleAuditWriter auditWriter;
39+
private final AuditWriter auditWriter;
4040

4141
private final ExecutionPlanner executionPlanner;
4242
private final String jwt;
4343

4444
CloudAuditPersistenceImpl(EnvironmentId environmentId,
4545
ServiceId serviceId,
4646
String jwt,
47-
LifecycleAuditWriter auditWriter,
47+
AuditWriter auditWriter,
4848
ExecutionPlanner executionPlanner,
4949
Runnable closer) {
5050
this.environmentId =environmentId;

cloud/flamingock-cloud/src/main/java/io/flamingock/cloud/CloudAuditStoreImpl.java

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -15,6 +15,7 @@
1515
*/
1616
package io.flamingock.cloud;
1717

18+
import io.flamingock.internal.common.core.audit.AuditWriter;
1819
import io.flamingock.internal.util.Constants;
1920
import io.flamingock.internal.util.JsonObjectMapper;
2021
import io.flamingock.internal.util.id.RunnerId;
@@ -35,7 +36,6 @@
3536
import io.flamingock.cloud.planner.CloudExecutionPlanner;
3637
import io.flamingock.cloud.planner.client.ExecutionPlannerClient;
3738
import io.flamingock.cloud.planner.client.HttpExecutionPlannerClient;
38-
import io.flamingock.internal.core.external.store.audit.LifecycleAuditWriter;
3939
import io.flamingock.internal.core.external.targets.TransactionalTargetSystem;
4040
import io.flamingock.internal.core.external.targets.TargetSystemManager;
4141
import io.flamingock.internal.core.external.targets.mark.TargetSystemAuditMarker;
@@ -107,7 +107,7 @@ private CloudAuditPersistenceImpl buildPersistence(RunnerId runnerId,
107107
EnvironmentId environmentId = EnvironmentId.fromLong(authResponse.getEnvironmentId());
108108
ServiceId serviceId = ServiceId.fromLong(authResponse.getServiceId());
109109

110-
LifecycleAuditWriter auditWriter = new HtttpAuditWriter(
110+
AuditWriter auditWriter = new HtttpAuditWriter(
111111
cloudConfiguration.getHost(),
112112
environmentId,
113113
serviceId,

cloud/flamingock-cloud/src/main/java/io/flamingock/cloud/audit/CloudAuditWriter.java

Lines changed: 2 additions & 15 deletions
Original file line numberDiff line numberDiff line change
@@ -15,25 +15,12 @@
1515
*/
1616
package io.flamingock.cloud.audit;
1717

18+
import io.flamingock.internal.common.core.audit.AuditWriter;
1819
import io.flamingock.internal.util.Result;
19-
import io.flamingock.internal.core.external.store.audit.LifecycleAuditWriter;
2020
import io.flamingock.internal.core.external.store.audit.domain.ExecutionAuditContextBundle;
2121
import io.flamingock.internal.core.external.store.audit.domain.RollbackAuditContextBundle;
2222
import io.flamingock.internal.core.external.store.audit.domain.StartExecutionAuditContextBundle;
2323

24-
public interface CloudAuditWriter extends LifecycleAuditWriter {
24+
public interface CloudAuditWriter extends AuditWriter {
2525

26-
default Result writeStartExecution(StartExecutionAuditContextBundle auditContextBundle) {
27-
return Result.OK();//TODO remove this
28-
// return writeEntry(AuditEntryMapper.map(auditContextBundle));
29-
}
30-
31-
32-
default Result writeExecution(ExecutionAuditContextBundle auditContextBundle) {
33-
return writeEntry(auditContextBundle.toAuditEntry());
34-
}
35-
36-
default Result writeRollback(RollbackAuditContextBundle auditContextBundle) {
37-
return writeEntry(auditContextBundle.toAuditEntry());
38-
}
3926
}

community/flamingock-couchbase-auditstore/src/main/java/io/flamingock/store/couchbase/internal/CouchbaseAuditor.java

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -25,10 +25,10 @@
2525
import com.couchbase.client.java.kv.UpsertOptions;
2626
import io.flamingock.internal.common.core.audit.AuditEntry;
2727
import io.flamingock.internal.common.core.audit.AuditReader;
28+
import io.flamingock.internal.common.core.audit.AuditWriter;
2829
import io.flamingock.internal.common.couchbase.CouchbaseAuditMapper;
2930
import io.flamingock.internal.common.couchbase.CouchbaseCollectionHelper;
3031
import io.flamingock.internal.common.couchbase.CouchbaseCollectionInitializator;
31-
import io.flamingock.internal.core.external.store.audit.LifecycleAuditWriter;
3232
import io.flamingock.internal.util.Result;
3333
import io.flamingock.internal.util.log.FlamingockLoggerFactory;
3434
import org.slf4j.Logger;
@@ -37,7 +37,7 @@
3737
import java.util.stream.Collectors;
3838

3939

40-
public class CouchbaseAuditor implements LifecycleAuditWriter, AuditReader {
40+
public class CouchbaseAuditor implements AuditWriter, AuditReader {
4141

4242
private static final Logger logger = FlamingockLoggerFactory.getLogger("CouchbaseAuditor");
4343

community/flamingock-dynamodb-auditstore/src/main/java/io/flamingock/store/dynamodb/internal/DynamoDBAuditor.java

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -15,7 +15,7 @@
1515
*/
1616
package io.flamingock.store.dynamodb.internal;
1717

18-
import io.flamingock.internal.core.external.store.audit.LifecycleAuditWriter;
18+
import io.flamingock.internal.common.core.audit.AuditWriter;
1919
import io.flamingock.internal.core.external.store.audit.community.CommunityAuditReader;
2020
import io.flamingock.internal.util.dynamodb.entities.AuditEntryEntity;
2121
import io.flamingock.internal.common.core.audit.AuditEntry;
@@ -36,7 +36,7 @@
3636

3737
import static java.util.Collections.emptyList;
3838

39-
public class DynamoDBAuditor implements LifecycleAuditWriter, CommunityAuditReader {
39+
public class DynamoDBAuditor implements AuditWriter, CommunityAuditReader {
4040

4141
private static final Logger logger = FlamingockLoggerFactory.getLogger("DynamoAuditor");
4242

community/flamingock-mongodb-reactive-auditstore/src/main/java/io/flamingock/store/mongodb/reactive/internal/MongoDBReactiveAuditor.java

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -25,11 +25,11 @@
2525
import com.mongodb.reactivestreams.client.MongoDatabase;
2626
import io.flamingock.internal.common.core.audit.AuditEntry;
2727
import io.flamingock.internal.common.core.audit.AuditReader;
28+
import io.flamingock.internal.common.core.audit.AuditWriter;
2829
import io.flamingock.internal.common.mongodb.CollectionInitializator;
2930
import io.flamingock.internal.common.mongodb.MongoDBAuditMapper;
3031
import io.flamingock.internal.common.mongodb.MongoDBReactiveCollectionHelper;
3132
import io.flamingock.internal.common.mongodb.MongoDBDocumentHelper;
32-
import io.flamingock.internal.core.external.store.audit.LifecycleAuditWriter;
3333
import io.flamingock.internal.util.Result;
3434
import io.flamingock.internal.util.log.FlamingockLoggerFactory;
3535
import io.flamingock.reactive.util.PublisherSync;
@@ -44,7 +44,7 @@
4444
import static io.flamingock.internal.util.constants.AuditEntryFieldConstants.KEY_EXECUTION_ID;
4545
import static io.flamingock.internal.util.constants.AuditEntryFieldConstants.KEY_STATE;
4646

47-
public class MongoDBReactiveAuditor implements LifecycleAuditWriter, AuditReader {
47+
public class MongoDBReactiveAuditor implements AuditWriter, AuditReader {
4848

4949
private static final Logger logger = FlamingockLoggerFactory.getLogger("MongoDBReactiveAuditor");
5050

community/flamingock-mongodb-sync-auditstore/src/main/java/io/flamingock/store/mongodb/sync/MongoDBSyncAuditStore.java

Lines changed: 20 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -32,6 +32,7 @@
3232
import io.flamingock.store.mongodb.sync.internal.MongoDBSyncLockService;
3333
import io.flamingock.externalsystem.mongodb.api.MongoDBExternalSystem;
3434

35+
import static io.flamingock.internal.common.mongodb.journal.JournalEventPersistenceConstants.DEFAULT_JOURNAL_STORE_NAME;
3536
import static io.flamingock.internal.util.constants.CommunityPersistenceConstants.DEFAULT_AUDIT_STORE_NAME;
3637
import static io.flamingock.internal.util.constants.CommunityPersistenceConstants.DEFAULT_LOCK_STORE_NAME;
3738

@@ -46,6 +47,7 @@ public class MongoDBSyncAuditStore implements CommunityAuditStore {
4647
private MongoDatabase database;
4748
private String auditRepositoryName = DEFAULT_AUDIT_STORE_NAME;
4849
private String lockRepositoryName = DEFAULT_LOCK_STORE_NAME;
50+
private String journalRepositoryName = DEFAULT_JOURNAL_STORE_NAME;
4951
private ReadConcern readConcern = ReadConcern.MAJORITY;
5052
private ReadPreference readPreference = ReadPreference.primary();
5153
private WriteConcern writeConcern = WriteConcern.MAJORITY.withJournal(true);
@@ -85,6 +87,11 @@ public MongoDBSyncAuditStore withLockRepositoryName(String lockRepositoryName) {
8587
return this;
8688
}
8789

90+
public MongoDBSyncAuditStore withJournalRepositoryName(String journalRepositoryName) {
91+
this.journalRepositoryName = journalRepositoryName;
92+
return this;
93+
}
94+
8895
public MongoDBSyncAuditStore withReadConcern(ReadConcern readConcern) {
8996
this.readConcern = readConcern;
9097
return this;
@@ -120,6 +127,7 @@ public synchronized CommunityAuditPersistence getPersistence() {
120127
communityConfiguration,
121128
database,
122129
auditRepositoryName,
130+
journalRepositoryName,
123131
readConcern,
124132
readPreference,
125133
writeConcern,
@@ -161,10 +169,22 @@ private void validate() {
161169
throw new FlamingockException("The 'lockRepositoryName' property is required.");
162170
}
163171

172+
if (journalRepositoryName == null || journalRepositoryName.trim().isEmpty()) {
173+
throw new FlamingockException("The 'journalRepositoryName' property is required.");
174+
}
175+
164176
if (auditRepositoryName.trim().equalsIgnoreCase(lockRepositoryName.trim())) {
165177
throw new FlamingockException("The 'auditRepositoryName' and 'lockRepositoryName' properties must not be the same.");
166178
}
167179

180+
if (journalRepositoryName.trim().equalsIgnoreCase(auditRepositoryName.trim())) {
181+
throw new FlamingockException("The 'journalRepositoryName' and 'auditRepositoryName' properties must not be the same.");
182+
}
183+
184+
if (journalRepositoryName.trim().equalsIgnoreCase(lockRepositoryName.trim())) {
185+
throw new FlamingockException("The 'journalRepositoryName' and 'lockRepositoryName' properties must not be the same.");
186+
}
187+
168188
if (readConcern == null) {
169189
throw new FlamingockException("The 'readConcern' property is required.");
170190
}

community/flamingock-mongodb-sync-auditstore/src/main/java/io/flamingock/store/mongodb/sync/internal/MongoDBSyncAuditPersistence.java

Lines changed: 7 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -34,8 +34,10 @@
3434
public class MongoDBSyncAuditPersistence extends AbstractCommunityAuditPersistence {
3535

3636
private MongoDBSyncAuditor auditor;
37+
private MongoDBSyncJournalEventStore journalEventStore;
3738
private final MongoDatabase database;
3839
private final String auditCollectionName;
40+
private final String journalCollectionName;
3941
private final ReadConcern readConcern;
4042
private final ReadPreference readPreference;
4143
private final WriteConcern writeConcern;
@@ -45,13 +47,15 @@ public class MongoDBSyncAuditPersistence extends AbstractCommunityAuditPersisten
4547
public MongoDBSyncAuditPersistence(CommunityConfigurable localConfiguration,
4648
MongoDatabase database,
4749
String auditCollectionName,
50+
String journalCollectionName,
4851
ReadConcern readConcern,
4952
ReadPreference readPreference,
5053
WriteConcern writeConcern,
5154
boolean autoCreate) {
5255
super(localConfiguration);
5356
this.database = database;
5457
this.auditCollectionName = auditCollectionName;
58+
this.journalCollectionName = journalCollectionName;
5559
this.readConcern = readConcern;
5660
this.readPreference = readPreference;
5761
this.writeConcern = writeConcern;
@@ -63,6 +67,9 @@ protected void doInitialize(RunnerId runnerId) {
6367
//Auditor
6468
auditor = new MongoDBSyncAuditor(database, auditCollectionName, readConcern, readPreference, writeConcern);
6569
auditor.initialize(autoCreate);
70+
//Journal
71+
journalEventStore = new MongoDBSyncJournalEventStore(database, journalCollectionName, readConcern, readPreference, writeConcern);
72+
journalEventStore.initialize(autoCreate);
6673
}
6774

6875
@Deprecated

community/flamingock-mongodb-sync-auditstore/src/main/java/io/flamingock/store/mongodb/sync/internal/MongoDBSyncAuditor.java

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -25,11 +25,11 @@
2525
import com.mongodb.client.result.UpdateResult;
2626
import io.flamingock.internal.common.core.audit.AuditEntry;
2727
import io.flamingock.internal.common.core.audit.AuditReader;
28+
import io.flamingock.internal.common.core.audit.AuditWriter;
2829
import io.flamingock.internal.common.mongodb.CollectionInitializator;
2930
import io.flamingock.internal.common.mongodb.MongoDBAuditMapper;
3031
import io.flamingock.internal.common.mongodb.MongoDBSyncCollectionHelper;
3132
import io.flamingock.internal.common.mongodb.MongoDBDocumentHelper;
32-
import io.flamingock.internal.core.external.store.audit.LifecycleAuditWriter;
3333
import io.flamingock.internal.util.Result;
3434
import io.flamingock.internal.util.log.FlamingockLoggerFactory;
3535
import org.bson.Document;
@@ -44,7 +44,7 @@
4444
import static io.flamingock.internal.util.constants.AuditEntryFieldConstants.KEY_EXECUTION_ID;
4545
import static io.flamingock.internal.util.constants.AuditEntryFieldConstants.KEY_STATE;
4646

47-
public class MongoDBSyncAuditor implements LifecycleAuditWriter, AuditReader {
47+
public class MongoDBSyncAuditor implements AuditWriter, AuditReader {
4848

4949
private static final Logger logger = FlamingockLoggerFactory.getLogger("MongoDBSyncAuditor");
5050

0 commit comments

Comments
 (0)