execution.py 20 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614
  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_AD, 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 _execute_pause_ad(
  173. item: dict[str, Any],
  174. *,
  175. client: TencentClient,
  176. now: datetime,
  177. ) -> None:
  178. """Pause the entire ad after an ad-level ROI row is approved."""
  179. account_id = int(item["account_id"])
  180. adgroup_id = int(item["adgroup_id"])
  181. ad = client.get_ad(account_id, adgroup_id)
  182. before_status = str(ad.get("configured_status") or "")
  183. if item["execution_status"] == "PREPARED" and before_status == SUSPEND_STATUS:
  184. update_action_item(
  185. int(item["id"]),
  186. before_status=item.get("before_status"),
  187. target_status=SUSPEND_STATUS,
  188. readback_status=before_status,
  189. execution_status="SUCCESS",
  190. readback_json=ad,
  191. executed_at=now,
  192. )
  193. return
  194. if before_status != ACTIVE_STATUS:
  195. update_action_item(
  196. int(item["id"]),
  197. before_status=before_status,
  198. target_status=SUSPEND_STATUS,
  199. readback_status=before_status,
  200. execution_status="SKIPPED_STATE_MISMATCH",
  201. skip_reason="ad is no longer active",
  202. pre_state_json=ad,
  203. executed_at=now,
  204. )
  205. return
  206. update_action_item(
  207. int(item["id"]),
  208. before_status=before_status,
  209. target_status=SUSPEND_STATUS,
  210. execution_status="PREPARED",
  211. pre_state_json=ad,
  212. )
  213. try:
  214. readback = client.update_ad(
  215. account_id,
  216. adgroup_id,
  217. target_status=SUSPEND_STATUS,
  218. )
  219. update_action_item(
  220. int(item["id"]),
  221. readback_status=str(readback.get("configured_status") or ""),
  222. execution_status="SUCCESS",
  223. readback_json=readback,
  224. executed_at=now,
  225. )
  226. except Exception as exc:
  227. update_action_item(
  228. int(item["id"]),
  229. execution_status=_failure_status(exc),
  230. error_message=str(exc),
  231. readback_json=(exc.actual if isinstance(exc, PostWriteVerificationError) else None),
  232. executed_at=now,
  233. )
  234. def _resume_prepared_scale(
  235. item: dict[str, Any],
  236. ad: dict[str, Any],
  237. *,
  238. client: TencentClient,
  239. now: datetime,
  240. ) -> None:
  241. details = json.loads(item.get("pre_state_json") or "{}")
  242. bid_field = str(item.get("bid_field") or "")
  243. current = current_bid_fen(ad, bid_field)
  244. before = int(item.get("before_bid_fen") or 0)
  245. target = int(item.get("target_bid_fen") or 0)
  246. if current == target:
  247. sync_roi_base_bid(
  248. account_id=int(item["account_id"]),
  249. adgroup_id=int(item["adgroup_id"]),
  250. adgroup_name=str(ad.get("adgroup_name") or ""),
  251. bid_field=bid_field,
  252. initial_base_bid_fen=int(item["initial_base_bid_fen"]),
  253. new_base_bid_fen=int(details["new_base_bid_fen"]),
  254. boosted_date=(
  255. now.date() if details.get("boosted_date_today") else None
  256. ),
  257. action_at=now,
  258. )
  259. update_action_item(
  260. int(item["id"]),
  261. execution_status="SUCCESS",
  262. readback_json=ad,
  263. executed_at=now,
  264. )
  265. return
  266. if current != before:
  267. update_action_item(
  268. int(item["id"]),
  269. execution_status="SKIPPED_STATE_MISMATCH",
  270. skip_reason=f"bid changed after preparation: before={before} current={current}",
  271. readback_json=ad,
  272. executed_at=now,
  273. )
  274. return
  275. readback = client.update_ad(
  276. int(item["account_id"]),
  277. int(item["adgroup_id"]),
  278. bid_field=bid_field,
  279. target_bid_fen=target,
  280. )
  281. sync_roi_base_bid(
  282. account_id=int(item["account_id"]),
  283. adgroup_id=int(item["adgroup_id"]),
  284. adgroup_name=str(ad.get("adgroup_name") or ""),
  285. bid_field=bid_field,
  286. initial_base_bid_fen=int(item["initial_base_bid_fen"]),
  287. new_base_bid_fen=int(details["new_base_bid_fen"]),
  288. boosted_date=(now.date() if details.get("boosted_date_today") else None),
  289. action_at=now,
  290. )
  291. update_action_item(
  292. int(item["id"]),
  293. execution_status="SUCCESS",
  294. readback_json=readback,
  295. executed_at=now,
  296. )
  297. def _execute_scale(
  298. item: dict[str, Any],
  299. account: dict[str, Any],
  300. *,
  301. client: TencentClient,
  302. config: RoiConfig,
  303. now: datetime,
  304. ) -> None:
  305. account_id = int(item["account_id"])
  306. adgroup_id = int(item["adgroup_id"])
  307. ad = client.get_ad(account_id, adgroup_id)
  308. if str(ad.get("configured_status") or "") != ACTIVE_STATUS:
  309. update_action_item(
  310. int(item["id"]),
  311. execution_status="SKIPPED_STATE_MISMATCH",
  312. skip_reason="ad is no longer active",
  313. pre_state_json=ad,
  314. executed_at=now,
  315. )
  316. return
  317. if item["execution_status"] == "PREPARED":
  318. try:
  319. _resume_prepared_scale(item, ad, client=client, now=now)
  320. except Exception as exc:
  321. update_action_item(
  322. int(item["id"]),
  323. execution_status=_failure_status(exc),
  324. error_message=str(exc),
  325. executed_at=now,
  326. )
  327. return
  328. states = load_ad_states(account_id)
  329. state = states.get(adgroup_id) or {}
  330. last_scaled = state.get("last_roi_scaled_at")
  331. if last_scaled and now < last_scaled + timedelta(days=config.scale_cooldown_days):
  332. update_action_item(
  333. int(item["id"]),
  334. execution_status="SKIPPED_COOLDOWN",
  335. skip_reason=f"last ROI scale at {last_scaled}",
  336. pre_state_json=ad,
  337. executed_at=now,
  338. )
  339. return
  340. bid_field = resolve_bid_field(ad, account.get("bid_scene"))
  341. current = current_bid_fen(ad, bid_field)
  342. if current is None or current <= 0:
  343. raise ValueError(f"ad has no valid {bid_field}")
  344. base = int(state.get("base_bid_fen") or current)
  345. initial = int(state.get("initial_base_bid_fen") or base)
  346. intraday_ratio = Decimal(os.getenv("RTC_BID_UP_RATIO", "1.05"))
  347. plan = plan_scale_bid(
  348. current_bid_fen=current,
  349. base_bid_fen=base,
  350. initial_base_bid_fen=initial,
  351. boosted_date=state.get("boosted_date"),
  352. control_date=now.date(),
  353. intraday_ratio=intraday_ratio,
  354. scale_ratio=config.scale_ratio,
  355. max_base_ratio=config.max_base_ratio,
  356. )
  357. if not plan["allowed"]:
  358. update_action_item(
  359. int(item["id"]),
  360. bid_field=bid_field,
  361. initial_base_bid_fen=initial,
  362. base_bid_fen=base,
  363. before_bid_fen=current,
  364. target_bid_fen=current,
  365. execution_status=(
  366. "SKIPPED_CAP_REACHED"
  367. if plan.get("cap_reached")
  368. else "SKIPPED_STATE_MISMATCH"
  369. ),
  370. skip_reason=str(plan["reason"]),
  371. pre_state_json=ad,
  372. executed_at=now,
  373. )
  374. return
  375. new_base = int(plan["new_base_bid_fen"])
  376. target = int(plan["target_bid_fen"])
  377. boosted_date_today = bool(plan["boosted_date_today"])
  378. prepared_state = {
  379. **ad,
  380. "managed_base_bid_fen": base,
  381. "new_base_bid_fen": new_base,
  382. "boosted_date_today": boosted_date_today,
  383. "preserve_intraday_boost": bool(plan["preserve_intraday_boost"]),
  384. }
  385. update_action_item(
  386. int(item["id"]),
  387. bid_field=bid_field,
  388. initial_base_bid_fen=initial,
  389. base_bid_fen=base,
  390. before_bid_fen=current,
  391. target_bid_fen=target,
  392. execution_status="PREPARED",
  393. pre_state_json=prepared_state,
  394. )
  395. try:
  396. readback = client.update_ad(
  397. account_id,
  398. adgroup_id,
  399. bid_field=bid_field,
  400. target_bid_fen=target,
  401. )
  402. sync_roi_base_bid(
  403. account_id=account_id,
  404. adgroup_id=adgroup_id,
  405. adgroup_name=str(ad.get("adgroup_name") or ""),
  406. bid_field=bid_field,
  407. initial_base_bid_fen=initial,
  408. new_base_bid_fen=new_base,
  409. boosted_date=now.date() if boosted_date_today else None,
  410. action_at=now,
  411. )
  412. update_action_item(
  413. int(item["id"]),
  414. execution_status="SUCCESS",
  415. readback_json=readback,
  416. executed_at=now,
  417. )
  418. except Exception as exc:
  419. update_action_item(
  420. int(item["id"]),
  421. execution_status=_failure_status(exc),
  422. error_message=str(exc),
  423. readback_json=(exc.actual if isinstance(exc, PostWriteVerificationError) else None),
  424. executed_at=now,
  425. )
  426. def execute_roi_batch(
  427. run_id: str,
  428. *,
  429. sender_open_id: str,
  430. now: datetime,
  431. lock_name: str,
  432. tencent: TencentClient | None = None,
  433. ) -> dict[str, Any]:
  434. config = RoiConfig.from_env()
  435. if not config.apply_enabled:
  436. raise RuntimeError("ROI_APPLY_ENABLED=0, ROI Tencent writes are disabled")
  437. with advisory_lock(lock_name) as acquired:
  438. if not acquired:
  439. raise RuntimeError("实时控制正在执行,请稍后再次确认 ROI 批次")
  440. run = claim_run_for_execution(
  441. run_id,
  442. sender_open_id=sender_open_id,
  443. now=now,
  444. )
  445. if run.get("status") in FINAL_STATUSES:
  446. return run
  447. accounts = _managed_accounts()
  448. client = tencent or TencentClient()
  449. for item in load_action_items(run_id):
  450. if item["execution_status"] not in {"PENDING", "PREPARED"}:
  451. continue
  452. account_id = int(item["account_id"])
  453. account = accounts.get(account_id)
  454. if account is None:
  455. update_action_item(
  456. int(item["id"]),
  457. execution_status="SKIPPED_NOT_MANAGED",
  458. skip_reason="account left automated management scope",
  459. executed_at=now,
  460. )
  461. continue
  462. try:
  463. if item["action_type"] == ACTION_PAUSE_CREATIVE:
  464. _execute_pause(item, client=client, now=now)
  465. elif item["action_type"] == ACTION_PAUSE_AD:
  466. _execute_pause_ad(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. "ROI action failed run=%s item=%s", run_id, item["id"]
  485. )
  486. update_action_item(
  487. int(item["id"]),
  488. execution_status=_failure_status(exc),
  489. error_message=str(exc),
  490. executed_at=now,
  491. )
  492. return {"run_id": run_id, **finalize_run(run_id, now=now)}
  493. def execute_approved_roi_actions(
  494. item_ids: list[int],
  495. *,
  496. now: datetime,
  497. lock_name: str,
  498. tencent: TencentClient | None = None,
  499. ) -> dict[str, Any]:
  500. """Execute only explicitly approved sheet rows without claiming the batch."""
  501. config = RoiConfig.from_env()
  502. if not config.apply_enabled:
  503. raise RuntimeError("ROI_APPLY_ENABLED=0, ROI Tencent writes are disabled")
  504. with advisory_lock(lock_name) as acquired:
  505. if not acquired:
  506. raise RuntimeError("实时控制正在执行,ROI表格审批稍后自动重试")
  507. accounts = _managed_accounts()
  508. client = tencent or TencentClient()
  509. run_ids: set[str] = set()
  510. processed = 0
  511. for item in load_actions_by_ids(item_ids):
  512. run_ids.add(str(item["run_id"]))
  513. if item.get("approval_status") != "APPROVED":
  514. continue
  515. if item["execution_status"] not in {"PENDING", "PREPARED"}:
  516. continue
  517. processed += 1
  518. account_id = int(item["account_id"])
  519. account = accounts.get(account_id)
  520. if account is None:
  521. update_action_item(
  522. int(item["id"]),
  523. execution_status="SKIPPED_NOT_MANAGED",
  524. skip_reason="account left automated management scope",
  525. executed_at=now,
  526. )
  527. continue
  528. try:
  529. if item["action_type"] == ACTION_PAUSE_CREATIVE:
  530. _execute_pause(item, client=client, now=now)
  531. elif item["action_type"] == ACTION_PAUSE_AD:
  532. _execute_pause_ad(item, client=client, now=now)
  533. elif item["action_type"] == ACTION_SCALE_BID:
  534. _execute_scale(
  535. item,
  536. account,
  537. client=client,
  538. config=config,
  539. now=now,
  540. )
  541. else:
  542. update_action_item(
  543. int(item["id"]),
  544. execution_status="SKIPPED_UNSUPPORTED",
  545. skip_reason=f"unsupported action={item['action_type']}",
  546. executed_at=now,
  547. )
  548. except Exception as exc:
  549. logger.exception(
  550. "Approved ROI action failed run=%s item=%s",
  551. item["run_id"],
  552. item["id"],
  553. )
  554. update_action_item(
  555. int(item["id"]),
  556. execution_status=_failure_status(exc),
  557. error_message=str(exc),
  558. executed_at=now,
  559. )
  560. final = {
  561. run_id: finalize_sheet_run_if_resolved(run_id, now=now)
  562. for run_id in run_ids
  563. }
  564. return {"processed": processed, "runs": final}
  565. def reject_roi_batch(
  566. run_id: str,
  567. *,
  568. sender_open_id: str,
  569. now: datetime,
  570. ) -> dict[str, Any]:
  571. return reject_run(run_id, sender_open_id=sender_open_id, now=now)
  572. def roi_batch_status(run_id: str) -> dict[str, Any]:
  573. run = load_run(run_id)
  574. if not run:
  575. raise ValueError(f"ROI批次不存在: {run_id}")
  576. return run