allthingssecurity opened a new pull request, #26911: URL: https://github.com/apache/camel/pull/26911
# Description [CAMEL-25036](https://issues.apache.org/jira/browse/CAMEL-25036) `EventDrivenPollingConsumer.receive()` and `receive(timeout)` wrap `beforePoll` / `queue.take()` (or `queue.poll(timeout)`) / `afterPoll` in `lock.lock()`. That `lock` is the inherited `BaseService` lifecycle lock, the one `start()`, `stop()`, `suspend()`, `resume()` and `shutdown()` take. It came in with commit ee18bef57f (part of the CAMEL-20199 virtual-threads work, released in 4.8.0), which replaced `synchronized (this)` with `lock.lock()` to avoid pinning virtual threads. Before that, the receive monitor (`this`) and the lifecycle monitor were different objects, so the lifecycle lock was never meant to guard `receive()`. So while a thread waits in `receive()` on an idle endpoint, the polling consumer cannot be stopped: - `pollEnrich` with the default timeout: stopping the route or the CamelContext never returns. The forced shutdown blocks in `PollEnricher.doStop` -> `DefaultConsumerCache.doStop` -> `ServicePool.stop` -> `BaseService.stop`. - `ConsumerTemplate.receive(uri)` in one thread: `consumerTemplate.stop()` and `camelContext.stop()` hang. - `receiveNoWait()` and `receive(uri, timeout)` from other threads, which share the cached polling consumer, wait until the blocked `receive()` gets a message. This part is older, since `synchronized (this)` serialized them too. Every endpoint that does not override `createPollingConsumer()` uses this class, 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. Affected: 4.8.x and later (the 4.7.0 class uses its own monitor and stops in a few ms). This change: - A dedicated `ReentrantLock` replaces the lifecycle lock, so there is still no `synchronized` block (the goal of CAMEL-20199). It is held only while `beforePoll`/`afterPoll` run and a count of the receive calls in progress is updated. It is never held while waiting for a message. - CAMEL-10215 serialized `beforePoll`, poll and `afterPoll` so that one caller's `afterPoll` (which suspends a scheduled delegate consumer) could not strand another caller that is still polling. Receive calls now wait on the queue concurrently, and `afterPoll` is only invoked by the last receive call in progress. A plain `tryLock` would have been simpler, but then a `receiveNoWait()` could return null while a file it had just polled sat in the queue behind another caller. - `receive()` waits in slices of 1 second and returns null once the consumer is no longer running. This is what the existing "Consumer is not running, so returning null" path intended. - `receive(timeout)` keeps its single timed poll, and its wait no longer depends on other callers. - With overlapping callers, `beforePoll` can now run more than once before the single `afterPoll` of the last caller. `PollingConsumerPollingStrategy` promises no pairing, and its only implementation, `ScheduledPollConsumer`, tolerates it (its `beforePoll` starts or resumes the scheduler, which is idempotent). - A `stop()` that races with `beforePoll` can, as before 4.8, see the scheduled delegate restarted and then suspended by the last caller's `afterPoll`; `doShutdown` stops it. Tightening that small window is possible but out of scope here. No upgrade guide entry: the visible change is that `stop()` works again and a waiting `receive()` returns null when stopped. Tests: new `EventDrivenPollingConsumerStopTest`. It stops the CamelContext while a pollEnrich with the default timeout waits on an idle endpoint, stops a ConsumerTemplate while its `receive()` waits, and calls `receiveNoWait()` and `receive(uri, 100)` while another thread's `receive()` waits. That last test also checks that `afterPoll` is not invoked while that `receive()` is still in progress. The tests are bounded with `assertTimeoutPreemptively` and Awaitility, with no sleeps. Without the main-code change all three fail: ``` testStopContextWhilePollEnrichWaits execution timed out after 20000 ms testStopConsumerTemplateWhileReceiveWaits execution timed out after 10000 ms testReceiveNoWaitWhileReceiveWaits execution timed out after 5000 ms ``` With the change, `*PollEnrich*,*PollingConsumer*,*ConsumerTemplate*,*Enricher*,*ConsumerCache*,*ScheduledPoll*` in camel-core and camel-support pass: 130 tests, 0 failures. Found with a TLA+ model of `receive()`, `receiveNoWait()`, `stop()` and the delegate consumer. For the current code `StopTerminates` and `NoWaitTerminates` are violated; the pre-4.8 variant satisfies `StopTerminates`, and the fixed variant satisfies all properties. I then reproduced the hang against the real classes and checked the 4.7.0 class as a control. See also #PR_111 (an interrupted `receive()` spins forever), which builds on this change. # Target - [x] I checked that the commit is targeting the correct branch (Camel 4 uses the `main` branch) # Tracking - [x] If this is a large change, bug fix, or code improvement, I checked there is a [JIRA issue](https://issues.apache.org/jira/browse/CAMEL) filed for the change (usually before you start working on it). # Apache Camel coding standards and style - [x] I checked that each commit in the pull request has a meaningful subject line and body. - [ ] I have run `mvn clean install -DskipTests` locally from root folder and I have committed all auto-generated changes. (I built and tested the affected modules, including the formatter and import-sort plugins. I did not run the full root build.) # AI-assisted contributions - [x] If this PR includes AI-generated code, commits have proper co-authorship attribution (e.g., `Co-authored-by` trailers) and the PR description identifies the AI tool used. This PR was prepared with Claude Code (Claude Opus 5.5). The commit carries a `Co-Authored-By` trailer. _Claude Code on behalf of allthingssecurity_ 🤖 Generated with [Claude Code](https://claude.com/claude-code) -- This is an automated message from the Apache Git Service. To respond to the message, please log on to GitHub and use the URL above to go to the specific comment. To unsubscribe, e-mail: [email protected] For queries about this service, please contact Infrastructure at: [email protected]
