新增 §12 定时下派通道:网关自己定期派 order_list / order_detail 查询,把账号
真实订单沉淀进新表 account_orders(回写网关),上游可直接查 GET /api/account/orders
拿到账号里实际有哪些订单,不必自己记 order_number。
- collector.py: OrderDiscoveryCollector 常驻后台任务(与 sweep 并列)。每隔
account_discovery_interval_seconds 派 order_list,收割结果后把「编目里没有的
订单」upsert 进编目,再逐笔派 order_detail 沉淀 delivery_status/order_state。
状态无痕:不新增编排跟踪表,仅内存 _pending + _detail_dispatched_this_day;
detail query_id 按天分段(discover-detail-<order>-<日期>),同天幂等防重复派。
边界:collector 只消费 worker 结果的规范化字段,绝不解析 raw/raw_pages。
- db.py: account_orders 表 + AccountOrderRow + upsert/single/list/stale 访问;
upsert 以 order_number 为主键,重复采集只刷新,列表重扫不抹详情节点。
- 路由: GET /api/account/orders(列编目)、GET /api/account/orders/{n}(6005)、
POST /api/account/discovery/trigger(手动立即派一轮 list)。
- config: RAKUTEN_ACCOUNT_DISCOVERY_ENABLED / _INTERVAL_SECONDS / _MAX_PAGES /
RAKUTEN_ACCOUNT_DETAIL_REFRESH_SECONDS。
- 既有网关测试夹具统一关 discovery(会启动即派扫描污染查询队列语义),
新表测试在 tests/test_gateway_discovery.py(13 用例,真实 QueryQueue+DB 装配)。
472 tests passed.
Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
291 lines
11 KiB
Python
291 lines
11 KiB
Python
"""定时下派通道测试(docs/order-gateway.md §12)
|
|
|
|
collector 靠 QueryQueue + GatewayDB 真实装配驱动(与 test_gateway_queries.py 同一套
|
|
思路),用「真实队列 + 假 worker 回结果」来模拟完整闭环:
|
|
|
|
派 order_list → lease → worker 回带订单的 result → collector 收割 → 写进
|
|
account_orders 编目(回写网关)→ 派 order_detail → 回 detail → 更新配送阶段
|
|
|
|
跑测试时关掉 collector 的常驻循环,只驱动单次 tick()/各方法,避免引入 sleep。
|
|
"""
|
|
from __future__ import annotations
|
|
|
|
from pathlib import Path
|
|
|
|
import pytest
|
|
|
|
from app.gateway.collector import OrderDiscoveryCollector, _parse_list_result
|
|
from app.gateway.db import GatewayDB
|
|
from app.gateway.query_queue import QueryQueue
|
|
from app.shared.config import get_settings
|
|
|
|
TOKEN = get_settings().bearer_token
|
|
AUTH = {"Authorization": f"Bearer {TOKEN}"}
|
|
|
|
ORDER = {
|
|
"order_number": "306087-20260813-0863947697",
|
|
"order_date": "2026-08-13T09:47:28.000Z",
|
|
"shop_id": 306087,
|
|
"shop_name": "テスト店舗",
|
|
"items": [],
|
|
}
|
|
|
|
|
|
@pytest.fixture
|
|
async def collector(tmp_path: Path, monkeypatch):
|
|
"""真实装配的 (collector, db, queue) 三元组,collector 的常驻循环不启动"""
|
|
monkeypatch.setenv("RAKUTEN_ACCOUNT_DISCOVERY_ENABLED", "true")
|
|
get_settings.cache_clear()
|
|
db = GatewayDB(tmp_path / "c.db")
|
|
await db.start()
|
|
queue = QueryQueue(
|
|
db, lease_ttl_seconds=60, query_ttl_seconds=600,
|
|
max_attempts=3, retention_seconds=7 * 24 * 3600,
|
|
)
|
|
c = OrderDiscoveryCollector(
|
|
settings=get_settings(), db=db, query_queue=queue,
|
|
)
|
|
try:
|
|
yield c, db, queue
|
|
finally:
|
|
await db.close()
|
|
get_settings.cache_clear()
|
|
|
|
|
|
async def _complete_query(queue, query_id: str, result: dict, *, success: bool = True):
|
|
"""模拟 worker:lease 这张单并回一个结果"""
|
|
leased = await queue.lease(worker_id="w-fake", wait=0, site=None, max_wait=60)
|
|
assert leased.query_id == query_id
|
|
await queue.submit_result(
|
|
query_id, worker_id="w-fake", success=success, result=result,
|
|
error_code=None if success else 1234,
|
|
error_message="" if success else "boom",
|
|
)
|
|
|
|
|
|
# ---- 纯函数:从列表结果里抽规范化订单 ----
|
|
|
|
def test_parse_list_result_extracts_valid_rows():
|
|
result = {
|
|
"orders": [
|
|
{**ORDER},
|
|
{"order_number": "x", "order_date": "t"},
|
|
{}, # 缺 order_number,应被丢弃
|
|
"not-a-dict", # 非 dict,应被丢弃
|
|
]
|
|
}
|
|
rows = _parse_list_result(result)
|
|
assert len(rows) == 2
|
|
assert rows[0]["order_number"] == ORDER["order_number"]
|
|
|
|
|
|
def test_parse_list_result_empty_on_no_orders():
|
|
assert _parse_list_result(None) == []
|
|
assert _parse_list_result({}) == []
|
|
assert _parse_list_result({"orders": "oops"}) == []
|
|
|
|
|
|
# ---- 回写:采集到的订单写进编目 ----
|
|
|
|
async def test_ingest_list_result_writes_back_new_order(collector):
|
|
c, db, queue = collector
|
|
await c._queue.submit(query_id="q1", site="rakuten", kind="order_list", params={})
|
|
# 模拟 collector 已派这张单并在等结果(_collect_results 只处理 _pending 里的单)
|
|
c._pending["q1"] = "order_list"
|
|
await _complete_query(queue, "q1", {"orders": [ORDER], "window_fully_covered": True})
|
|
|
|
await c._collect_results() # 收割 q1 的结果,回写编目
|
|
|
|
row = await db.get_account_order(ORDER["order_number"])
|
|
assert row is not None
|
|
assert row.shop_name == "テスト店舗"
|
|
assert row.delivery_status is None # 列表结果不带配送阶段,回到编目后仍待详情
|
|
assert row.detail_fetched_at is None
|
|
|
|
|
|
async def test_reingest_same_order_no_duplicate(collector):
|
|
"""同一订单反复被采集,编目里只有一行(upsert 以 order_number 为主键)"""
|
|
c, db, _ = collector
|
|
for _ in range(2):
|
|
await c._ingest_list_result({"orders": [ORDER]}, query_id="q")
|
|
rows, total = await db.list_account_orders()
|
|
assert total == 1
|
|
assert rows[0].order_number == ORDER["order_number"]
|
|
|
|
|
|
async def test_list_reingest_preserves_detail_data(collector):
|
|
"""订单已取过详情后,再来一轮 list 扫描不得把配送阶段/详情时间抹掉"""
|
|
c, db, _ = collector
|
|
# 先经 detail 沉淀配送阶段
|
|
await c._ingest_list_result({"orders": [ORDER]}, query_id="q")
|
|
await c._ingest_detail_result(
|
|
{
|
|
"order_number": ORDER["order_number"], "found": True,
|
|
"delivery_status": "CHECKING_ORDER", "order_state": "created",
|
|
}
|
|
)
|
|
before = await db.get_account_order(ORDER["order_number"])
|
|
assert before.detail_fetched_at is not None
|
|
|
|
# 再来一轮列表扫描
|
|
await c._ingest_list_result({"orders": [ORDER]}, query_id="q2")
|
|
|
|
after = await db.get_account_order(ORDER["order_number"])
|
|
assert after.delivery_status == "CHECKING_ORDER"
|
|
assert after.detail_fetched_at == before.detail_fetched_at
|
|
|
|
|
|
async def test_found_false_detail_does_not_wipe_status(collector):
|
|
"""详情 found=False(订单还没反映出来)时,不覆盖已沉淀的配送阶段"""
|
|
c, db, _ = collector
|
|
await c._ingest_list_result({"orders": [ORDER]}, query_id="q")
|
|
await c._ingest_detail_result(
|
|
{
|
|
"order_number": ORDER["order_number"], "found": True,
|
|
"delivery_status": "CHECKING_ORDER", "order_state": "created",
|
|
}
|
|
)
|
|
|
|
# 一瞬抖动:这单暂时查不到(found=False,结果里配送字段为 null)
|
|
await c._ingest_detail_result(
|
|
{"order_number": ORDER["order_number"], "found": False,
|
|
"delivery_status": None, "order_state": None}
|
|
)
|
|
|
|
after = await db.get_account_order(ORDER["order_number"])
|
|
assert after.delivery_status == "CHECKING_ORDER" # 没被抹掉
|
|
assert after.detail_fetched_at is not None
|
|
|
|
|
|
# ---- 完整闭环:派 list → 消化 → 派 detail → 更新 ----
|
|
|
|
async def test_full_discovery_loop_populates_catalog(collector):
|
|
"""一次 tick 派 list,回来后被写进编目,且触发 detail 下派与更新"""
|
|
c, db, queue = collector
|
|
|
|
# 第一次 tick:派出一张 order_list 扫描
|
|
await c._dispatch_list_if_due()
|
|
assert c._pending # 有一张在飞
|
|
|
|
# writer worker 领走并回一个有订单的列表
|
|
(qid,) = c._pending
|
|
await _complete_query(queue, qid, {"orders": [ORDER], "window_fully_covered": True})
|
|
|
|
# 收割 + 派 stale detail
|
|
await c.tick()
|
|
|
|
# 编目里有了订单
|
|
row = await db.get_account_order(ORDER["order_number"])
|
|
assert row is not None and row.detail_fetched_at is None
|
|
|
|
# detail 下派已产生一张 order_detail 单
|
|
pending_detail = [q for q, k in c._pending.items() if k == "order_detail"]
|
|
assert pending_detail
|
|
|
|
# worker 回 detail 结果
|
|
dqid = pending_detail[0]
|
|
await _complete_query(
|
|
queue, dqid,
|
|
{
|
|
"order_number": ORDER["order_number"],
|
|
"found": True,
|
|
"delivery_status": "CHECKING_ORDER",
|
|
"order_state": "created",
|
|
},
|
|
)
|
|
await c._collect_results()
|
|
|
|
after = await db.get_account_order(ORDER["order_number"])
|
|
assert after.delivery_status == "CHECKING_ORDER"
|
|
assert after.detail_fetched_at is not None
|
|
|
|
|
|
async def test_detail_dispatch_is_same_day_idempotent(collector):
|
|
"""同一天内轮询到同一订单,detail 单只派一张(幂等 submit 命中既有单)"""
|
|
c, db, _ = collector
|
|
await c._ingest_list_result({"orders": [ORDER]}, query_id="q")
|
|
|
|
await c._dispatch_one_detail(ORDER["order_number"])
|
|
first = set(c._pending)
|
|
await c._dispatch_one_detail(ORDER["order_number"])
|
|
second = set(c._pending)
|
|
|
|
assert first == second
|
|
assert len(first) == 1
|
|
|
|
|
|
async def test_stale_detail_not_redispatched_within_day(collector):
|
|
"""详情 found=False(订单还没反映)时,当天不再重复派同一订单的详情"""
|
|
c, db, _ = collector
|
|
await c._ingest_list_result({"orders": [ORDER]}, query_id="q")
|
|
|
|
# 第一轮:派详情 → 回 found=False → 订单仍算「未取过详情」
|
|
await c._dispatch_stale_details()
|
|
(dqid,) = [q for q, k in c._pending.items() if k == "order_detail"]
|
|
await _complete_query(
|
|
c._queue, dqid,
|
|
{"order_number": ORDER["order_number"], "found": False},
|
|
)
|
|
await c._collect_results()
|
|
assert (await db.get_account_order(ORDER["order_number"])).detail_fetched_at is None
|
|
|
|
# 第二轮:不应再为同一订单派详情(当天已派过)
|
|
before = set(c._pending)
|
|
await c._dispatch_stale_details()
|
|
assert set(c._pending) == before
|
|
|
|
|
|
async def test_no_list_redispatch_while_pending(collector):
|
|
"""有在飞的 list 单时,即使间隔到了也不重复派"""
|
|
c, _, _ = collector
|
|
await c._dispatch_list_if_due()
|
|
assert len(c._pending) == 1
|
|
# 模拟时间流逝:把 last_list_submitted 拨回到很久以前
|
|
c._last_list_submitted = None
|
|
await c._dispatch_list_if_due()
|
|
# 幂等?list 用的是自动生成的 q- 前缀 id,重新 submit 会再新建——
|
|
# 但「有在飞 list」会挡住。此时 pending 里那张 q- 单还没到终态,不应再派。
|
|
assert len([k for k in c._pending.values() if k == "order_list"]) == 1
|
|
|
|
|
|
# ---- HTTP 层:编目读取与手动触发 ----
|
|
|
|
@pytest.fixture
|
|
def gateway_client(tmp_path: Path, monkeypatch):
|
|
monkeypatch.setenv("RAKUTEN_GATEWAY_DB_PATH", str(tmp_path / "gw.db"))
|
|
monkeypatch.setenv("RAKUTEN_ACCOUNT_DISCOVERY_ENABLED", "false")
|
|
get_settings.cache_clear()
|
|
try:
|
|
from app.gateway.main import create_app
|
|
from fastapi.testclient import TestClient
|
|
|
|
app = create_app()
|
|
with TestClient(app) as client:
|
|
yield client
|
|
finally:
|
|
get_settings.cache_clear()
|
|
|
|
|
|
def test_get_cataloged_order_not_found_6005(gateway_client):
|
|
response = gateway_client.get("/api/account/orders/nope", headers=AUTH)
|
|
assert response.status_code == 404
|
|
assert response.json()["code"] == 6005
|
|
|
|
|
|
def test_trigger_discovery_dispatch_list_sweep(gateway_client):
|
|
"""手动触发应派一张 order_list 单(走真实队列)"""
|
|
# 触发里 collector 真实存在;手动派单不应因 discovery disabled 而失效——
|
|
# 但这个夹具关了 discovery,trigger 只是 submit 一张单,仍应正常。
|
|
body = gateway_client.post("/api/account/discovery/trigger", headers=AUTH).json()
|
|
assert body["success"] is True
|
|
assert body["data"]["query_id"].startswith("q-")
|
|
# 这张单真的进了查询队列
|
|
detail = gateway_client.get(f"/api/account/queries/{body['data']['query_id']}", headers=AUTH).json()
|
|
assert detail["data"]["kind"] == "order_list"
|
|
|
|
|
|
def test_catalog_endpoints_reject_missing_token(gateway_client):
|
|
assert gateway_client.get("/api/account/orders").status_code == 401
|
|
assert gateway_client.get("/api/account/orders/x").status_code == 401
|
|
assert gateway_client.post("/api/account/discovery/trigger").status_code == 401
|