f1147dced3
- app/pay 独立支付子应用(config/wxpay/models/repository/service/routers),便于后期拆分微服务 - 充值订单表 0026_compute_recharge;套餐支持折扣(discount)与上架开关(enabled),DB system_configs 优先、env 兜底 - 统一小程序支付:桌面端出小程序码→扫码进小程序确认页按 openid 发起 JSAPI→回调幂等到账(引擎 adjust_user_quota+镜像回写) - 内嵌 wechatpayv3(含 async_),响应验签失败降级告警;新增 aiofiles 依赖 Co-Authored-By: Claude <noreply@anthropic.com>
131 lines
5.6 KiB
Python
131 lines
5.6 KiB
Python
# -*- coding: utf-8 -*-
|
||
"""支付子应用接口层:微信回调 + 算力充值(C 端自服务)。
|
||
|
||
- ``POST /opc/pay/notify`` 微信支付回调(公网、无登录鉴权,靠 V3 验签)。
|
||
- ``/opc/compute/recharge/*`` 用户充值(套餐/下单/状态/记录),require_roles("opc_member")。
|
||
"""
|
||
from __future__ import annotations
|
||
|
||
import json
|
||
import logging
|
||
|
||
from fastapi import APIRouter, Depends, HTTPException, Request, Response
|
||
from pydantic import BaseModel, Field
|
||
|
||
from ..api.dependencies import get_db
|
||
from ..infrastructure.repositories import Database
|
||
from ..rbac import require_roles
|
||
from . import service, wxpay
|
||
from .repository import ComputeRechargeOrderRepository
|
||
|
||
logger = logging.getLogger("pay.api")
|
||
|
||
router = APIRouter(tags=["pay"])
|
||
|
||
|
||
def _repo(db: Database) -> ComputeRechargeOrderRepository:
|
||
"""订单仓储(复用请求级 session;支付域自持仓储,便于整域拆分微服务)。"""
|
||
return ComputeRechargeOrderRepository(db.session)
|
||
|
||
|
||
# ── 微信支付回调(微信服务器调用,验签解密后到账)─────────────────────────────
|
||
@router.post("/opc/pay/notify", summary="微信支付回调(无登录鉴权,V3 验签)")
|
||
async def pay_notify(request: Request, db: Database = Depends(get_db)):
|
||
body = await request.body()
|
||
result = await wxpay.verify_and_decrypt(request.headers, body)
|
||
if not result or not isinstance(result, dict):
|
||
raise HTTPException(status_code=400, detail="回调验签/解密失败")
|
||
# 展平 resource 到顶层(out_trade_no / transaction_id / trade_state ...)
|
||
resource = result.pop("resource", {})
|
||
if isinstance(resource, dict):
|
||
result.update(resource)
|
||
order_no = result.get("out_trade_no", "")
|
||
logger.info("微信支付回调: order=%s trade_state=%s", order_no, result.get("trade_state"))
|
||
try:
|
||
ok = await service.handle_notify(db, result)
|
||
except Exception as exc: # noqa: BLE001
|
||
logger.error("回调处理异常 order=%s: %s", order_no, exc, exc_info=True)
|
||
raise HTTPException(status_code=500, detail="处理失败,等待微信重试") from exc
|
||
if not ok:
|
||
raise HTTPException(status_code=500, detail="处理失败,等待微信重试")
|
||
return {"code": "SUCCESS", "message": "成功"}
|
||
|
||
|
||
# ── 用户充值(C 端自服务)─────────────────────────────────────────────────────
|
||
class RechargeOrderRequest(BaseModel):
|
||
"""充值下单请求:package_id(套餐)或 amount_yuan(自由金额)二选一。"""
|
||
client_type: str = Field(default="native", description="native(桌面扫码) | jsapi(小程序)")
|
||
package_id: str = ""
|
||
amount_yuan: float | None = Field(default=None, ge=0)
|
||
|
||
|
||
@router.get("/opc/compute/recharge/packages", summary="充值套餐与支付开关")
|
||
async def recharge_packages(
|
||
db: Database = Depends(get_db),
|
||
_u: dict = Depends(require_roles("opc_member")),
|
||
):
|
||
return await service.packages(db)
|
||
|
||
|
||
@router.post("/opc/compute/recharge/orders", summary="创建充值订单(native=二维码 / jsapi=小程序支付参数)")
|
||
async def recharge_create(
|
||
body: RechargeOrderRequest,
|
||
db: Database = Depends(get_db),
|
||
user: dict = Depends(require_roles("opc_member")),
|
||
):
|
||
try:
|
||
order = await service.create_order(
|
||
db, user, client_type=body.client_type,
|
||
package_id=body.package_id, amount_yuan=body.amount_yuan,
|
||
)
|
||
except ValueError as exc:
|
||
raise HTTPException(status_code=400, detail=str(exc)) from exc
|
||
except RuntimeError as exc:
|
||
raise HTTPException(status_code=503, detail=str(exc)) from exc
|
||
return order
|
||
|
||
|
||
@router.get("/opc/compute/recharge/orders/{order_no}/pay-params", summary="小程序扫码确认支付(归属校验 + JSAPI 参数)")
|
||
async def recharge_pay_params(
|
||
order_no: str,
|
||
db: Database = Depends(get_db),
|
||
user: dict = Depends(require_roles("opc_member")),
|
||
):
|
||
"""微信扫桌面端小程序码 → 打开小程序本接口取 wx.requestPayment 参数。"""
|
||
try:
|
||
return await service.build_pay_params(db, user, order_no)
|
||
except ValueError as exc:
|
||
raise HTTPException(status_code=400, detail=str(exc)) from exc
|
||
except RuntimeError as exc:
|
||
raise HTTPException(status_code=503, detail=str(exc)) from exc
|
||
|
||
|
||
@router.post("/opc/compute/recharge/orders/{order_no}/status", summary="查询充值订单状态(含主动对账补单)")
|
||
async def recharge_status(
|
||
order_no: str,
|
||
db: Database = Depends(get_db),
|
||
user: dict = Depends(require_roles("opc_member")),
|
||
):
|
||
repo = _repo(db)
|
||
order = await repo.get_by_order_no(order_no)
|
||
if order is None or order["user_id"] != user["id"]:
|
||
raise HTTPException(status_code=404, detail="订单不存在")
|
||
order = await service.reconcile(db, order)
|
||
return {
|
||
"order_no": order["order_no"], "status": order["status"],
|
||
"paid": order["status"] in ("paid", "credited"),
|
||
"credited": order["status"] == "credited",
|
||
"amount_fen": order["amount_fen"], "quota_micro": order["quota_micro"],
|
||
"bonus_quota_micro": order["bonus_quota_micro"],
|
||
"transaction_id": order["transaction_id"], "expires_at": order["expires_at"],
|
||
"paid_at": order.get("paid_at", ""), "credited_at": order.get("credited_at", ""),
|
||
}
|
||
|
||
|
||
@router.get("/opc/compute/recharge/orders", summary="我的充值记录")
|
||
async def recharge_orders(
|
||
db: Database = Depends(get_db),
|
||
user: dict = Depends(require_roles("opc_member")),
|
||
):
|
||
return {"items": await _repo(db).list_by_user(user["id"])}
|