[ 
https://issues.apache.org/jira/browse/CAMEL-25036?page=com.atlassian.jira.plugin.system.issuetabpanels:all-tabpanel
 ]

Claus Ibsen resolved CAMEL-25036.
---------------------------------
    Resolution: Fixed

Fixed by https://github.com/apache/camel/pull/26911 (merged to main for 4.23.0).

_Claude Code on behalf of davsclaus_

> 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)

Reply via email to