diff --git a/alembic/versions/coin_transaction_trace_id.py b/alembic/versions/coin_transaction_trace_id.py index 1258c63..b914688 100644 --- a/alembic/versions/coin_transaction_trace_id.py +++ b/alembic/versions/coin_transaction_trace_id.py @@ -1,7 +1,7 @@ """coin_transaction.trace_id (金币记录按会话聚合比价/领券看广告金币) Revision ID: coin_transaction_trace_id -Revises: comparison_updated_at +Revises: savings_record_trace_id Create Date: 2026-08-07 比价/领券信息流发奖时把本场 trace_id 一并写入 coin_transaction;金币变动记录接口按 @@ -16,7 +16,7 @@ import sqlalchemy as sa # revision identifiers, used by Alembic. revision: str = 'coin_transaction_trace_id' -down_revision: Union[str, Sequence[str], None] = 'comparison_updated_at' +down_revision: Union[str, Sequence[str], None] = 'savings_record_trace_id' branch_labels: Union[str, Sequence[str], None] = None depends_on: Union[str, Sequence[str], None] = None diff --git a/alembic/versions/savings_record_trace_id.py b/alembic/versions/savings_record_trace_id.py new file mode 100644 index 0000000..8ff2fd8 --- /dev/null +++ b/alembic/versions/savings_record_trace_id.py @@ -0,0 +1,34 @@ +"""savings_record.trace_id (下单归因到的比价 trace_id) + +Revision ID: savings_record_trace_id +Revises: comparison_updated_at +Create Date: 2026-08-07 00:00:00.000000 + +「已下单」从店级改单次级:下单上报带上本次比价的 trace_id,落这一列,读取时按 +trace_id 精确对齐 comparison_record.trace_id —— 同一家店比价多次,只有真正下单的 +那一条标「已下单」。可空:demo 行 / 老客户端 / 历史订单没有 trace_id(→ 不进任何 +记录的「已下单」,不做回填)。见 repositories.comparison._ordered_trace_id_select。 +""" +from typing import Sequence, Union + +from alembic import op +import sqlalchemy as sa + + +# revision identifiers, used by Alembic. +revision: str = 'savings_record_trace_id' +down_revision: Union[str, Sequence[str], None] = 'comparison_updated_at' +branch_labels: Union[str, Sequence[str], None] = None +depends_on: Union[str, Sequence[str], None] = None + + +def upgrade() -> None: + with op.batch_alter_table('savings_record', schema=None) as batch_op: + batch_op.add_column(sa.Column('trace_id', sa.String(length=64), nullable=True)) + batch_op.create_index(batch_op.f('ix_savings_record_trace_id'), ['trace_id'], unique=False) + + +def downgrade() -> None: + with op.batch_alter_table('savings_record', schema=None) as batch_op: + batch_op.drop_index(batch_op.f('ix_savings_record_trace_id')) + batch_op.drop_column('trace_id') diff --git a/app/admin/repositories/queries.py b/app/admin/repositories/queries.py index 28d10ed..ae455ba 100644 --- a/app/admin/repositories/queries.py +++ b/app/admin/repositories/queries.py @@ -76,30 +76,30 @@ def _device_marketing_name(model: str | None) -> str | None: def _attach_comparison_order_status(db: Session, items: list[ComparisonRecord]) -> None: - """按 C 端既有口径给比价记录批量补充是否真实下单。""" + """按 C 端既有口径给比价记录批量补充是否真实下单(按 trace_id 精确对齐真实下单上报)。""" user_ids = {item.user_id for item in items if item.user_id is not None} - shop_names = {item.store_name for item in items if item.store_name} + trace_ids = {item.trace_id for item in items if item.trace_id} ordered_pairs: set[tuple[int, str]] = set() - if user_ids and shop_names: + if user_ids and trace_ids: rows = db.execute( - select(SavingsRecord.user_id, SavingsRecord.shop_name) + select(SavingsRecord.user_id, SavingsRecord.trace_id) .where( SavingsRecord.user_id.in_(user_ids), SavingsRecord.source == "compare", - SavingsRecord.shop_name.in_(shop_names), + SavingsRecord.trace_id.in_(trace_ids), ) .distinct() ).all() ordered_pairs = { - (row.user_id, row.shop_name) + (row.user_id, row.trace_id) for row in rows - if row.shop_name is not None + if row.trace_id is not None } for item in items: item.ordered = bool( item.user_id is not None - and item.store_name - and (item.user_id, item.store_name) in ordered_pairs + and item.trace_id + and (item.user_id, item.trace_id) in ordered_pairs ) diff --git a/app/admin/repositories/stats.py b/app/admin/repositories/stats.py index 94b8beb..d2fbb3f 100644 --- a/app/admin/repositories/stats.py +++ b/app/admin/repositories/stats.py @@ -376,15 +376,14 @@ def dashboard_overview( .where( SavingsRecord.user_id == ComparisonRecord.user_id, SavingsRecord.source == "compare", - SavingsRecord.shop_name.is_not(None), - SavingsRecord.shop_name == ComparisonRecord.store_name, + SavingsRecord.trace_id.is_not(None), + SavingsRecord.trace_id == ComparisonRecord.trace_id, ) .exists() ) period_ordered_count = _count( ComparisonRecord, *period_comparison_conds, - ComparisonRecord.store_name.is_not(None), ordered_exists, ) diff --git a/app/models/savings.py b/app/models/savings.py index 0a85c3f..656f13a 100644 --- a/app/models/savings.py +++ b/app/models/savings.py @@ -56,6 +56,9 @@ class SavingsRecord(Base): source_deeplink: Mapped[str | None] = mapped_column(String(512), nullable=True) # 客户端幂等键(UUID);demo 行为 NULL client_event_id: Mapped[str | None] = mapped_column(String(64), nullable=True) + # 本次下单归因到的那次比价的 trace_id(客户端归因时从比价会话缓存带来;demo/老客户端/历史行为 NULL)。 + # 「已下单」按它精确对齐 comparison_record.trace_id —— 同店多次比价只标真正下单的那一条。 + trace_id: Mapped[str | None] = mapped_column(String(64), nullable=True, index=True) device_id: Mapped[str | None] = mapped_column(String(128), nullable=True) created_at: Mapped[datetime] = mapped_column( diff --git a/app/repositories/comparison.py b/app/repositories/comparison.py index 2703f92..01d5cfa 100644 --- a/app/repositories/comparison.py +++ b/app/repositories/comparison.py @@ -737,18 +737,18 @@ def harvest_abort( return rec -def _ordered_shop_name_select(user_id: int): - """该用户「真实下单」(source='compare')覆盖到的店名 select,给「已下单」筛选当子查询。 +def _ordered_trace_id_select(user_id: int): + """该用户「真实下单」(source='compare')覆盖到的比价 trace_id 的 select,给「已下单」筛选当子查询。 - 口径与 [_ordered_shop_names] 完全一致,只是时机不同:那边是**拿到本页之后**按 candidates + 口径与 [_ordered_trace_ids] 完全一致,只是时机不同:那边是**拿到本页之后**按 candidates 反查打标;这边是**分页之前**就要过滤,拿不到 candidates,只能整段下推成子查询。 - 没有先捞成集合再展开 IN (...) 字面量 —— 重度用户下单过的店名可能上千,展开会撞 SQLite + 没有先捞成集合再展开 IN (...) 字面量 —— 重度用户下单过的 trace_id 可能上千,展开会撞 SQLite 的绑定变量上限,而且又变回了那个「随下单量线性变慢」的老写法。 """ - return select(SavingsRecord.shop_name).where( + return select(SavingsRecord.trace_id).where( SavingsRecord.user_id == user_id, SavingsRecord.source == "compare", - SavingsRecord.shop_name.is_not(None), + SavingsRecord.trace_id.is_not(None), ) @@ -760,27 +760,27 @@ def _like_escape(kw: str) -> str: return kw.replace("\\", "\\\\").replace("%", "\\%").replace("_", "\\_") -def _ordered_shop_names(db: Session, user_id: int, candidates: set[str]) -> set[str]: - """[candidates] 里哪些店名被该用户「真实下单」(source='compare')覆盖过,用来打「已下单」。 +def _ordered_trace_ids(db: Session, user_id: int, candidates: set[str]) -> set[str]: + """[candidates] 里哪些 trace_id 被该用户「真实下单」(source='compare')覆盖过,用来打「已下单」。 - 只认 compare(归因命中后真实上报),demo 演示数据不算。下单上报不带 trace_id, - 只能按店名对齐——两边店名同源(都来自比价意图识别阶段的门店名 query),精确相等即视为同店。 - 语义=店级:同一家店比价过多次,这些记录会一并标「已下单」。 + 只认 compare(归因命中后真实上报),demo 演示数据不算。下单上报带上本次比价的 trace_id → + 按 trace_id 精确对齐 comparison_record.trace_id:同一家店比价多次,只有真正下单的那一条会被 + 标「已下单」。没带 trace_id 的下单(历史 / 老客户端)对齐不上任何记录 → 不进「已下单」。 - ⚠️ 只查**本页出现过的店名**(candidates ≤ limit 条),不再把该用户全部下单店名捞回内存: - 老写法随下单量线性增长,重度用户几千行全读一遍只为跟 50 条记录取交集。空集合直接返回 + ⚠️ 只查**本页出现过的 trace_id**(candidates ≤ limit 条),不再把该用户全部下单 trace_id 捞回 + 内存:老写法随下单量线性增长,重度用户几千行全读一遍只为跟 50 条记录取交集。空集合直接返回 (避免 IN () 非法)。 """ if not candidates: return set() rows = db.execute( - select(SavingsRecord.shop_name).where( + select(SavingsRecord.trace_id).where( SavingsRecord.user_id == user_id, SavingsRecord.source == "compare", - SavingsRecord.shop_name.in_(candidates), + SavingsRecord.trace_id.in_(candidates), ).distinct() ).scalars().all() - return {s for s in rows if s} + return {t for t in rows if t} def _ad_coins_by_trace(db: Session, user_id: int, trace_ids: list[str]) -> dict[str, int]: @@ -832,7 +832,7 @@ def list_records( ordered: bool | None = None, keyword: str | None = None, ) -> tuple[list[ComparisonRecord], int | None]: - """比价记录分页(按创建时间倒序、id 兜底,游标式)。附「已下单」店级标记 + 「看广告赚的金币」(瞬态,不写库)。""" + """比价记录分页(按创建时间倒序、id 兜底,游标式)。附「已下单」标记(按 trace_id 精确对齐真实下单)+ 「看广告赚的金币」(瞬态,不写库)。""" stmt = ( select(ComparisonRecord) .where(ComparisonRecord.user_id == user_id) @@ -844,7 +844,7 @@ def list_records( # 分页之后一页里可能一条都不命中,列表看着就是空的/卡住的,得翻很多页才蹦出一条。 if ordered: stmt = stmt.where( - ComparisonRecord.store_name.in_(_ordered_shop_name_select(user_id)) + ComparisonRecord.trace_id.in_(_ordered_trace_id_select(user_id)) ) kw = (keyword or "").strip() if kw: @@ -865,16 +865,16 @@ def list_records( items = list(db.execute(stmt).scalars().all()) next_cursor = items[-1].id if len(items) == limit else None - # 「已下单」标记:本页记录的 store_name 若落在该用户真实下单的店名集合里即 True。 + # 「已下单」标记:本页记录的 trace_id 若落在该用户真实下单(带 trace_id)的集合里即 True。 # ordered / ad_coins_earned 均非 ORM 列,仅挂实例上供 ComparisonRecordOut(from_attributes) 读出,不持久化。 - page_shops = {it.store_name for it in items if it.store_name} - # ordered=True 时上面已按同一口径(_ordered_shop_name_select)筛过,本页必然全是已下单, - # 省掉这次反查;其余情况照旧按本页店名反查 savings。 - ordered_shops = page_shops if ordered else _ordered_shop_names(db, user_id, page_shops) + page_traces = {it.trace_id for it in items} + # ordered=True 时上面已按同一口径(_ordered_trace_id_select)筛过,本页必然全是已下单, + # 省掉这次反查;其余情况按本页 trace_id 反查 savings。 + ordered_traces = page_traces if ordered else _ordered_trace_ids(db, user_id, page_traces) # 「本次比价看广告赚的金币」:按本页 trace_id 一次性聚合(同 ordered 范式)。 ad_coins = _ad_coins_by_trace(db, user_id, [it.trace_id for it in items]) for it in items: - it.ordered = bool(it.store_name and it.store_name in ordered_shops) + it.ordered = it.trace_id in ordered_traces it.ad_coins_earned = ad_coins.get(it.trace_id, 0) return items, next_cursor diff --git a/app/repositories/savings.py b/app/repositories/savings.py index 5c8574d..1d3f53d 100644 --- a/app/repositories/savings.py +++ b/app/repositories/savings.py @@ -177,6 +177,7 @@ def create_from_report( source_platform_name=req.source_platform_name, source_deeplink=req.source_deeplink, client_event_id=req.client_event_id, + trace_id=req.trace_id, device_id=req.device_id, source="compare", # created_at 显式存 naive 北京 wall-clock(与 demo 行、聚合 _local_date 的 naive 分支一致)。 diff --git a/app/schemas/order.py b/app/schemas/order.py index 9d7af17..3eafaac 100644 --- a/app/schemas/order.py +++ b/app/schemas/order.py @@ -12,6 +12,9 @@ class OrderReportRequest(BaseModel): paid_amount_cents: int = Field(..., ge=0, description="实际支付金额(分)") device_id: str | None = Field(None, max_length=128) # ===== 比价时携带的记账信息(客户端从意图识别阶段缓存而来;旧版客户端可能不传,故全部可空)===== + # 本次下单归因到的那次比价的 trace_id。服务端据此把「已下单」精确对齐到该条比价记录 + # (同店多次比价只标真正下单的那一条)。旧客户端不传 → None → 该单不进任何记录的「已下单」。 + trace_id: str | None = Field(None, max_length=64, description="归因到的比价 trace_id") shop_name: str | None = Field(None, max_length=128, description="门店名,如 肯德基宅急送(天北路店)") dishes: list[str] = Field(default_factory=list, description="菜品名列表") original_price_cents: int | None = Field(None, ge=0, description="源平台原价(分),省额=原价−实付") diff --git a/tests/test_admin_read.py b/tests/test_admin_read.py index 6b86c2e..8630a82 100644 --- a/tests/test_admin_read.py +++ b/tests/test_admin_read.py @@ -906,6 +906,7 @@ def test_comparison_records_show_real_order_status( order_amount_cents=1800, saved_amount_cents=300, shop_name="真实下单店", + trace_id="comparison-ordered-shop", source="compare", client_event_id="admin-comparison-real-order", ), diff --git a/tests/test_compare_record.py b/tests/test_compare_record.py index f16798c..c452981 100644 --- a/tests/test_compare_record.py +++ b/tests/test_compare_record.py @@ -223,10 +223,10 @@ def test_stats_compare_count_and_saved(client) -> None: def test_records_ordered_flag(client) -> None: - """「已下单」店级标记:店名命中该用户 source='compare' 的下单记录才 True。 + """「已下单」按 trace_id 标记:下单上报带上本次比价的 trace_id,精确命中该条记录才 True。 - 覆盖 list_records 只按**本页店名**反查 savings 的写法(原来是把该用户全部下单店名捞回内存 - 再取交集,随下单量线性变慢)——两种写法结果必须一致,故这里按店名逐条断言。 + 覆盖 list_records 只按**本页 trace_id**反查 savings 的写法(不把该用户全部下单 trace_id 捞回 + 内存再取交集)——两种写法结果必须一致,故这里逐条断言。 """ token = _login(client, "13800002010") @@ -236,14 +236,15 @@ def test_records_ordered_flag(client) -> None: other["store_name"] = "没下过单的店" client.post("/api/v1/compare/record", json=other, headers=_auth(token)) - # 下单前:两条都不该带「已下单」 - items = client.get("/api/v1/compare/records", headers=_auth(token)).json()["items"] - assert {it["store_name"]: it["ordered"] for it in items} == { - "海底捞(朝阳店)": False, - "没下过单的店": False, + # 下单前:两条都不该带「已下单」。按 trace_id 断言本测试自己的两条(不受同库其它用例数据干扰)。 + flags = { + it["trace_id"]: it["ordered"] + for it in client.get("/api/v1/compare/records", headers=_auth(token)).json()["items"] } + assert flags["ord-1"] is False + assert flags["ord-2"] is False - # 对海底捞真实下单一笔(order/report 写 source='compare' 的 savings_record) + # 对海底捞那次比价真实下单一笔:上报带上该次比价的 trace_id(order/report 写 source='compare') r = client.post( "/api/v1/order/report", json={ @@ -255,23 +256,28 @@ def test_records_ordered_flag(client) -> None: "paid_amount_cents": 12350, "shop_name": "海底捞(朝阳店)", "original_price_cents": 12850, + "trace_id": "ord-1", }, headers=_auth(token), ) assert r.status_code == 200, r.text - # 下单后:只有同店名那条翻成 True,另一条不受影响 - items = client.get("/api/v1/compare/records", headers=_auth(token)).json()["items"] - assert {it["store_name"]: it["ordered"] for it in items} == { - "海底捞(朝阳店)": True, - "没下过单的店": False, + # 下单后:只有命中 trace_id(ord-1)那条翻成 True,同店的另一次比价(ord-2)不受影响 + flags = { + it["trace_id"]: it["ordered"] + for it in client.get("/api/v1/compare/records", headers=_auth(token)).json()["items"] } + assert flags["ord-1"] is True + assert flags["ord-2"] is False - # 别人的下单不该影响本人标记(_ordered_shop_names 按 user_id 过滤) + # 别人的下单不该影响本人标记(_ordered_trace_ids 按 user_id 过滤) token_b = _login(client, "13800002011") client.post("/api/v1/compare/record", json=_food_payload("ord-b"), headers=_auth(token_b)) - items_b = client.get("/api/v1/compare/records", headers=_auth(token_b)).json()["items"] - assert [it["ordered"] for it in items_b] == [False] + flags_b = { + it["trace_id"]: it["ordered"] + for it in client.get("/api/v1/compare/records", headers=_auth(token_b)).json()["items"] + } + assert flags_b["ord-b"] is False def test_records_list_omits_raw_payload(client) -> None: @@ -316,20 +322,24 @@ def test_records_ordered_filter(client) -> None: q["store_name"] = "没下过单的店" client.post("/api/v1/compare/record", json=q, headers=_auth(token)) - client.post( - "/api/v1/order/report", - json={ - "client_event_id": "evt-ordered-filter", - "platform": "美团", - "platform_package": "com.sankuai.meituan", - "pay_channel": "wechat", - "compared_price_cents": 12350, - "paid_amount_cents": 12350, - "shop_name": "下过单的店", - "original_price_cents": 12850, - }, - headers=_auth(token), - ) + # 对这 3 条「下过单的店」比价分别真实下单(各带自己的 trace_id);2 条「没下过单的店」不下单 + for i in range(3): + r = client.post( + "/api/v1/order/report", + json={ + "client_event_id": f"evt-ordered-filter-{i}", + "platform": "美团", + "platform_package": "com.sankuai.meituan", + "pay_channel": "wechat", + "compared_price_cents": 12350, + "paid_amount_cents": 12350, + "shop_name": "下过单的店", + "original_price_cents": 12850, + "trace_id": f"of-ordered-{i}", + }, + headers=_auth(token), + ) + assert r.status_code == 200, r.text # 不传 ordered:5 条全出(「全部记录」tab 口径不变) assert len(client.get("/api/v1/compare/records", headers=_auth(token)).json()["items"]) == 5 @@ -439,3 +449,84 @@ def test_trace_id_required(client) -> None: token = _login(client, "13800002009") r = client.post("/api/v1/compare/record", json={"business_type": "food"}, headers=_auth(token)) assert r.status_code == 422 + + +def test_records_ordered_by_trace_id(client) -> None: + """「已下单」精确到单次比价:同一家店比价多次,只有真实下单的那一次(下单上报带 trace_id)才算已下单。 + + 下单上报带上本次比价的 trace_id,服务端按 trace_id 对齐 —— 同店其它比价(哪怕成功)不再被一并 + 标「已下单」。这是从"店级"改"单次级"的核心契约。 + """ + token = _login(client, "13800002016") + + # 同一家店(默认 payload 店名都是「海底捞(朝阳店)」)比价两次,不同 trace_id + client.post("/api/v1/compare/record", json=_food_payload("ord-trace-a"), headers=_auth(token)) + client.post("/api/v1/compare/record", json=_food_payload("ord-trace-b"), headers=_auth(token)) + + # 只对其中一次(trace-a)真实下单,上报带上该次比价的 trace_id + r = client.post( + "/api/v1/order/report", + json={ + "client_event_id": "evt-trace-a", + "platform": "美团", + "platform_package": "com.sankuai.meituan", + "pay_channel": "wechat", + "compared_price_cents": 12350, + "paid_amount_cents": 12350, + "shop_name": "海底捞(朝阳店)", + "original_price_cents": 12850, + "trace_id": "ord-trace-a", + }, + headers=_auth(token), + ) + assert r.status_code == 200, r.text + + # ordered=true 只出下单的那条(trace-a);同店没下单的 trace-b 不进来 + page = client.get("/api/v1/compare/records?ordered=true", headers=_auth(token)).json() + assert [it["trace_id"] for it in page["items"]] == ["ord-trace-a"] + + # 全部记录里,只有 trace-a 带 ordered=true,同店的 trace-b 仍是 False + items = client.get("/api/v1/compare/records", headers=_auth(token)).json()["items"] + assert {it["trace_id"]: it["ordered"] for it in items} == { + "ord-trace-a": True, + "ord-trace-b": False, + } + + +def test_records_ordered_ignores_order_without_trace_id(client) -> None: + """无 trace_id 的下单(老客户端/历史)即使店名相同也不进「已下单」—— 锁定「无店名回退」契约。 + + 下单上报不带 trace_id → savings_record.trace_id 为空 → 对齐不上任何比价记录。老逻辑靠店名 + 会把同店比价一并标已下单,这条正是要防它被改回去。 + """ + token = _login(client, "13800002017") + + # 一条比价记录(默认店名「海底捞(朝阳店)」) + client.post("/api/v1/compare/record", json=_food_payload("no-trace-order"), headers=_auth(token)) + + # 同店名、但**不带 trace_id** 的下单(模拟老客户端) + r = client.post( + "/api/v1/order/report", + json={ + "client_event_id": "evt-no-trace", + "platform": "美团", + "platform_package": "com.sankuai.meituan", + "pay_channel": "wechat", + "compared_price_cents": 12350, + "paid_amount_cents": 12350, + "shop_name": "海底捞(朝阳店)", + "original_price_cents": 12850, + }, + headers=_auth(token), + ) + assert r.status_code == 200, r.text + + # 店名虽同,但订单没 trace_id → 该记录不该被标「已下单」,也不进 ordered=true(无店名回退) + flags = { + it["trace_id"]: it["ordered"] + for it in client.get("/api/v1/compare/records", headers=_auth(token)).json()["items"] + } + assert flags["no-trace-order"] is False + assert client.get( + "/api/v1/compare/records?ordered=true", headers=_auth(token) + ).json()["items"] == []