[
https://issues.apache.org/jira/browse/CAMEL-25036?page=com.atlassian.jira.plugin.system.issuetabpanels:all-tabpanel
]
Claus Ibsen updated CAMEL-25036:
--------------------------------
Fix Version/s: 4.23.0
> camel-support - EventDrivenPollingConsumer.receive() holds the service
> lifecycle lock while it waits, so stopping a ConsumerTemplate, a pollEnrich
> route or the CamelContext hangs forever (regression in 4.8)
> --------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------
>
> Key: CAMEL-25036
> URL: https://issues.apache.org/jira/browse/CAMEL-25036
> Project: Camel
> Issue Type: Bug
> Components: camel-core
> Reporter: shashank
> Priority: Minor
> Fix For: 4.23.0
>
>
> {{EventDrivenPollingConsumer.receive()}} and {{receive(timeout)}}
> ({{core/camel-support/.../EventDrivenPollingConsumer.java:128-153}},
> {{:156-178}}) wrap {{beforePoll}} / {{queue.take()}} (or
> {{queue.poll(timeout)}}) / {{afterPoll}} in {{lock.lock()}} ...
> {{lock.unlock()}}. That {{lock}} is not a field of the class: it is the
> inherited {{BaseService.lock}} ({{core/camel-api/.../BaseService.java:58}}),
> the lock that {{BaseService.start()}}, {{stop()}}, {{suspend()}},
> {{resume()}} and {{shutdown()}} take ({{:114}}, {{:159}}, {{:195}}, {{:227}},
> ...).
> So while a thread waits in {{receive()}} for a message, nobody can stop the
> polling consumer: {{stop()}} blocks on the lock until a message arrives. On
> an idle endpoint that is forever.
> This came in with commit ee18bef57f (part of the CAMEL-20199 virtual-threads
> work, released in 4.8.0), which replaced {{synchronized (this)}} in
> {{receive}} with {{lock.lock()}}. Before that the receive monitor ({{this}})
> and the lifecycle monitor ({{BaseService.lock}}, then a plain {{Object}})
> were different, so {{stop()}} did not wait for {{receive()}}. Affected:
> 4.8.x, 4.14.x, 4.18.x, main.
> Every endpoint that does not override {{createPollingConsumer()}} uses
> {{EventDrivenPollingConsumer}}, as do consumers that implement
> {{PollingConsumerPollingStrategy}}. {{GenericFilePollingConsumer}} (file,
> ftp, sftp, smb) extends it but only takes the lock once a file is already
> queued, so it does not hang this way. Seda has its own
> {{SedaPollingConsumer}} and is not affected.
> Consequences:
> # {{pollEnrich}} with the default timeout (-1, {{receive()}}) on an endpoint
> that has no message: stopping the route or the CamelContext never returns.
> The graceful shutdown times out, and the forced shutdown then blocks in
> {{PollEnricher.doStop}} -> {{DefaultConsumerCache.doStop}} ->
> {{ServicePool.stop}} -> {{EventDrivenPollingConsumer.stop()}}.
> # {{ConsumerTemplate.receive(uri)}} in one thread:
> {{consumerTemplate.stop()}} and {{camelContext.stop()}} hang.
> # The same lock also makes {{receiveNoWait()}} and {{receive(uri, timeout)}}
> of other threads wait until the blocked {{receive()}} gets a message, because
> they share the polling consumer through the consumer cache.
> {{receiveNoWait()}} does not return at once, and {{receive(timeout)}} ignores
> its timeout. This part is older, since {{synchronized (this)}} serialized
> them too, and it has the same fix.
> *Reproduction* (4.23.0-SNAPSHOT, {{direct:idle}} = a DefaultEndpoint, so
> EventDrivenPollingConsumer):
> {noformat}
> template: thread A: consumerTemplate.receive("direct:idle") -> waiting in
> EventDrivenPollingConsumer.receive:141
> consumerTemplate.stop() STILL BLOCKED after 4000 ms at
> BaseService.stop:159 <- ServicePool$SinglePool.doStop:283 <-
> DefaultConsumerCache.doStop:271
> camelContext.stop() STILL BLOCKED after 4000 ms
> route: from("direct:start").pollEnrich("direct:idle"), shutdown timeout 2
> s, one exchange waiting in pollEnrich
> camelContext.stop() STILL BLOCKED after 10000 ms at
> BaseService.stop:159 <- ... <- DefaultConsumerCache.doStop:271 <-
> PollEnricher.doStop:667
> nowait: thread A blocked in receive("direct:idle"); on the same
> ConsumerTemplate:
> receiveNoWait("direct:idle") STILL BLOCKED after 4000 ms (returned
> null after 8022 ms, once a message was sent for A)
> receive("direct:idle", 500) STILL BLOCKED after 4000 ms
> {noformat}
> The same harness with the 4.7.0 {{EventDrivenPollingConsumer}} (control):
> {{consumerTemplate.stop()}} returns after 2 ms, {{camelContext.stop()}} after
> 8 ms, and the route case stops after the 2 s shutdown timeout. With the fix
> below all three return promptly: stop in 1 ms, {{receiveNoWait}} null in 0
> ms, {{receive(500)}} null in 505 ms, and the route context stops after 2017
> ms.
> A TLA+ model with receive(), receiveNoWait(), stop() and the delegate
> consumer. For the current code, {{StopTerminates}} is violated
> ({{pc_shared_stop_live}}: RLoop -> RAcq, i.e. waiting in take, then stop
> waits forever) and so is {{NoWaitTerminates}} ({{pc_shared_nowait_live}}).
> The pre-4.8 variant satisfies {{StopTerminates}}. The fix variant satisfies
> all properties ({{pc_fix_all}}, {{pc_fix_interrupt}}).
> *Proposed fix:*
> * Use a dedicated lock instead of the inherited lifecycle lock, held only
> while {{beforePoll}}/{{afterPoll}} run and a count of the receive calls in
> progress is updated, never while waiting for a message. {{afterPoll}} is only
> run by the last receive call in progress, which keeps the CAMEL-10215
> guarantee (one caller's {{afterPoll}} must not strand another caller that is
> still polling). This restores the pre-4.8 behaviour for stop.
> * Make {{receive()}} return when the consumer stops: wait in slices (for
> example {{queue.poll(1s)}} in the {{while (isRunAllowed())}} loop), so that a
> stopped consumer returns null, which is what the existing "Consumer is not
> running, so returning null" path intends.
> * {{receive(timeout)}} keeps its single timed poll, and no longer waits
> behind other callers.
> Suggested test: a ConsumerTemplate receive() on an idle endpoint in a thread,
> then {{context.stop()}} must return; a pollEnrich with the default timeout on
> an idle endpoint, then stopping the context must return; and receiveNoWait()
> from a second thread must return at once.
> A companion ticket, filed together with this one, covers an interrupted
> {{receive()}} that spins forever.
> _Filed with Claude Code on behalf of allthingssecurity._
--
This message was sent by Atlassian Jira
(v8.20.10#820010)