Metadata-Version: 2.5
Name: ys-data-client
Version: 0.1.0
Summary: YS 数据中心 Python SDK：可靠事件发布订阅与跨系统 RPC
Requires-Python: >=3.13
Requires-Dist: websockets<18,>=15
Description-Content-Type: text/markdown

# ys-data-client

YS 数据中心 Python SDK。业务系统只需配置数据中心地址、端口和系统 token，即可通过装饰器完成两类跨系统能力：

| 需求                                   | 机制  | SDK 用法                                   |
| -------------------------------------- | ----- | ------------------------------------------ |
| 状态通知、待办、最终状态变化           | Event | `@client.event(...)` / `client.publish()`  |
| 短时查询、需要直接返回值的命令         | RPC   | `@client.rpc(...)` / `client.call()`       |
| 长耗时或高风险命令（付款、审批）       | 组合  | RPC 返回“已受理 + 任务 ID”，再用 Event 通知最终状态 |

SDK 封装了 `/system` WebSocket 协议、认证、订阅、方法注册、心跳、指数退避重连、请求关联和 ACK，业务代码不需要手工发送任何 JSON。

## 安装

要求 Python 3.13+。在其他 uv 项目中以本地路径引用（开发阶段推荐可编辑安装）：

```toml
[project]
dependencies = ["ys-data-client"]

[tool.uv.sources]
ys-data-client = { path = "../ys_data_client", editable = true }  # 使用路径进行引入
```

然后执行 `uv sync`。也可以直接 `uv add --editable ../ys_data_client`。

## 三项连接配置

```python
import os

from ys_data_client import YsDataClient

client = YsDataClient(
    url=os.environ.get("YS_DATA_URL", "localhost"),
    port=int(os.environ.get("YS_DATA_PORT", "9101")),
    token=os.environ["YS_DATA_TOKEN"],
)
```

- `url` 接受 `host`、`host:port`、`http://`、`https://`、`ws://`、`wss://`，以及已带 `/system` 的地址；SDK 自动补 `/system`，`http(s)` 自动映射为 `ws(s)`。
- 显式传入的 `port` 优先于 URL 内的端口。
- token 只通过认证消息发送，禁止放进 URL query string；系统编码 `client_id` 由数据中心根据 token 识别，不需要配置。兼容旧协议时可选传 `client_id=`。

## 最小可运行示例

```python
import asyncio
import os

from ys_data_client import EventContext, RpcBusinessError, RpcContext, YsDataClient

client = YsDataClient(url="localhost", port=9101, token=os.environ["YS_DATA_TOKEN"])


@client.event("demo.order.created")
async def on_order_created(payload: dict, context: EventContext) -> None:
    # 正常返回即自动 ACK；至少一次投递，必须按 context.event_id 做业务幂等
    print("order created", payload, context.event_id)


@client.rpc("demo.inventory.reserve")
async def reserve_inventory(params: dict, context: RpcContext) -> dict:
    if params.get("quantity", 0) <= 0:
        raise RpcBusinessError("INVALID_QUANTITY", "数量必须大于 0", {"quantity": params.get("quantity")})
    return {"reserved": True, "reservation_id": f"RES-{params['order_id']}"}


async def main() -> None:
    async with client:
        await client.wait_until_ready(timeout=10)

        result = await client.call(
            "demo.inventory.reserve",
            target="ys-test02",
            params={"order_id": "ORDER-001", "quantity": 2},
            timeout=10,
            idempotency_key="order:ORDER-001:reserve",
        )
        print(result)  # 目标系统 handler 的真实返回值

        await client.publish(
            "demo.order.created",
            payload={"order_id": "ORDER-001"},
            event_id="order:ORDER-001:created",
        )


asyncio.run(main())
```

事件订阅和 RPC 方法由装饰器自动收集，认证成功后 SDK 自动发送 `rpc.register` 与 `subscribe`，重连后自动恢复。

## ys_base AppServer

公司业务系统使用 `ys_base.server.AppServer` 时，把 SDK 挂到它的生命周期钩子即可，不需要 FastAPI lifespan：

```python
server = AppServer(settings, app_version=__version__)
server.on_started(client.start)  # 启动完成后连接数据中心；需要等待就绪可再包一层 wait_until_ready
server.on_stopping(client.close)  # 停止前优雅断开
```

业务入口（`EventRouter` 的 WS 事件或 `server.add_router` 挂载的 HTTP 路由）里直接 `await client.call(...)` / `await client.publish(...)`。

## FastAPI lifespan

```python
from fastapi import FastAPI

app = FastAPI(lifespan=client.lifespan(ready_timeout=5))
```

启动时后台连接数据中心并最多等待 `ready_timeout` 秒；等不到也不会阻止应用启动，SDK 在后台继续重连。应用关闭时会等待进行中的 handler 收尾，再关闭连接。需要自定义 lifespan 时，用 `async with client:` 或 `await client.start()` / `await client.close()` 包裹即可。

## 提供 RPC：`@client.rpc`

```python
@client.rpc("demo.inventory.reserve")
async def reserve(params: dict, context: RpcContext) -> dict: ...
```

- handler 可以是 `async def` 或普通 `def`；返回值必须可 JSON 序列化，原样成为调用方的结果。
- 抛出 `RpcBusinessError(code, message, details)` 会作为业务错误返回给调用方；`code` 必须是稳定的业务错误码。
- 其他异常会被消毒为通用 `INTERNAL_ERROR` 返回，细节只保留在本地日志。
- `context.idempotency_key` 是调用方提供的幂等键，有副作用的方法应据此去重；`context.remaining()` 返回距调用方超时还剩多少秒。
- 收到时已过期或已被数据中心取消的请求不会再执行。

数据中心只会把请求转发给已在 ACL `serve_methods` 中授权、且当前在线注册了该方法的连接。同一系统同一方法只允许一个在线连接提供，重复注册会被拒绝并记录日志。

## 调用 RPC：`await client.call`

```python
result = await client.call(method, target=..., params={...}, timeout=10, idempotency_key="...")
```

`target` 和 `method` 是本次调用的路由参数，由业务代码明确表达。异常：

| 异常                  | 含义                                         | 远端是否可能已执行 |
| --------------------- | -------------------------------------------- | ------------------ |
| `RpcRemoteError`      | 目标 handler 返回业务错误，含 `code/details` | 是（业务已决定）   |
| `RpcRejectedError`    | 数据中心拒绝：无调用权限、目标无提供权限、参数无效、待处理过多 | 否 |
| `RpcUnavailableError` | 目标离线、未注册方法、处理期间断线、SDK 未连接 | 看 `maybe_executed` |
| `RpcTimeoutError`     | 超时未收到目标响应                           | 是                 |

所有 `RpcError` 都带 `maybe_executed`。为 True 时不能盲目重试有副作用的调用，应复用同一个 `idempotency_key` 或转人工核对。SDK 不会在结果未知时自动重试，也不会在重连后重发断线前的请求。

`timeout` 默认 10 秒，可逐次覆盖；数据中心会把它夹在 0.5 到 60 秒之间，到期后由数据中心返回超时，并通知目标方取消。超时不代表目标业务未执行；取消等待也不代表远端已取消。付款、长耗时审批等操作应让 RPC 只返回受理结果和任务 ID，再用 Event 通知最终状态。

## 订阅事件：`@client.event`

```python
@client.event("demo.inventory.changed")
async def on_changed(payload: dict, context: EventContext) -> None: ...
```

- handler 正常返回，SDK 自动回复 `acknowledged`。
- 抛出 `EventRejected("原因")`，SDK 回复 `business_rejected`；该 event_id 之后不可重放，发布方需要修正载荷并使用新的 event_id。
- 抛出 `EventRetry("原因")` 或任何其他异常，SDK 回复 `retry`：数据中心保留投递并按退避重投，达到失败上限后进入死信。接收循环不会因 handler 异常退出。
- `context` 提供 `event_id`、`event_type`、`delivery_id`、`sequence`、`status_version`、`publisher_client_id`、`published_at`、`metadata`。

事件是至少一次投递：断线重连、并发投递、数据中心重试都可能让同一个 `event_id` 重复到达，业务必须用持久存储按 `event_id` 幂等。SDK 不提供进程内去重，进程内 `set` 不能替代生产级幂等。

## 发布事件：`await client.publish`

```python
result = await client.publish(
    "demo.order.created",
    payload={"order_id": "ORDER-001"},
    event_id="order:ORDER-001:created",
    metadata={"trace_id": "..."},
    status_version=None,
)
result.accepted, result.duplicate, result.target_client_ids
```

返回值只表示数据中心已接受并持久化事件，并已为订阅方创建投递记录，不代表订阅方已处理。`event_id` 是传输幂等键，必须稳定，重复发布相同内容返回 `duplicate=True`；被 ACL、event_id 冲突或 status_version 倒退拒绝时抛出 `PublishRejectedError`。

## 其他异常

- `ConfigurationError`：地址、token 或 handler 注册无效。
- `AuthenticationError`：token 被拒绝。SDK 会记录错误并继续退避重连，需要人工检查系统身份是否启用、过期或凭证已轮换。
- `ConnectionClosedError`：未连接或等待期间断开。
- `RequestTimeoutError` / `RequestRejectedError`：publish、ack 等协议请求超时或被拒绝。

## 优雅关闭

`await client.close()`（或退出 `async with` / lifespan）会：停止重连；等待进行中的 handler 收尾（默认 5 秒，可传 `drain_timeout`），期间仍可回 ACK 和 RPC 结果；关闭连接；让所有等待中的请求以明确异常结束；不遗留后台任务。

## 日志

SDK 使用标准库 `logging`，logger 名为 `ys_data_client`。日志不会输出 token 和完整业务载荷。

## 开发

```text
uv sync
uv run pytest -q
uv run ruff check .
uv run pyright src
```

测试使用内置的最小协议模拟服务，不需要真实数据中心。
