Metadata-Version: 2.4
Name: lzm-edsm
Version: 0.1.3
Summary: Event-Driven State Machine - 去中心化的事件驱动状态机引擎，pip install 即装即用
Author: Lzm
License: MIT
Project-URL: Homepage, https://github.com/lzm/lzm-edsm
Project-URL: Issues, https://github.com/lzm/lzm-edsm/issues
Project-URL: Repository, https://github.com/lzm/lzm-edsm
Keywords: state-machine,event-driven,event-bus,state,asyncio,edsm,lzm,domain-driven-design
Classifier: Development Status :: 4 - Beta
Classifier: Framework :: AsyncIO
Classifier: Intended Audience :: Developers
Classifier: License :: OSI Approved :: MIT License
Classifier: Operating System :: OS Independent
Classifier: Programming Language :: Python :: 3
Classifier: Programming Language :: Python :: 3.11
Classifier: Programming Language :: Python :: 3.12
Classifier: Topic :: Software Development :: Libraries :: Python Modules
Classifier: Typing :: Typed
Requires-Python: >=3.11
Description-Content-Type: text/markdown
Requires-Dist: aiosqlite>=0.17.0
Provides-Extra: redis
Requires-Dist: redis>=5.0; extra == "redis"
Provides-Extra: rabbitmq
Requires-Dist: aio-pika>=9.0; extra == "rabbitmq"
Provides-Extra: mqtt
Requires-Dist: asyncio-mqtt>=0.16.0; extra == "mqtt"
Provides-Extra: kafka
Requires-Dist: aiokafka>=0.10; extra == "kafka"
Provides-Extra: nats
Requires-Dist: nats-py>=2.0; extra == "nats"
Provides-Extra: zmq
Requires-Dist: pyzmq>=24.0; extra == "zmq"
Provides-Extra: all
Requires-Dist: lzm-edsm[redis]; extra == "all"
Requires-Dist: lzm-edsm[rabbitmq]; extra == "all"
Requires-Dist: lzm-edsm[mqtt]; extra == "all"
Requires-Dist: lzm-edsm[kafka]; extra == "all"
Requires-Dist: lzm-edsm[nats]; extra == "all"
Requires-Dist: lzm-edsm[zmq]; extra == "all"

# lzm-edsm

> 编码：utf8 | 作者：Lzm | 日期：2026-07-15

**Event-Driven State Machine** — 去中心化的事件驱动状态机引擎。

一句话：**状态转移 = 事件发布，事件消费 = 状态响应。**

lzm-edsm 将状态机、事件总线、蓝图合并成一个统一的 Python 包，`pip install` 即装即用，不绑定任何框架或中间件。

---

## 安装

```bash
# 核心（含 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` 是状态、转移、事件的统一容器：

```python
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：

```python
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. 多监听器并发

```python
@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 消费事件 | 任务分发、竞态敏感的状态机转移 |

```python
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（默认）

零依赖，进程内事件广播：

```python
from lzm.edsm import Engine

engine = Engine()  # 默认 InProcessTransport
```

### Redis

跨进程事件传输，支持双模式：

```python
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",
)
```

安装：

```bash
pip install lzm-edsm[redis]
```

### RabbitMQ

跨进程事件传输，支持双模式：

```python
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",
)
```

安装：

```bash
pip install lzm-edsm[rabbitmq]
```

### MQTT

轻量级 IoT 协议，支持双模式：

```python
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",
)
```

安装：

```bash
pip install lzm-edsm[mqtt]
```

### Kafka

高吞吐持久化，Consumer Group 天然支持竞争消费：

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

安装：

```bash
pip install lzm-edsm[kafka]
```

### NATS

云原生轻量级传输，支持双模式：

```python
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,
)
```

安装：

```bash
pip install lzm-edsm[nats]
```

### ZeroMQ

无中间件，进程间直接通信，支持双模式：

```python
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",
)
```

安装：

```bash
pip install lzm-edsm[zmq]
```

---

## 事件溯源审计（SQLiteStore）

`SQLiteStore` 是内置的事件持久化组件，用于审计和溯源（已包含在核心依赖中）。

使用示例：

```python
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

---

## 完整示例

```python
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` 提供域的管理/自省能力，无需额外安装。

### 域概览

```python
# 查看所有已注册域
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 | 监听器数 |

### 域详情

```python
# 查看指定域的完整结构
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"]}
```

### 事件查询

```python
# 所有已注册事件
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": ""}]
```

### 蓝图导出

```python
# 导出完整蓝图（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

```python
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

```python
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

```python
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)    # 反序列化
```

### 异常

```python
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/
```

---

## 相关文档

- [开发进度时间线](docs/timeline.md)
- [常见错误归档](docs/common-errors.md)
- [编码规范](docs/coding-standards.md)
- [TODO 管理](docs/todo.md)
- [逻辑记录](docs/logic-records/)

---

## 许可证

MIT License

---

## 标签

#event-driven #state-machine #asyncio #python #edsm #lzm
