"""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: 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]: out: set[str] = set() for values in sets: out.update(str(v) for v in values if str(v)) return out