#!/usr/bin/env python3
# -*- coding: utf-8 -*-

"""
Скрипт-агрегатор и распределитель данных анализов.
Запускается с параметром --server для запуска WebSocket-сервера и TCP-приёмника.
При обычном запуске читает JSON из stdin и отправляет его на TCP-порт сервера для обновления хранилища и рассылки клиентам.
"""

import asyncio
import websockets
import socket
import threading
import json
import sys
import argparse
import logging
from datetime import datetime

# Настройка логирования
logging.basicConfig(level=logging.INFO, format='%(asctime)s - %(levelname)s - %(message)s')
logger = logging.getLogger(__name__)

# Конфигурация портов
WEBSOCKET_PORT = 8765
TCP_PORT = 9050
TCP_HOST = 'localhost'

# Множество подключенных WebSocket-клиентов
connected_clients = set()

# Глобальное хранилище данных по инструментам
# Структура: {
#   "BTC": {
#       "BTC_M5": { "trend": {...}, "volatility": {...}, "candle": {...}, "levels": [...] },
#       "BTC_H1": { ... },
#       "BTC_D1": { ... }
#   },
#   ...
# }
data_store = {}

# Блокировка для доступа к хранилищу (хотя asyncio однопоточное, но на всякий случай)
store_lock = asyncio.Lock()

# ------------------------------------------------------------
# WebSocket сервер
# ------------------------------------------------------------
async def websocket_handler(websocket):
    connected_clients.add(websocket)
    client_id = id(websocket)
    logger.info(f"WebSocket клиент {client_id} подключился. Всего клиентов: {len(connected_clients)}")
    try:
        # --- ОТПРАВКА ТЕКУЩЕГО СОСТОЯНИЯ ---
        async with store_lock:
            if data_store:
                for instrument, inst_data in data_store.items():
                    payload = {
                        'instrument': instrument,
                        'data': inst_data,
                        'timestamp': datetime.now().isoformat()
                    }
                    await websocket.send(json.dumps(payload, default=str))
        # -------------------------------------
        await websocket.wait_closed()
    finally:
        connected_clients.remove(websocket)
        logger.info(f"WebSocket клиент {client_id} отключился. Осталось: {len(connected_clients)}")

async def broadcast_data(instrument_data):
    """Расслыает данные по конкретному инструменту всем подключенным WebSocket-клиентам."""
    if not connected_clients:
        return
    message = json.dumps(instrument_data, ensure_ascii=False, default=str)
    # Создаем список задач для отправки
    tasks = [asyncio.create_task(client.send(message)) for client in connected_clients]
    if tasks:
        await asyncio.gather(*tasks, return_exceptions=True)
        logger.debug(f"Отправлено сообщение {len(tasks)} клиентам")

async def run_websocket_server():
    """Запускает WebSocket-сервер."""
    async with websockets.serve(websocket_handler, "0.0.0.0", WEBSOCKET_PORT):
        logger.info(f"WebSocket сервер запущен на порту {WEBSOCKET_PORT}")
        await asyncio.Future()  # работает вечно

# ------------------------------------------------------------
# Функции обновления хранилища
# ------------------------------------------------------------
def extract_instrument_from_results(results):
    """Извлекает имя инструмента из первого успешного результата."""
    for res in results:
        if res.get('success'):
            return res.get('instrument')
    return None

async def update_store_from_json(json_data):
    """
    Обновляет глобальное хранилище данными из JSON, полученного от одного из скриптов аналитики.
    Определяет тип данных по наличию специфических полей.
    """
    global data_store
    if not isinstance(json_data, dict):
        logger.error("Получен некорректный JSON (не словарь)")
        return

    success = json_data.get('success', False)
    if not success:
        logger.warning("JSON помечен как неуспешный, пропускаем")
        return

    results = json_data.get('results', [])
    if not results:
        logger.warning("JSON не содержит результатов")
        return

    logger.info(f"Получен JSON: {json_data.get('message', '')}, results: {len(results)}")

    instrument = extract_instrument_from_results(results)
    logger.info(f"Извлечён инструмент: {instrument}")
    
    if not instrument:
        logger.error("Не удалось определить инструмент из результатов")
        return

    # Определяем тип анализа по наличию характерных полей в первом успешном результате
    sample = next((r for r in results if r.get('success')), None)
    if sample:
        logger.info(f"Sample keys: {list(sample.keys())}")
        if 'wave_count' in sample:
            logger.info("Найден ключ wave_count")
        if 'wto_value' in sample:
            logger.info("Найден ключ wto_value")
            
    if not sample:
        logger.warning("Нет успешных результатов в данных")
        return

    analysis_type = None
    if 'trend' in sample or 'trend_name' in sample:
        analysis_type = 'trend'
    elif 'atr' in sample or 'bb_width' in sample:
        analysis_type = 'volatility'
    elif 'bullish_count' in sample or 'bearish_count' in sample:
        analysis_type = 'candle'
    elif 'levels' in sample or any('type' in k for k in sample if isinstance(k, dict)):  # упрощённо
        analysis_type = 'levels'
    elif 'wave_count' in sample or 'wto_value' in sample:
        analysis_type = 'wave'
    else:
        logger.error("Не удалось определить тип анализа по полям")
        return
    
    logger.info(f"Получены данные типа '{analysis_type}' для инструмента {instrument}")

    async with store_lock:
        # Инициализируем запись для инструмента, если её нет
        if instrument not in data_store:
            data_store[instrument] = {}

        # Обрабатываем каждый таймфрейм из results
        for res in results:
            if not res.get('success'):
                continue
            table_name = res.get('table_name')
            if not table_name:
                continue

            # Инициализируем запись для таблицы, если её нет
            if table_name not in data_store[instrument]:
                data_store[instrument][table_name] = {
                    'trend': None,
                    'volatility': None,
                    'candle': None,
                    'levels': None,
                    'wave': None
                }

            # Обновляем соответствующий раздел
            if analysis_type == 'trend':
                # Сохраняем все поля, кроме служебных
                trend_data = {k: v for k, v in res.items() if k not in ['instrument', 'table_name', 'success', 'message']}
                data_store[instrument][table_name]['trend'] = trend_data
            elif analysis_type == 'volatility':
                vol_data = {k: v for k, v in res.items() if k not in ['instrument', 'table_name', 'success', 'message']}
                data_store[instrument][table_name]['volatility'] = vol_data
            elif analysis_type == 'candle':
                candle_data = {k: v for k, v in res.items() if k not in ['instrument', 'table_name', 'success', 'message']}
                # --- НОВЫЙ БЛОК: сортируем и ограничиваем историю паттернов ---
                if 'pattern_history' in candle_data and isinstance(candle_data['pattern_history'], list):
                    # Сортируем по убыванию end_time (самые свежие первыми)
                    sorted_history = sorted(
                        candle_data['pattern_history'],
                        key=lambda x: x.get('end_time', ''),
                        reverse=True
                    )
                    # Оставляем только первые 7
                    candle_data['pattern_history'] = sorted_history[:7]
                # -------------------------------------------------------------
                data_store[instrument][table_name]['candle'] = candle_data
            elif analysis_type == 'levels':
                # Для уровней сохраняем массив levels
                levels_data = res.get('levels', [])
                data_store[instrument][table_name]['levels'] = levels_data
            elif analysis_type == 'wave':
                wave_data = {k: v for k, v in res.items() if k not in ['instrument', 'table_name', 'success', 'message']}
                logger.info(f"Сохранение wave для {instrument} {table_name}: {wave_data}")
                data_store[instrument][table_name]['wave'] = wave_data

        # Формируем объект для отправки клиентам (только данные этого инструмента)
        instrument_payload = {
            'instrument': instrument,
            'data': data_store[instrument],
            'timestamp': datetime.now().isoformat()
        }
        return instrument_payload

# ------------------------------------------------------------
# TCP сервер для приёма данных от локальных отправителей
# ------------------------------------------------------------
def tcp_server(loop):
    """Запускает TCP-сервер в отдельном потоке."""
    with socket.socket(socket.AF_INET, socket.SOCK_STREAM) as s:
        s.setsockopt(socket.SOL_SOCKET, socket.SO_REUSEADDR, 1)
        s.bind((TCP_HOST, TCP_PORT))
        s.listen()
        logger.info(f"TCP сервер запущен на {TCP_HOST}:{TCP_PORT}")
        while True:
            conn, addr = s.accept()
            # Обрабатываем соединение в новом потоке
            thread = threading.Thread(target=handle_tcp_connection, args=(conn, addr, loop))
            thread.daemon = True
            thread.start()

def handle_tcp_connection(conn, addr, loop):
    """Обрабатывает входящее TCP-соединение: читает JSON, обновляет хранилище и отправляет в WebSocket."""
    logger.info(f"TCP соединение от {addr}")
    try:
        data = b""
        while True:
            chunk = conn.recv(4096)
            if not chunk:
                break
            data += chunk
        if data:
            try:
                json_data = json.loads(data.decode('utf-8'))
                logger.info(f"Получены данные от {addr}, размер: {len(data)} байт")
                # Обновляем хранилище и получаем данные для рассылки (корутина)
                future = asyncio.run_coroutine_threadsafe(update_and_broadcast(json_data), loop)
                # Можно дождаться результата, но не обязательно
                future.add_done_callback(lambda f: logger.debug("Обновление завершено"))
            except json.JSONDecodeError:
                logger.error(f"Получены некорректные JSON от {addr}: {data[:200]}")
    except Exception as e:
        logger.error(f"Ошибка обработки TCP соединения: {e}")
    finally:
        conn.close()

async def update_and_broadcast(json_data):
    """Обновляет хранилище и отправляет обновлённые данные клиентам."""
    payload = await update_store_from_json(json_data)
    if payload:
        await broadcast_data(payload)

# ------------------------------------------------------------
# Функция отправки данных на TCP-порт сервера (используется при запуске без --server)
# ------------------------------------------------------------
def send_data_via_tcp(data_json):
    """Подключается к локальному TCP-порту и отправляет JSON."""
    try:
        with socket.socket(socket.AF_INET, socket.SOCK_STREAM) as s:
            s.connect((TCP_HOST, TCP_PORT))
            s.sendall(data_json.encode('utf-8'))
            logger.info("Данные отправлены на TCP-сервер")
    except ConnectionRefusedError:
        logger.error("Не удалось подключиться к TCP-серверу. Возможно, сервер не запущен.")
        sys.exit(1)

# ------------------------------------------------------------
# Главная функция
# ------------------------------------------------------------
def main():
    parser = argparse.ArgumentParser(description='Агрегатор и распределитель данных анализов')
    parser.add_argument('--server', action='store_true', help='Запустить сервер (WebSocket + TCP)')
    args = parser.parse_args()

    if args.server:
        # Запуск сервера
        loop = asyncio.new_event_loop()
        asyncio.set_event_loop(loop)

        # Запускаем TCP сервер в отдельном потоке, передавая ему loop
        tcp_thread = threading.Thread(target=tcp_server, args=(loop,), daemon=True)
        tcp_thread.start()

        # Запускаем WebSocket сервер в основном потоке asyncio
        try:
            loop.run_until_complete(run_websocket_server())
        except KeyboardInterrupt:
            logger.info("Сервер остановлен пользователем")
        finally:
            loop.close()
    else:
        # Обычный режим: читаем JSON из stdin и отправляем на TCP-порт
        data = sys.stdin.read()
        if not data:
            logger.error("Нет данных в stdin")
            sys.exit(1)
        send_data_via_tcp(data)

if __name__ == "__main__":
    main()