Skip to content

[feat]: New Durable RAG support - #1205

Merged
yaron2 merged 8 commits into
dapr:mainfrom
yaron2:rag
Oct 1, 2026
Merged

yaron2 merged 8 commits into
dapr:mainfrom
yaron2:rag

Conversation

@yaron2

@yaron2 yaron2 commented Sep 11, 2026

Copy link
Copy Markdown
Member

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:

  • Sources: S3, Azure Blob. Embedders: OpenAI, Azure OpenAI. Vector stores: pgvector, Pinecone, Azure AI Search.
  • Versioned activation (build a new index version alongside the live one, atomically switch on completion).
  • Event-driven ingestion via Dapr pub/sub, with dedup + debounced reconciliation.
  • An Azure-native flagship path (Blob → Event Grid → Service Bus → Workflow → Azure OpenAI → Azure AI Search alias → RAG query API with grounded answers + citations), plus optional Foundry IQ knowledge-source registration.
  • FailureInjector test helper for deterministically demonstrating the crash/resume behavior.

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).

Signed-off-by: yaron2 <schneider.yaron@live.com>
@yaron2
yaron2 requested review from a team as code owners September 11, 2026 21:14
@yaron2 yaron2 changed the title [feat]: New DurableRag support [feat]: New Durable RAG support Sep 11, 2026
Signed-off-by: yaron2 <schneider.yaron@live.com>
@codecov

codecov Bot commented Sep 11, 2026 •

Copy link
Copy Markdown

Codecov Report

❌ Patch coverage is 91.40777% with 177 lines in your changes missing coverage. Please review.
✅ Project coverage is 85.46%. Comparing base (a8887e3) to head (adc8b12).

Files with missing lines Patch % Lines
dapr/ext/rag/vector_stores/azure_ai_search.py 80.59% 59 Missing ⚠️
dapr/ext/rag/pipeline.py 92.82% 28 Missing ⚠️
dapr/ext/rag/sources/s3.py 78.35% 21 Missing ⚠️
dapr/ext/rag/vector_stores/pgvector.py 81.48% 20 Missing ⚠️
dapr/ext/rag/sources/azure_blob.py 85.04% 16 Missing ⚠️
dapr/ext/rag/vector_stores/pinecone.py 90.38% 10 Missing ⚠️
dapr/ext/rag/embedding/_openai_common.py 87.87% 8 Missing ⚠️
dapr/ext/rag/embedding/openai.py 84.21% 6 Missing ⚠️
dapr/ext/rag/models.py 98.58% 4 Missing ⚠️
dapr/ext/rag/embedding/azure_openai.py 93.54% 2 Missing ⚠️
... and 3 more
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.
📢 Have feedback on the report? Share it here.

🚀 New features to boost your workflow:
  • ❄️ Test Analytics: Detect flaky tests, report on failures, and find test suite problems.

Signed-off-by: yaron2 <schneider.yaron@live.com>
Signed-off-by: yaron2 <schneider.yaron@live.com>

@CasperGN CasperGN left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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

  1. dapr[all] now pulls in PyTorch and CUDA for everyone. unstructured[pdf] is in all, which the dev group installs, so every CI job and every contributor's uv sync downloads 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 keep rag-unstructured opt-in and out of all. That also makes the onnxruntime constraint and the invoke-custom-data/proto → custom_data_proto rename unnecessary; both exist only because of this dependency.
  2. fossa-scan reports 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

  1. The example rebuilds the live version in place, so the "separate, inactive version" promise doesn't hold. ReconciliationTrigger names versions by day (reconciliation_workflow.py L56), and start() doesn't refuse the active version. A second event on the same day writes straight into the index queries are reading.
  2. 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/rag calls delete_document, deleted source documents keep their chunks, and validation checks actual >= expected (pipeline.py L977). It also counts completion records left by an earlier run of the same version. Please make start() 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.
  3. 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. DocumentChangedError is 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 after start() succeeds, record a "dirty" flag to re-run after an in-flight run, and make DocumentChangedError non-retryable.
  4. Chunk IDs collide across pages. TextSplitter.split restarts chunk_ordinal at 0 for each parsed page (splitting.py L106-L113), and compute_chunk_id doesn'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.
  5. Pinecone upserts will be rejected. Metadata includes None values on the default paths (embedding_deployment, source_version_id, page_number) (pinecone.py L196-L210), and Pinecone rejects null metadata values. PineconeApiException and UnauthorizedException are on the transient list (L36-L43), so the 400, and a bad API key, are retried until retries run out. Please drop None values, as the Azure adapter already does, and classify errors by status code.
  6. Validation reads eventually consistent counts straight after the last write. describe_index_stats (Pinecone) and get_document_count (Azure AI Search) lag recent upserts, and an undercount returns valid=False rather 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 when actual < expected.

Non-blocking

  • Run counters carry over when a version is re-run or update_status is redelivered, so totals double-count and old failures linger. The failures list also grows without bound and is returned in workflow history.
  • An empty manifest page while cursor < total_documents makes continue_as_new loop 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_document stops after 1,000 chunks.
  • query_api.py resolves the active version twice. The S3 sequencer dedup key should include bucket and key, and an Azure event without id maps 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.py and test_parsing_langchain.py should use pytest.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>
@yaron2
yaron2 requested a review from CasperGN September 29, 2026 21:15

@CasperGN CasperGN left a comment •

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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:

  • unstructured is out of all, 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.
  • DocumentChangedError is non-retryable.
  • Pinecone drops None metadata and classifies errors by status code.
  • An undercount raises a transient error.

Minor

  1. The failures list 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.

    failures = list(current.failures)
    for outcome in outcomes:
    if outcome.status == DocumentOutcomeStatus.COMPLETED.value:
    completed += 1
    elif outcome.status == DocumentOutcomeStatus.SKIPPED.value:
    skipped += 1
    else:
    failed += 1
    failures.append(
    DocumentFailure(
    document_id=outcome.document_id,
    error_type=outcome.error_type or 'Unknown',
    error_message=outcome.error_message or '',
    retryable=bool(outcome.retryable),
    )
    )
    total_chunks += outcome.chunk_count
    embedded_chunks += outcome.embedded_chunk_count

  2. 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.

    if self._pending_timer is not None:
    self._pending_timer.cancel()
    self._pending_prefix = prefix if prefix == self._pending_prefix else None
    self._pending_event_ids.append(event.event_id)
    timer = threading.Timer(self._debounce_seconds, self._reconcile)
    timer.daemon = True
    self._pending_timer = timer
    timer.start()

  3. Small leftovers from the earlier non-blocking list:

    • test_sources_s3.py imports boto3 at 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 databricks extra and the pyspark mypy override aren't RAG work, and would sit better in their own PR.

    import boto3
    from botocore.stub import Stubber

    current, etag = self._state.read_activation(pipeline_id)
    if (
    current is not None
    and current.active_version == version
    and current.manifest_hash == manifest_hash
    ):
    return to_wire(current) # already active under this exact manifest: idempotent no-op
    previous_version = current.active_version if current is not None else None
    # Store-native activation (e.g. AzureAISearchVectorStore's alias switch) runs
    # before the Dapr-state write and is a no-op for stores without one (see
    # VectorIndex.activate_version's default). This ordering makes a retry after a
    # partial failure safe: re-running finds the store's own state already correct
    # and just completes the Dapr-state write, rather than switching twice.
    self._vector_store.activate_version(version, previous_version=previous_version)
    record = ActivationRecord(
    pipeline_id=pipeline_id,
    active_version=version,
    previous_version=previous_version,
    manifest_hash=manifest_hash,
    activated_at=_utcnow_iso(),
    workflow_instance_id=workflow_instance_id,
    )
    self._state.write_activation(record, etag=etag)

    python-sdk/pyproject.toml

    Lines 98 to 105 in adc8b12

    # databricks is also empty on purpose: it schedules Dapr Workflow executions
    # (via dapr.ext.workflow, already in core) from inside a Databricks Lakeflow
    # pipeline. pyspark.pipelines is provided by that Databricks runtime itself —
    # depending on a pip-installable "pyspark" here would be both unnecessary for
    # real usage and insufficient for local imports, since pyspark.pipelines is
    # not part of the open-source PyPI pyspark package. dapr.ext.databricks
    # imports it lazily and raises a clear error if it's unavailable.
    databricks = []

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.

@CasperGN
CasperGN dismissed their stale review September 30, 2026 11:11

Withdrawn: the blocking chunk-ID finding was wrong (5f9fcc3 already renumbers ordinals across the document). The remaining points are minor and non-blocking.

@yaron2

yaron2 commented Sep 30, 2026

Copy link
Copy Markdown
Member Author

@CasperGN is this good to merge?

@CasperGN CasperGN left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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:

  • unstructured is out of all, 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.
  • DocumentChangedError is non-retryable. The example marks an event seen only after start() 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 None metadata 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

  1. The failures list 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.

    failures = list(current.failures)
    for outcome in outcomes:
    if outcome.status == DocumentOutcomeStatus.COMPLETED.value:
    completed += 1
    elif outcome.status == DocumentOutcomeStatus.SKIPPED.value:
    skipped += 1
    else:
    failed += 1
    failures.append(
    DocumentFailure(
    document_id=outcome.document_id,
    error_type=outcome.error_type or 'Unknown',
    error_message=outcome.error_message or '',
    retryable=bool(outcome.retryable),
    )
    )
    total_chunks += outcome.chunk_count
    embedded_chunks += outcome.embedded_chunk_count

  2. 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.

    if self._pending_timer is not None:
    self._pending_timer.cancel()
    self._pending_prefix = prefix if prefix == self._pending_prefix else None
    self._pending_event_ids.append(event.event_id)
    timer = threading.Timer(self._debounce_seconds, self._reconcile)
    timer.daemon = True
    self._pending_timer = timer
    timer.start()

  3. Small leftovers:

    • test_sources_s3.py imports boto3 at 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_document fetches at most 1,000 chunk IDs, so a larger document leaves chunks behind (validation then fails rather than serving them).

    import boto3
    from botocore.stub import Stubber

    current, etag = self._state.read_activation(pipeline_id)
    if (
    current is not None
    and current.active_version == version
    and current.manifest_hash == manifest_hash
    ):
    return to_wire(current) # already active under this exact manifest: idempotent no-op
    previous_version = current.active_version if current is not None else None
    # Store-native activation (e.g. AzureAISearchVectorStore's alias switch) runs
    # before the Dapr-state write and is a no-op for stores without one (see
    # VectorIndex.activate_version's default). This ordering makes a retry after a
    # partial failure safe: re-running finds the store's own state already correct
    # and just completes the Dapr-state write, rather than switching twice.
    self._vector_store.activate_version(version, previous_version=previous_version)
    record = ActivationRecord(
    pipeline_id=pipeline_id,
    active_version=version,
    previous_version=previous_version,
    manifest_hash=manifest_hash,
    activated_at=_utcnow_iso(),
    workflow_instance_id=workflow_instance_id,
    )
    self._state.write_activation(record, etag=etag)

    matches = search_client.search(
    search_text='*',
    filter=f"source_document_id eq '{_escape_odata_string(document_id)}'",
    select=['id'],
    top=1000,
    )
    keys = [match['id'] for match in matches]
    if keys:
    search_client.delete_documents(documents=[{'id': key} for key in keys])

Reviewed at adc8b12.
Disagree, or want to talk it through? Reply /human and Casper will pick it up.

@yaron2
yaron2 added this pull request to the merge queue Oct 1, 2026
@github-merge-queue
github-merge-queue Bot removed this pull request from the merge queue due to failed status checks Oct 1, 2026
@yaron2
yaron2 added this pull request to the merge queue Oct 1, 2026
Merged via the queue into dapr:main with commit 3c4cbef Oct 1, 2026
27 of 29 checks passed
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

2 participants