🐝 Fase 2: Orquestador 1-dron-por-archivo (bkps orquestar)
- proc_xen.rb: listar() scanner + ejecutar_uno() atómico - bkps.rb: --archivo flag + subcomando orquestar - Patrón: scan → lanza drones secuenciales → 1 VM por dron - Filosofía: Menos es Más
This commit is contained in:
+119
-10
@@ -49,6 +49,8 @@ module ADN
|
|||||||
cmd_list
|
cmd_list
|
||||||
when 'run'
|
when 'run'
|
||||||
cmd_run(@args)
|
cmd_run(@args)
|
||||||
|
when 'orquestar'
|
||||||
|
cmd_orquestar(@args)
|
||||||
when 'menu'
|
when 'menu'
|
||||||
cmd_menu
|
cmd_menu
|
||||||
when 'status'
|
when 'status'
|
||||||
@@ -199,6 +201,7 @@ module ADN
|
|||||||
opts.banner = "Uso: ./adn/tools/run bkps run <T1|C1|ID> [opciones]"
|
opts.banner = "Uso: ./adn/tools/run bkps run <T1|C1|ID> [opciones]"
|
||||||
opts.on("--batch", "Modo batch (sin interacción, para cron)") { options[:batch] = true }
|
opts.on("--batch", "Modo batch (sin interacción, para cron)") { options[:batch] = true }
|
||||||
opts.on("--dry-run", "Simulación sin ejecutar") { options[:dry_run] = true }
|
opts.on("--dry-run", "Simulación sin ejecutar") { options[:dry_run] = true }
|
||||||
|
opts.on("--archivo NOMBRE", "Procesar solo este archivo (1-dron-por-archivo)") { |a| options[:archivo] = a }
|
||||||
end.parse!(args)
|
end.parse!(args)
|
||||||
|
|
||||||
ref = args.shift
|
ref = args.shift
|
||||||
@@ -359,6 +362,7 @@ module ADN
|
|||||||
options = { dias: retencion['dias_default'] || 6 }
|
options = { dias: retencion['dias_default'] || 6 }
|
||||||
OptionParser.new do |opts|
|
OptionParser.new do |opts|
|
||||||
opts.banner = "Uso: ./adn/tools/run bkps sanear [opciones]"
|
opts.banner = "Uso: ./adn/tools/run bkps sanear [opciones]"
|
||||||
|
opts.on("--batch", "Modo batch (sin confirmación)") { options[:batch] = true }
|
||||||
opts.on("--dias N", Integer, "Días de retención (default: #{options[:dias]})") { |d| options[:dias] = d }
|
opts.on("--dias N", Integer, "Días de retención (default: #{options[:dias]})") { |d| options[:dias] = d }
|
||||||
opts.on("--tier N", Integer, "Solo tier N (1=proxmox, 2=local, 3=nube)") { |t| options[:tier] = t }
|
opts.on("--tier N", Integer, "Solo tier N (1=proxmox, 2=local, 3=nube)") { |t| options[:tier] = t }
|
||||||
opts.on("--dry-run", "Solo listar, no borrar") { options[:dry_run] = true }
|
opts.on("--dry-run", "Solo listar, no borrar") { options[:dry_run] = true }
|
||||||
@@ -408,14 +412,15 @@ module ADN
|
|||||||
storage = tier['storage']
|
storage = tier['storage']
|
||||||
hosts = tier['hosts'] || []
|
hosts = tier['hosts'] || []
|
||||||
ruta_dump = "/mnt/pve/#{storage}/dump"
|
ruta_dump = "/mnt/pve/#{storage}/dump"
|
||||||
|
candado_clave = tier['candado'] || 'srv-dasu:root'
|
||||||
|
|
||||||
system("ruby #{candados_path} authorize > /dev/null 2>&1")
|
system("ruby #{candados_path} authorize > /dev/null 2>&1")
|
||||||
|
|
||||||
output_ls = ""
|
output_ls = ""
|
||||||
hosts.each do |h|
|
hosts.each do |h|
|
||||||
cmd_ls = "ls -1 --full-time #{ruta_dump}/vzdump-*"
|
cmd_ls = "ls -1 --full-time #{ruta_dump}/vzdump-*"
|
||||||
res = `ruby #{candados_path} run admindasu SSHPASS 'sshpass -e ssh -o StrictHostKeyChecking=no root@#{h['ip']} "#{cmd_ls}"' 2>&1`
|
res = `ruby #{candados_path} run #{candado_clave} SSHPASS 'sshpass -e ssh -o StrictHostKeyChecking=no root@#{h['ip']} "#{cmd_ls}"' 2>&1`
|
||||||
unless res.include?('No such file') || res.include?('Connection refused')
|
if $?.success? || res.include?('No such file')
|
||||||
output_ls = res
|
output_ls = res
|
||||||
break
|
break
|
||||||
end
|
end
|
||||||
@@ -455,9 +460,10 @@ module ADN
|
|||||||
|
|
||||||
# Borrar via SSH
|
# Borrar via SSH
|
||||||
if elim > 0 && !options[:dry_run]
|
if elim > 0 && !options[:dry_run]
|
||||||
todos.select { |fp, _| !protegido?(fp, por_nodo, todos) && todos[fp][:fecha] < limite }.each_value do |info|
|
candidatos = calcular_candidatos(todos, por_nodo, limite)
|
||||||
|
candidatos.each_value do |info|
|
||||||
info[:archivos].each do |ruta|
|
info[:archivos].each do |ruta|
|
||||||
system("ruby #{candados_path} run admindasu SSHPASS 'sshpass -e ssh -o StrictHostKeyChecking=no root@#{hosts.first['ip']} \"rm -f '#{ruta}'\"' > /dev/null 2>&1")
|
system("ruby #{candados_path} run #{candado_clave} SSHPASS 'sshpass -e ssh -o StrictHostKeyChecking=no root@#{hosts.first['ip']} \"rm -f '#{ruta}'\"' > /dev/null 2>&1")
|
||||||
end
|
end
|
||||||
end
|
end
|
||||||
end
|
end
|
||||||
@@ -598,11 +604,15 @@ module ADN
|
|||||||
return [0, protegidos]
|
return [0, protegidos]
|
||||||
end
|
end
|
||||||
|
|
||||||
print "\n¿Borrar todos? (S/N): "
|
if options[:batch]
|
||||||
resp = STDIN.gets&.chomp&.strip&.downcase
|
puts "\n#{Color::YELLOW}[BATCH] Borrando automáticamente sin pedir confirmación.#{Color::RESET}"
|
||||||
unless resp == 's'
|
else
|
||||||
puts "Cancelado."
|
print "\n¿Borrar todos? (S/N): "
|
||||||
return [0, protegidos]
|
resp = STDIN.gets&.chomp&.strip&.downcase
|
||||||
|
unless resp == 's'
|
||||||
|
puts "Cancelado."
|
||||||
|
return [0, protegidos]
|
||||||
|
end
|
||||||
end
|
end
|
||||||
|
|
||||||
lista.each_with_index do |b, i|
|
lista.each_with_index do |b, i|
|
||||||
@@ -738,7 +748,15 @@ module ADN
|
|||||||
t_inicio = Time.now
|
t_inicio = Time.now
|
||||||
@logger.info("▶ #{tarea[:texto]} (#{tipo})")
|
@logger.info("▶ #{tarea[:texto]} (#{tipo})")
|
||||||
procesador = factory.call(@logger)
|
procesador = factory.call(@logger)
|
||||||
resultado = procesador.ejecutar(tarea)
|
|
||||||
|
# Modo atómico: procesar un solo archivo
|
||||||
|
resultado = if options[:archivo] && procesador.respond_to?(:ejecutar_uno)
|
||||||
|
@logger.info("🎯 Modo atómico: #{options[:archivo]}")
|
||||||
|
procesador.ejecutar_uno(tarea, options[:archivo])
|
||||||
|
else
|
||||||
|
procesador.ejecutar(tarea)
|
||||||
|
end
|
||||||
|
|
||||||
duracion = (Time.now - t_inicio).to_i
|
duracion = (Time.now - t_inicio).to_i
|
||||||
mins = duracion / 60
|
mins = duracion / 60
|
||||||
segs = duracion % 60
|
segs = duracion % 60
|
||||||
@@ -746,6 +764,97 @@ module ADN
|
|||||||
resultado
|
resultado
|
||||||
end
|
end
|
||||||
|
|
||||||
|
# ─── Orquestador: 1 dron por archivo ──────────────────────────
|
||||||
|
|
||||||
|
def cmd_orquestar(args)
|
||||||
|
options = {}
|
||||||
|
OptionParser.new do |opts|
|
||||||
|
opts.banner = "Uso: ./adn/tools/run bkps orquestar <T1|ID> [opciones]"
|
||||||
|
opts.on("--timeout N", Integer, "Timeout por dron en segundos (default: 7200)") { |t| options[:timeout] = t }
|
||||||
|
opts.on("--dry-run", "Solo listar archivos sin lanzar drones") { options[:dry_run] = true }
|
||||||
|
end.parse!(args)
|
||||||
|
|
||||||
|
ref = args.shift
|
||||||
|
unless ref
|
||||||
|
puts "#{Color::RED}✗ Falta referencia de tarea (T1, T2, o ID)#{Color::RESET}"
|
||||||
|
return
|
||||||
|
end
|
||||||
|
|
||||||
|
# Resolver tarea
|
||||||
|
tarea = case ref.upcase
|
||||||
|
when /^T(\d+)$/ then @config.tareas[$1.to_i - 1]
|
||||||
|
else @config.find_tarea(ref)
|
||||||
|
end
|
||||||
|
|
||||||
|
unless tarea
|
||||||
|
puts "#{Color::RED}✗ Tarea '#{ref}' no encontrada#{Color::RESET}"
|
||||||
|
return
|
||||||
|
end
|
||||||
|
|
||||||
|
tipo = tarea[:tipo].to_s
|
||||||
|
factory = TIPOS[tipo]
|
||||||
|
unless factory
|
||||||
|
puts "#{Color::RED}✗ Tipo '#{tipo}' no soporta orquestación#{Color::RESET}"
|
||||||
|
return
|
||||||
|
end
|
||||||
|
|
||||||
|
procesador = factory.call(@logger)
|
||||||
|
unless procesador.respond_to?(:listar)
|
||||||
|
puts "#{Color::RED}✗ Procesador '#{tipo}' no soporta listar (orquestación no disponible)#{Color::RESET}"
|
||||||
|
return
|
||||||
|
end
|
||||||
|
|
||||||
|
# Escanear archivos pendientes
|
||||||
|
archivos = procesador.listar(tarea)
|
||||||
|
pendientes = archivos.select { |a| a[:pendiente] }
|
||||||
|
|
||||||
|
puts "\n#{Color::BOLD}#{Color::CYAN}🐝 Orquestador AtomicDrones — #{tarea[:texto]}#{Color::RESET}"
|
||||||
|
puts "═" * 60
|
||||||
|
puts " Total archivos: #{archivos.length}"
|
||||||
|
puts " Pendientes: #{pendientes.length}"
|
||||||
|
puts " Ya procesados: #{archivos.length - pendientes.length}"
|
||||||
|
puts ""
|
||||||
|
|
||||||
|
if pendientes.empty?
|
||||||
|
puts "#{Color::GREEN}✔ Todos los archivos ya están procesados.#{Color::RESET}"
|
||||||
|
return
|
||||||
|
end
|
||||||
|
|
||||||
|
adn_run = File.join(ADN::PROJECT_ROOT, 'adn', 'tools', 'run')
|
||||||
|
timeout = options[:timeout] || 7200
|
||||||
|
|
||||||
|
pendientes.each_with_index do |arch, i|
|
||||||
|
tamaño_gb = (arch[:tamaño].to_f / 1024**3).round(1)
|
||||||
|
puts " #{Color::BOLD}[#{i+1}/#{pendientes.length}]#{Color::RESET} #{arch[:vm]} — #{arch[:archivo]} (#{tamaño_gb}GB)"
|
||||||
|
|
||||||
|
if options[:dry_run]
|
||||||
|
puts " #{Color::YELLOW}[DRY-RUN] Lanzaría dron para: #{arch[:archivo]}#{Color::RESET}"
|
||||||
|
next
|
||||||
|
end
|
||||||
|
|
||||||
|
# Lanzar dron atómico para este archivo
|
||||||
|
nota = "Comprimir #{arch[:vm]} (#{tamaño_gb}GB)"
|
||||||
|
cmd_dron = "#{adn_run} dron lanzar" \
|
||||||
|
" --nota \"#{nota}\"" \
|
||||||
|
" --timeout #{timeout}" \
|
||||||
|
" -- #{adn_run} bkps run #{ref} --batch --archivo #{Shellwords.escape(arch[:archivo])}"
|
||||||
|
|
||||||
|
puts " 🛸 Lanzando dron..."
|
||||||
|
resultado = system(cmd_dron)
|
||||||
|
|
||||||
|
if resultado
|
||||||
|
puts " #{Color::GREEN}✔ Dron completado#{Color::RESET}"
|
||||||
|
else
|
||||||
|
puts " #{Color::RED}✖ Dron falló — deteniendo orquestación#{Color::RESET}"
|
||||||
|
break
|
||||||
|
end
|
||||||
|
|
||||||
|
puts ""
|
||||||
|
end
|
||||||
|
|
||||||
|
puts "#{Color::BOLD}🐝 Orquestación finalizada.#{Color::RESET}"
|
||||||
|
end
|
||||||
|
|
||||||
def ejecutar_comando(comando, options = {})
|
def ejecutar_comando(comando, options = {})
|
||||||
t_inicio = Time.now
|
t_inicio = Time.now
|
||||||
@logger.info("=== #{comando[:texto]} ===")
|
@logger.info("=== #{comando[:texto]} ===")
|
||||||
|
|||||||
@@ -3,6 +3,11 @@
|
|||||||
# ==========================================================
|
# ==========================================================
|
||||||
# adn/tools/bkps/lib/proc_xen.rb — Procesador XEN (.xva)
|
# adn/tools/bkps/lib/proc_xen.rb — Procesador XEN (.xva)
|
||||||
# Migrado desde dtic-BKPs v6.0 → ADN::BKPs
|
# Migrado desde dtic-BKPs v6.0 → ADN::BKPs
|
||||||
|
#
|
||||||
|
# Fase 2: Soporta orquestación 1-dron-por-archivo
|
||||||
|
# listar(tarea) → lista archivos pendientes
|
||||||
|
# ejecutar_uno(tarea, arch) → comprime UN solo archivo
|
||||||
|
# ejecutar(tarea) → comprime TODOS (legacy)
|
||||||
# ==========================================================
|
# ==========================================================
|
||||||
|
|
||||||
require 'shellwords'
|
require 'shellwords'
|
||||||
@@ -18,6 +23,51 @@ module ADN
|
|||||||
@pv = system('which pv > /dev/null 2>&1')
|
@pv = system('which pv > /dev/null 2>&1')
|
||||||
end
|
end
|
||||||
|
|
||||||
|
# Scanner: devuelve lista de archivos .xva pendientes de procesar
|
||||||
|
def listar(tarea)
|
||||||
|
archivos = Dir.glob(File.join(tarea[:origen], '*.xva'))
|
||||||
|
archivos.map do |ruta_xva|
|
||||||
|
nombre_xva = File.basename(ruta_xva)
|
||||||
|
nombre_vm = nombre_xva.split('_').first
|
||||||
|
dir_dest = File.join(tarea[:destino], nombre_vm)
|
||||||
|
nombre_sin_ext = File.basename(nombre_xva, '.xva')
|
||||||
|
ruta_tar = File.join(dir_dest, "#{nombre_sin_ext}.tar.gz")
|
||||||
|
|
||||||
|
{
|
||||||
|
archivo: nombre_xva,
|
||||||
|
ruta: ruta_xva,
|
||||||
|
vm: nombre_vm,
|
||||||
|
destino: ruta_tar,
|
||||||
|
tamaño: File.size(ruta_xva),
|
||||||
|
pendiente: !File.exist?(ruta_tar) || tarea[:sobrescribir]
|
||||||
|
}
|
||||||
|
end
|
||||||
|
end
|
||||||
|
|
||||||
|
# Atómico: comprime UN solo archivo .xva
|
||||||
|
def ejecutar_uno(tarea, nombre_archivo)
|
||||||
|
ruta_xva = File.join(tarea[:origen], nombre_archivo)
|
||||||
|
unless File.exist?(ruta_xva)
|
||||||
|
@log.error("Archivo no encontrado: #{nombre_archivo}")
|
||||||
|
return false
|
||||||
|
end
|
||||||
|
|
||||||
|
nombre_vm = nombre_archivo.split('_').first
|
||||||
|
dir_dest = File.join(tarea[:destino], nombre_vm)
|
||||||
|
FileUtils.mkdir_p(dir_dest)
|
||||||
|
|
||||||
|
nombre_sin_ext = File.basename(nombre_archivo, '.xva')
|
||||||
|
ruta_tar = File.join(dir_dest, "#{nombre_sin_ext}.tar.gz")
|
||||||
|
|
||||||
|
if File.exist?(ruta_tar) && !tarea[:sobrescribir]
|
||||||
|
@log.info("#{File.basename(ruta_tar)} existe. Saltando.")
|
||||||
|
return true
|
||||||
|
end
|
||||||
|
|
||||||
|
comprimir(tarea, ruta_xva, nombre_archivo, ruta_tar)
|
||||||
|
end
|
||||||
|
|
||||||
|
# Legacy: procesa TODOS los archivos (compatibilidad con sistema previo)
|
||||||
def ejecutar(tarea)
|
def ejecutar(tarea)
|
||||||
@log.info("--- #{tarea[:texto]} ---")
|
@log.info("--- #{tarea[:texto]} ---")
|
||||||
@log.info("#{tarea[:id]}: Origen: #{tarea[:origen]}")
|
@log.info("#{tarea[:id]}: Origen: #{tarea[:origen]}")
|
||||||
@@ -49,39 +99,44 @@ module ADN
|
|||||||
end
|
end
|
||||||
end
|
end
|
||||||
|
|
||||||
begin
|
procesados += 1 if comprimir(tarea, ruta_xva, nombre_xva, ruta_tar)
|
||||||
@log.info("Comprimiendo #{nombre_xva}")
|
|
||||||
|
|
||||||
if @pv
|
|
||||||
cmd_tar = ['tar', '-czf', '-', '-C', tarea[:origen], nombre_xva]
|
|
||||||
cmd_pv = ['pv', '-s', File.size(ruta_xva).to_s]
|
|
||||||
File.open(ruta_tar, 'wb') do |f|
|
|
||||||
Open3.pipeline(cmd_tar, cmd_pv, out: f).each_with_index do |s, i|
|
|
||||||
raise "Falló #{%w[tar pv][i]} (#{s.exitstatus})" unless s.success?
|
|
||||||
end
|
|
||||||
end
|
|
||||||
else
|
|
||||||
cmd = "tar -czf #{Shellwords.escape(ruta_tar)} -C #{Shellwords.escape(tarea[:origen])} #{Shellwords.escape(nombre_xva)}"
|
|
||||||
system(cmd) or raise "tar falló (#{$?.exitstatus})"
|
|
||||||
end
|
|
||||||
|
|
||||||
@log.info("✔ Compresión OK: #{File.basename(ruta_tar)}")
|
|
||||||
|
|
||||||
if tarea[:eliminar_origen]
|
|
||||||
FileUtils.rm_f(ruta_xva)
|
|
||||||
@log.info("XVA original eliminado.")
|
|
||||||
end
|
|
||||||
|
|
||||||
procesados += 1
|
|
||||||
rescue => e
|
|
||||||
@log.error("#{tarea[:id]}: Error comprimiendo #{nombre_xva}: #{e.message}")
|
|
||||||
FileUtils.rm_f(ruta_tar)
|
|
||||||
end
|
|
||||||
end
|
end
|
||||||
|
|
||||||
@log.info("✔ #{procesados} archivos .xva procesados.")
|
@log.info("✔ #{procesados} archivos .xva procesados.")
|
||||||
true
|
true
|
||||||
end
|
end
|
||||||
|
|
||||||
|
private
|
||||||
|
|
||||||
|
def comprimir(tarea, ruta_xva, nombre_xva, ruta_tar)
|
||||||
|
@log.info("Comprimiendo #{nombre_xva}")
|
||||||
|
|
||||||
|
if @pv
|
||||||
|
cmd_tar = ['tar', '-czf', '-', '-C', tarea[:origen], nombre_xva]
|
||||||
|
cmd_pv = ['pv', '-s', File.size(ruta_xva).to_s]
|
||||||
|
File.open(ruta_tar, 'wb') do |f|
|
||||||
|
Open3.pipeline(cmd_tar, cmd_pv, out: f).each_with_index do |s, i|
|
||||||
|
raise "Falló #{%w[tar pv][i]} (#{s.exitstatus})" unless s.success?
|
||||||
|
end
|
||||||
|
end
|
||||||
|
else
|
||||||
|
cmd = "tar -czf #{Shellwords.escape(ruta_tar)} -C #{Shellwords.escape(tarea[:origen])} #{Shellwords.escape(nombre_xva)}"
|
||||||
|
system(cmd) or raise "tar falló (#{$?.exitstatus})"
|
||||||
|
end
|
||||||
|
|
||||||
|
@log.info("✔ Compresión OK: #{File.basename(ruta_tar)}")
|
||||||
|
|
||||||
|
if tarea[:eliminar_origen]
|
||||||
|
FileUtils.rm_f(ruta_xva)
|
||||||
|
@log.info("XVA original eliminado.")
|
||||||
|
end
|
||||||
|
|
||||||
|
true
|
||||||
|
rescue => e
|
||||||
|
@log.error("#{tarea[:id]}: Error comprimiendo #{nombre_xva}: #{e.message}")
|
||||||
|
FileUtils.rm_f(ruta_tar)
|
||||||
|
false
|
||||||
|
end
|
||||||
end
|
end
|
||||||
end
|
end
|
||||||
end
|
end
|
||||||
|
|||||||
Reference in New Issue
Block a user