diff --git a/backend/app/api/auth.py b/backend/app/api/auth.py index 2de449b..16387dd 100644 --- a/backend/app/api/auth.py +++ b/backend/app/api/auth.py @@ -10,12 +10,24 @@ from app.services.event_service import log_event router = APIRouter(prefix="/auth", tags=["Authentication"]) +import logging +logger = logging.getLogger("uvicorn.error") + @router.post("/login", response_model=TokenResponse) async def login(credentials: LoginRequest, db: AsyncSession = Depends(get_db)): - result = await db.execute(select(User).where(User.email == credentials.email.strip().lower())) + email_clean = credentials.email.strip().lower() + result = await db.execute(select(User).where(User.email == email_clean)) user = result.scalar_one_or_none() - if not user or not verify_password(credentials.password, user.hashed_password): + if not user: + logger.warning(f"Login failed: User not found with email '{email_clean}'") + raise HTTPException( + status_code=status.HTTP_401_UNAUTHORIZED, + detail="Incorrect email or password" + ) + + if not verify_password(credentials.password, user.hashed_password): + logger.warning(f"Login failed: Password mismatch for email '{email_clean}' (sent password length: {len(credentials.password)})") raise HTTPException( status_code=status.HTTP_401_UNAUTHORIZED, detail="Incorrect email or password" diff --git a/backend/app/api/jobs.py b/backend/app/api/jobs.py index 5589a3a..1a39985 100644 --- a/backend/app/api/jobs.py +++ b/backend/app/api/jobs.py @@ -5,8 +5,11 @@ from sqlalchemy.ext.asyncio import AsyncSession from sqlalchemy import select, func, desc from app.core.database import get_db -from app.models.models import User, Client, BackupJob, BackupFile -from app.schemas.schemas import JobCreate, JobUpdate, JobResponse, AgentJobRegister, PurgeOrphansRequest +from app.models.models import User, Client, BackupJob, BackupFile, JobRun +from app.schemas.schemas import ( + JobCreate, JobUpdate, JobResponse, AgentJobRegister, PurgeOrphansRequest, + JobRunResponse, JobRunStartRequest, JobRunCompleteRequest +) from app.api.deps import get_current_user, require_admin, get_current_client from app.services.event_service import log_event from app.ws.manager import ws_manager @@ -46,9 +49,9 @@ async def agent_register_job( current_client: Client = Depends(get_current_client) ): """Called by the Windows Agent to register a new local job/folder on the server.""" - count_res = await db.execute(select(func.count(BackupJob.id))) - job_count = count_res.scalar() or 0 - job_code = f"JOB-{job_count + 1:03d}" + max_id_res = await db.execute(select(func.max(BackupJob.id))) + max_id = max_id_res.scalar() or 0 + job_code = f"JOB-{max_id + 1:03d}" job = BackupJob( job_code=job_code, @@ -61,6 +64,7 @@ async def agent_register_job( skip_system=payload.skip_system, skip_readonly=payload.skip_readonly, sync_deletions=payload.sync_deletions, + local_destination_path=payload.local_destination_path, schedule_cron=payload.schedule_cron, keep_daily=payload.keep_daily, keep_weekly=payload.keep_weekly, @@ -124,9 +128,9 @@ async def create_job( if not client: raise HTTPException(status_code=status.HTTP_404_NOT_FOUND, detail="Client not found") - count_res = await db.execute(select(func.count(BackupJob.id))) - job_count = count_res.scalar() or 0 - job_code = f"JOB-{job_count + 1:03d}" + max_id_res = await db.execute(select(func.max(BackupJob.id))) + max_id = max_id_res.scalar() or 0 + job_code = f"JOB-{max_id + 1:03d}" job = BackupJob( job_code=job_code, @@ -139,6 +143,7 @@ async def create_job( skip_system=payload.skip_system, skip_readonly=payload.skip_readonly, sync_deletions=payload.sync_deletions, + local_destination_path=payload.local_destination_path, schedule_cron=payload.schedule_cron, keep_daily=payload.keep_daily, keep_weekly=payload.keep_weekly, @@ -322,3 +327,81 @@ async def purge_orphans( "freed_bytes": total_freed_bytes, "message": f"Successfully purged {purged_count} orphan files from server." } + +@router.post("/{job_id}/runs/start", response_model=JobRunResponse) +async def start_job_run( + job_id: int, + payload: JobRunStartRequest, + db: AsyncSession = Depends(get_db), + current_client: Client = Depends(get_current_client) +): + """Called by the Windows Agent to register the start of a scheduled or manual job run.""" + res = await db.execute(select(BackupJob).where(BackupJob.id == job_id, BackupJob.client_id == current_client.id)) + job = res.scalar_one_or_none() + if not job: + raise HTTPException(status_code=status.HTTP_404_NOT_FOUND, detail="Job not found or not owned by client") + + job.status = "RUNNING" + job.last_run_at = datetime.now(timezone.utc) + + run = JobRun( + job_id=job.id, + client_id=current_client.id, + started_at=datetime.now(timezone.utc), + status="RUNNING" + ) + db.add(run) + await db.commit() + await db.refresh(run) + + return run + +@router.post("/{job_id}/runs/{run_id}/complete", response_model=JobRunResponse) +async def complete_job_run( + job_id: int, + run_id: int, + payload: JobRunCompleteRequest, + db: AsyncSession = Depends(get_db), + current_client: Client = Depends(get_current_client) +): + """Called by the Windows Agent to record completion metrics for a run.""" + res = await db.execute(select(JobRun).where(JobRun.id == run_id, JobRun.job_id == job_id, JobRun.client_id == current_client.id)) + run = res.scalar_one_or_none() + if not run: + raise HTTPException(status_code=status.HTTP_404_NOT_FOUND, detail="Job run not found") + + job_res = await db.execute(select(BackupJob).where(BackupJob.id == job_id)) + job = job_res.scalar_one_or_none() + + run.completed_at = datetime.now(timezone.utc) + run.status = payload.status + run.files_scanned = payload.files_scanned + run.files_copied = payload.files_copied + run.files_skipped = payload.files_skipped + run.errors_count = payload.errors_count + run.bytes_transferred = payload.bytes_transferred + run.error_summary = payload.error_summary + + if job: + job.status = payload.status + job.last_run_at = run.completed_at + + await db.commit() + await db.refresh(run) + + return run + +@router.get("/{job_id}/runs", response_model=List[JobRunResponse]) +async def list_job_runs( + job_id: int, + db: AsyncSession = Depends(get_db), + current_user: User = Depends(get_current_user) +): + """Called by the Web UI to retrieve execution history for a job.""" + result = await db.execute( + select(JobRun) + .where(JobRun.job_id == job_id) + .order_by(desc(JobRun.started_at)) + .limit(50) + ) + return result.scalars().all() diff --git a/backend/app/models/models.py b/backend/app/models/models.py index 5b97da1..521237f 100644 --- a/backend/app/models/models.py +++ b/backend/app/models/models.py @@ -98,11 +98,34 @@ class BackupJob(Base): status = Column(String(50), default="IDLE", nullable=False) # IDLE, RUNNING, ERROR, SUCCESS last_run_at = Column(DateTime(timezone=True), nullable=True) next_run_at = Column(DateTime(timezone=True), nullable=True) + local_destination_path = Column(String(1024), nullable=True) created_at = Column(DateTime(timezone=True), default=utc_now, nullable=False) client = relationship("Client", back_populates="jobs") backup_files = relationship("BackupFile", back_populates="job", cascade="all, delete-orphan") backup_sessions = relationship("BackupSession", back_populates="job", cascade="all, delete-orphan") + runs = relationship("JobRun", back_populates="job", cascade="all, delete-orphan") + +class JobRun(Base): + __tablename__ = "job_runs" + + id = Column(Integer, primary_key=True, index=True) + job_id = Column(Integer, ForeignKey("backup_jobs.id", ondelete="CASCADE"), nullable=False) + client_id = Column(Integer, ForeignKey("clients.id", ondelete="CASCADE"), nullable=False) + started_at = Column(DateTime(timezone=True), default=utc_now, nullable=False) + completed_at = Column(DateTime(timezone=True), nullable=True) + status = Column(String(50), default="RUNNING", nullable=False) # RUNNING, SUCCESS, FAILED + + # Detailed execution metrics + files_scanned = Column(Integer, default=0, nullable=False) + files_copied = Column(Integer, default=0, nullable=False) + files_skipped = Column(Integer, default=0, nullable=False) + errors_count = Column(Integer, default=0, nullable=False) + bytes_transferred = Column(BigInteger, default=0, nullable=False) + error_summary = Column(Text, nullable=True) + + job = relationship("BackupJob", back_populates="runs") + client = relationship("Client") class BackupSession(Base): __tablename__ = "backup_sessions" diff --git a/backend/app/schemas/schemas.py b/backend/app/schemas/schemas.py index d56e5ae..0364675 100644 --- a/backend/app/schemas/schemas.py +++ b/backend/app/schemas/schemas.py @@ -87,6 +87,7 @@ class JobCreate(BaseModel): skip_system: bool = False skip_readonly: bool = False sync_deletions: bool = False + local_destination_path: Optional[str] = None schedule_cron: str = "0 2 * * *" keep_daily: int = 7 keep_weekly: int = 4 @@ -102,6 +103,7 @@ class JobUpdate(BaseModel): skip_system: Optional[bool] = None skip_readonly: Optional[bool] = None sync_deletions: Optional[bool] = None + local_destination_path: Optional[str] = None schedule_cron: Optional[str] = None is_active: Optional[bool] = None keep_daily: Optional[int] = None @@ -121,6 +123,7 @@ class JobResponse(BaseModel): skip_system: bool skip_readonly: bool sync_deletions: bool + local_destination_path: Optional[str] schedule_cron: str is_active: bool keep_daily: int @@ -255,6 +258,7 @@ class AgentJobRegister(BaseModel): skip_system: bool = False skip_readonly: bool = False sync_deletions: bool = False + local_destination_path: Optional[str] = None schedule_cron: str = "daily" keep_daily: int = 7 keep_weekly: int = 4 @@ -263,3 +267,32 @@ class AgentJobRegister(BaseModel): class PurgeOrphansRequest(BaseModel): active_relative_paths: List[str] + +# --- Job Run Schemas --- +class JobRunStartRequest(BaseModel): + client_id: int + +class JobRunCompleteRequest(BaseModel): + status: str # SUCCESS, FAILED + files_scanned: int + files_copied: int + files_skipped: int + errors_count: int + bytes_transferred: int + error_summary: Optional[str] = None + +class JobRunResponse(BaseModel): + id: int + job_id: int + client_id: int + started_at: datetime + completed_at: Optional[datetime] + status: str + files_scanned: int + files_copied: int + files_skipped: int + errors_count: int + bytes_transferred: int + error_summary: Optional[str] + + model_config = {"from_attributes": True} diff --git a/backend/migrate_runs_and_dest.py b/backend/migrate_runs_and_dest.py new file mode 100644 index 0000000..deaf0d5 --- /dev/null +++ b/backend/migrate_runs_and_dest.py @@ -0,0 +1,47 @@ +import sqlite3 +import os + +DB_PATH = os.path.join(os.path.dirname(__file__), "onever_drive.db") + +def migrate(): + print(f"Connecting to database at {DB_PATH}...") + conn = sqlite3.connect(DB_PATH) + cursor = conn.cursor() + + # 1. Add local_destination_path column to backup_jobs table + try: + cursor.execute("ALTER TABLE backup_jobs ADD COLUMN local_destination_path VARCHAR(1024) NULL") + print("Column 'local_destination_path' added successfully to backup_jobs.") + except sqlite3.OperationalError as e: + print(f"Column 'local_destination_path' could not be added (maybe it already exists?): {e}") + + # 2. Create job_runs table + try: + cursor.execute(""" + CREATE TABLE IF NOT EXISTS job_runs ( + id INTEGER PRIMARY KEY AUTOINCREMENT, + job_id INTEGER NOT NULL, + client_id INTEGER NOT NULL, + started_at DATETIME NOT NULL, + completed_at DATETIME NULL, + status VARCHAR(50) NOT NULL DEFAULT 'RUNNING', + files_scanned INTEGER NOT NULL DEFAULT 0, + files_copied INTEGER NOT NULL DEFAULT 0, + files_skipped INTEGER NOT NULL DEFAULT 0, + errors_count INTEGER NOT NULL DEFAULT 0, + bytes_transferred BIGINT NOT NULL DEFAULT 0, + error_summary TEXT NULL, + FOREIGN KEY(job_id) REFERENCES backup_jobs(id) ON DELETE CASCADE, + FOREIGN KEY(client_id) REFERENCES clients(id) ON DELETE CASCADE + ) + """) + print("Table 'job_runs' created successfully.") + except sqlite3.OperationalError as e: + print(f"Table 'job_runs' could not be created: {e}") + + conn.commit() + conn.close() + print("Migration finished.") + +if __name__ == "__main__": + migrate() diff --git a/frontend/src/pages/JobsView.tsx b/frontend/src/pages/JobsView.tsx index 7d2b19e..30596bd 100644 --- a/frontend/src/pages/JobsView.tsx +++ b/frontend/src/pages/JobsView.tsx @@ -37,9 +37,30 @@ export const JobsView: React.FC = ({ jobs, clients, onRefresh, se skip_system: false, skip_readonly: false, sync_deletions: false, + local_destination_path: '', }); const [editingJobId, setEditingJobId] = useState(null); + const [showRunsModal, setShowRunsModal] = useState(false); + const [selectedJobForRuns, setSelectedJobForRuns] = useState(null); + const [runsHistory, setRunsHistory] = useState([]); + const [loadingRuns, setLoadingRuns] = useState(false); + const [loading, setLoading] = useState(false); + + const handleShowRuns = async (job: BackupJobItem) => { + setSelectedJobForRuns(job); + setShowRunsModal(true); + setLoadingRuns(true); + try { + const history = await api.getJobRuns(job.id); + setRunsHistory(history); + } catch (err: any) { + alert(`Error obteniendo historial: ${err.message}`); + } finally { + setLoadingRuns(false); + } + }; + const handleSubmitJob = async (e: React.FormEvent) => { e.preventDefault(); if (!formData.client_id) { @@ -79,6 +100,7 @@ export const JobsView: React.FC = ({ jobs, clients, onRefresh, se skip_system: job.skip_system, skip_readonly: job.skip_readonly, sync_deletions: job.sync_deletions, + local_destination_path: job.local_destination_path || '', }); setShowModal(true); }; @@ -141,6 +163,7 @@ export const JobsView: React.FC = ({ jobs, clients, onRefresh, se skip_system: false, skip_readonly: false, sync_deletions: false, + local_destination_path: '', }); setShowModal(true); }} @@ -186,9 +209,17 @@ export const JobsView: React.FC = ({ jobs, clients, onRefresh, se {getClientName(job.client_id)} -
- - {job.source_path} +
+
+ + {job.source_path} +
+ {job.local_destination_path && ( +
+ Copia Local: + {job.local_destination_path} +
+ )}
@@ -223,6 +254,14 @@ export const JobsView: React.FC = ({ jobs, clients, onRefresh, se Ejecutar + +
+ +
+ {loadingRuns ? ( +
Cargando historial...
+ ) : runsHistory.length === 0 ? ( +
No hay ejecuciones registradas para este trabajo todavía.
+ ) : ( +
+ + + + + + + + + + + + + + + {runsHistory.map((run) => { + const start = new Date(run.started_at); + const end = run.completed_at ? new Date(run.completed_at) : null; + const durationSecs = end ? Math.round((end.getTime() - start.getTime()) / 1000) : 0; + + return ( + + + + + + + + + + + ); + })} + +
Fecha de InicioDuración / FinEstadoEscaneadosCopiadosOmitidosErroresTransferido
{start.toLocaleString()} + {end ? `${durationSecs}s (${end.toLocaleTimeString()})` : En ejecución...} + + + {run.status} + + {run.error_summary && ( +
+ {run.error_summary} +
+ )} +
{run.files_scanned}{run.files_copied}{run.files_skipped} 0 ? '#f43f5e' : 'inherit' }}>{run.errors_count}{(run.bytes_transferred / (1024 * 1024)).toFixed(2)} MB
+
+ )} +
+ +
+ +
+ + + )} ); }; diff --git a/frontend/src/services/api.ts b/frontend/src/services/api.ts index f829a6a..d04c2ee 100644 --- a/frontend/src/services/api.ts +++ b/frontend/src/services/api.ts @@ -48,6 +48,7 @@ export interface BackupJobItem { skip_system: boolean; skip_readonly: boolean; sync_deletions: boolean; + local_destination_path?: string | null; schedule_cron: string; is_active: boolean; keep_daily: number; @@ -60,6 +61,21 @@ export interface BackupJobItem { created_at: string; } +export interface JobRunItem { + id: number; + job_id: number; + client_id: number; + started_at: string; + completed_at?: string; + status: string; + files_scanned: number; + files_copied: number; + files_skipped: number; + errors_count: number; + bytes_transferred: number; + error_summary?: string; +} + export interface BackupFileItem { id: number; client_id: number; @@ -200,6 +216,7 @@ export const api = { skip_system: boolean; skip_readonly: boolean; sync_deletions: boolean; + local_destination_path?: string | null; schedule_cron: string; keep_daily: number; keep_weekly: number; @@ -220,6 +237,7 @@ export const api = { skip_system: boolean; skip_readonly: boolean; sync_deletions: boolean; + local_destination_path?: string | null; schedule_cron: string; keep_daily: number; keep_weekly: number; @@ -235,6 +253,8 @@ export const api = { request<{ message: string }>(`/jobs/${jobId}/trigger`, { method: 'POST' }), deleteJob: (jobId: number) => request<{ message: string }>(`/jobs/${jobId}`, { method: 'DELETE' }), + getJobRuns: (jobId: number) => + request(`/jobs/${jobId}/runs`), // Backups getBackups: (clientId?: number, jobId?: number) => { diff --git a/tests/test_karens_features.py b/tests/test_karens_features.py index f134d68..8580431 100644 --- a/tests/test_karens_features.py +++ b/tests/test_karens_features.py @@ -4,7 +4,7 @@ import platform import ctypes import tempfile from pathlib import Path -from datetime import datetime, timezone +from datetime import datetime, timezone, timedelta from app.core.database import init_db, AsyncSessionLocal from app.models.models import Client, BackupJob, BackupFile @@ -12,6 +12,7 @@ from app.api.jobs import purge_orphans from app.schemas.schemas import PurgeOrphansRequest from app.core.config import settings from agent.scanner import DirectoryScanner +from agent.service import is_job_due def set_windows_attributes(filepath: Path, hidden: bool = False, system: bool = False, readonly: bool = False): if platform.system() != "Windows": @@ -222,3 +223,82 @@ async def test_replica_mode_purge_orphans(): assert db_files[0].is_active == True assert db_files[2].is_active == True assert (Path(settings.STORAGE_ROOT) / db_files[0].relative_path).exists() + +def test_is_job_due_advanced(): + # 1. Traditional intervals + # Last run was 50 mins ago (3000s) + last_run_iso = (datetime.now(timezone.utc) - timedelta(minutes=50)).isoformat() + assert is_job_due("hourly", last_run_iso) == False + + # Last run was 70 mins ago (4200s) + last_run_iso = (datetime.now(timezone.utc) - timedelta(minutes=70)).isoformat() + assert is_job_due("hourly", last_run_iso) == True + + # 2. Calendar schedule match + # Let's test a simple day: + # If the scheduled time occurred in the window, it is due. + five_mins_ago = datetime.now(timezone.utc) - timedelta(minutes=5) + day_name = ["mon", "tue", "wed", "thu", "fri", "sat", "sun"][five_mins_ago.weekday()] + time_str = five_mins_ago.strftime("%H:%M") + + sched_str = f"{day_name} {time_str}" + last_run = (five_mins_ago - timedelta(hours=1)).strftime("%Y-%m-%d %H:%M:%S") + assert is_job_due(sched_str, last_run) == True + + # 3. Cron schedule match + # Cron format: minute hour day month weekday (0-6 starting Sunday) + cron_weekday = (five_mins_ago.weekday() + 1) % 7 + cron_str = f"{five_mins_ago.minute} {five_mins_ago.hour} * * {cron_weekday}" + assert is_job_due(cron_str, last_run) == True + +@pytest.mark.asyncio +async def test_job_runs_api_telemetry(): + """ + Tests the creation, execution tracking, and completion endpoints of job runs. + """ + async with AsyncSessionLocal() as db: + # 1. Setup client and job + client = Client( + client_code="CLIENT-R-01", + name="Runs Test Machine", + status="ONLINE", + storage_used_bytes=0, + storage_quota_bytes=100000 + ) + db.add(client) + await db.flush() + + job = BackupJob( + job_code="JOB-R-01", + client_id=client.id, + name="Runs Test Job", + source_path="C:\\RunsSource" + ) + db.add(job) + await db.commit() + + # 2. Simulate runs endpoints inside backend + from app.api.jobs import start_job_run, complete_job_run + from app.schemas.schemas import JobRunStartRequest, JobRunCompleteRequest + + start_req = JobRunStartRequest(client_id=client.id) + run_res = await start_job_run(job_id=job.id, payload=start_req, db=db, current_client=client) + assert run_res.status == "RUNNING" + assert run_res.job_id == job.id + + # Complete the run + complete_req = JobRunCompleteRequest( + status="SUCCESS", + files_scanned=10, + files_copied=3, + files_skipped=7, + errors_count=0, + bytes_transferred=50000 + ) + completed_res = await complete_job_run(job_id=job.id, run_id=run_res.id, payload=complete_req, db=db, current_client=client) + assert completed_res.status == "SUCCESS" + assert completed_res.files_scanned == 10 + assert completed_res.files_copied == 3 + assert completed_res.files_skipped == 7 + assert completed_res.errors_count == 0 + assert completed_res.bytes_transferred == 50000 diff --git a/windows-agent/agent/config.py b/windows-agent/agent/config.py index d6ea4be..e124d8f 100644 --- a/windows-agent/agent/config.py +++ b/windows-agent/agent/config.py @@ -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 diff --git a/windows-agent/agent/service.py b/windows-agent/agent/service.py index 1d1c413..4802a3c 100644 --- a/windows-agent/agent/service.py +++ b/windows-agent/agent/service.py @@ -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}") diff --git a/windows-agent/agent_app_pyqt.py b/windows-agent/agent_app_pyqt.py index 5db8e80..18ac1b5 100644 --- a/windows-agent/agent_app_pyqt.py +++ b/windows-agent/agent_app_pyqt.py @@ -210,7 +210,7 @@ class AddFolderDialog(QDialog): def __init__(self, parent=None): super().__init__(parent) self.setWindowTitle("Añadir Carpeta de Backup — OnEver Drive") - self.resize(500, 470) + self.resize(500, 520) self.setStyleSheet(DARK_QSS) layout = QVBoxLayout(self) @@ -269,6 +269,15 @@ class AddFolderDialog(QDialog): self.chk_sync_deletions.setChecked(False) form.addRow("", self.chk_sync_deletions) + dest_layout = QHBoxLayout() + self.txt_dest = QLineEdit() + self.txt_dest.setPlaceholderText("Ej: D:\\LocalBackup (Opcional)") + btn_browse_dest = QPushButton("Explorar...") + btn_browse_dest.clicked.connect(self._browse_dest_folder) + dest_layout.addWidget(self.txt_dest) + dest_layout.addWidget(btn_browse_dest) + form.addRow("Copia Local Adicional:", dest_layout) + layout.addLayout(form) btn_layout = QHBoxLayout() @@ -290,6 +299,11 @@ class AddFolderDialog(QDialog): if not self.txt_name.text(): self.txt_name.setText(Path(folder).name) + def _browse_dest_folder(self): + folder = QFileDialog.getExistingDirectory(self, "Seleccionar carpeta destino local") + if folder: + self.txt_dest.setText(folder) + def _validate_and_accept(self): if not self.txt_path.text() or not os.path.exists(self.txt_path.text()): QMessageBox.warning(self, "Ruta Inválida", "Por favor selecciona una carpeta existente en Windows.") @@ -308,6 +322,7 @@ class AddFolderDialog(QDialog): skip_system=self.chk_skip_system.isChecked(), skip_readonly=self.chk_skip_readonly.isChecked(), sync_deletions=self.chk_sync_deletions.isChecked(), + local_destination_path=self.txt_dest.text().strip() or None, schedule_cron=self.cmb_schedule.currentText(), min_stable_seconds=self.spin_stable.value() )