Metadata-Version: 2.4
Name: ei-pipe-sdk
Version: 0.2.7
Summary: Standalone SDK for ei-pipeline executors: ei_infer (model inference) and ei_ois (OIS object-storage transfer) workers
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
Provides-Extra: profile
Requires-Dist: fastapi>=0.115.0; extra == "profile"
Requires-Dist: memray>=1.14; extra == "profile"
Requires-Dist: py-spy>=0.4.0; extra == "profile"

# ei-pipe-sdk

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

## 安装

```bash
pip install ei-pipe-sdk
# 或先在仓库 ei_pipe_sdk/ 下本地构建
uv build && pip install dist/ei_pipe_sdk-*.whl
```

## 三步接入

```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 已恢复好的绝对路径
        #   = {EI_INFER_DATA_ROOT}/{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
# 在算法部署里跑（pod 内需挂载 PFS data_root 到 EI_INFER_DATA_ROOT，
# 并给出调度服务的 service_id：静态给定或经 scheduler_url 解析）
EI_INFER_DATA_ROOT=/lpai/pvc/ei-autolabel-bd-ga-infer \
EI_INFER_SCHEDULER_URL=http://ei-pipeline:8000 python -m ei_infer.worker
```

## 协议

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

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

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

## 配置（环境变量）

| 变量 | 默认 | 说明 |
|---|---|---|
| `EI_INFER_REDIS_URL` | 回退 `EI_ARQ_REDIS_URL` → `EI_REDIS_URL` → `redis://localhost:6379/0` | 消费队列的 Redis |
| `EI_INFER_SERVICE_ID` | 空 | 调度服务的 service_id（h8），静态给定时优先生效 |
| `EI_INFER_SCHEDULER_URL` | 空 | 未给 `EI_INFER_SERVICE_ID` 时，向该调度请求 `GET /api/status/service-id` 解析（不可达每分钟重试）。两者必须二选一 |
| `EI_INFER_MAX_JOBS` | `4` | 每 worker 并发任务数（arq `max_jobs`） |
| `EI_INFER_MAX_THREADS` | `max(4, max_jobs)` | 跑同步推理的线程池大小 |
| `EI_INFER_DATA_ROOT` | `/lpai/pvc/ei-autolabel-bd-ga-infer` | PFS data_root 在 pod 内的挂载根 |

## 队列

每个注册模型一条队列，统一文法（详见 ei-pipeline `docs/arq_queue_naming.md`）：

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

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

## 多模型

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

## 失败模式

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

## 本地联调

```bash
# 起一个本地 redis 后（redis-server）
EI_INFER_SERVICE_ID=f7436895 python -m ei_infer.worker
# 每个注册模型一条队列：arq:ei:f7436895:ei_infer:{model_name}
```

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