from typing import TYPE_CHECKING
if TYPE_CHECKING:
from uuid import UUID
__all__ = (
"JobCancelledError",
"MissingDependencyError",
"NonRetryableError",
"QueueConfigurationError",
"QueueDispatchError",
"QueueError",
"QueueEventBufferFull",
"QueueWarning",
"TaskIdentityError",
"TaskIdentityTooLargeError",
"job_cancelled",
"non_retryable",
)
[docs]
class QueueError(Exception):
"""Base exception for litestar-queues errors."""
[docs]
class QueueWarning(UserWarning):
"""Base class for litestar-queues warnings."""
[docs]
class QueueConfigurationError(QueueError):
"""Raised when queue backend configuration is invalid."""
[docs]
class QueueDispatchError(QueueError):
"""Raised when a persisted record could not be handed to its transport.
``committed`` is the part callers act on. A committed record is durable and
a repair sweep can retry its dispatch; an uncommitted one never reached
storage, so the caller owns retrying the whole enqueue.
"""
[docs]
def __init__(self, message: "str", *, task_id: "UUID", committed: "bool") -> "None":
"""Initialize dispatch error.
Args:
message: Human-readable description of the dispatch failure.
task_id: Identifier of the record that failed to dispatch.
committed: Whether the record is durably persisted.
"""
super().__init__(message)
self.task_id = task_id
self.committed = committed
[docs]
class TaskIdentityError(QueueError):
"""Raised when task uniqueness identity cannot be derived.
Signals that ``unique_by="arguments"`` was requested for a call whose bound
arguments cannot be represented by the package's canonical JSON identity
contract (for example non-finite floats or non-JSON objects). Uniqueness
identity never falls back to pickle or ``repr()``; the caller must supply an
explicit key or pass identity-friendly arguments instead.
"""
[docs]
class QueueEventBufferFull(QueueError): # noqa: N818
"""Raised when queue event buffering cannot accept another event."""
[docs]
class NonRetryableError(QueueError):
"""Raised by a task to mark the current failure as permanent."""
[docs]
class JobCancelledError(QueueError):
"""Raised by a task to cooperatively mark itself cancelled."""
[docs]
def non_retryable(message: "str") -> "None":
"""Raise a non-retryable task failure.
Raises:
NonRetryableError: Always raised with the provided message.
"""
raise NonRetryableError(message)
[docs]
def job_cancelled(message: "str" = "Task cancelled") -> "None":
"""Raise a cooperative task cancellation.
Raises:
JobCancelledError: Always raised with the provided message.
"""
raise JobCancelledError(message)
[docs]
class MissingDependencyError(QueueError, ImportError):
"""Raised when a required optional dependency is not installed."""
[docs]
def __init__(self, package: "str", install_package: "str | None" = None) -> "None":
"""Initialize missing dependency error.
Args:
package: The missing import package.
install_package: Optional package or extra to install.
"""
install_name = install_package or package
super().__init__(
f"Package {package!r} is not installed but required. Install it with "
f"'pip install litestar-queues[{install_name}]' or install {install_name!r} separately."
)