oscerd commented on code in PR #27477: URL: https://github.com/apache/camel/pull/27477#discussion_r4204674749
########## 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: Good catch — fixed in 4abae810a1d9. The output format is now resolved once at submit time (the `CamelDoclingOutputFormat` header, else the submitting endpoint's `outputFormat`), stored on `AsyncTaskEntry`, and used for extraction in both the consumer (`emitResult`) and the producer's `checkLocalAsyncTask` — so neither re-derives it from its own endpoint. That closes the pre-existing `checkLocalAsyncTask` gap too, as you noted. Added `DoclingConsumerTest.extractsWithTheSubmittedFormatNotTheConsumerFormat`: a task submitted with `html` and consumed from a markdown-default `docling:onComplete` endpoint now emits the HTML body instead of an empty string. _Claude Code on behalf of oscerd_ ########## 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: Documented in 4abae810a1d9. The async-consumer section now carries a NOTE stating that the consumer drains the whole component's tasks, so it shouldn't be mixed with manual `CHECK_CONVERSION_STATUS` polling on the same component — once the consumer claims a task, a later status check for that id no longer finds it locally and falls back to the server, which doesn't know the local id. I kept the draining behaviour (it's the point of the consumer) and documented the constraint, rather than tracking consumed ids on the fallback — that would add state for a usage pattern the docs now steer away from. Happy to revisit the code-side option if you'd prefer it. _Claude Code on behalf of oscerd_ -- 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]
