Repository navigation
fix(streaming): stop writing a response cache nothing can read (#150) - #234
Conversation
The durable cache is keyed (session_id, row_index) and is the source of truth for resume. Streaming builds a fresh Pipeline per chunk, and each new ExecutionContext takes a new uuid4 session — so every row a chunk wrote was unreadable by the next run, which generated new ids again. Measured on a 6-row stream over 3 chunks, run twice: 6 LLM calls each time, and responses.db grew to 12 rows across 6 dead sessions. Nothing was ever read back. Those writes were not merely wasted. Up to max_pending_chunks chunks wrote that one SQLite file concurrently, which is the contention #147 papered over with busy_timeout. This removes the writes rather than the symptom. Sub-pipelines — streaming chunks and the auto-retry pass — no longer attach a durable cache. The non-streaming path is untouched: it still writes one session per run, and resume still skips rows already answered. This answers the question #150 asked. Streaming has no resume: none of execute_stream(), execute_stream_async() or execute_stream_pipelined() accept resume_from, so there is no way to ask for one, and the shared cache was accidental rather than load-bearing. The checkpointing guide now says so plainly, including the workaround — shard the data and run one resumable pipeline per shard — rather than leaving readers to infer it from a missing parameter. Tests pin both halves: streaming writes nothing, and resume still reuses cached rows. Only asserting the first would be satisfied by disabling the cache everywhere, which would make every resume re-pay for completed work. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
📝 WalkthroughWalkthroughThe pipeline now skips durable response caches for streaming chunks and retry sub-pipelines. Standard executions retain SQLite caching and resume support. Tests verify these behaviors, and the checkpointing guide documents streaming limitations. ChangesStreaming response-cache ownership
Estimated code review effort: 3 (Moderate) | ~20 minutes Sequence Diagram(s)sequenceDiagram
participant Caller
participant TopLevelPipeline
participant ChunkPipeline
participant SqliteResponseCache
Caller->>TopLevelPipeline: execute_stream_pipelined
TopLevelPipeline->>ChunkPipeline: create streaming chunk sub-pipeline
ChunkPipeline->>SqliteResponseCache: skip durable cache creation
Caller->>TopLevelPipeline: execute or execute_async with resume_from
TopLevelPipeline->>SqliteResponseCache: read and reuse cached rows
Possibly related PRs
🚥 Pre-merge checks | ✅ 5✅ Passed checks (5 passed)
✨ Finishing Touches📝 Generate docstrings
🧪 Generate unit tests (beta)
Thanks for using CodeRabbit! It's free for OSS, and your support helps us grow. If you like it, consider giving us a shout-out. Comment |
There was a problem hiding this comment.
Actionable comments posted: 2
🤖 Prompt for all review comments with AI agents
Verify each finding against current code. Fix only still-valid issues, skip the
rest with a brief reason, keep changes minimal, and validate.
Inline comments:
In `@docs/guides/checkpointing.md`:
- Around line 60-62: Update the checkpointing documentation’s duplicate-cost
statement to clarify that only rows processed before the failure are charged
again on restart; replace “every row is paid for again” with wording such as
“previously processed rows are paid for again” or “the restarted run reprocesses
all rows.”
- Around line 72-74: Update the streaming checkpointing guidance in the section
around line 284 to remove the claim that checkpoints support chunk-level resume.
State that streaming does not support resuming, or link readers to the existing
large-dataset guidance near “execute()” and checkpointing, while preserving the
rest of the documentation.
🪄 Autofix
Fix all unresolved CodeRabbit comments on this PR:
- Push a commit to this branch (recommended)
- Create a new PR with the fixes
ℹ️ Review info
⚙️ Run configuration
Configuration used: defaults
Review profile: CHILL
Plan: Pro Plus
Run ID: 79d3cf00-bdc8-4d0c-80db-d6b08a0eaf65
📒 Files selected for processing (3)
docs/guides/checkpointing.mdondine/api/pipeline.pytests/unit/test_streaming_no_orphan_cache.py
| None of them accept `resume_from`, so there is no way to ask for one. A | ||
| streamed run that dies must be restarted from the beginning, and every row is | ||
| paid for again. |
There was a problem hiding this comment.
🎯 Functional Correctness | 🟡 Minor | ⚡ Quick win
Qualify the duplicate-cost statement.
After a failure at row N, only previously processed rows were charged in the failed run. Restarting repeats those rows. Later rows incur their first charge. Replace “every row is paid for again” with “previously processed rows are paid for again” or “the restarted run reprocesses all rows.”
🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the
rest with a brief reason, keep changes minimal, and validate.
In `@docs/guides/checkpointing.md` around lines 60 - 62, Update the checkpointing
documentation’s duplicate-cost statement to clarify that only rows processed
before the failure are charged again on restart; replace “every row is paid for
again” with wording such as “previously processed rows are paid for again” or
“the restarted run reprocesses all rows.”
| **If you need resume on a large dataset**, prefer `execute()` with | ||
| checkpointing over streaming, or split the data yourself and run one pipeline | ||
| per shard so each shard has a session you can resume. |
There was a problem hiding this comment.
🎯 Functional Correctness | 🟡 Minor | ⚡ Quick win
Remove the conflicting chunk-level checkpointing guidance.
Line 284 still says that streaming checkpointing belongs at the chunk level. This conflicts with the new statement that streaming has no resume support. Update that section to describe the no-resume limitation or link to this section.
🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the
rest with a brief reason, keep changes minimal, and validate.
In `@docs/guides/checkpointing.md` around lines 72 - 74, Update the streaming
checkpointing guidance in the section around line 284 to remove the claim that
checkpoints support chunk-level resume. State that streaming does not support
resuming, or link readers to the existing large-dataset guidance near
“execute()” and checkpointing, while preserving the rest of the documentation.
Closes #150.
The question the issue asked
It does not, and it never could. None of
execute_stream(),execute_stream_async()orexecute_stream_pipelined()acceptresume_from, so there is no way to ask for one.Measured, not reasoned
The durable cache is keyed
(session_id, row_index). Streaming builds a freshPipelineper chunk, and each newExecutionContexttakes a newuuid4session — so every row a chunk writes is unreadable by the next run, which generates new ids again.A 6-row stream over 3 chunks, run twice against a counting client:
Second run re-called every row.
responses.dbgrew to 12 rows across 6 dead sessions, none of which anything can read.That is outcome (2) in the issue: the shared cache is accidental, not load-bearing.
Why it mattered beyond waste
Up to
max_pending_chunkschunk pipelines wrote that single SQLite file concurrently. That is thedatabase is lockedcontention #147 patched withbusy_timeout. This removes the writes rather than the symptom — the mitigation stays, but it now has nothing to mitigate on this path.The change
Sub-pipelines — streaming chunks and the auto-retry pass — no longer attach a durable cache.
After:
The non-streaming path is untouched, and resume still skips completed rows — verified by killing a run after 3 of 6 rows and resuming: 3 new calls, not 6.
Documented, since the absence is the surprising part
docs/guides/checkpointing.mdnow states plainly that streaming does not resume, why it is structural rather than an oversight, and the workaround: shard the data and run one resumable pipeline per shard. A missing parameter is not documentation.Tests
Both halves are pinned, deliberately:
responses.dbAsserting only the first would be satisfied by disabling the cache everywhere — which would silently make every resume re-pay for work already done. The first test fails when the sub-pipeline skip is removed.
Full suite: 1164 passing, mypy clean.
Not fixed here
Streaming still cannot resume. Making it resumable is a feature — chunks would need to share one session, or the top level would need to own a cache keyed across them — not something to smuggle into a cleanup. This makes the current behaviour honest instead of half-implemented.
ai assistance: i directed this work with help from claude code.
Summary by CodeRabbit
Bug Fixes
Documentation