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]