"""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") ) 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_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) 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 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