Compare commits

..

3 Commits

Author SHA1 Message Date
unknown 2004e8c910 feat(cps): add reconciliation audit logs 2026-07-23 11:35:55 +08:00
linkeyu b7cfcf7495 功能:美团和京东 CPS 每日自动对账 (#162)
## 需求

- 保留现有后台手动对账逻辑不变
- 每天北京时间 05:00 自动刷新 CPS 对账
- 当前仅处理美团和京东
- 每次按更新时间回拉近 3 天,覆盖延迟更新和订单状态变化

## 实现

- 新增进程内 CPS 自动对账 worker,并接入应用生命周期
- 美团使用更新时间查询类型 2,京东使用更新时间查询类型 3
- 复用现有仓储层对账及 order_id 幂等更新逻辑
- 美团和京东独立会话、独立异常处理,单个平台失败不阻塞另一平台
- 增加单实例锁、开关、执行小时、回拉天数和轮询间隔配置
- 服务在 05:00 后重启时会补跑当天任务

## 验证

- `pytest tests/test_cps_reconcile_worker.py tests/test_cps_admin.py tests/test_admin_read.py -q`:27 passed
- 相关文件 Ruff 检查通过
- `git diff --check` 通过

---------

Co-authored-by: unknown <798648091@qq.com>
Reviewed-on: #162
Co-authored-by: linkeyu <linkeyu@wonderable.ai>
Co-committed-by: linkeyu <linkeyu@wonderable.ai>
2026-07-23 10:51:20 +08:00
linkeyu 77f772f47c 性能:比价和领券分位数改为 PostgreSQL 聚合 (#161)
## 背景
- 比价记录页和领券记录页前端已经只拉当前页,并直接展示后端 summary。
- 原后端仍会将筛选区间内的全部耗时值取回 Python 计算分位数,生产数据量增大后会放大数据库读取和应用内存开销。

## 修改内容
- 比价记录:成功耗时 P5/P50/P95/P99、平均耗时以及中途退出耗时 P5/P50/P95 改用 PostgreSQL percentile_cont/AVG 聚合。
- 领券记录:发起数、完成数、平均耗时和完成耗时 P5/P50/P95/P99 合并为一条 PostgreSQL 聚合查询。
- 日期、用户、环境、状态、店铺和商品等筛选条件继续与列表共用,统计口径不变。
- SQLite 不支持 percentile_cont,仅在本地和测试环境回退读取耗时单列;不加载完整业务记录。
- API 字段与前端展示保持不变,无需前端改动。

## 验证
- 比价/领券及关联广告收益、点位、按券统计测试:26 passed。
- 本次涉及文件 ruff 检查通过。
- PostgreSQL SQL 编译测试确认使用 ordered-set percentile_cont 聚合。
- 全量测试:497 passed;8 个现有失败集中在邀请奖励、提现档位和代理转发等无关模块。

---------

Co-authored-by: unknown <798648091@qq.com>
Reviewed-on: #161
Co-authored-by: linkeyu <linkeyu@wonderable.ai>
Co-committed-by: linkeyu <linkeyu@wonderable.ai>
2026-07-23 10:35:41 +08:00
7 changed files with 685 additions and 86 deletions
+79 -18
View File
@@ -1,7 +1,7 @@
"""admin「领券数据」看板聚合:发起/完成数、耗时均值与分位、按天/小时趋势、逐条明细。
数据源 coupon_session(一次领券一行,客户端 /api/v1/coupon/session 两段上报)。量级不大,全量拉
区间数据后 Python 聚合(分位 SQLite 无 percentile,统一 Python 算,PG 上也一致)
数据源 coupon_session(一次领券一行,客户端 /api/v1/coupon/session 两段上报)。生产 PostgreSQL
使用 percentile_cont 聚合耗时分位SQLite 本地/测试环境回退读取耗时单列计算
- 发起数 = 区间内全部 session(含 started/completed/failed/abandoned),= 流失统计的基数。
- 完成数 / 耗时均值 / 分位 = 仅 status==completed 子集(成功跑完才有可比的"领券耗时")。
- summary/daily/hourly/total 在全量上算,不受分页;items 为排序后当前页。
@@ -45,6 +45,75 @@ def _percentile(sorted_vals: list[int], q: float) -> int | None:
return round(sorted_vals[lo] * (1 - frac) + sorted_vals[hi] * frac)
def _round_duration_ms(value) -> int | None:
"""将数据库聚合结果按既有 Python round 口径转为整数毫秒。"""
if value is None:
return None
return int(round(value))
def _coupon_summary_aggregate_stmt(conditions: list):
"""PostgreSQL 汇总卡聚合语句;计数、均值与四个分位一次返回。"""
completed = CouponSession.status == "completed"
completed_elapsed = completed & CouponSession.elapsed_ms.is_not(None)
return select(
func.count(CouponSession.id),
func.sum(case((completed, 1), else_=0)),
func.avg(CouponSession.elapsed_ms).filter(completed_elapsed),
*(
func.percentile_cont(q)
.within_group(CouponSession.elapsed_ms)
.filter(completed_elapsed)
for q in (0.05, 0.5, 0.95, 0.99)
),
).where(*conditions)
def _coupon_summary_aggregates(db: Session, conditions: list) -> dict:
"""汇总卡基础指标;生产 PG 全部在数据库内完成,SQLite 仅作测试回退。"""
if db.bind is not None and db.bind.dialect.name == "postgresql":
row = db.execute(_coupon_summary_aggregate_stmt(conditions)).one()
return {
"started_count": int(row[0] or 0),
"completed_count": int(row[1] or 0),
"avg_elapsed_ms": _round_duration_ms(row[2]),
"p5_ms": _round_duration_ms(row[3]),
"p50_ms": _round_duration_ms(row[4]),
"p95_ms": _round_duration_ms(row[5]),
"p99_ms": _round_duration_ms(row[6]),
}
counts = db.execute(
select(
func.count(CouponSession.id),
func.sum(case((CouponSession.status == "completed", 1), else_=0)),
).where(*conditions)
).one()
# SQLite 没有 percentile_cont;本地/测试只回退读取耗时单列,不加载完整记录。
completed_elapsed = list(
db.execute(
select(CouponSession.elapsed_ms)
.where(
*conditions,
CouponSession.status == "completed",
CouponSession.elapsed_ms.is_not(None),
)
.order_by(CouponSession.elapsed_ms)
).scalars()
)
return {
"started_count": int(counts[0] or 0),
"completed_count": int(counts[1] or 0),
"avg_elapsed_ms": _round_duration_ms(
sum(completed_elapsed) / len(completed_elapsed)
) if completed_elapsed else None,
"p5_ms": _percentile(completed_elapsed, 5),
"p50_ms": _percentile(completed_elapsed, 50),
"p95_ms": _percentile(completed_elapsed, 95),
"p99_ms": _percentile(completed_elapsed, 99),
}
def _avg(vals: list[int]) -> int | None:
return round(sum(vals) / len(vals)) if vals else None
@@ -215,30 +284,22 @@ def coupon_data_report(
if not user_ids:
return _empty_result()
stmt = select(CouponSession).where(
conditions = [
CouponSession.started_date >= d_from,
CouponSession.started_date <= d_to,
)
]
if app_env is not None:
stmt = stmt.where(CouponSession.app_env == app_env)
conditions.append(CouponSession.app_env == app_env)
if statuses:
stmt = stmt.where(CouponSession.status.in_(statuses))
conditions.append(CouponSession.status.in_(statuses))
if user_ids is not None:
stmt = stmt.where(CouponSession.user_id.in_(user_ids))
conditions.append(CouponSession.user_id.in_(user_ids))
stmt = select(CouponSession).where(*conditions)
rows = list(db.execute(stmt).scalars())
# ── 汇总卡 ──
completed_elapsed = sorted(
r.elapsed_ms for r in rows if r.status == "completed" and r.elapsed_ms is not None
)
summary = {
"started_count": len(rows),
"completed_count": sum(1 for r in rows if r.status == "completed"),
"avg_elapsed_ms": _avg(completed_elapsed),
"p5_ms": _percentile(completed_elapsed, 5),
"p50_ms": _percentile(completed_elapsed, 50),
"p95_ms": _percentile(completed_elapsed, 95),
"p99_ms": _percentile(completed_elapsed, 99),
**_coupon_summary_aggregates(db, conditions),
**_success_rates(rows),
}
@@ -323,7 +384,7 @@ def coupon_data_report(
"summary": summary,
"daily": daily,
"hourly": hourly,
"total": len(rows),
"total": summary["started_count"],
"items": items,
}
+71 -16
View File
@@ -5,6 +5,7 @@
"""
from __future__ import annotations
import logging
from datetime import datetime, timedelta, timezone
from decimal import Decimal, InvalidOperation
from typing import Any
@@ -15,12 +16,14 @@ from sqlalchemy.orm import Session
from app.admin.repositories.queries import _as_utc, offset_paginate
from app.integrations import jd_union, meituan
from app.repositories import cps_link as cps_link_repo
from app.models.cps_activity import CpsActivity
from app.models.cps_group import CpsGroup
from app.models.cps_link import CpsClick, CpsLink
from app.models.cps_order import CpsOrder
from app.models.cps_wx_user import CpsWxUser
from app.repositories import cps_link as cps_link_repo
logger = logging.getLogger("shagua.cps_reconcile")
# 美团订单状态:取消(4)/风控(5)不计佣金;结算(6)为佣金真正到账
_INVALID_STATUS = {"4", "5"}
@@ -378,18 +381,35 @@ def effective_commission_cents(order: CpsOrder) -> int:
def reconcile_orders(
db: Session, *, start_time: int, end_time: int,
query_time_type: int = 1, sid: str | None = None, max_pages: int = 200,
audit_context: dict[str, Any] | None = None,
) -> dict:
"""调美团 query_order 分页拉单 → 按 order_id upsert。返回 {fetched, inserted, updated, pages}。
订单状态会随时间变(付款→完成→结算/退款),重复拉同一单则更新。max_pages 防异常死循环。
"""
fetched = inserted = updated = pages = 0
fetched = inserted = updated = pages = api_requests = 0
page = 1
while page <= max_pages:
resp = meituan.query_order(
sid=sid, start_time=start_time, end_time=end_time,
query_time_type=query_time_type, page=page, limit=100,
)
api_requests += 1
try:
resp = meituan.query_order(
sid=sid, start_time=start_time, end_time=end_time,
query_time_type=query_time_type, page=page, limit=100,
)
except Exception: # noqa: BLE001 - 记录失败页后保持原异常类型继续抛出
if audit_context is not None:
logger.exception(
"CPS reconcile upstream request failed platform=meituan page=%s",
page,
extra={
**audit_context,
"event": "cps_reconcile.request_failed",
"platform": "meituan",
"failed_page": page,
"api_request_number": api_requests,
},
)
raise
rows = ((resp.get("data") or {}).get("dataList")) or []
if not rows:
break
@@ -419,30 +439,58 @@ def reconcile_orders(
break
page += 1
db.commit()
return {"fetched": fetched, "inserted": inserted, "updated": updated, "pages": pages}
return {
"fetched": fetched,
"inserted": inserted,
"updated": updated,
"pages": pages,
"api_requests": api_requests,
}
def reconcile_jd_orders(
db: Session, *, start_time: datetime, end_time: datetime,
query_time_type: int = 3, max_pages: int = 100,
audit_context: dict[str, Any] | None = None,
) -> dict:
"""调京东 order.row.query 拉单 → 按订单行 upsert。
京东单次查询窗口最多 1 小时,这里按北京自然时间切窗并逐页拉取。
"""
fetched = inserted = updated = pages = 0
fetched = inserted = updated = pages = api_requests = windows = 0
cur = start_time
while cur < end_time:
win_end = min(cur + timedelta(hours=1), end_time)
windows += 1
page = 1
while page <= max_pages:
resp = jd_union.query_order_rows(
start_time=cur,
end_time=win_end,
query_time_type=query_time_type,
page_index=page,
page_size=200,
)
api_requests += 1
try:
resp = jd_union.query_order_rows(
start_time=cur,
end_time=win_end,
query_time_type=query_time_type,
page_index=page,
page_size=200,
)
except Exception: # noqa: BLE001 - 记录失败窗口后保持原异常类型继续抛出
if audit_context is not None:
logger.exception(
"CPS reconcile upstream request failed platform=jd window=%s..%s page=%s",
cur.isoformat(),
win_end.isoformat(),
page,
extra={
**audit_context,
"event": "cps_reconcile.request_failed",
"platform": "jd",
"failed_window_start": cur.isoformat(),
"failed_window_end": win_end.isoformat(),
"failed_page": page,
"api_request_number": api_requests,
},
)
raise
rows = resp.get("rows") or []
has_more = bool(resp.get("has_more"))
if not rows:
@@ -469,7 +517,14 @@ def reconcile_jd_orders(
page += 1
cur = win_end
db.commit()
return {"fetched": fetched, "inserted": inserted, "updated": updated, "pages": pages}
return {
"fetched": fetched,
"inserted": inserted,
"updated": updated,
"pages": pages,
"api_requests": api_requests,
"windows": windows,
}
def list_orders(
+73 -25
View File
@@ -6,6 +6,7 @@
from __future__ import annotations
from datetime import date, datetime, time, timedelta, timezone
from decimal import ROUND_HALF_UP, Decimal
from zoneinfo import ZoneInfo
from sqlalchemy import Select, asc, case, desc, func, or_, select
@@ -297,6 +298,58 @@ def _comparison_percentile(sorted_values: list[int], q: float) -> int | None:
return int(value + 0.5)
def _round_duration_ms(value) -> int | None:
"""将数据库聚合结果按既有口径四舍五入为整数毫秒。"""
if value is None:
return None
return int(Decimal(str(value)).quantize(Decimal("1"), rounding=ROUND_HALF_UP))
def _comparison_duration_aggregate_stmt(conditions: list, status: str, quantiles: tuple[float, ...]):
"""PostgreSQL 耗时聚合语句;每种状态只返回一行。"""
return select(
func.avg(ComparisonRecord.total_ms),
*(
func.percentile_cont(q).within_group(ComparisonRecord.total_ms)
for q in quantiles
),
).where(
*conditions,
ComparisonRecord.status == status,
ComparisonRecord.total_ms.is_not(None),
)
def _comparison_duration_aggregates(
db: Session,
*,
conditions: list,
status: str,
quantiles: tuple[float, ...],
) -> list[int | None]:
"""返回平均值和各分位数;生产 PG 在数据库内聚合,SQLite 仅作测试回退。"""
if db.bind is not None and db.bind.dialect.name == "postgresql":
row = db.execute(
_comparison_duration_aggregate_stmt(conditions, status, quantiles)
).one()
return [_round_duration_ms(value) for value in row]
# SQLite 没有 percentile_cont;本地/测试只回退读取耗时单列,不加载完整记录。
values = list(
db.execute(
select(ComparisonRecord.total_ms)
.where(
*conditions,
ComparisonRecord.status == status,
ComparisonRecord.total_ms.is_not(None),
)
.order_by(ComparisonRecord.total_ms)
).scalars()
)
average = _round_duration_ms(sum(values) / len(values)) if values else None
return [average, *(_comparison_percentile(values, q) for q in quantiles)]
def comparison_records_summary(
db: Session,
*,
@@ -332,20 +385,18 @@ def comparison_records_summary(
success = int(row[2] or 0)
lower_price = int(row[4] or 0)
cancelled = int(row[5] or 0)
success_durations = sorted(db.execute(
select(ComparisonRecord.total_ms).where(
*conditions,
ComparisonRecord.status == "success",
ComparisonRecord.total_ms.is_not(None),
)
).scalars().all())
cancelled_durations = sorted(db.execute(
select(ComparisonRecord.total_ms).where(
*conditions,
ComparisonRecord.status == "cancelled",
ComparisonRecord.total_ms.is_not(None),
)
).scalars().all())
success_duration_stats = _comparison_duration_aggregates(
db,
conditions=conditions,
status="success",
quantiles=(0.05, 0.5, 0.95, 0.99),
)
cancelled_duration_stats = _comparison_duration_aggregates(
db,
conditions=conditions,
status="cancelled",
quantiles=(0.05, 0.5, 0.95),
)
success_rate_denominator = started - cancelled
return {
"started": started,
@@ -354,19 +405,16 @@ def comparison_records_summary(
"success_rate": success / success_rate_denominator if success_rate_denominator else None,
"avg_token_cost": float(row[3]) if row[3] is not None else None,
"lower_price_rate": lower_price / success if success else None,
"avg_duration_ms": (
int(sum(success_durations) / len(success_durations) + 0.5)
if success_durations else None
),
"p5_duration_ms": _comparison_percentile(success_durations, 0.05),
"p50_duration_ms": _comparison_percentile(success_durations, 0.5),
"p95_duration_ms": _comparison_percentile(success_durations, 0.95),
"p99_duration_ms": _comparison_percentile(success_durations, 0.99),
"avg_duration_ms": success_duration_stats[0],
"p5_duration_ms": success_duration_stats[1],
"p50_duration_ms": success_duration_stats[2],
"p95_duration_ms": success_duration_stats[3],
"p99_duration_ms": success_duration_stats[4],
"cancelled": cancelled,
"cancelled_rate": cancelled / started if started else None,
"cancelled_p5_ms": _comparison_percentile(cancelled_durations, 0.05),
"cancelled_p50_ms": _comparison_percentile(cancelled_durations, 0.5),
"cancelled_p95_ms": _comparison_percentile(cancelled_durations, 0.95),
"cancelled_p5_ms": cancelled_duration_stats[1],
"cancelled_p50_ms": cancelled_duration_stats[2],
"cancelled_p95_ms": cancelled_duration_stats[3],
}
+276 -23
View File
@@ -9,15 +9,17 @@ import asyncio
import contextlib
import logging
import os
import socket
import time
from collections.abc import Iterator
from collections.abc import Callable, Iterator
from datetime import date, datetime, timedelta
from pathlib import Path
from uuid import uuid4
from app.admin.repositories import cps as cps_repo
from app.core.config import settings
from app.core.rewards import CN_TZ
from app.db.session import SessionLocal
from app.db.session import SessionLocal, engine
logger = logging.getLogger("shagua.cps_reconcile")
_LOCK_PATH = Path(__file__).resolve().parents[2] / "data" / "cps_reconcile.lock"
@@ -27,6 +29,21 @@ def _cn_now() -> datetime:
return datetime.now(CN_TZ)
def _new_run_id(now: datetime) -> str:
return f"{now:%Y%m%d-%H%M%S}-{uuid4().hex[:8]}"
def _next_run_at(now: datetime, run_hour: int) -> datetime:
scheduled = now.replace(hour=run_hour, minute=0, second=0, microsecond=0)
return scheduled if now < scheduled else scheduled + timedelta(days=1)
def _trigger_for_run(worker_started_at: datetime, now: datetime, run_hour: int) -> str:
if worker_started_at.date() == now.date() and worker_started_at.hour >= run_hour:
return "startup_catchup"
return "scheduled"
def _touch_lock() -> None:
with contextlib.suppress(FileNotFoundError):
os.utime(_LOCK_PATH, None)
@@ -67,52 +84,240 @@ def _single_instance_lock(stale_after_sec: int) -> Iterator[bool]:
def _empty_result() -> dict:
return {"fetched": 0, "inserted": 0, "updated": 0, "pages": 0}
return {
"fetched": 0,
"inserted": 0,
"updated": 0,
"pages": 0,
"api_requests": 0,
}
def _reconcile_once(now: datetime | None = None) -> dict:
def _platform_log_fields(platform: str, result: dict) -> dict:
fields = {
"platform": platform,
"platform_status": result["status"],
"duration_ms": result["duration_ms"],
}
for key in ("fetched", "inserted", "updated", "pages", "api_requests", "windows"):
if key in result:
fields[key] = result[key]
return fields
def _run_platform(
*,
platform: str,
query_time_type: int,
common: dict,
reconcile: Callable[[], dict],
) -> tuple[dict, str | None]:
started = time.perf_counter()
logger.info(
"CPS auto reconcile platform started run_id=%s platform=%s",
common["run_id"],
platform,
extra={
**common,
"event": "cps_reconcile.platform_started",
"platform": platform,
"query_time_type": query_time_type,
},
)
try:
raw_result = reconcile()
except Exception as exc: # noqa: BLE001 - 单平台失败不能阻断另一平台
result = {
"status": "failed",
"duration_ms": round((time.perf_counter() - started) * 1000, 3),
"error_type": type(exc).__name__,
"error_summary": str(exc)[:500],
}
logger.exception(
"CPS auto reconcile platform failed run_id=%s platform=%s error_type=%s",
common["run_id"],
platform,
type(exc).__name__,
extra={
**common,
"event": "cps_reconcile.platform_failed",
**_platform_log_fields(platform, result),
"error_type": type(exc).__name__,
"error_summary": str(exc)[:500],
},
)
return result, str(exc)
result = {
**raw_result,
"status": "success",
"duration_ms": round((time.perf_counter() - started) * 1000, 3),
}
logger.info(
"CPS auto reconcile platform completed run_id=%s platform=%s fetched=%s inserted=%s updated=%s",
common["run_id"],
platform,
result.get("fetched", 0),
result.get("inserted", 0),
result.get("updated", 0),
extra={
**common,
"event": "cps_reconcile.platform_completed",
**_platform_log_fields(platform, result),
},
)
return result, None
def _skip_platform(platform: str, common: dict) -> dict:
result = {
**_empty_result(),
"status": "skipped",
"skipped": "not_configured",
"duration_ms": 0.0,
}
logger.warning(
"CPS auto reconcile platform skipped run_id=%s platform=%s reason=not_configured",
common["run_id"],
platform,
extra={
**common,
"event": "cps_reconcile.platform_skipped",
**_platform_log_fields(platform, result),
"skip_reason": "not_configured",
},
)
return result
def _reconcile_once(
now: datetime | None = None,
*,
run_id: str | None = None,
trigger: str = "scheduled",
) -> dict:
"""独立拉取美团和京东;单平台异常只记日志,不影响另一平台。"""
end = (now or _cn_now()).astimezone(CN_TZ)
run_id = run_id or _new_run_id(end)
lookback_days = max(1, int(settings.CPS_AUTO_RECONCILE_LOOKBACK_DAYS))
start = end - timedelta(days=lookback_days)
task_started = time.perf_counter()
common = {
"run_id": run_id,
"trigger": trigger,
"scheduled_date": end.date().isoformat(),
"app_env": settings.APP_ENV,
"db_dialect": engine.dialect.name,
"hostname": socket.gethostname(),
"pid": os.getpid(),
"window_start": start.isoformat(),
"window_end": end.isoformat(),
"lookback_days": lookback_days,
}
result = {
"run_id": run_id,
"trigger": trigger,
"window_start": start.isoformat(),
"window_end": end.isoformat(),
"meituan": None,
"jd": None,
"errors": {},
}
logger.info(
"CPS auto reconcile started run_id=%s trigger=%s window=%s..%s",
run_id,
trigger,
result["window_start"],
result["window_end"],
extra={**common, "event": "cps_reconcile.started"},
)
if settings.mt_cps_configured:
try:
def reconcile_meituan() -> dict:
with SessionLocal() as db:
result["meituan"] = cps_repo.reconcile_orders(
return cps_repo.reconcile_orders(
db,
start_time=int(start.timestamp()),
end_time=int(end.timestamp()),
query_time_type=2,
audit_context=common,
)
except Exception as exc: # noqa: BLE001 - 单平台失败不能阻断另一平台
result["errors"]["meituan"] = str(exc)
logger.exception("CPS auto reconcile failed platform=meituan")
result["meituan"], error = _run_platform(
platform="meituan",
query_time_type=2,
common=common,
reconcile=reconcile_meituan,
)
if error is not None:
result["errors"]["meituan"] = error
else:
result["meituan"] = {**_empty_result(), "skipped": "not_configured"}
result["meituan"] = _skip_platform("meituan", common)
_touch_lock()
if settings.jd_union_configured:
try:
def reconcile_jd() -> dict:
with SessionLocal() as db:
result["jd"] = cps_repo.reconcile_jd_orders(
return cps_repo.reconcile_jd_orders(
db,
start_time=start,
end_time=end,
query_time_type=3,
audit_context=common,
)
except Exception as exc: # noqa: BLE001 - 单平台失败不能阻断另一平台
result["errors"]["jd"] = str(exc)
logger.exception("CPS auto reconcile failed platform=jd")
result["jd"], error = _run_platform(
platform="jd",
query_time_type=3,
common=common,
reconcile=reconcile_jd,
)
if error is not None:
result["errors"]["jd"] = error
else:
result["jd"] = {**_empty_result(), "skipped": "not_configured"}
result["jd"] = _skip_platform("jd", common)
successes = sum(
platform_result.get("status") == "success"
for platform_result in (result["meituan"], result["jd"])
)
if result["errors"]:
task_status = "partial_success" if successes else "failed"
else:
task_status = "success"
completed_at = _cn_now()
duration_ms = round((time.perf_counter() - task_started) * 1000, 3)
next_run_at = _next_run_at(
completed_at,
min(max(int(settings.CPS_AUTO_RECONCILE_RUN_HOUR), 0), 23),
).isoformat()
result.update({
"status": task_status,
"duration_ms": duration_ms,
"manual_retry_required": bool(result["errors"]),
"next_run_at": next_run_at,
})
completion_fields = {
**common,
"event": "cps_reconcile.completed",
"task_status": task_status,
"duration_ms": duration_ms,
"manual_retry_required": result["manual_retry_required"],
"failed_platforms": sorted(result["errors"]),
"next_run_at": next_run_at,
"meituan_result": result["meituan"],
"jd_result": result["jd"],
}
log_method = logger.info if task_status == "success" else logger.warning
log_method(
"CPS auto reconcile completed run_id=%s status=%s duration_ms=%s failed_platforms=%s next_run_at=%s",
run_id,
task_status,
duration_ms,
",".join(sorted(result["errors"])) or "-",
next_run_at,
extra=completion_fields,
)
return result
@@ -127,42 +332,90 @@ async def _run_loop() -> None:
lock_stale_after = max(interval * 3, 1800)
with _single_instance_lock(lock_stale_after) as lock_acquired:
if not lock_acquired:
logger.warning("CPS auto reconcile skipped: another worker owns lock")
logger.warning(
"CPS auto reconcile worker skipped: another worker owns lock",
extra={
"event": "cps_reconcile.worker_skipped",
"skip_reason": "lock_not_acquired",
"pid": os.getpid(),
},
)
return
await _run_locked_loop(interval, run_hour)
async def _run_locked_loop(interval: int, run_hour: int) -> None:
worker_started_at = _cn_now()
logger.info(
"CPS auto reconcile worker started run_hour=%s interval=%ss lookback_days=%s",
run_hour,
interval,
settings.CPS_AUTO_RECONCILE_LOOKBACK_DAYS,
extra={
"event": "cps_reconcile.worker_started",
"pid": os.getpid(),
"hostname": socket.gethostname(),
"app_env": settings.APP_ENV,
"db_dialect": engine.dialect.name,
"run_hour": run_hour,
"interval_sec": interval,
"lookback_days": settings.CPS_AUTO_RECONCILE_LOOKBACK_DAYS,
"next_run_at": _next_run_at(worker_started_at, run_hour).isoformat(),
},
)
last_run: date | None = None
try:
while True:
run_id: str | None = None
try:
_touch_lock()
now = _cn_now()
if _should_run(last_run, now, run_hour):
result = await asyncio.to_thread(_reconcile_once, now)
run_id = _new_run_id(now)
trigger = _trigger_for_run(worker_started_at, now, run_hour)
await asyncio.to_thread(
_reconcile_once,
now,
run_id=run_id,
trigger=trigger,
)
last_run = now.date()
logger.info("CPS auto reconcile done date=%s result=%s", last_run, result)
except Exception: # noqa: BLE001 - 后台任务不能因单次异常退出
logger.exception("CPS auto reconcile unexpected error")
logger.exception(
"CPS auto reconcile unexpected worker error",
extra={
"event": "cps_reconcile.unexpected_failed",
"run_id": run_id or "unassigned",
"pid": os.getpid(),
},
)
await asyncio.sleep(interval)
except asyncio.CancelledError:
logger.info("CPS auto reconcile worker stopped")
logger.info(
"CPS auto reconcile worker stopped",
extra={"event": "cps_reconcile.worker_stopped", "pid": os.getpid()},
)
raise
def start_cps_reconcile_worker() -> asyncio.Task | None:
if not settings.CPS_AUTO_RECONCILE_ENABLED:
logger.info("CPS auto reconcile disabled")
logger.info(
"CPS auto reconcile disabled",
extra={
"event": "cps_reconcile.worker_skipped",
"skip_reason": "disabled",
},
)
return None
if not settings.mt_cps_configured and not settings.jd_union_configured:
logger.warning("CPS auto reconcile not started: Meituan and JD credentials are missing")
logger.warning(
"CPS auto reconcile not started: Meituan and JD credentials are missing",
extra={
"event": "cps_reconcile.worker_skipped",
"skip_reason": "all_platform_credentials_missing",
},
)
return None
return asyncio.create_task(_run_loop(), name="cps-auto-reconcile")
+16
View File
@@ -4,12 +4,28 @@ from __future__ import annotations
from datetime import UTC, date, datetime
import pytest
from sqlalchemy.dialects import postgresql
from app.admin.repositories import queries
from app.db.session import SessionLocal
from app.models.comparison import ComparisonRecord
def test_postgresql_duration_summary_uses_ordered_set_aggregates() -> None:
stmt = queries._comparison_duration_aggregate_stmt(
[], "success", (0.05, 0.5, 0.95, 0.99)
)
sql = str(
stmt.compile(
dialect=postgresql.dialect(),
compile_kwargs={"literal_binds": True},
)
)
assert sql.count("percentile_cont") == 4
assert "comparison_record.status = 'success'" in sql
def test_summary_uses_only_success_durations_and_filters_beijing_date() -> None:
db = SessionLocal()
try:
+55 -2
View File
@@ -7,8 +7,12 @@ from __future__ import annotations
from datetime import UTC, date, datetime
from sqlalchemy import delete, func, select
from sqlalchemy.dialects import postgresql
from app.admin.repositories.coupon_data import coupon_data_report
from app.admin.repositories.coupon_data import (
_coupon_summary_aggregate_stmt,
coupon_data_report,
)
from app.db.session import SessionLocal
from app.models.coupon_state import CouponSession
from app.repositories.coupon_state import (
@@ -19,7 +23,12 @@ from app.repositories.coupon_state import (
def _agg_session(
trace: str, platforms, platform_success, *, status: str = "completed"
trace: str,
platforms,
platform_success,
*,
status: str = "completed",
elapsed_ms: int | None = None,
) -> CouponSession:
"""构造一条聚合测试用 session(started_date 固定 2020-01-02、app_env=prod,不 commit)。"""
return CouponSession(
@@ -31,9 +40,53 @@ def _agg_session(
platform_success=platform_success,
started_at=datetime(2020, 1, 2, tzinfo=UTC),
started_date=date(2020, 1, 2),
elapsed_ms=elapsed_ms,
)
def test_postgresql_coupon_summary_uses_ordered_set_aggregates() -> None:
sql = str(
_coupon_summary_aggregate_stmt([]).compile(
dialect=postgresql.dialect(),
compile_kwargs={"literal_binds": True},
)
)
assert sql.count("percentile_cont") == 4
assert "FILTER (WHERE coupon_session.status = 'completed'" in sql
def test_coupon_duration_summary_uses_only_completed_rows() -> None:
db = SessionLocal()
try:
db.add_all([
_agg_session("duration-a", [], [], elapsed_ms=1000),
_agg_session("duration-b", [], [], elapsed_ms=3000),
_agg_session(
"duration-failed", [], [], status="failed", elapsed_ms=100_000
),
])
db.flush()
summary = coupon_data_report(
db,
date_from="2020-01-02",
date_to="2020-01-02",
app_env="prod",
)["summary"]
assert summary["started_count"] == 3
assert summary["completed_count"] == 2
assert summary["avg_elapsed_ms"] == 2000
assert summary["p5_ms"] == 1100
assert summary["p50_ms"] == 2000
assert summary["p95_ms"] == 2900
assert summary["p99_ms"] == 2980
finally:
db.rollback()
db.close()
def _make_session(db, trace_id: str, **kw) -> CouponSession:
row = CouponSession(
trace_id=trace_id,
+115 -2
View File
@@ -1,11 +1,15 @@
"""美团、京东 CPS 每日自动对账 worker。"""
from __future__ import annotations
import logging
from datetime import date, datetime, timedelta
import pytest
from app.core import cps_reconcile_worker as worker
from app.core.config import settings
from app.core.rewards import CN_TZ
from app.db.session import SessionLocal
from app.integrations.meituan import MeituanCpsError
@@ -23,11 +27,24 @@ def test_reconcile_once_pulls_meituan_and_jd_by_update_time(monkeypatch) -> None
def fake_meituan(db, **kwargs):
calls["meituan"] = kwargs
return {"fetched": 2, "inserted": 1, "updated": 1, "pages": 1}
return {
"fetched": 2,
"inserted": 1,
"updated": 1,
"pages": 1,
"api_requests": 1,
}
def fake_jd(db, **kwargs):
calls["jd"] = kwargs
return {"fetched": 3, "inserted": 2, "updated": 1, "pages": 2}
return {
"fetched": 3,
"inserted": 2,
"updated": 1,
"pages": 2,
"api_requests": 72,
"windows": 72,
}
monkeypatch.setattr(worker.cps_repo, "reconcile_orders", fake_meituan)
monkeypatch.setattr(worker.cps_repo, "reconcile_jd_orders", fake_jd)
@@ -45,6 +62,10 @@ def test_reconcile_once_pulls_meituan_and_jd_by_update_time(monkeypatch) -> None
assert calls["meituan"]["start_time"] == int((now - timedelta(days=3)).timestamp())
assert calls["jd"]["start_time"] == now - timedelta(days=3)
assert calls["jd"]["end_time"] == now
assert calls["meituan"]["audit_context"]["run_id"] == result["run_id"]
assert calls["jd"]["audit_context"]["run_id"] == result["run_id"]
assert result["status"] == "success"
assert result["manual_retry_required"] is False
def test_meituan_failure_does_not_block_jd(monkeypatch) -> None:
@@ -68,6 +89,9 @@ def test_meituan_failure_does_not_block_jd(monkeypatch) -> None:
assert result["errors"]["meituan"] == "temporary failure"
assert jd_called is True
assert result["jd"]["fetched"] == 1
assert result["meituan"]["status"] == "failed"
assert result["status"] == "partial_success"
assert result["manual_retry_required"] is True
def test_should_run_once_after_five_am() -> None:
@@ -79,6 +103,95 @@ def test_should_run_once_after_five_am() -> None:
assert worker._should_run(date(2026, 7, 22), at_five, 5) is False
def test_trigger_and_next_run_are_beijing_time() -> None:
before_five = datetime(2026, 7, 22, 4, 30, tzinfo=CN_TZ)
after_five = datetime(2026, 7, 22, 6, 0, tzinfo=CN_TZ)
assert worker._trigger_for_run(before_five, after_five, 5) == "scheduled"
assert worker._trigger_for_run(after_five, after_five, 5) == "startup_catchup"
assert worker._next_run_at(before_five, 5) == datetime(2026, 7, 22, 5, 0, tzinfo=CN_TZ)
assert worker._next_run_at(after_five, 5) == datetime(2026, 7, 23, 5, 0, tzinfo=CN_TZ)
def test_reconcile_emits_structured_audit_logs(monkeypatch, caplog) -> None:
_configure_platforms(monkeypatch)
monkeypatch.setattr(
worker.cps_repo,
"reconcile_orders",
lambda db, **kwargs: {
"fetched": 2,
"inserted": 1,
"updated": 1,
"pages": 1,
"api_requests": 1,
},
)
monkeypatch.setattr(
worker.cps_repo,
"reconcile_jd_orders",
lambda db, **kwargs: {
"fetched": 3,
"inserted": 2,
"updated": 1,
"pages": 1,
"api_requests": 72,
"windows": 72,
},
)
monkeypatch.setattr(worker, "_touch_lock", lambda: None)
monkeypatch.setattr(
worker,
"_cn_now",
lambda: datetime(2026, 7, 22, 5, 1, tzinfo=CN_TZ),
)
with caplog.at_level(logging.INFO, logger=worker.logger.name):
result = worker._reconcile_once(
datetime(2026, 7, 22, 5, 0, tzinfo=CN_TZ),
run_id="audit-run-1",
trigger="scheduled",
)
records = {record.event: record for record in caplog.records if hasattr(record, "event")}
assert records["cps_reconcile.started"].run_id == "audit-run-1"
assert records["cps_reconcile.started"].lookback_days == 3
assert records["cps_reconcile.started"].scheduled_date == "2026-07-22"
assert records["cps_reconcile.started"].hostname
assert records["cps_reconcile.completed"].task_status == "success"
assert records["cps_reconcile.completed"].manual_retry_required is False
assert records["cps_reconcile.completed"].meituan_result["fetched"] == 2
assert records["cps_reconcile.completed"].jd_result["api_requests"] == 72
assert result["next_run_at"] == "2026-07-23T05:00:00+08:00"
def test_jd_request_failure_logs_exact_window(monkeypatch, caplog) -> None:
start = datetime(2026, 7, 20, 5, 0, tzinfo=CN_TZ)
def fail_query(**kwargs):
raise RuntimeError("upstream timeout")
monkeypatch.setattr(worker.cps_repo.jd_union, "query_order_rows", fail_query)
with SessionLocal() as db, caplog.at_level(logging.ERROR, logger=worker.logger.name):
with pytest.raises(RuntimeError, match="upstream timeout"):
worker.cps_repo.reconcile_jd_orders(
db,
start_time=start,
end_time=start + timedelta(hours=2),
audit_context={"run_id": "audit-run-2", "trigger": "scheduled"},
)
record = next(
item
for item in caplog.records
if getattr(item, "event", "") == "cps_reconcile.request_failed"
)
assert record.run_id == "audit-run-2"
assert record.platform == "jd"
assert record.failed_window_start == "2026-07-20T05:00:00+08:00"
assert record.failed_window_end == "2026-07-20T06:00:00+08:00"
assert record.failed_page == 1
def test_start_worker_respects_auto_switch(monkeypatch) -> None:
monkeypatch.setattr(settings, "CPS_AUTO_RECONCILE_ENABLED", False)
assert worker.start_cps_reconcile_worker() is None