| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544 |
- """Execute an approved ROI batch with Tencent read-back and state coordination."""
- from __future__ import annotations
- import json
- import logging
- import os
- from datetime import date, datetime, timedelta
- from decimal import Decimal, ROUND_HALF_UP
- from typing import Any
- from storage import (
- advisory_lock,
- load_ad_states,
- load_managed_accounts,
- sync_roi_base_bid,
- )
- from tencent_client import (
- ACTIVE_STATUS,
- SUSPEND_STATUS,
- PostWriteVerificationError,
- TencentClient,
- TencentWriteOutcomeUnknownError,
- current_bid_fen,
- resolve_bid_field,
- )
- from .config import RoiConfig
- from .policy import ACTION_PAUSE_CREATIVE, ACTION_SCALE_BID
- from .repository import (
- FINAL_STATUSES,
- claim_run_for_execution,
- finalize_run,
- finalize_sheet_run_if_resolved,
- load_actions_by_ids,
- load_action_items,
- load_run,
- reject_run,
- update_action_item,
- )
- logger = logging.getLogger("auto_put_ad_mini.roi_execution")
- def _round_fen(value: Decimal) -> int:
- return int(value.quantize(Decimal("1"), rounding=ROUND_HALF_UP))
- def _managed_accounts() -> dict[int, dict[str, Any]]:
- return {
- int(row["account_id"]): row
- for row in load_managed_accounts()
- }
- def _failure_status(exc: Exception) -> str:
- if isinstance(exc, PostWriteVerificationError):
- return "VERIFY_FAILED"
- if isinstance(exc, TencentWriteOutcomeUnknownError):
- return "OUTCOME_UNKNOWN"
- return "FAILED"
- def plan_scale_bid(
- *,
- current_bid_fen: int,
- base_bid_fen: int,
- initial_base_bid_fen: int,
- boosted_date: date | None,
- control_date: date,
- intraday_ratio: Decimal,
- scale_ratio: Decimal,
- max_base_ratio: Decimal,
- ) -> dict[str, Any]:
- """Plan a permanent base increase without reapplying a restored boost."""
- boosted_date_today = boosted_date == control_date
- expected_boosted = _round_fen(Decimal(base_bid_fen) * intraday_ratio)
- actively_boosted = boosted_date_today and current_bid_fen == expected_boosted
- if current_bid_fen != base_bid_fen and not actively_boosted:
- return {
- "allowed": False,
- "reason": (
- f"current bid={current_bid_fen} does not match managed "
- f"base={base_bid_fen} or intraday boost={expected_boosted}"
- ),
- }
- max_base = _round_fen(Decimal(initial_base_bid_fen) * max_base_ratio)
- new_base = min(
- _round_fen(Decimal(base_bid_fen) * scale_ratio),
- max_base,
- )
- if new_base <= base_bid_fen:
- return {
- "allowed": False,
- "reason": f"permanent base reached cap={max_base}",
- "cap_reached": True,
- "max_base_bid_fen": max_base,
- }
- target = (
- _round_fen(Decimal(new_base) * intraday_ratio)
- if actively_boosted
- else new_base
- )
- return {
- "allowed": True,
- "new_base_bid_fen": new_base,
- "target_bid_fen": target,
- "boosted_date_today": boosted_date_today,
- "preserve_intraday_boost": actively_boosted,
- }
- def _execute_pause(
- item: dict[str, Any],
- *,
- client: TencentClient,
- now: datetime,
- ) -> None:
- account_id = int(item["account_id"])
- adgroup_id = int(item["adgroup_id"])
- creative_id = int(item["dynamic_creative_id"])
- creative = client.get_dynamic_creative(account_id, creative_id)
- actual_adgroup_id = int(creative.get("adgroup_id") or 0)
- before_status = str(creative.get("configured_status") or "")
- if actual_adgroup_id != adgroup_id:
- update_action_item(
- int(item["id"]),
- execution_status="SKIPPED_MAPPING_MISMATCH",
- skip_reason=(
- f"creative belongs to adgroup={actual_adgroup_id}, expected={adgroup_id}"
- ),
- pre_state_json=creative,
- executed_at=now,
- )
- return
- if item["execution_status"] == "PREPARED" and before_status == SUSPEND_STATUS:
- update_action_item(
- int(item["id"]),
- before_status=item.get("before_status"),
- target_status=SUSPEND_STATUS,
- readback_status=before_status,
- execution_status="SUCCESS",
- readback_json=creative,
- executed_at=now,
- )
- return
- if before_status != ACTIVE_STATUS:
- update_action_item(
- int(item["id"]),
- before_status=before_status,
- target_status=SUSPEND_STATUS,
- readback_status=before_status,
- execution_status="SKIPPED_STATE_MISMATCH",
- skip_reason="creative is no longer active",
- pre_state_json=creative,
- executed_at=now,
- )
- return
- update_action_item(
- int(item["id"]),
- before_status=before_status,
- target_status=SUSPEND_STATUS,
- execution_status="PREPARED",
- pre_state_json=creative,
- )
- try:
- readback = client.update_dynamic_creative_status(
- account_id,
- creative_id,
- SUSPEND_STATUS,
- )
- update_action_item(
- int(item["id"]),
- readback_status=str(readback.get("configured_status") or ""),
- execution_status="SUCCESS",
- readback_json=readback,
- executed_at=now,
- )
- except Exception as exc:
- update_action_item(
- int(item["id"]),
- execution_status=_failure_status(exc),
- error_message=str(exc),
- readback_json=(exc.actual if isinstance(exc, PostWriteVerificationError) else None),
- executed_at=now,
- )
- def _resume_prepared_scale(
- item: dict[str, Any],
- ad: dict[str, Any],
- *,
- client: TencentClient,
- now: datetime,
- ) -> None:
- details = json.loads(item.get("pre_state_json") or "{}")
- bid_field = str(item.get("bid_field") or "")
- current = current_bid_fen(ad, bid_field)
- before = int(item.get("before_bid_fen") or 0)
- target = int(item.get("target_bid_fen") or 0)
- if current == target:
- sync_roi_base_bid(
- account_id=int(item["account_id"]),
- adgroup_id=int(item["adgroup_id"]),
- adgroup_name=str(ad.get("adgroup_name") or ""),
- bid_field=bid_field,
- initial_base_bid_fen=int(item["initial_base_bid_fen"]),
- new_base_bid_fen=int(details["new_base_bid_fen"]),
- boosted_date=(
- now.date() if details.get("boosted_date_today") else None
- ),
- action_at=now,
- )
- update_action_item(
- int(item["id"]),
- execution_status="SUCCESS",
- readback_json=ad,
- executed_at=now,
- )
- return
- if current != before:
- update_action_item(
- int(item["id"]),
- execution_status="SKIPPED_STATE_MISMATCH",
- skip_reason=f"bid changed after preparation: before={before} current={current}",
- readback_json=ad,
- executed_at=now,
- )
- return
- readback = client.update_ad(
- int(item["account_id"]),
- int(item["adgroup_id"]),
- bid_field=bid_field,
- target_bid_fen=target,
- )
- sync_roi_base_bid(
- account_id=int(item["account_id"]),
- adgroup_id=int(item["adgroup_id"]),
- adgroup_name=str(ad.get("adgroup_name") or ""),
- bid_field=bid_field,
- initial_base_bid_fen=int(item["initial_base_bid_fen"]),
- new_base_bid_fen=int(details["new_base_bid_fen"]),
- boosted_date=(now.date() if details.get("boosted_date_today") else None),
- action_at=now,
- )
- update_action_item(
- int(item["id"]),
- execution_status="SUCCESS",
- readback_json=readback,
- executed_at=now,
- )
- def _execute_scale(
- item: dict[str, Any],
- account: dict[str, Any],
- *,
- client: TencentClient,
- config: RoiConfig,
- now: datetime,
- ) -> None:
- account_id = int(item["account_id"])
- adgroup_id = int(item["adgroup_id"])
- ad = client.get_ad(account_id, adgroup_id)
- if str(ad.get("configured_status") or "") != ACTIVE_STATUS:
- update_action_item(
- int(item["id"]),
- execution_status="SKIPPED_STATE_MISMATCH",
- skip_reason="ad is no longer active",
- pre_state_json=ad,
- executed_at=now,
- )
- return
- if item["execution_status"] == "PREPARED":
- try:
- _resume_prepared_scale(item, ad, client=client, now=now)
- except Exception as exc:
- update_action_item(
- int(item["id"]),
- execution_status=_failure_status(exc),
- error_message=str(exc),
- executed_at=now,
- )
- return
- states = load_ad_states(account_id)
- state = states.get(adgroup_id) or {}
- last_scaled = state.get("last_roi_scaled_at")
- if last_scaled and now < last_scaled + timedelta(days=config.scale_cooldown_days):
- update_action_item(
- int(item["id"]),
- execution_status="SKIPPED_COOLDOWN",
- skip_reason=f"last ROI scale at {last_scaled}",
- pre_state_json=ad,
- executed_at=now,
- )
- return
- bid_field = resolve_bid_field(ad, account.get("bid_scene"))
- current = current_bid_fen(ad, bid_field)
- if current is None or current <= 0:
- raise ValueError(f"ad has no valid {bid_field}")
- base = int(state.get("base_bid_fen") or current)
- initial = int(state.get("initial_base_bid_fen") or base)
- intraday_ratio = Decimal(os.getenv("RTC_BID_UP_RATIO", "1.05"))
- plan = plan_scale_bid(
- current_bid_fen=current,
- base_bid_fen=base,
- initial_base_bid_fen=initial,
- boosted_date=state.get("boosted_date"),
- control_date=now.date(),
- intraday_ratio=intraday_ratio,
- scale_ratio=config.scale_ratio,
- max_base_ratio=config.max_base_ratio,
- )
- if not plan["allowed"]:
- update_action_item(
- int(item["id"]),
- bid_field=bid_field,
- initial_base_bid_fen=initial,
- base_bid_fen=base,
- before_bid_fen=current,
- target_bid_fen=current,
- execution_status=(
- "SKIPPED_CAP_REACHED"
- if plan.get("cap_reached")
- else "SKIPPED_STATE_MISMATCH"
- ),
- skip_reason=str(plan["reason"]),
- pre_state_json=ad,
- executed_at=now,
- )
- return
- new_base = int(plan["new_base_bid_fen"])
- target = int(plan["target_bid_fen"])
- boosted_date_today = bool(plan["boosted_date_today"])
- prepared_state = {
- **ad,
- "managed_base_bid_fen": base,
- "new_base_bid_fen": new_base,
- "boosted_date_today": boosted_date_today,
- "preserve_intraday_boost": bool(plan["preserve_intraday_boost"]),
- }
- update_action_item(
- int(item["id"]),
- bid_field=bid_field,
- initial_base_bid_fen=initial,
- base_bid_fen=base,
- before_bid_fen=current,
- target_bid_fen=target,
- execution_status="PREPARED",
- pre_state_json=prepared_state,
- )
- try:
- readback = client.update_ad(
- account_id,
- adgroup_id,
- bid_field=bid_field,
- target_bid_fen=target,
- )
- sync_roi_base_bid(
- account_id=account_id,
- adgroup_id=adgroup_id,
- adgroup_name=str(ad.get("adgroup_name") or ""),
- bid_field=bid_field,
- initial_base_bid_fen=initial,
- new_base_bid_fen=new_base,
- boosted_date=now.date() if boosted_date_today else None,
- action_at=now,
- )
- update_action_item(
- int(item["id"]),
- execution_status="SUCCESS",
- readback_json=readback,
- executed_at=now,
- )
- except Exception as exc:
- update_action_item(
- int(item["id"]),
- execution_status=_failure_status(exc),
- error_message=str(exc),
- readback_json=(exc.actual if isinstance(exc, PostWriteVerificationError) else None),
- executed_at=now,
- )
- def execute_roi_batch(
- run_id: str,
- *,
- sender_open_id: str,
- now: datetime,
- lock_name: str,
- tencent: TencentClient | None = None,
- ) -> dict[str, Any]:
- config = RoiConfig.from_env()
- if not config.apply_enabled:
- raise RuntimeError("ROI_APPLY_ENABLED=0, ROI Tencent writes are disabled")
- with advisory_lock(lock_name) as acquired:
- if not acquired:
- raise RuntimeError("实时控制正在执行,请稍后再次确认 ROI 批次")
- run = claim_run_for_execution(
- run_id,
- sender_open_id=sender_open_id,
- now=now,
- )
- if run.get("status") in FINAL_STATUSES:
- return run
- accounts = _managed_accounts()
- client = tencent or TencentClient()
- for item in load_action_items(run_id):
- if item["execution_status"] not in {"PENDING", "PREPARED"}:
- continue
- account_id = int(item["account_id"])
- account = accounts.get(account_id)
- if account is None:
- update_action_item(
- int(item["id"]),
- execution_status="SKIPPED_NOT_MANAGED",
- skip_reason="account left automated management scope",
- executed_at=now,
- )
- continue
- try:
- if item["action_type"] == ACTION_PAUSE_CREATIVE:
- _execute_pause(item, client=client, now=now)
- elif item["action_type"] == ACTION_SCALE_BID:
- _execute_scale(
- item,
- account,
- client=client,
- config=config,
- now=now,
- )
- else:
- update_action_item(
- int(item["id"]),
- execution_status="SKIPPED_UNSUPPORTED",
- skip_reason=f"unsupported action={item['action_type']}",
- executed_at=now,
- )
- except Exception as exc:
- logger.exception(
- "ROI action failed run=%s item=%s", run_id, item["id"]
- )
- update_action_item(
- int(item["id"]),
- execution_status=_failure_status(exc),
- error_message=str(exc),
- executed_at=now,
- )
- return {"run_id": run_id, **finalize_run(run_id, now=now)}
- def execute_approved_roi_actions(
- item_ids: list[int],
- *,
- now: datetime,
- lock_name: str,
- tencent: TencentClient | None = None,
- ) -> dict[str, Any]:
- """Execute only explicitly approved sheet rows without claiming the batch."""
- config = RoiConfig.from_env()
- if not config.apply_enabled:
- raise RuntimeError("ROI_APPLY_ENABLED=0, ROI Tencent writes are disabled")
- with advisory_lock(lock_name) as acquired:
- if not acquired:
- raise RuntimeError("实时控制正在执行,ROI表格审批稍后自动重试")
- accounts = _managed_accounts()
- client = tencent or TencentClient()
- run_ids: set[str] = set()
- processed = 0
- for item in load_actions_by_ids(item_ids):
- run_ids.add(str(item["run_id"]))
- if item.get("approval_status") != "APPROVED":
- continue
- if item["execution_status"] not in {"PENDING", "PREPARED"}:
- continue
- processed += 1
- account_id = int(item["account_id"])
- account = accounts.get(account_id)
- if account is None:
- update_action_item(
- int(item["id"]),
- execution_status="SKIPPED_NOT_MANAGED",
- skip_reason="account left automated management scope",
- executed_at=now,
- )
- continue
- try:
- if item["action_type"] == ACTION_PAUSE_CREATIVE:
- _execute_pause(item, client=client, now=now)
- elif item["action_type"] == ACTION_SCALE_BID:
- _execute_scale(
- item,
- account,
- client=client,
- config=config,
- now=now,
- )
- else:
- update_action_item(
- int(item["id"]),
- execution_status="SKIPPED_UNSUPPORTED",
- skip_reason=f"unsupported action={item['action_type']}",
- executed_at=now,
- )
- except Exception as exc:
- logger.exception(
- "Approved ROI action failed run=%s item=%s",
- item["run_id"],
- item["id"],
- )
- update_action_item(
- int(item["id"]),
- execution_status=_failure_status(exc),
- error_message=str(exc),
- executed_at=now,
- )
- final = {
- run_id: finalize_sheet_run_if_resolved(run_id, now=now)
- for run_id in run_ids
- }
- return {"processed": processed, "runs": final}
- def reject_roi_batch(
- run_id: str,
- *,
- sender_open_id: str,
- now: datetime,
- ) -> dict[str, Any]:
- return reject_run(run_id, sender_open_id=sender_open_id, now=now)
- def roi_batch_status(run_id: str) -> dict[str, Any]:
- run = load_run(run_id)
- if not run:
- raise ValueError(f"ROI批次不存在: {run_id}")
- return run
|