149 lines
4.3 KiB
Markdown
149 lines
4.3 KiB
Markdown
# event_bus
|
|
|
|
Un bus d'événements asynchrone en mémoire pour Tokio, entièrement basé sur le typage fort en Rust (`TypeId`).
|
|
|
|
---
|
|
|
|
## 🎯 Problématique résolue
|
|
|
|
Dans l'implémentation initiale (`event_bus`), le modèle était inspiré de JavaScript (topics basés sur des chaînes de caractères) :
|
|
- Les événements étaient transportés via un pointeur générique `AnyEvent` (`Arc<dyn Any + Send + Sync>`).
|
|
- L'émission imposait de spécifier un topic string (`bus.emit("topic", event)`).
|
|
- La réception nécessitait de spécifier à la fois le type et le topic string (`bus.on_async::<Event, _, _>("topic", ...)`), puis d'effectuer un déréférencement / downcast dynamique (`downcast_ref::<T>()` ou macro `match_event!`) sur chaque message reçu.
|
|
|
|
**`event_bus` résout entièrement cette complexité :**
|
|
- **Typage fort natif** : le routage est directement effectué par l'identifiant de type (`std::any::TypeId`), sans nom de topic requis.
|
|
- **Zéro downcast / déréférencement à la réception** : le callback reçoit directement la structure d'événement typée.
|
|
- **Syntaxe ergonomique** : support complet de la syntaxe turbofish demandée `bus.on_async::<MessageUpdatedEvent>(|event| async move { ... })` ainsi que de l'inférence automatique `bus.on_async(|event: MessageUpdatedEvent| async move { ... })`.
|
|
- **Source unique de vérité** : plus besoin de `Scope` externe ; le contexte (ex: `channel_id`, `server_id`, `caller_id`) est directement transporté dans les champs de la structure typée.
|
|
- **Targeted wake-up** : canaux Tokio `broadcast` isolés par type d'événement, garantissant des performances optimales sans réveil inutile de tâches.
|
|
|
|
---
|
|
|
|
## 🚀 Utilisation
|
|
|
|
### 1. Définir des événements
|
|
|
|
N'importe quelle structure Rust implémentant `Clone + Send + Sync + 'static` est automatiquement un `Event` valide (aucun macro derive supplémentaire nécessaire) :
|
|
|
|
```rust
|
|
use uuid::Uuid;
|
|
|
|
#[derive(Clone, Debug)]
|
|
pub struct MessageCreatedEvent {
|
|
pub server_id: Option<Uuid>,
|
|
pub channel_id: Uuid,
|
|
pub content: String,
|
|
}
|
|
|
|
#[derive(Clone, Debug)]
|
|
pub struct MessageUpdatedEvent {
|
|
pub id: u64,
|
|
pub content: String,
|
|
}
|
|
```
|
|
|
|
---
|
|
|
|
### 2. Émission d'événements
|
|
|
|
```rust
|
|
use event_bus::EventBus;
|
|
|
|
let bus = EventBus::new();
|
|
|
|
// Émission typée directe
|
|
bus.emit(MessageUpdatedEvent {
|
|
id: 42,
|
|
content: "Nouveau message".into(),
|
|
});
|
|
```
|
|
|
|
---
|
|
|
|
### 3. Réception asynchrone (`on_async`)
|
|
|
|
Syntaxe turbofish exacte demandée :
|
|
|
|
```rust
|
|
bus.on_async::<MessageUpdatedEvent>(|event| async move {
|
|
// `event` est directement de type MessageUpdatedEvent
|
|
println!("Message {} mis à jour : {}", event.id, event.content);
|
|
});
|
|
```
|
|
|
|
Ou avec inférence sur l'argument de fermeture :
|
|
|
|
```rust
|
|
bus.on_async(|event: MessageUpdatedEvent| async move {
|
|
println!("Contenu : {}", event.content);
|
|
});
|
|
```
|
|
|
|
---
|
|
|
|
### 4. Réception synchrone (`on`)
|
|
|
|
```rust
|
|
bus.on::<MessageCreatedEvent>(|event| {
|
|
println!("Nouveau message créé sur le salon : {:?}", event.channel_id);
|
|
});
|
|
```
|
|
|
|
---
|
|
|
|
### 5. Utilisation avec contexte injecté (`on_async_with`)
|
|
|
|
Pratique pour passer des services ou repositories sans clones manuels répétés :
|
|
|
|
```rust
|
|
bus.on_async_with::<MessageCreatedEvent, _>(router, |router, event| async move {
|
|
router.gateway.send(...);
|
|
});
|
|
```
|
|
|
|
---
|
|
|
|
### 6. Écouteurs One-Shot (`wait_next` et `wait_for`)
|
|
|
|
Permet d'attendre un événement de manière linéaire avec une `Future` (sans boucle manuelle ni fuite de souscription) :
|
|
|
|
```rust
|
|
// Attend le tout prochain événement de ce type
|
|
let event = bus.wait_next::<MessageCreatedEvent>().await?;
|
|
|
|
// Ou attend un événement répondant à une condition précise
|
|
let confirmed = bus.wait_for::<MessageSavedEvent>(|e| e.id == target_id).await?;
|
|
```
|
|
|
|
---
|
|
|
|
### 7. Flux direct / Récepteur sans callback (`subscribe`)
|
|
|
|
Si vous préférez consommer les événements dans votre propre boucle de streaming :
|
|
|
|
```rust
|
|
let mut rx = bus.subscribe::<MessageUpdatedEvent>();
|
|
|
|
tokio::spawn(async move {
|
|
while let Ok(event) = rx.recv().await {
|
|
// `event` est directement MessageUpdatedEvent, aucun `match_event!` requis !
|
|
println!("Reçu : {}", event.content);
|
|
}
|
|
});
|
|
```
|
|
|
|
---
|
|
|
|
## 🧪 Tests et Benchmarks
|
|
|
|
Exécuter les tests du crate :
|
|
```bash
|
|
cargo test --manifest-path event_bus/Cargo.toml
|
|
```
|
|
|
|
Exécuter les benchmarks Criterion :
|
|
```bash
|
|
cargo bench --manifest-path event_bus/Cargo.toml
|
|
```
|