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
9 changes: 9 additions & 0 deletions .changeset/postgres-optional-schema-creation.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,9 @@
---
"@langchain/langgraph-checkpoint-postgres": minor
---

Add a `schemaSetup` option to `PostgresSaver` and `PostgresStore` for use alongside a custom `schema`. When `"create"` (the default), `setup()` runs `CREATE SCHEMA IF NOT EXISTS` as before. When `"verify"`, `setup()` instead checks the schema already exists and throws if it does not, without attempting to create it — useful for least-privilege database roles that are not permitted to create schemas. Table migrations still run either way.

Add an `ensureTables` option to `PostgresSaver` (matching `PostgresStore`). When `true` (the default), the first database operation runs `setup()` automatically, so calling `setup()` explicitly is now optional. When `false`, no auto-setup occurs and `setup()` must be called explicitly before use. This is backward-compatible: explicit `setup()` remains idempotent.

Fix a first-operation concurrency race in both `PostgresSaver` and `PostgresStore`: on a fresh instance, two operations starting concurrently could each trigger a full migration run, risking duplicate migration execution. `setup()` is now single-flighted so concurrent callers share one run, and a failed run is not cached (a corrected condition can retry).
93 changes: 86 additions & 7 deletions libs/checkpoint-postgres/src/index.ts
Original file line number Diff line number Diff line change
Expand Up @@ -20,15 +20,40 @@ import {
type SQL_TYPES,
getSQLStatements,
getTablesWithSchema,
assertSchemaExists,
} from "./sql.js";

/** @inline */
interface PostgresSaverOptions {
schema: string;
/**
* What `setup()` does about the schema.
*
* When `"create"` (default), `setup()` runs `CREATE SCHEMA IF NOT EXISTS`.
* When `"verify"`, `setup()` instead checks the schema already exists and
* throws if it does not — useful for least-privilege roles that are not
* permitted to create schemas. Table migrations still run either way.
*
* @default "create"
*/
schemaSetup: "create" | "verify";
/**
* Whether the first database operation should run `setup()` automatically.
*
* When `true` (default), the checkpointer lazily runs `setup()` the first
* time it touches the database, so calling `setup()` explicitly is optional.
* When `false`, no auto-setup occurs and `setup()` must be called explicitly
* before use — useful when migrations are run out-of-band.
*
* @default true
*/
ensureTables: boolean;
}

const _defaultOptions: PostgresSaverOptions = {
schema: "public",
schemaSetup: "create",
ensureTables: true,
};

const _ensureCompleteOptions = (
Expand All @@ -37,6 +62,8 @@ const _ensureCompleteOptions = (
return {
...options,
schema: options?.schema ?? _defaultOptions.schema,
schemaSetup: options?.schemaSetup ?? _defaultOptions.schemaSetup,
ensureTables: options?.ensureTables ?? _defaultOptions.ensureTables,
};
};

Expand All @@ -57,11 +84,15 @@ const { Pool } = pg;
* "postgresql://user:password@localhost:5432/db",
* // optional configuration object
* {
* schema: "custom_schema" // defaults to "public"
* schema: "custom_schema", // defaults to "public"
* schemaSetup: "verify", // defaults to "create"; "verify" checks the schema exists instead of creating it
* ensureTables: false // defaults to true; when false, you must call setup() explicitly before use
* }
* );
*
* // NOTE: you need to call .setup() the first time you're using your checkpointer
* // With the default `ensureTables: true`, setup runs automatically on first
* // use. Call setup() explicitly only when `ensureTables` is false, or to
* // provision tables ahead of time:
* await checkpointer.setup();
*
* const graph = createReactAgent({
Expand Down Expand Up @@ -90,6 +121,8 @@ export class PostgresSaver extends BaseCheckpointSaver {

protected isSetup: boolean;

private setupPromise?: Promise<void>;

constructor(
pool: pg.Pool,
serde?: SerializerProtocol,
Expand All @@ -114,6 +147,7 @@ export class PostgresSaver extends BaseCheckpointSaver {
* const checkpointer = PostgresSaver.fromConnString(connString, {
* schema: "custom_schema" // defaults to "public"
* });
* // setup() is optional with the default `ensureTables: true`.
* await checkpointer.setup();
*/
static fromConnString(
Expand All @@ -128,16 +162,54 @@ export class PostgresSaver extends BaseCheckpointSaver {
* Set up the checkpoint database asynchronously.
*
* This method creates the necessary tables in the Postgres database if they don't
* already exist and runs database migrations. It MUST be called directly by the user
* the first time checkpointer is used.
* already exist and runs database migrations.
*
* Calling this explicitly is optional: by default (`ensureTables: true`) the first
* database operation runs `setup()` automatically. Call it explicitly when
* `ensureTables` is `false`, or to provision tables ahead of first use. It is safe
* to call concurrently — the underlying migration run is single-flighted, so racing
* operations share one run rather than each replaying migrations, and a failed run
* is not cached (the next call retries).
*
* By default the target schema is created via `CREATE SCHEMA IF NOT EXISTS`. If the
* `schemaSetup` option was set to `"verify"` (e.g. for least-privilege roles that may
* not create schemas), this method instead verifies the schema already exists and
* throws if it does not. Table migrations run either way.
*/
async setup(): Promise<void> {
if (this.isSetup) return;

this.setupPromise ??= this.runSetupOnce();
try {
await this.setupPromise;
} catch (error) {
this.setupPromise = undefined;
throw error;
}
}

/**
* Run lazy auto-setup on the first database operation when `ensureTables`
* is enabled. A no-op once setup has completed or when `ensureTables` is
* `false` (in which case the user must call `setup()` explicitly).
*/
private async ensureSetup(): Promise<void> {
if (this.options.ensureTables && !this.isSetup) {
await this.setup();
}
}

private async runSetupOnce(): Promise<void> {
const client = await this.pool.connect();
const SCHEMA_TABLES = getTablesWithSchema(this.options.schema);
try {
await client.query(
`CREATE SCHEMA IF NOT EXISTS "${this.options.schema}"`
);
if (this.options.schemaSetup === "create") {
await client.query(
`CREATE SCHEMA IF NOT EXISTS "${this.options.schema}"`
);
} else {
await assertSchemaExists(client, this.options.schema);
}
let version = -1;
const MIGRATIONS = getMigrations(this.options.schema);

Expand Down Expand Up @@ -170,6 +242,8 @@ export class PostgresSaver extends BaseCheckpointSaver {
[v]
);
}

this.isSetup = true;
} finally {
client.release();
}
Expand Down Expand Up @@ -360,6 +434,7 @@ export class PostgresSaver extends BaseCheckpointSaver {
* @returns The retrieved checkpoint tuple, or undefined.
*/
async getTuple(config: RunnableConfig): Promise<CheckpointTuple | undefined> {
await this.ensureSetup();
const {
thread_id,
checkpoint_ns = "",
Expand Down Expand Up @@ -443,6 +518,7 @@ export class PostgresSaver extends BaseCheckpointSaver {
config: RunnableConfig,
options?: CheckpointListOptions
): AsyncGenerator<CheckpointTuple> {
await this.ensureSetup();
const { filter, before, limit } = options ?? {};
const [where, args] = this._searchWhere(config, filter, before);
let query = `${this.SQL_STATEMENTS.SELECT_SQL}${where} ORDER BY checkpoint_id DESC`;
Expand Down Expand Up @@ -561,6 +637,7 @@ export class PostgresSaver extends BaseCheckpointSaver {
metadata: CheckpointMetadata,
newVersions: ChannelVersions
): Promise<RunnableConfig> {
await this.ensureSetup();
if (config.configurable === undefined) {
throw new Error(`Missing "configurable" field in "config" param`);
}
Expand Down Expand Up @@ -633,6 +710,7 @@ export class PostgresSaver extends BaseCheckpointSaver {
writes: PendingWrite[],
taskId: string
): Promise<void> {
await this.ensureSetup();
const query = writes.every((w) => w[0] in WRITES_IDX_MAP)
? this.SQL_STATEMENTS.UPSERT_CHECKPOINT_WRITES_SQL
: this.SQL_STATEMENTS.INSERT_CHECKPOINT_WRITES_SQL;
Expand Down Expand Up @@ -664,6 +742,7 @@ export class PostgresSaver extends BaseCheckpointSaver {
}

async deleteThread(threadId: string): Promise<void> {
await this.ensureSetup();
const client = await this.pool.connect();
try {
await client.query("BEGIN");
Expand Down
30 changes: 29 additions & 1 deletion libs/checkpoint-postgres/src/sql.ts
Original file line number Diff line number Diff line change
@@ -1,3 +1,4 @@
import type pg from "pg";
import { type Checkpoint, TASKS } from "@langchain/langgraph-checkpoint";

export interface SQL_STATEMENTS {
Expand Down Expand Up @@ -130,8 +131,35 @@ export const getSQLStatements = (schema: string): SQL_STATEMENTS => {
export const tableExistsSQL = (schema: string, table: string) => {
const tableWithoutSchema = table.split(".")[1];
return `SELECT EXISTS (
SELECT FROM information_schema.tables
SELECT FROM information_schema.tables
WHERE table_schema = '${schema}'
AND table_name = '${tableWithoutSchema}'
);`;
};

export const schemaExistsSQL = `SELECT EXISTS (
SELECT FROM information_schema.schemata
WHERE schema_name = $1
);`;

/**
* Verify that the given schema already exists, throwing a clear error if not.
*
* Used by `setup()` when `schemaSetup` is `"verify"`: the schema must be
* provisioned out-of-band (e.g. by a DBA) because the connecting role is not
* permitted to create schemas. Shared by both the checkpointer and the store
* so the check and its error message stay in one place.
*/
export const assertSchemaExists = async (
client: pg.PoolClient,
schema: string
): Promise<void> => {
const result = await client.query(schemaExistsSQL, [schema]);
if (!result.rows[0]?.exists) {
throw new Error(
`Schema "${schema}" does not exist (or is not visible to the current role). ` +
`It was expected to already exist because "schemaSetup" is "verify". ` +
`Create the schema (or grant access to it) out-of-band, or set "schemaSetup" to "create".`
);
}
};
75 changes: 66 additions & 9 deletions libs/checkpoint-postgres/src/store/index.ts
Original file line number Diff line number Diff line change
Expand Up @@ -28,11 +28,39 @@ import {
StoreMigrationConfig,
} from "./store-migrations.js";
import { getStoreTablesWithSchema } from "./sql.js";
import { assertSchemaExists } from "../sql.js";

export type * from "./modules/types.js";

const { Pool } = pg;

/**
* `PostgresStoreConfig` with every default-able field resolved to a concrete
* value. `connectionOptions` is required from the caller, so it has no default
* and stays mandatory.
*/
type ResolvedPostgresStoreConfig = PostgresStoreConfig & {
schema: string;
schemaSetup: "create" | "verify";
ensureTables: boolean;
textSearchLanguage: string;
};

/**
* Resolve a user-supplied config into a complete config by filling in defaults
* for any omitted default-able field. Centralizes defaults in one place,
* mirroring `PostgresSaver`'s `_ensureCompleteOptions`.
*/
const _ensureCompleteConfig = (
config: PostgresStoreConfig
): ResolvedPostgresStoreConfig => ({
...config,
schema: config.schema ?? "public",
schemaSetup: config.schemaSetup ?? "create",
ensureTables: config.ensureTables ?? true,
textSearchLanguage: config.textSearchLanguage ?? "english",
});

/**
* PostgreSQL implementation of the BaseStore interface.
* This is now a lightweight orchestrator that delegates to specialized modules.
Expand All @@ -50,34 +78,41 @@ export class PostgresStore extends BaseStore {

private isSetup: boolean = false;

private setupPromise?: Promise<void>;

private isClosed: boolean = false;

private ensureTables: boolean;

private schemaSetup: "create" | "verify";

constructor(config: PostgresStoreConfig) {
super();

const resolved = _ensureCompleteConfig(config);

// Create connection pool
const pool =
typeof config.connectionOptions === "string"
? new Pool({ connectionString: config.connectionOptions })
: new Pool(config.connectionOptions);
typeof resolved.connectionOptions === "string"
? new Pool({ connectionString: resolved.connectionOptions })
: new Pool(resolved.connectionOptions);

// Initialize core and modules
this.core = new DatabaseCore(
pool,
config.schema || "public",
config.ttl,
config.index,
config.textSearchLanguage
resolved.schema,
resolved.ttl,
resolved.index,
resolved.textSearchLanguage
);

this.vectorOps = new VectorOperations(this.core);
this.crudOps = new CrudOperations(this.core, this.vectorOps);
this.searchOps = new SearchOperations(this.core, this.vectorOps);
this.ttlManager = new TTLManager(this.core);

this.ensureTables = config.ensureTables ?? true;
this.ensureTables = resolved.ensureTables;
this.schemaSetup = resolved.schemaSetup;
}

/**
Expand Down Expand Up @@ -196,10 +231,26 @@ export class PostgresStore extends BaseStore {

/**
* Initialize the store by running migrations to create necessary tables and indexes.
*
* Safe to call concurrently: the underlying migration run is single-flighted,
* so multiple operations racing on a fresh store share one `setup()` rather
* than each replaying migrations. A failed run is not cached — the next call
* retries, so a corrected condition (e.g. a schema created out-of-band after
* a `schemaSetup: "verify"` failure) can succeed.
*/
async setup(): Promise<void> {
if (this.isSetup) return;

this.setupPromise ??= this.runSetupOnce();
try {
await this.setupPromise;
} catch (error) {
this.setupPromise = undefined;
throw error;
}
}

private async runSetupOnce(): Promise<void> {
await this.runStoreMigrations();
this.isSetup = true;

Expand All @@ -217,7 +268,13 @@ export class PostgresStore extends BaseStore {
const STORE_TABLES = getStoreTablesWithSchema(this.core.schema);

try {
await client.query(`CREATE SCHEMA IF NOT EXISTS "${this.core.schema}"`);
if (this.schemaSetup === "create") {
await client.query(`CREATE SCHEMA IF NOT EXISTS "${this.core.schema}"`);
} else {
// The option is read from the constructor (not setup()), so the store's
// lazy auto-setup paths honor it too.
await assertSchemaExists(client, this.core.schema);
}

let version = -1;

Expand Down
4 changes: 2 additions & 2 deletions libs/checkpoint-postgres/src/store/modules/database-core.ts
Original file line number Diff line number Diff line change
Expand Up @@ -21,13 +21,13 @@ export class DatabaseCore {
schema: string,
ttlConfig?: TTLConfig,
indexConfig?: IndexConfig,
textSearchLanguage?: string
textSearchLanguage: string = "english"
) {
this.pool = pool;
this.schema = schema;
this.ttlConfig = ttlConfig;
this.indexConfig = indexConfig;
this.textSearchLanguage = textSearchLanguage || "english";
this.textSearchLanguage = textSearchLanguage;
}

async withClient<T>(
Expand Down
Loading
Loading