# Events (v1 — JSON wire)
Typed publish/subscribe over **user-owned transports** (Kafka, NATS, etc.). Noema provides compile-time wiring and JSON glue; **you** own the consumer loop and broker connections.
## Features
Enable in `Cargo.toml`:
```toml
noema = { version = "0.3", features = ["events"] }
```
Depends on `core` + `di`. Wire format v1: **JSON** via `serde` only.
## Quick overview
| **Publish** | `EventPublishRaw::publish_raw(name, bytes)` | `EventPublisher<E>::publish(event)` → JSON + name |
| **Subscribe** | Consumer loop → `transport.dispatch(name, bytes).await` | Registry lookup → deserialize → `EventListener` handlers |
There is **no** crate-root `publish()`. Inject `Arc<dyn EventPublisher<E> + Send + Sync>` (or use your transport directly).
Noema does **not** call `start()` on your transport.
## Define an event
```rust
use noema::events::Event;
use noema::event;
use serde::{Deserialize, Serialize};
#[event(name = "order.created", description = "Order was created")]
#[derive(Serialize, Deserialize)]
struct OrderCreated {
order_id: u64,
}
```
`name = "..."` sets `OrderCreated::WIRE_NAME` — the stable identifier used in registry and `publish_raw`. Required for external transports.
In integration tests or other crates, import the attribute: `use noema::event;` (or `#[noema::event(...)]`).
## Transport struct
```rust
struct OrdersKafka {
spawner: Arc<dyn BackgroundSpawner + Send + Sync>,
errors: Arc<dyn BackgroundErrorHandler + Send + Sync>,
// your Kafka producer/consumer fields…
}
#[async_trait::async_trait]
impl EventPublishRaw for OrdersKafka {
async fn publish_raw(&self, name: &str, payload: &[u8]) -> noema::events::NoemaResult<()> {
// send to Kafka…
Ok(())
}
}
impl EventDispatcherContext for OrdersKafka {
fn dispatch_context(&self) -> DispatchContext {
DispatchContext::new(self.spawner.clone(), self.errors.clone())
}
}
```
### Spawner and error handler (receive only)
These fields are required for **`subscribe!` / dispatch**, not for publish.
| **`BackgroundSpawner`** | Runs each `EventListener::handle` in the background when `invoke_mode` is `Spawn` (default). |
| **`invoke_mode`** | `Spawn` (return before handlers finish) or `Await` (wait for `handle`). Set on `DispatchContext`. |
| **`BackgroundErrorHandler`** | `Spawn`: called from `EventListener::on_error` when `handle` returns `Err`, and on unknown wire names. `Await`: handler `Err` is returned from `dispatch` (the caller decides — e.g. a WebSocket reply). Invalid JSON in `receive` also returns `Err` from `dispatch`. |
`EventDispatcherContext` is **receive-side only**. Implement it on your transport; `subscribe!` checks the impl at compile time.
## DI — register the transport
Use standard DI (no events-specific macro):
```rust
noema::dependency!(singleton, OrdersKafka);
```
### Publish side (producer / API)
```rust
let publisher: Arc<dyn EventPublisher<OrderCreated> + Send + Sync> = resolve::<OrdersKafka>();
publisher.publish(OrderCreated { order_id: 1 }).await?;
```
### Receive side (consumer worker)
```rust
let dispatcher: Arc<dyn EventDispatch + Send + Sync> = resolve::<OrdersKafka>();
dispatcher.dispatch(&wire_name, &wire_payload).await?;
```
Both can coerce from the same `resolve::<OrdersKafka>()` singleton.
## Register publishers and subscribers
```rust
noema::publisher!(OrdersKafka: [OrderCreated, OrderShipped]);
noema::subscribe!(
OrdersKafka,
OrderCreated: [SendEmailHandler, AuditHandler],
OrderShipped: [NotifyHandler],
);
```
Rules:
- **One** `subscribe!(OrdersKafka, …)` batch per transport type per crate (second call → compile error).
- All events in one batch share the same transport type’s static registry.
- Handlers must implement `EventListener<E>` and `Injectable<Container>`.
- Transport must implement `EventDispatcherContext` before `subscribe!` (compile-time check).
## Subscribing
In your consumer loop:
```rust
loop {
let (name, payload) = kafka.recv().await?;
kafka.dispatch(&name, &payload).await?;
}
```
## Multiple transport instances
Use **newtypes** so each instance has its own registry:
```rust
struct OrdersKafka;
struct AnalyticsKafka;
noema::subscribe!(OrdersKafka, OrderCreated: [H1]);
noema::subscribe!(AnalyticsKafka, OrderCreated: [H2]);
```
## Migration from `dispatcher`
| `dispatcher` feature | `events` feature |
| `dispatch(event)` | `EventPublisher::publish` / `EventDispatch::dispatch` |
| `listen::<E, H>()` | `subscribe!(Transport, E: [H, …])` |
| `dispatcher_config()` | Per-transport spawner/error handler via `EventDispatcherContext` |
| Runtime handler `RwLock` | Compile-time static registry |
## Non-goals (v1)
- InProcess transport
- Protobuf / multi-codec
- Framework-owned consumer loops
- `inventory` / runtime registration
- `config_transports!` (use `dependency!` directly)