from contextlib import asynccontextmanager from datetime import datetime import asyncio from fastapi import FastAPI, WebSocket, WebSocketDisconnect from fastapi.middleware.cors import CORSMiddleware from sqlalchemy import select from app.core.config import settings from app.core.database import init_db, AsyncSessionLocal from app.core.security import get_password_hash from app.models.models import User from app.api.auth import router as auth_router from app.api.clients import router as clients_router from app.api.jobs import router as jobs_router from app.api.upload import router as upload_router from app.api.backups import router as backups_router from app.api.events import router as events_router from app.api.stats import router as stats_router from app.api.settings import router as settings_router from app.api.users import router as users_router from app.api.script_jobs import router as script_jobs_router from app.ws.manager import ws_manager async def script_scheduler_daemon(): """ Background daemon that runs every minute to execute scheduled script jobs. """ from app.models.models import ScriptJob from app.api.script_jobs import execute_script_in_background from app.core.database import AsyncSessionLocal from sqlalchemy import select print("[*] Script Scheduler Daemon started.") def match_cron(cron_expr: str, dt: datetime) -> bool: try: fields = cron_expr.strip().split() if len(fields) < 5: return False minute, hour, dom, month, dow = fields def match_field(val: int, field: str, dow_check=False) -> bool: if field == '*': return True if '/' in field: base, step = field.split('/') step = int(step) if base == '*': return val % step == 0 return (val - int(base)) % step == 0 if ',' in field: parts = field.split(',') return any(match_field(val, p, dow_check) for p in parts) if '-' in field: start, end = map(int, field.split('-')) return start <= val <= end if dow_check: cron_dow = (dt.weekday() + 1) % 7 if int(field) == 7 and cron_dow == 0: return True return cron_dow == int(field) return val == int(field) return ( match_field(dt.minute, minute) and match_field(dt.hour, hour) and match_field(dt.day, dom) and match_field(dt.month, month) and match_field(dt.weekday(), dow, dow_check=True) ) except Exception: return False while True: try: now_sec = datetime.now().second sleep_time = 60 - now_sec await asyncio.sleep(sleep_time) now = datetime.now() async with AsyncSessionLocal() as db: result = await db.execute(select(ScriptJob).where(ScriptJob.is_active == True)) jobs = result.scalars().all() for job in jobs: if match_cron(job.schedule_cron, now): print(f"[*] Scheduler: Script Job '{job.name}' is due. Triggering execution...") asyncio.create_task(execute_script_in_background(job.id)) except asyncio.CancelledError: break except Exception as e: print(f"[!] Scheduler error: {e}") await asyncio.sleep(5) @asynccontextmanager async def lifespan(app: FastAPI): # Initialize database tables await init_db() # Start scheduler daemon task scheduler_task = asyncio.create_task(script_scheduler_daemon()) # Seed default administrator and initialize settings async with AsyncSessionLocal() as session: # Load and apply settings to storage_provider from app.services.settings_service import get_all_settings settings_dict = await get_all_settings(session) # Apply storage_root from app.storage.local import storage_provider from pathlib import Path import os storage_root = settings_dict.get("storage_root") if storage_root: path = Path(storage_root).resolve() os.makedirs(path, exist_ok=True) storage_provider.root_dir = path result = await session.execute(select(User).where(User.email == "admin@oneverdrive.local")) admin = result.scalar_one_or_none() if not admin: admin_user = User( email="admin@oneverdrive.local", hashed_password=get_password_hash("Admin1234!"), full_name="System Administrator", role="ADMIN", is_active=True ) session.add(admin_user) await session.commit() print(">> [OnEver Drive] Default admin created: admin@oneverdrive.local / Admin1234!") else: admin.hashed_password = get_password_hash("Admin1234!") admin.is_active = True session.add(admin) await session.commit() print(">> [OnEver Drive] Default admin verified & password reset to: Admin1234!") try: yield finally: scheduler_task.cancel() try: await scheduler_task except asyncio.CancelledError: pass print("[*] Script Scheduler Daemon stopped.") app = FastAPI( title=settings.PROJECT_NAME, version=settings.VERSION, description="Centralized Backup & Sync Platform for Windows on Proxmox VE", lifespan=lifespan ) # CORS Middleware to allow Web UI connections app.add_middleware( CORSMiddleware, allow_origins=["*"], allow_credentials=True, allow_methods=["*"], allow_headers=["*"], ) # Register API Routers app.include_router(auth_router, prefix=settings.API_V1_PREFIX) app.include_router(clients_router, prefix=settings.API_V1_PREFIX) app.include_router(jobs_router, prefix=settings.API_V1_PREFIX) app.include_router(upload_router, prefix=settings.API_V1_PREFIX) app.include_router(backups_router, prefix=settings.API_V1_PREFIX) app.include_router(events_router, prefix=settings.API_V1_PREFIX) app.include_router(stats_router, prefix=settings.API_V1_PREFIX) app.include_router(settings_router, prefix=settings.API_V1_PREFIX) app.include_router(users_router, prefix=settings.API_V1_PREFIX) app.include_router(script_jobs_router, prefix=settings.API_V1_PREFIX) @app.websocket("/ws/telemetry") async def websocket_telemetry(websocket: WebSocket): """WebSocket endpoint for real-time dashboard telemetry and live upload meters.""" await ws_manager.connect(websocket) try: while True: # Keep connection open and accept incoming ping/pong or messages data = await websocket.receive_text() if data == "ping": await websocket.send_text("pong") except WebSocketDisconnect: ws_manager.disconnect(websocket) except Exception: ws_manager.disconnect(websocket) @app.get("/health") async def health(): return { "status": "healthy", "service": settings.PROJECT_NAME, "version": settings.VERSION }