Metadata-Version: 2.4
Name: ei-pipe-sdk
Version: 0.3.0
Summary: Standalone SDK for ei-pipeline executors: ei_infer (model inference) and ei_ois (OIS object-storage transfer) workers, plus shared arq naming, logging, path, profiling and redis (global KV / distributed lock) helpers
Requires-Python: <3.15,>=3.11
Description-Content-Type: text/markdown
Requires-Dist: arq>=0.26
Requires-Dist: redis>=8.0.1
Requires-Dist: structlog>=25.4.0
Requires-Dist: ois3-sdk-python>=3.2.31
Requires-Dist: pyyaml>=6.0.2

# ei-pipe-sdk

面向算法同学的极简模型接入 SDK。把你的模型写成一个 `EiInfer` 的两个方法，ei-pipeline
就能调度它推理：**你不用关心数据怎么进来、结果怎么提交，也不用关心部署**。

## 安装

```bash
pip install ei-pipe-sdk
# 或在 sdk/ 目录下本地构建后安装
uv build && pip install dist/ei_pipe_sdk-*.whl
```

依赖：`arq` + `redis`（队列）、`structlog`（日志）、`ois3-sdk-python`（OIS 传输）、
`pyyaml`（配置文件）。指标上报不在依赖里：SDK 只暴露 backend 接口，由宿主进程注入
自己的指标实现。

## 三步接入

```python
# my_model.py
from ei_infer import EiInfer, register


class MyModel(EiInfer):
    def init(self, model_cfg: dict | None = None) -> None:
        # 进程启动时调用一次：加载权重、tokenizer、选设备。失败会让 worker 启动即崩溃。
        self.model = load_weights(...)

    def inference(self, input_path: str, output_path: str, node_info: dict) -> dict:
        # input_path / output_path: worker 已恢复好的绝对路径
        #   = {data_root}/{service_id}/buckets/{time_bucket}/{pipeline_id}/<相对路径>
        # input_path 是输入（文件或目录）；output_path 是允许写结果文件的目录（可能为空）。
        # 是否落盘由模型自己决定；返回的 dict 会被 worker 回传给 pipeline。
        return {"detections": self.model(open(input_path))}


register("my-model", MyModel())   # 注册名 = 管道节点 params.name 的值
```

```bash
pip install ei-pipe-sdk
```

```python
# worker 入口：参数经构造参数显式给定（也可经配置文件，见下文），
# ei_infer 本身不读环境变量。每个注册模型各消费一条自己的队列。
from ei_infer import InferNode

node = InferNode(
    redis_url="redis://scheduler:6379/0",
    service_id="3fa2b1c9",  # 调度服务 h8；或经配置文件给 scheduler_url 解析
    data_root="/lpai/pvc/ei-autolabel-bd-ga-infer",  # PFS 挂载点裸路径
)
node.run()
```

## 协议

调度端（ei-pipeline 的 `ei_infer` 节点）会按统一协议给你发任务：

| 字段 | 类型 | 说明 |
|---|---|---|
| `name` | str | 模型注册名，由你 `register("name", ...)` 时命名 |
| `input_path` | str | 输入路径（相对 pipe 共享目录，如 `data/`），worker 恢复为绝对路径 |
| `output_path` | str? | 允许写结果文件的目录（相对 pipe 共享目录，可选） |
| `version` | str? | 版本提示，透传保留 |

worker 用 `node_info.service_id + time_bucket + pipeline_id` 与 `InferNode`
配置的 `data_root` 把 `input_path` / `output_path` 恢复成绝对路径
（`{data_root}/{service_id}/buckets/{time_bucket}/{pipeline_id}/<相对路径>`），
调用 `inference(input_path, output_path, node_info[, abort_event])`。
**是否把结果落盘由模型自己决定**，worker 不替模型落盘，只把模型返回的 dict 原样回传。
整个链路里**你的函数只看绝对路径 input_path / output_path 和 node_info**，其余（数据就位、
并发、轮询）全部由 SDK 兜底。

## 协作式取消

Python 线程杀不掉。job 被取消（`job.abort()`）或超时时，worker 会激活该 job
专属的 `abort_event`；`inference` 声明第 4 个参数即可收到它，在推理间隙检查并
提前退出：

```python
from ei_infer import EiInfer, InferenceAborted


class MyModel(EiInfer):
    def inference(self, input_path, output_path, node_info, abort_event=None):
        for item in load_items(input_path):
            if abort_event.is_set():  # 推理间隙检查取消信号
                raise InferenceAborted("job cancelled")
            ...
```

- 事件按 job 生成，并发执行时各 job 各持独立事件，互不影响。
- 老的 3 参签名不受影响（拿不到取消信号，跑到自然结束），完全向后兼容。
- 事件未激活时自发抛 `InferenceAborted` 按普通 job 失败处理。

## 配置

worker 参数一律显式给定（`InferNode` 构造参数，或经配置文件 + `build_node`），
`ei_infer` / `ei_ois` 不读环境变量。例外只有身份与 Redis 共享件（`ei_redis`、
`ei_identity`），见下文「全局 KV 与分布式锁」。

构造参数 / 配置文件字段：

| 字段 | 默认 | 说明 |
|---|---|---|
| `redis_url` | 必填 | 消费队列的 Redis（与调度服务同一套） |
| `service_id` | 二选一 | 调度服务的 service_id（h8），静态给定时优先生效 |
| `scheduler_url` | 二选一 | 未给 `service_id` 时，向该调度请求 `GET /api/status/service-id` 解析（不可达每 10 秒重试） |
| `data_root` | 必填 | PFS 挂载点裸路径；构建时自动并入 `{service_id}/buckets` 隔离段 |
| `max_jobs` | `4` | 每模型队列的并发任务数（arq `max_jobs`） |
| `max_threads` | `max(4, max_jobs)` | 跑同步推理的线程池大小（仅构造参数） |
| `result_max_bytes` | `65536` | job result 进 Redis 前的尺寸上限（字节）；超限 job 失败（`ei_arq.result_guard`） |

配置文件（yaml/json，`{{变量名}}` 占位由 `secrets` 映射在载入时注入）::

    service_id: '{{service_id}}'
    # scheduler_url: http://scheduler:8000
    redis_url: redis://localhost:6379/0
    data_root: /lpai/pvc/ei-autolabel-bd-ga-infer
    max_jobs: 4

```python
from ei_infer.worker import build_node

node = build_node("infer.yaml", secrets={"service_id": "3fa2b1c9"})
node.run()
```

完整示例见 `example/`（`infer_exa.py` + `run.sh`）。示例只需三个变量
（redis、data_root、service_id/scheduler_url），直接构造 `InferNode`，
不走配置文件；配置文件方式用于生产部署。

## 队列

每个注册模型一条队列，统一文法（实现见 `ei_arq/naming.py`）：

```
arq:ei:{service_id}:ei_infer:{model_name}
```

worker 进程为每个注册模型起一个 arq Worker 消费对应队列；`service_id` 与
调度侧表名/Redis 前缀同源，隔离共用一套 Redis 的多个调度服务。

## 多模型

多个模型可在同一个 worker 进程里 `register` 多个名字；进程会为每个模型各起
一个 arq Worker，各自消费自己的队列，互不抢任务。

## 失败模式

- `inference` 返回非 dict → 任务失败（错误信息含模型名）。
- `name` 未注册 → 任务失败，节点错误信息会指出未知模型名。
- `init` 抛异常 → worker 启动即退出，方便在你自己的日志里早发现。
- `inference` 抛 `InferenceAborted`：abort 事件已激活 → 任务取消；未激活 → 普通任务失败。

## 全局 KV 与分布式锁

`ei_redis` 是从调度服务（ei-pipeline）迁过来的两件共享件：跨进程共享的 **global
KV** 与 **分布式锁**。多套调度服务和它们的 worker 复用同一套 Redis，所以本包写的
键**一律自动带 `{service_id}_` 前缀**，调用方只管逻辑名，不要自己拼前缀：

| 用途 | 实际键 |
|---|---|
| global KV | `{service_id}_kv:{key}` |
| 分布式锁 | `{service_id}_lock:{name}` |

要往同一个 Redis 里放自己的键，也用 `ei_redis.redis_key("my_ns", "name")` 拼，别手写
前缀——前缀段的校验（单段、挡掉 `.` / `..` 这类可穿越段）都在那个函数里。

> arq 自己的键（队列、`arq:result:*`、心跳）不走这个前缀：队列文法
> `arq:ei:{service_id}:...` 与调度侧互为镜像，改前缀会直接破坏互操作。

```python
from ei_redis import create_distributed_lock, get_global_kv

kv = get_global_kv()  # 单例：EI_REDIS_URL 与 service_id 都按环境变量解析
kv.set_json("module:state", {"phase": "running"}, ttl_sec=300)  # 不填 TTL 默认 7 天，最长 7 天
state = kv.get_json("module:state")  # 键不存在回 None

lock = create_distributed_lock("leader", ttl_sec=30, auto_renew=True)
if lock.acquire():  # 拿不到就是 False，不抛
    try:
        ...  # 临界区
    finally:
        lock.release()  # 按 owner token 释放，不会误删别人的锁
```

不想让包碰环境变量，就自己建连接、显式给 `service_id`（只有传 `None` 才回落环境变量）：

```python
import redis
from ei_redis import RedisDistributedLock, RedisGlobalKv

client = redis.Redis.from_url("redis://scheduler:6379/0", decode_responses=True)
kv = RedisGlobalKv(client, service_id="3fa2b1c9")
lock = RedisDistributedLock(client, "leader", ttl_sec=30, service_id="3fa2b1c9")
```

失败模式：

- `EI_REDIS_URL` 未配置（非 local 环境）：KV 每次操作抛 `GlobalKvDisabledError`，
  锁在创建时抛 `RedisNotConfiguredError`——不会静默退化成进程内的内存态。
- Redis 报错：KV 一律包成 `GlobalKvUnavailableError` 抛出；锁的 `acquire()` 返回
  `False`（宁可当成没拿到），`is_lock_held()` 返回 `False`（探活失败即放行）。
- `auto_renew=True` 时按 `ttl_sec / 3` 起一个 daemon 线程续期；续期一失败就停手，
  让 TTL 决定锁的生死，`release()` 最多再等它 5 秒退出。
- 指标默认不上报，宿主进程用 `install_kv_metrics_backend(backend)` 注入实现（与
  `ei_ois` 的指标后端同一种做法），SDK 自身不依赖任何指标库。

### 相关环境变量

身份与 Redis 共享件是 SDK 里唯一读环境变量的部分，worker 参数不受影响：

| 变量 | 默认 | 用途 |
|---|---|---|
| `EI_REDIS_URL` | local 环境为 `redis://localhost:6379/0`，其它环境不配即禁用 KV / 锁报错 | global KV 与锁共用的 Redis |
| `EI_REDIS_SOCKET_TIMEOUT_SEC` | `1` | 每次 Redis 读写的 socket 超时，非正数直接报错 |
| `LI_VDC` / `LI_ENV` / `APP_NAME` / `COMP_NAME` | 空 / `local` / 空 / `ei-pipe-sdk` | 平台四元组拼成 version（超 32 字符先整体 md5），再取 `md5[:8]` 得到本进程 `service_id` |
| `EI_VERSION_HASH_OVERRIDE` | 空（用 md5 结果） | 直接指定 h8，只接受字母数字下划线，本地/测试固定命名空间用 |
| `EI_ENV` | 取 `LI_ENV`，都没有则 `local` | 环境名，决定 local 的 `EI_REDIS_URL` 兜底是否生效 |

想知道本进程落在哪个命名空间，直接问 SDK（和上面前缀的那一段是同一个值）：

```python
from ei_identity import service_id

print(service_id())  # '3fa2b1c9' -> 键形如 3fa2b1c9_kv:module:state
```

`COMP_NAME` 默认取 `ei-pipe-sdk` 而不是 `ei-pipeline`：SDK 进程不是调度组件，落到
调度侧默认身份上会和调度服务共用同一个隔离段，正是这层前缀要防的碰撞。

> 本包的命名空间和调度服务侧（`ei:{env}:{h8}:kv:*`）**不互认**：前缀格式不同，传入同一个
> `service_id` 也拼不出对方的键名。worker 与调度之间要传数据走任务参数和队列，不要靠共享
> Redis 键。

## 本地联调

```bash
# 起一个本地 redis 后（redis-server），在 example/ 下：
bash run.sh   # infer_exa.py（前台）+ ois_exa.py（后台）
# 每个注册模型一条队列：arq:ei:{service_id}:ei_infer:{model_name}
```

然后在 ei-pipeline 用 `ei_infer` 节点提交一条任务即可走通。

## 仓库结构

可发布项目整体在 `sdk/` 下（仓库根只放 CI 配置、钩子和说明文档），以下路径相对 `sdk/`：

```
src/ei_arq/      # 队列文法、worker 心跳、job result 尺寸守卫
src/ei_identity/   # service_id：本进程身份（环境变量推导）与调度身份（静态 / 调度查询）
src/ei_infer/    # 模型推理 worker（EiInfer / InferNode / register）
src/ei_log/      # kwargs 风格结构化日志（structlog）
src/ei_ois/      # OIS 对象存储传输 worker（OisNode / client / transfer）
src/ei_path/     # pipe 共享目录路径解析（time bucket / 绝对路径还原）
src/ei_redis/    # 跨进程共享件：global KV 与 Redis 分布式锁（键自动带 service_id 前缀）
tests/           # 全量单测，不依赖真实 redis / OIS / 调度服务
example/         # 可跑的最小 worker 示例
publish.sh       # 构建 + 上传到 PyPI 兼容索引
```

## 开发与发布

命令都在仓库根执行，用 `--project sdk` 指定项目（CI 与 pre-commit 钩子同样如此）：

```bash
uv sync --project sdk                                            # 装运行依赖 + dev 组
uv run --no-sync --project sdk pytest sdk/tests -q               # 全量测试
uv run --no-sync --project sdk ruff check sdk                    # lint
uv run --no-sync --project sdk black --target-version py311 --check sdk
uv run --project sdk pre-commit install                          # 提交时自动跑 ruff + black
```

版本号手工维护：发布前改 `sdk/pyproject.toml` 的 `version`（`publish.sh` 不自动递增），
并跑 `uv lock --project sdk` 同步 `sdk/uv.lock`。

```bash
bash sdk/publish.sh --skip-upload   # 只构建，检查 sdk/dist/ 产物
bash sdk/publish.sh --internal      # 上传内部 Artifactory（索引取自 pyproject 的 [[tool.uv.index]]）
bash sdk/publish.sh                 # 上传公网 PyPI，需要 UV_PUBLISH_TOKEN
```
