Cambios: - ejecutor.rb: Corregir lectura de output con Open3 - bitacora_db.rb: Agregar métodos crear_evento/actualizar_evento (public) - dron_db.rb: Corregir conversión de PG::Result a Hash (.first.to_h) - vigilante.rb: Corregir formato de fecha en dashboard - migrations/004: Crear tabla bitacoras.events para Bitácora Web Pruebas: - dron lanzar --evento AUTO --nota "Test" -- echo "hola" ✅ - dron flota ✅ - dron salud ✅ Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>
347 lines
12 KiB
Ruby
347 lines
12 KiB
Ruby
# frozen_string_literal: true
|
|
|
|
require_relative 'bitacora_db'
|
|
require 'securerandom'
|
|
|
|
module ADN
|
|
module DB
|
|
# DronDB - Acceso a datos para el sistema de AtomicDrones
|
|
# Gestiona registro y consulta de drones en tablas internas
|
|
class DronDB
|
|
# Estados posibles de un dron
|
|
ESTADOS = {
|
|
pending: 'pending',
|
|
running: 'running',
|
|
completed: 'completed',
|
|
failed: 'failed',
|
|
zombie: 'zombie'
|
|
}.freeze
|
|
|
|
# Tipos de drones
|
|
TIPOS = {
|
|
ejecutor: 'ejecutor',
|
|
vigilante: 'vigilante',
|
|
sanador: 'sanador',
|
|
planificador: 'planificador',
|
|
mensajero: 'mensajero',
|
|
orquestador: 'orquestador'
|
|
}.freeze
|
|
|
|
# Tiempo máximo sin heartbeat antes de considerar zombie (5 minutos)
|
|
ZOMBIE_TIMEOUT = 300
|
|
|
|
class << self
|
|
# Generar ID único para un dron
|
|
# Formato: dron_HHMMSS_PID_RANDOM
|
|
def generar_dron_id
|
|
timestamp = Time.now.strftime('%H%M%S')
|
|
pid = Process.pid
|
|
random = SecureRandom.hex(3)
|
|
"dron_#{timestamp}_#{pid}_#{random}"
|
|
end
|
|
|
|
# Registrar inicio de un dron
|
|
# @param tipo [String] Tipo de dron (ejecutor, vigilante, etc.)
|
|
# @param cmd [String] Comando a ejecutar
|
|
# @param flujo_id [String, nil] ID de orquestación (si aplica)
|
|
# @param metadata [Hash] Metadata adicional
|
|
# @return [Hash] Datos del dron registrado
|
|
def registrar_inicio(tipo:, cmd:, flujo_id: nil, metadata: {})
|
|
sql = <<-SQL
|
|
INSERT INTO bitacoras.dron_logs (dron_id, tipo, estado, cmd, flujo_id, metadata)
|
|
VALUES ($1, $2, $3, $4, $5, $6)
|
|
RETURNING *
|
|
SQL
|
|
|
|
dron_id = generar_dron_id
|
|
params = [dron_id, tipo, ESTADOS[:running], cmd, flujo_id, metadata.to_json]
|
|
|
|
resultado = BitacorasDB::BitacoraDB.with_connection { |db| db.execute(sql, params) }.first.to_h
|
|
|
|
# Actualizar heartbeat inicial
|
|
actualizar_heartbeat(dron_id)
|
|
|
|
{
|
|
dron_id: dron_id,
|
|
estado: ESTADOS[:running],
|
|
started_at: resultado['started_at'],
|
|
mensaje: "Dron #{dron_id} registrado como #{tipo}"
|
|
}
|
|
end
|
|
|
|
# Registrar finalización de un dron
|
|
# @param dron_id [String] ID del dron
|
|
# @param exit_code [Integer] Código de salida
|
|
# @param output [String] Salida del comando
|
|
# @return [Hash] Resultado de la operación
|
|
def registrar_fin(dron_id, exit_code, output = nil)
|
|
estado = exit_code == 0 ? ESTADOS[:completed] : ESTADOS[:failed]
|
|
|
|
sql = <<-SQL
|
|
UPDATE bitacoras.dron_logs
|
|
SET estado = $1,
|
|
exit_code = $2,
|
|
output = $3,
|
|
completed_at = NOW(),
|
|
updated_at = NOW()
|
|
WHERE dron_id = $4
|
|
RETURNING *
|
|
SQL
|
|
|
|
resultado = BitacorasDB::BitacoraDB.with_connection { |db| db.execute(sql, [estado, exit_code, output, dron_id]) }.first.to_h
|
|
|
|
# Actualizar métricas diarias
|
|
actualizar_metricas(dron_id, estado)
|
|
|
|
{
|
|
dron_id: dron_id,
|
|
estado: estado,
|
|
exit_code: exit_code,
|
|
completed_at: resultado['completed_at']
|
|
}
|
|
end
|
|
|
|
# Actualizar heartbeat de un dron
|
|
# @param dron_id [String] ID del dron
|
|
def actualizar_heartbeat(dron_id)
|
|
sql = <<-SQL
|
|
UPDATE bitacoras.dron_logs
|
|
SET heartbeat = NOW(), updated_at = NOW()
|
|
WHERE dron_id = $1
|
|
SQL
|
|
|
|
BitacorasDB::BitacoraDB.with_connection { |db| db.execute(sql, [dron_id]) }
|
|
end
|
|
|
|
# Obtener información de un dron
|
|
# @param dron_id [String] ID del dron
|
|
# @return [Hash, nil] Datos del dron o nil si no existe
|
|
def obtener_dron(dron_id)
|
|
sql = 'SELECT * FROM bitacoras.dron_logs WHERE dron_id = $1'
|
|
resultado = BitacorasDB::BitacoraDB.with_connection { |db| db.execute(sql, [dron_id]) }.first.to_h
|
|
resultado.first
|
|
end
|
|
|
|
# Obtener drones por estado
|
|
# @param estado [String] Estado a filtrar
|
|
# @return [Array<Hash>] Lista de drones
|
|
def obtener_por_estado(estado)
|
|
sql = 'SELECT * FROM bitacoras.dron_logs WHERE estado = $1 ORDER BY started_at DESC'
|
|
BitacorasDB::BitacoraDB.with_connection { |db| db.execute(sql, [estado]) }
|
|
end
|
|
|
|
# Obtener todos los drones activos (running o pending)
|
|
# @return [Array<Hash>] Lista de drones activos
|
|
def obtener_activos
|
|
sql = <<-SQL
|
|
SELECT * FROM bitacoras.dron_logs
|
|
WHERE estado IN ('running', 'pending')
|
|
ORDER BY started_at DESC
|
|
SQL
|
|
|
|
BitacorasDB::BitacoraDB.with_connection { |db| db.execute(sql).to_a }
|
|
end
|
|
|
|
# Detectar drones zombies (heartbeat > 5min)
|
|
# @return [Array<Hash>] Lista de drones zombies
|
|
def detectar_zombies
|
|
sql = <<-SQL
|
|
SELECT * FROM bitacoras.vw_dron_zombies
|
|
ORDER BY heartbeat ASC
|
|
SQL
|
|
|
|
BitacorasDB::BitacoraDB.with_connection { |db| db.execute(sql).to_a }
|
|
end
|
|
|
|
# Marcar dron como zombie
|
|
# @param dron_id [String] ID del dron
|
|
def marcar_zombie(dron_id)
|
|
sql = <<-SQL
|
|
UPDATE bitacoras.dron_logs
|
|
SET estado = $1, updated_at = NOW()
|
|
WHERE dron_id = $2
|
|
SQL
|
|
|
|
BitacorasDB::BitacoraDB.with_connection { |db| db.execute(sql, [ESTADOS[:zombie], dron_id]) }
|
|
end
|
|
|
|
# Incrementar contador de reintentos
|
|
# @param dron_id [String] ID del dron
|
|
# @return [Integer] Nuevo contador de reintentos
|
|
def incrementar_reintentos(dron_id)
|
|
sql = <<-SQL
|
|
UPDATE bitacoras.dron_logs
|
|
SET reintentos = reintentos + 1, updated_at = NOW()
|
|
WHERE dron_id = $1
|
|
RETURNING reintentos
|
|
SQL
|
|
|
|
resultado = BitacorasDB::BitacoraDB.with_connection { |db| db.execute(sql, [dron_id]) }.first.to_h
|
|
resultado.first['reintentos'].to_i
|
|
end
|
|
|
|
# Registrar avance de un paso en flujo multi-paso
|
|
# @param dron_id [String] ID del dron
|
|
# @param paso [Integer] Número de paso
|
|
# @param descripcion [String] Descripción del paso
|
|
# @param estado [String] Estado del paso
|
|
# @param output [String, nil] Salida del paso
|
|
# @param flujo_id [String, nil] ID de orquestación
|
|
def registrar_avance(dron_id, paso:, descripcion:, estado:, output: nil, flujo_id: nil)
|
|
sql = <<-SQL
|
|
INSERT INTO bitacoras.dron_avances (dron_id, flujo_id, paso, descripcion, estado, output)
|
|
VALUES ($1, $2, $3, $4, $5, $6)
|
|
RETURNING *
|
|
SQL
|
|
|
|
BitacorasDB::BitacoraDB.with_connection { |db| db.execute(sql, [dron_id, flujo_id, paso, descripcion, estado, output]) }
|
|
end
|
|
|
|
# Obtener avances de un dron
|
|
# @param dron_id [String] ID del dron
|
|
# @return [Array<Hash>] Lista de avances
|
|
def obtener_avances(dron_id)
|
|
sql = <<-SQL
|
|
SELECT * FROM bitacoras.dron_avances
|
|
WHERE dron_id = $1
|
|
ORDER BY paso ASC
|
|
SQL
|
|
|
|
BitacorasDB::BitacoraDB.with_connection { |db| db.execute(sql, [dron_id]) }
|
|
end
|
|
|
|
# Obtener avances de un flujo
|
|
# @param flujo_id [String] ID del flujo
|
|
# @return [Array<Hash>] Lista de avances
|
|
def obtener_avances_flujo(flujo_id)
|
|
sql = <<-SQL
|
|
SELECT * FROM bitacoras.dron_avances
|
|
WHERE flujo_id = $1
|
|
ORDER BY paso ASC
|
|
SQL
|
|
|
|
BitacorasDB::BitacoraDB.with_connection { |db| db.execute(sql, [flujo_id]) }
|
|
end
|
|
|
|
# Obtener métricas por fecha y tipo
|
|
# @param fecha [Date] Fecha de las métricas
|
|
# @param tipo [String] Tipo de dron
|
|
# @return [Hash, nil] Métricas o nil si no existen
|
|
def obtener_metricas(fecha:, tipo:)
|
|
sql = 'SELECT * FROM bitacoras.dron_metricas WHERE fecha = $1 AND dron_tipo = $2'
|
|
resultado = BitacorasDB::BitacoraDB.with_connection { |db| db.execute(sql, [fecha, tipo]) }.first.to_h
|
|
resultado.first
|
|
end
|
|
|
|
# Actualizar métricas diarias (llamado automáticamente al finalizar un dron)
|
|
# @param dron_id [String] ID del dron
|
|
# @param estado [String] Estado final del dron
|
|
def actualizar_metricas(dron_id, estado)
|
|
dron = obtener_dron(dron_id)
|
|
return unless dron
|
|
|
|
fecha = dron['started_at'].to_date
|
|
tipo = dron['tipo']
|
|
|
|
# Calcular duración
|
|
started = dron['started_at']
|
|
completed = dron['completed_at'] || Time.now
|
|
duracion = completed - started
|
|
|
|
# Upsert de métricas
|
|
sql = <<-SQL
|
|
INSERT INTO bitacoras.dron_metricas (fecha, dron_tipo, total_ejecuciones, completados, fallidos, zombies, tiempo_promedio, reintentos_total)
|
|
VALUES ($1, $2, 1, $3, $4, $5, $6, $7)
|
|
ON CONFLICT (fecha, dron_tipo) DO UPDATE SET
|
|
total_ejecuciones = bitacoras.dron_metricas.total_ejecuciones + 1,
|
|
completados = bitacoras.dron_metricas.completados + $3,
|
|
fallidos = bitacoras.dron_metricas.fallidos + $4,
|
|
zombies = bitacoras.dron_metricas.zombies + $5,
|
|
tiempo_promedio = (
|
|
(bitacoras.dron_metricas.tiempo_promedio * bitacoras.dron_metricas.total_ejecuciones + $6) /
|
|
(bitacoras.dron_metricas.total_ejecuciones + 1)
|
|
),
|
|
reintentos_total = bitacoras.dron_metricas.reintentos_total + $7,
|
|
updated_at = NOW()
|
|
SQL
|
|
|
|
es_completado = estado == ESTADOS[:completed] ? 1 : 0
|
|
es_fallido = estado == ESTADOS[:failed] ? 1 : 0
|
|
es_zombie = estado == ESTADOS[:zombie] ? 1 : 0
|
|
reintentos = dron['reintentos'].to_i
|
|
|
|
BitacorasDB::BitacoraDB.with_connection { |db| db.execute(sql, [fecha, tipo, es_completado, es_fallido, es_zombie, duracion, reintentos]) }
|
|
end
|
|
|
|
# Obtener resumen de estado de la flota
|
|
# @return [Hash] Resumen de la flota
|
|
def resumen_flota
|
|
sql = <<-SQL
|
|
SELECT
|
|
estado,
|
|
COUNT(*) as cantidad
|
|
FROM bitacoras.dron_logs
|
|
GROUP BY estado
|
|
SQL
|
|
|
|
resultados = BitacorasDB::BitacoraDB.with_connection { |db| db.execute(sql).to_a }
|
|
resumen = { total: 0 }
|
|
|
|
resultados.each do |fila|
|
|
resumen[fila['estado'].to_sym] = fila['cantidad'].to_i
|
|
resumen[:total] += fila['cantidad'].to_i
|
|
end
|
|
|
|
resumen
|
|
end
|
|
|
|
# Obtener drones por flujo
|
|
# @param flujo_id [String] ID del flujo
|
|
# @return [Array<Hash>] Lista de drones del flujo
|
|
def obtener_por_flujo(flujo_id)
|
|
sql = <<-SQL
|
|
SELECT * FROM bitacoras.dron_logs
|
|
WHERE flujo_id = $1
|
|
ORDER BY started_at ASC
|
|
SQL
|
|
|
|
BitacorasDB::BitacoraDB.with_connection { |db| db.execute(sql, [flujo_id]) }
|
|
end
|
|
|
|
# Consolidar estado de un flujo
|
|
# @param flujo_id [String] ID del flujo
|
|
# @return [String] Estado consolidado (running, completed, failed)
|
|
def consolidar_estado_flujo(flujo_id)
|
|
drones = obtener_por_flujo(flujo_id)
|
|
return nil if drones.empty?
|
|
|
|
# Si todos completados → completed
|
|
# Si alguno failed → failed
|
|
# Si no → running
|
|
if drones.all? { |d| d['estado'] == ESTADOS[:completed] }
|
|
ESTADOS[:completed]
|
|
elsif drones.any? { |d| d['estado'] == ESTADOS[:failed] }
|
|
ESTADOS[:failed]
|
|
else
|
|
ESTADOS[:running]
|
|
end
|
|
end
|
|
|
|
# Limpiar drones antiguos (más de N días)
|
|
# @param dias [Integer] Días de antigüedad
|
|
# @return [Integer] Cantidad de drones eliminados
|
|
def limpiar_antiguos(dias: 7)
|
|
sql = <<-SQL
|
|
DELETE FROM bitacoras.dron_logs
|
|
WHERE estado IN ('completed', 'failed', 'zombie')
|
|
AND completed_at < NOW() - INTERVAL '#{dias} days'
|
|
SQL
|
|
|
|
resultado = BitacorasDB::BitacoraDB.with_connection { |db| db.execute(sql).to_a }
|
|
resultado.first['count'].to_i
|
|
end
|
|
end
|
|
end
|
|
end
|
|
end
|