Files
group_xinghuo_jinrong/app/repository/convert_repository.py
T

381 lines
18 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.
"""基金转换(convert)代理侧仓储 · T-4(agent 库 risk_convert_detail · 语义收窄)。
**本仓储 = `risk_convert_detail` 唯一写入口**(全仓 grep 确认;PRD §4.1 唯一写入方)。
T+1 模型下 `risk_convert_detail.status` 的 5 值**只表示 agent 侧写入进度**
(PRD §4.1 表定义,业务状态一律以 `core_convert_request.status` 6 态为准):
- `pending` = agent 侧详情尚未写入(受理已落 Core,镜像未写)
- `completed` = 详情 + 主审计已写入成功
- `failed` = 附加写入失败,待补偿(触发者之一,架构 §5.4)
- `cancelled` / `expired` = 保留值(免 ALTER,本版不使用)
唯一写入口 = 受理 / 确认 / 撤单 / 补偿 4 处(T-6/T-7/T-9/T-12 落地后由
convert_service 调用,每处写调用都有注释标注对应任务)。
**S2(评审):清理标记不硬删**——`mark_expired` 置 `status='expired'`,绝不
DELETE,留痕供对账(与 cleanup_pending_convert.py 的口径一致)。
⚠️ v1.0「阶段零占位 + 阶段一实时扣减 + 阶段二回写」三阶段流属**已废弃模型**
(D23:全量重写 T+1,无双模型开关),本文件 `insert_placeholder/complete_convert/
mark_failed` 三方法为 v1.0 遗留,**语义收窄后不再承载状态迁移**,仅供 T-7 完工前
过渡期兼容;T-7 完成后由新镜像写入路径替代(旧方法整体删除)。
"""
from __future__ import annotations
import logging
from datetime import date, datetime, timedelta
from decimal import Decimal
from typing import Any
from sqlalchemy import text
from sqlalchemy.engine import Engine
from sqlalchemy.exc import IntegrityError
from app.config.settings import settings
from app.utils.db import get_engine
logger = logging.getLogger(__name__)
def _to_bind(value: Any) -> Any:
"""Decimal 转 float 再绑定(sqlite 不支持直接绑定 Decimal;MySQL DECIMAL 列自动收口)。
与 T-3 测试教训一致:插入/更新 Decimal 一律转 float,读回侧用 `Decimal(str(...))` 还原。
"""
return float(value) if isinstance(value, Decimal) else value
_STATUS_PENDING = "pending"
_STATUS_COMPLETED = "completed"
_STATUS_FAILED = "failed"
_STATUS_EXPIRED = "expired"
class ConvertRepository:
"""risk_convert_detail 读写(进度镜像 / 查询 / 清理)。
T+1 模型下本表状态 = agent 写入进度(PRD §4.1 语义),唯一写入口 `sync_mirror`
(受理/确认/撤单/补偿 4 处调用,T-6/T-7/T-9/T-12);v1.0 遗留方法标注 Deprecated。
"""
def __init__(self, engine: Engine | None = None) -> None:
self._engine = engine or get_engine(settings.mysql_database, "rw")
# ---------- T+1 模型:进度镜像(唯一写入口) ----------
def sync_mirror(
self,
group_id: str,
status: str,
*,
client_request_id: str | None = None,
out_trade_id: str | None = None,
in_trade_id: str | None = None,
related_trade_id: str | None = None,
nav: Decimal | None = None,
nav_date: date | None = None,
fee_amount: Decimal | None = None,
hold_days_min: int | None = None,
hold_days_max: int | None = None,
nav_stale: bool = False,
cancelled_at: datetime | None = None,
) -> None:
"""T+1 模型下 `risk_convert_detail` **唯一进度镜像写入口**(T-4;PRD §4.1)。
`status` 只允许 agent 写入进度三值(`pending`/`completed`/`failed`)——
**业务状态一律以 `core_convert_request.status` 6 态为准**(cancelled/expired
为建表保留值,本方法不直接写,传则 ValueError)。
三步法(R-a:先查再 INSERT 或 UPDATE,弃方言 UPSERT 兼容 sqlite/MySQL):
- 行不存在 → INSERT 镜像行(`status` = 传入进度,`estimated` = 0 真实请求);
- 行存在 → UPDATE **幂等重写**(同 group_id 重复调用收敛为最新进度,不炸)。
**4 个调用点契约**(开发计划 T-4 DoD 1;每处调用须带对应任务注释):
1. 受理(T-6):`sync_mirror(gid, 'pending', client_request_id=...)`
受理单落 Core 后标记「待写入」(详情确认后才有,受理时全 NULL);
2. 确认(T-7):`sync_mirror(gid, 'completed', out_trade_id=, in_trade_id=,
related_trade_id=, nav=, nav_date=, fee_amount=, hold_days_min=,
hold_days_max=, nav_stale=)` —— 详情 + 主审计写入完成;
3. 撤单(T-7/T-9):`sync_mirror(gid, 原进度, cancelled_at=...)` 幂等重写 ——
撤单只动 Core 状态(→cancelled),镜像**不推 cancelled**(保留值;PRD §4.1
注释),但把撤单时刻写进 `cancelled_at`(Core 的 `core_convert_request`
无该列,撤单时刻只在 agent 侧留痕);
4. 补偿(T-12):附加写入失败 `failed` → 补偿成功回 `completed`。
**不做状态迁移判定**:`status` 由调用方(读 `core_convert_request` 权威 6 态
后)传入,本仓储只落值,不「从 X 到 Y」条件垫片(那是
`convert_request_repository.transition_status` 的职责,T-4 DoD 2)。
"""
if status not in (_STATUS_PENDING, _STATUS_COMPLETED, _STATUS_FAILED):
raise ValueError(
f"risk_convert_detail.status 只允许 agent 进度值"
f"(pending/completed/failed),实际 {status!r};"
f"cancelled/expired 为建表保留值,业务状态请写 core_convert_request"
)
params = {
"gid": group_id,
"status": status,
"cr": client_request_id,
"out": out_trade_id,
"in": in_trade_id,
"rel": related_trade_id,
"nav": _to_bind(nav),
"nav_date": nav_date,
"fee": _to_bind(fee_amount),
"hmin": hold_days_min,
"hmax": hold_days_max,
"stale": 1 if nav_stale else 0,
"cancel_at": cancelled_at,
}
with self._engine.begin() as conn:
existing = conn.execute(
text("SELECT 1 FROM risk_convert_detail WHERE convert_group_id = :gid"),
{"gid": group_id},
).first()
if existing is not None:
conn.execute(
text(
"""
UPDATE risk_convert_detail
SET status = :status,
client_request_id = COALESCE(:cr, client_request_id),
out_trade_id = COALESCE(:out, out_trade_id),
in_trade_id = COALESCE(:in, in_trade_id),
related_trade_id = COALESCE(:rel, related_trade_id),
nav = COALESCE(:nav, nav),
nav_date = COALESCE(:nav_date, nav_date),
fee_amount = COALESCE(:fee, fee_amount),
hold_days_min = COALESCE(:hmin, hold_days_min),
hold_days_max = COALESCE(:hmax, hold_days_max),
nav_stale = :stale,
cancelled_at = COALESCE(:cancel_at, cancelled_at)
WHERE convert_group_id = :gid
"""
),
params,
)
else:
conn.execute(
text(
"""
INSERT INTO risk_convert_detail
(convert_group_id, client_request_id, status, estimated,
out_trade_id, in_trade_id, related_trade_id, nav, nav_date,
fee_amount, hold_days_min, hold_days_max, nav_stale, cancelled_at)
VALUES (:gid, :cr, :status, 0,
:out, :in, :rel, :nav, :nav_date,
:fee, :hmin, :hmax, :stale, :cancel_at)
"""
),
params,
)
# ---------- v1.0 遗留:阶段零占位 / 阶段二回写(T-7 后删除) ----------
def insert_placeholder(self, group_id: str, client_request_id: str | None) -> bool:
"""⛔ **Deprecated(v1.0 遗留,T-7 全量重写后删除)**:阶段零占位。
新模型由 `sync_mirror`(受理 ④ pending 镜像)取代;本方法仅供 T-7 完工前
test_convert_service / test_convert_concurrency 的 v1.0 流回归使用。
"""
"""阶段零占位:插一行 `status='pending'`;**已存在则置回 pending**。
返回 **True = 本笔持有该占位,可以继续**;**False = 该 `client_request_id`
已被另一个 `group_id` 占住**(同键并发,本笔必须让路,由调用方回 202)。
`client_request_id` 为 None 时绑 NULL——MySQL / sqlite 的 UNIQUE 约束均允许多个 NULL,
故「无幂等键的请求」可重复占位、互不冲突(uk_idem 仅对非空键兜底)。
`estimated=0`:convert 占位是真实请求,非风控预估单(与 risk_alert 语义区分)。
**三步法(R-a:弃用方言 UPSERT,改「先查再 INSERT 或 UPDATE」)**。
为什么必须容错「行已存在」:阶段一失败(典型是 `LotConflict` 409)时占位已被
`mark_failed` 置为 `failed`,而架构 §8.3 要求调用方带**同一** `client_request_id`
退避重试(≤3 次、100/200/400ms);重试会走「复用原 group_id 重跑」分支再次进入
阶段零——此时 `uk_group` 与 `uk_idem` 都已被那一行占用,朴素的 INSERT 必撞唯一键,
把**可重试的 409 升级成 `IdempotencyUnavailable`(503)**,且是**确定性的**(重试永不成功),
与 errors.py 里「瞬时状态、恢复后重试即可成功」的注释相反。(T-13 压测前置修复)
置回 `pending` 是正确语义:同一 `group_id` 的这一次尝试正在进行中。
到达此处时既有行的状态只可能是 `pending`(上一轮中途崩溃)/ `failed`(阶段一或阶段二失败);
`completed` 已在幂等前置分支返回,不会走到这里。
**为什么还要 catch IntegrityError**:`convert:idem:{cid}` 锁只包住幂等判定
(出块即释放),两笔同键请求可能**都判定为"无占位"**、各自生成了不同的 `group_id`
(T-13 真库实测:8 路并发下偶发);后插入的那笔撞 `uk_idem` → 此前会直穿 503。
这里把它收敛为 False(让路),交由调用方回 202(架构 §9「同键并发 → 202」)。
"""
try:
with self._engine.begin() as conn:
existing = conn.execute(
text("SELECT 1 FROM risk_convert_detail WHERE convert_group_id = :gid"),
{"gid": group_id},
).first()
if existing is not None:
conn.execute(
text(
"UPDATE risk_convert_detail"
" SET status = :s WHERE convert_group_id = :gid"
),
{"s": _STATUS_PENDING, "gid": group_id},
)
return True
conn.execute(
text(
"""
INSERT INTO risk_convert_detail
(convert_group_id, client_request_id, status, estimated)
VALUES (:gid, :cid_req, :status, 0)
"""
),
{"gid": group_id, "cid_req": client_request_id, "status": _STATUS_PENDING},
)
except IntegrityError:
# uk_idem 竞态:同键的另一笔刚刚占位成功 → 本笔不持有,让路
logger.info(
"幂等键已被并发请求占用(本次让路):cid=%s gid=%s",
client_request_id,
group_id,
)
return False
return True
# ---------- 阶段二:回写 completed + 详情 ----------
def complete_convert(
self,
group_id: str,
*,
out_trade_id: str,
in_trade_id: str,
related_trade_id: str | None = None,
nav: Decimal | None = None,
nav_date: date | None = None,
fee_amount: Decimal | None = None,
hold_days_min: int | None = None,
hold_days_max: int | None = None,
nav_stale: bool = False,
) -> None:
"""⛔ **Deprecated(v1.0 遗留,T-7 全量重写后删除)**:阶段二回写 completed。
新模型由 `sync_mirror`(确认 ② completed 镜像)取代;本方法仅供 T-7 完工前
v1.0 流回归使用。
"""
params = {
"status": _STATUS_COMPLETED,
"out": out_trade_id,
"in": in_trade_id,
"rel": related_trade_id,
"nav": _to_bind(nav),
"nav_date": nav_date,
"fee": _to_bind(fee_amount),
"hmin": hold_days_min,
"hmax": hold_days_max,
"stale": 1 if nav_stale else 0,
"gid": group_id,
"cid_req": None,
}
with self._engine.begin() as conn:
existing = conn.execute(
text("SELECT 1 FROM risk_convert_detail WHERE convert_group_id = :gid"),
{"gid": group_id},
).first()
if existing is not None:
conn.execute(
text(
"""
UPDATE risk_convert_detail
SET status = :status,
out_trade_id = :out,
in_trade_id = :in,
related_trade_id = :rel,
nav = :nav,
nav_date = :nav_date,
fee_amount = :fee,
hold_days_min = :hmin,
hold_days_max = :hmax,
nav_stale = :stale
WHERE convert_group_id = :gid
"""
),
params,
)
else:
conn.execute(
text(
"""
INSERT INTO risk_convert_detail
(convert_group_id, client_request_id, status, estimated,
out_trade_id, in_trade_id, related_trade_id, nav, nav_date,
fee_amount, hold_days_min, hold_days_max, nav_stale)
VALUES (:gid, :cid_req, :status, 0,
:out, :in, :rel, :nav, :nav_date,
:fee, :hmin, :hmax, :stale)
"""
),
params,
)
def mark_failed(self, group_id: str) -> None:
"""⛔ **Deprecated(v1.0 遗留,T-7 后删除)**:阶段一失败 → 占位置 failed。
新模型由 `sync_mirror(gid, 'failed')`(补偿 ④)取代。
"""
with self._engine.begin() as conn:
conn.execute(
text("UPDATE risk_convert_detail SET status = :s WHERE convert_group_id = :gid"),
{"s": _STATUS_FAILED, "gid": group_id},
)
# ---------- 查询 ----------
def get_by_group_id(self, group_id: str) -> dict[str, Any] | None:
"""按 convert_group_id 读单行(重试判定 / 阶段二补跑读取)。"""
with self._engine.connect() as conn:
row = conn.execute(
text("SELECT * FROM risk_convert_detail WHERE convert_group_id = :gid"),
{"gid": group_id},
).mappings().first()
return dict(row) if row else None
def get_by_client_request_id(self, client_request_id: str | None) -> dict[str, Any] | None:
"""按 client_request_id 读单行(幂等命中读取);None 直接返回 None(WHERE = NULL 永不命中)。"""
if client_request_id is None:
return None
with self._engine.connect() as conn:
row = conn.execute(
text("SELECT * FROM risk_convert_detail WHERE client_request_id = :cid_req"),
{"cid_req": client_request_id},
).mappings().first()
return dict(row) if row else None
# ---------- 清理(S2:标记不硬删) ----------
def list_expired_candidates(self, hours: int) -> list[dict[str, Any]]:
"""取超 `hours` 小时的 `pending` 孤儿(供 cleanup_pending_convert.py 巡检)。
cutoff = 当前本地时间 - hours;sqlite 默认以 localtime 落 created_at,
与 Python datetime.now()(本地)口径一致(B5 评审 P3-4 同口径)。
"""
cutoff = datetime.now() - timedelta(hours=hours)
with self._engine.connect() as conn:
return [
dict(r)
for r in conn.execute(
text(
"SELECT * FROM risk_convert_detail "
"WHERE status = :s AND created_at < :cutoff "
"ORDER BY created_at ASC"
),
{"s": _STATUS_PENDING, "cutoff": cutoff},
).mappings()
]
def mark_expired(self, group_id: str) -> None:
"""超时 pending 孤儿 → 置 `status='expired'`(标记不硬删,留痕供对账,S2)。"""
with self._engine.begin() as conn:
conn.execute(
text("UPDATE risk_convert_detail SET status = :s WHERE convert_group_id = :gid"),
{"s": _STATUS_EXPIRED, "gid": group_id},
)