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


##########
components/camel-reactive-streams/src/main/java/org/apache/camel/component/reactive/streams/ReactiveStreamsConsumer.java:
##########
@@ -61,17 +75,45 @@ protected void doStart() throws Exception {
                     getEndpoint().getEndpointUri(), poolSize);
         }
 
-        this.service.attachCamelConsumer(endpoint.getStream(), this);
+        // the items left queued by a stop that timed out are routed now
+        scheduleQueuedItems();
+
+        this.subscriber = 
this.service.attachCamelConsumer(endpoint.getStream(), this);
+    }
+
+    @Override
+    protected void doSuspend() throws Exception {
+        // nothing to wait for: while the consumer is suspended, its 
subscriber requests no more items from the stream
+        // and the queued items are not routed (see routeQueuedItem); the 
exchanges being routed complete normally
+    }
+
+    @Override
+    protected void doResume() throws Exception {
+        if (executor == null) {
+            // suspended while it was not started (before its start or after a 
stop)
+            doStart();

Review Comment:
   Added the comment in c92970f2. One correction to the reasoning: `resume()` 
sets the status to `STARTING` before calling `doResume()`, so the transition 
does happen, and that is exactly why `start()` cannot be used here (it returns 
early when the status is `STARTING`); the comment says that instead. Module 
suite: 75 tests, 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