实现下单任务网关与本地 worker
按 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>
This commit is contained in:
@@ -0,0 +1,12 @@
|
||||
"""下单任务网关:第三个部署单元
|
||||
|
||||
抓取服务(app.scraping)匿名、无状态、可多开;交易服务(app.trading)带账号、
|
||||
有状态、单实例;本网关是任务队列与状态镜像,部署在服务器侧,**独立于两者**:
|
||||
|
||||
- 与抓取服务分进程的原因:抓取实例可以多开,而任务队列有状态,多实例会抢同一批
|
||||
任务(同一账号的写操作必须串行)。
|
||||
- 与交易服务分进程的原因:交易在本地(NAT 后无公网入口),网关在服务器,本地
|
||||
通过出站长轮询从这里领任务。
|
||||
|
||||
完整规格见 docs/order-gateway.md。
|
||||
"""
|
||||
@@ -0,0 +1 @@
|
||||
"""网关 API 包"""
|
||||
@@ -0,0 +1 @@
|
||||
"""网关路由集合"""
|
||||
@@ -0,0 +1,23 @@
|
||||
"""网关健康检查路由"""
|
||||
from fastapi import APIRouter, Depends
|
||||
|
||||
from app.gateway.container import GatewayContainer
|
||||
from app.gateway.models import GatewayHealthData
|
||||
from app.shared.api import ApiResponse, get_container
|
||||
|
||||
router = APIRouter(tags=["health"])
|
||||
|
||||
|
||||
@router.get("/health", response_model=ApiResponse[GatewayHealthData])
|
||||
async def health(
|
||||
container: GatewayContainer = Depends(get_container),
|
||||
) -> ApiResponse[GatewayHealthData]:
|
||||
"""网关健康状态
|
||||
|
||||
只做规格 §4.7 的两条兜底告警:worker 失联、任务长时间无人领。
|
||||
付款期限监控不在网关——本地机 7×24 在线,那套逻辑放本地。
|
||||
"""
|
||||
data = await container.task_queue.health_snapshot()
|
||||
return ApiResponse[GatewayHealthData](
|
||||
success=True, msg="success", data=data, code=0
|
||||
)
|
||||
@@ -0,0 +1,181 @@
|
||||
"""下单任务网关路由
|
||||
|
||||
承担规格 §4 的全部接口:
|
||||
|
||||
- POST /api/orders 上游提交下单意图(幂等)
|
||||
- GET /api/orders/lease 本地长轮询领取
|
||||
- POST /api/orders/{id}/renew 续租
|
||||
- POST /api/orders/{id}/report 本地回报状态
|
||||
- POST /api/orders/{id}/reclaim 把 stale 任务重新租给 worker(绝不自动重投)
|
||||
- GET /api/orders/{id} 任务详情 + 完整状态历史
|
||||
- GET /api/orders 任务列表(运维与上游对账)
|
||||
|
||||
handler 风格与 app/scraping/api/routes/scrape.py 保持一致:不在 handler 里写
|
||||
try/except,业务异常通过 AppError 自动套进 register_exception_handlers。
|
||||
"""
|
||||
from __future__ import annotations
|
||||
|
||||
from fastapi import APIRouter, Depends, Query
|
||||
|
||||
from app.gateway.container import GatewayContainer
|
||||
from app.gateway.models import (
|
||||
LeaseData,
|
||||
RenewData,
|
||||
ReclaimData,
|
||||
ReportData,
|
||||
SubmitOrderData,
|
||||
SubmitOrderRequest,
|
||||
TaskDetail,
|
||||
TaskListData,
|
||||
RenewRequest,
|
||||
ReclaimRequest,
|
||||
ReportRequest,
|
||||
)
|
||||
from app.shared.api import ApiResponse, get_container, require_bearer_token
|
||||
|
||||
router = APIRouter(prefix="/api/orders", tags=["orders"])
|
||||
|
||||
|
||||
@router.post(
|
||||
"",
|
||||
response_model=ApiResponse[SubmitOrderData],
|
||||
dependencies=[Depends(require_bearer_token)],
|
||||
)
|
||||
async def submit_order(
|
||||
payload: SubmitOrderRequest,
|
||||
container: GatewayContainer = Depends(get_container),
|
||||
) -> ApiResponse[SubmitOrderData]:
|
||||
"""上游提交下单意图
|
||||
|
||||
重复提交同一 task_id 不新建任务,返回既有任务且 created=false——上游重发
|
||||
不会变成两单。
|
||||
"""
|
||||
data = await container.task_queue.submit(
|
||||
task_id=payload.task_id, site=payload.site, intent=payload.intent
|
||||
)
|
||||
return ApiResponse[SubmitOrderData](
|
||||
success=True, msg="success", data=data, code=0
|
||||
)
|
||||
|
||||
|
||||
@router.get(
|
||||
"/lease",
|
||||
response_model=ApiResponse[LeaseData | None],
|
||||
dependencies=[Depends(require_bearer_token)],
|
||||
)
|
||||
async def lease_order(
|
||||
worker_id: str = Query(...),
|
||||
wait: int = Query(default=30, ge=0, le=300),
|
||||
site: str | None = Query(default=None),
|
||||
container: GatewayContainer = Depends(get_container),
|
||||
) -> ApiResponse[LeaseData | None]:
|
||||
"""本地长轮询领取
|
||||
|
||||
无可领任务时挂起到 wait 秒后返回 data: null,HTTP 仍 200。已有 leased/running
|
||||
任务时立即返回空(全局并发度 1)。
|
||||
"""
|
||||
data = await container.task_queue.lease(
|
||||
worker_id=worker_id,
|
||||
wait=wait,
|
||||
site=site,
|
||||
max_wait=container.settings.lease_max_wait_seconds,
|
||||
)
|
||||
return ApiResponse[LeaseData | None](
|
||||
success=True, msg="success", data=data, code=0
|
||||
)
|
||||
|
||||
|
||||
@router.post(
|
||||
"/{task_id}/renew",
|
||||
response_model=ApiResponse[RenewData],
|
||||
dependencies=[Depends(require_bearer_token)],
|
||||
)
|
||||
async def renew_order(
|
||||
task_id: str,
|
||||
payload: RenewRequest,
|
||||
container: GatewayContainer = Depends(get_container),
|
||||
) -> ApiResponse[RenewData]:
|
||||
"""续租。worker 在长任务中每 60 秒调一次,避免租约 TTL 误判"""
|
||||
data = await container.task_queue.renew(task_id, payload.worker_id)
|
||||
return ApiResponse[RenewData](success=True, msg="success", data=data, code=0)
|
||||
|
||||
|
||||
@router.post(
|
||||
"/{task_id}/report",
|
||||
response_model=ApiResponse[ReportData],
|
||||
dependencies=[Depends(require_bearer_token)],
|
||||
)
|
||||
async def report_order(
|
||||
task_id: str,
|
||||
payload: ReportRequest,
|
||||
container: GatewayContainer = Depends(get_container),
|
||||
) -> ApiResponse[ReportData]:
|
||||
"""本地回报订单状态
|
||||
|
||||
同一 (task_id, state) 重复上报幂等。terminal=true 时释放租约并推到终态。
|
||||
"""
|
||||
data = await container.task_queue.report(
|
||||
task_id,
|
||||
worker_id=payload.worker_id,
|
||||
state=payload.state,
|
||||
payable_yen=payload.payable_yen,
|
||||
pay_deadline=payload.pay_deadline,
|
||||
site_order_id=payload.site_order_id,
|
||||
evidence_ref=payload.evidence_ref,
|
||||
detail=payload.detail,
|
||||
terminal=payload.terminal,
|
||||
terminal_status=payload.terminal_status.value if payload.terminal_status else None,
|
||||
)
|
||||
return ApiResponse[ReportData](success=True, msg="success", data=data, code=0)
|
||||
|
||||
|
||||
@router.post(
|
||||
"/{task_id}/reclaim",
|
||||
response_model=ApiResponse[ReclaimData],
|
||||
dependencies=[Depends(require_bearer_token)],
|
||||
)
|
||||
async def reclaim_order(
|
||||
task_id: str,
|
||||
payload: ReclaimRequest,
|
||||
container: GatewayContainer = Depends(get_container),
|
||||
) -> ApiResponse[ReclaimData]:
|
||||
"""把 stale 任务重新租给 worker(恢复领取)
|
||||
|
||||
返回体的 lease_count 必然 > 1,worker 必须先核对站点订单列表再决定是否
|
||||
继续执行——核对不出结论时报 needs_human,绝不重新提交。
|
||||
"""
|
||||
data = await container.task_queue.reclaim(task_id, payload.worker_id)
|
||||
return ApiResponse[ReclaimData](success=True, msg="success", data=data, code=0)
|
||||
|
||||
|
||||
@router.get(
|
||||
"/{task_id}",
|
||||
response_model=ApiResponse[TaskDetail],
|
||||
dependencies=[Depends(require_bearer_token)],
|
||||
)
|
||||
async def get_order(
|
||||
task_id: str,
|
||||
container: GatewayContainer = Depends(get_container),
|
||||
) -> ApiResponse[TaskDetail]:
|
||||
"""任务详情 + 完整状态历史"""
|
||||
data = await container.task_queue.get_task_detail(task_id)
|
||||
return ApiResponse[TaskDetail](success=True, msg="success", data=data, code=0)
|
||||
|
||||
|
||||
@router.get(
|
||||
"",
|
||||
response_model=ApiResponse[TaskListData],
|
||||
dependencies=[Depends(require_bearer_token)],
|
||||
)
|
||||
async def list_orders(
|
||||
status: str | None = Query(default=None),
|
||||
site: str | None = Query(default=None),
|
||||
limit: int = Query(default=50, ge=1, le=500),
|
||||
offset: int = Query(default=0, ge=0),
|
||||
container: GatewayContainer = Depends(get_container),
|
||||
) -> ApiResponse[TaskListData]:
|
||||
"""任务列表(运维与上游对账用)"""
|
||||
data = await container.task_queue.list_tasks(
|
||||
status=status, site=site, limit=limit, offset=offset
|
||||
)
|
||||
return ApiResponse[TaskListData](success=True, msg="success", data=data, code=0)
|
||||
@@ -0,0 +1,29 @@
|
||||
"""网关容器:集中管理网关侧服务实例,用于依赖注入"""
|
||||
from __future__ import annotations
|
||||
|
||||
import asyncio
|
||||
import logging
|
||||
from dataclasses import dataclass
|
||||
|
||||
from app.gateway.db import GatewayDB
|
||||
from app.gateway.task_queue import TaskQueue
|
||||
from app.shared.config import Settings
|
||||
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
|
||||
@dataclass(slots=True)
|
||||
class GatewayContainer:
|
||||
"""网关容器
|
||||
|
||||
与抓取/交易容器同样的依赖注入风格,但持有的是任务队列与 SQLite 连接,
|
||||
生命周期由 main.lifespan 管理:start 时打开 DB,close 时关闭。
|
||||
|
||||
`sweep_task` 是常驻后台扫描,把过期的 leased/running 推到 stale。
|
||||
即使没有 lease 请求,过期的任务也会被及时发现(规格 §5 的关键约束)。
|
||||
"""
|
||||
|
||||
settings: Settings
|
||||
db: GatewayDB
|
||||
task_queue: TaskQueue
|
||||
sweep_task: asyncio.Task | None = None
|
||||
@@ -0,0 +1,356 @@
|
||||
"""网关 SQLite 访问层:三张表 + 原生 SQL,无业务逻辑
|
||||
|
||||
业务规则(状态机、并发度、幂等)放在 `task_queue.py`,本模块只负责持久化与查询,
|
||||
返回 dataclass 行。所有时间戳以 ISO8601 UTC 字符串存储(以 `Z` 结尾),便于
|
||||
跨进程对账;时间戳运算在 task_queue 层完成。
|
||||
|
||||
写操作由 task_queue 的 asyncio.Lock 串行化(详见该模块),本层不重复加锁,
|
||||
因此**调用方必须确保写操作在外层锁的保护下进行**。
|
||||
"""
|
||||
from __future__ import annotations
|
||||
|
||||
import json
|
||||
import logging
|
||||
from dataclasses import dataclass
|
||||
from pathlib import Path
|
||||
from typing import Any
|
||||
|
||||
import aiosqlite
|
||||
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
SCHEMA = """
|
||||
CREATE TABLE IF NOT EXISTS tasks (
|
||||
task_id TEXT PRIMARY KEY,
|
||||
site TEXT NOT NULL,
|
||||
intent_json TEXT NOT NULL,
|
||||
status TEXT NOT NULL,
|
||||
lease_owner TEXT,
|
||||
lease_expires_at TEXT,
|
||||
lease_count INTEGER NOT NULL DEFAULT 0,
|
||||
created_at TEXT NOT NULL,
|
||||
updated_at TEXT NOT NULL
|
||||
);
|
||||
CREATE INDEX IF NOT EXISTS idx_tasks_pending ON tasks(status, created_at);
|
||||
|
||||
CREATE TABLE IF NOT EXISTS task_reports (
|
||||
task_id TEXT NOT NULL,
|
||||
state TEXT NOT NULL,
|
||||
payable_yen INTEGER,
|
||||
pay_deadline TEXT,
|
||||
site_order_id TEXT,
|
||||
evidence_ref TEXT,
|
||||
detail TEXT,
|
||||
reported_at TEXT NOT NULL,
|
||||
PRIMARY KEY (task_id, state)
|
||||
);
|
||||
|
||||
CREATE TABLE IF NOT EXISTS workers (
|
||||
worker_id TEXT PRIMARY KEY,
|
||||
last_seen_at TEXT NOT NULL
|
||||
);
|
||||
"""
|
||||
|
||||
|
||||
@dataclass(slots=True)
|
||||
class TaskRow:
|
||||
"""tasks 表一行的强类型视图"""
|
||||
|
||||
task_id: str
|
||||
site: str
|
||||
intent_json: str
|
||||
status: str
|
||||
lease_owner: str | None
|
||||
lease_expires_at: str | None
|
||||
lease_count: int
|
||||
created_at: str
|
||||
updated_at: str
|
||||
|
||||
@property
|
||||
def intent(self) -> dict[str, Any]:
|
||||
return json.loads(self.intent_json)
|
||||
|
||||
|
||||
@dataclass(slots=True)
|
||||
class ReportRow:
|
||||
"""task_reports 表一行的强类型视图"""
|
||||
|
||||
task_id: str
|
||||
state: str
|
||||
payable_yen: int | None
|
||||
pay_deadline: str | None
|
||||
site_order_id: str | None
|
||||
evidence_ref: str | None
|
||||
detail: str
|
||||
reported_at: str
|
||||
|
||||
|
||||
@dataclass(slots=True)
|
||||
class WorkerRow:
|
||||
"""workers 表一行的强类型视图"""
|
||||
|
||||
worker_id: str
|
||||
last_seen_at: str
|
||||
|
||||
|
||||
def _row_to_task(row: aiosqlite.Row) -> TaskRow:
|
||||
return TaskRow(
|
||||
task_id=row["task_id"],
|
||||
site=row["site"],
|
||||
intent_json=row["intent_json"],
|
||||
status=row["status"],
|
||||
lease_owner=row["lease_owner"],
|
||||
lease_expires_at=row["lease_expires_at"],
|
||||
lease_count=row["lease_count"],
|
||||
created_at=row["created_at"],
|
||||
updated_at=row["updated_at"],
|
||||
)
|
||||
|
||||
|
||||
def _row_to_report(row: aiosqlite.Row) -> ReportRow:
|
||||
return ReportRow(
|
||||
task_id=row["task_id"],
|
||||
state=row["state"],
|
||||
payable_yen=row["payable_yen"],
|
||||
pay_deadline=row["pay_deadline"],
|
||||
site_order_id=row["site_order_id"],
|
||||
evidence_ref=row["evidence_ref"],
|
||||
detail=row["detail"],
|
||||
reported_at=row["reported_at"],
|
||||
)
|
||||
|
||||
|
||||
def _row_to_worker(row: aiosqlite.Row) -> WorkerRow:
|
||||
return WorkerRow(worker_id=row["worker_id"], last_seen_at=row["last_seen_at"])
|
||||
|
||||
|
||||
class GatewayDB:
|
||||
"""网关 SQLite 访问对象
|
||||
|
||||
单连接 + aiosqlite 的内部线程。生命周期由 `GatewayContainer` 管理:
|
||||
`start()` 在 lifespan 启动时调用,`close()` 在关闭时调用。
|
||||
"""
|
||||
|
||||
def __init__(self, db_path: Path):
|
||||
self._db_path = db_path
|
||||
self._conn: aiosqlite.Connection | None = None
|
||||
|
||||
async def start(self) -> None:
|
||||
"""打开连接并初始化 schema(IF NOT EXISTS,可重复执行)"""
|
||||
self._db_path.parent.mkdir(parents=True, exist_ok=True)
|
||||
self._conn = await aiosqlite.connect(str(self._db_path))
|
||||
self._conn.row_factory = aiosqlite.Row
|
||||
await self._conn.executescript(SCHEMA)
|
||||
await self._conn.commit()
|
||||
logger.info("网关 DB 已就绪:%s", self._db_path)
|
||||
|
||||
async def close(self) -> None:
|
||||
if self._conn is not None:
|
||||
await self._conn.close()
|
||||
self._conn = None
|
||||
|
||||
@property
|
||||
def conn(self) -> aiosqlite.Connection:
|
||||
if self._conn is None:
|
||||
raise RuntimeError("GatewayDB 未启动:先调用 start()")
|
||||
return self._conn
|
||||
|
||||
# ---- tasks ----
|
||||
|
||||
async def get_task(self, task_id: str) -> TaskRow | None:
|
||||
async with self.conn.execute("SELECT * FROM tasks WHERE task_id = ?", (task_id,)) as cur:
|
||||
row = await cur.fetchone()
|
||||
return _row_to_task(row) if row else None
|
||||
|
||||
async def insert_task(self, row: TaskRow) -> bool:
|
||||
"""插入新任务。返回 True=新建,False=task_id 已存在(幂等命中)"""
|
||||
try:
|
||||
await self.conn.execute(
|
||||
"INSERT INTO tasks (task_id, site, intent_json, status, "
|
||||
"lease_owner, lease_expires_at, lease_count, created_at, updated_at) "
|
||||
"VALUES (?, ?, ?, ?, NULL, NULL, 0, ?, ?)",
|
||||
(
|
||||
row.task_id,
|
||||
row.site,
|
||||
row.intent_json,
|
||||
row.status,
|
||||
row.created_at,
|
||||
row.updated_at,
|
||||
),
|
||||
)
|
||||
await self.conn.commit()
|
||||
return True
|
||||
except aiosqlite.IntegrityError:
|
||||
return False
|
||||
|
||||
async def update_task(
|
||||
self,
|
||||
task_id: str,
|
||||
*,
|
||||
status: str | None = None,
|
||||
lease_owner: str | None = None,
|
||||
lease_expires_at: str | None = None,
|
||||
lease_count: int | None = None,
|
||||
updated_at: str | None = None,
|
||||
clear_lease: bool = False,
|
||||
) -> None:
|
||||
"""更新任务字段。clear_lease=True 时把 lease_owner/expires_at 置 NULL"""
|
||||
sets: list[str] = []
|
||||
params: list[Any] = []
|
||||
if status is not None:
|
||||
sets.append("status = ?")
|
||||
params.append(status)
|
||||
if lease_owner is not None:
|
||||
sets.append("lease_owner = ?")
|
||||
params.append(lease_owner)
|
||||
if lease_expires_at is not None:
|
||||
sets.append("lease_expires_at = ?")
|
||||
params.append(lease_expires_at)
|
||||
if lease_count is not None:
|
||||
sets.append("lease_count = ?")
|
||||
params.append(lease_count)
|
||||
if updated_at is not None:
|
||||
sets.append("updated_at = ?")
|
||||
params.append(updated_at)
|
||||
if clear_lease:
|
||||
sets.append("lease_owner = NULL")
|
||||
sets.append("lease_expires_at = NULL")
|
||||
if not sets:
|
||||
return
|
||||
params.append(task_id)
|
||||
await self.conn.execute(
|
||||
f"UPDATE tasks SET {', '.join(sets)} WHERE task_id = ?",
|
||||
params,
|
||||
)
|
||||
await self.conn.commit()
|
||||
|
||||
async def list_tasks(
|
||||
self,
|
||||
*,
|
||||
status: str | None = None,
|
||||
site: str | None = None,
|
||||
limit: int = 50,
|
||||
offset: int = 0,
|
||||
) -> tuple[list[TaskRow], int]:
|
||||
"""分页列出任务,按 created_at 升序"""
|
||||
where = []
|
||||
params: list[Any] = []
|
||||
if status:
|
||||
where.append("status = ?")
|
||||
params.append(status)
|
||||
if site:
|
||||
where.append("site = ?")
|
||||
params.append(site)
|
||||
clause = f"WHERE {' AND '.join(where)}" if where else ""
|
||||
|
||||
async with self.conn.execute(
|
||||
f"SELECT COUNT(*) FROM tasks {clause}",
|
||||
params,
|
||||
) as cur:
|
||||
total = (await cur.fetchone())[0]
|
||||
|
||||
sql = f"SELECT * FROM tasks {clause} ORDER BY created_at ASC LIMIT ? OFFSET ?"
|
||||
async with self.conn.execute(sql, [*params, limit, offset]) as cur:
|
||||
rows = await cur.fetchall()
|
||||
return [_row_to_task(r) for r in rows], total
|
||||
|
||||
async def list_tasks_in_statuses(self, statuses: tuple[str, ...]) -> list[TaskRow]:
|
||||
"""取处于给定状态集合的全部任务(不分页,用于健康检查的 active 列表)"""
|
||||
if not statuses:
|
||||
return []
|
||||
placeholders = ", ".join("?" for _ in statuses)
|
||||
sql = (
|
||||
f"SELECT * FROM tasks WHERE status IN ({placeholders}) "
|
||||
"ORDER BY created_at ASC"
|
||||
)
|
||||
async with self.conn.execute(sql, statuses) as cur:
|
||||
rows = await cur.fetchall()
|
||||
return [_row_to_task(r) for r in rows]
|
||||
|
||||
async def pick_queued_task(self, site: str | None) -> TaskRow | None:
|
||||
"""取最早的 queued 任务,可选按站点过滤"""
|
||||
if site:
|
||||
sql = "SELECT * FROM tasks WHERE status = ? AND site = ? ORDER BY created_at ASC LIMIT 1"
|
||||
params: tuple[Any, ...] = ("queued", site)
|
||||
else:
|
||||
sql = "SELECT * FROM tasks WHERE status = ? ORDER BY created_at ASC LIMIT 1"
|
||||
params = ("queued",)
|
||||
async with self.conn.execute(sql, params) as cur:
|
||||
row = await cur.fetchone()
|
||||
return _row_to_task(row) if row else None
|
||||
|
||||
# ---- task_reports ----
|
||||
|
||||
async def upsert_report(self, row: ReportRow) -> bool:
|
||||
"""写入一条 report。同一 (task_id, state) 已存在则覆盖(幂等)。
|
||||
|
||||
返回 True=本次新写入,False=覆盖了既有行。
|
||||
|
||||
先 SELECT 再 INSERT OR REPLACE,避免依赖 SQLite rowcount——后者在
|
||||
INSERT OR REPLACE 上即使命中既有行也报 1,无法区分。
|
||||
"""
|
||||
async with self.conn.execute(
|
||||
"SELECT 1 FROM task_reports WHERE task_id = ? AND state = ?",
|
||||
(row.task_id, row.state),
|
||||
) as cur:
|
||||
existed = await cur.fetchone() is not None
|
||||
|
||||
await self.conn.execute(
|
||||
"INSERT OR REPLACE INTO task_reports "
|
||||
"(task_id, state, payable_yen, pay_deadline, site_order_id, "
|
||||
"evidence_ref, detail, reported_at) VALUES (?, ?, ?, ?, ?, ?, ?, ?)",
|
||||
(
|
||||
row.task_id,
|
||||
row.state,
|
||||
row.payable_yen,
|
||||
row.pay_deadline,
|
||||
row.site_order_id,
|
||||
row.evidence_ref,
|
||||
row.detail,
|
||||
row.reported_at,
|
||||
),
|
||||
)
|
||||
await self.conn.commit()
|
||||
return not existed
|
||||
|
||||
async def get_reports(self, task_id: str) -> list[ReportRow]:
|
||||
async with self.conn.execute(
|
||||
"SELECT * FROM task_reports WHERE task_id = ? ORDER BY reported_at ASC",
|
||||
(task_id,),
|
||||
) as cur:
|
||||
rows = await cur.fetchall()
|
||||
return [_row_to_report(r) for r in rows]
|
||||
|
||||
async def get_latest_state(self, task_id: str) -> str | None:
|
||||
"""取该任务最近一次 report 的 state,无 report 返回 None"""
|
||||
async with self.conn.execute(
|
||||
"SELECT state FROM task_reports WHERE task_id = ? "
|
||||
"ORDER BY reported_at DESC LIMIT 1",
|
||||
(task_id,),
|
||||
) as cur:
|
||||
row = await cur.fetchone()
|
||||
return row["state"] if row else None
|
||||
|
||||
# ---- workers ----
|
||||
|
||||
async def upsert_worker(self, worker_id: str, last_seen_at: str) -> None:
|
||||
await self.conn.execute(
|
||||
"INSERT INTO workers (worker_id, last_seen_at) VALUES (?, ?) "
|
||||
"ON CONFLICT(worker_id) DO UPDATE SET last_seen_at = excluded.last_seen_at",
|
||||
(worker_id, last_seen_at),
|
||||
)
|
||||
await self.conn.commit()
|
||||
|
||||
async def list_workers(self) -> list[WorkerRow]:
|
||||
async with self.conn.execute("SELECT * FROM workers ORDER BY last_seen_at DESC") as cur:
|
||||
rows = await cur.fetchall()
|
||||
return [_row_to_worker(r) for r in rows]
|
||||
|
||||
# ---- 计数 ----
|
||||
|
||||
async def count_by_status(self, status: str) -> int:
|
||||
async with self.conn.execute(
|
||||
"SELECT COUNT(*) FROM tasks WHERE status = ?",
|
||||
(status,),
|
||||
) as cur:
|
||||
return (await cur.fetchone())[0]
|
||||
@@ -0,0 +1,124 @@
|
||||
"""下单任务网关入口: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.health import router as health_router
|
||||
from app.gateway.api.routes.orders import router as orders_router
|
||||
from app.gateway.container import GatewayContainer
|
||||
from app.gateway.db import GatewayDB
|
||||
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,
|
||||
)
|
||||
return GatewayContainer(settings=settings, db=db, task_queue=task_queue)
|
||||
|
||||
|
||||
async def _sweep_loop(container: GatewayContainer) -> None:
|
||||
"""常驻后台任务:每 60 秒扫一次过期租约,把 leased/running 推到 stale
|
||||
|
||||
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)
|
||||
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)
|
||||
|
||||
container.sweep_task = asyncio.create_task(
|
||||
_sweep_loop(container), name="gateway-sweep"
|
||||
)
|
||||
|
||||
try:
|
||||
yield
|
||||
finally:
|
||||
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)
|
||||
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,
|
||||
)
|
||||
@@ -0,0 +1,194 @@
|
||||
"""网关 API 数据模型:请求体与响应体
|
||||
|
||||
intent 字段刻意保留成 `dict[str, Any]`——网关不解释下单意图,结构由 trading 侧
|
||||
定义。网关只负责把它存下来、原样吐给 worker,避免业务规则悄悄渗进任务队列。
|
||||
"""
|
||||
from __future__ import annotations
|
||||
|
||||
from typing import Any
|
||||
|
||||
from pydantic import BaseModel, Field
|
||||
|
||||
from app.shared.task_state import OrderState, TaskStatus
|
||||
|
||||
|
||||
# ---- POST /api/orders ----
|
||||
|
||||
|
||||
class SubmitOrderRequest(BaseModel):
|
||||
"""上游提交下单意图
|
||||
|
||||
task_id 可选:上游自带的幂等键。不传则服务端生成。同一个 task_id 重复提交
|
||||
不新建任务,返回既有任务且 created=false。
|
||||
"""
|
||||
|
||||
task_id: str | None = None
|
||||
site: str
|
||||
intent: dict[str, Any]
|
||||
|
||||
|
||||
class SubmitOrderData(BaseModel):
|
||||
"""提交响应"""
|
||||
|
||||
task_id: str
|
||||
status: TaskStatus
|
||||
created: bool # True=本次新建,False=命中既有任务(幂等)
|
||||
|
||||
|
||||
# ---- GET /api/orders/lease ----
|
||||
|
||||
|
||||
class LeaseData(BaseModel):
|
||||
"""lease 响应
|
||||
|
||||
无任务可领时整个 data 为 null(HTTP 仍 200)。lease_count > 1 表示这是恢复
|
||||
领取(stale → reclaim),worker 必须先核对站点订单列表,见 runner 主循环。
|
||||
"""
|
||||
|
||||
task_id: str
|
||||
site: str
|
||||
intent: dict[str, Any]
|
||||
lease_expires_at: str # ISO8601 UTC
|
||||
lease_count: int
|
||||
known_state: OrderState | None # 之前上报过的最新订单状态;首次领取为 null
|
||||
|
||||
|
||||
# ---- POST /api/orders/{id}/renew ----
|
||||
|
||||
|
||||
class RenewRequest(BaseModel):
|
||||
"""续租请求"""
|
||||
|
||||
worker_id: str
|
||||
|
||||
|
||||
class RenewData(BaseModel):
|
||||
"""续租响应"""
|
||||
|
||||
task_id: str
|
||||
lease_expires_at: str
|
||||
lease_count: int
|
||||
|
||||
|
||||
# ---- POST /api/orders/{id}/report ----
|
||||
|
||||
|
||||
class ReportRequest(BaseModel):
|
||||
"""本地回报订单状态
|
||||
|
||||
同一 (task_id, state) 重复上报是幂等的——网络抖动导致 worker 重发时覆盖同一行,
|
||||
不产生第二条记录。terminal=true 时释放租约并把任务推到终态。
|
||||
"""
|
||||
|
||||
worker_id: str
|
||||
state: OrderState
|
||||
payable_yen: int | None = None
|
||||
pay_deadline: str | None = None # ISO8601
|
||||
site_order_id: str | None = None
|
||||
evidence_ref: str | None = None # 本地相对路径,不含扩展名
|
||||
detail: str = ""
|
||||
terminal: bool = False
|
||||
terminal_status: TaskStatus | None = None # terminal=true 时指定终态,缺省按 state 推断
|
||||
|
||||
|
||||
class ReportData(BaseModel):
|
||||
"""回报响应"""
|
||||
|
||||
task_id: str
|
||||
status: TaskStatus
|
||||
recorded: bool # False=同 (task_id, state) 已存在,本次为幂等覆盖;True=新写入
|
||||
|
||||
|
||||
# ---- POST /api/orders/{id}/reclaim ----
|
||||
|
||||
|
||||
class ReclaimRequest(BaseModel):
|
||||
"""把 stale 任务重新租给 worker"""
|
||||
|
||||
worker_id: str
|
||||
|
||||
|
||||
class ReclaimData(BaseModel):
|
||||
"""reclaim 响应,结构与 LeaseData 一致,但 lease_count 必然 > 1"""
|
||||
|
||||
task_id: str
|
||||
site: str
|
||||
intent: dict[str, Any]
|
||||
lease_expires_at: str
|
||||
lease_count: int
|
||||
known_state: OrderState | None
|
||||
|
||||
|
||||
# ---- GET /api/orders/{id} 与 GET /api/orders ----
|
||||
|
||||
|
||||
class ReportEntry(BaseModel):
|
||||
"""task_reports 单行"""
|
||||
|
||||
state: OrderState
|
||||
payable_yen: int | None = None
|
||||
pay_deadline: str | None = None
|
||||
site_order_id: str | None = None
|
||||
evidence_ref: str | None = None
|
||||
detail: str = ""
|
||||
reported_at: str
|
||||
|
||||
|
||||
class TaskDetail(BaseModel):
|
||||
"""单个任务详情"""
|
||||
|
||||
task_id: str
|
||||
site: str
|
||||
intent: dict[str, Any]
|
||||
status: TaskStatus
|
||||
lease_owner: str | None = None
|
||||
lease_expires_at: str | None = None
|
||||
lease_count: int = 0
|
||||
created_at: str
|
||||
updated_at: str
|
||||
latest_state: OrderState | None = None
|
||||
reports: list[ReportEntry] = Field(default_factory=list)
|
||||
|
||||
|
||||
class TaskListData(BaseModel):
|
||||
"""任务列表"""
|
||||
|
||||
items: list[TaskDetail]
|
||||
total: int
|
||||
limit: int
|
||||
offset: int
|
||||
|
||||
|
||||
# ---- GET /health ----
|
||||
|
||||
|
||||
class WorkerHealthEntry(BaseModel):
|
||||
"""单个 worker 的健康指标"""
|
||||
|
||||
worker_id: str
|
||||
last_seen_at: str
|
||||
last_seen_seconds: int
|
||||
|
||||
|
||||
class QueuedAlertEntry(BaseModel):
|
||||
"""长时间无人领的任务告警项"""
|
||||
|
||||
task_id: str
|
||||
site: str
|
||||
created_at: str
|
||||
age_seconds: int
|
||||
|
||||
|
||||
class GatewayHealthData(BaseModel):
|
||||
"""网关健康状态
|
||||
|
||||
只做规格 §4.7 的两条兜底告警:worker 失联、任务长时间无人领。付款期限监控
|
||||
不在网关——本地机 7×24 在线,那套逻辑放本地。
|
||||
"""
|
||||
|
||||
status: str # "ok" 或 "degraded"
|
||||
queued_count: int
|
||||
active_tasks: list[TaskDetail] = Field(default_factory=list)
|
||||
workers: list[WorkerHealthEntry] = Field(default_factory=list)
|
||||
offline_workers: list[WorkerHealthEntry] = Field(default_factory=list)
|
||||
stale_queued_tasks: list[QueuedAlertEntry] = Field(default_factory=list)
|
||||
@@ -0,0 +1,505 @@
|
||||
"""任务队列业务逻辑:状态机、租约、并发度、幂等
|
||||
|
||||
所有写操作通过 `self._lock` 串行化。该锁同时保护「检查 + 写入」的原子序列
|
||||
(如 lease 时「检查是否已有 active 任务 → 取 queued → 改为 leased」),避免两个
|
||||
lease 请求并发取到同一个任务。
|
||||
|
||||
长轮询通过 `self._cond`(基于同一把锁的 Condition)唤醒:submit / report(terminal)
|
||||
/ reclaim 失败等任何可能让「队列前进」的事件都 notify_all。lease 在无可领任务时
|
||||
wait_for,超时返回 None。
|
||||
|
||||
绝不自动重投(docs/order-gateway.md §5):sweep 把过期任务置 stale 后**不**回到
|
||||
queued,普通 lease 取不到它;恢复只能走 reclaim。
|
||||
"""
|
||||
from __future__ import annotations
|
||||
|
||||
import asyncio
|
||||
import logging
|
||||
import secrets
|
||||
from datetime import datetime, timedelta, timezone
|
||||
|
||||
from app.shared.task_state import (
|
||||
ACTIVE_STATUSES,
|
||||
LEASABLE_STATUSES,
|
||||
RECLAIMABLE_STATUSES,
|
||||
TERMINAL_STATUSES,
|
||||
TaskStatus,
|
||||
)
|
||||
from app.gateway.db import GatewayDB, ReportRow, TaskRow
|
||||
from app.gateway.models import (
|
||||
GatewayHealthData,
|
||||
LeaseData,
|
||||
QueuedAlertEntry,
|
||||
ReportEntry,
|
||||
RenewData,
|
||||
ReclaimData,
|
||||
ReportData,
|
||||
SubmitOrderData,
|
||||
TaskDetail,
|
||||
TaskListData,
|
||||
WorkerHealthEntry,
|
||||
)
|
||||
from app.shared.errors import (
|
||||
InvalidTaskStateError,
|
||||
LeaseInvalidError,
|
||||
TaskNotFoundError,
|
||||
)
|
||||
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
|
||||
def utcnow() -> datetime:
|
||||
return datetime.now(timezone.utc)
|
||||
|
||||
|
||||
def to_iso(dt: datetime) -> str:
|
||||
"""统一时间戳格式:ISO8601 UTC,秒精度,Z 结尾"""
|
||||
return dt.astimezone(timezone.utc).replace(microsecond=0).isoformat().replace("+00:00", "Z")
|
||||
|
||||
|
||||
def parse_iso(s: str) -> datetime:
|
||||
"""解析 to_iso 产出的字符串(也兼容带 +00:00 的变体)"""
|
||||
if s.endswith("Z"):
|
||||
s = s[:-1] + "+00:00"
|
||||
return datetime.fromisoformat(s)
|
||||
|
||||
|
||||
def _generate_task_id() -> str:
|
||||
"""服务端生成的 task_id:日期 + 8 字节随机 hex,避免与上游自带的冲突"""
|
||||
return f"po-{to_iso(utcnow())[:10].replace('-', '')}-{secrets.token_hex(4)}"
|
||||
|
||||
|
||||
def _task_to_detail(task: TaskRow, latest_state: str | None, reports: list[ReportEntry]) -> TaskDetail:
|
||||
return TaskDetail(
|
||||
task_id=task.task_id,
|
||||
site=task.site,
|
||||
intent=task.intent,
|
||||
status=TaskStatus(task.status),
|
||||
lease_owner=task.lease_owner,
|
||||
lease_expires_at=task.lease_expires_at,
|
||||
lease_count=task.lease_count,
|
||||
created_at=task.created_at,
|
||||
updated_at=task.updated_at,
|
||||
latest_state=latest_state,
|
||||
reports=reports,
|
||||
)
|
||||
|
||||
|
||||
class TaskQueue:
|
||||
"""任务队列业务逻辑
|
||||
|
||||
持有一个 GatewayDB 实例。所有写操作与「检查 + 写入」原子序列都在 `self._lock`
|
||||
保护下;长轮询等待通过 `self._cond` 唤醒。
|
||||
"""
|
||||
|
||||
def __init__(self, db: GatewayDB, *, lease_ttl_seconds: int, worker_offline_alert_seconds: int):
|
||||
self._db = db
|
||||
self._lease_ttl = lease_ttl_seconds
|
||||
self._worker_offline_alert_seconds = worker_offline_alert_seconds
|
||||
self._lock = asyncio.Lock()
|
||||
self._cond = asyncio.Condition(self._lock)
|
||||
|
||||
# ---- 提交 ----
|
||||
|
||||
async def submit(self, *, task_id: str | None, site: str, intent: dict) -> SubmitOrderData:
|
||||
"""上游提交下单意图。task_id 缺省时服务端生成;重复提交幂等"""
|
||||
import json
|
||||
|
||||
tid = task_id or _generate_task_id()
|
||||
now = to_iso(utcnow())
|
||||
row = TaskRow(
|
||||
task_id=tid,
|
||||
site=site,
|
||||
intent_json=json.dumps(intent, ensure_ascii=False),
|
||||
status=TaskStatus.QUEUED.value,
|
||||
lease_owner=None,
|
||||
lease_expires_at=None,
|
||||
lease_count=0,
|
||||
created_at=now,
|
||||
updated_at=now,
|
||||
)
|
||||
async with self._lock:
|
||||
created = await self._db.insert_task(row)
|
||||
if created:
|
||||
self._cond.notify_all()
|
||||
else:
|
||||
existing = await self._db.get_task(tid)
|
||||
# 既不可能创建失败又拿不到既有行:除非被并发删除,按罕见错误处理
|
||||
if existing is None:
|
||||
raise RuntimeError(f"任务 {tid} 既未新建也无法读取,状态异常")
|
||||
row = existing
|
||||
return SubmitOrderData(task_id=row.task_id, status=TaskStatus(row.status), created=created)
|
||||
|
||||
# ---- 领取 ----
|
||||
|
||||
async def lease(
|
||||
self,
|
||||
*,
|
||||
worker_id: str,
|
||||
wait: int,
|
||||
site: str | None,
|
||||
max_wait: int,
|
||||
) -> LeaseData | None:
|
||||
"""本地长轮询领取
|
||||
|
||||
- 全局并发度 1:已有 leased/running 任务时立即返回 None
|
||||
- 无可领任务时挂起最多 `min(wait, max_wait)` 秒,超时返回 None
|
||||
- 拿到任务时 queued → leased,写入 lease_owner / lease_expires_at / lease_count
|
||||
- 顺带刷新 worker 心跳(last_seen_at),这就是 lease 兼任心跳的设计
|
||||
"""
|
||||
effective_wait = max(0, min(wait, max_wait))
|
||||
deadline = utcnow() + timedelta(seconds=effective_wait)
|
||||
|
||||
async with self._lock:
|
||||
# 先把过期租约扫到 stale,避免「卡死的任务」堵住新任务
|
||||
await self._sweep_locked()
|
||||
|
||||
await self._db.upsert_worker(worker_id, to_iso(utcnow()))
|
||||
|
||||
while True:
|
||||
result = await self._try_pick_locked(worker_id, site)
|
||||
if result is not None:
|
||||
return result
|
||||
|
||||
remaining = (deadline - utcnow()).total_seconds()
|
||||
if remaining <= 0:
|
||||
return None
|
||||
try:
|
||||
await asyncio.wait_for(self._cond.wait(), timeout=remaining)
|
||||
except asyncio.TimeoutError:
|
||||
return None
|
||||
|
||||
async def _try_pick_locked(self, worker_id: str, site: str | None) -> LeaseData | None:
|
||||
"""锁内尝试领取一次。返回 None 表示此刻无可领任务"""
|
||||
# 全局并发度 1:已有 leased/running 立即返回空
|
||||
active = await self._db.list_tasks_in_statuses(tuple(s.value for s in ACTIVE_STATUSES))
|
||||
if active:
|
||||
return None
|
||||
|
||||
task = await self._db.pick_queued_task(site)
|
||||
if task is None:
|
||||
return None
|
||||
|
||||
return await self._lease_to_locked(task, worker_id, known_state=None)
|
||||
|
||||
async def _lease_to_locked(
|
||||
self,
|
||||
task: TaskRow,
|
||||
worker_id: str,
|
||||
*,
|
||||
known_state: str | None,
|
||||
) -> LeaseData:
|
||||
"""把一个 queued 或 stale 任务租给 worker。调用方必须持有 self._lock"""
|
||||
expires = to_iso(utcnow() + timedelta(seconds=self._lease_ttl))
|
||||
new_count = task.lease_count + 1
|
||||
await self._db.update_task(
|
||||
task.task_id,
|
||||
status=TaskStatus.LEASED.value,
|
||||
lease_owner=worker_id,
|
||||
lease_expires_at=expires,
|
||||
lease_count=new_count,
|
||||
updated_at=to_iso(utcnow()),
|
||||
)
|
||||
# 已知状态:取最近一次 report 的 state(首次领取为 None)
|
||||
if known_state is None:
|
||||
known_state = await self._db.get_latest_state(task.task_id)
|
||||
return LeaseData(
|
||||
task_id=task.task_id,
|
||||
site=task.site,
|
||||
intent=task.intent,
|
||||
lease_expires_at=expires,
|
||||
lease_count=new_count,
|
||||
known_state=known_state, # type: ignore[arg-type]
|
||||
)
|
||||
|
||||
# ---- 续租 ----
|
||||
|
||||
async def renew(self, task_id: str, worker_id: str) -> RenewData:
|
||||
async with self._lock:
|
||||
task = await self._db.get_task(task_id)
|
||||
if task is None:
|
||||
raise TaskNotFoundError(task_id)
|
||||
if task.lease_owner != worker_id:
|
||||
raise LeaseInvalidError(f"worker {worker_id} 不是任务 {task_id} 的租约持有者")
|
||||
if TaskStatus(task.status) in TERMINAL_STATUSES:
|
||||
raise LeaseInvalidError(f"任务 {task_id} 已终结,无法续租")
|
||||
if TaskStatus(task.status) not in ACTIVE_STATUSES:
|
||||
raise LeaseInvalidError(f"任务 {task_id} 状态 {task.status},无可续租约")
|
||||
|
||||
expires = to_iso(utcnow() + timedelta(seconds=self._lease_ttl))
|
||||
await self._db.update_task(
|
||||
task_id,
|
||||
lease_expires_at=expires,
|
||||
updated_at=to_iso(utcnow()),
|
||||
)
|
||||
return RenewData(task_id=task_id, lease_expires_at=expires, lease_count=task.lease_count)
|
||||
|
||||
# ---- 回报 ----
|
||||
|
||||
async def report(
|
||||
self,
|
||||
task_id: str,
|
||||
*,
|
||||
worker_id: str,
|
||||
state: str,
|
||||
payable_yen: int | None,
|
||||
pay_deadline: str | None,
|
||||
site_order_id: str | None,
|
||||
evidence_ref: str | None,
|
||||
detail: str,
|
||||
terminal: bool,
|
||||
terminal_status: str | None,
|
||||
) -> ReportData:
|
||||
async with self._lock:
|
||||
task = await self._db.get_task(task_id)
|
||||
if task is None:
|
||||
raise TaskNotFoundError(task_id)
|
||||
if task.lease_owner != worker_id:
|
||||
raise LeaseInvalidError(f"worker {worker_id} 不是任务 {task_id} 的租约持有者")
|
||||
status = TaskStatus(task.status)
|
||||
|
||||
if status in TERMINAL_STATUSES:
|
||||
# 任务已终结:仍然追加 report(规格 §6 末尾「任务已 terminal 的仍可上报,
|
||||
# gateway 追加到 task_reports」),但不允许改任务状态
|
||||
if terminal and terminal_status is not None and terminal_status != status.value:
|
||||
raise InvalidTaskStateError(
|
||||
f"任务 {task_id} 已是终态 {status.value},不能改为 {terminal_status}"
|
||||
)
|
||||
elif status == TaskStatus.LEASED:
|
||||
# 首次 report:leased → running
|
||||
await self._db.update_task(
|
||||
task_id, status=TaskStatus.RUNNING.value, updated_at=to_iso(utcnow())
|
||||
)
|
||||
|
||||
now = to_iso(utcnow())
|
||||
inserted = await self._db.upsert_report(
|
||||
ReportRow(
|
||||
task_id=task_id,
|
||||
state=state,
|
||||
payable_yen=payable_yen,
|
||||
pay_deadline=pay_deadline,
|
||||
site_order_id=site_order_id,
|
||||
evidence_ref=evidence_ref,
|
||||
detail=detail,
|
||||
reported_at=now,
|
||||
)
|
||||
)
|
||||
|
||||
final_status = status.value
|
||||
if terminal:
|
||||
final_status = self._resolve_terminal_status(
|
||||
current=status, explicit=terminal_status, state=state
|
||||
)
|
||||
await self._db.update_task(
|
||||
task_id,
|
||||
status=final_status,
|
||||
clear_lease=True,
|
||||
updated_at=now,
|
||||
)
|
||||
# 任务终结可能让并发度 1 的闸门放开,唤醒等着的 lease
|
||||
self._cond.notify_all()
|
||||
|
||||
return ReportData(task_id=task_id, status=TaskStatus(final_status), recorded=inserted)
|
||||
|
||||
@staticmethod
|
||||
def _resolve_terminal_status(
|
||||
*, current: TaskStatus, explicit: str | None, state: str
|
||||
) -> str:
|
||||
"""terminal=true 时决定任务的终态。
|
||||
|
||||
显式指定优先;否则按订单 state 推断:paid → succeeded,cancelled → failed,
|
||||
其余(如 awaiting_payment 永久搁置)→ needs_human。
|
||||
|
||||
注意 cancelled 既可能来自明确失败,也可能是站点侧自动取消(如付款期限过期),
|
||||
统一记 failed 让上游看到需要处理。
|
||||
"""
|
||||
if explicit is not None:
|
||||
if explicit not in {s.value for s in TERMINAL_STATUSES}:
|
||||
raise InvalidTaskStateError(f"terminal_status={explicit} 不是终态")
|
||||
return explicit
|
||||
if state == "paid":
|
||||
return TaskStatus.SUCCEEDED.value
|
||||
if state == "cancelled":
|
||||
return TaskStatus.FAILED.value
|
||||
# awaiting_payment 等无明确结局的状态被标 terminal 时,多半是 worker 检测到
|
||||
# 3DS 之类人工环节,按 needs_human 处理
|
||||
return TaskStatus.NEEDS_HUMAN.value
|
||||
|
||||
# ---- 恢复(绝不自动重投,见规格 §5)----
|
||||
|
||||
async def reclaim(self, task_id: str, worker_id: str) -> ReclaimData:
|
||||
"""把 stale 任务重新租给 worker。lease_count 必然 > 1"""
|
||||
async with self._lock:
|
||||
task = await self._db.get_task(task_id)
|
||||
if task is None:
|
||||
raise TaskNotFoundError(task_id)
|
||||
if TaskStatus(task.status) not in RECLAIMABLE_STATUSES:
|
||||
raise InvalidTaskStateError(
|
||||
f"任务 {task_id} 状态 {task.status},不能 reclaim(只能 reclaim stale 任务)"
|
||||
)
|
||||
|
||||
known_state = await self._db.get_latest_state(task_id)
|
||||
lease = await self._lease_to_locked(task, worker_id, known_state=known_state)
|
||||
# 复用 LeaseData 的字段构造 ReclaimData(结构一致)
|
||||
return ReclaimData(
|
||||
task_id=lease.task_id,
|
||||
site=lease.site,
|
||||
intent=lease.intent,
|
||||
lease_expires_at=lease.lease_expires_at,
|
||||
lease_count=lease.lease_count,
|
||||
known_state=lease.known_state,
|
||||
)
|
||||
|
||||
# ---- 过期清扫 ----
|
||||
|
||||
async def sweep(self) -> int:
|
||||
"""把过期的 leased/running 任务置 stale。返回清扫条数"""
|
||||
async with self._lock:
|
||||
return await self._sweep_locked()
|
||||
|
||||
async def _sweep_locked(self) -> int:
|
||||
"""锁内执行 sweep。调用方必须持有 self._lock"""
|
||||
active = await self._db.list_tasks_in_statuses(tuple(s.value for s in ACTIVE_STATUSES))
|
||||
now = utcnow()
|
||||
swept = 0
|
||||
for task in active:
|
||||
if task.lease_expires_at and parse_iso(task.lease_expires_at) < now:
|
||||
await self._db.update_task(
|
||||
task.task_id,
|
||||
status=TaskStatus.STALE.value,
|
||||
clear_lease=True,
|
||||
updated_at=to_iso(now),
|
||||
)
|
||||
swept += 1
|
||||
logger.warning(
|
||||
"租约过期,任务置 stale(不自动重投):task_id=%s site=%s",
|
||||
task.task_id,
|
||||
task.site,
|
||||
)
|
||||
if swept:
|
||||
# stale 后腾出了「全局并发度 1」的槽位,唤醒等待者
|
||||
# (虽然 stale 任务不会被普通 lease 取到,但队列里可能还有 queued)
|
||||
self._cond.notify_all()
|
||||
return swept
|
||||
|
||||
# ---- 查询 ----
|
||||
|
||||
async def get_task_detail(self, task_id: str) -> TaskDetail:
|
||||
task = await self._db.get_task(task_id)
|
||||
if task is None:
|
||||
raise TaskNotFoundError(task_id)
|
||||
reports = await self._db.get_reports(task_id)
|
||||
latest = reports[-1].state if reports else None
|
||||
return _task_to_detail(
|
||||
task,
|
||||
latest_state=latest,
|
||||
reports=[
|
||||
ReportEntry(
|
||||
state=r.state,
|
||||
payable_yen=r.payable_yen,
|
||||
pay_deadline=r.pay_deadline,
|
||||
site_order_id=r.site_order_id,
|
||||
evidence_ref=r.evidence_ref,
|
||||
detail=r.detail,
|
||||
reported_at=r.reported_at,
|
||||
)
|
||||
for r in reports
|
||||
],
|
||||
)
|
||||
|
||||
async def list_tasks(
|
||||
self,
|
||||
*,
|
||||
status: str | None,
|
||||
site: str | None,
|
||||
limit: int,
|
||||
offset: int,
|
||||
) -> TaskListData:
|
||||
rows, total = await self._db.list_tasks(
|
||||
status=status, site=site, limit=limit, offset=offset
|
||||
)
|
||||
items: list[TaskDetail] = []
|
||||
for row in rows:
|
||||
reports = await self._db.get_reports(row.task_id)
|
||||
latest = reports[-1].state if reports else None
|
||||
items.append(
|
||||
_task_to_detail(
|
||||
row,
|
||||
latest_state=latest,
|
||||
reports=[
|
||||
ReportEntry(
|
||||
state=r.state,
|
||||
payable_yen=r.payable_yen,
|
||||
pay_deadline=r.pay_deadline,
|
||||
site_order_id=r.site_order_id,
|
||||
evidence_ref=r.evidence_ref,
|
||||
detail=r.detail,
|
||||
reported_at=r.reported_at,
|
||||
)
|
||||
for r in reports
|
||||
],
|
||||
)
|
||||
)
|
||||
return TaskListData(items=items, total=total, limit=limit, offset=offset)
|
||||
|
||||
# ---- 健康检查 ----
|
||||
|
||||
async def health_snapshot(self) -> GatewayHealthData:
|
||||
"""规格 §4.7:只做 worker 失联 + 任务长时间无人领两条兜底告警"""
|
||||
async with self._lock:
|
||||
queued_count = await self._db.count_by_status(TaskStatus.QUEUED.value)
|
||||
active_rows = await self._db.list_tasks_in_statuses(
|
||||
tuple(s.value for s in ACTIVE_STATUSES)
|
||||
)
|
||||
queued_rows = await self._db.list_tasks_in_statuses(
|
||||
(TaskStatus.QUEUED.value,)
|
||||
)
|
||||
workers = await self._db.list_workers()
|
||||
|
||||
now = utcnow()
|
||||
offline_workers: list[WorkerHealthEntry] = []
|
||||
worker_entries: list[WorkerHealthEntry] = []
|
||||
for w in workers:
|
||||
last_seen = parse_iso(w.last_seen_at)
|
||||
age = int((now - last_seen).total_seconds())
|
||||
entry = WorkerHealthEntry(
|
||||
worker_id=w.worker_id,
|
||||
last_seen_at=w.last_seen_at,
|
||||
last_seen_seconds=age,
|
||||
)
|
||||
worker_entries.append(entry)
|
||||
if age > self._worker_offline_alert_seconds:
|
||||
offline_workers.append(entry)
|
||||
|
||||
stale_queued_threshold = utcnow() - timedelta(seconds=self._worker_offline_alert_seconds)
|
||||
stale_queued: list[QueuedAlertEntry] = []
|
||||
for t in queued_rows:
|
||||
created = parse_iso(t.created_at)
|
||||
if created < stale_queued_threshold:
|
||||
stale_queued.append(
|
||||
QueuedAlertEntry(
|
||||
task_id=t.task_id,
|
||||
site=t.site,
|
||||
created_at=t.created_at,
|
||||
age_seconds=int((now - created).total_seconds()),
|
||||
)
|
||||
)
|
||||
|
||||
active_details = [
|
||||
_task_to_detail(
|
||||
t,
|
||||
latest_state=await self._db.get_latest_state(t.task_id),
|
||||
reports=[],
|
||||
)
|
||||
for t in active_rows
|
||||
]
|
||||
|
||||
degraded = bool(offline_workers or stale_queued)
|
||||
return GatewayHealthData(
|
||||
status="degraded" if degraded else "ok",
|
||||
queued_count=queued_count,
|
||||
active_tasks=active_details,
|
||||
workers=worker_entries,
|
||||
offline_workers=offline_workers,
|
||||
stale_queued_tasks=stale_queued,
|
||||
)
|
||||
Reference in New Issue
Block a user