diff --git a/alembic/versions/coupon_claim_event.py b/alembic/versions/coupon_claim_event.py new file mode 100644 index 0000000..8161330 --- /dev/null +++ b/alembic/versions/coupon_claim_event.py @@ -0,0 +1,93 @@ +"""add per-session coupon claim event table + +Revision ID: coupon_claim_event +Revises: 8e04cc13a211 +Create Date: 2026-07-23 +""" + +from collections.abc import Sequence + +import sqlalchemy as sa +from sqlalchemy.dialects import postgresql + +from alembic import op + +revision: str = "coupon_claim_event" +down_revision: str | Sequence[str] | None = "8e04cc13a211" +branch_labels: str | Sequence[str] | None = None +depends_on: str | Sequence[str] | None = None + +_JSON = sa.JSON().with_variant(postgresql.JSONB(), "postgresql") + + +def upgrade() -> None: + op.create_table( + "coupon_claim_event", + sa.Column("id", sa.Integer(), autoincrement=True, nullable=False), + sa.Column("trace_id", sa.String(length=64), nullable=False), + sa.Column("device_id", sa.String(length=64), nullable=False), + sa.Column("user_id", sa.Integer(), nullable=True), + sa.Column("coupon_id", sa.String(length=64), nullable=False), + sa.Column("claim_date", sa.Date(), nullable=False), + sa.Column("status", sa.String(length=24), nullable=False), + sa.Column("app_env", sa.String(length=16), nullable=True), + sa.Column("vendor", sa.String(length=48), nullable=True), + sa.Column("coupon_name", sa.String(length=128), nullable=True), + sa.Column("claimed_count", sa.Integer(), nullable=True), + sa.Column("reason", sa.String(length=255), nullable=True), + sa.Column("extra", _JSON, nullable=True), + sa.Column("created_at", sa.DateTime(timezone=True), server_default=sa.func.now(), nullable=False), + sa.Column("updated_at", sa.DateTime(timezone=True), server_default=sa.func.now(), nullable=False), + sa.PrimaryKeyConstraint("id"), + sa.UniqueConstraint( + "trace_id", "coupon_id", + name="uq_coupon_claim_event_trace_coupon", + ), + ) + op.create_index( + "ix_coupon_claim_event_date_env", + "coupon_claim_event", + ["claim_date", "app_env"], + unique=False, + ) + op.create_index( + op.f("ix_coupon_claim_event_app_env"), + "coupon_claim_event", + ["app_env"], + unique=False, + ) + op.create_index( + op.f("ix_coupon_claim_event_trace_id"), + "coupon_claim_event", + ["trace_id"], + unique=False, + ) + op.create_index( + op.f("ix_coupon_claim_event_user_id"), + "coupon_claim_event", + ["user_id"], + unique=False, + ) + + # 旧表只能回填当前仍保留的 trace;历史上已被每日去重覆盖的关联无法恢复。 + op.execute( + """ + INSERT INTO coupon_claim_event ( + trace_id, device_id, user_id, coupon_id, claim_date, status, app_env, + vendor, coupon_name, claimed_count, reason, extra, created_at, updated_at + ) + SELECT + trace_id, device_id, user_id, coupon_id, claim_date, status, app_env, + vendor, coupon_name, claimed_count, reason, extra, created_at, updated_at + FROM coupon_claim_record + WHERE trace_id IS NOT NULL + """ + ) + + +def downgrade() -> None: + op.drop_index(op.f("ix_coupon_claim_event_user_id"), table_name="coupon_claim_event") + op.drop_index(op.f("ix_coupon_claim_event_trace_id"), table_name="coupon_claim_event") + op.drop_index(op.f("ix_coupon_claim_event_app_env"), table_name="coupon_claim_event") + op.drop_index("ix_coupon_claim_event_date_env", table_name="coupon_claim_event") + op.drop_table("coupon_claim_event") diff --git a/alembic/versions/d8dd2106e438_merge_guide_video_seq_uq_and_coupon_.py b/alembic/versions/d8dd2106e438_merge_guide_video_seq_uq_and_coupon_.py new file mode 100644 index 0000000..a8032c9 --- /dev/null +++ b/alembic/versions/d8dd2106e438_merge_guide_video_seq_uq_and_coupon_.py @@ -0,0 +1,26 @@ +"""merge guide_video seq_uq and coupon_claim_event heads + +Revision ID: d8dd2106e438 +Revises: coupon_claim_event, guide_video_user_seq_uq +Create Date: 2026-07-24 11:52:18.290731 + +""" +from typing import Sequence, Union + +from alembic import op +import sqlalchemy as sa + + +# revision identifiers, used by Alembic. +revision: str = 'd8dd2106e438' +down_revision: Union[str, Sequence[str], None] = ('coupon_claim_event', 'guide_video_user_seq_uq') +branch_labels: Union[str, Sequence[str], None] = None +depends_on: Union[str, Sequence[str], None] = None + + +def upgrade() -> None: + pass + + +def downgrade() -> None: + pass diff --git a/alembic/versions/d9c03cc3ea07_merge_guide_video_and_main_alembic_heads.py b/alembic/versions/d9c03cc3ea07_merge_guide_video_and_main_alembic_heads.py new file mode 100644 index 0000000..0834cae --- /dev/null +++ b/alembic/versions/d9c03cc3ea07_merge_guide_video_and_main_alembic_heads.py @@ -0,0 +1,26 @@ +"""merge guide_video and main alembic heads + +Revision ID: d9c03cc3ea07 +Revises: 8e04cc13a211, guide_video_play_table +Create Date: 2026-07-23 22:57:40.998161 + +""" +from typing import Sequence, Union + +from alembic import op +import sqlalchemy as sa + + +# revision identifiers, used by Alembic. +revision: str = 'd9c03cc3ea07' +down_revision: Union[str, Sequence[str], None] = ('8e04cc13a211', 'guide_video_play_table') +branch_labels: Union[str, Sequence[str], None] = None +depends_on: Union[str, Sequence[str], None] = None + + +def upgrade() -> None: + pass + + +def downgrade() -> None: + pass diff --git a/alembic/versions/guide_video_play_table.py b/alembic/versions/guide_video_play_table.py new file mode 100644 index 0000000..3c6a62e --- /dev/null +++ b/alembic/versions/guide_video_play_table.py @@ -0,0 +1,50 @@ +"""新手引导视频播放记录表(领券浮层前 N 次替代广告) + +见 app/models/guide_video.py:按账号计次(开播即计数)、play_token 幂等发币。 +配置(开关 / 视频地址 / 次数 / 金币)复用既有 app_config 表,无需建表。 + +Revision ID: guide_video_play_table +Revises: meituan_coupon_feed_indexes +Create Date: 2026-07-23 12:00:00.000000 + +""" +from typing import Sequence, Union + +import sqlalchemy as sa +from alembic import op + +# revision identifiers, used by Alembic. +revision: str = "guide_video_play_table" +down_revision: Union[str, Sequence[str], None] = "meituan_coupon_feed_indexes" +branch_labels: Union[str, Sequence[str], None] = None +depends_on: Union[str, Sequence[str], None] = None + + +def upgrade() -> None: + op.create_table( + "guide_video_play", + sa.Column("id", sa.Integer(), autoincrement=True, nullable=False), + sa.Column("user_id", sa.Integer(), nullable=False), + sa.Column("play_token", sa.String(length=64), nullable=False), + sa.Column("scene", sa.String(length=16), nullable=False, server_default="coupon"), + sa.Column("seq", sa.Integer(), nullable=False, server_default="1"), + sa.Column("video_url", sa.String(length=512), nullable=True), + sa.Column("coin", sa.Integer(), nullable=False, server_default="0"), + sa.Column("status", sa.String(length=16), nullable=False, server_default="playing"), + sa.Column("completed", sa.Integer(), nullable=False, server_default="0"), + sa.Column( + "started_at", sa.DateTime(timezone=True), server_default=sa.func.now(), nullable=False + ), + sa.Column("granted_at", sa.DateTime(timezone=True), nullable=True), + sa.ForeignKeyConstraint(["user_id"], ["user.id"]), + sa.PrimaryKeyConstraint("id"), + sa.UniqueConstraint("play_token", name="uq_guide_video_play_token"), + ) + op.create_index("ix_guide_video_play_user_id", "guide_video_play", ["user_id"]) + op.create_index("ix_guide_video_play_started_at", "guide_video_play", ["started_at"]) + + +def downgrade() -> None: + op.drop_index("ix_guide_video_play_started_at", table_name="guide_video_play") + op.drop_index("ix_guide_video_play_user_id", table_name="guide_video_play") + op.drop_table("guide_video_play") diff --git a/alembic/versions/guide_video_play_user_seq_unique.py b/alembic/versions/guide_video_play_user_seq_unique.py new file mode 100644 index 0000000..e4c105e --- /dev/null +++ b/alembic/versions/guide_video_play_user_seq_unique.py @@ -0,0 +1,40 @@ +"""guide_video_play 加 (user_id, seq) 唯一约束:堵住并发 /start 绕过次数上限 + +start_play 是无锁 check-then-insert(读 COUNT(*) 算 seq=used+1 再插一行),N 个并发 +/start 会都读到同一个 used、算出同一个 seq、各插一行拿到各自的 play_token,于是 3 次 +上限被绕过、每个 token 都能换 120 金币。加唯一键后并发同 seq 必撞,start_play 捕获 +IntegrityError 降级为 should_play=false(客户端照旧放广告)。 + +用 unique index 而不是 batch_alter_table 加 UniqueConstraint:SQLite 加约束要整表重建, +而 CREATE UNIQUE INDEX 两边都原生支持,回滚也干净。 + +注:若库里已有并发产生的重复 (user_id, seq),建索引会失败 —— 本功能尚未上线,表通常是空的; +真撞上了先按 seq 去重(留 id 最小的一行,多发的金币按 scripts/reset_guide_video.py 的口径退)。 + +Revision ID: guide_video_user_seq_uq +Revises: d9c03cc3ea07 +Create Date: 2026-07-24 10:00:00.000000 + +""" +from typing import Sequence, Union + +from alembic import op + +# revision identifiers, used by Alembic. +revision: str = "guide_video_user_seq_uq" +down_revision: Union[str, Sequence[str], None] = "d9c03cc3ea07" +branch_labels: Union[str, Sequence[str], None] = None +depends_on: Union[str, Sequence[str], None] = None + + +def upgrade() -> None: + op.create_index( + "uq_guide_video_play_user_seq", + "guide_video_play", + ["user_id", "seq"], + unique=True, + ) + + +def downgrade() -> None: + op.drop_index("uq_guide_video_play_user_seq", table_name="guide_video_play") diff --git a/alembic/versions/meituan_coupon_feed_indexes.py b/alembic/versions/meituan_coupon_feed_indexes.py new file mode 100644 index 0000000..c120ca6 --- /dev/null +++ b/alembic/versions/meituan_coupon_feed_indexes.py @@ -0,0 +1,52 @@ +"""meituan_coupon 首页 feed 分页复合索引(销量最高 / 智能推荐) + +「销量最高」「智能推荐」两个 tab 都是 + WHERE city_id = ? [+ 过滤] → DISTINCT ON (dedup_key) ORDER BY dedup_key, <排序键> DESC +的形状。列顺序对齐后 Postgres 可以顺着索引流式去重,免掉「每翻一页就把该城全部券重排一遍」, +这是首页下滑到底越来越慢的根因之一(另一半在 app 层:见 api/v1/meituan.py 的 _paged_dedup_ids)。 + +⚠️ 本文件同时是一个 **merge 迁移**:主干此前有 3 个并行 head +(comparison_user_created_idx / monitoring_audit_rbac / notification_table), +`alembic upgrade head` 会因 multiple heads 报错。这里一并收敛回单 head。 + +Revision ID: meituan_coupon_feed_indexes +Revises: comparison_user_created_idx, monitoring_audit_rbac, notification_table +Create Date: 2026-07-23 10:00:00.000000 + +""" +from typing import Sequence, Union + +import sqlalchemy as sa +from alembic import op + +# revision identifiers, used by Alembic. +revision: str = "meituan_coupon_feed_indexes" +down_revision: Union[str, Sequence[str], None] = ( + "comparison_user_created_idx", + "monitoring_audit_rbac", + "notification_table", +) +branch_labels: Union[str, Sequence[str], None] = None +depends_on: Union[str, Sequence[str], None] = None + + +def upgrade() -> None: + # 销量最高:WHERE city_id=? AND sale_volume_num IS NOT NULL + # ORDER BY dedup_key, sale_volume_num DESC, commission_percent DESC + op.create_index( + "ix_meituan_coupon_city_dedup_sales", + "meituan_coupon", + ["city_id", "dedup_key", sa.text("sale_volume_num DESC"), sa.text("commission_percent DESC")], + ) + # 智能推荐:WHERE city_id=? AND commission_percent>=3.0 + # ORDER BY dedup_key, commission_percent DESC + op.create_index( + "ix_meituan_coupon_city_dedup_comm", + "meituan_coupon", + ["city_id", "dedup_key", sa.text("commission_percent DESC")], + ) + + +def downgrade() -> None: + op.drop_index("ix_meituan_coupon_city_dedup_comm", table_name="meituan_coupon") + op.drop_index("ix_meituan_coupon_city_dedup_sales", table_name="meituan_coupon") diff --git a/app/admin/main.py b/app/admin/main.py index c869487..918b238 100644 --- a/app/admin/main.py +++ b/app/admin/main.py @@ -30,6 +30,7 @@ from app.admin.routers.analytics_health import router as analytics_health_router from app.admin.routers.event_logs import router as event_logs_router from app.admin.routers.feedback import router as feedback_router from app.admin.routers.feedback_qr import router as feedback_qr_router +from app.admin.routers.guide_video import router as guide_video_router from app.admin.routers.huawei_review import router as huawei_review_router from app.admin.routers.onboarding import router as onboarding_router from app.admin.routers.ops_marquee_seed import router as ops_marquee_seed_router @@ -40,6 +41,7 @@ from app.admin.routers.wallet import router as wallet_router from app.admin.routers.withdraw import router as withdraw_router from app.core.config import settings from app.core.logging import setup_logging +from app.integrations import meituan as mt_meituan setup_logging(debug=settings.APP_DEBUG) logger = logging.getLogger("shagua.admin") @@ -53,6 +55,8 @@ async def lifespan(_: FastAPI) -> AsyncIterator[None]: settings.DATABASE_URL.split("://", 1)[0], ) yield + # CPS 后台页会打美团(routers/cps.py),那条共享 client 若被建过要在这里关掉连接池 + mt_meituan.close_client() logger.info("admin app shutting down") @@ -101,6 +105,7 @@ admin_app.include_router(feedback_router) admin_app.include_router(event_logs_router) admin_app.include_router(analytics_health_router) admin_app.include_router(feedback_qr_router) +admin_app.include_router(guide_video_router) admin_app.include_router(admins_router) admin_app.include_router(roles_router) admin_app.include_router(audit_router) diff --git a/app/admin/repositories/coupon_data.py b/app/admin/repositories/coupon_data.py index 5ae22d4..f7f8179 100644 --- a/app/admin/repositories/coupon_data.py +++ b/app/admin/repositories/coupon_data.py @@ -16,7 +16,7 @@ from sqlalchemy import case, func, or_, select from sqlalchemy.orm import Session from app.core import rewards -from app.models.coupon_state import CouponClaimRecord, CouponSession +from app.models.coupon_state import CouponClaimEvent, CouponClaimRecord, CouponSession from app.models.user import User from app.repositories import ad_ecpm as crud_ecpm from app.repositories.coupon_state import DEFAULT_PLATFORMS, coupon_id_to_platform @@ -193,18 +193,18 @@ def _point_scores_by_trace(db: Session, trace_ids: list[str]) -> dict[str, dict[ """聚合查询批量返回逐场点位分数,不加载逐券明细。""" if not trace_ids: return {} - succeeded = func.sum(case((CouponClaimRecord.status.in_(_SLOT_OK), 1), else_=0)) + succeeded = func.sum(case((CouponClaimEvent.status.in_(_SLOT_OK), 1), else_=0)) rows = db.execute( select( - CouponClaimRecord.trace_id, + CouponClaimEvent.trace_id, succeeded.label("succeeded"), func.count().label("tried"), ) .where( - CouponClaimRecord.trace_id.in_(trace_ids), - CouponClaimRecord.status.in_(_SLOT_TRIED), + CouponClaimEvent.trace_id.in_(trace_ids), + CouponClaimEvent.status.in_(_SLOT_TRIED), ) - .group_by(CouponClaimRecord.trace_id) + .group_by(CouponClaimEvent.trace_id) ).all() return { trace_id: {"succeeded": int(success_count or 0), "tried": int(tried or 0)} @@ -217,13 +217,13 @@ def coupon_point_details(db: Session, *, trace_id: str) -> list[dict]: """按单个 trace 查询逐券结果;仅在后台用户点击分数时调用。""" rows = db.execute( select( - CouponClaimRecord.coupon_id, - CouponClaimRecord.coupon_name, - CouponClaimRecord.status, - CouponClaimRecord.reason, + CouponClaimEvent.coupon_id, + CouponClaimEvent.coupon_name, + CouponClaimEvent.status, + CouponClaimEvent.reason, ) - .where(CouponClaimRecord.trace_id == trace_id) - .order_by(CouponClaimRecord.id) + .where(CouponClaimEvent.trace_id == trace_id) + .order_by(CouponClaimEvent.id) ).all() return [ { diff --git a/app/admin/repositories/queries.py b/app/admin/repositories/queries.py index 33b510d..75550e7 100644 --- a/app/admin/repositories/queries.py +++ b/app/admin/repositories/queries.py @@ -1171,24 +1171,30 @@ def user_reward_stats( acc = db.get(CoinAccount, user_id) # 现金余额:当前快照,不随窗口 cash_balance = acc.cash_balance_cents if acc else 0 - rv = list(db.execute( - select(AdRewardRecord).where( + # 只投影本统计实际使用的列。避免滚动发布或旧本地库尚未补齐无关新列时, + # SQLAlchemy 因 select(ORM) 自动展开整表字段而让提现详情整体 500。 + rv = db.execute( + select(AdRewardRecord.ecpm_raw, AdRewardRecord.coin).where( AdRewardRecord.user_id == user_id, AdRewardRecord.reward_scene == "reward_video", AdRewardRecord.status == "granted", *_window_conds(AdRewardRecord.created_at, date_from, date_to), ) - ).scalars()) + ).all() rv_ecpms = [rewards.parse_ecpm_fen(r.ecpm_raw) for r in rv if r.ecpm_raw] rv_coins = sum(r.coin for r in rv) - feed = list(db.execute( - select(AdFeedRewardRecord).where( + feed = db.execute( + select( + AdFeedRewardRecord.unit_count, + AdFeedRewardRecord.ecpm_raw, + AdFeedRewardRecord.coin, + ).where( AdFeedRewardRecord.user_id == user_id, AdFeedRewardRecord.status == "granted", *_window_conds(AdFeedRewardRecord.created_at, date_from, date_to), ) - ).scalars()) + ).all() feed_ecpms = [rewards.parse_ecpm_fen(f.ecpm_raw) for f in feed if f.ecpm_raw] feed_coins = sum(f.coin for f in feed) @@ -1246,8 +1252,14 @@ def user_coin_records( signin_from = date_from + timedelta(hours=8) if date_from is not None else None signin_to = date_to + timedelta(hours=8) if date_to is not None else None + # 三类来源都只取页面需要的列,避免无关 ORM 新列造成旧库查询失败。 for rec in db.execute( - select(AdRewardRecord) + select( + AdRewardRecord.reward_scene, + AdRewardRecord.created_at, + AdRewardRecord.ecpm_raw, + AdRewardRecord.coin, + ) .where( AdRewardRecord.user_id == user_id, AdRewardRecord.status == "granted", @@ -1255,7 +1267,7 @@ def user_coin_records( ) .order_by(AdRewardRecord.created_at.desc()) .limit(fetch) - ).scalars(): + ).all(): is_video = rec.reward_scene == "reward_video" rows.append({ "source": rec.reward_scene, @@ -1266,7 +1278,12 @@ def user_coin_records( }) for rec in db.execute( - select(AdFeedRewardRecord) + select( + AdFeedRewardRecord.feed_scene, + AdFeedRewardRecord.created_at, + AdFeedRewardRecord.ecpm_raw, + AdFeedRewardRecord.coin, + ) .where( AdFeedRewardRecord.user_id == user_id, AdFeedRewardRecord.status == "granted", @@ -1274,7 +1291,7 @@ def user_coin_records( ) .order_by(AdFeedRewardRecord.created_at.desc()) .limit(fetch) - ).scalars(): + ).all(): rows.append({ "source": "feed", "source_label": _FEED_SCENE_LABEL.get(rec.feed_scene, "信息流广告"), @@ -1284,7 +1301,7 @@ def user_coin_records( }) for rec in db.execute( - select(CoinTransaction) + select(CoinTransaction.created_at, CoinTransaction.amount) .where( CoinTransaction.user_id == user_id, CoinTransaction.biz_type == "signin", @@ -1292,7 +1309,7 @@ def user_coin_records( ) .order_by(CoinTransaction.created_at.desc()) .limit(fetch) - ).scalars(): + ).all(): rows.append({ "source": "signin", "source_label": "签到", diff --git a/app/admin/repositories/stats.py b/app/admin/repositories/stats.py index e82f5b6..26e13df 100644 --- a/app/admin/repositories/stats.py +++ b/app/admin/repositories/stats.py @@ -30,18 +30,17 @@ from app.models.user import User from app.models.wallet import CoinTransaction, WithdrawOrder _BEIJING = timezone(timedelta(hours=8)) -REWARD_VIDEO_BIZ_TYPES = ("reward_video", "ad_reward") +REWARD_VIDEO_BIZ_TYPES = ("reward_video", "ad_reward", "signin_boost") # 领券/比价奖励金币的真实来源是信息流广告发奖(ad_feed_reward_record,按 feed_scene 分场景); # coin_transaction 里只有扁平的 feed_ad_reward、biz_type 不分 coupon/comparison,故这俩桶历史从未 -# 被写入,仅留作未来兜底,实际金额在下方按 feed_scene 汇总 ad_feed_reward_record 得出。reward_video/ -# ad_reward 是激励视频,单独成桶、不再混进领券奖励(历史误并会把激励视频金币双计进领券)。 +# 被写入,仅留作未来兜底,实际金额在下方按 feed_scene 汇总 ad_feed_reward_record 得出。 +# reward_video/ad_reward 及历史 signin_boost 均归看视频桶,不再混进领券奖励或常规任务。 COUPON_REWARD_BIZ_TYPES = ("coupon", "coupon_reward") COMPARISON_REWARD_BIZ_TYPES = ("comparison", "compare_reward", "comparison_reward") # 常规任务必须按明确来源相加;不能从全部正向流水反减排除项,否则新增广告/运营 # biz_type 时会在排除清单更新前自动混入该桶。task_ 前缀在查询处单独覆盖现有及未来任务。 REGULAR_TASK_EXACT_BIZ_TYPES = ( "signin", - "signin_boost", "price_report_reward", "feedback_reward", ) diff --git a/app/admin/routers/feedback.py b/app/admin/routers/feedback.py index 987c76b..28b5b71 100644 --- a/app/admin/routers/feedback.py +++ b/app/admin/routers/feedback.py @@ -12,6 +12,10 @@ from app.admin.repositories import mutations, queries from app.admin.schemas.common import CursorPage, OkResponse from app.admin.schemas.feedback import ( FeedbackApproveRequest, + FeedbackBulkApproveRequest, + FeedbackBulkItemResult, + FeedbackBulkRejectRequest, + FeedbackBulkResult, FeedbackOut, FeedbackRejectRequest, FeedbackSummary, @@ -33,6 +37,123 @@ def _ensure_pending(fb: Feedback) -> None: raise HTTPException(status_code=400, detail="反馈已审核") +def _approve_feedback( + db: AdminDb, + admin: AdminUser, + feedback_id: int, + payload: FeedbackApproveRequest | FeedbackBulkApproveRequest, + ip: str, + *, + bulk: bool = False, +) -> FeedbackOut: + fb = db.get(Feedback, feedback_id, with_for_update=True) + if fb is None: + raise HTTPException(status_code=404, detail="反馈不存在") + _ensure_pending(fb) + + before = fb.status + mutations.review_feedback( + db, + fb, + status="adopted", + reward_coins=payload.reward_coins, + review_note=payload.note, + admin_reply=payload.reply, + reviewed_by_admin_id=admin.id, + commit=False, + ) + wallet_repo.grant_coins( + db, + fb.user_id, + payload.reward_coins, + biz_type="feedback_reward", + ref_id=str(fb.id), + remark="意见反馈被采纳", + ) + detail = { + "before": before, + "after": "adopted", + "reward_coins": payload.reward_coins, + "note": payload.note, + "reply": payload.reply, + } + if bulk: + detail["bulk"] = True + write_audit( + db, + admin, + action="feedback.approve", + target_type="feedback", + target_id=feedback_id, + detail=detail, + ip=ip, + commit=False, + ) + db.commit() + db.refresh(fb) + out = FeedbackOut.model_validate(fb) + notification_events.notify_feedback_reward(db, fb) + return out + + +def _reject_feedback( + db: AdminDb, + admin: AdminUser, + feedback_id: int, + payload: FeedbackRejectRequest | FeedbackBulkRejectRequest, + ip: str, + *, + bulk: bool = False, +) -> FeedbackOut: + fb = db.get(Feedback, feedback_id, with_for_update=True) + if fb is None: + raise HTTPException(status_code=404, detail="反馈不存在") + _ensure_pending(fb) + + before = fb.status + mutations.review_feedback( + db, + fb, + status="rejected", + reject_reason=payload.reason, + review_note=payload.note, + admin_reply=payload.reply, + reviewed_by_admin_id=admin.id, + commit=False, + ) + detail = { + "before": before, + "after": "rejected", + "reason": payload.reason, + "note": payload.note, + "reply": payload.reply, + } + if bulk: + detail["bulk"] = True + write_audit( + db, + admin, + action="feedback.reject", + target_type="feedback", + target_id=feedback_id, + detail=detail, + ip=ip, + commit=False, + ) + db.commit() + db.refresh(fb) + out = FeedbackOut.model_validate(fb) + notification_events.notify_feedback_reply(db, fb) + return out + + +def _bulk_result(items: list[FeedbackBulkItemResult]) -> FeedbackBulkResult: + success = sum(1 for item in items if item.ok) + return FeedbackBulkResult( + total=len(items), success=success, failed=len(items) - success, items=items, + ) + + @router.get("", response_model=CursorPage[FeedbackOut], summary="反馈工单列表") def list_feedbacks( db: AdminDb, @@ -73,6 +194,50 @@ def feedback_summary(db: AdminDb) -> FeedbackSummary: return FeedbackSummary.model_validate(queries.feedback_summary(db)) +@router.post("/bulk/approve", response_model=FeedbackBulkResult, summary="批量采纳反馈并发金币") +def bulk_approve_feedbacks( + body: FeedbackBulkApproveRequest, + request: Request, + admin: Annotated[AdminUser, Depends(require_role("operator"))], + db: AdminDb, +) -> FeedbackBulkResult: + results: list[FeedbackBulkItemResult] = [] + ip = get_client_ip(request) + for feedback_id in body.ids: + try: + out = _approve_feedback(db, admin, feedback_id, body, ip, bulk=True) + results.append(FeedbackBulkItemResult(id=feedback_id, ok=True, status=out.status)) + except HTTPException as exc: + db.rollback() + results.append(FeedbackBulkItemResult(id=feedback_id, ok=False, error=str(exc.detail))) + except Exception: # noqa: BLE001 - 单笔失败不打断整批 + db.rollback() + results.append(FeedbackBulkItemResult(id=feedback_id, ok=False, error="系统异常")) + return _bulk_result(results) + + +@router.post("/bulk/reject", response_model=FeedbackBulkResult, summary="批量拒绝采纳反馈") +def bulk_reject_feedbacks( + body: FeedbackBulkRejectRequest, + request: Request, + admin: Annotated[AdminUser, Depends(require_role("operator"))], + db: AdminDb, +) -> FeedbackBulkResult: + results: list[FeedbackBulkItemResult] = [] + ip = get_client_ip(request) + for feedback_id in body.ids: + try: + out = _reject_feedback(db, admin, feedback_id, body, ip, bulk=True) + results.append(FeedbackBulkItemResult(id=feedback_id, ok=True, status=out.status)) + except HTTPException as exc: + db.rollback() + results.append(FeedbackBulkItemResult(id=feedback_id, ok=False, error=str(exc.detail))) + except Exception: # noqa: BLE001 - 单笔失败不打断整批 + db.rollback() + results.append(FeedbackBulkItemResult(id=feedback_id, ok=False, error="系统异常")) + return _bulk_result(results) + + @router.post("/{feedback_id}/handle", response_model=OkResponse, summary="标记反馈已处理") def handle_feedback( feedback_id: int, @@ -93,53 +258,7 @@ def approve_feedback( admin: Annotated[AdminUser, Depends(require_role("operator"))], db: AdminDb, ) -> FeedbackOut: - fb = db.get(Feedback, feedback_id) - if fb is None: - raise HTTPException(status_code=404, detail="反馈不存在") - _ensure_pending(fb) - - before = fb.status - mutations.review_feedback( - db, - fb, - status="adopted", - reward_coins=payload.reward_coins, - review_note=payload.note, - admin_reply=payload.reply, - reviewed_by_admin_id=admin.id, - commit=False, - ) - wallet_repo.grant_coins( - db, - fb.user_id, - payload.reward_coins, - biz_type="feedback_reward", - ref_id=str(fb.id), - remark="意见反馈被采纳", - ) - write_audit( - db, - admin, - action="feedback.approve", - target_type="feedback", - target_id=feedback_id, - detail={ - "before": before, - "after": "adopted", - "reward_coins": payload.reward_coins, - "note": payload.note, - "reply": payload.reply, - }, - ip=get_client_ip(request), - commit=False, - ) - db.commit() - db.refresh(fb) - out = FeedbackOut.model_validate(fb) - # PRD #10 反馈奖励:采纳发金币后通知用户(站内 + push,必带官方留言)。 - # 业务已 commit,通知失败只 log 不影响审核结果。 - notification_events.notify_feedback_reward(db, fb) - return out + return _approve_feedback(db, admin, feedback_id, payload, get_client_ip(request)) @router.post("/{feedback_id}/reject", response_model=FeedbackOut, summary="拒绝采纳反馈") @@ -150,41 +269,4 @@ def reject_feedback( admin: Annotated[AdminUser, Depends(require_role("operator"))], db: AdminDb, ) -> FeedbackOut: - fb = db.get(Feedback, feedback_id) - if fb is None: - raise HTTPException(status_code=404, detail="反馈不存在") - _ensure_pending(fb) - - before = fb.status - mutations.review_feedback( - db, - fb, - status="rejected", - reject_reason=payload.reason, - review_note=payload.note, - admin_reply=payload.reply, - reviewed_by_admin_id=admin.id, - commit=False, - ) - write_audit( - db, - admin, - action="feedback.reject", - target_type="feedback", - target_id=feedback_id, - detail={ - "before": before, - "after": "rejected", - "reason": payload.reason, - "note": payload.note, - "reply": payload.reply, - }, - ip=get_client_ip(request), - commit=False, - ) - db.commit() - db.refresh(fb) - out = FeedbackOut.model_validate(fb) - # PRD #9 官方回复:未采纳也回复了用户(原因/留言用户端可见),通知去反馈历史页查看。 - notification_events.notify_feedback_reply(db, fb) - return out + return _reject_feedback(db, admin, feedback_id, payload, get_client_ip(request)) diff --git a/app/admin/routers/guide_video.py b/app/admin/routers/guide_video.py new file mode 100644 index 0000000..59f75ae --- /dev/null +++ b/app/admin/routers/guide_video.py @@ -0,0 +1,101 @@ +"""admin 新手引导视频配置:读 / 改开关次数金币 / 上传视频 / 删视频(带审计)。 + +整份配置存通用 app_config 表(见 app/repositories/guide_video.py),App 领券等候浮层 +每次展示前调 POST /api/v1/guide-video/start 同步。权限:operator 可改(运营维护), +super 恒可;读为只读(任意已登录 admin)。 + +⚠️ 视频上限 100MB(settings.GUIDE_VIDEO_MAX_BYTES),已在 admin nginx 为本接口单独放宽 +client_max_body_size,见 shaguabijia-admin-web/deploy/nginx/admin.shaguabijia.com.conf。 +""" +from __future__ import annotations + +from typing import Annotated + +from fastapi import APIRouter, Depends, File, HTTPException, Request, UploadFile + +from app.admin.audit import write_audit +from app.admin.deps import AdminDb, get_client_ip, get_current_admin, require_role +from app.admin.schemas.guide_video import GuideVideoConfigOut, GuideVideoConfigUpdate +from app.core import media +from app.models.admin import AdminUser +from app.repositories import guide_video + +router = APIRouter( + prefix="/admin/api/guide-video", + tags=["admin-guide-video"], + dependencies=[Depends(get_current_admin)], +) + + +def _out(db: AdminDb) -> GuideVideoConfigOut: + """配置 + 播放统计合成响应(四个写接口都以最新状态返回,前端一次同步到位)。""" + return GuideVideoConfigOut(**guide_video.get_config(db), **guide_video.play_stats(db)) + + +@router.get("", response_model=GuideVideoConfigOut, summary="新手引导视频配置(领券浮层)") +def get_config(db: AdminDb) -> GuideVideoConfigOut: + return _out(db) + + +@router.patch("", response_model=GuideVideoConfigOut, summary="改开关/次数/金币(带审计)") +def update_config( + body: GuideVideoConfigUpdate, + request: Request, + admin: Annotated[AdminUser, Depends(require_role("operator"))], + db: AdminDb, +) -> GuideVideoConfigOut: + before, after = guide_video.update_config( + db, + enabled=body.enabled, + max_plays=body.max_plays, + reward_coin=body.reward_coin, + admin_id=admin.id, + commit=False, + ) + write_audit( + db, admin, action="guide_video.update", target_type="guide_video", target_id=None, + detail={"before": before, "after": after}, ip=get_client_ip(request), commit=False, + ) + db.commit() + return _out(db) + + +@router.post("/video", response_model=GuideVideoConfigOut, summary="上传新手引导视频(MP4,带审计)") +async def upload_video( + request: Request, + admin: Annotated[AdminUser, Depends(require_role("operator"))], + db: AdminDb, + file: UploadFile = File(...), +) -> GuideVideoConfigOut: + data = await file.read() + try: + url = media.save_guide_video(data) + except media.MediaError as e: + raise HTTPException(status_code=400, detail=str(e)) from e + before, after = guide_video.set_video(db, url, admin_id=admin.id, commit=False) + write_audit( + db, admin, action="guide_video.set_video", target_type="guide_video", target_id=None, + detail={"before": before.get("video_url"), "after": url, "bytes": len(data)}, + ip=get_client_ip(request), commit=False, + ) + db.commit() + # 提交成功后再删旧片,避免新片没落库就把旧片丢了 + media.delete_guide_video(before.get("video_url")) + return _out(db) + + +@router.delete("/video", response_model=GuideVideoConfigOut, summary="移除新手引导视频(带审计)") +def delete_video( + request: Request, + admin: Annotated[AdminUser, Depends(require_role("operator"))], + db: AdminDb, +) -> GuideVideoConfigOut: + """移除后 /guide-video/start 一律返回 should_play=false,领券浮层回到「只放广告」。""" + before, after = guide_video.set_video(db, None, admin_id=admin.id, commit=False) + write_audit( + db, admin, action="guide_video.delete_video", target_type="guide_video", target_id=None, + detail={"before": before.get("video_url")}, ip=get_client_ip(request), commit=False, + ) + db.commit() + media.delete_guide_video(before.get("video_url")) + return _out(db) diff --git a/app/admin/routers/price_report.py b/app/admin/routers/price_report.py index 9f65545..6f656a7 100644 --- a/app/admin/routers/price_report.py +++ b/app/admin/routers/price_report.py @@ -17,6 +17,10 @@ from app.admin.deps import AdminDb, get_client_ip, get_current_admin, require_ro from app.admin.repositories import mutations, queries from app.admin.schemas.common import CursorPage, OkResponse from app.admin.schemas.price_report import ( + PriceReportBulkItemResult, + PriceReportBulkRejectRequest, + PriceReportBulkRequest, + PriceReportBulkResult, PriceReportOut, PriceReportRejectRequest, PriceReportSummary, @@ -34,6 +38,59 @@ router = APIRouter( ) +def _approve_price_report( + db: AdminDb, admin: AdminUser, report_id: int, ip: str, *, bulk: bool = False +) -> PriceReport: + rep = db.get(PriceReport, report_id, with_for_update=True) + if rep is None: + raise HTTPException(status_code=404, detail="上报记录不存在") + if rep.status != "pending": + raise HTTPException(status_code=400, detail=f"该上报已审核过(当前 {rep.status}),不可重复操作") + coins = PRICE_REPORT_REWARD_COINS + mutations.review_price_report(db, rep, status="approved", reward_coins=coins, commit=False) + wallet_repo.grant_coins( + db, rep.user_id, coins, + biz_type="price_report_reward", ref_id=str(rep.id), remark="上报更低价审核通过", + ) + detail = {"reward_coins": coins, "user_id": rep.user_id} + if bulk: + detail["bulk"] = True + write_audit( + db, admin, action="price_report.approve", target_type="price_report", target_id=report_id, + detail=detail, ip=ip, commit=False, + ) + db.commit() + notification_events.notify_report_approved(db, rep) + return rep + + +def _reject_price_report( + db: AdminDb, admin: AdminUser, report_id: int, reason: str, ip: str, *, bulk: bool = False +) -> PriceReport: + rep = db.get(PriceReport, report_id, with_for_update=True) + if rep is None: + raise HTTPException(status_code=404, detail="上报记录不存在") + if rep.status != "pending": + raise HTTPException(status_code=400, detail=f"该上报已审核过(当前 {rep.status}),不可重复操作") + mutations.review_price_report(db, rep, status="rejected", reject_reason=reason, commit=False) + detail = {"reason": reason, "user_id": rep.user_id} + if bulk: + detail["bulk"] = True + write_audit( + db, admin, action="price_report.reject", target_type="price_report", target_id=report_id, + detail=detail, ip=ip, commit=False, + ) + db.commit() + return rep + + +def _bulk_result(items: list[PriceReportBulkItemResult]) -> PriceReportBulkResult: + success = sum(1 for item in items if item.ok) + return PriceReportBulkResult( + total=len(items), success=success, failed=len(items) - success, items=items, + ) + + @router.get("", response_model=CursorPage[PriceReportOut], summary="上报更低价列表(筛选+分页)") def list_price_reports( db: AdminDb, @@ -60,6 +117,50 @@ def price_report_summary(db: AdminDb) -> PriceReportSummary: return PriceReportSummary.model_validate(queries.price_report_summary(db)) +@router.post("/bulk/approve", response_model=PriceReportBulkResult, summary="批量通过上报(发固定金币)") +def bulk_approve_price_reports( + body: PriceReportBulkRequest, + request: Request, + admin: Annotated[AdminUser, Depends(require_role("operator"))], + db: AdminDb, +) -> PriceReportBulkResult: + results: list[PriceReportBulkItemResult] = [] + ip = get_client_ip(request) + for report_id in body.ids: + try: + rep = _approve_price_report(db, admin, report_id, ip, bulk=True) + results.append(PriceReportBulkItemResult(id=report_id, ok=True, status=rep.status)) + except HTTPException as exc: + db.rollback() + results.append(PriceReportBulkItemResult(id=report_id, ok=False, error=str(exc.detail))) + except Exception: # noqa: BLE001 - 单笔失败不打断整批 + db.rollback() + results.append(PriceReportBulkItemResult(id=report_id, ok=False, error="系统异常")) + return _bulk_result(results) + + +@router.post("/bulk/reject", response_model=PriceReportBulkResult, summary="批量拒绝上报") +def bulk_reject_price_reports( + body: PriceReportBulkRejectRequest, + request: Request, + admin: Annotated[AdminUser, Depends(require_role("operator"))], + db: AdminDb, +) -> PriceReportBulkResult: + results: list[PriceReportBulkItemResult] = [] + ip = get_client_ip(request) + for report_id in body.ids: + try: + rep = _reject_price_report(db, admin, report_id, body.reason, ip, bulk=True) + results.append(PriceReportBulkItemResult(id=report_id, ok=True, status=rep.status)) + except HTTPException as exc: + db.rollback() + results.append(PriceReportBulkItemResult(id=report_id, ok=False, error=str(exc.detail))) + except Exception: # noqa: BLE001 - 单笔失败不打断整批 + db.rollback() + results.append(PriceReportBulkItemResult(id=report_id, ok=False, error="系统异常")) + return _bulk_result(results) + + @router.post("/{report_id}/approve", response_model=OkResponse, summary="通过上报(发固定金币)") def approve_price_report( report_id: int, @@ -67,27 +168,7 @@ def approve_price_report( admin: Annotated[AdminUser, Depends(require_role("operator"))], db: AdminDb, ) -> OkResponse: - # 行锁(SELECT FOR UPDATE):并发/连点双请求会都读到 pending → 各发一次金币双倍发奖, - # 锁住该行串行化,第二个请求拿锁后看到 approved → 走 400。SQLite 下 FOR UPDATE 为 no-op。 - rep = db.get(PriceReport, report_id, with_for_update=True) - if rep is None: - raise HTTPException(status_code=404, detail="上报记录不存在") - if rep.status != "pending": - raise HTTPException(status_code=400, detail=f"该上报已审核过(当前 {rep.status}),不可重复操作") - coins = PRICE_REPORT_REWARD_COINS - # 改状态 + 发金币 + 审计同一事务(commit=False),最后一起 commit:改了就有痕、发了就留账 - mutations.review_price_report(db, rep, status="approved", reward_coins=coins, commit=False) - wallet_repo.grant_coins( - db, rep.user_id, coins, - biz_type="price_report_reward", ref_id=str(rep.id), remark="上报更低价审核通过", - ) - write_audit( - db, admin, action="price_report.approve", target_type="price_report", target_id=report_id, - detail={"reward_coins": coins, "user_id": rep.user_id}, ip=get_client_ip(request), commit=False, - ) - db.commit() - # PRD #11 爆料审核通过:发金币后通知用户(站内 + push)。业务已 commit,通知失败只 log。 - notification_events.notify_report_approved(db, rep) + _approve_price_report(db, admin, report_id, get_client_ip(request)) return OkResponse() @@ -99,16 +180,5 @@ def reject_price_report( admin: Annotated[AdminUser, Depends(require_role("operator"))], db: AdminDb, ) -> OkResponse: - rep = db.get(PriceReport, report_id, with_for_update=True) # 行锁,同 approve(防并发重复审核) - if rep is None: - raise HTTPException(status_code=404, detail="上报记录不存在") - if rep.status != "pending": - raise HTTPException(status_code=400, detail=f"该上报已审核过(当前 {rep.status}),不可重复操作") - reason = body.reason.strip() - mutations.review_price_report(db, rep, status="rejected", reject_reason=reason, commit=False) - write_audit( - db, admin, action="price_report.reject", target_type="price_report", target_id=report_id, - detail={"reason": reason, "user_id": rep.user_id}, ip=get_client_ip(request), commit=False, - ) - db.commit() + _reject_price_report(db, admin, report_id, body.reason, get_client_ip(request)) return OkResponse() diff --git a/app/admin/schemas/coupon_data.py b/app/admin/schemas/coupon_data.py index f295650..6b19fd7 100644 --- a/app/admin/schemas/coupon_data.py +++ b/app/admin/schemas/coupon_data.py @@ -79,10 +79,10 @@ class CouponDataRow(BaseModel): started_at: datetime = Field(..., description="发起时刻(明细「时间」列)") claimed_count: int | None = None point_success_count: int | None = Field( - None, description="本次成功券点位数(success+already_claimed);无逐券埋点为空" + None, description="本次成功单券数(success+already_claimed);无逐券事件为空" ) point_total_count: int | None = Field( - None, description="本次尝试券点位数(success+already_claimed+failed,不含 skipped);无逐券埋点为空" + None, description="本次尝试单券数(success+already_claimed+failed,不含 skipped);无逐券事件为空" ) trace_url: str | None = Field(None, description="pricebot 公网 trace 链接(仅 completed 有);admin 渲染可点链接,无则显示可复制 trace_id") ad_revenue_yuan: float = Field( diff --git a/app/admin/schemas/feedback.py b/app/admin/schemas/feedback.py index bb588cf..dde6229 100644 --- a/app/admin/schemas/feedback.py +++ b/app/admin/schemas/feedback.py @@ -3,7 +3,7 @@ from __future__ import annotations from datetime import datetime -from pydantic import BaseModel, ConfigDict, Field +from pydantic import BaseModel, ConfigDict, Field, field_validator from app.core.rewards import FEEDBACK_REWARD_MAX_COINS @@ -55,6 +55,54 @@ class FeedbackRejectRequest(BaseModel): reply: str | None = Field(default=None, max_length=256, description="给用户的回复留言,用户端可见") +class FeedbackBulkRequest(BaseModel): + ids: list[int] = Field(min_length=1, max_length=50, description="待审核反馈 ID 列表") + + @field_validator("ids") + @classmethod + def _ids_must_be_unique(cls, ids: list[int]) -> list[int]: + if len(ids) != len(set(ids)): + raise ValueError("反馈 ID 不能重复") + return ids + + +class FeedbackBulkApproveRequest(FeedbackBulkRequest): + reward_coins: int = Field( + ge=1, + le=FEEDBACK_REWARD_MAX_COINS, + description="每条采纳反馈发放的金币数", + ) + note: str | None = Field(default=None, max_length=256, description="采纳要点/审核备注(内部)") + reply: str | None = Field(default=None, max_length=256, description="给用户的回复留言,用户端可见") + + +class FeedbackBulkRejectRequest(FeedbackBulkRequest): + reason: str = Field(min_length=1, max_length=256, description="批量未采纳原因,用户端可见") + note: str | None = Field(default=None, max_length=256, description="运营内部审核备注") + reply: str | None = Field(default=None, max_length=256, description="给用户的回复留言,用户端可见") + + @field_validator("reason") + @classmethod + def _reason_not_blank(cls, value: str) -> str: + if not value.strip(): + raise ValueError("未采纳原因不能为空") + return value.strip() + + +class FeedbackBulkItemResult(BaseModel): + id: int + ok: bool + status: str | None = None + error: str | None = None + + +class FeedbackBulkResult(BaseModel): + total: int + success: int + failed: int + items: list[FeedbackBulkItemResult] + + class FeedbackSummary(BaseModel): """审核台顶部各状态计数(pending 含历史 new 态)。""" diff --git a/app/admin/schemas/guide_video.py b/app/admin/schemas/guide_video.py new file mode 100644 index 0000000..721f0a6 --- /dev/null +++ b/app/admin/schemas/guide_video.py @@ -0,0 +1,25 @@ +"""admin 新手引导视频配置 schemas(开关 / 视频地址 / 前几次 / 每次金币)。""" +from __future__ import annotations + +from pydantic import BaseModel, Field + +from app.repositories.guide_video import MAX_PLAYS_LIMIT, REWARD_COIN_LIMIT + + +class GuideVideoConfigOut(BaseModel): + enabled: bool + video_url: str | None = None # 相对地址 /media/guide_video/xxx.mp4;未配片 = None + max_plays: int + reward_coin: int + updated_at: str | None = None + # 只读统计,后台展示用:已有多少次播放、其中已发币多少次。 + total_plays: int = 0 + granted_plays: int = 0 + + +class GuideVideoConfigUpdate(BaseModel): + """部分更新:只改传入(非 None)字段。视频文件走 /video 上传接口。""" + + enabled: bool | None = None + max_plays: int | None = Field(default=None, ge=0, le=MAX_PLAYS_LIMIT) + reward_coin: int | None = Field(default=None, ge=0, le=REWARD_COIN_LIMIT) diff --git a/app/admin/schemas/price_report.py b/app/admin/schemas/price_report.py index abb89d9..d67208a 100644 --- a/app/admin/schemas/price_report.py +++ b/app/admin/schemas/price_report.py @@ -56,6 +56,42 @@ class PriceReportRejectRequest(BaseModel): return v.strip() +class PriceReportBulkRequest(BaseModel): + ids: list[int] = Field(min_length=1, max_length=50, description="待审核上报 ID 列表") + + @field_validator("ids") + @classmethod + def _ids_must_be_unique(cls, ids: list[int]) -> list[int]: + if len(ids) != len(set(ids)): + raise ValueError("上报 ID 不能重复") + return ids + + +class PriceReportBulkRejectRequest(PriceReportBulkRequest): + reason: str = Field(min_length=1, max_length=256, description="批量拒绝理由,用户端记录页会看到") + + @field_validator("reason") + @classmethod + def _reason_not_blank(cls, value: str) -> str: + if not value.strip(): + raise ValueError("拒绝理由不能为空") + return value.strip() + + +class PriceReportBulkItemResult(BaseModel): + id: int + ok: bool + status: str | None = None + error: str | None = None + + +class PriceReportBulkResult(BaseModel): + total: int + success: int + failed: int + items: list[PriceReportBulkItemResult] + + class PriceReportSummary(BaseModel): """审核台顶部各状态计数。""" diff --git a/app/api/v1/compare_record.py b/app/api/v1/compare_record.py index f1d67d6..5afc8ac 100644 --- a/app/api/v1/compare_record.py +++ b/app/api/v1/compare_record.py @@ -20,6 +20,8 @@ from app.db.session import SessionLocal from app.models.comparison import ComparisonRecord from app.repositories import comparison as crud_compare from app.schemas.compare_record import ( + CompareStartReserveIn, + CompareStartReserveOut, CompareStatsOut, ComparisonRecordCreatedOut, ComparisonRecordDetailOut, @@ -35,6 +37,41 @@ logger = logging.getLogger("shagua.compare_record") router = APIRouter(prefix="/api/v1/compare", tags=["compare-record"]) +@router.post( + "/start", + response_model=CompareStartReserveOut, + summary="预占一次当日比价发起次数(每人每天最多100次)", +) +def reserve_compare_start( + payload: CompareStartReserveIn, + user: CurrentUser, + db: DbSession, +) -> CompareStartReserveOut: + try: + _, used = crud_compare.reserve_daily_start( + db, + user_id=user.id, + trace_id=payload.trace_id, + business_type=payload.business_type, + device_id=payload.device_id, + ) + except crud_compare.DailyCompareStartLimitExceeded: + raise HTTPException( + status_code=status.HTTP_429_TOO_MANY_REQUESTS, + detail="今日已比价超过100次,请明天再试", + ) from None + except crud_compare.ComparisonTraceOwnershipError: + raise HTTPException( + status_code=status.HTTP_409_CONFLICT, + detail="比价任务标识冲突,请重新发起", + ) from None + return CompareStartReserveOut( + limit=crud_compare.DAILY_COMPARE_START_LIMIT, + used=used, + remaining=max(crud_compare.DAILY_COMPARE_START_LIMIT - used, 0), + ) + + @router.post( "/record", response_model=ComparisonRecordCreatedOut, diff --git a/app/api/v1/coupon.py b/app/api/v1/coupon.py index bd2127a..ecd40a8 100644 --- a/app/api/v1/coupon.py +++ b/app/api/v1/coupon.py @@ -81,7 +81,7 @@ def _record_claims_blocking( device_id: str, user_id: int | None, trace_id: str | None, results: list[dict] ) -> None: with SessionLocal() as db: - # 取本次 session 环境,给 coupon_claim_record 打 app_env 标(每券成功率表按它过滤;设计 §13)。 + # 取本次 session 环境,给每日资产和逐次事件同时打环境标。 app_env = coupon_repo.session_app_env(db, trace_id) coupon_repo.record_claims(db, device_id, user_id, trace_id, results, app_env=app_env) # 顺带把本帧「成功平台」并入 coupon_session.platform_success(admin 领券数据 ②整单/③点位成功率; @@ -176,7 +176,7 @@ async def coupon_step( resp_json = resp.json() - # 领券结果沉淀(资产):中间帧 last_coupon_result + done 帧 coupon_results 幂等写库。 + # 领券结果沉淀:每日资产 + 逐次事件;中间帧和 done 全量帧均幂等写库。 # 当前只记录、不参与"要不要领"判断(MVP 先不去重)。写库失败不影响返回。 if device_id: results = _extract_coupon_results(resp_json) diff --git a/app/api/v1/guide_video.py b/app/api/v1/guide_video.py new file mode 100644 index 0000000..c38251b --- /dev/null +++ b/app/api/v1/guide_video.py @@ -0,0 +1,64 @@ +"""新手引导视频(领券等候浮层前 N 次替代广告)。 + +路由前缀 `/api/v1/guide-video`(均需 Bearer): + POST /start 这次浮层放引导视频还是放广告?命中则**当场计次**并下发 play_token + POST /reward 播完 / 中途关闭都调,按 play_token 幂等发固定金币 + +发币额度以**服务端配置**为准(运营后台可改),客户端只报"播完/关闭",报不了金额, +所以被破解也刷不到超额金币;次数上限由 guide_video_play 行数(按账号)硬卡。 +""" +from __future__ import annotations + +import logging + +from fastapi import APIRouter, Depends + +from app.api.deps import CurrentUser, DbSession +from app.core.ratelimit import rate_limit +from app.repositories import guide_video as crud_guide +from app.schemas.guide_video import ( + GuideVideoRewardIn, + GuideVideoRewardOut, + GuideVideoStartIn, + GuideVideoStartOut, +) + +logger = logging.getLogger("shagua.guide_video") + +router = APIRouter(prefix="/api/v1/guide-video", tags=["guide-video"]) + + +@router.post( + "/start", + response_model=GuideVideoStartOut, + summary="领券浮层是否放新手引导视频(命中即计次)", + dependencies=[Depends(rate_limit(60, 60, "guide-video-start"))], +) +def start(payload: GuideVideoStartIn, user: CurrentUser, db: DbSession) -> GuideVideoStartOut: + """开播即计数:返回 should_play=True 时服务端已写下这一次,客户端必须真的播。 + + 没配视频 / 开关关 / 次数用完 → should_play=False,客户端照旧走广告链路(行为不变)。 + """ + result = crud_guide.start_play(db, user.id, scene=payload.scene or "coupon") + logger.info( + "guide video start user_id=%d scene=%s should_play=%s seq=%d remaining=%d", + user.id, payload.scene, result["should_play"], result["seq"], result["remaining"], + ) + return GuideVideoStartOut(**result) + + +@router.post( + "/reward", + response_model=GuideVideoRewardOut, + summary="引导视频发金币(播完/中途关闭都发,play_token 幂等)", + dependencies=[Depends(rate_limit(60, 60, "guide-video-reward"))], +) +def reward(payload: GuideVideoRewardIn, user: CurrentUser, db: DbSession) -> GuideVideoRewardOut: + result = crud_guide.grant_play( + db, user.id, play_token=payload.play_token, completed=payload.completed + ) + logger.info( + "guide video reward user_id=%d token=%s completed=%s granted=%s coin=%d", + user.id, payload.play_token[:12], payload.completed, result["granted"], result["coin"], + ) + return GuideVideoRewardOut(**result) diff --git a/app/api/v1/meituan.py b/app/api/v1/meituan.py index c1d07a1..b3cf1c6 100644 --- a/app/api/v1/meituan.py +++ b/app/api/v1/meituan.py @@ -5,11 +5,12 @@ from __future__ import annotations import logging -from concurrent.futures import ThreadPoolExecutor, as_completed +from concurrent.futures import ThreadPoolExecutor +from typing import TYPE_CHECKING, Any from fastapi import APIRouter, Depends, HTTPException from sqlalchemy import nullslast, select -from sqlalchemy.orm import Session, aliased +from sqlalchemy.orm import Session from app.core.config import settings from app.db.session import get_db @@ -25,8 +26,14 @@ from app.schemas.meituan import ( ReferralLinkResponse, TopSalesRequest, ) +from app.utils import mt_search_cursor from app.utils.meituan_city import get_meituan_city +if TYPE_CHECKING: # 仅供类型标注(本模块已开 from __future__ import annotations) + from collections.abc import Callable + + from sqlalchemy import ColumnElement + logger = logging.getLogger("shagua.meituan") @@ -109,6 +116,86 @@ def _commission_pct(card: CouponCard) -> float: return 0.0 +# ────────────── 离线库分页(智能推荐 / 销量最高 共用) ────────────── +# 去重+排序阶段**只投影这几列**:够 DISTINCT ON 分组、够排序、够回表定位,且全是定长小字段。 +# ⚠️ 关键性能点:`raw` 是整条美团原始返回(JSONB,每行数 KB)。原实现用 select(MeituanCoupon) +# 做子查询,等于把整城几千行连 raw 一起塞进两次排序(DISTINCT ON 一次 + 分页一次), +# 体量轻松超过 work_mem → Postgres 落盘做外部归并排序,而且**每翻一页都要重来一遍**。 +# 拆成「先在小列上排出本页 id,再按 id 回表取 raw」后,排序数据量降到原来的百分之几, +# JSONB 只解析当前页 ~20 行。 +_DEDUP_COLS = ( + MeituanCoupon.id, + MeituanCoupon.dedup_key, + MeituanCoupon.sale_volume_num, + MeituanCoupon.commission_percent, +) + + +def _paged_dedup_ids( + db: Session, + *, + conds: list[ColumnElement[bool]], + dedup_order: list[ColumnElement], + page_order: Callable[[Any], list[ColumnElement]], + page: int, + page_size: int, +) -> tuple[list[int], bool]: + """DISTINCT ON(dedup_key) 跨源去重 → 整体排序 → 分页,返回 (本页 id 列表, 是否还有下一页)。 + + - `dedup_order`:同一个 dedup_key 的多条里留哪条(如销量最高/佣金最高)。 + - `page_order`:接收去重子查询的列集合(`sub.c`),返回去重后的整体排序。 + 多取 1 条用于判断 has_next。 + """ + deduped = ( + select(*_DEDUP_COLS) + .where(*conds) + .distinct(MeituanCoupon.dedup_key) + .order_by(MeituanCoupon.dedup_key, *dedup_order) + .subquery() + ) + ids = db.execute( + select(deduped.c.id) + .order_by(*page_order(deduped.c)) + .offset((page - 1) * page_size) + .limit(page_size + 1) + ).scalars().all() + return list(ids[:page_size]), len(ids) > page_size + + +def _load_raws(db: Session, ids: list[int]) -> list[dict]: + """按给定 id 顺序取 raw(只回表本页 ~20 行)。缺行(被 ETL 清掉)静默跳过。""" + if not ids: + return [] + raw_by_id = { + row_id: raw + for row_id, raw in db.execute( + select(MeituanCoupon.id, MeituanCoupon.raw).where(MeituanCoupon.id.in_(ids)) + ).all() + } + return [raw_by_id[i] for i in ids if i in raw_by_id] + + +def _cards_from_raws(raws: list[dict], *, hide_distance: bool) -> list[CouponCard]: + """raw → CouponCard;解析失败的单条跳过,不整页失败。 + + hide_distance:离线库里的距离是相对「城市默认点」算的,对用户无意义且误导 —— 智能推荐 / + 销量最高两个 tab 一律置空,前端「距离 店名」那行只剩店名、自动顶到最左。 + """ + cards: list[CouponCard] = [] + for raw in raws: + try: + card = CouponCard.from_raw(raw or {}) + except Exception: # noqa: BLE001 + continue + if not card.product_view_sign: + continue + if hide_distance: + card.distance_text = None + card.distance_meters = None + cards.append(card) + return cards + + @router.post("/feed", response_model=FeedResponse, summary="混合feed(外卖+到店交叉);tab=rec智能推荐/distance距离最近") def feed(req: FeedRequest, db: Session = Depends(get_db)) -> FeedResponse: lon, lat = req.longitude, req.latitude @@ -133,17 +220,21 @@ def feed(req: FeedRequest, db: Session = Depends(get_db)) -> FeedResponse: return [], True # 距离最近:搜索召回(外卖搜"外卖" + 到店搜"美食",都 sortField=6 离我最近)一页页拉。 - # 搜索翻页必须用 searchId(pageNo 翻不动),所以每个 feed 页顺序翻到第 N 页;两路并行、page 1 最快。 - # 无状态、不改 APP(传页码即可);按你位置实时算距离(库里没存 POI 经纬度,只能实时)。 + # 搜索翻页必须用 searchId(pageNo 翻不动),而接口是无状态的(客户端只传页码)—— 原实现因此 + # 每次都从第 1 页顺序重放到第 N 页,取第 N 页要向美团发 N 次请求,越往下滑越慢。 + # 现在把沿途 searchId 记进 [mt_search_cursor],稳态下每翻一页恒定 1 次请求;两路仍并行。 + # 按你位置实时算距离(库里没存 POI 经纬度,只能实时)。 if tab == "distance": lon_i, lat_i = int(lon * 1_000_000), int(lat * 1_000_000) - def _search_page_n(platform: int, biz_line: int | None, keyword: str, n: int) -> tuple[list[dict], bool, bool]: - """顺序翻到第 n 页(搜索须 searchId 续页),返回(第 n 页 items, 是否还有下一页, 是否调用失败)。""" - sid: str | None = None + def _replay( + platform: int, biz_line: int | None, keyword: str, + key: mt_search_cursor.RouteKey, start: int, sid: str | None, n: int, + ) -> tuple[list[dict], bool, bool]: + """从第 start 页(用 sid 取)顺序翻到第 n 页。start==1 时 sid 应为 None(走 pageNo=1)。""" data: list[dict] = [] has_next = False - for pg in range(1, n + 1): + for pg in range(start, n + 1): body: dict = { "platform": platform, "searchText": keyword, "sortField": 6, "longitude": lon_i, "latitude": lat_i, "pageSize": 20, @@ -161,10 +252,27 @@ def feed(req: FeedRequest, db: Session = Depends(get_db)) -> FeedResponse: data = r.get("data") or [] sid = r.get("searchId") has_next = bool(r.get("hasNext")) and bool(data) + # 记下「下一页要用哪个 searchId」;没有下一页就别记,免得存进死游标。 + if sid and has_next: + mt_search_cursor.remember(key, pg + 1, sid) if not data or (not has_next and pg < n): return [], False, False # 没那么多页了(非错误) return data, has_next, False + def _search_page_n(platform: int, biz_line: int | None, keyword: str, n: int) -> tuple[list[dict], bool, bool]: + """取第 n 页,返回(第 n 页 items, 是否还有下一页, 是否调用失败)。 + + 优先用缓存游标一发直达;缓存未命中/过期才从最近的已知页往后重放,并把沿途游标补进缓存。 + """ + key = mt_search_cursor.route_key(lat, lon, platform, keyword) + start, sid = mt_search_cursor.lookup(key, n) + data, has_next, failed = _replay(platform, biz_line, keyword, key, start, sid, n) + # 用缓存游标却打不通,多半是上游 searchId 过期:作废整条路线,回到第 1 页重放一次。 + if failed and start > 1: + mt_search_cursor.drop(key) + data, has_next, failed = _replay(platform, biz_line, keyword, key, 1, None, n) + return data, has_next, failed + with ThreadPoolExecutor(max_workers=2) as pool: f_wm = pool.submit(_search_page_n, 1, None, "外卖", req.page) f_dd = pool.submit(_search_page_n, 2, 1, "美食", req.page) @@ -194,39 +302,25 @@ def feed(req: FeedRequest, db: Session = Depends(get_db)) -> FeedResponse: return FeedResponse(items=[], has_next=False, page=req.page, status="degraded") PAGE = 20 try: - base = select(MeituanCoupon).where( - MeituanCoupon.commission_percent >= 3.0, - MeituanCoupon.city_id == city_id, - ) - deduped = base.distinct(MeituanCoupon.dedup_key).order_by( - MeituanCoupon.dedup_key, - MeituanCoupon.commission_percent.desc(), - ).subquery() - m = aliased(MeituanCoupon, deduped) - start = (req.page - 1) * PAGE - rows = db.execute( - select(m) + ids, has_next = _paged_dedup_ids( + db, + conds=[ + MeituanCoupon.commission_percent >= 3.0, + MeituanCoupon.city_id == city_id, + ], + # 同一去重键留佣金最高那条 + dedup_order=[MeituanCoupon.commission_percent.desc()], # 销量高的优先(无销量档排后),同档佣金高优先,id 兜底稳定分页 - .order_by(nullslast(m.sale_volume_num.desc()), m.commission_percent.desc(), m.id) - .offset(start) - .limit(PAGE + 1) - ).scalars().all() + page_order=lambda c: [ + nullslast(c.sale_volume_num.desc()), c.commission_percent.desc(), c.id, + ], + page=req.page, page_size=PAGE, + ) + raws = _load_raws(db, ids) except Exception: # noqa: BLE001 logger.exception("[feed] rec 库查询失败,降级返空") return FeedResponse(items=[], has_next=False, page=req.page, status="degraded") - has_next = len(rows) > PAGE - cards: list[CouponCard] = [] - for row in rows[:PAGE]: - try: - card = CouponCard.from_raw(row.raw or {}) - except Exception: # noqa: BLE001 - continue - if card.product_view_sign: - # 智能推荐不显示距离:库里的距离是相对城市默认点的(对用户无意义、且误导)。 - # 置空后前端"距离 店名"那行只剩店名、自动顶到最左(店名移到原距离的位置)。 - card.distance_text = None - card.distance_meters = None - cards.append(card) + cards = _cards_from_raws(raws, hide_distance=True) if not cards and req.page == 1: # 命中城市却 0 券:该城确无 ≥3% 券,或 ETL 灌的 city_id 与 city_dict 口径不一致。 logger.info("[feed] rec city_id=%s 命中 0 券(该城确无券?或 ETL/city_dict 的 city_id 口径不一致)", city_id) @@ -282,51 +376,36 @@ def top_sales(req: TopSalesRequest, db: Session = Depends(get_db)) -> CouponList if not city_id: return CouponListResponse(items=[], has_next=False, search_id=None, status="degraded") - # 去重 + 排序 + 分页全在 SQL 做,每页只取并解析当前页 ~20 条。 - # (之前实现每翻一页都全表拉取 + 全量 from_raw 解析,翻页慢 → 客户端滑动卡顿/翻不动。) + # 去重 + 排序 + 分页全在 SQL 做,每页只回表并解析当前页 ~20 条(见 _paged_dedup_ids 的性能说明)。 # 库为空(prod 刚部署 / ETL 未跑完)时返空 + status=empty,不崩;库查询异常降级 degraded。 + conds = [ + MeituanCoupon.sale_volume_num.isnot(None), + MeituanCoupon.city_id == city_id, + ] + if req.platform is not None: + conds.append(MeituanCoupon.platform == req.platform) try: - # 1) DISTINCT ON (dedup_key):每个去重键(品牌|名|价)只留销量最高那条(同销量再按佣金) - base = select(MeituanCoupon).where( - MeituanCoupon.sale_volume_num.isnot(None), - MeituanCoupon.city_id == city_id, - ) - if req.platform is not None: - base = base.where(MeituanCoupon.platform == req.platform) - deduped = base.distinct(MeituanCoupon.dedup_key).order_by( - MeituanCoupon.dedup_key, - MeituanCoupon.sale_volume_num.desc(), - MeituanCoupon.commission_percent.desc(), - ).subquery() - - # 2) 对去重结果按销量降序分页;多取 1 条判断 has_next,只对本页做 from_raw - m = aliased(MeituanCoupon, deduped) - start = (req.page - 1) * req.page_size - rows = db.execute( - select(m) + ids, has_next = _paged_dedup_ids( + db, + conds=conds, + # 每个去重键(品牌|名|价)只留销量最高那条(同销量再按佣金) + dedup_order=[ + MeituanCoupon.sale_volume_num.desc(), + MeituanCoupon.commission_percent.desc(), + ], # 加 id 作稳定 tiebreaker:同销量同佣金的并列项排序确定,避免跨页重复/漏项 - .order_by(m.sale_volume_num.desc(), m.commission_percent.desc(), m.id) - .offset(start) - .limit(req.page_size + 1) - ).scalars().all() + page_order=lambda c: [ + c.sale_volume_num.desc(), c.commission_percent.desc(), c.id, + ], + page=req.page, page_size=req.page_size, + ) + raws = _load_raws(db, ids) except Exception: # noqa: BLE001 logger.exception("[top-sales] 库查询失败,降级返空") return CouponListResponse(items=[], has_next=False, search_id=None, status="degraded") - has_next = len(rows) > req.page_size - cards: list[CouponCard] = [] - for row in rows[:req.page_size]: - try: - card = CouponCard.from_raw(row.raw or {}) - except Exception: # noqa: BLE001 - continue - if card.product_view_sign: - # 不显示距离:库里的距离是相对城市默认点的(对用户无意义、且误导)。 - # 置空后前端"距离 店名"那行只剩店名、自动顶到最左(店名移到原距离的位置)。 - # 逻辑与推荐流保持一致 - card.distance_text = None - card.distance_meters = None - cards.append(card) + # 不显示距离:库里的距离是相对城市默认点的(对用户无意义、且误导),与推荐流口径一致。 + cards = _cards_from_raws(raws, hide_distance=True) if not cards and req.page == 1: # 命中城市却 0 券:可能该城确无券,也可能 ETL 灌的 city_id 与 city_dict 口径不一致(静默降级的隐患)。 logger.info("[top-sales] city_id=%s 命中 0 券(该城确无券?或 ETL/city_dict 的 city_id 口径不一致)", city_id) diff --git a/app/api/v1/wallet.py b/app/api/v1/wallet.py index 454ff7c..7b0dac0 100644 --- a/app/api/v1/wallet.py +++ b/app/api/v1/wallet.py @@ -192,7 +192,7 @@ def withdraw_info( wechat_bound=bool(u and u.wechat_openid), wechat_nickname=u.wechat_nickname if u else None, wechat_avatar_url=u.wechat_avatar_url if u else None, - transfer_auth_enabled=bool(auth and auth.state == "active"), + transfer_auth_enabled=bool(auth and auth.state == "active" and auth.authorization_id), tiers=[WithdrawTierOut(**t) for t in crud_wallet.withdraw_tier_states(db, user.id, source)], ) @@ -335,7 +335,8 @@ def open_transfer_auth(user: CurrentUser, db: DbSession) -> TransferAuthResultOu def transfer_auth_status(user: CurrentUser, db: DbSession) -> TransferAuthStatusOut: auth = crud_wallet.sync_transfer_auth(db, user.id) state = auth.state if auth else "none" - return TransferAuthStatusOut(state=state, enabled=(state == "active")) + enabled = bool(auth and auth.state == "active" and auth.authorization_id) + return TransferAuthStatusOut(state=state, enabled=enabled) @router.post( diff --git a/app/core/config.py b/app/core/config.py index afae3ee..81149db 100644 --- a/app/core/config.py +++ b/app/core/config.py @@ -384,6 +384,9 @@ class Settings(BaseSettings): MEDIA_ROOT: str = "./data/media" MEDIA_URL_PREFIX: str = "/media" AVATAR_MAX_BYTES: int = 5 * 1024 * 1024 # 头像最大 5MB + # 运营后台上传的新手引导视频上限。视频比图片大一个量级,单独一档; + # ⚠️ 改大时同步放宽网关 client_max_body_size(实测 QA 4MiB / prod 32MiB),否则 nginx 先挡下。 + GUIDE_VIDEO_MAX_BYTES: int = 100 * 1024 * 1024 # 引导视频最大 100MB # ===== 邀请好友 ===== # 分享落地页(二维码 / 分享链接指向这里;扫码 → 落地页 → 引导浏览器下载 APK)。 diff --git a/app/core/logging.py b/app/core/logging.py index cb91b5c..c3833bd 100644 --- a/app/core/logging.py +++ b/app/core/logging.py @@ -22,13 +22,13 @@ from __future__ import annotations import json import logging import os +import shutil import sys from contextvars import ContextVar from datetime import datetime from logging.handlers import RotatingFileHandler from pathlib import Path - # 请求级 trace_id:入口(如 compare.py 透传壳)set 之后, 本请求上下文(含 run_in_threadpool # 拷贝出去的线程)内所有日志自动带上。默认空串 = 非请求上下文(启动/后台 worker)。 trace_id_ctx: ContextVar[str] = ContextVar("trace_id", default="") @@ -91,6 +91,55 @@ class TextFormatter(logging.Formatter): return f"{base} trace={tid}" if tid else base +class SafeRotatingFileHandler(RotatingFileHandler): + """Windows 下不会被外部句柄卡死的 RotatingFileHandler。 + + stdlib 轮转靠 rename 活动文件(app-server.log → .1);Windows 只要有别的句柄(IDE 索引、 + app.admin.main 第二进程、残留 --reload worker、杀软扫描)开着它, rename 就 WinError 32, + 轮转永久卡死——文件停在 maxBytes、之后每条日志被丢。这里 Windows 改用 copytruncate:把活动 + 文件拷进备份、再通过自己的句柄原地清空, 从不 rename 活动文件, 故外部句柄开着也能转。 + POSIX(生产 Linux)rename 打开中的文件本就合法, 保留 stdlib 的原子轮转不变。 + + 代价:copytruncate 在“拷贝→清空”极窄窗口内并发写可能丢几行(仅跨进程;同进程 emit 有 + handler 锁串行, 无此问题)。对本地开发日志可接受。 + """ + + def doRollover(self) -> None: + if os.name != "nt": + super().doRollover() + return + if self.stream is None: + self.stream = self._open() + else: + self.stream.flush() + try: + self._copytruncate_backups() + except OSError: + # 备份腾挪是尽力而为:任一备份被占用也绝不能挡住下面的清空, 否则活动文件继续涨、 + # 轮转又卡死——那就白改了。 + pass + # 通过自己独占的句柄原地清空:不涉及 rename, 外部只读句柄不受影响。 + self.stream.seek(0) + self.stream.truncate() + self.stream.flush() + + def _copytruncate_backups(self) -> None: + """把 .N-1→.N 逐级腾挪, 再把活动文件拷到 .1(不动活动文件本身)。""" + if self.backupCount <= 0: + return + for i in range(self.backupCount - 1, 0, -1): + sfn = self.rotation_filename(f"{self.baseFilename}.{i}") + dfn = self.rotation_filename(f"{self.baseFilename}.{i + 1}") + if os.path.exists(sfn): + if os.path.exists(dfn): + os.remove(dfn) + os.replace(sfn, dfn) + dfn = self.rotation_filename(f"{self.baseFilename}.1") + if os.path.exists(dfn): + os.remove(dfn) + shutil.copyfile(self.baseFilename, dfn) + + _CONFIGURED = False @@ -126,7 +175,7 @@ def setup_logging(debug: bool = False) -> None: Path(os.getenv("LOG_DIR", "logs")) / "app-server.log" ) Path(log_file).parent.mkdir(parents=True, exist_ok=True) - file_handler = RotatingFileHandler( + file_handler = SafeRotatingFileHandler( log_file, maxBytes=10 * 1024 * 1024, backupCount=5, encoding="utf-8", ) file_handler.setFormatter(JsonFormatter(service)) diff --git a/app/core/media.py b/app/core/media.py index 192b186..a6e84e4 100644 --- a/app/core/media.py +++ b/app/core/media.py @@ -77,6 +77,35 @@ def save_feedback_qr(data: bytes) -> str: return _save_named("feedback_qr", "qr", data) +def _sniff_video_ext(data: bytes) -> str | None: + """按魔数判定视频类型,返回扩展名;非支持类型返回 None。 + + 只认 MP4 家族(ISO BMFF):`....ftyp` 在偏移 4。Android ExoPlayer 与浏览器