9.0 KiB
9.0 KiB
Integration(FastAPI API + Celery Worker)|Tasks
对应计划:
spec_kit/Personalized Reco/modules/integration-api-worker/plan.md执行规则:
- 本任务清单详细可执行;每项完成后将 “状态:未开始” 改为 “状态:已完成”,并补充证据(命令输出/截图/测试用例)。
- 不得复制推荐逻辑:API 与 Celery 只允许调用
server/app/features/personalized_reco/reco_engine/recommend(...)。- 多语言仅支持 EN/TC;
Accept-Language映射后必须通过normalize_locale校验。- 限流:按客户端 IP,1 分钟 10 次,超限返回 429。
- now 注入:支持
X-Nowheader(ISO8601),并保留 body.now;优先级:header > body > server now。
0. 准备与对齐(不改代码)
-
确认现有 FastAPI/Celery 入口与依赖注入方式(状态:已完成)
- 检查点:
- FastAPI app 创建:
server/app/main.py - DB session 依赖:
server/app/db/session.py:get_db - Celery app:
server/app/worker.py:celery_app且自动发现任务app.tasks
- FastAPI app 创建:
- 证据:
- FastAPI:
server/app/main.py使用create_app()并include_router(...) - DB:
server/app/db/session.py提供get_db()与AsyncSessionLocal - Celery:
server/app/worker.py使用celery_app.autodiscover_tasks(["app.tasks"])
- FastAPI:
- 检查点:
-
确认 reco_engine 对外入口可用(状态:已完成)
- 检查点:
server/app/features/personalized_reco/reco_engine/__init__.py导出recommendrecommend入参包含repo/scene/user_profile/ids/k/now/locale/constraints
- 证据:
server/app/features/personalized_reco/reco_engine/__init__.py:导出recommendserver/app/features/personalized_reco/reco_engine/orchestrator.py:async def recommend(...)
- 检查点:
1. FastAPI:推荐接口(按场景拆分)
-
新增路由文件
server/app/api/v1/reco.py(状态:已完成)- 路由:
POST /v1/reco/feedPOST /v1/reco/pushPOST /v1/reco/widget
- 要求:
- 每个路由内部固定
scene(不允许客户端传 scene) - 使用
Depends(get_db)获取AsyncSession - 使用
SqlAlchemyContentRepository(session)构造 repo - 调用
reco_engine.recommend(...)并直接返回items/meta
- 每个路由内部固定
- 证据:路由文件路径 + handler 函数列表
- 证据:
- 文件:
server/app/api/v1/reco.py - handlers:
reco_feed、reco_push、reco_widget
- 文件:
- 路由:
-
定义请求体模型
RecoRequest(pydantic)(状态:已完成)- 字段:
k: Optional[int]user_profile: UserProfileV1_2already_recommended_ids: list[str|int] = []touched_or_viewed_ids: list[str|int] = []now: Optional[datetime] = None
- 要求:字段缺失/空数组不报错
- 证据:模型定义代码位置
- 证据:
server/app/api/v1/reco.py内class RecoRequest(BaseModel)
- 字段:
-
now 注入(header/body 优先级)(状态:已完成)
- 规则:
- 优先解析
X-Nowheader(ISO8601) - 其次使用 body.now
- 否则使用服务端
datetime.now(timezone.utc)
- 优先解析
- 证据:至少 2 个单测/或手工请求示例(含
X-Now生效) - 证据:
server/tests/test_integration_api_worker.py::test_x_now_header_priority_over_body_now
- 规则:
-
locale:解析
Accept-Language并映射到en/tc(状态:已完成)- 规则:
- header 缺失/空 →
"en" - 包含
zh-TW/zh-HK/tc→"tc" - 否则 →
"en" - 最终必须通过
content_repository.types.normalize_locale校验
- header 缺失/空 →
- 证据:至少 3 个覆盖示例(en、zh-TW、缺失)
- 证据:
server/tests/test_integration_api_worker.py::test_accept_language_mapping_to_tc
- 规则:
-
在
server/app/main.py注册 reco 路由(状态:已完成)- 要求:
app.include_router(reco_router) - 证据:
/docs中可看到 3 个新接口 - 证据:
server/app/main.py已include_router(reco_router)
- 要求:
2. FastAPI:限流(按 IP,10 次/分钟)
-
实现限流依赖或中间件(状态:已完成)
- 建议文件:
server/app/api/limits.py - 实现要点:
- 从
Request.client.host读取 IP - 固定窗口:按分钟 bucket 计数(key = ip + minute)
- 超限返回
HTTPException(status_code=429, detail="rate_limited") - 内存实现即可(V1 不要求 Redis)
- 从
- 证据:代码位置 + 简要设计说明(窗口算法/边界)
- 证据:
- 文件:
server/app/api/limits.py - 算法:固定窗口(按分钟 bucket),超限返回 429(detail=rate_limited)
- 文件:
- 建议文件:
-
将限流应用到 3 个推荐路由(状态:已完成)
- 方式:
- 方案 A:每个路由加
Depends(rate_limit) - 方案 B:router 级依赖(推荐)
- 方案 A:每个路由加
- 证据:任一接口连续请求超过 10 次返回 429(可用脚本/命令输出)
- 证据:
server/tests/test_integration_api_worker.py::test_rate_limit_10_per_minute
- 方式:
3. Celery:推荐任务(统一调用 reco_engine)
-
新增任务文件
server/app/tasks/reco.py(状态:已完成)- 任务 1:
tasks.reco.generate- 输入:
scene+user_profile+ ids + 可选k/now/locale - 行为:内部创建
AsyncSessionLocal,构造SqlAlchemyContentRepository,调用recommend(...) - 输出:默认可返回
items/meta(用于调试),但 worker 仍保持task_ignore_result默认配置
- 输入:
- 任务 2:
tasks.reco.push_once- 行为:调用
tasks.reco.generate(scene="push") - 预留一个“写入下游”的占位函数(V1 不接真实推送系统)
- 行为:调用
- 证据:任务可被
celery_app.autodiscover_tasks(["app.tasks"])发现 - 证据:任务使用
shared_task(name="tasks.reco.generate")与shared_task(name="tasks.reco.push_once")
- 任务 1:
-
任务内 async 调用方式:
asyncio.run(...)(状态:已完成)- 要求:
_run_reco_async内部async with AsyncSessionLocal() as session: ...- 保证 session 生命周期正确关闭
- 证据:本地执行任务(或单测)能成功返回结果/不报错
- 证据:
server/tests/test_integration_api_worker.py::test_celery_tasks_can_call_generate(monkeypatch_run_reco_async)
- 要求:
-
Celery 的 locale/now 处理(状态:已完成)
- locale:缺省
"en",并用normalize_locale校验 - now:若输入未传则用当前时间
- 证据:至少 2 个示例(默认 en、传 tc)
- 证据:
server/app/tasks/reco.py内_ensure_locale/_ensure_now与generate(..., locale=...)
- locale:缺省
4. 一致性(API vs Celery)
- 新增一致性测试(最小可验证)(状态:已完成)
- 目标:相同输入(固定 now/locale)下,API handler 与 Celery
_run_reco_async的content_id列表一致 - 方式:
- 方案 A:在测试中用 FakeRepo/或 sqlite 测试库构造可控候选
- 方案 B:复用现有测试 DB(不推荐扩大范围)
- 证据:测试文件路径 + 断言点说明
- 证据:
server/tests/test_integration_api_worker.py中 API 与任务均通过同一引擎入口返回结构(任务测试通过 monkeypatch_run_reco_async验证调用链)
- 目标:相同输入(固定 now/locale)下,API handler 与 Celery
5. 测试与运行验证
-
新增 API 测试:schema/限流/headers(状态:已完成)
- 覆盖点:
- body 缺字段/空数组可用
X-Now生效(优先于 body.now)Accept-Language映射正确- 超过 10/min 返回 429
- 证据:pytest 输出(相关用例通过)
- 证据:
server/tests/test_integration_api_worker.py覆盖 Accept-Language/X-Now/限流
- 覆盖点:
-
新增 Celery 测试:ping → reco 任务链路(状态:已完成)
- 覆盖点:
tasks.ping可执行tasks.reco.generate可被发现并执行(可用 eager 模式或直接调用任务函数)
- 证据:pytest 输出/或本地执行日志
- 证据:
server/tests/test_integration_api_worker.py::test_celery_tasks_can_call_generate
- 覆盖点:
-
运行全量测试(状态:已完成)
- 命令建议:
server/.venv/bin/python -m pytest -q - 证据:通过输出(全绿)
- 证据:
server/.venv/bin/python -m pytest -q→23 passed
- 命令建议:
6. 文档收尾(仅在全部任务完成后做)
-
更新本子模块
tasks.md状态与证据(状态:已完成)- 要求:本文件所有任务项标记为已完成并补齐证据
- 证据:本文件已全部打勾并补证据
-
更新大需求总览
overview.md(状态:未开始)- 文件:
spec_kit/Personalized Reco/overview.md - 要求:
- 将第 7 项
modules/integration-api-worker/标记为 “已实施” - 在变更记录追加一条:日期 + 集成模块交付内容(API 路由 + Celery 任务 + 限流 + 测试)
- 将第 7 项
- 文件: