Files
server-core/app/pay/service.py
T
Pine 0fd210ac1f feat(pay): 微信服务商分账+资金托管改造——OPC个人/商户双绑定收款、验收95/5分账、回调确认、escrow资金托管字段扩展
- profitsharing.py:服务商分账客户端(添加接收方/发起分账/查分账/退款/回调验签)
- payment_bindings + settlement_logs:OPC 收款绑定(personal openid / merchant 子商户号)与分账流水镜像
- settlement_service.release_escrow:验收后优先发起微信分账(95%接单者/5%平台),未配置时降级台账结算(channel=manual)
- handle_split_notify:分账回调确认 escrow released(幂等)
- 路由:分账回调 + 绑定创建/列表/停用
- 迁移:alembic 0039(新表 + escrows 幂等加列)+ 兜底迁移脚本 + 冒烟测试
2026-09-01 18:39:41 +08:00

434 lines
20 KiB
Python
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
# -*- coding: utf-8 -*-
"""充值业务服务:下单 / 回调到账 / 对账补单 / 幂等。
到账链路(与全平台算力统一口径一致):
支付成功 → compute_engine users.quota 增加(1 元 = 1,000,000 微元)→ 平台镜像回写。
幂等三重保证:Redis SETNX 锁(30s+ ``mark_paid`` 单向状态机 + 订单终态判定。
事务边界:引擎到账是跨服务 HTTP,无法与本地 DB 同事务 → 先本地置 paid 终态,
再引擎加款;引擎失败置 failed,由 status 对账接口补单(不重复加款)。
"""
from __future__ import annotations
import json
import logging
import secrets
import time
from datetime import datetime, timedelta, timezone
from .. import config
from ..infrastructure.repositories import Database, utcnow_iso
from ..services import compute_client
from . import config as pay_config
from . import profitsharing, wxpay
from .repository import PaymentBindingRepository
logger = logging.getLogger("pay.service")
# 订单号前缀(回调按前缀分发,为未来其他支付业务留扩展位)
ORDER_PREFIX = "CR_" # 算力充值
ORDER_PREFIX_EVENT = "EV_" # 活动报名(课程/沙龙定价,见 event_service
# 套餐缓存(60s:DB 配置优先,env 兜底)
_pkg_cache: tuple[float, list[dict]] | None = None
_PKG_TTL = 60.0
def _iso_in(seconds: int) -> str:
return datetime.fromtimestamp(
datetime.now(timezone.utc).timestamp() + seconds, tz=timezone.utc,
).isoformat()
# ---------------------------------------------------------------------------
# 套餐
# ---------------------------------------------------------------------------
def _normalize_packages(raw) -> list[dict]:
"""套餐数据清洗:[{id,amount(元),bonus(元),discount(%),enabled,label}],金额全整数分口径。
- discount:充值折扣(默认 100=无折扣),实付 = amount * discount / 100
- enabled:false 时不下发(运营端可临时下架)
"""
items: list[dict] = []
if not isinstance(raw, list):
return items
for it in raw:
if isinstance(it, dict) and it.get("enabled") is False:
continue
try:
amount = int(round(float(it.get("amount", 0)) * 100))
bonus = int(round(float(it.get("bonus", 0) or 0) * 100))
discount = int(it.get("discount", 100) or 100)
except (TypeError, ValueError):
continue
if amount <= 0:
continue
discount = max(1, min(100, discount))
items.append({
"id": str(it.get("id") or f"p{amount}"),
"amount": amount, # 分(原价)
"bonus": bonus, # 分(赠送折算金额)
"discount": discount, # %100=原价)
"pay_fen": amount * discount // 100, # 分(实付)
"label": str(it.get("label") or f"{amount // 100}"),
})
return items
async def packages(db: Database) -> dict:
"""充值套餐(system_configs.recharge_packages 优先,env 兜底,60s 缓存)+ 支付开关。"""
global _pkg_cache
now = time.monotonic()
items: list[dict] | None = None
if _pkg_cache and now - _pkg_cache[0] < _PKG_TTL:
items = _pkg_cache[1]
else:
try:
raw = await db.config.get("recharge_packages")
if raw:
items = _normalize_packages(json.loads(raw))
except (ValueError, TypeError): # noqa: PERF203
items = None
if not items:
try:
items = _normalize_packages(json.loads(pay_config.DEFAULT_PACKAGES_JSON))
except (ValueError, TypeError):
items = []
_pkg_cache = (now, items)
enabled = pay_config.pay_enabled()
return {
"enabled": enabled,
"items": [
{
"id": p["id"], "amount": p["amount"] // 100, "bonus": p["bonus"] // 100,
"discount": p["discount"], "pay": p["pay_fen"] // 100,
"label": p["label"],
"quota_micro": (p["pay_fen"] + p["bonus"]) // 100 * pay_config.QUOTA_PER_YUAN,
}
for p in items
] if enabled else [],
"min_yuan": pay_config.RECHARGE_MIN_YUAN,
"max_yuan": pay_config.RECHARGE_MAX_YUAN,
}
# ---------------------------------------------------------------------------
# 下单
# ---------------------------------------------------------------------------
async def create_order(db: Database, user: dict, *, client_type: str,
package_id: str = "", amount_yuan: float | None = None) -> dict:
"""创建充值订单。
client_type
- "mp"(小程序支付,**统一入口**):不在下单时调微信——桌面端展示小程序码,
微信扫码自动打开小程序「确认支付」页,由小程序侧按 openid 发起 JSAPI 支付
(与「小程序扫码登录」同构:桌面出码 → 小程序内确认/支付 → 桌面轮询状态)。
- "jsapi"(小程序内直充):下单即返回 wx.requestPayment 参数。
- "native"(保留:桌面直接出微信收款码,需公众号 appid 绑定商户号)。
"""
if client_type not in ("mp", "native", "jsapi"):
raise ValueError("client_type 仅支持 mp/jsapi/native")
if not pay_config.pay_enabled():
raise RuntimeError("支付未配置,暂不可用")
# 1) 金额:套餐固定价,或自由金额(整数分口径)
plist = (await packages(db))["items"]
if package_id:
pkg = next((p for p in plist if p["id"] == package_id), None)
if pkg is None:
raise ValueError("套餐不存在或已下架")
# 实付按折扣计算(discount 100=原价);packages() 返回的 amount/bonus 单位为元
amount_fen = pkg["amount"] * pkg["discount"] # 元 × % → 分
bonus_quota_micro = pkg["bonus"] * pay_config.QUOTA_PER_YUAN
else:
amount_fen = int(round(float(amount_yuan or 0) * 100))
bonus_quota_micro = 0
if amount_fen <= 0:
raise ValueError("金额必须大于 0")
min_fen = int(round(pay_config.RECHARGE_MIN_YUAN * 100))
max_fen = int(round(pay_config.RECHARGE_MAX_YUAN * 100))
if not (min_fen <= amount_fen <= max_fen):
raise ValueError(f"金额需在 {pay_config.RECHARGE_MIN_YUAN:g} ~ {pay_config.RECHARGE_MAX_YUAN:g} 元之间")
# 2) 复用未过期同金额 pending 单(防重复扫码)
now_iso = utcnow_iso()
dup = await db.compute_recharges.find_pending_same_amount(user["id"], amount_fen, now_iso)
if dup and dup.get("code_url"):
return _order_view(dup)
# 3) 引擎用户(幂等建号)并解析 engine_user_id
username = user.get("username", "")
await compute_client.ensure_user(username)
engine_user_id = await _resolve_engine_user_id(username)
# 4) 下单:"mp" 只落库(支付由小程序扫码后发起);native/jsapi 此处即调微信
order_no = f"{ORDER_PREFIX}{int(time.time())}_{secrets.token_hex(4).upper()}"
expires_at = _iso_in(pay_config.RECHARGE_EXPIRE_MINUTES * 60)
wx_result: dict = {}
if client_type != "mp":
try:
wx_result = await wxpay.create_order(
client_type=client_type, out_trade_no=order_no, total_fen=amount_fen,
description=f"算力充值 - {amount_fen // 100}", openid=user.get("wx_mini_openid", ""),
)
except RuntimeError as exc:
logger.error("充值下单失败 user=%s: %s", username, exc)
raise ValueError(f"微信下单失败:{exc}") from exc
order = await db.compute_recharges.create({
"order_no": order_no, "user_id": user["id"], "username": username,
"engine_user_id": engine_user_id, "package_id": package_id,
"amount_fen": amount_fen,
"quota_micro": amount_fen // 100 * pay_config.QUOTA_PER_YUAN,
"bonus_quota_micro": bonus_quota_micro,
"client_type": client_type, "appid": wx_result.get("appid", ""),
"prepay_id": wx_result.get("prepay_id", ""), "code_url": wx_result.get("code_url", ""),
"expires_at": expires_at,
})
view = _order_view(order)
if client_type == "jsapi":
view["pay_params"] = wx_result.get("pay_params") # wx.requestPayment 直接参数
if client_type == "mp":
# 小程序码:微信扫码自动打开小程序「确认支付」页并携带 scene=order_no
# order_no 22 字符,满足 wxacode scene ≤32 限制)
from ..services import wechat
try:
png = await wechat.get_wxacode(order_no, page="pages/pay/index")
import base64 as _b64
view["qr_image"] = f"data:image/png;base64,{_b64.b64encode(png).decode('ascii')}"
except wechat.WechatError as exc:
logger.error("生成小程序码失败 order=%s: %s", order_no, exc)
view["qr_image"] = ""
return view
async def build_pay_params(db: Database, user: dict, order_no: str) -> dict:
"""小程序扫码进入「确认支付」页:校验归属后按当前小程序用户 openid 发起 JSAPI 下单,
返回订单摘要 + wx.requestPayment 参数。"""
order = await db.compute_recharges.get_by_order_no(order_no)
if order is None:
raise ValueError("订单不存在")
if order["user_id"] != user["id"]:
raise ValueError("请使用下单账号登录的小程序扫码(订单归属不一致)")
if order["status"] != "pending":
return {"order": _order_view(order), "pay_params": None, "finished": True}
if order["expires_at"] and order["expires_at"] < utcnow_iso():
await db.compute_recharges.mark_status(order_no, "closed")
return {"order": await db.compute_recharges.get_by_order_no(order_no), "pay_params": None, "finished": True}
openid = user.get("wx_mini_openid", "")
if not openid:
raise ValueError("当前账号未绑定小程序微信身份,请先用微信登录小程序")
try:
result = await wxpay.create_order(
client_type="jsapi", out_trade_no=order_no, total_fen=int(order["amount_fen"]),
description=f"算力充值 - {int(order['amount_fen']) // 100}", openid=openid,
)
except RuntimeError as exc:
logger.error("小程序拉起支付失败 order=%s: %s", order_no, exc)
raise ValueError(f"微信下单失败:{exc}") from exc
# 记录 prepay/appid(不改变 pending 状态)
await db.compute_recharges.set_prepay(order_no, result.get("appid", ""), result.get("prepay_id", ""))
return {
"order": _order_view(order),
"pay_params": result.get("pay_params"),
"finished": False,
}
def _order_view(order: dict) -> dict:
"""对外的订单视图(不含内部载荷)。"""
return {
"order_no": order["order_no"], "amount_fen": order["amount_fen"],
"quota_micro": order["quota_micro"], "bonus_quota_micro": order["bonus_quota_micro"],
"status": order["status"], "code_url": order["code_url"],
"prepay_id": order.get("prepay_id", ""), "expires_at": order["expires_at"],
"created_at": order["created_at"], "paid_at": order.get("paid_at", ""),
}
async def _resolve_engine_user_id(username: str) -> int:
"""按用户名在引擎侧解析用户 id(失败返回 0,到账时再试)。"""
try:
for page in (1, 2):
for it in await compute_client.list_engine_users(page=page, page_size=100):
if str(it.get("username") or "") == username:
return int(it.get("id") or 0)
except Exception as exc: # noqa: BLE001
logger.warning("解析引擎用户 id 失败 %s: %s", username, exc)
return 0
# ---------------------------------------------------------------------------
# 回调 / 到账
# ---------------------------------------------------------------------------
async def handle_notify(db: Database, result: dict) -> bool:
"""微信回调分发:CR_=充值到账,EV_=活动报名确认,MK_=市场购买;未知前缀忽略。"""
out_trade_no = result.get("out_trade_no", "")
if out_trade_no.startswith(ORDER_PREFIX):
return await process_paid_order(db, out_trade_no, result)
if out_trade_no.startswith(ORDER_PREFIX_EVENT):
from .event_service import process_paid_event_order
return await process_paid_event_order(db, out_trade_no, result)
if out_trade_no.startswith("MK_"):
from ..market import service as market_service
return await market_service.handle_notify(db, out_trade_no, result)
logger.warning("未知订单类型回调: %s", out_trade_no)
return True
async def process_paid_order(db: Database, order_no: str, wechat_data: dict) -> bool:
"""支付成功 → 引擎加款(幂等)。返回 False = 处理失败(触发微信重试)。"""
order = await db.compute_recharges.get_by_order_no(order_no)
if order is None:
logger.warning("充值订单不存在: %s", order_no)
return True
if order["status"] in ("paid", "credited"):
return True # 已处理(回调重放/并发),幂等成功
# 幂等排他:仅 pending → paid 成功者继续到账
payload_json = json.dumps(wechat_data, ensure_ascii=False)[:8000]
got = await db.compute_recharges.mark_paid(order_no, wechat_data.get("transaction_id", ""), payload_json)
if not got:
return True
total_micro = int(order["quota_micro"]) + int(order["bonus_quota_micro"])
try:
engine_user_id = int(order["engine_user_id"] or 0)
if not engine_user_id:
engine_user_id = await _resolve_engine_user_id(order["username"])
if not engine_user_id:
raise RuntimeError(f"引擎用户未解析: {order['username']}")
await compute_client.adjust_user_quota(engine_user_id, total_micro, "add")
except Exception as exc: # noqa: BLE001
logger.error("充值到账失败 order=%s: %s", order_no, exc, exc_info=True)
await db.compute_recharges.mark_status(order_no, "failed")
await db.audit.add(action="compute.recharge_failed", resource="compute_recharge_order",
resource_id=order_no, detail=f"{total_micro} micro: {exc}",
user_id=order["user_id"])
return False
await db.compute_recharges.mark_credited(order_no)
await compute_client.sync_user_mirror(db, engine_user_id) # best-effort 镜像回写
await db.audit.add(action="compute.recharge_credited", resource="compute_recharge_order",
resource_id=order_no, detail=f"{total_micro} micro, engine_user={engine_user_id}",
user_id=order["user_id"])
logger.info("充值到账成功 order=%s user=%s micro=%s", order_no, order["username"], total_micro)
return True
# ---------------------------------------------------------------------------
# 对账 / 过期
# ---------------------------------------------------------------------------
async def reconcile(db: Database, order: dict) -> dict:
"""pending 单对账:查微信侧状态,已支付则补单;超时未付置 closed。"""
if order["status"] != "pending":
return order
if order["expires_at"] and order["expires_at"] < utcnow_iso():
await db.compute_recharges.mark_status(order["order_no"], "closed")
return await db.compute_recharges.get_by_order_no(order["order_no"]) or order
data = await wxpay.query_order(order["order_no"])
state = (data or {}).get("trade_state", "")
if state == "SUCCESS":
if not await process_paid_order(db, order["order_no"], data):
logger.error("对账补单失败 order=%s", order["order_no"])
elif state in ("CLOSED", "PAY_ERROR", "REVOKED"):
await db.compute_recharges.mark_status(order["order_no"], "closed")
return await db.compute_recharges.get_by_order_no(order["order_no"]) or order
# ---------------------------------------------------------------------------
# 分账结果回调(微信服务商分账)
# ---------------------------------------------------------------------------
async def handle_profitsharing_notify(db: Database, result: dict) -> bool:
"""分账回调分发 → 结算服务确认 escrow released(幂等)。"""
from ..services.settlement_service import SettlementService
return await SettlementService(db).handle_split_notify(result)
# ---------------------------------------------------------------------------
# OPC 收款绑定(个人 openid / 商户子商户号 双路径)
# ---------------------------------------------------------------------------
def _bindings(db: Database) -> PaymentBindingRepository:
return PaymentBindingRepository(db.session)
async def bind_personal(db: Database, user: dict, real_name: str = "") -> dict:
"""个人绑定:把 OPC 小程序 openid 添加为 PERSONAL_OPENID 分账接收方。
前置:用户已用小程序微信登录(users.wx_mini_openid 已存)。
"""
openid = user.get("wx_mini_openid", "")
if not openid:
raise ValueError("当前账号未绑定小程序微信身份,请先用微信登录小程序")
if not pay_config.profitsharing_enabled():
raise RuntimeError("服务商分账未配置,暂不可绑定")
name = real_name or user.get("nickname", "")
try:
await profitsharing.add_receiver(
account_type="PERSONAL_OPENID", account=openid, name=name,
)
except RuntimeError as exc:
raise ValueError(f"微信添加分账接收方失败:{exc}") from exc
repo = _bindings(db)
existing = await repo.get_active(user["id"], "personal")
if existing is not None:
return await repo.update(existing["id"], {"real_name": name, "status": "active"})
return await repo.create({
"user_id": user["id"], "bind_type": "personal",
"openid": openid, "real_name": name, "status": "active",
})
async def bind_merchant(db: Database, user: dict, sub_mchid: str = "",
applyment_id: str = "") -> dict:
"""商户绑定:记录特约商户进件(子商户号)并把 sub_mchid 添加为 MERCHANT_ID 接收方。
- 进件为异步微信审核:有 applyment_id 时记录 applying,子商户号回填后置 active;
- 有 sub_mchid 时直接调微信添加接收方并置 active。
"""
if not sub_mchid and not applyment_id:
raise ValueError("sub_mchid 或 applyment_id 至少提供一个")
repo = _bindings(db)
binding = await repo.create({
"user_id": user["id"], "bind_type": "merchant",
"sub_mchid": sub_mchid, "applyment_id": applyment_id,
"status": "active" if sub_mchid else "applying",
})
if sub_mchid and pay_config.profitsharing_enabled():
try:
await profitsharing.add_receiver(
account_type="MERCHANT_ID", account=sub_mchid,
)
except RuntimeError as exc:
# 接收方添加失败不影响进件记录;置 failed 待人工处理
return await repo.set_status(binding["id"], "failed",
detail=str(exc)[:300])
return binding
async def update_merchant_applyment(db: Database, user_id: str, applyment_id: str,
sub_mchid: str) -> dict | None:
"""进件审核通过后回填子商户号并激活(运营/回调调用)。"""
repo = _bindings(db)
rows = await repo.list_by_user(user_id)
b = next((x for x in rows if x["applyment_id"] == applyment_id), None)
if b is None:
return None
return await repo.update(b["id"], {"sub_mchid": sub_mchid, "status": "active"})
async def list_bindings(db: Database, user: dict) -> list[dict]:
"""我的收款绑定列表。"""
return await _bindings(db).list_by_user(user["id"])
async def disable_binding(db: Database, user: dict, binding_id: str) -> dict:
"""停用绑定(结清后可停用,避免误分账)。"""
repo = _bindings(db)
b = await repo.get(binding_id)
if b is None or b["user_id"] != user["id"]:
raise ValueError("绑定不存在")
return await repo.set_status(binding_id, "disabled")