service.py 19 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561
  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 tencent_client import ACTIVE_STATUS, TencentClient
  13. from .agency_delivery import agency_route_key, publish_agency_reports
  14. from .config import AgencyWebhookConfig, RoiConfig
  15. from .data_source import (
  16. ODPSClient,
  17. date_window,
  18. fetch_ad_age,
  19. fetch_daily_data,
  20. resolve_end_date,
  21. )
  22. from .feishu import RoiFeishuPublisher
  23. from .metrics import ENTITY_SELF, METRIC_RUN_SUFFIX, METRIC_VERSION
  24. from .fission_multiplier import (
  25. DEFAULT_FISSION_PARAMETER_VERSION,
  26. FissionMultiplierParameters,
  27. parameters_from_database,
  28. )
  29. from .policy import annotate_execution
  30. from .reporting import (
  31. REPORT_RUN_SUFFIX,
  32. REPORT_VERSION,
  33. write_agency_workbooks,
  34. write_workbook,
  35. )
  36. from .rules import (
  37. POLICY_RUN_SUFFIX,
  38. POLICY_VERSION,
  39. RuleConfig,
  40. evaluate_rules,
  41. )
  42. from .repository import (
  43. FINAL_STATUSES,
  44. create_or_load_run,
  45. load_agency_deliveries,
  46. load_fission_parameter_release,
  47. mark_failed,
  48. mark_published,
  49. replace_run_results,
  50. )
  51. SHANGHAI = ZoneInfo("Asia/Shanghai")
  52. logger = logging.getLogger("auto_put_ad_mini.roi_control")
  53. def _redact_agency_webhooks(
  54. error: Exception,
  55. config: AgencyWebhookConfig,
  56. ) -> str:
  57. message = str(error)
  58. for webhook_url in (config.webhooks or {}).values():
  59. message = message.replace(webhook_url, "<redacted>")
  60. return message
  61. def _internal_source_revision(config: RoiConfig, now: datetime) -> str:
  62. if config.internal_test_dedup_enabled:
  63. return "internal_test"
  64. return f"it_{now.strftime('%H%M%S%f')}"
  65. def _build_internal_reports(
  66. *,
  67. agency_reports: list[dict[str, object]],
  68. ) -> list[dict[str, object]]:
  69. reports = [
  70. {
  71. **report,
  72. "title": f"内部测试_{Path(str(report['report'])).stem}",
  73. }
  74. for report in agency_reports
  75. if agency_route_key(str(report["agency_name"])) == "自动化投放"
  76. ]
  77. if len(reports) != 1:
  78. raise RuntimeError(
  79. "Internal ROI test requires exactly one 自动化投放 agency report"
  80. )
  81. return reports
  82. def _internal_webhook_config(
  83. reports: list[dict[str, object]],
  84. webhook_url: str,
  85. ) -> AgencyWebhookConfig:
  86. return AgencyWebhookConfig(
  87. enabled=True,
  88. webhooks={
  89. agency_route_key(str(report["agency_name"])): webhook_url
  90. for report in reports
  91. },
  92. )
  93. def _annotate_current_creative_status(
  94. rows,
  95. tencent: TencentClient | None = None,
  96. ):
  97. """Read current Tencent status only for creative-level stop decisions."""
  98. result = rows.copy()
  99. result["当前创意状态"] = ""
  100. mask = result["entity_type"].eq(ENTITY_SELF) & result["动作"].eq("关停")
  101. if not mask.any():
  102. return result
  103. client = tencent or TencentClient()
  104. cache: dict[tuple[int, int], str] = {}
  105. try:
  106. for index, row in result.loc[mask].iterrows():
  107. try:
  108. account_id = int(row["账号id"])
  109. creative_id = int(row["创意id"])
  110. except (TypeError, ValueError):
  111. result.at[index, "当前创意状态"] = "读取失败"
  112. continue
  113. key = (account_id, creative_id)
  114. if key not in cache:
  115. try:
  116. creative = client.get_dynamic_creative(account_id, creative_id)
  117. status = str(creative.get("configured_status") or "")
  118. cache[key] = (
  119. "正常"
  120. if status == ACTIVE_STATUS
  121. else "已停止" if status else "读取失败"
  122. )
  123. except Exception as exc:
  124. logger.warning(
  125. "Failed to read creative status account=%s creative=%s: %s",
  126. account_id,
  127. creative_id,
  128. exc,
  129. )
  130. cache[key] = "读取失败"
  131. result.at[index, "当前创意状态"] = cache[key]
  132. finally:
  133. if tencent is None:
  134. client.session.close()
  135. return result
  136. def _rule_config(config: RoiConfig) -> RuleConfig:
  137. return RuleConfig(
  138. self_stop_min_age=config.self_stop_min_age,
  139. self_up_min_age=config.self_up_min_age,
  140. self_min_daily_uv=config.self_min_daily_uv,
  141. partner_min_daily_uv=config.partner_min_daily_uv,
  142. observe_min_latest_uv=config.observe_min_latest_uv,
  143. one_day_min_uv=config.one_day_min_uv,
  144. one_day_p30_min_uv=config.one_day_p30_min_uv,
  145. one_day_hard_stop_roi=config.one_day_hard_stop_roi,
  146. one_day_stop_quantile=config.one_day_stop_quantile,
  147. stop_quantile=config.stop_quantile,
  148. up_quantile=config.up_quantile,
  149. )
  150. def _dates(start_date: str) -> list[str]:
  151. start = datetime.strptime(start_date, "%Y%m%d")
  152. return [
  153. (start + timedelta(days=offset)).strftime("%Y%m%d")
  154. for offset in range(3)
  155. ]
  156. def _run_identity(
  157. end_date: str,
  158. fission_parameters: FissionMultiplierParameters,
  159. source_revision: str | None = None,
  160. ) -> tuple[str, str]:
  161. if source_revision and not re.fullmatch(
  162. r"[a-z0-9][a-z0-9_-]{0,63}", source_revision
  163. ):
  164. raise ValueError(
  165. "source_revision must use lowercase letters, numbers, '_' or '-'"
  166. )
  167. release = fission_parameters.release
  168. run_suffixes = [release.run_suffix]
  169. run_key_versions = [release.version]
  170. run_id = (
  171. f"roi_{end_date}_{METRIC_RUN_SUFFIX}_{POLICY_RUN_SUFFIX}_"
  172. f"{'_'.join(run_suffixes)}_{REPORT_RUN_SUFFIX}"
  173. )
  174. run_key = (
  175. f"{METRIC_VERSION}:{POLICY_VERSION}:{':'.join(run_key_versions)}:"
  176. f"{REPORT_VERSION}:{end_date}"
  177. )
  178. if source_revision:
  179. run_id = f"{run_id}_{source_revision}"
  180. revision_hash = hashlib.sha256(source_revision.encode("utf-8")).hexdigest()
  181. run_key = f"{run_key}:sr:{revision_hash[:16]}"
  182. if len(run_id) > 64:
  183. raise ValueError("ROI run_id exceeds database limit")
  184. if len(run_key) > 128:
  185. raise ValueError("ROI run_key exceeds database limit")
  186. return run_id, run_key
  187. def run_daily_roi(
  188. *,
  189. requested_end_date: str | None = None,
  190. output_dir: Path,
  191. send_feishu: bool,
  192. now: datetime | None = None,
  193. source_revision: str | None = None,
  194. internal_test: bool = False,
  195. ) -> dict[str, Any]:
  196. """Compute, snapshot, report, and optionally publish one ROI batch."""
  197. config = RoiConfig.from_env()
  198. agency_webhook_config = AgencyWebhookConfig.from_env()
  199. if internal_test and send_feishu:
  200. raise ValueError("internal_test cannot publish formal Feishu notifications")
  201. if internal_test and not (agency_webhook_config.webhooks or {}).get("内部"):
  202. raise ValueError("Internal ROI test requires agency webhook route: 内部")
  203. effective_now = now or datetime.now(SHANGHAI)
  204. if effective_now.tzinfo is None:
  205. effective_now = effective_now.replace(tzinfo=SHANGHAI)
  206. if internal_test:
  207. if source_revision:
  208. raise ValueError("internal_test manages source_revision automatically")
  209. source_revision = _internal_source_revision(config, effective_now)
  210. client = ODPSClient(project=os.getenv("ODPS_PROJECT", "loghubods"))
  211. end_date = resolve_end_date(
  212. client,
  213. requested_end_date,
  214. now=effective_now,
  215. )
  216. start_date, end_date = date_window(end_date)
  217. expected_dates = _dates(start_date)
  218. initialize_schema()
  219. fission_version = os.getenv(
  220. "ROI_FISSION_PARAMETER_VERSION",
  221. DEFAULT_FISSION_PARAMETER_VERSION,
  222. )
  223. fission_parameters = parameters_from_database(
  224. *load_fission_parameter_release(fission_version)
  225. )
  226. run_config = config.snapshot()
  227. run_config["fission_multiplier"] = fission_parameters.snapshot()
  228. run_config["report_version"] = REPORT_VERSION
  229. run_config["agency_webhook"] = agency_webhook_config.snapshot()
  230. run_config["internal_test"] = internal_test
  231. if source_revision:
  232. run_config["source_revision"] = source_revision
  233. run_id, run_key = _run_identity(
  234. end_date,
  235. fission_parameters,
  236. source_revision,
  237. )
  238. run = create_or_load_run(
  239. {
  240. "run_id": run_id,
  241. "run_key": run_key,
  242. "metric_version": METRIC_VERSION,
  243. "policy_version": POLICY_VERSION,
  244. "fission_parameter_version": fission_parameters.release.version,
  245. "fission_cohort_date": datetime.strptime(
  246. fission_parameters.release.cohort_date, "%Y%m%d"
  247. ).date(),
  248. "start_date": datetime.strptime(start_date, "%Y%m%d").date(),
  249. "end_date": datetime.strptime(end_date, "%Y%m%d").date(),
  250. "config": run_config,
  251. }
  252. )
  253. reusable_statuses = FINAL_STATUSES | {"PENDING_APPROVAL", "EXECUTING"}
  254. publish_requested = send_feishu or internal_test
  255. if run.get("status") in reusable_statuses or (
  256. run.get("status") == "COMPUTED" and not publish_requested
  257. ):
  258. logger.info("Reuse ROI run=%s status=%s", run_id, run.get("status"))
  259. reused_result: dict[str, Any] = {
  260. "run_id": run_id,
  261. "status": run.get("status"),
  262. "reused": True,
  263. "sheet_url": run.get("sheet_url"),
  264. }
  265. if send_feishu and agency_webhook_config.enabled:
  266. stored_deliveries = load_agency_deliveries(run_id)
  267. if stored_deliveries:
  268. retry_reports = [
  269. {
  270. "agency_name": row["agency_name"],
  271. "report_version": row["agency_report_version"],
  272. "report": row["file_path"],
  273. "creative_rows": row.get("creative_rows") or 0,
  274. "ad_rows": row.get("ad_rows") or 0,
  275. }
  276. for row in stored_deliveries
  277. ]
  278. publisher = RoiFeishuPublisher()
  279. try:
  280. try:
  281. reused_result["agency_deliveries"] = publish_agency_reports(
  282. run_id=run_id,
  283. reports=retry_reports,
  284. config=agency_webhook_config,
  285. publisher=publisher,
  286. now=now,
  287. )
  288. except Exception as exc:
  289. safe_error = _redact_agency_webhooks(
  290. exc,
  291. agency_webhook_config,
  292. )
  293. logger.error("Agency ROI retry failed: %s", safe_error)
  294. reused_result["agency_delivery_error"] = safe_error
  295. finally:
  296. publisher.close()
  297. return reused_result
  298. try:
  299. logger.info(
  300. "ROI %s loading ODPS daily data %s..%s fission_parameter=%s",
  301. run_id,
  302. start_date,
  303. end_date,
  304. fission_parameters.release.version,
  305. )
  306. daily = fetch_daily_data(client, start_date, end_date)
  307. ad_age = fetch_ad_age(client, end_date)
  308. _, thresholds, summary = evaluate_rules(
  309. daily,
  310. expected_dates,
  311. ad_age,
  312. _rule_config(config),
  313. fission_parameters=fission_parameters,
  314. )
  315. managed_ids = (
  316. set()
  317. if internal_test
  318. else {
  319. int(row["account_id"])
  320. for row in load_managed_accounts()
  321. }
  322. )
  323. annotated, snapshots, actions = annotate_execution(
  324. summary,
  325. managed_ids,
  326. run_id,
  327. )
  328. annotated = _annotate_current_creative_status(annotated)
  329. fission_match_summary = (
  330. summary.groupby(
  331. ["entity_type", "传播裂变系数匹配层级"],
  332. as_index=False,
  333. dropna=False,
  334. )
  335. .size()
  336. .rename(
  337. columns={
  338. "entity_type": "实体类型",
  339. "传播裂变系数匹配层级": "匹配层级",
  340. "size": "实体数",
  341. }
  342. )
  343. )
  344. channel_totals = fission_match_summary.groupby("实体类型")[
  345. "实体数"
  346. ].transform("sum")
  347. fission_match_summary["渠道实体数"] = channel_totals
  348. fission_match_summary["匹配率"] = (
  349. fission_match_summary["实体数"] / channel_totals
  350. )
  351. fission_match_summary["参数版本"] = fission_parameters.release.version
  352. recommendation_rows = annotated[annotated["动作"].ne("")].copy()
  353. report_rows = annotated[
  354. annotated["阈值样本状态"].isin(
  355. [
  356. "进入三日统一阈值样本池",
  357. "广告级三日合格_不进入阈值样本池",
  358. "单日补充决策_昨日UV>200",
  359. "补充观察_昨日UV>200",
  360. ]
  361. )
  362. ].copy()
  363. for row in snapshots:
  364. row["run_id"] = run_id
  365. for row in actions:
  366. row["run_id"] = run_id
  367. threshold_record = {
  368. "统计窗口": f"{start_date} 至 {end_date}",
  369. "整体三日关停线": thresholds.to_dict("records"),
  370. }
  371. replace_run_results(
  372. run_id,
  373. snapshots=snapshots,
  374. actions=actions,
  375. thresholds=threshold_record,
  376. )
  377. output_dir.mkdir(parents=True, exist_ok=True)
  378. batch_name = f"ROI调控_{start_date}-{end_date}"
  379. if source_revision:
  380. batch_name = f"{batch_name}_{source_revision}"
  381. output_path = output_dir / f"{batch_name}.xlsx"
  382. if not internal_test:
  383. write_workbook(
  384. report_rows,
  385. thresholds,
  386. expected_dates,
  387. output_path,
  388. run_config,
  389. fission_match_summary,
  390. )
  391. agency_report_date = effective_now.strftime("%Y%m%d")
  392. agency_reports = write_agency_workbooks(
  393. report_rows,
  394. output_dir / f"{agency_report_date}_调控建议",
  395. agency_report_date,
  396. agency_names={"自动化投放"} if internal_test else None,
  397. )
  398. internal_reports = (
  399. _build_internal_reports(agency_reports=agency_reports)
  400. if internal_test
  401. else []
  402. )
  403. result_report = (
  404. str(internal_reports[0]["report"])
  405. if internal_test
  406. else str(output_path)
  407. )
  408. result: dict[str, Any] = {
  409. "run_id": run_id,
  410. "batch_name": batch_name,
  411. "status": "COMPUTED",
  412. "reused": False,
  413. "start_date": start_date,
  414. "end_date": end_date,
  415. "source_revision": source_revision,
  416. "entity_count": len(snapshots),
  417. "candidate_count": len(recommendation_rows),
  418. "actionable_count": len(actions),
  419. "thresholds": threshold_record,
  420. "fission_multiplier_matches": fission_match_summary.to_dict(
  421. "records"
  422. ),
  423. "report": result_report,
  424. "agency_reports": agency_reports,
  425. }
  426. if not publish_requested:
  427. return result
  428. if internal_test:
  429. internal_url = (agency_webhook_config.webhooks or {})["内部"]
  430. internal_config = _internal_webhook_config(
  431. internal_reports,
  432. internal_url,
  433. )
  434. publisher = RoiFeishuPublisher(require_chat_ids=False)
  435. try:
  436. deliveries = publish_agency_reports(
  437. run_id=run_id,
  438. reports=internal_reports,
  439. config=internal_config,
  440. publisher=publisher,
  441. now=effective_now,
  442. )
  443. finally:
  444. publisher.close()
  445. result["internal_deliveries"] = deliveries
  446. failed = [row for row in deliveries if row.get("status") != "SENT"]
  447. if failed:
  448. raise RuntimeError(
  449. f"Internal ROI test delivery failed for {len(failed)} report(s)"
  450. )
  451. main_delivery = next(
  452. row for row in deliveries if row["agency_name"] == "自动化投放"
  453. )
  454. mark_published(
  455. run_id,
  456. sheet_token=str(main_delivery["sheet_token"]),
  457. sheet_url=str(main_delivery["sheet_url"]),
  458. message_id="internal-webhook",
  459. expires_at=None,
  460. requires_approval=False,
  461. )
  462. result.update(
  463. {
  464. "status": "COMPLETED",
  465. "sheet_url": main_delivery["sheet_url"],
  466. }
  467. )
  468. return result
  469. counts = recommendation_rows.groupby("动作").size().to_dict()
  470. summary_text = (
  471. f"统计窗口:{start_date} - {end_date}\n"
  472. f"关停建议:{counts.get('关停', 0)}\n"
  473. f"观察:{counts.get('观察', 0)}\n"
  474. f"可执行动作:{len(actions)}\n"
  475. f"审批有效期:发送后 {config.approval_ttl_minutes} 分钟"
  476. )
  477. publisher = RoiFeishuPublisher()
  478. try:
  479. published = publisher.publish(
  480. output_path,
  481. run_id=run_id,
  482. batch_name=batch_name,
  483. summary=summary_text,
  484. requires_approval=bool(actions),
  485. )
  486. expires_at = (
  487. effective_now
  488. + timedelta(minutes=config.approval_ttl_minutes)
  489. if actions
  490. else None
  491. )
  492. mark_published(
  493. run_id,
  494. sheet_token=published["sheet_token"],
  495. sheet_url=published["url"],
  496. message_id=published["message_id"],
  497. expires_at=expires_at,
  498. requires_approval=bool(actions),
  499. )
  500. try:
  501. result["agency_deliveries"] = publish_agency_reports(
  502. run_id=run_id,
  503. reports=agency_reports,
  504. config=agency_webhook_config,
  505. publisher=publisher,
  506. now=now,
  507. )
  508. except Exception as exc:
  509. safe_error = _redact_agency_webhooks(
  510. exc,
  511. agency_webhook_config,
  512. )
  513. logger.error("Agency ROI delivery phase failed: %s", safe_error)
  514. result["agency_delivery_error"] = safe_error
  515. finally:
  516. publisher.close()
  517. result.update(
  518. {
  519. "status": "PENDING_APPROVAL" if actions else "COMPLETED",
  520. "sheet_url": published["url"],
  521. "expires_at": expires_at.isoformat() if expires_at else None,
  522. }
  523. )
  524. return result
  525. except Exception as exc:
  526. mark_failed(run_id, str(exc))
  527. raise