Skip to content

Stream Queues

StreamOffset dataclass

Stream queue offset specification.

Use class methods to create

StreamOffset.first() StreamOffset.last() StreamOffset.next() StreamOffset.offset(42) StreamOffset.timestamp(datetime(...))

Source code in src/rabbitkit/streams.py
@dataclass(frozen=True, slots=True)
class StreamOffset:
    """Stream queue offset specification.

    Use class methods to create:
        StreamOffset.first()
        StreamOffset.last()
        StreamOffset.next()
        StreamOffset.offset(42)
        StreamOffset.timestamp(datetime(...))
    """

    type: StreamOffsetType = StreamOffsetType.NEXT
    value: int | datetime | None = None

    @classmethod
    def first(cls) -> StreamOffset:
        """Start from the beginning of the stream."""
        return cls(type=StreamOffsetType.FIRST)

    @classmethod
    def last(cls) -> StreamOffset:
        """Start from the end (new messages only)."""
        return cls(type=StreamOffsetType.LAST)

    @classmethod
    def next(cls) -> StreamOffset:
        """Start from the next unconsumed message (default)."""
        return cls(type=StreamOffsetType.NEXT)

    @classmethod
    def offset(cls, value: int) -> StreamOffset:
        """Start from a specific numeric offset."""
        if value < 0:
            msg = "Stream offset must be non-negative"
            raise ValueError(msg)
        return cls(type=StreamOffsetType.OFFSET, value=value)

    @classmethod
    def timestamp(cls, value: datetime) -> StreamOffset:
        """Start from messages published after the given timestamp."""
        return cls(type=StreamOffsetType.TIMESTAMP, value=value)

    def to_consume_arguments(self) -> dict[str, Any]:
        """Convert to RabbitMQ consume arguments (x-stream-offset).

        Returns dict suitable for merging into basic_consume arguments.
        """
        if self.type == StreamOffsetType.FIRST:
            return {"x-stream-offset": "first"}
        if self.type == StreamOffsetType.LAST:
            return {"x-stream-offset": "last"}
        if self.type == StreamOffsetType.NEXT:
            return {"x-stream-offset": "next"}
        if self.type == StreamOffsetType.OFFSET:
            return {"x-stream-offset": self.value}
        if self.type == StreamOffsetType.TIMESTAMP:
            if not isinstance(self.value, datetime):
                raise TypeError("TIMESTAMP stream offset requires a datetime value")
            # RabbitMQ x-stream-offset by time expects a Unix timestamp in seconds.
            return {"x-stream-offset": int(self.value.timestamp())}
        return {}  # pragma: no cover

    def __repr__(self) -> str:
        if self.value is not None:
            return f"StreamOffset({self.type.value}={self.value})"
        return f"StreamOffset({self.type.value})"

Methods:

first() -> StreamOffset classmethod

Start from the beginning of the stream.

Source code in src/rabbitkit/streams.py
@classmethod
def first(cls) -> StreamOffset:
    """Start from the beginning of the stream."""
    return cls(type=StreamOffsetType.FIRST)

last() -> StreamOffset classmethod

Start from the end (new messages only).

Source code in src/rabbitkit/streams.py
@classmethod
def last(cls) -> StreamOffset:
    """Start from the end (new messages only)."""
    return cls(type=StreamOffsetType.LAST)

next() -> StreamOffset classmethod

Start from the next unconsumed message (default).

Source code in src/rabbitkit/streams.py
@classmethod
def next(cls) -> StreamOffset:
    """Start from the next unconsumed message (default)."""
    return cls(type=StreamOffsetType.NEXT)

offset(value: int) -> StreamOffset classmethod

Start from a specific numeric offset.

Source code in src/rabbitkit/streams.py
@classmethod
def offset(cls, value: int) -> StreamOffset:
    """Start from a specific numeric offset."""
    if value < 0:
        msg = "Stream offset must be non-negative"
        raise ValueError(msg)
    return cls(type=StreamOffsetType.OFFSET, value=value)

timestamp(value: datetime) -> StreamOffset classmethod

Start from messages published after the given timestamp.

Source code in src/rabbitkit/streams.py
@classmethod
def timestamp(cls, value: datetime) -> StreamOffset:
    """Start from messages published after the given timestamp."""
    return cls(type=StreamOffsetType.TIMESTAMP, value=value)

to_consume_arguments() -> dict[str, Any]

Convert to RabbitMQ consume arguments (x-stream-offset).

Returns dict suitable for merging into basic_consume arguments.

Source code in src/rabbitkit/streams.py
def to_consume_arguments(self) -> dict[str, Any]:
    """Convert to RabbitMQ consume arguments (x-stream-offset).

    Returns dict suitable for merging into basic_consume arguments.
    """
    if self.type == StreamOffsetType.FIRST:
        return {"x-stream-offset": "first"}
    if self.type == StreamOffsetType.LAST:
        return {"x-stream-offset": "last"}
    if self.type == StreamOffsetType.NEXT:
        return {"x-stream-offset": "next"}
    if self.type == StreamOffsetType.OFFSET:
        return {"x-stream-offset": self.value}
    if self.type == StreamOffsetType.TIMESTAMP:
        if not isinstance(self.value, datetime):
            raise TypeError("TIMESTAMP stream offset requires a datetime value")
        # RabbitMQ x-stream-offset by time expects a Unix timestamp in seconds.
        return {"x-stream-offset": int(self.value.timestamp())}
    return {}  # pragma: no cover

StreamOffsetType

Bases: str, Enum

Stream offset specification types.

Source code in src/rabbitkit/streams.py
class StreamOffsetType(str, enum.Enum):
    """Stream offset specification types."""

    FIRST = "first"
    LAST = "last"
    NEXT = "next"
    OFFSET = "offset"
    TIMESTAMP = "timestamp"

StreamConsumerConfig dataclass

Configuration for stream queue consumers.

Extends basic consumer behavior with stream-specific options.

Source code in src/rabbitkit/streams.py
@dataclass(frozen=True, slots=True)
class StreamConsumerConfig:
    """Configuration for stream queue consumers.

    Extends basic consumer behavior with stream-specific options.
    """

    offset: StreamOffset = field(default_factory=StreamOffset.next)
    consumer_name: str | None = None  # x-stream-consumer-name for single-active-consumer

    def to_consume_arguments(self) -> dict[str, Any]:
        """Build consume arguments for stream queue subscription."""
        args = self.offset.to_consume_arguments()
        if self.consumer_name is not None:
            args["x-stream-consumer-name"] = self.consumer_name
        return args

Methods:

to_consume_arguments() -> dict[str, Any]

Build consume arguments for stream queue subscription.

Source code in src/rabbitkit/streams.py
def to_consume_arguments(self) -> dict[str, Any]:
    """Build consume arguments for stream queue subscription."""
    args = self.offset.to_consume_arguments()
    if self.consumer_name is not None:
        args["x-stream-consumer-name"] = self.consumer_name
    return args