| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240 |
- """Upload per-agency ROI workbooks and notify exact Feishu bot routes."""
- from __future__ import annotations
- import hashlib
- import logging
- from datetime import datetime
- from pathlib import Path
- from typing import Any
- from zoneinfo import ZoneInfo
- import httpx
- from .config import AgencyWebhookConfig
- from .feishu import RoiFeishuPublisher
- from .repository import upsert_agency_delivery, update_agency_delivery
- SHANGHAI = ZoneInfo("Asia/Shanghai")
- logger = logging.getLogger("auto_put_ad_mini.roi_control.agency_delivery")
- def agency_route_key(agency_name: str) -> str:
- """Resolve the configured short agency key without fuzzy matching."""
- normalized = "".join(str(agency_name).split())
- for prefix in ("小程序-代投-", "小程序-代投"):
- if normalized.startswith(prefix):
- return normalized[len(prefix) :]
- return normalized
- def _file_sha256(path: Path) -> str:
- digest = hashlib.sha256()
- with path.open("rb") as handle:
- for chunk in iter(lambda: handle.read(1024 * 1024), b""):
- digest.update(chunk)
- return digest.hexdigest()
- def _route_fingerprint(url: str) -> str:
- return hashlib.sha256(url.encode("utf-8")).hexdigest()
- class AgencyWebhookNotifier:
- def __init__(self, timeout: float = 20.0) -> None:
- self.client = httpx.Client(timeout=timeout)
- def close(self) -> None:
- self.client.close()
- def send(
- self,
- webhook_url: str,
- *,
- title: str,
- sheet_url: str,
- creative_rows: int,
- ad_rows: int,
- ) -> str:
- card = {
- "config": {"wide_screen_mode": True},
- "header": {
- "template": "blue",
- "title": {"tag": "plain_text", "content": title},
- },
- "elements": [
- {
- "tag": "div",
- "text": {
- "tag": "lark_md",
- "content": (
- f"创意建议:**{creative_rows}** 条\n"
- f"广告建议:**{ad_rows}** 条\n"
- "请点击下方按钮查看调控建议。"
- ),
- },
- },
- {
- "tag": "action",
- "actions": [
- {
- "tag": "button",
- "type": "primary",
- "text": {"tag": "plain_text", "content": "查看调控建议"},
- "url": sheet_url,
- }
- ],
- },
- ],
- }
- response = self.client.post(
- webhook_url,
- json={"msg_type": "interactive", "card": card},
- )
- response.raise_for_status()
- payload = response.json()
- code = payload.get("code", payload.get("StatusCode"))
- if code != 0:
- message = payload.get("msg", payload.get("StatusMessage", "unknown"))
- raise RuntimeError(f"Feishu agency webhook rejected card: {message}")
- return str(code)
- def publish_agency_reports(
- *,
- run_id: str,
- reports: list[dict[str, object]],
- config: AgencyWebhookConfig,
- publisher: RoiFeishuPublisher,
- notifier: AgencyWebhookNotifier | None = None,
- now: datetime | None = None,
- ) -> list[dict[str, object]]:
- """Publish configured agency reports independently and return safe outcomes."""
- if not config.enabled:
- return []
- routes = config.webhooks or {}
- owned_notifier = notifier is None
- sender = notifier or AgencyWebhookNotifier()
- outcomes: list[dict[str, object]] = []
- try:
- for report in reports:
- agency_name = str(report["agency_name"])
- route_key = agency_route_key(agency_name)
- webhook_url = routes.get(route_key)
- path = Path(str(report["report"]))
- delivery_id: int | None = None
- try:
- if not path.is_file():
- raise FileNotFoundError(path)
- delivery = upsert_agency_delivery(
- {
- "run_id": run_id,
- "agency_name": agency_name,
- "agency_report_version": str(report["report_version"]),
- "file_path": str(path),
- "file_sha256": _file_sha256(path),
- "creative_rows": int(report.get("creative_rows") or 0),
- "ad_rows": int(report.get("ad_rows") or 0),
- "route_fingerprint": (
- _route_fingerprint(webhook_url) if webhook_url else None
- ),
- }
- )
- delivery_id = int(delivery["id"])
- if delivery.get("status") == "SENT":
- outcomes.append(
- {
- "agency_name": agency_name,
- "status": "SENT",
- "sheet_url": str(delivery.get("sheet_url") or ""),
- "sheet_token": str(delivery.get("sheet_token") or ""),
- "reused": True,
- }
- )
- continue
- if not webhook_url:
- update_agency_delivery(
- delivery_id,
- status="SKIPPED",
- error_message="未配置代理机器人",
- )
- outcomes.append(
- {
- "agency_name": agency_name,
- "status": "SKIPPED",
- "reason": "未配置代理机器人",
- }
- )
- continue
- sheet_url = str(delivery.get("sheet_url") or "")
- sheet_token = str(delivery.get("sheet_token") or "")
- if not sheet_url or not sheet_token:
- imported = publisher.upload_workbook(path)
- sheet_url = imported["url"]
- sheet_token = imported["sheet_token"]
- update_agency_delivery(
- delivery_id,
- status="UPLOADED",
- sheet_token=sheet_token,
- sheet_url=sheet_url,
- )
- response_code = sender.send(
- webhook_url,
- title=str(report.get("title") or path.stem),
- sheet_url=sheet_url,
- creative_rows=int(report.get("creative_rows") or 0),
- ad_rows=int(report.get("ad_rows") or 0),
- )
- update_agency_delivery(
- delivery_id,
- status="SENT",
- response_code=response_code,
- sent_at=now or datetime.now(SHANGHAI),
- increment_attempt=True,
- )
- outcomes.append(
- {
- "agency_name": agency_name,
- "status": "SENT",
- "sheet_url": sheet_url,
- "sheet_token": sheet_token,
- "reused": False,
- }
- )
- except Exception as exc:
- safe_error = str(exc)
- if webhook_url:
- safe_error = safe_error.replace(webhook_url, "<redacted>")
- if delivery_id is not None:
- try:
- update_agency_delivery(
- delivery_id,
- status="FAILED",
- error_message=safe_error,
- increment_attempt=True,
- )
- except Exception as audit_exc:
- logger.error(
- "Agency ROI delivery audit failed agency=%s: %s",
- agency_name,
- str(audit_exc),
- )
- logger.error(
- "Agency ROI delivery failed agency=%s: %s",
- agency_name,
- safe_error,
- )
- outcomes.append(
- {
- "agency_name": agency_name,
- "status": "FAILED",
- "error": safe_error,
- }
- )
- finally:
- if owned_notifier:
- sender.close()
- return outcomes
|