Metadata-Version: 2.4
Name: scienith-task-queue-django
Version: 0.1.3
Summary: Reusable Django app for durable Task Queue jobs.
Requires-Python: >=3.11
Description-Content-Type: text/markdown
Requires-Dist: Django>=4.2
Requires-Dist: celery>=5.3
Requires-Dist: scienith-task-queue-contract==0.1.3
Provides-Extra: test
Requires-Dist: pytest>=8; extra == "test"

# Task Queue Backend

本仓库发布 Task Queue 的单一可安装 Django app。它提供共享任务队列 runtime 的后端能力，包括 DB 状态真源、event/artifact 持久化、dispatch outbox、Celery task、dispatcher、reconciler 和 handler registry。

本仓库不绑定任何具体业务应用，也不依赖某个业务项目才能运行测试。业务项目安装本 Django app 后，只负责注册自己的 handler、暴露自己的业务 API 或 facade，并把业务对象关联到通用 `JobRun`。

## 仓库边界

- 只发布一个 Django app package：`scienith-task-queue-django`。
- 不额外拆分 core library、service library 或 companion package。
- 共享 API shape 和任务状态类型来自 `task_queue_contract` 的 Python contract package。
- host 项目后续通过 `INSTALLED_APPS`、URL include、settings 和 handler registry 接入。
- 本仓库不包含任何业务 app 的真实计算任务、领域 handler、业务 run facade 或业务权限逻辑。
- `queue_name`、`task_type`、`business_owner_type`、`business_owner_id` 是接入方自定义的命名空间，必须由 host app 自己治理，不能在共享 app 中 hardcode。

## 本地开发

```bash
source .venv/bin/activate
pip install -e ../task_queue_contract/packages/python
pip install -e ".[test]"
python -m pytest
```

## 本地 RabbitMQ

```bash
mkdir -p .tmp/rabbitmq
chmod 777 .tmp/rabbitmq
docker compose up -d rabbitmq
docker compose exec -T rabbitmq rabbitmq-diagnostics -q ping
```

停止：

```bash
docker compose down
```

## Demo Project

完整说明见 `examples/demo_project/README.md`。该 demo project 是宿主 Django 项目示例，不属于可安装 app 发布内容；它注册 `demo.long_task` handler，并提供 `/demo/long-tasks` 用 `sleep` 模拟长耗时任务。

```bash
source .venv/bin/activate
pip install -e ../task_queue_contract/packages/python
pip install -e ".[test]"
python manage.py migrate
python manage.py runserver 127.0.0.1:8010
```

Worker：

```bash
source .venv/bin/activate
celery -A examples.demo_project.celery_app worker -Q task_queue.default --pool=solo --loglevel=INFO
```

Dispatcher：

```bash
source .venv/bin/activate
python manage.py task_queue_dispatcher
```

Dispatcher 会扫描未发布的 dispatch outbox，并只在满足以下条件时发布：

- job 状态为 `queued` 或未完成发布的 `dispatching`。重试等待仍使用 `queued`，通过 `attempt_count` 和 `next_run_at` 判断是否为后续尝试。
- `next_run_at` 为空或已经到期。
- queue 未禁用。
- 当前 queue 的已占用槽位数量小于 `QueueConfig.max_concurrency`；已占用槽位包括 `dispatching` 和 `running` job，当前正在补发的同一个 job 会被排除，避免 broker 短暂失败后无法自重试。

`enqueue_job` 默认只写 DB 和 outbox，由 dispatcher 补发。Host 如果希望事务提交后立即尝试发布，可以传入 `publish_on_commit=True`，或设置 `TASK_QUEUE_PUBLISH_ON_COMMIT=True`；发布失败不会丢 job，错误会记录在 outbox 的 `last_error`，后续 dispatcher 仍可补发。

Reconciler：

```bash
source .venv/bin/activate
python manage.py task_queue_reconciler
```

Reconciler 会处理：

- `running` 且 `lock_until` 过期的 stale job：未超过 `max_attempts` 时重新排队，超过后标记 `timed_out`。
- 停留在 `dispatching` 过久且尚未进入 running 的 job：重新回到 `queued`。这同时覆盖 broker publish 失败、publish 后 dispatcher/worker 竞态失败，以及进程在标记 published 后、worker 接管前中断的恢复路径。
- 停留在 `cancel_requested` 过久的 job：标记 `cancelled`。

可通过 settings 调整：

- `TASK_QUEUE_DISPATCH_TIMEOUT_SECONDS`，默认 `300`。
- `TASK_QUEUE_CANCEL_TIMEOUT_SECONDS`，默认 `300`。

如果 `QueueConfig.retry_policy` 配置了 `backoff_seconds` 与 `backoff_multiplier`，stale job 重新排队时会写入 `next_run_at`，dispatcher 会等到 backoff 到期后再发布。可选的 `jitter_seconds` 或 `jitter` 会在指数退避基础上追加随机抖动。

长时间运行的 handler 可以调用 `context.heartbeat()` 延长 `lock_until`。默认延长 `QueueConfig.heartbeat_seconds`，也可以传 `lease_seconds` 覆盖本次 heartbeat。

提交 demo long task：

```bash
curl -s -X POST http://127.0.0.1:8010/demo/long-tasks \
  -H 'content-type: application/json' \
  -d '{"steps": 8, "sleep_seconds": 1, "business_owner_id": "demo-001"}'
```

## 调试与运维

查看 stuck / active job：

```bash
source .venv/bin/activate
python manage.py shell -c "from scienith_task_queue_django.models import JobRun; print(list(JobRun.objects.exclude(status__in=['succeeded','failed','cancelled','timed_out']).values('id','status','queue_name','task_type','attempt_count','lock_until','celery_task_id')))"
```

处理 stale running job：

```bash
source .venv/bin/activate
python manage.py task_queue_reconciler --once
```

命令输出格式：

```text
requeued=<count> failed=<count> cancelled=<count>
```

手动取消 job：

```bash
source .venv/bin/activate
python manage.py shell -c "from scienith_task_queue_django.models import JobRun; from scienith_task_queue_django.public import request_cancel; request_cancel(JobRun.objects.get(id='<job_id>'))"
```

手动重试 failed / timed_out / cancelled job：

```bash
source .venv/bin/activate
python manage.py shell -c "from scienith_task_queue_django.models import JobRun; from scienith_task_queue_django.public import retry_job; retry_job(JobRun.objects.get(id='<job_id>'))"
python manage.py task_queue_dispatcher --once
```

HTTP API 同样提供：

- `POST /task-queue/jobs/<job_id>/cancel`
- `POST /task-queue/jobs/<job_id>/retry`

## 验证

```bash
source .venv/bin/activate
pip install -e ../task_queue_contract/packages/python
pip install build -e ".[test]"
python manage.py makemigrations --check --dry-run
python -m pytest
python -m build --wheel
```

## Host 接入边界

- Host 项目通过 `INSTALLED_APPS` 安装 `scienith_task_queue_django`。
- Host 项目通过 URL include 暴露 Task Queue API。
- Host 项目注册自己的 `task_type -> handler`。
- 本仓库不包含任何业务应用的真实 handler；demo handler 只用于本地演示，不是可安装 app 的业务能力。
- 本仓库不额外发布 core library 或 service library。

最小接入示例：

```python
INSTALLED_APPS = [
    "django.contrib.contenttypes",
    "scienith_task_queue_django",
]
```

```python
from django.urls import include, path

urlpatterns = [
    path("", include("scienith_task_queue_django.urls")),
]
```

```python
from scienith_task_queue_django.public import register_handler


def my_task_handler(context):
    context.emit_progress("started")
    return {"ok": True}


register_handler("my.task", my_task_handler)
```
