diff --git a/alembic/versions/comparison_below_minimum_as_success.py b/alembic/versions/comparison_below_minimum_as_success.py new file mode 100644 index 0000000..e985fea --- /dev/null +++ b/alembic/versions/comparison_below_minimum_as_success.py @@ -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) + ) diff --git a/alembic/versions/deepseek_v4_flash_price.py b/alembic/versions/deepseek_v4_flash_price.py new file mode 100644 index 0000000..c60213b --- /dev/null +++ b/alembic/versions/deepseek_v4_flash_price.py @@ -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) diff --git a/app/admin/repositories/queries.py b/app/admin/repositories/queries.py index a2158ee..c085445 100644 --- a/app/admin/repositories/queries.py +++ b/app/admin/repositories/queries.py @@ -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)), diff --git a/app/admin/routers/comparison.py b/app/admin/routers/comparison.py index 8876dc0..33e37f1 100644 --- a/app/admin/routers/comparison.py +++ b/app/admin/routers/comparison.py @@ -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, diff --git a/app/admin/schemas/comparison.py b/app/admin/schemas/comparison.py index c9f1d43..71bdd9d 100644 --- a/app/admin/schemas/comparison.py +++ b/app/admin/schemas/comparison.py @@ -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 = [] # 全量环境 diff --git a/app/core/config_schema.py b/app/core/config_schema.py index b2289b3..6532204 100644 --- a/app/core/config_schema.py +++ b/app/core/config_schema.py @@ -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", }, diff --git a/app/models/comparison.py b/app/models/comparison.py index 7edbc99..51956d4 100644 --- a/app/models/comparison.py +++ b/app/models/comparison.py @@ -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…"; # 失败:具体原因(如"美团、京东外卖均未找到该商品")。前端在比价失败时当原因展示。 diff --git a/app/repositories/comparison.py b/app/repositories/comparison.py index 1ad23bd..5af04f1 100644 --- a/app/repositories/comparison.py +++ b/app/repositories/comparison.py @@ -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 diff --git a/app/schemas/compare_record.py b/app/schemas/compare_record.py index 74fde79..87c768f 100644 --- a/app/schemas/compare_record.py +++ b/app/schemas/compare_record.py @@ -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(不单列), diff --git a/tests/test_admin_read.py b/tests/test_admin_read.py index 806faed..081cea9 100644 --- a/tests/test_admin_read.py +++ b/tests/test_admin_read.py @@ -720,6 +720,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 +757,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( diff --git a/tests/test_compare_harvest.py b/tests/test_compare_harvest.py index 5a42e3a..22adef2 100644 --- a/tests/test_compare_harvest.py +++ b/tests/test_compare_harvest.py @@ -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 → 展示京东那条)。""" diff --git a/tests/test_comparison_admin_summary.py b/tests/test_comparison_admin_summary.py index 2f86e24..cbf43c4 100644 --- a/tests/test_comparison_admin_summary.py +++ b/tests/test_comparison_admin_summary.py @@ -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() diff --git a/tests/test_comparison_below_minimum_migration.py b/tests/test_comparison_below_minimum_migration.py new file mode 100644 index 0000000..9bb2322 --- /dev/null +++ b/tests/test_comparison_below_minimum_migration.py @@ -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", "商品未找到"), + ] diff --git a/tests/test_deepseek_price_migration.py b/tests/test_deepseek_price_migration.py new file mode 100644 index 0000000..5cba327 --- /dev/null +++ b/tests/test_deepseek_price_migration.py @@ -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 diff --git a/tests/test_llm_cost.py b/tests/test_llm_cost.py index dbf95e0..ae938aa 100644 --- a/tests/test_llm_cost.py +++ b/tests/test_llm_cost.py @@ -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",