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]
# 带 MQTT 传输
pip install lzm-edsm[mqtt]
# 带 Kafka 传输
pip install lzm-edsm[kafka]
# 带 NATS 传输
pip install lzm-edsm[nats]
# 带 ZeroMQ 传输
pip install lzm-edsm[zmq]
# 组合安装(如 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)。
传输层(Transport)
lzm-edsm 支持多种传输后端,进程内零依赖,跨进程可选 Redis/RabbitMQ/MQTT/Kafka/NATS/ZeroMQ。
InProcess(默认)
零依赖,进程内事件广播:
from lzm.edsm import Engine
engine = Engine() # 默认 InProcessTransport
Redis
跨进程事件广播,使用 Redis Pub/Sub:
from lzm.edsm import Engine
from lzm.edsm.transports.redis import RedisTransport
transport = RedisTransport(
url="redis://localhost:6379",
channel_prefix="edsm",
seen_ttl=5, # 去重 TTL(秒)
)
engine = Engine(transport=transport)
await engine.start()
安装:
pip install lzm-edsm[redis]
RabbitMQ
跨进程事件广播,使用 RabbitMQ Fanout Exchange:
from lzm.edsm import Engine
from lzm.edsm.transports.rabbitmq import RabbitMQTransport
transport = RabbitMQTransport(
url="amqp://guest:guest@localhost:5672/",
exchange_prefix="edsm",
seen_ttl=5,
)
engine = Engine(transport=transport)
await engine.start()
安装:
pip install lzm-edsm[rabbitmq]
MQTT
轻量级 IoT 协议,适合设备事件:
from lzm.edsm import Engine
from lzm.edsm.transports.mqtt import MQTTTransport
transport = MQTTTransport(
url="mqtt://localhost:1883",
topic_prefix="edsm",
qos=1,
seen_ttl=5,
)
engine = Engine(transport=transport)
await engine.start()
安装:
pip install lzm-edsm[mqtt]
Kafka
高吞吐持久化,按域聚合主题:
from lzm.edsm import Engine
from lzm.edsm.transports.kafka import KafkaTransport
transport = KafkaTransport(
bootstrap_servers="localhost:9092",
topic_prefix="edsm",
group_id="edsm-consumer",
seen_ttl=5,
)
engine = Engine(transport=transport)
await engine.start()
安装:
pip install lzm-edsm[kafka]
NATS
云原生轻量级传输:
from lzm.edsm import Engine
from lzm.edsm.transports.nats import NATSTransport
transport = NATSTransport(
url="nats://localhost:4222",
subject_prefix="edsm",
queue_group="workers", # 可选,负载均衡
seen_ttl=5,
)
engine = Engine(transport=transport)
await engine.start()
安装:
pip install lzm-edsm[nats]
ZeroMQ
无中间件,进程间直接通信:
from lzm.edsm import Engine
from lzm.edsm.transports.zmq import ZMQTransport
transport = ZMQTransport(
pub_bind="tcp://*:5555",
sub_connect="tcp://localhost:5555",
topic_prefix="edsm",
seen_ttl=5,
)
engine = Engine(transport=transport)
await engine.start()
安装:
pip install lzm-edsm[zmq]
事件溯源审计(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())
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()
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.trace_id # 链路追踪
# 序列化
event.to_dict() # dict
Event.from_dict(data) # 反序列化
异常
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(可选) |
| 传输:MQTT | asyncio-mqtt | >= 0.16.0(可选) |
| 传输:Kafka | aiokafka | >= 0.10(可选) |
| 传输:NATS | nats-py | >= 2.0(可选) |
| 传输:ZeroMQ | pyzmq | >= 24.0(可选) |
| 存储:SQLite | aiosqlite | 内置 |
| 测试 | pytest | >= 8.0 |
目录结构
lzm-edsm/
├── pyproject.toml
├── src/lzm/edsm/
│ ├── __init__.py # 导出 Domain, Engine, Event
│ ├── core/
│ │ ├── domain.py # Domain 类
│ │ ├── engine.py # Engine 类(Pipeline)
│ │ ├── event.py # Event 数据结构
│ │ └── exceptions.py # 异常定义
│ ├── transports/
│ │ ├── base.py # Transport 抽象基类
│ │ ├── inprocess.py # 进程内(默认)
│ │ ├── redis.py # Redis Pub/Sub
│ │ ├── rabbitmq.py # RabbitMQ
│ │ ├── mqtt.py # MQTT
│ │ ├── kafka.py # Kafka
│ │ ├── nats.py # NATS
│ │ └── zmq.py # ZeroMQ
│ └── 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.1.1.tar.gz.
File metadata
- Download URL: lzm_edsm-0.1.1.tar.gz
- Upload date:
- Size: 34.1 kB
- Tags: Source
- Uploaded using Trusted Publishing? No
- Uploaded via: twine/6.2.0 CPython/3.12.10
File hashes
| Algorithm | Hash digest | |
|---|---|---|
| SHA256 |
da828fd76a9f92e4fb68e60aef943c13415466a3ed58f439612b5c30a3274c0a
|
|
| MD5 |
03634a5f11c98c9123ebe81d5fdb7da2
|
|
| BLAKE2b-256 |
21d5e90ceceef3987311bb740433997c210783b9bda13b26b23a940afe5bac92
|
File details
Details for the file lzm_edsm-0.1.1-py3-none-any.whl.
File metadata
- Download URL: lzm_edsm-0.1.1-py3-none-any.whl
- Upload date:
- Size: 37.2 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 |
05bf0e60045456a96314d05e11345d581fd840c67adba7e2dad8815107c7d754
|
|
| MD5 |
e09f1b379f3102940da0fd57b43ddd1d
|
|
| BLAKE2b-256 |
4cffd8065341cb97b971b21f75ffda0d7d4b0d41cf22e960fdb912a8573d0b3e
|