jiriOndrusek-agent commented on code in PR #9290:
URL: https://github.com/apache/camel-quarkus/pull/9290#discussion_r4218001134
##########
extensions/langchain4j-ingest/runtime/src/main/java/org/apache/camel/quarkus/component/langchain4j/ingest/IngestRoutes.java:
##########
@@ -143,59 +201,344 @@ private IngestPipelineDefinition
builderPipeline(IngestBuilderPipelines.Entry en
IngestPipeline definition = builderPipelines.definition(entry);
String uri = "file".equals(definition.sourceType()) ? null :
definition.sourceUri();
- return definition(name, uri, definition.asRunTimeConfig(),
+ compositionRoute(name, uri, definition.asRunTimeConfig(),
resolveStore(name,
definition.embeddingStoreName().orElse(null)),
resolveModel(name,
definition.embeddingModelName().orElse(null)),
definition.maxSegmentSize(),
definition.maxOverlapSize(),
definition.embeddingBatchSize(),
definition.maxDocumentSize(),
definition.documentSplitterName().orElse(null),
- definition.parser().orElse(null));
+ definition.parser().orElse(null),
+ definition.modality(),
+ definition.contentType().orElse(null));
}
/**
* One creation path for both declaration styles, mirroring how a builder
pipeline provides
* the same runtime-configuration view the configuration path reads.
*/
- private IngestPipelineDefinition definition(String name, String uri,
+ private void compositionRoute(String name, String uri,
IngestRunTimeConfig.PipelineRunTimeConfig runtime,
EmbeddingStore<TextSegment> store, EmbeddingModel model,
int maxSegmentSize, int maxOverlapSize, int embeddingBatchSize,
int maxDocumentSize,
- String documentSplitterName, String parser) {
+ String documentSplitterName, String parser, String modality,
String contentType) {
+
+ // the name is substituted into Kamelet URIs and registry references,
so both declaration
+ // styles are held to one charset; @Ingest names are already checked
at build time
+ if (!name.matches("[A-Za-z0-9._-]+")) {
+ throw new IllegalStateException(
+ "Ingestion pipeline name '" + name + "' may only contain
letters, digits, '.', '_' and '-'");
+ }
+ // the component validates modality and content type; that a media
document is never
+ // parsed is a rule of this composition, where the parser action sits
before the sink.
+ // A configured pipeline fails the build on both already
+ boolean media = "media".equalsIgnoreCase(modality);
+ if (media && parser != null) {
+ throw new IllegalStateException("Ingestion pipeline '" + name
+ + "' sets modality 'media' together with a parser. A media
document is embedded whole and"
+ + " never parsed; remove one of them.");
+ }
+ String storeRef = bindInstance(name, "store", store);
+ String modelRef = bindInstance(name, "model", model);
+ String documentId = runtime == null ? null :
runtime.source().documentId().orElse(null);
+ String repositoryRef = repositoryRef(name, runtime, uri == null);
+
+ // options at their defaults are passed all the same, here and in the
other steps: an
+ // application-wide camel.kamelet.<kamelet>.* property would otherwise
fill them in too.
+ // An unset string option goes as "", which the component reads as
unset. An unset bean
+ // option (documentFilter, documentSplitter, idempotentRepository) is
left out on purpose:
+ // "" converts to no bean and fails the start, so such a property
still reaches it, as
+ // usage.adoc warns
+ Map<String, Object> sink = new LinkedHashMap<>();
+ sink.put("pipelineName", name);
+ sink.put("documentIdHeader", LangChain4jIngestHeaders.DOCUMENT_ID);
+ sink.put("maxSegmentSize", String.valueOf(maxSegmentSize));
+ sink.put("maxOverlapSize", String.valueOf(maxOverlapSize));
+ sink.put("embeddingBatchSize", String.valueOf(embeddingBatchSize));
+ sink.put("modality", modality);
+ sink.put("contentType", contentType == null ? "" : contentType);
+ // 0, no limit, is the endpoint's default too
+ sink.put("maxDocumentSize", String.valueOf(maxDocumentSize));
+ sink.put("minDocumentSize", "0");
+ if (documentSplitterName != null) {
+ sink.put("documentSplitter", "#bean:" + documentSplitterName);
+ }
+ // filter.* comes from configuration for both declaration styles, the
builder having no
+ // filter API, and the component enforces it
+ IngestRunTimeConfig.PipelineRunTimeConfig configured =
runTimeConfig.pipelines().get(name);
+ String includeId = configured == null ? "" :
configured.filter().includeId().orElse("");
+ String excludeId = configured == null ? "" :
configured.filter().excludeId().orElse("");
+ sink.put("includeId", includeId);
+ sink.put("excludeId", excludeId);
+ if (configured != null) {
+ if (configured.filter().minDocumentSize() != 0) {
+ sink.put("minDocumentSize",
String.valueOf(configured.filter().minDocumentSize()));
+ }
+ configured.filter().documentFilter().ifPresent(bean ->
sink.put("documentFilter", "#bean:" + bean));
+ }
+ sink.put("embeddingStore", "#bean:" + storeRef);
+ sink.put("embeddingModel", "#bean:" + modelRef);
- IngestPipelineDefinition definition;
if (uri == null) {
String directory = required(name, runtime == null ? null :
runtime.source().directory().orElse(null),
"source.directory");
- definition = IngestPipelineDefinition.directory(name, directory)
- .recursive(runtime.source().recursive());
+ // substituted into the file-source Kamelet's endpoint URI: a '?'
or '#' could
+ // inject consumer options - delete=true would consume the user's
documents
+ if (directory.contains("?") || directory.contains("#")) {
+ throw new IllegalStateException("Ingestion pipeline '" + name
+ + "': the directory must not contain '?' or '#' (got
'" + directory + "')");
+ }
+ Map<String, Object> source = new LinkedHashMap<>();
+ source.put("directory", directory);
+ source.put("recursive",
String.valueOf(runtime.source().recursive()));
+ // the file consumer's own default poll delay
+ source.put("delay", "500");
+ if (parser == null && !media) {
+ // text is read as UTF-8; a parser or a media model receives
the raw bytes instead -
+ // the format is its business, and a charset conversion would
corrupt a binary document
+ source.put("charset", "UTF-8");
+ }
+ source.put("idempotentRepository", "#bean:" + repositoryRef);
+
+ ProcessorDefinition<?> route = from(kameletUri(name,
FILE_SOURCE_KAMELET, "source", source))
+ .routeId("langchain4j-ingest-" + name);
+ if (documentId != null) {
+ // override the source's file-name default; captured before
any further step
+ route = route.setHeader(LangChain4jIngestHeaders.DOCUMENT_ID,
documentIdExpression(documentId));
+ }
+ // no duplicate pre-check: the source register already filtered
duplicates out
+ route = parseSteps(route, name, parser, maxDocumentSize, true,
null, includeId, excludeId);
+ // the file consumer discards the reply and the source register
already keeps the
+ // same file version from being ingested twice, so no repository
goes to the sink.
+ // The discarded reply would hide an empty outcome, so it is
logged: warned with a
+ // parser, since a parse to nothing typically means a missing Tika
parser module or an
+ // image-only document, and the file's register key is committed.
A filtered one is
+ // logged at DEBUG
+ route.to(kameletUri(name, SINK_KAMELET, "sink",
sink)).process(exchange -> {
+ IngestResult result =
exchange.getMessage().getBody(IngestResult.class);
+ if (result != null && result.outcome() ==
IngestResult.Outcome.FILTERED) {
+ logFiltered(name, result.documentId());
+ }
+ if (result == null || result.outcome() !=
IngestResult.Outcome.EMPTY) {
+ return;
+ }
+ if (parser != null) {
+ LOG.warnf("Ingestion pipeline '%s': document '%s' parsed
to no text and was skipped; its key is"
+ + " committed, so it is not retried until the file
changes (missing parser module?"
+ + " image-only document?)", name,
result.documentId());
+ } else {
+ LOG.debugf("Ingestion pipeline '%s': document '%s'
contained no text, nothing was written",
+ name, result.documentId());
+ }
+ });
+ LOG.infof("Ingestion pipeline '%s': source=file:%s", name,
directory);
} else {
- definition = IngestPipelineDefinition.consumer(name, uri);
+ if (!uri.contains(":")) {
+ throw new IllegalStateException(
+ "Ingestion pipeline '" + name + "': '" + uri + "' is
not a consumer URI");
+ }
+ ProcessorDefinition<?> route =
from(uri).routeId("langchain4j-ingest-" + name);
+ if (documentId != null &&
!IngestHeaders.DOCUMENT_ID.equals(documentId)) {
+ // copied into the canonical header the actions read, as on a
directory pipeline:
+ // evaluated against the exchange the consumer delivered,
before any parse, and a
+ // plain header name is read directly, never substituted into
an expression
+ route = route.setHeader(LangChain4jIngestHeaders.DOCUMENT_ID,
documentIdExpression(documentId));
+ if (!isSimpleExpression(documentId)) {
+ // the sink's missing-id error names the header the route
reads, as before; the
+ // parser steps keep the canonical one
+ sink.put("documentIdHeader", documentId);
+ }
+ } else {
+ // unset, or the component's canonical header (the 3.40
rename, #9162): read with the
+ // previous fallbacks, its exchange property and the 3.39 name
+ route = route.process(documentIdFallback(name));
+ }
+ // captured into the exchange property before any parse, as by the
previous builder: the
+ // sink reads the property first, so the id resolved here wins
over any header
+ route = route.process(exchange ->
exchange.setProperty(LangChain4jIngest.DOCUMENT_ID_PROPERTY,
+
exchange.getMessage().getHeader(LangChain4jIngestHeaders.DOCUMENT_ID)));
+ IdempotentRepository register = repositoryRef == null
+ ? null
+ :
getContext().getRegistry().lookupByNameAndType(repositoryRef,
IdempotentRepository.class);
+ route = parseSteps(route, name, parser, maxDocumentSize, false,
register, includeId, excludeId);
+ if (repositoryRef != null) {
+ // deduplication by document id happens inside the sink's
producer: a duplicate
+ // is answered SKIPPED, a blank delivery releases its claim
+ sink.put("idempotentRepository", "#bean:" + repositoryRef);
+ }
+ route.to(kameletUri(name, SINK_KAMELET, "sink", sink));
+ LOG.infof("Ingestion pipeline '%s': source=%s", name,
URISupport.sanitizeUri(uri));
}
- definition.embeddingStore(store)
- .embeddingModel(model)
- .splitter(maxSegmentSize, maxOverlapSize)
- .embeddingBatchSize(embeddingBatchSize);
- if (maxDocumentSize > 0) {
- definition.maxDocumentSize(maxDocumentSize);
+ }
+
+ /**
+ * The optional parse stage: the id guard, the advisory id-pattern and
duplicate checks and the
+ * raw-size guard, then the parser action Kamelet, which captures the
document id from the
+ * canonical header into the exchange property before the parse. The route
is returned
+ * unchanged when the pipeline has no parser.
+ */
+ private static ProcessorDefinition<?> parseSteps(ProcessorDefinition<?>
route, String name, String parser,
+ int maxDocumentSize, boolean directory, IdempotentRepository
register, String includeId, String excludeId) {
+ if (parser == null) {
Review Comment:
Partly true: an oversized media payload is still refused, because
`MediaIngestService` re-checks `bytes.length`. But that happens only after the
whole body is read, and a forged `CamelFileLength` passes the up-front check.
`maxDocumentSize` is the component's option, so the fix belongs there rather
than in a media copy of the guard here:
[CAMEL-25463](https://issues.apache.org/jira/browse/CAMEL-25463). With it, a
body not yet in memory will be read only one byte (media) or one character
(text) past the cap, then refused. `CamelFileLength` will still reject early,
but it will no longer decide how much is read. That will cover every user of
the component, not only Camel Quarkus, and keep another copy of the size logic
out of the extension.
I'll add a test here to verify it, so this PR will wait for the fix. That's
no extra delay: the fix targets Camel 4.23.0, which this PR already waits for.
Why that will be enough:
- A directory pipeline is already safe: the file consumer's
`CamelFileLength` is checked before the file is read.
- On a consumer-fed pipeline, a guard in the extension could not do more.
Stream caching (on by default, in memory) reads a stream whole before any of
our steps run. That needs a limit at the consumer (e.g.
`quarkus.http.limits.max-body-size`) or spooling, as the component docs in the
fix will say.
The extension keeps its own guard only for parser pipelines, since the parse
runs before the sink.
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]