execution.py 15 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469
  1. """Execute an approved ROI batch with Tencent read-back and state coordination."""
  2. from __future__ import annotations
  3. import json
  4. import logging
  5. import os
  6. from datetime import date, datetime, timedelta
  7. from decimal import Decimal, ROUND_HALF_UP
  8. from typing import Any
  9. from storage import (
  10. advisory_lock,
  11. load_ad_states,
  12. load_managed_accounts,
  13. sync_roi_base_bid,
  14. )
  15. from tencent_client import (
  16. ACTIVE_STATUS,
  17. SUSPEND_STATUS,
  18. PostWriteVerificationError,
  19. TencentClient,
  20. TencentWriteOutcomeUnknownError,
  21. current_bid_fen,
  22. resolve_bid_field,
  23. )
  24. from .config import RoiConfig
  25. from .policy import ACTION_PAUSE_CREATIVE, ACTION_SCALE_BID
  26. from .repository import (
  27. FINAL_STATUSES,
  28. claim_run_for_execution,
  29. finalize_run,
  30. load_action_items,
  31. load_run,
  32. reject_run,
  33. update_action_item,
  34. )
  35. logger = logging.getLogger("auto_put_ad_mini.roi_execution")
  36. def _round_fen(value: Decimal) -> int:
  37. return int(value.quantize(Decimal("1"), rounding=ROUND_HALF_UP))
  38. def _managed_accounts() -> dict[int, dict[str, Any]]:
  39. return {
  40. int(row["account_id"]): row
  41. for row in load_managed_accounts()
  42. }
  43. def _failure_status(exc: Exception) -> str:
  44. if isinstance(exc, PostWriteVerificationError):
  45. return "VERIFY_FAILED"
  46. if isinstance(exc, TencentWriteOutcomeUnknownError):
  47. return "OUTCOME_UNKNOWN"
  48. return "FAILED"
  49. def plan_scale_bid(
  50. *,
  51. current_bid_fen: int,
  52. base_bid_fen: int,
  53. initial_base_bid_fen: int,
  54. boosted_date: date | None,
  55. control_date: date,
  56. intraday_ratio: Decimal,
  57. scale_ratio: Decimal,
  58. max_base_ratio: Decimal,
  59. ) -> dict[str, Any]:
  60. """Plan a permanent base increase without reapplying a restored boost."""
  61. boosted_date_today = boosted_date == control_date
  62. expected_boosted = _round_fen(Decimal(base_bid_fen) * intraday_ratio)
  63. actively_boosted = boosted_date_today and current_bid_fen == expected_boosted
  64. if current_bid_fen != base_bid_fen and not actively_boosted:
  65. return {
  66. "allowed": False,
  67. "reason": (
  68. f"current bid={current_bid_fen} does not match managed "
  69. f"base={base_bid_fen} or intraday boost={expected_boosted}"
  70. ),
  71. }
  72. max_base = _round_fen(Decimal(initial_base_bid_fen) * max_base_ratio)
  73. new_base = min(
  74. _round_fen(Decimal(base_bid_fen) * scale_ratio),
  75. max_base,
  76. )
  77. if new_base <= base_bid_fen:
  78. return {
  79. "allowed": False,
  80. "reason": f"permanent base reached cap={max_base}",
  81. "cap_reached": True,
  82. "max_base_bid_fen": max_base,
  83. }
  84. target = (
  85. _round_fen(Decimal(new_base) * intraday_ratio)
  86. if actively_boosted
  87. else new_base
  88. )
  89. return {
  90. "allowed": True,
  91. "new_base_bid_fen": new_base,
  92. "target_bid_fen": target,
  93. "boosted_date_today": boosted_date_today,
  94. "preserve_intraday_boost": actively_boosted,
  95. }
  96. def _execute_pause(
  97. item: dict[str, Any],
  98. *,
  99. client: TencentClient,
  100. now: datetime,
  101. ) -> None:
  102. account_id = int(item["account_id"])
  103. adgroup_id = int(item["adgroup_id"])
  104. creative_id = int(item["dynamic_creative_id"])
  105. creative = client.get_dynamic_creative(account_id, creative_id)
  106. actual_adgroup_id = int(creative.get("adgroup_id") or 0)
  107. before_status = str(creative.get("configured_status") or "")
  108. if actual_adgroup_id != adgroup_id:
  109. update_action_item(
  110. int(item["id"]),
  111. execution_status="SKIPPED_MAPPING_MISMATCH",
  112. skip_reason=(
  113. f"creative belongs to adgroup={actual_adgroup_id}, expected={adgroup_id}"
  114. ),
  115. pre_state_json=creative,
  116. executed_at=now,
  117. )
  118. return
  119. if item["execution_status"] == "PREPARED" and before_status == SUSPEND_STATUS:
  120. update_action_item(
  121. int(item["id"]),
  122. before_status=item.get("before_status"),
  123. target_status=SUSPEND_STATUS,
  124. readback_status=before_status,
  125. execution_status="SUCCESS",
  126. readback_json=creative,
  127. executed_at=now,
  128. )
  129. return
  130. if before_status != ACTIVE_STATUS:
  131. update_action_item(
  132. int(item["id"]),
  133. before_status=before_status,
  134. target_status=SUSPEND_STATUS,
  135. readback_status=before_status,
  136. execution_status="SKIPPED_STATE_MISMATCH",
  137. skip_reason="creative is no longer active",
  138. pre_state_json=creative,
  139. executed_at=now,
  140. )
  141. return
  142. update_action_item(
  143. int(item["id"]),
  144. before_status=before_status,
  145. target_status=SUSPEND_STATUS,
  146. execution_status="PREPARED",
  147. pre_state_json=creative,
  148. )
  149. try:
  150. readback = client.update_dynamic_creative_status(
  151. account_id,
  152. creative_id,
  153. SUSPEND_STATUS,
  154. )
  155. update_action_item(
  156. int(item["id"]),
  157. readback_status=str(readback.get("configured_status") or ""),
  158. execution_status="SUCCESS",
  159. readback_json=readback,
  160. executed_at=now,
  161. )
  162. except Exception as exc:
  163. update_action_item(
  164. int(item["id"]),
  165. execution_status=_failure_status(exc),
  166. error_message=str(exc),
  167. readback_json=(exc.actual if isinstance(exc, PostWriteVerificationError) else None),
  168. executed_at=now,
  169. )
  170. def _resume_prepared_scale(
  171. item: dict[str, Any],
  172. ad: dict[str, Any],
  173. *,
  174. client: TencentClient,
  175. now: datetime,
  176. ) -> None:
  177. details = json.loads(item.get("pre_state_json") or "{}")
  178. bid_field = str(item.get("bid_field") or "")
  179. current = current_bid_fen(ad, bid_field)
  180. before = int(item.get("before_bid_fen") or 0)
  181. target = int(item.get("target_bid_fen") or 0)
  182. if current == target:
  183. sync_roi_base_bid(
  184. account_id=int(item["account_id"]),
  185. adgroup_id=int(item["adgroup_id"]),
  186. adgroup_name=str(ad.get("adgroup_name") or ""),
  187. bid_field=bid_field,
  188. initial_base_bid_fen=int(item["initial_base_bid_fen"]),
  189. new_base_bid_fen=int(details["new_base_bid_fen"]),
  190. boosted_date=(
  191. now.date() if details.get("boosted_date_today") else None
  192. ),
  193. action_at=now,
  194. )
  195. update_action_item(
  196. int(item["id"]),
  197. execution_status="SUCCESS",
  198. readback_json=ad,
  199. executed_at=now,
  200. )
  201. return
  202. if current != before:
  203. update_action_item(
  204. int(item["id"]),
  205. execution_status="SKIPPED_STATE_MISMATCH",
  206. skip_reason=f"bid changed after preparation: before={before} current={current}",
  207. readback_json=ad,
  208. executed_at=now,
  209. )
  210. return
  211. readback = client.update_ad(
  212. int(item["account_id"]),
  213. int(item["adgroup_id"]),
  214. bid_field=bid_field,
  215. target_bid_fen=target,
  216. )
  217. sync_roi_base_bid(
  218. account_id=int(item["account_id"]),
  219. adgroup_id=int(item["adgroup_id"]),
  220. adgroup_name=str(ad.get("adgroup_name") or ""),
  221. bid_field=bid_field,
  222. initial_base_bid_fen=int(item["initial_base_bid_fen"]),
  223. new_base_bid_fen=int(details["new_base_bid_fen"]),
  224. boosted_date=(now.date() if details.get("boosted_date_today") else None),
  225. action_at=now,
  226. )
  227. update_action_item(
  228. int(item["id"]),
  229. execution_status="SUCCESS",
  230. readback_json=readback,
  231. executed_at=now,
  232. )
  233. def _execute_scale(
  234. item: dict[str, Any],
  235. account: dict[str, Any],
  236. *,
  237. client: TencentClient,
  238. config: RoiConfig,
  239. now: datetime,
  240. ) -> None:
  241. account_id = int(item["account_id"])
  242. adgroup_id = int(item["adgroup_id"])
  243. ad = client.get_ad(account_id, adgroup_id)
  244. if str(ad.get("configured_status") or "") != ACTIVE_STATUS:
  245. update_action_item(
  246. int(item["id"]),
  247. execution_status="SKIPPED_STATE_MISMATCH",
  248. skip_reason="ad is no longer active",
  249. pre_state_json=ad,
  250. executed_at=now,
  251. )
  252. return
  253. if item["execution_status"] == "PREPARED":
  254. try:
  255. _resume_prepared_scale(item, ad, client=client, now=now)
  256. except Exception as exc:
  257. update_action_item(
  258. int(item["id"]),
  259. execution_status=_failure_status(exc),
  260. error_message=str(exc),
  261. executed_at=now,
  262. )
  263. return
  264. states = load_ad_states(account_id)
  265. state = states.get(adgroup_id) or {}
  266. last_scaled = state.get("last_roi_scaled_at")
  267. if last_scaled and now < last_scaled + timedelta(days=config.scale_cooldown_days):
  268. update_action_item(
  269. int(item["id"]),
  270. execution_status="SKIPPED_COOLDOWN",
  271. skip_reason=f"last ROI scale at {last_scaled}",
  272. pre_state_json=ad,
  273. executed_at=now,
  274. )
  275. return
  276. bid_field = resolve_bid_field(ad, account.get("bid_scene"))
  277. current = current_bid_fen(ad, bid_field)
  278. if current is None or current <= 0:
  279. raise ValueError(f"ad has no valid {bid_field}")
  280. base = int(state.get("base_bid_fen") or current)
  281. initial = int(state.get("initial_base_bid_fen") or base)
  282. intraday_ratio = Decimal(os.getenv("RTC_BID_UP_RATIO", "1.25"))
  283. plan = plan_scale_bid(
  284. current_bid_fen=current,
  285. base_bid_fen=base,
  286. initial_base_bid_fen=initial,
  287. boosted_date=state.get("boosted_date"),
  288. control_date=now.date(),
  289. intraday_ratio=intraday_ratio,
  290. scale_ratio=config.scale_ratio,
  291. max_base_ratio=config.max_base_ratio,
  292. )
  293. if not plan["allowed"]:
  294. update_action_item(
  295. int(item["id"]),
  296. bid_field=bid_field,
  297. initial_base_bid_fen=initial,
  298. base_bid_fen=base,
  299. before_bid_fen=current,
  300. target_bid_fen=current,
  301. execution_status=(
  302. "SKIPPED_CAP_REACHED"
  303. if plan.get("cap_reached")
  304. else "SKIPPED_STATE_MISMATCH"
  305. ),
  306. skip_reason=str(plan["reason"]),
  307. pre_state_json=ad,
  308. executed_at=now,
  309. )
  310. return
  311. new_base = int(plan["new_base_bid_fen"])
  312. target = int(plan["target_bid_fen"])
  313. boosted_date_today = bool(plan["boosted_date_today"])
  314. prepared_state = {
  315. **ad,
  316. "managed_base_bid_fen": base,
  317. "new_base_bid_fen": new_base,
  318. "boosted_date_today": boosted_date_today,
  319. "preserve_intraday_boost": bool(plan["preserve_intraday_boost"]),
  320. }
  321. update_action_item(
  322. int(item["id"]),
  323. bid_field=bid_field,
  324. initial_base_bid_fen=initial,
  325. base_bid_fen=base,
  326. before_bid_fen=current,
  327. target_bid_fen=target,
  328. execution_status="PREPARED",
  329. pre_state_json=prepared_state,
  330. )
  331. try:
  332. readback = client.update_ad(
  333. account_id,
  334. adgroup_id,
  335. bid_field=bid_field,
  336. target_bid_fen=target,
  337. )
  338. sync_roi_base_bid(
  339. account_id=account_id,
  340. adgroup_id=adgroup_id,
  341. adgroup_name=str(ad.get("adgroup_name") or ""),
  342. bid_field=bid_field,
  343. initial_base_bid_fen=initial,
  344. new_base_bid_fen=new_base,
  345. boosted_date=now.date() if boosted_date_today else None,
  346. action_at=now,
  347. )
  348. update_action_item(
  349. int(item["id"]),
  350. execution_status="SUCCESS",
  351. readback_json=readback,
  352. executed_at=now,
  353. )
  354. except Exception as exc:
  355. update_action_item(
  356. int(item["id"]),
  357. execution_status=_failure_status(exc),
  358. error_message=str(exc),
  359. readback_json=(exc.actual if isinstance(exc, PostWriteVerificationError) else None),
  360. executed_at=now,
  361. )
  362. def execute_roi_batch(
  363. run_id: str,
  364. *,
  365. sender_open_id: str,
  366. now: datetime,
  367. lock_name: str,
  368. tencent: TencentClient | None = None,
  369. ) -> dict[str, Any]:
  370. config = RoiConfig.from_env()
  371. if not config.apply_enabled:
  372. raise RuntimeError("ROI_APPLY_ENABLED=0, ROI Tencent writes are disabled")
  373. with advisory_lock(lock_name) as acquired:
  374. if not acquired:
  375. raise RuntimeError("实时控制正在执行,请稍后再次确认 ROI 批次")
  376. run = claim_run_for_execution(
  377. run_id,
  378. sender_open_id=sender_open_id,
  379. now=now,
  380. )
  381. if run.get("status") in FINAL_STATUSES:
  382. return run
  383. accounts = _managed_accounts()
  384. client = tencent or TencentClient()
  385. for item in load_action_items(run_id):
  386. if item["execution_status"] not in {"PENDING", "PREPARED"}:
  387. continue
  388. account_id = int(item["account_id"])
  389. account = accounts.get(account_id)
  390. if not account:
  391. update_action_item(
  392. int(item["id"]),
  393. execution_status="SKIPPED_NOT_MANAGED",
  394. skip_reason="account left automated management scope",
  395. executed_at=now,
  396. )
  397. continue
  398. try:
  399. if item["action_type"] == ACTION_PAUSE_CREATIVE:
  400. _execute_pause(item, client=client, now=now)
  401. elif item["action_type"] == ACTION_SCALE_BID:
  402. _execute_scale(
  403. item,
  404. account,
  405. client=client,
  406. config=config,
  407. now=now,
  408. )
  409. else:
  410. update_action_item(
  411. int(item["id"]),
  412. execution_status="SKIPPED_UNSUPPORTED",
  413. skip_reason=f"unsupported action={item['action_type']}",
  414. executed_at=now,
  415. )
  416. except Exception as exc:
  417. logger.exception(
  418. "ROI action failed run=%s item=%s", run_id, item["id"]
  419. )
  420. update_action_item(
  421. int(item["id"]),
  422. execution_status=_failure_status(exc),
  423. error_message=str(exc),
  424. executed_at=now,
  425. )
  426. return {"run_id": run_id, **finalize_run(run_id, now=now)}
  427. def reject_roi_batch(
  428. run_id: str,
  429. *,
  430. sender_open_id: str,
  431. now: datetime,
  432. ) -> dict[str, Any]:
  433. return reject_run(run_id, sender_open_id=sender_open_id, now=now)
  434. def roi_batch_status(run_id: str) -> dict[str, Any]:
  435. run = load_run(run_id)
  436. if not run:
  437. raise ValueError(f"ROI批次不存在: {run_id}")
  438. return run