feat: implementar Fases 3, 4 y 5 del roadmap (Planificacion Avanzada, Historial de Corridas y Copia Local Duplicada)

This commit is contained in:
Carlos Tello
2026-08-15 00:58:05 -03:00
parent 2c2cfeb867
commit 929f9affd7
11 changed files with 677 additions and 57 deletions
+1
View File
@@ -26,6 +26,7 @@ class LocalFolderJob(BaseModel):
skip_system: bool = False
skip_readonly: bool = False
sync_deletions: bool = False
local_destination_path: Optional[str] = None
schedule_cron: str = "daily"
schedule_interval_minutes: int = 60
min_stable_seconds: int = 60
+221 -41
View File
@@ -3,7 +3,9 @@ import socket
import platform
import threading
import logging
from datetime import datetime, timezone
import shutil
import re
from datetime import datetime, timezone, timedelta
from pathlib import Path
from typing import Optional, Callable, Dict, Any, List
import httpx
@@ -27,7 +29,6 @@ def is_job_due(schedule_str: str, last_run_str: Optional[str]) -> bool:
return True
try:
# Try parsing ISO (from server) or standard YYYY-MM-DD HH:MM:SS (local)
if "T" in last_run_str:
last_run = datetime.fromisoformat(last_run_str.replace("Z", "+00:00"))
else:
@@ -36,19 +37,84 @@ def is_job_due(schedule_str: str, last_run_str: Optional[str]) -> bool:
return True
now = datetime.now(timezone.utc)
delta = now - last_run
if last_run >= now:
return False
sched = schedule_str.lower().strip()
# 1. Traditional intervals
if sched == "hourly":
return delta.total_seconds() >= 3600
return (now - last_run).total_seconds() >= 3600
elif sched == "daily":
return delta.total_seconds() >= 86400
return (now - last_run).total_seconds() >= 86400
elif sched == "weekly":
return delta.total_seconds() >= 86400 * 7
return (now - last_run).total_seconds() >= 86400 * 7
elif sched == "monthly":
return delta.total_seconds() >= 86400 * 30
else:
return True
return (now - last_run).total_seconds() >= 86400 * 30
# 2. Parse Proxmox-like calendar string (e.g. "mon..fri 22:00", "sat,sun 18:00", "03:00")
match = re.match(r"^(?:([a-z\.,\s]+)\s+)?(\d{1,2}):(\d{2})$", sched)
if match:
day_spec, hour_str, min_str = match.groups()
target_hour = int(hour_str)
target_min = int(min_str)
allowed_weekdays = set(range(7)) # Default: all days
if day_spec:
day_spec = day_spec.strip()
if day_spec == "mon..fri":
allowed_weekdays = {0, 1, 2, 3, 4}
elif day_spec == "sat..sun" or day_spec == "sat,sun":
allowed_weekdays = {5, 6}
elif any(d in day_spec for d in ["mon", "tue", "wed", "thu", "fri", "sat", "sun"]):
day_map = {"mon": 0, "tue": 1, "wed": 2, "thu": 3, "fri": 4, "sat": 5, "sun": 6}
allowed_weekdays = {day_map[d.strip()] for d in day_spec.split(",") if d.strip() in day_map}
curr = last_run + timedelta(minutes=1)
if (now - curr).days > 7:
curr = now - timedelta(days=7)
while curr <= now:
if curr.hour == target_hour and curr.minute == target_min:
if curr.weekday() in allowed_weekdays:
return True
curr += timedelta(minutes=1)
return False
# 3. Simple 5-field cron parsing (minute hour day_of_month month day_of_week)
fields = sched.split()
if len(fields) == 5:
curr = last_run + timedelta(minutes=1)
if (now - curr).days > 7:
curr = now - timedelta(days=7)
def match_field(val: int, field: str) -> bool:
if field == "*":
return True
if "," in field:
return any(match_field(val, f) for f in field.split(","))
if "-" in field:
start, end = map(int, field.split("-"))
return start <= val <= end
if field.startswith("*/"):
step = int(field[2:])
return val % step == 0
return int(field) == val
while curr <= now:
cron_weekday = (curr.weekday() + 1) % 7 # 0=Sunday, 1=Monday... 6=Saturday
if (match_field(curr.minute, fields[0]) and
match_field(curr.hour, fields[1]) and
match_field(curr.day, fields[2]) and
match_field(curr.month, fields[3]) and
match_field(cron_weekday, fields[4])):
return True
curr += timedelta(minutes=1)
return False
return True
class AgentDaemon:
"""Background service worker for Windows: handles heartbeats, job polling and scheduled backups."""
@@ -172,8 +238,10 @@ class AgentDaemon:
lj.skip_system = sj.get("skip_system", lj.skip_system)
lj.skip_readonly = sj.get("skip_readonly", lj.skip_readonly)
lj.sync_deletions = sj.get("sync_deletions", lj.sync_deletions)
lj.local_destination_path = sj.get("local_destination_path", lj.local_destination_path)
lj.schedule_cron = sj.get("schedule_cron", lj.schedule_cron)
lj.min_stable_seconds = sj.get("min_stable_time_seconds", lj.min_stable_seconds)
lj.is_active = sj.get("is_active", lj.is_active)
# Also sync last_run_at from server if available and newer
if sj.get("last_run_at"):
@@ -186,7 +254,7 @@ class AgentDaemon:
config_changed = True
self.config.local_folders = local_jobs_to_keep
# Add server jobs that are missing locally
local_job_ids = {lj.job_id for lj in self.config.local_folders if lj.job_id is not None}
for sj in server_jobs:
@@ -201,6 +269,8 @@ class AgentDaemon:
skip_system=sj.get("skip_system", False),
skip_readonly=sj.get("skip_readonly", False),
sync_deletions=sj.get("sync_deletions", False),
local_destination_path=sj.get("local_destination_path"),
is_active=sj.get("is_active", True),
schedule_cron=sj["schedule_cron"],
min_stable_seconds=sj["min_stable_time_seconds"],
last_status=sj.get("status", "En espera")
@@ -224,6 +294,7 @@ class AgentDaemon:
"skip_system": lj.skip_system,
"skip_readonly": lj.skip_readonly,
"sync_deletions": lj.sync_deletions,
"local_destination_path": lj.local_destination_path,
"schedule_cron": lj.schedule_cron,
"min_stable_time_seconds": lj.min_stable_seconds
}
@@ -258,6 +329,23 @@ class AgentDaemon:
logger.warning(f"Source path {source_path} for job {job_name} does not exist. Skipping.")
continue
# Start run session on server
run_id = None
if job.job_id is not None:
try:
base_url = self.config.server_url.rstrip("/")
resp = httpx.post(
f"{base_url}/api/jobs/{job.job_id}/runs/start",
headers=self._get_headers(),
json={"client_id": 0},
timeout=10.0
)
if resp.status_code == 200:
run_id = resp.json().get("id")
logger.info(f"Started job run {run_id} on server.")
except Exception as e:
logger.warning(f"Could not start job run on server: {e}")
scanner = DirectoryScanner(
source_path=source_path,
file_patterns=file_patterns,
@@ -269,9 +357,20 @@ class AgentDaemon:
)
files = scanner.scan()
# Stats trackers
files_scanned = len(files)
files_copied = 0
files_skipped = 0
errors_count = 0
bytes_transferred = 0
error_log_messages = []
# Track files successfully backed up in this run
active_relative_paths = []
dest_dir = getattr(job, "local_destination_path", None)
for filepath in files:
rel_p = None
try:
rel_p = str(filepath.relative_to(Path(source_path).resolve())).replace("\\", "/")
active_relative_paths.append(rel_p)
@@ -280,42 +379,73 @@ class AgentDaemon:
if not is_file_stable(filepath, min_stable_seconds=min_stable):
logger.warning(f"File {filepath.name} is currently locked or growing. Skipping.")
files_skipped += 1
continue
current_sha = compute_file_sha256(filepath)
if state_db.is_file_already_backed_up(str(filepath), current_sha):
continue
file_size_bytes = filepath.stat().st_size
logger.info(f"Starting backup for file: {filepath.name} ({file_size_bytes / (1024*1024):.2f} MB)")
if self.on_started:
self.on_started(filepath.name, file_size_bytes)
cloud_uploaded = False
def on_chunk_progress(done, total, pct):
if self.on_progress:
self.on_progress(filepath.name, done, total, pct)
# 2A. Cloud Upload
if state_db.is_file_already_backed_up(str(filepath), current_sha):
files_skipped += 1
cloud_uploaded = True
else:
logger.info(f"Starting backup for file: {filepath.name} ({file_size_bytes / (1024*1024):.2f} MB)")
if self.on_started:
self.on_started(filepath.name, file_size_bytes)
try:
res = self.uploader.upload_file(
filepath,
job_id=job.job_id,
source_path=source_path,
progress_callback=on_chunk_progress
)
logger.info(f"Successfully backed up {filepath.name}!")
if self.on_completed:
self.on_completed(filepath.name, res.get("sha256", ""), file_size_bytes)
def on_chunk_progress(done, total, pct):
if self.on_progress:
self.on_progress(filepath.name, done, total, pct)
except Exception as ex:
logger.error(f"Failed to backup {filepath.name}: {str(ex)}")
job.last_status = "Error"
save_config(self.config)
if self.on_error:
self.on_error(filepath.name, str(ex))
try:
res = self.uploader.upload_file(
filepath,
job_id=job.job_id,
source_path=source_path,
progress_callback=on_chunk_progress
)
logger.info(f"Successfully backed up {filepath.name} to cloud!")
files_copied += 1
bytes_transferred += file_size_bytes
cloud_uploaded = True
if self.on_completed:
self.on_completed(filepath.name, res.get("sha256", ""), file_size_bytes)
except Exception as ex:
logger.error(f"Failed to backup {filepath.name} to cloud: {str(ex)}")
errors_count += 1
error_log_messages.append(f"Cloud upload error for {filepath.name}: {ex}")
if self.on_error:
self.on_error(filepath.name, str(ex))
# Enforce mirror replica mode (purge orphans on server)
# 2B. Duplicate Local Copy (Phase 5)
if dest_dir and rel_p:
try:
dest_path = Path(dest_dir) / rel_p
should_copy_local = True
if dest_path.exists():
try:
dest_stat = dest_path.stat()
src_stat = filepath.stat()
if dest_stat.st_size == src_stat.st_size and abs(dest_stat.st_mtime - src_stat.st_mtime) < 2:
should_copy_local = False
except Exception:
pass
if should_copy_local:
dest_path.parent.mkdir(parents=True, exist_ok=True)
shutil.copy2(filepath, dest_path)
logger.info(f"Successfully copied {filepath.name} to local dest {dest_path}")
if not cloud_uploaded:
files_copied += 1
bytes_transferred += file_size_bytes
except Exception as ex:
logger.error(f"Failed to copy {filepath.name} to local path {dest_dir}: {str(ex)}")
errors_count += 1
error_log_messages.append(f"Local copy error for {filepath.name}: {ex}")
# Enforce mirror replica mode on Cloud (purge orphans on server)
if job.sync_deletions and job.job_id is not None:
try:
base_url = self.config.server_url.rstrip("/")
@@ -336,7 +466,57 @@ class AgentDaemon:
except Exception as e:
logger.error(f"Error purging orphan files from server: {e}")
# Update job state in config after checking directory
# Enforce mirror replica mode on Local Destination (Phase 5)
if job.sync_deletions and dest_dir and Path(dest_dir).exists():
try:
dest_root = Path(dest_dir).resolve()
for root, dirs, files_in_dir in os.walk(dest_root):
for file_in_dir in files_in_dir:
full_path = Path(root) / file_in_dir
try:
rel_to_dest = str(full_path.relative_to(dest_root)).replace("\\", "/")
if rel_to_dest not in active_relative_paths:
logger.info(f"Removing local orphan file: {full_path}")
full_path.unlink()
except Exception as ex:
logger.error(f"Error removing local orphan file {full_path}: {ex}")
for root, dirs, files_in_dir in os.walk(dest_root, topdown=False):
for dir_name in dirs:
dir_path = Path(root) / dir_name
try:
if not os.listdir(dir_path):
logger.info(f"Removing empty local directory: {dir_path}")
dir_path.rmdir()
except Exception:
pass
except Exception as e:
logger.error(f"Error purging local orphans in {dest_dir}: {e}")
# Update job state in config
job.last_backup_at = datetime.now(timezone.utc).strftime("%Y-%m-%d %H:%M:%S")
job.last_status = "Backup Exitoso" if job.last_status != "Error" else "Error"
job_status = "SUCCESS" if errors_count == 0 else "FAILED"
job.last_status = "Backup Exitoso" if job_status == "SUCCESS" else "Error"
save_config(self.config)
# Complete run session on server
if run_id is not None:
try:
base_url = self.config.server_url.rstrip("/")
complete_payload = {
"status": job_status,
"files_scanned": files_scanned,
"files_copied": files_copied,
"files_skipped": files_skipped,
"errors_count": errors_count,
"bytes_transferred": bytes_transferred,
"error_summary": "\n".join(error_log_messages) if error_log_messages else None
}
httpx.post(
f"{base_url}/api/jobs/{job.job_id}/runs/{run_id}/complete",
headers=self._get_headers(),
json=complete_payload,
timeout=10.0
)
logger.info(f"Completed job run {run_id} on server with status {job_status}.")
except Exception as e:
logger.warning(f"Could not complete job run on server: {e}")