This is an automated email from the ASF dual-hosted git repository.
exceptionfactory pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/nifi.git
The following commit(s) were added to refs/heads/main by this push:
new 96737bbd588 NIFI-16181 Adjust MockProcessSession.commitAsync() to
allow null callbacks #11524)
96737bbd588 is described below
commit 96737bbd588dd6fea2b6624c9cffae55557fce72
Author: Alaksiej Ščarbaty <[email protected]>
AuthorDate: Mon Aug 10 16:09:42 2026 +0200
NIFI-16181 Adjust MockProcessSession.commitAsync() to allow null callbacks
#11524)
Signed-off-by: David Handermann <[email protected]>
---
.../org/apache/nifi/util/MockProcessSession.java | 8 ++-
.../apache/nifi/util/TestMockProcessSession.java | 64 ++++++++++++++++++++++
2 files changed, 70 insertions(+), 2 deletions(-)
diff --git
a/nifi-mock/src/main/java/org/apache/nifi/util/MockProcessSession.java
b/nifi-mock/src/main/java/org/apache/nifi/util/MockProcessSession.java
index 1f4146ac8c0..ee2b10a6f65 100644
--- a/nifi-mock/src/main/java/org/apache/nifi/util/MockProcessSession.java
+++ b/nifi-mock/src/main/java/org/apache/nifi/util/MockProcessSession.java
@@ -370,11 +370,15 @@ public class MockProcessSession implements ProcessSession
{
commitInternal();
} catch (final Throwable t) {
rollback();
- onFailure.accept(t);
+ if (onFailure != null) {
+ onFailure.accept(t);
+ }
throw t;
}
- onSuccess.run();
+ if (onSuccess != null) {
+ onSuccess.run();
+ }
}
/**
diff --git
a/nifi-mock/src/test/java/org/apache/nifi/util/TestMockProcessSession.java
b/nifi-mock/src/test/java/org/apache/nifi/util/TestMockProcessSession.java
index 4e51132b201..4eba9e0ca43 100644
--- a/nifi-mock/src/test/java/org/apache/nifi/util/TestMockProcessSession.java
+++ b/nifi-mock/src/test/java/org/apache/nifi/util/TestMockProcessSession.java
@@ -40,6 +40,7 @@ import java.util.Collection;
import java.util.List;
import java.util.Map;
import java.util.Set;
+import java.util.concurrent.atomic.AtomicBoolean;
import java.util.concurrent.atomic.AtomicLong;
import java.util.regex.Pattern;
@@ -602,6 +603,69 @@ public class TestMockProcessSession {
}
}
+ @Nested
+ class RegardingCommitAsync {
+
+ @Test
+ void invokesOnSuccess() {
+ final AtomicBoolean successInvoked = new AtomicBoolean();
+
+ session.commitAsync(() -> successInvoked.set(true), failure ->
fail("onFailure should not be invoked"));
+
+ assertTrue(successInvoked.get());
+ session.assertCommitted();
+ }
+
+ @Test
+ void allowsNullCallbacks() {
+ assertDoesNotThrow(() -> session.commitAsync(null, null));
+
+ session.assertCommitted();
+ }
+
+ @Test
+ void invokesOnFailure() {
+ final MockProcessSession failingSession =
MockProcessSession.builder(sharedState, processor)
+ .stateManager(stateManager)
+ .failCommit()
+ .build();
+ final AtomicBoolean failureInvoked = new AtomicBoolean();
+
+ final FlowFileHandlingException thrown =
assertThrows(FlowFileHandlingException.class,
+ () -> failingSession.commitAsync(
+ () -> fail("onSuccess should not be invoked"),
+ failure -> failureInvoked.set(true)));
+
+ assertTrue(failureInvoked.get());
+ assertEquals("Cannot commit session because the session was
requested to fail by a test", thrown.getMessage());
+ failingSession.assertNotCommitted();
+ }
+
+ @Test
+ void throwsWithNullOnFailure() {
+ final MockProcessSession failingSession =
MockProcessSession.builder(sharedState, processor)
+ .stateManager(stateManager)
+ .failCommit()
+ .build();
+
+ assertThrows(FlowFileHandlingException.class, () ->
failingSession.commitAsync(() -> { }, null));
+
+ failingSession.assertNotCommitted();
+ }
+
+ @Test
+ void throwsWithSingleCallback() {
+ final MockProcessSession failingSession =
MockProcessSession.builder(sharedState, processor)
+ .stateManager(stateManager)
+ .failCommit()
+ .build();
+
+ assertThrows(FlowFileHandlingException.class, () ->
failingSession.commitAsync(() -> { }));
+
+ failingSession.assertNotCommitted();
+ }
+ }
+
private MockProcessSession createMockProcessSession() {
return createMockProcessSession(new TestProcessor());
}