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 701eb46e477e CAMEL-24161: camel-salesforce - fix medium-severity bugs 
from code review
701eb46e477e is described below

commit 701eb46e477efd0e983c28c7d6832db6e52df172
Author: Claus Ibsen <[email protected]>
AuthorDate: Fri Jul 17 17:28:27 2026 +0200

    CAMEL-24161: camel-salesforce - fix medium-severity bugs from code review
    
    Co-Authored-By: Claude Opus 4.6 <[email protected]>
    Signed-off-by: Claus Ibsen <[email protected]>
---
 .../salesforce/SalesforceLoginConfig.java          |  2 +-
 .../internal/client/AbstractClientBase.java        |  5 ++++
 .../internal/client/DefaultRawClient.java          |  4 +--
 .../internal/client/PubSubApiClient.java           | 25 ++++++++++--------
 .../internal/client/SalesforceSecurityHandler.java | 30 ++++++++++++++--------
 .../internal/processor/AbstractRestProcessor.java  | 20 ++++++++++++---
 .../internal/processor/CompositeApiProcessor.java  |  9 +++++++
 .../internal/processor/JsonRestProcessor.java      |  4 ++-
 .../internal/streaming/SubscriptionHelper.java     |  4 ++-
 .../camel/component/salesforce/PubSubApiTest.java  | 15 +++++------
 10 files changed, 79 insertions(+), 39 deletions(-)

diff --git 
a/components/camel-salesforce/camel-salesforce-component/src/main/java/org/apache/camel/component/salesforce/SalesforceLoginConfig.java
 
b/components/camel-salesforce/camel-salesforce-component/src/main/java/org/apache/camel/component/salesforce/SalesforceLoginConfig.java
index e1e68ab2cfc5..0f86c911dfbd 100644
--- 
a/components/camel-salesforce/camel-salesforce-component/src/main/java/org/apache/camel/component/salesforce/SalesforceLoginConfig.java
+++ 
b/components/camel-salesforce/camel-salesforce-component/src/main/java/org/apache/camel/component/salesforce/SalesforceLoginConfig.java
@@ -270,7 +270,7 @@ public class SalesforceLoginConfig {
     public String toString() {
         return "SalesforceLoginConfig[" + "instanceUrl= '" + instanceUrl + "', 
loginUrl='" + loginUrl + '\'' + ","
                + "clientId='" + clientId + '\'' + ", clientSecret='********'"
-               + ", refreshToken='" + refreshToken + '\'' + ", userName='" + 
userName + '\'' + ", password=********'"
+               + ", refreshToken='********'" + ", userName='" + userName + 
'\'' + ", password=********'"
                + ", keystore=********', audience='" + jwtAudience + '\'' + ","
                + ", lazyLogin=" + lazyLogin + ']';
     }
diff --git 
a/components/camel-salesforce/camel-salesforce-component/src/main/java/org/apache/camel/component/salesforce/internal/client/AbstractClientBase.java
 
b/components/camel-salesforce/camel-salesforce-component/src/main/java/org/apache/camel/component/salesforce/internal/client/AbstractClientBase.java
index 546b6977385e..af142cf6773f 100644
--- 
a/components/camel-salesforce/camel-salesforce-component/src/main/java/org/apache/camel/component/salesforce/internal/client/AbstractClientBase.java
+++ 
b/components/camel-salesforce/camel-salesforce-component/src/main/java/org/apache/camel/component/salesforce/internal/client/AbstractClientBase.java
@@ -190,6 +190,11 @@ public abstract class AbstractClientBase extends 
ServiceSupport
             final List<ByteBuffer> buffers = new ArrayList<>();
             while (true) {
                 Content.Chunk chunk = inputStreamRequestContent.read();
+                if (Content.Chunk.isFailure(chunk)) {
+                    chunk.release();
+                    throw new RuntimeException(
+                            new IOException("Failed to buffer request 
content", chunk.getFailure()));
+                }
                 if (chunk.isLast()) {
                     break;
                 } else {
diff --git 
a/components/camel-salesforce/camel-salesforce-component/src/main/java/org/apache/camel/component/salesforce/internal/client/DefaultRawClient.java
 
b/components/camel-salesforce/camel-salesforce-component/src/main/java/org/apache/camel/component/salesforce/internal/client/DefaultRawClient.java
index 610b8c320e63..27f59e4a9cff 100644
--- 
a/components/camel-salesforce/camel-salesforce-component/src/main/java/org/apache/camel/component/salesforce/internal/client/DefaultRawClient.java
+++ 
b/components/camel-salesforce/camel-salesforce-component/src/main/java/org/apache/camel/component/salesforce/internal/client/DefaultRawClient.java
@@ -48,8 +48,8 @@ public class DefaultRawClient extends AbstractClientBase 
implements RawClient {
     @Override
     protected void setAccessToken(Request request) {
         // replace old token
-        request.headers(h -> h.add(BULK_TOKEN_HEADER, accessToken));
-        request.headers(h -> h.add(REST_TOKEN_HEADER, TOKEN_PREFIX + 
accessToken));
+        request.headers(h -> h.put(BULK_TOKEN_HEADER, accessToken));
+        request.headers(h -> h.put(REST_TOKEN_HEADER, TOKEN_PREFIX + 
accessToken));
     }
 
     @Override
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..edff9c8c854a 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
@@ -101,10 +101,6 @@ public class PubSubApiClient extends ServiceSupport {
     private ManagedChannel channel;
     private boolean usePlainTextConnection = false;
 
-    private ReplayPreset initialReplayPreset;
-    private String initialReplayId;
-    private boolean fallbackToLatestReplayId;
-
     public PubSubApiClient(SalesforceSession session, SalesforceLoginConfig 
loginConfig, String pubSubHost,
                            int pubSubPort, long backoffIncrement, long 
maxBackoff, boolean allowUseProxyServer) {
         this.session = session;
@@ -151,9 +147,6 @@ public class PubSubApiClient extends ServiceSupport {
     public void subscribe(
             PubSubApiConsumer consumer, ReplayPreset replayPreset, String 
initialReplayId, boolean fallbackToLatestReplayId) {
         LOG.debug("Starting subscribe {}", consumer.getTopic());
-        this.initialReplayPreset = replayPreset;
-        this.initialReplayId = initialReplayId;
-        this.fallbackToLatestReplayId = fallbackToLatestReplayId;
         if (replayPreset == ReplayPreset.CUSTOM && initialReplayId == null) {
             throw new RuntimeException("initialReplayId is required for 
ReplayPreset.CUSTOM");
         }
@@ -164,7 +157,8 @@ public class PubSubApiClient extends ServiceSupport {
             replayId = base64DecodeToByteString(initialReplayId);
         }
         LOG.info("Subscribing to topic: {}.", topic);
-        final FetchResponseObserver responseObserver = new 
FetchResponseObserver(consumer);
+        final FetchResponseObserver responseObserver
+                = new FetchResponseObserver(consumer, replayPreset, 
initialReplayId, fallbackToLatestReplayId);
         StreamObserver<FetchRequest> serverStream = 
asyncStub.subscribe(responseObserver);
         LOG.info("Subscribe successful.");
         responseObserver.setServerStream(serverStream);
@@ -292,13 +286,20 @@ public class PubSubApiClient extends ServiceSupport {
         private final PubSubApiConsumer consumer;
         private final Map<String, Class<?>> eventClassMap;
         private final Class<?> pojoClass;
+        private final String initialReplayId;
+        private final boolean fallbackToLatestReplayId;
+        private ReplayPreset initialReplayPreset;
         private String replayId;
         private StreamObserver<FetchRequest> serverStream;
 
-        public FetchResponseObserver(PubSubApiConsumer consumer) {
+        public FetchResponseObserver(PubSubApiConsumer consumer, ReplayPreset 
initialReplayPreset,
+                                     String initialReplayId, boolean 
fallbackToLatestReplayId) {
             this.consumer = consumer;
             this.eventClassMap = consumer.getEventClassMap();
             this.pojoClass = consumer.getPojoClass();
+            this.initialReplayPreset = initialReplayPreset;
+            this.initialReplayId = initialReplayId;
+            this.fallbackToLatestReplayId = fallbackToLatestReplayId;
         }
 
         @Override
@@ -315,7 +316,9 @@ public class PubSubApiClient extends ServiceSupport {
                 try {
                     processEvent(fetchResponse.getRpcId(), ce);
                 } catch (Exception e) {
-                    LOG.error(e.toString(), e);
+                    LOG.error("Failed to process event on topic {}: {}", 
topic, e.getMessage(), e);
+                    consumer.getExceptionHandler().handleException(
+                            "Failed to process Pub/Sub event on topic " + 
topic, e);
                 }
             }
             replayId = 
base64EncodeByteString(fetchResponse.getLatestReplayId());
@@ -366,6 +369,8 @@ public class PubSubApiClient extends ServiceSupport {
                                 
consumer.getExceptionHandler().handleException(new InvalidReplayIdException(
                                         "Corrupt replay id: " + currReplayId,
                                         currReplayId));
+                                replayId = null;
+                                return;
                             }
                         }
                         default -> LOG.error("unexpected errorCode: {}", 
errorCode);
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..054168e21e0a 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
@@ -186,9 +186,8 @@ public class SalesforceSecurityHandler implements 
ProtocolHandler {
                 // REST token expiry
                 LOG.warn("Retrying on Salesforce authentication error [{}]: 
[{}]", status, reason);
 
-                // remember original request and send a relogin request in
-                // current conversation
-                retryLogin(request, retries);
+                // coordinated re-login and retry via session.login() on 
worker pool
+                retryLogin(request, retries, client);
 
             } else if (status < HttpStatus.OK_200 || status >= 
HttpStatus.MULTIPLE_CHOICES_300) {
 
@@ -202,7 +201,7 @@ public class SalesforceSecurityHandler implements 
ProtocolHandler {
 
                     // retry Bulk API call
                     LOG.warn("Retrying on Bulk API Salesforce authentication 
error [{}]: [{}]", status, reason);
-                    retryLogin(request, retries);
+                    retryLogin(request, retries, client);
 
                 } else {
 
@@ -236,13 +235,21 @@ public class SalesforceSecurityHandler implements 
ProtocolHandler {
                     && 
"InvalidSessionId".equals(e.getErrors().get(0).getErrorCode());
         }
 
-        private void retryLogin(HttpRequest request, Integer retries) {
-
+        private void retryLogin(HttpRequest request, Integer retries, 
AbstractClientBase client) {
             final HttpConversation conversation = request.getConversation();
-            // remember the original request to resend
-            conversation.setAttribute(AUTHENTICATION_REQUEST_ATTRIBUTE, 
request);
 
-            retryRequest((HttpRequest) session.getLoginRequest(conversation), 
null, retries, conversation, false);
+            // use session.login() on worker pool for single-flight 
coordination:
+            // concurrent 401s share one login instead of each firing its own
+            httpClient.getWorkerPool().execute(() -> {
+                try {
+                    session.login(session.getAccessToken());
+                } catch (SalesforceException e) {
+                    LOG.error("Login failed during authentication retry", e);
+                    forwardFailureComplete(request, null, null, e);
+                    return;
+                }
+                retryRequest(request, client, retries, conversation, true);
+            });
         }
 
         private void retryRequest(
@@ -255,8 +262,9 @@ public class SalesforceSecurityHandler implements 
ProtocolHandler {
             if (copy) {
                 newRequest = httpClient.copyRequest(request, request.getURI());
                 final Request.Content body = newRequest.getBody();
-                if (body != null) {
-                    body.rewind();
+                if (body != null && !body.rewind()) {
+                    LOG.warn("Request body cannot be replayed for 
authentication retry (content type: {})",
+                            body.getContentType());
                 }
                 newRequest.method(request.getMethod());
                 newRequest.headers(headers -> {
diff --git 
a/components/camel-salesforce/camel-salesforce-component/src/main/java/org/apache/camel/component/salesforce/internal/processor/AbstractRestProcessor.java
 
b/components/camel-salesforce/camel-salesforce-component/src/main/java/org/apache/camel/component/salesforce/internal/processor/AbstractRestProcessor.java
index 3158c2772ef5..aff931270648 100644
--- 
a/components/camel-salesforce/camel-salesforce-component/src/main/java/org/apache/camel/component/salesforce/internal/processor/AbstractRestProcessor.java
+++ 
b/components/camel-salesforce/camel-salesforce-component/src/main/java/org/apache/camel/component/salesforce/internal/processor/AbstractRestProcessor.java
@@ -479,7 +479,11 @@ public abstract class AbstractRestProcessor extends 
AbstractSalesforceProcessor
         final AbstractSObjectBase sObjectBase = 
exchange.getIn().getBody(AbstractSObjectBase.class);
         if (sObjectBase != null) {
             oldValue = getAndClearPropertyValue(sObjectBase, sObjectExtIdName);
-            sObjectExtIdValue = oldValue.toString();
+            if (oldValue != null) {
+                sObjectExtIdValue = oldValue.toString();
+            } else {
+                sObjectExtIdValue = getParameter(SOBJECT_EXT_ID_VALUE, 
exchange, USE_BODY, NOT_OPTIONAL);
+            }
         } else {
             sObjectExtIdValue = getParameter(SOBJECT_EXT_ID_VALUE, exchange, 
USE_BODY, NOT_OPTIONAL);
         }
@@ -508,7 +512,11 @@ public abstract class AbstractRestProcessor extends 
AbstractSalesforceProcessor
         final AbstractSObjectBase sObjectBase = 
exchange.getIn().getBody(AbstractSObjectBase.class);
         if (sObjectBase != null) {
             oldValue = getAndClearPropertyValue(sObjectBase, sObjectExtIdName);
-            sObjectExtIdValue = oldValue.toString();
+            if (oldValue != null) {
+                sObjectExtIdValue = oldValue.toString();
+            } else {
+                sObjectExtIdValue = getParameter(SOBJECT_EXT_ID_VALUE, 
exchange, IGNORE_BODY, NOT_OPTIONAL);
+            }
             // clear base object fields, which cannot be updated
             sObjectBase.clearBaseFields();
         } else {
@@ -536,7 +544,11 @@ public abstract class AbstractRestProcessor extends 
AbstractSalesforceProcessor
         String sObjectExtIdValue;
         if (sObjectBase != null) {
             oldValue = getAndClearPropertyValue(sObjectBase, sObjectExtIdName);
-            sObjectExtIdValue = oldValue.toString();
+            if (oldValue != null) {
+                sObjectExtIdValue = oldValue.toString();
+            } else {
+                sObjectExtIdValue = getParameter(SOBJECT_EXT_ID_VALUE, 
exchange, USE_BODY, NOT_OPTIONAL);
+            }
         } else {
             sObjectExtIdValue = getParameter(SOBJECT_EXT_ID_VALUE, exchange, 
USE_BODY, NOT_OPTIONAL);
         }
@@ -546,8 +558,8 @@ public abstract class AbstractRestProcessor extends 
AbstractSalesforceProcessor
                 new RestClient.ResponseCallback() {
                     @Override
                     public void onResponse(InputStream response, Map<String, 
String> headers, SalesforceException exception) {
-                        processResponse(exchange, response, headers, 
exception, callback);
                         restoreFields(exchange, sObjectBase, null, 
sObjectExtIdName, finalOldValue);
+                        processResponse(exchange, response, headers, 
exception, callback);
                     }
                 });
     }
diff --git 
a/components/camel-salesforce/camel-salesforce-component/src/main/java/org/apache/camel/component/salesforce/internal/processor/CompositeApiProcessor.java
 
b/components/camel-salesforce/camel-salesforce-component/src/main/java/org/apache/camel/component/salesforce/internal/processor/CompositeApiProcessor.java
index eb28d8722df9..584448170db6 100644
--- 
a/components/camel-salesforce/camel-salesforce-component/src/main/java/org/apache/camel/component/salesforce/internal/processor/CompositeApiProcessor.java
+++ 
b/components/camel-salesforce/camel-salesforce-component/src/main/java/org/apache/camel/component/salesforce/internal/processor/CompositeApiProcessor.java
@@ -115,6 +115,9 @@ public final class CompositeApiProcessor extends 
AbstractSalesforceProcessor {
 
                 out.copyFromWithNewBody(in, response);
                 out.getHeaders().putAll(headers);
+                if (exception != null) {
+                    exchange.setException(exception);
+                }
             }
         } finally {
             // notify callback that exchange is done
@@ -136,6 +139,9 @@ public final class CompositeApiProcessor extends 
AbstractSalesforceProcessor {
 
                 out.copyFromWithNewBody(in, response);
                 out.getHeaders().putAll(headers);
+                if (exception != null) {
+                    exchange.setException(exception);
+                }
             }
         } finally {
             // notify callback that exchange is done
@@ -157,6 +163,9 @@ public final class CompositeApiProcessor extends 
AbstractSalesforceProcessor {
 
                 out.copyFromWithNewBody(in, response);
                 out.getHeaders().putAll(headers);
+                if (exception != null) {
+                    exchange.setException(exception);
+                }
             }
         } finally {
             // notify callback that exchange is done
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..2542ee8d5e1a 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
@@ -219,8 +219,10 @@ public class JsonRestProcessor extends 
AbstractRestProcessor {
             String msg = "Error parsing JSON response: " + e.getMessage();
             exchange.setException(new SalesforceException(msg, e));
         } finally {
-            // cleanup temporary exchange headers
+            // cleanup temporary exchange properties
             exchange.removeProperty(RESPONSE_CLASS);
+            exchange.removeProperty(RESPONSE_CLASS_DEFERRED);
+            exchange.removeProperty(RESPONSE_CLASS_PREFIX);
             exchange.removeProperty(RESPONSE_TYPE);
 
             // consume response entity
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..3e8cef329cde 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
@@ -203,7 +203,8 @@ public class SubscriptionHelper extends ServiceSupport {
                 if (toSubscribe != null) {
                     LOG.info("Subscribing to channels: {}", toSubscribe);
                     for (var channelName : toSubscribe) {
-                        for (var consumer : 
channelToConsumers.get(channelName)) {
+                        var consumers = 
channelToConsumers.getOrDefault(channelName, emptySet());
+                        for (var consumer : consumers) {
                             subscribe(consumer);
                         }
                     }
@@ -592,6 +593,7 @@ public class SubscriptionHelper extends ServiceSupport {
                 consumers.remove(consumer);
                 if (consumers.isEmpty()) {
                     channelToConsumers.remove(channelName);
+                    channelsToSubscribe.remove(channelName);
                 }
             }
             final ClientSessionChannel.MessageListener listener = 
consumerToListener.remove(consumer);
diff --git 
a/components/camel-salesforce/camel-salesforce-component/src/test/java/org/apache/camel/component/salesforce/PubSubApiTest.java
 
b/components/camel-salesforce/camel-salesforce-component/src/test/java/org/apache/camel/component/salesforce/PubSubApiTest.java
index 9ca33f2d46d9..cb191c89c303 100644
--- 
a/components/camel-salesforce/camel-salesforce-component/src/test/java/org/apache/camel/component/salesforce/PubSubApiTest.java
+++ 
b/components/camel-salesforce/camel-salesforce-component/src/test/java/org/apache/camel/component/salesforce/PubSubApiTest.java
@@ -204,17 +204,14 @@ public class PubSubApiTest {
         final String replayId = encodeReplayId("123");
         client.subscribe(consumer, ReplayPreset.CUSTOM, replayId, false);
 
-        Thread.sleep(1000);
-
-        InOrder inOrder = Mockito.inOrder(client);
-        inOrder.verify(client, timeout(5000).times(3)).subscribe(consumer, 
ReplayPreset.CUSTOM, replayId, false);
-        inOrder.verify(client, never()).subscribe(consumer, 
ReplayPreset.LATEST, null, false);
+        // with the fix, corrupted replay ID with fallback disabled stops 
after one attempt
+        // (no infinite retry with the same corrupt ID)
+        verify(client, timeout(5000).times(1)).subscribe(consumer, 
ReplayPreset.CUSTOM, replayId, false);
+        verify(client, never()).subscribe(consumer, ReplayPreset.LATEST, null, 
false);
 
         ArgumentCaptor<InvalidReplayIdException> captor = 
ArgumentCaptor.forClass(InvalidReplayIdException.class);
-        verify(exceptionHandler, 
timeout(5000).times(3)).handleException(captor.capture());
-        for (InvalidReplayIdException exception : captor.getAllValues()) {
-            Assertions.assertEquals(replayId, exception.getReplayId());
-        }
+        verify(exceptionHandler, 
timeout(5000).times(1)).handleException(captor.capture());
+        Assertions.assertEquals(replayId, captor.getValue().getReplayId());
     }
 
     private String encodeReplayId(String replayId) {

Reply via email to