Merge pull request 'Ответ оператора не исчезает при сбое БД, и реплай на собственный ответ снова маршрутизируется' (#3479) from fix/3471-bridge-db-failure-reply-loss into main
Some checks failed
Deploy Trade-In / perimeter-smoke (push) Blocked by required conditions
Deploy Trade-In / deploy-status (push) Blocked by required conditions
Deploy Trade-In / changes (push) Successful in 14s
Deploy Trade-In / build-frontend (push) Has been skipped
Deploy Trade-In / build-browser (push) Has been skipped
Deploy Trade-In / test (push) Successful in 6m48s
Deploy Trade-In / build-backend (push) Successful in 1m13s
Deploy Trade-In / deploy (push) Has been cancelled
Some checks failed
Deploy Trade-In / perimeter-smoke (push) Blocked by required conditions
Deploy Trade-In / deploy-status (push) Blocked by required conditions
Deploy Trade-In / changes (push) Successful in 14s
Deploy Trade-In / build-frontend (push) Has been skipped
Deploy Trade-In / build-browser (push) Has been skipped
Deploy Trade-In / test (push) Successful in 6m48s
Deploy Trade-In / build-backend (push) Successful in 1m13s
Deploy Trade-In / deploy (push) Has been cancelled
This commit is contained in:
commit
2edaae148f
3 changed files with 316 additions and 23 deletions
|
|
@ -223,7 +223,13 @@ class BridgeStorage(Protocol):
|
||||||
) -> int | None: ...
|
) -> int | None: ...
|
||||||
|
|
||||||
def record_web_out_message(
|
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,
|
||||||
|
support_chat_id: int | None = None,
|
||||||
) -> None: ...
|
) -> None: ...
|
||||||
|
|
||||||
|
|
||||||
|
|
@ -376,14 +382,20 @@ class SqlBridgeStorage:
|
||||||
со ЧУЖИМ (не NULL, не текущим) support_chat_id — исторический артефакт
|
со ЧУЖИМ (не NULL, не текущим) support_chat_id — исторический артефакт
|
||||||
ротации support-группы, не валидный маршрут сегодня. NULL (строки до
|
ротации support-группы, не валидный маршрут сегодня. NULL (строки до
|
||||||
миграции 188, если есть) — лениентный wildcard-матч (единственный
|
миграции 188, если есть) — лениентный wildcard-матч (единственный
|
||||||
действовавший чат на тот момент)."""
|
действовавший чат на тот момент).
|
||||||
|
|
||||||
|
БЕЗ фильтра по direction (#3471 P0, было `AND direction = 'in'`): с тех
|
||||||
|
пор как `record_message` на исходящем ответе тоже сохраняет
|
||||||
|
`topic_message_id` (id сообщения оператора В ТОПИКЕ), реплай оператора
|
||||||
|
на СВОЙ предыдущий ответ обязан резолвиться так же, как реплай на
|
||||||
|
зеркало клиента — иначе продолжение диалога без повторного цитирования
|
||||||
|
клиента тихо проваливалось в orphan-check."""
|
||||||
row = self._db.execute(
|
row = self._db.execute(
|
||||||
text(
|
text(
|
||||||
"""
|
"""
|
||||||
SELECT chat_id
|
SELECT chat_id
|
||||||
FROM tg_support_messages
|
FROM tg_support_messages
|
||||||
WHERE topic_message_id = CAST(:topic_message_id AS bigint)
|
WHERE topic_message_id = CAST(:topic_message_id AS bigint)
|
||||||
AND direction = 'in'
|
|
||||||
AND (support_chat_id = CAST(:support_chat_id AS bigint)
|
AND (support_chat_id = CAST(:support_chat_id AS bigint)
|
||||||
OR support_chat_id IS NULL)
|
OR support_chat_id IS NULL)
|
||||||
ORDER BY created_at DESC
|
ORDER BY created_at DESC
|
||||||
|
|
@ -413,13 +425,21 @@ class SqlBridgeStorage:
|
||||||
)
|
)
|
||||||
|
|
||||||
def record_web_out_message(
|
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,
|
||||||
|
support_chat_id: int | None = None,
|
||||||
) -> None:
|
) -> None:
|
||||||
web_support_storage.record_outbound(
|
web_support_storage.record_outbound(
|
||||||
self._db,
|
self._db,
|
||||||
thread_id=thread_id,
|
thread_id=thread_id,
|
||||||
text_body=text_body,
|
text_body=text_body,
|
||||||
operator_tg_id=operator_tg_id,
|
operator_tg_id=operator_tg_id,
|
||||||
|
topic_message_id=topic_message_id,
|
||||||
|
support_chat_id=support_chat_id,
|
||||||
)
|
)
|
||||||
|
|
||||||
|
|
||||||
|
|
@ -448,7 +468,7 @@ async def _notify_topic(
|
||||||
text: str,
|
text: str,
|
||||||
reply_to_message_id: int | None,
|
reply_to_message_id: int | None,
|
||||||
context: str,
|
context: str,
|
||||||
) -> None:
|
) -> bool:
|
||||||
"""Служебное уведомление оператору в support-топик (вторичное действие).
|
"""Служебное уведомление оператору в support-топик (вторичное действие).
|
||||||
|
|
||||||
Два свойства, которых не было у прямых `client.send_message` вызовов:
|
Два свойства, которых не было у прямых `client.send_message` вызовов:
|
||||||
|
|
@ -459,6 +479,12 @@ async def _notify_topic(
|
||||||
и не решает судьбу апдейта.
|
и не решает судьбу апдейта.
|
||||||
`context` — только технические идентификаторы (chat_id/thread_id), НЕ текст
|
`context` — только технические идентификаторы (chat_id/thread_id), НЕ текст
|
||||||
переписки: логи моста принципиально не содержат ПДн.
|
переписки: логи моста принципиально не содержат ПДн.
|
||||||
|
|
||||||
|
Возвращает True, если уведомление реально ушло, False — если само
|
||||||
|
уведомление тоже упало (напр. Telegram недоступен). Вызывающий, для
|
||||||
|
которого проваленное уведомление означает ПОЛНУЮ тишину (ни клиенту, ни
|
||||||
|
оператору), обязан на False залогировать `logger.error` с идентификаторами
|
||||||
|
(#3471 P0) — иначе единственный след остаётся только в этом WARNING.
|
||||||
"""
|
"""
|
||||||
try:
|
try:
|
||||||
await client.send_message(
|
await client.send_message(
|
||||||
|
|
@ -470,6 +496,7 @@ async def _notify_topic(
|
||||||
max_retries=_NOTIFY_SEND_MAX_RETRIES,
|
max_retries=_NOTIFY_SEND_MAX_RETRIES,
|
||||||
max_backoff=_NOTIFY_SEND_MAX_BACKOFF_S,
|
max_backoff=_NOTIFY_SEND_MAX_BACKOFF_S,
|
||||||
)
|
)
|
||||||
|
return True
|
||||||
except Exception:
|
except Exception:
|
||||||
logger.warning(
|
logger.warning(
|
||||||
"tgbot bridge: не удалось отправить уведомление оператору в топик (%s) — "
|
"tgbot bridge: не удалось отправить уведомление оператору в топик (%s) — "
|
||||||
|
|
@ -477,6 +504,7 @@ async def _notify_topic(
|
||||||
context,
|
context,
|
||||||
exc_info=True,
|
exc_info=True,
|
||||||
)
|
)
|
||||||
|
return False
|
||||||
|
|
||||||
|
|
||||||
# ── Update routing ────────────────────────────────────────────────────────────
|
# ── Update routing ────────────────────────────────────────────────────────────
|
||||||
|
|
@ -649,10 +677,22 @@ async def _handle_group_reply(
|
||||||
chat_id=target_chat_id,
|
chat_id=target_chat_id,
|
||||||
direction="out",
|
direction="out",
|
||||||
tg_message_id=tg_message_id,
|
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),
|
kind=_infer_kind(message),
|
||||||
text_body=message.get("text") or message.get("caption"),
|
text_body=message.get("text") or message.get("caption"),
|
||||||
operator_tg_id=operator_id,
|
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
|
return
|
||||||
|
|
||||||
|
|
@ -686,11 +726,64 @@ async def _handle_group_reply(
|
||||||
|
|
||||||
operator = message.get("from") or {}
|
operator = message.get("from") or {}
|
||||||
operator_id = operator.get("id")
|
operator_id = operator.get("id")
|
||||||
storage.record_web_out_message(
|
try:
|
||||||
thread_id=web_thread_id,
|
storage.record_web_out_message(
|
||||||
text_body=text_body,
|
thread_id=web_thread_id,
|
||||||
operator_tg_id=operator_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,
|
||||||
|
# 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).
|
||||||
|
# Откат без уведомления означал бы: оператор уверен, что ответил,
|
||||||
|
# клиент ждёт молча. 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 дальше не упал на
|
||||||
|
# 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
|
return
|
||||||
|
|
||||||
# Обычная болтовня в топике (реплай на чьё-то ещё сообщение) — не логируем,
|
# Обычная болтовня в топике (реплай на чьё-то ещё сообщение) — не логируем,
|
||||||
|
|
@ -741,7 +834,12 @@ async def process_update(
|
||||||
`PendingRollbackError`, `process_update` вылетит без сохранения offset'а,
|
`PendingRollbackError`, `process_update` вылетит без сохранения offset'а,
|
||||||
следующая итерация получит СТАРЫЙ offset от `get_offset()` и переиграет
|
следующая итерация получит СТАРЫЙ offset от `get_offset()` и переиграет
|
||||||
тот же апдейт — copyMessage задублирует зеркало клиента в топике на
|
тот же апдейт — copyMessage задублирует зеркало клиента в топике на
|
||||||
каждый повтор поллинга (#3 review, воспроизведено).
|
каждый повтор поллинга (#3 review, воспроизведено). Для веб-ветки
|
||||||
|
(`_handle_group_reply` → `record_web_out_message`) этот `SQLAlchemyError`
|
||||||
|
перехватывается ЛОКАЛЬНО, до этого места: там запись в БД И ЕСТЬ
|
||||||
|
доставка клиенту, поэтому rollback сопровождается уведомлением оператору
|
||||||
|
в топике, что ответ НЕ доставлен (#3471 P0) — сюда, на верхний уровень,
|
||||||
|
это исключение уже не долетает.
|
||||||
- любое прочее исключение (в т.ч. `TelegramApiError` — площадка ОТВЕТИЛА
|
- любое прочее исключение (в т.ч. `TelegramApiError` — площадка ОТВЕТИЛА
|
||||||
отказом, повтор ничего не изменит) — offset двигаем, «ядовитый» апдейт
|
отказом, повтор ничего не изменит) — offset двигаем, «ядовитый» апдейт
|
||||||
не блокирует поток.
|
не блокирует поток.
|
||||||
|
|
|
||||||
|
|
@ -106,8 +106,17 @@ def record_inbound(
|
||||||
def find_thread_by_topic_message(
|
def find_thread_by_topic_message(
|
||||||
db: Session, topic_message_id: int, support_chat_id: int
|
db: Session, topic_message_id: int, support_chat_id: int
|
||||||
) -> int | None:
|
) -> int | None:
|
||||||
"""Резолвит id зеркала (сообщения оператора reply_to) в thread_id — только
|
"""Резолвит id зеркала/ответа (reply_to) в thread_id.
|
||||||
среди direction='in' записей, зеркало-конвенция как в tg_support_messages (186).
|
|
||||||
|
БЕЗ фильтра по 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): строка со
|
Скоупим к ТЕКУЩЕМУ `support_chat_id` (#tgsupport-web review M1): строка со
|
||||||
ЧУЖИМ (не NULL и не текущим) support_chat_id — это исторический артефакт
|
ЧУЖИМ (не NULL и не текущим) support_chat_id — это исторический артефакт
|
||||||
|
|
@ -120,7 +129,6 @@ def find_thread_by_topic_message(
|
||||||
SELECT thread_id
|
SELECT thread_id
|
||||||
FROM web_support_messages
|
FROM web_support_messages
|
||||||
WHERE topic_message_id = CAST(:topic_message_id AS bigint)
|
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)
|
AND (support_chat_id = CAST(:support_chat_id AS bigint) OR support_chat_id IS NULL)
|
||||||
ORDER BY created_at DESC
|
ORDER BY created_at DESC
|
||||||
LIMIT 1
|
LIMIT 1
|
||||||
|
|
@ -132,18 +140,40 @@ def find_thread_by_topic_message(
|
||||||
|
|
||||||
|
|
||||||
def record_outbound(
|
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,
|
||||||
|
support_chat_id: int | None = None,
|
||||||
) -> int | None:
|
) -> int | None:
|
||||||
"""Записывает ответ оператора (реплай на веб-зеркало) как direction='out'.
|
"""Записывает ответ оператора (реплай на веб-зеркало) как 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-фильтр там снят).
|
||||||
|
|
||||||
|
`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(
|
row = db.execute(
|
||||||
text(
|
text(
|
||||||
"""
|
"""
|
||||||
INSERT INTO web_support_messages
|
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
|
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(:support_chat_id AS bigint),
|
||||||
CAST(:operator_tg_id AS bigint), NOW())
|
CAST(:operator_tg_id AS bigint), NOW())
|
||||||
RETURNING id
|
RETURNING id
|
||||||
"""
|
"""
|
||||||
|
|
@ -151,6 +181,8 @@ def record_outbound(
|
||||||
{
|
{
|
||||||
"thread_id": thread_id,
|
"thread_id": thread_id,
|
||||||
"text_body": text_body,
|
"text_body": text_body,
|
||||||
|
"topic_message_id": topic_message_id,
|
||||||
|
"support_chat_id": support_chat_id,
|
||||||
"operator_tg_id": operator_tg_id,
|
"operator_tg_id": operator_tg_id,
|
||||||
},
|
},
|
||||||
).fetchone()
|
).fetchone()
|
||||||
|
|
|
||||||
|
|
@ -80,6 +80,10 @@ class FakeBridgeStorage:
|
||||||
self._next_id = 1
|
self._next_id = 1
|
||||||
self.clock_s: float = 0.0
|
self.clock_s: float = 0.0
|
||||||
self.fail_next_record_message = False
|
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 ->
|
# #tgsupport-web: web_support_messages-эквивалент, topic_message_id ->
|
||||||
# (thread_id, support_chat_id) — второй элемент моделирует колонку
|
# (thread_id, support_chat_id) — второй элемент моделирует колонку
|
||||||
# web_support_messages.support_chat_id (review M1); None = легаси wildcard.
|
# 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:
|
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-скоуп (review M1): запись со ЧУЖИМ (не None, не текущим)
|
||||||
support_chat_id не матчится — None (легаси/дефолт) матчится всегда."""
|
support_chat_id не матчится — None (легаси/дефолт) матчится всегда.
|
||||||
|
БЕЗ фильтра по direction (#3471 P0) — реплай на СВОЙ предыдущий ответ
|
||||||
|
(direction='out') резолвится так же, как реплай на зеркало клиента."""
|
||||||
for m in reversed(self.messages):
|
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
|
continue
|
||||||
entry_chat_id = m.get("support_chat_id")
|
entry_chat_id = m.get("support_chat_id")
|
||||||
if entry_chat_id is not None and entry_chat_id != support_chat_id:
|
if entry_chat_id is not None and entry_chat_id != support_chat_id:
|
||||||
|
|
@ -185,15 +191,32 @@ class FakeBridgeStorage:
|
||||||
return thread_id
|
return thread_id
|
||||||
|
|
||||||
def record_web_out_message(
|
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,
|
||||||
|
support_chat_id: int | None = 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(
|
self.web_out_messages.append(
|
||||||
{
|
{
|
||||||
"thread_id": thread_id,
|
"thread_id": thread_id,
|
||||||
"text_body": text_body,
|
"text_body": text_body,
|
||||||
"operator_tg_id": operator_tg_id,
|
"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) ───────────
|
# ── httpx mocking helpers (mirrors tests/services/test_dadata.py) ───────────
|
||||||
|
|
@ -855,6 +878,146 @@ async def test_group_reply_refuses_delivery_when_both_tg_and_web_match(
|
||||||
assert storage.get_offset() == 62
|
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 всё равно сдвигается (та же
|
||||||
|
политика, что у любого другого `SQLAlchemyError` в `process_update` — сбой
|
||||||
|
БД не переигрывается, human-in-the-loop retry заменяет технический) →
|
||||||
|
исключение наружу НЕ улетает."""
|
||||||
|
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"]
|
||||||
|
# Сбой БД не переигрывается (общая политика SQLAlchemyError) — 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 # сбой БД не переигрывается — offset сдвинут
|
||||||
|
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.
|
||||||
|
|
||||||
|
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()
|
||||||
|
# Симулируем уже сохранённый предыдущий ответ оператора (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
|
||||||
|
# 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.
|
||||||
|
|
||||||
|
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()
|
||||||
|
# Предыдущий ответ оператора клиенту, зафиксированный с 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
|
||||||
|
# Deep review PR #3479: НЕ NULL — иначе эта строка стала бы вечным wildcard.
|
||||||
|
assert storage.messages[-1]["support_chat_id"] == SUPPORT_CHAT_ID
|
||||||
|
|
||||||
|
|
||||||
# ── C) дедуп ──────────────────────────────────────────────────────────────────
|
# ── C) дедуп ──────────────────────────────────────────────────────────────────
|
||||||
|
|
||||||
|
|
||||||
|
|
|
||||||
Loading…
Add table
Reference in a new issue