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 69a5e032209 NIFI-16315 Reject oversized PutUDP FlowFiles before 
copying content (#11642)
69a5e032209 is described below

commit 69a5e032209616b537526db8ade40cb86267e242
Author: Joe Witt <[email protected]>
AuthorDate: Wed Sep 9 20:25:09 2026 -0700

    NIFI-16315 Reject oversized PutUDP FlowFiles before copying content (#11642)
    
    Signed-off-by: David Handermann <[email protected]>
---
 .../apache/nifi/processors/standard/PutUDP.java    | 75 +++++++++++++++++++++-
 .../nifi/processors/standard/TestPutUDP.java       | 15 +++--
 2 files changed, 82 insertions(+), 8 deletions(-)

diff --git 
a/nifi-extension-bundles/nifi-standard-bundle/nifi-standard-processors/src/main/java/org/apache/nifi/processors/standard/PutUDP.java
 
b/nifi-extension-bundles/nifi-standard-bundle/nifi-standard-processors/src/main/java/org/apache/nifi/processors/standard/PutUDP.java
index 55d224697b5..beb9d52c17a 100644
--- 
a/nifi-extension-bundles/nifi-standard-bundle/nifi-standard-processors/src/main/java/org/apache/nifi/processors/standard/PutUDP.java
+++ 
b/nifi-extension-bundles/nifi-standard-bundle/nifi-standard-processors/src/main/java/org/apache/nifi/processors/standard/PutUDP.java
@@ -24,6 +24,7 @@ import 
org.apache.nifi.annotation.documentation.CapabilityDescription;
 import org.apache.nifi.annotation.documentation.DeprecationNotice;
 import org.apache.nifi.annotation.documentation.SeeAlso;
 import org.apache.nifi.annotation.documentation.Tags;
+import org.apache.nifi.annotation.lifecycle.OnScheduled;
 import org.apache.nifi.event.transport.configuration.TransportProtocol;
 import org.apache.nifi.event.transport.netty.ByteArrayNettyEventSenderFactory;
 import org.apache.nifi.event.transport.netty.NettyEventSenderFactory;
@@ -37,11 +38,15 @@ import org.apache.nifi.util.StopWatch;
 
 import java.io.IOException;
 import java.io.InputStream;
+import java.net.Inet6Address;
+import java.net.InetAddress;
+import java.net.UnknownHostException;
 import java.util.concurrent.TimeUnit;
 
-@CapabilityDescription("The PutUDP processor receives a FlowFile and packages 
the FlowFile content into a single UDP datagram packet which is then 
transmitted to the configured UDP server."
-        + " The user must ensure that the FlowFile content being fed to this 
processor is not larger than the maximum size for the underlying UDP transport. 
The maximum transport size will "
-        + "vary based on the platform setup but is generally just under 64KB. 
FlowFiles will be marked as failed if their content is larger than the maximum 
transport size.")
+@CapabilityDescription("The PutUDP processor receives a FlowFile and packages 
the FlowFile content into a single UDP datagram packet which is then 
transmitted to the configured UDP server. "
+        + "A FlowFile larger than the destination address family maximum 
(65,507 bytes for IPv4, 65,527 bytes for IPv6) cannot be sent as one datagram "
+        + "and is routed to failure without reading content. "
+        + "The local UDP stack may still reject smaller datagrams; those 
FlowFiles are also marked as failed.")
 @InputRequirement(Requirement.INPUT_REQUIRED)
 @SeeAlso({ListenUDP.class, PutTCP.class})
 @Tags({ "remote", "egress", "put", "udp" })
@@ -49,6 +54,24 @@ import java.util.concurrent.TimeUnit;
 @DeprecationNotice(reason = "NIFI-16323: Limited transport size and lack of 
application protocol semantics")
 public class PutUDP extends AbstractPutEventProcessor<byte[]> {
 
+    /**
+     * Maximum UDP payload for IPv4: 65,535 byte IP packet minus 20 byte IPv4 
header minus 8 byte UDP header.
+     */
+    static final int MAX_IPV4_UDP_PAYLOAD_LENGTH = 65_535 - 20 - 8;
+
+    /**
+     * Maximum UDP payload for IPv6: 16-bit UDP Length covers the 8 byte UDP 
header and payload; the IPv6 header is not included.
+     */
+    static final int MAX_IPV6_UDP_PAYLOAD_LENGTH = 65_535 - 8;
+
+    private volatile int maxUdpPayloadLength = MAX_IPV4_UDP_PAYLOAD_LENGTH;
+
+    @OnScheduled
+    public void resolveMaxUdpPayloadLength(final ProcessContext context) {
+        final String hostname = 
context.getProperty(HOSTNAME).evaluateAttributeExpressions().getValue();
+        maxUdpPayloadLength = maxPayloadLength(hostname);
+    }
+
     @Override
     public void onTrigger(final ProcessContext context, final 
ProcessSessionFactory sessionFactory) throws ProcessException {
         final ProcessSession session = sessionFactory.createSession();
@@ -57,6 +80,14 @@ public class PutUDP extends 
AbstractPutEventProcessor<byte[]> {
             return;
         }
 
+        if (flowFile.getSize() > maxUdpPayloadLength) {
+            getLogger().error("Cannot send {} as a UDP datagram: size {} 
exceeds the {} maximum payload of {} bytes",
+                    flowFile, flowFile.getSize(), 
addressFamilyLabel(maxUdpPayloadLength), maxUdpPayloadLength);
+            session.transfer(session.penalize(flowFile), REL_FAILURE);
+            session.commitAsync();
+            return;
+        }
+
         final StopWatch stopWatch = new StopWatch(true);
         try {
             final byte[] content = readContent(session, flowFile);
@@ -83,6 +114,44 @@ public class PutUDP extends 
AbstractPutEventProcessor<byte[]> {
         return new ByteArrayNettyEventSenderFactory(getLogger(), hostname, 
port, TransportProtocol.UDP);
     }
 
+    /**
+     * Returns the UDP payload ceiling for the destination hostname. When any 
resolved address is a non-mapped IPv6
+     * address the IPv6 maximum is used so dual-stack names are not rejected 
at the IPv4 ceiling. Unknown hosts use
+     * the IPv6 maximum as a copy cap so IPv6-legal sizes are not rejected 
before send.
+     */
+    static int maxPayloadLength(final String hostname) {
+        try {
+            for (final InetAddress address : 
InetAddress.getAllByName(hostname)) {
+                if (isUnmappedIPv6Address(address)) {
+                    return MAX_IPV6_UDP_PAYLOAD_LENGTH;
+                }
+            }
+            return MAX_IPV4_UDP_PAYLOAD_LENGTH;
+        } catch (final UnknownHostException e) {
+            return MAX_IPV6_UDP_PAYLOAD_LENGTH;
+        }
+    }
+
+    static boolean isUnmappedIPv6Address(final InetAddress address) {
+        return address instanceof Inet6Address && 
!isIPv4MappedAddress(address.getAddress());
+    }
+
+    private static boolean isIPv4MappedAddress(final byte[] address) {
+        if (address.length != 16) {
+            return false;
+        }
+        for (int i = 0; i < 10; i++) {
+            if (address[i] != 0) {
+                return false;
+            }
+        }
+        return address[10] == (byte) 0xff && address[11] == (byte) 0xff;
+    }
+
+    private static String addressFamilyLabel(final int maxPayloadLength) {
+        return maxPayloadLength == MAX_IPV6_UDP_PAYLOAD_LENGTH ? "IPv6" : 
"IPv4";
+    }
+
     private byte[] readContent(final ProcessSession session, final FlowFile 
flowFile) throws IOException {
         try (final InputStream inputStream = session.read(flowFile)) {
             return IOUtils.toByteArray(inputStream);
diff --git 
a/nifi-extension-bundles/nifi-standard-bundle/nifi-standard-processors/src/test/java/org/apache/nifi/processors/standard/TestPutUDP.java
 
b/nifi-extension-bundles/nifi-standard-bundle/nifi-standard-processors/src/test/java/org/apache/nifi/processors/standard/TestPutUDP.java
index aee554d3498..6efe6e8906a 100644
--- 
a/nifi-extension-bundles/nifi-standard-bundle/nifi-standard-processors/src/test/java/org/apache/nifi/processors/standard/TestPutUDP.java
+++ 
b/nifi-extension-bundles/nifi-standard-bundle/nifi-standard-processors/src/test/java/org/apache/nifi/processors/standard/TestPutUDP.java
@@ -39,6 +39,7 @@ import java.util.concurrent.LinkedBlockingQueue;
 import static org.junit.jupiter.api.Assertions.assertArrayEquals;
 import static org.junit.jupiter.api.Assertions.assertNotNull;
 import static org.junit.jupiter.api.Assertions.assertNull;
+import static org.junit.jupiter.api.Assertions.assertTrue;
 
 @Timeout(10)
 public class TestPutUDP {
@@ -49,7 +50,7 @@ public class TestPutUDP {
     private static final Charset CHARSET = StandardCharsets.UTF_8;
     private static final int MAX_FRAME_LENGTH = 32800;
     private static final int VALID_LARGE_FILE_SIZE = 32768;
-    private static final int INVALID_LARGE_FILE_SIZE = 1_000_000;
+    private static final int OVERSIZE_UDP_PAYLOAD = 
PutUDP.MAX_IPV4_UDP_PAYLOAD_LENGTH + 1;
     private static final char CONTENT_CHAR = 'x';
     private static final int DATA_WAIT_PERIOD = 50;
     private static final String[] EMPTY_FILE = {""};
@@ -100,13 +101,17 @@ public class TestPutUDP {
     }
 
     @Test
-    public void testSendLargeFileInvalid() throws Exception {
+    public void testSendLargerThanUdpPayloadLimit() throws Exception {
         configureProperties();
-        String[] testData = createContent(INVALID_LARGE_FILE_SIZE);
-        sendMessages(testData);
-        checkRelationships(0, testData.length);
+        runner.enqueue(new byte[OVERSIZE_UDP_PAYLOAD]);
+        runner.run();
+
+        checkRelationships(0, 1);
         checkNoDataReceived();
         runner.assertQueueEmpty();
+        assertTrue(runner.getLogger().getErrorMessages().stream()
+                .anyMatch(message -> message.getMsg().contains("exceeds the 
IPv4 maximum payload of "
+                        + PutUDP.MAX_IPV4_UDP_PAYLOAD_LENGTH)));
     }
 
     @Test

Reply via email to