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 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 ) 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) 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) 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("/") 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"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_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}")