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