From ad106c93572b8aee7a2cc2502c875e29ffd6043e Mon Sep 17 00:00:00 2001 From: bercianor Date: Fri, 18 Sep 2026 17:19:29 +0200 Subject: [PATCH 1/3] fix(sql): close non-transactional execution connections --- .../store/sql/SqlAuditStoreTest.java | 47 +++++++++++----- .../targets/AbstractTargetSystem.java | 20 ++++++- .../targetsystem/sql/SqlTargetSystem.java | 11 ++++ .../SqlTargetSystemSharedTxManagerTest.java | 54 +++++++++++++++++++ 4 files changed, 118 insertions(+), 14 deletions(-) diff --git a/community/flamingock-sql-auditstore/src/test/java/io/flamingock/store/sql/SqlAuditStoreTest.java b/community/flamingock-sql-auditstore/src/test/java/io/flamingock/store/sql/SqlAuditStoreTest.java index 984fdc8ee..1d70b7d0c 100644 --- a/community/flamingock-sql-auditstore/src/test/java/io/flamingock/store/sql/SqlAuditStoreTest.java +++ b/community/flamingock-sql-auditstore/src/test/java/io/flamingock/store/sql/SqlAuditStoreTest.java @@ -94,6 +94,27 @@ void tearDown() throws SQLException { } } + @Test + @DisplayName("container-backed setup stops the container when datasource creation fails") + void containerBackedSetupStopsContainerWhenDataSourceCreationFails() { + JdbcDatabaseContainer container = org.mockito.Mockito.mock(JdbcDatabaseContainer.class); + RuntimeException setupFailure = new IllegalStateException("datasource creation failed"); + org.mockito.Mockito.when(container.isRunning()).thenReturn(true); + + try (MockedStatic helper = org.mockito.Mockito.mockStatic(SqlAuditTestHelper.class)) { + helper.when(() -> SqlAuditTestHelper.createContainer("informix")).thenReturn(container); + helper.when(() -> SqlAuditTestHelper.createDataSource(container)).thenThrow(setupFailure); + + RuntimeException thrown = assertThrows(RuntimeException.class, + () -> setupTest(SqlDialect.INFORMIX, "informix")); + + assertSame(setupFailure, thrown); + helper.verify(() -> SqlAuditTestHelper.createDataSource(container)); + org.mockito.Mockito.verify(container).start(); + org.mockito.Mockito.verify(container).stop(); + } + } + private TestContext setupTest(SqlDialect sqlDialect, String dialectName) throws SQLException { if ("h2".equals(dialectName)) { HikariConfig config = new HikariConfig(); @@ -145,22 +166,24 @@ private TestContext setupTest(SqlDialect sqlDialect, String dialectName) throws JdbcDatabaseContainer container = SqlAuditTestHelper.createContainer(dialectName); container.start(); - HikariConfig config = new HikariConfig(); - config.setJdbcUrl(container.getJdbcUrl()); - config.setUsername(container.getUsername()); - config.setPassword(container.getPassword()); - config.setDriverClassName(container.getDriverClassName()); - DataSource dataSource = new HikariDataSource(config); - TestContext testContext = new TestContext(dataSource, container, sqlDialect); - + DataSource dataSource = null; + TestContext testContext = null; try { + dataSource = SqlAuditTestHelper.createDataSource(container); + testContext = new TestContext(dataSource, container, sqlDialect); SqlAuditTestHelper.createTables(dataSource, sqlDialect); - } catch (SQLException exception) { - testContext.cleanup(); + return testContext; + } catch (SQLException | RuntimeException exception) { + TestContext cleanupContext = testContext != null + ? testContext + : new TestContext(dataSource, container, sqlDialect); + try { + cleanupContext.cleanup(); + } catch (SQLException cleanupFailure) { + exception.addSuppressed(cleanupFailure); + } throw exception; } - - return testContext; } private Class[] getChangeClasses(String dialectName, String scenario) { diff --git a/core/flamingock-core/src/main/java/io/flamingock/internal/core/external/targets/AbstractTargetSystem.java b/core/flamingock-core/src/main/java/io/flamingock/internal/core/external/targets/AbstractTargetSystem.java index 248143528..2c2bc04f5 100644 --- a/core/flamingock-core/src/main/java/io/flamingock/internal/core/external/targets/AbstractTargetSystem.java +++ b/core/flamingock-core/src/main/java/io/flamingock/internal/core/external/targets/AbstractTargetSystem.java @@ -93,7 +93,11 @@ public String getId() { */ public final T applyChange(Function changeApplier, ExecutionRuntime executionRuntime) { enhanceExecutionRuntime(executionRuntime, false); - return changeApplier.apply(executionRuntime); + try { + return changeApplier.apply(executionRuntime); + } finally { + cleanupExecutionRuntime(executionRuntime); + } } /** @@ -110,10 +114,22 @@ public final T applyChange(Function changeApplier, Exec */ public final T rollbackChange(Function changeRollbacker, ExecutionRuntime executionRuntime) { enhanceExecutionRuntime(executionRuntime, false); - return changeRollbacker.apply(executionRuntime); + try { + return changeRollbacker.apply(executionRuntime); + } finally { + cleanupExecutionRuntime(executionRuntime); + } } + /** + * Hook for cleaning up session-scoped dependencies after non-transactional execution. + * + * @param executionRuntime the runtime whose session-scoped dependencies should be cleaned up + */ + protected void cleanupExecutionRuntime(RuntimeContext executionRuntime) { + } + /** * Hook for injecting session-scoped dependencies into the execution runtime. *

diff --git a/core/target-systems/flamingock-sql-targetsystem/src/main/java/io/flamingock/targetsystem/sql/SqlTargetSystem.java b/core/target-systems/flamingock-sql-targetsystem/src/main/java/io/flamingock/targetsystem/sql/SqlTargetSystem.java index 133b00f35..cc7a83108 100644 --- a/core/target-systems/flamingock-sql-targetsystem/src/main/java/io/flamingock/targetsystem/sql/SqlTargetSystem.java +++ b/core/target-systems/flamingock-sql-targetsystem/src/main/java/io/flamingock/targetsystem/sql/SqlTargetSystem.java @@ -92,6 +92,17 @@ protected void enhanceExecutionRuntime(RuntimeContext executionRuntime, boolean } + @Override + protected void cleanupExecutionRuntime(RuntimeContext executionRuntime) { + executionRuntime.getContext().getDependencyValue(Connection.class).ifPresent(connection -> { + try { + connection.close(); + } catch (SQLException e) { + throw new FlamingockException(e); + } + }); + } + private SqlTxWrapper createTxWrapper(TransactionManager txManager) { return new SqlTxWrapper(txManager); } diff --git a/core/target-systems/flamingock-sql-targetsystem/src/test/java/io/flamingock/targetsystem/sql/SqlTargetSystemSharedTxManagerTest.java b/core/target-systems/flamingock-sql-targetsystem/src/test/java/io/flamingock/targetsystem/sql/SqlTargetSystemSharedTxManagerTest.java index 4a940f4eb..72f05fd3b 100644 --- a/core/target-systems/flamingock-sql-targetsystem/src/test/java/io/flamingock/targetsystem/sql/SqlTargetSystemSharedTxManagerTest.java +++ b/core/target-systems/flamingock-sql-targetsystem/src/test/java/io/flamingock/targetsystem/sql/SqlTargetSystemSharedTxManagerTest.java @@ -17,6 +17,7 @@ import io.flamingock.internal.common.core.context.ContextResolver; import io.flamingock.internal.core.builder.FlamingockEdition; +import io.flamingock.internal.core.runtime.ExecutionRuntime; import io.flamingock.internal.core.transaction.TransactionManager; import org.junit.jupiter.api.DisplayName; import org.junit.jupiter.api.Test; @@ -33,6 +34,7 @@ import static org.mockito.ArgumentMatchers.any; import static org.mockito.ArgumentMatchers.anyString; import static org.mockito.Mockito.mock; +import static org.mockito.Mockito.verify; import static org.mockito.Mockito.when; class SqlTargetSystemSharedTxManagerTest { @@ -58,6 +60,58 @@ void txWrapperAndAuditMarkerShouldShareSameTxManager() throws Exception { "SqlTxWrapper and SqlTargetSystemAuditMarker must share the same TransactionManager instance"); } + @Test + @DisplayName("Should close SQL connection after non-transactional apply") + void shouldCloseConnectionAfterNonTransactionalApply() throws Exception { + Connection connection = mock(Connection.class); + DataSource dataSource = mock(DataSource.class); + when(dataSource.getConnection()).thenReturn(connection); + + SqlTargetSystem targetSystem = new SqlTargetSystem("test-sql", dataSource); + targetSystem.applyChange(runtime -> null, mockRuntimeWith(connection)); + + verify(connection).close(); + } + + @Test + @DisplayName("Should close SQL connection after non-transactional rollback") + void shouldCloseConnectionAfterNonTransactionalRollback() throws Exception { + Connection connection = mock(Connection.class); + DataSource dataSource = mock(DataSource.class); + when(dataSource.getConnection()).thenReturn(connection); + + SqlTargetSystem targetSystem = new SqlTargetSystem("test-sql", dataSource); + targetSystem.rollbackChange(runtime -> null, mockRuntimeWith(connection)); + + verify(connection).close(); + } + + @Test + @DisplayName("Should close SQL connection when non-transactional callback fails") + void shouldCloseConnectionWhenNonTransactionalCallbackFails() throws Exception { + Connection connection = mock(Connection.class); + DataSource dataSource = mock(DataSource.class); + when(dataSource.getConnection()).thenReturn(connection); + RuntimeException callbackFailure = new RuntimeException("callback failed"); + + SqlTargetSystem targetSystem = new SqlTargetSystem("test-sql", dataSource); + RuntimeException thrown = assertThrows(RuntimeException.class, + () -> targetSystem.applyChange(runtime -> { + throw callbackFailure; + }, mockRuntimeWith(connection))); + + assertSame(callbackFailure, thrown); + verify(connection).close(); + } + + private static ExecutionRuntime mockRuntimeWith(Connection connection) { + ExecutionRuntime executionRuntime = mock(ExecutionRuntime.class); + ContextResolver contextResolver = mock(ContextResolver.class); + when(contextResolver.getDependencyValue(Connection.class)).thenReturn(Optional.of(connection)); + when(executionRuntime.getContext()).thenReturn(contextResolver); + return executionRuntime; + } + private static DataSource mockDataSource() throws Exception { ResultSet emptyResultSet = mock(ResultSet.class); when(emptyResultSet.next()).thenReturn(false); From 2941d47b7d48fa970b9d2854eb7ea7ab45812556 Mon Sep 17 00:00:00 2001 From: bercianor Date: Sat, 19 Sep 2026 13:44:55 +0200 Subject: [PATCH 2/3] fix(sql): scope non-transactional connection cleanup --- .../targets/AbstractTargetSystem.java | 38 +++++++++---------- .../targetsystem/sql/SqlTargetSystem.java | 37 +++++++++--------- .../SqlTargetSystemSharedTxManagerTest.java | 3 ++ 3 files changed, 39 insertions(+), 39 deletions(-) diff --git a/core/flamingock-core/src/main/java/io/flamingock/internal/core/external/targets/AbstractTargetSystem.java b/core/flamingock-core/src/main/java/io/flamingock/internal/core/external/targets/AbstractTargetSystem.java index 2c2bc04f5..29b97d496 100644 --- a/core/flamingock-core/src/main/java/io/flamingock/internal/core/external/targets/AbstractTargetSystem.java +++ b/core/flamingock-core/src/main/java/io/flamingock/internal/core/external/targets/AbstractTargetSystem.java @@ -83,8 +83,8 @@ public String getId() { * Applies a change operation with session-scoped dependency injection. *

* This method is the entry point for non-transactional change execution. - * It calls {@link #enhanceExecutionRuntime(RuntimeContext, boolean)} to allow - * subclasses to inject session-scoped dependencies before executing the change. + * It delegates lifecycle ownership to {@link #nonTxWrapper(Function, ExecutionRuntime)}, + * which enhances the runtime before invoking the change. * * @param the return type of the change operation * @param changeApplier the function that executes the actual change @@ -92,20 +92,15 @@ public String getId() { * @return the result of the change operation */ public final T applyChange(Function changeApplier, ExecutionRuntime executionRuntime) { - enhanceExecutionRuntime(executionRuntime, false); - try { - return changeApplier.apply(executionRuntime); - } finally { - cleanupExecutionRuntime(executionRuntime); - } + return nonTxWrapper(changeApplier, executionRuntime); } /** * Rolls back (reverts) a previously applied change with session-scoped dependency injection. *

* This method is the entry point for non-transactional rollback execution. - * It calls {@link #enhanceExecutionRuntime(RuntimeContext, boolean)} to allow - * subclasses to inject session-scoped dependencies before executing the rollback. + * It delegates lifecycle ownership to {@link #nonTxWrapper(Function, ExecutionRuntime)}, + * which enhances the runtime before invoking the rollback. * * @param the return type of the rollback operation * @param changeRollbacker the function that executes the actual rollback @@ -113,23 +108,28 @@ public final T applyChange(Function changeApplier, Exec * @return the result of the rollback operation */ public final T rollbackChange(Function changeRollbacker, ExecutionRuntime executionRuntime) { - enhanceExecutionRuntime(executionRuntime, false); - try { - return changeRollbacker.apply(executionRuntime); - } finally { - cleanupExecutionRuntime(executionRuntime); - } + return nonTxWrapper(changeRollbacker, executionRuntime); } /** - * Hook for cleaning up session-scoped dependencies after non-transactional execution. + * Executes a non-transactional callback and owns its runtime lifecycle. + *

+ * The default implementation enhances the runtime once before executing the callback. + * Subclasses that manage resources for non-transactional execution must enhance the runtime + * and release those resources within their override. * - * @param executionRuntime the runtime whose session-scoped dependencies should be cleaned up + * @param changeFunc the callback to execute + * @param executionRuntime the runtime to enhance and pass to the callback + * @param the callback return type + * @return the callback result */ - protected void cleanupExecutionRuntime(RuntimeContext executionRuntime) { + protected T nonTxWrapper(Function changeFunc, ExecutionRuntime executionRuntime) { + enhanceExecutionRuntime(executionRuntime, false); + return changeFunc.apply(executionRuntime); } + /** * Hook for injecting session-scoped dependencies into the execution runtime. *

diff --git a/core/target-systems/flamingock-sql-targetsystem/src/main/java/io/flamingock/targetsystem/sql/SqlTargetSystem.java b/core/target-systems/flamingock-sql-targetsystem/src/main/java/io/flamingock/targetsystem/sql/SqlTargetSystem.java index cc7a83108..5f633e7c5 100644 --- a/core/target-systems/flamingock-sql-targetsystem/src/main/java/io/flamingock/targetsystem/sql/SqlTargetSystem.java +++ b/core/target-systems/flamingock-sql-targetsystem/src/main/java/io/flamingock/targetsystem/sql/SqlTargetSystem.java @@ -22,12 +22,14 @@ import io.flamingock.internal.core.builder.FlamingockEdition; import io.flamingock.internal.core.external.targets.TransactionalTargetSystem; import io.flamingock.internal.core.external.targets.mark.NoOpTargetSystemAuditMarker; +import io.flamingock.internal.core.runtime.ExecutionRuntime; import io.flamingock.internal.core.transaction.TransactionManager; import io.flamingock.internal.common.core.transaction.TransactionWrapper; import javax.sql.DataSource; import java.sql.Connection; import java.sql.SQLException; +import java.util.function.Function; import static io.flamingock.internal.core.builder.FlamingockEdition.COMMUNITY; @@ -79,28 +81,23 @@ public TransactionWrapper getTxWrapper() { return txWrapper; } + /** + * Acquires one connection for a non-transactional callback, injects it into the runtime, + * and closes that same connection after the callback completes or fails. + * + * @param changeFunc the callback to execute + * @param executionRuntime the runtime to receive the connection dependency + * @param the callback return type + * @return the callback result + */ @Override - protected void enhanceExecutionRuntime(RuntimeContext executionRuntime, boolean isTransactional) { - //if transactional, the connection is injected in the wrapInTransaction - if (!isTransactional) { - try { - executionRuntime.addDependency(dataSource.getConnection()); - } catch (SQLException e) { - throw new FlamingockException(e); - } + protected T nonTxWrapper(Function changeFunc, ExecutionRuntime executionRuntime) { + try (Connection connection = dataSource.getConnection()) { + executionRuntime.addDependency(connection); + return changeFunc.apply(executionRuntime); + } catch (SQLException e) { + throw new FlamingockException(e); } - - } - - @Override - protected void cleanupExecutionRuntime(RuntimeContext executionRuntime) { - executionRuntime.getContext().getDependencyValue(Connection.class).ifPresent(connection -> { - try { - connection.close(); - } catch (SQLException e) { - throw new FlamingockException(e); - } - }); } private SqlTxWrapper createTxWrapper(TransactionManager txManager) { diff --git a/core/target-systems/flamingock-sql-targetsystem/src/test/java/io/flamingock/targetsystem/sql/SqlTargetSystemSharedTxManagerTest.java b/core/target-systems/flamingock-sql-targetsystem/src/test/java/io/flamingock/targetsystem/sql/SqlTargetSystemSharedTxManagerTest.java index 72f05fd3b..a1b1c963b 100644 --- a/core/target-systems/flamingock-sql-targetsystem/src/test/java/io/flamingock/targetsystem/sql/SqlTargetSystemSharedTxManagerTest.java +++ b/core/target-systems/flamingock-sql-targetsystem/src/test/java/io/flamingock/targetsystem/sql/SqlTargetSystemSharedTxManagerTest.java @@ -70,6 +70,7 @@ void shouldCloseConnectionAfterNonTransactionalApply() throws Exception { SqlTargetSystem targetSystem = new SqlTargetSystem("test-sql", dataSource); targetSystem.applyChange(runtime -> null, mockRuntimeWith(connection)); + verify(dataSource).getConnection(); verify(connection).close(); } @@ -83,6 +84,7 @@ void shouldCloseConnectionAfterNonTransactionalRollback() throws Exception { SqlTargetSystem targetSystem = new SqlTargetSystem("test-sql", dataSource); targetSystem.rollbackChange(runtime -> null, mockRuntimeWith(connection)); + verify(dataSource).getConnection(); verify(connection).close(); } @@ -101,6 +103,7 @@ void shouldCloseConnectionWhenNonTransactionalCallbackFails() throws Exception { }, mockRuntimeWith(connection))); assertSame(callbackFailure, thrown); + verify(dataSource).getConnection(); verify(connection).close(); } From 22b766fdbd688c43d36694f6173ec8a8b8a29d2f Mon Sep 17 00:00:00 2001 From: Antonio Perez Dieppa Date: Mon, 21 Sep 2026 20:54:48 +0100 Subject: [PATCH 3/3] refactor: nonTxWrapper --- .../cloud/utils/TestCloudTargetSystem.java | 8 +- .../core/utils/EmptyTransactionWrapper.java | 6 +- .../internal/DynamoDBAuditPersistence.java | 9 +- .../DynamoDBExternalSystemContractTest.java | 2 +- .../DynamoDBJournalEventStoreTest.java | 2 +- .../reactive/MongoDBReactiveAuditStore.java | 4 +- .../MongoDBReactiveAuditPersistence.java | 12 +- .../MongoDBReactiveAuditPersistenceTest.java | 4 +- .../mongodb/sync/MongoDBSyncAuditStore.java | 6 +- .../internal/MongoDBSyncAuditPersistence.java | 14 +-- .../sql/internal/SqlAuditPersistence.java | 8 +- .../sql/internal/SqlJournalEventStore.java | 8 +- .../SqlJournalEventStoreJdbcTest.java | 8 +- .../core/external/ExecutionWrapper.java | 88 +++++++++++++ .../TransactionalExternalSystem.java | 14 ++- .../core/transaction/TransactionWrapper.java | 117 ------------------ .../core/context/BasicRuntimeContext.java | 5 +- .../targets/AbstractTargetSystem.java | 74 +++++++---- .../targets/TransactionalTargetSystem.java | 13 +- .../pipeline/execution/StageExecutor.java | 7 +- .../core/context/BasicRuntimeContextTest.java | 12 +- .../targets/TargetSystemManagerTest.java | 4 +- .../couchbase/CouchbaseTargetSystem.java | 4 +- .../couchbase/CouchbaseTxWrapper.java | 6 +- .../dynamodb/api/DynamoDBExternalSystem.java | 2 +- .../dynamodb/DynamoDBTargetSystem.java | 4 +- .../dynamodb/DynamoDBTxWrapper.java | 6 +- .../mongodb/api/MongoDBExternalSystem.java | 3 +- .../api/MongoDBReactiveExternalSystem.java | 2 +- .../reactive/MongoDBReactiveTargetSystem.java | 4 +- .../reactive/MongoDBReactiveTxWrapper.java | 6 +- ...MongoDBSpringDataReactiveTargetSystem.java | 4 +- .../MongoDBSpringDataReactiveTxWrapper.java | 6 +- .../MongoDBSpringDataTargetSystem.java | 4 +- .../MongoDBSpringDataTxWrapper.java | 6 +- .../mongodb/sync/MongoDBSyncTargetSystem.java | 4 +- .../mongodb/sync/MongoDBSyncTxWrapper.java | 6 +- .../sql/api/SqlExternalSystem.java | 2 +- .../targetsystem/sql/SqlTargetSystem.java | 45 ++++--- .../targetsystem/sql/SqlTxWrapper.java | 6 +- .../sql/SqlTxWrapperLifecycleTest.java | 14 +-- docs/CHANGE-STEP-NAVIGATION.md | 4 +- docs/EXECUTION_FLOW_GUIDE.md | 2 +- 43 files changed, 291 insertions(+), 274 deletions(-) create mode 100644 core/flamingock-core-commons/src/main/java/io/flamingock/internal/common/core/external/ExecutionWrapper.java rename core/flamingock-core-commons/src/main/java/io/flamingock/internal/common/core/{transaction => external}/TransactionalExternalSystem.java (69%) delete mode 100644 core/flamingock-core-commons/src/main/java/io/flamingock/internal/common/core/transaction/TransactionWrapper.java diff --git a/cloud/flamingock-cloud/src/test/java/io/flamingock/core/cloud/utils/TestCloudTargetSystem.java b/cloud/flamingock-cloud/src/test/java/io/flamingock/core/cloud/utils/TestCloudTargetSystem.java index c2ecb7b91..5c8ac8e85 100644 --- a/cloud/flamingock-cloud/src/test/java/io/flamingock/core/cloud/utils/TestCloudTargetSystem.java +++ b/cloud/flamingock-cloud/src/test/java/io/flamingock/core/cloud/utils/TestCloudTargetSystem.java @@ -20,7 +20,7 @@ import io.flamingock.internal.core.external.targets.mark.TargetSystemAuditMark; import io.flamingock.internal.core.external.targets.mark.TargetSystemAuditMarker; import io.flamingock.internal.core.external.targets.TransactionalTargetSystem; -import io.flamingock.internal.common.core.transaction.TransactionWrapper; +import io.flamingock.internal.common.core.external.ExecutionWrapper; import org.mockito.Mockito; import java.util.Arrays; @@ -45,7 +45,7 @@ public TargetSystemAuditMarker getAuditMarker() { } @Override - public TransactionWrapper getTxWrapper() { + public ExecutionWrapper getTxWrapper() { return txWrapper; } @@ -60,10 +60,10 @@ protected TestCloudTargetSystem getSelf() { } - public static class TestCloudTxWrapper implements TransactionWrapper { + public static class TestCloudTxWrapper implements ExecutionWrapper { @Override - public RESULT wrapInTransaction(CONTEXT executionContext, Function operation) { + public RESULT wrapExecution(CONTEXT executionContext, Function operation) { return operation.apply(executionContext); } diff --git a/cloud/flamingock-cloud/src/test/java/io/flamingock/core/utils/EmptyTransactionWrapper.java b/cloud/flamingock-cloud/src/test/java/io/flamingock/core/utils/EmptyTransactionWrapper.java index d38821021..2d3b76cac 100644 --- a/cloud/flamingock-cloud/src/test/java/io/flamingock/core/utils/EmptyTransactionWrapper.java +++ b/cloud/flamingock-cloud/src/test/java/io/flamingock/core/utils/EmptyTransactionWrapper.java @@ -16,11 +16,11 @@ package io.flamingock.core.utils; import io.flamingock.internal.common.core.context.RuntimeContext; -import io.flamingock.internal.common.core.transaction.TransactionWrapper; +import io.flamingock.internal.common.core.external.ExecutionWrapper; import java.util.function.Function; -public class EmptyTransactionWrapper implements TransactionWrapper { +public class EmptyTransactionWrapper implements ExecutionWrapper { private boolean called = false; @@ -31,7 +31,7 @@ public boolean isCalled() { @Override - public RESULT wrapInTransaction(CONTEXT executionContext, Function operation) { + public RESULT wrapExecution(CONTEXT executionContext, Function operation) { called = true; return operation.apply(executionContext); } diff --git a/community/flamingock-dynamodb-auditstore/src/main/java/io/flamingock/store/dynamodb/internal/DynamoDBAuditPersistence.java b/community/flamingock-dynamodb-auditstore/src/main/java/io/flamingock/store/dynamodb/internal/DynamoDBAuditPersistence.java index 1efbfaad5..68a65e0dd 100644 --- a/community/flamingock-dynamodb-auditstore/src/main/java/io/flamingock/store/dynamodb/internal/DynamoDBAuditPersistence.java +++ b/community/flamingock-dynamodb-auditstore/src/main/java/io/flamingock/store/dynamodb/internal/DynamoDBAuditPersistence.java @@ -19,12 +19,11 @@ 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.common.core.external.ExecutionWrapper; 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; @@ -37,7 +36,7 @@ public class DynamoDBAuditPersistence extends AbstractCommunityAuditPersistence private final DynamoDBAuditRepository auditRepository; private final DynamoDBJournalEventStore journalEventStore; private JournalEventSequencer journalEventSequencer; - private final TransactionWrapper txWrapper; + private final ExecutionWrapper txWrapper; private final boolean autoCreate; /** @@ -54,7 +53,7 @@ public DynamoDBAuditPersistence(CommunityConfigurable localConfiguration, DynamoDBAuditRepository auditRepository, DynamoDBJournalEventStore journalEventStore, JournalEventSequencer journalEventSequencer, - TransactionWrapper txWrapper, + ExecutionWrapper txWrapper, boolean autoCreate) { super(localConfiguration); this.auditRepository = auditRepository; @@ -81,7 +80,7 @@ public List getAuditHistory() { public Result writeEntry(AuditEntry auditEntry) { if (isJournalEventsEnabled()) { RuntimeContext baseContext = new BasicRuntimeContext("write-changeState-" + auditEntry.getChangeId()); - Result result = txWrapper.wrapInTransaction(baseContext, runtimeContext -> { + Result result = txWrapper.wrapExecution(baseContext, runtimeContext -> { TransactWriteItemsEnhancedRequest.Builder builder = runtimeContext.getContext() .getRequiredDependencyValue(TransactWriteItemsEnhancedRequest.Builder.class); JournalEvent journalEvent = journalEventSequencer.newEvent(auditEntry); diff --git a/community/flamingock-dynamodb-auditstore/src/test/java/io/flamingock/store/dynamodb/DynamoDBExternalSystemContractTest.java b/community/flamingock-dynamodb-auditstore/src/test/java/io/flamingock/store/dynamodb/DynamoDBExternalSystemContractTest.java index 90e281841..9b51d9435 100644 --- a/community/flamingock-dynamodb-auditstore/src/test/java/io/flamingock/store/dynamodb/DynamoDBExternalSystemContractTest.java +++ b/community/flamingock-dynamodb-auditstore/src/test/java/io/flamingock/store/dynamodb/DynamoDBExternalSystemContractTest.java @@ -16,7 +16,7 @@ package io.flamingock.store.dynamodb; import io.flamingock.externalsystem.dynamodb.api.DynamoDBExternalSystem; -import io.flamingock.internal.common.core.transaction.TransactionalExternalSystem; +import io.flamingock.internal.common.core.external.TransactionalExternalSystem; import org.junit.jupiter.api.DisplayName; import org.junit.jupiter.api.Test; diff --git a/community/flamingock-dynamodb-auditstore/src/test/java/io/flamingock/store/dynamodb/internal/DynamoDBJournalEventStoreTest.java b/community/flamingock-dynamodb-auditstore/src/test/java/io/flamingock/store/dynamodb/internal/DynamoDBJournalEventStoreTest.java index 536fe187f..0cee48f3c 100644 --- a/community/flamingock-dynamodb-auditstore/src/test/java/io/flamingock/store/dynamodb/internal/DynamoDBJournalEventStoreTest.java +++ b/community/flamingock-dynamodb-auditstore/src/test/java/io/flamingock/store/dynamodb/internal/DynamoDBJournalEventStoreTest.java @@ -503,7 +503,7 @@ private void commit(List> events) { DynamoDBTxWrapper txWrapper = new DynamoDBTxWrapper( client, new TransactionManager<>(TransactWriteItemsEnhancedRequest::builder)); - txWrapper.wrapInTransaction(new BasicRuntimeContext("session-" + UUID.randomUUID()), ctx -> { + txWrapper.wrapExecution(new BasicRuntimeContext("session-" + UUID.randomUUID()), ctx -> { TransactWriteItemsEnhancedRequest.Builder builder = ctx.getContext() .getRequiredDependencyValue(TransactWriteItemsEnhancedRequest.Builder.class); for (JournalEvent event : events) { diff --git a/community/flamingock-mongodb-reactive-auditstore/src/main/java/io/flamingock/store/mongodb/reactive/MongoDBReactiveAuditStore.java b/community/flamingock-mongodb-reactive-auditstore/src/main/java/io/flamingock/store/mongodb/reactive/MongoDBReactiveAuditStore.java index 23bb4cc8c..91c3c33bc 100644 --- a/community/flamingock-mongodb-reactive-auditstore/src/main/java/io/flamingock/store/mongodb/reactive/MongoDBReactiveAuditStore.java +++ b/community/flamingock-mongodb-reactive-auditstore/src/main/java/io/flamingock/store/mongodb/reactive/MongoDBReactiveAuditStore.java @@ -41,7 +41,7 @@ import io.flamingock.store.mongodb.reactive.internal.MongoDBReactiveJournalEventStore; import io.flamingock.store.mongodb.reactive.internal.MongoDBReactiveLockService; -import io.flamingock.internal.common.core.transaction.TransactionWrapper; +import io.flamingock.internal.common.core.external.ExecutionWrapper; import java.util.Collections; import java.util.HashSet; @@ -162,7 +162,7 @@ public AuditPersistenceFactory getPersistenceFactory( } JournalEventSequencer journalEventSequencer = journalEventSequencerFactory.forStream(stageId); boolean supportsTransactions = mongoDBTargetSystem.supportsTransactions(); - TransactionWrapper txWrapper = supportsTransactions ? mongoDBTargetSystem.getTxWrapper() : null; + ExecutionWrapper txWrapper = supportsTransactions ? mongoDBTargetSystem.getTxWrapper() : null; MongoDBReactiveAuditPersistence stagePersistence = new MongoDBReactiveAuditPersistence( communityConfiguration, auditRepository, diff --git a/community/flamingock-mongodb-reactive-auditstore/src/main/java/io/flamingock/store/mongodb/reactive/internal/MongoDBReactiveAuditPersistence.java b/community/flamingock-mongodb-reactive-auditstore/src/main/java/io/flamingock/store/mongodb/reactive/internal/MongoDBReactiveAuditPersistence.java index ae34c35c5..eed5d1072 100644 --- a/community/flamingock-mongodb-reactive-auditstore/src/main/java/io/flamingock/store/mongodb/reactive/internal/MongoDBReactiveAuditPersistence.java +++ b/community/flamingock-mongodb-reactive-auditstore/src/main/java/io/flamingock/store/mongodb/reactive/internal/MongoDBReactiveAuditPersistence.java @@ -20,7 +20,7 @@ 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.common.core.external.ExecutionWrapper; import io.flamingock.internal.core.context.BasicRuntimeContext; import io.flamingock.internal.core.configuration.community.CommunityConfigurable; import io.flamingock.internal.core.external.store.audit.community.AbstractCommunityAuditPersistence; @@ -38,7 +38,7 @@ public class MongoDBReactiveAuditPersistence extends AbstractCommunityAuditPersi private final MongoDBReactiveJournalEventStore journalEventStore; private final JournalEventSequencer journalEventSequencer; private final boolean supportsTransactions; - private final TransactionWrapper txWrapper; + private final ExecutionWrapper txWrapper; private final boolean autoCreate; /** @@ -55,7 +55,7 @@ public MongoDBReactiveAuditPersistence(CommunityConfigurable localConfiguration, MongoDBReactiveJournalEventStore journalEventStore, JournalEventSequencer journalEventSequencer, boolean supportsTransactions, - TransactionWrapper txWrapper, + ExecutionWrapper txWrapper, boolean autoCreate) { super(localConfiguration); this.auditRepository = auditRepository; @@ -66,8 +66,8 @@ public MongoDBReactiveAuditPersistence(CommunityConfigurable localConfiguration, this.autoCreate = autoCreate; } - private static TransactionWrapper validateTransactionWrapper(boolean supportsTransactions, - TransactionWrapper txWrapper) { + private static ExecutionWrapper validateTransactionWrapper(boolean supportsTransactions, + ExecutionWrapper txWrapper) { if (supportsTransactions) { return Objects.requireNonNull( txWrapper, @@ -112,7 +112,7 @@ public Result writeEntry(AuditEntry auditEntry) { private Result writeJournalAndAuditInTransaction(AuditEntry auditEntry) { RuntimeContext baseContext = new BasicRuntimeContext("write-changeState-" + auditEntry.getChangeId()); - Result result = txWrapper.wrapInTransaction(baseContext, runtimeContext -> { + Result result = txWrapper.wrapExecution(baseContext, runtimeContext -> { ClientSession clientSession = runtimeContext.getContext().getRequiredDependencyValue(ClientSession.class); JournalEvent journalEvent = journalEventSequencer.newEvent(auditEntry); journalEventStore.append(clientSession, journalEvent); diff --git a/community/flamingock-mongodb-reactive-auditstore/src/test/java/io/flamingock/store/mongodb/reactive/internal/MongoDBReactiveAuditPersistenceTest.java b/community/flamingock-mongodb-reactive-auditstore/src/test/java/io/flamingock/store/mongodb/reactive/internal/MongoDBReactiveAuditPersistenceTest.java index a97a6d083..039917a1e 100644 --- a/community/flamingock-mongodb-reactive-auditstore/src/test/java/io/flamingock/store/mongodb/reactive/internal/MongoDBReactiveAuditPersistenceTest.java +++ b/community/flamingock-mongodb-reactive-auditstore/src/test/java/io/flamingock/store/mongodb/reactive/internal/MongoDBReactiveAuditPersistenceTest.java @@ -27,7 +27,7 @@ import io.flamingock.internal.common.core.feature.Features; import io.flamingock.internal.core.configuration.community.CommunityConfigurable; import io.flamingock.internal.core.journal.JournalEventSequencer; -import io.flamingock.internal.common.core.transaction.TransactionWrapper; +import io.flamingock.internal.common.core.external.ExecutionWrapper; import io.flamingock.internal.util.FeatureFlag; import io.flamingock.internal.util.id.RunnerId; import io.flamingock.reactive.util.PublisherSync; @@ -79,7 +79,7 @@ void beforeEach() { journalEventStore, mock(JournalEventSequencer.class), true, - mock(TransactionWrapper.class), + mock(ExecutionWrapper.class), true); persistence.initialize(RunnerId.fromString("runner-1")); } diff --git a/community/flamingock-mongodb-sync-auditstore/src/main/java/io/flamingock/store/mongodb/sync/MongoDBSyncAuditStore.java b/community/flamingock-mongodb-sync-auditstore/src/main/java/io/flamingock/store/mongodb/sync/MongoDBSyncAuditStore.java index d677607a9..4a9336890 100644 --- a/community/flamingock-mongodb-sync-auditstore/src/main/java/io/flamingock/store/mongodb/sync/MongoDBSyncAuditStore.java +++ b/community/flamingock-mongodb-sync-auditstore/src/main/java/io/flamingock/store/mongodb/sync/MongoDBSyncAuditStore.java @@ -20,13 +20,12 @@ 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.common.core.feature.Features; -import io.flamingock.internal.common.core.transaction.TransactionWrapper; +import io.flamingock.internal.common.core.external.ExecutionWrapper; 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; @@ -46,7 +45,6 @@ import java.util.Collections; import java.util.HashSet; -import java.util.List; import java.util.Set; import org.slf4j.Logger; @@ -163,7 +161,7 @@ public AuditPersistenceFactory getPersistenceFactory( return stageId -> { JournalEventSequencer journalEventSequencer = journalEventSequencerFactory.forStream(stageId); boolean supportsTransactions = mongoDBTargetSystem.supportsTransactions(); - TransactionWrapper txWrapper = supportsTransactions ? mongoDBTargetSystem.getTxWrapper() : null; + ExecutionWrapper txWrapper = supportsTransactions ? mongoDBTargetSystem.getTxWrapper() : null; persistence = new MongoDBSyncAuditPersistence( communityConfiguration, auditRepository, diff --git a/community/flamingock-mongodb-sync-auditstore/src/main/java/io/flamingock/store/mongodb/sync/internal/MongoDBSyncAuditPersistence.java b/community/flamingock-mongodb-sync-auditstore/src/main/java/io/flamingock/store/mongodb/sync/internal/MongoDBSyncAuditPersistence.java index 928967f06..fc3d976fd 100644 --- a/community/flamingock-mongodb-sync-auditstore/src/main/java/io/flamingock/store/mongodb/sync/internal/MongoDBSyncAuditPersistence.java +++ b/community/flamingock-mongodb-sync-auditstore/src/main/java/io/flamingock/store/mongodb/sync/internal/MongoDBSyncAuditPersistence.java @@ -20,7 +20,7 @@ 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.common.core.external.ExecutionWrapper; 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; @@ -38,7 +38,7 @@ public class MongoDBSyncAuditPersistence extends AbstractCommunityAuditPersisten private final MongoDBSyncJournalEventStore journalEventStore; private final JournalEventSequencer journalEventSequencer; private final boolean supportsTransactions; - private final TransactionWrapper txWrapper; + private final ExecutionWrapper txWrapper; private final boolean autoCreate; /** @@ -55,7 +55,7 @@ public MongoDBSyncAuditPersistence(CommunityConfigurable localConfiguration, MongoDBSyncJournalEventStore journalEventStore, JournalEventSequencer journalEventSequencer, boolean supportsTransactions, - TransactionWrapper txWrapper, + ExecutionWrapper txWrapper, boolean autoCreate) { super(localConfiguration); this.auditRepository = auditRepository; @@ -66,8 +66,8 @@ public MongoDBSyncAuditPersistence(CommunityConfigurable localConfiguration, this.autoCreate = autoCreate; } - private static TransactionWrapper validateTransactionWrapper(boolean supportsTransactions, - TransactionWrapper txWrapper) { + private static ExecutionWrapper validateTransactionWrapper(boolean supportsTransactions, + ExecutionWrapper txWrapper) { if (supportsTransactions) { return Objects.requireNonNull( txWrapper, @@ -113,14 +113,14 @@ public Result writeEntry(AuditEntry auditEntry) { private Result writeJournalAndAuditInTransaction(AuditEntry auditEntry) { RuntimeContext baseContext = new BasicRuntimeContext("write-changeState-" + auditEntry.getChangeId()); - Result result = txWrapper.wrapInTransaction(baseContext, runtimeContext -> { + Result result = txWrapper.wrapExecution(baseContext, runtimeContext -> { ClientSession clientSession = runtimeContext.getContext().getRequiredDependencyValue(ClientSession.class); JournalEvent journalEvent = journalEventSequencer.newEvent(auditEntry); journalEventStore.write(clientSession, journalEvent); return auditRepository.save(clientSession, auditEntry); }); // Spends the stream position, and only a committed transaction may reach this line. In general a - // normal return from wrapInTransaction does NOT mean commit — a FailedStep result is returned + // normal return from wrapExecution does NOT mean commit — a FailedStep result is returned // after a rollback, without an exception. It is sound here because this operation returns a // Result, which can never be a FailedStep, so the commit branch is the only graceful path; a // failing commit is caught and rethrown as DatabaseTransactionException. Keep that true: an diff --git a/community/flamingock-sql-auditstore/src/main/java/io/flamingock/store/sql/internal/SqlAuditPersistence.java b/community/flamingock-sql-auditstore/src/main/java/io/flamingock/store/sql/internal/SqlAuditPersistence.java index 43bc30b08..91fe9730a 100644 --- a/community/flamingock-sql-auditstore/src/main/java/io/flamingock/store/sql/internal/SqlAuditPersistence.java +++ b/community/flamingock-sql-auditstore/src/main/java/io/flamingock/store/sql/internal/SqlAuditPersistence.java @@ -18,7 +18,7 @@ import io.flamingock.internal.common.core.audit.AuditEntry; import io.flamingock.internal.common.core.context.RuntimeContext; import io.flamingock.internal.common.core.journal.JournalEvent; -import io.flamingock.internal.common.core.transaction.TransactionWrapper; +import io.flamingock.internal.common.core.external.ExecutionWrapper; 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; @@ -34,7 +34,7 @@ public class SqlAuditPersistence extends AbstractCommunityAuditPersistence { private final SqlAuditRepository auditRepository; private final SqlJournalEventStore journalEventStore; private final JournalEventSequencer journalEventSequencer; - private final TransactionWrapper txWrapper; + private final ExecutionWrapper txWrapper; private final boolean journalEventsEnabled; /** @@ -51,7 +51,7 @@ public SqlAuditPersistence(CommunityConfigurable localConfiguration, SqlAuditRepository auditRepository, SqlJournalEventStore journalEventStore, JournalEventSequencer journalEventSequencer, - TransactionWrapper txWrapper, + ExecutionWrapper txWrapper, boolean journalEventsEnabled) { super(localConfiguration); this.auditRepository = auditRepository; @@ -86,7 +86,7 @@ public synchronized Result writeEntry(AuditEntry auditEntry) { } RuntimeContext baseContext = new BasicRuntimeContext("write-changeState-" + auditEntry.getChangeId()); - Result result = txWrapper.wrapInTransaction(baseContext, runtimeContext -> { + Result result = txWrapper.wrapExecution(baseContext, runtimeContext -> { Connection connection = runtimeContext.getContext().getRequiredDependencyValue(Connection.class); JournalEvent journalEvent = journalEventSequencer.newEvent(auditEntry); journalEventStore.append(connection, journalEvent); diff --git a/community/flamingock-sql-auditstore/src/main/java/io/flamingock/store/sql/internal/SqlJournalEventStore.java b/community/flamingock-sql-auditstore/src/main/java/io/flamingock/store/sql/internal/SqlJournalEventStore.java index 02981c8cf..b68e4a331 100644 --- a/community/flamingock-sql-auditstore/src/main/java/io/flamingock/store/sql/internal/SqlJournalEventStore.java +++ b/community/flamingock-sql-auditstore/src/main/java/io/flamingock/store/sql/internal/SqlJournalEventStore.java @@ -18,7 +18,7 @@ import io.flamingock.internal.common.core.audit.AuditEntry; import io.flamingock.internal.common.core.context.RuntimeContext; import io.flamingock.internal.common.core.journal.JournalEvent; -import io.flamingock.internal.common.core.transaction.TransactionWrapper; +import io.flamingock.internal.common.core.external.ExecutionWrapper; import io.flamingock.internal.common.sql.SqlDialect; import io.flamingock.internal.common.sql.SqlDialectFactory; import io.flamingock.internal.core.context.BasicRuntimeContext; @@ -54,7 +54,7 @@ public class SqlJournalEventStore implements JournalEventStore { private final DataSource dataSource; private final String tableName; - private final TransactionWrapper txWrapper; + private final ExecutionWrapper txWrapper; private SqlJournalEventMapper mapper; private SqlJournalDialectHelper dialectHelper; @@ -66,7 +66,7 @@ public class SqlJournalEventStore implements JournalEventStore { * @param tableName journal table name * @param txWrapper SQL transaction wrapper that owns Journal writes */ - public SqlJournalEventStore(DataSource dataSource, String tableName, TransactionWrapper txWrapper) { + public SqlJournalEventStore(DataSource dataSource, String tableName, ExecutionWrapper txWrapper) { if (dataSource == null) { throw new IllegalArgumentException("dataSource must not be null"); } @@ -183,7 +183,7 @@ public long acknowledgeEvents(Collection eventIds) { for (String eventId : validEventIds) { RuntimeContext baseContext = new BasicRuntimeContext( "acknowledge-journal-event-" + UUID.randomUUID()); - acknowledged += txWrapper.wrapInTransaction(baseContext, runtimeContext -> { + acknowledged += txWrapper.wrapExecution(baseContext, runtimeContext -> { Connection connection = runtimeContext.getContext().getRequiredDependencyValue(Connection.class); return acknowledgeEvents(connection, Collections.singleton(eventId)); }); diff --git a/community/flamingock-sql-auditstore/src/test/java/io/flamingock/store/sql/internal/SqlJournalEventStoreJdbcTest.java b/community/flamingock-sql-auditstore/src/test/java/io/flamingock/store/sql/internal/SqlJournalEventStoreJdbcTest.java index 3d4b112b3..f86a469a1 100644 --- a/community/flamingock-sql-auditstore/src/test/java/io/flamingock/store/sql/internal/SqlJournalEventStoreJdbcTest.java +++ b/community/flamingock-sql-auditstore/src/test/java/io/flamingock/store/sql/internal/SqlJournalEventStoreJdbcTest.java @@ -24,7 +24,7 @@ import io.flamingock.internal.common.core.error.DatabaseTransactionException; import io.flamingock.internal.common.core.journal.JournalEvent; import io.flamingock.internal.common.core.journal.JournalEventType; -import io.flamingock.internal.common.core.transaction.TransactionWrapper; +import io.flamingock.internal.common.core.external.ExecutionWrapper; import io.flamingock.internal.common.sql.SqlDialect; import io.flamingock.internal.core.transaction.TransactionManager; import io.flamingock.targetsystem.sql.SqlTxWrapper; @@ -233,12 +233,12 @@ void acknowledgesEachEventInItsOwnTransaction() throws Exception { append(event("stage-2", 1L, "second-event", false)); int[] transactionCount = {0}; - TransactionWrapper failingAfterSecondTransaction = new TransactionWrapper() { + ExecutionWrapper failingAfterSecondTransaction = new ExecutionWrapper() { @Override - public RESULT wrapInTransaction( + public RESULT wrapExecution( CONTEXT runtimeContext, Function operation) { int transactionNumber = ++transactionCount[0]; - return txWrapper.wrapInTransaction(runtimeContext, context -> { + return txWrapper.wrapExecution(runtimeContext, context -> { RESULT result = operation.apply(context); if (transactionNumber == 2) { throw new IllegalStateException("forced second acknowledgement failure"); diff --git a/core/flamingock-core-commons/src/main/java/io/flamingock/internal/common/core/external/ExecutionWrapper.java b/core/flamingock-core-commons/src/main/java/io/flamingock/internal/common/core/external/ExecutionWrapper.java new file mode 100644 index 000000000..f2c55e51a --- /dev/null +++ b/core/flamingock-core-commons/src/main/java/io/flamingock/internal/common/core/external/ExecutionWrapper.java @@ -0,0 +1,88 @@ +/* + * Copyright 2023 Flamingock (https://www.flamingock.io) + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package io.flamingock.internal.common.core.external; + +import io.flamingock.internal.common.core.context.RuntimeContext; +import io.flamingock.internal.common.core.error.DatabaseTransactionException; + +import java.util.function.Function; + +/** + * Runs an operation inside an execution boundary: acquire scoped resources, publish their handles into + * the {@link RuntimeContext}, apply the operation, release the resources — on every outcome. The caller + * supplies only the work; it reaches the handles by resolving them from the context it gets back. + *

+ * There are two kinds, and the method name that returns a wrapper is what tells them apart: + *

    + *
  • Non-transactional ({@code AbstractTargetSystem#getNonTxWrapper()}) — scopes + * resources, promises no atomicity. Often just applies the operation; SQL borrows a + * {@code Connection} and closes it.
  • + *
  • Transactional ({@link TransactionalExternalSystem#getTxWrapper()}) — also opens + * a transaction and commits or rolls back. Used to apply a change transactionally, and by audit + * stores to make a group of their own writes atomic.
  • + *
+ * Holding an {@code ExecutionWrapper} therefore tells you nothing about whether your work is + * transactional — that was decided by which wrapper you were handed. + * + * @see RuntimeContext + * @see DatabaseTransactionException + */ +public interface ExecutionWrapper { + + /** + * Executes {@code operation} within this wrapper's boundary and returns its result. + * + *

The one thing to get right

+ * A normal return does not mean the work committed. An operation can report failure + * two ways, and they are handled differently: + *
    + *
  • By value — it returns a {@code FailedStep}. A transactional wrapper rolls + * back and returns that same failed step, without throwing. This is what lets the + * engine audit the failure and drive recovery from a value instead of a stack unwind.
  • + *
  • By exception — a transactional wrapper rolls back and rethrows wrapped in + * {@link DatabaseTransactionException}, which carries the diagnostics, including whether the + * rollback itself succeeded + * ({@link DatabaseTransactionException.RollbackStatus RollbackStatus}).
  • + *
+ * So callers that need to know the work is durable — to advance a counter, publish, acknowledge — + * must check the result, or use an operation whose return type cannot be a failed step. A + * non-transactional wrapper has nothing staged to roll back: results and exceptions pass through + * as-is. + * + *

Implementing one

+ *
    + *
  • Publish handles only once they are usable, and release them in a {@code finally} / + * try-with-resources so no path leaks them.
  • + *
  • Return the operation's result untouched — never rewrite it.
  • + *
  • Do not call {@code enhanceExecutionRuntime}: the caller already did, and + * doing it again double-injects. Contribute only what your own boundary owns.
  • + *
  • Transactional only: key the session on {@link RuntimeContext#getSessionId()} (the change id, + * for change execution), and surface every transactional failure as + * {@link DatabaseTransactionException} so callers have one exception type to handle.
  • + *
+ * + * @param the concrete runtime context type, passed to the operation as-is + * @param the type produced by the operation + * @param runtimeContext identifies the execution scope and receives the scoped dependencies + * @param operation the work to execute within the boundary + * @return the operation's result, untouched — including a failed step returned after a rollback + * @throws DatabaseTransactionException if a transactional operation throws, or the transaction cannot + * be started, committed or rolled back. Non-transactional + * wrappers propagate the original exception unwrapped + */ + RESULT wrapExecution(CONTEXT runtimeContext, Function operation); + +} diff --git a/core/flamingock-core-commons/src/main/java/io/flamingock/internal/common/core/transaction/TransactionalExternalSystem.java b/core/flamingock-core-commons/src/main/java/io/flamingock/internal/common/core/external/TransactionalExternalSystem.java similarity index 69% rename from core/flamingock-core-commons/src/main/java/io/flamingock/internal/common/core/transaction/TransactionalExternalSystem.java rename to core/flamingock-core-commons/src/main/java/io/flamingock/internal/common/core/external/TransactionalExternalSystem.java index 96d1ad594..07646f112 100644 --- a/core/flamingock-core-commons/src/main/java/io/flamingock/internal/common/core/transaction/TransactionalExternalSystem.java +++ b/core/flamingock-core-commons/src/main/java/io/flamingock/internal/common/core/external/TransactionalExternalSystem.java @@ -13,7 +13,7 @@ * See the License for the specific language governing permissions and * limitations under the License. */ -package io.flamingock.internal.common.core.transaction; +package io.flamingock.internal.common.core.external; import io.flamingock.api.external.ExternalSystem; @@ -33,12 +33,18 @@ default boolean supportsTransactions() { } /** - * Returns the transaction wrapper for this target system. + * Returns the transactional {@link ExecutionWrapper} for this external system. *

* The wrapper is responsible for starting, committing, and rolling back transactions, * as well as injecting transaction-scoped dependencies into the execution runtime. + *

+ * It is one of two wrappers a target system exposes — the other, + * {@code AbstractTargetSystem#getNonTxWrapper()}, serves changes that run outside a transaction. + * Both are {@code ExecutionWrapper}s, so the method name, not the type, is what states the intent. + * This one is only used when {@link #supportsTransactions()} is {@code true} and the change + * itself is declared transactional. * - * @return the transaction wrapper instance + * @return the transactional wrapper instance */ - TransactionWrapper getTxWrapper(); + ExecutionWrapper getTxWrapper(); } diff --git a/core/flamingock-core-commons/src/main/java/io/flamingock/internal/common/core/transaction/TransactionWrapper.java b/core/flamingock-core-commons/src/main/java/io/flamingock/internal/common/core/transaction/TransactionWrapper.java deleted file mode 100644 index ec42c2ec0..000000000 --- a/core/flamingock-core-commons/src/main/java/io/flamingock/internal/common/core/transaction/TransactionWrapper.java +++ /dev/null @@ -1,117 +0,0 @@ -/* - * Copyright 2023 Flamingock (https://www.flamingock.io) - * - * Licensed under the Apache License, Version 2.0 (the "License"); - * you may not use this file except in compliance with the License. - * You may obtain a copy of the License at - * - * http://www.apache.org/licenses/LICENSE-2.0 - * - * Unless required by applicable law or agreed to in writing, software - * distributed under the License is distributed on an "AS IS" BASIS, - * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. - * See the License for the specific language governing permissions and - * limitations under the License. - */ -package io.flamingock.internal.common.core.transaction; - -import io.flamingock.internal.common.core.context.RuntimeContext; -import io.flamingock.internal.common.core.error.DatabaseTransactionException; - -import java.util.function.Function; - -/** - * Runs an operation inside a target system's transaction boundary. - *

- * A wrapper owns the whole transaction lifecycle: it opens the transaction, publishes the - * transaction-scoped handle (a JDBC {@code Connection}, a MongoDB {@code ClientSession}, a DynamoDB - * write-request builder, …) into the supplied {@link RuntimeContext}, runs the operation, and then - * commits or rolls back according to how that operation ended — releasing the underlying resources - * either way. Callers contribute only the work to be run; they reach the handle, if they need it, by - * resolving it from the context the wrapper handed back to them. - *

- * Implementations are used in two distinct places: - *

    - *
  • Transactional target systems — to apply a change within the target system's - * own transaction.
  • - *
  • Audit stores — to make a group of the store's own writes atomic (for example, - * a change-state transition and its journal event).
  • - *
- * - * @see RuntimeContext - * @see DatabaseTransactionException - */ -public interface TransactionWrapper { - - /** - * Executes {@code operation} within a transaction and returns its result. - * - *

Lifecycle

- *
    - *
  1. Open a transaction, keyed by {@link RuntimeContext#getSessionId()}.
  2. - *
  3. Publish the transaction-scoped handle into {@code runtimeContext}, making it resolvable by - * the operation.
  4. - *
  5. Apply the operation.
  6. - *
  7. Commit or roll back, according to the outcome below.
  8. - *
  9. Release the transaction resources — always, on every outcome.
  10. - *
- * - *

Outcomes

- * An operation can end in three ways, and the two failure paths are deliberately different: - *
    - *
  • Success. The operation returns normally and the returned value is not a - * failed step. The transaction is committed and the value is returned to the caller - * unchanged — this method never rewrites the operation's result.
  • - * - *
  • Failure reported as a value. The operation returns a - * {@code io.flamingock.internal.core.change.navigation.step.FailedStep}, meaning the work - * failed but the failure has already been captured as part of the result. The transaction is - * rolled back (or, where writes are only staged until commit, simply never committed) and - * that same failed step is returned to the caller. No exception is thrown. - * This is what allows the execution engine to audit the failure and drive recovery from a - * returned value rather than by unwinding a stack.
  • - * - *
  • Failure reported as an exception. The operation throws. The transaction is - * rolled back and the original exception is rethrown wrapped in a - * {@link DatabaseTransactionException}, which carries the diagnostic context of the attempt — - * transaction state, duration, connection details and, importantly, whether the rollback - * itself succeeded - * ({@link DatabaseTransactionException.RollbackStatus RollbackStatus}). A failed rollback is - * reported through that status, not by swallowing the original failure.
  • - *
- * Note that a rolled-back transaction is therefore not in itself an error signal to the - * caller: whether it learns about the failure by value or by exception depends on how the operation - * chose to report it. - * - *

Implementation contract

- *
    - *
  • Use {@link RuntimeContext#getSessionId()} as the key for the underlying session. For change - * execution this is the change id; other callers scope it as they see fit, but it must - * identify one transaction scope and no other.
  • - *
  • Inject transaction-scoped dependencies only after the transaction has started, so - * the operation cannot observe a handle that is not yet usable.
  • - *
  • Return the operation's result untouched on success and on a failed step.
  • - *
  • Surface every transactional failure as a {@link DatabaseTransactionException}, so callers - * have one exception type to reason about regardless of the target system.
  • - *
  • Release resources in a {@code finally} block, so neither a commit failure nor a rollback - * failure can leak a session.
  • - *
- * - * @param the concrete runtime context type, returned to the operation as-is so it - * keeps its static type - * @param the type produced by the operation - * @param runtimeContext the context identifying the transaction scope and receiving the - * transaction-scoped dependencies - * @param operation the work to execute within the transaction - * @return whatever the operation returned — including a failed step, when the operation reported - * its failure that way. A normal return does not imply the transaction committed: - * a failed step is returned after a rollback, without an exception. Callers that must know the - * work is durable — to advance a counter, publish, or acknowledge — cannot infer it from a - * normal return alone; they either check the result themselves, or rely on an operation whose - * return type cannot be a failed step. - * @throws DatabaseTransactionException if the operation throws, or if the transaction itself cannot - * be started, committed or rolled back - */ - RESULT wrapInTransaction(CONTEXT runtimeContext, Function operation); - -} diff --git a/core/flamingock-core/src/main/java/io/flamingock/internal/core/context/BasicRuntimeContext.java b/core/flamingock-core/src/main/java/io/flamingock/internal/core/context/BasicRuntimeContext.java index 8e33d7ee3..d8057b741 100644 --- a/core/flamingock-core/src/main/java/io/flamingock/internal/core/context/BasicRuntimeContext.java +++ b/core/flamingock-core/src/main/java/io/flamingock/internal/core/context/BasicRuntimeContext.java @@ -19,13 +19,14 @@ import io.flamingock.internal.common.core.context.ContextResolver; import io.flamingock.internal.common.core.context.Dependency; import io.flamingock.internal.common.core.context.RuntimeContext; +import io.flamingock.internal.common.core.external.ExecutionWrapper; import java.util.Collection; /** * Minimal {@link RuntimeContext}: a session id plus a layered, writable dependency context. *

- * It exists for the cases where a {@link io.flamingock.internal.common.core.transaction.TransactionWrapper} + * It exists for the cases where an {@link ExecutionWrapper} * is needed outside change execution — the wrapper contract requires a {@code RuntimeContext} to publish * its transaction-scoped dependencies into, but the caller has no change to run and therefore no use for * the reflection machinery of {@code ExecutionRuntime}. @@ -36,7 +37,7 @@ * *

{@code
  * BasicRuntimeContext runtimeContext = new BasicRuntimeContext(sessionId);
- * txWrapper.wrapInTransaction(runtimeContext, ctx -> {
+ * txWrapper.wrapExecution(runtimeContext, ctx -> {
  *     ClientSession session = ctx.getContext().getRequiredDependencyValue(ClientSession.class);
  *     auditRepository.write(session, auditEntry);
  *     journalRepository.append(session, journalEvent);
diff --git a/core/flamingock-core/src/main/java/io/flamingock/internal/core/external/targets/AbstractTargetSystem.java b/core/flamingock-core/src/main/java/io/flamingock/internal/core/external/targets/AbstractTargetSystem.java
index 29b97d496..f98caaadf 100644
--- a/core/flamingock-core/src/main/java/io/flamingock/internal/core/external/targets/AbstractTargetSystem.java
+++ b/core/flamingock-core/src/main/java/io/flamingock/internal/core/external/targets/AbstractTargetSystem.java
@@ -17,6 +17,7 @@
 
 import io.flamingock.api.external.TargetSystem;
 import io.flamingock.internal.common.core.context.*;
+import io.flamingock.internal.common.core.external.ExecutionWrapper;
 import io.flamingock.internal.core.context.SimpleContext;
 import io.flamingock.internal.core.runtime.ExecutionRuntime;
 import io.flamingock.internal.util.Property;
@@ -51,9 +52,15 @@
  * including dependencies and properties that will be available to changes
  * during execution.
  * 

- * Subclasses should override {@link #enhanceExecutionRuntime(RuntimeContext, boolean)} - * to inject session-scoped dependencies (e.g., database connections, client sessions) - * that are obtained fresh for each change execution. + * Subclasses have two distinct extension points, and the split matters: + *

    + *
  • {@link #enhanceExecutionRuntime(RuntimeContext, boolean)} — injects session-scoped dependencies + * (database connections, client sessions) obtained fresh for each change execution. Called by the + * framework on both the transactional and the non-transactional path, before any wrapper runs.
  • + *
  • {@link #getNonTxWrapper()} — owns the resource lifecycle of non-transactional execution, when + * something must be acquired for the call and released afterwards. Its transactional counterpart + * is {@code TransactionalExternalSystem#getTxWrapper()}.
  • + *
* * @param the concrete target system type for fluent API support */ @@ -62,6 +69,18 @@ public abstract class AbstractTargetSystem { + + /** + * Pass-through wrapper: applies the operation and nothing else. Stateless, so one instance serves + * every target system that has no resource to acquire for non-transactional execution. + */ + private static final ExecutionWrapper PASS_THROUGH_WRAPPER = new ExecutionWrapper() { + @Override + public RESULT wrapExecution(CONTEXT runtimeContext, Function operation) { + return operation.apply(runtimeContext); + } + }; + private final String id; protected final Context targetSystemContext = new SimpleContext(); @@ -82,9 +101,9 @@ public String getId() { /** * Applies a change operation with session-scoped dependency injection. *

- * This method is the entry point for non-transactional change execution. - * It delegates lifecycle ownership to {@link #nonTxWrapper(Function, ExecutionRuntime)}, - * which enhances the runtime before invoking the change. + * This method is the entry point for non-transactional change execution. It enhances the runtime + * with the target system's session-scoped dependencies, then delegates execution to + * {@link #getNonTxWrapper()}, which owns whatever resources that execution needs. * * @param the return type of the change operation * @param changeApplier the function that executes the actual change @@ -92,15 +111,16 @@ public String getId() { * @return the result of the change operation */ public final T applyChange(Function changeApplier, ExecutionRuntime executionRuntime) { - return nonTxWrapper(changeApplier, executionRuntime); + enhanceExecutionRuntime(executionRuntime, false); + return getNonTxWrapper().wrapExecution(executionRuntime, changeApplier); } /** * Rolls back (reverts) a previously applied change with session-scoped dependency injection. *

- * This method is the entry point for non-transactional rollback execution. - * It delegates lifecycle ownership to {@link #nonTxWrapper(Function, ExecutionRuntime)}, - * which enhances the runtime before invoking the rollback. + * This method is the entry point for non-transactional rollback execution. It enhances the runtime + * with the target system's session-scoped dependencies, then delegates execution to + * {@link #getNonTxWrapper()}, which owns whatever resources that execution needs. * * @param the return type of the rollback operation * @param changeRollbacker the function that executes the actual rollback @@ -108,25 +128,33 @@ public final T applyChange(Function changeApplier, Exec * @return the result of the rollback operation */ public final T rollbackChange(Function changeRollbacker, ExecutionRuntime executionRuntime) { - return nonTxWrapper(changeRollbacker, executionRuntime); + enhanceExecutionRuntime(executionRuntime, false); + return getNonTxWrapper().wrapExecution(executionRuntime, changeRollbacker); } - /** - * Executes a non-transactional callback and owns its runtime lifecycle. + * Returns the wrapper used for non-transactional execution — the counterpart to + * {@code TransactionalExternalSystem#getTxWrapper()}, which is used when the change runs in a + * transaction. Both are {@link ExecutionWrapper}s; the two method names are what distinguish the + * intent, since the type alone no longer says whether a transaction is involved. + *

+ * The default implementation is a pass-through: it simply applies the operation, because a target + * system with nothing to acquire has no boundary to own. *

- * The default implementation enhances the runtime once before executing the callback. - * Subclasses that manage resources for non-transactional execution must enhance the runtime - * and release those resources within their override. + * Override this when non-transactional execution needs a resource for the duration of the call — a + * pooled connection, a client session — and that resource must be released afterwards. Acquire it, + * publish it into the supplied context, apply the operation, and release it on every outcome + * (try-with-resources, or a {@code finally} block). + *

+ * Do not enhance the runtime in an override. {@link #applyChange} and + * {@link #rollbackChange} already call {@link #enhanceExecutionRuntime(RuntimeContext, boolean)} + * with {@code isTransactional=false} before invoking the wrapper; enhancing again would inject the + * session-scoped dependencies twice. An override contributes only the handles its own boundary owns. * - * @param changeFunc the callback to execute - * @param executionRuntime the runtime to enhance and pass to the callback - * @param the callback return type - * @return the callback result + * @return the wrapper owning resource lifecycle for non-transactional apply and rollback */ - protected T nonTxWrapper(Function changeFunc, ExecutionRuntime executionRuntime) { - enhanceExecutionRuntime(executionRuntime, false); - return changeFunc.apply(executionRuntime); + protected ExecutionWrapper getNonTxWrapper() { + return PASS_THROUGH_WRAPPER; } diff --git a/core/flamingock-core/src/main/java/io/flamingock/internal/core/external/targets/TransactionalTargetSystem.java b/core/flamingock-core/src/main/java/io/flamingock/internal/core/external/targets/TransactionalTargetSystem.java index fa338e6af..443c362f2 100644 --- a/core/flamingock-core/src/main/java/io/flamingock/internal/core/external/targets/TransactionalTargetSystem.java +++ b/core/flamingock-core/src/main/java/io/flamingock/internal/core/external/targets/TransactionalTargetSystem.java @@ -19,11 +19,11 @@ import io.flamingock.internal.common.core.audit.AuditReaderType; import io.flamingock.internal.common.core.context.ContextInitializable; import io.flamingock.internal.common.core.context.RuntimeContext; -import io.flamingock.internal.common.core.transaction.TransactionalExternalSystem; +import io.flamingock.internal.common.core.external.TransactionalExternalSystem; import io.flamingock.internal.core.runtime.ExecutionRuntime; import io.flamingock.internal.core.external.targets.mark.NoOpTargetSystemAuditMarker; import io.flamingock.internal.core.external.targets.mark.TargetSystemAuditMarker; -import io.flamingock.internal.common.core.transaction.TransactionWrapper; +import io.flamingock.internal.common.core.external.ExecutionWrapper; import java.util.Optional; import java.util.function.Function; @@ -36,7 +36,10 @@ *

* Subclasses must provide: *

    - *
  • A {@link TransactionWrapper} for managing transactions
  • + *
  • An {@link ExecutionWrapper}, returned from {@link #getTxWrapper()}, that runs the change inside + * the target system's transaction. It is the transactional counterpart of + * {@link #getNonTxWrapper()}, which this class inherits and which still serves changes declared + * non-transactional — a transactional target system uses both paths.
  • *
  • An audit marker for tracking execution state (optional in Community Edition)
  • *
* @@ -65,7 +68,7 @@ public boolean hasMarker() { *
    *
  1. Calling {@link #enhanceExecutionRuntime(RuntimeContext, boolean)} with * {@code isTransactional=true} for session-scoped dependency injection
  2. - *
  3. Delegating to the {@link TransactionWrapper} for transaction management + *
  4. Delegating to the {@link ExecutionWrapper} for transaction management * and potential injection of transaction-scoped dependencies
  5. *
* @@ -76,7 +79,7 @@ public boolean hasMarker() { */ public final T applyChangeTransactional(Function changeApplier, ExecutionRuntime executionRuntime) { enhanceExecutionRuntime(executionRuntime, true); - return getTxWrapper().wrapInTransaction(executionRuntime, changeApplier); + return getTxWrapper().wrapExecution(executionRuntime, changeApplier); } /** diff --git a/core/flamingock-core/src/main/java/io/flamingock/internal/core/pipeline/execution/StageExecutor.java b/core/flamingock-core/src/main/java/io/flamingock/internal/core/pipeline/execution/StageExecutor.java index d22759a94..8767082a9 100644 --- a/core/flamingock-core/src/main/java/io/flamingock/internal/core/pipeline/execution/StageExecutor.java +++ b/core/flamingock-core/src/main/java/io/flamingock/internal/core/pipeline/execution/StageExecutor.java @@ -23,7 +23,6 @@ import io.flamingock.internal.common.core.pipeline.StageDescriptor; import io.flamingock.internal.common.core.response.data.StageResult; import io.flamingock.internal.core.context.PriorityContext; -import io.flamingock.internal.core.external.store.audit.community.CommunityAuditPersistence; import io.flamingock.internal.core.external.store.lock.Lock; import io.flamingock.internal.core.external.targets.TargetSystemManager; import io.flamingock.internal.core.operation.result.StageResultBuilder; @@ -32,7 +31,7 @@ import io.flamingock.internal.core.change.navigation.navigator.ChangeProcessResult; import io.flamingock.internal.core.change.navigation.navigator.ChangeProcessStrategy; import io.flamingock.internal.core.change.navigation.navigator.ChangeProcessStrategyFactory; -import io.flamingock.internal.common.core.transaction.TransactionWrapper; +import io.flamingock.internal.common.core.external.ExecutionWrapper; import io.flamingock.internal.util.log.FlamingockLoggerFactory; import org.slf4j.Logger; @@ -52,13 +51,13 @@ public class StageExecutor { private final ContextResolver baseDependencyContext; private final Set> nonGuardedTypes; private final TargetSystemManager targetSystemManager; - protected final TransactionWrapper auditStoreTxWrapper; + protected final ExecutionWrapper auditStoreTxWrapper; public StageExecutor(ContextResolver dependencyContext, Set> nonGuardedTypes, AuditPersistenceFactory auditPersistenceFactory, TargetSystemManager targetSystemManager, - TransactionWrapper auditStoreTxWrapper) { + ExecutionWrapper auditStoreTxWrapper) { this.baseDependencyContext = dependencyContext; this.nonGuardedTypes = nonGuardedTypes; this.auditPersistenceFactory = auditPersistenceFactory; diff --git a/core/flamingock-core/src/test/java/io/flamingock/internal/core/context/BasicRuntimeContextTest.java b/core/flamingock-core/src/test/java/io/flamingock/internal/core/context/BasicRuntimeContextTest.java index 9a83f3903..e732d192f 100644 --- a/core/flamingock-core/src/test/java/io/flamingock/internal/core/context/BasicRuntimeContextTest.java +++ b/core/flamingock-core/src/test/java/io/flamingock/internal/core/context/BasicRuntimeContextTest.java @@ -16,7 +16,7 @@ package io.flamingock.internal.core.context; import io.flamingock.internal.common.core.context.Dependency; -import io.flamingock.internal.common.core.transaction.TransactionWrapper; +import io.flamingock.internal.common.core.external.ExecutionWrapper; import io.flamingock.internal.common.core.context.RuntimeContext; import org.junit.jupiter.api.DisplayName; import org.junit.jupiter.api.Test; @@ -29,7 +29,7 @@ /** * Tests for {@link BasicRuntimeContext} — the minimal {@link RuntimeContext} used when a - * {@link TransactionWrapper} is needed outside change execution. + * {@link ExecutionWrapper} is needed outside change execution. */ class BasicRuntimeContextTest { @@ -117,10 +117,10 @@ void shouldResolveBulkInjectedDependencies() { void shouldCarrySessionInjectedByTransactionWrapper() { Session opened = new Session("opened-by-wrapper"); // Stands in for a real wrapper: starts a "transaction", publishes its session, runs the operation - TransactionWrapper txWrapper = new TransactionWrapper() { + ExecutionWrapper txWrapper = new ExecutionWrapper() { @Override - public RESULT wrapInTransaction(CONTEXT runtimeContext, - Function operation) { + public RESULT wrapExecution(CONTEXT runtimeContext, + Function operation) { runtimeContext.addDependency(new Dependency(Session.class, opened, false)); return operation.apply(runtimeContext); } @@ -130,7 +130,7 @@ public RESULT wrapInTransaction(CONTEXT assertFalse(runtimeContext.getContext().getDependency(Session.class).isPresent(), "The session must not be resolvable before the wrapper opens the transaction"); - Session seenByOperation = txWrapper.wrapInTransaction(runtimeContext, + Session seenByOperation = txWrapper.wrapExecution(runtimeContext, ctx -> ctx.getContext().getRequiredDependencyValue(Session.class)); assertEquals(opened, seenByOperation); diff --git a/core/flamingock-core/src/test/java/io/flamingock/internal/core/external/targets/TargetSystemManagerTest.java b/core/flamingock-core/src/test/java/io/flamingock/internal/core/external/targets/TargetSystemManagerTest.java index 2efda8ece..c6dccfea4 100644 --- a/core/flamingock-core/src/test/java/io/flamingock/internal/core/external/targets/TargetSystemManagerTest.java +++ b/core/flamingock-core/src/test/java/io/flamingock/internal/core/external/targets/TargetSystemManagerTest.java @@ -17,7 +17,7 @@ import io.flamingock.internal.common.core.context.RuntimeContext; import io.flamingock.internal.common.core.targets.OperationType; -import io.flamingock.internal.common.core.transaction.TransactionWrapper; +import io.flamingock.internal.common.core.external.ExecutionWrapper; import org.junit.jupiter.api.DisplayName; import org.junit.jupiter.api.Test; @@ -101,7 +101,7 @@ private static class StubTransactionalTargetSystem extends TransactionalTargetSy } @Override protected StubTransactionalTargetSystem getSelf() { return this; } @Override public boolean supportsTransactions() { return transactionsSupported; } - @Override public TransactionWrapper getTxWrapper() { return null; } + @Override public ExecutionWrapper getTxWrapper() { return null; } @Override protected void enhanceExecutionRuntime(RuntimeContext rt, boolean tx) {} @Override public void initialize(io.flamingock.internal.common.core.context.ContextResolver ctx) {} } diff --git a/core/target-systems/flamingock-couchbase-targetsystem/src/main/java/io/flamingock/targetsystem/couchbase/CouchbaseTargetSystem.java b/core/target-systems/flamingock-couchbase-targetsystem/src/main/java/io/flamingock/targetsystem/couchbase/CouchbaseTargetSystem.java index ec3b9a8f1..bdc209413 100644 --- a/core/target-systems/flamingock-couchbase-targetsystem/src/main/java/io/flamingock/targetsystem/couchbase/CouchbaseTargetSystem.java +++ b/core/target-systems/flamingock-couchbase-targetsystem/src/main/java/io/flamingock/targetsystem/couchbase/CouchbaseTargetSystem.java @@ -29,7 +29,7 @@ import io.flamingock.internal.core.external.targets.TransactionalTargetSystem; import io.flamingock.internal.core.external.targets.mark.NoOpTargetSystemAuditMarker; import io.flamingock.internal.core.transaction.TransactionManager; -import io.flamingock.internal.common.core.transaction.TransactionWrapper; +import io.flamingock.internal.common.core.external.ExecutionWrapper; import java.util.function.Supplier; import java.util.Objects; @@ -111,7 +111,7 @@ protected CouchbaseTargetSystem getSelf() { } @Override - public TransactionWrapper getTxWrapper() { + public ExecutionWrapper getTxWrapper() { return txWrapper; } diff --git a/core/target-systems/flamingock-couchbase-targetsystem/src/main/java/io/flamingock/targetsystem/couchbase/CouchbaseTxWrapper.java b/core/target-systems/flamingock-couchbase-targetsystem/src/main/java/io/flamingock/targetsystem/couchbase/CouchbaseTxWrapper.java index 50a91fa65..77c71a8ff 100644 --- a/core/target-systems/flamingock-couchbase-targetsystem/src/main/java/io/flamingock/targetsystem/couchbase/CouchbaseTxWrapper.java +++ b/core/target-systems/flamingock-couchbase-targetsystem/src/main/java/io/flamingock/targetsystem/couchbase/CouchbaseTxWrapper.java @@ -22,12 +22,12 @@ import io.flamingock.internal.common.core.context.RuntimeContext; import io.flamingock.internal.core.transaction.TransactionManager; import io.flamingock.internal.core.change.navigation.step.FailedStep; -import io.flamingock.internal.common.core.transaction.TransactionWrapper; +import io.flamingock.internal.common.core.external.ExecutionWrapper; import java.util.concurrent.atomic.AtomicReference; import java.util.function.Function; -public class CouchbaseTxWrapper implements TransactionWrapper { +public class CouchbaseTxWrapper implements ExecutionWrapper { private final Cluster cluster; private final TransactionManager txManager; @@ -42,7 +42,7 @@ public CouchbaseTxWrapper(Cluster cluster, TransactionManager RESULT wrapInTransaction(CONTEXT executionContext, Function operation) { + public RESULT wrapExecution(CONTEXT executionContext, Function operation) { String sessionId = executionContext.getSessionId(); AtomicReference resultRef = new AtomicReference<>(); diff --git a/core/target-systems/flamingock-dynamodb-externalsystem-api/src/main/java/io/flamingock/externalsystem/dynamodb/api/DynamoDBExternalSystem.java b/core/target-systems/flamingock-dynamodb-externalsystem-api/src/main/java/io/flamingock/externalsystem/dynamodb/api/DynamoDBExternalSystem.java index 641fdf7eb..ececeefb9 100644 --- a/core/target-systems/flamingock-dynamodb-externalsystem-api/src/main/java/io/flamingock/externalsystem/dynamodb/api/DynamoDBExternalSystem.java +++ b/core/target-systems/flamingock-dynamodb-externalsystem-api/src/main/java/io/flamingock/externalsystem/dynamodb/api/DynamoDBExternalSystem.java @@ -15,7 +15,7 @@ */ package io.flamingock.externalsystem.dynamodb.api; -import io.flamingock.internal.common.core.transaction.TransactionalExternalSystem; +import io.flamingock.internal.common.core.external.TransactionalExternalSystem; import software.amazon.awssdk.services.dynamodb.DynamoDbClient; public interface DynamoDBExternalSystem extends TransactionalExternalSystem { diff --git a/core/target-systems/flamingock-dynamodb-targetsystem/src/main/java/io/flamingock/targetsystem/dynamodb/DynamoDBTargetSystem.java b/core/target-systems/flamingock-dynamodb-targetsystem/src/main/java/io/flamingock/targetsystem/dynamodb/DynamoDBTargetSystem.java index 237bd4fb6..f69e85f76 100644 --- a/core/target-systems/flamingock-dynamodb-targetsystem/src/main/java/io/flamingock/targetsystem/dynamodb/DynamoDBTargetSystem.java +++ b/core/target-systems/flamingock-dynamodb-targetsystem/src/main/java/io/flamingock/targetsystem/dynamodb/DynamoDBTargetSystem.java @@ -25,7 +25,7 @@ import io.flamingock.internal.core.external.targets.TransactionalTargetSystem; import io.flamingock.internal.core.external.targets.mark.NoOpTargetSystemAuditMarker; import io.flamingock.internal.core.transaction.TransactionManager; -import io.flamingock.internal.common.core.transaction.TransactionWrapper; +import io.flamingock.internal.common.core.external.ExecutionWrapper; import software.amazon.awssdk.enhanced.dynamodb.model.TransactWriteItemsEnhancedRequest; import software.amazon.awssdk.services.dynamodb.DynamoDbClient; @@ -85,7 +85,7 @@ protected DynamoDBTargetSystem getSelf() { } @Override - public TransactionWrapper getTxWrapper() { + public ExecutionWrapper getTxWrapper() { return txWrapper; } diff --git a/core/target-systems/flamingock-dynamodb-targetsystem/src/main/java/io/flamingock/targetsystem/dynamodb/DynamoDBTxWrapper.java b/core/target-systems/flamingock-dynamodb-targetsystem/src/main/java/io/flamingock/targetsystem/dynamodb/DynamoDBTxWrapper.java index 73233e690..48bacf231 100644 --- a/core/target-systems/flamingock-dynamodb-targetsystem/src/main/java/io/flamingock/targetsystem/dynamodb/DynamoDBTxWrapper.java +++ b/core/target-systems/flamingock-dynamodb-targetsystem/src/main/java/io/flamingock/targetsystem/dynamodb/DynamoDBTxWrapper.java @@ -20,7 +20,7 @@ import io.flamingock.internal.common.core.error.DatabaseTransactionException; import io.flamingock.internal.core.transaction.TransactionManager; import io.flamingock.internal.core.change.navigation.step.FailedStep; -import io.flamingock.internal.common.core.transaction.TransactionWrapper; +import io.flamingock.internal.common.core.external.ExecutionWrapper; import io.flamingock.internal.util.dynamodb.DynamoDBUtil; import io.flamingock.internal.util.log.FlamingockLoggerFactory; import org.slf4j.Logger; @@ -33,7 +33,7 @@ import java.util.function.Function; import java.util.stream.Collectors; -public class DynamoDBTxWrapper implements TransactionWrapper { +public class DynamoDBTxWrapper implements ExecutionWrapper { private static final Logger logger = FlamingockLoggerFactory.getLogger("DynamoTx"); private final TransactionManager txManager; @@ -72,7 +72,7 @@ private String formatDuration(Duration duration) { } @Override - public RESULT wrapInTransaction(CONTEXT executionContext, Function operation) { + public RESULT wrapExecution(CONTEXT executionContext, Function operation) { LocalDateTime transactionStart = LocalDateTime.now(); String sessionId = executionContext.getSessionId(); TransactWriteItemsEnhancedRequest.Builder writeRequestBuilder = txManager.startSession(sessionId); diff --git a/core/target-systems/flamingock-mongodb-externalsystem-api/src/main/java/io/flamingock/externalsystem/mongodb/api/MongoDBExternalSystem.java b/core/target-systems/flamingock-mongodb-externalsystem-api/src/main/java/io/flamingock/externalsystem/mongodb/api/MongoDBExternalSystem.java index 1f284c536..c2de633bb 100644 --- a/core/target-systems/flamingock-mongodb-externalsystem-api/src/main/java/io/flamingock/externalsystem/mongodb/api/MongoDBExternalSystem.java +++ b/core/target-systems/flamingock-mongodb-externalsystem-api/src/main/java/io/flamingock/externalsystem/mongodb/api/MongoDBExternalSystem.java @@ -16,8 +16,7 @@ package io.flamingock.externalsystem.mongodb.api; import com.mongodb.client.MongoDatabase; -import io.flamingock.api.external.ExternalSystem; -import io.flamingock.internal.common.core.transaction.TransactionalExternalSystem; +import io.flamingock.internal.common.core.external.TransactionalExternalSystem; public interface MongoDBExternalSystem extends TransactionalExternalSystem { diff --git a/core/target-systems/flamingock-mongodb-reactive-externalsystem-api/src/main/java/io/flamingock/externalsystem/mongodb/reactive/api/MongoDBReactiveExternalSystem.java b/core/target-systems/flamingock-mongodb-reactive-externalsystem-api/src/main/java/io/flamingock/externalsystem/mongodb/reactive/api/MongoDBReactiveExternalSystem.java index 459fb3b70..e482dd945 100644 --- a/core/target-systems/flamingock-mongodb-reactive-externalsystem-api/src/main/java/io/flamingock/externalsystem/mongodb/reactive/api/MongoDBReactiveExternalSystem.java +++ b/core/target-systems/flamingock-mongodb-reactive-externalsystem-api/src/main/java/io/flamingock/externalsystem/mongodb/reactive/api/MongoDBReactiveExternalSystem.java @@ -16,7 +16,7 @@ package io.flamingock.externalsystem.mongodb.reactive.api; import com.mongodb.reactivestreams.client.MongoDatabase; -import io.flamingock.internal.common.core.transaction.TransactionalExternalSystem; +import io.flamingock.internal.common.core.external.TransactionalExternalSystem; public interface MongoDBReactiveExternalSystem extends TransactionalExternalSystem { diff --git a/core/target-systems/flamingock-mongodb-reactive-targetsystem/src/main/java/io/flamingock/targetsystem/mongodb/reactive/MongoDBReactiveTargetSystem.java b/core/target-systems/flamingock-mongodb-reactive-targetsystem/src/main/java/io/flamingock/targetsystem/mongodb/reactive/MongoDBReactiveTargetSystem.java index 8a89d074b..8eac793d2 100644 --- a/core/target-systems/flamingock-mongodb-reactive-targetsystem/src/main/java/io/flamingock/targetsystem/mongodb/reactive/MongoDBReactiveTargetSystem.java +++ b/core/target-systems/flamingock-mongodb-reactive-targetsystem/src/main/java/io/flamingock/targetsystem/mongodb/reactive/MongoDBReactiveTargetSystem.java @@ -27,7 +27,7 @@ import io.flamingock.internal.common.core.audit.AuditReaderType; import io.flamingock.internal.common.core.context.ContextResolver; import io.flamingock.internal.common.core.error.FlamingockException; -import io.flamingock.internal.common.core.transaction.TransactionWrapper; +import io.flamingock.internal.common.core.external.ExecutionWrapper; import io.flamingock.internal.core.builder.FlamingockEdition; import io.flamingock.internal.core.external.targets.TransactionalTargetSystem; import io.flamingock.internal.core.external.targets.mark.NoOpTargetSystemAuditMarker; @@ -174,7 +174,7 @@ protected MongoDBReactiveTargetSystem getSelf() { } @Override - public TransactionWrapper getTxWrapper() { + public ExecutionWrapper getTxWrapper() { return getReactiveTxWrapper(); } diff --git a/core/target-systems/flamingock-mongodb-reactive-targetsystem/src/main/java/io/flamingock/targetsystem/mongodb/reactive/MongoDBReactiveTxWrapper.java b/core/target-systems/flamingock-mongodb-reactive-targetsystem/src/main/java/io/flamingock/targetsystem/mongodb/reactive/MongoDBReactiveTxWrapper.java index 8b929d02a..5563919d8 100644 --- a/core/target-systems/flamingock-mongodb-reactive-targetsystem/src/main/java/io/flamingock/targetsystem/mongodb/reactive/MongoDBReactiveTxWrapper.java +++ b/core/target-systems/flamingock-mongodb-reactive-targetsystem/src/main/java/io/flamingock/targetsystem/mongodb/reactive/MongoDBReactiveTxWrapper.java @@ -20,7 +20,7 @@ import io.flamingock.internal.common.core.context.Dependency; import io.flamingock.internal.common.core.context.RuntimeContext; import io.flamingock.internal.common.core.error.DatabaseTransactionException; -import io.flamingock.internal.common.core.transaction.TransactionWrapper; +import io.flamingock.internal.common.core.external.ExecutionWrapper; import io.flamingock.internal.core.change.navigation.step.FailedStep; import io.flamingock.internal.core.transaction.TransactionManager; import io.flamingock.internal.util.log.FlamingockLoggerFactory; @@ -31,7 +31,7 @@ import java.time.LocalDateTime; import java.util.function.Function; -public class MongoDBReactiveTxWrapper implements TransactionWrapper { +public class MongoDBReactiveTxWrapper implements ExecutionWrapper { private static final Logger logger = FlamingockLoggerFactory.getLogger("MongoReactiveTx"); private final TransactionManager sessionManager; @@ -45,7 +45,7 @@ TransactionManager getTxManager() { } @Override - public RESULT wrapInTransaction(CONTEXT executionContext, Function operation) { + public RESULT wrapExecution(CONTEXT executionContext, Function operation) { LocalDateTime transactionStart = LocalDateTime.now(); String sessionId = executionContext.getSessionId(); ClientSession clientSession = sessionManager.startSession(sessionId); diff --git a/core/target-systems/flamingock-mongodb-springdata-reactive-targetsystem/src/main/java/io/flamingock/targetsystem/mongodb/springdata/reactive/MongoDBSpringDataReactiveTargetSystem.java b/core/target-systems/flamingock-mongodb-springdata-reactive-targetsystem/src/main/java/io/flamingock/targetsystem/mongodb/springdata/reactive/MongoDBSpringDataReactiveTargetSystem.java index 1eb4afefc..9e92c77ae 100644 --- a/core/target-systems/flamingock-mongodb-springdata-reactive-targetsystem/src/main/java/io/flamingock/targetsystem/mongodb/springdata/reactive/MongoDBSpringDataReactiveTargetSystem.java +++ b/core/target-systems/flamingock-mongodb-springdata-reactive-targetsystem/src/main/java/io/flamingock/targetsystem/mongodb/springdata/reactive/MongoDBSpringDataReactiveTargetSystem.java @@ -25,7 +25,7 @@ import io.flamingock.internal.common.core.audit.AuditReaderType; import io.flamingock.internal.common.core.context.ContextResolver; import io.flamingock.internal.common.core.error.FlamingockException; -import io.flamingock.internal.common.core.transaction.TransactionWrapper; +import io.flamingock.internal.common.core.external.ExecutionWrapper; import io.flamingock.internal.core.builder.FlamingockEdition; import io.flamingock.internal.core.external.targets.TransactionalTargetSystem; import io.flamingock.internal.core.external.targets.mark.NoOpTargetSystemAuditMarker; @@ -159,7 +159,7 @@ protected MongoDBSpringDataReactiveTargetSystem getSelf() { } @Override - public TransactionWrapper getTxWrapper() { + public ExecutionWrapper getTxWrapper() { if (!supportsTransactions()) { throw new FlamingockException("Transaction wrapper requested for a MongoDB target that does not support transactions."); } diff --git a/core/target-systems/flamingock-mongodb-springdata-reactive-targetsystem/src/main/java/io/flamingock/targetsystem/mongodb/springdata/reactive/MongoDBSpringDataReactiveTxWrapper.java b/core/target-systems/flamingock-mongodb-springdata-reactive-targetsystem/src/main/java/io/flamingock/targetsystem/mongodb/springdata/reactive/MongoDBSpringDataReactiveTxWrapper.java index c97bdc59e..2835c1f68 100644 --- a/core/target-systems/flamingock-mongodb-springdata-reactive-targetsystem/src/main/java/io/flamingock/targetsystem/mongodb/springdata/reactive/MongoDBSpringDataReactiveTxWrapper.java +++ b/core/target-systems/flamingock-mongodb-springdata-reactive-targetsystem/src/main/java/io/flamingock/targetsystem/mongodb/springdata/reactive/MongoDBSpringDataReactiveTxWrapper.java @@ -24,7 +24,7 @@ import io.flamingock.internal.common.core.context.Dependency; import io.flamingock.internal.common.core.context.RuntimeContext; import io.flamingock.internal.common.core.error.DatabaseTransactionException; -import io.flamingock.internal.common.core.transaction.TransactionWrapper; +import io.flamingock.internal.common.core.external.ExecutionWrapper; import io.flamingock.internal.core.change.navigation.step.FailedStep; import io.flamingock.internal.util.log.FlamingockLoggerFactory; import org.slf4j.Logger; @@ -34,7 +34,7 @@ import java.time.LocalDateTime; import java.util.function.Function; -public class MongoDBSpringDataReactiveTxWrapper implements TransactionWrapper { +public class MongoDBSpringDataReactiveTxWrapper implements ExecutionWrapper { private static final Logger logger = FlamingockLoggerFactory.getLogger("SpringMongoReactiveTx"); private final ReactiveMongoTemplate mongoTemplate; @@ -46,7 +46,7 @@ private MongoDBSpringDataReactiveTxWrapper(ReactiveMongoTemplate mongoTemplate, } @Override - public RESULT wrapInTransaction(CONTEXT executionContext, Function operation) { + public RESULT wrapExecution(CONTEXT executionContext, Function operation) { LocalDateTime transactionStart = LocalDateTime.now(); ClientSession clientSession = null; try { diff --git a/core/target-systems/flamingock-mongodb-springdata-targetsystem/src/main/java/io/flamingock/targetsystem/mongodb/springdata/MongoDBSpringDataTargetSystem.java b/core/target-systems/flamingock-mongodb-springdata-targetsystem/src/main/java/io/flamingock/targetsystem/mongodb/springdata/MongoDBSpringDataTargetSystem.java index b21f0085b..57bdf1936 100644 --- a/core/target-systems/flamingock-mongodb-springdata-targetsystem/src/main/java/io/flamingock/targetsystem/mongodb/springdata/MongoDBSpringDataTargetSystem.java +++ b/core/target-systems/flamingock-mongodb-springdata-targetsystem/src/main/java/io/flamingock/targetsystem/mongodb/springdata/MongoDBSpringDataTargetSystem.java @@ -27,7 +27,7 @@ import io.flamingock.internal.core.builder.FlamingockEdition; import io.flamingock.internal.core.external.targets.mark.NoOpTargetSystemAuditMarker; import io.flamingock.internal.core.external.targets.TransactionalTargetSystem; -import io.flamingock.internal.common.core.transaction.TransactionWrapper; +import io.flamingock.internal.common.core.external.ExecutionWrapper; import io.flamingock.externalsystem.mongodb.api.MongoDBExternalSystem; import org.springframework.data.mongodb.core.MongoTemplate; @@ -160,7 +160,7 @@ protected MongoDBSpringDataTargetSystem getSelf() { } @Override - public TransactionWrapper getTxWrapper() { + public ExecutionWrapper getTxWrapper() { if (!supportsTransactions()) { throw new FlamingockException("Transaction wrapper requested for a MongoDB target that does not support transactions."); } diff --git a/core/target-systems/flamingock-mongodb-springdata-targetsystem/src/main/java/io/flamingock/targetsystem/mongodb/springdata/MongoDBSpringDataTxWrapper.java b/core/target-systems/flamingock-mongodb-springdata-targetsystem/src/main/java/io/flamingock/targetsystem/mongodb/springdata/MongoDBSpringDataTxWrapper.java index 4efe04424..84298fa3f 100644 --- a/core/target-systems/flamingock-mongodb-springdata-targetsystem/src/main/java/io/flamingock/targetsystem/mongodb/springdata/MongoDBSpringDataTxWrapper.java +++ b/core/target-systems/flamingock-mongodb-springdata-targetsystem/src/main/java/io/flamingock/targetsystem/mongodb/springdata/MongoDBSpringDataTxWrapper.java @@ -22,7 +22,7 @@ import io.flamingock.internal.common.core.context.RuntimeContext; import io.flamingock.internal.common.core.error.DatabaseTransactionException; import io.flamingock.internal.core.change.navigation.step.FailedStep; -import io.flamingock.internal.common.core.transaction.TransactionWrapper; +import io.flamingock.internal.common.core.external.ExecutionWrapper; import io.flamingock.internal.util.log.FlamingockLoggerFactory; import org.slf4j.Logger; import org.springframework.data.mongodb.MongoDatabaseFactory; @@ -35,7 +35,7 @@ import java.time.LocalDateTime; import java.util.function.Function; -public class MongoDBSpringDataTxWrapper implements TransactionWrapper { +public class MongoDBSpringDataTxWrapper implements ExecutionWrapper { private static final Logger logger = FlamingockLoggerFactory.getLogger("SpringMongoTx"); private final TransactionTemplate txTemplate; @@ -55,7 +55,7 @@ public static Builder builder() { } @Override - public RESULT wrapInTransaction(CONTEXT executionContext, Function operation) { + public RESULT wrapExecution(CONTEXT executionContext, Function operation) { LocalDateTime transactionStart = LocalDateTime.now(); try { diff --git a/core/target-systems/flamingock-mongodb-sync-targetsystem/src/main/java/io/flamingock/targetsystem/mongodb/sync/MongoDBSyncTargetSystem.java b/core/target-systems/flamingock-mongodb-sync-targetsystem/src/main/java/io/flamingock/targetsystem/mongodb/sync/MongoDBSyncTargetSystem.java index 8e8b80c92..ab59a1120 100644 --- a/core/target-systems/flamingock-mongodb-sync-targetsystem/src/main/java/io/flamingock/targetsystem/mongodb/sync/MongoDBSyncTargetSystem.java +++ b/core/target-systems/flamingock-mongodb-sync-targetsystem/src/main/java/io/flamingock/targetsystem/mongodb/sync/MongoDBSyncTargetSystem.java @@ -31,7 +31,7 @@ import io.flamingock.internal.core.external.targets.TransactionalTargetSystem; import io.flamingock.internal.core.external.targets.mark.NoOpTargetSystemAuditMarker; import io.flamingock.internal.core.transaction.TransactionManager; -import io.flamingock.internal.common.core.transaction.TransactionWrapper; +import io.flamingock.internal.common.core.external.ExecutionWrapper; import java.util.Objects; import java.util.Optional; @@ -199,7 +199,7 @@ protected MongoDBSyncTargetSystem getSelf() { } @Override - public TransactionWrapper getTxWrapper() { + public ExecutionWrapper getTxWrapper() { return getSyncTxWrapper(); } diff --git a/core/target-systems/flamingock-mongodb-sync-targetsystem/src/main/java/io/flamingock/targetsystem/mongodb/sync/MongoDBSyncTxWrapper.java b/core/target-systems/flamingock-mongodb-sync-targetsystem/src/main/java/io/flamingock/targetsystem/mongodb/sync/MongoDBSyncTxWrapper.java index 87ce6afe5..3b3c8ea27 100644 --- a/core/target-systems/flamingock-mongodb-sync-targetsystem/src/main/java/io/flamingock/targetsystem/mongodb/sync/MongoDBSyncTxWrapper.java +++ b/core/target-systems/flamingock-mongodb-sync-targetsystem/src/main/java/io/flamingock/targetsystem/mongodb/sync/MongoDBSyncTxWrapper.java @@ -21,7 +21,7 @@ import io.flamingock.internal.common.core.context.RuntimeContext; import io.flamingock.internal.common.core.error.DatabaseTransactionException; import io.flamingock.internal.core.change.navigation.step.FailedStep; -import io.flamingock.internal.common.core.transaction.TransactionWrapper; +import io.flamingock.internal.common.core.external.ExecutionWrapper; import io.flamingock.internal.core.transaction.TransactionManager; import io.flamingock.internal.util.log.FlamingockLoggerFactory; import org.slf4j.Logger; @@ -30,7 +30,7 @@ import java.time.LocalDateTime; import java.util.function.Function; -public class MongoDBSyncTxWrapper implements TransactionWrapper { +public class MongoDBSyncTxWrapper implements ExecutionWrapper { private static final Logger logger = FlamingockLoggerFactory.getLogger("MongoTx"); private final TransactionManager sessionManager; @@ -65,7 +65,7 @@ private String formatDuration(Duration duration) { } @Override - public RESULT wrapInTransaction(CONTEXT executionContext, Function operation) { + public RESULT wrapExecution(CONTEXT executionContext, Function operation) { LocalDateTime transactionStart = LocalDateTime.now(); String sessionId = executionContext.getSessionId(); Dependency clienteSessionDependency; diff --git a/core/target-systems/flamingock-sql-externalsystem-api/src/main/java/io/flamingock/externalsystem/sql/api/SqlExternalSystem.java b/core/target-systems/flamingock-sql-externalsystem-api/src/main/java/io/flamingock/externalsystem/sql/api/SqlExternalSystem.java index e37bd332b..6812915bc 100644 --- a/core/target-systems/flamingock-sql-externalsystem-api/src/main/java/io/flamingock/externalsystem/sql/api/SqlExternalSystem.java +++ b/core/target-systems/flamingock-sql-externalsystem-api/src/main/java/io/flamingock/externalsystem/sql/api/SqlExternalSystem.java @@ -15,7 +15,7 @@ */ package io.flamingock.externalsystem.sql.api; -import io.flamingock.internal.common.core.transaction.TransactionalExternalSystem; +import io.flamingock.internal.common.core.external.TransactionalExternalSystem; import javax.sql.DataSource; diff --git a/core/target-systems/flamingock-sql-targetsystem/src/main/java/io/flamingock/targetsystem/sql/SqlTargetSystem.java b/core/target-systems/flamingock-sql-targetsystem/src/main/java/io/flamingock/targetsystem/sql/SqlTargetSystem.java index 5f633e7c5..12d26a182 100644 --- a/core/target-systems/flamingock-sql-targetsystem/src/main/java/io/flamingock/targetsystem/sql/SqlTargetSystem.java +++ b/core/target-systems/flamingock-sql-targetsystem/src/main/java/io/flamingock/targetsystem/sql/SqlTargetSystem.java @@ -24,7 +24,7 @@ import io.flamingock.internal.core.external.targets.mark.NoOpTargetSystemAuditMarker; import io.flamingock.internal.core.runtime.ExecutionRuntime; import io.flamingock.internal.core.transaction.TransactionManager; -import io.flamingock.internal.common.core.transaction.TransactionWrapper; +import io.flamingock.internal.common.core.external.ExecutionWrapper; import javax.sql.DataSource; import java.sql.Connection; @@ -37,11 +37,24 @@ public class SqlTargetSystem extends TransactionalTargetSystem private final DataSource dataSource; + private final ExecutionWrapper nonTxWrapper; + private SqlTxWrapper txWrapper; public SqlTargetSystem(String id, DataSource dataSource) { super(id); this.dataSource = dataSource; + this.nonTxWrapper = new ExecutionWrapper() { + @Override + public RESULT wrapExecution(CONTEXT runtimeContext, Function operation) { + try (Connection connection = dataSource.getConnection()) { + runtimeContext.addDependency(connection); + return operation.apply(runtimeContext); + } catch (SQLException e) { + throw new FlamingockException(e); + } + } + }; } @Override @@ -77,27 +90,27 @@ protected SqlTargetSystem getSelf() { } @Override - public TransactionWrapper getTxWrapper() { + public ExecutionWrapper getTxWrapper() { return txWrapper; } /** - * Acquires one connection for a non-transactional callback, injects it into the runtime, - * and closes that same connection after the callback completes or fails. - * - * @param changeFunc the callback to execute - * @param executionRuntime the runtime to receive the connection dependency - * @param the callback return type - * @return the callback result + * Non-transactional execution needs a JDBC {@link Connection} of its own: the change still has SQL to + * run, it just runs outside a transaction. This wrapper borrows one connection from the + * {@link DataSource} for the duration of the call, publishes it into the runtime so the change can + * resolve it, and closes it afterwards — including when the operation throws. + *

+ * Auto-commit is left untouched, at whatever the {@code DataSource} hands back — so under the usual + * auto-commit default each statement commits on its own. That is the point of this path: a change + * declared non-transactional gets no rollback, and a failure partway through leaves the statements + * that already ran in place. + *

+ * The transactional path does not go through here — see {@link #getTxWrapper()}, whose connection is + * owned by the {@code TransactionManager} and spans the whole transaction. */ @Override - protected T nonTxWrapper(Function changeFunc, ExecutionRuntime executionRuntime) { - try (Connection connection = dataSource.getConnection()) { - executionRuntime.addDependency(connection); - return changeFunc.apply(executionRuntime); - } catch (SQLException e) { - throw new FlamingockException(e); - } + protected ExecutionWrapper getNonTxWrapper() { + return nonTxWrapper; } private SqlTxWrapper createTxWrapper(TransactionManager txManager) { diff --git a/core/target-systems/flamingock-sql-targetsystem/src/main/java/io/flamingock/targetsystem/sql/SqlTxWrapper.java b/core/target-systems/flamingock-sql-targetsystem/src/main/java/io/flamingock/targetsystem/sql/SqlTxWrapper.java index 74edf2195..48f393a54 100644 --- a/core/target-systems/flamingock-sql-targetsystem/src/main/java/io/flamingock/targetsystem/sql/SqlTxWrapper.java +++ b/core/target-systems/flamingock-sql-targetsystem/src/main/java/io/flamingock/targetsystem/sql/SqlTxWrapper.java @@ -20,7 +20,7 @@ import io.flamingock.internal.common.core.error.DatabaseTransactionException; import io.flamingock.internal.core.transaction.TransactionManager; import io.flamingock.internal.core.change.navigation.step.FailedStep; -import io.flamingock.internal.common.core.transaction.TransactionWrapper; +import io.flamingock.internal.common.core.external.ExecutionWrapper; import io.flamingock.internal.util.log.FlamingockLoggerFactory; import org.slf4j.Logger; @@ -30,7 +30,7 @@ import java.time.LocalDateTime; import java.util.function.Function; -public class SqlTxWrapper implements TransactionWrapper { +public class SqlTxWrapper implements ExecutionWrapper { private static final Logger logger = FlamingockLoggerFactory.getLogger("SqlTx"); private final TransactionManager txManager; @@ -69,7 +69,7 @@ private String formatDuration(Duration duration) { } @Override - public RESULT wrapInTransaction(CONTEXT executionContext, Function operation) { + public RESULT wrapExecution(CONTEXT executionContext, Function operation) { LocalDateTime transactionStart = LocalDateTime.now(); String sessionId = executionContext.getSessionId(); diff --git a/core/target-systems/flamingock-sql-targetsystem/src/test/java/io/flamingock/targetsystem/sql/SqlTxWrapperLifecycleTest.java b/core/target-systems/flamingock-sql-targetsystem/src/test/java/io/flamingock/targetsystem/sql/SqlTxWrapperLifecycleTest.java index c35971450..df5df99fd 100644 --- a/core/target-systems/flamingock-sql-targetsystem/src/test/java/io/flamingock/targetsystem/sql/SqlTxWrapperLifecycleTest.java +++ b/core/target-systems/flamingock-sql-targetsystem/src/test/java/io/flamingock/targetsystem/sql/SqlTxWrapperLifecycleTest.java @@ -60,7 +60,7 @@ void setUp() throws SQLException { void commitsSuccessfulCallback() throws Exception { BasicRuntimeContext context = new BasicRuntimeContext("success"); - String result = txWrapper.wrapInTransaction(context, runtimeContext -> { + String result = txWrapper.wrapExecution(context, runtimeContext -> { insert(runtimeContext, 1, "committed"); return "success"; }); @@ -75,7 +75,7 @@ void rollsBackFailedStepValue() throws Exception { FailedStep failedStep = mock(FailedStep.class); BasicRuntimeContext context = new BasicRuntimeContext("failed-value"); - FailedStep result = txWrapper.wrapInTransaction(context, runtimeContext -> { + FailedStep result = txWrapper.wrapExecution(context, runtimeContext -> { insert(runtimeContext, 2, "rolled-back-value"); return failedStep; }); @@ -90,7 +90,7 @@ void rollsBackCallbackException() throws Exception { BasicRuntimeContext context = new BasicRuntimeContext("exception"); DatabaseTransactionException exception = assertThrows(DatabaseTransactionException.class, - () -> txWrapper.wrapInTransaction(context, runtimeContext -> { + () -> txWrapper.wrapExecution(context, runtimeContext -> { insert(runtimeContext, 3, "rolled-back-exception"); throw new IllegalStateException("callback failed"); })); @@ -105,11 +105,11 @@ void reusesSessionIdWithFreshConnection() throws Exception { BasicRuntimeContext firstContext = new BasicRuntimeContext("reused-session"); BasicRuntimeContext secondContext = new BasicRuntimeContext("reused-session"); - txWrapper.wrapInTransaction(firstContext, runtimeContext -> { + txWrapper.wrapExecution(firstContext, runtimeContext -> { insert(runtimeContext, 4, "first-transaction"); return "first"; }); - txWrapper.wrapInTransaction(secondContext, runtimeContext -> { + txWrapper.wrapExecution(secondContext, runtimeContext -> { insert(runtimeContext, 5, "second-transaction"); return "second"; }); @@ -124,11 +124,11 @@ void reusesSessionIdAfterFailedStepRollback() throws Exception { BasicRuntimeContext successfulContext = new BasicRuntimeContext("failed-then-reused"); FailedStep failedStep = mock(FailedStep.class); - txWrapper.wrapInTransaction(failedContext, runtimeContext -> { + txWrapper.wrapExecution(failedContext, runtimeContext -> { insert(runtimeContext, 6, "discarded"); return failedStep; }); - txWrapper.wrapInTransaction(successfulContext, runtimeContext -> { + txWrapper.wrapExecution(successfulContext, runtimeContext -> { insert(runtimeContext, 7, "committed-after-failure"); return "success"; }); diff --git a/docs/CHANGE-STEP-NAVIGATION.md b/docs/CHANGE-STEP-NAVIGATION.md index 8354c7186..3601a659f 100644 --- a/docs/CHANGE-STEP-NAVIGATION.md +++ b/docs/CHANGE-STEP-NAVIGATION.md @@ -32,7 +32,7 @@ graph TB %% Transactional Path CHECK_TRANSACTIONAL -->|Yes| CLOUD_CHECK{Cloud Edition?} CLOUD_CHECK -->|Yes| SET_ONGOING[Set OngoingStatus.EXECUTION] - CLOUD_CHECK -->|No| TRANSACTION_WRAPPER[TransactionWrapper] + CLOUD_CHECK -->|No| TRANSACTION_WRAPPER[ExecutionWrapper - getTxWrapper] SET_ONGOING --> TRANSACTION_WRAPPER TRANSACTION_WRAPPER --> EXECUTE_IN_TRANSACTION[executeChange in Transaction] @@ -131,7 +131,7 @@ graph TB 3. **ExecutionStep** → **AfterExecutionAuditStep** variants (via `applyAuditResult()`) ### Transaction Handling -- **Transactional changes** go through `TransactionWrapper` +- **Transactional changes** go through `ExecutionWrapper` - **Cloud edition** tracks ongoing status during execution - **Auto-rollback** occurs for transactional failures - **Manual rollback** handles non-transactional failures diff --git a/docs/EXECUTION_FLOW_GUIDE.md b/docs/EXECUTION_FLOW_GUIDE.md index 8affe08d1..9d70af4ef 100644 --- a/docs/EXECUTION_FLOW_GUIDE.md +++ b/docs/EXECUTION_FLOW_GUIDE.md @@ -218,7 +218,7 @@ Flamingock's transaction handling is sophisticated and context-aware. #### Transaction Decision Logic A change executes within a transaction when **all** conditions are met: -1. **TransactionWrapper available** - AuditStore supports transactions +1. **Transactional `ExecutionWrapper` available** - the target system exposes `getTxWrapper()` and reports `supportsTransactions()` 2. **Change configured as transactional** - `@Change(transactional = true)` (default) 3. **Database supports transactions** - Not all databases/operations are transactional