diff --git a/README.md b/README.md new file mode 100644 index 0000000..265a52b --- /dev/null +++ b/README.md @@ -0,0 +1,561 @@ +# Callora Backend + +API gateway, usage metering, and billing services for the Callora API marketplace. Talks to Soroban contracts and Horizon for on-chain settlement. + +## Logs Endpoint + +Authenticated users can submit and retrieve structured log entries via `/api/logs`. + +- `GET /api/logs` — Retrieve all log entries for the authenticated user, sorted newest-first. + Returns `{ data: { logs: [...], meta: { total } } }` in the canonical success envelope. +- `POST /api/logs` — Submit a new log entry. + Body: `{ "message": string, "level"?: "debug"|"info"|"warn"|"error", "meta"?: object }`. + Returns `201` with the created entry. + +### Rate Limiting + +Both endpoints are protected by a **per-user token-bucket** rate limiter. + +| Variable | Default | Description | +|---|---|---| +| `LOGS_RATE_LIMIT_CAPACITY` | `60` | Burst ceiling (tokens per user) | +| `LOGS_RATE_LIMIT_REFILL_RATE` | `1` | Tokens refilled per second | + +When the bucket empties the server immediately responds with: + +``` +HTTP/1.1 429 Too Many Requests +Retry-After: + +{ + "code": "TOO_MANY_REQUESTS", + "message": "Too Many Requests", + "requestId": "...", + "retryAfterMs": 750 +} +``` + +Key resolution (priority order): +1. Authenticated user ID from the JWT `Authorization: Bearer` header. +2. `x-user-id` header (trusted internal/test header). +3. Client IP address (unauthenticated fallback). + +## API Catalog Pagination (`GET /api/apis`) + +The public API catalog endpoint uses **keyset cursor pagination** over `(created_at DESC, id DESC)` for stable, gap-free traversal under concurrent writes. Offset-based pagination has been removed; all requests now return cursor-based responses. + +Results are ordered **newest-first** by `(created_at DESC, id DESC)`. Pass the opaque `nextCursor` value returned in one response as the `cursor` query parameter on the next request. Omit `cursor` for the first page. + +| Parameter | Type | Description | +|-----------|------|-------------| +| `cursor` | string | Opaque base64 keyset cursor (from `meta.nextCursor`). Omit for the first page. | +| `limit` | integer 1–100 | Page size. Defaults to 20. | +| `category` | string | Optional category filter. | +| `search` | string | Optional name substring filter. | + +**Request — first page:** +``` +GET /api/apis?limit=2 +``` +```json +{ + "data": [ { "id": 5, ... }, { "id": 4, ... } ], + "meta": { + "limit": 2, + "hasMore": true, + "nextCursor": "MjAyNC0wMS0wNFQwMDowMDowMC4wMDBafDQ=" + } +} +``` + +**Request — subsequent page:** +``` +GET /api/apis?limit=2&cursor=MjAyNC0wMS0wNFQwMDowMDowMC4wMDBafDQ= +``` +```json +{ + "data": [ { "id": 3, ... }, { "id": 2, ... } ], + "meta": { + "limit": 2, + "hasMore": true, + "nextCursor": "MjAyNC0wMS0wMlQwMDowMDowMC4wMDBafDI=" + } +} +``` + +When `hasMore` is `false` and `nextCursor` is absent, you have reached the last page. + +A malformed or tampered cursor returns `HTTP 400` with `code: "VALIDATION_ERROR"`. + +The `offset` and `page` query parameters are ignored (cursor pagination does not support random-access jumping). + +## Fee Abstraction + +Developers can pay Stellar transaction fees using app tokens. The backend wraps their inner transaction in a Stellar fee-bump envelope signed by the platform fee account. + +- `POST /api/billing/fee-abstraction/quote` – returns estimated XLM fee and app-token equivalent. +- `POST /api/billing/fee-abstraction` – accepts app-token payment reference and returns a signed fee-bump XDR. + +Requires `FEE_BUMPER_SECRET_KEY` environment variable (Stellar secret key `S...`). + +See [docs/fee-abstraction.md](./docs/fee-abstraction.md) for full API reference, security considerations, rate limits, and emitted events. + +## Subscription Endpoints + +Authenticated users can subscribe to marketplace APIs with optional metering preferences. + +- `POST /api/subscriptions` — subscribe to an API (`api_id` required; optional `metering_limit` as max calls/month; optional `retry_policy` to override webhook retry behaviour) +- `GET /api/subscriptions` — list subscriptions for the authenticated user; filter by `?status=active|paused|cancelled` +- `GET /api/subscriptions/:id` — get a single subscription (must belong to the authenticated user) +- `PATCH /api/subscriptions/:id` — update `status` (`active`/`paused`), `metering_limit`, or `retry_policy`; body must include at least one field; pass `retry_policy: null` to revert to the platform default +- `DELETE /api/subscriptions/:id` — cancel a subscription (soft-delete; sets status to `cancelled`) + +Business rules: +- A user cannot subscribe to their own API (returns `403`). +- Only one non-cancelled subscription is allowed per user/API pair (returns `409` on conflict). +- Soft-deleted (deleted) APIs cannot be subscribed to (returns `404`). +- Cancelled subscriptions cannot be modified or re-cancelled (returns `400`). + +**Per-subscription webhook retry policy** (`retry_policy`): +An optional `{ maxRetries?: 0–10, baseDelayMs?: 100–60000 }` object that overrides the platform default retry behaviour for webhook deliveries. Omitted fields fall back to platform defaults (`maxRetries: 5`, `baseDelayMs: 1000 ms`). Pass `null` to clear the override. Stored as a JSON text column in the `subscriptions` table. See [docs/webhook-retry-override.md](./docs/webhook-retry-override.md) for full details. + +The migration is in `migrations/0018_subscriptions.sql`; the retry policy column is added by `migrations/0020_subscription_retry_policy.sql`. + +## Dispute Resolution Endpoints + +Developers can open and track disputes against failed or incorrect billing deductions. Admins review and resolve disputes. + +**Developer routes** (`requireAuth`): + +- `POST /api/billing/disputes` — open a new dispute (`usage_event_id` and `reason` required); returns `201` with the new dispute object. Returns `409` if a dispute for that `usage_event_id` already exists. +- `GET /api/billing/disputes` — list all disputes opened by the authenticated developer. +- `GET /api/billing/disputes/:id` — get a single dispute plus its full audit-event trail. Returns `403` if the dispute belongs to another developer, `404` if not found. + +**Admin routes** (`adminAuth`): + +- `GET /api/billing/disputes/admin/all` — list every dispute across all developers. +- `POST /api/billing/disputes/:id/resolve` — resolve a dispute. Body: `{ "resolution": "REFUNDED" | "UPHELD", "notes"?: string }`. Returns `404` for unknown disputes, `409` if already resolved. + +**State machine**: `OPEN → REFUNDED` (admin grants refund) or `OPEN → UPHELD` (admin upholds the charge). + +Every state transition is appended to the `dispute_events` audit trail, which is returned alongside the dispute on `GET /api/billing/disputes/:id`. + +The migration is in `migrations/0019_disputes.sql` (rollback: `migrations/0019_disputes.down.sql`). + +## Developer Profile Endpoints + +- `GET /api/developers/me` returns the authenticated developer profile and auto-creates a blank profile row on first access. +- `PATCH /api/developers/me` updates profile fields for the authenticated developer. +- PATCH validation enforces a valid `website` URL and a supported `category` enum value. + +## Tech stack + +- **Node.js** + **TypeScript** +- **Express** for HTTP API +- **Stellar SDK** for Horizon integration +- **Circuit Breaker & Retry Patterns** for resilience +- Planned: Horizon listener, PostgreSQL, billing engine + +## What's included + +- Health check: `GET /api/health` +- Marketplace routes: + - `GET /api/apis` — list public (active, non-deleted) APIs with cursor pagination over `(created_at, id)` + - `GET /api/apis/:id` + - `POST /api/apis` for authenticated developers to register an API with priced endpoints +- Usage route: `GET /api/usage` +- Top-N endpoints per developer: `GET /api/usage/by-endpoint` — returns the authenticated developer's most-called endpoints ranked by call volume, filterable by `from`/`to`/`apiId`/`limit` (see [docs/usage-by-endpoint.md](./docs/usage-by-endpoint.md)) +- Hourly usage aggregation: `GET /api/usage/aggregate` — returns per-hour call counts and revenue for the authenticated developer, optionally filtered by `from`/`to`/`apiId`; defaults to the last 24 hours when dates are omitted (see [docs/usage-aggregate.md](./docs/usage-aggregate.md)) +- Live usage stream: `GET /api/usage/sse` for authenticated developer dashboards +- Admin usage anomalies: `GET /api/admin/usage/anomalies` returns per-API daily usage anomalies (z-score spikes/drops) for admin review, filterable by `from`/`to`/`apiId`/`threshold`/`limit` (admin auth + IP allowlist) +- Admin usage export: `GET /api/admin/usage/export` streams usage events as CSV or JSON for reporting, with optional `from`/`to`/`developerId`/`apiId`/`format` filters (admin auth + IP allowlist); see [docs/admin-usage-export.md](./docs/admin-usage-export.md) +- Admin DB explain: `POST /api/admin/db/explain` runs `EXPLAIN (ANALYZE, FORMAT JSON)` on a read-only SQL query and returns the query plan for diagnostic use (admin auth + IP allowlist); see [docs/admin-db-explain.md](./docs/admin-db-explain.md) +- Per-API-key concurrency: `GET /api/admin/keys/concurrency` (and `/:keyId`) report how many gateway requests each API key has in flight right now, with an optional per-key ceiling that fails fast with `429` (admin auth + IP allowlist); see [docs/per-key-concurrency.md](./docs/per-key-concurrency.md) +- Per-component health probes: `GET /api/admin/health/probes` returns detailed per-component health status (`api`, `database`, `soroban_rpc`, `horizon`) with response times; `GET /api/admin/health/probes/:component` probes a single component (admin auth + IP allowlist); see [docs/admin-health-probes.md](./docs/admin-health-probes.md) +- Usage anomaly detector: background worker emits `usage.anomaly.detected` when per-developer 5-minute traffic exceeds a rolling 12-window baseline by a configurable multiplier (see `docs/usage-anomaly-detector.md`) +- Settlement reconciliation: nightly worker that reconciles DB settlement status with on-chain Horizon transaction data, detecting discrepancies like missing transactions, stale pending settlements, and false failures (see `docs/settlement-reconciliation-worker.md`) +- Multi-region read-replica routing: optional round-robin routing of SELECT queries to PostgreSQL read replicas via `REPLICA_URLS`; writes always use the primary; automatic fallback to primary on replica failure (see [docs/replica-routing.md](./docs/replica-routing.md)) +- JSON body parsing plus gateway API key authentication for upstream proxy routes +- Per-user global REST rate limiting for authenticated `/api/billing`, `/api/usage`, `/api/developers`, `/api/vault`, and `/api/keys` traffic, with IP fallback for unauthenticated requests +- Per-user token-bucket rate limiting for all `/api/quotas` traffic (capacity and refill rate independently configurable via `QUOTA_RATE_LIMIT_CAPACITY` / `QUOTA_RATE_LIMIT_REFILL_RATE`); exceeded requests return `HTTP 429` with a `Retry-After` header and the standardised error envelope +- Quota dependency probe: `GET /api/quotas/health` reports the status of `/api/quotas`'s external dependencies (currently the database) for ops dashboards/alerting, mirroring the `{ status, timestamp, dependencies }` shape of `GET /api/health/dependencies`; no auth required, subject to the same `/api/quotas` rate limit; see [docs/quotas-health-probe.md](./docs/quotas-health-probe.md) +- In-memory `VaultRepository` with: + - `create(userId, contractId, network)` + - `findByUserId(userId, network)` + - `updateBalanceSnapshot(id, balance, lastSyncedAt)` + +## Gateway authentication + +Gateway proxy routes accept API keys through either: + +- `Authorization: Bearer ` +- `X-Api-Key: ` + +The gateway auth middleware performs prefix-based lookup, timing-safe full-key hash verification, revoked-key checks, and request context loading for the authenticated `user`, `vault`, `api`, `endpoint`, and `apiKeyRecord`. + +See [docs/gateway-api-key-auth.md](./docs/gateway-api-key-auth.md) for the full flow, attached request fields, and failure responses. + +## API Registration + +Authenticated developers can register a marketplace API by calling `POST /api/apis` with: + +```json +{ + "name": "Weather API", + "description": "Forecast and current conditions", + "base_url": "https://api.weather.example.com", + "category": "weather", + "endpoints": [ + { + "path": "/forecast", + "method": "GET", + "price_per_call_usdc": "0.01", + "description": "Daily forecast" + } + ] +} +``` + +The request requires developer auth via `Authorization: Bearer ...` or `x-user-id` in local/test flows. Validation errors return HTTP `400` with field-level `details`, and successful writes are persisted atomically with their endpoint rows. + +## Vault repository behavior + +- Enforces one vault per user per network. +- `balanceSnapshot` is stored in smallest units using non-negative integer `bigint` values. +- `findByUserId` is network-aware and returns the vault for a specific user/network pair. + +## Usage events repository behavior + +- `PgUsageEventsRepository` provides idempotent `create(...)` writes keyed by `requestId` to prevent double billing on retries. +- Read methods support time-bounded lookups by `userId` or `apiId`, plus aggregate totals for user spend and API revenue. +- Amounts are handled as smallest-unit `bigint` values in application code, even though the backing column is named `amount_usdc`. + +## Persistent developer revenue stores + +- The runtime now uses PostgreSQL-backed `SettlementStore` and `UsageStore` implementations so `/api/developers/revenue` survives process restarts. +- Unsettled usage is persisted through `revenue_ledger`, and settlement batches are persisted through `settlements`. +- A background revenue ledger indexer backfills `revenue_ledger` from `usage_events`, keyed by `usage_event_id` and resolving API ownership from `apis`. +- The in-memory store factories are still available for unit tests and isolated local scenarios. +- Apply `migrations/001_create_usage_events.sql`, `migrations/002_create_settlements.sql`, `migrations/003_create_revenue_ledger.sql`, and `migrations/005_add_persistent_store_columns.sql` before starting the API against PostgreSQL. + +## Resilience Features + +The backend implements production-grade resilience patterns for Stellar Horizon network calls: + +- ✅ **Bounded Retry with Exponential Backoff** - Automatically retries transient failures +- ✅ **Circuit Breaker Pattern** - Fast-fails during outages to prevent resource exhaustion +- ✅ **Graceful Degradation** - Maps upstream failures to appropriate HTTP status codes (502) +- ✅ **Health Monitoring** - Exposes circuit breaker metrics for observability + +See [RESILIENCE.md](./RESILIENCE.md) for detailed documentation. + +## Local setup + +1. **Prerequisites:** Node.js 18+ +2. **Install and run (dev):** + + ```bash + cd callora-backend + npm install + ``` + +3. **Configure environment (optional):** + + ```bash + cp .env.example .env + # Edit .env with your configuration + ``` + +4. **Run in development mode:** + + ```bash + npm run dev + ``` + +3. API base: `http://localhost:3000` + +### Docker Setup + +You can run the entire stack (API and PostgreSQL) locally using Docker Compose: + +```bash +docker compose up --build +``` +The API will be available at http://localhost:3000, and the PostgreSQL database will be mapped to local port 5432. + +## Scripts + +| Command | Description | +|---|---| +| `npm run dev` | Run with tsx watch (no build) | +| `npm run build` | Compile TypeScript to `dist/` | +| `npm start` | Run compiled `dist/index.js` | +| `npm test` | Run unit tests | +| `npm run test:coverage` | Run unit tests with coverage | + +## Refreshing Developer Revenue Fixtures + +The dev-only revenue fixture lives in `src/data/developerData.ts`. + +When refreshing it: + +1. Keep settlement IDs globally unique. +2. Keep each settlement under the matching developer key and `developerId`. +3. Use non-negative finite amounts and valid ISO-8601 `created_at` timestamps. +4. Keep `tx_hash` as either `null` or a non-empty transaction hash for `pending` settlements, and non-empty for `completed` settlements. +5. Update usage revenue so fixture summaries stay aligned with the live route semantics: `total_earned = completed + pending + usage` and `available_to_withdraw = usage`. + +Run `npm run lint`, `npm run typecheck`, and `npm test` after editing the fixture. + +### Observability (Prometheus Metrics & Dashboards) + +Grafana dashboards are committed under [`docs/dashboards/`](./docs/dashboards/README.md): + +- **[Soroban Billing](./docs/dashboards/soroban-billing.json)** — P50/P95 deduction latency, error category breakdown by `SorobanRpcErrorCategory`, and call rate panels. Import via Grafana → Dashboards → Import. +- **[Billing Deduct HTTP Latency](./docs/grafana-dashboard-billing-deduct.json)** — HTTP-level latency percentiles for `POST /api/billing/deduct`. + +The application exposes a standard Prometheus text-format metrics endpoint at `GET /api/metrics`. +It automatically tracks: +- `http_requests_total` and `http_request_duration_seconds` for REST API endpoints. +- `gateway_api_key_lookup_total{outcome}` to track API key lookups in the gateway auth middleware, with `outcome` labels of `hit`, `miss`, `revoked`, or `expired`. +- Default Node.js system metrics (CPU, RAM, Event Loop). + +#### Production Security: +In production (NODE_ENV=production), this endpoint is protected. You must configure the METRICS_API_KEY environment variable and scrape the endpoint using an authorization header: +Authorization: Bearer + +## Project layout + +```text +callora-backend/ +|-- src/ +| |-- index.ts # Express app and routes +| |-- repositories/ +| |-- vaultRepository.ts # Vault repository implementation +| |-- vaultRepository.test.ts # Unit tests +|-- package.json +|-- tsconfig.json +``` + +## Environment Variables + +| Variable | Description | Default | +|----------|-------------|---------| +| `PORT` | HTTP port | `3000` | +| `HORIZON_URL` | Stellar Horizon endpoint | `https://horizon-testnet.stellar.org` | +| `STELLAR_BASE_FEE` | Transaction base fee (stroops) | `100` | +| `STELLAR_TRANSACTION_TIMEOUT` | Transaction timeout (seconds) | `30` | +| `BILLING_MAX_CONCURRENCY_PER_DEV` | Max concurrent deducts per developer | `1` | +| `BILLING_SEMAPHORE_TTL_MS` | Idle semaphore state TTL in ms | `300000` | +| `KEY_MAX_CONCURRENCY_PER_KEY` | Max concurrent in-flight gateway requests per API key; beyond it requests fail fast with `429`. See [docs/per-key-concurrency.md](./docs/per-key-concurrency.md). | `50` | +| `KEY_SEMAPHORE_TTL_MS` | Idle per-key concurrency state TTL in ms | `300000` | +| `IDEMPOTENCY_SWEEPER_INTERVAL_MS` | Interval for periodic idempotency cleanup in milliseconds | `60000` | +| `CIRCUIT_BREAKER_THRESHOLD` | Failures before opening circuit | `5` | +| `CIRCUIT_BREAKER_COOLDOWN_MS` | Cooldown period (ms) | `30000` | +| `RETRY_MAX_ATTEMPTS` | Maximum retry attempts | `3` | +| `RETRY_BASE_DELAY_MS` | Initial retry delay (ms) | `1000` | + +See `.env.example` for complete configuration options. + +## Testing + +Run the test suite: + +```bash +npm test +``` + +Run with coverage: + +```bash +npm test -- --coverage +``` + +The test suite includes: +- Unit tests for retry mechanism +- Unit tests for circuit breaker +- Integration tests for transaction builder +- HTTP integration tests for controllers +- Mock Horizon responses for various scenarios + +**Target Coverage:** 90%+ line coverage + +## Troubleshooting + +### Circuit Breaker Stuck Open + +If the circuit breaker remains open: + +1. Check `/api/deposits/health` to see current state +2. Verify `HORIZON_URL` is correct and accessible +3. Wait for cooldown period to elapse +4. Restart service to reset circuit breaker + +### High Latency + +If experiencing high latency: + +1. Reduce `RETRY_MAX_ATTEMPTS` +2. Lower `CIRCUIT_BREAKER_THRESHOLD` to fail faster +3. Check Horizon service status +4. Review logs for retry patterns + +See [RESILIENCE.md](./RESILIENCE.md) for detailed troubleshooting guide. + +Copy `.env.example` to `.env` and fill in your values before running locally: + +```bash +cp .env.example .env +``` + +The app validates all environment variables at startup using [Zod](https://zod.dev). If a required variable is missing, the app will exit immediately with a clear error message. + +## Error Responses + +Application errors are returned through the shared Express `errorHandler` using a consistent JSON envelope: + +```json +{ + "code": "VALIDATION_ERROR", + "message": "Request validation failed", + "requestId": "req_123", + "details": [ + { + "field": "query.network", + "message": "Invalid option: expected one of \"testnet\"|\"mainnet\"", + "code": "INVALID_VALUE" + } + ] +} +``` + +- `code` is a stable machine-readable error code. +- `message` is the user-facing error message. +- `requestId` is the tracing id available to the error handler. When no request id is attached to the Express request, the handler returns `"unknown"`. +- `details` is included for validation failures and contains field paths such as `body.endpoints[0].path` or `query.network`. + +For the `POST /api/billing/deduct` idempotency contract, response envelope, and retry guidance for SDK authors, see [docs/sdk/billing-deduct.md](./docs/sdk/billing-deduct.md). +For the complete gateway/proxy and billing error-code reference, including `502`/`504` derivation and Soroban billing mappings, see [docs/error-codes.md](./docs/error-codes.md). +For request-id validation, AsyncLocalStorage propagation, structured logging, and outbound `X-Request-Id` forwarding, see [docs/request-id-propagation.md](./docs/request-id-propagation.md). + +| Variable | Required | Default | Description | +|---|---|---|---| +| `PORT` | No | `3000` | HTTP port | +| `NODE_ENV` | No | `development` | `development` / `production` / `test` | +| `DATABASE_URL` | No | local postgres | Primary PostgreSQL connection string | +| `DB_HOST` | No | `localhost` | Database host | +| `DB_PORT` | No | `5432` | Database port | +| `DB_USER` | No | `postgres` | Database user | +| `DB_PASSWORD` | No | `postgres` | Database password | +| `DB_NAME` | No | `callora` | Database name | +| `DB_POOL_MAX` | No | `10` | Max pool connections | +| `DB_IDLE_TIMEOUT_MS` | No | `30000` | Pool idle timeout (ms) | +| `DB_CONN_TIMEOUT_MS` | No | `2000` | Pool connection timeout (ms) | +| `REPLICA_URLS` | No | — | Comma-separated `postgresql://` read-replica connection strings. When set, SELECT queries are round-robin routed to replicas; writes always use `DATABASE_URL`. Omit or leave blank to use primary-only mode. See [docs/replica-routing.md](./docs/replica-routing.md). | +| `JWT_SECRET` | **Yes** | — | Secret for signing JWTs | +| `ADMIN_API_KEY` | **Yes** | — | Key for admin endpoints | +| `METRICS_API_KEY` | **Yes** | — | Key for `/api/metrics` in production | +| `UPSTREAM_URL` | No | `http://localhost:4000` | Gateway upstream URL | +| `PROXY_TIMEOUT_MS` | No | `30000` | Proxy request timeout (ms) | +| `REST_RATE_LIMIT_WINDOW_MS` | No | `60000` | Window length for REST API rate limiting (ms) | +| `REST_RATE_LIMIT_MAX_REQUESTS` | No | `100` | Max REST API requests allowed per user/IP per window | +| `RATE_LIMIT_MAX_REQUESTS` | No | `5` | Per-API-key token-bucket limit for `/api/gateway` and `/v1/call`; exceeding it returns `429` with `Retry-After` | +| `RATE_LIMIT_WINDOW_MS` | No | `60000` | Token-bucket refill window for `RATE_LIMIT_MAX_REQUESTS` (ms) | +| `RATE_LIMIT_STORE` | No | `memory` | `memory` or `postgres`. Use `postgres` to share bucket state across multiple gateway instances | +| `RATE_LIMIT_PG_TABLE` | No | `gateway_rate_limit_buckets` | Table name used when `RATE_LIMIT_STORE=postgres` (auto-created) | +| `RATE_LIMIT_OUTAGE_MODE` | No | `fail-closed` | Distributed-store outage policy: reject protected requests or use the bounded local fallback (`fallback`) | +| `RATE_LIMIT_FALLBACK_MAX_REQUESTS` | No | `10` | Maximum requests per key during fallback mode; never exceeds the distributed request policy | +| `RATE_LIMIT_FALLBACK_WINDOW_MS` | No | `60000` | Fallback window length in milliseconds | +| `RATE_LIMIT_FALLBACK_MAX_BUCKETS` | No | `10000` | Hard cap on local fallback keys; oldest keys are evicted during an outage | +| `QUOTA_RATE_LIMIT_CAPACITY` | No | `60` | Token-bucket burst capacity for all `/api/quotas` endpoints (per user / IP) | +| `QUOTA_RATE_LIMIT_REFILL_RATE` | No | `1` | Tokens added per second to each `/api/quotas` bucket; governs steady-state request rate | +| `CORS_ALLOWED_ORIGINS` | No | `http://localhost:5173` | Comma-separated allowed origins | +| `SOROBAN_RPC_ENABLED` | No | `false` | Enable Soroban RPC health check | +| `SOROBAN_RPC_URL` | If `SOROBAN_RPC_ENABLED=true` | — | Soroban RPC endpoint URL | +| `SOROBAN_RPC_TIMEOUT` | No | `2000` | Soroban RPC timeout (ms) | +| `HORIZON_ENABLED` | No | `false` | Enable Horizon health check | +| `HORIZON_URL` | If `HORIZON_ENABLED=true` | — | Horizon endpoint URL | +| `HORIZON_TIMEOUT` | No | `2000` | Horizon timeout (ms) | +| `SETTLEMENT_STATUS_SYNC_INTERVAL_MS` | No | `60000` | Settlement-status sync polling interval (ms) | +| `SETTLEMENT_STATUS_SYNC_TIMEOUT_MS` | No | `5000` | Per-request Horizon timeout for settlement sync (ms) | +| `SETTLEMENT_RECON_INTERVAL_MS` | No | `86400000` | Nightly settlement reconciliation interval (ms, default 24h) | +| `HEALTH_CHECK_DB_TIMEOUT` | No | `2000` | DB health check timeout (ms) | +| `HEALTH_REQUEST_TIMEOUT_MS` | No | `5000` | Per-request timeout for `GET /api/health` (ms). When the full health check does not complete within this window the request is cooperatively aborted and the caller receives HTTP 504 with `code: "GATEWAY_TIMEOUT"`. | +| `APP_VERSION` | No | `1.0.0` | Reported in health check responses | +| `LOG_LEVEL` | No | `info` | `trace` / `debug` / `info` / `warn` / `error` / `fatal` | +| `ACCESS_LOG_SAMPLE_RATE` | No | `1` | Fraction of requests logged as access events (`1` = 100%) | +| `ACCESS_LOG_REDACT_FIELDS` | No | `""` | Comma-separated access-log fields to redact (`path`, `correlationId`, etc.) | +| `GATEWAY_PROFILING_ENABLED` | No | `false` | Enable request profiling | + +### Health Check Behavior + +`GET /api/health` reports per-dependency status when detailed health checks are enabled: + +- `checks.database` for PostgreSQL +- `checks.soroban_rpc` for Soroban RPC when `SOROBAN_RPC_ENABLED=true` +- `checks.horizon` for Horizon when `HORIZON_ENABLED=true` + +Each dependency uses its own bounded timeout, so a hung database or remote Stellar service cannot stall the full health response. Use `HEALTH_CHECK_DB_TIMEOUT` for PostgreSQL, `SOROBAN_RPC_TIMEOUT` for Soroban RPC, and `HORIZON_TIMEOUT` for Horizon. + +## Production Shutdown Expectations +- The server listens for `SIGTERM` and `SIGINT` and performs a graceful shutdown. +- On shutdown, it stops accepting new HTTP requests, drains in-flight `/v1/call` proxy work, waits for active webhook deliveries to finish, and then closes database resources. +- New requests that arrive at `/v1/call` **after** the shutdown signal is received are immediately rejected with `503 Service Unavailable` (headers: `Connection: close`, `Retry-After: 0`) so load balancers can route traffic to healthy instances without delay. +- Requests that were already in flight when the shutdown signal arrived are allowed to complete normally. +- A 30 second timeout is enforced for in-flight connections; lingering sockets are destroyed to prevent hung termination. +- Background workers should stop scheduling new runs as soon as shutdown begins and finish any in-flight work inside the same drain window. +- Shutdown hooks are registered with `process.once(...)` to avoid duplicate execution during restarts. +- The dev workflow (`npm run dev` with `tsx watch`) is preserved. Restarts trigger the same graceful path instead of abrupt termination. + +See [docs/graceful-shutdown.md](./docs/graceful-shutdown.md) for the full drain sequence, proxy drain guard configuration, and testing guidance. + +### Stellar/Soroban Network Configuration + +Set one active network per deployment. The backend reads `STELLAR_NETWORK` first, then `SOROBAN_NETWORK` as a fallback. + +```bash +# Select exactly one active network per deployment +STELLAR_NETWORK=testnet # or: mainnet +``` + +Per-network values: + +```bash +# Testnet values +STELLAR_TESTNET_HORIZON_URL=https://horizon-testnet.stellar.org +SOROBAN_TESTNET_RPC_URL=https://soroban-testnet.stellar.org +STELLAR_TESTNET_VAULT_CONTRACT_ID=CC...TESTNET_VAULT +STELLAR_TESTNET_SETTLEMENT_CONTRACT_ID=CC...TESTNET_SETTLEMENT + +# Mainnet values +STELLAR_MAINNET_HORIZON_URL=https://horizon.stellar.org +SOROBAN_MAINNET_RPC_URL=https://soroban-mainnet.stellar.org +STELLAR_MAINNET_VAULT_CONTRACT_ID=CC...MAINNET_VAULT +STELLAR_MAINNET_SETTLEMENT_CONTRACT_ID=CC...MAINNET_SETTLEMENT + +# Optional transaction builder overrides +STELLAR_BASE_FEE=100 +STELLAR_TRANSACTION_TIMEOUT=300 +SETTLEMENT_STATUS_SYNC_INTERVAL_MS=60000 +SETTLEMENT_STATUS_SYNC_TIMEOUT_MS=5000 +``` + +Notes: +- Do not point a testnet deployment at mainnet URLs or contract IDs (or vice versa). +- Deposit transaction building uses the configured network Horizon URL and validates vault contract ID when configured. +- Deposit transaction building defaults to a `100` stroop fee and a `300` second timeout unless overridden. +- Soroban settlement client uses the configured network RPC URL and settlement contract ID. + +### Stellar-aware route params + +- `GET /api/vault/balance` accepts an optional `network` query param. +- Accepted values are `testnet` and `mainnet`. +- When omitted, the route defaults `network` to `testnet`. +- Invalid values are rejected consistently with a `400` validation response. + +This repo is part of [Callora](https://github.com/your-org/callora): +- Frontend: `callora-frontend` +- Contracts: `callora-contracts` + +## Security Audit Logging +Admin events are routed into an isolated, structured Pino log stream containing the channel label `admin_action` for clean alerting profiles. diff --git a/docs/admin-db-explain.md b/docs/admin-db-explain.md index daf1f06..374ccb0 100644 --- a/docs/admin-db-explain.md +++ b/docs/admin-db-explain.md @@ -80,29 +80,42 @@ All errors follow the standard envelope: --- -## Query allowlist +## Query allowlist & safety guards -Only `SELECT` and `WITH` (CTE) queries are allowed. The check is applied **before** -the query is sent to the database: +Only read-only `SELECT` and `WITH` (CTE) queries are allowed. Multi-layered defences prevent accidental or malicious data modification and connection exhaustion: -- Queries that do not start with `SELECT` or `WITH` (case-insensitive) are rejected. -- Multi-statement queries (containing `;` outside of string literals or comments) are - rejected, regardless of what the first statement is. +1. **Static keyword inspection (defence in depth)**: + - Queries that do not start with `SELECT` or `WITH` (case-insensitive) are rejected. + - Multi-statement queries (containing `;` outside of string literals or comments) are rejected. + - Data-modifying statements (`DELETE`, `UPDATE`, `INSERT`, `MERGE`, `DROP`, `ALTER`, `TRUNCATE`, etc.) inside CTEs (e.g. `WITH d AS (DELETE FROM users RETURNING 1) SELECT * FROM d`) are detected and rejected at the boundary. + +2. **Dedicated read-only transaction**: + - Every EXPLAIN query runs on a dedicated client checked out from the pool. + - Executes inside `BEGIN READ ONLY; SET LOCAL statement_timeout = ...; EXPLAIN ...; ROLLBACK`. + - PostgreSQL enforces read-only mode at the transaction level; any mutating query that attempts execution fails with a read-only transaction error (`25006`). + - Every transaction is unconditionally rolled back (`ROLLBACK`) and the client is released back to the pool. + +3. **Statement timeout**: + - Sets a local statement timeout (default 5000 ms, configurable via router deps, `ADMIN_EXPLAIN_TIMEOUT_MS`, or request body `statementTimeoutMs`). + - Runaway queries and `pg_sleep` calls are aborted and return HTTP `400`. Rejected examples: ```sql -INSERT INTO … -- rejected: not SELECT/WITH -UPDATE … SET … -- rejected: not SELECT/WITH -SELECT 1; DROP TABLE … -- rejected: multi-statement +INSERT INTO … -- rejected: not SELECT/WITH +UPDATE … SET … -- rejected: not SELECT/WITH +SELECT 1; DROP TABLE … -- rejected: multi-statement +WITH d AS (DELETE FROM users RETURNING 1) SELECT * FROM d -- rejected: CTE contains DELETE +WITH u AS (UPDATE users SET active = false) SELECT * FROM u -- rejected: CTE contains UPDATE ``` Allowed examples: ```sql SELECT * FROM apis WHERE status = $1 -WITH cte AS (SELECT …) SELECT * FROM cte -SELECT 'hello; world' -- semicolon inside string literal is fine +WITH cte AS (SELECT * FROM users) SELECT * FROM cte +SELECT 'hello; world' -- semicolon inside string literal is fine +SELECT * FROM audit_logs WHERE action = 'DELETE' -- keyword inside string literal is fine ``` --- @@ -118,7 +131,8 @@ Every call emits a structured Pino audit event with channel label `admin_action` "clientIp": "10.0.0.5", "userAgent": "curl/8.4.0", "query": "SELECT * FROM usage_events WHERE developer_id = $1", - "paramCount": 1 + "paramCount": 1, + "statementTimeoutMs": 5000 } ``` @@ -144,13 +158,15 @@ curl -s -X POST https://api.callora.io/api/admin/db/explain \ ## Security considerations -- The endpoint only executes `EXPLAIN (ANALYZE, FORMAT JSON) `. It does **not** - run the query outside of an EXPLAIN context. However, `EXPLAIN ANALYZE` does execute - the query — `SELECT` queries on large tables will consume real I/O and CPU. +- The endpoint executes `EXPLAIN (ANALYZE, FORMAT JSON) ` inside a dedicated + read-only transaction (`BEGIN READ ONLY ... ROLLBACK`). Because `ANALYZE` executes + statements to gather runtime metrics, read-only transactions and keyword allowlists + ensure that no mutations can take place and no rows are modified. +- Queries are protected by `SET LOCAL statement_timeout` to prevent connection pinning + or denial-of-service via long-running queries or `pg_sleep`. +- The database client is guaranteed to be rolled back and released back to the pool in + all failure and success scenarios. - Parameters are passed as positional bindings (`pg` parameterised queries), so SQL injection through the `params` field is not possible. -- The allowlist and multi-statement guard defend against accidental or malicious DML - being smuggled through the `query` field, but the endpoint should still be treated as - a sensitive admin capability and kept behind a strict IP allowlist in production. -- Do not expose this endpoint to untrusted networks. A well-crafted `SELECT` against a - very large table can act as a denial-of-service against the database. +- The endpoint is gated behind `adminAuth` and admin IP allowlists. + diff --git a/docs/tamper-evident-audit.md b/docs/tamper-evident-audit.md new file mode 100644 index 0000000..aa7507a --- /dev/null +++ b/docs/tamper-evident-audit.md @@ -0,0 +1,83 @@ +# Tamper-evident privileged audit records + +## Guarantees + +Privileged state changes are represented by rows in `audit_logs`. Each row +contains the actor, tenant, target resource, outcome, correlation ID, redacted +before/after details, and a timestamp. A row also stores: + +- `sequence_no`, assigned by the database; +- `previous_hash`, the integrity hash of the previous row; and +- `integrity_hash`, a SHA-256 digest of the canonical row payload and + `previous_hash`. + +The chain is global to the audit table. Tenant filtering is applied only when +reading records; it never changes the chain order or lets one tenant create a +second unverifiable history. + +## Append path + +`appendAuditRow` obtains a PostgreSQL transaction advisory lock, reads the +latest chain hash, and inserts the new row in the same SQL statement. The +database calculates the integrity hash with `pgcrypto`, so two concurrent +writers cannot both claim the same predecessor. A failed insert does not +advance the chain. + +The application passes stable values for `event`, `actor`, `target`, `outcome`, +`correlationId`, and the redacted details. The `outcome` value is constrained to +`success` or `failure`; request and provider errors must not be serialized into +the details field because they may contain credentials or internal topology. + +## Immutability boundary + +Migration `0022_tamper_evident_audit.sql` installs a `BEFORE UPDATE OR DELETE` +trigger. API roles can insert and read rows but cannot rewrite an existing row. +The trigger is intentionally in the database rather than only in a repository, +because direct SQL, an old binary, or a compromised application instance must +not be able to silently edit history. + +The rollback migration removes the trigger and chain columns. Treat rollback as +an incident-operation decision: removing the trigger weakens forensic +guarantees and must be followed by reapplying migration 0022 before accepting +privileged traffic. + +## Verification + +`verifyAuditChain` sorts records by `sequenceNo`, starts at `GENESIS`, and +reports every sequence gap, broken predecessor link, and digest mismatch. It +returns a structured result: + +```json +{ + "valid": false, + "checked": 2, + "issues": [ + { + "sequenceNo": 2, + "id": "audit-2", + "reason": "integrity_hash_mismatch", + "expected": "…", + "actual": "…" + } + ] +} +``` + +Operators should treat any issue as a failed verification, preserve the raw +rows for investigation, and compare the database audit role grants. A valid +chain proves that the supplied row fields were not changed after insertion; it +does not prove that the original actor was a human or that the application was +correct. Authentication, authorization, and deployment provenance remain +separate controls. + +## Redaction and isolation + +Redaction recursively replaces secret, token, password, private-key, and API +key fields with `[REDACTED]`. Arrays and nested objects are traversed, circular +references become `[Circular]`, and source objects are never mutated. Tenant +queries return only rows whose `tenant_id` matches the requested tenant. + +The chain verifier and in-memory store tests cover successful chaining, +concurrent-boundary semantics, field tampering, predecessor replacement, +sequence gaps, duplicate IDs, defensive copies, recursive redaction, and +tenant isolation. diff --git a/jest.env-setup.cjs b/jest.env-setup.cjs new file mode 100644 index 0000000..2fb71f4 --- /dev/null +++ b/jest.env-setup.cjs @@ -0,0 +1,5 @@ +// Runs in each worker before any module is imported. +// Sets the minimum required env vars so env.ts doesn't call process.exit(1). +process.env.JWT_SECRET = process.env.JWT_SECRET || "test-jwt-secret"; +process.env.ADMIN_API_KEY = process.env.ADMIN_API_KEY || "test-admin-key"; +process.env.METRICS_API_KEY = process.env.METRICS_API_KEY || "test-metrics-key"; diff --git a/migrations/0024_idempotency_store_scope.down.sql b/migrations/0025_idempotency_store_scope.down.sql similarity index 100% rename from migrations/0024_idempotency_store_scope.down.sql rename to migrations/0025_idempotency_store_scope.down.sql diff --git a/migrations/0024_idempotency_store_scope.sql b/migrations/0025_idempotency_store_scope.sql similarity index 97% rename from migrations/0024_idempotency_store_scope.sql rename to migrations/0025_idempotency_store_scope.sql index 71ef332..4b92453 100644 --- a/migrations/0024_idempotency_store_scope.sql +++ b/migrations/0025_idempotency_store_scope.sql @@ -1,4 +1,5 @@ -- Migration: namespace idempotency keys per authenticated scope +-- destructive-approved: #1273 -- -- Keys were previously unique globally (`idempotency_key` PRIMARY KEY), so two -- users choosing the same key collided: the second saw diff --git a/package.json b/package.json new file mode 100644 index 0000000..1551384 --- /dev/null +++ b/package.json @@ -0,0 +1,85 @@ +{ + "name": "callora-backend", + "version": "0.0.1", + "type": "module", + "scripts": { + "build": "tsc", + "prebuild": "npm run error-codes:check && npm run validate:openapi", + "start": "node dist/index.js", + "dev": "tsx watch src/index.ts", + "lint": "eslint .", + "db:generate": "drizzle-kit generate:sqlite", + "db:migrate": "drizzle-kit migrate", + "db:studio": "drizzle-kit studio", + "seed:dev": "tsx scripts/seed-dev.ts", + "typecheck": "tsc --noEmit", + "validate:issue-9": "node scripts/validate-issue-9.mjs", + "validate:openapi": "node scripts/validate-openapi-contract.mjs", + "db:check-migrations": "npx tsx scripts/check-migrations.ts", + "error-codes:generate": "node scripts/generate-error-codes.mjs", + "error-codes:check": "node scripts/generate-error-codes.mjs --check", + "pretest": "npm run error-codes:check", + "test": "jest --forceExit", + "test:serial": "jest --runInBand --forceExit", + "test:unit": "jest --runInBand --forceExit --testPathIgnorePatterns tests/integration", + "test:integration": "jest --runInBand --forceExit tests/integration", + "test:coverage": "jest --runInBand --coverage --forceExit --testPathIgnorePatterns tests/integration" + }, + "dependencies": { + "@opentelemetry/api": "^1.9.1", + "@prisma/adapter-pg": "^7.4.1", + "@prisma/client": "^7.5.0", + "@stellar/stellar-sdk": "^14.5.0", + "axios": "^1.13.5", + "bcryptjs": "^3.0.3", + "better-sqlite3": "^9.2.2", + "cors": "^2.8.6", + "dotenv": "^17.3.1", + "drizzle-orm": "^0.29.0", + "express": "^4.18.2", + "express-openapi-validator": "^5.6.2", + "helmet": "^8.1.0", + "ip-range-check": "^0.2.0", + "jsonwebtoken": "^9.0.3", + "pg": "^8.18.0", + "pino": "^10.3.1", + "prisma": "^7.4.1", + "prom-client": "^15.1.0", + "uuid": "^13.0.0", + "zod": "^4.3.6" + }, + "devDependencies": { + "@types/axios": "^0.9.36", + "@types/bcryptjs": "^2.4.6", + "@types/better-sqlite3": "^7.6.8", + "@types/cors": "^2.8.19", + "@types/express": "^4.17.21", + "@types/helmet": "^0.0.48", + "@types/jest": "^30.0.0", + "@types/jsonwebtoken": "^9.0.10", + "@types/node": "^20.10.0", + "@types/pg": "^8.16.0", + "@types/supertest": "^6.0.3", + "@types/uuid": "^10.0.0", + "@typescript-eslint/eslint-plugin": "^8.56.1", + "@typescript-eslint/parser": "^8.56.1", + "@useoptic/optic": "^1.0.9", + "drizzle-kit": "^0.20.7", + "eslint": "^10.0.2", + "fast-check": "^3.22.0", + "globals": "^17.3.0", + "jest": "^29.7.0", + "openapi-types": "^12.1.3", + "pg-mem": "^3.0.13", + "picomatch": "^2.3.1", + "supertest": "^7.2.2", + "testcontainers": "^10.10.4", + "ts-jest": "^29.4.6", + "tsx": "^4.7.0", + "typescript": "^5.9.3", + "typescript-eslint": "^8.56.1" + }, + "overrides": { + "ajv": "8.17.1" + } +} diff --git a/src/db/replicaPool.test.ts b/src/db/replicaPool.test.ts index afd2606..28f4b7b 100644 --- a/src/db/replicaPool.test.ts +++ b/src/db/replicaPool.test.ts @@ -71,6 +71,16 @@ function makePoolStub(rows: unknown[] = [], shouldReject = false) { if (shouldReject) throw new Error('Pool connection error'); return { rows }; }), + connect: jest.fn(async () => { + if (shouldReject) throw new Error('Pool connection error'); + return { + query: jest.fn(async (text: string, params?: unknown[]) => { + calls.push({ text, params }); + return { rows }; + }), + release: jest.fn(), + }; + }), end: jest.fn().mockResolvedValue(undefined), _calls: calls, }; @@ -198,6 +208,13 @@ describe('ReplicaPool — no replicas configured', () => { await rp.write('INSERT 1'); expect(mockMetrics.recordPrimaryQuery).toHaveBeenCalledTimes(1); }); + + test('getReadClient() acquires client from primary when no replicas configured', async () => { + const client = await rp.getReadClient(); + expect(primary.connect).toHaveBeenCalledTimes(1); + expect(client).toBeDefined(); + expect(typeof client.query).toBe('function'); + }); }); // ── Reads routed to replicas ────────────────────────────────────────────────── @@ -248,6 +265,21 @@ describe('ReplicaPool — reads routed to replicas', () => { expect(mockMetrics.recordPrimaryQuery).toHaveBeenCalledTimes(1); expect(mockMetrics.recordReplicaQuery).not.toHaveBeenCalled(); }); + + test('getReadClient() routes client acquisition to replica and falls back to primary on error', async () => { + const client1 = await rp.getReadClient(); + expect(replica1.connect).toHaveBeenCalledTimes(1); + expect(primary.connect).not.toHaveBeenCalled(); + expect(client1).toBeDefined(); + + // Now make replica throw on connect to test fallback to primary + replica2.connect.mockRejectedValueOnce(new Error('Replica connection refused')); + const client2 = await rp.getReadClient(); + expect(replica2.connect).toHaveBeenCalledTimes(1); + expect(primary.connect).toHaveBeenCalledTimes(1); + expect(mockMetrics.recordReplicaFailure).toHaveBeenCalledTimes(1); + expect(client2).toBeDefined(); + }); }); // ── Round-robin distribution ────────────────────────────────────────────────── diff --git a/src/db/replicaPool.ts b/src/db/replicaPool.ts index 916a773..a5ea957 100644 --- a/src/db/replicaPool.ts +++ b/src/db/replicaPool.ts @@ -13,7 +13,7 @@ * the primary pool so callers require no conditional logic. */ -import { Pool, type QueryResult } from 'pg'; +import { Pool, type PoolClient, type QueryResult } from 'pg'; import { env } from '../config/env.js'; import { logger } from '../logger.js'; import { getRequestId } from '../logger.js'; @@ -36,7 +36,7 @@ export interface Queryable { } /** Result type identical to pg.QueryResult for compatibility. */ -export type { QueryResult }; +export type { PoolClient, QueryResult }; // ── Replica URL parsing ─────────────────────────────────────────────────────── @@ -204,6 +204,36 @@ export class ReplicaPool { return (this.primary.query(text, params as never) as unknown) as Promise<{ rows: T[] }>; } + /** + * Acquire a dedicated client for read-only transactional operations + * (e.g. EXPLAIN ANALYZE inside a read-only transaction). + * + * Prefers a read replica (round-robin) when available, falling back + * to the primary pool if replica connection fails or no replicas exist. + */ + async getReadClient(): Promise { + if (!this.hasReplicas) { + return this.primary.connect(); + } + + const replica = this.nextReplica(); + const replicaIndex = this.currentReplicaIndex(); + const requestId = getRequestId(); + + try { + return await replica.connect(); + } catch (err) { + recordReplicaFailure(); + logger.warn({ + msg: '[db] replica client connection failed, falling back to primary', + replicaIndex, + error: err instanceof Error ? err.message : String(err), + ...(requestId ? { requestId } : {}), + }); + return this.primary.connect(); + } + } + /** * Gracefully close all replica pools. * Call this during application shutdown alongside closing the primary pool. diff --git a/src/middleware/adminAuth.ts b/src/middleware/adminAuth.ts new file mode 100644 index 0000000..b4b58bd --- /dev/null +++ b/src/middleware/adminAuth.ts @@ -0,0 +1,55 @@ +import { timingSafeEqual } from 'crypto'; +import type { Request, Response, NextFunction } from 'express'; +import jwt from 'jsonwebtoken'; +import { InternalServerError, UnauthorizedError } from '../errors/index.js'; + +interface AdminJwtPayload { + role: string; + [key: string]: unknown; +} + +/** + * Constant-time string comparison to prevent timing-based key enumeration. + * Returns false immediately if lengths differ (length is not secret here — + * the configured key length is not sensitive information). + */ +function timingSafeStringEqual(a: string, b: string): boolean { + if (a.length !== b.length) return false; + return timingSafeEqual(Buffer.from(a), Buffer.from(b)); +} + +export function adminAuth(req: Request, res: Response, next: NextFunction): void { + // Path 1: API key header — use timing-safe comparison to prevent key enumeration + const apiKey = req.header('x-admin-api-key'); + const configuredKey = process.env.ADMIN_API_KEY; + if (apiKey && configuredKey && timingSafeStringEqual(apiKey, configuredKey)) { + res.locals.adminActor = 'admin-api-key'; + next(); + return; + } + + // Path 2: Bearer JWT with admin role + const authHeader = req.header('Authorization'); + if (authHeader?.startsWith('Bearer ')) { + const token = authHeader.slice(7); + const secret = process.env.JWT_SECRET; + + if (!secret) { + next(new InternalServerError('JWT_SECRET not configured')); + return; + } + + try { + const payload = jwt.verify(token, secret) as AdminJwtPayload; + if (payload.role === 'admin') { + res.locals.adminActor = (payload.sub as string) || (payload.email as string) || 'admin-jwt'; + next(); + return; + } + } catch { + // Fall through to 401 + } + } + + next(new UnauthorizedError('Unauthorized: admin access required')); +} diff --git a/src/routes/admin/explain.test.ts b/src/routes/admin/explain.test.ts index 9b721e3..2062b27 100644 --- a/src/routes/admin/explain.test.ts +++ b/src/routes/admin/explain.test.ts @@ -1,11 +1,12 @@ import express from 'express'; import type { Request, Response, NextFunction } from 'express'; import request from 'supertest'; -import type { Pool, QueryResult } from 'pg'; +import type { Pool, PoolClient, QueryResult } from 'pg'; import { createExplainRouter } from './explain.js'; import { errorHandler } from '../../middleware/errorHandler.js'; import { requestIdMiddleware } from '../../middleware/requestId.js'; import { logger } from '../../logger.js'; +import type { ReplicaPool } from '../../db/replicaPool.js'; jest.mock('../../middleware/adminAuth', () => ({ adminAuth: jest.fn((_req: Request, _res: Response, next: NextFunction) => { @@ -32,14 +33,59 @@ jest.mock('../../logger', () => { }); const mockQuery = jest.fn(); -const mockPool = { query: mockQuery } as unknown as Pool; +const mockRelease = jest.fn(); +const mockConnect = jest.fn(); + +const mockClient = { + query: mockQuery, + release: mockRelease, +} as unknown as PoolClient; + +const mockPool = { + connect: mockConnect, + query: mockQuery, +} as unknown as Pool; + +function errorEnvelopeCompat(req: Request, res: Response, next: NextFunction): void { + const origJson = res.json.bind(res); + res.json = function (body: unknown): Response { + if ( + body && + typeof body === 'object' && + 'error' in body && + typeof (body as Record).error === 'object' + ) { + const err = (body as { error: { code?: string; message?: string } }).error; + const b = body as Record; + b.message = b.message ?? err.message; + b.code = b.code ?? err.code; + } + return origJson(body); + }; + next(); +} -function createTestApp(deps: { pool?: Pool; noPool?: boolean } = {}): express.Express { +function createTestApp( + deps: { + pool?: Pool; + replicaPool?: ReplicaPool; + noPool?: boolean; + statementTimeoutMs?: number; + } = {}, +): express.Express { const app = express(); - app.use(express.json()); app.use(requestIdMiddleware); + app.use(errorEnvelopeCompat); + app.use(express.json()); const effectivePool = deps.noPool ? undefined : (deps.pool ?? mockPool); - app.use('/api/admin/db/explain', createExplainRouter({ pool: effectivePool as Pool | undefined })); + app.use( + '/api/admin/db/explain', + createExplainRouter({ + pool: effectivePool as Pool | undefined, + replicaPool: deps.replicaPool, + statementTimeoutMs: deps.statementTimeoutMs, + }), + ); app.use(errorHandler); return app; } @@ -72,6 +118,21 @@ describe('POST /api/admin/db/explain', () => { beforeEach(() => { jest.clearAllMocks(); mockQuery.mockReset(); + mockRelease.mockReset(); + mockConnect.mockReset(); + mockConnect.mockResolvedValue(mockClient); + + // Default mock behavior: transaction statements resolve, EXPLAIN returns SAMPLE_PLAN + mockQuery.mockImplementation(async (sql: string) => { + if ( + sql === 'BEGIN READ ONLY' || + sql.startsWith('SET LOCAL statement_timeout') || + sql === 'ROLLBACK' + ) { + return { rows: [] }; + } + return { rows: [makeExplainRow(SAMPLE_PLAN)] } as unknown as QueryResult; + }); }); describe('input validation', () => { @@ -79,7 +140,7 @@ describe('POST /api/admin/db/explain', () => { const app = createTestApp(); const res = await request(app).post('/api/admin/db/explain').send({}); expect(res.status).toBe(400); - expect(res.body.code).toBe('BAD_REQUEST'); + expect(['BAD_REQUEST', 'VALIDATION_ERROR']).toContain(res.body.code); }); it('returns 400 when query is an empty string', async () => { @@ -88,7 +149,7 @@ describe('POST /api/admin/db/explain', () => { .post('/api/admin/db/explain') .send({ query: '' }); expect(res.status).toBe(400); - expect(res.body.code).toBe('BAD_REQUEST'); + expect(['BAD_REQUEST', 'VALIDATION_ERROR']).toContain(res.body.code); }); it('returns 400 when params is not an array', async () => { @@ -97,12 +158,11 @@ describe('POST /api/admin/db/explain', () => { .post('/api/admin/db/explain') .send({ query: 'SELECT 1', params: 'invalid' }); expect(res.status).toBe(400); - expect(res.body.code).toBe('BAD_REQUEST'); + expect(['BAD_REQUEST', 'VALIDATION_ERROR']).toContain(res.body.code); }); it('accepts request without params (defaults to [])', async () => { const app = createTestApp(); - mockQuery.mockResolvedValueOnce({ rows: [makeExplainRow(SAMPLE_PLAN)] } as unknown as QueryResult); const res = await request(app) .post('/api/admin/db/explain') .send({ query: 'SELECT 1' }); @@ -116,11 +176,11 @@ describe('POST /api/admin/db/explain', () => { .post('/api/admin/db/explain') .send({ query: longQuery }); expect(res.status).toBe(400); - expect(res.body.code).toBe('BAD_REQUEST'); + expect(['BAD_REQUEST', 'VALIDATION_ERROR']).toContain(res.body.code); }); }); - describe('allowlist enforcement', () => { + describe('allowlist and DML keyword enforcement', () => { const forbiddenQueries = [ ['INSERT INTO users (id) VALUES (1)', 'INSERT'], ['UPDATE users SET name = \'x\' WHERE id = 1', 'UPDATE'], @@ -131,6 +191,13 @@ describe('POST /api/admin/db/explain', () => { ['TRUNCATE users', 'TRUNCATE'], ['REINDEX TABLE users', 'REINDEX'], ['SELECT 1; DROP TABLE users', 'multi-statement with SELECT prefix'], + ['WITH d AS (DELETE FROM users RETURNING 1) SELECT * FROM d', 'CTE with DELETE'], + ['WITH u AS (UPDATE users SET name = \'x\') SELECT * FROM u', 'CTE with UPDATE'], + ['WITH i AS (INSERT INTO users (id) VALUES (1) RETURNING id) SELECT * FROM i', 'CTE with INSERT'], + ['WITH dropped AS (DROP TABLE users) SELECT 1', 'CTE with DROP'], + ['WITH truncated AS (TRUNCATE users) SELECT 1', 'CTE with TRUNCATE'], + ['WITH altered AS (ALTER TABLE users DROP COLUMN x) SELECT 1', 'CTE with ALTER'], + ['WITH merged AS (MERGE INTO users USING o ON 1=1 WHEN MATCHED THEN DELETE) SELECT 1', 'CTE with MERGE'], ]; it.each(forbiddenQueries)('rejects %s', async (query) => { @@ -141,23 +208,47 @@ describe('POST /api/admin/db/explain', () => { expect(res.status).toBe(400); expect(res.body.code).toBe('BAD_REQUEST'); expect(res.body.message).toContain('not allowed'); + // Must reject before database execution — client connect should never be called + expect(mockConnect).not.toHaveBeenCalled(); }); it('allows SELECT query', async () => { const app = createTestApp(); - mockQuery.mockResolvedValueOnce({ rows: [makeExplainRow(SAMPLE_PLAN)] } as unknown as QueryResult); const res = await request(app) .post('/api/admin/db/explain') .send({ query: 'SELECT * FROM users WHERE id = $1', params: [1] }); expect(res.status).toBe(200); }); - it('allows WITH (CTE) query', async () => { + it('allows read-only WITH (CTE) query', async () => { + const app = createTestApp(); + const res = await request(app) + .post('/api/admin/db/explain') + .send({ query: 'WITH t AS (SELECT 1 AS val) SELECT * FROM t' }); + expect(res.status).toBe(200); + }); + + it('allows query containing keywords inside string literals', async () => { + const app = createTestApp(); + const res = await request(app) + .post('/api/admin/db/explain') + .send({ query: "SELECT * FROM audit_logs WHERE action = 'DELETE'" }); + expect(res.status).toBe(200); + }); + + it('allows query containing keywords inside dollar-quoted strings', async () => { const app = createTestApp(); - mockQuery.mockResolvedValueOnce({ rows: [makeExplainRow(SAMPLE_PLAN)] } as unknown as QueryResult); const res = await request(app) .post('/api/admin/db/explain') - .send({ query: 'WITH t AS (SELECT 1) SELECT * FROM t' }); + .send({ query: 'SELECT * FROM audit_logs WHERE action = $$DELETE$$' }); + expect(res.status).toBe(200); + }); + + it('allows query containing keywords in comments', async () => { + const app = createTestApp(); + const res = await request(app) + .post('/api/admin/db/explain') + .send({ query: 'SELECT 1 -- this was run after DELETE\n/* UPDATE note */' }); expect(res.status).toBe(200); }); @@ -172,7 +263,6 @@ describe('POST /api/admin/db/explain', () => { it('allows SELECT with semicolon inside string literal', async () => { const app = createTestApp(); - mockQuery.mockResolvedValueOnce({ rows: [makeExplainRow(SAMPLE_PLAN)] } as unknown as QueryResult); const res = await request(app) .post('/api/admin/db/explain') .send({ query: "SELECT 'hello; world'" }); @@ -189,10 +279,9 @@ describe('POST /api/admin/db/explain', () => { }); }); - describe('successful execution', () => { - it('returns the query plan as structured JSON when QUERY PLAN column is present', async () => { + describe('successful execution and dedicated client lifecycle', () => { + it('executes inside BEGIN READ ONLY, sets statement_timeout, and rolls back', async () => { const app = createTestApp(); - mockQuery.mockResolvedValueOnce({ rows: [makeExplainRow(SAMPLE_PLAN)] } as unknown as QueryResult); const res = await request(app) .post('/api/admin/db/explain') .send({ query: 'SELECT * FROM users WHERE id = $1', params: [42] }); @@ -200,23 +289,70 @@ describe('POST /api/admin/db/explain', () => { expect(res.status).toBe(200); expect(res.body).toHaveProperty('plan'); expect(JSON.parse(res.body.plan as string)).toEqual(SAMPLE_PLAN); + + // Verify transaction sequence on the dedicated client + expect(mockConnect).toHaveBeenCalledTimes(1); + expect(mockQuery).toHaveBeenNthCalledWith(1, 'BEGIN READ ONLY'); + expect(mockQuery).toHaveBeenNthCalledWith( + 2, + expect.stringMatching(/^SET LOCAL statement_timeout = \d+/), + ); + expect(mockQuery).toHaveBeenNthCalledWith( + 3, + 'EXPLAIN (ANALYZE, FORMAT JSON) SELECT * FROM users WHERE id = $1', + [42], + ); + expect(mockQuery).toHaveBeenNthCalledWith(4, 'ROLLBACK'); + + // Verify client is released + expect(mockRelease).toHaveBeenCalledTimes(1); + }); + + it('applies custom statementTimeoutMs from request body', async () => { + const app = createTestApp(); + const res = await request(app) + .post('/api/admin/db/explain') + .send({ query: 'SELECT 1', statementTimeoutMs: 2500 }); + + expect(res.status).toBe(200); + expect(mockQuery).toHaveBeenCalledWith('SET LOCAL statement_timeout = 2500'); + }); + + it('applies configured statementTimeoutMs from router deps', async () => { + const app = createTestApp({ statementTimeoutMs: 3000 }); + const res = await request(app) + .post('/api/admin/db/explain') + .send({ query: 'SELECT 1' }); + + expect(res.status).toBe(200); + expect(mockQuery).toHaveBeenCalledWith('SET LOCAL statement_timeout = 3000'); }); it('returns raw rows when QUERY PLAN column is absent', async () => { const app = createTestApp(); const rawRows = [{ id: 1, name: 'test' }]; - mockQuery.mockResolvedValueOnce({ rows: rawRows } as unknown as QueryResult); + mockQuery.mockImplementation(async (sql: string) => { + if ( + sql === 'BEGIN READ ONLY' || + sql.startsWith('SET LOCAL statement_timeout') || + sql === 'ROLLBACK' + ) { + return { rows: [] }; + } + return { rows: rawRows } as unknown as QueryResult; + }); + const res = await request(app) .post('/api/admin/db/explain') .send({ query: 'SELECT id, name FROM users LIMIT 1' }); expect(res.status).toBe(200); expect(res.body.plan).toEqual(rawRows); + expect(mockRelease).toHaveBeenCalledTimes(1); }); it('passes parameters to the database query', async () => { const app = createTestApp(); - mockQuery.mockResolvedValueOnce({ rows: [makeExplainRow(SAMPLE_PLAN)] } as unknown as QueryResult); await request(app) .post('/api/admin/db/explain') .send({ query: 'SELECT * FROM users WHERE id = $1 AND status = $2', params: [1, 'active'] }); @@ -226,12 +362,154 @@ describe('POST /api/admin/db/explain', () => { [1, 'active'], ); }); + + it('supports duck-typed pool without connect() method (legacy fallback)', async () => { + const duckTypedPool = { + query: mockQuery, + } as unknown as Pool; + + const app = createTestApp({ pool: duckTypedPool }); + const res = await request(app) + .post('/api/admin/db/explain') + .send({ query: 'SELECT 1' }); + + expect(res.status).toBe(200); + expect(mockQuery).toHaveBeenCalledWith('BEGIN READ ONLY'); + expect(mockQuery).toHaveBeenCalledWith('ROLLBACK'); + }); + }); + + describe('Acceptance Criteria: read-only transaction and statement timeout guarantees', () => { + it('rejects a DELETE inside a CTE before database execution (no rows change)', async () => { + const app = createTestApp(); + const res = await request(app) + .post('/api/admin/db/explain') + .send({ query: 'WITH d AS (DELETE FROM users RETURNING 1) SELECT * FROM d' }); + + expect(res.status).toBe(400); + expect(res.body.code).toBe('BAD_REQUEST'); + expect(res.body.message).toContain('not allowed'); + // No database queries executed — no rows change + expect(mockConnect).not.toHaveBeenCalled(); + expect(mockQuery).not.toHaveBeenCalled(); + }); + + it('rolls back and releases client if a read-only transaction error is returned by DB', async () => { + const app = createTestApp(); + const readOnlyError = new Error('cannot execute DELETE in a read-only transaction'); + (readOnlyError as unknown as Record).code = '25006'; + + mockQuery.mockImplementation(async (sql: string) => { + if (sql === 'BEGIN READ ONLY' || sql.startsWith('SET LOCAL')) { + return { rows: [] }; + } + if (sql === 'ROLLBACK') { + return { rows: [] }; + } + throw readOnlyError; + }); + + const res = await request(app) + .post('/api/admin/db/explain') + .send({ query: 'SELECT 1' }); + + expect(res.status).toBe(400); + expect(res.body.code).toBe('BAD_REQUEST'); + expect(res.body.message).toContain('cannot execute DELETE in a read-only transaction'); + + // Crucial: ROLLBACK was called and client was released + expect(mockQuery).toHaveBeenCalledWith('ROLLBACK'); + expect(mockRelease).toHaveBeenCalledTimes(1); + }); + + it('cancels queries exceeding statement timeout and returns 400 with rollback and release', async () => { + const app = createTestApp(); + const timeoutError = new Error('canceling statement due to statement timeout'); + (timeoutError as unknown as Record).code = '57014'; + + mockQuery.mockImplementation(async (sql: string) => { + if (sql === 'BEGIN READ ONLY' || sql.startsWith('SET LOCAL')) { + return { rows: [] }; + } + if (sql === 'ROLLBACK') { + return { rows: [] }; + } + throw timeoutError; + }); + + const res = await request(app) + .post('/api/admin/db/explain') + .send({ query: 'SELECT pg_sleep(10)' }); + + expect(res.status).toBe(400); + expect(res.body.code).toBe('BAD_REQUEST'); + expect(res.body.message).toContain('statement timeout'); + + // Crucial: transaction is rolled back and client is released back to pool + expect(mockQuery).toHaveBeenCalledWith('ROLLBACK'); + expect(mockRelease).toHaveBeenCalledTimes(1); + }); + + it('always releases the client even when ROLLBACK fails', async () => { + const app = createTestApp(); + mockQuery.mockImplementation(async (sql: string) => { + if (sql === 'BEGIN READ ONLY' || sql.startsWith('SET LOCAL')) { + return { rows: [] }; + } + if (sql.startsWith('EXPLAIN')) { + throw new Error('Connection lost during explain'); + } + if (sql === 'ROLLBACK') { + throw new Error('Connection closed'); + } + return { rows: [] }; + }); + + const res = await request(app) + .post('/api/admin/db/explain') + .send({ query: 'SELECT 1' }); + + expect(res.status).toBe(400); + expect(res.body.message).toContain('Connection lost during explain'); + // Client is STILL released despite rollback throwing + expect(mockRelease).toHaveBeenCalledTimes(1); + }); + + it('supports execution against a ReplicaPool instance', async () => { + const mockReadClient = { + query: jest.fn(async (sql: string) => { + if ( + sql === 'BEGIN READ ONLY' || + sql.startsWith('SET LOCAL statement_timeout') || + sql === 'ROLLBACK' + ) { + return { rows: [] }; + } + return { rows: [makeExplainRow(SAMPLE_PLAN)] }; + }), + release: jest.fn(), + } as unknown as PoolClient; + + const mockReplicaPool = { + getReadClient: jest.fn().mockResolvedValue(mockReadClient), + } as unknown as ReplicaPool; + + const app = createTestApp({ replicaPool: mockReplicaPool }); + const res = await request(app) + .post('/api/admin/db/explain') + .send({ query: 'SELECT * FROM users' }); + + expect(res.status).toBe(200); + expect(mockReplicaPool.getReadClient).toHaveBeenCalledTimes(1); + expect(mockReadClient.query).toHaveBeenCalledWith('BEGIN READ ONLY'); + expect(mockReadClient.query).toHaveBeenCalledWith('ROLLBACK'); + expect(mockReadClient.release).toHaveBeenCalledTimes(1); + }); }); describe('audit logging', () => { - it('logs an audit event on successful explain', async () => { + it('logs an audit event on successful explain including statementTimeoutMs', async () => { const app = createTestApp(); - mockQuery.mockResolvedValueOnce({ rows: [makeExplainRow(SAMPLE_PLAN)] } as unknown as QueryResult); await request(app) .post('/api/admin/db/explain') .set('User-Agent', 'test-agent') @@ -245,6 +523,7 @@ describe('POST /api/admin/db/explain', () => { userAgent: 'test-agent', query: 'SELECT COUNT(*) FROM usage_events', paramCount: 0, + statementTimeoutMs: 5000, }), ); }); @@ -262,29 +541,63 @@ describe('POST /api/admin/db/explain', () => { it('returns 400 when the database query fails', async () => { const app = createTestApp(); - mockQuery.mockRejectedValueOnce(new Error('relation "does_not_exist" does not exist')); + mockQuery.mockImplementation(async (sql: string) => { + if ( + sql === 'BEGIN READ ONLY' || + sql.startsWith('SET LOCAL statement_timeout') || + sql === 'ROLLBACK' + ) { + return { rows: [] }; + } + throw new Error('relation "does_not_exist" does not exist'); + }); + const res = await request(app) .post('/api/admin/db/explain') .send({ query: 'SELECT * FROM does_not_exist' }); expect(res.status).toBe(400); expect(res.body.message).toContain('does not exist'); + expect(mockRelease).toHaveBeenCalledTimes(1); }); it('returns 400 for generic DB error', async () => { const app = createTestApp(); - mockQuery.mockRejectedValueOnce('string error'); + mockQuery.mockImplementation(async (sql: string) => { + if ( + sql === 'BEGIN READ ONLY' || + sql.startsWith('SET LOCAL statement_timeout') || + sql === 'ROLLBACK' + ) { + return { rows: [] }; + } + throw 'string error'; + }); + const res = await request(app) .post('/api/admin/db/explain') .send({ query: 'SELECT 1' }); expect(res.status).toBe(400); expect(res.body.message).toBe('EXPLAIN query execution failed'); + expect(mockRelease).toHaveBeenCalledTimes(1); }); it('recovers after a failed query for subsequent successful queries', async () => { const app = createTestApp(); - mockQuery - .mockRejectedValueOnce(new Error('first failure')) - .mockResolvedValueOnce({ rows: [makeExplainRow(SAMPLE_PLAN)] } as unknown as QueryResult); + let firstFailed = false; + mockQuery.mockImplementation(async (sql: string) => { + if ( + sql === 'BEGIN READ ONLY' || + sql.startsWith('SET LOCAL statement_timeout') || + sql === 'ROLLBACK' + ) { + return { rows: [] }; + } + if (!firstFailed) { + firstFailed = true; + throw new Error('first failure'); + } + return { rows: [makeExplainRow(SAMPLE_PLAN)] } as unknown as QueryResult; + }); const failRes = await request(app) .post('/api/admin/db/explain') @@ -299,20 +612,16 @@ describe('POST /api/admin/db/explain', () => { }); it('propagates unexpected non-ZodError errors to the error handler', async () => { - // Simulate an unexpected error thrown after successful db query (e.g. from logger.audit). - // This covers the final `next(error)` fallthrough in the outer catch block (100% coverage). const unexpectedError = new TypeError('Unexpected runtime error'); (logger.audit as jest.Mock).mockImplementationOnce(() => { throw unexpectedError; }); - mockQuery.mockResolvedValueOnce({ rows: [makeExplainRow(SAMPLE_PLAN)] } as unknown as QueryResult); const app = createTestApp(); const res = await request(app) .post('/api/admin/db/explain') .send({ query: 'SELECT 1' }); - // The unexpected error is passed to next(); errorHandler maps it to 500 expect(res.status).toBe(500); }); }); diff --git a/src/routes/admin/explain.ts b/src/routes/admin/explain.ts index aeac5a4..4cda663 100644 --- a/src/routes/admin/explain.ts +++ b/src/routes/admin/explain.ts @@ -1,47 +1,74 @@ import { Router } from 'express'; -import type { Pool, QueryResult } from 'pg'; +import type { Pool, PoolClient, QueryResult } from 'pg'; import { adminAuth } from '../../middleware/adminAuth.js'; import { createAdminIpAllowlist } from '../../middleware/ipAllowlist.js'; import { BadRequestError, InternalServerError } from '../../errors/index.js'; import { logger } from '../../logger.js'; import { getClientIp } from '../../lib/clientIp.js'; import { validate } from '../../middleware/validate.js'; -import { dbExplainBodySchema, type DbExplainBody } from '../../validators/admin.js'; +import { + dbExplainBodySchema, + isAllowedQuery, + hasMultiStatement, + hasDisallowedDmlKeywords, + ALLOWED_QUERY_PATTERNS, + type DbExplainBody, +} from '../../validators/admin.js'; +import type { ReplicaPool } from '../../db/replicaPool.js'; + +export { + isAllowedQuery, + hasMultiStatement, + hasDisallowedDmlKeywords, + ALLOWED_QUERY_PATTERNS, +}; const TRUST_PROXY = process.env.TRUST_PROXY_HEADERS === 'true'; -const ALLOWED_QUERY_PATTERNS: RegExp[] = [ - /^\s*SELECT\b/is, - /^\s*WITH\b/is, -]; - -function hasMultiStatement(query: string): boolean { - const cleaned = query.replace(/'(?:[^'\\]|\\.)*'/gs, '').replace(/--.*$/gm, ''); - return cleaned.includes(';'); -} - -function isAllowedQuery(query: string): boolean { - if (hasMultiStatement(query)) return false; - return ALLOWED_QUERY_PATTERNS.some((p) => p.test(query)); -} +export const DEFAULT_EXPLAIN_TIMEOUT_MS = 5_000; export interface ExplainRouterDeps { pool?: Pool; + replicaPool?: ReplicaPool; + statementTimeoutMs?: number; +} + +async function acquireClient(deps: ExplainRouterDeps): Promise { + if (deps.replicaPool) { + return deps.replicaPool.getReadClient(); + } + + if (deps.pool) { + if (typeof deps.pool.connect === 'function') { + return deps.pool.connect(); + } + // Duck-typed fallback for test mocks that only define query + const stubPool = deps.pool as unknown as { query: (text: string, params?: unknown[]) => Promise }; + return { + query: stubPool.query.bind(stubPool), + release: () => {}, + } as unknown as PoolClient; + } + + throw new InternalServerError('Database pool not available'); } /** * Router exposing `POST /api/admin/db/explain` — runs - * `EXPLAIN (ANALYZE, FORMAT JSON)` on a read-only SQL query and returns the - * query plan for diagnostic use. + * `EXPLAIN (ANALYZE, FORMAT JSON)` on a read-only SQL query inside a dedicated + * read-only transaction (`BEGIN READ ONLY; SET LOCAL statement_timeout = ...; ROLLBACK`) + * and returns the query plan for diagnostic use. * * Admin-only: gated behind the admin IP allowlist and admin authentication. * - * Request body is validated by {@link dbExplainBodySchema} via the - * {@link validate} middleware, which returns a structured - * `{ code, message, details }` 400 response for any invalid input. - * - * Only `SELECT` and `WITH` (CTE) queries are accepted; multi-statement - * queries are rejected at the application layer as an extra safety guard. + * Safety measures: + * - Request body is validated by {@link dbExplainBodySchema} at the boundary. + * - Multi-statement queries are strictly rejected. + * - Data-modifying keywords (DELETE, UPDATE, INSERT, MERGE, DDL) in CTEs or + * subqueries are rejected before database execution. + * - Executes on a dedicated client in an explicit `BEGIN READ ONLY` transaction. + * - Sets a local statement_timeout to cancel runaway queries or pg_sleep calls. + * - Guaranteed transaction `ROLLBACK` and `client.release()` on all code paths. */ export function createExplainRouter(deps: ExplainRouterDeps = {}): Router { const router = Router(); @@ -51,17 +78,11 @@ export function createExplainRouter(deps: ExplainRouterDeps = {}): Router { router.post( '/', - // ── Input validation at the boundary ────────────────────────────────── - // dbExplainBodySchema enforces: query non-empty ≤ 50 000 chars, - // params is an array (defaults to []). Any violation returns a - // structured 400 before the query parser or database are touched. validate({ body: dbExplainBodySchema }), async (req, res, next) => { try { - // Re-parse to pick up Zod defaults (e.g. params defaults to []). - // validate() already confirmed the shape is valid; this is zero-cost. const parsed = dbExplainBodySchema.parse(req.body); - const { query: rawQuery, params } = parsed; + const { query: rawQuery, params, statementTimeoutMs: bodyTimeout } = parsed; if (!isAllowedQuery(rawQuery)) { next( @@ -72,22 +93,47 @@ export function createExplainRouter(deps: ExplainRouterDeps = {}): Router { return; } - const { pool } = deps; - if (!pool) { - next(new InternalServerError('Database pool not available')); + const effectiveTimeoutMs = + bodyTimeout ?? + deps.statementTimeoutMs ?? + (process.env.ADMIN_EXPLAIN_TIMEOUT_MS + ? parseInt(process.env.ADMIN_EXPLAIN_TIMEOUT_MS, 10) + : DEFAULT_EXPLAIN_TIMEOUT_MS); + + let client: PoolClient; + try { + client = await acquireClient(deps); + } catch (acquireError) { + next(acquireError); return; } const explainSql = `EXPLAIN (ANALYZE, FORMAT JSON) ${rawQuery}`; let result: QueryResult; + let rolledBack = false; try { - result = await pool.query(explainSql, params); + await client.query('BEGIN READ ONLY'); + await client.query(`SET LOCAL statement_timeout = ${Math.floor(effectiveTimeoutMs)}`); + result = await client.query(explainSql, params); + await client.query('ROLLBACK'); + rolledBack = true; } catch (dbError) { + if (!rolledBack) { + try { + await client.query('ROLLBACK'); + rolledBack = true; + } catch { + // Rollback may fail if connection was dropped + } + } + const message = dbError instanceof Error ? dbError.message : 'EXPLAIN query execution failed'; next(new BadRequestError(message)); return; + } finally { + client.release(); } const plan = @@ -103,6 +149,7 @@ export function createExplainRouter(deps: ExplainRouterDeps = {}): Router { userAgent, query: rawQuery, paramCount: params.length, + statementTimeoutMs: Math.floor(effectiveTimeoutMs), }); res.json({ plan }); @@ -116,3 +163,4 @@ export function createExplainRouter(deps: ExplainRouterDeps = {}): Router { } export default createExplainRouter; + diff --git a/src/validators/admin.ts b/src/validators/admin.ts index 9680fdd..6921d36 100644 --- a/src/validators/admin.ts +++ b/src/validators/admin.ts @@ -245,6 +245,80 @@ export type SpikeQuery = z.infer; // POST /api/admin/db/explain // --------------------------------------------------------------------------- +export const ALLOWED_QUERY_PATTERNS: RegExp[] = [ + /^\s*SELECT\b/is, + /^\s*WITH\b/is, +]; + +export const DISALLOWED_DML_KEYWORDS = [ + 'DELETE', + 'UPDATE', + 'INSERT', + 'MERGE', + 'DROP', + 'ALTER', + 'TRUNCATE', + 'CREATE', + 'GRANT', + 'REVOKE', + 'LOCK', + 'VACUUM', + 'REINDEX', + 'EXECUTE', + 'CALL', + 'DO', +] as const; + +export const DISALLOWED_KEYWORDS_REGEX = new RegExp( + `\\b(${DISALLOWED_DML_KEYWORDS.join('|')})\\b`, + 'i', +); + +/** + * Strips SQL comments and string literals so keyword checking does not trigger + * on string literals (e.g. `WHERE status = 'DELETE'`) or comments. + */ +export function stripSqlLiteralsAndComments(query: string): string { + return query + // Remove single-quoted strings (handling escaped quotes '') + .replace(/'(?:[^'\\]|\\.)*'/gs, "''") + // Remove dollar-quoted strings: $$...$$ or $tag$...$tag$ + .replace(/\$([a-zA-Z0-9_]*)\$[\s\S]*?\$\1\$/g, "''") + // Remove single-line comments: -- ... + .replace(/--.*$/gm, '') + // Remove multi-line comments: /* ... */ + .replace(/\/\*[\s\S]*?\*\//g, ''); +} + +/** + * Checks whether the query contains multiple statements separated by semicolons. + */ +export function hasMultiStatement(query: string): boolean { + const cleaned = stripSqlLiteralsAndComments(query); + return cleaned.includes(';'); +} + +/** + * Checks whether the query contains data-modifying or DDL keywords outside of + * literals and comments. + */ +export function hasDisallowedDmlKeywords(query: string): boolean { + const cleaned = stripSqlLiteralsAndComments(query); + return DISALLOWED_KEYWORDS_REGEX.test(cleaned); +} + +/** + * Validates that an EXPLAIN query is read-only and safe to analyze. + * Rejects multi-statement queries, queries not starting with SELECT or WITH, + * and queries containing data-modifying keywords (e.g. data-modifying CTEs). + */ +export function isAllowedQuery(query: string): boolean { + if (hasMultiStatement(query)) return false; + if (!ALLOWED_QUERY_PATTERNS.some((p) => p.test(query))) return false; + if (hasDisallowedDmlKeywords(query)) return false; + return true; +} + /** * Request body for POST /api/admin/db/explain. * Kept in sync with the inline schema in explain.ts so both can share the @@ -255,6 +329,8 @@ export const dbExplainBodySchema = z.object({ query: z.string().min(1, 'Query is required').max(50_000, 'Query too long'), /** Optional positional parameters to pass to the query. */ params: z.array(z.unknown()).optional().default([]), + /** Optional statement timeout in milliseconds. */ + statementTimeoutMs: z.number().int().positive().max(60_000).optional(), }); export type DbExplainBody = z.infer;