Compare commits

...

3 Commits

Author SHA1 Message Date
guke 89b2d5ed0c test(compare): 锁定「订单无 trace_id → 不进已下单」的无回退契约
老客户端/历史订单不带 trace_id,即使店名相同也不该标已下单。补这条回归测试,
防止将来有人把店名回退悄悄加回来。

Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
2026-08-07 13:51:19 +08:00
guke d699c7c89c fix(compare): 「已下单」从店级改按 trace_id 精确对齐
下单上报带上本次比价的 trace_id 落 savings_record.trace_id,「已下单」筛选/打标改为
按 trace_id 精确对齐 comparison_record.trace_id —— 同店多次比价只标真正下单那条,失败/
未下单记录不再混入。没有 trace_id 的订单(历史/老客户端)不进「已下单」,不做回退/回填。

- savings_record 加 trace_id 列 + 索引 + Alembic 迁移
- OrderReportRequest 加 trace_id;create_from_report 落库
- comparison.list_records 过滤/打标 + helpers 改按 trace_id(删店名匹配)
- admin 列表/详情/看板下单数同口径改 trace_id
- 迁移既有测试到新语义 + 新增 test_records_ordered_by_trace_id

Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
2026-08-07 11:30:07 +08:00
guke a6b660486f fix(auto-exchange): 「当天已兑」标记持久化到 app_config,防非0点重启重复补扫 (#223)
概览
把 0 点自动兑金币 worker 的「当天是否已兑」判定,从内存变量 last_run 改为持久化到 app_config(key auto_exchange.last_run_date)。修掉「进程重启即忘 → 每次非 0 点部署都全量补扫、把 0 点后才达标用户在非 0 点兑现金」的 bug。改 1 个实现文件 + 1 个文档 + 新增 1 个测试(+171 / −11)。

---------

Co-authored-by: guke <guke@autohome.com.cn>
Reviewed-on: #223
2026-08-06 18:48:38 +08:00
12 changed files with 370 additions and 78 deletions
@@ -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,
)
+53 -10
View File
@@ -5,9 +5,14 @@
`wallet.daily_auto_exchange`
健壮性:
- **逐用户幂等**:当天已有 exchange_in 流水的用户跳过( wallet._has_exchange_in_on),
故启动补跑 / 多次唤醒 / 进程重启都安全,不会重复兑
- **当天首跑即补**:进程起来时若当天还没兑过,立即兑一轮(等价原 timer Persistent 补跑)
- **持久化当天已兑标记**(app_config `auto_exchange.last_run_date`):当天已兑则整轮跳过,
**标记跨进程重启不丢** 同一北京日内多次启动/部署不再重复补扫这是修 bug 的关键:原来用
内存变量记今天跑过没,进程重启即归零,导致每次非 0 点部署都全量补扫一遍 0 点后才达标
的用户在** 0 **兑成现金(用户反馈的 0 点也出现金币转现金记录)
- **漏了 0 点即补**:标记 < 今天( 0 点服务器宕机)时标记 != today 照常补跑一轮(保留原
timer Persistent 语义)这类补跑确实落非 0 ,但仅限真漏了 0 的罕见场景,非常态
- **逐用户幂等**:当天已有 exchange_in 流水的用户跳过( wallet._has_exchange_in_on),做兜底
防同机多进程 / 标记写入失败的竞态,不会重复兑
- **同机多进程互斥**:文件锁保证多 worker 只有一个实际跑(防跨进程并发导致 TOCTOU 双兑)
- **开关**:settings.AUTO_EXCHANGE_ENABLED=false 时不启动(与脚本/ timer 同一开关)
"""
@@ -27,10 +32,14 @@ from sqlalchemy.exc import SQLAlchemyError
from app.core import rewards
from app.core.config import settings
from app.db.session import SessionLocal
from app.models.app_config import AppConfig
from app.repositories import wallet as wallet_repo
logger = logging.getLogger("shagua.daily_exchange")
_LOCK_PATH = Path(__file__).resolve().parents[2] / "data" / "daily_exchange.lock"
# 「上次成功自动兑的北京日」标记,持久化在 app_config(仿 compare_alert 水位,不进 CONFIG_DEFS——
# 它是 worker 内部运行状态,非运营可配项)。用它替代内存 last_run:重启不丢 → 同日不重复补扫。
LAST_RUN_KEY = "auto_exchange.last_run_date"
def _touch_lock() -> None:
@@ -72,9 +81,44 @@ def _single_instance_lock(stale_after_sec: int) -> Iterator[bool]:
_LOCK_PATH.unlink()
def _exchange_once() -> dict:
def _read_last_run(db) -> date | None:
"""读持久化的「上次成功自动兑的北京日」标记(app_config,跨进程重启不丢)。脏值当作没跑过。"""
row = db.get(AppConfig, LAST_RUN_KEY)
if row is None or not row.value:
return None
try:
return date.fromisoformat(row.value)
except (ValueError, TypeError):
return None
def _write_last_run(db, day: date) -> None:
"""落库「当天已兑」标记。用专用 key 直接写 AppConfig(仿 compare_alert 水位,不走 set_value)。"""
iso = day.isoformat()
row = db.get(AppConfig, LAST_RUN_KEY)
if row is None:
db.add(AppConfig(key=LAST_RUN_KEY, value=iso, updated_by_admin_id=None))
else:
row.value = iso
db.commit()
def _exchange_if_due(db, today: date) -> dict | None:
"""当天未兑过(持久化标记 != today)才跑一轮并记标记;已兑过返回 None(整轮跳过)。
标记写在 daily_auto_exchange 之后:即便中途崩,标记仍是旧值 下次重跑,逐用户幂等会跳过已兑的
补完剩下的(安全)真漏了 0 (标记 < today)时标记 != today,仍会补跑,保留原当天首跑即补
"""
if _read_last_run(db) == today:
return None
result = wallet_repo.daily_auto_exchange(db)
_write_last_run(db, today)
return result
def _exchange_once(today: date) -> dict | None:
with SessionLocal() as db:
return wallet_repo.daily_auto_exchange(db)
return _exchange_if_due(db, today)
async def _run_loop() -> None:
@@ -89,16 +133,15 @@ async def _run_loop() -> None:
async def _run_locked_loop(interval: int) -> None:
logger.info("daily auto-exchange worker started interval=%ss", interval)
# 本进程上次跑过的北京日;None=尚未跑过本进程(启动即补当天)。
last_run: date | None = None
try:
while True:
try:
_touch_lock()
today = rewards.cn_today()
if last_run != today:
result = await asyncio.to_thread(_exchange_once)
last_run = today
# 「今天是否已兑」以持久化标记为准(见 _exchange_if_due),不再用内存变量 →
# 进程重启不会把当天当「没跑过」重复补扫;已兑当天返回 None(整轮跳过)。
result = await asyncio.to_thread(_exchange_once, today)
if result is not None:
logger.info("daily auto-exchange done date=%s result=%s", today, result)
except SQLAlchemyError:
logger.exception("daily auto-exchange db error")
+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="源平台原价(分),省额=原价−实付")
+5 -1
View File
@@ -5,9 +5,13 @@
## 现状(2026-06 起):已改为 App 进程内任务,不再需要 systemd timer
「0 点自动兑换」现由 **App 进程内后台任务** `app.core.daily_exchange_worker` 负责:App 一起来就
每 10 分钟检查、跨过北京 0 点自动跑一轮(逐用户幂等、文件锁互斥、重启补跑当天遗漏)。**无需再装 / 启用
每 10 分钟检查、跨过北京 0 点自动跑一轮(逐用户幂等、文件锁互斥)。**无需再装 / 启用
`daily-exchange.timer`**。
- **「当天已兑」标记持久化在 app_config**(key=`auto_exchange.last_run_date`):当天兑过后,同一北京日
内进程重启 / 部署**不再重复补扫**——避免把 0 点后才达标的用户在非 0 点兑现金。只有真漏了 0 点
(标记 < 今天,如 0 点服务器宕机)才会在重启后补跑一轮。
- 开关仍是 `.env``AUTO_EXCHANGE_ENABLED`(false → worker 不启动);间隔由 `AUTO_EXCHANGE_CHECK_INTERVAL_SEC`(默认 600s=10min)控。
- 下方 systemd timer / 脚本属**遗留 + 手动应急**:逻辑同一套且幂等,可手动 `--once` 补跑;
但**不要再 `enable` timer 与进程内 worker 并存**(虽幂等不会双兑,纯属多余)。
+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"] == []
+113
View File
@@ -0,0 +1,113 @@
"""0 点自动兑金币 worker:持久化「当天已兑」标记,防同日重启重复补扫。
背景(路径 B bug):worker 原用内存变量 last_run 今天跑过没,进程重启即归零 每次
0 点部署/重启都会全量补扫一遍, 0 点后才达标的用户在** 0 **兑成现金(用户反馈的
0 点也出现金币转现金记录)修复:标记持久化到 app_config,跨重启不丢
"""
from __future__ import annotations
from datetime import date
import pytest
from sqlalchemy import delete, select, update
from app.core import daily_exchange_worker as w
from app.db.session import SessionLocal
from app.models.app_config import AppConfig
from app.models.user import User
from app.models.wallet import CashTransaction, CoinAccount
from app.repositories import wallet as wallet_repo
_PHONE_SEQ = [0]
@pytest.fixture(autouse=True)
def _isolate_exchange_state():
"""daily_auto_exchange 全表扫描 + app_config 标记会跨用例泄漏(SQLite 测试库 session 级共享,
commit rollback no-op),故每个用例先清零所有余额 + 清掉自动兑标记,保证干净起步"""
db = SessionLocal()
try:
db.execute(update(CoinAccount).values(
coin_balance=0, cash_balance_cents=0, invite_cash_balance_cents=0))
db.execute(delete(AppConfig).where(AppConfig.key == w.LAST_RUN_KEY))
db.commit()
finally:
db.close()
yield
def _new_user(db, *, coin: int) -> int:
"""建一个 User + 指定金币余额的 CoinAccount,返回 user_id。"""
_PHONE_SEQ[0] += 1
# 198 段避免撞其他测试文件的固定手机号 / username UNIQUE
u = User(phone=f"198{_PHONE_SEQ[0]:08d}", username=f"ax{_PHONE_SEQ[0]}", status="active")
db.add(u)
db.flush()
acc = wallet_repo.get_or_create_account(db, u.id, commit=False)
acc.coin_balance = coin
acc.total_coin_earned = coin
db.flush()
return u.id
def test_first_run_converts_and_records_marker() -> None:
"""首轮(标记为空)→ 到分全额兑 + 落库当天标记。"""
db = SessionLocal()
try:
uid = _new_user(db, coin=300)
db.commit()
today = date(2026, 8, 6)
result = w._exchange_if_due(db, today)
assert result is not None and result["converted"] >= 1
acc = db.get(CoinAccount, uid)
assert acc.coin_balance == 0 and acc.cash_balance_cents > 0 # 300 为整分,无零头
assert w._read_last_run(db) == today # 标记落库
finally:
db.rollback()
db.close()
def test_same_day_restart_does_not_resweep() -> None:
"""路径 B 回归:0 点兑过后,同日进程重启(下午部署)不得再补扫当天新达标用户。"""
db = SessionLocal()
try:
today = date(2026, 8, 6)
first = _new_user(db, coin=200)
db.commit()
assert w._exchange_if_due(db, today) is not None
assert db.get(CoinAccount, first).coin_balance == 0 # 首轮兑掉
# 0 点后才达标的用户(白天攒够;或 0 点是零头、白天赚够)
late = _new_user(db, coin=500)
db.commit()
# 模拟同一北京日内进程重启 → 持久化标记 == today → 整轮跳过,不碰 late
assert w._exchange_if_due(db, today) is None
acc = db.get(CoinAccount, late)
assert acc.coin_balance == 500 and acc.cash_balance_cents == 0
assert db.execute(select(CashTransaction).where(
CashTransaction.user_id == late,
CashTransaction.biz_type == "exchange_in",
)).first() is None # late 无任何兑现金流水
finally:
db.rollback()
db.close()
def test_new_day_runs_again() -> None:
"""跨到北京新的一天(或真漏了 0 点,标记 < today)→ 照常补跑。"""
db = SessionLocal()
try:
w._write_last_run(db, date(2026, 8, 6)) # 昨天已兑
late = _new_user(db, coin=500)
db.commit()
result = w._exchange_if_due(db, date(2026, 8, 7))
assert result is not None and result["converted"] >= 1
assert db.get(CoinAccount, late).coin_balance == 0
assert w._read_last_run(db) == date(2026, 8, 7)
finally:
db.rollback()
db.close()