import type { MastraMessageContentV2 } from '@mastra/core/agent'; import type { MastraDBMessage, StorageThreadType } from '@mastra/core/memory'; import { MemoryStorage } from '@mastra/core/storage'; /** * Columns added to the OM table after its initial release. * Used in `alterTable({ ifNotExists })` so that databases created on older * versions get the new columns automatically. * * When you add a column to OBSERVATIONAL_MEMORY_SCHEMA in @mastra/core, * you MUST also add it here — the unit test `om-migration-columns.test.ts` * will fail otherwise. */ export declare const OM_MIGRATION_COLUMNS: string[]; import type { StorageResourceType, StorageListMessagesInput, StorageListMessagesByResourceIdInput, StorageListMessagesOutput, StorageListThreadsInput, StorageListThreadsOutput, CreateIndexOptions, StorageCloneThreadInput, StorageCopyThreadOutput, ObservationalMemoryRecord, ObservationalMemoryHistoryOptions, CreateObservationalMemoryInput, UpdateActiveObservationsInput, UpdateBufferedObservationsInput, SwapBufferedToActiveInput, SwapBufferedToActiveResult, UpdateBufferedReflectionInput, SwapBufferedReflectionToActiveInput, CreateReflectionGenerationInput, UpdateObservationalMemoryConfigInput, PruneOptions, PruneResult, RetentionTablesDescriptor, TableRetentionPolicy } from '@mastra/core/storage'; import type { PgDomainConfig } from '../../db/index.js'; export declare class MemoryPG extends MemoryStorage { #private; readonly supportsPartialThreadUpdate = true; readonly supportsObservationalMemory = true; /** * Retention-eligible tables. `threads`, `messages`, and `resources` all anchor * on the timezone-aware `createdAtZ` mirror column (kept in sync by triggers), * and are indexed for fast batched deletes. Cascade order is enforced in * `prune()` (children before threads), not here. Observational memory has no * timestamp anchor and is deliberately excluded. */ static readonly retentionTables: RetentionTablesDescriptor; /** Tables managed by this domain */ static readonly MANAGED_TABLES: readonly ["mastra_threads", "mastra_messages", "mastra_resources", "mastra_observational_memory"]; constructor(config: PgDomainConfig); init(): Promise; /** * Lazily ensures a btree index exists on each configured policy's retention * anchor column so age-based `prune()` deletes stay fast on large tables. * Called from the prune path (not init) so only deployments that configure * retention pay the index's write/disk overhead. Best-effort: failures are * logged and pruning proceeds (correct, just slower). * Created even with `skipDefaultIndexes` — retention is an explicit opt-in, * so its supporting index is not part of the default index set. */ private ensureRetentionIndexes; /** * Returns default index definitions for the memory domain tables. * @param schemaPrefix - Prefix for index names (e.g. "my_schema_" or "") */ static getDefaultIndexDefs(schemaPrefix: string): CreateIndexOptions[]; /** * Returns all DDL statements for this domain: tables (threads, messages, resources, OM), indexes. * Used by exportSchemas to produce a complete, reproducible schema export. */ static getExportDDL(schemaName?: string): string[]; /** * Returns default index definitions for this instance's schema. */ getDefaultIndexDefinitions(): CreateIndexOptions[]; /** * Creates default indexes for optimal query performance. */ createDefaultIndexes(): Promise; /** * Creates custom user-defined indexes for this domain's tables. */ createCustomIndexes(): Promise; dangerouslyClearAll(): Promise; /** * Deletes rows older than the configured `maxAge` per table, in bounded, * batched, cancellable chunks. Tables are pruned children-first (messages and * resources before threads) since PostgreSQL has no FK cascade in this schema. * Unset tables are kept forever. * * When a `messages` policy is set, semantic-recall embeddings for pruned * messages are also swept from same-schema `memory_messages*` vector tables * (best-effort, mirroring `deleteThread`). Embeddings held in an external * vector store are out of reach and must be pruned by the operator. */ prune(policies: Record, options?: PruneOptions): Promise; /** * Best-effort sweep of semantic-recall vector rows whose source message no * longer exists (e.g. it was just pruned), so recall doesn't keep returning * embeddings that resolve to nothing. Only same-schema default vector tables * (`memory_messages*`) are covered — the same set `deleteThread` cleans up. * Failures are logged, never thrown: vector cleanup must not fail the prune. */ private pruneOrphanedVectorRows; /** * Normalizes message row from database by applying createdAtZ fallback */ private normalizeMessageRow; getThreadById({ threadId, resourceId, }: { threadId: string; resourceId?: string; }): Promise; /** * Atomically reassign a thread and all of its messages to a different resource. * * Runs inside a single transaction and takes a `SELECT ... FOR UPDATE` row lock on the * thread, so overlapping transfers of the same thread serialize and can never interleave * the thread update with the message update. Either both the thread and every message move * to the new resource, or neither does — there is no split-ownership window. The thread's * `createdAt` is preserved. Callers are responsible for authorizing the reassignment. */ updateThreadResourceId({ threadId, resourceId, }: { threadId: string; resourceId: string; }): Promise; listThreads(args: StorageListThreadsInput): Promise; saveThread({ thread }: { thread: StorageThreadType; }): Promise; updateThread({ id, title, metadata, }: { id: string; title?: string; metadata?: Record; }): Promise; deleteThread({ threadId }: { threadId: string; }): Promise; /** * Fetches messages around target messages using cursor-based pagination. * * This replaces the previous ROW_NUMBER() approach which caused severe performance * issues on large tables (see GitHub issue #11150). The old approach required * scanning and sorting ALL messages in a thread to assign row numbers. * * The current approach uses two phases for optimal performance: * 1. Batch-fetch all target messages' metadata (thread_id, createdAt) in one query * 2. Build cursor subqueries using "createdAt" directly (not COALESCE) so that * the existing (thread_id, createdAt DESC) index can be used for index scans * instead of sequential scans. This fixes GitHub issue #11702 where semantic * recall latency scaled linearly with message count (~30s for 7.4k messages). */ private _sortMessages; /** * Fetches included messages by ID, discovering their thread automatically. * This handles cross-thread includes where the include item doesn't specify a threadId. * When a resourceId is given, both the target lookup and the surrounding window stay * inside that resource, so an include never leaks another resource's messages. */ private _getIncludedMessages; private parseRow; listMessagesById({ messageIds }: { messageIds: string[]; }): Promise<{ messages: MastraDBMessage[]; }>; listMessages(args: StorageListMessagesInput): Promise; listMessagesByResourceId(args: StorageListMessagesByResourceIdInput): Promise; saveMessages({ messages }: { messages: MastraDBMessage[]; }): Promise<{ messages: MastraDBMessage[]; }>; updateMessages({ messages, }: { messages: (Partial> & { id: string; content?: { metadata?: MastraMessageContentV2['metadata']; content?: MastraMessageContentV2['content']; }; })[]; }): Promise; deleteMessages(messageIds: string[]): Promise; getResourceById({ resourceId }: { resourceId: string; }): Promise; saveResource({ resource }: { resource: StorageResourceType; }): Promise; updateResource({ resourceId, workingMemory, metadata, }: { resourceId: string; workingMemory?: string; metadata?: Record; }): Promise; copyThread(args: StorageCloneThreadInput): Promise; private getOMKey; private parseOMRow; getObservationalMemory(threadId: string | null, resourceId: string): Promise; getObservationalMemoryHistory(threadId: string | null, resourceId: string, limit?: number, options?: ObservationalMemoryHistoryOptions): Promise; initializeObservationalMemory(input: CreateObservationalMemoryInput): Promise; insertObservationalMemoryRecord(record: ObservationalMemoryRecord): Promise; updateActiveObservations(input: UpdateActiveObservationsInput): Promise; createReflectionGeneration(input: CreateReflectionGenerationInput): Promise; setReflectingFlag(id: string, isReflecting: boolean): Promise; setObservingFlag(id: string, isObserving: boolean): Promise; setBufferingObservationFlag(id: string, isBuffering: boolean, lastBufferedAtTokens?: number): Promise; setBufferingReflectionFlag(id: string, isBuffering: boolean): Promise; clearObservationalMemory(threadId: string | null, resourceId: string): Promise; setPendingMessageTokens(id: string, tokenCount: number): Promise; updateObservationalMemoryConfig(input: UpdateObservationalMemoryConfigInput): Promise; updateBufferedObservations(input: UpdateBufferedObservationsInput): Promise; swapBufferedToActive(input: SwapBufferedToActiveInput): Promise; updateBufferedReflection(input: UpdateBufferedReflectionInput): Promise; swapBufferedReflectionToActive(input: SwapBufferedReflectionToActiveInput): Promise; } //# sourceMappingURL=index.d.ts.map