Event-Driven State Machine - 去中心化的事件驱动状态机引擎,pip install 即装即用
Project description
lzm-edsm
编码:utf8 | 作者:Lzm | 日期:2026-07-15
Event-Driven State Machine — 去中心化的事件驱动状态机引擎。
一句话:状态转移 = 事件发布,事件消费 = 状态响应。
lzm-edsm 将状态机、事件总线、蓝图合并成一个统一的 Python 包,pip install 即装即用,不绑定任何框架或中间件。
安装
# 核心(含 SQLiteStore 事件溯源)
pip install lzm-edsm
# 带 Redis 传输
pip install lzm-edsm[redis]
# 带 RabbitMQ 传输
pip install lzm-edsm[rabbitmq]
# 组合安装(如 RabbitMQ + 默认 SQLite)
pip install lzm-edsm[rabbitmq]
# 全量
pip install lzm-edsm[all]
说明:
aiosqlite是核心依赖,安装后自动获得 SQLiteStore 事件溯源审计功能。
快速开始
1. 定义域(Domain)
Domain 是状态、转移、事件的统一容器:
from lzm.edsm import Domain
# 创建域
user = Domain("user_status", initial="PENDING", description="用户状态管理")
# 注册状态
user.state("PENDING", label="待审批")
user.state("ACTIVE", label="正常")
user.state("DISABLED", label="停用")
user.state("REJECTED", label="拒绝")
# 注册转移(自动派生事件名:user_status.approve)
user.transition("approve", "PENDING", "ACTIVE")
user.transition("reject", "PENDING", "REJECTED")
user.transition("disable", "ACTIVE", "DISABLED")
user.transition("enable", "DISABLED", "ACTIVE")
# 注册守卫(可选)
@user.transition("approve", "PENDING", "ACTIVE")
async def approve_guard(payload, ctx):
"""只有 admin 可以审批"""
if payload.get("role") != "admin":
return False, "无操作权限"
return True, None
# 注册事件监听器
@user.on("approve")
async def send_welcome_notification(event):
await send_email(event.payload["user_id"], "欢迎加入!")
2. 创建引擎(Engine)
Engine 是状态机执行器,核心是 7 步 Pipeline:
from lzm.edsm import Engine
engine = Engine() # 默认 InProcessTransport
engine.register(user)
await engine.start()
# 触发状态转移
new_state = await engine.trigger(
domain="user_status",
current="PENDING",
event="approve",
actor="admin:1001",
payload={"user_id": 123, "role": "admin"},
)
print(new_state) # → "ACTIVE"
await engine.stop()
3. 多监听器并发
@user.on("approve")
async def sync_to_search_engine(event):
await es.index("user", event.payload)
@user.on("approve")
async def log_audit(event):
await db.insert("audit_log", {"event": event.name, "payload": event.payload})
所有 @domain.on("approve") 的监听器会并发执行(asyncio.gather)。
功能清单
🧠 核心引擎
| 功能 | 说明 |
|---|---|
| Event-Driven State Machine | 状态转移即事件发布,事件消费即状态响应 |
| Pipeline 8 步管线 | 校验 → Guard → Before Hook → 转移 → After Hook → Middleware → 发布事件 → 远程分发 |
| 多 Domain 隔离 | 同一 Engine 注册多个 Domain,命名空间和事件名完全隔离 |
| 热重载 | reload_config() + ConfigWatcher 文件轮询,保留现有 guard/handler |
| 异步去重 | 基于 event_id + 5s TTL 的本地缓存,防止远程事件重复消费 |
| 可选 Result 返回 | Handler 可选返回 Result(allowed, reason, data),None 完全向后兼容 |
🗺️ 域定义(Domain)
| 功能 | 说明 |
|---|---|
| 状态注册 | domain.state(name, label, data_default, **meta),支持默认上下文数据 |
| 转移注册 | domain.transition(event, source, target, description),支持多源状态 |
| 守卫(Guard) | @domain.transition(event, source, target) 装饰器注册,返回 (bool, reason) |
| 前置/后置钩子(Hook) | Guard 之后/转移之前,转移之后/事件发布之前 |
| 事件监听器(Handler) | @domain.on(event, include_data) 注册,支持三态数据控制 |
| 热更新 API | to_dict() / from_dict() / update_from_dict() 序列化与增量更新 |
| 状态/转移管理 | get_state() / get_transition() / remove_state() / remove_transition() |
📦 事件系统(Event)
| 功能 | 说明 |
|---|---|
| UUID v7 ID | 时间有序,DB 索引友好,符合 RFC 9562 |
| 毫秒时间戳 | 自动生成(东八区) |
| 序列化/反序列化 | to_dict() / from_dict() 跨进程传输 |
| 状态机上下文注入 | event.metadata 自动注入 domain/current_state/target_state |
🛡️ 数据流控制
| 功能 | 说明 |
|---|---|
| 全局开关 | Engine(with_data=True/False) 控制 payload/data 透传 |
| Handler 级覆盖 | @domain.on("event", include_data=True/False/None) 三态控制 |
| 状态级默认数据 | domain.state("S", data_default={...}) 提供上下文数据回退链 |
🔌 传输层(Transport)
| 功能 | 说明 |
|---|---|
| InProcess(默认) | 进程内,零依赖,纳秒级延迟 |
| Redis 广播 | Pub/Sub 模式,所有 Worker 全量接收 |
| Redis 竞争 | Stream + Consumer Group 模式,仅一个 Worker 消费 |
| RabbitMQ 广播 | Fanout Exchange 模式,独占队列 |
| RabbitMQ 竞争 | Topic Exchange 模式,共享队列 + routing_key # |
| 消费模式枚举 | ConsumeMode.BROADCAST / ConsumeMode.COMPETING |
⚙️ 动态配置
| 功能 | 说明 |
|---|---|
| 字典批量注册 | engine.register_dict(config) |
| 配置热重载 | engine.reload_config(config) 返回 {added, updated, events_changed} |
| 文件监听 | ConfigWatcher(engine, file_path) 自动轮询文件 mtime |
| 自定义适配器 | ConfigAdapter 支持 JSON/YAML/自定义格式 |
🔍 管理自省(Management)
| 功能 | 说明 |
|---|---|
| 域概览 | list_domains() 返回 name/state_count/transition_count/handler_count |
| 域详情 | inspect(name) 返回 states/transitions/handlers 完整结构 |
| 事件列表 | list_events() 全部已注册事件名(排序) |
| 可用事件 | allowed_events(domain, state) 指定状态下允许的事件 |
| 蓝图导出 | export() 完整蓝图 JSON 序列化(含版本号/时间戳) |
📝 事件溯源(SQLiteStore)
| 功能 | 说明 |
|---|---|
| 自动存储 | store.save(event) 异步写入,不阻塞主流程 |
| 按条件查询 | query(domain, event_name, actor, time_range, limit, offset) |
| 实体溯源 | get_entity_history(entity_id) payload 模糊匹配 |
| 事件统计 | count(domain, event_name) 数量统计 |
| TTL 自动清理 | ttl_days=90 可选,写入后自动清理过期事件 |
传输层(Transport)
lzm-edsm 支持 3 种传输后端,进程内零依赖,跨进程可选 Redis/RabbitMQ。
消费模式
每个 Transport 支持两种消费模式(v0.1.3+):
| 模式 | 枚举值 | 说明 | 适用场景 |
|---|---|---|---|
| 广播 | ConsumeMode.BROADCAST |
所有 Worker 全量接收事件(默认) | 配置变更、通知推送、状态同步 |
| 竞争 | ConsumeMode.COMPETING |
仅一个 Worker 消费事件 | 任务分发、竞态敏感的状态机转移 |
from lzm.edsm.transports.base import ConsumeMode
各传输层的竞争模式实现机制:
| 传输层 | 广播实现 | 竞争实现 |
|---|---|---|
| InProcess | 内存广播 | — |
| Redis | Pub/Sub | Stream + Consumer Group |
| RabbitMQ | Fanout Exchange | Direct Exchange + 共享队列 |
InProcess(默认)
零依赖,进程内事件广播:
from lzm.edsm import Engine
engine = Engine() # 默认 InProcessTransport
Redis
跨进程事件传输,支持双模式:
from lzm.edsm import Engine
from lzm.edsm.transports.redis import RedisTransport
from lzm.edsm.transports.base import ConsumeMode
# 广播模式(Pub/Sub,默认)—— 所有 Worker 全量接收
transport = RedisTransport(
url="redis://localhost:6379",
channel_prefix="edsm",
seen_ttl=5, # 去重 TTL(秒)
)
engine = Engine(transport=transport)
await engine.start()
# 竞争模式(Stream + Consumer Group)—— 仅一个 Worker 消费
transport = RedisTransport(
url="redis://localhost:6379",
mode=ConsumeMode.COMPETING,
channel_prefix="edsm",
consumer_group="edsm-workers",
consumer_id="worker-1",
)
安装:
pip install lzm-edsm[redis]
RabbitMQ
跨进程事件传输,支持双模式:
from lzm.edsm import Engine
from lzm.edsm.transports.rabbitmq import RabbitMQTransport, RabbitMode
# 广播模式(Fanout Exchange,默认)
transport = RabbitMQTransport(
url="amqp://guest:guest@localhost:5672/",
exchange_prefix="edsm",
seen_ttl=5,
)
# 竞争模式(Direct Exchange + 共享队列)
transport = RabbitMQTransport(
url="amqp://guest:guest@localhost:5672/",
mode=RabbitMode.COMPETING,
exchange_prefix="edsm",
)
安装:
pip install lzm-edsm[rabbitmq]
事件溯源审计(SQLiteStore)
SQLiteStore 是内置的事件持久化组件,用于审计和溯源(已包含在核心依赖中)。
使用示例:
from lzm.edsm import Engine
from lzm.edsm.contrib.store import SQLiteStore
# 初始化存储
store = SQLiteStore(
db_path="events.db", # 或 ":memory:"
ttl_days=90, # 保留 90 天
auto_cleanup=True, # 自动清理过期事件
)
await store.start()
# 在监听器中保存事件
@user.on("approve")
async def save_for_audit(event):
await store.save(event)
# 查询事件
events = await store.query(domain="user_status", limit=10)
# 实体溯源
history = await store.get_entity_history("user:123", limit=50)
# 统计
count = await store.count(domain="user_status")
await store.stop()
注意:SQLiteStore 不会无限增长:
- 默认保留 90 天(
ttl_days=90) - 每次写入后自动清理过期事件
- 可通过
ttl_days=None禁用 TTL
完整示例
import asyncio
from lzm.edsm import Domain, Engine
# 1. 定义域
order = Domain("order_status", initial="CREATED", description="订单状态")
order.state("CREATED", label="已创建")
order.state("PAID", label="已支付")
order.state("SHIPPED", label="已发货")
order.state("COMPLETED", label="已完成")
order.state("CANCELLED", label="已取消")
order.transition("pay", "CREATED", "PAID")
order.transition("ship", "PAID", "SHIPPED")
order.transition("complete", "SHIPPED", "COMPLETED")
order.transition("cancel", ["CREATED", "PAID"], "CANCELLED")
# 监听支付事件
@order.on("pay")
async def on_paid(event):
print(f"订单 {event.payload['order_id']} 已支付,金额:{event.payload['amount']}")
@order.on("pay")
async def notify_warehouse(event):
print(f"通知仓库发货:订单 {event.payload['order_id']}")
# 2. 创建引擎
engine = Engine()
engine.register(order)
await engine.start()
# 3. 触发转移
async def main():
# 支付
state = await engine.trigger(
domain="order_status",
current="CREATED",
event="pay",
actor="user:1001",
payload={"order_id": "ORD-001", "amount": 299.0},
)
print(f"当前状态:{state}") # PAID
# 发货
state = await engine.trigger(
domain="order_status",
current="PAID",
event="ship",
actor="system",
payload={"order_id": "ORD-001"},
)
print(f"当前状态:{state}") # SHIPPED
await engine.stop()
asyncio.run(main())
管理接口(Management)
engine.management 提供域的管理/自省能力,无需额外安装。
域概览
# 查看所有已注册域
info_list = engine.management.list_domains()
for info in info_list:
print(info.domain, info.state_count, info.transition_count, info.handler_count)
# 输出:
# user_status 4 3 1
# connector_task 3 2 1
DomainInfo 字段:
| 字段 | 类型 | 说明 |
|---|---|---|
domain |
str | 域标识 |
description |
str | 描述 |
initial |
str | 初始状态 |
state_count |
int | 状态数 |
transition_count |
int | 转移数 |
handler_count |
int | 监听器数 |
域详情
# 查看指定域的完整结构
detail = engine.management.inspect("user_status")
print(detail.states)
# → {"PENDING": {"label": "待审批"}, "ACTIVE": {"label": "正常"}, ...}
print(detail.transitions)
# → {"approve": {"source": "PENDING", "target": "ACTIVE", "has_guard": True}, ...}
print(detail.handlers)
# → {"approve": ["send_welcome_notification", "sync_to_search_engine"]}
事件查询
# 所有已注册事件
events = engine.management.list_events()
# → ["connector_task.complete", "user_status.approve", "user_status.reject", ...]
# 查看某个状态的可用事件
allowed = engine.management.allowed_events("user_status", "PENDING")
# → [{"event": "approve", "source": "PENDING", "target": "ACTIVE", "description": ""},
# {"event": "reject", "source": "PENDING", "target": "REJECTED", "description": ""}]
蓝图导出
# 导出完整蓝图(JSON 可序列化)
blueprint = engine.management.export()
# → {
# "version": "0.1.3",
# "generated_at": 1719294000.123,
# "domains": [
# {
# "domain": "user_status",
# "description": "用户状态管理",
# "initial": "PENDING",
# "states": ["PENDING", "ACTIVE", "DISABLED", "REJECTED"],
# "transitions": [
# {"event": "approve", "source": "PENDING", "target": "ACTIVE", "has_guard": True},
# ...
# ],
# "handler_count": 2,
# },
# ...
# ],
# "total_domains": 2,
# "total_events": 7,
# }
API 参考
Domain
Domain(name: str, initial: str, description: str = "")
# 注册状态
domain.state(name: str, label: str = "", **metadata)
# 注册转移
domain.transition(event: str, source: str | list[str], target: str)
# 注册守卫(装饰器)
@domain.transition(event, source, target)
async def guard(payload: dict, ctx: dict) -> tuple[bool, str | None]
# 注册事件监听器
@domain.on(event: str)
async def handler(event: Event)
Engine
Engine(transport: Transport | None = None)
# 注册域
engine.register(domain: Domain)
# 注销域
engine.unregister(domain_name: str)
# 启动
await engine.start()
# 触发转移
await engine.trigger(
domain: str,
current: str,
event: str,
actor: str = "system",
payload: dict = {},
trace_id: str = "",
ctx: dict = {},
) -> str # 返回新状态
# 停止
await engine.stop()
# 管理接口
engine.management # → Management 实例
engine.management.list_domains() # 域概览
engine.management.inspect(name) # 域详情
engine.management.list_events() # 事件列表
engine.management.allowed_events(...) # 可用事件
engine.management.export() # 蓝图导出
Event
Event(
actor: str, # 触发者("type:id")
name: str, # 事件名(自动生成或手动指定)
payload: dict = {}, # 业务数据
trace_id: str = "", # 链路追踪
)
# 属性
event.id # UUID v7(时间有序)
event.name # 事件名(含 domain 前缀)
event.timestamp # 毫秒时间戳
event.actor # 触发者
event.payload # 业务数据
event.data # 状态上下文数据(with_data=True 时有效)
event.trace_id # 链路追踪
event.metadata # 状态机上下文(Engine 自动注入)
# 序列化
event.to_dict() # dict
Event.from_dict(data) # 反序列化
Result
Result(
allowed: bool = True, # 是否允许继续
reason: str | None = None, # 拒绝/失败原因
data: dict | None = None, # 处理输出数据
)
# 用法:Handler 可选返回 Result
@domain.on("publish")
async def handler(event: Event) -> Result:
approval_id = await process(event.data)
return Result(data={"approval_id": approval_id})
异常
from lzm.edsm.core.exceptions import (
EDsmError, # 基类
DomainNotFound, # 域未注册
TransitionNotAllowed, # 非法状态转移
GuardRejected, # 守卫拦截
EventPublishError, # 事件发布失败
)
技术栈
| 层 | 技术 | 版本 |
|---|---|---|
| 语言 | Python | >= 3.11 |
| 核心(零依赖) | 纯 asyncio | — |
| 传输:Redis | redis-py | >= 5.0(可选) |
| 传输:RabbitMQ | aio-pika | >= 9.0(可选) |
| 存储:SQLite | aiosqlite | 内置 |
| 测试 | pytest | >= 8.0 |
目录结构
lzm-edsm/
├── pyproject.toml
├── src/lzm/edsm/
│ ├── __init__.py # 导出 Domain, Engine, Event, Result
│ ├── core/
│ │ ├── domain.py # Domain 类
│ │ ├── engine.py # Engine 类(Pipeline)
│ │ ├── event.py # Event 数据结构
│ │ ├── result.py # Result 数据类
│ │ └── exceptions.py # 异常定义
│ ├── transports/
│ │ ├── base.py # Transport 抽象基类 + ConsumeMode 枚举
│ │ ├── inprocess.py # 进程内(默认)
│ │ ├── redis.py # Redis Pub/Sub(广播)+ Stream(竞争)
│ │ └── rabbitmq.py # RabbitMQ(双模式)
│ └── contrib/
│ ├── store.py # SQLiteStore 事件溯源
│ └── _uuid7.py # UUID v7 生成器
├── tests/
├── examples/
└── docs/
相关文档
许可证
MIT License
标签
#event-driven #state-machine #asyncio #python #edsm #lzm
Project details
Download files
Download the file for your platform. If you're not sure which to choose, learn more about installing packages.
Source Distribution
Built Distribution
Filter files by name, interpreter, ABI, and platform.
If you're not sure about the file name format, learn more about wheel file names.
Copy a direct link to the current filters
File details
Details for the file lzm_edsm-0.2.0.tar.gz.
File metadata
- Download URL: lzm_edsm-0.2.0.tar.gz
- Upload date:
- Size: 56.3 kB
- Tags: Source
- Uploaded using Trusted Publishing? No
- Uploaded via: twine/6.2.0 CPython/3.12.10
File hashes
| Algorithm | Hash digest | |
|---|---|---|
| SHA256 |
2de497506e16c0dba7c010c1d80c5ae01c78415bc29be6aa859db38838f5c918
|
|
| MD5 |
3598e37c8ef4c11ffe090427a623948e
|
|
| BLAKE2b-256 |
e70eae3934ae774334989bc872e168182a8dd353f7fb34fd24fccefdb750e6aa
|
File details
Details for the file lzm_edsm-0.2.0-py3-none-any.whl.
File metadata
- Download URL: lzm_edsm-0.2.0-py3-none-any.whl
- Upload date:
- Size: 39.1 kB
- Tags: Python 3
- Uploaded using Trusted Publishing? No
- Uploaded via: twine/6.2.0 CPython/3.12.10
File hashes
| Algorithm | Hash digest | |
|---|---|---|
| SHA256 |
2ee847ce772317aae1b91a5bbbb26aa5da4d4986fea16b9a11a5a4694211dcc7
|
|
| MD5 |
e3f9f246da74753f81fe60d553af7062
|
|
| BLAKE2b-256 |
99987a2761dd4b072f29a8f13acef26207237b79d2eea49850a2731c4982c1e5
|