Metadata-Version: 2.4
Name: scheduler-sdk
Version: 0.2.0
Summary: 内网任务调度系统的 Worker SDK
License-Expression: MIT
License-File: LICENSE
Keywords: nats,rpc,scheduler,worker
Classifier: Development Status :: 4 - Beta
Classifier: Intended Audience :: Developers
Classifier: Operating System :: Microsoft :: Windows
Classifier: Operating System :: POSIX :: Linux
Classifier: Programming Language :: Python :: 3 :: Only
Classifier: Programming Language :: Python :: 3.10
Classifier: Programming Language :: Python :: 3.11
Classifier: Programming Language :: Python :: 3.12
Classifier: Programming Language :: Python :: 3.13
Classifier: Typing :: Typed
Requires-Python: >=3.10
Requires-Dist: nats-py>=2.9
Requires-Dist: python-dotenv>=1.0
Description-Content-Type: text/markdown

# Scheduler Worker SDK

当前实现包含任务声明与 Worker Runtime。一个脚本可以组合并运行多个任务函数，
基础设施配置全部来自当前目录的 `.env`。

## 安装

发布包是纯 Python 通用 wheel，同时支持 Windows 和 Linux：

```bash
python -m pip install scheduler-sdk==0.2.0
```

需要 Python 3.10 或更高版本。

```python
import logging

from scheduler_sdk import TaskExecution, TaskRegistry, Worker

logging.basicConfig(level=logging.INFO)
tasks = TaskRegistry("assets")


@tasks.task("report.import")
async def import_report(
    task_context: str,
    execution: TaskExecution,
    overwrite: bool = False,
) -> dict:
    source = execution.files[0].path
    return {"source": str(source), "overwrite": overwrite}


Worker(tasks).run()
```

也可以使用 `tasks.register(name, function)` 接入已有函数，使用
`tasks.include(other_registry)` 显式组合多个业务模块。固定任务字段按同名参数注入，
`kv_arg` 按 Python 关键字参数规则展开；SDK 不执行隐式类型转换。

## 运行配置

在 Worker 运行目录创建 `.env`；源码包内提供 `.env.example` 模板，支持以下参数：

- `NATS_URL`：NATS 内网地址，必填。
- `NATS_TOKEN`：NATS 令牌，可留空。
- `WORKER_CLIENT_ID`：Worker 唯一标识，必填。
- `WORKER_CHANNEL`：Worker 通道，默认 `default`。
- `WORKER_CONCURRENCY`：进程内全局并发数，默认 `4`。
- `WORKER_HEARTBEAT_INTERVAL`：心跳间隔秒数，默认 `10`。
- `SCHEDULER_HTTP_URL`：任务包服务完整 HTTP 根地址；不配置则关闭自更新。
- `WORKER_FILE_HTTP_TIMEOUT`：任务文件下载超时秒数，默认 `30`。
- `WORKER_FILE_MAX_BYTES`：单个任务文件下载上限，默认 `100MB`。
- `WORKER_HEALTH_FILE`：Daemon 读取的健康文件路径，默认使用系统临时目录。
- `WORKER_HEALTH_INTERVAL`：健康文件刷新间隔秒数，默认 `5`。
- `WORKER_BUNDLE_ROOT`：本地任务包状态目录，默认 `.scheduler-worker`。
- `WORKER_BUNDLE_UPDATE_INTERVAL`：运行中检查间隔秒数，默认 `60`。
- `WORKER_BUNDLE_HTTP_TIMEOUT`：HTTP 超时秒数，默认 `30`。
- `WORKER_BUNDLE_MAX_BYTES`：任务包 ZIP 最大字节数，默认 `20MB`。

进程启动后，任务注册表会在连接 NATS 前冻结。同步任务在线程中执行，异步任务直接
运行在事件循环中，两者共享同一个全局并发限制。SDK 不设置任务执行超时，也不读取
`kv_arg.timeout`；同步函数真正结束前始终占用并发槽。

## 任务文件

服务端通过 multipart 接收文件后，会在任务消息中注入文件元数据。配置
`SCHEDULER_HTTP_URL` 后，SDK 在调用任务函数前下载并核对文件大小，将本地临时路径
放入 `TaskExecution.files`。`files` 是 SDK 保留输入，不会作为普通 `kv_arg` 参数展开。

```python
from scheduler_sdk import TaskExecution

@tasks.task("report.import")
def import_report(execution: TaskExecution) -> dict:
    task_file = execution.files[0]
    return {
        "name": task_file.name,
        "content": task_file.path.read_text(encoding="utf-8"),
    }
```

本地文件仅在任务函数执行期间有效，函数返回或抛出异常后 SDK 会清理任务临时目录。
缺少 HTTP 地址、下载失败、超过大小上限或下载大小与元数据不一致时，SDK 返回
`task_file_*` 结构化错误，并且不会调用任务函数。

## 错误返回

任务函数可以抛出 `TaskError` 返回可判断的业务错误：

```python
from scheduler_sdk import TaskError

raise TaskError(
    "报表格式不支持",
    code="report_invalid",
    details={"line": 3},
)
```

`TaskError` 会返回 `code/message/details`。其他异常会返回异常类型与消息，完整 traceback
只写入 Worker 日志，不通过 NATS 返回。

## 任务包

任务包使用固定结构，不携带或安装 Python 依赖：

```text
tasks.zip
├── bundle.json
└── task_bundle/
    ├── __init__.py
    └── ...
```

`bundle.json` 只要求非空版本号：

```json
{"version": "2026.07.27.1"}
```

`task_bundle/__init__.py` 必须导出名为 `tasks` 的 `TaskRegistry`。包内可以继续使用
相对导入拆分多个任务模块。管理员通过受信网络发布 ZIP 并同时更新通道目标：

```bash
curl -F "bundle=@tasks.zip" \
  http://192.0.2.10:17002/example/scheduler/admin/task-bundles/stable
```

Worker 启动前会下载目标包、核对 SHA-256、安全解压并验证 Python 入口，然后原子更新
激活指针。运行中发现新版本时，后台只完成下载与结构校验，不执行新任务函数代码；随后
停止接收新任务、等待在途任务结束并重新执行当前脚本。当前版本无法加载时自动回滚上一
版本，本地只保留当前与上一版本。任务包状态与错误原因会通过 Heartbeat v2 上报。
