diff --git a/README.md b/README.md index 3e8fd98..b588018 100644 --- a/README.md +++ b/README.md @@ -61,6 +61,49 @@ If `userId` is omitted, memory is stored in the shared default scope. That is us - `memory.learn: false`: disables learning for a single call. - `memory.recall: false`: disables recall for a single call. +## Data Migration + +Use `agent.memory.migrate()` when an app already has memory-like data and wants to seed the configured persistence layer. The SDK does not fetch or map source data; pass normalized memories or events from your own import code. + +```ts +await agent.memory.migrate({ + userId: "user_123", + data: [ + { + type: "preference", + content: "User prefers concise weekly reports.", + confidence: 0.9, + importance: 0.8 + } + ] +}) +``` + +Migration validates the input, resolves the target scope, writes through the configured memory store, links source events when provided, creates a synthetic source event for bare memory rows, and stores embeddings so imported memory can be recalled immediately. Use `mode: "skipExisting"` to preserve an existing memory with the same ID or canonical key. + +```ts +await agent.memory.migrate({ + userId: "user_123", + mode: "skipExisting", + events: [ + { + id: "evt_legacy_1", + role: "user", + content: "Legacy note: the customer portal was renamed to Atlas." + } + ], + memories: [ + { + id: "mem_legacy_1", + type: "fact", + content: "The customer portal was renamed to Atlas.", + canonicalKey: "fact:customer-portal-renamed", + sourceEventIds: ["evt_legacy_1"] + } + ] +}) +``` + ## Storage The public `agent-memory-sdk` package defaults to local JSON persistence: diff --git a/packages/agent-memory/README.md b/packages/agent-memory/README.md index c5503b4..b0b850b 100644 --- a/packages/agent-memory/README.md +++ b/packages/agent-memory/README.md @@ -39,3 +39,23 @@ console.log(result.text) First-party helpers include `openai()`, `anthropic()`, `gemini()`, and `xai()`. Use `openAICompatible()` for custom chat-completions endpoints. If `userId` is omitted, memory is stored in the shared default scope. Local memory defaults to `.memory/memory.json`; use `sqliteMemory()` for `.memory/memory.sqlite` or `postgresMemory()` for Postgres with automatic pgvector migrations. + +## Data Migration + +Seed existing app data by passing normalized memories or events to the configured memory store: + +```ts +await agent.memory.migrate({ + userId: "user_123", + data: [ + { + type: "preference", + content: "User prefers concise weekly reports.", + confidence: 0.9, + importance: 0.8 + } + ] +}) +``` + +The SDK only validates and stores mapped input. Your app owns where the data comes from and how it is mapped. Use `events` and `sourceEventIds` when you want imported memories linked to imported history, and `mode: "skipExisting"` to leave matching memories unchanged. diff --git a/packages/agent-memory/test/agent.test.mjs b/packages/agent-memory/test/agent.test.mjs index 6d66c70..b2de76e 100644 --- a/packages/agent-memory/test/agent.test.mjs +++ b/packages/agent-memory/test/agent.test.mjs @@ -101,6 +101,121 @@ test("generate without userId writes to the default scope", async () => { assert.match(memories[0]?.content ?? "", /prefers? short answers/i) }) +test("memory migration imports caller-mapped data and makes it recallable", async () => { + const seenRequests = [] + const agent = createAgent({ + model: customModel({ + id: "test-model", + generate: async (request) => { + seenRequests.push(request) + return { text: "ok" } + } + }), + memory: createMemoryStore() + }) + + const report = await agent.memory.migrate({ + userId: "migrated_user", + data: [ + { + type: "preference", + content: "User prefers dashboard summaries as bullet lists.", + confidence: 0.92, + importance: 0.81 + } + ] + }) + + assert.equal(report.memories.created, 1) + assert.equal(report.memories.updated, 0) + assert.equal(report.memories.skipped, 0) + assert.equal(report.memories.failed, 0) + assert.equal(report.events.imported, 1) + assert.deepEqual(report.failures, []) + + const memories = await agent.memory.list({ userId: "migrated_user" }) + assert.equal(memories.length, 1) + assert.equal(memories[0]?.scopeKey, "user:migrated_user") + assert.equal(memories[0]?.type, "preference") + assert.match(memories[0]?.canonicalKey ?? "", /^preference:/) + + await agent.generate({ + userId: "migrated_user", + messages: [ + { role: "user", content: "How should I format the dashboard summary?" } + ], + debug: true + }) + + const recalledRequest = seenRequests.at(-1) + assert.equal(recalledRequest.messages[0]?.role, "system") + assert.match(recalledRequest.messages[0]?.content ?? "", /dashboard summaries as bullet lists/i) +}) + +test("memory migration imports mapped events and can skip existing memories", async () => { + const agent = createAgent({ + model: customModel({ + id: "test-model", + generate: async () => ({ text: "ok" }) + }), + memory: createMemoryStore() + }) + + const first = await agent.memory.migrate({ + userId: "legacy_user", + threadId: "legacy_thread", + operationId: "legacy_import", + events: [ + { + id: "evt_legacy_1", + role: "user", + content: "Legacy note: the customer portal was renamed to Atlas.", + metadata: { source: "legacy-export" }, + createdAt: "2026-01-01T00:00:00.000Z" + } + ], + memories: [ + { + id: "mem_legacy_1", + type: "fact", + content: "The customer portal was renamed to Atlas.", + canonicalKey: "fact:customer-portal-renamed", + sourceEventIds: ["evt_legacy_1"], + createdAt: "2026-01-01T00:00:00.000Z", + updatedAt: "2026-01-01T00:00:00.000Z" + } + ] + }) + + assert.equal(first.events.imported, 1) + assert.equal(first.memories.created, 1) + assert.equal(first.memories.updated, 0) + + const second = await agent.memory.migrate({ + userId: "legacy_user", + mode: "skipExisting", + memories: [ + { + type: "fact", + content: "The customer portal has a different migrated description.", + canonicalKey: "fact:customer-portal-renamed" + } + ] + }) + + assert.equal(second.memories.created, 0) + assert.equal(second.memories.updated, 0) + assert.equal(second.memories.skipped, 1) + + const memories = await agent.memory.list({ userId: "legacy_user" }) + const exported = await agent.memory.export({ userId: "legacy_user" }) + + assert.equal(memories.length, 1) + assert.equal(memories[0]?.id, "mem_legacy_1") + assert.equal(memories[0]?.content, "The customer portal was renamed to Atlas.") + assert.equal(exported.events.some((event) => event.id === "evt_legacy_1"), true) +}) + test("stream returns text and commits memory after consumption", async () => { const agent = createAgent({ model: customModel({ diff --git a/packages/core/src/agent.ts b/packages/core/src/agent.ts index 7c2f7ad..477d3f4 100644 --- a/packages/core/src/agent.ts +++ b/packages/core/src/agent.ts @@ -1,6 +1,7 @@ import { randomUUID } from "node:crypto" import { createDefaultCompiler, canonicalKeyFor, validateMemoryPatch } from "./compiler.js" import { createMemoryStore, makeEmbedding, makeEvent } from "./memory-store.js" +import { migrateMemory } from "./migration.js" import { retrieveMemoryContext } from "./retrieval.js" import { resolveScope, targetScopeKeys } from "./scopes.js" import type { @@ -10,6 +11,7 @@ import type { CompilerProvider, GenerateInput, GenerateResult, + MemoryMigrationInput, MemoryPatch, MemoryRecord, MemoryScopeInput, @@ -221,6 +223,11 @@ export function createAgent(config: AgentConfig): Agent { }) }, + async migrate(input: MemoryMigrationInput) { + ensureStore(store) + return migrateMemory({ store, migration: input }) + }, + async delete(memoryId: string) { ensureStore(store) return store.deleteMemory(memoryId) diff --git a/packages/core/src/index.ts b/packages/core/src/index.ts index 96dd677..97f48a8 100644 --- a/packages/core/src/index.ts +++ b/packages/core/src/index.ts @@ -1,6 +1,7 @@ export { createAgent } from "./agent.js" export { createDefaultCompiler, canonicalKeyFor, validateMemoryPatch } from "./compiler.js" export { createMemoryStore, makeEmbedding, makeEvent, tokenize } from "./memory-store.js" +export { migrateMemory } from "./migration.js" export { customModel } from "./providers.js" export { retrieveMemoryContext } from "./retrieval.js" export { resolveScope, targetScopeKeys } from "./scopes.js" diff --git a/packages/core/src/migration.ts b/packages/core/src/migration.ts new file mode 100644 index 0000000..2ac0378 --- /dev/null +++ b/packages/core/src/migration.ts @@ -0,0 +1,275 @@ +import { randomUUID } from "node:crypto" +import { canonicalKeyFor } from "./compiler.js" +import { makeEmbedding } from "./memory-store.js" +import type { + AgentRole, + MemoryEvent, + MemoryMigrationEvent, + MemoryMigrationInput, + MemoryMigrationItem, + MemoryMigrationMode, + MemoryMigrationResult, + MemoryRecord, + MemoryStore +} from "./types.js" + +const EXISTING_LOOKUP_LIMIT = 10000 +const DEFAULT_CONFIDENCE = 0.85 +const DEFAULT_IMPORTANCE = 0.6 +const allowedModes = new Set(["upsert", "skipExisting"]) +const allowedRoles = new Set(["system", "user", "assistant", "tool"]) +const allowedSensitivity = new Set(["normal", "sensitive", "secret"]) +const allowedStatus = new Set(["active", "superseded", "deleted"]) + +export async function migrateMemory(input: { + store: MemoryStore + migration: MemoryMigrationInput +}): Promise { + const mode = normalizeMode(input.migration.mode) + const result: MemoryMigrationResult = { + memories: { + created: 0, + updated: 0, + skipped: 0, + failed: 0 + }, + events: { + imported: 0, + skipped: 0, + failed: 0 + }, + failures: [] + } + + const events = input.migration.events ?? [] + for (let index = 0; index < events.length; index += 1) { + try { + const event = normalizeEvent(input.migration, events[index]) + if (mode === "skipExisting" && await hasExistingEvent(input.store, event)) { + result.events.skipped += 1 + continue + } + + await input.store.writeEvent(event) + result.events.imported += 1 + } catch (error) { + result.events.failed += 1 + result.failures.push({ + target: "event", + index, + reason: errorReason(error) + }) + } + } + + const memories = [...(input.migration.data ?? []), ...(input.migration.memories ?? [])] + for (let index = 0; index < memories.length; index += 1) { + try { + const memory = normalizeMemory(input.migration, memories[index]) + const existing = await findExistingMemory(input.store, memory) + if (mode === "skipExisting" && existing) { + result.memories.skipped += 1 + continue + } + + if (memory.sourceEventIds.length === 0) { + const sourceEvent = syntheticSourceEvent(input.migration, memory) + await input.store.writeEvent(sourceEvent) + result.events.imported += 1 + memory.sourceEventIds = [sourceEvent.id] + } + + const stored = await input.store.upsertMemory(memory) + for (const eventId of memory.sourceEventIds) { + await input.store.linkSource({ + memoryId: stored.id, + eventId, + reason: "Imported memory source." + }) + } + + await input.store.upsertEmbedding({ + ownerType: "memory", + ownerId: stored.id, + model: "local-hash-v1", + dimensions: 32, + vector: makeEmbedding(stored.content), + createdAt: stored.updatedAt + }) + + if (existing) { + result.memories.updated += 1 + } else { + result.memories.created += 1 + } + } catch (error) { + result.memories.failed += 1 + result.failures.push({ + target: "memory", + index, + reason: errorReason(error) + }) + } + } + + return result +} + +function normalizeMode(mode: MemoryMigrationMode | undefined): MemoryMigrationMode { + if (!mode) return "upsert" + if (!allowedModes.has(mode)) { + throw new Error(`Unsupported memory migration mode: ${mode}`) + } + return mode +} + +function normalizeEvent(scope: MemoryMigrationInput, event: MemoryMigrationEvent | undefined): MemoryEvent { + if (!event) throw new Error("Migration event is missing.") + + const role = event.role ?? "user" + if (!allowedRoles.has(role)) { + throw new Error(`Migration event role is unsupported: ${role}`) + } + + return { + id: optionalString(event.id) ?? `evt_${randomUUID()}`, + scopeKey: migrationScopeKey(scope, event.scopeKey), + threadId: optionalString(event.threadId) ?? optionalString(scope.threadId), + operationId: optionalString(event.operationId) ?? optionalString(scope.operationId), + role, + content: requiredString(event.content, "Migration event content"), + metadata: normalizeMetadata(event.metadata), + createdAt: isoDate(event.createdAt, "Migration event createdAt") + } +} + +function normalizeMemory(scope: MemoryMigrationInput, item: MemoryMigrationItem | undefined): MemoryRecord { + if (!item) throw new Error("Migration memory is missing.") + + const type = optionalString(item.type) ?? "fact" + const content = requiredString(item.content, "Migration memory content") + const sensitivity = item.sensitivity ?? "normal" + const status = item.status ?? "active" + + if (!allowedSensitivity.has(sensitivity)) { + throw new Error(`Migration memory sensitivity is unsupported: ${sensitivity}`) + } + + if (!allowedStatus.has(status)) { + throw new Error(`Migration memory status is unsupported: ${status}`) + } + + const createdAt = isoDate(item.createdAt, "Migration memory createdAt") + + return { + id: optionalString(item.id) ?? `mem_${randomUUID()}`, + scopeKey: migrationScopeKey(scope, item.scopeKey), + type, + content, + canonicalKey: optionalString(item.canonicalKey) ?? canonicalKeyFor(type, content), + confidence: score(item.confidence, DEFAULT_CONFIDENCE, "Migration memory confidence"), + importance: score(item.importance, DEFAULT_IMPORTANCE, "Migration memory importance"), + sensitivity, + status, + sourceEventIds: uniqueStrings(item.sourceEventIds), + createdAt, + updatedAt: isoDate(item.updatedAt, "Migration memory updatedAt", createdAt) + } +} + +function syntheticSourceEvent(scope: MemoryMigrationInput, memory: MemoryRecord): MemoryEvent { + return { + id: `evt_${randomUUID()}`, + scopeKey: memory.scopeKey, + threadId: optionalString(scope.threadId), + operationId: optionalString(scope.operationId), + role: "user", + content: memory.content, + metadata: { + kind: "migration", + memoryId: memory.id, + memoryType: memory.type + }, + createdAt: memory.createdAt + } +} + +async function hasExistingEvent(store: MemoryStore, event: MemoryEvent): Promise { + const existing = await store.listRecentEvents({ + scopeKeys: [event.scopeKey], + limit: EXISTING_LOOKUP_LIMIT + }) + return existing.some((item) => item.id === event.id) +} + +async function findExistingMemory(store: MemoryStore, memory: MemoryRecord): Promise { + const existing = await store.listMemories({ + scopeKeys: [memory.scopeKey], + limit: EXISTING_LOOKUP_LIMIT + }) + return existing.find((item) => { + return item.id === memory.id || + (Boolean(item.canonicalKey) && item.canonicalKey === memory.canonicalKey) + }) +} + +function migrationScopeKey(scope: MemoryMigrationInput, explicitScopeKey?: string): string { + const explicit = optionalString(explicitScopeKey) + if (explicit) return explicit + + const userId = optionalString(scope.userId) + if (userId) return `user:${userId}` + + const orgId = optionalString(scope.orgId) + if (orgId) return `org:${orgId}` + + const threadId = optionalString(scope.threadId) + if (threadId) return `thread:${threadId}` + + return "default" +} + +function normalizeMetadata(metadata: Record | undefined): Record { + if (metadata === undefined) return {} + if (!metadata || typeof metadata !== "object" || Array.isArray(metadata)) { + throw new Error("Migration event metadata must be an object.") + } + return { ...metadata } +} + +function requiredString(value: string, label: string): string { + const normalized = optionalString(value) + if (!normalized) throw new Error(`${label} must be a non-empty string.`) + return normalized +} + +function optionalString(value: string | undefined): string | undefined { + if (typeof value !== "string") return undefined + const normalized = value.trim() + return normalized.length > 0 ? normalized : undefined +} + +function uniqueStrings(values: string[] | undefined): string[] { + return [...new Set((values ?? []).map(optionalString).filter((value): value is string => Boolean(value)))] +} + +function score(value: number | undefined, fallback: number, label: string): number { + if (value === undefined) return fallback + if (!Number.isFinite(value) || value < 0 || value > 1) { + throw new Error(`${label} must be a number between 0 and 1.`) + } + return value +} + +function isoDate(value: string | undefined, label: string, fallback = new Date().toISOString()): string { + if (value === undefined) return fallback + const timestamp = Date.parse(value) + if (!Number.isFinite(timestamp)) { + throw new Error(`${label} must be a valid date string.`) + } + return new Date(timestamp).toISOString() +} + +function errorReason(error: unknown): string { + return error instanceof Error ? error.message : String(error) +} diff --git a/packages/core/src/types.ts b/packages/core/src/types.ts index 9a95aca..ccfcc6e 100644 --- a/packages/core/src/types.ts +++ b/packages/core/src/types.ts @@ -114,6 +114,7 @@ export type Agent = { memory: { list(query?: MemoryListQuery): Promise search(query?: MemorySearchQuery): Promise + migrate(input: MemoryMigrationInput): Promise delete(memoryId: string): Promise forget(query?: MemoryScopeInput): Promise export(query?: MemoryScopeInput): Promise @@ -198,6 +199,62 @@ export type MemoryExport = { events: MemoryEvent[] } +export type MemoryMigrationMode = "upsert" | "skipExisting" + +export type MemoryMigrationInput = MemoryScopeInput & { + data?: MemoryMigrationItem[] + memories?: MemoryMigrationItem[] + events?: MemoryMigrationEvent[] + mode?: MemoryMigrationMode +} + +export type MemoryMigrationItem = { + id?: string + scopeKey?: string + type?: MemoryRecordType + content: string + canonicalKey?: string + confidence?: number + importance?: number + sensitivity?: MemoryRecord["sensitivity"] + status?: MemoryRecord["status"] + sourceEventIds?: string[] + createdAt?: string + updatedAt?: string +} + +export type MemoryMigrationEvent = { + id?: string + scopeKey?: string + threadId?: string + operationId?: string + role?: AgentRole + content: string + metadata?: Record + createdAt?: string +} + +export type MemoryMigrationFailure = { + target: "memory" | "event" + index: number + reason: string +} + +export type MemoryMigrationResult = { + memories: { + created: number + updated: number + skipped: number + failed: number + } + events: { + imported: number + skipped: number + failed: number + } + failures: MemoryMigrationFailure[] +} + export type MemoryStore = { kind?: string writeEvent(event: MemoryEvent): Promise