feat(dron): Migración a AtomicDrones con arquitectura DB-First
Implementación del sistema AtomicDrones basado en "Menos es Más": - Refactorización de dron.rb a dispatcher + módulos atómicos - Nuevos módulos: ejecutor, vigilante, sanador, bitacora - DronDB para acceso a datos en PostgreSQL - Tablas: dron_logs, dron_avances, dron_metricas - Views: vw_dron_estado, vw_dron_zombies, vw_dron_metricas_semanal - Health check con detección de zombies (heartbeat > 5min) - Dashboard compacto de la flota - Saneamiento con reintentos (backoff exponencial) y limpieza Comandos soportados: - dron lanzar: Ejecutar comando como dron - dron flota: Dashboard de la flota - dron salud: Health check activo - dron estado <ID>: Estado detallado - dron sanear: Saneamiento completo - dron limpiar: Limpieza de antiguos - dron reintentar: Re-intentar fallidos Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>
This commit is contained in:
co-authored by
Claude Opus 4.6
parent
b5e89144e3
commit
3543e62263
@@ -0,0 +1,346 @@
|
||||
# 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) }
|
||||
|
||||
# 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]) }
|
||||
|
||||
# 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]) }
|
||||
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]) }
|
||||
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]) }
|
||||
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
|
||||
Reference in New Issue
Block a user