Skip to main content

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

lzm_edsm-0.2.0.tar.gz (56.3 kB view details)

Uploaded Source

Built Distribution

If you're not sure about the file name format, learn more about wheel file names.

lzm_edsm-0.2.0-py3-none-any.whl (39.1 kB view details)

Uploaded Python 3

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

Hashes for lzm_edsm-0.2.0.tar.gz
Algorithm Hash digest
SHA256 2de497506e16c0dba7c010c1d80c5ae01c78415bc29be6aa859db38838f5c918
MD5 3598e37c8ef4c11ffe090427a623948e
BLAKE2b-256 e70eae3934ae774334989bc872e168182a8dd353f7fb34fd24fccefdb750e6aa

See more details on using hashes here.

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

Hashes for lzm_edsm-0.2.0-py3-none-any.whl
Algorithm Hash digest
SHA256 2ee847ce772317aae1b91a5bbbb26aa5da4d4986fea16b9a11a5a4694211dcc7
MD5 e3f9f246da74753f81fe60d553af7062
BLAKE2b-256 99987a2761dd4b072f29a8f13acef26207237b79d2eea49850a2731c4982c1e5

See more details on using hashes here.

Supported by

AWS Cloud computing and Security Sponsor Datadog Monitoring Depot Continuous Integration Fastly CDN Google Download Analytics Pingdom Monitoring Sentry Error logging StatusPage Status page