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
2 changes: 1 addition & 1 deletion build.gradle.kts
Original file line number Diff line number Diff line change
Expand Up @@ -22,7 +22,7 @@ allprojects {
val declaredVersion = "1.5.0-SNAPSHOT"
version = VersionManager.resolveVersion(declaredVersion, project.hasProperty("release"))

extra["generalUtilVersion"] = "1.5.3"
extra["generalUtilVersion"] = "1.6.0"
extra["templateApiVersion"] = "1.3.4"
extra["coreApiVersion"] = "1.3.3"
extra["sqlVersion"] = "1.3.2"
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -54,20 +54,7 @@ public final class CloudAuditPersistenceImpl implements CloudAuditPersistence, C
this.executionPlanner = executionPlanner;
}

@Override
public EnvironmentId getEnvironmentId() {
return environmentId;
}

@Override
public ServiceId getServiceId() {
return serviceId;
}

@Override
public String getJwt() {
return jwt;
}

@Override
public ExecutionPlanner getExecutionPlanner() {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -18,6 +18,7 @@
import com.mongodb.ReadConcern;
import com.mongodb.ReadPreference;
import com.mongodb.WriteConcern;
import com.mongodb.reactivestreams.client.ClientSession;
import com.mongodb.reactivestreams.client.MongoDatabase;
import io.flamingock.externalsystem.mongodb.reactive.api.MongoDBReactiveExternalSystem;
import io.flamingock.internal.common.core.context.ContextResolver;
Expand All @@ -32,6 +33,10 @@
import io.flamingock.store.mongodb.reactive.internal.MongoDBReactiveAuditPersistence;
import io.flamingock.store.mongodb.reactive.internal.MongoDBReactiveLockService;

import java.util.Collections;
import java.util.HashSet;
import java.util.Set;

import static io.flamingock.internal.util.constants.CommunityPersistenceConstants.DEFAULT_AUDIT_STORE_NAME;
import static io.flamingock.internal.util.constants.CommunityPersistenceConstants.DEFAULT_LOCK_STORE_NAME;

Expand Down Expand Up @@ -176,4 +181,9 @@ private void validate() {
throw new FlamingockException("The 'writeConcern' property is required.");
}
}

@Override
public Set<Class<?>> getNonGuardedTypes() {
return new HashSet<>(Collections.singletonList(ClientSession.class));
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -65,11 +65,6 @@ protected void doInitialize(RunnerId runnerId) {
auditor.initialize(autoCreate);
}

@Deprecated
@Override
public Set<Class<?>> getNonGuardedTypes() {
return new HashSet<>(Collections.singletonList(ClientSession.class));
}

@Override
public List<AuditEntry> getAuditHistory() {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -18,20 +18,33 @@
import com.mongodb.ReadConcern;
import com.mongodb.ReadPreference;
import com.mongodb.WriteConcern;
import com.mongodb.client.ClientSession;
import com.mongodb.client.MongoDatabase;
import io.flamingock.internal.common.core.audit.AuditEntry;
import io.flamingock.internal.common.core.audit.AuditPersistenceFactory;
import io.flamingock.internal.common.core.audit.AuditReader;
import io.flamingock.internal.common.core.context.ContextResolver;
import io.flamingock.internal.common.core.error.FlamingockException;
import io.flamingock.internal.core.configuration.community.CommunityConfigurable;
import io.flamingock.internal.core.external.store.CommunityAuditStore;
import io.flamingock.internal.core.external.store.audit.community.CommunityAuditPersistence;
import io.flamingock.internal.core.external.store.lock.community.CommunityLockService;
import io.flamingock.internal.core.journal.JournalEventSequencer;
import io.flamingock.internal.core.journal.JournalEventSequencerFactory;
import io.flamingock.internal.util.Constants;
import io.flamingock.internal.util.TimeService;
import io.flamingock.internal.util.id.RunnerId;
import io.flamingock.store.mongodb.sync.internal.MongoDBSyncAuditPersistence;
import io.flamingock.store.mongodb.sync.internal.MongoDBSyncAuditRepository;
import io.flamingock.store.mongodb.sync.internal.MongoDBSyncJournalEventStore;
import io.flamingock.store.mongodb.sync.internal.MongoDBSyncLockService;
import io.flamingock.externalsystem.mongodb.api.MongoDBExternalSystem;

import java.util.Collections;
import java.util.HashSet;
import java.util.List;
import java.util.Set;

import static io.flamingock.internal.common.mongodb.journal.JournalEventPersistenceConstants.DEFAULT_JOURNAL_STORE_NAME;
import static io.flamingock.internal.util.constants.CommunityPersistenceConstants.DEFAULT_AUDIT_STORE_NAME;
import static io.flamingock.internal.util.constants.CommunityPersistenceConstants.DEFAULT_LOCK_STORE_NAME;
Expand All @@ -52,6 +65,9 @@ public class MongoDBSyncAuditStore implements CommunityAuditStore {
private ReadPreference readPreference = ReadPreference.primary();
private WriteConcern writeConcern = WriteConcern.MAJORITY.withJournal(true);
private boolean autoCreate = true;
private MongoDBSyncAuditRepository auditRepository;
private MongoDBSyncJournalEventStore journalEventStore;
private JournalEventSequencerFactory journalEventSequencerFactory;


private MongoDBSyncAuditStore(MongoDBExternalSystem mongoDBTargetSystem) {
Expand Down Expand Up @@ -117,41 +133,52 @@ public void initialize(ContextResolver baseContext) {
runnerId = baseContext.getRequiredDependencyValue(RunnerId.class);
communityConfiguration = baseContext.getRequiredDependencyValue(CommunityConfigurable.class);
database = mongoDBTargetSystem.getMongoDatabase();
auditRepository = new MongoDBSyncAuditRepository(database, auditRepositoryName, readConcern, readPreference, writeConcern);
journalEventStore = new MongoDBSyncJournalEventStore(database, journalRepositoryName, readConcern, readPreference, writeConcern);
journalEventSequencerFactory = new JournalEventSequencerFactory(journalEventStore);

lockService = new MongoDBSyncLockService(
database,
lockRepositoryName,
readConcern,
readPreference,
writeConcern,
TimeService.getDefault()
);
lockService.initialize(autoCreate);
this.validate();
}

@Override
public synchronized CommunityAuditPersistence getPersistence() {
if (persistence == null) {
public AuditPersistenceFactory<CommunityAuditPersistence> getPersistenceFactory() {
return stageId -> {
JournalEventSequencer journalEventSequencer = journalEventSequencerFactory.forStream(stageId);
persistence = new MongoDBSyncAuditPersistence(
communityConfiguration,
database,
auditRepositoryName,
journalRepositoryName,
readConcern,
readPreference,
writeConcern,
auditRepository,
journalEventStore,
journalEventSequencer,
mongoDBTargetSystem.getTxWrapper(),
autoCreate
);
persistence.initialize(runnerId);
}
return persistence;
return persistence;
};
}

@Override
public synchronized CommunityLockService getLockService() {
if (lockService == null) {
lockService = new MongoDBSyncLockService(
database,
lockRepositoryName,
readConcern,
readPreference,
writeConcern,
TimeService.getDefault()
);
lockService.initialize(autoCreate);
public AuditReader getAuditReader() {
return () -> auditRepository.getAuditHistory();
}

}

@Override
public CommunityAuditPersistence getPersistence() {
throw new RuntimeException("getPersistence shouldn´t be called at MongodbSync ");
}
Comment thread
dieppa marked this conversation as resolved.

@Override
public synchronized CommunityLockService getLockService() {
return lockService;
}

Expand Down Expand Up @@ -197,4 +224,9 @@ private void validate() {
throw new FlamingockException("The 'writeConcern' property is required.");
}
}

@Override
public Set<Class<?>> getNonGuardedTypes() {
return new HashSet<>(Collections.singletonList(ClientSession.class));
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -15,77 +15,81 @@
*/
package io.flamingock.store.mongodb.sync.internal;

import com.mongodb.ReadConcern;
import com.mongodb.ReadPreference;
import com.mongodb.WriteConcern;
import com.mongodb.client.ClientSession;
import com.mongodb.client.MongoDatabase;
import io.flamingock.internal.common.core.audit.AuditEntry;
import io.flamingock.internal.common.core.context.RuntimeContext;
import io.flamingock.internal.common.core.feature.Features;
import io.flamingock.internal.common.core.journal.JournalEvent;
import io.flamingock.internal.common.core.transaction.TransactionWrapper;
import io.flamingock.internal.core.configuration.community.CommunityConfigurable;
import io.flamingock.internal.core.context.BasicRuntimeContext;
import io.flamingock.internal.core.external.store.audit.community.AbstractCommunityAuditPersistence;
import io.flamingock.internal.core.journal.JournalEventSequencer;
import io.flamingock.internal.util.FeatureFlag;
import io.flamingock.internal.util.Result;
import io.flamingock.internal.util.id.RunnerId;

import java.util.Collections;
import java.util.HashSet;
import java.util.List;
import java.util.Set;

public class MongoDBSyncAuditPersistence extends AbstractCommunityAuditPersistence {

private MongoDBSyncAuditor auditor;
private MongoDBSyncJournalEventStore journalEventStore;
private final MongoDatabase database;
private final String auditCollectionName;
private final String journalCollectionName;
private final ReadConcern readConcern;
private final ReadPreference readPreference;
private final WriteConcern writeConcern;
private final MongoDBSyncAuditRepository auditRepository;
private final MongoDBSyncJournalEventStore journalEventStore;
private final JournalEventSequencer journalEventSequencer;
private final TransactionWrapper txWrapper;
private final boolean autoCreate;


public MongoDBSyncAuditPersistence(CommunityConfigurable localConfiguration,
MongoDatabase database,
String auditCollectionName,
String journalCollectionName,
ReadConcern readConcern,
ReadPreference readPreference,
WriteConcern writeConcern,
boolean autoCreate) {
MongoDBSyncAuditRepository auditRepository,
MongoDBSyncJournalEventStore journalEventStore,
JournalEventSequencer journalEventSequencer,
TransactionWrapper txWrapper,
boolean autoCreate) {
super(localConfiguration);
this.database = database;
this.auditCollectionName = auditCollectionName;
this.journalCollectionName = journalCollectionName;
this.readConcern = readConcern;
this.readPreference = readPreference;
this.writeConcern = writeConcern;
this.auditRepository = auditRepository;
this.journalEventStore = journalEventStore;
this.journalEventSequencer = journalEventSequencer;
this.txWrapper = txWrapper;
this.autoCreate = autoCreate;
}

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

@Deprecated
@Override
public Set<Class<?>> getNonGuardedTypes() {
return new HashSet<>(Collections.singletonList(ClientSession.class));
}

@Override
public List<AuditEntry> getAuditHistory() {
return auditor.getAuditHistory();
return auditRepository.getAuditHistory();
}

@Override
public Result writeEntry(AuditEntry auditEntry) {
return auditor.writeEntry(auditEntry);
RuntimeContext baseContext = new BasicRuntimeContext("write-changeState-" + auditEntry.getChangeId());
if (FeatureFlag.isEnabled(Features.JOURNAL_EVENTS)) {
return txWrapper.wrapInTransaction(baseContext, runtimeContext -> {
ClientSession clientSession = runtimeContext.getContext().getRequiredDependencyValue(ClientSession.class);
// Read once rather than per branch: the journal append and the audit write shape are two halves of
// one model. With events, the audit record is the change's current state and the journal is the
// history; without them, the audit record set is itself the history.
JournalEvent<AuditEntry> journalEvent = journalEventSequencer.newEvent(auditEntry);
journalEventStore.write(clientSession, journalEvent);
return auditRepository.save(clientSession, auditEntry);

});
} else {
return auditRepository.append(auditEntry);
}



}

}
Loading
Loading