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
199
200
201
202
203
204
205
206
207
208
209
210
211
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
222
223
224
225
226
227
228
229
230
231
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
233
234
235
236
237
238
239
240
241
242
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
244
245
246
247
248
249
250
251
252
253
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
255
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
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
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
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:
        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
320
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
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:
        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
370
371
372
373
374
375
376
377
378
379
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
447
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
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
484
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
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
516
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
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
548
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
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 = block indefinitely)

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
595
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
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 = block indefinitely)
        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

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
698
699
700
701
702
703
704
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
718
719
720
721
722
723
724
725
726
727
728
729
730
731
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
733
734
735
736
737
738
739
740
741
742
743
744
745
746
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
748
749
750
751
752
753
754
755
756
757
758
759
760
761
762
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
764
765
766
767
768
769
770
771
772
773
774
775
776
777
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
779
780
781
782
783
784
785
786
787
788
789
790
791
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
817
818
819
820
821
822
823
824
825
826
827
828
829
830
831
832
833
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
835
836
837
838
839
840
841
842
843
844
845
846
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
852
853
854
855
856
857
858
859
860
861
862
863
864
865
866
867
868
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
890
891
892
893
894
895
896
897
898
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
900
901
902
903
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
905
906
907
908
909
910
911
912
913
914
915
916
917
918
919
920
921
922
923
924
925
926
927
928
929
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
931
932
933
934
935
936
937
938
939
940
941
942
943
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
945
946
947
948
949
950
951
952
953
954
955
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
957
958
959
960
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
962
963
964
965
966
967
968
969
970
971
972
973
974
975
976
977
978
979
980
981
982
983
984
985
986
987
988
989
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
 991
 992
 993
 994
 995
 996
 997
 998
 999
1000
1001
1002
1003
1004
1005
1006
1007
1008
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
1010
1011
1012
1013
1014
1015
1016
1017
1018
1019
1020
1021
1022
1023
1024
1025
1026
1027
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
1029
1030
1031
1032
1033
1034
1035
1036
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
1038
1039
1040
1041
1042
1043
1044
1045
1046
1047
1048
1049
1050
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}"
    )