diff --git a/.env.example b/.env.example index 3de88fa..1c69849 100644 --- a/.env.example +++ b/.env.example @@ -18,8 +18,9 @@ ENCRYPTION_KEY= # whose username and bcrypt password hash are supplied via env. There # is no users table; rotating the credential means redeploying with a # new ADMIN_PASSWORD_HASH. -# Generate the hash with: -# python -c "from app.auth.hashing import hash_password; print(hash_password('your-password'))" +# Generate the hash from backend/ with its environment activated. The prompt +# keeps the real password out of this file and your shell history: +# python -c "from getpass import getpass; from app.auth.hashing import hash_password; print(hash_password(getpass('Admin password: ')))" ADMIN_USERNAME= ADMIN_PASSWORD_HASH= diff --git a/.github/copilot-instructions.md b/.github/copilot-instructions.md index 95d4b64..ee21d54 100644 --- a/.github/copilot-instructions.md +++ b/.github/copilot-instructions.md @@ -19,7 +19,7 @@ These instructions apply repository-wide. Prefer local instructions under `.gith - Use PyJWT + bcrypt for auth. - Use structlog structured logging. - Use Vite + React SPA for frontend. -- Include rich SQL editor in wizard flows (Monaco or CodeMirror 6). +- Preserve the existing CodeMirror 6 SQL editor in endpoint wizard flows. ## Build and Test Commands - Backend setup: `cd backend && python -m venv .venv && . .venv/bin/activate && pip install -r requirements.txt`. @@ -33,6 +33,13 @@ These instructions apply repository-wide. Prefer local instructions under `.gith ## Mandatory Safety Rules - SQL must be parameterized with bind params only (`:param_name`). - Never generate SQL using string concatenation from user input. +- Treat bind markers inside single-quoted SQL literals as text, not parameters. +- Enforce required parameters for both live and snapshot HTTP requests, even when endpoint + defaults exist. +- Keep scheduler bindings independent from endpoint request defaults. Every scheduled bind must + use a validated declarative schedule source. +- Never serve a parameterized snapshot without complete cached-column mappings, coverage + validation, and typed row filtering. Snapshot filters do not provide tenant authorization. - Never store secrets in code, tests, fixtures, or docs. - Never break existing `/api/v1/*` contracts without version bump + migration notes. - Never change an applied Alembic revision; add a new revision. @@ -43,6 +50,8 @@ These instructions apply repository-wide. Prefer local instructions under `.gith - Include docs updates for API/config/workflow changes. - Keep changes minimal and scoped to request. - Explain risks when touching auth, SQL execution, migrations, or scheduler logic. +- Update `docs/scheduler_parameter_bindings.md` whenever the endpoint, schedule, or snapshot + parameter contract changes. ## Stop Conditions - If API contract is unclear, inspect existing `/api/v1` routers and docs. diff --git a/.github/instructions/backend.instructions.md b/.github/instructions/backend.instructions.md index f30bec5..f0432be 100644 --- a/.github/instructions/backend.instructions.md +++ b/.github/instructions/backend.instructions.md @@ -11,11 +11,19 @@ - Emit structlog events with required fields. - Validate SQL bind parameters via typed schemas before execution. - Use SQLAlchemy `text()` and binds for user SQL. +- Enforce required HTTP parameters with `build_param_model(..., enforce_required=True)` for both + live and snapshot data paths. +- Resolve scheduled SQL binds only from schedule-owned declarative bindings. +- Require complete snapshot filter mappings, prove retained-run coverage, and filter cached rows + with the declared typed `eq`/`gte`/`lte` operators. ## Do Not - Do not add unversioned API routes. - Do not use `python-jose` or `passlib`. - Do not interpolate SQL strings with user input. +- Do not treat quoted `':param'` text as a bind placeholder. +- Do not reuse endpoint defaults as implicit schedule inputs or serve an unfiltered parameterized + snapshot. - Do not edit old migration revisions after merge. - Do not return secrets in API responses or logs. @@ -36,6 +44,8 @@ - `endpoint` - `status` - `duration_ms` +- `method` +- `client_ip` - `event` ## Stop Conditions diff --git a/.github/instructions/docker.instructions.md b/.github/instructions/docker.instructions.md index 223ae37..8a99286 100644 --- a/.github/instructions/docker.instructions.md +++ b/.github/instructions/docker.instructions.md @@ -1,8 +1,8 @@ # Docker Instructions ## Do -- Keep Docker assets in `docker/` and `docker-compose.yml`. -- Ensure compose supports `api`, `web`, `db` services. +- Keep Docker assets in `docker/`, `docker-compose.yml`, and `compose.production.yml`. +- Ensure compose supports `db`, one-shot `migrate`, `api`, and `web` services. - Keep optional local Oracle service/profile isolated and documented. - Use deterministic base image tags. - Ensure backend and frontend images build in CI. @@ -21,8 +21,10 @@ - `docker compose logs --tail=200 db` ## Runtime Expectations -- Backend starts only after DB readiness. -- Migrations are documented and run deterministically. +- The `migrate` service starts after DB readiness and must complete successfully before the API + starts. +- Migrations are applied deterministically with `alembic upgrade head`; do not add an independent + API-startup migration path that can race the Compose service. - Config comes from environment variables compatible with Pydantic Settings. ## Stop Conditions diff --git a/.github/instructions/frontend.instructions.md b/.github/instructions/frontend.instructions.md index f77c514..0c73258 100644 --- a/.github/instructions/frontend.instructions.md +++ b/.github/instructions/frontend.instructions.md @@ -5,8 +5,11 @@ - Use React + TypeScript + Tailwind + shadcn/ui. - Keep API clients aligned to `/api/v1/admin/*` and `/api/v1/data/*`. - Implement wizard UX for Module 2 with a rich SQL editor. -- Use Monaco (`@monaco-editor/react`) or CodeMirror 6 (`@uiw/react-codemirror`) for SQL authoring. +- Use the existing CodeMirror 6 integration (`@uiw/react-codemirror`) for SQL authoring. - Provide explicit validation errors for bind params and auth setup. +- Keep preview samples separate from persisted endpoint defaults and schedule bindings. +- For snapshot endpoints, require a cached output column and `eq`/`gte`/`lte` operator for every + request parameter in create and edit flows. - Write component tests for wizard steps and critical forms. ## Do Not @@ -26,6 +29,8 @@ - Wizard must enforce bind variable awareness (`:param_name`). - Endpoint creation UI must expose auth assignment and data strategy selection. - Show clear status for live-query vs scheduled-snapshot behavior. +- Explain that schedule bindings control what Oracle loads, while snapshot filter mappings control + which cached rows an authenticated data request returns; neither replaces authentication. - Surface backend validation errors verbatim when safe. ## Stop Conditions diff --git a/.github/instructions/testing.instructions.md b/.github/instructions/testing.instructions.md index 650227b..b10393f 100644 --- a/.github/instructions/testing.instructions.md +++ b/.github/instructions/testing.instructions.md @@ -16,11 +16,18 @@ - SQL bind parameter validation and rejection of unsafe SQL composition. - Dynamic route resolution under `/api/v1/data/*`. - Scheduler job creation/execution logging and snapshot cache behavior. +- Required-parameter enforcement for both live and snapshot requests, including supported date + formats and optional SQL `NULL` behavior. +- Schedule binding completeness, logical-date/window resolution, retained-snapshot coverage, + reversed ranges, mapped-column failures, and typed cached-row filtering. - Alembic migration upgrade/downgrade sanity. ## Frontend Test Focus + - Wizard step transitions and validation. - SQL editor integration behavior and parameter UX. +- Separation of preview samples, endpoint defaults, schedule bindings, and snapshot filter + mappings in endpoint create/edit flows. - Connection/auth/schedule/settings form validation. - Error rendering for API failures. diff --git a/AGENTS.md b/AGENTS.md index 1911151..75dcfac 100644 --- a/AGENTS.md +++ b/AGENTS.md @@ -31,6 +31,39 @@ Build and maintain a secure, testable monorepo for dynamic SQL-to-API exposure w - Parameter mapping source: validated request inputs -> typed schema -> bind dict. - Reject queries containing interpolated values from raw strings. - Execute user SQL through SQLAlchemy Core `text()` with bound params. +- Bind markers inside single-quoted SQL literals are text, not parameters; never quote a bind + placeholder (use `column = :value`, not `column = ':value'`). + +## Parameter Ownership Contract + +- SQL preview values are temporary samples only; they are never persisted as endpoint or schedule + defaults. +- Live and snapshot HTTP requests enforce every descriptor marked `required`, even when the + endpoint descriptor contains a default. +- Optional live-request parameters may use a typed literal default, explicit SQL `NULL`, or the + supported dynamic date defaults `today` and `yesterday`. +- Scheduled snapshot execution never reads endpoint defaults. Every SQL bind is owned by the + schedule through exactly one validated binding source: `literal`, `null`, `run_date`, + `relative_date`, `window_start`, or `window_end`. +- Date inputs accept `YYYY-MM-DD` and `DD-MM-YYYY`. Schedule calendar math is evaluated from the + persisted nominal run time in the schedule's IANA timezone. + +## Snapshot Request Contract + +- Every parameterized snapshot endpoint must map every request parameter to a cached output + column and one operator: `eq`, `gte`, or `lte`. +- Mappings target the final cached column name after `column_map` renaming. They filter rows; they + are not tenant authorization. Authentication remains mandatory and independent. +- Select the newest retained snapshot whose persisted resolved schedule parameters cover the + request, then apply typed row filtering. Never return an unfiltered parameterized snapshot. +- Missing/invalid required parameters, reversed ranges, incomplete mappings, unavailable mapped + columns, and out-of-coverage requests return explicit HTTP 422 responses. No retained snapshot + returns HTTP 503; an in-coverage request with no matching rows returns HTTP 200 with `data: []`. +- `null_means_all` is valid only for an optional `eq` mapping and means a scheduled SQL `NULL` + covers every requested value for that parameter. +- Snapshot mappings are stored in the existing endpoint parameter JSON and do not require a + relational migration. Schedule-owned bindings and logical-run audit fields are relational and + are covered by Alembic revision `e4a6c2d9f801`. ## Required Log Fields - `request_id` @@ -52,6 +85,8 @@ Build and maintain a secure, testable monorepo for dynamic SQL-to-API exposure w - Run formatting/lint/tests for changed areas. - Add/update Alembic migration when schema changed. - Update docs when contracts, settings, or workflows changed. +- Keep `README.md`, `docs/architecture.md`, `docs/scheduler_parameter_bindings.md`, and relevant + agent instructions aligned when parameter or snapshot behavior changes. - Verify no secrets/tokens/credentials are committed. ## Stop Conditions diff --git a/AI_WORKFLOW.md b/AI_WORKFLOW.md index 44feab5..9fd40ea 100644 --- a/AI_WORKFLOW.md +++ b/AI_WORKFLOW.md @@ -18,6 +18,9 @@ 4. If API changes: preserve `/api/v1/*` compatibility or add version bump. 5. After editing: run lint/tests/build for touched areas. 6. Update docs/changelog notes when behavior or contract changes. +7. For parameter behavior, keep preview samples, endpoint request defaults, schedule bindings, + and snapshot request filters separate; verify the canonical contract in + `docs/scheduler_parameter_bindings.md`. ## Verification Checklist - Backend checks pass: `ruff`, `mypy`, `pytest`. @@ -26,6 +29,8 @@ - Migrations included for DB changes. - No secrets in code, logs, tests, or docs. - No breaking API change without explicit versioning plan. +- Required live/snapshot parameters remain enforced and parameterized snapshots cannot fall back + to unfiltered cached data. ## PR Checklist - Scope is clear and minimal. diff --git a/CLAUDE.md b/CLAUDE.md index bb0de63..77c964f 100644 --- a/CLAUDE.md +++ b/CLAUDE.md @@ -7,11 +7,12 @@ QueryGateway is a monorepo for creating secure, dynamic REST endpoints from Orac - Python 3.14+, FastAPI, Pydantic Settings v2. - PostgreSQL app DB, SQLAlchemy 2.0, Alembic. - Oracle connectivity via `python-oracledb`. -- APScheduler 3.x with persistent PostgreSQL job store. +- APScheduler 3.x with an in-memory job store; schedule definitions are persisted in PostgreSQL + and active jobs are restored on API startup. - Auth: `PyJWT` + `bcrypt`. - Logging: `logging` + `structlog`. - Frontend: Vite + React SPA + TypeScript + shadcn/ui + Tailwind. -- Module 2 requires rich SQL editor: Monaco (`@monaco-editor/react`) or CodeMirror 6 (`@uiw/react-codemirror`). +- The endpoint wizard uses CodeMirror 6 through `@uiw/react-codemirror` for SQL authoring. ## Repo Map - `backend/`: API, models, migrations, scheduler, auth, SQL execution. @@ -45,21 +46,43 @@ QueryGateway is a monorepo for creating secure, dynamic REST endpoints from Orac - Validate and coerce params with typed schemas before execution. - Use SQLAlchemy `text()` + bind dict. - Never concatenate request values into SQL strings. +- Bind markers inside single-quoted SQL literals are not parameters and must not be used. + +## Endpoint and Scheduler Parameter Conventions + +- Preview inputs are temporary samples and are not persisted. +- Every required HTTP parameter remains required for live and snapshot requests, regardless of + configured endpoint defaults. +- Endpoint defaults apply to omitted optional live requests only. Schedules own every SQL bind + through `literal`, `null`, `run_date`, `relative_date`, `window_start`, or `window_end`. +- Parameterized snapshot endpoints map every request parameter to a post-rename cached column via + `eq`, `gte`, or `lte`. Coverage is checked against persisted job-run parameters before cached + rows are filtered. +- Snapshot filter mappings select data; they never replace endpoint authentication or authorize a + tenant. +- The full contract and stable error codes are documented in + `docs/scheduler_parameter_bindings.md`. ## Logging Conventions + - Use structured logging everywhere. -- Mandatory fields: `request_id`, `user`, `endpoint`, `status`, `duration_ms`, `event`. +- Mandatory fields: `request_id`, `user`, `endpoint`, `status`, `duration_ms`, `method`, + `client_ip`, `event`. - Add scheduler fields for jobs: `job_id`, `run_id`, `row_count`, `success`. ## Run Commands (Expected) -- Backend setup: `cd backend && python -m venv .venv && . .venv/bin/activate` (Windows: `.venv\Scripts\activate`) then `pip install -r requirements.txt`. + +- Backend setup: `cd backend && python3.14 -m venv .venv && . .venv/bin/activate` + (Windows: `py -3.14 -m venv .venv`, then `.venv\Scripts\activate`) and install + `requirements.txt`. - Backend dev: `cd backend && uvicorn app.main:app --reload`. - Backend checks: `cd backend && ruff check . && mypy . && pytest`. - Frontend setup: `cd frontend && npm install`. - Frontend dev: `cd frontend && npm run dev`. - Frontend checks: `cd frontend && npm run eslint && npm run prettier:check && npm run test`. - Docker build: `docker compose build`. -- Docker run: `docker compose up -d`. +- Docker run: `docker compose up -d --build` (the one-shot `migrate` service must complete before + the API starts). ## Working Protocol - Before edits: summarize target files and planned commands. diff --git a/README.md b/README.md index 09e4f3e..4e31866 100644 --- a/README.md +++ b/README.md @@ -1,6 +1,6 @@ # QueryGateway -> A self-hosted platform that turns Oracle SQL queries into secure, versioned REST API endpoints — no application code required. Author a parameterized query in a guided wizard, attach authentication, choose a data-freshness strategy, and publish a live endpoint that other systems can consume. +> A self-hosted platform that turns Oracle SQL queries into secure, versioned REST API endpoints — no application code required. Author a parameterized query in a guided wizard, attach authentication, choose a data-freshness strategy, and publish a live or snapshot-backed endpoint that other systems can consume. [![License: GPL v3](https://img.shields.io/badge/License-GPLv3-blue.svg)](LICENSE.txt) [![Python](https://img.shields.io/badge/python-3.14%2B-blue.svg)](https://www.python.org/) @@ -61,9 +61,9 @@ Define connection ─▶ Author SQL (with :bind params) ─▶ Attach auth ─ ``` 1. **Connect** to your Oracle database with securely stored, encrypted credentials. -2. **Author** a `SELECT` query using named bind parameters (`:param_name`) in a rich SQL editor. The wizard detects parameters as you type, requests temporary sample values before previewing, and supports optional endpoint defaults for live requests. Required live parameters remain mandatory; callers can supply dates as either `YYYY-MM-DD` or `DD-MM-YYYY`. +2. **Author** a `SELECT` query using unquoted named bind parameters (`:param_name`) in a rich SQL editor. The wizard detects parameters as you type and requests temporary sample values before previewing; preview samples are never persisted. Optional live-request parameters can have typed defaults, explicit SQL `NULL`, or `today`/`yesterday` date defaults. Required parameters remain mandatory for both live and snapshot requests even when a default exists, and dates accept `YYYY-MM-DD` or `DD-MM-YYYY`. 3. **Secure** the endpoint by attaching a Bearer token, Basic Auth, or API key policy. -4. **Choose** a data strategy: serve results **live** on each request, or from a **scheduled snapshot** cache. Snapshot schedules own their parameter values: bind a fixed value, explicit SQL `NULL`, logical run date, relative date, or a reusable calendar window, then preview the next three resolved runs before creating the schedule. +4. **Choose** a data strategy: serve results **live** on each request, or from a **scheduled snapshot** cache. Snapshot endpoints explicitly map each request parameter to a cached output column and comparison. Schedules independently own the SQL values that define cache coverage: bind a fixed value, explicit SQL `NULL`, logical run date, relative date, or a reusable calendar window, then preview the next three resolved runs before creating the schedule. 5. **Publish** a versioned endpoint under `/api/v1/data/*` that resolves dynamically — no service restart needed. ## Features @@ -73,9 +73,9 @@ QueryGateway is organized into five admin modules, all driven from the React adm | Module | What it does | |--------|--------------| | **Connections** | Create, edit, test, and delete Oracle database connections. Credentials are encrypted at rest; pool sizing and timeouts are configurable. Uses `python-oracledb`; thin mode needs no native client, while the backend Docker image includes Oracle Instant Client 19.32 for thick mode. | -| **API Creation Wizard** | A multi-step wizard that turns a parameterized SQL query into a deployable GET endpoint: pick a connection, author SQL with a rich editor, supply preview-only sample values for detected bind parameters, configure fixed values, explicit SQL `NULL` values for optional binds, or dynamic date defaults, preview sample rows and inferred schema, map/rename output columns, attach an auth method, and select a data strategy. | +| **API Creation Wizard** | A multi-step wizard that turns a parameterized SQL query into a deployable GET endpoint: pick a connection, author SQL with a rich editor, supply preview-only sample values, configure optional live-request defaults, preview rows and inferred schema, map/rename output columns, attach an auth method, select a data strategy, and map every snapshot request parameter to a cached output column and comparison. | | **Authentication** | Manage per-endpoint auth methods — Bearer token (JWT), Basic Auth, and API key. Tokens are issued/verified with `PyJWT`; credentials are hashed with `bcrypt`. Every `/api/v1/data/*` request is authenticated; endpoints without a dedicated method require the platform admin Bearer token. | -| **Scheduling & Snapshots** | Schedule query refreshes with friendly hourly, daily, weekly, or monthly calendar controls, an advanced custom-cron option, or a fixed interval. Each schedule has an IANA timezone and explicit per-parameter sources: fixed value, SQL `NULL`, logical run date, relative date, or inclusive calendar-window boundaries such as previous day, last N complete days, week/month to date, previous week, and previous month. Preview the next three resolved runs, run now, pause/resume, and inspect persisted logical-date, window, and resolved-bind audit context. Results are cached as PostgreSQL JSONB snapshots. Active APScheduler jobs are restored on API startup; deletion preserves job-run history. | +| **Scheduling & Snapshots** | Schedule query refreshes with friendly hourly, daily, weekly, or monthly calendar controls, an advanced custom-cron option, or a fixed interval. Each schedule has an IANA timezone and explicit per-parameter sources: fixed value, SQL `NULL`, logical run date, relative date, or inclusive calendar-window boundaries such as previous day, last N complete days, week/month to date, previous week, and previous month. Authenticated requests select the newest retained snapshot that covers their typed parameter values, then apply explicit equality or inclusive range filters to cached rows. Preview the next three resolved runs, run now, pause/resume, and inspect persisted logical-date, window, and resolved-bind audit context. Results are cached as PostgreSQL JSONB snapshots. Active APScheduler jobs are restored on API startup; deletion preserves job-run history. | | **Settings & Health** | Configure runtime settings (base URL/port, logging level, query timeouts, CORS/rate-limit inputs) and view a health dashboard covering API, PostgreSQL, Oracle connectivity, scheduler status, and recent job outcomes. | ### Security by Default @@ -83,7 +83,7 @@ QueryGateway is organized into five admin modules, all driven from the React adm - **SQL injection resistant** — user-defined SQL runs only through SQLAlchemy `text()` with named bind parameters. Request values are never concatenated into SQL strings, and bind values are validated through typed schemas before execution. - **Encrypted credentials** — Oracle connection secrets are encrypted at rest using an environment-provided key. - **Mandatory data authentication** — attach a Bearer token, Basic Auth, or API key policy to an endpoint. If no dedicated method is attached, the endpoint requires the platform admin Bearer token; anonymous data access is never allowed. -- **Structured, redacted logging** — `structlog` emits JSON logs with correlation fields (`request_id`, `user`, `endpoint`, `status`, `duration_ms`); credentials and tokens are redacted before emission. +- **Structured, redacted logging** — `structlog` emits JSON logs with correlation fields (`request_id`, `user`, `endpoint`, `status`, `duration_ms`, `method`, `client_ip`, `event`); credentials and tokens are redacted before emission. ### Two API Surfaces @@ -112,7 +112,7 @@ make setup make check # Start with Docker -cp .env.example .env # Set JWT_SECRET_KEY and ENCRYPTION_KEY (see deployment.md for generation commands) +cp .env.example .env # Set all required secrets and admin credentials (see deployment.md) make docker-up ``` @@ -149,7 +149,7 @@ py -3.14 -m venv .venv pip install --upgrade pip pip install -r requirements.txt Copy-Item .env.example .env -# edit .env — set ENCRYPTION_KEY and JWT_SECRET_KEY before continuing +# edit .env — set every required secret and admin credential before continuing # Step 3 — Run database migrations alembic upgrade head @@ -181,7 +181,7 @@ source .venv/bin/activate pip install --upgrade pip pip install -r requirements.txt cp .env.example .env -# edit .env — set ENCRYPTION_KEY and JWT_SECRET_KEY before continuing +# edit .env — set every required secret and admin credential before continuing # Step 3 — Run database migrations alembic upgrade head @@ -248,7 +248,7 @@ QueryGateway currently ships the five core admin modules (Connections, API Creat - **Advanced access control** — finer-grained RBAC beyond the baseline admin role. - **GraphQL** — an alternative query surface alongside REST. -These are not commitments or a release schedule — they're open directions. If one interests you, open an issue to discuss before starting work. See [`project_plan.md`](project_plan.md) for the full scope and non-goals. +These are not commitments or a release schedule — they're open directions. If one interests you, open an issue to discuss before starting work. See the [implementation plan](docs/project_plan.md) for the full scope and non-goals. ## Contributing & Community @@ -262,7 +262,8 @@ Contributions are welcome — whether that's code, docs, bug reports, or ideas. ## Documentation - [Architecture](docs/architecture.md) — system components, design decisions, directory layout -- [Implementation Plan](project_plan.md) — modules, scope, and phased delivery +- [Endpoint, scheduler, and snapshot parameter contracts](docs/scheduler_parameter_bindings.md) — request defaults, schedule bindings, coverage, filtering, and error semantics +- [Implementation Plan](docs/project_plan.md) — modules, scope, and phased delivery - [Deployment](docs/deployment.md) — self-hosted setup and secret generation - [Operations](docs/operations.md) — backup/restore, monitoring, troubleshooting - [Security Checklist](docs/security_checklist.md) — security validation controls diff --git a/SECURITY.md b/SECURITY.md index c24191d..cc95eca 100644 --- a/SECURITY.md +++ b/SECURITY.md @@ -67,19 +67,14 @@ item-by-item list) plus the standards every change is expected to uphold. ### 3.1 Authentication & Authorization - Authentication on data endpoints (`/api/v1/data/*`) is configured **per - endpoint**. An endpoint is protected only when an auth method is attached to - it. Serving an endpoint **without** an auth method (public/unauthenticated) - is allowed but must be a **deliberate, explicit choice**: the admin API - rejects a create or update that would leave an endpoint with no auth method - unless `allow_unauthenticated` is set to `true` (`422` otherwise). The data - plane enforces the same invariant at request time — an endpoint with no auth - method **and** no opt-in (e.g. after its referenced auth method is deleted, - which nulls the reference) **default-denies with `401`** instead of serving — - so an endpoint can never become silently public by omission or out-of-band. - Each request to a genuinely public (opted-in) endpoint is logged at `WARNING` - with `event="public_endpoint_served"`; the default-deny path logs - `event="unauthenticated_endpoint_denied"`. Review public endpoints - periodically for unintended exposure. + endpoint** and anonymous data access is never allowed. When an endpoint auth + method is attached, the data plane enforces that Bearer, Basic, or API-key + policy. When no dedicated method is attached, the legacy-named + `allow_unauthenticated=true` value explicitly selects the platform-admin + Bearer fallback; it does **not** make the endpoint public. A configuration + with neither path is rejected by the admin API and default-denied with `401` + by the data plane. The denied path logs + `event="unauthenticated_endpoint_denied"` with the required request context. - When an auth method **is** attached but is missing or inactive at request time, the endpoint **default-denies** with `401` (it never silently falls open). Expired or malformed credentials likewise return `401`, never `500`. @@ -132,12 +127,21 @@ item-by-item list) plus the standards every change is expected to uphold. stack traces, Oracle connection strings, and internal schema details are never leaked to clients. - Enforce string-length limits and enum validation on all user-supplied fields. +- Enforce every endpoint parameter marked `required` for both live and + snapshot HTTP requests, even if the endpoint stores a default. Preview + samples are temporary, optional live defaults are request behavior, and + schedule bindings are a separate persisted contract. +- Parameterized snapshot endpoints must map every request parameter to a final + cached output column and an allowlisted typed comparison. The data plane + proves coverage from persisted job-run parameters before filtering rows; it + never falls back to returning an unfiltered snapshot. Snapshot mappings are + data selection, not tenant authorization. ### 3.5 Logging, Auditing & Privacy - Use **structured logging** everywhere. Mandatory fields: `request_id`, - `user`, `endpoint`, `status`, `duration_ms`, `event`; scheduler jobs add - `job_id`, `run_id`, `row_count`, `success`. + `user`, `endpoint`, `status`, `duration_ms`, `method`, `client_ip`, `event`; + scheduler jobs add `job_id`, `run_id`, `row_count`, `success`. - Redact secrets and high-risk fields at logging boundaries; log the minimum PII required for the operation. - Every data-endpoint request is access-logged with path, method, principal, @@ -147,8 +151,12 @@ item-by-item list) plus the standards every change is expected to uphold. - Scheduled jobs use stored, encrypted connection credentials — never request-time credentials. -- Job execution errors are recorded internally and never surfaced to data-API - consumers (the data endpoint returns `503`). +- Scheduled jobs resolve every SQL bind from schedule-owned declarative + bindings (`literal`, `null`, logical/relative run date, or window boundary). + They never inherit endpoint request defaults. +- Job execution errors are recorded internally and never expose Oracle details to data-API + consumers. Snapshot requests use another retained covering run when available and return `503` + when no retained snapshot exists. - Concurrency is bounded (`max_instances=1` per job, `max_job_concurrency`) and snapshot retention is capped to prevent resource and storage exhaustion. diff --git a/SECURITY_AI.md b/SECURITY_AI.md index fa23c59..6bcd304 100644 --- a/SECURITY_AI.md +++ b/SECURITY_AI.md @@ -4,6 +4,8 @@ - Never commit secrets, API keys, JWT signing keys, DB credentials, or private certs. - Never log plaintext credentials or token bodies. - Never bypass authentication for `/api/v1/data/*`. +- The legacy `allow_unauthenticated` field only opts into platform-admin Bearer authentication; + it never authorizes anonymous data access. - Never allow raw SQL string interpolation with user input. - Never weaken token expiry validation. @@ -17,11 +19,19 @@ - Enforce `:param_name` bind style. - Validate and coerce bind values through typed schemas. - Reject unsafe query composition patterns. +- Never quote a bind placeholder inside a SQL string literal. +- Enforce required parameters on both live and snapshot requests, regardless of endpoint defaults. +- Scheduled execution must resolve every SQL bind from the schedule's own declarative bindings; + never read endpoint request defaults at run time. +- Parameterized snapshots must prove coverage from persisted job-run parameters and apply typed, + allowlisted cached-column filters before returning rows. Snapshot filters are not authorization. - Prefer least-privilege Oracle credentials for query execution. ## Logging and Privacy + - Use structured logs with minimal necessary PII. -- Required fields: `request_id`, `user`, `endpoint`, `status`, `duration_ms`, `event`. +- Required fields: `request_id`, `user`, `endpoint`, `status`, `duration_ms`, `method`, + `client_ip`, `event`. - Redact secrets and high-risk fields at logging boundaries. ## Migration and Release Safety diff --git a/backend/.env.example b/backend/.env.example index 9b16489..c72c9fa 100644 --- a/backend/.env.example +++ b/backend/.env.example @@ -30,8 +30,9 @@ JWT_ACCESS_TOKEN_EXPIRE_MINUTES=60 # Seeded admin — REQUIRED. App fails to start if either is unset or # empty. There is no users table; rotation means redeploying with a # new ADMIN_PASSWORD_HASH. -# Generate the hash with: -# python -c "from app.auth.hashing import hash_password; print(hash_password('your-password'))" +# Generate the hash with an interactive prompt so the real password does not +# appear in this file or your shell history: +# python -c "from getpass import getpass; from app.auth.hashing import hash_password; print(hash_password(getpass('Admin password: ')))" ADMIN_USERNAME= ADMIN_PASSWORD_HASH= diff --git a/backend/app/models/snapshot.py b/backend/app/models/snapshot.py index 705fc16..55aa262 100644 --- a/backend/app/models/snapshot.py +++ b/backend/app/models/snapshot.py @@ -13,8 +13,9 @@ class Snapshot(UUIDPrimaryKeyMixin, Base): """Immutable cache of a single scheduler query execution result. - snapshot-strategy endpoints read from the latest Snapshot for their - endpoint_id, falling back to the previous one on job failure. + Parameterized snapshot endpoints select the newest retained snapshot whose + associated job-run inputs cover the request, then apply typed row filters. + Unparameterized endpoints use the newest retained snapshot. """ __tablename__ = "snapshots" @@ -30,7 +31,7 @@ class Snapshot(UUIDPrimaryKeyMixin, Base): ) # Full query result stored as a JSON array of row objects. - data: Mapped[dict[str, object]] = mapped_column(JSONB, nullable=False) + data: Mapped[list[dict[str, object]]] = mapped_column(JSONB, nullable=False) row_count: Mapped[int] = mapped_column(Integer, nullable=False) created_at: Mapped[datetime] = mapped_column( diff --git a/backend/app/repositories/job_run.py b/backend/app/repositories/job_run.py index 4123c8b..ec913e0 100644 --- a/backend/app/repositories/job_run.py +++ b/backend/app/repositories/job_run.py @@ -32,6 +32,13 @@ async def get_by_schedule_and_scheduled_for( ) return result.scalar_one_or_none() + async def get_by_ids(self, job_run_ids: Sequence[uuid.UUID]) -> Sequence[JobRun]: + """Return all job runs matching the supplied IDs in one query.""" + if not job_run_ids: + return [] + result = await self._db.execute(select(JobRun).where(JobRun.id.in_(job_run_ids))) + return result.scalars().all() + async def get_all( self, *, diff --git a/backend/app/schemas/endpoint.py b/backend/app/schemas/endpoint.py index 0df0cb6..8d839ce 100644 --- a/backend/app/schemas/endpoint.py +++ b/backend/app/schemas/endpoint.py @@ -1,7 +1,7 @@ """Pydantic schemas for API endpoint management (Phase 4). Public contract rules: -- ``sql_text`` must use named bind parameters (``:`param_name``). +- ``sql_text`` must use named bind parameters (``:param_name``). - ``path`` must be a valid URL segment (no leading slash, no whitespace). - ``param_schema_json`` maps parameter names to type/required/default descriptors. - ``column_map_json`` maps source column names to output names (optional). @@ -12,7 +12,7 @@ from datetime import datetime from typing import Literal, Self -from pydantic import BaseModel, Field, field_validator, model_validator +from pydantic import BaseModel, Field, ValidationError, field_validator, model_validator from app.models.endpoint import DataStrategy @@ -68,6 +68,34 @@ def validate_sql_safety(sql: str) -> list[str]: return errors +class SnapshotFilter(BaseModel): + """Explicit mapping from one request parameter to one cached output column.""" + + column: str = Field(..., min_length=1, max_length=255) + operator: Literal["eq", "gte", "lte"] + null_means_all: bool = Field( + False, + description=( + "For equality filters only, a scheduled SQL NULL means the snapshot covers all " + "values for this parameter." + ), + ) + + @field_validator("column") + @classmethod + def normalize_column(cls, value: str) -> str: + normalized = value.strip() + if not normalized: + raise ValueError("Snapshot filter column cannot be blank.") + return normalized + + @model_validator(mode="after") + def validate_null_coverage(self) -> Self: + if self.null_means_all and self.operator != "eq": + raise ValueError("null_means_all is supported only for the eq operator.") + return self + + class ParamDescriptor(BaseModel): """Schema for a single bind parameter.""" @@ -98,9 +126,17 @@ class ParamDescriptor(BaseModel): ge=1, description="Maximum allowed length for string parameters.", ) + snapshot_filter: SnapshotFilter | None = Field( + None, + description=( + "Required for snapshot endpoints. Maps the request value to a cached output column." + ), + ) @model_validator(mode="after") def optional_must_have_default(self) -> Self: + if self.snapshot_filter and self.snapshot_filter.null_means_all and self.required: + raise ValueError("null_means_all is supported only for optional parameters.") configured_defaults = sum( ( self.default is not None, @@ -121,12 +157,58 @@ def optional_must_have_default(self) -> Self: model = build_param_model({"value": self.model_dump()}) model.model_validate({}) except (TypeError, ValueError) as exc: - raise ValueError( - f"Invalid default for parameter type '{self.type}'." - ) from exc + raise ValueError(f"Invalid default for parameter type '{self.type}'.") from exc return self +def require_snapshot_filter_mappings( + data_strategy: DataStrategy, + param_schema: dict[str, ParamDescriptor] | dict[str, object], +) -> None: + """Require complete mappings and compatible paired range-filter types.""" + if data_strategy != DataStrategy.snapshot or not param_schema: + return + + missing: list[str] = [] + parsed_descriptors: dict[str, ParamDescriptor] = {} + for name, descriptor in param_schema.items(): + try: + parsed = ( + descriptor + if isinstance(descriptor, ParamDescriptor) + else ParamDescriptor.model_validate(descriptor) + ) + except ValidationError as exc: + raise SnapshotConfigurationError("Endpoint has an invalid parameter schema.") from exc + parsed_descriptors[name] = parsed + if parsed.snapshot_filter is None: + missing.append(name) + + if missing: + names = ", ".join(f":{name}" for name in sorted(missing)) + raise SnapshotConfigurationError( + f"Snapshot endpoints require explicit snapshot filter mappings for: {names}." + ) + + range_types: dict[str, dict[str, set[str]]] = {} + for descriptor in parsed_descriptors.values(): + mapping = descriptor.snapshot_filter + if mapping is None or mapping.operator not in {"gte", "lte"}: + continue + range_types.setdefault(mapping.column, {}).setdefault(mapping.operator, set()).add( + descriptor.type + ) + + for column, bounds in range_types.items(): + if "gte" not in bounds or "lte" not in bounds: + continue + if len(bounds["gte"] | bounds["lte"]) > 1: + raise SnapshotConfigurationError( + "Snapshot lower and upper bounds for " + f"'{column}' must use the same declared parameter type." + ) + + class EndpointCreate(BaseModel): """Payload for POST /api/v1/admin/endpoints.""" @@ -197,6 +279,7 @@ def bind_params_match_schema(self) -> Self: 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_filter_mappings(self.data_strategy, self.param_schema) return self diff --git a/backend/app/services/data.py b/backend/app/services/data.py index bbe70a4..efd0dd0 100644 --- a/backend/app/services/data.py +++ b/backend/app/services/data.py @@ -36,8 +36,17 @@ from app.repositories.auth_method import AuthMethodRepository from app.repositories.connection import ConnectionRepository from app.repositories.endpoint import EndpointRepository +from app.repositories.job_run import JobRunRepository from app.repositories.snapshot import SnapshotRepository +from app.schemas.endpoint import SnapshotConfigurationError, require_snapshot_filter_mappings from app.services.auth_method import AuthMethodService +from app.services.snapshot_filtering import ( + compile_snapshot_filters, + filter_snapshot_rows, + snapshot_covers_request, + unavailable_snapshot_filter_columns, + validate_snapshot_parameter_ranges, +) from app.sql.executor import SqlExecutionError, execute_query from app.sql.param_models import build_param_model @@ -122,7 +131,13 @@ async def serve(self, path: str, request: Request) -> DataServiceResult: principal = await self._enforce_platform_auth(request, endpoint, started_at, path) if endpoint.data_strategy.value == "snapshot": - response = await self._serve_snapshot(endpoint, path, principal) + response = await self._serve_snapshot( + endpoint, + request, + path, + principal, + started_at, + ) else: response = await self._serve_live(endpoint, request, path, principal) @@ -189,9 +204,7 @@ async def _resolve_endpoint(self, path: str) -> ApiEndpoint: # ── Auth (per-endpoint) ───────────────────────────────────────────────── - async def _enforce_auth( - self, request: Request, auth_method_id: uuid.UUID - ) -> str: + async def _enforce_auth(self, request: Request, auth_method_id: uuid.UUID) -> str: svc = AuthMethodService(AuthMethodRepository(self._db)) auth_method = await svc.get_auth_method(auth_method_id) if auth_method is None or not auth_method.is_active: @@ -270,27 +283,143 @@ async def _enforce_auth( # ── Snapshot mode ─────────────────────────────────────────────────────── async def _serve_snapshot( - self, endpoint: ApiEndpoint, path: str, principal: str | None + self, + endpoint: ApiEndpoint, + request: Request, + path: str, + principal: str | None, + started_at: float, ) -> JSONResponse: - snapshot = await SnapshotRepository(self._db).get_latest_by_endpoint(endpoint.id) - if snapshot is None: + try: + params = self._coerce_params(endpoint, request) + except ValidationError as exc: + return self._parameter_validation_response(exc) + + param_schema = endpoint.param_schema_json or {} + try: + require_snapshot_filter_mappings(endpoint.data_strategy, param_schema) + except SnapshotConfigurationError: + self._log_snapshot_rejection( + request=request, + endpoint=endpoint, + path=path, + principal=principal, + started_at=started_at, + event="snapshot_filter_not_configured", + ) + return JSONResponse( + status_code=status.HTTP_422_UNPROCESSABLE_ENTITY, + content={ + "code": "snapshot_filter_not_configured", + "detail": ("Snapshot request filtering is not configured for every parameter."), + }, + ) + filters = compile_snapshot_filters(param_schema) + try: + validate_snapshot_parameter_ranges( + filters=filters, + request_params=params, + ) + except ValueError: + self._log_snapshot_rejection( + request=request, + endpoint=endpoint, + path=path, + principal=principal, + started_at=started_at, + event="invalid_snapshot_parameter_range", + ) + return JSONResponse( + status_code=status.HTTP_422_UNPROCESSABLE_ENTITY, + content={ + "code": "invalid_parameter_range", + "detail": "Snapshot filter lower bound must not exceed its upper bound.", + }, + ) + + snapshots = await SnapshotRepository(self._db).get_by_endpoint(endpoint.id, limit=100) + if not snapshots: return JSONResponse( status_code=status.HTTP_503_SERVICE_UNAVAILABLE, + content={"detail": "No snapshot available yet. Wait for the scheduled job to run."}, + ) + + snapshot = None + if not param_schema: + snapshot = snapshots[0] + else: + candidate_run_ids = list( + dict.fromkeys( + candidate.job_run_id + for candidate in snapshots + if candidate.job_run_id is not None + ) + ) + job_runs = await JobRunRepository(self._db).get_by_ids(candidate_run_ids) + job_runs_by_id = {job_run.id: job_run for job_run in job_runs} + for candidate in snapshots: + if candidate.job_run_id is None: + continue + job_run = job_runs_by_id.get(candidate.job_run_id) + if job_run is not None and snapshot_covers_request( + filters=filters, + request_params=params, + resolved_params=job_run.resolved_params_json or {}, + ): + snapshot = candidate + break + + if snapshot is None: + self._log_snapshot_rejection( + request=request, + endpoint=endpoint, + path=path, + principal=principal, + started_at=started_at, + event="snapshot_out_of_coverage", + ) + return JSONResponse( + status_code=status.HTTP_422_UNPROCESSABLE_ENTITY, content={ - "detail": "No snapshot available yet. " - "Wait for the scheduled job to run." + "code": "snapshot_out_of_coverage", + "detail": "Requested parameters are outside retained snapshot coverage.", }, ) snapshot_data: list[dict[str, object]] = ( snapshot.data if isinstance(snapshot.data, list) else [] ) + if unavailable_snapshot_filter_columns( + rows=snapshot_data, + filters=filters, + ): + self._log_snapshot_rejection( + request=request, + endpoint=endpoint, + path=path, + principal=principal, + started_at=started_at, + event="snapshot_filter_column_unavailable", + ) + return JSONResponse( + status_code=status.HTTP_422_UNPROCESSABLE_ENTITY, + content={ + "code": "snapshot_filter_column_unavailable", + "detail": "A configured snapshot filter column is unavailable.", + }, + ) + filtered_data = filter_snapshot_rows( + rows=snapshot_data, + filters=filters, + request_params=params, + ) return JSONResponse( status_code=200, content={ - "data": snapshot_data, + "data": filtered_data, "meta": { - "row_count": snapshot.row_count, + "row_count": len(filtered_data), + "snapshot_row_count": snapshot.row_count, "endpoint": path, "version": endpoint.version, "data_strategy": "snapshot", @@ -300,6 +429,29 @@ async def _serve_snapshot( headers=_deprecation_headers(endpoint), ) + @staticmethod + def _log_snapshot_rejection( + *, + request: Request, + endpoint: ApiEndpoint, + path: str, + principal: str | None, + started_at: float, + event: str, + ) -> None: + """Emit the required request context before a snapshot HTTP 422 response.""" + log.warning( + event, + request_id=resolve_request_id(request), + user=principal or "anonymous", + endpoint=path, + endpoint_id=str(endpoint.id), + status=status.HTTP_422_UNPROCESSABLE_ENTITY, + duration_ms=round((time.perf_counter() - started_at) * 1000, 2), + method=request.method, + client_ip=request.client.host if request.client else None, + ) + # ── Live mode ─────────────────────────────────────────────────────────── async def _serve_live( @@ -312,17 +464,7 @@ async def _serve_live( try: params = self._coerce_params(endpoint, request) except ValidationError as exc: - # Pull the first field error so the message stays readable. - first = exc.errors()[0] - field = ".".join(str(part) for part in first.get("loc", ())) or "?" - return JSONResponse( - status_code=status.HTTP_422_UNPROCESSABLE_ENTITY, - content={ - "detail": ( - f"Invalid value for parameter '{field}': {first.get('msg')}" - ) - }, - ) + return self._parameter_validation_response(exc) connection = await ConnectionRepository(self._db).get_by_id(endpoint.connection_id) if connection is None or not connection.is_active: @@ -377,13 +519,23 @@ async def _serve_live( headers=_deprecation_headers(endpoint), ) + @staticmethod + def _parameter_validation_response(exc: ValidationError) -> JSONResponse: + """Return a stable public error for typed query-parameter failures.""" + first = exc.errors()[0] + field = ".".join(str(part) for part in first.get("loc", ())) or "?" + return JSONResponse( + status_code=status.HTTP_422_UNPROCESSABLE_ENTITY, + content={"detail": f"Invalid value for parameter '{field}': {first.get('msg')}"}, + ) + @staticmethod def _coerce_params(endpoint: ApiEndpoint, request: Request) -> dict[str, Any]: param_schema = endpoint.param_schema_json or {} - # Scheduler defaults make snapshot refreshes autonomous, but they do - # not make required HTTP query parameters optional. Live callers must - # supply every descriptor marked required; configured defaults remain - # available to the scheduler and to omitted optional request fields. + # Schedule-owned bindings make snapshot refreshes autonomous, but they + # do not make required HTTP query parameters optional. Live and + # snapshot callers must supply every descriptor marked required; + # endpoint defaults apply only to omitted optional request fields. Model = build_param_model(param_schema, enforce_required=True) # Pull only declared params from the query string; ignore unknowns # so the legacy loop's behavior is preserved. Filter on diff --git a/backend/app/services/endpoint.py b/backend/app/services/endpoint.py index 4b9db85..a50036b 100644 --- a/backend/app/services/endpoint.py +++ b/backend/app/services/endpoint.py @@ -31,6 +31,7 @@ SqlPreviewRequest, SqlPreviewResponse, extract_bind_params, + require_snapshot_filter_mappings, ) from app.services.schedule_bindings import ScheduleBindingError, resolve_schedule_parameters from app.sql.executor import SqlExecutionError, execute_query @@ -205,9 +206,7 @@ async def update_endpoint( # Handle param_schema serialization if "param_schema" in payload.model_fields_set and payload.param_schema is not None: - effective_param_schema = { - k: v.model_dump() for k, v in payload.param_schema.items() - } + effective_param_schema = {k: v.model_dump() for k, v in payload.param_schema.items()} changes["param_schema_json"] = effective_param_schema changes.pop("param_schema", None) @@ -216,10 +215,11 @@ async def update_endpoint( changes["column_map_json"] = dict(payload.column_map) changes.pop("column_map", None) + effective_strategy = payload.data_strategy or obj.data_strategy + if self._schedule_repo is not None: schedule = await self._schedule_repo.get_by_endpoint_id(endpoint_id) if schedule is not None: - effective_strategy = payload.data_strategy or obj.data_strategy if effective_strategy != DataStrategy.snapshot: raise ScheduleBindingError( "Delete the attached schedule before changing this endpoint to live data." @@ -242,6 +242,8 @@ async def update_endpoint( window=schedule.window_config_json, ) + require_snapshot_filter_mappings(effective_strategy, effective_param_schema) + obj = await self._repo.update(obj, changes) log.info( diff --git a/backend/app/services/snapshot_filtering.py b/backend/app/services/snapshot_filtering.py new file mode 100644 index 0000000..223eb04 --- /dev/null +++ b/backend/app/services/snapshot_filtering.py @@ -0,0 +1,221 @@ +"""Typed request filtering for persisted snapshot rows.""" + +import json +from dataclasses import dataclass +from datetime import datetime +from functools import lru_cache +from typing import Any, Literal + +from pydantic import BaseModel, ValidationError + +from app.sql.param_models import build_param_model + +SnapshotFilterOperator = Literal["eq", "gte", "lte"] + + +@dataclass(frozen=True) +class CompiledSnapshotFilter: + """One validated request parameter to persisted-row comparison.""" + + parameter: str + column: str + operator: SnapshotFilterOperator + null_means_all: bool + param_type: str + coercion_key: str + value_model: type[BaseModel] + + def coerce_value(self, value: object) -> Any: + """Coerce one scheduled or cached value using the parameter's declared type.""" + return self.value_model.model_validate({"value": value}).model_dump()["value"] + + +def _coercion_descriptor(descriptor: dict[str, object]) -> dict[str, object]: + """Return only descriptor fields that affect validation of a supplied value.""" + coercion_descriptor: dict[str, object] = { + "type": descriptor.get("type", "string"), + "required": True, + } + if "max_length" in descriptor: + coercion_descriptor["max_length"] = descriptor["max_length"] + return coercion_descriptor + + +@lru_cache(maxsize=256) +def _cached_value_model(coercion_key: str) -> type[BaseModel]: + """Build each distinct snapshot value model once per process.""" + descriptor = json.loads(coercion_key) + return build_param_model({"value": descriptor}, enforce_required=True) + + +def compile_snapshot_filters( + param_schema: dict[str, object], +) -> tuple[CompiledSnapshotFilter, ...]: + """Compile persisted filter mappings once for one snapshot request.""" + compiled: list[CompiledSnapshotFilter] = [] + for name, descriptor in param_schema.items(): + if not isinstance(descriptor, dict): + continue + mapping = descriptor.get("snapshot_filter") + if not isinstance(mapping, dict): + continue + column = mapping.get("column") + operator = mapping.get("operator") + if not isinstance(column, str) or operator not in {"eq", "gte", "lte"}: + continue + coercion_key = json.dumps( + _coercion_descriptor(descriptor), + sort_keys=True, + separators=(",", ":"), + ) + compiled.append( + CompiledSnapshotFilter( + parameter=name, + column=column, + operator=operator, + null_means_all=mapping.get("null_means_all") is True, + param_type=str(descriptor.get("type", "string")), + coercion_key=coercion_key, + value_model=_cached_value_model(coercion_key), + ) + ) + return tuple(compiled) + + +def _coerce_cached_row_value(item: CompiledSnapshotFilter, value: object) -> Any: + """Normalize Oracle DATE/TIMESTAMP strings before typed comparison.""" + if item.param_type == "date" and isinstance(value, str): + try: + value = datetime.fromisoformat(value).date() + except ValueError: + # The shared date coercer still accepts the documented DD-MM-YYYY form. + pass + return item.coerce_value(value) + + +def snapshot_covers_request( + *, + filters: tuple[CompiledSnapshotFilter, ...], + request_params: dict[str, object], + resolved_params: dict[str, object], +) -> bool: + """Return whether one snapshot job run contains the requested selection.""" + for item in filters: + requested = request_params.get(item.parameter) + if requested is None: + # An omitted optional request is the same SQL NULL input used by a + # scheduled run. A fixed-value run is only a subset regardless of + # whether that query interprets NULL as all values or as a literal. + if item.parameter not in resolved_params or resolved_params[item.parameter] is not None: + return False + continue + if item.parameter not in resolved_params: + return False + resolved = resolved_params[item.parameter] + if resolved is None: + if item.operator == "eq" and item.null_means_all: + continue + return False + try: + coverage_value = item.coerce_value(resolved) + except ValidationError, ValueError, TypeError: + return False + if item.operator == "eq" and coverage_value != requested: + return False + # A scheduled lower bound must start on or before the requested lower bound. + if item.operator == "gte" and coverage_value > requested: + return False + # A scheduled upper bound must end on or after the requested upper bound. + if item.operator == "lte" and coverage_value < requested: + return False + return True + + +def validate_snapshot_parameter_ranges( + *, + filters: tuple[CompiledSnapshotFilter, ...], + request_params: dict[str, object], +) -> None: + """Reject an inclusive lower bound that is later than its upper bound.""" + bounds: dict[str, dict[str, list[Any]]] = {} + for item in filters: + if item.operator not in {"gte", "lte"}: + continue + value = request_params.get(item.parameter) + if value is not None: + bounds.setdefault(item.column, {}).setdefault(item.operator, []).append(value) + + for column, values in bounds.items(): + lower_values = values.get("gte", []) + upper_values = values.get("lte", []) + try: + lower = max(lower_values) if lower_values else None + upper = min(upper_values) if upper_values else None + except TypeError as exc: + raise ValueError( + f"Snapshot filter bounds have incompatible types for '{column}'." + ) from exc + if lower is not None and upper is not None: + try: + reversed_range = lower > upper + except TypeError as exc: + raise ValueError( + f"Snapshot filter bounds have incompatible types for '{column}'." + ) from exc + if reversed_range: + raise ValueError(f"Snapshot filter lower bound exceeds upper bound for '{column}'.") + + +def unavailable_snapshot_filter_columns( + *, + rows: list[dict[str, object]], + filters: tuple[CompiledSnapshotFilter, ...], +) -> list[str]: + """Return configured output columns absent from a non-empty snapshot.""" + if not rows: + return [] + available = {column for row in rows for column in row} + return sorted({item.column for item in filters} - available) + + +def filter_snapshot_rows( + *, + rows: list[dict[str, object]], + filters: tuple[CompiledSnapshotFilter, ...], + request_params: dict[str, object], +) -> list[dict[str, object]]: + """Apply explicitly configured, typed comparisons to cached rows.""" + filtered: list[dict[str, object]] = [] + for row in rows: + matches = True + coerced_values: dict[tuple[str, str], Any] = {} + for item in filters: + requested = request_params.get(item.parameter) + if requested is None: + continue + if item.column not in row: + matches = False + break + try: + cache_key = (item.column, item.coercion_key) + if cache_key not in coerced_values: + coerced_values[cache_key] = _coerce_cached_row_value(item, row[item.column]) + row_value = coerced_values[cache_key] + except ValidationError, ValueError, TypeError: + matches = False + break + if row_value is None: + matches = False + break + if item.operator == "eq" and row_value != requested: + matches = False + break + if item.operator == "gte" and row_value < requested: + matches = False + break + if item.operator == "lte" and row_value > requested: + matches = False + break + if matches: + filtered.append(row) + return filtered diff --git a/backend/tests/test_endpoints.py b/backend/tests/test_endpoints.py index 166f03b..8ff21fe 100644 --- a/backend/tests/test_endpoints.py +++ b/backend/tests/test_endpoints.py @@ -137,9 +137,7 @@ def test_endpoint_create_rejects_incompatible_static_default() -> None: path="invalid-default", connection_id=uuid.uuid4(), sql_text="SELECT * FROM stores WHERE id = :store_id", - param_schema={ - "store_id": {"type": "integer", "required": True, "default": "abc"} - }, + param_schema={"store_id": {"type": "integer", "required": True, "default": "abc"}}, allow_unauthenticated=True, ) @@ -147,9 +145,7 @@ def test_endpoint_create_rejects_incompatible_static_default() -> None: def test_endpoint_update_rejects_incompatible_static_default() -> None: with pytest.raises(ValueError, match="Invalid default"): EndpointUpdate( - param_schema={ - "store_id": {"type": "integer", "required": True, "default": "abc"} - } + param_schema={"store_id": {"type": "integer", "required": True, "default": "abc"}} ) @@ -160,8 +156,16 @@ def test_snapshot_endpoint_allows_schedule_to_own_parameter_values() -> None: 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}, + "start_date": { + "type": "date", + "required": True, + "snapshot_filter": {"column": "business_date", "operator": "gte"}, + }, + "end_date": { + "type": "date", + "required": True, + "snapshot_filter": {"column": "business_date", "operator": "lte"}, + }, }, allow_unauthenticated=True, data_strategy="snapshot", @@ -182,11 +186,13 @@ def test_snapshot_endpoint_accepts_dynamic_defaults() -> None: "type": "date", "required": True, "default_expression": "yesterday", + "snapshot_filter": {"column": "business_date", "operator": "gte"}, }, "end_date": { "type": "date", "required": True, "default_expression": "today", + "snapshot_filter": {"column": "business_date", "operator": "lte"}, }, }, allow_unauthenticated=True, @@ -206,6 +212,11 @@ def test_snapshot_endpoint_accepts_explicit_null_default() -> None: "type": "string", "required": False, "default_is_null": True, + "snapshot_filter": { + "column": "id", + "operator": "eq", + "null_means_all": True, + }, } }, allow_unauthenticated=True, @@ -502,7 +513,16 @@ async def test_update_live_endpoint_to_snapshot_leaves_values_to_schedule( "path": f"snapshot-update-{uuid.uuid4().hex[:8]}", "connection_id": connection.json()["id"], "sql_text": "SELECT * FROM orders WHERE business_date = :business_date", - "param_schema": {"business_date": {"type": "date", "required": True}}, + "param_schema": { + "business_date": { + "type": "date", + "required": True, + "snapshot_filter": { + "column": "business_date", + "operator": "eq", + }, + } + }, "allow_unauthenticated": True, "data_strategy": "live", }, @@ -548,7 +568,12 @@ async def test_update_snapshot_with_invalid_stored_schema_returns_422( "connection_id": connection.json()["id"], "sql_text": "SELECT * FROM stores WHERE id = :store_id", "param_schema": { - "store_id": {"type": "integer", "required": True, "default": 1} + "store_id": { + "type": "integer", + "required": True, + "default": 1, + "snapshot_filter": {"column": "id", "operator": "eq"}, + } }, "allow_unauthenticated": True, "data_strategy": "snapshot", @@ -561,9 +586,7 @@ async def test_update_snapshot_with_invalid_stored_schema_returns_422( update(ApiEndpoint) .where(ApiEndpoint.id == endpoint_id) .values( - param_schema_json={ - "store_id": {"type": "integer", "required": True, "default": "abc"} - } + param_schema_json={"store_id": {"type": "integer", "required": True, "default": "abc"}} ) ) await session.flush() diff --git a/backend/tests/test_schedules.py b/backend/tests/test_schedules.py index e7dd4d2..be0d4ff 100644 --- a/backend/tests/test_schedules.py +++ b/backend/tests/test_schedules.py @@ -264,8 +264,22 @@ async def _create_snapshot_endpoint_with_date_range(client: object) -> str: "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}, + "start_date": { + "type": "date", + "required": True, + "snapshot_filter": { + "column": "business_date", + "operator": "gte", + }, + }, + "end_date": { + "type": "date", + "required": True, + "snapshot_filter": { + "column": "business_date", + "operator": "lte", + }, + }, }, "data_strategy": "snapshot", }, @@ -347,8 +361,22 @@ async def test_endpoint_update_cannot_invalidate_attached_schedule( f"/api/v1/admin/endpoints/{endpoint_id}", json={ "param_schema": { - "start_date": {"type": "string", "required": True}, - "end_date": {"type": "date", "required": True}, + "start_date": { + "type": "string", + "required": True, + "snapshot_filter": { + "column": "business_date", + "operator": "gte", + }, + }, + "end_date": { + "type": "date", + "required": True, + "snapshot_filter": { + "column": "business_date", + "operator": "lte", + }, + }, } }, ) @@ -359,8 +387,22 @@ async def test_endpoint_update_cannot_invalidate_attached_schedule( f"/api/v1/admin/endpoints/{endpoint_id}", json={ "param_schema": { - "start_date": {"type": "date", "required": True}, - "replacement_end_date": {"type": "date", "required": True}, + "start_date": { + "type": "date", + "required": True, + "snapshot_filter": { + "column": "business_date", + "operator": "gte", + }, + }, + "replacement_end_date": { + "type": "date", + "required": True, + "snapshot_filter": { + "column": "business_date", + "operator": "lte", + }, + }, } }, ) @@ -640,7 +682,12 @@ async def test_create_schedule_with_invalid_stored_schema_returns_422( "connection_id": connection.json()["id"], "sql_text": "SELECT * FROM stores WHERE id = :store_id", "param_schema": { - "store_id": {"type": "integer", "required": True, "default": 1} + "store_id": { + "type": "integer", + "required": True, + "default": 1, + "snapshot_filter": {"column": "id", "operator": "eq"}, + } }, "allow_unauthenticated": True, "data_strategy": "snapshot", @@ -653,9 +700,7 @@ async def test_create_schedule_with_invalid_stored_schema_returns_422( update(ApiEndpoint) .where(ApiEndpoint.id == endpoint_id) .values( - param_schema_json={ - "store_id": {"type": "integer", "required": True, "default": "abc"} - } + param_schema_json={"store_id": {"type": "integer", "required": True, "default": "abc"}} ) ) await session.flush() diff --git a/backend/tests/test_snapshot_request_filtering.py b/backend/tests/test_snapshot_request_filtering.py new file mode 100644 index 0000000..8708342 --- /dev/null +++ b/backend/tests/test_snapshot_request_filtering.py @@ -0,0 +1,614 @@ +"""Public snapshot endpoints enforce and apply declared request parameters.""" + +import uuid +from collections.abc import Sequence +from datetime import UTC, datetime, timedelta + +import pytest +from app.models.endpoint import ApiEndpoint, DataStrategy +from app.models.job_run import JobRun, JobRunStatus +from app.models.snapshot import Snapshot +from app.repositories.job_run import JobRunRepository +from app.schemas.endpoint import EndpointCreate, ParamDescriptor +from app.services.snapshot_filtering import ( + compile_snapshot_filters, + validate_snapshot_parameter_ranges, +) +from httpx import AsyncClient +from sqlalchemy import select +from sqlalchemy.ext.asyncio import AsyncSession +from structlog.testing import capture_logs + + +def test_snapshot_filter_contract_is_preserved_in_parameter_schema() -> None: + descriptor = ParamDescriptor( + type="integer", + required=False, + default_is_null=True, + snapshot_filter={ + "column": "store_id", + "operator": "eq", + "null_means_all": True, + }, + ) + + assert descriptor.snapshot_filter is not None + assert descriptor.snapshot_filter.column == "store_id" + assert descriptor.snapshot_filter.null_means_all is True + + +def test_snapshot_endpoint_rejects_parameter_without_filter_mapping() -> None: + with pytest.raises(ValueError, match="snapshot filter mappings"): + EndpointCreate( + name="unmapped-snapshot", + path="unmapped-snapshot", + connection_id=uuid.uuid4(), + sql_text="SELECT * FROM stores WHERE store_id = :store_id", + param_schema={"store_id": {"type": "integer", "required": True}}, + allow_unauthenticated=True, + data_strategy="snapshot", + ) + + +def test_null_means_all_is_rejected_for_range_filter() -> None: + with pytest.raises(ValueError, match="null_means_all"): + ParamDescriptor( + type="date", + snapshot_filter={ + "column": "business_date", + "operator": "gte", + "null_means_all": True, + }, + ) + + +def test_null_means_all_is_rejected_for_required_parameter() -> None: + with pytest.raises(ValueError, match="optional parameters"): + ParamDescriptor( + type="integer", + required=True, + snapshot_filter={ + "column": "store_id", + "operator": "eq", + "null_means_all": True, + }, + ) + + +def test_snapshot_range_mappings_require_matching_parameter_types() -> None: + with pytest.raises(ValueError, match="same declared parameter type"): + EndpointCreate( + name="mixed-range-types", + path="mixed-range-types", + connection_id=uuid.uuid4(), + sql_text=( + "SELECT * FROM orders " + "WHERE business_date >= :start_date AND business_date <= :end_date" + ), + param_schema={ + "start_date": { + "type": "integer", + "snapshot_filter": { + "column": "business_date", + "operator": "gte", + }, + }, + "end_date": { + "type": "string", + "snapshot_filter": { + "column": "business_date", + "operator": "lte", + }, + }, + }, + allow_unauthenticated=True, + data_strategy="snapshot", + ) + + +def test_snapshot_range_validation_retains_duplicate_directional_bounds() -> None: + filters = compile_snapshot_filters( + { + "strict_start": { + "type": "integer", + "snapshot_filter": {"column": "sequence", "operator": "gte"}, + }, + "loose_start": { + "type": "integer", + "snapshot_filter": {"column": "sequence", "operator": "gte"}, + }, + "end": { + "type": "integer", + "snapshot_filter": {"column": "sequence", "operator": "lte"}, + }, + } + ) + + with pytest.raises(ValueError, match="lower bound exceeds upper bound"): + validate_snapshot_parameter_ranges( + filters=filters, + request_params={"strict_start": 10, "loose_start": 5, "end": 7}, + ) + + +async def _seed_snapshot_endpoint( + client: AsyncClient, + session: AsyncSession, +) -> str: + connection = await client.post( + "/api/v1/admin/connections/", + json={ + "name": f"snapshot-filter-conn-{uuid.uuid4().hex[:8]}", + "host": "oracle.example.com", + "service_name": "SVC", + "username": "hr", + "password": "secret", + }, + ) + assert connection.status_code == 201 + + path = f"snapshot-filter-{uuid.uuid4().hex[:8]}" + endpoint = ApiEndpoint( + name=f"snapshot-filter-{uuid.uuid4().hex[:8]}", + path=path, + connection_id=uuid.UUID(connection.json()["id"]), + sql_text=( + "SELECT * FROM orders " + "WHERE business_date BETWEEN :start_date AND :end_date " + "AND (:store_id IS NULL OR store_id = :store_id)" + ), + param_schema_json={ + "start_date": { + "type": "date", + "required": True, + "snapshot_filter": {"column": "business_date", "operator": "gte"}, + }, + "end_date": { + "type": "date", + "required": True, + "snapshot_filter": {"column": "business_date", "operator": "lte"}, + }, + "store_id": { + "type": "integer", + "required": False, + "default_is_null": True, + "snapshot_filter": { + "column": "store_id", + "operator": "eq", + "null_means_all": True, + }, + }, + }, + column_map_json={}, + allow_unauthenticated=True, + data_strategy=DataStrategy.snapshot, + ) + session.add(endpoint) + await session.flush() + + run = JobRun( + endpoint_id=endpoint.id, + started_at=datetime.now(UTC), + finished_at=datetime.now(UTC), + status=JobRunStatus.success, + row_count=2, + resolved_params_json={ + "start_date": "2026-08-01", + "end_date": "2026-08-31", + "store_id": None, + }, + trigger_source="schedule", + ) + session.add(run) + await session.flush() + + session.add( + Snapshot( + endpoint_id=endpoint.id, + job_run_id=run.id, + data=[ + {"business_date": "2026-08-10", "store_id": 1, "amount": 10}, + {"business_date": "2026-08-20", "store_id": 2, "amount": 20}, + ], + row_count=2, + ) + ) + await session.flush() + return path + + +@pytest.mark.integration +async def test_snapshot_requires_declared_required_parameters( + async_client: object, + db_session: AsyncSession, +) -> None: + client: AsyncClient = async_client # type: ignore[assignment] + path = await _seed_snapshot_endpoint(client, db_session) + + response = await client.get(f"/api/v1/data/{path}") + + assert response.status_code == 422 + assert "Field required" in response.json()["detail"] + assert any(name in response.json()["detail"] for name in ("start_date", "end_date")) + + +@pytest.mark.integration +async def test_snapshot_filters_dates_and_store_as_data_parameters( + async_client: object, + db_session: AsyncSession, +) -> None: + client: AsyncClient = async_client # type: ignore[assignment] + path = await _seed_snapshot_endpoint(client, db_session) + + response = await client.get( + f"/api/v1/data/{path}", + params={ + "start_date": "2026-08-15", + "end_date": "2026-08-25", + "store_id": "2", + }, + ) + + assert response.status_code == 200 + assert response.json()["data"] == [{"business_date": "2026-08-20", "store_id": 2, "amount": 20}] + assert response.json()["meta"]["row_count"] == 1 + + +@pytest.mark.integration +async def test_snapshot_filters_oracle_datetime_strings_as_dates( + async_client: object, + db_session: AsyncSession, +) -> None: + client: AsyncClient = async_client # type: ignore[assignment] + path = await _seed_snapshot_endpoint(client, db_session) + endpoint = ( + await db_session.execute(select(ApiEndpoint).where(ApiEndpoint.path == path)) + ).scalar_one() + snapshot = ( + await db_session.execute(select(Snapshot).where(Snapshot.endpoint_id == endpoint.id)) + ).scalar_one() + snapshot.data = [ + {"business_date": "2026-08-10 00:00:00", "store_id": 1, "amount": 10}, + {"business_date": "2026-08-20 17:45:00", "store_id": 2, "amount": 20}, + ] + await db_session.flush() + + response = await client.get( + f"/api/v1/data/{path}", + params={ + "start_date": "2026-08-15", + "end_date": "2026-08-25", + "store_id": "2", + }, + ) + + assert response.status_code == 200 + assert response.json()["data"] == [ + {"business_date": "2026-08-20 17:45:00", "store_id": 2, "amount": 20} + ] + + +@pytest.mark.integration +async def test_snapshot_rejects_request_outside_retained_coverage( + async_client: object, + db_session: AsyncSession, +) -> None: + client: AsyncClient = async_client # type: ignore[assignment] + path = await _seed_snapshot_endpoint(client, db_session) + + response = await client.get( + f"/api/v1/data/{path}", + params={"start_date": "2026-09-01", "end_date": "2026-09-07"}, + ) + + assert response.status_code == 422 + assert response.json()["code"] == "snapshot_out_of_coverage" + + +@pytest.mark.integration +async def test_snapshot_rejects_reversed_parameter_range( + async_client: object, + db_session: AsyncSession, +) -> None: + client: AsyncClient = async_client # type: ignore[assignment] + path = await _seed_snapshot_endpoint(client, db_session) + + response = await client.get( + f"/api/v1/data/{path}", + params={"start_date": "2026-08-25", "end_date": "2026-08-15"}, + ) + + assert response.status_code == 422 + assert response.json()["code"] == "invalid_parameter_range" + + +@pytest.mark.integration +async def test_snapshot_rejection_log_contains_required_request_context( + async_client: object, + db_session: AsyncSession, +) -> None: + client: AsyncClient = async_client # type: ignore[assignment] + path = await _seed_snapshot_endpoint(client, db_session) + + with capture_logs() as logs: + response = await client.get( + f"/api/v1/data/{path}", + params={"start_date": "2026-08-25", "end_date": "2026-08-15"}, + headers={"X-Request-ID": "snapshot-review-log"}, + ) + + assert response.status_code == 422 + rejection = next( + entry for entry in logs if entry.get("event") == "invalid_snapshot_parameter_range" + ) + assert rejection["request_id"] == "snapshot-review-log" + assert rejection["endpoint"] == path + assert rejection["status"] == 422 + assert rejection["method"] == "GET" + assert rejection["user"] + assert "client_ip" in rejection + assert rejection["duration_ms"] >= 0 + + +@pytest.mark.integration +async def test_snapshot_returns_empty_data_for_no_matches_inside_coverage( + async_client: object, + db_session: AsyncSession, +) -> None: + client: AsyncClient = async_client # type: ignore[assignment] + path = await _seed_snapshot_endpoint(client, db_session) + + response = await client.get( + f"/api/v1/data/{path}", + params={ + "start_date": "2026-08-21", + "end_date": "2026-08-22", + "store_id": "1", + }, + ) + + assert response.status_code == 200 + assert response.json()["data"] == [] + assert response.json()["meta"]["row_count"] == 0 + + +@pytest.mark.integration +async def test_snapshot_selects_newest_retained_snapshot_that_covers_request( + async_client: object, + db_session: AsyncSession, +) -> None: + client: AsyncClient = async_client # type: ignore[assignment] + path = await _seed_snapshot_endpoint(client, db_session) + endpoint = ( + await db_session.execute(select(ApiEndpoint).where(ApiEndpoint.path == path)) + ).scalar_one() + + newer_run = JobRun( + endpoint_id=endpoint.id, + started_at=datetime.now(UTC), + finished_at=datetime.now(UTC), + status=JobRunStatus.success, + row_count=1, + resolved_params_json={ + "start_date": "2026-08-25", + "end_date": "2026-08-31", + "store_id": None, + }, + trigger_source="schedule", + ) + db_session.add(newer_run) + await db_session.flush() + db_session.add( + Snapshot( + endpoint_id=endpoint.id, + job_run_id=newer_run.id, + data=[{"business_date": "2026-08-30", "store_id": 2, "amount": 30}], + row_count=1, + created_at=datetime.now(UTC) + timedelta(minutes=1), + ) + ) + await db_session.flush() + + response = await client.get( + f"/api/v1/data/{path}", + params={ + "start_date": "2026-08-20", + "end_date": "2026-08-20", + "store_id": "2", + }, + ) + + assert response.status_code == 200 + assert response.json()["data"] == [{"business_date": "2026-08-20", "store_id": 2, "amount": 20}] + + +@pytest.mark.integration +async def test_snapshot_candidate_job_runs_are_loaded_in_one_batch( + async_client: object, + db_session: AsyncSession, + monkeypatch: pytest.MonkeyPatch, +) -> None: + client: AsyncClient = async_client # type: ignore[assignment] + path = await _seed_snapshot_endpoint(client, db_session) + original_get_by_ids = JobRunRepository.get_by_ids + batch_call_count = 0 + + async def tracked_get_by_ids( + repository: JobRunRepository, + job_run_ids: Sequence[uuid.UUID], + ) -> Sequence[JobRun]: + nonlocal batch_call_count + batch_call_count += 1 + return await original_get_by_ids(repository, job_run_ids) + + monkeypatch.setattr(JobRunRepository, "get_by_ids", tracked_get_by_ids) + + response = await client.get( + f"/api/v1/data/{path}", + params={"start_date": "2026-08-20", "end_date": "2026-08-20"}, + ) + + assert response.status_code == 200 + assert batch_call_count == 1 + + +@pytest.mark.integration +async def test_omitted_all_value_filter_skips_fixed_value_snapshot( + async_client: object, + db_session: AsyncSession, +) -> None: + client: AsyncClient = async_client # type: ignore[assignment] + path = await _seed_snapshot_endpoint(client, db_session) + endpoint = ( + await db_session.execute(select(ApiEndpoint).where(ApiEndpoint.path == path)) + ).scalar_one() + + fixed_store_run = JobRun( + endpoint_id=endpoint.id, + started_at=datetime.now(UTC), + finished_at=datetime.now(UTC), + status=JobRunStatus.success, + row_count=1, + resolved_params_json={ + "start_date": "2026-08-01", + "end_date": "2026-08-31", + "store_id": 1, + }, + trigger_source="schedule", + ) + db_session.add(fixed_store_run) + await db_session.flush() + db_session.add( + Snapshot( + endpoint_id=endpoint.id, + job_run_id=fixed_store_run.id, + data=[{"business_date": "2026-08-20", "store_id": 1, "amount": 99}], + row_count=1, + created_at=datetime.now(UTC) + timedelta(minutes=1), + ) + ) + await db_session.flush() + + response = await client.get( + f"/api/v1/data/{path}", + params={"start_date": "2026-08-20", "end_date": "2026-08-20"}, + ) + + assert response.status_code == 200 + assert response.json()["data"] == [{"business_date": "2026-08-20", "store_id": 2, "amount": 20}] + + +@pytest.mark.integration +async def test_omitted_optional_filter_requires_a_null_resolved_snapshot( + async_client: object, + db_session: AsyncSession, +) -> None: + client: AsyncClient = async_client # type: ignore[assignment] + path = await _seed_snapshot_endpoint(client, db_session) + endpoint = ( + await db_session.execute(select(ApiEndpoint).where(ApiEndpoint.path == path)) + ).scalar_one() + stored_descriptor = endpoint.param_schema_json["store_id"] + assert isinstance(stored_descriptor, dict) + store_descriptor = dict(stored_descriptor) + stored_mapping = store_descriptor["snapshot_filter"] + assert isinstance(stored_mapping, dict) + store_mapping = dict(stored_mapping) + store_mapping["null_means_all"] = False + store_descriptor["snapshot_filter"] = store_mapping + endpoint.param_schema_json = { + **endpoint.param_schema_json, + "store_id": store_descriptor, + } + + fixed_store_run = JobRun( + endpoint_id=endpoint.id, + started_at=datetime.now(UTC), + finished_at=datetime.now(UTC), + status=JobRunStatus.success, + row_count=1, + resolved_params_json={ + "start_date": "2026-08-01", + "end_date": "2026-08-31", + "store_id": 1, + }, + trigger_source="schedule", + ) + db_session.add(fixed_store_run) + await db_session.flush() + db_session.add( + Snapshot( + endpoint_id=endpoint.id, + job_run_id=fixed_store_run.id, + data=[{"business_date": "2026-08-20", "store_id": 1, "amount": 99}], + row_count=1, + created_at=datetime.now(UTC) + timedelta(minutes=1), + ) + ) + await db_session.flush() + + response = await client.get( + f"/api/v1/data/{path}", + params={"start_date": "2026-08-20", "end_date": "2026-08-20"}, + ) + + assert response.status_code == 200 + assert response.json()["data"] == [{"business_date": "2026-08-20", "store_id": 2, "amount": 20}] + + +@pytest.mark.integration +async def test_legacy_snapshot_cannot_silently_ignore_unmapped_parameters( + async_client: object, + db_session: AsyncSession, +) -> None: + client: AsyncClient = async_client # type: ignore[assignment] + path = await _seed_snapshot_endpoint(client, db_session) + endpoint = ( + await db_session.execute(select(ApiEndpoint).where(ApiEndpoint.path == path)) + ).scalar_one() + endpoint.param_schema_json = { + name: {key: value for key, value in descriptor.items() if key != "snapshot_filter"} + for name, descriptor in endpoint.param_schema_json.items() + if isinstance(descriptor, dict) + } + await db_session.flush() + + response = await client.get( + f"/api/v1/data/{path}", + params={"start_date": "2026-08-15", "end_date": "2026-08-25"}, + ) + + assert response.status_code == 422 + assert response.json()["code"] == "snapshot_filter_not_configured" + + +@pytest.mark.integration +async def test_snapshot_rejects_unavailable_cached_filter_column( + async_client: object, + db_session: AsyncSession, +) -> None: + client: AsyncClient = async_client # type: ignore[assignment] + path = await _seed_snapshot_endpoint(client, db_session) + endpoint = ( + await db_session.execute(select(ApiEndpoint).where(ApiEndpoint.path == path)) + ).scalar_one() + store_descriptor = endpoint.param_schema_json["store_id"] + assert isinstance(store_descriptor, dict) + store_descriptor["snapshot_filter"] = { + "column": "missing_store_column", + "operator": "eq", + "null_means_all": True, + } + await db_session.flush() + + response = await client.get( + f"/api/v1/data/{path}", + params={ + "start_date": "2026-08-15", + "end_date": "2026-08-25", + "store_id": "2", + }, + ) + + assert response.status_code == 422 + assert response.json()["code"] == "snapshot_filter_column_unavailable" diff --git a/docs/PROJECT_ANALYSIS.md b/docs/PROJECT_ANALYSIS.md index 713513c..c7db2a4 100644 --- a/docs/PROJECT_ANALYSIS.md +++ b/docs/PROJECT_ANALYSIS.md @@ -1,5 +1,13 @@ # QueryGateway: Project Analysis & Manifesto +> **Historical analysis with current implementation notes (updated 2026-08-31).** This document +> records the product and technology evaluation that preceded implementation. The authoritative +> runtime contract is now [Architecture](architecture.md), and endpoint/scheduler/snapshot +> parameter behavior is defined in [Parameter contracts](scheduler_parameter_bindings.md). +> QueryGateway currently runs Python 3.14+, a Vite + React admin SPA, an in-process APScheduler +> whose PostgreSQL schedule definitions are restored at startup, schedule-owned SQL bindings, and +> coverage-aware typed snapshot filtering. + ## 1. Executive Summary **Project Name:** QueryGateway @@ -329,7 +337,7 @@ A step-by-step wizard flow: | **Oracle Connectivity** | python-oracledb | Official Oracle Python driver; thin mode (no Oracle Client needed) | | **ORM / Query Layer** | SQLAlchemy 2.0 | Database abstraction for future multi-DB support | | **App Database** | PostgreSQL | Stores connections, endpoints, schedules, logs, cached results | -| **Task Scheduling** | APScheduler | Python-native scheduler with cron triggers; persistent job store in PostgreSQL | +| **Task Scheduling** | APScheduler | Python-native in-process scheduler; definitions persist in PostgreSQL and active jobs are restored on startup | | **Frontend** | Next.js 14+ (App Router) | User preference; SSR/SSG capabilities | | **UI Components** | shadcn/ui + Tailwind CSS | User preference; modern, accessible component library | | **Testing** | pytest + pytest-asyncio | User preference; comprehensive async test support | @@ -501,7 +509,8 @@ The proposed stack from Section 3.4 is evaluated below on a component-by-compone - *Risk:* FastAPI's dependency injection system can become complex in large applications. *Mitigation:* Establish clear patterns early (service layer, repository pattern) and document conventions. - *Risk:* Python GIL limits true CPU parallelism. *Mitigation:* Not a concern — this project is I/O-bound (DB queries, HTTP serving), not CPU-bound. -**Recommendation:** Python 3.12+ (instead of 3.11+) should be targeted. Python 3.12 brings significant performance improvements (10-15% faster), better error messages, and will have longer support. Python 3.11 EOL is October 2027; Python 3.12 EOL is October 2028. +**Implemented target:** Python 3.14+. The repository's Ruff, mypy, dependencies, Docker image, and +developer setup all use Python 3.14 or newer. --- @@ -563,7 +572,7 @@ The proposed stack from Section 3.4 is evaluated below on a component-by-compone | Criterion | Assessment | |-----------|-----------| | **Cron Triggers** | Full cron expression support matches the scheduling requirements in Module 4. | -| **Persistent Job Store** | PostgreSQL-backed job store ensures scheduled jobs survive application restarts. | +| **Restart restoration** | Schedule definitions and next-run metadata persist in PostgreSQL; the single API process restores active jobs into APScheduler's in-memory job store at startup. | | **Python-Native** | Runs in-process, no external dependencies (no Redis, no RabbitMQ, no Celery). | | **Dynamic Jobs** | Jobs can be added, modified, paused, and removed at runtime — essential for the wizard-based scheduling flow. | diff --git a/docs/architecture.md b/docs/architecture.md index 1df3447..89548a7 100644 --- a/docs/architecture.md +++ b/docs/architecture.md @@ -41,25 +41,66 @@ ### SQL Safety -All user-defined SQL is executed via SQLAlchemy `text()` with named bind parameters (`:param_name`). String interpolation of user input is prohibited at all layers. Bind values are validated through typed Pydantic schemas before reaching the query executor. - -Live data requests enforce every parameter marked `required`. Endpoint defaults apply only to live request/preview behavior; they are not scheduler configuration. Date query parameters accept `YYYY-MM-DD` and `DD-MM-YYYY`; both formats are normalized to Python `date` values before database binding. +All user-defined SQL is executed via SQLAlchemy `text()` with named bind parameters +(`:param_name`). String interpolation of user input is prohibited at all layers. A bind marker +inside a single-quoted SQL literal is text and is intentionally ignored. Bind values are validated +through typed Pydantic schemas before reaching the query executor. + +Live and snapshot data requests enforce every parameter marked `required`, even if an endpoint +default exists. Preview inputs are temporary. Endpoint defaults resolve only omitted optional live +request values; scheduled queries instead require a separate binding for every SQL parameter. Date +query parameters accept `YYYY-MM-DD` and `DD-MM-YYYY`; both formats are normalized to Python +`date` values before database binding. See [Endpoint, scheduler, and snapshot parameter +contracts](scheduler_parameter_bindings.md). ### Authentication -- Admin API: session-based or JWT Bearer (TBD per Phase 1). -- Data endpoints: every `/api/v1/data/*` request is authenticated. Middleware enforces the endpoint's Bearer token, Basic Auth, or API key method when configured; otherwise it requires the platform admin Bearer token. Anonymous data access is never allowed. +- Admin API: an environment-seeded administrator authenticates through `/api/v1/auth/login`; the + protected admin API uses the resulting JWT Bearer token. +- Data endpoints: every `/api/v1/data/*` request is authenticated. The data service enforces the + endpoint's Bearer token, Basic Auth, or API key method when configured; otherwise it requires the + platform admin Bearer token. Anonymous data access is never allowed. - Credentials are hashed with `bcrypt`; tokens are issued/verified with `PyJWT`. ### Scheduler -APScheduler 3.x runs in-process with an in-memory job store. Schedule definitions are persisted in PostgreSQL, and every active schedule is registered when it is created, updated, resumed, or restored during API startup. Each schedule owns its timezone, parameter bindings, and optional date-window preset. The database `next_run_at` is captured as the nominal `scheduled_for` time, so delayed execution still resolves the same logical date; manual runs can provide an explicit logical date. A unique `(schedule_id, scheduled_for)` constraint makes a logical run idempotent. Execution telemetry includes start/finish times, logical date, window boundaries, resolved parameters, trigger source, binding hash, row count, status, and errors. Deleting a schedule sets the historical job run's `schedule_id` to `NULL`, preserving both audit history and snapshots. Deleting an endpoint removes its schedule and cached snapshots, unregisters its in-memory job after the database commit, and sets the historical job run's `endpoint_id` (and cascaded `schedule_id`) to `NULL` so the audit record remains available. +APScheduler 3.x runs in-process with an in-memory job store. Schedule definitions are persisted in +PostgreSQL, and every active schedule is registered when it is created, updated, resumed, or +restored during API startup. Each schedule owns its timezone, parameter bindings, and optional +date-window preset. The database `next_run_at` is captured as the nominal `scheduled_for` time, so +delayed execution still resolves the same logical date; manual runs can provide an explicit logical +date. A unique `(schedule_id, scheduled_for)` constraint makes a logical run idempotent. Execution +telemetry includes start/finish times, logical date, window boundaries, resolved parameters, +trigger source, binding hash, row count, status, and errors. Deleting a schedule sets the +historical job run's `schedule_id` to `NULL`, preserving both audit history and snapshots. Deleting +an endpoint removes its schedule and cached snapshots, unregisters its in-memory job after the +database commit, and sets the historical job run's `endpoint_id` (and cascaded `schedule_id`) to +`NULL` so the audit record remains available. The scheduler is intentionally single-process. Run one API process/replica unless distributed scheduler coordination is added; otherwise each process would register and execute the same persisted schedules. ### Snapshot Cache -Scheduled endpoints can serve results from a PostgreSQL JSONB snapshot rather than executing live queries. Schedule creation requires exact coverage of the endpoint's SQL binds. The declarative binding sources are fixed literal, explicit SQL `NULL` for an optional bind, logical run date, relative logical date, and inclusive window start/end. Supported windows are previous day, last N complete days, week to date, previous week, month to date, and previous month. Arbitrary Python, JavaScript, and SQL expressions are intentionally unsupported. See [Scheduler parameter bindings](scheduler_parameter_bindings.md) for the API contract and date semantics. +Scheduled endpoints can serve results from a PostgreSQL JSONB snapshot rather than executing live +queries. Schedule creation requires exact coverage of the endpoint's SQL binds. The declarative +binding sources are fixed literal, explicit SQL `NULL` for an optional bind, logical run date, +relative logical date, and inclusive window start/end. Supported windows are previous day, last N +complete days, week to date, previous week, month to date, and previous month. + +Each parameterized snapshot endpoint also declares an explicit final cached-output-column mapping +and one whitelisted operator (`eq`, `gte`, or `lte`). The data plane validates required fields and +types, rejects reversed ranges, selects the newest retained job snapshot whose persisted resolved +values cover the request, validates mapped columns, and filters the cached rows. It never falls +back to returning the entire snapshot when mappings or request parameters are missing. A covered +request with no business rows returns an empty `data` array; an out-of-coverage request is an +explicit HTTP 422, while the absence of any retained snapshot is HTTP 503. These mappings are data +selection only; tenant authorization continues to be owned by the endpoint authentication method. +Arbitrary Python, JavaScript, and SQL expressions are intentionally unsupported. + +Snapshot filter mappings extend the endpoint's existing JSON parameter document and therefore do +not change the relational schema. Schedule-owned parameter bindings and logical-run audit fields +are relational and are managed by Alembic. See [Endpoint, scheduler, and snapshot parameter +contracts](scheduler_parameter_bindings.md) for the full contract and date semantics. ### Configuration @@ -67,7 +108,9 @@ Pydantic Settings v2 loads all configuration from environment variables (`.env` ### Logging -`structlog` emits structured JSON with mandatory correlation fields (`request_id`, `user`, `endpoint`, `status`, `duration_ms`, `event`). Sensitive fields are redacted at middleware level before emission. +`structlog` emits structured JSON with mandatory correlation fields (`request_id`, `user`, +`endpoint`, `status`, `duration_ms`, `method`, `client_ip`, `event`). Sensitive fields are redacted +at middleware level before emission. ## Directory Layout @@ -79,10 +122,9 @@ QueryGateway/ │ │ ├── config.py # Pydantic Settings │ │ ├── models/ # SQLAlchemy models (Phase 1+) │ │ ├── routers/ # Route handlers (Phase 1+) -│ │ ├── services/ # Business logic layer (Phase 1+) +│ │ ├── services/ # Business logic, scheduler, snapshot filtering │ │ ├── repositories/ # DB access layer (Phase 1+) │ │ ├── auth/ # JWT + bcrypt utilities (Phase 3+) -│ │ ├── scheduler/ # APScheduler integration (Phase 5+) │ │ └── sql/ # SQL execution and validation (Phase 4+) │ ├── alembic/ # Migration environment (Phase 1+) │ ├── tests/ # Pytest test suite @@ -104,12 +146,13 @@ QueryGateway/ │ ├── Dockerfile.frontend │ └── nginx.conf ├── docs/ -│ ├── architecture.md # This file -│ ├── conventions.md # Coding standards -│ ├── contributing.md # Onboarding guide -│ ├── deployment.md # Deployment runbook (Phase 7) -│ ├── operations.md # Backup/restore, monitoring, troubleshooting (Phase 7) -│ └── security_checklist.md # Security validation checklist (Phase 7) +│ ├── architecture.md # This file +│ ├── scheduler_parameter_bindings.md # Parameter and snapshot contracts +│ ├── conventions.md # Coding standards +│ ├── contributing.md # Onboarding guide +│ ├── deployment.md # Deployment runbook +│ ├── operations.md # Backup/restore and troubleshooting +│ └── security_checklist.md # Security validation checklist ├── .github/ │ ├── workflows/ │ │ ├── backend.yml # Backend CI diff --git a/docs/code_refacoring_plan.md b/docs/code_refacoring_plan.md index b3d69eb..b701b33 100644 --- a/docs/code_refacoring_plan.md +++ b/docs/code_refacoring_plan.md @@ -1,5 +1,11 @@ # QueryGateway Phased Refactor Plan +> **Historical implementation plan.** The refactor described here has been superseded by the +> current repository structure. Use [Architecture](architecture.md), +> [Implementation progress](progress.md), and [Endpoint, scheduler, and snapshot parameter +> contracts](scheduler_parameter_bindings.md) as the current source of truth. Line numbers and +> pre-refactor helper names below are retained only as execution history. + ## Context A code-health audit of the QueryGateway monorepo (FastAPI + React, Oracle query gateway) scored **66/100**. Findings include one critical security gap — `/api/v1/admin/*` is fully unauthenticated, with no User model, login endpoint, dependency, or middleware — plus architecture leakage (`routers/data.py` bypasses the service layer), per-request Oracle client init, hand-rolled param coercion, ~6 near-identical repositories and routers, monolithic frontend pages (350-450 LOC), no frontend login UI, and 1 frontend test for 4,173 LOC. Existing utilities (JWT helpers, bcrypt helpers, ServiceContext patterns, FastAPI lifespan) are reusable and should be the foundation. @@ -320,4 +326,4 @@ cd backend && alembic upgrade head && pytest tests/test_endpoints.py -v | 5 | 87 | duplication down | | 6 | 89 | frontend duplication down, testable | | 7 | 91 | coverage + CI thresholds enforced | -| 8 | 92 | defense-in-depth on data plane | \ No newline at end of file +| 8 | 92 | defense-in-depth on data plane | diff --git a/docs/contributing.md b/docs/contributing.md index 159ed77..54e86ee 100644 --- a/docs/contributing.md +++ b/docs/contributing.md @@ -16,7 +16,7 @@ ```sh git clone -cd DB2API-Exposure +cd QueryGateway ``` ### 2. Backend @@ -50,9 +50,18 @@ npm install ```sh # Copy root env example cp .env.example .env -# Edit .env to set JWT_SECRET_KEY -docker compose up -d +# Generate the required secrets (run with the backend environment activated) +python -c "import secrets; print(secrets.token_urlsafe(48))" +python -c "from cryptography.fernet import Fernet; print(Fernet.generate_key().decode())" +cd backend +python -c "from getpass import getpass; from app.auth.hashing import hash_password; print(hash_password(getpass('Admin password: ')))" +cd .. + +# Edit .env to set JWT_SECRET_KEY, ENCRYPTION_KEY, ADMIN_USERNAME, and +# ADMIN_PASSWORD_HASH. Paste the generated bcrypt hash; never store plaintext. + +docker compose up -d --build ``` Compose runs the one-shot `migrate` service before the API starts, so a fresh local `db_data` volume is initialized automatically. @@ -127,6 +136,10 @@ make docker-build 6. Update `docs/` if you changed API contracts, config, or significant behavior. 7. Open a PR against `main`. +For endpoint parameter, schedule binding, or snapshot filtering changes, update +[the canonical parameter contract](scheduler_parameter_bindings.md) and keep the related root and +`.github/instructions/` agent guidance aligned. + ## What CI Checks | Job | Checks | diff --git a/docs/conventions.md b/docs/conventions.md index 135cc48..abf1b8a 100644 --- a/docs/conventions.md +++ b/docs/conventions.md @@ -35,6 +35,8 @@ - Never edit an already-applied migration file; add a new revision. - Verify rollback: `alembic downgrade -1` must succeed in CI. - Migration files must be committed alongside the model change. +- JSON contract extensions inside existing JSONB columns do not require a relational migration; + document and test the API compatibility impact instead. ## Security Constraints @@ -46,6 +48,21 @@ - Never store secrets in code, fixtures, tests, or docs. Use environment variables. - Never return plaintext credentials or token bodies in API responses or logs. +## Parameter Contract + +- Keep SQL preview samples, optional live-request defaults, schedule-owned bindings, and snapshot + request filters as separate concepts. +- Required descriptors remain required for live and snapshot HTTP requests even if they contain a + default. +- Schedules must bind every endpoint SQL parameter explicitly and must never read endpoint + defaults during execution. +- Parameterized snapshot endpoints must map every request parameter to a final cached output + column with `eq`, `gte`, or `lte`, prove retained-snapshot coverage, and filter cached rows. +- Snapshot filter mappings are row selection only; authorization belongs to the endpoint auth + method. +- The canonical behavior and error codes live in + [Endpoint, scheduler, and snapshot parameter contracts](scheduler_parameter_bindings.md). + ## Logging Standards - Use `structlog` for all structured JSON logging. diff --git a/docs/deployment.md b/docs/deployment.md index 4bf8b97..5ce2b7b 100644 --- a/docs/deployment.md +++ b/docs/deployment.md @@ -19,10 +19,13 @@ Create a `.env` file from `.env.example` and configure: |----------|----------|---------|-------------| | `DATABASE_URL` | Yes | — | PostgreSQL connection string (`postgresql+asyncpg://user:pass@host:5432/db`) | | `ENCRYPTION_KEY` | Yes | — | Fernet encryption key for credential storage (generate: `python -c "from cryptography.fernet import Fernet; print(Fernet.generate_key().decode())"`) | +| `JWT_SECRET_KEY` | Yes | — | High-entropy platform JWT signing key with at least 32 characters | +| `ADMIN_USERNAME` | Yes | — | Seeded platform administrator username | +| `ADMIN_PASSWORD_HASH` | Yes | — | Seeded administrator bcrypt hash; generate with `app.auth.hashing.hash_password()` | | `APP_ENV` | No | `development` | Environment identifier (`development`, `staging`, `production`) | | `DEBUG` | No | `false` | Enable debug mode (never `true` in production) | | `LOG_LEVEL` | No | `INFO` | Logging level (`DEBUG`, `INFO`, `WARNING`, `ERROR`) | -| `CORS_ORIGINS` | No | `*` | Comma-separated allowed CORS origins | +| `CORS_ORIGINS` | No | `http://localhost:5173` | Comma-separated allowed CORS origins; wildcard is rejected | ### Generating an Encryption Key @@ -38,7 +41,7 @@ Store this key securely. If the key is lost, all encrypted credentials (Oracle p ```bash # Clone the repository -git clone && cd DB2API-Exposure +git clone && cd QueryGateway # Configure environment cp .env.example .env @@ -241,13 +244,24 @@ alembic downgrade -1 ## Post-Deployment Verification -1. **Health check**: `curl http://localhost:8000/api/v1/admin/health/live` -2. **Database ready**: `curl http://localhost:8000/api/v1/admin/health/ready` -3. **Admin UI**: Open `http://localhost:3000` in a browser +1. **Health check**: use `curl http://localhost/api/v1/admin/health/live` for Docker, or + `curl http://localhost:8000/api/v1/admin/health/live` for a bare-metal backend. +2. **Database ready**: use `curl http://localhost/api/v1/admin/health/ready` for Docker, or + `curl http://localhost:8000/api/v1/admin/health/ready` for bare metal. +3. **Admin UI**: Open `http://localhost` for Docker, or `http://localhost:5173` for Vite development 4. **Create first connection**: Use the Connections page to add an Oracle data source 5. **Test connection**: Click "Test" to verify Oracle connectivity 6. **Create first endpoint**: Use the API Endpoints page wizard -7. **Verify data endpoint**: `curl -H "Authorization: Bearer " http://localhost:8000/api/v1/data/` +7. **Verify data endpoint**: use + `curl -H "Authorization: Bearer " "http://localhost/api/v1/data/?required_param=value"` + for Docker, or change the origin to `http://localhost:8000` for bare metal. +8. **Verify required inputs**: omit a required live and snapshot parameter and confirm HTTP 422 +9. **Verify snapshot selection**: preview and run the schedule, confirm an in-coverage request is + filtered, and confirm an out-of-coverage request returns `snapshot_out_of_coverage` +10. **Verify scheduler restoration**: for Docker, run `docker compose restart api`; for bare metal, + restart the single API process with its process manager (for example, + `sudo systemctl restart querygateway-api`). Then confirm active jobs are registered again on + the health dashboard. ## Security Hardening for Production diff --git a/docs/operations.md b/docs/operations.md index 271d3e5..1c0dfe8 100644 --- a/docs/operations.md +++ b/docs/operations.md @@ -154,6 +154,8 @@ Logs are emitted as structured JSON via structlog. Key fields: | `endpoint` | API path | | `status` | HTTP status code | | `duration_ms` | Request duration | +| `method` | HTTP method | +| `client_ip` | Request source after trusted-proxy resolution | | `job_id` | Scheduler job identifier | | `run_id` | Job run identifier | | `row_count` | Query result row count | @@ -178,7 +180,8 @@ cat app.log | jq 'select(.status == 401)' ### Service Won't Start -1. **Check environment variables**: Ensure `DATABASE_URL` and `ENCRYPTION_KEY` are set +1. **Check environment variables**: Ensure `DATABASE_URL`, `ENCRYPTION_KEY`, and + `JWT_SECRET_KEY` are set and valid 2. **Check PostgreSQL**: `pg_isready -h localhost` 3. **Check migrations**: `alembic current` should show the latest revision 4. **Check logs**: Look for startup errors in container logs @@ -218,6 +221,12 @@ docker compose logs api --tail 50 - Oracle connection timeout - Large result set exceeding memory - Concurrent job limit reached + - A missing or extra schedule parameter binding + - A date window source without a configured window preset + - An invalid IANA timezone or a logical-date replay outside the intended business calendar + +Use `POST /api/v1/admin/schedules/preview` with the proposed timing and binding payload to inspect +the next resolved logical dates and typed SQL values before saving a schedule. ### Auth Token Issues @@ -229,10 +238,39 @@ docker compose logs api --tail 50 ### Data Endpoint Returns Unexpected Results 1. **Check endpoint config**: `GET /api/v1/admin/endpoints/{id}` -2. **Verify parameters**: Required params must be provided in query string -3. **Check column mapping**: `column_map_json` may rename output columns -4. **Snapshot staleness**: For snapshot endpoints, check the `snapshot_created_at` in response -5. **SQL syntax**: Use the SQL preview feature to test queries interactively +2. **Verify parameters**: Required parameters must be present in the query string for both live + and snapshot requests. Endpoint defaults do not satisfy a missing required HTTP value. +3. **Check parameter ownership**: Preview samples are temporary, live defaults apply only to + omitted optional requests, and schedule execution uses the schedule's own bindings. +4. **Check snapshot mappings and error codes**: + - Snapshot endpoints require a cached-column mapping for every request parameter. + - The mapped name is the final output name after `column_map_json` renaming. + - `snapshot_filter_not_configured` means an older endpoint must be updated in the endpoint edit + dialog or admin API. + - `invalid_parameter_range` means a lower request bound is later than its upper bound. + - `snapshot_out_of_coverage` means no retained job run covers the complete request; inspect + schedule bindings, window, timezone, and snapshot retention. + - `snapshot_filter_column_unavailable` means a configured mapped column is absent from the + non-empty cached payload. + - A covered request with no matching business rows is successful and returns `data: []`. + - No retained snapshot returns HTTP 503 rather than a coverage error. +5. **Check snapshot staleness**: For snapshot endpoints, compare `snapshot_created_at`, + `snapshot_row_count`, and the filtered `row_count` response metadata. +6. **Check SQL syntax**: Use the SQL preview feature with temporary sample values. Ensure bind + placeholders are not enclosed in single quotes. + +See [Endpoint, scheduler, and snapshot parameter contracts](scheduler_parameter_bindings.md) for +the complete selection and coverage rules. + +### Schedule and Endpoint Deletion + +- Deleting a schedule unregisters the in-memory job after the database commit and preserves its + immutable job-run history. Historical `schedule_id` values become `NULL`; retained snapshots + continue to reference their job runs. +- Deleting an endpoint cascades its active schedule and cached snapshots, unregisters the schedule + job, and preserves historical job runs with nullable endpoint/schedule references. +- Back up PostgreSQL before bulk deletion when job-run audit history or cached payloads are part of + an operational retention requirement. ## Upgrade Procedures diff --git a/docs/phase2.md b/docs/phase2.md index b70716f..bae8b80 100644 --- a/docs/phase2.md +++ b/docs/phase2.md @@ -4,6 +4,22 @@ Extend QueryGateway from Oracle-only to a multi-database platform. Users will be able to register connections to different database engines, test them, and expose SQL queries as REST endpoints — exactly as today, but engine-agnostic. +## Current compatibility baseline + +Any future database adapter must preserve the current v1 parameter behavior: + +- unquoted named binds, typed validation, and required-parameter enforcement for live and snapshot + requests; +- separation of temporary preview samples, optional live defaults, and schedule-owned bindings; +- timezone-aware logical schedule dates and declarative date windows; and +- complete snapshot mappings, persisted-run coverage checks, and typed cached-row filtering before + data is returned. + +The canonical behavior is documented in +[Endpoint, scheduler, and snapshot parameter contracts](scheduler_parameter_bindings.md). Engine +adapters may translate bind syntax internally only after validation; they must not introduce raw +string interpolation or weaken the public `/api/v1` contract. + --- ## Proposed Database Types diff --git a/docs/production_deployment.md b/docs/production_deployment.md index 2c487f2..0c0bc9b 100644 --- a/docs/production_deployment.md +++ b/docs/production_deployment.md @@ -275,7 +275,7 @@ python3 -m venv .venv . .venv/bin/activate pip install --upgrade pip pip install -r requirements.txt -python -c "from app.auth.hashing import hash_password; print(hash_password('replace-with-strong-password'))" +python -c "from getpass import getpass; from app.auth.hashing import hash_password; print(hash_password(getpass('Admin password: ')))" ``` Important: @@ -381,6 +381,16 @@ For a controlled preflight migration before starting the API, run: docker compose -f docker-compose.yml -f compose.production.yml up --build --force-recreate migrate ``` +After the API starts, verify the container reports the current Alembic head: + +```bash +docker compose -f docker-compose.yml -f compose.production.yml exec -T api alembic current +``` + +Schedule-owned bindings and logical-run audit fields require the migration that introduced them; +snapshot filter mappings themselves live in the existing endpoint JSON and do not add a relational +migration. + ### 7. Start Application ```bash @@ -438,7 +448,11 @@ Before go-live: - Attach authentication to the endpoint. - Call one `/api/v1/data/*` route with valid credentials. - Confirm unauthenticated calls fail for protected data routes. -- Confirm scheduler health if scheduled snapshots are used. +- Confirm required parameters return HTTP 422 when omitted from both live and snapshot routes. +- If scheduled snapshots are used, preview the schedule's resolved bindings, run it once, and + confirm an in-coverage request returns only filtered rows while an out-of-coverage request returns + `code=snapshot_out_of_coverage`. +- Confirm scheduler health and verify active jobs return after an API-container restart. ## Startup After Windows Reboot diff --git a/docs/progress.md b/docs/progress.md index e5b9258..42f24ff 100644 --- a/docs/progress.md +++ b/docs/progress.md @@ -1,5 +1,7 @@ # Implementation Progress +**Last updated:** 2026-08-31 + ## Phase 0: Repository Scaffolding, CI/CD, Docker, Conventions — COMPLETE **Completed:** 2026-03-04 @@ -253,10 +255,13 @@ Deliver the core wizard that converts parameterized SQL into deployable versione ## Phase 5: Module 4 - Scheduling + Snapshot Cache End-to-End — COMPLETE -**Completed:** 2026-03-06 +**Initial completion:** 2026-03-06 + +**Scheduler/parameter hardening updated:** 2026-08-31 ### Goals -Enable scheduled data refresh with persistent jobs and cached response serving. +Enable scheduled data refresh with persisted definitions, restart restoration, and filtered cached +response serving. ### Delivered @@ -270,7 +275,8 @@ Enable scheduled data refresh with persistent jobs and cached response serving. | `backend/app/services/scheduler.py` | APScheduler 3.x AsyncIOScheduler integration: lifecycle (start/stop), job execution (Oracle query + snapshot persistence), job management (add/remove/pause/resume) | | `backend/app/services/schedule.py` | Business logic: CRUD, one-schedule-per-endpoint uniqueness, run-now, pause/resume, job run queries, snapshot queries | | `backend/app/routers/schedules.py` | REST CRUD + `/run` + `/pause` + `/resume` + job runs + snapshots under `/api/v1/admin/schedules/*` | -| `backend/app/routers/data.py` | **Updated**: Snapshot mode serving — snapshot-strategy endpoints return cached JSONB with freshness metadata | +| `backend/app/services/data.py` | Snapshot orchestration: required-parameter validation, retained-run coverage selection, typed cached-row filtering, and stable error responses | +| `backend/app/services/snapshot_filtering.py` | Compiled `eq`/`gte`/`lte` filters, coverage checks, range validation, cached date normalization, and row filtering | | `backend/app/main.py` | **Updated**: APScheduler lifecycle (start on startup, stop on shutdown), schedule router registration | | `backend/tests/test_schedules.py` | Schema validation unit tests + API integration tests | @@ -293,21 +299,37 @@ Enable scheduled data refresh with persistent jobs and cached response serving. #### Scheduling Features - APScheduler 3.x with AsyncIOScheduler integration - Cron (5-field) and interval (seconds) schedule types +- Friendly hourly/daily/weekly/monthly cron controls plus advanced custom cron +- IANA timezone, logical run date, and inclusive reusable calendar windows +- Explicit schedule-owned sources for every SQL bind: literal, SQL `NULL`, run date, relative + date, window start, or window end +- Preview of upcoming nominal runs, logical dates, windows, and resolved typed parameters - Job coalescing (max 1 instance per job, 60s misfire grace time) -- Scheduler lifecycle tied to FastAPI startup/shutdown +- Scheduler lifecycle tied to FastAPI startup/shutdown; active database definitions are restored on + startup in the documented single-process topology - One schedule per endpoint uniqueness constraint #### Snapshot Cache Features - JSONB snapshot storage in PostgreSQL - Automatic snapshot retention (keeps latest 5 per endpoint) -- Snapshot-mode data endpoints serve cached results with `snapshot_created_at` metadata +- Parameterized snapshot endpoints require a final cached-column mapping and allowlisted operator + for every request parameter +- Data requests enforce required parameters, choose the newest retained snapshot whose job-run + inputs cover the complete request, and return only typed filtered rows +- Covered requests with no matching rows return `data: []`; invalid ranges or requests outside + retained coverage return HTTP 422 +- Snapshot responses include filtered `row_count`, original `snapshot_row_count`, and + `snapshot_created_at` metadata - Fallback: 503 if no snapshot available yet - Column mapping applied during job execution #### Job Execution -- Immutable job run audit records (started_at, finished_at, status, row_count, error_detail) +- Immutable job run audit records including scheduled/logical time, window boundaries, resolved + parameters, trigger source, binding hash, row count, status, and error detail - Status tracking: running → success/failed/timeout -- Default parameter values used for scheduled queries +- Schedule-owned bindings are resolved for scheduled queries; endpoint request defaults are never + read by the scheduler +- `(schedule_id, scheduled_for)` uniqueness prevents duplicate logical runs - Structured logging with job_id, run_id, row_count, duration_ms, success fields #### Frontend @@ -317,11 +339,13 @@ Enable scheduled data refresh with persistent jobs and cached response serving. | `frontend/src/lib/api.ts` | Axios-based `schedulesApi` client (CRUD + run/pause/resume + jobs + snapshots) | | `frontend/src/lib/queryClient.ts` | Schedule query key factories | | `frontend/src/pages/SchedulesPage.tsx` | List table + create/delete dialogs + run now/pause/resume controls + job runs viewer | +| `frontend/src/components/schedules/ScheduleParameterBindings.tsx` | Schedule-local binding sources, calendar window controls, and resolved-run preview | +| `frontend/src/components/endpoints/SnapshotFilterMappings.tsx` | Required create/edit mappings from request parameters to cached output columns | | `frontend/src/pages/DashboardPage.tsx` | Updated with schedules count card (4-column grid) | | `frontend/src/components/Layout.tsx` | Updated with Schedules nav item, version bumped to v0.5.0 | | `frontend/src/App.tsx` | Updated with /schedules route | -#### Checks +#### Phase completion checks (2026-03-06) - `ruff check .` — clean - `mypy .` — clean (58 files, 0 errors) - `pytest -k "not integration"` — 69 passed @@ -368,7 +392,7 @@ Provide centralized operational controls and health visibility. |-----|------|---------|-----------------| | `log_level` | enum (DEBUG/INFO/WARNING/ERROR) | INFO | Yes | | `query_timeout_seconds` | integer (1-300) | 30 | No | -| `cors_origins` | string (comma-separated) | * | Yes | +| `cors_origins` | string (comma-separated) | `http://localhost:5173` | Yes (wildcard rejected) | | `snapshot_retention_count` | integer (1-100) | 5 | No | | `max_job_concurrency` | integer (1-20) | 3 | Yes | diff --git a/docs/project_plan.md b/docs/project_plan.md index 34f7727..0d0f10f 100644 --- a/docs/project_plan.md +++ b/docs/project_plan.md @@ -1,6 +1,10 @@ # QueryGateway Implementation Plan -QueryGateway will be delivered as a self-hosted platform that lets teams expose Oracle query results as secure REST endpoints through a guided wizard, while preserving security-by-default, reproducible infrastructure, and operational reliability from day one. This plan translates the approved architecture and recommendations into phased, testable implementation work using Python 3.12+, FastAPI, PostgreSQL, SQLAlchemy 2.0 + Alembic, APScheduler 3.x, Pydantic Settings, API versioning, Vite + React SPA, shadcn/ui + Tailwind, and Dockerized CI/CD. +QueryGateway is delivered as a self-hosted platform that lets teams expose Oracle query results as +secure REST endpoints through a guided wizard while preserving security-by-default, reproducible +infrastructure, and operational reliability. The implemented baseline uses Python 3.14+, FastAPI, +PostgreSQL, SQLAlchemy 2.0 + Alembic, APScheduler 3.x, Pydantic Settings, API versioning, a Vite + +React SPA, shadcn/ui + Tailwind, and Dockerized CI/CD. ## Scope and Non-Goals @@ -47,7 +51,8 @@ QueryGateway will be delivered as a self-hosted platform that lets teams expose - FastAPI service responsibilities: - Validate auth for every data endpoint request. - Resolve endpoint metadata and selected freshness strategy. - - Execute parameterized SQL against Oracle (live mode) or return latest snapshot (scheduled mode). + - Execute parameterized SQL against Oracle (live mode) or select a covering retained snapshot + and apply typed cached-row filters (snapshot mode). - Manage scheduler jobs and persist execution telemetry. ### Versioning and Deprecation Rules @@ -74,7 +79,7 @@ QueryGateway will be delivered as a self-hosted platform that lets teams expose - Shared coding standards and contribution conventions. ### Key Implementation Tasks -- Initialize backend Python project for 3.12+ with dependency management and lint/test/type scripts. +- Initialize backend Python project for 3.14+ with dependency management and lint/test/type scripts. - Initialize frontend Vite + React + TypeScript with Tailwind, shadcn/ui base setup, eslint/prettier, and vitest (or Jest). - Add root-level tooling docs and Makefile/task runner commands for consistent local workflows. - Create `docker-compose.yml` with `api`, `web`, `db` and optional `oracle` profile. @@ -212,7 +217,7 @@ QueryGateway will be delivered as a self-hosted platform that lets teams expose ### Deliverables - Multi-step wizard UI and backend orchestration APIs. -- Rich SQL editor integration (Monaco via `@monaco-editor/react` or CodeMirror 6 via `@uiw/react-codemirror`). +- Rich SQL editor integration with CodeMirror 6 via `@uiw/react-codemirror`. - SQL preview engine with bind parameter extraction/validation rules. - Endpoint registration pipeline and dynamic data router. @@ -223,6 +228,9 @@ QueryGateway will be delivered as a self-hosted platform that lets teams expose - Only named bind variables (e.g., `:param_name`). - Reject string-concatenated SQL interpolation patterns. - Validate input types/coercion using explicit schemas. + - Treat preview samples as temporary values, never persisted defaults. + - Enforce required descriptors for both live and snapshot requests regardless of defaults. + - Keep optional live defaults independent from schedule-owned SQL bindings. - Build preview execution endpoint to return sample JSON and inferred schema. - Implement result mapping layer (rename/select columns and output shaping). - Build endpoint definition step (path, GET method for MVP, auth method, data strategy). @@ -247,10 +255,12 @@ QueryGateway will be delivered as a self-hosted platform that lets teams expose ## Phase 5: Module 4 - Scheduling + Snapshot Cache End-to-End ### Goals -- Enable scheduled data refresh with persistent jobs and cached response serving. +- Enable scheduled data refresh with persisted definitions, restart restoration, and filtered + cached response serving. ### Deliverables -- APScheduler 3.x integration with PostgreSQL-backed job store. +- APScheduler 3.x integration with an in-memory job store whose active schedule definitions are + persisted in PostgreSQL and restored on API startup. - Schedule CRUD and control actions (run now, pause/resume, enable/disable). - Snapshot cache storage in PostgreSQL JSONB. - Job execution logging dashboard and APIs. @@ -260,14 +270,18 @@ QueryGateway will be delivered as a self-hosted platform that lets teams expose - Model schedules, job state, run history, and snapshot payload metadata. - Implement admin APIs under `/api/v1/admin/schedules/*` and `/api/v1/admin/jobs/*`. - Create execution worker logic for scheduled query runs and cache replacement strategy. +- Require schedule-owned parameter bindings, timezone-aware logical dates, and declarative + calendar windows; persist resolved values for audit and coverage selection. - Add staleness tracking and fallback handling when latest snapshot fails. -- Wire data endpoints in snapshot mode to read from cache with freshness metadata. +- Wire data endpoints in snapshot mode to require cached-column mappings, select the newest + retained run that covers the request, and apply typed equality/range filters. - Build frontend schedule management UI and execution log views. ### Acceptance Criteria -- Scheduled jobs persist across process restarts. +- Active scheduled jobs are restored exactly once after a single API process restarts. - Manual run and pause/resume actions work reliably from UI and API. -- Snapshot-mode endpoints return cached JSON and expose last-refresh timestamps. +- Snapshot-mode endpoints enforce required request parameters, reject requests outside retained + coverage, return only filtered cached rows, and expose last-refresh metadata. - Job logs capture start time, duration, row count, status, and error details. ### Dependencies @@ -359,7 +373,7 @@ QueryGateway will be delivered as a self-hosted platform that lets teams expose GitHub Actions workflows will run on pull requests and protected branches with separate jobs: - Backend job - - Setup Python 3.12+ + - Setup Python 3.14+ - Install backend dependencies - Run `ruff`, `mypy`, `pytest` - Frontend job diff --git a/docs/scheduler_parameter_bindings.md b/docs/scheduler_parameter_bindings.md index df0af9e..45f8f79 100644 --- a/docs/scheduler_parameter_bindings.md +++ b/docs/scheduler_parameter_bindings.md @@ -1,6 +1,41 @@ -# Scheduler parameter bindings +# Endpoint, scheduler, and snapshot parameter contracts + +QueryGateway treats the same SQL bind name differently depending on where it is used. Keeping these contexts separate prevents a preview value or request default from silently changing a recurring snapshot window. + +| Context | Value source | Persisted? | Purpose | +|---|---|---|---| +| SQL preview | Temporary sample entered in the endpoint wizard | No | Execute a safe sample query and discover output columns | +| Live data request | Authenticated HTTP query string, with optional endpoint default | Endpoint default only | Bind and execute the Oracle query for this request | +| Snapshot schedule | Schedule-owned declarative parameter binding | Yes, on the schedule | Decide which rows Oracle loads into each snapshot | +| Snapshot data request | Authenticated HTTP query string plus endpoint `snapshot_filter` mappings | Mapping only | Select a covering retained snapshot and filter its cached rows | + +## SQL preview and live request rules + +- A bind placeholder must be outside SQL string quotes: use `store_id = :store_id`, not + `store_id = ':store_id'`. Text inside single quotes is a SQL literal and is not detected as a + bind parameter. +- Every parameter marked `required` must be supplied by both live and snapshot callers. A stored + endpoint default never weakens that request contract. +- An omitted optional live parameter may use a typed literal default, explicit SQL `NULL`, or the + dynamic date default `today` or `yesterday`. An optional parameter with no default resolves to + `NULL`. +- Date requests accept `YYYY-MM-DD` and `DD-MM-YYYY` and normalize to a Python `date` before + binding. Boolean requests accept `true`, `false`, `1`, `0`, `yes`, or `no`. +- Preview sample values are request-local and never become endpoint defaults or schedule + bindings. + +For example, a required date range and optional store filter are supplied as ordinary query +parameters: -Snapshot schedules own their SQL bind values. Endpoint defaults remain useful for live requests and SQL preview, but scheduled execution never reads them. This separation prevents a request default from silently changing a recurring data window. +```text +GET /api/v1/data/store-orders?start_date=2026-08-01&end_date=31-08-2026&store_id=10 +``` + +Invalid or missing declared values return HTTP 422 before Oracle execution or snapshot lookup. + +## Schedule-owned parameter bindings + +Snapshot schedules own their SQL bind values. Scheduled execution never reads endpoint defaults. ## Binding sources @@ -60,11 +95,104 @@ Before creating it, send the same timing and binding fields to: POST /api/v1/admin/schedules/preview ``` -The response contains the next one to ten nominal fire times (three by default), logical dates, window boundaries, and resolved typed parameters. +The response contains the next one to ten nominal fire times (three by default), logical dates, +window boundaries, and resolved typed parameters. + +## Snapshot request filters and coverage + +Schedule bindings decide which rows Oracle loads into a snapshot. Authenticated data-request +parameters decide which rows are returned from that cache. Every parameterized snapshot endpoint +must explicitly map each request parameter to a cached output column after `column_map` renaming +and one whitelisted comparison: + +| Operator | Row selection | Coverage requirement | +|---|---|---| +| `eq` | Cached column equals the request value | Scheduled value equals the request value | +| `gte` | Cached column is greater than or equal to the request value | Scheduled lower bound is less than or equal to the request lower bound | +| `lte` | Cached column is less than or equal to the request value | Scheduled upper bound is greater than or equal to the request upper bound | + +Example endpoint parameter schema: + +```json +{ + "start_date": { + "type": "date", + "required": true, + "snapshot_filter": { "column": "business_date", "operator": "gte" } + }, + "end_date": { + "type": "date", + "required": true, + "snapshot_filter": { "column": "business_date", "operator": "lte" } + }, + "store_id": { + "type": "integer", + "required": false, + "default_is_null": true, + "snapshot_filter": { + "column": "store_id", + "operator": "eq", + "null_means_all": true + } + } +} +``` + +`null_means_all` is valid only for `eq` on an optional parameter. It means a schedule run whose +resolved value is SQL `NULL` covers every requested value for that parameter. Whenever an optional +request parameter is omitted, the selector requires a retained run where that parameter also +resolved to SQL `NULL`; a fixed-value snapshot is only a subset and cannot satisfy the request. +Without `null_means_all`, that NULL-resolved snapshot represents the query's exact NULL semantics +rather than all possible values. This flag does not make a required HTTP parameter optional. A +parameter such as `store_id` is an ordinary row filter here; it is not tenant authorization. +Authentication and authorization remain the responsibility of the endpoint's configured auth +method. + +The data plane checks retained snapshots newest first and selects the newest snapshot whose +persisted job-run parameters cover the complete request. It then applies every configured mapping +to the cached rows using the parameter's declared type. Date columns containing Oracle +DATE/TIMESTAMP ISO strings are normalized to dates before comparison. Behavior is explicit: + +- Missing or invalid required parameters return HTTP 422 with the field in `detail`. +- Lower and upper mappings for the same cached column must declare the same parameter type. +- When multiple parameters provide bounds for the same cached column, validation retains every + value and uses `max(gte)` as the effective lower bound and `min(lte)` as the effective upper + bound. If the effective lower bound is greater, the request returns HTTP 422 with + `code=invalid_parameter_range`. For example, two `gte` values of `5` and `10` plus an `lte` + value of `7` resolve to the invalid effective range `10..7`. +- A valid request outside all retained coverage returns HTTP 422 with `code=snapshot_out_of_coverage`. +- A request inside coverage with no matching business rows returns HTTP 200 with `data: []`. +- An endpoint created before this contract without complete mappings returns HTTP 422 with `code=snapshot_filter_not_configured`; add mappings through the endpoint edit dialog or admin update API. +- A mapping that does not exist in a non-empty cached row returns HTTP 422 with `code=snapshot_filter_column_unavailable`. +- No retained snapshot still returns HTTP 503. + +Representative stable error bodies are: + +```json +{ + "code": "snapshot_out_of_coverage", + "detail": "Requested parameters are outside retained snapshot coverage." +} +``` + +```json +{ + "code": "invalid_parameter_range", + "detail": "Snapshot filter lower bound must not exceed its upper bound." +} +``` + +The mapping is stored inside the endpoint's existing JSON parameter schema, so it requires no +relational database migration. Schedule-owned bindings, timezone, logical-run fields, resolved +parameter audit data, and the `(schedule_id, scheduled_for)` idempotency constraint were added by +Alembic revision `e4a6c2d9f801`. ## Logical time, retries, and manual runs -Cron expressions and run dates are evaluated in the schedule's IANA timezone. Normal execution uses the persisted nominal `next_run_at` as `scheduled_for`, not the wall-clock time at which a delayed job starts. Job runs store the resolved context and a binding-configuration hash. A unique `(schedule_id, scheduled_for)` key prevents duplicate execution of the same logical run. +Cron expressions and run dates are evaluated in the schedule's IANA timezone. Normal execution +uses the persisted nominal `next_run_at` as `scheduled_for`, not the wall-clock time at which a +delayed job starts. Job runs store the resolved context and a binding-configuration hash. A unique +`(schedule_id, scheduled_for)` key prevents duplicate execution of the same logical run. `Run now` derives its logical date from the current time by default. An administrator can replay a particular business date without editing the schedule: @@ -80,7 +208,10 @@ The supplied date is interpreted at midnight in the schedule timezone and is rec ## Endpoint changes -An attached schedule and its endpoint form one validated configuration. While a schedule exists, endpoint updates must keep the endpoint in snapshot mode and preserve SQL/parameter names and types that the stored bindings can resolve. Incompatible updates return HTTP 422 without changing the endpoint. Delete or update the schedule first when intentionally changing that contract. +An attached schedule and its endpoint form one validated configuration. While a schedule exists, +endpoint updates must keep the endpoint in snapshot mode and preserve SQL/parameter names and +types that the stored bindings can resolve. Incompatible updates return HTTP 422 without changing +the endpoint. Delete or update the schedule first when intentionally changing that contract. ## Existing schedules @@ -91,4 +222,5 @@ The Alembic migration converts existing endpoint defaults into schedule-local bi - Explicit SQL `NULL` remains `null`. - Fixed defaults become `literal`. -A legacy schedule parameter with no resolvable default is left unbound and fails clearly until an administrator selects a schedule source. The migration never invents a value. +A legacy schedule parameter with no resolvable default is left unbound and fails clearly until an +administrator selects a schedule source. The migration never invents a value. diff --git a/docs/security_checklist.md b/docs/security_checklist.md index 5e9fb1c..5c4f179 100644 --- a/docs/security_checklist.md +++ b/docs/security_checklist.md @@ -40,7 +40,7 @@ Comprehensive security validation for QueryGateway production deployments. All i | 22 | Template interpolation rejected (`${`, `{var}`) | Verified | Pattern included in safety validation | | 23 | SQL validation runs on both create and update | Verified | `EndpointCreate` and `EndpointUpdate` share the validator | | 24 | SQL preview also validates before execution | Verified | `SqlPreviewRequest` includes the same validator | -| 25 | Parameters coerced through typed schemas before SQL execution | Verified | `_coerce_param()` with `ParamDescriptor` type validation | +| 25 | Parameters coerced through typed schemas before SQL execution | Verified | `build_param_model()` creates the Pydantic request model; `DataService` enforces required fields before binding | | 26 | SQLAlchemy `text()` with bind dict used for execution | Verified | `sql/executor.py` uses parameterized execution | ## Input Validation @@ -61,7 +61,7 @@ Comprehensive security validation for QueryGateway production deployments. All i | 33 | 500 errors return generic message, not stack traces | Verified | `unhandled_exception_handler` returns `{"detail": "Internal server error"}` | | 34 | Validation errors do not leak internal schema details | Verified | Pydantic errors are structured but safe | | 35 | Oracle connection strings not exposed in error responses | Verified | SQL execution errors logged server-side; generic message returned | -| 36 | Access logs record all data endpoint requests | Verified | `_write_access_log()` called on every code path in data router | +| 36 | Access logs record all data endpoint requests | Verified | The data router wraps `DataService.serve()` in the access-log context manager for success and failure paths | | 37 | Access logs include: path, method, principal, IP, status, duration, request_id | Verified | `AccessLog` model captures all audit fields | ## Transport Security @@ -116,10 +116,21 @@ Comprehensive security validation for QueryGateway production deployments. All i | 63 | Branch protection with required, non-bypassable checks | Action Required | Repo-admin only — settings documented in [repository_governance.md](repository_governance.md) | | 64 | Container image signing + provenance | Action Required | Images scanned + SBOM'd, not yet signed (cosign/Sigstore) | +## Snapshot Parameter and Data Isolation + +| # | Check | Status | Notes | +|---|-------|--------|-------| +| 65 | Required endpoint parameters are enforced for live and snapshot HTTP requests | Verified | `DataService._coerce_params()` uses `build_param_model(..., enforce_required=True)`; endpoint defaults cannot weaken required fields | +| 66 | Scheduled queries never inherit endpoint request defaults | Verified | Schedule creation requires exactly one validated binding per SQL parameter; execution resolves only `parameter_bindings_json` | +| 67 | Every parameterized snapshot request has an allowlisted cached-column mapping | Verified | Endpoint create/update requires a post-rename column plus `eq`, `gte`, or `lte` for every parameter | +| 68 | Cached data is served only after retained-run coverage and typed row filtering | Verified | `DataService._serve_snapshot()` selects the newest covering job run, validates mapped columns, and calls `filter_snapshot_rows()`; no unfiltered fallback exists | +| 69 | Range and SQL `NULL` coverage semantics fail closed | Verified | Reversed bounds return 422; range types must match; `null_means_all` is limited to optional equality mappings | +| 70 | Snapshot filter mappings cannot substitute for authorization | Verified | Data authentication runs before snapshot selection; `store_id` and other mapped values are ordinary row filters only | + ## Summary -- **Verified items**: 58/64 -- **Action required**: 6/64 +- **Verified items**: 64/70 +- **Action required**: 6/70 - **High-severity unresolved findings**: 0 - **All code-level security controls validated through automated tests** diff --git a/docs/v0.2.0-security-performance-plan.md b/docs/v0.2.0-security-performance-plan.md index eaba77a..645945a 100644 --- a/docs/v0.2.0-security-performance-plan.md +++ b/docs/v0.2.0-security-performance-plan.md @@ -1,5 +1,12 @@ # QueryGateway v0.2.0 — Security & Performance Hardening Plan +> **Planning baseline, not the current runtime contract.** Several findings below have since been +> implemented, including active schedule restoration from persisted definitions. The current +> scheduler also owns explicit logical parameter bindings, and parameterized snapshot requests use +> coverage-aware typed cached-row filtering. See [Architecture](architecture.md), +> [Implementation progress](progress.md), and [Parameter contracts](scheduler_parameter_bindings.md) +> before acting on a historical finding or file/line reference in this plan. + ## Goal Ship v0.2.0: close the remaining security gaps (rate limiting, scheduler persistence, edge hardening, audit retention) and the biggest performance costs (per-request Oracle connection setup, per-request bcrypt/API-key cost, per-request model rebuilds, single-bundle frontend), plus small aligned features. Multi-database support (docs/phase2.md) is explicitly deferred to v0.3.0. @@ -24,7 +31,7 @@ Ship v0.2.0: close the remaining security gaps (rate limiting, scheduler persist | S3 | High | Scheduler jobs lost on restart: APScheduler uses the in-memory job store and `start_scheduler()` never re-registers rows from `schedules`; after any restart/deploy all snapshot refreshes silently stop until an admin edits each schedule | `backend/app/services/scheduler.py:210-228`, `backend/app/services/schedule.py:106-113` (jobs added only on create/update/resume) | | S4 | Medium | nginx edge lacks security headers (HSTS, X-Content-Type-Options, X-Frame-Options, CSP, Referrer-Policy), compression, explicit body-size cap, and `Cache-Control` for API responses; no TLS enforcement story in compose | `docker/nginx.conf` (29 lines, proxy + SPA fallback only); production_deployment.md gap #4 | | S5 | Medium | No retention for `access_logs` and `job_runs` — unbounded table growth (snapshots already have `snapshot_retention_count`) | `backend/app/services/scheduler.py:137-148` (retention exists for snapshots only) | -| S6 | Low | No audit-friendly listing of public (`allow_unauthenticated`) endpoints to support the go-live audit in production_deployment.md gap #7 | `backend/app/services/data.py:113-134` | +| S6 | Low | No audit-friendly listing of endpoints using the legacy-named `allow_unauthenticated` platform-admin Bearer fallback; the field does not permit public access | `backend/app/services/data.py` authentication dispatch | ### Performance @@ -58,7 +65,9 @@ Ship v0.2.0: close the remaining security gaps (rate limiting, scheduler persist - Startup test: exactly one APScheduler job per active schedule after restore; a multi-process simulation proves no duplicate job registration/execution. **A3. Public-endpoint audit surface (S6)** -- Add list of active endpoints with `allow_unauthenticated=true` (id, path, name, updated_at) to the admin-gated health dashboard response. +- Add a list of active endpoints using `allow_unauthenticated=true` (id, path, name, updated_at) to + the admin-gated health dashboard and label them as platform-admin Bearer fallback endpoints, not + public endpoints. **A4. Audit retention (S5)** - New seeded settings: `access_log_retention_days` (default 90), `job_run_retention_days` (default 30) — follow the `snapshot_retention_count` pattern in `SettingsRepository`. diff --git a/frontend/src/components/endpoints/EndpointWizard.test.tsx b/frontend/src/components/endpoints/EndpointWizard.test.tsx index 6f669b7..1412225 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("allows snapshot review when its schedule will own parameter values", async () => { + it("requires snapshot row mappings before review", async () => { renderWizard(); await advanceToParameters(); @@ -121,6 +121,10 @@ describe("EndpointWizard preview coordination", () => { fireEvent.change(configSelectors[1], { target: { value: "snapshot" } }); fireEvent.click(screen.getByRole("checkbox")); + expect(screen.getByRole("button", { name: "Next" })).toBeDisabled(); + fireEvent.change(screen.getByLabelText("Snapshot column for customer_id"), { + target: { value: "customer_id" }, + }); expect(screen.getByRole("button", { name: "Next" })).toBeEnabled(); expect( screen.getByText(/Scheduled values are configured with the schedule/i), diff --git a/frontend/src/components/endpoints/EndpointWizard.tsx b/frontend/src/components/endpoints/EndpointWizard.tsx index defa9bb..15ab32c 100644 --- a/frontend/src/components/endpoints/EndpointWizard.tsx +++ b/frontend/src/components/endpoints/EndpointWizard.tsx @@ -146,7 +146,12 @@ 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; - return !!state.name.trim() && !!state.path.trim() && authOk; + const snapshotMappingsOk = + state.data_strategy !== "snapshot" || + Object.values(state.param_schema).every( + (descriptor) => !!descriptor.snapshot_filter?.column.trim(), + ); + return !!state.name.trim() && !!state.path.trim() && authOk && snapshotMappingsOk; } return true; }; @@ -218,7 +223,12 @@ export function EndpointWizard({ onSuccess, onCancel }: EndpointWizardProps) { )} {currentStep === "Parameters" && } {currentStep === "Auth & Config" && ( - + )} {currentStep === "Review" && ( diff --git a/frontend/src/components/endpoints/SnapshotFilterMappings.tsx b/frontend/src/components/endpoints/SnapshotFilterMappings.tsx new file mode 100644 index 0000000..5214a99 --- /dev/null +++ b/frontend/src/components/endpoints/SnapshotFilterMappings.tsx @@ -0,0 +1,114 @@ +import { Input } from "@/components/ui/input"; +import { Label } from "@/components/ui/label"; +import { Select } from "@/components/ui/select"; +import type { ParamDescriptor, SnapshotFilter, SnapshotFilterOperator } from "@/types/endpoint"; + +interface SnapshotFilterMappingsProps { + paramSchema: Record; + onChange: (paramSchema: Record) => void; + previewColumns?: string[]; +} + +export function SnapshotFilterMappings({ + paramSchema, + onChange, + previewColumns = [], +}: SnapshotFilterMappingsProps) { + const updateSnapshotFilter = ( + name: string, + patch: Partial & { operator?: SnapshotFilterOperator }, + ) => { + const descriptor = paramSchema[name]; + const current = descriptor.snapshot_filter; + const next: SnapshotFilter = { + column: current?.column ?? "", + operator: current?.operator ?? "eq", + null_means_all: current?.null_means_all ?? false, + ...patch, + }; + if (next.operator !== "eq") next.null_means_all = false; + onChange({ + ...paramSchema, + [name]: { ...descriptor, snapshot_filter: next }, + }); + }; + + if (Object.keys(paramSchema).length === 0) return null; + + return ( +
+
+

Snapshot request filters

+

+ Map every request parameter to a cached output column. These mappings select rows; they do + not grant tenant access. +

+
+ {Object.entries(paramSchema).map(([name, descriptor]) => { + const mapping = descriptor.snapshot_filter; + const columnInputId = `snapshot-column-${name}`; + const columnListId = `snapshot-columns-${name}`; + const operatorInputId = `snapshot-operator-${name}`; + return ( +
+

:{name}

+
+
+ + updateSnapshotFilter(name, { column: event.target.value })} + placeholder="business_date" + /> + + {previewColumns.map((column) => ( + +
+
+ + +
+
+ {!descriptor.required && (mapping?.operator ?? "eq") === "eq" && ( + + )} +
+ ); + })} +
+ ); +} diff --git a/frontend/src/components/endpoints/wizard/ConfigStep.test.tsx b/frontend/src/components/endpoints/wizard/ConfigStep.test.tsx index baf7746..a2a8d4b 100644 --- a/frontend/src/components/endpoints/wizard/ConfigStep.test.tsx +++ b/frontend/src/components/endpoints/wizard/ConfigStep.test.tsx @@ -88,4 +88,80 @@ describe("ConfigStep snapshot scheduling guidance", () => { screen.getByText(/Scheduled values are configured with the schedule/i), ).toBeInTheDocument(); }); + + it("configures an explicit cached column and operator for every parameter", () => { + const update = vi.fn(); + render( + , + ); + + expect(screen.getByText(/Snapshot request filters/i)).toBeInTheDocument(); + fireEvent.change(screen.getByLabelText("Snapshot column for start_date"), { + target: { value: "business_date" }, + }); + + expect(update).toHaveBeenCalledWith({ + param_schema: { + start_date: { + type: "date", + required: true, + default: null, + snapshot_filter: { + column: "business_date", + operator: "eq", + null_means_all: false, + }, + }, + }, + }); + }); + + it("allows an optional equality filter to declare scheduled NULL as all values", () => { + const update = vi.fn(); + render( + , + ); + + fireEvent.click(screen.getByLabelText("Scheduled NULL covers all store_id values")); + + expect(update).toHaveBeenCalledWith({ + param_schema: { + store_id: expect.objectContaining({ + snapshot_filter: { + column: "store_id", + operator: "eq", + null_means_all: true, + }, + }), + }, + }); + }); }); diff --git a/frontend/src/components/endpoints/wizard/ConfigStep.tsx b/frontend/src/components/endpoints/wizard/ConfigStep.tsx index 38ab230..862543b 100644 --- a/frontend/src/components/endpoints/wizard/ConfigStep.tsx +++ b/frontend/src/components/endpoints/wizard/ConfigStep.tsx @@ -6,15 +6,18 @@ import { Textarea } from "@/components/ui/textarea"; import type { AuthMethod } from "@/types/auth_method"; import type { DataStrategy } from "@/types/endpoint"; +import { SnapshotFilterMappings } from "../SnapshotFilterMappings"; + import type { WizardState, WizardUpdate } from "./types"; interface ConfigStepProps { state: WizardState; update: WizardUpdate; authMethods: AuthMethod[]; + previewColumns?: string[]; } -export function ConfigStep({ state, update, authMethods }: ConfigStepProps) { +export function ConfigStep({ state, update, authMethods, previewColumns = [] }: ConfigStepProps) { return (

Endpoint Configuration

@@ -92,13 +95,21 @@ export function ConfigStep({ state, update, authMethods }: ConfigStepProps) {
{state.data_strategy === "snapshot" && ( - - Scheduled values are configured with the schedule - - 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. - - + <> + + Scheduled values are configured with the schedule + + 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. + + + + update({ param_schema: paramSchema })} + /> + )} {!state.auth_method_id && ( diff --git a/frontend/src/components/endpoints/wizard/ReviewStep.tsx b/frontend/src/components/endpoints/wizard/ReviewStep.tsx index 1b7a05a..246d3b9 100644 --- a/frontend/src/components/endpoints/wizard/ReviewStep.tsx +++ b/frontend/src/components/endpoints/wizard/ReviewStep.tsx @@ -39,6 +39,12 @@ export function ReviewStep({ state, connections, authMethods }: ReviewStepProps) :{name} — {descriptor.type}; default:{" "} {describeParameterDefault(descriptor)} + {state.data_strategy === "snapshot" && descriptor.snapshot_filter && ( + <> + ; snapshot: {descriptor.snapshot_filter.column}{" "} + {descriptor.snapshot_filter.operator} + + )} ))} diff --git a/frontend/src/pages/EndpointsPage.test.tsx b/frontend/src/pages/EndpointsPage.test.tsx index a638ced..371b3b0 100644 --- a/frontend/src/pages/EndpointsPage.test.tsx +++ b/frontend/src/pages/EndpointsPage.test.tsx @@ -4,6 +4,7 @@ import { beforeEach, describe, expect, it, vi } from "vitest"; const listMock = vi.fn(); const deleteMock = vi.fn(); +const updateMock = vi.fn(); vi.mock("@/components/endpoints/EndpointWizard", () => ({ EndpointWizard: () =>
Endpoint wizard
, @@ -13,7 +14,7 @@ vi.mock("@/lib/api", () => ({ endpointsApi: { list: (...args: unknown[]) => listMock(...args), delete: (...args: unknown[]) => deleteMock(...args), - update: vi.fn(), + update: (...args: unknown[]) => updateMock(...args), }, getApiError: (error: unknown) => (error instanceof Error ? error.message : String(error)), getPublicApiBaseUrl: () => "http://localhost:8000", @@ -36,6 +37,7 @@ describe("EndpointsPage", () => { beforeEach(() => { listMock.mockReset(); deleteMock.mockReset(); + updateMock.mockReset(); listMock.mockResolvedValue([ { id: "0ff4f30e-6799-48af-9ba7-bb31471de215", @@ -74,4 +76,58 @@ describe("EndpointsPage", () => { expect(await screen.findByRole("alert")).toHaveTextContent("Endpoint could not be deleted."); expect(screen.getByRole("heading", { name: "Delete Endpoint" })).toBeInTheDocument(); }); + + it("lets an existing snapshot endpoint add missing request filter mappings", async () => { + listMock.mockResolvedValueOnce([ + { + id: "0ff4f30e-6799-48af-9ba7-bb31471de215", + name: "store_orders", + description: null, + path: "store-orders", + connection_id: "fd421223-0866-45fd-8937-7a374ba8eb06", + sql_text: "SELECT * FROM orders WHERE store_id = :store_id", + param_schema: { + store_id: { type: "integer", required: true }, + }, + column_map: {}, + auth_method_id: null, + allow_unauthenticated: true, + data_strategy: "snapshot", + version: "v1", + is_active: true, + is_deprecated: false, + deprecation_note: null, + created_at: "2026-08-30T00:00:00Z", + updated_at: "2026-08-30T00:00:00Z", + }, + ]); + updateMock.mockResolvedValue({}); + renderPage(); + + expect(await screen.findByText("store_orders")).toBeInTheDocument(); + fireEvent.click(screen.getByTitle("Edit endpoint")); + expect(screen.getByRole("button", { name: "Save" })).toBeDisabled(); + fireEvent.change(screen.getByLabelText("Snapshot column for store_id"), { + target: { value: "store_id" }, + }); + expect(screen.getByRole("button", { name: "Save" })).toBeEnabled(); + fireEvent.click(screen.getByRole("button", { name: "Save" })); + + await waitFor(() => + expect(updateMock).toHaveBeenCalledWith( + "0ff4f30e-6799-48af-9ba7-bb31471de215", + expect.objectContaining({ + param_schema: { + store_id: expect.objectContaining({ + snapshot_filter: { + column: "store_id", + operator: "eq", + null_means_all: false, + }, + }), + }, + }), + ), + ); + }); }); diff --git a/frontend/src/pages/EndpointsPage.tsx b/frontend/src/pages/EndpointsPage.tsx index 06fb914..8503e92 100644 --- a/frontend/src/pages/EndpointsPage.tsx +++ b/frontend/src/pages/EndpointsPage.tsx @@ -17,6 +17,7 @@ import { Input } from "@/components/ui/input"; import { Label } from "@/components/ui/label"; import { Textarea } from "@/components/ui/textarea"; import { EndpointWizard } from "@/components/endpoints/EndpointWizard"; +import { SnapshotFilterMappings } from "@/components/endpoints/SnapshotFilterMappings"; import { endpointsApi, getApiError, getPublicApiBaseUrl } from "@/lib/api"; import { queryKeys } from "@/lib/queryClient"; import type { Endpoint, EndpointUpdate } from "@/types/endpoint"; @@ -70,11 +71,18 @@ export function EndpointsPage() { is_active: ep.is_active, is_deprecated: ep.is_deprecated, deprecation_note: ep.deprecation_note ?? "", + param_schema: ep.param_schema, }); setEditError(""); }; const baseUrl = getPublicApiBaseUrl(); + const editParamSchema = editForm.param_schema ?? editEndpoint?.param_schema ?? {}; + const editSnapshotMappingsOk = + editEndpoint?.data_strategy !== "snapshot" || + Object.values(editParamSchema).every( + (descriptor) => !!descriptor.snapshot_filter?.column.trim(), + ); if (showWizard) { return ( @@ -190,7 +198,12 @@ export function EndpointsPage() { > - diff --git a/frontend/src/types/endpoint.ts b/frontend/src/types/endpoint.ts index 01ea5f6..544dbf5 100644 --- a/frontend/src/types/endpoint.ts +++ b/frontend/src/types/endpoint.ts @@ -2,6 +2,13 @@ export type DataStrategy = "live" | "snapshot"; export type DateDefaultExpression = "today" | "yesterday"; +export type SnapshotFilterOperator = "eq" | "gte" | "lte"; + +export interface SnapshotFilter { + column: string; + operator: SnapshotFilterOperator; + null_means_all?: boolean; +} export interface ParamDescriptor { type: "string" | "integer" | "float" | "date" | "boolean"; @@ -10,6 +17,7 @@ export interface ParamDescriptor { default_is_null?: boolean; default_expression?: DateDefaultExpression | null; description?: string; + snapshot_filter?: SnapshotFilter | null; } export interface Endpoint {