Compare commits

..

1 Commits

Author SHA1 Message Date
unknown 7f28355378 perf(admin): aggregate record duration percentiles in PostgreSQL 2026-07-22 19:37:30 +08:00
7 changed files with 16 additions and 742 deletions
-7
View File
@@ -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,26 +0,0 @@
"""merge notification/comparison_user_idx/monitoring_audit heads
Revision ID: 8e04cc13a211
Revises: comparison_user_created_idx, monitoring_audit_rbac, notification_table
Create Date: 2026-07-23 15:37:26.967540
"""
from typing import Sequence, Union
from alembic import op
import sqlalchemy as sa
# revision identifiers, used by Alembic.
revision: str = '8e04cc13a211'
down_revision: Union[str, Sequence[str], None] = ('comparison_user_created_idx', 'monitoring_audit_rbac', 'notification_table')
branch_labels: Union[str, Sequence[str], None] = None
depends_on: Union[str, Sequence[str], None] = None
def upgrade() -> None:
pass
def downgrade() -> None:
pass
+16 -71
View File
@@ -5,7 +5,6 @@
"""
from __future__ import annotations
import logging
from datetime import datetime, timedelta, timezone
from decimal import Decimal, InvalidOperation
from typing import Any
@@ -16,14 +15,12 @@ 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"}
@@ -381,35 +378,18 @@ 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 = api_requests = 0
fetched = inserted = updated = pages = 0
page = 1
while page <= max_pages:
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
resp = meituan.query_order(
sid=sid, start_time=start_time, end_time=end_time,
query_time_type=query_time_type, page=page, limit=100,
)
rows = ((resp.get("data") or {}).get("dataList")) or []
if not rows:
break
@@ -439,58 +419,30 @@ def reconcile_orders(
break
page += 1
db.commit()
return {
"fetched": fetched,
"inserted": inserted,
"updated": updated,
"pages": pages,
"api_requests": api_requests,
}
return {"fetched": fetched, "inserted": inserted, "updated": updated, "pages": pages}
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 = api_requests = windows = 0
fetched = inserted = updated = pages = 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:
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
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,
)
rows = resp.get("rows") or []
has_more = bool(resp.get("has_more"))
if not rows:
@@ -517,14 +469,7 @@ def reconcile_jd_orders(
page += 1
cur = win_end
db.commit()
return {
"fetched": fetched,
"inserted": inserted,
"updated": updated,
"pages": pages,
"api_requests": api_requests,
"windows": windows,
}
return {"fetched": fetched, "inserted": inserted, "updated": updated, "pages": pages}
def list_orders(
-7
View File
@@ -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 移动应用,用于微信支付);这是【已认证服务号】。
-428
View File
@@ -1,428 +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 socket
import time
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, engine
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 _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)
@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,
"api_requests": 0,
}
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:
def reconcile_meituan() -> dict:
with SessionLocal() as db:
return cps_repo.reconcile_orders(
db,
start_time=int(start.timestamp()),
end_time=int(end.timestamp()),
query_time_type=2,
audit_context=common,
)
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"] = _skip_platform("meituan", common)
_touch_lock()
if settings.jd_union_configured:
def reconcile_jd() -> dict:
with SessionLocal() as db:
return cps_repo.reconcile_jd_orders(
db,
start_time=start,
end_time=end,
query_time_type=3,
audit_context=common,
)
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"] = _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
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 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):
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()
except Exception: # noqa: BLE001 - 后台任务不能因单次异常退出
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",
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",
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",
extra={
"event": "cps_reconcile.worker_skipped",
"skip_reason": "all_platform_credentials_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
-6
View File
@@ -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)
-197
View File
@@ -1,197 +0,0 @@
"""美团、京东 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
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,
"api_requests": 1,
}
def fake_jd(db, **kwargs):
calls["jd"] = kwargs
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)
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
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:
_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
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:
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_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