Maiko is built around a small set of core abstractions that work together:
- Events are messages that flow through the system
- Topics determine which actors receive which events (actors subscribe to topics)
- Actors are independent units that process events and maintain state
- Context allows actors to send events and interact with the runtime
- Supervisor manages actor lifecycles and coordinates the system
- Envelopes wrap events with metadata for tracing and parent-child linking
This document covers each abstraction in detail.
Events are messages that flow through the system. They must implement the Event trait (there is derive macro for convenience):
#[derive(Event, Debug)]
enum NetworkEvent {
PacketReceived(Vec<u8>),
ConnectionClosed(u32),
Error(String),
}Events are:
- Cloneable - events may be delivered to multiple actors
- Send + Sync - as they travel between tasks and threads.
Events should be:
- Debuggable - for logging and diagnostics
Topics group events and route them to interested actors. Each event maps to exactly one topic, determined by the Topic::from_event() implementation.
Define custom topics for fine-grained routing:
#[derive(Debug, Hash, Eq, PartialEq, Clone)]
enum NetworkTopic {
Ingress,
Egress,
Control,
}
impl Topic<NetworkEvent> for NetworkTopic {
fn from_event(event: &NetworkEvent) -> Self {
match event {
NetworkEvent::PacketReceived(_) => NetworkTopic::Ingress,
NetworkEvent::ConnectionClosed(_) => NetworkTopic::Control,
NetworkEvent::Error(_) => NetworkTopic::Control,
}
}
}Each topic defines what happens when a subscriber's channel is full, via overflow_policy():
impl Topic<NetworkEvent> for NetworkTopic {
fn from_event(event: &NetworkEvent) -> Self { /* ... */ }
fn overflow_policy(&self) -> OverflowPolicy {
match self {
NetworkTopic::Ingress => OverflowPolicy::Block, // wait for space
NetworkTopic::Control => OverflowPolicy::Block, // commands must arrive
NetworkTopic::Egress => OverflowPolicy::Drop, // discard if slow
}
}
}The default is Fail - the subscriber's channel is closed and the actor terminates. This surfaces problems immediately. See OverflowPolicy for details.
Use DefaultTopic when you don't need routing - all events go to all subscribed actors:
sup.add_actor("processor", factory, &[DefaultTopic])?;Actors are independent units that process or produce events. They implement the Actor trait with two core methods:
handle_event- Process incoming eventsstep- Produce events or perform periodic work
struct PacketProcessor {
ctx: Context<NetworkEvent>,
buffer: Vec<u8>,
}
impl Actor for PacketProcessor {
type Event = NetworkEvent;
async fn handle_event(&mut self, envelope: &Envelope<Self::Event>) -> Result {
match envelope.event() {
NetworkEvent::PacketReceived(data) => {
self.buffer.extend(data);
}
_ => {}
}
Ok(())
}
async fn step(&mut self) -> Result<StepAction> {
if self.buffer.len() > 1000 {
self.flush_buffer().await?;
}
Ok(StepAction::AwaitEvent)
}
}Concurrent tasks typically fall into one of three roles: sources that produce data, sinks that consume it, and processors that do both. Maiko doesn't have separate types for these - every actor implements the same Actor trait. The role emerges from which methods you use and what you subscribe to:
| Role | step() |
handle_event() |
Subscriptions |
|---|---|---|---|
| Source | Produces events | No-op | Subscribe::none() |
| Sink | Never (default) |
Consumes events | Topics it cares about |
| Processor | Optional | Receives events, sends new ones | Selective topics |
A temperature sensor is a source - it uses step() to emit readings and subscribes to nothing. A logger is a sink - it handles events but never sends. An alerter is a processor - it receives readings and emits alerts.
This is a deliberate design choice. A single trait keeps the API small and lets actors evolve. A sink that later needs to emit events just adds a Context field - no type change, no rewiring.
Actors can implement optional lifecycle methods:
on_start- Called once when the actor startson_shutdown- Called during graceful shutdownon_error- Handle errors (swallow or propagate)
async fn on_start(&mut self) -> Result {
println!("Actor starting...");
Ok(())
}
fn on_error(&mut self, error: Error) -> Result {
eprintln!("Error: {}", error);
Ok(()) // Swallow error, continue running
}The Context provides actors with capabilities to interact with the system. Context is optional - actors that only consume events (without sending) don't need to store it:
// Pure consumer - no context needed
struct Logger;
impl Actor for Logger {
type Event = MyEvent;
async fn handle_event(&mut self, envelope: &Envelope<Self::Event>) -> Result {
println!("Received: {:?}", envelope.event());
Ok(())
}
}
// Producer/processor - stores context to send events
struct Processor {
ctx: Context<MyEvent>,
}Context capabilities:
// Send events
ctx.send(NetworkEvent::PacketReceived(data)).await?;
// Send child event (linked to parent for tracing)
ctx.send_child_event(ResponseEvent::Ok, envelope.id()).await?;
// Stop this actor (other actors continue)
ctx.stop();
// Shut down the entire runtime
ctx.stop_runtime();
// Get actor's name
let name = ctx.actor_name();The Supervisor manages actor lifecycles and provides registration APIs.
let mut sup = Supervisor::<NetworkEvent, NetworkTopic>::default();
sup.add_actor("ingress", |ctx| IngressActor::new(ctx), &[NetworkTopic::Ingress])?;
sup.add_actor("egress", |ctx| EgressActor::new(ctx), &[NetworkTopic::Egress])?;
sup.add_actor("monitor", MonitorActor::new, Subscribe::all())?; // no closure neededUse build_actor when an actor needs non-default configuration (e.g. a larger channel):
sup.build_actor("writer", |ctx| Writer::new(ctx))
.topics(&[NetworkTopic::Ingress])
.channel_capacity(512)
.build()?;See Advanced Topics - Per-Actor Config for details.
The terminal methods (run, join, stop) consume the supervisor, preventing
use-after-shutdown at compile time.
// Send events before shutdown (Supervisor as actor)
sup.send(MyEvent::Data(42)).await?;
// Start all actors (non-blocking) — borrows &mut self
sup.start().await?;
// Wait for completion — consumes the supervisor
sup.join().await?;
// Or combine start + join in one call
sup.run().await?;
// Graceful shutdown — consumes the supervisor
sup.stop().await?;The step() method returns StepAction to control scheduling:
| Action | Behavior |
|---|---|
StepAction::Continue |
Run step again immediately |
StepAction::Yield |
Yield to runtime, then run again |
StepAction::AwaitEvent |
Pause until next event arrives |
StepAction::Backoff(Duration) |
Sleep, then run again |
StepAction::Never |
Disable step permanently (default) |
Event producer with interval:
async fn step(&mut self) -> Result<StepAction> {
self.ctx.send(HeartbeatEvent).await?;
Ok(StepAction::Backoff(Duration::from_secs(5)))
}External I/O source:
async fn step(&mut self) -> Result<StepAction> {
let data = self.websocket.read().await?;
self.ctx.send(DataEvent(data)).await?;
Ok(StepAction::Continue)
}Pure event processor (default):
async fn step(&mut self) -> Result<StepAction> {
Ok(StepAction::Never)
}ActorId uniquely identifies a registered actor. It is returned by Supervisor::add_actor() and used throughout the system for:
- Event metadata - every envelope carries the sender's
ActorId - Test assertions - verify which actors sent/received events
- Parent tracking - trace event causality between actors
// Returned when registering an actor
let producer: ActorId = sup.add_actor("producer", |ctx| Producer::new(ctx), &[DefaultTopic])?;
// Access sender from event metadata
async fn handle_event(&mut self, envelope: &Envelope<Self::Event>) -> Result {
let sender: &ActorId = envelope.meta().actor_id();
println!("Event from: {}", sender.as_str());
Ok(())
}
// Use in test harness
assert!(test.event(id).was_delivered_to(&consumer));ActorId supports serialization (with the serde feature) and works correctly across IPC boundaries. Equality uses string comparison with a fast-path for pointer equality when IDs share the same memory allocation:
// Same-process: pointer comparison (O(1))
// Cross-process (after deserialization): string comparison
assert_eq!(local_id, deserialized_id); // Works correctly in both casesThis makes ActorId suitable for distributed scenarios where actor references need to survive serialization.
Events arrive wrapped in an Envelope containing metadata:
async fn handle_event(&mut self, envelope: &Envelope<Self::Event>) -> Result {
// Access the event
let event = envelope.event();
// Access metadata
let sender = envelope.meta().actor_name();
let event_id = envelope.meta().id();
let parent = envelope.meta().parent_id();
// Send child event linked to this one
self.ctx.send_child_event(ResponseEvent::Ok, envelope.id()).await?;
Ok(())
}Parent tracking enables tracing event causality through the system.