anycms-event
A type-safe, async event bus system for AnyCMS, built on tokio broadcast channels.
类型安全的异步事件总线,基于 tokio broadcast channels 实现,支持本地进程内通信与 Redis 跨进程分布式通信。
Features
- Type-safe — 编译期保证事件类型正确,无需手动类型转换
- Async-first — 基于 tokio 构建的异步 API,handler 以独立 tokio task 运行
- Thread-safe —
EventBus内部使用Arc<RwLock<>>,可安全跨 task/线程共享 - Derive Macro —
#[derive(Event)]自动实现 Event trait,零样板代码 event_bus!Macro — 一键定义事件结构体 + topic 分组 + 类型化总线- Wildcard Topics — 支持
*(单段匹配)和**(多段匹配)通配符 - Telemetry — 可插拔的遥测中间件,内置 tracing 和自定义实现
- Event Registry — 事件注册表,支持事件发现、查询和元数据管理
- Execution Log — 执行日志,追踪事件发布和 Handler 执行历史
- Trigger Rule Engine — 触发规则引擎,动态配置事件到动作的映射
- Testing Utilities —
EventCollector消除测试中的sleep()等待 - SSE Streaming — 通过
anycms-event-sse将事件实时推送到前端 - Redis Transport — 通过
anycms-event-redis实现跨进程事件传递 - Framework Integrations —
anycms-event-axum/anycms-event-actix简化框架集成
Workspace Structure
anycms-event/
├── src/ # 核心事件总线库
├── crates/
│ ├── anycms-event-derive/ # 过程宏(#[derive(Event)] + event_bus!)
│ ├── anycms-event-redis/ # Redis 传输层(分布式支持)
│ ├── anycms-event-axum/ # Axum 集成辅助
│ ├── anycms-event-actix/ # Actix-web 集成辅助
│ └── anycms-event-sse/ # SSE 实时推流
├── examples/ # 示例代码
└── tests/ # 集成测试
Quick Start
添加依赖到 Cargo.toml:
[]
= "0.1"
= { = "1", = ["full"] }
= { = "1", = ["derive"] }
手动定义事件
use *;
async
使用 #[derive(Event)] 宏
use *;
// 自动生成: event_name() -> "user.created", topic() -> "user"
通过属性自定义名称:
使用 event_bus! 宏
一键定义事件结构体、topic 分组和类型化事件总线:
use event_bus;
event_bus!
async
Core API
Event Trait
所有事件必须实现此 trait:
EventBus
| 方法 | 说明 |
|---|---|
EventBus::new() |
创建默认容量(1024)的事件总线 |
EventBus::builder() |
使用 Builder 模式创建,支持配置容量、遥测、注册表和执行日志 |
bus.publish(event).await |
发布事件(无订阅者时为空操作) |
bus.subscribe(handler).await |
订阅特定事件类型,handler 在独立 tokio task 中运行 |
bus.subscribe_pattern(pattern, handler).await |
使用通配符订阅,支持 * 和 ** |
bus.registry() |
获取事件注册表引用 |
bus.execution_log() |
获取执行日志引用(如果已配置) |
Builder 模式
use TracingTelemetry;
use ;
use Arc;
let bus = builder
.capacity
.telemetry
.execution_log
.build;
Topic 通配符
| 模式 | 含义 | 示例 |
|---|---|---|
user.* |
匹配单段 | 匹配 user.created,不匹配 user.foo.bar |
user.** |
匹配多段 | 匹配 user.created 和 user.foo.bar |
user.created |
精确匹配 | 仅匹配 user.created |
System Management 系统管理
anycms-event 提供三大系统管理模块,支持事件发现、执行追踪和动态触发规则配置。
P1: Event Registry 事件注册表
自动跟踪已注册事件,支持事件发现和元数据查询:
// 事件在 publish/subscribe 时自动注册
bus.publish.await.unwrap;
let registry = bus.registry;
// 列出所有已注册事件
for desc in registry.list_all
// 按条件查询
let user_events = registry.query;
// 手动注册带完整元数据的事件
registry.register;
P2: Execution Log 执行日志
记录和查询事件发布和 Handler 执行历史:
use ;
// 创建共享存储,Telemetry 和查询接口使用同一个后端
let log = new;
let bus = builder
.telemetry
.execution_log
.build;
// 发布事件后查询执行日志
bus.publish.await.unwrap;
let log = bus.execution_log.unwrap;
// 查询所有发布记录
let publishes = log.query;
// 查询失败的 Handler 执行
let failures = log.query;
P3: Trigger Rule Engine 触发规则引擎
动态配置事件到动作的映射规则,可与 WorkflowEngine 集成:
use ;
let trigger_engine = new;
// 注册动作处理器
trigger_engine.register_action;
trigger_engine.register_action;
// 动态配置触发规则
trigger_engine.add_rule;
// 带条件过滤的规则
trigger_engine.add_rule;
// 运行时管理规则
trigger_engine.disable_rule; // 禁用
trigger_engine.enable_rule; // 启用
trigger_engine.remove_rule; // 删除
trigger_engine.list_rules; // 列出所有
// 处理事件
let results = trigger_engine.process_event.await;
条件匹配操作符:
| 操作符 | 说明 | 示例 |
|---|---|---|
$eq |
等于 | {"status": {"$eq": "published"}} |
$ne |
不等于 | {"status": {"$ne": "draft"}} |
$gt / $gte |
大于 / 大于等于 | {"amount": {"$gt": 100}} |
$lt / $lte |
小于 / 小于等于 | {"amount": {"$lt": 1000}} |
$in |
包含在列表中 | {"category": {"$in": ["books", "tech"]}} |
$contains |
字符串包含 | {"title": {"$contains": "Rust"}} |
Telemetry 遥测
可插拔的遥测中间件,监控事件总线的发布/订阅生命周期:
use Telemetry;
use Duration;
// 自定义遥测实现
;
// 使用自定义遥测
let bus = builder
.telemetry
.build;
内置实现:
TracingTelemetry— 基于 tracing 的结构化日志NoopTelemetry— 空操作,用于禁用遥测ExecutionLogTelemetry— 将事件生命周期记录到执行日志
Testing 测试工具
EventCollector 消除测试中的 sleep() 等待,提供类型安全的事件断言。
启用方式:
[]
= { = "0.1", = ["testing"] }
use EventCollector;
use Duration;
async
可用的断言方法:
| 方法 | 说明 |
|---|---|
collect_now() |
返回当前已收集事件的快照 |
wait_for(count, timeout) |
异步等待直到收集到指定数量的事件 |
assert_count(expected) |
断言事件数量 |
assert_contains(predicate) |
断言包含满足条件的事件 |
assert_not_contains(predicate) |
断言不包含满足条件的事件 |
SSE 实时推流
通过 anycms-event-sse 将 EventBus 事件通过 Server-Sent Events 推送到前端:
[]
= "0.1"
use ;
// 创建 SSE 桥接器,注册需要推送的事件类型
let = new
.
.
.with_filter // 可选:只推送匹配的事件
.into_stream
.await;
// 在 Axum handler 中使用
use ;
async
内置过滤器:
AllowFilter— 白名单过滤DenyFilter— 黑名单过滤PatternFilter— 通配符过滤(复用 topic 匹配规则)
Framework Integration
Actix-Web
use ;
async
async
Axum
use ;
async
async
Framework Extractor Crates
anycms-event-axum 和 anycms-event-actix 提供 HasEventBus trait,统一框架集成模式:
use HasEventBus;
// 为你的类型化总线实现 trait
Redis 分布式通信
通过 anycms-event-redis 实现跨进程事件传递:
[]
= "0.1"
use RedisTransport;
async
完整示例见 examples/redis_distributed.rs。
Examples
| 示例 | 说明 | 运行命令 |
|---|---|---|
| 基础用法 | 手动定义事件和 Event trait | cargo run --example basic_usage |
| Actix-Web 集成 | 在 Actix-Web 中共享 EventBus | cargo run --example actix_integration |
| Axum 集成 | 在 Axum 中共享 EventBus | cargo run --example axum_integration |
| Redis 分布式 | 跨进程事件通信(需 Redis) | cargo run --example redis_distributed |
| Topic 通配符 | 通配符订阅示例 | cargo run --example topic_subscription |
| 错误处理 | Handler 错误处理示例 | cargo run --example error_handling |
| 测试工具 | EventCollector 收集和断言事件 | cargo run --example testing_collector --features testing |
| 遥测中间件 | 自定义 Telemetry + Builder 模式 | cargo run --example telemetry |
| SSE + Axum | SSE 实时推流 + Axum 集成 | cargo run --example sse_axum |
| SSE + Actix | SSE 实时推流 + Actix-web 集成 | cargo run --example sse_actix |
| 触发规则引擎 | 系统管理功能:Registry + Trigger Engine | cargo run --example trigger_workflow |
Error Handling
Handler 返回错误时会通过 tracing 记录日志,不会影响其他订阅者。如果配置了 Telemetry,错误信息会通过 on_handler_complete 回调上报。
use EventBusError;
match bus.publish.await
Architecture
┌──────────────────────────────────────────────────────────┐
│ Trigger Rule Engine (P3) │
│ 规则 CRUD / 事件模式匹配 / JSON 条件过滤 / 动作执行 │
│ register_action("workflow", |ctx| WorkflowEngine.emit()) │
├──────────────────────────────────────────────────────────┤
│ Event Registry (P1) Execution Log (P2) │
│ 事件发现/查询/搜索/元数据 执行记录/状态/耗时追踪 │
│ bus.registry().query(...) bus.execution_log() │
├──────────────────────────────────────────────────────────┤
│ Core EventBus │
│ publish / subscribe / pattern matching / telemetry │
│ derive macros / SSE / Redis transport │
└──────────────────────────────────────────────────────────┘
License
MIT