Repository navigation
[feat]: New Durable RAG support - #1205
Conversation
Signed-off-by: yaron2 <schneider.yaron@live.com>
Codecov Report❌ Patch coverage is Additional details and impacted files@@ Coverage Diff @@
## main #1205 +/- ##
==========================================
+ Coverage 84.30% 85.46% +1.16%
==========================================
Files 132 162 +30
Lines 10534 12594 +2060
==========================================
+ Hits 8881 10764 +1883
- Misses 1653 1830 +177 ☔ View full report in Codecov by Harness. 🚀 New features to boost your workflow:
|
CasperGN
left a comment
There was a problem hiding this comment.
Thanks for this. The workflow mechanics are solid: I found no determinism problems in the orchestrators, large payloads stay out of workflow history, and crash-resume within a single run looks right. The unit tests pass (269; 0.7 s on the CI invocation), and ruff and mypy are clean. There are problems at the project level and in the version/event semantics that need fixing before this merges.
Blocking: project level
dapr[all]now pulls in PyTorch and CUDA for everyone.unstructured[pdf]is inall, which thedevgroup installs, so every CI job and every contributor'suv syncdownloads torch and the full NVIDIA CUDA stack. The lock grows from 127 to 284 packages. Estimated download for a Linux job goes from about 82 MB to about 3.5 GB (7–8 GB on disk), and macOS gets torch too. No unit test uses any of it. Please keeprag-unstructuredopt-in and out ofall. That also makes theonnxruntimeconstraint and theinvoke-custom-data/proto→custom_data_protorename unnecessary; both exist only because of this dependency.fossa-scanreports 23 license issues. CI only has a push-only key, so the details are in the FOSSA web app. They need to be resolved or triaged, and I'd expect most to come from the same dependency tree.
Blocking: behaviour
- The example rebuilds the live version in place, so the "separate, inactive version" promise doesn't hold.
ReconciliationTriggernames versions by day (reconciliation_workflow.py L56), andstart()doesn't refuse the active version. A second event on the same day writes straight into the index queries are reading. - Stale chunks are never removed, and validation doesn't notice. Chunk IDs include the content hash, so a changed document adds new chunks next to the old ones. Nothing in
dapr/ext/ragcallsdelete_document, deleted source documents keep their chunks, and validation checksactual >= expected(pipeline.py L977). It also counts completion records left by an earlier run of the same version. Please makestart()refuse the active version, delete a document's old chunks when its content changes, drop documents missing from the manifest, and validate with equality against this run's records. - Event-driven changes can be lost. When a run is already in flight,
start()returns the existing instance, and that run's manifest predates the change, so nothing picks the change up later.DocumentChangedErroris retryable (errors.py L78), but the manifest etag never changes within a run, so it fails until retries run out. In the pub/sub example the event is marked seen before an in-memory 30 s timer fires (pubsub_trigger.py L61-L70), so a crash or an exception in that window loses it. Suggestion: debounce inside a workflow (timer plus external event), mark an event seen only afterstart()succeeds, record a "dirty" flag to re-run after an in-flight run, and makeDocumentChangedErrornon-retryable. - Chunk IDs collide across pages.
TextSplitter.splitrestartschunk_ordinalat 0 for each parsed page (splitting.py L106-L113), andcompute_chunk_iddoesn't include the page. Two pages with identical text (blank or boilerplate pages) overwrite each other, the store holds fewer chunks than expected, and the whole version fails validation. - Pinecone upserts will be rejected. Metadata includes
Nonevalues on the default paths (embedding_deployment,source_version_id,page_number) (pinecone.py L196-L210), and Pinecone rejects null metadata values.PineconeApiExceptionandUnauthorizedExceptionare on the transient list (L36-L43), so the 400, and a bad API key, are retried until retries run out. Please dropNonevalues, as the Azure adapter already does, and classify errors by status code. - Validation reads eventually consistent counts straight after the last write.
describe_index_stats(Pinecone) andget_document_count(Azure AI Search) lag recent upserts, and an undercount returnsvalid=Falserather than raising, so the retry policy never covers it and a correct run is marked FAILED. Please poll up to a bounded timeout, or raise a transient error whenactual < expected.
Non-blocking
- Run counters carry over when a version is re-run or
update_statusis redelivered, so totals double-count and old failures linger. Thefailureslist also grows without bound and is returned in workflow history. - An empty manifest page while
cursor < total_documentsmakescontinue_as_newloop forever. - On an activation conflict the activity retries and re-activates its own version, so an older build can win. The first activation also has no etag.
- Azure OData filters escape values but not field names. Version strings aren't checked against Azure index naming rules until
create_index.delete_documentstops after 1,000 chunks. query_api.pyresolves the active version twice. The S3sequencerdedup key should include bucket and key, and an Azure event withoutidmaps to one shared key.- The integration tests are opt-in, and no CI job runs them, so the crash-resume behaviour this PR is built around isn't exercised in CI.
test_sources_s3.pyandtest_parsing_langchain.pyshould usepytest.importorskip. - Given the size (about 9k lines of source plus docs, Bicep infra and examples), splitting it into a core pipeline PR and per-adapter PRs would make it much easier to review and land.
Signed-off-by: yaron2 <schneider.yaron@live.com>
Signed-off-by: yaron2 <schneider.yaron@live.com>
There was a problem hiding this comment.
Claude review · not a review sign-off.
Correction: this replaces my earlier review, which requested changes over a chunk-ID collision across pages. That was wrong. 5f9fcc3 renumbers chunk ordinals across the whole document in pipeline.py before compute_chunk_id, and its comment covers exactly that case. I had checked the splitter and the ID function, but not their caller. Sorry for the noise, and thanks for the fix. I also withdrew a point about the langchain test, which already imports lazily. Nothing I'd block on is left:
unstructuredis out ofall, and fossa-scan and every check are green.start()refuses the active version.- Changed and dropped documents lose their old chunks, and validation checks equality.
DocumentChangedErroris non-retryable.- Pinecone drops
Nonemetadata and classifies errors by status code. - An undercount raises a transient error.
Minor
-
The
failureslist is copied forward and appended to without a bound, and it's returned in the status result, so a large run with many failures keeps growing it. Capping it (for example the last N failures plus a count) would keep that in check.
python-sdk/dapr/ext/rag/pipeline.py
Lines 1272 to 1290 in adc8b12
-
The example debounces with an in-process
threading.Timer, so events pending at a restart are lost after the handler has acked them. The module doc says so, which is fine for an example; worth repeating in the README.
python-sdk/examples/rag/reconciliation_workflow.py
Lines 118 to 126 in adc8b12
-
Small leftovers from the earlier non-blocking list:
test_sources_s3.pyimportsboto3at module level, so collection fails without the S3 extra;pytest.importorskip('boto3')would skip it instead.- The first activation write has no etag.
- The new
databricksextra and the pyspark mypy override aren't RAG work, and would sit better in their own PR.
python-sdk/tests/ext/rag/test_sources_s3.py
Lines 20 to 21 in adc8b12
python-sdk/dapr/ext/rag/pipeline.py
Lines 1154 to 1178 in adc8b12
Lines 98 to 105 in adc8b12
Checked: 5f9fcc3 against the earlier review, reading the code at this head. Tests were not run locally, because of the torch-sized RAG extras. The RAG unit tests do run in CI (pytest tests -m "not e2e").
Reviewed at adc8b12.
Disagree, or want to talk it through? Reply /human and Casper will pick it up.
Withdrawn: the blocking chunk-ID finding was wrong (5f9fcc3 already renumbers ordinals across the document). The remaining points are minor and non-blocking.
|
@CasperGN is this good to merge? |
CasperGN
left a comment
There was a problem hiding this comment.
Claude review · approved on Casper's behalf.
Yes, good to go. I rechecked every point from both earlier rounds at this head, and each blocking one is fixed:
unstructuredis out ofall, and fossa-scan and every other check are green.start()refuses the active version, and the example picks a version other than the active one.- Changed and dropped documents lose their old chunks, and validation checks equality.
DocumentChangedErroris non-retryable. The example marks an event seen only afterstart()returns, and defers changes that arrive while a run is in flight.- Chunk ordinals are renumbered across the whole document before the ID is computed.
- Pinecone drops
Nonemetadata and classifies errors by status code. - An undercount raises a transient error.
The minor points below are left over from earlier rounds and are still open. None of them blocks the merge; a follow-up is fine. I've dropped last round's point about the databricks extra and the pyspark mypy override. It was wrong: both are already on main, and this PR doesn't add them.
Minor
-
The
failureslist is copied forward and appended to without a bound, and it's returned in the status result. Capping it (for example the last N failures plus a count) would keep a large failing run in check.
python-sdk/dapr/ext/rag/pipeline.py
Lines 1272 to 1290 in adc8b12
-
The example debounces with an in-process
threading.Timer, so events still pending at a restart are lost. The module doc says so; worth repeating in the README.
python-sdk/examples/rag/reconciliation_workflow.py
Lines 118 to 126 in adc8b12
-
Small leftovers:
test_sources_s3.pyimportsboto3at module level, so collection fails without the S3 extra;pytest.importorskip('boto3')would skip it instead.- The first activation write has no etag, and activation has no ordering check, so an older build that finishes after a newer one still activates over it.
- Azure AI Search
delete_documentfetches at most 1,000 chunk IDs, so a larger document leaves chunks behind (validation then fails rather than serving them).
python-sdk/tests/ext/rag/test_sources_s3.py
Lines 20 to 21 in adc8b12
python-sdk/dapr/ext/rag/pipeline.py
Lines 1154 to 1178 in adc8b12
python-sdk/dapr/ext/rag/vector_stores/azure_ai_search.py
Lines 335 to 343 in adc8b12
Reviewed at adc8b12.
Disagree, or want to talk it through? Reply /human and Casper will pick it up.
Adds
dapr.ext.rag- a DurableRAGPipeline extension that runs RAG ingestion (discover → parse/chunk → embed → write to a vector index) as a Dapr Workflow instead of a plain script.Value prop: ingestion pipelines that embed documents are slow and expensive to redo. Today, a crash midway through means either re-embedding everything or hand-rolled checkpointing. This makes crash recovery a property of the
runtime: kill the worker process at document 9 of 10, restart it, and Dapr Workflow resumes from document 9 automatically - no re-embedding, no lost progress, and no half-built index ever exposed to queries (a version only
becomes active after it's fully validated).
Included:
Testing: full unit suite (mocked I/O) per module, plus a real, opt-in e2e profile against LocalStack S3 + real pgvector + a real Dapr sidecar that proves: idempotent re-runs (unchanged content isn't re-embedded), and the literal
headline case — a real worker process hard-killed mid-run, with a second process resuming the same in-flight workflow instance. Manually verified end-to-end against a real OpenAI key as well (ingest → embed → activate → resolve
→ query, correct top-ranked result).