Compare commits
5 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| b06d8716fa | |||
| 1a61cb5a65 | |||
| 67ac2dcbbb | |||
| 08a49504fa | |||
| ab2de6ec79 |
@@ -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)
|
||||
)
|
||||
@@ -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)
|
||||
@@ -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)),
|
||||
|
||||
@@ -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,
|
||||
|
||||
@@ -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 = []
|
||||
# 全量环境
|
||||
|
||||
@@ -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()
|
||||
|
||||
@@ -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
@@ -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
@@ -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",
|
||||
)
|
||||
|
||||
@@ -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…";
|
||||
# 失败:具体原因(如"美团、京东外卖均未找到该商品")。前端在比价失败时当原因展示。
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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(不单列),
|
||||
|
||||
@@ -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. **平台不可用时手动恢复** —— 登录服务器,交互式选择备份恢复,与平台走**同一套脚本**。
|
||||
|
||||
数据规模 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 认证可用。
|
||||
@@ -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",
|
||||
|
||||
|
||||
@@ -801,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()
|
||||
@@ -828,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(
|
||||
|
||||
@@ -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 → 展示京东那条)。"""
|
||||
|
||||
@@ -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", "商品未找到"),
|
||||
]
|
||||
@@ -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
|
||||
@@ -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",
|
||||
|
||||
Reference in New Issue
Block a user