import pandas as pd
import numpy as np
from sqlalchemy import create_engine, MetaData, Table, Column, Integer, String, DateTime, Float, text, inspect
from sqlalchemy.dialects.mysql import TINYINT, VARCHAR, DATETIME, FLOAT as mysqlFloat, TIMESTAMP as MySQL_TIMESTAMP
from sqlalchemy.types import Float, TIMESTAMP
from sqlalchemy.exc import SQLAlchemyError
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',
}

# ИЗМЕНЕНО: вместо M5 теперь W1
TIMEFRAME_SUFFIXES = ['W1', 'H1', 'D1']
LIMIT = 6000
LEVEL_CLUSTER_THRESHOLD = 0.0005   # 0.5% от последней цены

# Настройки алгоритма (окна для разных таймфреймов)
WINDOW_BY_SUFFIX = {
    'W1': 5,   # <-- ДОБАВЛЕНО
    'H1': 5,
    'D1': 5
}
MIN_STRENGTH = 5.0                # минимальная сила уровня (может быть дробной)

# Флаги функциональности
INCLUDE_PIVOT = True               # добавлять ли уровни Pivot Points
USE_VOLUME = False                  # учитывать ли объёмы (пока нет, но код готов)

# ------------------------------------------------------------
# Подключение к БД
# ------------------------------------------------------------
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)

# ------------------------------------------------------------
# Создание / проверка таблицы levels (с колонкой instrument)
# ------------------------------------------------------------
def ensure_levels_table_exists():
    metadata = MetaData()
    inspector = inspect(engine)

    # Желаемая структура (добавлена колонка instrument)
    target_columns = [
        Column('id', Integer, primary_key=True, autoincrement=True),
        Column('table_name', VARCHAR(50), nullable=False),
        Column('instrument', VARCHAR(50), nullable=False),          # новая колонка
        Column('level_type', VARCHAR(20), nullable=False),
        Column('price', Float, nullable=False),
        Column('strength', Float, nullable=False, default=1.0),
        Column('last_updated', DateTime, nullable=False),
        Column('created_at', TIMESTAMP, server_default=text('CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP'))
    ]
    target_table = Table('levels', metadata, *target_columns,
                         mysql_engine='InnoDB',
                         mysql_default_charset='utf8mb4')

    if not inspector.has_table('levels'):
        metadata.create_all(engine)
        eprint("Таблица levels создана (с колонкой instrument).")
        return

    # Проверка структуры существующей таблицы
    existing_columns = {col['name']: col for col in inspector.get_columns('levels')}
    expected_columns = {col.name: col for col in target_columns}

    missing_columns = set(expected_columns.keys()) - set(existing_columns.keys())
    extra_columns = set(existing_columns.keys()) - set(expected_columns.keys())
    type_mismatch = False
    for col_name, expected_col in expected_columns.items():
        if col_name in existing_columns:
            if str(expected_col.type) != str(existing_columns[col_name]['type']):
                type_mismatch = True
                break

    if missing_columns or extra_columns or type_mismatch:
        eprint("Обнаружено несоответствие структуры таблицы. Пересоздаём levels...")
        with engine.connect() as conn:
            conn.execute(text("DROP TABLE IF EXISTS levels"))
            conn.commit()
        metadata.create_all(engine)
        eprint("Таблица levels пересоздана (с колонкой instrument).")
        return

    # Добавляем индекс для ускорения запросов по instrument (опционально)
    indexes = inspector.get_indexes('levels')
    index_columns = [idx['column_names'] for idx in indexes]
    if ['instrument'] not in index_columns and ['table_name'] not in index_columns:
        with engine.connect() as conn:
            conn.execute(text("ALTER TABLE levels ADD INDEX idx_instrument (instrument)"))
            conn.commit()
        eprint("Добавлен индекс на instrument.")

# ------------------------------------------------------------
# Загрузка данных из таблицы (используем timestamp)
# ------------------------------------------------------------
def load_last_rows(table_name, limit=LIMIT):
    """
    Загружает последние limit записей из таблицы, сортирует по возрастанию времени.
    Использует колонку timestamp (BIGINT 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

# ------------------------------------------------------------
# Взвешенная кластеризация цен
# ------------------------------------------------------------
def cluster_prices_weighted(price_weight_list, threshold):
    """
    price_weight_list: список кортежей (price, weight)
    threshold: абсолютный порог схождения
    возвращает список словарей [{'price': avg, 'strength': total_weight}]
    """
    if not price_weight_list:
        return []
    sorted_items = sorted(price_weight_list, key=lambda x: x[0])
    clusters = []
    current_prices = [sorted_items[0][0]]
    current_weight = sorted_items[0][1]
    for price, weight in sorted_items[1:]:
        if price - current_prices[-1] <= threshold:
            current_prices.append(price)
            current_weight += weight
        else:
            avg_price = sum(current_prices) / len(current_prices)
            clusters.append({'price': avg_price, 'strength': current_weight})
            current_prices = [price]
            current_weight = weight
    # последний кластер
    avg_price = sum(current_prices) / len(current_prices)
    clusters.append({'price': avg_price, 'strength': current_weight})
    return clusters

# ------------------------------------------------------------
# Поиск уровней (экстремумы + опционально Pivot + объёмы)
# ------------------------------------------------------------
def find_support_resistance_levels(df, table_name,
                                   cluster_threshold=LEVEL_CLUSTER_THRESHOLD,
                                   include_pivot=INCLUDE_PIVOT,
                                   use_volume=USE_VOLUME,
                                   min_strength=MIN_STRENGTH):
    suffix = table_name.split('_')[-1] if '_' in table_name else ''
    window = WINDOW_BY_SUFFIX.get(suffix, 5)   # теперь включает и 'W1'
    if len(df) < window * 2 + 1:
        return []

    highs = df['High'].values
    lows = df['Low'].values
    closes = df['Close'].values

    high_indices = []
    low_indices = []

    for i in range(window, len(highs) - window):
        # Локальный максимум по High
        if all(highs[i] > highs[i - j] for j in range(1, window + 1)) and \
           all(highs[i] > highs[i + j] for j in range(1, window + 1)):
            high_indices.append(i)
        # Локальный минимум по Low
        if all(lows[i] < lows[i - j] for j in range(1, window + 1)) and \
           all(lows[i] < lows[i + j] for j in range(1, window + 1)):
            low_indices.append(i)

    last_price = closes[-1]
    abs_threshold = cluster_threshold * last_price

    # Учёт объёмов
    if use_volume and 'Volume' in df.columns:
        median_volume = df['Volume'].median()
        if median_volume == 0:
            median_volume = 1
    else:
        use_volume = False

    def get_weight(idx):
        if use_volume:
            vol = df['Volume'].iloc[idx]
            bonus = min((vol / median_volume) * 0.5, 2.0)
            return 1.0 + bonus
        return 1.0

    all_items = []
    for idx in high_indices:
        all_items.append((highs[idx], get_weight(idx)))
    for idx in low_indices:
        all_items.append((lows[idx], get_weight(idx)))

    if include_pivot:
        last = df.iloc[-1]
        high, low, close = last['High'], last['Low'], last['Close']
        pp = (high + low + close) / 3
        r1 = 2 * pp - low
        s1 = 2 * pp - high
        r2 = pp + (high - low)
        s2 = pp - (high - low)
        for p in [pp, r1, s1, r2, s2]:
            all_items.append((p, 1.0))

    if not all_items:
        return []

    clustered = cluster_prices_weighted(all_items, abs_threshold)

    final_levels = []
    for lev in clustered:
        price = lev['price']
        strength = lev['strength']
        if price < last_price - abs_threshold:
            lev_type = 'support'
        elif price > last_price + abs_threshold:
            lev_type = 'resistance'
        else:
            continue
        if strength >= min_strength:
            final_levels.append({'price': price, 'strength': strength, 'type': lev_type})

    final_levels.sort(key=lambda x: x['price'])
    return final_levels

# ------------------------------------------------------------
# Обновление таблицы levels (с учётом instrument)
# ------------------------------------------------------------
def update_levels_table(table_name, instrument, levels):
    delete_stmt = text("DELETE FROM levels WHERE table_name = :table_name")
    insert_stmt = text("""
        INSERT INTO levels (table_name, instrument, level_type, price, strength, last_updated)
        VALUES (:table_name, :instrument, :level_type, :price, :strength, :last_updated)
    """)
    now = pd.Timestamp.now()
    with engine.connect() as conn:
        conn.execute(delete_stmt, {"table_name": table_name})
        for lev in levels:
            conn.execute(insert_stmt, {
                "table_name": table_name,
                "instrument": instrument,
                "level_type": lev['type'],
                "price": float(lev['price']),
                "strength": float(lev['strength']),
                "last_updated": now
            })
        conn.commit()

# ------------------------------------------------------------
# Проверка времени обработки (как во втором коде)
# ------------------------------------------------------------
def should_process_table(table_name, current_time):
    """
    Определяет, нужно ли обрабатывать таблицу в зависимости от системного времени.
    Для W1: только в понедельник с 00:00 до 00:05.
    Для H1: первые 5 минут часа.
    Для D1: первые 5 минут дня.
    """
    # ИЗМЕНЕНО: добавлено условие для W1
    if table_name.endswith('_W1'):
        # Анализ недельных данных выполняется только в понедельник с 00:00 до 00:05
        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

# ------------------------------------------------------------
# Основной цикл с аргументом --inst и поддержкой JSON
# ------------------------------------------------------------
def main():
    parser = argparse.ArgumentParser(description='Поиск уровней поддержки/сопротивления по свечным таблицам')
    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

    # Формируем список таблиц для данного инструмента
    tables_to_process = [f"{instrument_clean}_{suffix}" for suffix in TIMEFRAME_SUFFIXES]

    # Проверяем/создаём таблицу levels с правильной структурой
    ensure_levels_table_exists()

    # Получаем inspector для проверки существования таблиц
    inspector = inspect(engine)
    current_time = datetime.now()

    results = []

    for table in tables_to_process:
        eprint(f"\n--- Обработка таблицы {table} ---")

        res = {
            "instrument": instrument_clean,
            "table_name": table,
            "success": False,
            "message": "",
            "levels": []  # массив уровней
        }

        # Проверка существования таблицы в базе данных
        if not inspector.has_table(table):
            msg = f"Таблица {table} не существует в базе данных. Пропускаем."
            eprint(f"  {msg}")
            res["message"] = msg
            # Удаляем возможные старые записи для этой таблицы (опционально)
            with engine.connect() as conn:
                conn.execute(text("DELETE FROM levels WHERE table_name = :table_name"), {"table_name": table})
                conn.commit()
            results.append(res)
            continue

        # Проверка времени
        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

        # Загружаем данные
        df = load_last_rows(table, LIMIT)

        if df.empty:
            msg = f"Нет данных в таблице {table}. Пропускаем."
            eprint(f"  {msg}")
            res["message"] = msg
            # Очищаем возможные старые уровни
            with engine.connect() as conn:
                conn.execute(text("DELETE FROM levels WHERE table_name = :table_name"), {"table_name": table})
                conn.commit()
            results.append(res)
            continue

        # Проверяем, достаточно ли данных для анализа (минимум 20 свечей, как в оригинале)
        if len(df) < 20:
            msg = f"Недостаточно данных (менее 20 записей). Пропускаем."
            eprint(f"  {msg}")
            res["message"] = msg
            with engine.connect() as conn:
                conn.execute(text("DELETE FROM levels WHERE table_name = :table_name"), {"table_name": table})
                conn.commit()
            results.append(res)
            continue

        # Поиск уровней
        levels = find_support_resistance_levels(
            df, table,
            cluster_threshold=LEVEL_CLUSTER_THRESHOLD,
            include_pivot=INCLUDE_PIVOT,
            use_volume=USE_VOLUME,
            min_strength=MIN_STRENGTH
        )

        # Формируем список уровней для вывода
        level_list = []
        for lev in levels:
            level_list.append({
                "price": float(lev['price']),
                "strength": float(lev['strength']),
                "type": lev['type']
            })
        res["levels"] = level_list

        if not levels:
            eprint("  Уровни не найдены.")
        else:
            eprint(f"  Найдено уровней: {len(levels)}")
            for lev in levels:
                eprint(f"    {lev['type'].capitalize()}: {lev['price']:.2f} (сила {lev['strength']:.2f})")

        # Обновляем таблицу уровней
        update_levels_table(table, instrument_clean, levels)

        # Помечаем как успешно обработанную
        res["success"] = True
        res["message"] = "OK"
        results.append(res)

    eprint("\nАнализ завершён.")

    if args.json:
        output = {
            "success": True,
            "message": "OK",
            "results": results
        }
        print(json.dumps(output, ensure_ascii=False, indent=2))

if __name__ == "__main__":
    main()