execution.py 18 KB

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