telegram_reactions.py 17 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393
  1. """Answer the post-print outcome prompt with a Telegram reaction (#3046).
  2. Reactions arrive as ``message_reaction`` updates, which the Bot API only
  3. hands out through ``getUpdates`` (long polling) or a webhook. Bambuddy
  4. polls: no inbound connectivity, no public URL, no certificate — the same
  5. outbound-only footing every other notification provider works from. One
  6. task per distinct bot token: several per-printer providers usually share a
  7. bot, and Telegram allows a single getUpdates consumer per bot (a second one
  8. gets a 409). The notifications routes resync the set whenever a provider
  9. changes.
  10. Reactions set by bots are never delivered by Telegram, and in groups the
  11. bot has to be an administrator to receive them at all — both documented Bot
  12. API behaviour, nothing to work around here.
  13. """
  14. import asyncio
  15. import json
  16. import logging
  17. import time
  18. from datetime import timedelta
  19. import httpx
  20. from sqlalchemy import delete, select
  21. from backend.app.models.archive import PrintArchive
  22. from backend.app.models.notification import NotificationProvider, TelegramPendingVerdict
  23. from backend.app.services.notification_service import _USER_AGENT, telegram_markdown_escape
  24. from backend.app.services.print_confirmation import apply_outcome_verdict
  25. from backend.app.utils.local_time import utcnow_naive
  26. logger = logging.getLogger(__name__)
  27. # Telegram holds a getUpdates call open for up to this long before answering
  28. # with an empty list. The HTTP read timeout below has to outlast it.
  29. GET_UPDATES_TIMEOUT = 50
  30. # A 409 means a webhook is set for the bot or a second poller is running; a
  31. # 401/404 means the bot token is wrong or the bot was deleted. None of them
  32. # clears by retrying, so the loop backs off well beyond the poll interval
  33. # instead of hammering the API every few seconds.
  34. CONFLICT_COOLDOWN = 300
  35. MAX_BACKOFF = 60
  36. PENDING_TTL = timedelta(days=7)
  37. PRUNE_INTERVAL = 3600
  38. REACTION_VERDICTS = {"\U0001f44d": "good", "\U0001f44e": "reject"}
  39. VERDICT_SUFFIX = {"good": "✅ marked as good", "reject": "❌ marked as reject"}
  40. def verdict_from_reaction(new_reaction: list) -> str | None:
  41. """Map the reaction list of a message_reaction update to a verdict.
  42. Only plain emoji reactions count; custom emoji and paid reactions are
  43. ignored. The first thumbs in the list wins when several are set.
  44. """
  45. for reaction in new_reaction or []:
  46. if not isinstance(reaction, dict) or reaction.get("type") != "emoji":
  47. continue
  48. verdict = REACTION_VERDICTS.get(reaction.get("emoji", ""))
  49. if verdict:
  50. return verdict
  51. return None
  52. def _provider_reaction_token(provider: NotificationProvider) -> str | None:
  53. """Bot token of a provider that should be polled, else None."""
  54. if provider.provider_type != "telegram" or not provider.enabled:
  55. return None
  56. if (provider.telegram_verdict_mode or "buttons") == "buttons":
  57. return None
  58. config = json.loads(provider.config) if isinstance(provider.config, str) else (provider.config or {})
  59. token = str(config.get("bot_token") or "").strip()
  60. return token or None
  61. def _bot_label(bot_token: str) -> str:
  62. """Log-safe name for a bot: the numeric id in front of the secret part."""
  63. return f"bot {bot_token.split(':', 1)[0]}"
  64. def prompt_is_live(pending: TelegramPendingVerdict, archive: PrintArchive) -> bool:
  65. """Whether a reaction on this prompt may still record a verdict.
  66. Exactly as long as the prompt's one-tap links would: the token the
  67. message went out with is still the archive's token and nothing has spent
  68. it. A reprint clears the token and the next completion mints a new one,
  69. so a reaction on an earlier run's message never answers a later run.
  70. """
  71. return (
  72. bool(pending.confirm_token)
  73. and pending.confirm_token == archive.confirm_token
  74. and archive.confirm_token_used_at is None
  75. )
  76. class TelegramReactionPoller:
  77. def __init__(self):
  78. # Everything is keyed by bot token, not provider: update_ids are a
  79. # per-bot sequence and Telegram serves one getUpdates consumer per
  80. # bot, so providers sharing a bot share one poll.
  81. self._tasks: dict[str, asyncio.Task] = {}
  82. # bot token -> ids of the providers whose prompts this poll answers.
  83. self._providers: dict[str, set[int]] = {}
  84. # bot token -> last update_id confirmed. Kept across task restarts
  85. # within the process so a resync does not re-fetch handled updates.
  86. self._offsets: dict[str, int] = {}
  87. self._http_client: httpx.AsyncClient | None = None
  88. self._session_factory = None
  89. # ------------------------------------------------------------------
  90. # Lifecycle
  91. # ------------------------------------------------------------------
  92. async def start(self):
  93. await self.sync()
  94. def stop(self) -> list[asyncio.Task]:
  95. """Cancel every poll. Returns the cancelled tasks so aclose() can wait for them."""
  96. tasks = list(self._tasks.values())
  97. for task in tasks:
  98. task.cancel()
  99. if tasks:
  100. logger.info("Telegram reaction poller stopped (%d bot(s))", len(tasks))
  101. self._tasks.clear()
  102. self._providers.clear()
  103. return tasks
  104. async def aclose(self):
  105. """Shutdown: stop the polls, let them unwind, and release the HTTP client."""
  106. tasks = self.stop()
  107. if tasks:
  108. await asyncio.gather(*tasks, return_exceptions=True)
  109. if self._http_client is not None and not self._http_client.is_closed:
  110. await self._http_client.aclose()
  111. self._http_client = None
  112. async def sync(self):
  113. """Reconcile the running tasks with the providers in the database.
  114. Called at startup and after every provider create/update/delete.
  115. Starts one task per bot token that at least one enabled provider in
  116. reactions/both mode uses, re-points a running task at the providers
  117. that now share its bot, and cancels the task of a bot no provider
  118. wants any more (deleted, disabled, token changed, or switched back
  119. to buttons). A running task never restarts on a provider edit alone.
  120. """
  121. async with self._sessions()() as db:
  122. result = await db.execute(
  123. select(NotificationProvider).where(NotificationProvider.provider_type == "telegram")
  124. )
  125. providers = list(result.scalars().all())
  126. wanted: dict[str, set[int]] = {}
  127. for provider in providers:
  128. token = _provider_reaction_token(provider)
  129. if token:
  130. wanted.setdefault(token, set()).add(provider.id)
  131. for token in list(self._tasks):
  132. if token not in wanted:
  133. self._tasks.pop(token).cancel()
  134. self._providers.pop(token, None)
  135. logger.info("Telegram reaction poll stopped for %s", _bot_label(token))
  136. for token, provider_ids in wanted.items():
  137. self._providers[token] = provider_ids
  138. if token not in self._tasks:
  139. self._tasks[token] = asyncio.create_task(
  140. self._run(token), name=f"telegram-reactions-{token.split(':', 1)[0]}"
  141. )
  142. logger.info(
  143. "Telegram reaction poll started for %s (provider(s) %s)",
  144. _bot_label(token),
  145. ", ".join(str(i) for i in sorted(provider_ids)),
  146. )
  147. def is_polling(self, provider_id: int) -> bool:
  148. return any(provider_id in ids for token, ids in self._providers.items() if token in self._tasks)
  149. # ------------------------------------------------------------------
  150. # Plumbing
  151. # ------------------------------------------------------------------
  152. def _sessions(self):
  153. if self._session_factory is not None:
  154. return self._session_factory
  155. from backend.app.core.database import async_session
  156. return async_session
  157. async def _get_client(self) -> httpx.AsyncClient:
  158. if self._http_client is None or self._http_client.is_closed:
  159. self._http_client = httpx.AsyncClient(
  160. timeout=httpx.Timeout(GET_UPDATES_TIMEOUT + 10.0, connect=5.0),
  161. headers={"User-Agent": _USER_AGENT},
  162. )
  163. return self._http_client
  164. async def _sleep(self, seconds: float):
  165. await asyncio.sleep(seconds)
  166. async def _run(self, bot_token: str):
  167. backoff = 1.0
  168. last_prune: float | None = None
  169. while True:
  170. try:
  171. if last_prune is None or time.monotonic() - last_prune > PRUNE_INTERVAL:
  172. await self.prune_stale()
  173. last_prune = time.monotonic()
  174. status = await self.poll_once(bot_token)
  175. except asyncio.CancelledError:
  176. raise
  177. except Exception as e:
  178. logger.warning("Telegram reaction poll failed for %s: %s", _bot_label(bot_token), e)
  179. await self._sleep(backoff)
  180. backoff = min(backoff * 2, MAX_BACKOFF)
  181. continue
  182. backoff = 1.0
  183. if status != "ok":
  184. await self._sleep(CONFLICT_COOLDOWN)
  185. # ------------------------------------------------------------------
  186. # One poll
  187. # ------------------------------------------------------------------
  188. async def poll_once(self, bot_token: str) -> str:
  189. """One getUpdates round trip.
  190. Returns "ok", or "conflict" / "rejected" for the permanent failures
  191. (409 webhook or second poller; 401/404 bad or deleted bot token) that
  192. are surfaced on the providers and cooled down instead of retried.
  193. """
  194. params: dict = {"timeout": GET_UPDATES_TIMEOUT, "allowed_updates": ["message_reaction"]}
  195. offset = self._offsets.get(bot_token)
  196. if offset is not None:
  197. params["offset"] = offset + 1
  198. client = await self._get_client()
  199. response = await client.post(f"https://api.telegram.org/bot{bot_token}/getUpdates", json=params)
  200. if response.status_code in (401, 404, 409):
  201. description = ""
  202. try:
  203. description = response.json().get("description") or ""
  204. except Exception:
  205. pass
  206. if response.status_code == 409:
  207. status = "conflict"
  208. message = (
  209. f"Telegram getUpdates conflict: {description or 'conflict'}. Reactions cannot be received "
  210. f"while a webhook is set for this bot or another poller is running"
  211. )
  212. else:
  213. status = "rejected"
  214. detail = f": {description}" if description else ""
  215. message = (
  216. f"Telegram rejected the bot token (HTTP {response.status_code}{detail}). "
  217. f"Check the provider's bot token"
  218. )
  219. await self._record_error(bot_token, f"{message}; retrying in {CONFLICT_COOLDOWN // 60} minutes.")
  220. return status
  221. if response.status_code != 200:
  222. raise RuntimeError(f"getUpdates HTTP {response.status_code}")
  223. payload = response.json()
  224. if not payload.get("ok"):
  225. raise RuntimeError(f"getUpdates failed: {payload.get('description', 'unknown error')}")
  226. for update in payload.get("result") or []:
  227. update_id = update.get("update_id")
  228. if isinstance(update_id, int):
  229. self._offsets[bot_token] = max(self._offsets.get(bot_token, update_id), update_id)
  230. reaction = update.get("message_reaction")
  231. if not isinstance(reaction, dict):
  232. continue
  233. try:
  234. await self._handle_reaction(bot_token, reaction)
  235. except Exception as e:
  236. logger.warning("Telegram reaction for %s could not be applied: %s", _bot_label(bot_token), e)
  237. return "ok"
  238. async def _record_error(self, bot_token: str, message: str):
  239. """Surface a permanent poll failure on every provider that uses the bot."""
  240. provider_ids = sorted(self._providers.get(bot_token, ()))
  241. logger.error("%s (provider(s) %s): %s", _bot_label(bot_token), provider_ids or "-", message)
  242. if not provider_ids:
  243. return
  244. try:
  245. async with self._sessions()() as db:
  246. result = await db.execute(select(NotificationProvider).where(NotificationProvider.id.in_(provider_ids)))
  247. for provider in result.scalars().all():
  248. provider.last_error = message
  249. provider.last_error_at = utcnow_naive()
  250. await db.commit()
  251. except Exception as e:
  252. logger.warning("Could not record the Telegram poll error on providers %s: %s", provider_ids, e)
  253. async def _handle_reaction(self, bot_token: str, reaction: dict):
  254. chat_id = str((reaction.get("chat") or {}).get("id", "")).strip()
  255. message_id = reaction.get("message_id")
  256. verdict = verdict_from_reaction(reaction.get("new_reaction") or [])
  257. provider_ids = self._providers.get(bot_token)
  258. if not chat_id or not isinstance(message_id, int) or verdict is None or not provider_ids:
  259. return
  260. async with self._sessions()() as db:
  261. # message_ids are unique per chat, but two bots talking to the
  262. # same person each see their own numbering, so the match stays
  263. # scoped to the providers of this bot.
  264. pending = await db.scalar(
  265. select(TelegramPendingVerdict).where(
  266. TelegramPendingVerdict.provider_id.in_(sorted(provider_ids)),
  267. TelegramPendingVerdict.chat_id == chat_id,
  268. TelegramPendingVerdict.message_id == message_id,
  269. )
  270. )
  271. if pending is None:
  272. return
  273. archive = await db.get(PrintArchive, pending.archive_id)
  274. if archive is None or not prompt_is_live(pending, archive):
  275. ignored = "the prompt belongs to an earlier run or was answered elsewhere"
  276. elif not await apply_outcome_verdict(db, archive, verdict, source="reaction"):
  277. ignored = "verdict already recorded"
  278. else:
  279. ignored = None
  280. has_caption = bool(pending.has_caption)
  281. text = pending.message_text
  282. archive_id = pending.archive_id
  283. # Whatever happened to the archive, this message is answered.
  284. await db.delete(pending)
  285. await db.commit()
  286. if ignored:
  287. logger.info("Telegram reaction on archive %s ignored: %s", archive_id, ignored)
  288. return
  289. logger.info("[#3046] Telegram reaction marked archive %s as '%s'", archive_id, verdict)
  290. await self._confirm_on_message(bot_token, chat_id, message_id, has_caption, text, verdict)
  291. async def _confirm_on_message(
  292. self, bot_token: str, chat_id: str, message_id: int, has_caption: bool, text: str | None, verdict: str
  293. ):
  294. """Append the verdict to the prompt so the chat shows it was taken.
  295. Editing also drops the inline keyboard in "both" mode — the links
  296. are dead once a verdict landed. The link preview stays off, as on the
  297. original send. Best effort: a failed edit is logged, the verdict
  298. itself is already committed.
  299. """
  300. suffix = VERDICT_SUFFIX[verdict]
  301. client = await self._get_client()
  302. try:
  303. if text and has_caption:
  304. # A photo's caption gets no link preview of its own.
  305. method = "editMessageCaption"
  306. payload = {
  307. "chat_id": chat_id,
  308. "message_id": message_id,
  309. "caption": telegram_markdown_escape(f"{text}\n\n{suffix}"),
  310. "parse_mode": "Markdown",
  311. }
  312. elif text:
  313. method = "editMessageText"
  314. payload = {
  315. "chat_id": chat_id,
  316. "message_id": message_id,
  317. "text": telegram_markdown_escape(f"{text}\n\n{suffix}"),
  318. "parse_mode": "Markdown",
  319. "disable_web_page_preview": True,
  320. }
  321. else:
  322. method = "sendMessage"
  323. payload = {
  324. "chat_id": chat_id,
  325. "text": suffix,
  326. "reply_parameters": {"message_id": message_id},
  327. "disable_web_page_preview": True,
  328. }
  329. response = await client.post(f"https://api.telegram.org/bot{bot_token}/{method}", json=payload)
  330. if response.status_code != 200 or not response.json().get("ok"):
  331. logger.warning("Telegram %s failed after reaction verdict: HTTP %s", method, response.status_code)
  332. except Exception as e:
  333. logger.warning("Telegram confirmation edit failed: %s", e)
  334. async def prune_stale(self):
  335. """Forget prompts nobody reacted to within a week."""
  336. async with self._sessions()() as db:
  337. await db.execute(
  338. delete(TelegramPendingVerdict).where(TelegramPendingVerdict.created_at < utcnow_naive() - PENDING_TTL)
  339. )
  340. await db.commit()
  341. telegram_reaction_poller = TelegramReactionPoller()