Metadata-Version: 2.4
Name: lzm-workers
Version: 0.1.0
Summary: 统一的计算节点框架 — 调度器(Scheduler)与执行器(Executor)于一体，支持水平扩展、高可用的任务分发与执行
Author: Lzm
License-Expression: MIT
Project-URL: Homepage, https://github.com/lzm/lzm-workers
Project-URL: Issues, https://github.com/lzm/lzm-workers/issues
Project-URL: Repository, https://github.com/lzm/lzm-workers
Keywords: task-scheduler,task-executor,distributed-computing,worker,lzm,edsm,plugin-system
Classifier: Development Status :: 3 - Alpha
Classifier: Framework :: AsyncIO
Classifier: Intended Audience :: Developers
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: Topic :: System :: Distributed Computing
Classifier: Typing :: Typed
Requires-Python: >=3.11
Description-Content-Type: text/markdown
License-File: LICENSE
Requires-Dist: lzm-edsm>=0.2
Requires-Dist: lzm-plugin>=0.2
Requires-Dist: redis>=5.0
Provides-Extra: mysql
Requires-Dist: asyncmy>=0.2; extra == "mysql"
Provides-Extra: edsm-redis
Requires-Dist: lzm-edsm[redis]; extra == "edsm-redis"
Provides-Extra: edsm-rabbitmq
Requires-Dist: lzm-edsm[rabbitmq]; extra == "edsm-rabbitmq"
Provides-Extra: edsm-all
Requires-Dist: lzm-edsm[all]; extra == "edsm-all"
Provides-Extra: all
Requires-Dist: lzm-workers[mysql]; extra == "all"
Requires-Dist: lzm-workers[edsm-all]; extra == "all"
Dynamic: license-file

# lzm-workers — 统一的计算节点框架

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

**调度与执行彻底分离** — `lzm-workers` 一个包，两种模式：
- **调度器（Scheduler）**：任务管理、分发、定时、重试、统计
- **执行器（Executor）**：任务消费、插件/脚本执行、结果回传

支持水平扩展、高可用部署，通过 `lzm-edsm` 事件驱动通信，实现**零轮询、高吞吐、低延迟**的任务分发。

---

## 安装

```bash
# 核心依赖
pip install lzm-workers

# 带 MySQL 支持（调度器持久化）
pip install lzm-workers[mysql]

# 带 EDSM 传输
pip install lzm-workers[edsm-redis]     # Redis 传输
pip install lzm-workers[edsm-rabbitmq]  # RabbitMQ 传输
pip install lzm-workers[edsm-all]       # 全部传输

# 全量
pip install lzm-workers[all]
```

---

## 快速开始

### 方式一：CLI 启动

```bash
# 启动执行器守护进程
lzm-workers executor --config config.yaml

# 启动调度器
lzm-workers scheduler --config config.yaml
```

### 方式二：嵌入代码

**调度器模式** — 可嵌入任意 Python 服务：

```python
import asyncio
from lzm.workers import Scheduler
from lzm.workers.common.config import Settings

async def main():
    scheduler = Scheduler(Settings(mode="scheduler"))
    await scheduler.start()

    # 创建一个插件任务
    record = await scheduler.create_task(
        task_type="plugin",
        task_name="heading_strict",
        payload={"text": "..."},
        priority=5,
        timeout=300,
    )
    print(f"任务已创建: {record.task_id}")

    # 保持运行
    await asyncio.Event().wait()

asyncio.run(main())
```

**执行器模式** — 独立守护进程：

```python
import asyncio
from lzm.workers import ExecutorRunner
from lzm.workers.common.config import Settings

async def main():
    runner = ExecutorRunner(Settings(mode="executor"))
    await runner.run_forever()

asyncio.run(main())
```

---

## 核心架构

```
┌── lzm-workers 调度器模式 (可嵌入任意服务) ──────────────────────┐
│  TaskManager → Dispatcher → Reaper → EDSM Event → Redis         │
└──────────────────────────────┬───────────────────────────────────┘
                               │
                    EDSM Event (COMPETING)
                    Redis 心跳 (SETEX)
                               │
┌── lzm-workers 执行器模式 (独立守护进程) ────────────────────────┐
│  EDSM Consumer → TaskRunner(租约) → ExecutorPool → 结果回写     │
│  HeartbeatReporter (每 3s) → HealthServer (/health, /metrics)  │
└──────────────────────────────────────────────────────────────────┘
```

### 3 道健壮防线

| 防线 | 机制 | 响应时间 |
|------|------|---------|
| **租约** | `Redis SET NX EX task:lease:{id}` | 任务超时 + 30s 后自动释放 |
| **心跳 + Reaper** | 每 3s 刷新 TTL=10s 心跳，Reaper 每 30s 巡检 | ~30s 回收失联执行器 |
| **调度器无状态** | 状态存 MySQL + Redis，重启自动恢复 | 瞬时恢复 |

### 性能特性

- **零轮询**：纯事件驱动，无任何轮询操作
- **批量回写**：结果攒 50 条或每 3 秒批量 flush
- **Redis O(1)**：心跳、租约、发现均为 Redis O(1) 操作
- **异步并发**：asyncio.Semaphore 控制并发，非阻塞调度

---

## 使用场景

| 场景 | 推荐模式 | 说明 |
|------|---------|------|
| Web 服务需要异步执行任务 | 调度器嵌入 Web 服务 | 请求进来 → 创建任务 → 立即返回 → 异步执行 |
| 独立计算节点集群 | 执行器守护进程 | 多台机器各自运行 Executor，水平扩展 |
| AI 流水线处理 | 调度器 + 执行器 | 调度器编排 DAG，执行器跑 Python 插件 |
| 定时/周期任务 | 调度器 | Cron 表达式触发，调度器自动创建任务 |
| 脚本批量执行 | 执行器 | 通过 ScriptRunner 执行 Shell/Python 脚本 |

---

## 配置参考

```yaml
# config.yaml
mode: executor                  # scheduler / executor

edsm:
  transport_type: redis         # redis / rabbitmq / inprocess
  dsn: "redis://localhost:6379/0"

redis:
  host: localhost
  port: 6379
  db: 1

mysql:                          # 调度器模式需要
  host: localhost
  port: 3306
  user: root
  password: ""
  database: lzm_workers

executor:
  max_concurrency: 10           # 最大并发任务数
  heartbeat_interval: 3         # 心跳间隔（秒）
  health_port: 9100             # 健康检查端口

scheduler:
  reaper_interval: 30           # Reaper 巡检间隔（秒）
```

---

## 项目目录

```
lzm-workers/
├── pyproject.toml
├── src/lzm/workers/
│   ├── __init__.py                  # 导出 Scheduler / ExecutorRunner / Task
│   ├── cli.py                       # CLI 入口
│   ├── common/                      # 公共模块
│   │   ├── config.py                # 统一配置模型
│   │   ├── models.py                # Task / TaskResult / TaskStatus
│   │   └── exceptions.py            # 异常体系
│   ├── executor/                    # 执行器板块
│   │   ├── runner.py                # 主循环
│   │   ├── task_runner.py           # 单任务引擎（租约+超时）
│   │   ├── pool.py                  # 执行器池（路由+并发）
│   │   ├── base.py                  # BaseExecutor 抽象基类
│   │   ├── plugin_executor.py       # PluginExecutor → lzm-plugin
│   │   ├── script_runner.py         # ScriptRunner → subprocess
│   │   ├── heartbeat.py             # 心跳上报
│   │   ├── result_reporter.py       # 批量回写
│   │   └── health_server.py         # 健康检查 HTTP
│   └── scheduler/                   # 调度器板块
│       ├── scheduler.py             # Scheduler 入口
│       ├── task_manager.py          # 任务 CRUD + 定时
│       ├── dispatcher.py            # 执行器发现 + 分发
│       ├── reaper.py                # 收割者
│       └── models.py                # 数据库模型
├── tests/                           # 26 个测试用例
└── docs/
    └── logic-records/               # 架构设计文档
```

---

## 依赖

- [lzm-edsm](https://pypi.org/project/lzm-edsm/) >= 0.2 — 事件驱动状态机引擎
- [lzm-plugin](https://pypi.org/project/lzm-plugin/) >= 0.2 — 通用插件化框架
- `redis` >= 5.0 — 心跳 + 租约 + 执行器发现

---

## 相关项目

| 项目 | 说明 |
|------|------|
| [lzm-edsm](https://pypi.org/project/lzm-edsm/) | 事件驱动状态机引擎 |
| [lzm-plugin](https://pypi.org/project/lzm-plugin/) | 通用插件化框架 |
| [lzm-space](https://pypi.org/project/lzm-space/) | 空间存储编排引擎 |
| lzm-workers | 👈 本包：调度与执行计算框架 |

---

## License

MIT © Lzm
