feat: script jobs scheduler and manager page for Fortigate, UniFi and Grandstream scripts
This commit is contained in:
+88
-1
@@ -1,4 +1,6 @@
|
||||
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
|
||||
@@ -16,12 +18,88 @@ 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:
|
||||
@@ -59,7 +137,15 @@ async def lifespan(app: FastAPI):
|
||||
await session.commit()
|
||||
print(">> [OnEver Drive] Default admin verified & password reset to: Admin1234!")
|
||||
|
||||
yield
|
||||
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,
|
||||
@@ -87,6 +173,7 @@ 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):
|
||||
|
||||
Reference in New Issue
Block a user