Compare commits
1 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| 25c65280e5 |
@@ -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
|
||||
@@ -1,7 +1,7 @@
|
||||
"""admin「领券数据」看板聚合:发起/完成数、耗时均值与分位、按天/小时趋势、逐条明细。
|
||||
|
||||
数据源 coupon_session(一次领券一行,客户端 /api/v1/coupon/session 两段上报)。生产 PostgreSQL
|
||||
使用 percentile_cont 聚合耗时分位;SQLite 本地/测试环境回退读取耗时单列计算。
|
||||
数据源 coupon_session(一次领券一行,客户端 /api/v1/coupon/session 两段上报)。量级不大,全量拉
|
||||
区间数据后 Python 聚合(分位 SQLite 无 percentile,统一 Python 算,PG 上也一致)。
|
||||
- 发起数 = 区间内全部 session(含 started/completed/failed/abandoned),= 流失统计的基数。
|
||||
- 完成数 / 耗时均值 / 分位 = 仅 status==completed 子集(成功跑完才有可比的"领券耗时")。
|
||||
- summary/daily/hourly/total 在全量上算,不受分页;items 为排序后当前页。
|
||||
@@ -45,75 +45,6 @@ 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
|
||||
|
||||
@@ -284,22 +215,30 @@ def coupon_data_report(
|
||||
if not user_ids:
|
||||
return _empty_result()
|
||||
|
||||
conditions = [
|
||||
stmt = select(CouponSession).where(
|
||||
CouponSession.started_date >= d_from,
|
||||
CouponSession.started_date <= d_to,
|
||||
]
|
||||
)
|
||||
if app_env is not None:
|
||||
conditions.append(CouponSession.app_env == app_env)
|
||||
stmt = stmt.where(CouponSession.app_env == app_env)
|
||||
if statuses:
|
||||
conditions.append(CouponSession.status.in_(statuses))
|
||||
stmt = stmt.where(CouponSession.status.in_(statuses))
|
||||
if user_ids is not None:
|
||||
conditions.append(CouponSession.user_id.in_(user_ids))
|
||||
stmt = select(CouponSession).where(*conditions)
|
||||
stmt = stmt.where(CouponSession.user_id.in_(user_ids))
|
||||
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 = {
|
||||
**_coupon_summary_aggregates(db, conditions),
|
||||
"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),
|
||||
**_success_rates(rows),
|
||||
}
|
||||
|
||||
@@ -384,7 +323,7 @@ def coupon_data_report(
|
||||
"summary": summary,
|
||||
"daily": daily,
|
||||
"hourly": hourly,
|
||||
"total": summary["started_count"],
|
||||
"total": len(rows),
|
||||
"items": items,
|
||||
}
|
||||
|
||||
|
||||
@@ -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(
|
||||
|
||||
@@ -6,7 +6,6 @@
|
||||
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
|
||||
@@ -298,58 +297,6 @@ 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,
|
||||
*,
|
||||
@@ -385,18 +332,20 @@ def comparison_records_summary(
|
||||
success = int(row[2] or 0)
|
||||
lower_price = int(row[4] or 0)
|
||||
cancelled = int(row[5] or 0)
|
||||
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_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_rate_denominator = started - cancelled
|
||||
return {
|
||||
"started": started,
|
||||
@@ -405,16 +354,19 @@ 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": 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],
|
||||
"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),
|
||||
"cancelled": cancelled,
|
||||
"cancelled_rate": cancelled / started if started else None,
|
||||
"cancelled_p5_ms": cancelled_duration_stats[1],
|
||||
"cancelled_p50_ms": cancelled_duration_stats[2],
|
||||
"cancelled_p95_ms": cancelled_duration_stats[3],
|
||||
"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),
|
||||
}
|
||||
|
||||
|
||||
|
||||
+85
-167
@@ -12,10 +12,6 @@ from app.admin.repositories import mutations, queries
|
||||
from app.admin.schemas.common import CursorPage, OkResponse
|
||||
from app.admin.schemas.feedback import (
|
||||
FeedbackApproveRequest,
|
||||
FeedbackBulkApproveRequest,
|
||||
FeedbackBulkItemResult,
|
||||
FeedbackBulkRejectRequest,
|
||||
FeedbackBulkResult,
|
||||
FeedbackOut,
|
||||
FeedbackRejectRequest,
|
||||
FeedbackSummary,
|
||||
@@ -37,123 +33,6 @@ def _ensure_pending(fb: Feedback) -> None:
|
||||
raise HTTPException(status_code=400, detail="反馈已审核")
|
||||
|
||||
|
||||
def _approve_feedback(
|
||||
db: AdminDb,
|
||||
admin: AdminUser,
|
||||
feedback_id: int,
|
||||
payload: FeedbackApproveRequest | FeedbackBulkApproveRequest,
|
||||
ip: str,
|
||||
*,
|
||||
bulk: bool = False,
|
||||
) -> FeedbackOut:
|
||||
fb = db.get(Feedback, feedback_id, with_for_update=True)
|
||||
if fb is None:
|
||||
raise HTTPException(status_code=404, detail="反馈不存在")
|
||||
_ensure_pending(fb)
|
||||
|
||||
before = fb.status
|
||||
mutations.review_feedback(
|
||||
db,
|
||||
fb,
|
||||
status="adopted",
|
||||
reward_coins=payload.reward_coins,
|
||||
review_note=payload.note,
|
||||
admin_reply=payload.reply,
|
||||
reviewed_by_admin_id=admin.id,
|
||||
commit=False,
|
||||
)
|
||||
wallet_repo.grant_coins(
|
||||
db,
|
||||
fb.user_id,
|
||||
payload.reward_coins,
|
||||
biz_type="feedback_reward",
|
||||
ref_id=str(fb.id),
|
||||
remark="意见反馈被采纳",
|
||||
)
|
||||
detail = {
|
||||
"before": before,
|
||||
"after": "adopted",
|
||||
"reward_coins": payload.reward_coins,
|
||||
"note": payload.note,
|
||||
"reply": payload.reply,
|
||||
}
|
||||
if bulk:
|
||||
detail["bulk"] = True
|
||||
write_audit(
|
||||
db,
|
||||
admin,
|
||||
action="feedback.approve",
|
||||
target_type="feedback",
|
||||
target_id=feedback_id,
|
||||
detail=detail,
|
||||
ip=ip,
|
||||
commit=False,
|
||||
)
|
||||
db.commit()
|
||||
db.refresh(fb)
|
||||
out = FeedbackOut.model_validate(fb)
|
||||
notification_events.notify_feedback_reward(db, fb)
|
||||
return out
|
||||
|
||||
|
||||
def _reject_feedback(
|
||||
db: AdminDb,
|
||||
admin: AdminUser,
|
||||
feedback_id: int,
|
||||
payload: FeedbackRejectRequest | FeedbackBulkRejectRequest,
|
||||
ip: str,
|
||||
*,
|
||||
bulk: bool = False,
|
||||
) -> FeedbackOut:
|
||||
fb = db.get(Feedback, feedback_id, with_for_update=True)
|
||||
if fb is None:
|
||||
raise HTTPException(status_code=404, detail="反馈不存在")
|
||||
_ensure_pending(fb)
|
||||
|
||||
before = fb.status
|
||||
mutations.review_feedback(
|
||||
db,
|
||||
fb,
|
||||
status="rejected",
|
||||
reject_reason=payload.reason,
|
||||
review_note=payload.note,
|
||||
admin_reply=payload.reply,
|
||||
reviewed_by_admin_id=admin.id,
|
||||
commit=False,
|
||||
)
|
||||
detail = {
|
||||
"before": before,
|
||||
"after": "rejected",
|
||||
"reason": payload.reason,
|
||||
"note": payload.note,
|
||||
"reply": payload.reply,
|
||||
}
|
||||
if bulk:
|
||||
detail["bulk"] = True
|
||||
write_audit(
|
||||
db,
|
||||
admin,
|
||||
action="feedback.reject",
|
||||
target_type="feedback",
|
||||
target_id=feedback_id,
|
||||
detail=detail,
|
||||
ip=ip,
|
||||
commit=False,
|
||||
)
|
||||
db.commit()
|
||||
db.refresh(fb)
|
||||
out = FeedbackOut.model_validate(fb)
|
||||
notification_events.notify_feedback_reply(db, fb)
|
||||
return out
|
||||
|
||||
|
||||
def _bulk_result(items: list[FeedbackBulkItemResult]) -> FeedbackBulkResult:
|
||||
success = sum(1 for item in items if item.ok)
|
||||
return FeedbackBulkResult(
|
||||
total=len(items), success=success, failed=len(items) - success, items=items,
|
||||
)
|
||||
|
||||
|
||||
@router.get("", response_model=CursorPage[FeedbackOut], summary="反馈工单列表")
|
||||
def list_feedbacks(
|
||||
db: AdminDb,
|
||||
@@ -194,50 +73,6 @@ def feedback_summary(db: AdminDb) -> FeedbackSummary:
|
||||
return FeedbackSummary.model_validate(queries.feedback_summary(db))
|
||||
|
||||
|
||||
@router.post("/bulk/approve", response_model=FeedbackBulkResult, summary="批量采纳反馈并发金币")
|
||||
def bulk_approve_feedbacks(
|
||||
body: FeedbackBulkApproveRequest,
|
||||
request: Request,
|
||||
admin: Annotated[AdminUser, Depends(require_role("operator"))],
|
||||
db: AdminDb,
|
||||
) -> FeedbackBulkResult:
|
||||
results: list[FeedbackBulkItemResult] = []
|
||||
ip = get_client_ip(request)
|
||||
for feedback_id in body.ids:
|
||||
try:
|
||||
out = _approve_feedback(db, admin, feedback_id, body, ip, bulk=True)
|
||||
results.append(FeedbackBulkItemResult(id=feedback_id, ok=True, status=out.status))
|
||||
except HTTPException as exc:
|
||||
db.rollback()
|
||||
results.append(FeedbackBulkItemResult(id=feedback_id, ok=False, error=str(exc.detail)))
|
||||
except Exception: # noqa: BLE001 - 单笔失败不打断整批
|
||||
db.rollback()
|
||||
results.append(FeedbackBulkItemResult(id=feedback_id, ok=False, error="系统异常"))
|
||||
return _bulk_result(results)
|
||||
|
||||
|
||||
@router.post("/bulk/reject", response_model=FeedbackBulkResult, summary="批量拒绝采纳反馈")
|
||||
def bulk_reject_feedbacks(
|
||||
body: FeedbackBulkRejectRequest,
|
||||
request: Request,
|
||||
admin: Annotated[AdminUser, Depends(require_role("operator"))],
|
||||
db: AdminDb,
|
||||
) -> FeedbackBulkResult:
|
||||
results: list[FeedbackBulkItemResult] = []
|
||||
ip = get_client_ip(request)
|
||||
for feedback_id in body.ids:
|
||||
try:
|
||||
out = _reject_feedback(db, admin, feedback_id, body, ip, bulk=True)
|
||||
results.append(FeedbackBulkItemResult(id=feedback_id, ok=True, status=out.status))
|
||||
except HTTPException as exc:
|
||||
db.rollback()
|
||||
results.append(FeedbackBulkItemResult(id=feedback_id, ok=False, error=str(exc.detail)))
|
||||
except Exception: # noqa: BLE001 - 单笔失败不打断整批
|
||||
db.rollback()
|
||||
results.append(FeedbackBulkItemResult(id=feedback_id, ok=False, error="系统异常"))
|
||||
return _bulk_result(results)
|
||||
|
||||
|
||||
@router.post("/{feedback_id}/handle", response_model=OkResponse, summary="标记反馈已处理")
|
||||
def handle_feedback(
|
||||
feedback_id: int,
|
||||
@@ -258,7 +93,53 @@ def approve_feedback(
|
||||
admin: Annotated[AdminUser, Depends(require_role("operator"))],
|
||||
db: AdminDb,
|
||||
) -> FeedbackOut:
|
||||
return _approve_feedback(db, admin, feedback_id, payload, get_client_ip(request))
|
||||
fb = db.get(Feedback, feedback_id)
|
||||
if fb is None:
|
||||
raise HTTPException(status_code=404, detail="反馈不存在")
|
||||
_ensure_pending(fb)
|
||||
|
||||
before = fb.status
|
||||
mutations.review_feedback(
|
||||
db,
|
||||
fb,
|
||||
status="adopted",
|
||||
reward_coins=payload.reward_coins,
|
||||
review_note=payload.note,
|
||||
admin_reply=payload.reply,
|
||||
reviewed_by_admin_id=admin.id,
|
||||
commit=False,
|
||||
)
|
||||
wallet_repo.grant_coins(
|
||||
db,
|
||||
fb.user_id,
|
||||
payload.reward_coins,
|
||||
biz_type="feedback_reward",
|
||||
ref_id=str(fb.id),
|
||||
remark="意见反馈被采纳",
|
||||
)
|
||||
write_audit(
|
||||
db,
|
||||
admin,
|
||||
action="feedback.approve",
|
||||
target_type="feedback",
|
||||
target_id=feedback_id,
|
||||
detail={
|
||||
"before": before,
|
||||
"after": "adopted",
|
||||
"reward_coins": payload.reward_coins,
|
||||
"note": payload.note,
|
||||
"reply": payload.reply,
|
||||
},
|
||||
ip=get_client_ip(request),
|
||||
commit=False,
|
||||
)
|
||||
db.commit()
|
||||
db.refresh(fb)
|
||||
out = FeedbackOut.model_validate(fb)
|
||||
# PRD #10 反馈奖励:采纳发金币后通知用户(站内 + push,必带官方留言)。
|
||||
# 业务已 commit,通知失败只 log 不影响审核结果。
|
||||
notification_events.notify_feedback_reward(db, fb)
|
||||
return out
|
||||
|
||||
|
||||
@router.post("/{feedback_id}/reject", response_model=FeedbackOut, summary="拒绝采纳反馈")
|
||||
@@ -269,4 +150,41 @@ def reject_feedback(
|
||||
admin: Annotated[AdminUser, Depends(require_role("operator"))],
|
||||
db: AdminDb,
|
||||
) -> FeedbackOut:
|
||||
return _reject_feedback(db, admin, feedback_id, payload, get_client_ip(request))
|
||||
fb = db.get(Feedback, feedback_id)
|
||||
if fb is None:
|
||||
raise HTTPException(status_code=404, detail="反馈不存在")
|
||||
_ensure_pending(fb)
|
||||
|
||||
before = fb.status
|
||||
mutations.review_feedback(
|
||||
db,
|
||||
fb,
|
||||
status="rejected",
|
||||
reject_reason=payload.reason,
|
||||
review_note=payload.note,
|
||||
admin_reply=payload.reply,
|
||||
reviewed_by_admin_id=admin.id,
|
||||
commit=False,
|
||||
)
|
||||
write_audit(
|
||||
db,
|
||||
admin,
|
||||
action="feedback.reject",
|
||||
target_type="feedback",
|
||||
target_id=feedback_id,
|
||||
detail={
|
||||
"before": before,
|
||||
"after": "rejected",
|
||||
"reason": payload.reason,
|
||||
"note": payload.note,
|
||||
"reply": payload.reply,
|
||||
},
|
||||
ip=get_client_ip(request),
|
||||
commit=False,
|
||||
)
|
||||
db.commit()
|
||||
db.refresh(fb)
|
||||
out = FeedbackOut.model_validate(fb)
|
||||
# PRD #9 官方回复:未采纳也回复了用户(原因/留言用户端可见),通知去反馈历史页查看。
|
||||
notification_events.notify_feedback_reply(db, fb)
|
||||
return out
|
||||
|
||||
@@ -17,10 +17,6 @@ from app.admin.deps import AdminDb, get_client_ip, get_current_admin, require_ro
|
||||
from app.admin.repositories import mutations, queries
|
||||
from app.admin.schemas.common import CursorPage, OkResponse
|
||||
from app.admin.schemas.price_report import (
|
||||
PriceReportBulkItemResult,
|
||||
PriceReportBulkRejectRequest,
|
||||
PriceReportBulkRequest,
|
||||
PriceReportBulkResult,
|
||||
PriceReportOut,
|
||||
PriceReportRejectRequest,
|
||||
PriceReportSummary,
|
||||
@@ -38,59 +34,6 @@ router = APIRouter(
|
||||
)
|
||||
|
||||
|
||||
def _approve_price_report(
|
||||
db: AdminDb, admin: AdminUser, report_id: int, ip: str, *, bulk: bool = False
|
||||
) -> PriceReport:
|
||||
rep = db.get(PriceReport, report_id, with_for_update=True)
|
||||
if rep is None:
|
||||
raise HTTPException(status_code=404, detail="上报记录不存在")
|
||||
if rep.status != "pending":
|
||||
raise HTTPException(status_code=400, detail=f"该上报已审核过(当前 {rep.status}),不可重复操作")
|
||||
coins = PRICE_REPORT_REWARD_COINS
|
||||
mutations.review_price_report(db, rep, status="approved", reward_coins=coins, commit=False)
|
||||
wallet_repo.grant_coins(
|
||||
db, rep.user_id, coins,
|
||||
biz_type="price_report_reward", ref_id=str(rep.id), remark="上报更低价审核通过",
|
||||
)
|
||||
detail = {"reward_coins": coins, "user_id": rep.user_id}
|
||||
if bulk:
|
||||
detail["bulk"] = True
|
||||
write_audit(
|
||||
db, admin, action="price_report.approve", target_type="price_report", target_id=report_id,
|
||||
detail=detail, ip=ip, commit=False,
|
||||
)
|
||||
db.commit()
|
||||
notification_events.notify_report_approved(db, rep)
|
||||
return rep
|
||||
|
||||
|
||||
def _reject_price_report(
|
||||
db: AdminDb, admin: AdminUser, report_id: int, reason: str, ip: str, *, bulk: bool = False
|
||||
) -> PriceReport:
|
||||
rep = db.get(PriceReport, report_id, with_for_update=True)
|
||||
if rep is None:
|
||||
raise HTTPException(status_code=404, detail="上报记录不存在")
|
||||
if rep.status != "pending":
|
||||
raise HTTPException(status_code=400, detail=f"该上报已审核过(当前 {rep.status}),不可重复操作")
|
||||
mutations.review_price_report(db, rep, status="rejected", reject_reason=reason, commit=False)
|
||||
detail = {"reason": reason, "user_id": rep.user_id}
|
||||
if bulk:
|
||||
detail["bulk"] = True
|
||||
write_audit(
|
||||
db, admin, action="price_report.reject", target_type="price_report", target_id=report_id,
|
||||
detail=detail, ip=ip, commit=False,
|
||||
)
|
||||
db.commit()
|
||||
return rep
|
||||
|
||||
|
||||
def _bulk_result(items: list[PriceReportBulkItemResult]) -> PriceReportBulkResult:
|
||||
success = sum(1 for item in items if item.ok)
|
||||
return PriceReportBulkResult(
|
||||
total=len(items), success=success, failed=len(items) - success, items=items,
|
||||
)
|
||||
|
||||
|
||||
@router.get("", response_model=CursorPage[PriceReportOut], summary="上报更低价列表(筛选+分页)")
|
||||
def list_price_reports(
|
||||
db: AdminDb,
|
||||
@@ -117,50 +60,6 @@ def price_report_summary(db: AdminDb) -> PriceReportSummary:
|
||||
return PriceReportSummary.model_validate(queries.price_report_summary(db))
|
||||
|
||||
|
||||
@router.post("/bulk/approve", response_model=PriceReportBulkResult, summary="批量通过上报(发固定金币)")
|
||||
def bulk_approve_price_reports(
|
||||
body: PriceReportBulkRequest,
|
||||
request: Request,
|
||||
admin: Annotated[AdminUser, Depends(require_role("operator"))],
|
||||
db: AdminDb,
|
||||
) -> PriceReportBulkResult:
|
||||
results: list[PriceReportBulkItemResult] = []
|
||||
ip = get_client_ip(request)
|
||||
for report_id in body.ids:
|
||||
try:
|
||||
rep = _approve_price_report(db, admin, report_id, ip, bulk=True)
|
||||
results.append(PriceReportBulkItemResult(id=report_id, ok=True, status=rep.status))
|
||||
except HTTPException as exc:
|
||||
db.rollback()
|
||||
results.append(PriceReportBulkItemResult(id=report_id, ok=False, error=str(exc.detail)))
|
||||
except Exception: # noqa: BLE001 - 单笔失败不打断整批
|
||||
db.rollback()
|
||||
results.append(PriceReportBulkItemResult(id=report_id, ok=False, error="系统异常"))
|
||||
return _bulk_result(results)
|
||||
|
||||
|
||||
@router.post("/bulk/reject", response_model=PriceReportBulkResult, summary="批量拒绝上报")
|
||||
def bulk_reject_price_reports(
|
||||
body: PriceReportBulkRejectRequest,
|
||||
request: Request,
|
||||
admin: Annotated[AdminUser, Depends(require_role("operator"))],
|
||||
db: AdminDb,
|
||||
) -> PriceReportBulkResult:
|
||||
results: list[PriceReportBulkItemResult] = []
|
||||
ip = get_client_ip(request)
|
||||
for report_id in body.ids:
|
||||
try:
|
||||
rep = _reject_price_report(db, admin, report_id, body.reason, ip, bulk=True)
|
||||
results.append(PriceReportBulkItemResult(id=report_id, ok=True, status=rep.status))
|
||||
except HTTPException as exc:
|
||||
db.rollback()
|
||||
results.append(PriceReportBulkItemResult(id=report_id, ok=False, error=str(exc.detail)))
|
||||
except Exception: # noqa: BLE001 - 单笔失败不打断整批
|
||||
db.rollback()
|
||||
results.append(PriceReportBulkItemResult(id=report_id, ok=False, error="系统异常"))
|
||||
return _bulk_result(results)
|
||||
|
||||
|
||||
@router.post("/{report_id}/approve", response_model=OkResponse, summary="通过上报(发固定金币)")
|
||||
def approve_price_report(
|
||||
report_id: int,
|
||||
@@ -168,7 +67,27 @@ def approve_price_report(
|
||||
admin: Annotated[AdminUser, Depends(require_role("operator"))],
|
||||
db: AdminDb,
|
||||
) -> OkResponse:
|
||||
_approve_price_report(db, admin, report_id, get_client_ip(request))
|
||||
# 行锁(SELECT FOR UPDATE):并发/连点双请求会都读到 pending → 各发一次金币双倍发奖,
|
||||
# 锁住该行串行化,第二个请求拿锁后看到 approved → 走 400。SQLite 下 FOR UPDATE 为 no-op。
|
||||
rep = db.get(PriceReport, report_id, with_for_update=True)
|
||||
if rep is None:
|
||||
raise HTTPException(status_code=404, detail="上报记录不存在")
|
||||
if rep.status != "pending":
|
||||
raise HTTPException(status_code=400, detail=f"该上报已审核过(当前 {rep.status}),不可重复操作")
|
||||
coins = PRICE_REPORT_REWARD_COINS
|
||||
# 改状态 + 发金币 + 审计同一事务(commit=False),最后一起 commit:改了就有痕、发了就留账
|
||||
mutations.review_price_report(db, rep, status="approved", reward_coins=coins, commit=False)
|
||||
wallet_repo.grant_coins(
|
||||
db, rep.user_id, coins,
|
||||
biz_type="price_report_reward", ref_id=str(rep.id), remark="上报更低价审核通过",
|
||||
)
|
||||
write_audit(
|
||||
db, admin, action="price_report.approve", target_type="price_report", target_id=report_id,
|
||||
detail={"reward_coins": coins, "user_id": rep.user_id}, ip=get_client_ip(request), commit=False,
|
||||
)
|
||||
db.commit()
|
||||
# PRD #11 爆料审核通过:发金币后通知用户(站内 + push)。业务已 commit,通知失败只 log。
|
||||
notification_events.notify_report_approved(db, rep)
|
||||
return OkResponse()
|
||||
|
||||
|
||||
@@ -180,5 +99,16 @@ def reject_price_report(
|
||||
admin: Annotated[AdminUser, Depends(require_role("operator"))],
|
||||
db: AdminDb,
|
||||
) -> OkResponse:
|
||||
_reject_price_report(db, admin, report_id, body.reason, get_client_ip(request))
|
||||
rep = db.get(PriceReport, report_id, with_for_update=True) # 行锁,同 approve(防并发重复审核)
|
||||
if rep is None:
|
||||
raise HTTPException(status_code=404, detail="上报记录不存在")
|
||||
if rep.status != "pending":
|
||||
raise HTTPException(status_code=400, detail=f"该上报已审核过(当前 {rep.status}),不可重复操作")
|
||||
reason = body.reason.strip()
|
||||
mutations.review_price_report(db, rep, status="rejected", reject_reason=reason, commit=False)
|
||||
write_audit(
|
||||
db, admin, action="price_report.reject", target_type="price_report", target_id=report_id,
|
||||
detail={"reason": reason, "user_id": rep.user_id}, ip=get_client_ip(request), commit=False,
|
||||
)
|
||||
db.commit()
|
||||
return OkResponse()
|
||||
|
||||
@@ -3,7 +3,7 @@ from __future__ import annotations
|
||||
|
||||
from datetime import datetime
|
||||
|
||||
from pydantic import BaseModel, ConfigDict, Field, field_validator
|
||||
from pydantic import BaseModel, ConfigDict, Field
|
||||
|
||||
from app.core.rewards import FEEDBACK_REWARD_MAX_COINS
|
||||
|
||||
@@ -55,54 +55,6 @@ class FeedbackRejectRequest(BaseModel):
|
||||
reply: str | None = Field(default=None, max_length=256, description="给用户的回复留言,用户端可见")
|
||||
|
||||
|
||||
class FeedbackBulkRequest(BaseModel):
|
||||
ids: list[int] = Field(min_length=1, max_length=50, description="待审核反馈 ID 列表")
|
||||
|
||||
@field_validator("ids")
|
||||
@classmethod
|
||||
def _ids_must_be_unique(cls, ids: list[int]) -> list[int]:
|
||||
if len(ids) != len(set(ids)):
|
||||
raise ValueError("反馈 ID 不能重复")
|
||||
return ids
|
||||
|
||||
|
||||
class FeedbackBulkApproveRequest(FeedbackBulkRequest):
|
||||
reward_coins: int = Field(
|
||||
ge=1,
|
||||
le=FEEDBACK_REWARD_MAX_COINS,
|
||||
description="每条采纳反馈发放的金币数",
|
||||
)
|
||||
note: str | None = Field(default=None, max_length=256, description="采纳要点/审核备注(内部)")
|
||||
reply: str | None = Field(default=None, max_length=256, description="给用户的回复留言,用户端可见")
|
||||
|
||||
|
||||
class FeedbackBulkRejectRequest(FeedbackBulkRequest):
|
||||
reason: str = Field(min_length=1, max_length=256, description="批量未采纳原因,用户端可见")
|
||||
note: str | None = Field(default=None, max_length=256, description="运营内部审核备注")
|
||||
reply: str | None = Field(default=None, max_length=256, description="给用户的回复留言,用户端可见")
|
||||
|
||||
@field_validator("reason")
|
||||
@classmethod
|
||||
def _reason_not_blank(cls, value: str) -> str:
|
||||
if not value.strip():
|
||||
raise ValueError("未采纳原因不能为空")
|
||||
return value.strip()
|
||||
|
||||
|
||||
class FeedbackBulkItemResult(BaseModel):
|
||||
id: int
|
||||
ok: bool
|
||||
status: str | None = None
|
||||
error: str | None = None
|
||||
|
||||
|
||||
class FeedbackBulkResult(BaseModel):
|
||||
total: int
|
||||
success: int
|
||||
failed: int
|
||||
items: list[FeedbackBulkItemResult]
|
||||
|
||||
|
||||
class FeedbackSummary(BaseModel):
|
||||
"""审核台顶部各状态计数(pending 含历史 new 态)。"""
|
||||
|
||||
|
||||
@@ -56,42 +56,6 @@ class PriceReportRejectRequest(BaseModel):
|
||||
return v.strip()
|
||||
|
||||
|
||||
class PriceReportBulkRequest(BaseModel):
|
||||
ids: list[int] = Field(min_length=1, max_length=50, description="待审核上报 ID 列表")
|
||||
|
||||
@field_validator("ids")
|
||||
@classmethod
|
||||
def _ids_must_be_unique(cls, ids: list[int]) -> list[int]:
|
||||
if len(ids) != len(set(ids)):
|
||||
raise ValueError("上报 ID 不能重复")
|
||||
return ids
|
||||
|
||||
|
||||
class PriceReportBulkRejectRequest(PriceReportBulkRequest):
|
||||
reason: str = Field(min_length=1, max_length=256, description="批量拒绝理由,用户端记录页会看到")
|
||||
|
||||
@field_validator("reason")
|
||||
@classmethod
|
||||
def _reason_not_blank(cls, value: str) -> str:
|
||||
if not value.strip():
|
||||
raise ValueError("拒绝理由不能为空")
|
||||
return value.strip()
|
||||
|
||||
|
||||
class PriceReportBulkItemResult(BaseModel):
|
||||
id: int
|
||||
ok: bool
|
||||
status: str | None = None
|
||||
error: str | None = None
|
||||
|
||||
|
||||
class PriceReportBulkResult(BaseModel):
|
||||
total: int
|
||||
success: int
|
||||
failed: int
|
||||
items: list[PriceReportBulkItemResult]
|
||||
|
||||
|
||||
class PriceReportSummary(BaseModel):
|
||||
"""审核台顶部各状态计数。"""
|
||||
|
||||
|
||||
@@ -9,17 +9,15 @@ import asyncio
|
||||
import contextlib
|
||||
import logging
|
||||
import os
|
||||
import socket
|
||||
import time
|
||||
from collections.abc import Callable, Iterator
|
||||
from collections.abc import 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
|
||||
from app.db.session import SessionLocal
|
||||
|
||||
logger = logging.getLogger("shagua.cps_reconcile")
|
||||
_LOCK_PATH = Path(__file__).resolve().parents[2] / "data" / "cps_reconcile.lock"
|
||||
@@ -29,21 +27,6 @@ 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)
|
||||
@@ -84,240 +67,52 @@ def _single_instance_lock(stale_after_sec: int) -> Iterator[bool]:
|
||||
|
||||
|
||||
def _empty_result() -> dict:
|
||||
return {
|
||||
"fetched": 0,
|
||||
"inserted": 0,
|
||||
"updated": 0,
|
||||
"pages": 0,
|
||||
"api_requests": 0,
|
||||
}
|
||||
return {"fetched": 0, "inserted": 0, "updated": 0, "pages": 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:
|
||||
def _reconcile_once(now: datetime | None = None) -> 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:
|
||||
try:
|
||||
with SessionLocal() as db:
|
||||
return cps_repo.reconcile_orders(
|
||||
result["meituan"] = 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
|
||||
except Exception as exc: # noqa: BLE001 - 单平台失败不能阻断另一平台
|
||||
result["errors"]["meituan"] = str(exc)
|
||||
logger.exception("CPS auto reconcile failed platform=meituan")
|
||||
else:
|
||||
result["meituan"] = _skip_platform("meituan", common)
|
||||
result["meituan"] = {**_empty_result(), "skipped": "not_configured"}
|
||||
|
||||
_touch_lock()
|
||||
if settings.jd_union_configured:
|
||||
def reconcile_jd() -> dict:
|
||||
try:
|
||||
with SessionLocal() as db:
|
||||
return cps_repo.reconcile_jd_orders(
|
||||
result["jd"] = 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
|
||||
except Exception as exc: # noqa: BLE001 - 单平台失败不能阻断另一平台
|
||||
result["errors"]["jd"] = str(exc)
|
||||
logger.exception("CPS auto reconcile failed platform=jd")
|
||||
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,
|
||||
)
|
||||
result["jd"] = {**_empty_result(), "skipped": "not_configured"}
|
||||
|
||||
return result
|
||||
|
||||
@@ -332,90 +127,42 @@ 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 worker skipped: another worker owns lock",
|
||||
extra={
|
||||
"event": "cps_reconcile.worker_skipped",
|
||||
"skip_reason": "lock_not_acquired",
|
||||
"pid": os.getpid(),
|
||||
},
|
||||
)
|
||||
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:
|
||||
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,
|
||||
)
|
||||
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 worker error",
|
||||
extra={
|
||||
"event": "cps_reconcile.unexpected_failed",
|
||||
"run_id": run_id or "unassigned",
|
||||
"pid": os.getpid(),
|
||||
},
|
||||
)
|
||||
logger.exception("CPS auto reconcile unexpected error")
|
||||
await asyncio.sleep(interval)
|
||||
except asyncio.CancelledError:
|
||||
logger.info(
|
||||
"CPS auto reconcile worker stopped",
|
||||
extra={"event": "cps_reconcile.worker_stopped", "pid": os.getpid()},
|
||||
)
|
||||
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",
|
||||
extra={
|
||||
"event": "cps_reconcile.worker_skipped",
|
||||
"skip_reason": "disabled",
|
||||
},
|
||||
)
|
||||
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",
|
||||
extra={
|
||||
"event": "cps_reconcile.worker_skipped",
|
||||
"skip_reason": "all_platform_credentials_missing",
|
||||
},
|
||||
)
|
||||
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")
|
||||
|
||||
|
||||
@@ -14,7 +14,6 @@ from app.core.security import hash_password
|
||||
from app.db.session import SessionLocal
|
||||
from app.models.admin import AdminAuditLog
|
||||
from app.models.feedback import Feedback
|
||||
from app.models.price_report import PriceReport
|
||||
from app.models.user import User
|
||||
from app.models.wallet import CoinAccount, CoinTransaction, WithdrawOrder
|
||||
from app.repositories import user as user_repo
|
||||
@@ -83,26 +82,6 @@ def _seed_feedback(phone: str) -> int:
|
||||
db.close()
|
||||
|
||||
|
||||
def _seed_price_report(phone: str) -> int:
|
||||
uid = _seed_user(phone)
|
||||
db = SessionLocal()
|
||||
try:
|
||||
report = PriceReport(
|
||||
user_id=uid,
|
||||
store_name="测试门店",
|
||||
reported_platform_id="eleme",
|
||||
reported_platform_name="饿了么",
|
||||
reported_price_cents=2990,
|
||||
images=[],
|
||||
status="pending",
|
||||
)
|
||||
db.add(report)
|
||||
db.commit()
|
||||
return report.id
|
||||
finally:
|
||||
db.close()
|
||||
|
||||
|
||||
# ===== 调金币 =====
|
||||
|
||||
def test_grant_coins_writes_txn_and_audit(admin_client: TestClient, finance_token: str) -> None:
|
||||
@@ -418,100 +397,6 @@ def test_feedback_review_stores_admin_reply(
|
||||
assert r.json()["admin_reply"] == "已收到,后续跟进"
|
||||
|
||||
|
||||
def test_bulk_approve_feedbacks_returns_per_item_results(
|
||||
admin_client: TestClient, operator_token: str
|
||||
) -> None:
|
||||
first_id = _seed_feedback("13900000031")
|
||||
second_id = _seed_feedback("13900000032")
|
||||
r = admin_client.post(
|
||||
"/admin/api/feedbacks/bulk/approve",
|
||||
json={"ids": [first_id, second_id, 999999], "reward_coins": 600, "note": "批量采纳"},
|
||||
headers=_auth(operator_token),
|
||||
)
|
||||
assert r.status_code == 200, r.text
|
||||
payload = r.json()
|
||||
assert payload["total"] == 3 and payload["success"] == 2 and payload["failed"] == 1
|
||||
assert payload["items"][-1] == {"id": 999999, "ok": False, "status": None, "error": "反馈不存在"}
|
||||
db = SessionLocal()
|
||||
try:
|
||||
for feedback_id in (first_id, second_id):
|
||||
feedback = db.get(Feedback, feedback_id)
|
||||
assert feedback is not None and feedback.status == "adopted"
|
||||
assert db.get(CoinAccount, feedback.user_id).coin_balance == 600
|
||||
log = db.execute(
|
||||
select(AdminAuditLog).where(
|
||||
AdminAuditLog.action == "feedback.approve",
|
||||
AdminAuditLog.target_id == str(feedback_id),
|
||||
)
|
||||
).scalar_one()
|
||||
assert log.detail["bulk"] is True
|
||||
finally:
|
||||
db.close()
|
||||
|
||||
|
||||
def test_bulk_approve_price_reports_returns_per_item_results(
|
||||
admin_client: TestClient, operator_token: str
|
||||
) -> None:
|
||||
first_id = _seed_price_report("13900000041")
|
||||
second_id = _seed_price_report("13900000042")
|
||||
r = admin_client.post(
|
||||
"/admin/api/price-reports/bulk/approve",
|
||||
json={"ids": [first_id, second_id, 999999]},
|
||||
headers=_auth(operator_token),
|
||||
)
|
||||
assert r.status_code == 200, r.text
|
||||
payload = r.json()
|
||||
assert payload["total"] == 3 and payload["success"] == 2 and payload["failed"] == 1
|
||||
assert payload["items"][-1] == {"id": 999999, "ok": False, "status": None, "error": "上报记录不存在"}
|
||||
db = SessionLocal()
|
||||
try:
|
||||
for report_id in (first_id, second_id):
|
||||
report = db.get(PriceReport, report_id)
|
||||
assert report is not None and report.status == "approved"
|
||||
assert report.reward_coins == 1000
|
||||
assert db.get(CoinAccount, report.user_id).coin_balance == 1000
|
||||
log = db.execute(
|
||||
select(AdminAuditLog).where(
|
||||
AdminAuditLog.action == "price_report.approve",
|
||||
AdminAuditLog.target_id == str(report_id),
|
||||
)
|
||||
).scalar_one()
|
||||
assert log.detail["bulk"] is True
|
||||
finally:
|
||||
db.close()
|
||||
|
||||
|
||||
def test_bulk_reject_review_requests_apply_shared_reason(
|
||||
admin_client: TestClient, operator_token: str
|
||||
) -> None:
|
||||
feedback_id = _seed_feedback("13900000051")
|
||||
report_id = _seed_price_report("13900000052")
|
||||
feedback_response = admin_client.post(
|
||||
"/admin/api/feedbacks/bulk/reject",
|
||||
json={"ids": [feedback_id], "reason": "信息不足", "reply": "请补充完整截图"},
|
||||
headers=_auth(operator_token),
|
||||
)
|
||||
report_response = admin_client.post(
|
||||
"/admin/api/price-reports/bulk/reject",
|
||||
json={"ids": [report_id], "reason": "截图无法核实"},
|
||||
headers=_auth(operator_token),
|
||||
)
|
||||
assert feedback_response.status_code == 200, feedback_response.text
|
||||
assert report_response.status_code == 200, report_response.text
|
||||
assert feedback_response.json()["success"] == 1
|
||||
assert report_response.json()["success"] == 1
|
||||
db = SessionLocal()
|
||||
try:
|
||||
feedback = db.get(Feedback, feedback_id)
|
||||
report = db.get(PriceReport, report_id)
|
||||
assert feedback is not None and feedback.status == "rejected"
|
||||
assert feedback.reject_reason == "信息不足" and feedback.admin_reply == "请补充完整截图"
|
||||
assert report is not None and report.status == "rejected"
|
||||
assert report.reject_reason == "截图无法核实"
|
||||
finally:
|
||||
db.close()
|
||||
|
||||
|
||||
# ===== admin 账号管理(super_admin) =====
|
||||
|
||||
def test_create_and_update_admin(admin_client: TestClient, super_token: str) -> None:
|
||||
|
||||
@@ -4,28 +4,12 @@ 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,12 +7,8 @@ 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_summary_aggregate_stmt,
|
||||
coupon_data_report,
|
||||
)
|
||||
from app.admin.repositories.coupon_data import coupon_data_report
|
||||
from app.db.session import SessionLocal
|
||||
from app.models.coupon_state import CouponSession
|
||||
from app.repositories.coupon_state import (
|
||||
@@ -23,12 +19,7 @@ from app.repositories.coupon_state import (
|
||||
|
||||
|
||||
def _agg_session(
|
||||
trace: str,
|
||||
platforms,
|
||||
platform_success,
|
||||
*,
|
||||
status: str = "completed",
|
||||
elapsed_ms: int | None = None,
|
||||
trace: str, platforms, platform_success, *, status: str = "completed"
|
||||
) -> CouponSession:
|
||||
"""构造一条聚合测试用 session(started_date 固定 2020-01-02、app_env=prod,不 commit)。"""
|
||||
return CouponSession(
|
||||
@@ -40,53 +31,9 @@ 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,15 +1,11 @@
|
||||
"""美团、京东 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
|
||||
|
||||
|
||||
@@ -27,24 +23,11 @@ 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,
|
||||
"api_requests": 1,
|
||||
}
|
||||
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,
|
||||
"api_requests": 72,
|
||||
"windows": 72,
|
||||
}
|
||||
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)
|
||||
@@ -62,10 +45,6 @@ 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:
|
||||
@@ -89,9 +68,6 @@ 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:
|
||||
@@ -103,95 +79,6 @@ 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
|
||||
|
||||
Reference in New Issue
Block a user