diff --git a/alembic/versions/0039_wx_profitsharing.py b/alembic/versions/0039_wx_profitsharing.py new file mode 100644 index 0000000..8364f0d --- /dev/null +++ b/alembic/versions/0039_wx_profitsharing.py @@ -0,0 +1,97 @@ +"""微信服务商分账:OPC 收款绑定表 + 分账流水镜像表 + escrows 资金托管扩展列。 + +- ``payment_bindings``:OPC 收款绑定(个人 openid / 商户子商户号 双路径) +- ``settlement_logs``:分账流水镜像(业务侧对账) +- ``escrows`` 增列:channel / transaction_id / share_order_no / share_status / receiver_binding_id + +Revision ID: 0039_wx_profitsharing +Revises: 0038_task_park_visibility +Create Date: 2026-09-01 +""" +from __future__ import annotations + +import sqlalchemy as sa +from alembic import op + +revision = "0039_wx_profitsharing" +down_revision = "0038_task_park_visibility" +branch_labels = None +depends_on = None + +S = sa.String +T = sa.Text +I = sa.Integer + + +def upgrade() -> None: + # 1) OPC 收款绑定(服务商分账接收方) + op.create_table( + "payment_bindings", + sa.Column("id", S(), primary_key=True), # pb_ + sa.Column("user_id", S(), nullable=True, server_default=""), # 平台 users.id(OPC 主体) + sa.Column("bind_type", S(), nullable=True, server_default="personal"), # personal | merchant + sa.Column("openid", S(), nullable=True, server_default=""), # personal:小程序 openid + sa.Column("real_name", S(), nullable=True, server_default=""), # personal:实名 + sa.Column("sub_mchid", S(), nullable=True, server_default=""), # merchant:子商户号 + sa.Column("applyment_id", S(), nullable=True, server_default=""), # merchant:进件申请单号 + sa.Column("status", S(), nullable=True, server_default="applying"), # applying|active|disabled|rejected|failed + sa.Column("created_at", S(), nullable=True, server_default=""), + sa.Column("updated_at", S(), nullable=True, server_default=""), + ) + op.create_index("ix_payment_bindings_user_id", "payment_bindings", ["user_id"]) + op.create_index("ix_payment_bindings_status", "payment_bindings", ["status"]) + op.create_index("ix_payment_bindings_sub_mchid", "payment_bindings", ["sub_mchid"]) + op.create_index("ix_payment_bindings_created_at", "payment_bindings", ["created_at"]) + + # 2) 分账流水镜像 + op.create_table( + "settlement_logs", + sa.Column("id", S(), primary_key=True), # sl_ + sa.Column("escrow_id", S(), nullable=True, server_default=""), + sa.Column("task_id", S(), nullable=True, server_default=""), + sa.Column("share_order_no", S(), nullable=True, server_default=""), # 分账单号(幂等键) + sa.Column("channel", S(), nullable=True, server_default="wx_split"), # wx_split|manual|... + sa.Column("amount", I(), nullable=True, server_default="0"), # 分账总额(分) + sa.Column("commission", I(), nullable=True, server_default="0"), # 平台佣金(分) + sa.Column("receivers_json", T(), nullable=True), # 接收方+金额快照 + sa.Column("status", S(), nullable=True, server_default="initiated"), # initiated|shared|failed|returned + sa.Column("detail", S(), nullable=True, server_default=""), + sa.Column("created_at", S(), nullable=True, server_default=""), + sa.Column("updated_at", S(), nullable=True, server_default=""), + ) + op.create_index("ix_settlement_logs_escrow_id", "settlement_logs", ["escrow_id"]) + op.create_index("ix_settlement_logs_share_order_no", "settlement_logs", ["share_order_no"]) + op.create_index("ix_settlement_logs_status", "settlement_logs", ["status"]) + op.create_index("ix_settlement_logs_created_at", "settlement_logs", ["created_at"]) + + # 3) escrows 资金托管扩展列(幂等:先查后加,SQLite/MySQL 通用) + bind = op.get_bind() + existing = {c["name"] for c in bind.dialect.get_columns(bind, "escrows")} + escrow_cols = [ + ("channel", "VARCHAR(32) NOT NULL DEFAULT 'wx_split'"), + ("transaction_id", "VARCHAR(64) NOT NULL DEFAULT ''"), + ("share_order_no", "VARCHAR(64) NOT NULL DEFAULT ''"), + ("share_status", "VARCHAR(32) NOT NULL DEFAULT 'pending'"), + ("receiver_binding_id", "VARCHAR(64) NOT NULL DEFAULT ''"), + ] + for name, ddl in escrow_cols: + if name not in existing: + op.execute(f"ALTER TABLE escrows ADD COLUMN {name} {ddl}") + + +def downgrade() -> None: + op.drop_index("ix_settlement_logs_created_at", table_name="settlement_logs") + op.drop_index("ix_settlement_logs_status", table_name="settlement_logs") + op.drop_index("ix_settlement_logs_share_order_no", table_name="settlement_logs") + op.drop_index("ix_settlement_logs_escrow_id", table_name="settlement_logs") + op.drop_table("settlement_logs") + op.drop_index("ix_payment_bindings_created_at", table_name="payment_bindings") + op.drop_index("ix_payment_bindings_sub_mchid", table_name="payment_bindings") + op.drop_index("ix_payment_bindings_status", table_name="payment_bindings") + op.drop_index("ix_payment_bindings_user_id", table_name="payment_bindings") + op.drop_table("payment_bindings") + # escrows 扩展列(降级时删除,仅 MySQL 支持 DROP COLUMN;SQLite 需重建,此处跳过) + bind = op.get_bind() + if bind.dialect.name == "mysql": + for name in ("receiver_binding_id", "share_status", "share_order_no", "transaction_id", "channel"): + op.execute(f"ALTER TABLE escrows DROP COLUMN {name}") diff --git a/app/infrastructure/models.py b/app/infrastructure/models.py index 4c3d407..dd5c529 100644 --- a/app/infrastructure/models.py +++ b/app/infrastructure/models.py @@ -595,6 +595,12 @@ class Escrow(Base): amount: Mapped[int] = mapped_column(Integer, default=0) commission: Mapped[int] = mapped_column(Integer, default=0) # 平台佣金 status: Mapped[str] = mapped_column(String, default="deposited") # deposited/frozen/released/refunded + # ── 微信服务商分账(资金托管)字段 ── + channel: Mapped[str] = mapped_column(String, default="wx_split") # wx_split | bank_escrow | direct_pay | manual + transaction_id: Mapped[str] = mapped_column(String, default="") # 微信支付单号(收单来源) + share_order_no: Mapped[str] = mapped_column(String, default="") # 分账单号(幂等键) + share_status: Mapped[str] = mapped_column(String, default="pending") # pending/sharing/shared/failed + receiver_binding_id: Mapped[str] = mapped_column(String, default="") # 接单者收款绑定 id created_at: Mapped[str] = mapped_column(String, default="") updated_at: Mapped[str] = mapped_column(String, default="") diff --git a/app/infrastructure/repositories.py b/app/infrastructure/repositories.py index e2788b5..f980934 100644 --- a/app/infrastructure/repositories.py +++ b/app/infrastructure/repositories.py @@ -2179,10 +2179,25 @@ class EscrowRepository: await self.session.commit() return self._to_dict(row) + async def update_fields(self, escrow_id: str, fields: dict) -> dict | None: + """更新台账扩展字段(channel / transaction_id / share_order_no / share_status / receiver_binding_id)。""" + row = await self.session.get(Escrow, escrow_id) + if row is None: + return None + for k, v in fields.items(): + if hasattr(Escrow, k) and k not in ("id", "task_id"): + setattr(row, k, v) + row.updated_at = utcnow_iso() + await self.session.commit() + return self._to_dict(row) + @staticmethod def _to_dict(e: Escrow) -> dict: return {"id": e.id, "task_id": e.task_id, "task_title": e.task_title, "amount": e.amount, - "commission": e.commission, "status": e.status, "created_at": e.created_at} + "commission": e.commission, "status": e.status, "channel": e.channel, + "transaction_id": e.transaction_id, "share_order_no": e.share_order_no, + "share_status": e.share_status, "receiver_binding_id": e.receiver_binding_id, + "created_at": e.created_at, "updated_at": e.updated_at} class ContractRepository: diff --git a/app/pay/config.py b/app/pay/config.py index 1bad99e..b6ef464 100644 --- a/app/pay/config.py +++ b/app/pay/config.py @@ -66,3 +66,26 @@ def pay_enabled() -> bool: WECHATPAY_MCHID and WECHATPAY_CERT_SERIAL_NO and WECHATPAY_APIV3_KEY and WECHATPAY_PRIVATE_KEY_PATH and WECHATPAY_NOTIFY_URL ) + + +# --------------------------------------------------------------------------- +# 微信支付服务商(分账 / 资金托管)配置 —— 与直连商户配置并存,分账以其为开关 +# --------------------------------------------------------------------------- +# 服务商商户号(10 位):分账调用 /v3/profitsharing/orders 的请求主体(sp_mchid) +WECHATPAY_SP_MCHID = os.environ.get("PINEAGENTS_WX_SP_MCHID", "") +# 服务商 AppID(小程序/公众号):PERSONAL_OPENID 接收方的 openid 归属该 appid +WECHATPAY_SP_APPID = os.environ.get("PINEAGENTS_WX_SP_APPID", "") +# 分账结果回调通知地址(公网 HTTPS,如 https://opc.pinesound.cn/opc/pay/profitsharing/notify) +WECHATPAY_SPLIT_NOTIFY_URL = os.environ.get("PINEAGENTS_WX_SPLIT_NOTIFY_URL", "") +# 平台佣金比例:分账给服务商商户号的部分(默认 5%,与 settlement COMMISSION_RATE 一致) +WECHATPAY_SPLIT_COMMISSION_RATE = float( + os.environ.get("PINEAGENTS_WX_SPLIT_COMMISSION_RATE", "0.05") +) + + +def profitsharing_enabled() -> bool: + """服务商分账能力是否可用(sp_mchid + 证书/密钥 + 回调地址齐备)。""" + return bool( + WECHATPAY_SP_MCHID and WECHATPAY_CERT_SERIAL_NO and WECHATPAY_APIV3_KEY + and WECHATPAY_PRIVATE_KEY_PATH and WECHATPAY_SPLIT_NOTIFY_URL + ) diff --git a/app/pay/models.py b/app/pay/models.py index 21d5be4..cc55286 100644 --- a/app/pay/models.py +++ b/app/pay/models.py @@ -1,5 +1,8 @@ # -*- coding: utf-8 -*- -"""支付子应用 ORM:算力充值订单 + 活动报名订单(落平台主库,alembic 统一迁移)。""" +"""支付子应用 ORM:算力充值订单 + 活动报名订单 + OPC 收款绑定 + 分账流水。 + +(落平台主库,alembic 统一迁移;新表由 create_all 幂等创建。) +""" from __future__ import annotations from sqlalchemy import Integer, String, Text @@ -67,3 +70,44 @@ class EventOrder(Base): notify_payload: Mapped[str] = mapped_column(Text, default="") created_at: Mapped[str] = mapped_column(String, default="", index=True) updated_at: Mapped[str] = mapped_column(String, default="") + + +class PaymentBinding(Base): + """OPC 收款绑定(微信服务商分账接收方)。 + + bind_type: + - personal:个人绑定 → PERSONAL_OPENID 分账到 OPC 微信零钱 + - merchant:商户绑定 → 特约商户进件生成子商户号 sub_mchid → MERCHANT_ID 接收方 + + 状态机:applying(商户进件中)→ active(可接收分账)→ disabled / rejected。 + """ + __tablename__ = "payment_bindings" + + id: Mapped[str] = mapped_column(String, primary_key=True) # pb_ + user_id: Mapped[str] = mapped_column(String, default="", index=True) # 平台 users.id(OPC 主体) + bind_type: Mapped[str] = mapped_column(String, default="personal") # personal | merchant + openid: Mapped[str] = mapped_column(String, default="") # personal:小程序 openid + real_name: Mapped[str] = mapped_column(String, default="") # personal:实名(微信校验) + sub_mchid: Mapped[str] = mapped_column(String, default="", index=True) # merchant:子商户号 + applyment_id: Mapped[str] = mapped_column(String, default="") # merchant:进件申请单号 + status: Mapped[str] = mapped_column(String, default="applying", index=True) # applying|active|disabled|rejected + created_at: Mapped[str] = mapped_column(String, default="", index=True) + updated_at: Mapped[str] = mapped_column(String, default="") + + +class SettlementLog(Base): + """分账流水镜像(业务侧对账用,真金在微信侧)。""" + __tablename__ = "settlement_logs" + + id: Mapped[str] = mapped_column(String, primary_key=True) # sl_ + escrow_id: Mapped[str] = mapped_column(String, default="", index=True) # escrows.id + task_id: Mapped[str] = mapped_column(String, default="", index=True) + share_order_no: Mapped[str] = mapped_column(String, default="", index=True) # 分账单号(幂等键) + channel: Mapped[str] = mapped_column(String, default="wx_split") # wx_split | manual | ... + amount: Mapped[int] = mapped_column(Integer, default=0) # 分账总额(分) + commission: Mapped[int] = mapped_column(Integer, default=0) # 平台佣金(分) + receivers_json: Mapped[str] = mapped_column(Text, default="[]") # 接收方+金额快照 + status: Mapped[str] = mapped_column(String, default="initiated", index=True) # initiated|shared|failed|returned + detail: Mapped[str] = mapped_column(String, default="") + created_at: Mapped[str] = mapped_column(String, default="", index=True) + updated_at: Mapped[str] = mapped_column(String, default="") diff --git a/app/pay/profitsharing.py b/app/pay/profitsharing.py new file mode 100644 index 0000000..ac52b14 --- /dev/null +++ b/app/pay/profitsharing.py @@ -0,0 +1,207 @@ +# -*- coding: utf-8 -*- +"""微信支付服务商 · 分账 / 资金托管客户端(懒加载单例)。 + +覆盖能力(服务商模式 partner_mode=True,mchid=服务商商户号): +- 添加分账接收方(MERCHANT_ID 子商户号 / PERSONAL_OPENID 个人微信零钱) +- 请求分账(按订单分账:OPC 收款 95% + 平台佣金 5%) +- 查询分账 / 分账回退 +- 分账结果回调验签解密(复用微信平台证书 + APIv3 密钥) + +设计原则: +- 与 ``wxpay.py`` 直连收单客户端并存:收单继续走现有商户,分账用服务商商户号发起; +- 懒加载:缺服务商配置时 ``profitsharing_enabled()==False``,所有调用返回不可用,业务降级; +- 所有失败统一抛 RuntimeError(携带微信报文摘要),由调用方决定降级/重试。 +""" +from __future__ import annotations + +import json +import logging +from typing import Optional + +from wechatpayv3.async_ import AsyncWeChatPay, WeChatPayType +from wechatpayv3.async_.utils import aes_decrypt + +from . import config as pay_config + +logger = logging.getLogger("pay.profitsharing") + +_sp: Optional[AsyncWeChatPay] = None +_init_lock = False +_init_error: str = "" + + +async def _ensure_split_pay() -> bool: + """确保服务商分账客户端已初始化(懒加载 + 双检锁)。""" + global _sp, _init_lock, _init_error + if _sp is not None: + return True + if not pay_config.profitsharing_enabled(): + _init_error = "服务商分账配置不完整(PINEAGENTS_WX_SP_MCHID / 证书 / APIv3 / notify)" + return False + if _init_lock: + import asyncio + for _ in range(50): + if _sp is not None: + return True + await asyncio.sleep(0.1) + return False + _init_lock = True + try: + with open(pay_config.WECHATPAY_PRIVATE_KEY_PATH, mode="r") as f: + private_key = f.read() + import os + public_key_path = pay_config.WECHATPAY_PUBLIC_KEY_PATH + # 服务商模式:mchid=服务商商户号,appid=服务商 AppID;证书/密钥与直连同商户侧体系 + _sp = AsyncWeChatPay( + wechatpay_type=WeChatPayType.NATIVE, + mchid=pay_config.WECHATPAY_SP_MCHID, + private_key=private_key, + cert_serial_no=pay_config.WECHATPAY_CERT_SERIAL_NO, + appid=pay_config.WECHATPAY_SP_APPID or pay_config.WECHATPAY_NATIVE_APPID, + apiv3_key=pay_config.WECHATPAY_APIV3_KEY, + notify_url=pay_config.WECHATPAY_SPLIT_NOTIFY_URL, + cert_dir=pay_config.WECHATPAY_CERT_DIR, + logger=logger, + partner_mode=True, + public_key=( + open(public_key_path).read() + if public_key_path and os.path.exists(public_key_path) + else None + ), + public_key_id=pay_config.WECHATPAY_PUBLIC_KEY_ID or None, + ) + await _sp.__aenter__() + logger.info("微信服务商分账客户端懒加载初始化成功 sp_mchid=%s", pay_config.WECHATPAY_SP_MCHID) + return True + except Exception as exc: # noqa: BLE001 + _init_error = str(exc) + logger.error("微信服务商分账客户端初始化失败: %s", exc, exc_info=True) + return False + finally: + _init_lock = False + + +def init_error() -> str: + """最近一次初始化失败原因(供接口 503 提示)。""" + return _init_error + + +def _parse(data, code: int, result) -> dict: + """统一把 wechatpayv3 返回解析为 dict;非 200 抛 RuntimeError。""" + if code != 200: + raise RuntimeError(f"微信分账接口失败 http={code}: {str(result)[:300]}") + return result if isinstance(result, dict) else json.loads(result) + + +# --------------------------------------------------------------------------- +# 分账接收方 +# --------------------------------------------------------------------------- +async def add_receiver(*, account_type: str, account: str, name: str = "", + relation_type: str = "SERVICE_PROVIDER") -> dict: + """添加分账接收方。 + + - account_type: MERCHANT_ID(子商户号) | PERSONAL_OPENID(个人微信零钱) + - account: 子商户号 或 小程序 openid + - name: 个人实名(PERSONAL_OPENID 时传,微信校验实名一致) + """ + if not await _ensure_split_pay(): + raise RuntimeError(f"服务商分账未就绪: {_init_error or '未配置'}") + code, result = await _sp.profitsharing_add_receiver( + account_type=account_type, account=account, + relation_type=relation_type, + name=name or None, + appid=pay_config.WECHATPAY_SP_APPID or pay_config.WECHATPAY_NATIVE_APPID, + ) + return _parse(result, code, result) + + +# --------------------------------------------------------------------------- +# 请求分账 / 查询 / 回退 +# --------------------------------------------------------------------------- +async def create_split(*, transaction_id: str, out_order_no: str, + receivers: list[dict], sub_mchid: str = "", + unfreeze_unsplit: bool = True) -> dict: + """请求分账。 + + - transaction_id: 微信支付单号(订单支付成功回调带回) + - out_order_no: 平台分账单号(幂等键,PS__) + - receivers: [{type, account, amount(分), description}] + - sub_mchid: 服务商模式下收单子商户号(分账方) + """ + if not await _ensure_split_pay(): + raise RuntimeError(f"服务商分账未就绪: {_init_error or '未配置'}") + code, result = await _sp.profitsharing_order( + transaction_id=transaction_id, + out_order_no=out_order_no, + receivers=receivers, + unfreeze_unsplit=unfreeze_unsplit, + sub_mchid=sub_mchid or None, + ) + return _parse(result, code, result) + + +async def query_split(*, transaction_id: str, out_order_no: str, + sub_mchid: str = "") -> Optional[dict]: + """查询分账结果;无结果返回 None。""" + if not await _ensure_split_pay(): + return None + try: + code, result = await _sp.profitsharing_order_query( + transaction_id=transaction_id, out_order_no=out_order_no, + sub_mchid=sub_mchid or None, + ) + except Exception as exc: # noqa: BLE001 + logger.warning("查询分账异常 order=%s: %s", out_order_no, exc) + return None + if code != 200: + return None + return result if isinstance(result, dict) else json.loads(result) + + +async def split_return(*, out_order_no: str, return_mchid: str, amount: int, + description: str = "分账回退", sub_mchid: str = "") -> dict: + """分账回退(退款/纠错场景)。""" + if not await _ensure_split_pay(): + raise RuntimeError(f"服务商分账未就绪: {_init_error or '未配置'}") + out_return_no = f"RT_{out_order_no}_{int(__import__('time').time())}" + code, result = await _sp.profitsharing_return( + out_return_no=out_return_no, return_mchid=return_mchid, amount=amount, + description=description, out_order_no=out_order_no, sub_mchid=sub_mchid or None, + ) + return _parse(result, code, result) + + +# --------------------------------------------------------------------------- +# 分账结果回调验签 + 解密(复用同一平台证书 / APIv3 密钥体系) +# --------------------------------------------------------------------------- +async def verify_split_notify(headers, body) -> Optional[dict]: + """验证微信分账回调签名并解密 resource;失败返回 None。 + + 分账回调与支付回调同用微信平台证书 + APIv3 密钥;若直连客户端已初始化, + 直接复用 wxpay.verify_and_decrypt,否则用服务商客户端验签。 + """ + from . import wxpay + if await _ensure_split_pay(): + body_str = body.decode("UTF-8") if isinstance(body, bytes) else body + if await _sp._core._verify_signature_async(headers, body_str): + data = json.loads(body_str) + resource = data.get("resource") + if not resource: + logger.error("分账回调缺少 resource 字段") + return None + try: + decrypted = aes_decrypt( + nonce=resource.get("nonce"), + ciphertext=resource.get("ciphertext"), + associated_data=resource.get("associated_data") or "", + apiv3_key=_sp._core._apiv3_key, + ) + except Exception as exc: # noqa: BLE001 + logger.error("分账回调解密异常: %s", exc) + return None + if not decrypted: + return None + data.update({"resource": json.loads(decrypted)}) + return data + # 回退:直连客户端验签 + return await wxpay.verify_and_decrypt(headers, body) diff --git a/app/pay/repository.py b/app/pay/repository.py index 9d77073..95116c7 100644 --- a/app/pay/repository.py +++ b/app/pay/repository.py @@ -5,7 +5,7 @@ from __future__ import annotations from sqlalchemy import select from ..infrastructure.repositories import new_id, utcnow_iso -from .models import ComputeRechargeOrder, EventOrder +from .models import ComputeRechargeOrder, EventOrder, PaymentBinding, SettlementLog def _to_dict(o: ComputeRechargeOrder) -> dict: @@ -210,3 +210,122 @@ def _to_event_dict(o: EventOrder) -> dict: "expires_at": o.expires_at, "paid_at": o.paid_at, "credited_at": o.credited_at, "notify_payload": o.notify_payload, "created_at": o.created_at, "updated_at": o.updated_at, } + + +# --------------------------------------------------------------------------- +# OPC 收款绑定(微信服务商分账接收方) +# --------------------------------------------------------------------------- +def _to_binding_dict(b: PaymentBinding) -> dict: + return { + "id": b.id, "user_id": b.user_id, "bind_type": b.bind_type, + "openid": b.openid, "real_name": b.real_name, "sub_mchid": b.sub_mchid, + "applyment_id": b.applyment_id, "status": b.status, + "created_at": b.created_at, "updated_at": b.updated_at, + } + + +class PaymentBindingRepository: + """OPC 收款绑定仓储(个人 openid / 商户子商户号双路径)。""" + + def __init__(self, session): + self.session = session + + async def create(self, fields: dict) -> dict: + now = utcnow_iso() + row = PaymentBinding( + id=new_id("pb"), created_at=now, updated_at=now, + **{k: v for k, v in fields.items() if hasattr(PaymentBinding, k) and k != "id"}, + ) + self.session.add(row) + await self.session.commit() + return _to_binding_dict(row) + + async def get(self, binding_id: str) -> dict | None: + row = await self.session.get(PaymentBinding, binding_id) + return _to_binding_dict(row) if row else None + + async def list_by_user(self, user_id: str) -> list[dict]: + rows = (await self.session.scalars( + select(PaymentBinding).where(PaymentBinding.user_id == user_id) + .order_by(PaymentBinding.created_at.desc()) + )).all() + return [_to_binding_dict(r) for r in rows] + + async def get_active(self, user_id: str, bind_type: str | None = None) -> dict | None: + """取该用户可接收分账的绑定(优先商户绑定→个人绑定;bind_type 指定时仅取该类型)。""" + order = ("merchant", "personal") if bind_type is None else (bind_type,) + for bt in order: + row = await self.session.scalar( + select(PaymentBinding) + .where( + PaymentBinding.user_id == user_id, + PaymentBinding.status == "active", + PaymentBinding.bind_type == bt, + ) + .order_by(PaymentBinding.created_at.desc()) + .limit(1) + ) + if row is not None: + return _to_binding_dict(row) + return None + + async def update(self, binding_id: str, fields: dict) -> dict | None: + row = await self.session.get(PaymentBinding, binding_id) + if row is None: + return None + for k, v in fields.items(): + if hasattr(PaymentBinding, k): + setattr(row, k, v) + row.updated_at = utcnow_iso() + await self.session.commit() + return _to_binding_dict(row) + + async def set_status(self, binding_id: str, status: str, **extra) -> dict | None: + return await self.update(binding_id, {"status": status, **extra}) + + +# --------------------------------------------------------------------------- +# 分账流水镜像(对账) +# --------------------------------------------------------------------------- +def _to_log_dict(l: SettlementLog) -> dict: + return { + "id": l.id, "escrow_id": l.escrow_id, "task_id": l.task_id, + "share_order_no": l.share_order_no, "channel": l.channel, + "amount": l.amount, "commission": l.commission, + "receivers_json": l.receivers_json, "status": l.status, "detail": l.detail, + "created_at": l.created_at, "updated_at": l.updated_at, + } + + +class SettlementLogRepository: + """分账流水镜像仓储(业务侧记录真金去向,供对账/审计)。""" + + def __init__(self, session): + self.session = session + + async def create(self, fields: dict) -> dict: + now = utcnow_iso() + row = SettlementLog( + id=new_id("sl"), created_at=now, updated_at=now, + **{k: v for k, v in fields.items() if hasattr(SettlementLog, k) and k != "id"}, + ) + self.session.add(row) + await self.session.commit() + return _to_log_dict(row) + + async def get_by_share_order_no(self, share_order_no: str) -> dict | None: + row = await self.session.scalar( + select(SettlementLog).where(SettlementLog.share_order_no == share_order_no) + ) + return _to_log_dict(row) if row else None + + async def set_status(self, log_id: str, status: str, detail: str = "") -> dict | None: + row = await self.session.get(SettlementLog, log_id) + if row is None: + return None + row.status = status + if detail: + row.detail = detail + row.updated_at = utcnow_iso() + await self.session.commit() + return _to_log_dict(row) diff --git a/app/pay/routers.py b/app/pay/routers.py index 635ccb6..693ee2a 100644 --- a/app/pay/routers.py +++ b/app/pay/routers.py @@ -15,7 +15,7 @@ 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 . import profitsharing, service, wxpay from .repository import ComputeRechargeOrderRepository logger = logging.getLogger("pay.api") @@ -128,3 +128,88 @@ async def recharge_orders( user: dict = Depends(require_roles("opc_member")), ): return {"items": await _repo(db).list_by_user(user["id"])} + + +# ── 微信服务商分账回调(微信服务器调用,无登录鉴权,V3 验签)─────────────────── +@router.post("/opc/pay/profitsharing/notify", summary="微信分账结果回调(无登录鉴权,V3 验签)") +async def profitsharing_notify( + request: Request, + db: Database = Depends(get_db), +): + body = await request.body() + result = await profitsharing.verify_split_notify(request.headers, body) + if not result or not isinstance(result, dict): + raise HTTPException(status_code=400, detail="回调验签/解密失败") + resource = result.pop("resource", {}) + if isinstance(resource, dict): + result.update(resource) + logger.info("微信分账回调: out_order_no=%s", result.get("out_order_no")) + try: + ok = await service.handle_profitsharing_notify(db, result) + except Exception as exc: # noqa: BLE001 + logger.error("分账回调处理异常: %s", 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": "成功"} + + +# ── OPC 收款绑定(个人 openid / 商户子商户号 双路径)────────────────────────── +class BindPersonalRequest(BaseModel): + """个人收款绑定:openid 取自登录用户小程序身份。""" + real_name: str = Field(default="", description="实名(用于微信校验,缺省用昵称)") + + +class BindMerchantRequest(BaseModel): + """商户收款绑定:进件申请单号 或 已回填的子商户号。""" + sub_mchid: str = Field(default="", description="特约商户子商户号(进件通过后回填)") + applyment_id: str = Field(default="", description="进件申请单号(进件中)") + + +@router.post("/opc/pay/bindings/personal", summary="个人收款绑定(分账到微信零钱)") +async def opc_bind_personal( + body: BindPersonalRequest, + db: Database = Depends(get_db), + user: dict = Depends(require_roles("opc_member")), +): + try: + return await service.bind_personal(db, user, real_name=body.real_name) + 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/pay/bindings/merchant", summary="商户收款绑定(进件 / 子商户号)") +async def opc_bind_merchant( + body: BindMerchantRequest, + db: Database = Depends(get_db), + user: dict = Depends(require_roles("opc_member")), +): + try: + return await service.bind_merchant( + db, user, sub_mchid=body.sub_mchid, applyment_id=body.applyment_id) + 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.get("/opc/pay/bindings", summary="我的收款绑定列表") +async def opc_bindings( + db: Database = Depends(get_db), + user: dict = Depends(require_roles("opc_member")), +): + return {"items": await service.list_bindings(db, user)} + + +@router.post("/opc/pay/bindings/{binding_id}/disable", summary="停用收款绑定") +async def opc_disable_binding( + binding_id: str, + db: Database = Depends(get_db), + user: dict = Depends(require_roles("opc_member")), +): + try: + return await service.disable_binding(db, user, binding_id) + except ValueError as exc: + raise HTTPException(status_code=404, detail=str(exc)) from exc diff --git a/app/pay/service.py b/app/pay/service.py index 9fba847..c2cfff7 100644 --- a/app/pay/service.py +++ b/app/pay/service.py @@ -19,7 +19,8 @@ from .. import config from ..infrastructure.repositories import Database, utcnow_iso from ..services import compute_client from . import config as pay_config -from . import wxpay +from . import profitsharing, wxpay +from .repository import PaymentBindingRepository logger = logging.getLogger("pay.service") @@ -334,3 +335,99 @@ async def reconcile(db: Database, order: dict) -> dict: 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") diff --git a/app/services/settlement_service.py b/app/services/settlement_service.py index 32d00f6..1b6962b 100644 --- a/app/services/settlement_service.py +++ b/app/services/settlement_service.py @@ -2,16 +2,33 @@ """业务层 · 结算/合同/争议/撮合服务。 资金相关规则(佣金、结算、争议联动)在业务层集中,接口层瘦身。 + +结算(release_escrow)· 验收通过后: + 1) 优先走「微信服务商分账」(资金托管): + 接单者 active 收款绑定 + 订单 transaction_id + 服务商配置齐备 + → 分账 95% 给接单者 / 5% 佣金给平台服务商商户号 → 等分账回调确认后置 released + 2) 降级(未接微信收单 / 无绑定 / 分账未配置): + 仅台账置 released(channel=manual),不阻塞业务,供后续对账补齐。 """ from __future__ import annotations +import json +import logging +import secrets +import time + from fastapi import HTTPException from ..infrastructure.repositories import Database +from ..pay import config as pay_config +from ..pay import profitsharing +from ..pay.repository import PaymentBindingRepository, SettlementLogRepository # 平台佣金比例 COMMISSION_RATE = 0.05 +logger = logging.getLogger("services.settlement") + class SettlementService: """结算托管、电子合同、争议、投资撮合。""" @@ -20,7 +37,7 @@ class SettlementService: self.db = db async def release_escrow(self, task_id: str) -> dict: - """验收通过后结算:任务须 completed;无托管则兜底创建(佣金 5%)。""" + """验收通过后结算:优先微信服务商分账,不满足条件时降级台账结算。""" task = await self.db.tasks.get(task_id) if task is None or task["status"] != "completed": raise HTTPException(status_code=400, detail="任务未完成,不能结算") @@ -30,7 +47,113 @@ class SettlementService: amount = max(task.get("budget_min", 0), task.get("budget_max", 0)) commission = int(amount * COMMISSION_RATE) esc = await self.db.escrows.create(task_id, task["title"], amount, commission) - return await self.db.escrows.set_status(esc["id"], "released") + + tx_id = esc.get("transaction_id") or "" + if not tx_id or not pay_config.profitsharing_enabled(): + # 降级:未接微信收单 / 服务商分账未配置 → 台账结算 + esc = await self.db.escrows.update_fields( + esc["id"], {"channel": "manual", "share_status": "manual"}) + esc = await self.db.escrows.set_status(esc["id"], "released") + await self._log(esc, "manual", "未接入微信收单或分账未配置,台账结算", "shared", + f"no_transaction_id={not tx_id}") + return esc + + receiver, reason = await self._resolve_receiver(task) + if receiver is None: + esc = await self.db.escrows.update_fields( + esc["id"], {"channel": "manual", "share_status": "manual"}) + esc = await self.db.escrows.set_status(esc["id"], "released") + await self._log(esc, "manual", "接单者无收款绑定,台账结算", "shared", reason) + return esc + + # 发起微信分账:金额(元→分);接单者 = 总额 - 佣金,平台佣金单独分给服务商商户号 + total_fen = int(esc["amount"]) * 100 + commission_fen = int(esc["commission"]) * 100 + opc_fen = total_fen - commission_fen + share_no = f"PS_{int(time.time())}_{secrets.token_hex(4).upper()}" + receivers = [ + { + "type": "MERCHANT_ID" if receiver["bind_type"] == "merchant" else "PERSONAL_OPENID", + "account": receiver.get("sub_mchid") or receiver.get("openid"), + "amount": max(opc_fen, 0), + "description": "任务服务费", + }, + { + "type": "MERCHANT_ID", + "account": pay_config.WECHATPAY_SP_MCHID, + "amount": commission_fen, + "description": "平台服务费", + }, + ] + try: + await profitsharing.create_split( + transaction_id=tx_id, out_order_no=share_no, receivers=receivers, + sub_mchid=receiver.get("sub_mchid") or "", + ) + except RuntimeError as exc: + esc = await self.db.escrows.update_fields(esc["id"], { + "channel": "wx_split", "share_order_no": share_no, + "share_status": "failed", "receiver_binding_id": receiver["id"], + }) + await self._log(esc, "wx_split", "分账发起失败", "failed", str(exc)[:300]) + raise HTTPException(status_code=502, detail=f"分账发起失败:{exc}") from exc + + esc = await self.db.escrows.update_fields(esc["id"], { + "channel": "wx_split", "share_order_no": share_no, + "share_status": "sharing", "receiver_binding_id": receiver["id"], + }) + await self._log(esc, "wx_split", "分账已发起,等待微信回调确认", "initiated", + receivers_json=receivers) + logger.info("分账已发起 task=%s share=%s receiver=%s", task_id, share_no, receiver["id"]) + return esc + + async def handle_split_notify(self, result: dict) -> bool: + """分账结果回调:微信确认分账成功 → escrow released + 流水 shared(幂等)。""" + share_no = result.get("out_order_no", "") + if not share_no: + logger.warning("分账回调缺少 out_order_no") + return True + repo = SettlementLogRepository(self.db.session) + log = await repo.get_by_share_order_no(share_no) + if log is None: + logger.warning("分账回调无对应流水: %s", share_no) + return True + if log["status"] == "shared": + return True # 已确认(回调重放) + esc = await self.db.escrows.get(log["escrow_id"]) + if esc is not None: + await self.db.escrows.update_fields(esc["id"], {"share_status": "shared"}) + await self.db.escrows.set_status(esc["id"], "released") + await repo.set_status(log["id"], "shared", "微信分账回调确认") + await self.db.audit.add(action="escrow.released_wx", resource="escrow", + resource_id=log["escrow_id"], + detail=f"share={share_no}", user_id="") + logger.info("分账回调确认 share=%s escrow=%s", share_no, log["escrow_id"]) + return True + + async def _resolve_receiver(self, task: dict) -> tuple[dict | None, str]: + """接单者(claimed_by)的可用收款绑定;无则返回 (None, 原因)。""" + claimer = task.get("claimed_by") + if not claimer: + return None, "任务无接单者(claimed_by 为空)" + repo = PaymentBindingRepository(self.db.session) + binding = await repo.get_active(claimer) + if binding is None: + return None, f"接单者无 active 收款绑定(user={claimer})" + return binding, "" + + async def _log(self, esc: dict, channel: str, note: str, status: str, + detail: str = "", receivers_json: list | None = None) -> dict: + """写分账流水镜像(金额统一为分,与微信口径一致)。""" + repo = SettlementLogRepository(self.db.session) + return await repo.create({ + "escrow_id": esc.get("id", ""), "task_id": esc.get("task_id", ""), + "share_order_no": esc.get("share_order_no", ""), "channel": channel, + "amount": int(esc.get("amount", 0)) * 100, + "commission": int(esc.get("commission", 0)) * 100, + "receivers_json": json.dumps(receivers_json or [], ensure_ascii=False), + "status": status, "detail": detail or note, + }) async def sign_contract(self, task_id: str, actor: dict, opc_id: str | None = None) -> dict: """签订电子合同(幂等:已签返回现有合同)。""" diff --git a/scripts/db/migrate_escrow_columns.py b/scripts/db/migrate_escrow_columns.py new file mode 100644 index 0000000..09c7757 --- /dev/null +++ b/scripts/db/migrate_escrow_columns.py @@ -0,0 +1,60 @@ +# -*- coding: utf-8 -*- +"""幂等迁移:为 ``escrows`` 表补充微信服务商分账扩展列(非运行态)。 + +用法:uv run python scripts/db/migrate_escrow_columns.py + +- 通过 SQLAlchemy Inspector 检查缺失列,对缺失列执行 ``ALTER TABLE ... ADD COLUMN``; +- 幂等:已存在列自动跳过,可重复执行;兼容 SQLite(默认)与 MySQL(生产)。 +- 新增表 ``payment_bindings`` / ``settlement_logs`` 由应用启动时 ``create_all`` 幂等创建, + 无需在此处理(本脚本只负责对既有表加列)。 +""" +from __future__ import annotations + +import sys +from pathlib import Path + +ROOT = Path(__file__).resolve().parent.parent.parent +sys.path.insert(0, str(ROOT)) + +from sqlalchemy import create_engine, inspect # noqa: E402 + +from app.infrastructure.db import sync_database_url # noqa: E402 + +# escrows 需补充的扩展列(列名, 方言通用 DDL) +ESCROW_COLUMNS = [ + ("channel", "VARCHAR(32) NOT NULL DEFAULT 'wx_split'"), + ("transaction_id", "VARCHAR(64) NOT NULL DEFAULT ''"), + ("share_order_no", "VARCHAR(64) NOT NULL DEFAULT ''"), + ("share_status", "VARCHAR(32) NOT NULL DEFAULT 'pending'"), + ("receiver_binding_id", "VARCHAR(64) NOT NULL DEFAULT ''"), +] + + +def main() -> None: + url = sync_database_url() + engine = create_engine(url) + insp = inspect(engine) + if "escrows" not in insp.get_table_names(): + print("escrows 表不存在(未建库),无需迁移") + return + existing = {c["name"] for c in insp.get_columns("escrows")} + added, skipped = [], [] + for name, ddl in ESCROW_COLUMNS: + if name in existing: + skipped.append(name) + continue + with engine.begin() as conn: + conn.exec_driver_sql(f"ALTER TABLE escrows ADD COLUMN {name} {ddl}") + added.append(name) + print(f"新增列: {added or '无'}") + print(f"已存在跳过: {skipped or '无'}") + engine.dispose() + if added: + print("迁移完成。") + else: + print("无新增列,escrows 已是最新。") + print("注:新表 payment_bindings / settlement_logs 由应用启动时 create_all 幂等创建。") + + +if __name__ == "__main__": + main() diff --git a/scripts/smoke_profitsharing.py b/scripts/smoke_profitsharing.py new file mode 100644 index 0000000..9ef6386 --- /dev/null +++ b/scripts/smoke_profitsharing.py @@ -0,0 +1,109 @@ +# -*- coding: utf-8 -*- +"""微信服务商分账改造 · 冒烟测试(临时 sqlite,不碰主库)。 + +覆盖: +1. release_escrow:未接微信收单/分账未配置 → 降级台账结算(channel=manual, released) +2. release_escrow:配置齐备 + 接单者 active 绑定 → 发起微信分账(mock create_split) +3. handle_split_notify:回调确认 → escrow released + log shared(幂等) +""" +from __future__ import annotations + +import asyncio +import os +import sys +import tempfile +from pathlib import Path + +# 临时库:优先 sqlite 内存/临时文件,避免触碰 .env 指定的 MySQL 主库 +_tmp = tempfile.NamedTemporaryFile(suffix=".db", delete=False) +_tmp.close() +os.environ["PINEAGENTS_DEMO_DATABASE_URL"] = f"sqlite+aiosqlite:///{_tmp.name}" + +ROOT = Path(__file__).resolve().parent +sys.path.insert(0, str(ROOT)) + +from unittest.mock import patch # noqa: E402 + +from app.infrastructure.db import Base # noqa: E402 +from app.infrastructure.repositories import Database # noqa: E402 +import app.pay.models # noqa: E402,F401 确保新模型注册 +import app.infrastructure.models # noqa: E402,F401 +from app.services.settlement_service import SettlementService # noqa: E402 + + +async def _setup(db: Database): + async with db._engine.begin() as conn: + await conn.run_sync(Base.metadata.create_all) + + +async def _mk_completed_task(db, claimer: str = "") -> dict: + return await db.tasks.create({ + "task_code": "smk001", "title": "冒烟任务", "status": "completed", + "budget_min": 1000, "budget_max": 1000, "claimed_by": claimer or None, + }) + + +async def main() -> None: + db = Database() + await _setup(db) + + # ── 1) 未配置分账 / 无收单 → 降级台账结算 ──────────────────────────────── + task = await _mk_completed_task(db) + svc = SettlementService(db) + r = await svc.release_escrow(task["id"]) + print("① 降级台账:", r["status"], r["channel"], r["share_status"]) + assert r["status"] == "released" and r["channel"] == "manual" + # 直接查库验证流水 + import sqlite3 + conn = sqlite3.connect(_tmp.name) + row = conn.execute("SELECT channel,status,detail FROM settlement_logs ORDER BY created_at DESC LIMIT 1").fetchone() + conn.close() + print("① 流水:", row) + assert row and row[0] == "manual" and row[1] == "shared" + + # ── 2) 配置齐备 + 接单者 active 绑定 → 发起分账 ───────────────────────── + from app.pay.repository import PaymentBindingRepository + # 造一个接单者用户 + 商户绑定 + user = await db.users.create("smk_opc", "pw123456", nickname="接单者", role="opc_member") + bind = await PaymentBindingRepository(db.session).create({ + "user_id": user["id"], "bind_type": "merchant", + "sub_mchid": "1900000109", "status": "active", + }) + task2 = await _mk_completed_task(db) + await db.tasks.claim(task2["id"], user["id"]) # 接单者=用户 + await db.tasks.set_status(task2["id"], "completed") # 接单后置完成(验收通过) + esc = await db.escrows.create(task2["id"], task2["title"], 1000, 50) + await db.escrows.update_fields(esc["id"], {"transaction_id": "4200001234"}) + with patch.object( + __import__("app.pay.config", fromlist=["profitsharing_enabled"]), + "profitsharing_enabled", return_value=True, + ), patch("app.pay.profitsharing.create_split") as mock_split: + r2 = await svc.release_escrow(task2["id"]) + mock_split.assert_called_once() + call = mock_split.call_args.kwargs + print("② 分账发起: channel=%s share_status=%s" % (r2["channel"], r2["share_status"])) + print("② 接收方:", [(x["type"], x["amount"]) for x in call["receivers"]]) + assert r2["channel"] == "wx_split" and r2["share_status"] == "sharing" + assert call["receivers"][0]["amount"] == 95000 and call["receivers"][1]["amount"] == 5000 + assert call["receivers"][0]["type"] == "MERCHANT_ID" # 商户绑定 → MERCHANT_ID + share_no = call["out_order_no"] + + # ── 3) 回调确认(幂等)─────────────────────────────────────────────────── + ok1 = await svc.handle_split_notify({"out_order_no": share_no}) + conn = sqlite3.connect(_tmp.name) + esc_state = conn.execute("SELECT status,share_status FROM escrows WHERE id=?", (esc["id"],)).fetchone() + log_state = conn.execute("SELECT status FROM settlement_logs WHERE share_order_no=?", (share_no,)).fetchone() + conn.close() + print("③ 回调后 escrow:", esc_state, "log:", log_state) + assert ok1 and esc_state == ("released", "shared") and log_state == ("shared",) + ok2 = await svc.handle_split_notify({"out_order_no": share_no}) # 重放 + assert ok2 # 幂等不报错 + print("③ 幂等重放 OK") + + print("\n✅ 冒烟测试全部通过") + await db.close() + os.unlink(_tmp.name) + + +if __name__ == "__main__": + asyncio.run(main())