"""网关 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")