jiriOndrusek-agent commented on code in PR #9290:
URL: https://github.com/apache/camel-quarkus/pull/9290#discussion_r4218001134


##########
extensions/langchain4j-ingest/runtime/src/main/java/org/apache/camel/quarkus/component/langchain4j/ingest/IngestRoutes.java:
##########
@@ -143,59 +201,344 @@ private IngestPipelineDefinition 
builderPipeline(IngestBuilderPipelines.Entry en
 
         IngestPipeline definition = builderPipelines.definition(entry);
         String uri = "file".equals(definition.sourceType()) ? null : 
definition.sourceUri();
-        return definition(name, uri, definition.asRunTimeConfig(),
+        compositionRoute(name, uri, definition.asRunTimeConfig(),
                 resolveStore(name, 
definition.embeddingStoreName().orElse(null)),
                 resolveModel(name, 
definition.embeddingModelName().orElse(null)),
                 definition.maxSegmentSize(),
                 definition.maxOverlapSize(),
                 definition.embeddingBatchSize(),
                 definition.maxDocumentSize(),
                 definition.documentSplitterName().orElse(null),
-                definition.parser().orElse(null));
+                definition.parser().orElse(null),
+                definition.modality(),
+                definition.contentType().orElse(null));
     }
 
     /**
      * One creation path for both declaration styles, mirroring how a builder 
pipeline provides
      * the same runtime-configuration view the configuration path reads.
      */
-    private IngestPipelineDefinition definition(String name, String uri,
+    private void compositionRoute(String name, String uri,
             IngestRunTimeConfig.PipelineRunTimeConfig runtime,
             EmbeddingStore<TextSegment> store, EmbeddingModel model,
             int maxSegmentSize, int maxOverlapSize, int embeddingBatchSize, 
int maxDocumentSize,
-            String documentSplitterName, String parser) {
+            String documentSplitterName, String parser, String modality, 
String contentType) {
+
+        // the name is substituted into Kamelet URIs and registry references, 
so both declaration
+        // styles are held to one charset; @Ingest names are already checked 
at build time
+        if (!name.matches("[A-Za-z0-9._-]+")) {
+            throw new IllegalStateException(
+                    "Ingestion pipeline name '" + name + "' may only contain 
letters, digits, '.', '_' and '-'");
+        }
+        // the component validates modality and content type; that a media 
document is never
+        // parsed is a rule of this composition, where the parser action sits 
before the sink.
+        // A configured pipeline fails the build on both already
+        boolean media = "media".equalsIgnoreCase(modality);
+        if (media && parser != null) {
+            throw new IllegalStateException("Ingestion pipeline '" + name
+                    + "' sets modality 'media' together with a parser. A media 
document is embedded whole and"
+                    + " never parsed; remove one of them.");
+        }
+        String storeRef = bindInstance(name, "store", store);
+        String modelRef = bindInstance(name, "model", model);
+        String documentId = runtime == null ? null : 
runtime.source().documentId().orElse(null);
+        String repositoryRef = repositoryRef(name, runtime, uri == null);
+
+        // options at their defaults are passed all the same, here and in the 
other steps: an
+        // application-wide camel.kamelet.<kamelet>.* property would otherwise 
fill them in too.
+        // An unset string option goes as "", which the component reads as 
unset. An unset bean
+        // option (documentFilter, documentSplitter, idempotentRepository) is 
left out on purpose:
+        // "" converts to no bean and fails the start, so such a property 
still reaches it, as
+        // usage.adoc warns
+        Map<String, Object> sink = new LinkedHashMap<>();
+        sink.put("pipelineName", name);
+        sink.put("documentIdHeader", LangChain4jIngestHeaders.DOCUMENT_ID);
+        sink.put("maxSegmentSize", String.valueOf(maxSegmentSize));
+        sink.put("maxOverlapSize", String.valueOf(maxOverlapSize));
+        sink.put("embeddingBatchSize", String.valueOf(embeddingBatchSize));
+        sink.put("modality", modality);
+        sink.put("contentType", contentType == null ? "" : contentType);
+        // 0, no limit, is the endpoint's default too
+        sink.put("maxDocumentSize", String.valueOf(maxDocumentSize));
+        sink.put("minDocumentSize", "0");
+        if (documentSplitterName != null) {
+            sink.put("documentSplitter", "#bean:" + documentSplitterName);
+        }
+        // filter.* comes from configuration for both declaration styles, the 
builder having no
+        // filter API, and the component enforces it
+        IngestRunTimeConfig.PipelineRunTimeConfig configured = 
runTimeConfig.pipelines().get(name);
+        String includeId = configured == null ? "" : 
configured.filter().includeId().orElse("");
+        String excludeId = configured == null ? "" : 
configured.filter().excludeId().orElse("");
+        sink.put("includeId", includeId);
+        sink.put("excludeId", excludeId);
+        if (configured != null) {
+            if (configured.filter().minDocumentSize() != 0) {
+                sink.put("minDocumentSize", 
String.valueOf(configured.filter().minDocumentSize()));
+            }
+            configured.filter().documentFilter().ifPresent(bean -> 
sink.put("documentFilter", "#bean:" + bean));
+        }
+        sink.put("embeddingStore", "#bean:" + storeRef);
+        sink.put("embeddingModel", "#bean:" + modelRef);
 
-        IngestPipelineDefinition definition;
         if (uri == null) {
             String directory = required(name, runtime == null ? null : 
runtime.source().directory().orElse(null),
                     "source.directory");
-            definition = IngestPipelineDefinition.directory(name, directory)
-                    .recursive(runtime.source().recursive());
+            // substituted into the file-source Kamelet's endpoint URI: a '?' 
or '#' could
+            // inject consumer options - delete=true would consume the user's 
documents
+            if (directory.contains("?") || directory.contains("#")) {
+                throw new IllegalStateException("Ingestion pipeline '" + name
+                        + "': the directory must not contain '?' or '#' (got 
'" + directory + "')");
+            }
+            Map<String, Object> source = new LinkedHashMap<>();
+            source.put("directory", directory);
+            source.put("recursive", 
String.valueOf(runtime.source().recursive()));
+            // the file consumer's own default poll delay
+            source.put("delay", "500");
+            if (parser == null && !media) {
+                // text is read as UTF-8; a parser or a media model receives 
the raw bytes instead -
+                // the format is its business, and a charset conversion would 
corrupt a binary document
+                source.put("charset", "UTF-8");
+            }
+            source.put("idempotentRepository", "#bean:" + repositoryRef);
+
+            ProcessorDefinition<?> route = from(kameletUri(name, 
FILE_SOURCE_KAMELET, "source", source))
+                    .routeId("langchain4j-ingest-" + name);
+            if (documentId != null) {
+                // override the source's file-name default; captured before 
any further step
+                route = route.setHeader(LangChain4jIngestHeaders.DOCUMENT_ID, 
documentIdExpression(documentId));
+            }
+            // no duplicate pre-check: the source register already filtered 
duplicates out
+            route = parseSteps(route, name, parser, maxDocumentSize, true, 
null, includeId, excludeId);
+            // the file consumer discards the reply and the source register 
already keeps the
+            // same file version from being ingested twice, so no repository 
goes to the sink.
+            // The discarded reply would hide an empty outcome, so it is 
logged: warned with a
+            // parser, since a parse to nothing typically means a missing Tika 
parser module or an
+            // image-only document, and the file's register key is committed. 
A filtered one is
+            // logged at DEBUG
+            route.to(kameletUri(name, SINK_KAMELET, "sink", 
sink)).process(exchange -> {
+                IngestResult result = 
exchange.getMessage().getBody(IngestResult.class);
+                if (result != null && result.outcome() == 
IngestResult.Outcome.FILTERED) {
+                    logFiltered(name, result.documentId());
+                }
+                if (result == null || result.outcome() != 
IngestResult.Outcome.EMPTY) {
+                    return;
+                }
+                if (parser != null) {
+                    LOG.warnf("Ingestion pipeline '%s': document '%s' parsed 
to no text and was skipped; its key is"
+                            + " committed, so it is not retried until the file 
changes (missing parser module?"
+                            + " image-only document?)", name, 
result.documentId());
+                } else {
+                    LOG.debugf("Ingestion pipeline '%s': document '%s' 
contained no text, nothing was written",
+                            name, result.documentId());
+                }
+            });
+            LOG.infof("Ingestion pipeline '%s': source=file:%s", name, 
directory);
         } else {
-            definition = IngestPipelineDefinition.consumer(name, uri);
+            if (!uri.contains(":")) {
+                throw new IllegalStateException(
+                        "Ingestion pipeline '" + name + "': '" + uri + "' is 
not a consumer URI");
+            }
+            ProcessorDefinition<?> route = 
from(uri).routeId("langchain4j-ingest-" + name);
+            if (documentId != null && 
!IngestHeaders.DOCUMENT_ID.equals(documentId)) {
+                // copied into the canonical header the actions read, as on a 
directory pipeline:
+                // evaluated against the exchange the consumer delivered, 
before any parse, and a
+                // plain header name is read directly, never substituted into 
an expression
+                route = route.setHeader(LangChain4jIngestHeaders.DOCUMENT_ID, 
documentIdExpression(documentId));
+                if (!isSimpleExpression(documentId)) {
+                    // the sink's missing-id error names the header the route 
reads, as before; the
+                    // parser steps keep the canonical one
+                    sink.put("documentIdHeader", documentId);
+                }
+            } else {
+                // unset, or the component's canonical header (the 3.40 
rename, #9162): read with the
+                // previous fallbacks, its exchange property and the 3.39 name
+                route = route.process(documentIdFallback(name));
+            }
+            // captured into the exchange property before any parse, as by the 
previous builder: the
+            // sink reads the property first, so the id resolved here wins 
over any header
+            route = route.process(exchange -> 
exchange.setProperty(LangChain4jIngest.DOCUMENT_ID_PROPERTY,
+                    
exchange.getMessage().getHeader(LangChain4jIngestHeaders.DOCUMENT_ID)));
+            IdempotentRepository register = repositoryRef == null
+                    ? null
+                    : 
getContext().getRegistry().lookupByNameAndType(repositoryRef, 
IdempotentRepository.class);
+            route = parseSteps(route, name, parser, maxDocumentSize, false, 
register, includeId, excludeId);
+            if (repositoryRef != null) {
+                // deduplication by document id happens inside the sink's 
producer: a duplicate
+                // is answered SKIPPED, a blank delivery releases its claim
+                sink.put("idempotentRepository", "#bean:" + repositoryRef);
+            }
+            route.to(kameletUri(name, SINK_KAMELET, "sink", sink));
+            LOG.infof("Ingestion pipeline '%s': source=%s", name, 
URISupport.sanitizeUri(uri));
         }
-        definition.embeddingStore(store)
-                .embeddingModel(model)
-                .splitter(maxSegmentSize, maxOverlapSize)
-                .embeddingBatchSize(embeddingBatchSize);
-        if (maxDocumentSize > 0) {
-            definition.maxDocumentSize(maxDocumentSize);
+    }
+
+    /**
+     * The optional parse stage: the id guard, the advisory id-pattern and 
duplicate checks and the
+     * raw-size guard, then the parser action Kamelet, which captures the 
document id from the
+     * canonical header into the exchange property before the parse. The route 
is returned
+     * unchanged when the pipeline has no parser.
+     */
+    private static ProcessorDefinition<?> parseSteps(ProcessorDefinition<?> 
route, String name, String parser,
+            int maxDocumentSize, boolean directory, IdempotentRepository 
register, String includeId, String excludeId) {
+        if (parser == null) {

Review Comment:
   Partly true: an oversized media payload is still refused, because 
`MediaIngestService` re-checks `bytes.length`. But that happens only after the 
whole body is read, and a forged `CamelFileLength` passes the up-front check.
   
   `maxDocumentSize` is the component's option, so the fix belongs there rather 
than in a media copy of the guard here: 
[CAMEL-25463](https://issues.apache.org/jira/browse/CAMEL-25463). With it, a 
body not yet in memory will be read only one byte (media) or one character 
(text) past the cap, then refused. `CamelFileLength` will still reject early, 
but it will no longer decide how much is read. That will cover every user of 
the component, not only Camel Quarkus, and keep another copy of the size logic 
out of the extension.
   
   I'll add a test here to verify it, so this PR will wait for the fix. That's 
no extra delay: the fix targets Camel 4.23.0, which this PR already waits for.
   
   Why that will be enough:
   - A directory pipeline is already safe: the file consumer's 
`CamelFileLength` is checked before the file is read.
   - On a consumer-fed pipeline, a guard in the extension could not do more. 
Stream caching (on by default, in memory) reads a stream whole before any of 
our steps run. That needs a limit at the consumer (e.g. 
`quarkus.http.limits.max-body-size`) or spooling, as the component docs in the 
fix will say.
   
   The extension keeps its own guard only for parser pipelines, since the parse 
runs before the sink.
   



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