Skip to content

Commit 33e73d7

Browse files
authored
refactor: Journal events foundation. MongoDB sync impl. Feature flags (#944)
1 parent 9876cdd commit 33e73d7

48 files changed

Lines changed: 1785 additions & 264 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.

build.gradle.kts

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -22,7 +22,7 @@ allprojects {
2222
val declaredVersion = "1.5.0-SNAPSHOT"
2323
version = VersionManager.resolveVersion(declaredVersion, project.hasProperty("release"))
2424

25-
extra["generalUtilVersion"] = "1.5.3"
25+
extra["generalUtilVersion"] = "1.6.0"
2626
extra["templateApiVersion"] = "1.3.4"
2727
extra["coreApiVersion"] = "1.3.3"
2828
extra["sqlVersion"] = "1.3.2"

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

Lines changed: 0 additions & 13 deletions
Original file line numberDiff line numberDiff line change
@@ -54,20 +54,7 @@ public final class CloudAuditPersistenceImpl implements CloudAuditPersistence, C
5454
this.executionPlanner = executionPlanner;
5555
}
5656

57-
@Override
58-
public EnvironmentId getEnvironmentId() {
59-
return environmentId;
60-
}
6157

62-
@Override
63-
public ServiceId getServiceId() {
64-
return serviceId;
65-
}
66-
67-
@Override
68-
public String getJwt() {
69-
return jwt;
70-
}
7158

7259
@Override
7360
public ExecutionPlanner getExecutionPlanner() {

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

Lines changed: 10 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -18,6 +18,7 @@
1818
import com.mongodb.ReadConcern;
1919
import com.mongodb.ReadPreference;
2020
import com.mongodb.WriteConcern;
21+
import com.mongodb.reactivestreams.client.ClientSession;
2122
import com.mongodb.reactivestreams.client.MongoDatabase;
2223
import io.flamingock.externalsystem.mongodb.reactive.api.MongoDBReactiveExternalSystem;
2324
import io.flamingock.internal.common.core.context.ContextResolver;
@@ -32,6 +33,10 @@
3233
import io.flamingock.store.mongodb.reactive.internal.MongoDBReactiveAuditPersistence;
3334
import io.flamingock.store.mongodb.reactive.internal.MongoDBReactiveLockService;
3435

36+
import java.util.Collections;
37+
import java.util.HashSet;
38+
import java.util.Set;
39+
3540
import static io.flamingock.internal.util.constants.CommunityPersistenceConstants.DEFAULT_AUDIT_STORE_NAME;
3641
import static io.flamingock.internal.util.constants.CommunityPersistenceConstants.DEFAULT_LOCK_STORE_NAME;
3742

@@ -176,4 +181,9 @@ private void validate() {
176181
throw new FlamingockException("The 'writeConcern' property is required.");
177182
}
178183
}
184+
185+
@Override
186+
public Set<Class<?>> getNonGuardedTypes() {
187+
return new HashSet<>(Collections.singletonList(ClientSession.class));
188+
}
179189
}

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

Lines changed: 0 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -65,11 +65,6 @@ protected void doInitialize(RunnerId runnerId) {
6565
auditor.initialize(autoCreate);
6666
}
6767

68-
@Deprecated
69-
@Override
70-
public Set<Class<?>> getNonGuardedTypes() {
71-
return new HashSet<>(Collections.singletonList(ClientSession.class));
72-
}
7368

7469
@Override
7570
public List<AuditEntry> getAuditHistory() {

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

Lines changed: 54 additions & 22 deletions
Original file line numberDiff line numberDiff line change
@@ -18,20 +18,33 @@
1818
import com.mongodb.ReadConcern;
1919
import com.mongodb.ReadPreference;
2020
import com.mongodb.WriteConcern;
21+
import com.mongodb.client.ClientSession;
2122
import com.mongodb.client.MongoDatabase;
23+
import io.flamingock.internal.common.core.audit.AuditEntry;
24+
import io.flamingock.internal.common.core.audit.AuditPersistenceFactory;
25+
import io.flamingock.internal.common.core.audit.AuditReader;
2226
import io.flamingock.internal.common.core.context.ContextResolver;
2327
import io.flamingock.internal.common.core.error.FlamingockException;
2428
import io.flamingock.internal.core.configuration.community.CommunityConfigurable;
2529
import io.flamingock.internal.core.external.store.CommunityAuditStore;
2630
import io.flamingock.internal.core.external.store.audit.community.CommunityAuditPersistence;
2731
import io.flamingock.internal.core.external.store.lock.community.CommunityLockService;
32+
import io.flamingock.internal.core.journal.JournalEventSequencer;
33+
import io.flamingock.internal.core.journal.JournalEventSequencerFactory;
2834
import io.flamingock.internal.util.Constants;
2935
import io.flamingock.internal.util.TimeService;
3036
import io.flamingock.internal.util.id.RunnerId;
3137
import io.flamingock.store.mongodb.sync.internal.MongoDBSyncAuditPersistence;
38+
import io.flamingock.store.mongodb.sync.internal.MongoDBSyncAuditRepository;
39+
import io.flamingock.store.mongodb.sync.internal.MongoDBSyncJournalEventStore;
3240
import io.flamingock.store.mongodb.sync.internal.MongoDBSyncLockService;
3341
import io.flamingock.externalsystem.mongodb.api.MongoDBExternalSystem;
3442

43+
import java.util.Collections;
44+
import java.util.HashSet;
45+
import java.util.List;
46+
import java.util.Set;
47+
3548
import static io.flamingock.internal.common.mongodb.journal.JournalEventPersistenceConstants.DEFAULT_JOURNAL_STORE_NAME;
3649
import static io.flamingock.internal.util.constants.CommunityPersistenceConstants.DEFAULT_AUDIT_STORE_NAME;
3750
import static io.flamingock.internal.util.constants.CommunityPersistenceConstants.DEFAULT_LOCK_STORE_NAME;
@@ -52,6 +65,9 @@ public class MongoDBSyncAuditStore implements CommunityAuditStore {
5265
private ReadPreference readPreference = ReadPreference.primary();
5366
private WriteConcern writeConcern = WriteConcern.MAJORITY.withJournal(true);
5467
private boolean autoCreate = true;
68+
private MongoDBSyncAuditRepository auditRepository;
69+
private MongoDBSyncJournalEventStore journalEventStore;
70+
private JournalEventSequencerFactory journalEventSequencerFactory;
5571

5672

5773
private MongoDBSyncAuditStore(MongoDBExternalSystem mongoDBTargetSystem) {
@@ -117,41 +133,52 @@ public void initialize(ContextResolver baseContext) {
117133
runnerId = baseContext.getRequiredDependencyValue(RunnerId.class);
118134
communityConfiguration = baseContext.getRequiredDependencyValue(CommunityConfigurable.class);
119135
database = mongoDBTargetSystem.getMongoDatabase();
136+
auditRepository = new MongoDBSyncAuditRepository(database, auditRepositoryName, readConcern, readPreference, writeConcern);
137+
journalEventStore = new MongoDBSyncJournalEventStore(database, journalRepositoryName, readConcern, readPreference, writeConcern);
138+
journalEventSequencerFactory = new JournalEventSequencerFactory(journalEventStore);
139+
140+
lockService = new MongoDBSyncLockService(
141+
database,
142+
lockRepositoryName,
143+
readConcern,
144+
readPreference,
145+
writeConcern,
146+
TimeService.getDefault()
147+
);
148+
lockService.initialize(autoCreate);
120149
this.validate();
121150
}
122151

123152
@Override
124-
public synchronized CommunityAuditPersistence getPersistence() {
125-
if (persistence == null) {
153+
public AuditPersistenceFactory<CommunityAuditPersistence> getPersistenceFactory() {
154+
return stageId -> {
155+
JournalEventSequencer journalEventSequencer = journalEventSequencerFactory.forStream(stageId);
126156
persistence = new MongoDBSyncAuditPersistence(
127157
communityConfiguration,
128-
database,
129-
auditRepositoryName,
130-
journalRepositoryName,
131-
readConcern,
132-
readPreference,
133-
writeConcern,
158+
auditRepository,
159+
journalEventStore,
160+
journalEventSequencer,
161+
mongoDBTargetSystem.getTxWrapper(),
134162
autoCreate
135163
);
136164
persistence.initialize(runnerId);
137-
}
138-
return persistence;
165+
return persistence;
166+
};
139167
}
140168

141169
@Override
142-
public synchronized CommunityLockService getLockService() {
143-
if (lockService == null) {
144-
lockService = new MongoDBSyncLockService(
145-
database,
146-
lockRepositoryName,
147-
readConcern,
148-
readPreference,
149-
writeConcern,
150-
TimeService.getDefault()
151-
);
152-
lockService.initialize(autoCreate);
170+
public AuditReader getAuditReader() {
171+
return () -> auditRepository.getAuditHistory();
172+
}
153173

154-
}
174+
175+
@Override
176+
public CommunityAuditPersistence getPersistence() {
177+
throw new RuntimeException("getPersistence shouldn´t be called at MongodbSync ");
178+
}
179+
180+
@Override
181+
public synchronized CommunityLockService getLockService() {
155182
return lockService;
156183
}
157184

@@ -197,4 +224,9 @@ private void validate() {
197224
throw new FlamingockException("The 'writeConcern' property is required.");
198225
}
199226
}
227+
228+
@Override
229+
public Set<Class<?>> getNonGuardedTypes() {
230+
return new HashSet<>(Collections.singletonList(ClientSession.class));
231+
}
200232
}

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

Lines changed: 45 additions & 41 deletions
Original file line numberDiff line numberDiff line change
@@ -15,77 +15,81 @@
1515
*/
1616
package io.flamingock.store.mongodb.sync.internal;
1717

18-
import com.mongodb.ReadConcern;
19-
import com.mongodb.ReadPreference;
20-
import com.mongodb.WriteConcern;
2118
import com.mongodb.client.ClientSession;
22-
import com.mongodb.client.MongoDatabase;
2319
import io.flamingock.internal.common.core.audit.AuditEntry;
20+
import io.flamingock.internal.common.core.context.RuntimeContext;
21+
import io.flamingock.internal.common.core.feature.Features;
22+
import io.flamingock.internal.common.core.journal.JournalEvent;
23+
import io.flamingock.internal.common.core.transaction.TransactionWrapper;
2424
import io.flamingock.internal.core.configuration.community.CommunityConfigurable;
25+
import io.flamingock.internal.core.context.BasicRuntimeContext;
2526
import io.flamingock.internal.core.external.store.audit.community.AbstractCommunityAuditPersistence;
27+
import io.flamingock.internal.core.journal.JournalEventSequencer;
28+
import io.flamingock.internal.util.FeatureFlag;
2629
import io.flamingock.internal.util.Result;
2730
import io.flamingock.internal.util.id.RunnerId;
2831

29-
import java.util.Collections;
30-
import java.util.HashSet;
3132
import java.util.List;
32-
import java.util.Set;
3333

3434
public class MongoDBSyncAuditPersistence extends AbstractCommunityAuditPersistence {
3535

36-
private MongoDBSyncAuditor auditor;
37-
private MongoDBSyncJournalEventStore journalEventStore;
38-
private final MongoDatabase database;
39-
private final String auditCollectionName;
40-
private final String journalCollectionName;
41-
private final ReadConcern readConcern;
42-
private final ReadPreference readPreference;
43-
private final WriteConcern writeConcern;
36+
private final MongoDBSyncAuditRepository auditRepository;
37+
private final MongoDBSyncJournalEventStore journalEventStore;
38+
private final JournalEventSequencer journalEventSequencer;
39+
private final TransactionWrapper txWrapper;
4440
private final boolean autoCreate;
4541

4642

4743
public MongoDBSyncAuditPersistence(CommunityConfigurable localConfiguration,
48-
MongoDatabase database,
49-
String auditCollectionName,
50-
String journalCollectionName,
51-
ReadConcern readConcern,
52-
ReadPreference readPreference,
53-
WriteConcern writeConcern,
54-
boolean autoCreate) {
44+
MongoDBSyncAuditRepository auditRepository,
45+
MongoDBSyncJournalEventStore journalEventStore,
46+
JournalEventSequencer journalEventSequencer,
47+
TransactionWrapper txWrapper,
48+
boolean autoCreate) {
5549
super(localConfiguration);
56-
this.database = database;
57-
this.auditCollectionName = auditCollectionName;
58-
this.journalCollectionName = journalCollectionName;
59-
this.readConcern = readConcern;
60-
this.readPreference = readPreference;
61-
this.writeConcern = writeConcern;
50+
this.auditRepository = auditRepository;
51+
this.journalEventStore = journalEventStore;
52+
this.journalEventSequencer = journalEventSequencer;
53+
this.txWrapper = txWrapper;
6254
this.autoCreate = autoCreate;
6355
}
6456

6557
@Override
6658
protected void doInitialize(RunnerId runnerId) {
67-
//Auditor
68-
auditor = new MongoDBSyncAuditor(database, auditCollectionName, readConcern, readPreference, writeConcern);
69-
auditor.initialize(autoCreate);
70-
//Journal
71-
journalEventStore = new MongoDBSyncJournalEventStore(database, journalCollectionName, readConcern, readPreference, writeConcern);
72-
journalEventStore.initialize(autoCreate);
59+
auditRepository.initialize(autoCreate);
60+
// Creating the indexes is what brings the journal collection into existence — there is no explicit
61+
// createCollection call — so skipping this keeps it from ever appearing. It must stay in step with the
62+
// append in writeEntry: skipping setup while still appending would let insertOne create the collection
63+
// implicitly and without indexes, voiding the unique (streamId, streamSequence) and eventId guarantees.
64+
FeatureFlag.ifEnabled(Features.JOURNAL_EVENTS, () -> journalEventStore.initialize(autoCreate));
7365
}
7466

75-
@Deprecated
76-
@Override
77-
public Set<Class<?>> getNonGuardedTypes() {
78-
return new HashSet<>(Collections.singletonList(ClientSession.class));
79-
}
8067

8168
@Override
8269
public List<AuditEntry> getAuditHistory() {
83-
return auditor.getAuditHistory();
70+
return auditRepository.getAuditHistory();
8471
}
8572

8673
@Override
8774
public Result writeEntry(AuditEntry auditEntry) {
88-
return auditor.writeEntry(auditEntry);
75+
RuntimeContext baseContext = new BasicRuntimeContext("write-changeState-" + auditEntry.getChangeId());
76+
if (FeatureFlag.isEnabled(Features.JOURNAL_EVENTS)) {
77+
return txWrapper.wrapInTransaction(baseContext, runtimeContext -> {
78+
ClientSession clientSession = runtimeContext.getContext().getRequiredDependencyValue(ClientSession.class);
79+
// Read once rather than per branch: the journal append and the audit write shape are two halves of
80+
// one model. With events, the audit record is the change's current state and the journal is the
81+
// history; without them, the audit record set is itself the history.
82+
JournalEvent<AuditEntry> journalEvent = journalEventSequencer.newEvent(auditEntry);
83+
journalEventStore.write(clientSession, journalEvent);
84+
return auditRepository.save(clientSession, auditEntry);
85+
86+
});
87+
} else {
88+
return auditRepository.append(auditEntry);
89+
}
90+
91+
92+
8993
}
9094

9195
}

0 commit comments

Comments
 (0)