Metadata-Version: 2.5
Name: ringo-task-queue
Version: 0.1.0.dev0
Summary: Python asyncio SDK for Ringo Task Queue with a bundled offline Go daemon
Author: Ringo Task Queue contributors
License-Expression: MIT
License-File: LICENSE
Keywords: asyncio,grpc,postgresql,sqlite,task-queue
Classifier: Development Status :: 3 - Alpha
Classifier: Framework :: AsyncIO
Classifier: Intended Audience :: Developers
Classifier: Operating System :: MacOS :: MacOS X
Classifier: Operating System :: Microsoft :: Windows
Classifier: Operating System :: POSIX :: Linux
Classifier: Programming Language :: Python :: 3
Classifier: Programming Language :: Python :: 3 :: Only
Classifier: Programming Language :: Python :: 3.11
Classifier: Programming Language :: Python :: 3.12
Classifier: Programming Language :: Python :: 3.13
Classifier: Topic :: Software Development :: Libraries :: Python Modules
Requires-Python: >=3.11
Requires-Dist: grpcio-health-checking<2,>=1.71.0
Requires-Dist: grpcio-status<2,>=1.71.0
Requires-Dist: grpcio<2,>=1.71.0
Requires-Dist: protobuf<7,>=5.29.5
Provides-Extra: dev
Requires-Dist: build<2,>=1; extra == 'dev'
Requires-Dist: grpcio-tools<2,>=1.71.0; extra == 'dev'
Requires-Dist: hatchling<2,>=1.27; extra == 'dev'
Requires-Dist: mypy<2,>=1.15; extra == 'dev'
Requires-Dist: pytest-asyncio>=0.23; extra == 'dev'
Requires-Dist: pytest>=8; extra == 'dev'
Requires-Dist: ruff<1,>=0.11; extra == 'dev'
Provides-Extra: generate
Requires-Dist: grpcio-tools<2,>=1.71.0; extra == 'generate'
Provides-Extra: test
Requires-Dist: pytest-asyncio>=0.23; extra == 'test'
Requires-Dist: pytest>=8; extra == 'test'
Description-Content-Type: text/markdown

# Ringo Task Queue — Python SDK

面向 Python 3.11+ 的 asyncio 任务队列 SDK。它支持 SDK 管理的本地 SQLite/PostgreSQL daemon 和远程 gRPC 服务，提供永久去重、至少一次执行、lease fencing、自动续租、重试和 DEAD/DLQ 语义。

## 安装后即可离线运行

```shell
pip install ringo-task-queue
```

正式发布包含五个同版本平台 wheel：

| 系统 | wheel tag | 内置 daemon |
|---|---|---|
| Windows x64 | `win_amd64` | Windows x64 |
| Linux x64 | `manylinux_2_17_x86_64` | Linux x64 |
| Linux ARM64 | `manylinux_2_17_aarch64` | Linux ARM64 |
| macOS Intel | `macosx_11_0_x86_64` | macOS Intel |
| macOS Apple Silicon | `macosx_11_0_arm64` | macOS ARM64 |

`pip` 会自动选择当前平台的 wheel。每个 wheel 已包含匹配的 Go daemon；本地模式不会在安装或运行时下载二进制，不要求安装 Go，也不需要设置 `RINGO_DAEMON_PATH`。PostgreSQL 服务本身不包含在包内。

开发版本需要允许预发布：

```shell
pip install --pre ringo-task-queue
```

## 快速开始

```python
import asyncio
from ringo_task_queue import Ringo


async def main() -> None:
    async with Ringo.local("./data/ringo") as client:
        queue = client.queue("email")
        handled = asyncio.Event()

        async def handle(task, ctx) -> None:
            print(task.id, task.payload)
            await ctx.report_progress(1, 1, "sent")
            handled.set()

        worker = queue.worker(
            handle,
            task_types=["send-email"],
            concurrency=4,
        )
        worker_task = asyncio.create_task(worker.run())

        result = await queue.enqueue(
            unique_key="welcome:user-42",
            task_type="send-email",
            payload={"user_id": 42},
        )
        print(result.task_id, result.created)

        await asyncio.wait_for(handled.wait(), timeout=10)
        await worker.close()
        await worker_task


asyncio.run(main())
```

`(queue, unique_key)` 在任务完整生命周期内永久唯一。重复 enqueue 返回原任务 ID 和 `created=False`，不会创建副本，也不会抛出重复异常。

## 本地模式

### SQLite

```python
async with Ringo.local("./data/ringo") as client:
    queue = client.queue("crawl")
```

SDK 自动完成 daemon 校验、随机 loopback 端口启动、健康检查、协议协商和退出回收。SQLite 数据保存在 `data_dir`；daemon 日志追加到 `data_dir/ringo-daemon.stderr.log`。同一时间不能有两个本地 daemon 占用同一个 data dir。

### PostgreSQL

```python
import os
from ringo_task_queue import PostgresStorage, Ringo


async with Ringo.local(
    "./data/ringo-process",
    storage=PostgresStorage(os.environ["RINGO_POSTGRES_DSN"]),
) as client:
    queue = client.queue("crawl")
```

PostgreSQL 必须由应用提供。SDK 不下载、启动或停止 PostgreSQL。DSN 通过子进程环境传递，不放进进程参数，`PostgresStorage` 的 `repr` 也会隐藏它。

## 远程模式

```python
async with Ringo.connect(
    "https://queue.internal.example:7233",
    token="service-token",
) as client:
    queue = client.queue("crawl")
```

`grpc://`/`http://` 使用明文，`grpcs://`/`https://` 使用 TLS；省略端口时默认 `7233`。

受限链路可以启用整个 gRPC channel 的 gzip：

```python
async with Ringo.connect(
    "https://queue.internal.example:7233",
    token="service-token",
    compression="gzip",
) as client:
    ...
```

`compression` 只接受 `"none"` 或 `"gzip"`。它覆盖 unary RPC 和 Work stream 并在重连后保持；本地模式始终不压缩。

## 提交任务

```python
from datetime import datetime, timedelta, timezone
from ringo_task_queue import RetryPolicy, RetryStrategy, TaskSpec

spec = TaskSpec(
    unique_key="page:10001",
    task_type="fetch-page",
    payload={"url": "https://example.com/10001"},
    priority=10,
    available_at=datetime.now(timezone.utc) + timedelta(seconds=5),
    retry_policy=RetryPolicy(
        max_attempts=5,
        strategy=RetryStrategy.EXPONENTIAL,
        initial_delay=timedelta(seconds=1),
        max_delay=timedelta(seconds=60),
        multiplier=2.0,
        jitter=0.2,
    ),
)
result = await queue.enqueue(spec)
```

也可以使用关键字参数：

```python
result = await queue.enqueue(
    unique_key="page:10001",
    task_type="fetch-page",
    payload={"url": "https://example.com/10001"},
)
```

payload 必须显式提供；显式 `None` 表示 JSON null。允许 JSON 对象、数组和标量，但：

- UTF-8 JSON 编码不超过 1 MiB；
- 不允许非有限浮点数；
- 整数必须位于 `[-2**53, 2**53]`；
- JSON 对象的键必须是字符串；
- 实际使用建议保持在 16–64 KiB 内，大对象只传外部引用。

### 批量提交

```python
batch = await queue.enqueue_many(
    [
        TaskSpec(
            unique_key=f"page:{i}",
            task_type="fetch-page",
            payload={"page": i},
        )
        for i in range(100)
    ]
)
print(batch.created_count, batch.duplicate_count)
```

每批 1–1,000 条，`batch.results` 保留逐项结果。

## 查询和手工重试

```python
task = await queue.get_by_id(result.task_id)
task = await queue.get_by_unique_key("page:10001")

# 始终按业务 unique_key 解释
reopened = await queue.retry("page:10001")

# 显式按服务端 task ID
reopened = await queue.retry_by_id(result.task_id)
```

状态包括 `SCHEDULED`、`READY`、`LEASED`、`SUCCEEDED`、`DEAD` 和 `CANCELLED`。`Task` 还提供 attempt 历史、最后错误、进度和当前 lease 摘要。手工 retry 保留历史和永久去重关系。

## Worker

```python
from datetime import timedelta
from ringo_task_queue import PermanentTaskError, RetryTaskError


async def handle(task, ctx) -> None:
    await ctx.report_progress(1, 2, "started")
    if should_retry(task):
        raise RetryTaskError("rate limited", delay=30)  # 秒
    if invalid(task):
        raise PermanentTaskError("invalid input")
    await process(task)


worker = queue.worker(
    handle,
    task_types=["fetch-page"],
    concurrency=32,
    lease_duration=timedelta(seconds=60),
    grace_period=timedelta(seconds=30),
)
await worker.run()
```

`run()` 一直运行到其他协程调用 `await worker.close()`，并且只能调用一次。通常用 `asyncio.create_task(worker.run())` 启动。

Handler 映射：

| 结果 | 动作 |
|---|---|
| 正常返回 | ACK / `SUCCEEDED` |
| 普通异常 | NACK / 按策略重试 |
| `RetryTaskError` | NACK，可覆盖本次延迟 |
| `PermanentTaskError` | Reject / `DEAD` |
| 取消 | 不发送终态，等待 lease 到期重投 |

Worker 持有 lease token，自动续租并严格控制并发。`TaskContext` 提供 `attempt`、`max_attempts`、`worker_id`、`cancelled`、`wait_cancelled()` 和 `report_progress()`。lease 丢失、stream 断开或强制关闭时，context 会被取消。

这是至少一次系统：handler 的外部副作用必须使用 `unique_key` 或任务 ID 实现幂等。

## 手工 Claim

```python
from datetime import timedelta

leased = await queue.claim(
    limit=10,
    task_types=["fetch-page"],
    wait_timeout=timedelta(seconds=5),
    lease_duration=timedelta(seconds=60),
    worker_id="manual-worker-1",
)

for item in leased:
    try:
        await process(item.task)
        await item.ack()
    except TimeoutError as exc:
        await item.nack(str(exc), delay=timedelta(seconds=10))
    except ValueError as exc:
        await item.reject(str(exc))
```

续租：

```python
await item.extend_lease(timedelta(seconds=60))
```

lease token 不公开。旧 attempt 或迟到 ACK 会被 fencing 拒绝并映射成 `LeaseLostError`。

## 生命周期和关闭

推荐：

```python
async with Ringo.local("./data/ringo") as client:
    ...
```

或手工：

```python
client = Ringo.local("./data/ringo")
await client.start()
try:
    ...
finally:
    await client.close()
```

`close()` 可重复调用，且调用者取消不会中断底层清理。它会先 drain worker，在 grace period 内继续续租并允许 handler 完成，再关闭 channel 和 daemon。被强制取消的 handler 不发送 ACK/NACK。

本地 daemon 意外退出时最多自动重启三次。结果不确定的写操作不会被盲目重放；只读查询可以在重新协商后重试。

## 错误处理

```python
from ringo_task_queue import NotFoundError, RingoError, UnavailableError

try:
    task = await queue.get_by_unique_key("missing")
except NotFoundError:
    ...
except UnavailableError as exc:
    logger.warning("temporarily unavailable: %s", exc)
except RingoError as exc:
    logger.exception(
        "queue error code=%s request_id=%s metadata=%r",
        exc.code,
        exc.request_id,
        exc.metadata,
    )
```

公共异常：

- `InvalidArgumentError`；
- `DuplicateError`（普通重复 enqueue 返回 `created=False`）；
- `NotFoundError`；
- `ConflictError`；
- `LeaseLostError`；
- `UnavailableError`；
- `DeadlineExceededError`；
- `IncompatibleVersionError`；
- `InternalError`。

## 默认值与限制

| 项目 | 默认值或限制 |
|---|---|
| 协议 | 1.0 |
| schema | 3 |
| lease | 60 秒 |
| worker concurrency | 1 |
| graceful shutdown | 30 秒 |
| retry attempts | 默认 3，最大 32 |
| retry | 指数退避，1 秒到 60 秒，multiplier 2，jitter 0.2 |
| batch | 最大 1,000 |
| payload | 最大 1 MiB UTF-8 JSON |
| error text | 最大 256 KiB UTF-8 |
| daemon 自动重启 | 最多 3 次 |

## wheel 完整性和二进制解析

本地 daemon 解析优先级：

1. `Ringo.local(..., daemon_path=...)`；
2. `RINGO_DAEMON_PATH`；
3. wheel 内置平台 daemon；
4. 仅用于开发 checkout 的搜索。

普通 PyPI 安装直接使用第 3 项。SDK 在启动前验证 daemon 是普通文件、长度和 SHA-256 与 manifest 一致，然后运行 `version --json` 校验产品名、版本、协议和 schema。校验失败时拒绝启动，不会静默使用错误二进制。

## 常见问题

### `No matching distribution found`

确认 Python 为 3.11+，对应平台 wheel 已上传；安装 `.dev0` 时使用 `--pre`。

### 找不到 daemon

正式 wheel 不应出现。使用 `python -m pip show -f ringo-task-queue` 检查是否误装了 editable checkout 或非正式产物。

### `ConflictError`

另一个进程通常已占用同一个 SQLite data dir。复用已有客户端或换一个目录。

### 任务重复执行

这是至少一次语义的正常边界。使用任务 ID/`unique_key` 为外部副作用建立幂等约束。

### Worker 关闭较慢

默认允许 handler 在 30 秒内完成。可以调整 `grace_period`，并让 handler 响应 cancellation 或 `ctx.cancelled`。

## 维护者构建平台 wheel

真实 daemon 不提交到源码树。必须先从同一源码版本生成并验证五平台 release manifest，再逐个平台构建 wheel：

```shell
go run ./tools/release build --out dist --version 0.1.0-dev --commit <revision> --date <RFC3339-UTC>
go run ./tools/release verify --manifest dist/manifest.json
```

然后从 `sdk/python` 运行：

```shell
python tools/build_platform_wheel.py \
  --binary ../../dist/ringo-task-queue-linux-amd64 \
  --release-manifest ../../dist/manifest.json \
  --platform linux-amd64 \
  --output-dir ../../package-dist/python
```

分别构建 `windows-amd64`、`linux-amd64`、`linux-arm64`、`macos-amd64`、`macos-arm64`。构建会记录根 manifest 的 SHA-256 来源，完成后清理暂存二进制。上传前运行：

```shell
python tools/verify_wheel_set.py ../../package-dist/python
python -m twine check ../../package-dist/python/*
```

没有经过完整五目标 manifest 暂存的 wheel 和 sdist 会被主动拒绝，避免发布不含 daemon 或平台集合不完整的包。

## 许可证

MIT License。
