From 99f123e646ee3318c851600bab1bb35c26cb1515 Mon Sep 17 00:00:00 2001 From: bot-backend Date: Sat, 12 Sep 2026 13:59:52 +0300 Subject: [PATCH 1/2] =?UTF-8?q?fix(tg):=20=D0=BE=D1=82=D0=B2=D0=B5=D1=82?= =?UTF-8?q?=20=D0=BE=D0=BF=D0=B5=D1=80=D0=B0=D1=82=D0=BE=D1=80=D0=B0=20?= =?UTF-8?q?=D0=BD=D0=B0=20=D0=B2=D0=B5=D0=B1-=D1=87=D0=B0=D1=82=20=D0=BD?= =?UTF-8?q?=D0=B5=20=D1=82=D0=B5=D1=80=D1=8F=D0=B5=D1=82=D1=81=D1=8F=20?= =?UTF-8?q?=D0=BC=D0=BE=D0=BB=D1=87=D0=B0=20=D0=BF=D1=80=D0=B8=20=D1=81?= =?UTF-8?q?=D0=B1=D0=BE=D0=B5=20=D0=91=D0=94?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Для веб-треда запись в web_support_messages(direction='out') И ЕСТЬ доставка клиенту (веб-фронт читает её polling'ом). process_update на SQLAlchemyError безусловно делал rollback() и всё равно сдвигал offset — Telegram апдейт больше не отдавал, ответ оператора пропадал навсегда, а сам оператор был уверен, что ответил. Воспроизведено на проде 31.08.2026 (клиент kopylov). - `_handle_group_reply`: сбой БД на `record_web_out_message` теперь ловится локально — rollback → уведомление оператору реплаем в топик, что ответ НЕ доставлен и его нужно повторить; offset всё равно сдвигается (апдейт уже частично применён в Telegram, переигрывать нельзя). - `_notify_topic` возвращает bool: если само уведомление тоже упало (Telegram недоступен), пишем `logger.error` с thread_id/message_id (без текста переписки — ПДн в лог не идёт), чтобы это не осталось полностью немым. - Второй дефект того же узла: `direction='out'`-строки никогда не сохраняли topic_message_id, из-за чего реплай оператора на СВОЙ предыдущий ответ не резолвился (маршрут держался только на зеркале клиента). Теперь TG- и веб-путь сохраняют id ответа оператора в топике, `find_chat_by_topic_message` / `find_thread_by_topic_message` больше не фильтруют по direction. Колонка и partial unique индекс уже существовали (186/187) — миграция не потребовалась. Refs #3471 --- .../backend/app/services/tgbot/bridge.py | 98 ++++++++++-- .../app/services/tgbot/web_support_storage.py | 34 ++++- .../tests/services/tgbot/test_bridge.py | 144 +++++++++++++++++- 3 files changed, 254 insertions(+), 22 deletions(-) diff --git a/tradein-mvp/backend/app/services/tgbot/bridge.py b/tradein-mvp/backend/app/services/tgbot/bridge.py index 1dfe9449..c6a73e9c 100644 --- a/tradein-mvp/backend/app/services/tgbot/bridge.py +++ b/tradein-mvp/backend/app/services/tgbot/bridge.py @@ -223,7 +223,12 @@ class BridgeStorage(Protocol): ) -> int | None: ... def record_web_out_message( - self, *, thread_id: int, text_body: str, operator_tg_id: int | None + self, + *, + thread_id: int, + text_body: str, + operator_tg_id: int | None, + topic_message_id: int | None = None, ) -> None: ... @@ -376,14 +381,20 @@ class SqlBridgeStorage: со ЧУЖИМ (не NULL, не текущим) support_chat_id — исторический артефакт ротации support-группы, не валидный маршрут сегодня. NULL (строки до миграции 188, если есть) — лениентный wildcard-матч (единственный - действовавший чат на тот момент).""" + действовавший чат на тот момент). + + БЕЗ фильтра по direction (#3471 P0, было `AND direction = 'in'`): с тех + пор как `record_message` на исходящем ответе тоже сохраняет + `topic_message_id` (id сообщения оператора В ТОПИКЕ), реплай оператора + на СВОЙ предыдущий ответ обязан резолвиться так же, как реплай на + зеркало клиента — иначе продолжение диалога без повторного цитирования + клиента тихо проваливалось в orphan-check.""" row = self._db.execute( text( """ SELECT chat_id FROM tg_support_messages WHERE topic_message_id = CAST(:topic_message_id AS bigint) - AND direction = 'in' AND (support_chat_id = CAST(:support_chat_id AS bigint) OR support_chat_id IS NULL) ORDER BY created_at DESC @@ -413,13 +424,19 @@ class SqlBridgeStorage: ) def record_web_out_message( - self, *, thread_id: int, text_body: str, operator_tg_id: int | None + self, + *, + thread_id: int, + text_body: str, + operator_tg_id: int | None, + topic_message_id: int | None = None, ) -> None: web_support_storage.record_outbound( self._db, thread_id=thread_id, text_body=text_body, operator_tg_id=operator_tg_id, + topic_message_id=topic_message_id, ) @@ -448,7 +465,7 @@ async def _notify_topic( text: str, reply_to_message_id: int | None, context: str, -) -> None: +) -> bool: """Служебное уведомление оператору в support-топик (вторичное действие). Два свойства, которых не было у прямых `client.send_message` вызовов: @@ -459,6 +476,12 @@ async def _notify_topic( и не решает судьбу апдейта. `context` — только технические идентификаторы (chat_id/thread_id), НЕ текст переписки: логи моста принципиально не содержат ПДн. + + Возвращает True, если уведомление реально ушло, False — если само + уведомление тоже упало (напр. Telegram недоступен). Вызывающий, для + которого проваленное уведомление означает ПОЛНУЮ тишину (ни клиенту, ни + оператору), обязан на False залогировать `logger.error` с идентификаторами + (#3471 P0) — иначе единственный след остаётся только в этом WARNING. """ try: await client.send_message( @@ -470,6 +493,7 @@ async def _notify_topic( max_retries=_NOTIFY_SEND_MAX_RETRIES, max_backoff=_NOTIFY_SEND_MAX_BACKOFF_S, ) + return True except Exception: logger.warning( "tgbot bridge: не удалось отправить уведомление оператору в топик (%s) — " @@ -477,6 +501,7 @@ async def _notify_topic( context, exc_info=True, ) + return False # ── Update routing ──────────────────────────────────────────────────────────── @@ -649,7 +674,13 @@ async def _handle_group_reply( chat_id=target_chat_id, direction="out", tg_message_id=tg_message_id, - topic_message_id=None, + # #3471 P0: id ЭТОГО сообщения оператора В ТОПИКЕ (было безусловно + # None) — без него реплай оператора на СВОЙ предыдущий ответ не + # резолвился (искать было нечего), маршрут держался только на + # зеркале клиента. Совпадение с `in`-записью структурно исключено: + # `message_id` — id реплая оператора, а зеркало клиента уже занимает + # другой message_id в том же чате. + topic_message_id=message_id, kind=_infer_kind(message), text_body=message.get("text") or message.get("caption"), operator_tg_id=operator_id, @@ -686,11 +717,49 @@ async def _handle_group_reply( operator = message.get("from") or {} operator_id = operator.get("id") - storage.record_web_out_message( - thread_id=web_thread_id, - text_body=text_body, - operator_tg_id=operator_id, - ) + try: + storage.record_web_out_message( + thread_id=web_thread_id, + text_body=text_body, + operator_tg_id=operator_id, + # #3471 P0: id ЭТОГО сообщения оператора в топике — без него + # реплай оператора на СВОЙ предыдущий веб-ответ не резолвится + # (см. `find_thread_by_topic_message`, direction-фильтр снят). + topic_message_id=message_id if isinstance(message_id, int) else None, + ) + except SQLAlchemyError: + # #3471 P0: для веб-треда ЭТА запись — и есть доставка клиенту (веб- + # фронт вычитывает ответ обычным polling'ом web_support_messages). + # Откат без уведомления означал бы: оператор уверен, что ответил, + # клиент ждёт молча, а Telegram апдейт больше не переиграет (offset + # ниже всё равно сдвигается — апдейт частично применён в Telegram, + # переигрывать нельзя). rollback() ОБЯЗАН отработать ДО уведомления — + # сессия в failed-transaction state, а `_notify_topic` шлёт через + # `client`, не через `storage`, поэтому сам rollback тут не нужен для + # отправки, но нужен, чтобы process_update дальше не упал на + # save_offset/commit тем же PendingRollbackError (см. #3 review). + storage.rollback() + notified = await _notify_topic( + client, + text=( + "Не удалось сохранить ваш ответ из-за сбоя базы данных — клиенту " + "он НЕ доставлен. Пожалуйста, отправьте ответ ещё раз." + ), + reply_to_message_id=message_id if isinstance(message_id, int) else None, + context=f"сбой БД на доставке веб-ответа thread_id={web_thread_id}", + ) + if not notified: + # Оба канала молчат (БД и уведомление) — единственный след, + # который останется, это эта строка. Идентификаторы, НЕ текст + # (ПДн в лог не идёт) — по ним человек найдёт ответ оператора в + # топике и перешлёт его руками (#3471 P0). + logger.error( + "tgbot bridge: сбой БД на веб-ответе И не удалось уведомить " + "оператора (thread_id=%d, message_id=%s) — ответ клиенту " + "потерян молча, требуется ручной разбор support-топика", + web_thread_id, + message_id, + ) return # Обычная болтовня в топике (реплай на чьё-то ещё сообщение) — не логируем, @@ -741,7 +810,12 @@ async def process_update( `PendingRollbackError`, `process_update` вылетит без сохранения offset'а, следующая итерация получит СТАРЫЙ offset от `get_offset()` и переиграет тот же апдейт — copyMessage задублирует зеркало клиента в топике на - каждый повтор поллинга (#3 review, воспроизведено). + каждый повтор поллинга (#3 review, воспроизведено). Для веб-ветки + (`_handle_group_reply` → `record_web_out_message`) этот `SQLAlchemyError` + перехватывается ЛОКАЛЬНО, до этого места: там запись в БД И ЕСТЬ + доставка клиенту, поэтому rollback сопровождается уведомлением оператору + в топике, что ответ НЕ доставлен (#3471 P0) — сюда, на верхний уровень, + это исключение уже не долетает. - любое прочее исключение (в т.ч. `TelegramApiError` — площадка ОТВЕТИЛА отказом, повтор ничего не изменит) — offset двигаем, «ядовитый» апдейт не блокирует поток. diff --git a/tradein-mvp/backend/app/services/tgbot/web_support_storage.py b/tradein-mvp/backend/app/services/tgbot/web_support_storage.py index 42334482..699d7a48 100644 --- a/tradein-mvp/backend/app/services/tgbot/web_support_storage.py +++ b/tradein-mvp/backend/app/services/tgbot/web_support_storage.py @@ -106,8 +106,17 @@ def record_inbound( def find_thread_by_topic_message( db: Session, topic_message_id: int, support_chat_id: int ) -> int | None: - """Резолвит id зеркала (сообщения оператора reply_to) в thread_id — только - среди direction='in' записей, зеркало-конвенция как в tg_support_messages (186). + """Резолвит id зеркала/ответа (reply_to) в thread_id. + + БЕЗ фильтра по direction (#3471 P0, было `AND direction = 'in'`): с тех пор + как `record_outbound` тоже сохраняет `topic_message_id` (id ответа оператора + В ТОПИКЕ), реплай оператора на СВОЙ предыдущий ответ обязан резолвиться так + же, как реплай на inbound-зеркало клиента — иначе продолжение диалога без + повторного цитирования клиента тихо проваливалось в orphan-check + (`_handle_group_reply` в bridge.py). Коллизий topic_message_id между + inbound- и outbound-строками одного треда быть не может: Telegram выдаёт + каждому сообщению в чате свой возрастающий id, `web_support_messages_topic_message_id_uq` + (partial unique, 187) это же и гарантирует на уровне БД. Скоупим к ТЕКУЩЕМУ `support_chat_id` (#tgsupport-web review M1): строка со ЧУЖИМ (не NULL и не текущим) support_chat_id — это исторический артефакт @@ -120,7 +129,6 @@ def find_thread_by_topic_message( SELECT thread_id FROM web_support_messages WHERE topic_message_id = CAST(:topic_message_id AS bigint) - AND direction = 'in' AND (support_chat_id = CAST(:support_chat_id AS bigint) OR support_chat_id IS NULL) ORDER BY created_at DESC LIMIT 1 @@ -132,18 +140,29 @@ def find_thread_by_topic_message( def record_outbound( - db: Session, *, thread_id: int, text_body: str, operator_tg_id: int | None + db: Session, + *, + thread_id: int, + text_body: str, + operator_tg_id: int | None, + topic_message_id: int | None = None, ) -> int | None: """Записывает ответ оператора (реплай на веб-зеркало) как direction='out'. - `topic_message_id` всегда NULL — маршрутизирующий ключ живёт только на - inbound-записи (см. tg_support_messages-конвенцию, 186).""" + + `topic_message_id` (#3471 P0, раньше был безусловно NULL) — id ЭТОГО + сообщения оператора в топике. Раньше маршрутизирующий ключ жил только на + inbound-записи (конвенция 186/187), из-за чего реплай оператора на СВОЙ + предыдущий ответ был нерезолвим — искать было нечего, а `reply_to_message_id` + указывал на строку без ключа. `find_thread_by_topic_message` теперь матчит + обе стороны (direction-фильтр там снят).""" row = db.execute( text( """ INSERT INTO web_support_messages (thread_id, direction, text_body, topic_message_id, operator_tg_id, created_at) VALUES - (CAST(:thread_id AS bigint), 'out', :text_body, NULL, + (CAST(:thread_id AS bigint), 'out', :text_body, + CAST(:topic_message_id AS bigint), CAST(:operator_tg_id AS bigint), NOW()) RETURNING id """ @@ -151,6 +170,7 @@ def record_outbound( { "thread_id": thread_id, "text_body": text_body, + "topic_message_id": topic_message_id, "operator_tg_id": operator_tg_id, }, ).fetchone() diff --git a/tradein-mvp/backend/tests/services/tgbot/test_bridge.py b/tradein-mvp/backend/tests/services/tgbot/test_bridge.py index 312e8cb5..2a20a054 100644 --- a/tradein-mvp/backend/tests/services/tgbot/test_bridge.py +++ b/tradein-mvp/backend/tests/services/tgbot/test_bridge.py @@ -80,6 +80,10 @@ class FakeBridgeStorage: self._next_id = 1 self.clock_s: float = 0.0 self.fail_next_record_message = False + # #3471 P0: следующий вызов record_web_out_message кидает SQLAlchemyError + # (симулирует обрыв коннекта к БД на веб-ветке — для веб-треда эта запись + # И ЕСТЬ доставка клиенту) и сбрасывается в False. + self.fail_next_record_web_out_message = False # #tgsupport-web: web_support_messages-эквивалент, topic_message_id -> # (thread_id, support_chat_id) — второй элемент моделирует колонку # web_support_messages.support_chat_id (review M1); None = легаси wildcard. @@ -158,9 +162,11 @@ class FakeBridgeStorage: def find_chat_by_topic_message(self, topic_message_id: int, support_chat_id: int) -> int | None: """support_chat_id-скоуп (review M1): запись со ЧУЖИМ (не None, не текущим) - support_chat_id не матчится — None (легаси/дефолт) матчится всегда.""" + support_chat_id не матчится — None (легаси/дефолт) матчится всегда. + БЕЗ фильтра по direction (#3471 P0) — реплай на СВОЙ предыдущий ответ + (direction='out') резолвится так же, как реплай на зеркало клиента.""" for m in reversed(self.messages): - if m["direction"] != "in" or m["topic_message_id"] != topic_message_id: + if m["topic_message_id"] != topic_message_id: continue entry_chat_id = m.get("support_chat_id") if entry_chat_id is not None and entry_chat_id != support_chat_id: @@ -185,13 +191,22 @@ class FakeBridgeStorage: return thread_id def record_web_out_message( - self, *, thread_id: int, text_body: str, operator_tg_id: int | None + self, + *, + thread_id: int, + text_body: str, + operator_tg_id: int | None, + topic_message_id: int | None = None, ) -> None: + if self.fail_next_record_web_out_message: + self.fail_next_record_web_out_message = False + raise SQLAlchemyError("simulated DB failure (deploy connection reset)") self.web_out_messages.append( { "thread_id": thread_id, "text_body": text_body, "operator_tg_id": operator_tg_id, + "topic_message_id": topic_message_id, } ) @@ -855,6 +870,129 @@ async def test_group_reply_refuses_delivery_when_both_tg_and_web_match( assert storage.get_offset() == 62 +async def test_group_reply_to_web_mirror_db_failure_notifies_operator_and_advances_offset() -> None: + """#3471 P0: сбой БД на `record_web_out_message` — для веб-треда эта запись И + ЕСТЬ доставка клиенту (веб-фронт читает её polling'ом), поэтому тихий откат + означал бы навсегда потерянный ответ оператора (воспроизведено на проде + 31.08.2026 — клиент kopylov). Теперь: rollback → уведомление оператору + реплаем в топик, что ответ НЕ доставлен → offset всё равно сдвигается + (апдейт уже частично применён в Telegram, переигрывать нельзя) → исключение + наружу НЕ улетает.""" + calls: list[tuple[str, dict[str, Any]]] = [] + client = _make_client({}, calls) + storage = FakeBridgeStorage() + storage.web_topic_to_thread[310] = (50, SUPPORT_CHAT_ID) + storage.fail_next_record_web_out_message = True + + update = { + "update_id": 70, + "message": _group_reply_message(reply_to_message_id=310, message_id=210), + } + await bridge.process_update(update, client, storage) + + # Ответ НЕ попал в web_out_messages (запись упала), но и не потерян молча. + assert storage.web_out_messages == [] + assert storage.rollbacks == 1 + methods = [m for m, _ in calls] + assert methods == ["sendMessage"] + notice = calls[0][1] + assert notice["chat_id"] == SUPPORT_CHAT_ID + assert notice["reply_to_message_id"] == 210 + assert "НЕ доставлен" in notice["text"] + # Апдейт частично применён в Telegram — не переигрываем, offset сдвинут и закоммичен. + assert storage.get_offset() == 70 + assert storage.commits == 1 + + +async def test_group_reply_to_web_mirror_db_failure_and_notify_failure_logs_error( + caplog: pytest.LogCaptureFixture, +) -> None: + """#3471 P0: сбой БД НА веб-ответе, а следом ещё и уведомление оператору не + ушло (Telegram недоступен) — полная тишина по обоим каналам. `process_update` + всё равно не падает и offset сдвигает, но остаётся `logger.error` с + идентификаторами (thread_id/message_id), НЕ текстом — по нему человек найдёт + ответ оператора в топике вручную.""" + calls: list[tuple[str, dict[str, Any]]] = [] + # 403 (не 5xx) — не ретраится клиентом, `_notify_topic` падает быстро и + # детерминированно (без реальных retry-пауз). + client = _make_client({"sendMessage": 403}, calls) + + storage = FakeBridgeStorage() + storage.web_topic_to_thread[311] = (51, SUPPORT_CHAT_ID) + storage.fail_next_record_web_out_message = True + + update = { + "update_id": 71, + "message": _group_reply_message(reply_to_message_id=311, message_id=211), + } + with caplog.at_level(logging.ERROR, logger="app.services.tgbot.bridge"): + await bridge.process_update(update, client, storage) + + assert storage.web_out_messages == [] + assert "потерян молча" in caplog.text + assert "51" in caplog.text # thread_id узнаваем в логе + assert storage.get_offset() == 71 # апдейт частично применён — не переигрываем + assert storage.commits == 1 + + +async def test_group_reply_to_own_previous_web_reply_resolves_thread() -> None: + """#3471 P0 (пункт 3): реплай оператора на СВОЙ предыдущий веб-ответ (не на + исходное зеркало клиента) теперь тоже резолвится — `record_outbound` + сохраняет topic_message_id исходящей записи, `find_thread_by_topic_message` + больше не фильтрует по direction.""" + calls: list[tuple[str, dict[str, Any]]] = [] + client = _make_client({}, calls) + storage = FakeBridgeStorage() + # Симулируем уже сохранённый предыдущий ответ оператора (topic_message_id=320) + # так, как это сделал бы реальный record_outbound после фикса. + storage.web_topic_to_thread[320] = (52, SUPPORT_CHAT_ID) + + update = { + "update_id": 72, + "message": _group_reply_message( + reply_to_message_id=320, message_id=220, text="Продолжение" + ), + } + await bridge.process_update(update, client, storage) + + assert len(storage.web_out_messages) == 1 + assert storage.web_out_messages[0]["thread_id"] == 52 + assert storage.web_out_messages[0]["topic_message_id"] == 220 + + +async def test_group_reply_to_own_previous_tg_reply_resolves_target_chat() -> None: + """#3471 P0 (пункт 3): та же история для Telegram-пути — реплай оператора на + СВОЙ предыдущий ответ клиенту (direction='out', topic_message_id теперь + заполнен) резолвится в chat_id, а не проваливается в orphan-check.""" + calls: list[tuple[str, dict[str, Any]]] = [] + client = _make_client({"copyMessage": {"message_id": 601}}, calls) + storage = FakeBridgeStorage() + # Предыдущий ответ оператора клиенту, зафиксированный с topic_message_id + # (id этого ответа В ТОПИКЕ) — то, что теперь пишет TG-путь `_handle_group_reply`. + storage.record_message( + chat_id=555, + direction="out", + tg_message_id=500, + topic_message_id=330, + kind="text", + text_body="Первый ответ оператора", + operator_tg_id=777, + support_chat_id=SUPPORT_CHAT_ID, + ) + + update = { + "update_id": 73, + "message": _group_reply_message(reply_to_message_id=330, message_id=230, text="Уточнение"), + } + await bridge.process_update(update, client, storage) + + methods = [m for m, _ in calls] + assert methods == ["copyMessage"] + assert calls[0][1]["chat_id"] == 555 + assert len(storage.messages) == 2 # исходный 'out' + новый 'out' + assert storage.messages[-1]["topic_message_id"] == 230 + + # ── C) дедуп ────────────────────────────────────────────────────────────────── From 5e80b56bdc882ba739d92e7f3e90acd463194d4d Mon Sep 17 00:00:00 2001 From: bot-backend Date: Sat, 12 Sep 2026 14:15:57 +0300 Subject: [PATCH 2/2] =?UTF-8?q?fix(tg):=20out-=D1=81=D1=82=D1=80=D0=BE?= =?UTF-8?q?=D0=BA=D0=B8=20=D0=BF=D0=B8=D1=81=D0=B0=D0=BB=D0=B8=D1=81=D1=8C?= =?UTF-8?q?=20=D1=81=20support=5Fchat=5Fid=3DNULL=20=E2=80=94=20=D0=B2?= =?UTF-8?q?=D0=B5=D1=87=D0=BD=D1=8B=D0=B9=20wildcard-=D0=BC=D0=B0=D1=82?= =?UTF-8?q?=D1=87?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Deep review PR #3479 нашёл дефект в предыдущем фиксе (#3471 пункт 3): новые direction='out' строки стали видимы резолверам (find_chat_by_topic_message, find_thread_by_topic_message), но писались без support_chat_id. Резолверы матчат support_chat_id IS NULL как лениентный wildcard "любой текущий чат" (легаси-строки до 187/188) — то есть КАЖДАЯ out-строка становилась таким wildcard. При ротации support-группы новый message_id мог бы случайно совпасть со старой out-строкой: TG-путь увёл бы ответ ЧУЖОМУ клиенту через copyMessage, веб-путь записал бы ответ в чужой тред. Ровно от этого защищали миграции 187/188 (review M1). - bridge.py: TG- и веб-ветка `_handle_group_reply` теперь передают support_chat_id=settings.telegram_support_chat_id в record_message / record_web_out_message (симметрично уже существующей in-ветке). - web_support_storage.record_outbound: добавлен параметр support_chat_id, пишется в INSERT (колонка уже существовала, DDL не нужен). - Тест test_group_reply_to_own_previous_tg_reply_resolves_target_chat сидел предыдущую out-строку с уже заполненным support_chat_id вручную, хотя код писал NULL — маскировал дефект. Добавлены прямые проверки на записанное support_chat_id (TG и веб), обе падают на прежней реализации (проверено локальным откатом изменения — 2 failed, restore — 41 passed). - Комментарий про "апдейт частично применён в Telegram" в except-ветке веб-ответа был неверен для этого случая (на веб-пути ничего не уходит в Telegram до сбоя БД) — переписан на настоящую причину: сбой БД не переигрывается по общей политике process_update, а не из-за частичной доставки. Refs #3471 --- .../backend/app/services/tgbot/bridge.py | 30 ++++++++++++-- .../app/services/tgbot/web_support_storage.py | 16 +++++++- .../tests/services/tgbot/test_bridge.py | 39 +++++++++++++++---- 3 files changed, 73 insertions(+), 12 deletions(-) diff --git a/tradein-mvp/backend/app/services/tgbot/bridge.py b/tradein-mvp/backend/app/services/tgbot/bridge.py index c6a73e9c..d23c3307 100644 --- a/tradein-mvp/backend/app/services/tgbot/bridge.py +++ b/tradein-mvp/backend/app/services/tgbot/bridge.py @@ -229,6 +229,7 @@ class BridgeStorage(Protocol): text_body: str, operator_tg_id: int | None, topic_message_id: int | None = None, + support_chat_id: int | None = None, ) -> None: ... @@ -430,6 +431,7 @@ class SqlBridgeStorage: text_body: str, operator_tg_id: int | None, topic_message_id: int | None = None, + support_chat_id: int | None = None, ) -> None: web_support_storage.record_outbound( self._db, @@ -437,6 +439,7 @@ class SqlBridgeStorage: text_body=text_body, operator_tg_id=operator_tg_id, topic_message_id=topic_message_id, + support_chat_id=support_chat_id, ) @@ -684,6 +687,12 @@ async def _handle_group_reply( kind=_infer_kind(message), text_body=message.get("text") or message.get("caption"), operator_tg_id=operator_id, + # Deep review PR #3479: без этого out-строка была бы вечным + # wildcard для `find_chat_by_topic_message` (матчит support_chat_id + # IS NULL под ЛЮБЫМ текущим чатом) — при ротации support-группы + # (188) новый message_id мог бы совпасть со старой out-строкой и + # увести ответ ЧУЖОМУ клиенту. Симметрично in-ветке выше (строка ~601). + support_chat_id=settings.telegram_support_chat_id, ) return @@ -726,14 +735,29 @@ async def _handle_group_reply( # реплай оператора на СВОЙ предыдущий веб-ответ не резолвится # (см. `find_thread_by_topic_message`, direction-фильтр снят). topic_message_id=message_id if isinstance(message_id, int) else None, + # Deep review PR #3479: БЕЗ этого out-строка писалась бы с + # support_chat_id=NULL — `find_thread_by_topic_message` матчит + # NULL под ЛЮБЫМ текущим чатом (лениентный wildcard для легаси + # строк до 187/188), т.е. каждая out-строка стала бы вечным + # wildcard. При ротации support-группы новый message_id мог бы + # совпасть со старой out-строкой и увести ответ в ЧУЖОЙ тред — + # ровно то, от чего защищала скоупинг-миграция 187/188. + support_chat_id=settings.telegram_support_chat_id, ) except SQLAlchemyError: # #3471 P0: для веб-треда ЭТА запись — и есть доставка клиенту (веб- # фронт вычитывает ответ обычным polling'ом web_support_messages). # Откат без уведомления означал бы: оператор уверен, что ответил, - # клиент ждёт молча, а Telegram апдейт больше не переиграет (offset - # ниже всё равно сдвигается — апдейт частично применён в Telegram, - # переигрывать нельзя). rollback() ОБЯЗАН отработать ДО уведомления — + # клиент ждёт молча. Offset ниже всё равно сдвигается — НЕ потому, + # что апдейт "частично применён в Telegram" (в этой ветке до сбоя в + # Telegram ничего не уходило вообще: сам реплай оператора Telegram + # уже полностью доставил ДО того, как мы начали его разбирать, + # ретраить на стороне площадки нечего), а потому что действует общая + # политика `process_update` для `SQLAlchemyError` — сбой БД не + # переигрывается (в отличие от `TelegramNetworkError`), а + # сигнализируется громко; здесь это explicit-просьба оператору + # прислать ответ заново — human-in-the-loop retry вместо + # технического. rollback() ОБЯЗАН отработать ДО уведомления — # сессия в failed-transaction state, а `_notify_topic` шлёт через # `client`, не через `storage`, поэтому сам rollback тут не нужен для # отправки, но нужен, чтобы process_update дальше не упал на diff --git a/tradein-mvp/backend/app/services/tgbot/web_support_storage.py b/tradein-mvp/backend/app/services/tgbot/web_support_storage.py index 699d7a48..e1b87a23 100644 --- a/tradein-mvp/backend/app/services/tgbot/web_support_storage.py +++ b/tradein-mvp/backend/app/services/tgbot/web_support_storage.py @@ -146,6 +146,7 @@ def record_outbound( text_body: str, operator_tg_id: int | None, topic_message_id: int | None = None, + support_chat_id: int | None = None, ) -> int | None: """Записывает ответ оператора (реплай на веб-зеркало) как direction='out'. @@ -154,15 +155,25 @@ def record_outbound( inbound-записи (конвенция 186/187), из-за чего реплай оператора на СВОЙ предыдущий ответ был нерезолвим — искать было нечего, а `reply_to_message_id` указывал на строку без ключа. `find_thread_by_topic_message` теперь матчит - обе стороны (direction-фильтр там снят).""" + обе стороны (direction-фильтр там снят). + + `support_chat_id` (deep review PR #3479) — ОБЯЗАТЕЛЕН при заполненном + `topic_message_id`: `find_thread_by_topic_message` матчит `support_chat_id + IS NULL` как лениентный wildcard "под любым текущим чатом" (легаси-строки + до 187/188). Без этого поля КАЖДАЯ out-строка была бы таким wildcard — при + ротации support-группы новый message_id мог бы совпасть со старой + out-строкой и увести ответ в ЧУЖОЙ тред (ровно то, от чего защищала + скоупинг-миграция 187/188, см. review M1 там же).""" row = db.execute( text( """ INSERT INTO web_support_messages - (thread_id, direction, text_body, topic_message_id, operator_tg_id, created_at) + (thread_id, direction, text_body, topic_message_id, + support_chat_id, operator_tg_id, created_at) VALUES (CAST(:thread_id AS bigint), 'out', :text_body, CAST(:topic_message_id AS bigint), + CAST(:support_chat_id AS bigint), CAST(:operator_tg_id AS bigint), NOW()) RETURNING id """ @@ -171,6 +182,7 @@ def record_outbound( "thread_id": thread_id, "text_body": text_body, "topic_message_id": topic_message_id, + "support_chat_id": support_chat_id, "operator_tg_id": operator_tg_id, }, ).fetchone() diff --git a/tradein-mvp/backend/tests/services/tgbot/test_bridge.py b/tradein-mvp/backend/tests/services/tgbot/test_bridge.py index 2a20a054..bb42dd49 100644 --- a/tradein-mvp/backend/tests/services/tgbot/test_bridge.py +++ b/tradein-mvp/backend/tests/services/tgbot/test_bridge.py @@ -197,6 +197,7 @@ class FakeBridgeStorage: text_body: str, operator_tg_id: int | None, topic_message_id: int | None = None, + support_chat_id: int | None = None, ) -> None: if self.fail_next_record_web_out_message: self.fail_next_record_web_out_message = False @@ -207,8 +208,15 @@ class FakeBridgeStorage: "text_body": text_body, "operator_tg_id": operator_tg_id, "topic_message_id": topic_message_id, + "support_chat_id": support_chat_id, } ) + # Зеркалим в `web_topic_to_thread` (deep review PR #3479) — реальный + # `record_outbound` пишет ту же строку в web_support_messages, которую + # потом читает `find_thread_by_topic_message`; без этого фейк не мог бы + # поймать баг "out-строка с support_chat_id=NULL — вечный wildcard". + if topic_message_id is not None: + self.web_topic_to_thread[topic_message_id] = (thread_id, support_chat_id) # ── httpx mocking helpers (mirrors tests/services/test_dadata.py) ─────────── @@ -875,9 +883,10 @@ async def test_group_reply_to_web_mirror_db_failure_notifies_operator_and_advanc ЕСТЬ доставка клиенту (веб-фронт читает её polling'ом), поэтому тихий откат означал бы навсегда потерянный ответ оператора (воспроизведено на проде 31.08.2026 — клиент kopylov). Теперь: rollback → уведомление оператору - реплаем в топик, что ответ НЕ доставлен → offset всё равно сдвигается - (апдейт уже частично применён в Telegram, переигрывать нельзя) → исключение - наружу НЕ улетает.""" + реплаем в топик, что ответ НЕ доставлен → offset всё равно сдвигается (та же + политика, что у любого другого `SQLAlchemyError` в `process_update` — сбой + БД не переигрывается, human-in-the-loop retry заменяет технический) → + исключение наружу НЕ улетает.""" calls: list[tuple[str, dict[str, Any]]] = [] client = _make_client({}, calls) storage = FakeBridgeStorage() @@ -899,7 +908,7 @@ async def test_group_reply_to_web_mirror_db_failure_notifies_operator_and_advanc assert notice["chat_id"] == SUPPORT_CHAT_ID assert notice["reply_to_message_id"] == 210 assert "НЕ доставлен" in notice["text"] - # Апдейт частично применён в Telegram — не переигрываем, offset сдвинут и закоммичен. + # Сбой БД не переигрывается (общая политика SQLAlchemyError) — offset сдвинут и закоммичен. assert storage.get_offset() == 70 assert storage.commits == 1 @@ -931,7 +940,7 @@ async def test_group_reply_to_web_mirror_db_failure_and_notify_failure_logs_erro assert storage.web_out_messages == [] assert "потерян молча" in caplog.text assert "51" in caplog.text # thread_id узнаваем в логе - assert storage.get_offset() == 71 # апдейт частично применён — не переигрываем + assert storage.get_offset() == 71 # сбой БД не переигрывается — offset сдвинут assert storage.commits == 1 @@ -939,7 +948,14 @@ async def test_group_reply_to_own_previous_web_reply_resolves_thread() -> None: """#3471 P0 (пункт 3): реплай оператора на СВОЙ предыдущий веб-ответ (не на исходное зеркало клиента) теперь тоже резолвится — `record_outbound` сохраняет topic_message_id исходящей записи, `find_thread_by_topic_message` - больше не фильтрует по direction.""" + больше не фильтрует по direction. + + Deep review PR #3479: новая out-строка ОБЯЗАНА писаться с ТЕКУЩИМ + `support_chat_id`, а не NULL — NULL матчится `find_thread_by_topic_message` + как лениентный wildcard "любой чат" (легаси до 187/188), т.е. NULL сделал бы + КАЖДУЮ out-строку вечным wildcard-совпадением при ротации support-группы. + Эта проверка падает на дефектной реализации (support_chat_id не передавался + в `record_web_out_message`), даже когда resolve выше внешне "работает".""" calls: list[tuple[str, dict[str, Any]]] = [] client = _make_client({}, calls) storage = FakeBridgeStorage() @@ -958,12 +974,19 @@ async def test_group_reply_to_own_previous_web_reply_resolves_thread() -> None: assert len(storage.web_out_messages) == 1 assert storage.web_out_messages[0]["thread_id"] == 52 assert storage.web_out_messages[0]["topic_message_id"] == 220 + # Deep review PR #3479: НЕ NULL — иначе эта строка стала бы вечным wildcard. + assert storage.web_out_messages[0]["support_chat_id"] == SUPPORT_CHAT_ID async def test_group_reply_to_own_previous_tg_reply_resolves_target_chat() -> None: """#3471 P0 (пункт 3): та же история для Telegram-пути — реплай оператора на СВОЙ предыдущий ответ клиенту (direction='out', topic_message_id теперь - заполнен) резолвится в chat_id, а не проваливается в orphan-check.""" + заполнен) резолвится в chat_id, а не проваливается в orphan-check. + + Deep review PR #3479: новая out-строка ОБЯЗАНА писаться с ТЕКУЩИМ + `support_chat_id` (симметрично in-ветке) — иначе она стала бы вечным + wildcard в `find_chat_by_topic_message` при ротации support-группы, и + reply мог бы увести ответ ЧУЖОМУ клиенту.""" calls: list[tuple[str, dict[str, Any]]] = [] client = _make_client({"copyMessage": {"message_id": 601}}, calls) storage = FakeBridgeStorage() @@ -991,6 +1014,8 @@ async def test_group_reply_to_own_previous_tg_reply_resolves_target_chat() -> No assert calls[0][1]["chat_id"] == 555 assert len(storage.messages) == 2 # исходный 'out' + новый 'out' assert storage.messages[-1]["topic_message_id"] == 230 + # Deep review PR #3479: НЕ NULL — иначе эта строка стала бы вечным wildcard. + assert storage.messages[-1]["support_chat_id"] == SUPPORT_CHAT_ID # ── C) дедуп ──────────────────────────────────────────────────────────────────