creative_material_usage.py 12 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344
  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. CREATE_USAGE_TABLE_SQL = """
  18. CREATE TABLE IF NOT EXISTS creative_material_usage (
  19. id BIGINT AUTO_INCREMENT PRIMARY KEY COMMENT '自增主键',
  20. account_id BIGINT NOT NULL COMMENT '腾讯广告账户ID',
  21. adgroup_id BIGINT DEFAULT NULL COMMENT '广告ID',
  22. crowd_package VARCHAR(200) NOT NULL COMMENT '投放人群包名称',
  23. landing_video_id BIGINT DEFAULT NULL COMMENT '承接视频ID,用于同人群包近期排重',
  24. material_id VARCHAR(100) NOT NULL COMMENT '内部素材ID',
  25. material_image_id VARCHAR(100) DEFAULT NULL COMMENT '腾讯图片ID',
  26. dynamic_creative_id BIGINT DEFAULT NULL COMMENT '腾讯动态创意ID',
  27. status VARCHAR(50) NOT NULL DEFAULT 'prepared' COMMENT 'prepared/submitted/failed',
  28. source VARCHAR(50) DEFAULT NULL COMMENT 'primary/hot',
  29. raw_record MEDIUMTEXT DEFAULT NULL COMMENT '准备记录JSON',
  30. created_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP COMMENT '创建时间',
  31. updated_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP COMMENT '更新时间',
  32. KEY idx_crowd_material_created (crowd_package, material_id, created_at),
  33. KEY idx_crowd_landing_created (crowd_package, landing_video_id, created_at),
  34. KEY idx_crowd_account_ad_material (crowd_package, account_id, adgroup_id, material_id),
  35. KEY idx_account_created (account_id, created_at),
  36. KEY idx_status_created (status, created_at)
  37. ) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COMMENT='创意素材使用历史/占位'
  38. """
  39. def ensure_usage_table() -> None:
  40. from db.connection import get_connection
  41. conn = get_connection()
  42. try:
  43. with conn.cursor() as cur:
  44. cur.execute(CREATE_USAGE_TABLE_SQL)
  45. cur.execute("SHOW INDEX FROM creative_material_usage WHERE Key_name = 'idx_crowd_landing_created'")
  46. if not cur.fetchall():
  47. cur.execute(
  48. """
  49. ALTER TABLE creative_material_usage
  50. ADD INDEX idx_crowd_landing_created (crowd_package, landing_video_id, created_at)
  51. """
  52. )
  53. conn.commit()
  54. finally:
  55. conn.close()
  56. def load_recent_used_material_ids(
  57. crowd_package: str,
  58. lookback_days: int = CREATIVE_MATERIAL_DEDUPE_LOOKBACK_DAYS,
  59. ) -> set[str]:
  60. """Load recent material reservations for the same crowd package."""
  61. if not crowd_package:
  62. return set()
  63. ensure_usage_table()
  64. from db.connection import get_connection
  65. conn = get_connection()
  66. try:
  67. with conn.cursor() as cur:
  68. cur.execute(
  69. """
  70. SELECT DISTINCT material_id
  71. FROM (
  72. SELECT material_id
  73. FROM creative_material_usage
  74. WHERE crowd_package=%s
  75. AND material_id IS NOT NULL
  76. AND material_id <> ''
  77. AND status IN ('posted_ok', 'rejected')
  78. AND created_at >= DATE_SUB(NOW(), INTERVAL %s DAY)
  79. UNION
  80. SELECT material_id
  81. FROM creative_creation_task
  82. WHERE material_id IS NOT NULL
  83. AND material_id <> ''
  84. AND review_status IN ('approved', 'rejected')
  85. AND submitted_at >= DATE_SUB(NOW(), INTERVAL %s DAY)
  86. AND JSON_UNQUOTE(JSON_EXTRACT(raw_record, '$.audience_tier'))=%s
  87. ) t
  88. """,
  89. (crowd_package, int(lookback_days), int(lookback_days), crowd_package),
  90. )
  91. rows = cur.fetchall() or []
  92. finally:
  93. conn.close()
  94. return {str(row["material_id"]) for row in rows if row.get("material_id")}
  95. def load_recent_used_landing_video_ids(
  96. crowd_package: str,
  97. lookback_days: int = CREATIVE_LANDING_DEDUPE_LOOKBACK_DAYS,
  98. ) -> set[int]:
  99. """Load recent landing-video reservations for the same crowd package."""
  100. if not crowd_package:
  101. return set()
  102. ensure_usage_table()
  103. from db.connection import get_connection
  104. conn = get_connection()
  105. try:
  106. with conn.cursor() as cur:
  107. cur.execute(
  108. """
  109. SELECT DISTINCT landing_video_id
  110. FROM (
  111. SELECT landing_video_id
  112. FROM creative_material_usage
  113. WHERE crowd_package=%s
  114. AND landing_video_id IS NOT NULL
  115. AND landing_video_id > 0
  116. AND created_at >= DATE_SUB(NOW(), INTERVAL %s DAY)
  117. UNION
  118. SELECT landing_video_id
  119. FROM creative_creation_task
  120. WHERE landing_video_id IS NOT NULL
  121. AND landing_video_id > 0
  122. AND submitted_at >= DATE_SUB(NOW(), INTERVAL %s DAY)
  123. AND JSON_UNQUOTE(JSON_EXTRACT(raw_record, '$.audience_tier'))=%s
  124. ) t
  125. """,
  126. (crowd_package, int(lookback_days), int(lookback_days), crowd_package),
  127. )
  128. rows = cur.fetchall() or []
  129. finally:
  130. conn.close()
  131. out: set[int] = set()
  132. for row in rows:
  133. try:
  134. out.add(int(row["landing_video_id"]))
  135. except (TypeError, ValueError):
  136. continue
  137. return out
  138. def _material_source_from_record(material_id: object, raw_record: object) -> str:
  139. material_id_text = str(material_id or "").strip()
  140. if material_id_text.startswith("ai:"):
  141. return "ai_generated"
  142. if raw_record:
  143. try:
  144. data = json.loads(str(raw_record))
  145. source = str(data.get("material_source") or "").strip()
  146. raw_material_id = str(data.get("_material_id") or "").strip()
  147. if source == "ai_generated" or raw_material_id.startswith("ai:"):
  148. return "ai_generated"
  149. except Exception:
  150. pass
  151. return "history"
  152. def load_recent_landing_usage_counts(
  153. crowd_package: str,
  154. lookback_days: int = CREATIVE_LANDING_DEDUPE_LOOKBACK_DAYS,
  155. ) -> dict[str, dict[int, int]]:
  156. """Load recent landing-video usage counts split by material source.
  157. The `source` column in creative_material_usage means primary/hot video source,
  158. so material source must be derived from material_id/raw_record.
  159. """
  160. counts: dict[str, dict[int, int]] = {
  161. "history": {},
  162. "ai_generated": {},
  163. }
  164. if not crowd_package:
  165. return counts
  166. ensure_usage_table()
  167. from db.connection import get_connection
  168. conn = get_connection()
  169. try:
  170. with conn.cursor() as cur:
  171. cur.execute(
  172. """
  173. SELECT account_id, adgroup_id, landing_video_id, material_id, raw_record
  174. FROM creative_material_usage
  175. WHERE crowd_package=%s
  176. AND landing_video_id IS NOT NULL
  177. AND landing_video_id > 0
  178. AND created_at >= DATE_SUB(NOW(), INTERVAL %s DAY)
  179. """,
  180. (crowd_package, int(lookback_days)),
  181. )
  182. usage_rows = cur.fetchall() or []
  183. cur.execute(
  184. """
  185. SELECT account_id, adgroup_id, landing_video_id, material_id, raw_record
  186. FROM creative_creation_task
  187. WHERE landing_video_id IS NOT NULL
  188. AND landing_video_id > 0
  189. AND submitted_at >= DATE_SUB(NOW(), INTERVAL %s DAY)
  190. AND JSON_UNQUOTE(JSON_EXTRACT(raw_record, '$.audience_tier'))=%s
  191. """,
  192. (int(lookback_days), crowd_package),
  193. )
  194. task_rows = cur.fetchall() or []
  195. finally:
  196. conn.close()
  197. seen: set[tuple[str, int, int, int, str]] = set()
  198. for row in list(usage_rows) + list(task_rows):
  199. try:
  200. landing_video_id = int(row["landing_video_id"])
  201. account_id = int(row.get("account_id") or 0)
  202. adgroup_id = int(row.get("adgroup_id") or 0)
  203. except (TypeError, ValueError):
  204. continue
  205. material_id = str(row.get("material_id") or "")
  206. source = _material_source_from_record(material_id, row.get("raw_record"))
  207. key = (source, account_id, adgroup_id, landing_video_id, material_id)
  208. if key in seen:
  209. continue
  210. seen.add(key)
  211. counts.setdefault(source, {})
  212. counts[source][landing_video_id] = counts[source].get(landing_video_id, 0) + 1
  213. return counts
  214. def record_prepared_material_usage(record: dict, status: str = "prepared") -> None:
  215. """Reserve a material once Phase 1 has produced a pending creative."""
  216. material_id = str(record.get("_material_id") or "").strip()
  217. crowd_package = str(record.get("audience_tier") or "").strip()
  218. if not material_id or not crowd_package:
  219. return
  220. ensure_usage_table()
  221. from db.connection import get_connection
  222. raw_record = json.dumps(record, ensure_ascii=False, default=str)
  223. conn = get_connection()
  224. try:
  225. with conn.cursor() as cur:
  226. cur.execute(
  227. """
  228. INSERT INTO creative_material_usage
  229. (account_id, adgroup_id, crowd_package, landing_video_id,
  230. material_id, material_image_id, dynamic_creative_id,
  231. status, source, raw_record)
  232. VALUES (%s, %s, %s, %s, %s, %s, %s, %s, %s, %s)
  233. """,
  234. (
  235. int(record["account_id"]),
  236. int(record.get("adgroup_id") or 0) or None,
  237. crowd_package,
  238. record.get("landing_video_id"),
  239. material_id,
  240. str(record.get("_material_image_id") or ""),
  241. record.get("dynamic_creative_id"),
  242. status,
  243. record.get("landing_source"),
  244. raw_record,
  245. ),
  246. )
  247. conn.commit()
  248. finally:
  249. conn.close()
  250. def update_material_usage_status(
  251. record: dict,
  252. status: str,
  253. dynamic_creative_id: int | str | None = None,
  254. error: str = "",
  255. ) -> None:
  256. """Update the latest usage reservation for this crowd_package + material_id."""
  257. material_id = str(record.get("_material_id") or "").strip()
  258. crowd_package = str(record.get("audience_tier") or "").strip()
  259. if not material_id or not crowd_package:
  260. return
  261. ensure_usage_table()
  262. from db.connection import get_connection
  263. normalized_status = {
  264. "reject": "rejected",
  265. "approve": "approved",
  266. "hold": "no_result",
  267. "skip": "no_result",
  268. }.get(status, status)
  269. raw_record = json.dumps(record, ensure_ascii=False, default=str)
  270. conn = get_connection()
  271. try:
  272. with conn.cursor() as cur:
  273. cur.execute(
  274. """
  275. UPDATE creative_material_usage
  276. SET status=%s,
  277. dynamic_creative_id=COALESCE(%s, dynamic_creative_id),
  278. raw_record=%s,
  279. updated_at=CURRENT_TIMESTAMP
  280. WHERE crowd_package=%s
  281. AND material_id=%s
  282. AND account_id=%s
  283. AND adgroup_id <=> %s
  284. ORDER BY id DESC
  285. LIMIT 1
  286. """,
  287. (
  288. normalized_status,
  289. int(dynamic_creative_id) if dynamic_creative_id else None,
  290. raw_record if not error else json.dumps(
  291. {**record, "usage_error": error[:2000]},
  292. ensure_ascii=False,
  293. default=str,
  294. ),
  295. crowd_package,
  296. material_id,
  297. int(record["account_id"]),
  298. int(record.get("adgroup_id") or 0) or None,
  299. ),
  300. )
  301. conn.commit()
  302. finally:
  303. conn.close()
  304. def merge_used_material_ids(*sets: Iterable[str]) -> set[str]:
  305. out: set[str] = set()
  306. for values in sets:
  307. out.update(str(v) for v in values if str(v))
  308. return out