Skip to content
Open
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 @@ -15,50 +15,63 @@
*/
package io.flamingock.store.dynamodb;

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.constants.CommunityPersistenceConstants;
import io.flamingock.internal.util.dynamodb.entities.journal.JournalEventFieldConstants;
import io.flamingock.internal.util.id.RunnerId;
import io.flamingock.store.dynamodb.internal.DynamoDBAuditPersistence;
import io.flamingock.store.dynamodb.internal.DynamoDBAuditRepository;
import io.flamingock.store.dynamodb.internal.DynamoDBJournalEventStore;
import io.flamingock.store.dynamodb.internal.DynamoDBLockService;
import io.flamingock.externalsystem.dynamodb.api.DynamoDBExternalSystem;
import software.amazon.awssdk.services.dynamodb.DynamoDbClient;

public class DynamoDBAuditStore implements CommunityAuditStore {

private final DynamoDbClient client;
private final DynamoDBExternalSystem targetSystem;

private RunnerId runnerId;
private CommunityConfigurable communityConfiguration;
private DynamoDBAuditPersistence persistence;
private DynamoDBLockService lockService;
private final DynamoDbClient client;
private String auditRepositoryName = CommunityPersistenceConstants.DEFAULT_AUDIT_STORE_NAME;
private String lockRepositoryName = CommunityPersistenceConstants.DEFAULT_LOCK_STORE_NAME;
private String journalRepositoryName = JournalEventFieldConstants.DEFAULT_JOURNAL_REPOSITORY_NAME;
private long readCapacityUnits = 5L;
private long writeCapacityUnits = 5L;
private boolean autoCreate = true;
private DynamoDBAuditRepository auditRepository;
private DynamoDBJournalEventStore journalEventStore;
private JournalEventSequencerFactory journalEventSequencerFactory;

private DynamoDBAuditStore(DynamoDbClient client) {
this.client = client;
private DynamoDBAuditStore(DynamoDBExternalSystem targetSystem) {
this.targetSystem = targetSystem;
this.client = targetSystem.getClient();
}

/**
* Creates a {@link DynamoDBAuditStore} using the same DynamoDB client
* configured in the given {@link DynamoDBExternalSystem}.
* <p>
* Only the underlying DynamoDB instance (client) is reused.
* No additional target-system configuration is carried over.
* The DynamoDB client and transaction wrapper are reused from the target system.
*
* @param targetSystem the target system from which to derive the client
* @return a new audit store bound to the same DynamoDB instance as the target system
*/
public static DynamoDBAuditStore from(DynamoDBExternalSystem targetSystem) {
return new DynamoDBAuditStore(targetSystem.getClient());
return new DynamoDBAuditStore(targetSystem);
}

@Override
Expand All @@ -76,6 +89,11 @@ public DynamoDBAuditStore withLockRepositoryName(String lockRepositoryName) {
return this;
}

public DynamoDBAuditStore withJournalRepositoryName(String journalRepositoryName) {
this.journalRepositoryName = journalRepositoryName;
return this;
}

public DynamoDBAuditStore withReadCapacityUnits(long readCapacityUnits) {
this.readCapacityUnits = readCapacityUnits;
return this;
Expand All @@ -95,34 +113,53 @@ public DynamoDBAuditStore withAutoCreate(boolean autoCreate) {
public void initialize(ContextResolver baseContext) {
runnerId = baseContext.getRequiredDependencyValue(RunnerId.class);
communityConfiguration = baseContext.getRequiredDependencyValue(CommunityConfigurable.class);
auditRepository = new DynamoDBAuditRepository(client);
journalEventStore = new DynamoDBJournalEventStore(
client,
journalRepositoryName,
readCapacityUnits,
writeCapacityUnits
);
journalEventSequencerFactory = new JournalEventSequencerFactory(journalEventStore);

lockService = new DynamoDBLockService(client, TimeService.getDefault());
lockService.initialize(
autoCreate,
lockRepositoryName,
readCapacityUnits,
writeCapacityUnits
);
this.validate();
}

@Override
public synchronized CommunityAuditPersistence getPersistence() {
if (persistence == null) {
public AuditPersistenceFactory<CommunityAuditPersistence> getPersistenceFactory() {
return stageId -> {
JournalEventSequencer journalEventSequencer = journalEventSequencerFactory.forStream(stageId);
persistence = new DynamoDBAuditPersistence(
client,
auditRepositoryName,
readCapacityUnits,
writeCapacityUnits,
autoCreate,
communityConfiguration);
communityConfiguration,
auditRepository,
journalEventStore,
journalEventSequencer,
targetSystem.getTxWrapper(),
auditRepositoryName,
readCapacityUnits,
writeCapacityUnits,
autoCreate
);
persistence.initialize(runnerId);
}
return persistence;
return persistence;
};
}

@Override
public AuditReader getAuditReader() {
auditRepository.initialize(autoCreate, auditRepositoryName, readCapacityUnits, writeCapacityUnits);
return () -> auditRepository.getAuditHistory();
}

@Override
public synchronized CommunityLockService getLockService() {
if (lockService == null) {
lockService = new DynamoDBLockService(client, TimeService.getDefault());
lockService.initialize(
autoCreate,
lockRepositoryName,
readCapacityUnits,
writeCapacityUnits);
}
return lockService;
}

Expand All @@ -140,8 +177,29 @@ private void validate() {
throw new FlamingockException("The 'lockRepositoryName' property is required.");
}

if (journalRepositoryName == null || journalRepositoryName.trim().isEmpty()) {
throw new FlamingockException("The 'journalRepositoryName' property is required.");
}

if (readCapacityUnits <= 0) {
throw new FlamingockException("The 'readCapacityUnits' property must be greater than zero.");
}

if (writeCapacityUnits <= 0) {
throw new FlamingockException("The 'writeCapacityUnits' property must be greater than zero.");
}

if (auditRepositoryName.trim().equalsIgnoreCase(lockRepositoryName.trim())) {
throw new FlamingockException("The 'auditRepositoryName' and 'lockRepositoryName' properties must not be the same.");
}

if (journalRepositoryName.trim().equalsIgnoreCase(auditRepositoryName.trim())) {
throw new FlamingockException("The 'journalRepositoryName' and 'auditRepositoryName' properties must not be the same.");
}

if (journalRepositoryName.trim().equalsIgnoreCase(lockRepositoryName.trim())) {
throw new FlamingockException("The 'journalRepositoryName' and 'lockRepositoryName' properties must not be the same.");
}
}

}
Original file line number Diff line number Diff line change
Expand Up @@ -16,32 +16,60 @@
package io.flamingock.store.dynamodb.internal;

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.core.journal.JournalEventSequencerFactory;
import io.flamingock.internal.util.FeatureFlag;
import io.flamingock.internal.util.Result;
import io.flamingock.internal.util.id.RunnerId;
import software.amazon.awssdk.services.dynamodb.DynamoDbClient;
import software.amazon.awssdk.enhanced.dynamodb.model.TransactWriteItemsEnhancedRequest;

import java.util.List;

public class DynamoDBAuditPersistence extends AbstractCommunityAuditPersistence {

private final DynamoDbClient client;
private final DynamoDBAuditRepository auditRepository;
private final DynamoDBJournalEventStore journalEventStore;
private JournalEventSequencer journalEventSequencer;
private final TransactionWrapper txWrapper;
private final String auditTableName;
private final long readCapacityUnits;
private final long writeCapacityUnits;
private final boolean autoCreate;

private DynamoDBAuditor auditor;

public DynamoDBAuditPersistence(DynamoDbClient client,
/**
* Creates a persistence over explicitly supplied audit, journal and transaction collaborators.
*
* @param localConfiguration community configuration
* @param auditRepository repository for audits
* @param journalEventStore journal store receiving staged events
* @param journalEventSequencer sequencer for the stage journal stream
* @param txWrapper transaction wrapper shared with the target system
* @param auditTableName audit table name
* @param readCapacityUnits audit and journal read capacity
* @param writeCapacityUnits audit and journal write capacity
* @param autoCreate whether missing tables may be created
*/
Comment thread
bercianor marked this conversation as resolved.
public DynamoDBAuditPersistence(CommunityConfigurable localConfiguration,
DynamoDBAuditRepository auditRepository,
DynamoDBJournalEventStore journalEventStore,
JournalEventSequencer journalEventSequencer,
TransactionWrapper txWrapper,
String auditTableName,
long readCapacityUnits,
long writeCapacityUnits,
boolean autoCreate,
CommunityConfigurable localConfiguration) {
boolean autoCreate) {
super(localConfiguration);
this.client = client;
this.auditRepository = auditRepository;
this.journalEventStore = journalEventStore;
this.journalEventSequencer = journalEventSequencer;
this.txWrapper = txWrapper;
this.auditTableName = auditTableName;
this.readCapacityUnits = readCapacityUnits;
this.writeCapacityUnits = writeCapacityUnits;
Expand All @@ -50,21 +78,44 @@ public DynamoDBAuditPersistence(DynamoDbClient client,

@Override
protected void doInitialize(RunnerId runnerId) {
auditor = new DynamoDBAuditor(client);
auditor.initialize(
auditRepository.initialize(
autoCreate,
auditTableName,
readCapacityUnits,
writeCapacityUnits);
if (isJournalEventsEnabled()) {
journalEventStore.initialize(autoCreate);
}
}

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

@Override
public Result writeEntry(AuditEntry auditEntry) {
return auditor.writeEntry(auditEntry);
if (isJournalEventsEnabled()) {
RuntimeContext baseContext = new BasicRuntimeContext("write-changeState-" + auditEntry.getChangeId());
Result result = txWrapper.wrapInTransaction(baseContext, runtimeContext -> {
TransactWriteItemsEnhancedRequest.Builder builder = runtimeContext.getContext()
.getRequiredDependencyValue(TransactWriteItemsEnhancedRequest.Builder.class);
JournalEvent<AuditEntry> journalEvent = journalEventSequencer.newEvent(auditEntry);
journalEventStore.contributeToTransaction(builder, journalEvent);
return auditRepository.contributeToTransaction(builder, auditEntry);
});
journalEventSequencer.confirm();
return result;
}
return auditRepository.writeEntry(auditEntry);
}

private static boolean isJournalEventsEnabled() {
try {
return FeatureFlag.isEnabled(Features.JOURNAL_EVENTS, false);
} catch (RuntimeException exception) {
return false;
}
}

}
Loading
Loading