Compare commits

...

6 Commits

Author SHA1 Message Date
guke b06d8716fa docs: 线上 PG 数据库备份/恢复方案设计
覆盖每日定时备份、平台手动备份、平台指定备份恢复到旁库、平台不可用时 SSH 交互式恢复四条路径;逻辑备份(pg_dump)方案,本地存储 + S3 预留升级路径。

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
2026-08-03 16:11:20 +08:00
linkeyu 1a61cb5a65 修复 DeepSeek V4 Flash TOKEN 成本高估 (#214)
## 改动
- 为 deepseek-v4-flash 配置 DashScope 华北 2 官方单价:输入 ¥1 / 输出 ¥2(每百万 Token)
- 定向重算历史上误用 3/15 兜底价冻结的成本和价格快照
- 保留已有人工单价及配置生效时间,避免影响历史缺失成本回填
- 增加迁移升级、降级和原配置保留测试

## 验证
- 相关测试:17 passed
- ruff check:通过

---------

Co-authored-by: guke <guke@wonderable.ai>
Co-authored-by: linkeyu <798648091@qq.com>
Co-authored-by: unknown <798648091@qq.com>
Reviewed-on: #214
Co-authored-by: linkeyu <linkeyu@wonderable.ai>
Co-committed-by: linkeyu <linkeyu@wonderable.ai>
2026-08-01 23:25:27 +08:00
linkeyu 67ac2dcbbb fix: 领券成功率剔除中途退出 (#213)
## 改动说明
- 领券成功率分母改为:发起数 - 中途退出数
- failed 与 started 仍保留在分母
- 接口新增 abandoned 和 success_denominator 字段
- 中途退出已有单券结果时返回真实成功/尝试数
- 中途退出且没有逐券终态时明确返回 0/0
- 新增 point_event_count,区分「只有 skipped、无有效计分结果」和「完全无逐券事件」
- 用户领券记录接口同步聚合逐券结果

## 验证
- 相关后端测试:21 passed
- Ruff:通过
- 线上数据库只读核对:空白记录确实没有 coupon_claim_event

配套前端:WonderableAI/shaguabijia-admin-web#101

---------

Co-authored-by: guke <guke@wonderable.ai>
Co-authored-by: linkeyu <798648091@qq.com>
Reviewed-on: #213
Co-authored-by: linkeyu <linkeyu@wonderable.ai>
Co-committed-by: linkeyu <linkeyu@wonderable.ai>
2026-08-01 23:24:54 +08:00
zuochenyong 08a49504fa fix(反馈): 优化图片处理与静态资源缓存 (#211)
Co-authored-by: guke <guke@wonderable.ai>
Co-authored-by: exinglang <exinglang@qq.com>
Reviewed-on: #211
Co-authored-by: zuochenyong <zuochenyong@wonderable.ai>
Co-committed-by: zuochenyong <zuochenyong@wonderable.ai>
2026-08-01 23:24:17 +08:00
linkeyu ab2de6ec79 修复:统一比价记录状态枚举口径 (#209)
## 变更说明

- below_minimum 归入成功,保留原始业务结局
- store_closed/store_not_found/items_not_found/no_delivery/unsupported 归入失败
- running 保持进行中生命周期状态
- 历史细分状态迁移为三态终态,迁移可逆
- 后台成功/失败筛选及汇总兼容迁移前历史值

## 验证

- 相关回归:48 passed
- Ruff:通过
- Alembic:单一 head

---------

Co-authored-by: guke <guke@wonderable.ai>
Co-authored-by: linkeyu <798648091@qq.com>
Co-authored-by: unknown <798648091@qq.com>
Reviewed-on: #209
Co-authored-by: linkeyu <linkeyu@wonderable.ai>
Co-committed-by: linkeyu <linkeyu@wonderable.ai>
2026-08-01 23:09:21 +08:00
linkeyu 84251770b4 修复:首页轮播排除测试账号数据 (#215)
## 背景

线上首页轮播候选记录中,测试账号 u33(11111111111)贡献约 51.8% 的真实记录,导致该用户频繁出现。

## 改动

- App 首页轮播真实记录查询排除全部已配置测试账号
- 后台首页轮播可展示记录同步采用相同过滤口径
- 同时兼容 TEST_ACCOUNT_PHONE 与 TEST_ACCOUNT_PHONES
- 测试账号配置变化时立即使轮播查询缓存失效
- 不删除历史业务数据,仅在展示查询中排除

## 验证

- ruff check:通过
- pytest tests/test_ops_marquee.py -q:1 passed

---------

Co-authored-by: linkeyu <798648091@qq.com>
Reviewed-on: #215
Co-authored-by: linkeyu <linkeyu@wonderable.ai>
Co-committed-by: linkeyu <linkeyu@wonderable.ai>
2026-08-01 21:35:06 +08:00
28 changed files with 1608 additions and 66 deletions
@@ -0,0 +1,76 @@
"""normalize granular comparison outcomes into terminal record statuses
Revision ID: comparison_below_min_success
Revises: limit_policy_global_bundle
"""
from __future__ import annotations
from collections.abc import Sequence
import sqlalchemy as sa
from sqlalchemy.dialects import postgresql
from alembic import op
revision: str = "comparison_below_min_success"
down_revision: str | Sequence[str] | None = "limit_policy_global_bundle"
branch_labels: str | Sequence[str] | None = None
depends_on: str | Sequence[str] | None = None
_JSON = sa.JSON().with_variant(postgresql.JSONB(), "postgresql")
_FAILED_OUTCOMES = (
"store_closed",
"store_not_found",
"items_not_found",
"no_delivery",
"unsupported",
)
def _comparison_record() -> sa.TableClause:
return sa.table(
"comparison_record",
sa.column("status", sa.String(16)),
sa.column("fail_reason", sa.String(256)),
sa.column("raw_payload", _JSON),
)
def upgrade() -> None:
comparison_record = _comparison_record()
op.execute(
comparison_record.update()
.where(comparison_record.c.status == "below_minimum")
.values(status="success", fail_reason=None)
)
op.execute(
comparison_record.update()
.where(comparison_record.c.status.in_(_FAILED_OUTCOMES))
.values(status="failed")
)
def downgrade() -> None:
comparison_record = _comparison_record()
# The write paths deliberately preserve the granular outcome. Restore only
# rows that this change normalized, without touching ordinary successes.
raw_outcome = sa.func.coalesce(
comparison_record.c.raw_payload["record_status"].as_string(),
comparison_record.c.raw_payload["status"].as_string(),
)
op.execute(
comparison_record.update()
.where(
comparison_record.c.status == "success",
raw_outcome == "below_minimum",
)
.values(status="below_minimum")
)
op.execute(
comparison_record.update()
.where(
comparison_record.c.status == "failed",
raw_outcome.in_(_FAILED_OUTCOMES),
)
.values(status=raw_outcome)
)
+211
View File
@@ -0,0 +1,211 @@
"""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)
+29 -10
View File
@@ -166,6 +166,18 @@ 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,
@@ -182,8 +194,9 @@ def _session_to_row(
"app_env": r.app_env,
"started_at": r.started_at,
"claimed_count": r.claimed_count,
"point_success_count": point_stats["succeeded"] if point_stats else None,
"point_total_count": point_stats["tried"] if point_stats else None,
"point_success_count": point_success_count,
"point_total_count": point_total_count,
"point_event_count": point_event_count,
"trace_url": r.trace_url,
"ad_revenue_yuan": ad_revenue_yuan,
}
@@ -194,21 +207,24 @@ 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"),
func.count().label("tried"),
)
.where(
CouponClaimEvent.trace_id.in_(trace_ids),
CouponClaimEvent.status.in_(_SLOT_TRIED),
tried.label("tried"),
func.count().label("events"),
)
.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 or 0)}
for trace_id, success_count, tried in rows
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
if trace_id is not None
}
@@ -400,12 +416,15 @@ 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()
rev_map = crud_ecpm.revenue_yuan_by_trace(db, [r.trace_id for r in rows])
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)
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
],
+34 -6
View File
@@ -39,6 +39,34 @@ from app.repositories import activity, ad_ecpm
# 折算成可提现现金时,非广告金币来源的排除集(广告单独统计、人工调整不算"赚取")
_NON_TASK_BIZ_TYPES = ("reward_video", "feed_ad_reward", "admin_grant", "admin_deduct")
# comparison_record historically persisted a few granular business outcomes as
# top-level statuses. Admin filters and metrics expose lifecycle buckets while
# retaining the raw values until the data migration has run everywhere.
_COMPARISON_STATUS_ALIASES = {
"success": ("success", "below_minimum"),
"failed": (
"failed",
"store_closed",
"store_not_found",
"items_not_found",
"no_delivery",
"unsupported",
),
"cancelled": ("cancelled",),
"running": ("running",),
}
_COMPARISON_SUCCESS_STATUSES = _COMPARISON_STATUS_ALIASES["success"]
_COMPARISON_FAILED_STATUSES = _COMPARISON_STATUS_ALIASES["failed"]
_COMPARISON_COMPLETED_STATUSES = (
*_COMPARISON_SUCCESS_STATUSES,
*_COMPARISON_FAILED_STATUSES,
)
def _comparison_status_condition(status: str):
values = _COMPARISON_STATUS_ALIASES.get(status, (status,))
return ComparisonRecord.status.in_(values)
# 信息流点位场景 → 金币记录「赚取途径」展示名;NULL/未知 = 历史未分类。
_FEED_SCENE_LABEL = {
"comparison": "比价信息流",
@@ -334,7 +362,7 @@ def _comparison_conditions(
)
)
if status:
conditions.append(ComparisonRecord.status == status)
conditions.append(_comparison_status_condition(status))
if business_type:
conditions.append(ComparisonRecord.business_type == business_type)
if store:
@@ -419,7 +447,7 @@ def _comparison_duration_aggregate_stmt(conditions: list, status: str, quantiles
),
).where(
*conditions,
ComparisonRecord.status == status,
_comparison_status_condition(status),
ComparisonRecord.total_ms.is_not(None),
)
@@ -444,7 +472,7 @@ def _comparison_duration_aggregates(
select(ComparisonRecord.total_ms)
.where(
*conditions,
ComparisonRecord.status == status,
_comparison_status_condition(status),
ComparisonRecord.total_ms.is_not(None),
)
.order_by(ComparisonRecord.total_ms)
@@ -474,11 +502,11 @@ def comparison_records_summary(
row = db.execute(
select(
func.count(ComparisonRecord.id),
func.sum(case((ComparisonRecord.status.in_(("success", "failed")), 1), else_=0)),
func.sum(case((ComparisonRecord.status == "success", 1), else_=0)),
func.sum(case((ComparisonRecord.status.in_(_COMPARISON_COMPLETED_STATUSES), 1), else_=0)),
func.sum(case((ComparisonRecord.status.in_(_COMPARISON_SUCCESS_STATUSES), 1), else_=0)),
func.avg(ComparisonRecord.llm_cost_yuan),
func.sum(case((
(ComparisonRecord.status == "success")
ComparisonRecord.status.in_(_COMPARISON_SUCCESS_STATUSES)
& (ComparisonRecord.saved_amount_cents > 0), 1
), else_=0)),
func.sum(case((ComparisonRecord.status == "cancelled", 1), else_=0)),
+9 -1
View File
@@ -566,6 +566,10 @@ 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
@@ -731,9 +735,13 @@ 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_started, 4) if coupon_started else None
round(coupon_all_success / coupon_success_denominator, 4)
if coupon_success_denominator
else None
),
"point_success": coupon_point_success,
"points_per_session": coupon_points_per_session,
+2 -2
View File
@@ -35,7 +35,7 @@ def list_comparison_records(
db: AdminDb,
user_id: Annotated[int | None, Query()] = None,
phone: Annotated[str | None, Query(description="手机号前缀")] = None,
status: Annotated[str | None, Query(pattern="^(success|failed|cancelled)$")] = None,
status: Annotated[str | None, Query(pattern="^(success|failed|cancelled|running)$")] = None,
business_type: Annotated[str | None, Query()] = None,
store: Annotated[str | None, Query(description="店名子串模糊匹配")] = None,
product: Annotated[str | None, Query(description="商品名子串模糊匹配")] = None,
@@ -66,7 +66,7 @@ def comparison_records_summary(
db: AdminDb,
user_id: Annotated[int | None, Query()] = None,
phone: Annotated[str | None, Query(description="手机号前缀")] = None,
status: Annotated[str | None, Query(pattern="^(success|failed|cancelled)$")] = None,
status: Annotated[str | None, Query(pattern="^(success|failed|cancelled|running)$")] = None,
business_type: Annotated[str | None, Query()] = None,
store: Annotated[str | None, Query(description="店名子串模糊匹配")] = None,
product: Annotated[str | None, Query(description="商品名子串模糊匹配")] = None,
+2 -1
View File
@@ -20,7 +20,7 @@ class AdminComparisonListItem(BaseModel):
trace_id: str
# admin 是 debug 工具,无条件下发 trace_url(不看 user.debug_trace_enabled)
trace_url: str | None = None
status: str
status: str # success / failed / cancelled / running;旧细分值由前端兼容映射
information: str | None = None
store_name: str | None = None
product_names: str | None = None # 下单商品名派生串(顿号分隔;「商品」列展示 + 商品搜索)
@@ -83,6 +83,7 @@ class AdminComparisonDetail(AdminComparisonListItem):
skipped_dish_count: int | None = None
device_id: str | None = None
items: list = []
platforms: list = [] # pricebot 渲染就绪的逐平台卡片模型(status=ok/业务失败细分)
comparison_results: list = [] # 逐平台对比(价格/rank/coupon/打烊...)
skipped_dish_names: list = []
# 全量环境
+8 -2
View File
@@ -79,10 +79,16 @@ class CouponDataRow(BaseModel):
started_at: datetime = Field(..., description="发起时刻(明细「时间」列)")
claimed_count: int | None = None
point_success_count: int | None = Field(
None, description="本次成功单券数(success+already_claimed);无逐券事件为空"
None,
description="本次成功单券数(success+already_claimed);中途退出且无逐券结果为0,其它无事件为空",
)
point_total_count: int | None = Field(
None, description="本次尝试单券数(success+already_claimed+failed,不含 skipped);无逐券事件为空"
None,
description="本次尝试单券数(success+already_claimed+failed,不含 skipped);中途退出且无逐券结果为0,其它无事件为空",
)
point_event_count: int = Field(
0,
description="本次全部逐券事件数(含 skipped);用于区分无有效计分事件与完全无事件",
)
trace_url: str | None = Field(None, description="pricebot 公网 trace 链接(仅 completed 有);admin 渲染可点链接,无则显示可复制 trace_id")
ad_revenue_yuan: float = Field(
+3
View File
@@ -70,6 +70,9 @@ 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
+4 -2
View File
@@ -54,11 +54,13 @@ def _app_status(db_status: str) -> str:
def _record_out(fb) -> FeedbackRecordOut:
images = fb.images or []
return FeedbackRecordOut(
id=fb.id,
content=fb.content,
scene=getattr(fb, "scene", None),
images=fb.images or [],
images=images,
image_thumbnails=[media.feedback_thumbnail_url(url) for url in images],
status=_app_status(fb.status),
reject_reason=getattr(fb, "reject_reason", None),
reward_coins=getattr(fb, "reward_coins", None),
@@ -84,7 +86,7 @@ async def submit_feedback(
device_model: str = Form(default=""),
rom_name: str = Form(default=""),
android_version: str = Form(default=""),
images: list[UploadFile] = File(default=[]),
images: list[UploadFile] = File(default=[]), # noqa: B008 - FastAPI dependency declaration
) -> FeedbackOut:
content = content.strip()
contact = contact.strip()
+5 -1
View File
@@ -227,7 +227,11 @@ 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}},
"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},
},
"default": {"input_per_1m": 3.0, "output_per_1m": 15.0},
"currency": "CNY", "unit": "per_1m_tokens",
},
+90 -2
View File
@@ -11,6 +11,7 @@
from __future__ import annotations
import json
import logging
import os
import secrets
import subprocess
@@ -18,8 +19,17 @@ from decimal import ROUND_HALF_UP, Decimal, InvalidOperation
from pathlib import Path
from typing import Any
from PIL import Image, ImageOps
from app.core.config import settings
logger = logging.getLogger("shagua.media")
_FEEDBACK_DIR = "feedback"
_FEEDBACK_THUMB_DIR = "feedback_thumbs"
_FEEDBACK_THUMB_MAX_PX = 256
_FEEDBACK_THUMB_QUALITY = 78
class MediaError(Exception):
"""上传文件不合法(类型/大小)。调用方转 400。"""
@@ -68,8 +78,86 @@ def save_avatar(user_id: int, data: bytes) -> str:
def save_feedback_image(user_id: int, data: bytes) -> str:
"""保存反馈截图,返回相对 URL(`/media/feedback/<file>`)"""
return _save_image("feedback", user_id, data)
"""保存反馈截图并预生成历史页缩略图,返回原图相对 URL"""
url = _save_image(_FEEDBACK_DIR, user_id, data)
# 缩略图失败不影响反馈受理;读取缩略图 URL 时会按需重试并回退原图。
ensure_feedback_thumbnail(url)
return url
def feedback_thumbnail_url(image_url: str) -> str:
"""把反馈原图 URL 映射成确定的缩略图 URL,不在 records 接口内做图片解码。
上传文件名由服务端生成且不会覆盖;旧数据在客户端真正请求可见图片时按需补图。
"""
paths = _feedback_thumbnail_paths(image_url)
return paths[2] if paths is not None else image_url
def _feedback_thumbnail_paths(image_url: str) -> tuple[Path, Path, str] | None:
prefix = f"{settings.MEDIA_URL_PREFIX}/{_FEEDBACK_DIR}/"
if not image_url.startswith(prefix):
return None
filename = image_url.removeprefix(prefix)
# 只接受当前目录下的单个文件名,避免数据库脏数据造成路径穿越。
if not filename or Path(filename).name != filename:
return None
source = _media_dir(_FEEDBACK_DIR) / filename
thumb_name = f"{Path(filename).stem}.jpg"
thumb = _media_dir(_FEEDBACK_THUMB_DIR) / thumb_name
thumb_url = f"{settings.MEDIA_URL_PREFIX}/{_FEEDBACK_THUMB_DIR}/{thumb_name}"
return source, thumb, thumb_url
def ensure_feedback_thumbnail(image_url: str) -> Path | None:
"""确保缩略图存在并返回文件;生成失败时回退原图,供动态缩略图路由使用。"""
paths = _feedback_thumbnail_paths(image_url)
if paths is None:
return None
source, thumb, _ = paths
if thumb.is_file():
return thumb
if not source.is_file():
return None
temp = thumb.with_name(f".{thumb.name}.{secrets.token_hex(4)}.tmp")
try:
with Image.open(source) as opened:
image = ImageOps.exif_transpose(opened)
image.thumbnail(
(_FEEDBACK_THUMB_MAX_PX, _FEEDBACK_THUMB_MAX_PX),
Image.Resampling.LANCZOS,
)
if image.mode != "RGB":
image = image.convert("RGB")
image.save(
temp,
format="JPEG",
quality=_FEEDBACK_THUMB_QUALITY,
optimize=True,
)
os.replace(temp, thumb)
return thumb
except (Image.DecompressionBombError, OSError, ValueError):
logger.warning("生成反馈缩略图失败: %s", source, exc_info=True)
return source
finally:
temp.unlink(missing_ok=True)
def feedback_thumbnail_file(filename: str) -> Path | None:
"""由缩略图文件名找到原反馈图并按需生成,非法/不存在返回 None。"""
if not filename or Path(filename).name != filename or Path(filename).suffix.lower() != ".jpg":
return None
stem = Path(filename).stem
for ext in (".jpg", ".png", ".webp"):
original = _media_dir(_FEEDBACK_DIR) / f"{stem}{ext}"
if original.is_file():
original_url = f"{settings.MEDIA_URL_PREFIX}/{_FEEDBACK_DIR}/{original.name}"
return ensure_feedback_thumbnail(original_url)
return None
def save_report_image(user_id: int, data: bytes) -> str:
+32 -2
View File
@@ -10,7 +10,7 @@ from collections.abc import AsyncIterator
from contextlib import asynccontextmanager
from pathlib import Path
from fastapi import FastAPI
from fastapi import FastAPI, HTTPException
from fastapi.middleware.cors import CORSMiddleware
from fastapi.responses import FileResponse
from fastapi.staticfiles import StaticFiles
@@ -44,6 +44,7 @@ from app.api.v1.tasks import router as tasks_router
from app.api.v1.user import router as user_router
from app.api.v1.wallet import router as wallet_router
from app.api.v1.wxpay import router as wxpay_router
from app.core import media
from app.core.config import settings
from app.core.cps_reconcile_worker import (
start_cps_reconcile_worker,
@@ -82,6 +83,19 @@ setup_logging(debug=settings.APP_DEBUG)
logger = logging.getLogger("shagua.main")
class FeedbackMediaStaticFiles(StaticFiles):
"""反馈原图/缩略图文件名不可变,可长期缓存,避免列表反复回源。"""
async def get_response(self, path: str, scope):
response = await super().get_response(path, scope)
media_path = path.replace("\\", "/").lstrip("/")
if response.status_code == 200 and media_path.startswith(
("feedback/", "feedback_thumbs/")
):
response.headers["Cache-Control"] = "public, max-age=31536000, immutable"
return response
@asynccontextmanager
async def lifespan(_: FastAPI) -> AsyncIterator[None]:
# 提示而非强制建表:生产用 alembic upgrade head,本地 dev 也建议先跑一次 migration。
@@ -212,8 +226,24 @@ def download_apk() -> FileResponse:
)
@app.get(
f"{settings.MEDIA_URL_PREFIX}/feedback_thumbs/{{filename}}",
tags=["feedback"],
include_in_schema=False,
)
def feedback_thumbnail(filename: str) -> FileResponse:
"""旧反馈图按首次可见请求补缩略图;新图上传时已预生成。"""
path = media.feedback_thumbnail_file(filename)
if path is None:
raise HTTPException(status_code=404, detail="图片不存在")
return FileResponse(
path,
headers={"Cache-Control": "public, max-age=31536000, immutable"},
)
app.mount(
settings.MEDIA_URL_PREFIX,
StaticFiles(directory=str(_media_root)),
FeedbackMediaStaticFiles(directory=str(_media_root)),
name="media",
)
+1 -1
View File
@@ -97,7 +97,7 @@ class ComparisonRecord(Base):
total_dish_count: Mapped[int | None] = mapped_column(Integer, nullable=True)
skipped_dish_count: Mapped[int | None] = mapped_column(Integer, nullable=True)
# success(拿到有效对比)/ failed(出错或没采到目标价)
# success(流程正常完成,含 below_minimum)/ failed(技术异常或未形成可比报价,含店铺打烊等)
status: Mapped[str] = mapped_column(String(16), nullable=False, default="success")
# done 帧 information 文案。成功:"在美团找到同店,到手价 ¥X…";
# 失败:具体原因(如"美团、京东外卖均未找到该商品")。前端在比价失败时当原因展示。
+39 -10
View File
@@ -76,6 +76,30 @@ _BIZ_STATUS_PRIORITY = (
)
def _normalize_record_status(status: str | None) -> str | None:
"""Map granular business outcomes onto the record lifecycle status.
``record_status`` describes the business outcome, while
``comparison_record.status`` is also the completed-comparison flag used by
milestones, stats and idempotent rewards. ``below_minimum`` is a completed
success because the target cart produced a trustworthy conclusion. Other
known target-side outcomes did not produce a comparable quote and belong
to the failed record bucket. The granular outcome remains in ``raw_payload``
and ``platform_results`` for result rendering.
"""
if status == "below_minimum":
return "success"
if status in {
"store_closed",
"store_not_found",
"items_not_found",
"no_delivery",
"unsupported",
}:
return "failed"
return status
def _store_closed_text(reason: str | None) -> str:
"""打烊/暂停营业/休息类 reason 常带脏店名元数据 → 只留结论,套简短模板。"""
r = reason or ""
@@ -164,9 +188,10 @@ def _derive(payload: ComparisonRecordIn) -> dict:
is_source_best = best.is_source if best is not None else None
# status:优先 pricebot record_status(区分 below_minimum/store_closed) → 客户端显式 status
# → 兜底"非源且有价"=success/否则 failed。record_status 让"未满起送"不再塌缩成 failed。
status = payload.record_status or payload.status
# status:优先 pricebot record_status → 客户端显式 status → 兜底派生。
# below_minimum 是已形成可信结论的正常完成态,记录级归 success;细分结局仍完整保留在
# raw_payload/platform_results,供结果卡展示"未满起送"。
status = _normalize_record_status(payload.record_status or payload.status)
if status is None:
has_valid_target = any(
(not r.is_source) and r.price is not None for r in results
@@ -201,7 +226,9 @@ def upsert_record(
# 单源派生: 与 harvest_done 一致, payload 带 platforms 时从它派生(唯一真相源
# _derive_from_platforms), 老客户端不带 platforms 时回退 _derive(从 comparison_results)。
if payload.platforms:
derived = _derive_from_platforms(payload.platforms, payload.record_status)
derived = _derive_from_platforms(
payload.platforms, payload.record_status or payload.status
)
# 对齐 _derive 返回键(#189 fail_reason): 两路径 fields 键集一致, 覆盖已有行时不残留旧值
derived["fail_reason"] = (
_derive_fail_display(payload.information, payload.platform_results or {})
@@ -356,10 +383,10 @@ def _derive_from_results(
"saved_amount_cents": saved_amount_cents,
"is_source_best": best.get("is_source") if best else None,
"store_name": (src_row or {}).get("store_name") or None,
# 记录级结局: 优先用 pricebot 下发的 record_status(区分 below_minimum/store_closed,
# 不再把"未满起送"塌缩成 failed → 记录页不再误报"网络开小差"); 旧 pricebot 未下发时
# 回退老的 success/failed 二态派生, 向后兼容。
"status": record_status or ("success" if has_valid_target else "failed"),
# below_minimum 已完成到购物车并形成可信结论,记录级计 success;细分结局仍在
# raw_payload/platform_results。旧 pricebot 未下发 record_status 时回退二态派生。
"status": _normalize_record_status(record_status)
or ("success" if has_valid_target else "failed"),
}
@@ -412,7 +439,8 @@ def _derive_from_platforms(
"saved_amount_cents": saved_amount_cents,
"is_source_best": (best.get("role") == "source") if best else None,
"store_name": store_name or None,
"status": record_status or ("success" if has_valid_target else "failed"),
"status": _normalize_record_status(record_status)
or ("success" if has_valid_target else "failed"),
}
@@ -584,7 +612,8 @@ def harvest_done(
行不存在(理论上帧0已建;防御)则新建。"""
results = done_params.get("comparison_results") or []
# 展示模型统一数组(pricebot 新增, 每平台一行自带 status/is_best): 原样存, 记录页据此直渲染。
# record_status: 记录级结局(success/below_minimum/store_closed/failed), 覆盖老二态派生
# record_status: 记录级业务结局(success/below_minimum/store_closed/failed)。其中
# below_minimum 是正常完成态,持久化 status 归 success,原值仍随 done_params 落 raw_payload。
platforms = done_params.get("platforms") or []
record_status = done_params.get("record_status")
# 单源派生: platforms(含 pricebot 权威 is_best)是唯一真相源, best_*/source_*/saved/status
+21 -7
View File
@@ -25,6 +25,7 @@ from datetime import datetime, timedelta
from sqlalchemy import delete, func, select, update
from sqlalchemy.orm import Session
from app.core.config import settings
from app.core.rewards import CN_TZ
from app.models.comparison import ComparisonRecord
from app.models.ops_marquee_seed import OpsMarqueeSeed
@@ -54,7 +55,7 @@ _SEED_MAX_CENTS = 100000
# 展示层随机(抽样/金额/时间/名字合成)仍每次重算,缓存只省查询;新记录最多晚 30s 进轮播,可接受。
_REAL_ROWS_TTL_SECONDS = 30
_REAL_ROWS_FETCH_CAP = 600 # 一次多取些,够 limit≤30 去重后取数;命中缓存后复用
_real_rows_cache: dict = {"at": None, "rows": None}
_real_rows_cache: dict = {"at": None, "rows": None, "test_phones": None}
# ===== 用户标识脱敏(对齐 PRD) + 种子无真实昵称时的假名合成 =====
# 脱敏规则(按字符数,中英文皆适用):有昵称→n≥5「首+***+末」、n=4「首+**+末」、n≤3「首+**」;
@@ -198,10 +199,16 @@ def _recent_real_rows(db: Session) -> list[tuple[int, int, str | None]]:
返回纯元组(脱离 session),可安全跨请求复用。极端并发下偶尔多查一次(无锁、幂等),纯门面无副作用。
"""
now = datetime.now(CN_TZ)
test_phones = tuple(sorted(settings.test_account_phones))
cached, at = _real_rows_cache["rows"], _real_rows_cache["at"]
if cached is not None and at is not None and (now - at).total_seconds() < _REAL_ROWS_TTL_SECONDS:
if (
cached is not None
and at is not None
and _real_rows_cache["test_phones"] == test_phones
and (now - at).total_seconds() < _REAL_ROWS_TTL_SECONDS
):
return cached
rows = db.execute(
stmt = (
select(ComparisonRecord.user_id, ComparisonRecord.saved_amount_cents, User.nickname)
.join(User, User.id == ComparisonRecord.user_id)
.where(
@@ -211,9 +218,12 @@ def _recent_real_rows(db: Session) -> list[tuple[int, int, str | None]]:
)
.order_by(ComparisonRecord.created_at.desc())
.limit(_REAL_ROWS_FETCH_CAP)
).all()
)
if test_phones:
stmt = stmt.where(User.phone.not_in(test_phones))
rows = db.execute(stmt).all()
out = [(int(uid), int(sc), nick) for uid, sc, nick in rows]
_real_rows_cache["rows"], _real_rows_cache["at"] = out, now
_real_rows_cache.update(rows=out, at=now, test_phones=test_phones)
return out
@@ -340,7 +350,7 @@ def list_real_records(
# pool: [(cluster_key, item)];cluster_key 供去连簇——真实=user_id、种子=各自唯一负数(互不聚簇)
pool: list[tuple[int, dict]] = []
if mode != "seed":
rows = db.execute(
stmt = (
select(
ComparisonRecord.user_id,
ComparisonRecord.saved_amount_cents,
@@ -355,7 +365,11 @@ def list_real_records(
)
.order_by(ComparisonRecord.created_at.desc())
.limit(_REAL_BROWSE_CAP)
).all()
)
test_phones = tuple(sorted(settings.test_account_phones))
if test_phones:
stmt = stmt.where(User.phone.not_in(test_phones))
rows = db.execute(stmt).all()
for uid, sc, nick, ca in rows:
pool.append((
int(uid),
+3 -2
View File
@@ -111,8 +111,9 @@ class ComparisonRecordIn(BaseModel):
# status/is_best/display/display_order,记录页据此直渲染。宽松 list[dict] 存(结构由
# pricebot 定,server 只原样落库),前端读它、老记录空时回退 comparison_results。
platforms: list[dict] = Field(default_factory=list)
# 记录级结局(pricebot 下发): success/below_minimum/store_closed/failed。让"未满起送"不再
# 被塌缩成 failed。_derive 优先用它、其次客户端 status、再兜底二态派生。
# 记录级业务结局(pricebot 下发): success/below_minimum/store_closed/failed。
# below_minimum 表示流程正常完成,持久化主状态归 success;store_closed/items_not_found 等
# 已知无报价结局归 failed。原值仍随 raw_payload 落库,admin/记录页从 platform_results 展示细分结论。
record_status: str | None = None
# 逐平台结局摘要(含失败平台的细分原因 status: store_not_found/items_not_found/below_minimum/
# unsupported/...)。来自 done.params.platform_results,客户端透传;落 raw_payload(不单列),
+2
View File
@@ -33,6 +33,8 @@ class FeedbackRecordOut(BaseModel):
# 比价反馈的问题场景(找错商品/优惠不对…);普通反馈为 None
scene: str | None = None
images: list[str] = Field(default_factory=list)
# 与 images 下标一一对应;生成失败时该项回退原图 URL,兼容历史数据。
image_thumbnails: list[str] = Field(default_factory=list)
status: str
reject_reason: str | None = None
reward_coins: int | None = None
@@ -0,0 +1,304 @@
# 线上 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. **平台不可用时手动恢复** —— 登录服务器,交互式选择备份恢复,与平台走**同一套脚本**。
数据规模 110GB → **逻辑备份(`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` | 平台 API6 个端点) |
| `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 认证可用。
+3
View File
@@ -35,6 +35,9 @@ dependencies = [
# multipart form (FastAPI 表单上传依赖)
"python-multipart>=0.0.9",
# 用户反馈截图缩略图,避免 App 历史页为 48dp 小图下载数 MB 原图
"pillow>=11.0.0",
# admin 后台账号密码 hash(用户侧是手机号+验证码登录,不需要密码;admin 才用)
"bcrypt>=4.0.0",
+94 -1
View File
@@ -1,7 +1,7 @@
"""Admin M2 读接口测试:大盘聚合 + 用户/流水/提现/反馈列表 + 鉴权拦截。"""
from __future__ import annotations
from datetime import UTC, datetime
from datetime import UTC, date, datetime
import pytest
from fastapi.testclient import TestClient
@@ -12,6 +12,7 @@ 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
@@ -120,6 +121,86 @@ 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))
@@ -720,6 +801,16 @@ def test_comparison_records_show_readable_device_and_rom_version(
rom_name="OriginOS",
rom_version=4,
android_version="14",
platforms=[
{
"platform_id": "meituan",
"platform_name": "美团",
"status": "ok",
"role": "target",
"price": 18.8,
"is_best": True,
}
],
)
db.add(record)
db.commit()
@@ -747,6 +838,8 @@ def test_comparison_records_show_readable_device_and_rom_version(
assert detail.status_code == 200, detail.text
assert detail.json()["device_model_name"] == "vivo Y77e"
assert detail.json()["rom_version"] == 4
assert detail.json()["platforms"][0]["status"] == "ok"
assert detail.json()["platforms"][0]["is_best"] is True
def test_comparison_records_show_real_order_status(
+131
View File
@@ -13,6 +13,7 @@ import uuid
from unittest.mock import MagicMock, patch
import httpx
import pytest
from sqlalchemy import select
from app.db.session import SessionLocal
@@ -25,6 +26,26 @@ def _tid() -> str:
return uuid.uuid4().hex
@pytest.mark.parametrize(
("raw_status", "record_status"),
[
("success", "success"),
("below_minimum", "success"),
("failed", "failed"),
("store_closed", "failed"),
("store_not_found", "failed"),
("items_not_found", "failed"),
("no_delivery", "failed"),
("unsupported", "failed"),
("cancelled", "cancelled"),
("running", "running"),
(None, None),
],
)
def test_record_status_normalization(raw_status, record_status) -> None:
assert crud._normalize_record_status(raw_status) == record_status
def _done_params() -> dict:
"""一份典型 done 帧 params:美团 25 元 vs 源淘宝闪购 30 元 → 省 5 元、success。"""
return {
@@ -111,6 +132,116 @@ def test_harvest_done_derives_and_newly_success_once(client) -> None:
assert newly2 is False
def test_harvest_done_below_minimum_counts_as_completed_success(client) -> None:
"""未达起送是可信业务结论:主状态/完成奖励归 success,细分结局仍留在 raw_payload。"""
tid = _tid()
done_below_minimum = {
"record_status": "below_minimum",
"comparison_results": [
{
"platform_id": "meituan",
"platform_name": "美团",
"package": "com.sankuai.meituan",
"price": 59.0,
"is_source": True,
"rank": 1,
"store_name": "测试店",
"items": [{"name": "红乌苏", "qty": 1}],
},
],
"platform_results": {
"meituan": {"is_source": True, "status": "source", "price": 59.0},
"taobao_flash": {
"is_source": False,
"status": "below_minimum",
"reason": "购物车未达起送门槛(差 ¥25.2",
},
},
"information": "淘宝未达起送门槛,可加菜凑单后下单",
}
with SessionLocal() as db:
crud.harvest_running(db, trace_id=tid, user_id=None)
rec, newly = crud.harvest_done(
db, trace_id=tid, user_id=None, done_params=done_below_minimum
)
assert rec.status == "success"
assert newly is True
assert rec.fail_reason is None
assert rec.raw_payload["record_status"] == "below_minimum"
assert (
rec.raw_payload["platform_results"]["taobao_flash"]["status"]
== "below_minimum"
)
def test_harvest_done_platforms_below_minimum_counts_as_success(client) -> None:
"""新 platforms 单源派生路径也必须执行同一 below_minimum → success 归一化。"""
tid = _tid()
with SessionLocal() as db:
rec, newly = crud.harvest_done(
db,
trace_id=tid,
user_id=None,
done_params={
"record_status": "below_minimum",
"platforms": [
{
"role": "source",
"platform_id": "meituan",
"platform_name": "美团",
"price": 59.0,
"store_name": "测试店",
"items": [{"name": "红乌苏", "qty": 1}],
},
{
"role": "target",
"platform_id": "taobao_flash",
"platform_name": "淘宝",
"status": "below_minimum",
"price": None,
},
],
"platform_results": {
"taobao_flash": {
"is_source": False,
"status": "below_minimum",
"reason": "购物车未达起送门槛",
}
},
},
)
assert rec.status == "success"
assert newly is True
assert rec.raw_payload["record_status"] == "below_minimum"
def test_legacy_upsert_below_minimum_counts_as_success(client) -> None:
"""灰度期客户端直报路径无论走 status 还是 record_status 都不能落第四种主状态。"""
tid = _tid()
payload = ComparisonRecordIn(
trace_id=tid,
business_type="food",
status="below_minimum",
comparison_results=[],
platform_results={
"taobao_flash": {
"is_source": False,
"status": "below_minimum",
"reason": "购物车未达起送门槛",
}
},
information="淘宝未达起送门槛,可加菜凑单后下单",
)
with SessionLocal() as db:
rec = crud.upsert_record(db, user_id=987654, payload=payload)
assert rec.status == "success"
assert rec.fail_reason is None
assert rec.raw_payload["status"] == "below_minimum"
def test_harvest_done_failed_derives_fail_reason(client) -> None:
"""failed 记录:记录级 information 笼统,但 fail_reason 从 platform_results 救出具体原因
(id 3030 :美团系统失败 + 京东 items_not_found 展示京东那条)"""
+34 -13
View File
@@ -23,7 +23,7 @@ def test_postgresql_duration_summary_uses_ordered_set_aggregates() -> None:
)
assert sql.count("percentile_cont") == 4
assert "comparison_record.status = 'success'" in sql
assert "comparison_record.status IN ('success', 'below_minimum')" in sql
def test_summary_uses_only_success_durations_and_filters_beijing_date() -> None:
@@ -34,6 +34,9 @@ def test_summary_uses_only_success_durations_and_filters_beijing_date() -> None:
("summary-success-b", "success", 3000, 2.0, 0),
("summary-failed", "failed", 100_000, 3.0, 0),
("summary-cancelled", "cancelled", 5000, 4.0, 0),
("summary-below-minimum", "below_minimum", 5000, None, 0),
("summary-store-closed", "store_closed", 200_000, None, 0),
("summary-running", "running", 4000, None, 0),
]
for trace_id, status, total_ms, cost, saved in rows:
db.add(ComparisonRecord(
@@ -56,26 +59,44 @@ def test_summary_uses_only_success_durations_and_filters_beijing_date() -> None:
db, date_from=date(2038, 1, 15), date_to=date(2038, 1, 15)
)
assert summary["started"] == 4
assert summary["completed"] == 3
assert summary["success"] == 2
assert summary["success_rate"] == pytest.approx(2 / 3)
assert summary["started"] == 7
assert summary["completed"] == 5
assert summary["success"] == 3
assert summary["success_rate"] == pytest.approx(3 / 6)
assert summary["avg_token_cost"] == pytest.approx(2.5)
assert summary["lower_price_rate"] == 0.5
assert summary["avg_duration_ms"] == 2000
assert summary["p5_duration_ms"] == 1100
assert summary["p50_duration_ms"] == 2000
assert summary["p95_duration_ms"] == 2900
assert summary["p99_duration_ms"] == 2980
assert summary["lower_price_rate"] == pytest.approx(1 / 3)
assert summary["avg_duration_ms"] == 3000
assert summary["p5_duration_ms"] == 1200
assert summary["p50_duration_ms"] == 3000
assert summary["p95_duration_ms"] == 4800
assert summary["p99_duration_ms"] == 4960
assert summary["cancelled"] == 1
assert summary["cancelled_rate"] == 0.25
assert summary["cancelled_rate"] == pytest.approx(1 / 7)
assert summary["cancelled_p50_ms"] == 5000
items, _next_cursor, total = queries.list_comparison_records(
db, date_from=date(2038, 1, 15), date_to=date(2038, 1, 15), limit=20
)
assert total == 4
assert total == 7
assert {item.trace_id for item in items} == {row[0] for row in rows}
success_items, _next_cursor, success_total = queries.list_comparison_records(
db, status="success", date_from=date(2038, 1, 15), date_to=date(2038, 1, 15)
)
assert success_total == 3
assert {item.status for item in success_items} == {"success", "below_minimum"}
failed_items, _next_cursor, failed_total = queries.list_comparison_records(
db, status="failed", date_from=date(2038, 1, 15), date_to=date(2038, 1, 15)
)
assert failed_total == 2
assert {item.status for item in failed_items} == {"failed", "store_closed"}
running_items, _next_cursor, running_total = queries.list_comparison_records(
db, status="running", date_from=date(2038, 1, 15), date_to=date(2038, 1, 15)
)
assert running_total == 1
assert running_items[0].status == "running"
finally:
db.rollback()
db.close()
@@ -0,0 +1,95 @@
"""Regression coverage for comparison terminal-status normalization."""
from __future__ import annotations
import importlib.util
from pathlib import Path
import sqlalchemy as sa
from alembic.migration import MigrationContext
from alembic.operations import Operations
def _migration_module():
path = (
Path(__file__).parents[1]
/ "alembic"
/ "versions"
/ "comparison_below_minimum_as_success.py"
)
spec = importlib.util.spec_from_file_location(path.stem, 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_business_status_migration_upgrade_and_downgrade() -> None:
engine = sa.create_engine("sqlite:///:memory:")
metadata = sa.MetaData()
records = sa.Table(
"comparison_record",
metadata,
sa.Column("id", sa.Integer, primary_key=True),
sa.Column("status", sa.String(16), nullable=False),
sa.Column("fail_reason", sa.String(256)),
sa.Column("raw_payload", sa.JSON, nullable=False),
)
metadata.create_all(engine)
module = _migration_module()
with engine.begin() as connection:
connection.execute(
records.insert(),
[
{
"status": "below_minimum",
"fail_reason": "旧失败原因",
"raw_payload": {"record_status": "below_minimum"},
},
{
"status": "success",
"fail_reason": None,
"raw_payload": {"record_status": "success"},
},
{
"status": "failed",
"fail_reason": "技术异常",
"raw_payload": {"status": "failed"},
},
{
"status": "store_closed",
"fail_reason": "店铺打烊",
"raw_payload": {"record_status": "store_closed"},
},
{
"status": "items_not_found",
"fail_reason": "商品未找到",
"raw_payload": {"status": "items_not_found"},
},
],
)
module.op = Operations(MigrationContext.configure(connection))
module.upgrade()
upgraded = connection.execute(
sa.select(records.c.status, records.c.fail_reason).order_by(records.c.id)
).all()
assert upgraded == [
("success", None),
("success", None),
("failed", "技术异常"),
("failed", "店铺打烊"),
("failed", "商品未找到"),
]
module.downgrade()
downgraded = connection.execute(
sa.select(records.c.status, records.c.fail_reason).order_by(records.c.id)
).all()
assert downgraded == [
("below_minimum", None),
("success", None),
("failed", "技术异常"),
("store_closed", "店铺打烊"),
("items_not_found", "商品未找到"),
]
+111 -3
View File
@@ -10,6 +10,7 @@ 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
@@ -38,6 +39,7 @@ 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"
@@ -48,8 +50,8 @@ def test_point_scores_by_trace() -> None:
db.close()
def test_skipped_detail_does_not_create_a_score() -> None:
"""仅有 skipped 时按需明细仍可查到,但列表没有虚假的 0/0 分数"""
def test_skipped_detail_is_distinguished_from_no_events() -> None:
"""仅有 skipped 时分数仍为0/0,但保留事件数供前端开放明细"""
db = SessionLocal()
trace = "point-score-skipped"
try:
@@ -63,7 +65,7 @@ def test_skipped_detail_does_not_create_a_score() -> None:
db.flush()
scores = _point_scores_by_trace(db, [trace, "missing-trace"])
assert trace not in scores
assert scores[trace] == {"succeeded": 0, "tried": 0, "events": 1}
assert "missing-trace" not in scores
assert coupon_point_details(db, trace_id=trace)[0]["status"] == "skipped"
finally:
@@ -114,6 +116,112 @@ 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()
+184
View File
@@ -0,0 +1,184 @@
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
+4
View File
@@ -108,6 +108,10 @@ 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",
+77
View File
@@ -0,0 +1,77 @@
from __future__ import annotations
from app.core.config import settings
from app.db.session import SessionLocal
from app.models.comparison import ComparisonRecord
from app.models.user import User
from app.repositories import ops_marquee
def _clear_real_rows_cache() -> None:
ops_marquee._real_rows_cache.update(at=None, rows=None, test_phones=None)
def test_real_marquee_records_exclude_configured_test_account(monkeypatch) -> None:
legacy_test_phone = "19900009991"
listed_test_phone = "19900009992"
real_phone = "19900009993"
monkeypatch.setattr(settings, "TEST_ACCOUNT_PHONE", legacy_test_phone)
monkeypatch.setattr(settings, "TEST_ACCOUNT_PHONES", listed_test_phone)
_clear_real_rows_cache()
try:
with SessionLocal() as db:
legacy_test_user = User(
phone=legacy_test_phone,
username="29900009991",
register_channel="sms",
)
listed_test_user = User(
phone=listed_test_phone,
username="29900009992",
register_channel="sms",
)
real_user = User(
phone=real_phone,
username="29900009993",
register_channel="sms",
)
db.add_all([legacy_test_user, listed_test_user, real_user])
db.flush()
db.add_all([
ComparisonRecord(
user_id=legacy_test_user.id,
trace_id="marquee-legacy-test-account-trace",
status="success",
saved_amount_cents=29_991,
items=[],
comparison_results=[],
),
ComparisonRecord(
user_id=listed_test_user.id,
trace_id="marquee-listed-test-account-trace",
status="success",
saved_amount_cents=29_992,
items=[],
comparison_results=[],
),
ComparisonRecord(
user_id=real_user.id,
trace_id="marquee-real-account-trace",
status="success",
saved_amount_cents=29_993,
items=[],
comparison_results=[],
),
])
db.flush()
app_rows = ops_marquee._recent_real_rows(db)
browse_rows, _ = ops_marquee.list_real_records(db, mode="real", limit=1_000)
finally:
_clear_real_rows_cache()
assert all(saved not in {29_991, 29_992} for _uid, saved, _nickname in app_rows)
assert any(saved == 29_993 for _uid, saved, _nickname in app_rows)
assert all(row["saved_amount_cents"] not in {29_991, 29_992} for row in browse_rows)
assert any(row["saved_amount_cents"] == 29_993 for row in browse_rows)