@j0nathan-ll0yd/database
Multi-provider Drizzle ORM client for AWS Lambda with connection pooling, observability, and automatic OCC retry.
Supported Providers
| Provider | Use Case |
|---|---|
aurora-dsql | AWS Aurora DSQL (serverless, zero-management) |
aurora-serverless-v2 | AWS Aurora Serverless v2 (Data API) |
neon | Neon Postgres (serverless, branching) |
Getting a Client
getDrizzleClient(config)
Get or create a Drizzle client. Caches the client for Lambda container reuse.
import { getDrizzleClient } from "@j0nathan-ll0yd/database";
import { getRequiredEnv } from "@j0nathan-ll0yd/env";
// Aurora DSQL
const db = await getDrizzleClient({
provider: "aurora-dsql",
endpoint: getRequiredEnv("DSQL_ENDPOINT"),
region: getRequiredEnv("AWS_REGION"),
});
// Aurora Serverless v2
const db = await getDrizzleClient({
provider: "aurora-serverless-v2",
connectionString: getRequiredEnv("DATABASE_URL"),
});
// Neon
const db = await getDrizzleClient({
provider: "neon",
connectionString: getRequiredEnv("DATABASE_URL"),
});closeDrizzleClient(config?)
Close a cached database client. Call during graceful shutdown. Pass the same config used to create the client, or omit to close the default client.
Transactions
withTransaction(db, fn)
Execute a function within a database transaction. Rolls back automatically on error.
import { withTransaction } from "@j0nathan-ll0yd/database";
await withTransaction(db, async (tx) => {
await tx.insert(users).values({ name: "Alice" });
await tx.insert(audit).values({ action: "user_created" });
});Query Instrumentation
withQueryMetrics(queryName, fn, options?)
Wrap a database query with automatic CloudWatch metrics, X-Ray tracing, structured error logging, and OCC retry.
import { withQueryMetrics } from "@j0nathan-ll0yd/database";
const users = await withQueryMetrics("Users.getActive", async () => {
return db.select().from(usersTable).where(eq(usersTable.active, true));
});Signature:
function withQueryMetrics<T>(
queryName: string,
fn: () => Promise<T>,
options?: { logInput?: unknown; retry?: DsqlRetryOptions | false },
): Promise<T>;queryName-- used as the CloudWatch dimension and X-Ray span name.logInput-- when provided, logs query parameters at DEBUG level (sanitized).retry-- OCC retry is enabled by default. Passfalseto disable, or aDsqlRetryOptionsobject to customise.
Throws DatabaseError when the underlying query fails.
onConnectionInvalidated(callback)
Register a callback for connection invalidation events (useful for Aurora DSQL token refresh).
OCC Retry
withDsqlRetry(fn, options?)
Execute an async function with retry logic for Aurora DSQL OCC conflict errors (40001, OC000, OC001). Non-OCC errors are re-thrown immediately.
withQueryMetrics uses this internally by default -- you only need withDsqlRetry directly for standalone database calls outside entity queries.
import { withDsqlRetry } from "@j0nathan-ll0yd/database";
const result = await withDsqlRetry(() => db.select().from(users).where(eq(users.id, id)), {
maxRetries: 5,
baseDelayMs: 50,
});DsqlRetryOptions
interface DsqlRetryOptions {
maxRetries?: number; // default: 3
baseDelayMs?: number; // default: 100
maxDelayMs?: number; // default: 5000
onRetry?: (error: Error, attempt: number) => void;
}Entity Queries
createQueryFactory(getDb)
Build a project-local query factory bound to one database client. Call it once per project; every query in the project then uses the returned defineQuery and definePreparedQuery. This is the current entity-query pattern -- it replaces the class plus @RequiresTable plus bind pattern for new code.
// src/db/defineQuery.ts
import { createQueryFactory } from "@j0nathan-ll0yd/database";
import { getDb } from "./client.js";
export const { defineQuery, definePreparedQuery } = createQueryFactory(getDb);Each query declares the tables and operations it touches through DefineQueryOptions, which mantle generate permissions extracts by AST analysis:
import { DatabaseOperation } from "@j0nathan-ll0yd/database";
import { eq } from "@j0nathan-ll0yd/database/orm";
import { defineQuery } from "../../db/defineQuery.js";
import { files } from "../../db/schema.js";
export const getFile = defineQuery(
{ tables: [{ table: files, operations: [DatabaseOperation.Select] }] },
async function getFile(db, fileId: string) {
const rows = await db.select().from(files).where(eq(files.fileId, fileId)).limit(1);
return rows[0] ?? null;
},
);Both factories wrap the query with withQueryMetrics, OCC retry, and -- when transaction: true is set -- withTransaction. definePreparedQuery additionally reuses a prepared statement across invocations in the same container.
See the Entity Queries Guide for import conventions and the full option list.
Entity Decorators
@RequiresTable
Stage 3 method decorator that declares table-level permissions on static entity query methods. Consumed by mantle generate permissions to auto-generate per-Lambda PostgreSQL roles.
import { RequiresTable, DatabaseOperation, withQueryMetrics } from "@j0nathan-ll0yd/database";
import { eq } from "@j0nathan-ll0yd/database/orm";
class UserQueries {
@RequiresTable([{ table: "users", operations: [DatabaseOperation.Select] }])
static async findById(db: DrizzleClient, id: string) {
return withQueryMetrics("Users.findById", async () => {
return db.select().from(users).where(eq(users.id, id));
});
}
}
export const findById = UserQueries.findById.bind(UserQueries);getTablePermissions(target)
Retrieve the TablePermission[] array attached by @RequiresTable decorators on a class. Used by the CLI permission extractor.
DatabaseOperation
enum DatabaseOperation {
Select = 'SELECT'
Insert = 'INSERT'
Update = 'UPDATE'
Delete = 'DELETE'
All = 'ALL'
}TablePermission
interface TablePermission {
table: string;
operations: DatabaseOperation[];
}Migrations
readMigrationFiles(folder)
Read and parse all migration SQL files from a Drizzle migrations folder. Returns MigrationFile[].
runMigrations(config)
Run all pending migrations against the configured database provider. Handles locking, tracking, and DSQL-specific statement sanitization.
import { runMigrations } from "@j0nathan-ll0yd/database";
const result = await runMigrations({
migrationsFolder: "./drizzle",
database: { provider: "aurora-dsql", endpoint: "...", region: "..." },
logger: console.log,
});
// result.applied, result.migrationsclassifyStatement(stmt)
Classify a SQL statement for DSQL compatibility. Returns { type, description, table? } where type is 'compatible', 'create_index', 'needs_recreation', or 'unsupported_strip'.
MigrateConfig
interface MigrateConfig {
migrationsFolder: string;
database: DatabaseConfig;
migrationsTable?: string; // default: '__drizzle_migrations'
migrationsSchema?: string; // default: 'drizzle'
useLock?: boolean; // default: true
logger?: (msg: string) => void;
}MigrateResult
interface MigrateResult {
applied: number;
migrations: Array<{ tag: string; hash: string; durationMs: number }>;
}MigrationFile
interface MigrationFile {
tag: string; // e.g. '0001_my_migration'
hash: string; // SHA-256 of raw SQL content
sql: string; // full raw SQL
statements: string[]; // split on Drizzle breakpoint marker
}markMigrationApplied(config)
Record a migration as applied without executing it. Use after applying SQL out of band -- for example when a migration was run manually to recover a stuck deploy -- so the tracking table and the database agree.
import { markMigrationApplied } from "@j0nathan-ll0yd/database";
const result = await markMigrationApplied({
migrationsFolder: "./drizzle",
database: { provider: "aurora-dsql", endpoint: "...", region: "..." },
tag: "0007_add_files_index",
});Backs mantle db mark-applied. Returns MarkAppliedResult; takes MarkAppliedConfig.
Permissions
applyPermissions(config)
Apply SQL permission files to the database using an admin connection. Handles env var substitution and idempotent error skipping.
import { applyPermissions } from "@j0nathan-ll0yd/database";
const result = await applyPermissions({
permissionsFolder: "./infra/permissions",
database: { provider: "aurora-dsql", endpoint: "...", region: "..." },
logger: console.log,
});
// result.applied, result.skipped, result.errorsPermissionsConfig
interface PermissionsConfig {
permissionsFolder: string;
database: DatabaseConfig;
logger?: (msg: string) => void;
}PermissionsResult
interface PermissionsResult {
applied: number;
skipped: number;
errors: string[];
}validatePermissions(config)
Compare the grants present in the database against the grants the permissions/ folder declares, and report the gaps. Read-only: it never issues a GRANT.
import { validatePermissions } from "@j0nathan-ll0yd/database";
const result = await validatePermissions({
permissionsFolder: "./infra/permissions",
database: { provider: "aurora-dsql", endpoint: "...", region: "..." },
});
// result.gaps: PermissionGap[]Takes ValidatePermissionsConfig, returns ValidatePermissionsResult whose gaps are PermissionGap entries. Backs mantle check permissions-fresh.
Types
DatabaseConfig
interface DatabaseConfig {
provider: DatabaseProvider;
connectionString?: string; // for Neon or Aurora Serverless v2
endpoint?: string; // Aurora DSQL cluster endpoint
region?: string; // AWS region for IAM signing
username?: string; // database username/role name
isAdmin?: boolean; // use admin token for DSQL
schema?: Record<string, unknown>; // Drizzle schema object
maxConnections?: number; // default: 1 (Lambda-optimized)
idleTimeout?: number; // seconds
connectTimeout?: number; // seconds
instrument?: boolean; // emit metrics during connect (default: true)
}DatabaseProvider
type DatabaseProvider = "aurora-dsql" | "aurora-serverless-v2" | "neon";TransactionClient
type TransactionClient<TSchema = Record<string, unknown>> = PgTransaction<...>Drizzle transaction client type for PostgreSQL providers. Generic over the schema type.
DSQL Translation Layer
Aurora DSQL is PostgreSQL-compatible but not PostgreSQL-complete. Everything that bridges the gap lives in packages/database/src/dsql/ and is re-exported from the package root. These modules never import @j0nathan-ll0yd/observability, which is what keeps the CLI migration path importable.
Statement classification and rewriting
| Function | Purpose |
|---|---|
classifyStatement(stmt) | Classify one SQL statement for DSQL compatibility. Returns StatementClassification. |
sanitizeForDsql(stmt) | Rewrite a PostgreSQL statement for DSQL, or return null when it must be stripped. |
adaptForStandardPg(sql, prefix?) | The inverse: rewrite DSQL SQL so it runs on standard PostgreSQL (used by mantle db clone). |
import { classifyStatement, sanitizeForDsql } from "@j0nathan-ll0yd/database";
const classification = classifyStatement("CREATE INDEX idx_files_user ON files (user_id)");
const rewritten = sanitizeForDsql("ALTER TABLE files ADD COLUMN size integer NOT NULL DEFAULT 0");Incompatibility registry
DSQL_INCOMPATIBILITIES is the queryable catalogue of known PostgreSQL-to-DSQL differences. Each entry is a DsqlIncompatibility carrying a DsqlCategory and the DsqlAction the framework takes.
import {
DSQL_INCOMPATIBILITIES,
DsqlAction,
DsqlCategory,
getByAction,
getByCategory,
getIncompatibility,
} from "@j0nathan-ll0yd/database";
getByCategory(DsqlCategory.DDL); // every DDL-related difference
getByAction(DsqlAction.Strip); // everything the sanitizer removes
getIncompatibility("alter-table-add-constraint"); // one entry by id, or undefinedDsqlCategory values: DDL, DML, DataType, Permission, Index, Transaction, SystemCatalog, Connection, Function. DsqlAction values: Compatible, Rewrite, Strip, Recreate, Error, Warn, AppLayer.
Table recreation
DSQL rejects several ALTER TABLE forms. For those, the migration runner rebuilds the table instead.
| Function | Purpose |
|---|---|
introspectTable(db, tableName) | Read the live TableSchema (columns, constraints, indexes). |
applyAlterToSchema(schema, stmt) | Apply an ALTER TABLE statement to a schema in memory. Pure. |
buildCreateTableSql(name, schema) | Render a CREATE TABLE statement from a TableSchema. Pure. |
getEstimatedRowCount(db, tableName) | Estimated row count, used to size and log the rebuild. |
generateRecreationPlan(current, desired, rows) | Produce the ordered RecreationPlan steps. Pure. |
executeRecreation(db, plan, logger?) | Run the plan one statement per transaction, as DSQL requires. |
Supporting types: TableSchema, TableColumn, TableConstraint, TableIndex, RecreationPlan.
Error codes
| Export | Meaning |
|---|---|
OCC_ERROR_CODES | SQLSTATE codes that mean an optimistic-concurrency conflict. |
IDEMPOTENT_DDL_CODES | Codes safe to swallow when re-running DDL (already exists, already dropped). |
IDEMPOTENT_TEST_CODES | The same idea for the test harness's setup and teardown SQL. |
isOccError(error) | True when the error is an OCC conflict and the operation should be retried. |
isIdempotentDdlError(err) | True when re-running the DDL is a no-op rather than a failure. |
Foreign-key enforcement
DSQL does not enforce foreign keys. Assert the parent row at the application layer instead.
import { assertRowExists, assertRowsExist } from "@j0nathan-ll0yd/database";
await assertRowExists(db, "users", "user_id", userId);
await assertRowsExist(db, "users", "user_id", userIds);Both throw ForeignKeyViolationError (HTTP 409) from @j0nathan-ll0yd/errors when a referenced row is absent.
Subpath: @j0nathan-ll0yd/database/orm
Re-exported Drizzle ORM operators and types for entity query building. Use this instead of importing drizzle-orm directly to avoid dual-instance type conflicts with linked packages.
Query Operators
import {
eq,
and,
or,
gt,
gte,
lt,
lte,
ne,
sql,
count,
asc,
desc,
inArray,
notInArray,
isNull,
isNotNull,
} from "@j0nathan-ll0yd/database/orm";Type Helpers
import type {
SelectModel,
InsertModel,
UpdateModel,
InferSelectModel,
InferInsertModel,
} from "@j0nathan-ll0yd/database/orm";
// Infer the select (read) model from a table
type UserRow = SelectModel<typeof usersTable>;
// Infer the insert model from a table
type NewUser = InsertModel<typeof usersTable>;
// Partial update type, excluding primary key and timestamps
type UserUpdate = UpdateModel<typeof usersTable, "id" | "createdAt">;