Skip to content

Inline Broker

The Inline broker is an in-memory message broker that keeps messages within the same process. With message_processing = "sync", it delivers each message to its subscribers as soon as it is published. It's the default broker in Protean and requires no external dependencies.

Overview

The Inline broker is designed for:

  • Development environments where simplicity is key
  • Testing scenarios where deterministic behavior is required
  • Small applications that don't need distributed messaging
  • Prototyping when you want to defer technology decisions

The Inline broker maintains messages in memory using Python data structures:

  • Messages are stored in dictionaries keyed by stream name
  • Consumer groups track message processing state
  • All data is lost when the process terminates

Configuration

[brokers.default]
provider = "inline"

# Optional configuration for retry behavior
max_retries = 3  # Maximum retry attempts for failed messages
retry_delay = 1.0  # Initial retry delay in seconds
backoff_multiplier = 2.0  # Exponential backoff multiplier
message_timeout = 300.0  # Message timeout in seconds (5 minutes default)
enable_dlq = true  # Enable dead letter queue for failed messages

Configuration Options

Option Default Description
provider Required Must be "inline" for Inline broker
max_retries 3 Maximum retry attempts for failed messages
retry_delay 1.0 Initial retry delay in seconds
backoff_multiplier 2.0 Multiplier for exponential backoff
message_timeout 300.0 Timeout for message processing (seconds)
enable_dlq true Enable dead letter queue

Capabilities

The Inline broker supports the following capabilities:

  • ✅ BASIC_PUBSUB - Fire-and-forget message publishing
  • ✅ SIMPLE_QUEUING - Consumer groups for message distribution
  • ✅ RELIABLE_MESSAGING - Message acknowledgment and rejection
  • ✅ DEAD_LETTER_QUEUE - Failed messages routed to DLQ for inspection and replay
  • ❌ ORDERED_MESSAGING - Not supported
  • ❌ ENTERPRISE_STREAMING - Not supported

Usage Examples

Basic Publishing and Subscribing

A subscriber is a class with a __call__ method that receives each message as a dict:

from protean import Domain

domain = Domain(name="Users")
domain.config["brokers"] = {"default": {"provider": "inline"}}

# Call subscribers in the publishing process, as soon as a message is published
domain.config["message_processing"] = "sync"

created = []


# Subscribing to messages
@domain.subscriber(stream="user-events")
class UserEventSubscriber:
    def __call__(self, payload: dict) -> None:
        if payload["type"] == "user.created":
            print(f"User created: {payload['name']}")
            created.append(payload["user_id"])


# Register subscribers before initializing the domain
domain.init(traverse=False)

with domain.domain_context():
    # Publishing messages
    domain.brokers.publish(
        stream="user-events",
        message={"type": "user.created", "user_id": "123", "name": "John Doe"},
    )

The subscriber prints User created: John Doe.

Testing with Inline Broker

The Inline broker is ideal for testing as it provides deterministic, synchronous behavior. Use Protean's DomainFixture to manage the domain lifecycle:

import pytest

from protean import Domain
from protean.integrations.pytest import DomainFixture

testing_domain = Domain(name="InlineTesting")
testing_domain.config["brokers"] = {"default": {"provider": "inline"}}
testing_domain.config["message_processing"] = "sync"

# Track processed messages
processed = []


@testing_domain.subscriber(stream="test-stream")
class RecordingSubscriber:
    def __call__(self, payload: dict) -> None:
        processed.append(payload)


@pytest.fixture(scope="session")
def app_fixture():
    fixture = DomainFixture(testing_domain)
    fixture.setup()
    yield fixture
    fixture.teardown()


@pytest.fixture(autouse=True)
def _ctx(app_fixture):
    with app_fixture.domain_context():
        yield


def test_message_processing():
    # Publish a message
    testing_domain.brokers.publish(
        stream="test-stream",
        message={"type": "test.event", "data": "test"},
    )

    # Message is processed synchronously
    assert len(processed) == 1
    assert processed[0]["data"] == "test"

Consumer Groups

Consumer groups read a stream independently. Each group receives every message, and within a group each message is handed out once:

with domain.domain_context():
    broker = domain.brokers["default"]
    broker.publish("orders", {"type": "order.created", "order_id": "A1"})

    # Each consumer group receives every message on the stream
    billing_id, billing_message = broker.get_next("orders", "billing")
    shipping_id, shipping_message = broker.get_next("orders", "shipping")

    # Within one group, each message is handed out once. Acknowledge it when done.
    broker.ack("orders", billing_id, "billing")
    broker.ack("orders", shipping_id, "shipping")

Limitations

  • No Persistence
    • Messages are lost on process restart
    • No durability guarantees
    • Cannot recover from crashes
  • No Distribution
    • Cannot scale across multiple processes
    • All processing happens in the same Python process
    • Not suitable for high-throughput scenarios
  • No Ordering Guarantees
    • Messages may be processed out of order
    • No support for partitioned delivery
    • Cannot ensure strict message sequencing
  • Limited Error Recovery
    • Basic retry support with configurable max retries
    • Dead letter queue stores failed messages in memory (lost on restart)
    • DLQ messages can be listed, inspected, replayed, and purged via protean dlq CLI or the Observatory dashboard

Migration Path

The Inline broker is designed to be easily replaced with production-ready brokers:

# Development configuration
[dev.brokers.default]
provider = "inline"

# Production configuration (same code works!)
[prod.brokers.default]
provider = "redis"
URI = "redis://localhost:6379/0"

Your application code remains unchanged when switching brokers, as long as you:

  1. Only use capabilities supported by both brokers
  2. Handle broker-specific errors appropriately
  3. Test with the production broker before deployment