Compare commits

..

3 Commits

Author SHA1 Message Date
OuYingJun1024 ea27fe0c01 feat(admin/users): 用户列表支持列排序 + 渠道/昵称/时间范围筛选
list_users 由 keyset 游标改为 offset 分页以支持任意列排序:
- 排序:sort_by ∈ {id, created_at, last_login_at} + sort_order(白名单,非法回落 id;同值 id 同向兜底)
- 筛选:新增 nickname 模糊、注册/最近登录时间范围(register_channel 本就支持,这次在 UI/文档补齐)
- 文档 admin-users-list.md 与 test_admin_read 同步(排序/渠道/昵称/时间范围/非法 sort_by 422)

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
2026-06-13 21:41:10 +08:00
OuYingJun1024 995908cd0b feat(rewards/admin): LT 因子改账号累计 + 审计 only_mismatch + 用户金币现金 set 模式
- 看视频奖励 LT 因子(因子2)按「账号累计第 N 次看视频」递减,不再按天重置:
  ad_reward 新增 _granted_cumulative、ad_feed_reward 单位序号去掉当天过滤、
  rewards 形参 today_count_after_this → count_after_this;每日次数上限/冷却仍按当日统计
- 广告金币审计:加 only_mismatch 筛选只看 ✗ 行;total/mismatch_count 改全量统计
  (不受 limit/筛选影响、截断前算)+ 新增 truncated 标记展示集是否被截断
- admin 用户金币/现金接口加 mode=delta|set:set=设为目标值(读余额算差值、仍写一笔流水,
  沿用原子/审计/扣负保护);新增 admin-user-cash 文档 + 更新 API 索引/coins 等文档 + 补 admin read/write 测试
2026-06-13 21:26:11 +08:00
OuYingJun1024 fb661814cb fix(marquee): 首页轮播脱敏名改为按 user_id 恒定 + 去 Faker 依赖
- _mask_real 无昵称分支用 random.Random(user_id) 播种合成假名 → 同一用户每次展示恒定
  (刷新/翻页不变脸)、不同用户各异;修复原全局 _rng 导致同一用户每请求换名(像假数据)的问题
- 去掉 faker 依赖,改用本地姓池(~100)×名字字池(~120)组合真名 + 中英网络昵称语料,组合空间上万、
  同屏几乎不撞名;_synth_full_name/_synth_masked_name 串入 rng 参数,种子/兜底仍用全局 _rng 出多样
- 同步 docs/api/platform-savings-feed.md、docs/database/ops_marquee_seed.md 脱敏描述
  (去手机尾号/Faker 表述,补 user_id 恒定 + JOIN user 仅活跃用户记录说明)

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
2026-06-13 17:27:05 +08:00
36 changed files with 123 additions and 1048 deletions
-40
View File
@@ -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')
-77
View File
@@ -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')
+3 -3
View File
@@ -50,7 +50,7 @@ def _reward_video_rows(
AdRewardRecord.reward_date == date, AdRewardRecord.reward_date == date,
AdRewardRecord.reward_scene == "reward_video", AdRewardRecord.reward_scene == "reward_video",
) )
.order_by(AdRewardRecord.user_id, AdRewardRecord.created_at, AdRewardRecord.id) .order_by(AdRewardRecord.user_id, AdRewardRecord.created_at)
) )
if user_id is not None: if user_id is not None:
stmt = stmt.where(AdRewardRecord.user_id == user_id) stmt = stmt.where(AdRewardRecord.user_id == user_id)
@@ -127,7 +127,7 @@ def _feed_rows(db: Session, *, date: str, user_id: int | None) -> list[dict]:
stmt = ( stmt = (
select(AdFeedRewardRecord) select(AdFeedRewardRecord)
.where(AdFeedRewardRecord.reward_date == date) .where(AdFeedRewardRecord.reward_date == date)
.order_by(AdFeedRewardRecord.user_id, AdFeedRewardRecord.created_at, AdFeedRewardRecord.id) .order_by(AdFeedRewardRecord.user_id, AdFeedRewardRecord.created_at)
) )
if user_id is not None: if user_id is not None:
stmt = stmt.where(AdFeedRewardRecord.user_id == user_id) stmt = stmt.where(AdFeedRewardRecord.user_id == user_id)
@@ -205,7 +205,7 @@ def ad_coin_audit(
rows.extend(_reward_video_rows(db, date=date, user_id=user_id)) rows.extend(_reward_video_rows(db, date=date, user_id=user_id))
if scene in (None, "feed"): if scene in (None, "feed"):
rows.extend(_feed_rows(db, date=date, user_id=user_id)) rows.extend(_feed_rows(db, date=date, user_id=user_id))
rows.sort(key=lambda r: (r["created_at"], r["record_id"]), reverse=True) rows.sort(key=lambda r: r["created_at"], reverse=True)
total = len(rows) total = len(rows)
mismatch_count = sum(1 for r in rows if not r["matched"]) mismatch_count = sum(1 for r in rows if not r["matched"])
+11 -14
View File
@@ -4,7 +4,7 @@
""" """
from __future__ import annotations from __future__ import annotations
from sqlalchemy import func, select from sqlalchemy import select
from sqlalchemy.orm import Session from sqlalchemy.orm import Session
from app.models.admin import AdminAuditLog from app.models.admin import AdminAuditLog
@@ -49,9 +49,8 @@ def list_audit_logs(
admin_id: int | None = None, admin_id: int | None = None,
limit: int = 50, limit: int = 50,
cursor: int | None = None, cursor: int | None = None,
) -> tuple[list[AdminAuditLog], int | None, int]: ) -> tuple[list[AdminAuditLog], int | None]:
"""offset 分页(id 倒序)+ total。cursor 即 offset((page-1)*pageSize),支持页码跳页。 """游标分页(id 倒序),与现有 list_* 约定一致。返回 (rows, next_cursor)。"""
返回 (rows, next_cursor, total)。"""
stmt = select(AdminAuditLog) stmt = select(AdminAuditLog)
if action: if action:
stmt = stmt.where(AdminAuditLog.action == action) stmt = stmt.where(AdminAuditLog.action == action)
@@ -59,15 +58,13 @@ def list_audit_logs(
stmt = stmt.where(AdminAuditLog.target_type == target_type) stmt = stmt.where(AdminAuditLog.target_type == target_type)
if admin_id is not None: if admin_id is not None:
stmt = stmt.where(AdminAuditLog.admin_id == admin_id) stmt = stmt.where(AdminAuditLog.admin_id == admin_id)
if cursor is not None:
total = int(db.execute(select(func.count()).select_from(stmt.subquery())).scalar_one()) stmt = stmt.where(AdminAuditLog.id < cursor)
offset = max(cursor or 0, 0) stmt = stmt.order_by(AdminAuditLog.id.desc())
rows = list( rows = list(db.execute(stmt.limit(limit + 1)).scalars().all())
db.execute(
stmt.order_by(AdminAuditLog.id.desc()).offset(offset).limit(limit + 1)
).scalars().all()
)
has_more = len(rows) > limit has_more = len(rows) > limit
items = rows[:limit] items = rows[:limit]
next_cursor = offset + limit if has_more else None # next_cursor 必须是"本页返回的最后一条"的 id(下一页查 id < 它),不能用 rows[limit]——
return items, next_cursor, total # rows[limit] 是探测下一页用的第 limit+1 条,它既不在本页也不在下页 → 每页边界丢一条。
next_cursor = items[-1].id if has_more else None
return items, next_cursor
+45 -84
View File
@@ -38,28 +38,6 @@ def cursor_paginate(
return items, next_cursor return items, next_cursor
def offset_paginate(
db: Session, stmt: Select, sort_clause: tuple, *, limit: int, cursor: int | None
) -> tuple[list, int | None, int]:
"""offset 分页 + 总数。stmt 只含 where/join,不要预先带 order_by/offset/limit。
cursor offset(页码分页:offset=(page-1)*pageSize)返回 (items, next_cursor, total):
- total:符合筛选条件的总条数( antd pagination 渲染页码/ N ),count P0 量级开销可忽略;
- next_cursor:下一页 offset(兼容加载更多),末页为 None
多取 1 条探测下一页sort_clause order_by 表达式元组(末位应含 id 保证稳定排序)"""
total = int(
db.execute(select(func.count()).select_from(stmt.subquery())).scalar_one()
)
offset = max(cursor or 0, 0)
rows = list(
db.execute(stmt.order_by(*sort_clause).offset(offset).limit(limit + 1)).scalars().all()
)
has_more = len(rows) > limit
items = rows[:limit]
next_cursor = offset + limit if has_more else None
return items, next_cursor, total
def list_users( def list_users(
db: Session, db: Session,
*, *,
@@ -75,11 +53,11 @@ def list_users(
sort_order: str = "desc", sort_order: str = "desc",
limit: int = 20, limit: int = 20,
cursor: int | None = None, cursor: int | None = None,
) -> tuple[list[User], int | None, int]: ) -> tuple[list[User], int | None]:
"""用户列表(admin 全量)。支持手机号前缀 / 渠道 / 状态 / 昵称模糊 / 注册·最近登录时间范围筛选, """用户列表(admin 全量)。支持手机号前缀 / 渠道 / 状态 / 昵称模糊 / 注册·最近登录时间范围筛选,
id·注册时间·最近登录排序**offset 分页**(cursor=offset):任意列排序下游标语义统一, id·注册时间·最近登录排序**offset 分页**(cursor=offset):任意列排序下游标语义统一,
代价是翻页期间数据变动可能错位一条admin 低频场景可接受( [list_all_withdraw_orders]) 代价是翻页期间数据变动可能错位一条admin 低频场景可接受( [list_all_withdraw_orders])
日期入参统一转 tz-aware UTC 比较(列为 timestamptz, _as_utc)""" 日期入参统一转 UTC naive 比较(User 时间均为 UTC naive, _as_utc_naive)"""
stmt = select(User) stmt = select(User)
if phone: if phone:
stmt = stmt.where(User.phone.like(f"{phone}%")) # 前缀匹配 stmt = stmt.where(User.phone.like(f"{phone}%")) # 前缀匹配
@@ -90,13 +68,13 @@ def list_users(
if nickname and nickname.strip(): if nickname and nickname.strip():
stmt = stmt.where(User.nickname.ilike(f"%{nickname.strip()}%")) stmt = stmt.where(User.nickname.ilike(f"%{nickname.strip()}%"))
if created_from is not None: if created_from is not None:
stmt = stmt.where(User.created_at >= _as_utc(created_from)) stmt = stmt.where(User.created_at >= _as_utc_naive(created_from))
if created_to is not None: if created_to is not None:
stmt = stmt.where(User.created_at <= _as_utc(created_to)) stmt = stmt.where(User.created_at <= _as_utc_naive(created_to))
if last_login_from is not None: if last_login_from is not None:
stmt = stmt.where(User.last_login_at >= _as_utc(last_login_from)) stmt = stmt.where(User.last_login_at >= _as_utc_naive(last_login_from))
if last_login_to is not None: if last_login_to is not None:
stmt = stmt.where(User.last_login_at <= _as_utc(last_login_to)) stmt = stmt.where(User.last_login_at <= _as_utc_naive(last_login_to))
sort_cols = { sort_cols = {
"id": User.id, "id": User.id,
@@ -106,7 +84,14 @@ def list_users(
sort_col = sort_cols.get(sort_by, User.id) sort_col = sort_cols.get(sort_by, User.id)
order_fn = asc if sort_order == "asc" else desc order_fn = asc if sort_order == "asc" else desc
id_order = asc(User.id) if sort_order == "asc" else desc(User.id) id_order = asc(User.id) if sort_order == "asc" else desc(User.id)
return offset_paginate(db, stmt, (order_fn(sort_col), id_order), limit=limit, cursor=cursor) stmt = stmt.order_by(order_fn(sort_col), id_order)
offset = max(cursor or 0, 0)
rows = list(db.execute(stmt.offset(offset).limit(limit + 1)).scalars().all())
has_more = len(rows) > limit
items = rows[:limit]
next_cursor = offset + limit if has_more else None
return items, next_cursor
def list_onboarding_devices(db: Session, *, limit: int = 500) -> list[dict]: def list_onboarding_devices(db: Session, *, limit: int = 500) -> list[dict]:
@@ -175,7 +160,7 @@ def list_all_withdraw_orders(
quick_filter: str | None = None, quick_filter: str | None = None,
limit: int = 20, limit: int = 20,
cursor: int | None = None, cursor: int | None = None,
) -> tuple[list[WithdrawOrder], int | None, int]: ) -> tuple[list[WithdrawOrder], int | None]:
stmt = select(WithdrawOrder) stmt = select(WithdrawOrder)
needs_user_join = bool(keyword and keyword.strip()) or quick_filter == "high_risk" needs_user_join = bool(keyword and keyword.strip()) or quick_filter == "high_risk"
if needs_user_join: if needs_user_join:
@@ -205,16 +190,16 @@ def list_all_withdraw_orders(
date_col = WithdrawOrder.updated_at if date_field == "updated_at" else WithdrawOrder.created_at date_col = WithdrawOrder.updated_at if date_field == "updated_at" else WithdrawOrder.created_at
if date_from is not None: if date_from is not None:
stmt = stmt.where(date_col >= _as_utc(date_from)) stmt = stmt.where(date_col >= _as_utc_naive(date_from))
if date_to is not None: if date_to is not None:
stmt = stmt.where(date_col <= _as_utc(date_to)) stmt = stmt.where(date_col <= _as_utc_naive(date_to))
# tz-aware:列为 timestamptz,比较绝对时刻、与 DB 会话时区无关(同 _as_utc / stats.py) now = datetime.now(timezone.utc).replace(tzinfo=None)
now = datetime.now(timezone.utc)
today_start = ( today_start = (
datetime.now(ZoneInfo("Asia/Shanghai")) datetime.now(ZoneInfo("Asia/Shanghai"))
.replace(hour=0, minute=0, second=0, microsecond=0) .replace(hour=0, minute=0, second=0, microsecond=0)
.astimezone(timezone.utc) .astimezone(timezone.utc)
.replace(tzinfo=None)
) )
if quick_filter == "abnormal": if quick_filter == "abnormal":
stmt = stmt.where( stmt = stmt.where(
@@ -259,53 +244,6 @@ def list_all_withdraw_orders(
sort_col = sort_cols.get(sort_by, WithdrawOrder.created_at) sort_col = sort_cols.get(sort_by, WithdrawOrder.created_at)
order_fn = asc if sort_order == "asc" else desc order_fn = asc if sort_order == "asc" else desc
id_order = asc(WithdrawOrder.id) if sort_order == "asc" else desc(WithdrawOrder.id) id_order = asc(WithdrawOrder.id) if sort_order == "asc" else desc(WithdrawOrder.id)
return offset_paginate(db, stmt, (order_fn(sort_col), id_order), limit=limit, cursor=cursor)
def _as_utc(value: datetime) -> datetime:
"""前端传 ISO 时间 → 统一成 tz-aware UTC 再比较。
所有时间列均为 `DateTime(timezone=True)`(Postgres timestamptz); tz-aware 绑定参数
比较的是绝对时刻, DB 会话时区无关恒正确曾用 naive UTC,正确性依赖会话 TimeZone=UTC,
生产会话非 UTC 时筛选边界会整体偏移故统一 tz-aware( stats.py / withdraw_summary 一致)
无时区入参按 UTC 解释"""
if value.tzinfo is None:
return value.replace(tzinfo=timezone.utc)
return value.astimezone(timezone.utc)
def list_feedbacks(
db: Session,
*,
status: str | None = None,
user_id: int | None = None,
content: str | None = None,
created_from: datetime | None = None,
created_to: datetime | None = None,
sort_by: str = "id",
sort_order: str = "desc",
limit: int = 20,
cursor: int | None = None,
) -> tuple[list[Feedback], int | None]:
"""反馈工单列表。支持 状态 / 用户ID / 内容模糊 / 提交时间范围 筛选,按 id·提交时间排序。
**offset 分页**(cursor=offset):任意列排序下游标语义统一( [list_users]),代价是翻页期间
数据变动可能错位一条admin 低频场景可接受created_at timestamptz,日期入参统一转 tz-aware UTC 比较"""
stmt = select(Feedback)
if status:
stmt = stmt.where(Feedback.status == status)
if user_id is not None:
stmt = stmt.where(Feedback.user_id == user_id)
if content and content.strip():
stmt = stmt.where(Feedback.content.ilike(f"%{content.strip()}%"))
if created_from is not None:
stmt = stmt.where(Feedback.created_at >= _as_utc(created_from))
if created_to is not None:
stmt = stmt.where(Feedback.created_at <= _as_utc(created_to))
sort_cols = {"id": Feedback.id, "created_at": Feedback.created_at}
sort_col = sort_cols.get(sort_by, Feedback.id)
order_fn = asc if sort_order == "asc" else desc
id_order = asc(Feedback.id) if sort_order == "asc" else desc(Feedback.id)
stmt = stmt.order_by(order_fn(sort_col), id_order) stmt = stmt.order_by(order_fn(sort_col), id_order)
offset = max(cursor or 0, 0) offset = max(cursor or 0, 0)
@@ -316,6 +254,29 @@ def list_feedbacks(
return items, next_cursor return items, next_cursor
def _as_utc_naive(value: datetime) -> datetime:
"""前端传 ISO 时间;DB 当前按 UTC naive 比较最稳(SQLite/本地开发一致)。"""
if value.tzinfo is None:
return value
return value.astimezone(timezone.utc).replace(tzinfo=None)
def list_feedbacks(
db: Session,
*,
status: str | None = None,
user_id: int | None = None,
limit: int = 20,
cursor: int | None = None,
) -> tuple[list[Feedback], int | None]:
stmt = select(Feedback)
if status:
stmt = stmt.where(Feedback.status == status)
if user_id is not None:
stmt = stmt.where(Feedback.user_id == user_id)
return cursor_paginate(db, stmt, Feedback.id, limit=limit, cursor=cursor)
def get_withdraw_by_out_bill_no(db: Session, out_bill_no: str) -> WithdrawOrder | None: def get_withdraw_by_out_bill_no(db: Session, out_bill_no: str) -> WithdrawOrder | None:
"""按商户单号查提现单(admin 重试打款先拿 user_id 用,M3)。""" """按商户单号查提现单(admin 重试打款先拿 user_id 用,M3)。"""
return db.execute( return db.execute(
@@ -502,14 +463,14 @@ def list_price_reports(
user_id: int | None = None, user_id: int | None = None,
limit: int = 20, limit: int = 20,
cursor: int | None = None, cursor: int | None = None,
) -> tuple[list[PriceReport], int | None, int]: ) -> tuple[list[PriceReport], int | None]:
"""上报更低价列表(admin 全量,可按状态/用户筛)。offset 分页 + total,id 倒序。""" """上报更低价列表(admin 全量,可按状态/用户筛)。游标同 feedback:id 倒序。"""
stmt = select(PriceReport) stmt = select(PriceReport)
if status: if status:
stmt = stmt.where(PriceReport.status == status) stmt = stmt.where(PriceReport.status == status)
if user_id is not None: if user_id is not None:
stmt = stmt.where(PriceReport.user_id == user_id) stmt = stmt.where(PriceReport.user_id == user_id)
return offset_paginate(db, stmt, (PriceReport.id.desc(),), limit=limit, cursor=cursor) return cursor_paginate(db, stmt, PriceReport.id, limit=limit, cursor=cursor)
def price_report_summary(db: Session) -> dict: def price_report_summary(db: Session) -> dict:
+5 -24
View File
@@ -17,12 +17,6 @@ router = APIRouter(
) )
def _active_super_count(db: AdminDb) -> int:
return sum(
1 for a in admin_repo.list_admins(db) if a.role == "super_admin" and a.status == "active"
)
@router.get("", response_model=list[AdminOut], summary="管理员列表") @router.get("", response_model=list[AdminOut], summary="管理员列表")
def list_admins(db: AdminDb) -> list[AdminOut]: def list_admins(db: AdminDb) -> list[AdminOut]:
return [AdminOut.model_validate(a) for a in admin_repo.list_admins(db)] return [AdminOut.model_validate(a) for a in admin_repo.list_admins(db)]
@@ -54,29 +48,16 @@ def update_admin(
if admin_id == admin.id and body.status == "disabled": if admin_id == admin.id and body.status == "disabled":
raise HTTPException(status_code=400, detail="不能禁用自己") raise HTTPException(status_code=400, detail="不能禁用自己")
# 防自锁:降级 / 禁用某个 super_admin 前,确认操作后仍至少剩 1 个 active super_admin,
# 否则会进入「零可用超管」死局——本路由仅 super 可进,只能改库恢复。
demotes_super = (
target.role == "super_admin"
and target.status == "active"
and (
(body.role is not None and body.role != "super_admin")
or body.status == "disabled"
)
)
if demotes_super and _active_super_count(db) <= 1:
raise HTTPException(status_code=400, detail="不能降级/禁用最后一个超级管理员")
changes: dict = {} changes: dict = {}
if body.role is not None and body.role != target.role: if body.role is not None:
changes["role"] = {"before": target.role, "after": body.role}
target.role = body.role target.role = body.role
if body.status is not None and body.status != target.status: changes["role"] = body.role
changes["status"] = {"before": target.status, "after": body.status} if body.status is not None:
target.status = body.status target.status = body.status
changes["status"] = body.status
if body.password is not None: if body.password is not None:
changes["password"] = "reset"
target.password_hash = hash_password(body.password) target.password_hash = hash_password(body.password)
changes["password"] = "reset"
if not changes: if not changes:
raise HTTPException(status_code=400, detail="无任何变更字段") raise HTTPException(status_code=400, detail="无任何变更字段")
db.commit() db.commit()
+2 -4
View File
@@ -26,11 +26,9 @@ def list_audit_logs(
limit: Annotated[int, Query(ge=1, le=100)] = 50, limit: Annotated[int, Query(ge=1, le=100)] = 50,
cursor: Annotated[int | None, Query()] = None, cursor: Annotated[int | None, Query()] = None,
) -> CursorPage[AdminAuditLogOut]: ) -> CursorPage[AdminAuditLogOut]:
items, next_cursor, total = audit_repo.list_audit_logs( items, next_cursor = audit_repo.list_audit_logs(
db, action=action, target_type=target_type, admin_id=admin_id, limit=limit, cursor=cursor, db, action=action, target_type=target_type, admin_id=admin_id, limit=limit, cursor=cursor,
) )
return CursorPage( return CursorPage(
items=[AdminAuditLogOut.model_validate(x) for x in items], items=[AdminAuditLogOut.model_validate(x) for x in items], next_cursor=next_cursor,
next_cursor=next_cursor,
total=total,
) )
+2 -17
View File
@@ -1,7 +1,6 @@
"""admin 反馈工单:列表(读,支持 状态/用户ID/内容/时间 筛选 + 排序)+ 标记已处理(写,带审计)。""" """admin 反馈工单:列表(读)+ 标记已处理(写,带审计)。"""
from __future__ import annotations from __future__ import annotations
from datetime import datetime
from typing import Annotated from typing import Annotated
from fastapi import APIRouter, Depends, HTTPException, Query, Request from fastapi import APIRouter, Depends, HTTPException, Query, Request
@@ -26,25 +25,11 @@ def list_feedbacks(
db: AdminDb, db: AdminDb,
status: Annotated[str | None, Query()] = None, status: Annotated[str | None, Query()] = None,
user_id: Annotated[int | None, Query()] = None, user_id: Annotated[int | None, Query()] = None,
content: Annotated[str | None, Query(max_length=100)] = None,
created_from: Annotated[datetime | None, Query()] = None,
created_to: Annotated[datetime | None, Query()] = None,
sort_by: Annotated[str, Query(pattern="^(id|created_at)$")] = "id",
sort_order: Annotated[str, Query(pattern="^(asc|desc)$")] = "desc",
limit: Annotated[int, Query(ge=1, le=100)] = 20, limit: Annotated[int, Query(ge=1, le=100)] = 20,
cursor: Annotated[int | None, Query()] = None, cursor: Annotated[int | None, Query()] = None,
) -> CursorPage[FeedbackOut]: ) -> CursorPage[FeedbackOut]:
items, next_cursor = queries.list_feedbacks( items, next_cursor = queries.list_feedbacks(
db, db, status=status, user_id=user_id, limit=limit, cursor=cursor,
status=status,
user_id=user_id,
content=content,
created_from=created_from,
created_to=created_to,
sort_by=sort_by,
sort_order=sort_order,
limit=limit,
cursor=cursor,
) )
return CursorPage( return CursorPage(
items=[FeedbackOut.model_validate(f) for f in items], next_cursor=next_cursor, items=[FeedbackOut.model_validate(f) for f in items], next_cursor=next_cursor,
+4 -8
View File
@@ -40,13 +40,11 @@ def list_price_reports(
limit: Annotated[int, Query(ge=1, le=100)] = 20, limit: Annotated[int, Query(ge=1, le=100)] = 20,
cursor: Annotated[int | None, Query()] = None, cursor: Annotated[int | None, Query()] = None,
) -> CursorPage[PriceReportOut]: ) -> CursorPage[PriceReportOut]:
items, next_cursor, total = queries.list_price_reports( items, next_cursor = queries.list_price_reports(
db, status=status, user_id=user_id, limit=limit, cursor=cursor, db, status=status, user_id=user_id, limit=limit, cursor=cursor,
) )
return CursorPage( return CursorPage(
items=[PriceReportOut.model_validate(r) for r in items], items=[PriceReportOut.model_validate(r) for r in items], next_cursor=next_cursor,
next_cursor=next_cursor,
total=total,
) )
@@ -62,9 +60,7 @@ def approve_price_report(
admin: Annotated[AdminUser, Depends(require_role("operator"))], admin: Annotated[AdminUser, Depends(require_role("operator"))],
db: AdminDb, db: AdminDb,
) -> OkResponse: ) -> OkResponse:
# 行锁(SELECT FOR UPDATE):并发/连点双请求会都读到 pending → 各发一次金币双倍发奖, rep = db.get(PriceReport, report_id)
# 锁住该行串行化,第二个请求拿锁后看到 approved → 走 400。SQLite 下 FOR UPDATE 为 no-op。
rep = db.get(PriceReport, report_id, with_for_update=True)
if rep is None: if rep is None:
raise HTTPException(status_code=404, detail="上报记录不存在") raise HTTPException(status_code=404, detail="上报记录不存在")
if rep.status != "pending": if rep.status != "pending":
@@ -92,7 +88,7 @@ def reject_price_report(
admin: Annotated[AdminUser, Depends(require_role("operator"))], admin: Annotated[AdminUser, Depends(require_role("operator"))],
db: AdminDb, db: AdminDb,
) -> OkResponse: ) -> OkResponse:
rep = db.get(PriceReport, report_id, with_for_update=True) # 行锁,同 approve(防并发重复审核) rep = db.get(PriceReport, report_id)
if rep is None: if rep is None:
raise HTTPException(status_code=404, detail="上报记录不存在") raise HTTPException(status_code=404, detail="上报记录不存在")
if rep.status != "pending": if rep.status != "pending":
+7 -12
View File
@@ -45,7 +45,7 @@ def list_users(
limit: Annotated[int, Query(ge=1, le=100)] = 20, limit: Annotated[int, Query(ge=1, le=100)] = 20,
cursor: Annotated[int | None, Query()] = None, cursor: Annotated[int | None, Query()] = None,
) -> CursorPage[AdminUserListItem]: ) -> CursorPage[AdminUserListItem]:
items, next_cursor, total = queries.list_users( items, next_cursor = queries.list_users(
db, phone=phone, register_channel=register_channel, status=status, db, phone=phone, register_channel=register_channel, status=status,
nickname=nickname, created_from=created_from, created_to=created_to, nickname=nickname, created_from=created_from, created_to=created_to,
last_login_from=last_login_from, last_login_to=last_login_to, last_login_from=last_login_from, last_login_to=last_login_to,
@@ -54,7 +54,6 @@ def list_users(
return CursorPage( return CursorPage(
items=[AdminUserListItem.model_validate(u) for u in items], items=[AdminUserListItem.model_validate(u) for u in items],
next_cursor=next_cursor, next_cursor=next_cursor,
total=total,
) )
@@ -127,8 +126,7 @@ def grant_user_coins(
if body.mode == "set": if body.mode == "set":
if body.amount < 0: if body.amount < 0:
raise HTTPException(status_code=400, detail="目标金币值不能为负") raise HTTPException(status_code=400, detail="目标金币值不能为负")
# lock=True:锁账户行,防连点/并发各读同一 before 算同一 delta 双写,余额错位 before = wallet_repo.get_or_create_account(db, user_id, commit=False).coin_balance
before = wallet_repo.get_or_create_account(db, user_id, commit=False, lock=True).coin_balance
delta = body.amount - before delta = body.amount - before
if delta == 0: if delta == 0:
raise HTTPException(status_code=400, detail=f"当前金币已为 {body.amount},无需调整") raise HTTPException(status_code=400, detail=f"当前金币已为 {body.amount},无需调整")
@@ -136,9 +134,9 @@ def grant_user_coins(
if body.amount == 0: if body.amount == 0:
raise HTTPException(status_code=400, detail="amount 不能为 0") raise HTTPException(status_code=400, detail="amount 不能为 0")
delta = body.amount delta = body.amount
# 负数扣减时不允许扣成负余额(运营误操作保护);lock=True 防并发扣穿 # 负数扣减时不允许扣成负余额(运营误操作保护)
if delta < 0: if delta < 0:
acc_now = wallet_repo.get_or_create_account(db, user_id, commit=False, lock=True) acc_now = wallet_repo.get_or_create_account(db, user_id, commit=False)
if acc_now.coin_balance + delta < 0: if acc_now.coin_balance + delta < 0:
raise HTTPException( raise HTTPException(
status_code=400, detail=f"扣减后金币为负(当前余额 {acc_now.coin_balance})" status_code=400, detail=f"扣减后金币为负(当前余额 {acc_now.coin_balance})"
@@ -176,10 +174,7 @@ def grant_user_cash(
if body.mode == "set": if body.mode == "set":
if body.amount_cents < 0: if body.amount_cents < 0:
raise HTTPException(status_code=400, detail="目标现金值不能为负") raise HTTPException(status_code=400, detail="目标现金值不能为负")
# lock=True:锁账户行,防连点/并发各读同一 before 算同一 delta 双写,余额错位 before = wallet_repo.get_or_create_account(db, user_id, commit=False).cash_balance_cents
before = wallet_repo.get_or_create_account(
db, user_id, commit=False, lock=True
).cash_balance_cents
delta = body.amount_cents - before delta = body.amount_cents - before
if delta == 0: if delta == 0:
raise HTTPException( raise HTTPException(
@@ -189,9 +184,9 @@ def grant_user_cash(
if body.amount_cents == 0: if body.amount_cents == 0:
raise HTTPException(status_code=400, detail="amount_cents 不能为 0") raise HTTPException(status_code=400, detail="amount_cents 不能为 0")
delta = body.amount_cents delta = body.amount_cents
# 负数扣减时不允许扣成负余额(运营误操作保护);lock=True 防并发扣穿 # 负数扣减时不允许扣成负余额(运营误操作保护)
if delta < 0: if delta < 0:
acc_now = wallet_repo.get_or_create_account(db, user_id, commit=False, lock=True) acc_now = wallet_repo.get_or_create_account(db, user_id, commit=False)
if acc_now.cash_balance_cents + delta < 0: if acc_now.cash_balance_cents + delta < 0:
raise HTTPException( raise HTTPException(
status_code=400, detail=f"扣减后现金为负(当前余额 {acc_now.cash_balance_cents} 分)" status_code=400, detail=f"扣减后现金为负(当前余额 {acc_now.cash_balance_cents} 分)"
+4 -11
View File
@@ -63,7 +63,7 @@ def list_withdraws(
limit: Annotated[int, Query(ge=1, le=100)] = 20, limit: Annotated[int, Query(ge=1, le=100)] = 20,
cursor: Annotated[int | None, Query()] = None, cursor: Annotated[int | None, Query()] = None,
) -> CursorPage[WithdrawOrderOut]: ) -> CursorPage[WithdrawOrderOut]:
items, next_cursor, total = queries.list_all_withdraw_orders( items, next_cursor = queries.list_all_withdraw_orders(
db, db,
user_id=user_id, user_id=user_id,
status=status, status=status,
@@ -78,9 +78,7 @@ def list_withdraws(
cursor=cursor, cursor=cursor,
) )
return CursorPage( return CursorPage(
items=[WithdrawOrderOut.model_validate(o) for o in items], items=[WithdrawOrderOut.model_validate(o) for o in items], next_cursor=next_cursor,
next_cursor=next_cursor,
total=total,
) )
@@ -89,12 +87,7 @@ def withdraws_summary(db: AdminDb) -> WithdrawSummaryOut:
return WithdrawSummaryOut(**queries.withdraw_summary(db)) return WithdrawSummaryOut(**queries.withdraw_summary(db))
@router.get( @router.get("/health-check", response_model=WxpayHealthCheckOut, summary="提现配置健康检查")
"/health-check",
response_model=WxpayHealthCheckOut,
summary="提现配置健康检查",
dependencies=[Depends(require_role("finance"))], # 暴露密钥路径/配置,限财务+super
)
def withdraw_health_check() -> WxpayHealthCheckOut: def withdraw_health_check() -> WxpayHealthCheckOut:
private_path = wxpay._resolve_config_path(settings.WXPAY_MCH_PRIVATE_KEY_PATH) # noqa: SLF001 private_path = wxpay._resolve_config_path(settings.WXPAY_MCH_PRIVATE_KEY_PATH) # noqa: SLF001
public_path = wxpay._resolve_config_path(settings.WXPAY_PUBLIC_KEY_PATH) # noqa: SLF001 public_path = wxpay._resolve_config_path(settings.WXPAY_PUBLIC_KEY_PATH) # noqa: SLF001
@@ -167,7 +160,7 @@ def withdraw_detail(out_bill_no: str, db: AdminDb) -> WithdrawDetailOut:
withdraw_success_cents=overview["withdraw_success_cents"], withdraw_success_cents=overview["withdraw_success_cents"],
) )
recent_withdraws, _, _ = queries.list_all_withdraw_orders( recent_withdraws, _ = queries.list_all_withdraw_orders(
db, user_id=order.user_id, limit=5, cursor=None, db, user_id=order.user_id, limit=5, cursor=None,
) )
recent_cash_transactions, _ = queries.list_all_cash_transactions( recent_cash_transactions, _ = queries.list_all_cash_transactions(
+1 -6
View File
@@ -9,15 +9,10 @@ T = TypeVar("T")
class CursorPage(BaseModel, Generic[T]): class CursorPage(BaseModel, Generic[T]):
"""分页响应:items + 下一页游标(next_cursor=None 表示末页)+ 可选 total """游标分页响应:items + 下一页游标(next_cursor=None 表示末页)。"""
next_cursor:offset 分页时即下一页 offset,加载更多;末页为 None
total:符合筛选条件的总条数,页码分页(antd pagination);不需要总数的接口可不传(None)
"""
items: list[T] items: list[T]
next_cursor: int | None = None next_cursor: int | None = None
total: int | None = None
class OkResponse(BaseModel): class OkResponse(BaseModel):
+1 -9
View File
@@ -3,7 +3,7 @@ from __future__ import annotations
from datetime import datetime from datetime import datetime
from pydantic import BaseModel, ConfigDict, Field, field_validator from pydantic import BaseModel, ConfigDict, Field
class PriceReportOut(BaseModel): class PriceReportOut(BaseModel):
@@ -36,14 +36,6 @@ class PriceReportOut(BaseModel):
class PriceReportRejectRequest(BaseModel): class PriceReportRejectRequest(BaseModel):
reason: str = Field(min_length=1, max_length=256, description="拒绝理由,用户端记录页会看到") reason: str = Field(min_length=1, max_length=256, description="拒绝理由,用户端记录页会看到")
@field_validator("reason")
@classmethod
def _reason_not_blank(cls, v: str) -> str:
# min_length=1 放过纯空白(" "),trim 后再校验非空,避免审计/用户端记录到空理由
if not v.strip():
raise ValueError("拒绝理由不能为空")
return v.strip()
class PriceReportSummary(BaseModel): class PriceReportSummary(BaseModel):
"""审核台顶部各状态计数。""" """审核台顶部各状态计数。"""
+1 -12
View File
@@ -4,7 +4,7 @@ from __future__ import annotations
from datetime import datetime from datetime import datetime
from typing import Literal from typing import Literal
from pydantic import BaseModel, ConfigDict, Field, field_validator from pydantic import BaseModel, ConfigDict, Field
class AdminUserListItem(BaseModel): class AdminUserListItem(BaseModel):
@@ -37,13 +37,6 @@ class AdminUserOverview(BaseModel):
feedback_total: int feedback_total: int
def _strip_reason(v: str) -> str:
# min_length=1 放过纯空白(" "),trim 后再校验非空,避免审计记到空原因
if not v.strip():
raise ValueError("操作原因不能为空")
return v.strip()
class GrantCoinsRequest(BaseModel): class GrantCoinsRequest(BaseModel):
mode: Literal["delta", "set"] = Field( mode: Literal["delta", "set"] = Field(
"delta", description="delta=增减(amount 为变动量) / set=设为(amount 为目标值,须≥0)" "delta", description="delta=增减(amount 为变动量) / set=设为(amount 为目标值,须≥0)"
@@ -54,8 +47,6 @@ class GrantCoinsRequest(BaseModel):
) )
reason: str = Field(..., min_length=1, max_length=128, description="操作原因(必填,入审计)") reason: str = Field(..., min_length=1, max_length=128, description="操作原因(必填,入审计)")
_v_reason = field_validator("reason")(_strip_reason)
class GrantCashRequest(BaseModel): class GrantCashRequest(BaseModel):
mode: Literal["delta", "set"] = Field( mode: Literal["delta", "set"] = Field(
@@ -67,8 +58,6 @@ class GrantCashRequest(BaseModel):
) )
reason: str = Field(..., min_length=1, max_length=128, description="操作原因(必填,入审计)") reason: str = Field(..., min_length=1, max_length=128, description="操作原因(必填,入审计)")
_v_reason = field_validator("reason")(_strip_reason)
class SetUserStatusRequest(BaseModel): class SetUserStatusRequest(BaseModel):
status: Literal["active", "disabled"] = Field( status: Literal["active", "disabled"] = Field(
-79
View File
@@ -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)
+5 -4
View File
@@ -1,7 +1,7 @@
"""帮助与反馈 endpoint。 """帮助与反馈 endpoint。
路由前缀 `/api/v1/feedback`, Bearer 鉴权(反馈绑到登录用户,便于回访) 路由前缀 `/api/v1/feedback`, Bearer 鉴权(反馈绑到登录用户,便于回访)
POST / 提交反馈(multipart:content 必填;contact 可选(原型改版后客户端已不再采集);images 可选 6 ) POST / 提交反馈(multipart:content / contact 必填,images 可选 4 )
截图复用 [app.core.media] 落盘到 /media/feedback/ 截图复用 [app.core.media] 落盘到 /media/feedback/
""" """
@@ -20,7 +20,7 @@ logger = logging.getLogger("shagua.feedback")
router = APIRouter(prefix="/api/v1/feedback", tags=["feedback"]) router = APIRouter(prefix="/api/v1/feedback", tags=["feedback"])
_MAX_IMAGES = 6 _MAX_IMAGES = 4
_CONTENT_MAX = 2000 _CONTENT_MAX = 2000
_CONTACT_MAX = 128 _CONTACT_MAX = 128
@@ -30,8 +30,7 @@ async def submit_feedback(
user: CurrentUser, user: CurrentUser,
db: DbSession, db: DbSession,
content: str = Form(...), content: str = Form(...),
# 原型改版后客户端不再采集联系方式;保留字段以兼容旧端 + 后续可能复用,默认空串。 contact: str = Form(...),
contact: str = Form(default=""),
images: list[UploadFile] = File(default=[]), images: list[UploadFile] = File(default=[]),
) -> FeedbackOut: ) -> FeedbackOut:
content = content.strip() content = content.strip()
@@ -40,6 +39,8 @@ async def submit_feedback(
raise HTTPException(status_code=400, detail="反馈内容不能为空") raise HTTPException(status_code=400, detail="反馈内容不能为空")
if len(content) > _CONTENT_MAX: if len(content) > _CONTENT_MAX:
raise HTTPException(status_code=400, detail="反馈内容过长") raise HTTPException(status_code=400, detail="反馈内容过长")
if not contact:
raise HTTPException(status_code=400, detail="联系方式不能为空")
if len(contact) > _CONTACT_MAX: if len(contact) > _CONTACT_MAX:
raise HTTPException(status_code=400, detail="联系方式过长") raise HTTPException(status_code=400, detail="联系方式过长")
+1 -29
View File
@@ -9,16 +9,11 @@ from functools import lru_cache
from pathlib import Path from pathlib import Path
from typing import Literal from typing import Literal
from pydantic import Field, model_validator from pydantic import Field
from pydantic_settings import BaseSettings, SettingsConfigDict from pydantic_settings import BaseSettings, SettingsConfigDict
_PROJECT_ROOT = Path(__file__).resolve().parent.parent.parent _PROJECT_ROOT = Path(__file__).resolve().parent.parent.parent
# 生产环境 JWT secret 的最小可接受长度(字节)。HS256 推荐高熵随机串;<16 视为弱密钥。
_MIN_PROD_SECRET_LEN = 16
# 已知的占位默认值(代码里写死的 default),prod 下绝不能沿用。
_INSECURE_SECRET_DEFAULTS = frozenset({"change-me", "change-me-admin", ""})
class Settings(BaseSettings): class Settings(BaseSettings):
model_config = SettingsConfigDict( model_config = SettingsConfigDict(
@@ -203,29 +198,6 @@ class Settings(BaseSettings):
def is_prod(self) -> bool: def is_prod(self) -> bool:
return self.APP_ENV == "prod" return self.APP_ENV == "prod"
@model_validator(mode="after")
def _enforce_prod_secrets(self) -> "Settings":
"""prod 下强校验 JWT secret,弱/默认/空即启动报错(fail-fast,挡住 token 被伪造)。
只校验两个签发凭证:App 用户的 JWT_SECRET_KEY后台的 ADMIN_JWT_SECRET它们沿用默认值
时任何人都能伪造 access/admin token 账号与后台失陷INTERNAL_API_SECRET 默认空 = 内部端点
关闭( 503),是安全的默认态,故不在此强制dev 不触发,便于本地直接起
"""
if not self.is_prod:
return self
weak: list[str] = []
for name in ("JWT_SECRET_KEY", "ADMIN_JWT_SECRET"):
value = getattr(self, name)
if value in _INSECURE_SECRET_DEFAULTS or len(value) < _MIN_PROD_SECRET_LEN:
weak.append(name)
if weak:
raise ValueError(
f"APP_ENV=prod 但检测到弱/默认密钥: {', '.join(weak)} —— 必须改成 "
f"{_MIN_PROD_SECRET_LEN} 位高熵随机串(否则 JWT 可被伪造 → 用户/后台账号失陷)。"
f"生成示例: python -c \"import secrets; print(secrets.token_urlsafe(48))\""
)
return self
@lru_cache(maxsize=1) @lru_cache(maxsize=1)
def get_settings() -> Settings: def get_settings() -> Settings:
+1 -12
View File
@@ -143,13 +143,6 @@ AD_LT_FACTOR_TABLE: tuple[tuple[float, int, int | None], ...] = (
(1.0, 11, None), (1.0, 11, None),
) )
# 客户端可影响的 eCPM 可信上限(分/千次展示):信息流广告一期由客户端上报 eCPM,伪造天价 eCPM
# 可铸出天量金币(见 calculate_ad_reward_coin)。真实 eCPM 一般 <¥100 CPM(=10000 分),档位表顶档
# 为 >¥400(=40000 分);取 ¥500 CPM=50000 分,留足真实头部余量又封死伪造值。钳在唯一计算口
# calculate_ad_reward_coin,故 feed 与 reward_video(回退客户端上报 eCPM 时)一并护住;阈值设在所有
# 真实值之上,不会少发正规奖励。
AD_ECPM_MAX_FEN: int = 50_000
def parse_ecpm_fen(ecpm: str | int | float | None) -> float: def parse_ecpm_fen(ecpm: str | int | float | None) -> float:
"""解析 eCPM 原始值(穿山甲 getEcpm 原值,单位=分/千次展示)。非法/缺失→0。""" """解析 eCPM 原始值(穿山甲 getEcpm 原值,单位=分/千次展示)。非法/缺失→0。"""
@@ -194,12 +187,8 @@ def calculate_ad_reward_coin(ecpm: str | int | float | None, count_after_this: i
eCPM 是穿山甲 getEcpm 原值,单位/千次展示; ÷100 转成元(因子判档 + 收益换算都用元) eCPM 是穿山甲 getEcpm 原值,单位/千次展示; ÷100 转成元(因子判档 + 收益换算都用元)
单次收益()= eCPM元 ÷ 1000(每千次单次) × 因子1(eCPM 元档) × 因子2(LT); 单次收益()= eCPM元 ÷ 1000(每千次单次) × 因子1(eCPM 元档) × 因子2(LT);
再按 1 =10000 金币取整count_after_this 为账号累计第 N 次看视频(LT 因子用,不按天重置) 再按 1 =10000 金币取整count_after_this 为账号累计第 N 次看视频(LT 因子用,不按天重置)
eCPM 在此先钳到 AD_ECPM_MAX_FEN(¥500 CPM):信息流广告一期 eCPM 由客户端上报,伪造天价值
会铸天量金币;钳在这唯一入口,feed reward_video 回退客户端 eCPM 的路径都护住,且阈值高于
所有真实值,不影响正规发奖
""" """
ecpm_yuan = min(parse_ecpm_yuan(ecpm), AD_ECPM_MAX_FEN / 100.0) ecpm_yuan = parse_ecpm_yuan(ecpm)
yuan = (ecpm_yuan / 1000.0) * ad_ecpm_factor(ecpm_yuan) * ad_lt_factor(count_after_this) yuan = (ecpm_yuan / 1000.0) * ad_ecpm_factor(ecpm_yuan) * ad_lt_factor(count_after_this)
return max(0, round(yuan * COIN_PER_YUAN)) return max(0, round(yuan * COIN_PER_YUAN))
+1 -3
View File
@@ -21,7 +21,6 @@ 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.compare_record import router as compare_record_router
from app.api.v1.coupon import router as coupon_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.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.feedback import router as feedback_router
from app.api.v1.invite import router as invite_router from app.api.v1.invite import router as invite_router
from app.api.v1.meituan import router as meituan_router from app.api.v1.meituan import router as meituan_router
@@ -103,9 +102,8 @@ app.include_router(savings_router)
app.include_router(ad_router) app.include_router(ad_router)
app.include_router(order_router) app.include_router(order_router)
app.include_router(report_router) app.include_router(report_router)
# 内部(server→server)端点:pricebot 上报价格观测 / 店铺映射,靠共享密钥头校验,不对客户端开放。 # 内部(server→server)端点:pricebot 上报价格观测,靠共享密钥头校验,不对客户端开放。
app.include_router(internal_price_router) app.include_router(internal_price_router)
app.include_router(internal_store_router)
app.include_router(platform_router) app.include_router(platform_router)
# 用户上传文件(头像)静态服务。生产可改由 nginx 直接 serve MEDIA_ROOT。 # 用户上传文件(头像)静态服务。生产可改由 nginx 直接 serve MEDIA_ROOT。
-1
View File
@@ -23,7 +23,6 @@ from app.models.price_observation import PriceObservation # noqa: F401
from app.models.price_report import PriceReport # noqa: F401 from app.models.price_report import PriceReport # noqa: F401
from app.models.savings import SavingsRecord # noqa: F401 from app.models.savings import SavingsRecord # noqa: F401
from app.models.signin import SigninBoostRecord, SigninRecord # 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.task import UserTask # noqa: F401
from app.models.user import User # noqa: F401 from app.models.user import User # noqa: F401
from app.models.wallet import ( # noqa: F401 from app.models.wallet import ( # noqa: F401
+1 -1
View File
@@ -94,7 +94,7 @@ class ComparisonRecord(Base):
# ===== 明细(JSON,越详细越好)===== # ===== 明细(JSON,越详细越好)=====
# 下单菜品 [{name, qty, specs?}] # 下单菜品 [{name, qty, specs?}]
items: Mapped[list] = mapped_column(_JSON, nullable=False, default=list) 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) comparison_results: Mapped[list] = mapped_column(_JSON, nullable=False, default=list)
# 目标平台未找到、跳过的菜名 # 目标平台未找到、跳过的菜名
skipped_dish_names: Mapped[list] = mapped_column(_JSON, nullable=False, default=list) skipped_dish_names: Mapped[list] = mapped_column(_JSON, nullable=False, default=list)
+2 -3
View File
@@ -1,8 +1,7 @@
"""用户反馈表(帮助与反馈)。 """用户反馈表(帮助与反馈)。
每条 = 用户一次提交content 必填;contact 原为必填(微信/QQ/手机),原型改版后客户端不再采集, 每条 = 用户一次提交content 必填,contact 必填(微信/QQ/手机,便于回访),images 为可选的
新数据存空串(列保持 NOT NULL,免迁移;历史数据仍有值);images 为可选的截图 URL 列表 截图 URL 列表(/media/feedback/...,JSON )status: new(待处理)/ handled(已处理)
(/media/feedback/...,JSON )status: new(待处理)/ handled(已处理)
""" """
from __future__ import annotations from __future__ import annotations
-117
View File
@@ -1,117 +0,0 @@
"""平台店铺表(store_mapping)—— 跨平台"同一家店"的 id/名 映射资产层。
每完成一次淘宝比价(在目标淘宝店通过 更多操作分享复制链接 拿到分享短链
HTTP 解析出 shopId ),pricebot serverserver 内部上报落这里一行**与登录无关
不依赖客户端鉴权**(比价透传链路当前不鉴权,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}>"
)
+2 -12
View File
@@ -16,10 +16,6 @@ from app.repositories import wallet as crud_wallet
FEED_REWARD_UNIT_SECONDS = 10 FEED_REWARD_UNIT_SECONDS = 10
# 单个 feed 事件的时长上限(秒):一期 duration_seconds 由客户端上报,伪造超长时长会刷份数
# (每 10 秒 1 份)。真实单条信息流视频远小于此;取 120s=12 份封顶,挡刷量、不影响正规单。
# 与 rewards.AD_ECPM_MAX_FEN(eCPM 钳顶)合起来,把单事件可铸金币锁进有限区间。
FEED_MAX_DURATION_SECONDS = 120
def _find_by_event(db: Session, client_event_id: str) -> AdFeedRewardRecord | None: def _find_by_event(db: Session, client_event_id: str) -> AdFeedRewardRecord | None:
@@ -70,19 +66,13 @@ def grant_feed_reward(
adn: str | None = None, adn: str | None = None,
slot_id: str | None = None, slot_id: str | None = None,
) -> AdFeedRewardRecord: ) -> AdFeedRewardRecord:
"""完成一条信息流广告后结算奖励。client_event_id 幂等,同号重试不重复发。 """完成一条信息流广告后结算奖励。client_event_id 幂等,同号重试不重复发。"""
一期 eCPM/时长均由客户端上报,故服务端两道硬闸防刷:时长钳到 FEED_MAX_DURATION_SECONDS
限单事件份数,eCPM rewards.calculate_ad_reward_coin 内钳到 AD_ECPM_MAX_FEN 限单份金额;
叠加每日 get_ad_daily_limit 条数上限,把单用户日产出锁进有限区间
"""
existing = _find_by_event(db, client_event_id) existing = _find_by_event(db, client_event_id)
if existing is not None: if existing is not None:
return existing return existing
today = cn_today().isoformat() today = cn_today().isoformat()
# 客户端上报时长先钳到 FEED_MAX_DURATION_SECONDS,防伪造超长时长刷份数(见常量注释)。 safe_duration = max(0, min(duration_seconds, 24 * 60 * 60))
safe_duration = max(0, min(duration_seconds, FEED_MAX_DURATION_SECONDS))
unit_count = safe_duration // FEED_REWARD_UNIT_SECONDS unit_count = safe_duration // FEED_REWARD_UNIT_SECONDS
if _granted_today(db, user_id, today) >= rewards.get_ad_daily_limit(db): if _granted_today(db, user_id, today) >= rewards.get_ad_daily_limit(db):
-172
View File
@@ -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
+3 -9
View File
@@ -78,15 +78,9 @@ class WithdrawNotReviewable(Exception):
"""提现单当前状态不可审核(非 reviewing,可能已被处理过)。""" """提现单当前状态不可审核(非 reviewing,可能已被处理过)。"""
def get_or_create_account( def get_or_create_account(db: Session, user_id: int, *, commit: bool = True) -> CoinAccount:
db: Session, user_id: int, *, commit: bool = True, lock: bool = False """取用户金币账户,不存在则建一个空账户。"""
) -> CoinAccount: acc = db.get(CoinAccount, user_id)
"""取用户金币账户,不存在则建一个空账户。
lock=True 时对已存在的账户行加 SELECT FOR UPDATE(--写余额的调用方串行化,防并发
双写余额错位, admin set 模式连点);默认 False 不改 C 端发奖行为SQLite 下为 no-op
"""
acc = db.get(CoinAccount, user_id, with_for_update=True) if lock else db.get(CoinAccount, user_id)
if acc is None: if acc is None:
acc = CoinAccount( acc = CoinAccount(
user_id=user_id, user_id=user_id,
-11
View File
@@ -24,13 +24,6 @@ class ComparisonItemIn(BaseModel):
specs: list[str] | None = None specs: list[str] | None = None
class AppliedCouponIn(BaseModel):
"""单笔已用优惠(来自 comparison_results[].applied_coupons)。amount 单位:元、正数。"""
name: str
amount: float
class ComparisonResultIn(BaseModel): class ComparisonResultIn(BaseModel):
"""逐平台对比项(来自 done.params.comparison_results)。price 单位:元。""" """逐平台对比项(来自 done.params.comparison_results)。price 单位:元。"""
@@ -47,10 +40,6 @@ class ComparisonResultIn(BaseModel):
# 优惠**来源名**(展示用, best-effort): 美团"外卖大额神券"/京东"百亿补贴"/淘宝"平台红包"。 # 优惠**来源名**(展示用, best-effort): 美团"外卖大额神券"/京东"百亿补贴"/淘宝"平台红包"。
# None=没抠到 → 前端走通用"红包"。同样必须显式声明否则上报边界被 pydantic 静默丢弃(pricebot#38 引入)。 # None=没抠到 → 前端走通用"红包"。同样必须显式声明否则上报边界被 pydantic 静默丢弃(pricebot#38 引入)。
coupon_name: str | None = None coupon_name: str | None = None
# 多券明细 [{name, amount}](全口径: 平台红包+商家券+满减+配送减免, amount 单位元正数)。
# 跟 coupon_saved 并存, 是更丰富的明细; 空=没抠到 → 前端回退单券路径。
# 必须显式声明: 落库走 model_dump(), pydantic 默认丢未知字段, 不声明这行会被悄悄吞掉。
applied_coupons: list[AppliedCouponIn] = Field(default_factory=list)
class ComparisonRecordIn(BaseModel): class ComparisonRecordIn(BaseModel):
-62
View File
@@ -1,62 +0,0 @@
"""平台店铺映射内部上报的收发模型。
pricebot 在淘宝比价拿到 shopId serverserver 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
+5 -11
View File
@@ -1,4 +1,4 @@
# GET /admin/api/feedbacks — 反馈工单列表(offset 分页 + 筛选/排序) # GET /admin/api/feedbacks — 反馈工单列表(游标分页)
> 所属:Admin·反馈 组(前缀 `/admin/api/feedbacks` | 鉴权:Bearer admin_token(角色:任意已登录管理员,无 `require_role`,仅 `get_current_admin` | [← 返回 API 索引](./README.md) > 所属:Admin·反馈 组(前缀 `/admin/api/feedbacks` | 鉴权:Bearer admin_token(角色:任意已登录管理员,无 `require_role`,仅 `get_current_admin` | [← 返回 API 索引](./README.md)
@@ -7,13 +7,8 @@
|---|---|---|---|---| |---|---|---|---|---|
| `status` | string | ❌ | null | 反馈状态,精确匹配:`new`(待处理) / `handled`(已处理);传空/不传则不筛 | | `status` | string | ❌ | null | 反馈状态,精确匹配:`new`(待处理) / `handled`(已处理);传空/不传则不筛 |
| `user_id` | int | ❌ | null | 按提交用户 id 精确筛 | | `user_id` | int | ❌ | null | 按提交用户 id 精确筛 |
| `content` | string | ❌ | null | 反馈内容模糊匹配(ilike,≤100 字) |
| `created_from` | datetime | ❌ | null | 提交时间 ≥(ISO,统一按 UTC 比较) |
| `created_to` | datetime | ❌ | null | 提交时间 ≤(ISO,统一按 UTC 比较) |
| `sort_by` | string | ❌ | `id` | 排序列:`id` / `created_at` |
| `sort_order` | string | ❌ | `desc` | `asc` / `desc` |
| `limit` | int | ❌ | 20 | 1100 | | `limit` | int | ❌ | 20 | 1100 |
| `cursor` | int | ❌ | null | 上一页 next_cursor(**offset 分页**,cursor=offset) | | `cursor` | int | ❌ | null | 上一页 next_cursor(按 feedback id 倒序,查 `id < cursor`) |
## 出参 ## 出参
响应 `200`:`{ items: FeedbackOut[], next_cursor: int|null }`(`next_cursor=null` 表示末页) 响应 `200`:`{ items: FeedbackOut[], next_cursor: int|null }`(`next_cursor=null` 表示末页)
@@ -31,10 +26,9 @@
## 错误码 ## 错误码
- `401` 未带/无效/过期 admin token、管理员被禁用(头带 `WWW-Authenticate: Bearer`) - `401` 未带/无效/过期 admin token、管理员被禁用(头带 `WWW-Authenticate: Bearer`)
- `422` `limit` 超出 1100 范围 / `sort_by`·`sort_order` 不在允许集 / 字段类型不合法 - `422` `limit` 超出 1–100 范围 / 字段类型不合法
## 说明 ## 说明
- **offset 分页**(同 [admin-users-list](./admin-users-list.md)):`cursor` 即 offset,传上一页返回的 `next_cursor`;`next_cursor=null` 即末页。改用 offset 是为了在任意列排序下游标语义统一,代价是翻页期间数据变动可能错位一条(admin 低频可接受)。 - 游标分页约定:结果按 feedback `id` 倒序;`cursor` 传上一页返回的 `next_cursor`;`next_cursor=null` 即末页。
- 排序:`sort_by`(id/created_at)× `sort_order`(asc/desc),恒以 `id` 同向兜底次序 - `status` / `user_id` 均为精确匹配,可叠加
- `status` / `user_id` 精确匹配、`content` 模糊、`created_from`/`created_to` 时间范围,均可叠加。
- 关联表 [feedback](../database/feedback.md);截图为相对路径,经 `GET /media/feedback/<file>` 静态读。 - 关联表 [feedback](../database/feedback.md);截图为相对路径,经 `GET /media/feedback/<file>` 静态读。
+4 -4
View File
@@ -7,8 +7,8 @@
| 字段 | 类型 | 必填 | 说明 | | 字段 | 类型 | 必填 | 说明 |
|---|---|---|---| |---|---|---|---|
| `content` | string | ✓ | 反馈正文,**1-2000 字**(strip 后) | | `content` | string | ✓ | 反馈正文,**1-2000 字**(strip 后) |
| `contact` | string | | 联系方式(微信/QQ/手机号),**128 字**。原型改版后客户端已不再采集、不传该字段(后端默认空串);保留字段兼容旧端 | | `contact` | string | | 联系方式(微信/QQ/手机号),**1-128 字**,便于回访 |
| `images` | file[] | ✗ | 截图,**最多 6 张**,每张走头像同款校验(JPEG/PNG/WebP,≤ 5 MB,魔数嗅探) | | `images` | file[] | ✗ | 截图,**最多 4 张**,每张走头像同款校验(JPEG/PNG/WebP,≤ 5 MB,魔数嗅探) |
## 出参 ## 出参
响应 `200`: 响应 `200`:
@@ -29,9 +29,9 @@
> 不返回上传的 image URL——这是给运营后台看的,客户端通常不需要。 > 不返回上传的 image URL——这是给运营后台看的,客户端通常不需要。
## 错误码 ## 错误码
- `400` 内容为空 / 内容超 2000 字 / 联系方式超 128 字 / 图片超 6 张 / 单图非法(空/过大/格式不对) - `400` 内容为空 / 内容超 2000 字 / 联系方式为空 / 联系方式超 128 字 / 图片超 4 张 / 单图非法(空/过大/格式不对)
- `401` 未带 token / token 无效或过期 / 用户被禁用 - `401` 未带 token / token 无效或过期 / 用户被禁用
- `422``content` 字段 - `422``content``contact` 字段
## 说明 ## 说明
- **反馈绑用户**:`feedback.user_id = current_user.id`,便于回访 - **反馈绑用户**:`feedback.user_id = current_user.id`,便于回访
+1 -19
View File
@@ -2,7 +2,7 @@
> 跨表视角。单表字段级细节看同目录 `<表名>.md`(索引见 [README](./README.md))。 > 跨表视角。单表字段级细节看同目录 `<表名>.md`(索引见 [README](./README.md))。
> 本文专门回答三件「跨表」的事:**① 每块 App 功能用到哪些表 ② 什么操作往哪张表写 ③ 表和表怎么连(join key,含没有外键约束、靠业务字段对齐的语义关联)**。 > 本文专门回答三件「跨表」的事:**① 每块 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 兜底 | | profile「累计省了 / 省钱战绩 / 省钱明细」 | [`savings_record`](./savings_record.md) | 真实下单归因(source=compare)+ 无真实数据时 demo 兜底 |
| 「上报更低价」提交 / 列表 | [`price_report`](./price_report.md) | 众包纠偏:用户举证某平台更便宜,人工审核发奖 | | 「上报更低价」提交 / 列表 | [`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 位置 / 动作 | 表 | 说明 | | App 位置 / 动作 | 表 | 说明 |
|---|---|---| |---|---|---|
@@ -84,13 +76,6 @@
| 首次进 profile 省钱页且无真实记录 | `savings_record`(C `source=demo`) | 懒种子,`ensure_seeded` 按 user 幂等 | | 首次进 profile 省钱页且无真实记录 | `savings_record`(C `source=demo`) | 懒种子,`ensure_seeded` 按 user 幂等 |
| 上报更低价 `POST /report` | `price_report`(C) | 读 `comparison_record.best_price_cents` 校验 | | 上报更低价 `POST /report` | `price_report`(C) | 读 `comparison_record.best_price_cents` 校验 |
| 提交反馈 `POST /feedback` | `feedback`(C) | | | 提交反馈 `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`) ### admin 端(管理员触发,均额外写一条 `admin_audit_log`)
| 后台操作 | 写入 | 操作 | | 后台操作 | 写入 | 操作 |
@@ -137,7 +122,6 @@
- **`comparison_record.store_name``savings_record.shop_name`**:无 id 关联,按**店名字符串相等**给比价记录打「已下单」标记(瞬态,不写库)。两边店名同源 = 比价意图识别阶段的门店 query,语义=**店级**(同店比价多次会一并标已下单)。 - **`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` 仍是旧版兼容统计,不逐条参与发奖。 - **广告流会话关联**:`ad_reward_record.ad_session_id` 可与 `ad_ecpm_record.ad_session_id` 对齐;`ad_watch_log` 仍是旧版兼容统计,不逐条参与发奖。
- **里程碑解锁进度不存库**:`comparison_milestone_claim` 只记「哪几档已领」;进度 = `comparison_record``status='success'``count` - **里程碑解锁进度不存库**:`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`)。登录读、走完引导写,决定是否再展示新手引导。 - **`onboarding_completion.(user_id, device_id)`**:`user_id` 语义关联 `user.id`(无硬 FK,同 `coupon_*` 设备表),`device_id` = 客户端硬件级 `ANDROID_ID`(≠ 领券 per-install `device_id`)。登录读、走完引导写,决定是否再展示新手引导。
### ER 关系(文字版) ### 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, 可空) comparison_record ─1:N─ price_report (comparison_record_id, 可空)
admin_user ─1:N─ admin_audit_log admin_user ─1:N─ admin_audit_log
app_config (独立, 无外键, key 为主键) app_config (独立, 无外键, key 为主键)
coupon_prompt_engagement / coupon_daily_completion / coupon_claim_record
(独立, 无硬 FK; 维度=device_id+北京日, user_id/trace_id 仅软关联)
``` ```
--- ---
+5 -9
View File
@@ -3,13 +3,16 @@
> 数据库:SQLite 起步(`data/app.db`),生产可切 PostgreSQL(改 `DATABASE_URL`)。 > 数据库:SQLite 起步(`data/app.db`),生产可切 PostgreSQL(改 `DATABASE_URL`)。
> ORM:SQLAlchemy 2.0(`app/models/`),迁移:Alembic(`alembic/versions/`,`render_as_batch` 兼容 SQLite)。 > ORM:SQLAlchemy 2.0(`app/models/`),迁移:Alembic(`alembic/versions/`,`render_as_batch` 兼容 SQLite)。
> 金额字段一律存**整数**:金币=个数,现金=**分**(`*_cents`)。时间列 `DateTime(timezone=True)` > 金额字段一律存**整数**:金币=个数,现金=**分**(`*_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」都在那;本页只做**单表索引**,点进每张表的详情看字段级说明。 > 🧭 **先看 [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) | | `savings_record` | 省钱记录(profile 省钱战绩源;真实下单归因 + demo) | `models/savings.py` | [详情](./savings_record.md) |
| `price_report` | 上报更低价(众包纠偏,人工审核发奖) | `models/price_report.py` | [详情](./price_report.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 券缓存 ### 美团 CPS 券缓存
| 表 | 用途 | 模型 | 文档 | | 表 | 用途 | 模型 | 文档 |
|---|---|---|---| |---|---|---|---|
-123
View File
@@ -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` | JSONPG 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))。
+3 -3
View File
@@ -2,10 +2,10 @@
> 模型 `app/models/feedback.py` · 仓库 `app/repositories/feedback.py` · 接口 [feedback](../api/feedback.md) · admin [admin-feedbacks-list](../api/admin-feedbacks-list.md) / [admin-feedback-handle](../api/admin-feedback-handle.md) · [← 索引](./README.md) · [总览](./OVERVIEW.md) > 模型 `app/models/feedback.py` · 仓库 `app/repositories/feedback.py` · 接口 [feedback](../api/feedback.md) · admin [admin-feedbacks-list](../api/admin-feedbacks-list.md) / [admin-feedback-handle](../api/admin-feedback-handle.md) · [← 索引](./README.md) · [总览](./OVERVIEW.md)
App「帮助与反馈」每次提交写一行。`content` 必填;`contact` 必填,**原型改版后客户端不再采集,新数据存空串**(列保持 NOT NULL、免迁移,历史数据仍有值);`images` 为可选截图(≤6 张)。后台人工处理后置 `handled`。与 `price_report`(结构化上报更低价)不同,本表是**自由文本**反馈。 App「帮助与反馈」每次提交写一行。`content` `contact` 必填,`images` 为可选截图。后台人工处理后置 `handled`。与 `price_report`(结构化上报更低价)不同,本表是**自由文本**反馈。
## 用在哪 / 增删改查 ## 用在哪 / 增删改查
- **C(插入)**:`POST /api/v1/feedback`(multipart:`content` + 可选 `contact` + 可选 `images`;`create_feedback`)。截图先经 `core.media``/media/feedback/` 拿相对路径,再随反馈写入,`status='new'` - **C(插入)**:`POST /api/v1/feedback`(multipart:`content` + `contact` + 可选 `images`;`create_feedback`)。截图先经 `core.media``/media/feedback/` 拿相对路径,再随反馈写入,`status='new'`
- **U(更新)**:admin 处理反馈 `update_feedback_status``status='handled'`(同事务写 `admin_audit_log`)。 - **U(更新)**:admin 处理反馈 `update_feedback_status``status='handled'`(同事务写 `admin_audit_log`)。
- **D**:无。 - **D**:无。
- **R**:admin 反馈列表(可按 `status` 筛)。C 端当前无"我的反馈列表"读接口。 - **R**:admin 反馈列表(可按 `status` 筛)。C 端当前无"我的反馈列表"读接口。
@@ -16,7 +16,7 @@ App「帮助与反馈」每次提交写一行。`content` 必填;`contact` 原
| `id` | Integer | PK, autoincrement | | | `id` | Integer | PK, autoincrement | |
| `user_id` | Integer | FK→user.id, index, NOT NULL | 提交用户 | | `user_id` | Integer | FK→user.id, index, NOT NULL | 提交用户 |
| `content` | Text | NOT NULL | 反馈正文 | | `content` | Text | NOT NULL | 反馈正文 |
| `contact` | String(128) | NOT NULL | 联系方式(微信/QQ/手机)。客户端改版后不再采集,新数据为空串;列仍 NOT NULL | | `contact` | String(128) | NOT NULL | 联系方式(微信/QQ/手机,便于回访) |
| `images` | JSON | nullable | 截图相对 URL 列表 `/media/feedback/...`;无图为 NULL | | `images` | JSON | nullable | 截图相对 URL 列表 `/media/feedback/...`;无图为 NULL |
| `status` | String(16) | NOT NULL, default `new` | 取值:`new`(待处理)/ `handled`(已处理) | | `status` | String(16) | NOT NULL, default `new` | 取值:`new`(待处理)/ `handled`(已处理) |
| `created_at` | DateTime(tz) | server_default now(), index | 提交时间 | | `created_at` | DateTime(tz) | server_default now(), index | 提交时间 |
+3 -4
View File
@@ -125,7 +125,7 @@ def test_long_password_does_not_crash(admin_client: TestClient) -> None:
def test_audit_log_pagination_no_gap() -> None: def test_audit_log_pagination_no_gap() -> None:
"""审计分页跨页不丢/不重(offset 分页:cursor offset,翻完覆盖全部)。""" """审计游标分页跨页不丢/不重(回归 next_cursor off-by-one)。"""
from app.admin.repositories import admin_user as admin_repo from app.admin.repositories import admin_user as admin_repo
from app.admin.repositories import audit_log as audit_repo from app.admin.repositories import audit_log as audit_repo
@@ -142,13 +142,12 @@ def test_audit_log_pagination_no_gap() -> None:
) )
created_ids.append(log.id) created_ids.append(log.id)
# limit=2 翻 5 条,收集所有 id,应正好覆盖创建的 5 条(无丢无重);total 恒为符合条件总数 # limit=2 翻 5 条,收集所有 id,应正好覆盖创建的 5 条(无丢无重)
seen: list[int] = [] seen: list[int] = []
cursor = None cursor = None
for _ in range(10): # 上限防死循环 for _ in range(10): # 上限防死循环
rows, cursor, total = audit_repo.list_audit_logs(db, action=action, limit=2, cursor=cursor) rows, cursor = audit_repo.list_audit_logs(db, action=action, limit=2, cursor=cursor)
seen.extend(r.id for r in rows) seen.extend(r.id for r in rows)
assert total == len(created_ids), f"total 应为 {len(created_ids)},得 {total}"
if cursor is None: if cursor is None:
break break
assert sorted(seen) == sorted(created_ids), f"分页丢/重: want={created_ids} got={seen}" assert sorted(seen) == sorted(created_ids), f"分页丢/重: want={created_ids} got={seen}"