#!/usr/bin/env python3
"""
sync_signals — синхронизация сигналов MoE v12 из summary_report в MySQL (current_signals).

Вызывается из monitor.py после завершения анализа, либо вручную:
    python -m utils.sync_signals

Записывает для каждого тикера: signal_type, confidence, p_long, p_short,
tp_price, sl_price, close_price, updated_at.
"""

import sys, os, re, time
from typing import List, Dict, Optional

import numpy as np
import pandas as pd

sys.path.insert(0, os.path.dirname(os.path.dirname(os.path.abspath(__file__))))
from db.connection import get_connection

# MOEX акции торгуются только в LONG — запрещаем SELL/p_short при записи в БД.
try:
    from config import is_moex_only_long as _is_moex_only_long
except ImportError:
    _is_moex_only_long = None


def signal_is_moex_long_only(ticker: str) -> bool:
    """True если тикер относится к MOEX в режиме only-LONG (no short)."""
    if _is_moex_only_long is None:
        return False
    try:
        return _is_moex_only_long(ticker)
    except Exception:
        return False


MONITOR_DIR = os.path.join(os.path.dirname(os.path.dirname(os.path.abspath(__file__))), 'monitoring')
SUMMARY_FILE = os.path.join(MONITOR_DIR, 'summary_report.md')

# Fallback-тикеры (если summary не найдём)
FALLBACK_TICKERS = ['ASTR', 'GAZP', 'LKOH', 'MOEX', 'MTSS', 'NSVZ', 'NVTK',
                    'PHOR', 'PLZL', 'ROSN', 'SBER', 'SNGSP', 'VTBR', 'X5',
                    'BITCOIN', 'BITCOINC', 'ZCASH', 'EURUSD']


def parse_summary_report(filepath: str = SUMMARY_FILE) -> List[Dict]:
    """Парсит summary_report.md и возвращает список сигналов по тикерам."""
    if not os.path.exists(filepath):
        print(f"[sync_signals] Файл не найден: {filepath}")
        return []

    with open(filepath, 'r', encoding='utf-8') as f:
        content = f.read()

    signals = []
    # Ищем строки таблицы: | TICKER | ok | SIGNAL | conf% | P_long% | P_short% | ... |
    table_pattern = re.compile(
        r'^\|\s*(\w+)\s*\|\s*(ok|stale|error|no_candle|no_model)\s*\|\s*(\w+)\s*\|\s*([\d.]+)%\s*\|\s*([\d.]+)%\s*\|\s*([\d.]+)%',
        re.MULTILINE
    )

    for match in table_pattern.finditer(content):
        ticker = match.group(1)
        status = match.group(2)
        signal = match.group(3)
        conf = float(match.group(4)) / 100.0
        p_long = float(match.group(5)) / 100.0
        p_short = float(match.group(6)) / 100.0

        if status == 'ok':
            signals.append({
                'ticker': ticker,
                'signal_type': signal if signal in ('BUY', 'SELL') else 'NEUTRAL',
                'confidence': conf,
                'p_long': p_long,
                'p_short': p_short,
                'tp_price': 0.0,
                'sl_price': 0.0,
                'close_price': 0.0,
            })

    print(f"[sync_signals] Распарсено {len(signals)} сигналов из {filepath}")
    return signals


def sync_to_db(signals: List[Dict], model_type: str = 'moe_v12') -> int:
    """Upsert сигналов в таблицу current_signals."""
    if not signals:
        return 0

    now_ts = int(time.time())
    updated = 0

    with get_connection() as conn:
        cur = conn.cursor()

        for s in signals:
            # MOEX-only-long: акции РФ не торгуются в шорт — p_short принудительно 0.
            # Защита от "протечки" старых p_short из stale summary_report в current_signals.
            if signal_is_moex_long_only(s['ticker']):
                s['p_short'] = 0.0
                if s.get('signal_type') == 'SELL':
                    s['signal_type'] = 'NEUTRAL'
            cur.execute("""
                INSERT INTO current_signals
                    (ticker, signal_type, confidence, p_long, p_short,
                     tp_price, sl_price, close_price, model_type, updated_at)
                VALUES (%s, %s, %s, %s, %s, %s, %s, %s, %s, %s)
                ON DUPLICATE KEY UPDATE
                    signal_type = VALUES(signal_type),
                    confidence = VALUES(confidence),
                    p_long = VALUES(p_long),
                    p_short = VALUES(p_short),
                    tp_price = VALUES(tp_price),
                    sl_price = VALUES(sl_price),
                    close_price = VALUES(close_price),
                    model_type = VALUES(model_type),
                    updated_at = VALUES(updated_at)
            """, (
                s['ticker'], s['signal_type'], s['confidence'],
                s['p_long'], s['p_short'],
                s['tp_price'], s['sl_price'], s['close_price'],
                model_type, now_ts
            ))
            updated += 1

    print(f"[sync_signals] Записано {updated} сигналов в current_signals")
    return updated


# Порог устаревания сигнала в секундах (6 часов).
# Если current_signals.updated_at старее этого порога, сигнал обнуляется в NEUTRAL.
# Это защищает веб-дашборд от «зависших» сигналов (PLZL BUY от 01.08 — кейс 2026-08-03).
STALE_SIGNAL_THRESHOLD_SEC = 6 * 3600


def cleanup_stale_signals(threshold_sec: int = STALE_SIGNAL_THRESHOLD_SEC) -> int:
    """Сбрасывает устаревшие сигналы в current_signals в NEUTRAL.

    Если updated_at < (now - threshold_sec) для тикера — обнуляет signal_type,
    confidence, tp_price, sl_price, выставляет updated_at=now.

    Возвращает количество обнулённых строк.
    """
    now_ts = int(time.time())
    cutoff = now_ts - threshold_sec
    reset = 0

    with get_connection() as conn:
        cur = conn.cursor(dictionary=True)
        # Находим все тикеры с устаревшим updated_at и не-NEUTRAL сигналом
        cur.execute(
            "SELECT ticker, signal_type, updated_at FROM current_signals "
            "WHERE updated_at < %s AND signal_type != 'NEUTRAL'",
            (cutoff,)
        )
        stale_rows = cur.fetchall()

        if not stale_rows:
            cur.close()
            return 0

        # Прологируем и сбросим
        for row in stale_rows:
            stale_sec = now_ts - int(row.get('updated_at') or 0)
            stale_h = stale_sec / 3600
            old_sig = row.get('signal_type', '?')
            print(f"[sync_signals] STALE {row['ticker']}: {old_sig} "
                  f"updated {stale_h:.1f}h ago → NEUTRAL")
            logger = None
            try:
                import logging
                logging.getLogger('AI_Strategy').info(
                    f"current_signals STALE: {row['ticker']} {old_sig} "
                    f"({stale_h:.1f}h ago) → NEUTRAL"
                )
            except Exception:
                pass

        cur.execute(
            "UPDATE current_signals "
            "SET signal_type='NEUTRAL', confidence=0, p_long=0, p_short=0, "
            "    tp_price=0, sl_price=0, close_price=0, updated_at=%s "
            "WHERE updated_at < %s AND signal_type != 'NEUTRAL'",
            (now_ts, cutoff)
        )
        reset = cur.rowcount
        conn.commit()
        cur.close()

    if reset > 0:
        print(f"[sync_signals] Сброшено устаревших сигналов: {reset}")
    return reset


def sync_ticker_prices(signals: List[Dict]) -> None:
    """Дополняет сигналы ценами close, tp, sl из individual monitor-файлов."""
    for s in signals:
        ticker = s['ticker']
        monitor_file = os.path.join(MONITOR_DIR, f'{ticker.lower()}_monitor.md')
        if not os.path.exists(monitor_file):
            continue

        try:
            with open(monitor_file, 'r', encoding='utf-8') as f:
                content = f.read()

            # Ищем последнюю строку таблицы сигналов
            lines = content.strip().split('\n')
            data_lines = [l for l in lines if l.startswith('| ') and l.count('|') >= 5]
            if data_lines:
                last_line = data_lines[-1]
                parts = [p.strip() for p in last_line.split('|')]
                if len(parts) >= 7:
                    # Парсим close
                    try:
                        s['close_price'] = float(parts[2].replace(',', ''))
                    except (ValueError, IndexError):
                        pass
        except Exception:
            pass

        # B1: TP/SL для BUY/SELL — из close + ATR (как TradeManager).
        # Раньше tp_price/sl_price всегда были 0.0, и веб-дашборд показывал
        # «TP 0 / SL 0» для всех MoE-сигналов.
        _compute_tp_sl(s)


def _compute_tp_sl(s: Dict) -> None:
    """Вычисляет TP/SL для сигнала BUY/SELL из последней H1-свечи + ATR.

    BUG-FIX (2026-07-31): множители берутся per-market из get_risk_params
    (MARKET_TARGET_CONFIGS). Ранее захардкожено 3/6 — это противоречило
    тренировочным целям для crypto (4/8) и forex (2/4). Также порог ATR
    валидации был global 5%/10%, что некорректно для crypto.

    Для NEUTRAL TP/SL = 0.
    """
    sig = s.get('signal_type', 'NEUTRAL')
    if sig not in ('BUY', 'SELL'):
        s['tp_price'] = 0.0
        s['sl_price'] = 0.0
        return

    try:
        from config import get_risk_params
        sl_mult, tp_mult, max_atr_ratio = get_risk_params(s['ticker'])

        from data.loader import load_dataframe
        from features.technical import engineer_features
        df = load_dataframe(s['ticker'], 'H1', limit=100, clean=True)
        if df is None or len(df) < 20:
            return
        df = engineer_features(df)

        atr_pct = float(df['atr_pct'].iloc[-1])
        close = float(df['Close'].iloc[-1])
        if not (atr_pct > 0) or close <= 0:
            return

        # Per-market проверка ATR на правдоподобность (как в
        # TradeManager._validate_price) —AINARNO: crypto разрешён до 15%,
        # MOEX до 5%, forex до 2%.
        if atr_pct > max_atr_ratio:
            return
        # Дополнительная защита от аномально большого всплеска (10× порога).
        if atr_pct > max_atr_ratio * 3:
            return

        if sig == 'BUY':
            s['tp_price'] = round(close * (1 + tp_mult * atr_pct), 8)
            s['sl_price'] = round(close * (1 - sl_mult * atr_pct), 8)
        else:
            s['tp_price'] = round(close * (1 - tp_mult * atr_pct), 8)
            s['sl_price'] = round(close * (1 + sl_mult * atr_pct), 8)
        s['close_price'] = round(close, 8)
    except Exception:
        pass


# ---------------------------------------------------------------------------
# Регрессионные (directional) сигналы
# ---------------------------------------------------------------------------

# Колонки directional сигналов
DIRECTIONAL_SIGNAL_COLS = [
    'rsi_signal', 'bb_signal', 'macd_signal',
    'trend_signal', 'momentum_signal', 'volume_signal',
    'directional_bias', 'signal_strength',
]


def sync_regression_signals(tickers: Optional[List[str]] = None) -> int:
    """Вычисляет directional сигналы для всех тикеров и сохраняет в regression_signal_data.
    
    Args:
        tickers: список тикеров. Если None — использует FALLBACK_TICKERS.
    
    Returns:
        количество записанных тикеров.
    """
    if tickers is None:
        tickers = FALLBACK_TICKERS

    from data.loader import load_dataframe
    from features.technical import engineer_features
    from features.directional import add_directional_signals

    now_ts = int(time.time())
    updated = 0

    for ticker in tickers:
        try:
            df = load_dataframe(ticker, 'H1', limit=200, clean=True)
            if df is None or len(df) < 50:
                continue

            # Feature engineering + directional signals
            df = engineer_features(df)
            df = add_directional_signals(df)

            last = df.iloc[-1]

            # Извлекаем 8 сигналов
            signals_data = {}
            for col in DIRECTIONAL_SIGNAL_COLS:
                val = float(last.get(col, 0.0))
                if pd.isna(val) or np.isnan(val):
                    val = 0.0
                signals_data[col] = round(val, 6)

            atr_pct = float(last.get('atr_pct', 0.0))
            close = float(last.get('Close', 0.0))

            with get_connection() as conn:
                cur = conn.cursor()
                cur.execute("""
                    INSERT INTO regression_signal_data
                        (ticker, rsi_signal, bb_signal, macd_signal,
                         trend_signal, momentum_signal, volume_signal,
                         directional_bias, signal_strength,
                         atr_pct, close_price, updated_at)
                    VALUES (%s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s)
                    ON DUPLICATE KEY UPDATE
                        rsi_signal = VALUES(rsi_signal),
                        bb_signal = VALUES(bb_signal),
                        macd_signal = VALUES(macd_signal),
                        trend_signal = VALUES(trend_signal),
                        momentum_signal = VALUES(momentum_signal),
                        volume_signal = VALUES(volume_signal),
                        directional_bias = VALUES(directional_bias),
                        signal_strength = VALUES(signal_strength),
                        atr_pct = VALUES(atr_pct),
                        close_price = VALUES(close_price),
                        updated_at = VALUES(updated_at)
                """, (
                    ticker,
                    signals_data['rsi_signal'],
                    signals_data['bb_signal'],
                    signals_data['macd_signal'],
                    signals_data['trend_signal'],
                    signals_data['momentum_signal'],
                    signals_data['volume_signal'],
                    signals_data['directional_bias'],
                    signals_data['signal_strength'],
                    round(atr_pct, 6),
                    round(close, 8),
                    now_ts,
                ))
                updated += 1

        except Exception as e:
            print(f"[sync_signals]  ✗ {ticker}: {e}")
            continue

    print(f"[sync_signals] Записано регрессионных сигналов: {updated}")
    return updated


def main():
    """Основной Entry Point."""
    import argparse
    parser = argparse.ArgumentParser(description='Синхронизация сигналов в MySQL')
    parser.add_argument('--regression', action='store_true',
                        help='Синхронизировать только регрессионные (directional) сигналы')
    parser.add_argument('--all', action='store_true',
                        help='Синхронизировать и основные, и регрессионные сигналы')
    args = parser.parse_args()

    if args.regression:
        sync_regression_signals()
        return

    if args.all:
        sync_regression_signals()

    signals = parse_summary_report()
    if not signals:
        print("[sync_signals] Нет сигналов для синхронизации.")
        return

    sync_ticker_prices(signals)
    sync_to_db(signals)

    # Показываем сводку
    buys = sum(1 for s in signals if s['signal_type'] == 'BUY')
    sells = sum(1 for s in signals if s['signal_type'] == 'SELL')
    print(f"[sync_signals] Итого: 🟢 {buys} BUY | 🔴 {sells} SELL | ⚪ {len(signals) - buys - sells} NEUTRAL")


if __name__ == '__main__':
    main()
