Implement pausing for submissions - #132
Conversation
a9f5556 to
bcd4161
Compare
|
TODO: How to deal with reserved chunks when pausing? |
b09f93b to
7e7052d
Compare
Discussed this with @ReinierMaas yesterday, we tried to make it work s.t. chunk completions for cancelled and paused submissions would still be processed. However, implementing that proved to be difficult:
After struggling with it for a while, I thought it would be easier to invoke the idempotency assumption we have for how consumers process chunks, and just ignore the completion in case the submission is cancelled or paused. This results in less edge-cases and simpler code in the end, at the cost of slightly higher compute cost. But I think it's worth the trade-off: there are enough edge-cases in the codebase as it is :). @ReinierMaas, let me know if you disagree ;). |
990a5d2 to
71ace4e
Compare
| pass | ||
|
|
||
|
|
||
| class ChunkNotFoundError(IncorrectUsageError): |
There was a problem hiding this comment.
Note that ChunkNotFoundError can be dropped since it is (and was) dead code.
| ) | ||
| .await | ||
| }) | ||
| .await; |
There was a problem hiding this comment.
Note that we were silently dropping errors here before, which went unnoticed because of the underscore prefix of _chunk_size.
e39c8a9 to
82397bd
Compare
Thought some more: if the assumption is that a paused job should not have any work being done, this will violate that. Maybe easier to just only allow pausing of not yet started submissions for now? Other option is to make pausing async, and only transition from |
After sleeping on this for a night, I'm leaning towards ignoring the problem, getting this PR merged, and opening a new ticket to resolve this in the future. Otherwise this PR will become too big IMO (and it is already quite hefty). Will discuss with @ReinierMaas. |
|
I thought a bit on it. I think the easiest way is to only allow a DAG: We can insert a submissions in either From From
If we need more complicated logic we can pick that up as a separate issue. To be abundantly clear how the DAG state transitions would look like: graph LR
%% Active States
P((Paused))
R((Running))
%% Completed States
S([Success])
F([Failed])
C([Cancelled])
%% Internal Active Transitions
P -- unpause --> R
%% Transitions to Completed States
P -- no chunks --> S
P -.-> F
P -- user cancellation --> C
R -- no chunks / all succeed --> S
R -- chunk retries exceeded --> F
R -- user cancellation --> C
|
I think this is impossible, it would always go via
This would indeed make it simpler. In practice, that would amount to removing the Note that this does mean that we will have to implement
But will double-check that with @radekchannable as well. |
Double-checked, we indeed don't have to pause anything when JM asks. |
This was wrong 😅, you can cancel a paused submission, then it would end up in failed. |
3161c6a to
9fc1444
Compare
| #[pyo3(signature = (chunk_contents, metadata=None, strategic_metadata=None, chunk_size=None, otel_trace_carrier=CarrierMap::default()))] | ||
| #[allow(clippy::result_large_err, clippy::type_complexity)] | ||
| /// Submit chunks and then stream the completed output chunks. | ||
| /// | ||
| /// # Errors | ||
| /// | ||
| /// Returns an error if upload, submission creation, or streaming fails. | ||
| pub fn run_submission_chunks( | ||
| &self, | ||
| py: Python<'_>, | ||
| chunk_contents: Py<PyIterator>, | ||
| metadata: Option<submission::Metadata>, | ||
| strategic_metadata: Option<StrategicMetadataMap>, | ||
| chunk_size: Option<i64>, | ||
| otel_trace_carrier: CarrierMap, | ||
| ) -> CPyResult< | ||
| PyChunksIter, | ||
| E![ | ||
| FatalPythonException, | ||
| errors::SubmissionFailed, | ||
| ChunksStorageError, | ||
| InternalProducerClientError, | ||
| ], | ||
| > { | ||
| let submission_id = self | ||
| .insert_submission_chunks( | ||
| py, | ||
| chunk_contents, | ||
| metadata, | ||
| strategic_metadata, | ||
| chunk_size, | ||
| otel_trace_carrier, | ||
| ) | ||
| .map_err(|CError(e)| { | ||
| CError(match e { | ||
| L(e) => L(e), | ||
| R(e) => R(R(e)), | ||
| }) | ||
| })?; | ||
| let res = self | ||
| .blocking_stream_completed_submission_chunks(py, submission_id) | ||
| .map_err(|CError(e)| { | ||
| CError(match e { | ||
| L(e) => L(e), | ||
| R(L(e)) => R(L(e)), | ||
| R(R(e)) => R(R(R(e))), | ||
| }) | ||
| })?; | ||
| Ok(res) | ||
| } | ||
|
|
There was a problem hiding this comment.
This was dead code, because it is re-implemented on the Python side.
There was a problem hiding this comment.
Shouldn't we then prefer the Rust side over the Python side as that is exposed to all downstream consumers.
There was a problem hiding this comment.
Note that this is in opsqueue_python, so it's not exposed to any downstream consumers ATM :).
There was a problem hiding this comment.
Pull request overview
Implements first-class “paused submissions” in opsqueue by storing paused submissions/chunks in separate tables, exposing unpause functionality via the producer API, and updating Rust + Python clients and tests to support pause/unpause and related status reporting.
Changes:
- Add
submissions_paused/chunks_pausedtables and extend submission status to includePaused. - Add producer endpoint + clients for unpausing, and adjust producer insertion to optionally start paused.
- Update metrics, error types, Python bindings, and roundtrip tests (including configurable timeouts).
Reviewed changes
Copilot reviewed 18 out of 19 changed files in this pull request and generated 2 comments.
Show a summary per file
| File | Description |
|---|---|
| opsqueue/src/prometheus.rs | Adds paused/unpaused counters and updates chunk backlog count typing. |
| opsqueue/src/producer/server.rs | Adds /submissions/unpause/{submission_id} endpoint and avoids notifying consumers for paused inserts. |
| opsqueue/src/producer/common.rs | Extends insert request payload with paused: bool. |
| opsqueue/src/producer/client.rs | Adds Rust producer client unpause_submission + updates tests for new count types/status. |
| opsqueue/src/consumer/strategy.rs | Updates tests to pass new paused parameter on submission creation. |
| opsqueue/src/consumer/client.rs | Updates tests to pass new paused parameter on submission creation. |
| opsqueue/src/common/submission.rs | Core implementation: paused submission type/status, insert paused, unpause/cancel paused, updated count APIs. |
| opsqueue/src/common/errors.rs | Removes ChunkNotFound error type. |
| opsqueue/src/common/chunk.rs | Adds paused chunk storage ops; changes chunk completion/retry logic and count types. |
| opsqueue/migrations/20260715143000_pausing.up.sql | Introduces submissions_paused and chunks_paused schema. |
| opsqueue/migrations/20260715143000_pausing.down.sql | Drops paused tables on rollback. |
| libs/opsqueue_python/tests/test_roundtrip.py | Adds timeouts to blocking calls + new tests for timeout and pause/unpause behavior. |
| libs/opsqueue_python/src/producer.rs | Exposes pause/unpause and adds timeout support for blocking streaming methods. |
| libs/opsqueue_python/src/lib.rs | Exports SubmissionPaused to Python module. |
| libs/opsqueue_python/src/errors.rs | Removes chunk-not-found mapping; maps tokio::time::Elapsed to Python TimeoutError. |
| libs/opsqueue_python/src/common.rs | Adds Python-side SubmissionStatus.Paused and SubmissionPaused model. |
| libs/opsqueue_python/python/opsqueue/producer.py | Adds paused + timeout parameters and exposes unpause_submission. |
| libs/opsqueue_python/python/opsqueue/exceptions.py | Removes ChunkNotFoundError exception. |
Comments suppressed due to low confidence (2)
opsqueue/src/common/submission.rs:559
- With the counter increment moved to
unpause_submission(after the transaction commits),unpause_submission_rawshould not also incrementSUBMISSIONS_UNPAUSED_COUNTER, otherwise successful unpauses will be double-counted.
Err(E::R(SubmissionNotFound(id)))
} else {
counter!(crate::prometheus::SUBMISSIONS_UNPAUSED_COUNTER).increment(1);
Ok(())
}
opsqueue/src/common/chunk.rs:318
complete_chunknow bubbles upSubmissionNotFound, which is expected when a submission is cancelled while a chunk is reserved (submission is no longer insubmissions). This turns an expected race into an error (and ends up logged as an error by the completer). Consider treatingSubmissionNotFoundas a non-fatal no-op here, while still propagating real database errors.
.await?;
counter!(crate::prometheus::CHUNKS_COMPLETED_COUNTER).increment(1);
Ok(())
}
💡 Add Copilot custom instructions for smarter, more guided reviews. Learn how to get started.
| conn.transaction(move |mut tx| { | ||
| Box::pin(async move { | ||
| unpause_submission_raw(id, &mut tx).await?; | ||
| super::chunk::db::restore_paused_chunks(id, &mut tx).await?; | ||
| Ok(()) | ||
| }) | ||
| }) | ||
| .await | ||
| } |
There was a problem hiding this comment.
This pattern is there for all of the other metrics as well. I would suggest going for consistency in this PR, and opening an issue for solving it separately.
| counter!(crate::prometheus::SUBMISSIONS_PAUSED_COUNTER).increment(1); | ||
| counter!(crate::prometheus::SUBMISSIONS_TOTAL_COUNTER).increment(1); | ||
| counter!(crate::prometheus::CHUNKS_TOTAL_COUNTER).increment(chunks_total); | ||
| res |
ReinierMaas
left a comment
There was a problem hiding this comment.
I think we are nearly there. I did find some changes we should still make but they are minor and not architectural.
| with pytest.raises(TimeoutError): | ||
| producer_client.run_submission( | ||
| [1], | ||
| chunk_size=1, | ||
| timeout=0.1, | ||
| ) |
There was a problem hiding this comment.
We always check something on the exception raised by pytest, i.e. from a bit higher in the file:
# We expect the intended attributes to be there:
assert isinstance(exc_info.value.failure, str)
assert isinstance(exc_info.value.submission, SubmissionFailed)
assert isinstance(exc_info.value.chunk, ChunkFailed)Just checking the instance of exc_info would be enough here, we have experienced pytest catching the wrong exceptions and giving the green checkmark incidentally
| with pytest.raises(TimeoutError): | |
| producer_client.run_submission( | |
| [1], | |
| chunk_size=1, | |
| timeout=0.1, | |
| ) | |
| with pytest.raises(TimeoutError) as exc_info: | |
| producer_client.run_submission( | |
| [1], | |
| chunk_size=1, | |
| timeout=0.1, | |
| ) | |
| assert isinstance(exc_info, TimeoutError) |
There was a problem hiding this comment.
Hmm, as far as I understand pytest.raises(TimeoutError) already does isinstance(exc_info, TimeoutError)?
As in the following should fail:
with pytest.raised(NotTheRaisedError):
raise AnotherError()There was a problem hiding this comment.
It does:
>>> import pytest
>>> with pytest.raises(TimeoutError):
... raise Exception("Not timeout error")
...
Traceback (most recent call last):
File "<stdin>", line 2, in <module>
Exception: Not timeout error
>>> with pytest.raises(TimeoutError):
... raise TimeoutError()
...
>>> | with pytest.raises(SubmissionNotFoundError): | ||
| producer_client.unpause_submission(submission_id) |
There was a problem hiding this comment.
Same pytest.raises remark here.
c9fc39e to
622b0fd
Compare
u63 forces us to wrap literals in `u63::new`, and we need to convert to u64 at actual usage sites anyway.
…ly completed, failed, or cancelled chunks Because of the idempotency assumption for processing chunks, nothing should break if we just ignore the error. Besides, we were already ignoring the error accidentally.
Introduce `submissions_paused` and `chunks_paused` tables
(alongside the existing `submissions_{completed,failed,cancelled}` and
`chunks_{completed,failed}` tables).
A submission can now be created in a Paused state. It's then stored in
`submissions_paused` and its chunks are stored in `chunks_paused`.
Because paused chunks are not in the `chunks` table, the consumer
dispatcher naturally skips them without any changes to the dispatch
query.
Unpausing moves the submission and the chunks to `submissions` and
`chunks` and notifies waiting consumers.
Paused submissions are cancellable; `cancel_submission` now handles
the case where the submission is found in `submissions_paused`.
We don't allow pausing submissions after creation. That proved to
have too many edge cases we would need to resolve.
622b0fd to
22c5d1e
Compare
There was a problem hiding this comment.
Pull request overview
Copilot reviewed 18 out of 19 changed files in this pull request and generated no new comments.
Suppressed comments (3)
opsqueue/src/common/submission.rs:810
submission_statuscan transiently returnOk(None)during an unpause: if the unpause transaction commits after thesubmissionslookup (line 809) but before the finalsubmissions_pausedlookup, both queries will see “not found” and the function returnsNoneeven though the submission exists (now InProgress). To avoid this race, querysubmissions_pausedbeforesubmissions(or run the lookups inside a single read transaction / single SQL statement so you get a consistent snapshot).
// NOTE: The order is important here; a concurrent writer could move a submission
// from InProgress to Completed/Failed in-between the queries.
let submission_row = submission_status_in_progress_query(id)
.fetch_optional(conn.get_inner())
opsqueue/src/common/submission.rs:530
- Unpausing a zero-chunk paused submission will move it into
submissionswithchunks_total = 0, but nothing will ever complete it (no chunks will be processed andmaybe_complete_submissionis never called on unpause). This leaves the submission stuck InProgress forever. Consider completing the submission during unpause whenchunks_done == chunks_total(which covers the zero-chunk case).
Box::pin(async move {
unpause_submission_raw(id, &mut tx).await?;
super::chunk::db::restore_paused_chunks(id, &mut tx).await?;
Ok(())
opsqueue/src/common/chunk.rs:311
complete_chunkalways callsmaybe_complete_submissioneven whencomplete_chunk_rawdidn’t actually move a chunk (e.g. the chunk was already completed/failed/cancelled). In those cases the submission may already be gone fromsubmissions(cancelled submissions are moved out), somaybe_complete_submissioncan returnSubmissionNotFoundand turn an “ignored duplicate completion” into a hard error. A robust fix is to havecomplete_chunk_rawreturn whether it moved a row and only callmaybe_complete_submissionwhen it did.
crate::common::submission::db::maybe_complete_submission(
chunk_id.submission_id,
&mut tx,
)
.await
417143a to
9c58c87
Compare
There was a problem hiding this comment.
Pull request overview
Copilot reviewed 18 out of 19 changed files in this pull request and generated no new comments.
Suppressed comments (1)
opsqueue/src/common/chunk.rs:326
complete_chunknow returnsOk(())even when the chunk row wasn’t moved (e.g., the chunk was already completed/failed/cancelled). In that case,CHUNKS_COMPLETED_COUNTERis still incremented unconditionally, which can overcount completions (notably when the function is called twice for the same chunk). Increment the counter only when the chunk was actually moved.
.await?;
counter!(crate::prometheus::CHUNKS_COMPLETED_COUNTER).increment(1);
Ok(())
9c58c87 to
0abcc15
Compare
There was a problem hiding this comment.
Pull request overview
Copilot reviewed 18 out of 19 changed files in this pull request and generated 1 comment.
Suppressed comments (1)
libs/opsqueue_python/tests/test_roundtrip.py:756
- This call can block indefinitely if the consumer fails to start or the submission never completes. Since this test suite now standardizes on
SUBMISSION_COMPLETED_TIMEOUT, pass it here as well to avoid hanging CI on failures.
with background_process(run_consumer):
producer_client.blocking_stream_completed_submission(submission_id)
assert isinstance(
0abcc15 to
3ee2454
Compare
Introduce
submissions_pausedandchunks_pausedtables(alongside the existing
submissions_{completed,failed,cancelled}andchunks_{completed,failed}tables).A submission can now be created in a Paused state. It's then stored in
submissions_pausedand its chunks are stored inchunks_paused.Because paused chunks are not in the
chunkstable, the consumerdispatcher naturally skips them without any changes to the dispatch
query.
Unpausing moves the submission and the chunks to
submissionsandchunksand notifies waiting consumers.Paused submissions are cancellable;
cancel_submissionnow handlesthe case where the submission is found in
submissions_paused.We don't allow pausing submissions after creation. That proved to
have too many edge cases we would need to resolve.
Also cleans up some code encountered along the way.