Skip to content

Feat: stream completed chunks - #157

Open
vincevannoort wants to merge 25 commits into
masterfrom
vince/agents/currently-results-can-only-be-obtained-once-the-44bc27f2
Open

vincevannoort wants to merge 25 commits into
masterfrom
vince/agents/currently-results-can-only-be-obtained-once-the-44bc27f2

Conversation

@vincevannoort

@vincevannoort vincevannoort commented Jul 29, 2026 •

Copy link
Copy Markdown
Contributor

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 be oldest, not random or newest. 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.

@vincevannoort
vincevannoort force-pushed the vince/agents/currently-results-can-only-be-obtained-once-the-44bc27f2 branch from 6e49065 to b7df6a4 Compare July 29, 2026 05:41
@vincevannoort
vincevannoort force-pushed the vince/agents/currently-results-can-only-be-obtained-once-the-44bc27f2 branch from b7df6a4 to 8e48718 Compare October 9, 2026 04:47
@vincevannoort vincevannoort self-assigned this Oct 9, 2026
@vincevannoort
vincevannoort marked this pull request as ready for review October 9, 2026 09:58
@vincevannoort

Copy link
Copy Markdown
Contributor Author

@ReinierMaas or @jerbaroo is one of you available to review this somewhere in the coming week(s)?

@jerbaroo

jerbaroo commented Oct 9, 2026

Copy link
Copy Markdown
Contributor

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

@jerbaroo
jerbaroo requested review from jerbaroo and a balanced review from Copilot October 9, 2026 10:30

Copilot AI left a comment

Copy link
Copy Markdown

Choose a reason for hiding this comment

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

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

Comment thread libs/opsqueue_python/src/producer.rs
Comment thread libs/opsqueue_python/src/producer.rs Outdated
Comment thread libs/opsqueue_python/src/producer.rs Outdated
Comment thread libs/opsqueue_python/src/producer.rs Outdated
@vincevannoort

Copy link
Copy Markdown
Contributor Author

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

Cool! If you have concerns or want to discuss the approach, feel free to send me message. 👍

@ReinierMaas ReinierMaas left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

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

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.

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

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

This folder should be .gitignored.

submission_id: SubmissionId,
strategy: Strategy,
) -> AsyncIterator[bytes]:
"""Stream chunks progressively; strategy must be Oldest for this submission."""

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

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

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"

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

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

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

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

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

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!

Comment on lines 112 to 123
// 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}",
))

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

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

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.

@ReinierMaas
ReinierMaas self-requested a review October 9, 2026 14:37
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.

4 participants