id: rag-knowledge-base-vector-reconciliation
namespace: company.ai
description: |
Automate synchronization and drift reconciliation between raw knowledge base documents
and production vector store indexes. Computes cryptographic content hashes (SHA-256) to
detect new, updated, and deleted source files, purges orphaned vectors, and selectively
embeds changed chunks with audit reporting.
triggers:
- id: scheduled_reconciliation_sync
type: io.kestra.plugin.core.trigger.Schedule
description: Nightly reconciliation sweep to catch out-of-band document drift.
Shipped disabled by default.
cron: "0 2 * * *"
disabled: true
- id: knowledge_repo_webhook
type: io.kestra.plugin.core.trigger.Webhook
description: Webhook triggered by Git push or CI/CD documentation release pipelines.
key: "{{ secret('DOCS_WEBHOOK_KEY') }}"
inputs:
- id: index_name
type: STRING
defaults: rag-knowledge-base
description: Target vector database index name to reconcile against source documents.
- id: batch_size
type: INT
defaults: 50
description: Maximum document chunk batch size for incremental embedding operations.
- id: purge_orphaned
type: BOOLEAN
defaults: true
description: Whether to permanently purge orphaned vectors associated with
deleted source documents.
- id: cost_per_1m_tokens
type: FLOAT
defaults: 0.02
description: Embedding model cost in USD per 1M tokens (e.g. OpenAI
text-embedding-3-small).
tasks:
- id: scan_source_documents
type: io.kestra.plugin.scripts.python.Script
description: Scans knowledge base documents, extracts metadata, and computes
cryptographic SHA-256 content hashes.
beforeCommands:
- pip install pyarrow
outputFiles:
- active_manifest.json
script: |
import os
import json
import hashlib
# Simulated enterprise knowledge base documents
# In production, this reads from an S3/GCS bucket or a cloned Git repository
corpus = [
{
"doc_id": "SEC-POL-001",
"title": "Corporate Information Security Policy",
"category": "security",
"content": "All production API keys and database credentials must be rotated every 90 days. Multi-factor authentication is mandatory across all enterprise cloud environments. Zero trust network architecture applies to all internal VPN traffic."
},
{
"doc_id": "OPS-RUN-042",
"title": "Production Incident Escalation Runbook",
"category": "operations",
"content": "P1 incidents require immediate PagerDuty page to secondary on-call engineer within 5 minutes. The incident commander must initiate a dedicated Slack triage channel and post 15-minute executive status updates."
},
{
"doc_id": "API-REF-109",
"title": "Customer Payment Gateway Integration Guide",
"category": "engineering",
"content": "All charge requests must provide an idempotency key header to prevent duplicate billing. Webhooks must verify HMAC-SHA256 signatures before processing payload events. TLS 1.3 is enforced."
},
{
"doc_id": "HR-BEN-088",
"title": "Employee Benefits and Remote Work Stipend Policy",
"category": "human_resources",
"content": "Full-time remote team members are eligible for an annual home office equipment stipend of $1,500. Health insurance enrollment opens annually in November for eligible dependents."
}
]
manifest = []
total_tokens_approx = 0
for doc in corpus:
content_bytes = doc["content"].encode("utf-8")
sha256_hash = hashlib.sha256(content_bytes).hexdigest()
char_count = len(doc["content"])
token_estimate = int(char_count / 4) # Standard 4 char/token approximation
total_tokens_approx += token_estimate
manifest.append({
"doc_id": doc["doc_id"],
"title": doc["title"],
"category": doc["category"],
"content_hash": sha256_hash,
"char_count": char_count,
"token_estimate": token_estimate
})
with open("active_manifest.json", "w") as f:
json.dump(manifest, f, indent=2)
print(f"Scanned {len(manifest)} active documents. Total corpus tokens: {total_tokens_approx}")
- id: reconcile_vector_drift
type: io.kestra.plugin.jdbc.duckdb.Queries
description: Adjudicates differences between source document hashes and vector
index catalog using DuckDB SQL.
url: "jdbc:duckdb:"
fetch: true
queries: |
-- 1. Create simulated vector store tracking catalog (representing existing vectors in database)
CREATE TABLE vector_store_catalog (
doc_id VARCHAR,
stored_hash VARCHAR,
vector_count INTEGER,
last_embedded_at TIMESTAMP
);
-- Seed existing state:
-- SEC-POL-001: Exists, hash unchanged (UP TO DATE)
-- OPS-RUN-042: Exists, but hash changed from previous version (MODIFIED / DRIFT)
-- LEGACY-DEPR-007: Exists in vector store, but removed from source corpus (ORPHANED / DELETED)
-- API-REF-109 & HR-BEN-088: Missing in vector catalog (NEW / TO INSERT)
INSERT INTO vector_store_catalog VALUES
('SEC-POL-001', '80ceca2ca0dc4fb92b49e1e12739266e7f781df0bbdf88d72dfa4f02888cf941', 3, CURRENT_TIMESTAMP - INTERVAL 10 DAY),
('OPS-RUN-042', 'stale_hash_from_previous_commit_version_0000000000000000000000000', 4, CURRENT_TIMESTAMP - INTERVAL 30 DAY),
('LEGACY-DEPR-007', 'obsolete_deprecated_policy_vectors_that_should_not_exist', 5, CURRENT_TIMESTAMP - INTERVAL 60 DAY);
-- 2. Read scanned active manifest from previous task
CREATE TABLE active_manifest AS
SELECT * FROM read_json_auto('{{ outputs.scan_source_documents.outputFiles["active_manifest.json"] }}');
-- 3. Calculate reconciliation set differences
CREATE TABLE reconciliation_summary AS
SELECT
-- Documents requiring new embedding (in active manifest, missing from vector store)
(SELECT COUNT(*) FROM active_manifest a LEFT JOIN vector_store_catalog v ON a.doc_id = v.doc_id WHERE v.doc_id IS NULL) AS to_insert_count,
-- Documents requiring re-embedding (in both, but content hash modified)
(SELECT COUNT(*) FROM active_manifest a JOIN vector_store_catalog v ON a.doc_id = v.doc_id WHERE a.content_hash <> v.stored_hash) AS to_update_count,
-- Orphaned documents requiring vector purge (in vector store, deleted from active manifest)
(SELECT COUNT(*) FROM vector_store_catalog v LEFT JOIN active_manifest a ON v.doc_id = a.doc_id WHERE a.doc_id IS NULL) AS to_delete_count,
-- Total active documents in corpus
(SELECT COUNT(*) FROM active_manifest) AS total_active_documents,
-- Total documents in vector index
(SELECT COUNT(*) FROM vector_store_catalog) AS total_stored_documents;
SELECT
to_insert_count,
to_update_count,
to_delete_count,
(to_insert_count + to_update_count + to_delete_count) AS total_drift_count,
(to_insert_count + to_update_count + to_delete_count > 0) AS drift_detected,
total_active_documents,
total_stored_documents
FROM reconciliation_summary;
- id: evaluate_drift_gate
type: io.kestra.plugin.core.flow.If
description: Branches based on whether document drift is detected between source
files and vector store.
condition: "{{ outputs.reconcile_vector_drift.rows[0].drift_detected == true }}"
then:
- id: execute_pruning_and_sync
type: io.kestra.plugin.scripts.python.Script
description: Simulates vector purge, chunks modified documents, and computes
token cost savings.
outputFiles:
- sync_summary.json
script: |
import json
to_insert = {{ outputs.reconcile_vector_drift.rows[0].to_insert_count }}
to_update = {{ outputs.reconcile_vector_drift.rows[0].to_update_count }}
to_delete = {{ outputs.reconcile_vector_drift.rows[0].to_delete_count }}
total_active = {{ outputs.reconcile_vector_drift.rows[0].total_active_documents }}
cost_per_1m = {{ inputs.cost_per_1m_tokens }}
# Naive full re-indexing would process 100% of all docs
# Incremental sync only processes (to_insert + to_update)
avg_tokens_per_doc = 350
reindexed_tokens = (to_insert + to_update) * avg_tokens_per_doc
full_reindex_tokens = total_active * avg_tokens_per_doc
saved_tokens = max(0, full_reindex_tokens - reindexed_tokens)
savings_percent = round((saved_tokens / full_reindex_tokens) * 100, 2) if full_reindex_tokens > 0 else 0.0
purged_vector_ids = ["LEGACY-DEPR-007#chunk_1", "LEGACY-DEPR-007#chunk_2", "LEGACY-DEPR-007#chunk_3", "LEGACY-DEPR-007#chunk_4", "LEGACY-DEPR-007#chunk_5"]
summary = {
"status": "RECONCILED",
"to_insert_count": to_insert,
"to_update_count": to_update,
"to_delete_count": to_delete,
"purged_vectors": len(purged_vector_ids),
"incremental_tokens_consumed": reindexed_tokens,
"tokens_saved": saved_tokens,
"embedding_cost_savings_pct": savings_percent,
"estimated_cost_usd": round((reindexed_tokens / 1_000_000) * cost_per_1m, 6)
}
with open("sync_summary.json", "w") as f:
json.dump(summary, f, indent=2)
print(f"Reconciliation complete: {to_insert} inserted, {to_update} updated, {to_delete} purged ({len(purged_vector_ids)} vectors). Saved {savings_percent}% tokens!")
- id: export_sync_audit_report
type: io.kestra.plugin.core.storage.Write
description: Persists reconciliation audit report in Kestra storage for
compliance verification.
fileName: "rag-sync-report-{{ execution.id }}.md"
content: |
# RAG Knowledge Base Vector Reconciliation Report
- **Execution ID:** `{{ execution.id }}`
- **Index Name:** `{{ inputs.index_name }}`
- **Reconciled At:** `{{ now() }}`
- **Total Drift Count:** `{{ outputs.reconcile_vector_drift.rows[0].total_drift_count }}`
- **Drift Detected:** `{{ outputs.reconcile_vector_drift.rows[0].drift_detected }}`
## Reconciliation Actions
- **New Documents Embedded (Insert):** `{{ outputs.reconcile_vector_drift.rows[0].to_insert_count }}`
- **Modified Documents Re-Embedded (Update):** `{{ outputs.reconcile_vector_drift.rows[0].to_update_count }}`
- **Orphaned Documents Purged (Delete):** `{{ outputs.reconcile_vector_drift.rows[0].to_delete_count }}`
- **Total Active Knowledge Base Documents:** `{{ outputs.reconcile_vector_drift.rows[0].total_active_documents }}`
- **Total Stored Index Documents:** `{{ outputs.reconcile_vector_drift.rows[0].total_stored_documents }}`
## Efficiency and Cost Optimization
- **Purged Orphaned Vectors:** `{{ outputs.execute_pruning_and_sync.outputFiles['sync_summary.json'] }}`
- **Embedding Optimization:** Incremental hashing prevented full-corpus re-embedding, preserving rate limits and reducing API costs.
- id: notify_slack_sync_complete
type: io.kestra.plugin.slack.notifications.SlackIncomingWebhook
description: Notifies team Slack channel of completed vector store reconciliation.
url: "{{ secret('SLACK_WEBHOOK_URL') }}"
payload: |
{
"text": "📚 *RAG Vector Reconciliation Complete: Index `{{ inputs.index_name }}`*\n• *Status:* `RECONCILED`\n• *New Documents:* `{{ outputs.reconcile_vector_drift.rows[0].to_insert_count }}`\n• *Updated Documents:* `{{ outputs.reconcile_vector_drift.rows[0].to_update_count }}`\n• *Purged Orphaned Documents:* `{{ outputs.reconcile_vector_drift.rows[0].to_delete_count }}`\n• *Execution ID:* `{{ execution.id }}`\n• *Audit Report:* `{{ outputs.export_sync_audit_report.uri }}`"
}
else:
- id: log_clean_state
type: io.kestra.plugin.core.log.Log
description: Logs that vector index is in sync with raw documents and zero drift
was detected.
message: "Vector store '{{ inputs.index_name }}' is fully synchronized with
active documentation corpus. Zero drift detected (Execution {{
execution.id }})."
errors:
- id: alert_reconciliation_failure
type: io.kestra.plugin.slack.notifications.SlackIncomingWebhook
description: Alerts Slack on-call if vector store reconciliation or pruning
encounters an error.
url: "{{ secret('SLACK_WEBHOOK_URL') }}"
payload: |
{
"text": "⚠️ *RAG Vector Reconciliation Flow Failure:* Execution {{ execution.id }} failed during index reconciliation on `{{ inputs.index_name }}`. Check Kestra execution logs."
}