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]