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