Защита вебхука от залпа: таймауты, дедупликация, предел параллелизма, починка 400 у Yandex

Причина аварии 18.09.2026 (разобрана замерами по логам и нагрузкой):
MAX доставляет один и тот же update повторно, каждый повтор запускал новый поток
runserver, а поток держал своё соединение к MySQL на всё время запроса — отсюда
117 одновременных соединений от одного процесса и 1040 Too many connections у
общего сервера БД. Залп запускался тем, что бот не мог ответить: user_query из
callback'а с числовым data уходил в Yandex как content-число, Yandex отвечал
400 'failed to parse request JSON', сообщение оставалось в истории сессии и ломало
все последующие вызовы этой сессии.

- max_bot/max_api.py: таймаут (3.05, 15) на все 12 вызовов MAX API и обёртка
  _max_request — таймаут даёт success=False, а не исключение (иначе вебхук отвечает
  500 и провоцирует повторную доставку).
- ai_agent/api_yandex_ai.py: таймаут 20 с и max_retries=0; приведение user_query
  к строке; откат неудачной попытки из истории; окно истории 20 сообщений;
  лимит 500 сессий в памяти процесса.
- ai_agent/api_common.py: user_query приводится к строке на входе во все AI-агенты.
- max_bot/views.py: дедупликация вебхука по идентификатору update (TTL 300 с,
  cache.add) и семафор на 12 одновременных обработок — лишние получают быстрый 200.

Проверено: юнит active, авторелоад, 0 новых ERROR за 3 часа, счётчик 1040 не вырос,
пик соединений процесса упал со 117 до 2, два живых числовых запроса из прода
получили нормальные ответы вместо ошибки.
main
pilot 6 days ago
parent eded036ee2
commit 9eaaacce4d

@ -246,6 +246,11 @@ def call_ai_agent(client: Client, session_id: str, user_query: str, contact_name
if client_prompt: if client_prompt:
system_prompt += f"\n\nДополнительная информация о ресторане:\n{client_prompt}" system_prompt += f"\n\nДополнительная информация о ресторане:\n{client_prompt}"
# user_query приходит из callback'а и может быть числом (кнопки с числовым data).
# Yandex принимает content только строкой: с числом он отвечает
# 400 'failed to parse request JSON' (воспроизведено 18.09.2026).
user_query = '' if user_query is None else str(user_query)
kwargs = { kwargs = {
'session_id': session_id, 'session_id': session_id,
'system_prompt': system_prompt, 'system_prompt': system_prompt,

@ -5,6 +5,21 @@ from ai_agent.api_common import _log_dialogue, parse_ai_response, GLOBAL_ERROR_M
_contexts: Dict[str, List[Dict[str, str]]] = {} _contexts: Dict[str, List[Dict[str, str]]] = {}
# Границы памяти и размера запроса. Процесс живёт неделями, без этих пределов
# _contexts растёт бесконечно, а запрос к Yandex — вместе с историей сессии.
MAX_HISTORY_MESSAGES = 20 # system-промт + последние N сообщений диалога
MAX_SESSIONS = 500 # сколько сессий держим в памяти процесса
def _session_context(session_id: str) -> List[Dict[str, str]]:
"""История сессии; число сессий в памяти ограничено (вытесняем самые старые)."""
ctx = _contexts.get(session_id)
if ctx is None:
if len(_contexts) >= MAX_SESSIONS:
_contexts.pop(next(iter(_contexts))) # dict хранит порядок вставки
ctx = _contexts[session_id] = []
return ctx
def yandex_chat_structured( def yandex_chat_structured(
session_id: str, session_id: str,
@ -19,14 +34,20 @@ def yandex_chat_structured(
""" """
Отправляет запрос к Yandex GPT через OpenAI-совместимый API. Отправляет запрос к Yandex GPT через OpenAI-совместимый API.
""" """
if session_id not in _contexts: # user_query может прийти числом (кнопка с числовым data): Yandex принимает content
_contexts[session_id] = [ # только строкой, иначе 400 'failed to parse request JSON' (воспроизведено 18.09.2026).
{"role": "system", "content": system_prompt} user_query = '' if user_query is None else str(user_query)
]
ctx = _session_context(session_id)
if not ctx:
ctx.append({"role": "system", "content": system_prompt})
_log_dialogue(session_id, 'SYSTEM', system_prompt) _log_dialogue(session_id, 'SYSTEM', system_prompt)
_log_dialogue(session_id, 'USER', user_query) _log_dialogue(session_id, 'USER', user_query)
_contexts[session_id].append({"role": "user", "content": user_query}) ctx.append({"role": "user", "content": user_query})
# историю не копим: оставляем system-промт и последние MAX_HISTORY_MESSAGES сообщений
if len(ctx) > MAX_HISTORY_MESSAGES + 1:
del ctx[1:len(ctx) - MAX_HISTORY_MESSAGES]
try: try:
if not folder_id: if not folder_id:
@ -43,7 +64,7 @@ def yandex_chat_structured(
response = client.chat.completions.create( response = client.chat.completions.create(
model=model_uri, # <-- используем полный URI model=model_uri, # <-- используем полный URI
messages=_contexts[session_id], messages=ctx,
max_tokens=max_tokens, max_tokens=max_tokens,
temperature=temperature, temperature=temperature,
# response_format не используем # response_format не используем
@ -52,11 +73,15 @@ def yandex_chat_structured(
assistant_content = response.choices[0].message.content assistant_content = response.choices[0].message.content
_log_dialogue(session_id, 'ASSISTANT', assistant_content) _log_dialogue(session_id, 'ASSISTANT', assistant_content)
_contexts[session_id].append({"role": "assistant", "content": assistant_content}) ctx.append({"role": "assistant", "content": assistant_content})
return parse_ai_response(assistant_content) return parse_ai_response(assistant_content)
except Exception as e: except Exception as e:
# Неудачную попытку не оставляем в истории: иначе одно сообщение, на котором
# Yandex отвечает 400, ломает все последующие вызовы этой сессии.
if ctx and ctx[-1].get("role") == "user":
ctx.pop()
error_msg = GLOBAL_ERROR_MSG error_msg = GLOBAL_ERROR_MSG
_log_dialogue(session_id, 'ERROR', f"Yandex GPT error: {e}") _log_dialogue(session_id, 'ERROR', f"Yandex GPT error: {e}")
return { return {

@ -12,6 +12,26 @@ CACERT_PATH = os.path.join(BASE_DIR, 'cacert.pem')
max_verify = CACERT_PATH max_verify = CACERT_PATH
url_max = "https://platform-api2.max.ru" url_max = "https://platform-api2.max.ru"
MAX_API_TIMEOUT = (3.05, 15)
def _max_request(method: str, url: str, parse_json: bool = False, **kwargs):
"""Один вызов MAX API: с таймаутом и без исключения наружу.
Возвращает тот же словарь, что и функции модуля.
Сетевая ошибка или таймаут дают success=False, а не исключение:
иначе вебхук отвечает 500 и MAX повторяет доставку — залп растёт.
"""
kwargs.setdefault("verify", max_verify)
kwargs["timeout"] = MAX_API_TIMEOUT
try:
response = requests.request(method, url, **kwargs)
data = response.json() if parse_json else response
return {"success": True, "error": "", "data": data}
except requests.RequestException as exc:
empty = {"messages": []} if parse_json else None
return {"success": False, "error": f"max api request failed: {exc}", "data": empty}
def maxbot_send_text_message(chat_id: str, max_token: str, message: str): def maxbot_send_text_message(chat_id: str, max_token: str, message: str):
url = f"{url_max}/messages?chat_id={chat_id}" url = f"{url_max}/messages?chat_id={chat_id}"
@ -25,8 +45,7 @@ def maxbot_send_text_message(chat_id: str, max_token: str, message: str):
'Content-Type': 'application/json' 'Content-Type': 'application/json'
} }
response = requests.request("POST", url, headers=headers, data=payload, verify=max_verify) rt = _max_request("POST", url, headers=headers, data=payload)
rt = {'success': True, 'error': '', 'data': response}
return rt return rt
@ -47,8 +66,7 @@ def maxbot_send_img_message(chat_id: str, max_token: str, message: str, img: str
'Authorization': f'{max_token}', 'Authorization': f'{max_token}',
'Content-Type': 'application/json' 'Content-Type': 'application/json'
} }
response = requests.request("POST", url, headers=headers, data=payload, verify=max_verify) rt = _max_request("POST", url, headers=headers, data=payload)
rt = {'success': True, 'error': '', 'data': response}
return rt return rt
@ -94,8 +112,7 @@ def maxbot_send_menu_button(client: Client, chat_id: str, max_token: str, messag
'Content-Type': 'application/json' 'Content-Type': 'application/json'
} }
response = requests.request("POST", url, headers=headers, data=payload, verify=max_verify) rt = _max_request("POST", url, headers=headers, data=payload)
rt = {'success': True, 'error': '', 'data': response}
return rt return rt
@ -108,8 +125,7 @@ def maxbot_get_all_messages(chat_id: str, max_token: str):
'Authorization': f'{max_token}', 'Authorization': f'{max_token}',
'Content-Type': 'application/json' 'Content-Type': 'application/json'
} }
response = requests.request("GET", url, headers=headers, data=payload, verify=max_verify) rt = _max_request("GET", url, headers=headers, data=payload, parse_json=True)
rt = {'success': True, 'error': '', 'data': response.json()}
return rt return rt
@ -122,8 +138,7 @@ def maxbot_del_messages(max_token: str, message_id: str):
'Authorization': f'{max_token}', 'Authorization': f'{max_token}',
'Content-Type': 'application/json' 'Content-Type': 'application/json'
} }
response = requests.request("DELETE", url, headers=headers, data=payload, verify=max_verify) rt = _max_request("DELETE", url, headers=headers, data=payload, parse_json=True)
rt = {'success': True, 'error': '', 'data': response.json()}
return rt return rt
@ -160,8 +175,7 @@ def maxbot_get_phone(chat_id: str, max_token: str, text: str):
'Content-Type': 'application/json' 'Content-Type': 'application/json'
} }
response = requests.request("POST", url, headers=headers, data=payload, verify=max_verify) rt = _max_request("POST", url, headers=headers, data=payload)
rt = {'success': True, 'error': '', 'data': response}
return rt return rt
@ -200,8 +214,7 @@ def maxbot_send_day_button(client: Client, chat_id: str, max_token: str, message
'Content-Type': 'application/json' 'Content-Type': 'application/json'
} }
response = requests.request("POST", url, headers=headers, data=payload, verify=max_verify) rt = _max_request("POST", url, headers=headers, data=payload)
rt = {'success': True, 'error': '', 'data': response}
return rt return rt
@ -244,8 +257,7 @@ def maxbot_send_month_button(client: Client, chat_id: str, max_token: str, messa
'Content-Type': 'application/json' 'Content-Type': 'application/json'
} }
response = requests.request("POST", url, headers=headers, data=payload, verify=max_verify) rt = _max_request("POST", url, headers=headers, data=payload)
rt = {'success': True, 'error': '', 'data': response}
return rt return rt
@ -297,8 +309,7 @@ def maxbot_send_feedback_button(client: Client, chat_id: str, max_token: str, me
'Content-Type': 'application/json' 'Content-Type': 'application/json'
} }
response = requests.request("POST", url, headers=headers, data=payload, verify=max_verify) rt = _max_request("POST", url, headers=headers, data=payload)
rt = {'success': True, 'error': '', 'data': response}
return rt return rt
@ -338,8 +349,7 @@ def maxbot_send_booking_count_people_button(client: Client, chat_id: str, max_to
'Content-Type': 'application/json' 'Content-Type': 'application/json'
} }
response = requests.request("POST", url, headers=headers, data=payload, verify=max_verify) rt = _max_request("POST", url, headers=headers, data=payload)
rt = {'success': True, 'error': '', 'data': response}
return rt return rt
@ -371,8 +381,7 @@ def maxbot_send_link_button(text, title_button, link, chat_id, max_token):
'Content-Type': 'application/json' 'Content-Type': 'application/json'
} }
response = requests.request("POST", url, headers=headers, data=payload, verify=max_verify) rt = _max_request("POST", url, headers=headers, data=payload)
rt = {'success': True, 'error': '', 'data': response}
def maxbot_set_commands(max_token: str, commands: list) -> requests.Response: def maxbot_set_commands(max_token: str, commands: list) -> requests.Response:
@ -394,5 +403,10 @@ def maxbot_set_commands(max_token: str, commands: list) -> requests.Response:
'Content-Type': 'application/json' 'Content-Type': 'application/json'
} }
payload = {"commands": commands} payload = {"commands": commands}
response = requests.patch(url, headers=headers, json=payload, verify=max_verify) try:
response = requests.patch(url, headers=headers, json=payload,
verify=max_verify, timeout=MAX_API_TIMEOUT)
except requests.RequestException as exc:
print(f"max api set_commands failed: {exc}")
return None
return response return response

@ -1,5 +1,9 @@
import hashlib
import json import json
import logging
import threading
from datetime import datetime from datetime import datetime
from django.core.cache import cache
from django.core.handlers.wsgi import WSGIRequest from django.core.handlers.wsgi import WSGIRequest
from django.http import JsonResponse, HttpResponse from django.http import JsonResponse, HttpResponse
from django.views.decorators.csrf import csrf_exempt from django.views.decorators.csrf import csrf_exempt
@ -16,6 +20,93 @@ from restoran_max_bot.utils import is_json, has_key
from django.shortcuts import render from django.shortcuts import render
# Дедупликация вебхуков MAX. Мессенджер доставляет один и тот же update повторно, пока не
# получит быстрый успешный ответ. Замер 18.09.2026: 2196 доставок в одну секунду (08:09:46),
# при этом уникальных update в секунду — не больше 3, за минуту — 13, за сутки — 479.
# Каждый повтор запускал новый поток runserver (а с ним новое соединение к MySQL), новый
# вызов AI и новые вызовы MAX API. Повторы гасим до обращения к БД.
WEBHOOK_DEDUP_TTL = 300 # секунды, сколько помним уже обработанный update
def _webhook_update_key(token: str, data: dict) -> str:
"""Ключ дедупликации: client + тип update + чат + сам update.
Признак повтора — идентификатор update (`timestamp`, мс). Замер 18.09.2026: MAX
переотправляет один и тот же update (85 доставок с одним и тем же timestamp, при этом
callback_id каждый раз новый — как ключ он не годится). У нового нажатия timestamp
другой, поэтому законные повторные нажатия не глушатся.
Fallback: mid/seq сообщения, затем отпечаток тела.
"""
message = data.get("message") or {}
body = message.get("body") or {}
chat_id = (message.get("recipient") or {}).get("chat_id") or data.get("chat_id")
uid = data.get("timestamp")
if data.get("update_type") == "message_created" and (body.get("mid") or body.get("seq")):
# У сообщения свой идентификатор стабилен, добавляем его: два разных сообщения
# в одну и ту же миллисекунду не склеятся.
uid = f"{body.get('mid') or body.get('seq')}:{uid}"
if uid is None:
uid = body.get("mid") or body.get("seq")
if uid is None:
uid = hashlib.sha256(
json.dumps(data, sort_keys=True, ensure_ascii=False).encode("utf-8")
).hexdigest()
return f"max_webhook:{token}:{data.get('update_type')}:{chat_id}:{uid}"
# Предел одновременных обработок вебхука. Один обработчик = один поток runserver и одно
# соединение к MySQL на всё время запроса (замер 18.09.2026), поэтому это и есть предел
# соединений к серверам БД — он не зависит от того, сколько запросов прислал мессенджер.
WEBHOOK_MAX_CONCURRENT = 12 # одновременных обработок
WEBHOOK_QUEUE_WAIT = 2.0 # сколько секунд ждать свободный слот, потом отвечаем сразу
_webhook_slots = threading.BoundedSemaphore(WEBHOOK_MAX_CONCURRENT)
_logger = logging.getLogger('django')
def webhook_concurrency_limit(func):
"""Не пускать в обработку больше WEBHOOK_MAX_CONCURRENT вебхуков одновременно.
Лишние получают быстрый 200 (чтобы MAX не начал повторять доставку) и не занимают
поток: без этого один залп снова съест соединения общего сервера БД.
Ожидание слота не держит соединение к БД — оно открывается позже, при первом запросе.
"""
def wrapper_limit(*args, **kwargs):
if not _webhook_slots.acquire(timeout=WEBHOOK_QUEUE_WAIT):
_logger.error('webhook: нет свободных слотов (N=%s, ожидание %ss), update отброшен',
WEBHOOK_MAX_CONCURRENT, WEBHOOK_QUEUE_WAIT)
return JsonResponse({'success': True, 'error': '', 'data': 'busy'}, status=200)
try:
return func(*args, **kwargs)
finally:
_webhook_slots.release()
return wrapper_limit
def webhook_dedupe(func):
"""Не обрабатывать один и тот же update дважды.
cache.add атомарен: первый запрос помечает update, повторы получают 200 сразу,
не открывая соединение к MySQL и не вызывая AI.
"""
def wrapper_dedupe(*args, **kwargs):
request = args[0]
token = request.headers.get("X-Max-Bot-Api-Secret", "")
try:
data = json.loads(request.body.decode("utf-8"))
except (UnicodeDecodeError, json.JSONDecodeError):
# Тело не разобрать — пусть дальше отвечает обычная проверка is_json.
return func(*args, **kwargs)
key = _webhook_update_key(token, data)
if not cache.add(key, 1, WEBHOOK_DEDUP_TTL):
rt = {"success": True, "error": "", "data": "duplicate"}
return JsonResponse(rt, status=200)
return func(*args, **kwargs)
return wrapper_dedupe
def api_decorator(func): def api_decorator(func):
def wrapper_api_decorator(*args, **kwargs): def wrapper_api_decorator(*args, **kwargs):
if 'X-Max-Bot-Api-Secret' in args[0].headers: if 'X-Max-Bot-Api-Secret' in args[0].headers:
@ -61,6 +152,8 @@ def api_decorator(func):
@csrf_exempt @csrf_exempt
@webhook_dedupe
@webhook_concurrency_limit
@api_decorator @api_decorator
def api_start_max_v1(request, token, **kwargs): def api_start_max_v1(request, token, **kwargs):
data = json.loads(request.body.decode('utf-8')) data = json.loads(request.body.decode('utf-8'))

Loading…
Cancel
Save