From a51fccef97c79b724ebf966be79bd63eea416bd0 Mon Sep 17 00:00:00 2001 From: Aditya Vaish Date: Thu, 9 Oct 2025 19:30:16 +0530 Subject: [PATCH] init commit syncing features --- .../23_add_connector_schedules_table.py | 96 ++++++++ surfsense_backend/app/db.py | 47 ++++ .../app/routes/connector_schedules_routes.py | 206 ++++++++++++++++++ surfsense_backend/app/schemas/__init__.py | 10 + .../app/schemas/connector_schedule.py | 65 ++++++ .../app/utils/schedule_helpers.py | 76 +++++++ 6 files changed, 500 insertions(+) create mode 100644 surfsense_backend/alembic/versions/23_add_connector_schedules_table.py create mode 100644 surfsense_backend/app/routes/connector_schedules_routes.py create mode 100644 surfsense_backend/app/schemas/connector_schedule.py create mode 100644 surfsense_backend/app/utils/schedule_helpers.py diff --git a/surfsense_backend/alembic/versions/23_add_connector_schedules_table.py b/surfsense_backend/alembic/versions/23_add_connector_schedules_table.py new file mode 100644 index 000000000..edd8a9cce --- /dev/null +++ b/surfsense_backend/alembic/versions/23_add_connector_schedules_table.py @@ -0,0 +1,96 @@ +"""Add connector schedules table + +Revision ID: 23 +Revises: 22 +""" + +from collections.abc import Sequence + +from sqlalchemy import inspect + +from alembic import op + +# revision identifiers, used by Alembic. +revision: str = "23" +down_revision: str | None = "22" +branch_labels: str | Sequence[str] | None = None +depends_on: str | Sequence[str] | None = None + + +def upgrade() -> None: + """Upgrade schema - add ScheduleType enum and connector_schedules table.""" + + # Create ScheduleType enum if it doesn't exist + op.execute( + """ + DO $$ + BEGIN + IF NOT EXISTS (SELECT 1 FROM pg_type WHERE typname = 'scheduletype') THEN + CREATE TYPE scheduletype AS ENUM ('HOURLY', 'DAILY', 'WEEKLY', 'CUSTOM'); + END IF; + END$$; + """ + ) + + # Create connector_schedules table if it doesn't exist + op.execute( + """ + CREATE TABLE IF NOT EXISTS connector_schedules ( + id SERIAL PRIMARY KEY, + created_at TIMESTAMPTZ NOT NULL DEFAULT NOW(), + connector_id INTEGER NOT NULL REFERENCES search_source_connectors(id) ON DELETE CASCADE, + search_space_id INTEGER NOT NULL REFERENCES searchspaces(id) ON DELETE CASCADE, + schedule_type scheduletype NOT NULL, + cron_expression VARCHAR(100), + is_active BOOLEAN NOT NULL DEFAULT TRUE, + last_run_at TIMESTAMPTZ, + next_run_at TIMESTAMPTZ, + CONSTRAINT uq_connector_search_space UNIQUE (connector_id, search_space_id) + ); + """ + ) + + # Get existing indexes + conn = op.get_bind() + inspector = inspect(conn) + existing_indexes = [ + idx["name"] for idx in inspector.get_indexes("connector_schedules") + ] + + # Create indexes only if they don't already exist + if "ix_connector_schedules_id" not in existing_indexes: + op.create_index("ix_connector_schedules_id", "connector_schedules", ["id"]) + if "ix_connector_schedules_created_at" not in existing_indexes: + op.create_index( + "ix_connector_schedules_created_at", "connector_schedules", ["created_at"] + ) + if "ix_connector_schedules_connector_id" not in existing_indexes: + op.create_index( + "ix_connector_schedules_connector_id", "connector_schedules", ["connector_id"] + ) + if "ix_connector_schedules_is_active" not in existing_indexes: + op.create_index( + "ix_connector_schedules_is_active", "connector_schedules", ["is_active"] + ) + if "ix_connector_schedules_next_run_at" not in existing_indexes: + op.create_index( + "ix_connector_schedules_next_run_at", "connector_schedules", ["next_run_at"] + ) + + +def downgrade() -> None: + """Downgrade schema - remove connector_schedules table and enum.""" + + # Drop indexes + op.drop_index("ix_connector_schedules_next_run_at", table_name="connector_schedules") + op.drop_index("ix_connector_schedules_is_active", table_name="connector_schedules") + op.drop_index("ix_connector_schedules_connector_id", table_name="connector_schedules") + op.drop_index("ix_connector_schedules_created_at", table_name="connector_schedules") + op.drop_index("ix_connector_schedules_id", table_name="connector_schedules") + + # Drop table + op.drop_table("connector_schedules") + + # Drop enum + op.execute("DROP TYPE IF EXISTS scheduletype") + diff --git a/surfsense_backend/app/db.py b/surfsense_backend/app/db.py index 99be9f277..a67c27a61 100644 --- a/surfsense_backend/app/db.py +++ b/surfsense_backend/app/db.py @@ -118,6 +118,13 @@ class LogStatus(str, Enum): FAILED = "FAILED" +class ScheduleType(str, Enum): + HOURLY = "HOURLY" + DAILY = "DAILY" + WEEKLY = "WEEKLY" + CUSTOM = "CUSTOM" + + class Base(DeclarativeBase): pass @@ -234,6 +241,12 @@ class SearchSpace(BaseModel, TimestampMixin): order_by="Log.id", cascade="all, delete-orphan", ) + connector_schedules = relationship( + "ConnectorSchedule", + back_populates="search_space", + order_by="ConnectorSchedule.id", + cascade="all, delete-orphan", + ) class SearchSourceConnector(BaseModel, TimestampMixin): @@ -252,6 +265,40 @@ class SearchSourceConnector(BaseModel, TimestampMixin): UUID(as_uuid=True), ForeignKey("user.id", ondelete="CASCADE"), nullable=False ) user = relationship("User", back_populates="search_source_connectors") + schedules = relationship( + "ConnectorSchedule", + back_populates="connector", + cascade="all, delete-orphan", + ) + + +class ConnectorSchedule(BaseModel, TimestampMixin): + __tablename__ = "connector_schedules" + __table_args__ = ( + UniqueConstraint( + "connector_id", "search_space_id", name="uq_connector_search_space" + ), + ) + + connector_id = Column( + Integer, + ForeignKey("search_source_connectors.id", ondelete="CASCADE"), + nullable=False, + index=True, + ) + search_space_id = Column( + Integer, + ForeignKey("searchspaces.id", ondelete="CASCADE"), + nullable=False, + ) + schedule_type = Column(SQLAlchemyEnum(ScheduleType), nullable=False) + cron_expression = Column(String(100), nullable=True) + is_active = Column(Boolean, nullable=False, default=True, index=True) + last_run_at = Column(TIMESTAMP(timezone=True), nullable=True) + next_run_at = Column(TIMESTAMP(timezone=True), nullable=True, index=True) + + connector = relationship("SearchSourceConnector", back_populates="schedules") + search_space = relationship("SearchSpace", back_populates="connector_schedules") class LLMConfig(BaseModel, TimestampMixin): diff --git a/surfsense_backend/app/routes/connector_schedules_routes.py b/surfsense_backend/app/routes/connector_schedules_routes.py new file mode 100644 index 000000000..eaa3333e8 --- /dev/null +++ b/surfsense_backend/app/routes/connector_schedules_routes.py @@ -0,0 +1,206 @@ +""" +ConnectorSchedule routes for CRUD operations: +POST /connector-schedules/ - Create a new schedule +GET /connector-schedules/ - List all schedules for the current user +GET /connector-schedules/{schedule_id} - Get a specific schedule +PUT /connector-schedules/{schedule_id} - Update a specific schedule +DELETE /connector-schedules/{schedule_id} - Delete a specific schedule +PATCH /connector-schedules/{schedule_id}/toggle - Activate/deactivate a schedule +""" + +import logging +from typing import Any + +from fastapi import APIRouter, Depends, HTTPException +from sqlalchemy.exc import IntegrityError +from sqlalchemy.ext.asyncio import AsyncSession +from sqlalchemy.future import select + +from app.db import ( + ConnectorSchedule, + SearchSourceConnector, + SearchSpace, + User, + get_async_session, +) +from app.schemas import ( + ConnectorScheduleCreate, + ConnectorScheduleRead, + ConnectorScheduleUpdate, +) +from app.users import current_active_user +from app.utils.check_ownership import check_ownership +from app.utils.schedule_helpers import calculate_next_run + +# Set up logging +logger = logging.getLogger(__name__) + +router = APIRouter() + + +@router.post("/connector-schedules/", response_model=ConnectorScheduleRead) +async def create_connector_schedule( + schedule: ConnectorScheduleCreate, + session: AsyncSession = Depends(get_async_session), + user: User = Depends(current_active_user), +): + """ + Create a new connector schedule. + + Each connector can have only one schedule per search space. + The schedule will automatically calculate the next run time based on the schedule type. + """ + try: + # Verify connector belongs to user + connector = await check_ownership( + session, SearchSourceConnector, schedule.connector_id, user + ) + + # Verify connector is indexable + if not connector.is_indexable: + raise HTTPException( + status_code=400, + detail=f"Connector {connector.name} is not indexable and cannot be scheduled", + ) + + # Verify search space belongs to user + await check_ownership(session, SearchSpace, schedule.search_space_id, user) + + # Check if schedule already exists for this connector-space pair + result = await session.execute( + select(ConnectorSchedule).filter( + ConnectorSchedule.connector_id == schedule.connector_id, + ConnectorSchedule.search_space_id == schedule.search_space_id, + ) + ) + existing_schedule = result.scalars().first() + if existing_schedule: + raise HTTPException( + status_code=409, + detail=f"A schedule already exists for connector {schedule.connector_id} and search space {schedule.search_space_id}", + ) + + # Calculate next run time + next_run_at = calculate_next_run( + schedule.schedule_type, schedule.cron_expression + ) + + # Create schedule + db_schedule = ConnectorSchedule( + **schedule.model_dump(), next_run_at=next_run_at + ) + session.add(db_schedule) + await session.commit() + await session.refresh(db_schedule) + + logger.info( + f"Created schedule {db_schedule.id} for connector {schedule.connector_id} (next run: {next_run_at})" + ) + return db_schedule + + except HTTPException: + await session.rollback() + raise + except IntegrityError as e: + await session.rollback() + raise HTTPException( + status_code=409, + detail=f"A schedule already exists for this connector and search space: {e!s}", + ) from e + except Exception as e: + await session.rollback() + logger.error(f"Failed to create connector schedule: {e!s}") + raise HTTPException( + status_code=500, detail=f"Failed to create connector schedule: {e!s}" + ) from e + + +@router.get("/connector-schedules/", response_model=list[ConnectorScheduleRead]) +async def read_connector_schedules( + skip: int = 0, + limit: int = 100, + connector_id: int | None = None, + search_space_id: int | None = None, + is_active: bool | None = None, + session: AsyncSession = Depends(get_async_session), + user: User = Depends(current_active_user), +): + """ + List all connector schedules for the current user. + + Optional filters: + - connector_id: Filter by specific connector + - search_space_id: Filter by specific search space + - is_active: Filter by active/inactive status + """ + try: + # Build query to get schedules for connectors owned by user + query = ( + select(ConnectorSchedule) + .join(SearchSourceConnector) + .filter(SearchSourceConnector.user_id == user.id) + ) + + # Apply filters + if connector_id is not None: + query = query.filter(ConnectorSchedule.connector_id == connector_id) + if search_space_id is not None: + query = query.filter(ConnectorSchedule.search_space_id == search_space_id) + if is_active is not None: + query = query.filter(ConnectorSchedule.is_active == is_active) + + result = await session.execute(query.offset(skip).limit(limit)) + return result.scalars().all() + + except Exception as e: + logger.error(f"Failed to fetch connector schedules: {e!s}") + raise HTTPException( + status_code=500, detail=f"Failed to fetch connector schedules: {e!s}" + ) from e + + +@router.get("/connector-schedules/{schedule_id}", response_model=ConnectorScheduleRead) +async def read_connector_schedule( + schedule_id: int, + session: AsyncSession = Depends(get_async_session), + user: User = Depends(current_active_user), +): + """Get a specific connector schedule by ID.""" + try: + # Get schedule + result = await session.execute( + select(ConnectorSchedule).filter(ConnectorSchedule.id == schedule_id) + ) + schedule = result.scalars().first() + + if not schedule: + raise HTTPException(status_code=404, detail="Schedule not found") + + # Verify schedule's connector belongs to user + await check_ownership(session, SearchSourceConnector, schedule.connector_id, user) + + return schedule + + except HTTPException: + raise + except Exception as e: + logger.error(f"Failed to fetch connector schedule: {e!s}") + raise HTTPException( + status_code=500, detail=f"Failed to fetch connector schedule: {e!s}" + ) from e + + +@router.put("/connector-schedules/{schedule_id}", response_model=ConnectorScheduleRead) +async def update_connector_schedule( + schedule_id: int, + schedule_update: ConnectorScheduleUpdate, + session: AsyncSession = Depends(get_async_session), + user: User = Depends(current_active_user), +): + """ + Update a connector schedule. + + Can update schedule_type, cron_expression, and is_active. + If schedule_type changes, next_run_at is recalculated automatically. + """ + \ No newline at end of file diff --git a/surfsense_backend/app/schemas/__init__.py b/surfsense_backend/app/schemas/__init__.py index 41b2ce23c..adcba71ef 100644 --- a/surfsense_backend/app/schemas/__init__.py +++ b/surfsense_backend/app/schemas/__init__.py @@ -8,6 +8,12 @@ from .chats import ( ChatUpdate, ) from .chunks import ChunkBase, ChunkCreate, ChunkRead, ChunkUpdate +from .connector_schedule import ( + ConnectorScheduleBase, + ConnectorScheduleCreate, + ConnectorScheduleRead, + ConnectorScheduleUpdate, +) from .documents import ( DocumentBase, DocumentRead, @@ -52,6 +58,10 @@ __all__ = [ "ChunkCreate", "ChunkRead", "ChunkUpdate", + "ConnectorScheduleBase", + "ConnectorScheduleCreate", + "ConnectorScheduleRead", + "ConnectorScheduleUpdate", "DocumentBase", "DocumentRead", "DocumentUpdate", diff --git a/surfsense_backend/app/schemas/connector_schedule.py b/surfsense_backend/app/schemas/connector_schedule.py new file mode 100644 index 000000000..6a4d84d92 --- /dev/null +++ b/surfsense_backend/app/schemas/connector_schedule.py @@ -0,0 +1,65 @@ +from datetime import datetime + +from pydantic import BaseModel, ConfigDict, field_validator + +from app.db import ScheduleType + +from .base import IDModel, TimestampModel + + +class ConnectorScheduleBase(BaseModel): + connector_id: int + search_space_id: int + schedule_type: ScheduleType + cron_expression: str | None = None + is_active: bool = True + + @field_validator("cron_expression") + @classmethod + def validate_cron_expression(cls, v: str | None, values: dict) -> str | None: + """Validate cron expression is provided when schedule_type is CUSTOM.""" + schedule_type = values.data.get("schedule_type") + if schedule_type == ScheduleType.CUSTOM and not v: + raise ValueError( + "cron_expression is required when schedule_type is CUSTOM" + ) + if schedule_type != ScheduleType.CUSTOM and v: + raise ValueError( + f"cron_expression should only be provided for CUSTOM schedule_type, got {schedule_type}" + ) + return v + + +class ConnectorScheduleCreate(ConnectorScheduleBase): + """Schema for creating a new connector schedule.""" + + pass + + +class ConnectorScheduleUpdate(BaseModel): + """Schema for updating an existing connector schedule.""" + + schedule_type: ScheduleType | None = None + cron_expression: str | None = None + is_active: bool | None = None + + @field_validator("cron_expression") + @classmethod + def validate_cron_expression_update(cls, v: str | None, values: dict) -> str | None: + """Validate cron expression for updates.""" + schedule_type = values.data.get("schedule_type") + if schedule_type == ScheduleType.CUSTOM and v is None: + raise ValueError( + "cron_expression is required when schedule_type is CUSTOM" + ) + return v + + +class ConnectorScheduleRead(ConnectorScheduleBase, IDModel, TimestampModel): + """Schema for reading a connector schedule.""" + + last_run_at: datetime | None = None + next_run_at: datetime | None = None + + model_config = ConfigDict(from_attributes=True) + diff --git a/surfsense_backend/app/utils/schedule_helpers.py b/surfsense_backend/app/utils/schedule_helpers.py new file mode 100644 index 000000000..b9970a699 --- /dev/null +++ b/surfsense_backend/app/utils/schedule_helpers.py @@ -0,0 +1,76 @@ +"""Helper utilities for calculating schedule next run times.""" + +from datetime import datetime, timedelta + +from croniter import croniter + +from app.db import ScheduleType + + +def calculate_next_run( + schedule_type: ScheduleType, cron_expression: str | None = None +) -> datetime: + """ + Calculate the next run time based on schedule type. + + Args: + schedule_type: The type of schedule (HOURLY, DAILY, WEEKLY, CUSTOM) + cron_expression: Optional cron expression for CUSTOM type + + Returns: + datetime: The next scheduled run time + + Raises: + ValueError: If schedule_type is CUSTOM but no cron_expression provided + """ + now = datetime.now() + + if schedule_type == ScheduleType.HOURLY: + # Run at the top of the next hour + return (now + timedelta(hours=1)).replace(minute=0, second=0, microsecond=0) + + elif schedule_type == ScheduleType.DAILY: + # Run at 2 AM next day (off-peak hours) + next_run = now.replace(hour=2, minute=0, second=0, microsecond=0) + if next_run <= now: + next_run += timedelta(days=1) + return next_run + + elif schedule_type == ScheduleType.WEEKLY: + # Run on Sunday at 2 AM + next_run = now.replace(hour=2, minute=0, second=0, microsecond=0) + days_until_sunday = (6 - now.weekday()) % 7 + if days_until_sunday == 0 and next_run <= now: + days_until_sunday = 7 + next_run += timedelta(days=days_until_sunday) + return next_run + + elif schedule_type == ScheduleType.CUSTOM: + if not cron_expression: + raise ValueError("cron_expression is required for CUSTOM schedule type") + try: + cron = croniter(cron_expression, now) + return cron.get_next(datetime) + except Exception as e: + raise ValueError(f"Invalid cron expression: {cron_expression}") from e + + else: + raise ValueError(f"Unknown schedule type: {schedule_type}") + + +def is_valid_cron_expression(expression: str) -> bool: + """ + Validate a cron expression. + + Args: + expression: The cron expression to validate + + Returns: + bool: True if valid, False otherwise + """ + try: + croniter(expression) + return True + except Exception: + return False +