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

davsclaus pushed a commit to branch fix/CAMEL-24160
in repository https://gitbox.apache.org/repos/asf/camel.git

commit 4d3f6d8f78c03f7fff627c2a58cbd617e1a4249a
Author: Claus Ibsen <[email protected]>
AuthorDate: Fri Jul 17 15:49:19 2026 +0200

    CAMEL-24160: camel-salesforce - fix high-severity bugs from code review
    
    Fixes 9 bugs found during a deep code review of camel-salesforce:
    
    - JsonRestProcessor: fix infinite loop on non-object JSON in deferred
      SObject type detection, and fix RuntimeException escaping IOException
      catch block causing silent data corruption
    - DefaultBulkApiV2Client: add missing return after error callback in
      changeJobState/changeQueryJobState to prevent double callback
    - SalesforceEndpointConfig: add allOrNone to toValueMap() so the URI
      option takes effect, and fix REPLAY_PRESET mapped to wrong field
    - PubSubApiClient: skip binary trailer keys in onError to prevent
      IllegalArgumentException killing reconnection, add stop/channel
      checks in resubscribeOnError to prevent infinite loop on shutdown,
      cap reconnect delay at maxBackoff, and restore thread interrupt
    - SubscriptionHelper: add backoff delay to handshake retry to prevent
      tight retry loop hammering the Salesforce API
    - SalesforceSecurityHandler: guard against null client in 401 handling
      to prevent NPE for CometD streaming requests
    - SalesforceSession: make latch volatile and remove dead spin loop to
      fix race condition in attemptLoginUntilSuccessful
    
    Co-Authored-By: Claude Opus 4.6 <[email protected]>
    Signed-off-by: Claus Ibsen <[email protected]>
---
 .../salesforce/SalesforceEndpointConfig.java       |  3 +-
 .../salesforce/internal/SalesforceSession.java     |  6 +---
 .../internal/client/DefaultBulkApiV2Client.java    |  2 ++
 .../internal/client/PubSubApiClient.java           | 21 +++++++++++---
 .../internal/client/SalesforceSecurityHandler.java | 33 +++++++++++++---------
 .../internal/processor/JsonRestProcessor.java      | 10 +++++--
 .../internal/streaming/SubscriptionHelper.java     | 22 +++++++++++++--
 7 files changed, 67 insertions(+), 30 deletions(-)

diff --git 
a/components/camel-salesforce/camel-salesforce-component/src/main/java/org/apache/camel/component/salesforce/SalesforceEndpointConfig.java
 
b/components/camel-salesforce/camel-salesforce-component/src/main/java/org/apache/camel/component/salesforce/SalesforceEndpointConfig.java
index 6b8f3f4d0778..843d7e42a04d 100644
--- 
a/components/camel-salesforce/camel-salesforce-component/src/main/java/org/apache/camel/component/salesforce/SalesforceEndpointConfig.java
+++ 
b/components/camel-salesforce/camel-salesforce-component/src/main/java/org/apache/camel/component/salesforce/SalesforceEndpointConfig.java
@@ -834,6 +834,7 @@ public class SalesforceEndpointConfig implements Cloneable {
         valueMap.put(APEX_URL, apexUrl);
         // apexQueryParams are handled explicitly in AbstractRestProcessor
         valueMap.put(COMPOSITE_METHOD, compositeMethod);
+        valueMap.put(ALL_OR_NONE, allOrNone);
         valueMap.put(LIMIT, limit);
         valueMap.put(APPROVAL, approval);
         valueMap.put(EVENT_NAME, eventName);
@@ -869,7 +870,7 @@ public class SalesforceEndpointConfig implements Cloneable {
         valueMap.put(INITIAL_REPLAY_ID_MAP, initialReplayIdMap);
 
         // add Pub/Sub API properties
-        valueMap.put(REPLAY_PRESET, initialReplayIdMap);
+        valueMap.put(REPLAY_PRESET, replayPreset);
         valueMap.put(PUB_SUB_DESERIALIZE_TYPE, pubSubDeserializeType);
         valueMap.put(PUB_SUB_POJO_CLASS, pubSubPojoClass);
 
diff --git 
a/components/camel-salesforce/camel-salesforce-component/src/main/java/org/apache/camel/component/salesforce/internal/SalesforceSession.java
 
b/components/camel-salesforce/camel-salesforce-component/src/main/java/org/apache/camel/component/salesforce/internal/SalesforceSession.java
index e810917ce066..fbbb9ca40649 100644
--- 
a/components/camel-salesforce/camel-salesforce-component/src/main/java/org/apache/camel/component/salesforce/internal/SalesforceSession.java
+++ 
b/components/camel-salesforce/camel-salesforce-component/src/main/java/org/apache/camel/component/salesforce/internal/SalesforceSession.java
@@ -91,7 +91,7 @@ public class SalesforceSession extends ServiceSupport {
 
     private final CamelContext camelContext;
     private final AtomicBoolean loggingIn = new AtomicBoolean();
-    private CountDownLatch latch = new CountDownLatch(1);
+    private volatile CountDownLatch latch = new CountDownLatch(1);
 
     public SalesforceSession(CamelContext camelContext, SalesforceHttpClient 
httpClient, long timeout,
                              SalesforceLoginConfig config) {
@@ -113,11 +113,7 @@ public class SalesforceSession extends ServiceSupport {
         // if another thread is logging in, we will just wait until it's 
successful
         if (!loggingIn.compareAndSet(false, true)) {
             LOG.debug("waiting on login from another thread");
-            // TODO: This is janky
             try {
-                while (latch == null) {
-                    Thread.sleep(100);
-                }
                 latch.await();
             } catch (InterruptedException ex) {
                 Thread.currentThread().interrupt();
diff --git 
a/components/camel-salesforce/camel-salesforce-component/src/main/java/org/apache/camel/component/salesforce/internal/client/DefaultBulkApiV2Client.java
 
b/components/camel-salesforce/camel-salesforce-component/src/main/java/org/apache/camel/component/salesforce/internal/client/DefaultBulkApiV2Client.java
index 3fcd74054bd7..7bbbe07d169f 100644
--- 
a/components/camel-salesforce/camel-salesforce-component/src/main/java/org/apache/camel/component/salesforce/internal/client/DefaultBulkApiV2Client.java
+++ 
b/components/camel-salesforce/camel-salesforce-component/src/main/java/org/apache/camel/component/salesforce/internal/client/DefaultBulkApiV2Client.java
@@ -113,6 +113,7 @@ public class DefaultBulkApiV2Client extends 
AbstractClientBase implements BulkAp
             public void onResponse(InputStream response, Map<String, String> 
headers, SalesforceException ex) {
                 if (ex != null) {
                     callback.onResponse(null, headers, ex);
+                    return;
                 }
                 Job responseJob = null;
                 try {
@@ -237,6 +238,7 @@ public class DefaultBulkApiV2Client extends 
AbstractClientBase implements BulkAp
             public void onResponse(InputStream response, Map<String, String> 
headers, SalesforceException ex) {
                 if (ex != null) {
                     callback.onResponse(null, headers, ex);
+                    return;
                 }
                 QueryJob responseJob = null;
                 try {
diff --git 
a/components/camel-salesforce/camel-salesforce-component/src/main/java/org/apache/camel/component/salesforce/internal/client/PubSubApiClient.java
 
b/components/camel-salesforce/camel-salesforce-component/src/main/java/org/apache/camel/component/salesforce/internal/client/PubSubApiClient.java
index 461e24e3bb9b..08b23b69b7de 100644
--- 
a/components/camel-salesforce/camel-salesforce-component/src/main/java/org/apache/camel/component/salesforce/internal/client/PubSubApiClient.java
+++ 
b/components/camel-salesforce/camel-salesforce-component/src/main/java/org/apache/camel/component/salesforce/internal/client/PubSubApiClient.java
@@ -339,8 +339,12 @@ public class PubSubApiClient extends ServiceSupport {
                 String errorCode = "";
                 LOG.error("Trailers:");
                 if (trailers != null) {
-                    trailers.keys().forEach(trailer -> LOG.error("Trailer: {}, 
Value: {}", trailer,
-                            trailers.get(Metadata.Key.of(trailer, 
Metadata.ASCII_STRING_MARSHALLER))));
+                    trailers.keys().forEach(trailer -> {
+                        if (!trailer.endsWith(Metadata.BINARY_HEADER_SUFFIX)) {
+                            LOG.error("Trailer: {}, Value: {}", trailer,
+                                    trailers.get(Metadata.Key.of(trailer, 
Metadata.ASCII_STRING_MARSHALLER)));
+                        }
+                    });
                     errorCode = trailers.get(Metadata.Key.of("error-code", 
Metadata.ASCII_STRING_MARSHALLER));
                 }
                 if (errorCode != null) {
@@ -378,12 +382,21 @@ public class PubSubApiClient extends ServiceSupport {
         }
 
         private void resubscribeOnError() {
+            if (isStoppingOrStopped() || channel == null || 
channel.isShutdown()) {
+                LOG.debug("Client is stopping or channel is shut down, 
skipping resubscribe");
+                return;
+            }
             try {
                 LOG.debug("Will attempt resubscribe in {} ms", reconnectDelay);
                 Thread.sleep(reconnectDelay);
-                reconnectDelay = reconnectDelay + backoffIncrement;
+                reconnectDelay = Math.min(reconnectDelay + backoffIncrement, 
maxBackoff);
             } catch (InterruptedException e) {
-                throw new RuntimeException(e);
+                Thread.currentThread().interrupt();
+                return;
+            }
+            if (isStoppingOrStopped() || channel == null || 
channel.isShutdown()) {
+                LOG.debug("Client is stopping or channel is shut down, 
skipping resubscribe");
+                return;
             }
             if (replayId != null) {
                 subscribe(consumer, ReplayPreset.CUSTOM, replayId, 
fallbackToLatestReplayId);
diff --git 
a/components/camel-salesforce/camel-salesforce-component/src/main/java/org/apache/camel/component/salesforce/internal/client/SalesforceSecurityHandler.java
 
b/components/camel-salesforce/camel-salesforce-component/src/main/java/org/apache/camel/component/salesforce/internal/client/SalesforceSecurityHandler.java
index e2d2d533633c..e1ac6f29ff80 100644
--- 
a/components/camel-salesforce/camel-salesforce-component/src/main/java/org/apache/camel/component/salesforce/internal/client/SalesforceSecurityHandler.java
+++ 
b/components/camel-salesforce/camel-salesforce-component/src/main/java/org/apache/camel/component/salesforce/internal/client/SalesforceSecurityHandler.java
@@ -171,16 +171,19 @@ public class SalesforceSecurityHandler implements 
ProtocolHandler {
                 // Salesforce will allow successful login with an expired 
password, but any subsequent
                 // API calls will fail with a 401 and message about expired 
password.
                 // It's fatal. User must reset password.
-                List<RestError> errors = Collections.emptyList();
-                try {
-                    errors = client.readErrorsFrom(getContentAsInputStream(), 
objectMapper);
-                } catch (IOException e) {
-                    LOG.warn("Unable to deserialize errors from response 
body.");
-                }
-                if (errors.stream().anyMatch(error -> 
EXPIRED_PASSWORD_CODE.equals(error.getErrorCode()))) {
-                    SalesforceException salesforceException = 
createSalesforceException(client, status);
-                    forwardFailureComplete(request, null, response, 
salesforceException);
-                    return;
+                // Note: client may be null for CometD streaming requests
+                if (client != null) {
+                    List<RestError> errors = Collections.emptyList();
+                    try {
+                        errors = 
client.readErrorsFrom(getContentAsInputStream(), objectMapper);
+                    } catch (IOException e) {
+                        LOG.warn("Unable to deserialize errors from response 
body.");
+                    }
+                    if (errors.stream().anyMatch(error -> 
EXPIRED_PASSWORD_CODE.equals(error.getErrorCode()))) {
+                        SalesforceException salesforceException = 
createSalesforceException(client, status);
+                        forwardFailureComplete(request, null, response, 
salesforceException);
+                        return;
+                    }
                 }
 
                 // REST token expiry
@@ -214,10 +217,12 @@ public class SalesforceSecurityHandler implements 
ProtocolHandler {
 
         private SalesforceException 
createSalesforceException(AbstractClientBase client, int statusCode) {
             List<RestError> restErrors = Collections.emptyList();
-            try {
-                restErrors = client.readErrorsFrom(getContentAsInputStream(), 
new ObjectMapper());
-            } catch (IOException e) {
-                LOG.warn("Unable to deserialize errors from response body.");
+            if (client != null) {
+                try {
+                    restErrors = 
client.readErrorsFrom(getContentAsInputStream(), new ObjectMapper());
+                } catch (IOException e) {
+                    LOG.warn("Unable to deserialize errors from response 
body.");
+                }
             }
             return new SalesforceException(restErrors, statusCode);
         }
diff --git 
a/components/camel-salesforce/camel-salesforce-component/src/main/java/org/apache/camel/component/salesforce/internal/processor/JsonRestProcessor.java
 
b/components/camel-salesforce/camel-salesforce-component/src/main/java/org/apache/camel/component/salesforce/internal/processor/JsonRestProcessor.java
index a06da2bfdcf9..7dc586be91e2 100644
--- 
a/components/camel-salesforce/camel-salesforce-component/src/main/java/org/apache/camel/component/salesforce/internal/processor/JsonRestProcessor.java
+++ 
b/components/camel-salesforce/camel-salesforce-component/src/main/java/org/apache/camel/component/salesforce/internal/processor/JsonRestProcessor.java
@@ -241,7 +241,8 @@ public class JsonRestProcessor extends 
AbstractRestProcessor {
         Class<?> responseClass;
         try (final JsonParser parser = new 
JsonFactory().createParser(responseEntity)) {
             String type = null;
-            while (parser.nextToken() != JsonToken.END_OBJECT) {
+            JsonToken token;
+            while ((token = parser.nextToken()) != null && token != 
JsonToken.END_OBJECT) {
                 String propName = parser.currentName();
                 if ("type".equals(propName)) {
                     parser.nextToken();
@@ -249,10 +250,13 @@ public class JsonRestProcessor extends 
AbstractRestProcessor {
                     break;
                 }
             }
+            if (type == null) {
+                return null;
+            }
             String prefix = exchange.getProperty(RESPONSE_CLASS_PREFIX, "", 
String.class);
             responseClass = getSObjectClass(prefix + type, null);
-        } catch (IOException | SalesforceException exc) {
-            throw new RuntimeException(exc);
+        } catch (SalesforceException exc) {
+            throw new IOException("Failed to detect SObject response class", 
exc);
         } finally {
             responseEntity.reset();
         }
diff --git 
a/components/camel-salesforce/camel-salesforce-component/src/main/java/org/apache/camel/component/salesforce/internal/streaming/SubscriptionHelper.java
 
b/components/camel-salesforce/camel-salesforce-component/src/main/java/org/apache/camel/component/salesforce/internal/streaming/SubscriptionHelper.java
index ab77195d963b..ce17a1d165ea 100644
--- 
a/components/camel-salesforce/camel-salesforce-component/src/main/java/org/apache/camel/component/salesforce/internal/streaming/SubscriptionHelper.java
+++ 
b/components/camel-salesforce/camel-salesforce-component/src/main/java/org/apache/camel/component/salesforce/internal/streaming/SubscriptionHelper.java
@@ -144,9 +144,25 @@ public class SubscriptionHelper extends ServiceSupport {
                         }
                     }
                 }
-                // failed, so keep trying
-                LOG.debug("Handshake failed, so try again.");
-                client.handshake();
+                // failed, so keep trying with backoff
+                final long backoff = 
handshakeBackoff.getAndAdd(backoffIncrement);
+                if (backoff > maxBackoff) {
+                    LOG.error("Handshake retry aborted after exceeding {} 
msecs backoff", maxBackoff);
+                } else {
+                    LOG.debug("Pausing for {} msecs before handshake retry", 
backoff);
+                    if (backoff > 0) {
+                        Tasks.foregroundTask()
+                                .withBudget(Budgets.iterationBudget()
+                                        .withMaxIterations(1)
+                                        
.withInitialDelay(Duration.ofMillis(backoff))
+                                        .withInterval(Duration.ZERO)
+                                        .build())
+                                .withName("SalesforceHandshakeRetryDelay")
+                                .build()
+                                .run(component.getCamelContext(), () -> true);
+                    }
+                    client.handshake();
+                }
             } else if (!channelToConsumers.isEmpty()) {
                 channelsLock.lock();
                 try {

Reply via email to