Compare commits
1 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| 7f28355378 |
@@ -113,13 +113,6 @@ JD_UNION_APP_SECRET=
|
||||
JD_UNION_SITE_ID=
|
||||
JD_UNION_AUTH_KEY=
|
||||
|
||||
# 美团 + 京东订单每天北京时间 05:00 自动对账;按更新时间回拉近 3 天,重叠防漏单并刷新状态。
|
||||
# 手动对账按钮不受该开关影响。通常保持开启;临时停自动任务时设为 false。
|
||||
CPS_AUTO_RECONCILE_ENABLED=true
|
||||
CPS_AUTO_RECONCILE_RUN_HOUR=5
|
||||
CPS_AUTO_RECONCILE_LOOKBACK_DAYS=3
|
||||
CPS_AUTO_RECONCILE_CHECK_INTERVAL_SEC=60
|
||||
|
||||
# ===== Pricebot 上游 (领券/比价业务透传目标) =====
|
||||
# 客户端调本服务的 /api/v1/coupon/step 等,我们透传到 pricebot-backend。
|
||||
# 本地开发用 localhost:8000。生产部署改成内网地址(如 http://pricebot.internal:8000)。
|
||||
|
||||
@@ -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,
|
||||
}
|
||||
|
||||
|
||||
@@ -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],
|
||||
}
|
||||
|
||||
|
||||
|
||||
@@ -188,13 +188,6 @@ class Settings(BaseSettings):
|
||||
"""京东联盟订单查询凭证齐全。"""
|
||||
return bool(self.JD_UNION_APP_KEY and self.JD_UNION_APP_SECRET)
|
||||
|
||||
# 美团 + 京东 CPS 订单自动对账:进程内 worker 每天北京时间 05:00 后跑一轮。
|
||||
# 按更新时间回拉近 N 天(重叠窗口防漏单并刷新状态),order_id 幂等更新;手动接口不受影响。
|
||||
CPS_AUTO_RECONCILE_ENABLED: bool = True
|
||||
CPS_AUTO_RECONCILE_RUN_HOUR: int = 5
|
||||
CPS_AUTO_RECONCILE_LOOKBACK_DAYS: int = 3
|
||||
CPS_AUTO_RECONCILE_CHECK_INTERVAL_SEC: int = 60
|
||||
|
||||
# ===== 微信服务号(网页授权) =====
|
||||
# CPS 落地页在微信内拿用户 openid(base 静默)/昵称头像(userinfo),做用户级群统计。
|
||||
# ⚠️ 区别于 WECHAT_APP_ID(那是 App 移动应用,用于微信支付);这是【已认证服务号】。
|
||||
|
||||
@@ -1,175 +0,0 @@
|
||||
"""美团、京东 CPS 订单每日自动对账任务。
|
||||
|
||||
每天北京时间 `CPS_AUTO_RECONCILE_RUN_HOUR`(默认 05:00)后执行一次,按更新时间
|
||||
回拉最近若干天订单并复用 admin CPS 仓储层的幂等 upsert。手动对账接口保持独立、不受影响。
|
||||
"""
|
||||
from __future__ import annotations
|
||||
|
||||
import asyncio
|
||||
import contextlib
|
||||
import logging
|
||||
import os
|
||||
import time
|
||||
from collections.abc import Iterator
|
||||
from datetime import date, datetime, timedelta
|
||||
from pathlib import Path
|
||||
|
||||
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
|
||||
|
||||
logger = logging.getLogger("shagua.cps_reconcile")
|
||||
_LOCK_PATH = Path(__file__).resolve().parents[2] / "data" / "cps_reconcile.lock"
|
||||
|
||||
|
||||
def _cn_now() -> datetime:
|
||||
return datetime.now(CN_TZ)
|
||||
|
||||
|
||||
def _touch_lock() -> None:
|
||||
with contextlib.suppress(FileNotFoundError):
|
||||
os.utime(_LOCK_PATH, None)
|
||||
|
||||
|
||||
@contextlib.contextmanager
|
||||
def _single_instance_lock(stale_after_sec: int) -> Iterator[bool]:
|
||||
"""同机多进程保护:同一时间只允许一个 CPS 自动对账 worker 运行。"""
|
||||
_LOCK_PATH.parent.mkdir(parents=True, exist_ok=True)
|
||||
fd: int | None = None
|
||||
try:
|
||||
try:
|
||||
fd = os.open(str(_LOCK_PATH), os.O_CREAT | os.O_EXCL | os.O_WRONLY)
|
||||
except FileExistsError:
|
||||
try:
|
||||
age = time.time() - _LOCK_PATH.stat().st_mtime
|
||||
except FileNotFoundError:
|
||||
age = stale_after_sec + 1
|
||||
if age > stale_after_sec:
|
||||
with contextlib.suppress(FileNotFoundError):
|
||||
_LOCK_PATH.unlink()
|
||||
try:
|
||||
fd = os.open(str(_LOCK_PATH), os.O_CREAT | os.O_EXCL | os.O_WRONLY)
|
||||
except FileExistsError:
|
||||
fd = None
|
||||
|
||||
if fd is None:
|
||||
yield False
|
||||
return
|
||||
|
||||
os.write(fd, f"pid={os.getpid()} started_at={int(time.time())}\n".encode("ascii"))
|
||||
yield True
|
||||
finally:
|
||||
if fd is not None:
|
||||
os.close(fd)
|
||||
with contextlib.suppress(FileNotFoundError):
|
||||
_LOCK_PATH.unlink()
|
||||
|
||||
|
||||
def _empty_result() -> dict:
|
||||
return {"fetched": 0, "inserted": 0, "updated": 0, "pages": 0}
|
||||
|
||||
|
||||
def _reconcile_once(now: datetime | None = None) -> dict:
|
||||
"""独立拉取美团和京东;单平台异常只记日志,不影响另一平台。"""
|
||||
end = (now or _cn_now()).astimezone(CN_TZ)
|
||||
lookback_days = max(1, int(settings.CPS_AUTO_RECONCILE_LOOKBACK_DAYS))
|
||||
start = end - timedelta(days=lookback_days)
|
||||
result = {
|
||||
"window_start": start.isoformat(),
|
||||
"window_end": end.isoformat(),
|
||||
"meituan": None,
|
||||
"jd": None,
|
||||
"errors": {},
|
||||
}
|
||||
|
||||
if settings.mt_cps_configured:
|
||||
try:
|
||||
with SessionLocal() as db:
|
||||
result["meituan"] = cps_repo.reconcile_orders(
|
||||
db,
|
||||
start_time=int(start.timestamp()),
|
||||
end_time=int(end.timestamp()),
|
||||
query_time_type=2,
|
||||
)
|
||||
except Exception as exc: # noqa: BLE001 - 单平台失败不能阻断另一平台
|
||||
result["errors"]["meituan"] = str(exc)
|
||||
logger.exception("CPS auto reconcile failed platform=meituan")
|
||||
else:
|
||||
result["meituan"] = {**_empty_result(), "skipped": "not_configured"}
|
||||
|
||||
_touch_lock()
|
||||
if settings.jd_union_configured:
|
||||
try:
|
||||
with SessionLocal() as db:
|
||||
result["jd"] = cps_repo.reconcile_jd_orders(
|
||||
db,
|
||||
start_time=start,
|
||||
end_time=end,
|
||||
query_time_type=3,
|
||||
)
|
||||
except Exception as exc: # noqa: BLE001 - 单平台失败不能阻断另一平台
|
||||
result["errors"]["jd"] = str(exc)
|
||||
logger.exception("CPS auto reconcile failed platform=jd")
|
||||
else:
|
||||
result["jd"] = {**_empty_result(), "skipped": "not_configured"}
|
||||
|
||||
return result
|
||||
|
||||
|
||||
def _should_run(last_run: date | None, now: datetime, run_hour: int) -> bool:
|
||||
return last_run != now.date() and now.hour >= run_hour
|
||||
|
||||
|
||||
async def _run_loop() -> None:
|
||||
interval = max(30, int(settings.CPS_AUTO_RECONCILE_CHECK_INTERVAL_SEC))
|
||||
run_hour = min(max(int(settings.CPS_AUTO_RECONCILE_RUN_HOUR), 0), 23)
|
||||
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")
|
||||
return
|
||||
await _run_locked_loop(interval, run_hour)
|
||||
|
||||
|
||||
async def _run_locked_loop(interval: int, run_hour: int) -> None:
|
||||
logger.info(
|
||||
"CPS auto reconcile worker started run_hour=%s interval=%ss lookback_days=%s",
|
||||
run_hour,
|
||||
interval,
|
||||
settings.CPS_AUTO_RECONCILE_LOOKBACK_DAYS,
|
||||
)
|
||||
last_run: date | None = None
|
||||
try:
|
||||
while True:
|
||||
try:
|
||||
_touch_lock()
|
||||
now = _cn_now()
|
||||
if _should_run(last_run, now, run_hour):
|
||||
result = await asyncio.to_thread(_reconcile_once, now)
|
||||
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")
|
||||
await asyncio.sleep(interval)
|
||||
except asyncio.CancelledError:
|
||||
logger.info("CPS auto reconcile worker stopped")
|
||||
raise
|
||||
|
||||
|
||||
def start_cps_reconcile_worker() -> asyncio.Task | None:
|
||||
if not settings.CPS_AUTO_RECONCILE_ENABLED:
|
||||
logger.info("CPS auto reconcile 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")
|
||||
return None
|
||||
return asyncio.create_task(_run_loop(), name="cps-auto-reconcile")
|
||||
|
||||
|
||||
async def stop_cps_reconcile_worker(task: asyncio.Task | None) -> None:
|
||||
if task is None:
|
||||
return
|
||||
task.cancel()
|
||||
with contextlib.suppress(asyncio.CancelledError):
|
||||
await task
|
||||
@@ -43,10 +43,6 @@ from app.api.v1.user import router as user_router
|
||||
from app.api.v1.wallet import router as wallet_router
|
||||
from app.api.v1.wxpay import router as wxpay_router
|
||||
from app.core.config import settings
|
||||
from app.core.cps_reconcile_worker import (
|
||||
start_cps_reconcile_worker,
|
||||
stop_cps_reconcile_worker,
|
||||
)
|
||||
from app.core.daily_exchange_worker import (
|
||||
start_daily_exchange_worker,
|
||||
stop_daily_exchange_worker,
|
||||
@@ -93,7 +89,6 @@ async def lifespan(_: FastAPI) -> AsyncIterator[None]:
|
||||
except Exception: # noqa: BLE001
|
||||
logger.exception("reverse_geocoder 预热失败(城市反查将在首个请求时懒加载)")
|
||||
reconcile_task = start_withdraw_reconcile_worker()
|
||||
cps_reconcile_task = start_cps_reconcile_worker()
|
||||
heartbeat_task = start_heartbeat_monitor()
|
||||
daily_exchange_task = start_daily_exchange_worker()
|
||||
observe_task = start_observe_worker()
|
||||
@@ -103,7 +98,6 @@ async def lifespan(_: FastAPI) -> AsyncIterator[None]:
|
||||
finally:
|
||||
await stop_heartbeat_monitor(heartbeat_task)
|
||||
await stop_withdraw_reconcile_worker(reconcile_task)
|
||||
await stop_cps_reconcile_worker(cps_reconcile_task)
|
||||
await stop_daily_exchange_worker(daily_exchange_task)
|
||||
await stop_observe_worker(observe_task)
|
||||
await stop_inactivity_reset_worker(inactivity_task)
|
||||
|
||||
@@ -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:
|
||||
|
||||
@@ -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,
|
||||
|
||||
@@ -1,84 +0,0 @@
|
||||
"""美团、京东 CPS 每日自动对账 worker。"""
|
||||
from __future__ import annotations
|
||||
|
||||
from datetime import date, datetime, timedelta
|
||||
|
||||
from app.core import cps_reconcile_worker as worker
|
||||
from app.core.config import settings
|
||||
from app.core.rewards import CN_TZ
|
||||
from app.integrations.meituan import MeituanCpsError
|
||||
|
||||
|
||||
def _configure_platforms(monkeypatch) -> None:
|
||||
monkeypatch.setattr(settings, "MT_CPS_APP_KEY", "mt-key")
|
||||
monkeypatch.setattr(settings, "MT_CPS_APP_SECRET", "mt-secret")
|
||||
monkeypatch.setattr(settings, "JD_UNION_APP_KEY", "jd-key")
|
||||
monkeypatch.setattr(settings, "JD_UNION_APP_SECRET", "jd-secret")
|
||||
monkeypatch.setattr(settings, "CPS_AUTO_RECONCILE_LOOKBACK_DAYS", 3)
|
||||
|
||||
|
||||
def test_reconcile_once_pulls_meituan_and_jd_by_update_time(monkeypatch) -> None:
|
||||
_configure_platforms(monkeypatch)
|
||||
calls: dict[str, dict] = {}
|
||||
|
||||
def fake_meituan(db, **kwargs):
|
||||
calls["meituan"] = kwargs
|
||||
return {"fetched": 2, "inserted": 1, "updated": 1, "pages": 1}
|
||||
|
||||
def fake_jd(db, **kwargs):
|
||||
calls["jd"] = kwargs
|
||||
return {"fetched": 3, "inserted": 2, "updated": 1, "pages": 2}
|
||||
|
||||
monkeypatch.setattr(worker.cps_repo, "reconcile_orders", fake_meituan)
|
||||
monkeypatch.setattr(worker.cps_repo, "reconcile_jd_orders", fake_jd)
|
||||
monkeypatch.setattr(worker, "_touch_lock", lambda: None)
|
||||
now = datetime(2026, 7, 22, 5, 0, tzinfo=CN_TZ)
|
||||
|
||||
result = worker._reconcile_once(now)
|
||||
|
||||
assert result["errors"] == {}
|
||||
assert result["meituan"]["fetched"] == 2
|
||||
assert result["jd"]["fetched"] == 3
|
||||
assert calls["meituan"]["query_time_type"] == 2
|
||||
assert calls["jd"]["query_time_type"] == 3
|
||||
assert calls["meituan"]["end_time"] == int(now.timestamp())
|
||||
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
|
||||
|
||||
|
||||
def test_meituan_failure_does_not_block_jd(monkeypatch) -> None:
|
||||
_configure_platforms(monkeypatch)
|
||||
jd_called = False
|
||||
|
||||
def fail_meituan(db, **kwargs):
|
||||
raise MeituanCpsError("temporary failure")
|
||||
|
||||
def fake_jd(db, **kwargs):
|
||||
nonlocal jd_called
|
||||
jd_called = True
|
||||
return {"fetched": 1, "inserted": 1, "updated": 0, "pages": 1}
|
||||
|
||||
monkeypatch.setattr(worker.cps_repo, "reconcile_orders", fail_meituan)
|
||||
monkeypatch.setattr(worker.cps_repo, "reconcile_jd_orders", fake_jd)
|
||||
monkeypatch.setattr(worker, "_touch_lock", lambda: None)
|
||||
|
||||
result = worker._reconcile_once(datetime(2026, 7, 22, 5, 0, tzinfo=CN_TZ))
|
||||
|
||||
assert result["errors"]["meituan"] == "temporary failure"
|
||||
assert jd_called is True
|
||||
assert result["jd"]["fetched"] == 1
|
||||
|
||||
|
||||
def test_should_run_once_after_five_am() -> None:
|
||||
before = datetime(2026, 7, 22, 4, 59, tzinfo=CN_TZ)
|
||||
at_five = datetime(2026, 7, 22, 5, 0, tzinfo=CN_TZ)
|
||||
|
||||
assert worker._should_run(None, before, 5) is False
|
||||
assert worker._should_run(None, at_five, 5) is True
|
||||
assert worker._should_run(date(2026, 7, 22), at_five, 5) is False
|
||||
|
||||
|
||||
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
|
||||
Reference in New Issue
Block a user