Files
OnEverDrive/windows-agent/agent/service.py
T

605 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
)
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}")