Compare commits
3 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| ea27fe0c01 | |||
| 995908cd0b | |||
| fb661814cb |
+1
-13
@@ -11,7 +11,7 @@ DATABASE_URL=sqlite:///./data/app.db
|
||||
|
||||
# ===== JWT =====
|
||||
# 生产部署务必改成随机长字符串,可用:python -c "import secrets; print(secrets.token_urlsafe(64))"
|
||||
JWT_SECRET_KEY=change-me-in-prod-please-use-a-long-random-string
|
||||
JWT_SECRET_KEY=
|
||||
JWT_ALGORITHM=HS256
|
||||
# access token 有效期(分钟),默认 2 小时
|
||||
JWT_ACCESS_TOKEN_EXPIRE_MINUTES=120
|
||||
@@ -27,18 +27,6 @@ JG_PRIVATE_KEY_PATH=./secrets/jverify_rsa_private.pem
|
||||
JG_VERIFY_ENDPOINT=https://api.verification.jpush.cn/v1/web/loginTokenVerify
|
||||
JG_REQUEST_TIMEOUT_SEC=15
|
||||
|
||||
# ===== 极光推送 JPush(无障碍保护掉线告警)=====
|
||||
# 客户端 push appkey 已知 = 966b451a8d9cfe12d173ea9d。若推送与一键登录/短信是同一个极光应用,
|
||||
# JPUSH_* 留空即自动回退到上面的 JG_APP_KEY/JG_MASTER_SECRET;否则单独填那个 push 应用的密钥。
|
||||
# 运维清单(厂商通道等)见 spec/accessibility-liveness-push.md §7。
|
||||
JPUSH_APP_KEY=
|
||||
JPUSH_MASTER_SECRET=
|
||||
JPUSH_PUSH_ENDPOINT=https://api.jpush.cn/v3/push
|
||||
# 无障碍保护存活监控后台任务
|
||||
HEARTBEAT_MONITOR_ENABLED=true
|
||||
HEARTBEAT_TIMEOUT_MINUTES=10
|
||||
HEARTBEAT_SCAN_INTERVAL_SEC=60
|
||||
|
||||
# ===== 短信 (mock 模式) =====
|
||||
# mock = true 时,任意 6 位数字均通过,且 /sms/send 不真发短信(只 log)。
|
||||
# 后续接阿里云/腾讯云短信时,改成 false 并填供应商相关 key。
|
||||
|
||||
@@ -1,26 +0,0 @@
|
||||
"""merge coupon_engage_app_pkg and coin_txn_task_ref_uq heads
|
||||
|
||||
Revision ID: 3a9941e76909
|
||||
Revises: coin_txn_task_ref_uq, coupon_engage_app_pkg
|
||||
Create Date: 2026-06-13 10:09:37.557466
|
||||
|
||||
"""
|
||||
from typing import Sequence, Union
|
||||
|
||||
from alembic import op
|
||||
import sqlalchemy as sa
|
||||
|
||||
|
||||
# revision identifiers, used by Alembic.
|
||||
revision: str = '3a9941e76909'
|
||||
down_revision: Union[str, Sequence[str], None] = ('coin_txn_task_ref_uq', 'coupon_engage_app_pkg')
|
||||
branch_labels: Union[str, Sequence[str], None] = None
|
||||
depends_on: Union[str, Sequence[str], None] = None
|
||||
|
||||
|
||||
def upgrade() -> None:
|
||||
pass
|
||||
|
||||
|
||||
def downgrade() -> None:
|
||||
pass
|
||||
@@ -1,47 +0,0 @@
|
||||
"""coupon_prompt_engagement 加 app_package(弹窗频控改按 App 为单位)
|
||||
|
||||
2026-06-12 产品确认:领券引导窗「一个 app 一天只能弹一次」,没领完进其他平台要再弹,
|
||||
彻底领完(coupon_daily_completion)才全局不弹。频控行从 (device, 日) 唯一改为
|
||||
(device, app_package, 日) 唯一;旧数据 app_package 回填 NULL(= 全局兜底行,
|
||||
按 App 查询时忽略,不带包名的旧式查询仍生效)。
|
||||
|
||||
⚠️ downgrade 有损:同设备同日多 App 各一行时,重建 (device, 日) 唯一约束会撞;
|
||||
仅开发/测试库可降级(频控行本就是当日临时数据,coupon_claim_record 资产不受影响)。
|
||||
|
||||
Revision ID: coupon_engage_app_pkg
|
||||
Revises: 9b894f5fff05
|
||||
Create Date: 2026-06-12 00:00:00.000000
|
||||
|
||||
"""
|
||||
|
||||
from collections.abc import Sequence
|
||||
|
||||
import sqlalchemy as sa
|
||||
|
||||
from alembic import op
|
||||
|
||||
# revision identifiers, used by Alembic.
|
||||
revision: str = "coupon_engage_app_pkg"
|
||||
down_revision: str | Sequence[str] | None = "9b894f5fff05"
|
||||
branch_labels: str | Sequence[str] | None = None
|
||||
depends_on: str | Sequence[str] | None = None
|
||||
|
||||
|
||||
def upgrade() -> None:
|
||||
# batch 模式兼容 SQLite(本地)与 PG(线上):SQLite 改约束要重建表,batch 自动处理。
|
||||
with op.batch_alter_table("coupon_prompt_engagement", schema=None) as batch_op:
|
||||
batch_op.add_column(sa.Column("app_package", sa.String(length=64), nullable=True))
|
||||
batch_op.drop_constraint("uq_coupon_engage_device_date", type_="unique")
|
||||
batch_op.create_unique_constraint(
|
||||
"uq_coupon_engage_device_app_date",
|
||||
["device_id", "app_package", "engage_date"],
|
||||
)
|
||||
|
||||
|
||||
def downgrade() -> None:
|
||||
with op.batch_alter_table("coupon_prompt_engagement", schema=None) as batch_op:
|
||||
batch_op.drop_constraint("uq_coupon_engage_device_app_date", type_="unique")
|
||||
batch_op.create_unique_constraint(
|
||||
"uq_coupon_engage_device_date", ["device_id", "engage_date"]
|
||||
)
|
||||
batch_op.drop_column("app_package")
|
||||
@@ -1,53 +0,0 @@
|
||||
"""device table (无障碍保护存活检测 + 极光推送)
|
||||
|
||||
Revision ID: device_liveness_table
|
||||
Revises: f3d0a16bb4c2
|
||||
Create Date: 2026-06-15 12:00:00.000000
|
||||
|
||||
"""
|
||||
from typing import Sequence, Union
|
||||
|
||||
from alembic import op
|
||||
import sqlalchemy as sa
|
||||
|
||||
|
||||
# revision identifiers, used by Alembic.
|
||||
revision: str = 'device_liveness_table'
|
||||
down_revision: Union[str, Sequence[str], None] = 'f3d0a16bb4c2'
|
||||
branch_labels: Union[str, Sequence[str], None] = None
|
||||
depends_on: Union[str, Sequence[str], None] = None
|
||||
|
||||
|
||||
def upgrade() -> None:
|
||||
op.create_table(
|
||||
'device',
|
||||
sa.Column('id', sa.Integer(), autoincrement=True, nullable=False),
|
||||
sa.Column('user_id', sa.Integer(), nullable=False),
|
||||
sa.Column('device_id', sa.String(length=128), nullable=False),
|
||||
sa.Column('registration_id', sa.String(length=64), nullable=True),
|
||||
sa.Column('platform', sa.String(length=16), nullable=False),
|
||||
sa.Column('app_version', sa.String(length=32), nullable=True),
|
||||
sa.Column('ever_protected', sa.Boolean(), nullable=False),
|
||||
sa.Column('last_heartbeat_at', sa.DateTime(timezone=True), nullable=True),
|
||||
sa.Column('last_report_protection_on', sa.Boolean(), nullable=False),
|
||||
sa.Column('liveness_state', sa.String(length=16), nullable=False),
|
||||
sa.Column('notified_at', sa.DateTime(timezone=True), nullable=True),
|
||||
sa.Column('created_at', sa.DateTime(timezone=True), server_default=sa.text('(CURRENT_TIMESTAMP)'), nullable=False),
|
||||
sa.Column('updated_at', sa.DateTime(timezone=True), server_default=sa.text('(CURRENT_TIMESTAMP)'), nullable=False),
|
||||
sa.ForeignKeyConstraint(['user_id'], ['user.id'], ),
|
||||
sa.PrimaryKeyConstraint('id'),
|
||||
sa.UniqueConstraint('user_id', 'device_id', name='uq_device_user_device'),
|
||||
)
|
||||
with op.batch_alter_table('device', schema=None) as batch_op:
|
||||
batch_op.create_index(batch_op.f('ix_device_user_id'), ['user_id'], unique=False)
|
||||
batch_op.create_index(batch_op.f('ix_device_device_id'), ['device_id'], unique=False)
|
||||
batch_op.create_index(batch_op.f('ix_device_last_heartbeat_at'), ['last_heartbeat_at'], unique=False)
|
||||
|
||||
|
||||
def downgrade() -> None:
|
||||
with op.batch_alter_table('device', schema=None) as batch_op:
|
||||
batch_op.drop_index(batch_op.f('ix_device_last_heartbeat_at'))
|
||||
batch_op.drop_index(batch_op.f('ix_device_device_id'))
|
||||
batch_op.drop_index(batch_op.f('ix_device_user_id'))
|
||||
|
||||
op.drop_table('device')
|
||||
@@ -1,26 +0,0 @@
|
||||
"""merge store_mapping cols branch into coupon/coin mergepoint
|
||||
|
||||
Revision ID: f3d0a16bb4c2
|
||||
Revises: 3a9941e76909, store_mapping_jd_cols
|
||||
Create Date: 2026-06-14 14:31:25.504510
|
||||
|
||||
"""
|
||||
from typing import Sequence, Union
|
||||
|
||||
from alembic import op
|
||||
import sqlalchemy as sa
|
||||
|
||||
|
||||
# revision identifiers, used by Alembic.
|
||||
revision: str = 'f3d0a16bb4c2'
|
||||
down_revision: Union[str, Sequence[str], None] = ('3a9941e76909', 'store_mapping_jd_cols')
|
||||
branch_labels: Union[str, Sequence[str], None] = None
|
||||
depends_on: Union[str, Sequence[str], None] = None
|
||||
|
||||
|
||||
def upgrade() -> None:
|
||||
pass
|
||||
|
||||
|
||||
def downgrade() -> None:
|
||||
pass
|
||||
@@ -1,40 +0,0 @@
|
||||
"""store_mapping 加京东原料列(jd_vender_id + 分享链/反查 URL/deeplink)
|
||||
|
||||
Revision ID: store_mapping_jd_cols
|
||||
Revises: store_mapping_meituan_cols
|
||||
Create Date: 2026-06-13 00:00:00.000000
|
||||
|
||||
京东秒送接入店内搜索 deeplink 链路: 3.cn 短链 → 跟随重定向反查 venderId+storeId →
|
||||
openapp.jdmobile:// deeplink。京东店铺身份是**两个**稳定数字 id: storeId 进 id_jd(稳定
|
||||
店主键, 同 taobao 的 shopId→id_taobao), venderId 进单列 jd_vender_id(deeplink 还需它)。
|
||||
其余三列与 taobao_*/meituan_* 原料列平行。
|
||||
"""
|
||||
from typing import Sequence, Union
|
||||
|
||||
from alembic import op
|
||||
import sqlalchemy as sa
|
||||
|
||||
|
||||
# revision identifiers, used by Alembic.
|
||||
revision: str = 'store_mapping_jd_cols'
|
||||
down_revision: Union[str, Sequence[str], None] = 'store_mapping_meituan_cols'
|
||||
branch_labels: Union[str, Sequence[str], None] = None
|
||||
depends_on: Union[str, Sequence[str], None] = None
|
||||
|
||||
|
||||
def upgrade() -> None:
|
||||
with op.batch_alter_table('store_mapping', schema=None) as batch_op:
|
||||
batch_op.add_column(sa.Column('jd_vender_id', sa.String(length=64), nullable=True))
|
||||
batch_op.add_column(sa.Column('jd_share_url', sa.String(length=256), nullable=True))
|
||||
batch_op.add_column(sa.Column('jd_resolved_url', sa.Text(), nullable=True))
|
||||
batch_op.add_column(sa.Column('jd_deeplink', sa.Text(), nullable=True))
|
||||
batch_op.create_index(batch_op.f('ix_store_mapping_jd_vender_id'), ['jd_vender_id'], 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_jd_vender_id'))
|
||||
batch_op.drop_column('jd_deeplink')
|
||||
batch_op.drop_column('jd_resolved_url')
|
||||
batch_op.drop_column('jd_share_url')
|
||||
batch_op.drop_column('jd_vender_id')
|
||||
@@ -1,39 +0,0 @@
|
||||
"""store_mapping 加美团原料列(poi_id_str + 分享链/反查 URL/deeplink)
|
||||
|
||||
Revision ID: store_mapping_meituan_cols
|
||||
Revises: store_mapping_table
|
||||
Create Date: 2026-06-13 00:00:00.000000
|
||||
|
||||
美团接入店内搜索 deeplink 链路: dpurl.cn 短链 → 302 反查 poi_id_str → imeituan:// deeplink。
|
||||
poi_id_str 每次分享重新加密、非稳定主键, 单列 meituan_poi_id_str 存(不占 id_meituan,
|
||||
后者留给将来 CPS API 的稳定数字 poi_id)。其余三列与 taobao_* 原料列平行。
|
||||
"""
|
||||
from typing import Sequence, Union
|
||||
|
||||
from alembic import op
|
||||
import sqlalchemy as sa
|
||||
|
||||
|
||||
# revision identifiers, used by Alembic.
|
||||
revision: str = 'store_mapping_meituan_cols'
|
||||
down_revision: Union[str, Sequence[str], None] = 'store_mapping_table'
|
||||
branch_labels: Union[str, Sequence[str], None] = None
|
||||
depends_on: Union[str, Sequence[str], None] = None
|
||||
|
||||
|
||||
def upgrade() -> None:
|
||||
with op.batch_alter_table('store_mapping', schema=None) as batch_op:
|
||||
batch_op.add_column(sa.Column('meituan_poi_id_str', sa.String(length=64), nullable=True))
|
||||
batch_op.add_column(sa.Column('meituan_share_url', sa.String(length=256), nullable=True))
|
||||
batch_op.add_column(sa.Column('meituan_resolved_url', sa.Text(), nullable=True))
|
||||
batch_op.add_column(sa.Column('meituan_deeplink', sa.Text(), nullable=True))
|
||||
batch_op.create_index(batch_op.f('ix_store_mapping_meituan_poi_id_str'), ['meituan_poi_id_str'], 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_meituan_poi_id_str'))
|
||||
batch_op.drop_column('meituan_deeplink')
|
||||
batch_op.drop_column('meituan_resolved_url')
|
||||
batch_op.drop_column('meituan_share_url')
|
||||
batch_op.drop_column('meituan_poi_id_str')
|
||||
@@ -1,77 +0,0 @@
|
||||
"""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')
|
||||
@@ -1,79 +0,0 @@
|
||||
"""平台店铺映射内部上报端点(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.get(
|
||||
"/store-mapping/lookup",
|
||||
summary="比价前按源平台店名反查各目标平台已沉淀的店铺 id(命中→pricebot 直接 deeplink)",
|
||||
)
|
||||
def lookup_store_mapping(
|
||||
source_platform: str,
|
||||
name: str,
|
||||
db: DbSession,
|
||||
lat: float | None = None,
|
||||
lng: float | None = None,
|
||||
x_internal_secret: Annotated[str | None, Header()] = None,
|
||||
) -> dict:
|
||||
_check_secret(x_internal_secret)
|
||||
result = repo.lookup_nearest(db, source_platform, name, lat, lng)
|
||||
if result:
|
||||
hits = ", ".join(
|
||||
f"{t}:row{v['row_id']}"
|
||||
f"{'(' + str(v['dist_km']) + 'km)' if 'dist_km' in v else ''}"
|
||||
f"→{v.get('deeplink') or '(无deeplink)'}"
|
||||
for t, v in result.items()
|
||||
)
|
||||
logger.info(
|
||||
"store_mapping lookup source=%s name=%r geo=(%s,%s) → 命中 %s",
|
||||
source_platform, name, lat, lng, hits,
|
||||
)
|
||||
else:
|
||||
logger.info(
|
||||
"store_mapping lookup source=%s name=%r geo=(%s,%s) → MISS",
|
||||
source_platform, name, lat, lng,
|
||||
)
|
||||
return result
|
||||
|
||||
|
||||
@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)
|
||||
created, row_id = repo.upsert(db, payload)
|
||||
logger.info(
|
||||
"store_mapping trace=%s %s row_id=%s source=%s "
|
||||
"taobao=(%s,%s) jd=(%s,%s) device=%s user=%s",
|
||||
payload.trace_id, "新建" if created else "合并", row_id, payload.source_platform,
|
||||
payload.id_taobao, payload.name_taobao, payload.id_jd, payload.name_jd,
|
||||
payload.source_device_id, payload.source_user_id,
|
||||
)
|
||||
return StoreMappingOut(inserted=created, row_id=row_id)
|
||||
+12
-44
@@ -27,7 +27,6 @@ from app.schemas.coupon_state import (
|
||||
CouponCompletedTodayOut,
|
||||
CouponPromptDismissIn,
|
||||
CouponPromptShouldShowOut,
|
||||
CouponPromptShownIn,
|
||||
)
|
||||
|
||||
logger = logging.getLogger("shagua.coupon")
|
||||
@@ -113,10 +112,8 @@ async def coupon_step(
|
||||
user_id = _to_int(meta.get("user_id")) # 登录态才带;判断不靠它,资产留痕用
|
||||
trace_id = meta.get("trace_id")
|
||||
|
||||
# 领券任务首帧(step=0)= 用户已发起领券 → 记一条今日 engagement(claim_started)。
|
||||
# 透传链路拿不到发起 App 的包名 → 写 app_package=NULL 的全局兜底行:只对旧式不带
|
||||
# package 的 should-show 查询生效;按 App 频控(2026-06-12)由客户端 /prompt/shown
|
||||
# 负责,本行不堵其他 App 的弹窗(领一半终止,其他平台还要弹)。写库失败绝不能
|
||||
# 领券任务首帧(step=0)= 用户已发起领券 → 记一条今日 engagement(claim_started),
|
||||
# 今天这台设备不再弹引导窗(对齐前台"点一键领取即 markEngaged")。写库失败绝不能
|
||||
# 连累领券主流程,整段吞掉。
|
||||
if device_id and meta.get("step") == 0:
|
||||
try:
|
||||
@@ -191,66 +188,37 @@ async def coupon_step(
|
||||
return resp_json
|
||||
|
||||
|
||||
@router.post("/prompt/shown", summary="领券引导窗已对某 App 弹出(按 App 频控主判据)")
|
||||
def coupon_prompt_shown(payload: CouponPromptShownIn, db: DbSession) -> dict[str, bool]:
|
||||
"""客户端弹窗一亮即调用 → 记一条今日该 App 的 engagement(shown)。
|
||||
|
||||
2026-06-12 频控按 App 为单位(方案文档「一个 app 一天只能弹一次,时机=每天第一次
|
||||
进入」):以"弹出"为频控锚点,用户领/拒/无视都只算这一次;dismiss/claim_started
|
||||
后续只升级同一行的 engage_type。MVP 不鉴权,按 device_id 记。
|
||||
"""
|
||||
coupon_repo.mark_engagement(
|
||||
db, payload.device_id, payload.user_id, "shown", payload.package
|
||||
)
|
||||
return {"ok": True}
|
||||
|
||||
|
||||
@router.post("/prompt/dismiss", summary="用户拒绝/关闭领券引导窗(记今日已 engage)")
|
||||
def coupon_prompt_dismiss(payload: CouponPromptDismissIn, db: DbSession) -> dict[str, bool]:
|
||||
"""客户端点关闭引导窗时调用 → 记一条今日 engagement(dismissed),今天该 App 不再弹。
|
||||
"""客户端点关闭引导窗时调用 → 记一条今日 engagement(dismissed),今天不再弹。
|
||||
|
||||
server 在透传链路里看不到"用户拒绝"(拒绝不发起领券),故必须客户端通知。
|
||||
频控以 /prompt/shown 为主判据,本端点把同一行升级成 dismissed(记录用);
|
||||
旧客户端不带 package → 写 NULL 全局兜底行(旧语义)。MVP 不鉴权,按 device_id 记。
|
||||
MVP 不鉴权,按 device_id 记。
|
||||
"""
|
||||
coupon_repo.mark_engagement(
|
||||
db, payload.device_id, payload.user_id, "dismissed", payload.package
|
||||
)
|
||||
coupon_repo.mark_engagement(db, payload.device_id, payload.user_id, "dismissed")
|
||||
return {"ok": True}
|
||||
|
||||
|
||||
@router.get(
|
||||
"/prompt/should-show",
|
||||
response_model=CouponPromptShouldShowOut,
|
||||
summary="切到某外卖 App 时是否还应弹领券引导窗(按 App 频控)",
|
||||
summary="切到外卖 App 时是否还应弹领券引导窗",
|
||||
)
|
||||
def coupon_prompt_should_show(
|
||||
device_id: str, db: DbSession, package: str | None = None
|
||||
device_id: str, db: DbSession
|
||||
) -> CouponPromptShouldShowOut:
|
||||
"""should_show = 今天没跑完整轮领券 AND 该 App 今天没弹过。
|
||||
|
||||
2026-06-12 按 App 为单位(产品确认):
|
||||
- 该 App 今天弹过(shown/dismissed/claim_started 任一)→ false;
|
||||
- **其他 App** 弹过/领了一半不影响本 App → 仍 true(没领完进其他平台要再弹);
|
||||
- 今天已跑完整轮(coupon_daily_completion,"彻底领完")→ 全局 false。
|
||||
package 不传 = 旧客户端,退回旧全局语义(任意一行算 engage)+ 同样吃 completion 闸。
|
||||
客户端据此决定弹不弹(纯后台判据,客户端不做前台 SP 缓存判断)。"""
|
||||
if coupon_repo.has_completed_today(db, device_id):
|
||||
return CouponPromptShouldShowOut(should_show=False)
|
||||
"""今天这台设备已 engage(领或拒)过 → should_show=false。客户端据此决定弹不弹
|
||||
(纯后台判据,客户端不再做前台 SP 缓存判断)。"""
|
||||
return CouponPromptShouldShowOut(
|
||||
should_show=not coupon_repo.has_engaged_today(db, device_id, package)
|
||||
should_show=not coupon_repo.has_engaged_today(db, device_id)
|
||||
)
|
||||
|
||||
|
||||
@router.post("/prompt/reset", summary="重置今日领券状态:弹窗 engagement + 完成记录(开发测试用)")
|
||||
@router.post("/prompt/reset", summary="重置今日领券引导窗 engagement(开发测频控用)")
|
||||
def coupon_prompt_reset(payload: CouponPromptDismissIn, db: DbSession) -> dict[str, bool]:
|
||||
"""删这台设备今天的 engagement + coupon_daily_completion:弹窗又能弹,且首页
|
||||
「去领取」卡/弹窗「一键自动领取」CTA 恢复可点(只删 engagement 的话, 重置后 CTA
|
||||
仍被"今日已跑完整轮"置灰, 2026-06-12)。配合客户端本地重置 = 等效重装。
|
||||
领券记录(coupon_claim_record)是资产沉淀、不参与任何门控判断, 不删。
|
||||
"""删这台设备今天的 engagement → has_engaged_today 变 false,今天又能弹。
|
||||
开发设置「重置今日领券弹窗状态」按钮调。MVP 不鉴权,按 device_id。"""
|
||||
coupon_repo.reset_today_engagement(db, payload.device_id)
|
||||
coupon_repo.reset_today_completion(db, payload.device_id)
|
||||
return {"ok": True}
|
||||
|
||||
|
||||
|
||||
@@ -1,66 +0,0 @@
|
||||
"""设备注册 / 心跳 endpoint(无障碍保护存活检测)。
|
||||
|
||||
路由前缀 /api/v1/device,需 Bearer 鉴权(设备绑登录用户)。
|
||||
POST /register 注册设备 / 更新 registration_id(App 前台、拿到 push token 时调)
|
||||
POST /heartbeat 上报心跳(无障碍服务存活时周期调,刷新存活)
|
||||
|
||||
后端 heartbeat_monitor_worker 据此发现心跳超时的设备并极光推送告警。
|
||||
见 spec: spec/accessibility-liveness-push.md。
|
||||
"""
|
||||
from __future__ import annotations
|
||||
|
||||
import logging
|
||||
|
||||
from fastapi import APIRouter
|
||||
|
||||
from app.api.deps import CurrentUser, DbSession
|
||||
from app.repositories import device as device_repo
|
||||
from app.schemas.device import (
|
||||
DeviceOut,
|
||||
DeviceRegisterRequest,
|
||||
HeartbeatRequest,
|
||||
OkResponse,
|
||||
)
|
||||
|
||||
logger = logging.getLogger("shagua.device")
|
||||
|
||||
router = APIRouter(prefix="/api/v1/device", tags=["device"])
|
||||
|
||||
|
||||
@router.post("/register", response_model=DeviceOut, summary="注册设备/更新推送token")
|
||||
def register_device(
|
||||
req: DeviceRegisterRequest,
|
||||
user: CurrentUser,
|
||||
db: DbSession,
|
||||
) -> DeviceOut:
|
||||
device = device_repo.register_or_update(
|
||||
db,
|
||||
user_id=user.id,
|
||||
device_id=req.device_id,
|
||||
registration_id=req.registration_id,
|
||||
platform=req.platform,
|
||||
app_version=req.app_version,
|
||||
)
|
||||
logger.info(
|
||||
"device register user_id=%d device_id=%s reg=%s",
|
||||
user.id,
|
||||
req.device_id,
|
||||
bool(req.registration_id),
|
||||
)
|
||||
return DeviceOut.model_validate(device)
|
||||
|
||||
|
||||
@router.post("/heartbeat", response_model=OkResponse, summary="上报心跳")
|
||||
def report_heartbeat(
|
||||
req: HeartbeatRequest,
|
||||
user: CurrentUser,
|
||||
db: DbSession,
|
||||
) -> OkResponse:
|
||||
device_repo.touch_heartbeat(
|
||||
db,
|
||||
user_id=user.id,
|
||||
device_id=req.device_id,
|
||||
accessibility_enabled=req.accessibility_enabled,
|
||||
registration_id=req.registration_id,
|
||||
)
|
||||
return OkResponse()
|
||||
@@ -61,32 +61,6 @@ class Settings(BaseSettings):
|
||||
JG_VERIFY_ENDPOINT: str = "https://api.verification.jpush.cn/v1/web/loginTokenVerify"
|
||||
JG_REQUEST_TIMEOUT_SEC: int = 15
|
||||
|
||||
# ===== 极光推送 JPush(无障碍保护掉线告警)=====
|
||||
# 客户端 push appkey 已知 = 966b451a8d9cfe12d173ea9d(android build.gradle manifestPlaceholder)。
|
||||
# 若推送与一键登录/短信是同一个极光应用(大概率),JPUSH_* 留空即自动回退到 JG_*。
|
||||
# 否则在 .env 单独配 JPUSH_APP_KEY / JPUSH_MASTER_SECRET(对应那个 push appkey)。
|
||||
JPUSH_APP_KEY: str = ""
|
||||
JPUSH_MASTER_SECRET: str = ""
|
||||
JPUSH_PUSH_ENDPOINT: str = "https://api.jpush.cn/v3/push"
|
||||
|
||||
# 无障碍保护存活监控后台任务
|
||||
HEARTBEAT_MONITOR_ENABLED: bool = True # 总开关
|
||||
HEARTBEAT_TIMEOUT_MINUTES: int = 10 # 多久没心跳算掉线(≈3 个客户端心跳周期)
|
||||
HEARTBEAT_SCAN_INTERVAL_SEC: int = 60 # 扫描周期
|
||||
|
||||
@property
|
||||
def jpush_app_key(self) -> str:
|
||||
return self.JPUSH_APP_KEY or self.JG_APP_KEY
|
||||
|
||||
@property
|
||||
def jpush_master_secret(self) -> str:
|
||||
return self.JPUSH_MASTER_SECRET or self.JG_MASTER_SECRET
|
||||
|
||||
@property
|
||||
def jpush_configured(self) -> bool:
|
||||
"""推送凭证齐全(缺则 monitor 只扫不发,不报错)。"""
|
||||
return bool(self.jpush_app_key and self.jpush_master_secret)
|
||||
|
||||
# ===== 短信 =====
|
||||
SMS_MOCK: bool = True
|
||||
SMS_CODE_TTL_SEC: int = 300
|
||||
|
||||
@@ -1,148 +0,0 @@
|
||||
"""无障碍保护存活监控后台任务。
|
||||
|
||||
周期扫描「曾经保护过、当前 alive、心跳超时」的设备 = App 被彻底杀掉/无障碍已停(心跳断了),
|
||||
**命中即在服务器终端打印告警**(本期先不接推送,工程量大,用终端打印代替真实通知);并把状态机
|
||||
推进到 notified 防每轮重复打印(心跳恢复时由 repositories.device.touch_heartbeat 重置回 alive)。
|
||||
结构仿 withdraw_reconcile_worker(单实例锁 + asyncio 轮询 + 优雅退出)。
|
||||
|
||||
见 spec: spec/accessibility-liveness-push.md。
|
||||
"""
|
||||
from __future__ import annotations
|
||||
|
||||
import asyncio
|
||||
import contextlib
|
||||
import logging
|
||||
import os
|
||||
import time
|
||||
from collections.abc import Iterator
|
||||
from datetime import datetime, timezone
|
||||
from pathlib import Path
|
||||
|
||||
from sqlalchemy.exc import SQLAlchemyError
|
||||
|
||||
from app.core.config import settings
|
||||
from app.db.session import SessionLocal
|
||||
from app.repositories import device as device_repo
|
||||
|
||||
logger = logging.getLogger("shagua.heartbeat_monitor")
|
||||
_LOCK_PATH = Path(__file__).resolve().parents[2] / "data" / "heartbeat_monitor.lock"
|
||||
|
||||
|
||||
def _touch_lock() -> None:
|
||||
with contextlib.suppress(FileNotFoundError):
|
||||
os.utime(_LOCK_PATH, None)
|
||||
|
||||
|
||||
@contextlib.contextmanager
|
||||
def _single_instance_lock(stale_after_sec: int) -> Iterator[bool]:
|
||||
"""同机多进程保护:同一时间只允许一个监控 worker 运行。"""
|
||||
_LOCK_PATH.parent.mkdir(parents=True, exist_ok=True)
|
||||
fd: int | None = None
|
||||
try:
|
||||
try:
|
||||
fd = os.open(str(_LOCK_PATH), os.O_CREAT | os.O_EXCL | os.O_WRONLY)
|
||||
except FileExistsError:
|
||||
try:
|
||||
age = time.time() - _LOCK_PATH.stat().st_mtime
|
||||
except FileNotFoundError:
|
||||
age = stale_after_sec + 1
|
||||
if age > stale_after_sec:
|
||||
with contextlib.suppress(FileNotFoundError):
|
||||
_LOCK_PATH.unlink()
|
||||
try:
|
||||
fd = os.open(str(_LOCK_PATH), os.O_CREAT | os.O_EXCL | os.O_WRONLY)
|
||||
except FileExistsError:
|
||||
fd = None
|
||||
|
||||
if fd is None:
|
||||
yield False
|
||||
return
|
||||
|
||||
os.write(fd, f"pid={os.getpid()} started_at={int(time.time())}\n".encode("ascii"))
|
||||
yield True
|
||||
finally:
|
||||
if fd is not None:
|
||||
os.close(fd)
|
||||
with contextlib.suppress(FileNotFoundError):
|
||||
_LOCK_PATH.unlink()
|
||||
|
||||
|
||||
def _silent_seconds(last: datetime | None) -> int | None:
|
||||
"""距上次心跳的秒数(兼容 sqlite 取回的 naive datetime)。"""
|
||||
if last is None:
|
||||
return None
|
||||
ref = datetime.now(timezone.utc) if last.tzinfo is not None else datetime.utcnow()
|
||||
return int((ref - last).total_seconds())
|
||||
|
||||
|
||||
def _scan_once(timeout_minutes: int) -> dict:
|
||||
"""扫描一轮:找出心跳超时(App 被彻底杀掉/无障碍已停)的设备,在**服务器终端打印**告警代替真实推送。
|
||||
|
||||
本期不接推送(极光/厂商通道工程量大),只做服务端掉线检测:命中即 logger.warning 打印到终端,
|
||||
并把状态机推进到 notified 防每轮重复打印(心跳恢复时 touch_heartbeat 会重置回 alive)。
|
||||
"""
|
||||
notified = 0
|
||||
with SessionLocal() as db:
|
||||
overdue = device_repo.list_overdue(db, timeout_minutes=timeout_minutes)
|
||||
for device in overdue:
|
||||
silent = _silent_seconds(device.last_heartbeat_at)
|
||||
logger.warning(
|
||||
"🔴 [掉线检测] user_id=%s device_id=%s 已 %s 秒无心跳(阈值 %d 分钟)"
|
||||
" → 判定 App 已被杀/无障碍已停。【本应推送通知提醒用户重开;推送暂未接,先终端打印代替】",
|
||||
device.user_id,
|
||||
device.device_id,
|
||||
silent if silent is not None else "?",
|
||||
timeout_minutes,
|
||||
)
|
||||
device_repo.mark_notified(db, device_id_pk=device.id)
|
||||
notified += 1
|
||||
return {"checked": len(overdue), "notified": notified}
|
||||
|
||||
|
||||
async def _run_loop() -> None:
|
||||
interval = max(10, int(settings.HEARTBEAT_SCAN_INTERVAL_SEC))
|
||||
timeout_minutes = max(1, int(settings.HEARTBEAT_TIMEOUT_MINUTES))
|
||||
lock_stale_after = max(interval * 3, 600)
|
||||
with _single_instance_lock(lock_stale_after) as lock_acquired:
|
||||
if not lock_acquired:
|
||||
logger.warning("heartbeat monitor skipped: another worker owns lock")
|
||||
return
|
||||
await _run_locked_loop(interval, timeout_minutes)
|
||||
|
||||
|
||||
async def _run_locked_loop(interval: int, timeout_minutes: int) -> None:
|
||||
logger.info(
|
||||
"heartbeat monitor started interval=%ss timeout=%sm",
|
||||
interval,
|
||||
timeout_minutes,
|
||||
)
|
||||
try:
|
||||
while True:
|
||||
try:
|
||||
_touch_lock()
|
||||
result = await asyncio.to_thread(_scan_once, timeout_minutes)
|
||||
if result["notified"]:
|
||||
logger.info("heartbeat monitor result=%s", result)
|
||||
except SQLAlchemyError:
|
||||
logger.exception("heartbeat monitor db error")
|
||||
except Exception: # noqa: BLE001 - 后台任务不能因单次异常退出
|
||||
logger.exception("heartbeat monitor unexpected error")
|
||||
await asyncio.sleep(interval)
|
||||
except asyncio.CancelledError:
|
||||
logger.info("heartbeat monitor stopped")
|
||||
raise
|
||||
|
||||
|
||||
def start_heartbeat_monitor() -> asyncio.Task | None:
|
||||
if not settings.HEARTBEAT_MONITOR_ENABLED:
|
||||
logger.info("heartbeat monitor disabled")
|
||||
return None
|
||||
return asyncio.create_task(_run_loop(), name="heartbeat-monitor")
|
||||
|
||||
|
||||
async def stop_heartbeat_monitor(task: asyncio.Task | None) -> None:
|
||||
if task is None:
|
||||
return
|
||||
task.cancel()
|
||||
with contextlib.suppress(asyncio.CancelledError):
|
||||
await task
|
||||
@@ -1,94 +0,0 @@
|
||||
"""极光推送 JPush(无障碍保护掉线告警)。
|
||||
|
||||
调极光 Push REST v3 /v3/push 给指定 registration_id 推一条通知。鉴权复用极光 Basic Auth
|
||||
(appKey:masterSecret),与一键登录/短信同模式(见 integrations/jiguang.py、sms.py)。
|
||||
|
||||
凭证:settings.jpush_app_key / jpush_master_secret(JPUSH_* 留空时自动回退到 JG_*,
|
||||
若推送与一键登录是同一极光应用)。客户端 push appkey 已知 = 966b451a8d9cfe12d173ea9d。
|
||||
厂商通道(App 被杀也能到达)由极光后台 + 客户端插件负责,本服务只管调 push 接口。
|
||||
"""
|
||||
from __future__ import annotations
|
||||
|
||||
import base64
|
||||
import logging
|
||||
|
||||
import httpx
|
||||
|
||||
from app.core.config import settings
|
||||
|
||||
logger = logging.getLogger("shagua.jpush")
|
||||
|
||||
|
||||
class JPushError(Exception):
|
||||
"""极光推送调用失败。"""
|
||||
|
||||
|
||||
class JPushNotConfiguredError(JPushError):
|
||||
"""缺 appKey / masterSecret。"""
|
||||
|
||||
|
||||
def push_to_registration_ids(
|
||||
registration_ids: list[str],
|
||||
*,
|
||||
title: str,
|
||||
alert: str,
|
||||
extras: dict | None = None,
|
||||
) -> dict:
|
||||
"""给一批 registration_id 推送通知 + 透传消息。失败抛 JPushError。
|
||||
|
||||
Returns: 极光响应 JSON(含 sendno / msg_id)。
|
||||
"""
|
||||
if not settings.jpush_configured:
|
||||
raise JPushNotConfiguredError(
|
||||
"JPush 未配置(缺 JPUSH_APP_KEY/JPUSH_MASTER_SECRET,且 JG_* 也为空)"
|
||||
)
|
||||
reg_ids = [r for r in registration_ids if r]
|
||||
if not reg_ids:
|
||||
raise JPushError("registration_ids 为空")
|
||||
|
||||
auth_b64 = base64.b64encode(
|
||||
f"{settings.jpush_app_key}:{settings.jpush_master_secret}".encode()
|
||||
).decode()
|
||||
payload = {
|
||||
"platform": ["android"],
|
||||
"audience": {"registration_id": reg_ids},
|
||||
"notification": {
|
||||
"android": {
|
||||
"alert": alert,
|
||||
"title": title,
|
||||
"priority": 1,
|
||||
"extras": extras or {},
|
||||
},
|
||||
},
|
||||
"message": {
|
||||
"msg_content": alert,
|
||||
"title": title,
|
||||
"content_type": "text",
|
||||
"extras": extras or {},
|
||||
},
|
||||
"options": {"time_to_live": 86400, "apns_production": True},
|
||||
}
|
||||
|
||||
try:
|
||||
resp = httpx.post(
|
||||
settings.JPUSH_PUSH_ENDPOINT,
|
||||
json=payload,
|
||||
headers={
|
||||
"Content-Type": "application/json",
|
||||
"Authorization": f"Basic {auth_b64}",
|
||||
},
|
||||
timeout=settings.JG_REQUEST_TIMEOUT_SEC,
|
||||
)
|
||||
except httpx.HTTPError as e:
|
||||
raise JPushError(f"jpush 网络错误: {e}") from e
|
||||
|
||||
if resp.status_code != 200:
|
||||
body = resp.text[:300]
|
||||
logger.error("[JPush] http=%s body=%s", resp.status_code, body)
|
||||
raise JPushError(f"jpush http {resp.status_code}")
|
||||
|
||||
data = resp.json()
|
||||
if not data.get("sendno") and not data.get("msg_id"):
|
||||
logger.error("[JPush] unexpected response: %s", data)
|
||||
raise JPushError(f"jpush 响应异常: {data}")
|
||||
return data
|
||||
+1
-11
@@ -20,9 +20,7 @@ from app.api.v1.compare import router as compare_router
|
||||
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.v1.device import router as device_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
|
||||
@@ -36,10 +34,6 @@ from app.api.v1.user import router as user_router
|
||||
from app.api.v1.wallet import router as wallet_router
|
||||
from app.api.v1.wxpay import router as wxpay_router
|
||||
from app.core.config import settings
|
||||
from app.core.heartbeat_monitor_worker import (
|
||||
start_heartbeat_monitor,
|
||||
stop_heartbeat_monitor,
|
||||
)
|
||||
from app.core.logging import setup_logging
|
||||
from app.core.withdraw_reconcile_worker import (
|
||||
start_withdraw_reconcile_worker,
|
||||
@@ -61,11 +55,9 @@ async def lifespan(_: FastAPI) -> AsyncIterator[None]:
|
||||
settings.DATABASE_URL.split("://", 1)[0],
|
||||
)
|
||||
reconcile_task = start_withdraw_reconcile_worker()
|
||||
heartbeat_task = start_heartbeat_monitor()
|
||||
try:
|
||||
yield
|
||||
finally:
|
||||
await stop_heartbeat_monitor(heartbeat_task)
|
||||
await stop_withdraw_reconcile_worker(reconcile_task)
|
||||
logger.info("shutting down")
|
||||
|
||||
@@ -98,7 +90,6 @@ app.include_router(user_router)
|
||||
app.include_router(feedback_router)
|
||||
app.include_router(invite_router)
|
||||
app.include_router(coupon_router)
|
||||
app.include_router(device_router)
|
||||
app.include_router(compare_router)
|
||||
app.include_router(compare_record_router)
|
||||
app.include_router(compare_milestone_router)
|
||||
@@ -111,9 +102,8 @@ 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。
|
||||
|
||||
@@ -7,7 +7,6 @@ from app.models.admin import AdminAuditLog, AdminUser # noqa: F401
|
||||
from app.models.app_config import AppConfig # noqa: F401
|
||||
from app.models.comparison import ComparisonRecord # noqa: F401
|
||||
from app.models.comparison_milestone import ComparisonMilestoneClaim # noqa: F401
|
||||
from app.models.device import Device # noqa: F401
|
||||
from app.models.coupon_state import ( # noqa: F401
|
||||
CouponClaimRecord,
|
||||
CouponDailyCompletion,
|
||||
@@ -24,7 +23,6 @@ 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
|
||||
|
||||
@@ -94,7 +94,7 @@ class ComparisonRecord(Base):
|
||||
# ===== 明细(JSON,越详细越好)=====
|
||||
# 下单菜品 [{name, qty, specs?}]
|
||||
items: Mapped[list] = mapped_column(_JSON, nullable=False, default=list)
|
||||
# 逐平台对比 [{platform_id, platform_name, package, price, is_source, rank, coupon_saved, coupon_name, applied_coupons}](price/coupon_saved 单位:元,原样存;coupon_name=优惠来源名;applied_coupons=[{name,amount}] 多券明细)
|
||||
# 逐平台对比 [{platform_id, platform_name, package, price, is_source, rank, coupon_saved, coupon_name}](price/coupon_saved 单位:元,原样存;coupon_name=优惠来源名)
|
||||
comparison_results: Mapped[list] = mapped_column(_JSON, nullable=False, default=list)
|
||||
# 目标平台未找到、跳过的菜名
|
||||
skipped_dish_names: Mapped[list] = mapped_column(_JSON, nullable=False, default=list)
|
||||
|
||||
@@ -137,23 +137,14 @@ class CouponDailyCompletion(Base):
|
||||
|
||||
|
||||
class CouponPromptEngagement(Base):
|
||||
"""按 (device, App, 自然日) 记"今天该 App 是否弹过/engage 过领券引导窗"——弹窗频控源。
|
||||
|
||||
2026-06-12 改为**按 App 为单位**(对齐方案文档「一个 app 一天只能弹出弹窗一次」+ 产品确认:
|
||||
没领完就进其他平台要再弹,彻底领完才全局不弹——"彻底领完"由 coupon_daily_completion 负责):
|
||||
- app_package 非空:该 App 当天的频控行(shown/dismissed 由客户端带包名上报)。
|
||||
- app_package 为 NULL:全局兜底行(step=0 透传链路自动记的 claim_started 拿不到包名 +
|
||||
旧版客户端不带包名)。**按 App 查询时忽略 NULL 行**(领券中途终止 ≠ 其他 App engage 过);
|
||||
不带包名的旧式查询仍把 NULL 行算作"今天 engage 过"(旧客户端行为不变)。
|
||||
"""
|
||||
"""按 (device, 自然日) 记"今天是否对领券引导窗表达过意向"——弹窗频控源。"""
|
||||
|
||||
__tablename__ = "coupon_prompt_engagement"
|
||||
__table_args__ = (
|
||||
# 一台设备、一个 App、一天一条(app_package=NULL 的全局兜底行不受唯一约束限制——
|
||||
# SQL 标准里 NULL 互不相等;写入走 select-first upsert,正常不会堆重复行)。
|
||||
# 一台设备一天一条:今天 engage 过(领或拒)就不再弹。
|
||||
UniqueConstraint(
|
||||
"device_id", "app_package", "engage_date",
|
||||
name="uq_coupon_engage_device_app_date",
|
||||
"device_id", "engage_date",
|
||||
name="uq_coupon_engage_device_date",
|
||||
),
|
||||
)
|
||||
|
||||
@@ -161,12 +152,10 @@ class CouponPromptEngagement(Base):
|
||||
|
||||
device_id: Mapped[str] = mapped_column(String(64), nullable=False)
|
||||
user_id: Mapped[int | None] = mapped_column(Integer, index=True, nullable=True)
|
||||
# 在哪个外卖 App 弹的窗(Android 包名,如 com.sankuai.meituan)。NULL = 全局兜底行(见类注释)。
|
||||
app_package: Mapped[str | None] = mapped_column(String(64), nullable=True)
|
||||
# Asia/Shanghai 自然日。
|
||||
engage_date: Mapped[date] = mapped_column(Date, nullable=False)
|
||||
# shown(弹窗已对该 App 弹出,频控主判据)/ claim_started(点了一键领取)/
|
||||
# dismissed(点了拒绝/关闭)。判断只看"今天该 App 有没有这条",type 仅记录区分。
|
||||
# claim_started(点了一键领取)/ dismissed(点了拒绝/关闭)。仅记录区分,
|
||||
# 判断只看"今天有没有这条",type 不影响弹不弹。
|
||||
engage_type: Mapped[str] = mapped_column(String(16), nullable=False)
|
||||
|
||||
created_at: Mapped[datetime] = mapped_column(
|
||||
|
||||
@@ -1,83 +0,0 @@
|
||||
"""设备表(无障碍保护存活检测 + 极光推送)。
|
||||
|
||||
每条 = 一个用户的一台设备(per-install,device_id 由客户端 DeviceId.get() 生成)。
|
||||
客户端的无障碍服务存活时周期上报心跳刷新 last_heartbeat_at;App 前台/登录时上报
|
||||
registration_id(极光推送目标)。后端 heartbeat_monitor_worker 扫描「曾经保护过、
|
||||
现在心跳超时」的设备,通过极光推送提醒用户重开无障碍。
|
||||
|
||||
liveness_state 状态机(防刷屏,一次掉线只推一条):
|
||||
unknown → alive(收到 service 心跳)→ silent/notified(扫描发现超时并已推送)
|
||||
心跳恢复时 handler 重置回 alive。
|
||||
见 spec: spec/accessibility-liveness-push.md。
|
||||
"""
|
||||
from __future__ import annotations
|
||||
|
||||
from datetime import datetime
|
||||
|
||||
from sqlalchemy import (
|
||||
Boolean,
|
||||
DateTime,
|
||||
ForeignKey,
|
||||
Integer,
|
||||
String,
|
||||
UniqueConstraint,
|
||||
func,
|
||||
)
|
||||
from sqlalchemy.orm import Mapped, mapped_column
|
||||
|
||||
from app.db.base import Base
|
||||
|
||||
|
||||
class Device(Base):
|
||||
__tablename__ = "device"
|
||||
__table_args__ = (
|
||||
UniqueConstraint("user_id", "device_id", name="uq_device_user_device"),
|
||||
)
|
||||
|
||||
id: Mapped[int] = mapped_column(Integer, primary_key=True, autoincrement=True)
|
||||
user_id: Mapped[int] = mapped_column(
|
||||
Integer, ForeignKey("user.id"), index=True, nullable=False
|
||||
)
|
||||
# 客户端 DeviceId.get() 生成的 per-install id(如 device_Pixel_ab12cd34)
|
||||
device_id: Mapped[str] = mapped_column(String(128), index=True, nullable=False)
|
||||
# 极光推送 registration id;拿到才填(JCollectionAuth 同意后才下发)
|
||||
registration_id: Mapped[str | None] = mapped_column(String(64), nullable=True)
|
||||
platform: Mapped[str] = mapped_column(String(16), nullable=False, default="android")
|
||||
app_version: Mapped[str | None] = mapped_column(String(32), nullable=True)
|
||||
|
||||
# 收到过 service 心跳即 true(=该设备开过无障碍,功能对它有意义)
|
||||
ever_protected: Mapped[bool] = mapped_column(
|
||||
Boolean, nullable=False, default=False
|
||||
)
|
||||
# 最近一次 service 心跳时间(存活证明);超时即视为保护掉线
|
||||
last_heartbeat_at: Mapped[datetime | None] = mapped_column(
|
||||
DateTime(timezone=True), index=True, nullable=True
|
||||
)
|
||||
# 最近一次上报的无障碍开关状态(观测用)
|
||||
last_report_protection_on: Mapped[bool] = mapped_column(
|
||||
Boolean, nullable=False, default=False
|
||||
)
|
||||
# unknown / alive / silent / notified
|
||||
liveness_state: Mapped[str] = mapped_column(
|
||||
String(16), nullable=False, default="unknown"
|
||||
)
|
||||
# 最近一次推送告警时间
|
||||
notified_at: Mapped[datetime | None] = mapped_column(
|
||||
DateTime(timezone=True), nullable=True
|
||||
)
|
||||
|
||||
created_at: Mapped[datetime] = mapped_column(
|
||||
DateTime(timezone=True), server_default=func.now(), nullable=False
|
||||
)
|
||||
updated_at: Mapped[datetime] = mapped_column(
|
||||
DateTime(timezone=True),
|
||||
server_default=func.now(),
|
||||
onupdate=func.now(),
|
||||
nullable=False,
|
||||
)
|
||||
|
||||
def __repr__(self) -> str: # pragma: no cover
|
||||
return (
|
||||
f"<Device id={self.id} user_id={self.user_id} "
|
||||
f"device_id={self.device_id} state={self.liveness_state}>"
|
||||
)
|
||||
@@ -1,117 +0,0 @@
|
||||
"""平台店铺表(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=shopId)、**美团**(meituan_poi_id_str,非稳定主键、单列存)、
|
||||
**京东**(id_jd=storeId + jd_vender_id,均稳定数字主键;3.cn 短链反查)。
|
||||
"""
|
||||
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 机制(poi_id_str 每次变,见下方 meituan_poi_id_str),id_meituan
|
||||
# 留给将来 CPS API 的稳定 poi_id;name-only 时仅 name_meituan(源平台店名来自 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)
|
||||
# 京东:已接通 share→id(3.cn 短链反查)。id_jd = storeId(门店稳定数字主键,同 taobao shopId→
|
||||
# id_taobao);venderId(deeplink 还需)单列存 jd_vender_id。name_jd = 店铺页店名。
|
||||
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
|
||||
|
||||
# ===== 美团原料(同淘宝;dpurl.cn 短链 → 302 反查 poi_id_str → imeituan:// deeplink)=====
|
||||
# ⚠️ poi_id_str 每次分享重新加密、非稳定主键(调研文档 §八), 故单列存"可复跳的一次性票据",
|
||||
# 不进 id_meituan —— 后者留给将来 CPS API 拿到的稳定数字 poi_id。
|
||||
meituan_poi_id_str: Mapped[str | None] = mapped_column(String(64), index=True, nullable=True)
|
||||
meituan_share_url: Mapped[str | None] = mapped_column(String(256), nullable=True) # dpurl.cn 短链
|
||||
meituan_resolved_url: Mapped[str | None] = mapped_column(Text, nullable=True) # 302 落地 menu URL(含 poi_id_str)
|
||||
meituan_deeplink: Mapped[str | None] = mapped_column(Text, nullable=True) # 拼好的 imeituan:// 店内搜索 deeplink
|
||||
|
||||
# ===== 京东原料(秒送;3.cn 短链 → 跟随重定向反查 venderId+storeId → openapp.jdmobile:// deeplink)=====
|
||||
# storeId 进 id_jd(稳定店主键);venderId 单列存(deeplink 模板 venderId+storeId 都要,且 venderId
|
||||
# 是商家维度、可跨门店,与门店 storeId 分开记)。其余三列与 taobao_*/meituan_* 平行。
|
||||
jd_vender_id: Mapped[str | None] = mapped_column(String(64), index=True, nullable=True)
|
||||
jd_share_url: Mapped[str | None] = mapped_column(String(256), nullable=True) # 3.cn 短链
|
||||
jd_resolved_url: Mapped[str | None] = mapped_column(Text, nullable=True) # 反查出的目标 openapp.jdmobile:// deeplink
|
||||
jd_deeplink: Mapped[str | None] = mapped_column(Text, nullable=True) # 拼好的 pages/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}>"
|
||||
)
|
||||
@@ -31,46 +31,26 @@ def today_cn() -> date:
|
||||
|
||||
# ===== 弹窗频控(coupon_prompt_engagement)=====
|
||||
|
||||
def has_engaged_today(
|
||||
db: Session, device_id: str, app_package: str | None = None
|
||||
) -> bool:
|
||||
"""这台设备今天是否已弹过/engage 过领券引导窗。有 = 不再弹。
|
||||
|
||||
2026-06-12 按 App 为单位:
|
||||
- app_package 给定(新客户端):只看**该 App** 当天的行;app_package=NULL 的全局兜底行
|
||||
(step=0 claim_started / 旧客户端写的)**不算**——领券中途终止不该堵死其他 App 的弹窗,
|
||||
"彻底领完才全局不弹"由 has_completed_today 在 should-show 端点单独把关。
|
||||
- app_package=None(旧客户端不带包名):保持旧全局语义,当天任意一行都算 engage 过。
|
||||
"""
|
||||
cond = [
|
||||
CouponPromptEngagement.device_id == device_id,
|
||||
CouponPromptEngagement.engage_date == today_cn(),
|
||||
]
|
||||
if app_package is not None:
|
||||
cond.append(CouponPromptEngagement.app_package == app_package)
|
||||
row = db.execute(select(CouponPromptEngagement.id).where(*cond)).first()
|
||||
def has_engaged_today(db: Session, device_id: str) -> bool:
|
||||
"""这台设备今天是否已对领券引导窗表达过意向(领或拒)。有 = 不再弹。"""
|
||||
row = db.execute(
|
||||
select(CouponPromptEngagement.id).where(
|
||||
CouponPromptEngagement.device_id == device_id,
|
||||
CouponPromptEngagement.engage_date == today_cn(),
|
||||
)
|
||||
).first()
|
||||
return row is not None
|
||||
|
||||
|
||||
def mark_engagement(
|
||||
db: Session,
|
||||
device_id: str,
|
||||
user_id: int | None,
|
||||
engage_type: str,
|
||||
app_package: str | None = None,
|
||||
db: Session, device_id: str, user_id: int | None, engage_type: str
|
||||
) -> None:
|
||||
"""记今日意向(shown / claim_started / dismissed)。(device, App, 今天) 唯一,幂等 upsert。
|
||||
|
||||
app_package=None = 全局兜底行(step=0 透传链路拿不到包名 / 旧客户端),与各 App 行互不覆盖。
|
||||
同 (device, App, 日) 重复上报走更新(shown → dismissed 升级 engage_type)。
|
||||
"""
|
||||
"""记今日意向(claim_started / dismissed)。(device, 今天) 唯一,幂等 upsert。"""
|
||||
today = today_cn()
|
||||
row = db.execute(
|
||||
select(CouponPromptEngagement).where(
|
||||
CouponPromptEngagement.device_id == device_id,
|
||||
CouponPromptEngagement.engage_date == today,
|
||||
# SQLAlchemy 的 == None 会生成 IS NULL,NULL 兜底行与 App 行各自独立 upsert。
|
||||
CouponPromptEngagement.app_package == app_package,
|
||||
)
|
||||
).scalar_one_or_none()
|
||||
if row is not None:
|
||||
@@ -79,20 +59,19 @@ def mark_engagement(
|
||||
row.user_id = user_id
|
||||
else:
|
||||
db.add(CouponPromptEngagement(
|
||||
device_id=device_id, user_id=user_id, app_package=app_package,
|
||||
device_id=device_id, user_id=user_id,
|
||||
engage_date=today, engage_type=engage_type,
|
||||
))
|
||||
try:
|
||||
db.commit()
|
||||
except IntegrityError:
|
||||
# 并发下另一请求刚插了同 (device, App, 日) → 唯一约束撞,回滚忽略(本就幂等)。
|
||||
# 并发下另一请求刚插了同 (device, 日) → 唯一约束撞,回滚忽略(本就幂等)。
|
||||
db.rollback()
|
||||
|
||||
|
||||
def reset_today_engagement(db: Session, device_id: str) -> int:
|
||||
"""删这台设备今天的 engagement(开发设置「重置今日领券弹窗状态」调,测频控用)。
|
||||
按 device+日 删,**所有 App 的频控行 + NULL 兜底行一并清**(频控按 App 拆行后语义不变:
|
||||
重置 = 等效重装,每个 App 今天都又能弹)。返回删除行数。"""
|
||||
删后 has_engaged_today → false,今天又能弹。返回删除行数。"""
|
||||
result = db.execute(
|
||||
delete(CouponPromptEngagement).where(
|
||||
CouponPromptEngagement.device_id == device_id,
|
||||
@@ -105,21 +84,6 @@ def reset_today_engagement(db: Session, device_id: str) -> int:
|
||||
|
||||
# ===== 今日跑完整轮(coupon_daily_completion)=====
|
||||
|
||||
def reset_today_completion(db: Session, device_id: str) -> int:
|
||||
"""删这台设备今天的「已跑完整轮」记录(开发设置「重置今日领券弹窗状态」调)。
|
||||
删后 has_completed_today → false:首页「去领取」卡恢复可点、领券弹窗 CTA 不再置灰
|
||||
(只清 engagement 不清这条的话, 重置后弹窗能弹但「一键自动领取」仍是灰的, 2026-06-12)。
|
||||
返回删除行数。"""
|
||||
result = db.execute(
|
||||
delete(CouponDailyCompletion).where(
|
||||
CouponDailyCompletion.device_id == device_id,
|
||||
CouponDailyCompletion.complete_date == today_cn(),
|
||||
)
|
||||
)
|
||||
db.commit()
|
||||
return result.rowcount or 0
|
||||
|
||||
|
||||
def has_completed_today(db: Session, device_id: str) -> bool:
|
||||
"""这台设备今天是否已跑完整轮领券(到 done 帧)。有 = 首页置灰、不能再领。"""
|
||||
row = db.execute(
|
||||
|
||||
@@ -1,106 +0,0 @@
|
||||
"""device 表读写(设备注册 / 心跳 / 超时扫描)。"""
|
||||
from __future__ import annotations
|
||||
|
||||
from datetime import datetime, timedelta, timezone
|
||||
|
||||
from sqlalchemy import select
|
||||
from sqlalchemy.orm import Session
|
||||
|
||||
from app.models.device import Device
|
||||
|
||||
|
||||
def _get(db: Session, *, user_id: int, device_id: str) -> Device | None:
|
||||
stmt = select(Device).where(
|
||||
Device.user_id == user_id, Device.device_id == device_id
|
||||
)
|
||||
return db.execute(stmt).scalar_one_or_none()
|
||||
|
||||
|
||||
def register_or_update(
|
||||
db: Session,
|
||||
*,
|
||||
user_id: int,
|
||||
device_id: str,
|
||||
registration_id: str | None,
|
||||
platform: str = "android",
|
||||
app_version: str | None = None,
|
||||
) -> Device:
|
||||
"""注册设备或更新其 registration_id / 元信息。upsert by (user_id, device_id)。"""
|
||||
device = _get(db, user_id=user_id, device_id=device_id)
|
||||
if device is None:
|
||||
device = Device(
|
||||
user_id=user_id,
|
||||
device_id=device_id,
|
||||
registration_id=registration_id,
|
||||
platform=platform or "android",
|
||||
app_version=app_version,
|
||||
)
|
||||
db.add(device)
|
||||
else:
|
||||
if registration_id:
|
||||
device.registration_id = registration_id
|
||||
if platform:
|
||||
device.platform = platform
|
||||
if app_version:
|
||||
device.app_version = app_version
|
||||
db.commit()
|
||||
db.refresh(device)
|
||||
return device
|
||||
|
||||
|
||||
def touch_heartbeat(
|
||||
db: Session,
|
||||
*,
|
||||
user_id: int,
|
||||
device_id: str,
|
||||
accessibility_enabled: bool,
|
||||
registration_id: str | None,
|
||||
) -> Device:
|
||||
"""处理一次心跳(心跳也能自注册)。
|
||||
|
||||
service 心跳或 accessibility_enabled=true 时,刷新存活并把状态机重置回 alive、
|
||||
清掉 notified_at(掉线恢复 → 下次再断才会再推一条)。
|
||||
"""
|
||||
now = datetime.now(timezone.utc)
|
||||
device = _get(db, user_id=user_id, device_id=device_id)
|
||||
if device is None:
|
||||
device = Device(user_id=user_id, device_id=device_id)
|
||||
db.add(device)
|
||||
|
||||
if registration_id:
|
||||
device.registration_id = registration_id
|
||||
device.last_report_protection_on = accessibility_enabled
|
||||
|
||||
if accessibility_enabled:
|
||||
device.last_heartbeat_at = now
|
||||
device.ever_protected = True
|
||||
device.liveness_state = "alive"
|
||||
device.notified_at = None
|
||||
|
||||
db.commit()
|
||||
db.refresh(device)
|
||||
return device
|
||||
|
||||
|
||||
def list_overdue(db: Session, *, timeout_minutes: int) -> list[Device]:
|
||||
"""掉线设备:曾经保护过、当前 alive、心跳超时。
|
||||
|
||||
本期只做终端打印检测、不推送 → 不再要求有 registration_id(没接极光 token 的设备也要检出)。
|
||||
"""
|
||||
cutoff = datetime.now(timezone.utc) - timedelta(minutes=timeout_minutes)
|
||||
stmt = select(Device).where(
|
||||
Device.ever_protected.is_(True),
|
||||
Device.liveness_state == "alive",
|
||||
Device.last_heartbeat_at.is_not(None),
|
||||
Device.last_heartbeat_at < cutoff,
|
||||
)
|
||||
return list(db.execute(stmt).scalars().all())
|
||||
|
||||
|
||||
def mark_notified(db: Session, *, device_id_pk: int) -> None:
|
||||
"""标记已推送告警(状态机进入 notified,避免重复推送)。"""
|
||||
device = db.get(Device, device_id_pk)
|
||||
if device is not None:
|
||||
device.liveness_state = "notified"
|
||||
device.notified_at = datetime.now(timezone.utc)
|
||||
db.commit()
|
||||
@@ -1,172 +0,0 @@
|
||||
"""平台店铺映射落库:一次比价(trace)一行,各目标平台解析出 id 就 upsert 进同一行。
|
||||
|
||||
合并键 = trace_id(唯一约束)。一次跨平台比价 = 一个 trace = 一个真实店铺:淘宝腿解析出
|
||||
shopId 先建行(填淘宝列),京东腿(将来)解析出 id 再 upsert 进**同一行**(填京东列)。
|
||||
|
||||
合并策略 = 填空(fill-the-blanks):只写该行当前为 NULL 的列,绝不覆盖已有非空值。保证后到
|
||||
的平台只填自己那几列、动不了先到平台的数据;共享列(geo / source / 溯源)先到先得。
|
||||
不依赖 ON CONFLICT,跨方言(PG / SQLite dev)都安全;并发撞唯一约束则回滚后转走合并路径。
|
||||
|
||||
⚠️ 一行里 id_taobao 与 id_jd 共存只表示"两条腿搜同一个源店名各自匹配到了某家店",是
|
||||
name-match 置信度、非已核实同一实体。作为 append-only 原始资产留存,跨平台精确匹配由下游做。
|
||||
"""
|
||||
from __future__ import annotations
|
||||
|
||||
import logging
|
||||
import math
|
||||
|
||||
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")
|
||||
|
||||
# 源平台 → 该平台店名所在列(缓存查询的匹配键)。镜像 pricebot reporter 的同名表。
|
||||
_SOURCE_NAME_COLUMN = {
|
||||
"meituan": "name_meituan",
|
||||
"meituan_waimai": "name_meituan",
|
||||
"jd_waimai": "name_jd",
|
||||
"jd_waimai_standalone": "name_jd",
|
||||
"taobao_flash": "name_taobao",
|
||||
}
|
||||
|
||||
# 源平台 → 它自己对应的目标 key(查缓存时排除"源平台自己", 不会复用源平台的店铺 id)。
|
||||
_SOURCE_TARGET_KEY = {
|
||||
"meituan": "meituan", "meituan_waimai": "meituan",
|
||||
"jd_waimai": "jd", "jd_waimai_standalone": "jd",
|
||||
"taobao_flash": "taobao",
|
||||
}
|
||||
|
||||
# upsert 填空时可写的列。不含: trace_id(合并键)/ business_type(非空默认)/
|
||||
# id(主键)/ created_at(server_default)。各平台只会带自己那几列非空, 其余为 None 不动。
|
||||
_MERGE_COLUMNS = (
|
||||
"source_platform",
|
||||
"id_taobao", "name_taobao", "id_meituan", "name_meituan", "id_jd", "name_jd",
|
||||
"city", "geohash", "lng", "lat", "taobao_address",
|
||||
"source_device_id", "source_user_id",
|
||||
"taobao_share_url", "taobao_resolved_url", "taobao_deeplink",
|
||||
"meituan_poi_id_str", "meituan_share_url", "meituan_resolved_url", "meituan_deeplink",
|
||||
"jd_vender_id", "jd_share_url", "jd_resolved_url", "jd_deeplink",
|
||||
"attrs",
|
||||
)
|
||||
|
||||
|
||||
def _find(db: Session, trace_id: str) -> StoreMapping | None:
|
||||
return db.execute(
|
||||
select(StoreMapping).where(StoreMapping.trace_id == trace_id)
|
||||
).scalar_one_or_none()
|
||||
|
||||
|
||||
def _merge_fill_blanks(existing: StoreMapping, payload: StoreMappingIn) -> list[str]:
|
||||
"""把 payload 里非空、且 existing 当前为 NULL 的列填进去。返回被填的列名(空=无变化)。"""
|
||||
filled: list[str] = []
|
||||
for col in _MERGE_COLUMNS:
|
||||
new = getattr(payload, col)
|
||||
if new is not None and getattr(existing, col) is None:
|
||||
setattr(existing, col, new)
|
||||
filled.append(col)
|
||||
return filled
|
||||
|
||||
|
||||
def upsert(db: Session, payload: StoreMappingIn) -> tuple[int, int | None]:
|
||||
"""落一行跨平台店铺映射,返回 (created, row_id)。
|
||||
created=1 新建该 trace 的行 / 0 合并进已存在行(填空,不覆盖)。"""
|
||||
existing = _find(db, payload.trace_id)
|
||||
if existing is None:
|
||||
# 首写: 把 payload 全部可合并列灌进去(逐列 setattr 而非硬编码构造器, 否则首写的是
|
||||
# 美团/京东腿时它们的列会漏 —— 不在硬编码列表里就丢)。trace_id/business_type 是键/默认, 显式给。
|
||||
row = StoreMapping(
|
||||
trace_id=payload.trace_id,
|
||||
business_type=payload.business_type,
|
||||
)
|
||||
for col in _MERGE_COLUMNS:
|
||||
setattr(row, col, getattr(payload, col))
|
||||
db.add(row)
|
||||
try:
|
||||
db.commit()
|
||||
except IntegrityError:
|
||||
# 并发: 另一个请求刚插了同 trace → 撞唯一约束。回滚后转合并路径填空。
|
||||
db.rollback()
|
||||
existing = _find(db, payload.trace_id)
|
||||
if existing is None:
|
||||
raise
|
||||
logger.warning("store_mapping 并发冲突 trace=%s, 转填空合并", payload.trace_id)
|
||||
else:
|
||||
db.refresh(row)
|
||||
return 1, row.id
|
||||
|
||||
# 已存在(或并发回退到此): 填空合并, 只写当前 NULL 的列
|
||||
filled = _merge_fill_blanks(existing, payload)
|
||||
if filled:
|
||||
db.commit()
|
||||
db.refresh(existing)
|
||||
logger.info("store_mapping 合并 trace=%s 填列=%s", payload.trace_id, filled)
|
||||
return 0, existing.id
|
||||
|
||||
|
||||
# ============================================================
|
||||
# 缓存查询: 比价前按"源平台店名"反查已沉淀的各目标平台店铺 id, 命中就让 pricebot 直接
|
||||
# deeplink 跳店内搜索, 省掉"开平台→进店→分享反查"整段。
|
||||
# ============================================================
|
||||
|
||||
def _haversine_km(lat1: float, lng1: float, lat2: float, lng2: float) -> float:
|
||||
"""两点球面距离(km)。仅用于同名候选里挑最近, 精度够用。"""
|
||||
r = 6371.0
|
||||
p1, p2 = math.radians(lat1), math.radians(lat2)
|
||||
dp = math.radians(lat2 - lat1)
|
||||
dl = math.radians(lng2 - lng1)
|
||||
a = math.sin(dp / 2) ** 2 + math.cos(p1) * math.cos(p2) * math.sin(dl / 2) ** 2
|
||||
return 2 * r * math.asin(math.sqrt(a))
|
||||
|
||||
|
||||
def _pick_best(rows: list[StoreMapping], lat: float | None, lng: float | None) -> StoreMapping:
|
||||
"""同名 + 含目标 id 的候选里挑一条:有入参 geo 且有候选带 geo → 取最近;否则取 created_at 最新。"""
|
||||
if lat is not None and lng is not None:
|
||||
geod = [r for r in rows if r.lat is not None and r.lng is not None]
|
||||
if geod:
|
||||
return min(geod, key=lambda r: _haversine_km(lat, lng, r.lat, r.lng))
|
||||
return max(rows, key=lambda r: r.created_at)
|
||||
|
||||
|
||||
# 目标 key → (该平台店铺 id 列, 组装返回 payload 的函数)。pricebot 拿 id 现拼 deeplink。
|
||||
_TARGETS = {
|
||||
"taobao": ("id_taobao", lambda r: {"shop_id": r.id_taobao, "deeplink": r.taobao_deeplink}),
|
||||
"jd": ("id_jd", lambda r: {"store_id": r.id_jd, "vender_id": r.jd_vender_id, "deeplink": r.jd_deeplink}),
|
||||
"meituan": ("meituan_poi_id_str", lambda r: {"poi_id_str": r.meituan_poi_id_str, "deeplink": r.meituan_deeplink}),
|
||||
}
|
||||
|
||||
|
||||
def lookup_nearest(
|
||||
db: Session, source_platform: str, store_name: str,
|
||||
lat: float | None = None, lng: float | None = None,
|
||||
) -> dict:
|
||||
"""按"源平台店名"反查各目标平台已沉淀的店铺 id。返回 {target_key: {id..., deeplink, row_id, ...}}。
|
||||
- 匹配键 = 源平台对应的 name 列 == store_name(精确)。
|
||||
- 每个目标**分别**取"含该目标 id 的同名候选里最近一条"(淘宝 id / 京东 id 可能在不同行)。
|
||||
- 排除源平台自己(不复用源平台的店铺 id)。命中为空 = 没缓存, pricebot 走现场反查老路。"""
|
||||
name_col = _SOURCE_NAME_COLUMN.get(source_platform)
|
||||
if not name_col or not store_name:
|
||||
return {}
|
||||
rows = db.execute(
|
||||
select(StoreMapping).where(getattr(StoreMapping, name_col) == store_name)
|
||||
).scalars().all()
|
||||
if not rows:
|
||||
return {}
|
||||
src_key = _SOURCE_TARGET_KEY.get(source_platform)
|
||||
out: dict = {}
|
||||
for tgt, (id_attr, make_payload) in _TARGETS.items():
|
||||
if tgt == src_key:
|
||||
continue # 不返回源平台自己
|
||||
cands = [r for r in rows if getattr(r, id_attr)]
|
||||
if not cands:
|
||||
continue
|
||||
best = _pick_best(cands, lat, lng)
|
||||
payload = make_payload(best)
|
||||
payload["row_id"] = best.id
|
||||
if best.lat is not None and best.lng is not None and lat is not None and lng is not None:
|
||||
payload["dist_km"] = round(_haversine_km(lat, lng, best.lat, best.lng), 3)
|
||||
out[tgt] = payload
|
||||
return out
|
||||
@@ -24,13 +24,6 @@ class ComparisonItemIn(BaseModel):
|
||||
specs: list[str] | None = None
|
||||
|
||||
|
||||
class AppliedCouponIn(BaseModel):
|
||||
"""单笔已用优惠(来自 comparison_results[].applied_coupons)。amount 单位:元、正数。"""
|
||||
|
||||
name: str
|
||||
amount: float
|
||||
|
||||
|
||||
class ComparisonResultIn(BaseModel):
|
||||
"""逐平台对比项(来自 done.params.comparison_results)。price 单位:元。"""
|
||||
|
||||
@@ -47,10 +40,6 @@ class ComparisonResultIn(BaseModel):
|
||||
# 优惠**来源名**(展示用, best-effort): 美团"外卖大额神券"/京东"百亿补贴"/淘宝"平台红包"。
|
||||
# None=没抠到 → 前端走通用"红包"。同样必须显式声明否则上报边界被 pydantic 静默丢弃(pricebot#38 引入)。
|
||||
coupon_name: str | None = None
|
||||
# 多券明细 [{name, amount}](全口径: 平台红包+商家券+满减+配送减免, amount 单位元正数)。
|
||||
# 跟 coupon_saved 并存, 是更丰富的明细; 空=没抠到 → 前端回退单券路径。
|
||||
# 必须显式声明: 落库走 model_dump(), pydantic 默认丢未知字段, 不声明这行会被悄悄吞掉。
|
||||
applied_coupons: list[AppliedCouponIn] = Field(default_factory=list)
|
||||
|
||||
|
||||
class ComparisonRecordIn(BaseModel):
|
||||
|
||||
@@ -7,34 +7,16 @@ from pydantic import BaseModel
|
||||
class CouponPromptDismissIn(BaseModel):
|
||||
"""客户端拒绝/关闭领券引导窗的通知体。
|
||||
|
||||
server 据此记一条今日 engagement(dismissed)。2026-06-12 频控按 App 为单位:
|
||||
package = 在哪个外卖 App 关的窗(Android 包名);旧客户端不带 → None(全局兜底行)。
|
||||
server 据此记一条今日 engagement(dismissed)→ 今天这台设备不再弹引导窗。
|
||||
MVP 不鉴权,按 device_id 判断;user_id 登录态带上就一并记(资产),可空。
|
||||
"""
|
||||
|
||||
device_id: str
|
||||
user_id: int | None = None
|
||||
package: str | None = None
|
||||
|
||||
|
||||
class CouponPromptShownIn(BaseModel):
|
||||
"""领券引导窗已对某 App 弹出(shown)的通知体——按 App 频控的主判据。
|
||||
|
||||
客户端弹窗一亮即上报:该 App 今天不再弹(「一个 app 一天只能弹一次,
|
||||
时机=每天第一次进入」);领/拒后续上报只升级 engage_type,不影响频控。
|
||||
"""
|
||||
|
||||
device_id: str
|
||||
user_id: int | None = None
|
||||
package: str | None = None
|
||||
|
||||
|
||||
class CouponPromptShouldShowOut(BaseModel):
|
||||
"""切到外卖 App 时是否还应弹领券引导窗。
|
||||
|
||||
false 的两种来源:该 App 今天已弹过(per-App engagement)/ 今天已跑完整轮领券
|
||||
(completion,彻底领完全局不弹)。
|
||||
"""
|
||||
"""切到外卖 App 时是否还应弹领券引导窗。今天已 engage(领或拒)过 → false。"""
|
||||
|
||||
should_show: bool
|
||||
|
||||
|
||||
@@ -1,36 +0,0 @@
|
||||
"""设备注册 / 心跳相关 schema(无障碍保护存活检测)。"""
|
||||
from __future__ import annotations
|
||||
|
||||
from datetime import datetime
|
||||
|
||||
from pydantic import BaseModel, ConfigDict
|
||||
|
||||
|
||||
class DeviceRegisterRequest(BaseModel):
|
||||
device_id: str
|
||||
registration_id: str | None = None
|
||||
platform: str = "android"
|
||||
app_version: str | None = None
|
||||
|
||||
|
||||
class HeartbeatRequest(BaseModel):
|
||||
device_id: str
|
||||
source: str = "service" # service | app
|
||||
accessibility_enabled: bool = True
|
||||
registration_id: str | None = None
|
||||
|
||||
|
||||
class DeviceOut(BaseModel):
|
||||
model_config = ConfigDict(from_attributes=True)
|
||||
|
||||
id: int
|
||||
device_id: str
|
||||
registration_id: str | None
|
||||
ever_protected: bool
|
||||
liveness_state: str
|
||||
last_heartbeat_at: datetime | None
|
||||
updated_at: datetime | None
|
||||
|
||||
|
||||
class OkResponse(BaseModel):
|
||||
ok: bool = True
|
||||
@@ -1,62 +0,0 @@
|
||||
"""平台店铺映射内部上报的收发模型。
|
||||
|
||||
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
|
||||
|
||||
# 美团原料(poi_id_str 非稳定主键, 单列存; 见 model 注释)
|
||||
meituan_poi_id_str: str | None = None
|
||||
meituan_share_url: str | None = None
|
||||
meituan_resolved_url: str | None = None
|
||||
meituan_deeplink: str | None = None
|
||||
|
||||
# 京东原料(venderId 单列存, storeId 进 id_jd; 见 model 注释)
|
||||
jd_vender_id: str | None = None
|
||||
jd_share_url: str | None = None
|
||||
jd_resolved_url: str | None = None
|
||||
jd_deeplink: str | None = None
|
||||
|
||||
attrs: dict | None = None
|
||||
|
||||
|
||||
class StoreMappingOut(BaseModel):
|
||||
"""上报结果。inserted=1 为新建该 trace 行,0 为合并进已存在行(填空,不覆盖);
|
||||
row_id 为该 trace 对应行 id。"""
|
||||
|
||||
inserted: int
|
||||
row_id: int | None = None
|
||||
@@ -2,7 +2,7 @@
|
||||
|
||||
> 跨表视角。单表字段级细节看同目录 `<表名>.md`(索引见 [README](./README.md))。
|
||||
> 本文专门回答三件「跨表」的事:**① 每块 App 功能用到哪些表 ② 什么操作往哪张表写 ③ 表和表怎么连(join key,含没有外键约束、靠业务字段对齐的语义关联)**。
|
||||
> **范围**:业务表全部在 `shaguabijia-app-server`(SQLAlchemy 2.0 + SQLite 开发 / PostgreSQL 生产)。`pricebot-backend`(比价/领券 Agent)是纯内存态、**无任何表**;Android 客户端只有 EncryptedSharedPreferences / SharedPreferences、**无关系库**。共 **28 张业务表** + `alembic_version`(框架的迁移版本指针)。领券联动的「今日状态」三张表(`coupon_*`)同理:领券过程在 pricebot 内存态跑、**不落库**,只有结果回到 app-server 才落这三张表。
|
||||
> **范围**:业务表全部在 `shaguabijia-app-server`(SQLAlchemy 2.0 + SQLite 开发 / PostgreSQL 生产)。`pricebot-backend`(比价/领券 Agent)是纯内存态、**无任何表**;Android 客户端只有 EncryptedSharedPreferences / SharedPreferences、**无关系库**。共 **23 张业务表** + `alembic_version`(框架的迁移版本指针)。
|
||||
|
||||
---
|
||||
|
||||
@@ -17,14 +17,6 @@
|
||||
| profile「累计省了 / 省钱战绩 / 省钱明细」 | [`savings_record`](./savings_record.md) | 真实下单归因(source=compare)+ 无真实数据时 demo 兜底 |
|
||||
| 「上报更低价」提交 / 列表 | [`price_report`](./price_report.md) | 众包纠偏:用户举证某平台更便宜,人工审核发奖 |
|
||||
|
||||
### 领券(每日领券联动 · 今日状态)
|
||||
| App 位置 / 动作 | 表 | 说明 |
|
||||
|---|---|---|
|
||||
| 领券**过程**(看屏→领券) | (无) | 在 pricebot-backend 内存态跑,**过程不落库**;结果回 app-server 才落下面三张表 |
|
||||
| 切外卖 App 时是否弹领券引导窗 | [`coupon_prompt_engagement`](./coupon_state.md) | 今天 engage 过(点领/点拒)就不再弹;判断维度 device_id |
|
||||
| 首页「去领取」卡是否置灰 | [`coupon_daily_completion`](./coupon_state.md) | 今天跑完整轮(到 done)就置灰;判断维度 device_id |
|
||||
| 每张券领取结果留痕 | [`coupon_claim_record`](./coupon_state.md) | 资产/画像/排查/CPS;当前**不参与**判断 |
|
||||
|
||||
### 钱包 / 福利(看广告赚钱闭环)
|
||||
| App 位置 / 动作 | 表 | 说明 |
|
||||
|---|---|---|
|
||||
@@ -84,13 +76,6 @@
|
||||
| 首次进 profile 省钱页且无真实记录 | `savings_record`(C `source=demo`) | 懒种子,`ensure_seeded` 按 user 幂等 |
|
||||
| 上报更低价 `POST /report` | `price_report`(C) | 读 `comparison_record.best_price_cents` 校验 |
|
||||
| 提交反馈 `POST /feedback` | `feedback`(C) | |
|
||||
| 领券首帧 `POST /api/v1/coupon/step`(step=0) | `coupon_prompt_engagement`(C/U `claim_started`) | `(device_id, 北京日)` 幂等;best-effort |
|
||||
| 领券每帧结果 `POST /api/v1/coupon/step` | `coupon_claim_record`(C/U) | `(device_id, coupon_id, 北京日)` 幂等;best-effort |
|
||||
| 领券跑完 `POST /api/v1/coupon/step`(action.command=done) | `coupon_daily_completion`(C/U) | `(device_id, 北京日)` 幂等;best-effort |
|
||||
| 拒绝领券引导窗 `POST /api/v1/coupon/prompt/dismiss` | `coupon_prompt_engagement`(C/U `dismissed`) | 同上;客户端通知(透传链路看不到拒绝) |
|
||||
| 重置今日弹窗 `POST /api/v1/coupon/prompt/reset`(开发) | `coupon_prompt_engagement`(**D** 今日条) | 删后今天又能弹 |
|
||||
|
||||
> 领券三表写库**全 best-effort**:`/coupon/step` 里写失败只 `logger.warning`、不连累领券返回;**判断只看 `device_id`**,`user_id` 可空旁路(资产留痕)。
|
||||
|
||||
### admin 端(管理员触发,均额外写一条 `admin_audit_log`)
|
||||
| 后台操作 | 写入 | 操作 |
|
||||
@@ -137,7 +122,6 @@
|
||||
- **`comparison_record.store_name` ≈ `savings_record.shop_name`**:无 id 关联,按**店名字符串相等**给比价记录打「已下单」标记(瞬态,不写库)。两边店名同源 = 比价意图识别阶段的门店 query,语义=**店级**(同店比价多次会一并标已下单)。
|
||||
- **广告流会话关联**:`ad_reward_record.ad_session_id` 可与 `ad_ecpm_record.ad_session_id` 对齐;`ad_watch_log` 仍是旧版兼容统计,不逐条参与发奖。
|
||||
- **里程碑解锁进度不存库**:`comparison_milestone_claim` 只记「哪几档已领」;进度 = `comparison_record` 里 `status='success'` 的 `count`。
|
||||
- **领券三表无硬 FK,全靠软关联**:`coupon_prompt_engagement` / `coupon_daily_completion` / `coupon_claim_record` 的 `user_id` **软指** `user.id`(可空、有登录态才记、不进唯一键、不阻塞判断);`trace_id` **软指** pricebot work_logs(排查回指);唯一键都以 `device_id` + 北京自然日为主(详见 [`coupon_state.md`](./coupon_state.md))。
|
||||
- **`onboarding_completion.(user_id, device_id)`**:`user_id` 语义关联 `user.id`(无硬 FK,同 `coupon_*` 设备表),`device_id` = 客户端硬件级 `ANDROID_ID`(≠ 领券 per-install `device_id`)。登录读、走完引导写,决定是否再展示新手引导。
|
||||
|
||||
### ER 关系(文字版)
|
||||
@@ -152,8 +136,6 @@ user ─1:N─ onboarding_completion (user_id, 无硬 FK; (user_id,d
|
||||
comparison_record ─1:N─ price_report (comparison_record_id, 可空)
|
||||
admin_user ─1:N─ admin_audit_log
|
||||
app_config (独立, 无外键, key 为主键)
|
||||
coupon_prompt_engagement / coupon_daily_completion / coupon_claim_record
|
||||
(独立, 无硬 FK; 维度=device_id+北京日, user_id/trace_id 仅软关联)
|
||||
```
|
||||
|
||||
---
|
||||
|
||||
@@ -3,13 +3,16 @@
|
||||
> 数据库:SQLite 起步(`data/app.db`),生产可切 PostgreSQL(改 `DATABASE_URL`)。
|
||||
> ORM:SQLAlchemy 2.0(`app/models/`),迁移:Alembic(`alembic/versions/`,`render_as_batch` 兼容 SQLite)。
|
||||
> 金额字段一律存**整数**:金币=个数,现金=**分**(`*_cents`)。时间列 `DateTime(timezone=True)`。
|
||||
> 最后更新:2026-06-11(合并:新增 3 张领券今日状态表 `coupon_*` + `onboarding_completion` 新手引导完成表;含 [OVERVIEW 总览](./OVERVIEW.md))
|
||||
|
||||
> 最后更新:2026-06-10(新增 `onboarding_completion` 新手引导完成表 → 23 张业务表)
|
||||
=======
|
||||
|
||||
|
||||
> 🧭 **先看 [OVERVIEW.md — 表 × 功能 × 关系](./OVERVIEW.md)**:跨表的「每块功能用哪些表 / 什么操作写哪张表 / 表间 join key」都在那;本页只做**单表索引**,点进每张表的详情看字段级说明。
|
||||
|
||||
---
|
||||
|
||||
## 表总览(28 张业务表 + `alembic_version` 框架表)
|
||||
## 表总览(23 张业务表 + `alembic_version` 框架表)
|
||||
|
||||
### 账号 / 反馈
|
||||
| 表 | 用途 | 模型 | 文档 |
|
||||
@@ -42,13 +45,6 @@
|
||||
| `savings_record` | 省钱记录(profile 省钱战绩源;真实下单归因 + demo) | `models/savings.py` | [详情](./savings_record.md) |
|
||||
| `price_report` | 上报更低价(众包纠偏,人工审核发奖) | `models/price_report.py` | [详情](./price_report.md) |
|
||||
|
||||
### 领券(每日领券联动 · 今日状态)
|
||||
| 表 | 用途 | 模型 | 文档 |
|
||||
|---|---|---|---|
|
||||
| `coupon_prompt_engagement` | 领券引导窗频控源(今日是否已 engage,按 device+日) | `models/coupon_state.py` | [详情](./coupon_state.md) |
|
||||
| `coupon_daily_completion` | 首页「去领取」置灰源(今日是否已跑完整轮) | `models/coupon_state.py` | [详情](./coupon_state.md) |
|
||||
| `coupon_claim_record` | 每张券领取结果沉淀(资产/画像/排查,不参与判断) | `models/coupon_state.py` | [详情](./coupon_state.md) |
|
||||
|
||||
### 美团 CPS 券缓存
|
||||
| 表 | 用途 | 模型 | 文档 |
|
||||
|---|---|---|---|
|
||||
|
||||
@@ -1,123 +0,0 @@
|
||||
# coupon_state — 领券今日状态三张表(弹窗频控 / 首页置灰 / 领券记录)
|
||||
|
||||
> 模型 `app/models/coupon_state.py` · 仓库 `app/repositories/coupon_state.py` · 接口 `app/api/v1/coupon.py`(prefix `/api/v1/coupon`) · [← 索引](./README.md) · [总览](./OVERVIEW.md)
|
||||
|
||||
领券(优惠券自动化)联动产生的三张「今日状态」表,都挂在领券透传端点 `POST /api/v1/coupon/step` 这条链路上(pricebot 跑领券,结果回 app-server 落库;**领券过程本身在 pricebot 内存态跑、不落库**)。三表各管一件事:
|
||||
|
||||
- **`coupon_prompt_engagement`** — 弹窗频控源。按 `(device, 自然日)` 记「今天是否对领券引导窗表达过**意向**」(点「一键领取」=`claim_started` / 点拒绝关闭=`dismissed` 都算)。切到外卖 App 时据此决定弹不弹:今天 engage 过就不再弹。
|
||||
- **`coupon_daily_completion`** — 首页置灰源。按 `(device, 自然日)` 记「今天是否已**跑完整轮**领券(到 done 帧)」。首页「去领取」卡据此置灰:今天跑完了就不能再领。
|
||||
- **`coupon_claim_record`** — 资产沉淀层。按 `(device, 券, 自然日)` 记每张券的领取结果(success/already_claimed/failed/skipped),**纯沉淀**(资产/画像/排查/CPS 归因),当前**不参与**「要不要领 / 弹不弹」的判断。
|
||||
|
||||
三表共同口径:
|
||||
- **判断维度是 `device_id`,不是 `user_id`**:券发到的是设备上登录的那个外卖账号,device 比 user 更贴近「哪个登录环境」,且 `device_id` 全链路现成、不依赖领券鉴权(领券 MVP 阶段 `/coupon/step` 不鉴权)。客户端 `getOrCreateDeviceId` 生成存 SP,**卸载重装会变 → 当新设备重新弹一次**(产品预期)。
|
||||
- **日期 = `Asia/Shanghai` 自然日**(`claim_date` / `engage_date` / `complete_date`,`repositories/coupon_state.today_cn()`)。每日可领的券(签到/天天红包)靠这天然每天一条。
|
||||
- **`user_id` 可空**:领券登录态有就记(资产/画像),可空、**不进唯一键、不阻塞判断**。
|
||||
- **`trace_id` 可空**:回指 pricebot work_logs,供排查(哪次任务领的)。
|
||||
- **engagement vs completion vs claim 的区别**:engagement = 用户**表达过意向**(点了领或拒,不管跑没跑完);completion = 这一轮**真跑到了 done**(整套流程走完);claim = **每张券一条**的结果留痕。
|
||||
|
||||
> 写库全部 **best-effort**:在 `/coupon/step` 里写库失败只 `logger.warning`、**绝不连累领券主流程/返回**(见 `app/api/v1/coupon.py`)。
|
||||
|
||||
---
|
||||
|
||||
## coupon_prompt_engagement — 弹窗频控(今日是否已对引导窗表达意向)
|
||||
|
||||
`(device_id, engage_date)` 唯一,一台设备一天一条;今天 engage 过(领或拒)就不再弹。
|
||||
|
||||
### 用在哪 / 增删改查
|
||||
- **C / U(幂等 upsert)**:`mark_engagement`。两条触发:
|
||||
- `POST /api/v1/coupon/step` 且 `step==0`(领券首帧=用户已发起领券)→ 记 `claim_started`;
|
||||
- `POST /api/v1/coupon/prompt/dismiss`(用户点关闭引导窗;server 在透传链路看不到「拒绝」,必须客户端通知)→ 记 `dismissed`。
|
||||
- 已有今天那条则覆盖 `engage_type`(并补 `user_id`),否则插入。
|
||||
- **D**:`POST /api/v1/coupon/prompt/reset`(`reset_today_engagement`)—— 删这台设备今天那条,开发设置「重置今日领券弹窗状态」按钮调,测频控用;删后今天又能弹。
|
||||
- **R**:`GET /api/v1/coupon/prompt/should-show?device_id=…`(`has_engaged_today`)→ `should_show = not 今天已 engage`。客户端切外卖 App 前查,纯后台判据。
|
||||
|
||||
### 字段
|
||||
| 列 | 类型 | 约束 / 默认 | 说明(取值 / join) |
|
||||
|---|---|---|---|
|
||||
| `id` | Integer | PK, autoincrement | |
|
||||
| `device_id` | String(64) | NOT NULL | 判断/聚合维度;客户端 `getOrCreateDeviceId`,重装会变 |
|
||||
| `user_id` | Integer | index, 可空 | 登录态有就记(资产);不进唯一键、不阻塞判断 |
|
||||
| `engage_date` | **Date** | NOT NULL | **北京时间**自然日(`today_cn()`) |
|
||||
| `engage_type` | String(16) | NOT NULL | `claim_started`(点一键领取)/ `dismissed`(点拒绝关闭);仅记录区分,**判断只看「今天有没有这条」,type 不影响弹不弹** |
|
||||
| `created_at` | DateTime(tz) | server_default now() | |
|
||||
| `updated_at` | DateTime(tz) | server_default now(), onupdate now() | |
|
||||
|
||||
### 索引与约束
|
||||
- PK `id`;index `user_id`;UNIQUE(`device_id`, `engage_date`) = `uq_coupon_engage_device_date`(一台设备一天一条)。
|
||||
|
||||
### 注意
|
||||
- `device_id` 重装会变 → 重装当新设备,今天重新弹一次(产品预期)。
|
||||
- 判断只看「今天这台设备有没有这条」,不看 `engage_type`(领或拒都算 engage 过、都不再弹)。
|
||||
|
||||
---
|
||||
|
||||
## coupon_daily_completion — 首页置灰(今日是否已跑完整轮领券)
|
||||
|
||||
`(device_id, complete_date)` 唯一,一台设备一天一条;今天跑完整轮(到 done 帧)就把首页「去领取」卡置灰。
|
||||
|
||||
### 用在哪 / 增删改查
|
||||
- **C / U(幂等 upsert)**:`mark_completed_today`,由 `POST /api/v1/coupon/step` 在 pricebot 返回 `action.command == "done"` 那帧调。pricebot 把中途单券 done 改写成 `wait+continue=true`,只有整套全跑完那帧才保留 `command=="done"`,故 **done 已等价「整轮完成」**。已有今天那条则补 `user_id`/`trace_id`,否则插入。
|
||||
- **U / D**:无业务删除。
|
||||
- **R**:`GET /api/v1/coupon/completed-today?device_id=…`(`has_completed_today`)→ `completed`。客户端据此把首页「去领取」卡置灰、不可点。
|
||||
|
||||
### 字段
|
||||
| 列 | 类型 | 约束 / 默认 | 说明(取值 / join) |
|
||||
|---|---|---|---|
|
||||
| `id` | Integer | PK, autoincrement | |
|
||||
| `device_id` | String(64) | NOT NULL | 判断维度,与 engagement/claim 一致;客户端两端都用 ANDROID_ID |
|
||||
| `user_id` | Integer | index, 可空 | 登录态有就记(资产) |
|
||||
| `complete_date` | **Date** | NOT NULL | **北京时间**自然日(`today_cn()`) |
|
||||
| `trace_id` | String(64) | index, 可空 | 哪次任务跑到 done,回指 pricebot work_logs / 排查 |
|
||||
| `created_at` | DateTime(tz) | server_default now() | |
|
||||
| `updated_at` | DateTime(tz) | server_default now(), onupdate now() | |
|
||||
|
||||
### 索引与约束
|
||||
- PK `id`;index `user_id`、`trace_id`;UNIQUE(`device_id`, `complete_date`) = `uq_coupon_completion_device_date`(一台设备一天一条)。
|
||||
|
||||
### 注意
|
||||
- 口径(用户决策 2026-06-10 A 方案):**到 done 即算完成,不管单券成败**——失败/跳过常是无障碍/环境问题,重复点也补不回来。
|
||||
- 与 engagement 区别:engagement 是「表达过意向」(点了就记,不管跑没跑完);completion 是「真跑到了 done」。
|
||||
|
||||
---
|
||||
|
||||
## coupon_claim_record — 领券记录(每张券一天一条,资产沉淀层)
|
||||
|
||||
`(device_id, coupon_id, claim_date)` 唯一,同设备同券同一天只一条;纯沉淀,**当前不参与判断**,留作以后按券去重 / CPS 归因 / 用户画像的数据源。
|
||||
|
||||
### 用在哪 / 增删改查
|
||||
- **C / U(幂等 upsert)**:`record_claims`,由 `POST /api/v1/coupon/step` 写入。一帧的券结果来自 pricebot 的 `last_coupon_result`(最后一张)+ `action.params.coupon_results`(全量)——**会重复带同一张券**,端点 `_extract_coupon_results` 先**按 `coupon_id` 去重**(全量覆盖单张),仓库再靠唯一键幂等:已有则更新 `status`/`reason`/`claimed_count`/`extra`(以最后一次为准),否则插入。
|
||||
- **U / D**:无业务删除。
|
||||
- **R**:**当前无读取端点**(纯写入沉淀,未来做去重/归因/画像时再用)。
|
||||
|
||||
### 字段
|
||||
| 列 | 类型 | 约束 / 默认 | 说明(取值 / join) |
|
||||
|---|---|---|---|
|
||||
| `id` | Integer | PK, autoincrement | |
|
||||
| `device_id` | String(64) | NOT NULL | 聚合维度;客户端 `getOrCreateDeviceId`,重装会变 |
|
||||
| `user_id` | Integer | index, 可空 | 登录态有就记(资产/画像);不进唯一键 |
|
||||
| `coupon_id` | String(64) | NOT NULL | 券标识(取自 pricebot 结果) |
|
||||
| `claim_date` | **Date** | NOT NULL | **北京时间**自然日(`today_cn()`);每日可领的券靠它天然每天一条 |
|
||||
| `status` | String(24) | NOT NULL | `success` / `already_claimed` / `failed` / `skipped`(原样取 pricebot coupon 结果) |
|
||||
| `vendor` | String(48) | 可空 | 券提供方 |
|
||||
| `coupon_name` | String(128) | 可空 | 取 pricebot `name` |
|
||||
| `claimed_count` | Integer | 可空 | 这张领到几张(pricebot `display_count`,给不出时 None;兼容 `claimed_count`) |
|
||||
| `trace_id` | String(64) | index, 可空 | 哪次任务领的,回指 pricebot work_logs / 排查 |
|
||||
| `reason` | String(255) | 可空 | failed / skipped 原因 |
|
||||
| `extra` | JSON(PG JSONB) | 可空 | 杂项兜底:券的结构化信息(面额/入口/关键节点摘要等),免得加字段就迁移;当前直接存这帧 pricebot 单券结果 dict |
|
||||
| `created_at` | DateTime(tz) | server_default now() | |
|
||||
| `updated_at` | DateTime(tz) | server_default now(), onupdate now() | |
|
||||
|
||||
### 索引与约束
|
||||
- PK `id`;index `user_id`、`trace_id`;UNIQUE(`device_id`, `coupon_id`, `claim_date`) = `uq_coupon_claim_device_coupon_date`;Index(`device_id`, `claim_date`) = `ix_coupon_claim_device_date`(按 device+日 取一天所有券)。
|
||||
|
||||
### 注意
|
||||
- **`extra` 别塞原始无障碍树**(几十 KB → 行膨胀);原始大树看 `trace_id` 指过去的 work_logs。
|
||||
- 同批去重很关键:`autoflush=False` 下同 `coupon_id` 两次 `add` 会撞唯一约束、`IntegrityError` 回滚整批(done/单券记录全丢),故端点 `_extract_coupon_results` + 仓库 `seen` 集合双重防御。
|
||||
- `extra` 是 `JSON().with_variant(JSONB(), "postgresql")`:PG 用 JSONB(可建 GIN 索引),SQLite 退化通用 JSON(同 `price_observation` / `comparison_record`)。
|
||||
|
||||
---
|
||||
|
||||
## 三表共性小结
|
||||
- 数据流向:客户端 → `POST /api/v1/coupon/step`(透传给 pricebot)→ 结果回写这三张表(best-effort,写库失败不影响领券)。
|
||||
- 唯一键都含 `device_id` + 某个北京自然日列;`user_id` 永远是可空旁路(资产留痕,不进唯一键、不阻塞判断)。
|
||||
- 无硬外键:`user_id` 软指 `user.id`、`trace_id` 软指 pricebot work_logs(详见 [OVERVIEW → 表间关系 & Join Key](./OVERVIEW.md))。
|
||||
Reference in New Issue
Block a user