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());
     }

Reply via email to