| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317 |
- """Orchestrate one idempotent daily ROI metric and approval batch."""
- from __future__ import annotations
- import hashlib
- import logging
- import os
- import re
- from datetime import datetime, timedelta
- from pathlib import Path
- from typing import Any
- from zoneinfo import ZoneInfo
- from storage import initialize_schema, load_managed_accounts
- from .config import RoiConfig
- from .data_source import (
- ODPSClient,
- date_window,
- fetch_ad_age,
- fetch_daily_data,
- resolve_end_date,
- )
- from .feishu import RoiFeishuPublisher
- from .metrics import ENTITY_QIWEI, METRIC_RUN_SUFFIX, METRIC_VERSION
- from .fission_multiplier import (
- DEFAULT_FISSION_PARAMETER_VERSION,
- FissionMultiplierParameters,
- QIWEI_REFERENCE_MULTIPLIER,
- QIWEI_REFERENCE_RUN_SUFFIX,
- QIWEI_REFERENCE_VERSION,
- parameters_from_database,
- )
- from .policy import annotate_execution
- from .reporting import REPORT_RUN_SUFFIX, REPORT_VERSION, write_workbook
- from .rules import (
- POLICY_RUN_SUFFIX,
- POLICY_VERSION,
- RuleConfig,
- evaluate_rules,
- )
- from .repository import (
- FINAL_STATUSES,
- create_or_load_run,
- load_fission_parameter_release,
- mark_failed,
- mark_published,
- replace_run_results,
- )
- SHANGHAI = ZoneInfo("Asia/Shanghai")
- logger = logging.getLogger("auto_put_ad_mini.roi_control")
- def _rule_config(config: RoiConfig) -> RuleConfig:
- return RuleConfig(
- self_stop_min_age=config.self_stop_min_age,
- self_up_min_age=config.self_up_min_age,
- self_min_avg_uv=config.self_min_avg_uv,
- partner_min_avg_uv=config.partner_min_avg_uv,
- min_daily_cost=config.min_daily_cost,
- stop_quantile=config.stop_quantile,
- up_quantile=config.up_quantile,
- gzh_adjust_rank_min=config.gzh_adjust_rank_min,
- gzh_adjust_rank_max=config.gzh_adjust_rank_max,
- )
- def _dates(start_date: str) -> list[str]:
- start = datetime.strptime(start_date, "%Y%m%d")
- return [
- (start + timedelta(days=offset)).strftime("%Y%m%d")
- for offset in range(3)
- ]
- def _run_identity(
- end_date: str,
- fission_parameters: FissionMultiplierParameters,
- source_revision: str | None = None,
- ) -> tuple[str, str]:
- if source_revision and not re.fullmatch(
- r"[a-z0-9][a-z0-9_-]{0,63}", source_revision
- ):
- raise ValueError(
- "source_revision must use lowercase letters, numbers, '_' or '-'"
- )
- release = fission_parameters.release
- run_suffixes = [release.run_suffix, QIWEI_REFERENCE_RUN_SUFFIX]
- run_key_versions = [release.version, QIWEI_REFERENCE_VERSION]
- run_id = (
- f"roi_{end_date}_{METRIC_RUN_SUFFIX}_{POLICY_RUN_SUFFIX}_"
- f"{'_'.join(run_suffixes)}_{REPORT_RUN_SUFFIX}"
- )
- run_key = (
- f"{METRIC_VERSION}:{POLICY_VERSION}:{':'.join(run_key_versions)}:"
- f"{REPORT_VERSION}:{end_date}"
- )
- if source_revision:
- run_id = f"{run_id}_{source_revision}"
- revision_hash = hashlib.sha256(source_revision.encode("utf-8")).hexdigest()
- run_key = f"{run_key}:sr:{revision_hash[:16]}"
- if len(run_id) > 64:
- raise ValueError("ROI run_id exceeds database limit")
- if len(run_key) > 128:
- raise ValueError("ROI run_key exceeds database limit")
- return run_id, run_key
- def run_daily_roi(
- *,
- requested_end_date: str | None = None,
- output_dir: Path,
- send_feishu: bool,
- now: datetime | None = None,
- source_revision: str | None = None,
- ) -> dict[str, Any]:
- """Compute, snapshot, report, and optionally publish one ROI batch."""
- config = RoiConfig.from_env()
- initialize_schema()
- fission_version = os.getenv(
- "ROI_FISSION_PARAMETER_VERSION",
- DEFAULT_FISSION_PARAMETER_VERSION,
- )
- fission_parameters = parameters_from_database(
- *load_fission_parameter_release(fission_version)
- )
- run_config = config.snapshot()
- run_config["fission_multiplier"] = fission_parameters.snapshot()
- run_config["qiwei_reference_multiplier"] = QIWEI_REFERENCE_MULTIPLIER
- run_config["qiwei_reference_version"] = QIWEI_REFERENCE_VERSION
- run_config["report_version"] = REPORT_VERSION
- if source_revision:
- run_config["source_revision"] = source_revision
- client = ODPSClient(project=os.getenv("ODPS_PROJECT", "loghubods"))
- end_date = resolve_end_date(client, requested_end_date)
- start_date, end_date = date_window(end_date)
- expected_dates = _dates(start_date)
- run_id, run_key = _run_identity(
- end_date,
- fission_parameters,
- source_revision,
- )
- run = create_or_load_run(
- {
- "run_id": run_id,
- "run_key": run_key,
- "metric_version": METRIC_VERSION,
- "policy_version": POLICY_VERSION,
- "fission_parameter_version": fission_parameters.release.version,
- "fission_cohort_date": datetime.strptime(
- fission_parameters.release.cohort_date, "%Y%m%d"
- ).date(),
- "start_date": datetime.strptime(start_date, "%Y%m%d").date(),
- "end_date": datetime.strptime(end_date, "%Y%m%d").date(),
- "config": run_config,
- }
- )
- reusable_statuses = FINAL_STATUSES | {"PENDING_APPROVAL", "EXECUTING"}
- if run.get("status") in reusable_statuses or (
- run.get("status") == "COMPUTED" and not send_feishu
- ):
- logger.info("Reuse ROI run=%s status=%s", run_id, run.get("status"))
- return {
- "run_id": run_id,
- "status": run.get("status"),
- "reused": True,
- "sheet_url": run.get("sheet_url"),
- }
- try:
- logger.info(
- "ROI %s loading ODPS daily data %s..%s fission_parameter=%s",
- run_id,
- start_date,
- end_date,
- fission_parameters.release.version,
- )
- daily = fetch_daily_data(client, start_date, end_date)
- ad_age = fetch_ad_age(client, end_date)
- _, thresholds, summary = evaluate_rules(
- daily,
- expected_dates,
- ad_age,
- _rule_config(config),
- fission_parameters=fission_parameters,
- )
- managed_ids = {
- int(row["account_id"])
- for row in load_managed_accounts()
- }
- annotated, snapshots, actions = annotate_execution(
- summary,
- managed_ids,
- run_id,
- )
- fission_match_summary = (
- summary.groupby(
- ["entity_type", "传播裂变系数匹配层级"],
- as_index=False,
- dropna=False,
- )
- .size()
- .rename(
- columns={
- "entity_type": "实体类型",
- "传播裂变系数匹配层级": "匹配层级",
- "size": "实体数",
- }
- )
- )
- channel_totals = fission_match_summary.groupby("实体类型")[
- "实体数"
- ].transform("sum")
- fission_match_summary["渠道实体数"] = channel_totals
- fission_match_summary["匹配率"] = (
- fission_match_summary["实体数"] / channel_totals
- )
- fission_match_summary["参数版本"] = fission_parameters.release.version
- actionable_candidates = annotated[annotated["动作"].ne("")].copy()
- report_rows = annotated[
- annotated["阈值样本状态"].eq("进入阈值样本池")
- | annotated["entity_type"].eq(ENTITY_QIWEI)
- ].copy()
- for row in snapshots:
- row["run_id"] = run_id
- for row in actions:
- row["run_id"] = run_id
- threshold_record = thresholds.iloc[0].to_dict()
- replace_run_results(
- run_id,
- snapshots=snapshots,
- actions=actions,
- thresholds=threshold_record,
- )
- output_dir.mkdir(parents=True, exist_ok=True)
- batch_name = f"ROI调控_{start_date}-{end_date}"
- if source_revision:
- batch_name = f"{batch_name}_{source_revision}"
- output_path = output_dir / f"{batch_name}.xlsx"
- write_workbook(
- report_rows,
- thresholds,
- expected_dates,
- output_path,
- run_config,
- fission_match_summary,
- )
- result: dict[str, Any] = {
- "run_id": run_id,
- "batch_name": batch_name,
- "status": "COMPUTED",
- "reused": False,
- "start_date": start_date,
- "end_date": end_date,
- "source_revision": source_revision,
- "entity_count": len(snapshots),
- "candidate_count": len(actionable_candidates),
- "actionable_count": len(actions),
- "thresholds": threshold_record,
- "fission_multiplier_matches": fission_match_summary.to_dict(
- "records"
- ),
- "report": str(output_path),
- }
- if not send_feishu:
- return result
- counts = actionable_candidates.groupby("动作").size().to_dict()
- summary_text = (
- f"统计窗口:{start_date} - {end_date}\n"
- f"关停建议:{counts.get('关停', 0)}\n"
- f"扩量建议:{counts.get('扩量', 0)}\n"
- f"素材调整建议:{counts.get('调整封面&落地页视频', 0)}\n"
- f"可执行动作:{len(actions)}\n"
- f"审批有效期:发送后 {config.approval_ttl_minutes} 分钟"
- )
- publisher = RoiFeishuPublisher()
- try:
- published = publisher.publish(
- output_path,
- run_id=run_id,
- batch_name=batch_name,
- summary=summary_text,
- requires_approval=bool(actions),
- )
- finally:
- publisher.close()
- expires_at = (
- (now or datetime.now(SHANGHAI))
- + timedelta(minutes=config.approval_ttl_minutes)
- if actions
- else None
- )
- mark_published(
- run_id,
- sheet_token=published["sheet_token"],
- sheet_url=published["url"],
- message_id=published["message_id"],
- expires_at=expires_at,
- requires_approval=bool(actions),
- )
- result.update(
- {
- "status": "PENDING_APPROVAL" if actions else "COMPLETED",
- "sheet_url": published["url"],
- "expires_at": expires_at.isoformat() if expires_at else None,
- }
- )
- return result
- except Exception as exc:
- mark_failed(run_id, str(exc))
- raise
|