From cdd49c64211eeb46e69e97bab0fa3bb5780c00fa Mon Sep 17 00:00:00 2001 From: linkeyu Date: Mon, 27 Jul 2026 10:24:32 +0800 Subject: [PATCH] =?UTF-8?q?=E4=BF=AE=E5=A4=8D=EF=BC=9A=E8=A1=A5=E9=BD=90?= =?UTF-8?q?=E6=94=B6=E7=9B=8A=E6=98=8E=E7=BB=86=E5=B9=BF=E5=91=8A=E7=BD=91?= =?UTF-8?q?=E7=BB=9C=E6=9D=A5=E6=BA=90=20(#179)?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit ## 本次改动 - 激励视频发奖记录按用户与 `ad_session_id` 回填展示侧 ADN - 信息流明细保留每条发奖自身上报的 ADN 与底层代码位 - 历史记录仅在 `用户 + trace_id + eCPM` 候选网络唯一时安全回填 - 多网络候选保持空值,避免错误归因 ## 验证 专项测试 4 passed,覆盖不同 ADN 明细及唯一/歧义回填场景。 --------- Co-authored-by: guke Co-authored-by: unknown <798648091@qq.com> Reviewed-on: https://gitea.shaguabijia.com/WonderableAI/shaguabijia-app-server/pulls/179 Co-authored-by: linkeyu Co-committed-by: linkeyu --- app/admin/repositories/ad_audit.py | 104 ++++++++++++++++++++++++- app/admin/repositories/ad_revenue.py | 10 ++- app/admin/schemas/ad_revenue.py | 2 + tests/test_admin_ad_revenue_scope.py | 110 +++++++++++++++++++++++++++ 4 files changed, 220 insertions(+), 6 deletions(-) diff --git a/app/admin/repositories/ad_audit.py b/app/admin/repositories/ad_audit.py index 888467c..0be72ce 100644 --- a/app/admin/repositories/ad_audit.py +++ b/app/admin/repositories/ad_audit.py @@ -12,10 +12,11 @@ """ from __future__ import annotations -from sqlalchemy import func, select +from sqlalchemy import func, or_, select from sqlalchemy.orm import Session from app.core import rewards +from app.models.ad_ecpm import AdEcpmRecord from app.models.ad_feed_reward import AdFeedRewardRecord from app.models.ad_reward import AdRewardRecord from app.repositories.ad_feed_reward import FEED_REWARD_UNIT_SECONDS @@ -55,10 +56,21 @@ def _reward_video_rows( if user_id is not None: stmt = stmt.where(AdRewardRecord.user_id == user_id) + records = list(db.execute(stmt).scalars()) + # S2S 发奖回调不携带实际填充 ADN;用相同用户和 ad_session_id 的展示记录回填。 + session_ids = {record.ad_session_id for record in records if record.ad_session_id} + impression_by_session = { + (record.user_id, record.ad_session_id): record + for record in db.execute( + select(AdEcpmRecord).where(AdEcpmRecord.ad_session_id.in_(session_ids)) + ).scalars() + } if session_ids else {} + # 用本日之前的累计份数做起点,当日 granted 在其上继续递增 → 与 _granted_cumulative+1 对齐 granted_n: dict[int, int] = _prior_granted_counts(db, date=date, user_id=user_id) rows: list[dict] = [] - for rec in db.execute(stmt).scalars(): + for rec in records: + impression = impression_by_session.get((rec.user_id, rec.ad_session_id)) if rec.status == "granted": nth = granted_n.get(rec.user_id, 0) + 1 granted_n[rec.user_id] = nth @@ -68,6 +80,8 @@ def _reward_video_rows( "record_id": rec.id, "user_id": rec.user_id, "ad_session_id": rec.ad_session_id, + "adn": impression.adn if impression is not None else None, + "slot_id": impression.slot_id if impression is not None else None, "app_env": rec.app_env, "our_code_id": rec.our_code_id, "created_at": rec.created_at, @@ -90,6 +104,8 @@ def _reward_video_rows( "record_id": rec.id, "user_id": rec.user_id, "ad_session_id": rec.ad_session_id, + "adn": impression.adn if impression is not None else None, + "slot_id": impression.slot_id if impression is not None else None, "app_env": rec.app_env, "our_code_id": rec.our_code_id, "created_at": rec.created_at, @@ -149,6 +165,78 @@ def _feed_scene_matches(rec: AdFeedRewardRecord, scene: str | None) -> bool: return True +def _nonblank(value: str | None) -> str | None: + value = value.strip() if value else None + return value or None + + +def _unique_ad_source(records: list[AdEcpmRecord]) -> tuple[str | None, str | None]: + """仅在候选展示记录指向唯一 ADN 时回填来源,避免错误归因。""" + adns = {_nonblank(record.adn) for record in records} + adns.discard(None) + if len(adns) != 1: + return None, None + + slots = {_nonblank(record.slot_id) for record in records} + slots.discard(None) + return next(iter(adns)), next(iter(slots)) if len(slots) == 1 else None + + +def _feed_source_fallbacks( + db: Session, records: list[AdFeedRewardRecord] +) -> tuple[ + dict[tuple[int, str], tuple[str | None, str | None]], + dict[tuple[int, str, str], tuple[str | None, str | None]], +]: + """为旧信息流发奖记录构建安全来源索引。""" + session_ids = {record.ad_session_id for record in records if record.ad_session_id} + trace_ids = {record.trace_id for record in records if record.trace_id} + if not session_ids and not trace_ids: + return {}, {} + + filters = [] + if session_ids: + filters.append(AdEcpmRecord.ad_session_id.in_(session_ids)) + if trace_ids: + filters.append(AdEcpmRecord.trace_id.in_(trace_ids)) + impressions = list(db.execute(select(AdEcpmRecord).where(or_(*filters))).scalars()) + + by_session: dict[tuple[int, str], list[AdEcpmRecord]] = {} + by_trace_ecpm: dict[tuple[int, str, str], list[AdEcpmRecord]] = {} + for impression in impressions: + if impression.ad_session_id: + by_session.setdefault((impression.user_id, impression.ad_session_id), []).append(impression) + if impression.trace_id: + by_trace_ecpm.setdefault( + (impression.user_id, impression.trace_id, impression.ecpm_raw), [] + ).append(impression) + + return ( + {key: _unique_ad_source(value) for key, value in by_session.items()}, + {key: _unique_ad_source(value) for key, value in by_trace_ecpm.items()}, + ) + + +def _feed_source( + record: AdFeedRewardRecord, + *, + by_session: dict[tuple[int, str], tuple[str | None, str | None]], + by_trace_ecpm: dict[tuple[int, str, str], tuple[str | None, str | None]], +) -> tuple[str | None, str | None]: + """返回本条发奖广告的来源;无唯一证据时保留原始空值。""" + adn, slot_id = _nonblank(record.adn), _nonblank(record.slot_id) + if adn and slot_id: + return adn, slot_id + + candidate = by_session.get((record.user_id, record.ad_session_id or "")) + if candidate is None and record.trace_id: + candidate = by_trace_ecpm.get((record.user_id, record.trace_id, record.ecpm_raw)) + if candidate is None: + return adn, slot_id + candidate_adn, candidate_slot_id = candidate + return adn or candidate_adn, slot_id or candidate_slot_id + + def _feed_rows( db: Session, *, date: str, user_id: int | None, scene: str | None = None ) -> list[dict]: @@ -167,11 +255,17 @@ def _feed_rows( if user_id is not None: stmt = stmt.where(AdFeedRewardRecord.user_id == user_id) + records = list(db.execute(stmt).scalars()) + by_session, by_trace_ecpm = _feed_source_fallbacks(db, records) + # 本日之前的累计**条数**做起点,与发奖侧 granted_unit_total(COUNT granted)对齐 granted_count: dict[int, int] = _feed_prior_granted_count(db, date=date, user_id=user_id) rows: list[dict] = [] - for rec in db.execute(stmt).scalars(): + for rec in records: keep = _feed_scene_matches(rec, scene) # 累计照常推进,这里只决定是否展示本行 + adn, slot_id = _feed_source( + rec, by_session=by_session, by_trace_ecpm=by_trace_ecpm + ) if rec.status == "granted": # 一条广告 = 1 份(与 grant_feed_reward 同口径:看满一份即发该条满额,不按 unit_count 累加)。 # nth = 账号累计第几**条**(含本日之前),与发奖侧 granted_unit_total+1 对齐;累计照常推进 @@ -188,6 +282,8 @@ def _feed_rows( "record_id": rec.id, "user_id": rec.user_id, "ad_session_id": rec.ad_session_id, + "adn": adn, + "slot_id": slot_id, "trace_id": rec.trace_id, "app_env": rec.app_env, "our_code_id": rec.our_code_id, @@ -214,6 +310,8 @@ def _feed_rows( "record_id": rec.id, "user_id": rec.user_id, "ad_session_id": rec.ad_session_id, + "adn": adn, + "slot_id": slot_id, "trace_id": rec.trace_id, "app_env": rec.app_env, "our_code_id": rec.our_code_id, diff --git a/app/admin/repositories/ad_revenue.py b/app/admin/repositories/ad_revenue.py index 8aef9c6..b8c33ae 100644 --- a/app/admin/repositories/ad_revenue.py +++ b/app/admin/repositories/ad_revenue.py @@ -96,7 +96,11 @@ _REWARD_DETAIL_KEYS = ( def _reward_detail(row: dict) -> dict: """从 audit 行抽出发奖复算明细(给前端展开行渲染因子1/因子2/份数/LT/应发实发)。""" - return {k: row[k] for k in _REWARD_DETAIL_KEYS} + detail = {key: row[key] for key in _REWARD_DETAIL_KEYS} + # 聚合父行可能包含多个 ADN,来源必须保留在每一条发奖明细上。 + detail["adn"] = row.get("adn") + detail["slot_id"] = row.get("slot_id") + return detail def ad_revenue_report( @@ -244,8 +248,8 @@ def ad_revenue_report( "impressions": 0, "ecpm": row["ecpm"], "revenue_yuan": 0.0, - "adn": None, - "slot_id": None, + "adn": row.get("adn"), + "slot_id": row.get("slot_id"), "has_reward": True, "status": row["status"], "expected_coin": int(row["expected_coin"]), diff --git a/app/admin/schemas/ad_revenue.py b/app/admin/schemas/ad_revenue.py index 8143eff..abeee4c 100644 --- a/app/admin/schemas/ad_revenue.py +++ b/app/admin/schemas/ad_revenue.py @@ -40,6 +40,8 @@ class AdRevenueRecord(BaseModel): expected_coin: int = Field(..., description="按公式复算应发金币") actual_coin: int = Field(..., description="实际入账金币") matched: bool = Field(..., description="复算与实发是否一致") + adn: str | None = Field(None, description="本条发奖对应的实际填充 ADN 子渠道") + slot_id: str | None = Field(None, description="本条发奖对应的底层 mediation rit") class AdRevenueDaily(BaseModel): diff --git a/tests/test_admin_ad_revenue_scope.py b/tests/test_admin_ad_revenue_scope.py index f495d66..396af9f 100644 --- a/tests/test_admin_ad_revenue_scope.py +++ b/tests/test_admin_ad_revenue_scope.py @@ -8,6 +8,7 @@ from sqlalchemy import delete from app.admin.repositories import ad_revenue from app.db.session import SessionLocal from app.models.ad_ecpm import AdEcpmRecord +from app.models.ad_feed_reward import AdFeedRewardRecord from app.models.ad_pangle_revenue import AdPangleDailyRevenue from app.models.ad_reward import AdRewardRecord from app.models.user import User @@ -208,3 +209,112 @@ def test_reward_video_incomplete_playback_has_zero_revenue() -> None: db.execute(delete(User).where(User.phone == phone)) db.commit() db.close() + + +def test_feed_reward_details_keep_each_record_adn() -> None: + db = SessionLocal() + phone = "18800009994" + detail_date = "2040-02-06" + try: + user = User(phone=phone, username="29999999994", register_channel="sms") + db.add(user) + db.flush() + db.add_all([ + AdFeedRewardRecord( + client_event_id="detail-adn-pangle", user_id=user.id, + reward_date=detail_date, duration_seconds=20, unit_count=1, + ecpm_raw="12000", adn="pangle", slot_id="rit-pangle", + ad_type="draw", feed_scene="coupon", trace_id="detail-adn-trace", + app_env="prod", our_code_id="104098712", coin=12, status="granted", + created_at=datetime(2040, 2, 6, 1, tzinfo=UTC), + ), + AdFeedRewardRecord( + client_event_id="detail-adn-gdt", user_id=user.id, + reward_date=detail_date, duration_seconds=20, unit_count=1, + ecpm_raw="25000", adn="gdt", slot_id="rit-gdt", + ad_type="draw", feed_scene="coupon", trace_id="detail-adn-trace", + app_env="prod", our_code_id="104098712", coin=25, status="granted", + created_at=datetime(2040, 2, 6, 2, tzinfo=UTC), + ), + ]) + db.commit() + + result = ad_revenue.ad_revenue_report( + db, date_from=detail_date, date_to=detail_date, + user_id=user.id, app_env="prod", revenue_scope="all", + ) + + item = next(row for row in result["items"] if row["event_key"].startswith("feedgrp-")) + assert [detail["adn"] for detail in item["sub_rewards"]] == ["pangle", "gdt"] + assert [detail["slot_id"] for detail in item["sub_rewards"]] == ["rit-pangle", "rit-gdt"] + finally: + db.rollback() + db.execute(delete(AdFeedRewardRecord).where(AdFeedRewardRecord.reward_date == detail_date)) + db.execute(delete(User).where(User.phone == phone)) + db.commit() + db.close() + + +def test_feed_reward_source_fallback_requires_unique_trace_and_ecpm() -> None: + db = SessionLocal() + phone = "18800009995" + fallback_date = "2040-02-07" + try: + user = User(phone=phone, username="29999999995", register_channel="sms") + db.add(user) + db.flush() + db.add_all([ + AdFeedRewardRecord( + client_event_id="source-fallback-unique", user_id=user.id, + reward_date=fallback_date, duration_seconds=3, unit_count=0, + ecpm_raw="4700", ad_session_id="flow-session", trace_id="source-trace", + app_env="prod", our_code_id="104098712", coin=0, status="too_short", + ad_type="draw", feed_scene="comparison", + created_at=datetime(2040, 2, 7, 1, tzinfo=UTC), + ), + AdFeedRewardRecord( + client_event_id="source-fallback-ambiguous", user_id=user.id, + reward_date=fallback_date, duration_seconds=3, unit_count=0, + ecpm_raw="4800", ad_session_id="flow-session", trace_id="source-trace", + app_env="prod", our_code_id="104098712", coin=0, status="too_short", + ad_type="draw", feed_scene="comparison", + created_at=datetime(2040, 2, 7, 2, tzinfo=UTC), + ), + AdEcpmRecord( + user_id=user.id, ad_type="draw", ad_session_id="impression-unique", + trace_id="source-trace", ecpm_raw="4700", adn="baidu", slot_id="rit-baidu", + app_env="prod", our_code_id="104098712", report_date=fallback_date, + created_at=datetime(2040, 2, 7, 1, tzinfo=UTC), + ), + AdEcpmRecord( + user_id=user.id, ad_type="draw", ad_session_id="impression-ambiguous-a", + trace_id="source-trace", ecpm_raw="4800", adn="baidu", slot_id="rit-baidu", + app_env="prod", our_code_id="104098712", report_date=fallback_date, + created_at=datetime(2040, 2, 7, 2, tzinfo=UTC), + ), + AdEcpmRecord( + user_id=user.id, ad_type="draw", ad_session_id="impression-ambiguous-b", + trace_id="source-trace", ecpm_raw="4800", adn="ks", slot_id="rit-ks", + app_env="prod", our_code_id="104098712", report_date=fallback_date, + created_at=datetime(2040, 2, 7, 2, 1, tzinfo=UTC), + ), + ]) + db.commit() + + result = ad_revenue.ad_revenue_report( + db, date_from=fallback_date, date_to=fallback_date, + user_id=user.id, app_env="prod", revenue_scope="all", + ) + item = next(row for row in result["items"] if row["event_key"].startswith("feedgrp-")) + details = {detail["ecpm"]: detail for detail in item["sub_rewards"]} + assert details["4700"]["adn"] == "baidu" + assert details["4700"]["slot_id"] == "rit-baidu" + assert details["4800"]["adn"] is None + assert details["4800"]["slot_id"] is None + finally: + db.rollback() + db.execute(delete(AdFeedRewardRecord).where(AdFeedRewardRecord.reward_date == fallback_date)) + db.execute(delete(AdEcpmRecord).where(AdEcpmRecord.report_date == fallback_date)) + db.execute(delete(User).where(User.phone == phone)) + db.commit() + db.close()