Implementing Custom Domain Events in the OpenHuman Event Bus
To implement custom domain events in OpenHuman, extend the DomainEvent enum in src/core/event_bus/events.rs, map the new variant to a domain string via the domain() method, register a subscriber using subscribe_global(), and publish events using publish_global() from any component with access to CoreContext.
OpenHuman relies on a typed publish-subscribe system centered around an event bus architecture that enables both fire-and-forget notifications and synchronous request-response patterns. Understanding how to extend this system with custom domain events is essential for integrating new functionality into domains like agent, memory, channel, or tool. This guide walks through the exact implementation steps using the Rust source code from the tinyhumansai/openhuman repository.
Understanding the Event Bus Architecture
The event bus implementation in src/core/event_bus/ provides two distinct communication mechanisms. The global publish/subscribe system allows any number of listeners to receive typed notifications asynchronously, utilizing types like DomainEvent, EventBus, and SubscriptionHandle. For synchronous operations requiring immediate responses, the native request/response system uses NativeRegistry and functions like request_native_global to execute typed handlers and return results.
Every domain defines its event variants inside the central DomainEvent enum located in src/core/event_bus/events.rs. When adding custom events, you maintain type safety through Rust’s enum system while enabling selective subscription based on domain names.
Extending the DomainEvent Enum
Step 1: Add a New Variant
First, define your event payload struct and add it as a variant to the DomainEvent enum in src/core/event_bus/events.rs:
// src/core/event_bus/events.rs
pub struct UserInviteEvent {
pub inviter_id: String,
pub invitee_email: String,
}
pub enum DomainEvent {
// …existing variants…
UserInvite(UserInviteEvent),
}
Step 2: Map the Variant to a Domain
Implement the domain() method to associate your variant with a domain identifier. This enables selective subscription where listeners filter for specific domain strings:
// src/core/event_bus/events.rs
impl DomainEvent {
pub fn domain(&self) -> &'static str {
match self {
// …existing matches…
DomainEvent::UserInvite(_) => "user_invite",
}
}
}
Registering Subscribers and Publishing Events
Step 3: Register a Global Subscriber
Register your handler during system initialization, typically within src/core/jsonrpc.rs or a domain-specific startup module. Use subscribe_global() with the domain string defined in the previous step:
// src/core/jsonrpc.rs – during CoreBuilder initialization
core_context.event_bus().subscribe_global(
"user_invite",
Arc::new(|event| {
if let DomainEvent::UserInvite(payload) = event {
log::info!(
"Processing invitation from {} to {}",
payload.inviter_id,
payload.invitee_email
);
}
}),
);
Step 4: Publish Events from Any Component
Emit your custom event from any module that has access to CoreContext using publish_global():
// Example usage inside a domain module
let ev = DomainEvent::UserInvite(UserInviteEvent {
inviter_id: user.id.clone(),
invitee_email: "alice@example.com".into(),
});
core_context.event_bus().publish_global(ev);
Implementing Native Request-Response Patterns
For synchronous operations requiring immediate return values, implement a native request/response handler instead of the fire-and-forget event system. This pattern uses fully typed structs for requests and responses.
Register the handler once during startup using register_native_global():
// Define request/response types
#[derive(Debug)]
pub struct CheckQuotaRequest {
pub user_id: String,
}
#[derive(Debug)]
pub struct CheckQuotaResponse {
pub remaining: u64,
}
// Register in initialization code
core_context.event_bus().register_native_global(
"quota.check",
Arc::new(|req: CheckQuotaRequest| -> CheckQuotaResponse {
// Synchronous logic – e.g., database lookup
CheckQuotaResponse { remaining: 42 }
}),
);
Invoke the handler from any component using request_native_global() with explicit type parameters:
let req = CheckQuotaRequest { user_id: "user123".into() };
let resp = core_context
.event_bus()
.request_native_global::<CheckQuotaRequest, CheckQuotaResponse>("quota.check", req);
println!("Quota remaining: {}", resp.remaining);
Key Source Files Reference
Understanding the layout of the event bus codebase helps navigate the implementation:
src/core/event_bus/events.rs– CentralDomainEventenum and domain mapping logicsrc/core/event_bus/bus.rs– CoreEventBusimplementation providingpublish_global()andsubscribe_global()src/core/event_bus/native_request.rs– Native request/response registry and dispatchersrc/core/jsonrpc.rs– Primary location for startup registration of subscribers and native handlerssrc/openhuman/web_chat/event_bus.rs– Concrete example of publishing domain-specific events for the chat subsystem
Summary
- Extend the DomainEvent enum in
src/core/event_bus/events.rsto add custom event payloads as new variants. - Implement the domain() method to map each variant to a unique domain string for selective subscription.
- Register subscribers using subscribe_global() during system initialization, typically in
src/core/jsonrpc.rs. - Publish events using publish_global() from any component holding a
CoreContextreference. - Use native request/response patterns via
register_native_global()andrequest_native_global()when synchronous round-trips are required. - Reference existing implementations in
src/openhuman/web_chat/event_bus.rsfor domain-specific patterns.
Frequently Asked Questions
What is the difference between global events and native requests in OpenHuman?
Global events follow a fire-and-forget publish-subscribe pattern where multiple subscribers can listen to DomainEvent variants asynchronously using subscribe_global(). Native requests provide synchronous, typed request-response interactions where a single handler processes the request and returns a concrete type immediately via request_native_global().
Where should I register my event subscribers?
Register subscribers during system startup, typically within src/core/jsonrpc.rs where the CoreBuilder initializes the CoreContext. This ensures all handlers are active before the system begins processing events. Domain-specific modules may also contain registration logic if they encapsulate their own event handling.
How does OpenHuman ensure type safety for custom events?
The system leverages Rust’s compile-time type checking through the DomainEvent enum and generic functions. When you extend the enum with a new variant and match on it in subscribers, the compiler verifies that you handle the correct payload type. Native requests use explicit type parameters on request_native_global<T, R>() to enforce correct request and response structs at compile time.
Can I emit custom events from any domain module?
Yes, any component that holds a reference to CoreContext can emit events using core_context.event_bus().publish_global(). This decoupled architecture allows domains like agent, memory, or tool to emit events without knowing which subscribers exist, maintaining clean separation of concerns while participating in system-wide observability and security gating layers.
Have a question about this repo?
These articles cover the highlights, but your codebase questions are specific. Give your agent direct access to the source. Share this with your agent to get started:
curl -s "https://instagit.com/install.md" Maintain an open-source project? Get it listed too →