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。

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

lzm_edsm-0.1.1.tar.gz (34.1 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.1-py3-none-any.whl (37.2 kB view details)

Uploaded Python 3

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

Hashes for lzm_edsm-0.1.1.tar.gz
Algorithm Hash digest
SHA256 da828fd76a9f92e4fb68e60aef943c13415466a3ed58f439612b5c30a3274c0a
MD5 03634a5f11c98c9123ebe81d5fdb7da2
BLAKE2b-256 21d5e90ceceef3987311bb740433997c210783b9bda13b26b23a940afe5bac92

See more details on using hashes here.

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

Hashes for lzm_edsm-0.1.1-py3-none-any.whl
Algorithm Hash digest
SHA256 05bf0e60045456a96314d05e11345d581fd840c67adba7e2dad8815107c7d754
MD5 e09f1b379f3102940da0fd57b43ddd1d
BLAKE2b-256 4cffd8065341cb97b971b21f75ffda0d7d4b0d41cf22e960fdb912a8573d0b3e

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