From 7391a347cac0fbd67d8479c31b40596c23964fab Mon Sep 17 00:00:00 2001 From: Aditya Vaish Date: Wed, 15 Oct 2025 14:54:15 +0530 Subject: [PATCH] issues resolved --- .../app/routes/connector_schedules_routes.py | 10 +++- .../app/routes/scheduler_routes.py | 28 +++++----- .../app/schemas/connector_schedule.py | 54 +++++++++++++++++++ .../services/connector_scheduler_service.py | 4 +- 4 files changed, 79 insertions(+), 17 deletions(-) diff --git a/surfsense_backend/app/routes/connector_schedules_routes.py b/surfsense_backend/app/routes/connector_schedules_routes.py index 0b73ae1bf..5922f4db7 100644 --- a/surfsense_backend/app/routes/connector_schedules_routes.py +++ b/surfsense_backend/app/routes/connector_schedules_routes.py @@ -242,10 +242,16 @@ async def update_connector_schedule( # Update fields that were provided update_data = schedule_update.model_dump(exclude_unset=True) + # Check if any time-related fields changed (requiring next_run recalc) + time_fields_changed = any( + field in update_data + for field in ["daily_time", "weekly_day", "weekly_time", "hourly_minute"] + ) + # If schedule_type is being updated, recalculate next_run_at - if "schedule_type" in update_data: + if "schedule_type" in update_data or time_fields_changed: # Use the new schedule_type and existing values for calculation - new_schedule_type = update_data["schedule_type"] + new_schedule_type = update_data.get("schedule_type", schedule.schedule_type) cron_expr = update_data.get("cron_expression", schedule.cron_expression) daily_time = update_data.get("daily_time", schedule.daily_time) weekly_day = update_data.get("weekly_day", schedule.weekly_day) diff --git a/surfsense_backend/app/routes/scheduler_routes.py b/surfsense_backend/app/routes/scheduler_routes.py index 16ed3c4fc..0821de37b 100644 --- a/surfsense_backend/app/routes/scheduler_routes.py +++ b/surfsense_backend/app/routes/scheduler_routes.py @@ -10,7 +10,7 @@ from sqlalchemy import select, Integer from sqlalchemy.ext.asyncio import AsyncSession from sqlalchemy.orm import selectinload -from app.db import get_async_session +from app.db import get_async_session, Log from app.schemas import ConnectorScheduleRead from app.services.connector_scheduler_service import get_scheduler from app.users import User, current_active_user @@ -161,35 +161,37 @@ async def get_recent_schedule_executions( Shows the execution history for monitoring and debugging. """ try: - from app.db import ConnectorSchedule, SearchSourceConnector, logs + from app.db import ConnectorSchedule, SearchSourceConnector # Get recent executions from the logs table query = ( select(Log) - .join(SearchSourceConnector, Log.metadata["connector_id"].astext.cast(Integer) == SearchSourceConnector.id) + .join(SearchSourceConnector, Log.log_metadata["connector_id"].astext.cast(Integer) == SearchSourceConnector.id) .filter( SearchSourceConnector.user_id == user.id, - Log.task_name.like("scheduled_sync_%"), + Log.message.like("Scheduled sync%"), ) .order_by(Log.created_at.desc()) .limit(limit) ) result = await session.execute(query) - logs = result.scalars().all() + log_rows = result.scalars().all() executions = [] - for log in logs: + for log in log_rows: executions.append({ "log_id": log.id, - "task_name": log.task_name, - "status": log.status, + "task_name": log.log_metadata.get("task_name") if log.log_metadata else None, + "status": log.status.value, + "level": log.level.value, + "message": log.message, + "source": log.source, "created_at": log.created_at, - "completed_at": log.completed_at, - "connector_id": log.metadata.get("connector_id") if log.metadata else None, - "schedule_id": log.metadata.get("schedule_id") if log.metadata else None, - "error_message": log.error_message, - "documents_processed": log.metadata.get("documents_processed") if log.metadata else None, + "search_space_id": log.search_space_id, + "connector_id": log.log_metadata.get("connector_id") if log.log_metadata else None, + "schedule_id": log.log_metadata.get("schedule_id") if log.log_metadata else None, + "documents_processed": log.log_metadata.get("documents_processed") if log.log_metadata else None, }) return executions diff --git a/surfsense_backend/app/schemas/connector_schedule.py b/surfsense_backend/app/schemas/connector_schedule.py index 559dcb338..9467bfdae 100644 --- a/surfsense_backend/app/schemas/connector_schedule.py +++ b/surfsense_backend/app/schemas/connector_schedule.py @@ -97,6 +97,12 @@ class ConnectorScheduleUpdate(BaseModel): schedule_type: ScheduleType | None = None cron_expression: str | None = None is_active: bool | None = None + + # Enhanced time selection options for updates + daily_time: Optional[time] = None # For DAILY schedules + weekly_day: Optional[int] = None # For WEEKLY schedules (0=Monday, 6=Sunday) + weekly_time: Optional[time] = None # For WEEKLY schedules + hourly_minute: Optional[int] = None # For HOURLY schedules (0-59) @field_validator("cron_expression") @classmethod @@ -108,6 +114,54 @@ class ConnectorScheduleUpdate(BaseModel): "cron_expression is required when schedule_type is CUSTOM" ) return v + + @field_validator("daily_time") + @classmethod + def validate_daily_time_update(cls, v: time | None, info: FieldValidationInfo) -> time | None: + """Validate daily_time is only provided for DAILY schedule type.""" + schedule_type = info.data.get("schedule_type") + if v is not None and schedule_type != ScheduleType.DAILY: + raise ValueError( + "daily_time should only be provided for DAILY schedule_type" + ) + return v + + @field_validator("weekly_day") + @classmethod + def validate_weekly_day_update(cls, v: int | None, info: FieldValidationInfo) -> int | None: + """Validate weekly_day is only provided for WEEKLY schedule type.""" + schedule_type = info.data.get("schedule_type") + if v is not None and schedule_type != ScheduleType.WEEKLY: + raise ValueError( + "weekly_day should only be provided for WEEKLY schedule_type" + ) + if v is not None and not (0 <= v <= 6): + raise ValueError("weekly_day must be between 0 (Monday) and 6 (Sunday)") + return v + + @field_validator("weekly_time") + @classmethod + def validate_weekly_time_update(cls, v: time | None, info: FieldValidationInfo) -> time | None: + """Validate weekly_time is only provided for WEEKLY schedule type.""" + schedule_type = info.data.get("schedule_type") + if v is not None and schedule_type != ScheduleType.WEEKLY: + raise ValueError( + "weekly_time should only be provided for WEEKLY schedule_type" + ) + return v + + @field_validator("hourly_minute") + @classmethod + def validate_hourly_minute_update(cls, v: int | None, info: FieldValidationInfo) -> int | None: + """Validate hourly_minute is only provided for HOURLY schedule type.""" + schedule_type = info.data.get("schedule_type") + if v is not None and schedule_type != ScheduleType.HOURLY: + raise ValueError( + "hourly_minute should only be provided for HOURLY schedule_type" + ) + if v is not None and not (0 <= v <= 59): + raise ValueError("hourly_minute must be between 0 and 59") + return v class ConnectorScheduleRead(ConnectorScheduleBase, IDModel, TimestampModel): diff --git a/surfsense_backend/app/services/connector_scheduler_service.py b/surfsense_backend/app/services/connector_scheduler_service.py index b119c122e..27a619472 100644 --- a/surfsense_backend/app/services/connector_scheduler_service.py +++ b/surfsense_backend/app/services/connector_scheduler_service.py @@ -126,7 +126,7 @@ class ConnectorSchedulerService: async def _get_due_schedules(self, session: AsyncSession) -> List[ConnectorSchedule]: """Get all schedules that are due for execution.""" - now = datetime.now(datetime.utc) + now = datetime.now(timezone.utc) query = ( select(ConnectorSchedule) @@ -224,7 +224,7 @@ class ConnectorSchedulerService: "%Y-%m-%d" ) - end_date = datetime.now(datetime.utc).strftime("%Y-%m-%d") + end_date = datetime.now(timezone.utc).strftime("%Y-%m-%d") # Execute the indexer function documents_processed, error_message = await indexer_func(