From 2004e8c9107917cb2d297490e82a09cb75eae2b5 Mon Sep 17 00:00:00 2001 From: unknown <798648091@qq.com> Date: Thu, 23 Jul 2026 11:35:55 +0800 Subject: [PATCH] feat(cps): add reconciliation audit logs --- app/admin/repositories/cps.py | 87 +++++++-- app/core/cps_reconcile_worker.py | 299 ++++++++++++++++++++++++++--- tests/test_cps_reconcile_worker.py | 117 ++++++++++- 3 files changed, 462 insertions(+), 41 deletions(-) diff --git a/app/admin/repositories/cps.py b/app/admin/repositories/cps.py index 419c62f..7fd177b 100644 --- a/app/admin/repositories/cps.py +++ b/app/admin/repositories/cps.py @@ -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( diff --git a/app/core/cps_reconcile_worker.py b/app/core/cps_reconcile_worker.py index 9c0e6d2..d7c9fe6 100644 --- a/app/core/cps_reconcile_worker.py +++ b/app/core/cps_reconcile_worker.py @@ -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") diff --git a/tests/test_cps_reconcile_worker.py b/tests/test_cps_reconcile_worker.py index 953533a..6ff0e54 100644 --- a/tests/test_cps_reconcile_worker.py +++ b/tests/test_cps_reconcile_worker.py @@ -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