Repository navigation
Feat: stream completed chunks - #157
vincevannoort wants to merge 25 commits into
Conversation
6e49065 to
b7df6a4
Compare
b7df6a4 to
8e48718
Compare
|
@ReinierMaas or @jerbaroo is one of you available to review this somewhere in the coming week(s)? |
|
Thanks for opening the PR @vincevannoort. I will do a review, but also a review from a more experienced Rustacean such as @ReinierMaas would be important |
There was a problem hiding this comment.
🟡 Changes recommended
Failure races can omit completed results, while per-chunk polling introduces substantial latency and request overhead.
4 open findings
What changed in this PR
Adds progressive chunk-result streaming so producers can consume contiguous completed chunks before submission completion.
Changes:
- Tracks ready chunk boundaries and enforces oldest-first ordering.
- Adds synchronous/asynchronous Python streaming APIs and integration tests.
- Bumps the workspace version to 0.42.0.
| File | Description |
|---|---|
opsqueue/src/consumer/strategy.rs |
Orders oldest chunks by index. |
opsqueue/src/common/submission.rs |
Calculates and exposes the ready boundary. |
libs/opsqueue_python/src/producer.rs |
Implements progressive streaming. |
libs/opsqueue_python/src/errors.rs |
Adds debug support for failure errors. |
libs/opsqueue_python/python/opsqueue/producer.py |
Exposes the new Python APIs. |
libs/opsqueue_python/tests/test_roundtrip.py |
Tests streaming behavior and validation. |
libs/opsqueue_python/tests/conftest.py |
Adds oldest-strategy fixtures. |
Cargo.toml |
Bumps the workspace version. |
Cargo.lock |
Updates locked package versions. |
.vscode/opsqueue.code-workspace |
Adds workspace configuration. |
🧠 Review effort: Balanced
💡 Add a code-review agent skill or configure MCP servers for context-aware, tailored reviews. Learn more in the docs.
Cool! If you have concerns or want to discuss the approach, feel free to send me message. 👍 |
ReinierMaas
left a comment
There was a problem hiding this comment.
I didn't get through all the Rust code but I think we can ensure that we always pick-up the lowest available chunk from the submission and thereby drop the constraint that this only works on a specific strategy. This would also ensure we fail submissions for which a single chunk has error quicker, i.e. we will pick-up the lowest chunk_index thus retry the available chunk that was tried before and if that hits the retry count we will mark the submission as failing.
There was a problem hiding this comment.
This folder should be .gitignored.
| submission_id: SubmissionId, | ||
| strategy: Strategy, | ||
| ) -> AsyncIterator[bytes]: | ||
| """Stream chunks progressively; strategy must be Oldest for this submission.""" |
There was a problem hiding this comment.
Do we need a strategy if we only accept Oldest? Can't we simply document that we will use Oldest as that is the only one that works with this endpoint?
Or is it about the inner most strategy needing to be Oldest, then the documentation needs a small change.
| "SELECT * FROM chunks WHERE {ffi_is_not_reserved} ORDER BY submission_id ASC, chunk_index ASC" | ||
| )), | ||
| Newest => qb.push(format!( | ||
| "SELECT * FROM chunks WHERE {ffi_is_not_reserved} ORDER BY submission_id DESC" |
There was a problem hiding this comment.
We can also deterministically pickup the chunks of the newest submission in lowest chunk_index order. Then that strategy would also work with the new endpoints.
| Newest => qb.push(format!( | ||
| "SELECT * FROM chunks WHERE {ffi_is_not_reserved} ORDER BY submission_id DESC" | ||
| )), | ||
| Random => Self::push_random_order_query(qb, "*", "chunks", Some(ffi_is_not_reserved)), |
There was a problem hiding this comment.
I am questioning whether we can switch to random order for submissions and also pickup the lowest chunk_index there. It think that is possible as well!
| // In SQLite, <foo> CROSS JOIN <bar> ON/WHERE does NOT produce N | ||
| // x M rows, it acts as an INNER JOIN but forces the query | ||
| // planner to use '<foo>' as the outer loop, preserving the | ||
| // underlying sort order. | ||
| // c.f. https://sqlite.org/optoverview.html#manual_control_of_query_plans_using_cross_join | ||
| qb.push(format!( | ||
| " SELECT chunks.* | ||
| FROM underlying_submission_ids | ||
| CROSS JOIN chunks | ||
| ON chunks.submission_id = underlying_submission_ids.submission_id | ||
| AND {ffi_is_not_reserved}", | ||
| )) |
There was a problem hiding this comment.
We should really investigate whether this query would support something along the lines of CROSS JOIN (SELECT * FROM chunks ORDER BY chunk_index ASC) without changing the order returned by the underlying_submissions_ids.


Description
This pull request implements streaming completed chunks before the whole submission is completed. The reason for adding this is that submissions can be large, and waiting for the whole submission to finish can take long. By streaming completed chunks, producers can get results before the whole submission is completed.
The completed chunks are detected by finding lowest chunk that is not ready yet, using
SELECT MIN(chunk_index) FROM chunks WHERE submission_id = submissions.id). This means that one constraint is added to make streaming useful, which is that the strategy for the producer and consumer should beoldest, notrandomornewest. By querying the lowest chunk that is not ready yet, the lookup stays performant without needing a scan over all chunks to find which are ready.