diff --git a/.changeset/soft-taxis-dress.md b/.changeset/soft-taxis-dress.md new file mode 100644 index 000000000..aac0b9a97 --- /dev/null +++ b/.changeset/soft-taxis-dress.md @@ -0,0 +1,72 @@ +--- +"@voltagent/serverless-hono": patch +"@voltagent/server-elysia": patch +"@voltagent/server-core": patch +"@voltagent/server-hono": patch +--- + +feat: add memory HTTP endpoints for conversations, messages, working memory, and search across server-core, Hono, Elysia, and serverless runtimes. + +### Endpoints + +- `GET /api/memory/conversations` +- `POST /api/memory/conversations` +- `GET /api/memory/conversations/:conversationId` +- `PATCH /api/memory/conversations/:conversationId` +- `DELETE /api/memory/conversations/:conversationId` +- `POST /api/memory/conversations/:conversationId/clone` +- `GET /api/memory/conversations/:conversationId/messages` +- `GET /api/memory/conversations/:conversationId/working-memory` +- `POST /api/memory/conversations/:conversationId/working-memory` +- `POST /api/memory/save-messages` +- `POST /api/memory/messages/delete` +- `GET /api/memory/search` + +Note: include `agentId` (query/body) when multiple agents are registered or no global memory is configured. + +### Examples + +Create a conversation: + +```bash +curl -X POST http://localhost:3141/api/memory/conversations \ + -H "Content-Type: application/json" \ + -d '{ + "userId": "user-123", + "resourceId": "assistant", + "title": "Support Chat", + "metadata": { "channel": "web" } + }' +``` + +Save messages into the conversation: + +```bash +curl -X POST http://localhost:3141/api/memory/save-messages \ + -H "Content-Type: application/json" \ + -d '{ + "userId": "user-123", + "conversationId": "conv-001", + "messages": [ + { "role": "user", "content": "Hi there" }, + { "role": "assistant", "content": "Hello!" } + ] + }' +``` + +Update working memory (append mode): + +```bash +curl -X POST http://localhost:3141/api/memory/conversations/conv-001/working-memory \ + -H "Content-Type: application/json" \ + -d '{ + "content": "Customer prefers email follow-ups.", + "mode": "append" + }' +``` + +Search memory (requires embedding + vector adapters): + +```bash +curl "http://localhost:3141/api/memory/search?searchQuery=refund%20policy&limit=5" +``` diff --git a/packages/cloudflare-d1/src/memory-adapter.ts b/packages/cloudflare-d1/src/memory-adapter.ts index 5fbb2943a..45e3f6d91 100644 --- a/packages/cloudflare-d1/src/memory-adapter.ts +++ b/packages/cloudflare-d1/src/memory-adapter.ts @@ -789,6 +789,23 @@ export class D1MemoryAdapter implements StorageAdapter { ); } + async deleteMessages( + messageIds: string[], + userId: string, + conversationId: string, + ): Promise { + await this.ensureInitialized(); + + if (messageIds.length === 0) { + return; + } + + const messagesTable = `${this.tablePrefix}_messages`; + const placeholders = messageIds.map(() => "?").join(","); + const sql = `DELETE FROM ${messagesTable} WHERE conversation_id = ? AND user_id = ? AND message_id IN (${placeholders})`; + await this.run(sql, [conversationId, userId, ...messageIds]); + } + // ========================================================================== // Conversation Operations // ========================================================================== @@ -962,6 +979,27 @@ export class D1MemoryAdapter implements StorageAdapter { })); } + async countConversations(options: ConversationQueryOptions): Promise { + await this.ensureInitialized(); + + const conversationsTable = `${this.tablePrefix}_conversations`; + let sql = `SELECT COUNT(*) as count FROM ${conversationsTable} WHERE 1=1`; + const args: unknown[] = []; + + if (options.userId) { + sql += " AND user_id = ?"; + args.push(options.userId); + } + + if (options.resourceId) { + sql += " AND resource_id = ?"; + args.push(options.resourceId); + } + + const rows = await this.all<{ count?: number }>(sql, args); + return rows[0]?.count ?? 0; + } + async updateConversation( id: string, updates: Partial>, diff --git a/packages/core/src/index.ts b/packages/core/src/index.ts index b57009ddb..2ed88b2b5 100644 --- a/packages/core/src/index.ts +++ b/packages/core/src/index.ts @@ -298,6 +298,7 @@ export type { ManagedMemoryAddMessagesInput, ManagedMemoryGetMessagesInput, ManagedMemoryClearMessagesInput, + ManagedMemoryDeleteMessagesInput, ManagedMemoryUpdateConversationInput, ManagedMemoryWorkingMemoryInput, ManagedMemorySetWorkingMemoryInput, diff --git a/packages/core/src/memory/adapters/storage/in-memory.ts b/packages/core/src/memory/adapters/storage/in-memory.ts index 7b943ec5e..f29f3c2e7 100644 --- a/packages/core/src/memory/adapters/storage/in-memory.ts +++ b/packages/core/src/memory/adapters/storage/in-memory.ts @@ -116,7 +116,7 @@ export class InMemoryStorageAdapter implements StorageAdapter { options?: GetMessagesOptions, _context?: OperationContext, ): Promise[]> { - const { limit = 100, before, after, roles } = options || {}; + const { limit, before, after, roles } = options || {}; // Get user's messages or return empty array const userMessages = this.storage[userId] || {}; @@ -208,6 +208,25 @@ export class InMemoryStorageAdapter implements StorageAdapter { })); } + /** + * Delete specific messages by ID for a conversation + */ + async deleteMessages( + messageIds: string[], + userId: string, + conversationId: string, + _context?: OperationContext, + ): Promise { + if (!this.storage[userId]?.[conversationId]) { + return; + } + + const ids = new Set(messageIds); + this.storage[userId][conversationId] = this.storage[userId][conversationId].filter( + (message) => !ids.has(message.id), + ); + } + /** * Clear messages for a user */ @@ -336,6 +355,23 @@ export class InMemoryStorageAdapter implements StorageAdapter { return conversations.map((c) => deepClone(c)); } + /** + * Count conversations matching query filters + */ + async countConversations(options: ConversationQueryOptions): Promise { + let conversations = Array.from(this.conversations.values()); + + if (options.userId) { + conversations = conversations.filter((c) => c.userId === options.userId); + } + + if (options.resourceId) { + conversations = conversations.filter((c) => c.resourceId === options.resourceId); + } + + return conversations.length; + } + /** * Update a conversation */ diff --git a/packages/core/src/memory/index.spec-d.ts b/packages/core/src/memory/index.spec-d.ts index 3045d730c..58ec7f869 100644 --- a/packages/core/src/memory/index.spec-d.ts +++ b/packages/core/src/memory/index.spec-d.ts @@ -25,6 +25,7 @@ describe("Memory V2 Type System", () => { addMessages: async () => {}, getMessages: async () => [], clearMessages: async () => {}, + deleteMessages: async () => {}, createConversation: async () => ({ id: "test", resourceId: "res", @@ -38,6 +39,7 @@ describe("Memory V2 Type System", () => { getConversations: async () => [], getConversationsByUserId: async () => [], queryConversations: async () => [], + countConversations: async () => 0, updateConversation: async () => ({ id: "test", resourceId: "res", @@ -152,6 +154,15 @@ describe("Memory V2 Type System", () => { expectTypeOf(adapter.getMessages).returns.toMatchTypeOf>(); }); + it("should enforce messageIds for deleteMessages", () => { + const adapter: StorageAdapter = mockStorageAdapter; + + expectTypeOf(adapter.deleteMessages).parameters.toMatchTypeOf< + [string[], string, string, OperationContext?] + >(); + expectTypeOf(adapter.deleteMessages).returns.toMatchTypeOf>(); + }); + it("should enforce Conversation type for createConversation", () => { const adapter: StorageAdapter = mockStorageAdapter; @@ -169,6 +180,15 @@ describe("Memory V2 Type System", () => { >(); expectTypeOf(adapter.queryConversations).returns.toMatchTypeOf>(); }); + + it("should return number from countConversations", () => { + const adapter: StorageAdapter = mockStorageAdapter; + + expectTypeOf(adapter.countConversations).parameters.toMatchTypeOf< + [ConversationQueryOptions] + >(); + expectTypeOf(adapter.countConversations).returns.toMatchTypeOf>(); + }); }); describe("EmbeddingAdapter Interface", () => { @@ -245,6 +265,18 @@ describe("Memory V2 Type System", () => { expectTypeOf(memory.clearMessages).returns.toMatchTypeOf>(); }); + + it("should return void for deleteMessages", () => { + const memory = new Memory({ storage: mockStorageAdapter }); + + expectTypeOf(memory.deleteMessages).returns.toMatchTypeOf>(); + }); + + it("should return number for countConversations", () => { + const memory = new Memory({ storage: mockStorageAdapter }); + + expectTypeOf(memory.countConversations).returns.toMatchTypeOf>(); + }); }); describe("Type Parameter Constraints", () => { diff --git a/packages/core/src/memory/index.ts b/packages/core/src/memory/index.ts index b91be0239..6bee01a7b 100644 --- a/packages/core/src/memory/index.ts +++ b/packages/core/src/memory/index.ts @@ -134,6 +134,31 @@ export class Memory { return this.storage.clearMessages(userId, conversationId, context); } + /** + * Delete specific messages by ID for a conversation + * Adapters should delete atomically when possible; otherwise a best-effort delete may be used. + */ + async deleteMessages( + messageIds: string[], + userId: string, + conversationId: string, + context?: OperationContext, + ): Promise { + await this.storage.deleteMessages(messageIds, userId, conversationId, context); + + if (this.vector && messageIds.length > 0) { + try { + const vectorIds = messageIds.map((id) => `msg_${conversationId}_${id}`); + await this.vector.deleteBatch(vectorIds); + } catch (error) { + console.warn( + `Failed to delete vectors for conversation ${conversationId} messages:`, + error, + ); + } + } + } + async getConversationSteps( userId: string, conversationId: string, @@ -180,6 +205,13 @@ export class Memory { return this.storage.queryConversations(options); } + /** + * Count conversations with the same filtering as queryConversations (ignores limit/offset) + */ + async countConversations(options: ConversationQueryOptions): Promise { + return this.storage.countConversations(options); + } + /** * Create a new conversation */ diff --git a/packages/core/src/memory/types.ts b/packages/core/src/memory/types.ts index 506e8f817..99cd4d465 100644 --- a/packages/core/src/memory/types.ts +++ b/packages/core/src/memory/types.ts @@ -358,6 +358,17 @@ export interface StorageAdapter { context?: OperationContext, ): Promise[]>; clearMessages(userId: string, conversationId?: string, context?: OperationContext): Promise; + /** + * Delete specific messages by ID for a conversation. + * Adapters should perform an atomic delete when possible. If atomic deletes or transactions + * are unavailable, a best-effort deletion (for example, clear + rehydrate) may be used. + */ + deleteMessages( + messageIds: string[], + userId: string, + conversationId: string, + context?: OperationContext, + ): Promise; // Conversation operations createConversation(input: CreateConversationInput): Promise; @@ -368,6 +379,10 @@ export interface StorageAdapter { options?: Omit, ): Promise; queryConversations(options: ConversationQueryOptions): Promise; + /** + * Count conversations matching query filters (limit/offset ignored). + */ + countConversations(options: ConversationQueryOptions): Promise; updateConversation( id: string, updates: Partial>, diff --git a/packages/core/src/voltops/client.ts b/packages/core/src/voltops/client.ts index 45bd37f1c..baa42f2bf 100644 --- a/packages/core/src/voltops/client.ts +++ b/packages/core/src/voltops/client.ts @@ -32,6 +32,7 @@ import type { ManagedMemoryCredentialCreateResult, ManagedMemoryCredentialListResult, ManagedMemoryDatabaseSummary, + ManagedMemoryDeleteMessagesInput, ManagedMemoryDeleteVectorsInput, ManagedMemoryGetConversationStepsInput, ManagedMemoryGetMessagesInput, @@ -439,6 +440,7 @@ export class VoltOpsClient implements IVoltOpsClient { addBatch: (databaseId, input) => this.addManagedMemoryMessages(databaseId, input), list: (databaseId, input) => this.getManagedMemoryMessages(databaseId, input), clear: (databaseId, input) => this.clearManagedMemoryMessages(databaseId, input), + delete: (databaseId, input) => this.deleteManagedMemoryMessages(databaseId, input), }, conversations: { create: (databaseId, input) => this.createManagedMemoryConversation(databaseId, input), @@ -600,6 +602,21 @@ export class VoltOpsClient implements IVoltOpsClient { } } + private async deleteManagedMemoryMessages( + databaseId: string, + input: ManagedMemoryDeleteMessagesInput, + ): Promise { + const payload = await this.request<{ success: boolean }>( + "POST", + `/managed-memory/projects/databases/${databaseId}/messages/delete`, + input, + ); + + if (!payload?.success) { + throw new Error("Failed to delete managed memory messages via VoltOps"); + } + } + private async storeManagedMemoryVector( databaseId: string, input: ManagedMemoryStoreVectorInput, diff --git a/packages/core/src/voltops/types.ts b/packages/core/src/voltops/types.ts index d0cbcddbb..002d1ccb5 100644 --- a/packages/core/src/voltops/types.ts +++ b/packages/core/src/voltops/types.ts @@ -1116,6 +1116,12 @@ export interface ManagedMemoryClearMessagesInput { conversationId?: string; } +export interface ManagedMemoryDeleteMessagesInput { + conversationId: string; + userId: string; + messageIds: string[]; +} + export interface ManagedMemoryGetConversationStepsInput { conversationId: string; userId: string; @@ -1178,6 +1184,7 @@ export interface ManagedMemoryMessagesClient { addBatch(databaseId: string, input: ManagedMemoryAddMessagesInput): Promise; list(databaseId: string, input: ManagedMemoryGetMessagesInput): Promise; clear(databaseId: string, input: ManagedMemoryClearMessagesInput): Promise; + delete(databaseId: string, input: ManagedMemoryDeleteMessagesInput): Promise; } export interface ManagedMemoryConversationsClient { diff --git a/packages/libsql/src/memory-core.ts b/packages/libsql/src/memory-core.ts index 43b41267c..aeeb7e188 100644 --- a/packages/libsql/src/memory-core.ts +++ b/packages/libsql/src/memory-core.ts @@ -735,6 +735,26 @@ export class LibSQLMemoryCore implements StorageAdapter { } } + async deleteMessages( + messageIds: string[], + userId: string, + conversationId: string, + ): Promise { + await this.initialize(); + + if (messageIds.length === 0) { + return; + } + + const messagesTable = `${this.tablePrefix}_messages`; + const placeholders = messageIds.map(() => "?").join(","); + const sql = `DELETE FROM ${messagesTable} WHERE conversation_id = ? AND user_id = ? AND message_id IN (${placeholders})`; + await this.client.execute({ + sql, + args: [conversationId, userId, ...messageIds], + }); + } + // ============================================================================ // Conversation Operations // ============================================================================ @@ -876,6 +896,28 @@ export class LibSQLMemoryCore implements StorageAdapter { })); } + async countConversations(options: ConversationQueryOptions): Promise { + await this.initialize(); + + const conversationsTable = `${this.tablePrefix}_conversations`; + let sql = `SELECT COUNT(*) as count FROM ${conversationsTable} WHERE 1=1`; + const args: any[] = []; + + if (options.userId) { + sql += " AND user_id = ?"; + args.push(options.userId); + } + + if (options.resourceId) { + sql += " AND resource_id = ?"; + args.push(options.resourceId); + } + + const result = await this.client.execute({ sql, args }); + const count = Number(result.rows[0]?.count ?? 0); + return Number.isNaN(count) ? 0 : count; + } + async updateConversation( id: string, updates: Partial>, diff --git a/packages/postgres/src/memory-adapter.ts b/packages/postgres/src/memory-adapter.ts index a9e44d3ae..e4dc17596 100644 --- a/packages/postgres/src/memory-adapter.ts +++ b/packages/postgres/src/memory-adapter.ts @@ -779,6 +779,33 @@ export class PostgreSQLMemoryAdapter implements StorageAdapter { } } + /** + * Delete specific messages by ID for a conversation + */ + async deleteMessages( + messageIds: string[], + userId: string, + conversationId: string, + ): Promise { + await this.initPromise; + + if (messageIds.length === 0) { + return; + } + + const client = await this.pool.connect(); + try { + const messagesTable = this.getTableName(`${this.tablePrefix}_messages`); + await client.query( + `DELETE FROM ${messagesTable} + WHERE conversation_id = $1 AND user_id = $2 AND message_id = ANY($3::text[])`, + [conversationId, userId, messageIds], + ); + } finally { + client.release(); + } + } + // ============================================================================ // Conversation Operations // ============================================================================ @@ -961,6 +988,39 @@ export class PostgreSQLMemoryAdapter implements StorageAdapter { } } + /** + * Count conversations with filters + */ + async countConversations(options: ConversationQueryOptions): Promise { + await this.initPromise; + + const client = await this.pool.connect(); + try { + const conversationsTable = this.getTableName(`${this.tablePrefix}_conversations`); + let sql = `SELECT COUNT(*) as count FROM ${conversationsTable} WHERE 1=1`; + const params: any[] = []; + let paramCount = 1; + + if (options.userId) { + sql += ` AND user_id = $${paramCount}`; + params.push(options.userId); + paramCount++; + } + + if (options.resourceId) { + sql += ` AND resource_id = $${paramCount}`; + params.push(options.resourceId); + paramCount++; + } + + const result = await client.query(sql, params); + const count = Number(result.rows[0]?.count ?? 0); + return Number.isNaN(count) ? 0 : count; + } finally { + client.release(); + } + } + /** * Update a conversation */ diff --git a/packages/server-core/src/auth/defaults.spec.ts b/packages/server-core/src/auth/defaults.spec.ts index d56ea2c68..37149f60f 100644 --- a/packages/server-core/src/auth/defaults.spec.ts +++ b/packages/server-core/src/auth/defaults.spec.ts @@ -127,6 +127,12 @@ describe("Auth Defaults", () => { expect(requiresAuth("DELETE", "/observability/spans/123")).toBe(true); }); + it("should require auth for memory endpoints via wildcard", () => { + expect(requiresAuth("GET", "/api/memory/conversations")).toBe(true); + expect(requiresAuth("GET", "/api/memory/conversations/conv-1/messages")).toBe(true); + expect(requiresAuth("POST", "/api/memory/save-messages")).toBe(true); + }); + it("should require auth for system update endpoints", () => { expect(requiresAuth("GET", "/updates")).toBe(true); expect(requiresAuth("POST", "/updates")).toBe(true); diff --git a/packages/server-core/src/auth/defaults.ts b/packages/server-core/src/auth/defaults.ts index 54541474d..80282df0b 100644 --- a/packages/server-core/src/auth/defaults.ts +++ b/packages/server-core/src/auth/defaults.ts @@ -110,6 +110,11 @@ export const PROTECTED_ROUTES = [ "POST /workflows/:id/executions/:executionId/resume", // Resume execution "POST /workflows/:id/executions/:executionId/cancel", // Cancel execution + // ======================================== + // MEMORY (User Data) + // ======================================== + "/api/memory/*", // All memory endpoints (GET/POST/PATCH/DELETE) + // ======================================== // OBSERVABILITY (Admin/Internal Tooling) // Covers all observability endpoints: diff --git a/packages/server-core/src/edge.ts b/packages/server-core/src/edge.ts index aa22b20d4..3c0077450 100644 --- a/packages/server-core/src/edge.ts +++ b/packages/server-core/src/edge.ts @@ -4,6 +4,7 @@ export { AGENT_ROUTES, WORKFLOW_ROUTES, A2A_ROUTES, + MEMORY_ROUTES, OBSERVABILITY_ROUTES, OBSERVABILITY_MEMORY_ROUTES, LOG_ROUTES, @@ -35,6 +36,21 @@ export { getWorkingMemoryHandler, } from "./handlers/memory-observability.handlers"; +export { + handleListMemoryConversations, + handleGetMemoryConversation, + handleListMemoryConversationMessages, + handleGetMemoryWorkingMemory, + handleSaveMemoryMessages, + handleCreateMemoryConversation, + handleUpdateMemoryConversation, + handleDeleteMemoryConversation, + handleCloneMemoryConversation, + handleUpdateMemoryWorkingMemory, + handleDeleteMemoryMessages, + handleSearchMemory, +} from "./handlers/memory.handlers"; + export { resolveAgentCard, executeA2ARequest, diff --git a/packages/server-core/src/handlers/memory.handlers.ts b/packages/server-core/src/handlers/memory.handlers.ts new file mode 100644 index 000000000..31401f9c1 --- /dev/null +++ b/packages/server-core/src/handlers/memory.handlers.ts @@ -0,0 +1,842 @@ +import { + AgentRegistry, + type Conversation, + ConversationAlreadyExistsError, + ConversationNotFoundError, + EmbeddingAdapterNotConfiguredError, + type Memory, + type ServerProviderDeps, + VectorAdapterNotConfiguredError, +} from "@voltagent/core"; +import { safeStringify } from "@voltagent/internal"; +import { type UIMessage, generateId } from "ai"; +import type { ApiResponse } from "../types"; + +type MemoryResolution = + | { + ok: true; + memory: Memory; + agentId?: string; + agentName?: string; + resourceId?: string; + } + | { + ok: false; + error: string; + httpStatus?: number; + }; + +function resolveMemory( + deps: { agentRegistry: ServerProviderDeps["agentRegistry"] }, + agentId?: string, +): MemoryResolution { + if (agentId) { + const agent = deps.agentRegistry.getAgent(agentId); + if (!agent) { + return { + ok: false, + error: `Agent ${agentId} not found`, + httpStatus: 404, + }; + } + + const memory = agent.getMemory(); + if (!memory) { + return { + ok: false, + error: `Memory not configured for agent ${agentId}`, + httpStatus: 400, + }; + } + + const state = agent.getFullState(); + return { + ok: true, + memory, + agentId: state.id, + agentName: state.name, + resourceId: state.id, + }; + } + + const registry = AgentRegistry.getInstance(); + const globalMemory = registry.getGlobalMemory(); + if (globalMemory) { + return { ok: true, memory: globalMemory }; + } + + const agents = deps.agentRegistry.getAllAgents(); + const agentsWithMemory = agents.filter((agent) => { + const memory = agent.getMemory(); + return typeof memory === "object" && memory !== null; + }); + + if (agentsWithMemory.length === 1) { + const agent = agentsWithMemory[0]; + const memory = agent.getMemory() as Memory; + const state = agent.getFullState(); + return { + ok: true, + memory, + agentId: state.id, + agentName: state.name, + resourceId: state.id, + }; + } + + if (agentsWithMemory.length > 1) { + return { + ok: false, + error: "agentId is required when multiple agents are configured", + httpStatus: 400, + }; + } + + return { + ok: false, + error: "Memory not configured", + httpStatus: 400, + }; +} + +function buildErrorResponse(error: unknown): ApiResponse { + return { + success: false, + error: error instanceof Error ? error.message : safeStringify(error), + }; +} + +export async function handleListMemoryConversations( + deps: ServerProviderDeps, + query: { + agentId?: string; + resourceId?: string; + userId?: string; + limit?: number; + offset?: number; + orderBy?: "created_at" | "updated_at" | "title"; + orderDirection?: "ASC" | "DESC"; + }, +): Promise< + ApiResponse<{ conversations: Conversation[]; total: number; limit: number; offset: number }> +> { + try { + const resolved = resolveMemory(deps, query.agentId); + if (!resolved.ok) { + return { + success: false, + error: resolved.error, + httpStatus: resolved.httpStatus, + }; + } + + const resourceId = query.resourceId ?? resolved.resourceId; + const [conversations, total] = await Promise.all([ + resolved.memory.queryConversations({ + userId: query.userId, + resourceId, + limit: query.limit, + offset: query.offset, + orderBy: query.orderBy, + orderDirection: query.orderDirection, + }), + resolved.memory.countConversations({ + userId: query.userId, + resourceId, + }), + ]); + + return { + success: true, + data: { + conversations, + total, + limit: query.limit ?? conversations.length, + offset: query.offset ?? 0, + }, + }; + } catch (error) { + return buildErrorResponse(error); + } +} + +export async function handleGetMemoryConversation( + deps: ServerProviderDeps, + conversationId: string, + query: { agentId?: string }, +): Promise> { + try { + const resolved = resolveMemory(deps, query.agentId); + if (!resolved.ok) { + return { + success: false, + error: resolved.error, + httpStatus: resolved.httpStatus, + }; + } + + const conversation = await resolved.memory.getConversation(conversationId); + if (!conversation) { + return { + success: false, + error: "Conversation not found", + httpStatus: 404, + }; + } + + return { + success: true, + data: { conversation }, + }; + } catch (error) { + return buildErrorResponse(error); + } +} + +export async function handleListMemoryConversationMessages( + deps: ServerProviderDeps, + conversationId: string, + query: { + agentId?: string; + limit?: number; + before?: Date; + after?: Date; + roles?: string[]; + userId?: string; + }, +): Promise> { + try { + const resolved = resolveMemory(deps, query.agentId); + if (!resolved.ok) { + return { + success: false, + error: resolved.error, + httpStatus: resolved.httpStatus, + }; + } + + const conversation = await resolved.memory.getConversation(conversationId); + if (!conversation) { + return { + success: false, + error: "Conversation not found", + httpStatus: 404, + }; + } + + const userId = query.userId ?? conversation.userId; + const messages = await resolved.memory.getMessages(userId, conversationId, { + limit: query.limit, + before: query.before, + after: query.after, + roles: query.roles, + }); + + return { + success: true, + data: { conversation, messages }, + }; + } catch (error) { + return buildErrorResponse(error); + } +} + +export async function handleGetMemoryWorkingMemory( + deps: ServerProviderDeps, + conversationId: string, + query: { agentId?: string; scope?: "conversation" | "user"; userId?: string }, +): Promise< + ApiResponse<{ + content: string | null; + format: "markdown" | "json" | null; + template: string | null; + scope: "conversation" | "user"; + }> +> { + try { + const resolved = resolveMemory(deps, query.agentId); + if (!resolved.ok) { + return { + success: false, + error: resolved.error, + httpStatus: resolved.httpStatus, + }; + } + + const scope = query.scope === "user" ? "user" : "conversation"; + let content: string | null = null; + let userId = query.userId; + + if (scope === "conversation") { + const conversation = await resolved.memory.getConversation(conversationId); + if (!conversation) { + return { + success: false, + error: "Conversation not found", + httpStatus: 404, + }; + } + + userId = userId ?? conversation.userId; + content = await resolved.memory.getWorkingMemory({ + conversationId, + userId, + }); + } else { + if (!userId) { + return { + success: false, + error: "userId is required for user-scoped working memory", + httpStatus: 400, + }; + } + + content = await resolved.memory.getWorkingMemory({ + userId, + }); + } + + if (content === null) { + return { + success: false, + error: "Working memory not found", + httpStatus: 404, + }; + } + + return { + success: true, + data: { + content, + scope, + format: resolved.memory.getWorkingMemoryFormat?.() ?? null, + template: resolved.memory.getWorkingMemoryTemplate?.() ?? null, + }, + }; + } catch (error) { + return buildErrorResponse(error); + } +} + +type SaveMessageEntry = + | (UIMessage & { userId?: string; conversationId?: string }) + | { message: UIMessage; userId?: string; conversationId?: string }; + +export async function handleSaveMemoryMessages( + deps: ServerProviderDeps, + body: { + agentId?: string; + userId?: string; + conversationId?: string; + messages?: SaveMessageEntry[]; + }, +): Promise> { + try { + const resolved = resolveMemory(deps, body.agentId); + if (!resolved.ok) { + return { + success: false, + error: resolved.error, + httpStatus: resolved.httpStatus, + }; + } + + if (!Array.isArray(body.messages) || body.messages.length === 0) { + return { + success: false, + error: "messages array is required", + httpStatus: 400, + }; + } + + const normalized = body.messages.map((entry) => { + const isWrapped = typeof entry === "object" && entry !== null && "message" in entry; + const message = isWrapped ? entry.message : (entry as UIMessage); + const conversationId = + (isWrapped ? entry.conversationId : entry.conversationId) ?? body.conversationId; + const userId = (isWrapped ? entry.userId : entry.userId) ?? body.userId; + + return { + message: { + ...message, + id: message.id || generateId(), + }, + conversationId, + userId, + }; + }); + + const missing = normalized.filter((item) => !item.conversationId || !item.userId); + if (missing.length > 0) { + return { + success: false, + error: "Each message must include conversationId and userId", + httpStatus: 400, + }; + } + + const conversationCache = new Map(); + for (const item of normalized) { + const conversationId = item.conversationId as string; + if (!conversationCache.has(conversationId)) { + const conversation = await resolved.memory.getConversation(conversationId); + if (!conversation) { + return { + success: false, + error: `Conversation not found: ${conversationId}`, + httpStatus: 404, + }; + } + conversationCache.set(conversationId, conversation); + } + } + + for (const item of normalized) { + const conversation = conversationCache.get(item.conversationId as string); + if (conversation && conversation.userId !== item.userId) { + return { + success: false, + error: `userId does not match conversation ${conversation.id}`, + httpStatus: 400, + }; + } + } + + const grouped = new Map< + string, + { userId: string; conversationId: string; messages: UIMessage[] } + >(); + for (const item of normalized) { + const conversationId = item.conversationId as string; + const userId = item.userId as string; + const key = `${userId}:${conversationId}`; + const entry = grouped.get(key) || { userId, conversationId, messages: [] }; + entry.messages.push(item.message); + grouped.set(key, entry); + } + + for (const entry of grouped.values()) { + await resolved.memory.addMessages(entry.messages, entry.userId, entry.conversationId); + } + + return { + success: true, + data: { saved: normalized.length }, + }; + } catch (error) { + return buildErrorResponse(error); + } +} + +export async function handleCreateMemoryConversation( + deps: ServerProviderDeps, + body: { + agentId?: string; + conversationId?: string; + resourceId?: string; + userId?: string; + title?: string; + metadata?: Record; + }, +): Promise> { + try { + const resolved = resolveMemory(deps, body.agentId); + if (!resolved.ok) { + return { + success: false, + error: resolved.error, + httpStatus: resolved.httpStatus, + }; + } + + if (!body.userId) { + return { + success: false, + error: "userId is required", + httpStatus: 400, + }; + } + + const resourceId = body.resourceId ?? resolved.resourceId; + if (!resourceId) { + return { + success: false, + error: "resourceId is required", + httpStatus: 400, + }; + } + + const conversationId = body.conversationId ?? generateId(); + const conversation = await resolved.memory.createConversation({ + id: conversationId, + resourceId, + userId: body.userId, + title: body.title ?? "", + metadata: body.metadata ?? {}, + }); + + return { + success: true, + data: { conversation }, + }; + } catch (error) { + if (error instanceof ConversationAlreadyExistsError) { + return { + success: false, + error: error.message, + httpStatus: 409, + }; + } + return buildErrorResponse(error); + } +} + +export async function handleUpdateMemoryConversation( + deps: ServerProviderDeps, + conversationId: string, + body: { + agentId?: string; + resourceId?: string; + userId?: string; + title?: string; + metadata?: Record; + }, +): Promise> { + try { + const resolved = resolveMemory(deps, body.agentId); + if (!resolved.ok) { + return { + success: false, + error: resolved.error, + httpStatus: resolved.httpStatus, + }; + } + + const updates: Partial> = {}; + if (body.resourceId !== undefined) { + updates.resourceId = body.resourceId; + } + if (body.userId !== undefined) { + updates.userId = body.userId; + } + if (body.title !== undefined) { + updates.title = body.title; + } + if (body.metadata !== undefined) { + updates.metadata = body.metadata; + } + + if (Object.keys(updates).length === 0) { + return { + success: false, + error: "No updates provided", + httpStatus: 400, + }; + } + + const conversation = await resolved.memory.updateConversation(conversationId, updates); + return { + success: true, + data: { conversation }, + }; + } catch (error) { + if (error instanceof ConversationNotFoundError) { + return { + success: false, + error: error.message, + httpStatus: 404, + }; + } + return buildErrorResponse(error); + } +} + +export async function handleDeleteMemoryConversation( + deps: ServerProviderDeps, + conversationId: string, + query: { agentId?: string }, +): Promise> { + try { + const resolved = resolveMemory(deps, query.agentId); + if (!resolved.ok) { + return { + success: false, + error: resolved.error, + httpStatus: resolved.httpStatus, + }; + } + + await resolved.memory.deleteConversation(conversationId); + return { + success: true, + data: { deleted: true }, + }; + } catch (error) { + if (error instanceof ConversationNotFoundError) { + return { + success: false, + error: error.message, + httpStatus: 404, + }; + } + return buildErrorResponse(error); + } +} + +export async function handleCloneMemoryConversation( + deps: ServerProviderDeps, + conversationId: string, + body: { + agentId?: string; + newConversationId?: string; + resourceId?: string; + userId?: string; + title?: string; + metadata?: Record; + includeMessages?: boolean; + }, +): Promise> { + try { + const resolved = resolveMemory(deps, body.agentId); + if (!resolved.ok) { + return { + success: false, + error: resolved.error, + httpStatus: resolved.httpStatus, + }; + } + + const source = await resolved.memory.getConversation(conversationId); + if (!source) { + return { + success: false, + error: "Conversation not found", + httpStatus: 404, + }; + } + + const clonedId = body.newConversationId ?? generateId(); + const conversation = await resolved.memory.createConversation({ + id: clonedId, + resourceId: body.resourceId ?? source.resourceId, + userId: body.userId ?? source.userId, + title: body.title ?? source.title, + metadata: body.metadata ?? source.metadata, + }); + + let messageCount = 0; + const includeMessages = body.includeMessages !== false; + if (includeMessages) { + const messages = await resolved.memory.getMessages(source.userId, conversationId); + if (messages.length > 0) { + await resolved.memory.addMessages(messages, conversation.userId, conversation.id); + messageCount = messages.length; + } + } + + return { + success: true, + data: { conversation, messageCount }, + }; + } catch (error) { + if (error instanceof ConversationAlreadyExistsError) { + return { + success: false, + error: error.message, + httpStatus: 409, + }; + } + return buildErrorResponse(error); + } +} + +export async function handleUpdateMemoryWorkingMemory( + deps: ServerProviderDeps, + conversationId: string, + body: { + agentId?: string; + userId?: string; + content?: string | Record; + mode?: "replace" | "append"; + }, +): Promise> { + try { + const resolved = resolveMemory(deps, body.agentId); + if (!resolved.ok) { + return { + success: false, + error: resolved.error, + httpStatus: resolved.httpStatus, + }; + } + + if (body.content === undefined) { + return { + success: false, + error: "content is required", + httpStatus: 400, + }; + } + + const conversation = await resolved.memory.getConversation(conversationId); + if (!conversation) { + return { + success: false, + error: "Conversation not found", + httpStatus: 404, + }; + } + + const userId = body.userId ?? conversation.userId; + if (body.userId && body.userId !== conversation.userId) { + return { + success: false, + error: `userId does not match conversation ${conversation.id}`, + httpStatus: 400, + }; + } + await resolved.memory.updateWorkingMemory({ + conversationId, + userId, + content: body.content, + options: body.mode ? { mode: body.mode } : undefined, + }); + + return { + success: true, + data: { updated: true }, + }; + } catch (error) { + return buildErrorResponse(error); + } +} + +export async function handleDeleteMemoryMessages( + deps: ServerProviderDeps, + body: { + agentId?: string; + conversationId?: string; + userId?: string; + messageIds?: string[]; + }, +): Promise> { + try { + const resolved = resolveMemory(deps, body.agentId); + if (!resolved.ok) { + return { + success: false, + error: resolved.error, + httpStatus: resolved.httpStatus, + }; + } + + if (!Array.isArray(body.messageIds) || body.messageIds.length === 0) { + return { + success: false, + error: "messageIds array is required", + httpStatus: 400, + }; + } + + if (!body.conversationId || !body.userId) { + return { + success: false, + error: "conversationId and userId are required", + httpStatus: 400, + }; + } + + const conversation = await resolved.memory.getConversation(body.conversationId); + if (!conversation) { + return { + success: false, + error: "Conversation not found", + httpStatus: 404, + }; + } + if (conversation.userId !== body.userId) { + return { + success: false, + error: `userId does not match conversation ${conversation.id}`, + httpStatus: 400, + }; + } + + const messages = await resolved.memory.getMessages(body.userId, body.conversationId); + const idsToDelete = new Set(body.messageIds); + const deleted = messages.filter((message) => idsToDelete.has(message.id)).length; + await resolved.memory.deleteMessages(body.messageIds, body.userId, body.conversationId); + return { + success: true, + data: { deleted }, + }; + } catch (error) { + return buildErrorResponse(error); + } +} + +export async function handleSearchMemory( + deps: ServerProviderDeps, + query: { + agentId?: string; + searchQuery?: string; + limit?: number; + threshold?: number; + conversationId?: string; + userId?: string; + }, +): Promise> { + try { + if (!query.searchQuery) { + return { + success: false, + error: "searchQuery is required", + httpStatus: 400, + }; + } + + const resolved = resolveMemory(deps, query.agentId); + if (!resolved.ok) { + return { + success: false, + error: resolved.error, + httpStatus: resolved.httpStatus, + }; + } + + const filter: Record = {}; + if (query.conversationId) { + filter.conversationId = query.conversationId; + } + if (query.userId) { + filter.userId = query.userId; + } + + const results = await resolved.memory.searchSimilar(query.searchQuery, { + limit: query.limit, + threshold: query.threshold, + filter: Object.keys(filter).length > 0 ? filter : undefined, + }); + + return { + success: true, + data: { + results, + count: results.length, + query: query.searchQuery, + }, + }; + } catch (error) { + if ( + error instanceof EmbeddingAdapterNotConfiguredError || + error instanceof VectorAdapterNotConfiguredError + ) { + return { + success: false, + error: error.message, + httpStatus: 400, + }; + } + return buildErrorResponse(error); + } +} diff --git a/packages/server-core/src/index.ts b/packages/server-core/src/index.ts index 1fe7e206a..464b7ca42 100644 --- a/packages/server-core/src/index.ts +++ b/packages/server-core/src/index.ts @@ -27,6 +27,7 @@ export * from "./handlers/log.handlers"; export * from "./handlers/update.handlers"; export * from "./handlers/observability.handlers"; export * from "./handlers/trigger.handlers"; +export * from "./handlers/memory.handlers"; export * from "./handlers/memory-observability.handlers"; export { setupObservabilityHandler } from "./handlers/observability-setup.handler"; diff --git a/packages/server-core/src/routes/definitions.ts b/packages/server-core/src/routes/definitions.ts index 1d270c6ec..ba7bc80ed 100644 --- a/packages/server-core/src/routes/definitions.ts +++ b/packages/server-core/src/routes/definitions.ts @@ -798,7 +798,7 @@ export const OBSERVABILITY_MEMORY_ROUTES = { description: "Retrieve conversations stored in memory with optional filtering by agent or user. Results are paginated and sorted by last update by default.", tags: ["Observability", "Memory"], - operationId: "listMemoryConversations", + operationId: "listObservabilityMemoryConversations", responses: { 200: { description: "Successfully retrieved conversations", @@ -932,6 +932,293 @@ export const TOOL_ROUTES = { }, } as const; +/** + * Memory route definitions + */ +export const MEMORY_ROUTES = { + listConversations: { + method: "get" as const, + path: "/api/memory/conversations", + summary: "List memory conversations", + description: + "Retrieve conversations stored in memory with optional filtering by resource or user.", + tags: ["Memory"], + operationId: "listMemoryConversations", + responses: { + 200: { + description: "Successfully retrieved memory conversations", + contentType: "application/json", + }, + 400: { + description: "Invalid query parameters", + contentType: "application/json", + }, + 500: { + description: "Failed to list memory conversations", + contentType: "application/json", + }, + }, + }, + getConversation: { + method: "get" as const, + path: "/api/memory/conversations/:conversationId", + summary: "Get conversation by ID", + description: "Retrieve a single conversation by ID from memory storage.", + tags: ["Memory"], + operationId: "getMemoryConversation", + responses: { + 200: { + description: "Successfully retrieved conversation", + contentType: "application/json", + }, + 404: { + description: "Conversation not found", + contentType: "application/json", + }, + 500: { + description: "Failed to retrieve conversation", + contentType: "application/json", + }, + }, + }, + listMessages: { + method: "get" as const, + path: "/api/memory/conversations/:conversationId/messages", + summary: "List conversation messages", + description: "Retrieve messages for a conversation with optional filtering.", + tags: ["Memory"], + operationId: "listMemoryConversationMessages", + responses: { + 200: { + description: "Successfully retrieved conversation messages", + contentType: "application/json", + }, + 404: { + description: "Conversation not found", + contentType: "application/json", + }, + 500: { + description: "Failed to retrieve conversation messages", + contentType: "application/json", + }, + }, + }, + getMemoryWorkingMemory: { + method: "get" as const, + path: "/api/memory/conversations/:conversationId/working-memory", + summary: "Get working memory", + description: "Retrieve working memory content for a conversation.", + tags: ["Memory"], + operationId: "getMemoryWorkingMemory", + responses: { + 200: { + description: "Successfully retrieved working memory", + contentType: "application/json", + }, + 404: { + description: "Working memory not found", + contentType: "application/json", + }, + 500: { + description: "Failed to retrieve working memory", + contentType: "application/json", + }, + }, + }, + saveMessages: { + method: "post" as const, + path: "/api/memory/save-messages", + summary: "Save messages", + description: "Persist new messages into memory storage.", + tags: ["Memory"], + operationId: "saveMemoryMessages", + responses: { + 200: { + description: "Successfully saved messages", + contentType: "application/json", + }, + 400: { + description: "Invalid request body", + contentType: "application/json", + }, + 500: { + description: "Failed to save messages", + contentType: "application/json", + }, + }, + }, + createConversation: { + method: "post" as const, + path: "/api/memory/conversations", + summary: "Create conversation", + description: "Create a new conversation in memory storage.", + tags: ["Memory"], + operationId: "createMemoryConversation", + responses: { + 200: { + description: "Successfully created conversation", + contentType: "application/json", + }, + 400: { + description: "Invalid request body", + contentType: "application/json", + }, + 409: { + description: "Conversation already exists", + contentType: "application/json", + }, + 500: { + description: "Failed to create conversation", + contentType: "application/json", + }, + }, + }, + updateConversation: { + method: "patch" as const, + path: "/api/memory/conversations/:conversationId", + summary: "Update conversation", + description: "Update an existing conversation in memory storage.", + tags: ["Memory"], + operationId: "updateMemoryConversation", + responses: { + 200: { + description: "Successfully updated conversation", + contentType: "application/json", + }, + 400: { + description: "Invalid request body", + contentType: "application/json", + }, + 404: { + description: "Conversation not found", + contentType: "application/json", + }, + 500: { + description: "Failed to update conversation", + contentType: "application/json", + }, + }, + }, + deleteConversation: { + method: "delete" as const, + path: "/api/memory/conversations/:conversationId", + summary: "Delete conversation", + description: "Delete a conversation and its messages from memory storage.", + tags: ["Memory"], + operationId: "deleteMemoryConversation", + responses: { + 200: { + description: "Successfully deleted conversation", + contentType: "application/json", + }, + 404: { + description: "Conversation not found", + contentType: "application/json", + }, + 500: { + description: "Failed to delete conversation", + contentType: "application/json", + }, + }, + }, + cloneConversation: { + method: "post" as const, + path: "/api/memory/conversations/:conversationId/clone", + summary: "Clone conversation", + description: "Create a copy of a conversation, optionally including messages.", + tags: ["Memory"], + operationId: "cloneMemoryConversation", + responses: { + 200: { + description: "Successfully cloned conversation", + contentType: "application/json", + }, + 404: { + description: "Conversation not found", + contentType: "application/json", + }, + 409: { + description: "Conversation already exists", + contentType: "application/json", + }, + 500: { + description: "Failed to clone conversation", + contentType: "application/json", + }, + }, + }, + updateWorkingMemory: { + method: "post" as const, + path: "/api/memory/conversations/:conversationId/working-memory", + summary: "Update working memory", + description: "Update working memory content for a conversation.", + tags: ["Memory"], + operationId: "updateMemoryWorkingMemory", + responses: { + 200: { + description: "Successfully updated working memory", + contentType: "application/json", + }, + 400: { + description: "Invalid request body", + contentType: "application/json", + }, + 404: { + description: "Conversation not found", + contentType: "application/json", + }, + 500: { + description: "Failed to update working memory", + contentType: "application/json", + }, + }, + }, + deleteMessages: { + method: "post" as const, + path: "/api/memory/messages/delete", + summary: "Delete messages", + description: "Delete specific messages from memory storage.", + tags: ["Memory"], + operationId: "deleteMemoryMessages", + responses: { + 200: { + description: "Successfully deleted messages", + contentType: "application/json", + }, + 400: { + description: "Invalid request body", + contentType: "application/json", + }, + 500: { + description: "Failed to delete messages", + contentType: "application/json", + }, + }, + }, + searchMemory: { + method: "get" as const, + path: "/api/memory/search", + summary: "Search memory", + description: "Search memory using semantic search when available.", + tags: ["Memory"], + operationId: "searchMemory", + responses: { + 200: { + description: "Successfully searched memory", + contentType: "application/json", + }, + 400: { + description: "Invalid query parameters", + contentType: "application/json", + }, + 500: { + description: "Failed to search memory", + contentType: "application/json", + }, + }, + }, +} as const; + /** * All route definitions combined */ @@ -941,6 +1228,7 @@ export const ALL_ROUTES = { ...TOOL_ROUTES, ...LOG_ROUTES, ...UPDATE_ROUTES, + ...MEMORY_ROUTES, ...OBSERVABILITY_ROUTES, ...OBSERVABILITY_MEMORY_ROUTES, } as const; diff --git a/packages/server-elysia/src/app-factory.ts b/packages/server-elysia/src/app-factory.ts index 3386abe91..bea629e07 100644 --- a/packages/server-elysia/src/app-factory.ts +++ b/packages/server-elysia/src/app-factory.ts @@ -14,6 +14,7 @@ import { registerAgentRoutes, registerLogRoutes, registerMcpRoutes, + registerMemoryRoutes, registerObservabilityRoutes, registerToolRoutes, registerTriggerRoutes, @@ -46,6 +47,7 @@ export async function createApp( logs: () => registerLogRoutes(app, deps, logger), updates: () => registerUpdateRoutes(app, deps, logger), observability: () => registerObservabilityRoutes(app, deps, logger), + memory: () => registerMemoryRoutes(app, deps, logger), tools: () => registerToolRoutes(app, deps, logger), triggers: () => registerTriggerRoutes(app, deps, logger), mcp: () => registerMcpRoutes(app, deps as any, logger), @@ -133,6 +135,7 @@ export async function createApp( routes.logs(); routes.updates(); routes.observability(); + routes.memory(); routes.triggers(); routes.mcp(); routes.a2a(); diff --git a/packages/server-elysia/src/routes/index.ts b/packages/server-elysia/src/routes/index.ts index acd565957..44e6e520f 100644 --- a/packages/server-elysia/src/routes/index.ts +++ b/packages/server-elysia/src/routes/index.ts @@ -7,3 +7,4 @@ export { registerA2ARoutes } from "./a2a.routes"; export { registerToolRoutes } from "./tool.routes"; export { registerTriggerRoutes } from "./trigger.routes"; export { registerObservabilityRoutes } from "./observability"; +export { registerMemoryRoutes } from "./memory.routes"; diff --git a/packages/server-elysia/src/routes/memory.routes.ts b/packages/server-elysia/src/routes/memory.routes.ts new file mode 100644 index 000000000..17e92db06 --- /dev/null +++ b/packages/server-elysia/src/routes/memory.routes.ts @@ -0,0 +1,210 @@ +import type { ServerProviderDeps } from "@voltagent/core"; +import type { Logger } from "@voltagent/internal"; +import { + MEMORY_ROUTES, + handleCloneMemoryConversation, + handleCreateMemoryConversation, + handleDeleteMemoryConversation, + handleDeleteMemoryMessages, + handleGetMemoryConversation, + handleGetMemoryWorkingMemory, + handleListMemoryConversationMessages, + handleListMemoryConversations, + handleSaveMemoryMessages, + handleSearchMemory, + handleUpdateMemoryConversation, + handleUpdateMemoryWorkingMemory, +} from "@voltagent/server-core"; +import type { Elysia } from "elysia"; + +function parseNumber(value?: string | number): number | undefined { + if (value === undefined || value === null) { + return undefined; + } + const parsed = typeof value === "number" ? value : Number.parseInt(value as string, 10); + return Number.isNaN(parsed) ? undefined : parsed; +} + +function parseFloatValue(value?: string | number): number | undefined { + if (value === undefined || value === null) { + return undefined; + } + const parsed = typeof value === "number" ? value : Number.parseFloat(value as string); + return Number.isNaN(parsed) ? undefined : parsed; +} + +function parseDate(value?: string): Date | undefined { + if (!value) { + return undefined; + } + const parsed = new Date(value); + return Number.isNaN(parsed.getTime()) ? undefined : parsed; +} + +type MemoryRoutesCompat = typeof MEMORY_ROUTES & { + getWorkingMemory?: { path: string }; +}; + +const memoryWorkingMemoryPath = + (MEMORY_ROUTES as MemoryRoutesCompat).getMemoryWorkingMemory?.path ?? + (MEMORY_ROUTES as MemoryRoutesCompat).getWorkingMemory?.path ?? + "/api/memory/conversations/:conversationId/working-memory"; + +/** + * Register memory routes + */ +export function registerMemoryRoutes(app: Elysia, deps: ServerProviderDeps, logger: Logger) { + app.get(MEMORY_ROUTES.listConversations.path, async ({ query, set }) => { + logger.trace("GET /api/memory/conversations - fetching conversations", { query }); + const response = await handleListMemoryConversations(deps, { + agentId: query.agentId as string | undefined, + resourceId: query.resourceId as string | undefined, + userId: query.userId as string | undefined, + limit: parseNumber(query.limit as string | number | undefined), + offset: parseNumber(query.offset as string | number | undefined), + orderBy: query.orderBy as "created_at" | "updated_at" | "title" | undefined, + orderDirection: query.orderDirection as "ASC" | "DESC" | undefined, + }); + set.status = response.success ? 200 : (response.httpStatus ?? 500); + return response; + }); + + app.get(MEMORY_ROUTES.getConversation.path, async ({ params, query, set }) => { + const conversationId = params.conversationId; + logger.trace(`GET /api/memory/conversations/${conversationId} - fetching conversation`); + const response = await handleGetMemoryConversation(deps, conversationId, { + agentId: query.agentId as string | undefined, + }); + set.status = response.success ? 200 : (response.httpStatus ?? 500); + return response; + }); + + app.get(MEMORY_ROUTES.listMessages.path, async ({ params, query, set }) => { + const conversationId = params.conversationId; + logger.trace(`GET /api/memory/conversations/${conversationId}/messages - fetching messages`, { + query, + }); + const response = await handleListMemoryConversationMessages(deps, conversationId, { + agentId: query.agentId as string | undefined, + limit: parseNumber(query.limit as string | number | undefined), + before: parseDate(query.before as string | undefined), + after: parseDate(query.after as string | undefined), + roles: query.roles ? String(query.roles).split(",") : undefined, + userId: query.userId as string | undefined, + }); + set.status = response.success ? 200 : (response.httpStatus ?? 500); + return response; + }); + + app.get(memoryWorkingMemoryPath, async ({ params, query, set }) => { + const conversationId = params.conversationId; + logger.trace( + `GET /api/memory/conversations/${conversationId}/working-memory - fetching working memory`, + { query }, + ); + const response = await handleGetMemoryWorkingMemory(deps, conversationId, { + agentId: query.agentId as string | undefined, + scope: query.scope === "user" ? "user" : "conversation", + userId: query.userId as string | undefined, + }); + set.status = response.success ? 200 : (response.httpStatus ?? 500); + return response; + }); + + app.post(MEMORY_ROUTES.saveMessages.path, async ({ body, query, set }) => { + const payload = body as Record | undefined; + logger.trace("POST /api/memory/save-messages - saving messages", { + messageCount: Array.isArray(payload?.messages) ? payload?.messages.length : 0, + }); + const response = await handleSaveMemoryMessages(deps, { + ...(payload ?? {}), + agentId: (payload?.agentId as string | undefined) ?? (query.agentId as string | undefined), + }); + set.status = response.success ? 200 : (response.httpStatus ?? 500); + return response; + }); + + app.post(MEMORY_ROUTES.createConversation.path, async ({ body, query, set }) => { + const payload = body as Record | undefined; + logger.trace("POST /api/memory/conversations - creating conversation"); + const response = await handleCreateMemoryConversation(deps, { + ...(payload ?? {}), + agentId: (payload?.agentId as string | undefined) ?? (query.agentId as string | undefined), + }); + set.status = response.success ? 200 : (response.httpStatus ?? 500); + return response; + }); + + app.patch(MEMORY_ROUTES.updateConversation.path, async ({ params, body, query, set }) => { + const conversationId = params.conversationId; + const payload = body as Record | undefined; + logger.trace(`PATCH /api/memory/conversations/${conversationId} - updating conversation`); + const response = await handleUpdateMemoryConversation(deps, conversationId, { + ...(payload ?? {}), + agentId: (payload?.agentId as string | undefined) ?? (query.agentId as string | undefined), + }); + set.status = response.success ? 200 : (response.httpStatus ?? 500); + return response; + }); + + app.delete(MEMORY_ROUTES.deleteConversation.path, async ({ params, query, set }) => { + const conversationId = params.conversationId; + logger.trace(`DELETE /api/memory/conversations/${conversationId} - deleting conversation`); + const response = await handleDeleteMemoryConversation(deps, conversationId, { + agentId: query.agentId as string | undefined, + }); + set.status = response.success ? 200 : (response.httpStatus ?? 500); + return response; + }); + + app.post(MEMORY_ROUTES.cloneConversation.path, async ({ params, body, query, set }) => { + const conversationId = params.conversationId; + const payload = body as Record | undefined; + logger.trace(`POST /api/memory/conversations/${conversationId}/clone - cloning conversation`); + const response = await handleCloneMemoryConversation(deps, conversationId, { + ...(payload ?? {}), + agentId: (payload?.agentId as string | undefined) ?? (query.agentId as string | undefined), + }); + set.status = response.success ? 200 : (response.httpStatus ?? 500); + return response; + }); + + app.post(MEMORY_ROUTES.updateWorkingMemory.path, async ({ params, body, query, set }) => { + const conversationId = params.conversationId; + const payload = body as Record | undefined; + logger.trace( + `POST /api/memory/conversations/${conversationId}/working-memory - updating working memory`, + ); + const response = await handleUpdateMemoryWorkingMemory(deps, conversationId, { + ...(payload ?? {}), + agentId: (payload?.agentId as string | undefined) ?? (query.agentId as string | undefined), + }); + set.status = response.success ? 200 : (response.httpStatus ?? 500); + return response; + }); + + app.post(MEMORY_ROUTES.deleteMessages.path, async ({ body, query, set }) => { + const payload = body as Record | undefined; + logger.trace("POST /api/memory/messages/delete - deleting messages"); + const response = await handleDeleteMemoryMessages(deps, { + ...(payload ?? {}), + agentId: (payload?.agentId as string | undefined) ?? (query.agentId as string | undefined), + }); + set.status = response.success ? 200 : (response.httpStatus ?? 500); + return response; + }); + + app.get(MEMORY_ROUTES.searchMemory.path, async ({ query, set }) => { + logger.trace("GET /api/memory/search - searching memory", { query }); + const response = await handleSearchMemory(deps, { + agentId: query.agentId as string | undefined, + searchQuery: query.searchQuery as string | undefined, + limit: parseNumber(query.limit as string | number | undefined), + threshold: parseFloatValue(query.threshold as string | number | undefined), + conversationId: query.conversationId as string | undefined, + userId: query.userId as string | undefined, + }); + set.status = response.success ? 200 : (response.httpStatus ?? 500); + return response; + }); +} diff --git a/packages/server-hono/src/app-factory.ts b/packages/server-hono/src/app-factory.ts index 51ef5d800..20b2d8312 100644 --- a/packages/server-hono/src/app-factory.ts +++ b/packages/server-hono/src/app-factory.ts @@ -14,6 +14,7 @@ import { registerAgentRoutes, registerLogRoutes, registerMcpRoutes, + registerMemoryRoutes, registerObservabilityRoutes, registerToolRoutes, registerTriggerRoutes, @@ -53,6 +54,7 @@ export async function createApp( logs: () => registerLogRoutes(app as any, resolvedDeps, logger), updates: () => registerUpdateRoutes(app as any, resolvedDeps, logger), observability: () => registerObservabilityRoutes(app as any, resolvedDeps, logger), + memory: () => registerMemoryRoutes(app as any, resolvedDeps, logger), tools: () => registerToolRoutes(app as any, resolvedDeps as any, logger), triggers: () => registerTriggerRoutes(app as any, resolvedDeps, logger), mcp: () => registerMcpRoutes(app as any, resolvedDeps as any, logger), @@ -121,6 +123,7 @@ export async function createApp( routes.logs(); routes.updates(); routes.observability(); + routes.memory(); routes.triggers(); routes.mcp(); routes.a2a(); diff --git a/packages/server-hono/src/routes/index.ts b/packages/server-hono/src/routes/index.ts index 221590716..9c3b547e7 100644 --- a/packages/server-hono/src/routes/index.ts +++ b/packages/server-hono/src/routes/index.ts @@ -48,6 +48,7 @@ export { registerMcpRoutes } from "./mcp.routes"; export { registerA2ARoutes } from "./a2a.routes"; export { registerToolRoutes } from "./tool.routes"; export { registerTriggerRoutes } from "./trigger.routes"; +export { registerMemoryRoutes } from "./memory.routes"; /** * Register agent routes diff --git a/packages/server-hono/src/routes/memory.routes.ts b/packages/server-hono/src/routes/memory.routes.ts new file mode 100644 index 000000000..44e4ff58e --- /dev/null +++ b/packages/server-hono/src/routes/memory.routes.ts @@ -0,0 +1,271 @@ +import type { ServerProviderDeps } from "@voltagent/core"; +import type { Logger } from "@voltagent/internal"; +import { + MEMORY_ROUTES, + handleCloneMemoryConversation, + handleCreateMemoryConversation, + handleDeleteMemoryConversation, + handleDeleteMemoryMessages, + handleGetMemoryConversation, + handleGetMemoryWorkingMemory, + handleListMemoryConversationMessages, + handleListMemoryConversations, + handleSaveMemoryMessages, + handleSearchMemory, + handleUpdateMemoryConversation, + handleUpdateMemoryWorkingMemory, +} from "@voltagent/server-core"; +import type { OpenAPIHonoType } from "../zod-openapi-compat"; + +function parseNumber(value?: string): number | undefined { + if (!value) { + return undefined; + } + const parsed = Number.parseInt(value, 10); + return Number.isNaN(parsed) ? undefined : parsed; +} + +function parseFloatValue(value?: string): number | undefined { + if (!value) { + return undefined; + } + const parsed = Number.parseFloat(value); + return Number.isNaN(parsed) ? undefined : parsed; +} + +function parseDate(value?: string): Date | undefined { + if (!value) { + return undefined; + } + const parsed = new Date(value); + return Number.isNaN(parsed.getTime()) ? undefined : parsed; +} + +const orderByAllowlist = new Set(["created_at", "updated_at", "title"]); + +function parseOrderBy(value?: string): "created_at" | "updated_at" | "title" | undefined { + if (!value || !orderByAllowlist.has(value)) { + return undefined; + } + return value as "created_at" | "updated_at" | "title"; +} + +function parseOrderDirection(value?: string): "ASC" | "DESC" | undefined { + if (!value) { + return undefined; + } + const normalized = value.toUpperCase(); + if (normalized === "ASC" || normalized === "DESC") { + return normalized; + } + return undefined; +} + +type MemoryRoutesCompat = typeof MEMORY_ROUTES & { + getWorkingMemory?: { path: string }; +}; + +const memoryWorkingMemoryPath = + (MEMORY_ROUTES as MemoryRoutesCompat).getMemoryWorkingMemory?.path ?? + (MEMORY_ROUTES as MemoryRoutesCompat).getWorkingMemory?.path ?? + "/api/memory/conversations/:conversationId/working-memory"; + +/** + * Register memory routes + */ +export function registerMemoryRoutes( + app: OpenAPIHonoType, + deps: ServerProviderDeps, + logger: Logger, +) { + app.get(MEMORY_ROUTES.listConversations.path, async (c) => { + const query = c.req.query(); + logger.trace("GET /api/memory/conversations - fetching conversations", { query }); + const response = await handleListMemoryConversations(deps, { + agentId: query.agentId, + resourceId: query.resourceId, + userId: query.userId, + limit: parseNumber(query.limit), + offset: parseNumber(query.offset), + orderBy: parseOrderBy(query.orderBy), + orderDirection: parseOrderDirection(query.orderDirection), + }); + + return c.json(response, response.success ? 200 : (response.httpStatus ?? 500)); + }); + + app.get(MEMORY_ROUTES.getConversation.path, async (c) => { + const conversationId = c.req.param("conversationId"); + const query = c.req.query(); + logger.trace(`GET /api/memory/conversations/${conversationId} - fetching conversation`); + const response = await handleGetMemoryConversation(deps, conversationId, { + agentId: query.agentId, + }); + return c.json(response, response.success ? 200 : (response.httpStatus ?? 500)); + }); + + app.get(MEMORY_ROUTES.listMessages.path, async (c) => { + const conversationId = c.req.param("conversationId"); + const query = c.req.query(); + logger.trace(`GET /api/memory/conversations/${conversationId}/messages - fetching messages`, { + query, + }); + const response = await handleListMemoryConversationMessages(deps, conversationId, { + agentId: query.agentId, + limit: parseNumber(query.limit), + before: parseDate(query.before), + after: parseDate(query.after), + roles: query.roles ? query.roles.split(",") : undefined, + userId: query.userId, + }); + return c.json(response, response.success ? 200 : (response.httpStatus ?? 500)); + }); + + app.get(memoryWorkingMemoryPath, async (c) => { + const conversationId = c.req.param("conversationId"); + const query = c.req.query(); + logger.trace( + `GET /api/memory/conversations/${conversationId}/working-memory - fetching working memory`, + { query }, + ); + const response = await handleGetMemoryWorkingMemory(deps, conversationId, { + agentId: query.agentId, + scope: query.scope === "user" ? "user" : "conversation", + userId: query.userId, + }); + return c.json(response, response.success ? 200 : (response.httpStatus ?? 500)); + }); + + app.post(MEMORY_ROUTES.saveMessages.path, async (c) => { + const query = c.req.query(); + let body: any; + try { + body = await c.req.json(); + } catch (error) { + logger.warn("Invalid JSON body for save messages", { error }); + return c.json({ success: false, error: "Invalid JSON body" }, 400); + } + logger.trace("POST /api/memory/save-messages - saving messages", { + messageCount: Array.isArray(body?.messages) ? body.messages.length : 0, + }); + const response = await handleSaveMemoryMessages(deps, { + ...body, + agentId: body?.agentId ?? query.agentId, + }); + return c.json(response, response.success ? 200 : (response.httpStatus ?? 500)); + }); + + app.post(MEMORY_ROUTES.createConversation.path, async (c) => { + const query = c.req.query(); + let body: any; + try { + body = await c.req.json(); + } catch (error) { + logger.warn("Invalid JSON body for create conversation", { error }); + return c.json({ success: false, error: "Invalid JSON body" }, 400); + } + logger.trace("POST /api/memory/conversations - creating conversation"); + const response = await handleCreateMemoryConversation(deps, { + ...body, + agentId: body?.agentId ?? query.agentId, + }); + return c.json(response, response.success ? 200 : (response.httpStatus ?? 500)); + }); + + app.patch(MEMORY_ROUTES.updateConversation.path, async (c) => { + const conversationId = c.req.param("conversationId"); + const query = c.req.query(); + let body: any; + try { + body = await c.req.json(); + } catch (error) { + logger.warn("Invalid JSON body for update conversation", { error, conversationId }); + return c.json({ success: false, error: "Invalid JSON body" }, 400); + } + logger.trace(`PATCH /api/memory/conversations/${conversationId} - updating conversation`); + const response = await handleUpdateMemoryConversation(deps, conversationId, { + ...body, + agentId: body?.agentId ?? query.agentId, + }); + return c.json(response, response.success ? 200 : (response.httpStatus ?? 500)); + }); + + app.delete(MEMORY_ROUTES.deleteConversation.path, async (c) => { + const conversationId = c.req.param("conversationId"); + const query = c.req.query(); + logger.trace(`DELETE /api/memory/conversations/${conversationId} - deleting conversation`); + const response = await handleDeleteMemoryConversation(deps, conversationId, { + agentId: query.agentId, + }); + return c.json(response, response.success ? 200 : (response.httpStatus ?? 500)); + }); + + app.post(MEMORY_ROUTES.cloneConversation.path, async (c) => { + const conversationId = c.req.param("conversationId"); + const query = c.req.query(); + let body: any; + try { + body = await c.req.json(); + } catch (error) { + logger.warn("Invalid JSON body for clone conversation", { error, conversationId }); + return c.json({ success: false, error: "Invalid JSON body" }, 400); + } + logger.trace(`POST /api/memory/conversations/${conversationId}/clone - cloning conversation`); + const response = await handleCloneMemoryConversation(deps, conversationId, { + ...body, + agentId: body?.agentId ?? query.agentId, + }); + return c.json(response, response.success ? 200 : (response.httpStatus ?? 500)); + }); + + app.post(MEMORY_ROUTES.updateWorkingMemory.path, async (c) => { + const conversationId = c.req.param("conversationId"); + const query = c.req.query(); + let body: any; + try { + body = await c.req.json(); + } catch (error) { + logger.warn("Invalid JSON body for update working memory", { error, conversationId }); + return c.json({ success: false, error: "Invalid JSON body" }, 400); + } + logger.trace( + `POST /api/memory/conversations/${conversationId}/working-memory - updating working memory`, + ); + const response = await handleUpdateMemoryWorkingMemory(deps, conversationId, { + ...body, + agentId: body?.agentId ?? query.agentId, + }); + return c.json(response, response.success ? 200 : (response.httpStatus ?? 500)); + }); + + app.post(MEMORY_ROUTES.deleteMessages.path, async (c) => { + const query = c.req.query(); + let body: any; + try { + body = await c.req.json(); + } catch (error) { + logger.warn("Invalid JSON body for delete messages", { error }); + return c.json({ success: false, error: "Invalid JSON body" }, 400); + } + logger.trace("POST /api/memory/messages/delete - deleting messages"); + const response = await handleDeleteMemoryMessages(deps, { + ...body, + agentId: body?.agentId ?? query.agentId, + }); + return c.json(response, response.success ? 200 : (response.httpStatus ?? 500)); + }); + + app.get(MEMORY_ROUTES.searchMemory.path, async (c) => { + const query = c.req.query(); + logger.trace("GET /api/memory/search - searching memory", { query }); + const response = await handleSearchMemory(deps, { + agentId: query.agentId, + searchQuery: query.searchQuery, + limit: parseNumber(query.limit), + threshold: parseFloatValue(query.threshold), + conversationId: query.conversationId, + userId: query.userId, + }); + return c.json(response, response.success ? 200 : (response.httpStatus ?? 500)); + }); +} diff --git a/packages/serverless-hono/src/app-factory.ts b/packages/serverless-hono/src/app-factory.ts index 270852cbc..f0338c78a 100644 --- a/packages/serverless-hono/src/app-factory.ts +++ b/packages/serverless-hono/src/app-factory.ts @@ -8,6 +8,7 @@ import { registerA2ARoutes, registerAgentRoutes, registerLogRoutes, + registerMemoryRoutes, registerObservabilityRoutes, registerToolRoutes, registerTriggerRoutes, @@ -81,6 +82,7 @@ export async function createServerlessApp(deps: ServerProviderDeps, config?: Ser registerToolRoutes(app, resolvedDeps, logger); registerLogRoutes(app, resolvedDeps, logger); registerUpdateRoutes(app, resolvedDeps, logger); + registerMemoryRoutes(app, resolvedDeps, logger); registerObservabilityRoutes(app, resolvedDeps, logger); registerTriggerRoutes(app, resolvedDeps, logger); registerA2ARoutes(app, resolvedDeps, logger); diff --git a/packages/serverless-hono/src/routes.ts b/packages/serverless-hono/src/routes.ts index 07105b9b5..30b4c49d4 100644 --- a/packages/serverless-hono/src/routes.ts +++ b/packages/serverless-hono/src/routes.ts @@ -26,6 +26,7 @@ import { type A2ARequestContext, A2A_ROUTES, AGENT_ROUTES, + MEMORY_ROUTES, OBSERVABILITY_MEMORY_ROUTES, OBSERVABILITY_ROUTES, TOOL_ROUTES, @@ -38,6 +39,10 @@ import { getConversationStepsHandler, handleChatStream, handleCheckUpdates, + handleCloneMemoryConversation, + handleCreateMemoryConversation, + handleDeleteMemoryConversation, + handleDeleteMemoryMessages, handleExecuteTool, handleExecuteWorkflow, handleGenerateObject, @@ -46,18 +51,26 @@ import { handleGetAgentHistory, handleGetAgents, handleGetLogs, + handleGetMemoryConversation, + handleGetMemoryWorkingMemory, handleGetWorkflow, handleGetWorkflowState, handleGetWorkflows, handleInstallUpdates, + handleListMemoryConversationMessages, + handleListMemoryConversations, handleListTools, handleListWorkflowRuns, handleResumeChatStream, handleResumeWorkflow, + handleSaveMemoryMessages, + handleSearchMemory, handleStreamObject, handleStreamText, handleStreamWorkflow, handleSuspendWorkflow, + handleUpdateMemoryConversation, + handleUpdateMemoryWorkingMemory, isErrorResponse, mapLogResponse, parseJsonRpcRequest, @@ -83,6 +96,39 @@ async function readJsonBody(c: any, logger: Logger): Promise { } } +function parseNumber(value?: string): number | undefined { + if (!value) { + return undefined; + } + const parsed = Number.parseInt(value, 10); + return Number.isNaN(parsed) ? undefined : parsed; +} + +function parseFloatValue(value?: string): number | undefined { + if (!value) { + return undefined; + } + const parsed = Number.parseFloat(value); + return Number.isNaN(parsed) ? undefined : parsed; +} + +function parseDate(value?: string): Date | undefined { + if (!value) { + return undefined; + } + const parsed = new Date(value); + return Number.isNaN(parsed.getTime()) ? undefined : parsed; +} + +type MemoryRoutesCompat = typeof MEMORY_ROUTES & { + getWorkingMemory?: { path: string }; +}; + +const memoryWorkingMemoryPath = + (MEMORY_ROUTES as MemoryRoutesCompat).getMemoryWorkingMemory?.path ?? + (MEMORY_ROUTES as MemoryRoutesCompat).getWorkingMemory?.path ?? + "/api/memory/conversations/:conversationId/working-memory"; + function extractHeaders( headers: Headers | NodeJS.Dict, ): Record { @@ -492,6 +538,159 @@ export function registerUpdateRoutes(app: Hono, deps: ServerProviderDeps, logger }); } +export function registerMemoryRoutes(app: Hono, deps: ServerProviderDeps, logger: Logger) { + app.get(MEMORY_ROUTES.listConversations.path, async (c) => { + const query = c.req.query(); + const response = await handleListMemoryConversations(deps, { + agentId: query.agentId, + resourceId: query.resourceId, + userId: query.userId, + limit: parseNumber(query.limit), + offset: parseNumber(query.offset), + orderBy: query.orderBy as "created_at" | "updated_at" | "title" | undefined, + orderDirection: query.orderDirection as "ASC" | "DESC" | undefined, + }); + return c.json(response, response.success ? 200 : (response.httpStatus ?? 500)); + }); + + app.get(MEMORY_ROUTES.getConversation.path, async (c) => { + const conversationId = c.req.param("conversationId"); + const query = c.req.query(); + const response = await handleGetMemoryConversation(deps, conversationId, { + agentId: query.agentId, + }); + return c.json(response, response.success ? 200 : (response.httpStatus ?? 500)); + }); + + app.get(MEMORY_ROUTES.listMessages.path, async (c) => { + const conversationId = c.req.param("conversationId"); + const query = c.req.query(); + const response = await handleListMemoryConversationMessages(deps, conversationId, { + agentId: query.agentId, + limit: parseNumber(query.limit), + before: parseDate(query.before), + after: parseDate(query.after), + roles: query.roles ? query.roles.split(",") : undefined, + userId: query.userId, + }); + return c.json(response, response.success ? 200 : (response.httpStatus ?? 500)); + }); + + app.get(memoryWorkingMemoryPath, async (c) => { + const conversationId = c.req.param("conversationId"); + const query = c.req.query(); + const response = await handleGetMemoryWorkingMemory(deps, conversationId, { + agentId: query.agentId, + scope: query.scope === "user" ? "user" : "conversation", + userId: query.userId, + }); + return c.json(response, response.success ? 200 : (response.httpStatus ?? 500)); + }); + + app.post(MEMORY_ROUTES.saveMessages.path, async (c) => { + const query = c.req.query(); + const body = await readJsonBody(c, logger); + if (!body) { + return c.json({ success: false, error: "Invalid JSON body" }, 400); + } + const response = await handleSaveMemoryMessages(deps, { + ...body, + agentId: (body.agentId as string | undefined) ?? query.agentId, + }); + return c.json(response, response.success ? 200 : (response.httpStatus ?? 500)); + }); + + app.post(MEMORY_ROUTES.createConversation.path, async (c) => { + const query = c.req.query(); + const body = await readJsonBody(c, logger); + if (!body) { + return c.json({ success: false, error: "Invalid JSON body" }, 400); + } + const response = await handleCreateMemoryConversation(deps, { + ...body, + agentId: (body.agentId as string | undefined) ?? query.agentId, + }); + return c.json(response, response.success ? 200 : (response.httpStatus ?? 500)); + }); + + app.patch(MEMORY_ROUTES.updateConversation.path, async (c) => { + const conversationId = c.req.param("conversationId"); + const query = c.req.query(); + const body = await readJsonBody(c, logger); + if (!body) { + return c.json({ success: false, error: "Invalid JSON body" }, 400); + } + const response = await handleUpdateMemoryConversation(deps, conversationId, { + ...body, + agentId: (body.agentId as string | undefined) ?? query.agentId, + }); + return c.json(response, response.success ? 200 : (response.httpStatus ?? 500)); + }); + + app.delete(MEMORY_ROUTES.deleteConversation.path, async (c) => { + const conversationId = c.req.param("conversationId"); + const query = c.req.query(); + const response = await handleDeleteMemoryConversation(deps, conversationId, { + agentId: query.agentId, + }); + return c.json(response, response.success ? 200 : (response.httpStatus ?? 500)); + }); + + app.post(MEMORY_ROUTES.cloneConversation.path, async (c) => { + const conversationId = c.req.param("conversationId"); + const query = c.req.query(); + const body = await readJsonBody(c, logger); + if (!body) { + return c.json({ success: false, error: "Invalid JSON body" }, 400); + } + const response = await handleCloneMemoryConversation(deps, conversationId, { + ...body, + agentId: (body.agentId as string | undefined) ?? query.agentId, + }); + return c.json(response, response.success ? 200 : (response.httpStatus ?? 500)); + }); + + app.post(MEMORY_ROUTES.updateWorkingMemory.path, async (c) => { + const conversationId = c.req.param("conversationId"); + const query = c.req.query(); + const body = await readJsonBody(c, logger); + if (!body) { + return c.json({ success: false, error: "Invalid JSON body" }, 400); + } + const response = await handleUpdateMemoryWorkingMemory(deps, conversationId, { + ...body, + agentId: (body.agentId as string | undefined) ?? query.agentId, + }); + return c.json(response, response.success ? 200 : (response.httpStatus ?? 500)); + }); + + app.post(MEMORY_ROUTES.deleteMessages.path, async (c) => { + const query = c.req.query(); + const body = await readJsonBody(c, logger); + if (!body) { + return c.json({ success: false, error: "Invalid JSON body" }, 400); + } + const response = await handleDeleteMemoryMessages(deps, { + ...body, + agentId: (body.agentId as string | undefined) ?? query.agentId, + }); + return c.json(response, response.success ? 200 : (response.httpStatus ?? 500)); + }); + + app.get(MEMORY_ROUTES.searchMemory.path, async (c) => { + const query = c.req.query(); + const response = await handleSearchMemory(deps, { + agentId: query.agentId, + searchQuery: query.searchQuery, + limit: parseNumber(query.limit), + threshold: parseFloatValue(query.threshold), + conversationId: query.conversationId, + userId: query.userId, + }); + return c.json(response, response.success ? 200 : (response.httpStatus ?? 500)); + }); +} + export function registerObservabilityRoutes(app: Hono, deps: ServerProviderDeps, logger: Logger) { app.post(OBSERVABILITY_ROUTES.setupObservability.path, (c) => c.json( diff --git a/packages/supabase/src/memory-adapter.ts b/packages/supabase/src/memory-adapter.ts index 9077c6f6d..8501d09ce 100644 --- a/packages/supabase/src/memory-adapter.ts +++ b/packages/supabase/src/memory-adapter.ts @@ -803,6 +803,33 @@ END OF MIGRATION SQL this.log(`Cleared messages for user ${userId}`); } + /** + * Delete specific messages by ID for a conversation + */ + async deleteMessages( + messageIds: string[], + userId: string, + conversationId: string, + ): Promise { + await this.initialize(); + + if (messageIds.length === 0) { + return; + } + + const messagesTable = `${this.baseTableName}_messages`; + const { error } = await this.client + .from(messagesTable) + .delete() + .eq("conversation_id", conversationId) + .eq("user_id", userId) + .in("message_id", messageIds); + + if (error) { + throw new Error(`Failed to delete messages: ${error.message}`); + } + } + // ============================================================================ // Conversation Operations // ============================================================================ @@ -977,6 +1004,32 @@ END OF MIGRATION SQL })); } + /** + * Count conversations with filters + */ + async countConversations(options: ConversationQueryOptions): Promise { + await this.initialize(); + + const conversationsTable = `${this.baseTableName}_conversations`; + let query = this.client.from(conversationsTable).select("id", { count: "exact", head: true }); + + if (options.userId) { + query = query.eq("user_id", options.userId); + } + + if (options.resourceId) { + query = query.eq("resource_id", options.resourceId); + } + + const { count, error } = await query; + + if (error) { + throw new Error(`Failed to count conversations: ${error.message}`); + } + + return count ?? 0; + } + /** * Update a conversation */ diff --git a/packages/voltagent-memory/src/index.ts b/packages/voltagent-memory/src/index.ts index 81b84edcd..1a8c46db3 100644 --- a/packages/voltagent-memory/src/index.ts +++ b/packages/voltagent-memory/src/index.ts @@ -260,6 +260,29 @@ export class ManagedMemoryAdapter implements StorageAdapter { }).then(() => undefined); } + deleteMessages( + messageIds: string[], + userId: string, + conversationId: string, + _context?: OperationContext, + ): Promise { + return this.withClientContext(async ({ client, database }) => { + if (messageIds.length === 0) { + return; + } + + this.log( + "Deleting managed memory messages", + safeStringify({ count: messageIds.length, userId, conversationId }), + ); + await client.managedMemory.messages.delete(database.id, { + messageIds, + userId, + conversationId, + }); + }).then(() => undefined); + } + createConversation(input: CreateConversationInput): Promise { return this.withClientContext(({ client, database }) => { this.log("Creating managed memory conversation", safeStringify({ conversationId: input.id })); @@ -298,6 +321,35 @@ export class ManagedMemoryAdapter implements StorageAdapter { }); } + countConversations(options: ConversationQueryOptions): Promise { + return this.withClientContext(async ({ client, database }) => { + const pageSize = 200; + let offset = 0; + let total = 0; + + while (true) { + const page = await client.managedMemory.conversations.query(database.id, { + userId: options.userId, + resourceId: options.resourceId, + orderBy: options.orderBy, + orderDirection: options.orderDirection, + limit: pageSize, + offset, + }); + + total += page.length; + + if (page.length < pageSize) { + break; + } + + offset += pageSize; + } + + return total; + }); + } + updateConversation( id: string, updates: Partial>, diff --git a/website/docs/agents/memory.md b/website/docs/agents/memory.md index 337e274b9..fbcfad929 100644 --- a/website/docs/agents/memory.md +++ b/website/docs/agents/memory.md @@ -60,6 +60,7 @@ await agent.generateText("What's my name?", { For detailed configuration, provider setup, and advanced features: - **[Memory Overview](./memory/overview.md)** - Full memory system documentation +- **[Memory API Endpoints](../api/endpoints/memory.md)** - HTTP endpoints for conversations and messages - **[Managed Memory](./memory/managed-memory.md)** - Production-ready hosted storage - **[Semantic Search](./memory/semantic-search.md)** - Vector-based message retrieval - **[Working Memory](./memory/working-memory.md)** - Compact context management diff --git a/website/docs/api/api-reference.md b/website/docs/api/api-reference.md index 38fd46b19..bd607b7c8 100644 --- a/website/docs/api/api-reference.md +++ b/website/docs/api/api-reference.md @@ -91,6 +91,23 @@ Default port is 3141, but may vary based on configuration. } ``` +## Memory Endpoints + +| Method | Path | Description | Auth | +| ------ | ---------------------------------------------------------- | --------------------- | ---- | +| GET | `/api/memory/conversations` | List conversations | Yes | +| GET | `/api/memory/conversations/:conversationId` | Get conversation | Yes | +| GET | `/api/memory/conversations/:conversationId/messages` | List messages | Yes | +| GET | `/api/memory/conversations/:conversationId/working-memory` | Get working memory | Yes | +| POST | `/api/memory/save-messages` | Save messages | Yes | +| POST | `/api/memory/conversations` | Create conversation | Yes | +| PATCH | `/api/memory/conversations/:conversationId` | Update conversation | Yes | +| DELETE | `/api/memory/conversations/:conversationId` | Delete conversation | Yes | +| POST | `/api/memory/conversations/:conversationId/clone` | Clone conversation | Yes | +| POST | `/api/memory/conversations/:conversationId/working-memory` | Update working memory | Yes | +| POST | `/api/memory/messages/delete` | Delete messages | Yes | +| GET | `/api/memory/search` | Search memory | Yes | + ## Logging & Observability | Method | Path | Description | Auth | @@ -284,6 +301,7 @@ Note: The server reads its port from the `honoServer({ port })` config. `PORT` i - **[Server Architecture](./server-architecture.md)** - Understanding server design - **[Agent Endpoints](./endpoints/agents.md)** - Detailed agent API - **[Workflow Endpoints](./endpoints/workflows.md)** - Workflow execution details +- **[Memory Endpoints](./endpoints/memory.md)** - Conversation and message APIs - **[Authentication](./authentication.md)** - Security and auth setup - **[Custom Endpoints](./custom-endpoints.md)** - Adding custom routes diff --git a/website/docs/api/endpoints/memory.md b/website/docs/api/endpoints/memory.md new file mode 100644 index 000000000..c697f520e --- /dev/null +++ b/website/docs/api/endpoints/memory.md @@ -0,0 +1,214 @@ +--- +title: Memory Endpoints +sidebar_label: Memory +--- + +# Memory API Endpoints + +VoltAgent exposes memory endpoints under `/api/memory/*` for managing conversations, messages, working memory, and semantic search results. + +**Auth:** Protected by default (see [Authentication](../authentication.md)). + +## Common Parameters + +- `agentId` (query/body): Optional. Required when multiple agents are registered or no global memory is configured. +- `resourceId` (query/body): Optional. Defaults to the agent ID when `agentId` is provided. +- `userId` (query/body): Required for creating conversations and saving/deleting messages. +- `conversationId` (path/body): Conversation identifier. It can be generated by the server if omitted when creating a conversation. + +## List Conversations + +**Endpoint:** `GET /api/memory/conversations` + +**Query Parameters:** `agentId`, `resourceId`, `userId`, `limit`, `offset`, `orderBy`, `orderDirection` + +```bash +curl "http://localhost:3141/api/memory/conversations?userId=user-123&limit=20" +``` + +### Response + +```json +{ + "success": true, + "data": { + "conversations": [ + { + "id": "conv-001", + "resourceId": "assistant", + "userId": "user-123", + "title": "Support Chat", + "metadata": {}, + "createdAt": "2025-01-01T12:00:00.000Z", + "updatedAt": "2025-01-01T12:05:00.000Z" + } + ], + "total": 1, + "limit": 20, + "offset": 0 + } +} +``` + +## Get Conversation + +**Endpoint:** `GET /api/memory/conversations/:conversationId` + +```bash +curl "http://localhost:3141/api/memory/conversations/conv-001" +``` + +## Create Conversation + +**Endpoint:** `POST /api/memory/conversations` + +### Request Body + +```json +{ + "userId": "user-123", + "resourceId": "assistant", + "title": "New Chat", + "metadata": { "source": "web" } +} +``` + +## Update Conversation + +**Endpoint:** `PATCH /api/memory/conversations/:conversationId` + +### Request Body + +```json +{ + "title": "Updated Title", + "metadata": { "priority": "high" } +} +``` + +## Delete Conversation + +**Endpoint:** `DELETE /api/memory/conversations/:conversationId` + +```bash +curl -X DELETE "http://localhost:3141/api/memory/conversations/conv-001" +``` + +## Clone Conversation + +**Endpoint:** `POST /api/memory/conversations/:conversationId/clone` + +### Request Body + +```json +{ + "newConversationId": "conv-002", + "title": "Clone of Support Chat", + "includeMessages": true +} +``` + +## List Messages + +**Endpoint:** `GET /api/memory/conversations/:conversationId/messages` + +**Query Parameters:** `agentId`, `limit`, `before`, `after`, `roles`, `userId` + +```bash +curl "http://localhost:3141/api/memory/conversations/conv-001/messages?limit=50" +``` + +Notes: + +- `roles` accepts a comma-separated list (e.g. `user,assistant,tool`). +- `before` and `after` expect ISO 8601 timestamps. + +## Save Messages + +**Endpoint:** `POST /api/memory/save-messages` + +### Request Body + +```json +{ + "userId": "user-123", + "conversationId": "conv-001", + "messages": [ + { + "role": "user", + "content": "Hi there" + }, + { + "message": { + "role": "assistant", + "content": "Hello!" + } + } + ] +} +``` + +Notes: + +- Each message must include `userId` and `conversationId`, either on the message entry or in the request body. +- Message IDs are generated when omitted. + +## Delete Messages + +**Endpoint:** `POST /api/memory/messages/delete` + +### Request Body + +```json +{ + "userId": "user-123", + "conversationId": "conv-001", + "messageIds": ["msg-1", "msg-2"] +} +``` + +## Get Working Memory + +**Endpoint:** `GET /api/memory/conversations/:conversationId/working-memory` + +**Query Parameters:** `agentId`, `scope`, `userId` + +```bash +curl "http://localhost:3141/api/memory/conversations/conv-001/working-memory?scope=conversation" +``` + +Notes: + +- `scope=user` requires `userId` in the query. + +## Update Working Memory + +**Endpoint:** `POST /api/memory/conversations/:conversationId/working-memory` + +### Request Body + +```json +{ + "content": "Customer prefers email follow-ups.", + "mode": "append" +} +``` + +Notes: + +- `userId` is optional, but if provided it must match the conversation owner. +- `content` can be a string or a JSON object when working memory is schema-based. + +## Search Memory + +**Endpoint:** `GET /api/memory/search` + +**Query Parameters:** `searchQuery`, `conversationId`, `userId`, `limit`, `threshold`, `agentId` + +```bash +curl "http://localhost:3141/api/memory/search?searchQuery=refund%20policy&limit=5" +``` + +Notes: + +- Requires embedding and vector adapters; otherwise the endpoint returns `400`. diff --git a/website/docs/api/overview.md b/website/docs/api/overview.md index 8d8790313..e4d414de5 100644 --- a/website/docs/api/overview.md +++ b/website/docs/api/overview.md @@ -114,6 +114,21 @@ Get the raw OpenAPI 3.1 spec at [`http://localhost:3141/doc`](http://localhost:3 - `GET /tools` - List all registered tools (across agents) - `POST /tools/:name/execute` - Execute a tool directly over HTTP +### Memory Endpoints + +- `GET /api/memory/conversations` - List conversations +- `GET /api/memory/conversations/:conversationId` - Get conversation +- `GET /api/memory/conversations/:conversationId/messages` - List messages +- `GET /api/memory/conversations/:conversationId/working-memory` - Get working memory +- `POST /api/memory/save-messages` - Save messages +- `POST /api/memory/conversations` - Create conversation +- `PATCH /api/memory/conversations/:conversationId` - Update conversation +- `DELETE /api/memory/conversations/:conversationId` - Delete conversation +- `POST /api/memory/conversations/:conversationId/clone` - Clone conversation +- `POST /api/memory/conversations/:conversationId/working-memory` - Update working memory +- `POST /api/memory/messages/delete` - Delete messages +- `GET /api/memory/search` - Search memory + ### Observability & Logs - `POST /setup-observability` - Configure `.env` with VoltAgent keys @@ -136,6 +151,7 @@ Get the raw OpenAPI 3.1 spec at [`http://localhost:3141/doc`](http://localhost:3 - **[Server Architecture](./server-architecture.md)** - Understanding the pluggable server design - **[Agent Endpoints](./endpoints/agents.md)** - Complete agent API reference with examples - **[Workflow Endpoints](./endpoints/workflows.md)** - Workflow execution and management +- **[Memory Endpoints](./endpoints/memory.md)** - Conversation and message storage APIs - **[Authentication](./authentication.md)** - Securing your API endpoints - **[Streaming](./streaming.md)** - Real-time features with SSE and WebSocket - **[Custom Endpoints](./custom-endpoints.md)** - Adding your own REST endpoints diff --git a/website/sidebars.ts b/website/sidebars.ts index 0931a5712..3da76beca 100644 --- a/website/sidebars.ts +++ b/website/sidebars.ts @@ -257,7 +257,12 @@ const sidebars: SidebarsConfig = { { type: "category", label: "Endpoints", - items: ["api/endpoints/agents", "api/endpoints/workflows", "api/endpoints/tools"], + items: [ + "api/endpoints/agents", + "api/endpoints/workflows", + "api/endpoints/memory", + "api/endpoints/tools", + ], }, ], },