Compare commits
3 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| 5871c02055 | |||
| 14a9473323 | |||
| 2a98e5147c |
@@ -1,211 +0,0 @@
|
||||
"""correct DeepSeek V4 Flash token price and frozen historical costs
|
||||
|
||||
Revision ID: deepseek_v4_flash_price
|
||||
Revises: comparison_below_min_success
|
||||
"""
|
||||
from __future__ import annotations
|
||||
|
||||
import json
|
||||
from collections.abc import Sequence
|
||||
from typing import Any
|
||||
|
||||
import sqlalchemy as sa
|
||||
from sqlalchemy.dialects import postgresql
|
||||
|
||||
from alembic import op
|
||||
|
||||
revision: str = "deepseek_v4_flash_price"
|
||||
down_revision: str | Sequence[str] | None = "comparison_below_min_success"
|
||||
branch_labels: str | Sequence[str] | None = None
|
||||
depends_on: str | Sequence[str] | None = None
|
||||
|
||||
_JSON = sa.JSON().with_variant(postgresql.JSONB(), "postgresql")
|
||||
_CONFIG_KEY = "llm_token_price"
|
||||
_MODEL = "deepseek-v4-flash"
|
||||
_OLD_INPUT_PRICE = 3.0
|
||||
_OLD_OUTPUT_PRICE = 15.0
|
||||
_NEW_INPUT_PRICE = 1.0
|
||||
_NEW_OUTPUT_PRICE = 2.0
|
||||
_CORRECTION_MARKER = "deepseek_v4_flash_price"
|
||||
_CONFIG_MARKER_KEY = "migration_deepseek_v4_flash_price"
|
||||
|
||||
|
||||
def _decode_object(value: Any) -> dict[str, Any] | None:
|
||||
if isinstance(value, dict):
|
||||
return value
|
||||
if isinstance(value, str):
|
||||
try:
|
||||
parsed = json.loads(value)
|
||||
except (TypeError, ValueError):
|
||||
return None
|
||||
return parsed if isinstance(parsed, dict) else None
|
||||
return None
|
||||
|
||||
|
||||
def _model_tokens(calls: Any) -> tuple[int, int]:
|
||||
if not isinstance(calls, list):
|
||||
return 0, 0
|
||||
input_tokens = 0
|
||||
output_tokens = 0
|
||||
for call in calls:
|
||||
if not isinstance(call, dict) or call.get("error") or call.get("model") != _MODEL:
|
||||
continue
|
||||
usage = _decode_object(call.get("usage"))
|
||||
if usage is None:
|
||||
continue
|
||||
input_tokens += int(usage.get("prompt_tokens") or 0)
|
||||
output_tokens += int(usage.get("completion_tokens") or 0)
|
||||
return input_tokens, output_tokens
|
||||
|
||||
|
||||
def _app_config_table() -> sa.TableClause:
|
||||
return sa.table(
|
||||
"app_config",
|
||||
sa.column("key", sa.String(64)),
|
||||
sa.column("value", _JSON),
|
||||
sa.column("updated_at", sa.DateTime(timezone=True)),
|
||||
)
|
||||
|
||||
|
||||
def _comparison_table() -> sa.TableClause:
|
||||
return sa.table(
|
||||
"comparison_record",
|
||||
sa.column("id", sa.Integer),
|
||||
sa.column("llm_calls", _JSON),
|
||||
sa.column("llm_cost_yuan", sa.Float),
|
||||
sa.column("llm_price_snapshot", _JSON),
|
||||
)
|
||||
|
||||
|
||||
def _update_config(conn, *, upgrade: bool) -> None:
|
||||
table = _app_config_table()
|
||||
row = conn.execute(
|
||||
sa.select(table.c.value).where(table.c.key == _CONFIG_KEY)
|
||||
).mappings().first()
|
||||
if row is None or not isinstance(row["value"], dict):
|
||||
return
|
||||
|
||||
config = dict(row["value"])
|
||||
per_model = dict(config.get("per_model") or {})
|
||||
current = per_model.get(_MODEL)
|
||||
if upgrade:
|
||||
# Preserve an operator's explicit model price. The production defect is specifically
|
||||
# the missing key falling through to the generic 3/15 price.
|
||||
if current is not None:
|
||||
return
|
||||
per_model[_MODEL] = {
|
||||
"input_per_1m": _NEW_INPUT_PRICE,
|
||||
"output_per_1m": _NEW_OUTPUT_PRICE,
|
||||
}
|
||||
conn.execute(table.insert().values(key=_CONFIG_MARKER_KEY, value=True))
|
||||
else:
|
||||
marker_exists = conn.execute(
|
||||
sa.select(table.c.key).where(table.c.key == _CONFIG_MARKER_KEY)
|
||||
).scalar_one_or_none()
|
||||
if marker_exists is None:
|
||||
return
|
||||
if current != {
|
||||
"input_per_1m": _NEW_INPUT_PRICE,
|
||||
"output_per_1m": _NEW_OUTPUT_PRICE,
|
||||
}:
|
||||
conn.execute(table.delete().where(table.c.key == _CONFIG_MARKER_KEY))
|
||||
return
|
||||
per_model.pop(_MODEL, None)
|
||||
conn.execute(table.delete().where(table.c.key == _CONFIG_MARKER_KEY))
|
||||
config["per_model"] = per_model
|
||||
conn.execute(
|
||||
table.update()
|
||||
.where(table.c.key == _CONFIG_KEY)
|
||||
# 这是对历史误配置的追溯修正,不是从部署时刻开始的新价格。保留原 updated_at,
|
||||
# 否则缺失成本回填会把部署前的记录全部排除。
|
||||
.values(value=config)
|
||||
)
|
||||
|
||||
|
||||
def _correct_frozen_costs(conn, *, upgrade: bool) -> None:
|
||||
table = _comparison_table()
|
||||
rows = conn.execute(
|
||||
sa.select(
|
||||
table.c.id,
|
||||
table.c.llm_calls,
|
||||
table.c.llm_cost_yuan,
|
||||
table.c.llm_price_snapshot,
|
||||
).where(
|
||||
table.c.llm_cost_yuan.is_not(None),
|
||||
table.c.llm_price_snapshot.is_not(None),
|
||||
)
|
||||
).mappings()
|
||||
|
||||
for row in rows:
|
||||
snapshot = row["llm_price_snapshot"]
|
||||
if not isinstance(snapshot, dict):
|
||||
continue
|
||||
prices = snapshot.get("prices")
|
||||
if not isinstance(prices, dict):
|
||||
continue
|
||||
model_price = prices.get(_MODEL)
|
||||
if not isinstance(model_price, dict):
|
||||
continue
|
||||
|
||||
if upgrade:
|
||||
if not (
|
||||
model_price.get("_source") == "default"
|
||||
and model_price.get("input_per_1m") == _OLD_INPUT_PRICE
|
||||
and model_price.get("output_per_1m") == _OLD_OUTPUT_PRICE
|
||||
):
|
||||
continue
|
||||
elif snapshot.get("pricing_correction") != _CORRECTION_MARKER:
|
||||
continue
|
||||
|
||||
input_tokens, output_tokens = _model_tokens(row["llm_calls"])
|
||||
if input_tokens == 0 and output_tokens == 0:
|
||||
continue
|
||||
if upgrade:
|
||||
delta = (
|
||||
input_tokens / 1_000_000 * (_OLD_INPUT_PRICE - _NEW_INPUT_PRICE)
|
||||
+ output_tokens / 1_000_000 * (_OLD_OUTPUT_PRICE - _NEW_OUTPUT_PRICE)
|
||||
)
|
||||
corrected_price = {
|
||||
"input_per_1m": _NEW_INPUT_PRICE,
|
||||
"output_per_1m": _NEW_OUTPUT_PRICE,
|
||||
"_source": "per_model",
|
||||
}
|
||||
snapshot["pricing_correction"] = _CORRECTION_MARKER
|
||||
new_cost = max(0.0, float(row["llm_cost_yuan"]) - delta)
|
||||
else:
|
||||
delta = (
|
||||
input_tokens / 1_000_000 * (_OLD_INPUT_PRICE - _NEW_INPUT_PRICE)
|
||||
+ output_tokens / 1_000_000 * (_OLD_OUTPUT_PRICE - _NEW_OUTPUT_PRICE)
|
||||
)
|
||||
corrected_price = {
|
||||
"input_per_1m": _OLD_INPUT_PRICE,
|
||||
"output_per_1m": _OLD_OUTPUT_PRICE,
|
||||
"_source": "default",
|
||||
}
|
||||
snapshot.pop("pricing_correction", None)
|
||||
new_cost = float(row["llm_cost_yuan"]) + delta
|
||||
|
||||
updated_prices = dict(prices)
|
||||
updated_prices[_MODEL] = corrected_price
|
||||
updated_snapshot = dict(snapshot)
|
||||
updated_snapshot["prices"] = updated_prices
|
||||
conn.execute(
|
||||
table.update()
|
||||
.where(table.c.id == row["id"])
|
||||
.values(
|
||||
llm_cost_yuan=round(new_cost, 6),
|
||||
llm_price_snapshot=updated_snapshot,
|
||||
)
|
||||
)
|
||||
|
||||
|
||||
def upgrade() -> None:
|
||||
conn = op.get_bind()
|
||||
_update_config(conn, upgrade=True)
|
||||
_correct_frozen_costs(conn, upgrade=True)
|
||||
|
||||
|
||||
def downgrade() -> None:
|
||||
conn = op.get_bind()
|
||||
_correct_frozen_costs(conn, upgrade=False)
|
||||
_update_config(conn, upgrade=False)
|
||||
@@ -166,18 +166,6 @@ def _session_to_row(
|
||||
point_stats: dict | None = None,
|
||||
) -> dict:
|
||||
"""CouponSession ORM → 明细行 dict(主表「领券数据」与「用户全部领券」抽屉共用)。"""
|
||||
# 中途退出可能发生在第一张券产生终态之前,此时没有逐券事件。
|
||||
# 明确返回 0/0,让前端区分「退出前无单券结果」与其它状态的埋点缺失。
|
||||
if point_stats is not None:
|
||||
point_success_count = point_stats["succeeded"]
|
||||
point_total_count = point_stats["tried"]
|
||||
elif r.status == "abandoned":
|
||||
point_success_count = 0
|
||||
point_total_count = 0
|
||||
else:
|
||||
point_success_count = None
|
||||
point_total_count = None
|
||||
point_event_count = point_stats["events"] if point_stats is not None else 0
|
||||
return {
|
||||
"id": r.id,
|
||||
"trace_id": r.trace_id,
|
||||
@@ -194,9 +182,8 @@ def _session_to_row(
|
||||
"app_env": r.app_env,
|
||||
"started_at": r.started_at,
|
||||
"claimed_count": r.claimed_count,
|
||||
"point_success_count": point_success_count,
|
||||
"point_total_count": point_total_count,
|
||||
"point_event_count": point_event_count,
|
||||
"point_success_count": point_stats["succeeded"] if point_stats else None,
|
||||
"point_total_count": point_stats["tried"] if point_stats else None,
|
||||
"trace_url": r.trace_url,
|
||||
"ad_revenue_yuan": ad_revenue_yuan,
|
||||
}
|
||||
@@ -207,24 +194,21 @@ def _point_scores_by_trace(db: Session, trace_ids: list[str]) -> dict[str, dict[
|
||||
if not trace_ids:
|
||||
return {}
|
||||
succeeded = func.sum(case((CouponClaimEvent.status.in_(_SLOT_OK), 1), else_=0))
|
||||
tried = func.sum(case((CouponClaimEvent.status.in_(_SLOT_TRIED), 1), else_=0))
|
||||
rows = db.execute(
|
||||
select(
|
||||
CouponClaimEvent.trace_id,
|
||||
succeeded.label("succeeded"),
|
||||
tried.label("tried"),
|
||||
func.count().label("events"),
|
||||
func.count().label("tried"),
|
||||
)
|
||||
.where(
|
||||
CouponClaimEvent.trace_id.in_(trace_ids),
|
||||
CouponClaimEvent.status.in_(_SLOT_TRIED),
|
||||
)
|
||||
.where(CouponClaimEvent.trace_id.in_(trace_ids))
|
||||
.group_by(CouponClaimEvent.trace_id)
|
||||
).all()
|
||||
return {
|
||||
trace_id: {
|
||||
"succeeded": int(success_count or 0),
|
||||
"tried": int(tried_count or 0),
|
||||
"events": int(event_count or 0),
|
||||
}
|
||||
for trace_id, success_count, tried_count, event_count in rows
|
||||
trace_id: {"succeeded": int(success_count or 0), "tried": int(tried or 0)}
|
||||
for trace_id, success_count, tried in rows
|
||||
if trace_id is not None
|
||||
}
|
||||
|
||||
@@ -416,15 +400,12 @@ def coupon_user_records(db: Session, *, user_id: int, limit: int = 100) -> dict:
|
||||
total = db.execute(
|
||||
select(func.count()).select_from(CouponSession).where(CouponSession.user_id == user_id)
|
||||
).scalar_one()
|
||||
trace_ids = [r.trace_id for r in rows]
|
||||
rev_map = crud_ecpm.revenue_yuan_by_trace(db, trace_ids)
|
||||
point_stats_map = _point_scores_by_trace(db, trace_ids)
|
||||
rev_map = crud_ecpm.revenue_yuan_by_trace(db, [r.trace_id for r in rows])
|
||||
return {
|
||||
"items": [
|
||||
_session_to_row(
|
||||
r,
|
||||
ad_revenue_yuan=rev_map.get(r.trace_id, 0.0),
|
||||
point_stats=point_stats_map.get(r.trace_id),
|
||||
)
|
||||
for r in rows
|
||||
],
|
||||
|
||||
@@ -566,10 +566,6 @@ def dashboard_overview(
|
||||
)
|
||||
).all()
|
||||
coupon_started = len(period_coupon_sessions)
|
||||
coupon_abandoned = sum(s.status == "abandoned" for s in period_coupon_sessions)
|
||||
# 用户主动中途退出不代表领券流程失败,不进入整场成功率样本。
|
||||
# started / failed 仍留在分母:前者是尚未形成终态的流失,后者是实际执行失败。
|
||||
coupon_success_denominator = coupon_started - coupon_abandoned
|
||||
coupon_completed_elapsed = sorted(
|
||||
s.elapsed_ms
|
||||
for s in period_coupon_sessions
|
||||
@@ -735,13 +731,9 @@ def dashboard_overview(
|
||||
},
|
||||
"coupon": {
|
||||
"started": coupon_started,
|
||||
"abandoned": coupon_abandoned,
|
||||
"success_denominator": coupon_success_denominator,
|
||||
"all_success": coupon_all_success,
|
||||
"success_rate": (
|
||||
round(coupon_all_success / coupon_success_denominator, 4)
|
||||
if coupon_success_denominator
|
||||
else None
|
||||
round(coupon_all_success / coupon_started, 4) if coupon_started else None
|
||||
),
|
||||
"point_success": coupon_point_success,
|
||||
"points_per_session": coupon_points_per_session,
|
||||
|
||||
@@ -79,16 +79,10 @@ class CouponDataRow(BaseModel):
|
||||
started_at: datetime = Field(..., description="发起时刻(明细「时间」列)")
|
||||
claimed_count: int | None = None
|
||||
point_success_count: int | None = Field(
|
||||
None,
|
||||
description="本次成功单券数(success+already_claimed);中途退出且无逐券结果为0,其它无事件为空",
|
||||
None, description="本次成功单券数(success+already_claimed);无逐券事件为空"
|
||||
)
|
||||
point_total_count: int | None = Field(
|
||||
None,
|
||||
description="本次尝试单券数(success+already_claimed+failed,不含 skipped);中途退出且无逐券结果为0,其它无事件为空",
|
||||
)
|
||||
point_event_count: int = Field(
|
||||
0,
|
||||
description="本次全部逐券事件数(含 skipped);用于区分无有效计分事件与完全无事件",
|
||||
None, description="本次尝试单券数(success+already_claimed+failed,不含 skipped);无逐券事件为空"
|
||||
)
|
||||
trace_url: str | None = Field(None, description="pricebot 公网 trace 链接(仅 completed 有);admin 渲染可点链接,无则显示可复制 trace_id")
|
||||
ad_revenue_yuan: float = Field(
|
||||
|
||||
@@ -70,9 +70,6 @@ class DashboardPeriodCoupon(BaseModel):
|
||||
成功口径 success+already_claimed(与「我的」页累计领券一致)。"""
|
||||
|
||||
started: int = 0
|
||||
# 用户主动中途退出,不计入整场成功率分母。
|
||||
abandoned: int = 0
|
||||
success_denominator: int = 0
|
||||
# 全部领成功的次数:completed 且当日该设备全部点位成功
|
||||
all_success: int = 0
|
||||
success_rate: float | None = None
|
||||
|
||||
@@ -227,11 +227,7 @@ CONFIG_DEFS: dict[str, dict[str, Any]] = {
|
||||
# 编辑框;set_value 不校验类型,嵌套 JSON 照存。
|
||||
"llm_token_price": {
|
||||
"default": {
|
||||
"per_model": {
|
||||
"qwen3.5-flash": {"input_per_1m": 0.8, "output_per_1m": 2.0},
|
||||
# DashScope 华北 2 公网调用原价;必须显式配置,不能落到 3/15 的未知模型兜底价。
|
||||
"deepseek-v4-flash": {"input_per_1m": 1.0, "output_per_1m": 2.0},
|
||||
},
|
||||
"per_model": {"qwen3.5-flash": {"input_per_1m": 0.8, "output_per_1m": 2.0}},
|
||||
"default": {"input_per_1m": 3.0, "output_per_1m": 15.0},
|
||||
"currency": "CNY", "unit": "per_1m_tokens",
|
||||
},
|
||||
|
||||
@@ -1,304 +0,0 @@
|
||||
# 线上 PostgreSQL 数据库备份 / 恢复方案 设计
|
||||
|
||||
- 日期:2026-08-03
|
||||
- 状态:设计已评审,待写实现计划
|
||||
- 相关代码(复用的现成范式):
|
||||
- [deploy/daily-exchange.service](../../../deploy/daily-exchange.service) / [deploy/daily-exchange.timer](../../../deploy/daily-exchange.timer)(systemd oneshot + timer 定时任务范式:`TZ=Asia/Shanghai`、`Persistent=true` 补跑、`ProtectSystem=strict` 加固、文件锁)
|
||||
- [app/core/heartbeat_monitor_worker.py](../../../app/core/heartbeat_monitor_worker.py)(单实例文件锁 + 优雅退出范式)
|
||||
- [app/integrations/notifier.py](../../../app/integrations/notifier.py)(可插拔通知器:`LogNotifier` 占位、协议不变后续替换)
|
||||
- [app/admin/deps.py](../../../app/admin/deps.py)(`require_role("super_admin")` 权限守卫、审计 IP)
|
||||
- [app/admin/routers/withdraw.py](../../../app/admin/routers/withdraw.py)(admin router 风格范式)
|
||||
- [scripts/init_postgres.py](../../../scripts/init_postgres.py)(PG 运维脚本范式)
|
||||
|
||||
## 1. 背景与目标
|
||||
|
||||
线上业务库 `shaguabijia`(PostgreSQL 16,当前数据量 1–10GB)目前没有系统化的备份/恢复方案。本方案提供覆盖「定时 + 手动 + 平台恢复 + 平台不可用兜底」四条路径的逻辑备份体系。
|
||||
|
||||
四条明确需求:
|
||||
|
||||
1. **每日定时备份** —— 无人值守,不依赖平台进程。
|
||||
2. **平台手动备份** —— admin 后台一键触发。
|
||||
3. **平台指定备份文件快速恢复** —— 在 admin 后台选一份备份,恢复到**旁库**供核对。
|
||||
4. **平台不可用时手动恢复** —— 登录服务器,交互式选择备份恢复,与平台走**同一套脚本**。
|
||||
|
||||
数据规模 1–10GB → **逻辑备份(`pg_dump`)足够**,不引入物理备份 / PITR 的复杂度(见 §2 预留路径)。
|
||||
|
||||
## 2. 非目标(本期不做 / 预留升级路径)
|
||||
|
||||
- **物理备份 + PITR**(`pg_basebackup` + WAL 归档):数据量涨到「凌晨全量 dump 也影响业务」时再上。§3 的架构不阻碍后续叠加。
|
||||
- **异地对象存储(S3 / 阿里云 OSS)**:本期只落**本地磁盘**。备份脚本预留 `upload_to_remote()` 可插拔 hook(`BACKUP_S3_ENABLED=false` 时直接返回),等 S3 服务确认后填充,主流程不改。
|
||||
- **从只读副本 dump**:搭流复制 standby、改从副本备份以对主库零影响,属后续升级;那个副本还能顺带承载 PITR。
|
||||
- **平台「一键覆盖生产」按钮**:刻意不做。扶正生产(旁库→主库)永远是人工 SSH 执行的高危脚本(§6)。
|
||||
- **平台下载备份文件**:整库备份含手机号 / 微信 / 提现等敏感数据,经浏览器下载放大泄露面且文件大。取文件走 SSH。
|
||||
- **真实推送告警渠道**:项目当前无真实推送能力(`notifier.py` 是 `LogNotifier` 占位、心跳 worker 也只打印)。本期告警走「日志 + openobserve + 平台新鲜度徽标」,推送做占位 hook(§8)。
|
||||
|
||||
## 3. 核心设计原则与整体架构
|
||||
|
||||
### 3.1 两条铁律
|
||||
|
||||
1. **执行逻辑沉到自包含脚本,平台与定时器都只是调用者。** —— 需求 4(平台不可用时手动恢复)因此不是另写一套,而是复用同一套脚本,逻辑不重复、兜底路径永远可用。
|
||||
2. **备份清单的权威来源是文件系统(备份目录 + 每目录内的 `manifest.json`),不是数据库任务表。** —— 任务表只记录「平台发起的操作过程」用于展示进度与审计;删掉任务表不影响任何一份备份的可用性与可恢复性。最需要恢复的时刻(DB 崩了)恰恰是任务表也读不到的时刻,而恢复根本不读任务表。
|
||||
|
||||
### 3.2 分层架构
|
||||
|
||||
```
|
||||
┌─────────────────────────────────────────┐
|
||||
│ 核心执行层(自包含,不依赖平台/DB 任务表) │
|
||||
│ scripts/db_backup.sh ← 生成备份 │
|
||||
│ scripts/db_restore.sh ← 从备份恢复到旁库 │
|
||||
│ scripts/db_promote.sh ← 旁库扶正生产(高危)│
|
||||
└─────────────────────────────────────────┘
|
||||
▲ ▲ ▲
|
||||
┌─────────────┘ │ └──────────────┐
|
||||
① 每日定时 ② 平台手动备份/恢复 ④ 平台挂了,人 SSH
|
||||
systemd timer admin API → setsid 起脚本 跑 db_restore.sh
|
||||
(不经过平台) (任务表记状态,web 不阻塞) (交互式选备份)
|
||||
```
|
||||
|
||||
### 3.3 平台执行方式:方案 B(轻量版)
|
||||
|
||||
平台触发的「几分钟重操作」不放进 web 请求,采用**任务表解耦 + 独立进程执行**:
|
||||
|
||||
| | 采用 | 说明 |
|
||||
|---|---|---|
|
||||
| 谁拉起脚本 | admin API 插任务后 `setsid` 起一个**脱离 web 进程组**的 subprocess 跑脚本 | 不新增常驻 worker 服务,最贴合项目「脚本 + systemd」风格;web 重启不影响已起的脚本 |
|
||||
| 并发控制 | `db_backup.sh` 的 `flock` 文件锁天然串行 | 防定时与手动同时跑、上一轮未完下一轮又起 |
|
||||
| 升级路径 | 任务量大、需排队/限流/重试时,升级为常驻 worker 轮询任务表 | 本期低频操作用不上 |
|
||||
|
||||
**被否决的方案 A**:admin API 里同步 `subprocess` 跑 dump/restore,请求一直挂着。否决原因:几分钟操作阻塞 web 进程、易请求超时;admin 进程需持有 DB 高权限;平台一挂手动备份路径也没了(解耦度反而更差)。
|
||||
|
||||
## 4. 备份设计
|
||||
|
||||
### 4.1 备份产物:一个自描述目录
|
||||
|
||||
```
|
||||
$BACKUP_DIR/20260803_030000_daily/
|
||||
├── shaguabijia.dump # pg_dump -Fc(自定义格式,内建压缩)
|
||||
├── globals.sql # pg_dumpall --globals-only --no-role-passwords
|
||||
├── manifest.json # 元信息(见下)
|
||||
└── SHA256SUMS # 上述两个文件的校验和
|
||||
```
|
||||
|
||||
目录名(= `backup_id`)格式:`YYYYMMDD_HHMMSS_<reason>`(北京时 `Asia/Shanghai`,纯数字 + 下划线),如 `20260803_030000_daily`,`reason ∈ {daily, manual}`。刻意不用连字符/冒号/字母,使其**同时是合法文件名与合法 PG 库名**——旁库名由它直接拼接(见 §5.1)。
|
||||
|
||||
`manifest.json` 字段:
|
||||
|
||||
| 字段 | 说明 |
|
||||
|---|---|
|
||||
| `backup_id` | = 目录名,全局唯一标识 |
|
||||
| `created_at` | ISO8601 北京时 |
|
||||
| `reason` | `daily` / `manual` |
|
||||
| `triggered_by` | `systemd-timer` / admin 用户名 / `ssh-manual` |
|
||||
| `pg_version` | 备份时 PG 版本 |
|
||||
| `database` | `shaguabijia` |
|
||||
| `format` | `custom` |
|
||||
| `alembic_version` | 备份时的 schema 版本(恢复时对齐迁移用;旁库 `alembic_version` 表亦可查) |
|
||||
| `files` | 每个文件的 `name` / `size_bytes` / `sha256` |
|
||||
| `dump_duration_sec` | 导出耗时 |
|
||||
| `status` | `success` / `failed`(**仅自检通过才写 success**) |
|
||||
|
||||
### 4.2 `db_backup.sh` 关键行为
|
||||
|
||||
| 环节 | 做法 | 理由 |
|
||||
|---|---|---|
|
||||
| 导出业务库 | `pg_dump -Fc` → `shaguabijia.dump` | 单文件、压缩,`pg_restore` 支持并行(`-j`)与按表选择性恢复 |
|
||||
| 导出全局角色 | `pg_dumpall --globals-only --no-role-passwords` → `globals.sql` | 恢复到新机器时角色/权限齐全;`--no-role-passwords` 免超级用户读 `pg_authid`、旁库核对也不需要密码 |
|
||||
| 资源降级 | `nice -n 19 ionice -c3` 包裹 pg_dump | 让业务优先,压制备份对 CPU/IO 的抢占 |
|
||||
| **备份后自检** | dump 完立刻 `pg_restore --list *.dump >/dev/null` 解析 TOC + 校验 sha256 | 不验证的备份是薛定谔的备份;能抓到文件截断/损坏,只有通过才写 `status=success` |
|
||||
| 并发锁 | `flock` 独占锁(仿 daily-exchange 30min 锁) | 串行化,防重入 |
|
||||
| 保留清理 | **仅自检通过后**执行(§4.4) | 保证先有新备份、再删旧的 |
|
||||
| 远程上传 | 末尾 `upload_to_remote()`,`BACKUP_S3_ENABLED=false` 时 return 0 | S3 可插拔扩展点 |
|
||||
| 通知 | 末尾 `notify()`(§8),成功/失败都打结构化日志 | 贴合现状、可插拔 |
|
||||
| 退出码 | 失败非 0 退出 + `SyslogIdentifier=pg-backup` | systemd 与 openobserve 可感知 |
|
||||
| 触发标签 | `--reason daily|manual`,写进 manifest | 定时与手动共用一个脚本 |
|
||||
|
||||
### 4.3 对线上读写的影响与缓解
|
||||
|
||||
`pg_dump` 基于 MVCC 一致性快照,**不阻塞正常增删改查**,导出的是某一时刻的一致性视图。要点:
|
||||
|
||||
- **唯一会互斥的例外是 DDL**:dump 持 `ACCESS SHARE` 锁,与 `ALTER TABLE` / `DROP` / `TRUNCATE` / `VACUUM FULL` / 非并发建索引冲突。规避:alembic 迁移仅在手动部署时跑,与凌晨备份天然错开。
|
||||
- **真正的影响是资源争抢**(1–10GB 下的关注点):磁盘 I/O(整库顺序读)、CPU(zlib 压缩)、缓存冲刷(热数据被挤出 shared_buffers,备份后短时命中率下降)、长事务期间 VACUUM 暂时回收不了死元组。
|
||||
- **缓解**:定时放凌晨低谷(03:00,与 0 点 daily-exchange 错开)+ `nice`/`ionice` 降级 + 控制 dump 时长在几分钟内。
|
||||
|
||||
### 4.4 保留策略
|
||||
|
||||
清理只在**备份成功且自检通过后**执行,且**只删本地、不碰 S3**(S3 用自身生命周期策略):
|
||||
|
||||
- `BACKUP_KEEP_DAILY_DAYS=14` —— 保留最近 14 天每日备份。
|
||||
- `BACKUP_KEEP_MONTHLY_COUNT=6` —— 更早的备份中,每月 1 号那份额外保留 6 个月(防「问题两周前就埋下」)。
|
||||
- `BACKUP_MIN_KEEP=3` —— **硬底线:无论配置如何,永远至少保留最近 3 份**,防时间跳变/配置错误把备份删光。
|
||||
- 清理前检查磁盘剩余空间,不足则打告警日志而非静默继续。
|
||||
|
||||
## 5. 恢复设计(恢复到旁库)
|
||||
|
||||
**铁律:平台与默认脚本永远只恢复到旁库,绝不自动碰生产。**
|
||||
|
||||
### 5.1 `db_restore.sh` —— 恢复到旁库(常规路径)
|
||||
|
||||
输入一个 `backup_id`,产出可供核对的旁库 `shaguabijia_restore_<backup_id>`(如 `shaguabijia_restore_20260803_030000_daily`):
|
||||
|
||||
| 步 | 动作 | 护栏 |
|
||||
|---|---|---|
|
||||
| 1 | **恢复前先验完整性**:校验 `SHA256SUMS` + `pg_restore --list` 确认可解析 | 不把损坏备份恢复到一半才发现 |
|
||||
| 2 | 创建旁库 `shaguabijia_restore_<backup_id>`;已存在则提示换名或显式 `--force` 重建 | 不撞库、不误删 |
|
||||
| 3 | 灌 `globals.sql`(角色已存在则跳过,幂等) | |
|
||||
| 4 | `pg_restore -j <并行度,默认 2>` 到旁库(custom 格式自动识别,无需 `-Fc`),`nice`/`ionice` 降级 | 往旁库写,**不锁生产表**,对生产仅轻微资源争抢 |
|
||||
| 5 | **恢复后自检报告**:表数量、关键业务表(user / wallet / withdraw 等)行数、`alembic_version`、最新记录时间戳 | 人凭报告判断「数据对不对、新不新」 |
|
||||
|
||||
**磁盘峰值提醒**:旁库与生产库共存于同一实例,恢复期间磁盘占用约为 2×(生产 + 旁库),恢复前脚本检查磁盘余量。
|
||||
|
||||
### 5.2 兜底:交互式选择(需求 4)
|
||||
|
||||
`db_restore.sh` **不带参数**运行时 → 扫描 `$BACKUP_DIR`,列表打印所有备份(`backup_id` / 时间 / 大小 / `status`),提示输入序号选一个恢复到旁库。这就是「平台挂了 SSH 登录手动选备份恢复」——与平台调用的是同一个脚本,零额外学习成本。
|
||||
|
||||
### 5.3 深度核对的边界
|
||||
|
||||
平台展示 §5.1 第 5 步的自检报告即可支撑「数据对不对」的判断。**逐行深度核对不在本方案范围**——需要时 DBA 直接连旁库跑只读 SQL,旁库为此存在。
|
||||
|
||||
## 6. 扶正生产:`db_promote.sh`(高危、人工、不进平台)
|
||||
|
||||
核对无误后把旁库扶正为生产。**故意不做成平台按钮**,必须人工 SSH 执行,因为不可逆且涉及短暂停机:
|
||||
|
||||
1. **二次确认**:要求手输生产库名 `shaguabijia` 才继续(仿 GitHub 删仓库确认)。
|
||||
2. **先自动备份当前生产库**(调 `db_backup.sh --reason manual`)—— 最关键的兜底,扶错了能退回来。
|
||||
3. **断开生产连接**:停 app 服务或 `pg_terminate_backend` 清连接(执行者自身不能连在待改名的库上)。**此步有短暂停机。**
|
||||
4. **rename 切换**(不用 drop):`shaguabijia` → `shaguabijia_old_<ts>`,`shaguabijia_restore_xxx` → `shaguabijia`。出错可立即换回,比删库重建安全得多。
|
||||
5. 重启 app 服务,人工验证。
|
||||
|
||||
**为何用脚本而非裸敲 SQL**:切换不是「只有 rename 一句」,而是一串不能漏的动作,其中「先备份当前生产库」裸敲最易漏、漏了就没退路。脚本把这串封装成带护栏的原子操作。极端情况下连脚本都不可用时,rename 本质是两条 `ALTER DATABASE` SQL,DBA 亦可手工执行,但会失去自动兜底备份保护,属下策。
|
||||
|
||||
## 7. 平台 API + 数据模型 + 权限
|
||||
|
||||
### 7.1 数据模型(两张任务表,均新建 + alembic 迁移)
|
||||
|
||||
`db_backup_job`:
|
||||
|
||||
| 字段 | 类型 | 说明 |
|
||||
|---|---|---|
|
||||
| `id` | PK | |
|
||||
| `triggered_by_admin_id` | FK admin_user, nullable | 手动才有 |
|
||||
| `status` | str | `pending` / `running` / `success` / `failed` |
|
||||
| `backup_id` | str, nullable | 成功后填产出目录名 |
|
||||
| `error_msg` | text, nullable | |
|
||||
| `created_at` / `started_at` / `finished_at` | datetime | 后两者 nullable |
|
||||
| `duration_sec` | int, nullable | |
|
||||
|
||||
`db_restore_job`:
|
||||
|
||||
| 字段 | 类型 | 说明 |
|
||||
|---|---|---|
|
||||
| `id` | PK | |
|
||||
| `backup_id` | str | 恢复哪个备份 |
|
||||
| `target_db` | str | 旁库名 |
|
||||
| `triggered_by_admin_id` | FK admin_user | |
|
||||
| `status` | str | `pending` / `running` / `success` / `failed` |
|
||||
| `sanity_report` | json, nullable | §5.1 第 5 步的自检报告 |
|
||||
| `error_msg` | text, nullable | |
|
||||
| `created_at` / `started_at` / `finished_at` | datetime | 后两者 nullable |
|
||||
|
||||
模型需在 [app/models/__init__.py](../../../app/models/__init__.py) 导入以便 Alembic 发现。
|
||||
|
||||
**备份清单不查任务表**:定时备份不写任务表(不经过平台),但其产出的备份**出现在备份列表**(列表来自扫盘)。平台因此有两个视图:备份文件列表(扫盘,含所有来源)、操作记录(任务表 + 审计)。
|
||||
|
||||
### 7.2 API 端点(新建 `app/admin/routers/db_backup.py`,注册进 `app/admin/main.py`)
|
||||
|
||||
| 方法 | 路径 | 作用 | 返回 |
|
||||
|---|---|---|---|
|
||||
| GET | `/admin/db-backup/backups` | 列所有备份(扫盘读 manifest) | `list[BackupItem]` |
|
||||
| GET | `/admin/db-backup/health` | 新鲜度:最近成功备份时间与年龄 | `{last_success_at, age_hours, healthy}` |
|
||||
| POST | `/admin/db-backup/backups` | 手动触发备份 | `202 {job_id}` |
|
||||
| GET | `/admin/db-backup/jobs/{id}` | 轮询备份任务状态 | `BackupJobStatus` |
|
||||
| POST | `/admin/db-backup/restores` | 恢复 `backup_id` 到旁库(body: `{backup_id}`) | `202 {restore_job_id}` |
|
||||
| GET | `/admin/db-backup/restores/{id}` | 轮询恢复状态 + 自检报告 | `RestoreJobStatus` |
|
||||
|
||||
`BackupItem`:`backup_id / created_at / reason / triggered_by / size_bytes / database / pg_version / alembic_version / status`。
|
||||
|
||||
契约放 `app/schemas/db_backup.py`(Pydantic)。POST 端点为**异步**语义:插任务表 → `setsid` 起脚本 → 立即 `202` 返回 `job_id`,前端轮询对应 GET 端点。
|
||||
|
||||
### 7.3 权限与审计
|
||||
|
||||
- **权限**:备份与恢复端点**均限 `super_admin`**(复用 `require_role("super_admin")`)。恢复能触及全库数据,不宜下放;将来给运维岗再引入 `require_page("db_backup")` 细分。
|
||||
- **审计**:每次备份 / 恢复写现有审计日志(触发人、`backup_id`、结果),与项目其它高危操作一致。
|
||||
- **不提供下载**(见 §2)。
|
||||
|
||||
## 8. 监控告警(贴合「日志 + openobserve + 平台徽标」现状)
|
||||
|
||||
三层,从被动到主动:
|
||||
|
||||
1. **结构化日志**:脚本成功/失败均打 `SyslogIdentifier=pg-backup` 日志,`journalctl` 可查、openobserve 可抓,失败为 ERROR 级 → 若已配 openobserve 告警规则即命中。
|
||||
2. **平台新鲜度徽标**(最直观、不依赖推送):admin 备份页顶部显示「最近成功备份:X 小时前」,超 `BACKUP_FRESH_MAX_HOURS`(默认 26h)红色高亮。数据源为 `GET /admin/db-backup/health`。
|
||||
3. **可插拔 `notify()` hook**:脚本内通知点,现在只打日志(照搬 `notifier.py` 的 `LogNotifier` 模式),将来推送能力(飞书/短信)就绪再填,主流程不改。
|
||||
|
||||
**新鲜度自检兜底**:`deploy/pg-backup-check.timer`(+service)每天 09:00 检查最新成功备份是否在 26h 内,过期打 ERROR 日志 + 调 `notify()`。这覆盖 `OnFailure` 抓不到的盲区(timer 被禁 / 宕机没补跑)——因为那种情况脚本根本没运行,靠「有没有新备份」反向判断。
|
||||
|
||||
## 9. 脚本连库身份与权限
|
||||
|
||||
备份/恢复脚本**以本机 `postgres` 超级用户走 Unix socket(peer 认证,无需密码)**执行,而非业务的 `DATABASE_URL`(TCP + 密码)。理由:
|
||||
|
||||
- `pg_dumpall --globals-only`、`CREATE DATABASE`(建旁库)、`ALTER DATABASE ... RENAME`、`pg_terminate_backend` 都需高权限;
|
||||
- 本机 socket peer 认证是运维脚本标准做法,权限最省心,且不用把超级用户密码写进任何配置文件。
|
||||
|
||||
脚本以 `root` 或 `postgres` 系统用户运行(systemd service 内 `User=postgres` 或经 `sudo -u postgres`)。SSH 手动执行时同理。
|
||||
|
||||
## 10. 配置项
|
||||
|
||||
新增到 [app/core/config.py](../../../app/core/config.py) 的 `Settings`(供平台侧 §8 新鲜度等读取),shell 脚本经 systemd `EnvironmentFile=.env` 注入、SSH 手动跑时脚本内有默认值兜底:
|
||||
|
||||
| 配置 | 默认 | 说明 |
|
||||
|---|---|---|
|
||||
| `BACKUP_DIR` | `/opt/pg_backups` | 备份根目录(建议独立数据盘,与 PG 数据文件不同物理盘) |
|
||||
| `BACKUP_KEEP_DAILY_DAYS` | `14` | 每日备份保留天数 |
|
||||
| `BACKUP_KEEP_MONTHLY_COUNT` | `6` | 月度长留份数 |
|
||||
| `BACKUP_MIN_KEEP` | `3` | 硬底线,永远至少保留份数 |
|
||||
| `BACKUP_FRESH_MAX_HOURS` | `26` | 新鲜度阈值 |
|
||||
| `BACKUP_PG_SUPERUSER` | `postgres` | 脚本连库超级用户 |
|
||||
| `BACKUP_PGHOST` | `/var/run/postgresql` | socket 目录(peer 认证) |
|
||||
| `BACKUP_S3_ENABLED` | `false` | 远程上传总开关(预留) |
|
||||
| `BACKUP_S3_*` | 空 | bucket / endpoint / 凭证(预留,S3 确认后填) |
|
||||
|
||||
同步更新 `.env.example`。
|
||||
|
||||
## 11. 完整文件清单
|
||||
|
||||
### 本仓(后端 + 脚本)
|
||||
|
||||
**新增:**
|
||||
|
||||
| 文件 | 作用 |
|
||||
|---|---|
|
||||
| `scripts/db_backup.sh` | 备份核心(dump + 自检 + 保留清理 + upload hook + notify) |
|
||||
| `scripts/db_restore.sh` | 恢复到旁库(含无参交互式选择 = 需求 4) |
|
||||
| `scripts/db_promote.sh` | 旁库扶正生产(高危、人工、二次确认) |
|
||||
| `app/models/db_backup_job.py` | `db_backup_job` + `db_restore_job` 模型 |
|
||||
| `app/admin/repositories/db_backup_job.py` | 任务表数据访问 + 扫盘读 manifest |
|
||||
| `app/admin/routers/db_backup.py` | 平台 API(6 个端点) |
|
||||
| `app/schemas/db_backup.py` | API 契约 |
|
||||
| `alembic/versions/xxxx_add_db_backup_jobs.py` | 建两张任务表 |
|
||||
| `deploy/pg-backup.service` / `pg-backup.timer` | 每日定时备份 |
|
||||
| `deploy/pg-backup-check.service` / `pg-backup-check.timer` | 新鲜度自检 |
|
||||
| `deploy/pg-backup.md` | 部署文档(仿 [deploy/daily-exchange.md](../../../deploy/daily-exchange.md)) |
|
||||
|
||||
**改动:** `app/models/__init__.py`(导模型)、`app/core/config.py`(§10 配置项)、`app/admin/main.py`(注册 router)、`.env.example`。
|
||||
|
||||
### 另一仓(`shaguabijia-admin-web`)配套
|
||||
|
||||
备份管理前端页:备份列表 + 一键备份 + 选备份恢复到旁库 + 恢复进度/自检报告展示 + 新鲜度徽标。本 spec 定义 API 契约(§7.2),前端据此实现,单独走该仓的开发流程。
|
||||
|
||||
## 12. 需求覆盖对照
|
||||
|
||||
| 需求 | 落地 |
|
||||
|---|---|
|
||||
| 每日定时备份 | `pg-backup.timer` → `db_backup.sh --reason daily`,凌晨 3 点,不依赖平台 |
|
||||
| 平台手动备份 | `POST /admin/db-backup/backups` → 插任务 + `setsid` 起 `db_backup.sh --reason manual` |
|
||||
| 平台指定备份快速恢复 | `POST /admin/db-backup/restores` → `db_restore.sh` 恢复到旁库 + 自检报告 |
|
||||
| 平台不可用手动恢复 | SSH 跑 `db_restore.sh`(交互式选备份),同一套脚本 |
|
||||
| (延伸)扶正生产 | 人工 `db_promote.sh`,二次确认 + 先备份当前 + rename 切换 |
|
||||
|
||||
## 13. 部署时需确认的运维参数
|
||||
|
||||
以下取决于服务器实际情况,部署或本 spec 复审时确认,不阻碍设计:
|
||||
|
||||
1. **备份目录位置**:默认 `/opt/pg_backups`;若有独立数据盘,建议放数据盘且与 PG 数据文件不同物理盘(避免一盘挂掉数据与备份同亡)。
|
||||
2. **保留量匹配磁盘余量**:1–10GB 压缩后每份约几百 MB~2GB,14 天约 3–30GB;核对与磁盘余量是否匹配。
|
||||
3. **备份时间**:默认凌晨 03:00(与 0 点 daily-exchange 错开);确认无其它凌晨任务撞车。
|
||||
4. **脚本运行用户**:确认以 `postgres`(或可 `sudo -u postgres`)运行、socket peer 认证可用。
|
||||
@@ -1,7 +1,7 @@
|
||||
"""Admin M2 读接口测试:大盘聚合 + 用户/流水/提现/反馈列表 + 鉴权拦截。"""
|
||||
from __future__ import annotations
|
||||
|
||||
from datetime import UTC, date, datetime
|
||||
from datetime import UTC, datetime
|
||||
|
||||
import pytest
|
||||
from fastapi.testclient import TestClient
|
||||
@@ -12,7 +12,6 @@ from app.admin.repositories import admin_user as admin_repo
|
||||
from app.admin.repositories import queries
|
||||
from app.db.session import SessionLocal, engine
|
||||
from app.models.comparison import ComparisonRecord
|
||||
from app.models.coupon_state import CouponClaimRecord, CouponSession
|
||||
from app.models.feedback import Feedback
|
||||
from app.models.invite import InviteRelation
|
||||
from app.models.savings import SavingsRecord
|
||||
@@ -121,86 +120,6 @@ def test_dashboard_period_comparison_is_aggregated_by_backend(
|
||||
assert comparison["token_cost_total_yuan"] == pytest.approx(1.0)
|
||||
|
||||
|
||||
def test_dashboard_coupon_success_rate_excludes_abandoned_sessions(
|
||||
admin_client: TestClient, admin_token: str
|
||||
) -> None:
|
||||
started_date = date(2038, 1, 16)
|
||||
started_at = datetime(2038, 1, 16, 8, tzinfo=UTC)
|
||||
sessions = [
|
||||
("coupon-rate-completed-1", "coupon-rate-device-1", "completed"),
|
||||
("coupon-rate-completed-2", "coupon-rate-device-2", "completed"),
|
||||
("coupon-rate-failed", "coupon-rate-device-3", "failed"),
|
||||
("coupon-rate-abandoned", "coupon-rate-device-4", "abandoned"),
|
||||
]
|
||||
db = SessionLocal()
|
||||
try:
|
||||
for trace_id, device_id, status in sessions:
|
||||
db.add(
|
||||
CouponSession(
|
||||
trace_id=trace_id,
|
||||
device_id=device_id,
|
||||
status=status,
|
||||
app_env="prod",
|
||||
platforms=["meituan-waimai"],
|
||||
started_at=started_at,
|
||||
started_date=started_date,
|
||||
)
|
||||
)
|
||||
for index, device_id in enumerate(("coupon-rate-device-1", "coupon-rate-device-2")):
|
||||
db.add(
|
||||
CouponClaimRecord(
|
||||
device_id=device_id,
|
||||
coupon_id=f"mt_dashboard_rate_{index}",
|
||||
claim_date=started_date,
|
||||
status="success",
|
||||
app_env="prod",
|
||||
)
|
||||
)
|
||||
db.commit()
|
||||
finally:
|
||||
db.close()
|
||||
|
||||
response = admin_client.get(
|
||||
"/admin/api/stats/overview",
|
||||
params={"date_from": "2038-01-16", "date_to": "2038-01-16"},
|
||||
headers=_auth(admin_token),
|
||||
)
|
||||
assert response.status_code == 200, response.text
|
||||
coupon = response.json()["period"]["coupon"]
|
||||
assert coupon["started"] == 4
|
||||
assert coupon["abandoned"] == 1
|
||||
assert coupon["success_denominator"] == 3
|
||||
assert coupon["all_success"] == 2
|
||||
assert coupon["success_rate"] == pytest.approx(2 / 3, abs=0.0001)
|
||||
|
||||
db = SessionLocal()
|
||||
try:
|
||||
db.add(
|
||||
CouponSession(
|
||||
trace_id="coupon-rate-only-abandoned",
|
||||
device_id="coupon-rate-device-only-abandoned",
|
||||
status="abandoned",
|
||||
app_env="prod",
|
||||
platforms=["meituan-waimai"],
|
||||
started_at=datetime(2038, 1, 17, 8, tzinfo=UTC),
|
||||
started_date=date(2038, 1, 17),
|
||||
)
|
||||
)
|
||||
db.commit()
|
||||
finally:
|
||||
db.close()
|
||||
|
||||
empty_denominator_response = admin_client.get(
|
||||
"/admin/api/stats/overview",
|
||||
params={"date_from": "2038-01-17", "date_to": "2038-01-17"},
|
||||
headers=_auth(admin_token),
|
||||
)
|
||||
assert empty_denominator_response.status_code == 200
|
||||
only_abandoned = empty_denominator_response.json()["period"]["coupon"]
|
||||
assert only_abandoned["success_denominator"] == 0
|
||||
assert only_abandoned["success_rate"] is None
|
||||
|
||||
|
||||
def test_user_list_and_detail(admin_client: TestClient, admin_token: str) -> None:
|
||||
uid = _seed_user_with_data("13800000002")
|
||||
r = admin_client.get("/admin/api/users", headers=_auth(admin_token))
|
||||
|
||||
@@ -10,7 +10,6 @@ from app.admin.repositories.coupon_data import (
|
||||
_point_scores_by_trace,
|
||||
coupon_data_report,
|
||||
coupon_point_details,
|
||||
coupon_user_records,
|
||||
)
|
||||
from app.admin.security import create_admin_token
|
||||
from app.db.session import SessionLocal
|
||||
@@ -39,7 +38,6 @@ def test_point_scores_by_trace() -> None:
|
||||
stats = _point_scores_by_trace(db, [trace])[trace]
|
||||
assert stats["succeeded"] == 2
|
||||
assert stats["tried"] == 3
|
||||
assert stats["events"] == 4
|
||||
details = coupon_point_details(db, trace_id=trace)
|
||||
assert [item["status"] for item in details] == [
|
||||
"success", "already_claimed", "failed", "skipped"
|
||||
@@ -50,8 +48,8 @@ def test_point_scores_by_trace() -> None:
|
||||
db.close()
|
||||
|
||||
|
||||
def test_skipped_detail_is_distinguished_from_no_events() -> None:
|
||||
"""仅有 skipped 时分数仍为0/0,但保留事件数供前端开放明细。"""
|
||||
def test_skipped_detail_does_not_create_a_score() -> None:
|
||||
"""仅有 skipped 时按需明细仍可查到,但列表没有虚假的 0/0 分数。"""
|
||||
db = SessionLocal()
|
||||
trace = "point-score-skipped"
|
||||
try:
|
||||
@@ -65,7 +63,7 @@ def test_skipped_detail_is_distinguished_from_no_events() -> None:
|
||||
db.flush()
|
||||
|
||||
scores = _point_scores_by_trace(db, [trace, "missing-trace"])
|
||||
assert scores[trace] == {"succeeded": 0, "tried": 0, "events": 1}
|
||||
assert trace not in scores
|
||||
assert "missing-trace" not in scores
|
||||
assert coupon_point_details(db, trace_id=trace)[0]["status"] == "skipped"
|
||||
finally:
|
||||
@@ -116,112 +114,6 @@ def test_coupon_data_report_returns_scores_without_embedding_details() -> None:
|
||||
db.close()
|
||||
|
||||
|
||||
def test_coupon_data_report_marks_abandoned_without_point_results() -> None:
|
||||
"""中途退出且没有逐券终态时返回0/0,其他状态缺埋点仍保持为空。"""
|
||||
db = SessionLocal()
|
||||
report_date = date(2020, 1, 6)
|
||||
user_id = 910006
|
||||
try:
|
||||
db.add_all([
|
||||
CouponSession(
|
||||
trace_id="point-score-abandoned-without-result",
|
||||
device_id="score-abandoned-device",
|
||||
user_id=user_id,
|
||||
status="abandoned",
|
||||
app_env="prod",
|
||||
platforms=["meituan-waimai"],
|
||||
started_at=datetime(2020, 1, 6, tzinfo=UTC),
|
||||
started_date=report_date,
|
||||
),
|
||||
CouponSession(
|
||||
trace_id="point-score-completed-without-result",
|
||||
device_id="score-completed-device",
|
||||
user_id=user_id,
|
||||
status="completed",
|
||||
app_env="prod",
|
||||
platforms=["meituan-waimai"],
|
||||
started_at=datetime(2020, 1, 6, 1, tzinfo=UTC),
|
||||
started_date=report_date,
|
||||
),
|
||||
CouponSession(
|
||||
trace_id="point-score-abandoned-with-result",
|
||||
device_id="score-abandoned-result-device",
|
||||
user_id=user_id,
|
||||
status="abandoned",
|
||||
app_env="prod",
|
||||
platforms=["meituan-waimai"],
|
||||
started_at=datetime(2020, 1, 6, 2, tzinfo=UTC),
|
||||
started_date=report_date,
|
||||
),
|
||||
CouponSession(
|
||||
trace_id="point-score-abandoned-skipped-only",
|
||||
device_id="score-abandoned-skipped-device",
|
||||
user_id=user_id,
|
||||
status="abandoned",
|
||||
app_env="prod",
|
||||
platforms=["meituan-waimai"],
|
||||
started_at=datetime(2020, 1, 6, 3, tzinfo=UTC),
|
||||
started_date=report_date,
|
||||
),
|
||||
])
|
||||
db.add_all([
|
||||
CouponClaimEvent(
|
||||
trace_id="point-score-abandoned-with-result",
|
||||
device_id="score-abandoned-result-device",
|
||||
coupon_id=f"mt-abandoned-{status}",
|
||||
claim_date=report_date,
|
||||
status=status,
|
||||
)
|
||||
for status in ("success", "failed")
|
||||
])
|
||||
db.add(CouponClaimEvent(
|
||||
trace_id="point-score-abandoned-skipped-only",
|
||||
device_id="score-abandoned-skipped-device",
|
||||
coupon_id="mt-abandoned-skipped",
|
||||
claim_date=report_date,
|
||||
status="skipped",
|
||||
))
|
||||
db.flush()
|
||||
|
||||
report = coupon_data_report(
|
||||
db,
|
||||
date_from=report_date.isoformat(),
|
||||
date_to=report_date.isoformat(),
|
||||
app_env="prod",
|
||||
)
|
||||
rows = {item["trace_id"]: item for item in report["items"]}
|
||||
abandoned = rows["point-score-abandoned-without-result"]
|
||||
assert abandoned["point_success_count"] == 0
|
||||
assert abandoned["point_total_count"] == 0
|
||||
assert abandoned["point_event_count"] == 0
|
||||
|
||||
abandoned_with_result = rows["point-score-abandoned-with-result"]
|
||||
assert abandoned_with_result["point_success_count"] == 1
|
||||
assert abandoned_with_result["point_total_count"] == 2
|
||||
assert abandoned_with_result["point_event_count"] == 2
|
||||
|
||||
abandoned_skipped = rows["point-score-abandoned-skipped-only"]
|
||||
assert abandoned_skipped["point_success_count"] == 0
|
||||
assert abandoned_skipped["point_total_count"] == 0
|
||||
assert abandoned_skipped["point_event_count"] == 1
|
||||
|
||||
completed = rows["point-score-completed-without-result"]
|
||||
assert completed["point_success_count"] is None
|
||||
assert completed["point_total_count"] is None
|
||||
|
||||
user_rows = {
|
||||
item["trace_id"]: item
|
||||
for item in coupon_user_records(db, user_id=user_id)["items"]
|
||||
}
|
||||
assert user_rows["point-score-abandoned-without-result"]["point_total_count"] == 0
|
||||
assert user_rows["point-score-abandoned-with-result"]["point_total_count"] == 2
|
||||
assert user_rows["point-score-abandoned-skipped-only"]["point_event_count"] == 1
|
||||
assert user_rows["point-score-completed-without-result"]["point_total_count"] is None
|
||||
finally:
|
||||
db.rollback()
|
||||
db.close()
|
||||
|
||||
|
||||
def test_coupon_point_details_endpoint() -> None:
|
||||
"""前端点击使用的接口按约定返回 trace_id 和逐券 items。"""
|
||||
db = SessionLocal()
|
||||
|
||||
@@ -1,184 +0,0 @@
|
||||
from __future__ import annotations
|
||||
|
||||
import importlib.util
|
||||
from datetime import datetime
|
||||
from pathlib import Path
|
||||
|
||||
import sqlalchemy as sa
|
||||
|
||||
|
||||
def _load_migration():
|
||||
path = (
|
||||
Path(__file__).parents[1]
|
||||
/ "alembic"
|
||||
/ "versions"
|
||||
/ "deepseek_v4_flash_price.py"
|
||||
)
|
||||
spec = importlib.util.spec_from_file_location("deepseek_v4_flash_price", path)
|
||||
assert spec is not None and spec.loader is not None
|
||||
module = importlib.util.module_from_spec(spec)
|
||||
spec.loader.exec_module(module)
|
||||
return module
|
||||
|
||||
|
||||
def test_migration_follows_comparison_status_normalization():
|
||||
migration = _load_migration()
|
||||
|
||||
assert migration.down_revision == "comparison_below_min_success"
|
||||
|
||||
|
||||
def test_migration_adds_price_and_corrects_only_mispriced_snapshot(monkeypatch):
|
||||
migration = _load_migration()
|
||||
engine = sa.create_engine("sqlite://")
|
||||
metadata = sa.MetaData()
|
||||
app_config = sa.Table(
|
||||
"app_config",
|
||||
metadata,
|
||||
sa.Column("key", sa.String(64), primary_key=True),
|
||||
sa.Column("value", sa.JSON, nullable=False),
|
||||
sa.Column("updated_at", sa.DateTime),
|
||||
)
|
||||
comparison = sa.Table(
|
||||
"comparison_record",
|
||||
metadata,
|
||||
sa.Column("id", sa.Integer, primary_key=True),
|
||||
sa.Column("llm_calls", sa.JSON),
|
||||
sa.Column("llm_cost_yuan", sa.Float),
|
||||
sa.Column("llm_price_snapshot", sa.JSON),
|
||||
)
|
||||
metadata.create_all(engine)
|
||||
|
||||
calls = [
|
||||
{
|
||||
"model": "deepseek-v4-flash",
|
||||
"error": None,
|
||||
"usage": {"prompt_tokens": 12031, "completion_tokens": 125},
|
||||
},
|
||||
{
|
||||
"model": "qwen3.5-flash",
|
||||
"error": None,
|
||||
"usage": {"prompt_tokens": 2416, "completion_tokens": 119},
|
||||
},
|
||||
]
|
||||
snapshot = {
|
||||
"mode": "per_model",
|
||||
"prices": {
|
||||
"deepseek-v4-flash": {
|
||||
"input_per_1m": 3.0,
|
||||
"output_per_1m": 15.0,
|
||||
"_source": "default",
|
||||
},
|
||||
"qwen3.5-flash": {
|
||||
"input_per_1m": 0.2,
|
||||
"output_per_1m": 2.0,
|
||||
"_source": "per_model",
|
||||
},
|
||||
},
|
||||
}
|
||||
|
||||
with engine.begin() as conn:
|
||||
original_updated_at = datetime(2026, 7, 13, 18, 20, 10)
|
||||
conn.execute(
|
||||
app_config.insert().values(
|
||||
key="llm_token_price",
|
||||
value={
|
||||
"per_model": {
|
||||
"qwen3.5-flash": {
|
||||
"input_per_1m": 0.2,
|
||||
"output_per_1m": 2.0,
|
||||
}
|
||||
},
|
||||
"default": {"input_per_1m": 3.0, "output_per_1m": 15.0},
|
||||
},
|
||||
updated_at=original_updated_at,
|
||||
)
|
||||
)
|
||||
conn.execute(
|
||||
comparison.insert().values(
|
||||
id=1,
|
||||
llm_calls=calls,
|
||||
llm_cost_yuan=0.038689,
|
||||
llm_price_snapshot=snapshot,
|
||||
)
|
||||
)
|
||||
monkeypatch.setattr(migration.op, "get_bind", lambda: conn)
|
||||
|
||||
migration.upgrade()
|
||||
|
||||
config = conn.execute(
|
||||
sa.select(app_config.c.value).where(
|
||||
app_config.c.key == "llm_token_price"
|
||||
)
|
||||
).scalar_one()
|
||||
assert config["per_model"]["deepseek-v4-flash"] == {
|
||||
"input_per_1m": 1.0,
|
||||
"output_per_1m": 2.0,
|
||||
}
|
||||
assert conn.execute(
|
||||
sa.select(app_config.c.updated_at).where(
|
||||
app_config.c.key == "llm_token_price"
|
||||
)
|
||||
).scalar_one() == original_updated_at
|
||||
corrected = conn.execute(sa.select(comparison)).mappings().one()
|
||||
assert corrected["llm_cost_yuan"] == 0.013002
|
||||
assert corrected["llm_price_snapshot"]["prices"]["deepseek-v4-flash"] == {
|
||||
"input_per_1m": 1.0,
|
||||
"output_per_1m": 2.0,
|
||||
"_source": "per_model",
|
||||
}
|
||||
assert corrected["llm_price_snapshot"]["pricing_correction"] == (
|
||||
"deepseek_v4_flash_price"
|
||||
)
|
||||
|
||||
migration.downgrade()
|
||||
|
||||
reverted_config = conn.execute(
|
||||
sa.select(app_config.c.value).where(
|
||||
app_config.c.key == "llm_token_price"
|
||||
)
|
||||
).scalar_one()
|
||||
assert "deepseek-v4-flash" not in reverted_config["per_model"]
|
||||
reverted = conn.execute(sa.select(comparison)).mappings().one()
|
||||
assert reverted["llm_cost_yuan"] == 0.038689
|
||||
assert "pricing_correction" not in reverted["llm_price_snapshot"]
|
||||
|
||||
|
||||
def test_downgrade_preserves_price_that_existed_before_upgrade(monkeypatch):
|
||||
migration = _load_migration()
|
||||
engine = sa.create_engine("sqlite://")
|
||||
metadata = sa.MetaData()
|
||||
app_config = sa.Table(
|
||||
"app_config",
|
||||
metadata,
|
||||
sa.Column("key", sa.String(64), primary_key=True),
|
||||
sa.Column("value", sa.JSON, nullable=False),
|
||||
sa.Column("updated_at", sa.DateTime),
|
||||
)
|
||||
sa.Table(
|
||||
"comparison_record",
|
||||
metadata,
|
||||
sa.Column("id", sa.Integer, primary_key=True),
|
||||
sa.Column("llm_calls", sa.JSON),
|
||||
sa.Column("llm_cost_yuan", sa.Float),
|
||||
sa.Column("llm_price_snapshot", sa.JSON),
|
||||
)
|
||||
metadata.create_all(engine)
|
||||
|
||||
explicit_price = {"input_per_1m": 1.0, "output_per_1m": 2.0}
|
||||
with engine.begin() as conn:
|
||||
conn.execute(
|
||||
app_config.insert().values(
|
||||
key="llm_token_price",
|
||||
value={
|
||||
"per_model": {"deepseek-v4-flash": explicit_price},
|
||||
"default": {"input_per_1m": 3.0, "output_per_1m": 15.0},
|
||||
},
|
||||
)
|
||||
)
|
||||
monkeypatch.setattr(migration.op, "get_bind", lambda: conn)
|
||||
|
||||
migration.upgrade()
|
||||
migration.downgrade()
|
||||
|
||||
config = conn.execute(sa.select(app_config.c.value)).scalar_one()
|
||||
assert config["per_model"]["deepseek-v4-flash"] == explicit_price
|
||||
@@ -108,10 +108,6 @@ def test_get_llm_prices_falls_back_to_default_then_uses_override():
|
||||
# 无 override → CONFIG_DEFS 默认(含 per_model / default)
|
||||
prices = get_llm_prices(db)
|
||||
assert "per_model" in prices and "default" in prices
|
||||
assert prices["per_model"]["deepseek-v4-flash"] == {
|
||||
"input_per_1m": 1.0,
|
||||
"output_per_1m": 2.0,
|
||||
}
|
||||
# 有 override → 用 DB 值
|
||||
app_config.set_value(
|
||||
db, "llm_token_price",
|
||||
|
||||
Reference in New Issue
Block a user