diff --git a/adn/tools/bkps/bkps.rb b/adn/tools/bkps/bkps.rb index 59bc96a0..b33bcd1e 100644 --- a/adn/tools/bkps/bkps.rb +++ b/adn/tools/bkps/bkps.rb @@ -49,6 +49,8 @@ module ADN cmd_list when 'run' cmd_run(@args) + when 'orquestar' + cmd_orquestar(@args) when 'menu' cmd_menu when 'status' @@ -199,6 +201,7 @@ module ADN opts.banner = "Uso: ./adn/tools/run bkps run [opciones]" 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("--archivo NOMBRE", "Procesar solo este archivo (1-dron-por-archivo)") { |a| options[:archivo] = a } end.parse!(args) ref = args.shift @@ -359,6 +362,7 @@ module ADN options = { dias: retencion['dias_default'] || 6 } OptionParser.new do |opts| 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("--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 } @@ -408,14 +412,15 @@ module ADN storage = tier['storage'] hosts = tier['hosts'] || [] ruta_dump = "/mnt/pve/#{storage}/dump" + candado_clave = tier['candado'] || 'srv-dasu:root' system("ruby #{candados_path} authorize > /dev/null 2>&1") output_ls = "" hosts.each do |h| 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` - unless res.include?('No such file') || res.include?('Connection refused') + res = `ruby #{candados_path} run #{candado_clave} SSHPASS 'sshpass -e ssh -o StrictHostKeyChecking=no root@#{h['ip']} "#{cmd_ls}"' 2>&1` + if $?.success? || res.include?('No such file') output_ls = res break end @@ -455,9 +460,10 @@ module ADN # Borrar via SSH 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| - 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 @@ -598,11 +604,15 @@ module ADN return [0, protegidos] end - print "\n¿Borrar todos? (S/N): " - resp = STDIN.gets&.chomp&.strip&.downcase - unless resp == 's' - puts "Cancelado." - return [0, protegidos] + if options[:batch] + puts "\n#{Color::YELLOW}[BATCH] Borrando automáticamente sin pedir confirmación.#{Color::RESET}" + else + print "\n¿Borrar todos? (S/N): " + resp = STDIN.gets&.chomp&.strip&.downcase + unless resp == 's' + puts "Cancelado." + return [0, protegidos] + end end lista.each_with_index do |b, i| @@ -738,7 +748,15 @@ module ADN t_inicio = Time.now @logger.info("▶ #{tarea[:texto]} (#{tipo})") 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 mins = duracion / 60 segs = duracion % 60 @@ -746,6 +764,97 @@ module ADN resultado end + # ─── Orquestador: 1 dron por archivo ────────────────────────── + + def cmd_orquestar(args) + options = {} + OptionParser.new do |opts| + opts.banner = "Uso: ./adn/tools/run bkps orquestar [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 = {}) t_inicio = Time.now @logger.info("=== #{comando[:texto]} ===") diff --git a/adn/tools/bkps/lib/proc_xen.rb b/adn/tools/bkps/lib/proc_xen.rb index b5d1f5ba..0f233542 100644 --- a/adn/tools/bkps/lib/proc_xen.rb +++ b/adn/tools/bkps/lib/proc_xen.rb @@ -3,6 +3,11 @@ # ========================================================== # adn/tools/bkps/lib/proc_xen.rb — Procesador XEN (.xva) # 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' @@ -18,6 +23,51 @@ module ADN @pv = system('which pv > /dev/null 2>&1') 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) @log.info("--- #{tarea[:texto]} ---") @log.info("#{tarea[:id]}: Origen: #{tarea[:origen]}") @@ -49,39 +99,44 @@ module ADN end end - begin - @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 + procesados += 1 if comprimir(tarea, ruta_xva, nombre_xva, ruta_tar) end @log.info("✔ #{procesados} archivos .xva procesados.") true 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