| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422 |
- """Creative material usage reservation for duplicate control.
- Usage is written when a pending creative is prepared, before Tencent POST.
- Review/task tables remain responsible for submitted Tencent creatives.
- """
- from __future__ import annotations
- import json
- import logging
- import os
- from typing import Iterable
- logger = logging.getLogger(__name__)
- CREATIVE_MATERIAL_DEDUPE_LOOKBACK_DAYS = int(
- os.getenv("CREATIVE_MATERIAL_DEDUPE_LOOKBACK_DAYS", "7")
- )
- CREATIVE_LANDING_DEDUPE_LOOKBACK_DAYS = int(
- os.getenv("CREATIVE_LANDING_DEDUPE_LOOKBACK_DAYS", "7")
- )
- CREATION_RECOVERY_LOOKBACK_HOURS = int(
- os.getenv("CREATION_RECOVERY_LOOKBACK_HOURS", "24")
- )
- CREATE_USAGE_TABLE_SQL = """
- CREATE TABLE IF NOT EXISTS creative_material_usage (
- id BIGINT AUTO_INCREMENT PRIMARY KEY COMMENT '自增主键',
- account_id BIGINT NOT NULL COMMENT '腾讯广告账户ID',
- adgroup_id BIGINT DEFAULT NULL COMMENT '广告ID',
- crowd_package VARCHAR(200) NOT NULL COMMENT '投放人群包名称',
- landing_video_id BIGINT DEFAULT NULL COMMENT '承接视频ID,用于同人群包近期排重',
- material_id VARCHAR(100) NOT NULL COMMENT '内部素材ID',
- material_image_id VARCHAR(100) DEFAULT NULL COMMENT '腾讯图片ID',
- dynamic_creative_id BIGINT DEFAULT NULL COMMENT '腾讯动态创意ID',
- status VARCHAR(50) NOT NULL DEFAULT 'prepared' COMMENT 'prepared/submitted/failed',
- source VARCHAR(50) DEFAULT NULL COMMENT 'primary/hot',
- raw_record MEDIUMTEXT DEFAULT NULL COMMENT '准备记录JSON',
- created_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP COMMENT '创建时间',
- updated_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP COMMENT '更新时间',
- KEY idx_crowd_material_created (crowd_package, material_id, created_at),
- KEY idx_crowd_landing_created (crowd_package, landing_video_id, created_at),
- KEY idx_crowd_account_ad_material (crowd_package, account_id, adgroup_id, material_id),
- KEY idx_account_created (account_id, created_at),
- KEY idx_status_created (status, created_at)
- ) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COMMENT='创意素材使用历史/占位'
- """
- def ensure_usage_table() -> None:
- """确保 creative_material_usage 表存在(不存在则创建)。"""
- from db.connection import get_connection
- conn = get_connection()
- try:
- with conn.cursor() as cur:
- cur.execute(CREATE_USAGE_TABLE_SQL)
- cur.execute("SHOW INDEX FROM creative_material_usage WHERE Key_name = 'idx_crowd_landing_created'")
- if not cur.fetchall():
- cur.execute(
- """
- ALTER TABLE creative_material_usage
- ADD INDEX idx_crowd_landing_created (crowd_package, landing_video_id, created_at)
- """
- )
- conn.commit()
- finally:
- conn.close()
- def load_recent_used_material_ids(
- crowd_package: str,
- lookback_days: int = CREATIVE_MATERIAL_DEDUPE_LOOKBACK_DAYS,
- ) -> set[str]:
- """Load recent material reservations for the same crowd package."""
- if not crowd_package:
- return set()
- ensure_usage_table()
- from db.connection import get_connection
- conn = get_connection()
- try:
- with conn.cursor() as cur:
- cur.execute(
- """
- SELECT DISTINCT material_id
- FROM (
- SELECT material_id
- FROM creative_material_usage
- WHERE crowd_package=%s
- AND material_id IS NOT NULL
- AND material_id <> ''
- AND status IN ('posted_ok', 'rejected')
- AND created_at >= DATE_SUB(NOW(), INTERVAL %s DAY)
- UNION
- SELECT material_id
- FROM creative_creation_task
- WHERE material_id IS NOT NULL
- AND material_id <> ''
- AND review_status IN ('approved', 'rejected')
- AND submitted_at >= DATE_SUB(NOW(), INTERVAL %s DAY)
- AND JSON_UNQUOTE(JSON_EXTRACT(raw_record, '$.audience_tier'))=%s
- ) t
- """,
- (crowd_package, int(lookback_days), int(lookback_days), crowd_package),
- )
- rows = cur.fetchall() or []
- finally:
- conn.close()
- return {str(row["material_id"]) for row in rows if row.get("material_id")}
- def load_recent_used_landing_video_ids(
- crowd_package: str,
- lookback_days: int = CREATIVE_LANDING_DEDUPE_LOOKBACK_DAYS,
- ) -> set[int]:
- """Load recent landing-video reservations for the same crowd package."""
- if not crowd_package:
- return set()
- ensure_usage_table()
- from db.connection import get_connection
- conn = get_connection()
- try:
- with conn.cursor() as cur:
- cur.execute(
- """
- SELECT DISTINCT landing_video_id
- FROM (
- SELECT landing_video_id
- FROM creative_material_usage
- WHERE crowd_package=%s
- AND landing_video_id IS NOT NULL
- AND landing_video_id > 0
- AND status IN ('posted_ok', 'submitted')
- AND created_at >= DATE_SUB(NOW(), INTERVAL %s DAY)
- UNION
- SELECT landing_video_id
- FROM creative_creation_task
- WHERE landing_video_id IS NOT NULL
- AND landing_video_id > 0
- AND submitted_at >= DATE_SUB(NOW(), INTERVAL %s DAY)
- AND JSON_UNQUOTE(JSON_EXTRACT(raw_record, '$.audience_tier'))=%s
- ) t
- """,
- (crowd_package, int(lookback_days), int(lookback_days), crowd_package),
- )
- rows = cur.fetchall() or []
- finally:
- conn.close()
- out: set[int] = set()
- for row in rows:
- try:
- out.add(int(row["landing_video_id"]))
- except (TypeError, ValueError):
- continue
- return out
- def _material_source_from_record(material_id: object, raw_record: object) -> str:
- material_id_text = str(material_id or "").strip()
- if material_id_text.startswith("external:"):
- return "external_recall"
- if material_id_text.startswith("ai:"):
- return "ai_generated"
- if raw_record:
- try:
- data = json.loads(str(raw_record))
- source = str(data.get("material_source") or "").strip()
- raw_material_id = str(data.get("_material_id") or "").strip()
- if source == "external_recall" or raw_material_id.startswith("external:"):
- return "external_recall"
- if source == "ai_generated" or raw_material_id.startswith("ai:"):
- return "ai_generated"
- except Exception:
- pass
- return "history"
- def load_recent_landing_usage_counts(
- crowd_package: str,
- lookback_days: int = CREATIVE_LANDING_DEDUPE_LOOKBACK_DAYS,
- ) -> dict[str, dict[int, int]]:
- """Load recent landing-video usage counts split by material source.
- The `source` column in creative_material_usage means primary/hot video source,
- so material source must be derived from material_id/raw_record.
- """
- counts: dict[str, dict[int, int]] = {
- "history": {},
- "ai_generated": {},
- "external_recall": {},
- }
- if not crowd_package:
- return counts
- ensure_usage_table()
- from db.connection import get_connection
- conn = get_connection()
- try:
- with conn.cursor() as cur:
- cur.execute(
- """
- SELECT account_id, adgroup_id, landing_video_id, material_id, raw_record
- FROM creative_material_usage
- WHERE crowd_package=%s
- AND landing_video_id IS NOT NULL
- AND landing_video_id > 0
- AND status IN ('posted_ok', 'submitted')
- AND created_at >= DATE_SUB(NOW(), INTERVAL %s DAY)
- """,
- (crowd_package, int(lookback_days)),
- )
- usage_rows = cur.fetchall() or []
- cur.execute(
- """
- SELECT account_id, adgroup_id, landing_video_id, material_id, raw_record
- FROM creative_creation_task
- WHERE landing_video_id IS NOT NULL
- AND landing_video_id > 0
- AND submitted_at >= DATE_SUB(NOW(), INTERVAL %s DAY)
- AND JSON_UNQUOTE(JSON_EXTRACT(raw_record, '$.audience_tier'))=%s
- """,
- (int(lookback_days), crowd_package),
- )
- task_rows = cur.fetchall() or []
- finally:
- conn.close()
- seen: set[tuple[str, int, int, int, str]] = set()
- for row in list(usage_rows) + list(task_rows):
- try:
- landing_video_id = int(row["landing_video_id"])
- account_id = int(row.get("account_id") or 0)
- adgroup_id = int(row.get("adgroup_id") or 0)
- except (TypeError, ValueError):
- continue
- material_id = str(row.get("material_id") or "")
- source = _material_source_from_record(material_id, row.get("raw_record"))
- key = (source, account_id, adgroup_id, landing_video_id, material_id)
- if key in seen:
- continue
- seen.add(key)
- counts.setdefault(source, {})
- counts[source][landing_video_id] = counts[source].get(landing_video_id, 0) + 1
- return counts
- def load_recoverable_prepared_records(
- *,
- account_id: int,
- adgroup_id: int,
- crowd_package: str,
- material_source: str,
- limit: int,
- lookback_hours: int = CREATION_RECOVERY_LOOKBACK_HOURS,
- ) -> list[dict]:
- """Load unsubmitted prepared records after an interrupted Phase 1 run."""
- if limit <= 0 or not crowd_package:
- return []
- ensure_usage_table()
- from db.connection import get_connection
- conn = get_connection()
- try:
- with conn.cursor() as cur:
- cur.execute(
- """
- SELECT raw_record
- FROM creative_material_usage
- WHERE account_id=%s
- AND adgroup_id <=> %s
- AND crowd_package=%s
- AND status='prepared'
- AND dynamic_creative_id IS NULL
- AND raw_record IS NOT NULL
- AND raw_record <> ''
- AND created_at >= DATE_SUB(NOW(), INTERVAL %s HOUR)
- ORDER BY id ASC
- LIMIT %s
- """,
- (
- int(account_id),
- int(adgroup_id) if adgroup_id else None,
- crowd_package,
- int(lookback_hours),
- max(int(limit) * 3, int(limit)),
- ),
- )
- rows = cur.fetchall() or []
- finally:
- conn.close()
- out: list[dict] = []
- seen_materials: set[str] = set()
- for row in rows:
- try:
- record = json.loads(str(row.get("raw_record") or ""))
- except Exception:
- continue
- material_id = str(record.get("_material_id") or "").strip()
- if not material_id or material_id in seen_materials:
- continue
- if _material_source_from_record(material_id, row.get("raw_record")) != material_source:
- continue
- if not record.get("_request_body"):
- continue
- seen_materials.add(material_id)
- out.append(record)
- if len(out) >= limit:
- break
- return out
- def record_prepared_material_usage(record: dict, status: str = "prepared") -> None:
- """Reserve a material once Phase 1 has produced a pending creative."""
- material_id = str(record.get("_material_id") or "").strip()
- crowd_package = str(record.get("audience_tier") or "").strip()
- if not material_id or not crowd_package:
- return
- ensure_usage_table()
- from db.connection import get_connection
- raw_record = json.dumps(record, ensure_ascii=False, default=str)
- conn = get_connection()
- try:
- with conn.cursor() as cur:
- cur.execute(
- """
- INSERT INTO creative_material_usage
- (account_id, adgroup_id, crowd_package, landing_video_id,
- material_id, material_image_id, dynamic_creative_id,
- status, source, raw_record)
- VALUES (%s, %s, %s, %s, %s, %s, %s, %s, %s, %s)
- """,
- (
- int(record["account_id"]),
- int(record.get("adgroup_id") or 0) or None,
- crowd_package,
- record.get("landing_video_id"),
- material_id,
- str(record.get("_material_image_id") or ""),
- record.get("dynamic_creative_id"),
- status,
- record.get("landing_source"),
- raw_record,
- ),
- )
- conn.commit()
- finally:
- conn.close()
- def update_material_usage_status(
- record: dict,
- status: str,
- dynamic_creative_id: int | str | None = None,
- error: str = "",
- ) -> None:
- """Update the latest usage reservation for this crowd_package + material_id."""
- material_id = str(record.get("_material_id") or "").strip()
- crowd_package = str(record.get("audience_tier") or "").strip()
- if not material_id or not crowd_package:
- return
- ensure_usage_table()
- from db.connection import get_connection
- normalized_status = {
- "reject": "rejected",
- "approve": "approved",
- "hold": "no_result",
- "skip": "no_result",
- }.get(status, status)
- raw_record = json.dumps(record, ensure_ascii=False, default=str)
- conn = get_connection()
- try:
- with conn.cursor() as cur:
- cur.execute(
- """
- UPDATE creative_material_usage
- SET status=%s,
- dynamic_creative_id=COALESCE(%s, dynamic_creative_id),
- raw_record=%s,
- updated_at=CURRENT_TIMESTAMP
- WHERE crowd_package=%s
- AND material_id=%s
- AND account_id=%s
- AND adgroup_id <=> %s
- ORDER BY id DESC
- LIMIT 1
- """,
- (
- normalized_status,
- int(dynamic_creative_id) if dynamic_creative_id else None,
- raw_record if not error else json.dumps(
- {**record, "usage_error": error[:2000]},
- ensure_ascii=False,
- default=str,
- ),
- crowd_package,
- material_id,
- int(record["account_id"]),
- int(record.get("adgroup_id") or 0) or None,
- ),
- )
- conn.commit()
- finally:
- conn.close()
- def merge_used_material_ids(*sets: Iterable[str]) -> set[str]:
- """合并多个素材 ID 集合,统一转字符串并去重(过滤空值)。"""
- out: set[str] = set()
- for values in sets:
- out.update(str(v) for v in values if str(v))
- return out
|