jiriOndrusek-agent commented on code in PR #9290:
URL: https://github.com/apache/camel-quarkus/pull/9290#discussion_r4218982387
##########
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:
True. I added a comment pointing to
`LangChain4jIngestProducer.parsePatterns`/`idAccepted` ("keep them in step").
The component stays authoritative; the copy only spares the parse. Exposing
them from the component would remove the copy, and I can raise that in Camel.
##########
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:
Agreed. The builder view now delegates to the configured `filter.*`, or
returns an empty filter when there is none. `IngestRoutes` reads
`runtime.filter()` for both declaration styles.
##########
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:
Partly. The misleading error is real, but the start fails, so the real
pipeline doesn't run without the filter.
The error now names the Java-declared pipelines: `Ingestion pipeline
'datasheet' has no source.directory, and no @Ingest method declares it
(declared in Java: [datasheets]). Fix the name, or set ...source.directory`.
Test: `IngestBuilderFilterTypoTest`.
##########
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:
Done. An unset `filter.document-filter` now goes as an always-true
predicate. The IT poison set gains
`camel.kamelet.langchain4j-ingest-sink.documentFilter=#bean:rejectAll`.
`documentSplitter` and `idempotentRepository` stay unpinned, since no neutral
bean exists for them.
##########
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:
Fixed for both paths: a value with nothing before `;` is rejected by
`IngestPipeline.contentType()` and, for `content-type`, at build time. Tests
extended.
##########
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:
Done: it is lower-cased before it goes to the Kamelet.
--
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]