Skip to content

BaseBroker

Message broker interface. All broker adapters (Inline, Redis Stream, Redis PubSub, etc.) implement this contract.

See Broker Adapters for concrete adapter configuration.

Base class for all broker implementations.

Provides shared behavior (UoW integration, connection recovery, capability checks, subscriber registration) via the Template Method pattern. Subclasses implement the abstract _underscore methods for broker-specific logic.

Source code in src/protean/port/broker.py
200
201
202
203
204
205
206
207
208
209
210
211
212
def __init__(
    self, name: str, domain: "Domain", conn_info: dict[str, str | bool]
) -> None:
    self.name = name
    self.domain = domain
    self.conn_info = conn_info

    self._subscribers: defaultdict[str, set[type[BaseSubscriber]]] = defaultdict(
        set
    )
    self._last_ping_time: float | None = None
    self._last_ping_success: bool | None = None
    self._start_time = time.time()

capabilities abstractmethod property

capabilities: BrokerCapabilities

Return the capabilities of this broker implementation.

RETURNS DESCRIPTION
BrokerCapabilities

The capabilities supported by this broker

TYPE: BrokerCapabilities

has_capability

has_capability(capability: BrokerCapabilities) -> bool

Check if broker has a specific capability.

PARAMETER DESCRIPTION
capability

The capability to check for

TYPE: BrokerCapabilities

RETURNS DESCRIPTION
bool

True if the broker has the capability, False otherwise

TYPE: bool

Source code in src/protean/port/broker.py
223
224
225
226
227
228
229
230
231
232
def has_capability(self, capability: BrokerCapabilities) -> bool:
    """Check if broker has a specific capability.

    Args:
        capability: The capability to check for

    Returns:
        bool: True if the broker has the capability, False otherwise
    """
    return capability in self.capabilities

has_all_capabilities

has_all_capabilities(
    capabilities: BrokerCapabilities,
) -> bool

Check if broker has all the specified capabilities.

PARAMETER DESCRIPTION
capabilities

The capabilities to check for

TYPE: BrokerCapabilities

RETURNS DESCRIPTION
bool

True if the broker has all capabilities, False otherwise

TYPE: bool

Source code in src/protean/port/broker.py
234
235
236
237
238
239
240
241
242
243
def has_all_capabilities(self, capabilities: BrokerCapabilities) -> bool:
    """Check if broker has all the specified capabilities.

    Args:
        capabilities: The capabilities to check for

    Returns:
        bool: True if the broker has all capabilities, False otherwise
    """
    return (self.capabilities & capabilities) == capabilities

has_any_capability

has_any_capability(
    capabilities: BrokerCapabilities,
) -> bool

Check if broker has any of the specified capabilities.

PARAMETER DESCRIPTION
capabilities

The capabilities to check for

TYPE: BrokerCapabilities

RETURNS DESCRIPTION
bool

True if the broker has any of the capabilities, False otherwise

TYPE: bool

Source code in src/protean/port/broker.py
245
246
247
248
249
250
251
252
253
254
def has_any_capability(self, capabilities: BrokerCapabilities) -> bool:
    """Check if broker has any of the specified capabilities.

    Args:
        capabilities: The capabilities to check for

    Returns:
        bool: True if the broker has any of the capabilities, False otherwise
    """
    return bool(self.capabilities & capabilities)

publish

publish(stream: str, message: dict[str, Any]) -> str | None

Publish a message to the broker.

PARAMETER DESCRIPTION
stream

The stream to which the message should be published

TYPE: str

message

The message payload to be published

TYPE: dict

RETURNS DESCRIPTION
str

The identifier of the message. The content of the identifier is broker-specific.

TYPE: str | None

str | None

All brokers are guaranteed to provide message identifiers.

RAISES DESCRIPTION
ValidationError

If message is an empty dict

Source code in src/protean/port/broker.py
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
def publish(self, stream: str, message: dict[str, Any]) -> str | None:
    """Publish a message to the broker.

    Args:
        stream (str): The stream to which the message should be published
        message (dict): The message payload to be published

    Returns:
        str: The identifier of the message. The content of the identifier is broker-specific.
        All brokers are guaranteed to provide message identifiers.

    Raises:
        ValidationError: If message is an empty dict
    """
    if not message:
        raise ValidationError({"message": ["Message cannot be empty"]})

    if current_uow:
        logger.debug(f"Recording message {message} in {current_uow} for dispatch")

        current_uow.register_message(stream, message, broker_name=self.name)
        return None
    else:
        try:
            identifier = self._publish(stream, message)
        except Exception as e:
            # Check if this is a connection-related error and attempt recovery
            if self._is_connection_error(e):
                logger.warning(f"Connection error during publish: {e}")
                if self._ensure_connection():
                    # Retry the operation once after reconnection
                    identifier = self._publish(stream, message)
                else:
                    raise
            else:
                raise

        if (
            self.domain.config["message_processing"] == Processing.SYNC.value
            and self._subscribers[stream]
        ):
            for subscriber_cls in self._subscribers[stream]:
                subscriber = subscriber_cls()
                subscriber(message)

        return identifier

ping

ping() -> bool

Test broker connectivity.

RETURNS DESCRIPTION
bool

True if broker is reachable and responsive, False otherwise

TYPE: bool

Source code in src/protean/port/broker.py
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
def ping(self) -> bool:
    """Test broker connectivity.

    Returns:
        bool: True if broker is reachable and responsive, False otherwise
    """
    try:
        start_time = time.time()
        result = self._ping()
        self._last_ping_time = time.time() - start_time
        self._last_ping_success = result
        return result
    except Exception as e:  # noqa: BLE001 - health probe: any adapter fault is "down"
        logger.debug(f"Ping failed for broker {self.name}: {e}")
        self._last_ping_time = None
        self._last_ping_success = False
        return False

health_stats

health_stats() -> dict[str, Any]

Get comprehensive health statistics for the broker.

RETURNS DESCRIPTION
dict

Health statistics with the following structure: { 'status': 'healthy' | 'degraded' | 'unhealthy', 'connected': bool, 'last_ping_ms': float | None, 'uptime_seconds': float, 'details': dict # Broker-specific details }

TYPE: dict[str, Any]

Source code in src/protean/port/broker.py
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
def health_stats(self) -> dict[str, Any]:
    """Get comprehensive health statistics for the broker.

    Returns:
        dict: Health statistics with the following structure:
            {
                'status': 'healthy' | 'degraded' | 'unhealthy',
                'connected': bool,
                'last_ping_ms': float | None,
                'uptime_seconds': float,
                'details': dict  # Broker-specific details
            }
    """
    try:
        # Get broker-specific health details
        broker_details = self._health_stats()

        # Perform a fresh ping to get current connectivity status
        is_connected = self.ping()

        # Determine overall health status
        if is_connected and broker_details.get("healthy", True):
            status = "healthy"
        elif is_connected:
            status = "degraded"  # Connected but some issues reported
        else:
            status = "unhealthy"

        # Calculate uptime since broker initialization
        uptime_seconds = time.time() - self._start_time

        return {
            "status": status,
            "connected": is_connected,
            "last_ping_ms": self._last_ping_time * 1000
            if self._last_ping_time is not None
            else None,
            "uptime_seconds": uptime_seconds,
            "details": broker_details,
        }
    except Exception as e:  # noqa: BLE001 - health probe: report the fault as unhealthy
        logger.error(f"Error gathering health stats for broker {self.name}: {e}")
        return {
            "status": "unhealthy",
            "connected": False,
            "last_ping_ms": None,
            "uptime_seconds": 0,
            "details": {"error": str(e)},
        }

ensure_connection

ensure_connection() -> bool

Ensure broker connection is healthy, attempt reconnection if needed.

This method can be called explicitly or is triggered automatically when connection-related exceptions are encountered.

RETURNS DESCRIPTION
bool

True if connection is healthy/restored, False otherwise

TYPE: bool

Source code in src/protean/port/broker.py
371
372
373
374
375
376
377
378
379
380
def ensure_connection(self) -> bool:
    """Ensure broker connection is healthy, attempt reconnection if needed.

    This method can be called explicitly or is triggered automatically
    when connection-related exceptions are encountered.

    Returns:
        bool: True if connection is healthy/restored, False otherwise
    """
    return self._ensure_connection()

get_next

get_next(
    stream: str, consumer_group: str
) -> tuple[str, dict[str, Any]] | None

Retrieve the next message to process from broker.

PARAMETER DESCRIPTION
stream

The stream from which to retrieve the message

TYPE: str

consumer_group

The consumer group identifier

TYPE: str

RETURNS DESCRIPTION
tuple[str, dict[str, Any]] | None

tuple[str, dict] | None: A tuple of (identifier, message) or None if no messages available

Source code in src/protean/port/broker.py
448
449
450
451
452
453
454
455
456
457
458
459
460
461
462
463
464
465
466
467
468
469
470
471
472
473
474
475
476
477
def get_next(
    self, stream: str, consumer_group: str
) -> tuple[str, dict[str, Any]] | None:
    """Retrieve the next message to process from broker.

    Args:
        stream (str): The stream from which to retrieve the message
        consumer_group (str): The consumer group identifier

    Returns:
        tuple[str, dict] | None: A tuple of (identifier, message) or None if no messages available
    """
    # Check if broker supports consumer groups
    if not self.has_capability(BrokerCapabilities.CONSUMER_GROUPS):
        logger.warning(f"Broker {self.name} does not support consumer groups")
        return None

    try:
        return self._get_next(stream, consumer_group)
    except Exception as e:
        # Check if this is a connection-related error and attempt recovery
        if self._is_connection_error(e):
            logger.warning(f"Connection error during get_next: {e}")
            if self._ensure_connection():
                # Retry the operation once after reconnection
                return self._get_next(stream, consumer_group)
            else:
                raise
        else:
            raise

read

read(
    stream: str, consumer_group: str, no_of_messages: int
) -> list[tuple[str, dict[str, Any]]]

Read messages from the broker.

PARAMETER DESCRIPTION
stream

The stream from which to read messages

TYPE: str

consumer_group

The consumer group identifier

TYPE: str

no_of_messages

The number of messages to read

TYPE: int

RETURNS DESCRIPTION
list[tuple[str, dict[str, Any]]]

list[tuple[str, dict]]: The list of (identifier, message) tuples

Source code in src/protean/port/broker.py
485
486
487
488
489
490
491
492
493
494
495
496
497
498
499
500
501
502
503
504
505
506
507
508
509
510
511
512
513
514
515
def read(
    self, stream: str, consumer_group: str, no_of_messages: int
) -> list[tuple[str, dict[str, Any]]]:
    """Read messages from the broker.

    Args:
        stream (str): The stream from which to read messages
        consumer_group (str): The consumer group identifier
        no_of_messages (int): The number of messages to read

    Returns:
        list[tuple[str, dict]]: The list of (identifier, message) tuples
    """
    # Check if broker supports consumer groups
    if not self.has_capability(BrokerCapabilities.CONSUMER_GROUPS):
        logger.warning(f"Broker {self.name} does not support consumer groups")
        return []

    try:
        return self._read(stream, consumer_group, no_of_messages)
    except Exception as e:
        # Check if this is a connection-related error and attempt recovery
        if self._is_connection_error(e):
            logger.warning(f"Connection error during read: {e}")
            if self._ensure_connection():
                # Retry the operation once after reconnection
                return self._read(stream, consumer_group, no_of_messages)
            else:
                raise
        else:
            raise

ack

ack(
    stream: str, identifier: str, consumer_group: str
) -> bool

Acknowledge successful processing of a message.

PARAMETER DESCRIPTION
stream

The stream from which the message was received

TYPE: str

identifier

The unique identifier of the message to acknowledge

TYPE: str

consumer_group

The consumer group that processed the message

TYPE: str

RETURNS DESCRIPTION
bool

True if the message was successfully acknowledged, False otherwise

TYPE: bool

Source code in src/protean/port/broker.py
517
518
519
520
521
522
523
524
525
526
527
528
529
530
531
532
533
534
535
536
537
538
539
540
541
542
543
544
545
546
547
def ack(self, stream: str, identifier: str, consumer_group: str) -> bool:
    """Acknowledge successful processing of a message.

    Args:
        stream (str): The stream from which the message was received
        identifier (str): The unique identifier of the message to acknowledge
        consumer_group (str): The consumer group that processed the message

    Returns:
        bool: True if the message was successfully acknowledged, False otherwise
    """
    # Check if broker supports acknowledgment
    if not self.has_capability(BrokerCapabilities.ACK_NACK):
        logger.warning(
            f"Broker {self.name} does not support message acknowledgment"
        )
        return False

    try:
        return self._ack(stream, identifier, consumer_group)
    except Exception as e:
        # Check if this is a connection-related error and attempt recovery
        if self._is_connection_error(e):
            logger.warning(f"Connection error during ack: {e}")
            if self._ensure_connection():
                # Retry the operation once after reconnection
                return self._ack(stream, identifier, consumer_group)
            else:
                raise
        else:
            raise

nack

nack(
    stream: str, identifier: str, consumer_group: str
) -> bool

Negative acknowledge - mark message for reprocessing.

PARAMETER DESCRIPTION
stream

The stream from which the message was received

TYPE: str

identifier

The unique identifier of the message to nack

TYPE: str

consumer_group

The consumer group that failed to process the message

TYPE: str

RETURNS DESCRIPTION
bool

True if the message was successfully marked for reprocessing, False otherwise

TYPE: bool

Source code in src/protean/port/broker.py
549
550
551
552
553
554
555
556
557
558
559
560
561
562
563
564
565
566
567
568
569
570
571
572
573
574
575
576
577
578
579
def nack(self, stream: str, identifier: str, consumer_group: str) -> bool:
    """Negative acknowledge - mark message for reprocessing.

    Args:
        stream (str): The stream from which the message was received
        identifier (str): The unique identifier of the message to nack
        consumer_group (str): The consumer group that failed to process the message

    Returns:
        bool: True if the message was successfully marked for reprocessing, False otherwise
    """
    # Check if broker supports negative acknowledgment
    if not self.has_capability(BrokerCapabilities.ACK_NACK):
        logger.warning(
            f"Broker {self.name} does not support message negative acknowledgment"
        )
        return False

    try:
        return self._nack(stream, identifier, consumer_group)
    except Exception as e:
        # Check if this is a connection-related error and attempt recovery
        if self._is_connection_error(e):
            logger.warning(f"Connection error during nack: {e}")
            if self._ensure_connection():
                # Retry the operation once after reconnection
                return self._nack(stream, identifier, consumer_group)
            else:
                raise
        else:
            raise

read_blocking

read_blocking(
    stream: str,
    consumer_group: str,
    consumer_name: str,
    timeout_ms: int = 5000,
    count: int = 1,
) -> list[tuple[str, dict[str, Any]]]

Read messages from the broker using blocking mode.

This is an optional method that brokers can implement to support efficient blocking reads for stream-based subscriptions.

PARAMETER DESCRIPTION
stream

The stream from which to read messages

TYPE: str

consumer_group

The consumer group identifier

TYPE: str

consumer_name

The unique consumer name within the group

TYPE: str

timeout_ms

Timeout in milliseconds to wait for messages (0 = return immediately, without waiting)

TYPE: int DEFAULT: 5000

count

Maximum number of messages to read

TYPE: int DEFAULT: 1

RETURNS DESCRIPTION
list[tuple[str, dict[str, Any]]]

list[tuple[str, dict]]: The list of (identifier, message) tuples

Source code in src/protean/port/broker.py
596
597
598
599
600
601
602
603
604
605
606
607
608
609
610
611
612
613
614
615
616
617
618
619
620
621
622
623
624
625
626
627
628
629
630
631
632
633
634
635
636
637
638
639
640
def read_blocking(
    self,
    stream: str,
    consumer_group: str,
    consumer_name: str,
    timeout_ms: int = 5000,
    count: int = 1,
) -> list[tuple[str, dict[str, Any]]]:
    """Read messages from the broker using blocking mode.

    This is an optional method that brokers can implement to support
    efficient blocking reads for stream-based subscriptions.

    Args:
        stream (str): The stream from which to read messages
        consumer_group (str): The consumer group identifier
        consumer_name (str): The unique consumer name within the group
        timeout_ms (int): Timeout in milliseconds to wait for messages (0 = return immediately, without waiting)
        count (int): Maximum number of messages to read

    Returns:
        list[tuple[str, dict]]: The list of (identifier, message) tuples
    """
    # Check if broker supports blocking reads
    if not self.has_capability(BrokerCapabilities.BLOCKING_READ):
        # Fall back to regular read for brokers that don't support blocking
        return self._read(stream, consumer_group, count)

    try:
        return self._read_blocking(
            stream, consumer_group, consumer_name, timeout_ms, count
        )
    except Exception as e:
        # Check if this is a connection-related error and attempt recovery
        if self._is_connection_error(e):
            logger.warning(f"Connection error during read_blocking: {e}")
            if self._ensure_connection():
                # Retry the operation once after reconnection
                return self._read_blocking(
                    stream, consumer_group, consumer_name, timeout_ms, count
                )
            else:
                raise
        else:
            raise

read_blocking_streams

read_blocking_streams(
    streams: Sequence[str],
    consumer_group: str,
    consumer_name: str,
    timeout_ms: int = 5000,
    count: int = 1,
) -> dict[str, list[tuple[str, dict[str, Any]]]]

Read messages from several streams in one call.

The result has one key per requested stream, in the order given, with an empty list for a stream that returned nothing. The default implementation waits only on the last stream. Brokers that can wait on all of them at once, such as Redis, override _read_blocking_streams. An empty streams returns an empty dict without reading.

PARAMETER DESCRIPTION
streams

The streams to read from, in priority order

TYPE: Sequence[str]

consumer_group

The consumer group identifier

TYPE: str

consumer_name

The unique consumer name within the group

TYPE: str

timeout_ms

Longest time to wait for messages, in milliseconds (0 = return immediately, without waiting)

TYPE: int DEFAULT: 5000

count

Maximum number of messages to read from each stream

TYPE: int DEFAULT: 1

RETURNS DESCRIPTION
dict[str, list[tuple[str, dict[str, Any]]]]

dict[str, list[tuple[str, dict]]]: The (identifier, message) tuples

dict[str, list[tuple[str, dict[str, Any]]]]

read from each stream, keyed by stream name

Source code in src/protean/port/broker.py
642
643
644
645
646
647
648
649
650
651
652
653
654
655
656
657
658
659
660
661
662
663
664
665
666
667
668
669
670
671
672
673
674
675
676
677
678
679
680
681
682
683
684
685
686
687
688
689
690
691
692
693
694
695
696
def read_blocking_streams(
    self,
    streams: Sequence[str],
    consumer_group: str,
    consumer_name: str,
    timeout_ms: int = 5000,
    count: int = 1,
) -> dict[str, list[tuple[str, dict[str, Any]]]]:
    """Read messages from several streams in one call.

    The result has one key per requested stream, in the order given, with
    an empty list for a stream that returned nothing. The default
    implementation waits only on the last stream. Brokers that can wait on
    all of them at once, such as Redis, override ``_read_blocking_streams``. An empty ``streams``
    returns an empty dict without reading.

    Args:
        streams (Sequence[str]): The streams to read from, in priority order
        consumer_group (str): The consumer group identifier
        consumer_name (str): The unique consumer name within the group
        timeout_ms (int): Longest time to wait for messages, in milliseconds
            (0 = return immediately, without waiting)
        count (int): Maximum number of messages to read from each stream

    Returns:
        dict[str, list[tuple[str, dict]]]: The (identifier, message) tuples
        read from each stream, keyed by stream name
    """
    if not streams:
        return {}

    # Check if broker supports blocking reads
    if not self.has_capability(BrokerCapabilities.BLOCKING_READ):
        # Fall back to a regular read on each stream in turn
        return {
            stream: self._read(stream, consumer_group, count) for stream in streams
        }

    try:
        return self._read_blocking_streams(
            streams, consumer_group, consumer_name, timeout_ms, count
        )
    except Exception as e:
        # Check if this is a connection-related error and attempt recovery
        if self._is_connection_error(e):
            logger.warning(f"Connection error during read_blocking_streams: {e}")
            if self._ensure_connection():
                # Retry the operation once after reconnection
                return self._read_blocking_streams(
                    streams, consumer_group, consumer_name, timeout_ms, count
                )
            else:
                raise
        else:
            raise

info

info() -> dict[str, Any]

Get information about consumer groups and consumers in each group.

RETURNS DESCRIPTION
dict

Information about consumer groups and their consumers

TYPE: dict[str, Any]

Source code in src/protean/port/broker.py
808
809
810
811
812
813
814
def info(self) -> dict[str, Any]:
    """Get information about consumer groups and consumers in each group.

    Returns:
        dict: Information about consumer groups and their consumers
    """
    return self._info()

dlq_list

dlq_list(
    dlq_streams: list[str], limit: int = 100
) -> list[DLQEntry]

List DLQ messages across specified DLQ streams.

PARAMETER DESCRIPTION
dlq_streams

List of DLQ stream names to query.

TYPE: list[str]

limit

Maximum total messages to return.

TYPE: int DEFAULT: 100

RETURNS DESCRIPTION
list[DLQEntry]

List of DLQEntry objects sorted by failure time (newest first).

Source code in src/protean/port/broker.py
828
829
830
831
832
833
834
835
836
837
838
839
840
841
def dlq_list(self, dlq_streams: list[str], limit: int = 100) -> list[DLQEntry]:
    """List DLQ messages across specified DLQ streams.

    Args:
        dlq_streams: List of DLQ stream names to query.
        limit: Maximum total messages to return.

    Returns:
        List of DLQEntry objects sorted by failure time (newest first).
    """
    if not self.has_capability(BrokerCapabilities.DEAD_LETTER_QUEUE):
        logger.warning(f"Broker {self.name} does not support DLQ management")
        return []
    return self._dlq_list(dlq_streams, limit)

dlq_inspect

dlq_inspect(
    dlq_stream: str, dlq_id: str
) -> DLQEntry | None

Inspect a specific DLQ message by its DLQ entry ID.

PARAMETER DESCRIPTION
dlq_stream

The DLQ stream to look in.

TYPE: str

dlq_id

The entry identifier within the DLQ stream.

TYPE: str

RETURNS DESCRIPTION
DLQEntry | None

DLQEntry if found, None otherwise.

Source code in src/protean/port/broker.py
843
844
845
846
847
848
849
850
851
852
853
854
855
856
def dlq_inspect(self, dlq_stream: str, dlq_id: str) -> DLQEntry | None:
    """Inspect a specific DLQ message by its DLQ entry ID.

    Args:
        dlq_stream: The DLQ stream to look in.
        dlq_id: The entry identifier within the DLQ stream.

    Returns:
        DLQEntry if found, None otherwise.
    """
    if not self.has_capability(BrokerCapabilities.DEAD_LETTER_QUEUE):
        logger.warning(f"Broker {self.name} does not support DLQ management")
        return None
    return self._dlq_inspect(dlq_stream, dlq_id)

dlq_replay

dlq_replay(
    dlq_stream: str, dlq_id: str, target_stream: str
) -> bool

Replay a single DLQ message back to its original stream.

PARAMETER DESCRIPTION
dlq_stream

The DLQ stream the message is in.

TYPE: str

dlq_id

The entry identifier within the DLQ stream.

TYPE: str

target_stream

The stream to re-publish the message to.

TYPE: str

RETURNS DESCRIPTION
bool

True if replayed successfully, False otherwise.

Source code in src/protean/port/broker.py
858
859
860
861
862
863
864
865
866
867
868
869
870
871
872
def dlq_replay(self, dlq_stream: str, dlq_id: str, target_stream: str) -> bool:
    """Replay a single DLQ message back to its original stream.

    Args:
        dlq_stream: The DLQ stream the message is in.
        dlq_id: The entry identifier within the DLQ stream.
        target_stream: The stream to re-publish the message to.

    Returns:
        True if replayed successfully, False otherwise.
    """
    if not self.has_capability(BrokerCapabilities.DEAD_LETTER_QUEUE):
        logger.warning(f"Broker {self.name} does not support DLQ management")
        return False
    return self._dlq_replay(dlq_stream, dlq_id, target_stream)

dlq_replay_all

dlq_replay_all(dlq_stream: str, target_stream: str) -> int

Replay all DLQ messages from a stream back to the original stream.

PARAMETER DESCRIPTION
dlq_stream

The DLQ stream to drain.

TYPE: str

target_stream

The stream to re-publish messages to.

TYPE: str

RETURNS DESCRIPTION
int

Number of messages replayed.

Source code in src/protean/port/broker.py
874
875
876
877
878
879
880
881
882
883
884
885
886
887
def dlq_replay_all(self, dlq_stream: str, target_stream: str) -> int:
    """Replay all DLQ messages from a stream back to the original stream.

    Args:
        dlq_stream: The DLQ stream to drain.
        target_stream: The stream to re-publish messages to.

    Returns:
        Number of messages replayed.
    """
    if not self.has_capability(BrokerCapabilities.DEAD_LETTER_QUEUE):
        logger.warning(f"Broker {self.name} does not support DLQ management")
        return 0
    return self._dlq_replay_all(dlq_stream, target_stream)

dlq_purge

dlq_purge(dlq_stream: str) -> int

Purge all messages from a DLQ stream.

PARAMETER DESCRIPTION
dlq_stream

The DLQ stream to purge.

TYPE: str

RETURNS DESCRIPTION
int

Number of messages purged.

Source code in src/protean/port/broker.py
889
890
891
892
893
894
895
896
897
898
899
900
901
def dlq_purge(self, dlq_stream: str) -> int:
    """Purge all messages from a DLQ stream.

    Args:
        dlq_stream: The DLQ stream to purge.

    Returns:
        Number of messages purged.
    """
    if not self.has_capability(BrokerCapabilities.DEAD_LETTER_QUEUE):
        logger.warning(f"Broker {self.name} does not support DLQ management")
        return 0
    return self._dlq_purge(dlq_stream)

dlq_trim

dlq_trim(dlq_stream: str, min_id: str) -> int

Trim DLQ messages older than min_id (time-based trimming).

This is an optional method used by the DLQ maintenance task to remove messages that have exceeded their retention period. Brokers that don't support time-based trimming can leave the default implementation which returns 0.

PARAMETER DESCRIPTION
dlq_stream

The DLQ stream to trim.

TYPE: str

min_id

Cutoff identifier (format varies by broker). Messages older than this cutoff will be removed.

TYPE: str

RETURNS DESCRIPTION
int

Number of messages trimmed.

Source code in src/protean/port/broker.py
927
928
929
930
931
932
933
934
935
936
937
938
939
940
941
942
943
def dlq_trim(self, dlq_stream: str, min_id: str) -> int:
    """Trim DLQ messages older than *min_id* (time-based trimming).

    This is an optional method used by the DLQ maintenance task to
    remove messages that have exceeded their retention period.
    Brokers that don't support time-based trimming can leave the
    default implementation which returns 0.

    Args:
        dlq_stream: The DLQ stream to trim.
        min_id: Cutoff identifier (format varies by broker).
            Messages older than this cutoff will be removed.

    Returns:
        Number of messages trimmed.
    """
    return 0

dlq_depth

dlq_depth(dlq_stream: str) -> int

Return the current number of messages in a DLQ stream.

Brokers that don't support depth queries return 0.

PARAMETER DESCRIPTION
dlq_stream

The DLQ stream to query.

TYPE: str

RETURNS DESCRIPTION
int

Number of messages in the DLQ stream.

Source code in src/protean/port/broker.py
945
946
947
948
949
950
951
952
953
954
955
956
def dlq_depth(self, dlq_stream: str) -> int:
    """Return the current number of messages in a DLQ stream.

    Brokers that don't support depth queries return 0.

    Args:
        dlq_stream: The DLQ stream to query.

    Returns:
        Number of messages in the DLQ stream.
    """
    return 0

trim

trim(stream: str, maxlen: int) -> int

Trim a subscription stream to at most maxlen entries.

Optional maintenance hook called after each processed batch when a subscription sets retention_maxlen. Brokers with a persistent stream (e.g. Redis Streams) override this to cap the stream's size in a way that never removes an entry a consumer group still needs. Brokers without a persistent stream leave the default, which is a no-op.

PARAMETER DESCRIPTION
stream

The stream to trim.

TYPE: str

maxlen

The target maximum number of entries to keep.

TYPE: int

RETURNS DESCRIPTION
int

Number of messages trimmed (0 by default).

Source code in src/protean/port/broker.py
962
963
964
965
966
967
968
969
970
971
972
973
974
975
976
977
978
def trim(self, stream: str, maxlen: int) -> int:
    """Trim a subscription stream to at most *maxlen* entries.

    Optional maintenance hook called after each processed batch when a
    subscription sets ``retention_maxlen``. Brokers with a persistent
    stream (e.g. Redis Streams) override this to cap the stream's size in a
    way that never removes an entry a consumer group still needs. Brokers
    without a persistent stream leave the default, which is a no-op.

    Args:
        stream: The stream to trim.
        maxlen: The target maximum number of entries to keep.

    Returns:
        Number of messages trimmed (0 by default).
    """
    return 0

record_partition

record_partition(category: str, key: str) -> None

Record key in the maintained partition index for category.

Called on publish to {category}:{key} so the partitioned consumer can discover the partition without scanning the keyspace (ADR-0028 decision 7). Idempotent: recording the same key twice is a no-op.

Source code in src/protean/port/broker.py
1000
1001
1002
1003
1004
1005
1006
1007
1008
def record_partition(self, category: str, key: str) -> None:
    """Record *key* in the maintained partition index for *category*.

    Called on publish to ``{category}:{key}`` so the partitioned consumer
    can discover the partition without scanning the keyspace (ADR-0028
    decision 7). Idempotent: recording the same key twice is a no-op.
    """
    self._require_partitioning("record_partition")
    self._record_partition(category, key)

partition_keys

partition_keys(category: str) -> set[str]

Return the set of live partition keys recorded for category.

Source code in src/protean/port/broker.py
1010
1011
1012
1013
def partition_keys(self, category: str) -> set[str]:
    """Return the set of live partition keys recorded for *category*."""
    self._require_partitioning("partition_keys")
    return self._partition_keys(category)

reap_partition

reap_partition(
    category: str,
    key: str,
    min_idle_ms: int,
    backfill_suffix: str | None = None,
) -> bool

Prune a cold partition from the index, atomically and race-safely.

Only reaps a partition when no consumer group on any of its streams has pending entries (each stream is shared by every handler on the category) and all of them have been idle for at least min_idle_ms (ADR-0028 decision 7). When priority lanes are enabled, pass backfill_suffix so the key's backfill lane ({category}:{key}:{backfill_suffix}) is checked and deleted alongside the main stream; otherwise a lane with unconsumed work would be stranded. Removes the index entry and the streams; the generation counter is left in place so a re-created partition keeps a monotonic fence. Returns True when the partition was reaped. The publish/reap race is tolerated: a publisher that re-adds the key right after a reap is re-discovered on the next cycle. Meant to be called by the partition's current owner once it has drained the partition.

Source code in src/protean/port/broker.py
1015
1016
1017
1018
1019
1020
1021
1022
1023
1024
1025
1026
1027
1028
1029
1030
1031
1032
1033
1034
1035
1036
1037
1038
1039
def reap_partition(
    self,
    category: str,
    key: str,
    min_idle_ms: int,
    backfill_suffix: str | None = None,
) -> bool:
    """Prune a cold partition from the index, atomically and race-safely.

    Only reaps a partition when **no** consumer group on any of its streams
    has pending entries (each stream is shared by every handler on the
    category) and all of them have been idle for at least *min_idle_ms*
    (ADR-0028 decision 7). When priority lanes are enabled, pass
    *backfill_suffix* so the key's backfill lane
    (``{category}:{key}:{backfill_suffix}``) is checked and deleted alongside
    the main stream; otherwise a lane with unconsumed work would be
    stranded. Removes the index entry and the streams; the generation counter
    is left in place so a re-created partition keeps a monotonic fence.
    Returns ``True`` when the partition was reaped. The publish/reap race is
    tolerated: a publisher that re-adds the key right after a reap is
    re-discovered on the next cycle. Meant to be called by the partition's
    current owner once it has drained the partition.
    """
    self._require_partitioning("reap_partition")
    return self._reap_partition(category, key, min_idle_ms, backfill_suffix)

acquire_partition_lease

acquire_partition_lease(
    lease_key: str,
    generation_key: str,
    owner_id: str,
    ttl_ms: int,
) -> int | None

Acquire the ownership lease for a partition, returning its generation.

Atomically (ADR-0028 decision 5): if the lease is unheld, increment the durable generation counter, take the lease at that generation for ttl_ms, and return the generation. If already held, return None.

Source code in src/protean/port/broker.py
1041
1042
1043
1044
1045
1046
1047
1048
1049
1050
1051
1052
1053
def acquire_partition_lease(
    self, lease_key: str, generation_key: str, owner_id: str, ttl_ms: int
) -> int | None:
    """Acquire the ownership lease for a partition, returning its generation.

    Atomically (ADR-0028 decision 5): if the lease is unheld, increment the
    durable generation counter, take the lease at that generation for
    *ttl_ms*, and return the generation. If already held, return ``None``.
    """
    self._require_partitioning("acquire_partition_lease")
    return self._acquire_partition_lease(
        lease_key, generation_key, owner_id, ttl_ms
    )

renew_partition_lease

renew_partition_lease(
    lease_key: str, fence_token: str, ttl_ms: int
) -> bool

Renew a held lease, extending its TTL. Returns False if lost.

fence_token is {owner_id}:{generation}. The renew succeeds only while the lease still holds that exact token, so an owner that lost the lease (expiry, takeover) cannot resurrect it.

Source code in src/protean/port/broker.py
1055
1056
1057
1058
1059
1060
1061
1062
1063
1064
1065
def renew_partition_lease(
    self, lease_key: str, fence_token: str, ttl_ms: int
) -> bool:
    """Renew a held lease, extending its TTL. Returns ``False`` if lost.

    *fence_token* is ``{owner_id}:{generation}``. The renew succeeds only
    while the lease still holds that exact token, so an owner that lost the
    lease (expiry, takeover) cannot resurrect it.
    """
    self._require_partitioning("renew_partition_lease")
    return self._renew_partition_lease(lease_key, fence_token, ttl_ms)

release_partition_lease

release_partition_lease(
    lease_key: str, fence_token: str
) -> bool

Release a held lease on graceful shutdown. Returns False if not held.

Source code in src/protean/port/broker.py
1067
1068
1069
1070
def release_partition_lease(self, lease_key: str, fence_token: str) -> bool:
    """Release a held lease on graceful shutdown. Returns ``False`` if not held."""
    self._require_partitioning("release_partition_lease")
    return self._release_partition_lease(lease_key, fence_token)

read_partition_fenced

read_partition_fenced(
    stream: str,
    consumer_group: str,
    consumer_name: str,
    lease_key: str,
    fence_token: str,
    count: int = 1,
    new_messages: bool = True,
) -> list[tuple[str, dict[str, Any]]]

Read from a partition stream, fenced by the ownership lease.

The lease check and the XREADGROUP are one atomic step (ADR-0028 decision 5): a caller that no longer holds the lease at fence_token reads nothing and raises LeaseLostError rather than advancing the partition. new_messages selects undelivered (>) vs this consumer's already-delivered pending (0) entries. Non-blocking.

Source code in src/protean/port/broker.py
1072
1073
1074
1075
1076
1077
1078
1079
1080
1081
1082
1083
1084
1085
1086
1087
1088
1089
1090
1091
1092
1093
1094
1095
1096
1097
1098
1099
def read_partition_fenced(
    self,
    stream: str,
    consumer_group: str,
    consumer_name: str,
    lease_key: str,
    fence_token: str,
    count: int = 1,
    new_messages: bool = True,
) -> list[tuple[str, dict[str, Any]]]:
    """Read from a partition stream, fenced by the ownership lease.

    The lease check and the ``XREADGROUP`` are one atomic step (ADR-0028
    decision 5): a caller that no longer holds the lease at *fence_token*
    reads nothing and raises `LeaseLostError` rather than advancing
    the partition. ``new_messages`` selects undelivered (``>``) vs this
    consumer's already-delivered pending (``0``) entries. Non-blocking.
    """
    self._require_partitioning("read_partition_fenced")
    return self._read_partition_fenced(
        stream,
        consumer_group,
        consumer_name,
        lease_key,
        fence_token,
        count,
        new_messages,
    )

ack_partition_fenced

ack_partition_fenced(
    stream: str,
    identifier: str,
    consumer_group: str,
    lease_key: str,
    fence_token: str,
) -> bool

Ack a partition message, fenced by the ownership lease.

Like read_partition_fenced, the lease check and the XACK are atomic: a fenced stale owner cannot ack out from under the new owner and raises LeaseLostError.

Source code in src/protean/port/broker.py
1101
1102
1103
1104
1105
1106
1107
1108
1109
1110
1111
1112
1113
1114
1115
1116
1117
1118
def ack_partition_fenced(
    self,
    stream: str,
    identifier: str,
    consumer_group: str,
    lease_key: str,
    fence_token: str,
) -> bool:
    """Ack a partition message, fenced by the ownership lease.

    Like [`read_partition_fenced`][protean.port.broker.BaseBroker.read_partition_fenced], the lease check and the ``XACK`` are
    atomic: a fenced stale owner cannot ack out from under the new owner and
    raises `LeaseLostError`.
    """
    self._require_partitioning("ack_partition_fenced")
    return self._ack_partition_fenced(
        stream, identifier, consumer_group, lease_key, fence_token
    )

reclaim_partition_pending

reclaim_partition_pending(
    stream: str,
    consumer_group: str,
    consumer_name: str,
    min_idle_ms: int = 0,
    count: int = 100,
) -> list[tuple[str, dict[str, Any]]]

Reclaim a dead owner's pending entries into consumer_name (XAUTOCLAIM).

On failover the new owner reclaims the previous owner's unacked entries so none are lost or skipped (ADR-0028 decision 5, crash reclaim). Returns the reclaimed (id, payload) entries, now owned by consumer_name.

Source code in src/protean/port/broker.py
1120
1121
1122
1123
1124
1125
1126
1127
1128
1129
1130
1131
1132
1133
1134
1135
1136
1137
def reclaim_partition_pending(
    self,
    stream: str,
    consumer_group: str,
    consumer_name: str,
    min_idle_ms: int = 0,
    count: int = 100,
) -> list[tuple[str, dict[str, Any]]]:
    """Reclaim a dead owner's pending entries into *consumer_name* (XAUTOCLAIM).

    On failover the new owner reclaims the previous owner's unacked entries
    so none are lost or skipped (ADR-0028 decision 5, crash reclaim). Returns
    the reclaimed ``(id, payload)`` entries, now owned by *consumer_name*.
    """
    self._require_partitioning("reclaim_partition_pending")
    return self._reclaim_partition_pending(
        stream, consumer_group, consumer_name, min_idle_ms, count
    )

close

close() -> None

Close the broker and release all connections.

Subclasses that hold external resources (connection pools, sockets, etc.) should override this to perform cleanup. The default implementation is a no-op so that adapters without external resources (e.g. the inline broker) work without changes.

Source code in src/protean/port/broker.py
1139
1140
1141
1142
1143
1144
1145
1146
def close(self) -> None:
    """Close the broker and release all connections.

    Subclasses that hold external resources (connection pools, sockets,
    etc.) should override this to perform cleanup.  The default
    implementation is a no-op so that adapters without external
    resources (e.g. the inline broker) work without changes.
    """

register

register(subscriber_cls: type[BaseSubscriber]) -> None

Register a subscriber to this broker against its stream.

PARAMETER DESCRIPTION
subscriber_cls

The subscriber class connected to the stream.

TYPE: type[BaseSubscriber]

Source code in src/protean/port/broker.py
1148
1149
1150
1151
1152
1153
1154
1155
1156
1157
1158
1159
1160
def register(self, subscriber_cls: type[BaseSubscriber]) -> None:
    """Register a subscriber to this broker against its stream.

    Args:
        subscriber_cls: The subscriber class connected to the stream.
    """
    stream = subscriber_cls.meta_.stream

    self._subscribers[stream].add(subscriber_cls)

    logger.debug(
        f"Broker {self.name}: Registered Subscriber {subscriber_cls.__name__} for stream {stream}"
    )