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:
objectCoordinate 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
Token used by
unsubscribe().Raises
Exception
Description
ValueError
event_typeis empty.TypeError
handleris 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
Token returned by
subscribe()oronce().Returns
Type
Description
bool
Truewhen a handler was removed; otherwiseFalse.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.
Nonereturns 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.
Noneuses 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
typestring.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.
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
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
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
Trueafter 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
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
Nonefor no limit.predicate
Callable[[EventDict], bool] | None
Additional acceptance condition.
Returns
Type
Description
EventDict | None
Accepted event, or
Nonewhen 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 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:
objectIdentify 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_typeis 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:
TypedDictDescribe the event mapping delivered to
EventBushandlers.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"])