service.py 10 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317
  1. """Orchestrate one idempotent daily ROI metric and approval batch."""
  2. from __future__ import annotations
  3. import hashlib
  4. import logging
  5. import os
  6. import re
  7. from datetime import datetime, timedelta
  8. from pathlib import Path
  9. from typing import Any
  10. from zoneinfo import ZoneInfo
  11. from storage import initialize_schema, load_managed_accounts
  12. from .config import RoiConfig
  13. from .data_source import (
  14. ODPSClient,
  15. date_window,
  16. fetch_ad_age,
  17. fetch_daily_data,
  18. resolve_end_date,
  19. )
  20. from .feishu import RoiFeishuPublisher
  21. from .metrics import ENTITY_QIWEI, METRIC_RUN_SUFFIX, METRIC_VERSION
  22. from .fission_multiplier import (
  23. DEFAULT_FISSION_PARAMETER_VERSION,
  24. FissionMultiplierParameters,
  25. QIWEI_REFERENCE_MULTIPLIER,
  26. QIWEI_REFERENCE_RUN_SUFFIX,
  27. QIWEI_REFERENCE_VERSION,
  28. parameters_from_database,
  29. )
  30. from .policy import annotate_execution
  31. from .reporting import REPORT_RUN_SUFFIX, REPORT_VERSION, write_workbook
  32. from .rules import (
  33. POLICY_RUN_SUFFIX,
  34. POLICY_VERSION,
  35. RuleConfig,
  36. evaluate_rules,
  37. )
  38. from .repository import (
  39. FINAL_STATUSES,
  40. create_or_load_run,
  41. load_fission_parameter_release,
  42. mark_failed,
  43. mark_published,
  44. replace_run_results,
  45. )
  46. SHANGHAI = ZoneInfo("Asia/Shanghai")
  47. logger = logging.getLogger("auto_put_ad_mini.roi_control")
  48. def _rule_config(config: RoiConfig) -> RuleConfig:
  49. return RuleConfig(
  50. self_stop_min_age=config.self_stop_min_age,
  51. self_up_min_age=config.self_up_min_age,
  52. self_min_avg_uv=config.self_min_avg_uv,
  53. partner_min_avg_uv=config.partner_min_avg_uv,
  54. min_daily_cost=config.min_daily_cost,
  55. stop_quantile=config.stop_quantile,
  56. up_quantile=config.up_quantile,
  57. gzh_adjust_rank_min=config.gzh_adjust_rank_min,
  58. gzh_adjust_rank_max=config.gzh_adjust_rank_max,
  59. )
  60. def _dates(start_date: str) -> list[str]:
  61. start = datetime.strptime(start_date, "%Y%m%d")
  62. return [
  63. (start + timedelta(days=offset)).strftime("%Y%m%d")
  64. for offset in range(3)
  65. ]
  66. def _run_identity(
  67. end_date: str,
  68. fission_parameters: FissionMultiplierParameters,
  69. source_revision: str | None = None,
  70. ) -> tuple[str, str]:
  71. if source_revision and not re.fullmatch(
  72. r"[a-z0-9][a-z0-9_-]{0,63}", source_revision
  73. ):
  74. raise ValueError(
  75. "source_revision must use lowercase letters, numbers, '_' or '-'"
  76. )
  77. release = fission_parameters.release
  78. run_suffixes = [release.run_suffix, QIWEI_REFERENCE_RUN_SUFFIX]
  79. run_key_versions = [release.version, QIWEI_REFERENCE_VERSION]
  80. run_id = (
  81. f"roi_{end_date}_{METRIC_RUN_SUFFIX}_{POLICY_RUN_SUFFIX}_"
  82. f"{'_'.join(run_suffixes)}_{REPORT_RUN_SUFFIX}"
  83. )
  84. run_key = (
  85. f"{METRIC_VERSION}:{POLICY_VERSION}:{':'.join(run_key_versions)}:"
  86. f"{REPORT_VERSION}:{end_date}"
  87. )
  88. if source_revision:
  89. run_id = f"{run_id}_{source_revision}"
  90. revision_hash = hashlib.sha256(source_revision.encode("utf-8")).hexdigest()
  91. run_key = f"{run_key}:sr:{revision_hash[:16]}"
  92. if len(run_id) > 64:
  93. raise ValueError("ROI run_id exceeds database limit")
  94. if len(run_key) > 128:
  95. raise ValueError("ROI run_key exceeds database limit")
  96. return run_id, run_key
  97. def run_daily_roi(
  98. *,
  99. requested_end_date: str | None = None,
  100. output_dir: Path,
  101. send_feishu: bool,
  102. now: datetime | None = None,
  103. source_revision: str | None = None,
  104. ) -> dict[str, Any]:
  105. """Compute, snapshot, report, and optionally publish one ROI batch."""
  106. config = RoiConfig.from_env()
  107. initialize_schema()
  108. fission_version = os.getenv(
  109. "ROI_FISSION_PARAMETER_VERSION",
  110. DEFAULT_FISSION_PARAMETER_VERSION,
  111. )
  112. fission_parameters = parameters_from_database(
  113. *load_fission_parameter_release(fission_version)
  114. )
  115. run_config = config.snapshot()
  116. run_config["fission_multiplier"] = fission_parameters.snapshot()
  117. run_config["qiwei_reference_multiplier"] = QIWEI_REFERENCE_MULTIPLIER
  118. run_config["qiwei_reference_version"] = QIWEI_REFERENCE_VERSION
  119. run_config["report_version"] = REPORT_VERSION
  120. if source_revision:
  121. run_config["source_revision"] = source_revision
  122. client = ODPSClient(project=os.getenv("ODPS_PROJECT", "loghubods"))
  123. end_date = resolve_end_date(client, requested_end_date)
  124. start_date, end_date = date_window(end_date)
  125. expected_dates = _dates(start_date)
  126. run_id, run_key = _run_identity(
  127. end_date,
  128. fission_parameters,
  129. source_revision,
  130. )
  131. run = create_or_load_run(
  132. {
  133. "run_id": run_id,
  134. "run_key": run_key,
  135. "metric_version": METRIC_VERSION,
  136. "policy_version": POLICY_VERSION,
  137. "fission_parameter_version": fission_parameters.release.version,
  138. "fission_cohort_date": datetime.strptime(
  139. fission_parameters.release.cohort_date, "%Y%m%d"
  140. ).date(),
  141. "start_date": datetime.strptime(start_date, "%Y%m%d").date(),
  142. "end_date": datetime.strptime(end_date, "%Y%m%d").date(),
  143. "config": run_config,
  144. }
  145. )
  146. reusable_statuses = FINAL_STATUSES | {"PENDING_APPROVAL", "EXECUTING"}
  147. if run.get("status") in reusable_statuses or (
  148. run.get("status") == "COMPUTED" and not send_feishu
  149. ):
  150. logger.info("Reuse ROI run=%s status=%s", run_id, run.get("status"))
  151. return {
  152. "run_id": run_id,
  153. "status": run.get("status"),
  154. "reused": True,
  155. "sheet_url": run.get("sheet_url"),
  156. }
  157. try:
  158. logger.info(
  159. "ROI %s loading ODPS daily data %s..%s fission_parameter=%s",
  160. run_id,
  161. start_date,
  162. end_date,
  163. fission_parameters.release.version,
  164. )
  165. daily = fetch_daily_data(client, start_date, end_date)
  166. ad_age = fetch_ad_age(client, end_date)
  167. _, thresholds, summary = evaluate_rules(
  168. daily,
  169. expected_dates,
  170. ad_age,
  171. _rule_config(config),
  172. fission_parameters=fission_parameters,
  173. )
  174. managed_ids = {
  175. int(row["account_id"])
  176. for row in load_managed_accounts()
  177. }
  178. annotated, snapshots, actions = annotate_execution(
  179. summary,
  180. managed_ids,
  181. run_id,
  182. )
  183. fission_match_summary = (
  184. summary.groupby(
  185. ["entity_type", "传播裂变系数匹配层级"],
  186. as_index=False,
  187. dropna=False,
  188. )
  189. .size()
  190. .rename(
  191. columns={
  192. "entity_type": "实体类型",
  193. "传播裂变系数匹配层级": "匹配层级",
  194. "size": "实体数",
  195. }
  196. )
  197. )
  198. channel_totals = fission_match_summary.groupby("实体类型")[
  199. "实体数"
  200. ].transform("sum")
  201. fission_match_summary["渠道实体数"] = channel_totals
  202. fission_match_summary["匹配率"] = (
  203. fission_match_summary["实体数"] / channel_totals
  204. )
  205. fission_match_summary["参数版本"] = fission_parameters.release.version
  206. actionable_candidates = annotated[annotated["动作"].ne("")].copy()
  207. report_rows = annotated[
  208. annotated["阈值样本状态"].eq("进入阈值样本池")
  209. | annotated["entity_type"].eq(ENTITY_QIWEI)
  210. ].copy()
  211. for row in snapshots:
  212. row["run_id"] = run_id
  213. for row in actions:
  214. row["run_id"] = run_id
  215. threshold_record = thresholds.iloc[0].to_dict()
  216. replace_run_results(
  217. run_id,
  218. snapshots=snapshots,
  219. actions=actions,
  220. thresholds=threshold_record,
  221. )
  222. output_dir.mkdir(parents=True, exist_ok=True)
  223. batch_name = f"ROI调控_{start_date}-{end_date}"
  224. if source_revision:
  225. batch_name = f"{batch_name}_{source_revision}"
  226. output_path = output_dir / f"{batch_name}.xlsx"
  227. write_workbook(
  228. report_rows,
  229. thresholds,
  230. expected_dates,
  231. output_path,
  232. run_config,
  233. fission_match_summary,
  234. )
  235. result: dict[str, Any] = {
  236. "run_id": run_id,
  237. "batch_name": batch_name,
  238. "status": "COMPUTED",
  239. "reused": False,
  240. "start_date": start_date,
  241. "end_date": end_date,
  242. "source_revision": source_revision,
  243. "entity_count": len(snapshots),
  244. "candidate_count": len(actionable_candidates),
  245. "actionable_count": len(actions),
  246. "thresholds": threshold_record,
  247. "fission_multiplier_matches": fission_match_summary.to_dict(
  248. "records"
  249. ),
  250. "report": str(output_path),
  251. }
  252. if not send_feishu:
  253. return result
  254. counts = actionable_candidates.groupby("动作").size().to_dict()
  255. summary_text = (
  256. f"统计窗口:{start_date} - {end_date}\n"
  257. f"关停建议:{counts.get('关停', 0)}\n"
  258. f"扩量建议:{counts.get('扩量', 0)}\n"
  259. f"素材调整建议:{counts.get('调整封面&落地页视频', 0)}\n"
  260. f"可执行动作:{len(actions)}\n"
  261. f"审批有效期:发送后 {config.approval_ttl_minutes} 分钟"
  262. )
  263. publisher = RoiFeishuPublisher()
  264. try:
  265. published = publisher.publish(
  266. output_path,
  267. run_id=run_id,
  268. batch_name=batch_name,
  269. summary=summary_text,
  270. requires_approval=bool(actions),
  271. )
  272. finally:
  273. publisher.close()
  274. expires_at = (
  275. (now or datetime.now(SHANGHAI))
  276. + timedelta(minutes=config.approval_ttl_minutes)
  277. if actions
  278. else None
  279. )
  280. mark_published(
  281. run_id,
  282. sheet_token=published["sheet_token"],
  283. sheet_url=published["url"],
  284. message_id=published["message_id"],
  285. expires_at=expires_at,
  286. requires_approval=bool(actions),
  287. )
  288. result.update(
  289. {
  290. "status": "PENDING_APPROVAL" if actions else "COMPLETED",
  291. "sheet_url": published["url"],
  292. "expires_at": expires_at.isoformat() if expires_at else None,
  293. }
  294. )
  295. return result
  296. except Exception as exc:
  297. mark_failed(run_id, str(exc))
  298. raise