606 lines
27 KiB
Python
606 lines
27 KiB
Python
import time
|
|
import socket
|
|
import platform
|
|
import threading
|
|
import logging
|
|
import shutil
|
|
import re
|
|
import os
|
|
from datetime import datetime, timezone, timedelta
|
|
from pathlib import Path
|
|
from typing import Optional, Callable, Dict, Any, List
|
|
import httpx
|
|
|
|
from agent.config import AgentConfig, load_config, save_config, LocalFolderJob
|
|
from agent.scanner import DirectoryScanner, is_file_stable
|
|
from agent.chunker import compute_file_sha256
|
|
from agent.uploader import ChunkUploader
|
|
from agent.state_db import state_db
|
|
|
|
logging.basicConfig(
|
|
level=logging.INFO,
|
|
format="%(asctime)s [%(levelname)s] [OnEver Agent] %(message)s"
|
|
)
|
|
logger = logging.getLogger("OnEverAgent")
|
|
|
|
def send_to_recycle_bin(path_str: str) -> bool:
|
|
"""
|
|
Sends a file or folder to the Windows Recycle Bin using native shell32 API.
|
|
"""
|
|
try:
|
|
import ctypes
|
|
from ctypes import wintypes
|
|
|
|
class SHFILEOPSTRUCTW(ctypes.Structure):
|
|
_fields_ = [
|
|
("hwnd", wintypes.HWND),
|
|
("wFunc", wintypes.UINT),
|
|
("pFrom", ctypes.c_wchar_p),
|
|
("pTo", ctypes.c_wchar_p),
|
|
("fFlags", ctypes.c_ushort),
|
|
("fAnyOperationsAborted", wintypes.BOOL),
|
|
("hNameMappings", wintypes.LPVOID),
|
|
("lpszProgressTitle", ctypes.c_wchar_p),
|
|
]
|
|
|
|
path_null = os.path.abspath(path_str) + "\0\0"
|
|
fileop = SHFILEOPSTRUCTW()
|
|
fileop.hwnd = None
|
|
fileop.wFunc = 3 # FO_DELETE
|
|
fileop.pFrom = path_null
|
|
fileop.pTo = None
|
|
# FOF_ALLOWUNDO (0x0040) sends to Recycle Bin. FOF_NOCONFIRMATION (0x0010) + FOF_NOERRORUI (0x0400) + FOF_SILENT (0x0004)
|
|
fileop.fFlags = 0x0040 | 0x0010 | 0x0400 | 0x0004
|
|
|
|
res = ctypes.windll.shell32.SHFileOperationW(ctypes.byref(fileop))
|
|
return res == 0
|
|
except Exception as e:
|
|
logger.error(f"Failed to recycle bin file {path_str} via ctypes: {e}")
|
|
return False
|
|
|
|
def resolve_path_tags(path: str) -> str:
|
|
"""
|
|
Resolves Karen-style datetime tags inside file/folder paths.
|
|
"""
|
|
now = datetime.now()
|
|
replacements = {
|
|
"<yyyy>": now.strftime("%Y"),
|
|
"<yy>": now.strftime("%y"),
|
|
"<year>": now.strftime("%Y"),
|
|
"<y>": now.strftime("%Y"),
|
|
"<month>": now.strftime("%b"),
|
|
"<mm>": now.strftime("%m"),
|
|
"<m>": str(now.month),
|
|
"<dd>": now.strftime("%d"),
|
|
"<d>": str(now.day),
|
|
"<dow>": str(now.weekday() + 1),
|
|
"<w>": str(now.weekday() + 1),
|
|
"<hour>": now.strftime("%H"),
|
|
"<hh>": now.strftime("%H"),
|
|
"<minute>": now.strftime("%M"),
|
|
"<min>": now.strftime("%M"),
|
|
}
|
|
resolved = path
|
|
for tag, val in replacements.items():
|
|
resolved = re.sub(re.escape(tag), val, resolved, flags=re.IGNORECASE)
|
|
return resolved
|
|
|
|
def is_job_due(schedule_str: str, last_run_str: Optional[str]) -> bool:
|
|
if not schedule_str:
|
|
return True
|
|
if not last_run_str:
|
|
return True
|
|
|
|
try:
|
|
if "T" in last_run_str:
|
|
last_run = datetime.fromisoformat(last_run_str.replace("Z", "+00:00"))
|
|
else:
|
|
last_run = datetime.strptime(last_run_str, "%Y-%m-%d %H:%M:%S").replace(tzinfo=timezone.utc)
|
|
except Exception:
|
|
return True
|
|
|
|
now = datetime.now(timezone.utc)
|
|
if last_run >= now:
|
|
return False
|
|
|
|
sched = schedule_str.lower().strip()
|
|
|
|
# 1. Traditional intervals
|
|
if sched == "hourly":
|
|
return (now - last_run).total_seconds() >= 3600
|
|
elif sched == "daily":
|
|
return (now - last_run).total_seconds() >= 86400
|
|
elif sched == "weekly":
|
|
return (now - last_run).total_seconds() >= 86400 * 7
|
|
elif sched == "monthly":
|
|
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."""
|
|
|
|
def __init__(
|
|
self,
|
|
config: Optional[AgentConfig] = None,
|
|
on_started: Optional[Callable[[str, int], None]] = None,
|
|
on_progress: Optional[Callable[[str, int, int, float], None]] = None,
|
|
on_completed: Optional[Callable[[str, str, int], None]] = None,
|
|
on_error: Optional[Callable[[str, str], None]] = None,
|
|
on_status: Optional[Callable[[str, str], None]] = None
|
|
):
|
|
self.config = config or load_config()
|
|
self.running = False
|
|
self.uploader = ChunkUploader(self.config)
|
|
self.last_server_jobs = []
|
|
self._heartbeat_thread: Optional[threading.Thread] = None
|
|
self._worker_thread: Optional[threading.Thread] = None
|
|
|
|
# Event callbacks
|
|
self.on_started = on_started
|
|
self.on_progress = on_progress
|
|
self.on_completed = on_completed
|
|
self.on_error = on_error
|
|
self.on_status = on_status
|
|
|
|
def start(self):
|
|
self.config = load_config()
|
|
if not self.config.device_id or not self.config.device_token:
|
|
logger.warning("Agent is not registered yet. Waiting for registration.")
|
|
if self.on_status:
|
|
self.on_status("UNREGISTERED", "El agente no está registrado en el servidor.")
|
|
return
|
|
|
|
self.running = True
|
|
logger.info(f"Starting OnEver Drive Windows Agent ({self.config.client_code} - {self.config.client_name})")
|
|
logger.info(f"Target Server: {self.config.server_url}")
|
|
|
|
if self.on_status:
|
|
self.on_status("ONLINE", f"Conectado a {self.config.server_url} ({self.config.client_code})")
|
|
|
|
self._heartbeat_thread = threading.Thread(target=self._heartbeat_loop, daemon=True)
|
|
self._worker_thread = threading.Thread(target=self._backup_worker_loop, daemon=True)
|
|
|
|
self._heartbeat_thread.start()
|
|
self._worker_thread.start()
|
|
|
|
def stop(self):
|
|
logger.info("Stopping agent daemon...")
|
|
self.running = False
|
|
if self.on_status:
|
|
self.on_status("PAUSED", "Servicio en pausa.")
|
|
|
|
def _get_headers(self):
|
|
return {
|
|
"X-Device-Id": self.config.device_id,
|
|
"X-Device-Token": self.config.device_token
|
|
}
|
|
|
|
def _heartbeat_loop(self):
|
|
while self.running:
|
|
try:
|
|
base_url = self.config.server_url.rstrip("/")
|
|
with httpx.Client(base_url=base_url, headers=self._get_headers(), timeout=10.0) as client:
|
|
resp = client.get("/api/jobs/agent/assigned")
|
|
if resp.status_code == 200:
|
|
logger.debug("Heartbeat acknowledged by server.")
|
|
except Exception as ex:
|
|
logger.warning(f"Heartbeat failed: {str(ex)}")
|
|
|
|
time.sleep(self.config.heartbeat_interval_seconds)
|
|
|
|
def _backup_worker_loop(self):
|
|
while self.running:
|
|
try:
|
|
self._run_backup_cycle()
|
|
except Exception as ex:
|
|
logger.error(f"Error during backup cycle: {str(ex)}")
|
|
|
|
time.sleep(30)
|
|
|
|
def _run_backup_cycle(self, force: bool = False):
|
|
self.config = load_config()
|
|
self.uploader.config = self.config
|
|
|
|
# 1. Fetch server-assigned jobs and run bidirectional sync
|
|
server_jobs = []
|
|
sync_success = False
|
|
try:
|
|
base_url = self.config.server_url.rstrip("/")
|
|
with httpx.Client(base_url=base_url, headers=self._get_headers(), timeout=10.0) as client:
|
|
resp = client.get("/api/jobs/agent/assigned")
|
|
if resp.status_code == 200:
|
|
server_jobs = resp.json()
|
|
self.last_server_jobs = server_jobs
|
|
sync_success = True
|
|
except Exception as e:
|
|
logger.warning(f"Could not fetch server jobs: {e}")
|
|
|
|
config_changed = False
|
|
|
|
if sync_success:
|
|
# A. Sync Server -> Local
|
|
server_job_ids = {sj["id"] for sj in server_jobs}
|
|
|
|
# Remove local jobs that have a job_id but are not on the server anymore (deleted on server)
|
|
local_jobs_to_keep = []
|
|
for lj in self.config.local_folders:
|
|
if lj.job_id is None:
|
|
# New local job, keep it so we register it next
|
|
local_jobs_to_keep.append(lj)
|
|
elif lj.job_id in server_job_ids:
|
|
# Keep it and update local properties from server
|
|
sj = next(x for x in server_jobs if x["id"] == lj.job_id)
|
|
lj.name = sj.get("name", lj.name)
|
|
lj.source_path = sj.get("source_path", lj.source_path)
|
|
lj.file_patterns = sj.get("file_patterns", lj.file_patterns)
|
|
lj.exclude_patterns = sj.get("exclude_patterns", lj.exclude_patterns)
|
|
lj.skip_hidden = sj.get("skip_hidden", lj.skip_hidden)
|
|
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"):
|
|
lj.last_backup_at = sj["last_run_at"].replace("T", " ")[:19]
|
|
lj.last_status = sj.get("status", lj.last_status)
|
|
|
|
local_jobs_to_keep.append(lj)
|
|
else:
|
|
# Deleted on server, don't keep it
|
|
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:
|
|
if sj["id"] not in local_job_ids:
|
|
new_job = LocalFolderJob(
|
|
job_id=sj["id"],
|
|
name=sj["name"],
|
|
source_path=sj["source_path"],
|
|
file_patterns=sj["file_patterns"],
|
|
exclude_patterns=sj.get("exclude_patterns", ""),
|
|
skip_hidden=sj.get("skip_hidden", False),
|
|
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")
|
|
)
|
|
if sj.get("last_run_at"):
|
|
new_job.last_backup_at = sj["last_run_at"].replace("T", " ")[:19]
|
|
self.config.local_folders.append(new_job)
|
|
config_changed = True
|
|
|
|
# B. Sync Local -> Server (Register new local folders on the server)
|
|
for lj in self.config.local_folders:
|
|
if lj.job_id is None:
|
|
try:
|
|
base_url = self.config.server_url.rstrip("/")
|
|
payload = {
|
|
"name": lj.name,
|
|
"source_path": lj.source_path,
|
|
"file_patterns": lj.file_patterns,
|
|
"exclude_patterns": lj.exclude_patterns,
|
|
"skip_hidden": lj.skip_hidden,
|
|
"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
|
|
}
|
|
resp = httpx.post(f"{base_url}/api/jobs/agent/register", headers=self._get_headers(), json=payload, timeout=10.0)
|
|
if resp.status_code == 200:
|
|
data = resp.json()
|
|
lj.job_id = data["id"]
|
|
config_changed = True
|
|
logger.info(f"Registered local job '{lj.name}' on server with ID {lj.job_id}")
|
|
except Exception as e:
|
|
logger.warning(f"Could not register local job '{lj.name}' on server: {e}")
|
|
|
|
if config_changed:
|
|
save_config(self.config)
|
|
|
|
# 2. Process active jobs
|
|
for job in self.config.local_folders:
|
|
if not job.is_active:
|
|
continue
|
|
|
|
# Check if job is due or forced
|
|
if not force and not is_job_due(job.schedule_cron, job.last_backup_at):
|
|
continue
|
|
|
|
source_path = job.source_path
|
|
file_patterns = job.file_patterns
|
|
min_stable = job.min_stable_seconds
|
|
job_name = job.name
|
|
|
|
# Skip if path does not exist
|
|
if not Path(source_path).exists():
|
|
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,
|
|
exclude_patterns=getattr(job, "exclude_patterns", ""),
|
|
skip_hidden=getattr(job, "skip_hidden", False),
|
|
skip_system=getattr(job, "skip_system", False),
|
|
skip_readonly=getattr(job, "skip_readonly", False),
|
|
min_stable_seconds=min_stable,
|
|
enable_global_exclusions=getattr(self.config, "enable_global_exclusions", True)
|
|
)
|
|
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 = resolve_path_tags(getattr(job, "local_destination_path", None)) if getattr(job, "local_destination_path", None) else 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)
|
|
except Exception:
|
|
pass
|
|
|
|
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)
|
|
file_size_bytes = filepath.stat().st_size
|
|
cloud_uploaded = False
|
|
|
|
# 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)
|
|
|
|
def on_chunk_progress(done, total, pct):
|
|
if self.on_progress:
|
|
self.on_progress(filepath.name, done, total, pct)
|
|
|
|
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))
|
|
|
|
# 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)
|
|
# Robust write pattern inspired by Karen's Replicator 3.5.0
|
|
temp_dest = dest_path.with_suffix(dest_path.suffix + ".tmp")
|
|
try:
|
|
shutil.copy2(filepath, temp_dest)
|
|
if dest_path.exists():
|
|
dest_path.unlink()
|
|
temp_dest.rename(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 copy_err:
|
|
if temp_dest.exists():
|
|
try:
|
|
temp_dest.unlink()
|
|
except Exception:
|
|
pass
|
|
raise copy_err
|
|
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("/")
|
|
purge_payload = {
|
|
"active_relative_paths": active_relative_paths
|
|
}
|
|
resp = httpx.post(
|
|
f"{base_url}/api/jobs/{job.job_id}/purge-orphans",
|
|
headers=self._get_headers(),
|
|
json=purge_payload,
|
|
timeout=30.0
|
|
)
|
|
if resp.status_code == 200:
|
|
purge_res = resp.json()
|
|
purged_count = purge_res.get("purged_count", 0)
|
|
if purged_count > 0:
|
|
logger.info(f"Purged {purged_count} orphan files from server for job '{job_name}'")
|
|
except Exception as e:
|
|
logger.error(f"Error purging orphan files from server: {e}")
|
|
|
|
# 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"Recycling local orphan file: {full_path}")
|
|
recycled = send_to_recycle_bin(str(full_path))
|
|
if not recycled:
|
|
logger.info(f"Recycle bin failed. Permanently deleting 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"Recycling empty local directory: {dir_path}")
|
|
recycled = send_to_recycle_bin(str(dir_path))
|
|
if not recycled:
|
|
logger.info(f"Recycle bin failed. Permanently 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_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}")
|