#!/usr/bin/env python3
# -*- coding: utf-8 -*-

"""
Анализ Price-Action и Объёмов для криптовалютных инструментов.
Загружает последние свечи из таблиц вида ИНСТРУМЕНТ_TF (W1, H1, D1),
вычисляет свечные паттерны, тренд, отношение объёма к среднему,
генерирует торговые сигналы и сохраняет результаты в таблицу pa_volume_analysis.
"""

import pandas as pd
import numpy as np
from sqlalchemy import create_engine, MetaData, Table, Column, Integer, Float, String, DateTime, TIMESTAMP, text, inspect
from sqlalchemy.dialects.mysql import VARCHAR
import warnings
import argparse
from datetime import datetime
import re
import json
import sys

warnings.filterwarnings('ignore')

# ===================== ВЫВОД В STDERR =====================
def eprint(*args, **kwargs):
    print(*args, file=sys.stderr, **kwargs)

# ===================== КОНФИГУРАЦИЯ =====================
DB_CONFIG = {
    'host': 'nlbotinterface.ru',
    'port': 3306,
    'database': 'bitcoin_tickers',
    'user': 'bitcoin',
    'password': 'g49020007',
}

# Суффиксы таймфреймов для формирования имён таблиц
TIMEFRAME_SUFFIXES = ['W1', 'H1', 'D1']

LIMIT = 500                     # сколько последних свечей загружаем
VOLUME_MA_PERIOD = 20           # период для средней объема
TREND_EMA_PERIOD = 20           # период EMA для определения тренда
TREND_THRESHOLD = 0.005         # 0.5% отклонения для определения бокового тренда
VOLUME_SPIKE_THRESHOLD = 1.5    # порог всплеска объёма (текущий / средний)

# Строка подключения к БД
DATABASE_URL = f"mysql+mysqlconnector://{DB_CONFIG['user']}:{DB_CONFIG['password']}@{DB_CONFIG['host']}:{DB_CONFIG['port']}/{DB_CONFIG['database']}"
engine = create_engine(DATABASE_URL, echo=False)

# ===================== ВСПОМОГАТЕЛЬНЫЕ ФУНКЦИИ =====================
def sanitize_table_name(name):
    """Очищает имя инструмента от недопустимых символов."""
    return re.sub(r'[^a-zA-Z0-9_]', '', name)

# ===================== РАБОТА С ТАБЛИЦЕЙ РЕЗУЛЬТАТОВ =====================
def ensure_pa_volume_table_exists():
    """Создаёт таблицу pa_volume_analysis, если её нет, или приводит структуру в соответствие."""
    metadata = MetaData()
    inspector = inspect(engine)

    # Желаемая структура таблицы
    target_columns = [
        Column('id', Integer, primary_key=True, autoincrement=True),
        Column('table_name', VARCHAR(50), nullable=False),
        Column('instrument', VARCHAR(50), nullable=False),
        Column('datetime', DateTime, nullable=False),
        Column('trend', VARCHAR(10), nullable=True, comment='UP/DOWN/SIDEWAYS'),
        Column('patterns', String(500), nullable=True, comment='Список обнаруженных свечных паттернов'),
        Column('volume_ratio', Float, nullable=True, comment='Текущий объём / средний объём'),
        Column('volume_spike', String(10), nullable=True, comment='HIGH/NORMAL/LOW'),
        Column('signals', String(500), nullable=True, comment='Текстовые сигналы через разделитель'),
        Column('created_at', TIMESTAMP, server_default=text('CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP'))
    ]
    target_table = Table('pa_volume_analysis', metadata, *target_columns,
                         mysql_engine='InnoDB', mysql_default_charset='utf8mb4')

    if not inspector.has_table('pa_volume_analysis'):
        metadata.create_all(engine)
        eprint("Таблица pa_volume_analysis создана.")
        return

    # Таблица существует — проверяем структуру
    existing_columns = {col['name'] for col in inspector.get_columns('pa_volume_analysis')}
    expected_columns = {col.name for col in target_columns if col.name != 'id'}
    expected_columns.add('id')

    missing = expected_columns - existing_columns
    extra = existing_columns - expected_columns

    if missing or extra:
        eprint(f"Несоответствие структуры. Пропущенные: {missing}, лишние: {extra}. Пересоздаём таблицу.")
        with engine.connect() as conn:
            conn.execute(text("DROP TABLE IF EXISTS pa_volume_analysis"))
            conn.commit()
        metadata.create_all(engine)
        eprint("Таблица pa_volume_analysis пересоздана.")
        return

    # Проверяем уникальный ключ на table_name
    indexes = inspector.get_indexes('pa_volume_analysis')
    has_unique = any(idx['column_names'] == ['table_name'] and idx.get('unique', False) for idx in indexes)
    if not has_unique:
        with engine.connect() as conn:
            # Удаляем дубликаты, если они есть
            conn.execute(text("""
                DELETE t1 FROM pa_volume_analysis t1
                INNER JOIN pa_volume_analysis t2
                WHERE t1.id > t2.id AND t1.table_name = t2.table_name
            """))
            conn.execute(text("DELETE FROM pa_volume_analysis WHERE table_name IS NULL"))
            conn.execute(text("ALTER TABLE pa_volume_analysis ADD UNIQUE INDEX (table_name)"))
            conn.commit()
        eprint("Добавлен уникальный ключ на table_name.")

# ===================== ЗАГРУЗКА ДАННЫХ =====================
def load_last_rows(table_name, limit=LIMIT):
    """
    Загружает последние limit записей из таблицы, используя колонку timestamp (Unix time).
    Возвращает DataFrame с колонками: datetime, Open, High, Low, Close, Volume.
    """
    query = f"""
        SELECT timestamp, Open, High, Low, Close, Volume
        FROM {table_name}
        ORDER BY timestamp DESC
        LIMIT {limit}
    """
    df = pd.read_sql(query, engine)

    if df.empty:
        return df

    # Преобразование timestamp в datetime
    df['datetime'] = pd.to_datetime(df['timestamp'], unit='s')
    df.drop('timestamp', axis=1, inplace=True)
    df.sort_values('datetime', inplace=True)
    df.reset_index(drop=True, inplace=True)

    for col in ['Open', 'High', 'Low', 'Close', 'Volume']:
        df[col] = pd.to_numeric(df[col], errors='coerce')
    df.dropna(inplace=True)
    return df

# ===================== РАСЧЁТ ПОКАЗАТЕЛЕЙ PRICE-ACTION и ОБЪЁМОВ =====================
def calculate_trend(df, period=TREND_EMA_PERIOD, threshold=TREND_THRESHOLD):
    """
    Определяет направление тренда на основе EMA.
    Возвращает 'UP', 'DOWN' или 'SIDEWAYS'.
    """
    if len(df) < period:
        return None
    ema = df['Close'].ewm(span=period, adjust=False).mean()
    current_close = df['Close'].iloc[-1]
    current_ema = ema.iloc[-1]
    if pd.isna(current_ema):
        return None
    diff = (current_close - current_ema) / current_ema
    if diff > threshold:
        return 'UP'
    elif diff < -threshold:
        return 'DOWN'
    else:
        return 'SIDEWAYS'

def detect_candle_patterns(df):
    """
    Анализирует последние две свечи на наличие популярных паттернов.
    Возвращает список строк с названиями паттернов.
    """
    if len(df) < 2:
        return []

    patterns = []
    last = df.iloc[-1]
    prev = df.iloc[-2]

    # Doji (тело очень маленькое)
    body = abs(last['Close'] - last['Open'])
    candle_range = last['High'] - last['Low']
    if candle_range > 0 and body / candle_range < 0.1:
        patterns.append('DOJI')

    # Пин-бар (маленькое тело, длинная тень)
    if candle_range > 0 and body / candle_range < 0.3:
        lower_shadow = min(last['Open'], last['Close']) - last['Low']
        upper_shadow = last['High'] - max(last['Open'], last['Close'])
        # Бычий пин-бар (длинная нижняя тень)
        if lower_shadow > 2 * body and upper_shadow < body:
            patterns.append('BULLISH_PINBAR')
        # Медвежий пин-бар (длинная верхняя тень)
        if upper_shadow > 2 * body and lower_shadow < body:
            patterns.append('BEARISH_PINBAR')

    # Внутренний бар (inside bar)
    if last['High'] <= prev['High'] and last['Low'] >= prev['Low']:
        patterns.append('INSIDE_BAR')

    # Внешний бар / поглощение (engulfing)
    prev_bull = prev['Close'] > prev['Open']
    prev_bear = prev['Close'] < prev['Open']
    curr_bull = last['Close'] > last['Open']
    curr_bear = last['Close'] < last['Open']

    # Бычье поглощение
    if prev_bear and curr_bull and last['Open'] <= prev['Close'] and last['Close'] >= prev['Open']:
        patterns.append('BULLISH_ENGULFING')
    # Медвежье поглощение
    if prev_bull and curr_bear and last['Open'] >= prev['Close'] and last['Close'] <= prev['Open']:
        patterns.append('BEARISH_ENGULFING')

    return patterns

def calculate_volume_metrics(df, period=VOLUME_MA_PERIOD, spike_threshold=VOLUME_SPIKE_THRESHOLD):
    """
    Рассчитывает отношение текущего объёма к среднему и уровень всплеска.
    Возвращает (volume_ratio, volume_spike_level)
    volume_spike_level: 'HIGH' если ratio > threshold, 'LOW' если ratio < 1/threshold, иначе 'NORMAL'.
    """
    if len(df) < period + 1:
        return None, None

    volume_ma = df['Volume'].rolling(window=period).mean()
    current_volume = df['Volume'].iloc[-1]
    current_ma = volume_ma.iloc[-1]

    if pd.isna(current_ma) or current_ma == 0:
        return None, None

    ratio = current_volume / current_ma

    if ratio > spike_threshold:
        spike = 'HIGH'
    elif ratio < 1.0 / spike_threshold:
        spike = 'LOW'
    else:
        spike = 'NORMAL'

    return float(ratio), spike

# ===================== ГЕНЕРАЦИЯ СООБЩЕНИЙ =====================
def generate_pa_volume_signals(trend, patterns, volume_ratio, volume_spike):
    """
    Формирует список текстовых сигналов на основе тренда, паттернов и объёма.
    """
    messages = []

    # Тренд
    if trend:
        if trend == 'UP':
            messages.append(f"📈 Восходящий тренд (цена выше EMA{TREND_EMA_PERIOD})")
        elif trend == 'DOWN':
            messages.append(f"📉 Нисходящий тренд (цена ниже EMA{TREND_EMA_PERIOD})")
        else:
            messages.append(f"⚖️ Боковой тренд (цена около EMA{TREND_EMA_PERIOD})")

    # Паттерны
    for p in patterns:
        if p == 'DOJI':
            messages.append("🕯 Дожи — неопределённость, возможен разворот")
        elif p == 'BULLISH_PINBAR':
            messages.append("🔼 Бычий пин-бар — отскок от поддержки")
        elif p == 'BEARISH_PINBAR':
            messages.append("🔽 Медвежий пин-бар — отскок от сопротивления")
        elif p == 'INSIDE_BAR':
            messages.append("📦 Внутренний бар — сжатие, ждём пробоя")
        elif p == 'BULLISH_ENGULFING':
            messages.append("💚 Бычье поглощение — сильный сигнал покупки")
        elif p == 'BEARISH_ENGULFING':
            messages.append("❤️ Медвежье поглощение — сильный сигнал продажи")

    # Объём
    if volume_ratio is not None:
        if volume_spike == 'HIGH':
            messages.append(f"🔥 Высокий объём ({volume_ratio:.2f}x от среднего) — подтверждение движения")
        elif volume_spike == 'LOW':
            messages.append(f"💧 Низкий объём ({volume_ratio:.2f}x от среднего) — слабость движения")
        else:
            messages.append(f"📊 Объём в норме ({volume_ratio:.2f}x от среднего)")

    # Комбинированные сигналы (примеры)
    if 'BULLISH_PINBAR' in patterns and volume_spike == 'HIGH' and trend == 'UP':
        messages.append("💰 СИГНАЛ К ПОКУПКЕ: бычий пин-бар на высоком объёме в восходящем тренде")
    if 'BEARISH_PINBAR' in patterns and volume_spike == 'HIGH' and trend == 'DOWN':
        messages.append("💰 СИГНАЛ К ПРОДАЖЕ: медвежий пин-бар на высоком объёме в нисходящем тренде")
    if 'BULLISH_ENGULFING' in patterns and volume_spike == 'HIGH':
        messages.append("💰 СИГНАЛ К ПОКУПКЕ: бычье поглощение с высоким объёмом")
    if 'BEARISH_ENGULFING' in patterns and volume_spike == 'HIGH':
        messages.append("💰 СИГНАЛ К ПРОДАЖЕ: медвежье поглощение с высоким объёмом")

    return messages

# ===================== РАБОТА С ЗАПИСЯМИ В ТАБЛИЦУ РЕЗУЛЬТАТОВ =====================
def get_last_pa_volume_datetime(table_name):
    """Возвращает datetime последней обработанной свечи для данной таблицы."""
    query = text("SELECT datetime FROM pa_volume_analysis WHERE table_name = :table_name")
    with engine.connect() as conn:
        result = conn.execute(query, {"table_name": table_name}).fetchone()
        return result[0] if result else None

def update_pa_volume_analysis(table_name, instrument, dt, trend, patterns, volume_ratio, volume_spike, signals_str):
    """Вставляет или обновляет запись в таблице pa_volume_analysis."""
    delete_stmt = text("DELETE FROM pa_volume_analysis WHERE table_name = :table_name")
    insert_stmt = text("""
        INSERT INTO pa_volume_analysis (table_name, instrument, datetime, trend, patterns, volume_ratio, volume_spike, signals)
        VALUES (:table_name, :instrument, :datetime, :trend, :patterns, :volume_ratio, :volume_spike, :signals)
        ON DUPLICATE KEY UPDATE
            instrument = VALUES(instrument),
            datetime = VALUES(datetime),
            trend = VALUES(trend),
            patterns = VALUES(patterns),
            volume_ratio = VALUES(volume_ratio),
            volume_spike = VALUES(volume_spike),
            signals = VALUES(signals)
    """)

    # Преобразуем список паттернов в строку через разделитель
    patterns_str = ' | '.join(patterns) if patterns else None

    with engine.connect() as conn:
        # Если все показатели пустые – удаляем запись (чтобы не засорять)
        if (trend is None and not patterns and volume_ratio is None and volume_spike is None and not signals_str):
            conn.execute(delete_stmt, {"table_name": table_name})
        else:
            conn.execute(insert_stmt, {
                "table_name": table_name,
                "instrument": instrument,
                "datetime": dt,
                "trend": trend,
                "patterns": patterns_str,
                "volume_ratio": volume_ratio,
                "volume_spike": volume_spike,
                "signals": signals_str
            })
        conn.commit()

# ===================== ПРОВЕРКА ВРЕМЕНИ ОБРАБОТКИ =====================
def should_process_table(table_name, current_time):
    """
    Определяет, нужно ли обрабатывать таблицу в зависимости от системного времени.
    Для W1: только в понедельник с 00:00 до 00:05.
    Для H1: первые 5 минут часа.
    Для D1: первые 5 минут дня.
    """
    if table_name.endswith('_W1'):
        return current_time.weekday() == 0 and current_time.hour == 0 and current_time.minute < 5
    elif table_name.endswith('_H1'):
        return current_time.minute < 5
    elif table_name.endswith('_D1'):
        return current_time.hour == 0 and current_time.minute < 5
    else:
        return False

# ===================== ОСНОВНАЯ ПРОГРАММА =====================
def main():
    parser = argparse.ArgumentParser(description='Анализ Price-Action и Объёмов')
    parser.add_argument('--inst', required=True, help='Название инструмента (например, BTC, ETH)')
    parser.add_argument('--force', action='store_true', help='Принудительно обработать все таблицы')
    parser.add_argument('--json', action='store_true', help='Вывод результатов в формате JSON (в stdout)')
    args = parser.parse_args()

    # Очищаем имя инструмента
    instrument_raw = args.inst
    instrument_clean = sanitize_table_name(instrument_raw)
    if not instrument_clean:
        msg = "ОШИБКА: имя инструмента после очистки пустое. Используйте буквы, цифры и подчёркивание."
        eprint(msg)
        if args.json:
            result = {
                "success": False,
                "message": msg,
                "results": []
            }
            print(json.dumps(result, ensure_ascii=False, indent=2))
        return

    # Формируем список таблиц для данного инструмента
    all_possible_tables = [f"{instrument_clean}_{suffix}" for suffix in TIMEFRAME_SUFFIXES]

    # Проверяем существование таблиц в базе данных
    inspector = inspect(engine)
    existing_tables = []
    for table in all_possible_tables:
        if inspector.has_table(table):
            existing_tables.append(table)
        else:
            eprint(f"Предупреждение: таблица {table} не найдена в базе данных. Пропускаем.")

    if not existing_tables:
        msg = f"ОШИБКА: для инструмента {instrument_clean} не найдено ни одной таблицы (суффиксы {TIMEFRAME_SUFFIXES})."
        eprint(msg)
        if args.json:
            result = {
                "success": False,
                "message": msg,
                "results": []
            }
            print(json.dumps(result, ensure_ascii=False, indent=2))
        return

    ensure_pa_volume_table_exists()
    current_time = datetime.now()

    results = []

    for table in existing_tables:
        eprint(f"\n{'='*60}")
        eprint(f"Таблица: {table} (инструмент: {instrument_clean})")
        eprint('='*60)

        res = {
            "instrument": instrument_clean,
            "table_name": table,
            "datetime": None,
            "trend": None,
            "patterns": None,
            "volume_ratio": None,
            "volume_spike": None,
            "signals": None,
            "success": False,
            "message": ""
        }

        # Проверка времени обработки
        if not args.force and not should_process_table(table, current_time):
            msg = f"Пропускаем (время {current_time.strftime('%H:%M')} не подходит)."
            eprint(f"  {msg}")
            res["message"] = msg
            results.append(res)
            continue

        # Последняя обработанная свеча
        last_analyzed_dt = get_last_pa_volume_datetime(table)

        # Загружаем данные
        df = load_last_rows(table, LIMIT)

        if df.empty:
            msg = f"Нет данных."
            eprint(f"  {msg}")
            res["message"] = msg
            # Удаляем запись, если она была (данных больше нет)
            if last_analyzed_dt is not None:
                update_pa_volume_analysis(table, instrument_clean, None, None, [], None, None, None)
            results.append(res)
            continue

        last_candle_dt = df.iloc[-1]['datetime']

        # Проверка актуальности (если свеча уже обработана)
        if not args.force and last_analyzed_dt is not None and last_candle_dt == last_analyzed_dt:
            msg = f"Данные актуальны (последняя свеча {last_candle_dt})."
            eprint(f"  {msg}")
            res["message"] = msg
            results.append(res)
            continue

        # Минимально необходимое количество свечей
        min_required = max(VOLUME_MA_PERIOD, TREND_EMA_PERIOD) + 1
        if len(df) < min_required:
            msg = f"Недостаточно данных (нужно {min_required}, имеется {len(df)})."
            eprint(f"  {msg}")
            update_pa_volume_analysis(table, instrument_clean, None, None, [], None, None, None)
            res["message"] = msg
            results.append(res)
            continue

        # Расчёт метрик
        trend = calculate_trend(df, TREND_EMA_PERIOD, TREND_THRESHOLD)
        patterns = detect_candle_patterns(df)
        volume_ratio, volume_spike = calculate_volume_metrics(df, VOLUME_MA_PERIOD, VOLUME_SPIKE_THRESHOLD)

        # Вывод в stderr
        eprint(f"  Последняя свеча: {last_candle_dt}")
        eprint(f"    Тренд: {trend}")
        eprint(f"    Паттерны: {patterns if patterns else 'нет'}")
        if volume_ratio is not None:
            eprint(f"    Объём: текущий/средний = {volume_ratio:.2f} ({volume_spike})")

        # Генерация сигналов
        signals = generate_pa_volume_signals(trend, patterns, volume_ratio, volume_spike)
        signals_str = ' | '.join(signals) if signals else None
        if signals:
            eprint("\n  📋 СИГНАЛЫ:")
            for msg in signals:
                eprint(f"    {msg}")
        else:
            eprint("\n  ➖ Нет выраженных сигналов.")

        # Сохраняем результат
        update_pa_volume_analysis(table, instrument_clean, last_candle_dt, trend, patterns, volume_ratio, volume_spike, signals_str)

        # Заполняем результат для JSON
        res["datetime"] = last_candle_dt.isoformat() if last_candle_dt else None
        res["trend"] = trend
        res["patterns"] = ' | '.join(patterns) if patterns else None
        res["volume_ratio"] = volume_ratio
        res["volume_spike"] = volume_spike
        res["signals"] = signals_str
        res["success"] = True
        res["message"] = "OK"
        results.append(res)

    eprint("\n✅ Анализ Price-Action и объёмов завершён.")

    if args.json:
        output = {
            "success": True,
            "message": "OK",
            "results": results
        }
        print(json.dumps(output, ensure_ascii=False, indent=2))

if __name__ == "__main__":
    main()