id: embedding-model-migration-check
namespace: company.team
description: |
Before switching embedding models, measure how much your retrieval changes.
The same documents are indexed with the current and the candidate model,
a fixed set of probe queries runs against both, and the migration is blocked
when the top-k results overlap less than your threshold.
triggers:
- id: model_change_webhook
type: io.kestra.plugin.core.trigger.Webhook
description: Call this from CI on a pull request that changes your embedding
model, so the check runs before the change merges.
key: embedding-migration-check
- id: weekly_check
type: io.kestra.plugin.core.trigger.Schedule
description: Re-run weekly when the candidate is a moving tag (for example a
"latest" model). Shipped disabled.
disabled: true
cron: "0 6 * * 1"
inputs:
- id: document_urls
type: ARRAY
itemType: STRING
defaults:
- https://raw.githubusercontent.com/kestra-io/kestra/develop/README.md
description: Markdown or text documents that make up the corpus. Use a
representative sample of what your application actually retrieves.
- id: probe_queries
type: ARRAY
itemType: STRING
defaults:
- How do I start Kestra locally with Docker?
- Can tasks run Python or other programming languages?
- How do triggers react to events from external systems?
- Which AI and LLM providers are supported?
- What changed in Kestra 2.0 for loops?
- How do I version flows in Git and deploy them with CI/CD?
- What is the No-Code editor?
- How can I write my own plugin?
- Where do I report a bug or ask for help?
- Which license is Kestra released under?
description: Questions your users really ask. Every query is run against both indexes.
- id: current_model
type: STRING
defaults: all-minilm
description: The embedding model in production today.
- id: candidate_model
type: STRING
defaults: nomic-embed-text
description: The embedding model you want to switch to.
- id: ollama_endpoint
type: STRING
defaults: http://localhost:11434
description: Ollama server that serves both models. Swap the provider blocks for
OpenAI, Gemini, or another AI plugin provider to compare hosted models.
- id: top_k
type: INT
defaults: 5
description: Number of chunks retrieved per query and compared between models.
- id: min_overlap
type: FLOAT
defaults: 0.6
description: Minimum average top-k overlap (0.0-1.0). Below it, the migration is
blocked.
- id: chunk_chars
type: INT
defaults: 600
description: Approximate chunk size. Paragraphs are packed into chunks up to
this many characters.
tasks:
- id: list_models
type: io.kestra.plugin.core.http.Request
description: Ask Ollama which models are installed, so a missing model fails in
seconds with a clear fix instead of halfway through indexing.
uri: "{{ inputs.ollama_endpoint }}/api/tags"
retry:
type: exponential
interval: PT2S
maxInterval: PT30S
maxAttempts: 3
- id: check_models
type: io.kestra.plugin.core.flow.If
description: Ollama lists installed models as "name":"<model>:<tag>", so a plain
text match works for both "all-minilm" and "all-minilm:latest".
condition: >-
{{ not (outputs.list_models.body contains ('"name":"' ~
inputs.current_model)) or not (outputs.list_models.body contains
('"name":"' ~ inputs.candidate_model)) }}
then:
- id: models_missing
type: io.kestra.plugin.core.execution.Fail
errorMessage: >-
Embedding model missing on {{ inputs.ollama_endpoint }}. Install it
with: {% for m in [inputs.current_model, inputs.candidate_model] %}{%
if not (outputs.list_models.body contains ('"name":"' ~ m)) %}ollama
pull {{ m }}; {% endif %}{% endfor %}
- id: build_indexes
type: io.kestra.plugin.core.flow.WorkingDirectory
description: Chunk the corpus once, then embed the identical chunks with both
models so the comparison only reflects the model.
tasks:
- id: chunk_documents
type: io.kestra.plugin.scripts.python.Script
description: Download the documents and write one file per chunk. Each chunk
starts with its id, so search results can be matched across models.
taskRunner:
type: io.kestra.plugin.core.runner.Process
env:
DOCUMENT_URLS: "{{ inputs.document_urls | toJson }}"
CHUNK_CHARS: "{{ inputs.chunk_chars }}"
script: |
import json
import os
import re
import urllib.request
limit = int(os.environ["CHUNK_CHARS"])
os.makedirs("chunks", exist_ok=True)
count = 0
for url in json.loads(os.environ["DOCUMENT_URLS"]):
text = urllib.request.urlopen(url, timeout=60).read().decode("utf-8", "replace")
text = re.sub(r"<[^>]+>|!\[[^\]]*\]\([^)]*\)", " ", text)
chunk = ""
for para in re.split(r"\n\s*\n", text):
para = " ".join(para.split())
if len(para) < 40:
continue
if chunk and len(chunk) + len(para) > limit:
count += 1
with open("chunks/c%04d.txt" % count, "w") as f:
f.write("[c%04d] %s" % (count, chunk))
chunk = ""
chunk = (chunk + " " + para).strip()
if chunk:
count += 1
with open("chunks/c%04d.txt" % count, "w") as f:
f.write("[c%04d] %s" % (count, chunk))
if count == 0:
raise SystemExit("no text found in document_urls")
print("wrote %d chunks" % count)
print('::{"outputs":{"chunks":%d}}::' % count)
- id: index_current
type: io.kestra.plugin.ai.rag.IngestDocument
provider:
type: io.kestra.plugin.ai.provider.Ollama
endpoint: "{{ inputs.ollama_endpoint }}"
modelName: "{{ inputs.current_model }}"
embeddings:
type: io.kestra.plugin.ai.embeddings.KestraKVStore
kvName: "{{ flow.id }}-current"
drop: true
fromPath: chunks
retry:
type: exponential
interval: PT2S
maxInterval: PT30S
maxAttempts: 3
- id: index_candidate
type: io.kestra.plugin.ai.rag.IngestDocument
provider:
type: io.kestra.plugin.ai.provider.Ollama
endpoint: "{{ inputs.ollama_endpoint }}"
modelName: "{{ inputs.candidate_model }}"
embeddings:
type: io.kestra.plugin.ai.embeddings.KestraKVStore
kvName: "{{ flow.id }}-candidate"
drop: true
fromPath: chunks
retry:
type: exponential
interval: PT2S
maxInterval: PT30S
maxAttempts: 3
- id: probe
type: io.kestra.plugin.core.flow.Loop
description: Run every probe query against both indexes in parallel.
values: "{{ inputs.probe_queries | toJson }}"
concurrencyLimit: 4
outputs:
- id: query
type: STRING
value: "{{ item.value }}"
- id: current
type: JSON
value: "{{ outputs.search_current.results | toJson }}"
- id: candidate
type: JSON
value: "{{ outputs.search_candidate.results | toJson }}"
tasks:
- id: search_both
type: io.kestra.plugin.core.flow.Parallel
tasks:
- id: search_current
type: io.kestra.plugin.ai.rag.Search
provider:
type: io.kestra.plugin.ai.provider.Ollama
endpoint: "{{ inputs.ollama_endpoint }}"
modelName: "{{ inputs.current_model }}"
embeddings:
type: io.kestra.plugin.ai.embeddings.KestraKVStore
kvName: "{{ flow.id }}-current"
query: "{{ item.value }}"
maxResults: "{{ inputs.top_k }}"
minScore: 0
fetchType: FETCH
retry:
type: exponential
interval: PT2S
maxInterval: PT30S
maxAttempts: 3
- id: search_candidate
type: io.kestra.plugin.ai.rag.Search
provider:
type: io.kestra.plugin.ai.provider.Ollama
endpoint: "{{ inputs.ollama_endpoint }}"
modelName: "{{ inputs.candidate_model }}"
embeddings:
type: io.kestra.plugin.ai.embeddings.KestraKVStore
kvName: "{{ flow.id }}-candidate"
query: "{{ item.value }}"
maxResults: "{{ inputs.top_k }}"
minScore: 0
fetchType: FETCH
retry:
type: exponential
interval: PT2S
maxInterval: PT30S
maxAttempts: 3
- id: compare
type: io.kestra.plugin.jdbc.duckdb.Queries
description: Match chunk ids between the two result lists and score each query.
overlap is the share of the top-k both models return; top1_same says
whether the best chunk is unchanged.
inputFiles:
probes.json: "{{ outputs.probe.outputs | toJson }}"
fetchType: FETCH
sql: |
CREATE TABLE hits AS
SELECT p.outputs.query AS query, 'current' AS model, r.pos AS rank, regexp_extract(r.hit, '^\[(c\d+)\]', 1) AS chunk
FROM read_json_auto('{{ workingDir }}/probes.json') p, unnest(p.outputs.current) WITH ORDINALITY AS r(hit, pos)
UNION ALL
SELECT p.outputs.query, 'candidate', r.pos, regexp_extract(r.hit, '^\[(c\d+)\]', 1)
FROM read_json_auto('{{ workingDir }}/probes.json') p, unnest(p.outputs.candidate) WITH ORDINALITY AS r(hit, pos);
CREATE TABLE per_query AS
SELECT query,
round(len(list_intersect(list(chunk) FILTER (WHERE model = 'current'),
list(chunk) FILTER (WHERE model = 'candidate'))) / {{ inputs.top_k }}, 2) AS overlap,
max(chunk) FILTER (WHERE model = 'current' AND rank = 1)
= max(chunk) FILTER (WHERE model = 'candidate' AND rank = 1) AS top1_same,
string_agg(chunk, ' ' ORDER BY rank) FILTER (WHERE model = 'current') AS current_top,
string_agg(chunk, ' ' ORDER BY rank) FILTER (WHERE model = 'candidate') AS candidate_top
FROM hits
GROUP BY query;
SELECT round(avg(overlap), 3) AS avg_overlap,
round(avg(CASE WHEN top1_same THEN 1 ELSE 0 END), 3) AS top1_agreement,
count(*) AS queries,
count(*) FILTER (WHERE overlap < {{ inputs.min_overlap }}) AS queries_below_threshold
FROM per_query;
SELECT query, overlap, top1_same, current_top, candidate_top
FROM per_query ORDER BY overlap, query;
- id: gate
type: io.kestra.plugin.core.flow.If
description: Block the migration when the models disagree too much on what to retrieve.
condition: "{{ outputs.compare.outputs[0].rows[0].avg_overlap < inputs.min_overlap }}"
then:
- id: alert_blocked
type: io.kestra.plugin.slack.notifications.SlackIncomingWebhook
url: "{{ secret('SLACK_WEBHOOK_URL') }}"
payload: |
{
"text": "Embedding migration {{ inputs.current_model }} -> {{ inputs.candidate_model }} BLOCKED: average top-{{ inputs.top_k }} overlap {{ outputs.compare.outputs[0].rows[0].avg_overlap }} is below {{ inputs.min_overlap }} (top-1 agreement {{ outputs.compare.outputs[0].rows[0].top1_agreement }}). Queries that change most:\n{% for r in outputs.compare.outputs[1].rows | slice(0, 3) %}- \"{{ r.query }}\": overlap {{ r.overlap }}\n{% endfor %}Re-test these queries with real users before switching. Execution {{ execution.id }}."
}
- id: blocked
type: io.kestra.plugin.core.execution.Exit
state: FAILED
else:
- id: log_passed
type: io.kestra.plugin.core.log.Log
message: |
Embedding migration {{ inputs.current_model }} -> {{ inputs.candidate_model }} passed: average top-{{ inputs.top_k }} overlap {{ outputs.compare.outputs[0].rows[0].avg_overlap }}, top-1 agreement {{ outputs.compare.outputs[0].rows[0].top1_agreement }} over {{ outputs.compare.outputs[0].rows[0].queries }} queries.
errors:
- id: alert_failure
type: io.kestra.plugin.slack.notifications.SlackIncomingWebhook
url: "{{ secret('SLACK_WEBHOOK_URL') }}"
messageText: >-
Embedding migration check {{ inputs.current_model }} -> {{
inputs.candidate_model }} could not run (execution {{ execution.id }},
task `{{ errorLogs()[0]['taskId'] ?? 'unknown' }}`): {{
(errorLogs()[0]['message'] ?? 'see the execution logs') | split(' \[\[') |
first }}
outputs:
- id: summary
type: JSON
description: Average overlap, top-1 agreement, and the number of queries below
the threshold.
value: "{{ outputs.compare.outputs[0].rows[0] ?? {} }}"
- id: per_query
type: JSON
description: Overlap and the top-k chunk ids from each model for every probe
query, most changed first.
value: "{{ outputs.compare.outputs[1].rows ?? [] }}"