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 {
