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

extra["generalUtilVersion"] = "1.6.0"
extra["generalUtilVersion"] = "1.7.0-SNAPSHOT"
extra["templateApiVersion"] = "1.3.4"
extra["coreApiVersion"] = "1.3.3"
extra["sqlVersion"] = "1.3.2"
extra["sqlVersion"] = "1.4.0-SNAPSHOT"
extra["mongodbTemplateVersion"] = "1.3.2"

repositories {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -19,18 +19,15 @@
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.common.core.feature.Features;
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.FeatureFlag;
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;
Expand All @@ -50,7 +47,7 @@ public class DynamoDBAuditStore implements CommunityAuditStore {
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 String journalRepositoryName = CommunityPersistenceConstants.DEFAULT_JOURNAL_STORE_NAME;
private long readCapacityUnits = 5L;
private long writeCapacityUnits = 5L;
private boolean autoCreate = true;
Expand Down Expand Up @@ -139,10 +136,7 @@ public void initialize(ContextResolver baseContext) {
public AuditPersistenceFactory<CommunityAuditPersistence> getPersistenceFactory() {
return stageId -> {
auditRepository.initialize(autoCreate);
if (FeatureFlag.isEnabled(Features.JOURNAL_EVENTS, false)) {
journalEventStore.initialize(autoCreate);
}
JournalEventSequencer journalEventSequencer = journalEventSequencerFactory.forStream(stageId);
JournalEventSequencer journalEventSequencer = journalEventSequencerFactory.initializeForStage(stageId, autoCreate);
persistence = new DynamoDBAuditPersistence(
communityConfiguration,
auditRepository,
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -114,6 +114,7 @@ public DynamoDBJournalEventStore(DynamoDbClient client,
*
* @param autoCreate whether to create the table when missing
*/
@Override
public synchronized void initialize(boolean autoCreate) {
if (!isJournalEventsEnabled() || table != null) {
return;
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -45,9 +45,9 @@
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;
import static io.flamingock.internal.util.constants.CommunityPersistenceConstants.DEFAULT_JOURNAL_STORE_NAME;

public class MongoDBSyncAuditStore implements CommunityAuditStore {

Expand Down Expand Up @@ -152,7 +152,7 @@ public void initialize(ContextResolver baseContext) {
@Override
public AuditPersistenceFactory<CommunityAuditPersistence> getPersistenceFactory() {
return stageId -> {
JournalEventSequencer journalEventSequencer = journalEventSequencerFactory.forStream(stageId);
JournalEventSequencer journalEventSequencer = journalEventSequencerFactory.initializeForStage(stageId, autoCreate);
persistence = new MongoDBSyncAuditPersistence(
communityConfiguration,
auditRepository,
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -84,7 +84,8 @@ public MongoDBSyncJournalEventStore(MongoDatabase database,
.withWriteConcern(writeConcern);
}

protected void initialize(boolean autoCreate) {
@Override
public void initialize(boolean autoCreate) {
CollectionInitializator<MongoDBDocumentHelper> initializer = new CollectionInitializator<>(
new MongoDBSyncCollectionHelper(collection),
() -> new MongoDBDocumentHelper(new Document()),
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -19,15 +19,14 @@
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.common.core.feature.Features;
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.common.sql.journal.SqlJournalConstants;
import io.flamingock.internal.util.Constants;
import io.flamingock.internal.util.FeatureFlag;
import io.flamingock.internal.util.constants.CommunityPersistenceConstants;
import io.flamingock.internal.util.id.RunnerId;
import io.flamingock.store.sql.internal.SqlAuditPersistence;
Expand All @@ -40,9 +39,6 @@

public class SqlAuditStore implements CommunityAuditStore {

private static final String SQL_IDENTIFIER_PATTERN = "[A-Za-z][A-Za-z0-9_]*";
private static final String DEFAULT_JOURNAL_REPOSITORY_NAME = "flamingockJournalEvents";

private final SqlExternalSystem targetSystem;
private final DataSource dataSource;
private CommunityConfigurable communityConfiguration;
Expand All @@ -53,7 +49,7 @@ public class SqlAuditStore implements CommunityAuditStore {
private SqlAuditRepository auditRepository;
private String auditRepositoryName = CommunityPersistenceConstants.DEFAULT_AUDIT_STORE_NAME;
private String lockRepositoryName = CommunityPersistenceConstants.DEFAULT_LOCK_STORE_NAME;
private String journalRepositoryName = DEFAULT_JOURNAL_REPOSITORY_NAME;
private String journalRepositoryName = CommunityPersistenceConstants.DEFAULT_JOURNAL_STORE_NAME;
private boolean autoCreate = true;

private SqlAuditStore(SqlExternalSystem targetSystem) {
Expand Down Expand Up @@ -120,12 +116,8 @@ public void initialize(ContextResolver baseContext) {
@Override
public AuditPersistenceFactory<CommunityAuditPersistence> getPersistenceFactory() {
return stageId -> {
boolean journalEventsEnabled = isJournalEventsEnabled();
JournalEventSequencer journalEventSequencer = null;
if (journalEventsEnabled) {
journalEventStore.initialize(autoCreate);
journalEventSequencer = journalEventSequencerFactory.forStream(stageId);
}
JournalEventSequencer journalEventSequencer = journalEventSequencerFactory.initializeForStage(stageId, autoCreate);
boolean journalEventsEnabled = journalEventSequencer != null;

SqlAuditPersistence persistence = new SqlAuditPersistence(
communityConfiguration,
Expand Down Expand Up @@ -157,31 +149,25 @@ private void validate() {
validateRepositoryName(auditRepositoryName, "auditRepositoryName");
validateRepositoryName(lockRepositoryName, "lockRepositoryName");
validateRepositoryName(journalRepositoryName, "journalRepositoryName");
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.");
}
validateDistinct(auditRepositoryName, "auditRepositoryName", lockRepositoryName, "lockRepositoryName");
validateDistinct(journalRepositoryName, "journalRepositoryName", auditRepositoryName, "auditRepositoryName");
validateDistinct(journalRepositoryName, "journalRepositoryName", lockRepositoryName, "lockRepositoryName");
}

private void validateRepositoryName(String repositoryName, String propertyName) {
if (repositoryName == null || repositoryName.trim().isEmpty()) {
throw new FlamingockException(propertyName + " must not be blank");
}
if (!repositoryName.matches(SQL_IDENTIFIER_PATTERN)) {
throw new FlamingockException(propertyName + " must be a simple SQL identifier");
try {
SqlJournalConstants.validateIdentifier(repositoryName, propertyName);
} catch (IllegalArgumentException exception) {
throw new FlamingockException(exception.getMessage());
}
}

private static boolean isJournalEventsEnabled() {
private void validateDistinct(String firstName, String firstField, String secondName, String secondField) {
try {
return FeatureFlag.isEnabled(Features.JOURNAL_EVENTS, false);
} catch (RuntimeException exception) {
return false;
SqlJournalConstants.validateDistinct(firstName, firstField, secondName, secondField);
} catch (IllegalArgumentException exception) {
throw new FlamingockException("The '" + firstField + "' and '" + secondField
+ "' properties must not be the same.");
}
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -18,32 +18,25 @@
import io.flamingock.api.RecoveryStrategy;
import io.flamingock.internal.common.core.audit.AuditEntry;
import io.flamingock.internal.common.core.audit.AuditTxType;
import io.flamingock.internal.common.sql.journal.SqlAuditColumnNames;

import java.sql.PreparedStatement;
import java.sql.ResultSet;
import java.sql.SQLException;
import java.sql.Timestamp;
import java.sql.Types;
import java.util.Arrays;
import java.util.Collections;
import java.util.List;

/**
* Binds and reads the typed, flattened SQL representation of an {@link AuditEntry}.
*/
final class AuditEntryMapper {

private static final List<String> COLUMN_NAMES = Collections.unmodifiableList(Arrays.asList(
"execution_id", "stage_id", "change_id", "author", "created_at", "state", "invoked_class",
"invoked_method", "source_file", "metadata", "execution_millis", "execution_hostname",
"error_trace", "type", "tx_strategy", "target_system_id", "change_order", "recovery_strategy",
"transaction_flag", "system_change"));

private AuditEntryMapper() {
}

static List<String> columnNames() {
return COLUMN_NAMES;
return SqlAuditColumnNames.columnNames();
}

static void bind(PreparedStatement statement, AuditEntry auditEntry, int firstColumn) throws SQLException {
Expand Down Expand Up @@ -107,7 +100,7 @@ static AuditEntry fromResultSet(ResultSet resultSet) throws SQLException {
}

private static String columnName(int index) {
return COLUMN_NAMES.get(index);
return SqlAuditColumnNames.columnNames().get(index);
}

private static void setNullableBoolean(PreparedStatement statement,
Expand Down

This file was deleted.

Original file line number Diff line number Diff line change
Expand Up @@ -18,6 +18,7 @@
import io.flamingock.internal.common.core.audit.AuditEntry;
import io.flamingock.internal.common.sql.SqlDialect;
import io.flamingock.internal.common.sql.dialectHelpers.SqlAuditorDialectHelper;
import io.flamingock.internal.common.sql.journal.SqlJournalConstants;
import io.flamingock.internal.util.Result;

import javax.sql.DataSource;
Expand All @@ -35,7 +36,7 @@ public SqlAuditRepository(DataSource dataSource, String auditTableName) {
if (dataSource == null) {
throw new IllegalArgumentException("dataSource must not be null");
}
JournalEventConstants.validateIdentifier(auditTableName, "auditTableName");
SqlJournalConstants.validateIdentifier(auditTableName, "auditTableName");
this.dataSource = dataSource;
this.auditTableName = auditTableName;
}
Expand Down Expand Up @@ -113,7 +114,7 @@ Result save(Connection connection, AuditEntry auditEntry) {
if (auditEntry == null) {
throw new IllegalArgumentException("auditEntry must not be null");
}
JournalEventConstants.validateIdentifier(auditTableName, "auditTableName");
SqlJournalConstants.validateIdentifier(auditTableName, "auditTableName");
if (auditEntry.getChangeId() == null || auditEntry.getChangeId().trim().isEmpty()) {
throw new IllegalArgumentException("changeId must not be blank");
}
Expand Down
Loading
Loading