From 728d67359eca2d1bc9172659690a8f8cfa1517d7 Mon Sep 17 00:00:00 2001 From: Vince van Noort Date: Wed, 29 Jul 2026 07:26:42 +0200 Subject: [PATCH 01/25] Feat: stream completed chunks --- .vscode/opsqueue.code-workspace | 13 ++ .../python/opsqueue/producer.py | 8 + libs/opsqueue_python/src/errors.rs | 1 + libs/opsqueue_python/src/producer.rs | 216 +++++++++++++++++- libs/opsqueue_python/tests/test_roundtrip.py | 96 +++++++- 5 files changed, 324 insertions(+), 10 deletions(-) create mode 100644 .vscode/opsqueue.code-workspace diff --git a/.vscode/opsqueue.code-workspace b/.vscode/opsqueue.code-workspace new file mode 100644 index 00000000..dc0b7ae7 --- /dev/null +++ b/.vscode/opsqueue.code-workspace @@ -0,0 +1,13 @@ +{ + "folders": [ + { + "path": "..", + "name": "opsqueue", + }, + { + "path": "../libs/opsqueue_python", + "name": "opsqueue_python", + } + ], + "settings": {} +} diff --git a/libs/opsqueue_python/python/opsqueue/producer.py b/libs/opsqueue_python/python/opsqueue/producer.py index bdb6b593..e2aee6d0 100644 --- a/libs/opsqueue_python/python/opsqueue/producer.py +++ b/libs/opsqueue_python/python/opsqueue/producer.py @@ -240,6 +240,14 @@ def run_submission_chunks( ) return self.blocking_stream_completed_submission_chunks(submission_id, timeout) + def stream_submission_chunks(self, submission_id: SubmissionId) -> Iterator[bytes]: + return self.inner.stream_submission_chunks(submission_id) # type: ignore[no-any-return] + + async def async_stream_submission_chunks( + self, submission_id: SubmissionId + ) -> AsyncIterator[bytes]: + return await self.inner.async_stream_submission_chunks(submission_id) # type: ignore[no-any-return] + async def async_run_submission_chunks( self, chunk_contents: Iterable[bytes], diff --git a/libs/opsqueue_python/src/errors.rs b/libs/opsqueue_python/src/errors.rs index 043ac11f..4681d1c7 100644 --- a/libs/opsqueue_python/src/errors.rs +++ b/libs/opsqueue_python/src/errors.rs @@ -151,6 +151,7 @@ impl From> for PyErr { } } +#[derive(Debug)] pub struct SubmissionFailed( pub crate::common::SubmissionFailed, pub crate::common::ChunkFailed, diff --git a/libs/opsqueue_python/src/producer.rs b/libs/opsqueue_python/src/producer.rs index 5aed787f..84e1e812 100644 --- a/libs/opsqueue_python/src/producer.rs +++ b/libs/opsqueue_python/src/producer.rs @@ -421,6 +421,146 @@ impl ProducerClient { }) } + /// Stream output chunks as soon as each consumer has completed them. + pub fn stream_submission_chunks(&self, submission_id: SubmissionId) -> PyChunksIter { + self.streaming_submission_chunks(submission_id) + } + + fn streaming_submission_chunks(&self, submission_id: SubmissionId) -> PyChunksIter { + let client = self.client.clone(); + let object_store_client = self.object_store_client.clone(); + let stream = futures::stream::unfold( + ( + client, + object_store_client, + submission_id, + u63::new(0), + None, + Duration::from_millis(10), + ), + |(client, object_store_client, submission_id, index, prefix, interval)| async move { + let mut interval = interval; + loop { + let status = match client.get_submission(submission_id.into()).await { + Ok(Some(status)) => status, + Ok(None) => { + return Some(( + Err(StreamingChunkError::SubmissionNotFound), + ( + client, + object_store_client, + submission_id, + index, + prefix, + interval, + ), + )); + } + Err(error) => { + return Some(( + Err(StreamingChunkError::Internal(error)), + ( + client, + object_store_client, + submission_id, + index, + prefix, + interval, + ), + )); + } + }; + + match status { + submission::SubmissionStatus::InProgress(submission) => { + let prefix = prefix.clone().or(submission.prefix); + if index < submission.chunks_done.into() { + let prefix = prefix + .expect("in-progress submissions have an object-store prefix"); + let result = object_store_client + .retrieve_chunk(&prefix, index.into(), ChunkType::Output) + .await + .map_err(StreamingChunkError::Retrieval); + return Some(( + result, + ( + client, + object_store_client, + submission_id, + index + u63::new(1), + Some(prefix), + interval, + ), + )); + } + } + submission::SubmissionStatus::Completed(submission) => { + let prefix = prefix.clone().or(submission.prefix); + if index < submission.chunks_total.into() { + let prefix = prefix + .expect("completed submissions have an object-store prefix"); + let result = object_store_client + .retrieve_chunk(&prefix, index.into(), ChunkType::Output) + .await + .map_err(StreamingChunkError::Retrieval); + return Some(( + result, + ( + client, + object_store_client, + submission_id, + index + u63::new(1), + Some(prefix), + interval, + ), + )); + } + return None; + } + submission::SubmissionStatus::Failed(submission, chunk) => { + let failure = + crate::common::ChunkFailed::from_internal(chunk, &submission); + return Some(( + Err(StreamingChunkError::Failed( + crate::errors::SubmissionFailed(submission.into(), failure), + )), + ( + client, + object_store_client, + submission_id, + index, + prefix, + interval, + ), + )); + } + submission::SubmissionStatus::Cancelled(_) => { + return Some(( + Err(StreamingChunkError::Cancelled), + ( + client, + object_store_client, + submission_id, + index, + prefix, + interval, + ), + )); + } + } + + tokio::time::sleep(interval).await; + if interval < SUBMISSION_POLLING_INTERVAL { + interval = (interval * 2).min(SUBMISSION_POLLING_INTERVAL); + } + } + }, + ) + .map(|item| item.map_err(CError)) + .boxed(); + PyChunksIter::from_stream(self, stream) + } + /// Blocks (and short-polls) until the submission is completed. /// /// We start with a small short-polling interval @@ -474,6 +614,30 @@ impl ProducerClient { }) } + /// Return an awaitable that resolves immediately to an async iterator of output chunks. + /// + /// The iterator polls submission progress and yields each output chunk as soon as it is ready. + /// + /// # Errors + /// + /// Returns a Python error if creating the awaitable fails. + pub fn async_stream_submission_chunks<'p>( + &self, + py: Python<'p>, + submission_id: SubmissionId, + ) -> PyResult> { + let me = self.clone(); + let _tokio_active_runtime_guard = me.runtime.enter(); + async_util::future_into_py( + py, + async_util::async_detach(Box::pin(async move { + Ok(PyChunksAsyncIter::from( + me.streaming_submission_chunks(submission_id), + )) + })), + ) + } + /// Return an awaitable that resolves to an async iterator of output chunks. /// /// # Errors @@ -583,7 +747,41 @@ impl ProducerClient { } } -pub type ChunksStream = BoxStream<'static, CPyResult, ChunkRetrievalError>>; +#[derive(Debug)] +enum StreamingChunkError { + Retrieval(ChunkRetrievalError), + Internal(InternalProducerClientError), + Failed(crate::errors::SubmissionFailed), + SubmissionNotFound, + Cancelled, +} + +impl std::fmt::Display for StreamingChunkError { + fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { + match self { + Self::Retrieval(error) => error.fmt(f), + Self::Internal(error) => error.fmt(f), + Self::Failed(_) => write!(f, "Submission failed"), + Self::SubmissionNotFound => write!(f, "Submission not found"), + Self::Cancelled => write!(f, "Submission cancelled"), + } + } +} + +impl std::error::Error for StreamingChunkError {} + +impl From> for PyErr { + fn from(value: CError) -> Self { + match value.0 { + StreamingChunkError::Retrieval(error) => CError(error).into(), + StreamingChunkError::Internal(error) => CError(error).into(), + StreamingChunkError::Failed(error) => CError(error).into(), + error => PyException::new_err(error.to_string()), + } + } +} + +type ChunksStream = BoxStream<'static, CPyResult, StreamingChunkError>>; #[pyclass(module = "opsqueue")] pub struct PyChunksIter { @@ -592,16 +790,20 @@ pub struct PyChunksIter { } impl PyChunksIter { + fn from_stream(client: &ProducerClient, stream: ChunksStream) -> Self { + Self { + stream: Arc::new(tokio::sync::Mutex::new(stream)), + runtime: client.runtime.clone(), + } + } + pub(crate) fn new(client: &ProducerClient, prefix: String, chunks_total: u63) -> Self { let stream = client .object_store_client .retrieve_chunks(prefix, chunks_total, ChunkType::Output) - .map_err(CError) + .map_err(|error| CError(StreamingChunkError::Retrieval(error))) .boxed(); - Self { - stream: Arc::new(tokio::sync::Mutex::new(stream)), - runtime: client.runtime.clone(), - } + Self::from_stream(client, stream) } } @@ -611,7 +813,7 @@ impl PyChunksIter { slf } - fn __next__(&self, py: Python<'_>) -> Option, ChunkRetrievalError>> { + fn __next__(&self, py: Python<'_>) -> Option, StreamingChunkError>> { // The only time we need the GIL is when turning the result back. // By unlocking here, we reduce the chance of deadlocks. py.detach(move || { diff --git a/libs/opsqueue_python/tests/test_roundtrip.py b/libs/opsqueue_python/tests/test_roundtrip.py index ecfd0064..990d4220 100644 --- a/libs/opsqueue_python/tests/test_roundtrip.py +++ b/libs/opsqueue_python/tests/test_roundtrip.py @@ -3,8 +3,8 @@ # - use `RUST_LOG="opsqueue=info"` (or `opsqueue=debug` or `debug` for even more verbosity), together with to the pytest option `-s` AKA `--capture=no`, to debug the opsqueue binary itself. import logging -import time from collections.abc import Iterator, Sequence +import asyncio import pytest from conftest import ( @@ -29,7 +29,6 @@ SubmissionNotCancellable, SubmissionNotCancellableError, TooManyMatchingSubmissionsError, - InitialSubmissionStatus, ) SUBMISSION_COMPLETED_TIMEOUT = 10.0 @@ -338,7 +337,6 @@ def test_async_producer(opsqueue: OpsqueueProcess) -> None: """ A simple sanity check to ensure the async API does its basic job """ - import asyncio def run_consumer() -> None: def increment_list(ints: Sequence[int], _chunk: Chunk) -> Sequence[int]: @@ -824,3 +822,95 @@ def test_cancel_paused(opsqueue: OpsqueueProcess) -> None: assert isinstance( producer_client.get_submission_status(submission_id), SubmissionStatus.Cancelled ) + + +def test_streams_completed_chunks_before_submission_finishes( + opsqueue: OpsqueueProcess, + any_consumer_strategy: StrategyDescription, +) -> None: + url = "file:///tmp/opsqueue/test_streaming_results" + producer_client = ProducerClient(f"localhost:{opsqueue.port}", url) + submission_id = producer_client.insert_submission_chunks( + [b"[1]", b"[2]"], chunk_size=1 + ) + + def complete_chunks( + _submission_id_value: int, + strategy: StrategyDescription, + ) -> None: + consumer_client = ConsumerClient(f"localhost:{opsqueue.port}", url) + chunks = sorted( + consumer_client.reserve_chunks( + max=2, + strategy=strategy_from_description(strategy), + ), + key=lambda chunk: chunk.chunk_index, + ) + consumer_client.complete_chunk( + chunks[0].submission_id, + chunks[0].submission_prefix, + chunks[0].chunk_index, + chunks[0].input_content, + ) + time.sleep(0.25) + consumer_client.complete_chunk( + chunks[1].submission_id, + chunks[1].submission_prefix, + chunks[1].chunk_index, + chunks[1].input_content, + ) + + with background_process( + complete_chunks, + args=(submission_id.id, any_consumer_strategy), + ): + results = producer_client.stream_submission_chunks(submission_id) + assert next(results) == b"[1]" + assert next(results) == b"[2]" + + +def test_async_streams_completed_chunks_before_submission_finishes( + opsqueue: OpsqueueProcess, + any_consumer_strategy: StrategyDescription, +) -> None: + url = "file:///tmp/opsqueue/test_async_streaming_results" + producer_client = ProducerClient(f"localhost:{opsqueue.port}", url) + submission_id = producer_client.insert_submission_chunks( + [b"[1]", b"[2]"], chunk_size=1 + ) + + def complete_chunks( + _submission_id_value: int, + strategy: StrategyDescription, + ) -> None: + consumer_client = ConsumerClient(f"localhost:{opsqueue.port}", url) + chunks = sorted( + consumer_client.reserve_chunks( + max=2, + strategy=strategy_from_description(strategy), + ), + key=lambda chunk: chunk.chunk_index, + ) + consumer_client.complete_chunk( + chunks[0].submission_id, + chunks[0].submission_prefix, + chunks[0].chunk_index, + chunks[0].input_content, + ) + time.sleep(0.25) + consumer_client.complete_chunk( + chunks[1].submission_id, + chunks[1].submission_prefix, + chunks[1].chunk_index, + chunks[1].input_content, + ) + + async def collect() -> list[bytes]: + results = await producer_client.async_stream_submission_chunks(submission_id) + return [chunk async for chunk in results] + + with background_process( + complete_chunks, + args=(submission_id.id, any_consumer_strategy), + ): + assert asyncio.run(collect()) == [b"[1]", b"[2]"] From 8e487184f679e05e538ee2b0f0796397a8a3690c Mon Sep 17 00:00:00 2001 From: Vince van Noort Date: Wed, 29 Jul 2026 07:41:03 +0200 Subject: [PATCH 02/25] Fix: address clippy errors --- libs/opsqueue_python/src/producer.rs | 10 ++++++---- 1 file changed, 6 insertions(+), 4 deletions(-) diff --git a/libs/opsqueue_python/src/producer.rs b/libs/opsqueue_python/src/producer.rs index 84e1e812..c10cf2b7 100644 --- a/libs/opsqueue_python/src/producer.rs +++ b/libs/opsqueue_python/src/producer.rs @@ -422,10 +422,12 @@ impl ProducerClient { } /// Stream output chunks as soon as each consumer has completed them. + #[must_use] pub fn stream_submission_chunks(&self, submission_id: SubmissionId) -> PyChunksIter { self.streaming_submission_chunks(submission_id) } + #[allow(clippy::too_many_lines)] fn streaming_submission_chunks(&self, submission_id: SubmissionId) -> PyChunksIter { let client = self.client.clone(); let object_store_client = self.object_store_client.clone(); @@ -521,9 +523,9 @@ impl ProducerClient { let failure = crate::common::ChunkFailed::from_internal(chunk, &submission); return Some(( - Err(StreamingChunkError::Failed( + Err(StreamingChunkError::Failed(Box::new( crate::errors::SubmissionFailed(submission.into(), failure), - )), + ))), ( client, object_store_client, @@ -751,7 +753,7 @@ impl ProducerClient { enum StreamingChunkError { Retrieval(ChunkRetrievalError), Internal(InternalProducerClientError), - Failed(crate::errors::SubmissionFailed), + Failed(Box), SubmissionNotFound, Cancelled, } @@ -775,7 +777,7 @@ impl From> for PyErr { match value.0 { StreamingChunkError::Retrieval(error) => CError(error).into(), StreamingChunkError::Internal(error) => CError(error).into(), - StreamingChunkError::Failed(error) => CError(error).into(), + StreamingChunkError::Failed(error) => CError(*error).into(), error => PyException::new_err(error.to_string()), } } From a7508bb9e15091e5fe5dc81864723cd126a16116 Mon Sep 17 00:00:00 2001 From: Vince van Noort Date: Fri, 9 Oct 2026 06:49:53 +0200 Subject: [PATCH 03/25] Fix: update imports --- libs/opsqueue_python/src/producer.rs | 1 + libs/opsqueue_python/tests/test_roundtrip.py | 10 +++++++++- 2 files changed, 10 insertions(+), 1 deletion(-) diff --git a/libs/opsqueue_python/src/producer.rs b/libs/opsqueue_python/src/producer.rs index c10cf2b7..1bd50552 100644 --- a/libs/opsqueue_python/src/producer.rs +++ b/libs/opsqueue_python/src/producer.rs @@ -536,6 +536,7 @@ impl ProducerClient { ), )); } + submission::SubmissionStatus::Paused(_) => {} submission::SubmissionStatus::Cancelled(_) => { return Some(( Err(StreamingChunkError::Cancelled), diff --git a/libs/opsqueue_python/tests/test_roundtrip.py b/libs/opsqueue_python/tests/test_roundtrip.py index 990d4220..4d1208ab 100644 --- a/libs/opsqueue_python/tests/test_roundtrip.py +++ b/libs/opsqueue_python/tests/test_roundtrip.py @@ -3,8 +3,8 @@ # - use `RUST_LOG="opsqueue=info"` (or `opsqueue=debug` or `debug` for even more verbosity), together with to the pytest option `-s` AKA `--capture=no`, to debug the opsqueue binary itself. import logging +import time from collections.abc import Iterator, Sequence -import asyncio import pytest from conftest import ( @@ -29,6 +29,7 @@ SubmissionNotCancellable, SubmissionNotCancellableError, TooManyMatchingSubmissionsError, + InitialSubmissionStatus, ) SUBMISSION_COMPLETED_TIMEOUT = 10.0 @@ -337,6 +338,7 @@ def test_async_producer(opsqueue: OpsqueueProcess) -> None: """ A simple sanity check to ensure the async API does its basic job """ + import asyncio def run_consumer() -> None: def increment_list(ints: Sequence[int], _chunk: Chunk) -> Sequence[int]: @@ -866,6 +868,10 @@ def complete_chunks( ): results = producer_client.stream_submission_chunks(submission_id) assert next(results) == b"[1]" + assert isinstance( + producer_client.get_submission_status(submission_id), + SubmissionStatus.InProgress, + ) assert next(results) == b"[2]" @@ -873,6 +879,8 @@ def test_async_streams_completed_chunks_before_submission_finishes( opsqueue: OpsqueueProcess, any_consumer_strategy: StrategyDescription, ) -> None: + import asyncio + url = "file:///tmp/opsqueue/test_async_streaming_results" producer_client = ProducerClient(f"localhost:{opsqueue.port}", url) submission_id = producer_client.insert_submission_chunks( From 12d4ce95ab581d58a75926fcc69d8367bc6180ad Mon Sep 17 00:00:00 2001 From: Vince van Noort Date: Fri, 9 Oct 2026 06:50:58 +0200 Subject: [PATCH 04/25] fixup! Fix: update imports --- libs/opsqueue_python/tests/test_roundtrip.py | 2 -- 1 file changed, 2 deletions(-) diff --git a/libs/opsqueue_python/tests/test_roundtrip.py b/libs/opsqueue_python/tests/test_roundtrip.py index 4d1208ab..5e19d18f 100644 --- a/libs/opsqueue_python/tests/test_roundtrip.py +++ b/libs/opsqueue_python/tests/test_roundtrip.py @@ -879,8 +879,6 @@ def test_async_streams_completed_chunks_before_submission_finishes( opsqueue: OpsqueueProcess, any_consumer_strategy: StrategyDescription, ) -> None: - import asyncio - url = "file:///tmp/opsqueue/test_async_streaming_results" producer_client = ProducerClient(f"localhost:{opsqueue.port}", url) submission_id = producer_client.insert_submission_chunks( From f4ff26f6b07743cbff7f6057c5aab9c6ef7cfa8d Mon Sep 17 00:00:00 2001 From: Vince van Noort Date: Fri, 9 Oct 2026 07:11:47 +0200 Subject: [PATCH 05/25] Fix: ensure strategy oldest --- .../python/opsqueue/producer.py | 17 +++- libs/opsqueue_python/src/producer.rs | 39 ++++++--- libs/opsqueue_python/tests/test_roundtrip.py | 79 ++++++++++------- opsqueue/src/common/submission.rs | 9 ++ opsqueue/src/consumer/strategy.rs | 87 ++++++++----------- 5 files changed, 133 insertions(+), 98 deletions(-) diff --git a/libs/opsqueue_python/python/opsqueue/producer.py b/libs/opsqueue_python/python/opsqueue/producer.py index e2aee6d0..08b79c1c 100644 --- a/libs/opsqueue_python/python/opsqueue/producer.py +++ b/libs/opsqueue_python/python/opsqueue/producer.py @@ -28,6 +28,7 @@ SubmissionNotCancellable, SubmissionPaused, InitialSubmissionStatus, + Strategy, ) __all__ = [ @@ -240,13 +241,21 @@ def run_submission_chunks( ) return self.blocking_stream_completed_submission_chunks(submission_id, timeout) - def stream_submission_chunks(self, submission_id: SubmissionId) -> Iterator[bytes]: - return self.inner.stream_submission_chunks(submission_id) # type: ignore[no-any-return] + def stream_submission_chunks( + self, submission_id: SubmissionId, strategy: Strategy + ) -> Iterator[bytes]: + """Stream chunks progressively; strategy must be Oldest for this submission.""" + return self.inner.stream_submission_chunks( # type: ignore[no-any-return] + submission_id, strategy + ) async def async_stream_submission_chunks( - self, submission_id: SubmissionId + self, submission_id: SubmissionId, strategy: Strategy ) -> AsyncIterator[bytes]: - return await self.inner.async_stream_submission_chunks(submission_id) # type: ignore[no-any-return] + """Stream chunks progressively; strategy must be Oldest for this submission.""" + return await self.inner.async_stream_submission_chunks( # type: ignore[no-any-return] + submission_id, strategy + ) async def async_run_submission_chunks( self, diff --git a/libs/opsqueue_python/src/producer.rs b/libs/opsqueue_python/src/producer.rs index 1bd50552..1cd19ad1 100644 --- a/libs/opsqueue_python/src/producer.rs +++ b/libs/opsqueue_python/src/producer.rs @@ -1,6 +1,6 @@ use pyo3::{ create_exception, - exceptions::{PyException, PyStopAsyncIteration}, + exceptions::{PyException, PyStopAsyncIteration, PyValueError}, prelude::*, types::PyIterator, }; @@ -10,7 +10,7 @@ use std::{future::IntoFuture, sync::Arc, time::Duration}; use crate::{ async_util, common::{ - InitialSubmissionStatus, SubmissionId, SubmissionStatus, run_unless_interrupted, + InitialSubmissionStatus, Strategy, SubmissionId, SubmissionStatus, run_unless_interrupted, start_runtime, }, errors::{self, CError, CPyResult, FatalPythonException}, @@ -422,13 +422,30 @@ impl ProducerClient { } /// Stream output chunks as soon as each consumer has completed them. + /// + /// `strategy` must be `Oldest`, and consumers processing this submission must + /// also reserve chunks using `Oldest`. #[must_use] - pub fn stream_submission_chunks(&self, submission_id: SubmissionId) -> PyChunksIter { - self.streaming_submission_chunks(submission_id) + pub fn stream_submission_chunks( + &self, + submission_id: SubmissionId, + strategy: &Strategy, + ) -> PyResult { + self.streaming_submission_chunks(submission_id, strategy) } #[allow(clippy::too_many_lines)] - fn streaming_submission_chunks(&self, submission_id: SubmissionId) -> PyChunksIter { + fn streaming_submission_chunks( + &self, + submission_id: SubmissionId, + strategy: &Strategy, + ) -> PyResult { + if !matches!(strategy, Strategy::Oldest()) { + return Err(PyValueError::new_err( + "streaming submission chunks requires Strategy.Oldest; consumers must also use Strategy.Oldest", + )); + } + let client = self.client.clone(); let object_store_client = self.object_store_client.clone(); let stream = futures::stream::unfold( @@ -476,7 +493,7 @@ impl ProducerClient { match status { submission::SubmissionStatus::InProgress(submission) => { let prefix = prefix.clone().or(submission.prefix); - if index < submission.chunks_done.into() { + if index < submission.chunks_ready.into() { let prefix = prefix .expect("in-progress submissions have an object-store prefix"); let result = object_store_client @@ -561,7 +578,7 @@ impl ProducerClient { ) .map(|item| item.map_err(CError)) .boxed(); - PyChunksIter::from_stream(self, stream) + Ok(PyChunksIter::from_stream(self, stream)) } /// Blocks (and short-polls) until the submission is completed. @@ -628,16 +645,14 @@ impl ProducerClient { &self, py: Python<'p>, submission_id: SubmissionId, + strategy: &Strategy, ) -> PyResult> { let me = self.clone(); + let stream = me.streaming_submission_chunks(submission_id, strategy)?; let _tokio_active_runtime_guard = me.runtime.enter(); async_util::future_into_py( py, - async_util::async_detach(Box::pin(async move { - Ok(PyChunksAsyncIter::from( - me.streaming_submission_chunks(submission_id), - )) - })), + async_util::async_detach(Box::pin(async move { Ok(PyChunksAsyncIter::from(stream)) })), ) } diff --git a/libs/opsqueue_python/tests/test_roundtrip.py b/libs/opsqueue_python/tests/test_roundtrip.py index 5e19d18f..6dd5dd75 100644 --- a/libs/opsqueue_python/tests/test_roundtrip.py +++ b/libs/opsqueue_python/tests/test_roundtrip.py @@ -5,7 +5,7 @@ import logging import time from collections.abc import Iterator, Sequence - +import asyncio import pytest from conftest import ( background_process, @@ -16,7 +16,7 @@ strategy_from_description, ) from opsqueue.common import SerializationFormat -from opsqueue.consumer import ConsumerClient, Chunk +from opsqueue.consumer import ConsumerClient, Chunk, Strategy from opsqueue.producer import ( SubmissionId, ProducerClient, @@ -828,7 +828,6 @@ def test_cancel_paused(opsqueue: OpsqueueProcess) -> None: def test_streams_completed_chunks_before_submission_finishes( opsqueue: OpsqueueProcess, - any_consumer_strategy: StrategyDescription, ) -> None: url = "file:///tmp/opsqueue/test_streaming_results" producer_client = ProducerClient(f"localhost:{opsqueue.port}", url) @@ -836,37 +835,36 @@ def test_streams_completed_chunks_before_submission_finishes( [b"[1]", b"[2]"], chunk_size=1 ) - def complete_chunks( - _submission_id_value: int, - strategy: StrategyDescription, - ) -> None: + def complete_chunks(_submission_id_value: int) -> None: consumer_client = ConsumerClient(f"localhost:{opsqueue.port}", url) chunks = sorted( consumer_client.reserve_chunks( max=2, - strategy=strategy_from_description(strategy), + strategy=Strategy.Oldest(), ), key=lambda chunk: chunk.chunk_index, ) - consumer_client.complete_chunk( - chunks[0].submission_id, - chunks[0].submission_prefix, - chunks[0].chunk_index, - chunks[0].input_content, - ) - time.sleep(0.25) consumer_client.complete_chunk( chunks[1].submission_id, chunks[1].submission_prefix, chunks[1].chunk_index, chunks[1].input_content, ) + time.sleep(0.25) + consumer_client.complete_chunk( + chunks[0].submission_id, + chunks[0].submission_prefix, + chunks[0].chunk_index, + chunks[0].input_content, + ) with background_process( complete_chunks, - args=(submission_id.id, any_consumer_strategy), + args=(submission_id.id,), ): - results = producer_client.stream_submission_chunks(submission_id) + results = producer_client.stream_submission_chunks( + submission_id, Strategy.Oldest() + ) assert next(results) == b"[1]" assert isinstance( producer_client.get_submission_status(submission_id), @@ -877,7 +875,6 @@ def complete_chunks( def test_async_streams_completed_chunks_before_submission_finishes( opsqueue: OpsqueueProcess, - any_consumer_strategy: StrategyDescription, ) -> None: url = "file:///tmp/opsqueue/test_async_streaming_results" producer_client = ProducerClient(f"localhost:{opsqueue.port}", url) @@ -885,38 +882,56 @@ def test_async_streams_completed_chunks_before_submission_finishes( [b"[1]", b"[2]"], chunk_size=1 ) - def complete_chunks( - _submission_id_value: int, - strategy: StrategyDescription, - ) -> None: + def complete_chunks(_submission_id_value: int) -> None: consumer_client = ConsumerClient(f"localhost:{opsqueue.port}", url) chunks = sorted( consumer_client.reserve_chunks( max=2, - strategy=strategy_from_description(strategy), + strategy=Strategy.Oldest(), ), key=lambda chunk: chunk.chunk_index, ) - consumer_client.complete_chunk( - chunks[0].submission_id, - chunks[0].submission_prefix, - chunks[0].chunk_index, - chunks[0].input_content, - ) - time.sleep(0.25) consumer_client.complete_chunk( chunks[1].submission_id, chunks[1].submission_prefix, chunks[1].chunk_index, chunks[1].input_content, ) + time.sleep(0.25) + consumer_client.complete_chunk( + chunks[0].submission_id, + chunks[0].submission_prefix, + chunks[0].chunk_index, + chunks[0].input_content, + ) async def collect() -> list[bytes]: - results = await producer_client.async_stream_submission_chunks(submission_id) + results = await producer_client.async_stream_submission_chunks( + submission_id, Strategy.Oldest() + ) return [chunk async for chunk in results] with background_process( complete_chunks, - args=(submission_id.id, any_consumer_strategy), + args=(submission_id.id,), ): assert asyncio.run(collect()) == [b"[1]", b"[2]"] + + +def test_stream_submission_chunks_requires_oldest_strategy( + opsqueue: OpsqueueProcess, +) -> None: + url = "file:///tmp/opsqueue/test_streaming_requires_oldest" + producer_client = ProducerClient(f"localhost:{opsqueue.port}", url) + submission_id = producer_client.insert_submission_chunks([b"[1]"], chunk_size=1) + + with pytest.raises(ValueError, match="requires Strategy.Oldest"): + producer_client.stream_submission_chunks(submission_id, Strategy.Random()) + + async def collect_with_newest() -> None: + await producer_client.async_stream_submission_chunks( + submission_id, Strategy.Newest() + ) + + with pytest.raises(ValueError, match="requires Strategy.Oldest"): + asyncio.run(collect_with_newest()) diff --git a/opsqueue/src/common/submission.rs b/opsqueue/src/common/submission.rs index a304c30a..64e5af5d 100644 --- a/opsqueue/src/common/submission.rs +++ b/opsqueue/src/common/submission.rs @@ -156,6 +156,7 @@ pub struct Submission { pub prefix: Option, pub chunks_total: ChunkCount, pub chunks_done: ChunkCount, + pub chunks_ready: ChunkCount, pub chunk_size: ChunkSize, pub metadata: Option, #[serde(default)] @@ -262,6 +263,7 @@ impl Submission { prefix: None, chunks_total: ChunkCount::zero(), chunks_done: ChunkCount::zero(), + chunks_ready: ChunkCount::zero(), chunk_size: ChunkSize::default(), metadata: None, strategic_metadata: StrategicMetadataMap::default(), @@ -283,6 +285,7 @@ impl Submission { prefix: None, chunks_total: len, chunks_done: ChunkCount::zero(), + chunks_ready: ChunkCount::zero(), chunk_size, metadata, strategic_metadata: StrategicMetadataMap::default(), @@ -617,6 +620,7 @@ pub mod db { prefix, chunks_total: len, chunks_done: ChunkCount::zero(), + chunks_ready: ChunkCount::zero(), chunk_size, metadata, strategic_metadata, @@ -659,6 +663,7 @@ pub mod db { , prefix , chunks_total AS "chunks_total: ChunkCount" , chunks_done AS "chunks_done: ChunkCount" + , COALESCE((SELECT MIN(chunk_index) FROM chunks WHERE submission_id = submissions.id), chunks_total) AS "chunks_ready: ChunkCount" , chunk_size AS "chunk_size!: ChunkSize" , metadata , ( SELECT json_group_object(metadata_key, metadata_value) @@ -679,6 +684,7 @@ pub mod db { prefix: row.prefix, chunks_total: row.chunks_total, chunks_done: row.chunks_done, + chunks_ready: row.chunks_ready, chunk_size: row.chunk_size, metadata: row.metadata, strategic_metadata: row.strategic_metadata.0, @@ -843,6 +849,7 @@ pub mod db { prefix: row.prefix, chunks_total: row.chunks_total, chunks_done: row.chunks_done, + chunks_ready: row.chunks_ready, chunk_size: row.chunk_size, metadata: row.metadata, strategic_metadata: row.strategic_metadata.0, @@ -916,6 +923,7 @@ pub mod db { prefix: Option, chunks_total: ChunkCount, chunks_done: ChunkCount, + chunks_ready: ChunkCount, chunk_size: ChunkSize, metadata: Option, strategic_metadata: sqlx::types::Json, @@ -939,6 +947,7 @@ pub mod db { , prefix , chunks_total AS "chunks_total: ChunkCount" , chunks_done AS "chunks_done: ChunkCount" + , COALESCE((SELECT MIN(chunk_index) FROM chunks WHERE submission_id = submissions.id), chunks_total) AS "chunks_ready: ChunkCount" , chunk_size AS "chunk_size!: ChunkSize" , metadata , ( SELECT json_group_object(metadata_key, metadata_value) diff --git a/opsqueue/src/consumer/strategy.rs b/opsqueue/src/consumer/strategy.rs index 3785e33b..916200b4 100644 --- a/opsqueue/src/consumer/strategy.rs +++ b/opsqueue/src/consumer/strategy.rs @@ -7,10 +7,14 @@ use sqlx::{QueryBuilder, Sqlite}; #[cfg(feature = "server-logic")] use crate::common::chunk::Chunk; +#[cfg(feature = "server-logic")] +const RANDOM_CHUNK_WINDOW: u64 = 4; + #[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] pub enum Strategy { Oldest, Newest, + /// Randomize chunks within a small, advancing index window per submission. Random, /// Perform a sort of submissions by metadata (ordered by in-flight counts), /// ties are broken by the underlying strategy. For example: if we have two @@ -98,12 +102,29 @@ impl Strategy { "opsqueue_is_reserved(chunks.submission_id, chunks.chunk_index) = FALSE"; match self { Oldest => qb.push(format!( - "SELECT * FROM chunks WHERE {ffi_is_not_reserved} ORDER BY submission_id ASC" - )), - Newest => qb.push(format!( - "SELECT * FROM chunks WHERE {ffi_is_not_reserved} ORDER BY submission_id DESC" + "SELECT * FROM chunks WHERE {ffi_is_not_reserved} ORDER BY submission_id ASC, chunk_index ASC" )), - Random => Self::push_random_order_query(qb, "*", "chunks", Some(ffi_is_not_reserved)), + Newest => { + let qb = qb.push("WITH newest_submission_ids AS MATERIALIZED ("); + let qb = qb.push("SELECT id as submission_id FROM submissions ORDER BY id DESC"); + qb.push(format!( + ") SELECT chunks.* + FROM newest_submission_ids + CROSS JOIN chunks + ON chunks.submission_id = newest_submission_ids.submission_id + AND {ffi_is_not_reserved}" + )) + } + Random => { + let condition = format!( + "{ffi_is_not_reserved} AND chunks.chunk_index < ( + SELECT MIN(previous.chunk_index) + {RANDOM_CHUNK_WINDOW} + FROM chunks AS previous + WHERE previous.submission_id = chunks.submission_id + )" + ); + Self::push_random_order_query(qb, "*", "chunks", Some(&condition)) + } PreferDistinct { .. } => { // Unique submission IDs from the underlying strategy. let qb = qb.push("WITH underlying_submission_ids AS MATERIALIZED ("); @@ -214,7 +235,7 @@ impl Strategy { qb: &'a mut QueryBuilder, columns: &'static str, table_name: &'static str, - condition: Option<&'static str>, + condition: Option<&str>, ) -> &'a mut QueryBuilder { let random_offset: u16 = rand::random(); let push_select = |qb: &mut QueryBuilder, operator: &str| { @@ -379,7 +400,8 @@ pub mod test { WHERE opsqueue_is_reserved(chunks.submission_id, chunks.chunk_index) = FALSE ORDER BY - submission_id ASC + submission_id ASC, + chunk_index ASC "); let explained = explain(qb, &mut conn).await; @@ -394,22 +416,15 @@ pub mod test { let mut qb = QueryBuilder::new(""); let qb = Strategy::Newest.build_query(&mut qb); - let options = FormatOptions::default(); - let formatted_query = format(qb.sql().as_str(), &QueryParams::None, &options); - insta::assert_snapshot!(formatted_query, @" - SELECT - * - FROM - chunks - WHERE - opsqueue_is_reserved(chunks.submission_id, chunks.chunk_index) = FALSE - ORDER BY - submission_id DESC - "); + assert!(qb.sql().as_str().contains("ORDER BY id DESC")); + assert!( + qb.sql() + .as_str() + .contains("chunks.submission_id = newest_submission_ids.submission_id") + ); let explained = explain(qb, &mut conn).await; - assert_streaming_query(qb, &explained); - assert_eq!(explained, "3, 0, SCAN chunks"); + assert_streaming_chunks(qb, &explained); } #[sqlx::test(migrator = "crate::MIGRATOR")] @@ -420,38 +435,10 @@ pub mod test { let qb = Strategy::Random.build_query(&mut qb); - let formatted_query = format( - qb.sql().as_str(), - &QueryParams::None, - &FormatOptions::default(), - ); - insta::assert_snapshot!(formatted_query, @" - SELECT - * - FROM - chunks - WHERE - random_order >= ? - AND opsqueue_is_reserved(chunks.submission_id, chunks.chunk_index) = FALSE - UNION ALL - SELECT - * - FROM - chunks - WHERE - random_order < ? - AND opsqueue_is_reserved(chunks.submission_id, chunks.chunk_index) = FALSE - "); + assert!(qb.sql().as_str().contains("MIN(previous.chunk_index) + 4")); let explained = explain(qb, &mut conn).await; assert_streaming_query(qb, &explained); - insta::assert_snapshot!(explained, @r" - 1, 0, COMPOUND QUERY - 2, 1, LEFT-MOST SUBQUERY - 5, 2, SEARCH chunks USING INDEX random_chunks_order (random_order>?) - 26, 1, UNION ALL - 29, 26, SEARCH chunks USING INDEX random_chunks_order (random_order Date: Fri, 9 Oct 2026 07:20:29 +0200 Subject: [PATCH 06/25] Fix: update completion order in test --- libs/opsqueue_python/tests/test_roundtrip.py | 14 +++++++------- 1 file changed, 7 insertions(+), 7 deletions(-) diff --git a/libs/opsqueue_python/tests/test_roundtrip.py b/libs/opsqueue_python/tests/test_roundtrip.py index 6dd5dd75..1d804b02 100644 --- a/libs/opsqueue_python/tests/test_roundtrip.py +++ b/libs/opsqueue_python/tests/test_roundtrip.py @@ -844,19 +844,19 @@ def complete_chunks(_submission_id_value: int) -> None: ), key=lambda chunk: chunk.chunk_index, ) - consumer_client.complete_chunk( - chunks[1].submission_id, - chunks[1].submission_prefix, - chunks[1].chunk_index, - chunks[1].input_content, - ) - time.sleep(0.25) consumer_client.complete_chunk( chunks[0].submission_id, chunks[0].submission_prefix, chunks[0].chunk_index, chunks[0].input_content, ) + time.sleep(0.25) + consumer_client.complete_chunk( + chunks[1].submission_id, + chunks[1].submission_prefix, + chunks[1].chunk_index, + chunks[1].input_content, + ) with background_process( complete_chunks, From b25a44cb659d89e086ad2c306565927d100cb5fc Mon Sep 17 00:00:00 2001 From: Vince van Noort Date: Fri, 9 Oct 2026 07:22:34 +0200 Subject: [PATCH 07/25] Fix: update query plan for submission status --- opsqueue/src/common/submission.rs | 6 ++++-- 1 file changed, 4 insertions(+), 2 deletions(-) diff --git a/opsqueue/src/common/submission.rs b/opsqueue/src/common/submission.rs index 64e5af5d..8e9059ea 100644 --- a/opsqueue/src/common/submission.rs +++ b/opsqueue/src/common/submission.rs @@ -1707,8 +1707,10 @@ pub mod test { assert_no_temporary_b_trees(explained.as_str()); insta::assert_snapshot!(explained, @" 3, 0, SEARCH submissions USING INDEX sqlite_autoindex_submissions_1 (id=?) - 17, 0, CORRELATED SCALAR SUBQUERY 1 - 22, 17, SEARCH submissions_metadata USING PRIMARY KEY (submission_id=?) + 15, 0, CORRELATED SCALAR SUBQUERY 1 + 20, 15, SEARCH chunks USING PRIMARY KEY (submission_id=?) + 40, 0, CORRELATED SCALAR SUBQUERY 2 + 45, 40, SEARCH submissions_metadata USING PRIMARY KEY (submission_id=?) "); } From 004b47260334a03b21f68ac0c9750aaa17a24cda Mon Sep 17 00:00:00 2001 From: Vince van Noort Date: Fri, 9 Oct 2026 07:30:13 +0200 Subject: [PATCH 08/25] Fix: revert other strategy changes --- opsqueue/src/consumer/strategy.rs | 77 ++++++++++++------------------- 1 file changed, 29 insertions(+), 48 deletions(-) diff --git a/opsqueue/src/consumer/strategy.rs b/opsqueue/src/consumer/strategy.rs index 916200b4..5f6af335 100644 --- a/opsqueue/src/consumer/strategy.rs +++ b/opsqueue/src/consumer/strategy.rs @@ -7,14 +7,10 @@ use sqlx::{QueryBuilder, Sqlite}; #[cfg(feature = "server-logic")] use crate::common::chunk::Chunk; -#[cfg(feature = "server-logic")] -const RANDOM_CHUNK_WINDOW: u64 = 4; - #[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] pub enum Strategy { Oldest, Newest, - /// Randomize chunks within a small, advancing index window per submission. Random, /// Perform a sort of submissions by metadata (ordered by in-flight counts), /// ties are broken by the underlying strategy. For example: if we have two @@ -104,27 +100,10 @@ impl Strategy { Oldest => qb.push(format!( "SELECT * FROM chunks WHERE {ffi_is_not_reserved} ORDER BY submission_id ASC, chunk_index ASC" )), - Newest => { - let qb = qb.push("WITH newest_submission_ids AS MATERIALIZED ("); - let qb = qb.push("SELECT id as submission_id FROM submissions ORDER BY id DESC"); - qb.push(format!( - ") SELECT chunks.* - FROM newest_submission_ids - CROSS JOIN chunks - ON chunks.submission_id = newest_submission_ids.submission_id - AND {ffi_is_not_reserved}" - )) - } - Random => { - let condition = format!( - "{ffi_is_not_reserved} AND chunks.chunk_index < ( - SELECT MIN(previous.chunk_index) + {RANDOM_CHUNK_WINDOW} - FROM chunks AS previous - WHERE previous.submission_id = chunks.submission_id - )" - ); - Self::push_random_order_query(qb, "*", "chunks", Some(&condition)) - } + 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)), PreferDistinct { .. } => { // Unique submission IDs from the underlying strategy. let qb = qb.push("WITH underlying_submission_ids AS MATERIALIZED ("); @@ -416,15 +395,10 @@ pub mod test { let mut qb = QueryBuilder::new(""); let qb = Strategy::Newest.build_query(&mut qb); - assert!(qb.sql().as_str().contains("ORDER BY id DESC")); - assert!( - qb.sql() - .as_str() - .contains("chunks.submission_id = newest_submission_ids.submission_id") - ); + assert!(qb.sql().as_str().contains("ORDER BY submission_id DESC")); let explained = explain(qb, &mut conn).await; - assert_streaming_chunks(qb, &explained); + assert_streaming_query(qb, &explained); } #[sqlx::test(migrator = "crate::MIGRATOR")] @@ -435,7 +409,8 @@ pub mod test { let qb = Strategy::Random.build_query(&mut qb); - assert!(qb.sql().as_str().contains("MIN(previous.chunk_index) + 4")); + assert!(qb.sql().as_str().contains("random_order >= ?")); + assert!(qb.sql().as_str().contains("random_order < ?")); let explained = explain(qb, &mut conn).await; assert_streaming_query(qb, &explained); @@ -823,9 +798,7 @@ pub mod test { #[sqlx::test(migrator = "crate::MIGRATOR")] /// Tests whether the 'cutting the deck' technique is working /// - /// We do this by checking whether two selects in a huge amount of available chunks - /// give a different result. - /// (There is a super tiny chance of this test flaking). + /// Repeated selects should eventually return a different ordering. pub async fn test_random_strategy_is_random(pool: sqlx::SqlitePool) { let db_pools = crate::db::DBPools::from_test_pool(&pool); @@ -845,24 +818,32 @@ pub mod test { let mut conn = db_pools.reader_conn().await.unwrap(); register_lookup_noops(conn.get_inner()).await; - let mut query_builder = QueryBuilder::default(); - let vals1: Vec = Strategy::Random - .build_query(&mut query_builder) + let mut first_query = QueryBuilder::default(); + let first_result: Vec = Strategy::Random + .build_query(&mut first_query) .build_query_as() .fetch(conn.get_inner()) .try_collect() .await .unwrap(); - let mut query_builder = QueryBuilder::default(); - let vals2: Vec = Strategy::Random - .build_query(&mut query_builder) - .build_query_as() - .fetch(conn.get_inner()) - .try_collect() - .await - .unwrap(); + let mut observed_different_order = false; + for _ in 0..32 { + let mut query_builder = QueryBuilder::default(); + let result: Vec = Strategy::Random + .build_query(&mut query_builder) + .build_query_as() + .fetch(conn.get_inner()) + .try_collect() + .await + .unwrap(); + + if result != first_result { + observed_different_order = true; + break; + } + } - assert_ne!(vals1, vals2); + assert!(observed_different_order); } } From fe6d0391af5565cd456ec6a1a2ace534dc675ac9 Mon Sep 17 00:00:00 2001 From: Vince van Noort Date: Fri, 9 Oct 2026 07:32:11 +0200 Subject: [PATCH 09/25] Fix: revert other test changes --- opsqueue/src/consumer/strategy.rs | 69 ++++++++++++++++++++----------- 1 file changed, 45 insertions(+), 24 deletions(-) diff --git a/opsqueue/src/consumer/strategy.rs b/opsqueue/src/consumer/strategy.rs index 5f6af335..58cfb99d 100644 --- a/opsqueue/src/consumer/strategy.rs +++ b/opsqueue/src/consumer/strategy.rs @@ -214,7 +214,7 @@ impl Strategy { qb: &'a mut QueryBuilder, columns: &'static str, table_name: &'static str, - condition: Option<&str>, + condition: Option<&'static str>, ) -> &'a mut QueryBuilder { let random_offset: u16 = rand::random(); let push_select = |qb: &mut QueryBuilder, operator: &str| { @@ -409,11 +409,38 @@ pub mod test { let qb = Strategy::Random.build_query(&mut qb); - assert!(qb.sql().as_str().contains("random_order >= ?")); - assert!(qb.sql().as_str().contains("random_order < ?")); + let formatted_query = format( + qb.sql().as_str(), + &QueryParams::None, + &FormatOptions::default(), + ); + insta::assert_snapshot!(formatted_query, @" + SELECT + * + FROM + chunks + WHERE + random_order >= ? + AND opsqueue_is_reserved(chunks.submission_id, chunks.chunk_index) = FALSE + UNION ALL + SELECT + * + FROM + chunks + WHERE + random_order < ? + AND opsqueue_is_reserved(chunks.submission_id, chunks.chunk_index) = FALSE + "); let explained = explain(qb, &mut conn).await; assert_streaming_query(qb, &explained); + insta::assert_snapshot!(explained, @r" + 1, 0, COMPOUND QUERY + 2, 1, LEFT-MOST SUBQUERY + 5, 2, SEARCH chunks USING INDEX random_chunks_order (random_order>?) + 26, 1, UNION ALL + 29, 26, SEARCH chunks USING INDEX random_chunks_order (random_order = Strategy::Random - .build_query(&mut first_query) + let mut query_builder = QueryBuilder::default(); + let vals1: Vec = Strategy::Random + .build_query(&mut query_builder) .build_query_as() .fetch(conn.get_inner()) .try_collect() .await .unwrap(); - let mut observed_different_order = false; - for _ in 0..32 { - let mut query_builder = QueryBuilder::default(); - let result: Vec = Strategy::Random - .build_query(&mut query_builder) - .build_query_as() - .fetch(conn.get_inner()) - .try_collect() - .await - .unwrap(); - - if result != first_result { - observed_different_order = true; - break; - } - } + let mut query_builder = QueryBuilder::default(); + let vals2: Vec = Strategy::Random + .build_query(&mut query_builder) + .build_query_as() + .fetch(conn.get_inner()) + .try_collect() + .await + .unwrap(); - assert!(observed_different_order); + assert_ne!(vals1, vals2); } } From ca7821bb5df92b211021bcac6c08a81a18d315a2 Mon Sep 17 00:00:00 2001 From: Vince van Noort Date: Fri, 9 Oct 2026 07:32:34 +0200 Subject: [PATCH 10/25] fixup! Fix: revert other test changes --- opsqueue/src/consumer/strategy.rs | 14 +++++++++++++- 1 file changed, 13 insertions(+), 1 deletion(-) diff --git a/opsqueue/src/consumer/strategy.rs b/opsqueue/src/consumer/strategy.rs index 58cfb99d..9d0f954a 100644 --- a/opsqueue/src/consumer/strategy.rs +++ b/opsqueue/src/consumer/strategy.rs @@ -395,10 +395,22 @@ pub mod test { let mut qb = QueryBuilder::new(""); let qb = Strategy::Newest.build_query(&mut qb); - assert!(qb.sql().as_str().contains("ORDER BY submission_id DESC")); + let options = FormatOptions::default(); + let formatted_query = format(qb.sql().as_str(), &QueryParams::None, &options); + insta::assert_snapshot!(formatted_query, @" + SELECT + * + FROM + chunks + WHERE + opsqueue_is_reserved(chunks.submission_id, chunks.chunk_index) = FALSE + ORDER BY + submission_id DESC + "); let explained = explain(qb, &mut conn).await; assert_streaming_query(qb, &explained); + assert_eq!(explained, "3, 0, SCAN chunks"); } #[sqlx::test(migrator = "crate::MIGRATOR")] From cdf3f62c1c112c82aca8d59cb3c7d1a1b98cd9b8 Mon Sep 17 00:00:00 2001 From: Vince van Noort Date: Fri, 9 Oct 2026 07:48:02 +0200 Subject: [PATCH 11/25] Fix: ensure backwards compatability --- opsqueue/src/common/submission.rs | 1 + 1 file changed, 1 insertion(+) diff --git a/opsqueue/src/common/submission.rs b/opsqueue/src/common/submission.rs index 8e9059ea..5efbf2e2 100644 --- a/opsqueue/src/common/submission.rs +++ b/opsqueue/src/common/submission.rs @@ -156,6 +156,7 @@ pub struct Submission { pub prefix: Option, pub chunks_total: ChunkCount, pub chunks_done: ChunkCount, + #[serde(default = "ChunkCount::zero")] pub chunks_ready: ChunkCount, pub chunk_size: ChunkSize, pub metadata: Option, From 2e7a53ddaa8d7beec2951f4e49331b71c31deeba Mon Sep 17 00:00:00 2001 From: Vince van Noort Date: Fri, 9 Oct 2026 07:49:51 +0200 Subject: [PATCH 12/25] Fix: extend test suite --- libs/opsqueue_python/tests/test_roundtrip.py | 82 ++++++++++++++++++++ 1 file changed, 82 insertions(+) diff --git a/libs/opsqueue_python/tests/test_roundtrip.py b/libs/opsqueue_python/tests/test_roundtrip.py index 1d804b02..1a29da06 100644 --- a/libs/opsqueue_python/tests/test_roundtrip.py +++ b/libs/opsqueue_python/tests/test_roundtrip.py @@ -918,6 +918,88 @@ async def collect() -> list[bytes]: assert asyncio.run(collect()) == [b"[1]", b"[2]"] +def test_streams_chunks_in_order_when_consumers_complete_out_of_order( + opsqueue: OpsqueueProcess, +) -> None: + url = "file:///tmp/opsqueue/test_streaming_out_of_order_consumers" + producer_client = ProducerClient(f"localhost:{opsqueue.port}", url) + submission_id = producer_client.insert_submission_chunks( + [b"[1]", b"[2]"], chunk_size=1 + ) + first_consumer = ConsumerClient(f"localhost:{opsqueue.port}", url) + second_consumer = ConsumerClient(f"localhost:{opsqueue.port}", url) + + [first_chunk] = first_consumer.reserve_chunks( + max=1, strategy=Strategy.Oldest() + ) + [second_chunk] = second_consumer.reserve_chunks( + max=1, strategy=Strategy.Oldest() + ) + assert (first_chunk.chunk_index, second_chunk.chunk_index) == (0, 1) + + second_consumer.complete_chunk( + second_chunk.submission_id, + second_chunk.submission_prefix, + second_chunk.chunk_index, + second_chunk.input_content, + ) + + def complete_first_chunk(_submission_id_value: int) -> None: + time.sleep(0.25) + consumer_client = ConsumerClient(f"localhost:{opsqueue.port}", url) + consumer_client.complete_chunk( + first_chunk.submission_id, + first_chunk.submission_prefix, + first_chunk.chunk_index, + first_chunk.input_content, + ) + + with background_process(complete_first_chunk, args=(submission_id.id,)): + results = producer_client.stream_submission_chunks( + submission_id, Strategy.Oldest() + ) + assert next(results) == b"[1]" + assert next(results) == b"[2]" + + +def test_stream_submission_chunks_yields_completed_results_before_failure( + opsqueue: OpsqueueProcess, +) -> None: + url = "file:///tmp/opsqueue/test_streaming_results_before_failure" + producer_client = ProducerClient(f"localhost:{opsqueue.port}", url) + submission_id = producer_client.insert_submission_chunks( + [b"[1]", b"[2]"], chunk_size=1 + ) + + def complete_then_fail(_submission_id_value: int) -> None: + consumer_client = ConsumerClient(f"localhost:{opsqueue.port}", url) + first_chunk, second_chunk = sorted( + consumer_client.reserve_chunks(max=2, strategy=Strategy.Oldest()), + key=lambda chunk: chunk.chunk_index, + ) + consumer_client.complete_chunk( + first_chunk.submission_id, + first_chunk.submission_prefix, + first_chunk.chunk_index, + first_chunk.input_content, + ) + time.sleep(0.25) + consumer_client.fail_chunk( + second_chunk.submission_id, + second_chunk.submission_prefix, + second_chunk.chunk_index, + "Simulated failure", + ) + + with background_process(complete_then_fail, args=(submission_id.id,)): + results = producer_client.stream_submission_chunks( + submission_id, Strategy.Oldest() + ) + assert next(results) == b"[1]" + with pytest.raises(SubmissionFailedError): + next(results) + + def test_stream_submission_chunks_requires_oldest_strategy( opsqueue: OpsqueueProcess, ) -> None: From 3743125203d4b6a48cf5cd0ed48a7521f39c29d9 Mon Sep 17 00:00:00 2001 From: Vince van Noort Date: Fri, 9 Oct 2026 07:50:11 +0200 Subject: [PATCH 13/25] fixup! Fix: extend test suite --- libs/opsqueue_python/tests/test_roundtrip.py | 8 ++------ 1 file changed, 2 insertions(+), 6 deletions(-) diff --git a/libs/opsqueue_python/tests/test_roundtrip.py b/libs/opsqueue_python/tests/test_roundtrip.py index 1a29da06..2197446c 100644 --- a/libs/opsqueue_python/tests/test_roundtrip.py +++ b/libs/opsqueue_python/tests/test_roundtrip.py @@ -929,12 +929,8 @@ def test_streams_chunks_in_order_when_consumers_complete_out_of_order( first_consumer = ConsumerClient(f"localhost:{opsqueue.port}", url) second_consumer = ConsumerClient(f"localhost:{opsqueue.port}", url) - [first_chunk] = first_consumer.reserve_chunks( - max=1, strategy=Strategy.Oldest() - ) - [second_chunk] = second_consumer.reserve_chunks( - max=1, strategy=Strategy.Oldest() - ) + [first_chunk] = first_consumer.reserve_chunks(max=1, strategy=Strategy.Oldest()) + [second_chunk] = second_consumer.reserve_chunks(max=1, strategy=Strategy.Oldest()) assert (first_chunk.chunk_index, second_chunk.chunk_index) == (0, 1) second_consumer.complete_chunk( From cfdb994e2bd7a039f2fea429ecbe3a95dd589256 Mon Sep 17 00:00:00 2001 From: Vince van Noort Date: Fri, 9 Oct 2026 08:01:50 +0200 Subject: [PATCH 14/25] Fix: document error and remove redundant must use --- libs/opsqueue_python/src/producer.rs | 5 ++++- 1 file changed, 4 insertions(+), 1 deletion(-) diff --git a/libs/opsqueue_python/src/producer.rs b/libs/opsqueue_python/src/producer.rs index 1cd19ad1..4e37e4da 100644 --- a/libs/opsqueue_python/src/producer.rs +++ b/libs/opsqueue_python/src/producer.rs @@ -425,7 +425,10 @@ impl ProducerClient { /// /// `strategy` must be `Oldest`, and consumers processing this submission must /// also reserve chunks using `Oldest`. - #[must_use] + /// + /// # Errors + /// + /// Returns `ValueError` if `strategy` is not `Oldest`. pub fn stream_submission_chunks( &self, submission_id: SubmissionId, From feae80c27e7ff4be4f544c68d2bf205ff6e1bc9f Mon Sep 17 00:00:00 2001 From: Vince van Noort Date: Fri, 9 Oct 2026 08:31:06 +0200 Subject: [PATCH 15/25] Fix: bump version --- Cargo.lock | 4 ++-- Cargo.toml | 2 +- 2 files changed, 3 insertions(+), 3 deletions(-) diff --git a/Cargo.lock b/Cargo.lock index 25de1869..51e92265 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -1995,7 +1995,7 @@ dependencies = [ [[package]] name = "opsqueue" -version = "0.41.0" +version = "0.42.0" dependencies = [ "anyhow", "arc-swap", @@ -2048,7 +2048,7 @@ dependencies = [ [[package]] name = "opsqueue_python" -version = "0.41.0" +version = "0.42.0" dependencies = [ "anyhow", "chrono", diff --git a/Cargo.toml b/Cargo.toml index a650f0b6..0021f7d9 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -14,7 +14,7 @@ hakari-package = "workspace-hack" [workspace.package] edition = "2024" -version = "0.41.0" +version = "0.42.0" [workspace.dependencies] anyhow = { version = "1.0.102", default-features = false } From 85bc8fa1707c40db68113ac2e7db5d71cbda81b0 Mon Sep 17 00:00:00 2001 From: Vince van Noort Date: Fri, 9 Oct 2026 08:49:34 +0200 Subject: [PATCH 16/25] Fix: use context manager for service --- libs/opsqueue_python/tests/test_roundtrip.py | 65 ++++++++++---------- 1 file changed, 32 insertions(+), 33 deletions(-) diff --git a/libs/opsqueue_python/tests/test_roundtrip.py b/libs/opsqueue_python/tests/test_roundtrip.py index 2197446c..cb61d047 100644 --- a/libs/opsqueue_python/tests/test_roundtrip.py +++ b/libs/opsqueue_python/tests/test_roundtrip.py @@ -958,42 +958,41 @@ def complete_first_chunk(_submission_id_value: int) -> None: assert next(results) == b"[2]" -def test_stream_submission_chunks_yields_completed_results_before_failure( - opsqueue: OpsqueueProcess, -) -> None: +def test_stream_submission_chunks_yields_completed_results_before_failure() -> None: url = "file:///tmp/opsqueue/test_streaming_results_before_failure" - producer_client = ProducerClient(f"localhost:{opsqueue.port}", url) - submission_id = producer_client.insert_submission_chunks( - [b"[1]", b"[2]"], chunk_size=1 - ) - - def complete_then_fail(_submission_id_value: int) -> None: - consumer_client = ConsumerClient(f"localhost:{opsqueue.port}", url) - first_chunk, second_chunk = sorted( - consumer_client.reserve_chunks(max=2, strategy=Strategy.Oldest()), - key=lambda chunk: chunk.chunk_index, - ) - consumer_client.complete_chunk( - first_chunk.submission_id, - first_chunk.submission_prefix, - first_chunk.chunk_index, - first_chunk.input_content, - ) - time.sleep(0.25) - consumer_client.fail_chunk( - second_chunk.submission_id, - second_chunk.submission_prefix, - second_chunk.chunk_index, - "Simulated failure", + with opsqueue_service(command_args=("--max-chunk-retries", "1")) as opsqueue: + producer_client = ProducerClient(f"localhost:{opsqueue.port}", url) + submission_id = producer_client.insert_submission_chunks( + [b"[1]", b"[2]"], chunk_size=1 ) - with background_process(complete_then_fail, args=(submission_id.id,)): - results = producer_client.stream_submission_chunks( - submission_id, Strategy.Oldest() - ) - assert next(results) == b"[1]" - with pytest.raises(SubmissionFailedError): - next(results) + def complete_then_fail(_submission_id_value: int) -> None: + consumer_client = ConsumerClient(f"localhost:{opsqueue.port}", url) + first_chunk, second_chunk = sorted( + consumer_client.reserve_chunks(max=2, strategy=Strategy.Oldest()), + key=lambda chunk: chunk.chunk_index, + ) + consumer_client.complete_chunk( + first_chunk.submission_id, + first_chunk.submission_prefix, + first_chunk.chunk_index, + first_chunk.input_content, + ) + time.sleep(0.25) + consumer_client.fail_chunk( + second_chunk.submission_id, + second_chunk.submission_prefix, + second_chunk.chunk_index, + "Simulated failure", + ) + + with background_process(complete_then_fail, args=(submission_id.id,)): + results = producer_client.stream_submission_chunks( + submission_id, Strategy.Oldest() + ) + assert next(results) == b"[1]" + with pytest.raises(SubmissionFailedError): + next(results) def test_stream_submission_chunks_requires_oldest_strategy( From 8882e438763609bcae3ff4f54c0aaace365ae518 Mon Sep 17 00:00:00 2001 From: Vince van Noort Date: Fri, 9 Oct 2026 08:57:27 +0200 Subject: [PATCH 17/25] Fix: update assert --- libs/opsqueue_python/tests/test_roundtrip.py | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/libs/opsqueue_python/tests/test_roundtrip.py b/libs/opsqueue_python/tests/test_roundtrip.py index cb61d047..35869693 100644 --- a/libs/opsqueue_python/tests/test_roundtrip.py +++ b/libs/opsqueue_python/tests/test_roundtrip.py @@ -931,7 +931,7 @@ def test_streams_chunks_in_order_when_consumers_complete_out_of_order( [first_chunk] = first_consumer.reserve_chunks(max=1, strategy=Strategy.Oldest()) [second_chunk] = second_consumer.reserve_chunks(max=1, strategy=Strategy.Oldest()) - assert (first_chunk.chunk_index, second_chunk.chunk_index) == (0, 1) + assert (first_chunk.chunk_index.id, second_chunk.chunk_index.id) == (0, 1) second_consumer.complete_chunk( second_chunk.submission_id, From 1d248bd6b34ca54f46dfde4cd1377ea4335b963a Mon Sep 17 00:00:00 2001 From: Vince van Noort Date: Fri, 9 Oct 2026 09:26:50 +0200 Subject: [PATCH 18/25] Fix: prevent complete pickle --- libs/opsqueue_python/tests/test_roundtrip.py | 33 +++++++++++++++----- 1 file changed, 25 insertions(+), 8 deletions(-) diff --git a/libs/opsqueue_python/tests/test_roundtrip.py b/libs/opsqueue_python/tests/test_roundtrip.py index 35869693..79df9f42 100644 --- a/libs/opsqueue_python/tests/test_roundtrip.py +++ b/libs/opsqueue_python/tests/test_roundtrip.py @@ -16,7 +16,7 @@ strategy_from_description, ) from opsqueue.common import SerializationFormat -from opsqueue.consumer import ConsumerClient, Chunk, Strategy +from opsqueue.consumer import ConsumerClient, Chunk, Strategy, opsqueue_internal from opsqueue.producer import ( SubmissionId, ProducerClient, @@ -940,17 +940,34 @@ def test_streams_chunks_in_order_when_consumers_complete_out_of_order( second_chunk.input_content, ) - def complete_first_chunk(_submission_id_value: int) -> None: + opsqueue_address = f"localhost:{opsqueue.port}" + + def complete_first_chunk( + _submission_id_value: int, + chunk_submission_id: int, + submission_prefix: str, + chunk_index: int, + input_content: bytes, + ) -> None: time.sleep(0.25) - consumer_client = ConsumerClient(f"localhost:{opsqueue.port}", url) + consumer_client = ConsumerClient(opsqueue_address, url) consumer_client.complete_chunk( - first_chunk.submission_id, - first_chunk.submission_prefix, - first_chunk.chunk_index, - first_chunk.input_content, + SubmissionId(chunk_submission_id), + submission_prefix, + opsqueue_internal.ChunkIndex(chunk_index), + input_content, ) - with background_process(complete_first_chunk, args=(submission_id.id,)): + with background_process( + complete_first_chunk, + args=( + submission_id.id, + first_chunk.submission_id.id, + first_chunk.submission_prefix, + first_chunk.chunk_index.id, + first_chunk.input_content, + ), + ): results = producer_client.stream_submission_chunks( submission_id, Strategy.Oldest() ) From 6b61b0f4139740c33ae2579fb2891620e1623eab Mon Sep 17 00:00:00 2001 From: Vince van Noort Date: Fri, 9 Oct 2026 09:34:02 +0200 Subject: [PATCH 19/25] Feat: support prefer distinct with underlying oldest strategy --- libs/opsqueue_python/src/producer.rs | 15 +++++++---- libs/opsqueue_python/tests/conftest.py | 26 ++++++++++++++++++++ libs/opsqueue_python/tests/test_roundtrip.py | 24 ++++++++++++------ 3 files changed, 53 insertions(+), 12 deletions(-) diff --git a/libs/opsqueue_python/src/producer.rs b/libs/opsqueue_python/src/producer.rs index 4e37e4da..23fd5bba 100644 --- a/libs/opsqueue_python/src/producer.rs +++ b/libs/opsqueue_python/src/producer.rs @@ -21,6 +21,7 @@ use opsqueue::{ common::errors::E::{self, L, R}, common::errors::{SubmissionNotCancellable, SubmissionNotFound, TooManyMatchingSubmissions}, common::{StrategicMetadataMap, chunk, submission}, + consumer::strategy::Strategy as ConsumerStrategy, object_store::{ChunkRetrievalError, ChunkType, ChunksStorageError, NewObjectStoreClientError}, producer::ChunkContents, producer::client::{Client as ActualClient, InternalProducerClientError}, @@ -423,12 +424,13 @@ impl ProducerClient { /// Stream output chunks as soon as each consumer has completed them. /// - /// `strategy` must be `Oldest`, and consumers processing this submission must - /// also reserve chunks using `Oldest`. + /// `strategy` must be `Oldest` or `PreferDistinct` with an underlying strategy + /// that eventually resolves to `Oldest`. Consumers processing this submission + /// must also reserve chunks using the same strategy. /// /// # Errors /// - /// Returns `ValueError` if `strategy` is not `Oldest`. + /// Returns `ValueError` if the strategy does not resolve to `Oldest`. pub fn stream_submission_chunks( &self, submission_id: SubmissionId, @@ -443,9 +445,12 @@ impl ProducerClient { submission_id: SubmissionId, strategy: &Strategy, ) -> PyResult { - if !matches!(strategy, Strategy::Oldest()) { + let internal_strategy = ConsumerStrategy::from(strategy); + let mut meta_keys = internal_strategy.meta_keys(); + meta_keys.by_ref().for_each(drop); + if !matches!(meta_keys.take(), ConsumerStrategy::Oldest) { return Err(PyValueError::new_err( - "streaming submission chunks requires Strategy.Oldest; consumers must also use Strategy.Oldest", + "streaming submission chunks requires Strategy.Oldest or Strategy.PreferDistinct ending in Strategy.Oldest; consumers must use the same strategy", )); } diff --git a/libs/opsqueue_python/tests/conftest.py b/libs/opsqueue_python/tests/conftest.py index e8c0f7ed..5c49e57b 100644 --- a/libs/opsqueue_python/tests/conftest.py +++ b/libs/opsqueue_python/tests/conftest.py @@ -366,6 +366,21 @@ def multiple_background_processes( ) +def _ends_with_oldest(strategy: StrategyDescription) -> bool: + match strategy: + case "Oldest": + return True + case ("PreferDistinct", _, underlying): + return _ends_with_oldest(underlying) + case _: + return False + + +oldest_strategies: tuple[StrategyDescription, ...] = tuple( + strategy for strategy in any_strategies if _ends_with_oldest(strategy) +) + + @pytest.fixture( scope="function", ids=lambda s: f"Strategy.{strategy_from_description(s)}", @@ -388,6 +403,17 @@ def any_consumer_strategy( yield request.param +@pytest.fixture( + scope="function", + ids=lambda s: f"Strategy.{strategy_from_description(s)}", + params=oldest_strategies, +) +def oldest_consumer_strategy( + request: pytest.FixtureRequest, +) -> Generator[StrategyDescription, None, None]: + yield request.param + + @pytest.fixture(scope="function", params=[json_as_bytes, cbor2, pickle]) def serialization_format( request: pytest.FixtureRequest, diff --git a/libs/opsqueue_python/tests/test_roundtrip.py b/libs/opsqueue_python/tests/test_roundtrip.py index 79df9f42..8e50c7aa 100644 --- a/libs/opsqueue_python/tests/test_roundtrip.py +++ b/libs/opsqueue_python/tests/test_roundtrip.py @@ -875,8 +875,10 @@ def complete_chunks(_submission_id_value: int) -> None: def test_async_streams_completed_chunks_before_submission_finishes( opsqueue: OpsqueueProcess, + oldest_consumer_strategy: StrategyDescription, ) -> None: url = "file:///tmp/opsqueue/test_async_streaming_results" + strategy = strategy_from_description(oldest_consumer_strategy) producer_client = ProducerClient(f"localhost:{opsqueue.port}", url) submission_id = producer_client.insert_submission_chunks( [b"[1]", b"[2]"], chunk_size=1 @@ -887,7 +889,7 @@ def complete_chunks(_submission_id_value: int) -> None: chunks = sorted( consumer_client.reserve_chunks( max=2, - strategy=Strategy.Oldest(), + strategy=strategy_from_description(oldest_consumer_strategy), ), key=lambda chunk: chunk.chunk_index, ) @@ -907,7 +909,7 @@ def complete_chunks(_submission_id_value: int) -> None: async def collect() -> list[bytes]: results = await producer_client.async_stream_submission_chunks( - submission_id, Strategy.Oldest() + submission_id, strategy ) return [chunk async for chunk in results] @@ -920,8 +922,10 @@ async def collect() -> list[bytes]: def test_streams_chunks_in_order_when_consumers_complete_out_of_order( opsqueue: OpsqueueProcess, + oldest_consumer_strategy: StrategyDescription, ) -> None: url = "file:///tmp/opsqueue/test_streaming_out_of_order_consumers" + strategy = strategy_from_description(oldest_consumer_strategy) producer_client = ProducerClient(f"localhost:{opsqueue.port}", url) submission_id = producer_client.insert_submission_chunks( [b"[1]", b"[2]"], chunk_size=1 @@ -929,8 +933,8 @@ def test_streams_chunks_in_order_when_consumers_complete_out_of_order( first_consumer = ConsumerClient(f"localhost:{opsqueue.port}", url) second_consumer = ConsumerClient(f"localhost:{opsqueue.port}", url) - [first_chunk] = first_consumer.reserve_chunks(max=1, strategy=Strategy.Oldest()) - [second_chunk] = second_consumer.reserve_chunks(max=1, strategy=Strategy.Oldest()) + [first_chunk] = first_consumer.reserve_chunks(max=1, strategy=strategy) + [second_chunk] = second_consumer.reserve_chunks(max=1, strategy=strategy) assert (first_chunk.chunk_index.id, second_chunk.chunk_index.id) == (0, 1) second_consumer.complete_chunk( @@ -968,9 +972,7 @@ def complete_first_chunk( first_chunk.input_content, ), ): - results = producer_client.stream_submission_chunks( - submission_id, Strategy.Oldest() - ) + results = producer_client.stream_submission_chunks(submission_id, strategy) assert next(results) == b"[1]" assert next(results) == b"[2]" @@ -1022,6 +1024,14 @@ def test_stream_submission_chunks_requires_oldest_strategy( with pytest.raises(ValueError, match="requires Strategy.Oldest"): producer_client.stream_submission_chunks(submission_id, Strategy.Random()) + with pytest.raises(ValueError, match="requires Strategy.Oldest"): + producer_client.stream_submission_chunks( + submission_id, + Strategy.PreferDistinct( + meta_key="company_id", underlying=Strategy.Newest() + ), + ) + async def collect_with_newest() -> None: await producer_client.async_stream_submission_chunks( submission_id, Strategy.Newest() From a2fdf7be475f60aed78c98d4c48711a48529a9a7 Mon Sep 17 00:00:00 2001 From: Vince van Noort Date: Fri, 9 Oct 2026 09:42:16 +0200 Subject: [PATCH 20/25] Feat: ignore attr defined --- libs/opsqueue_python/tests/test_roundtrip.py | 3 ++- 1 file changed, 2 insertions(+), 1 deletion(-) diff --git a/libs/opsqueue_python/tests/test_roundtrip.py b/libs/opsqueue_python/tests/test_roundtrip.py index 8e50c7aa..28970017 100644 --- a/libs/opsqueue_python/tests/test_roundtrip.py +++ b/libs/opsqueue_python/tests/test_roundtrip.py @@ -16,7 +16,8 @@ strategy_from_description, ) from opsqueue.common import SerializationFormat -from opsqueue.consumer import ConsumerClient, Chunk, Strategy, opsqueue_internal +from opsqueue.consumer import ConsumerClient, Chunk, Strategy +from opsqueue.consumer import opsqueue_internal # type: ignore[attr-defined] from opsqueue.producer import ( SubmissionId, ProducerClient, From 7fe4ade2596c72a22be48eff9309c93fbb141475 Mon Sep 17 00:00:00 2001 From: Vince van Noort Date: Fri, 9 Oct 2026 11:07:41 +0200 Subject: [PATCH 21/25] Fix: format over multiple lines --- libs/opsqueue_python/python/opsqueue/producer.py | 8 ++++++-- 1 file changed, 6 insertions(+), 2 deletions(-) diff --git a/libs/opsqueue_python/python/opsqueue/producer.py b/libs/opsqueue_python/python/opsqueue/producer.py index 08b79c1c..8363d187 100644 --- a/libs/opsqueue_python/python/opsqueue/producer.py +++ b/libs/opsqueue_python/python/opsqueue/producer.py @@ -242,7 +242,9 @@ def run_submission_chunks( return self.blocking_stream_completed_submission_chunks(submission_id, timeout) def stream_submission_chunks( - self, submission_id: SubmissionId, strategy: Strategy + self, + submission_id: SubmissionId, + strategy: Strategy, ) -> Iterator[bytes]: """Stream chunks progressively; strategy must be Oldest for this submission.""" return self.inner.stream_submission_chunks( # type: ignore[no-any-return] @@ -250,7 +252,9 @@ def stream_submission_chunks( ) async def async_stream_submission_chunks( - self, submission_id: SubmissionId, strategy: Strategy + self, + submission_id: SubmissionId, + strategy: Strategy, ) -> AsyncIterator[bytes]: """Stream chunks progressively; strategy must be Oldest for this submission.""" return await self.inner.async_stream_submission_chunks( # type: ignore[no-any-return] From 8cb53bc6f67d0060497735d00951991c80280ea8 Mon Sep 17 00:00:00 2001 From: Vince van Noort Date: Fri, 9 Oct 2026 11:48:08 +0200 Subject: [PATCH 22/25] Fix: remove trailing whitespace --- libs/opsqueue_python/python/opsqueue/producer.py | 8 ++++---- 1 file changed, 4 insertions(+), 4 deletions(-) diff --git a/libs/opsqueue_python/python/opsqueue/producer.py b/libs/opsqueue_python/python/opsqueue/producer.py index 8363d187..fbe21dbe 100644 --- a/libs/opsqueue_python/python/opsqueue/producer.py +++ b/libs/opsqueue_python/python/opsqueue/producer.py @@ -242,8 +242,8 @@ def run_submission_chunks( return self.blocking_stream_completed_submission_chunks(submission_id, timeout) def stream_submission_chunks( - self, - submission_id: SubmissionId, + self, + submission_id: SubmissionId, strategy: Strategy, ) -> Iterator[bytes]: """Stream chunks progressively; strategy must be Oldest for this submission.""" @@ -252,8 +252,8 @@ def stream_submission_chunks( ) async def async_stream_submission_chunks( - self, - submission_id: SubmissionId, + self, + submission_id: SubmissionId, strategy: Strategy, ) -> AsyncIterator[bytes]: """Stream chunks progressively; strategy must be Oldest for this submission.""" From 4dc87fad8cc8996e7b3ab6f609e691737ba1433b Mon Sep 17 00:00:00 2001 From: Vince van Noort Date: Fri, 9 Oct 2026 14:44:56 +0200 Subject: [PATCH 23/25] Fix: reduce polling overhead in chunk streaming --- libs/opsqueue_python/src/producer.rs | 159 ++++++++----------- libs/opsqueue_python/tests/test_roundtrip.py | 131 ++++++++------- 2 files changed, 144 insertions(+), 146 deletions(-) diff --git a/libs/opsqueue_python/src/producer.rs b/libs/opsqueue_python/src/producer.rs index 23fd5bba..8180faf1 100644 --- a/libs/opsqueue_python/src/producer.rs +++ b/libs/opsqueue_python/src/producer.rs @@ -33,6 +33,7 @@ use ux::u63; create_exception!(opsqueue_internal, ProducerClientError, PyException); const SUBMISSION_POLLING_INTERVAL: Duration = Duration::from_secs(5); +const INITIAL_SUBMISSION_POLLING_INTERVAL: Duration = Duration::from_millis(10); // NOTE: ProducerClient is reasonably cheap to clone, as most of its fields are behind Arcs. #[pyclass(from_py_object, module = "opsqueue")] @@ -457,129 +458,89 @@ impl ProducerClient { let client = self.client.clone(); let object_store_client = self.object_store_client.clone(); let stream = futures::stream::unfold( - ( + StreamingChunkState { client, object_store_client, submission_id, - u63::new(0), - None, - Duration::from_millis(10), - ), - |(client, object_store_client, submission_id, index, prefix, interval)| async move { - let mut interval = interval; + index: u63::new(0), + prefix: None, + ready_until: u63::new(0), + finished: false, + interval: INITIAL_SUBMISSION_POLLING_INTERVAL, + }, + |mut state| async move { loop { - let status = match client.get_submission(submission_id.into()).await { + if state.index < state.ready_until { + let prefix = state + .prefix + .as_deref() + .expect("ready submissions have an object-store prefix"); + let result = state + .object_store_client + .retrieve_chunk(prefix, state.index.into(), ChunkType::Output) + .await + .map_err(StreamingChunkError::Retrieval); + state.index = state.index + u63::new(1); + state.interval = INITIAL_SUBMISSION_POLLING_INTERVAL; + return Some((result, state)); + } + + if state.finished { + return None; + } + + let status = match state + .client + .get_submission(state.submission_id.into()) + .await + { Ok(Some(status)) => status, Ok(None) => { return Some(( - Err(StreamingChunkError::SubmissionNotFound), - ( - client, - object_store_client, - submission_id, - index, - prefix, - interval, - ), + Err(StreamingChunkError::SubmissionNotFound(SubmissionNotFound( + state.submission_id.into(), + ))), + state, )); } Err(error) => { - return Some(( - Err(StreamingChunkError::Internal(error)), - ( - client, - object_store_client, - submission_id, - index, - prefix, - interval, - ), - )); + return Some((Err(StreamingChunkError::Internal(error)), state)); } }; match status { submission::SubmissionStatus::InProgress(submission) => { - let prefix = prefix.clone().or(submission.prefix); - if index < submission.chunks_ready.into() { - let prefix = prefix - .expect("in-progress submissions have an object-store prefix"); - let result = object_store_client - .retrieve_chunk(&prefix, index.into(), ChunkType::Output) - .await - .map_err(StreamingChunkError::Retrieval); - return Some(( - result, - ( - client, - object_store_client, - submission_id, - index + u63::new(1), - Some(prefix), - interval, - ), - )); - } + state.prefix = state.prefix.take().or(submission.prefix); + state.ready_until = submission.chunks_ready.into(); } submission::SubmissionStatus::Completed(submission) => { - let prefix = prefix.clone().or(submission.prefix); - if index < submission.chunks_total.into() { - let prefix = prefix - .expect("completed submissions have an object-store prefix"); - let result = object_store_client - .retrieve_chunk(&prefix, index.into(), ChunkType::Output) - .await - .map_err(StreamingChunkError::Retrieval); - return Some(( - result, - ( - client, - object_store_client, - submission_id, - index + u63::new(1), - Some(prefix), - interval, - ), - )); - } - return None; + state.prefix = state.prefix.take().or(submission.prefix); + state.ready_until = submission.chunks_total.into(); + state.finished = true; } submission::SubmissionStatus::Failed(submission, chunk) => { let failure = crate::common::ChunkFailed::from_internal(chunk, &submission); + state.finished = true; return Some(( Err(StreamingChunkError::Failed(Box::new( crate::errors::SubmissionFailed(submission.into(), failure), ))), - ( - client, - object_store_client, - submission_id, - index, - prefix, - interval, - ), + state, )); } submission::SubmissionStatus::Paused(_) => {} submission::SubmissionStatus::Cancelled(_) => { - return Some(( - Err(StreamingChunkError::Cancelled), - ( - client, - object_store_client, - submission_id, - index, - prefix, - interval, - ), - )); + return Some((Err(StreamingChunkError::Cancelled), state)); } } - tokio::time::sleep(interval).await; - if interval < SUBMISSION_POLLING_INTERVAL { - interval = (interval * 2).min(SUBMISSION_POLLING_INTERVAL); + if state.index < state.ready_until || state.finished { + continue; + } + tokio::time::sleep(state.interval).await; + if state.interval < SUBMISSION_POLLING_INTERVAL { + state.interval = (state.interval * 2).min(SUBMISSION_POLLING_INTERVAL); } } }, @@ -773,12 +734,23 @@ impl ProducerClient { } } +struct StreamingChunkState { + client: ActualClient, + object_store_client: opsqueue::object_store::ObjectStoreClient, + submission_id: SubmissionId, + index: u63, + prefix: Option, + ready_until: u63, + finished: bool, + interval: Duration, +} + #[derive(Debug)] enum StreamingChunkError { Retrieval(ChunkRetrievalError), Internal(InternalProducerClientError), Failed(Box), - SubmissionNotFound, + SubmissionNotFound(SubmissionNotFound), Cancelled, } @@ -788,7 +760,7 @@ impl std::fmt::Display for StreamingChunkError { Self::Retrieval(error) => error.fmt(f), Self::Internal(error) => error.fmt(f), Self::Failed(_) => write!(f, "Submission failed"), - Self::SubmissionNotFound => write!(f, "Submission not found"), + Self::SubmissionNotFound(error) => error.fmt(f), Self::Cancelled => write!(f, "Submission cancelled"), } } @@ -802,6 +774,7 @@ impl From> for PyErr { StreamingChunkError::Retrieval(error) => CError(error).into(), StreamingChunkError::Internal(error) => CError(error).into(), StreamingChunkError::Failed(error) => CError(*error).into(), + StreamingChunkError::SubmissionNotFound(error) => CError(error).into(), error => PyException::new_err(error.to_string()), } } diff --git a/libs/opsqueue_python/tests/test_roundtrip.py b/libs/opsqueue_python/tests/test_roundtrip.py index 28970017..275cd13f 100644 --- a/libs/opsqueue_python/tests/test_roundtrip.py +++ b/libs/opsqueue_python/tests/test_roundtrip.py @@ -2,35 +2,40 @@ # - use pytest's `--log-cli-level=info` (or `=debug`) argument to get more detailed logs from the producer/consumer clients # - use `RUST_LOG="opsqueue=info"` (or `opsqueue=debug` or `debug` for even more verbosity), together with to the pytest option `-s` AKA `--capture=no`, to debug the opsqueue binary itself. +import asyncio import logging import time from collections.abc import Iterator, Sequence -import asyncio + import pytest from conftest import ( + OpsqueueProcess, + StrategyDescription, background_process, multiple_background_processes, - OpsqueueProcess, opsqueue_service, - StrategyDescription, strategy_from_description, ) from opsqueue.common import SerializationFormat -from opsqueue.consumer import ConsumerClient, Chunk, Strategy -from opsqueue.consumer import opsqueue_internal # type: ignore[attr-defined] +from opsqueue.consumer import ( + Chunk, + ConsumerClient, + Strategy, + opsqueue_internal, # type: ignore[attr-defined] +) from opsqueue.producer import ( - SubmissionId, + ChunkFailed, + InitialSubmissionStatus, ProducerClient, SubmissionCompleted, SubmissionFailed, - ChunkFailed, - SubmissionStatus, SubmissionFailedError, - SubmissionNotFoundError, + SubmissionId, SubmissionNotCancellable, SubmissionNotCancellableError, + SubmissionNotFoundError, + SubmissionStatus, TooManyMatchingSubmissionsError, - InitialSubmissionStatus, ) SUBMISSION_COMPLETED_TIMEOUT = 10.0 @@ -59,7 +64,7 @@ def run_consumer() -> None: consumer_client.run_each_op(increment, strategy=strategy) with background_process(run_consumer) as _consumer: - input_iter = range(0, 100) + input_iter = range(100) output_iter: Iterator[int] = producer_client.run_submission( input_iter, @@ -131,7 +136,7 @@ def run_consumer(_consumer_id: int) -> None: ) with multiple_background_processes(run_consumer, n_consumers) as _consumers: - input_iter = range(0, n_ops) + input_iter = range(n_ops) output_iter: Iterator[int] = producer_client.run_submission( input_iter, @@ -141,7 +146,7 @@ def run_consumer(_consumer_id: int) -> None: ) res = sum(output_iter) - assert res == sum(range(0, n_ops)) + assert res == sum(range(n_ops)) def test_empty_submission(opsqueue: OpsqueueProcess) -> None: @@ -191,7 +196,7 @@ def run_consumer() -> None: ) with background_process(run_consumer) as _consumer: - input_iter = range(0, 100) + input_iter = range(100) output_iter: Iterator[int] = producer_client.run_submission( input_iter, @@ -237,7 +242,7 @@ def broken_increment(input: int) -> float: with background_process(run_consumer) as consumer: logging.error(f"Opsqueue: {opsqueue}") logging.error(f"Consumer: {consumer}") - input_iter = range(0, 100) + input_iter = range(100) with pytest.raises(SubmissionFailedError) as exc_info: producer_client.run_submission( @@ -281,7 +286,7 @@ def increment_list(ints: Sequence[int], _chunk: Chunk) -> Sequence[int]: consumer_client.run_each_chunk(increment_list, strategy=strategy) with background_process(run_consumer) as _consumer: - input_iter = map(lambda i: cbor2.dumps([i, i, i]), range(0, 10)) + input_iter = map(lambda i: cbor2.dumps([i, i, i]), range(10)) output_iter: Iterator[list[int]] = map( lambda c: cbor2.loads(c), producer_client.run_submission_chunks( @@ -324,7 +329,7 @@ def run_consumer(consumer_id: int) -> None: n_consumers = 16 with multiple_background_processes(run_consumer, n_consumers) as _consumers: - input_iter = range(0, 1000) + input_iter = range(1000) output_iter: Iterator[int] = producer_client.run_submission( input_iter, chunk_size=100, @@ -357,7 +362,7 @@ def increment_list(ints: Sequence[int], _chunk: Chunk) -> Sequence[int]: async def run_one_submission(top: int) -> int: logging.debug(f"Running submission {top}") - input_iter = range(0, top) + input_iter = range(top) output_iter = await producer_client.async_run_submission( input_iter, chunk_size=1000 ) @@ -750,13 +755,12 @@ def process_op(x: int) -> int: consumer_client.run_each_op(process_op) - with background_process(run_consumer) as _consumer: - with pytest.raises(TimeoutError): - producer_client.run_submission( - [1], - chunk_size=1, - timeout=0.1, - ) + with background_process(run_consumer) as _consumer, pytest.raises(TimeoutError): + producer_client.run_submission( + [1], + chunk_size=1, + timeout=0.1, + ) def test_unpause_and_complete(opsqueue: OpsqueueProcess) -> None: @@ -978,41 +982,62 @@ def complete_first_chunk( assert next(results) == b"[2]" -def test_stream_submission_chunks_yields_completed_results_before_failure() -> None: - url = "file:///tmp/opsqueue/test_streaming_results_before_failure" +def test_stream_submission_chunks_fails_if_submission_failed_before_read() -> None: + url = "file:///tmp/opsqueue/test_streaming_failed_before_read" with opsqueue_service(command_args=("--max-chunk-retries", "1")) as opsqueue: producer_client = ProducerClient(f"localhost:{opsqueue.port}", url) submission_id = producer_client.insert_submission_chunks( [b"[1]", b"[2]"], chunk_size=1 ) + results = producer_client.stream_submission_chunks( + submission_id, Strategy.Oldest() + ) + consumer_client = ConsumerClient(f"localhost:{opsqueue.port}", url) + first_chunk, second_chunk = sorted( + consumer_client.reserve_chunks(max=2, strategy=Strategy.Oldest()), + key=lambda chunk: chunk.chunk_index, + ) + consumer_client.complete_chunk( + first_chunk.submission_id, + first_chunk.submission_prefix, + first_chunk.chunk_index, + first_chunk.input_content, + ) + consumer_client.fail_chunk( + second_chunk.submission_id, + second_chunk.submission_prefix, + second_chunk.chunk_index, + "Simulated failure", + ) + assert isinstance( + producer_client.get_submission_status(submission_id), + SubmissionStatus.Failed, + ) + with pytest.raises(SubmissionFailedError): + next(results) + with pytest.raises(StopIteration): + next(results) - def complete_then_fail(_submission_id_value: int) -> None: - consumer_client = ConsumerClient(f"localhost:{opsqueue.port}", url) - first_chunk, second_chunk = sorted( - consumer_client.reserve_chunks(max=2, strategy=Strategy.Oldest()), - key=lambda chunk: chunk.chunk_index, - ) - consumer_client.complete_chunk( - first_chunk.submission_id, - first_chunk.submission_prefix, - first_chunk.chunk_index, - first_chunk.input_content, - ) - time.sleep(0.25) - consumer_client.fail_chunk( - second_chunk.submission_id, - second_chunk.submission_prefix, - second_chunk.chunk_index, - "Simulated failure", - ) - with background_process(complete_then_fail, args=(submission_id.id,)): - results = producer_client.stream_submission_chunks( - submission_id, Strategy.Oldest() - ) - assert next(results) == b"[1]" - with pytest.raises(SubmissionFailedError): - next(results) +def test_stream_submission_chunks_not_found(opsqueue: OpsqueueProcess) -> None: + url = "file:///tmp/opsqueue/test_streaming_submission_not_found" + producer_client = ProducerClient(f"localhost:{opsqueue.port}", url) + submission_id = SubmissionId(0) + + results = producer_client.stream_submission_chunks(submission_id, Strategy.Oldest()) + with pytest.raises(SubmissionNotFoundError) as exc_info: + next(results) + assert exc_info.value.submission_id == submission_id.id + + async def read_missing_chunk() -> None: + results = await producer_client.async_stream_submission_chunks( + submission_id, Strategy.Oldest() + ) + await results.__anext__() + + with pytest.raises(SubmissionNotFoundError) as exc_info: + asyncio.run(read_missing_chunk()) + assert exc_info.value.submission_id == submission_id.id def test_stream_submission_chunks_requires_oldest_strategy( From 3d3021d7bc471cd40e2da9d0cc7b8642e295b7a3 Mon Sep 17 00:00:00 2001 From: Vince van Noort Date: Fri, 9 Oct 2026 15:12:34 +0200 Subject: [PATCH 24/25] Fix: tests and lints --- libs/opsqueue_python/src/producer.rs | 2 +- libs/opsqueue_python/tests/test_roundtrip.py | 8 ++++---- 2 files changed, 5 insertions(+), 5 deletions(-) diff --git a/libs/opsqueue_python/src/producer.rs b/libs/opsqueue_python/src/producer.rs index 8180faf1..0157e553 100644 --- a/libs/opsqueue_python/src/producer.rs +++ b/libs/opsqueue_python/src/producer.rs @@ -775,7 +775,7 @@ impl From> for PyErr { StreamingChunkError::Internal(error) => CError(error).into(), StreamingChunkError::Failed(error) => CError(*error).into(), StreamingChunkError::SubmissionNotFound(error) => CError(error).into(), - error => PyException::new_err(error.to_string()), + error @ StreamingChunkError::Cancelled => PyException::new_err(error.to_string()), } } } diff --git a/libs/opsqueue_python/tests/test_roundtrip.py b/libs/opsqueue_python/tests/test_roundtrip.py index 275cd13f..7c945c96 100644 --- a/libs/opsqueue_python/tests/test_roundtrip.py +++ b/libs/opsqueue_python/tests/test_roundtrip.py @@ -1009,10 +1009,10 @@ def test_stream_submission_chunks_fails_if_submission_failed_before_read() -> No second_chunk.chunk_index, "Simulated failure", ) - assert isinstance( - producer_client.get_submission_status(submission_id), - SubmissionStatus.Failed, - ) + with pytest.raises(SubmissionFailedError): + producer_client.blocking_stream_completed_submission_chunks( + submission_id, timeout=SUBMISSION_COMPLETED_TIMEOUT + ) with pytest.raises(SubmissionFailedError): next(results) with pytest.raises(StopIteration): From b6d884c1b3f61ead560a95d9e87e66f1c1dc517d Mon Sep 17 00:00:00 2001 From: Vince van Noort Date: Fri, 9 Oct 2026 15:30:58 +0200 Subject: [PATCH 25/25] Fix: types --- libs/opsqueue_python/tests/test_roundtrip.py | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/libs/opsqueue_python/tests/test_roundtrip.py b/libs/opsqueue_python/tests/test_roundtrip.py index 7c945c96..d96fb525 100644 --- a/libs/opsqueue_python/tests/test_roundtrip.py +++ b/libs/opsqueue_python/tests/test_roundtrip.py @@ -17,11 +17,11 @@ strategy_from_description, ) from opsqueue.common import SerializationFormat -from opsqueue.consumer import ( +from opsqueue.consumer import ( # type: ignore[attr-defined] Chunk, ConsumerClient, Strategy, - opsqueue_internal, # type: ignore[attr-defined] + opsqueue_internal, ) from opsqueue.producer import ( ChunkFailed,