按 docs/order-gateway.md 落地:第三个部署单元 app.gateway(:31109)承担任务队列 + 状态镜像;本地 worker 在 app.trading.worker 内,按 RAKUTEN_ORDER_GATEWAY_URL 决定是否启动。规格 §5 最关键约束已守:租约过期绝不自动重投,恢复只能 reclaim, worker 收到 lease_count>1 时先核对站点订单。 站点交互(加购/下单/付款/订单列表反查)按规格 §10 留接口缝,site_interact.py 全部 NotImplementedError,verify.py 恒返回 unknown——等真实账号实测后再填, 不写猜测的提交逻辑。 310 个测试全绿,覆盖规格 §9 验收清单 12 条;架构测试守住三方互不 import。 Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>
369 lines
12 KiB
Python
369 lines
12 KiB
Python
"""网关 DB 与任务队列层测试
|
|
|
|
直接构造 GatewayDB + TaskQueue,绕开 HTTP 层,覆盖规格 §9 验收清单里能在队列层
|
|
独立验证的条目:幂等提交、全局并发度 1、租约过期变 stale(且不重投)、
|
|
reclaim 恢复、(task_id, state) 幂等 report。
|
|
|
|
长轮询与并发 worker 抢任务的时序在 test_gateway_leasing.py 单独覆盖。
|
|
"""
|
|
from __future__ import annotations
|
|
|
|
from pathlib import Path
|
|
|
|
import pytest
|
|
|
|
from app.gateway.db import GatewayDB
|
|
from app.gateway.task_queue import TaskQueue
|
|
from app.shared.task_state import OrderState, TaskStatus
|
|
|
|
|
|
@pytest.fixture
|
|
async def queue(tmp_path: Path) -> TaskQueue:
|
|
db = GatewayDB(tmp_path / "gw.db")
|
|
await db.start()
|
|
q = TaskQueue(db, lease_ttl_seconds=300, worker_offline_alert_seconds=300)
|
|
yield q
|
|
await db.close()
|
|
|
|
|
|
# ---- 幂等提交(规格 §9 第 1 条)----
|
|
|
|
|
|
async def test_submit_with_same_task_id_is_idempotent(queue: TaskQueue):
|
|
"""同 task_id 提交两次只产生一个任务,第二次 created=False"""
|
|
r1 = await queue.submit(task_id="t1", site="rakuten", intent={"k": "v"})
|
|
r2 = await queue.submit(task_id="t1", site="rakuten", intent={"k": "v"})
|
|
|
|
assert r1.task_id == "t1"
|
|
assert r1.created is True
|
|
assert r2.task_id == "t1"
|
|
assert r2.created is False
|
|
assert r1.status == r2.status == TaskStatus.QUEUED
|
|
|
|
|
|
async def test_submit_generates_task_id_when_missing(queue: TaskQueue):
|
|
r = await queue.submit(task_id=None, site="rakuten", intent={})
|
|
assert r.task_id.startswith("po-")
|
|
assert r.created is True
|
|
|
|
|
|
# ---- 全局并发度 1(规格 §9 第 2、3 条)----
|
|
|
|
|
|
async def test_lease_returns_none_when_active_task_exists(queue: TaskQueue):
|
|
"""已有 leased/running 时,lease 立即返回 None"""
|
|
await queue.submit(task_id="t1", site="rakuten", intent={})
|
|
first = await queue.lease(worker_id="w1", wait=0, site=None, max_wait=60)
|
|
assert first is not None and first.task_id == "t1"
|
|
|
|
second = await queue.lease(worker_id="w2", wait=0, site=None, max_wait=60)
|
|
assert second is None
|
|
|
|
|
|
async def test_lease_picks_queued_in_fifo_order(queue: TaskQueue):
|
|
await queue.submit(task_id="t1", site="rakuten", intent={})
|
|
await queue.submit(task_id="t2", site="rakuten", intent={})
|
|
|
|
first = await queue.lease(worker_id="w1", wait=0, site=None, max_wait=60)
|
|
assert first.task_id == "t1"
|
|
|
|
# 终结第一个后才能领第二个
|
|
await queue.report(
|
|
"t1",
|
|
worker_id="w1",
|
|
state=OrderState.CREATED,
|
|
payable_yen=None,
|
|
pay_deadline=None,
|
|
site_order_id=None,
|
|
evidence_ref=None,
|
|
detail="",
|
|
terminal=True,
|
|
terminal_status=TaskStatus.SUCCEEDED,
|
|
)
|
|
|
|
second = await queue.lease(worker_id="w1", wait=0, site=None, max_wait=60)
|
|
assert second is not None and second.task_id == "t2"
|
|
|
|
|
|
# ---- 站点过滤 ----
|
|
|
|
|
|
async def test_lease_with_site_filter(queue: TaskQueue):
|
|
await queue.submit(task_id="t1", site="rakuten", intent={})
|
|
await queue.submit(task_id="t2", site="rakuma", intent={})
|
|
|
|
leased = await queue.lease(worker_id="w1", wait=0, site="rakuma", max_wait=60)
|
|
assert leased is not None and leased.task_id == "t2" and leased.site == "rakuma"
|
|
|
|
|
|
# ---- 续租 ----
|
|
|
|
|
|
async def test_renew_requires_lease_owner(queue: TaskQueue):
|
|
await queue.submit(task_id="t1", site="rakuten", intent={})
|
|
await queue.lease(worker_id="w1", wait=0, site=None, max_wait=60)
|
|
|
|
from app.shared.errors import LeaseInvalidError
|
|
|
|
with pytest.raises(LeaseInvalidError):
|
|
await queue.renew("t1", "w2")
|
|
|
|
|
|
async def test_renew_extends_lease(queue: TaskQueue):
|
|
await queue.submit(task_id="t1", site="rakuten", intent={})
|
|
await queue.lease(worker_id="w1", wait=0, site=None, max_wait=60)
|
|
renewed = await queue.renew("t1", "w1")
|
|
assert renewed.task_id == "t1"
|
|
assert renewed.lease_count == 1
|
|
|
|
|
|
# ---- report 幂等(规格 §9 第 5 条)----
|
|
|
|
|
|
async def test_report_same_state_is_idempotent(queue: TaskQueue):
|
|
"""同一 (task_id, state) 重复上报不产生第二条记录"""
|
|
await queue.submit(task_id="t1", site="rakuten", intent={})
|
|
await queue.lease(worker_id="w1", wait=0, site=None, max_wait=60)
|
|
|
|
kwargs = dict(
|
|
worker_id="w1",
|
|
state=OrderState.IN_CART,
|
|
payable_yen=100,
|
|
pay_deadline=None,
|
|
site_order_id=None,
|
|
evidence_ref="t1/01-cart-add",
|
|
detail="已加购",
|
|
terminal=False,
|
|
terminal_status=None,
|
|
)
|
|
r1 = await queue.report("t1", **kwargs)
|
|
r2 = await queue.report("t1", **kwargs)
|
|
|
|
assert r1.recorded is True # 新写入
|
|
assert r2.recorded is False # 幂等覆盖
|
|
|
|
detail = await queue.get_task_detail("t1")
|
|
assert len(detail.reports) == 1
|
|
|
|
|
|
async def test_first_report_moves_task_to_running(queue: TaskQueue):
|
|
await queue.submit(task_id="t1", site="rakuten", intent={})
|
|
await queue.lease(worker_id="w1", wait=0, site=None, max_wait=60)
|
|
|
|
await queue.report(
|
|
"t1",
|
|
worker_id="w1",
|
|
state=OrderState.IN_CART,
|
|
payable_yen=None,
|
|
pay_deadline=None,
|
|
site_order_id=None,
|
|
evidence_ref=None,
|
|
detail="",
|
|
terminal=False,
|
|
terminal_status=None,
|
|
)
|
|
detail = await queue.get_task_detail("t1")
|
|
assert detail.status == TaskStatus.RUNNING
|
|
|
|
|
|
async def test_terminal_report_releases_lease(queue: TaskQueue):
|
|
await queue.submit(task_id="t1", site="rakuten", intent={})
|
|
await queue.lease(worker_id="w1", wait=0, site=None, max_wait=60)
|
|
|
|
await queue.report(
|
|
"t1",
|
|
worker_id="w1",
|
|
state=OrderState.PAID,
|
|
payable_yen=9800,
|
|
pay_deadline=None,
|
|
site_order_id="ord-1",
|
|
evidence_ref="t1/05-payment",
|
|
detail="付款完成",
|
|
terminal=True,
|
|
terminal_status=TaskStatus.SUCCEEDED,
|
|
)
|
|
detail = await queue.get_task_detail("t1")
|
|
assert detail.status == TaskStatus.SUCCEEDED
|
|
assert detail.lease_owner is None
|
|
assert detail.lease_expires_at is None
|
|
assert detail.latest_state == OrderState.PAID
|
|
|
|
|
|
async def test_terminal_without_explicit_status_infers_from_state(queue: TaskQueue):
|
|
"""terminal=true 但不传 terminal_status 时,按 state 推断:paid → succeeded"""
|
|
await queue.submit(task_id="t1", site="rakuten", intent={})
|
|
await queue.lease(worker_id="w1", wait=0, site=None, max_wait=60)
|
|
|
|
await queue.report(
|
|
"t1",
|
|
worker_id="w1",
|
|
state=OrderState.PAID,
|
|
payable_yen=None,
|
|
pay_deadline=None,
|
|
site_order_id=None,
|
|
evidence_ref=None,
|
|
detail="",
|
|
terminal=True,
|
|
terminal_status=None,
|
|
)
|
|
assert (await queue.get_task_detail("t1")).status == TaskStatus.SUCCEEDED
|
|
|
|
|
|
async def test_terminal_with_awaiting_payment_state_becomes_needs_human(queue: TaskQueue):
|
|
"""3DS 等人工环节 → worker 报 awaiting_payment + terminal → 任务转 needs_human"""
|
|
await queue.submit(task_id="t1", site="rakuten", intent={})
|
|
await queue.lease(worker_id="w1", wait=0, site=None, max_wait=60)
|
|
|
|
await queue.report(
|
|
"t1",
|
|
worker_id="w1",
|
|
state=OrderState.AWAITING_PAYMENT,
|
|
payable_yen=None,
|
|
pay_deadline=None,
|
|
site_order_id=None,
|
|
evidence_ref=None,
|
|
detail="检测到 3DS 验证",
|
|
terminal=True,
|
|
terminal_status=None,
|
|
)
|
|
assert (await queue.get_task_detail("t1")).status == TaskStatus.NEEDS_HUMAN
|
|
|
|
|
|
# ---- 租约过期与绝不自动重投(规格 §5、§9 第 6 条)----
|
|
|
|
|
|
async def test_expired_lease_becomes_stale_and_is_not_re_leasable(tmp_path: Path):
|
|
"""租约过期 → stale → 普通 lease 取不到(绝不自动重投)"""
|
|
db = GatewayDB(tmp_path / "gw.db")
|
|
await db.start()
|
|
try:
|
|
q = TaskQueue(db, lease_ttl_seconds=0, worker_offline_alert_seconds=300)
|
|
await q.submit(task_id="t1", site="rakuten", intent={})
|
|
leased = await q.lease(worker_id="w1", wait=0, site=None, max_wait=60)
|
|
assert leased is not None
|
|
|
|
# TTL=0:立即过期。下一次 lease 内部会先扫到 stale
|
|
empty = await q.lease(worker_id="w2", wait=0, site=None, max_wait=60)
|
|
assert empty is None
|
|
|
|
detail = await q.get_task_detail("t1")
|
|
assert detail.status == TaskStatus.STALE
|
|
finally:
|
|
await db.close()
|
|
|
|
|
|
async def test_sweep_returns_count_of_expired_tasks(queue: TaskQueue):
|
|
await queue.submit(task_id="t1", site="rakuten", intent={})
|
|
# 直接造一个过期任务:lease TTL 设为 0
|
|
queue._lease_ttl = 0 # type: ignore[attr-defined]
|
|
await queue.lease(worker_id="w1", wait=0, site=None, max_wait=60)
|
|
|
|
swept = await queue.sweep()
|
|
assert swept == 1
|
|
assert (await queue.get_task_detail("t1")).status == TaskStatus.STALE
|
|
|
|
|
|
# ---- reclaim(规格 §5、§9 第 7 条)----
|
|
|
|
|
|
async def test_reclaim_leases_stale_task_with_count_increment(tmp_path: Path):
|
|
"""reclaim 领回 stale 任务,lease_count > 1,带 known_state"""
|
|
db = GatewayDB(tmp_path / "gw.db")
|
|
await db.start()
|
|
try:
|
|
q = TaskQueue(db, lease_ttl_seconds=0, worker_offline_alert_seconds=300)
|
|
await q.submit(task_id="t1", site="rakuten", intent={})
|
|
|
|
first = await q.lease(worker_id="w1", wait=0, site=None, max_wait=60)
|
|
assert first is not None and first.lease_count == 1
|
|
|
|
# 在过期之前先 report 一次,让 known_state 有值
|
|
await q.report(
|
|
"t1",
|
|
worker_id="w1",
|
|
state=OrderState.IN_CART,
|
|
payable_yen=None,
|
|
pay_deadline=None,
|
|
site_order_id=None,
|
|
evidence_ref=None,
|
|
detail="",
|
|
terminal=False,
|
|
terminal_status=None,
|
|
)
|
|
# 走一次 lease 触发 sweep,把过期任务置 stale
|
|
await q.lease(worker_id="w2", wait=0, site=None, max_wait=60)
|
|
assert (await q.get_task_detail("t1")).status == TaskStatus.STALE
|
|
|
|
reclaimed = await q.reclaim("t1", "w2")
|
|
assert reclaimed.task_id == "t1"
|
|
assert reclaimed.lease_count == 2 # > 1:worker 必须先核对
|
|
assert reclaimed.known_state == OrderState.IN_CART
|
|
|
|
# reclaim 后任务状态变 leased(被新 worker 持有)
|
|
detail = await q.get_task_detail("t1")
|
|
assert detail.status == TaskStatus.LEASED
|
|
assert detail.lease_owner == "w2"
|
|
finally:
|
|
await db.close()
|
|
|
|
|
|
async def test_reclaim_rejects_non_stale_task(queue: TaskQueue):
|
|
"""只能 reclaim stale 任务;其他状态报 6003"""
|
|
from app.shared.errors import InvalidTaskStateError
|
|
|
|
await queue.submit(task_id="t1", site="rakuten", intent={})
|
|
|
|
with pytest.raises(InvalidTaskStateError):
|
|
await queue.reclaim("t1", "w1")
|
|
|
|
|
|
async def test_normal_lease_does_not_pick_stale(tmp_path: Path):
|
|
"""普通 lease 不会取到 stale 任务(绝不自动重投)"""
|
|
db = GatewayDB(tmp_path / "gw.db")
|
|
await db.start()
|
|
try:
|
|
q = TaskQueue(db, lease_ttl_seconds=0, worker_offline_alert_seconds=300)
|
|
await q.submit(task_id="stale-task", site="rakuten", intent={})
|
|
await q.lease(worker_id="w1", wait=0, site=None, max_wait=60)
|
|
# 触发 sweep
|
|
await q.lease(worker_id="w2", wait=0, site=None, max_wait=60)
|
|
assert (await q.get_task_detail("stale-task")).status == TaskStatus.STALE
|
|
|
|
# 再入队一个新任务
|
|
await q.submit(task_id="fresh", site="rakuten", intent={})
|
|
picked = await q.lease(worker_id="w3", wait=0, site=None, max_wait=60)
|
|
assert picked is not None and picked.task_id == "fresh"
|
|
finally:
|
|
await db.close()
|
|
|
|
|
|
# ---- 错误映射 ----
|
|
|
|
|
|
async def test_report_wrong_worker_raises(queue: TaskQueue):
|
|
from app.shared.errors import LeaseInvalidError
|
|
|
|
await queue.submit(task_id="t1", site="rakuten", intent={})
|
|
await queue.lease(worker_id="w1", wait=0, site=None, max_wait=60)
|
|
|
|
with pytest.raises(LeaseInvalidError):
|
|
await queue.report(
|
|
"t1",
|
|
worker_id="w2",
|
|
state=OrderState.IN_CART,
|
|
payable_yen=None,
|
|
pay_deadline=None,
|
|
site_order_id=None,
|
|
evidence_ref=None,
|
|
detail="",
|
|
terminal=False,
|
|
terminal_status=None,
|
|
)
|
|
|
|
|
|
async def test_get_unknown_task_raises(queue: TaskQueue):
|
|
from app.shared.errors import TaskNotFoundError
|
|
|
|
with pytest.raises(TaskNotFoundError):
|
|
await queue.get_task_detail("does-not-exist")
|