From 972474ad458b5c45cf277fb19b79eab25c2de65e Mon Sep 17 00:00:00 2001 From: gitadmin Date: Sat, 10 Jan 2026 15:58:18 +0300 Subject: [PATCH] =?UTF-8?q?[=D0=92=D0=9E=D0=9B=D0=9A]=20=D0=B2=D0=B2=D0=BE?= =?UTF-8?q?=D0=B4=D0=B8=D0=BC=20ANYIO=20BACKEND=20=D0=B4=D0=BB=D1=8F=20?= =?UTF-8?q?=D1=81=D0=B8=D0=BD=D1=85=D1=80=D0=BE=D0=BD=D0=BD=D0=BE=D0=B3?= =?UTF-8?q?=D0=BE=20=D0=BF=D0=BE=D1=82=D0=BE=D0=BA=D0=B0=20=D0=B7=D0=B0?= =?UTF-8?q?=D0=B3=D1=80=D1=83=D0=B7=D0=BA=D0=B8=20=D0=BA=D0=B0=D1=80=D1=82?= =?UTF-8?q?=D0=B8=D0=BD=D0=BE=D0=BA?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- db_edit.py | 18 ++++++ tg_publish.py | 160 +++++++++++++++++++++++++++++++++----------------- 2 files changed, 123 insertions(+), 55 deletions(-) diff --git a/db_edit.py b/db_edit.py index b753502..3b03ae4 100644 --- a/db_edit.py +++ b/db_edit.py @@ -1032,10 +1032,26 @@ def api_publish_post(post_id): # Функция для запуска в отдельном потоке def run_async_func(): try: + # Устанавливаем правильную политику event loop для работы с anyio + # Это критично для правильной работы httpx в отдельном потоке + if sys.platform == 'win32': + # На Windows используем ProactorEventLoop + policy = asyncio.WindowsProactorEventLoopPolicy() + else: + # На Linux используем DefaultEventLoopPolicy + policy = asyncio.DefaultEventLoopPolicy() + + asyncio.set_event_loop_policy(policy) + # Создаем новый event loop для этого потока # Это важно при работе с gunicorn, который может иметь свой event loop loop = asyncio.new_event_loop() asyncio.set_event_loop(loop) + + # Устанавливаем переменную окружения для anyio backend + import os + os.environ['ANYIO_BACKEND'] = 'asyncio' + try: result = loop.run_until_complete(publish_to_tg(post['vk_post_id'])) return result @@ -1054,6 +1070,8 @@ def api_publish_post(post_id): pass finally: loop.close() + # Сбрасываем политику event loop + asyncio.set_event_loop_policy(None) except Exception as e: import traceback log_event(f"Ошибка в run_async_func: {str(e)}\n{traceback.format_exc()}") diff --git a/tg_publish.py b/tg_publish.py index 531328a..44aa858 100644 --- a/tg_publish.py +++ b/tg_publish.py @@ -9,10 +9,22 @@ import re import httpx import traceback import io +import sys from datetime import datetime, timezone from dotenv import load_dotenv from formatter import get_event_text +# Настройка anyio для правильной работы в отдельном потоке +# Устанавливаем правильный backend для anyio (используем asyncio) +try: + import anyio + # Устанавливаем backend для anyio явно + if not hasattr(anyio, '_backend'): + # Используем asyncio backend + os.environ.setdefault('ANYIO_BACKEND', 'asyncio') +except ImportError: + pass + # Загрузка переменных окружения load_dotenv() @@ -248,68 +260,106 @@ async def publish_to_tg(vk_post_id): write_timeout = max(120.0, 60.0 + (file_size_mb * 10.0)) logger.info(f"Размер файла: {file_size_mb:.2f} МБ, таймаут на запись: {write_timeout:.1f} секунд") - # Создаем клиент с увеличенным таймаутом на запись для больших файлов - upload_timeout_config = httpx.Timeout( - connect=10.0, - read=120.0, - write=write_timeout, # Динамический таймаут на запись - pool=10.0 - ) - upload_client = httpx.AsyncClient( - http2=False, - timeout=upload_timeout_config, - follow_redirects=True - ) + # Для больших файлов используем синхронный клиент, чтобы избежать проблем с event loop + # Синхронный клиент более стабилен при работе в отдельном потоке с gunicorn + use_sync_client = file_size_mb > 0.5 # Для файлов больше 0.5 МБ используем синхронный клиент - try: - # Логируем время начала запроса - request_start_time = time.time() - logger.info(f"Начало запроса отправки фото в {datetime.now(timezone.utc).isoformat()}") + # Логируем время начала запроса + request_start_time = time.time() + logger.info(f"Начало запроса отправки фото в {datetime.now(timezone.utc).isoformat()}") + + if use_sync_client: + logger.info(f"Используем синхронный httpx клиент для файла размером {file_size_mb:.2f} МБ") - # Отправляем изображение как файл через multipart/form-data - # Создаем BytesIO объект для файла - 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() - } + # Используем синхронный httpx клиент для больших файлов + import httpx as httpx_sync + upload_timeout_config = httpx_sync.Timeout( + connect=10.0, + read=120.0, + write=write_timeout, # Динамический таймаут на запись + pool=10.0 + ) - photo_response = await upload_client.post(send_photo_url, files=files, data=data) - request_duration = time.time() - request_start_time - logger.info(f"Запрос отправки фото завершен за {request_duration:.2f} секунд") + with httpx_sync.Client( + http2=False, + timeout=upload_timeout_config, + 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) + else: + # Для маленьких файлов используем асинхронный клиент + # Создаем клиент с увеличенным таймаутом на запись для больших файлов + upload_timeout_config = httpx.Timeout( + connect=10.0, + read=120.0, + write=write_timeout, # Динамический таймаут на запись + pool=10.0 + ) + upload_client = httpx.AsyncClient( + http2=False, + timeout=upload_timeout_config, + follow_redirects=True + ) - if photo_response.status_code == 200: - photo_result = photo_response.json() - if photo_result.get('ok'): - message_data = photo_result.get('result', {}) - message_id = message_data.get('message_id') - message_has_image = True # Успешно отправили изображение - else: - # Ошибка в ответе API - логируем полную информацию - log_telegram_error(photo_response, f"[Отправка фото для VK ID {vk_post_id}]") - error_code = photo_result.get('error_code') - if error_code == 429: - retry_after = photo_result.get('parameters', {}).get('retry_after', 60) - raise RetryAfterException(retry_after) - raise Exception(f"API error: {photo_result.get('description', 'Unknown error')}") + try: + # Отправляем изображение как файл через multipart/form-data + # Создаем BytesIO объект для файла + 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 = await upload_client.post(send_photo_url, files=files, data=data) + finally: + await upload_client.aclose() + + request_duration = time.time() - request_start_time + logger.info(f"Запрос отправки фото завершен за {request_duration:.2f} секунд") + + # Обработка ответа (общая для синхронного и асинхронного клиентов) + if photo_response.status_code == 200: + photo_result = photo_response.json() + if photo_result.get('ok'): + message_data = photo_result.get('result', {}) + message_id = message_data.get('message_id') + message_has_image = True # Успешно отправили изображение else: + # Ошибка в ответе API - логируем полную информацию log_telegram_error(photo_response, f"[Отправка фото для VK ID {vk_post_id}]") - if photo_response.status_code == 429: - retry_after = 60 - try: - error_data = photo_response.json() - retry_after = error_data.get('parameters', {}).get('retry_after', 60) - except: - pass + error_code = photo_result.get('error_code') + if error_code == 429: + retry_after = photo_result.get('parameters', {}).get('retry_after', 60) raise RetryAfterException(retry_after) - raise Exception(f"HTTP {photo_response.status_code}: {photo_response.text}") - finally: - await upload_client.aclose() + raise Exception(f"API error: {photo_result.get('description', 'Unknown error')}") + else: + log_telegram_error(photo_response, f"[Отправка фото для VK ID {vk_post_id}]") + if photo_response.status_code == 429: + retry_after = 60 + try: + error_data = photo_response.json() + retry_after = error_data.get('parameters', {}).get('retry_after', 60) + except: + pass + raise RetryAfterException(retry_after) + raise Exception(f"HTTP {photo_response.status_code}: {photo_response.text}") except RetryAfterException: raise except (httpx.TimeoutException, httpx.WriteTimeout, httpx.ReadTimeout) as e: