/**
 * ╔═══════════════════════════════════════════════════════════════╗
 * ║  DeviceService — управление устройствами ИНФУД v6.4           ║
 * ╚═══════════════════════════════════════════════════════════════╝
 *
 * Архитектура взаимодействия с устройствами:
 *
 *   JSON (data/devices.json)           ← persistence
 *        ↕
 *   DeviceService                      ← бизнес-логика
 *        ↕
 *   FeatureFlagsService                ← проверка MODBUS/OPCUA/MQTT
 *        ↕
 *   LkvService.setSensorField()        ← обновление in-memory снимка
 *        ↕
 *   WebSocketGateway.broadcast()       ← push на фронтенд
 *
 * Реальные протоколы (Modbus/OPC-UA/MQTT) подключаются как
 * опциональные npm-пакеты. Если пакет не установлен —
 * FeatureFlagGuard вернёт 503 для управляющих эндпоинтов,
 * а данные продолжат обновляться из симуляции.
 *
 * Совместимость с infud-server-v4.js:
 *   GET  /api/devices      → массив DeviceRecord
 *   POST /api/devices      → { ok, id }
 *   PUT  /api/devices/:id  → { ok }
 *   DELETE /api/devices/:id→ { ok }
 *   GET  /api/catalog      → массив каталога (с фильтром q и cat)
 *   POST /api/catalog      → { ok, id }
 */

import {
  Injectable, Logger, OnModuleInit,
  NotFoundException, BadRequestException,
} from '@nestjs/common';
import * as fs     from 'fs';
import * as path   from 'path';
import * as crypto from 'crypto';

import { FeatureFlagsService } from '../../config/feature-flags.service';
import { LkvService }          from '../store/lkv.service';
import { CreateDeviceDto, DeviceRecord } from '../../dto/device/device.dto';

// ── Каталог устройств по умолчанию ──────────────────────────────
// Минимальный inline-каталог (13 типов).
// Полный каталог (75+ типов) загружается из devices_catalog.json.
const CATALOG_INLINE: CatalogItem[] = [
  { id:'temp_pt100',   cat:'temperature', name:'PT100/PT1000',           vendor:'Universal',      proto:['modbus_tcp','modbus_rtu','analog_420'], unit:'°C',    scale:0.1,  minNorm:-50,  maxNorm:200,  widget:'kpi'   },
  { id:'temp_pasteur', cat:'temperature', name:'T° пастеризации E+H',    vendor:'Endress+Hauser', proto:['modbus_tcp'],                           unit:'°C',    scale:0.1,  minNorm:72,   maxNorm:85,   widget:'gauge' },
  { id:'temp_ecograph',cat:'temperature', name:'E+H Ecograph RSG35',     vendor:'Endress+Hauser', proto:['modbus_tcp'],                           unit:'°C',    scale:0.01, model:'RSG35',               widget:'kpi'   },
  { id:'temp_jumo',    cat:'temperature', name:'Jumo Logoscreen 600',    vendor:'Jumo',           proto:['modbus_tcp'],                           unit:'°C',    scale:0.01, model:'Logoscreen 600',      widget:'kpi'   },
  { id:'ph_inline',    cat:'quality',     name:'pH inline Mettler-Toledo',vendor:'Mettler-Toledo', proto:['modbus_tcp','analog_420'],              unit:'pH',    scale:0.01, minNorm:3.2,  maxNorm:4.5,  widget:'gauge' },
  { id:'brix_kpatents',cat:'quality',     name:'K-Patents PR-23',        vendor:'K-Patents',      proto:['modbus_tcp'],                           unit:'°Brix', model:'PR-23',                           widget:'gauge' },
  { id:'flow_promag',  cat:'flow',        name:'E+H Proline Promag 10W', vendor:'Endress+Hauser', proto:['modbus_tcp'],                           unit:'м³/ч',  model:'Promag 10W',                      widget:'kpi'   },
  { id:'flow_krohne',  cat:'flow',        name:'KROHNE IFS 4000 F',      vendor:'KROHNE',         proto:['modbus_tcp'],                           unit:'м³/ч',  model:'IFS 4000 F',                      widget:'kpi'   },
  { id:'level_vega',   cat:'level',       name:'VEGA VEGAPULS 69',       vendor:'VEGA',           proto:['modbus_tcp'],                           unit:'м',     model:'VEGAPULS 69',                     widget:'level' },
  { id:'mb210_gateway',cat:'gateway',     name:'Owen MB210',             vendor:'ОВЕН',           proto:['modbus_tcp'],                           unit:'—',     model:'MB210',                           widget:'status'},
  { id:'plc_s7_1200',  cat:'gateway',     name:'Siemens S7-1200',        vendor:'Siemens',        proto:['s7comm','profinet','modbus_tcp'],        unit:'—',     model:'S7-1200',                         widget:'status'},
  { id:'oee_line',     cat:'oee',         name:'OEE линии (%)',          vendor:'Calculated',     proto:['internal'],                             unit:'%',     minNorm:85, maxNorm:100,                  widget:'gauge' },
  { id:'manual_any',   cat:'manual',      name:'Ручной ввод',            vendor:'—',              proto:['manual'],                               unit:'любая',                                            widget:'kpi'   },
];

/** Тип записи каталога устройств */
export interface CatalogItem {
  id:       string;
  cat:      string;
  name:     string;
  vendor:   string;
  proto:    string[];
  unit:     string;
  scale?:   number;
  minNorm?: number;
  maxNorm?: number;
  model?:   string;
  widget:   string;
  custom?:  boolean;
}

/** Состояния опроса Modbus/OPC-UA */
interface PollingEntry {
  device:    DeviceRecord;
  connected: boolean;
  retries:   number;
  client:    any | null;
  lastValue: number | null;
  lastSeen:  string | null;
}

@Injectable()
export class DeviceService implements OnModuleInit {
  private readonly logger = new Logger(DeviceService.name);

  /** Активные Modbus-клиенты: deviceId → состояние */
  private readonly modbusClients = new Map<string, PollingEntry>();
  /** Активные OPC-UA сессии: deviceId → состояние */
  private readonly opcuaSessions = new Map<string, PollingEntry>();

  constructor(
    private readonly features: FeatureFlagsService,
    private readonly lkv:      LkvService,
  ) {}

  // ─────────────────────────────────────────────────────────────
  // LIFECYCLE
  // ─────────────────────────────────────────────────────────────

  onModuleInit(): void {
    // Запускаем опрос реальных устройств каждые 5 сек
    setInterval(() => this.pollAllDevices(), 5_000);
    // Инициализируем соединения для существующих устройств
    this.initExistingDevices();
  }

  /**
   * При старте проходим по сохранённым устройствам
   * и инициализируем соединения для Modbus/OPC-UA.
   */
  private initExistingDevices(): void {
    const devices = this.readDevices();
    for (const d of devices) {
      this.initDevice(d).catch(e =>
        this.logger.warn(`Не удалось инициализировать ${d.name}: ${e.message}`)
      );
    }
    this.logger.log(`Инициализированы ${devices.length} устройств из хранилища`);
  }

  // ─────────────────────────────────────────────────────────────
  // DEVICES CRUD
  // ─────────────────────────────────────────────────────────────

  /** GET /api/devices → все устройства */
  getAll(): DeviceRecord[] {
    return this.readDevices();
  }

  /**
   * POST /api/devices — добавить устройство.
   * После добавления автоматически инициализирует Modbus/OPC-UA.
   */
  async create(dto: CreateDeviceDto): Promise<{ ok: boolean; id: string }> {
    const devices = this.readDevices();
    const record: DeviceRecord = {
      ...dto,
      id:        'd_' + Date.now(),
      status:    'pending',
      lastValue: null,
      lastSeen:  null,
      createdAt: new Date().toISOString(),
    };
    devices.push(record);
    await this.writeDevices(devices);

    // Асинхронно инициализируем протокол
    this.initDevice(record).catch(e =>
      this.logger.warn(`Инициализация устройства ${record.name}: ${e.message}`)
    );

    return { ok: true, id: record.id };
  }

  /**
   * PUT /api/devices/:id — обновить параметры устройства.
   * При изменении proto/ip/port — переинициализирует соединение.
   */
  async update(id: string, dto: Partial<CreateDeviceDto>): Promise<void> {
    const devices = this.readDevices();
    const idx     = devices.findIndex(d => d.id === id);
    if (idx < 0) throw new NotFoundException(`Устройство не найдено: ${id}`);

    const oldProto = devices[idx].proto;
    devices[idx]   = { ...devices[idx], ...dto };
    await this.writeDevices(devices);

    // Переподключаем если изменился протокол или адрес
    if (dto.proto || dto.ip || dto.port) {
      this.disconnectDevice(id);
      this.initDevice(devices[idx]).catch(() => {});
    }
  }

  /** DELETE /api/devices/:id */
  async remove(id: string): Promise<void> {
    let devices = this.readDevices();
    const found = devices.find(d => d.id === id);
    if (!found) throw new NotFoundException(`Устройство не найдено: ${id}`);
    this.disconnectDevice(id);
    devices = devices.filter(d => d.id !== id);
    await this.writeDevices(devices);
  }

  // ─────────────────────────────────────────────────────────────
  // CATALOG
  // ─────────────────────────────────────────────────────────────

  /**
   * GET /api/catalog?q=&cat=
   * Поиск по имени/вендору/модели и фильтр по категории.
   */
  getCatalog(q?: string, cat?: string): CatalogItem[] {
    let catalog = this.readCatalog();
    if (q) {
      const ql = q.toLowerCase();
      catalog = catalog.filter(d =>
        (d.name + (d.vendor ?? '') + (d.model ?? '') + d.id).toLowerCase().includes(ql)
      );
    }
    if (cat) catalog = catalog.filter(d => d.cat === cat);
    return catalog;
  }

  /** POST /api/catalog — добавить тип устройства в каталог */
  async addToCatalog(item: Partial<CatalogItem>): Promise<{ ok: boolean; id: string }> {
    const catalog = this.readCatalog();
    const newItem: CatalogItem = {
      id:     'custom_' + Date.now(),
      cat:    item.cat    ?? 'manual',
      name:   item.name   ?? 'Custom Device',
      vendor: item.vendor ?? '—',
      proto:  item.proto  ?? ['manual'],
      unit:   item.unit   ?? '—',
      widget: item.widget ?? 'kpi',
      custom: true,
      ...item,
    };
    catalog.push(newItem);
    await this.writeCatalog(catalog);
    return { ok: true, id: newItem.id };
  }

  // ─────────────────────────────────────────────────────────────
  // POLLING ENGINE
  // ─────────────────────────────────────────────────────────────

  /**
   * Основной цикл опроса всех устройств.
   * Вызывается каждые 5 сек через setInterval.
   */
  private async pollAllDevices(): Promise<void> {
    const devices = this.readDevices();
    for (const d of devices) {
      try {
        await this.pollDevice(d);
      } catch (e) {
        // Логируем только изредка чтобы не спамить
        if (Math.random() < 0.01) {
          this.logger.debug(`Poll ${d.name}: ${(e as Error).message}`);
        }
      }
    }
  }

  /** Опрос одного устройства в зависимости от протокола */
  private async pollDevice(device: DeviceRecord): Promise<void> {
    const proto = device.proto ?? '';

    if ((proto === 'modbus_tcp' || proto === 'modbus_rtu') && this.features.isEnabled('MODBUS')) {
      await this.pollModbus(device);
    } else if (proto === 'opcua' && this.features.isEnabled('OPCUA')) {
      await this.pollOpcua(device);
    }
    // MQTT — пассивный (push), опроса нет
  }

  // ─────────────────────────────────────────────────────────────
  // MODBUS ENGINE
  // ─────────────────────────────────────────────────────────────

  /**
   * Инициализация Modbus TCP/RTU соединения.
   * Exponential backoff при повторных попытках подключения.
   */
  private async initModbus(device: DeviceRecord): Promise<void> {
    if (!this.features.isEnabled('MODBUS')) return;
    let ModbusRTU: any;
    try { ModbusRTU = require('modbus-serial'); } catch (_) { return; }

    const key     = device.id;
    const existing = this.modbusClients.get(key);
    if (existing?.connected) return; // уже подключён

    const client = new ModbusRTU();
    const entry: PollingEntry = {
      device, connected: false, retries: 0, client, lastValue: null, lastSeen: null,
    };
    this.modbusClients.set(key, entry);

    const connect = async () => {
      try {
        if (device.proto === 'modbus_tcp') {
          await client.connectTCP(device.ip!, { port: parseInt(device.port ?? '502') });
        } else {
          await client.connectRTUBuffered(device.ip!, { baudRate: 9600 });
        }
        client.setID(parseInt(device.unit ?? '1'));
        client.setTimeout(3_000);
        entry.connected = true;
        entry.retries   = 0;
        this.logger.log(`Modbus подключён: ${device.name} @ ${device.ip}:${device.port}`);
        // Обновляем статус в JSON
        await this.updateDeviceStatus(device.id, 'ok');
      } catch (e) {
        entry.connected = false;
        entry.retries++;
        const delay = Math.min(30_000, 2_000 * Math.pow(2, entry.retries));
        this.logger.warn(
          `Modbus ${device.name} — ошибка подключения: ${(e as Error).message}. ` +
          `Retry через ${delay / 1000}s (попытка ${entry.retries})`
        );
        setTimeout(connect, delay);
      }
    };

    await connect();
  }

  /**
   * Чтение Holding Register и обновление LKV.
   * Применяет scale и offset из конфигурации устройства.
   */
  private async pollModbus(device: DeviceRecord): Promise<void> {
    const entry = this.modbusClients.get(device.id);
    if (!entry?.connected) { await this.initModbus(device); return; }

    try {
      // Регистр: 40001 → адрес 0, 40010 → адрес 9
      const reg    = parseInt(device.reg ?? '40001');
      const addr   = reg > 40000 ? reg - 40001 : reg;
      const result = await entry.client.readHoldingRegisters(addr, 2);
      const raw    = result.data[0] as number;

      // Применяем масштабирование: value = raw * scale + offset
      const scale  = parseFloat(device.scale  ?? '0.1');
      const offset = parseFloat(device.offset ?? '0');
      const value  = +(raw * scale + offset).toFixed(4);

      entry.lastValue = value;
      entry.lastSeen  = new Date().toISOString();

      // Записываем в LKV если устройство привязано к KPI-полю
      this.applyToLkv(device, value, 'modbus');

      // Обновляем статус в JSON-хранилище
      await this.updateDeviceStatus(device.id, 'ok', value, entry.lastSeen);
    } catch (e) {
      entry.connected = false;
      await this.updateDeviceStatus(device.id, 'error');
      this.logger.warn(`Modbus poll ${device.name}: ${(e as Error).message}`);
      // Переподключение через 5 сек
      setTimeout(() => this.initModbus(device), 5_000);
    }
  }

  // ─────────────────────────────────────────────────────────────
  // OPC-UA ENGINE
  // ─────────────────────────────────────────────────────────────

  /** Инициализация OPC-UA сессии */
  private async initOpcua(device: DeviceRecord): Promise<void> {
    if (!this.features.isEnabled('OPCUA') || device.proto !== 'opcua') return;
    let opcua: any;
    try { opcua = require('node-opcua'); } catch (_) { return; }

    try {
      const endpointUrl = `opc.tcp://${device.ip}:${device.port ?? 4840}`;
      const client = opcua.OPCUAClient.create({
        applicationName:    'ИНФУД Dashboard',
        endpointMustExist:  false,
        connectionStrategy: { maxRetry: 10, initialDelay: 2_000, maxDelay: 30_000 },
        securityMode:       opcua.MessageSecurityMode.None,
        securityPolicy:     opcua.SecurityPolicy.None,
      });
      await client.connect(endpointUrl);
      const session = await client.createSession();
      this.opcuaSessions.set(device.id, {
        device, connected: true, retries: 0,
        client: { client, session }, lastValue: null, lastSeen: null,
      });
      this.logger.log(`OPC-UA подключён: ${device.name} @ ${endpointUrl}`);
    } catch (e) {
      this.logger.warn(`OPC-UA ${device.name} — ошибка: ${(e as Error).message}`);
      setTimeout(() => this.initOpcua(device), 10_000);
    }
  }

  private async pollOpcua(device: DeviceRecord): Promise<void> {
    const entry = this.opcuaSessions.get(device.id);
    if (!entry?.connected) { await this.initOpcua(device); return; }

    try {
      const { session } = entry.client;
      const nodeId  = device.reg ?? `ns=2;s=${device.name}`;
      const dv      = await session.readVariableValue(nodeId);
      const raw     = dv?.value?.value;
      if (raw == null) return;

      const scale  = parseFloat(device.scale  ?? '1');
      const offset = parseFloat(device.offset ?? '0');
      const value  = +(raw * scale + offset).toFixed(4);

      entry.lastValue = value;
      entry.lastSeen  = new Date().toISOString();
      this.applyToLkv(device, value, 'opcua');
      await this.updateDeviceStatus(device.id, 'ok', value, entry.lastSeen);
    } catch (e) {
      entry.connected = false;
      this.opcuaSessions.delete(device.id);
      await this.updateDeviceStatus(device.id, 'error');
      setTimeout(() => this.initOpcua(device), 8_000);
    }
  }

  // ─────────────────────────────────────────────────────────────
  // LKV MAPPING
  // ─────────────────────────────────────────────────────────────

  /**
   * Если устройство привязано к KPI-полю (device.kpi = 'kv-ps'),
   * обновляем соответствующее поле в LkvService.
   *
   * Маппинг: 'kv-ps' → 'past_steklo', 'kv-ph' → 'ph', и т.д.
   */
  private applyToLkv(
    device: DeviceRecord,
    value:  number,
    source: string,
  ): void {
    if (!device.kpi) return;

    const kpiToField: Record<string, string> = {
      'kv-ps':  'past_steklo',
      'kv-sig': 'past_sig',
      'kv-tp':  'past_tetra',
      'kv-ph':  'ph',
      'kv-brix':'brix',
      'kv-oee': 'oee',
      'kv-out': 'output',
      'kv-pw':  'power_kw',
      'kv-wt':  'water_today',  // будет в resources
      'kv-el':  'power_today',
      'kv-st':  'steam_today',
      'kv-air': 'air_bar',
    };

    const field = kpiToField[device.kpi];
    if (field) {
      this.lkv.setSensorField(
        field as any,
        value,
        source,
      );
      this.logger.debug(`LKV обновлён: ${field} = ${value} (${source} — ${device.name})`);
    }
  }

  // ─────────────────────────────────────────────────────────────
  // PRIVATE HELPERS
  // ─────────────────────────────────────────────────────────────

  /** Инициализировать устройство по типу протокола */
  private async initDevice(device: DeviceRecord): Promise<void> {
    const proto = device.proto ?? '';
    if (proto === 'modbus_tcp' || proto === 'modbus_rtu') {
      if (this.features.isEnabled('MODBUS') && device.ip) {
        await this.initModbus(device);
      }
    } else if (proto === 'opcua') {
      if (this.features.isEnabled('OPCUA') && device.ip) {
        await this.initOpcua(device);
      }
    }
    // MQTT-устройства подключаются через MqttService (подписка на topics)
  }

  /** Закрыть соединение с устройством */
  private disconnectDevice(id: string): void {
    const mb = this.modbusClients.get(id);
    if (mb?.client) {
      try { mb.client.close(); } catch (_) {}
      this.modbusClients.delete(id);
    }
    const opc = this.opcuaSessions.get(id);
    if (opc?.client) {
      try { opc.client.session?.close(); opc.client.client?.disconnect(); } catch (_) {}
      this.opcuaSessions.delete(id);
    }
  }

  /** Обновить статус устройства в JSON-хранилище (неблокирующий) */
  private async updateDeviceStatus(
    id:        string,
    status:    'ok' | 'error' | 'pending',
    lastValue?: number | null,
    lastSeen?:  string | null,
  ): Promise<void> {
    const devices = this.readDevices();
    const d       = devices.find(x => x.id === id);
    if (!d) return;
    d.status = status;
    if (lastValue !== undefined) d.lastValue = lastValue;
    if (lastSeen  !== undefined) d.lastSeen  = lastSeen;
    await this.writeDevices(devices);
  }

  // ─────────────────────────────────────────────────────────────
  // JSON PERSISTENCE
  // ─────────────────────────────────────────────────────────────

  private get dataDir(): string {
    return process.env.DATA_DIR ?? path.join(process.cwd(), 'data');
  }

  private readJson<T>(file: string, fallback: T): T {
    const p = path.join(this.dataDir, file + '.json');
    try {
      return fs.existsSync(p) ? JSON.parse(fs.readFileSync(p, 'utf8')) : fallback;
    } catch (_) { return fallback; }
  }

  private async writeJson(file: string, data: unknown): Promise<void> {
    fs.mkdirSync(this.dataDir, { recursive: true });
    const p   = path.join(this.dataDir, file + '.json');
    const tmp = p + '.tmp';
    fs.writeFileSync(tmp, JSON.stringify(data, null, 2), 'utf8');
    fs.renameSync(tmp, p);
  }

  private readDevices():  DeviceRecord[]  { return this.readJson<DeviceRecord[]>('devices', []); }
  private writeDevices(d: DeviceRecord[]) { return this.writeJson('devices', d); }
  private readCatalog():  CatalogItem[]   { return this.readJson<CatalogItem[]>('devices_catalog', CATALOG_INLINE); }
  private writeCatalog(c: CatalogItem[])  { return this.writeJson('devices_catalog', c); }
}
