creative_review.py 28 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726727728729730731732733734735736
  1. """腾讯动态创意审核结果扫描器。
  2. 本模块跟踪创建链路已提交的创意,并轮询腾讯官方审核结果 API,不执行预审。
  3. """
  4. from __future__ import annotations
  5. import json
  6. import logging
  7. from dataclasses import dataclass, field
  8. from datetime import datetime, timezone
  9. from typing import Iterable, Iterator, Optional
  10. logger = logging.getLogger(__name__)
  11. CREATE_REVIEW_TABLES_SQL = [
  12. """
  13. CREATE TABLE IF NOT EXISTS creative_creation_task (
  14. id BIGINT AUTO_INCREMENT PRIMARY KEY COMMENT '自增主键',
  15. account_id BIGINT NOT NULL COMMENT '腾讯广告账户ID',
  16. adgroup_id BIGINT DEFAULT NULL COMMENT '广告ID',
  17. dynamic_creative_id BIGINT NOT NULL COMMENT '动态创意ID',
  18. dynamic_creative_name VARCHAR(200) DEFAULT NULL COMMENT '动态创意名称',
  19. landing_video_id BIGINT DEFAULT NULL COMMENT '承接视频ID',
  20. material_id VARCHAR(100) DEFAULT NULL COMMENT '内部素材ID',
  21. material_image_id VARCHAR(100) DEFAULT NULL COMMENT '腾讯图片ID',
  22. review_status VARCHAR(50) NOT NULL DEFAULT 'submitted' COMMENT '审核状态',
  23. submit_status VARCHAR(50) NOT NULL DEFAULT 'submitted' COMMENT '提交状态',
  24. submitted_at TIMESTAMP NOT NULL DEFAULT CURRENT_TIMESTAMP COMMENT '提交时间',
  25. last_review_checked_at TIMESTAMP NULL DEFAULT NULL COMMENT '最近审核扫描时间',
  26. review_finished_at TIMESTAMP NULL DEFAULT NULL COMMENT '审核完成时间',
  27. error TEXT DEFAULT NULL COMMENT '提交或扫描错误',
  28. raw_record MEDIUMTEXT DEFAULT NULL COMMENT '创建记录JSON',
  29. created_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP COMMENT '创建时间',
  30. updated_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP COMMENT '更新时间',
  31. UNIQUE KEY uk_account_creative (account_id, dynamic_creative_id),
  32. KEY idx_review_status (review_status),
  33. KEY idx_submitted_at (submitted_at),
  34. KEY idx_landing_submitted (landing_video_id, submitted_at),
  35. KEY idx_last_review_checked_at (last_review_checked_at),
  36. KEY idx_account_status (account_id, review_status)
  37. ) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COMMENT='创意创建与审核扫描任务'
  38. """,
  39. """
  40. CREATE TABLE IF NOT EXISTS creative_review_result (
  41. id BIGINT AUTO_INCREMENT PRIMARY KEY COMMENT '自增主键',
  42. account_id BIGINT NOT NULL COMMENT '腾讯广告账户ID',
  43. dynamic_creative_id BIGINT NOT NULL COMMENT '动态创意ID',
  44. review_status VARCHAR(50) NOT NULL COMMENT '解析后的审核状态',
  45. reject_messages MEDIUMTEXT DEFAULT NULL COMMENT '拒绝原因JSON数组',
  46. delay_messages MEDIUMTEXT DEFAULT NULL COMMENT '延迟审核原因JSON数组',
  47. raw_result MEDIUMTEXT NOT NULL COMMENT '腾讯审核结果原文JSON',
  48. checked_at TIMESTAMP NOT NULL DEFAULT CURRENT_TIMESTAMP COMMENT '扫描时间',
  49. created_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP COMMENT '创建时间',
  50. updated_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP COMMENT '更新时间',
  51. UNIQUE KEY uk_account_creative (account_id, dynamic_creative_id),
  52. KEY idx_review_status (review_status),
  53. KEY idx_checked_at (checked_at)
  54. ) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COMMENT='腾讯动态创意审核结果'
  55. """,
  56. """
  57. CREATE TABLE IF NOT EXISTS creative_rejection_fact (
  58. id BIGINT AUTO_INCREMENT PRIMARY KEY COMMENT '自增主键',
  59. account_id BIGINT NOT NULL COMMENT '腾讯广告账户ID',
  60. dynamic_creative_id BIGINT NOT NULL COMMENT '动态创意ID',
  61. fact_type VARCHAR(50) NOT NULL COMMENT '拒绝事实类型:site/element/compose',
  62. fact_key VARCHAR(255) DEFAULT NULL COMMENT '元素/版位/组件标识',
  63. reason TEXT NOT NULL COMMENT '拒绝原因',
  64. site_set VARCHAR(100) DEFAULT NULL COMMENT '影响版位',
  65. element_type VARCHAR(100) DEFAULT NULL COMMENT '元素类型',
  66. component_type VARCHAR(100) DEFAULT NULL COMMENT '组件类型',
  67. raw_detail MEDIUMTEXT DEFAULT NULL COMMENT '原始明细JSON',
  68. created_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP COMMENT '创建时间',
  69. KEY idx_account_creative (account_id, dynamic_creative_id),
  70. KEY idx_fact_type (fact_type),
  71. KEY idx_fact_key (fact_key(100))
  72. ) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COMMENT='创意审核拒绝事实'
  73. """,
  74. ]
  75. FINAL_REVIEW_STATUSES = {"approved", "rejected"}
  76. PENDING_REVIEW_STATUSES = {"submitted", "pending", "unknown"}
  77. # 面向业务报表的腾讯 API 枚举说明;未知值保持原样,便于识别和诊断新增枚举。
  78. STATUS_DESCRIPTIONS = {
  79. "AD_STATUS_NORMAL": "投放中",
  80. "AD_STATUS_SUSPEND": "暂停",
  81. "AD_STATUS_DELETED": "已删除",
  82. "AD_STATUS_DENIED": "审核拒绝",
  83. "CREATIVE_SET_APPROVAL_STATUS_UNKNOWN": "未知",
  84. "CREATIVE_SET_APPROVAL_STATUS_PENDING": "审核中",
  85. "CREATIVE_SET_APPROVAL_STATUS_NORMAL": "投放中",
  86. "CREATIVE_SET_APPROVAL_STATUS_DENIED": "审核拒绝",
  87. "CREATIVE_SET_APPROVAL_STATUS_PARTIAL_NORMAL": "部分投放中",
  88. "DYNAMIC_CREATIVE_STATUS_DENIED": "审核拒绝",
  89. "DYNAMIC_CREATIVE_STATUS_PENDING": "审核中",
  90. "REVIEW_STATUS_REJECTED": "审核拒绝",
  91. "REVIEW_STATUS_APPROVED": "审核通过",
  92. "REVIEW_STATUS_PENDING": "审核中",
  93. }
  94. @dataclass
  95. class RejectionFact:
  96. fact_type: str
  97. reason: str
  98. fact_key: str = ""
  99. site_set: str = ""
  100. element_type: str = ""
  101. component_type: str = ""
  102. raw_detail: dict = field(default_factory=dict)
  103. @dataclass
  104. class ParsedReviewResult:
  105. dynamic_creative_id: Optional[int]
  106. review_status: str
  107. reject_messages: list[str] = field(default_factory=list)
  108. delay_messages: list[str] = field(default_factory=list)
  109. rejection_facts: list[RejectionFact] = field(default_factory=list)
  110. def _display_key(value: object, fallback: str) -> str:
  111. key = str(value or "").strip()
  112. return key or fallback
  113. def _review_item_label(
  114. item: dict,
  115. fallback: str,
  116. *,
  117. site: bool = False,
  118. include_ids: bool = False,
  119. ) -> str:
  120. if site:
  121. base = _display_key(item.get("site_set"), fallback)
  122. id_keys = ("site_set_id", "site_id", "id")
  123. else:
  124. base = _display_key(
  125. item.get("element_name")
  126. or item.get("element_type"),
  127. fallback,
  128. )
  129. id_keys = ("element_id", "image_id", "video_id", "id")
  130. if not include_ids:
  131. return base
  132. identifiers = []
  133. for key in id_keys:
  134. value = str(item.get(key) or "").strip()
  135. if value and value not in identifiers:
  136. identifiers.append(value)
  137. return f"{base}(id={','.join(identifiers)})" if identifiers else base
  138. def _is_normal_review_status(value: object) -> bool:
  139. return str(value or "").strip() == "AD_STATUS_NORMAL"
  140. def status_desc(value: object) -> str:
  141. """返回 API 的中文说明,未知枚举值保持原样。"""
  142. raw = str(value or "").strip()
  143. return STATUS_DESCRIPTIONS.get(raw, raw)
  144. def _labeled_values(
  145. values: Iterable[tuple[str, object]],
  146. *,
  147. describe_status: bool = False,
  148. ) -> str:
  149. rendered = []
  150. for label, raw_value in values:
  151. value = str(raw_value or "").strip()
  152. if value:
  153. if describe_status:
  154. value = status_desc(value)
  155. rendered.append(f"{label}: {value}")
  156. return "\n".join(_dedupe(rendered))
  157. def review_granularity_fields(result: dict | None) -> dict[str, str]:
  158. """将腾讯元素和站点审核结果渲染为业务报表字段。"""
  159. raw = result if isinstance(result, dict) else {}
  160. element_statuses: list[tuple[str, object]] = []
  161. element_reasons: list[tuple[str, object]] = []
  162. for index, element in enumerate(raw.get("element_result_list") or [], start=1):
  163. if not isinstance(element, dict):
  164. continue
  165. status = element.get("review_status") or element.get("system_status")
  166. if _is_normal_review_status(status):
  167. continue
  168. details = element.get("element_reject_detail_info") or []
  169. reasons = [
  170. detail.get("reason")
  171. for detail in details
  172. if isinstance(detail, dict) and detail.get("reason")
  173. ]
  174. if not reasons and element.get("reason"):
  175. reasons = [element.get("reason")]
  176. label = _review_item_label(
  177. element,
  178. f"元素{index}",
  179. include_ids=_status_is_rejected(str(status or "")) or bool(reasons),
  180. )
  181. element_statuses.append((label, status))
  182. element_reasons.extend((label, reason) for reason in reasons)
  183. site_statuses: list[tuple[str, object]] = []
  184. site_reasons: list[tuple[str, object]] = []
  185. for index, site in enumerate(raw.get("site_set_result_list") or [], start=1):
  186. if not isinstance(site, dict):
  187. continue
  188. status = site.get("system_status") or site.get("review_status")
  189. if _is_normal_review_status(status):
  190. continue
  191. reasons = []
  192. if site.get("reject_message"):
  193. reasons.append(site.get("reject_message"))
  194. for detail in site.get("element_reject_detail_info") or []:
  195. if isinstance(detail, dict) and detail.get("reason"):
  196. reasons.append(detail.get("reason"))
  197. label = _review_item_label(
  198. site,
  199. f"版位{index}",
  200. site=True,
  201. include_ids=_status_is_rejected(str(status or "")) or bool(reasons),
  202. )
  203. site_statuses.append((label, status))
  204. site_reasons.extend((label, reason) for reason in reasons)
  205. return {
  206. "element_review_status": _labeled_values(
  207. element_statuses,
  208. describe_status=True,
  209. ),
  210. "element_reject_reason": _labeled_values(element_reasons),
  211. "site_review_status": _labeled_values(
  212. site_statuses,
  213. describe_status=True,
  214. ),
  215. "site_reject_reason": _labeled_values(site_reasons),
  216. }
  217. def _dedupe(values: Iterable[str]) -> list[str]:
  218. """字符串列表去重,保留首次出现顺序,跳过空字符串。"""
  219. out: list[str] = []
  220. seen: set[str] = set()
  221. for raw in values:
  222. value = str(raw or "").strip()
  223. if not value or value in seen:
  224. continue
  225. seen.add(value)
  226. out.append(value)
  227. return out
  228. def _status_is_rejected(status: str) -> bool:
  229. """判断审核状态是否为拒绝(REJECT 或 DENIED)。"""
  230. upper = (status or "").upper()
  231. return "REJECT" in upper or "DENIED" in upper
  232. def _status_is_pending(status: str) -> bool:
  233. """判断审核状态是否为待审(PENDING 或 REVIEWING)。"""
  234. upper = (status or "").upper()
  235. return "PENDING" in upper or "REVIEWING" in upper
  236. def _status_is_approved(status: str) -> bool:
  237. """判断审核状态是否为通过(PASS 或 APPROVED 或 NORMAL)。"""
  238. upper = (status or "").upper()
  239. return "PASS" in upper or "APPROVED" in upper or "NORMAL" in upper
  240. def _component_type(detail: dict) -> str:
  241. """从审核明细中提取组件类型(component_type)。"""
  242. component = detail.get("component_info") or {}
  243. if isinstance(component, dict):
  244. return str(component.get("component_type") or "")
  245. return ""
  246. def _collect_element_facts(result: dict) -> tuple[list[str], list[RejectionFact], bool, bool]:
  247. """从审核结果中提取元素级拒绝事实,返回 (拒绝消息列表, 拒绝事实列表, 是否有待审, 是否有通过)。"""
  248. messages: list[str] = []
  249. facts: list[RejectionFact] = []
  250. has_pending = False
  251. has_approved = False
  252. for element in result.get("element_result_list") or []:
  253. if not isinstance(element, dict):
  254. continue
  255. status = str(element.get("review_status") or "")
  256. has_pending = has_pending or _status_is_pending(status)
  257. has_approved = has_approved or _status_is_approved(status)
  258. if not _status_is_rejected(status):
  259. continue
  260. details = element.get("element_reject_detail_info") or []
  261. if not details:
  262. details = [{"reason": element.get("reason") or ""}]
  263. for detail in details:
  264. if not isinstance(detail, dict):
  265. continue
  266. reason = detail.get("reason") or element.get("reason") or ""
  267. if reason:
  268. messages.append(str(reason))
  269. facts.append(RejectionFact(
  270. fact_type="element",
  271. fact_key=str(
  272. element.get("image_id")
  273. or element.get("video_id")
  274. or element.get("element_name")
  275. or ""
  276. ),
  277. reason=str(reason),
  278. element_type=str(element.get("element_type") or ""),
  279. component_type=_component_type(element),
  280. raw_detail={**element, "element_reject_detail": detail},
  281. ))
  282. return messages, facts, has_pending, has_approved
  283. def _collect_site_facts(result: dict) -> tuple[list[str], list[RejectionFact], bool, bool]:
  284. """从审核结果中提取版位级拒绝事实,返回 (拒绝消息列表, 拒绝事实列表, 是否有待审, 是否有通过)。"""
  285. messages: list[str] = []
  286. facts: list[RejectionFact] = []
  287. has_pending = False
  288. has_approved = False
  289. for site in result.get("site_set_result_list") or []:
  290. if not isinstance(site, dict):
  291. continue
  292. status = str(site.get("system_status") or site.get("review_status") or "")
  293. has_pending = has_pending or _status_is_pending(status)
  294. has_approved = has_approved or _status_is_approved(status)
  295. reject_message = site.get("reject_message") or ""
  296. if _status_is_rejected(status) or reject_message:
  297. if reject_message:
  298. messages.append(str(reject_message))
  299. facts.append(RejectionFact(
  300. fact_type="site",
  301. fact_key=str(site.get("site_set") or ""),
  302. site_set=str(site.get("site_set") or ""),
  303. reason=str(reject_message),
  304. raw_detail=site,
  305. ))
  306. for detail in site.get("element_reject_detail_info") or []:
  307. if not isinstance(detail, dict):
  308. continue
  309. reason = detail.get("reason") or ""
  310. if reason:
  311. messages.append(str(reason))
  312. facts.append(RejectionFact(
  313. fact_type="element",
  314. fact_key=str(detail.get("element_name") or ""),
  315. site_set=str(site.get("site_set") or ""),
  316. reason=str(reason),
  317. element_type=str(detail.get("element_type") or ""),
  318. component_type=_component_type(detail),
  319. raw_detail={**detail, "site_set": site.get("site_set")},
  320. ))
  321. return messages, facts, has_pending, has_approved
  322. def _collect_compose_facts(result: dict) -> tuple[list[str], list[RejectionFact]]:
  323. """从审核结果中提取组件组合级拒绝事实,返回 (拒绝消息列表, 拒绝事实列表)。"""
  324. messages: list[str] = []
  325. facts: list[RejectionFact] = []
  326. for item in result.get("reject_component_compose_info_list") or []:
  327. if not isinstance(item, dict):
  328. continue
  329. reason = item.get("reject_message") or ""
  330. if reason:
  331. messages.append(str(reason))
  332. facts.append(RejectionFact(
  333. fact_type="compose",
  334. fact_key="component_compose",
  335. reason=str(reason),
  336. raw_detail=item,
  337. ))
  338. return messages, facts
  339. def parse_review_result(result: dict) -> ParsedReviewResult:
  340. """将腾讯动态创意审核结果解析为稳定的内部状态。"""
  341. reject_messages: list[str] = []
  342. delay_messages = [str(v) for v in (result.get("delay_message_list") or []) if v]
  343. facts: list[RejectionFact] = []
  344. has_pending = bool(result.get("is_all_component_compose_pending"))
  345. has_approved = False
  346. reject_messages.extend(str(v) for v in (result.get("reject_message_list") or []) if v)
  347. if result.get("reject_component_compose_count"):
  348. reject_messages.append("组件组合审核拒绝")
  349. if result.get("pass_component_compose_count"):
  350. has_approved = True
  351. if result.get("total_component_compose_count") and result.get("reject_component_compose_count") == 0:
  352. has_approved = True
  353. for collector in (_collect_site_facts, _collect_element_facts):
  354. messages, new_facts, pending, approved = collector(result)
  355. reject_messages.extend(messages)
  356. facts.extend(new_facts)
  357. has_pending = has_pending or pending
  358. has_approved = has_approved or approved
  359. messages, new_facts = _collect_compose_facts(result)
  360. reject_messages.extend(messages)
  361. facts.extend(new_facts)
  362. reject_messages = _dedupe(reject_messages)
  363. delay_messages = _dedupe(delay_messages)
  364. if reject_messages or facts or int(result.get("reject_component_compose_count") or 0) > 0:
  365. status = "rejected"
  366. elif has_pending or delay_messages:
  367. status = "pending"
  368. elif has_approved:
  369. status = "approved"
  370. else:
  371. status = "unknown"
  372. dynamic_creative_id = result.get("dynamic_creative_id")
  373. try:
  374. dynamic_creative_id = int(dynamic_creative_id) if dynamic_creative_id is not None else None
  375. except (TypeError, ValueError):
  376. dynamic_creative_id = None
  377. return ParsedReviewResult(
  378. dynamic_creative_id=dynamic_creative_id,
  379. review_status=status,
  380. reject_messages=reject_messages,
  381. delay_messages=delay_messages,
  382. rejection_facts=facts,
  383. )
  384. def group_review_tasks(
  385. tasks: Iterable[dict],
  386. batch_size: int = 100,
  387. ) -> Iterator[tuple[int, list[int]]]:
  388. """按腾讯 API 限制分批生成账户 ID 与动态创意 ID 列表。"""
  389. by_account: dict[int, list[int]] = {}
  390. for task in tasks:
  391. try:
  392. account_id = int(task["account_id"])
  393. creative_id = int(task["dynamic_creative_id"])
  394. except (KeyError, TypeError, ValueError):
  395. continue
  396. by_account.setdefault(account_id, []).append(creative_id)
  397. for account_id in sorted(by_account):
  398. ids = by_account[account_id]
  399. for i in range(0, len(ids), batch_size):
  400. yield account_id, ids[i:i + batch_size]
  401. def ensure_review_tables() -> None:
  402. """确保审核相关的数据库表(creative_creation_task / creative_review_result / creative_rejection_fact)已创建。"""
  403. from db.connection import get_connection
  404. conn = get_connection()
  405. try:
  406. with conn.cursor() as cur:
  407. for sql in CREATE_REVIEW_TABLES_SQL:
  408. cur.execute(sql)
  409. finally:
  410. conn.close()
  411. def record_creation_submission(record: dict, dynamic_creative_id: int) -> None:
  412. """将成功提交的创意保存为待审核扫描状态。"""
  413. ensure_review_tables()
  414. from db.connection import get_connection
  415. body = record.get("_request_body") or {}
  416. raw_record = json.dumps(record, ensure_ascii=False, default=str)
  417. conn = get_connection()
  418. try:
  419. with conn.cursor() as cur:
  420. cur.execute(
  421. """
  422. INSERT INTO creative_creation_task
  423. (account_id, adgroup_id, dynamic_creative_id, dynamic_creative_name,
  424. landing_video_id, material_id, material_image_id,
  425. review_status, submit_status, submitted_at, raw_record)
  426. VALUES (%s, %s, %s, %s, %s, %s, %s, 'submitted', 'submitted', NOW(), %s)
  427. ON DUPLICATE KEY UPDATE
  428. adgroup_id=VALUES(adgroup_id),
  429. dynamic_creative_name=VALUES(dynamic_creative_name),
  430. landing_video_id=VALUES(landing_video_id),
  431. material_id=VALUES(material_id),
  432. material_image_id=VALUES(material_image_id),
  433. submit_status='submitted',
  434. review_status=IF(review_status IN ('approved','rejected'), review_status, 'submitted'),
  435. raw_record=VALUES(raw_record),
  436. updated_at=CURRENT_TIMESTAMP
  437. """,
  438. (
  439. int(record["account_id"]),
  440. int(record.get("adgroup_id") or body.get("adgroup_id") or 0) or None,
  441. int(dynamic_creative_id),
  442. record.get("creative_name") or body.get("dynamic_creative_name"),
  443. record.get("landing_video_id"),
  444. str(record.get("_material_id") or ""),
  445. str(record.get("_material_image_id") or ""),
  446. raw_record,
  447. ),
  448. )
  449. finally:
  450. conn.close()
  451. def mark_creation_submit_failed(record: dict, error: str) -> None:
  452. """持久化提交失败信息,供任务追踪与诊断。"""
  453. ensure_review_tables()
  454. from db.connection import get_connection
  455. body = record.get("_request_body") or {}
  456. raw_record = json.dumps(record, ensure_ascii=False, default=str)
  457. conn = get_connection()
  458. try:
  459. with conn.cursor() as cur:
  460. cur.execute(
  461. """
  462. INSERT INTO creative_creation_task
  463. (account_id, adgroup_id, dynamic_creative_id, dynamic_creative_name,
  464. landing_video_id, material_id, material_image_id,
  465. review_status, submit_status, error, raw_record)
  466. VALUES (%s, %s, 0, %s, %s, %s, %s, 'submit_failed', 'failed', %s, %s)
  467. """,
  468. (
  469. int(record["account_id"]),
  470. int(record.get("adgroup_id") or body.get("adgroup_id") or 0) or None,
  471. record.get("creative_name") or body.get("dynamic_creative_name"),
  472. record.get("landing_video_id"),
  473. str(record.get("_material_id") or ""),
  474. str(record.get("_material_image_id") or ""),
  475. error[:2000],
  476. raw_record,
  477. ),
  478. )
  479. except Exception as e:
  480. logger.warning("[creative_review] 记录提交失败状态异常:%s", e)
  481. finally:
  482. conn.close()
  483. def load_pending_review_tasks(lookback_hours: int = 72, limit: int = 1000) -> list[dict]:
  484. """从数据库中加载待审核的创意任务列表(最近 N 小时内提交且状态为 submitted/pending/unknown)。"""
  485. ensure_review_tables()
  486. from db.connection import get_connection
  487. conn = get_connection()
  488. try:
  489. with conn.cursor() as cur:
  490. cur.execute(
  491. """
  492. SELECT account_id, dynamic_creative_id
  493. FROM creative_creation_task
  494. WHERE dynamic_creative_id > 0
  495. AND review_status IN ('submitted','pending','unknown')
  496. AND submitted_at >= DATE_SUB(NOW(), INTERVAL %s HOUR)
  497. ORDER BY COALESCE(last_review_checked_at, submitted_at) ASC
  498. LIMIT %s
  499. """,
  500. (int(lookback_hours), int(limit)),
  501. )
  502. return list(cur.fetchall())
  503. finally:
  504. conn.close()
  505. def fetch_dynamic_creative_review_results(
  506. account_id: int,
  507. creative_ids: list[int],
  508. ) -> list[dict]:
  509. """调腾讯 API 批量查询动态创意的审核结果。"""
  510. if not creative_ids:
  511. return []
  512. from tools.ad_api import _check, _get
  513. resp = _get(
  514. "/dynamic_creative_review_results/get",
  515. {
  516. "account_id": account_id,
  517. "dynamic_creative_id_list": [int(v) for v in creative_ids],
  518. },
  519. )
  520. data = _check(resp, "dynamic_creative_review_results/get")
  521. return data.get("list") or []
  522. def save_review_result(account_id: int, raw_result: dict) -> ParsedReviewResult:
  523. """解析并持久化一条创意的审核结果,同步更新 creation_task 状态与 rejection_fact 明细。"""
  524. ensure_review_tables()
  525. parsed = parse_review_result(raw_result)
  526. if parsed.dynamic_creative_id is None:
  527. raise ValueError("审核结果缺 dynamic_creative_id")
  528. raw_json = json.dumps(raw_result, ensure_ascii=False, default=str)
  529. reject_json = json.dumps(parsed.reject_messages, ensure_ascii=False)
  530. delay_json = json.dumps(parsed.delay_messages, ensure_ascii=False)
  531. finished_at_sql = "NOW()" if parsed.review_status in FINAL_REVIEW_STATUSES else "NULL"
  532. from db.connection import get_connection
  533. conn = get_connection()
  534. try:
  535. with conn.cursor() as cur:
  536. cur.execute(
  537. """
  538. INSERT INTO creative_review_result
  539. (account_id, dynamic_creative_id, review_status, reject_messages, delay_messages, raw_result, checked_at)
  540. VALUES (%s, %s, %s, %s, %s, %s, NOW())
  541. ON DUPLICATE KEY UPDATE
  542. review_status=VALUES(review_status),
  543. reject_messages=VALUES(reject_messages),
  544. delay_messages=VALUES(delay_messages),
  545. raw_result=VALUES(raw_result),
  546. checked_at=NOW(),
  547. updated_at=CURRENT_TIMESTAMP
  548. """,
  549. (
  550. int(account_id),
  551. int(parsed.dynamic_creative_id),
  552. parsed.review_status,
  553. reject_json,
  554. delay_json,
  555. raw_json,
  556. ),
  557. )
  558. cur.execute(
  559. f"""
  560. UPDATE creative_creation_task
  561. SET review_status=%s,
  562. last_review_checked_at=NOW(),
  563. review_finished_at={finished_at_sql},
  564. updated_at=CURRENT_TIMESTAMP
  565. WHERE account_id=%s AND dynamic_creative_id=%s
  566. """,
  567. (parsed.review_status, int(account_id), int(parsed.dynamic_creative_id)),
  568. )
  569. cur.execute(
  570. """
  571. DELETE FROM creative_rejection_fact
  572. WHERE account_id=%s AND dynamic_creative_id=%s
  573. """,
  574. (int(account_id), int(parsed.dynamic_creative_id)),
  575. )
  576. if parsed.rejection_facts:
  577. cur.executemany(
  578. """
  579. INSERT INTO creative_rejection_fact
  580. (account_id, dynamic_creative_id, fact_type, fact_key, reason,
  581. site_set, element_type, component_type, raw_detail)
  582. VALUES (%s, %s, %s, %s, %s, %s, %s, %s, %s)
  583. """,
  584. [
  585. (
  586. int(account_id),
  587. int(parsed.dynamic_creative_id),
  588. fact.fact_type,
  589. fact.fact_key[:255] if fact.fact_key else None,
  590. fact.reason,
  591. fact.site_set or None,
  592. fact.element_type or None,
  593. fact.component_type or None,
  594. json.dumps(fact.raw_detail, ensure_ascii=False, default=str),
  595. )
  596. for fact in parsed.rejection_facts
  597. ],
  598. )
  599. finally:
  600. conn.close()
  601. return parsed
  602. def scan_pending_reviews(
  603. lookback_hours: int = 72,
  604. limit: int = 1000,
  605. batch_size: int = 100,
  606. ) -> dict:
  607. """扫描待审核的已提交创意,并保存官方审核结果。"""
  608. tasks = load_pending_review_tasks(lookback_hours=lookback_hours, limit=limit)
  609. from tools.ad_api import prefetch_access_tokens
  610. prefetched_tokens = prefetch_access_tokens(
  611. int(task["account_id"]) for task in tasks
  612. )
  613. summary = {
  614. "run_started": datetime.now(timezone.utc).isoformat(),
  615. "tasks": len(tasks),
  616. "tokens_prefetched": len(prefetched_tokens),
  617. "batches": 0,
  618. "results": 0,
  619. "approved": 0,
  620. "rejected": 0,
  621. "pending": 0,
  622. "unknown": 0,
  623. "errors": [],
  624. }
  625. for account_id, creative_ids in group_review_tasks(tasks, batch_size=batch_size):
  626. summary["batches"] += 1
  627. try:
  628. results = fetch_dynamic_creative_review_results(account_id, creative_ids)
  629. except Exception as e:
  630. err = f"account={account_id} ids={len(creative_ids)} error={e}"
  631. logger.exception("[creative_review] 查询审核结果失败:%s", err)
  632. summary["errors"].append(err)
  633. continue
  634. for raw in results:
  635. try:
  636. parsed = save_review_result(account_id, raw)
  637. except Exception as e:
  638. err = f"account={account_id} creative={raw.get('dynamic_creative_id')} error={e}"
  639. logger.exception("[creative_review] 保存审核结果失败:%s", err)
  640. summary["errors"].append(err)
  641. continue
  642. summary["results"] += 1
  643. summary[parsed.review_status] = summary.get(parsed.review_status, 0) + 1
  644. logger.info(
  645. "[creative_review] account=%d creative=%s status=%s reject=%d delay=%d",
  646. account_id,
  647. parsed.dynamic_creative_id,
  648. parsed.review_status,
  649. len(parsed.reject_messages),
  650. len(parsed.delay_messages),
  651. )
  652. summary["run_finished"] = datetime.now(timezone.utc).isoformat()
  653. return summary