davsclaus commented on code in PR #27477:
URL: https://github.com/apache/camel/pull/27477#discussion_r4204456861


##########
components/camel-ai/camel-docling/src/main/docs/docling-component.adoc:
##########
@@ -1219,6 +1219,33 @@ YAML::
 ----
 ====
 
+== Consuming asynchronous conversions
+
+Instead of submitting a conversion with `SUBMIT_ASYNC_CONVERSION` and then 
polling `CHECK_CONVERSION_STATUS` by hand, a
+`docling:` consumer can emit each async conversion event-driven, as soon as it 
finishes:
+
+[source,java]
+----
+// submit conversions asynchronously; each returns a task id and runs in the 
background
+from("direct:submit")
+    
.to("docling:convert?useDoclingServe=true&operation=SUBMIT_ASYNC_CONVERSION");
+
+// the consumer emits one exchange per completed task (the path, here 
"onComplete", is just a label)
+from("docling:onComplete?useDoclingServe=true&delay=5000")
+    .to("mock:done");
+----
+
+The consumer scans the tasks submitted via `SUBMIT_ASYNC_CONVERSION` on the 
same component and emits one exchange as

Review Comment:
   **The consumer changes what `CHECK_CONVERSION_STATUS` returns.** The 
consumer removes every completed task from the component-wide map. After that, 
`CHECK_CONVERSION_STATUS` for a local id such as `task-3` no longer finds it 
locally and falls back to `pollTaskStatus` on docling-serve. The server doesn't 
know local ids, so it returns `FAILED` with an error message. In a context that 
has a consumer, polling by hand therefore reports successful conversions as 
failed.
   
   Could the docs say that the consumer and manual `CHECK_CONVERSION_STATUS` 
polling shouldn't be mixed on the same component? Alternatively, the fallback 
could know which ids the consumer already took.



##########
components/camel-ai/camel-docling/src/main/java/org/apache/camel/component/docling/DoclingConsumer.java:
##########
@@ -0,0 +1,106 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements.  See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License.  You may obtain a copy of the License at
+ *
+ *      http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+package org.apache.camel.component.docling;
+
+import java.util.Map;
+import java.util.concurrent.CancellationException;
+import java.util.concurrent.CompletableFuture;
+import java.util.concurrent.CompletionException;
+
+import ai.docling.serve.api.convert.response.ConvertDocumentResponse;
+import org.apache.camel.Exchange;
+import org.apache.camel.Processor;
+import org.apache.camel.support.ScheduledPollConsumer;
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
+
+/**
+ * Event-driven consumer for docling async conversions.
+ * <p>
+ * It polls the component-wide registry of tasks submitted via {@code 
SUBMIT_ASYNC_CONVERSION} and, each time a task
+ * completes, emits one exchange whose body is the converted content (the same 
shape the synchronous {@code CONVERT_*}
+ * operations produce for the configured {@code outputFormat}), carrying the 
task id in the
+ * {@link DoclingHeaders#TASK_ID} header. A task that failed completes the 
exchange with the underlying exception, so
+ * the route's {@code onException} machinery applies. The poll interval is the 
standard scheduled-consumer {@code delay}
+ * option.
+ */
+public class DoclingConsumer extends ScheduledPollConsumer {
+
+    private static final Logger LOG = 
LoggerFactory.getLogger(DoclingConsumer.class);
+
+    public DoclingConsumer(DoclingEndpoint endpoint, Processor processor) {
+        super(endpoint, processor);
+    }
+
+    @Override
+    public DoclingEndpoint getEndpoint() {
+        return (DoclingEndpoint) super.getEndpoint();
+    }
+
+    @Override
+    protected int poll() throws Exception {
+        Map<String, AsyncTaskEntry> pendingAsyncTasks
+                = ((DoclingComponent) 
getEndpoint().getComponent()).getPendingAsyncTasks();
+
+        int polled = 0;
+        for (Map.Entry<String, AsyncTaskEntry> entry : 
pendingAsyncTasks.entrySet()) {
+            String taskId = entry.getKey();
+            AsyncTaskEntry task = entry.getValue();
+            if (!task.getFuture().isDone()) {
+                continue;
+            }
+            // claim the completed task atomically, so a concurrent 
CHECK_CONVERSION_STATUS or another consumer cannot
+            // emit the same completion twice
+            if (!pendingAsyncTasks.remove(taskId, task)) {
+                continue;
+            }
+
+            polled++;
+            Exchange exchange = createExchange(false);
+            try {
+                exchange.getIn().setHeader(DoclingHeaders.TASK_ID, taskId);
+                emitResult(exchange, taskId, task.getFuture());
+                getProcessor().process(exchange);
+            } catch (Exception e) {
+                exchange.setException(e);
+            } finally {
+                if (exchange.getException() != null) {
+                    getExceptionHandler().handleException(
+                            "Error processing docling async task " + taskId, 
exchange, exchange.getException());
+                }
+                releaseExchange(exchange, false);
+            }
+        }
+        return polled;
+    }
+
+    private void emitResult(Exchange exchange, String taskId, 
CompletableFuture<ConvertDocumentResponse> future) {
+        try {
+            ConvertDocumentResponse response = future.join();
+            exchange.getIn().setBody(
+                    DoclingContentExtractor.extract(response, 
getEndpoint().getConfiguration().getOutputFormat()));

Review Comment:
   **Output format mismatch, which can give an empty body.** The producer 
submits the async request with the format from the `CamelDoclingOutputFormat` 
header, falling back to the producer endpoint's `outputFormat` 
(`processSubmitAsyncConversion`). Here the consumer extracts with *its own* 
endpoint's `outputFormat` (default markdown).
   
   Example: a task submitted with header `OutputFormat=html`, or from 
`docling:convert?outputFormat=html&operation=SUBMIT_ASYNC_CONVERSION`, and 
consumed by `docling:onComplete?useDoclingServe=true`. docling-serve returns 
only HTML, `getMarkdownContent()` is `null`, and the extractor returns `""`. 
The route gets an empty body with no error.
   
   Suggestion: store the format on the task when it is submitted, and use it 
here:
   
   ```java
   // AsyncTaskEntry: new field + getter
   private final String outputFormat;
   
   // DoclingProducer.processSubmitAsyncConversion
   AsyncTaskEntry taskEntry = new AsyncTaskEntry(taskId, asyncResult, 
outputFormat);
   
   // DoclingConsumer.emitResult
   DoclingContentExtractor.extract(response, task.getOutputFormat())
   ```
   
   `checkLocalAsyncTask` in the producer has the same pre-existing gap (it uses 
`configuration.getOutputFormat()`), so the same field fixes both. A unit test 
that submits with one format and consumes from an endpoint with the default 
format would cover it.



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