import type { ConnectionOptions } from 'node:tls'; import { MastraBase } from '@mastra/core/base'; import type { StorageColumn, TABLE_NAMES, CreateIndexOptions, IndexInfo, StorageIndexStats } from '@mastra/core/storage'; import { Pool } from 'pg'; import type { DbClient } from '../client.js'; export type { DbClient } from '../client.js'; /** * Configuration for standalone domain usage. * Accepts either: * 1. An existing database client (Pool or PoolAdapter) * 2. Config to create a new pool internally */ export type PgDomainConfig = PgDomainClientConfig | PgDomainPoolConfig | PgDomainRestConfig; /** * Pass an existing database client (DbClient) */ export interface PgDomainClientConfig { /** The writer database client */ client: DbClient; /** Optional reader database client. Falls back to `client`. */ readClient?: DbClient; /** Optional schema name (defaults to 'public') */ schemaName?: string; /** When true, default indexes will not be created during initialization */ skipDefaultIndexes?: boolean; /** Custom indexes to create for this domain's tables */ indexes?: CreateIndexOptions[]; } /** * Pass an existing pg.Pool */ export interface PgDomainPoolConfig { /** Pre-configured writer pg.Pool */ pool: Pool; /** Optional reader pg.Pool. Falls back to `pool`. */ readPool?: Pool; /** Optional schema name (defaults to 'public') */ schemaName?: string; /** When true, default indexes will not be created during initialization */ skipDefaultIndexes?: boolean; /** Custom indexes to create for this domain's tables */ indexes?: CreateIndexOptions[]; } /** * Pass config to create a new pg.Pool internally */ export type PgDomainRestConfig = { /** Optional schema name (defaults to 'public') */ schemaName?: string; /** When true, default indexes will not be created during initialization */ skipDefaultIndexes?: boolean; /** Custom indexes to create for this domain's tables */ indexes?: CreateIndexOptions[]; } & ({ host: string; port: number; database: string; user: string; password: string; ssl?: boolean | ConnectionOptions; } | { connectionString: string; ssl?: boolean | ConnectionOptions; }); /** * Resolves PgDomainConfig to a database client and schema. * Handles creating a new pool if config is provided. */ export declare function resolvePgConfig(config: PgDomainConfig): { client: DbClient; readClient: DbClient; schemaName?: string; skipDefaultIndexes?: boolean; indexes?: CreateIndexOptions[]; }; export declare function getSchemaName(schema?: string): string; export declare function getTableName({ indexName, schemaName }: { indexName: string; schemaName?: string; }): string; export declare function generateTableSQL({ tableName, schema, schemaName, compositePrimaryKey, includeAllConstraints, }: { tableName: TABLE_NAMES; schema: Record; schemaName?: string; compositePrimaryKey?: string[]; /** When true, includes all constraints in the SQL (for exports). When false, some constraints are added at runtime after data migration. */ includeAllConstraints?: boolean; }): string; /** * Generates a CREATE INDEX SQL statement from index options. * Used by exportSchemas to produce index DDL without a database connection. */ export declare function generateIndexSQL(options: CreateIndexOptions, schemaName?: string): string; /** * Generates the SQL for a timestamp trigger function and trigger on a table. * Returns the DDL string without executing it. */ export declare function generateTimestampTriggerSQL(tableName: string, schemaName?: string): string; /** * Internal config for PgDB - accepts already-resolved client */ export interface PgDBInternalConfig { client: DbClient; readClient?: DbClient; schemaName?: string; skipDefaultIndexes?: boolean; } export declare class PgDB extends MastraBase { client: DbClient; readClient: DbClient; schemaName?: string; skipDefaultIndexes?: boolean; /** Cache of actual table columns: tableName -> Set */ private tableColumnsCache; /** Cache of column Postgres data types: tableName -> columnName -> data_type */ private columnTypeCache; constructor(config: PgDBInternalConfig); /** * Catalog snapshot for the current init window, or `null` outside it. * * When non-null, the init-path methods below answer existence questions from * it instead of round-tripping to the server, and record the objects they * create so later callers in the same init see them. See * {@link SchemaSnapshot} for why it is scoped to init only. */ private get schemaSnapshot(); /** * Whether the snapshot proves `generateTableSQL` would be a no-op for this * table — i.e. the CREATE statement can be skipped. * * For most tables that is just "the table exists". `workflow_snapshot` is the * exception: its generated SQL also carries a DO block that back-fills the * `(workflow_name, run_id)` unique constraint and promotes it to the table's * replica identity, so a table created by an older version still needs the * statement to run. */ private snapshotShowsTableConverged; /** Column set for `tableName` in the snapshot, created empty if absent. */ private snapshotColumns; /** * Records an out-of-band `ALTER TABLE … RENAME TO` in the init snapshot. * * Init-time migrations that issue raw DDL on `this.client` (instead of going * through createTable/alterTable/createIndex, which maintain the snapshot * themselves) MUST report it through these `note*` methods. A snapshot that * still lists a renamed-away table makes a later createTable() in the same * init skip the rebuild the migration depends on — stranding data. No-op * outside the init window. * * Indexes riding along with a rename keep their names, so the snapshot's * index set stays accurate without changes here. */ noteTableRenamed(oldName: string, newName: string): void; /** * Records an out-of-band `DROP TABLE` in the init snapshot. See * {@link noteTableRenamed} for why raw-DDL migrations must call this. * * The dropped table's indexes vanish with it, but the snapshot's flat index * set cannot map names back to tables. Stale entries only make a later * createIndex() skip a recreate until the next init re-reads the catalog — * the same self-healing bound the rest of the snapshot design accepts. */ noteTableDropped(tableName: string): void; /** * Records an out-of-band `ALTER TABLE … ADD COLUMN` in the init snapshot. * See {@link noteTableRenamed} for why raw-DDL migrations must call this. */ noteColumnAdded(tableName: string, column: string): void; /** * Gets the set of column names that actually exist in the database table. * Results are cached; the cache is invalidated when alterTable() adds new columns. */ private getTableColumns; /** * Filters a record to only include columns that exist in the actual database table. * Unknown columns are silently dropped to ensure forward compatibility when newer * domain packages add fields that haven't been migrated yet. */ private filterRecordToKnownColumns; hasColumn(table: string, column: string): Promise; /** * Returns the Postgres data type of a column (e.g. `jsonb`, `json`, `text`), * or null when the table or column does not exist. * * Answered from the init snapshot when one is installed, so a warm `init()` * issues no catalog probe. Outside init, results are cached per instance and * the cache is invalidated alongside {@link tableColumnsCache} whenever DDL * changes a table. */ getColumnType(table: string, column: string): Promise; /** * Prepares values for insertion, handling JSONB columns by stringifying them */ private prepareValuesForInsert; /** * Adds timestamp Z columns to a record if timestamp columns exist */ private addTimestampZColumns; /** * Prepares a value for database operations */ private prepareValue; private setupSchema; protected getDefaultValue(type: StorageColumn['type']): string; private executeInsert; insert({ tableName, record }: { tableName: TABLE_NAMES; record: Record; }): Promise; clearTable({ tableName }: { tableName: TABLE_NAMES; }): Promise; createTable({ tableName, schema, compositePrimaryKey, }: { tableName: TABLE_NAMES; schema: Record; compositePrimaryKey?: string[]; }): Promise; private setupTimestampTriggers; /** * Migrates the spans table schema from OLD_SPAN_SCHEMA to current SPAN_SCHEMA. * This adds new columns that don't exist in old schema. */ private migrateSpansTable; /** * Deduplicates spans in the mastra_ai_spans table before adding the PRIMARY KEY constraint. * Keeps spans based on priority: completed (endedAt NOT NULL) > most recent updatedAt > most recent createdAt. * * Note: This prioritizes migration completion over perfect data preservation. * Old trace data may be lost, which is acceptable for this use case. */ private deduplicateSpans; /** * Checks for duplicate (traceId, spanId) combinations in the spans table. * Returns information about duplicates for logging/CLI purposes. */ private checkForDuplicateSpans; /** * Checks if the PRIMARY KEY constraint on (traceId, spanId) already exists on the spans table. * Used to skip deduplication when the constraint already exists (migration already complete). */ private spansPrimaryKeyExists; /** Live-catalog variant of {@link spansPrimaryKeyExists}, bypassing the snapshot. */ private spansPrimaryKeyExistsLive; /** * Adds the PRIMARY KEY constraint on (traceId, spanId) to the spans table. * Should be called AFTER deduplication to ensure no duplicate key violations. */ private addSpansPrimaryKey; /** * Manually run the spans migration to deduplicate and add the unique constraint. * This is intended to be called from the CLI when duplicates are detected. * * @returns Migration result with status and details */ migrateSpans(): Promise<{ success: boolean; alreadyMigrated: boolean; duplicatesRemoved: number; message: string; }>; /** * Check migration status for the spans table. * Returns information about whether migration is needed. */ checkSpansMigrationStatus(): Promise<{ needsMigration: boolean; hasDuplicates: boolean; duplicateCount: number; constraintExists: boolean; tableName: string; }>; /** * Alters table schema to add columns if they don't exist * @param tableName Name of the table * @param schema Schema of the table * @param ifNotExists Array of column names to add if they don't exist */ alterTable({ tableName, schema, ifNotExists, }: { tableName: TABLE_NAMES; schema: Record; ifNotExists: string[]; }): Promise; load({ tableName, keys }: { tableName: TABLE_NAMES; keys: Record; }): Promise; batchInsert({ tableName, records }: { tableName: TABLE_NAMES; records: Record[]; }): Promise; dropTable({ tableName }: { tableName: TABLE_NAMES; }): Promise; createIndex(options: CreateIndexOptions): Promise; /** * Runs a caller-built `CREATE INDEX IF NOT EXISTS` statement, unless the init * snapshot already proves `indexName` exists. * * `createIndex` covers the indexes described by {@link CreateIndexOptions}; * this is for the two init paths that hand-write their statement (a partial * or otherwise non-standard index) and would otherwise send a no-op DDL on * every warm init. */ createIndexFromStatement(indexName: string, sql: string): Promise; dropIndex(indexName: string): Promise; listIndexes(tableName?: string): Promise; describeIndex(indexName: string): Promise; update({ tableName, keys, data, }: { tableName: TABLE_NAMES; keys: Record; data: Record; }): Promise; batchUpdate({ tableName, updates, }: { tableName: TABLE_NAMES; updates: Array<{ keys: Record; data: Record; }>; }): Promise; batchDelete({ tableName, keys }: { tableName: TABLE_NAMES; keys: Record[]; }): Promise; /** * Delete all data from a table (alias for clearTable for consistency with other stores) */ deleteData({ tableName }: { tableName: TABLE_NAMES; }): Promise; /** * Deletes up to `limit` rows from `tableName` whose `column` value is strictly * older than `cutoff`, in a single bounded statement. Returns the number of * rows deleted so the caller's batch loop can decide whether the table is * drained. * * PostgreSQL has no `DELETE ... LIMIT`, so this targets a bounded set of * physical rows via the `ctid` system column (PG's row identity, analogous to * SQLite's `rowid`). `cutoff` is bound as a parameter — a `Date`/ISO-8601 * string compared against a `timestamptz` anchor column, or a `number` * compared against a `bigint` epoch-ms anchor column. */ pruneBatch({ tableName, column, cutoff, limit, }: { tableName: TABLE_NAMES; column: string; cutoff: Date | string | number; limit: number; }): Promise; /** * Deletes up to `limit` aged parent rows *and* their child rows together, in * a single transaction (used by whole-unit pruning such as experiments → * experiment_results). The aged parent IDs are selected first and both * deletes target that exact ID set, so a bound or abort between batches never * leaves a parent hollow (kept, but with its children gone) or children * orphaned. */ pruneUnitsBatch({ parentTable, parentKey, parentColumn, childTable, childForeignKey, cutoff, limit, }: { parentTable: TABLE_NAMES; parentKey: string; parentColumn: string; childTable: TABLE_NAMES; childForeignKey: string; cutoff: Date | string | number; limit: number; }): Promise<{ parents: number; children: number; }>; /** * Creates a btree index on `column` for `tableName` if it does not already * exist, so age-based prune deletes stay fast. Delegates to {@link createIndex} * (which is a no-op when the index is present). * * The name is lowercased and truncated to Postgres' 63-byte identifier limit * (schema-prefixed names can exceed it), mirroring {@link buildConstraintName}. */ ensureIndex({ indexName, tableName, column, }: { indexName: string; tableName: TABLE_NAMES; column: string; }): Promise; private executeUpdate; } //# sourceMappingURL=index.d.ts.map