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

Reply via email to