Files
lzf_0626 a337cca31a 修复 T002 下单的两个 P0:无幂等保护、读账户/持仓无行锁
## P0-1 下单没有幂等保护(实测复现过重复扣款)

`app/api/controllers/trading.py` 的 T002 此前**没有任何幂等保护** —— 全仓 7 个
controller 共 29 处声明了 `Idempotency-Key`,**唯独 trading.py 一处都没有**,
而 `docs/05` §5.2 要求写端点必须带。

**修复前实测**(同一 key 连发两次):

    两个不同 order_no:SO...FD48A488 / SO...2FB182C0
    委托 +2、可用现金 -910.50(2×455.25)、510300 持仓 3500 -> 3700
    ⇒ 重复扣款 + 重复建仓

**改法**:接 `ApiTransactionService.execute_in` —— 业务写入与幂等回执**同一个事务**
(`docs/05` §5.2),同键同 body 的第二次请求直接回放上次响应。

- `trade_service.submit_order` 新增 `commit: bool = True`:包装层调用时传
  `commit=False`,由 `execute_in` 统一提交。**不做这一步就会"内层提交外层事务"**,
  幂等回执与业务写入分处两个事务,回执写失败时业务已落库,重放失去意义。
- 响应经 `model_dump(mode="json")` 统一 JSON 化后再 `model_validate` 还原 ——
  保证"首次"与"重放"两条路径返回**同一形状**,且对外结构不变
  (金额仍是字符串化 Decimal)。

**修复后实测**(同一 key 连发两次):

    两次 order_no 完全相同:SO202609141216512C7166D4
    两次响应逐字段相同(真正的回放)
    可用现金 -455.25(单次)、510300 持仓 +100(单次)
    T003 核对:43 条委托里只多出 1 条

## P0-3 下单读账户/持仓无行锁(并发可扣穿余额)

`trade_service.py` 全文 **零 `with_for_update`**,而项目其它 15 个 service 共 38 处
用了它 —— 规范早已建立,这里是遗漏。

**改法**:`_load_account` / `_load_holding` 增加 `for_update` 开关(默认关,
只读路径不加锁、不牺牲并发),`submit_order` 以 `for_update=True` 调用。
**加锁顺序固定「账户 → 持仓」**:并发事务按同一顺序取锁才不会成环,
这一点写在两处 docstring 里,改顺序前必须先想清楚。

**实测**(真并发 2 笔,各 45520 元,合计 91040 > 余额 49103):

    第 1 笔:201 成交 SO...59568811
    第 2 笔:422 INSUFFICIENT_FUNDS「可用余额 3578.91 不足」
    成交 1/2,最终余额 3578.91 >= 0

**关键证据**:被拒那笔读到的是 **3578.91(已扣减后)**而不是初始的 49103.46 ——
证明两个事务被行锁串行化了。无锁时两笔都会读到 49103.46 而双双通过,
余额会变成 -41936。

## 实测汇总

- `pytest tests/unit tests/contract` -> **1427 passed, 2 skipped, 2 failed**
  (那 2 个是既有的:投顾页面被替换、docs/05 §19 分组行,均与本提交无关)
- `ruff check` -> All checks passed
- 两次实测的完整证据见上;两份验证脚本在 %TEMP%(未入库)

## 顺带修正一个我自己脚本的 bug

验证脚本里用 `GET /users/me/orders?limit=200` 读委托数,而 T003 的 `limit` 上限是
**100**(`Query(le=100)`)→ 422 → 读到 0 条,一度让我误判"委托没增加"。
改用 `limit=100` 后确认委托确实只 +1。
2026-09-14 20:19:55 +08:00

199 lines
8.4 KiB
Python
Raw Permalink 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.
"""§T 用户自助场内基金模拟交易 controller(`docs/05` §19 T 段)。
端点与权限码(9 个端点 / 5 个权限码):
| § | 端点 | 权限码 | 摘要 |
|---|---|---|---|
| T001 | `GET /api/v1/users/me/account/dashboard` | `account:read:self` | 我的账户看板 |
| T002 | `POST /api/v1/users/me/orders` | `trade:order:create` | 提交委托(首版市价立即成交) |
| T003 | `GET /api/v1/users/me/orders` | `trade:order:read` | 委托列表 |
| T004 | `GET /api/v1/users/me/orders/{order_no}` | `trade:order:read` | 委托详情 |
| T005 | `POST /api/v1/users/me/orders/{order_no}/cancellations` | `trade:order:cancel` | 撤单 |
| T006 | `GET /api/v1/users/me/holdings` | `holding:read:self` | 持仓列表 |
| T007 | `GET /api/v1/users/me/transactions` | `trade:txn:read` | 成交记录列表 |
| T008 | `GET /api/v1/users/me/transactions/{txn_no}` | `trade:txn:read` | 成交详情 |
| T009 | `GET /api/v1/users/me/cash-ledger` | `account:read:self` | 资金明细 |
设计要点:
- 全部走 `build_request_context`(与 memory / portfolio 一致),数据范围 `self`。
- 不走限流依赖(`enforce_rate_limit`)——场内交易为低频,由底座网关层限流。
- 信封用 `envelope` / `list_envelope`,与 §3.3 一致。
"""
from __future__ import annotations
from fastapi import APIRouter, Depends, Header, Query, status
from sqlalchemy.ext.asyncio import AsyncSession
from app.api.dependencies.auth import build_request_context
from app.api.dependencies.database import get_session
from app.api.schemas.trading import OrderCreateRequest, OrderCreateResponse
from app.api.views.envelope import envelope, list_envelope
from app.core.contracts import RequestContext
from app.service.api_transaction_service import ApiTransactionService
from app.service.authorization_service import AuthorizationService
from app.service.suitability_service import SuitabilityService
from app.service.trade_service import TradeService
router = APIRouter(prefix="/api/v1/users/me", tags=["trading"])
#: T002 的幂等作用域。与权限码同名,便于审计时一眼对上是哪个写操作。
ORDER_SCOPE = "trade:order:create"
def _service(session: AsyncSession, context: RequestContext) -> TradeService:
return TradeService(session, suitability_evaluator=SuitabilityService())
async def _authorize(context: RequestContext, permission: str) -> None:
"""Enforce the endpoint permission declared in docs/05 before DB work."""
await AuthorizationService.require(context, permission)
# T001 账户看板
@router.get("/account/dashboard")
async def get_account_dashboard(
context: RequestContext = Depends(build_request_context), # noqa: B008
session: AsyncSession = Depends(get_session), # noqa: B008
) -> dict[str, object]:
await _authorize(context, "account:read:self")
data = await _service(session, context).get_account_dashboard(context)
return envelope(data, context)
# T002 提交委托(市价立即成交)
#
# ⚠️ 幂等:本端点此前**没有任何幂等保护** —— 全仓 7 个 controller 共 29 处声明了
# `Idempotency-Key`,唯独 trading.py 一处都没有(`docs/05` §5.2 要求写端点必须带)。
# 后果是**实测复现**过的:同一笔意愿在网关超时后被客户端按标准重试重发,
# 两次请求各自生成新的 `order_no`、各扣一次款、各建一次仓
# (实测:可用现金 -910.50 = 2×455.25,持仓 3500 -> 3700)。
#
# 现在接 `ApiTransactionService.execute_in`:业务写入与幂等回执**同一个事务**,
# 同键同 body 的第二次请求直接回放上次响应,不再重复成交。
@router.post("/orders", status_code=status.HTTP_201_CREATED)
async def submit_order(
payload: OrderCreateRequest,
context: RequestContext = Depends(build_request_context), # noqa: B008
session: AsyncSession = Depends(get_session), # noqa: B008
key: str | None = Header(default=None, alias="Idempotency-Key"),
) -> dict[str, object]:
await _authorize(context, ORDER_SCOPE)
async def action(inner: AsyncSession) -> dict[str, object]:
# `commit=False`:提交由 `execute_in` 统一做,业务写入与幂等回执同事务。
response = await _service(inner, context).submit_order(payload, context, commit=False)
# 统一 JSON 化。`execute_in` 回放的是**已落库的 dict**,首次返回值必须是
# 同一形状,否则"重放"与"首次"两种路径返回的字段类型会不一致。
return response.model_dump(mode="json")
data = await ApiTransactionService().execute_in(
session,
context,
ORDER_SCOPE,
key,
payload.model_dump(mode="json"),
action,
)
# 还原成响应模型再套信封:保持与改动前**完全一致**的响应结构
# (金额仍是字符串化 Decimal,见 docs/05 §3.3)。
return envelope(OrderCreateResponse.model_validate(data), context)
# T003 委托列表
@router.get("/orders")
async def list_orders(
limit: int = Query(default=20, ge=1, le=100),
cursor: str | None = Query(default=None),
context: RequestContext = Depends(build_request_context), # noqa: B008
session: AsyncSession = Depends(get_session), # noqa: B008
) -> dict[str, object]:
await _authorize(context, "trade:order:read")
cursor_id = int(cursor) if cursor else None
items, next_cursor = await _service(session, context).list_orders(
context, limit=limit, cursor=cursor_id
)
return list_envelope(
{"items": items, "next_cursor": next_cursor, "has_more": next_cursor is not None},
context,
)
# T004 委托详情
@router.get("/orders/{order_no}")
async def get_order(
order_no: str,
context: RequestContext = Depends(build_request_context), # noqa: B008
session: AsyncSession = Depends(get_session), # noqa: B008
) -> dict[str, object]:
await _authorize(context, "trade:order:read")
data = await _service(session, context).get_order(order_no, context)
return envelope(data, context)
# T005 撤单
@router.post("/orders/{order_no}/cancellations", status_code=status.HTTP_200_OK)
async def cancel_order(
order_no: str,
context: RequestContext = Depends(build_request_context), # noqa: B008
session: AsyncSession = Depends(get_session), # noqa: B008
) -> dict[str, object]:
await _authorize(context, "trade:order:cancel")
order = await _service(session, context).cancel_order(order_no, context)
return envelope(order, context)
# T006 持仓列表
@router.get("/holdings")
async def list_holdings(
context: RequestContext = Depends(build_request_context), # noqa: B008
session: AsyncSession = Depends(get_session), # noqa: B008
) -> dict[str, object]:
await _authorize(context, "holding:read:self")
data = await _service(session, context).list_holdings(context)
return envelope(data, context)
# T007 成交记录列表
@router.get("/transactions")
async def list_transactions(
limit: int = Query(default=20, ge=1, le=100),
cursor: str | None = Query(default=None),
context: RequestContext = Depends(build_request_context), # noqa: B008
session: AsyncSession = Depends(get_session), # noqa: B008
) -> dict[str, object]:
await _authorize(context, "trade:txn:read")
cursor_id = int(cursor) if cursor else None
data = await _service(session, context).list_transactions(
context, limit=limit, cursor=cursor_id
)
return envelope(data, context)
# T008 成交详情
@router.get("/transactions/{txn_no}")
async def get_transaction(
txn_no: str,
context: RequestContext = Depends(build_request_context), # noqa: B008
session: AsyncSession = Depends(get_session), # noqa: B008
) -> dict[str, object]:
await _authorize(context, "trade:txn:read")
item = await _service(session, context).get_transaction(txn_no, context)
return envelope(item, context)
# T009 资金明细
@router.get("/cash-ledger")
async def list_cash_ledger(
limit: int = Query(default=20, ge=1, le=100),
cursor: str | None = Query(default=None),
context: RequestContext = Depends(build_request_context), # noqa: B008
session: AsyncSession = Depends(get_session), # noqa: B008
) -> dict[str, object]:
await _authorize(context, "account:read:self")
cursor_id = int(cursor) if cursor else None
data = await _service(session, context).list_cash_ledger(
context, limit=limit, cursor=cursor_id
)
return envelope(data, context)