From 18f28bbf41bf0ade3fbd2b59af00fcee57ccb0ba Mon Sep 17 00:00:00 2001 From: Badry Date: Sun, 30 Aug 2026 22:29:33 +0300 Subject: [PATCH 1/7] feat(scheduler): add logical parameter bindings --- ...2d9f801_add_schedule_parameter_bindings.py | 132 ++++++++ backend/app/models/job_run.py | 19 +- backend/app/models/schedule.py | 19 +- backend/app/repositories/job_run.py | 12 + backend/app/routers/schedules.py | 83 +++-- backend/app/schemas/endpoint.py | 56 +--- backend/app/schemas/schedule.py | 140 +++++++- backend/app/services/endpoint.py | 28 +- backend/app/services/schedule.py | 217 ++++++++++-- backend/app/services/schedule_bindings.py | 219 ++++++++++++ backend/app/services/scheduler.py | 203 ++++++++--- backend/app/sql/param_models.py | 16 +- backend/tests/test_endpoints.py | 62 ++-- backend/tests/test_schedule_bindings.py | 263 +++++++++++++++ backend/tests/test_scheduler_restore.py | 51 ++- backend/tests/test_schedules.py | 315 +++++++++++++++++- 16 files changed, 1582 insertions(+), 253 deletions(-) create mode 100644 backend/alembic/versions/e4a6c2d9f801_add_schedule_parameter_bindings.py create mode 100644 backend/app/services/schedule_bindings.py create mode 100644 backend/tests/test_schedule_bindings.py diff --git a/backend/alembic/versions/e4a6c2d9f801_add_schedule_parameter_bindings.py b/backend/alembic/versions/e4a6c2d9f801_add_schedule_parameter_bindings.py new file mode 100644 index 0000000..8376fee --- /dev/null +++ b/backend/alembic/versions/e4a6c2d9f801_add_schedule_parameter_bindings.py @@ -0,0 +1,132 @@ +"""Add schedule-owned parameter bindings and logical run audit context. + +Revision ID: e4a6c2d9f801 +Revises: c7e91a4f2d60 +Create Date: 2026-08-30 +""" + +from collections.abc import Sequence + +import sqlalchemy as sa +from alembic import op +from sqlalchemy.dialects import postgresql + +revision: str = "e4a6c2d9f801" +down_revision: str | None = "c7e91a4f2d60" +branch_labels: str | Sequence[str] | None = None +depends_on: str | Sequence[str] | None = None + + +def upgrade() -> None: + op.add_column( + "schedules", + sa.Column("timezone", sa.String(length=64), server_default="UTC", nullable=False), + ) + op.add_column( + "schedules", + sa.Column( + "parameter_bindings_json", + postgresql.JSONB(astext_type=sa.Text()), + server_default=sa.text("'{}'::jsonb"), + nullable=False, + ), + ) + op.add_column( + "schedules", + sa.Column( + "window_config_json", + postgresql.JSONB(astext_type=sa.Text()), + nullable=True, + ), + ) + + # Preserve the behavior of existing schedules by copying endpoint defaults + # into explicit schedule-owned bindings. Any legacy parameter without a + # resolvable default remains absent and will be reported when the schedule + # is edited, previewed, or executed rather than silently inventing a value. + op.execute( + """ + UPDATE schedules AS schedule + SET parameter_bindings_json = COALESCE(migrated.bindings, '{}'::jsonb) + FROM ( + SELECT + endpoint.id AS endpoint_id, + jsonb_object_agg( + parameter.name, + CASE + WHEN parameter.descriptor->>'default_expression' = 'today' + THEN jsonb_build_object('source', 'run_date') + WHEN parameter.descriptor->>'default_expression' = 'yesterday' + THEN jsonb_build_object( + 'source', 'relative_date', 'offset_days', -1 + ) + WHEN COALESCE( + (parameter.descriptor->>'default_is_null')::boolean, + false + ) + THEN jsonb_build_object('source', 'null') + WHEN parameter.descriptor ? 'default' + AND parameter.descriptor->'default' <> 'null'::jsonb + THEN jsonb_build_object( + 'source', 'literal', + 'value', parameter.descriptor->'default' + ) + ELSE NULL + END + ) FILTER ( + WHERE parameter.descriptor->>'default_expression' IN ('today', 'yesterday') + OR COALESCE( + (parameter.descriptor->>'default_is_null')::boolean, + false + ) + OR ( + parameter.descriptor ? 'default' + AND parameter.descriptor->'default' <> 'null'::jsonb + ) + ) AS bindings + FROM endpoints AS endpoint + CROSS JOIN LATERAL jsonb_each( + COALESCE(endpoint.param_schema_json, '{}'::jsonb) + ) AS parameter(name, descriptor) + GROUP BY endpoint.id + ) AS migrated + WHERE schedule.endpoint_id = migrated.endpoint_id + """ + ) + + op.add_column( + "job_runs", + sa.Column("scheduled_for", sa.DateTime(timezone=True), nullable=True), + ) + op.add_column("job_runs", sa.Column("logical_date", sa.Date(), nullable=True)) + op.add_column("job_runs", sa.Column("window_start", sa.Date(), nullable=True)) + op.add_column("job_runs", sa.Column("window_end", sa.Date(), nullable=True)) + op.add_column( + "job_runs", + sa.Column( + "resolved_params_json", + postgresql.JSONB(astext_type=sa.Text()), + nullable=True, + ), + ) + op.add_column("job_runs", sa.Column("trigger_source", sa.String(length=20), nullable=True)) + op.add_column("job_runs", sa.Column("binding_hash", sa.String(length=64), nullable=True)) + op.create_unique_constraint( + "uq_job_runs_schedule_scheduled_for", + "job_runs", + ["schedule_id", "scheduled_for"], + ) + + +def downgrade() -> None: + op.drop_constraint("uq_job_runs_schedule_scheduled_for", "job_runs", type_="unique") + op.drop_column("job_runs", "binding_hash") + op.drop_column("job_runs", "trigger_source") + op.drop_column("job_runs", "resolved_params_json") + op.drop_column("job_runs", "window_end") + op.drop_column("job_runs", "window_start") + op.drop_column("job_runs", "logical_date") + op.drop_column("job_runs", "scheduled_for") + op.drop_column("schedules", "window_config_json") + op.drop_column("schedules", "parameter_bindings_json") + op.drop_column("schedules", "timezone") diff --git a/backend/app/models/job_run.py b/backend/app/models/job_run.py index 27d9f4f..2b1d416 100644 --- a/backend/app/models/job_run.py +++ b/backend/app/models/job_run.py @@ -1,11 +1,12 @@ """JobRun model — immutable execution audit record for each scheduler invocation.""" import uuid -from datetime import datetime +from datetime import date, datetime from enum import StrEnum -from sqlalchemy import DateTime, ForeignKey, Integer, String, func +from sqlalchemy import Date, DateTime, ForeignKey, Integer, String, UniqueConstraint, func from sqlalchemy import Enum as SAEnum +from sqlalchemy.dialects.postgresql import JSONB from sqlalchemy.orm import Mapped, mapped_column from app.models.base import Base, UUIDPrimaryKeyMixin @@ -27,6 +28,13 @@ class JobRun(UUIDPrimaryKeyMixin, Base): """ __tablename__ = "job_runs" + __table_args__ = ( + UniqueConstraint( + "schedule_id", + "scheduled_for", + name="uq_job_runs_schedule_scheduled_for", + ), + ) id: Mapped[uuid.UUID] = mapped_column(primary_key=True, default=uuid.uuid4) @@ -47,6 +55,13 @@ class JobRun(UUIDPrimaryKeyMixin, Base): row_count: Mapped[int | None] = mapped_column(Integer, nullable=True) # Truncated error detail; full stack trace goes to structured log. error_detail: Mapped[str | None] = mapped_column(String(5000), nullable=True) + scheduled_for: Mapped[datetime | None] = mapped_column(DateTime(timezone=True), nullable=True) + logical_date: Mapped[date | None] = mapped_column(Date, nullable=True) + window_start: Mapped[date | None] = mapped_column(Date, nullable=True) + window_end: Mapped[date | None] = mapped_column(Date, nullable=True) + resolved_params_json: Mapped[dict[str, object] | None] = mapped_column(JSONB, nullable=True) + trigger_source: Mapped[str | None] = mapped_column(String(20), nullable=True) + binding_hash: Mapped[str | None] = mapped_column(String(64), nullable=True) created_at: Mapped[datetime] = mapped_column( DateTime(timezone=True), server_default=func.now(), nullable=False diff --git a/backend/app/models/schedule.py b/backend/app/models/schedule.py index 9385122..8fb49cf 100644 --- a/backend/app/models/schedule.py +++ b/backend/app/models/schedule.py @@ -6,6 +6,7 @@ from sqlalchemy import Boolean, DateTime, ForeignKey, Integer, String, text from sqlalchemy import Enum as SAEnum +from sqlalchemy.dialects.postgresql import JSONB from sqlalchemy.orm import Mapped, mapped_column from app.models.base import Base, TimestampMixin, UUIDPrimaryKeyMixin @@ -38,13 +39,19 @@ class Schedule(UUIDPrimaryKeyMixin, TimestampMixin, Base): cron_expression: Mapped[str | None] = mapped_column(String(100), nullable=True) # Used when schedule_type == 'interval'. Positive integer seconds. interval_seconds: Mapped[int | None] = mapped_column(Integer, nullable=True) + timezone: Mapped[str] = mapped_column( + String(64), nullable=False, default="UTC", server_default="UTC" + ) + parameter_bindings_json: Mapped[dict[str, object]] = mapped_column( + JSONB, + nullable=False, + default=dict, + server_default=text("'{}'::jsonb"), + ) + window_config_json: Mapped[dict[str, object] | None] = mapped_column(JSONB, nullable=True) is_active: Mapped[bool] = mapped_column( Boolean, default=True, nullable=False, server_default=text("true") ) - last_run_at: Mapped[datetime | None] = mapped_column( - DateTime(timezone=True), nullable=True - ) - next_run_at: Mapped[datetime | None] = mapped_column( - DateTime(timezone=True), nullable=True - ) + last_run_at: Mapped[datetime | None] = mapped_column(DateTime(timezone=True), nullable=True) + next_run_at: Mapped[datetime | None] = mapped_column(DateTime(timezone=True), nullable=True) diff --git a/backend/app/repositories/job_run.py b/backend/app/repositories/job_run.py index 50f61a7..4123c8b 100644 --- a/backend/app/repositories/job_run.py +++ b/backend/app/repositories/job_run.py @@ -2,6 +2,7 @@ import uuid from collections.abc import Sequence +from datetime import datetime from sqlalchemy import select @@ -20,6 +21,17 @@ class JobRunRepository(BaseCrudRepository[JobRun]): model = JobRun + async def get_by_schedule_and_scheduled_for( + self, schedule_id: uuid.UUID, scheduled_for: datetime + ) -> JobRun | None: + result = await self._db.execute( + select(JobRun).where( + JobRun.schedule_id == schedule_id, + JobRun.scheduled_for == scheduled_for, + ) + ) + return result.scalar_one_or_none() + async def get_all( self, *, diff --git a/backend/app/routers/schedules.py b/backend/app/routers/schedules.py index 90ce0dd..cd9d570 100644 --- a/backend/app/routers/schedules.py +++ b/backend/app/routers/schedules.py @@ -5,10 +5,11 @@ Routes: GET / — List all schedules POST / — Create a schedule + POST /preview — Preview logical runs and resolved bindings GET /{id} — Get a single schedule PUT /{id} — Update a schedule DELETE /{id} — Delete a schedule - POST /{id}/run — Run schedule now + POST /{id}/run — Run now, optionally for a logical date POST /{id}/pause — Pause a schedule POST /{id}/resume — Resume a schedule GET /jobs/ — List job runs @@ -29,16 +30,19 @@ from app.repositories.job_run import JobRunRepository from app.repositories.schedule import ScheduleRepository from app.repositories.snapshot import SnapshotRepository -from app.schemas.endpoint import SnapshotConfigurationError from app.schemas.schedule import ( JobRunResponse, ScheduleCreate, + SchedulePreviewRequest, + SchedulePreviewResponse, ScheduleResponse, + ScheduleRunRequest, ScheduleUpdate, SnapshotDetailResponse, SnapshotResponse, ) from app.services.schedule import ScheduleService +from app.services.schedule_bindings import ScheduleBindingError from app.services.scheduler import remove_schedule_job log = structlog.get_logger() @@ -87,18 +91,35 @@ async def create_schedule( ) -> ScheduleResponse: try: result = await svc.create_schedule(payload) - except SnapshotConfigurationError as exc: + except ScheduleBindingError as exc: raise HTTPException( status_code=status.HTTP_422_UNPROCESSABLE_ENTITY, detail=str(exc) ) from exc except ValueError as exc: - raise HTTPException( - status_code=status.HTTP_409_CONFLICT, detail=str(exc) - ) from exc + raise HTTPException(status_code=status.HTTP_409_CONFLICT, detail=str(exc)) from exc await db.commit() return result +@router.post( + "/preview", + response_model=SchedulePreviewResponse, + summary="Preview resolved schedule runs", +) +async def preview_schedule( + payload: SchedulePreviewRequest, + svc: ScheduleService = Depends(_service), +) -> SchedulePreviewResponse: + try: + return await svc.preview_schedule(payload) + except ScheduleBindingError as exc: + raise HTTPException( + status_code=status.HTTP_422_UNPROCESSABLE_ENTITY, detail=str(exc) + ) from exc + except ValueError as exc: + raise HTTPException(status_code=status.HTTP_404_NOT_FOUND, detail=str(exc)) from exc + + @router.get( "/{schedule_id}", response_model=ScheduleResponse, @@ -110,9 +131,7 @@ async def get_schedule( ) -> ScheduleResponse: result = await svc.get_schedule(schedule_id) if result is None: - raise HTTPException( - status_code=status.HTTP_404_NOT_FOUND, detail="Schedule not found." - ) + raise HTTPException(status_code=status.HTTP_404_NOT_FOUND, detail="Schedule not found.") return result @@ -129,14 +148,14 @@ async def update_schedule( ) -> ScheduleResponse: try: result = await svc.update_schedule(schedule_id, payload) - except ValueError as exc: + except ScheduleBindingError as exc: raise HTTPException( - status_code=status.HTTP_409_CONFLICT, detail=str(exc) + status_code=status.HTTP_422_UNPROCESSABLE_ENTITY, detail=str(exc) ) from exc + except ValueError as exc: + raise HTTPException(status_code=status.HTTP_409_CONFLICT, detail=str(exc)) from exc if result is None: - raise HTTPException( - status_code=status.HTTP_404_NOT_FOUND, detail="Schedule not found." - ) + raise HTTPException(status_code=status.HTTP_404_NOT_FOUND, detail="Schedule not found.") await db.commit() return result @@ -153,9 +172,7 @@ async def delete_schedule( ) -> None: deleted = await svc.delete_schedule(schedule_id) if not deleted: - raise HTTPException( - status_code=status.HTTP_404_NOT_FOUND, detail="Schedule not found." - ) + raise HTTPException(status_code=status.HTTP_404_NOT_FOUND, detail="Schedule not found.") await db.commit() remove_schedule_job(schedule_id) @@ -171,14 +188,17 @@ async def delete_schedule( ) async def run_now( schedule_id: uuid.UUID, + payload: ScheduleRunRequest | None = None, svc: ScheduleService = Depends(_service), ) -> dict[str, str]: try: - await svc.run_now(schedule_id) - except ValueError as exc: + await svc.run_now(schedule_id, payload) + except ScheduleBindingError as exc: raise HTTPException( - status_code=status.HTTP_404_NOT_FOUND, detail=str(exc) + status_code=status.HTTP_422_UNPROCESSABLE_ENTITY, detail=str(exc) ) from exc + except ValueError as exc: + raise HTTPException(status_code=status.HTTP_404_NOT_FOUND, detail=str(exc)) from exc return {"status": "executed"} @@ -194,9 +214,7 @@ async def pause_schedule( ) -> ScheduleResponse: result = await svc.pause(schedule_id) if result is None: - raise HTTPException( - status_code=status.HTTP_404_NOT_FOUND, detail="Schedule not found." - ) + raise HTTPException(status_code=status.HTTP_404_NOT_FOUND, detail="Schedule not found.") await db.commit() return result @@ -213,9 +231,7 @@ async def resume_schedule( ) -> ScheduleResponse: result = await svc.resume(schedule_id) if result is None: - raise HTTPException( - status_code=status.HTTP_404_NOT_FOUND, detail="Schedule not found." - ) + raise HTTPException(status_code=status.HTTP_404_NOT_FOUND, detail="Schedule not found.") await db.commit() return result @@ -255,9 +271,7 @@ async def get_job_run( repo = JobRunRepository(db) obj = await repo.get_by_id(job_run_id) if obj is None: - raise HTTPException( - status_code=status.HTTP_404_NOT_FOUND, detail="Job run not found." - ) + raise HTTPException(status_code=status.HTTP_404_NOT_FOUND, detail="Job run not found.") return JobRunResponse( id=obj.id, schedule_id=obj.schedule_id, @@ -267,6 +281,13 @@ async def get_job_run( status=obj.status, row_count=obj.row_count, error_detail=obj.error_detail, + scheduled_for=obj.scheduled_for, + logical_date=obj.logical_date, + window_start=obj.window_start, + window_end=obj.window_end, + resolved_parameters=obj.resolved_params_json, + trigger_source=obj.trigger_source, + binding_hash=obj.binding_hash, created_at=obj.created_at, ) @@ -298,7 +319,5 @@ async def get_snapshot( ) -> SnapshotDetailResponse: result = await svc.get_snapshot(snapshot_id) if result is None: - raise HTTPException( - status_code=status.HTTP_404_NOT_FOUND, detail="Snapshot not found." - ) + raise HTTPException(status_code=status.HTTP_404_NOT_FOUND, detail="Snapshot not found.") return result diff --git a/backend/app/schemas/endpoint.py b/backend/app/schemas/endpoint.py index 728fda2..0df0cb6 100644 --- a/backend/app/schemas/endpoint.py +++ b/backend/app/schemas/endpoint.py @@ -9,7 +9,6 @@ import re import uuid -from collections.abc import Mapping from datetime import datetime from typing import Literal, Self @@ -110,9 +109,7 @@ def optional_must_have_default(self) -> Self: ) ) if configured_defaults > 1: - raise ValueError( - "Declare only one of default, default_is_null, or default_expression." - ) + raise ValueError("Declare only one of default, default_is_null, or default_expression.") if self.default_is_null and self.required: raise ValueError("A NULL default is supported only for optional parameters.") if self.default_expression is not None and self.type != "date": @@ -130,44 +127,6 @@ def optional_must_have_default(self) -> Self: return self -def missing_snapshot_defaults( - param_schema: Mapping[str, object], -) -> list[str]: - """Return snapshot bind names that cannot be resolved without a request.""" - missing: list[str] = [] - for name, raw_descriptor in param_schema.items(): - if isinstance(raw_descriptor, ParamDescriptor): - descriptor = raw_descriptor - elif isinstance(raw_descriptor, dict): - descriptor = ParamDescriptor.model_validate(raw_descriptor) - else: - missing.append(name) - continue - - if ( - descriptor.default is None - and not descriptor.default_is_null - and descriptor.default_expression is None - ): - missing.append(name) - return sorted(missing) - - -def require_snapshot_defaults( - data_strategy: DataStrategy, - param_schema: Mapping[str, object], -) -> None: - """Reject snapshot endpoints whose binds require caller-supplied values.""" - if data_strategy != DataStrategy.snapshot: - return - missing = missing_snapshot_defaults(param_schema) - if missing: - names = ", ".join(f":{name}" for name in missing) - raise SnapshotConfigurationError( - f"Snapshot endpoints require a default for every parameter. Missing: {names}." - ) - - class EndpointCreate(BaseModel): """Payload for POST /api/v1/admin/endpoints.""" @@ -235,14 +194,9 @@ def bind_params_match_schema(self) -> Self: undeclared = sql_params - schema_params unused = schema_params - sql_params if undeclared: - raise ValueError( - f"SQL references params not declared in schema: {sorted(undeclared)}" - ) + raise ValueError(f"SQL references params not declared in schema: {sorted(undeclared)}") if unused: - raise ValueError( - f"Schema declares params not referenced in SQL: {sorted(unused)}" - ) - require_snapshot_defaults(self.data_strategy, self.param_schema) + raise ValueError(f"Schema declares params not referenced in SQL: {sorted(unused)}") return self @@ -307,9 +261,7 @@ def bind_params_match_schema(self) -> Self: f"SQL references params not declared in schema: {sorted(undeclared)}" ) if unused: - raise ValueError( - f"Schema declares params not referenced in SQL: {sorted(unused)}" - ) + raise ValueError(f"Schema declares params not referenced in SQL: {sorted(unused)}") return self @model_validator(mode="after") diff --git a/backend/app/schemas/schedule.py b/backend/app/schemas/schedule.py index c0f33bd..3410070 100644 --- a/backend/app/schemas/schedule.py +++ b/backend/app/schemas/schedule.py @@ -1,13 +1,76 @@ """Pydantic schemas for schedule, job run, and snapshot resources.""" import uuid -from datetime import datetime +from datetime import date, datetime +from typing import Literal, Self +from zoneinfo import ZoneInfo, ZoneInfoNotFoundError -from pydantic import BaseModel, Field, field_validator +from pydantic import BaseModel, Field, field_validator, model_validator # ── Schedule schemas ───────────────────────────────────────────────────────── +BindingSource = Literal[ + "literal", + "null", + "run_date", + "relative_date", + "window_start", + "window_end", +] +WindowPreset = Literal[ + "previous_day", + "last_n_complete_days", + "week_to_date", + "previous_week", + "month_to_date", + "previous_month", +] + + +class ScheduleParameterBinding(BaseModel): + """Declarative source for one SQL bind during scheduled execution.""" + + source: BindingSource + value: str | int | float | bool | None = None + offset_days: int | None = Field(None, ge=-36500, le=36500) + + @model_validator(mode="after") + def validate_source_options(self) -> Self: + if self.source == "literal" and self.value is None: + raise ValueError("Literal bindings require a value; use source='null' for SQL NULL.") + if self.source == "relative_date" and self.offset_days is None: + raise ValueError("Relative-date bindings require offset_days.") + if self.source != "literal" and self.value is not None: + raise ValueError("A value is supported only for literal bindings.") + if self.source != "relative_date" and self.offset_days is not None: + raise ValueError("offset_days is supported only for relative-date bindings.") + return self + + +class ScheduleWindow(BaseModel): + """Reusable inclusive date window resolved from a schedule's logical date.""" + + preset: WindowPreset + days: int | None = Field(None, ge=1, le=3660) + + @model_validator(mode="after") + def validate_days(self) -> Self: + if self.preset == "last_n_complete_days" and self.days is None: + raise ValueError("last_n_complete_days requires days.") + if self.preset != "last_n_complete_days" and self.days is not None: + raise ValueError("days is supported only for last_n_complete_days.") + return self + + +def _validate_timezone(value: str) -> str: + try: + ZoneInfo(value) + except (ZoneInfoNotFoundError, ValueError) as exc: + raise ValueError(f"Unknown IANA timezone: {value}") from exc + return value + + def _validate_cron_expression(value: str | None) -> str | None: if value is None: return None @@ -30,8 +93,16 @@ class ScheduleCreate(BaseModel): schedule_type: str = Field(..., pattern=r"^(cron|interval)$") cron_expression: str | None = None interval_seconds: int | None = Field(None, ge=10) + timezone: str = "UTC" + parameter_bindings: dict[str, ScheduleParameterBinding] = Field(default_factory=dict) + window: ScheduleWindow | None = None is_active: bool = True + @field_validator("timezone") + @classmethod + def validate_timezone(cls, value: str) -> str: + return _validate_timezone(value) + @field_validator("cron_expression") @classmethod def validate_cron(cls, v: str | None) -> str | None: @@ -56,12 +127,20 @@ class ScheduleUpdate(BaseModel): schedule_type: str | None = Field(None, pattern=r"^(cron|interval)$") cron_expression: str | None = None interval_seconds: int | None = Field(None, ge=10) + timezone: str | None = None + parameter_bindings: dict[str, ScheduleParameterBinding] | None = None + window: ScheduleWindow | None = None is_active: bool | None = None + @field_validator("timezone") + @classmethod + def validate_timezone(cls, value: str | None) -> str | None: + return _validate_timezone(value) if value is not None else None + @field_validator("cron_expression") @classmethod - def validate_cron(cls, v: str | None) -> str | None: - return _validate_cron_expression(v) + def validate_cron(cls, value: str | None) -> str | None: + return _validate_cron_expression(value) class ScheduleResponse(BaseModel): @@ -70,6 +149,9 @@ class ScheduleResponse(BaseModel): schedule_type: str cron_expression: str | None interval_seconds: int | None + timezone: str + parameter_bindings: dict[str, ScheduleParameterBinding] + window: ScheduleWindow | None is_active: bool last_run_at: datetime | None next_run_at: datetime | None @@ -79,6 +161,49 @@ class ScheduleResponse(BaseModel): model_config = {"from_attributes": True} +class SchedulePreviewRequest(BaseModel): + endpoint_id: uuid.UUID + schedule_type: str = Field(..., pattern=r"^(cron|interval)$") + cron_expression: str | None = None + interval_seconds: int | None = Field(None, ge=10) + timezone: str = "UTC" + parameter_bindings: dict[str, ScheduleParameterBinding] = Field(default_factory=dict) + window: ScheduleWindow | None = None + count: int = Field(3, ge=1, le=10) + + @field_validator("timezone") + @classmethod + def validate_timezone(cls, value: str) -> str: + return _validate_timezone(value) + + @field_validator("cron_expression") + @classmethod + def validate_cron(cls, value: str | None) -> str | None: + return _validate_cron_expression(value) + + def model_post_init(self, __context: object) -> None: + if self.schedule_type == "cron" and not self.cron_expression: + raise ValueError("cron_expression is required when schedule_type is 'cron'.") + if self.schedule_type == "interval" and not self.interval_seconds: + raise ValueError("interval_seconds is required when schedule_type is 'interval'.") + + +class ScheduleRunPreview(BaseModel): + scheduled_for: datetime + logical_date: date + window_start: date | None + window_end: date | None + resolved_parameters: dict[str, object] + + +class SchedulePreviewResponse(BaseModel): + runs: list[ScheduleRunPreview] + + +class ScheduleRunRequest(BaseModel): + logical_date: date | None = None + + # ── JobRun schemas ─────────────────────────────────────────────────────────── @@ -91,6 +216,13 @@ class JobRunResponse(BaseModel): status: str row_count: int | None error_detail: str | None + scheduled_for: datetime | None = None + logical_date: date | None = None + window_start: date | None = None + window_end: date | None = None + resolved_parameters: dict[str, object] | None = None + trigger_source: str | None = None + binding_hash: str | None = None created_at: datetime model_config = {"from_attributes": True} diff --git a/backend/app/services/endpoint.py b/backend/app/services/endpoint.py index 60279c8..a9e37db 100644 --- a/backend/app/services/endpoint.py +++ b/backend/app/services/endpoint.py @@ -28,7 +28,6 @@ SqlPreviewRequest, SqlPreviewResponse, extract_bind_params, - require_snapshot_defaults, ) from app.sql.executor import SqlExecutionError, execute_query @@ -80,9 +79,7 @@ def __init__( self._repo = repo self._conn_repo = conn_repo - async def list_endpoints( - self, *, active_only: bool = False - ) -> Sequence[EndpointResponse]: + async def list_endpoints(self, *, active_only: bool = False) -> Sequence[EndpointResponse]: rows = await self._repo.get_all(active_only=active_only) return [_to_response(r) for r in rows] @@ -161,8 +158,7 @@ async def update_endpoint( "deprecation_note", } changes: dict[str, object] = { - field: getattr(payload, field) - for field in payload.model_fields_set & _updatable + field: getattr(payload, field) for field in payload.model_fields_set & _updatable } # Always preserve an authentication path in the merged state: either a @@ -182,14 +178,6 @@ async def update_endpoint( if "data_strategy" in payload.model_fields_set and payload.data_strategy is None: raise SnapshotConfigurationError("data_strategy cannot be null.") - effective_strategy = payload.data_strategy or obj.data_strategy - effective_param_schema: dict[str, ParamDescriptor] | dict[str, object] - if "param_schema" in payload.model_fields_set and payload.param_schema is not None: - effective_param_schema = payload.param_schema - else: - effective_param_schema = obj.param_schema_json or {} - require_snapshot_defaults(effective_strategy, effective_param_schema) - # Uniqueness check on name change if payload.name is not None and payload.name != obj.name: conflict = await self._repo.get_by_name(payload.name) @@ -200,9 +188,7 @@ async def update_endpoint( if payload.path is not None and payload.path != obj.path: conflict = await self._repo.get_by_path(payload.path) if conflict: - raise ValueError( - f"An endpoint with path '{payload.path}' already exists." - ) + raise ValueError(f"An endpoint with path '{payload.path}' already exists.") # Handle param_schema serialization if "param_schema" in payload.model_fields_set and payload.param_schema is not None: @@ -227,9 +213,7 @@ async def update_endpoint( ) return _to_response(obj) - async def delete_endpoint( - self, endpoint_id: uuid.UUID, *, actor: str = "system" - ) -> bool: + async def delete_endpoint(self, endpoint_id: uuid.UUID, *, actor: str = "system") -> bool: obj = await self._repo.get_by_id(endpoint_id) if obj is None: return False @@ -245,9 +229,7 @@ async def delete_endpoint( ) return True - async def preview_sql( - self, payload: SqlPreviewRequest - ) -> SqlPreviewResponse: + async def preview_sql(self, payload: SqlPreviewRequest) -> SqlPreviewResponse: """Execute SQL in preview mode and return sample results.""" if not self._conn_repo: raise ValueError("Connection repository not available for preview.") diff --git a/backend/app/services/schedule.py b/backend/app/services/schedule.py index 43d273d..088cd7f 100644 --- a/backend/app/services/schedule.py +++ b/backend/app/services/schedule.py @@ -10,23 +10,35 @@ import uuid from collections.abc import Sequence +from datetime import UTC, datetime, time +from zoneinfo import ZoneInfo import structlog +from app.models.endpoint import DataStrategy from app.models.schedule import Schedule from app.repositories.endpoint import EndpointRepository from app.repositories.job_run import JobRunRepository from app.repositories.schedule import ScheduleRepository from app.repositories.snapshot import SnapshotRepository -from app.schemas.endpoint import require_snapshot_defaults from app.schemas.schedule import ( JobRunResponse, ScheduleCreate, + ScheduleParameterBinding, + SchedulePreviewRequest, + SchedulePreviewResponse, ScheduleResponse, + ScheduleRunRequest, ScheduleUpdate, + ScheduleWindow, SnapshotDetailResponse, SnapshotResponse, ) +from app.services.schedule_bindings import ( + ScheduleBindingError, + preview_schedule_runs, + resolve_schedule_parameters, +) from app.services.scheduler import ( add_schedule_job, execute_scheduled_job, @@ -39,12 +51,22 @@ def _to_response(obj: Schedule) -> ScheduleResponse: + parameter_bindings = { + name: ScheduleParameterBinding.model_validate(binding) + for name, binding in (obj.parameter_bindings_json or {}).items() + } + window = ( + ScheduleWindow.model_validate(obj.window_config_json) if obj.window_config_json else None + ) return ScheduleResponse( id=obj.id, endpoint_id=obj.endpoint_id, schedule_type=obj.schedule_type, cron_expression=obj.cron_expression, interval_seconds=obj.interval_seconds, + timezone=obj.timezone, + parameter_bindings=parameter_bindings, + window=window, is_active=obj.is_active, last_run_at=obj.last_run_at, next_run_at=obj.next_run_at, @@ -53,6 +75,47 @@ def _to_response(obj: Schedule) -> ScheduleResponse: ) +def _bindings_to_json( + bindings: dict[str, ScheduleParameterBinding], +) -> dict[str, object]: + return {name: binding.model_dump(mode="json") for name, binding in bindings.items()} + + +def _window_to_json(window: ScheduleWindow | None) -> dict[str, object] | None: + return window.model_dump(mode="json") if window is not None else None + + +def _validate_snapshot_bindings( + *, + endpoint_strategy: DataStrategy, + param_schema: dict[str, object], + parameter_bindings: dict[str, ScheduleParameterBinding] | dict[str, object], + timezone: str, + window: ScheduleWindow | dict[str, object] | None, +) -> None: + if endpoint_strategy != DataStrategy.snapshot: + raise ScheduleBindingError("Schedules are supported only for snapshot endpoints.") + resolve_schedule_parameters( + param_schema=param_schema, + parameter_bindings=parameter_bindings, + timezone_name=timezone, + scheduled_for=datetime.now(UTC), + window=window, + ) + + +def _validate_schedule_timing( + *, + schedule_type: str, + cron_expression: str | None, + interval_seconds: int | None, +) -> None: + if schedule_type == "cron" and not cron_expression: + raise ValueError("cron_expression is required when schedule_type is 'cron'.") + if schedule_type == "interval" and not interval_seconds: + raise ValueError("interval_seconds is required when schedule_type is 'interval'.") + + class ScheduleService: """Business logic layer for schedule management.""" @@ -68,9 +131,7 @@ def __init__( self._job_repo = job_repo self._snap_repo = snap_repo - async def list_schedules( - self, *, active_only: bool = False - ) -> Sequence[ScheduleResponse]: + async def list_schedules(self, *, active_only: bool = False) -> Sequence[ScheduleResponse]: rows = await self._repo.get_all(active_only=active_only) return [_to_response(r) for r in rows] @@ -86,33 +147,43 @@ async def create_schedule( ep = await self._ep_repo.get_by_id(payload.endpoint_id) if ep is None: raise ValueError(f"Endpoint '{payload.endpoint_id}' not found.") - require_snapshot_defaults(ep.data_strategy, ep.param_schema_json or {}) + _validate_snapshot_bindings( + endpoint_strategy=ep.data_strategy, + param_schema=ep.param_schema_json or {}, + parameter_bindings=payload.parameter_bindings, + timezone=payload.timezone, + window=payload.window, + ) # Check uniqueness — one schedule per endpoint existing = await self._repo.get_by_endpoint_id(payload.endpoint_id) if existing: - raise ValueError( - f"A schedule already exists for endpoint '{payload.endpoint_id}'." - ) + raise ValueError(f"A schedule already exists for endpoint '{payload.endpoint_id}'.") obj = Schedule( endpoint_id=payload.endpoint_id, schedule_type=payload.schedule_type, cron_expression=payload.cron_expression, interval_seconds=payload.interval_seconds, + timezone=payload.timezone, + parameter_bindings_json=_bindings_to_json(payload.parameter_bindings), + window_config_json=_window_to_json(payload.window), is_active=payload.is_active, ) obj = await self._repo.create(obj) # Register with APScheduler if active if obj.is_active: - add_schedule_job( + next_run_at = add_schedule_job( schedule_id=obj.id, endpoint_id=obj.endpoint_id, schedule_type=obj.schedule_type, cron_expression=obj.cron_expression, interval_seconds=obj.interval_seconds, + timezone_name=obj.timezone, ) + if next_run_at is not None: + obj = await self._repo.update(obj, {"next_run_at": next_run_at}) log.info( "schedule_created", @@ -138,24 +209,79 @@ async def update_schedule( "schedule_type", "cron_expression", "interval_seconds", + "timezone", + "parameter_bindings", + "window", "is_active", } - changes: dict[str, object] = { - field: getattr(payload, field) - for field in payload.model_fields_set & _updatable - } + for field in {"schedule_type", "timezone", "parameter_bindings", "is_active"}: + if field in payload.model_fields_set and getattr(payload, field) is None: + raise ValueError(f"{field} cannot be null.") + + effective_schedule_type = payload.schedule_type or obj.schedule_type + effective_cron_expression = ( + payload.cron_expression + if "cron_expression" in payload.model_fields_set + else obj.cron_expression + ) + effective_interval_seconds = ( + payload.interval_seconds + if "interval_seconds" in payload.model_fields_set + else obj.interval_seconds + ) + _validate_schedule_timing( + schedule_type=effective_schedule_type, + cron_expression=effective_cron_expression, + interval_seconds=effective_interval_seconds, + ) + + changes: dict[str, object] = {} + for field in payload.model_fields_set & _updatable: + if field == "parameter_bindings": + changes["parameter_bindings_json"] = _bindings_to_json( + payload.parameter_bindings or {} + ) + elif field == "window": + changes["window_config_json"] = _window_to_json(payload.window) + else: + changes[field] = getattr(payload, field) + + if self._ep_repo: + endpoint = await self._ep_repo.get_by_id(obj.endpoint_id) + if endpoint is None: + raise ValueError(f"Endpoint '{obj.endpoint_id}' not found.") + effective_bindings: dict[str, ScheduleParameterBinding] | dict[str, object] + if "parameter_bindings" in payload.model_fields_set: + effective_bindings = payload.parameter_bindings or {} + else: + effective_bindings = obj.parameter_bindings_json or {} + effective_window: ScheduleWindow | dict[str, object] | None + if "window" in payload.model_fields_set: + effective_window = payload.window + else: + effective_window = obj.window_config_json + _validate_snapshot_bindings( + endpoint_strategy=endpoint.data_strategy, + param_schema=endpoint.param_schema_json or {}, + parameter_bindings=effective_bindings, + timezone=payload.timezone or obj.timezone, + window=effective_window, + ) obj = await self._repo.update(obj, changes) # Sync with APScheduler if obj.is_active: - add_schedule_job( + next_run_at = add_schedule_job( schedule_id=obj.id, endpoint_id=obj.endpoint_id, schedule_type=obj.schedule_type, cron_expression=obj.cron_expression, interval_seconds=obj.interval_seconds, + timezone_name=obj.timezone, ) + if next_run_at is not None: + obj = await self._repo.update(obj, {"next_run_at": next_run_at}) else: remove_schedule_job(obj.id) @@ -167,9 +293,7 @@ async def update_schedule( ) return _to_response(obj) - async def delete_schedule( - self, schedule_id: uuid.UUID, *, actor: str = "system" - ) -> bool: + async def delete_schedule(self, schedule_id: uuid.UUID, *, actor: str = "system") -> bool: obj = await self._repo.get_by_id(schedule_id) if obj is None: return False @@ -184,14 +308,53 @@ async def delete_schedule( ) return True - async def run_now(self, schedule_id: uuid.UUID) -> None: + async def run_now( + self, schedule_id: uuid.UUID, payload: ScheduleRunRequest | None = None + ) -> None: """Trigger immediate execution of a schedule's job.""" obj = await self._repo.get_by_id(schedule_id) if obj is None: raise ValueError("Schedule not found.") log.info("schedule_run_now", schedule_id=str(schedule_id)) - await execute_scheduled_job(str(schedule_id), str(obj.endpoint_id)) + scheduled_for = None + if payload is not None and payload.logical_date is not None: + scheduled_for = datetime.combine( + payload.logical_date, + time.min, + tzinfo=ZoneInfo(obj.timezone), + ) + await execute_scheduled_job( + str(schedule_id), + str(obj.endpoint_id), + scheduled_for=scheduled_for, + trigger_source="manual", + ) + + async def preview_schedule(self, payload: SchedulePreviewRequest) -> SchedulePreviewResponse: + if not self._ep_repo: + raise ValueError("Endpoint repository not available.") + endpoint = await self._ep_repo.get_by_id(payload.endpoint_id) + if endpoint is None: + raise ValueError(f"Endpoint '{payload.endpoint_id}' not found.") + _validate_snapshot_bindings( + endpoint_strategy=endpoint.data_strategy, + param_schema=endpoint.param_schema_json or {}, + parameter_bindings=payload.parameter_bindings, + timezone=payload.timezone, + window=payload.window, + ) + runs = preview_schedule_runs( + schedule_type=payload.schedule_type, + cron_expression=payload.cron_expression, + interval_seconds=payload.interval_seconds, + timezone_name=payload.timezone, + param_schema=endpoint.param_schema_json or {}, + parameter_bindings=payload.parameter_bindings, + window=payload.window, + count=payload.count, + ) + return SchedulePreviewResponse(runs=runs) async def pause(self, schedule_id: uuid.UUID) -> ScheduleResponse | None: """Pause a schedule.""" @@ -215,13 +378,16 @@ async def resume(self, schedule_id: uuid.UUID) -> ScheduleResponse | None: resume_schedule_job(obj.id) # Re-register the job to ensure it's in APScheduler - add_schedule_job( + next_run_at = add_schedule_job( schedule_id=obj.id, endpoint_id=obj.endpoint_id, schedule_type=obj.schedule_type, cron_expression=obj.cron_expression, interval_seconds=obj.interval_seconds, + timezone_name=obj.timezone, ) + if next_run_at is not None: + obj = await self._repo.update(obj, {"next_run_at": next_run_at}) log.info("schedule_resumed", schedule_id=str(schedule_id)) return _to_response(obj) @@ -252,6 +418,13 @@ async def list_job_runs( status=r.status, row_count=r.row_count, error_detail=r.error_detail, + scheduled_for=r.scheduled_for, + logical_date=r.logical_date, + window_start=r.window_start, + window_end=r.window_end, + resolved_parameters=r.resolved_params_json, + trigger_source=r.trigger_source, + binding_hash=r.binding_hash, created_at=r.created_at, ) for r in rows @@ -276,9 +449,7 @@ async def list_snapshots( for r in rows ] - async def get_snapshot( - self, snapshot_id: uuid.UUID - ) -> SnapshotDetailResponse | None: + async def get_snapshot(self, snapshot_id: uuid.UUID) -> SnapshotDetailResponse | None: if not self._snap_repo: raise ValueError("Snapshot repository not available.") snap = await self._snap_repo.get_by_id(snapshot_id) diff --git a/backend/app/services/schedule_bindings.py b/backend/app/services/schedule_bindings.py new file mode 100644 index 0000000..1a7f429 --- /dev/null +++ b/backend/app/services/schedule_bindings.py @@ -0,0 +1,219 @@ +"""Resolve declarative schedule bindings against a deterministic logical date.""" + +from dataclasses import dataclass +from datetime import UTC, date, datetime, timedelta +from typing import Any +from zoneinfo import ZoneInfo + +from pydantic import ValidationError + +from app.schemas.schedule import ( + ScheduleParameterBinding, + ScheduleRunPreview, + ScheduleWindow, +) +from app.sql.param_models import build_param_model + + +class ScheduleBindingError(ValueError): + """Raised when a schedule cannot resolve every endpoint SQL bind safely.""" + + +@dataclass(frozen=True) +class ResolvedScheduleContext: + parameters: dict[str, Any] + logical_date: date + window_start: date | None + window_end: date | None + + +def _normalise_bindings( + bindings: dict[str, ScheduleParameterBinding] | dict[str, object], +) -> dict[str, ScheduleParameterBinding]: + return { + name: ( + binding + if isinstance(binding, ScheduleParameterBinding) + else ScheduleParameterBinding.model_validate(binding) + ) + for name, binding in bindings.items() + } + + +def _resolve_window(logical_date: date, window: ScheduleWindow) -> tuple[date, date]: + if window.preset == "previous_day": + day = logical_date - timedelta(days=1) + return day, day + if window.preset == "last_n_complete_days": + if window.days is None: # Defensive; schema validation normally catches this. + raise ScheduleBindingError("last_n_complete_days requires days.") + end = logical_date - timedelta(days=1) + return end - timedelta(days=window.days - 1), end + if window.preset == "week_to_date": + return logical_date - timedelta(days=logical_date.weekday()), logical_date + if window.preset == "previous_week": + current_week_start = logical_date - timedelta(days=logical_date.weekday()) + end = current_week_start - timedelta(days=1) + return end - timedelta(days=6), end + if window.preset == "month_to_date": + return logical_date.replace(day=1), logical_date + if window.preset == "previous_month": + end = logical_date.replace(day=1) - timedelta(days=1) + return end.replace(day=1), end + raise ScheduleBindingError(f"Unsupported window preset: {window.preset}") + + +def resolve_schedule_parameters( + *, + param_schema: dict[str, Any], + parameter_bindings: dict[str, ScheduleParameterBinding] | dict[str, object], + timezone_name: str, + scheduled_for: datetime, + window: ScheduleWindow | dict[str, object] | None = None, +) -> ResolvedScheduleContext: + """Resolve and type-check every SQL bind for one logical schedule run.""" + if scheduled_for.tzinfo is None: + raise ScheduleBindingError("scheduled_for must include a timezone.") + + bindings = _normalise_bindings(parameter_bindings) + expected = set(param_schema) + supplied = set(bindings) + missing = sorted(expected - supplied) + unknown = sorted(supplied - expected) + if missing: + names = ", ".join(f":{name}" for name in missing) + raise ScheduleBindingError(f"Missing schedule bindings: {names}") + if unknown: + names = ", ".join(f":{name}" for name in unknown) + raise ScheduleBindingError(f"Unknown schedule bindings: {names}") + + logical_date = scheduled_for.astimezone(ZoneInfo(timezone_name)).date() + parsed_window = ( + window + if isinstance(window, ScheduleWindow) + else ScheduleWindow.model_validate(window) + if window is not None + else None + ) + uses_window = any( + binding.source in {"window_start", "window_end"} for binding in bindings.values() + ) + if uses_window and parsed_window is None: + raise ScheduleBindingError("Window configuration is required for window bindings.") + + window_start: date | None = None + window_end: date | None = None + if parsed_window is not None: + window_start, window_end = _resolve_window(logical_date, parsed_window) + + raw_parameters: dict[str, object] = {} + for name, binding in bindings.items(): + descriptor = param_schema.get(name) + if not isinstance(descriptor, dict): + raise ScheduleBindingError(f"Invalid parameter descriptor for :{name}.") + param_type = descriptor.get("type", "string") + is_required = bool(descriptor.get("required", True)) + + if binding.source == "null": + if is_required: + raise ScheduleBindingError(f":{name} cannot use SQL NULL because it is required.") + raw_parameters[name] = None + continue + + if ( + binding.source + in { + "run_date", + "relative_date", + "window_start", + "window_end", + } + and param_type != "date" + ): + raise ScheduleBindingError(f":{name} must be a date parameter to use {binding.source}.") + + if binding.source == "literal": + raw_parameters[name] = binding.value + elif binding.source == "run_date": + raw_parameters[name] = logical_date + elif binding.source == "relative_date": + raw_parameters[name] = logical_date + timedelta(days=binding.offset_days or 0) + elif binding.source == "window_start": + raw_parameters[name] = window_start + elif binding.source == "window_end": + raw_parameters[name] = window_end + + try: + ParamModel = build_param_model(param_schema) + parameters = ParamModel.model_validate(raw_parameters).model_dump() + except (ValidationError, ValueError) as exc: + raise ScheduleBindingError(f"Invalid resolved schedule parameters: {exc}") from exc + + return ResolvedScheduleContext( + parameters=parameters, + logical_date=logical_date, + window_start=window_start, + window_end=window_end, + ) + + +def preview_schedule_runs( + *, + schedule_type: str, + cron_expression: str | None, + interval_seconds: int | None, + timezone_name: str, + param_schema: dict[str, Any], + parameter_bindings: dict[str, ScheduleParameterBinding] | dict[str, object], + window: ScheduleWindow | dict[str, object] | None, + from_time: datetime | None = None, + count: int = 3, +) -> list[ScheduleRunPreview]: + """Return upcoming logical runs and their resolved parameters.""" + from apscheduler.triggers.cron import CronTrigger + from apscheduler.triggers.interval import IntervalTrigger + + zone = ZoneInfo(timezone_name) + cursor = from_time or datetime.now(UTC) + if cursor.tzinfo is None: + raise ScheduleBindingError("from_time must include a timezone.") + + try: + if schedule_type == "cron" and cron_expression: + trigger = CronTrigger.from_crontab(cron_expression, timezone=zone) + elif schedule_type == "interval" and interval_seconds: + trigger = IntervalTrigger( + seconds=interval_seconds, + start_date=(cursor + timedelta(seconds=interval_seconds)).astimezone(zone), + timezone=zone, + ) + else: + raise ScheduleBindingError("Invalid schedule timing configuration.") + except ValueError as exc: + raise ScheduleBindingError(f"Invalid schedule timing configuration: {exc}") from exc + + previews: list[ScheduleRunPreview] = [] + previous_fire_time: datetime | None = None + for _ in range(count): + next_fire_time = trigger.get_next_fire_time(previous_fire_time, cursor) + if next_fire_time is None: + break + context = resolve_schedule_parameters( + param_schema=param_schema, + parameter_bindings=parameter_bindings, + timezone_name=timezone_name, + scheduled_for=next_fire_time, + window=window, + ) + previews.append( + ScheduleRunPreview( + scheduled_for=next_fire_time, + logical_date=context.logical_date, + window_start=context.window_start, + window_end=context.window_end, + resolved_parameters=context.parameters, + ) + ) + previous_fire_time = next_fire_time + cursor = next_fire_time + return previews diff --git a/backend/app/services/scheduler.py b/backend/app/services/scheduler.py index 20da67e..e3252e8 100644 --- a/backend/app/services/scheduler.py +++ b/backend/app/services/scheduler.py @@ -10,11 +10,16 @@ rehydrated from the application database whenever the process starts. """ +import hashlib +import json import uuid from datetime import UTC, datetime from typing import Any +from zoneinfo import ZoneInfo import structlog +from fastapi.encoders import jsonable_encoder +from sqlalchemy.exc import IntegrityError from app.config import settings from app.database import AsyncSessionLocal @@ -25,8 +30,9 @@ from app.repositories.job_run import JobRunRepository from app.repositories.schedule import ScheduleRepository from app.repositories.snapshot import SnapshotRepository -from app.sql.executor import SqlExecutionError, execute_query -from app.sql.param_models import build_param_model +from app.schemas.schedule import ScheduleWindow +from app.services.schedule_bindings import resolve_schedule_parameters +from app.sql.executor import execute_query log = structlog.get_logger().bind( request_id=None, @@ -55,13 +61,39 @@ def get_scheduler() -> Any: return _scheduler -def resolve_scheduled_params(param_schema: dict[str, Any]) -> dict[str, Any]: - """Resolve every scheduled bind default, retaining explicit SQL NULLs.""" - ParamModel = build_param_model(param_schema) - return ParamModel.model_validate({}).model_dump() +def _binding_hash( + *, + timezone_name: str, + parameter_bindings: dict[str, object], + window: dict[str, object] | None, +) -> str: + payload = { + "timezone": timezone_name, + "parameter_bindings": parameter_bindings, + "window": window, + } + encoded = json.dumps(payload, sort_keys=True, separators=(",", ":")) + return hashlib.sha256(encoded.encode()).hexdigest() + + +def _get_next_run_at(schedule_id: uuid.UUID) -> datetime | None: + if _scheduler is None: + return None + try: + job = _scheduler.get_job(str(schedule_id)) + next_run_time = getattr(job, "next_run_time", None) + return next_run_time if isinstance(next_run_time, datetime) else None + except Exception: # noqa: BLE001 + return None -async def execute_scheduled_job(schedule_id: str, endpoint_id: str) -> None: +async def execute_scheduled_job( + schedule_id: str, + endpoint_id: str, + *, + scheduled_for: datetime | None = None, + trigger_source: str = "scheduled", +) -> None: """Execute a single scheduled job — query Oracle, save snapshot. This function is called by APScheduler in a thread. We create our own @@ -86,16 +118,62 @@ async def execute_scheduled_job(schedule_id: str, endpoint_id: str) -> None: ep_repo = EndpointRepository(db) conn_repo = ConnectionRepository(db) - # Create a running job record + schedule = await sched_repo.get_by_id(sid) + if schedule is None: + log.error( + "scheduled_job_missing_schedule", + job_id=schedule_id, + endpoint_id=endpoint_id, + ) + return + + logical_run_time = ( + scheduled_for + or (schedule.next_run_at if trigger_source == "scheduled" else started_at) + or started_at + ) + if logical_run_time.tzinfo is None: + logical_run_time = logical_run_time.replace(tzinfo=UTC) + + existing_run = await job_repo.get_by_schedule_and_scheduled_for(sid, logical_run_time) + if existing_run is not None: + log.info( + "scheduled_job_duplicate_skipped", + job_id=schedule_id, + existing_run_id=str(existing_run.id), + scheduled_for=logical_run_time.isoformat(), + ) + return + + binding_hash = _binding_hash( + timezone_name=schedule.timezone, + parameter_bindings=schedule.parameter_bindings_json or {}, + window=schedule.window_config_json, + ) + + # Create a running job record before external I/O. The unique logical + # run key makes retries/multiple workers idempotent for a given fire time. job_run = JobRun( id=run_id, schedule_id=sid, endpoint_id=eid, started_at=started_at, status=JobRunStatus.running, + scheduled_for=logical_run_time, + trigger_source=trigger_source, + binding_hash=binding_hash, ) - await job_repo.create(job_run) - await db.commit() + try: + await job_repo.create(job_run) + await db.commit() + except IntegrityError: + await db.rollback() + log.info( + "scheduled_job_duplicate_skipped", + job_id=schedule_id, + scheduled_for=logical_run_time.isoformat(), + ) + return try: # Load endpoint and connection @@ -106,19 +184,36 @@ async def execute_scheduled_job(schedule_id: str, endpoint_id: str) -> None: if not endpoint.is_active: raise ValueError(f"Endpoint {eid} is not active.") + param_schema = endpoint.param_schema_json or {} + window = ( + ScheduleWindow.model_validate(schedule.window_config_json) + if schedule.window_config_json + else None + ) + context = resolve_schedule_parameters( + param_schema=param_schema, + parameter_bindings=schedule.parameter_bindings_json or {}, + timezone_name=schedule.timezone, + scheduled_for=logical_run_time, + window=window, + ) + resolved_parameters = jsonable_encoder(context.parameters) + await job_repo.update( + job_run, + { + "logical_date": context.logical_date, + "window_start": context.window_start, + "window_end": context.window_end, + "resolved_params_json": resolved_parameters, + }, + ) + await db.commit() + params = context.parameters + connection = await conn_repo.get_by_id(endpoint.connection_id) if connection is None or not connection.is_active: raise ValueError("Data source connection is unavailable.") - # Scheduled snapshots have no request inputs, so validate and - # resolve every configured static or dynamic default through the - # same typed model used by live data requests. - param_schema = endpoint.param_schema_json or {} - # Keep explicit NULL defaults in the bind dictionary. Omitting a - # None-valued entry makes Oracle see a missing bind variable rather - # than a bind whose value is SQL NULL. - params = resolve_scheduled_params(param_schema) - columns, rows, duration_ms = await execute_query( connection=connection, sql=endpoint.sql_text, @@ -171,14 +266,15 @@ async def execute_scheduled_job(schedule_id: str, endpoint_id: str) -> None: ) # Update schedule last_run_at - schedule = await sched_repo.get_by_id(sid) - if schedule: - await sched_repo.update( - schedule, - { - "last_run_at": finished_at, - }, - ) + current_schedule = await sched_repo.get_by_id(sid) + if current_schedule: + schedule_changes: dict[str, object] = { + "last_run_at": finished_at, + } + next_run_at = _get_next_run_at(sid) + if next_run_at is not None: + schedule_changes["next_run_at"] = next_run_at + await sched_repo.update(current_schedule, schedule_changes) await db.commit() @@ -189,10 +285,13 @@ async def execute_scheduled_job(schedule_id: str, endpoint_id: str) -> None: endpoint_id=endpoint_id, row_count=len(rows), duration_ms=duration_ms, + scheduled_for=logical_run_time.isoformat(), + logical_date=context.logical_date.isoformat(), + binding_hash=binding_hash, success=True, ) - except (SqlExecutionError, ValueError, Exception) as exc: + except Exception as exc: # noqa: BLE001 finished_at = datetime.now(UTC) error_detail = str(exc)[:5000] @@ -210,14 +309,15 @@ async def execute_scheduled_job(schedule_id: str, endpoint_id: str) -> None: ) # Update schedule last_run_at even on failure - schedule = await sched_repo.get_by_id(sid) - if schedule: - await sched_repo.update( - schedule, - { - "last_run_at": finished_at, - }, - ) + current_schedule = await sched_repo.get_by_id(sid) + if current_schedule: + schedule_changes = { + "last_run_at": finished_at, + } + next_run_at = _get_next_run_at(sid) + if next_run_at is not None: + schedule_changes["next_run_at"] = next_run_at + await sched_repo.update(current_schedule, schedule_changes) await db.commit() @@ -226,6 +326,8 @@ async def execute_scheduled_job(schedule_id: str, endpoint_id: str) -> None: job_id=schedule_id, run_id=str(run_id), endpoint_id=endpoint_id, + scheduled_for=logical_run_time.isoformat(), + binding_hash=binding_hash, error=error_detail, success=False, ) @@ -236,18 +338,23 @@ async def restore_active_schedules() -> int: restored = 0 failed_schedule_ids: list[str] = [] async with AsyncSessionLocal() as db: - schedules = await ScheduleRepository(db).get_all(active_only=True) + schedule_repo = ScheduleRepository(db) + schedules = await schedule_repo.get_all(active_only=True) + updated_next_run = False for schedule in schedules: try: - registered = add_schedule_job( + next_run_at = add_schedule_job( schedule_id=schedule.id, endpoint_id=schedule.endpoint_id, schedule_type=schedule.schedule_type, cron_expression=schedule.cron_expression, interval_seconds=schedule.interval_seconds, + timezone_name=getattr(schedule, "timezone", "UTC"), ) - if not registered: + if not isinstance(next_run_at, datetime): raise ValueError("Schedule configuration could not be registered.") + await schedule_repo.update(schedule, {"next_run_at": next_run_at}) + updated_next_run = True restored += 1 except Exception as exc: # noqa: BLE001 failed_schedule_ids.append(str(schedule.id)) @@ -256,6 +363,8 @@ async def restore_active_schedules() -> int: schedule_id=str(schedule.id), error=str(exc), ) + if updated_next_run: + await db.commit() if failed_schedule_ids: failed = ", ".join(failed_schedule_ids) raise RuntimeError(f"Failed to restore active schedules: {failed}") @@ -272,6 +381,7 @@ async def start_scheduler() -> None: from apscheduler.schedulers.asyncio import AsyncIOScheduler # noqa: PLC0415 scheduler = AsyncIOScheduler( + timezone=ZoneInfo("UTC"), job_defaults={ "coalesce": True, "max_instances": 1, @@ -312,11 +422,12 @@ def add_schedule_job( schedule_type: str, cron_expression: str | None = None, interval_seconds: int | None = None, -) -> bool: + timezone_name: str = "UTC", +) -> datetime | None: """Register a job in APScheduler.""" if _scheduler is None: log.warning("scheduler_not_running", action="add_job") - return False + return None from apscheduler.schedulers.asyncio import AsyncIOScheduler # noqa: PLC0415 @@ -345,20 +456,24 @@ def add_schedule_job( kwargs["day"] = parts[2] kwargs["month"] = parts[3] kwargs["day_of_week"] = parts[4] + kwargs["timezone"] = ZoneInfo(timezone_name) elif schedule_type == "interval" and interval_seconds: kwargs["trigger"] = "interval" kwargs["seconds"] = interval_seconds + kwargs["timezone"] = ZoneInfo(timezone_name) else: log.warning("invalid_schedule_config", schedule_id=str(schedule_id)) - return False + return None - scheduler.add_job(**kwargs) + job = scheduler.add_job(**kwargs) log.info( "scheduler_job_added", job_id=job_id, schedule_type=schedule_type, + timezone=timezone_name, ) - return True + next_run_time = getattr(job, "next_run_time", None) + return next_run_time if isinstance(next_run_time, datetime) else None def remove_schedule_job(schedule_id: uuid.UUID) -> None: diff --git a/backend/app/sql/param_models.py b/backend/app/sql/param_models.py index 3461a63..3ad1ed1 100644 --- a/backend/app/sql/param_models.py +++ b/backend/app/sql/param_models.py @@ -116,14 +116,19 @@ def _build_field( if isinstance(max_length, int) and max_length >= 1: annotation = Annotated[str, Field(max_length=max_length)] # type: ignore[assignment] - # Scheduled execution must resolve every configured default without - # request input. Live requests use ``enforce_required=True`` so a - # scheduler default never weakens their public required-parameter - # contract. Optional request fields continue to use their defaults. + # ``required`` governs whether the caller must supply the bind, while an + # optional parameter may still be supplied explicitly as SQL NULL. Keep + # that nullability independent from whichever endpoint default is stored. + if not required or default_is_null: + annotation = annotation | None # type: ignore[assignment] + + # Live requests use ``enforce_required=True`` so an endpoint default never + # weakens their public required-parameter contract. Optional request fields + # continue to use their endpoint defaults. Schedules resolve their separate + # binding configuration before validating the resulting values here. if enforce_required and required: field_default = ... elif default_is_null: - annotation = annotation | None # type: ignore[assignment] field_default = None elif default is not None: field_default = default @@ -137,7 +142,6 @@ def _build_field( # to ``T | None`` here. This path matches legacy behavior: # ``_coerce_param`` skipped optional params entirely when they # weren't supplied. - annotation = annotation | None # type: ignore[assignment] field_default = None return annotation, field_default diff --git a/backend/tests/test_endpoints.py b/backend/tests/test_endpoints.py index fb2d0f0..5bc786e 100644 --- a/backend/tests/test_endpoints.py +++ b/backend/tests/test_endpoints.py @@ -17,7 +17,6 @@ ParamDescriptor, SqlPreviewRequest, extract_bind_params, - require_snapshot_defaults, validate_sql_safety, ) @@ -155,28 +154,22 @@ def test_endpoint_update_rejects_incompatible_static_default() -> None: ) -def test_snapshot_default_validation_rejects_incompatible_stored_default() -> None: - with pytest.raises(ValueError, match="Invalid default"): - require_snapshot_defaults( - DataStrategy.snapshot, - {"store_id": {"type": "integer", "required": True, "default": "abc"}}, - ) - +def test_snapshot_endpoint_allows_schedule_to_own_parameter_values() -> None: + payload = EndpointCreate( + name="snapshot-without-defaults", + path="snapshot-without-defaults", + connection_id=uuid.uuid4(), + sql_text=("SELECT * FROM orders WHERE business_date BETWEEN :start_date AND :end_date"), + param_schema={ + "start_date": {"type": "date", "required": True}, + "end_date": {"type": "date", "required": True}, + }, + allow_unauthenticated=True, + data_strategy="snapshot", + ) -def test_snapshot_endpoint_requires_defaults_for_all_parameters() -> None: - with pytest.raises(ValueError, match=r"Missing: :end_date, :start_date"): - EndpointCreate( - name="snapshot-without-defaults", - path="snapshot-without-defaults", - connection_id=uuid.uuid4(), - sql_text=("SELECT * FROM orders WHERE business_date BETWEEN :start_date AND :end_date"), - param_schema={ - "start_date": {"type": "date", "required": True}, - "end_date": {"type": "date", "required": True}, - }, - allow_unauthenticated=True, - data_strategy="snapshot", - ) + assert payload.param_schema["start_date"].default is None + assert payload.param_schema["end_date"].default is None def test_snapshot_endpoint_accepts_dynamic_defaults() -> None: @@ -489,7 +482,7 @@ async def test_update_endpoint(async_client: object) -> None: @pytest.mark.integration -async def test_update_live_endpoint_to_snapshot_requires_merged_defaults( +async def test_update_live_endpoint_to_snapshot_leaves_values_to_schedule( async_client: object, ) -> None: from httpx import AsyncClient @@ -517,28 +510,13 @@ async def test_update_live_endpoint_to_snapshot_requires_merged_defaults( ) assert endpoint.status_code == 201 - invalid = await client.put( + updated = await client.put( f"/api/v1/admin/endpoints/{endpoint.json()['id']}", json={"data_strategy": "snapshot"}, ) - assert invalid.status_code == 422 - assert ":business_date" in invalid.json()["detail"] - - valid = await client.put( - f"/api/v1/admin/endpoints/{endpoint.json()['id']}", - json={ - "data_strategy": "snapshot", - "param_schema": { - "business_date": { - "type": "date", - "required": True, - "default_expression": "today", - } - }, - }, - ) - assert valid.status_code == 200 - assert valid.json()["param_schema"]["business_date"]["default_expression"] == "today" + assert updated.status_code == 200 + assert updated.json()["data_strategy"] == "snapshot" + assert updated.json()["param_schema"]["business_date"]["default"] is None @pytest.mark.integration diff --git a/backend/tests/test_schedule_bindings.py b/backend/tests/test_schedule_bindings.py new file mode 100644 index 0000000..39485ce --- /dev/null +++ b/backend/tests/test_schedule_bindings.py @@ -0,0 +1,263 @@ +"""Behavioral tests for schedule-owned parameter bindings.""" + +from datetime import UTC, date, datetime + +import pytest +from app.schemas.schedule import ScheduleParameterBinding, ScheduleWindow +from app.services.schedule_bindings import ( + ScheduleBindingError, + preview_schedule_runs, + resolve_schedule_parameters, +) + +DATE_RANGE_SCHEMA = { + "start_date": {"type": "date", "required": True}, + "end_date": {"type": "date", "required": True}, + "store_id": { + "type": "string", + "required": False, + "default_is_null": True, + }, +} + + +def test_resolve_schedule_parameters_uses_logical_date_in_schedule_timezone() -> None: + context = resolve_schedule_parameters( + param_schema=DATE_RANGE_SCHEMA, + parameter_bindings={ + "start_date": ScheduleParameterBinding(source="relative_date", offset_days=-7), + "end_date": ScheduleParameterBinding(source="run_date"), + "store_id": ScheduleParameterBinding(source="null"), + }, + timezone_name="Asia/Riyadh", + scheduled_for=datetime(2026, 8, 30, 21, 30, tzinfo=UTC), + ) + + assert context.logical_date == date(2026, 8, 31) + assert context.parameters == { + "start_date": date(2026, 8, 24), + "end_date": date(2026, 8, 31), + "store_id": None, + } + assert context.window_start is None + assert context.window_end is None + + +@pytest.mark.parametrize( + ("window", "scheduled_for", "expected_start", "expected_end"), + [ + ( + ScheduleWindow(preset="previous_day"), + datetime(2026, 8, 31, 3, tzinfo=UTC), + date(2026, 8, 30), + date(2026, 8, 30), + ), + ( + ScheduleWindow(preset="last_n_complete_days", days=7), + datetime(2026, 8, 31, 3, tzinfo=UTC), + date(2026, 8, 24), + date(2026, 8, 30), + ), + ( + ScheduleWindow(preset="week_to_date"), + datetime(2026, 9, 2, 3, tzinfo=UTC), + date(2026, 8, 31), + date(2026, 9, 2), + ), + ( + ScheduleWindow(preset="previous_week"), + datetime(2026, 9, 2, 3, tzinfo=UTC), + date(2026, 8, 24), + date(2026, 8, 30), + ), + ( + ScheduleWindow(preset="month_to_date"), + datetime(2024, 2, 29, 3, tzinfo=UTC), + date(2024, 2, 1), + date(2024, 2, 29), + ), + ( + ScheduleWindow(preset="previous_month"), + datetime(2024, 3, 1, 3, tzinfo=UTC), + date(2024, 2, 1), + date(2024, 2, 29), + ), + ], +) +def test_resolve_schedule_parameters_supports_calendar_windows( + window: ScheduleWindow, + scheduled_for: datetime, + expected_start: date, + expected_end: date, +) -> None: + context = resolve_schedule_parameters( + param_schema={ + "start_date": {"type": "date", "required": True}, + "end_date": {"type": "date", "required": True}, + }, + parameter_bindings={ + "start_date": ScheduleParameterBinding(source="window_start"), + "end_date": ScheduleParameterBinding(source="window_end"), + }, + window=window, + timezone_name="UTC", + scheduled_for=scheduled_for, + ) + + assert context.parameters == { + "start_date": expected_start, + "end_date": expected_end, + } + assert context.window_start == expected_start + assert context.window_end == expected_end + + +def test_resolve_schedule_parameters_coerces_typed_literals() -> None: + context = resolve_schedule_parameters( + param_schema={ + "limit": {"type": "integer", "required": True}, + "enabled": {"type": "boolean", "required": True}, + "as_of": {"type": "date", "required": True}, + }, + parameter_bindings={ + "limit": ScheduleParameterBinding(source="literal", value="25"), + "enabled": ScheduleParameterBinding(source="literal", value="yes"), + "as_of": ScheduleParameterBinding(source="literal", value="31-08-2026"), + }, + timezone_name="UTC", + scheduled_for=datetime(2026, 8, 31, tzinfo=UTC), + ) + + assert context.parameters == { + "limit": 25, + "enabled": True, + "as_of": date(2026, 8, 31), + } + + +def test_schedule_can_override_optional_endpoint_default_with_sql_null() -> None: + context = resolve_schedule_parameters( + param_schema={ + "store_id": { + "type": "string", + "required": False, + "default": "ALL", + } + }, + parameter_bindings={ + "store_id": ScheduleParameterBinding(source="null"), + }, + timezone_name="UTC", + scheduled_for=datetime(2026, 8, 31, tzinfo=UTC), + ) + + assert context.parameters == {"store_id": None} + + +@pytest.mark.parametrize( + ("bindings", "message"), + [ + ({}, "Missing schedule bindings: :end_date, :start_date, :store_id"), + ( + { + "start_date": ScheduleParameterBinding(source="run_date"), + "end_date": ScheduleParameterBinding(source="run_date"), + "store_id": ScheduleParameterBinding(source="null"), + "unknown": ScheduleParameterBinding(source="literal", value="x"), + }, + "Unknown schedule bindings: :unknown", + ), + ], +) +def test_resolve_schedule_parameters_requires_exact_bind_coverage( + bindings: dict[str, ScheduleParameterBinding], message: str +) -> None: + with pytest.raises(ScheduleBindingError, match=message): + resolve_schedule_parameters( + param_schema=DATE_RANGE_SCHEMA, + parameter_bindings=bindings, + timezone_name="UTC", + scheduled_for=datetime(2026, 8, 31, tzinfo=UTC), + ) + + +def test_resolve_schedule_parameters_rejects_null_for_required_parameter() -> None: + with pytest.raises(ScheduleBindingError, match=":start_date cannot use SQL NULL"): + resolve_schedule_parameters( + param_schema=DATE_RANGE_SCHEMA, + parameter_bindings={ + "start_date": ScheduleParameterBinding(source="null"), + "end_date": ScheduleParameterBinding(source="run_date"), + "store_id": ScheduleParameterBinding(source="null"), + }, + timezone_name="UTC", + scheduled_for=datetime(2026, 8, 31, tzinfo=UTC), + ) + + +def test_resolve_schedule_parameters_requires_window_for_window_sources() -> None: + with pytest.raises(ScheduleBindingError, match="Window configuration is required"): + resolve_schedule_parameters( + param_schema=DATE_RANGE_SCHEMA, + parameter_bindings={ + "start_date": ScheduleParameterBinding(source="window_start"), + "end_date": ScheduleParameterBinding(source="window_end"), + "store_id": ScheduleParameterBinding(source="null"), + }, + timezone_name="UTC", + scheduled_for=datetime(2026, 8, 31, tzinfo=UTC), + ) + + +def test_preview_schedule_runs_returns_next_three_logical_contexts() -> None: + previews = preview_schedule_runs( + schedule_type="cron", + cron_expression="0 6 * * *", + interval_seconds=None, + timezone_name="Asia/Riyadh", + param_schema={"run_date": {"type": "date", "required": True}}, + parameter_bindings={ + "run_date": ScheduleParameterBinding(source="run_date"), + }, + window=None, + from_time=datetime(2026, 8, 30, 20, tzinfo=UTC), + count=3, + ) + + assert [preview.scheduled_for.isoformat() for preview in previews] == [ + "2026-08-31T06:00:00+03:00", + "2026-09-01T06:00:00+03:00", + "2026-09-02T06:00:00+03:00", + ] + assert [preview.logical_date for preview in previews] == [ + date(2026, 8, 31), + date(2026, 9, 1), + date(2026, 9, 2), + ] + assert [preview.resolved_parameters["run_date"] for preview in previews] == [ + date(2026, 8, 31), + date(2026, 9, 1), + date(2026, 9, 2), + ] + + +def test_preview_schedule_runs_preserves_local_cron_time_across_dst() -> None: + previews = preview_schedule_runs( + schedule_type="cron", + cron_expression="0 1 * * *", + interval_seconds=None, + timezone_name="America/New_York", + param_schema={"run_date": {"type": "date", "required": True}}, + parameter_bindings={ + "run_date": ScheduleParameterBinding(source="run_date"), + }, + window=None, + from_time=datetime(2026, 3, 7, tzinfo=UTC), + count=3, + ) + + assert [preview.scheduled_for.isoformat() for preview in previews] == [ + "2026-03-07T01:00:00-05:00", + "2026-03-08T01:00:00-05:00", + "2026-03-09T01:00:00-04:00", + ] diff --git a/backend/tests/test_scheduler_restore.py b/backend/tests/test_scheduler_restore.py index 2abb05c..11c4158 100644 --- a/backend/tests/test_scheduler_restore.py +++ b/backend/tests/test_scheduler_restore.py @@ -1,6 +1,7 @@ """Unit coverage for restoring in-memory scheduler jobs after restart.""" import uuid +from datetime import UTC, datetime from types import SimpleNamespace from typing import cast from unittest.mock import AsyncMock, MagicMock @@ -9,22 +10,6 @@ from structlog.testing import capture_logs -def test_scheduled_params_retain_explicit_null_bind() -> None: - from app.services.scheduler import resolve_scheduled_params - - params = resolve_scheduled_params( - { - "str_id": { - "type": "string", - "required": False, - "default_is_null": True, - } - } - ) - - assert params == {"str_id": None} - - @pytest.mark.asyncio async def test_restore_active_schedules_registers_each_row( monkeypatch: pytest.MonkeyPatch, @@ -47,10 +32,11 @@ async def test_restore_active_schedules_registers_each_row( interval_seconds=300, ), ] + db = SimpleNamespace(commit=AsyncMock()) class FakeSessionContext: async def __aenter__(self) -> object: - return object() + return db async def __aexit__(self, *args: object) -> None: return None @@ -63,7 +49,11 @@ async def get_all(self, *, active_only: bool = False) -> list[object]: assert active_only is True return cast(list[object], schedules) - add_job = MagicMock() + async def update(self, schedule: object, values: dict[str, object]) -> None: + setattr(schedule, "next_run_at", values["next_run_at"]) + + next_run = datetime(2026, 8, 31, 6, tzinfo=UTC) + add_job = MagicMock(return_value=next_run) monkeypatch.setattr(scheduler_service, "AsyncSessionLocal", FakeSessionContext) monkeypatch.setattr(scheduler_service, "ScheduleRepository", FakeScheduleRepository) monkeypatch.setattr(scheduler_service, "add_schedule_job", add_job) @@ -73,6 +63,8 @@ async def get_all(self, *, active_only: bool = False) -> list[object]: assert restored == 2 assert add_job.call_count == 2 assert add_job.call_args_list[0].kwargs["schedule_id"] == schedules[0].id + assert all(schedule.next_run_at == next_run for schedule in schedules) + db.commit.assert_awaited_once() @pytest.mark.asyncio @@ -141,3 +133,26 @@ async def test_start_scheduler_cleans_up_and_propagates_restore_failure( assert failure["duration_ms"] is None assert failure["method"] == "SCHEDULE" assert failure["client_ip"] is None + + +def test_add_schedule_job_uses_schedule_timezone_and_returns_next_run( + monkeypatch: pytest.MonkeyPatch, +) -> None: + from app.services import scheduler as scheduler_service + + next_run = datetime(2026, 8, 31, 6, tzinfo=UTC) + fake_scheduler = MagicMock() + fake_scheduler.add_job.return_value = SimpleNamespace(next_run_time=next_run) + monkeypatch.setattr(scheduler_service, "_scheduler", fake_scheduler) + + result = scheduler_service.add_schedule_job( + schedule_id=uuid.uuid4(), + endpoint_id=uuid.uuid4(), + schedule_type="cron", + cron_expression="0 6 * * *", + timezone_name="Asia/Riyadh", + ) + + assert result == next_run + kwargs = fake_scheduler.add_job.call_args.kwargs + assert str(kwargs["timezone"]) == "Asia/Riyadh" diff --git a/backend/tests/test_schedules.py b/backend/tests/test_schedules.py index 6668180..3e05aac 100644 --- a/backend/tests/test_schedules.py +++ b/backend/tests/test_schedules.py @@ -7,14 +7,17 @@ """ import uuid -from datetime import UTC, datetime +from datetime import UTC, date, datetime import pytest from app.schemas.schedule import ( JobRunResponse, ScheduleCreate, + ScheduleParameterBinding, + SchedulePreviewResponse, ScheduleResponse, ScheduleUpdate, + ScheduleWindow, SnapshotDetailResponse, SnapshotResponse, ) @@ -42,6 +45,34 @@ def test_schedule_create_interval_valid() -> None: assert payload.interval_seconds == 300 +def test_schedule_create_accepts_timezone_and_declarative_bindings() -> None: + payload = ScheduleCreate( + endpoint_id=uuid.uuid4(), + schedule_type="cron", + cron_expression="0 6 * * *", + timezone="Asia/Riyadh", + parameter_bindings={ + "start_date": ScheduleParameterBinding(source="window_start"), + "end_date": ScheduleParameterBinding(source="window_end"), + }, + window=ScheduleWindow(preset="last_n_complete_days", days=7), + ) + + assert payload.timezone == "Asia/Riyadh" + assert payload.parameter_bindings["start_date"].source == "window_start" + assert payload.window == ScheduleWindow(preset="last_n_complete_days", days=7) + + +def test_schedule_create_rejects_unknown_timezone() -> None: + with pytest.raises(ValueError, match="Unknown IANA timezone"): + ScheduleCreate( + endpoint_id=uuid.uuid4(), + schedule_type="cron", + cron_expression="0 6 * * *", + timezone="Mars/Olympus_Mons", + ) + + def test_schedule_create_cron_requires_expression() -> None: with pytest.raises(ValueError, match="cron_expression is required"): ScheduleCreate( @@ -128,6 +159,14 @@ def test_schedule_response_fields() -> None: assert "is_active" in fields assert "last_run_at" in fields assert "next_run_at" in fields + assert "timezone" in fields + assert "parameter_bindings" in fields + assert "window" in fields + + +def test_schedule_preview_response_contains_resolved_run_context() -> None: + fields = SchedulePreviewResponse.model_fields + assert "runs" in fields def test_job_run_response_fields() -> None: @@ -140,6 +179,13 @@ def test_job_run_response_fields() -> None: assert "status" in fields assert "row_count" in fields assert "error_detail" in fields + assert "scheduled_for" in fields + assert "logical_date" in fields + assert "window_start" in fields + assert "window_end" in fields + assert "resolved_parameters" in fields + assert "trigger_source" in fields + assert "binding_hash" in fields def test_job_run_response_allows_deleted_schedule_history() -> None: @@ -192,6 +238,273 @@ def test_snapshot_detail_response_fields() -> None: # ── API integration tests (require PostgreSQL) ────────────────────────────── +async def _create_snapshot_endpoint_with_date_range(client: object) -> str: + from httpx import AsyncClient + + typed_client: AsyncClient = client # type: ignore[assignment] + connection = await typed_client.post( + "/api/v1/admin/connections/", + json={ + "name": f"schedule-binding-conn-{uuid.uuid4().hex[:8]}", + "host": "oracle.example.com", + "service_name": "ORCLPDB", + "username": "hr", + "password": "secret", + }, + ) + assert connection.status_code == 201 + endpoint = await typed_client.post( + "/api/v1/admin/endpoints/", + json={ + "name": f"schedule-binding-ep-{uuid.uuid4().hex[:8]}", + "path": f"schedule-binding-path-{uuid.uuid4().hex[:8]}", + "connection_id": connection.json()["id"], + "allow_unauthenticated": True, + "sql_text": ( + "SELECT * FROM orders WHERE business_date BETWEEN :start_date AND :end_date" + ), + "param_schema": { + "start_date": {"type": "date", "required": True}, + "end_date": {"type": "date", "required": True}, + }, + "data_strategy": "snapshot", + }, + ) + assert endpoint.status_code == 201 + return str(endpoint.json()["id"]) + + +@pytest.mark.integration +async def test_schedule_bindings_are_required_at_schedule_creation( + async_client: object, +) -> None: + from httpx import AsyncClient + + client: AsyncClient = async_client # type: ignore[assignment] + endpoint_id = await _create_snapshot_endpoint_with_date_range(client) + + missing = await client.post( + "/api/v1/admin/schedules/", + json={ + "endpoint_id": endpoint_id, + "schedule_type": "cron", + "cron_expression": "0 6 * * *", + "timezone": "Asia/Riyadh", + "parameter_bindings": { + "end_date": {"source": "run_date"}, + }, + }, + ) + assert missing.status_code == 422 + assert "Missing schedule bindings: :start_date" in missing.json()["detail"] + + created = await client.post( + "/api/v1/admin/schedules/", + json={ + "endpoint_id": endpoint_id, + "schedule_type": "cron", + "cron_expression": "0 6 * * *", + "timezone": "Asia/Riyadh", + "parameter_bindings": { + "start_date": {"source": "relative_date", "offset_days": -7}, + "end_date": {"source": "run_date"}, + }, + }, + ) + assert created.status_code == 201 + assert created.json()["timezone"] == "Asia/Riyadh" + assert created.json()["parameter_bindings"]["start_date"] == { + "source": "relative_date", + "value": None, + "offset_days": -7, + } + + +@pytest.mark.integration +async def test_schedule_rejects_live_endpoint(async_client: object) -> None: + from httpx import AsyncClient + + client: AsyncClient = async_client # type: ignore[assignment] + connection = await client.post( + "/api/v1/admin/connections/", + json={ + "name": f"live-schedule-conn-{uuid.uuid4().hex[:8]}", + "host": "oracle.example.com", + "service_name": "ORCLPDB", + "username": "hr", + "password": "secret", + }, + ) + endpoint = await client.post( + "/api/v1/admin/endpoints/", + json={ + "name": f"live-schedule-ep-{uuid.uuid4().hex[:8]}", + "path": f"live-schedule-path-{uuid.uuid4().hex[:8]}", + "connection_id": connection.json()["id"], + "allow_unauthenticated": True, + "sql_text": "SELECT 1 FROM dual", + "data_strategy": "live", + }, + ) + + response = await client.post( + "/api/v1/admin/schedules/", + json={ + "endpoint_id": endpoint.json()["id"], + "schedule_type": "interval", + "interval_seconds": 60, + }, + ) + + assert response.status_code == 422 + assert response.json()["detail"] == "Schedules are supported only for snapshot endpoints." + + +@pytest.mark.integration +async def test_preview_schedule_resolves_next_runs(async_client: object) -> None: + from httpx import AsyncClient + + client: AsyncClient = async_client # type: ignore[assignment] + endpoint_id = await _create_snapshot_endpoint_with_date_range(client) + + response = await client.post( + "/api/v1/admin/schedules/preview", + json={ + "endpoint_id": endpoint_id, + "schedule_type": "cron", + "cron_expression": "0 6 * * *", + "timezone": "Asia/Riyadh", + "parameter_bindings": { + "start_date": {"source": "window_start"}, + "end_date": {"source": "window_end"}, + }, + "window": {"preset": "last_n_complete_days", "days": 7}, + "count": 3, + }, + ) + + assert response.status_code == 200 + runs = response.json()["runs"] + assert len(runs) == 3 + assert all(run["scheduled_for"].endswith("+03:00") for run in runs) + assert all(run["resolved_parameters"]["start_date"] for run in runs) + assert all(run["resolved_parameters"]["end_date"] for run in runs) + + +@pytest.mark.integration +async def test_scheduled_execution_persists_logical_context_and_is_idempotent( + async_client: object, monkeypatch: pytest.MonkeyPatch +) -> None: + from unittest.mock import AsyncMock + + from app.services import scheduler as scheduler_service + from httpx import AsyncClient + + client: AsyncClient = async_client # type: ignore[assignment] + endpoint_id = await _create_snapshot_endpoint_with_date_range(client) + created = await client.post( + "/api/v1/admin/schedules/", + json={ + "endpoint_id": endpoint_id, + "schedule_type": "cron", + "cron_expression": "0 6 * * *", + "timezone": "Asia/Riyadh", + "parameter_bindings": { + "start_date": {"source": "window_start"}, + "end_date": {"source": "window_end"}, + }, + "window": {"preset": "last_n_complete_days", "days": 7}, + }, + ) + assert created.status_code == 201 + schedule_id = created.json()["id"] + + execute_query = AsyncMock( + return_value=( + ["ORDER_ID"], + [{"ORDER_ID": 42}], + 12, + ) + ) + monkeypatch.setattr(scheduler_service, "execute_query", execute_query) + scheduled_for = datetime(2026, 8, 31, 3, tzinfo=UTC) + + await scheduler_service.execute_scheduled_job( + schedule_id, + endpoint_id, + scheduled_for=scheduled_for, + ) + await scheduler_service.execute_scheduled_job( + schedule_id, + endpoint_id, + scheduled_for=scheduled_for, + ) + + runs_response = await client.get( + "/api/v1/admin/schedules/jobs/", + params={"schedule_id": schedule_id}, + ) + assert runs_response.status_code == 200 + runs = runs_response.json() + assert len(runs) == 1 + assert runs[0]["status"] == "success" + assert runs[0]["scheduled_for"] == "2026-08-31T03:00:00Z" + assert runs[0]["logical_date"] == "2026-08-31" + assert runs[0]["window_start"] == "2026-08-24" + assert runs[0]["window_end"] == "2026-08-30" + assert runs[0]["resolved_parameters"] == { + "start_date": "2026-08-24", + "end_date": "2026-08-30", + } + assert runs[0]["trigger_source"] == "scheduled" + assert len(runs[0]["binding_hash"]) == 64 + execute_query.assert_awaited_once() + assert execute_query.await_args is not None + assert execute_query.await_args.kwargs["params"] == { + "start_date": date(2026, 8, 24), + "end_date": date(2026, 8, 30), + } + + +@pytest.mark.integration +async def test_run_now_accepts_an_explicit_logical_date( + async_client: object, monkeypatch: pytest.MonkeyPatch +) -> None: + from unittest.mock import AsyncMock + + from app.services import schedule as schedule_service + from httpx import AsyncClient + + client: AsyncClient = async_client # type: ignore[assignment] + endpoint_id = await _create_snapshot_endpoint_with_date_range(client) + created = await client.post( + "/api/v1/admin/schedules/", + json={ + "endpoint_id": endpoint_id, + "schedule_type": "cron", + "cron_expression": "0 6 * * *", + "timezone": "Asia/Riyadh", + "parameter_bindings": { + "start_date": {"source": "relative_date", "offset_days": -1}, + "end_date": {"source": "run_date"}, + }, + }, + ) + execute_job = AsyncMock() + monkeypatch.setattr(schedule_service, "execute_scheduled_job", execute_job) + + response = await client.post( + f"/api/v1/admin/schedules/{created.json()['id']}/run", + json={"logical_date": "2026-08-30"}, + ) + + assert response.status_code == 202 + execute_job.assert_awaited_once() + assert execute_job.await_args is not None + assert execute_job.await_args.kwargs["trigger_source"] == "manual" + assert execute_job.await_args.kwargs["scheduled_for"].isoformat() == "2026-08-30T00:00:00+03:00" + + @pytest.mark.integration async def test_create_schedule(async_client: object) -> None: from httpx import AsyncClient From 0bfd889f204cb40bdf778fc039629610424a1169 Mon Sep 17 00:00:00 2001 From: Badry Date: Sun, 30 Aug 2026 22:30:08 +0300 Subject: [PATCH 2/7] feat(schedules): add binding editor and run preview --- .../endpoints/EndpointWizard.test.tsx | 7 +- .../components/endpoints/EndpointWizard.tsx | 6 +- .../endpoints/wizard/ConfigStep.test.tsx | 17 +- .../endpoints/wizard/ConfigStep.tsx | 14 +- .../endpoints/wizard/ParamsStep.tsx | 6 +- .../wizard/parameterDefaults.test.ts | 31 --- .../endpoints/wizard/parameterDefaults.ts | 7 - .../ScheduleParameterBindings.test.tsx | 68 ++++++ .../schedules/ScheduleParameterBindings.tsx | 226 ++++++++++++++++++ .../schedules/scheduleBindings.test.ts | 56 +++++ .../components/schedules/scheduleBindings.ts | 49 ++++ frontend/src/lib/api.ts | 12 +- frontend/src/pages/SchedulesPage.tsx | 184 +++++++++++++- frontend/src/types/schedule.ts | 68 ++++++ 14 files changed, 672 insertions(+), 79 deletions(-) create mode 100644 frontend/src/components/schedules/ScheduleParameterBindings.test.tsx create mode 100644 frontend/src/components/schedules/ScheduleParameterBindings.tsx create mode 100644 frontend/src/components/schedules/scheduleBindings.test.ts create mode 100644 frontend/src/components/schedules/scheduleBindings.ts diff --git a/frontend/src/components/endpoints/EndpointWizard.test.tsx b/frontend/src/components/endpoints/EndpointWizard.test.tsx index 1a51a58..a368248 100644 --- a/frontend/src/components/endpoints/EndpointWizard.test.tsx +++ b/frontend/src/components/endpoints/EndpointWizard.test.tsx @@ -106,7 +106,7 @@ describe("EndpointWizard preview coordination", () => { expect(previewMock.mock.calls[0][0].params.customer_id).toMatch(/^\d{4}-\d{2}-\d{2}$/); }); - it("blocks snapshot review when a parameter has no schedule-time default", async () => { + it("allows snapshot review when its schedule will own parameter values", async () => { renderWizard(); await advanceToParameters(); @@ -121,7 +121,8 @@ describe("EndpointWizard preview coordination", () => { fireEvent.change(configSelectors[1], { target: { value: "snapshot" } }); fireEvent.click(screen.getByRole("checkbox")); - expect(screen.getByRole("button", { name: "Next" })).toBeDisabled(); - expect(screen.getByText(/:customer_id/)).toBeInTheDocument(); + expect(screen.getByRole("button", { name: "Next" })).toBeEnabled(); + expect(screen.getByText(/Scheduled values are configured with the schedule/i)).toBeInTheDocument(); + expect(screen.getByText(/logical run date/i)).toBeInTheDocument(); }); }); diff --git a/frontend/src/components/endpoints/EndpointWizard.tsx b/frontend/src/components/endpoints/EndpointWizard.tsx index 06be0f8..defa9bb 100644 --- a/frontend/src/components/endpoints/EndpointWizard.tsx +++ b/frontend/src/components/endpoints/EndpointWizard.tsx @@ -23,7 +23,6 @@ import { ReviewStep } from "./wizard/ReviewStep"; import { SqlStep } from "./wizard/SqlStep"; import { extractBindParams, reconcileParamSchema } from "./wizard/bindParams"; import { - missingSnapshotDefaults, resolvePreviewParameterDefault, updateParameterDescriptor, } from "./wizard/parameterDefaults"; @@ -147,10 +146,7 @@ export function EndpointWizard({ onSuccess, onCancel }: EndpointWizardProps) { // An endpoint without a dedicated method requires an explicit opt-in to // platform-admin Bearer authentication, mirroring the server-side 422. const authOk = !!state.auth_method_id || state.allow_unauthenticated; - const snapshotDefaultsOk = - state.data_strategy !== "snapshot" || - missingSnapshotDefaults(state.param_schema).length === 0; - return !!state.name.trim() && !!state.path.trim() && authOk && snapshotDefaultsOk; + return !!state.name.trim() && !!state.path.trim() && authOk; } return true; }; diff --git a/frontend/src/components/endpoints/wizard/ConfigStep.test.tsx b/frontend/src/components/endpoints/wizard/ConfigStep.test.tsx index 7b9b69f..baf7746 100644 --- a/frontend/src/components/endpoints/wizard/ConfigStep.test.tsx +++ b/frontend/src/components/endpoints/wizard/ConfigStep.test.tsx @@ -47,8 +47,8 @@ describe("ConfigStep platform authentication fallback", () => { }); }); -describe("ConfigStep snapshot defaults", () => { - it("lists snapshot parameters that have no default", () => { +describe("ConfigStep snapshot scheduling guidance", () => { + it("explains that scheduled values are configured separately", () => { render( { />, ); - expect(screen.getByText(/Snapshot defaults required/i)).toBeInTheDocument(); - expect(screen.getByText(/:end_date/)).toBeInTheDocument(); - expect(screen.queryByText(/:start_date/)).not.toBeInTheDocument(); + expect( + screen.getByText(/Scheduled values are configured with the schedule/i), + ).toBeInTheDocument(); + expect(screen.getByText(/logical run date/i)).toBeInTheDocument(); }); - it("accepts false and zero as configured snapshot defaults", () => { + it("shows the guidance even when endpoint defaults exist", () => { render( { />, ); - expect(screen.queryByText(/Snapshot defaults required/i)).not.toBeInTheDocument(); + expect( + screen.getByText(/Scheduled values are configured with the schedule/i), + ).toBeInTheDocument(); }); }); diff --git a/frontend/src/components/endpoints/wizard/ConfigStep.tsx b/frontend/src/components/endpoints/wizard/ConfigStep.tsx index eb0efc1..38ab230 100644 --- a/frontend/src/components/endpoints/wizard/ConfigStep.tsx +++ b/frontend/src/components/endpoints/wizard/ConfigStep.tsx @@ -7,7 +7,6 @@ import type { AuthMethod } from "@/types/auth_method"; import type { DataStrategy } from "@/types/endpoint"; import type { WizardState, WizardUpdate } from "./types"; -import { missingSnapshotDefaults } from "./parameterDefaults"; interface ConfigStepProps { state: WizardState; @@ -16,8 +15,6 @@ interface ConfigStepProps { } export function ConfigStep({ state, update, authMethods }: ConfigStepProps) { - const missingDefaults = missingSnapshotDefaults(state.param_schema); - return (

Endpoint Configuration

@@ -94,13 +91,12 @@ export function ConfigStep({ state, update, authMethods }: ConfigStepProps) {
- {state.data_strategy === "snapshot" && missingDefaults.length > 0 && ( - - Snapshot defaults required + {state.data_strategy === "snapshot" && ( + + Scheduled values are configured with the schedule - Add a fixed, dynamic, or NULL default for{" "} - {missingDefaults.map((name) => `:${name}`).join(", ")} in the Parameters step. Scheduled - snapshots have no request values to bind. + After publishing this endpoint, create its schedule and choose a fixed value, SQL NULL, + logical run date, relative date, or calendar window for every SQL parameter. )} diff --git a/frontend/src/components/endpoints/wizard/ParamsStep.tsx b/frontend/src/components/endpoints/wizard/ParamsStep.tsx index 74f1352..27dc6ed 100644 --- a/frontend/src/components/endpoints/wizard/ParamsStep.tsx +++ b/frontend/src/components/endpoints/wizard/ParamsStep.tsx @@ -171,9 +171,9 @@ export function ParamsStep({ state, onUpdateParam }: ParamsStepProps) { Define types and defaults for bind parameters detected in your query.

{state.data_strategy === "snapshot" && entries.length > 0 && ( -

- Snapshot endpoints require a fixed, NULL, or dynamic default for every parameter. Dynamic - dates are evaluated from the application server date whenever the snapshot runs. +

+ Endpoint defaults remain useful for live requests and previews. Scheduled values are + configured separately when you create the snapshot schedule.

)} {entries.length === 0 ? ( diff --git a/frontend/src/components/endpoints/wizard/parameterDefaults.test.ts b/frontend/src/components/endpoints/wizard/parameterDefaults.test.ts index 6cb80f1..d6f9185 100644 --- a/frontend/src/components/endpoints/wizard/parameterDefaults.test.ts +++ b/frontend/src/components/endpoints/wizard/parameterDefaults.test.ts @@ -2,42 +2,11 @@ import { describe, expect, it } from "vitest"; import { describeParameterDefault, - hasParameterDefault, - missingSnapshotDefaults, resolvePreviewParameterDefault, updateParameterDescriptor, } from "./parameterDefaults"; describe("parameter default helpers", () => { - it("counts false, zero, empty strings, and dynamic expressions as defaults", () => { - expect(hasParameterDefault({ type: "boolean", required: true, default: false })).toBe(true); - expect(hasParameterDefault({ type: "integer", required: true, default: 0 })).toBe(true); - expect(hasParameterDefault({ type: "string", required: true, default: "" })).toBe(true); - expect( - hasParameterDefault({ - type: "date", - required: true, - default_expression: "today", - }), - ).toBe(true); - expect( - hasParameterDefault({ - type: "string", - required: false, - default_is_null: true, - }), - ).toBe(true); - }); - - it("returns only parameter names with no default", () => { - expect( - missingSnapshotDefaults({ - start_date: { type: "date", required: true, default_expression: "yesterday" }, - end_date: { type: "date", required: true, default: null }, - }), - ).toEqual(["end_date"]); - }); - it("describes dynamic defaults in server-date terms", () => { expect( describeParameterDefault({ diff --git a/frontend/src/components/endpoints/wizard/parameterDefaults.ts b/frontend/src/components/endpoints/wizard/parameterDefaults.ts index 1bdfee0..40d260b 100644 --- a/frontend/src/components/endpoints/wizard/parameterDefaults.ts +++ b/frontend/src/components/endpoints/wizard/parameterDefaults.ts @@ -7,13 +7,6 @@ export function hasParameterDefault(descriptor: ParamDescriptor): boolean { : descriptor.default_expression === "today" || descriptor.default_expression === "yesterday"; } -export function missingSnapshotDefaults(paramSchema: Record): string[] { - return Object.entries(paramSchema) - .filter(([, descriptor]) => !hasParameterDefault(descriptor)) - .map(([name]) => name) - .sort(); -} - function formatLocalDate(value: Date): string { const year = value.getFullYear(); const month = String(value.getMonth() + 1).padStart(2, "0"); diff --git a/frontend/src/components/schedules/ScheduleParameterBindings.test.tsx b/frontend/src/components/schedules/ScheduleParameterBindings.test.tsx new file mode 100644 index 0000000..1ee22f0 --- /dev/null +++ b/frontend/src/components/schedules/ScheduleParameterBindings.test.tsx @@ -0,0 +1,68 @@ +import { fireEvent, render, screen } from "@testing-library/react"; +import { describe, expect, it, vi } from "vitest"; + +import { ScheduleParameterBindings } from "./ScheduleParameterBindings"; + +describe("ScheduleParameterBindings", () => { + it("offers dynamic date sources and SQL NULL only for optional parameters", () => { + render( + , + ); + + const dateSource = screen.getByLabelText("Value source for run_date"); + expect(dateSource).toHaveTextContent("Run date"); + expect(dateSource).toHaveTextContent("Relative to run date"); + expect(dateSource).toHaveTextContent("Window start"); + expect(dateSource).not.toHaveTextContent("SQL NULL"); + + expect(screen.getByLabelText("Value source for store_id")).toHaveTextContent("SQL NULL"); + }); + + it("emits a relative-date binding with an editable offset", () => { + const onBindingsChange = vi.fn(); + render( + , + ); + + fireEvent.change(screen.getByLabelText("Value source for run_date"), { + target: { value: "relative_date" }, + }); + expect(onBindingsChange).toHaveBeenLastCalledWith({ + run_date: { source: "relative_date", offset_days: -1 }, + }); + }); + + it("shows window presets when a parameter uses a window boundary", () => { + const onWindowChange = vi.fn(); + render( + , + ); + + expect(screen.getByLabelText("Date window preset")).toBeInTheDocument(); + fireEvent.change(screen.getByLabelText("Date window preset"), { + target: { value: "last_n_complete_days" }, + }); + expect(onWindowChange).toHaveBeenCalledWith({ + preset: "last_n_complete_days", + days: 7, + }); + }); +}); diff --git a/frontend/src/components/schedules/ScheduleParameterBindings.tsx b/frontend/src/components/schedules/ScheduleParameterBindings.tsx new file mode 100644 index 0000000..c73d98b --- /dev/null +++ b/frontend/src/components/schedules/ScheduleParameterBindings.tsx @@ -0,0 +1,226 @@ +import { Input } from "@/components/ui/input"; +import { Label } from "@/components/ui/label"; +import { Select } from "@/components/ui/select"; +import type { ParamDescriptor } from "@/types/endpoint"; +import type { + ScheduleBindingSource, + ScheduleParameterBinding, + ScheduleWindow, + ScheduleWindowPreset, +} from "@/types/schedule"; + +import { bindingsUseWindow } from "./scheduleBindings"; + +interface ScheduleParameterBindingsProps { + paramSchema: Record; + bindings: Record; + window?: ScheduleWindow | null; + onBindingsChange: (bindings: Record) => void; + onWindowChange: (window: ScheduleWindow | null) => void; +} + +function initialBinding(source: ScheduleBindingSource, descriptor: ParamDescriptor) { + if (source === "literal") { + if (descriptor.type === "boolean") return { source, value: true } as const; + return { source, value: descriptor.default ?? "" } as const; + } + if (source === "relative_date") return { source, offset_days: -1 } as const; + return { source } as ScheduleParameterBinding; +} + +function WindowEditor({ + window, + onChange, +}: { + window?: ScheduleWindow | null; + onChange: (window: ScheduleWindow) => void; +}) { + const current = window ?? { preset: "previous_day" as const }; + return ( +
+ + + {current.preset === "last_n_complete_days" && ( +
+ + + onChange({ + preset: "last_n_complete_days", + days: Math.max(1, Number.parseInt(event.target.value, 10) || 1), + }) + } + /> +
+ )} +

+ Window start and end are inclusive and are evaluated from each logical run date. +

+
+ ); +} + +export function ScheduleParameterBindings({ + paramSchema, + bindings, + window, + onBindingsChange, + onWindowChange, +}: ScheduleParameterBindingsProps) { + const updateBinding = (name: string, binding: ScheduleParameterBinding) => + onBindingsChange({ ...bindings, [name]: binding }); + + if (Object.keys(paramSchema).length === 0) { + return ( +
+ This endpoint has no SQL parameters. Each run will execute the same query. +
+ ); + } + + return ( +
+
+

Scheduled parameter values

+

+ Values are owned by this schedule and resolved from its logical run date. +

+
+ {Object.entries(paramSchema).map(([name, descriptor]) => { + const binding = bindings[name]; + return ( +
+
+ + +
+
+ {binding?.source === "literal" && descriptor.type === "boolean" && ( + <> + + + + )} + {binding?.source === "literal" && descriptor.type !== "boolean" && ( + <> + + { + const raw = event.target.value; + const value = + descriptor.type === "integer" || descriptor.type === "float" + ? raw === "" + ? "" + : Number(raw) + : raw; + updateBinding(name, { source: "literal", value }); + }} + /> + + )} + {binding?.source === "relative_date" && ( + <> + + + updateBinding(name, { + source: "relative_date", + offset_days: Number.parseInt(event.target.value, 10) || 0, + }) + } + /> + + )} + {binding && !["literal", "relative_date"].includes(binding.source) && ( +

+ {binding.source === "null" + ? "Oracle receives an explicit SQL NULL bind." + : "Resolved separately for every logical run."} +

+ )} +
+
+ ); + })} + {bindingsUseWindow(bindings) && } +
+ ); +} diff --git a/frontend/src/components/schedules/scheduleBindings.test.ts b/frontend/src/components/schedules/scheduleBindings.test.ts new file mode 100644 index 0000000..4d13e6d --- /dev/null +++ b/frontend/src/components/schedules/scheduleBindings.test.ts @@ -0,0 +1,56 @@ +import { describe, expect, it } from "vitest"; + +import type { ParamDescriptor } from "@/types/endpoint"; +import type { ScheduleParameterBinding } from "@/types/schedule"; + +import { scheduleBindingsComplete, suggestScheduleBindings } from "./scheduleBindings"; + +describe("schedule binding helpers", () => { + it("translates endpoint defaults into editable schedule-owned suggestions", () => { + const schema: Record = { + start_date: { type: "date", required: true, default_expression: "yesterday" }, + end_date: { type: "date", required: true, default_expression: "today" }, + store_id: { type: "string", required: false, default_is_null: true }, + limit: { type: "integer", required: true, default: 100 }, + region: { type: "string", required: true }, + }; + + expect(suggestScheduleBindings(schema)).toEqual({ + start_date: { source: "relative_date", offset_days: -1 }, + end_date: { source: "run_date" }, + store_id: { source: "null" }, + limit: { source: "literal", value: 100 }, + }); + }); + + it("requires one binding per endpoint parameter", () => { + const schema: Record = { + run_date: { type: "date", required: true }, + store_id: { type: "string", required: true }, + }; + const incomplete: Record = { + run_date: { source: "run_date" }, + }; + + expect(scheduleBindingsComplete(schema, incomplete, undefined)).toBe(false); + expect( + scheduleBindingsComplete( + schema, + { ...incomplete, store_id: { source: "literal", value: "101" } }, + undefined, + ), + ).toBe(true); + }); + + it("requires a window preset when a binding reads a window boundary", () => { + const schema: Record = { + start_date: { type: "date", required: true }, + }; + const bindings: Record = { + start_date: { source: "window_start" }, + }; + + expect(scheduleBindingsComplete(schema, bindings, undefined)).toBe(false); + expect(scheduleBindingsComplete(schema, bindings, { preset: "previous_day" })).toBe(true); + }); +}); diff --git a/frontend/src/components/schedules/scheduleBindings.ts b/frontend/src/components/schedules/scheduleBindings.ts new file mode 100644 index 0000000..1fed488 --- /dev/null +++ b/frontend/src/components/schedules/scheduleBindings.ts @@ -0,0 +1,49 @@ +import type { ParamDescriptor } from "@/types/endpoint"; +import type { ScheduleParameterBinding, ScheduleWindow } from "@/types/schedule"; + +export function suggestScheduleBindings( + paramSchema: Record, +): Record { + const bindings: Record = {}; + for (const [name, descriptor] of Object.entries(paramSchema)) { + if (descriptor.default_is_null) { + bindings[name] = { source: "null" }; + } else if (descriptor.default_expression === "today") { + bindings[name] = { source: "run_date" }; + } else if (descriptor.default_expression === "yesterday") { + bindings[name] = { source: "relative_date", offset_days: -1 }; + } else if (descriptor.default !== null && descriptor.default !== undefined) { + bindings[name] = { source: "literal", value: descriptor.default }; + } + } + return bindings; +} + +export function bindingsUseWindow(bindings: Record): boolean { + return Object.values(bindings).some( + (binding) => binding.source === "window_start" || binding.source === "window_end", + ); +} + +export function scheduleBindingsComplete( + paramSchema: Record, + bindings: Record, + window: ScheduleWindow | null | undefined, +): boolean { + const parameterNames = Object.keys(paramSchema); + if (parameterNames.some((name) => !bindings[name])) return false; + if (Object.keys(bindings).some((name) => !paramSchema[name])) return false; + if (bindingsUseWindow(bindings) && !window) return false; + + return parameterNames.every((name) => { + const binding = bindings[name]; + const descriptor = paramSchema[name]; + if (binding.source === "null") return !descriptor.required; + if (binding.source === "literal") { + if (binding.value === null || binding.value === undefined) return false; + if (descriptor.type !== "string" && binding.value === "") return false; + } + if (binding.source === "relative_date" && binding.offset_days == null) return false; + return true; + }); +} diff --git a/frontend/src/lib/api.ts b/frontend/src/lib/api.ts index 792c7dc..7796f21 100644 --- a/frontend/src/lib/api.ts +++ b/frontend/src/lib/api.ts @@ -26,6 +26,9 @@ import type { JobRun, Schedule, ScheduleCreate, + SchedulePreviewRequest, + SchedulePreviewResponse, + ScheduleRunRequest, ScheduleUpdate, SnapshotSummary, } from "@/types/schedule"; @@ -203,14 +206,19 @@ export const schedulesApi = { create: (payload: ScheduleCreate): Promise => http.post("/api/v1/admin/schedules/", payload).then((r) => r.data), + preview: (payload: SchedulePreviewRequest): Promise => + http + .post("/api/v1/admin/schedules/preview", payload) + .then((r) => r.data), + update: (id: string, payload: ScheduleUpdate): Promise => http.put(`/api/v1/admin/schedules/${id}`, payload).then((r) => r.data), delete: (id: string): Promise => http.delete(`/api/v1/admin/schedules/${id}`).then(() => undefined), - runNow: (id: string): Promise<{ status: string }> => - http.post<{ status: string }>(`/api/v1/admin/schedules/${id}/run`).then((r) => r.data), + runNow: (id: string, payload?: ScheduleRunRequest): Promise<{ status: string }> => + http.post<{ status: string }>(`/api/v1/admin/schedules/${id}/run`, payload).then((r) => r.data), pause: (id: string): Promise => http.post(`/api/v1/admin/schedules/${id}/pause`).then((r) => r.data), diff --git a/frontend/src/pages/SchedulesPage.tsx b/frontend/src/pages/SchedulesPage.tsx index 3bfe0a1..23805db 100644 --- a/frontend/src/pages/SchedulesPage.tsx +++ b/frontend/src/pages/SchedulesPage.tsx @@ -5,6 +5,7 @@ import { Pause, Play, Plus, RefreshCw, Trash2 } from "lucide-react"; import { Badge } from "@/components/ui/badge"; import { Button } from "@/components/ui/button"; import { CronScheduleBuilder } from "@/components/schedules/CronScheduleBuilder"; +import { ScheduleParameterBindings } from "@/components/schedules/ScheduleParameterBindings"; import { buildCronExpression, describeCronExpression, @@ -12,6 +13,11 @@ import { isValidCronExpression, type CronBuilderValue, } from "@/components/schedules/cronSchedule"; +import { + bindingsUseWindow, + scheduleBindingsComplete, + suggestScheduleBindings, +} from "@/components/schedules/scheduleBindings"; import { Dialog, DialogContent, @@ -26,7 +32,29 @@ import { Select } from "@/components/ui/select"; import { endpointsApi, getApiError, schedulesApi } from "@/lib/api"; import { queryKeys } from "@/lib/queryClient"; import type { Endpoint } from "@/types/endpoint"; -import type { JobRun, Schedule, ScheduleCreate } from "@/types/schedule"; +import type { JobRun, Schedule, ScheduleCreate, SchedulePreviewRequest } from "@/types/schedule"; + +const browserTimezone = Intl.DateTimeFormat().resolvedOptions().timeZone || "UTC"; +const commonTimezones = Array.from( + new Set([ + "UTC", + browserTimezone, + "Asia/Riyadh", + "Asia/Dubai", + "Europe/London", + "America/New_York", + ]), +); + +function initialCreateForm(): ScheduleCreate { + return { + endpoint_id: "", + schedule_type: "interval", + interval_seconds: 300, + timezone: browserTimezone, + parameter_bindings: {}, + }; +} function formatDate(iso: string | null): string { if (!iso) return "—"; @@ -46,13 +74,10 @@ export function SchedulesPage() { const [showCreate, setShowCreate] = useState(false); const [deleteSchedule, setDeleteSchedule] = useState(null); const [viewJobsFor, setViewJobsFor] = useState(null); - const [createForm, setCreateForm] = useState({ - endpoint_id: "", - schedule_type: "interval", - interval_seconds: 300, - }); + const [createForm, setCreateForm] = useState(initialCreateForm); const [cronBuilder, setCronBuilder] = useState(INITIAL_CRON_BUILDER); const [createError, setCreateError] = useState(""); + const [previewError, setPreviewError] = useState(""); const [deleteError, setDeleteError] = useState(""); const { @@ -83,8 +108,9 @@ export function SchedulesPage() { qc.invalidateQueries({ queryKey: queryKeys.schedules.all }); setShowCreate(false); setCreateError(""); - setCreateForm({ endpoint_id: "", schedule_type: "interval", interval_seconds: 300 }); + setCreateForm(initialCreateForm()); setCronBuilder(INITIAL_CRON_BUILDER); + setPreviewError(""); }, onError: (err) => setCreateError(getApiError(err)), }); @@ -99,6 +125,12 @@ export function SchedulesPage() { onError: (err) => setDeleteError(getApiError(err)), }); + const previewMutation = useMutation({ + mutationFn: (payload: SchedulePreviewRequest) => schedulesApi.preview(payload), + onSuccess: () => setPreviewError(""), + onError: (err) => setPreviewError(getApiError(err)), + }); + const runNowMutation = useMutation({ mutationFn: (id: string) => schedulesApi.runNow(id), onSuccess: () => { @@ -119,8 +151,28 @@ export function SchedulesPage() { // Filter endpoints that don't already have a schedule const scheduledEndpointIds = new Set(schedules.map((s) => s.endpoint_id)); const availableEndpoints = endpoints.filter( - (e) => e.is_active && !scheduledEndpointIds.has(e.id), + (e) => e.is_active && e.data_strategy === "snapshot" && !scheduledEndpointIds.has(e.id), ); + const selectedEndpoint = endpointMap.get(createForm.endpoint_id); + const parameterBindings = createForm.parameter_bindings ?? {}; + const bindingsComplete = selectedEndpoint + ? scheduleBindingsComplete(selectedEndpoint.param_schema, parameterBindings, createForm.window) + : false; + const timingComplete = + createForm.schedule_type === "interval" + ? (createForm.interval_seconds ?? 0) >= 10 + : isValidCronExpression(createForm.cron_expression ?? ""); + + const previewPayload = (): SchedulePreviewRequest => ({ + endpoint_id: createForm.endpoint_id, + schedule_type: createForm.schedule_type, + cron_expression: createForm.cron_expression, + interval_seconds: createForm.interval_seconds, + timezone: createForm.timezone ?? "UTC", + parameter_bindings: parameterBindings, + window: createForm.window, + count: 3, + }); return (
@@ -198,6 +250,7 @@ export function SchedulesPage() { {s.cron_expression} +

{s.timezone}

) : ( `Every ${s.interval_seconds}s` @@ -271,7 +324,7 @@ export function SchedulesPage() { {/* Create Dialog */} setShowCreate(false)}> - + Create Schedule @@ -289,7 +342,18 @@ export function SchedulesPage() { +

+ Only active snapshot endpoints without an existing schedule are shown. +

@@ -349,6 +416,97 @@ export function SchedulesPage() { }} /> )} +
+ + +

+ Cron times and logical dates are evaluated in this timezone. +

+
+ {selectedEndpoint && ( + { + setCreateForm((form) => ({ + ...form, + parameter_bindings: bindings, + window: bindingsUseWindow(bindings) + ? (form.window ?? { preset: "previous_day" }) + : undefined, + })); + previewMutation.reset(); + }} + onWindowChange={(window) => { + setCreateForm((form) => ({ ...form, window })); + previewMutation.reset(); + }} + /> + )} + {previewError && ( +
+ {previewError} +
+ )} + {selectedEndpoint && ( +
+
+
+

Resolved run preview

+

+ Verify the next three logical dates and bound values before creating. +

+
+ +
+ {previewMutation.data && ( +
+ + + + + + + + + + {previewMutation.data.runs.map((run) => ( + + + + + + ))} + +
Scheduled forLogical dateResolved parameters
{formatDate(run.scheduled_for)}{run.logical_date} + {JSON.stringify(run.resolved_parameters)} +
+
+ )} +
+ )}