Compare commits

..

6 Commits

Author SHA1 Message Date
guke 2461004aa5 Merge branch 'main' into feature/direct-vendor-push-main-20260702 2026-07-22 15:42:15 +08:00
左辰勇 0a8cab9e7e 合并 main 并解决 API 文档冲突 2026-07-20 23:58:31 +08:00
左辰勇 dd15c5dc97 'fix/实现五大厂商推送,与通知中心接口功能' 2026-07-17 11:03:17 +08:00
左辰勇 d6016c12f9 OPPO 推送补新消息分类(channel_id/category/notify_level)
2024-11-20 后创建的 OPPO 应用发送通知必须携带 category(新规),
配置项 OPPO_PUSH_CHANNEL_ID / OPPO_PUSH_CATEGORY / OPPO_PUSH_NOTIFY_LEVEL(0=不传),
均留空时 payload 与原先完全一致。凭据已到位:荣耀/小米(含 channel 154219)/OPPO/vivo
四家配齐(.env,不入库),仅剩华为待补。

Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
2026-07-14 21:22:25 +08:00
左辰勇 a8ac5dc0c7 消息通知中心 3 接口(mock 阶段) + 五厂商推送(补华为/通用发送/测试三件套)
- GET /api/v1/notifications(分页,时间倒序不分组)/unread-count /read(ids|all),
  对外 camelCase(PRD 前端契约);数据为内存 mock(13 类型全覆盖,重启复位),
  接真实数据只换 repositories/notification_mock.py
- vendor_push 补华为 Push Kit(OAuth+messages:send),抽通用 send_notification
  (任意文案+extras+mock 模式),send_accessibility_disabled 改薄封装行为不变
- 新增 /api/v1/push/{vendors,templates,test}:凭据状态/13 类 PRD 文案模板/测试发送
  (默认 mock,mock=false 真发,可联动插站内 mock 消息闭环验证已读)
- .env.example 补 HUAWEI_* 键;docs/api 新增 notifications.md + push-vendor-test.md
- alembic merge 1a924c274fce 收拢 direct_vendor_push_fields/feedback_type_reply
  双 head,恢复 upgrade head 可用
- tests: test_notifications + test_push_center 共 35 例

Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
2026-07-14 18:57:04 +08:00
lowmaster-chen 2236a8b3ee 基于 main 接入各厂商直推服务端
改动:新增厂商推送配置、设备 push_vendor/push_token 字段、device push-test 接口、心跳超时厂商直推发送逻辑和对应测试。

验证:python -m pytest tests/test_device_push.py tests/test_auth.py tests/test_health.py 通过。
2026-07-02 23:01:20 +08:00
25 changed files with 78 additions and 869 deletions
@@ -1,52 +0,0 @@
"""add composite index (user_id, created_at, id) on comparison_record
C 端「我的比价记录」列表(GET /api/v1/compare/records)是
`WHERE user_id=? ORDER BY created_at DESC, id DESC LIMIT n` —— 原来只有单列 user_id 索引,
过滤完还要把该用户的**全部**记录取出来排序才能拿前 n 条,重度用户随记录数线性变慢。
本复合索引的反向扫恰好等于 (created_at DESC, id DESC),规划器直接取前 n 条、免排序。
列序 (user_id, created_at, id) 与查询一一对应,不要调整。
Revision ID: comparison_user_created_idx
Revises: merge_active_phone
Create Date: 2026-07-21
"""
from __future__ import annotations
from alembic import op
revision = "comparison_user_created_idx"
down_revision = "merge_active_phone"
branch_labels = None
depends_on = None
INDEX_NAME = "ix_comparison_user_created"
COLUMNS = ["user_id", "created_at", "id"]
def upgrade() -> None:
bind = op.get_bind()
if bind.dialect.name == "postgresql":
# 线上 comparison_record 已有数据量,普通 CREATE INDEX 持表写锁会阻塞比价 harvest 写入;
# 用 CONCURRENTLY 不锁表(须脱离事务,autocommit_block 切到自动提交)。
# 同 comparison_status_created_idx 的做法。
with op.get_context().autocommit_block():
op.create_index(
INDEX_NAME, "comparison_record", COLUMNS,
unique=False, postgresql_concurrently=True,
)
else:
op.create_index(INDEX_NAME, "comparison_record", COLUMNS, unique=False)
def downgrade() -> None:
bind = op.get_bind()
if bind.dialect.name == "postgresql":
with op.get_context().autocommit_block():
op.drop_index(
INDEX_NAME, table_name="comparison_record",
postgresql_concurrently=True,
)
else:
op.drop_index(INDEX_NAME, table_name="comparison_record")
-56
View File
@@ -1,56 +0,0 @@
"""补齐监控审计页面权限。
Revision ID: monitoring_audit_rbac
Revises: merge_signin_boost_main
Create Date: 2026-07-22 00:00:00.000000
"""
from collections.abc import Sequence
import sqlalchemy as sa
from sqlalchemy.dialects import postgresql
from alembic import op
revision: str = "monitoring_audit_rbac"
down_revision: str | Sequence[str] | None = "merge_signin_boost_main"
branch_labels: str | Sequence[str] | None = None
depends_on: str | Sequence[str] | None = None
_JSON = sa.JSON().with_variant(postgresql.JSONB(), "postgresql")
_PAGE = "analytics-health"
def _role_table() -> sa.TableClause:
return sa.table(
"admin_role",
sa.column("name", sa.String),
sa.column("pages", _JSON),
)
def upgrade() -> None:
role = _role_table()
conn = op.get_bind()
pages = conn.execute(
sa.select(role.c.pages).where(role.c.name == "tech")
).scalar_one_or_none()
if pages is not None and _PAGE not in pages:
conn.execute(
role.update()
.where(role.c.name == "tech")
.values(pages=[*pages, _PAGE])
)
def downgrade() -> None:
role = _role_table()
conn = op.get_bind()
pages = conn.execute(
sa.select(role.c.pages).where(role.c.name == "tech")
).scalar_one_or_none()
if pages is not None and _PAGE in pages:
conn.execute(
role.update()
.where(role.c.name == "tech")
.values(pages=[page for page in pages if page != _PAGE])
)
-29
View File
@@ -10,8 +10,6 @@ from fastapi import Depends, HTTPException, Request, status
from fastapi.security import HTTPAuthorizationCredentials, HTTPBearer
from sqlalchemy.orm import Session
from app.admin.permissions import ALL_PAGE_KEYS, CUSTOM_ROLE, SUPER_ADMIN_ROLE, sanitize_pages
from app.admin.repositories import admin_role as role_repo
from app.admin.repositories import admin_user as admin_repo
from app.admin.security import AdminTokenError, decode_admin_token
from app.db.session import get_db
@@ -74,33 +72,6 @@ def require_role(*roles: str):
return _checker
def require_page(page: str):
"""页面权限守卫依赖工厂。
左侧导航隐藏只是 UI,这个守卫确保直接调用 API 也必须持有对应页面权限。
super_admin 恒通过;custom 读个人 pages_override;其余角色读 admin_role.pages。
"""
if page not in ALL_PAGE_KEYS:
raise ValueError(f"unknown admin page permission: {page}")
def _checker(admin: CurrentAdmin, db: AdminDb) -> AdminUser:
if admin.role == SUPER_ADMIN_ROLE:
return admin
pages = (
sanitize_pages(admin.pages_override)
if admin.role == CUSTOM_ROLE
else role_repo.effective_pages_of(db, admin.role)
)
if page not in pages:
raise HTTPException(
status_code=status.HTTP_403_FORBIDDEN,
detail=f"page '{page}' not allowed",
)
return admin
return _checker
def get_client_ip(request: Request) -> str:
"""取客户端 IP(审计日志用)。生产经 nginx 反代,优先 X-Forwarded-For 第一段;否则直连 IP。
+4 -7
View File
@@ -21,6 +21,7 @@ PERMISSION_CATALOG: list[dict] = [
{"key": "ad-revenue-report", "label": "广告收益"},
{"key": "comparison-records", "label": "比价记录"},
{"key": "cps", "label": "CPS收益"},
{"key": "device-liveness", "label": "设备存活"},
]},
{"group": "奖励审核", "pages": [
{"key": "withdraws", "label": "提现审核"},
@@ -33,14 +34,10 @@ PERMISSION_CATALOG: list[dict] = [
{"key": "huawei-review", "label": "华为审核开关"},
{"key": "users", "label": "用户管理"},
]},
{"group": "监控审计", "pages": [
{"key": "device-liveness", "label": "设备存活"},
{"key": "analytics-health", "label": "埋点成功率"},
{"key": "event-logs", "label": "埋点日志"},
{"key": "audit-logs", "label": "审计日志"},
]},
{"group": "其他", "pages": [
{"key": "admins", "label": "权限管理"},
{"key": "event-logs", "label": "埋点日志"},
{"key": "audit-logs", "label": "审计日志"},
]},
]
@@ -61,7 +58,7 @@ BUILTIN_ROLES: list[dict] = [
"dashboard", "ad-revenue-report", "cps", "withdraws",
]},
{"name": "tech", "label": "技术", "pages": [
"dashboard", "device-liveness", "analytics-health", "config", "ad-revenue", "huawei-review",
"dashboard", "device-liveness", "config", "ad-revenue", "huawei-review",
"event-logs", "audit-logs",
]},
]
-9
View File
@@ -81,10 +81,6 @@ def _date_range(date_from: str, date_to: str) -> list[str]:
# ad_feed_reward_record,由 audit 内部按 ad_type 区分(feed 含历史 NULL,draw 仅 ad_type=="draw")。
_AUDIT_SCENES = {"reward_video", "feed", "draw"}
# 激励视频未满足有效播放条件时不计客户端预估收益。客户端仍会在 onAdShow
# 上报 eCPM,随后才在关闭时补报以下终态,因此必须在展示/发奖合并后修正收益。
_ZERO_REVENUE_REWARD_VIDEO_STATUSES = frozenset({"closed_early", "too_short"})
# 发奖复算明细字段(展开下钻看「金币怎么算出来的」)——从 audit 行原样取这些 key。
_REWARD_DETAIL_KEYS = (
@@ -206,11 +202,6 @@ def ad_revenue_report(
"matched": bool(rwd["matched"]),
"reward_detail": _reward_detail(rwd),
})
if (
rec.ad_type == "reward_video"
and rwd["status"] in _ZERO_REVENUE_REWARD_VIDEO_STATUSES
):
ev["revenue_yuan"] = 0.0
else:
# 纯展示(信息流逐条展示、激励视频缺发奖记录):不计对账,matched=True。
ev.update({
+18 -79
View File
@@ -1,7 +1,7 @@
"""admin「领券数据」看板聚合:发起/完成数、耗时均值与分位、按天/小时趋势、逐条明细。
数据源 coupon_session(一次领券一行,客户端 /api/v1/coupon/session 两段上报)。生产 PostgreSQL
使用 percentile_cont 聚合耗时分位SQLite 本地/测试环境回退读取耗时单列计算
数据源 coupon_session(一次领券一行,客户端 /api/v1/coupon/session 两段上报)。量级不大,全量拉
区间数据后 Python 聚合(分位 SQLite 无 percentile,统一 Python 算,PG 上也一致)
- 发起数 = 区间内全部 session(含 started/completed/failed/abandoned),= 流失统计的基数。
- 完成数 / 耗时均值 / 分位 = 仅 status==completed 子集(成功跑完才有可比的"领券耗时")。
- summary/daily/hourly/total 在全量上算,不受分页;items 为排序后当前页。
@@ -45,75 +45,6 @@ def _percentile(sorted_vals: list[int], q: float) -> int | None:
return round(sorted_vals[lo] * (1 - frac) + sorted_vals[hi] * frac)
def _round_duration_ms(value) -> int | None:
"""将数据库聚合结果按既有 Python round 口径转为整数毫秒。"""
if value is None:
return None
return int(round(value))
def _coupon_summary_aggregate_stmt(conditions: list):
"""PostgreSQL 汇总卡聚合语句;计数、均值与四个分位一次返回。"""
completed = CouponSession.status == "completed"
completed_elapsed = completed & CouponSession.elapsed_ms.is_not(None)
return select(
func.count(CouponSession.id),
func.sum(case((completed, 1), else_=0)),
func.avg(CouponSession.elapsed_ms).filter(completed_elapsed),
*(
func.percentile_cont(q)
.within_group(CouponSession.elapsed_ms)
.filter(completed_elapsed)
for q in (0.05, 0.5, 0.95, 0.99)
),
).where(*conditions)
def _coupon_summary_aggregates(db: Session, conditions: list) -> dict:
"""汇总卡基础指标;生产 PG 全部在数据库内完成,SQLite 仅作测试回退。"""
if db.bind is not None and db.bind.dialect.name == "postgresql":
row = db.execute(_coupon_summary_aggregate_stmt(conditions)).one()
return {
"started_count": int(row[0] or 0),
"completed_count": int(row[1] or 0),
"avg_elapsed_ms": _round_duration_ms(row[2]),
"p5_ms": _round_duration_ms(row[3]),
"p50_ms": _round_duration_ms(row[4]),
"p95_ms": _round_duration_ms(row[5]),
"p99_ms": _round_duration_ms(row[6]),
}
counts = db.execute(
select(
func.count(CouponSession.id),
func.sum(case((CouponSession.status == "completed", 1), else_=0)),
).where(*conditions)
).one()
# SQLite 没有 percentile_cont;本地/测试只回退读取耗时单列,不加载完整记录。
completed_elapsed = list(
db.execute(
select(CouponSession.elapsed_ms)
.where(
*conditions,
CouponSession.status == "completed",
CouponSession.elapsed_ms.is_not(None),
)
.order_by(CouponSession.elapsed_ms)
).scalars()
)
return {
"started_count": int(counts[0] or 0),
"completed_count": int(counts[1] or 0),
"avg_elapsed_ms": _round_duration_ms(
sum(completed_elapsed) / len(completed_elapsed)
) if completed_elapsed else None,
"p5_ms": _percentile(completed_elapsed, 5),
"p50_ms": _percentile(completed_elapsed, 50),
"p95_ms": _percentile(completed_elapsed, 95),
"p99_ms": _percentile(completed_elapsed, 99),
}
def _avg(vals: list[int]) -> int | None:
return round(sum(vals) / len(vals)) if vals else None
@@ -284,22 +215,30 @@ def coupon_data_report(
if not user_ids:
return _empty_result()
conditions = [
stmt = select(CouponSession).where(
CouponSession.started_date >= d_from,
CouponSession.started_date <= d_to,
]
)
if app_env is not None:
conditions.append(CouponSession.app_env == app_env)
stmt = stmt.where(CouponSession.app_env == app_env)
if statuses:
conditions.append(CouponSession.status.in_(statuses))
stmt = stmt.where(CouponSession.status.in_(statuses))
if user_ids is not None:
conditions.append(CouponSession.user_id.in_(user_ids))
stmt = select(CouponSession).where(*conditions)
stmt = stmt.where(CouponSession.user_id.in_(user_ids))
rows = list(db.execute(stmt).scalars())
# ── 汇总卡 ──
completed_elapsed = sorted(
r.elapsed_ms for r in rows if r.status == "completed" and r.elapsed_ms is not None
)
summary = {
**_coupon_summary_aggregates(db, conditions),
"started_count": len(rows),
"completed_count": sum(1 for r in rows if r.status == "completed"),
"avg_elapsed_ms": _avg(completed_elapsed),
"p5_ms": _percentile(completed_elapsed, 5),
"p50_ms": _percentile(completed_elapsed, 50),
"p95_ms": _percentile(completed_elapsed, 95),
"p99_ms": _percentile(completed_elapsed, 99),
**_success_rates(rows),
}
@@ -384,7 +323,7 @@ def coupon_data_report(
"summary": summary,
"daily": daily,
"hourly": hourly,
"total": summary["started_count"],
"total": len(rows),
"items": items,
}
+25 -73
View File
@@ -6,7 +6,6 @@
from __future__ import annotations
from datetime import date, datetime, time, timedelta, timezone
from decimal import ROUND_HALF_UP, Decimal
from zoneinfo import ZoneInfo
from sqlalchemy import Select, asc, case, desc, func, or_, select
@@ -298,58 +297,6 @@ def _comparison_percentile(sorted_values: list[int], q: float) -> int | None:
return int(value + 0.5)
def _round_duration_ms(value) -> int | None:
"""将数据库聚合结果按既有口径四舍五入为整数毫秒。"""
if value is None:
return None
return int(Decimal(str(value)).quantize(Decimal("1"), rounding=ROUND_HALF_UP))
def _comparison_duration_aggregate_stmt(conditions: list, status: str, quantiles: tuple[float, ...]):
"""PostgreSQL 耗时聚合语句;每种状态只返回一行。"""
return select(
func.avg(ComparisonRecord.total_ms),
*(
func.percentile_cont(q).within_group(ComparisonRecord.total_ms)
for q in quantiles
),
).where(
*conditions,
ComparisonRecord.status == status,
ComparisonRecord.total_ms.is_not(None),
)
def _comparison_duration_aggregates(
db: Session,
*,
conditions: list,
status: str,
quantiles: tuple[float, ...],
) -> list[int | None]:
"""返回平均值和各分位数;生产 PG 在数据库内聚合,SQLite 仅作测试回退。"""
if db.bind is not None and db.bind.dialect.name == "postgresql":
row = db.execute(
_comparison_duration_aggregate_stmt(conditions, status, quantiles)
).one()
return [_round_duration_ms(value) for value in row]
# SQLite 没有 percentile_cont;本地/测试只回退读取耗时单列,不加载完整记录。
values = list(
db.execute(
select(ComparisonRecord.total_ms)
.where(
*conditions,
ComparisonRecord.status == status,
ComparisonRecord.total_ms.is_not(None),
)
.order_by(ComparisonRecord.total_ms)
).scalars()
)
average = _round_duration_ms(sum(values) / len(values)) if values else None
return [average, *(_comparison_percentile(values, q) for q in quantiles)]
def comparison_records_summary(
db: Session,
*,
@@ -385,18 +332,20 @@ def comparison_records_summary(
success = int(row[2] or 0)
lower_price = int(row[4] or 0)
cancelled = int(row[5] or 0)
success_duration_stats = _comparison_duration_aggregates(
db,
conditions=conditions,
status="success",
quantiles=(0.05, 0.5, 0.95, 0.99),
)
cancelled_duration_stats = _comparison_duration_aggregates(
db,
conditions=conditions,
status="cancelled",
quantiles=(0.05, 0.5, 0.95),
)
success_durations = sorted(db.execute(
select(ComparisonRecord.total_ms).where(
*conditions,
ComparisonRecord.status == "success",
ComparisonRecord.total_ms.is_not(None),
)
).scalars().all())
cancelled_durations = sorted(db.execute(
select(ComparisonRecord.total_ms).where(
*conditions,
ComparisonRecord.status == "cancelled",
ComparisonRecord.total_ms.is_not(None),
)
).scalars().all())
success_rate_denominator = started - cancelled
return {
"started": started,
@@ -405,16 +354,19 @@ def comparison_records_summary(
"success_rate": success / success_rate_denominator if success_rate_denominator else None,
"avg_token_cost": float(row[3]) if row[3] is not None else None,
"lower_price_rate": lower_price / success if success else None,
"avg_duration_ms": success_duration_stats[0],
"p5_duration_ms": success_duration_stats[1],
"p50_duration_ms": success_duration_stats[2],
"p95_duration_ms": success_duration_stats[3],
"p99_duration_ms": success_duration_stats[4],
"avg_duration_ms": (
int(sum(success_durations) / len(success_durations) + 0.5)
if success_durations else None
),
"p5_duration_ms": _comparison_percentile(success_durations, 0.05),
"p50_duration_ms": _comparison_percentile(success_durations, 0.5),
"p95_duration_ms": _comparison_percentile(success_durations, 0.95),
"p99_duration_ms": _comparison_percentile(success_durations, 0.99),
"cancelled": cancelled,
"cancelled_rate": cancelled / started if started else None,
"cancelled_p5_ms": cancelled_duration_stats[1],
"cancelled_p50_ms": cancelled_duration_stats[2],
"cancelled_p95_ms": cancelled_duration_stats[3],
"cancelled_p5_ms": _comparison_percentile(cancelled_durations, 0.05),
"cancelled_p50_ms": _comparison_percentile(cancelled_durations, 0.5),
"cancelled_p95_ms": _comparison_percentile(cancelled_durations, 0.95),
}
+2 -2
View File
@@ -6,7 +6,7 @@ from typing import Annotated
from fastapi import APIRouter, Depends, Query
from app.admin.deps import AdminDb, require_page
from app.admin.deps import AdminDb, get_current_admin
from app.admin.repositories import analytics_health as repo
from app.admin.schemas.analytics_health import (
HealthBreakdownRow,
@@ -17,7 +17,7 @@ from app.admin.schemas.analytics_health import (
router = APIRouter(
prefix="/admin/api/analytics-health",
tags=["admin-analytics-health"],
dependencies=[Depends(require_page("analytics-health"))],
dependencies=[Depends(get_current_admin)],
)
+3 -3
View File
@@ -1,11 +1,11 @@
"""admin 操作审计日志查询(需要 audit-logs 页面权限)。"""
"""admin 操作审计日志查询(所有 admin 可看:谁在何时对什么做了什么)。"""
from __future__ import annotations
from typing import Annotated
from fastapi import APIRouter, Depends, Query
from app.admin.deps import AdminDb, require_page
from app.admin.deps import AdminDb, get_current_admin
from app.admin.repositories import audit_log as audit_repo
from app.admin.schemas.admin import AdminAuditLogOut
from app.admin.schemas.common import CursorPage
@@ -13,7 +13,7 @@ from app.admin.schemas.common import CursorPage
router = APIRouter(
prefix="/admin/api/audit-logs",
tags=["admin-audit"],
dependencies=[Depends(require_page("audit-logs"))],
dependencies=[Depends(get_current_admin)],
)
+3 -3
View File
@@ -2,7 +2,7 @@
数据源 device_liveness (心跳 last_heartbeat_at + liveness_state + kill_alert_pending,
app/models/device.py)在线/掉线掉线时长由 repo HEARTBEAT_TIMEOUT_MINUTES 阈值派生
纯读:无写无审计需要 device-liveness 页面权限
纯读:无写无审计任意登录管理员可看(同大盘/设备管理,无角色门)
"""
from __future__ import annotations
@@ -10,7 +10,7 @@ from typing import Annotated
from fastapi import APIRouter, Depends, Query
from app.admin.deps import AdminDb, require_page
from app.admin.deps import AdminDb, get_current_admin
from app.admin.repositories import queries
from app.admin.schemas.common import CursorPage
from app.admin.schemas.device import DeviceLivenessItem, DeviceLivenessStats
@@ -18,7 +18,7 @@ from app.admin.schemas.device import DeviceLivenessItem, DeviceLivenessStats
router = APIRouter(
prefix="/admin/api/device-liveness",
tags=["admin-device-liveness"],
dependencies=[Depends(require_page("device-liveness"))],
dependencies=[Depends(get_current_admin)],
)
+2 -2
View File
@@ -6,7 +6,7 @@ from typing import Annotated
from fastapi import APIRouter, Depends, Query
from app.admin.deps import AdminDb, require_page
from app.admin.deps import AdminDb, get_current_admin
from app.admin.repositories import queries
from app.admin.schemas.analytics import AnalyticsEventOut
from app.admin.schemas.common import CursorPage
@@ -14,7 +14,7 @@ from app.admin.schemas.common import CursorPage
router = APIRouter(
prefix="/admin/api/event-logs",
tags=["admin-event-logs"],
dependencies=[Depends(require_page("event-logs"))],
dependencies=[Depends(get_current_admin)],
)
+4 -10
View File
@@ -26,10 +26,7 @@ class AdRevenueRecord(BaseModel):
record_id: int
created_at: datetime
status: str = Field(
...,
description="granted / capped / ecpm_missing / closed_early / too_short",
)
status: str = Field(..., description="granted / capped / ecpm_missing")
ecpm: str | None = Field(None, description="本次采用的 eCPM 原始值(分/千次展示)")
ecpm_factor: float | None = Field(None, description="因子1(eCPM 档);非 granted 为空")
units: int = Field(..., description="折算份数:激励视频恒 1;信息流 = 满 10 秒份数")
@@ -47,7 +44,7 @@ class AdRevenueDaily(BaseModel):
date: str = Field(..., description="北京时间 YYYY-MM-DD")
impressions: int = Field(..., description="当天展示条数合计")
revenue_yuan: float = Field(..., description="当天客户端有效预估收益合计(元;eCPM 折算)")
revenue_yuan: float = Field(..., description="当天客户端预估收益合计(元;eCPM 折算)")
pangle_revenue_yuan: float | None = Field(
None, description="当天穿山甲后台预估收益(元;GroMore revenue);非全量视图/无数据为空"
)
@@ -96,10 +93,7 @@ class AdRevenueRow(BaseModel):
has_impression: bool = Field(..., description="是否有广告展示(信息流逐条展示=True,纯发奖行=False)")
impressions: int = Field(..., description="本行展示条数:有展示=1 / 纯发奖=0(供日汇总、趋势图复用)")
ecpm: str | None = Field(None, description="eCPM 原始值(分/千次);展示行取展示值,纯发奖行取发奖采用值")
revenue_yuan: float = Field(
...,
description="本次有效展示预估收益(元)= eCPM元 ÷ 1000;纯发奖、激励视频提前关闭/时长不足=0",
)
revenue_yuan: float = Field(..., description="本次展示预估收益(元)= eCPM元 ÷ 1000;纯发奖行=0")
row_revenue_yuan: float | None = Field(
None,
description="主表逐行展示用的预估收益(元):一次比价/领券聚合行=该次发奖广告 eCPM 折算之和;"
@@ -156,7 +150,7 @@ class AdRevenueReportOut(BaseModel):
total: int = Field(..., description="广告事件总数(全量,不受分页影响;= 当前筛选下的分页总条数)")
truncated: bool = Field(..., description="当前页之后是否还有更多事件(len(events) > offset + limit)")
total_impressions: int = Field(..., description="全量展示条数合计")
total_revenue_yuan: float = Field(..., description="全量客户端有效预估收益合计(元;eCPM 折算)")
total_revenue_yuan: float = Field(..., description="全量客户端预估收益合计(元;eCPM 折算)")
total_pangle_revenue_yuan: float | None = Field(
None,
description="全量穿山甲后台预估收益合计(元;GroMore revenue)。穿山甲无用户/类型/场景维度,"
+1 -10
View File
@@ -115,22 +115,13 @@ def list_records(
db: DbSession,
limit: int = Query(20, ge=1, le=100),
cursor: int | None = Query(None, description="上一页末条 id"),
ordered: bool | None = Query(
None,
description="true=只看「已下单」(店名命中本人真实下单)的记录;不传=全部",
),
keyword: str | None = Query(
None,
max_length=64,
description="按店名 / 菜名模糊搜索,忽略大小写;空白串等同不传",
),
include_trace: bool = Query(
False,
description="客户端开了本机 agent 调试模式时带 true,放行本人记录的 trace_url",
),
) -> ComparisonRecordPage:
items, next_cursor = crud_compare.list_records(
db, user.id, limit=limit, cursor=cursor, ordered=ordered, keyword=keyword
db, user.id, limit=limit, cursor=cursor
)
outs = [ComparisonRecordOut.model_validate(it) for it in items]
# 权限闸:未开 debug_trace_enabled 的用户不下发 trace_url(列表页「复制调试链接」靠它)。
-4
View File
@@ -45,10 +45,6 @@ class ComparisonRecord(Base):
# 首页轮播 / 省钱战绩聚合都按 status='success' 过滤 + created_at 近期排序;
# 复合索引避免随数据量增大退化成全表扫(单列 created_at 索引不含 status)。
Index("ix_comparison_status_created", "status", "created_at"),
# C 端「我的比价记录」列表:WHERE user_id=? ORDER BY created_at DESC, id DESC LIMIT n。
# 单列 user_id 索引只能过滤,排序仍要把该用户全部记录取出来排一遍;这条复合索引的**反向扫**
# 恰好等于 (created_at DESC, id DESC),PG 直接取前 n 条、免排序。列序不能动。
Index("ix_comparison_user_created", "user_id", "created_at", "id"),
)
id: Mapped[int] = mapped_column(Integer, primary_key=True, autoincrement=True)
+8 -83
View File
@@ -7,8 +7,8 @@ from __future__ import annotations
from datetime import datetime
from sqlalchemy import func, or_, select
from sqlalchemy.orm import Session, defer
from sqlalchemy import func, select
from sqlalchemy.orm import Session
from app.core.rewards import CN_TZ
from app.models.ad_feed_reward import AdFeedRewardRecord
@@ -375,48 +375,19 @@ def harvest_abort(
return rec
def _ordered_shop_name_select(user_id: int):
"""该用户「真实下单」(source='compare')覆盖到的店名 select,给「已下单」筛选当子查询
口径与 [_ordered_shop_names] 完全一致,只是时机不同:那边是**拿到本页之后** candidates
反查打标;这边是**分页之前**就要过滤,拿不到 candidates,只能整段下推成子查询
没有先捞成集合再展开 IN (...) 字面量 重度用户下单过的店名可能上千,展开会撞 SQLite
的绑定变量上限,而且又变回了那个随下单量线性变慢的老写法
"""
return select(SavingsRecord.shop_name).where(
SavingsRecord.user_id == user_id,
SavingsRecord.source == "compare",
SavingsRecord.shop_name.is_not(None),
)
def _like_escape(kw: str) -> str:
"""转义 LIKE 通配符(百分号 / 下划线 / 反斜杠),让用户输入只按字面量匹配(配合 escape 参数)。
不转义的话搜一个%就等于把整表拉回来
"""
return kw.replace("\\", "\\\\").replace("%", "\\%").replace("_", "\\_")
def _ordered_shop_names(db: Session, user_id: int, candidates: set[str]) -> set[str]:
"""[candidates] 里哪些店名被该用户「真实下单」(source='compare')覆盖过,用来打「已下单」。
def _ordered_shop_names(db: Session, user_id: int) -> set[str]:
"""该用户「真实下单」(source='compare')覆盖到的店名集合,用来给比价记录打「已下单」
只认 compare(归因命中后真实上报),demo 演示数据不算下单上报不带 trace_id,
只能按店名对齐两边店名同源(都来自比价意图识别阶段的门店名 query),精确相等即视为同店
语义=店级:同一家店比价过多次,这些记录会一并标已下单
只查**本页出现过的店名**(candidates limit ),不再把该用户全部下单店名捞回内存:
老写法随下单量线性增长,重度用户几千行全读一遍只为跟 50 条记录取交集空集合直接返回
(避免 IN () 非法)
"""
if not candidates:
return set()
rows = db.execute(
select(SavingsRecord.shop_name).where(
SavingsRecord.user_id == user_id,
SavingsRecord.source == "compare",
SavingsRecord.shop_name.in_(candidates),
).distinct()
SavingsRecord.shop_name.is_not(None),
)
).scalars().all()
return {s for s in rows if s}
@@ -444,60 +415,17 @@ def _ad_coins_by_trace(db: Session, user_id: int, trace_ids: list[str]) -> dict[
return {tid: int(coin) for tid, coin in rows if tid}
# 列表出参(ComparisonRecordOut)根本不读、但 select(ORM) 默认会一并捞回来的重型 JSON 列:
# - raw_payload:done.params 上报体全量,**每条记录都有**(harvest 与 POST 两条写路径都落)。
# 单条几 KB~几十 KB,一页 50 条就是稳定几百 KB~几 MB 的白读 + 白反序列化。
# - llm_calls:每次 LLM 调用的 input_messages + output 全文。只有走老客户端 POST /compare/record
# 的记录才有(_backfill_llm_calls 回填;harvest 路径不落),但有的时候单条就能到 MB 级 —— 一页里
# 混进几条这种记录,整个请求就被它们拖住。
# - llm_price_snapshot:逐模型单价快照,同样只在回填时落。
# 三列全部读出来再被 pydantic 丢掉,是「比价记录/全部记录」页慢的主要来源。
# ⚠️ defer 的列一旦在别处被读到会触发**逐行**懒加载(N+1);列表这条链路(ComparisonRecordOut
# 不声明这三个字段 → 不会 getattr 到)是安全的。详情接口 get_record 不 defer,raw_payload 照常返回。
_LIST_DEFERRED = (
ComparisonRecord.raw_payload,
ComparisonRecord.llm_calls,
ComparisonRecord.llm_price_snapshot,
)
def list_records(
db: Session,
user_id: int,
*,
limit: int = 20,
cursor: int | None = None,
ordered: bool | None = None,
keyword: str | None = None,
) -> tuple[list[ComparisonRecord], int | None]:
"""比价记录分页(按创建时间倒序、id 兜底,游标式)。附「已下单」店级标记 + 「看广告赚的金币」(瞬态,不写库)。"""
stmt = (
select(ComparisonRecord)
.where(ComparisonRecord.user_id == user_id)
.options(*(defer(col) for col in _LIST_DEFERRED))
)
stmt = select(ComparisonRecord).where(ComparisonRecord.user_id == user_id)
if cursor is not None:
stmt = stmt.where(ComparisonRecord.id < cursor)
# 「已下单」tab 与搜索框的过滤都下推到这里,不能留给客户端对整页结果 filter ——
# 分页之后一页里可能一条都不命中,列表看着就是空的/卡住的,得翻很多页才蹦出一条。
if ordered:
stmt = stmt.where(
ComparisonRecord.store_name.in_(_ordered_shop_name_select(user_id))
)
kw = (keyword or "").strip()
if kw:
# product_names 是写路径从 items[].name 派生的普通文本列(items 本身是 JSON,SQLite 下
# 中文被 ensure_ascii 转义,没法直接 LIKE)—— 搜「菜名」靠的就是它。
# ilike:PG 原生 ILIKE,SQLite 渲染成 lower() LIKE lower(),两边都忽略大小写。
pattern = f"%{_like_escape(kw)}%"
stmt = stmt.where(
or_(
ComparisonRecord.store_name.ilike(pattern, escape="\\"),
ComparisonRecord.product_names.ilike(pattern, escape="\\"),
)
)
# 排序与 ix_comparison_user_created(user_id, created_at, id)对齐 —— DESC/DESC 正好是该索引的
# 反向扫,PG 免排序直接取前 limit 条。改排序方向前先想清楚索引还吃不吃得上。
stmt = stmt.order_by(ComparisonRecord.created_at.desc(), ComparisonRecord.id.desc()).limit(limit)
items = list(db.execute(stmt).scalars().all())
@@ -505,10 +433,7 @@ def list_records(
# 「已下单」标记:本页记录的 store_name 若落在该用户真实下单的店名集合里即 True。
# ordered / ad_coins_earned 均非 ORM 列,仅挂实例上供 ComparisonRecordOut(from_attributes) 读出,不持久化。
page_shops = {it.store_name for it in items if it.store_name}
# ordered=True 时上面已按同一口径(_ordered_shop_name_select)筛过,本页必然全是已下单,
# 省掉这次反查;其余情况照旧按本页店名反查 savings。
ordered_shops = page_shops if ordered else _ordered_shop_names(db, user_id, page_shops)
ordered_shops = _ordered_shop_names(db, user_id)
# 「本次比价看广告赚的金币」:按本页 trace_id 一次性聚合(同 ordered 范式)。
ad_coins = _ad_coins_by_trace(db, user_id, [it.trace_id for it in items])
for it in items:
-12
View File
@@ -25,18 +25,6 @@ server {
# (纯文字反馈体积小、不受影响 → 呈现为「时好时坏」)。根治仍需客户端上传前压缩。
client_max_body_size 32m;
# JSON 响应压缩。nginx 默认 gzip off,且就算 on 了 gzip_types 也只含 text/html、
# gzip_proxied 默认 off(反代来的响应一律不压)—— 三个默认值凑一起 = 我们所有接口都在裸奔。
# 比价记录列表这种一次 50 条、字段名 + 中文店名/菜名高度重复的 JSON,gzip 压缩比稳定在 8~10 倍
# (几百 KB → 几十 KB),弱网下省的就是首屏那几秒。
# 只压 JSON:APK 直链(/media/shaguabijia.apk)、图片本身已是压缩格式,再压纯浪费 CPU。
gzip on;
gzip_proxied any; # 反代响应也压(默认 off = 对我们这套反代等于没开)
gzip_types application/json;
gzip_min_length 1024; # 小响应压了反而更大(gzip 头开销),不值当
gzip_comp_level 5; # 5 是体积/CPU 的常用折中点,再往上收益递减
gzip_vary on; # 给 CDN/中间缓存正确按 Accept-Encoding 分桶
location / {
proxy_pass http://127.0.0.1:8770;
proxy_http_version 1.1;
+2 -3
View File
@@ -1,6 +1,6 @@
# GET /admin/api/audit-logs — 审计日志(谁改了什么,游标分页)
> 所属:Admin·Audit 组(前缀 `/admin/api/audit-logs` | 鉴权:Bearer admin_token + `audit-logs` 页面权限 | [← 返回 API 索引](../README.md)
> 所属:Admin·Audit 组(前缀 `/admin/api/audit-logs` | 鉴权:Bearer admin_token(角色:任意已登录 admin | [← 返回 API 索引](../README.md)
## 入参(query
| 字段 | 类型 | 必填 | 默认 | 说明 |
@@ -29,8 +29,7 @@
## 错误码
- `401` 未带 admin token / token 无效或过期 / 管理员被禁用
- `403` 当前管理员没有 `audit-logs` 页面权限
## 说明
- 整组(`/admin/api/audit-logs`)守卫为 `require_page("audit-logs")`,默认仅超级管理员和技术角色可查看,也可由超管给自定义角色授权
- 整组(`/admin/api/audit-logs`)守卫为 `get_current_admin`,任意已登录 admin 均可查看,无角色限制
- 审计日志只增不改不删,任何写操作经 `write_audit` 落一条。数据表见 [admin_audit_log](../database/admin_audit_log.md)。
+1 -2
View File
@@ -1,6 +1,6 @@
# /admin/api/device-liveness — 设备存活监控(#80)
> 所属:Admin 子应用(前缀 `/admin/api`) | 鉴权:admin + `device-liveness` 页面权限 | 表 [device_liveness](../../database/device_liveness.md) | [← 返回 API 索引](../README.md)
> 所属:Admin 子应用(前缀 `/admin/api`) | 鉴权:admin | 表 [device_liveness](../../database/device_liveness.md) | [← 返回 API 索引](../README.md)
无障碍保护存活的后台视角:哪些设备开过保护(`ever_protected`)、现在在线还是掉线(心跳超时,#107 起阈值 1 小时)、首次开启时间(`first_protected_at`)。
@@ -13,4 +13,3 @@
## 说明
- 「在线」= `last_heartbeat_at` 距今 < 超时阈值;掉线召回链路(worker 置 `kill_alert_pending` → 客户端 pull)见表文档。
- 无 `device-liveness` 页面权限时返回 `403`
+1 -2
View File
@@ -1,6 +1,6 @@
# /admin/api/event-logs — 埋点日志(#83)
> 所属:Admin 子应用(前缀 `/admin/api`) | 鉴权:admin + `event-logs` 页面权限 | 表 [analytics_event](../../database/analytics_event.md) | [← 返回 API 索引](../README.md)
> 所属:Admin 子应用(前缀 `/admin/api`) | 鉴权:admin | 表 [analytics_event](../../database/analytics_event.md) | [← 返回 API 索引](../README.md)
客户端埋点(`POST /api/v1/analytics/events` 批量上报)的后台检索页。
@@ -13,4 +13,3 @@
## 说明
- 纯只读;无聚合报表(要分析导出后自己算)。
- 时间轴用 `client_ts`(事件真实发生时刻),入库时间受客户端攒批影响。
- 无 `event-logs` 页面权限时返回 `403`
-9
View File
@@ -10,12 +10,6 @@
|---|---|---|---|---|
| `limit` | int | ❌ | 20 | 1100 |
| `cursor` | int | ❌ | null | 上一页末条 `id`,首页不传 |
| `ordered` | bool | ❌ | null | `true`=只出「已下单」(店名命中本人真实下单)的记录;不传=不筛 |
| `keyword` | string | ❌ | null | 按店名 / 菜名模糊搜索,忽略大小写,≤64 字符;纯空白等同不传 |
| `include_trace` | bool | ❌ | false | 客户端开了本机 agent 调试模式时带 `true`,放行**本人**记录的 `trace_url` |
`ordered` / `keyword` 都在服务端过滤后再分页,客户端不要拿一页结果自己 filter ——
分页之后一页里可能一条都不命中,列表会看着像空的。
## 出参
响应 `200``{ items: ComparisonRecordOut[], next_cursor: int|null }`(分页见 [索引#游标分页约定](./README.md#游标分页约定)
@@ -44,9 +38,6 @@
| `items` | object[] | 下单菜品 `{name, qty, specs?}` |
| `comparison_results` | object[] | 逐平台对比(price 单位元,已按 rank 升序) |
| `skipped_dish_names` | string[] | 被跳过的菜名 |
| `ordered` | bool | 「已下单」店级标记:店名命中本人 `source='compare'` 的下单记录即 `true`。**瞬态字段,不在表里**,每次查询现算 |
| `ad_coins_earned` | int | 本次比价看信息流广告实发的金币(按 `trace_id` 聚合)。同为瞬态字段 |
| `trace_url` | string \| null | pricebot 调试链接。未开 `debug_trace_enabled` 且未带 `include_trace=true` 时为 `null` |
| `created_at` | datetime | 时间 |
## 错误
-81
View File
@@ -9,11 +9,9 @@ from app.admin.repositories import ad_revenue
from app.db.session import SessionLocal
from app.models.ad_ecpm import AdEcpmRecord
from app.models.ad_pangle_revenue import AdPangleDailyRevenue
from app.models.ad_reward import AdRewardRecord
from app.models.user import User
REPORT_DATE = "2040-02-03"
PLAYBACK_DATE = "2040-02-04"
def test_business_scope_filters_client_and_pangle_by_env_and_code(monkeypatch) -> None:
@@ -129,82 +127,3 @@ def test_business_scope_filters_client_and_pangle_by_env_and_code(monkeypatch) -
db.execute(delete(User).where(User.phone == "18800009991"))
db.commit()
db.close()
def test_reward_video_incomplete_playback_has_zero_revenue() -> None:
db = SessionLocal()
phone = "18800009992"
sessions = {
"closed_early": "rv-zero-closed",
"too_short": "rv-zero-short",
"capped": "rv-keep-capped",
"granted": "rv-keep-granted",
}
try:
user = User(phone=phone, username="29999999992", register_channel="sms")
db.add(user)
db.flush()
for index, (status, session_id) in enumerate(sessions.items(), start=1):
created_at = datetime(2040, 2, 4, index, tzinfo=UTC)
db.add(AdEcpmRecord(
user_id=user.id,
ad_type="reward_video",
ad_session_id=session_id,
app_env="prod",
our_code_id="prod-reward",
ecpm_raw="10000",
report_date=PLAYBACK_DATE,
created_at=created_at,
))
db.add(AdRewardRecord(
trans_id=f"{session_id}-trans",
user_id=user.id,
coin=0,
status=status,
reward_scene="reward_video",
ad_session_id=session_id,
app_env="prod",
our_code_id="prod-reward",
ecpm_raw="10000",
reward_date=PLAYBACK_DATE,
created_at=created_at,
))
db.commit()
result = ad_revenue.ad_revenue_report(
db,
date_from=PLAYBACK_DATE,
date_to=PLAYBACK_DATE,
user_id=user.id,
ad_type="reward_video",
app_env="prod",
revenue_scope="all",
granularity="hour",
)
revenue_by_status = {row["status"]: row["revenue_yuan"] for row in result["items"]}
assert revenue_by_status == {
"closed_early": 0.0,
"too_short": 0.0,
"capped": 0.1,
"granted": 0.1,
}
assert result["total_impressions"] == 4
assert result["total_revenue_yuan"] == 0.2
assert len(result["daily"]) == 1
assert result["daily"][0]["date"] == PLAYBACK_DATE
assert result["daily"][0]["impressions"] == 4
assert result["daily"][0]["revenue_yuan"] == 0.2
assert sum(row["revenue_yuan"] for row in result["hourly"]) == 0.2
assert result["type_stats"]["reward_video"] == {
"impressions": 4,
"revenue_yuan": 0.2,
}
finally:
db.rollback()
db.execute(delete(AdRewardRecord).where(AdRewardRecord.reward_date == PLAYBACK_DATE))
db.execute(delete(AdEcpmRecord).where(AdEcpmRecord.report_date == PLAYBACK_DATE))
db.execute(delete(User).where(User.phone == phone))
db.commit()
db.close()
+2 -63
View File
@@ -44,78 +44,17 @@ def operator_token() -> str:
return _token("r_operator", "operator")
@pytest.fixture()
def tech_token() -> str:
return _token("r_tech", "tech")
def _auth(t: str) -> dict:
return {"Authorization": f"Bearer {t}"}
def test_super_pages_all_operator_limited(admin_client, super_token, operator_token) -> None:
su = admin_client.get("/admin/api/auth/me", headers=_auth(super_token)).json()
assert "admins" in su["pages"] and "analytics-health" in su["pages"] # 超管全页
assert "admins" in su["pages"] and "dashboard" in su["pages"] # 超管全页
op = admin_client.get("/admin/api/auth/me", headers=_auth(operator_token)).json()
assert "dashboard" in op["pages"] and "admins" not in op["pages"] # 运营看不到管理员页
def test_monitoring_audit_catalog_and_api_permissions(
admin_client, super_token, operator_token, tech_token
) -> None:
catalog = admin_client.get(
"/admin/api/roles/catalog", headers=_auth(super_token)
).json()
monitoring = next(group for group in catalog if group["group"] == "监控审计")
assert [page["key"] for page in monitoring["pages"]] == [
"device-liveness", "analytics-health", "event-logs", "audit-logs",
]
# 运营默认只能查设备存活,不能绕过导航直调技术/审计接口。
assert admin_client.get(
"/admin/api/device-liveness/stats", headers=_auth(operator_token)
).status_code == 200
for path in (
"/admin/api/analytics-health/overview?date_from=2026-07-01T00:00:00Z&date_to=2026-07-02T00:00:00Z",
"/admin/api/event-logs",
"/admin/api/audit-logs",
):
assert admin_client.get(path, headers=_auth(operator_token)).status_code == 403
# 技术角色默认拥有监控审计组全部四项权限。
for path in (
"/admin/api/device-liveness/stats",
"/admin/api/analytics-health/overview?date_from=2026-07-01T00:00:00Z&date_to=2026-07-02T00:00:00Z",
"/admin/api/event-logs",
"/admin/api/audit-logs",
):
assert admin_client.get(path, headers=_auth(tech_token)).status_code == 200
def test_custom_admin_api_permission_uses_pages_override(admin_client) -> None:
db = SessionLocal()
try:
admin = admin_repo.get_by_username(db, "r_monitoring_custom")
if admin is None:
admin = admin_repo.create_admin(
db, username="r_monitoring_custom", password="pass1234", role="custom"
)
admin.password_hash = hash_password("pass1234")
admin.role = "custom"
admin.pages_override = ["event-logs"]
admin.status = "active"
db.commit()
finally:
db.close()
token = admin_client.post(
"/admin/api/auth/login",
json={"username": "r_monitoring_custom", "password": "pass1234"},
).json()["access_token"]
assert admin_client.get("/admin/api/event-logs", headers=_auth(token)).status_code == 200
assert admin_client.get("/admin/api/audit-logs", headers=_auth(token)).status_code == 403
def test_roles_endpoints_super_only(admin_client, super_token, operator_token) -> None:
assert admin_client.get("/admin/api/roles", headers=_auth(super_token)).status_code == 200
assert admin_client.get("/admin/api/roles", headers=_auth(operator_token)).status_code == 403
@@ -184,7 +123,7 @@ def test_builtin_roles_labels_and_pages(admin_client, super_token) -> None:
# 页集对齐 Prototypes/dashboard/permissions.md 的 ROLES
assert set(roles["finance"]["pages"]) == {"dashboard", "ad-revenue-report", "cps", "withdraws"}
assert set(roles["tech"]["pages"]) == {
"dashboard", "device-liveness", "analytics-health", "config", "ad-revenue", "huawei-review",
"dashboard", "device-liveness", "config", "ad-revenue", "huawei-review",
"event-logs", "audit-logs",
}
-204
View File
@@ -222,210 +222,6 @@ def test_stats_compare_count_and_saved(client) -> None:
assert s2["compare_count"] == 2 # 仍 2(failed 不计)
def test_records_ordered_flag(client) -> None:
"""「已下单」店级标记:店名命中该用户 source='compare' 的下单记录才 True。
覆盖 list_records 只按**本页店名**反查 savings 的写法(原来是把该用户全部下单店名捞回内存
再取交集,随下单量线性变慢)两种写法结果必须一致,故这里按店名逐条断言
"""
token = _login(client, "13800002010")
# 两条比价记录:一条海底捞(稍后会有对应下单),一条没下过单的店
client.post("/api/v1/compare/record", json=_food_payload("ord-1"), headers=_auth(token))
other = _food_payload("ord-2")
other["store_name"] = "没下过单的店"
client.post("/api/v1/compare/record", json=other, headers=_auth(token))
# 下单前:两条都不该带「已下单」
items = client.get("/api/v1/compare/records", headers=_auth(token)).json()["items"]
assert {it["store_name"]: it["ordered"] for it in items} == {
"海底捞(朝阳店)": False,
"没下过单的店": False,
}
# 对海底捞真实下单一笔(order/report 写 source='compare' 的 savings_record)
r = client.post(
"/api/v1/order/report",
json={
"client_event_id": "evt-ordered-flag",
"platform": "美团",
"platform_package": "com.sankuai.meituan",
"pay_channel": "wechat",
"compared_price_cents": 12350,
"paid_amount_cents": 12350,
"shop_name": "海底捞(朝阳店)",
"original_price_cents": 12850,
},
headers=_auth(token),
)
assert r.status_code == 200, r.text
# 下单后:只有同店名那条翻成 True,另一条不受影响
items = client.get("/api/v1/compare/records", headers=_auth(token)).json()["items"]
assert {it["store_name"]: it["ordered"] for it in items} == {
"海底捞(朝阳店)": True,
"没下过单的店": False,
}
# 别人的下单不该影响本人标记(_ordered_shop_names 按 user_id 过滤)
token_b = _login(client, "13800002011")
client.post("/api/v1/compare/record", json=_food_payload("ord-b"), headers=_auth(token_b))
items_b = client.get("/api/v1/compare/records", headers=_auth(token_b)).json()["items"]
assert [it["ordered"] for it in items_b] == [False]
def test_records_list_omits_raw_payload(client) -> None:
"""列表出参不含 raw_payload(仓库层 defer 掉了重型 JSON 列);详情接口照常返回。
defer 的列一旦被 ORM 实例读到会触发逐行懒加载(N+1),而列表 schema 本就不该带 raw_payload
这条同时守住列表不泄露上报体全量没人不小心把它加回出参
"""
token = _login(client, "13800002012")
rid = client.post(
"/api/v1/compare/record", json=_food_payload("no-raw"), headers=_auth(token)
).json()["id"]
items = client.get("/api/v1/compare/records", headers=_auth(token)).json()["items"]
assert len(items) == 1
assert "raw_payload" not in items[0]
# 概要字段照常齐全(defer 没误伤列表要用的列)
assert items[0]["store_name"] == "海底捞(朝阳店)"
assert items[0]["best_platform_id"] == "meituan"
assert items[0]["comparison_results"] and items[0]["items"]
# 详情不 defer:raw_payload 全量还在
d = client.get(f"/api/v1/compare/records/{rid}", headers=_auth(token)).json()
assert d["raw_payload"]["trace_id"] == "no-raw"
def test_records_ordered_filter(client) -> None:
"""ordered=true 只出「已下单」的记录,且过滤结果自身能翻页。
这个筛选必须在服务端做:客户端早先是对已经拉回来的那一页 filter,分页之后一页里
很可能一条已下单都没有 已下单tab 就会看着像空的,得手动翻很多页才蹦出一条
"""
token = _login(client, "13800002013")
# 3 条「下过单的店」+ 2 条没下过单的店,交错写入,确保过滤不是靠顺序碰巧对上
for i in range(3):
p = _food_payload(f"of-ordered-{i}")
p["store_name"] = "下过单的店"
client.post("/api/v1/compare/record", json=p, headers=_auth(token))
if i < 2:
q = _food_payload(f"of-plain-{i}")
q["store_name"] = "没下过单的店"
client.post("/api/v1/compare/record", json=q, headers=_auth(token))
client.post(
"/api/v1/order/report",
json={
"client_event_id": "evt-ordered-filter",
"platform": "美团",
"platform_package": "com.sankuai.meituan",
"pay_channel": "wechat",
"compared_price_cents": 12350,
"paid_amount_cents": 12350,
"shop_name": "下过单的店",
"original_price_cents": 12850,
},
headers=_auth(token),
)
# 不传 ordered:5 条全出(「全部记录」tab 口径不变)
assert len(client.get("/api/v1/compare/records", headers=_auth(token)).json()["items"]) == 5
# ordered=true:只出那 3 条,且每条都自带 ordered=True
page = client.get("/api/v1/compare/records?ordered=true", headers=_auth(token)).json()
assert [it["store_name"] for it in page["items"]] == ["下过单的店"] * 3
assert all(it["ordered"] for it in page["items"])
assert page["next_cursor"] is None
# 游标只在「已下单」集合内走 —— 不会把没下单的记录算进一页的 limit 里
p1 = client.get(
"/api/v1/compare/records?ordered=true&limit=2", headers=_auth(token)
).json()
assert len(p1["items"]) == 2
assert p1["next_cursor"] is not None
p2 = client.get(
f"/api/v1/compare/records?ordered=true&limit=2&cursor={p1['next_cursor']}",
headers=_auth(token),
).json()
assert [it["store_name"] for it in p2["items"]] == ["下过单的店"]
# 两页不重叠,合起来正好 3 条
assert len({it["id"] for it in p1["items"] + p2["items"]}) == 3
# 别人的下单不该让本人记录进「已下单」
token_b = _login(client, "13800002014")
pb = _food_payload("of-b")
pb["store_name"] = "下过单的店"
client.post("/api/v1/compare/record", json=pb, headers=_auth(token_b))
assert client.get(
"/api/v1/compare/records?ordered=true", headers=_auth(token_b)
).json()["items"] == []
def test_records_keyword_search(client) -> None:
"""keyword 按店名 / 菜名模糊搜(忽略大小写),LIKE 通配符只当字面量;搜索结果也能翻页。
菜名走写路径派生的 product_names 文本列 items JSON,SQLite 下中文被 ensure_ascii
转义,直接 LIKE 搜不到
"""
token = _login(client, "13800002015")
b = _food_payload("kw-b")
b["store_name"] = "Pizza Hut"
b["items"] = [{"name": "榴莲比萨", "qty": 1}]
client.post("/api/v1/compare/record", json=b, headers=_auth(token))
c = _food_payload("kw-c")
c["store_name"] = "100%纯牛肉汉堡"
c["items"] = [{"name": "双层牛肉堡", "qty": 1}]
client.post("/api/v1/compare/record", json=c, headers=_auth(token))
# 默认 payload 的店名是「海底捞(朝阳店)」
client.post("/api/v1/compare/record", json=_food_payload("kw-a"), headers=_auth(token))
def _search(kw: str, **extra) -> list[str]:
r = client.get(
"/api/v1/compare/records",
params={"keyword": kw, **extra},
headers=_auth(token),
)
assert r.status_code == 200, r.text
return [it["store_name"] for it in r.json()["items"]]
assert _search("海底捞") == ["海底捞(朝阳店)"] # 店名命中
assert _search("榴莲") == ["Pizza Hut"] # 菜名命中(product_names)
assert _search("pizza") == ["Pizza Hut"] # 忽略大小写
assert _search("PIZZA") == ["Pizza Hut"]
assert _search("不存在的店") == [] # 没命中就是空
# 通配符只当普通字符:搜 % 不该把整表拉回来,搜 _ 也不该匹配任意单字符
assert _search("%") == ["100%纯牛肉汉堡"]
assert _search("_") == []
# 纯空白等同不传 → 不过滤
assert len(_search(" ")) == 3
# 搜索结果自身可翻页
for i in range(3):
p = _food_payload(f"kw-page-{i}")
p["store_name"] = "连锁烤鱼店"
client.post("/api/v1/compare/record", json=p, headers=_auth(token))
p1 = client.get(
"/api/v1/compare/records",
params={"keyword": "烤鱼", "limit": 2},
headers=_auth(token),
).json()
assert len(p1["items"]) == 2
assert p1["next_cursor"] is not None
p2 = client.get(
"/api/v1/compare/records",
params={"keyword": "烤鱼", "limit": 2, "cursor": p1["next_cursor"]},
headers=_auth(token),
).json()
assert [it["store_name"] for it in p2["items"]] == ["连锁烤鱼店"]
assert len({it["id"] for it in p1["items"] + p2["items"]}) == 3
def test_requires_auth(client) -> None:
"""不带 token 统一 401。"""
assert client.post("/api/v1/compare/record", json={"trace_id": "t"}).status_code == 401
-16
View File
@@ -4,28 +4,12 @@ from __future__ import annotations
from datetime import UTC, date, datetime
import pytest
from sqlalchemy.dialects import postgresql
from app.admin.repositories import queries
from app.db.session import SessionLocal
from app.models.comparison import ComparisonRecord
def test_postgresql_duration_summary_uses_ordered_set_aggregates() -> None:
stmt = queries._comparison_duration_aggregate_stmt(
[], "success", (0.05, 0.5, 0.95, 0.99)
)
sql = str(
stmt.compile(
dialect=postgresql.dialect(),
compile_kwargs={"literal_binds": True},
)
)
assert sql.count("percentile_cont") == 4
assert "comparison_record.status = 'success'" in sql
def test_summary_uses_only_success_durations_and_filters_beijing_date() -> None:
db = SessionLocal()
try:
+2 -55
View File
@@ -7,12 +7,8 @@ from __future__ import annotations
from datetime import UTC, date, datetime
from sqlalchemy import delete, func, select
from sqlalchemy.dialects import postgresql
from app.admin.repositories.coupon_data import (
_coupon_summary_aggregate_stmt,
coupon_data_report,
)
from app.admin.repositories.coupon_data import coupon_data_report
from app.db.session import SessionLocal
from app.models.coupon_state import CouponSession
from app.repositories.coupon_state import (
@@ -23,12 +19,7 @@ from app.repositories.coupon_state import (
def _agg_session(
trace: str,
platforms,
platform_success,
*,
status: str = "completed",
elapsed_ms: int | None = None,
trace: str, platforms, platform_success, *, status: str = "completed"
) -> CouponSession:
"""构造一条聚合测试用 session(started_date 固定 2020-01-02、app_env=prod,不 commit)。"""
return CouponSession(
@@ -40,53 +31,9 @@ def _agg_session(
platform_success=platform_success,
started_at=datetime(2020, 1, 2, tzinfo=UTC),
started_date=date(2020, 1, 2),
elapsed_ms=elapsed_ms,
)
def test_postgresql_coupon_summary_uses_ordered_set_aggregates() -> None:
sql = str(
_coupon_summary_aggregate_stmt([]).compile(
dialect=postgresql.dialect(),
compile_kwargs={"literal_binds": True},
)
)
assert sql.count("percentile_cont") == 4
assert "FILTER (WHERE coupon_session.status = 'completed'" in sql
def test_coupon_duration_summary_uses_only_completed_rows() -> None:
db = SessionLocal()
try:
db.add_all([
_agg_session("duration-a", [], [], elapsed_ms=1000),
_agg_session("duration-b", [], [], elapsed_ms=3000),
_agg_session(
"duration-failed", [], [], status="failed", elapsed_ms=100_000
),
])
db.flush()
summary = coupon_data_report(
db,
date_from="2020-01-02",
date_to="2020-01-02",
app_env="prod",
)["summary"]
assert summary["started_count"] == 3
assert summary["completed_count"] == 2
assert summary["avg_elapsed_ms"] == 2000
assert summary["p5_ms"] == 1100
assert summary["p50_ms"] == 2000
assert summary["p95_ms"] == 2900
assert summary["p99_ms"] == 2980
finally:
db.rollback()
db.close()
def _make_session(db, trace_id: str, **kw) -> CouponSession:
row = CouponSession(
trace_id=trace_id,