allthingssecurity commented on code in PR #27002:
URL: https://github.com/apache/camel/pull/27002#discussion_r4131019091


##########
core/camel-core-processor/src/main/java/org/apache/camel/processor/OnCompletionProcessor.java:
##########
@@ -188,27 +197,65 @@ public boolean process(Exchange exchange, AsyncCallback 
callback) {
     }
 
     /**
-     * Submits the onCompletion task to the thread pool (parallel processing). 
The task is counted as pending from when
-     * it is submitted until it is done, so a graceful shutdown waits for it.
+     * Submits the onCompletion task of the given copy to the thread pool 
(parallel processing). The task is counted as
+     * pending from when it is submitted until it is done, so a graceful 
shutdown waits for it.
+     * <p>
+     * The copy may hold its own reference to a stream cache (see {@link 
#prepareExchange(Exchange)}), which is released
+     * when the copy is done. When the task never runs (the thread pool 
rejects it, discards it because it is shut down,
+     * or drops it when the processor shuts its thread pool down), it is no 
longer counted as pending and the reference
+     * is released instead, as otherwise a spooled file would be kept until 
the stream caching strategy is stopped.
      */
     @SuppressWarnings("deprecation")
-    private void submitTask(Runnable task) {
+    private void submitTask(Exchange copy, Runnable task) {
+        ParallelTask parallelTask = new ParallelTask(copy, task);
         taskCount.increment();
-        Runnable counted = () -> {
-            try {
-                task.run();
-            } finally {
-                taskCount.decrement();
-            }
-        };
+        pendingTasks.add(parallelTask);
         try {
             // Deprecated since 4.19.0
-            executorService.submit(prepareMDCParallelTask(camelContext, 
counted));
+            executorService.submit(prepareMDCParallelTask(camelContext, 
parallelTask));
         } catch (RuntimeException e) {
             // the task will not run
-            taskCount.decrement();
+            parallelTask.discard();
             throw e;
         }
+        if (executorService.isShutdown()) {

Review Comment:
   Good catch, fixed in 23854821e85f.
   
   I didn't use the `getQueue()` check, because a task that isn't in the queue 
can still run: a `ThreadPoolExecutor` can hand the task to a new worker as its 
first task, or a worker may have already taken it from the queue. Discarding in 
either case would lose the onCompletion the same way.
   
   Instead, `submitTask` now checks `isShutdown()` before `submit`. A pool that 
is already shut down rejects every task, so only that case discards after 
`submit` (a no-op if the rejection policy ran the task). A task accepted before 
a graceful `shutdown()` is left alone and runs. With `shutdownNow` it is still 
discarded in `doShutdown`.
   
   What's left is a smaller window: the pool is shut down while `submit` is 
running and silently rejects the task (the TPE re-check after the enqueue). 
That task stays pending until the processor shuts its own pool down, and the 
Javadoc says so. I chose this so that a task that may still run is never 
dropped.
   
   New test `testParallelTaskAcceptedBeforeShutdown` uses a pool that shuts 
itself down right after it accepts the task. It fails on the previous head (the 
task is discarded and never runs) and passes now. The camel-core 
`*OnCompletion*` tests pass: 91 run, 0 failures.
   
   _Claude Code on behalf of allthingssecurity_



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

Reply via email to