Files
rakuten-api/app/gateway/main.py
T
q792602257andClaude Opus 5 f6c3976c0a feat(gateway): 定时下派通道——周期性盘点账号订单并回写网关编目
新增 §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>
2026-08-16 23:35:51 +08:00

165 lines
6.1 KiB
Python

"""下单任务网关入口:FastAPI 应用创建与生命周期管理
部署在服务器侧,与抓取服务(:31107)、交易服务(:31108)并列。本地 worker 通过
出站长轮询从这里领任务;上游业务系统通过 POST /api/orders 提交下单意图。
为什么必须独立部署单元(不能塞进抓取服务):抓取无状态可多开,任务队列有状态,
多实例会抢同一批任务,同一账号的写操作必须串行(详见 docs/order-gateway.md §2)。
"""
from __future__ import annotations
import asyncio
import logging
from contextlib import asynccontextmanager
from fastapi import FastAPI
from app.gateway.api.routes.account import router as account_router
from app.gateway.api.routes.health import router as health_router
from app.gateway.api.routes.orders import router as orders_router
from app.gateway.api.routes.queries import router as queries_router
from app.gateway.collector import OrderDiscoveryCollector
from app.gateway.container import GatewayContainer
from app.gateway.db import GatewayDB
from app.gateway.query_queue import QueryQueue
from app.gateway.task_queue import TaskQueue
from app.shared.api import register_exception_handlers
from app.shared.config import get_settings
from app.shared.logging_setup import configure_logging
logger = logging.getLogger(__name__)
SWEEP_INTERVAL_SECONDS = 60
def build_container() -> GatewayContainer:
"""构建网关容器:DB + 下单任务队列 + 账号只读查询队列"""
settings = get_settings()
db = GatewayDB(settings.gateway_db_path_resolved)
task_queue = TaskQueue(
db,
lease_ttl_seconds=settings.lease_ttl_seconds,
worker_offline_alert_seconds=settings.worker_offline_alert_seconds,
)
query_queue = QueryQueue(
db,
lease_ttl_seconds=settings.query_lease_ttl_seconds,
query_ttl_seconds=settings.query_ttl_seconds,
max_attempts=settings.query_max_attempts,
retention_seconds=settings.query_retention_seconds,
)
collector = OrderDiscoveryCollector(
settings=settings, db=db, query_queue=query_queue
)
return GatewayContainer(
settings=settings, db=db, task_queue=task_queue,
query_queue=query_queue, collector=collector,
)
async def _sweep_loop(container: GatewayContainer) -> None:
"""常驻后台任务:每 60 秒扫一次过期租约
下单侧把 leased/running 推到 stale;查询侧重投超时的、置 expired 超 TTL 的、
清掉过保留期的结果。
lease 请求本身也会顺带扫一次,但 worker 退场后没人来 lease,必须靠这个
兜底——否则过期的任务永远停在 active 状态,健康检查看不到,stale 任务也
取不出来 reclaim。
"""
while True:
try:
await asyncio.sleep(SWEEP_INTERVAL_SECONDS)
swept = await container.task_queue.sweep()
if swept:
logger.info("sweep 把 %s 个过期任务置为 stale", swept)
query_swept = await container.query_queue.sweep()
if query_swept:
logger.info("sweep 处理了 %s 张超时/过期的查询单", query_swept)
except asyncio.CancelledError:
raise
except Exception: # noqa: BLE001
# 后台任务不能因为偶发错误退出,否则过期任务再也无人清理
logger.exception("sweep 后台任务出错,将继续重试")
@asynccontextmanager
async def lifespan(app: FastAPI):
"""应用生命周期:启动 DB、起 sweep 后台任务、关闭时反序释放"""
container = build_container()
app.state.container = container
configure_logging(container.settings)
logger.info(
"网关启动:%s:%s", container.settings.gateway_host, container.settings.gateway_port
)
logger.info("当前环境:%s", container.settings.app_env)
logger.info("DB 路径:%s", container.settings.gateway_db_path_resolved)
await container.db.start()
# 启动时先扫一次:上次进程退出时可能留有 leased/running 的过期任务
startup_swept = await container.task_queue.sweep()
if startup_swept:
logger.warning("启动时把 %s 个遗留过期任务置为 stale", startup_swept)
startup_queries = await container.query_queue.sweep()
if startup_queries:
logger.warning("启动时处理了 %s 张遗留的超时/过期查询单", startup_queries)
container.sweep_task = asyncio.create_task(
_sweep_loop(container), name="gateway-sweep"
)
# 定时下派通道(§12):被启用时 collector.run() 里自带常驻循环,禁用则直接返回
assert container.collector is not None
container.collector_task = asyncio.create_task(
container.collector.run(), name="gateway-collector"
)
try:
yield
finally:
if container.collector_task is not None:
container.collector_task.cancel()
try:
await container.collector_task
except asyncio.CancelledError:
pass
container.collector_task = None
if container.sweep_task is not None:
container.sweep_task.cancel()
try:
await container.sweep_task
except asyncio.CancelledError:
pass
container.sweep_task = None
await container.db.close()
def create_app() -> FastAPI:
"""创建 FastAPI 应用实例,注册路由和异常处理器"""
app = FastAPI(title="Rakuten Order Gateway", lifespan=lifespan)
app.include_router(health_router)
app.include_router(orders_router)
app.include_router(queries_router)
app.include_router(account_router)
register_exception_handlers(app)
return app
app = create_app()
if __name__ == "__main__":
import uvicorn
settings = get_settings()
configure_logging(settings)
uvicorn.run(
"app.gateway.main:app",
host=settings.gateway_host,
port=settings.gateway_port,
log_config=None,
timeout_keep_alive=120,
# 单进程:SQLite 单连接 + 全局并发度 1,多进程会抢同一个 DB
workers=1,
)