diff --git a/apps/api/management/commands/run_image_tasks.py b/apps/api/management/commands/run_image_tasks.py index af485a7..d5ec8c7 100644 --- a/apps/api/management/commands/run_image_tasks.py +++ b/apps/api/management/commands/run_image_tasks.py @@ -4,6 +4,7 @@ import uuid from django.conf import settings from django.core.management.base import BaseCommand +from apps.api.models import ImageGenerationTask from apps.api.image_tasks import reap_stale_image_tasks, run_one_image_task @@ -56,7 +57,7 @@ class Command(BaseCommand): task = run_one_image_task(worker_id=worker_id) if task is not None: - self.stdout.write(f"task={task.task_id} status={task.status}") + self.stdout.write(format_task_log_line(task, started_at=now)) if once: return continue @@ -64,3 +65,18 @@ class Command(BaseCommand): if once: return time.sleep(sleep_seconds) + + +def format_task_log_line(task: ImageGenerationTask, *, started_at: float) -> str: + duration_ms = max(0, int((time.monotonic() - started_at) * 1000)) + alias = str((task.request_payload or {}).get("model") or "") + fields = { + "event": "image_task_processed", + "task_id": str(task.task_id), + "status": task.status, + "alias": alias, + "duration_ms": str(duration_ms), + } + if task.status in {ImageGenerationTask.Status.FAILED, ImageGenerationTask.Status.EXPIRED}: + fields["error_code"] = task.error_code or "upstream_error" + return " ".join(f"{key}={value}" for key, value in fields.items()) diff --git a/apps/api/tests.py b/apps/api/tests.py index f627d33..f7b2798 100644 --- a/apps/api/tests.py +++ b/apps/api/tests.py @@ -1,5 +1,6 @@ import uuid import base64 +import io import json import tempfile from datetime import timedelta @@ -13,6 +14,7 @@ from django.conf import settings from django.contrib import admin from django.contrib.auth import get_user_model from django.core.cache import cache +from django.core.management import call_command from django.test import TestCase, override_settings from django.urls import path from django.utils import timezone @@ -1585,6 +1587,34 @@ class GenerateApiTests(TestCase): self.assertEqual(poll.data["status"], ImageGenerationTask.Status.FAILED) self.assertEqual(poll.data["error"]["code"], "upstream_timeout") + def test_run_image_tasks_logs_failed_task_alias_error_and_duration(self): + self.provider.image_error = requests.Timeout("image upstream deadline exceeded") + response = self.post_with_provider( + "/api/v1/generate/image/tasks", + {"prompt": "生成图片", "model": self.image_alias, "resolution": "1K"}, + ) + out = io.StringIO() + + with patch("apps.api.generation.get_provider", return_value=self.provider): + call_command( + "run_image_tasks", + "--once", + "--worker-id", + "worker-log", + stdout=out, + ) + + task = ImageGenerationTask.objects.get(task_id=response.data["task_id"]) + output = out.getvalue() + self.assertEqual(task.status, ImageGenerationTask.Status.FAILED) + self.assertIn("event=image_task_processed", output) + self.assertIn(f"task_id={task.task_id}", output) + self.assertIn(f"alias={self.image_alias}", output) + self.assertIn("status=failed", output) + self.assertIn("error_code=upstream_timeout", output) + self.assertRegex(output, r"duration_ms=\d+") + self.assertNotIn("生成图片", output) + def test_async_image_reaper_fails_stale_running_task_and_refunds(self): response = self.post_with_provider( "/api/v1/generate/image/tasks", diff --git a/docs/current-state.md b/docs/current-state.md index 0878190..9a1badf 100644 --- a/docs/current-state.md +++ b/docs/current-state.md @@ -11,13 +11,14 @@ ## 当前快照 -- 日期:2026-07-08 +- 日期:2026-07-09 - 阶段:Phase 6 增强(MVP 后);Phase 3 对外 API 与充值已完成到 T-306,Phase 4 用户端 T-501 注册 / 登录(allauth)、T-502 API Key 自助管理页、T-503 个人中心 / 记录页、T-504 充值页与 T-505 用户端审核优化已完成,Phase 5 T-401 运营后台完善、T-402 MVP 完整验收与 T-403 部署 / 运行文档已完成,Phase 6 T-601 可用别名发现、T-602 django-admin 中文化第 1-3 层、T-603 django-admin 字段级中文化、T-604 中文敏感词本地过滤、T-605 免邮箱验证策略落地、T-606 公开首页 + 客户端下载入口、T-607 桌面端最新版本检查接口、T-608 新用户注册赠送 100 点试用点数、T-609 桌面端版本检查接口增加强制更新标记、T-610 首页导入模板下载入口、T-611 用户端品牌名统一为虾皮圈、T-612 生图同步接口止血、T-613 抽生成核心 service、T-614 生图异步任务化接口与 T-615 旧同步生图接口遥测 / 弃用口径已完成;后续仍需处理真实支付回调到账闭环、客户端下载包发布和生产侧旧同步接口用量观察 - 技术栈:系统 Python 3.12.3 + Django 5.2.15 + DRF 3.16.1 + django-allauth 65.18.0 + PyMySQL 1.1.3 + cryptography 49.0.0 + requests 2.34.2 + ahocorapy 1.6.2 + wechatpayv3 2.0.2 + python-alipay-sdk 3.4.0 + django-admin;MySQL 8.4 已接入 settings,并支持 `MYSQL_CONNECT_TIMEOUT` / `MYSQL_READ_TIMEOUT` / `MYSQL_WRITE_TIMEOUT`;用户端已用 Django 模板 SSR + Bootstrap + allauth 落地注册登录;生产部署口径为 VPS / 宝塔 + Nginx + Gunicorn(gthread) + systemd;详见 `03-tech-stack.md` 与 `deployment.md` - 生产代码:已有最小 Django 工程骨架:`manage.py`、`config/`;T-002 已创建 `apps/users|portal|billing|ai|api`;T-003 已把自定义 `User` 注册进 django-admin;T-004 已完成 email 唯一性、init 版本断言、app 顺序、`.env.example` 与 `pyproject.toml`;T-101 已新增 `apps/ai/providers/`(Provider 接口、注册表、chat/gemini/images/images_edits 适配器);T-102 已新增 `AiModel` / `ModelAlias`、Fernet 加密密钥存储、别名解析、admin 配置页、`import_ai_models` 导入命令;T-103 已新增 `AiConfigAuditLog` 审计表、admin 只读页面和后台保存/删除审计 hook;T-104/T-105 已完成录制 title/image smoke 与审核修补;T-201 已新增 `UserWallet` / `ApiKey`、`PointsLedger` / `CallRecord`、对应 admin 与迁移;T-202 已新增 `PricingRule` / `ExchangeRate`、`apps.billing.pricing` 计费计算函数、admin 配置页与迁移;T-203 已新增 `apps.billing.services`,实现并发安全预扣、成功确认与幂等失败退点;T-204 已新增 `billing.0003_pointsledger_unique_ledger_change_type_per_call`,用 MySQL 可落地的 `ref_call + change_type` 复合唯一约束兜底防重复 refund;T-301 已新增 `apps.api.authentication.ApiKeyAuthentication` 与 `ExternalApiView`;T-302 已新增生成接口编排、序列化器、图片本地存储和 `/api/v1/generate/title|image` 路由;T-303 已新增 `apps.billing.services.get_balance_snapshot()` 与 `/api/v1/balance` 余额查询接口;T-304 已新增 `RechargeOrder`、充值回调验签适配器、幂等入账服务、微信/支付宝回调路由与迁移 `billing.0004_rechargeorder_and_more`;T-305 已新增 `create_recharge_order()`、微信/支付宝扫码下单 mock/SDK 入口、`/api/v1/recharge/create` 与 `/api/v1/recharge/status`;T-306 已新增 `apps.api.throttles`、`apps.api.exceptions`、`REST_FRAMEWORK` 安全默认认证、生成/认证失败限流、`image_url` SSRF 防护与响应大小上限、充值单笔金额上限;T-501/T-608 已接入 allauth 注册登录路径,注册成功后 adapter 调用 `grant_signup_bonus()` 经 billing 一次性发放 100 点并写 `signup_bonus` 流水,新增 `SignupBonusGrant(user UNIQUE)` 幂等标记、admin 只读检索和 `billing.0007` 迁移;T-502 已新增 `/apikeys`、API Key 创建表单、列表页和删除(吊销)动作,生成后明文只显示一次,列表只显示 prefix;T-503/T-608 已扩展 `/dashboard` 为个人中心汇总,并新增 `/records/recharge` 充值记录与 `/records/usage` 点数记录,只读展示当前用户数据和注册赠点 / 消费 / 退款流水;T-504 已新增 `/recharge` 页面、`RechargeCreateForm`、充值导航入口和轮询脚本,页面创建 pending 订单、展示二维码票据、轮询 `/api/v1/recharge/status`,订单 paid 后刷新余额;T-505 已把 Bootstrap 5 CSS 与 qrcode.js vendoring 到 `apps/portal/static/portal/vendor/`,页面不再依赖 jsdelivr,并把充值记录 / 点数记录改为 Django `Paginator` 分页;T-401 已新增 `adjust_wallet_points()` 手工调点服务、钱包 admin 专用调点表单与模板,后台可管理/检索用户、钱包、API Key(脱敏)、计费规则、汇率、充值订单、点数流水、注册赠点记录和调用记录,流水/订单/调用记录保持只读;T-402 已新增 `docs/mvp-acceptance.md`,按 P0 验收矩阵记录 MVP 完整验收结论、测试证据和已知限制;T-403 已新增 `docs/deployment.md` 与 `requirements-production.txt`,并在 settings 中补齐 `STATIC_ROOT`、`CSRF_TRUSTED_ORIGINS`、共享 `CACHES`、HTTPS cookie、proxy SSL、HSTS 环境变量与 `ACCOUNT_SIGNUP_RATE_LIMIT` 注册限流配置;T-601 已新增 `apps.ai.catalog.get_public_model_catalog()`、`GET /api/v1/models` 与 portal `/models` 只读页面,只展示 active 可调用别名、能力、是否需要原图和点数单价,不解密 provider key,不暴露底层 SKU / URL / key / `extra_body`;T-604 已新增 `apps.moderation`、`SensitiveWord` 模型/admin/迁移、keyword provider、归一化管线和共享 cache 版本失效,生成接口已改为 prompt 先审再读取图片/计费/扣点/调上游;T-606 已新增公开首页 `/`、`DownloadRelease` 模型/admin/迁移、首页 SSR 模板、共享 `portal/brand.css`,并把现有 portal 页面套入同一套品牌 token;T-607/T-609 已新增 `ClientLatestReleaseView` 与 `/api/v1/client/releases/latest`,公开匿名返回当前客户端版本 JSON,`release.force_update` 表示该版本是否强制升级;`portal.0002_downloadrelease_force_update` 已给 `DownloadRelease` 增加 `force_update` 字段,admin 可编辑和筛选;T-610 已新增 `ImportTemplate` 模型/admin/迁移 `portal.0003_importtemplate`,首页读取当前模板并在“下载客户端”旁展示“下载导入模板”,本地文件 URL 转为当前站点绝对 URL,`external_url` 优先;T-611 已把用户端 portal 可见品牌名统一为“虾皮圈”,包括页面标题、顶部导航、首页 H1、用户端“虾皮圈 API Key”文案和 allauth 邮件模板;T-612 已新增 `AI_IMAGE_UPSTREAM_DEADLINE_SECONDS`,生图 Provider 上游请求和上游返回图片 URL 下载会按 `min(AiModel.timeout_seconds 或分辨率默认值, 硬截止)` 控制读取超时,超时返回 `upstream_timeout` 并走既有失败退点路径;T-613 已把旧同步生成链路抽成 `GenerationInput`、`prepare_generation()`、`precharge_generation()`、`execute_precharged_generation()` 与 `GenerationResult`,旧 view 只负责 serializer 和异常转 HTTP,后续异步 worker 可复用已预扣执行 / 确认 / 退点阶段。 - T-614 生产代码补充:已新增 `ImageGenerationTask` 与 `api.0001_initial`,新增 `POST /api/v1/generate/image/tasks`、`GET /api/v1/generate/image/tasks/{task_id}`、`apps.api.image_tasks` 任务服务、`run_image_tasks` management command 和只读 admin;异步提交支持 `Idempotency-Key` 去重 / 冲突检测,worker 使用 DB 任务表、租约、心跳与 reaper,成功返回 cmhub 托管 URL,失败 / 超时 / 僵任务走计费层幂等退款。 - T-614 实现取舍补充:worker 执行前会基于任务快照复跑 `prepare_generation()`,即复审 prompt 并重解析别名 / Provider / 定价;账务使用已预扣 `CallRecord.points_cost`,不会重复扣点。短队列下这是偏安全取舍,词库变更后排队任务仍可被拦截并退款;若后续队列积压或频繁切换模型,应单独做提交时模型配置快照。异步 submit 对 `image_base64` 只解码并保存输入文件,桌面端主链路推荐继续使用;`image_url` 会在 submit 阶段完成 SSRF 校验、远程下载和大小限制,可能阻塞提交请求,属于边缘兼容路径。 - T-615 生产代码补充:已新增 `apps.api.telemetry`,旧同步 `POST /api/v1/generate/image` 与新异步提交 `POST /api/v1/generate/image/tasks` 都写 `cmhub.api.generation_usage` 结构化日志事件 `generation_route_usage`;字段白名单为 `route_type`、API Key ID/前缀、`user_id`、`X-Client-Version`、别名、状态、耗时、错误码和 HTTP 状态,不记录 API Key 明文、prompt、图片 base64 或 provider raw。 +- 异步生图 worker 排障补充:`run_image_tasks` 每处理一个任务会输出 `event=image_task_processed task_id=... status=... alias=... duration_ms=...`,失败 / 过期任务额外输出 `error_code`;日志不包含 prompt、`image_base64`、provider raw 或密钥。生产扩容 worker 应按 2、4、8、16 逐级观察队列长度、耗时、错误码、MySQL 连接数和上游失败率,不建议直接扩到 100。 - 最新验证:T-615 旧同步生图接口遥测已验证 `py -3.12 -m py_compile apps\api\telemetry.py apps\api\views.py apps\api\tests.py` 通过;新增 3 条目标测试通过,覆盖旧同步成功日志、新异步提交成功日志、异步余额不足错误日志,确认可按 client version / api_key / route_type 查询且不泄露 prompt、base64、完整 API Key 或 provider raw;`py -3.12 manage.py check` 通过,0 issues;`py -3.12 manage.py makemigrations --check --dry-run` 通过,No changes detected;`py -3.12 manage.py test apps.api.tests.GenerateApiTests --keepdb --noinput --verbosity 2` 通过,31 tests OK;`.\init.ps1` 通过;`git diff --check` 通过,仅 Windows CRLF 提示。测试期仍保留 allauth 在 MySQL 条件唯一约束上的既有 `models.W036` 警告。 - 用户端导航:顶部导航 active 状态已修复,`portal/base.html` 基于 `request.resolver_match.url_name` 高亮当前页面入口,并用 `aria-current="page"` 标记;「充值」不再在非充值页固定深色高亮。 - 最新验证:T-614 生图异步任务化接口已验证 `py -3.12 manage.py makemigrations api` 生成 `api.0001_initial`;`py -3.12 manage.py test apps.api.tests.GenerateApiTests -v 2 --keepdb` 通过,28 tests OK,覆盖敏感词拦截不建任务/不扣点、余额不足 402、Idempotency-Key 去重和冲突、worker 成功轮询同一 URL、跨用户拒绝、失败/超时退款、reaper 僵任务退款、重复 worker 幂等和 reaper 退款后的迟到 worker 不可改回成功。首次不带 `--keepdb` 运行时因已有 `test_cmhub` 触发交互式删除确认导致 EOF 中断;重跑 `--keepdb` 通过。测试期仍保留 allauth 在 MySQL 条件唯一约束上的既有 `models.W036` 警告。 diff --git a/docs/deployment.md b/docs/deployment.md index 73dc5ac..4d11008 100644 --- a/docs/deployment.md +++ b/docs/deployment.md @@ -281,6 +281,10 @@ python3.12 manage.py run_image_tasks \ 建议单独托管为 `cmhub-image-worker.service`。该 worker 从数据库 `image_generation_task` 表抢 `queued` 任务,使用 MySQL `select_for_update(skip_locked)` 标记 `running`,执行成功后写 `succeeded` 和稳定 `result_url`;失败或上游超时会调用计费层退点并写 `failed`。worker 循环会按 `IMAGE_TASK_REAPER_INTERVAL_SECONDS` 扫描租约或心跳过期的 `running` 任务,默认判失败并幂等退点,不默认重排队。 +worker 每处理一个任务会向 stdout 输出一行结构化日志,形如 `event=image_task_processed task_id=... status=failed alias=... duration_ms=... error_code=upstream_timeout`。失败日志必须用于区分「还在 queued 未提交给上游」和「已 running 但上游超时 / 失败」;日志不得包含 prompt、`image_base64`、provider raw 或密钥。 + +多 worker 可以并行运行同一命令,只要 `--worker-id` 不同即可;MySQL 8.4 会通过 `select_for_update(skip_locked)` 避免重复抢同一任务。生产扩容应按 2、4、8、16 逐级观察 `queued` 长度、`duration_ms`、`error_code`、MySQL 连接数、VPS CPU/内存/磁盘写入和上游失败率。不要直接扩到 100 个 worker:这会同时放大 MySQL 连接、上游请求、图片下载和本地写文件压力;如果上游已经频繁 `upstream_timeout`,100 并发通常只会把失败更快放大。 + 当前 worker 执行前会基于任务快照复跑一次生成准备逻辑,包括 prompt 复审、别名 / Provider 解析和定价检查;账务仍使用 submit 阶段已预扣的 `CallRecord.points_cost`,不会重复扣点。这是短队列下偏安全的取舍:敏感词库变更后,排队任务仍可在执行前被拦截并退款。若生产出现明显排队或频繁切换别名 / 模型,应单独开发“提交时模型配置快照”,让 worker 使用提交时确认的模型执行。 输入方式对 submit 耗时有直接影响:桌面端批量生图应优先传 `image_base64`,submit 阶段只解码并写入输入文件;`image_url` 会在 submit 阶段完成 SSRF 校验、远程下载和大小限制,再保存为输入文件引用,因此可能阻塞提交请求。`image_url` 的好处是任务进入队列后自包含,worker 不再访问调用方外部 URL;生产排查 submit 慢时,应先确认是否有客户端批量使用 `image_url`。 diff --git a/progress.md b/progress.md index fde0d0c..ecc2f0d 100644 --- a/progress.md +++ b/progress.md @@ -1774,3 +1774,20 @@ - 验证:仅文档更新;使用 `git diff --check` 检查格式。 - 决策:当前不改代码。`image_base64` 是桌面端主路径,submit 仍为短请求;`image_url` 阻塞属于边缘兼容路径。模型配置快照不作为热修,待队列积压或多模型价差扩大后单独立任务。 - 下一步:无需立即编码;继续按生产优先级处理支付回调闭环、客户端发布和线上遥测观察。 + +## 2026-07-09 热修:异步生图 worker 失败日志与扩容口径 + +- 状态:DONE。 +- 背景:线上 admin 中多条生图任务停留在“待处理”。排查确认 `POST /api/v1/generate/image/tasks` 已返回 `202`,任务进入 `queued`;当前只有 1 个 `cmhub-image-worker` 串行消费,部分任务在上游 180 秒硬截止后 `upstream_timeout`,导致后续 queued 积压。 +- 代码变更: + - `apps/api/management/commands/run_image_tasks.py`:worker 每处理一个任务输出结构化 stdout 行 `event=image_task_processed task_id=... status=... alias=... duration_ms=...`;失败 / 过期任务额外输出 `error_code`。 + - `apps/api/tests.py`:新增目标测试,覆盖失败任务日志包含 `task_id`、别名、`error_code` 与 `duration_ms`,且不输出 prompt。 +- 文档变更: + - `docs/deployment.md`:补充 worker 任务日志、排障口径和扩容策略;明确不要直接扩到 100 个 worker,应按 2、4、8、16 逐级观察。 + - `docs/current-state.md`:同步当前 worker 日志字段和扩容口径。 +- 验证: + - `py -3.12 -m py_compile apps\api\management\commands\run_image_tasks.py apps\api\tests.py`:通过。 + - `py -3.12 manage.py check`:通过,0 issues。 + - `py -3.12 manage.py test apps.api.tests.GenerateApiTests.test_run_image_tasks_logs_failed_task_alias_error_and_duration --keepdb --noinput --verbosity 2`:通过,1 test OK。 + - `py -3.12 manage.py test apps.api.tests.GenerateApiTests --keepdb --noinput --verbosity 1`:通过,32 tests OK。 +- 决策:暂不把线上 `cmhub-image-worker` 直接扩到 100。100 个进程会同时放大 MySQL 连接、上游请求、图片下载和本地写文件压力;在上游已经出现 `upstream_timeout` 时,直接 100 并发更可能放大失败率。建议先用新日志观测后按 2、4、8、16 逐级扩容。