creative_material_usage.py 15 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422
  1. """Creative material usage reservation for duplicate control.
  2. Usage is written when a pending creative is prepared, before Tencent POST.
  3. Review/task tables remain responsible for submitted Tencent creatives.
  4. """
  5. from __future__ import annotations
  6. import json
  7. import logging
  8. import os
  9. from typing import Iterable
  10. logger = logging.getLogger(__name__)
  11. CREATIVE_MATERIAL_DEDUPE_LOOKBACK_DAYS = int(
  12. os.getenv("CREATIVE_MATERIAL_DEDUPE_LOOKBACK_DAYS", "7")
  13. )
  14. CREATIVE_LANDING_DEDUPE_LOOKBACK_DAYS = int(
  15. os.getenv("CREATIVE_LANDING_DEDUPE_LOOKBACK_DAYS", "7")
  16. )
  17. CREATION_RECOVERY_LOOKBACK_HOURS = int(
  18. os.getenv("CREATION_RECOVERY_LOOKBACK_HOURS", "24")
  19. )
  20. CREATE_USAGE_TABLE_SQL = """
  21. CREATE TABLE IF NOT EXISTS creative_material_usage (
  22. id BIGINT AUTO_INCREMENT PRIMARY KEY COMMENT '自增主键',
  23. account_id BIGINT NOT NULL COMMENT '腾讯广告账户ID',
  24. adgroup_id BIGINT DEFAULT NULL COMMENT '广告ID',
  25. crowd_package VARCHAR(200) NOT NULL COMMENT '投放人群包名称',
  26. landing_video_id BIGINT DEFAULT NULL COMMENT '承接视频ID,用于同人群包近期排重',
  27. material_id VARCHAR(100) NOT NULL COMMENT '内部素材ID',
  28. material_image_id VARCHAR(100) DEFAULT NULL COMMENT '腾讯图片ID',
  29. dynamic_creative_id BIGINT DEFAULT NULL COMMENT '腾讯动态创意ID',
  30. status VARCHAR(50) NOT NULL DEFAULT 'prepared' COMMENT 'prepared/submitted/failed',
  31. source VARCHAR(50) DEFAULT NULL COMMENT 'primary/hot',
  32. raw_record MEDIUMTEXT DEFAULT NULL COMMENT '准备记录JSON',
  33. created_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP COMMENT '创建时间',
  34. updated_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP COMMENT '更新时间',
  35. KEY idx_crowd_material_created (crowd_package, material_id, created_at),
  36. KEY idx_crowd_landing_created (crowd_package, landing_video_id, created_at),
  37. KEY idx_crowd_account_ad_material (crowd_package, account_id, adgroup_id, material_id),
  38. KEY idx_account_created (account_id, created_at),
  39. KEY idx_status_created (status, created_at)
  40. ) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COMMENT='创意素材使用历史/占位'
  41. """
  42. def ensure_usage_table() -> None:
  43. """确保 creative_material_usage 表存在(不存在则创建)。"""
  44. from db.connection import get_connection
  45. conn = get_connection()
  46. try:
  47. with conn.cursor() as cur:
  48. cur.execute(CREATE_USAGE_TABLE_SQL)
  49. cur.execute("SHOW INDEX FROM creative_material_usage WHERE Key_name = 'idx_crowd_landing_created'")
  50. if not cur.fetchall():
  51. cur.execute(
  52. """
  53. ALTER TABLE creative_material_usage
  54. ADD INDEX idx_crowd_landing_created (crowd_package, landing_video_id, created_at)
  55. """
  56. )
  57. conn.commit()
  58. finally:
  59. conn.close()
  60. def load_recent_used_material_ids(
  61. crowd_package: str,
  62. lookback_days: int = CREATIVE_MATERIAL_DEDUPE_LOOKBACK_DAYS,
  63. ) -> set[str]:
  64. """Load recent material reservations for the same crowd package."""
  65. if not crowd_package:
  66. return set()
  67. ensure_usage_table()
  68. from db.connection import get_connection
  69. conn = get_connection()
  70. try:
  71. with conn.cursor() as cur:
  72. cur.execute(
  73. """
  74. SELECT DISTINCT material_id
  75. FROM (
  76. SELECT material_id
  77. FROM creative_material_usage
  78. WHERE crowd_package=%s
  79. AND material_id IS NOT NULL
  80. AND material_id <> ''
  81. AND status IN ('posted_ok', 'rejected')
  82. AND created_at >= DATE_SUB(NOW(), INTERVAL %s DAY)
  83. UNION
  84. SELECT material_id
  85. FROM creative_creation_task
  86. WHERE material_id IS NOT NULL
  87. AND material_id <> ''
  88. AND review_status IN ('approved', 'rejected')
  89. AND submitted_at >= DATE_SUB(NOW(), INTERVAL %s DAY)
  90. AND JSON_UNQUOTE(JSON_EXTRACT(raw_record, '$.audience_tier'))=%s
  91. ) t
  92. """,
  93. (crowd_package, int(lookback_days), int(lookback_days), crowd_package),
  94. )
  95. rows = cur.fetchall() or []
  96. finally:
  97. conn.close()
  98. return {str(row["material_id"]) for row in rows if row.get("material_id")}
  99. def load_recent_used_landing_video_ids(
  100. crowd_package: str,
  101. lookback_days: int = CREATIVE_LANDING_DEDUPE_LOOKBACK_DAYS,
  102. ) -> set[int]:
  103. """Load recent landing-video reservations for the same crowd package."""
  104. if not crowd_package:
  105. return set()
  106. ensure_usage_table()
  107. from db.connection import get_connection
  108. conn = get_connection()
  109. try:
  110. with conn.cursor() as cur:
  111. cur.execute(
  112. """
  113. SELECT DISTINCT landing_video_id
  114. FROM (
  115. SELECT landing_video_id
  116. FROM creative_material_usage
  117. WHERE crowd_package=%s
  118. AND landing_video_id IS NOT NULL
  119. AND landing_video_id > 0
  120. AND status IN ('posted_ok', 'submitted')
  121. AND created_at >= DATE_SUB(NOW(), INTERVAL %s DAY)
  122. UNION
  123. SELECT landing_video_id
  124. FROM creative_creation_task
  125. WHERE landing_video_id IS NOT NULL
  126. AND landing_video_id > 0
  127. AND submitted_at >= DATE_SUB(NOW(), INTERVAL %s DAY)
  128. AND JSON_UNQUOTE(JSON_EXTRACT(raw_record, '$.audience_tier'))=%s
  129. ) t
  130. """,
  131. (crowd_package, int(lookback_days), int(lookback_days), crowd_package),
  132. )
  133. rows = cur.fetchall() or []
  134. finally:
  135. conn.close()
  136. out: set[int] = set()
  137. for row in rows:
  138. try:
  139. out.add(int(row["landing_video_id"]))
  140. except (TypeError, ValueError):
  141. continue
  142. return out
  143. def _material_source_from_record(material_id: object, raw_record: object) -> str:
  144. material_id_text = str(material_id or "").strip()
  145. if material_id_text.startswith("external:"):
  146. return "external_recall"
  147. if material_id_text.startswith("ai:"):
  148. return "ai_generated"
  149. if raw_record:
  150. try:
  151. data = json.loads(str(raw_record))
  152. source = str(data.get("material_source") or "").strip()
  153. raw_material_id = str(data.get("_material_id") or "").strip()
  154. if source == "external_recall" or raw_material_id.startswith("external:"):
  155. return "external_recall"
  156. if source == "ai_generated" or raw_material_id.startswith("ai:"):
  157. return "ai_generated"
  158. except Exception:
  159. pass
  160. return "history"
  161. def load_recent_landing_usage_counts(
  162. crowd_package: str,
  163. lookback_days: int = CREATIVE_LANDING_DEDUPE_LOOKBACK_DAYS,
  164. ) -> dict[str, dict[int, int]]:
  165. """Load recent landing-video usage counts split by material source.
  166. The `source` column in creative_material_usage means primary/hot video source,
  167. so material source must be derived from material_id/raw_record.
  168. """
  169. counts: dict[str, dict[int, int]] = {
  170. "history": {},
  171. "ai_generated": {},
  172. "external_recall": {},
  173. }
  174. if not crowd_package:
  175. return counts
  176. ensure_usage_table()
  177. from db.connection import get_connection
  178. conn = get_connection()
  179. try:
  180. with conn.cursor() as cur:
  181. cur.execute(
  182. """
  183. SELECT account_id, adgroup_id, landing_video_id, material_id, raw_record
  184. FROM creative_material_usage
  185. WHERE crowd_package=%s
  186. AND landing_video_id IS NOT NULL
  187. AND landing_video_id > 0
  188. AND status IN ('posted_ok', 'submitted')
  189. AND created_at >= DATE_SUB(NOW(), INTERVAL %s DAY)
  190. """,
  191. (crowd_package, int(lookback_days)),
  192. )
  193. usage_rows = cur.fetchall() or []
  194. cur.execute(
  195. """
  196. SELECT account_id, adgroup_id, landing_video_id, material_id, raw_record
  197. FROM creative_creation_task
  198. WHERE landing_video_id IS NOT NULL
  199. AND landing_video_id > 0
  200. AND submitted_at >= DATE_SUB(NOW(), INTERVAL %s DAY)
  201. AND JSON_UNQUOTE(JSON_EXTRACT(raw_record, '$.audience_tier'))=%s
  202. """,
  203. (int(lookback_days), crowd_package),
  204. )
  205. task_rows = cur.fetchall() or []
  206. finally:
  207. conn.close()
  208. seen: set[tuple[str, int, int, int, str]] = set()
  209. for row in list(usage_rows) + list(task_rows):
  210. try:
  211. landing_video_id = int(row["landing_video_id"])
  212. account_id = int(row.get("account_id") or 0)
  213. adgroup_id = int(row.get("adgroup_id") or 0)
  214. except (TypeError, ValueError):
  215. continue
  216. material_id = str(row.get("material_id") or "")
  217. source = _material_source_from_record(material_id, row.get("raw_record"))
  218. key = (source, account_id, adgroup_id, landing_video_id, material_id)
  219. if key in seen:
  220. continue
  221. seen.add(key)
  222. counts.setdefault(source, {})
  223. counts[source][landing_video_id] = counts[source].get(landing_video_id, 0) + 1
  224. return counts
  225. def load_recoverable_prepared_records(
  226. *,
  227. account_id: int,
  228. adgroup_id: int,
  229. crowd_package: str,
  230. material_source: str,
  231. limit: int,
  232. lookback_hours: int = CREATION_RECOVERY_LOOKBACK_HOURS,
  233. ) -> list[dict]:
  234. """Load unsubmitted prepared records after an interrupted Phase 1 run."""
  235. if limit <= 0 or not crowd_package:
  236. return []
  237. ensure_usage_table()
  238. from db.connection import get_connection
  239. conn = get_connection()
  240. try:
  241. with conn.cursor() as cur:
  242. cur.execute(
  243. """
  244. SELECT raw_record
  245. FROM creative_material_usage
  246. WHERE account_id=%s
  247. AND adgroup_id <=> %s
  248. AND crowd_package=%s
  249. AND status='prepared'
  250. AND dynamic_creative_id IS NULL
  251. AND raw_record IS NOT NULL
  252. AND raw_record <> ''
  253. AND created_at >= DATE_SUB(NOW(), INTERVAL %s HOUR)
  254. ORDER BY id ASC
  255. LIMIT %s
  256. """,
  257. (
  258. int(account_id),
  259. int(adgroup_id) if adgroup_id else None,
  260. crowd_package,
  261. int(lookback_hours),
  262. max(int(limit) * 3, int(limit)),
  263. ),
  264. )
  265. rows = cur.fetchall() or []
  266. finally:
  267. conn.close()
  268. out: list[dict] = []
  269. seen_materials: set[str] = set()
  270. for row in rows:
  271. try:
  272. record = json.loads(str(row.get("raw_record") or ""))
  273. except Exception:
  274. continue
  275. material_id = str(record.get("_material_id") or "").strip()
  276. if not material_id or material_id in seen_materials:
  277. continue
  278. if _material_source_from_record(material_id, row.get("raw_record")) != material_source:
  279. continue
  280. if not record.get("_request_body"):
  281. continue
  282. seen_materials.add(material_id)
  283. out.append(record)
  284. if len(out) >= limit:
  285. break
  286. return out
  287. def record_prepared_material_usage(record: dict, status: str = "prepared") -> None:
  288. """Reserve a material once Phase 1 has produced a pending creative."""
  289. material_id = str(record.get("_material_id") or "").strip()
  290. crowd_package = str(record.get("audience_tier") or "").strip()
  291. if not material_id or not crowd_package:
  292. return
  293. ensure_usage_table()
  294. from db.connection import get_connection
  295. raw_record = json.dumps(record, ensure_ascii=False, default=str)
  296. conn = get_connection()
  297. try:
  298. with conn.cursor() as cur:
  299. cur.execute(
  300. """
  301. INSERT INTO creative_material_usage
  302. (account_id, adgroup_id, crowd_package, landing_video_id,
  303. material_id, material_image_id, dynamic_creative_id,
  304. status, source, raw_record)
  305. VALUES (%s, %s, %s, %s, %s, %s, %s, %s, %s, %s)
  306. """,
  307. (
  308. int(record["account_id"]),
  309. int(record.get("adgroup_id") or 0) or None,
  310. crowd_package,
  311. record.get("landing_video_id"),
  312. material_id,
  313. str(record.get("_material_image_id") or ""),
  314. record.get("dynamic_creative_id"),
  315. status,
  316. record.get("landing_source"),
  317. raw_record,
  318. ),
  319. )
  320. conn.commit()
  321. finally:
  322. conn.close()
  323. def update_material_usage_status(
  324. record: dict,
  325. status: str,
  326. dynamic_creative_id: int | str | None = None,
  327. error: str = "",
  328. ) -> None:
  329. """Update the latest usage reservation for this crowd_package + material_id."""
  330. material_id = str(record.get("_material_id") or "").strip()
  331. crowd_package = str(record.get("audience_tier") or "").strip()
  332. if not material_id or not crowd_package:
  333. return
  334. ensure_usage_table()
  335. from db.connection import get_connection
  336. normalized_status = {
  337. "reject": "rejected",
  338. "approve": "approved",
  339. "hold": "no_result",
  340. "skip": "no_result",
  341. }.get(status, status)
  342. raw_record = json.dumps(record, ensure_ascii=False, default=str)
  343. conn = get_connection()
  344. try:
  345. with conn.cursor() as cur:
  346. cur.execute(
  347. """
  348. UPDATE creative_material_usage
  349. SET status=%s,
  350. dynamic_creative_id=COALESCE(%s, dynamic_creative_id),
  351. raw_record=%s,
  352. updated_at=CURRENT_TIMESTAMP
  353. WHERE crowd_package=%s
  354. AND material_id=%s
  355. AND account_id=%s
  356. AND adgroup_id <=> %s
  357. ORDER BY id DESC
  358. LIMIT 1
  359. """,
  360. (
  361. normalized_status,
  362. int(dynamic_creative_id) if dynamic_creative_id else None,
  363. raw_record if not error else json.dumps(
  364. {**record, "usage_error": error[:2000]},
  365. ensure_ascii=False,
  366. default=str,
  367. ),
  368. crowd_package,
  369. material_id,
  370. int(record["account_id"]),
  371. int(record.get("adgroup_id") or 0) or None,
  372. ),
  373. )
  374. conn.commit()
  375. finally:
  376. conn.close()
  377. def merge_used_material_ids(*sets: Iterable[str]) -> set[str]:
  378. """合并多个素材 ID 集合,统一转字符串并去重(过滤空值)。"""
  379. out: set[str] = set()
  380. for values in sets:
  381. out.update(str(v) for v in values if str(v))
  382. return out