Repository navigation
Conversation
Adds an Onnx provider that computes embeddings inside the worker with ONNX Runtime (CPU), from a BERT-style .onnx model and its Hugging Face tokenizer.json supplied as nsfile://, kestra:// or allowed file:// URIs. Embeddings only. Loaded models are cached per worker by sha256 of model + tokenizer + pooling mode: langchain4j never closes the ONNX Runtime session, so reloading on every task run leaked ~115 MB of native memory per run. shadowJar excludes onnxruntime debug symbols (*.pdb, *.dSYM), cutting the added size from ~115 MB to ~51 MB. close kestra-io#354
|
Hi @fdelbrayelle requesting your review here,
also wanted to discuss about things that may pop up in future.
|
jymaire
left a comment
There was a problem hiding this comment.
Thanks a lot for this contribution @akhil838, and for the detailed write-up on the JAR size!
I went through it against the plan agreed in #354, and the scope matches: embeddings only, no bundled model, MEAN/CLS pooling, and it plugs into the existing RAG tasks without any change to them. It also follows the conventions of the other providers (UnsupportedOperationException for unsupported model types, internalStorageURI properties read through URIFetcher, icon, and the doc entry).
What I checked locally on 159574e:
./gradlew test --tests 'io.kestra.plugin.ai.provider.OnnxTest': 2/2 passed.shadowJar: 270.4 MiB. The onnxruntime.pdb/.dSYMdebug symbols are gone, and the onnxruntime and DJL tokenizers native libraries are still there for linux-x64, linux-aarch64, osx-aarch64, osx-x64 and win-x64. So the +51 MB figure from the issue holds.
I have one point, about the model cache. See the inline comment.
…sions The provider now owns each OrtSession and keeps at most max-loaded-models (default 2) loaded per worker, unloading the least recently used idle model. A model is never unloaded while an embedding call uses it, and close() waits for segments still running.
jymaire
left a comment
There was a problem hiding this comment.
Kestra Plugin Code Review
Business Requirements — met
The provider still matches the plan agreed in #354 (embeddings only, bring-your-own model and tokenizer, MEAN/CLS pooling), and the model cache is now bounded.
Previous Review Follow-up — 1 of 1 addressed
Baseline: 159574e. Delta reviewed: 9a00790 (the merge 538bc17 only brings in main's version bumps).
- ✅ [MEDIUM] Unbounded static model cache (
Onnx.java):ModelCacheis now an LRU capped at 2 models by default (max-loaded-modelsplugin configuration, validated >= 1). The plugin owns theOrtSessionand closes it on eviction. A model is never closed while an embedding call uses it (activeUsescounter), and close waits for in-flight segments. A failed load removes its entry, and a close failure is logged without failing the call. The docs describe the behaviour, andOnnxModelCacheTestcovers reuse, LRU eviction, in-use protection, concurrent first load, close failure and interrupted-call close. - Threads resolved: the thread on the cache (comment 4177379687).
Kestra Guidelines — 1 finding
- 🟢
Onnx.java:253[Guidelines] Multi-line javadoc narration
Security (OWASP Top 10:2025 + KPS) — 0 findings
✅ No issues found. The build.gradle changes only add langchain4j-embeddings, shadow exclusions and debug-symbol excludes; .github/ is untouched.
Performance — 0 findings
✅ No issues found. Loading happens outside the cache lock, and closing happens outside the lock too.
Verdict: APPROVE (recommendation only, pending CI on 538bc17; the fork workflow run was approved and has not finished yet). The plugin is on kestraVersion=2.0.2, so the 1.3/2.0 compatibility check does not apply.
QA report: Onnx in-process embedding provider (head
|
| # | Flow | Covers | Status |
|---|---|---|---|
| 1 | onnx_rag |
@Example 1 as written: HF Download, then IngestDocument and Search with kestra:// URIs |
✅ SUCCESS |
| 2 | onnx_bge_ingest + onnx_bge_search |
@Example 2 as written (nsfile://, poolingMode: CLS), with a seeding ingest |
✅ SUCCESS |
| 3 | onnx_lru_sequence |
3 distinct models in sequence (default limit 2), reuse of the evicted model and of a cached one, file:// |
✅ SUCCESS |
| 4 | onnx_rag_chat |
rag.ChatCompletion with an OpenAI chat provider (mock) and the Onnx embedding provider | ✅ SUCCESS |
| 5 | onnx_parallel_same_model ×3 concurrent |
18 concurrent tasks sharing one loaded model | ✅ SUCCESS |
| 6 | onnx_parallel_churn ×12 (4 concurrent at a time) |
3 models in use at the same time with limit 2 | ✅ SUCCESS |
| 7 | onnx_no_evict_control ×2 |
Control: 60 calls on 2 models, no eviction | ✅ SUCCESS, memory flat |
| 8 | onnx_evict_cycle ×3 |
30 reloads/evictions per run (A,B,C,A,B,C…) | ❌ runs 1–2 SUCCESS, container OOM-killed during run 3 |
| 9 | onnx_errors |
7 error paths (allowFailure) | ✅ WARNING (every error is clear) |
| – | Embedding correctness | Vectors read from the KV store compared with a Python onnxruntime + tokenizers reference |
✅ 384/384/768-d, L2 norm 1.0, cosine 1.000000, max abs diff ≤ 2e-7 |
| – | UI/docs | Plugin page, icon, editor validation, No-code form | ✅ (notes below) |
| – | Shadow JAR / natives | Size, native libs loading on the image's arch (x86_64) | ✅ loads; size note below |
Findings
MEDIUM (new): eviction does not bound worker memory on the Kestra image, so a worker that keeps swapping models in and out is OOM-killed
src/main/java/io/kestra/plugin/ai/provider/Onnx.java: LoadedModel.load() (L296–300, environment.createSession(...)) runs on whichever Kestra worker thread runs the task, and ModelCache.close() (L452) → LoadedModel.close() (L336) frees the session on whichever thread evicts it.
Repro
- Run
onnx_evict_cycle(flow 8): 3 models, defaultmax-loaded-models: 2, used in round-robin. - Each use reloads one model and evicts another, so each run does 30 load/close cycles.
| Step | anon + swap of the container |
|---|---|
| After restart | 776 MiB |
3 models loaded once (onnx_lru_sequence) |
1460 MiB |
| Control: 60 calls on 2 models, no eviction | 1587 MiB (flat) |
onnx_evict_cycle run 1 |
4505 MiB |
onnx_evict_cycle run 2 |
5929 MiB (2 GiB already in swap) |
onnx_evict_cycle run 3 |
container OOM-killed (OOMKilled=true, exit 137) |
Same 3 runs with -e MALLOC_ARENA_MAX=2 |
2076 → 2050 → 2291 MiB (bounded) |
Cause. The sessions are closed. The JVM thread count stayed at 220 → 256 across 30 reloads, while a leaked session would leave its ONNX Runtime intra-op threads behind. The memory they free stays in glibc's per-thread malloc arenas instead of going back to the OS. Each model is allocated on a different long-lived Kestra worker thread, and the worker pool grew to 65 threads here, so every arena keeps its own high-water mark. The JAR alone reproduces this outside Kestra: 40 cycles of LoadedModel.load → embedAll → close spread over 32 long-lived threads take the RSS from 303 to 4382 MiB. The same cycles on one thread stay flat at about 650 MiB, and with MALLOC_ARENA_MAX=2 they stay at about 850 MiB.
Impact. The new LRU exists to bound memory, but here it does the opposite. With more distinct models than max-loaded-models (2 by default, for example 3 RAG flows using different models on one worker), every reload costs about 50–100 MiB that is never returned, until the worker is killed with all its running tasks.
Suggestion. Create and close the OrtSession (and tokenizer) on one dedicated long-lived thread, for example a static single-thread executor owned by ModelCache, so that model memory always comes from the same arena and gets reused. The single-thread run above supports this direction, but I did not test it inside Kestra. At minimum, document the behaviour and the MALLOC_ARENA_MAX workaround. An INFO/DEBUG log line on eviction would also make thrashing visible: today only Loading ONNX embedding model is logged.
LOW
- Misleading error for an incompatible tokenizer. MiniLM paired with the GPT-2
tokenizer.jsonfails withCannot embed empty or whitespace-only text. It is thrown by langchain4j'sDimensionAwareEmbeddingModel.dimension()probe, throughOnnx$CachedEmbeddingModel.dimension(Onnx.java L268). The task fails cleanly and nothing crashes, but the message does not point to the tokenizer. - Inherited fields that Onnx ignores. The docs page and the No-code form show the inherited
baseUrl,caPemandclientPem("Connection" group) on Onnx, and they do nothing there. These come fromModelProvider, so this is cosmetic.
Info / maintainer decisions
- Shadow JAR size. It is 271 MB, against 233 MB for released
plugin-ai2.2.3 (+38 MB, +16%). It ships ONNX Runtime and DJL tokenizers natives for linux-x64, linux-aarch64, osx-x64, osx-aarch64 and win-x64, about 150 MB uncompressed..pdband.dSYMare excluded as intended. Keeping only the Linux natives would save most of the increase. That is a maintainer call, since Kestra workers run on Linux. - The natives load fine on the
kestra/kestra:v2.0.5image (x86_64). aarch64 was not tested. - Not tested: the
max-loaded-modelsplugin configuration (it needs a server-level plugin config), and real Hugging Face models larger than about 90 MB. - No-code
1 Error(s)on the provider sub-form. The message isclass …provider.Onnx cannot be cast to class io.kestra.core.models.tasks.Task, and the existingOpenAIprovider shows the same thing (control). This is pre-existing Kestra UI behaviour, not this PR. The flow itself shows Valid.
Flow 1: onnx_rag (✅ SUCCESS), @Example 1 as written
Flow YAML
id: onnx_rag
namespace: company.ai
tasks:
- id: model
type: io.kestra.plugin.core.http.Download
uri: https://huggingface.co/sentence-transformers/all-MiniLM-L6-v2/resolve/main/onnx/model.onnx
- id: tokenizer
type: io.kestra.plugin.core.http.Download
uri: https://huggingface.co/sentence-transformers/all-MiniLM-L6-v2/resolve/main/tokenizer.json
- id: ingest
type: io.kestra.plugin.ai.rag.IngestDocument
provider:
type: io.kestra.plugin.ai.provider.Onnx
modelName: all-MiniLM-L6-v2
modelUri: "{{ outputs.model.uri }}"
tokenizerUri: "{{ outputs.tokenizer.uri }}"
embeddings:
type: io.kestra.plugin.ai.embeddings.KestraKVStore
drop: true
fromDocuments:
- content: Kestra is an open-source orchestration platform for data and AI workflows.
- content: PostgreSQL is a relational database.
- id: search
type: io.kestra.plugin.ai.rag.Search
provider:
type: io.kestra.plugin.ai.provider.Onnx
modelName: all-MiniLM-L6-v2
modelUri: "{{ outputs.model.uri }}"
tokenizerUri: "{{ outputs.tokenizer.uri }}"
embeddings:
type: io.kestra.plugin.ai.embeddings.KestraKVStore
query: Which tool orchestrates workflows?
maxResults: 1
minScore: 0.3
fetchType: FETCHGantt (screenshot)
| Task | Status | Duration |
|---|---|---|
model (Download 90 MB from HF) |
SUCCESS | 13.9s |
tokenizer |
SUCCESS | 0.4s |
ingest |
SUCCESS | 3.6s (first model load) |
search |
SUCCESS | 0.2s |
| Total | SUCCESS | 18.4s |
Logs synthesis: Download logged the HF 302 redirect and downloaded with size '90405214'. ingest logged Loading ONNX embedding model (sha256 6fd5d72f…) once. search reused the cached model and logged no second load.
Outputs synthesis (screenshot): ingest: ingestedDocuments: 2, inputTokenCount: 28, kvName: onnx_rag-embedding-store. search: results: ["Kestra is an open-source orchestration platform for data and AI workflows."], size: 1. The KV store holds 2 × 384-d vectors with L2 norm 1.0000.
Flow 2: onnx_bge_ingest + onnx_bge_search (✅ SUCCESS), @Example 2 as written
onnx_bge_ingest seeds the KV onnx_bge_search-embedding-store with 4 documents, using the same BGE files and CLS.
Flow YAML
id: onnx_bge_ingest
namespace: company.ai
description: Seeds the KV store read by the onnx_bge_search example (same model files, CLS pooling)
tasks:
- id: ingest
type: io.kestra.plugin.ai.rag.IngestDocument
provider:
type: io.kestra.plugin.ai.provider.Onnx
modelName: bge-small-en-v1.5
modelUri: nsfile:///models/bge-small-en-v1.5/model.onnx
tokenizerUri: nsfile:///models/bge-small-en-v1.5/tokenizer.json
poolingMode: CLS
embeddings:
type: io.kestra.plugin.ai.embeddings.KestraKVStore
kvName: onnx_bge_search-embedding-store
drop: true
fromDocuments:
- content: Kestra is an open-source orchestration platform for data and AI workflows.
- content: PostgreSQL is a relational database.
- content: Bananas are a yellow tropical fruit rich in potassium.
- content: Apache Kafka is a distributed event streaming platform.Flow YAML
id: onnx_bge_search
namespace: company.ai
inputs:
- id: query
type: STRING
tasks:
- id: search
type: io.kestra.plugin.ai.rag.Search
provider:
type: io.kestra.plugin.ai.provider.Onnx
modelName: bge-small-en-v1.5
modelUri: nsfile:///models/bge-small-en-v1.5/model.onnx
tokenizerUri: nsfile:///models/bge-small-en-v1.5/tokenizer.json
poolingMode: CLS
embeddings:
type: io.kestra.plugin.ai.embeddings.KestraKVStore
query: "{{ inputs.query }}"
maxResults: 3
minScore: 0.5
fetchType: FETCH| Task | Status | Duration |
|---|---|---|
onnx_bge_ingest.ingest |
SUCCESS | 0.30s |
onnx_bge_search.search |
SUCCESS | 0.08s |
Outputs synthesis (screenshot): for the query "Which tool orchestrates workflows?" the results are Kestra, then Kafka, then PostgreSQL (size: 3). For "What fruit is yellow?" the bananas document ranks first. CLS pooling therefore gives sensible rankings.
Flow 3: onnx_lru_sequence (✅ SUCCESS), 3 models in sequence with default limit 2
Flow YAML
id: onnx_lru_sequence
namespace: company.ai
description: 3 distinct models in sequence in one worker (default cache limit 2), then reuse the evicted one and a still-cached one; albert loaded via file URI
tasks:
- id: minilm_1
type: io.kestra.plugin.ai.rag.IngestDocument
provider:
type: io.kestra.plugin.ai.provider.Onnx
modelName: all-MiniLM-L6-v2
modelUri: nsfile:///models/all-MiniLM-L6-v2/model.onnx
tokenizerUri: nsfile:///models/all-MiniLM-L6-v2/tokenizer.json
embeddings:
type: io.kestra.plugin.ai.embeddings.KestraKVStore
kvName: dim-minilm
drop: true
fromDocuments:
- content: Kestra is an open-source orchestration platform for data and AI workflows.
- content: PostgreSQL is a relational database.
- id: bge_2
type: io.kestra.plugin.ai.rag.IngestDocument
provider:
type: io.kestra.plugin.ai.provider.Onnx
modelName: bge-small-en-v1.5
modelUri: nsfile:///models/bge-small-en-v1.5/model.onnx
tokenizerUri: nsfile:///models/bge-small-en-v1.5/tokenizer.json
poolingMode: CLS
embeddings:
type: io.kestra.plugin.ai.embeddings.KestraKVStore
kvName: dim-bge
drop: true
fromDocuments:
- content: Kestra is an open-source orchestration platform for data and AI workflows.
- content: PostgreSQL is a relational database.
- id: albert_3
type: io.kestra.plugin.ai.rag.IngestDocument
provider:
type: io.kestra.plugin.ai.provider.Onnx
modelName: paraphrase-albert-small-v2
modelUri: file:///models/albert/model.onnx
tokenizerUri: file:///models/albert/tokenizer.json
embeddings:
type: io.kestra.plugin.ai.embeddings.KestraKVStore
kvName: dim-albert
drop: true
fromDocuments:
- content: Kestra is an open-source orchestration platform for data and AI workflows.
- content: PostgreSQL is a relational database.
- id: minilm_again_evicted
type: io.kestra.plugin.ai.rag.Search
provider:
type: io.kestra.plugin.ai.provider.Onnx
modelName: all-MiniLM-L6-v2
modelUri: nsfile:///models/all-MiniLM-L6-v2/model.onnx
tokenizerUri: nsfile:///models/all-MiniLM-L6-v2/tokenizer.json
embeddings:
type: io.kestra.plugin.ai.embeddings.KestraKVStore
kvName: dim-minilm
query: "Which tool orchestrates workflows?"
maxResults: 1
minScore: 0.0
fetchType: FETCH
- id: albert_again_cached
type: io.kestra.plugin.ai.rag.Search
provider:
type: io.kestra.plugin.ai.provider.Onnx
modelName: paraphrase-albert-small-v2
modelUri: file:///models/albert/model.onnx
tokenizerUri: file:///models/albert/tokenizer.json
embeddings:
type: io.kestra.plugin.ai.embeddings.KestraKVStore
kvName: dim-albert
query: "Which tool orchestrates workflows?"
maxResults: 1
minScore: 0.0
fetchType: FETCHGantt (screenshot)
| Task | Status | Duration |
|---|---|---|
minilm_1 |
SUCCESS | 0.30s (cache hit from flow 1: same sha256) |
bge_2 |
SUCCESS | 0.26s (load) |
albert_3 (file://) |
SUCCESS | 0.49s (load, evicts MiniLM) |
minilm_again_evicted |
SUCCESS | 0.32s (reload: MiniLM had been evicted) |
albert_again_cached |
SUCCESS | 0.09s (cache hit, no load) |
| Total | SUCCESS | 1.6s |
Logs synthesis (screenshot): Loading ONNX embedding model appears for bge_2, albert_3 and minilm_again_evicted, and not for albert_again_cached. This is exactly the LRU order: the first model is evicted, and nothing crashes. The cache key is content-addressed, so the same bytes from kestra:// (flow 1) and nsfile:// hit the same entry. Apart from the DJL warning maxLength is not explicitly specified, use modelMaxLength: 512 (from the ALBERT tokenizer), there are no server warnings or errors.
Dimension correctness (KV stores dim-minilm, dim-bge and dim-albert compared with a Python onnxruntime 1.31 + tokenizers reference using masked-mean or CLS pooling and L2 normalization):
| Model | Pooling | Dim | cos(Kestra, reference) | max abs diff |
|---|---|---|---|---|
| all-MiniLM-L6-v2 | MEAN | 384 | 1.000000 | 1.3e-7 |
| bge-small-en-v1.5 (int8) | CLS | 384 | 1.000000 | 8.6e-8 |
| paraphrase-albert-small-v2 | MEAN | 768 | 1.000000 | 2.0e-7 |
Flow 4: onnx_rag_chat (✅ SUCCESS), full RAG with a mocked chat model
Flow YAML
id: onnx_rag_chat
namespace: company.ai
description: Full RAG with Onnx embeddings and a local mock OpenAI-compatible chat endpoint
tasks:
- id: ingest
type: io.kestra.plugin.ai.rag.IngestDocument
provider:
type: io.kestra.plugin.ai.provider.Onnx
modelName: all-MiniLM-L6-v2
modelUri: nsfile:///models/all-MiniLM-L6-v2/model.onnx
tokenizerUri: nsfile:///models/all-MiniLM-L6-v2/tokenizer.json
embeddings:
type: io.kestra.plugin.ai.embeddings.KestraKVStore
drop: true
fromDocuments:
- content: Kestra 2.0 introduced plugin artifacts and a new topology view.
- content: PostgreSQL is a relational database.
- content: Bananas are a yellow tropical fruit rich in potassium.
- id: chat_with_rag
type: io.kestra.plugin.ai.rag.ChatCompletion
chatProvider:
type: io.kestra.plugin.ai.provider.OpenAI
modelName: mock-gpt
apiKey: not-a-real-key
baseUrl: http://127.0.0.1:8089/v1
embeddingProvider:
type: io.kestra.plugin.ai.provider.Onnx
modelName: all-MiniLM-L6-v2
modelUri: nsfile:///models/all-MiniLM-L6-v2/model.onnx
tokenizerUri: nsfile:///models/all-MiniLM-L6-v2/tokenizer.json
embeddings:
type: io.kestra.plugin.ai.embeddings.KestraKVStore
contentRetrieverConfiguration:
maxResults: 1
minScore: 0.3
systemMessage: You are a helpful assistant.
prompt: What did Kestra 2.0 introduce?Gantt (screenshot)
| Task | Status | Duration |
|---|---|---|
ingest |
SUCCESS | 0.38s |
chat_with_rag |
SUCCESS | 0.22s |
| Total | SUCCESS | 0.67s |
Outputs synthesis (screenshot): the mock echoes the prompt it received, What did Kestra 2.0 introduce? Answer using the following information: Kestra 2.0 introduced plugin artifacts and a new topology view. sources[0] is the Kestra 2.0 document. The Onnx embedding retrieved the right chunk and the augmented prompt reached the chat model.
Flow 5: onnx_parallel_same_model ×3 concurrent (✅ SUCCESS)
Each document below is 300 random words. There are 30 documents per task, abbreviated here.
Flow YAML
id: onnx_parallel_same_model
namespace: company.ai
description: 6 parallel tasks sharing one model (bge), 30 x 300-word documents each
tasks:
- id: parallel
type: io.kestra.plugin.core.flow.Parallel
tasks:
- id: bge_0
type: io.kestra.plugin.ai.rag.IngestDocument
provider:
type: io.kestra.plugin.ai.provider.Onnx
modelName: bge-small-en-v1.5
modelUri: nsfile:///models/bge-small-en-v1.5/model.onnx
tokenizerUri: nsfile:///models/bge-small-en-v1.5/tokenizer.json
poolingMode: CLS
embeddings:
type: io.kestra.plugin.ai.embeddings.KestraKVStore
kvName: par-bge-0
drop: true
fromDocuments:
- content: "kestra workflow ... (300 random words from a 20-word vocabulary)"
# ... 29 more documents of the same shape
- id: bge_1
type: io.kestra.plugin.ai.rag.IngestDocument
provider:
type: io.kestra.plugin.ai.provider.Onnx
modelName: bge-small-en-v1.5
modelUri: nsfile:///models/bge-small-en-v1.5/model.onnx
tokenizerUri: nsfile:///models/bge-small-en-v1.5/tokenizer.json
poolingMode: CLS
embeddings:
type: io.kestra.plugin.ai.embeddings.KestraKVStore
kvName: par-bge-1
drop: true
fromDocuments:
- content: "kestra workflow ... (300 random words from a 20-word vocabulary)"
# ... 29 more documents of the same shape
- id: bge_2
type: io.kestra.plugin.ai.rag.IngestDocument
provider:
type: io.kestra.plugin.ai.provider.Onnx
modelName: bge-small-en-v1.5
modelUri: nsfile:///models/bge-small-en-v1.5/model.onnx
tokenizerUri: nsfile:///models/bge-small-en-v1.5/tokenizer.json
poolingMode: CLS
embeddings:
type: io.kestra.plugin.ai.embeddings.KestraKVStore
kvName: par-bge-2
drop: true
fromDocuments:
- content: "kestra workflow ... (300 random words from a 20-word vocabulary)"
# ... 29 more documents of the same shape
- id: bge_3
type: io.kestra.plugin.ai.rag.IngestDocument
provider:
type: io.kestra.plugin.ai.provider.Onnx
modelName: bge-small-en-v1.5
modelUri: nsfile:///models/bge-small-en-v1.5/model.onnx
tokenizerUri: nsfile:///models/bge-small-en-v1.5/tokenizer.json
poolingMode: CLS
embeddings:
type: io.kestra.plugin.ai.embeddings.KestraKVStore
kvName: par-bge-3
drop: true
fromDocuments:
- content: "kestra workflow ... (300 random words from a 20-word vocabulary)"
# ... 29 more documents of the same shape
- id: bge_4
type: io.kestra.plugin.ai.rag.IngestDocument
provider:
type: io.kestra.plugin.ai.provider.Onnx
modelName: bge-small-en-v1.5
modelUri: nsfile:///models/bge-small-en-v1.5/model.onnx
tokenizerUri: nsfile:///models/bge-small-en-v1.5/tokenizer.json
poolingMode: CLS
embeddings:
type: io.kestra.plugin.ai.embeddings.KestraKVStore
kvName: par-bge-4
drop: true
fromDocuments:
- content: "kestra workflow ... (300 random words from a 20-word vocabulary)"
# ... 29 more documents of the same shape
- id: bge_5
type: io.kestra.plugin.ai.rag.IngestDocument
provider:
type: io.kestra.plugin.ai.provider.Onnx
modelName: bge-small-en-v1.5
modelUri: nsfile:///models/bge-small-en-v1.5/model.onnx
tokenizerUri: nsfile:///models/bge-small-en-v1.5/tokenizer.json
poolingMode: CLS
embeddings:
type: io.kestra.plugin.ai.embeddings.KestraKVStore
kvName: par-bge-5
drop: true
fromDocuments:
- content: "kestra workflow ... (300 random words from a 20-word vocabulary)"
# ... 29 more documents of the same shape
- id: search_after
type: io.kestra.plugin.ai.rag.Search
provider:
type: io.kestra.plugin.ai.provider.Onnx
modelName: bge-small-en-v1.5
modelUri: nsfile:///models/bge-small-en-v1.5/model.onnx
tokenizerUri: nsfile:///models/bge-small-en-v1.5/tokenizer.json
poolingMode: CLS
embeddings:
type: io.kestra.plugin.ai.embeddings.KestraKVStore
kvName: par-bge-0
query: "Which tool orchestrates workflows?"
maxResults: 1
minScore: 0.0
fetchType: FETCHGantt (screenshot, first execution)
| Task | Status | Duration |
|---|---|---|
bge_0 … bge_5 (parallel) |
SUCCESS ×6 | 2.6–9.7s |
search_after |
SUCCESS | 0.29s |
| Total | SUCCESS | 10.2s |
Logs synthesis: 3 concurrent executions ran 18 tasks at the same time on one cached model. All 3 executions ended SUCCESS, with no closed/OrtException errors and no server warnings. Peak container memory was 2.33 GiB.
Flow 6: onnx_parallel_churn ×12 (✅ SUCCESS), 3 models at the same time with limit 2
Flow YAML
id: onnx_parallel_churn
namespace: company.ai
description: 3 distinct models used at the same time with the default limit of 2 (eviction pressure while in use)
tasks:
- id: parallel
type: io.kestra.plugin.core.flow.Parallel
tasks:
- id: minilm
type: io.kestra.plugin.ai.rag.IngestDocument
provider:
type: io.kestra.plugin.ai.provider.Onnx
modelName: all-MiniLM-L6-v2
modelUri: nsfile:///models/all-MiniLM-L6-v2/model.onnx
tokenizerUri: nsfile:///models/all-MiniLM-L6-v2/tokenizer.json
embeddings:
type: io.kestra.plugin.ai.embeddings.KestraKVStore
kvName: churn-minilm-{{ execution.id }}
drop: true
fromDocuments:
- content: "kestra workflow ... (300 random words from a 20-word vocabulary)"
# ... 29 more documents of the same shape
- id: bge
type: io.kestra.plugin.ai.rag.IngestDocument
provider:
type: io.kestra.plugin.ai.provider.Onnx
modelName: bge-small-en-v1.5
modelUri: nsfile:///models/bge-small-en-v1.5/model.onnx
tokenizerUri: nsfile:///models/bge-small-en-v1.5/tokenizer.json
poolingMode: CLS
embeddings:
type: io.kestra.plugin.ai.embeddings.KestraKVStore
kvName: churn-bge-{{ execution.id }}
drop: true
fromDocuments:
- content: "kestra workflow ... (300 random words from a 20-word vocabulary)"
# ... 29 more documents of the same shape
- id: albert
type: io.kestra.plugin.ai.rag.IngestDocument
provider:
type: io.kestra.plugin.ai.provider.Onnx
modelName: paraphrase-albert-small-v2
modelUri: file:///models/albert/model.onnx
tokenizerUri: file:///models/albert/tokenizer.json
embeddings:
type: io.kestra.plugin.ai.embeddings.KestraKVStore
kvName: churn-albert-{{ execution.id }}
drop: true
fromDocuments:
- content: "kestra workflow ... (300 random words from a 20-word vocabulary)"
# ... 29 more documents of the same shapeGantt (screenshot)
| Executions | Status | Duration |
|---|---|---|
| 3 rounds × 4 concurrent executions (36 tasks) | SUCCESS ×12 | 4.5–12.5s each |
Logs synthesis: no model was evicted while in use, and there were no closed-session errors. The cache goes over the limit while calls are in flight and shrinks back afterwards, as documented. 4 reloads happened across the 12 executions. The server log has no WARN/ERROR related to ONNX. Peak memory was 2.83 GiB / 4 GiB.
Flow 7: onnx_no_evict_control ×2 (✅ SUCCESS), control
Flow YAML
id: onnx_no_evict_control
namespace: company.ai
description: Control for onnx_evict_cycle - same Loop but only 2 models (fits the default limit), so no reload/eviction
tasks:
- id: loop
type: io.kestra.plugin.core.flow.Loop
values: [1,2,3,4,5,6,7,8,9,10,11,12,13,14,15]
concurrencyLimit: 1
tasks:
- id: a_minilm
type: io.kestra.plugin.ai.rag.Search
provider:
type: io.kestra.plugin.ai.provider.Onnx
modelName: all-MiniLM-L6-v2
modelUri: nsfile:///models/all-MiniLM-L6-v2/model.onnx
tokenizerUri: nsfile:///models/all-MiniLM-L6-v2/tokenizer.json
embeddings:
type: io.kestra.plugin.ai.embeddings.KestraKVStore
kvName: dim-minilm
query: "Which tool orchestrates workflows?"
maxResults: 1
minScore: 0.0
fetchType: FETCH
- id: b_bge
type: io.kestra.plugin.ai.rag.Search
provider:
type: io.kestra.plugin.ai.provider.Onnx
modelName: bge-small-en-v1.5
modelUri: nsfile:///models/bge-small-en-v1.5/model.onnx
tokenizerUri: nsfile:///models/bge-small-en-v1.5/tokenizer.json
poolingMode: CLS
embeddings:
type: io.kestra.plugin.ai.embeddings.KestraKVStore
kvName: dim-bge
query: "Which tool orchestrates workflows?"
maxResults: 1
minScore: 0.0
fetchType: FETCHGantt (screenshot): loop SUCCESS. Memory went 1460 → 1553 → 1587 MiB over 60 calls with no reload.
Flow 8: onnx_evict_cycle ×3 (❌ OOM-killed on the 3rd run), see the MEDIUM finding
Flow YAML
id: onnx_evict_cycle
namespace: company.ai
description: Cycles 3 models sequentially (A,B,C,A,B,C...) with the default limit of 2, so every use reloads a model and evicts (closes) another
tasks:
- id: loop
type: io.kestra.plugin.core.flow.Loop
values: [1,2,3,4,5,6,7,8,9,10]
concurrencyLimit: 1
tasks:
- id: a_minilm
type: io.kestra.plugin.ai.rag.Search
provider:
type: io.kestra.plugin.ai.provider.Onnx
modelName: all-MiniLM-L6-v2
modelUri: nsfile:///models/all-MiniLM-L6-v2/model.onnx
tokenizerUri: nsfile:///models/all-MiniLM-L6-v2/tokenizer.json
embeddings:
type: io.kestra.plugin.ai.embeddings.KestraKVStore
kvName: dim-minilm
query: "Which tool orchestrates workflows?"
maxResults: 1
minScore: 0.0
fetchType: FETCH
- id: b_bge
type: io.kestra.plugin.ai.rag.Search
provider:
type: io.kestra.plugin.ai.provider.Onnx
modelName: bge-small-en-v1.5
modelUri: nsfile:///models/bge-small-en-v1.5/model.onnx
tokenizerUri: nsfile:///models/bge-small-en-v1.5/tokenizer.json
poolingMode: CLS
embeddings:
type: io.kestra.plugin.ai.embeddings.KestraKVStore
kvName: dim-bge
query: "Which tool orchestrates workflows?"
maxResults: 1
minScore: 0.0
fetchType: FETCH
- id: c_albert
type: io.kestra.plugin.ai.rag.Search
provider:
type: io.kestra.plugin.ai.provider.Onnx
modelName: paraphrase-albert-small-v2
modelUri: file:///models/albert/model.onnx
tokenizerUri: file:///models/albert/tokenizer.json
embeddings:
type: io.kestra.plugin.ai.embeddings.KestraKVStore
kvName: dim-albert
query: "Which tool orchestrates workflows?"
maxResults: 1
minScore: 0.0
fetchType: FETCHGantt (run, the run in flight when the container was killed, resumed after restart): each run finishes loop SUCCESS in about 16s. The memory table is in the MEDIUM finding above and in the gallery.
Flow 9: onnx_errors (✅ WARNING, all errors are actionable)
Flow YAML
id: onnx_errors
namespace: company.ai
description: Error paths of the Onnx provider (each task allowed to fail)
tasks:
- id: errors
type: io.kestra.plugin.core.flow.Parallel
tasks:
- id: missing_model_file
type: io.kestra.plugin.ai.rag.Search
allowFailure: true
provider:
type: io.kestra.plugin.ai.provider.Onnx
modelName: missing
modelUri: nsfile:///models/does-not-exist/model.onnx
tokenizerUri: nsfile:///models/all-MiniLM-L6-v2/tokenizer.json
embeddings:
type: io.kestra.plugin.ai.embeddings.KestraKVStore
kvName: dim-minilm
query: test
minScore: 0.0
maxResults: 1
- id: http_url_as_model
type: io.kestra.plugin.ai.rag.Search
allowFailure: true
provider:
type: io.kestra.plugin.ai.provider.Onnx
modelName: http
modelUri: https://huggingface.co/sentence-transformers/all-MiniLM-L6-v2/resolve/main/onnx/model.onnx
tokenizerUri: https://huggingface.co/sentence-transformers/all-MiniLM-L6-v2/resolve/main/tokenizer.json
embeddings:
type: io.kestra.plugin.ai.embeddings.KestraKVStore
kvName: dim-minilm
query: test
minScore: 0.0
maxResults: 1
- id: file_not_allowed
type: io.kestra.plugin.ai.rag.Search
allowFailure: true
provider:
type: io.kestra.plugin.ai.provider.Onnx
modelName: notallowed
modelUri: file:///etc/passwd
tokenizerUri: file:///etc/hostname
embeddings:
type: io.kestra.plugin.ai.embeddings.KestraKVStore
kvName: dim-minilm
query: test
minScore: 0.0
maxResults: 1
- id: not_an_onnx_file
type: io.kestra.plugin.ai.rag.Search
allowFailure: true
provider:
type: io.kestra.plugin.ai.provider.Onnx
modelName: bad
modelUri: nsfile:///models/all-MiniLM-L6-v2/tokenizer.json
tokenizerUri: nsfile:///models/all-MiniLM-L6-v2/tokenizer.json
embeddings:
type: io.kestra.plugin.ai.embeddings.KestraKVStore
kvName: dim-minilm
query: test
minScore: 0.0
maxResults: 1
- id: invalid_tokenizer_json
type: io.kestra.plugin.ai.rag.Search
allowFailure: true
provider:
type: io.kestra.plugin.ai.provider.Onnx
modelName: badtok
modelUri: nsfile:///models/all-MiniLM-L6-v2/model.onnx
tokenizerUri: nsfile:///models/bogus.json
embeddings:
type: io.kestra.plugin.ai.embeddings.KestraKVStore
kvName: dim-minilm
query: test
minScore: 0.0
maxResults: 1
- id: incompatible_tokenizer_gpt2
type: io.kestra.plugin.ai.rag.Search
allowFailure: true
provider:
type: io.kestra.plugin.ai.provider.Onnx
modelName: gpt2tok
modelUri: nsfile:///models/all-MiniLM-L6-v2/model.onnx
tokenizerUri: nsfile:///models/gpt2/tokenizer.json
embeddings:
type: io.kestra.plugin.ai.embeddings.KestraKVStore
kvName: dim-minilm
query: "A query whose GPT-2 token ids exceed the BERT vocabulary: Ωmega ünïcödé tokenization"
minScore: 0.0
maxResults: 1
- id: onnx_as_chat_provider
type: io.kestra.plugin.ai.completion.ChatCompletion
allowFailure: true
provider:
type: io.kestra.plugin.ai.provider.Onnx
modelName: all-MiniLM-L6-v2
modelUri: nsfile:///models/all-MiniLM-L6-v2/model.onnx
tokenizerUri: nsfile:///models/all-MiniLM-L6-v2/tokenizer.json
messages:
- type: USER
content: helloGantt (screenshot): every task FAILED as expected (allowFailure → WARNING). The total was 3.9s.
Logs synthesis (screenshot)
| Task | Error |
|---|---|
missing_model_file |
Unable to read ONNX file from uri nsfile:///models/does-not-exist/model.onnx + File not found for URI … |
http_url_as_model |
Scheme not supported: https. Supported schemes are: [kestra, file, nsfile] (as documented; use http.Download first) |
file_not_allowed |
The path /etc/passwd is not authorized … (the allowed-paths guard works) |
not_an_onnx_file |
Unable to load the ONNX model and tokenizer. + ORT_INVALID_PROTOBUF … Protobuf parsing failed |
invalid_tokenizer_json |
Unable to load the ONNX model and tokenizer. + expected , or } at line 1 column 7 |
incompatible_tokenizer_gpt2 |
Cannot embed empty or whitespace-only text (LOW: misleading) |
onnx_as_chat_provider |
Onnx only supports embedding models. |
After these failures the cache still served later executions, for example onnx_rag_chat re-ran with SUCCESS, so a failed load leaves no stuck entry.
UI / docs
- The plugin page (screenshot) shows the Onnx icon, the title and the description, which covers the 2-model limit and
max-loaded-models.modelUri/tokenizerUridescriptions document thensfile/kestra/fileschemes. The provider appears in the Provider (21) list. The main plugin doc listsOnnx. - In the editor, the example flow shows Valid (screenshot), and the No-code provider form renders
modelName,modelUriandtokenizerUriwith their help text (screenshot).
Timeout notes
- Every scenario finished within its window except the deliberate memory test (flow 8, run 3, OOM).
- Flow 1 took 18s because of the 90 MB HF download. That time is network, not a timeout.
jymaire
left a comment
There was a problem hiding this comment.
Kestra Plugin Code Review
Follow-up to the QA report on Kestra v2.0.5, which failed on one MEDIUM. Re-verified in the code at bf2b524.
Business Requirements — met
The in-process Onnx embedding provider is implemented; the one QA failure is the native-memory issue below.
Kestra Guidelines — 1 finding
- 🟠
Onnx.java:300[Performance] Native memory unbounded across evict/reload cycles (details in Performance)
Low-severity notes, non-blocking and for maintainer awareness:
- An incompatible tokenizer surfaces as "Cannot embed empty or whitespace-only text", which hides the real cause. Consider translating it into a tokenizer/model mismatch message.
- The inherited
baseUrl,caPemandclientPemfields appear on this provider but do nothing. Hide them or document that they are ignored. - The shadow JAR grows from 233 MB to 271 MB because the native libs for all OSes are bundled. This is a maintainer decision.
Security (OWASP Top 10:2025 + KPS) — 0 findings
✅ No issues found
Performance — 1 finding
- 🟠
Onnx.java:300[Performance] OrtSession created and closed on arbitrary threads, so glibc arenas retain native memory
Verdict: REQUEST CHANGES (recommendation only, posted as COMMENT)
| var environment = OrtEnvironment.getEnvironment(); | ||
| OrtSession session = null; | ||
| try { | ||
| session = environment.createSession(modelPath.toString()); |
There was a problem hiding this comment.
🟠 [Performance] OrtSession created and closed on arbitrary threads
Problem: createSession runs on whichever Kestra worker thread executes the task (via ModelCache.loaded -> loader.get()), and eviction closes the session from another caller's thread (stopUsing/startUsing -> ModelCache.close -> LoadedModel.close(), L336-346). Sessions are closed correctly, but glibc keeps per-thread malloc arenas, so the freed native memory is not returned to the OS. QA on v2.0.5: 3 models round-robin with the default cap of 2 grew container RSS+swap 1.6 GB -> 4.5 GB -> 5.9 GB and the container was OOM-killed (exit 137) on the 3rd of 30 reload/evict cycles; with MALLOC_ARENA_MAX=2 it stays at ~2.1-2.3 GB, and a standalone repro spread over 32 threads grows 303 MB -> 4.4 GB vs ~650 MB single-threaded.
Fix: Create and close sessions on one dedicated single-thread executor owned by ModelCache (e.g. loadExecutor.submit(loader::get).get() and the same for close()), so all load/unload allocations share one arena. Also document the MALLOC_ARENA_MAX=2 workaround in the @Schema/plugin doc, and consider a smaller default max-loaded-models or less eager eviction. The fix is untested in Kestra, so please verify it with the 3-model round-robin scenario.


What changes are being made and why?
closes #354
Adds an
Onnxprovider that runs a sentence-embedding model inside the worker with ONNX Runtime, so you can do RAG without an API key or a model server. It only does embeddings, so it works anywhere an embedding provider is used (IngestDocument,Search,rag.ChatCompletion,EmbeddingStoreRetriever).chatModelandimageModelthrowUnsupportedOperationException.modelUriandtokenizerUripoint to the.onnxfile and itstokenizer.json. They are read withURIFetcher, sonsfile://,kestra://and allowedfile://all work. No model is bundled in the plugin.poolingModedefaults toMEAN(all-MiniLM, E5). BGE models needCLS.langchain4j-embeddingstakes the shaded JAR from 233 MB to 348 MB, but 64 MB of that is onnxruntime debug symbols (.pdb/.dSYM). I excluded them, so the JAR ends up at 284 MB (+51 MB) with all the native libs still there.How the changes have been QAed?
OnnxTestuses the all-MiniLM files that are already a test dependency, so it needs no network or container:I also tested it manually on a local Kestra with the plugin built from this branch. This flow downloads all-MiniLM, ingests three short documents and searches them;
searchreturns the Kestra one:I also checked loading the model from namespace files in another namespace (
nsfile://company.models/...) and ingesting a file uploaded through aFILEinput.I didn't run the full test suite, since most of it needs containers or API keys.
spotlessJavaCheckalready fails onmain, so I only formatted the two new Java files.Contributor Checklist ✅