This is an automated email from the ASF dual-hosted git repository.

davsclaus pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/camel.git


The following commit(s) were added to refs/heads/main by this push:
     new 6e562069d5af CAMEL-25280: camel-snmp - concurrent exchanges on one 
endpoint must not share the request PDU (#27309)
6e562069d5af is described below

commit 6e562069d5afdf5f345c8dbe310e796aa09f0e12
Author: allthingssecurity <[email protected]>
AuthorDate: Sat Oct 3 12:18:54 2026 +0530

    CAMEL-25280: camel-snmp - concurrent exchanges on one endpoint must not 
share the request PDU (#27309)
    
    Follow-up of #27242 (CAMEL-25245), as suggested in its review.
    `SnmpProducer` kept one request PDU in a field, created in `doStart`. The 
producer of an endpoint is shared by all its concurrent exchanges, and the 
`GET_NEXT` walk changes that PDU for every request (`clear()`, then the next 
OID). SNMP4J also encodes the same PDU object again when it retries after a 
timeout. So walks running at the same time on one endpoint changed each other's 
requests: a walk whose request was retried while another walk ran got the 
answer for the other walk's last  [...]
    
    Co-Authored-By: Claude Opus 5.5 <[email protected]>
---
 .../apache/camel/component/snmp/SnmpProducer.java  |  38 +++---
 .../camel/component/snmp/ConcurrentWalkTest.java   | 147 +++++++++++++++++++++
 2 files changed, 169 insertions(+), 16 deletions(-)

diff --git 
a/components/camel-snmp/src/main/java/org/apache/camel/component/snmp/SnmpProducer.java
 
b/components/camel-snmp/src/main/java/org/apache/camel/component/snmp/SnmpProducer.java
index 57cc2553dd18..89a0de8ab2f3 100644
--- 
a/components/camel-snmp/src/main/java/org/apache/camel/component/snmp/SnmpProducer.java
+++ 
b/components/camel-snmp/src/main/java/org/apache/camel/component/snmp/SnmpProducer.java
@@ -54,7 +54,6 @@ public class SnmpProducer extends DefaultProducer {
     private USM usm;
     private Target target;
     private SnmpActionType actionType;
-    private PDU pdu;
 
     public SnmpProducer(SnmpEndpoint endpoint, SnmpActionType actionType) {
         super(endpoint);
@@ -70,27 +69,34 @@ public class SnmpProducer extends DefaultProducer {
         LOG.debug("targetAddress: {}", targetAddress);
 
         this.usm = SnmpHelper.createAndSetUSM(endpoint);
-        this.pdu = SnmpHelper.createPDU(endpoint);
         this.target = SnmpHelper.createTarget(endpoint);
+    }
+
+    /**
+     * Creates the request PDU of one exchange: the producer is shared by 
concurrent exchanges, and a walk changes its
+     * PDU for every request (SNMP4J also sends the same PDU again on a retry).
+     */
+    private PDU createPdu() {
+        PDU pdu = SnmpHelper.createPDU(endpoint);
 
         // in here,only POLL do set the oids
         if (this.actionType == SnmpActionType.POLL) {
             for (OID oid : this.endpoint.getOids()) {
-                this.pdu.add(new VariableBinding(oid));
+                pdu.add(new VariableBinding(oid));
             }
         }
-        this.pdu.setErrorIndex(0);
-        this.pdu.setErrorStatus(0);
+        pdu.setErrorIndex(0);
+        pdu.setErrorStatus(0);
         if (endpoint.getSnmpVersion() > SnmpConstants.version1) {
-            this.pdu.setMaxRepetitions(0);
+            pdu.setMaxRepetitions(0);
         }
         // support POLL and GET_NEXT
         if (this.actionType == SnmpActionType.GET_NEXT) {
-            this.pdu.setType(PDU.GETNEXT);
+            pdu.setType(PDU.GETNEXT);
         } else {
-            this.pdu.setType(PDU.GET);
+            pdu.setType(PDU.GET);
         }
-
+        return pdu;
     }
 
     @Override
@@ -105,7 +111,6 @@ public class SnmpProducer extends DefaultProducer {
             this.targetAddress = null;
             this.usm = null;
             this.target = null;
-            this.pdu = null;
         }
     }
 
@@ -133,16 +138,17 @@ public class SnmpProducer extends DefaultProducer {
 
             snmp.listen();
 
+            PDU pdu = createPdu();
             if (this.actionType == SnmpActionType.GET_NEXT) {
                 // snmp walk
                 List<SnmpMessage> smLst = new ArrayList<>();
                 for (OID oid : this.endpoint.getOids()) {
-                    this.pdu.clear();
-                    this.pdu.add(new VariableBinding(oid));
+                    pdu.clear();
+                    pdu.add(new VariableBinding(oid));
 
                     boolean matched = true;
                     while (matched) {
-                        ResponseEvent responseEvent = snmp.send(this.pdu, 
this.target);
+                        ResponseEvent responseEvent = snmp.send(pdu, 
this.target);
                         if (responseEvent == null || 
responseEvent.getResponse() == null) {
                             throw new TimeoutException("SNMP Producer 
Timeout");
                         }
@@ -155,7 +161,7 @@ public class SnmpProducer extends DefaultProducer {
                             throw new CamelExchangeException(
                                     "SNMP walk of " + oid + " failed: " + 
response.getErrorStatusText(), exchange);
                         }
-                        OID requestedOid = this.pdu.get(0).getOid();
+                        OID requestedOid = pdu.get(0).getOid();
                         VariableBinding next = null;
                         for (VariableBinding variableBinding : 
response.getVariableBindings()) {
                             // compare the OIDs, not their strings: 
1.3.6.1.4.1.20 is not in the subtree of 1.3.6.1.4.1.2
@@ -177,7 +183,7 @@ public class SnmpProducer extends DefaultProducer {
                                                              + next.getOid(),
                                     exchange);
                         }
-                        this.pdu.clear();
+                        pdu.clear();
                         pdu.add(new VariableBinding(next.getOid()));
                         smLst.add(new 
SnmpMessage(getEndpoint().getCamelContext(), response));
                     }
@@ -185,7 +191,7 @@ public class SnmpProducer extends DefaultProducer {
                 exchange.getIn().setBody(smLst);
             } else {
                 // snmp get
-                ResponseEvent responseEvent = snmp.send(this.pdu, this.target);
+                ResponseEvent responseEvent = snmp.send(pdu, this.target);
 
                 LOG.debug("Snmp: sended");
 
diff --git 
a/components/camel-snmp/src/test/java/org/apache/camel/component/snmp/ConcurrentWalkTest.java
 
b/components/camel-snmp/src/test/java/org/apache/camel/component/snmp/ConcurrentWalkTest.java
new file mode 100644
index 000000000000..86842c33816b
--- /dev/null
+++ 
b/components/camel-snmp/src/test/java/org/apache/camel/component/snmp/ConcurrentWalkTest.java
@@ -0,0 +1,147 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements.  See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License.  You may obtain a copy of the License at
+ *
+ *      http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+package org.apache.camel.component.snmp;
+
+import java.io.IOException;
+import java.util.List;
+import java.util.Map;
+import java.util.concurrent.CompletableFuture;
+import java.util.concurrent.CountDownLatch;
+import java.util.concurrent.TimeUnit;
+import java.util.concurrent.atomic.AtomicBoolean;
+
+import org.apache.camel.builder.RouteBuilder;
+import org.junit.jupiter.api.AfterAll;
+import org.junit.jupiter.api.BeforeAll;
+import org.junit.jupiter.api.Test;
+import org.junit.jupiter.api.TestInstance;
+import org.snmp4j.CommandResponder;
+import org.snmp4j.CommandResponderEvent;
+import org.snmp4j.MessageException;
+import org.snmp4j.PDU;
+import org.snmp4j.Snmp;
+import org.snmp4j.mp.StatusInformation;
+import org.snmp4j.smi.OID;
+import org.snmp4j.smi.OctetString;
+import org.snmp4j.smi.UdpAddress;
+import org.snmp4j.smi.VariableBinding;
+import org.snmp4j.transport.DefaultUdpTransportMapping;
+
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertTrue;
+
+/**
+ * Two walks on the same endpoint at the same time: each must get the 
variables of the walked subtrees, also when one of
+ * its requests has to be sent again (SNMP4J sends the request PDU again after 
the timeout).
+ */
+@TestInstance(TestInstance.Lifecycle.PER_CLASS)
+public class ConcurrentWalkTest extends SnmpTestSupport {
+
+    private static final OID ALPHA = new OID("1.3.6.1.4.1.9999.1");
+    private static final OID BETA = new OID("1.3.6.1.4.1.9999.2");
+
+    // the agent's MIB: the answer to GETNEXT of each OID
+    private static final Map<OID, VariableBinding> NEXT = Map.of(
+            ALPHA, new VariableBinding(new OID("1.3.6.1.4.1.9999.1.1"), new 
OctetString("a1")),
+            new OID("1.3.6.1.4.1.9999.1.1"), new VariableBinding(new 
OID("1.3.6.1.4.1.9999.1.2"), new OctetString("a2")),
+            new OID("1.3.6.1.4.1.9999.1.2"), new VariableBinding(new 
OID("1.3.6.1.4.1.9999.2.1"), new OctetString("b1")),
+            BETA, new VariableBinding(new OID("1.3.6.1.4.1.9999.2.1"), new 
OctetString("b1")),
+            new OID("1.3.6.1.4.1.9999.2.1"), new VariableBinding(new 
OID("1.3.6.1.4.1.9999.2.2"), new OctetString("b2")),
+            new OID("1.3.6.1.4.1.9999.2.2"), new VariableBinding(new 
OID("1.3.6.1.4.1.9999.3.1"), new OctetString("c1")));
+
+    private final AtomicBoolean firstRequest = new AtomicBoolean(true);
+    private final CountDownLatch firstRequestDropped = new CountDownLatch(1);
+    private Snmp agent;
+    private String agentAddress;
+
+    @BeforeAll
+    public void startAgent() throws IOException {
+        DefaultUdpTransportMapping transport = new 
DefaultUdpTransportMapping(new UdpAddress("127.0.0.1/0"));
+        agent = new Snmp(transport);
+        agent.addCommandResponder(new Agent());
+        agent.listen();
+        agentAddress = 
transport.getListenAddress().toString().replaceFirst("/", ":");
+    }
+
+    @AfterAll
+    public void stopAgent() throws IOException {
+        if (agent != null) {
+            agent.close();
+        }
+    }
+
+    @Test
+    public void testConcurrentWalks() throws Exception {
+        // the agent does not answer the first request of the first walk, 
which SNMP4J sends again after the timeout;
+        // the second walk on the same endpoint runs completely in the meantime
+        CompletableFuture<List<?>> first
+                = CompletableFuture.supplyAsync(() -> 
template.requestBody("direct:walk", null, List.class));
+        assertTrue(firstRequestDropped.await(10, TimeUnit.SECONDS));
+
+        List<?> second = template.requestBody("direct:walk", null, List.class);
+
+        assertEquals(List.of("a1", "a2", "b1", "b2"), values(second));
+        assertEquals(List.of("a1", "a2", "b1", "b2"), values(first.get(20, 
TimeUnit.SECONDS)));
+    }
+
+    private static List<String> values(List<?> messages) {
+        return messages.stream()
+                .map(message -> ((SnmpMessage) 
message).getSnmpMessage().get(0).getVariable().toString())
+                .toList();
+    }
+
+    @Override
+    protected RouteBuilder createRouteBuilder() {
+        return new RouteBuilder() {
+            public void configure() {
+                from("direct:walk")
+                        .to("snmp:" + agentAddress + 
"?protocol=udp&type=GET_NEXT&timeout=2000&retries=1&oids="
+                            + ALPHA + "," + BETA);
+            }
+        };
+    }
+
+    private final class Agent implements CommandResponder {
+
+        @Override
+        public synchronized void processPdu(CommandResponderEvent event) {
+            PDU request = event.getPDU();
+            if (firstRequest.getAndSet(false)) {
+                firstRequestDropped.countDown();
+                return;
+            }
+            VariableBinding next = NEXT.get(request.get(0).getOid());
+            PDU response = (PDU) request.clone();
+            response.setType(PDU.RESPONSE);
+            if (next != null) {
+                response.set(0, next);
+            } else {
+                // end of the MIB view (SNMPv1)
+                response.setErrorStatus(PDU.noSuchName);
+                response.setErrorIndex(1);
+            }
+            try {
+                event.getMessageDispatcher().returnResponsePdu(
+                        event.getMessageProcessingModel(), 
event.getSecurityModel(), event.getSecurityName(),
+                        event.getSecurityLevel(), response, 
event.getMaxSizeResponsePDU(), event.getStateReference(),
+                        new StatusInformation());
+            } catch (MessageException e) {
+                throw new IllegalStateException(e);
+            }
+        }
+    }
+}

Reply via email to