| 1234567891011121314151617181920212223242526272829303132333435363738394041424344454647484950515253545556575859606162636465666768697071727374757677787980818283848586878889909192939495969798991001011021031041051061071081091101111121131141151161171181191201211221231241251261271281291301311321331341351361371381391401411421431441451461471481491501511521531541551561571581591601611621631641651661671681691701711721731741751761771781791801811821831841851861871881891901911921931941951961971981992002012022032042052062072082092102112122132142152162172182192202212222232242252262272282292302312322332342352362372382392402412422432442452462472482492502512522532542552562572582592602612622632642652662672682692702712722732742752762772782792802812822832842852862872882892902912922932942952962972982993003013023033043053063073083093103113123133143153163173183193203213223233243253263273283293303313323333343353363373383393403413423433443453463473483493503513523533543553563573583593603613623633643653663673683693703713723733743753763773783793803813823833843853863873883893903913923933943953963973983994004014024034044054064074084094104114124134144154164174184194204214224234244254264274284294304314324334344354364374384394404414424434444454464474484494504514524534544554564574584594604614624634644654664674684694704714724734744754764774784794804814824834844854864874884894904914924934944954964974984995005015025035045055065075085095105115125135145155165175185195205215225235245255265275285295305315325335345355365375385395405415425435445455465475485495505515525535545555565575585595605615625635645655665675685695705715725735745755765775785795805815825835845855865875885895905915925935945955965975985996006016026036046056066076086096106116126136146156166176186196206216226236246256266276286296306316326336346356366376386396406416426436446456466476486496506516526536546556566576586596606616626636646656666676686696706716726736746756766776786796806816826836846856866876886896906916926936946956966976986997007017027037047057067077087097107117127137147157167177187197207217227237247257267277287297307317327337347357367377387397407417427437447457467477487497507517527537547557567577587597607617627637647657667677687697707717727737747757767777787797807817827837847857867877887897907917927937947957967977987998008018028038048058068078088098108118128138148158168178188198208218228238248258268278288298308318328338348358368378388398408418428438448458468478488498508518528538548558568578588598608618628638648658668678688698708718728738748758768778788798808818828838848858868878888898908918928938948958968978988999009019029039049059069079089099109119129139149159169179189199209219229239249259269279289299309319329339349359369379389399409419429439449459469479489499509519529539549559569579589599609619629639649659669679689699709719729739749759769779789799809819829839849859869879889899909919929939949959969979989991000100110021003100410051006100710081009101010111012101310141015101610171018101910201021102210231024102510261027102810291030103110321033103410351036103710381039104010411042104310441045104610471048104910501051105210531054105510561057105810591060106110621063106410651066106710681069107010711072107310741075107610771078107910801081108210831084108510861087108810891090109110921093109410951096109710981099110011011102110311041105110611071108110911101111111211131114111511161117111811191120112111221123112411251126112711281129113011311132113311341135113611371138113911401141114211431144114511461147114811491150115111521153115411551156115711581159116011611162116311641165116611671168116911701171117211731174117511761177117811791180118111821183118411851186118711881189119011911192119311941195119611971198119912001201120212031204120512061207120812091210121112121213121412151216121712181219122012211222122312241225122612271228122912301231123212331234123512361237123812391240124112421243124412451246124712481249125012511252125312541255125612571258125912601261126212631264126512661267126812691270127112721273127412751276127712781279128012811282128312841285128612871288128912901291129212931294129512961297129812991300130113021303130413051306130713081309131013111312131313141315131613171318131913201321132213231324132513261327132813291330133113321333133413351336133713381339134013411342134313441345134613471348134913501351135213531354135513561357135813591360136113621363136413651366136713681369137013711372137313741375137613771378137913801381138213831384138513861387138813891390139113921393139413951396139713981399140014011402140314041405140614071408140914101411141214131414141514161417141814191420142114221423142414251426142714281429143014311432143314341435143614371438143914401441144214431444144514461447144814491450145114521453145414551456145714581459146014611462146314641465146614671468146914701471147214731474147514761477147814791480148114821483148414851486148714881489149014911492149314941495149614971498149915001501150215031504150515061507150815091510151115121513151415151516151715181519152015211522152315241525152615271528152915301531153215331534153515361537153815391540154115421543154415451546154715481549155015511552155315541555155615571558155915601561156215631564156515661567156815691570157115721573157415751576157715781579158015811582158315841585158615871588158915901591159215931594159515961597159815991600160116021603160416051606160716081609161016111612161316141615161616171618161916201621162216231624162516261627162816291630163116321633163416351636163716381639164016411642164316441645164616471648164916501651165216531654165516561657165816591660166116621663166416651666166716681669167016711672167316741675167616771678167916801681168216831684168516861687168816891690169116921693169416951696169716981699170017011702170317041705170617071708170917101711171217131714171517161717171817191720172117221723172417251726172717281729173017311732173317341735173617371738173917401741174217431744174517461747174817491750175117521753175417551756175717581759176017611762176317641765176617671768176917701771177217731774177517761777177817791780178117821783178417851786178717881789179017911792179317941795179617971798179918001801180218031804180518061807180818091810181118121813181418151816181718181819182018211822182318241825182618271828182918301831183218331834183518361837183818391840184118421843184418451846184718481849185018511852185318541855185618571858185918601861186218631864186518661867186818691870187118721873187418751876187718781879188018811882188318841885188618871888188918901891189218931894189518961897189818991900190119021903190419051906190719081909191019111912191319141915191619171918191919201921192219231924192519261927192819291930193119321933193419351936193719381939194019411942194319441945194619471948194919501951195219531954195519561957195819591960196119621963196419651966196719681969197019711972197319741975197619771978197919801981198219831984198519861987198819891990199119921993199419951996199719981999200020012002200320042005200620072008200920102011201220132014201520162017201820192020202120222023202420252026202720282029203020312032203320342035203620372038203920402041204220432044204520462047204820492050205120522053205420552056205720582059206020612062206320642065206620672068206920702071207220732074207520762077207820792080208120822083208420852086208720882089209020912092209320942095209620972098209921002101210221032104210521062107210821092110211121122113211421152116211721182119212021212122212321242125212621272128212921302131213221332134213521362137213821392140214121422143214421452146 |
- """Apply creative-review cleanup rules and notify agency groups."""
- from __future__ import annotations
- import hashlib
- import json
- import logging
- import os
- import threading
- from collections import defaultdict
- from concurrent.futures import ThreadPoolExecutor, as_completed
- from datetime import date, datetime, timedelta
- from pathlib import Path
- from typing import Any, Callable
- from zoneinfo import ZoneInfo
- import pandas as pd
- from openpyxl import Workbook
- from openpyxl.styles import Alignment, Font, PatternFill
- from openpyxl.utils import get_column_letter
- from db.connection import get_connection
- from roi_control.agency_delivery import publish_agency_reports, resolve_agency_webhook
- from roi_control.config import AgencyWebhookConfig
- from roi_control.feishu import RoiFeishuPublisher
- from storage import advisory_lock, initialize_schema
- from tools.creative_review import (
- fetch_dynamic_creative_review_results,
- parse_review_result,
- review_granularity_fields,
- status_desc,
- )
- logger = logging.getLogger(__name__)
- SHANGHAI = ZoneInfo("Asia/Shanghai")
- REPORT_VERSION = "creative_rejection_cleanup_v11"
- OPERATOR_SUMMARY_ROUTE = "投放调控汇总"
- DENIED_SYSTEM_STATUS = "DYNAMIC_CREATIVE_STATUS_DENIED"
- DELETED_STATUS = "AD_STATUS_DELETED"
- CREATIVE_DENIED_STATUS = "CREATIVE_SET_APPROVAL_STATUS_DENIED"
- CREATIVE_PARTIAL_NORMAL_STATUS = "CREATIVE_SET_APPROVAL_STATUS_PARTIAL_NORMAL"
- DELETE_CREATIVE = "DELETE_CREATIVE"
- ALERT_ONLY = "ALERT_ONLY"
- DEFAULT_PARTIAL_CREATIVE_COST_THRESHOLD_YUAN = 50.0
- DEFAULT_WECHAT_MINI_PROGRAM_COST_THRESHOLD_YUAN = 100.0
- DEFAULT_DELETE_CLAIM_STALE_MINUTES = 30
- AGENCY_REPORT_COLUMNS = (
- "代理名称",
- "账户ID",
- "账户名称",
- "广告ID",
- "广告名称",
- "创意ID",
- "创意名称",
- "近3天累计历史消耗(元)",
- "配置状态",
- "创意审核状态",
- "审核不通过原因",
- "执行操作",
- )
- OPERATOR_REPORT_COLUMNS = (
- "代理名称",
- "账户ID",
- "账户名称",
- "广告ID",
- "广告名称",
- "创意ID",
- "创意名称",
- "近3天累计历史消耗(元)",
- "消耗日期范围",
- "配置状态",
- "创意审核状态",
- "元素粒度审核状态",
- "元素粒度审核不通过原因",
- "版位粒度审核状态",
- "版位粒度审核不通过原因",
- "审核不通过原因",
- "检查时间",
- "执行操作",
- "操作判断原因",
- )
- # Compatibility for callers that treat the agency report as the default report.
- REPORT_COLUMNS = AGENCY_REPORT_COLUMNS
- def _json(value: Any) -> str:
- return json.dumps(value, ensure_ascii=False, default=str)
- def _is_rejected_status(value: Any) -> bool:
- upper = str(value or "").strip().upper()
- return "REJECT" in upper or "DENIED" in upper
- def has_rejected_wechat_mini_program_element(raw_result: dict | None) -> bool:
- """Return whether a rejected element is explicitly named 微信小程序."""
- raw = raw_result if isinstance(raw_result, dict) else {}
- for element in raw.get("element_result_list") or []:
- if not isinstance(element, dict):
- continue
- element_name = "".join(str(element.get("element_name") or "").split())
- if element_name != "微信小程序":
- continue
- if any(
- _is_rejected_status(element.get(field))
- for field in ("review_status", "system_status")
- ):
- return True
- return False
- def partial_creative_cost_threshold_fen() -> int:
- raw = os.getenv(
- "DAILY_PARTIAL_CREATIVE_COST_THRESHOLD_YUAN",
- str(DEFAULT_PARTIAL_CREATIVE_COST_THRESHOLD_YUAN),
- )
- try:
- yuan = float(raw)
- except (TypeError, ValueError) as exc:
- raise ValueError(
- "DAILY_PARTIAL_CREATIVE_COST_THRESHOLD_YUAN must be numeric"
- ) from exc
- if yuan < 0:
- raise ValueError(
- "DAILY_PARTIAL_CREATIVE_COST_THRESHOLD_YUAN must not be negative"
- )
- return int(round(yuan * 100))
- def wechat_mini_program_cost_threshold_fen() -> int:
- raw = os.getenv(
- "DAILY_WECHAT_MINI_PROGRAM_CREATIVE_COST_THRESHOLD_YUAN",
- str(DEFAULT_WECHAT_MINI_PROGRAM_COST_THRESHOLD_YUAN),
- )
- try:
- yuan = float(raw)
- except (TypeError, ValueError) as exc:
- raise ValueError(
- "DAILY_WECHAT_MINI_PROGRAM_CREATIVE_COST_THRESHOLD_YUAN must be numeric"
- ) from exc
- if yuan < 0:
- raise ValueError(
- "DAILY_WECHAT_MINI_PROGRAM_CREATIVE_COST_THRESHOLD_YUAN must not be negative"
- )
- return int(round(yuan * 100))
- def determine_cleanup_action(
- creative: dict[str, Any],
- raw_result: dict | None,
- *,
- recent_cost_fen: int | None = None,
- cost_threshold_fen: int = 5000,
- wechat_cost_threshold_fen: int = 10000,
- spend_error: str | None = None,
- ) -> dict[str, Any] | None:
- """Choose one of the three whole-creative delete rules or manual alert."""
- approval_status = str(creative.get("creative_set_approval_status") or "")
- if approval_status == CREATIVE_DENIED_STATUS:
- return {
- "cleanup_action": DELETE_CREATIVE,
- "component_ids": [],
- "element_ids": [],
- "action_reason": "创意审核状态为审核拒绝",
- "recent_cost_fen": recent_cost_fen,
- }
- if approval_status == CREATIVE_PARTIAL_NORMAL_STATUS:
- if has_rejected_wechat_mini_program_element(raw_result):
- if spend_error or recent_cost_fen is None:
- reason = (
- "部分投放中且微信小程序元素审核拒绝,"
- "近3天历史消耗读取失败,需人工判断"
- )
- if spend_error:
- reason = f"{reason}:{spend_error}"
- return {
- "cleanup_action": ALERT_ONLY,
- "component_ids": [],
- "element_ids": [],
- "action_reason": reason,
- "recent_cost_fen": None,
- }
- if recent_cost_fen >= wechat_cost_threshold_fen:
- return {
- "cleanup_action": ALERT_ONLY,
- "component_ids": [],
- "element_ids": [],
- "action_reason": (
- "部分投放中且微信小程序元素审核拒绝,近3天历史消耗"
- f"{recent_cost_fen / 100:.2f}元不低于"
- f"{wechat_cost_threshold_fen / 100:.2f}元,需人工判断是否删除"
- ),
- "recent_cost_fen": recent_cost_fen,
- }
- return {
- "cleanup_action": DELETE_CREATIVE,
- "component_ids": [],
- "element_ids": [],
- "action_reason": (
- "部分投放中且微信小程序元素审核拒绝,近3天历史消耗"
- f"{recent_cost_fen / 100:.2f}元低于"
- f"{wechat_cost_threshold_fen / 100:.2f}元"
- ),
- "recent_cost_fen": recent_cost_fen,
- }
- if spend_error or recent_cost_fen is None:
- reason = "部分投放中,近3天历史消耗读取失败,需人工判断"
- if spend_error:
- reason = f"{reason}:{spend_error}"
- return {
- "cleanup_action": ALERT_ONLY,
- "component_ids": [],
- "element_ids": [],
- "action_reason": reason,
- "recent_cost_fen": None,
- }
- if recent_cost_fen < cost_threshold_fen:
- return {
- "cleanup_action": DELETE_CREATIVE,
- "component_ids": [],
- "element_ids": [],
- "action_reason": (
- f"部分投放中且近3天历史消耗{recent_cost_fen / 100:.2f}元"
- f"低于{cost_threshold_fen / 100:.2f}元"
- ),
- "recent_cost_fen": recent_cost_fen,
- }
- return {
- "cleanup_action": ALERT_ONLY,
- "component_ids": [],
- "element_ids": [],
- "action_reason": (
- f"部分投放中且近3天历史消耗{recent_cost_fen / 100:.2f}元,"
- "需人工判断是否删除"
- ),
- "recent_cost_fen": recent_cost_fen,
- }
- return None
- def _as_int(value: Any) -> int | None:
- try:
- number = int(value)
- except (TypeError, ValueError):
- return None
- return number if number > 0 else None
- def _agency_name(value: Any) -> str:
- if value is None or pd.isna(value):
- return ""
- return "".join(str(value or "").split())
- def _env_flag(name: str, default: bool = False) -> bool:
- raw = os.getenv(name)
- if raw is None:
- return default
- return raw.strip().lower() in {"1", "true", "yes", "on"}
- def resolve_end_date(client, requested=None, *, now=None):
- from roi_control.data_source import resolve_end_date as resolve
- return resolve(client, requested, now=now)
- def date_window(end_date: str) -> tuple[str, str]:
- end = datetime.strptime(end_date, "%Y%m%d")
- return (end - timedelta(days=2)).strftime("%Y%m%d"), end_date
- def fetch_daily_data(client, start_date: str, end_date: str) -> pd.DataFrame:
- from roi_control.data_source import fetch_daily_data as fetch
- return fetch(client, start_date, end_date)
- def fetch_recent_spend_accounts(client, start_date: str, end_date: str) -> list[dict]:
- from roi_control.data_source import fetch_recent_spend_accounts as fetch
- return fetch(client, start_date, end_date)
- def fetch_account_agency_fallbacks(client, account_ids: list[int]) -> dict[int, str]:
- from roi_control.data_source import fetch_account_agency_fallbacks as fetch
- return fetch(client, account_ids)
- def prefetch_account_access_tokens(account_ids: list[int]) -> dict[int, str]:
- from tools.ad_api import prefetch_access_tokens
- return prefetch_access_tokens(account_ids)
- def build_agency_context(daily: pd.DataFrame) -> dict[str, dict]:
- """Build latest-date creative mappings and unambiguous account fallbacks."""
- creative_agency_sets: dict[tuple[int, int], set[str]] = defaultdict(set)
- creative_agency_dates: dict[tuple[int, int], str] = {}
- account_agencies: dict[int, set[str]] = defaultdict(set)
- account_agency_dates: dict[int, str] = {}
- account_names: dict[int, str] = {}
- account_name_dates: dict[int, str] = {}
- if daily.empty:
- return {
- "creative_agencies": {},
- "account_agencies": {},
- "fallback_account_agencies": {},
- "account_names": account_names,
- }
- miniapp = daily[daily["entity_type"].eq("self")]
- for _, row in miniapp.iterrows():
- account_id = _as_int(row.get("账号id"))
- creative_id = _as_int(row.get("创意id"))
- agency = _agency_name(row.get("代理名称"))
- row_date = str(row.get("dt") or "")
- if account_id is None:
- continue
- account_name = str(row.get("账号名称") or "").strip()
- if account_name and row_date >= account_name_dates.get(account_id, ""):
- account_names[account_id] = account_name
- account_name_dates[account_id] = row_date
- if agency:
- account_date = account_agency_dates.get(account_id, "")
- if row_date > account_date:
- account_agencies[account_id] = {agency}
- account_agency_dates[account_id] = row_date
- elif row_date == account_date:
- account_agencies[account_id].add(agency)
- if creative_id is not None:
- key = (account_id, creative_id)
- current_date = creative_agency_dates.get(key, "")
- if row_date > current_date:
- creative_agency_sets[key] = {agency}
- creative_agency_dates[key] = row_date
- elif row_date == current_date:
- creative_agency_sets[key].add(agency)
- return {
- "creative_agencies": {
- key: next(iter(agencies))
- for key, agencies in creative_agency_sets.items()
- if len(agencies) == 1
- },
- "account_agencies": {
- account_id: next(iter(agencies))
- for account_id, agencies in account_agencies.items()
- if len(agencies) == 1
- },
- "fallback_account_agencies": {},
- "account_names": account_names,
- }
- def _resolve_agency(
- context: dict[str, dict],
- account_id: int,
- creative_id: int,
- ) -> str:
- return str(
- context["creative_agencies"].get((account_id, creative_id))
- or context["account_agencies"].get(account_id)
- or context.get("fallback_account_agencies", {}).get(account_id)
- or ""
- )
- def upsert_cleanup_candidate(record: dict[str, Any]) -> dict[str, Any]:
- component_ids_json = _json(record.get("component_ids") or [])
- element_ids_json = _json(record.get("element_ids") or [])
- review_result_json = _json(record.get("review_result") or {})
- pre_state_json = _json(record.get("pre_state") or {})
- cleanup_status = (
- "ALERT_PENDING"
- if record["cleanup_action"] == ALERT_ONLY
- else "DISCOVERED"
- )
- insert_values = (
- record["account_id"],
- record.get("account_name"),
- record.get("agency_name"),
- record["adgroup_id"],
- record.get("adgroup_name"),
- record["dynamic_creative_id"],
- record.get("dynamic_creative_name"),
- record["check_date"],
- record["cleanup_action"],
- component_ids_json,
- element_ids_json,
- record.get("recent_cost_fen"),
- record.get("cost_start_date"),
- record.get("cost_end_date"),
- record.get("action_reason") or record["reject_reason"],
- record["reject_reason"],
- review_result_json,
- pre_state_json,
- cleanup_status,
- )
- connection = get_connection()
- try:
- with connection.cursor() as cursor:
- cursor.execute(
- """
- INSERT INTO creative_rejection_cleanup_item
- (account_id, account_name, agency_name, adgroup_id,
- adgroup_name, dynamic_creative_id, dynamic_creative_name,
- check_date,
- cleanup_action, target_component_ids_json,
- target_element_ids_json, recent_cost_fen,
- cost_start_date, cost_end_date, action_reason, reject_reason,
- review_result_json, pre_state_json, cleanup_status)
- VALUES (%s,%s,%s,%s,%s,%s,%s,%s,%s,%s,%s,%s,%s,%s,%s,%s,%s,%s,%s)
- ON DUPLICATE KEY UPDATE id=LAST_INSERT_ID(id)
- """,
- insert_values,
- )
- cursor.execute(
- """
- UPDATE creative_rejection_cleanup_item AS item
- JOIN (
- SELECT
- %s AS account_id, %s AS account_name, %s AS agency_name,
- %s AS adgroup_id, %s AS adgroup_name,
- %s AS dynamic_creative_id, %s AS dynamic_creative_name,
- %s AS check_date, %s AS cleanup_action,
- %s AS target_component_ids_json,
- %s AS target_element_ids_json, %s AS recent_cost_fen,
- %s AS cost_start_date, %s AS cost_end_date,
- %s AS action_reason, %s AS reject_reason,
- %s AS review_result_json, %s AS pre_state_json,
- %s AS cleanup_status
- ) AS incoming
- ON incoming.account_id=item.account_id
- AND incoming.dynamic_creative_id=item.dynamic_creative_id
- AND incoming.check_date=item.check_date
- SET
- item.account_name=COALESCE(
- NULLIF(incoming.account_name,''), item.account_name
- ),
- item.agency_name=COALESCE(
- NULLIF(incoming.agency_name,''), item.agency_name
- ),
- item.adgroup_id=incoming.adgroup_id,
- item.adgroup_name=COALESCE(
- NULLIF(incoming.adgroup_name,''), item.adgroup_name
- ),
- item.dynamic_creative_name=COALESCE(
- NULLIF(incoming.dynamic_creative_name,''),
- item.dynamic_creative_name
- ),
- item.action_reason=incoming.action_reason,
- item.reject_reason=incoming.reject_reason,
- item.review_result_json=incoming.review_result_json,
- item.pre_state_json=incoming.pre_state_json,
- notified_at=CASE
- WHEN NOT (item.cleanup_action <=> incoming.cleanup_action)
- OR NOT (item.target_component_ids_json <=> incoming.target_component_ids_json)
- OR NOT (item.target_element_ids_json <=> incoming.target_element_ids_json)
- OR NOT (item.recent_cost_fen <=> incoming.recent_cost_fen)
- OR NOT (item.cost_end_date <=> incoming.cost_end_date)
- THEN NULL ELSE item.notified_at
- END,
- item.agency_notified_at=CASE
- WHEN NOT (item.cleanup_action <=> incoming.cleanup_action)
- OR NOT (item.target_component_ids_json <=> incoming.target_component_ids_json)
- OR NOT (item.target_element_ids_json <=> incoming.target_element_ids_json)
- OR NOT (item.recent_cost_fen <=> incoming.recent_cost_fen)
- OR NOT (item.cost_end_date <=> incoming.cost_end_date)
- THEN NULL ELSE item.agency_notified_at
- END,
- item.operator_notified_at=CASE
- WHEN NOT (item.cleanup_action <=> incoming.cleanup_action)
- OR NOT (item.target_component_ids_json <=> incoming.target_component_ids_json)
- OR NOT (item.target_element_ids_json <=> incoming.target_element_ids_json)
- OR NOT (item.recent_cost_fen <=> incoming.recent_cost_fen)
- OR NOT (item.cost_end_date <=> incoming.cost_end_date)
- THEN NULL ELSE item.operator_notified_at
- END,
- item.deleted_at=CASE
- WHEN NOT (item.cleanup_action <=> incoming.cleanup_action)
- OR NOT (item.target_component_ids_json <=> incoming.target_component_ids_json)
- OR NOT (item.target_element_ids_json <=> incoming.target_element_ids_json)
- OR NOT (item.recent_cost_fen <=> incoming.recent_cost_fen)
- OR NOT (item.cost_end_date <=> incoming.cost_end_date)
- THEN NULL ELSE item.deleted_at
- END,
- item.readback_json=CASE
- WHEN NOT (item.cleanup_action <=> incoming.cleanup_action)
- OR NOT (item.target_component_ids_json <=> incoming.target_component_ids_json)
- OR NOT (item.target_element_ids_json <=> incoming.target_element_ids_json)
- OR NOT (item.recent_cost_fen <=> incoming.recent_cost_fen)
- OR NOT (item.cost_end_date <=> incoming.cost_end_date)
- THEN NULL ELSE item.readback_json
- END,
- item.cleanup_status=CASE
- WHEN NOT (item.cleanup_action <=> incoming.cleanup_action)
- OR NOT (item.target_component_ids_json <=> incoming.target_component_ids_json)
- OR NOT (item.target_element_ids_json <=> incoming.target_element_ids_json)
- OR NOT (item.recent_cost_fen <=> incoming.recent_cost_fen)
- OR NOT (item.cost_end_date <=> incoming.cost_end_date)
- THEN incoming.cleanup_status
- WHEN item.cleanup_status='SKIPPED_REVIEW_NOT_RECONFIRMED'
- THEN 'DISCOVERED'
- ELSE item.cleanup_status
- END,
- item.target_component_ids_json=incoming.target_component_ids_json,
- item.target_element_ids_json=incoming.target_element_ids_json,
- item.recent_cost_fen=incoming.recent_cost_fen,
- item.cost_start_date=incoming.cost_start_date,
- item.cost_end_date=incoming.cost_end_date,
- item.cleanup_action=incoming.cleanup_action
- WHERE item.cleanup_status NOT IN ('DELETING','CREATIVE_DELETED')
- """,
- insert_values,
- )
- cursor.execute(
- """
- SELECT * FROM creative_rejection_cleanup_item
- WHERE account_id=%s AND dynamic_creative_id=%s AND check_date=%s
- """,
- (
- record["account_id"],
- record["dynamic_creative_id"],
- record["check_date"],
- ),
- )
- return cursor.fetchone()
- finally:
- connection.close()
- def load_retryable_cleanup_items() -> list[dict[str, Any]]:
- stale_minutes = int(
- os.getenv(
- "DAILY_REJECTED_CREATIVE_DELETE_CLAIM_STALE_MINUTES",
- str(DEFAULT_DELETE_CLAIM_STALE_MINUTES),
- )
- )
- if stale_minutes <= 0:
- raise ValueError(
- "DAILY_REJECTED_CREATIVE_DELETE_CLAIM_STALE_MINUTES must be positive"
- )
- connection = get_connection()
- try:
- with connection.cursor() as cursor:
- cursor.execute(
- """
- SELECT item.*
- FROM creative_rejection_cleanup_item item
- JOIN (
- SELECT account_id, dynamic_creative_id, MAX(check_date) AS check_date
- FROM creative_rejection_cleanup_item
- GROUP BY account_id, dynamic_creative_id
- ) latest
- ON latest.account_id=item.account_id
- AND latest.dynamic_creative_id=item.dynamic_creative_id
- AND latest.check_date=item.check_date
- WHERE (
- item.cleanup_status IN
- ('DISCOVERED','DEFERRED','FAILED','WRITE_OUTCOME_UNKNOWN')
- OR (
- item.cleanup_status='DELETING'
- AND item.updated_at < DATE_SUB(NOW(), INTERVAL %s MINUTE)
- )
- )
- AND item.cleanup_action='DELETE_CREATIVE'
- ORDER BY item.id
- """,
- (stale_minutes,),
- )
- return list(cursor.fetchall())
- finally:
- connection.close()
- def claim_cleanup_item(item_id: int) -> bool:
- """Atomically claim a retryable deletion row for this run."""
- stale_minutes = int(
- os.getenv(
- "DAILY_REJECTED_CREATIVE_DELETE_CLAIM_STALE_MINUTES",
- str(DEFAULT_DELETE_CLAIM_STALE_MINUTES),
- )
- )
- if stale_minutes <= 0:
- raise ValueError(
- "DAILY_REJECTED_CREATIVE_DELETE_CLAIM_STALE_MINUTES must be positive"
- )
- connection = get_connection()
- try:
- with connection.cursor() as cursor:
- cursor.execute(
- """
- UPDATE creative_rejection_cleanup_item
- SET cleanup_status='DELETING', error_message=NULL, updated_at=NOW()
- WHERE id=%s
- AND cleanup_action='DELETE_CREATIVE'
- AND (
- cleanup_status IN
- ('DISCOVERED','DEFERRED','FAILED','WRITE_OUTCOME_UNKNOWN')
- OR (
- cleanup_status='DELETING'
- AND updated_at < DATE_SUB(NOW(), INTERVAL %s MINUTE)
- )
- )
- """,
- (item_id, stale_minutes),
- )
- return cursor.rowcount == 1
- finally:
- connection.close()
- def update_cleanup_item(item_id: int, **values: Any) -> bool:
- expected_cleanup_status = values.pop("_expected_cleanup_status", None)
- allowed = {
- "agency_name",
- "cleanup_action",
- "target_component_ids_json",
- "target_element_ids_json",
- "recent_cost_fen",
- "cost_start_date",
- "cost_end_date",
- "action_reason",
- "reject_reason",
- "review_result_json",
- "cleanup_status",
- "error_message",
- "pre_state_json",
- "readback_json",
- "deleted_at",
- "agency_notified_at",
- "operator_notified_at",
- "notified_at",
- }
- unknown = set(values) - allowed
- if unknown:
- raise ValueError(f"Unsupported cleanup fields: {sorted(unknown)}")
- if not values:
- return False
- assignments = ", ".join(f"{name}=%s" for name in values)
- connection = get_connection()
- try:
- with connection.cursor() as cursor:
- where = "WHERE id=%s"
- params = [*values.values(), item_id]
- if expected_cleanup_status is not None:
- if isinstance(expected_cleanup_status, (tuple, list, set, frozenset)):
- statuses = list(expected_cleanup_status)
- if not statuses:
- return False
- placeholders = ",".join(["%s"] * len(statuses))
- where += f" AND cleanup_status IN ({placeholders})"
- params.extend(statuses)
- else:
- where += " AND cleanup_status=%s"
- params.append(expected_cleanup_status)
- cursor.execute(
- f"UPDATE creative_rejection_cleanup_item SET {assignments} {where}",
- params,
- )
- return cursor.rowcount > 0
- finally:
- connection.close()
- def load_pending_notification_items(
- *,
- include_discovered: bool = False,
- ) -> list[dict[str, Any]]:
- connection = get_connection()
- try:
- with connection.cursor() as cursor:
- statuses = "'CREATIVE_DELETED','ALERT_PENDING'"
- if include_discovered:
- statuses += ",'DISCOVERED'"
- cursor.execute(
- f"""
- SELECT * FROM creative_rejection_cleanup_item
- WHERE cleanup_status IN ({statuses})
- AND (
- agency_notified_at IS NULL
- OR operator_notified_at IS NULL
- )
- ORDER BY check_date, agency_name, account_id, adgroup_id,
- dynamic_creative_id
- """
- )
- return list(cursor.fetchall())
- finally:
- connection.close()
- def load_unnotified_deleted_items(
- *,
- include_discovered: bool = False,
- ) -> list[dict[str, Any]]:
- """Compatibility entry point for pending cleanup notifications."""
- return load_pending_notification_items(include_discovered=include_discovered)
- def _mark_cleanup_channel_notified(
- item_ids: list[int],
- notified_at: datetime,
- *,
- channel: str,
- ) -> None:
- if not item_ids:
- return
- placeholders = ",".join(["%s"] * len(item_ids))
- if channel not in {"agency", "operator"}:
- raise ValueError(f"Unsupported cleanup notification channel: {channel}")
- channel_column = f"{channel}_notified_at"
- other_column = (
- "operator_notified_at" if channel == "agency" else "agency_notified_at"
- )
- connection = get_connection()
- try:
- with connection.cursor() as cursor:
- cursor.execute(
- f"""
- UPDATE creative_rejection_cleanup_item
- SET {channel_column}=%s,
- notified_at=CASE
- WHEN {other_column} IS NOT NULL THEN %s
- ELSE NULL
- END
- WHERE id IN ({placeholders}) AND {channel_column} IS NULL
- """,
- [notified_at, notified_at, *item_ids],
- )
- finally:
- connection.close()
- def mark_cleanup_items_agency_notified(
- item_ids: list[int], notified_at: datetime
- ) -> None:
- _mark_cleanup_channel_notified(item_ids, notified_at, channel="agency")
- def mark_cleanup_items_operator_notified(
- item_ids: list[int], notified_at: datetime
- ) -> None:
- _mark_cleanup_channel_notified(item_ids, notified_at, channel="operator")
- def mark_cleanup_items_notified(
- item_ids: list[int], notified_at: datetime
- ) -> None:
- """Compatibility entry point for the agency notification channel."""
- mark_cleanup_items_agency_notified(item_ids, notified_at)
- def _update_owned_cleanup_item(item_id: int, **values: Any) -> bool:
- """Update a row only while this run owns its DELETING claim."""
- return update_cleanup_item(
- item_id,
- _expected_cleanup_status="DELETING",
- **values,
- )
- def upsert_cleanup_delivery(record: dict[str, Any]) -> dict[str, Any]:
- connection = get_connection()
- try:
- with connection.cursor() as cursor:
- cursor.execute(
- """
- INSERT INTO creative_rejection_delivery
- (run_id, agency_name, agency_report_version, file_path,
- file_sha256, creative_rows, ad_rows,
- route_fingerprint, status)
- VALUES (%s,%s,%s,%s,%s,%s,%s,%s,'PENDING')
- ON DUPLICATE KEY UPDATE
- status=CASE
- WHEN status='SENT' THEN status
- WHEN NOT (route_fingerprint <=> VALUES(route_fingerprint))
- THEN 'PENDING'
- ELSE status
- END,
- error_message=CASE
- WHEN NOT (route_fingerprint <=> VALUES(route_fingerprint))
- THEN NULL ELSE error_message
- END,
- file_path=IF(status='SENT', file_path, VALUES(file_path)),
- file_sha256=IF(status='SENT', file_sha256, VALUES(file_sha256)),
- creative_rows=VALUES(creative_rows),
- ad_rows=VALUES(ad_rows),
- route_fingerprint=VALUES(route_fingerprint)
- """,
- (
- record["run_id"],
- record["agency_name"],
- record["agency_report_version"],
- record["file_path"],
- record["file_sha256"],
- record["creative_rows"],
- record["ad_rows"],
- record.get("route_fingerprint"),
- ),
- )
- cursor.execute(
- """
- SELECT * FROM creative_rejection_delivery
- WHERE run_id=%s AND agency_name=%s AND agency_report_version=%s
- """,
- (
- record["run_id"],
- record["agency_name"],
- record["agency_report_version"],
- ),
- )
- return cursor.fetchone()
- finally:
- connection.close()
- def update_cleanup_delivery(delivery_id: int, **values: Any) -> None:
- allowed = {
- "status",
- "sheet_token",
- "sheet_url",
- "response_code",
- "error_message",
- "sent_at",
- }
- increment_attempt = bool(values.pop("increment_attempt", False))
- unknown = set(values) - allowed
- if unknown:
- raise ValueError(f"Unsupported delivery fields: {sorted(unknown)}")
- assignments = [f"{name}=%s" for name in values]
- params = list(values.values())
- if increment_attempt:
- assignments.append("attempt_count=attempt_count+1")
- if not assignments:
- return
- connection = get_connection()
- try:
- with connection.cursor() as cursor:
- cursor.execute(
- f"UPDATE creative_rejection_delivery SET {', '.join(assignments)} WHERE id=%s",
- [*params, delivery_id],
- )
- finally:
- connection.close()
- def _file_sha256(path: Path) -> str:
- digest = hashlib.sha256()
- with path.open("rb") as handle:
- for chunk in iter(lambda: handle.read(1024 * 1024), b""):
- digest.update(chunk)
- return digest.hexdigest()
- def publish_cleanup_operator_summary(
- *,
- run_id: str,
- report: dict[str, object],
- chat_id: str,
- publisher: RoiFeishuPublisher,
- now: datetime,
- ) -> dict[str, object]:
- """Idempotently upload and send the all-agency report to the operator group."""
- path = Path(str(report["report"]))
- delivery_id: int | None = None
- try:
- if not path.is_file():
- raise FileNotFoundError(path)
- target_chat_id = str(chat_id or "").strip()
- if not target_chat_id:
- raise RuntimeError("FEISHU_AD_PROJECT_CHAT_ID 未配置")
- delivery = upsert_cleanup_delivery(
- {
- "run_id": run_id,
- "agency_name": OPERATOR_SUMMARY_ROUTE,
- "agency_report_version": str(report["report_version"]),
- "file_path": str(path),
- "file_sha256": _file_sha256(path),
- "creative_rows": int(report.get("creative_rows") or 0),
- "ad_rows": 0,
- "route_fingerprint": hashlib.sha256(
- target_chat_id.encode("utf-8")
- ).hexdigest(),
- }
- )
- delivery_id = int(delivery["id"])
- if delivery.get("status") == "SENT":
- return {
- "route": OPERATOR_SUMMARY_ROUTE,
- "status": "SENT",
- "sheet_url": str(delivery.get("sheet_url") or ""),
- "reused": True,
- }
- sheet_url = str(delivery.get("sheet_url") or "")
- sheet_token = str(delivery.get("sheet_token") or "")
- if not sheet_url or not sheet_token:
- imported = publisher.upload_workbook(path)
- sheet_url = imported["url"]
- sheet_token = imported["sheet_token"]
- update_cleanup_delivery(
- delivery_id,
- status="UPLOADED",
- sheet_token=sheet_token,
- sheet_url=sheet_url,
- )
- message_id = publisher.send_report_card(
- title=str(report.get("title") or path.stem),
- content=(
- (
- f"本次预演共 **{int(report.get('creative_rows') or 0)}** 条创意,"
- "未执行任何删除,包含建议删除项及需人工判断项。"
- )
- if report.get("dry_run")
- else (
- f"本批次共 **{int(report.get('creative_rows') or 0)}** 条创意,"
- "包含各代理自动删除及需人工判断的完整汇总。"
- )
- ),
- sheet_url=sheet_url,
- chat_id=target_chat_id,
- button_text="查看全部处理明细",
- )
- update_cleanup_delivery(
- delivery_id,
- status="SENT",
- response_code=message_id,
- sent_at=now,
- increment_attempt=True,
- )
- return {
- "route": OPERATOR_SUMMARY_ROUTE,
- "status": "SENT",
- "sheet_url": sheet_url,
- "reused": False,
- }
- except Exception as exc:
- if delivery_id is not None:
- try:
- update_cleanup_delivery(
- delivery_id,
- status="FAILED",
- error_message=str(exc),
- increment_attempt=True,
- )
- except Exception as audit_exc:
- logger.error("operator summary audit failed: %s", audit_exc)
- logger.error("operator summary delivery failed: %s", exc)
- return {
- "route": OPERATOR_SUMMARY_ROUTE,
- "status": "FAILED",
- "error": str(exc),
- }
- def _reject_reason(raw_result: dict[str, Any] | None, system_status: str) -> str:
- if raw_result:
- parsed = parse_review_result(raw_result)
- reasons = list(parsed.reject_messages)
- reasons.extend(fact.reason for fact in parsed.rejection_facts)
- unique = list(dict.fromkeys(reason.strip() for reason in reasons if reason.strip()))
- if unique:
- return ";".join(unique)
- if system_status == DENIED_SYSTEM_STATUS:
- return "腾讯正式审核未通过(接口未返回具体原因)"
- return "腾讯正式审核未通过"
- def _cleanup_reason(
- action: dict[str, Any],
- raw_result: dict[str, Any] | None,
- system_status: str,
- ) -> str:
- if action["cleanup_action"] == DELETE_CREATIVE:
- return "创意审核状态为 CREATIVE_SET_APPROVAL_STATUS_DENIED"
- return _reject_reason(raw_result, system_status)
- def _json_object(value: Any) -> dict[str, Any]:
- if isinstance(value, dict):
- return value
- if not value:
- return {}
- try:
- parsed = json.loads(str(value))
- except (TypeError, ValueError, json.JSONDecodeError):
- return {}
- return parsed if isinstance(parsed, dict) else {}
- def _display_date(value: Any) -> str:
- if isinstance(value, (date, datetime)):
- return value.strftime("%Y-%m-%d")
- return str(value or "")[:10]
- def _write_report(
- path: Path,
- rows: list[dict[str, Any]],
- *,
- columns: tuple[str, ...],
- ) -> None:
- workbook = Workbook()
- sheet = workbook.active
- sheet.title = "审核不通过创意清理"
- sheet.append(list(columns))
- for row in rows:
- raw_result = _json_object(
- row.get("review_result") or row.get("review_result_json")
- )
- pre_state = _json_object(row.get("pre_state") or row.get("pre_state_json"))
- granular = review_granularity_fields(raw_result)
- action = str(row.get("cleanup_action") or "")
- if (
- action == DELETE_CREATIVE
- and row.get("cleanup_status") == "DISCOVERED"
- ):
- execution_action = "建议删除创意(未执行)"
- elif action == DELETE_CREATIVE:
- execution_action = "删除创意"
- else:
- execution_action = "需人工判断"
- checked_at = (
- row.get("checked_at")
- or row.get("created_at")
- or row.get("updated_at")
- or row.get("deleted_at")
- )
- recent_cost_fen = row.get("recent_cost_fen")
- recent_cost_yuan = (
- ""
- if recent_cost_fen is None
- else f"{int(recent_cost_fen) / 100:.2f}"
- )
- cost_start = _display_date(row.get("cost_start_date"))
- cost_end = _display_date(row.get("cost_end_date"))
- values = {
- "代理名称": row.get("agency_name") or "",
- "账户ID": str(row["account_id"]),
- "账户名称": row.get("account_name") or "",
- "广告ID": str(row["adgroup_id"]),
- "广告名称": row.get("adgroup_name") or "",
- "创意ID": str(row["dynamic_creative_id"]),
- "创意名称": row.get("dynamic_creative_name") or "",
- "近3天累计历史消耗(元)": recent_cost_yuan,
- "消耗日期范围": (
- f"{cost_start} ~ {cost_end}" if cost_start and cost_end else ""
- ),
- "执行操作": execution_action,
- "操作判断原因": row.get("action_reason") or "",
- "配置状态": status_desc(
- pre_state.get("configured_status")
- or row.get("configured_status")
- ),
- "创意审核状态": status_desc(
- pre_state.get("creative_set_approval_status")
- or row.get("creative_set_approval_status")
- ),
- "元素粒度审核状态": granular["element_review_status"],
- "元素粒度审核不通过原因": granular["element_reject_reason"],
- "版位粒度审核状态": granular["site_review_status"],
- "版位粒度审核不通过原因": granular["site_reject_reason"],
- "审核不通过原因": row["reject_reason"],
- "检查时间": (
- checked_at.strftime("%Y-%m-%d %H:%M:%S")
- if isinstance(checked_at, (date, datetime))
- else str(checked_at or "")
- ),
- }
- sheet.append([values[column] for column in columns])
- header_fill = PatternFill("solid", fgColor="C65911")
- for cell in sheet[1]:
- cell.fill = header_fill
- cell.font = Font(color="FFFFFF", bold=True)
- cell.alignment = Alignment(horizontal="center", vertical="center")
- widths = {
- "代理名称": 22,
- "账户ID": 14,
- "账户名称": 22,
- "广告ID": 14,
- "广告名称": 30,
- "创意ID": 16,
- "创意名称": 30,
- "近3天累计历史消耗(元)": 22,
- "消耗日期范围": 24,
- "执行操作": 16,
- "操作判断原因": 60,
- "配置状态": 20,
- "创意审核状态": 24,
- "元素粒度审核状态": 40,
- "元素粒度审核不通过原因": 60,
- "版位粒度审核状态": 40,
- "版位粒度审核不通过原因": 60,
- "审核不通过原因": 60,
- "检查时间": 20,
- }
- for index, column in enumerate(columns, start=1):
- sheet.column_dimensions[get_column_letter(index)].width = widths[column]
- for row in sheet.iter_rows(min_row=2):
- for cell in row:
- cell.alignment = Alignment(vertical="top", wrap_text=True)
- for id_column in ("账户ID", "广告ID", "创意ID"):
- row[columns.index(id_column)].number_format = "@"
- sheet.freeze_panes = "A2"
- sheet.auto_filter.ref = sheet.dimensions
- path.parent.mkdir(parents=True, exist_ok=True)
- workbook.save(path)
- def write_cleanup_reports(
- rows: list[dict[str, Any]],
- output_dir: Path,
- report_date: str,
- ) -> tuple[str, list[dict[str, object]], dict[str, list[int]]]:
- grouped: dict[str, list[dict[str, Any]]] = defaultdict(list)
- for row in rows:
- grouped[str(row["agency_name"])].append(row)
- def digest_for(report_rows: list[dict[str, Any]], destination: str) -> str:
- digest_input = "|".join(
- ":".join(
- [
- str(row["id"]),
- str(row.get("cleanup_action") or ""),
- str(row.get("cleanup_status") or ""),
- str(row.get("recent_cost_fen")),
- _display_date(row.get("cost_end_date")),
- str(row.get("reject_reason") or ""),
- ]
- )
- for row in sorted(report_rows, key=lambda value: int(value["id"]))
- )
- return hashlib.sha256(
- f"{report_date}|{destination}|{digest_input}".encode("utf-8")
- ).hexdigest()[:12]
- digest = digest_for(rows, "all")
- run_id = f"reject_{report_date}_{digest}"
- reports: list[dict[str, object]] = []
- item_ids: dict[str, list[int]] = {}
- for agency, agency_rows in sorted(grouped.items()):
- agency_digest = digest_for(agency_rows, agency)
- safe_agency = agency.replace("/", "_").replace("\\", "_")
- path = output_dir / (
- f"{report_date}_{safe_agency}_创意审核异常处理_{agency_digest}.xlsx"
- )
- dry_run = any(
- row.get("cleanup_status") == "DISCOVERED"
- for row in agency_rows
- )
- _write_report(path, agency_rows, columns=AGENCY_REPORT_COLUMNS)
- reports.append(
- {
- "agency_name": agency,
- "report_version": REPORT_VERSION,
- "report": str(path),
- "title": (
- f"{report_date}_{agency}_创意审核异常处理预演通知"
- if dry_run
- else f"{report_date}_{agency}_创意审核异常处理通知"
- ),
- "creative_rows": len(agency_rows),
- "ad_rows": 0,
- "notification_type": (
- "creative_rejection_dry_run"
- if dry_run
- else "creative_rejection_cleanup"
- ),
- "run_id": f"reject_{report_date}_{agency_digest}",
- }
- )
- item_ids[agency] = [int(row["id"]) for row in agency_rows]
- return run_id, reports, item_ids
- def write_cleanup_operator_summary(
- rows: list[dict[str, Any]],
- output_dir: Path,
- report_date: str,
- run_id: str,
- ) -> dict[str, object]:
- digest_input = "|".join(
- ":".join(
- [
- str(row["id"]),
- str(row.get("cleanup_action") or ""),
- str(row.get("cleanup_status") or ""),
- str(row.get("recent_cost_fen")),
- _display_date(row.get("cost_end_date")),
- str(row.get("reject_reason") or ""),
- ]
- )
- for row in sorted(rows, key=lambda value: int(value["id"]))
- )
- digest = hashlib.sha256(
- f"{report_date}|{OPERATOR_SUMMARY_ROUTE}|{digest_input}".encode("utf-8")
- ).hexdigest()[:12]
- path = output_dir / (
- f"{report_date}_投放调控_创意审核异常处理汇总_{digest}.xlsx"
- )
- dry_run = any(row.get("cleanup_status") == "DISCOVERED" for row in rows)
- _write_report(path, rows, columns=OPERATOR_REPORT_COLUMNS)
- return {
- "report_version": f"{REPORT_VERSION}_operator_summary",
- "report": str(path),
- "title": (
- f"{report_date}_创意审核异常处理预演汇总通知"
- if dry_run
- else f"{report_date}_创意审核异常处理汇总通知"
- ),
- "creative_rows": len(rows),
- "run_id": f"reject_{report_date}_{digest}",
- "dry_run": dry_run,
- }
- def _chunks(values: list[int], size: int) -> list[list[int]]:
- return [values[offset : offset + size] for offset in range(0, len(values), size)]
- def _scan_one_account(
- account: dict[str, Any],
- *,
- tencent,
- review_fetcher: Callable[[int, list[int]], list[dict]],
- spend_start_date: date,
- spend_end_date: date,
- ) -> tuple[
- list[dict],
- dict[int, dict],
- dict[int, dict],
- dict[int, int],
- str | None,
- int,
- str | None,
- ]:
- account_id = int(account["account_id"])
- try:
- creatives = tencent.get_dynamic_creatives(account_id)
- ads = {
- int(ad["adgroup_id"]): ad
- for ad in tencent.get_ads(account_id)
- if _as_int(ad.get("adgroup_id")) is not None
- }
- ids = [
- int(row["dynamic_creative_id"])
- for row in creatives
- if _as_int(row.get("dynamic_creative_id")) is not None
- and row.get("configured_status") != DELETED_STATUS
- ]
- raw_by_id: dict[int, dict] = {}
- for batch in _chunks(ids, 100):
- for raw in review_fetcher(account_id, batch):
- creative_id = _as_int(raw.get("dynamic_creative_id"))
- if creative_id is not None:
- raw_by_id[creative_id] = raw
- cost_by_id: dict[int, int] = {}
- spend_error = None
- partial_ids = [
- int(row["dynamic_creative_id"])
- for row in creatives
- if _as_int(row.get("dynamic_creative_id")) is not None
- and row.get("configured_status") != DELETED_STATUS
- and row.get("creative_set_approval_status")
- == CREATIVE_PARTIAL_NORMAL_STATUS
- ]
- if partial_ids:
- try:
- cost_by_id = tencent.get_dynamic_creative_costs(
- account_id,
- partial_ids,
- spend_start_date,
- spend_end_date,
- )
- except Exception as exc:
- spend_error = str(exc)
- logger.exception(
- "creative cost scan failed account=%d start=%s end=%s",
- account_id,
- spend_start_date,
- spend_end_date,
- )
- return creatives, ads, raw_by_id, cost_by_id, spend_error, len(ids), None
- except Exception as exc:
- return [], {}, {}, {}, None, 0, f"account={account_id} scan failed: {exc}"
- def cleanup_precondition_failure(
- item: dict[str, Any],
- scanned_accounts: set[int],
- confirmed_actions: dict[tuple[int, int], dict[str, Any]],
- ) -> tuple[str, str] | None:
- """Fail closed unless this run reconfirmed the exact cleanup action."""
- account_id = int(item["account_id"])
- creative_id = int(item["dynamic_creative_id"])
- if account_id not in scanned_accounts:
- return "DEFERRED", "本轮账户审核结果扫描失败,未执行删除"
- confirmed = confirmed_actions.get((account_id, creative_id))
- if not confirmed:
- return (
- "SKIPPED_REVIEW_NOT_RECONFIRMED",
- "本轮未按新规则再次确认清理动作",
- )
- item_action = str(item.get("cleanup_action") or "")
- if item_action != DELETE_CREATIVE:
- return "SKIPPED_REVIEW_NOT_RECONFIRMED", "当前规则不允许删除拒审元素组件"
- if item_action != confirmed["cleanup_action"]:
- return "SKIPPED_REVIEW_NOT_RECONFIRMED", "本轮清理动作与候选记录不一致"
- return None
- def _same_cleanup_action(
- expected_action: str,
- actual: dict[str, Any] | None,
- ) -> bool:
- return bool(
- expected_action == DELETE_CREATIVE
- and actual
- and actual.get("cleanup_action") == DELETE_CREATIVE
- )
- def run_rejected_creative_cleanup(
- *,
- output_dir: Path,
- now: datetime | None = None,
- tencent=None,
- odps=None,
- review_fetcher: Callable[[int, list[int]], list[dict]] | None = None,
- publisher: RoiFeishuPublisher | None = None,
- notifier=None,
- ) -> dict[str, Any]:
- """Scan recent-spend accounts and apply the current creative-review rules."""
- effective_now = now or datetime.now(SHANGHAI)
- if effective_now.tzinfo is None:
- effective_now = effective_now.replace(tzinfo=SHANGHAI)
- apply_enabled = _env_flag("DAILY_REJECTED_CREATIVE_APPLY_ENABLED")
- webhook_config = AgencyWebhookConfig.from_env()
- if apply_enabled and not webhook_config.enabled:
- raise RuntimeError(
- "DAILY_REJECTED_CREATIVE_APPLY_ENABLED=1 requires "
- "ROI_AGENCY_WEBHOOK_ENABLED=1"
- )
- if apply_enabled and not os.getenv("FEISHU_AD_PROJECT_CHAT_ID", "").strip():
- raise RuntimeError(
- "DAILY_REJECTED_CREATIVE_APPLY_ENABLED=1 requires "
- "FEISHU_AD_PROJECT_CHAT_ID"
- )
- initialize_schema()
- if odps is None:
- from roi_control.odps_client import ODPSClient
- odps_client = ODPSClient(project=os.getenv("ODPS_PROJECT", "loghubods"))
- else:
- odps_client = odps
- end_date = resolve_end_date(odps_client, now=effective_now)
- start_date, end_date = date_window(end_date)
- spend_end_date = effective_now.date() - timedelta(days=1)
- spend_start_date = spend_end_date - timedelta(days=2)
- cost_threshold_fen = partial_creative_cost_threshold_fen()
- wechat_cost_threshold_fen = wechat_mini_program_cost_threshold_fen()
- daily = fetch_daily_data(odps_client, start_date, end_date)
- context = build_agency_context(daily)
- accounts = fetch_recent_spend_accounts(odps_client, start_date, end_date)
- unresolved_account_ids = [
- int(account["account_id"])
- for account in accounts
- if int(account["account_id"]) not in context["account_agencies"]
- ]
- try:
- context["fallback_account_agencies"] = fetch_account_agency_fallbacks(
- odps_client,
- unresolved_account_ids,
- )
- except Exception:
- logger.exception(
- "account agency fallback query failed accounts=%d",
- len(unresolved_account_ids),
- )
- owned_tencent = tencent is None
- if tencent is None:
- from tencent_client import TencentClient
- client = TencentClient()
- else:
- client = tencent
- account_ids = [int(account["account_id"]) for account in accounts]
- prefetched_tokens = prefetch_account_access_tokens(account_ids)
- seed_tokens = getattr(client, "seed_access_tokens", None)
- if callable(seed_tokens):
- seed_tokens(prefetched_tokens)
- fetch_reviews = review_fetcher or fetch_dynamic_creative_review_results
- discovered = 0
- scanned = 0
- scanned_accounts: set[int] = set()
- confirmed_actions: dict[tuple[int, int], dict[str, Any]] = {}
- scan_errors: list[str] = []
- try:
- scan_workers = int(os.getenv("TENCENT_AD_ACCOUNT_SCAN_WORKERS", "8"))
- if scan_workers < 1:
- raise ValueError("TENCENT_AD_ACCOUNT_SCAN_WORKERS must be at least 1")
- workers = min(scan_workers, len(accounts), 32) if accounts else 1
- def run_scan(account: dict[str, Any]):
- if owned_tencent:
- from tencent_client import TencentClient
- scan_client = TencentClient()
- scan_client.seed_access_tokens(prefetched_tokens)
- else:
- scan_client = client
- try:
- return _scan_one_account(
- account,
- tencent=scan_client,
- review_fetcher=fetch_reviews,
- spend_start_date=spend_start_date,
- spend_end_date=spend_end_date,
- )
- finally:
- if scan_client is not client:
- scan_client.session.close()
- logger.info("account scan started accounts=%d workers=%d", len(accounts), workers)
- scan_results = []
- with ThreadPoolExecutor(
- max_workers=workers,
- thread_name_prefix="creative-scan",
- ) as executor:
- futures = {executor.submit(run_scan, account): account for account in accounts}
- for completed, future in enumerate(as_completed(futures), start=1):
- account = futures[future]
- account_id = int(account["account_id"])
- try:
- (
- creatives,
- ads,
- raw_by_id,
- cost_by_id,
- spend_error,
- account_scanned,
- error,
- ) = future.result()
- except Exception as exc:
- creatives, ads, raw_by_id, cost_by_id = [], {}, {}, {}
- spend_error, account_scanned = None, 0
- error = f"account={account_id} scan failed: {exc}"
- scanned += account_scanned
- if error:
- logger.error(error)
- scan_errors.append(error)
- else:
- scanned_accounts.add(account_id)
- scan_results.append(
- (
- account,
- creatives,
- ads,
- raw_by_id,
- cost_by_id,
- spend_error,
- )
- )
- logger.info(
- "account scan progress=%d/%d account=%d creatives=%d error=%s",
- completed,
- len(accounts),
- account_id,
- account_scanned,
- bool(error),
- )
- def process_creative(task):
- (
- account,
- account_id,
- creative,
- ads,
- raw_by_id,
- cost_by_id,
- spend_error,
- ) = task
- creative_id = _as_int(creative.get("dynamic_creative_id"))
- adgroup_id = _as_int(creative.get("adgroup_id"))
- if creative_id is None or adgroup_id is None:
- return None
- raw_result = raw_by_id.get(creative_id)
- system_status = str(creative.get("system_status") or "")
- is_partial = (
- creative.get("creative_set_approval_status")
- == CREATIVE_PARTIAL_NORMAL_STATUS
- )
- action = determine_cleanup_action(
- creative,
- raw_result,
- recent_cost_fen=cost_by_id.get(creative_id),
- cost_threshold_fen=cost_threshold_fen,
- wechat_cost_threshold_fen=wechat_cost_threshold_fen,
- spend_error=(spend_error if is_partial else None),
- )
- if action is None:
- return None
- ad = ads.get(adgroup_id) or {}
- upsert_cleanup_candidate(
- {
- "account_id": account_id,
- "account_name": context["account_names"].get(account_id)
- or account.get("account_name")
- or "",
- "agency_name": _resolve_agency(
- context, account_id, creative_id
- ),
- "adgroup_id": adgroup_id,
- "adgroup_name": ad.get("adgroup_name") or "",
- "dynamic_creative_id": creative_id,
- "dynamic_creative_name": creative.get(
- "dynamic_creative_name"
- )
- or "",
- "check_date": effective_now.date(),
- **action,
- "action_reason": action.get("action_reason")
- or _cleanup_reason(action, raw_result, system_status),
- "reject_reason": _reject_reason(raw_result, system_status),
- "cost_start_date": spend_start_date,
- "cost_end_date": spend_end_date,
- "review_result": raw_result or {},
- "pre_state": creative,
- }
- )
- return account_id, creative_id, action
- tasks = [
- (
- account,
- int(account["account_id"]),
- creative,
- ads,
- raw_by_id,
- cost_by_id,
- spend_error,
- )
- for account, creatives, ads, raw_by_id, cost_by_id, spend_error
- in scan_results
- for creative in creatives
- ]
- process_workers = (
- min(
- int(os.getenv("TENCENT_AD_CREATIVE_PROCESS_WORKERS", "8")),
- len(tasks),
- 32,
- )
- if tasks
- else 1
- )
- logger.info(
- "creative processing started creatives=%d workers=%d",
- len(tasks),
- process_workers,
- )
- with ThreadPoolExecutor(
- max_workers=process_workers,
- thread_name_prefix="creative-process",
- ) as executor:
- futures = {
- executor.submit(process_creative, task): task for task in tasks
- }
- for completed, future in enumerate(as_completed(futures), start=1):
- task = futures[future]
- account_id = task[1]
- creative_id = _as_int(task[2].get("dynamic_creative_id"))
- try:
- result = future.result()
- except Exception as exc:
- error = (
- "creative processing failed "
- f"account={account_id} creative={creative_id}: {exc}"
- )
- scan_errors.append(error)
- logger.exception(error)
- continue
- if result is None:
- continue
- account_id, creative_id, action = result
- confirmed_actions[(account_id, creative_id)] = action
- discovered += 1
- if completed % 500 == 0:
- logger.info(
- "creative processing progress=%d/%d confirmed=%d",
- completed,
- len(tasks),
- discovered,
- )
- deleted = 0
- deferred = 0
- delete_errors: list[str] = []
- write_lock_name = os.getenv(
- "RTC_DB_LOCK_NAME", "tencent_realtime_control"
- )
- retryable_items = load_retryable_cleanup_items()
- # 阶段一(锁外):纯前置判断,无腾讯写。通过者进入 deletable_items。
- deletable_items: list[dict[str, Any]] = []
- for item in retryable_items if apply_enabled else []:
- item_id = int(item["id"])
- account_id = int(item["account_id"])
- creative_id = int(item["dynamic_creative_id"])
- snapshot_status = str(item.get("cleanup_status") or "")
- if snapshot_status == "DELETING":
- # A stale claim is recovered only after obtaining the global
- # Tencent write lock; do not mutate a potentially live owner.
- deletable_items.append(item)
- continue
- if item.get("cleanup_action") != DELETE_CREATIVE:
- update_cleanup_item(
- item_id,
- _expected_cleanup_status=snapshot_status,
- cleanup_status="SKIPPED_REVIEW_NOT_RECONFIRMED",
- error_message="当前规则仅允许删除整条创意",
- )
- continue
- agency = str(item.get("agency_name") or "") or _resolve_agency(
- context, account_id, creative_id
- )
- if agency and agency != item.get("agency_name"):
- update_cleanup_item(
- item_id,
- _expected_cleanup_status=snapshot_status,
- agency_name=agency,
- )
- item["agency_name"] = agency
- webhook_url = resolve_agency_webhook(
- agency, webhook_config.webhooks or {}
- )
- if not agency or not webhook_url:
- reason = (
- "代理商归属为空,禁止自动删除"
- if not agency
- else f"代理商 {agency} 未配置通知群,禁止自动删除"
- )
- updated = update_cleanup_item(
- item_id,
- _expected_cleanup_status=snapshot_status,
- cleanup_status="DEFERRED",
- error_message=reason,
- )
- if updated is not False:
- deferred += 1
- continue
- if item.get("cleanup_status") not in {
- "WRITE_OUTCOME_UNKNOWN",
- "DELETING",
- }:
- precondition_failure = cleanup_precondition_failure(
- item,
- scanned_accounts,
- confirmed_actions,
- )
- if precondition_failure:
- status, reason = precondition_failure
- updated = update_cleanup_item(
- item_id,
- _expected_cleanup_status=snapshot_status,
- cleanup_status=status,
- error_message=reason,
- )
- if status == "DEFERRED" and updated is not False:
- deferred += 1
- continue
- deletable_items.append(item)
- # 阶段二(锁内并发):整个删除批次持锁一次,锁内并发回读/复审/删除。
- # 每 worker 用独立 TencentClient(requests.Session 非线程安全),
- # 避免共享 session 并发导致不可预测行为(与扫描阶段一致)。
- if deletable_items:
- configured_delete_workers = int(
- os.getenv("TENCENT_AD_DELETE_WORKERS", "4")
- )
- if configured_delete_workers < 1:
- raise ValueError("TENCENT_AD_DELETE_WORKERS must be at least 1")
- delete_workers = min(
- configured_delete_workers,
- len(deletable_items),
- 16,
- )
- if not owned_tencent:
- # A caller-supplied client may wrap a requests.Session and is
- # not assumed to be thread-safe.
- delete_workers = 1
- delete_worker_local = threading.local()
- delete_worker_clients = []
- delete_worker_clients_lock = threading.Lock()
- def delete_one(item, delete_client):
- """锁内单条创意:回读 → 复审 → 删除;返回 (deleted, deferred, error)。"""
- item_id = int(item["id"])
- account_id = int(item["account_id"])
- creative_id = int(item["dynamic_creative_id"])
- try:
- try:
- before = delete_client.get_dynamic_creative(
- account_id, creative_id
- )
- except Exception as read_exc:
- if (
- item.get("cleanup_action") == DELETE_CREATIVE
- and str(read_exc).startswith(
- "Dynamic creative not found:"
- )
- ):
- updated = _update_owned_cleanup_item(
- item_id,
- cleanup_status="CREATIVE_DELETED",
- error_message=None,
- readback_json=_json({"deleted_from_listing": True}),
- deleted_at=effective_now,
- )
- if updated is not False:
- return 1, 0, None
- return 0, 0, (
- f"account={account_id} creative={creative_id}: "
- "delete result ignored because claim ownership was lost"
- )
- raise
- if item.get("cleanup_status") in {
- "WRITE_OUTCOME_UNKNOWN",
- "DELETING",
- }:
- precondition_failure = cleanup_precondition_failure(
- item,
- scanned_accounts,
- confirmed_actions,
- )
- if precondition_failure:
- status, reason = precondition_failure
- _update_owned_cleanup_item(
- item_id,
- cleanup_status=status,
- error_message=reason,
- )
- if status == "DEFERRED":
- return 0, 1, None
- return 0, 0, None
- action = str(item.get("cleanup_action") or "")
- approval_status = str(
- before.get("creative_set_approval_status") or ""
- )
- fresh_raw = None
- if approval_status == CREATIVE_DENIED_STATUS:
- fresh_action = determine_cleanup_action(before, None)
- else:
- fresh_results = fetch_reviews(account_id, [creative_id])
- fresh_raw = next(
- (
- result
- for result in fresh_results
- if _as_int(result.get("dynamic_creative_id"))
- == creative_id
- ),
- None,
- )
- fresh_cost_fen = None
- fresh_spend_error = None
- if approval_status == CREATIVE_PARTIAL_NORMAL_STATUS:
- try:
- fresh_cost_fen = (
- delete_client.get_dynamic_creative_costs(
- account_id,
- [creative_id],
- spend_start_date,
- spend_end_date,
- ).get(creative_id, 0)
- )
- except Exception as spend_exc:
- fresh_spend_error = str(spend_exc)
- fresh_action = determine_cleanup_action(
- before,
- fresh_raw,
- recent_cost_fen=fresh_cost_fen,
- cost_threshold_fen=cost_threshold_fen,
- wechat_cost_threshold_fen=wechat_cost_threshold_fen,
- spend_error=fresh_spend_error,
- )
- if not _same_cleanup_action(action, fresh_action):
- if (
- fresh_action
- and fresh_action.get("cleanup_action") == ALERT_ONLY
- ):
- _update_owned_cleanup_item(
- item_id,
- cleanup_action=ALERT_ONLY,
- target_component_ids_json="[]",
- target_element_ids_json="[]",
- recent_cost_fen=fresh_action.get("recent_cost_fen"),
- cost_start_date=spend_start_date,
- cost_end_date=spend_end_date,
- action_reason=fresh_action["action_reason"],
- reject_reason=_reject_reason(
- fresh_raw,
- str(before.get("system_status") or ""),
- ),
- review_result_json=_json(fresh_raw or {}),
- cleanup_status="ALERT_PENDING",
- error_message=None,
- pre_state_json=_json(before),
- readback_json=None,
- deleted_at=None,
- notified_at=None,
- )
- return 0, 0, None
- updated = _update_owned_cleanup_item(
- item_id,
- cleanup_status="SKIPPED_REVIEW_NOT_RECONFIRMED",
- error_message="写锁内回读发现整创意删除条件已变化",
- pre_state_json=_json(before),
- )
- return 0, 0, None
- if action == DELETE_CREATIVE:
- readback = delete_client.delete_dynamic_creative(
- account_id, creative_id
- )
- updated = _update_owned_cleanup_item(
- item_id,
- cleanup_status="CREATIVE_DELETED",
- error_message=None,
- pre_state_json=_json(before),
- readback_json=_json(readback),
- deleted_at=effective_now,
- )
- if updated is not False:
- return 1, 0, None
- return 0, 0, (
- f"account={account_id} creative={creative_id}: "
- "delete result ignored because claim ownership was lost"
- )
- return 0, 0, None
- except Exception as exc:
- from tencent_client import (
- PostWriteVerificationError,
- TencentWriteOutcomeUnknownError,
- TencentWriteRateLimitedError,
- )
- outcome_unknown = isinstance(
- exc,
- (
- TencentWriteOutcomeUnknownError,
- PostWriteVerificationError,
- ),
- )
- rate_limited = isinstance(exc, TencentWriteRateLimitedError)
- error = f"account={account_id} creative={creative_id}: {exc}"
- try:
- _update_owned_cleanup_item(
- item_id,
- cleanup_status=(
- "DEFERRED"
- if rate_limited
- else (
- "WRITE_OUTCOME_UNKNOWN"
- if outcome_unknown
- else "FAILED"
- )
- ),
- error_message=str(exc)[:4000],
- )
- except Exception as update_exc:
- logger.exception(
- "creative delete failure status update failed "
- "account=%d creative=%d",
- account_id,
- creative_id,
- )
- error += f"; status_update_failed={update_exc}"
- return 0, int(rate_limited), error
- def run_delete(item):
- if owned_tencent:
- delete_client = getattr(delete_worker_local, "client", None)
- if delete_client is None:
- from tencent_client import TencentClient
- delete_client = TencentClient()
- delete_client.seed_access_tokens(prefetched_tokens)
- delete_worker_local.client = delete_client
- with delete_worker_clients_lock:
- delete_worker_clients.append(delete_client)
- else:
- delete_client = client
- return delete_one(item, delete_client)
- with advisory_lock(write_lock_name) as acquired:
- if not acquired:
- for item in deletable_items:
- snapshot_status = str(item.get("cleanup_status") or "")
- if snapshot_status == "DELETING":
- continue
- updated = update_cleanup_item(
- int(item["id"]),
- _expected_cleanup_status=snapshot_status,
- cleanup_status="DEFERRED",
- error_message="腾讯写锁被实时调控占用",
- )
- if updated is not False:
- deferred += 1
- else:
- claimed_items = []
- for item in deletable_items:
- try:
- if claim_cleanup_item(int(item["id"])):
- claimed_items.append(item)
- except Exception as exc:
- error = (
- "creative delete claim failed "
- f"account={item['account_id']} "
- f"creative={item['dynamic_creative_id']}: {exc}"
- )
- delete_errors.append(error)
- logger.exception(error)
- executable_items = []
- for item in claimed_items:
- agency = str(item.get("agency_name") or "") or _resolve_agency(
- context,
- int(item["account_id"]),
- int(item["dynamic_creative_id"]),
- )
- webhook_url = resolve_agency_webhook(
- agency, webhook_config.webhooks or {}
- )
- if agency and webhook_url:
- if agency != item.get("agency_name"):
- _update_owned_cleanup_item(
- int(item["id"]),
- agency_name=agency,
- )
- item["agency_name"] = agency
- executable_items.append(item)
- continue
- reason = (
- "代理商归属为空,禁止自动删除"
- if not agency
- else f"代理商 {agency} 未配置通知群,禁止自动删除"
- )
- updated = _update_owned_cleanup_item(
- int(item["id"]),
- cleanup_status="DEFERRED",
- error_message=reason,
- )
- if updated is not False:
- deferred += 1
- logger.info(
- "creative delete started candidates=%d claimed=%d workers=%d",
- len(deletable_items),
- len(executable_items),
- delete_workers,
- )
- if executable_items:
- try:
- with ThreadPoolExecutor(
- max_workers=min(delete_workers, len(executable_items)),
- thread_name_prefix="creative-delete",
- ) as executor:
- futures = {
- executor.submit(run_delete, item): item
- for item in executable_items
- }
- for future in as_completed(futures):
- item = futures[future]
- try:
- d_deleted, d_deferred, d_error = future.result()
- except Exception as exc:
- d_deleted = 0
- d_deferred = 0
- d_error = (
- "creative delete worker failed "
- f"account={item['account_id']} "
- f"creative={item['dynamic_creative_id']}: {exc}"
- )
- logger.exception(d_error)
- try:
- _update_owned_cleanup_item(
- int(item["id"]),
- cleanup_status="FAILED",
- error_message=str(exc)[:4000],
- )
- except Exception as update_exc:
- logger.exception(
- "creative delete worker failure status "
- "update failed account=%s creative=%s",
- item["account_id"],
- item["dynamic_creative_id"],
- )
- d_error += (
- f"; status_update_failed={update_exc}"
- )
- deleted += d_deleted
- deferred += d_deferred
- if d_error:
- delete_errors.append(d_error)
- finally:
- for delete_worker_client in delete_worker_clients:
- delete_worker_client.session.close()
- deliveries: list[dict[str, object]] = []
- operator_deliveries: list[dict[str, object]] = []
- notification_errors: list[str] = []
- def publish_pending_notifications(
- pending_notifications: list[dict[str, Any]],
- ) -> None:
- owned_publisher = publisher is None
- sheet_publisher = publisher or RoiFeishuPublisher(require_chat_ids=False)
- try:
- rows_by_date: dict[str, list[dict[str, Any]]] = defaultdict(list)
- for row in pending_notifications:
- rows_by_date[_display_date(row.get("check_date"))].append(row)
- for check_date, daily_rows in sorted(rows_by_date.items()):
- report_date = check_date.replace("-", "")
- report_dir = output_dir / report_date
- agency_rows = [
- row
- for row in daily_rows
- if row.get("agency_notified_at") is None
- and str(row.get("agency_name") or "").strip()
- ]
- run_id = f"reject_{report_date}_{REPORT_VERSION}"
- if agency_rows:
- run_id, reports, report_item_ids = write_cleanup_reports(
- agency_rows,
- report_dir,
- report_date,
- )
- for report in reports:
- daily_deliveries = publish_agency_reports(
- run_id=str(report.get("run_id") or run_id),
- reports=[report],
- config=webhook_config,
- publisher=sheet_publisher,
- notifier=notifier,
- now=effective_now,
- upsert_delivery=upsert_cleanup_delivery,
- update_delivery=update_cleanup_delivery,
- )
- deliveries.extend(daily_deliveries)
- for outcome in daily_deliveries:
- agency_name = str(outcome["agency_name"])
- if outcome.get("status") == "SENT":
- mark_cleanup_items_notified(
- report_item_ids.get(agency_name, []),
- effective_now,
- )
- else:
- notification_errors.append(
- f"agency={agency_name}: "
- f"{outcome.get('error') or outcome.get('reason') or outcome.get('status')}"
- )
- operator_rows = [
- row
- for row in daily_rows
- if row.get("operator_notified_at") is None
- ]
- if operator_rows:
- operator_report = write_cleanup_operator_summary(
- operator_rows,
- report_dir,
- report_date,
- run_id,
- )
- operator_outcome = publish_cleanup_operator_summary(
- run_id=str(operator_report["run_id"]),
- report=operator_report,
- chat_id=os.getenv("FEISHU_AD_PROJECT_CHAT_ID", ""),
- publisher=sheet_publisher,
- now=effective_now,
- )
- operator_deliveries.append(operator_outcome)
- if operator_outcome.get("status") == "SENT":
- mark_cleanup_items_operator_notified(
- [int(row["id"]) for row in operator_rows],
- effective_now,
- )
- else:
- notification_errors.append(
- f"operator={check_date}: "
- f"{operator_outcome.get('error') or operator_outcome.get('status')}"
- )
- finally:
- if owned_publisher:
- sheet_publisher.close()
- pending_notification_probe = load_unnotified_deleted_items(
- include_discovered=not apply_enabled,
- )
- if pending_notification_probe and webhook_config.enabled:
- notification_lock_name = os.getenv(
- "DAILY_REJECTED_CREATIVE_NOTIFICATION_LOCK_NAME",
- "ad_rejected_creative_notification",
- )
- with advisory_lock(notification_lock_name) as acquired:
- if not acquired:
- logger.info(
- "creative cleanup notification skipped: lock busy name=%s",
- notification_lock_name,
- )
- else:
- pending_notifications = load_unnotified_deleted_items(
- include_discovered=not apply_enabled,
- )
- if pending_notifications:
- publish_pending_notifications(pending_notifications)
- return {
- "apply_enabled": apply_enabled,
- "account_scope": "opengid_recent_3d_spend",
- "account_scope_start_date": start_date,
- "account_scope_end_date": end_date,
- "accounts": len(accounts),
- "account_ids": account_ids,
- "tokens_prefetched": len(prefetched_tokens),
- "creatives_scanned": scanned,
- "rejected_discovered": discovered,
- "pending_cleanup": len(retryable_items) if not apply_enabled else 0,
- "deleted": deleted,
- "deferred": deferred,
- "scan_errors": scan_errors,
- "delete_errors": delete_errors,
- "notification_errors": notification_errors,
- "deliveries": deliveries,
- "operator_deliveries": operator_deliveries,
- }
- finally:
- if owned_tencent:
- client.session.close()
|