Se implementó filtrado de output para que la Bitácora Web muestre solo las líneas relevantes del progreso del dron, no los ~40 logs repetitivos de conexión a PostgreSQL. Cambios: - filtrar_lineas_relevantes(): Filtra líneas de PostgreSQL en tiempo real - filtrar_output_relevante(): Filtra y mantiene últimas 20 líneas para heartbeat - Heartbeat web: Ahora muestra solo progreso del dron, no conexiones DB Beneficio: La Bitácora Web ahora muestra información útil y legible en lugar de cientos de líneas de logs técnicos repetitivos. Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>
345 lines
14 KiB
Ruby
345 lines
14 KiB
Ruby
# 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 líneas relevantes para heartbeat web
|
||
exit_code = nil
|
||
progreso_visible = !ENV['BATCH_MODE']
|
||
|
||
# Mostrar encabezado de progreso (solo CLI)
|
||
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 con output filtrado (solo líneas relevantes)
|
||
if evento_id
|
||
output_filtrado = filtrar_output_relevante(accumulated_out)
|
||
actualizar_heartbeat_web(evento_id, dron_id, nota, cmd, inicio, tick, output_filtrado.join("\n"))
|
||
end
|
||
|
||
# Actualizar barra de progreso local (solo CLI)
|
||
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
|
||
# Filtrar y acumular solo líneas relevantes (no PostgreSQL)
|
||
lineas_relevantes = filtrar_lineas_relevantes(lineas)
|
||
accumulated_out.concat(lineas_relevantes) unless lineas_relevantes.empty?
|
||
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
|
||
|
||
# ─── Filtrado de output para Bitácora Web ─────────────────────
|
||
# Filtra líneas repetitivas de PostgreSQL y deja solo lo relevante
|
||
|
||
def filtrar_lineas_relevantes(lineas)
|
||
lineas.reject do |linea|
|
||
# Filtrar logs de conexión a PostgreSQL (son ~40 líneas repetitivas por dron)
|
||
linea.include?('Conectando a PostgreSQL') ||
|
||
linea.include?('Conexión a PostgreSQL establecida') ||
|
||
linea.include?('Conexión a PostgreSQL cerrada') ||
|
||
linea.include?('ℹ Conectando') ||
|
||
linea.include?('ℹ Conexión')
|
||
end
|
||
end
|
||
|
||
def filtrar_output_relevante(accumulated_out)
|
||
# Filtrar y dejar solo las últimas 20 líneas relevantes
|
||
filtrado = accumulated_out.reject do |linea|
|
||
linea.include?('Conectando a PostgreSQL') ||
|
||
linea.include?('Conexión a PostgreSQL') ||
|
||
linea.include?('ℹ Conectando') ||
|
||
linea.include?('ℹ Conexión')
|
||
end
|
||
|
||
# Mantener últimas 20 líneas para no saturar la bitácora
|
||
filtrado.last(20) || []
|
||
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
|