Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
4 changes: 2 additions & 2 deletions Cargo.lock

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

4 changes: 2 additions & 2 deletions Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -14,7 +14,7 @@ hakari-package = "workspace-hack"

[workspace.package]
edition = "2024"
version = "0.40.1"
version = "0.41.0"

[workspace.dependencies]
anyhow = { version = "1.0.102", default-features = false }
Expand All @@ -28,7 +28,7 @@ ciborium = { version = "0.2.2", default-features = false, features = ["std"] }
clap = { version = "4.5.60", default-features = false, features = ["derive", "std"] }
criterion = { version = "0.8", default-features = false, features = ["async_tokio"] }
crossbeam-skiplist = { version = "0.1.3", default-features = false, features = ["std"] }
dashmap = { version = "6.1.0", default-features = false, features = ["serde"]}
dashmap = { version = "6.1.0", default-features = false, features = ["serde"] }
either = { version = "1.13.0", default-features = false }
futures = { version = "0.3.32", default-features = false }
http = { version = "1.4.0", default-features = false }
Expand Down
5 changes: 4 additions & 1 deletion justfile
Original file line number Diff line number Diff line change
Expand Up @@ -113,7 +113,10 @@ semver:
set -euo pipefail
export DATABASE_URL="sqlite://{{justfile_directory()}}/opsqueue/opsqueue_example_database_schema.db"
# We select the latest git tag, not to be confused with the latest Cargo version.
cargo semver-checks --workspace --target x86_64-unknown-linux-gnu --baseline-rev "$(git tag -l --sort=-version:refname | head -1)"
cargo semver-checks --workspace --target x86_64-unknown-linux-gnu --baseline-rev "$(git tag -l --sort=-version:refname | head -1)" || {
echo "Semver checks failed. Please bump the version in Cargo.toml"
exit 1
}
Comment thread
SemMulder marked this conversation as resolved.

# Rust static analysis
[group('lint')]
Expand Down
9 changes: 0 additions & 9 deletions libs/opsqueue_python/python/opsqueue/exceptions.py
Original file line number Diff line number Diff line change
Expand Up @@ -92,15 +92,6 @@ class TryFromIntError(IncorrectUsageError):
pass


class ChunkNotFoundError(IncorrectUsageError):

@SemMulder SemMulder Jul 17, 2026

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

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

Note that ChunkNotFoundError can be dropped since it is (and was) dead code.

"""
Raised when a method is used to look up information about a chunk
but the chunk doesn't exist within the Opsqueue.
"""

pass


class SubmissionNotFoundError(IncorrectUsageError):
"""
Raised when a method is used to look up information about a submission
Expand Down
56 changes: 42 additions & 14 deletions libs/opsqueue_python/python/opsqueue/producer.py
Original file line number Diff line number Diff line change
@@ -1,46 +1,49 @@
from __future__ import annotations
from collections.abc import Iterable, Iterator, AsyncIterator
from typing import Any, cast

import itertools
from collections.abc import Iterable, Iterator, AsyncIterator
from typing import Any, cast

from opentelemetry import trace

from opsqueue.common import (
SerializationFormat,
encode_chunk,
decode_chunk,
DEFAULT_SERIALIZATION_FORMAT,
)
from . import opsqueue_internal
from . import tracing
from opsqueue.exceptions import (
SubmissionFailedError,
SubmissionNotCancellableError,
SubmissionNotFoundError,
TooManyMatchingSubmissionsError,
)
from . import opsqueue_internal
from . import tracing
from .opsqueue_internal import ( # type: ignore[import-not-found]
SubmissionId,
SubmissionStatus,
SubmissionCompleted,
SubmissionFailed,
ChunkFailed,
SubmissionNotCancellable,
SubmissionPaused,
InitialSubmissionStatus,
)

__all__ = [
"ChunkFailed",
"InitialSubmissionStatus",
"ProducerClient",
"SubmissionId",
"SubmissionStatus",
"SubmissionCompleted",
"SubmissionFailedError",
"SubmissionFailed",
"SubmissionFailedError",
"SubmissionId",
"SubmissionNotCancellable",
"SubmissionNotCancellableError",
"SubmissionNotFoundError",
"SubmissionPaused",
"SubmissionStatus",
"TooManyMatchingSubmissionsError",
"ChunkFailed",
]


Expand Down Expand Up @@ -96,6 +99,7 @@ def run_submission(
serialization_format: SerializationFormat = DEFAULT_SERIALIZATION_FORMAT,
metadata: None | bytes = None,
strategic_metadata: None | dict[str, int] = None,
timeout: float | None = None,
) -> Iterator[Any]:
"""
Inserts a submission into the queue, and blocks until it is completed.
Expand All @@ -116,6 +120,7 @@ def run_submission(
metadata=metadata,
strategic_metadata=strategic_metadata,
chunk_size=chunk_size,
timeout=timeout,
)
return _unchunk_iterator(results_iter, serialization_format)

Expand Down Expand Up @@ -146,6 +151,7 @@ def insert_submission(
serialization_format: SerializationFormat = DEFAULT_SERIALIZATION_FORMAT,
metadata: None | bytes = None,
strategic_metadata: None | dict[str, int] = None,
initial_status: InitialSubmissionStatus = InitialSubmissionStatus.InProgress,
) -> SubmissionId:
"""
Inserts a submission into the queue,
Expand All @@ -162,13 +168,15 @@ def insert_submission(
metadata=metadata,
strategic_metadata=strategic_metadata,
chunk_size=chunk_size,
initial_status=initial_status,
)

def blocking_stream_completed_submission(
self,
submission_id: SubmissionId,
*,
serialization_format: SerializationFormat = DEFAULT_SERIALIZATION_FORMAT,
timeout: float | None = None,
) -> Iterator[Any]:
"""
Blocks until the submission is completed.
Expand All @@ -181,7 +189,7 @@ def blocking_stream_completed_submission(
(after retrying a consumer kept failing on one of the chunks)
"""
return _unchunk_iterator(
self.blocking_stream_completed_submission_chunks(submission_id),
self.blocking_stream_completed_submission_chunks(submission_id, timeout),
serialization_format,
)

Expand Down Expand Up @@ -211,6 +219,7 @@ def run_submission_chunks(
metadata: None | bytes = None,
strategic_metadata: None | dict[str, int] = None,
chunk_size: None | int = None,
timeout: float | None = None,
) -> Iterator[bytes]:
"""
Inserts an already-chunked submission into the queue, and blocks until it is completed.
Expand All @@ -229,7 +238,7 @@ def run_submission_chunks(
strategic_metadata=strategic_metadata,
chunk_size=chunk_size,
)
return self.blocking_stream_completed_submission_chunks(submission_id)
return self.blocking_stream_completed_submission_chunks(submission_id, timeout)

async def async_run_submission_chunks(
self,
Expand Down Expand Up @@ -259,6 +268,7 @@ def insert_submission_chunks(
metadata: None | bytes = None,
strategic_metadata: None | dict[str, int] = None,
chunk_size: None | int = None,
initial_status: InitialSubmissionStatus = InitialSubmissionStatus.InProgress,
) -> SubmissionId:
"""
Inserts an already-chunked submission into the queue,
Expand All @@ -275,10 +285,13 @@ def insert_submission_chunks(
strategic_metadata=strategic_metadata,
chunk_size=chunk_size,
otel_trace_carrier=otel_trace_carrier,
initial_status=initial_status,
)

def blocking_stream_completed_submission_chunks(
self, submission_id: SubmissionId
self,
submission_id: SubmissionId,
timeout: float | None = None,
) -> Iterator[bytes]:
"""
Blocks until the submission is completed, and returns an iterator that lazily
Expand All @@ -289,7 +302,9 @@ def blocking_stream_completed_submission_chunks(
- `SubmissionFailedError` if the submission failed permanently
(after retrying a consumer kept failing on one of the chunks)
"""
return self.inner.blocking_stream_completed_submission_chunks(submission_id) # type: ignore[no-any-return]
return self.inner.blocking_stream_completed_submission_chunks( # type: ignore[no-any-return]
submission_id, timeout
)

async def async_stream_completed_submission_chunks(
self, submission_id: SubmissionId
Expand Down Expand Up @@ -326,7 +341,7 @@ def count_submissions(self) -> int:

def cancel_submission(self, submission_id: SubmissionId) -> None:
"""
Cancel a specific submission that is in progress.
Cancel a specific submission that is in progress or paused.

Returns None if the submission was successfully cancelled.

Expand All @@ -337,6 +352,19 @@ def cancel_submission(self, submission_id: SubmissionId) -> None:
"""
self.inner.cancel_submission(submission_id)

def unpause_submission(self, submission_id: SubmissionId) -> None:
"""
Unpause a specific submission that is currently paused,
making it available to consumers.

Returns None if the submission was successfully unpaused.

Raises:
- `SubmissionNotFoundError` if the submission is not currently paused.
- `InternalProducerClientError` if there is a low-level internal error.
"""
self.inner.unpause_submission(submission_id)

def get_submission_status(
self, submission_id: SubmissionId
) -> SubmissionStatus | None:
Expand Down
62 changes: 61 additions & 1 deletion libs/opsqueue_python/src/common.rs
Original file line number Diff line number Diff line change
Expand Up @@ -350,6 +350,24 @@ impl From<opsqueue::common::submission::SubmissionCancelled> for SubmissionCance
}
}

#[pyclass(from_py_object, eq, eq_int)]
#[derive(Default, Debug, Clone, PartialEq, Eq)]
pub enum InitialSubmissionStatus {
Paused,
#[default]
InProgress,
}

impl From<InitialSubmissionStatus> for opsqueue::common::submission::InitialSubmissionStatus {
fn from(value: InitialSubmissionStatus) -> Self {
use opsqueue::common::submission::InitialSubmissionStatus::{InProgress, Paused};
match value {
InitialSubmissionStatus::Paused => Paused,
InitialSubmissionStatus::InProgress => InProgress,
}
}
}

#[pyclass(from_py_object, frozen, module = "opsqueue")]
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum SubmissionStatus {
Expand All @@ -366,12 +384,15 @@ pub enum SubmissionStatus {
Cancelled {
submission: SubmissionCancelled,
},
Paused {
submission: SubmissionPaused,
},
}

impl From<opsqueue::common::submission::SubmissionStatus> for SubmissionStatus {
fn from(value: opsqueue::common::submission::SubmissionStatus) -> Self {
use opsqueue::common::submission::SubmissionStatus::{
Cancelled, Completed, Failed, InProgress,
Cancelled, Completed, Failed, InProgress, Paused,
};
match value {
InProgress(s) => SubmissionStatus::InProgress {
Expand All @@ -388,6 +409,9 @@ impl From<opsqueue::common::submission::SubmissionStatus> for SubmissionStatus {
Cancelled(s) => SubmissionStatus::Cancelled {
submission: s.into(),
},
Paused(s) => SubmissionStatus::Paused {
submission: s.into(),
},
}
}
}
Expand Down Expand Up @@ -512,6 +536,42 @@ pub struct SubmissionCancelled {
pub cancelled_at: DateTime<Utc>,
}

#[pyclass(from_py_object, frozen, get_all, module = "opsqueue")]
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct SubmissionPaused {
pub id: SubmissionId,
pub chunks_total: u64,
pub chunks_done: u64,
pub metadata: Option<submission::Metadata>,
pub strategic_metadata: StrategicMetadataMap,
}

impl From<opsqueue::common::submission::SubmissionPaused> for SubmissionPaused {
fn from(value: opsqueue::common::submission::SubmissionPaused) -> Self {
Self {
id: value.id.into(),
chunks_total: value.chunks_total.into(),
chunks_done: value.chunks_done.into(),
metadata: value.metadata,
strategic_metadata: value.strategic_metadata,
}
}
}

#[pymethods]
impl SubmissionPaused {
fn __repr__(&self) -> String {
format!(
"SubmissionPaused(id={0}, chunks_total={1}, chunks_done={2}, metadata={3:?}, strategic_metadata={4:?})",
self.id.__repr__(),
self.chunks_total,
self.chunks_done,
self.metadata,
self.strategic_metadata
)
}
}

/// Submission could not be cancelled because it was already completed, failed
/// or cancelled.
#[pyclass(from_py_object, frozen, module = "opsqueue")]
Expand Down
Loading