ZCore LogoZCore
How to

How to listen to and dispatch domain events

Subscribe to event channels using @on_event decorators, register listeners in plugins, and dispatch events safely.

ZCore's EventDispatcher enables completely decoupled communication between domains. Services can publish domain occurrences, while other services listen and react asynchronously without direct coupling.

1. Decorate Service Methods with @on_event

Decorate any asynchronous method in your service with @on_event("channel_name"):

# tasks/services.py
from zcore import BaseService, on_event
from .models import Task
from .repositories import TaskRepo

class TaskService(BaseService[Task]):
    def __init__(self, repo: TaskRepo):
        super().__init__(model=Task, repository=repo)

    @on_event("user.registered")
    async def on_user_registered(self, payload: dict):
        # Automatically creates an onboarding task when a new user registers
        user_id = payload.get("user_id")
        user_email = payload.get("email")
        
        print(f"Creating welcome task for user {user_email} ({user_id})")
        # await self.create(...)

2. Register Service Listeners in the Plugin

In your domain's plugin.py, resolve the global EventDispatcher and register the service class using register_listeners:

# tasks/plugin.py
from fastapi import FastAPI
from zcore import Plugin, container, EventDispatcher
from .routers import router_instance
from .services import TaskService

class TasksPlugin(Plugin):
    name = "tasks"
    version = "0.1.0"
    dependencies = []

    def setup(self, app: FastAPI) -> None:
        app.include_router(router_instance.router)

        # Register all @on_event listeners on TaskService into the global dispatcher
        dispatcher = container.resolve(EventDispatcher)
        dispatcher.register_listeners(TaskService, container)

Dynamic Auto-Wiring: When an event is dispatched, EventDispatcher resolves TaskService from the IoCContainer dynamically. All dependencies (such as repositories and database sessions) are injected automatically.

3. Dispatching Events

You can dispatch events in two ways:

Option A: Post-Commit via UnitOfWork (Recommended for DB operations)

Ensures events only fire after the database transaction commits successfully:

async with UnitOfWork(self.repository.db, self.dispatcher) as uow:
    user = await self.repository.create(user_data)
    
    # Buffered in memory and dispatched only after commit succeeds!
    uow.register_event("user.registered", {"user_id": str(user.id), "email": user.email})

Option B: Direct Dispatching (Immediate)

For immediate, non-transactional triggers:

from zcore import container, EventDispatcher

dispatcher = container.resolve(EventDispatcher)
await dispatcher.dispatch("user.registered", {"user_id": "123", "email": "[email protected]"})

Concurrent & Resilient: All asynchronous handlers listening to the same event channel execute concurrently via asyncio.gather. If one listener raises an error, it is caught and logged without aborting other listeners.

On this page