medi-rs 2.3.0

A lightweight async mediator library for Rust with dependency injection and derive macros
Documentation

medi-rs

medi-rs is a static async mediator for Rust. Applications declare commands, events, resources, and handlers in feature-local modules; mediator! then combines those manifests into one concrete mediator. Command dispatch and resource injection are generated at compile time—there is no runtime handler registry or type-based resource lookup.

Documentation

Runtime features

Choose one runtime adapter for event processing:

Feature Runtime Notes
tokio Tokio Hosted applications and the runnable Tokio examples.
wasm wasm_bindgen_futures WebAssembly event workers use spawn_local.
embassy Embassy no_std embedded applications; queue capacity must be a const expression.

The runtime features are mutually exclusive. Command-only mediators need no runtime feature. Events and runtime tasks require exactly one adapter feature.

[dependencies]
medi-rs = { version = "1", features = ["tokio"] }

Event configuration

Event mediators use a bounded queue. event_queue_capacity and event_workers must both be greater than zero. publish waits while the queue is full; try_publish never waits and returns the event when the queue is full or closed. Tokio and WebAssembly accept any usize expression for the capacity. Embassy uses the capacity as a const generic, so its value must be a const expression such as a literal or named const.

Quick start

A command derives MediCommand; its handler is marked with #[medi_handler]. A medi_module! manifest declares the route, and mediator! creates the application mediator.

use medi_rs::{MediCommand, Result, medi_handler, medi_module, mediator};

#[derive(MediCommand)]
#[medi_command(return_type = String, error_type = medi_rs::Error)]
struct Greet {
    name: String,
}

#[medi_handler]
async fn greet(command: Greet) -> Result<String> {
    Ok(format!("Hello, {}!", command.name))
}

medi_module! {
    manifest greeting;
    commands { Greet => greet; }
}

mediator! {
    pub struct AppMediator {
        event_queue_capacity: 16;
        event_workers: 1;
        modules: [greeting];
    }
}

async fn run() -> Result<()> {
    let greeting = AppMediator::new()
        .send(Greet { name: "Rust".into() })
        .await?;
    assert_eq!(greeting, "Hello, Rust!");
    Ok(())
}

MediCommand defaults to a () response and core::convert::Infallible error. Specify return_type and error_type when the handler returns other types.

Static-dispatch pattern

The registration graph is fixed where mediator! is expanded:

  1. Define a message type and derive MediCommand for commands.
  2. Mark each async handler with #[medi_handler]; its final parameter is the command or event it receives.
  3. In each feature-local module, use medi_module! to list its resources and routes.
  4. At the application boundary, select those manifests with mediator!.

mediator! generates one concrete mediator type. For every command route it implements a route for that command type and mediator type; send therefore calls the selected handler directly. Resources live in a typed nested tuple and are cloned into handler parameters by their compile-time tuple position. There is no runtime TypeId lookup, boxed handler registry, or handler selection at runtime. A command or resource registered twice is rejected while expanding the composition, and requesting an undeclared handler resource fails type checking.

Macro reference

#[derive(MediCommand)] and #[medi_command(...)]

Derive MediCommand on a command or query type. It implements Command, which supplies the response and handler error types:

#[derive(medi_rs::MediCommand)]
#[medi_command(return_type = User, error_type = CreateUserError)]
struct CreateUser { /* fields */ }

Both options are optional: return_type defaults to () and error_type defaults to core::convert::Infallible. return_type and error_type accept Rust type syntax, so application-specific result and error types work without conversion to a framework error.

#[medi_handler]

Apply this attribute to an async function. It retains the function and creates a typed, crate-visible internal invoker used by generated routes, so the handler function itself can remain private to its feature module. The last parameter is always the message. Value parameters before it are resources, which must be listed in the composed manifest and implement Clone. Optionally, the first parameter can be &AppMediator so a handler can send another command or publish an event.

#[medi_handler]
async fn create_user(
    mediator: &AppMediator,
    repository: UserRepository,
    command: CreateUser,
) -> Result<User, CreateUserError> {
    // `repository` is cloned from AppMediator's declared resources.
    mediator.send(RecordAudit).await?;
    repository.create(command).await
}

The handler's return type must match the command metadata for a command route. For event routes it should return Result<()>; event-handler errors are ignored after the handler completes. To wrap an invocation with cross-cutting behavior, define middleware functions that receive the command and a next continuation, then list them in declaration order:

use medi_rs::DecoratorNext;

async fn logging(
    command: CreateUser,
    next: impl DecoratorNext<CreateUser, Response = User, Error = CreateUserError>,
) -> Result<User, CreateUserError> {
    println!("creating a user");
    next.call(command).await
}

#[medi_handler(decorators = [logging, validation])]
async fn create_user(command: CreateUser) -> Result<User, CreateUserError> {
    // `logging` wraps `validation`, which wraps this handler.
    todo!()
}

next is the remaining decorator pipeline and handler. Calling next.call(command).await forwards the command; a decorator can reject or modify the command before forwarding it, and can run behavior after it returns. The generated continuation is inferred automatically; only the DecoratorNext<Command> parameter type must be declared.

To apply a decorator to every command and event handler in one mediator, add it to the mediator composition:

mediator! {
    struct AppMediator {
        event_queue_capacity: 16;
        event_workers: 1;
        modules: [users];
        decorators: [logging];
    }
}

medi_module!

Declare a reusable, feature-local manifest. It contains zero or more resources, commands, and events sections in any order. With a runtime feature, manifests may additionally contain a tasks section. Commands have one handler; events have one or more handlers. The macro creates the named manifest for inclusion by mediator!.

medi_module! {
    manifest users;
    resources { UserRepository; Clock; }
    commands { CreateUser => create_user; }
    events { UserCreated => [send_welcome_email, write_audit_log]; }
}

Use semicolons between resource and command entries, and commas between event handlers. When a handler is private in a feature module, use its crate-qualified path (for example, crate::users::create_user) in the manifest; the generated invoker remains crate-visible while the handler stays private. The manifest contains declarations only: it does not construct a mediator or register anything dynamically.

mediator!

Compose one or more manifests into the concrete application mediator. Its explicit modules list is the routing boundary and its order determines the order of resource arguments accepted by new.

mediator! {
    pub struct AppMediator {
        event_queue_capacity: 16;
        event_workers: 1;
        modules: [users, audit];
    }
}

The generated type has new, send, and, when an event route exists, publish, try_publish, start, and shutdown. It uses the runtime selected by the enabled tokio, wasm, or embassy feature. The event configuration rules are described in Event configuration. start requires 'static mediator storage (start(spawner) for Embassy).

To observe asynchronous handler failures, configure a unit-struct EventFailureReporter type. The reporter runs once per failed handler and does not prevent subsequent handlers from running. Because event routes may have unrelated error types, EventHandlerFailure intentionally provides event and handler names, not the concrete error value.

mediator! {
    struct AppMediator {
        event_queue_capacity: 16;
        event_workers: 1;
        modules: [users];
        event_failure_reporter: LogEventFailure;
    }
}

__medi_rs_finalize_composition! is exported only for internal macro expansion. It is not an application-facing API; use mediator! instead.

Resources

Resources are ordinary Clone values. List each resource in a module manifest, pass the values to the generated mediator constructor in declaration order, and request them as handler parameters before the command or event.

use medi_rs::{MediCommand, Result, medi_handler, medi_module, mediator};

#[derive(Clone)]
struct UserRepository;

#[derive(MediCommand)]
#[medi_command(error_type = medi_rs::Error)]
struct CreateUser;

#[medi_handler]
async fn create_user(_: UserRepository, _: CreateUser) -> Result<()> {
    Ok(())
}

medi_module! {
    manifest users;
    resources { UserRepository; }
    commands { CreateUser => create_user; }
}

mediator! {
    struct AppMediator {
        event_queue_capacity: 16;
        event_workers: 1;
        modules: [users];
    }
}

async fn run() -> Result<()> {
    AppMediator::new(UserRepository).send(CreateUser).await?;
    Ok(())
}

A missing or duplicate resource is a compile-time error. Resource derive macros are not required.

Request-scoped data

Resources are fixed when the mediator is constructed; send accepts only the command, so it cannot inject a resource for one call. Model request-scoped dependencies such as authentication, tenant information, and correlation IDs as command data instead. A command is moved into send and then into its handler, so its fields are not cloned during normal command dispatch.

struct AuthContext {
    user_id: String,
}

#[derive(MediCommand)]
#[medi_command(error_type = medi_rs::Error)]
struct CreateInvoice {
    customer_id: String,
    // `None` explicitly represents an anonymous request.
    auth: Option<AuthContext>,
}

async fn create_invoice(request: HttpRequest, mediator: &AppMediator) -> Result<()> {
    let auth = authenticate(&request).await?;

    mediator
        .send(CreateInvoice {
            customer_id: request.customer_id().to_owned(),
            auth: Some(auth),
        })
        .await
}

Use Option<T> in resources { ... } only when a dependency is optional for the lifetime of the mediator. For request-scoped data, keep the command field small (or use Arc inside it when sharing larger data is necessary).

Runtime tasks

With a runtime feature, #[medi_task] creates a task with the same typed resource injection as a handler. It may take &AppMediator as its first parameter when it needs mediator access. It may then take &ShutdownSignal; this is injected by the mediator and is cancelled by mediator.shutdown(). Remaining value parameters are declared resources. Register it in a tasks section. mediator.start(spawner) starts tasks on Embassy; mediator.start() does so on Tokio and Wasm.

#[medi_task]
async fn watch_button(mediator: &AppMediator, signal: &ShutdownSignal, board: BoardApi) {
    loop {
        // Select the runtime-specific work future against signal.cancelled().
        if signal.is_shutdown() {
            break;
        }
        board.wait_for_button_a().await;
        let _ = mediator.send(ButtonPressed).await;
    }
}

medi_module! {
    manifest buttons;
    resources { BoardApi; }
    tasks { watch_button; }
}

Events

Events are plain Clone + Send + 'static values. List each event route in a module manifest, create a 'static mediator, and call start before publishing. start starts workers and #[medi_task] tasks exactly once; later calls return StartError::AlreadyStarted and do not spawn additional work. Use is_started to inspect that state. Command-only mediators have no startup work, so is_started remains false. Each generated worker dispatches an event to every registered handler. publish waits when the configured bounded queue is full. Use try_publish when the caller must not wait: it returns TryPublishError::Full(event) when capacity is exhausted and TryPublishError::Closed(event) after shutdown or when workers are unavailable. Both errors retain the event for retry, persistence, or disposal. Handler failures never make publish fail and never stop the other handlers. Without event_failure_reporter, they are discarded for backwards compatibility. With a reporter, each failure is observed once through EventHandlerFailure metadata; concrete handler errors are not exposed because routes need not share an error type.

use medi_rs::{Result, medi_handler, medi_module, mediator};

#[derive(Clone)]
struct UserRegistered;

#[medi_handler]
async fn send_welcome_email(_: UserRegistered) -> Result<()> {
    Ok(())
}

medi_module! { manifest users; events { UserRegistered => [send_welcome_email]; } }
mediator! { struct AppMediator { event_queue_capacity: 16; event_workers: 1; modules: [users]; } }

async fn run() -> Result<()> {
    let mediator = Box::leak(Box::new(AppMediator::new()));
    mediator.start().expect("mediator must start");
    mediator.publish(UserRegistered).await?;
    Ok(())
}

For example, an HTTP endpoint can apply its own overload policy without awaiting queue capacity:

use medi_rs::TryPublishError;

let event = UserRegistered;
match mediator.try_publish(event) {
    Ok(()) => metrics::counter!("events.enqueued").increment(1),
    Err(TryPublishError::Full(event)) => persist_for_retry(event).await?,
    Err(TryPublishError::Closed(_event)) => return Err(AppError::ShuttingDown),
}

Shutdown

shutdown().await closes the event queue, so subsequent publish calls return Error::EventPublishingError. It then drains every event whose publish already returned Ok(()) and waits until all event workers and registered #[medi_task] tasks have returned. Tasks are cooperative: a task that runs indefinitely must return for shutdown to complete.

Tokio and Wasm spawn workers when start() is called; await shutdown() before dropping the runtime or application state. Embassy starts workers with start(spawner) from 'static storage and uses the same awaited shutdown API.

Call shutdown as part of application teardown, after stopping external sources of work such as HTTP servers, message consumers, or timers. This prevents the application from attempting publishes that shutdown will reject:

// First stop accepting new requests from outside the application.
server.stop_accepting().await;

// Then reject new mediator events and wait for accepted work to finish.
mediator.shutdown().await?;

// Mediator workers and runtime tasks have returned; dependencies may now close.
repository.close().await?;

Do not call shutdown from an event handler or a #[medi_task] belonging to the same mediator: it waits for that handler or task to return and would deadlock. Long-running tasks should declare signal: &ShutdownSignal and return when signal.cancelled().await completes (or check signal.is_shutdown() between short work units); otherwise shutdown().await continues waiting for them.

For Embassy, initialize the mediator in a StaticCell and call mediator.start(spawner). The Embassy integration requires embassy-executor 0.10 (with the platform feature appropriate for the target, such as platform-cortex-m). See the micro:bit example below.

Examples

Development

The repository uses mise for its Rust toolchain and commands:

mise install
mise run check-format
mise run lint
mise run test
mise run check-examples
mise run run-examples
mise run check-docs