feat: implementar Fase 2 del roadmap (Modo Replica Exacta / Mirroring)
This commit is contained in:
+83
-2
@@ -5,11 +5,14 @@ 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
|
||||
from app.schemas.schemas import JobCreate, JobUpdate, JobResponse, AgentJobRegister
|
||||
from app.models.models import User, Client, BackupJob, BackupFile
|
||||
from app.schemas.schemas import JobCreate, JobUpdate, JobResponse, AgentJobRegister, PurgeOrphansRequest
|
||||
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
|
||||
import logging
|
||||
|
||||
logger = logging.getLogger("uvicorn.error")
|
||||
|
||||
router = APIRouter(prefix="/jobs", tags=["Backup Jobs"])
|
||||
|
||||
@@ -53,6 +56,11 @@ async def agent_register_job(
|
||||
name=payload.name,
|
||||
source_path=payload.source_path,
|
||||
file_patterns=payload.file_patterns,
|
||||
exclude_patterns=payload.exclude_patterns,
|
||||
skip_hidden=payload.skip_hidden,
|
||||
skip_system=payload.skip_system,
|
||||
skip_readonly=payload.skip_readonly,
|
||||
sync_deletions=payload.sync_deletions,
|
||||
schedule_cron=payload.schedule_cron,
|
||||
keep_daily=payload.keep_daily,
|
||||
keep_weekly=payload.keep_weekly,
|
||||
@@ -126,6 +134,11 @@ async def create_job(
|
||||
name=payload.name,
|
||||
source_path=payload.source_path,
|
||||
file_patterns=payload.file_patterns,
|
||||
exclude_patterns=payload.exclude_patterns,
|
||||
skip_hidden=payload.skip_hidden,
|
||||
skip_system=payload.skip_system,
|
||||
skip_readonly=payload.skip_readonly,
|
||||
sync_deletions=payload.sync_deletions,
|
||||
schedule_cron=payload.schedule_cron,
|
||||
keep_daily=payload.keep_daily,
|
||||
keep_weekly=payload.keep_weekly,
|
||||
@@ -241,3 +254,71 @@ async def trigger_job(
|
||||
})
|
||||
|
||||
return {"message": f"Job {job.job_code} triggered"}
|
||||
|
||||
@router.post("/{job_id}/purge-orphans")
|
||||
async def purge_orphans(
|
||||
job_id: int,
|
||||
payload: PurgeOrphansRequest,
|
||||
db: AsyncSession = Depends(get_db),
|
||||
current_client: Client = Depends(get_current_client)
|
||||
):
|
||||
"""
|
||||
Called by the Windows Agent to delete backup files on the server that no longer exist
|
||||
in the client's local directory (orphans), enforcing replica mode.
|
||||
"""
|
||||
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 this client"
|
||||
)
|
||||
|
||||
if not job.sync_deletions:
|
||||
return {"purged_count": 0, "message": "Replica mode (sync_deletions) is not enabled for this job."}
|
||||
|
||||
files_res = await db.execute(
|
||||
select(BackupFile)
|
||||
.where(
|
||||
BackupFile.job_id == job_id,
|
||||
BackupFile.client_id == current_client.id,
|
||||
BackupFile.is_active == True
|
||||
)
|
||||
)
|
||||
db_files = files_res.scalars().all()
|
||||
|
||||
client_paths = {p.strip().replace("\\", "/").lower() for p in payload.active_relative_paths}
|
||||
|
||||
purged_count = 0
|
||||
total_freed_bytes = 0
|
||||
from app.storage import storage_provider
|
||||
|
||||
for db_file in db_files:
|
||||
normalized_db_path = db_file.client_relative_path.strip().replace("\\", "/").lower()
|
||||
if normalized_db_path not in client_paths:
|
||||
try:
|
||||
await storage_provider.delete_backup_file(db_file.relative_path)
|
||||
db_file.is_active = False
|
||||
total_freed_bytes += db_file.file_size
|
||||
purged_count += 1
|
||||
|
||||
await log_event(
|
||||
db=db,
|
||||
event_type="FILE_PURGED",
|
||||
message=f"Purged orphan file '{db_file.filename}' from server storage (Replica Mode).",
|
||||
severity="WARNING",
|
||||
client_id=current_client.id,
|
||||
job_id=job.id
|
||||
)
|
||||
except Exception as e:
|
||||
logger.error(f"Failed to purge orphan file {db_file.filename}: {e}")
|
||||
|
||||
if purged_count > 0:
|
||||
current_client.storage_used_bytes = max(0, current_client.storage_used_bytes - total_freed_bytes)
|
||||
await db.commit()
|
||||
|
||||
return {
|
||||
"purged_count": purged_count,
|
||||
"freed_bytes": total_freed_bytes,
|
||||
"message": f"Successfully purged {purged_count} orphan files from server."
|
||||
}
|
||||
|
||||
@@ -36,7 +36,8 @@ async def init_session(
|
||||
file_size=payload.file_size,
|
||||
sha256_full=payload.sha256,
|
||||
chunk_size=payload.chunk_size,
|
||||
job_id=payload.job_id
|
||||
job_id=payload.job_id,
|
||||
client_relative_path=payload.client_relative_path
|
||||
)
|
||||
|
||||
return {
|
||||
|
||||
@@ -81,7 +81,14 @@ class BackupJob(Base):
|
||||
file_patterns = Column(String(255), default="*.bak,*.mdf", nullable=False)
|
||||
schedule_cron = Column(String(100), default="0 2 * * *", nullable=False) # default 02:00 AM daily
|
||||
is_active = Column(Boolean, default=True, nullable=False)
|
||||
sync_deletions = Column(Boolean, default=False, nullable=False)
|
||||
|
||||
# Advanced Filtering (Karen's Replicator Style)
|
||||
exclude_patterns = Column(String(255), default="", nullable=False)
|
||||
skip_hidden = Column(Boolean, default=False, nullable=False)
|
||||
skip_system = Column(Boolean, default=False, nullable=False)
|
||||
skip_readonly = Column(Boolean, default=False, nullable=False)
|
||||
|
||||
# Retention policies
|
||||
keep_daily = Column(Integer, default=7, nullable=False)
|
||||
keep_weekly = Column(Integer, default=4, nullable=False)
|
||||
@@ -106,6 +113,7 @@ class BackupSession(Base):
|
||||
job_id = Column(Integer, ForeignKey("backup_jobs.id", ondelete="CASCADE"), nullable=True)
|
||||
|
||||
filename = Column(String(512), nullable=False)
|
||||
client_relative_path = Column(String(1024), nullable=True)
|
||||
file_size = Column(BigInteger, nullable=False)
|
||||
chunk_size = Column(Integer, default=4 * 1024 * 1024, nullable=False)
|
||||
total_chunks = Column(Integer, nullable=False)
|
||||
@@ -152,6 +160,7 @@ class BackupFile(Base):
|
||||
|
||||
filename = Column(String(512), nullable=False)
|
||||
relative_path = Column(String(1024), nullable=False)
|
||||
client_relative_path = Column(String(1024), default="", nullable=False)
|
||||
file_size = Column(BigInteger, nullable=False)
|
||||
sha256 = Column(String(64), nullable=False)
|
||||
retention_tag = Column(String(50), default="DAILY", nullable=False) # DAILY, WEEKLY, MONTHLY, MANUAL
|
||||
|
||||
@@ -82,6 +82,11 @@ class JobCreate(BaseModel):
|
||||
name: str
|
||||
source_path: str
|
||||
file_patterns: str = "*.bak,*.mdf"
|
||||
exclude_patterns: str = ""
|
||||
skip_hidden: bool = False
|
||||
skip_system: bool = False
|
||||
skip_readonly: bool = False
|
||||
sync_deletions: bool = False
|
||||
schedule_cron: str = "0 2 * * *"
|
||||
keep_daily: int = 7
|
||||
keep_weekly: int = 4
|
||||
@@ -92,6 +97,11 @@ class JobUpdate(BaseModel):
|
||||
name: Optional[str] = None
|
||||
source_path: Optional[str] = None
|
||||
file_patterns: Optional[str] = None
|
||||
exclude_patterns: Optional[str] = None
|
||||
skip_hidden: Optional[bool] = None
|
||||
skip_system: Optional[bool] = None
|
||||
skip_readonly: Optional[bool] = None
|
||||
sync_deletions: Optional[bool] = None
|
||||
schedule_cron: Optional[str] = None
|
||||
is_active: Optional[bool] = None
|
||||
keep_daily: Optional[int] = None
|
||||
@@ -106,6 +116,11 @@ class JobResponse(BaseModel):
|
||||
name: str
|
||||
source_path: str
|
||||
file_patterns: str
|
||||
exclude_patterns: str
|
||||
skip_hidden: bool
|
||||
skip_system: bool
|
||||
skip_readonly: bool
|
||||
sync_deletions: bool
|
||||
schedule_cron: str
|
||||
is_active: bool
|
||||
keep_daily: int
|
||||
@@ -122,6 +137,7 @@ class JobResponse(BaseModel):
|
||||
# --- Upload Session & Chunk Schemas ---
|
||||
class UploadSessionInitRequest(BaseModel):
|
||||
filename: str
|
||||
client_relative_path: Optional[str] = None
|
||||
file_size: int
|
||||
sha256: str
|
||||
chunk_size: int = 4 * 1024 * 1024
|
||||
@@ -234,8 +250,16 @@ class AgentJobRegister(BaseModel):
|
||||
name: str
|
||||
source_path: str
|
||||
file_patterns: str = "*.bak,*.mdf"
|
||||
exclude_patterns: str = ""
|
||||
skip_hidden: bool = False
|
||||
skip_system: bool = False
|
||||
skip_readonly: bool = False
|
||||
sync_deletions: bool = False
|
||||
schedule_cron: str = "daily"
|
||||
keep_daily: int = 7
|
||||
keep_weekly: int = 4
|
||||
keep_monthly: int = 12
|
||||
min_stable_time_seconds: int = 60
|
||||
|
||||
class PurgeOrphansRequest(BaseModel):
|
||||
active_relative_paths: List[str]
|
||||
|
||||
@@ -16,7 +16,8 @@ async def create_or_resume_session(
|
||||
file_size: int,
|
||||
sha256_full: str,
|
||||
chunk_size: int = 4 * 1024 * 1024,
|
||||
job_id: Optional[int] = None
|
||||
job_id: Optional[int] = None,
|
||||
client_relative_path: Optional[str] = None
|
||||
) -> Tuple[BackupSession, List[int]]:
|
||||
"""
|
||||
Initializes a new upload session or resumes an existing incomplete session
|
||||
@@ -63,6 +64,8 @@ async def create_or_resume_session(
|
||||
# Resume existing session
|
||||
received_chunks = await storage_provider.get_received_chunks(session.session_code)
|
||||
session.received_chunks_count = len(received_chunks)
|
||||
if client_relative_path:
|
||||
session.client_relative_path = client_relative_path
|
||||
await db.commit()
|
||||
await db.refresh(session)
|
||||
return session, received_chunks
|
||||
@@ -72,6 +75,7 @@ async def create_or_resume_session(
|
||||
client_id=client.id,
|
||||
job_id=job_id,
|
||||
filename=filename,
|
||||
client_relative_path=client_relative_path or filename,
|
||||
file_size=file_size,
|
||||
chunk_size=chunk_size,
|
||||
total_chunks=total_chunks,
|
||||
@@ -242,6 +246,7 @@ async def complete_session(
|
||||
session_id=session.id,
|
||||
filename=session.filename,
|
||||
relative_path=rel_path,
|
||||
client_relative_path=session.client_relative_path or session.filename,
|
||||
file_size=total_bytes,
|
||||
sha256=final_sha256,
|
||||
retention_tag="DAILY",
|
||||
|
||||
@@ -0,0 +1,29 @@
|
||||
import sqlite3
|
||||
import os
|
||||
|
||||
DB_PATH = os.path.join(os.path.dirname(__file__), "onever_drive.db")
|
||||
|
||||
def migrate():
|
||||
print(f"Connecting to {DB_PATH}...")
|
||||
conn = sqlite3.connect(DB_PATH)
|
||||
cursor = conn.cursor()
|
||||
|
||||
columns_to_add = [
|
||||
("exclude_patterns", "VARCHAR(255) DEFAULT '' NOT NULL"),
|
||||
("skip_hidden", "BOOLEAN DEFAULT 0 NOT NULL"),
|
||||
("skip_system", "BOOLEAN DEFAULT 0 NOT NULL"),
|
||||
("skip_readonly", "BOOLEAN DEFAULT 0 NOT NULL")
|
||||
]
|
||||
|
||||
for col_name, col_def in columns_to_add:
|
||||
try:
|
||||
cursor.execute(f"ALTER TABLE backup_jobs ADD COLUMN {col_name} {col_def}")
|
||||
print(f"Column '{col_name}' added successfully to backup_jobs.")
|
||||
except sqlite3.OperationalError as e:
|
||||
print(f"Column '{col_name}' could not be added (maybe it already exists?): {e}")
|
||||
|
||||
conn.commit()
|
||||
conn.close()
|
||||
|
||||
if __name__ == "__main__":
|
||||
migrate()
|
||||
@@ -0,0 +1,36 @@
|
||||
import sqlite3
|
||||
import os
|
||||
|
||||
DB_PATH = os.path.join(os.path.dirname(__file__), "onever_drive.db")
|
||||
|
||||
def migrate():
|
||||
print(f"Connecting to {DB_PATH}...")
|
||||
conn = sqlite3.connect(DB_PATH)
|
||||
cursor = conn.cursor()
|
||||
|
||||
# Add sync_deletions to backup_jobs
|
||||
try:
|
||||
cursor.execute("ALTER TABLE backup_jobs ADD COLUMN sync_deletions BOOLEAN DEFAULT 0 NOT NULL")
|
||||
print("Column 'sync_deletions' added to backup_jobs.")
|
||||
except sqlite3.OperationalError as e:
|
||||
print(f"Could not add 'sync_deletions' (maybe it already exists?): {e}")
|
||||
|
||||
# Add client_relative_path to backup_sessions
|
||||
try:
|
||||
cursor.execute("ALTER TABLE backup_sessions ADD COLUMN client_relative_path VARCHAR(1024) NULL")
|
||||
print("Column 'client_relative_path' added to backup_sessions.")
|
||||
except sqlite3.OperationalError as e:
|
||||
print(f"Could not add 'client_relative_path' to backup_sessions (maybe it already exists?): {e}")
|
||||
|
||||
# Add client_relative_path to backup_files
|
||||
try:
|
||||
cursor.execute("ALTER TABLE backup_files ADD COLUMN client_relative_path VARCHAR(1024) DEFAULT '' NOT NULL")
|
||||
print("Column 'client_relative_path' added to backup_files.")
|
||||
except sqlite3.OperationalError as e:
|
||||
print(f"Could not add 'client_relative_path' to backup_files (maybe it already exists?): {e}")
|
||||
|
||||
conn.commit()
|
||||
conn.close()
|
||||
|
||||
if __name__ == "__main__":
|
||||
migrate()
|
||||
Reference in New Issue
Block a user