ddp_utils.events

import ddp_utils.events

Provide a thread-safe synchronous and threaded publish-subscribe bus.

The bus supports exact and wildcard subscriptions, priorities, decorated handlers, one-shot listeners, blocking waits, middleware, and an optional bounded replay buffer.

Examples

Publish one synchronous event:

from ddp_utils.events import EventBus

with EventBus() as bus:
    bus.subscribe("job.done", lambda event: print(event["data"]))
    bus.emit("job.done", data={"id": 42})
class ddp_utils.events.EventBus(*, logger: LoggerProtocol | None = None, max_workers: int = 4, use_priority: bool = True)

Bases: object

Coordinate prioritized event delivery across threads safely.

Synchronous delivery calls handlers serially. Asynchronous delivery submits each handler to a managed thread pool. Handler exceptions are contained and sent to the optional logger.

Examples

Guarantee executor shutdown with a context manager:

with EventBus(max_workers=2) as bus:
    bus.emit_async("job.started")

Initialize handler storage and the asynchronous executor.

Parameters

Name

Type

Description

logger

Optional[LoggerProtocol]

Optional logger compatible with LoggerProtocol.

max_workers

int

Worker count for asynchronous handler delivery.

use_priority

bool

Sort lower numeric priorities before higher values.

Examples

Create a four-worker prioritized bus:

bus = EventBus(max_workers=4, use_priority=True)
subscribe(event_type: str, handler: Callable[[EventDict], Any], priority: int = 0) → SubscriptionToken

Subscribe one callable to an exact or wildcard event type.

Parameters

Name

Type

Description

event_type

str

Non-empty event name or "*" wildcard.

handler

Callable[[EventDict], Any]

Callable accepting one EventDict.

priority

int

Lower values run earlier when priority sorting is enabled.

Returns

Type

Description

SubscriptionToken

Token used by unsubscribe().

Raises

Exception

Description

ValueError

event_type is empty.

TypeError

handler is not callable.

RuntimeError

The bus is closed.

Examples

Subscribe and retain the removal token:

token = bus.subscribe("job.done", handle_done, priority=-10)
unsubscribe(token: SubscriptionToken) → bool

Remove the subscription identified by a token.

Parameters

Name

Type

Description

token

SubscriptionToken

Token returned by subscribe() or once().

Returns

Type

Description

bool

True when a handler was removed; otherwise False.

Examples

Remove a previously registered handler:

removed = bus.unsubscribe(token)
unsubscribe_all() → None

Remove every exact and wildcard subscription atomically.

Examples

Reset handlers while keeping the bus open:

bus.unsubscribe_all()
classmethod collect_decorated(module_globals: Dict[str, Any] | None = None) → List[Tuple[str, Callable[[EventDict], Any], Dict[str, Any]]]

Collect handlers marked by on().

Parameters

Name

Type

Description

module_globals

Dict[str, Any] | None

Optional namespace to scan explicitly. None returns a copy of the module-level decorated registry.

Returns

Type

Description

List[Tuple[str, Callable[[EventDict], Any], Dict[str, Any]]]

(event_type, handler, options) triples.

Examples

Collect only handlers from the current module:

handlers = EventBus.collect_decorated(globals())
register_decorated(handlers: Iterable[Tuple[str, Callable[[EventDict], Any], Dict[str, Any]]] | None = None) → List[SubscriptionToken]

Subscribe a collection of decorated handlers to this bus.

Parameters

Name

Type

Description

handlers

Iterable[Tuple[str, Callable[[EventDict], Any], Dict[str, Any]]] | None

Optional iterable of event, callable, and options triples. None uses the module-level registry.

Returns

Type

Description

List[SubscriptionToken]

Subscription tokens in input order.

Examples

Register handlers discovered in a module:

tokens = bus.register_decorated(
    EventBus.collect_decorated(globals()),
)
event(event_type: str, *, data: Dict[str, Any] | None = None, meta: Dict[str, Any] | None = None, source: str = 'business_client') → EventDict

Build an event mapping without publishing it.

Parameters

Name

Type

Description

event_type

str

Value converted to the event type string.

data

Dict[str, Any] | None

Optional business payload; falsey values become a new dict.

meta

Dict[str, Any] | None

Optional metadata; falsey values become a new dict.

source

str

Publisher identifier.

Returns

Type

Description

EventDict

New timestamped EventDict.

Examples

Prepare an event for inspection:

event = bus.event("job.done", data={"id": 42})
emit(event_type: str, *, data: Dict[str, Any] | None = None, meta: Dict[str, Any] | None = None, source: str = 'business_client') → EventDict

Publish an event synchronously through middleware and handlers.

Parameters

Name

Type

Description

event_type

str

Event name.

data

Dict[str, Any] | None

Optional business payload.

meta

Dict[str, Any] | None

Optional transport metadata.

source

str

Publisher identifier.

Returns

Type

Description

EventDict

Published event after all synchronous handlers finish.

Examples

Publish a completion event:

event = bus.emit("job.done", data={"id": 42})
emit_async(event_type: str, *, data: Dict[str, Any] | None = None, meta: Dict[str, Any] | None = None, source: str = 'business_client') → EventDict

Submit event handlers to the thread pool and return immediately.

Parameters

Name

Type

Description

event_type

str

Event name.

data

Dict[str, Any] | None

Optional business payload.

meta

Dict[str, Any] | None

Optional transport metadata.

source

str

Publisher identifier.

Returns

Type

Description

EventDict

Created event. A closed bus returns an undispatched event after attempting to log the drop.

Examples

Queue background delivery:

event = bus.emit_async("job.started", data={"id": 42})
close(wait: bool = True) → None

Close subscriptions and shut down the asynchronous executor.

Parameters

Name

Type

Description

wait

bool

Wait for submitted tasks. When False, queued futures that have not started are cancelled.

Examples

Close without waiting for queued handlers:

bus.close(wait=False)
property closed: bool

Return whether close() has completed its state transition.

Returns

True after the bus has been closed.

Examples

Avoid subscribing after shutdown:

if not bus.closed:
    bus.subscribe("ready", handle_ready)
once(event_type: str, handler: Callable[[EventDict], Any], priority: int = 0) → SubscriptionToken

Subscribe a handler that removes itself before its first invocation.

Parameters

Name

Type

Description

event_type

str

Exact or wildcard event name.

handler

Callable[[EventDict], Any]

Callable invoked once.

priority

int

Subscription priority.

Returns

Type

Description

SubscriptionToken

Token for optional removal before the event occurs.

Examples

Run startup logic once:

token = bus.once("app.ready", initialize)
await_event(event_type: str, timeout: float | None = None, predicate: Callable[[EventDict], bool] | None = None) → EventDict | None

Block until a matching event passes an optional predicate.

Parameters

Name

Type

Description

event_type

str

Exact event name or "*".

timeout

float | None

Maximum seconds, or None for no limit.

predicate

Callable[[EventDict], bool] | None

Additional acceptance condition.

Returns

Type

Description

EventDict | None

Accepted event, or None when waiting times out.

Examples

Wait for one specific order payment:

event = bus.await_event(
    "order.paid",
    timeout=30,
    predicate=lambda item: item["data"].get("id") == 42,
)
enable_replay(max_events: int = 100) → None

Capture future events in a bounded replay deque.

Parameters

Name

Type

Description

max_events

int

Maximum retained events. The implementation installs a low-priority wildcard subscriber each time it is called.

Examples

Retain the most recent fifty events:

bus.enable_replay(max_events=50)
get_replay(event_type: str | None = None) → List[EventDict]

Return a snapshot of retained replay events.

Parameters

Name

Type

Description

event_type

str | None

Optional exact event filter.

Returns

Type

Description

List[EventDict]

New list of all buffered events or those with the requested type. An unconfigured buffer behaves as empty.

Examples

Retrieve only completion events:

completed = bus.get_replay("job.done")
use(middleware: Callable[[EventDict, Callable[[EventDict], None]], None]) → None

Append middleware to the event-delivery chain.

Parameters

Name

Type

Description

middleware

Callable[[EventDict, Callable[[EventDict], None]], None]

Callable receiving an event and continuation. It must invoke next_fn(event) to continue delivery.

Examples

Add simple event observation:

def log_event(event, next_fn):
    print(event["type"])
    next_fn(event)

bus.use(log_event)
class ddp_utils.events.SubscriptionToken(id: int, event_type: str)

Bases: object

Identify one subscription for precise removal.

Variables

Name

Type

Description

id

int

Monotonically assigned subscription identifier.

event_type

str

Exact event bucket containing the subscription.

Examples

Remove the associated subscription:

bus.unsubscribe(token)
ddp_utils.events.on(event_type: str, **options: Any) → Callable[[Callable[[EventDict], Any]], Callable[[EventDict], Any]]

Mark a function for later registration as an event handler.

Parameters

Name

Type

Description

event_type

str

Non-empty event name.

**options

Any

Registration metadata such as priority.

Returns

Type

Description

Callable[[Callable[[EventDict], Any]], Callable[[EventDict], Any]]

Decorator that records event attributes on the original function and appends it to the module-level decorated-handler registry.

Raises

Exception

Description

ValueError

event_type is empty after normalization.

Examples

Declare a prioritized handler:

@on("order.created", priority=-10)
def validate_order(event):
    validate(event["data"])
class ddp_utils.events.EventDict

Bases: TypedDict

Describe the event mapping delivered to EventBus handlers.

Variables

Name

Type

Description

type

str

Event name such as "user.created".

ts

float

Unix timestamp assigned at publication.

source

str

Publisher or component identifier.

data

Dict[str, Any]

Arbitrary business payload.

meta

Dict[str, Any]

Correlation, retry, or other transport metadata.

Examples

Access an event in a handler:

def handle(event: EventDict) -> None:
    print(event["type"], event["data"])