Commit aef9e37c by 沈彬

feat(account): 账户层级与分桶限额、分账月账、对外接口文档,并修复审查三项与并发对账竞态

应用(app/)
- 新增 account_admission.py:可用性查询与真实调用准入共用同一套主体判定(账户状态、网关绑定、
  赠送门闩、未核清调用、主体级联停用、生效额度、当月用量、扣费桶余额)
- 号段分配前移到业务锁之前(_mutate / 开户 provisioning / 赠送入账),发号器改用独立连接池
  (db_generator_pool_size、db_generator_max_overflow;web 与 worker 双装配)
- 凭证列表改分页 + status/subAccountId/scopeId 筛选,回执补 scopeId/subAccountId/scopeType
- 并发同 requestId 对账:命中重放时重读业务行、版本不一致时在新事务复查幂等行,
  消除"已成功却返回过期 version/phase"与误判版本冲突
- 门店/员工 scope 与员工自然月额度、分账维度索引、月账单直读、超额告警

测试(tests/)
- 新增:review_findings_repro(三项修复护栏)、test_account_scope、test_account_sub_account、
  test_account_usage_month、test_account_ledger_scope、test_account_overage_alerts、
  test_account_scope_e2e
- test_account_gifts 补并发对账重放护栏用例

文档(docs/、README)
- external-api-reference.md(35 个操作全覆盖,路由自应用导出机械校验)、external-api-flow.html(含 SSE)
- architecture-code-quality-risks-20260930.md(含修复记录 §8);设计文档 / p0 契约 / 开发计划同步
- DDL 与对账 SQL 归档至 scripts/sql/,新增 scripts/p1_mysql_instance.sh 隔离 MySQL 实例脚本

回归(python3 -m pytest -q -W error::sqlalchemy.exc.SAWarning)
- 隔离 MySQL 8.0.43:2211 用例 / 2181 通过 / 0 失败 / 30 跳过 / 207s
- SQLite:2211 用例 / 1677 通过 / 0 失败 / 534 跳过 / 69s
parent f12d2992
......@@ -6,7 +6,7 @@
P1 数据底座、P2 开户/凭证、P3 赠送/账单与 P4 模型调用/结算补偿已完成开发和本地验证,目标环境验收待补。独立入口提供固定角色权限、开户恢复、凭证生命周期、来源授权、账户停启用、精确赠送、受控核清、账单查询、模型调用与租约补偿;P5 起的数据导入/切流未开始。
- `app/account_models.py`:11 张 `t_computing_*` 独立表,无 SaaS 主数据依赖;数据库默认时间使用 UTC 毫秒表达式,不依赖导入会话时区。
- `app/account_models.py`:15 张 `t_computing_*` 独立表(含子账号、门店/员工及月额度),无 SaaS 主数据依赖;数据库默认时间使用 UTC 毫秒表达式,不依赖导入会话时区。
- `app/account_repository.py`:账户/来源作用域查询、账户锁和未清调用查询。
- `app/db.py:create_mysql_engine`:显式 READ COMMITTED、UTC、严格模式;由调用方显式传入连接,尚未改部署配置。
- `app/id_generator.py`:新库调用方显式选择独立 IdSegment;种子只向前推进,fork 子进程重建发号锁并丢弃余段。调用方必须在子进程内创建 Engine/Session 工厂,不复用父进程连接池或 Session。
......@@ -17,7 +17,7 @@ P1 数据底座、P2 开户/凭证、P3 赠送/账单与 P4 模型调用/结算
`app/account_main.py` 是独立入口,仅加载新的账户接口,不存在 auth_mode/disabled 或可信 X-App-Id。默认监听 loopback;实际管理端网络隔离、TLS/受控加密通道和外置文件权限由部署方配置,本轮未修改。
外置 JSON 文件必填 `database_url`、`newapi_base_url`(HTTPS)、`newapi_pat`、`newapi_model`;可选 `db_pool_size`、`db_max_overflow`、`management_timeout_seconds`、`management_concurrency`、`connect_timeout_seconds`、`read_timeout_seconds`。不从环境变量补值,未知字段或缺失配置拒绝启动,真实配置文件必须位于项目目录外。
外置 JSON 文件必填 `database_url`、`newapi_base_url`(HTTPS)、`newapi_pat`、`newapi_model`;可选 `db_pool_size`、`db_max_overflow`、`db_generator_pool_size`、`db_generator_max_overflow`、`management_timeout_seconds`、`management_concurrency`、`connect_timeout_seconds`、`read_timeout_seconds`。发号器用独立小池(默认 2+2):号段分配不占用业务连接池,写事务不会为取号去抢自己的连接。不从环境变量补值,未知字段或缺失配置拒绝启动,真实配置文件必须位于项目目录外。
```bash
python -m app.account_main --config /path/outside-project/account.json
......@@ -57,23 +57,25 @@ P3 复用已有物理表,无新增 SQL,无需重跑建表;未修改真实
## 文档
- **[对外接口文档](docs/external-api-reference.md)**:给外部接入方(SaaS/B 端)的完整接口参考——鉴权、幂等、错误码、35 个操作、层级与限额语义、对账口径、对接纪律。**对外交付用这一份**。
- **[接入全链路流程图](docs/external-api-flow.html)**:开户→子账号→门店/员工→签发 Key→调用(含 SSE 与失败分支)可视化,附 OpenAI 协议符合性矩阵。可独立打开或随文档交付。
- [P0 契约](docs/p0-contract-freeze.md):字段、状态、权限与幂等域,含 P1–P4 细化。
- [账户层级与门店/员工限额设计](docs/account-hierarchy-and-scope-limits-design.md):段 1–3 的设计与取舍(子账号/门店/员工/额度/分账/月账单)。
- [独立化开发计划](docs/independent-account-development-plan.md):P0–P8、验收、迁移和单写切流。
- [旧设计](docs/design.md)、[旧 RSA 协议](docs/authentication.md):历史参考,商户耦合和公网路线已被取代。
## P1 SQL 交付
## SQL 文件与建库
全部脚本由用户审阅并手动执行。本轮未连接业务数据库,未执行迁移、授予权限或设置真实配置。
系统仍在开发、尚未上线。SQL 统一放在 [`scripts/sql/`](scripts/sql/README.md),新环境只使用当前全量基线,不再叠加历史增量。
| 文件 | 执行位置与用途 |
| --- | --- |
| `scripts/20260923_independent_account_ddl.sql` | 新独立 schema,创建 11 张表;不建库、不改权限、不导入数据、不初始化号段 |
| `scripts/20260923_account_source_precheck.sql` | 旧 SaaS 库,只读检查实际结构、来源、孤儿关联、余额、未清状态与高水位 |
| `scripts/20260923_account_target_check.sql` | 新独立库,只读检查余额、调用/消费/流水关联、门闩和号段 |
| `scripts/sql/schema.sql` | 空的独立算力库,创建当前 15 张表及索引、约束;由 ORM 生成 |
| `scripts/sql/check.sql` | 当前独立库,只读核对账本、余额桶、员工月用量、关联和号段 |
| `scripts/sql/saas/` | SaaS 接入/迁移材料,按实际对接方案单独核对;不随网关初始化执行 |
| `scripts/sql/archive/` | 历史基线、增量及旧检查原文,仅供追溯 |
MySQL 最低 8.0.16,目标版本上线前须确认。DDL 不使用 IF NOT EXISTS,部分失败必须先检查,不使用客户端 `--force` 继续执行。迁移需停旧分配器,复制旧已提交高水位;只有真正空的新安装才允许从 100000000000000000 初始化。不能用 MAX(id)+1,不能在未知数据情况下直接执行低位种子。
本阶段不提供自动导入/一键切流;实际历史数据和 SaaS 映射必须先经过源库预检。已经执行过的 `20260922_computing_stage0_ddl.sql` 不改动。
建库、号段初始化、结果判读和开发变更流程见 [SQL 使用说明](scripts/sql/README.md)。最低 MySQL 8.0.16;全量 DDL 不用于升级已有数据库,不使用 `--force` 跳过错误。脚本不自动建库、授权、导入数据或重置号段。
## 开发验证
......@@ -81,10 +83,10 @@ MySQL 最低 8.0.16,目标版本上线前须确认。DDL 不使用 IF NOT EXIS
```bash
python -m pytest -q
python scripts/render_account_ddl.py --output scripts/20260923_independent_account_ddl.sql
python scripts/render_account_ddl.py --output scripts/sql/schema.sql
```
可选 `--account-mysql-socket=<临时实例socket>` 验证真实 MySQL;测试仅接受父目录名以 `computing-p1-mysql.` 开头的专用实例,为每个测试创建随机独立库并清理,不应指向已有业务实例。未提供时 MySQL 分支跳过,SQLite 结果不能替代行锁验收。
可选 `--account-mysql-socket=<临时实例socket>` 验证真实 MySQL;测试仅接受父目录名以 `computing-p1-mysql.` 开头的专用实例,为每个测试创建随机独立库并清理,不应指向已有业务实例。未提供时 MySQL 分支跳过,SQLite 结果不能替代行锁验收。`scripts/p1_mysql_instance.sh start|stop|status` 负责起停这套隔离实例(unix socket、不开端口、READ COMMITTED + UTC + 严格模式),start 会打印可直接复制的回归命令。
2026-09-23 本地结果:Python 3.9.6 + 无网络监听的临时 MySQL 9.6.0,168 项通过、14 项跳过(均为 SQLite 不适用的 MySQL 专项分支,对应 MySQL 分支通过)。MySQL 测试直接执行交付 DDL,覆盖三元幂等、关联约束、UTC 默认时间、账户行锁、并发号段与持锁 fork、负余额及目标对账异常。Python 3.12、实际部署 MySQL 版本与运行账号无 SaaS 权限仍需补验,不能据此认定 P1 最终验收或迁移上线通过。
......
"""调用主体准入判定:可用性查询与真实调用共用同一套规则。
修复前 `AccountService._account_view`(可用性)与 `CallService._admit`(调用准入)各写一套:
可用性只读账户行(含**账户未分配桶**余额),准入按 Key 主体取桶(子账号桶、门店/员工状态、
员工月度额度),于是同一个 Key 会得到互相矛盾的结论:子账号没余额但总账号有余额时返回
available=true、调用失败;反过来返回不可用、调用成功。
这里把判定抽成一份实现,可用性与准入共用。判定顺序与调用准入原实现保持一致,错误码与
data 不变:
账户状态 → 网关绑定(由调用侧传入)→ 赠送门 → 未决调用 → 主体(门店/员工)→ 扣费桶余额
两侧的区别只在**取桶方式**:调用侧 `for_update=True`(锁桶后再扣费),可用性查询走普通读。
"""
from dataclasses import dataclass, field
from typing import List, Optional
from . import account_repository as repository
from .account_security import AccountError
from .billing import EXPECTED_QUOTA_PER_UNIT
ACCOUNT_BUCKET = "ACCOUNT"
SUB_ACCOUNT_BUCKET = "SUB_ACCOUNT"
@dataclass
class Admission:
"""判定结果:`reasons` 供 `blockedReasons` 出参,`error` 供调用侧直接抛。"""
reasons: List[str] = field(default_factory=list)
error: Optional[AccountError] = None
bucket_type: str = ACCOUNT_BUCKET
bucket_id: Optional[int] = None
balance_point_units: int = 0
@property
def allowed(self) -> bool:
return not self.reasons
def _block(self, reason, error):
self.reasons.append(reason)
if self.error is None:
self.error = error
return self
def subject_of(session, principal, month):
"""解析凭证绑定的门店/员工主体;主体悬空按停用处理(与调用准入一致)。"""
if principal.scope_id is None:
return None
subject = repository.resolve_call_subject(session, principal.scope_id, month)
if subject is None or subject.account_id != principal.account_id:
raise AccountError("SCOPE_DISABLED", data={"scopeId": str(principal.scope_id)})
return subject
def subject_bucket(subject):
"""主体扣费桶:门店用自身归属(NULL = 直接落在账户未分配桶),员工走当前门店。"""
return subject.store_bucket if subject.scope_type == "STORE" else subject.parent_bucket
def _judge_subject(subject, month):
"""门店/员工主体判定,返回 (reason, error);无阻断时返回 (None, None)。"""
if subject.scope_type == "STORE":
if subject.scope_status != "ACTIVE":
return "SCOPE_DISABLED", AccountError("SCOPE_DISABLED",
data={"scopeId": str(subject.scope_id)})
return None, None
# 员工只经当前门店可达:门店停用即全员下线。
if (subject.scope_status != "ACTIVE" or subject.parent_id is None
or subject.parent_type != "STORE"):
return "SCOPE_DISABLED", AccountError("SCOPE_DISABLED",
data={"scopeId": str(subject.scope_id)})
if subject.parent_status != "ACTIVE":
return "SCOPE_DISABLED", AccountError("SCOPE_DISABLED",
data={"scopeId": str(subject.scope_id),
"parentScopeId": str(subject.parent_id)})
quota = subject.current_quota
if quota is not None and subject.used_point_units >= quota:
return "EMPLOYEE_QUOTA_EXCEEDED", AccountError("EMPLOYEE_QUOTA_EXCEEDED", data={
"scopeId": str(subject.scope_id), "quotaMonth": month,
"limitPointUnits": str(quota), "usedPointUnits": str(subject.used_point_units)})
return None, None
def binding_error(account, settings, model=None):
"""网关绑定校验:provision 状态、令牌、网关身份、计量口径(可选再校验请求模型)。
调用侧传具体模型;可用性查询只校验"账户已绑定可用令牌且计量口径一致",因此查询侧
不传 model —— 两边用的是同一条规则。
"""
snapshot = account.gateway_snapshot or {}
if (account.provision_status != "READY" or not account.gateway_token_id
or not account.gateway_token_name
or snapshot.get("gatewayId") != settings.gateway_identity
or snapshot.get("quotaPerUnit") != EXPECTED_QUOTA_PER_UNIT):
return AccountError("DEPENDENCY_UNAVAILABLE")
if model is None:
return None if snapshot.get("model") else AccountError("DEPENDENCY_UNAVAILABLE")
return None if snapshot.get("model") == model else AccountError("DEPENDENCY_UNAVAILABLE")
def evaluate(session, account, *, principal=None, sub_account_id=None, subject=None, month=None,
for_update=False, uncertain=None, binding_error=None):
"""按调用主体口径判定账户能否调用,并给出真正扣费的那个桶及其余额。
`binding_error` 是调用侧算好的网关绑定错误(`_admit` 在账户状态之后、赠送门之前
校验绑定,顺序保持原样);可用性查询不需要绑定校验,传 None。
`principal` 供查询侧直接传凭证主体:内部解析门店/员工,主体悬空收成
SCOPE_DISABLED 原因(查询不该因为主体悬空直接报错)。
"""
admission = Admission()
subject_error = None
if principal is not None:
if sub_account_id is None:
sub_account_id = principal.sub_account_id
if subject is None and principal.scope_id is not None:
try:
subject = subject_of(session, principal, month)
except AccountError as exc:
subject_error = exc
if account.status != "ACTIVE":
admission._block("ACCOUNT_DISABLED", AccountError("ACCOUNT_DISABLED"))
if binding_error is not None:
admission._block("DEPENDENCY_UNAVAILABLE", binding_error)
if account.gift_gate is not None:
admission._block("GIFT_GATE_BUSY", AccountError("GIFT_GATE_BUSY"))
if uncertain is None:
uncertain = repository.count_uncertain_calls(session, account.id) > 0
if uncertain:
admission._block("UNRESOLVED_CALL",
AccountError("ACCOUNT_BLOCKED", reason="UNRESOLVED_CALL"))
if subject_error is not None:
admission._block("SCOPE_DISABLED", subject_error)
elif subject is not None:
reason, error = _judge_subject(subject, month)
if reason is not None:
admission._block(reason, error)
bucket_id = subject_bucket(subject) if subject is not None else sub_account_id
admission.bucket_id = bucket_id
if bucket_id is None:
admission.bucket_type = ACCOUNT_BUCKET
admission.balance_point_units = account.balance_point_units
if account.balance_point_units <= 0:
admission._block("INSUFFICIENT_BALANCE",
AccountError("ACCOUNT_BLOCKED", reason="INSUFFICIENT_BALANCE"))
return admission
admission.bucket_type = SUB_ACCOUNT_BUCKET
# A sub-account key never draws from the unallocated bucket; a zero
# unallocated balance on the parent must not block it.
if for_update:
sub = repository.lock_sub_account(session, bucket_id)
else:
sub = repository.find_sub_account(session, account.id, bucket_id)
if sub is None or sub.account_id != account.id:
return admission._block("SUB_ACCOUNT_NOT_FOUND",
AccountError("ACCOUNT_BLOCKED", reason="SUB_ACCOUNT_NOT_FOUND"))
admission.balance_point_units = sub.balance_point_units
if sub.status != "ACTIVE":
return admission._block("SUB_ACCOUNT_DISABLED", AccountError("SUB_ACCOUNT_DISABLED"))
if sub.balance_point_units <= 0:
return admission._block("INSUFFICIENT_BALANCE", AccountError(
"ACCOUNT_BLOCKED", reason="INSUFFICIENT_BALANCE", data={"subAccountId": str(sub.id)}))
return admission
......@@ -7,12 +7,13 @@ import anyio
from sqlalchemy.exc import IntegrityError
from . import account_repository as repository
from . import account_admission as admission
from .account_config import _valid_model
from .account_execution import run_database
from .account_gateway import GatewayError
from .account_model_gateway import StreamInterrupted
from .account_models import Call
from .account_security import AccountError, authenticate, utc_now, validate_request_id
from .account_security import AccountError, authenticate, quota_month, utc_now, validate_request_id
from .account_service import _time, fingerprint
from .account_settlement import SettlementService
from .billing import EXPECTED_QUOTA_PER_UNIT, format_points
......@@ -28,6 +29,26 @@ def _string(value):
return str(value) if value is not None else None
def _quota_month(now):
"""Registration month in the business timezone: the monthly usage bucket has
to flip at Beijing midnight, not at UTC midnight."""
return quota_month(now)
def _snapshot_bucket(principal, subject):
if subject is not None:
return admission.subject_bucket(subject)
return principal.sub_account_id
def _snapshot_store(subject):
"""The store a call is attributed to: the store itself for a store key, the
store the employee belonged to at registration for an employee key."""
if subject is None:
return None
return subject.scope_id if subject.scope_type == "STORE" else subject.parent_id
def call_view(row, *, include_content=True):
available = row.result_json is not None and row.create_time > utc_now() - timedelta(days=7)
result = json.loads(row.result_json) if available and include_content else None
......@@ -96,21 +117,17 @@ class CallService:
"tokenName": account.gateway_token_name, "model": model,
"quotaPerUnit": snapshot.get("quotaPerUnit")}
def _admit(self, session, account, model):
if account.status != "ACTIVE":
raise AccountError("ACCOUNT_DISABLED")
def _subject(self, session, principal, month):
"""Resolve the credential's scope, if any, into the calling subject."""
return admission.subject_of(session, principal, month)
def _admit(self, session, account, model, sub_account_id=None, subject=None, month=None):
binding = self._binding(account, model)
if (account.provision_status != "READY" or not binding["tokenId"] or not binding["tokenName"]
or binding["gatewayId"] != self.settings.gateway_identity
or binding["quotaPerUnit"] != EXPECTED_QUOTA_PER_UNIT
or (account.gateway_snapshot or {}).get("model") != model):
raise AccountError("DEPENDENCY_UNAVAILABLE")
if account.gift_gate is not None:
raise AccountError("GIFT_GATE_BUSY")
if repository.count_uncertain_calls(session, account.id):
raise AccountError("ACCOUNT_BLOCKED", reason="UNRESOLVED_CALL")
if account.balance_point_units <= 0:
raise AccountError("ACCOUNT_BLOCKED", reason="INSUFFICIENT_BALANCE")
verdict = admission.evaluate(session, account, sub_account_id=sub_account_id,
subject=subject, month=month, for_update=True,
binding_error=admission.binding_error(account, self.settings, model))
if verdict.error is not None:
raise verdict.error
return binding
def _state(self, row):
......@@ -118,15 +135,20 @@ class CallService:
"version": row.version, "binding": row.gateway_error_snapshot["binding"],
"view": call_view(row)}
def _prepare(self, secret, request_id, body, deadline, gateway_params=None, stream_options=None):
def _prepare(self, secret, request_id, body, deadline=None, gateway_params=None, stream_options=None):
values = {"modelSelection": body.model, "businessCode": body.business_code,
"businessRef": body.business_ref, "gatewayParams": gateway_params,
"messages": [message.model_dump() for message in body.messages]}
if stream_options is not None:
values["stream"] = stream_options
digest = fingerprint(values)
with self.factory() as session:
principal = self._principal(session, secret)
if principal.sub_account_id is not None or principal.scope_id is not None:
# Only a keyed subject extends the fingerprint: account-level keys
# keep the exact legacy digest so historical replays stay compatible.
values["subject"] = {"subAccountId": principal.sub_account_id,
"scopeId": principal.scope_id}
digest = fingerprint(values)
original = repository.find_call(session, principal.account_id, principal.client_id, request_id)
if original is not None:
if original.fingerprint != digest:
......@@ -147,12 +169,19 @@ class CallService:
if (not _valid_model(model) or model != self.settings.newapi_model
or body.business_code is not None and body.business_code not in BUSINESS_CODES):
raise AccountError("INVALID_ARGUMENT")
binding = self._admit(session, account, model)
registered_at = utc_now()
month = _quota_month(registered_at)
subject = self._subject(session, principal, month)
binding = self._admit(session, account, model, principal.sub_account_id,
subject=subject, month=month)
remaining = _remaining(deadline)
row = Call(id=call_id, account_id=principal.account_id, client_id=principal.client_id,
request_id=request_id, fingerprint=digest, requested_model=body.model,
model=model, business_code=body.business_code, business_ref=body.business_ref,
dispatch_deadline=utc_now() + timedelta(seconds=remaining),
sub_account_id=_snapshot_bucket(principal, subject),
scope_id=principal.scope_id, store_scope_id=_snapshot_store(subject),
quota_month=month,
dispatch_deadline=registered_at + timedelta(seconds=remaining),
gateway_error_snapshot={"binding": binding})
session.add(row)
session.flush()
......@@ -175,7 +204,9 @@ class CallService:
raise AccountError("REQUEST_NOT_FOUND")
if row.dispatch_phase != "REGISTERED" or row.billing_status != "PROCESSING" or row.version != state["version"]:
return None
binding = self._admit(session, account, row.model)
subject = self._subject(session, principal, row.quota_month)
binding = self._admit(session, account, row.model, row.sub_account_id,
subject=subject, month=row.quota_month)
if binding != state["binding"]:
raise AccountError("DEPENDENCY_UNAVAILABLE")
_remaining(deadline)
......
......@@ -84,6 +84,8 @@ class AccountSettings:
newapi_model: str = field(repr=False)
db_pool_size: int = 5
db_max_overflow: int = 0
db_generator_pool_size: int = 2
db_generator_max_overflow: int = 2
management_timeout_seconds: float = 30.0
management_concurrency: int = 8
connect_timeout_seconds: float = 10.0
......@@ -115,6 +117,11 @@ class AccountSettings:
raise ValueError
if type(self.db_max_overflow) is not int or self.db_max_overflow < 0:
raise ValueError
if type(self.db_generator_pool_size) is not int or self.db_generator_pool_size <= 0:
raise ValueError
if (type(self.db_generator_max_overflow) is not int
or self.db_generator_max_overflow < 0):
raise ValueError
if (type(self.management_concurrency) is not int
or not 0 < self.management_concurrency <= 100):
raise ValueError
......
......@@ -10,7 +10,10 @@ from starlette.concurrency import run_in_threadpool
from . import account_repository as repository
from .account_gateway import GatewayError
from .account_ledger import gift_view
from .account_models import AccountClient, Client, Credential, Gift, ManagementRequest, OperationAudit, PointRecord
from .account_models import (
Account, AccountClient, Client, Credential, Gift, ManagementRequest, OperationAudit,
PointRecord,
)
from .account_security import AccountError, Principal, authenticate, utc_now, validate_request_id
from .account_service import fingerprint, _time
from .billing import (
......@@ -112,6 +115,17 @@ class GiftService:
"operator_note": body.operator_note})
for attempt in range(2):
try:
# 号段在业务锁之前领取:先只读确认,命中的重放直接回执(不依赖发号器)。
with self.factory() as probe:
account = probe.get(Account, int(body.account_id))
actor = self._platform(probe, secret)
self._target(probe, account, body.client_id, authorized=True)
replayed = repository.find_gift(probe, account.id, body.client_id, request_id)
if replayed is not None:
if replayed.fingerprint != digest:
raise AccountError("FINGERPRINT_MISMATCH")
return self._state(probe, replayed), False
ids = self._ids(2)
with self.factory.begin() as session:
account = repository.lock_account(session, int(body.account_id))
actor = self._platform(session, secret)
......@@ -121,9 +135,6 @@ class GiftService:
if row.fingerprint != digest:
raise AccountError("FINGERPRINT_MISMATCH")
return self._state(session, row), False
# Replay must not depend on the ID segment; allocate only now
# that a new gift row is actually required.
ids = self._ids(2)
self._admit(session, account, delta * POINT_UNITS_PER_QUOTA)
row = Gift(id=ids[0], account_id=account.id, client_id=body.client_id,
actor_credential_id=actor.credential_id, request_id=request_id,
......@@ -327,18 +338,44 @@ class GiftService:
def _reconcile_state(self, secret, request_id, gift_id, body):
with self.factory() as session:
actor = self._platform(session, secret)
row = session.get(Gift, gift_id)
if row is None or (row.account_id, row.client_id) != (int(body.account_id), body.client_id):
raise AccountError("REQUEST_NOT_FOUND")
# 幂等行先查、gift 行后读:并发同 requestId 的对账里,另一个事务可能刚提交
# 了这笔对账,而镜像里 gift 行还是提交前的版本。先读 gift 会把"已成功"的
# 重放回执配上过期视图(version 与 phase 倒退),这里改为命中重放后重读。
previous = repository.find_management_request(session, actor.client_id, "RECONCILE", request_id)
if previous is not None:
if previous.fingerprint != self._reconcile_digest(gift_id, body):
raise AccountError("FINGERPRINT_MISMATCH")
row = session.get(Gift, gift_id, populate_existing=True)
if row is None or (row.account_id, row.client_id) != (int(body.account_id), body.client_id):
raise AccountError("REQUEST_NOT_FOUND")
return None, dict(self._view(row), operationStatus=previous.status)
row = session.get(Gift, gift_id)
if row is None or (row.account_id, row.client_id) != (int(body.account_id), body.client_id):
raise AccountError("REQUEST_NOT_FOUND")
state = self._state(session, row)
if state["version"] != body.expected_version:
replayed = self._replay_if_committed(gift_id, body, request_id, actor.client_id)
if replayed is not None:
return None, replayed
self._check_reconcile(state, body)
return state, None
def _replay_if_committed(self, gift_id, body, request_id, client_id):
"""并发同 requestId 的对账刚提交时的兜底重放。
版本先于本次预检的快照变化、而幂等行在同一提交里,说明重放行已经落库:
在新事务里复查幂等行(提交过就一定能查到),命中则按重放回执返回最新视图,
避免既抛"版本冲突"又给出过期视图。
"""
with self.factory() as fresh:
previous = repository.find_management_request(fresh, client_id, "RECONCILE", request_id)
if previous is None:
return None
if previous.fingerprint != self._reconcile_digest(gift_id, body):
raise AccountError("FINGERPRINT_MISMATCH")
row = fresh.get(Gift, gift_id)
return dict(self._view(row), operationStatus=previous.status)
def _check_reconcile(self, state, body):
if state["version"] != body.expected_version:
raise AccountError("ACCOUNT_BLOCKED", reason="VERSION_CONFLICT")
......
......@@ -84,6 +84,8 @@ def openai_error(error: AccountError) -> JSONResponse:
error_type = "api_error"
elif error.code == "RATE_LIMITED":
error_type = "rate_limit_error"
elif error.code == "EMPLOYEE_QUOTA_EXCEEDED":
status, error_type = 429, "insufficient_quota"
elif error.code == "ACCOUNT_BLOCKED" and reason == "INSUFFICIENT_BALANCE":
status, error_type = 429, "insufficient_quota"
else:
......
from sqlalchemy import func, select
from sqlalchemy import and_, func, select
from sqlalchemy.orm import aliased
from .account_models import Account, AccountClient, Call, Consumption, Gift, ManagementRequest
from .account_models import (
Account,
AccountClient,
Call,
Consumption,
Gift,
ManagementRequest,
Scope,
ScopeMonthUsage,
ScopeQuota,
SubAccount,
)
from .account_security import utc_now
def get_account(session, account_id, client_id):
......@@ -24,6 +37,185 @@ def lock_account(session, account_id):
)
def find_sub_account(session, account_id, sub_account_id):
return session.scalar(
select(SubAccount).where(
SubAccount.id == sub_account_id,
SubAccount.account_id == account_id,
)
)
def lock_sub_account(session, sub_account_id):
return session.scalar(
select(SubAccount)
.where(SubAccount.id == sub_account_id)
.with_for_update()
.execution_options(populate_existing=True)
)
def page_sub_accounts(session, account_id, *, page=1, size=20):
if page < 1 or not 1 <= size <= 200:
raise ValueError("分页范围无效")
conditions = [SubAccount.account_id == account_id]
total = session.scalar(select(func.count()).select_from(SubAccount).where(*conditions))
rows = []
if (page - 1) * size < total:
rows = session.scalars(
select(SubAccount)
.where(*conditions)
.order_by(SubAccount.create_time.desc(), SubAccount.id.desc())
.offset((page - 1) * size)
.limit(size)
).all()
return total, list(rows)
def find_scope(session, account_id, scope_id):
return session.scalar(
select(Scope).where(Scope.id == scope_id, Scope.account_id == account_id)
)
def find_scope_by_key(session, scope_type, scope_key):
"""Scope keys are globally unique: one store key can only ever exist once,
which is what makes "a store belongs to exactly one account" structural."""
return session.scalar(
select(Scope).where(Scope.scope_type == scope_type, Scope.scope_key == scope_key)
)
def lock_scope(session, scope_id):
return session.scalar(
select(Scope)
.where(Scope.id == scope_id)
.with_for_update()
.execution_options(populate_existing=True)
)
def page_scopes(session, account_id, *, scope_type=None, sub_account_id=None, page=1, size=20):
if page < 1 or not 1 <= size <= 200:
raise ValueError("分页范围无效")
conditions = [Scope.account_id == account_id]
if scope_type is not None:
conditions.append(Scope.scope_type == scope_type)
if sub_account_id is not None:
conditions.append(Scope.sub_account_id == sub_account_id)
total = session.scalar(select(func.count()).select_from(Scope).where(*conditions))
rows = []
if (page - 1) * size < total:
rows = session.scalars(
select(Scope)
.where(*conditions)
.order_by(Scope.create_time.desc(), Scope.id.desc())
.offset((page - 1) * size)
.limit(size)
).all()
return total, list(rows)
def find_scope_quota(session, scope_id, store_scope_id):
return session.scalar(
select(ScopeQuota).where(
ScopeQuota.scope_id == scope_id, ScopeQuota.store_scope_id == store_scope_id
)
)
def find_month_usage(session, scope_id, month):
"""当月已用(缺失 = 0,与 UPSERT 的起点一致)。"""
value = session.scalar(
select(ScopeMonthUsage.used_point_units).where(
ScopeMonthUsage.scope_id == scope_id, ScopeMonthUsage.month == month
)
)
return 0 if value is None else value
def resolve_call_subject(session, scope_id, month):
"""Resolve a store/employee call subject in a single query: the scope row,
its parent store, the charging bucket and the effective monthly quota.
Effective quota comes from the (employee, current store) row and nothing
else: the limit is configured per store, so a store without a row of its own
simply has no limit for this employee, and value NULL means "explicitly
unlimited". There is deliberately no fallback to the employee's other rows -
the store the employee belongs to now decides, which is what "the limit
follows the store" means; usage still accumulates across stores.
"""
store = aliased(Scope)
current = aliased(ScopeQuota)
usage = aliased(ScopeMonthUsage)
stmt = (
select(
Scope.id.label("scope_id"),
Scope.account_id.label("account_id"),
Scope.scope_type.label("scope_type"),
Scope.status.label("scope_status"),
Scope.sub_account_id.label("store_bucket"),
store.id.label("parent_id"),
store.scope_type.label("parent_type"),
store.status.label("parent_status"),
store.sub_account_id.label("parent_bucket"),
current.monthly_quota_point_units.label("current_quota"),
func.coalesce(usage.used_point_units, 0).label("used_point_units"),
)
.select_from(Scope)
.outerjoin(store, store.id == Scope.parent_scope_id)
.outerjoin(
current,
and_(current.scope_id == Scope.id, current.store_scope_id == Scope.parent_scope_id),
)
.outerjoin(usage, and_(usage.scope_id == Scope.id, usage.month == month))
.where(Scope.id == scope_id)
)
return session.execute(stmt).first()
def upsert_month_usage(session, scope_id, month, units, *, now=None):
"""Atomically add ``units`` to an employee's month bucket.
The month is a state key, so a new month starts from zero without any reset
job; the atomic UPSERT (never read-then-insert) keeps two concurrent
settlements that both land in a brand new month from racing the unique key.
"""
if units <= 0:
return
now = now or utc_now()
table = ScopeMonthUsage.__table__
dialect = session.get_bind().dialect.name
if dialect == "mysql":
from sqlalchemy.dialects.mysql import insert as dialect_insert
stmt = dialect_insert(table).values(
scope_id=scope_id, month=month, used_point_units=units,
create_time=now, last_update_time=now,
)
stmt = stmt.on_duplicate_key_update(
used_point_units=table.c.used_point_units + stmt.inserted.used_point_units,
last_update_time=stmt.inserted.last_update_time,
)
elif dialect == "sqlite":
from sqlalchemy.dialects.sqlite import insert as dialect_insert
stmt = dialect_insert(table).values(
scope_id=scope_id, month=month, used_point_units=units,
create_time=now, last_update_time=now,
)
stmt = stmt.on_conflict_do_update(
index_elements=[table.c.scope_id, table.c.month],
set_={
"used_point_units": table.c.used_point_units + stmt.excluded.used_point_units,
"last_update_time": stmt.excluded.last_update_time,
},
)
else:
raise RuntimeError("unsupported dialect: %s" % dialect)
session.execute(stmt)
def find_call(session, account_id, client_id, request_id):
return session.scalar(
select(Call).where(
......
......@@ -50,6 +50,14 @@ class CreateAccountRequest(WriteRequest):
class IssueCredentialRequest(WriteRequest):
client_id: Optional[ClientId] = None
sub_account_id: Optional[Identifier] = None
scope_id: Optional[Identifier] = None
@model_validator(mode="after")
def exclusive_subject(self):
if self.sub_account_id is not None and self.scope_id is not None:
raise ValueError("subAccountId and scopeId are mutually exclusive")
return self
class ActivateCredentialRequest(WriteRequest):
......@@ -80,6 +88,66 @@ class SetAccountStatusRequest(WriteRequest):
reason: str = Field(min_length=1, max_length=200)
class CreateSubAccountRequest(WriteRequest):
name: str = Field(min_length=1, max_length=100)
remark: Optional[str] = Field(default=None, max_length=500)
@field_validator("name")
@classmethod
def nonblank_sub_account_name(cls, value):
if not value.strip():
raise ValueError("Blank name")
return value
class SetSubAccountStatusRequest(WriteRequest):
status: Literal["ACTIVE", "DISABLED"]
reason: str = Field(min_length=1, max_length=200)
class TransferRequest(WriteRequest):
points: int = Field(ge=1, le=1000000000)
direction: Literal["IN", "OUT"]
operator_note: Optional[str] = Field(default=None, max_length=200)
ScopeKey = Annotated[str, StringConstraints(strict=True, pattern=r"^[a-z0-9-]{1,30}:[A-Za-z0-9._:-]{1,69}$")]
class BindScopeRequest(WriteRequest):
scope_type: Literal["STORE", "EMPLOYEE"]
scope_key: ScopeKey
sub_account_id: Optional[Identifier] = None
parent_scope_key: Optional[ScopeKey] = None
@model_validator(mode="after")
def scope_shape(self):
if self.scope_type == "STORE":
if self.parent_scope_key is not None:
raise ValueError("Store scopes have no parent scope")
elif self.parent_scope_key is None:
raise ValueError("Employee scopes require a parent store key")
return self
class SetScopeStatusRequest(WriteRequest):
status: Literal["ACTIVE", "DISABLED"]
reason: str = Field(min_length=1, max_length=200)
class MoveStoreRequest(WriteRequest):
sub_account_id: Optional[Identifier] = None
class MoveEmployeeRequest(WriteRequest):
parent_scope_id: Identifier
class SetScopeQuotaRequest(WriteRequest):
store_scope_id: Identifier
monthly_quota_points: Optional[int] = Field(default=None, ge=0, le=1000000000)
class AccountQueryRequest(AccountRequest):
account_ids: Optional[List[Identifier]] = Field(default=None, min_length=1, max_length=1000)
page: int = Field(default=1, ge=1)
......@@ -123,12 +191,23 @@ class ReconcileGiftRequest(WriteRequest):
return value
class ScopeMonthPageRequest(AccountQueryRequest):
owner_client_id: Optional[ClientId] = None
month: Optional[str] = None
sub_account_ids: Optional[List[Identifier]] = Field(default=None, min_length=1, max_length=1000)
scope_ids: Optional[List[Identifier]] = Field(default=None, min_length=1, max_length=1000)
store_scope_ids: Optional[List[Identifier]] = Field(default=None, min_length=1, max_length=1000)
class LedgerPageRequest(AccountQueryRequest):
owner_client_id: Optional[ClientId] = None
client_id: Optional[ClientId] = None
request_id: Optional[RequestId] = None
start: Optional[datetime] = None
end: Optional[datetime] = None
sub_account_ids: Optional[List[Identifier]] = Field(default=None, min_length=1, max_length=1000)
scope_ids: Optional[List[Identifier]] = Field(default=None, min_length=1, max_length=1000)
store_scope_ids: Optional[List[Identifier]] = Field(default=None, min_length=1, max_length=1000)
@field_validator("start", "end", mode="before")
@classmethod
......
......@@ -5,6 +5,7 @@ import secrets
from dataclasses import dataclass
from datetime import datetime, timezone
from typing import Optional
from zoneinfo import ZoneInfo
from .account_models import Account, AccountClient, Client, Credential
from .errors import AppError
......@@ -18,26 +19,37 @@ ERRORS = {
"CREDENTIAL_REVOKED": (401, "凭证已撤销"),
"ACCOUNT_ACCESS_REVOKED": (403, "账户来源授权已撤销"),
"ACCOUNT_DISABLED": (403, "账户已停用"),
"SUB_ACCOUNT_DISABLED": (403, "子账号已停用"),
"SCOPE_DISABLED": (403, "门店或员工已停用"),
"PERMISSION_DENIED": (403, "无权执行此操作"),
"ACCOUNT_NOT_FOUND": (404, "账户不存在或不可见"),
"SUB_ACCOUNT_NOT_FOUND": (404, "子账号不存在或不可见"),
"REQUEST_NOT_FOUND": (404, "请求或凭证不存在或不可见"),
"FINGERPRINT_MISMATCH": (409, "相同请求号的参数与原请求不一致"),
"ACCOUNT_BLOCKED": (409, "当前状态不允许此操作"),
"SCOPE_CONFLICT": (409, "门店或员工标识已存在"),
"GIFT_GATE_BUSY": (409, "账户正被其他赠送占用"),
"RATE_LIMITED": (429, "请求数量超过处理容量"),
"EMPLOYEE_QUOTA_EXCEEDED": (429, "员工当月额度已用尽"),
"DEPENDENCY_UNAVAILABLE": (503, "依赖尚未就绪"),
"SERVICE_UNAVAILABLE": (503, "服务暂不可用,请查询原请求"),
}
KEY_PATTERN = re.compile(r"ck-([1-9][0-9]{17})-([A-Za-z0-9_-]{43})\Z")
REQUEST_PATTERN = re.compile(r"[!-~]{8,100}\Z")
SCOPE_KEY_PATTERN = re.compile(r"[a-z0-9-]{1,30}:[A-Za-z0-9._:-]{1,69}\Z")
class AccountError(AppError):
def __init__(self, code, *, reason=None):
def __init__(self, code, *, reason=None, data=None):
status, message = ERRORS[code]
super().__init__(code, message)
self.http_status = status
self.data = {"reason": reason} if reason else None
if data is not None:
self.data = dict(data)
if reason:
self.data.setdefault("reason", reason)
else:
self.data = {"reason": reason} if reason else None
@dataclass(frozen=True)
......@@ -46,12 +58,20 @@ class Principal:
client_id: str
category: str
account_id: Optional[int]
sub_account_id: Optional[int] = None
scope_id: Optional[int] = None
def utc_now():
return datetime.now(timezone.utc)
def quota_month(now=None):
"""业务月键 '\''YYYY-MM'\'':登记时间按 Asia/Shanghai 派生(P0 §1 唯一时区派生键,
月界 = 北京 0 点 = UTC 前一日 16:00)。月用量、审计与告警必须共用这一个实现。"""
return (now or utc_now()).astimezone(ZoneInfo("Asia/Shanghai")).strftime("%Y-%m")
def issue_secret(credential_id):
secret = "ck-%s-%s" % (credential_id, secrets.token_urlsafe(32))
return secret, hashlib.sha256(secret.encode("ascii")).hexdigest(), secret[:4] + "****" + secret[-4:]
......@@ -63,6 +83,15 @@ def validate_request_id(value):
return value
def validate_scope_key(value):
"""Scope keys are opaque source identifiers that must carry a source
namespace prefix (``source:subject``); the prefix is what keeps two
source systems from colliding on the same key."""
if not isinstance(value, str) or not SCOPE_KEY_PATTERN.fullmatch(value):
raise AccountError("INVALID_ARGUMENT", reason="INVALID_SCOPE_KEY")
return value
def authenticate(session, secret, *, now=None):
match = KEY_PATTERN.fullmatch(secret) if isinstance(secret, str) else None
if match is None:
......@@ -89,7 +118,8 @@ def authenticate(session, secret, *, now=None):
access = session.get(AccountClient, (credential.account_id, credential.client_id))
if account is None or access is None or access.status != "AUTHORIZED":
raise AccountError("ACCOUNT_ACCESS_REVOKED")
return Principal(credential.id, credential.client_id, credential.category, credential.account_id)
return Principal(credential.id, credential.client_id, credential.category, credential.account_id,
credential.sub_account_id, credential.scope_id)
def require_management(principal):
......
......@@ -46,6 +46,29 @@ def _clear_lease(call):
call.version += 1
def _employee_scope(call):
"""Only an employee call keeps monthly usage. The registration snapshot
separates the subjects: a store call snapshots scope_id == store_scope_id,
an employee call snapshots its own scope plus the store it belonged to."""
if call.scope_id is None or call.store_scope_id is None or call.scope_id == call.store_scope_id:
return None
return call.scope_id
def _warn_overage(session, call, scope_id, units):
"""S31:软限额只保证「准入时未超」,并发的在途结算仍可能把当月用量推过限额。
只在**跨越那一刻**记一条结构化告警(准入端已拦后续调用),不逐笔刷屏。"""
usage = repository.find_month_usage(session, scope_id, call.quota_month)
row = repository.find_scope_quota(session, scope_id, call.store_scope_id)
limit = row.monthly_quota_point_units if row is not None else None
if limit is None or usage <= limit or usage - units >= limit:
return
logger.warning("scope month quota exceeded, accountId=%s scopeId=%s storeScopeId=%s quotaMonth=%s "
"limitPointUnits=%s usedPointUnits=%s overPointUnits=%s callId=%s",
call.account_id, scope_id, call.store_scope_id, call.quota_month,
limit, usage, usage - limit, call.id)
def apply_settlement(session, account, call, evidence, consumption_id, point_id,
*, source="GATEWAY_LOG", now):
binding = _binding(call)
......@@ -71,7 +94,18 @@ def apply_settlement(session, account, call, evidence, consumption_id, point_id,
raise AccountError("ACCOUNT_BLOCKED", reason="INVALID_SETTLEMENT_EVIDENCE")
try:
units = require_bigint(quota * POINT_UNITS_PER_QUOTA)
before = require_bigint(account.balance_point_units)
except (ValueError, OverflowError):
raise AccountError("ACCOUNT_BLOCKED", reason="SETTLEMENT_AMOUNT_INVALID") from None
# The registration snapshot decides the charging bucket; later reassignment
# of scopes or sub-accounts never moves an in-flight call.
bucket = None
if call.sub_account_id is not None:
bucket = repository.lock_sub_account(session, call.sub_account_id)
if bucket is None or bucket.account_id != call.account_id:
raise AccountError("ACCOUNT_BLOCKED", reason="INVALID_SETTLEMENT_EVIDENCE")
try:
holder = bucket if bucket is not None else account
before = require_bigint(holder.balance_point_units)
after = require_bigint(before - units)
except (ValueError, OverflowError):
raise AccountError("ACCOUNT_BLOCKED", reason="SETTLEMENT_AMOUNT_INVALID") from None
......@@ -84,6 +118,7 @@ def apply_settlement(session, account, call, evidence, consumption_id, point_id,
output_tokens=call.output_tokens, total_tokens=call.total_tokens,
consumed_quota=quota, point_units_per_quota=POINT_UNITS_PER_QUOTA, consumed_point_units=units,
settlement_status="SUCCESS", settlement_source=source, settled_time=now,
sub_account_id=call.sub_account_id, scope_id=call.scope_id, store_scope_id=call.store_scope_id,
)
session.add(consumption)
session.flush()
......@@ -93,9 +128,18 @@ def apply_settlement(session, account, call, evidence, consumption_id, point_id,
point_units=-units, balance_before_units=before, balance_after_units=after,
point_units_per_quota=POINT_UNITS_PER_QUOTA, call_id=call.id,
consumption_record_id=consumption_id, request_id=call.request_id,
sub_account_id=call.sub_account_id,
))
account.balance_point_units = after
account.version += 1
holder.balance_point_units = after
holder.version += 1
month_scope = _employee_scope(call)
if month_scope is not None and call.quota_month:
# Monthly usage is keyed by (employee, month): a new month starts from
# zero with no reset job, and the employee carries what it consumed into
# whichever store it moved to. An in-flight call always lands in the
# month it was registered in, never in the month it settles in.
repository.upsert_month_usage(session, month_scope, call.quota_month, units, now=now)
_warn_overage(session, call, month_scope, units)
call.gateway_request_id = request_id
call.consumed_quota = quota
call.point_units_per_quota = POINT_UNITS_PER_QUOTA
......@@ -215,10 +259,20 @@ class SettlementService:
call.error_code, call.error_message = code, None
if permanent or call.retry_count >= 24:
call.billing_status, call.complete_time = "SETTLE_FAILED", now
logger.warning("settlement stopped, accountId=%s callId=%s code=%s retryCount=%s "
"status=SETTLE_FAILED scopeId=%s storeScopeId=%s quotaMonth=%s",
call.account_id, call.id, code, call.retry_count, call.scope_id,
call.store_scope_id, call.quota_month)
else:
call.billing_status = "SETTLE_PENDING"
minutes = _RETRY_MINUTES[min(max(call.retry_count - 1, 0), len(_RETRY_MINUTES) - 1)]
call.next_retry_time = now + timedelta(minutes=minutes)
# S30:这段窗口里调用已放行、用量未落账,软限额会短暂放宽;把重试节奏
# 记成结构化字段,运维才能对「结算延迟窗口」设阈值。
logger.warning("settlement retry scheduled, accountId=%s callId=%s code=%s retryCount=%s "
"nextRetryMinutes=%s scopeId=%s storeScopeId=%s quotaMonth=%s",
call.account_id, call.id, code, call.retry_count, minutes, call.scope_id,
call.store_scope_id, call.quota_month)
return False
def _retain_negative(self, state, evidence):
......
......@@ -11,7 +11,7 @@ from .account_model_gateway import ModelGateway
from .account_models import IdSegment
from .account_service import AccountService
from .account_settlement import SettlementService
from .db import create_mysql_engine
from .db import create_generator_engine, create_mysql_engine
from .id_generator import SegmentIDGenerator
......@@ -22,12 +22,17 @@ async def run(settings, batch_size):
settings.validate()
engine = create_mysql_engine(settings.database_url, pool_size=settings.db_pool_size,
max_overflow=settings.db_max_overflow, bounded_operations=True)
# 发号器独立小池:结算等写事务不得为号段分配去抢业务连接池。
generator_engine = create_generator_engine(
settings.database_url, pool_size=settings.db_generator_pool_size,
max_overflow=settings.db_generator_max_overflow)
gateway = None
try:
await asyncio.to_thread(_database_ready, engine)
factory = sessionmaker(engine, expire_on_commit=False)
generator_factory = sessionmaker(generator_engine, expire_on_commit=False)
gateway = ModelGateway(settings)
accounts = AccountService(factory, SegmentIDGenerator(factory, segment_model=IdSegment), gateway, settings)
accounts = AccountService(factory, SegmentIDGenerator(generator_factory, segment_model=IdSegment), gateway, settings)
settlements = SettlementService(accounts)
settled = await settlements.run_once(batch_size)
cleared = await run_database(accounts, settlements.clear_results, batch_size)
......@@ -37,6 +42,7 @@ async def run(settings, batch_size):
if gateway is not None:
await gateway.aclose()
engine.dispose()
generator_engine.dispose()
def main():
......
......@@ -55,6 +55,17 @@ def create_mysql_engine(database_url: str, *, pool_size=5, max_overflow=0, bound
return engine
def create_generator_engine(database_url: str, *, pool_size=2, max_overflow=2):
"""发号器专用连接池:号段分配不得占用业务连接池。
写请求在业务事务里申请新号段时,若与业务共用连接池,池内无空闲连接就会等到
`pool_timeout` 再抛 503(`db_pool_size=1`、或池被在途事务占满时必现)。号段一次
取 1000 个、单次分配只是"锁一行 + 写一行 + 提交",独立小池足够,也不会挤压业务池。
"""
return create_mysql_engine(database_url, pool_size=pool_size, max_overflow=max_overflow,
bounded_operations=True)
def init_engine(settings: Settings) -> None:
"""按环境配置初始化全局 Engine。重复调用忽略(uvicorn 多次加载保护)。"""
global _engine, _session_factory
......
# 独立算力账户服务开发与迁移计划
> SQL 整理更新(2026-09-30):系统尚未上线,当前初始化使用 `scripts/sql/schema.sql`,检查使用 `scripts/sql/check.sql`。SaaS 材料及历史脚本已分别移入 `scripts/sql/saas/`、`scripts/sql/archive/`,执行流程以 `scripts/sql/README.md` 为准。本文保留原阶段计划用于追溯。
- 日期:2026-09-23
- 状态:待实施的开发计划,不代表功能已完成或已授权上线
- 范围:`mei1-computing-service`、`mei1-saas`、`business-saas/mwcloud`、`opt`
......
#!/bin/sh
# 隔离 MySQL 实例:给 MySQL 专项用例提供真机语义(行锁、原子 UPSERT、CHECK)。
#
# 用法:
# scripts/p1_mysql_instance.sh start # 初始化(首次)并后台启动
# scripts/p1_mysql_instance.sh stop # 关闭
# scripts/p1_mysql_instance.sh status # 探活 + 打印回归命令
#
# 约束:socket 所在目录名必须保持 `computing-p1-mysql.*` 前缀,
# tests/test_account_foundation.py 用它把临时实例与业务库严格区分开;
# 实例只监听 unix socket(--skip-networking),不使用 3306。
set -eu
DIR="${P1_MYSQL_DIR:-/tmp/computing-p1-mysql.dev}"
MYSQLD="${P1_MYSQLD:-/usr/local/mysql/bin/mysqld}"
MYSQLADMIN="${P1_MYSQLADMIN:-/usr/local/mysql/bin/mysqladmin}"
SOCKET="$DIR/mysqld.sock"
BASEDIR="$(cd "$(dirname "$MYSQLD")/.." && pwd)"
start() {
if [ ! -d "$DIR/data/mysql" ]; then
echo "初始化数据目录 $DIR/data"
mkdir -p "$DIR/data"
"$MYSQLD" --no-defaults --initialize-insecure --datadir="$DIR/data" \
--basedir="$BASEDIR" --log-error="$DIR/init.log"
fi
if [ -S "$SOCKET" ] && "$MYSQLADMIN" --no-defaults --socket="$SOCKET" -uroot ping >/dev/null 2>&1; then
echo "已在运行:$SOCKET"
else
# 与 P0 要求一致:READ COMMITTED、UTC、严格 SQL 模式(MySQL >= 8.0.16)。
nohup "$MYSQLD" --no-defaults --datadir="$DIR/data" --basedir="$BASEDIR" \
--socket="$SOCKET" --skip-networking --skip-mysqlx \
--pid-file="$DIR/mysqld.pid" --log-error="$DIR/error.log" \
--transaction-isolation=READ-COMMITTED --default-time-zone=+00:00 \
--sql-mode=STRICT_TRANS_TABLES,NO_ENGINE_SUBSTITUTION \
>/dev/null 2>&1 &
i=0
while [ "$i" -lt 60 ]; do
if "$MYSQLADMIN" --no-defaults --socket="$SOCKET" -uroot ping >/dev/null 2>&1; then break; fi
i=$((i + 1)); sleep 1
done
fi
"$MYSQLADMIN" --no-defaults --socket="$SOCKET" -uroot ping >/dev/null 2>&1 \
|| { echo "启动失败,见 $DIR/error.log"; exit 1; }
echo "就绪:$SOCKET"
echo "回归:python3 -m pytest --account-mysql-socket=$SOCKET -q --junitxml=/tmp/junit-mysql.xml"
}
stop() {
if [ -S "$SOCKET" ]; then
"$MYSQLADMIN" --no-defaults --socket="$SOCKET" -uroot shutdown || true
echo "已关闭 $SOCKET"
else
echo "未运行"
fi
}
status() {
if [ -S "$SOCKET" ] && "$MYSQLADMIN" --no-defaults --socket="$SOCKET" -uroot ping >/dev/null 2>&1; then
echo "运行中:$SOCKET"
echo "SQLite+MySQL 全量:python3 -m pytest --account-mysql-socket=$SOCKET -q --junitxml=/tmp/junit-mysql.xml"
else
echo "未运行(start 可启动)"
fi
}
case "${1:-status}" in
start) start ;;
stop) stop ;;
status) status ;;
*) echo "用法:$0 start|stop|status" >&2; exit 2 ;;
esac
......@@ -9,7 +9,9 @@ sys.path.insert(0, str(Path(__file__).resolve().parents[1]))
from app.account_models import Base
HEADER = """-- P1 independent account schema; generated by scripts/render_account_ddl.py.
HEADER = """-- Current independent account schema; generated by scripts/render_account_ddl.py.
-- Canonical output: scripts/sql/schema.sql (development baseline, not an upgrade).
-- Includes account hierarchy, scopes, monthly quotas and ledger indexes.
-- Target: a NEW, explicitly selected computing schema, MySQL >= 8.0.16.
-- Review and execute manually. Do not run on the existing SaaS schema.
-- No CREATE DATABASE/USE, no account grants, no source data changes, no ID seed.
......
# SQL 使用说明
系统处于开发阶段、尚未上线。**新建独立算力库只使用当前全量 `schema.sql`,不再按日期叠加历史增量。** 日常维护只需要关注本目录的两个 SQL 文件。
## 目录与目标库
| 路径 | 目标库 | 用途 |
| --- | --- | --- |
| [schema.sql](schema.sql) | 新建、空的独立算力库 | 当前 15 张表及全部索引、约束;唯一建表基线 |
| [check.sql](check.sql) | 当前独立算力库 | 只读检查账本、余额桶、调用关联、主体、员工月用量及号段 |
| [saas/source_precheck.sql](saas/source_precheck.sql) | 旧 SaaS 库 | 接入/迁移前的只读预检;不是网关初始化步骤 |
| [saas/binding_and_meiji_request_id.sql](saas/binding_and_meiji_request_id.sql) | SaaS 库 | 旧接入方案交付材料;实施 SaaS 对接时重新核对,不随网关建库执行 |
| [archive/](archive/) | 历史对应的库 | 旧基线、增量及检查脚本,原文保留供追溯;不参与当前初始化或测试 |
禁止递归批量执行整个目录。`saas/` 与 `archive/` 的脚本不是 `schema.sql` 的后续步骤。
## 新环境初始化
1. 在项目外准备独立数据库及运行账号,显式选择目标数据库。最低 MySQL 8.0.16;服务会使用 READ COMMITTED、UTC 和严格模式。
2. 在空的独立数据库执行 `schema.sql` 一次。脚本不包含建库、`USE`、授权、数据导入或号段初始化;现有表会让执行失败,不用 `--force` 跳过错误。
3. 准备号段。仅对于真正空的新安装,在开始发号前可显式执行下列语句;不使用 `INSERT IGNORE` 或覆盖式更新来掩盖已有状态:
```sql
INSERT INTO t_computing_id_segment (generator_key, next_id)
VALUES ('computing-global', 100000000000000000);
```
若涉及历史数据导入,必须停止旧分配器并继承其已提交高水位,不能使用上述空库种子,也不能用 `MAX(id)+1` 代替。
4. 执行 `check.sql`。然后按项目 README 配置外置文件、初始化管理凭证并启动服务;建表不等于完成凭证和网关配置。
已有开发数据库不会被这些文件自动升级、删除或重建。可以对确认不需要保留数据的环境重新建一个空库;要保留现有数据时,应依据其实际结构制定单独升级操作。本轮整理没有操作已有开发数据或业务库。
## 检查结果口径
- **带 `anomaly` 列**的结果集:表示一致性异常,在一致快照上应为空。
- 其他结果集:数据库信息、表清单、未终结调用、号段信息和授权状态,返回行不自动等于数据错误。
- 对账时暂停并发写入或使用一致快照,避免不同语句看到不同阶段的数据。空异常结果不代表源库与目标库迁移金额完全一致。
- 账户余额只核对 `sub_account_id IS NULL` 的账户桶流水;子账号各自核对自己的桶。
- 员工月用量只统计 EMPLOYEE 主体,按调用登记月归类;门店消费不产生员工月用量。
- 消费与流水按 `account_id + sub_account_id` 分组,防止不同账户的未分配桶相互抵消。
- 转账按来源、请求号、账户核对配对;相同请求号在不同来源下相互独立。
- 号段检查覆盖当前所有具有数值 `id` 的独立表,包括子账号、scope 和额度表;月用量表使用员工与月份联合主键,不参与发号。
## 开发期间如何变更
1. 修改 `app/account_models.py` 的结构定义。
2. 在项目根目录生成当前全量基线:
```bash
python scripts/render_account_ddl.py --output scripts/sql/schema.sql
```
3. 涉及账务或主体关系变化时,同步更新 `check.sql` 及相关契约。
4. MySQL 测试直接执行当前 `schema.sql`;基础检查测试与分账 E2E 执行交付的 `check.sql`,而非仅用 ORM 模拟其查询。
5. 未上线期间不再为每次结构修改增加一套日期基线和增量。正式上线后,冻结首个发布基线,再为需要保留数据的已发布数据库维护版本化增量。
本轮合并了早期检查和层级检查,校正余额桶、员工月用量、跨账户聚合及转账幂等域的口径;不改变运行服务的计费规则。
## 整理验证(2026-09-30)
- 确认当前全量 DDL 的结构语句与整理前的最新全量文件一致,仅更新说明与路径。
- 六份归档 SQL 逐文件核对,原文未变;当前测试和执行说明不再引用旧入口。
- SQLite + 专用临时 MySQL 全量回归:**2160 passed、27 skipped**(Python 3.9.6,200.36 秒)。目标 Python 3.12 尚未在本轮验证。
- 直接执行交付检查 SQL,覆盖正常分账、异常关联、跨账户桶差额及不同来源同请求号转账;未操作已有开发数据或业务库。
## 旧文件去向
| 原 `scripts/` 文件名 | 当前位置 |
| --- | --- |
| `20260929_independent_account_ddl.sql` | `schema.sql`,以后由 ORM 生成 |
| `20260923_account_source_precheck.sql` | `saas/source_precheck.sql` |
| `20260923_saas_binding_and_meiji_request_id.sql` | `saas/binding_and_meiji_request_id.sql` |
| `20260923_account_target_check.sql` | 原文放入 `archive/`,有效检查已合并进 `check.sql` |
| `20260929_account_hierarchy_target_check.sql` | 原文放入 `archive/`,有效检查已合并进 `check.sql` |
| `20260923_independent_account_ddl.sql` | `archive/` |
| `20260929_account_hierarchy_and_scopes_ddl.sql` | `archive/` |
| `20260922_computing_stage0_ddl.sql` | `archive/` |
| `reconciliation_check.sql` | `archive/` |
-- P6 hierarchy incremental DDL: account sub-accounts and store/employee scopes.
-- Generated alongside scripts/20260929_independent_account_ddl.sql (full schema,
-- models in app/account_models.py). Review and execute manually on the target
-- schema that already ran scripts/20260923_independent_account_ddl.sql.
-- No CREATE DATABASE/USE, no account grants, no data changes, no ID seed.
-- New tables are additive: no backfill required; existing call/consumption rows
-- keep NULL hierarchy snapshots (account bucket, no scope) which is the correct
-- pre-hierarchy semantics. t_computing_id_segment is untouched; the new tables
-- share the existing computing-global high water mark.
-- No IF NOT EXISTS: existing objects must stop execution; never continue with
-- --force. Retrying after partial failure requires inspecting the target, not
-- overwriting it. Order matters: new tables first, then ALTERs that reference them.
-- Foreign keys added to t_computing_credential and t_computing_point_record carry
-- explicit names (fk_..._sub_account / fk_..._scope); a fresh install from the
-- full DDL instead gets engine-generated names for the same constraints, which
-- is cosmetic only.
-- Rollback: this script is forward-only. Before any hierarchy rows exist, the
-- added columns could be dropped and the new tables discarded; once IDs or
-- transfers exist, reconcile first and never rewind IDs.
-- Design: docs/account-hierarchy-and-scope-limits-design.md (v0.2); contract:
-- docs/p0-contract-freeze.md §2.10.
CREATE TABLE t_computing_sub_account (
id BIGINT UNSIGNED NOT NULL,
account_id BIGINT UNSIGNED NOT NULL,
name VARCHAR(100) NOT NULL,
remark VARCHAR(500),
status ENUM('ACTIVE','DISABLED') NOT NULL DEFAULT 'ACTIVE',
balance_point_units BIGINT NOT NULL DEFAULT 0,
version INTEGER NOT NULL DEFAULT 0,
create_time DATETIME(3) NOT NULL DEFAULT (UTC_TIMESTAMP(3)),
last_update_time DATETIME(3) NOT NULL DEFAULT (UTC_TIMESTAMP(3)),
PRIMARY KEY (id),
CONSTRAINT uq_computing_sub_account_name UNIQUE (account_id, name),
FOREIGN KEY(account_id) REFERENCES t_computing_account (id)
)ENGINE=InnoDB CHARSET=utf8mb4 COLLATE utf8mb4_bin;
CREATE INDEX idx_computing_sub_account_created ON t_computing_sub_account (account_id, create_time, id);
CREATE TABLE t_computing_scope (
id BIGINT UNSIGNED NOT NULL,
account_id BIGINT UNSIGNED NOT NULL,
sub_account_id BIGINT UNSIGNED,
scope_type ENUM('STORE','EMPLOYEE') NOT NULL,
scope_key VARCHAR(100) CHARACTER SET ascii COLLATE ascii_bin NOT NULL,
parent_scope_id BIGINT UNSIGNED,
status ENUM('ACTIVE','DISABLED') NOT NULL DEFAULT 'ACTIVE',
version INTEGER NOT NULL DEFAULT 0,
create_time DATETIME(3) NOT NULL DEFAULT (UTC_TIMESTAMP(3)),
last_update_time DATETIME(3) NOT NULL DEFAULT (UTC_TIMESTAMP(3)),
PRIMARY KEY (id),
CONSTRAINT uq_computing_scope_identity UNIQUE (scope_type, scope_key),
CONSTRAINT ck_computing_scope_shape CHECK ((scope_type = 'STORE' AND parent_scope_id IS NULL) OR (scope_type = 'EMPLOYEE' AND parent_scope_id IS NOT NULL AND sub_account_id IS NULL)),
CONSTRAINT fk_computing_scope_account FOREIGN KEY(account_id) REFERENCES t_computing_account (id),
CONSTRAINT fk_computing_scope_sub_account FOREIGN KEY(sub_account_id) REFERENCES t_computing_sub_account (id),
CONSTRAINT fk_computing_scope_parent FOREIGN KEY(parent_scope_id) REFERENCES t_computing_scope (id)
)ENGINE=InnoDB CHARSET=utf8mb4 COLLATE utf8mb4_bin;
CREATE INDEX idx_computing_scope_parent ON t_computing_scope (parent_scope_id);
CREATE INDEX idx_computing_scope_sub_account ON t_computing_scope (sub_account_id, scope_type, status);
CREATE TABLE t_computing_scope_quota (
id BIGINT UNSIGNED NOT NULL,
scope_id BIGINT UNSIGNED NOT NULL,
store_scope_id BIGINT UNSIGNED NOT NULL,
monthly_quota_point_units BIGINT,
version INTEGER NOT NULL DEFAULT 0,
create_time DATETIME(3) NOT NULL DEFAULT (UTC_TIMESTAMP(3)),
last_update_time DATETIME(3) NOT NULL DEFAULT (UTC_TIMESTAMP(3)),
PRIMARY KEY (id),
CONSTRAINT uq_computing_scope_quota_pair UNIQUE (scope_id, store_scope_id),
CONSTRAINT ck_computing_scope_quota_nonnegative CHECK (monthly_quota_point_units IS NULL OR monthly_quota_point_units >= 0),
CONSTRAINT fk_computing_scope_quota_scope FOREIGN KEY(scope_id) REFERENCES t_computing_scope (id),
CONSTRAINT fk_computing_scope_quota_store FOREIGN KEY(store_scope_id) REFERENCES t_computing_scope (id)
)ENGINE=InnoDB CHARSET=utf8mb4 COLLATE utf8mb4_bin;
CREATE TABLE t_computing_scope_month_usage (
scope_id BIGINT UNSIGNED NOT NULL,
month CHAR(7) CHARACTER SET ascii COLLATE ascii_bin NOT NULL,
used_point_units BIGINT NOT NULL DEFAULT 0,
create_time DATETIME(3) NOT NULL DEFAULT (UTC_TIMESTAMP(3)),
last_update_time DATETIME(3) NOT NULL DEFAULT (UTC_TIMESTAMP(3)),
PRIMARY KEY (scope_id, month),
CONSTRAINT fk_computing_scope_month_usage_scope FOREIGN KEY(scope_id) REFERENCES t_computing_scope (id),
CONSTRAINT ck_computing_scope_month_usage_nonnegative CHECK (used_point_units >= 0)
)ENGINE=InnoDB CHARSET=utf8mb4 COLLATE utf8mb4_bin;
ALTER TABLE t_computing_management_request
MODIFY COLUMN operation_type ENUM('PROVISION_ACCOUNT','ISSUE_CALL_KEY','ACTIVATE_CREDENTIAL','REVOKE_CREDENTIAL','GRANT_CLIENT','REVOKE_CLIENT','SET_ACCOUNT_STATUS','RECONCILE','CREATE_SUB_ACCOUNT','SET_SUB_ACCOUNT_STATUS','TRANSFER','BIND_SCOPE','SET_SCOPE_STATUS','MOVE_STORE','MOVE_EMPLOYEE','SET_SCOPE_QUOTA') NOT NULL;
ALTER TABLE t_computing_credential
ADD COLUMN sub_account_id BIGINT UNSIGNED AFTER account_id,
ADD COLUMN scope_id BIGINT UNSIGNED AFTER sub_account_id,
DROP CHECK ck_computing_credential_scope,
ADD CONSTRAINT ck_computing_credential_scope CHECK ((category = 'CALL' AND account_id IS NOT NULL AND issued_by IS NOT NULL AND NOT (sub_account_id IS NOT NULL AND scope_id IS NOT NULL)) OR (category IN ('INTEGRATION', 'PLATFORM') AND account_id IS NULL AND sub_account_id IS NULL AND scope_id IS NULL)),
ADD CONSTRAINT fk_computing_credential_sub_account FOREIGN KEY(sub_account_id) REFERENCES t_computing_sub_account (id),
ADD CONSTRAINT fk_computing_credential_scope_ref FOREIGN KEY(scope_id) REFERENCES t_computing_scope (id);
ALTER TABLE t_computing_call
ADD COLUMN sub_account_id BIGINT UNSIGNED,
ADD COLUMN scope_id BIGINT UNSIGNED,
ADD COLUMN store_scope_id BIGINT UNSIGNED,
ADD COLUMN quota_month CHAR(7) CHARACTER SET ascii COLLATE ascii_bin;
ALTER TABLE t_computing_consumption
ADD COLUMN sub_account_id BIGINT UNSIGNED,
ADD COLUMN scope_id BIGINT UNSIGNED,
ADD COLUMN store_scope_id BIGINT UNSIGNED;
ALTER TABLE t_computing_point_record
ADD COLUMN sub_account_id BIGINT UNSIGNED AFTER consumption_record_id,
MODIFY COLUMN type ENUM('GIFT','CONSUME','TRANSFER_OUT','TRANSFER_IN') NOT NULL,
DROP CHECK ck_computing_point_record_source,
ADD CONSTRAINT ck_computing_point_record_source CHECK ((type = 'GIFT' AND point_units > 0 AND call_id IS NULL AND consumption_record_id IS NULL AND sub_account_id IS NULL AND (legacy_record = 1 OR gift_id IS NOT NULL)) OR (type = 'CONSUME' AND point_units < 0 AND gift_id IS NULL AND consumption_record_id IS NOT NULL AND (legacy_record = 1 OR call_id IS NOT NULL)) OR (type = 'TRANSFER_OUT' AND point_units < 0 AND gift_id IS NULL AND call_id IS NULL AND consumption_record_id IS NULL) OR (type = 'TRANSFER_IN' AND point_units > 0 AND gift_id IS NULL AND call_id IS NULL AND consumption_record_id IS NULL)),
ADD CONSTRAINT fk_computing_point_record_sub_account FOREIGN KEY(sub_account_id) REFERENCES t_computing_sub_account (id);
CREATE INDEX idx_computing_point_record_bucket ON t_computing_point_record (account_id, sub_account_id, type);
-- 段 3e(2026-09-30):分账报表索引。实测(隔离 MySQL 8.0.43;单账户 10 万行、门店选择性 1%、员工 0.05%、
-- 桶 80%,见 docs/account-hierarchy-and-scope-limits-design.md §8):
-- 门店账 LIMIT 20:无维度索引时顺 account+create_time 逆序读 3,983 行 / 2.80ms;加索引后 20 行 / 0.02ms
-- 员工账 LIMIT 20:无索引读 40,000 行 / 18.8ms;加索引后 20 行 / 0.01ms
-- 子账号桶(占行数 80%,现有索引顺扫即命中):0.73ms → 0.30ms,收益不足抵 44B/行写放大 → 不加
-- 月账单页:瓶颈在 ORDER BY used 排序而非 scope 扫描(6,300 scope 表全扫 1.65ms),scope 表比
-- consumption 小两个数量级 → 不加;scope 数到 10^4 量级或月账单变慢时再评估
-- 代价:每索引 ≈44B/行(30 万行 ≈12.6MB);call/consumption 各 2 条,插入与结算更新各多写 4 项。
CREATE INDEX idx_computing_call_store_scope ON t_computing_call (account_id, store_scope_id, create_time, id);
CREATE INDEX idx_computing_call_scope_subject ON t_computing_call (account_id, scope_id, create_time, id);
CREATE INDEX idx_computing_consumption_store_scope ON t_computing_consumption (account_id, store_scope_id, create_time, id);
CREATE INDEX idx_computing_consumption_scope_subject ON t_computing_consumption (account_id, scope_id, create_time, id);
-- P6 hierarchy target checks: read-only, run manually on the independent schema
-- AFTER scripts/20260929_account_hierarchy_and_scopes_ddl.sql has been applied.
-- Does not replace 20260923_account_target_check.sql (frozen, already executed);
-- run both. Empty anomaly results do not prove source-to-target equality.
-- Sections returning anomaly rows must be reviewed before relying on the data.
SELECT DATABASE() AS target_database, VERSION() AS mysql_version,
@@session.transaction_isolation AS isolation_level, @@session.time_zone AS time_zone;
SELECT table_name, table_rows
FROM information_schema.tables
WHERE table_schema = DATABASE() AND table_name IN (
't_computing_sub_account', 't_computing_scope',
't_computing_scope_quota', 't_computing_scope_month_usage'
)
ORDER BY table_name;
-- 1. Bucket identity: account balance equals the sum of its account-bucket
-- records (sub_account_id IS NULL), exactly one row per account when any
-- record exists; residual non-zero is an anomaly.
SELECT a.id AS account_id, a.balance_point_units,
COALESCE(p.net_units, 0) AS account_bucket_net_units,
a.balance_point_units - COALESCE(p.net_units, 0) AS residual_units
FROM t_computing_account a
LEFT JOIN (
SELECT account_id, SUM(point_units) AS net_units
FROM t_computing_point_record
WHERE sub_account_id IS NULL
GROUP BY account_id
) p ON p.account_id = a.id
WHERE a.balance_point_units <> COALESCE(p.net_units, 0);
-- 2. Bucket identity: each sub-account balance equals its bucket sum.
SELECT s.id AS sub_account_id, s.account_id, s.balance_point_units,
COALESCE(p.net_units, 0) AS sub_bucket_net_units,
s.balance_point_units - COALESCE(p.net_units, 0) AS residual_units
FROM t_computing_sub_account s
LEFT JOIN (
SELECT sub_account_id, SUM(point_units) AS net_units
FROM t_computing_point_record
WHERE sub_account_id IS NOT NULL
GROUP BY sub_account_id
) p ON p.sub_account_id = s.id
WHERE s.balance_point_units <> COALESCE(p.net_units, 0);
-- 3. Transfer conservation: every transfer requestId has exactly one OUT and
-- one IN, equal magnitude (net zero), one row on the account bucket and one
-- on a sub-account bucket of the same account.
SELECT 'transfer_pairing' AS anomaly, request_id, account_id,
SUM(type = 'TRANSFER_OUT') AS out_rows, SUM(type = 'TRANSFER_IN') AS in_rows,
SUM(point_units) AS net_units,
SUM(sub_account_id IS NULL) AS account_bucket_rows
FROM t_computing_point_record
WHERE type IN ('TRANSFER_OUT', 'TRANSFER_IN')
GROUP BY request_id, account_id
HAVING out_rows <> 1 OR in_rows <> 1 OR net_units <> 0 OR account_bucket_rows <> 1;
SELECT 'transfer_cross_account' AS anomaly, o.request_id, o.account_id AS out_account, i.account_id AS in_account
FROM t_computing_point_record o
JOIN t_computing_point_record i
ON i.request_id = o.request_id AND i.type = 'TRANSFER_IN' AND o.type = 'TRANSFER_OUT'
WHERE o.account_id <> i.account_id;
-- 4. GIFT records must stay on the account bucket.
SELECT 'gift_not_account_bucket' AS anomaly, p.id, p.gift_id
FROM t_computing_point_record p
WHERE p.type = 'GIFT' AND p.sub_account_id IS NOT NULL;
-- 5. Scope shape: employees must point at a STORE of the same account; stores
-- must not carry a parent; employee rows must not carry sub_account_id.
-- ACTIVE is deliberately NOT checked (a disabled store with live employees
-- is a legal steady state).
SELECT 'scope_shape' AS anomaly, s.id, s.scope_type, s.parent_scope_id, s.account_id
FROM t_computing_scope s
LEFT JOIN t_computing_scope p ON p.id = s.parent_scope_id
WHERE (s.scope_type = 'EMPLOYEE' AND (s.parent_scope_id IS NULL OR p.id IS NULL
OR p.scope_type <> 'STORE' OR p.account_id <> s.account_id
OR s.sub_account_id IS NOT NULL))
OR (s.scope_type = 'STORE' AND s.parent_scope_id IS NOT NULL);
-- 6. Store sub_account binding must belong to the same account.
SELECT 'scope_sub_account_mismatch' AS anomaly, s.id, s.scope_key, s.account_id, sa.account_id AS bound_account
FROM t_computing_scope s
JOIN t_computing_sub_account sa ON sa.id = s.sub_account_id
WHERE s.scope_type = 'STORE' AND sa.account_id <> s.account_id;
-- 7. scope_quota rows: employee x store of the same account.
SELECT 'scope_quota_shape' AS anomaly, q.id, q.scope_id, q.store_scope_id
FROM t_computing_scope_quota q
JOIN t_computing_scope e ON e.id = q.scope_id
JOIN t_computing_scope st ON st.id = q.store_scope_id
WHERE e.scope_type <> 'EMPLOYEE' OR st.scope_type <> 'STORE'
OR e.account_id <> st.account_id;
-- 8. Credential hierarchy: only CALL keys carry a subject; the subject must
-- belong to the credential's own account; subject columns are exclusive.
SELECT 'credential_subject_mismatch' AS anomaly, c.id, c.account_id, c.sub_account_id, c.scope_id
FROM t_computing_credential c
LEFT JOIN t_computing_sub_account sa ON sa.id = c.sub_account_id
LEFT JOIN t_computing_scope sc ON sc.id = c.scope_id
WHERE c.category <> 'CALL' AND (c.sub_account_id IS NOT NULL OR c.scope_id IS NOT NULL)
OR (c.sub_account_id IS NOT NULL AND c.scope_id IS NOT NULL)
OR (c.sub_account_id IS NOT NULL AND sa.account_id <> c.account_id)
OR (c.scope_id IS NOT NULL AND sc.account_id <> c.account_id);
-- 9. Call snapshots equal consumption snapshots for settled calls.
SELECT 'call_consumption_snapshot_mismatch' AS anomaly, c.id, c.account_id
FROM t_computing_call c
JOIN t_computing_consumption r ON r.id = c.consumption_record_id
WHERE c.billing_status = 'SUCCESS' AND NOT (
r.sub_account_id <=> c.sub_account_id
AND r.scope_id <=> c.scope_id
AND r.store_scope_id <=> c.store_scope_id
);
-- 10. Consume point records land on the same bucket as their consumption.
SELECT 'consume_bucket_mismatch' AS anomaly, p.id, p.consumption_record_id
FROM t_computing_point_record p
JOIN t_computing_consumption r ON r.id = p.consumption_record_id
WHERE p.type = 'CONSUME' AND NOT (p.sub_account_id <=> r.sub_account_id);
-- 11. Monthly usage: scope_month_usage equals the settled consumption grouped
-- by the registration-month snapshot carried on the call.
SELECT 'month_usage_mismatch' AS anomaly, u.scope_id, u.month,
u.used_point_units, COALESCE(agg.units, 0) AS aggregated_units
FROM t_computing_scope_month_usage u
LEFT JOIN (
SELECT c.scope_id, c.quota_month, SUM(r.consumed_point_units) AS units
FROM t_computing_consumption r
JOIN t_computing_call c ON c.id = r.call_id
WHERE r.settlement_status = 'SUCCESS' AND c.scope_id IS NOT NULL
GROUP BY c.scope_id, c.quota_month
) agg ON agg.scope_id = u.scope_id AND agg.quota_month = u.month
WHERE u.used_point_units <> COALESCE(agg.units, 0);
SELECT 'month_usage_orphan' AS anomaly, agg.scope_id, agg.quota_month, agg.units
FROM (
SELECT c.scope_id, c.quota_month, SUM(r.consumed_point_units) AS units
FROM t_computing_consumption r
JOIN t_computing_call c ON c.id = r.call_id
WHERE r.settlement_status = 'SUCCESS' AND c.scope_id IS NOT NULL
GROUP BY c.scope_id, c.quota_month
) agg
LEFT JOIN t_computing_scope_month_usage u
ON u.scope_id = agg.scope_id AND u.month = agg.quota_month
WHERE u.scope_id IS NULL;
-- 12. Every registration carries its business month (段 3 修正:quota_month 是**登记月**,
-- 所有调用都有;它只在员工主体上被读来累加月用量,账户级/子账号级调用留着不参与统计)。
-- 上一版规则把「无主体却带 quota_month」当异常,会对每一笔账户级调用误报,已作废。
SELECT 'quota_month_shape' AS anomaly, c.id, c.account_id, c.scope_id, c.store_scope_id, c.quota_month
FROM t_computing_call c
WHERE c.quota_month IS NULL
OR c.quota_month NOT REGEXP '^[0-9]{4}-(0[1-9]|1[0-2])$';
-- 13. 分账维度形状(段 3a):主体快照成对出现——门店/员工调用 scope_id 与 store_scope_id
-- 都有,账户级/子账号级调用两者都无;否则门店聚合会漏行(§8 恒等式的前提)。
SELECT 'ledger_snapshot_pairing' AS anomaly, c.id, c.account_id, c.sub_account_id, c.scope_id, c.store_scope_id
FROM t_computing_call c
WHERE (c.scope_id IS NULL) <> (c.store_scope_id IS NULL);
-- 14. 分账快照的桶必须属于同一账户(历史改派允许与 scope 当前绑定不同,但不允许跨账户)。
SELECT 'ledger_bucket_cross_account' AS anomaly, c.id, c.account_id, c.sub_account_id
FROM t_computing_call c
JOIN t_computing_sub_account sa ON sa.id = c.sub_account_id
WHERE sa.account_id <> c.account_id;
-- 15. 分账聚合 = 桶流水(段 3d):消费按 sub_account_id 聚合必须与 CONSUME 流水一致;
-- 账户桶(NULL)与子账号桶各自守恒。门店/员工维度不产生独立流水,由 #9 快照一致
-- 与 #11 月用量核对覆盖。
SELECT 'consumption_bucket_mismatch' AS anomaly, cr.bucket_id, cr.units AS consumption_units, pr.units AS point_record_units
FROM (SELECT COALESCE(sub_account_id, 0) AS bucket_id, SUM(consumed_point_units) AS units
FROM t_computing_consumption WHERE settlement_status = 'SUCCESS'
GROUP BY COALESCE(sub_account_id, 0)) cr
LEFT JOIN (SELECT COALESCE(sub_account_id, 0) AS bucket_id, SUM(-point_units) AS units
FROM t_computing_point_record WHERE type = 'CONSUME'
GROUP BY COALESCE(sub_account_id, 0)) pr ON pr.bucket_id = cr.bucket_id
WHERE COALESCE(pr.units, 0) <> cr.units
UNION ALL
SELECT 'consumption_bucket_orphan', pr.bucket_id, 0, pr.units
FROM (SELECT COALESCE(sub_account_id, 0) AS bucket_id, SUM(-point_units) AS units
FROM t_computing_point_record WHERE type = 'CONSUME'
GROUP BY COALESCE(sub_account_id, 0)) pr
LEFT JOIN (SELECT COALESCE(sub_account_id, 0) AS bucket_id, SUM(consumed_point_units) AS units
FROM t_computing_consumption WHERE settlement_status = 'SUCCESS'
GROUP BY COALESCE(sub_account_id, 0)) cr ON cr.bucket_id = pr.bucket_id
WHERE cr.bucket_id IS NULL;
# 历史 SQL,仅供追溯
这些文件原文保留,不再参与新环境初始化、当前数据库检查或自动化测试。文件头中的相对路径、执行顺序和阶段说明属于历史上下文,不应作为当前操作指引。
新库使用上一级的 [schema.sql](../schema.sql),当前独立库检查使用 [check.sql](../check.sql),完整说明见 [SQL 使用说明](../README.md)。
| 文件 | 历史用途 |
| --- | --- |
| `20260922_computing_stage0_ddl.sql` | 旧 SaaS 共享库增量 |
| `reconciliation_check.sql` | 旧 SaaS 共享库对账 |
| `20260923_independent_account_ddl.sql` | 11 张表的旧独立库基线 |
| `20260929_account_hierarchy_and_scopes_ddl.sql` | 旧独立库基线上的层级增量 |
| `20260923_account_target_check.sql` | 账户层级引入前的检查口径 |
| `20260929_account_hierarchy_target_check.sql` | 层级检查历史版本,当前合并版已校正部分口径 |
# SaaS 接入与迁移材料
这里的 SQL 只针对 **SaaS 数据库**,不在独立算力库执行,也不属于网关初始化步骤。
- [source_precheck.sql](source_precheck.sql):只读检查旧算力表、关联、未清记录和号段,实施迁移前使用。
- [binding_and_meiji_request_id.sql](binding_and_meiji_request_id.sql):保留原接入方案的交付内容;当前独立服务已经支持门店/员工 Key,实际对接时应重新核对凭证映射设计及 SaaS 目标结构,不视为已验收的最终迁移方案。
原接入脚本的 `computing_request_id` 增量与旧阶段 0 脚本存在重叠。执行前检查列及索引是否已存在,明确选择需要执行的语句;不要运行后忽略重复列错误或使用 `--force`。原阶段 0 文件现位于 `../archive/20260922_computing_stage0_ddl.sql`。
本次仅整理交付文件,未执行 SaaS 建表、改表或数据迁移。
-- P5 SaaS 侧交付(2026-09-23):手动执行,不由 Python 服务自动执行。
-- SaaS 接入方案材料(原 P5,2026-09-23),不是独立网关初始化脚本。
-- 当前门店/员工凭证映射需在正式对接时另行核对;参见本目录 README.md。
-- 1) 商户算力服务绑定表:SaaS 商户与 Python 独立算力账户(18 位数值 ID)一对一绑定,
-- 保存开户恢复请求号与 CALL Key AES-GCM 密文(base64,密钥来自部署 Secret,不入库明文)。
-- 2) 美际分析表补充 computing_request_id:稳定业务请求号 meiji-skin-analysis:{analysisId},
-- 用于按原请求号回查 /api/v1/calls/{request_id},不换号重发。
-- 注意:阶段 0 脚本 20260922_computing_stage0_ddl.sql 已含同一列增量;若该脚本已在目标库执行,
-- 本段 ALTER 会因列重复失败,跳过即可(仅执行第 1 段建表)。
-- 注意:../archive/20260922_computing_stage0_ddl.sql 已含同一列增量。
-- 执行前检查目标列和索引,选择尚未执行的语句;不要忽略重复列错误或使用 --force。
-- 遵循算力迁移 ID 规则:业务数值 ID 为 18 位、显式插入、无自增;本表主键沿用 SaaS 既有发号器生成的 BIGINT。
CREATE TABLE `t_saas_ai_computing_service_binding`
......
......@@ -149,7 +149,9 @@ def calling(management):
m.account_id = open_account(m)
m.gateway = CallGateway(m)
m.service.gateway = m.gateway
m.credential = issue(m, m.account_id)
# An account-level CALL key is a platform-only facility (2026-09-30): the key
# still belongs to client-a, which is what these call tests exercise.
m.credential = issue(m, m.account_id, source="platform", client_id="client-a")
activate(m, m.credential)
m.secret = m.credential["secret"]
m.calls = CallService(m.service)
......@@ -546,7 +548,7 @@ def test_rotation_and_changed_default_replay_original_snapshot_before_admission(
assert original.status_code == 200, original.text
data = original.json()["data"]
row = stored_call(m)
new = issue(m, m.account_id, request_id="rotate-p4-issue")
new = issue(m, m.account_id, request_id="rotate-p4-issue", source="platform", client_id="client-a")
activate(m, new, request_id="rotate-p4-activate", replaces=m.credential["credentialId"])
m.service.settings = replace(m.service.settings, newapi_model="replacement-model")
with m.factory.begin() as session:
......@@ -608,7 +610,7 @@ def test_results_and_request_identity_are_account_and_client_scoped(funded, scop
assert original.status_code == 200, original.text
if scope == "account":
other_account = open_account(m, request_id="other-p4-account")
other = issue(m, other_account, request_id="other-p4-key")
other = issue(m, other_account, request_id="other-p4-key", source="platform", client_id="client-a")
activate(m, other, request_id="other-p4-activate")
assert gift(m, request_id="other-p4-gift", accountId=other_account).status_code == 200
else:
......
......@@ -10,7 +10,7 @@ from sqlalchemy.orm import Session, sessionmaker
from app import account_repository as repository
from app.account_models import (
Account, AccountClient, Base, Call, Client, Consumption, Credential, Gift,
IdSegment, ManagementRequest, PointRecord,
IdSegment, ManagementRequest, PointRecord, SubAccount,
)
from app.db import create_mysql_engine
from app.id_generator import GENERATOR_KEY, ID_LOW, SEGMENT_SIZE, SegmentIDGenerator, seed_segment
......@@ -48,13 +48,15 @@ def account_engine(request):
"mysql+pymysql://root@localhost/" + database + "?unix_socket=" + str(socket_path)
)
try:
script = Path(__file__).resolve().parents[1] / "scripts/20260923_independent_account_ddl.sql"
ddl = "\n".join(line for line in script.read_text(encoding="utf-8").splitlines()
if not line.lstrip().startswith("--"))
# Fresh installations and MySQL tests use the same current full schema.
scripts = [Path(__file__).resolve().parents[1] / "scripts/sql/schema.sql"]
with engine.begin() as connection:
for statement in ddl.split(";"):
if statement.strip():
connection.execute(text(statement))
for script in scripts:
ddl = "\n".join(line for line in script.read_text(encoding="utf-8").splitlines()
if not line.lstrip().startswith("--"))
for statement in ddl.split(";"):
if statement.strip():
connection.execute(text(statement))
yield engine
finally:
engine.dispose()
......@@ -104,7 +106,7 @@ def gift(id, account_id=A, client_id="mei1-saas", request_id="same-request"):
def test_schema_has_no_saas_dependencies(account_engine):
names = inspect(account_engine).get_table_names()
assert len(names) == 11
assert len(names) == 15
assert all(name.startswith("t_computing_") for name in names)
for table in Base.metadata.tables.values():
assert not {"merchant_id", "store_id", "employee_id", "source_app"} & set(table.columns.keys())
......@@ -300,7 +302,7 @@ def test_mysql_isolation_utc_and_schema_constraints(account_engine):
def test_generated_ddl_matches_models():
from scripts.render_account_ddl import render_ddl
script = Path(__file__).resolve().parents[1] / "scripts/20260923_independent_account_ddl.sql"
script = Path(__file__).resolve().parents[1] / "scripts/sql/schema.sql"
assert script.read_text(encoding="utf-8") == render_ddl()
assert "AUTO_INCREMENT" not in render_ddl()
assert "DEFAULT (UTC_TIMESTAMP(3))" in render_ddl()
......@@ -484,14 +486,14 @@ def test_mysql_fork_while_generator_locked(account_engine):
("failed-call-success-consumption", {"consumption_call_mismatch"}),
("wrong-point-call", {"point_consumption_mismatch"}),
("null-point-call", {"point_consumption_mismatch"}),
("pending-consumption-point", {"point_consumption_mismatch"}),
("pending-consumption-point", {"point_consumption_mismatch", "consumption_bucket_orphan"}),
])
def test_mysql_target_checks_detect_inconsistent_links(account_session, account_engine, case, expected):
if account_engine.dialect.name != "mysql":
pytest.skip("目标对账 SQL 使用 MySQL NULL-safe 比较")
s = account_session
c = call(101)
s.add_all([c, call(102, request_id="other-call", billing_status="FAILED")])
c = call(101, quota_month="2026-09")
s.add_all([c, call(102, request_id="other-call", billing_status="FAILED", quota_month="2026-09")])
s.flush()
legacy = case in {"valid-legacy", "success-call-legacy-consumption"}
r = consumption(201, None if legacy else 101, legacy_record=legacy)
......@@ -520,18 +522,66 @@ def test_mysql_target_checks_detect_inconsistent_links(account_session, account_
s.get(Account, A).balance_point_units = -146
seed_segment(s, ID_LOW + 10000, segment_model=IdSegment)
s.commit()
script = Path(__file__).resolve().parents[1] / "scripts/20260923_account_target_check.sql"
assert target_check_anomalies(s) == expected
def target_check_anomalies(session):
"""Execute the delivered read-only SQL, not an ORM approximation of it."""
script = Path(__file__).resolve().parents[1] / "scripts/sql/check.sql"
sql = "\n".join(line for line in script.read_text(encoding="utf-8").splitlines()
if not line.lstrip().startswith("--"))
actual = set()
for statement in sql.split(";"):
if statement.strip():
result = s.execute(text(statement))
result = session.execute(text(statement))
if "anomaly" in result.keys():
actual.update(row.anomaly for row in result)
else:
result.all()
assert actual == expected
return actual
def test_mysql_target_checks_keep_account_buckets_separate(account_session, account_engine):
if account_engine.dialect.name != "mysql":
pytest.skip("目标对账 SQL 使用 MySQL NULL-safe 比较")
s = account_session
# Equal and opposite errors must not cancel across the two NULL buckets.
for offset, account_id, consumed, recorded in [(0, A, 146, 292), (1, B, 292, 146)]:
row = consumption(201 + offset, None, account_id=account_id,
request_id="bucket-check", legacy_record=True)
row.consumed_quota, row.consumed_point_units = consumed // 146, consumed
s.add(row)
s.flush()
s.add(PointRecord(id=301 + offset, account_id=account_id, client_id="mei1-saas",
type="CONSUME", consumption_record_id=row.id, legacy_record=True,
point_units=-recorded, balance_before_units=0, balance_after_units=-recorded))
s.get(Account, account_id).balance_point_units = -recorded
seed_segment(s, ID_LOW + 10000, segment_model=IdSegment)
s.commit()
assert "consumption_bucket_mismatch" in target_check_anomalies(s)
def test_mysql_target_checks_separate_transfer_clients(account_session, account_engine):
if account_engine.dialect.name != "mysql":
pytest.skip("目标对账 SQL 使用 MySQL NULL-safe 比较")
s = account_session
s.add(SubAccount(id=401, account_id=A, name="transfer-bucket", balance_point_units=20000))
s.flush()
for offset, client_id in enumerate(["mei1-saas", "octop"]):
before = offset * 10000
s.add_all([
PointRecord(id=301 + offset * 2, account_id=A, client_id=client_id,
type="TRANSFER_OUT", request_id="shared-transfer-id", point_units=-10000,
balance_before_units=-before, balance_after_units=-before - 10000),
PointRecord(id=302 + offset * 2, account_id=A, client_id=client_id,
type="TRANSFER_IN", request_id="shared-transfer-id", point_units=10000,
balance_before_units=before, balance_after_units=before + 10000,
sub_account_id=401),
])
s.get(Account, A).balance_point_units = -20000
seed_segment(s, ID_LOW + 10000, segment_model=IdSegment)
s.commit()
assert target_check_anomalies(s) == set()
def test_database_ready_accepts_seeded_history(account_engine, account_session):
......
......@@ -162,7 +162,7 @@ def test_gift_role_target_and_revocation_precedence(gifting):
m = gifting
assert post(m, "/internal/v1/gifts", body(m)).status_code == 403
assert gift(m, clientId="client-b").status_code == 404
key = issue(m, m.account_id)
key = issue(m, m.account_id, source="platform", client_id="client-a")
activate(m, key)
assert post(m, "/internal/v1/gifts", body(m), key=key["secret"]).status_code == 403
data = gift(m).json()["data"]
......@@ -466,6 +466,36 @@ def test_mysql_concurrent_reconcile_is_single_credit(gifting, monkeypatch):
assert m.gateway.writes == [68493]
def test_mysql_concurrent_reconcile_replay_never_returns_stale_view(gifting, monkeypatch):
"""同 requestId 的并发对账:败者必须回执最新视图,不能给过期 version/phase。
修复前 `_reconcile_state` 先读 gift 行、后查幂等行:另一个线程刚提交这笔对账时,
败者会拿到"operationStatus=SUCCEEDED + 提交前的 version/phase"(HTTP 202 202 的
过期回执),或被自己过期的预检版本判成版本冲突。这里连跑 3 轮做护栏。
"""
m = gifting
if m.engine.dialect.name != "mysql":
pytest.skip("requires MySQL row locks")
for round_index in range(3):
m.gateway.gift_failure = "lost-after"
data = gift(m, request_id="gift-request-%03d" % round_index).json()["data"]
# 每轮都要越过该笔赠送的派发窗口:时钟按轮次递增(DISPATCH_WINDOW_OPEN 会拦未过窗的)
later = utc_now() + timedelta(minutes=10 * (round_index + 1))
monkeypatch.setattr(gifts_module, "utc_now", lambda later=later: later)
m.gateway.check_transactions = False
request_id = "reconcile-request-%03d" % round_index
def confirm(_):
return reconcile(m, data, "CONFIRM_APPLIED", request_id=request_id)
with ThreadPoolExecutor(max_workers=3) as pool:
responses = list(pool.map(confirm, range(3)))
assert all(response.status_code == 200 for response in responses), [
(response.status_code, response.text[:200]) for response in responses]
assert len({response.json()["data"]["version"] for response in responses}) == 1
assert {response.json()["data"]["phase"] for response in responses} == {"CONFIRMED"}
def test_mysql_different_gifts_contend_on_durable_gate(gifting):
import threading
m = gifting
......@@ -499,7 +529,7 @@ def test_page_routes_are_readonly_and_validate_time_scope(gifting):
for values in [{"start": "2026-09-23T00:00:00"}, {"start": "2026-09-24T00:00:00Z", "end": "2026-09-23T00:00:00Z"},
{"page": 0}, {"size": 201}, {"accountIds": [m.account_id] * 1001}]:
assert post(m, "/internal/v1/gifts/page", values).status_code == 400
key = issue(m, m.account_id)
key = issue(m, m.account_id, source="platform", client_id="client-a")
activate(m, key)
assert post(m, "/api/v1/consumptions/page", {}, key=key["secret"]).status_code == 200
assert post(m, "/api/v1/consumptions/page", {"accountIds": [m.account_id]}, key=key["secret"]).status_code == 403
......@@ -316,7 +316,7 @@ def test_platform_requires_explicit_scope_validates_all_ids_and_audits_reads(led
def test_call_scope_is_fixed_and_gifts_forbidden(ledger):
m = ledger.m
credential = issue(m, str(ledger.a))
credential = issue(m, str(ledger.a), source="platform", client_id="client-a")
activate(m, credential)
secret = credential["secret"]
result = page_consumptions(m.service, secret)
......@@ -392,7 +392,7 @@ def test_queries_select_only_financial_columns_and_page_in_sql(ledger, function,
secret = m.keys["platform"] if role == "PLATFORM" else m.keys["client-a"]
parameters = {"account_ids": [str(ledger.a)]} if role == "PLATFORM" else {}
if role == "CALL":
credential = issue(m, str(ledger.a))
credential = issue(m, str(ledger.a), source="platform", client_id="client-a")
activate(m, credential)
secret = credential["secret"]
statements = []
......
"""分账(段 3a):消费投影带出计费主体,且可按门店/员工/子账号维度分页。
Rows are inserted directly (same style as tests/test_account_ledger.py): the
projection and the filters are what is under test, not the gateway plumbing.
"""
import pytest
from app.account_ledger import page_consumptions, page_gifts
from app.account_security import AccountError
from tests.test_account_foundation import account_engine # noqa: F401 (fixture plumbing)
from tests.test_account_ledger import ID, NOW, _call, _consumption
from tests.test_account_management import management, open_account
from tests.test_account_scope import SCOPE_PATH, bind_employee, bind_store, issue_scope_key
from tests.test_account_sub_account import create_sub_account
def _world(m):
"""One account, one sub-account bucket, two stores (keyed / direct), one
employee, and one consumption row per charging subject."""
account_id = open_account(m)
sub_id = int(create_sub_account(m, account_id, "ledger-budget",
request_id="ledger-sub").json()["data"]["subAccountId"])
store_a = int(bind_store(m, account_id, "mei1:ledger-store-a", sub_id,
request_id="ledger-store-a").json()["data"]["scopeId"])
store_b = int(bind_store(m, account_id, "mei1:ledger-store-b",
request_id="ledger-store-b").json()["data"]["scopeId"])
employee = int(bind_employee(m, account_id, "mei1:ledger-emp", "mei1:ledger-store-a",
request_id="ledger-emp").json()["data"]["scopeId"])
with m.factory.begin() as session:
session.add_all([
_call(201, account_id, request_id="ledger-call-201", quota_month="2026-09"),
_call(202, account_id, request_id="ledger-call-202", quota_month="2026-09",
sub_account_id=sub_id, scope_id=store_a, store_scope_id=store_a),
_call(203, account_id, request_id="ledger-call-203", quota_month="2026-09",
sub_account_id=sub_id, scope_id=employee, store_scope_id=store_a),
_call(204, account_id, request_id="ledger-call-204", quota_month="2026-09",
scope_id=store_b, store_scope_id=store_b),
_consumption(301, account_id, scope_id=store_b, store_scope_id=store_b),
_consumption(302, account_id),
])
return account_id, sub_id, store_a, store_b, employee
def _page(m, key, **kwargs):
return page_consumptions(m.service, key, **kwargs)
def _ids(result):
return {row["id"] for row in result["list"]}
def test_consumption_projection_exposes_the_charging_subject(management):
m = management
account_id, sub_id, store_a, store_b, employee = _world(m)
result = _page(m, m.keys["client-a"], account_ids=[account_id], size=200)
rows = {row["id"]: row for row in result["list"]}
assert set(rows) == {str(ID + offset) for offset in [201, 202, 203, 204, 301, 302]}
# Account-level money answers to the account bucket and no subject at all.
assert (rows[str(ID + 201)]["subAccountId"], rows[str(ID + 201)]["scopeId"],
rows[str(ID + 201)]["storeScopeId"]) == (None, None, None)
assert rows[str(ID + 202)]["subAccountId"] == str(sub_id)
assert rows[str(ID + 202)]["scopeId"] == str(store_a)
assert rows[str(ID + 202)]["storeScopeId"] == str(store_a)
# An employee row carries its own scope and the store it belonged to.
assert rows[str(ID + 203)]["scopeId"] == str(employee)
assert rows[str(ID + 203)]["storeScopeId"] == str(store_a)
# A store keyed straight to the account charges no bucket but keeps its scope.
assert rows[str(ID + 204)]["subAccountId"] is None
assert rows[str(ID + 204)]["scopeId"] == rows[str(ID + 204)]["storeScopeId"] == str(store_b)
# Imported rows go through the same projection.
assert rows[str(ID + 301)]["storeScopeId"] == str(store_b)
assert rows[str(ID + 302)]["scopeId"] is None
def test_store_ledger_merges_the_store_and_its_employees(management):
m = management
account_id, sub_id, store_a, store_b, employee = _world(m)
result = _page(m, m.keys["client-a"], account_ids=[account_id], store_scope_ids=[store_a])
assert result["total"] == 2
assert _ids(result) == {str(ID + 202), str(ID + 203)}
other = _page(m, m.keys["client-a"], account_ids=[account_id], store_scope_ids=[store_b])
assert _ids(other) == {str(ID + 204), str(ID + 301)}
both = _page(m, m.keys["client-a"], account_ids=[account_id], store_scope_ids=[store_a, store_b])
assert both["total"] == 4
def test_employee_and_sub_account_ledgers_use_their_own_dimension(management):
m = management
account_id, sub_id, store_a, store_b, employee = _world(m)
employee_page = _page(m, m.keys["client-a"], account_ids=[account_id], scope_ids=[employee])
assert _ids(employee_page) == {str(ID + 203)}
# The bucket dimension is not the store dimension: store B is keyed straight
# to the account, so it never shows up in the sub-account's ledger.
bucket_page = _page(m, m.keys["client-a"], account_ids=[account_id], sub_account_ids=[sub_id])
assert _ids(bucket_page) == {str(ID + 202), str(ID + 203)}
combos = _page(m, m.keys["client-a"], account_ids=[account_id], scope_ids=[employee],
sub_account_ids=[sub_id], store_scope_ids=[store_a])
assert _ids(combos) == {str(ID + 203)}
def test_scope_filters_page_consistently_and_never_leak_across_accounts(management):
m = management
account_id, sub_id, store_a, store_b, employee = _world(m)
collected = []
for page in range(1, 3):
result = _page(m, m.keys["client-a"], account_ids=[account_id], store_scope_ids=[store_a],
page=page, size=1)
assert result["total"] == 2
collected.extend(row["id"] for row in result["list"])
assert collected == [str(ID + 203), str(ID + 202)]
# A scope of another account simply cannot widen the account scope.
other_account = open_account(m, request_id="ledger-other-account")
other_store = int(bind_store(m, other_account, "mei1:ledger-store-x",
request_id="ledger-store-x").json()["data"]["scopeId"])
empty = _page(m, m.keys["client-a"], account_ids=[account_id],
store_scope_ids=[store_a, other_store])
assert empty["total"] == 2
unrelated = _page(m, m.keys["client-a"], account_ids=[account_id], store_scope_ids=[other_store])
assert (unrelated["total"], unrelated["list"]) == (0, [])
def test_ledger_scope_filters_are_refused_to_call_keys(management):
m = management
account_id, sub_id, store_a, store_b, employee = _world(m)
secret = issue_scope_key(m, account_id, store_a, request_id="ledger-store-key")["secret"]
for kwargs in [{"store_scope_ids": [store_a]}, {"scope_ids": [employee]},
{"sub_account_ids": [sub_id]}]:
with pytest.raises(AccountError) as error:
page_consumptions(m.service, secret, **kwargs)
assert error.value.code == "PERMISSION_DENIED"
# The key can still read its own account's ledger with no identity override.
allowed = page_consumptions(m.service, secret)
assert allowed["total"] == 6
def test_gift_page_rejects_subject_filters(management):
m = management
account_id, sub_id, store_a, _, _ = _world(m)
with pytest.raises(AccountError) as error:
page_gifts(m.service, m.keys["client-a"], account_ids=[account_id], store_scope_ids=[store_a])
assert error.value.code == "INVALID_ARGUMENT"
......@@ -134,11 +134,11 @@ def test_independent_happy_path_and_one_time_delivery(management):
m = management
account_id = open_account(m)
assert open_account(m) == account_id and m.gateway.creates == 1
created = issue(m, account_id)
created = issue(m, account_id, source="platform", client_id="client-a")
assert created["status"] == "PENDING" and created["secretAvailable"]
denied = m.client.get("/api/v1/account", headers=headers(m, key=created["secret"]))
assert denied.status_code == 401 and denied.json()["code"] == "CREDENTIAL_PENDING"
replay = issue(m, account_id)
replay = issue(m, account_id, source="platform", client_id="client-a")
assert replay["credentialId"] == created["credentialId"] and not replay["secretAvailable"] and "secret" not in replay
activate(m, created)
response = m.client.get("/api/v1/account", headers=headers(m, key=created["secret"]))
......@@ -169,7 +169,7 @@ def test_roles_scope_and_forged_identity(management):
response = m.client.post("/internal/v1/accounts/%s/status" % account_id,
json={"status": "DISABLED", "reason": "test"}, headers=forged)
assert response.status_code == 403
created = issue(m, account_id)
created = issue(m, account_id, source="platform", client_id="client-a")
activate(m, created)
assert post(m, "/internal/v1/accounts", {"name": "Forbidden"}, key=created["secret"]).status_code == 403
assert post(m, "/internal/v1/credentials/%s/activate" % created["credentialId"], {}, source="client-b").status_code == 404
......@@ -179,15 +179,20 @@ def test_roles_scope_and_forged_identity(management):
def test_revocation_precedes_replay_and_new_issuance(management):
m = management
account_id = open_account(m)
created = issue(m, account_id)
sub_id = post(m, "/internal/v1/accounts/%s/sub-accounts" % account_id, {"name": "budget"},
source="platform", request_id="create-sub-revoke").json()["data"]["subAccountId"]
created = post(m, "/internal/v1/accounts/%s/credentials" % account_id,
{"subAccountId": str(sub_id)}, request_id="issue-key-001").json()["data"]
activate(m, created)
revoked = post(m, "/internal/v1/accounts/%s/clients" % account_id,
{"clientId": "client-a", "status": "REVOKED", "reason": "revoke access"}, source="platform")
assert revoked.status_code == 200
assert m.client.get("/api/v1/account", headers=headers(m, key=created["secret"])).status_code == 403
assert post(m, "/internal/v1/accounts/%s/credentials" % account_id, {}, request_id="issue-key-001").status_code == 403
assert post(m, "/internal/v1/accounts/%s/credentials" % account_id, {"subAccountId": str(sub_id)},
request_id="issue-key-001").status_code == 403
assert m.client.get("/internal/v1/requests/ISSUE_CALL_KEY/issue-key-001", headers=headers(m)).status_code == 403
assert post(m, "/internal/v1/accounts/%s/credentials" % account_id, {}, request_id="issue-key-new").status_code == 403
assert post(m, "/internal/v1/accounts/%s/credentials" % account_id, {"subAccountId": str(sub_id)},
request_id="issue-key-new").status_code == 403
def test_platform_explicit_target_and_grant(management):
......@@ -207,9 +212,9 @@ def test_platform_explicit_target_and_grant(management):
def test_rotation_grace_does_not_extend_and_revoke_is_immediate(management):
m = management
account_id = open_account(m)
old = issue(m, account_id)
old = issue(m, account_id, source="platform", client_id="client-a")
activate(m, old)
new = issue(m, account_id, request_id="issue-key-002")
new = issue(m, account_id, request_id="issue-key-002", source="platform", client_id="client-a")
activate(m, new, request_id="rotate-key-001", replaces=old["credentialId"])
with m.factory() as s:
deadline = s.get(Credential, int(old["credentialId"])).valid_until
......@@ -228,12 +233,12 @@ def test_rotation_grace_does_not_extend_and_revoke_is_immediate(management):
def test_expiry_boundaries_and_disabled_client(management):
m = management
account_id = open_account(m)
created = issue(m, account_id)
created = issue(m, account_id, source="platform", client_id="client-a")
with m.factory.begin() as s:
s.get(Credential, int(created["credentialId"])).pending_expires_at = utc_now() - timedelta(seconds=1)
response = post(m, "/internal/v1/credentials/%s/activate" % created["credentialId"], {})
assert response.json()["code"] == "CREDENTIAL_EXPIRED"
newer = issue(m, account_id, request_id="newer-key-001")
newer = issue(m, account_id, request_id="newer-key-001", source="platform", client_id="client-a")
activate(m, newer)
with m.factory.begin() as s:
s.get(Credential, int(newer["credentialId"])).valid_until = utc_now() - timedelta(seconds=1)
......@@ -336,9 +341,10 @@ def segment_outage(monkeypatch, generator):
def test_credential_replay_survives_segment_failure(management, monkeypatch):
m = management
account_id = open_account(m)
created = issue(m, account_id, request_id="segment-001")
created = issue(m, account_id, request_id="segment-001", source="platform", client_id="client-a")
segment_outage(monkeypatch, m.service.generator)
replay = post(m, "/internal/v1/accounts/%s/credentials" % account_id, {}, request_id="segment-001")
replay = post(m, "/internal/v1/accounts/%s/credentials" % account_id, {"clientId": "client-a"},
request_id="segment-001", source="platform")
assert replay.status_code == 200, replay.text
assert replay.json()["data"]["credentialId"] == created["credentialId"]
assert "secret" not in replay.json()["data"]
......@@ -369,15 +375,19 @@ def test_cancelled_provision_preserves_dispatch_evidence(management):
def test_disabled_account_allows_reads_and_does_not_reset_provision(management):
m = management
account_id = open_account(m)
created = issue(m, account_id)
activate(m, created)
sub_id = post(m, "/internal/v1/accounts/%s/sub-accounts" % account_id, {"name": "budget"},
source="platform", request_id="create-sub-disabled").json()["data"]["subAccountId"]
created = post(m, "/internal/v1/accounts/%s/credentials" % account_id, {"subAccountId": str(sub_id)},
request_id="issue-sub-disabled").json()["data"]
activate(m, created, request_id="activate-sub-disabled")
response = post(m, "/internal/v1/accounts/%s/status" % account_id,
{"status": "DISABLED", "reason": "maintenance"}, source="platform")
assert response.status_code == 200
response = m.client.get("/api/v1/account", headers=headers(m, key=created["secret"]))
assert response.status_code == 200 and "ACCOUNT_DISABLED" in response.json()["data"]["blockedReasons"]
assert open_account(m) == account_id
assert post(m, "/internal/v1/accounts/%s/credentials" % account_id, {}, request_id="disabled-issue-01").status_code == 403
assert post(m, "/internal/v1/accounts/%s/credentials" % account_id, {"subAccountId": str(sub_id)},
request_id="disabled-issue-01").status_code == 403
def test_bootstrap_is_idempotent_and_callback_failure_rolls_back(management):
......@@ -479,9 +489,9 @@ def test_same_client_management_roles_cannot_replay_privileged_results(managemen
def test_disabled_account_rejects_new_activation_but_allows_original_replay(management):
m = management
account_id = open_account(m)
old = issue(m, account_id)
old = issue(m, account_id, source="platform", client_id="client-a")
activate(m, old)
new = issue(m, account_id, request_id="issue-while-active")
new = issue(m, account_id, request_id="issue-while-active", source="platform", client_id="client-a")
assert post(m, "/internal/v1/accounts/%s/status" % account_id,
{"status": "DISABLED", "reason": "maintenance"}, source="platform").status_code == 200
response = post(m, "/internal/v1/credentials/%s/activate" % new["credentialId"],
......@@ -611,7 +621,7 @@ def test_mysql_concurrent_provision_and_issuance_are_once(management):
account_id = int(outcomes[0]["accountId"])
def issue_once(_):
return m.service.issue_credential(m.keys["client-a"], "concurrent-issue-01", account_id)
return m.service.issue_credential(m.keys["platform"], "concurrent-issue-01", account_id, "client-a")
with ThreadPoolExecutor(max_workers=4) as executor:
outcomes = list(executor.map(issue_once, range(4)))
......
"""超额/窗口可观测(段 3c):S31 结算越限、S30 结算重试、S16 调低越限。
三条都必须给出**结构化** warning(key=value 字段),运维才能据此设阈值;且只在
「跨越那一刻」记一条,不在超额期间逐笔刷屏。
"""
import asyncio
import logging
from datetime import timedelta
from app.account_ledger import current_month
from app.account_models import Call, ScopeMonthUsage
from tests.test_account_foundation import account_engine # noqa: F401 (fixture plumbing)
from tests.test_account_management import management, open_account
from tests.test_account_scope import bind_employee, bind_store, month_usage, set_quota
from tests.test_account_settlement import evidence_for, new_call, process, settlements # noqa: F401
def _world(m, tag):
account_id = m.account_id
store_id = int(bind_store(m, account_id, "mei1:alert-store-" + tag,
request_id="alert-store-" + tag).json()["data"]["scopeId"])
employee_id = int(bind_employee(m, account_id, "mei1:alert-emp-" + tag, "mei1:alert-store-" + tag,
request_id="alert-emp-" + tag).json()["data"]["scopeId"])
return store_id, employee_id
def _messages(caplog, marker):
return [record.getMessage() for record in caplog.records if marker in record.getMessage()]
def test_settlement_warns_once_when_the_month_usage_crosses_the_limit(settlements, caplog):
m = settlements
store_id, employee_id = _world(m, "cross")
assert set_quota(m, m.account_id, employee_id, store_id, 1,
request_id="alert-quota-cross").status_code == 200
with m.factory.begin() as session:
session.add(ScopeMonthUsage(scope_id=employee_id, month="2026-09", used_point_units=9900))
call_id = new_call(m, scope_id=employee_id, store_scope_id=store_id, quota_month="2026-09")
evidence_for(m, call_id, quota=1)
with caplog.at_level(logging.WARNING, logger="app.account_settlement"):
assert process(m, call_id)
assert month_usage(m, employee_id, "2026-09") == 10046
over = _messages(caplog, "scope month quota exceeded")
assert len(over) == 1
for token in ["accountId=%s" % m.account_id, "scopeId=%s" % employee_id,
"storeScopeId=%s" % store_id, "quotaMonth=2026-09", "limitPointUnits=10000",
"usedPointUnits=10046", "overPointUnits=46", "callId=%s" % call_id]:
assert token in over[0], (token, over[0])
def test_settlement_inside_the_limit_and_repeat_overage_stay_quiet(settlements, caplog):
m = settlements
store_id, employee_id = _world(m, "quiet")
assert set_quota(m, m.account_id, employee_id, store_id, 1,
request_id="alert-quota-quiet").status_code == 200
inside = new_call(m, scope_id=employee_id, store_scope_id=store_id, quota_month="2026-09")
evidence_for(m, inside, quota=1)
with caplog.at_level(logging.WARNING, logger="app.account_settlement"):
assert process(m, inside)
assert _messages(caplog, "scope month quota exceeded") == []
# Already over the limit before this settlement: the crossing was reported
# once already, so the stream stays quiet while the state persists.
m.clock.now += timedelta(minutes=5)
over = new_call(m, scope_id=employee_id, store_scope_id=store_id, quota_month="2026-09")
with m.factory.begin() as session:
session.get(Call, over).create_time = m.clock.now
session.get(ScopeMonthUsage, (employee_id, "2026-09")).used_point_units = 50000
evidence_for(m, over, quota=1)
caplog.clear()
with caplog.at_level(logging.WARNING, logger="app.account_settlement"):
assert process(m, over)
assert month_usage(m, employee_id, "2026-09") == 50146
assert _messages(caplog, "scope month quota exceeded") == []
def test_settlement_retry_and_giving_up_are_logged(settlements, caplog):
m = settlements
call_id = new_call(m)
with caplog.at_level(logging.WARNING, logger="app.account_settlement"):
assert asyncio.run(m.settlements.process(m.settlements.claim(call_id))) is False
retry = _messages(caplog, "settlement retry scheduled")
assert len(retry) == 1
for token in ["accountId=%s" % m.account_id, "callId=%s" % call_id, "code=SETTLEMENT_LOG_PENDING",
"retryCount=1", "nextRetryMinutes=1"]:
assert token in retry[0], (token, retry[0])
# 24 次之后放弃:同样要有一条能设告警的终态日志。
m.clock.now += timedelta(minutes=1)
with m.factory.begin() as session:
session.get(Call, call_id).retry_count = 24
caplog.clear()
with caplog.at_level(logging.WARNING, logger="app.account_settlement"):
assert asyncio.run(m.settlements.process(m.settlements.claim(call_id))) is False
stopped = _messages(caplog, "settlement stopped")
assert len(stopped) == 1
for token in ["callId=%s" % call_id, "code=SETTLEMENT_LOG_PENDING", "retryCount=25",
"status=SETTLE_FAILED"]:
assert token in stopped[0], (token, stopped[0])
def test_lowering_the_quota_below_this_month_usage_is_logged(management, caplog):
m = management
m.account_id = int(open_account(m))
store_id, employee_id = _world(m, "lower")
assert set_quota(m, m.account_id, employee_id, store_id, 5,
request_id="alert-quota-lower-1").status_code == 200
with m.factory.begin() as session:
session.add(ScopeMonthUsage(scope_id=employee_id, month=current_month(), used_point_units=60000))
with caplog.at_level(logging.WARNING, logger="app.account_service"):
assert set_quota(m, m.account_id, employee_id, store_id, 1,
request_id="alert-quota-lower-2").status_code == 200
notes = _messages(caplog, "scope quota below current month usage")
assert len(notes) == 1
for token in ["scopeId=%s" % employee_id, "storeScopeId=%s" % store_id,
"quotaMonth=%s" % current_month(), "limitPointUnits=10000", "usedPointUnits=60000"]:
assert token in notes[0], (token, notes[0])
# Raising it again is not a risk: no log.
caplog.clear()
with caplog.at_level(logging.WARNING, logger="app.account_service"):
assert set_quota(m, m.account_id, employee_id, store_id, 50,
request_id="alert-quota-lower-3").status_code == 200
assert _messages(caplog, "scope quota below current month usage") == []
"""分账守恒 E2E(段 3d):真实调用穿过结算后,账本/月账/桶流水必须互相对得上。
覆盖 `scripts/sql/check.sql` B 部分第 9/11/13/15 节;MySQL 分支同时执行
交付的完整检查 SQL,防止 ORM 断言通过但 SQL 口径仍然误报。
"""
from sqlalchemy import func, select
from app.account_ledger import page_consumptions, page_month_usage
from app.account_models import Consumption, PointRecord
from tests.test_account_calls import calling, chat, chat_body, gift # noqa: F401 (fixture/helpers)
from tests.test_account_foundation import account_engine, target_check_anomalies # noqa: F401
from tests.test_account_management import management # noqa: F401 (fixture plumbing)
from tests.test_account_scope import (bind_employee, bind_store, buckets,
issue_scope_key, month_usage, set_quota)
from tests.test_account_sub_account import create_sub_account, transfer
def _charge(m, request_id, key):
response = chat(m, request_id=request_id, key=key, body=chat_body(businessRef=request_id))
assert response.status_code == 200, response.text
return int(response.json()["data"]["consumedPointUnits"])
def _ledger(m, account_id, **filters):
return page_consumptions(m.service, m.keys["client-a"], account_ids=[account_id], size=200, **filters)
def _sum(result):
return sum(int(row["consumedPointUnits"]) for row in result["list"])
def test_scoped_money_matches_between_ledger_month_bucket_and_point_records(calling):
m = calling
account_id = m.account_id
# Use the gift flow so the delivered SQL can verify the opening balance too.
response = gift(m)
assert response.status_code == 200, response.text
sub_response = create_sub_account(m, account_id, "e2e-budget", request_id="e2e-sub-0001")
assert sub_response.status_code == 200, sub_response.text
sub_id = int(sub_response.json()["data"]["subAccountId"])
assert transfer(m, account_id, sub_id, 500, "OUT", request_id="e2e-transfer").status_code == 200
store_a = int(bind_store(m, account_id, "mei1:e2e-store-a", sub_id,
request_id="e2e-store-a").json()["data"]["scopeId"])
store_b = int(bind_store(m, account_id, "mei1:e2e-store-b",
request_id="e2e-store-b").json()["data"]["scopeId"])
employee = int(bind_employee(m, account_id, "mei1:e2e-emp", "mei1:e2e-store-a",
request_id="e2e-emp-0001").json()["data"]["scopeId"])
assert set_quota(m, account_id, employee, store_a, 100, request_id="e2e-quota").status_code == 200
employee_key = issue_scope_key(m, account_id, employee, request_id="e2e-emp-key")["secret"]
store_a_key = issue_scope_key(m, account_id, store_a, request_id="e2e-store-a-key")["secret"]
store_b_key = issue_scope_key(m, account_id, store_b, request_id="e2e-store-b-key")["secret"]
account_before, bucket_before = buckets(m, sub_id)
employee_units = [_charge(m, "e2e-emp-%s" % tag, employee_key) for tag in ("01", "02")]
store_a_units = _charge(m, "e2e-store-a-01", store_a_key)
store_b_units = _charge(m, "e2e-store-b-01", store_b_key)
account_units = _charge(m, "e2e-account-01", m.secret)
account_after, bucket_after = buckets(m, sub_id)
# 1) 每笔消费带自己的计费主体(分账口径)。
rows = {row["requestId"]: row for row in _ledger(m, account_id)["list"]}
assert len(rows) == 5
for tag in ("01", "02"):
row = rows["e2e-emp-%s" % tag]
assert row["scopeId"] == str(employee) and row["storeScopeId"] == str(store_a)
assert row["subAccountId"] == str(sub_id)
assert rows["e2e-store-a-01"]["scopeId"] == rows["e2e-store-a-01"]["storeScopeId"] == str(store_a)
assert rows["e2e-store-b-01"]["subAccountId"] is None
assert rows["e2e-account-01"]["scopeId"] is None and rows["e2e-account-01"]["storeScopeId"] is None
# 2) 门店账 = 该店自身 + 其下员工;直挂门店与账户级调用不进这一档。
store_ledger = _ledger(m, account_id, store_scope_ids=[store_a])
assert store_ledger["total"] == 3
assert _sum(store_ledger) == sum(employee_units) + store_a_units
assert _sum(_ledger(m, account_id, scope_ids=[employee])) == sum(employee_units)
# 3) 月账单 = 员工当月真实消耗,且带生效额度;额度与实际口径同源。
month = _ledger(m, account_id)["list"][0]["createTime"][:7].replace("-", "-")
usage = page_month_usage(m.service, m.keys["client-a"], account_ids=[account_id], size=200)
assert usage["total"] == 1
row = usage["list"][0]
assert row["scopeId"] == str(employee) and row["storeScopeId"] == str(store_a)
assert row["usedPointUnits"] == str(sum(employee_units)) == str(month_usage(m, employee))
assert row["limitPointUnits"] == "1000000" and row["unlimited"] is False
assert month == row["month"]
# 4) 桶守恒:有桶的调用扣子账号桶,直挂门店与账户级调用扣账户未分配余额。
assert bucket_before - bucket_after == sum(employee_units) + store_a_units
assert account_before - account_after == store_b_units + account_units
# 5) 对账 #15:消费按桶聚合 = CONSUME 流水的同桶聚合(账户桶与子账号桶各自成立)。
with m.factory() as session:
consumption = dict(session.execute(
select(Consumption.sub_account_id, func.sum(Consumption.consumed_point_units))
.where(Consumption.settlement_status == "SUCCESS").group_by(Consumption.sub_account_id)).all())
records = dict(session.execute(
select(PointRecord.sub_account_id, func.sum(-PointRecord.point_units))
.where(PointRecord.type == "CONSUME").group_by(PointRecord.sub_account_id)).all())
assert consumption == records
assert consumption[None] == store_b_units + account_units
assert consumption[sub_id] == sum(employee_units) + store_a_units
if m.engine.dialect.name == "mysql":
with m.factory() as session:
assert target_check_anomalies(session) == set()
"""月账单直读(段 3b):`scope_month_usage` 分页 + 生效额度 + 门店/员工/子账号维度。
Rows are inserted directly, in the style of tests/test_account_ledger_scope.py.
"""
import pytest
from app.account_ledger import page_month_usage
from app.account_security import AccountError
from tests.test_account_foundation import account_engine # noqa: F401 (fixture plumbing)
from tests.test_account_management import management, open_account
from tests.test_account_scope import bind_employee, bind_store, month_usage, set_quota
from tests.test_account_sub_account import create_sub_account
CURRENT = "2026-09"
EARLIER = "2026-08"
def _usage(m, scope_id, month, units):
from app.account_models import ScopeMonthUsage
with m.factory.begin() as session:
session.add(ScopeMonthUsage(scope_id=scope_id, month=month, used_point_units=units))
assert month_usage(m, scope_id, month) == units
def _world(m):
"""Two stores under one bucket, three employees (one unlimited), usage rows
in the current and in an earlier month, and a store keyed straight to the
account so the bucket filter has something to exclude."""
account_id = open_account(m)
sub_id = int(create_sub_account(m, account_id, "usage-budget",
request_id="usage-sub").json()["data"]["subAccountId"])
store_a = int(bind_store(m, account_id, "mei1:usage-store-a", sub_id,
request_id="usage-store-a").json()["data"]["scopeId"])
store_b = int(bind_store(m, account_id, "mei1:usage-store-b",
request_id="usage-store-b").json()["data"]["scopeId"])
employees = [int(bind_employee(m, account_id, "mei1:usage-emp-%s" % tag, store,
request_id="usage-emp-" + tag).json()["data"]["scopeId"])
for tag, store in (("1", "mei1:usage-store-a"), ("2", "mei1:usage-store-a"),
("3", "mei1:usage-store-b"))]
assert set_quota(m, account_id, employees[0], store_a, 100,
request_id="usage-quota-1").status_code == 200
_usage(m, employees[0], CURRENT, 25000000)
_usage(m, employees[0], EARLIER, 10000000)
_usage(m, employees[1], CURRENT, 45000000)
_usage(m, employees[2], CURRENT, 15000000)
return account_id, sub_id, store_a, store_b, employees
def _page(m, **kwargs):
return page_month_usage(m.service, m.keys["client-a"], **kwargs)
def test_month_usage_page_reports_usage_scope_and_effective_limit(management):
m = management
account_id, sub_id, store_a, store_b, employees = _world(m)
result = _page(m, account_ids=[account_id], month=CURRENT, size=200)
assert (result["total"], result["month"]) == (3, CURRENT)
# Highest consumer first, then by id: deterministic paging without a sort key.
assert [row["scopeId"] for row in result["list"]] == [str(value) for value in
[employees[1], employees[0], employees[2]]]
rows = {row["scopeId"]: row for row in result["list"]}
budgeted = rows[str(employees[0])]
assert budgeted["scopeType"] == "EMPLOYEE" and budgeted["scopeKey"] == "mei1:usage-emp-1"
assert budgeted["storeScopeId"] == str(store_a) and budgeted["storeScopeKey"] == "mei1:usage-store-a"
assert budgeted["subAccountId"] == str(sub_id)
assert budgeted["month"] == CURRENT
assert budgeted["usedPointUnits"] == "25000000" and budgeted["usedPoints"] == "2500.0000"
assert budgeted["limitPointUnits"] == "1000000" and budgeted["limitPoints"] == "100.0000"
assert budgeted["unlimited"] is False
assert budgeted["scopeStatus"] == "ACTIVE" and budgeted["storeStatus"] == "ACTIVE"
# Never configured at this store: unlimited, and the bucket is the account's.
unbudgeted = rows[str(employees[1])]
assert (unbudgeted["limitPointUnits"], unbudgeted["limitPoints"], unbudgeted["unlimited"]) == (None, None, True)
assert unbudgeted["usedPointUnits"] == "45000000"
direct = rows[str(employees[2])]
assert direct["subAccountId"] is None and direct["storeScopeId"] == str(store_b)
def test_month_usage_keeps_months_apart_and_leaves_history_limitless(management):
m = management
account_id, sub_id, store_a, store_b, employees = _world(m)
earlier = _page(m, account_ids=[account_id], month=EARLIER, size=200)
assert earlier["total"] == 1 and earlier["month"] == EARLIER
row = earlier["list"][0]
assert row["scopeId"] == str(employees[0]) and row["usedPointUnits"] == "10000000"
# A past month's effective quota can only be replayed from the audit trail
# (design §3.3), so this endpoint reports no limit for history at all.
assert (row["limitPointUnits"], row["limitPoints"], row["unlimited"]) == (None, None, None)
# The current month is a different page with its own rows.
assert _page(m, account_ids=[account_id], month=CURRENT)["total"] == 3
def test_month_usage_filters_by_store_employee_and_bucket(management):
m = management
account_id, sub_id, store_a, store_b, employees = _world(m)
store_page = _page(m, account_ids=[account_id], month=CURRENT, store_scope_ids=[store_a])
assert {row["scopeId"] for row in store_page["list"]} == {str(employees[0]), str(employees[1])}
single = _page(m, account_ids=[account_id], month=CURRENT, scope_ids=[employees[2]])
assert [row["scopeId"] for row in single["list"]] == [str(employees[2])]
bucket = _page(m, account_ids=[account_id], month=CURRENT, sub_account_ids=[sub_id])
assert {row["scopeId"] for row in bucket["list"]} == {str(employees[0]), str(employees[1])}
both = _page(m, account_ids=[account_id], month=CURRENT, store_scope_ids=[store_a],
scope_ids=[employees[0]])
assert [row["scopeId"] for row in both["list"]] == [str(employees[0])]
def test_month_usage_pages_without_gaps_and_stays_inside_the_account(management):
m = management
account_id, sub_id, store_a, store_b, employees = _world(m)
seen = []
for page in range(1, 3):
result = _page(m, account_ids=[account_id], month=CURRENT, page=page, size=2)
assert (result["total"], result["page"], result["size"]) == (3, page, 2)
seen.extend(row["scopeId"] for row in result["list"])
assert seen == [str(value) for value in [employees[1], employees[0], employees[2]]]
outside = _page(m, account_ids=[account_id], month=CURRENT, store_scope_ids=[store_b],
scope_ids=[employees[0]])
assert (outside["total"], outside["list"]) == (0, [])
def test_month_usage_refuses_call_keys_and_bad_months(management):
m = management
account_id, sub_id, store_a, store_b, employees = _world(m)
from tests.test_account_scope import issue_scope_key
secret = issue_scope_key(m, account_id, store_a, request_id="usage-store-key")["secret"]
with pytest.raises(AccountError) as error:
page_month_usage(m.service, secret, month=CURRENT)
assert error.value.code == "PERMISSION_DENIED"
with pytest.raises(AccountError) as error:
_page(m, account_ids=[account_id], month="2026-13")
assert error.value.code == "INVALID_ARGUMENT"
with pytest.raises(AccountError) as error:
page_month_usage(m.service, m.keys["platform"], month=CURRENT)
assert error.value.code == "INVALID_ARGUMENT"
# INTEGRATION without an explicit target reads the accounts it owns - same
# rule as the ledger pages, and it cannot see past that ownership.
own = _page(m, month=CURRENT)
assert own["total"] == 3
assert {row["scopeId"] for row in own["list"]} == {str(value) for value in employees}
Markdown is supported
0% or
You are about to add 0 people to the discussion. Proceed with caution.
Finish editing this message first!
Please register or to comment