Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
561 changes: 561 additions & 0 deletions README.md

Large diffs are not rendered by default.

56 changes: 36 additions & 20 deletions docs/admin-db-explain.md
Original file line number Diff line number Diff line change
Expand Up @@ -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
```

---
Expand All @@ -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
}
```

Expand All @@ -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) <query>`. 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) <query>` 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.

83 changes: 83 additions & 0 deletions docs/tamper-evident-audit.md
Original file line number Diff line number Diff line change
@@ -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.
5 changes: 5 additions & 0 deletions jest.env-setup.cjs
Original file line number Diff line number Diff line change
@@ -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";
Original file line number Diff line number Diff line change
@@ -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
Expand Down
85 changes: 85 additions & 0 deletions package.json
Original file line number Diff line number Diff line change
@@ -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"
}
}
32 changes: 32 additions & 0 deletions src/db/replicaPool.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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,
};
Expand Down Expand Up @@ -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 ──────────────────────────────────────────────────
Expand Down Expand Up @@ -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 ──────────────────────────────────────────────────
Expand Down
34 changes: 32 additions & 2 deletions src/db/replicaPool.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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';
Expand All @@ -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 ───────────────────────────────────────────────────────

Expand Down Expand Up @@ -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<PoolClient> {
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.
Expand Down
Loading
Loading