Metadata-Version: 2.5
Name: scheduler-sdk
Version: 0.2.2
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

> 面向内网任务调度系统的 Python Worker SDK。通过 NATS 接收任务，将普通 Python
> 函数注册为远程任务，并向调度服务回传执行结果、心跳和运行状态。

| 项目 | 支持范围 |
|---|---|
| Python | 3.10 及以上 |
| 操作系统 | Windows、Linux |
| 任务类型 | 同步函数、异步函数 |
| 消息中间件 | NATS |
| 分发格式 | 纯 Python 通用 wheel（`py3-none-any`） |
| 配置方式 | 当前运行目录下的 `.env` + 进程环境变量 |

## 目录

- [核心能力](#核心能力)
- [安装](#安装)
- [五分钟快速开始](#五分钟快速开始)
- [注册任务](#注册任务)
- [任务参数注入](#任务参数注入)
- [同步与异步任务](#同步与异步任务)
- [接收任务文件](#接收任务文件)
- [错误处理](#错误处理)
- [组合多个业务模块](#组合多个业务模块)
- [完整配置说明](#完整配置说明)
- [并发、超时与关闭](#并发超时与关闭)
- [任务包更新与回滚](#任务包更新与回滚)
- [日志与运行状态](#日志与运行状态)
- [Windows 与 Linux 部署](#windows-与-linux-部署)
- [内网离线安装](#内网离线安装)
- [公开 API 速查](#公开-api-速查)
- [常见问题](#常见问题)
- [开发与验证](#开发与验证)

## 核心能力

- 使用装饰器或显式注册方式，把普通 Python 函数发布为远程任务。
- 同时支持同步函数和 `async def` 异步函数。
- 自动注入任务编号、任务名称、业务上下文、时间范围和 Worker 标识。
- 将 `kv_arg` 中的业务参数按 Python 关键字参数规则传给任务函数。
- 自动下载任务附件，在任务执行结束后清理临时文件。
- 使用统一并发限制保护 Worker，避免任务无限制堆积执行。
- 定时上报 Heartbeat v2 和本地健康文件。
- 支持按通道获取任务包、安全校验、原子激活、失败回滚和自动重启。
- 支持 Windows 和 Linux 信号处理，可通过 `Worker.stop()` 主动停止。

## 安装

### 使用 pip

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

固定生产版本：

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

### 使用 uv

```bash
uv add scheduler-sdk
```

或者在现有虚拟环境中安装：

```bash
uv pip install scheduler-sdk==0.2.2
```

安装完成后验证：

```bash
python -c "from scheduler_sdk import TaskRegistry, Worker; print('SDK 安装成功')"
```

## 五分钟快速开始

### 1. 创建目录

```text
my-worker/
├── .env
└── worker.py
```

### 2. 创建 `.env`

```dotenv
# NATS 必须填写 Worker 实际能够访问的地址。
NATS_URL=nats://192.0.2.10:4222
NATS_TOKEN=change-me-token

# 每个 Worker 进程必须使用唯一标识。
WORKER_CLIENT_ID=worker-A1
WORKER_CHANNEL=default

# 进程内所有任务共享的并发上限。
WORKER_CONCURRENCY=4
WORKER_HEARTBEAT_INTERVAL=10
```

> Docker 容器中的 `127.0.0.1` 和 `localhost` 指向容器自身。Worker 运行在容器中时，
> 应填写 NATS 所在主机的实际内网 IP、容器服务名或内网 DNS。

### 3. 创建 `worker.py`

```python
import logging

from scheduler_sdk import TaskRegistry, Worker


logging.basicConfig(
    level=logging.INFO,
    format="%(asctime)s %(levelname)s %(name)s %(message)s",
)

tasks = TaskRegistry("examples")


@tasks.task("用户查询")
def query_user(task_context: str, include_disabled: bool = False) -> dict:
    logging.info(
        "查询用户: user_id=%s include_disabled=%s",
        task_context,
        include_disabled,
    )
    return {
        "user_id": task_context,
        "include_disabled": include_disabled,
    }


if __name__ == "__main__":
    Worker(tasks).run()
```

### 4. 启动 Worker

Linux：

```bash
python worker.py
```

Windows PowerShell：

```powershell
python .\worker.py
```

启动成功后，日志会显示 Worker 标识、通道和已注册任务列表。

## 注册任务

### 装饰器注册

适合直接编写新任务：

```python
from scheduler_sdk import TaskRegistry

tasks = TaskRegistry("assets")


@tasks.task("资产查询")
def query_asset(task_context: str) -> dict:
    return {"asset_id": task_context}
```

装饰器不会替换原函数，因此该函数仍然可以在本地直接调用和测试。

### 显式注册已有函数

适合把现有业务函数接入 SDK：

```python
from scheduler_sdk import TaskRegistry


def calculate_total(amount: float, tax_rate: float = 0.06) -> dict:
    return {"total": amount * (1 + tax_rate)}


tasks = TaskRegistry("billing")
tasks.register("计算含税金额", calculate_total)
```

### 注册约束

- 同一个注册表内不能出现重复任务名。
- `Worker.run()` 启动后注册表会被冻结，不能继续注册或组合任务。
- 不支持 positional-only 参数，例如 `def task(value, /)`。
- 不支持 `*args`；可以使用 `**kwargs` 接收动态业务参数。
- 一个任务函数最多声明一个 `TaskExecution` 参数。

## 任务参数注入

服务端任务请求由固定字段和 `kv_arg` 组成。SDK 根据函数参数名自动注入固定字段，
再把 `kv_arg` 展开为关键字参数。

### 固定字段

| 参数名 | 类型 | 含义 |
|---|---|---|
| `task_id` | `str` | 当前任务唯一编号 |
| `task_name` | `str` | 当前任务名称 |
| `task_context` | `str` | 业务上下文，通常是对象编号或查询条件 |
| `time_start` | `str` | 开始时间；请求未提供时为空字符串 |
| `time_end` | `str` | 结束时间；请求未提供时为空字符串 |
| `client_id` | `str` | 实际执行任务的 Worker 标识 |
| `kv_arg` | 只读映射 | 原始业务参数映射，不包含任务文件元数据 |

任务函数只需声明自己需要的固定字段：

```python
@tasks.task("风险查询")
def query_risk(
    task_id: str,
    task_context: str,
    time_start: str,
    time_end: str,
    client_id: str,
) -> dict:
    return {
        "task_id": task_id,
        "area_id": task_context,
        "range": [time_start, time_end],
        "executed_by": client_id,
    }
```

### 展开 `kv_arg`

假设服务端发送：

```json
{
  "task_name": "报表生成",
  "task_context": "customer-001",
  "kv_arg": {
    "format": "xlsx",
    "overwrite": true
  }
}
```

任务函数可以直接声明同名参数：

```python
@tasks.task("报表生成")
def generate_report(
    task_context: str,
    format: str,
    overwrite: bool = False,
) -> dict:
    return {
        "customer": task_context,
        "format": format,
        "overwrite": overwrite,
    }
```

也可以使用 `**arguments` 接收不固定的业务参数：

```python
from typing import Any


@tasks.task("动态查询")
def dynamic_query(task_context: str, **arguments: Any) -> dict:
    return {"context": task_context, "filters": arguments}
```

SDK 不执行隐式类型转换。服务端传入字符串时，任务函数收到的仍是字符串；业务代码应
自行校验和转换类型。若缺少必填参数、出现未知参数或业务参数与固定字段冲突，SDK 会
返回 `task_input_invalid`。

以下名称为保留参数名，不能出现在 `kv_arg` 中：

```text
task_id, task_name, task_context, time_start, time_end, client_id, kv_arg
```

## 同步与异步任务

### 同步任务

普通 `def` 函数在线程中执行，适合已有阻塞代码、数据库驱动或同步 HTTP 客户端：

```python
import time


@tasks.task("同步计算")
def calculate(value: int) -> dict:
    time.sleep(1)
    return {"result": value * 2}
```

### 异步任务

`async def` 函数直接在 Worker 事件循环中执行：

```python
import asyncio


@tasks.task("异步计算")
async def calculate_async(value: int) -> dict:
    await asyncio.sleep(1)
    return {"result": value * 2}
```

同步和异步任务共享 `WORKER_CONCURRENCY` 并发上限。异步任务中不要直接执行长时间
阻塞操作；应改用异步库，或通过 `asyncio.to_thread()` 转移阻塞调用。

任务返回值需要能够被 JSON 序列化。推荐返回 `dict`、`list`、字符串、数字、布尔值或
`None`，不要直接返回数据库连接、文件句柄、自定义类实例等对象。

## 接收任务文件

服务端通过 multipart 接收附件后，会把文件元数据放入任务消息。Worker 配置
`SCHEDULER_HTTP_URL` 后，SDK 会在调用任务函数前：

1. 从调度服务下载文件。
2. 校验声明大小和实际下载大小。
3. 检查单文件大小上限。
4. 创建仅供当前任务使用的临时目录。
5. 通过 `TaskExecution.files` 提供本地路径。
6. 在任务成功或失败后清理临时目录。

配置示例：

```dotenv
# 必须是完整 HTTP 根地址；有 Nginx 子路径时必须包含 APP_ROOT_PATH。
SCHEDULER_HTTP_URL=http://192.0.2.10:17002/example/scheduler
WORKER_FILE_HTTP_TIMEOUT=30
WORKER_FILE_MAX_BYTES=104857600
```

任务示例：

```python
from scheduler_sdk import TaskExecution


@tasks.task("导入报表")
def import_report(execution: TaskExecution, overwrite: bool = False) -> dict:
    if not execution.files:
        return {"imported": False, "reason": "未上传文件"}

    source = execution.files[0]
    content = source.path.read_text(encoding="utf-8")
    return {
        "task_id": execution.task_id,
        "original_name": source.name,
        "size": source.size,
        "characters": len(content),
        "overwrite": overwrite,
    }
```

`TaskExecution` 和 `TaskFile` 均为只读数据对象：

```python
@dataclass(frozen=True)
class TaskExecution:
    task_id: str
    task_name: str
    files: tuple[TaskFile, ...]

@dataclass(frozen=True)
class TaskFile:
    name: str
    size: int
    path: Path
```

> `TaskFile.path` 只在任务函数执行期间有效。不要把这个临时路径保存到数据库后供其他
> 进程读取；需要长期保留时，应在任务函数返回前复制到业务存储目录。

常见文件错误码：

| 错误码 | 含义 |
|---|---|
| `task_file_unavailable` | 任务包含文件，但 Worker 未配置 `SCHEDULER_HTTP_URL` |
| `task_file_download_failed` | HTTP 下载失败或超时 |
| `task_file_too_large` | 文件超过 Worker 配置的大小上限 |
| `task_file_size_mismatch` | 实际下载大小与服务端元数据不一致 |

## 错误处理

### 返回可识别的业务错误

任务函数可以抛出 `TaskError`：

```python
from scheduler_sdk import TaskError


@tasks.task("校验报表")
def validate_report(row_count: int) -> dict:
    if row_count <= 0:
        raise TaskError(
            "报表没有有效数据",
            code="report_empty",
            details={"row_count": row_count},
        )
    return {"valid": True}
```

调用方收到的错误包含：

```json
{
  "code": "report_empty",
  "message": "报表没有有效数据",
  "details": {
    "row_count": 0
  }
}
```

### 未预期异常

其他异常会统一转换为 `task_execution_failed`。响应中包含异常类型和消息，完整 traceback
只写入 Worker 日志，不通过 NATS 返回，以避免泄露运行环境和内部代码细节。

建议：

- 对调用方需要判断的业务失败使用 `TaskError`。
- 对网络、数据库等异常保留原始异常，让 Worker 记录完整 traceback。
- 不要在错误消息或 `details` 中返回密码、令牌、连接串等敏感信息。

## 组合多个业务模块

大型 Worker 可以按业务域拆分注册表，再由入口统一组合。

目录结构：

```text
my-worker/
├── .env
├── worker.py
└── task_modules/
    ├── __init__.py
    ├── assets.py
    └── metering.py
```

`task_modules/assets.py`：

```python
from scheduler_sdk import TaskRegistry

assets_tasks = TaskRegistry("assets")


@assets_tasks.task("资产查询")
def query_asset(task_context: str) -> dict:
    return {"asset_id": task_context}
```

`task_modules/metering.py`：

```python
from scheduler_sdk import TaskRegistry

metering_tasks = TaskRegistry("metering")


@metering_tasks.task("电量查询")
def query_energy(task_context: str, month: str) -> dict:
    return {"meter_id": task_context, "month": month}
```

`worker.py`：

```python
import logging

from scheduler_sdk import TaskRegistry, Worker
from task_modules.assets import assets_tasks
from task_modules.metering import metering_tasks


logging.basicConfig(level=logging.INFO)

tasks = TaskRegistry("worker-entry")
tasks.include(assets_tasks, metering_tasks)


if __name__ == "__main__":
    Worker(tasks).run()
```

如果不同注册表中存在同名任务，`include()` 会立即抛出 `ValueError`，避免运行时出现
不确定的任务路由。

## 完整配置说明

SDK 先读取当前运行目录下的 `.env`，再使用进程环境变量覆盖同名配置。

| 配置项 | 必填 | 默认值 | 说明 |
|---|:---:|---|---|
| `NATS_URL` | 是 | 无 | Worker 可访问的 NATS 地址 |
| `NATS_TOKEN` | 否 | 空 | NATS Token；服务未启用认证时留空 |
| `WORKER_CLIENT_ID` | 是 | 无 | Worker 唯一标识 |
| `WORKER_CHANNEL` | 否 | `default` | 任务包发布通道 |
| `WORKER_CONCURRENCY` | 否 | `4` | 同步和异步任务共享的并发上限 |
| `WORKER_HEARTBEAT_INTERVAL` | 否 | `10` | NATS 心跳间隔，单位秒 |
| `WORKER_HEALTH_FILE` | 否 | 系统临时目录 | 本地健康文件路径 |
| `WORKER_HEALTH_INTERVAL` | 否 | `5` | 健康文件刷新间隔，单位秒 |
| `SCHEDULER_HTTP_URL` | 否 | 空 | 调度服务完整 HTTP 根地址；也用于文件下载和任务包更新 |
| `WORKER_FILE_HTTP_TIMEOUT` | 否 | `30` | 单个任务文件下载超时，单位秒 |
| `WORKER_FILE_MAX_BYTES` | 否 | `104857600` | 单个任务文件最大字节数，默认 100MB |
| `WORKER_BUNDLE_ROOT` | 否 | `.scheduler-worker` | 任务包、状态和回滚版本目录 |
| `WORKER_BUNDLE_UPDATE_INTERVAL` | 否 | `60` | 检查任务包更新的间隔，单位秒 |
| `WORKER_BUNDLE_HTTP_TIMEOUT` | 否 | `30` | 任务包查询和下载超时，单位秒 |
| `WORKER_BUNDLE_MAX_BYTES` | 否 | `20971520` | 任务包 ZIP 最大字节数，默认 20MB |

所有间隔、并发数和大小限制都必须大于 `0`。配置缺失或格式错误时，Worker 会在启动阶段
直接报错，不会带着无效配置继续运行。

推荐的 Linux 健康文件路径：

```dotenv
WORKER_HEALTH_FILE=/tmp/scheduler-worker-health.json
```

推荐的 Windows 相对路径：

```dotenv
WORKER_HEALTH_FILE=.scheduler-worker\worker-health.json
```

经过多层 Nginx 转发时，`SCHEDULER_HTTP_URL` 必须包含完整子路径，例如：

```dotenv
SCHEDULER_HTTP_URL=http://192.0.2.10:17002/company/platform/scheduler
```

## 并发、超时与关闭

### 全局并发限制

同一 Worker 进程中的所有任务共享一个并发信号量。达到 `WORKER_CONCURRENCY` 后，新任务
会等待并发槽位。健康文件会把满负载状态写为 `busy`，不会把正常满负载误判为卡死。

### 任务超时

SDK 不会强制中断任务，也不会读取 `kv_arg.timeout`。HTTP 同步调用超时只表示调用方停止
等待，不代表 Worker 中的任务已经取消；任务最终完成后仍会回传结果。

需要业务超时时，应在任务函数使用对应库的超时参数，例如 HTTP 请求超时、数据库语句
超时，或在异步任务中使用 `asyncio.timeout()`。

### 优雅关闭

Worker 支持以下停止方式：

- Linux：`SIGINT`、`SIGTERM`。
- Windows：控制台中断以及 Python 支持的终止信号。
- 代码调用：`worker.stop()`，可以从其他线程安全调用。

```python
import threading

from scheduler_sdk import TaskRegistry, Worker

tasks = TaskRegistry("scheduled-stop")
worker = Worker(tasks)

# 示例：由其他线程发出停止请求。
threading.Timer(60, worker.stop).start()
worker.run()
```

收到停止请求后，Worker 停止订阅新任务，等待已进入执行集合的任务完成，写入
`shutting_down` 健康状态，然后排空 NATS 连接。

## 任务包更新与回滚

任务包用于在不重新制作 Worker 安装包的情况下更新任务代码。启用条件是配置
`SCHEDULER_HTTP_URL`。

### 固定目录结构

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

`bundle.json`：

```json
{
  "version": "2026.07.29.1"
}
```

`task_bundle/__init__.py` 必须导出名为 `tasks` 的 `TaskRegistry`：

```python
from scheduler_sdk import TaskRegistry

tasks = TaskRegistry("remote-bundle")


@tasks.task("远程报表")
def build_report(task_context: str) -> dict:
    return {"report_id": task_context}
```

管理员通过受信网络发布任务包，例如：

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

路径末尾的 `stable` 是目标通道，需要与 Worker 的 `WORKER_CHANNEL` 对应。

### 更新流程

1. Worker 按 `WORKER_CHANNEL` 查询目标任务包。
2. 下载 ZIP，并检查压缩包大小和 SHA-256。
3. 拒绝绝对路径、`..`、反斜杠路径和符号链接，避免目录穿越。
4. 校验 `bundle.json` 和 Python 入口。
5. 启动阶段立即加载可用目标版本。
6. 运行中发现新版本时，先下载并校验，再停止接收任务。
7. 等待在途任务结束后重新执行当前 Worker 脚本。
8. 新版本无法加载时尝试回滚上一版本。

本地只保留当前版本和上一版本。任务包状态会通过 Heartbeat v2 上报。

> 任务包不会安装第三方依赖。任务包需要的库必须提前安装在 Worker 环境中。

## 日志与运行状态

SDK 使用 Python 标准库 `logging`，不会强制应用日志格式。建议入口统一配置：

```python
import logging

logging.basicConfig(
    level=logging.INFO,
    format="%(asctime)s [%(process)d] %(levelname)s %(name)s %(message)s",
)
```

主要日志包括：

- Worker 连接成功、停止和重启。
- 注册表冻结及任务列表。
- 任务开始、完成和失败。
- 心跳发送失败与下周期重试。
- 文件下载完成和临时目录清理。
- 任务包下载、激活、回滚和清理失败。
- 健康文件写入失败。

健康文件采用原子替换写入，主要字段包括：

```json
{
  "version": 1,
  "pid": 12345,
  "timestamp": "2026-07-29T08:00:00Z",
  "status": "healthy",
  "uptime_seconds": 120.5,
  "nats_connected": true,
  "workers": {
    "active_threads": 1,
    "max_threads": 4,
    "available_slots": 3
  }
}
```

常见状态：

| 状态 | 含义 |
|---|---|
| `healthy` | NATS 已连接且仍有可用并发槽 |
| `busy` | NATS 已连接，但并发槽已满 |
| `degraded` | NATS 当前未连接或正在重连 |
| `shutting_down` | Worker 正在停止 |

## Windows 与 Linux 部署

SDK 是纯 Python 包，不需要分别下载 Windows wheel 和 Linux wheel。文件名中的
`py3-none-any` 表示该 wheel 不绑定特定 Python ABI 和操作系统平台。

### Linux

```bash
python3 -m venv .venv
source .venv/bin/activate
python -m pip install scheduler-sdk==0.2.2
python worker.py
```

后台运行时建议交给 systemd、Supervisor 或容器编排系统管理，不建议仅使用 `nohup`。

### Windows PowerShell

```powershell
py -3.12 -m venv .venv
.\.venv\Scripts\Activate.ps1
python -m pip install scheduler-sdk==0.2.2
python .\worker.py
```

Windows CMD：

```batch
py -3.12 -m venv .venv
.venv\Scripts\activate.bat
python -m pip install scheduler-sdk==0.2.2
python worker.py
```

## 内网离线安装

内网机器无法访问 PyPI 时，需要在可联网机器上提前下载 SDK 及全部依赖。

### 方式一：目标平台分别准备 wheels

#### Windows x64 / Python 3.12

```bash
python -m pip download \
  --only-binary=:all: \
  --platform win_amd64 \
  --python-version 312 \
  --implementation cp \
  --dest wheels-windows \
  scheduler-sdk==0.2.2
```

#### Linux x64 / Python 3.12

```bash
python -m pip download \
  --only-binary=:all: \
  --platform manylinux2014_x86_64 \
  --python-version 312 \
  --implementation cp \
  --dest wheels-linux \
  scheduler-sdk==0.2.2
```

当前 SDK 及直接依赖均提供纯 Python wheel，但仍建议按目标平台分别保存目录，便于以后
增加包含原生扩展的业务依赖。

把对应 wheels 目录和 Worker 代码复制到内网后安装：

```bash
python -m pip install \
  --no-index \
  --find-links wheels-linux \
  scheduler-sdk==0.2.2
```

Windows：

```batch
python -m pip install --no-index --find-links wheels-windows scheduler-sdk==0.2.2
```

### 方式二：使用 uv 离线安装

先创建虚拟环境，再从本地 wheel 仓库安装：

```bash
uv venv --python 3.12
uv pip install --offline --no-index --find-links wheels-linux scheduler-sdk==0.2.2
```

Windows：

```batch
uv venv --python 3.12
uv pip install --offline --no-index --find-links wheels-windows scheduler-sdk==0.2.2
```

验证整个离线目录是否完整：

```bash
python -c "from scheduler_sdk import TaskExecution, TaskRegistry, Worker; print('离线安装成功')"
```

## 公开 API 速查

### `TaskRegistry(name)`

任务注册表。建议一个业务域使用一个注册表。

| 方法 | 说明 |
|---|---|
| `task(name)` | 装饰器形式注册任务 |
| `register(name, function)` | 注册已有函数，并返回原函数 |
| `include(*registries)` | 合并其他注册表 |

### `Worker(tasks)`

阻塞运行 Worker。

| 方法 | 说明 |
|---|---|
| `run()` | 在当前线程运行，直到收到停止请求 |
| `stop()` | 线程安全地请求停止 |

同一个 `Worker` 实例不能同时重复调用 `run()`。

### `TaskExecution`

当前任务执行信息：

| 属性 | 类型 | 说明 |
|---|---|---|
| `task_id` | `str` | 当前任务编号 |
| `task_name` | `str` | 当前任务名称 |
| `files` | `tuple[TaskFile, ...]` | 当前任务的临时文件 |

### `TaskFile`

| 属性 | 类型 | 说明 |
|---|---|---|
| `name` | `str` | 上传时的原始文件名 |
| `size` | `int` | 实际文件字节数 |
| `path` | `pathlib.Path` | 任务执行期间有效的本地临时路径 |

### `TaskError(message, *, code, details=None)`

向调用方返回结构化业务错误。

```python
raise TaskError(
    "数据不存在",
    code="data_not_found",
    details={"record_id": "A001"},
)
```

## 常见问题

### 启动时报“缺少 Worker 环境配置”

检查运行命令的当前目录是否存在 `.env`，并确认至少配置：

```dotenv
NATS_URL=nats://实际地址:4222
WORKER_CLIENT_ID=唯一标识
```

SDK 读取的是当前工作目录，而不是 `worker.py` 文件所在目录。建议启动前显式切换目录。

### Docker 中无法连接 NATS 或调度服务

不要使用 `127.0.0.1` 或 `localhost`。填写 Docker Compose 服务名、宿主机实际 IP 或
容器能够解析的内网 DNS。

### 任务始终没有被调用

依次检查：

1. Worker 日志中是否出现“Worker 已连接”。
2. 启动日志中的任务名是否与服务端提交的 `task_name` 完全一致。
3. `NATS_TOKEN` 是否与 NATS 服务端一致。
4. 多个 Worker 是否使用了预期的通道和唯一 `WORKER_CLIENT_ID`。

### 任务返回 `task_input_invalid`

常见原因：

- `kv_arg` 缺少任务函数必填参数。
- `kv_arg` 包含函数不接受的参数。
- `kv_arg` 使用了固定字段的保留名称。
- 函数声明了 positional-only 参数或 `*args`。

### 文件任务返回 `task_file_unavailable`

任务带有附件，但没有配置 `SCHEDULER_HTTP_URL`。该地址必须是 Worker 能够访问的完整
HTTP 根地址，经过 Nginx 子路径部署时还必须包含 `APP_ROOT_PATH`。

### Worker 满负载是否会被认为异常

不会。并发槽已满时健康状态为 `busy`，Heartbeat 仍会持续上报。

### 修改任务代码后为什么没有生效

本地任务代码需要重启 Worker。服务端任务包更新会由 SDK 自动检查，并在在途任务完成后
重新执行当前 Worker 脚本。

### 能否从任务函数中强制取消另一个任务

SDK 不提供任务强制中断接口。同步线程无法安全强杀；业务代码应使用超时、幂等设计和
可中断的外部调用。

## 开发与验证

在源码目录执行：

```bash
cd packages/scheduler-sdk
uv sync
uv run pytest -q
uv run mypy src
```

构建 wheel 和源码包：

```bash
uv build --wheel --sdist
```

构建产物位于 `dist/`。发布前建议执行：

```bash
uvx --from twine twine check dist/*
```

## 许可证

本项目采用 MIT License，详见 `LICENSE`。
