# frozen_string_literal: true # Dron::Ejecutor — Módulo atómico para ejecución de comandos # # Responsabilidad única: Ejecutar un comando y reportar resultado # # Fase 1.5: Reescrito para visibilidad en Bitácora Web # - Heartbeat actualiza events.descripcion (visible en web) # - Flujo simplificado: crear evento → heartbeat → cierre limpio # - Timeout implementado con Timeout.timeout() # - Errores de bitácora logeados explícitamente (sin rescue silencioso) require_relative 'base' require_relative '../../db/core/dron_db' require_relative '../../core/colores' require 'open3' require 'timeout' require 'time' require 'json' module Dron class Ejecutor DEFAULT_TIMEOUT = 3600 HEARTBEAT_INTERVAL = 30 class << self def lanzar(cmd:, nota: nil, evento_id: nil, timeout: DEFAULT_TIMEOUT, flujo_id: nil, nodo: nil) Base.logger.info("🛠️ Dron Ejecutor iniciado: #{Process.pid}") inicio = Time.now # 1. Registrar en tablas internas (dron_logs) resultado_inicio = ADN::DB::DronDB.registrar_inicio( tipo: 'ejecutor', cmd: cmd, flujo_id: flujo_id, metadata: { nota: nota, evento_id: evento_id, timeout: timeout, nodo: nodo } ) dron_id = resultado_inicio[:dron_id] ADN::DB::DronDB.registrar_avance( dron_id, paso: 1, descripcion: "Despegue inicial. Ejecutando comando...", estado: 'running', flujo_id: flujo_id ) # 2. Crear/vincular evento en Bitácora Web (visible) evento_id = crear_o_vincular_evento_web(dron_id, nota, evento_id, cmd, nodo) # 3. Ejecutar comando con heartbeat que actualiza AMBAS capas resultado = ejecutar_comando(cmd, dron_id, timeout, evento_id, nota, inicio) # 4. Registrar fin en tablas internas ADN::DB::DronDB.registrar_fin(dron_id, resultado[:exit_code], resultado[:output]) estado_fin = resultado[:exit_code] == 0 ? 'completed' : 'failed' desc_fin = resultado[:exit_code] == 0 ? "Aterrizaje exitoso. Comando finalizado." : "Aterrizaje con errores (Code #{resultado[:exit_code]})." ADN::DB::DronDB.registrar_avance( dron_id, paso: 2, descripcion: desc_fin, estado: estado_fin, flujo_id: flujo_id ) # 5. Cerrar evento en Bitácora Web (visible) cerrar_evento_web(evento_id, dron_id, nota, cmd, resultado, inicio) Base.logger.info("🛠️ Dron Ejecutor completado: #{dron_id} - #{resultado[:exit_code] == 0 ? '✅' : '❌'}") { dron_id: dron_id, exit_code: resultado[:exit_code], duration: resultado[:duration], evento_id: evento_id } rescue StandardError => e Base.logger.error("🛠️ Dron Ejecutor falló: #{e.message}") Base.logger.error(" #{e.backtrace&.first(3)&.join("\n ")}") { dron_id: nil, exit_code: 1, output: e.message, duration: 0, error: e.message } end private # ─── Ejecución con heartbeat dual ───────────────────────────── def ejecutar_comando(cmd, dron_id, timeout, evento_id, nota, inicio) start_time = Time.now output = [] last_out = "" accumulated_out = [] # Acumular todas las líneas de output para el heartbeat exit_code = nil progreso_visible = !ENV['BATCH_MODE'] # Solo mostrar progreso si no es batch # Mostrar encabezado de progreso if progreso_visible puts "\n#{Color::BOLD}#{Color::CYAN}⏳ Ejecutando...#{Color::RESET}" print " [" $stdout.flush end heartbeat_thread = Thread.new do tick = 0 start = Time.now loop do sleep(HEARTBEAT_INTERVAL) tick += 1 elapsed = (Time.now - start).to_i # Actualizar tabla interna (dron_logs.heartbeat) ADN::DB::DronDB.actualizar_heartbeat(dron_id) # Actualizar Bitácora Web (events.descripcion) — VISIBLE para el usuario actualizar_heartbeat_web(evento_id, dron_id, nota, cmd, inicio, tick, accumulated_out.join("\n")) if evento_id # Actualizar barra de progreso local if progreso_visible mins = elapsed / 60 segs = elapsed % 60 print "\r [#{Color::CYAN}#{mins}m #{segs}s#{Color::RESET}] ejecutando..." $stdout.flush end rescue StandardError => e Base.logger.warn("⚠️ Error en heartbeat ##{tick}: #{e.message}") end end begin Timeout.timeout(timeout) do Open3.popen3(cmd) do |_stdin, stdout, stderr, wait_thr| out_thread = Thread.new do begin while chunk = stdout.readpartial(1024) output << chunk lineas = chunk.split(/[\r\n]/).reject(&:empty?) last_out = lineas.last&.strip || last_out accumulated_out.concat(lineas) # Acumular líneas para heartbeat end rescue EOFError, IOError end end err_thread = Thread.new { output << stderr.read } out_thread.join err_thread.join exit_code = wait_thr.value.exitstatus end end rescue Timeout::Error exit_code = 124 output << "\n[TIMEOUT] #{timeout}s excedidos" rescue StandardError => e exit_code = 1 output << "\n[ERROR] #{e.message}" ensure heartbeat_thread.exit end # Finalizar barra de progreso if progreso_visible estado = exit_code == 0 ? "#{Color::GREEN}✅#{Color::RESET}" : "#{Color::RED}❌#{Color::RESET}" print "\r #{estado} " puts "#{Color::BOLD}Completado en #{Base.formato_duracion(Time.now - start_time)}#{Color::RESET}" end raw_output = output.join("\n").strip safe_output = raw_output.dup.force_encoding('UTF-8').scrub('?') { exit_code: exit_code || 1, output: safe_output, duration: (Time.now - start_time).round(2) } end # ─── Bitácora Web: Crear/Vincular evento ───────────────────── def crear_o_vincular_evento_web(dron_id, nota, evento_id, cmd, nodo) return nil unless nota # Sin nota = sin registro en web nodo_id = detectar_nodo_id(nodo) BitacorasDB::BitacoraDB.with_connection do |db| if evento_id.to_i > 0 # Vincular a evento existente: agregar nota de despegue entrada = db.execute("SELECT descripcion FROM bitacoras.entradas WHERE id = $1", [evento_id]).first target_table = entrada ? 'entradas' : 'events' target = entrada || db.execute("SELECT descripcion FROM bitacoras.events WHERE id = $1", [evento_id]).first if target nueva_desc = target['descripcion'].to_s + "\n\n🛸 **[#{Time.now.strftime('%H:%M')}] Vuelo Iniciado** (`#{dron_id}`): #{nota}" + "\n├─ Comando: `#{cmd}`" if target_table == 'entradas' db.update_entrada(evento_id, descripcion: nueva_desc) else db.actualizar_evento(evento_id, descripcion: nueva_desc) end end evento_id else # Crear evento nuevo en bitácora web descripcion = "🛸 #{nota} (`#{dron_id}`)\n├─ Comando: `#{cmd}`\n└─ ⏳ Iniciando..." bitacora_hoy = db.execute("SELECT id FROM bitacoras.bitacoras WHERE fecha::date = CURRENT_DATE LIMIT 1").first if bitacora_hoy entrada = db.create_entrada( inicio: Time.now.strftime('%H:%M'), descripcion: descripcion, estado: '⏳', modo: 'P', bitacora_id: bitacora_hoy['id'], nodo_id: nodo_id ) entrada['id'] else evento = db.crear_evento( nodo_id: nodo_id, descripcion: descripcion, inicio: Time.now.strftime('%H:%M'), estado: '⏳', modo: 'P', metadata: { dron_id: dron_id } ) evento['id'] end end end rescue StandardError => e Base.logger.error("🛸 Error creando evento web: #{e.message}") Base.logger.error(" #{e.backtrace&.first(3)&.join("\n ")}") nil end # ─── Bitácora Web: Heartbeat visible ───────────────────────── def actualizar_heartbeat_web(evento_id, dron_id, nota, cmd, inicio, tick, accumulated_out = "") duracion = Base.formato_duracion(Time.now - inicio) descripcion = "🛸 #{nota} (`#{dron_id}`)\n" \ "├─ Comando: `#{cmd}`\n" \ "└─ ❤️ #{duracion} — heartbeat ##{tick}" unless accumulated_out.nil? || accumulated_out.empty? # Limpiar caracteres no imprimibles y mostrar todo el output acumulado clean_out = accumulated_out.gsub(/[^[:print:]\n]/, '').strip # Mostrar últimas 15 líneas para no saturar la bitácora lineas = clean_out.split("\n") lineas_mostrar = lineas.last(15) descripcion += "\n\n```text\n#{lineas_mostrar.join("\n")}\n```" end BitacorasDB::BitacoraDB.with_connection do |db| # Intentar actualizar como entrada primero, luego como evento entrada = db.execute("SELECT id FROM bitacoras.entradas WHERE id = $1", [evento_id]).first if entrada db.update_entrada(evento_id, descripcion: descripcion) else db.actualizar_evento(evento_id, descripcion: descripcion) end end rescue StandardError => e Base.logger.warn("⚠️ Heartbeat web ##{tick} falló: #{e.message}") end # ─── Bitácora Web: Cierre de evento ────────────────────────── def cerrar_evento_web(evento_id, dron_id, nota, cmd, resultado, inicio) return unless evento_id.to_i > 0 duracion = Base.formato_duracion(resultado[:duration]) estado = resultado[:exit_code] == 0 ? '✅' : '❌' aterrizaje = resultado[:exit_code] == 0 ? 'Aterrizaje limpio' : "Fallo (code #{resultado[:exit_code]})" # Construir descripción final con output resumido output_resumen = Base.slice_safe(resultado[:output], 500) descripcion = "🛸 #{nota} (`#{dron_id}`) — #{estado} #{duracion}\n" \ "├─ Comando: `#{cmd}`\n" \ "└─ #{aterrizaje}" descripcion += "\n\n```text\n#{output_resumen}\n```" unless output_resumen.empty? BitacorasDB::BitacoraDB.with_connection do |db| entrada = db.execute("SELECT id FROM bitacoras.entradas WHERE id = $1", [evento_id]).first if entrada db.update_entrada(evento_id, descripcion: descripcion, estado: estado, fin: Time.now.strftime('%H:%M') ) else db.actualizar_evento(evento_id, descripcion: descripcion, estado: estado, fin: Time.now.strftime('%H:%M') ) end end Base.logger.info("🛸 Evento ##{evento_id} cerrado (#{estado})") rescue StandardError => e Base.logger.error("🛸 Error cerrando evento web ##{evento_id}: #{e.message}") Base.logger.error(" #{e.backtrace&.first(3)&.join("\n ")}") end # ─── Helpers ───────────────────────────────────────────────── def detectar_nodo_id(nodo) return nodo if nodo.is_a?(Integer) require 'socket' hostname = Socket.gethostname.downcase resultado = nil BitacorasDB::BitacoraDB.with_connection do |db| resultado = db.execute('SELECT id FROM nodos WHERE LOWER(nombre) = $1 LIMIT 1', [hostname]) if resultado.nil? || resultado.ntuples == 0 resultado = db.execute('SELECT id FROM nodos WHERE LOWER(nombre) LIKE $1 LIMIT 1', ["#{hostname}%"]) end end (resultado&.any?) ? resultado.first['id'] : 1 rescue StandardError => e Base.logger.warn("⚠️ Error detectando nodo: #{e.message}, usando default (1)") 1 end end end end