import os import json import time import traceback import socket from datetime import datetime from urllib.request import urlopen from urllib.error import URLError, HTTPError from dotenv import load_dotenv import mysql.connector from mysql.connector import Error from ensure_db import ensure_database_structure # Загрузка переменных окружения load_dotenv() VOLK_CATEGORY_LIST = os.getenv('VOLK_CATEGORY_LIST') VOLK_CATEGORY_DESC_PREFIX = os.getenv('VOLK_CATEGORY_DESC_PREFIX') VOLK_EVENT_LIST_PREFIX = os.getenv('VOLK_EVENT_LIST_PREFIX') VOLK_API_VERSION = os.getenv('VOLK_API_VERSION') MDB_HOST = os.getenv('MDB_HOST') MDB_USER = os.getenv('MDB_USER') MDB_PW = os.getenv('MDB_PW') MDBASE = os.getenv('MDBASE') LOG_FILE = os.getenv('LOG_FILE') def log_message(message, category="VOLK_loader"): """Функция для записи в лог-файл с повторными попытками при блокировке""" timestamp = datetime.now().strftime("%Y-%m-%d %H:%M:%S") log_entry = f"[{timestamp}] [{category}] {message}\n" # До 5 попыток записи с увеличением задержки for attempt in range(5): try: with open(LOG_FILE, 'a', encoding='utf-8') as f: f.write(log_entry) return True except (IOError, OSError) as e: if attempt < 4: # Не последняя попытка time.sleep(0.4) # Задержка 400ms между попытками else: # После 5 неудачных попыток просто возвращаем False # Не используем print() чтобы избежать реентерабельных вызовов с gunicorn return False return False def safe_log_direct(message, category="VOLK_loader"): """Безопасная функция для прямой записи в лог-файл без обработки исключений Используется в критических ситуациях, когда обычный log_message может не сработать""" try: timestamp = datetime.now().strftime("%Y-%m-%d %H:%M:%S") log_entry = f"[{timestamp}] [{category}] {message}\n" # Прямая запись без контекстного менеджера для максимальной надежности f = open(LOG_FILE, 'a', encoding='utf-8') try: f.write(log_entry) f.flush() # Принудительно сбрасываем буфер finally: f.close() except: # Игнорируем все исключения - эта функция должна быть максимально безопасной pass def load_categories(): """Экспортируемая функция для загрузки категорий из JSON в базу данных""" inserted_count = 0 updated_count = 0 try: log_message("Начало работы") # Проверка наличия URL if not VOLK_CATEGORY_LIST: error_msg = "Переменная окружения VOLK_CATEGORY_LIST не установлена" log_message(error_msg, "ERROR") return # Чтение JSON из URL log_message(f"Чтение данных из {VOLK_CATEGORY_LIST}") with urlopen(VOLK_CATEGORY_LIST) as response: data = json.loads(response.read().decode()) # Проверка версии API api_version = data.get('api_version') if api_version != VOLK_API_VERSION: error_msg = f"Несовпадение версии API: получена версия '{api_version}', ожидается '{VOLK_API_VERSION}'. Обработка остановлена." log_message(error_msg, "ERROR") return log_message(f"Версия API соответствует: {api_version}") # Проверка наличия блока data if 'data' not in data: error_msg = "В JSON отсутствует блок 'data'" log_message(error_msg, "ERROR") return categories_data = data['data'] log_message(f"Получено {len(categories_data)} записей для обработки") # Подключение к базе данных connection = mysql.connector.connect( host=MDB_HOST, user=MDB_USER, password=MDB_PW, database=MDBASE ) # Установка временной зоны соединения в UTC cursor_temp = connection.cursor() cursor_temp.execute("SET time_zone = '+00:00'") cursor_temp.close() cursor = connection.cursor(dictionary=True) # Проверка структуры БД log_message("Проверка структуры базы данных") ensure_database_structure() # Список категорий для последующей загрузки описаний loaded_categories = [] for item in categories_data: # Проверка обязательного поля ID if 'id' not in item or not item['id']: log_message(f"Пропущена запись без id или с пустым id: {item}", "WARNING") continue # Подготовка значений category_id = item['id'] # TITLE может отсутствовать или быть пустым - в этом случае используем пустую строку title = item.get('title', '') or '' # Проверка существующей записи cursor.execute( "SELECT AUTO_ID, ID, TITLE FROM categories WHERE ID = %s", (category_id,) ) existing = cursor.fetchone() if not existing: # Вставка новой записи cursor.execute( "INSERT INTO categories (ID, TITLE, description) VALUES (%s, %s, %s)", (category_id, title, None) ) inserted_count += 1 log_message(f"Добавление: Категория {category_id}, Название: {title if title else '(пусто)'}") else: # Обновление существующей записи (всегда обновляем TITLE, даже если оно пустое) if existing['TITLE'] != title: cursor.execute( "UPDATE categories SET TITLE = %s WHERE ID = %s", (title, category_id) ) updated_count += 1 log_message(f"Обновление: Категория {category_id}, Новое название: {title if title else '(пусто)'}") else: log_message(f"Категория {category_id} без изменений") connection.commit() loaded_categories.append(category_id) log_message(f"Загружено категорий: {len(loaded_categories)}. Добавлено записей: {inserted_count}. Обновлено записей: {updated_count}.") # Загрузка описаний категорий if VOLK_CATEGORY_DESC_PREFIX and loaded_categories: log_message(f"Начало загрузки описаний для {len(loaded_categories)} категорий") desc_updated_count = 0 for category_id in loaded_categories: try: desc_url = VOLK_CATEGORY_DESC_PREFIX + category_id log_message(f"Запрос описания для категории {category_id} из {desc_url}") log_message("отладочное сообщение для проверки вывода в лог-файл") desc_response = None try: # Добавляем таймаут для избежания зависаний (30 секунд) from urllib.request import Request req = Request(desc_url) safe_log_direct("run urlopen", "DEBUG") response = urlopen(req, timeout=30) safe_log_direct("done urlopen", "ERROR") log_message("urlopen успешно выполнен") except (URLError, HTTPError, OSError, socket.gaierror, socket.timeout, ConnectionError) as e: # Используем безопасное логирование для критических ошибок safe_log_direct("ПЕРЕХВАЧЕНА СЕТЕВАЯ ОШИБКА urlopen", "ERROR") safe_log_direct(f"Категория: {category_id}, URL: {desc_url}", "ERROR") safe_log_direct(f"Тип ошибки: {type(e).__name__}, Сообщение: {str(e)}", "ERROR") try: traceback_str = traceback.format_exc() safe_log_direct(f"Трассировка:\n{traceback_str}", "ERROR") except: pass # Также пробуем обычное логирование try: error_msg = f"Ошибка сетевого подключения при запросе описания категории {category_id} из {desc_url}: {type(e).__name__}: {e}" log_message(error_msg, "ERROR") except: pass # Пропускаем эту категорию и продолжаем обработку continue except Exception as e: # Перехватываем любые другие исключения, которые могут возникнуть при urlopen safe_log_direct("ПЕРЕХВАЧЕНА НЕОЖИДАННАЯ ОШИБКА urlopen", "ERROR") safe_log_direct(f"Категория: {category_id}, URL: {desc_url}", "ERROR") safe_log_direct(f"Тип ошибки: {type(e).__name__}, Сообщение: {str(e)}", "ERROR") try: traceback_str = traceback.format_exc() safe_log_direct(f"Трассировка:\n{traceback_str}", "ERROR") except: pass # Также пробуем обычное логирование try: error_msg = f"Неожиданная ошибка при запросе описания категории {category_id} из {desc_url}: {type(e).__name__}: {e}" log_message(error_msg, "ERROR") except: pass # Пропускаем эту категорию и продолжаем обработку continue try: with response: # Явно указываем кодировку UTF-8 при декодировании raw_data = response.read() log_message("raw data read") try: decoded_data = raw_data.decode('utf-8') log_message("data decoded") except UnicodeDecodeError as e: error_msg = f"Ошибка декодирования ответа для категории {category_id}: {e}. Попытка декодирования как latin-1." log_message(error_msg, "WARNING") decoded_data = raw_data.decode('latin-1') try: desc_response = json.loads(decoded_data) log_message("json loaded") except json.JSONDecodeError as e: error_msg = f"Ошибка парсинга JSON для категории {category_id}: {e}. Первые 500 символов ответа: {decoded_data[:500]}" log_message(error_msg, "ERROR") # Пропускаем эту категорию и продолжаем обработку continue except Exception as e: error_msg = f"Ошибка при чтении ответа для категории {category_id}: {type(e).__name__}: {e}" log_message(error_msg, "ERROR") traceback_str = traceback.format_exc() log_message(f"Трассировка ошибки чтения для категории {category_id}:\n{traceback_str}", "ERROR") # Пропускаем эту категорию и продолжаем обработку continue # Проверяем, что desc_response был успешно получен if desc_response is None: error_msg = f"Не удалось получить ответ для категории {category_id}" log_message(error_msg, "ERROR") continue # Проверка версии API api_version = desc_response.get('api_version') if api_version != VOLK_API_VERSION: error_msg = f"Несовпадение версии API при загрузке описания категории {category_id}: получена версия '{api_version}', ожидается '{VOLK_API_VERSION}'. Пропуск категории." log_message(error_msg, "ERROR") # Пропускаем эту категорию и продолжаем обработку continue # Проверка наличия блока data if 'data' not in desc_response: error_msg = f"В JSON ответе для категории {category_id} отсутствует блок 'data'. Пропуск категории." log_message(error_msg, "ERROR") # Пропускаем эту категорию и продолжаем обработку continue desc_data = desc_response['data'] # Проверка типа desc_data - должен быть словарём if not isinstance(desc_data, dict): error_msg = f"Блок 'data' для категории {category_id} не является словарём (тип: {type(desc_data)}). Пропуск категории." log_message(error_msg, "ERROR") # Пропускаем эту категорию и продолжаем обработку continue # Получаем description из блока data description = desc_data.get('description') # Если description пустое, None или отсутствует, обнуляем поле в БД if not description: description = None # Обновляем описание в базе данных cursor.execute( "UPDATE categories SET description = %s WHERE ID = %s", (description, category_id) ) connection.commit() desc_updated_count += 1 if description: log_message(f"Обновлено описание для категории {category_id}") else: log_message(f"Описание для категории {category_id} обнулено (пустое или отсутствует)") except Exception as e: error_msg = f"Ошибка при загрузке описания для категории {category_id}: {type(e).__name__}: {e}" log_message(error_msg, "ERROR") traceback_str = traceback.format_exc() log_message(f"Трассировка ошибки для категории {category_id}:\n{traceback_str}", "ERROR") # Пропускаем эту категорию и продолжаем обработку остальных continue log_message(f"Завершение загрузки описаний. Обновлено описаний: {desc_updated_count}.") log_message(f"Завершение работы. Добавлено записей: {inserted_count}. Обновлено записей: {updated_count}.") except Error as e: error_msg = f"Ошибка базы данных: {e}" log_message(error_msg, "ERROR") if 'connection' in locals() and connection.is_connected(): connection.rollback() except Exception as e: error_msg = f"Общая ошибка: {e}" log_message(error_msg, "ERROR") finally: if 'connection' in locals() and connection.is_connected(): cursor.close() connection.close() def load_events_by_category(category_id): """Экспортируемая функция для загрузки событий категории из JSON в базу данных""" inserted_count = 0 updated_count = 0 try: log_message(f"Начало загрузки событий для категории {category_id}") # Проверка наличия URL префикса if not VOLK_EVENT_LIST_PREFIX: error_msg = "Переменная окружения VOLK_EVENT_LIST_PREFIX не установлена" log_message(error_msg, "ERROR") return # Проверка наличия category_id if not category_id: error_msg = "ID категории не указан" log_message(error_msg, "ERROR") return # Формирование URL для запроса event_url = VOLK_EVENT_LIST_PREFIX + category_id log_message(f"Чтение данных из {event_url}") # Чтение JSON из URL with urlopen(event_url) as response: data = json.loads(response.read().decode()) # Проверка версии API api_version = data.get('api_version') if api_version != VOLK_API_VERSION: error_msg = f"Несовпадение версии API: получена версия '{api_version}', ожидается '{VOLK_API_VERSION}'. Обработка остановлена." log_message(error_msg, "ERROR") return log_message(f"Версия API соответствует: {api_version}") # Проверка наличия блока data if 'data' not in data: error_msg = "В JSON отсутствует блок 'data'" log_message(error_msg, "ERROR") return events_data = data['data'] log_message(f"Получено {len(events_data)} событий для обработки") # Подключение к базе данных connection = mysql.connector.connect( host=MDB_HOST, user=MDB_USER, password=MDB_PW, database=MDBASE ) # Установка временной зоны соединения в UTC cursor_temp = connection.cursor() cursor_temp.execute("SET time_zone = '+00:00'") cursor_temp.close() cursor = connection.cursor(dictionary=True) # Проверка структуры БД log_message("Проверка структуры базы данных") ensure_database_structure() # Получаем максимальный number для всей таблицы, чтобы продолжить нумерацию cursor.execute( "SELECT MAX(number) as max_number FROM events" ) result = cursor.fetchone() max_number = result['max_number'] or 0 start_number = max_number + 1 # Обработка событий current_number = start_number for item in events_data: # Проверка обязательного поля id if 'id' not in item or item['id'] is None: log_message(f"Пропущено событие без id: {item}", "WARNING") continue # Проверка обязательного поля title if 'title' not in item or not item['title']: log_message(f"Пропущено событие без title (id: {item.get('id', 'N/A')}): {item}", "WARNING") continue # Подготовка значений с маппингом полей # Преобразуем id в число, если это необходимо try: id_event = int(item['id']) # Уникальный идентификатор события из JSON except (ValueError, TypeError): log_message(f"Пропущено событие с некорректным id (не число): {item.get('id', 'N/A')}", "WARNING") continue name = item.get('title', '') about = item.get('announce', '') or '' authors = item.get('authors', '') or '' about_social_picture = item.get('image', '') or '' unit_name = category_id # Проверка существующей записи по id_event (уникальный идентификатор) cursor.execute( "SELECT id, id_event, name, authors, about, about_social_picture FROM events WHERE id_event = %s", (id_event,) ) existing = cursor.fetchone() if not existing: # Вставка новой записи # Используем id из JSON для id_event, а для number используем автоинкремент cursor.execute( """INSERT INTO events (id_event, number, unit_name, name, about, about_social_picture, authors) VALUES (%s, %s, %s, %s, %s, %s, %s)""", (id_event, current_number, unit_name, name, about, about_social_picture, authors) ) inserted_count += 1 log_message(f"Добавление: Событие '{name}', id_event: {id_event}, номер: {current_number}, категория: {category_id}") current_number += 1 else: # Обновление существующей записи needs_update = False update_fields = [] if existing['name'] != name: needs_update = True update_fields.append('name') if existing['authors'] != authors: needs_update = True update_fields.append('authors') if existing['about'] != about: needs_update = True update_fields.append('about') if existing['about_social_picture'] != about_social_picture: needs_update = True update_fields.append('about_social_picture') if needs_update: cursor.execute( """UPDATE events SET name = %s, authors = %s, about = %s, about_social_picture = %s WHERE id_event = %s""", (name, authors, about, about_social_picture, id_event) ) updated_count += 1 log_message(f"Обновление: Событие id_event={id_event} '{name}', обновлены поля: {', '.join(update_fields)}") else: log_message(f"Событие id_event={id_event} '{name}' без изменений") connection.commit() log_message(f"Завершение работы. Добавлено событий: {inserted_count}. Обновлено событий: {updated_count}.") except Error as e: error_msg = f"Ошибка базы данных: {e}" log_message(error_msg, "ERROR") if 'connection' in locals() and connection.is_connected(): connection.rollback() except Exception as e: error_msg = f"Общая ошибка: {e}" log_message(error_msg, "ERROR") finally: if 'connection' in locals() and connection.is_connected(): cursor.close() connection.close() def load_all_events(): """Экспортируемая функция для загрузки всех событий по всем категориям""" try: log_message("Начало загрузки всех событий по всем категориям") # Подключение к базе данных connection = mysql.connector.connect( host=MDB_HOST, user=MDB_USER, password=MDB_PW, database=MDBASE ) cursor = connection.cursor(dictionary=True) # Получаем список всех категорий cursor.execute("SELECT ID FROM categories ORDER BY ID") categories = cursor.fetchall() if not categories: log_message("В базе данных нет категорий для загрузки событий", "WARNING") cursor.close() connection.close() return log_message(f"Найдено {len(categories)} категорий для обработки") # Закрываем соединение перед вызовами load_events_by_category # (так как она создает свое соединение) cursor.close() connection.close() # Загружаем события для каждой категории processed_count = 0 for category in categories: category_id = category['ID'] try: log_message(f"Обработка категории {category_id} ({processed_count + 1}/{len(categories)})") load_events_by_category(category_id) processed_count += 1 except Exception as e: error_msg = f"Ошибка при загрузке событий для категории {category_id}: {e}" log_message(error_msg, "ERROR") # Продолжаем обработку других категорий continue log_message(f"Завершение загрузки всех событий. Обработано категорий: {processed_count}/{len(categories)}") except Error as e: error_msg = f"Ошибка базы данных при загрузке всех событий: {e}" log_message(error_msg, "ERROR") except Exception as e: error_msg = f"Общая ошибка при загрузке всех событий: {e}" log_message(error_msg, "ERROR") if __name__ == "__main__": load_categories()