Metadata-Version: 2.4
Name: ei-pipe-sdk
Version: 0.4.1
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, redis (global KV / distributed lock) and dba (generic DDL registry / schema verify / dba-mode fallback) 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
Provides-Extra: dba
Requires-Dist: sqlalchemy>=2.0; extra == "dba"
Requires-Dist: pymysql>=1.1; extra == "dba"
Requires-Dist: fastapi>=0.110; extra == "dba"

# ei-pipe-sdk

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

## 安装

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

# 需要数据库运维件（ei_dba）时装可选 extra：
pip install 'ei-pipe-sdk[dba]'   # 额外拉 SQLAlchemy / pymysql / FastAPI
```

依赖：`arq` + `redis`（队列）、`structlog`（日志）、`ois3-sdk-python`（OIS 传输）、
`pyyaml`（配置文件）。指标上报不在依赖里：SDK 只暴露 backend 接口，由宿主进程注入
自己的指标实现。`ei_dba`（数据库运维）的 SQLAlchemy / pymysql / FastAPI 走 `[dba]`
可选 extra，核心 SDK 不依赖它们——不装也能用其余所有包。

## 三步接入

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


@EI_REGISTER("my-model")            # 注册名 = 管道节点 params.name 的值
class MyModel(EiInfer):
    def __init__(self, name):       # EI_REGISTER 用注册名作为唯一位置参实例化
        self.name = name

    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))}
```

```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 | 模型注册名，由你 `@EI_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`）。
**整个 SDK 不读任何环境变量**：身份（`ei_identity.IdentityConfig`）、Redis 共享件
（`ei_redis.RedisConfig`）、日志（`ei_log.configure_logging`）、OIS（`ei_ois.OisOptions`）
的取值全部经构造参数 / 配置对象注入。唯一读环境变量的地方是 `example/` 下的入口
脚本——那是部署边缘，把环境变量翻译成配置参数后再交给 SDK。

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

| 字段 | 默认 | 说明 |
|---|---|---|
| `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 进程里 `@EI_REGISTER` 多个名字；进程会为每个模型各起
一个 arq Worker，各自消费自己的队列，互不抢任务。

## 失败模式

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

## 全局 KV 与分布式锁

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

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

- `business`：业务简写（如 `ei` / `gcp`），与 arq 队列文法的业务域同义，但两边各自
  独立配置、不互相 import；
- `service_id`：本进程隔离段，取值来自 `ei_identity.IdentityConfig.service_id`。

两段都必填、都按"单段安全 key"校验（字母数字 `._-`、至少一个字母数字，**冒号是段
分隔符，不允许出现在任一段里**）。要往同一个 Redis 里放自己的键，也用
`ei_redis.redis_key("my_ns", "name", business, service_id)` 拼，别手写前缀。

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

本包不读环境变量，也不留进程级单例：Redis 地址、命名空间段、socket 超时全部经
`RedisConfig` 注入，连接池由调用方持有和关闭。

```python
from ei_redis import RedisConfig, create_global_kv

config = RedisConfig(
    business="ei",
    service_id="3fa2b1c9",
    redis_url="redis://scheduler:6379/0",  # 留 None 则 KV 降级为 DisabledGlobalKv
    socket_timeout_sec=1,                  # 一次读写最坏 1 秒返回错误，绝不无界等待
)

kv = create_global_kv(config)              # client 由它新建，用完 kv.close()
kv.set_json("module:state", {"phase": "running"}, ttl_sec=300)  # 不填 TTL 默认 7 天，最长 7 天
state = kv.get_json("module:state")        # 键不存在回 None
```

global KV 与分布式锁通常共用同一个 client（连接池归调用方）：

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

client = config.create_client()            # 调用方负责 client.close()
kv = RedisGlobalKv(client, business="ei", service_id="3fa2b1c9")
lock = RedisDistributedLock(
    client, "leader", ttl_sec=30, business="ei", service_id="3fa2b1c9", auto_renew=True
)

if lock.acquire():                         # 拿不到就是 False，不抛
    try:
        ...                                # 临界区
    finally:
        lock.release()                     # 按 owner token 释放，不会误删别人的锁

is_lock_held(client, "leader", "ei", "3fa2b1c9")  # 只读探活
```

失败模式：

- `RedisConfig.redis_url` 留 `None`：`create_global_kv` 返回 `DisabledGlobalKv` 并记
  一条 warning，每次读写显式抛 `GlobalKvDisabledError`——不会静默退化成进程内内存态，
  也不替调用方猜一个地址。
- Redis 报错：KV 一律包成 `GlobalKvUnavailableError` 抛出；锁的 `acquire()` 返回
  `False`（宁可当成没拿到），`is_lock_held()` 返回 `False`（探活失败即放行）。
- `auto_renew=True` 时按 `ttl_sec / 3`（最小 1 秒）起一个 daemon 线程续期；**首次续期
  失败就停手**，让 TTL 决定锁的生死，`release()` 最多再等它 5 秒退出。
- 指标默认不上报，宿主进程用 `install_kv_metrics_backend(backend)` 注入实现（与
  `ei_ois` 的指标后端同一种做法），SDK 自身不依赖任何指标库。

### 身份：service_id 从哪来

`service_id` 就是上面前缀的那一段，由 `ei_identity.IdentityConfig` 从平台四元组推导
（**不读环境变量**，取值由调用方注入）：

```python
from ei_identity import IdentityConfig

identity = IdentityConfig(
    env="prod",                  # 默认 local
    vdc="cnhbnp01",              # 默认空
    app_name="ei",               # 默认空
    comp_name="runner",          # 默认 ei-pipe-sdk
    version_hash_override=None,  # 给了就直接当 h8（只接受字母数字下划线）
)
identity.version       # 'cnhbnp01:prod:ei:runner'（超 32 字符先整体 md5）
identity.service_id    # md5(version)[:8]，如 '3fa2b1c9'；override 优先；结果带缓存
```

| 字段 | 默认 | 用途 |
|---|---|---|
| `env` | `local` | 环境名 |
| `vdc` / `app_name` | 空 | 平台四元组的另两段 |
| `comp_name` | `ei-pipe-sdk` | 组件名。默认取 `ei-pipe-sdk` 而非 `ei-pipeline`：SDK 进程不是调度组件，落到调度侧默认身份上会和调度服务共用同一个隔离段，正是这层前缀要防的碰撞 |
| `version_hash_override` | `None` | 直接指定 h8，只接受字母数字下划线，本地/测试固定命名空间用 |

worker 要服务的**调度服务** h8 是另一个身份，用
`ei_identity.resolve_service_id(node, service_id, scheduler_url)` 解析：静态给定优先，
否则向调度 `GET /api/status/service-id` 查询（不可达每 10 秒重试）。

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

## 数据库运维（ei_dba）

`ei_dba` 把"建表 / 校验 schema / DB 不 ready 时降级到 dba 模式"做成与具体业务表无关的
能力（从 release/main 的 runner DBA 沉淀、通用化而来）。**表清单不写死**：宿主用
`TableSpec` 声明要管哪些表；连接地址、连接池、超时全部经配置参数注入，本包不读任何
环境变量。需要装 `[dba]` 可选 extra（SQLAlchemy / pymysql / FastAPI），核心 SDK 不依赖。

两条入口语义严格分开，这是"DB 不 ready 时降级"的地基：

- `init_db` **只校验、绝不建表**：缺表抛 `MissingTablesError`（可运行时修复，不缓存）、
  缺列抛 `RuntimeError`（永久 schema 错误，缓存后不再重试）；
- `create_tables` 是**唯一**发 `CREATE TABLE` 的幂等路径，多 pod 并发建表靠 `has_table`
  跳过 + 容忍 "already exists" 兜底，不依赖分布式锁。

```python
from ei_dba import DatabaseConfig, DbaClient, DbaConfig, TableSpec, mysql_url

INFO_DDL = (
    "CREATE TABLE `pipe_info` (\n"
    "  `id` BIGINT PRIMARY KEY AUTO_INCREMENT,\n"
    "  `pipe_id` VARCHAR(64) NOT NULL,\n"
    "  KEY `idx_pipe_id` (`pipe_id`)\n"
    ") ENGINE=InnoDB DEFAULT CHARSET=utf8mb4"
)

config = DbaConfig(
    database=DatabaseConfig(
        url=mysql_url(host="db.internal", user="app", password="s3cr3t", database="ei_pipe")
    ),
    tables=(TableSpec("pipe_info", INFO_DDL, role="info"),),  # ddl 也可给 (name, now) -> str 回调
)

client = DbaClient(config)
client.create_tables()   # 幂等建出缺失的表
client.init_db()         # 只校验：缺表 → MissingTablesError，缺列 → RuntimeError
```

**DB 不 ready 时自动降级到 dba 模式**：宿主用 `dba_guarded_lifespan` 接到
`FastAPI(lifespan=...)`，`init_db` 缺表就摘掉业务路由、挂上 `/api/dba/*`、跳过 worker
启动，等运维 `POST /api/dba/init-db` 建好表后重启进程回到正常模式，而不是崩溃重启：

```python
from functools import partial

from fastapi import FastAPI

from ei_dba import dba_guarded_lifespan, dba_healthz_payload

app = FastAPI(
    lifespan=partial(
        dba_guarded_lifespan,
        client=client,
        business_routers=(pipe_router,),  # 缺表降级时按 router 身份整块摘掉
        startup=start_workers,            # 缺表时绝不跑
        shutdown=stop_workers,
    )
)
app.include_router(pipe_router)


@app.get("/healthz")
async def healthz():
    return dba_healthz_payload(app)  # dba 模式回显缺失表清单，否则 mode=ready
```

DBA 路由只放"看 + 建"的非破坏性子集：`POST /api/dba/init-db`（建缺失表）、
`GET /api/dba/table`（列注册的表 / 角色 / 是否存在）、`GET /api/dba/ddl?table_name=`
（回显注册的 DDL，backend 无关）。刻意**不**含执行任意 SQL、truncate/drop、分区维护、
维护模式锁。字段表、失败模式与完整样例见 `src/ei_dba/README.md`。

## 本地联调

```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/`。
**每个包目录下都有自己的 `README.md`**，讲该包的用法、特性与样例：

```
src/ei_arq/      # 队列文法、worker 心跳、job result 尺寸守卫
src/ei_dba/      # 数据库运维件：通用 DDL 注册 / init_db 校验 / create_tables / dba 模式降级（[dba] extra）
src/ei_identity/   # service_id：本进程身份（IdentityConfig 参数推导）与调度身份（静态 / 调度查询）
src/ei_infer/    # 模型推理 worker（EiInfer / InferNode / EI_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 分布式锁（键自动带 {business}:{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
```
