Compare commits

...

3 Commits

Author SHA1 Message Date
unknown 9f51d51658 feat(admin): add bulk reward review APIs 2026-07-23 18:11:52 +08:00
marco ceceeb3458 fix(db): 合并 3 个 alembic 迁移 head(发车 0.4.4 前置,no-op merge 节点) 2026-07-23 17:55:35 +08:00
linkeyu 31f61f6aad 增强:CPS 每日自动对账增加可回查日志 (#163)
## 背景

PR #162 已合并。本 PR 为其后续日志增强,方便人工按一次任务完整回查美团、京东每日自动对账。

## 日志内容

- 每次运行生成唯一 `run_id`,记录 `scheduled` / `startup_catchup` 触发来源
- 记录北京时间计划日期、近 3 天回拉窗口、环境、数据库方言、主机名和 PID
- 分平台记录开始、成功、跳过、失败及执行耗时
- 成功结果记录 fetched / inserted / updated / pages / api_requests,京东额外记录小时窗口数
- 京东上游失败记录具体小时窗口、页码和请求序号;美团记录失败页码
- 最终汇总记录 success / partial_success / failed、失败平台、是否需要人工补跑和下次执行时间
- 日志使用现有 `extra` 结构化字段写入 JSON 日志,便于 SLS/人工检索

## 兼容性

- 手动对账接口、事务和返回模型不变
- 不记录密钥、Token、完整上游响应或订单明细
- 仓储层仅增加可选审计上下文和请求计数;不改变拉取及 upsert 逻辑

## 验证

- `pytest tests/test_cps_reconcile_worker.py tests/test_cps_admin.py tests/test_admin_read.py tests/test_observe.py -q`:44 passed
- 相关文件 Ruff 检查通过(忽略文件原有 UP017 提示)
- `git diff --check` 通过

---------

Co-authored-by: unknown <798648091@qq.com>
Reviewed-on: #163
Co-authored-by: linkeyu <linkeyu@wonderable.ai>
Co-committed-by: linkeyu <linkeyu@wonderable.ai>
2026-07-23 11:57:37 +08:00
9 changed files with 958 additions and 160 deletions
@@ -0,0 +1,26 @@
"""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
+71 -16
View File
@@ -5,6 +5,7 @@
"""
from __future__ import annotations
import logging
from datetime import datetime, timedelta, timezone
from decimal import Decimal, InvalidOperation
from typing import Any
@@ -15,12 +16,14 @@ from sqlalchemy.orm import Session
from app.admin.repositories.queries import _as_utc, offset_paginate
from app.integrations import jd_union, meituan
from app.repositories import cps_link as cps_link_repo
from app.models.cps_activity import CpsActivity
from app.models.cps_group import CpsGroup
from app.models.cps_link import CpsClick, CpsLink
from app.models.cps_order import CpsOrder
from app.models.cps_wx_user import CpsWxUser
from app.repositories import cps_link as cps_link_repo
logger = logging.getLogger("shagua.cps_reconcile")
# 美团订单状态:取消(4)/风控(5)不计佣金;结算(6)为佣金真正到账
_INVALID_STATUS = {"4", "5"}
@@ -378,18 +381,35 @@ def effective_commission_cents(order: CpsOrder) -> int:
def reconcile_orders(
db: Session, *, start_time: int, end_time: int,
query_time_type: int = 1, sid: str | None = None, max_pages: int = 200,
audit_context: dict[str, Any] | None = None,
) -> dict:
"""调美团 query_order 分页拉单 → 按 order_id upsert。返回 {fetched, inserted, updated, pages}。
订单状态会随时间变(付款→完成→结算/退款),重复拉同一单则更新。max_pages 防异常死循环。
"""
fetched = inserted = updated = pages = 0
fetched = inserted = updated = pages = api_requests = 0
page = 1
while page <= max_pages:
resp = meituan.query_order(
sid=sid, start_time=start_time, end_time=end_time,
query_time_type=query_time_type, page=page, limit=100,
)
api_requests += 1
try:
resp = meituan.query_order(
sid=sid, start_time=start_time, end_time=end_time,
query_time_type=query_time_type, page=page, limit=100,
)
except Exception: # noqa: BLE001 - 记录失败页后保持原异常类型继续抛出
if audit_context is not None:
logger.exception(
"CPS reconcile upstream request failed platform=meituan page=%s",
page,
extra={
**audit_context,
"event": "cps_reconcile.request_failed",
"platform": "meituan",
"failed_page": page,
"api_request_number": api_requests,
},
)
raise
rows = ((resp.get("data") or {}).get("dataList")) or []
if not rows:
break
@@ -419,30 +439,58 @@ def reconcile_orders(
break
page += 1
db.commit()
return {"fetched": fetched, "inserted": inserted, "updated": updated, "pages": pages}
return {
"fetched": fetched,
"inserted": inserted,
"updated": updated,
"pages": pages,
"api_requests": api_requests,
}
def reconcile_jd_orders(
db: Session, *, start_time: datetime, end_time: datetime,
query_time_type: int = 3, max_pages: int = 100,
audit_context: dict[str, Any] | None = None,
) -> dict:
"""调京东 order.row.query 拉单 → 按订单行 upsert。
京东单次查询窗口最多 1 小时,这里按北京自然时间切窗并逐页拉取。
"""
fetched = inserted = updated = pages = 0
fetched = inserted = updated = pages = api_requests = windows = 0
cur = start_time
while cur < end_time:
win_end = min(cur + timedelta(hours=1), end_time)
windows += 1
page = 1
while page <= max_pages:
resp = jd_union.query_order_rows(
start_time=cur,
end_time=win_end,
query_time_type=query_time_type,
page_index=page,
page_size=200,
)
api_requests += 1
try:
resp = jd_union.query_order_rows(
start_time=cur,
end_time=win_end,
query_time_type=query_time_type,
page_index=page,
page_size=200,
)
except Exception: # noqa: BLE001 - 记录失败窗口后保持原异常类型继续抛出
if audit_context is not None:
logger.exception(
"CPS reconcile upstream request failed platform=jd window=%s..%s page=%s",
cur.isoformat(),
win_end.isoformat(),
page,
extra={
**audit_context,
"event": "cps_reconcile.request_failed",
"platform": "jd",
"failed_window_start": cur.isoformat(),
"failed_window_end": win_end.isoformat(),
"failed_page": page,
"api_request_number": api_requests,
},
)
raise
rows = resp.get("rows") or []
has_more = bool(resp.get("has_more"))
if not rows:
@@ -469,7 +517,14 @@ def reconcile_jd_orders(
page += 1
cur = win_end
db.commit()
return {"fetched": fetched, "inserted": inserted, "updated": updated, "pages": pages}
return {
"fetched": fetched,
"inserted": inserted,
"updated": updated,
"pages": pages,
"api_requests": api_requests,
"windows": windows,
}
def list_orders(
+167 -85
View File
@@ -12,6 +12,10 @@ 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,
@@ -33,6 +37,123 @@ 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,
@@ -73,6 +194,50 @@ 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,
@@ -93,53 +258,7 @@ def approve_feedback(
admin: Annotated[AdminUser, Depends(require_role("operator"))],
db: AdminDb,
) -> FeedbackOut:
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
return _approve_feedback(db, admin, feedback_id, payload, get_client_ip(request))
@router.post("/{feedback_id}/reject", response_model=FeedbackOut, summary="拒绝采纳反馈")
@@ -150,41 +269,4 @@ def reject_feedback(
admin: Annotated[AdminUser, Depends(require_role("operator"))],
db: AdminDb,
) -> FeedbackOut:
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
return _reject_feedback(db, admin, feedback_id, payload, get_client_ip(request))
+103 -33
View File
@@ -17,6 +17,10 @@ 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,
@@ -34,6 +38,59 @@ 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,
@@ -60,6 +117,50 @@ 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,
@@ -67,27 +168,7 @@ def approve_price_report(
admin: Annotated[AdminUser, Depends(require_role("operator"))],
db: AdminDb,
) -> OkResponse:
# 行锁(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)
_approve_price_report(db, admin, report_id, get_client_ip(request))
return OkResponse()
@@ -99,16 +180,5 @@ def reject_price_report(
admin: Annotated[AdminUser, Depends(require_role("operator"))],
db: AdminDb,
) -> OkResponse:
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()
_reject_price_report(db, admin, report_id, body.reason, get_client_ip(request))
return OkResponse()
+49 -1
View File
@@ -3,7 +3,7 @@ from __future__ import annotations
from datetime import datetime
from pydantic import BaseModel, ConfigDict, Field
from pydantic import BaseModel, ConfigDict, Field, field_validator
from app.core.rewards import FEEDBACK_REWARD_MAX_COINS
@@ -55,6 +55,54 @@ 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 态)。"""
+36
View File
@@ -56,6 +56,42 @@ 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):
"""审核台顶部各状态计数。"""
+276 -23
View File
@@ -9,15 +9,17 @@ import asyncio
import contextlib
import logging
import os
import socket
import time
from collections.abc import Iterator
from collections.abc import Callable, Iterator
from datetime import date, datetime, timedelta
from pathlib import Path
from uuid import uuid4
from app.admin.repositories import cps as cps_repo
from app.core.config import settings
from app.core.rewards import CN_TZ
from app.db.session import SessionLocal
from app.db.session import SessionLocal, engine
logger = logging.getLogger("shagua.cps_reconcile")
_LOCK_PATH = Path(__file__).resolve().parents[2] / "data" / "cps_reconcile.lock"
@@ -27,6 +29,21 @@ def _cn_now() -> datetime:
return datetime.now(CN_TZ)
def _new_run_id(now: datetime) -> str:
return f"{now:%Y%m%d-%H%M%S}-{uuid4().hex[:8]}"
def _next_run_at(now: datetime, run_hour: int) -> datetime:
scheduled = now.replace(hour=run_hour, minute=0, second=0, microsecond=0)
return scheduled if now < scheduled else scheduled + timedelta(days=1)
def _trigger_for_run(worker_started_at: datetime, now: datetime, run_hour: int) -> str:
if worker_started_at.date() == now.date() and worker_started_at.hour >= run_hour:
return "startup_catchup"
return "scheduled"
def _touch_lock() -> None:
with contextlib.suppress(FileNotFoundError):
os.utime(_LOCK_PATH, None)
@@ -67,52 +84,240 @@ def _single_instance_lock(stale_after_sec: int) -> Iterator[bool]:
def _empty_result() -> dict:
return {"fetched": 0, "inserted": 0, "updated": 0, "pages": 0}
return {
"fetched": 0,
"inserted": 0,
"updated": 0,
"pages": 0,
"api_requests": 0,
}
def _reconcile_once(now: datetime | None = None) -> dict:
def _platform_log_fields(platform: str, result: dict) -> dict:
fields = {
"platform": platform,
"platform_status": result["status"],
"duration_ms": result["duration_ms"],
}
for key in ("fetched", "inserted", "updated", "pages", "api_requests", "windows"):
if key in result:
fields[key] = result[key]
return fields
def _run_platform(
*,
platform: str,
query_time_type: int,
common: dict,
reconcile: Callable[[], dict],
) -> tuple[dict, str | None]:
started = time.perf_counter()
logger.info(
"CPS auto reconcile platform started run_id=%s platform=%s",
common["run_id"],
platform,
extra={
**common,
"event": "cps_reconcile.platform_started",
"platform": platform,
"query_time_type": query_time_type,
},
)
try:
raw_result = reconcile()
except Exception as exc: # noqa: BLE001 - 单平台失败不能阻断另一平台
result = {
"status": "failed",
"duration_ms": round((time.perf_counter() - started) * 1000, 3),
"error_type": type(exc).__name__,
"error_summary": str(exc)[:500],
}
logger.exception(
"CPS auto reconcile platform failed run_id=%s platform=%s error_type=%s",
common["run_id"],
platform,
type(exc).__name__,
extra={
**common,
"event": "cps_reconcile.platform_failed",
**_platform_log_fields(platform, result),
"error_type": type(exc).__name__,
"error_summary": str(exc)[:500],
},
)
return result, str(exc)
result = {
**raw_result,
"status": "success",
"duration_ms": round((time.perf_counter() - started) * 1000, 3),
}
logger.info(
"CPS auto reconcile platform completed run_id=%s platform=%s fetched=%s inserted=%s updated=%s",
common["run_id"],
platform,
result.get("fetched", 0),
result.get("inserted", 0),
result.get("updated", 0),
extra={
**common,
"event": "cps_reconcile.platform_completed",
**_platform_log_fields(platform, result),
},
)
return result, None
def _skip_platform(platform: str, common: dict) -> dict:
result = {
**_empty_result(),
"status": "skipped",
"skipped": "not_configured",
"duration_ms": 0.0,
}
logger.warning(
"CPS auto reconcile platform skipped run_id=%s platform=%s reason=not_configured",
common["run_id"],
platform,
extra={
**common,
"event": "cps_reconcile.platform_skipped",
**_platform_log_fields(platform, result),
"skip_reason": "not_configured",
},
)
return result
def _reconcile_once(
now: datetime | None = None,
*,
run_id: str | None = None,
trigger: str = "scheduled",
) -> dict:
"""独立拉取美团和京东;单平台异常只记日志,不影响另一平台。"""
end = (now or _cn_now()).astimezone(CN_TZ)
run_id = run_id or _new_run_id(end)
lookback_days = max(1, int(settings.CPS_AUTO_RECONCILE_LOOKBACK_DAYS))
start = end - timedelta(days=lookback_days)
task_started = time.perf_counter()
common = {
"run_id": run_id,
"trigger": trigger,
"scheduled_date": end.date().isoformat(),
"app_env": settings.APP_ENV,
"db_dialect": engine.dialect.name,
"hostname": socket.gethostname(),
"pid": os.getpid(),
"window_start": start.isoformat(),
"window_end": end.isoformat(),
"lookback_days": lookback_days,
}
result = {
"run_id": run_id,
"trigger": trigger,
"window_start": start.isoformat(),
"window_end": end.isoformat(),
"meituan": None,
"jd": None,
"errors": {},
}
logger.info(
"CPS auto reconcile started run_id=%s trigger=%s window=%s..%s",
run_id,
trigger,
result["window_start"],
result["window_end"],
extra={**common, "event": "cps_reconcile.started"},
)
if settings.mt_cps_configured:
try:
def reconcile_meituan() -> dict:
with SessionLocal() as db:
result["meituan"] = cps_repo.reconcile_orders(
return cps_repo.reconcile_orders(
db,
start_time=int(start.timestamp()),
end_time=int(end.timestamp()),
query_time_type=2,
audit_context=common,
)
except Exception as exc: # noqa: BLE001 - 单平台失败不能阻断另一平台
result["errors"]["meituan"] = str(exc)
logger.exception("CPS auto reconcile failed platform=meituan")
result["meituan"], error = _run_platform(
platform="meituan",
query_time_type=2,
common=common,
reconcile=reconcile_meituan,
)
if error is not None:
result["errors"]["meituan"] = error
else:
result["meituan"] = {**_empty_result(), "skipped": "not_configured"}
result["meituan"] = _skip_platform("meituan", common)
_touch_lock()
if settings.jd_union_configured:
try:
def reconcile_jd() -> dict:
with SessionLocal() as db:
result["jd"] = cps_repo.reconcile_jd_orders(
return cps_repo.reconcile_jd_orders(
db,
start_time=start,
end_time=end,
query_time_type=3,
audit_context=common,
)
except Exception as exc: # noqa: BLE001 - 单平台失败不能阻断另一平台
result["errors"]["jd"] = str(exc)
logger.exception("CPS auto reconcile failed platform=jd")
result["jd"], error = _run_platform(
platform="jd",
query_time_type=3,
common=common,
reconcile=reconcile_jd,
)
if error is not None:
result["errors"]["jd"] = error
else:
result["jd"] = {**_empty_result(), "skipped": "not_configured"}
result["jd"] = _skip_platform("jd", common)
successes = sum(
platform_result.get("status") == "success"
for platform_result in (result["meituan"], result["jd"])
)
if result["errors"]:
task_status = "partial_success" if successes else "failed"
else:
task_status = "success"
completed_at = _cn_now()
duration_ms = round((time.perf_counter() - task_started) * 1000, 3)
next_run_at = _next_run_at(
completed_at,
min(max(int(settings.CPS_AUTO_RECONCILE_RUN_HOUR), 0), 23),
).isoformat()
result.update({
"status": task_status,
"duration_ms": duration_ms,
"manual_retry_required": bool(result["errors"]),
"next_run_at": next_run_at,
})
completion_fields = {
**common,
"event": "cps_reconcile.completed",
"task_status": task_status,
"duration_ms": duration_ms,
"manual_retry_required": result["manual_retry_required"],
"failed_platforms": sorted(result["errors"]),
"next_run_at": next_run_at,
"meituan_result": result["meituan"],
"jd_result": result["jd"],
}
log_method = logger.info if task_status == "success" else logger.warning
log_method(
"CPS auto reconcile completed run_id=%s status=%s duration_ms=%s failed_platforms=%s next_run_at=%s",
run_id,
task_status,
duration_ms,
",".join(sorted(result["errors"])) or "-",
next_run_at,
extra=completion_fields,
)
return result
@@ -127,42 +332,90 @@ async def _run_loop() -> None:
lock_stale_after = max(interval * 3, 1800)
with _single_instance_lock(lock_stale_after) as lock_acquired:
if not lock_acquired:
logger.warning("CPS auto reconcile skipped: another worker owns lock")
logger.warning(
"CPS auto reconcile worker skipped: another worker owns lock",
extra={
"event": "cps_reconcile.worker_skipped",
"skip_reason": "lock_not_acquired",
"pid": os.getpid(),
},
)
return
await _run_locked_loop(interval, run_hour)
async def _run_locked_loop(interval: int, run_hour: int) -> None:
worker_started_at = _cn_now()
logger.info(
"CPS auto reconcile worker started run_hour=%s interval=%ss lookback_days=%s",
run_hour,
interval,
settings.CPS_AUTO_RECONCILE_LOOKBACK_DAYS,
extra={
"event": "cps_reconcile.worker_started",
"pid": os.getpid(),
"hostname": socket.gethostname(),
"app_env": settings.APP_ENV,
"db_dialect": engine.dialect.name,
"run_hour": run_hour,
"interval_sec": interval,
"lookback_days": settings.CPS_AUTO_RECONCILE_LOOKBACK_DAYS,
"next_run_at": _next_run_at(worker_started_at, run_hour).isoformat(),
},
)
last_run: date | None = None
try:
while True:
run_id: str | None = None
try:
_touch_lock()
now = _cn_now()
if _should_run(last_run, now, run_hour):
result = await asyncio.to_thread(_reconcile_once, now)
run_id = _new_run_id(now)
trigger = _trigger_for_run(worker_started_at, now, run_hour)
await asyncio.to_thread(
_reconcile_once,
now,
run_id=run_id,
trigger=trigger,
)
last_run = now.date()
logger.info("CPS auto reconcile done date=%s result=%s", last_run, result)
except Exception: # noqa: BLE001 - 后台任务不能因单次异常退出
logger.exception("CPS auto reconcile unexpected error")
logger.exception(
"CPS auto reconcile unexpected worker error",
extra={
"event": "cps_reconcile.unexpected_failed",
"run_id": run_id or "unassigned",
"pid": os.getpid(),
},
)
await asyncio.sleep(interval)
except asyncio.CancelledError:
logger.info("CPS auto reconcile worker stopped")
logger.info(
"CPS auto reconcile worker stopped",
extra={"event": "cps_reconcile.worker_stopped", "pid": os.getpid()},
)
raise
def start_cps_reconcile_worker() -> asyncio.Task | None:
if not settings.CPS_AUTO_RECONCILE_ENABLED:
logger.info("CPS auto reconcile disabled")
logger.info(
"CPS auto reconcile disabled",
extra={
"event": "cps_reconcile.worker_skipped",
"skip_reason": "disabled",
},
)
return None
if not settings.mt_cps_configured and not settings.jd_union_configured:
logger.warning("CPS auto reconcile not started: Meituan and JD credentials are missing")
logger.warning(
"CPS auto reconcile not started: Meituan and JD credentials are missing",
extra={
"event": "cps_reconcile.worker_skipped",
"skip_reason": "all_platform_credentials_missing",
},
)
return None
return asyncio.create_task(_run_loop(), name="cps-auto-reconcile")
+115
View File
@@ -14,6 +14,7 @@ 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
@@ -82,6 +83,26 @@ 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:
@@ -397,6 +418,100 @@ 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:
+115 -2
View File
@@ -1,11 +1,15 @@
"""美团、京东 CPS 每日自动对账 worker。"""
from __future__ import annotations
import logging
from datetime import date, datetime, timedelta
import pytest
from app.core import cps_reconcile_worker as worker
from app.core.config import settings
from app.core.rewards import CN_TZ
from app.db.session import SessionLocal
from app.integrations.meituan import MeituanCpsError
@@ -23,11 +27,24 @@ def test_reconcile_once_pulls_meituan_and_jd_by_update_time(monkeypatch) -> None
def fake_meituan(db, **kwargs):
calls["meituan"] = kwargs
return {"fetched": 2, "inserted": 1, "updated": 1, "pages": 1}
return {
"fetched": 2,
"inserted": 1,
"updated": 1,
"pages": 1,
"api_requests": 1,
}
def fake_jd(db, **kwargs):
calls["jd"] = kwargs
return {"fetched": 3, "inserted": 2, "updated": 1, "pages": 2}
return {
"fetched": 3,
"inserted": 2,
"updated": 1,
"pages": 2,
"api_requests": 72,
"windows": 72,
}
monkeypatch.setattr(worker.cps_repo, "reconcile_orders", fake_meituan)
monkeypatch.setattr(worker.cps_repo, "reconcile_jd_orders", fake_jd)
@@ -45,6 +62,10 @@ def test_reconcile_once_pulls_meituan_and_jd_by_update_time(monkeypatch) -> None
assert calls["meituan"]["start_time"] == int((now - timedelta(days=3)).timestamp())
assert calls["jd"]["start_time"] == now - timedelta(days=3)
assert calls["jd"]["end_time"] == now
assert calls["meituan"]["audit_context"]["run_id"] == result["run_id"]
assert calls["jd"]["audit_context"]["run_id"] == result["run_id"]
assert result["status"] == "success"
assert result["manual_retry_required"] is False
def test_meituan_failure_does_not_block_jd(monkeypatch) -> None:
@@ -68,6 +89,9 @@ def test_meituan_failure_does_not_block_jd(monkeypatch) -> None:
assert result["errors"]["meituan"] == "temporary failure"
assert jd_called is True
assert result["jd"]["fetched"] == 1
assert result["meituan"]["status"] == "failed"
assert result["status"] == "partial_success"
assert result["manual_retry_required"] is True
def test_should_run_once_after_five_am() -> None:
@@ -79,6 +103,95 @@ def test_should_run_once_after_five_am() -> None:
assert worker._should_run(date(2026, 7, 22), at_five, 5) is False
def test_trigger_and_next_run_are_beijing_time() -> None:
before_five = datetime(2026, 7, 22, 4, 30, tzinfo=CN_TZ)
after_five = datetime(2026, 7, 22, 6, 0, tzinfo=CN_TZ)
assert worker._trigger_for_run(before_five, after_five, 5) == "scheduled"
assert worker._trigger_for_run(after_five, after_five, 5) == "startup_catchup"
assert worker._next_run_at(before_five, 5) == datetime(2026, 7, 22, 5, 0, tzinfo=CN_TZ)
assert worker._next_run_at(after_five, 5) == datetime(2026, 7, 23, 5, 0, tzinfo=CN_TZ)
def test_reconcile_emits_structured_audit_logs(monkeypatch, caplog) -> None:
_configure_platforms(monkeypatch)
monkeypatch.setattr(
worker.cps_repo,
"reconcile_orders",
lambda db, **kwargs: {
"fetched": 2,
"inserted": 1,
"updated": 1,
"pages": 1,
"api_requests": 1,
},
)
monkeypatch.setattr(
worker.cps_repo,
"reconcile_jd_orders",
lambda db, **kwargs: {
"fetched": 3,
"inserted": 2,
"updated": 1,
"pages": 1,
"api_requests": 72,
"windows": 72,
},
)
monkeypatch.setattr(worker, "_touch_lock", lambda: None)
monkeypatch.setattr(
worker,
"_cn_now",
lambda: datetime(2026, 7, 22, 5, 1, tzinfo=CN_TZ),
)
with caplog.at_level(logging.INFO, logger=worker.logger.name):
result = worker._reconcile_once(
datetime(2026, 7, 22, 5, 0, tzinfo=CN_TZ),
run_id="audit-run-1",
trigger="scheduled",
)
records = {record.event: record for record in caplog.records if hasattr(record, "event")}
assert records["cps_reconcile.started"].run_id == "audit-run-1"
assert records["cps_reconcile.started"].lookback_days == 3
assert records["cps_reconcile.started"].scheduled_date == "2026-07-22"
assert records["cps_reconcile.started"].hostname
assert records["cps_reconcile.completed"].task_status == "success"
assert records["cps_reconcile.completed"].manual_retry_required is False
assert records["cps_reconcile.completed"].meituan_result["fetched"] == 2
assert records["cps_reconcile.completed"].jd_result["api_requests"] == 72
assert result["next_run_at"] == "2026-07-23T05:00:00+08:00"
def test_jd_request_failure_logs_exact_window(monkeypatch, caplog) -> None:
start = datetime(2026, 7, 20, 5, 0, tzinfo=CN_TZ)
def fail_query(**kwargs):
raise RuntimeError("upstream timeout")
monkeypatch.setattr(worker.cps_repo.jd_union, "query_order_rows", fail_query)
with SessionLocal() as db, caplog.at_level(logging.ERROR, logger=worker.logger.name):
with pytest.raises(RuntimeError, match="upstream timeout"):
worker.cps_repo.reconcile_jd_orders(
db,
start_time=start,
end_time=start + timedelta(hours=2),
audit_context={"run_id": "audit-run-2", "trigger": "scheduled"},
)
record = next(
item
for item in caplog.records
if getattr(item, "event", "") == "cps_reconcile.request_failed"
)
assert record.run_id == "audit-run-2"
assert record.platform == "jd"
assert record.failed_window_start == "2026-07-20T05:00:00+08:00"
assert record.failed_window_end == "2026-07-20T06:00:00+08:00"
assert record.failed_page == 1
def test_start_worker_respects_auto_switch(monkeypatch) -> None:
monkeypatch.setattr(settings, "CPS_AUTO_RECONCILE_ENABLED", False)
assert worker.start_cps_reconcile_worker() is None