test_telegram_reactions.py 40 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726727728729730731732733734735736737738739740741742743744745746747748749750751752753754755756757758759760761762763764765766767768769770771772773774775776777778779780781782783784785786787788789790791792793794795796797798799800801802803804805806807808809810811812813814815816817818819820821822823824825826827828829830831832833834835836837838839840841842843844845846847848849850851852853854855856857858859860861862863864865866867868869870871872873874875876877878879880881882883884885886887888889890891892893894895896897898899900901902903904905906907908909910911912913914915916917918919920921922923924925926927928929930931932933934935936937938939940
  1. """Answering the outcome prompt with a Telegram reaction (#3046).
  2. The Bot API is stood in for at the httpx layer (the same ``_http_client``
  3. swap the forum-topic tests use), the database is the real test engine: the
  4. poller opens its own sessions, so the tests hand it the test session
  5. factory instead of the module-level one.
  6. """
  7. from datetime import timedelta
  8. from unittest.mock import AsyncMock, patch
  9. import httpx
  10. import pytest
  11. from httpx import AsyncClient
  12. from sqlalchemy import select, text
  13. from sqlalchemy.ext.asyncio import AsyncSession, async_sessionmaker
  14. from backend.app.models.archive import PrintArchive
  15. from backend.app.models.notification import TelegramPendingVerdict
  16. from backend.app.models.print_log import PrintLogEntry
  17. from backend.app.services.notification_service import TELEGRAM_REACTION_HINT, NotificationService, notification_service
  18. from backend.app.services.print_confirmation import apply_outcome_verdict, one_tap_url
  19. from backend.app.services.telegram_reactions import (
  20. CONFLICT_COOLDOWN,
  21. MAX_BACKOFF,
  22. TelegramReactionPoller,
  23. verdict_from_reaction,
  24. )
  25. from backend.app.utils.local_time import utcnow_naive
  26. BOT_TOKEN = "123456:AAbbCC"
  27. CHAT_ID = "-1002520100736"
  28. TELEGRAM_CONFIG = {"bot_token": BOT_TOKEN, "chat_id": CHAT_ID}
  29. THUMBS_UP = "\U0001f44d"
  30. THUMBS_DOWN = "\U0001f44e"
  31. class _FakeBotApi:
  32. """Stand-in for httpx.AsyncClient with scripted responses, in call order."""
  33. def __init__(self, responses=None):
  34. self.is_closed = False
  35. self.calls: list[dict] = []
  36. self.responses = list(responses or [])
  37. async def post(self, url, json=None, data=None, files=None):
  38. self.calls.append({"url": url, "json": json, "data": data, "files": files})
  39. if self.responses:
  40. nxt = self.responses.pop(0)
  41. if isinstance(nxt, Exception):
  42. raise nxt
  43. return nxt
  44. return httpx.Response(200, json={"ok": True, "result": []})
  45. async def aclose(self):
  46. self.is_closed = True
  47. def urls(self) -> list[str]:
  48. return [c["url"].rsplit("/", 1)[-1] for c in self.calls]
  49. def _updates(*updates: dict) -> httpx.Response:
  50. return httpx.Response(200, json={"ok": True, "result": list(updates)})
  51. def _reaction(update_id: int, message_id: int, emoji: str | None, chat_id: str = CHAT_ID, kind: str = "emoji") -> dict:
  52. new_reaction = [] if emoji is None else [{"type": kind, "emoji": emoji}]
  53. return {
  54. "update_id": update_id,
  55. "message_reaction": {
  56. "chat": {"id": int(chat_id)},
  57. "message_id": message_id,
  58. "user": {"id": 42},
  59. "date": 1_700_000_000,
  60. "old_reaction": [],
  61. "new_reaction": new_reaction,
  62. },
  63. }
  64. @pytest.fixture
  65. async def poller(test_engine):
  66. p = TelegramReactionPoller()
  67. p._session_factory = async_sessionmaker(test_engine, class_=AsyncSession, expire_on_commit=False)
  68. p._http_client = _FakeBotApi()
  69. yield p
  70. await p.aclose()
  71. def _track(poller: TelegramReactionPoller, *providers, token: str = BOT_TOKEN):
  72. """What sync() would record for these providers, without starting a task."""
  73. poller._providers.setdefault(token, set()).update(p.id for p in providers)
  74. @pytest.fixture
  75. def telegram_provider(notification_provider_factory):
  76. async def _create(**kwargs):
  77. defaults = {
  78. "name": "Telegram",
  79. "provider_type": "telegram",
  80. "config": TELEGRAM_CONFIG,
  81. "telegram_verdict_mode": "reactions",
  82. }
  83. defaults.update(kwargs)
  84. return await notification_provider_factory(**defaults)
  85. return _create
  86. @pytest.fixture
  87. def pending_factory(db_session):
  88. async def _create(provider_id: int, archive_id: int, message_id: int = 777, **kwargs):
  89. # A prompt only goes out with a live confirm token (the dispatch mints
  90. # one), and the send records it on the row: do the same here.
  91. if "confirm_token" not in kwargs:
  92. archive = await db_session.get(PrintArchive, archive_id)
  93. if archive.confirm_token is None:
  94. archive.confirm_token = f"tok-{archive_id}-{message_id}"
  95. kwargs["confirm_token"] = archive.confirm_token
  96. defaults = {
  97. "provider_id": provider_id,
  98. "chat_id": CHAT_ID,
  99. "message_id": message_id,
  100. "archive_id": archive_id,
  101. "has_caption": False,
  102. "message_text": "*Outcome?*\nTest_Print on X1C",
  103. }
  104. defaults.update(kwargs)
  105. row = TelegramPendingVerdict(**defaults)
  106. db_session.add(row)
  107. await db_session.commit()
  108. await db_session.refresh(row)
  109. return row
  110. return _create
  111. async def _pending_count(db_session) -> int:
  112. return len((await db_session.execute(select(TelegramPendingVerdict))).scalars().all())
  113. # ---------------------------------------------------------------------------
  114. # Shared verdict helper
  115. # ---------------------------------------------------------------------------
  116. class TestApplyOutcomeVerdict:
  117. @pytest.mark.asyncio
  118. @pytest.mark.integration
  119. async def test_first_verdict_wins_and_mirrors(self, db_session, printer_factory, archive_factory):
  120. printer = await printer_factory()
  121. archive = await archive_factory(printer.id, confirm_requested=True, confirm_token="tok")
  122. assert await apply_outcome_verdict(db_session, archive, "reject", source="reaction") is True
  123. await db_session.commit()
  124. assert archive.user_verdict == "reject"
  125. assert archive.user_verdict_source == "reaction"
  126. assert archive.user_verdict_at is not None
  127. # The capability is spent, not destroyed: the one-tap route can still
  128. # tell whoever taps an old button that this print was answered.
  129. assert archive.confirm_token == "tok"
  130. assert archive.confirm_token_used_at is not None
  131. entry = await db_session.scalar(
  132. select(PrintLogEntry).where(PrintLogEntry.archive_id == archive.id).order_by(PrintLogEntry.id.desc())
  133. )
  134. assert entry.user_verdict == "reject"
  135. # A later verdict from any path is a no-op.
  136. assert await apply_outcome_verdict(db_session, archive, "good", source="reaction") is False
  137. assert archive.user_verdict == "reject"
  138. @pytest.mark.asyncio
  139. @pytest.mark.integration
  140. async def test_rejects_unknown_verdict(self, db_session, printer_factory, archive_factory):
  141. printer = await printer_factory()
  142. archive = await archive_factory(printer.id)
  143. with pytest.raises(ValueError):
  144. await apply_outcome_verdict(db_session, archive, "meh", source="reaction")
  145. with pytest.raises(ValueError):
  146. await apply_outcome_verdict(db_session, archive, "good", source="carrier-pigeon")
  147. # ---------------------------------------------------------------------------
  148. # Sending the prompt
  149. # ---------------------------------------------------------------------------
  150. class _SendCapture:
  151. def __init__(self, message_id=555):
  152. self.is_closed = False
  153. self.calls: list[dict] = []
  154. self.message_id = message_id
  155. async def post(self, url, data=None, files=None, json=None):
  156. self.calls.append({"url": url, "data": data, "files": files, "json": json})
  157. result = {"message_id": self.message_id, "chat": {"id": int(CHAT_ID)}}
  158. return httpx.Response(200, json={"ok": True, "result": result})
  159. CONFIRM_VARIABLES = {"good_url": "http://bambuddy.local/good", "reject_url": "http://bambuddy.local/reject"}
  160. class TestSendingThePrompt:
  161. async def _send(self, provider, db_session, archive_id, image_data=None):
  162. service = NotificationService()
  163. client = _SendCapture()
  164. service._http_client = client
  165. ok, _ = await service._send_to_provider(
  166. provider,
  167. "Outcome?",
  168. "How did it go",
  169. db_session,
  170. image_data=image_data,
  171. event_type="print_confirm_request",
  172. variables={**CONFIRM_VARIABLES, "archive_id": archive_id},
  173. )
  174. assert ok
  175. return client
  176. @pytest.mark.asyncio
  177. @pytest.mark.integration
  178. async def test_reactions_mode_drops_buttons_and_records_message(
  179. self, db_session, printer_factory, archive_factory, telegram_provider
  180. ):
  181. printer = await printer_factory()
  182. archive = await archive_factory(printer.id, confirm_requested=True, confirm_token="live-token")
  183. provider = await telegram_provider(telegram_verdict_mode="reactions")
  184. client = await self._send(provider, db_session, archive.id)
  185. body = client.calls[0]["json"]
  186. assert "reply_markup" not in body
  187. assert TELEGRAM_REACTION_HINT in body["text"]
  188. row = await db_session.scalar(select(TelegramPendingVerdict))
  189. assert row is not None
  190. assert row.provider_id == provider.id
  191. assert row.archive_id == archive.id
  192. assert row.message_id == 555
  193. assert row.chat_id == CHAT_ID
  194. assert row.confirm_token == "live-token", "the token the prompt's links carry"
  195. assert row.has_caption is False
  196. assert row.message_text and TELEGRAM_REACTION_HINT in row.message_text
  197. @pytest.mark.asyncio
  198. @pytest.mark.integration
  199. async def test_both_mode_keeps_buttons_and_records_message(
  200. self, db_session, printer_factory, archive_factory, telegram_provider
  201. ):
  202. printer = await printer_factory()
  203. archive = await archive_factory(printer.id, confirm_requested=True, confirm_token="live-token")
  204. provider = await telegram_provider(telegram_verdict_mode="both")
  205. client = await self._send(provider, db_session, archive.id, image_data=b"\x89PNG")
  206. call = client.calls[0]
  207. assert call["url"].endswith("/sendPhoto")
  208. assert "inline_keyboard" in call["data"]["reply_markup"]
  209. row = await db_session.scalar(select(TelegramPendingVerdict))
  210. assert row is not None and row.has_caption is True
  211. @pytest.mark.asyncio
  212. @pytest.mark.integration
  213. async def test_buttons_mode_is_unchanged(self, db_session, printer_factory, archive_factory, telegram_provider):
  214. printer = await printer_factory()
  215. archive = await archive_factory(printer.id, confirm_requested=True)
  216. provider = await telegram_provider(telegram_verdict_mode="buttons")
  217. client = await self._send(provider, db_session, archive.id)
  218. body = client.calls[0]["json"]
  219. assert body["reply_markup"] == {
  220. "inline_keyboard": [
  221. [
  222. {"text": f"{THUMBS_UP} Good", "url": one_tap_url(CONFIRM_VARIABLES["good_url"])},
  223. {"text": f"{THUMBS_DOWN} Reject", "url": one_tap_url(CONFIRM_VARIABLES["reject_url"])},
  224. ]
  225. ]
  226. }
  227. assert TELEGRAM_REACTION_HINT not in body["text"]
  228. assert await _pending_count(db_session) == 0
  229. @pytest.mark.asyncio
  230. @pytest.mark.integration
  231. async def test_prompt_without_a_live_token_records_nothing(
  232. self, db_session, printer_factory, archive_factory, telegram_provider
  233. ):
  234. """The prompt's links are dead without a live token, and a reaction
  235. must not outlive them: nothing to remember."""
  236. printer = await printer_factory()
  237. archive = await archive_factory(
  238. printer.id, confirm_requested=True, confirm_token="spent", confirm_token_used_at=utcnow_naive()
  239. )
  240. provider = await telegram_provider(telegram_verdict_mode="reactions")
  241. await self._send(provider, db_session, archive.id)
  242. assert await _pending_count(db_session) == 0
  243. @pytest.mark.asyncio
  244. @pytest.mark.integration
  245. async def test_missing_message_id_records_nothing(
  246. self, db_session, printer_factory, archive_factory, telegram_provider
  247. ):
  248. printer = await printer_factory()
  249. archive = await archive_factory(printer.id, confirm_requested=True)
  250. provider = await telegram_provider(telegram_verdict_mode="reactions")
  251. service = NotificationService()
  252. client = _SendCapture(message_id=0)
  253. service._http_client = client
  254. ok, _ = await service._send_to_provider(
  255. provider,
  256. "Outcome?",
  257. "How did it go",
  258. db_session,
  259. event_type="print_confirm_request",
  260. variables={**CONFIRM_VARIABLES, "archive_id": archive.id},
  261. )
  262. assert ok
  263. assert await _pending_count(db_session) == 0
  264. # ---------------------------------------------------------------------------
  265. # Poller
  266. # ---------------------------------------------------------------------------
  267. class TestVerdictFromReaction:
  268. def test_thumbs_map_and_others_ignored(self):
  269. assert verdict_from_reaction([{"type": "emoji", "emoji": THUMBS_UP}]) == "good"
  270. assert verdict_from_reaction([{"type": "emoji", "emoji": THUMBS_DOWN}]) == "reject"
  271. assert verdict_from_reaction([{"type": "emoji", "emoji": "❤"}]) is None
  272. assert verdict_from_reaction([{"type": "custom_emoji", "custom_emoji_id": "1"}]) is None
  273. assert verdict_from_reaction([]) is None
  274. assert verdict_from_reaction([{"type": "emoji", "emoji": "❤"}, {"type": "emoji", "emoji": THUMBS_UP}]) == "good"
  275. class TestPoller:
  276. @pytest.mark.asyncio
  277. @pytest.mark.integration
  278. async def test_thumbs_up_marks_good_deletes_row_and_edits_message(
  279. self, poller, db_session, printer_factory, archive_factory, telegram_provider, pending_factory
  280. ):
  281. printer = await printer_factory()
  282. archive = await archive_factory(printer.id, confirm_requested=True, confirm_token="tok")
  283. provider = await telegram_provider()
  284. _track(poller, provider)
  285. await pending_factory(provider.id, archive.id, message_id=777)
  286. poller._http_client = _FakeBotApi([_updates(_reaction(10, 777, THUMBS_UP))])
  287. assert await poller.poll_once(BOT_TOKEN) == "ok"
  288. await db_session.refresh(archive)
  289. assert archive.user_verdict == "good"
  290. assert archive.user_verdict_source == "reaction"
  291. assert archive.confirm_token_used_at is not None
  292. entry = await db_session.scalar(
  293. select(PrintLogEntry).where(PrintLogEntry.archive_id == archive.id).order_by(PrintLogEntry.id.desc())
  294. )
  295. assert entry.user_verdict == "good"
  296. assert await _pending_count(db_session) == 0
  297. assert poller._http_client.urls() == ["getUpdates", "editMessageText"]
  298. edit = poller._http_client.calls[1]["json"]
  299. assert edit["chat_id"] == CHAT_ID and edit["message_id"] == 777
  300. assert edit["text"].endswith("✅ marked as good")
  301. assert "Test\\_Print" in edit["text"], "stored body is re-escaped for Markdown"
  302. assert "reply_markup" not in edit
  303. assert edit["disable_web_page_preview"] is True, "the edit keeps the preview off, like the send"
  304. # The getUpdates that follows confirms the update we handled.
  305. assert poller._offsets[BOT_TOKEN] == 10
  306. @pytest.mark.asyncio
  307. @pytest.mark.integration
  308. async def test_getupdates_request_shape(self, poller, telegram_provider):
  309. provider = await telegram_provider()
  310. _track(poller, provider)
  311. poller._offsets[BOT_TOKEN] = 41
  312. await poller.poll_once(BOT_TOKEN)
  313. call = poller._http_client.calls[0]
  314. assert call["url"] == f"https://api.telegram.org/bot{BOT_TOKEN}/getUpdates"
  315. assert call["json"] == {"timeout": 50, "allowed_updates": ["message_reaction"], "offset": 42}
  316. @pytest.mark.asyncio
  317. @pytest.mark.integration
  318. async def test_offset_belongs_to_the_bot_not_the_provider(self, poller, telegram_provider):
  319. """update_ids are a per-bot sequence: after a provider is re-pointed at
  320. another bot, the first getUpdates for it must not carry the old bot's
  321. offset, or the new bot's lower ids are confirmed away unseen."""
  322. provider = await telegram_provider()
  323. _track(poller, provider)
  324. poller._offsets[BOT_TOKEN] = 900_000_123
  325. _track(poller, provider, token="999:newtoken")
  326. await poller.poll_once("999:newtoken")
  327. call = poller._http_client.calls[0]
  328. assert call["url"].startswith("https://api.telegram.org/bot999:newtoken/")
  329. assert "offset" not in call["json"]
  330. @pytest.mark.asyncio
  331. @pytest.mark.integration
  332. async def test_thumbs_down_rejects(
  333. self, poller, db_session, printer_factory, archive_factory, telegram_provider, pending_factory
  334. ):
  335. printer = await printer_factory()
  336. archive = await archive_factory(printer.id, confirm_requested=True)
  337. provider = await telegram_provider()
  338. _track(poller, provider)
  339. await pending_factory(provider.id, archive.id, message_id=778)
  340. poller._http_client = _FakeBotApi([_updates(_reaction(11, 778, THUMBS_DOWN))])
  341. await poller.poll_once(BOT_TOKEN)
  342. await db_session.refresh(archive)
  343. assert archive.user_verdict == "reject"
  344. assert poller._http_client.calls[1]["json"]["text"].endswith("❌ marked as reject")
  345. @pytest.mark.asyncio
  346. @pytest.mark.integration
  347. async def test_caption_message_is_edited_with_editMessageCaption(
  348. self, poller, db_session, printer_factory, archive_factory, telegram_provider, pending_factory
  349. ):
  350. printer = await printer_factory()
  351. archive = await archive_factory(printer.id, confirm_requested=True)
  352. provider = await telegram_provider()
  353. _track(poller, provider)
  354. await pending_factory(provider.id, archive.id, message_id=779, has_caption=True)
  355. poller._http_client = _FakeBotApi([_updates(_reaction(12, 779, THUMBS_UP))])
  356. await poller.poll_once(BOT_TOKEN)
  357. assert poller._http_client.urls() == ["getUpdates", "editMessageCaption"]
  358. assert poller._http_client.calls[1]["json"]["caption"].endswith("✅ marked as good")
  359. @pytest.mark.asyncio
  360. @pytest.mark.integration
  361. async def test_unknown_message_is_ignored(
  362. self, poller, db_session, printer_factory, archive_factory, telegram_provider, pending_factory
  363. ):
  364. printer = await printer_factory()
  365. archive = await archive_factory(printer.id, confirm_requested=True)
  366. provider = await telegram_provider()
  367. _track(poller, provider)
  368. await pending_factory(provider.id, archive.id, message_id=780)
  369. poller._http_client = _FakeBotApi(
  370. [_updates(_reaction(13, 999, THUMBS_UP), _reaction(14, 780, THUMBS_UP, chat_id="-100999"))]
  371. )
  372. await poller.poll_once(BOT_TOKEN)
  373. await db_session.refresh(archive)
  374. assert archive.user_verdict is None
  375. assert await _pending_count(db_session) == 1
  376. assert poller._http_client.urls() == ["getUpdates"]
  377. assert poller._offsets[BOT_TOKEN] == 14
  378. @pytest.mark.asyncio
  379. @pytest.mark.integration
  380. async def test_other_emoji_and_removed_reaction_are_ignored(
  381. self, poller, db_session, printer_factory, archive_factory, telegram_provider, pending_factory
  382. ):
  383. printer = await printer_factory()
  384. archive = await archive_factory(printer.id, confirm_requested=True)
  385. provider = await telegram_provider()
  386. _track(poller, provider)
  387. await pending_factory(provider.id, archive.id, message_id=781)
  388. poller._http_client = _FakeBotApi([_updates(_reaction(15, 781, "❤"), _reaction(16, 781, None))])
  389. await poller.poll_once(BOT_TOKEN)
  390. await db_session.refresh(archive)
  391. assert archive.user_verdict is None
  392. assert await _pending_count(db_session) == 1, "the prompt stays open for a later thumbs"
  393. @pytest.mark.asyncio
  394. @pytest.mark.integration
  395. async def test_reaction_after_a_verdict_is_a_noop(
  396. self, poller, db_session, printer_factory, archive_factory, telegram_provider, pending_factory
  397. ):
  398. """Someone already answered in the web UI: the late reaction neither
  399. flips the verdict nor edits the message, but the row is retired."""
  400. printer = await printer_factory()
  401. archive = await archive_factory(printer.id, confirm_requested=True, user_verdict="reject")
  402. provider = await telegram_provider()
  403. _track(poller, provider)
  404. await pending_factory(provider.id, archive.id, message_id=782)
  405. poller._http_client = _FakeBotApi([_updates(_reaction(17, 782, THUMBS_UP))])
  406. await poller.poll_once(BOT_TOKEN)
  407. await db_session.refresh(archive)
  408. assert archive.user_verdict == "reject"
  409. assert await _pending_count(db_session) == 0
  410. assert poller._http_client.urls() == ["getUpdates"]
  411. @pytest.mark.asyncio
  412. @pytest.mark.integration
  413. async def test_an_earlier_runs_prompt_cannot_answer_a_reprint(
  414. self, poller, db_session, printer_factory, archive_factory, telegram_provider, monkeypatch
  415. ):
  416. """Printing an archive again resets its verdict and confirm token, so
  417. run 1's links stop working. Run 1's unanswered message has to stop with
  418. them: a thumbs on it while run 2 prints records nothing, and run 2
  419. still asks its own question, which its own message then answers."""
  420. from backend.app.main import dispatch_outcome_confirmation
  421. printer = await printer_factory()
  422. archive = await archive_factory(printer.id, confirm_requested=True)
  423. provider = await telegram_provider()
  424. _track(poller, provider)
  425. pending_columns = select(TelegramPendingVerdict.message_id, TelegramPendingVerdict.confirm_token)
  426. # Run 1 completes and asks; nobody answers.
  427. monkeypatch.setattr(notification_service, "_http_client", _SendCapture(message_id=901))
  428. assert await dispatch_outcome_confirmation(db_session, printer.id, printer.name, {}, archive.id, None)
  429. run1_token = archive.confirm_token
  430. assert (await db_session.execute(pending_columns)).one() == (901, run1_token)
  431. # Run 2 starts: the reprint reset in main.py (the expected_archive_id path).
  432. archive.status = "printing"
  433. archive.user_verdict = None
  434. archive.user_verdict_source = None
  435. archive.user_verdict_at = None
  436. archive.confirm_token = None
  437. archive.confirm_token_used_at = None
  438. await db_session.commit()
  439. # A thumbs-up on run 1's message while run 2 is printing.
  440. poller._http_client = _FakeBotApi([_updates(_reaction(30, 901, THUMBS_UP))])
  441. await poller.poll_once(BOT_TOKEN)
  442. await db_session.refresh(archive)
  443. assert archive.user_verdict is None
  444. assert poller._http_client.urls() == ["getUpdates"], "the stale message is not edited"
  445. assert await _pending_count(db_session) == 0
  446. # Run 2 completes and asks with a token of its own.
  447. archive.status = "completed"
  448. await db_session.commit()
  449. monkeypatch.setattr(notification_service, "_http_client", _SendCapture(message_id=902))
  450. assert await dispatch_outcome_confirmation(db_session, printer.id, printer.name, {}, archive.id, None)
  451. assert archive.confirm_token not in (None, run1_token)
  452. assert (await db_session.execute(pending_columns)).one() == (902, archive.confirm_token)
  453. # Its own message answers it.
  454. poller._http_client = _FakeBotApi([_updates(_reaction(31, 902, THUMBS_DOWN))])
  455. await poller.poll_once(BOT_TOKEN)
  456. await db_session.refresh(archive)
  457. assert archive.user_verdict == "reject"
  458. assert archive.user_verdict_source == "reaction"
  459. @pytest.mark.asyncio
  460. @pytest.mark.integration
  461. async def test_a_prompt_whose_link_was_used_is_retired(
  462. self, poller, db_session, printer_factory, archive_factory, telegram_provider, pending_factory
  463. ):
  464. """A reaction stops working exactly when a link would: once the token
  465. is spent, or replaced by a newer prompt's, the message is done."""
  466. printer = await printer_factory()
  467. archive = await archive_factory(printer.id, confirm_requested=True, confirm_token="old")
  468. provider = await telegram_provider()
  469. _track(poller, provider)
  470. await pending_factory(provider.id, archive.id, message_id=785)
  471. archive.confirm_token = "newer"
  472. await db_session.commit()
  473. poller._http_client = _FakeBotApi([_updates(_reaction(32, 785, THUMBS_UP))])
  474. await poller.poll_once(BOT_TOKEN)
  475. await db_session.refresh(archive)
  476. assert archive.user_verdict is None
  477. assert await _pending_count(db_session) == 0
  478. assert poller._http_client.urls() == ["getUpdates"]
  479. @pytest.mark.asyncio
  480. @pytest.mark.integration
  481. async def test_second_reaction_on_same_message_is_a_noop(
  482. self, poller, db_session, printer_factory, archive_factory, telegram_provider, pending_factory
  483. ):
  484. printer = await printer_factory()
  485. archive = await archive_factory(printer.id, confirm_requested=True)
  486. provider = await telegram_provider()
  487. _track(poller, provider)
  488. await pending_factory(provider.id, archive.id, message_id=783)
  489. poller._http_client = _FakeBotApi(
  490. [_updates(_reaction(18, 783, THUMBS_UP)), _updates(_reaction(19, 783, THUMBS_DOWN))]
  491. )
  492. await poller.poll_once(BOT_TOKEN)
  493. await poller.poll_once(BOT_TOKEN)
  494. await db_session.refresh(archive)
  495. assert archive.user_verdict == "good"
  496. assert poller._http_client.urls() == ["getUpdates", "editMessageText", "getUpdates"]
  497. @pytest.mark.asyncio
  498. @pytest.mark.integration
  499. async def test_failed_edit_does_not_lose_the_verdict(
  500. self, poller, db_session, printer_factory, archive_factory, telegram_provider, pending_factory
  501. ):
  502. printer = await printer_factory()
  503. archive = await archive_factory(printer.id, confirm_requested=True)
  504. provider = await telegram_provider()
  505. _track(poller, provider)
  506. await pending_factory(provider.id, archive.id, message_id=784)
  507. poller._http_client = _FakeBotApi(
  508. [_updates(_reaction(20, 784, THUMBS_UP)), httpx.ConnectError("telegram gone")]
  509. )
  510. assert await poller.poll_once(BOT_TOKEN) == "ok"
  511. await db_session.refresh(archive)
  512. assert archive.user_verdict == "good"
  513. assert await _pending_count(db_session) == 0
  514. @pytest.mark.asyncio
  515. @pytest.mark.integration
  516. async def test_conflict_is_surfaced_and_cooled_down(self, poller, db_session, telegram_provider):
  517. """409 = webhook set or a second poller: report it on the provider and
  518. wait the cooldown instead of looping on the next getUpdates."""
  519. provider = await telegram_provider()
  520. _track(poller, provider)
  521. conflict = httpx.Response(
  522. 409,
  523. json={"ok": False, "error_code": 409, "description": "Conflict: terminated by other getUpdates request"},
  524. )
  525. poller._http_client = _FakeBotApi([conflict, conflict])
  526. assert await poller.poll_once(BOT_TOKEN) == "conflict"
  527. await db_session.refresh(provider)
  528. assert "Conflict: terminated by other getUpdates request" in provider.last_error
  529. assert "webhook" in provider.last_error
  530. assert provider.last_error_at is not None
  531. sleeps: list[float] = []
  532. async def _sleep(seconds):
  533. sleeps.append(seconds)
  534. raise CancelledForTest
  535. poller._sleep = _sleep
  536. with pytest.raises(CancelledForTest):
  537. await poller._run(BOT_TOKEN)
  538. assert sleeps == [CONFLICT_COOLDOWN]
  539. assert len(poller._http_client.calls) == 2, "one getUpdates per cooldown, not a tight loop"
  540. @pytest.mark.asyncio
  541. @pytest.mark.integration
  542. async def test_bad_bot_token_is_surfaced_and_cooled_down(self, poller, db_session, telegram_provider):
  543. """401/404 = wrong or deleted bot token: it never fixes itself, so it
  544. lands on the provider card and waits the cooldown instead of a warning
  545. every minute forever."""
  546. provider = await telegram_provider()
  547. _track(poller, provider)
  548. unauthorized = httpx.Response(401, json={"ok": False, "error_code": 401, "description": "Unauthorized"})
  549. poller._http_client = _FakeBotApi([unauthorized, unauthorized])
  550. assert await poller.poll_once(BOT_TOKEN) == "rejected"
  551. await db_session.refresh(provider)
  552. assert "HTTP 401" in provider.last_error
  553. assert "Unauthorized" in provider.last_error
  554. assert "bot token" in provider.last_error
  555. assert provider.last_error_at is not None
  556. sleeps: list[float] = []
  557. async def _sleep(seconds):
  558. sleeps.append(seconds)
  559. raise CancelledForTest
  560. poller._sleep = _sleep
  561. with pytest.raises(CancelledForTest):
  562. await poller._run(BOT_TOKEN)
  563. assert sleeps == [CONFLICT_COOLDOWN]
  564. assert len(poller._http_client.calls) == 2
  565. @pytest.mark.asyncio
  566. @pytest.mark.integration
  567. async def test_conflict_is_recorded_on_every_provider_of_the_bot(self, poller, db_session, telegram_provider):
  568. first = await telegram_provider(name="X1C")
  569. second = await telegram_provider(name="P1S")
  570. _track(poller, first, second)
  571. conflict = httpx.Response(409, json={"ok": False, "error_code": 409, "description": "Conflict"})
  572. poller._http_client = _FakeBotApi([conflict])
  573. assert await poller.poll_once(BOT_TOKEN) == "conflict"
  574. for provider in (first, second):
  575. await db_session.refresh(provider)
  576. assert provider.last_error and "Conflict" in provider.last_error
  577. @pytest.mark.asyncio
  578. @pytest.mark.integration
  579. async def test_network_errors_back_off_exponentially(self, poller, telegram_provider):
  580. provider = await telegram_provider()
  581. _track(poller, provider)
  582. poller._http_client = _FakeBotApi([httpx.ConnectError("down")] * 10)
  583. sleeps: list[float] = []
  584. async def _sleep(seconds):
  585. sleeps.append(seconds)
  586. if len(sleeps) == 8:
  587. raise CancelledForTest
  588. poller._sleep = _sleep
  589. with pytest.raises(CancelledForTest):
  590. await poller._run(BOT_TOKEN)
  591. assert sleeps == [1, 2, 4, 8, 16, 32, 60, 60]
  592. assert max(sleeps) == MAX_BACKOFF
  593. @pytest.mark.asyncio
  594. @pytest.mark.integration
  595. async def test_prune_drops_rows_older_than_a_week(
  596. self, poller, db_session, printer_factory, archive_factory, telegram_provider, pending_factory
  597. ):
  598. printer = await printer_factory()
  599. archive = await archive_factory(printer.id, confirm_requested=True)
  600. provider = await telegram_provider()
  601. _track(poller, provider)
  602. await pending_factory(provider.id, archive.id, message_id=1, created_at=utcnow_naive() - timedelta(days=8))
  603. fresh = await pending_factory(provider.id, archive.id, message_id=2)
  604. await poller.prune_stale()
  605. rows = (await db_session.execute(select(TelegramPendingVerdict))).scalars().all()
  606. assert [r.id for r in rows] == [fresh.id]
  607. @pytest.mark.asyncio
  608. @pytest.mark.integration
  609. async def test_sync_follows_provider_mode_enabled_and_token(self, poller, db_session, telegram_provider):
  610. never = _NeverAnswers()
  611. poller._http_client = never
  612. reactions = await telegram_provider(name="R", telegram_verdict_mode="reactions")
  613. buttons = await telegram_provider(name="B", telegram_verdict_mode="buttons")
  614. disabled = await telegram_provider(name="D", telegram_verdict_mode="both", enabled=False)
  615. await poller.sync()
  616. assert poller.is_polling(reactions.id)
  617. assert not poller.is_polling(buttons.id)
  618. assert not poller.is_polling(disabled.id)
  619. assert set(poller._tasks) == {BOT_TOKEN}
  620. # The bot keeps its single poll while any provider wants it; the
  621. # provider set behind it follows the edits.
  622. first_task = poller._tasks[BOT_TOKEN]
  623. reactions.telegram_verdict_mode = "buttons"
  624. buttons.telegram_verdict_mode = "both"
  625. await db_session.commit()
  626. await poller.sync()
  627. assert not poller.is_polling(reactions.id)
  628. assert poller.is_polling(buttons.id)
  629. assert poller._tasks[BOT_TOKEN] is first_task
  630. assert poller._providers[BOT_TOKEN] == {buttons.id}
  631. # A token change moves the provider to a new bot: the old bot's poll
  632. # stops because nobody uses it any more, a fresh one starts.
  633. buttons.config = '{"bot_token": "999:newtoken", "chat_id": "1"}'
  634. await db_session.commit()
  635. await poller.sync()
  636. assert first_task.cancelled() or first_task.cancelling()
  637. assert set(poller._tasks) == {"999:newtoken"}
  638. assert poller.is_polling(buttons.id)
  639. await poller.aclose()
  640. assert poller._tasks == {}
  641. assert never.is_closed
  642. @pytest.mark.asyncio
  643. @pytest.mark.integration
  644. async def test_providers_sharing_a_bot_share_one_poll(
  645. self, poller, db_session, printer_factory, archive_factory, telegram_provider, pending_factory
  646. ):
  647. """The usual farm layout: one bot, one provider per printer. Telegram
  648. serves a single getUpdates consumer per bot, so both providers ride
  649. one task and a reaction to either prompt lands on its archive."""
  650. poller._http_client = _NeverAnswers()
  651. x1c = await printer_factory(name="X1C")
  652. p1s = await printer_factory(name="P1S")
  653. first = await telegram_provider(name="X1C", printer_id=x1c.id, config={"bot_token": BOT_TOKEN, "chat_id": "1"})
  654. second = await telegram_provider(name="P1S", printer_id=p1s.id, config={"bot_token": BOT_TOKEN, "chat_id": "2"})
  655. await poller.sync()
  656. assert len(poller._tasks) == 1
  657. assert poller.is_polling(first.id) and poller.is_polling(second.id)
  658. poller.stop()
  659. _track(poller, first, second)
  660. archive_a = await archive_factory(x1c.id, confirm_requested=True)
  661. archive_b = await archive_factory(p1s.id, confirm_requested=True)
  662. await pending_factory(first.id, archive_a.id, message_id=100, chat_id="1")
  663. await pending_factory(second.id, archive_b.id, message_id=100, chat_id="2")
  664. poller._http_client = _FakeBotApi(
  665. [_updates(_reaction(30, 100, THUMBS_DOWN, chat_id="2"), _reaction(31, 100, THUMBS_UP, chat_id="1"))]
  666. )
  667. await poller.poll_once(BOT_TOKEN)
  668. await db_session.refresh(archive_a)
  669. await db_session.refresh(archive_b)
  670. assert archive_a.user_verdict == "good"
  671. assert archive_b.user_verdict == "reject"
  672. assert await _pending_count(db_session) == 0
  673. assert poller._http_client.urls() == ["getUpdates", "editMessageText", "editMessageText"]
  674. @pytest.mark.asyncio
  675. @pytest.mark.integration
  676. async def test_reaction_for_another_bots_provider_is_not_matched(
  677. self, poller, db_session, printer_factory, archive_factory, telegram_provider, pending_factory
  678. ):
  679. """Two bots in private chat with the same person number their
  680. messages independently, so a match is scoped to the bot's providers."""
  681. printer = await printer_factory()
  682. archive = await archive_factory(printer.id, confirm_requested=True)
  683. other = await telegram_provider(name="other bot", config={"bot_token": "999:other", "chat_id": CHAT_ID})
  684. mine = await telegram_provider(name="mine")
  685. _track(poller, mine)
  686. _track(poller, other, token="999:other")
  687. await pending_factory(other.id, archive.id, message_id=5)
  688. poller._http_client = _FakeBotApi([_updates(_reaction(40, 5, THUMBS_UP))])
  689. await poller.poll_once(BOT_TOKEN)
  690. await db_session.refresh(archive)
  691. assert archive.user_verdict is None
  692. assert await _pending_count(db_session) == 1
  693. class CancelledForTest(BaseException):
  694. """Escape hatch out of the endless poll loop without tripping its handlers."""
  695. class _NeverAnswers:
  696. """getUpdates that stays open forever, like a quiet chat would."""
  697. def __init__(self):
  698. self.is_closed = False
  699. async def post(self, url, json=None, data=None, files=None):
  700. import asyncio
  701. await asyncio.Event().wait()
  702. async def aclose(self):
  703. self.is_closed = True
  704. # ---------------------------------------------------------------------------
  705. # Provider API round trip
  706. # ---------------------------------------------------------------------------
  707. class TestProviderApi:
  708. @pytest.mark.asyncio
  709. @pytest.mark.integration
  710. async def test_mode_round_trips_and_defaults_to_buttons(self, async_client: AsyncClient):
  711. with patch("backend.app.api.routes.notifications.telegram_reaction_poller.sync", new=AsyncMock()) as sync:
  712. created = await async_client.post(
  713. "/api/v1/notifications/",
  714. json={
  715. "name": "TG",
  716. "provider_type": "telegram",
  717. "config": TELEGRAM_CONFIG,
  718. "telegram_verdict_mode": "reactions",
  719. },
  720. )
  721. assert created.status_code == 200, created.text
  722. assert created.json()["telegram_verdict_mode"] == "reactions"
  723. provider_id = created.json()["id"]
  724. assert sync.await_count == 1
  725. listed = await async_client.get("/api/v1/notifications/")
  726. assert [p["telegram_verdict_mode"] for p in listed.json() if p["id"] == provider_id] == ["reactions"]
  727. patched = await async_client.patch(
  728. f"/api/v1/notifications/{provider_id}", json={"telegram_verdict_mode": "both"}
  729. )
  730. assert patched.status_code == 200
  731. assert patched.json()["telegram_verdict_mode"] == "both"
  732. assert sync.await_count == 2
  733. invalid = await async_client.patch(
  734. f"/api/v1/notifications/{provider_id}", json={"telegram_verdict_mode": "webhook"}
  735. )
  736. assert invalid.status_code == 422
  737. other = await async_client.post(
  738. "/api/v1/notifications/",
  739. json={"name": "ntfy", "provider_type": "ntfy", "config": {"topic": "t"}},
  740. )
  741. assert other.json()["telegram_verdict_mode"] == "buttons"
  742. deleted = await async_client.delete(f"/api/v1/notifications/{provider_id}")
  743. assert deleted.status_code == 200
  744. assert sync.await_count == 4
  745. @pytest.mark.asyncio
  746. @pytest.mark.integration
  747. async def test_legacy_null_mode_reads_as_buttons(
  748. self, async_client: AsyncClient, notification_provider_factory, db_session
  749. ):
  750. provider = await notification_provider_factory(name="Legacy")
  751. # The ORM default fills a None at INSERT; a pre-#3046 row holds a real NULL.
  752. await db_session.execute(
  753. text("UPDATE notification_providers SET telegram_verdict_mode = NULL WHERE id = :id"), {"id": provider.id}
  754. )
  755. await db_session.commit()
  756. stored = await db_session.scalar(
  757. text("SELECT telegram_verdict_mode FROM notification_providers WHERE id = :id"), {"id": provider.id}
  758. )
  759. assert stored is None, "row under test must actually hold NULL"
  760. response = await async_client.get(f"/api/v1/notifications/{provider.id}")
  761. assert response.status_code == 200
  762. assert response.json()["telegram_verdict_mode"] == "buttons"
  763. @pytest.mark.asyncio
  764. @pytest.mark.integration
  765. async def test_deleting_a_provider_cascades_its_pending_rows(
  766. self,
  767. async_client: AsyncClient,
  768. db_session,
  769. printer_factory,
  770. archive_factory,
  771. telegram_provider,
  772. pending_factory,
  773. ):
  774. printer = await printer_factory()
  775. archive = await archive_factory(printer.id, confirm_requested=True)
  776. provider = await telegram_provider()
  777. await pending_factory(provider.id, archive.id)
  778. with patch("backend.app.api.routes.notifications.telegram_reaction_poller.sync", new=AsyncMock()):
  779. assert (await async_client.delete(f"/api/v1/notifications/{provider.id}")).status_code == 200
  780. assert await _pending_count(db_session) == 0