From 4844e491611325f1e540196027a4c013d1fd5523 Mon Sep 17 00:00:00 2001 From: marco Date: Sat, 13 Jun 2026 20:34:53 +0800 Subject: [PATCH] =?UTF-8?q?feat(store=5Fmapping):=20=E8=B7=A8=E5=B9=B3?= =?UTF-8?q?=E5=8F=B0=E5=BA=97=E9=93=BA=E6=98=A0=E5=B0=84=E8=A1=A8=20+=20pr?= =?UTF-8?q?icebot=20=E5=86=85=E9=83=A8=E4=B8=8A=E6=8A=A5=E7=AB=AF=E7=82=B9?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 新增「平台店铺表」资产层: 淘宝比价拿到 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 --- alembic/versions/store_mapping_table.py | 77 ++++++++++++++++++++ app/api/internal/store.py | 46 ++++++++++++ app/main.py | 4 +- app/models/__init__.py | 1 + app/models/store_mapping.py | 97 +++++++++++++++++++++++++ app/repositories/store_mapping.py | 65 +++++++++++++++++ app/schemas/store_mapping.py | 49 +++++++++++++ 7 files changed, 338 insertions(+), 1 deletion(-) create mode 100644 alembic/versions/store_mapping_table.py create mode 100644 app/api/internal/store.py create mode 100644 app/models/store_mapping.py create mode 100644 app/repositories/store_mapping.py create mode 100644 app/schemas/store_mapping.py diff --git a/alembic/versions/store_mapping_table.py b/alembic/versions/store_mapping_table.py new file mode 100644 index 0000000..7e3e709 --- /dev/null +++ b/alembic/versions/store_mapping_table.py @@ -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') diff --git a/app/api/internal/store.py b/app/api/internal/store.py new file mode 100644 index 0000000..91ef749 --- /dev/null +++ b/app/api/internal/store.py @@ -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) diff --git a/app/main.py b/app/main.py index fe296f2..cea4568 100644 --- a/app/main.py +++ b/app/main.py @@ -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。 diff --git a/app/models/__init__.py b/app/models/__init__.py index 4ed0d1b..a089a8a 100644 --- a/app/models/__init__.py +++ b/app/models/__init__.py @@ -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 diff --git a/app/models/store_mapping.py b/app/models/store_mapping.py new file mode 100644 index 0000000..632f5b5 --- /dev/null +++ b/app/models/store_mapping.py @@ -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"" + ) diff --git a/app/repositories/store_mapping.py b/app/repositories/store_mapping.py new file mode 100644 index 0000000..f6566e8 --- /dev/null +++ b/app/repositories/store_mapping.py @@ -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 diff --git a/app/schemas/store_mapping.py b/app/schemas/store_mapping.py new file mode 100644 index 0000000..4143fe2 --- /dev/null +++ b/app/schemas/store_mapping.py @@ -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