From 4fad34858d9aaf28232baf382a4e63de16eebc3a Mon Sep 17 00:00:00 2001 From: Jimoh Abdullahi Date: Tue, 29 Sep 2026 17:20:59 +0100 Subject: [PATCH 1/4] fix(migrations): resolve 0024 prefix collision and add destructive-approved tag --- ...tore_scope.down.sql => 0025_idempotency_store_scope.down.sql} | 0 ...mpotency_store_scope.sql => 0025_idempotency_store_scope.sql} | 1 + 2 files changed, 1 insertion(+) rename migrations/{0024_idempotency_store_scope.down.sql => 0025_idempotency_store_scope.down.sql} (100%) rename migrations/{0024_idempotency_store_scope.sql => 0025_idempotency_store_scope.sql} (97%) 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 From 23afba0f92517a338b4384f7f063e946a120ca1b Mon Sep 17 00:00:00 2001 From: Jimoh Abdullahi Date: Tue, 29 Sep 2026 17:33:15 +0100 Subject: [PATCH 2/4] chore: restore inadvertently deleted package.json, README.md, and jest.env-setup.cjs --- README.md | 561 +++++++++++++++++++++++++++++++++++++++++++++ jest.env-setup.cjs | 5 + package.json | 85 +++++++ 3 files changed, 651 insertions(+) create mode 100644 README.md create mode 100644 jest.env-setup.cjs create mode 100644 package.json 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/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/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" + } +} From 07d81a949f30cf8af2f9e8e01045bb9878ab379a Mon Sep 17 00:00:00 2001 From: Jimoh Abdullahi Date: Tue, 29 Sep 2026 17:34:20 +0100 Subject: [PATCH 3/4] chore: restore inadvertently deleted adminAuth.ts and tamper-evident-audit.md --- docs/tamper-evident-audit.md | 83 ++++++++++++++++++++++++++++++++++++ src/middleware/adminAuth.ts | 55 ++++++++++++++++++++++++ 2 files changed, 138 insertions(+) create mode 100644 docs/tamper-evident-audit.md create mode 100644 src/middleware/adminAuth.ts 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/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')); +} From 2ebac44472a4d176c35fe22ddba2ff661c1c2f75 Mon Sep 17 00:00:00 2001 From: Jimoh Abdullahi Date: Tue, 29 Sep 2026 17:41:17 +0100 Subject: [PATCH 4/4] feat(webhooks): disable redirects and revalidate hosts for webhook dispatch Closes #1262 - Pass redirect: 'manual' in fetch options to prevent following 3xx redirects to internal or cloud metadata endpoints - Re-run validateWebhookUrl before each dispatch to prevent SSRF and DNS rebinding to private IP ranges - Cap response body reads via consumeCappedResponseBody to prevent memory exhaustion - Record failed deliveries in WebhookStore with clear explanatory reasons - Add tests with local HTTP server returning redirects (301, 302, 307) and verifying they are not followed - Add tests for private IP refusal at dispatch time and response body capping --- src/webhooks/webhook.dispatcher.test.ts | 227 +++++++++++++++++- src/webhooks/webhook.dispatcher.ts | 122 +++++++++- src/webhooks/webhook.validator.test.ts | 125 +++++++++- src/webhooks/webhook.validator.ts | 42 +++- .../webhook-dispatch-pipeline.test.ts | 8 + tests/integration/webhooks.test.ts | 1 + 6 files changed, 505 insertions(+), 20 deletions(-) diff --git a/src/webhooks/webhook.dispatcher.test.ts b/src/webhooks/webhook.dispatcher.test.ts index 3a5ef99..06b32bc 100644 --- a/src/webhooks/webhook.dispatcher.test.ts +++ b/src/webhooks/webhook.dispatcher.test.ts @@ -1,7 +1,30 @@ -import { dispatchWebhook, dispatchToAll, resetWebhookDispatcherForTests, stopWebhookDispatching } from './webhook.dispatcher.js'; +import http from 'http'; +import { AddressInfo } from 'net'; +import { + dispatchWebhook, + dispatchToAll, + resetWebhookDispatcherForTests, + stopWebhookDispatching, + consumeCappedResponseBody, + MAX_WEBHOOK_RESPONSE_BYTES, +} from './webhook.dispatcher.js'; import { WebhookStore } from './webhook.store.js'; import type { WebhookConfig, WebhookPayload } from './webhook.types.js'; +// Mock DNS lookup so URL validation resolves deterministically with fake timers +// eslint-disable-next-line no-var +var mockDnsLookup = jest.fn().mockImplementation(async (hostname: string) => { + if (hostname === '127.0.0.1' || hostname === 'localhost') { + return [{ address: '127.0.0.1', family: 4 }]; + } + return [{ address: '93.184.216.34', family: 4 }]; +}); + +jest.mock('dns/promises', () => { + const lookupFn = (...args: unknown[]) => mockDnsLookup(...args); + return { __esModule: true, default: { lookup: lookupFn }, lookup: lookupFn }; +}); + describe('Webhook Dispatcher', () => { let originalFetch: typeof global.fetch; @@ -352,4 +375,206 @@ describe('Webhook Dispatcher', () => { expect(fetchMock).toHaveBeenCalledTimes(1); }); }); + + describe('SSRF Protection & Redirect Refusal (Issue #1262)', () => { + let server: http.Server; + let serverUrl: string; + let receivedRequests: Array<{ method?: string; url?: string; headers: http.IncomingHttpHeaders }>; + let originalEnv: string | undefined; + + beforeEach(async () => { + originalEnv = process.env.NODE_ENV; + jest.useRealTimers(); + WebhookStore.clearFailedDeliveries(); + receivedRequests = []; + + await new Promise((resolve) => { + server = http.createServer((req, res) => { + receivedRequests.push({ method: req.method, url: req.url, headers: req.headers }); + + if (req.url === '/redirect-metadata-302') { + res.writeHead(302, { Location: 'http://169.254.169.254/latest/meta-data' }); + res.end('Redirecting to cloud metadata'); + return; + } + + if (req.url === '/redirect-internal-301') { + res.writeHead(301, { Location: 'http://127.0.0.1:8080/admin/secrets' }); + res.end('Redirecting to internal admin'); + return; + } + + if (req.url === '/redirect-307') { + res.writeHead(307, { Location: 'http://10.0.0.1/private' }); + res.end('Temporary redirect'); + return; + } + + if (req.url === '/large-response') { + res.writeHead(200, { 'Content-Type': 'text/plain' }); + const chunk = 'A'.repeat(16 * 1024); + for (let i = 0; i < 16; i++) { + res.write(chunk); + } + res.end(); + return; + } + + if (req.url === '/success') { + res.writeHead(200, { 'Content-Type': 'application/json' }); + res.end(JSON.stringify({ received: true })); + return; + } + + res.writeHead(404); + res.end('Not found'); + }); + + server.listen(0, '127.0.0.1', () => { + const address = server.address() as AddressInfo; + serverUrl = `http://127.0.0.1:${address.port}`; + resolve(); + }); + }); + }); + + afterEach(async () => { + process.env.NODE_ENV = originalEnv; + if (server) { + await new Promise((resolve) => server.close(() => resolve())); + } + }); + + it('does not follow 302 redirect to an internal metadata address and records failure reason', async () => { + const redirectConfig: WebhookConfig = { + developerId: 'dev_ssrf_test', + url: `${serverUrl}/redirect-metadata-302`, + events: ['new_api_call'], + createdAt: new Date(), + }; + + await dispatchWebhook(redirectConfig, payload); + + expect(receivedRequests.length).toBe(1); + expect(receivedRequests[0].url).toBe('/redirect-metadata-302'); + + const failures = WebhookStore.getRecentFailures(); + const failure = failures.find((f) => f.url === redirectConfig.url); + expect(failure).toBeDefined(); + expect(failure?.lastError).toContain('HTTP 302'); + expect(failure?.lastError).toContain('http://169.254.169.254/latest/meta-data'); + expect(failure?.lastError).toContain('redirects are not followed'); + }); + + it('does not follow 301 redirect to internal service and records failure', async () => { + const redirectConfig: WebhookConfig = { + developerId: 'dev_ssrf_test', + url: `${serverUrl}/redirect-internal-301`, + events: ['new_api_call'], + createdAt: new Date(), + }; + + await dispatchWebhook(redirectConfig, payload); + + expect(receivedRequests.length).toBe(1); + expect(receivedRequests[0].url).toBe('/redirect-internal-301'); + + const failures = WebhookStore.getRecentFailures(); + const failure = failures.find((f) => f.url === redirectConfig.url); + expect(failure).toBeDefined(); + expect(failure?.lastError).toContain('HTTP 301'); + expect(failure?.lastError).toContain('redirects are not followed'); + }); + + it('does not follow 307 temporary redirect', async () => { + const redirectConfig: WebhookConfig = { + developerId: 'dev_ssrf_test', + url: `${serverUrl}/redirect-307`, + events: ['new_api_call'], + createdAt: new Date(), + }; + + await dispatchWebhook(redirectConfig, payload); + + expect(receivedRequests.length).toBe(1); + const failures = WebhookStore.getRecentFailures(); + const failure = failures.find((f) => f.url === redirectConfig.url); + expect(failure).toBeDefined(); + expect(failure?.lastError).toContain('HTTP 307'); + }); + + it('delivers successfully to 200 OK endpoint on local server', async () => { + const successConfig: WebhookConfig = { + developerId: 'dev_ssrf_test', + url: `${serverUrl}/success`, + events: ['new_api_call'], + createdAt: new Date(), + }; + + await dispatchWebhook(successConfig, payload); + + expect(receivedRequests.length).toBe(1); + expect(receivedRequests[0].url).toBe('/success'); + const failures = WebhookStore.getRecentFailures(); + expect(failures.find((f) => f.url === successConfig.url)).toBeUndefined(); + }); + + it('caps response body reads to MAX_WEBHOOK_RESPONSE_BYTES (64 KB)', async () => { + const largeBodyConfig: WebhookConfig = { + developerId: 'dev_ssrf_test', + url: `${serverUrl}/large-response`, + events: ['new_api_call'], + createdAt: new Date(), + }; + + await dispatchWebhook(largeBodyConfig, payload); + expect(receivedRequests.length).toBe(1); + + const response = await fetch(`${serverUrl}/large-response`); + const consumed = await consumeCappedResponseBody(response, MAX_WEBHOOK_RESPONSE_BYTES); + expect(Buffer.byteLength(consumed, 'utf8')).toBeLessThanOrEqual(MAX_WEBHOOK_RESPONSE_BYTES); + }); + + it('refuses dispatch at dispatch time when DNS resolves to private range', async () => { + process.env.NODE_ENV = 'production'; + + mockDnsLookup.mockResolvedValueOnce([ + { address: '169.254.169.254', family: 4 }, + ]); + + const privateDnsConfig: WebhookConfig = { + developerId: 'dev_ssrf_test', + url: 'https://dynamic-rebind.example.com/webhook', + events: ['new_api_call'], + createdAt: new Date(), + }; + + await dispatchWebhook(privateDnsConfig, payload); + + const failures = WebhookStore.getRecentFailures(); + const failure = failures.find((f) => f.url === privateDnsConfig.url); + expect(failure).toBeDefined(); + expect(failure?.lastError).toContain('resolves to a private/internal IP address (169.254.169.254)'); + expect(failure?.attempts).toBe(0); + }); + + it('refuses dispatch at dispatch time when hostname DNS fails to resolve', async () => { + mockDnsLookup.mockRejectedValueOnce(new Error('getaddrinfo ENOTFOUND invalid.domain')); + + const invalidDnsConfig: WebhookConfig = { + developerId: 'dev_ssrf_test', + url: 'https://invalid.domain/webhook', + events: ['new_api_call'], + createdAt: new Date(), + }; + + await dispatchWebhook(invalidDnsConfig, payload); + + const failures = WebhookStore.getRecentFailures(); + const failure = failures.find((f) => f.url === invalidDnsConfig.url); + expect(failure).toBeDefined(); + expect(failure?.lastError).toContain('Could not resolve webhook hostname.'); + }); + }); }); + diff --git a/src/webhooks/webhook.dispatcher.ts b/src/webhooks/webhook.dispatcher.ts index 00ffcbb..2352cc5 100644 --- a/src/webhooks/webhook.dispatcher.ts +++ b/src/webhooks/webhook.dispatcher.ts @@ -4,6 +4,70 @@ import { WebhookStore } from './webhook.store.js'; import { logger } from '../logger.js'; import { getCorrelationId, getRequestId } from '../utils/asyncContext.js'; import { getEffectiveRetryPolicy } from '../services/webhookRetry.js'; +import { validateWebhookUrl, WebhookValidationError } from './webhook.validator.js'; + +export const MAX_WEBHOOK_RESPONSE_BYTES = 64 * 1024; + +export async function consumeCappedResponseBody( + response: Response, + maxBytes: number = MAX_WEBHOOK_RESPONSE_BYTES +): Promise { + if (!response.body) { + if (typeof (response as any).text === 'function') { + try { + const text = await (response as any).text(); + return typeof text === 'string' && text.length > maxBytes + ? text.slice(0, maxBytes) + : text || ''; + } catch { + return ''; + } + } + return ''; + } + + try { + const reader = response.body.getReader(); + let receivedBytes = 0; + const chunks: Uint8Array[] = []; + + try { + while (true) { + const { done, value } = await reader.read(); + if (done) break; + if (value) { + if (receivedBytes + value.byteLength > maxBytes) { + const remaining = maxBytes - receivedBytes; + if (remaining > 0) { + chunks.push(value.subarray(0, remaining)); + } + try { + await reader.cancel('Response body exceeded maximum allowed size'); + } catch { + // ignore cancel error + } + break; + } + chunks.push(value); + receivedBytes += value.byteLength; + } + } + } catch { + // ignore stream read error + } finally { + try { + reader.releaseLock(); + } catch { + // ignore release error + } + } + + return Buffer.concat(chunks).toString('utf8'); + } catch { + return ''; + } +} + let acceptingDispatches = true; const inFlightDispatches = new Set>(); @@ -83,17 +147,70 @@ export async function dispatchWebhook( } let lastError: unknown; + let attemptsMade = 0; + + try { + await validateWebhookUrl(config.url); + } catch (err) { + const failedAt = new Date().toISOString(); + const lastErrorMessage = + err instanceof Error ? err.message : String(err); + + logger.warn( + `[webhook] Pre-dispatch validation failed for ${config.url}:`, + lastErrorMessage + ); + logger.error( + `[webhook] ✗ Failed to deliver ${payload.event} to ${config.url} after 0 attempts.`, + err + ); + + WebhookStore.recordFailedDelivery({ + deliveryId, + developerId: config.developerId, + event: payload.event, + url: config.url, + failedAt, + lastError: lastErrorMessage, + attempts: 0, + }); + return; + } for (let attempt = 0; attempt < maxRetries; attempt++) { + attemptsMade = attempt + 1; try { const response = await fetch(config.url, { method: 'POST', body, headers, + redirect: 'manual', signal: AbortSignal.timeout(10_000), // 10s timeout per attempt }); + const isRedirect = + (response.status >= 300 && response.status < 400) || + response.type === 'opaqueredirect'; + + if (isRedirect) { + if (response.body) { + await consumeCappedResponseBody(response); + } + const location = response.headers.get('location') || ''; + const redirectMsg = location + ? `Webhook redirect to "${location}" refused (HTTP ${response.status}): redirects are not followed` + : `Webhook HTTP ${response.status} redirect refused: redirects are not followed`; + lastError = new Error(redirectMsg); + logger.warn( + `[webhook] ${redirectMsg} for ${config.url}, attempt ${attempt + 1}` + ); + break; + } + if (response.ok) { + if (response.body) { + await consumeCappedResponseBody(response); + } logger.info( `[webhook] ✓ Delivered ${payload.event} to ${config.url}`, `attempt ${attempt + 1}` @@ -101,6 +218,9 @@ export async function dispatchWebhook( return; } + if (response.body) { + await consumeCappedResponseBody(response); + } lastError = new Error(`HTTP ${response.status} ${response.statusText}`); logger.warn( `[webhook] Non-2xx response (${response.status}) for ${config.url}`, @@ -138,7 +258,7 @@ export async function dispatchWebhook( url: config.url, failedAt, lastError: lastErrorMessage, - attempts: maxRetries, + attempts: attemptsMade || maxRetries, }); })()); } diff --git a/src/webhooks/webhook.validator.test.ts b/src/webhooks/webhook.validator.test.ts index 56f93ab..1db3438 100644 --- a/src/webhooks/webhook.validator.test.ts +++ b/src/webhooks/webhook.validator.test.ts @@ -1,11 +1,120 @@ -// Webhook URL validation is tested via integration tests in tests/integration/webhooks.test.ts. -// This file is intentionally minimal — it exists to satisfy the project test file convention -// for the webhook.validator module. +import dns from 'dns/promises'; +import { + validateWebhookUrl, + WebhookValidationError, + BLOCKED_RANGES, +} from './webhook.validator.js'; describe('webhook.validator module', () => { - it('exists and is importable', async () => { - const mod = await import('./webhook.validator.js'); - expect(mod).toBeDefined(); - expect(typeof mod.validateWebhookUrl).toBe('function'); - }); + let originalEnv: string | undefined; + + beforeEach(() => { + originalEnv = process.env.NODE_ENV; + delete process.env.WEBHOOK_ENFORCE_PRIVATE_IP_CHECK; + }); + + afterEach(() => { + process.env.NODE_ENV = originalEnv; + delete process.env.WEBHOOK_ENFORCE_PRIVATE_IP_CHECK; + jest.restoreAllMocks(); + }); + + it('exists and is importable', async () => { + const mod = await import('./webhook.validator.js'); + expect(mod).toBeDefined(); + expect(typeof mod.validateWebhookUrl).toBe('function'); + expect(Array.isArray(BLOCKED_RANGES)).toBe(true); + }); + + it('rejects invalid URL strings', async () => { + await expect(validateWebhookUrl('not-a-valid-url')).rejects.toThrow( + WebhookValidationError + ); + await expect(validateWebhookUrl('not-a-valid-url')).rejects.toThrow( + 'Invalid URL format.' + ); + }); + + it('rejects unsupported protocols (e.g. ftp, file, gopher)', async () => { + await expect(validateWebhookUrl('ftp://example.com/hook')).rejects.toThrow( + 'Webhook URL must use HTTP or HTTPS protocol.' + ); + await expect(validateWebhookUrl('file:///etc/passwd')).rejects.toThrow( + 'Webhook URL must use HTTP or HTTPS protocol.' + ); + }); + + it('requires HTTPS in production', async () => { + process.env.NODE_ENV = 'production'; + await expect(validateWebhookUrl('http://example.com/webhook')).rejects.toThrow( + 'Webhook URL must use HTTPS in production.' + ); + }); + + it('rejects non-standard ports in production', async () => { + process.env.NODE_ENV = 'production'; + jest.spyOn(dns, 'lookup').mockResolvedValue([ + { address: '93.184.216.34', family: 4 }, + ] as any); + + await expect( + validateWebhookUrl('https://example.com:8443/webhook') + ).rejects.toThrow('Only ports 80 and 443 are allowed in production.'); + }); + + it('throws when hostname fails DNS resolution', async () => { + jest.spyOn(dns, 'lookup').mockRejectedValue(new Error('ENOTFOUND')); + await expect( + validateWebhookUrl('https://unresolvable.invalid/webhook') + ).rejects.toThrow('Could not resolve webhook hostname.'); + }); + + describe('SSRF / Private IP blocking', () => { + const privateIps = [ + '10.0.0.1', + '172.16.0.1', + '172.31.255.254', + '192.168.1.1', + '127.0.0.1', + '169.254.169.254', + '100.64.0.1', + '::1', + 'fc00::1', + ]; + + it.each(privateIps)('rejects private IP %s in production', async (ip) => { + process.env.NODE_ENV = 'production'; + jest.spyOn(dns, 'lookup').mockResolvedValue([ + { address: ip, family: ip.includes(':') ? 6 : 4 }, + ] as any); + + await expect( + validateWebhookUrl('https://receiver.example.com/webhook') + ).rejects.toThrow(`Webhook URL resolves to a private/internal IP address (${ip}), which is not allowed.`); + }); + + it('rejects private IPs when enforcePrivateIpCheck option is passed even in non-production', async () => { + process.env.NODE_ENV = 'development'; + jest.spyOn(dns, 'lookup').mockResolvedValue([ + { address: '169.254.169.254', family: 4 }, + ] as any); + + await expect( + validateWebhookUrl('http://receiver.example.com/webhook', { + enforcePrivateIpCheck: true, + }) + ).rejects.toThrow('resolves to a private/internal IP address (169.254.169.254)'); + }); + + it('allows public IP address in production', async () => { + process.env.NODE_ENV = 'production'; + jest.spyOn(dns, 'lookup').mockResolvedValue([ + { address: '93.184.216.34', family: 4 }, + ] as any); + + await expect( + validateWebhookUrl('https://public-receiver.example.com/webhook') + ).resolves.toBeUndefined(); + }); + }); }); diff --git a/src/webhooks/webhook.validator.ts b/src/webhooks/webhook.validator.ts index 4b1675f..a1297ec 100644 --- a/src/webhooks/webhook.validator.ts +++ b/src/webhooks/webhook.validator.ts @@ -2,7 +2,7 @@ import { URL } from 'url'; import dns from 'dns/promises'; import ipRangeCheck from 'ip-range-check'; -const BLOCKED_RANGES = [ +export const BLOCKED_RANGES = [ '10.0.0.0/8', '172.16.0.0/12', '192.168.0.0/16', @@ -16,9 +16,21 @@ const BLOCKED_RANGES = [ '240.0.0.0/4', ]; -export class WebhookValidationError extends Error {} +export class WebhookValidationError extends Error { + constructor(message: string) { + super(message); + this.name = 'WebhookValidationError'; + } +} + +export interface WebhookValidationOptions { + enforcePrivateIpCheck?: boolean; +} -export async function validateWebhookUrl(rawUrl: string): Promise { +export async function validateWebhookUrl( + rawUrl: string, + options?: WebhookValidationOptions +): Promise { let parsed: URL; // 1. Must be a valid URL @@ -28,6 +40,10 @@ export async function validateWebhookUrl(rawUrl: string): Promise { throw new WebhookValidationError('Invalid URL format.'); } + if (parsed.protocol !== 'http:' && parsed.protocol !== 'https:') { + throw new WebhookValidationError('Webhook URL must use HTTP or HTTPS protocol.'); + } + // 2. Must use HTTPS in production const isProduction = process.env.NODE_ENV === 'production'; if (isProduction && parsed.protocol !== 'https:') { @@ -43,13 +59,18 @@ export async function validateWebhookUrl(rawUrl: string): Promise { throw new WebhookValidationError('Could not resolve webhook hostname.'); } - if (isProduction) { + const shouldCheckPrivateIps = + isProduction || + options?.enforcePrivateIpCheck === true || + process.env.WEBHOOK_ENFORCE_PRIVATE_IP_CHECK === 'true'; + + if (shouldCheckPrivateIps) { for (const ip of addresses) { - if (ipRangeCheck(ip, BLOCKED_RANGES)) { - throw new WebhookValidationError( - `Webhook URL resolves to a private/internal IP address (${ip}), which is not allowed.` - ); - } + if (ipRangeCheck(ip, BLOCKED_RANGES)) { + throw new WebhookValidationError( + `Webhook URL resolves to a private/internal IP address (${ip}), which is not allowed.` + ); + } } } @@ -57,4 +78,5 @@ export async function validateWebhookUrl(rawUrl: string): Promise { if (isProduction && parsed.port && !['80', '443'].includes(parsed.port)) { throw new WebhookValidationError('Only ports 80 and 443 are allowed in production.'); } - } +} + diff --git a/tests/integration/webhook-dispatch-pipeline.test.ts b/tests/integration/webhook-dispatch-pipeline.test.ts index 097a470..65b9d95 100644 --- a/tests/integration/webhook-dispatch-pipeline.test.ts +++ b/tests/integration/webhook-dispatch-pipeline.test.ts @@ -1,6 +1,14 @@ import { WebhookStore } from '../../src/webhooks/webhook.store.js'; import { calloraEvents } from '../../src/events/event.emitter.js'; +jest.mock('dns/promises', () => ({ + __esModule: true, + default: { + lookup: jest.fn().mockResolvedValue([{ address: '93.184.216.34', family: 4 }]), + }, + lookup: jest.fn().mockResolvedValue([{ address: '93.184.216.34', family: 4 }]), +})); + async function flushAsyncEventHandlers(): Promise { await Promise.resolve(); await Promise.resolve(); diff --git a/tests/integration/webhooks.test.ts b/tests/integration/webhooks.test.ts index 0c22253..a2abd3c 100644 --- a/tests/integration/webhooks.test.ts +++ b/tests/integration/webhooks.test.ts @@ -591,6 +591,7 @@ describe('Webhook Signature Verification Tests', () => { 'X-Callora-Timestamp': testPayload.timestamp, 'X-Callora-Signature': `sha256=${expectedSignature}`, }), + redirect: 'manual', signal: expect.any(AbortSignal), }); });