This is an automated email from the ASF dual-hosted git repository.
oscerd pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/camel-kamelets.git
The following commit(s) were added to refs/heads/main by this push:
new c0f203e55 CAMEL-24753: Add filter support to the langchain4j-ingest
Kamelets (#3033)
c0f203e55 is described below
commit c0f203e5550295b65b55ef7ed3232716eb995a9a
Author: Jiří Ondrušek <[email protected]>
AuthorDate: Wed Sep 23 09:49:32 2026 +0200
CAMEL-24753: Add filter support to the langchain4j-ingest Kamelets (#3033)
File source: include/exclude Ant globs forwarded to the file consumer's
antInclude/antExclude, so a filtered file is never read. Sink: passthrough
of the component's filter options (includeId, excludeId, minDocumentSize,
documentFilter); rejected deliveries answer outcome=filtered and never
keep a dedup claim. Citrus tests cover both: the file test asserts the
aggregated delivery list is exactly [keep.md], the sink test drives a
non-matching id and a too-short document through the filters.
Co-authored-by: Claude Fable 5 <[email protected]>
---
.../langchain4j-ingest-file-source.kamelet.yaml | 13 +++
kamelets/langchain4j-ingest-sink.kamelet.yaml | 30 ++++++-
.../langchain4j-ingest-file-source.kamelet.yaml | 13 +++
.../kamelets/langchain4j-ingest-sink.kamelet.yaml | 30 ++++++-
...chain4j-ingest-file-filter-route.citrus.it.yaml | 67 +++++++++++++++
.../langchain4j-ingest-file-filter-route.yaml | 56 +++++++++++++
...chain4j-ingest-sink-filter-route.citrus.it.yaml | 95 ++++++++++++++++++++++
.../langchain4j-ingest-sink-filter-route.yaml | 77 ++++++++++++++++++
8 files changed, 379 insertions(+), 2 deletions(-)
diff --git a/kamelets/langchain4j-ingest-file-source.kamelet.yaml
b/kamelets/langchain4j-ingest-file-source.kamelet.yaml
index 5bcfc8908..3709e6a32 100644
--- a/kamelets/langchain4j-ingest-file-source.kamelet.yaml
+++ b/kamelets/langchain4j-ingest-file-source.kamelet.yaml
@@ -56,6 +56,17 @@ spec:
description: Whether subdirectories are ingested too.
type: boolean
default: true
+ include:
+ title: Include
+ description: Comma-separated Ant-style patterns for files to ingest,
matched on
+ the path relative to the directory, for example `**/*.pdf,**/*.md`.
When
+ unset, every file is ingested.
+ type: string
+ exclude:
+ title: Exclude
+ description: Comma-separated Ant-style patterns for files to skip, for
example
+ `**/draft-*`. Exclusion wins over inclusion.
+ type: string
delay:
title: Delay
description: Milliseconds between directory polls; the file endpoint's
default is
@@ -90,6 +101,8 @@ spec:
# an edited file gets a new key and re-ingests
idempotentKey: "${file:absolute.path}:${file:modified}:${file:size}"
recursive: "{{recursive}}"
+ antInclude: "{{?include}}"
+ antExclude: "{{?exclude}}"
readLock: "changed"
delay: "{{?delay}}"
charset: "{{?charset}}"
diff --git a/kamelets/langchain4j-ingest-sink.kamelet.yaml
b/kamelets/langchain4j-ingest-sink.kamelet.yaml
index c490cf65c..f0f1bc1e6 100644
--- a/kamelets/langchain4j-ingest-sink.kamelet.yaml
+++ b/kamelets/langchain4j-ingest-sink.kamelet.yaml
@@ -43,7 +43,9 @@ spec:
`embeddingStore` or `embeddingModel` is not set, the single registry
bean of that type is
used. With `idempotentRepository` set, the first write per document id
wins - a
re-delivered edited document is skipped, so streams that carry updates
need a
- version-aware id.
+ version-aware id. The filter options (`includeId`, `excludeId`,
`minDocumentSize`,
+ `documentFilter`) answer rejected deliveries with a filtered outcome
that never keeps
+ a dedup claim.
type: object
properties:
pipelineName:
@@ -80,6 +82,28 @@ spec:
pipeline holds a document in memory whole, so set the cap when the
source can
deliver oversized payloads. An oversized document fails the exchange
cleanly.
type: integer
+ includeId:
+ title: Include Id
+ description: Comma-separated Ant-style patterns the document id must
match to be
+ ingested, for example `docs/**,*.md`. A non-matching delivery is
answered filtered,
+ before the dedup claim and without reading the body.
+ type: string
+ excludeId:
+ title: Exclude Id
+ description: Comma-separated Ant-style patterns for document ids to
skip, for example
+ `**/draft-*`. Exclusion wins over includeId.
+ type: string
+ minDocumentSize:
+ title: Min Document Size
+ description: Minimum size of one document in characters; unset means
no minimum. A
+ shorter document is answered filtered and releases its dedup claim.
+ type: integer
+ documentFilter:
+ title: Document Filter
+ description: A Predicate bean deciding whether a delivery is ingested,
as a
+ `#bean:name` reference; a rejected delivery is answered filtered and
releases its
+ dedup claim.
+ type: string
documentSplitter:
title: Document Splitter
description: A DocumentSplitter bean replacing the default recursive
splitting, as a
@@ -117,6 +141,10 @@ spec:
maxOverlapSize: "{{maxOverlapSize}}"
embeddingBatchSize: "{{embeddingBatchSize}}"
maxDocumentSize: "{{?maxDocumentSize}}"
+ includeId: "{{?includeId}}"
+ excludeId: "{{?excludeId}}"
+ minDocumentSize: "{{?minDocumentSize}}"
+ documentFilter: "{{?documentFilter}}"
documentSplitter: "{{?documentSplitter}}"
embeddingStore: "{{?embeddingStore}}"
embeddingModel: "{{?embeddingModel}}"
diff --git
a/library/camel-kamelets/src/main/resources/kamelets/langchain4j-ingest-file-source.kamelet.yaml
b/library/camel-kamelets/src/main/resources/kamelets/langchain4j-ingest-file-source.kamelet.yaml
index 5bcfc8908..3709e6a32 100644
---
a/library/camel-kamelets/src/main/resources/kamelets/langchain4j-ingest-file-source.kamelet.yaml
+++
b/library/camel-kamelets/src/main/resources/kamelets/langchain4j-ingest-file-source.kamelet.yaml
@@ -56,6 +56,17 @@ spec:
description: Whether subdirectories are ingested too.
type: boolean
default: true
+ include:
+ title: Include
+ description: Comma-separated Ant-style patterns for files to ingest,
matched on
+ the path relative to the directory, for example `**/*.pdf,**/*.md`.
When
+ unset, every file is ingested.
+ type: string
+ exclude:
+ title: Exclude
+ description: Comma-separated Ant-style patterns for files to skip, for
example
+ `**/draft-*`. Exclusion wins over inclusion.
+ type: string
delay:
title: Delay
description: Milliseconds between directory polls; the file endpoint's
default is
@@ -90,6 +101,8 @@ spec:
# an edited file gets a new key and re-ingests
idempotentKey: "${file:absolute.path}:${file:modified}:${file:size}"
recursive: "{{recursive}}"
+ antInclude: "{{?include}}"
+ antExclude: "{{?exclude}}"
readLock: "changed"
delay: "{{?delay}}"
charset: "{{?charset}}"
diff --git
a/library/camel-kamelets/src/main/resources/kamelets/langchain4j-ingest-sink.kamelet.yaml
b/library/camel-kamelets/src/main/resources/kamelets/langchain4j-ingest-sink.kamelet.yaml
index c490cf65c..f0f1bc1e6 100644
---
a/library/camel-kamelets/src/main/resources/kamelets/langchain4j-ingest-sink.kamelet.yaml
+++
b/library/camel-kamelets/src/main/resources/kamelets/langchain4j-ingest-sink.kamelet.yaml
@@ -43,7 +43,9 @@ spec:
`embeddingStore` or `embeddingModel` is not set, the single registry
bean of that type is
used. With `idempotentRepository` set, the first write per document id
wins - a
re-delivered edited document is skipped, so streams that carry updates
need a
- version-aware id.
+ version-aware id. The filter options (`includeId`, `excludeId`,
`minDocumentSize`,
+ `documentFilter`) answer rejected deliveries with a filtered outcome
that never keeps
+ a dedup claim.
type: object
properties:
pipelineName:
@@ -80,6 +82,28 @@ spec:
pipeline holds a document in memory whole, so set the cap when the
source can
deliver oversized payloads. An oversized document fails the exchange
cleanly.
type: integer
+ includeId:
+ title: Include Id
+ description: Comma-separated Ant-style patterns the document id must
match to be
+ ingested, for example `docs/**,*.md`. A non-matching delivery is
answered filtered,
+ before the dedup claim and without reading the body.
+ type: string
+ excludeId:
+ title: Exclude Id
+ description: Comma-separated Ant-style patterns for document ids to
skip, for example
+ `**/draft-*`. Exclusion wins over includeId.
+ type: string
+ minDocumentSize:
+ title: Min Document Size
+ description: Minimum size of one document in characters; unset means
no minimum. A
+ shorter document is answered filtered and releases its dedup claim.
+ type: integer
+ documentFilter:
+ title: Document Filter
+ description: A Predicate bean deciding whether a delivery is ingested,
as a
+ `#bean:name` reference; a rejected delivery is answered filtered and
releases its
+ dedup claim.
+ type: string
documentSplitter:
title: Document Splitter
description: A DocumentSplitter bean replacing the default recursive
splitting, as a
@@ -117,6 +141,10 @@ spec:
maxOverlapSize: "{{maxOverlapSize}}"
embeddingBatchSize: "{{embeddingBatchSize}}"
maxDocumentSize: "{{?maxDocumentSize}}"
+ includeId: "{{?includeId}}"
+ excludeId: "{{?excludeId}}"
+ minDocumentSize: "{{?minDocumentSize}}"
+ documentFilter: "{{?documentFilter}}"
documentSplitter: "{{?documentSplitter}}"
embeddingStore: "{{?embeddingStore}}"
embeddingModel: "{{?embeddingModel}}"
diff --git
a/tests/camel-kamelets-itest/src/test/resources/langchain4j-ingest/langchain4j-ingest-file-filter-route.citrus.it.yaml
b/tests/camel-kamelets-itest/src/test/resources/langchain4j-ingest/langchain4j-ingest-file-filter-route.citrus.it.yaml
new file mode 100644
index 000000000..34101cf79
--- /dev/null
+++
b/tests/camel-kamelets-itest/src/test/resources/langchain4j-ingest/langchain4j-ingest-file-filter-route.citrus.it.yaml
@@ -0,0 +1,67 @@
+# ---------------------------------------------------------------------------
+# Licensed to the Apache Software Foundation (ASF) under one or more
+# contributor license agreements. See the NOTICE file distributed with
+# this work for additional information regarding copyright ownership.
+# The ASF licenses this file to You under the Apache License, Version 2.0
+# (the "License"); you may not use this file except in compliance with
+# the License. You may obtain a copy of the License at
+#
+# http://www.apache.org/licenses/LICENSE-2.0
+#
+# Unless required by applicable law or agreed to in writing, software
+# distributed under the License is distributed on an "AS IS" BASIS,
+# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+# See the License for the specific language governing permissions and
+# limitations under the License.
+# ---------------------------------------------------------------------------
+
+name: langchain4j-ingest-file-filter-route-test
+actions:
+ # drain any request a previous test may have left on the shared server
+ - purge:
+ endpoints:
+ - name: "httpServer"
+ timeout: 100
+ - createVariables:
+ variables:
+ - name: "report.url"
+ value: "http://localhost:${http.server.port}/result/file-filter"
+ # a fresh directory per test run: a still-terminating integration from
a
+ # previous local run must not re-consume this run's seed files
+ - name: "ingest.dir"
+ value: "ingest-filter-citrus:randomNumber(5)"
+
+ # Run a route seeding three files; the source's include/exclude globs let
+ # only keep.md through, and the aggregated id list is reported back
+ - camel:
+ cli:
+ run:
+ waitForRunningState: false
+ args:
+ - "--max-seconds=60"
+ integration:
+ file:
"langchain4j-ingest/langchain4j-ingest-file-filter-route.yaml"
+ systemProperties:
+ properties:
+ - name: "report.url"
+ value: "${report.url}"
+ - name: "ingest.dir"
+ value: "${ingest.dir}"
+
+ # Exactly one delivery: a leaked file changes the count, an over-filter
times out
+ - http:
+ server: "httpServer"
+ receiveRequest:
+ POST:
+ path: "/result/file-filter"
+ type: "plaintext"
+ body:
+ data: "1:keep.md"
+
+ - http:
+ server: "httpServer"
+ sendResponse:
+ response:
+ status: 200
+ reasonPhrase: "OK"
+ version: "HTTP/1.1"
diff --git
a/tests/camel-kamelets-itest/src/test/resources/langchain4j-ingest/langchain4j-ingest-file-filter-route.yaml
b/tests/camel-kamelets-itest/src/test/resources/langchain4j-ingest/langchain4j-ingest-file-filter-route.yaml
new file mode 100644
index 000000000..53070e81a
--- /dev/null
+++
b/tests/camel-kamelets-itest/src/test/resources/langchain4j-ingest/langchain4j-ingest-file-filter-route.yaml
@@ -0,0 +1,56 @@
+# ---------------------------------------------------------------------------
+# Licensed to the Apache Software Foundation (ASF) under one or more
+# contributor license agreements. See the NOTICE file distributed with
+# this work for additional information regarding copyright ownership.
+# The ASF licenses this file to You under the Apache License, Version 2.0
+# (the "License"); you may not use this file except in compliance with
+# the License. You may obtain a copy of the License at
+#
+# http://www.apache.org/licenses/LICENSE-2.0
+#
+# Unless required by applicable law or agreed to in writing, software
+# distributed under the License is distributed on an "AS IS" BASIS,
+# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+# See the License for the specific language governing permissions and
+# limitations under the License.
+# ---------------------------------------------------------------------------
+
+# The include/exclude filters: three files are seeded, only keep.md passes
+# (skip.txt fails the include glob, draft-notes.md hits the exclude glob).
+# Every delivered id is aggregated into one report carrying the count and the
+# first id, so a leaked file changes the count and fails the assertion -
+# regardless of delivery order.
+- route:
+ id: seed-files
+ from:
+ uri: "timer:seed"
+ parameters:
+ repeatCount: 1
+ steps:
+ - setBody:
+ constant: "kept document"
+ - to: "file:{{ingest.dir}}?fileName=keep.md"
+ - setBody:
+ constant: "wrong extension"
+ - to: "file:{{ingest.dir}}?fileName=skip.txt"
+ - setBody:
+ constant: "excluded draft"
+ - to: "file:{{ingest.dir}}?fileName=draft-notes.md"
+- route:
+ id: langchain4j-ingest-file-filter-route
+ from:
+ uri:
"kamelet:langchain4j-ingest-file-source?directory={{ingest.dir}}&include=**/*.md&exclude=**/draft-*&charset=UTF-8"
+ steps:
+ - log: "FILE-FILTER delivered
id=${header.CamelLangChain4jIngestDocumentId}"
+ - setBody:
+ simple: "${header.CamelLangChain4jIngestDocumentId}"
+ - aggregate:
+ aggregationStrategy:
"#class:org.apache.camel.processor.aggregate.GroupedBodyAggregationStrategy"
+ completionTimeout: 3000
+ correlationExpression:
+ constant: "all"
+ steps:
+ # the grouped list's own toString is opaque, so the report is
built explicitly
+ - setBody:
+ simple: "${body.size()}:${body[0]}"
+ - to: "{{report.url}}"
diff --git
a/tests/camel-kamelets-itest/src/test/resources/langchain4j-ingest/langchain4j-ingest-sink-filter-route.citrus.it.yaml
b/tests/camel-kamelets-itest/src/test/resources/langchain4j-ingest/langchain4j-ingest-sink-filter-route.citrus.it.yaml
new file mode 100644
index 000000000..bcde02cc4
--- /dev/null
+++
b/tests/camel-kamelets-itest/src/test/resources/langchain4j-ingest/langchain4j-ingest-sink-filter-route.citrus.it.yaml
@@ -0,0 +1,95 @@
+# ---------------------------------------------------------------------------
+# Licensed to the Apache Software Foundation (ASF) under one or more
+# contributor license agreements. See the NOTICE file distributed with
+# this work for additional information regarding copyright ownership.
+# The ASF licenses this file to You under the Apache License, Version 2.0
+# (the "License"); you may not use this file except in compliance with
+# the License. You may obtain a copy of the License at
+#
+# http://www.apache.org/licenses/LICENSE-2.0
+#
+# Unless required by applicable law or agreed to in writing, software
+# distributed under the License is distributed on an "AS IS" BASIS,
+# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+# See the License for the specific language governing permissions and
+# limitations under the License.
+# ---------------------------------------------------------------------------
+
+name: langchain4j-ingest-sink-filter-route-test
+actions:
+ # drain any request a previous test may have left on the shared server
+ - purge:
+ endpoints:
+ - name: "httpServer"
+ timeout: 100
+ - createVariables:
+ variables:
+ - name: "report.url"
+ value: "http://localhost:${http.server.port}/result/sink-filter"
+
+ # Run a route driving three documents through a sink with includeId and
+ # minDocumentSize; each IngestResult is reported back in order
+ - camel:
+ cli:
+ run:
+ waitForRunningState: false
+ args:
+ - "--dep=camel:groovy,camel:langchain4j-ingest"
+ - "--max-seconds=60"
+ integration:
+ file:
"langchain4j-ingest/langchain4j-ingest-sink-filter-route.yaml"
+ systemProperties:
+ properties:
+ - name: "report.url"
+ value: "${report.url}"
+
+ # 1: the id misses the include patterns
+ - http:
+ server: "httpServer"
+ receiveRequest:
+ POST:
+ path: "/result/sink-filter"
+ type: "plaintext"
+ body:
+ data: "@contains('documentId=notes.txt, segmentsWritten=0,
outcome=FILTERED')@"
+ - http:
+ server: "httpServer"
+ sendResponse:
+ response:
+ status: 200
+ reasonPhrase: "OK"
+ version: "HTTP/1.1"
+
+ # 2: the document is shorter than minDocumentSize
+ - http:
+ server: "httpServer"
+ receiveRequest:
+ POST:
+ path: "/result/sink-filter"
+ type: "plaintext"
+ body:
+ data: "@contains('documentId=stub.md, segmentsWritten=0,
outcome=FILTERED')@"
+ - http:
+ server: "httpServer"
+ sendResponse:
+ response:
+ status: 200
+ reasonPhrase: "OK"
+ version: "HTTP/1.1"
+
+ # 3: matching id, well-sized document
+ - http:
+ server: "httpServer"
+ receiveRequest:
+ POST:
+ path: "/result/sink-filter"
+ type: "plaintext"
+ body:
+ data: "@contains('documentId=camels.md, segmentsWritten=1,
outcome=INGESTED')@"
+ - http:
+ server: "httpServer"
+ sendResponse:
+ response:
+ status: 200
+ reasonPhrase: "OK"
+ version: "HTTP/1.1"
diff --git
a/tests/camel-kamelets-itest/src/test/resources/langchain4j-ingest/langchain4j-ingest-sink-filter-route.yaml
b/tests/camel-kamelets-itest/src/test/resources/langchain4j-ingest/langchain4j-ingest-sink-filter-route.yaml
new file mode 100644
index 000000000..84dbbc5e1
--- /dev/null
+++
b/tests/camel-kamelets-itest/src/test/resources/langchain4j-ingest/langchain4j-ingest-sink-filter-route.yaml
@@ -0,0 +1,77 @@
+# ---------------------------------------------------------------------------
+# Licensed to the Apache Software Foundation (ASF) under one or more
+# contributor license agreements. See the NOTICE file distributed with
+# this work for additional information regarding copyright ownership.
+# The ASF licenses this file to You under the Apache License, Version 2.0
+# (the "License"); you may not use this file except in compliance with
+# the License. You may obtain a copy of the License at
+#
+# http://www.apache.org/licenses/LICENSE-2.0
+#
+# Unless required by applicable law or agreed to in writing, software
+# distributed under the License is distributed on an "AS IS" BASIS,
+# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+# See the License for the specific language governing permissions and
+# limitations under the License.
+# ---------------------------------------------------------------------------
+
+# The sink's filter passthrough: a non-matching id and a too-short document
+# answer outcome=FILTERED, a matching well-sized one ingests. Three sequential
+# sends give the test a deterministic report order.
+- beans:
+ - name: store
+ type: dev.langchain4j.store.embedding.inmemory.InMemoryEmbeddingStore
+ - name: model
+ type: dev.langchain4j.model.embedding.EmbeddingModel
+ scriptLanguage: groovy
+ script: |
+ def embed = { String text ->
+ def random = new java.util.Random(text.hashCode())
+ float[] vector = new float[8]
+ for (int i = 0; i < 8; i++) {
+ vector[i] = random.nextFloat()
+ }
+ return dev.langchain4j.data.embedding.Embedding.from(vector)
+ }
+ return [
+ embedAll: { segments ->
+ dev.langchain4j.model.output.Response.from(segments.collect {
embed(it.text()) })
+ }
+ ] as dev.langchain4j.model.embedding.EmbeddingModel
+- route:
+ id: langchain4j-ingest-sink-filter-route
+ from:
+ uri: "timer:ingest"
+ parameters:
+ repeatCount: 1
+ steps:
+ - setBody:
+ constant: "A text file the id patterns must reject."
+ - setHeader:
+ name: "CamelLangChain4jIngestDocumentId"
+ constant: "notes.txt"
+ - to:
"kamelet:langchain4j-ingest-sink?includeId=**/*.md,*.md&minDocumentSize=25&embeddingStore=#bean:store&embeddingModel=#bean:model"
+ - log: "SINK-FILTER ${body}"
+ - convertBodyTo:
+ type: "java.lang.String"
+ - to: "{{report.url}}"
+ - setBody:
+ constant: "stub"
+ - setHeader:
+ name: "CamelLangChain4jIngestDocumentId"
+ constant: "stub.md"
+ - to:
"kamelet:langchain4j-ingest-sink?includeId=**/*.md,*.md&minDocumentSize=25&embeddingStore=#bean:store&embeddingModel=#bean:model"
+ - log: "SINK-FILTER ${body}"
+ - convertBodyTo:
+ type: "java.lang.String"
+ - to: "{{report.url}}"
+ - setBody:
+ constant: "Camels are resilient desert animals with two rows of
eyelashes."
+ - setHeader:
+ name: "CamelLangChain4jIngestDocumentId"
+ constant: "camels.md"
+ - to:
"kamelet:langchain4j-ingest-sink?includeId=**/*.md,*.md&minDocumentSize=25&embeddingStore=#bean:store&embeddingModel=#bean:model"
+ - log: "SINK-FILTER ${body}"
+ - convertBodyTo:
+ type: "java.lang.String"
+ - to: "{{report.url}}"