external_recalled_material.py 23 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670
  1. """模块 B 的外部召回素材处理。
  2. 仅对配置 ``material_source=external_recall`` 的账户启用此链路;历史召回与
  3. 既有文生图流程保持不变。
  4. """
  5. from __future__ import annotations
  6. import hashlib
  7. import json
  8. import logging
  9. import os
  10. from dataclasses import dataclass
  11. from datetime import datetime, timedelta
  12. from pathlib import Path
  13. from typing import Any, Iterable
  14. from tools.ai_generated_material import (
  15. OPENROUTER_IMAGE_MODEL,
  16. build_ai_image_object_key,
  17. generate_image_bytes,
  18. upload_image_to_oss,
  19. )
  20. from tools.ai_material_review import MaterialReviewResult, review_generated_material
  21. from tools.material_recall import Material, recall_materials_for_video
  22. from tools.video_recall import LandingVideo
  23. logger = logging.getLogger(__name__)
  24. EXTERNAL_RECALL_SOURCE_LABEL = os.getenv(
  25. "EXTERNAL_RECALL_SOURCE_LABEL", "外部合作"
  26. ).strip()
  27. EXTERNAL_RECALL_CANDIDATE_LIMIT = int(
  28. os.getenv("EXTERNAL_RECALL_CANDIDATE_LIMIT", "300")
  29. )
  30. EXTERNAL_RECALL_EDIT_LIMIT_PER_LANDING = int(
  31. os.getenv("EXTERNAL_RECALL_EDIT_LIMIT_PER_LANDING", "3")
  32. )
  33. EXTERNAL_RECALL_UV_WINDOW_DAYS = int(
  34. os.getenv("EXTERNAL_RECALL_UV_WINDOW_DAYS", "90")
  35. )
  36. EXTERNAL_RECALL_MIN_UV = int(
  37. os.getenv("EXTERNAL_RECALL_MIN_UV", "500")
  38. )
  39. EXTERNAL_IMAGE_MODEL = os.getenv(
  40. "EXTERNAL_IMAGE_MODEL", OPENROUTER_IMAGE_MODEL
  41. ).strip()
  42. EXTERNAL_CLEANUP_PROMPT_PATH = (
  43. Path(__file__).resolve().parents[1] / "prompts" / "external_material_cleanup.md"
  44. )
  45. CREATE_EXTERNAL_MATERIAL_TABLE_SQL = """
  46. CREATE TABLE IF NOT EXISTS external_recalled_material (
  47. id BIGINT AUTO_INCREMENT PRIMARY KEY COMMENT '自增主键',
  48. account_id BIGINT NOT NULL COMMENT '腾讯广告账户ID',
  49. adgroup_id BIGINT NOT NULL COMMENT '广告ID',
  50. crowd_package VARCHAR(200) NOT NULL COMMENT '投放人群包名称',
  51. landing_video_id BIGINT NOT NULL COMMENT '承接视频ID',
  52. source_material_id VARCHAR(255) NOT NULL COMMENT '外部召回materialId/card_cover_id',
  53. source_image_url VARCHAR(1000) NOT NULL COMMENT '外部素材原图URL',
  54. source_title VARCHAR(1000) DEFAULT NULL COMMENT '外部素材标题',
  55. similarity_score DOUBLE NOT NULL DEFAULT 0 COMMENT '向量召回相似度',
  56. visit_uv_30d BIGINT NOT NULL DEFAULT 0 COMMENT '配置窗口内访问UV之和(兼容历史字段名)',
  57. uv_window_start VARCHAR(8) DEFAULT NULL COMMENT 'UV窗口开始分区',
  58. uv_window_end VARCHAR(8) DEFAULT NULL COMMENT 'UV窗口结束分区',
  59. recall_json MEDIUMTEXT DEFAULT NULL COMMENT '召回命中与原始素材JSON',
  60. edit_prompt MEDIUMTEXT NOT NULL COMMENT '去播放按钮图生图prompt',
  61. model VARCHAR(200) NOT NULL COMMENT '图生图模型',
  62. oss_object_key VARCHAR(500) DEFAULT NULL COMMENT '派生图OSS key',
  63. oss_url VARCHAR(1000) DEFAULT NULL COMMENT '派生图公网URL',
  64. status VARCHAR(50) NOT NULL DEFAULT 'generated' COMMENT 'generated/prepared/approved/rejected/hold/skip/posted_ok/post_failed/error',
  65. approval_status VARCHAR(50) DEFAULT NULL COMMENT '人工审批状态',
  66. tencent_image_id VARCHAR(100) DEFAULT NULL COMMENT '腾讯图片ID',
  67. dynamic_creative_id BIGINT DEFAULT NULL COMMENT '腾讯动态创意ID',
  68. ai_review_status VARCHAR(50) DEFAULT NULL COMMENT 'AI审核状态',
  69. ai_review_score INT DEFAULT NULL COMMENT 'AI审核评分',
  70. ai_review_reason TEXT DEFAULT NULL COMMENT 'AI审核原因',
  71. ai_review_json MEDIUMTEXT DEFAULT NULL COMMENT 'AI审核JSON',
  72. error TEXT DEFAULT NULL COMMENT '处理错误',
  73. raw_response MEDIUMTEXT DEFAULT NULL COMMENT '图生图模型响应摘要',
  74. created_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP COMMENT '创建时间',
  75. updated_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP COMMENT '更新时间',
  76. UNIQUE KEY uk_account_ad_landing_source
  77. (account_id, adgroup_id, landing_video_id, source_material_id),
  78. KEY idx_account_ad_status (account_id, adgroup_id, status),
  79. KEY idx_source_created (source_material_id, created_at),
  80. KEY idx_status_created (status, created_at)
  81. ) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COMMENT='外部召回素材AI清理派生图'
  82. """
  83. @dataclass(frozen=True)
  84. class ExternalMaterialAsset:
  85. id: int
  86. source_material_id: str
  87. source_image_url: str
  88. source_title: str
  89. similarity_score: float
  90. visit_uv_30d: int
  91. uv_window_start: str
  92. uv_window_end: str
  93. oss_url: str
  94. recall_json: dict[str, Any]
  95. def to_material(self) -> Material:
  96. recall = self.recall_json
  97. return Material(
  98. material_id=f"external:{self.source_material_id}",
  99. score=self.similarity_score,
  100. title=self.source_title,
  101. cover=self.oss_url,
  102. cost=None,
  103. ctr=None,
  104. cvr=None,
  105. roi=None,
  106. impressions=None,
  107. quality_score=None,
  108. visit_uv_30d=self.visit_uv_30d,
  109. recall_strategy=str(recall.get("recall_strategy") or ""),
  110. recall_query_text=str(recall.get("recall_query_text") or ""),
  111. recall_config_code=str(recall.get("recall_config_code") or ""),
  112. recall_element_dimension=str(
  113. recall.get("recall_element_dimension") or ""
  114. ),
  115. recall_point_type=str(recall.get("recall_point_type") or ""),
  116. recall_standard_element=str(
  117. recall.get("recall_standard_element") or ""
  118. ),
  119. recall_hit_queries=list(recall.get("recall_hit_queries") or []),
  120. raw={
  121. "external_recalled_material_id": self.id,
  122. "source_material_id": self.source_material_id,
  123. "source_image_url": self.source_image_url,
  124. "visit_uv_30d": self.visit_uv_30d,
  125. "uv_window_start": self.uv_window_start,
  126. "uv_window_end": self.uv_window_end,
  127. "oss_url": self.oss_url,
  128. },
  129. )
  130. def ensure_external_material_table() -> None:
  131. from db.connection import get_connection
  132. conn = get_connection()
  133. try:
  134. with conn.cursor() as cur:
  135. cur.execute(CREATE_EXTERNAL_MATERIAL_TABLE_SQL)
  136. conn.commit()
  137. finally:
  138. conn.close()
  139. def _source_image_url(material: Material) -> str:
  140. if material.cover:
  141. return material.cover
  142. images = (material.raw or {}).get("imageList") or []
  143. return str(images[0]) if images else ""
  144. def _quote_sql(value: str) -> str:
  145. return "'" + value.replace("\\", "\\\\").replace("'", "\\'") + "'"
  146. def _get_odps_client():
  147. from tools.odps_module import get_odps_client
  148. client = get_odps_client(project="loghubods")
  149. if client is None:
  150. raise RuntimeError("无法初始化ODPS客户端")
  151. return client
  152. def _latest_partition_dt(client) -> str:
  153. table = client._odps.get_table("cooperate_top_cards_daily_v1")
  154. partition = table.get_max_partition()
  155. if partition is None:
  156. raise RuntimeError("cooperate_top_cards_daily_v1没有可用分区")
  157. name = getattr(partition, "name", "") or str(partition)
  158. for part in name.split(","):
  159. if part.strip().startswith("dt="):
  160. return part.split("=", 1)[1].strip(" '\"")
  161. raise RuntimeError(f"无法解析cooperate_top_cards_daily_v1最新分区:{name}")
  162. def _query_visit_uv_30d(
  163. material_ids: Iterable[str],
  164. ) -> tuple[dict[str, int], str, str]:
  165. ids = sorted({str(value).strip() for value in material_ids if str(value).strip()})
  166. if not ids:
  167. return {}, "", ""
  168. client = _get_odps_client()
  169. end_dt = _latest_partition_dt(client)
  170. end_date = datetime.strptime(end_dt, "%Y%m%d").date()
  171. start_dt = (end_date - timedelta(days=EXTERNAL_RECALL_UV_WINDOW_DAYS - 1)).strftime(
  172. "%Y%m%d"
  173. )
  174. totals: dict[str, int] = {}
  175. for offset in range(0, len(ids), 500):
  176. batch = ids[offset:offset + 500]
  177. quoted = ",".join(_quote_sql(value) for value in batch)
  178. sql = f"""
  179. SELECT card_cover_id, SUM(`访问uv`) AS visit_uv_30d
  180. FROM loghubods.cooperate_top_cards_daily_v1
  181. WHERE dt BETWEEN '{start_dt}' AND '{end_dt}'
  182. AND card_cover_id IN ({quoted})
  183. GROUP BY card_cover_id
  184. """
  185. instance = client._odps.execute_sql(sql)
  186. instance.wait_for_success()
  187. with instance.open_reader(tunnel=False) as reader:
  188. for row in reader:
  189. totals[str(row[0])] = int(row[1] or 0)
  190. return totals, start_dt, end_dt
  191. def rank_external_materials_by_visit_uv(
  192. materials: list[Material],
  193. min_uv: int | None = None,
  194. ) -> list[Material]:
  195. """补充配置窗口访问 UV,按 min_uv 过滤,再按 UV 和相似度排序。"""
  196. if not materials:
  197. return []
  198. totals, start_dt, end_dt = _query_visit_uv_30d(
  199. material.material_id for material in materials
  200. )
  201. for material in materials:
  202. material.visit_uv_30d = totals.get(material.material_id, 0)
  203. material.raw["visit_uv_30d"] = material.visit_uv_30d
  204. material.raw["uv_window_start"] = start_dt
  205. material.raw["uv_window_end"] = end_dt
  206. before = len(materials)
  207. if min_uv is not None:
  208. materials = [
  209. m for m in materials if (m.visit_uv_30d or 0) >= min_uv
  210. ]
  211. dropped_uv = before - len(materials)
  212. ranked = sorted(
  213. materials,
  214. key=lambda material: (material.visit_uv_30d or 0, material.score or 0),
  215. reverse=True,
  216. )
  217. logger.info(
  218. "[external_material] UV排序 candidates=%d dropped_uv=%d min_uv=%s days=%d window=%s-%s top=%s",
  219. before, dropped_uv, min_uv, EXTERNAL_RECALL_UV_WINDOW_DAYS, start_dt, end_dt,
  220. [
  221. (material.material_id, material.visit_uv_30d, round(material.score, 4))
  222. for material in ranked[:5]
  223. ],
  224. )
  225. return ranked
  226. def recall_external_materials_for_video(
  227. landing: LandingVideo,
  228. *,
  229. element_features: Iterable,
  230. ) -> list[Material]:
  231. # 外部合作素材绕过历史成本门槛(apply_cost_filter=False):
  232. # 素材多为合作方提供、无本账户历史投放成本,以访问 UV 门槛把关。
  233. materials = recall_materials_for_video(
  234. landing,
  235. final_top_n=None,
  236. source_labels=[EXTERNAL_RECALL_SOURCE_LABEL],
  237. element_features=element_features,
  238. apply_cover_blacklist=False,
  239. apply_cost_filter=False,
  240. )
  241. ranked = rank_external_materials_by_visit_uv(
  242. materials,
  243. min_uv=EXTERNAL_RECALL_MIN_UV,
  244. )
  245. return ranked[:max(1, EXTERNAL_RECALL_CANDIDATE_LIMIT)]
  246. def _load_cleanup_prompt() -> str:
  247. text = EXTERNAL_CLEANUP_PROMPT_PATH.read_text(encoding="utf-8").strip()
  248. if not text:
  249. raise RuntimeError(f"外部素材清理prompt为空:{EXTERNAL_CLEANUP_PROMPT_PATH}")
  250. return text
  251. def _recall_json(material: Material) -> dict[str, Any]:
  252. return {
  253. "recall_strategy": material.recall_strategy,
  254. "recall_query_text": material.recall_query_text,
  255. "recall_config_code": material.recall_config_code,
  256. "recall_element_dimension": material.recall_element_dimension,
  257. "recall_point_type": material.recall_point_type,
  258. "recall_standard_element": material.recall_standard_element,
  259. "recall_hit_queries": material.recall_hit_queries,
  260. "source_raw": material.raw,
  261. }
  262. def _asset_from_row(row: dict[str, Any]) -> ExternalMaterialAsset:
  263. try:
  264. recall = json.loads(row.get("recall_json") or "{}")
  265. except Exception:
  266. recall = {}
  267. return ExternalMaterialAsset(
  268. id=int(row["id"]),
  269. source_material_id=str(row["source_material_id"]),
  270. source_image_url=str(row.get("source_image_url") or ""),
  271. source_title=str(row.get("source_title") or ""),
  272. similarity_score=float(row.get("similarity_score") or 0),
  273. visit_uv_30d=int(row.get("visit_uv_30d") or 0),
  274. uv_window_start=str(row.get("uv_window_start") or ""),
  275. uv_window_end=str(row.get("uv_window_end") or ""),
  276. oss_url=str(row.get("oss_url") or ""),
  277. recall_json=recall,
  278. )
  279. def _load_existing_row(
  280. account_id: int,
  281. adgroup_id: int,
  282. landing_video_id: int,
  283. source_material_id: str,
  284. ) -> dict[str, Any] | None:
  285. ensure_external_material_table()
  286. from db.connection import get_connection
  287. conn = get_connection()
  288. try:
  289. with conn.cursor() as cur:
  290. cur.execute(
  291. """
  292. SELECT *
  293. FROM external_recalled_material
  294. WHERE account_id=%s AND adgroup_id=%s AND landing_video_id=%s
  295. AND source_material_id=%s
  296. LIMIT 1
  297. """,
  298. (account_id, adgroup_id, landing_video_id, source_material_id),
  299. )
  300. return cur.fetchone()
  301. finally:
  302. conn.close()
  303. def _insert_generated_asset(
  304. *,
  305. account_id: int,
  306. adgroup_id: int,
  307. crowd_package: str,
  308. landing: LandingVideo,
  309. material: Material,
  310. source_image_url: str,
  311. prompt: str,
  312. object_key: str,
  313. oss_url: str,
  314. raw_response: dict[str, Any],
  315. ) -> int:
  316. ensure_external_material_table()
  317. from db.connection import get_connection
  318. conn = get_connection()
  319. try:
  320. with conn.cursor() as cur:
  321. cur.execute(
  322. """
  323. INSERT INTO external_recalled_material
  324. (account_id, adgroup_id, crowd_package, landing_video_id,
  325. source_material_id, source_image_url, source_title, similarity_score,
  326. visit_uv_30d, uv_window_start, uv_window_end, recall_json,
  327. edit_prompt, model, oss_object_key, oss_url, raw_response)
  328. VALUES (%s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s)
  329. """,
  330. (
  331. account_id,
  332. adgroup_id,
  333. crowd_package,
  334. landing.video_id,
  335. material.material_id,
  336. source_image_url,
  337. material.title,
  338. material.score,
  339. int(material.visit_uv_30d or 0),
  340. str(material.raw.get("uv_window_start") or ""),
  341. str(material.raw.get("uv_window_end") or ""),
  342. json.dumps(_recall_json(material), ensure_ascii=False, default=str),
  343. prompt,
  344. EXTERNAL_IMAGE_MODEL,
  345. object_key,
  346. oss_url,
  347. json.dumps(raw_response, ensure_ascii=False, default=str)[:16000000],
  348. ),
  349. )
  350. asset_id = int(cur.lastrowid)
  351. conn.commit()
  352. return asset_id
  353. finally:
  354. conn.close()
  355. def _update_review(asset_id: int, review) -> None:
  356. from db.connection import get_connection
  357. status = "generated" if review.status == "pass" else review.status
  358. conn = get_connection()
  359. try:
  360. with conn.cursor() as cur:
  361. cur.execute(
  362. """
  363. UPDATE external_recalled_material
  364. SET status=%s, ai_review_status=%s, ai_review_score=%s,
  365. ai_review_reason=%s, ai_review_json=%s, error=%s,
  366. updated_at=CURRENT_TIMESTAMP
  367. WHERE id=%s
  368. """,
  369. (
  370. status,
  371. review.status,
  372. review.score,
  373. review.reason,
  374. json.dumps(review.raw, ensure_ascii=False, default=str),
  375. review.reason if review.status == "error" else None,
  376. asset_id,
  377. ),
  378. )
  379. conn.commit()
  380. finally:
  381. conn.close()
  382. def _review_asset(
  383. *,
  384. asset_id: int,
  385. oss_url: str,
  386. prompt: str,
  387. landing: LandingVideo,
  388. material: Material,
  389. ):
  390. try:
  391. review = review_generated_material(
  392. image_url=oss_url,
  393. prompt_type="external_cleanup",
  394. prompt_text=prompt,
  395. feature_hits=[
  396. {
  397. "landing_video_id": landing.video_id,
  398. "landing_title": landing.title,
  399. "source_material_id": material.material_id,
  400. "source_title": material.title,
  401. "visit_uv_30d": int(material.visit_uv_30d or 0),
  402. "recall_hit_queries": material.recall_hit_queries,
  403. }
  404. ],
  405. )
  406. except Exception as exc:
  407. logger.exception(
  408. "[external_material] AI审核异常 asset=%d source=%s:%s",
  409. asset_id, material.material_id, exc,
  410. )
  411. review = MaterialReviewResult(
  412. status="error",
  413. score=0,
  414. reason=f"AI审核异常:{exc}",
  415. risk_tags=["review_error"],
  416. ocr_text="",
  417. raw={"error": str(exc)},
  418. )
  419. _update_review(asset_id, review)
  420. return review
  421. def get_or_create_external_asset(
  422. *,
  423. account_id: int,
  424. adgroup_id: int,
  425. crowd_package: str,
  426. landing: LandingVideo,
  427. material: Material,
  428. ) -> Material | None:
  429. existing = _load_existing_row(
  430. account_id,
  431. adgroup_id,
  432. landing.video_id,
  433. material.material_id,
  434. )
  435. if existing:
  436. if (
  437. existing.get("oss_url")
  438. and existing.get("status") in {"generated", "error"}
  439. and existing.get("ai_review_status") in {None, "", "error"}
  440. ):
  441. review = _review_asset(
  442. asset_id=int(existing["id"]),
  443. oss_url=str(existing["oss_url"]),
  444. prompt=str(existing.get("edit_prompt") or _load_cleanup_prompt()),
  445. landing=landing,
  446. material=material,
  447. )
  448. existing = _load_existing_row(
  449. account_id,
  450. adgroup_id,
  451. landing.video_id,
  452. material.material_id,
  453. )
  454. if review.passed and existing:
  455. return _asset_from_row(existing).to_material()
  456. if (
  457. existing.get("ai_review_status") == "pass"
  458. and existing.get("status") in {"generated", "prepared"}
  459. and existing.get("oss_url")
  460. ):
  461. logger.info(
  462. "[external_material] 复用派生图 account=%d adgroup=%d landing=%d source=%s id=%s",
  463. account_id, adgroup_id, landing.video_id, material.material_id,
  464. existing["id"],
  465. )
  466. return _asset_from_row(existing).to_material()
  467. logger.info(
  468. "[external_material] 已有不可用派生记录 source=%s status=%s review=%s",
  469. material.material_id, existing.get("status"),
  470. existing.get("ai_review_status"),
  471. )
  472. return None
  473. source_image_url = _source_image_url(material)
  474. if not source_image_url:
  475. logger.warning(
  476. "[external_material] source=%s 没有可编辑图片URL", material.material_id
  477. )
  478. return None
  479. prompt = _load_cleanup_prompt()
  480. image_bytes, content_type, raw_response = generate_image_bytes(
  481. prompt,
  482. model=EXTERNAL_IMAGE_MODEL,
  483. input_image_url=source_image_url,
  484. )
  485. source_hash = hashlib.sha256(material.material_id.encode("utf-8")).hexdigest()[:12]
  486. object_key = build_ai_image_object_key(
  487. account_id=account_id,
  488. landing_video_id=landing.video_id,
  489. prompt_type=f"external_cleanup_{source_hash}",
  490. extension="jpg",
  491. )
  492. oss_url = upload_image_to_oss(image_bytes, content_type, object_key)
  493. asset_id = _insert_generated_asset(
  494. account_id=account_id,
  495. adgroup_id=adgroup_id,
  496. crowd_package=crowd_package,
  497. landing=landing,
  498. material=material,
  499. source_image_url=source_image_url,
  500. prompt=prompt,
  501. object_key=object_key,
  502. oss_url=oss_url,
  503. raw_response=raw_response,
  504. )
  505. review = _review_asset(
  506. asset_id=asset_id,
  507. oss_url=oss_url,
  508. prompt=prompt,
  509. landing=landing,
  510. material=material,
  511. )
  512. logger.info(
  513. "[external_material] 派生图完成 account=%d adgroup=%d landing=%d source=%s "
  514. "asset=%d uv_window=%d review=%s score=%d url=%s",
  515. account_id, adgroup_id, landing.video_id, material.material_id,
  516. asset_id, int(material.visit_uv_30d or 0), review.status, review.score, oss_url,
  517. )
  518. if not review.passed:
  519. return None
  520. row = _load_existing_row(
  521. account_id,
  522. adgroup_id,
  523. landing.video_id,
  524. material.material_id,
  525. )
  526. return _asset_from_row(row).to_material() if row else None
  527. def prepare_external_material_for_landing(
  528. *,
  529. account_id: int,
  530. adgroup_id: int,
  531. crowd_package: str,
  532. landing: LandingVideo,
  533. element_features: Iterable,
  534. excluded_material_ids: set[str],
  535. ) -> Material | None:
  536. ranked = recall_external_materials_for_video(
  537. landing,
  538. element_features=element_features,
  539. )
  540. attempted = 0
  541. for material in ranked:
  542. external_id = f"external:{material.material_id}"
  543. if external_id in excluded_material_ids:
  544. continue
  545. attempted += 1
  546. try:
  547. processed = get_or_create_external_asset(
  548. account_id=account_id,
  549. adgroup_id=adgroup_id,
  550. crowd_package=crowd_package,
  551. landing=landing,
  552. material=material,
  553. )
  554. except Exception as exc:
  555. logger.exception(
  556. "[external_material] 图生图失败 landing=%d source=%s:%s",
  557. landing.video_id, material.material_id, exc,
  558. )
  559. processed = None
  560. if processed is not None:
  561. return processed
  562. if attempted >= max(1, EXTERNAL_RECALL_EDIT_LIMIT_PER_LANDING):
  563. break
  564. logger.info(
  565. "[external_material] landing=%d 无可用派生素材 recalled=%d attempted=%d",
  566. landing.video_id, len(ranked), attempted,
  567. )
  568. return None
  569. def update_external_material_status(
  570. record: dict,
  571. status: str,
  572. *,
  573. dynamic_creative_id: int | str | None = None,
  574. tencent_image_id: str = "",
  575. error: str = "",
  576. ) -> None:
  577. asset_id = record.get("_external_recalled_material_id")
  578. if not asset_id:
  579. return
  580. ensure_external_material_table()
  581. from db.connection import get_connection
  582. normalized_status = {
  583. "approve": "approved",
  584. "reject": "rejected",
  585. }.get(status, status)
  586. approval_status = {
  587. "approve": "approved",
  588. "reject": "rejected",
  589. "hold": "hold",
  590. "skip": "skip",
  591. }.get(status, normalized_status)
  592. conn = get_connection()
  593. try:
  594. with conn.cursor() as cur:
  595. cur.execute(
  596. """
  597. UPDATE external_recalled_material
  598. SET status=%s, approval_status=%s,
  599. dynamic_creative_id=COALESCE(%s, dynamic_creative_id),
  600. tencent_image_id=COALESCE(NULLIF(%s, ''), tencent_image_id),
  601. error=COALESCE(NULLIF(%s, ''), error), updated_at=CURRENT_TIMESTAMP
  602. WHERE id=%s
  603. """,
  604. (
  605. normalized_status,
  606. approval_status,
  607. int(dynamic_creative_id) if dynamic_creative_id else None,
  608. tencent_image_id,
  609. error[:2000] if error else "",
  610. int(asset_id),
  611. ),
  612. )
  613. conn.commit()
  614. finally:
  615. conn.close()