feishu_command_service.py 9.1 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240
  1. """Feishu WebSocket entry for deterministic advertising control commands."""
  2. from __future__ import annotations
  3. import logging
  4. import os
  5. from concurrent.futures import ThreadPoolExecutor
  6. from datetime import datetime
  7. from typing import Any
  8. from zoneinfo import ZoneInfo
  9. from agent.tools.builtin.feishu.feishu_client import (
  10. ChatType,
  11. FeishuClient,
  12. FeishuMessageEvent,
  13. )
  14. from operator_commands import (
  15. ACTION_CANCEL,
  16. ACTION_CONFIRM,
  17. ACTION_DAY_PAUSE,
  18. ACTION_RESUME,
  19. ACTION_REJECT,
  20. ACTION_STATUS,
  21. ACTION_STOP,
  22. parse_command,
  23. )
  24. from operator_control import (
  25. cancel_command,
  26. execute_confirmed_command,
  27. pause_status_summary,
  28. preview_write_command,
  29. )
  30. from realtime_config import RealtimeControlConfig
  31. from roi_control.config import RoiConfig
  32. from roi_control.execution import execute_roi_batch, reject_roi_batch
  33. SHANGHAI = ZoneInfo("Asia/Shanghai")
  34. logger = logging.getLogger("tencent_realtime_control.feishu_commands")
  35. _ACTION_LABELS = {
  36. ACTION_DAY_PAUSE: "仅暂停今天",
  37. ACTION_STOP: "持续停止",
  38. ACTION_RESUME: "恢复投放",
  39. }
  40. def _split_ids(raw: str) -> set[str]:
  41. normalized = str(raw or "").replace(",", ",").replace(" ", ",")
  42. return {value.strip() for value in normalized.split(",") if value.strip()}
  43. class FeishuCommandService:
  44. def __init__(self, *, apply: bool) -> None:
  45. app_id = os.getenv("FEISHU_APP_ID", "").strip()
  46. app_secret = os.getenv("FEISHU_APP_SECRET", "").strip()
  47. if not app_id or not app_secret:
  48. raise RuntimeError("FEISHU_APP_ID and FEISHU_APP_SECRET are required")
  49. self.allowed_chat_id = (
  50. os.getenv("RTC_COMMAND_CHAT_ID", "").strip()
  51. or os.getenv("FEISHU_AD_PROJECT_CHAT_ID", "").strip()
  52. )
  53. self.allowed_open_ids = _split_ids(
  54. os.getenv("RTC_COMMAND_ALLOWED_OPEN_IDS", "")
  55. or os.getenv("FEISHU_OPERATOR_OPEN_ID", "")
  56. )
  57. if not self.allowed_chat_id:
  58. raise RuntimeError(
  59. "RTC_COMMAND_CHAT_ID/FEISHU_AD_PROJECT_CHAT_ID is required"
  60. )
  61. if not self.allowed_open_ids:
  62. raise RuntimeError(
  63. "RTC_COMMAND_ALLOWED_OPEN_IDS/FEISHU_OPERATOR_OPEN_ID is required"
  64. )
  65. self.apply = apply
  66. self.config = RealtimeControlConfig.from_env()
  67. self.roi_config = RoiConfig.from_env()
  68. self.confirmation_ttl_minutes = int(
  69. os.getenv("RTC_COMMAND_CONFIRM_TTL_MINUTES", "10")
  70. )
  71. if self.confirmation_ttl_minutes < 1:
  72. raise ValueError("RTC_COMMAND_CONFIRM_TTL_MINUTES must be positive")
  73. self.client = FeishuClient(app_id=app_id, app_secret=app_secret)
  74. self.executor = ThreadPoolExecutor(
  75. max_workers=int(os.getenv("RTC_COMMAND_WORKERS", "2")),
  76. thread_name_prefix="feishu-ad-command",
  77. )
  78. def start(self) -> Any:
  79. return self.client.start_websocket(
  80. on_message=self._enqueue_message,
  81. blocking=False,
  82. )
  83. def _enqueue_message(self, event: FeishuMessageEvent) -> None:
  84. self.executor.submit(self.handle_message, event)
  85. def _reply(self, event: FeishuMessageEvent, text: str) -> None:
  86. self.client.send_message(
  87. to=event.chat_id,
  88. text=text,
  89. reply_to_message_id=event.message_id,
  90. )
  91. def _authorized(self, event: FeishuMessageEvent) -> bool:
  92. if event.sender_open_id not in self.allowed_open_ids:
  93. return False
  94. if event.chat_type == ChatType.GROUP:
  95. return (
  96. event.chat_id == self.allowed_chat_id
  97. and event.mentioned_bot
  98. )
  99. return event.chat_type == ChatType.P2P
  100. def handle_message(self, event: FeishuMessageEvent) -> None:
  101. if event.content_type not in {"text", "post"}:
  102. return
  103. if not self._authorized(event):
  104. return
  105. try:
  106. parsed = parse_command(event.content)
  107. if parsed is None:
  108. return
  109. now = datetime.now(SHANGHAI)
  110. if parsed.action == ACTION_STATUS:
  111. summary = pause_status_summary()
  112. accounts = ",".join(str(value) for value in summary["accounts"]) or "无"
  113. self._reply(
  114. event,
  115. "当前运营暂停状态\n"
  116. f"- 暂停广告:{summary['total']} 条\n"
  117. f"- 仅暂停今天:{summary['until_next_delivery']} 条\n"
  118. f"- 持续停止:{summary['until_manual']} 条\n"
  119. f"- 涉及账户:{accounts}",
  120. )
  121. return
  122. if parsed.action == ACTION_CANCEL:
  123. if (parsed.command_id or "").startswith("roi_"):
  124. raise ValueError("ROI批次请使用“拒绝 roi_<run_id>”")
  125. command = cancel_command(
  126. parsed.command_id or "",
  127. event.sender_open_id,
  128. now,
  129. )
  130. message = (
  131. f"命令 {command['command_id']} 已取消"
  132. if command["status"] == "CANCELLED"
  133. else (
  134. f"命令 {command['command_id']} 未取消,"
  135. f"当前状态:{command['status']}"
  136. )
  137. )
  138. self._reply(
  139. event,
  140. message,
  141. )
  142. return
  143. if parsed.action == ACTION_REJECT:
  144. run = reject_roi_batch(
  145. parsed.command_id or "",
  146. sender_open_id=event.sender_open_id,
  147. now=now,
  148. )
  149. self._reply(
  150. event,
  151. f"ROI批次 {run['run_id']} 已处理\n- 状态:{run['status']}",
  152. )
  153. return
  154. if parsed.action == ACTION_CONFIRM:
  155. if (parsed.command_id or "").startswith("roi_"):
  156. if not self.roi_config.apply_enabled:
  157. self._reply(
  158. event,
  159. "ROI_APPLY_ENABLED=0,当前只生成报告,不执行腾讯写操作。",
  160. )
  161. return
  162. result = execute_roi_batch(
  163. parsed.command_id or "",
  164. sender_open_id=event.sender_open_id,
  165. now=now,
  166. lock_name=self.config.lock_name,
  167. )
  168. self._reply(
  169. event,
  170. f"ROI批次 {result['run_id']} 执行完成\n"
  171. f"- 状态:{result['status']}\n"
  172. f"- 成功:{result.get('successes', 0)} 条\n"
  173. f"- 跳过:{result.get('skipped', 0)} 条\n"
  174. f"- 失败:{result.get('failures', 0)} 条",
  175. )
  176. return
  177. if not self.apply:
  178. self._reply(event, "当前服务为 dry-run,禁止执行腾讯写操作。")
  179. return
  180. command = execute_confirmed_command(
  181. parsed.command_id or "",
  182. sender_open_id=event.sender_open_id,
  183. now=now,
  184. start_hour=self.config.start_hour,
  185. lock_name=self.config.lock_name,
  186. )
  187. self._reply(
  188. event,
  189. f"命令 {command['command_id']} 执行完成\n"
  190. f"- 状态:{command['status']}\n"
  191. f"- 成功:{command.get('successes', 0)} 条\n"
  192. f"- 失败:{command.get('failures', 0)} 条",
  193. )
  194. return
  195. command = preview_write_command(
  196. parsed,
  197. now=now,
  198. source_message_id=event.message_id,
  199. chat_id=event.chat_id,
  200. sender_open_id=event.sender_open_id,
  201. sender_name=event.sender_name,
  202. confirmation_ttl_minutes=self.confirmation_ttl_minutes,
  203. start_hour=self.config.start_hour,
  204. )
  205. resume_text = (
  206. command["resume_at"].strftime("%Y-%m-%d %H:%M")
  207. if command.get("resume_at")
  208. else "仅人工明确恢复"
  209. )
  210. self._reply(
  211. event,
  212. f"操作预览:{_ACTION_LABELS[command['action']]}\n"
  213. f"- 账户:{command['preview_account_count']} 个\n"
  214. f"- 预计涉及广告:{command['preview_ad_count']} 条\n"
  215. f"- 恢复时间:{resume_text}\n"
  216. f"- 命令ID:{command['command_id']}\n"
  217. f"请在 {self.confirmation_ttl_minutes} 分钟内回复:"
  218. f"确认 {command['command_id']}\n"
  219. f"取消命令请回复:取消 {command['command_id']}",
  220. )
  221. except Exception as exc:
  222. logger.exception("Feishu command failed")
  223. self._reply(event, f"命令处理失败:{exc}")