jamesnetherton commented on code in PR #9290:
URL: https://github.com/apache/camel-quarkus/pull/9290#discussion_r4208218032


##########
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:
   With `modality=media` there is no parser, so `parseSteps` returns here and 
the bounded raw-size read (max + 1 bytes, not trusting `CamelFileLength`) never 
runs for media pipelines, including consumer-fed ones (`source.uri` / 
`Source.endpoint`).
   
   The component's own check only trusts the `CamelFileLength` header before 
calling `getMandatoryBody(byte[])`. So a large body sent with a forged 
`CamelFileLength` gets past `max-document-size` and is loaded into the heap in 
full. Should the size guard also apply when `media` is true?



##########
extensions/langchain4j-ingest/runtime/src/main/doc/usage.adoc:
##########
@@ -115,6 +136,26 @@ A missing parser extension fails the build naming the 
artifact, the same way a m
 The Java twin is `IngestPipeline.parser("tika")`, and a consumer-fed pipeline 
parses the same way — the payload the consumer delivers goes to the parser 
first.
 With a parser, `max-document-size` also rejects a raw payload larger than the 
cap (in bytes) before it reaches the parser, on top of capping the extracted 
text (in characters).
 
+=== Ingesting media: images, audio, video and PDF
+
+`modality=media` embeds each document whole, as one vector, instead of reading 
it as text: the bytes go to a multimodal embedding model as audio, an image, 
video or a PDF, told apart by the MIME type, and the vector is stored with a 
placeholder segment whose text is the document id, stamped with the usual 
identity metadata, so retrieval cites it like any other.
+The model must declare the matching content type in its 
`supportedContentTypes()`: the pipeline fails to start with a text-only model, 
and a document whose medium the model does not declare fails before its bytes 
are read.
+`parser` and `document-splitter` must not be set, the splitter sizes and 
`embedding-batch-size` do not apply, and `max-document-size` and 
`filter.min-document-size` count bytes.
+
+[source,properties]
+----
+quarkus.camel.langchain4j.ingest.photos.source.directory=/var/data/photos
+quarkus.camel.langchain4j.ingest.photos.modality=media
+quarkus.camel.langchain4j.ingest.photos.embedding-model=my-multimodal-model
+quarkus.camel.langchain4j.ingest.photos.embedding-store=photos
+quarkus.camel.langchain4j.ingest.photos.filter.include-id=**/*.png,**/*.jpg
+----
+
+The MIME type handed to the model comes from `content-type` or, when unset, 
from the file extension through Camel's MIME table (wav, mp3, flac, ogg, opus, 
m4a, aac and aiff for audio; png, jpg, gif and webp for images; mp4, mov and 
webm for video; pdf); a file whose type cannot be determined fails the 
exchange, and `content-type` without `modality=media` is rejected like the 
parser and splitter conflicts.
+On a directory pipeline such a file — `notes.txt`, `Thumbs.db` — fails again 
on every poll, so restrict the ids with `filter.include-id`, as above: it is 
checked before the media type.

Review Comment:
   On a directory pipeline, a file whose media type can't be determined fails 
on every poll, for ever, unless `filter.include-id` is set. Since retrying can 
never succeed, could it be committed in the register as filtered/rejected (like 
an empty parse is) instead of only being documented?



##########
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("");

Review Comment:
   `filter.*` adds the pipeline name to the runtime map. A filter property for 
a misspelled `@Ingest` name (e.g. `datasheet.filter.exclude-id` when the 
pipeline is `datasheets`) creates a phantom configured pipeline. Startup then 
fails with a missing `source.directory` error rather than "unknown pipeline", 
and the real pipeline runs without the filter. Could a name that only has 
`filter.*` set be detected and reported?



##########
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) {
+            return route;
         }
-        if (documentSplitterName != null) {
-            definition.documentSplitter(documentSplitterName);
+        String documentIdHeader = LangChain4jIngestHeaders.DOCUMENT_ID;
+        // a parser must never run without a captured id: the header the 
action reads would be
+        // absent, and headers written by a parsed document could take its 
place - the previous
+        // engine failed such a delivery, and so does this guard
+        route = route.process(exchange -> {
+            String id = exchange.getMessage().getHeader(documentIdHeader, 
String.class);
+            if (id == null || id.isBlank()) {
+                throw new IllegalArgumentException("Ingestion pipeline '" + 
name
+                        + "': the document id resolved to nothing, and a 
parser pipeline requires it before"
+                        + " the parse - a parsed document must not supply its 
own identity");
+            }
+        });
+        String[] include = patterns(includeId);
+        String[] exclude = patterns(excludeId);
+        if (include != null || exclude != null) {
+            // advisory too: the sink's id patterns, applied before the parse 
so an excluded
+            // document is never parsed, nor fails the parse or the size guard 
on every poll
+            route = route.choice()
+                    .when(exchange -> {
+                        String id = 
exchange.getMessage().getHeader(documentIdHeader, String.class);
+                        return AntPathMatcher.INSTANCE.anyMatch(exclude, id)
+                                || (include != null && 
!AntPathMatcher.INSTANCE.anyMatch(include, id));
+                    })
+                    .process(exchange -> {
+                        String id = 
exchange.getMessage().getHeader(documentIdHeader, String.class);
+                        if (directory) {
+                            logFiltered(name, id);
+                        }
+                        exchange.getMessage().setBody(new IngestResult(name, 
id, 0, IngestResult.Outcome.FILTERED));
+                    })
+                    .stop()
+                    .end();
+        }
+        if (register != null) {
+            // advisory: a known duplicate is answered SKIPPED before the 
(possibly remote) parse
+            // is paid for; the sink producer's eager claim stays 
authoritative, so a duplicate
+            // racing this check is still caught there
+            route = route.choice()
+                    .when(exchange -> 
register.contains(exchange.getMessage().getHeader(documentIdHeader, 
String.class)))
+                    .process(exchange -> exchange.getMessage().setBody(new 
IngestResult(name,
+                            exchange.getMessage().getHeader(documentIdHeader, 
String.class), 0,
+                            IngestResult.Outcome.SKIPPED)))
+                    .stop()
+                    .end();
         }
-        if (parser != null) {
-            definition.parser(parser);
+        if (maxDocumentSize > 0) {
+            // the endpoint's own cap counts extracted characters, which 
protects the splitter and
+            // the model but not the parse: this guard rejects the raw payload 
first, before tika
+            // or docling materialize it
+            route = route.process(rawSizeGuard(name, maxDocumentSize, 
directory, documentIdHeader));
         }
+        Map<String, Object> action = new LinkedHashMap<>();
+        action.put("documentIdHeader", documentIdHeader);
+        // the guard above caps the raw payload; the action's own cap would 
read each file whole
+        // to measure it again, so it is pinned off
+        action.put("maxDocumentSize", "0");
+        return route.to(kameletUri(name, 
Parser.valueOf(parser.toUpperCase(Locale.ROOT)).actionKamelet, "parser", 
action));
+    }
 
-        String documentId = runtime == null ? null : 
runtime.source().documentId().orElse(null);
-        if (documentId != null) {
-            definition.documentId(documentId);
-        } else if (uri != null) {
-            // the component's own default since the 3.40 rename (#9162); the 
3.39 name is still
-            // read as a fallback in the route builder's documentIdExpression
-            definition.documentId(IngestHeaders.DOCUMENT_ID);
+    private static Processor rawSizeGuard(String name, int maxDocumentSize, 
boolean trustDeclaredLength,
+            String documentIdHeader) {
+        // the declared-length header spares even the read, but only the 
directory pipeline's own
+        // file consumer is trusted to have set it: on a consumer pipeline 
every header may be
+        // attacker-supplied along with the payload, so its body is always 
measured - a forged
+        // CamelFileLength must not talk an oversized payload past the guard 
and into the parser
+        return exchange -> {
+            Long declared = trustDeclaredLength
+                    ? exchange.getMessage().getHeader(Exchange.FILE_LENGTH, 
Long.class)
+                    : null;
+            long size;
+            byte[] bounded = null;
+            if (declared != null) {
+                size = declared;
+            } else {
+                // every body is measured through a stream - a GenericFile 
streams from disk, a
+                // byte[] merely wraps - so an attacker-sized payload is 
rejected after
+                // maxDocumentSize + 1 bytes instead of being materialized 
whole in the heap just
+                // to be measured; an accepted stream is consumed here, so the 
bytes replace it
+                // as the body. A null body carries no bytes to guard; it 
flows on and becomes
+                // the EMPTY outcome
+                InputStream stream = 
exchange.getMessage().getBody(InputStream.class);
+                if (stream == null) {
+                    size = 0;
+                } else {
+                    int limit = maxDocumentSize == Integer.MAX_VALUE ? 
Integer.MAX_VALUE : maxDocumentSize + 1;
+                    bounded = stream.readNBytes(limit);
+                    size = bounded.length;
+                }
+            }
+            if (size > maxDocumentSize) {
+                throw new IllegalArgumentException(
+                        "Ingestion pipeline '" + name + "': document '"
+                                + 
exchange.getMessage().getHeader(documentIdHeader, String.class)
+                                + "' exceeds maxDocumentSize (" + size + " > " 
+ maxDocumentSize + " bytes)");
+            }
+            if (bounded != null) {
+                exchange.getMessage().setBody(bounded);
+            }
+        };
+    }
+
+    /** Comma-separated id patterns, parsed as the component parses them; null 
when none remain. */
+    private static String[] patterns(String value) {

Review Comment:
   `patterns()` and the `AntPathMatcher` check before parsing re-implement the 
component's private `parsePatterns`/`idAccepted`. If the component's matching 
changes (case, a leading `/`, splitting, where the id comes from), the advisory 
filter and the authoritative one could disagree. Worth a comment linking the 
two, or a test that checks they agree?



##########
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);

Review Comment:
   The modality is forwarded exactly as typed (`"Media"`, `"MEDIA"`, both used 
in the tests), so this depends on the component's `IngestModality` enum 
conversion ignoring case. `modality.toLowerCase(Locale.ROOT)` here would remove 
that dependency.



##########
extensions/langchain4j-ingest/runtime/src/main/java/org/apache/camel/quarkus/component/langchain4j/ingest/IngestPipeline.java:
##########
@@ -177,6 +218,12 @@ public boolean enabled() {
             public SourceRunTimeConfig source() {
                 return sourceConfig;
             }
+
+            @Override
+            public FilterRunTimeConfig filter() {

Review Comment:
   `filter()` returns null, so this `PipelineRunTimeConfig` is only partly 
implemented. Any later code that reads `runtime.filter()` will hit an NPE, and 
only for `@Ingest` pipelines. Could it return an empty `FilterRunTimeConfig`, 
or delegate to the configured pipeline's filter?



##########
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));

Review Comment:
   When `filter.document-filter` is unset it isn't passed, so an 
application-wide `camel.kamelet.langchain4j-ingest-sink.documentFilter` still 
applies to every pipeline. The docs do say not to set those properties, but 
binding a neutral always-true `Predicate` as the default would pin the option, 
as the PR intends.



##########
extensions/langchain4j-ingest/runtime/src/main/java/org/apache/camel/quarkus/component/langchain4j/ingest/IngestPipeline.java:
##########
@@ -123,6 +125,33 @@ public IngestPipeline documentSplitter(String beanName) {
         return this;
     }
 
+    /**
+     * What the consumed payload is: {@code text}, the default, is split and 
embedded segment by
+     * segment; {@code media} (audio, an image, video or a PDF, told apart by 
the MIME type) is
+     * embedded whole, as one vector, by a model that declares the matching 
content type. The
+     * twin of the {@code modality} configuration property.
+     */
+    public IngestPipeline modality(String modality) {
+        // the same rule the configuration path is held to at build time
+        if (!"text".equalsIgnoreCase(modality) && 
!"media".equalsIgnoreCase(modality)) {
+            throw new IllegalArgumentException("modality must be 'text' or 
'media' (got '" + modality + "')");
+        }
+        this.modality = modality;
+        return this;
+    }
+
+    /**
+     * MIME type of a media payload, such as {@code audio/wav}; unset, it is 
derived from the
+     * document id's file extension. The twin of the {@code content-type} 
configuration property.
+     */
+    public IngestPipeline contentType(String contentType) {
+        if (contentType == null || contentType.isBlank()) {

Review Comment:
   This only rejects blank values. The component drops everything after `;` and 
treats an empty result as unset, so `contentType("; codecs=opus")` passes here 
and then silently falls back to the file extension. Could this apply the same 
normalization before validating?



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