feat(store_mapping): 跨平台店铺映射表 + pricebot 内部上报端点

新增「平台店铺表」资产层: 淘宝比价拿到 shopId 后, pricebot server→server 把跨平台
店铺身份(各平台 id/名 + 地理 + 来源)落库, 作为未来"我见过这家店→跳过重搜/匹配"的源头。
与 price_observation(价格事实)平行、独立。

- models/store_mapping.py: store_mapping 表 22 列 —— 跨平台身份(id/name_taobao/meituan/jd)
  + 地理(city/geohash/lng/lat/taobao_address) + 溯源(source_platform/trace_id/device/user)
  + 淘宝原料(share_url/resolved_url/deeplink) + attrs(JSONB) + created_at。
- schemas/store_mapping.py + repositories/store_mapping.py: append-only, trace_id 幂等
  (pricebot 重试/replay 不重复写; 并发 IntegrityError 兜底返已存在行), 跨方言安全。
- api/internal/store.py: POST /internal/store-mapping(复用 price.py 共享密钥 X-Internal-Secret 校验)。
- 注册 model(__init__) + router(main.py); 迁移 store_mapping_table 接现 head coin_txn_task_ref_uq。

验证: alembic 单 head + 零模型/迁移漂移; TestClient 全链(无密钥401/有密钥inserted=1/幂等inserted=0/错密钥401)。

Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
This commit is contained in:
2026-06-13 20:34:53 +08:00
parent 8d7b91219a
commit 4844e49161
7 changed files with 338 additions and 1 deletions
+77
View File
@@ -0,0 +1,77 @@
"""store_mapping table (平台店铺表:跨平台同店 id/名 映射,server 侧无条件落库)
Revision ID: store_mapping_table
Revises: coin_txn_task_ref_uq
Create Date: 2026-06-13 00:00:00.000000
"""
from typing import Sequence, Union
from alembic import op
import sqlalchemy as sa
# revision identifiers, used by Alembic.
revision: str = 'store_mapping_table'
down_revision: Union[str, Sequence[str], None] = 'coin_txn_task_ref_uq'
branch_labels: Union[str, Sequence[str], None] = None
depends_on: Union[str, Sequence[str], None] = None
def upgrade() -> None:
op.create_table(
'store_mapping',
sa.Column('id', sa.Integer(), autoincrement=True, nullable=False),
# 跨平台身份
sa.Column('id_taobao', sa.String(length=64), nullable=True),
sa.Column('name_taobao', sa.String(length=128), nullable=True),
sa.Column('id_meituan', sa.String(length=64), nullable=True),
sa.Column('name_meituan', sa.String(length=128), nullable=True),
sa.Column('id_jd', sa.String(length=64), nullable=True),
sa.Column('name_jd', sa.String(length=128), nullable=True),
# 地理
sa.Column('city', sa.String(length=64), nullable=True),
sa.Column('geohash', sa.String(length=16), nullable=True),
sa.Column('lng', sa.Float(), nullable=True),
sa.Column('lat', sa.Float(), nullable=True),
sa.Column('taobao_address', sa.String(length=256), nullable=True),
# 溯源
sa.Column('source_platform', sa.String(length=32), nullable=True),
sa.Column('business_type', sa.String(length=16), nullable=False),
sa.Column('trace_id', sa.String(length=64), nullable=False),
sa.Column('source_device_id', sa.String(length=64), nullable=True),
sa.Column('source_user_id', sa.Integer(), nullable=True),
# 淘宝原料(URL 可能很长 → Text)
sa.Column('taobao_share_url', sa.String(length=256), nullable=True),
sa.Column('taobao_resolved_url', sa.Text(), nullable=True),
sa.Column('taobao_deeplink', sa.Text(), nullable=True),
# PG 上为 JSONB,其它(SQLite)为 JSON——与模型层 with_variant 对齐
sa.Column('attrs', sa.JSON().with_variant(sa.dialects.postgresql.JSONB(), 'postgresql'), nullable=True),
sa.Column('created_at', sa.DateTime(timezone=True), server_default=sa.text('(CURRENT_TIMESTAMP)'), nullable=False),
sa.PrimaryKeyConstraint('id'),
# 一次比价一行,trace_id 幂等去重(防 pricebot 重试 / replay 重复写)
sa.UniqueConstraint('trace_id', name='uq_store_mapping_trace'),
)
with op.batch_alter_table('store_mapping', schema=None) as batch_op:
batch_op.create_index(batch_op.f('ix_store_mapping_id_taobao'), ['id_taobao'], unique=False)
batch_op.create_index(batch_op.f('ix_store_mapping_id_meituan'), ['id_meituan'], unique=False)
batch_op.create_index(batch_op.f('ix_store_mapping_id_jd'), ['id_jd'], unique=False)
batch_op.create_index(batch_op.f('ix_store_mapping_geohash'), ['geohash'], unique=False)
batch_op.create_index(batch_op.f('ix_store_mapping_source_platform'), ['source_platform'], unique=False)
batch_op.create_index(batch_op.f('ix_store_mapping_source_device_id'), ['source_device_id'], unique=False)
batch_op.create_index(batch_op.f('ix_store_mapping_source_user_id'), ['source_user_id'], unique=False)
batch_op.create_index(batch_op.f('ix_store_mapping_created_at'), ['created_at'], unique=False)
def downgrade() -> None:
with op.batch_alter_table('store_mapping', schema=None) as batch_op:
batch_op.drop_index(batch_op.f('ix_store_mapping_created_at'))
batch_op.drop_index(batch_op.f('ix_store_mapping_source_user_id'))
batch_op.drop_index(batch_op.f('ix_store_mapping_source_device_id'))
batch_op.drop_index(batch_op.f('ix_store_mapping_source_platform'))
batch_op.drop_index(batch_op.f('ix_store_mapping_geohash'))
batch_op.drop_index(batch_op.f('ix_store_mapping_id_jd'))
batch_op.drop_index(batch_op.f('ix_store_mapping_id_meituan'))
batch_op.drop_index(batch_op.f('ix_store_mapping_id_taobao'))
op.drop_table('store_mapping')
+46
View File
@@ -0,0 +1,46 @@
"""平台店铺映射内部上报端点(pricebot → app-server)。
pricebot 在淘宝比价拿到 shopId 后,把这一行跨平台店铺映射 POST 到这里落库。
**不是给客户端的接口**:不走用户 JWT,靠 server 间共享密钥头 `X-Internal-Secret` 校验
(复用 price.py 的 _check_secret,与 price-observation 同一密钥)。
与 price.py 的 /internal/price-observation 平行:那个落价格事实,这个落店铺身份映射。
"""
from __future__ import annotations
import logging
from typing import Annotated
from fastapi import APIRouter, Header
from app.api.deps import DbSession
from app.api.internal.price import _check_secret
from app.repositories import store_mapping as repo
from app.schemas.store_mapping import StoreMappingIn, StoreMappingOut
logger = logging.getLogger("shagua.internal.store")
router = APIRouter(prefix="/internal", tags=["internal"])
@router.post(
"/store-mapping",
response_model=StoreMappingOut,
summary="平台店铺映射内部上报(pricebot→app-server,落 store_mapping)",
)
def report_store_mapping(
payload: StoreMappingIn,
db: DbSession,
x_internal_secret: Annotated[str | None, Header()] = None,
) -> StoreMappingOut:
_check_secret(x_internal_secret)
inserted, row_id = repo.insert(db, payload)
logger.info(
"store_mapping trace=%s inserted=%d row_id=%s source=%s "
"taobao=(%s,%s) device=%s user=%s",
payload.trace_id, inserted, row_id, payload.source_platform,
payload.id_taobao, payload.name_taobao,
payload.source_device_id, payload.source_user_id,
)
return StoreMappingOut(inserted=inserted, row_id=row_id)
+3 -1
View File
@@ -21,6 +21,7 @@ from app.api.v1.compare_milestone import router as compare_milestone_router
from app.api.v1.compare_record import router as compare_record_router
from app.api.v1.coupon import router as coupon_router
from app.api.internal.price import router as internal_price_router
from app.api.internal.store import router as internal_store_router
from app.api.v1.feedback import router as feedback_router
from app.api.v1.invite import router as invite_router
from app.api.v1.meituan import router as meituan_router
@@ -102,8 +103,9 @@ app.include_router(savings_router)
app.include_router(ad_router)
app.include_router(order_router)
app.include_router(report_router)
# 内部(server→server)端点:pricebot 上报价格观测,靠共享密钥头校验,不对客户端开放。
# 内部(server→server)端点:pricebot 上报价格观测 / 店铺映射,靠共享密钥头校验,不对客户端开放。
app.include_router(internal_price_router)
app.include_router(internal_store_router)
app.include_router(platform_router)
# 用户上传文件(头像)静态服务。生产可改由 nginx 直接 serve MEDIA_ROOT。
+1
View File
@@ -23,6 +23,7 @@ from app.models.price_observation import PriceObservation # noqa: F401
from app.models.price_report import PriceReport # noqa: F401
from app.models.savings import SavingsRecord # noqa: F401
from app.models.signin import SigninBoostRecord, SigninRecord # noqa: F401
from app.models.store_mapping import StoreMapping # noqa: F401
from app.models.task import UserTask # noqa: F401
from app.models.user import User # noqa: F401
from app.models.wallet import ( # noqa: F401
+97
View File
@@ -0,0 +1,97 @@
"""平台店铺表(store_mapping)—— 跨平台"同一家店"的 id/名 映射资产层。
每完成一次淘宝比价(在目标淘宝店通过 更多操作→分享→复制链接 拿到分享短链、
HTTP 解析出 shopId 后),pricebot server→server 内部上报落这里一行。**与登录无关、
不依赖客户端鉴权**(比价透传链路当前不鉴权,user_id 客户端带上时一并记)。
与 price_observation 的区别:
- price_observation:平台/门店视角的**价格事实**(某店这单多少钱)。
- store_mapping:平台/门店视角的**身份映射**(同一家物理店在 淘宝/美团/京东 各自的
店铺 id 与店名)。是未来"我见过这家店→跳过重新搜索/匹配"的源头。两表独立。
「先存下来、用法后说」:列尽量铺全(各平台 id/名 + 地理 + 溯源 + 淘宝原料 URL),
attrs(JSONB)兜底存灵活明细,免得每多记一个字段就迁移 schema。
⚠️ 数据质量:跨平台"同一家店"的连接来自 agent 的 LLM 店铺匹配,匹配错则一行里连错店。
本表是 append-only 原始记录(每比价一行、trace_id 幂等防重试重复),清洗/归一二期再做。
当前只接通**淘宝**(id_taobao 才有真值;美团/京东多为 name-only 或空)。
"""
from __future__ import annotations
from datetime import datetime
from sqlalchemy import (
JSON,
DateTime,
Float,
Integer,
String,
Text,
UniqueConstraint,
func,
)
from sqlalchemy.dialects.postgresql import JSONB
from sqlalchemy.orm import Mapped, mapped_column
from app.db.base import Base
# PG 上用 JSONB,SQLite(本地/测试)退化为通用 JSON(同 price_observation / comparison_record)。
_JSON = JSON().with_variant(JSONB(), "postgresql")
class StoreMapping(Base):
__tablename__ = "store_mapping"
__table_args__ = (
# 一次比价(trace)只记一条:pricebot 重试 / 客户端 replay 重复上报时幂等去重。
# 一次淘宝比价 = 一个目标淘宝店 → 一行映射(源 + 各平台身份压在同一行)。
UniqueConstraint("trace_id", name="uq_store_mapping_trace"),
)
id: Mapped[int] = mapped_column(Integer, primary_key=True, autoincrement=True)
# ===== 跨平台店铺身份(同一家物理店在各平台的 id/名;按比价角色稀疏填充)=====
# id_taobao = 分享短链解析出的 shopId(淘宝当目标、走完取 id 流程才有);
# name_taobao = 店铺页 a11y content_desc "店铺标题:xxx" 剥前缀。
id_taobao: Mapped[str | None] = mapped_column(String(64), index=True, nullable=True)
name_taobao: Mapped[str | None] = mapped_column(String(128), nullable=True)
# 美团/京东:当前无同款 share→id 机制,多为 name-only(源平台店名来自 intent / agent 匹配名)。
id_meituan: Mapped[str | None] = mapped_column(String(64), index=True, nullable=True)
name_meituan: Mapped[str | None] = mapped_column(String(128), nullable=True)
id_jd: Mapped[str | None] = mapped_column(String(64), index=True, nullable=True)
name_jd: Mapped[str | None] = mapped_column(String(128), nullable=True)
# ===== 地理(同名店异地区分 / 地理分桶匹配的主要燃料)=====
city: Mapped[str | None] = mapped_column(String(64), nullable=True)
geohash: Mapped[str | None] = mapped_column(String(16), index=True, nullable=True)
lng: Mapped[float | None] = mapped_column(Float, nullable=True)
lat: Mapped[float | None] = mapped_column(Float, nullable=True)
# 淘宝门店地址(店铺页 a11y 抓到才有;比经纬度更利于人工/LLM 匹配)
taobao_address: Mapped[str | None] = mapped_column(String(256), nullable=True)
# ===== 溯源 / 用户画像 =====
# 源平台(发起比价那家:meituan / taobao_flash / jd_waimai ...)
source_platform: Mapped[str | None] = mapped_column(String(32), index=True, nullable=True)
business_type: Mapped[str] = mapped_column(String(16), nullable=False, default="food")
# pricebot 侧 trace_id:回指原始 trace(溯源)+ 幂等去重键(uq_store_mapping_trace)
trace_id: Mapped[str] = mapped_column(String(64), nullable=False)
source_device_id: Mapped[str | None] = mapped_column(String(64), index=True, nullable=True)
# user_id 当前比价链路不鉴权拿不到,客户端带上时才有;先可空。
source_user_id: Mapped[int | None] = mapped_column(Integer, index=True, nullable=True)
# ===== 淘宝原料(可复跳 / 可重解析 / 调试;URL 可能很长 → Text)=====
taobao_share_url: Mapped[str | None] = mapped_column(String(256), nullable=True) # m.tb.cn 短链
taobao_resolved_url: Mapped[str | None] = mapped_column(Text, nullable=True) # 解析出的目标 URL(含 shopId)
taobao_deeplink: Mapped[str | None] = mapped_column(Text, nullable=True) # 拼好的 et-store/search deeplink
# 灵活字段兜底(免得加字段就迁移)
attrs: Mapped[dict | None] = mapped_column(_JSON, nullable=True)
created_at: Mapped[datetime] = mapped_column(
DateTime(timezone=True), server_default=func.now(), index=True, nullable=False
)
def __repr__(self) -> str: # pragma: no cover
return (
f"<StoreMapping id={self.id} taobao=({self.id_taobao!r},{self.name_taobao!r}) "
f"source={self.source_platform} trace_id={self.trace_id}>"
)
+65
View File
@@ -0,0 +1,65 @@
"""平台店铺映射落库:一次比价一行,trace_id 幂等。
幂等键 = trace_id。pricebot 重试 / replay 会重复上报同一 trace,这里先查该 trace 是否
已存在,存在即跳过(返回已存在行 id)——跨方言(PG / SQLite dev)都安全,不依赖 ON CONFLICT。
极端并发下若仍撞唯一约束,IntegrityError 回滚后再查一次返回已存在行(本就幂等,不重试)。
"""
from __future__ import annotations
import logging
from sqlalchemy import select
from sqlalchemy.exc import IntegrityError
from sqlalchemy.orm import Session
from app.models.store_mapping import StoreMapping
from app.schemas.store_mapping import StoreMappingIn
logger = logging.getLogger("shagua.store_mapping")
def _existing_id(db: Session, trace_id: str) -> int | None:
return db.execute(
select(StoreMapping.id).where(StoreMapping.trace_id == trace_id)
).scalar_one_or_none()
def insert(db: Session, payload: StoreMappingIn) -> tuple[int, int | None]:
"""写入一行映射,返回 (inserted, row_id)。inserted=1 新写入 / 0 幂等跳过。"""
existing = _existing_id(db, payload.trace_id)
if existing is not None:
return 0, existing
row = StoreMapping(
trace_id=payload.trace_id,
business_type=payload.business_type,
source_platform=payload.source_platform,
id_taobao=payload.id_taobao,
name_taobao=payload.name_taobao,
id_meituan=payload.id_meituan,
name_meituan=payload.name_meituan,
id_jd=payload.id_jd,
name_jd=payload.name_jd,
city=payload.city,
geohash=payload.geohash,
lng=payload.lng,
lat=payload.lat,
taobao_address=payload.taobao_address,
source_device_id=payload.source_device_id,
source_user_id=payload.source_user_id,
taobao_share_url=payload.taobao_share_url,
taobao_resolved_url=payload.taobao_resolved_url,
taobao_deeplink=payload.taobao_deeplink,
attrs=payload.attrs,
)
db.add(row)
try:
db.commit()
except IntegrityError:
# 并发下另一个请求刚插了同 trace → 唯一约束撞了。回滚后返回已存在行(幂等)。
db.rollback()
logger.warning("store_mapping 并发幂等冲突 trace=%s,返回已存在行", payload.trace_id)
return 0, _existing_id(db, payload.trace_id)
db.refresh(row)
return 1, row.id
+49
View File
@@ -0,0 +1,49 @@
"""平台店铺映射内部上报的收发模型。
pricebot 在淘宝比价拿到 shopId 后 server→server POST 一行映射。所有字段可空(按比价
角色稀疏填充),server 端只校验共享密钥 + 幂等(trace_id)落库。
"""
from __future__ import annotations
from pydantic import BaseModel
class StoreMappingIn(BaseModel):
"""一次比价的跨平台店铺映射上报体(一行)。"""
trace_id: str
business_type: str = "food"
source_platform: str | None = None
# 跨平台身份(按角色稀疏填充)
id_taobao: str | None = None
name_taobao: str | None = None
id_meituan: str | None = None
name_meituan: str | None = None
id_jd: str | None = None
name_jd: str | None = None
# 地理
city: str | None = None
geohash: str | None = None
lng: float | None = None
lat: float | None = None
taobao_address: str | None = None
# 溯源
source_device_id: str | None = None
source_user_id: int | None = None
# 淘宝原料(可复跳 / 可重解析 / 调试)
taobao_share_url: str | None = None
taobao_resolved_url: str | None = None
taobao_deeplink: str | None = None
attrs: dict | None = None
class StoreMappingOut(BaseModel):
"""上报结果。inserted=1 为新写入,0 为幂等跳过(trace 已存在);row_id 为该 trace 对应行 id。"""
inserted: int
row_id: int | None = None