Files
dtic-DIIAA/adn/tools/cli/dron/ejecutor.rb
T
Ricardo MonlaandClaude Opus 4.6 94eeb8d470 [FIX] Heartbeat web: Filtrar PostgreSQL antes de mostrar
Se corrigió actualizar_heartbeat_web() para que filtre el output
antes de actualizar la Bitácora Web.

Antes: Mostraba ~400 líneas con logs de PostgreSQL
Ahora: Muestra solo ~15 líneas de progreso relevante del dron

El filtrado se aplica en dos lugares:
1. Al acumular líneas (filtrar_lineas_relevantes)
2. Al actualizar heartbeat web (filtrar_output_relevante)

Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>
2026-04-15 11:31:50 -03:00

345 lines
14 KiB
Ruby
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
# 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)
# Filtrar output: quitar logs de PostgreSQL y dejar solo progreso relevante
output_limpio = filtrar_output_relevante(accumulated_out.split("\n"))
descripcion = "🛸 #{nota} (`#{dron_id}`)\n" \
"├─ Comando: `#{cmd}`\n" \
"└─ ❤️ #{duracion} — heartbeat ##{tick}"
unless output_limpio.nil? || output_limpio.empty?
# Mostrar últimas 15 líneas de progreso relevante
lineas_mostrar = output_limpio.last(15)
descripcion += "\n\n```text\n#{lineas_mostrar.join("\n")}\n```"
end
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)
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