- POST /api/orders 新增可选 callback_url(仅 http/https,其余 422) - 仅终结类事件各通知一次:terminal report 推入终态(succeeded/failed/ needs_human)、租约过期被 sweep 置 stale;中间态与终结后的监控上报不通知 - 幂等重发不更新既有任务的回调地址;投递 best-effort 单次尝试,失败只记日志 - CallbackNotifier 发后不管(create_task + 在途任务强引用),关停等在途发完; 生产路径按回调地址逐次构造客户端,满足统一出站代理策略(test_proxy.py) - tasks 表加 callback_url 列,GatewayDB.start() 内置迁移兼容既有库 - 新增 RAKUTEN_CALLBACK_TIMEOUT_SECONDS(默认 10);文档补 §4.8; openapi.json 重新导出(gitignore 未跟踪);新增 17 条测试,全量 492 通过
174 lines
6.6 KiB
Python
174 lines
6.6 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.callback import CallbackNotifier
|
|
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)
|
|
# 终结类事件回调(§4.8):best-effort 单次投递,失败只记日志
|
|
notifier = CallbackNotifier(
|
|
settings=settings, timeout_seconds=settings.callback_timeout_seconds
|
|
)
|
|
task_queue = TaskQueue(
|
|
db,
|
|
lease_ttl_seconds=settings.lease_ttl_seconds,
|
|
worker_offline_alert_seconds=settings.worker_offline_alert_seconds,
|
|
notifier=notifier,
|
|
)
|
|
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, notifier=notifier,
|
|
)
|
|
|
|
|
|
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
|
|
if container.notifier is not None:
|
|
# 等在途回调发完(各次发送有超时兜底),避免关停时静默丢通知
|
|
await container.notifier.aclose()
|
|
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,
|
|
)
|