agency_delivery.py 8.6 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240
  1. """Upload per-agency ROI workbooks and notify exact Feishu bot routes."""
  2. from __future__ import annotations
  3. import hashlib
  4. import logging
  5. from datetime import datetime
  6. from pathlib import Path
  7. from typing import Any
  8. from zoneinfo import ZoneInfo
  9. import httpx
  10. from .config import AgencyWebhookConfig
  11. from .feishu import RoiFeishuPublisher
  12. from .repository import upsert_agency_delivery, update_agency_delivery
  13. SHANGHAI = ZoneInfo("Asia/Shanghai")
  14. logger = logging.getLogger("auto_put_ad_mini.roi_control.agency_delivery")
  15. def agency_route_key(agency_name: str) -> str:
  16. """Resolve the configured short agency key without fuzzy matching."""
  17. normalized = "".join(str(agency_name).split())
  18. for prefix in ("小程序-代投-", "小程序-代投"):
  19. if normalized.startswith(prefix):
  20. return normalized[len(prefix) :]
  21. return normalized
  22. def _file_sha256(path: Path) -> str:
  23. digest = hashlib.sha256()
  24. with path.open("rb") as handle:
  25. for chunk in iter(lambda: handle.read(1024 * 1024), b""):
  26. digest.update(chunk)
  27. return digest.hexdigest()
  28. def _route_fingerprint(url: str) -> str:
  29. return hashlib.sha256(url.encode("utf-8")).hexdigest()
  30. class AgencyWebhookNotifier:
  31. def __init__(self, timeout: float = 20.0) -> None:
  32. self.client = httpx.Client(timeout=timeout)
  33. def close(self) -> None:
  34. self.client.close()
  35. def send(
  36. self,
  37. webhook_url: str,
  38. *,
  39. title: str,
  40. sheet_url: str,
  41. creative_rows: int,
  42. ad_rows: int,
  43. ) -> str:
  44. card = {
  45. "config": {"wide_screen_mode": True},
  46. "header": {
  47. "template": "blue",
  48. "title": {"tag": "plain_text", "content": title},
  49. },
  50. "elements": [
  51. {
  52. "tag": "div",
  53. "text": {
  54. "tag": "lark_md",
  55. "content": (
  56. f"创意建议:**{creative_rows}** 条\n"
  57. f"广告建议:**{ad_rows}** 条\n"
  58. "请点击下方按钮查看调控建议。"
  59. ),
  60. },
  61. },
  62. {
  63. "tag": "action",
  64. "actions": [
  65. {
  66. "tag": "button",
  67. "type": "primary",
  68. "text": {"tag": "plain_text", "content": "查看调控建议"},
  69. "url": sheet_url,
  70. }
  71. ],
  72. },
  73. ],
  74. }
  75. response = self.client.post(
  76. webhook_url,
  77. json={"msg_type": "interactive", "card": card},
  78. )
  79. response.raise_for_status()
  80. payload = response.json()
  81. code = payload.get("code", payload.get("StatusCode"))
  82. if code != 0:
  83. message = payload.get("msg", payload.get("StatusMessage", "unknown"))
  84. raise RuntimeError(f"Feishu agency webhook rejected card: {message}")
  85. return str(code)
  86. def publish_agency_reports(
  87. *,
  88. run_id: str,
  89. reports: list[dict[str, object]],
  90. config: AgencyWebhookConfig,
  91. publisher: RoiFeishuPublisher,
  92. notifier: AgencyWebhookNotifier | None = None,
  93. now: datetime | None = None,
  94. ) -> list[dict[str, object]]:
  95. """Publish configured agency reports independently and return safe outcomes."""
  96. if not config.enabled:
  97. return []
  98. routes = config.webhooks or {}
  99. owned_notifier = notifier is None
  100. sender = notifier or AgencyWebhookNotifier()
  101. outcomes: list[dict[str, object]] = []
  102. try:
  103. for report in reports:
  104. agency_name = str(report["agency_name"])
  105. route_key = agency_route_key(agency_name)
  106. webhook_url = routes.get(route_key)
  107. path = Path(str(report["report"]))
  108. delivery_id: int | None = None
  109. try:
  110. if not path.is_file():
  111. raise FileNotFoundError(path)
  112. delivery = upsert_agency_delivery(
  113. {
  114. "run_id": run_id,
  115. "agency_name": agency_name,
  116. "agency_report_version": str(report["report_version"]),
  117. "file_path": str(path),
  118. "file_sha256": _file_sha256(path),
  119. "creative_rows": int(report.get("creative_rows") or 0),
  120. "ad_rows": int(report.get("ad_rows") or 0),
  121. "route_fingerprint": (
  122. _route_fingerprint(webhook_url) if webhook_url else None
  123. ),
  124. }
  125. )
  126. delivery_id = int(delivery["id"])
  127. if delivery.get("status") == "SENT":
  128. outcomes.append(
  129. {
  130. "agency_name": agency_name,
  131. "status": "SENT",
  132. "sheet_url": str(delivery.get("sheet_url") or ""),
  133. "sheet_token": str(delivery.get("sheet_token") or ""),
  134. "reused": True,
  135. }
  136. )
  137. continue
  138. if not webhook_url:
  139. update_agency_delivery(
  140. delivery_id,
  141. status="SKIPPED",
  142. error_message="未配置代理机器人",
  143. )
  144. outcomes.append(
  145. {
  146. "agency_name": agency_name,
  147. "status": "SKIPPED",
  148. "reason": "未配置代理机器人",
  149. }
  150. )
  151. continue
  152. sheet_url = str(delivery.get("sheet_url") or "")
  153. sheet_token = str(delivery.get("sheet_token") or "")
  154. if not sheet_url or not sheet_token:
  155. imported = publisher.upload_workbook(path)
  156. sheet_url = imported["url"]
  157. sheet_token = imported["sheet_token"]
  158. update_agency_delivery(
  159. delivery_id,
  160. status="UPLOADED",
  161. sheet_token=sheet_token,
  162. sheet_url=sheet_url,
  163. )
  164. response_code = sender.send(
  165. webhook_url,
  166. title=str(report.get("title") or path.stem),
  167. sheet_url=sheet_url,
  168. creative_rows=int(report.get("creative_rows") or 0),
  169. ad_rows=int(report.get("ad_rows") or 0),
  170. )
  171. update_agency_delivery(
  172. delivery_id,
  173. status="SENT",
  174. response_code=response_code,
  175. sent_at=now or datetime.now(SHANGHAI),
  176. increment_attempt=True,
  177. )
  178. outcomes.append(
  179. {
  180. "agency_name": agency_name,
  181. "status": "SENT",
  182. "sheet_url": sheet_url,
  183. "sheet_token": sheet_token,
  184. "reused": False,
  185. }
  186. )
  187. except Exception as exc:
  188. safe_error = str(exc)
  189. if webhook_url:
  190. safe_error = safe_error.replace(webhook_url, "<redacted>")
  191. if delivery_id is not None:
  192. try:
  193. update_agency_delivery(
  194. delivery_id,
  195. status="FAILED",
  196. error_message=safe_error,
  197. increment_attempt=True,
  198. )
  199. except Exception as audit_exc:
  200. logger.error(
  201. "Agency ROI delivery audit failed agency=%s: %s",
  202. agency_name,
  203. str(audit_exc),
  204. )
  205. logger.error(
  206. "Agency ROI delivery failed agency=%s: %s",
  207. agency_name,
  208. safe_error,
  209. )
  210. outcomes.append(
  211. {
  212. "agency_name": agency_name,
  213. "status": "FAILED",
  214. "error": safe_error,
  215. }
  216. )
  217. finally:
  218. if owned_notifier:
  219. sender.close()
  220. return outcomes