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]

# 带 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。

消费模式

每个 Transport 支持两种消费模式(v0.1.3+):

模式 枚举值 说明 适用场景
广播 ConsumeMode.BROADCAST 所有 Worker 全量接收事件(默认) 配置变更、通知推送、状态同步
竞争 ConsumeMode.COMPETING 仅一个 Worker 消费事件 任务分发、竞态敏感的状态机转移
from lzm.edsm.transports.base import ConsumeMode

各传输层的竞争模式实现机制:

传输层 广播实现 竞争实现
Redis Pub/Sub Stream + Consumer Group
RabbitMQ Fanout Exchange Direct Exchange + 共享队列
MQTT # 通配订阅 $share/group/# 共享订阅
Kafka Consumer Group(原生竞争)
NATS 普通订阅 Queue Group
ZeroMQ PUB/SUB PUSH/PULL

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]

MQTT

轻量级 IoT 协议,支持双模式:

from lzm.edsm import Engine
from lzm.edsm.transports.mqtt import MQTTTransport
from lzm.edsm.transports.base import ConsumeMode

# 广播模式(通配订阅,默认)
transport = MQTTTransport(
    url="mqtt://localhost:1883",
    topic_prefix="edsm",
    qos=1,
    seen_ttl=5,
)

# 竞争模式(共享订阅 $share/edsm-group/edsm/#)
transport = MQTTTransport(
    url="mqtt://localhost:1883",
    mode=ConsumeMode.COMPETING,
    topic_prefix="edsm",
    shared_group="edsm-workers",
)

安装:

pip install lzm-edsm[mqtt]

Kafka

高吞吐持久化,Consumer Group 天然支持竞争消费:

from lzm.edsm import Engine
from lzm.edsm.transports.kafka import KafkaTransport

# Kafka 始终使用 Consumer Group(竞争模式)
# 同一 group_id 内的消费者自动负载均衡
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",
    seen_ttl=5,
)

# 竞争模式(Queue Group)
transport = NATSTransport(
    url="nats://localhost:4222",
    subject_prefix="edsm",
    queue_group="workers",  # 同组内仅一个 Worker 消费
    seen_ttl=5,
)

安装:

pip install lzm-edsm[nats]

ZeroMQ

无中间件,进程间直接通信,支持双模式:

from lzm.edsm import Engine
from lzm.edsm.transports.zmq import ZMQTransport
from lzm.edsm.transports.base import ConsumeMode

# 广播模式(PUB/SUB,默认)
transport = ZMQTransport(
    pub_bind="tcp://*:5555",
    sub_connect="tcp://localhost:5555",
    topic_prefix="edsm",
    seen_ttl=5,
)

# 竞争模式(PUSH/PULL,ZMQ 自动轮询分发)
transport = ZMQTransport(
    mode=ConsumeMode.COMPETING,
    push_bind="tcp://*:5556",
    pull_connect="tcp://localhost:5556",
    topic_prefix="edsm",
)

安装:

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())

管理接口(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.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 抽象基类 + ConsumeMode 枚举
│   │   ├── inprocess.py     # 进程内(默认)
│   │   ├── redis.py         # Redis Pub/Sub(广播)+ Stream(竞争)
│   │   ├── rabbitmq.py      # RabbitMQ(双模式)
│   │   ├── mqtt.py          # MQTT(双模式)
│   │   ├── kafka.py         # Kafka(Consumer Group)
│   │   ├── nats.py          # NATS(双模式)
│   │   └── zmq.py           # ZeroMQ PUB/SUB(广播)+ PUSH/PULL(竞争)
│   └── 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.1.3.tar.gz (42.5 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.1.3-py3-none-any.whl (42.5 kB view details)

Uploaded Python 3

File details

Details for the file lzm_edsm-0.1.3.tar.gz.

File metadata

  • Download URL: lzm_edsm-0.1.3.tar.gz
  • Upload date:
  • Size: 42.5 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.1.3.tar.gz
Algorithm Hash digest
SHA256 dcdbbe3302466b16f08df62b38b1a331f16429634914fb74007f27b22d5ac899
MD5 42fa9ef1e5f4442eeace96a5ea15ddc3
BLAKE2b-256 b46a619f837d0e7274a06f8f842961697a0436ca0babd6d3ed4f2cdb740497d8

See more details on using hashes here.

File details

Details for the file lzm_edsm-0.1.3-py3-none-any.whl.

File metadata

  • Download URL: lzm_edsm-0.1.3-py3-none-any.whl
  • Upload date:
  • Size: 42.5 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.1.3-py3-none-any.whl
Algorithm Hash digest
SHA256 84f0fe063dd2fc5ff1f74fa34d5c75ec064e8773e36e1adb8768d3b5915e87cd
MD5 0ea637c0d690906c65f24e62b0a5aa49
BLAKE2b-256 2279afe796a87830931a7659c27642965f477b9528534a51f0b583e7c5d8decf

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