This commit is contained in:
+123
-55
@@ -7,6 +7,7 @@ import html
|
|||||||
import json
|
import json
|
||||||
import re
|
import re
|
||||||
import httpx
|
import httpx
|
||||||
|
import requests
|
||||||
import traceback
|
import traceback
|
||||||
import io
|
import io
|
||||||
import sys
|
import sys
|
||||||
@@ -103,16 +104,21 @@ def has_org_tag(text):
|
|||||||
return bool(re.search(pattern, text, re.IGNORECASE))
|
return bool(re.search(pattern, text, re.IGNORECASE))
|
||||||
|
|
||||||
def log_telegram_error(response, context=""):
|
def log_telegram_error(response, context=""):
|
||||||
"""Логирует полную информацию об ошибке от Telegram API"""
|
"""Логирует полную информацию об ошибке от Telegram API
|
||||||
|
Работает как с httpx.Response, так и с requests.Response
|
||||||
|
"""
|
||||||
try:
|
try:
|
||||||
if response is not None:
|
if response is not None:
|
||||||
logger.error(f"{context} Статус код: {response.status_code}")
|
logger.error(f"{context} Статус код: {response.status_code}")
|
||||||
logger.error(f"{context} Заголовки ответа: {dict(response.headers)}")
|
logger.error(f"{context} Заголовки ответа: {dict(response.headers)}")
|
||||||
try:
|
try:
|
||||||
|
# Работает для обоих типов ответов (httpx и requests)
|
||||||
error_data = response.json()
|
error_data = response.json()
|
||||||
logger.error(f"{context} Полный ответ от Telegram API: {json.dumps(error_data, ensure_ascii=False, indent=2)}")
|
logger.error(f"{context} Полный ответ от Telegram API: {json.dumps(error_data, ensure_ascii=False, indent=2)}")
|
||||||
except:
|
except:
|
||||||
logger.error(f"{context} Текст ответа: {response.text}")
|
# Для requests используем text, для httpx тоже text
|
||||||
|
response_text = getattr(response, 'text', str(response.content[:500]) if hasattr(response, 'content') else 'N/A')
|
||||||
|
logger.error(f"{context} Текст ответа: {response_text}")
|
||||||
else:
|
else:
|
||||||
logger.error(f"{context} Ответ от Telegram API отсутствует (None)")
|
logger.error(f"{context} Ответ от Telegram API отсутствует (None)")
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
@@ -136,6 +142,57 @@ async def download_image(image_url, httpx_client):
|
|||||||
logger.error(f"Ошибка при скачивании изображения {image_url}: {e}")
|
logger.error(f"Ошибка при скачивании изображения {image_url}: {e}")
|
||||||
raise
|
raise
|
||||||
|
|
||||||
|
def _send_photo_sync(send_photo_url, image_data, caption, chat_id, disable_notification, timeout):
|
||||||
|
"""
|
||||||
|
Синхронная функция для отправки фото через requests.
|
||||||
|
Вызывается из отдельного потока для избежания проблем с event loop в gunicorn.
|
||||||
|
|
||||||
|
Args:
|
||||||
|
send_photo_url: URL для отправки фото в Telegram API
|
||||||
|
image_data: Байты изображения
|
||||||
|
caption: Подпись к фото
|
||||||
|
chat_id: ID чата/канала
|
||||||
|
disable_notification: Отключить уведомления
|
||||||
|
timeout: Кортеж (connect_timeout, read_timeout) или число для общего таймаута
|
||||||
|
|
||||||
|
Returns:
|
||||||
|
requests.Response объект
|
||||||
|
"""
|
||||||
|
try:
|
||||||
|
# Используем requests для стабильной работы в gunicorn
|
||||||
|
# requests более надежен для синхронных операций в контексте WSGI
|
||||||
|
image_file = io.BytesIO(image_data)
|
||||||
|
files = {
|
||||||
|
'photo': ('image.jpg', image_file, 'image/jpeg')
|
||||||
|
}
|
||||||
|
data = {
|
||||||
|
'chat_id': chat_id,
|
||||||
|
'caption': caption,
|
||||||
|
'parse_mode': 'HTML',
|
||||||
|
'disable_notification': str(disable_notification).lower()
|
||||||
|
}
|
||||||
|
|
||||||
|
# Отправляем запрос с увеличенным таймаутом
|
||||||
|
# requests.post принимает timeout как кортеж (connect, read) или число
|
||||||
|
response = requests.post(
|
||||||
|
send_photo_url,
|
||||||
|
files=files,
|
||||||
|
data=data,
|
||||||
|
timeout=timeout,
|
||||||
|
stream=False # Отключаем потоковую передачу для стабильности
|
||||||
|
)
|
||||||
|
|
||||||
|
return response
|
||||||
|
except requests.exceptions.Timeout as e:
|
||||||
|
# Пробрасываем таймаут как есть, чтобы его можно было обработать выше
|
||||||
|
raise
|
||||||
|
except requests.exceptions.RequestException as e:
|
||||||
|
# Пробрасываем другие ошибки requests как есть
|
||||||
|
raise
|
||||||
|
except Exception as e:
|
||||||
|
# Пробрасываем остальные исключения как есть
|
||||||
|
raise
|
||||||
|
|
||||||
async def publish_to_tg(vk_post_id):
|
async def publish_to_tg(vk_post_id):
|
||||||
"""Публикация одной записи по VK post ID"""
|
"""Публикация одной записи по VK post ID"""
|
||||||
logger.info(f"Запуск публикации для записи VK ID {vk_post_id}")
|
logger.info(f"Запуск публикации для записи VK ID {vk_post_id}")
|
||||||
@@ -260,79 +317,89 @@ async def publish_to_tg(vk_post_id):
|
|||||||
write_timeout = max(120.0, 60.0 + (file_size_mb * 10.0))
|
write_timeout = max(120.0, 60.0 + (file_size_mb * 10.0))
|
||||||
logger.info(f"Размер файла: {file_size_mb:.2f} МБ, таймаут на запись: {write_timeout:.1f} секунд")
|
logger.info(f"Размер файла: {file_size_mb:.2f} МБ, таймаут на запись: {write_timeout:.1f} секунд")
|
||||||
|
|
||||||
# Используем синхронный httpx клиент для отправки файлов
|
# Используем requests для отправки файлов через отдельный поток
|
||||||
# Это необходимо для стабильной работы с gunicorn, который использует синхронные воркеры
|
# Это необходимо для стабильной работы с gunicorn, который использует синхронные воркеры
|
||||||
# Синхронный клиент не зависит от event loop и работает стабильно в отдельном потоке
|
# requests более надежен для синхронных операций в контексте WSGI
|
||||||
logger.info(f"Используем синхронный httpx клиент для файла размером {file_size_mb:.2f} МБ")
|
logger.info(f"Используем requests для отправки файла размером {file_size_mb:.2f} МБ через отдельный поток")
|
||||||
|
|
||||||
# Логируем время начала запроса
|
# Логируем время начала запроса
|
||||||
request_start_time = time.time()
|
request_start_time = time.time()
|
||||||
logger.info(f"Начало запроса отправки фото в {datetime.now(timezone.utc).isoformat()}")
|
logger.info(f"Начало запроса отправки фото в {datetime.now(timezone.utc).isoformat()}")
|
||||||
|
|
||||||
# Используем синхронный httpx клиент
|
# Вызываем синхронную функцию через asyncio.to_thread()
|
||||||
import httpx as httpx_sync
|
# Это позволяет избежать проблем с event loop в gunicorn
|
||||||
upload_timeout_config = httpx_sync.Timeout(
|
|
||||||
connect=10.0,
|
|
||||||
read=120.0,
|
|
||||||
write=write_timeout, # Динамический таймаут на запись
|
|
||||||
pool=10.0
|
|
||||||
)
|
|
||||||
|
|
||||||
try:
|
try:
|
||||||
with httpx_sync.Client(
|
# Используем кортеж для таймаута: (connect, read)
|
||||||
http2=False,
|
# requests использует один таймаут для всех операций, поэтому используем максимальный
|
||||||
timeout=upload_timeout_config,
|
timeout_tuple = (10.0, write_timeout) # (connect, read/write)
|
||||||
follow_redirects=True
|
|
||||||
) as sync_client:
|
|
||||||
# Отправляем изображение как файл через multipart/form-data
|
|
||||||
image_file = io.BytesIO(image_data)
|
|
||||||
files = {
|
|
||||||
'photo': ('image.jpg', image_file, 'image/jpeg')
|
|
||||||
}
|
|
||||||
data = {
|
|
||||||
'chat_id': CHANNEL_ID,
|
|
||||||
'caption': text,
|
|
||||||
'parse_mode': 'HTML',
|
|
||||||
'disable_notification': str(PUBLISH_SILENTLY).lower()
|
|
||||||
}
|
|
||||||
|
|
||||||
photo_response = sync_client.post(send_photo_url, files=files, data=data)
|
photo_response = await asyncio.to_thread(
|
||||||
except (httpx_sync.TimeoutException, httpx_sync.WriteTimeout, httpx_sync.ReadTimeout) as sync_timeout:
|
_send_photo_sync,
|
||||||
# Обработка таймаутов от синхронного клиента
|
send_photo_url,
|
||||||
timeout_type = type(sync_timeout).__name__
|
image_data,
|
||||||
|
text,
|
||||||
|
CHANNEL_ID,
|
||||||
|
PUBLISH_SILENTLY,
|
||||||
|
timeout_tuple
|
||||||
|
)
|
||||||
|
except Exception as sync_error:
|
||||||
|
# Обработка ошибок от синхронной функции
|
||||||
request_duration = time.time() - request_start_time
|
request_duration = time.time() - request_start_time
|
||||||
logger.error(f"ТАЙМАУТ ({timeout_type}) при отправке изображения (синхронный клиент) для записи VK ID {vk_post_id}")
|
error_msg = str(sync_error)
|
||||||
logger.error(f"Время до таймаута: {request_duration:.2f} секунд")
|
|
||||||
logger.error(f"Тип исключения: {timeout_type}")
|
# Проверяем, является ли это таймаутом
|
||||||
logger.error(f"Сообщение об ошибке: {str(sync_timeout)}")
|
is_timeout = (
|
||||||
logger.error(f"Полная информация об исключении: {repr(sync_timeout)}")
|
'Timeout' in error_msg or
|
||||||
logger.error(f"URL изображения: {image_url}")
|
isinstance(sync_error, requests.exceptions.Timeout) or
|
||||||
logger.error(f"Размер файла: {file_size_mb:.2f} МБ")
|
'timeout' in error_msg.lower()
|
||||||
if hasattr(sync_timeout, 'request'):
|
)
|
||||||
logger.error(f"Запрос, вызвавший таймаут: {sync_timeout.request.method} {sync_timeout.request.url if hasattr(sync_timeout.request, 'url') else 'N/A'}")
|
|
||||||
logger.error(f"Трассировка стека:\n{traceback.format_exc()}")
|
if is_timeout:
|
||||||
raise # Пробрасываем таймаут наверх
|
timeout_type = type(sync_error).__name__
|
||||||
|
logger.error(f"ТАЙМАУТ ({timeout_type}) при отправке изображения (requests) для записи VK ID {vk_post_id}")
|
||||||
|
logger.error(f"Время до таймаута: {request_duration:.2f} секунд")
|
||||||
|
logger.error(f"Тип исключения: {timeout_type}")
|
||||||
|
logger.error(f"Сообщение об ошибке: {error_msg}")
|
||||||
|
logger.error(f"Полная информация об исключении: {repr(sync_error)}")
|
||||||
|
logger.error(f"URL изображения: {image_url}")
|
||||||
|
logger.error(f"Размер файла: {file_size_mb:.2f} МБ")
|
||||||
|
logger.error(f"Запрос, вызвавший таймаут: POST {send_photo_url}")
|
||||||
|
logger.error(f"Трассировка стека:\n{traceback.format_exc()}")
|
||||||
|
raise # Пробрасываем таймаут наверх
|
||||||
|
else:
|
||||||
|
# Другие ошибки
|
||||||
|
logger.error(f"Ошибка при отправке изображения (requests) для записи VK ID {vk_post_id}: {error_msg}")
|
||||||
|
logger.error(f"Трассировка стека:\n{traceback.format_exc()}")
|
||||||
|
raise
|
||||||
|
|
||||||
request_duration = time.time() - request_start_time
|
request_duration = time.time() - request_start_time
|
||||||
logger.info(f"Запрос отправки фото завершен за {request_duration:.2f} секунд")
|
logger.info(f"Запрос отправки фото завершен за {request_duration:.2f} секунд")
|
||||||
|
|
||||||
# Обработка ответа (общая для синхронного и асинхронного клиентов)
|
# Обработка ответа от requests
|
||||||
if photo_response.status_code == 200:
|
if photo_response.status_code == 200:
|
||||||
photo_result = photo_response.json()
|
try:
|
||||||
|
photo_result = photo_response.json()
|
||||||
|
except Exception as json_error:
|
||||||
|
logger.error(f"Ошибка при парсинге JSON ответа для записи VK ID {vk_post_id}: {json_error}")
|
||||||
|
logger.error(f"Текст ответа: {photo_response.text[:500]}")
|
||||||
|
raise Exception(f"Ошибка парсинга JSON ответа: {json_error}")
|
||||||
|
|
||||||
if photo_result.get('ok'):
|
if photo_result.get('ok'):
|
||||||
message_data = photo_result.get('result', {})
|
message_data = photo_result.get('result', {})
|
||||||
message_id = message_data.get('message_id')
|
message_id = message_data.get('message_id')
|
||||||
message_has_image = True # Успешно отправили изображение
|
message_has_image = True # Успешно отправили изображение
|
||||||
else:
|
else:
|
||||||
# Ошибка в ответе API - логируем полную информацию
|
# Ошибка в ответе API - логируем полную информацию
|
||||||
log_telegram_error(photo_response, f"[Отправка фото для VK ID {vk_post_id}]")
|
logger.error(f"[Отправка фото для VK ID {vk_post_id}] Статус код: {photo_response.status_code}")
|
||||||
|
logger.error(f"[Отправка фото для VK ID {vk_post_id}] Полный ответ от Telegram API: {json.dumps(photo_result, ensure_ascii=False, indent=2)}")
|
||||||
error_code = photo_result.get('error_code')
|
error_code = photo_result.get('error_code')
|
||||||
if error_code == 429:
|
if error_code == 429:
|
||||||
retry_after = photo_result.get('parameters', {}).get('retry_after', 60)
|
retry_after = photo_result.get('parameters', {}).get('retry_after', 60)
|
||||||
raise RetryAfterException(retry_after)
|
raise RetryAfterException(retry_after)
|
||||||
raise Exception(f"API error: {photo_result.get('description', 'Unknown error')}")
|
raise Exception(f"API error: {photo_result.get('description', 'Unknown error')}")
|
||||||
else:
|
else:
|
||||||
log_telegram_error(photo_response, f"[Отправка фото для VK ID {vk_post_id}]")
|
logger.error(f"[Отправка фото для VK ID {vk_post_id}] Статус код: {photo_response.status_code}")
|
||||||
|
logger.error(f"[Отправка фото для VK ID {vk_post_id}] Текст ответа: {photo_response.text[:500]}")
|
||||||
if photo_response.status_code == 429:
|
if photo_response.status_code == 429:
|
||||||
retry_after = 60
|
retry_after = 60
|
||||||
try:
|
try:
|
||||||
@@ -341,15 +408,17 @@ async def publish_to_tg(vk_post_id):
|
|||||||
except:
|
except:
|
||||||
pass
|
pass
|
||||||
raise RetryAfterException(retry_after)
|
raise RetryAfterException(retry_after)
|
||||||
raise Exception(f"HTTP {photo_response.status_code}: {photo_response.text}")
|
raise Exception(f"HTTP {photo_response.status_code}: {photo_response.text[:500]}")
|
||||||
except RetryAfterException:
|
except RetryAfterException:
|
||||||
raise
|
raise
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
# Проверяем, является ли это таймаутом (может быть от синхронного или асинхронного клиента)
|
# Проверяем, является ли это таймаутом
|
||||||
import httpx as httpx_sync
|
# Может быть от requests или httpx
|
||||||
is_timeout = (
|
is_timeout = (
|
||||||
isinstance(e, (httpx.TimeoutException, httpx.WriteTimeout, httpx.ReadTimeout)) or
|
isinstance(e, (httpx.TimeoutException, httpx.WriteTimeout, httpx.ReadTimeout)) or
|
||||||
isinstance(e, (httpx_sync.TimeoutException, httpx_sync.WriteTimeout, httpx_sync.ReadTimeout))
|
isinstance(e, requests.exceptions.Timeout) or
|
||||||
|
'Timeout' in str(e) or
|
||||||
|
'timeout' in str(e).lower()
|
||||||
)
|
)
|
||||||
|
|
||||||
if is_timeout:
|
if is_timeout:
|
||||||
@@ -364,8 +433,7 @@ async def publish_to_tg(vk_post_id):
|
|||||||
logger.error(f"URL изображения: {image_url}")
|
logger.error(f"URL изображения: {image_url}")
|
||||||
if image_data is not None:
|
if image_data is not None:
|
||||||
logger.error(f"Размер файла: {len(image_data) / (1024 * 1024):.2f} МБ")
|
logger.error(f"Размер файла: {len(image_data) / (1024 * 1024):.2f} МБ")
|
||||||
if hasattr(e, 'request'):
|
logger.error(f"Запрос, вызвавший таймаут: POST {send_photo_url}")
|
||||||
logger.error(f"Запрос, вызвавший таймаут: {e.request.method} {e.request.url if hasattr(e.request, 'url') else 'N/A'}")
|
|
||||||
logger.error(f"Трассировка стека:\n{traceback.format_exc()}")
|
logger.error(f"Трассировка стека:\n{traceback.format_exc()}")
|
||||||
raise # Пробрасываем таймаут наверх
|
raise # Пробрасываем таймаут наверх
|
||||||
|
|
||||||
|
|||||||
Reference in New Issue
Block a user