Merge origin/main: 消双 head — coin_transaction_trace_id 重挂父到 savings_record_trace_id

origin/main(#224) 的 savings_record_trace_id 与本分支 coin_transaction_trace_id 都以
comparison_updated_at 为父,汇合后 alembic 双 head(upgrade head 报 multiple heads)。
把本分支迁移重挂到 savings_record_trace_id 之上,线性化为单 head:
comparison_updated_at → savings_record_trace_id → coin_transaction_trace_id。
两迁移分别改 savings_record / coin_transaction,互不依赖,串行安全。

验证:alembic heads 单 head;scratch DB 从 base 全链 upgrade head 通过;
test_welfare / compare_record / admin_read / order_savings 共 66 项测试通过。

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
This commit is contained in:
guke
2026-08-07 18:53:55 +08:00
10 changed files with 201 additions and 69 deletions
@@ -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
@@ -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')
+9 -9
View File
@@ -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
)
+2 -3
View File
@@ -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,
)
+3
View File
@@ -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(
+24 -24
View File
@@ -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
+1
View File
@@ -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 分支一致)。
+3
View File
@@ -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="源平台原价(分),省额=原价−实付")
+1
View File
@@ -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",
),
+122 -31
View File
@@ -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"] == []