Custom stores ​
Implement MemoryStore when the official adapters do not fit the application's storage layer. The core contract is deliberately small: load ordered messages, append one completed batch, and clear one scoped conversation.
1. Implement the required methods ​
import { createMemoryScopeKey, type Message } from '@anvia/core'
import type {
MemoryAppendOptions,
MemoryScope,
MemoryStore,
} from '@anvia/core/memory'
const scopeKey = (scope: MemoryScope) => createMemoryScopeKey({ scope, metadataKeys: ['tenantId'] })
export class ProductMemoryStore implements MemoryStore {
async load({ scope }: { scope: MemoryScope }): Promise<Message[]> {
const rows = await db.memoryMessages.findMany({
where: { scopeKey: scopeKey(scope) },
orderBy: { position: 'asc' },
})
return rows.map((row) => parseStoredMessage(row.message))
}
async append(input: MemoryAppendOptions): Promise<void> {
if (input.messages.length === 0) return
await appendOrderedMessagesAtomically({
scopeKey: scopeKey(input.scope),
runId: input.runId,
turn: input.turn,
messages: input.messages,
})
}
async clear({ scope }: { scope: MemoryScope }): Promise<void> {
await deleteConversation({ scopeKey: scopeKey(scope) })
}
}load(), append(), and clear() are required. Preserve every normalized message field and return canonical history in model order. Compaction must not change what load() returns.
The core runtime passes the same context shape to every method:
type MemoryScope = {
sessionId: string
userId?: string
metadata?: JsonObject
}Define one deterministic storage key from the complete scope required by the product. Apply it identically to reads, writes, deletion, errors, inspection, and compaction.
2. Make appends safe ​
Each append receives runId, turn, and an ordered Message[] batch. Use a transaction, lock, compare-and-swap, or database constraint so concurrent writers cannot assign duplicate positions or interleave a batch.
Use runId and turn as idempotency information when the storage system can replay writes. Do not silently discard a persistence failure; reject so the agent run can follow its failure path.
3. Record failures separately ​
recordError() is optional:
import type {
MemoryErrorOptions,
} from '@anvia/core/memory'
async recordError(input: MemoryErrorOptions): Promise<void> {
await db.memoryErrors.create({
data: {
scopeKey: scopeKey(input.scope),
runId: input.runId,
error: serializeSafeError(input.error),
messages: input.messages,
},
})
}Error records are operational diagnostics. load() must not mix them into normal conversation context.
4. Add read-only inspection when needed ​
Expose an optional inspector with:
listConversations({ limit, userId? })returning conversation summaries; andgetConversation(ref)returning one summary plus ordered message records.
The ref value is an opaque store-owned identifier used to retrieve that exact conversation. This capability lets Studio and internal tooling inspect memory without teaching core about the database schema.
5. Add atomic compaction when needed ​
Automatic compaction requires an optional compaction capability:
import type { MemoryCompactionCapability } from '@anvia/core/memory'
const compaction: MemoryCompactionCapability = {
async snapshot({ scope }) {
return loadRevisionAndModelProjection({ scopeKey: scopeKey(scope) })
},
async replacePrefix(input) {
return advanceCompactionCheckpointIfRevisionMatches({
scopeKey: scopeKey(input.scope),
expectedRevision: input.revision,
messageCount: input.messageCount,
replacement: input.replacement,
runId: input.runId,
})
},
}snapshot() returns an opaque revision and the ordered model-context projection. With no checkpoint, that projection is the canonical history. With a checkpoint, it is the stored summary followed by canonical messages after the summarized boundary.
replacePrefix() must atomically compare the revision and advance a separate checkpoint so the requested projected prefix is represented by replacement. Its messageCount includes the prior summary when a repeated compaction covers that summary. Return { status: 'committed' } or { status: 'conflict' } without deleting or rewriting canonical message rows. Keep the checkpoint's summary, canonical boundary, and generation or equivalent conflict state together.
Appending messages must preserve the current checkpoint. Clearing a conversation must remove both canonical messages and its checkpoint. Validate persisted summaries and reject a checkpoint whose boundary no longer exists in canonical history.
Never overwrite a newer transcript after a conflict. Let the core runtime reload and retry according to the configured conflict limit.
Test empty histories, repeated compactions, append-after-compaction, clear-after-compaction, concurrent writers, failed transactions, full-scope deletion, error isolation, malformed checkpoint state, compaction conflicts, canonical replay, model projection, and message-order preservation before using a custom store in production.
6. Reuse the official scope-key helper ​
createMemoryScopeKey({ scope, includeUserId?, metadataKeys? }) is exported from @anvia/core and @anvia/core/memory. It JSON-encodes an ordered array containing sessionId, userId by default (null when absent), and each selected metadata value. Unselected metadata is ignored. Dotted paths such as organization.id traverse plain objects; missing values become null.
The following complete adapter factory demonstrates one policy at every storage boundary. The backend is application-owned: it must implement atomic ordered/idempotent appends, protected error storage, canonical-history preservation, and revision-checked compaction as described above.
import { createMemoryScopeKey, type Message } from '@anvia/core'
import type {
MemoryAppendOptions, MemoryErrorOptions, MemoryScope, MemoryStore,
MemoryInspector, MemoryCompactionReplacePrefixOptions,
MemoryCompactionReplacePrefixResult, MemoryCompactionSnapshot,
} from '@anvia/core/memory'
type Keyed<T> = Omit<T, 'scope'> & { key: string }
type Backend = {
load(key: string): Promise<Message[]>
append(input: Keyed<MemoryAppendOptions>): Promise<void>
clear(key: string): Promise<void> // clear canonical rows AND the compaction checkpoint
recordError(input: Keyed<MemoryErrorOptions>): Promise<void>
snapshot(key: string): Promise<MemoryCompactionSnapshot>
replacePrefix(input: Keyed<MemoryCompactionReplacePrefixOptions>): Promise<MemoryCompactionReplacePrefixResult>
// Return summaries with ref equal to the same stored key; filter in authorized tooling.
listConversations: MemoryInspector['listConversations']
getConversation: MemoryInspector['getConversation']
}
export const keyFor = (scope: MemoryScope) => createMemoryScopeKey({ scope, metadataKeys: ['tenantId'] })
export function createProductMemoryStore(backend: Backend): MemoryStore {
return {
load: ({ scope }) => backend.load(keyFor(scope)),
append: ({ scope, ...input }) => backend.append({ ...input, key: keyFor(scope) }),
clear: ({ scope }) => backend.clear(keyFor(scope)),
recordError: ({ scope, ...input }) => backend.recordError({ ...input, key: keyFor(scope) }),
inspector: {
listConversations: (options) => backend.listConversations(options),
getConversation: ({ ref }) => backend.getConversation({ ref }),
},
compaction: {
snapshot: ({ scope }) => backend.snapshot(keyFor(scope)),
replacePrefix: ({ scope, ...input }) => backend.replacePrefix({ ...input, key: keyFor(scope) }),
},
}
}
const demoKey = keyFor({ sessionId: 'session-1', userId: 'user-1', metadata: { tenantId: 'tenant-1' } })
console.log(demoKey) // ["session-1","user-1","tenant-1"]The inspector's ref addresses the already stored key; it is not a user-supplied scope or proof of access. For a backend with separate opaque references, map them to the stored key server-side.
includeUserId: false removes the user element and deliberately shares a session across users with otherwise equal selected values. metadataKeys order matters, including nested JSON value serialization order. Keep selected values stable scalars when possible. Changing the key policy (including order or user inclusion) addresses different storage; plan migration of history, checkpoints, errors, and references before changing it for existing conversations.
Key derivation does not validate that tenant metadata exists or authorize any caller. If tenancy is required, reject a missing/blank tenant before constructing the scope. Missing and explicit null metadata values map to the same null key element by design. Retain a custom MemoryScopeKeyResolver for an application-specific policy, and apply it consistently at all boundaries.