How the TinyBus Event Bus Facilitates Typed Pub/Sub and Request/Response Across Domains
The tinybus event bus enables type-safe publish/subscribe and synchronous request/response communication across domain boundaries by wrapping tinybus primitives in global singleton functions that filter events based on DomainSet membership.
OpenHuman relies on the tinybus crate to decouple core services and native modules while preserving compile-time type safety. The runtime implements thin wrappers around tinybus that support both asynchronous event broadcasting and blocking RPC calls, routing every message through domain-specific filters defined in src/core/events.rs and src/core/bus.rs.
Global Singleton Initialization
At process startup, the core initializes a process-global singleton via tinybus::init_global(). Subsequent components obtain a handle to the bus through tinybus::global(); if the bus is not yet initialized, the codebase logs a warning (observed in src/openhuman/platform/service/shutdown.rs). This singleton guarantees a single, authoritative event source for the entire application, ensuring that publishers and subscribers reference the same underlying queue.
Typed Publish/Subscribe Architecture
Defining Domain Events
Domain events are Rust enums defined in src/core/events.rs (e.g., DomainEvent::SessionExpired, DomainEvent::ChannelInbound, DomainEvent::MemorySyncComplete). Each variant implements the tinybus::EventHandler trait, which exposes two critical methods:
name()– returns a stable string identifier such as"core::session_expired".domains()– returns the set ofDomainGroups authorized to receive the event.
This metadata allows the bus to enforce domain boundaries at runtime without sacrificing compile-time type safety.
Publishing Events Across Domains
Components emit events through the core::bus::publish_global() wrapper, which forwards the payload to all subscribers whose domains() intersect with the event’s domain set. For example, the Web-Chat bridge in src/openhuman/web_chat/event_bus.rs exposes a convenience function:
// src/openhuman/web_chat/event_bus.rs
pub fn publish_web_channel_event(event: DomainEvent) {
core::bus::publish_global(event);
}
A caller can signal a memory synchronization completion like this:
use crate::core::events::DomainEvent;
use crate::core::bus::publish_global;
fn signal_memory_sync_complete() {
let ev = DomainEvent::MemorySyncComplete { workspace_id: "ws1".into() };
publish_global(ev);
}
Subscribing with Domain Filters
Subscribers register interest via core::bus::subscribe_global() or domain-specific wrappers such as subscribe_web_channel_events(). The function returns a tinybus::SubscriptionHandle that remains active until dropped:
// src/openhuman/web_chat/event_bus.rs
pub fn subscribe_web_channel_events() -> tinybus::SubscriptionHandle {
tinybus::subscribe_global::<crate::core::events::DomainEvent>()
}
During startup, the core inspects the enabled DomainSet (managed in src/core/jsonrpc.rs) and attaches only the relevant subscribers, logging selections such as:
[event_bus] register_domain_subscribers: domains=[Agent, Memory] plan=…
Disabled domains skip registration entirely, ensuring that stray events never reach unauthorized handlers.
Request/Response for Native Modules
Interface Registration and Calling
Native modules (dynamic cdylib libraries) expose typed interfaces using the #[tinybus::interface] attribute. The core invokes these methods synchronously via core::bus::request_native_global(), which routes the call to the module’s object path and returns a typed tinybus::Result<T>:
let result: tinybus::Result<T> =
core::bus::request_native_global::<ModuleInterface, MethodSignature>(args);
This pattern is implemented concretely in src/openhuman/modules/wallet.rs, where wallet operations are exposed as tinybus interfaces and called from the core runtime.
Error Mapping and Type Safety
When a native call fails, tinybus returns tinybus::Error::MethodFailed. The OpenHuman wrappers map these low-level errors to domain-specific enums—such as WalletCallError in src/openhuman/modules/wallet.rs—preserving type safety for callers while abstracting transport details.
Module Lifecycle and Domain Registration
The src/openhuman/modules/registry.rs handles dynamic loading, digest verification, and bus attachment. After verifying a module’s integrity, the registry attaches it to the global tinybus instance, making its interfaces available to request_native_global callers. This ensures that only verified modules participate in the request/response fabric.
Summary
- tinybus provides the underlying typed message bus, while
src/core/bus.rsexposes thin global wrappers (publish_global,subscribe_global,request_native_global) that enforce domain scoping. - Domain events are defined in
src/core/events.rsand implementtinybus::EventHandlerto declare routing metadata via thedomains()method. - Pub/sub messages are filtered at runtime by comparing the event’s domain set against the subscriber’s registered
DomainSet, preventing cross-domain leakage. - Request/response calls use generic typed interfaces (
ModuleInterface,MethodSignature) to invoke native module methods with compile-time guarantees. - Error propagation maps
tinybus::Errorvariants to domain-specific enums, maintaining type safety across the FFI boundary.
Frequently Asked Questions
How does tinybus ensure type safety across domain boundaries?
tinybus leverages Rust’s type system through the EventHandler trait and generic request_native_global<T> signatures. Because event payloads and RPC return types are concrete Rust types rather than raw bytes, the compiler verifies correctness at build time, eliminating serialization mismatches before runtime.
What is the difference between publish_global and request_native_global?
publish_global implements a fire-and-forget broadcast pattern: it delivers a DomainEvent to every subscriber whose domain set intersects with the event’s metadata, allowing one-to-many communication. In contrast, request_native_global implements synchronous RPC: it targets a specific native module interface, blocks until completion, and returns a typed tinybus::Result<T> to the caller.
How are domain-specific subscribers filtered at runtime?
During initialization (handled in src/core/jsonrpc.rs), the core reads the active DomainSet and registers only those subscribers whose domains() method returns a set overlapping with the enabled domains. This is visible in logs such as [event_bus] register_domain_subscribers: domains=[Agent, Memory], confirming that disabled domains never receive events.
Where are domain events defined in the OpenHuman codebase?
Core domain events are enumerated in src/core/events.rs as the DomainEvent type. Each variant (e.g., SessionExpired, MemorySyncComplete) implements the tinybus::EventHandler trait to supply stable names and domain tags, while concrete usage examples appear in src/openhuman/web_chat/event_bus.rs and src/openhuman/modules/wallet.rs.
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 →