creative_review.py 28 KB

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