[
https://issues.apache.org/jira/browse/CAMEL-24958?page=com.atlassian.jira.plugin.system.issuetabpanels:all-tabpanel
]
Claus Ibsen resolved CAMEL-24958.
---------------------------------
Fix Version/s: 4.23.0
Resolution: Fixed
The fix is merged on main, so it is in Camel 4.23.0:
* 69c06886d5be CAMEL-24958: camel-core - Release the producers of recipients
the Recipient List did not send to
* 805e314dde29 CAMEL-24958: camel-core - Emit no exchange events for a
recipient the Recipient List skipped
* ce84a7632898 CAMEL-24958: camel-core - Make the new test class package-private
Resolving, as the ticket was not updated when the PR was merged.
_Claude Code on behalf of Claus Ibsen_
> Recipient List does not release the producers (and prototype endpoints) of
> recipients that were not sent to (stopOnException, timeout)
> --------------------------------------------------------------------------------------------------------------------------------------
>
> Key: CAMEL-24958
> URL: https://issues.apache.org/jira/browse/CAMEL-24958
> Project: Camel
> Issue Type: Bug
> Components: camel-core
> Reporter: shashank
> Priority: Minor
> Fix For: 4.23.0
>
>
> {{RecipientListProcessor.createProcessorExchangePairs}} acquires a producer
> for every recipient up front ({{producerCache.acquireProducer}}). The
> producer is released, and a prototype endpoint ({{cacheSize=-1}}) is stopped,
> only in {{RecipientProcessorExchangePair.done()}}, which
> {{MulticastProcessor.afterSend}} calls for recipients that were actually sent
> to. When the list stops early, the remaining recipients are never sent to and
> their producers are never released:
> * {{stopOnException}} and a recipient fails;
> * the timeout fires before a recipient is sent;
> * the task is rejected.
> {{MulticastProcessor.doDone}} only releases the pair exchanges. So:
> * non-singleton producers (ftp/sftp/ftps, smb, ssh, mina sync, ...): every
> failed exchange starts a new producer, which is never returned to the pool;
> * {{cacheSize=-1}}: the prototype endpoint and its producer are never
> stopped, not even when the CamelContext stops.
> Reproduced against 4.23.0-SNAPSHOT with a test component whose endpoint has
> {{isSingletonProducer()=false}} and counts starts/stops (50 exchanges per
> line, {{recipientList(header("to")).stopOnException()}}):
> {noformat}
> [pooled:a, direct:boom] (all sent) producers started=1
> stopped=0
> [direct:boom, pooled:b] (b not sent) producers started=51
> stopped=0 <- 50 extra producers
> cacheSize(-1) [pooled:dN] (all sent) endpoints started=+50
> stopped=+50 (control)
> cacheSize(-1) [direct:boom, pooled:cN] (c not sent) endpoints started=+50
> stopped=+0
> after CamelContext.stop() producers started=151
> stopped=101, endpoints started=102 stopped=52
> {noformat}
> Note: {{timeout}} requires {{parallelProcessing}} on a Recipient List
> ("Timeout is used but ParallelProcessing has not been enabled"), so the
> timeout case is a parallel one.
> Proposed fix (all in {{RecipientListProcessor}}):
> * {{RecipientProcessorExchangePair}} gets a state: NEW, BEGUN, DONE or
> RELEASED, changed only by compare-and-set.
> ** {{begin()}} claims NEW -> BEGUN. If the claim fails, the pair was already
> released, and {{getProcessor()}} then returns a no-op processor, so the
> recipient is not sent.
> ** {{done()}} releases only on BEGUN -> DONE, which makes it idempotent.
> ** A new {{releaseIfNotBegun()}} releases on NEW -> RELEASED.
> ** The release itself is unchanged: {{producerCache.releaseProducer(...)}},
> and stopping the prototype endpoint.
> * {{RecipientListProcessor}} overrides {{doDone(...)}}. It calls
> {{releaseIfNotBegun()}} on every pair (the pairs of a Recipient List are
> always a {{List}}), then {{super.doDone(...)}}. So the producers are released
> before the callback continues the route.
> * {{MulticastProcessor}}, {{Splitter}} and the {{ProcessorExchangePair}}
> interface are unchanged.
> Behaviour change: with {{parallelProcessing}}, a recipient whose task had not
> started yet when the Recipient List completed (stopOnException, timeout) is
> now skipped. Before, it was still sent afterwards. Recipients that had
> already started keep running, as before. There is an upgrade-guide entry.
> A PR with the fix follows, with regression test
> {{RecipientListReleaseUnsentProducerTest}}. It uses a component with a
> non-singleton producer that counts producer starts and endpoint starts/stops.
> Test methods:
> * {{testStopOnException}} (sequential)
> * {{testStopOnExceptionParallel}} (single-thread pool)
> * {{testPrototypeEndpointStopped}} ({{cacheSize(-1)}})
> * {{testTimeout}} (parallel, the single pool thread is blocked by the first
> recipient)
> Without the fix they fail with {{expected: <1> but was: <5>}}, {{expected:
> <1> but was: <5>}}, {{expected: <5> but was: <0>}} and {{expected: <1> but
> was: <2>}}.
> Found with a TLA+ model of the multicast task, then reproduced against
> 4.23.0-SNAPSHOT.
--
This message was sent by Atlassian Jira
(v8.20.10#820010)