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 b83e7b396d2 NIFI-16215 Fixed including current transaction queues in 
Connection Status Events (#11557)
b83e7b396d2 is described below

commit b83e7b396d261c66f5ee8daa70e745b5ce26974a
Author: boguszj <[email protected]>
AuthorDate: Mon Aug 17 20:48:28 2026 +0200

    NIFI-16215 Fixed including current transaction queues in Connection Status 
Events (#11557)
    
    Signed-off-by: David Handermann <[email protected]>
---
 .../apache/nifi/controller/repository/StandardProcessSession.java   | 5 +++--
 .../nifi/controller/repository/StandardProcessSessionTest.java      | 6 ++++++
 2 files changed, 9 insertions(+), 2 deletions(-)

diff --git 
a/nifi-framework-bundle/nifi-framework/nifi-framework-components/src/main/java/org/apache/nifi/controller/repository/StandardProcessSession.java
 
b/nifi-framework-bundle/nifi-framework/nifi-framework-components/src/main/java/org/apache/nifi/controller/repository/StandardProcessSession.java
index d7061d7d9b1..bb4a9ddf0a9 100644
--- 
a/nifi-framework-bundle/nifi-framework/nifi-framework-components/src/main/java/org/apache/nifi/controller/repository/StandardProcessSession.java
+++ 
b/nifi-framework-bundle/nifi-framework/nifi-framework-components/src/main/java/org/apache/nifi/controller/repository/StandardProcessSession.java
@@ -669,6 +669,9 @@ public class StandardProcessSession implements 
ProcessSession, ProvenanceEventEn
                 entry.getKey().putAll(entry.getValue());
             }
 
+            // Record ConnectionStatusEvents after FlowFiles are enqueued so 
that queue metadata is updated
+            recordConnectionStatusEvents(checkpoint);
+
             final long enqueueFlowFileFinishNanos = System.nanoTime();
             final long enqueueFlowFileNanos = enqueueFlowFileFinishNanos - 
updateEventRepositoryFinishNanos;
 
@@ -815,8 +818,6 @@ public class StandardProcessSession implements 
ProcessSession, ProvenanceEventEn
                 
context.getFlowFileEventRepository().updateRepository(connectionSessionEvent);
                 context.recordProcessSessionEvent(connectionSessionEvent);
             }
-
-            recordConnectionStatusEvents(checkpoint);
         } catch (final IOException ioe) {
             LOG.error("FlowFile Event Repository failed to update", ioe);
         }
diff --git 
a/nifi-framework-bundle/nifi-framework/nifi-framework-components/src/test/java/org/apache/nifi/controller/repository/StandardProcessSessionTest.java
 
b/nifi-framework-bundle/nifi-framework/nifi-framework-components/src/test/java/org/apache/nifi/controller/repository/StandardProcessSessionTest.java
index 1516507e894..1a57bf49b3f 100644
--- 
a/nifi-framework-bundle/nifi-framework/nifi-framework-components/src/test/java/org/apache/nifi/controller/repository/StandardProcessSessionTest.java
+++ 
b/nifi-framework-bundle/nifi-framework/nifi-framework-components/src/test/java/org/apache/nifi/controller/repository/StandardProcessSessionTest.java
@@ -64,6 +64,7 @@ import static org.mockito.ArgumentMatchers.anySet;
 import static org.mockito.ArgumentMatchers.eq;
 import static org.mockito.ArgumentMatchers.isA;
 import static org.mockito.ArgumentMatchers.isNull;
+import static org.mockito.Mockito.doAnswer;
 import static org.mockito.Mockito.mock;
 import static org.mockito.Mockito.never;
 import static org.mockito.Mockito.times;
@@ -79,6 +80,8 @@ class StandardProcessSessionTest {
 
     private static final long EXPECTED_BYTES = 32;
 
+    private static final int EXPECTED_FLOWFILES = 1;
+
     private static final byte[] CONTENT = new byte[]{2};
 
     private static final long BYTES_READ = CONTENT.length;
@@ -205,6 +208,7 @@ class StandardProcessSessionTest {
         final QueueSize outputQueueSize = mock(QueueSize.class);
         when(outputFlowFileQueue.size()).thenReturn(outputQueueSize);
         
when(outputFlowFileQueue.getFlowFileAvailability()).thenReturn(FlowFileAvailability.FLOWFILE_AVAILABLE);
+        doAnswer(inv -> when(outputFlowFileQueue.size()).thenReturn(new 
QueueSize(EXPECTED_FLOWFILES, 
EXPECTED_BYTES))).when(outputFlowFileQueue).putAll(any());
 
         final Relationship relationship = new 
Relationship.Builder().name(Relationship.class.getSimpleName()).build();
         
when(repositoryContext.getConnections(eq(relationship))).thenReturn(List.of(outputConnection));
@@ -238,6 +242,8 @@ class StandardProcessSessionTest {
         final ComponentMetricContext secondComponentMetricContext = 
secondConnectionStatusEvent.getComponentMetricContext();
         assertEquals(OUTPUT_CONNECTION_ID, secondComponentMetricContext.id());
         assertSourceDestinationFound(secondConnectionStatusEvent);
+        assertEquals(EXPECTED_FLOWFILES, 
secondConnectionStatusEvent.getQueuedCount());
+        assertEquals(EXPECTED_BYTES, 
secondConnectionStatusEvent.getQueuedBytes());
     }
 
     @Test

Reply via email to