diff --git a/app/src/utils/__tests__/tauriCommandsMemory.test.ts b/app/src/utils/__tests__/tauriCommandsMemory.test.ts index e66324a1d..73107ee7e 100644 --- a/app/src/utils/__tests__/tauriCommandsMemory.test.ts +++ b/app/src/utils/__tests__/tauriCommandsMemory.test.ts @@ -49,7 +49,7 @@ describe('memoryGraphQuery', () => { const result = await memoryGraphQuery('team', 'Alice', 'OWNS'); expect(mockCallCoreRpc).toHaveBeenCalledWith({ - method: 'memory.graph.query', + method: 'openhuman.memory_graph_query', params: { namespace: 'team', subject: 'Alice', predicate: 'OWNS' }, }); expect(result).toEqual(mockRelations); @@ -62,7 +62,7 @@ describe('memoryGraphQuery', () => { await memoryGraphQuery(); expect(mockCallCoreRpc).toHaveBeenCalledWith({ - method: 'memory.graph.query', + method: 'openhuman.memory_graph_query', params: { namespace: undefined, subject: undefined, predicate: undefined }, }); }); @@ -100,7 +100,7 @@ describe('memoryDocIngest', () => { const result = await memoryDocIngest(params); - expect(mockCallCoreRpc).toHaveBeenCalledWith({ method: 'memory.doc.ingest', params }); + expect(mockCallCoreRpc).toHaveBeenCalledWith({ method: 'openhuman.memory_doc_ingest', params }); expect(result).toEqual(ingestResult); }); @@ -111,6 +111,6 @@ describe('memoryDocIngest', () => { const params = { namespace: 'ns', key: 'k', title: 't', content: 'c' }; await memoryDocIngest(params); - expect(mockCallCoreRpc).toHaveBeenCalledWith({ method: 'memory.doc.ingest', params }); + expect(mockCallCoreRpc).toHaveBeenCalledWith({ method: 'openhuman.memory_doc_ingest', params }); }); }); diff --git a/app/src/utils/tauriCommands.ts b/app/src/utils/tauriCommands.ts index 9f0d67921..2112c1587 100644 --- a/app/src/utils/tauriCommands.ts +++ b/app/src/utils/tauriCommands.ts @@ -199,7 +199,7 @@ export async function syncMemoryClientToken(token: string): Promise { } try { console.debug('[memory] syncMemoryClientToken: payload → memory.init'); - await callCoreRpc({ method: 'memory.init', params: { jwt_token: token } }); + await callCoreRpc({ method: 'openhuman.memory_init', params: { jwt_token: token } }); console.info('[memory] syncMemoryClientToken: exit — ok'); } catch (err) { console.warn('[memory] syncMemoryClientToken: exit — error:', err); @@ -217,14 +217,17 @@ export async function memoryListDocuments(namespace?: string): Promise if (!isTauri()) { throw new Error('Not running in Tauri'); } - return await callCoreRpc({ method: 'memory.list_documents', params: { namespace } }); + return await callCoreRpc({ + method: 'openhuman.memory_list_documents', + params: { namespace }, + }); } export async function memoryListNamespaces(): Promise { if (!isTauri()) { throw new Error('Not running in Tauri'); } - return await callCoreRpc({ method: 'memory.list_namespaces' }); + return await callCoreRpc({ method: 'openhuman.memory_list_namespaces' }); } export async function memoryDeleteDocument( @@ -235,7 +238,7 @@ export async function memoryDeleteDocument( throw new Error('Not running in Tauri'); } return await callCoreRpc({ - method: 'memory.delete_document', + method: 'openhuman.memory_delete_document', params: { document_id: documentId, namespace }, }); } @@ -249,7 +252,7 @@ export async function memoryQueryNamespace( throw new Error('Not running in Tauri'); } return await callCoreRpc({ - method: 'memory.query_namespace', + method: 'openhuman.memory_query_namespace', params: { namespace, query, max_chunks: maxChunks }, }); } @@ -262,7 +265,7 @@ export async function memoryRecallNamespace( throw new Error('Not running in Tauri'); } return await callCoreRpc({ - method: 'memory.recall_namespace', + method: 'openhuman.memory_recall_context', params: { namespace, max_chunks: maxChunks }, }); } @@ -289,7 +292,7 @@ export async function memoryGraphQuery( throw new Error('Not running in Tauri'); } return await callCoreRpc({ - method: 'memory.graph.query', + method: 'openhuman.memory_graph_query', params: { namespace, subject, predicate }, }); } @@ -310,7 +313,7 @@ export async function memoryDocIngest(params: { if (!isTauri()) { throw new Error('Not running in Tauri'); } - return await callCoreRpc({ method: 'memory.doc.ingest', params }); + return await callCoreRpc({ method: 'openhuman.memory_doc_ingest', params }); } export async function aiListMemoryFiles(relativeDir = 'memory'): Promise { @@ -318,7 +321,7 @@ export async function aiListMemoryFiles(relativeDir = 'memory'): Promise({ - method: 'ai.list_memory_files', + method: 'openhuman.memory_list_files', params: { relative_dir: relativeDir }, }); } @@ -328,7 +331,7 @@ export async function aiReadMemoryFile(relativePath: string): Promise { throw new Error('Not running in Tauri'); } return await callCoreRpc({ - method: 'ai.read_memory_file', + method: 'openhuman.memory_read_file', params: { relative_path: relativePath }, }); } @@ -338,7 +341,7 @@ export async function aiWriteMemoryFile(relativePath: string, content: string): throw new Error('Not running in Tauri'); } await callCoreRpc({ - method: 'ai.write_memory_file', + method: 'openhuman.memory_write_file', params: { relative_path: relativePath, content }, }); } diff --git a/src/core/all.rs b/src/core/all.rs index 4a6a1d098..731e1aa20 100644 --- a/src/core/all.rs +++ b/src/core/all.rs @@ -62,6 +62,7 @@ fn build_registered_controllers() -> Vec { controllers.extend(crate::openhuman::skills::all_skills_registered_controllers()); controllers.extend(crate::openhuman::workspace::all_workspace_registered_controllers()); controllers.extend(crate::openhuman::tools::all_tools_registered_controllers()); + controllers.extend(crate::openhuman::memory::all_memory_registered_controllers()); controllers } @@ -89,6 +90,7 @@ fn build_declared_controller_schemas() -> Vec { schemas.extend(crate::openhuman::skills::all_skills_controller_schemas()); schemas.extend(crate::openhuman::workspace::all_workspace_controller_schemas()); schemas.extend(crate::openhuman::tools::all_tools_controller_schemas()); + schemas.extend(crate::openhuman::memory::all_memory_controller_schemas()); schemas } @@ -122,6 +124,7 @@ pub fn namespace_description(namespace: &str) -> Option<&'static str> { "service" => Some("Desktop service lifecycle management."), "skills" => Some("Skill registry, runtime lifecycle, setup, tools, and sync."), "socket" => Some("Skills runtime socket bridge controls."), + "memory" => Some("Document storage, vector search, key-value store, and knowledge graph."), _ => None, } } diff --git a/src/core/jsonrpc.rs b/src/core/jsonrpc.rs index c8ce4b5b1..0f86fe9fe 100644 --- a/src/core/jsonrpc.rs +++ b/src/core/jsonrpc.rs @@ -478,46 +478,46 @@ mod tests { #[tokio::test] async fn invoke_memory_init_missing_required_param_fails() { - let err = invoke_method(default_state(), "memory.init", json!({})) + let err = invoke_method(default_state(), "openhuman.memory_init", json!({})) .await .expect_err("missing jwt_token should fail"); - assert!(err.contains("missing field `jwt_token`") || err.contains("jwt_token")); + assert!(err.contains("jwt_token")); } #[tokio::test] async fn invoke_memory_list_namespaces_rejects_unknown_param() { let err = invoke_method( default_state(), - "memory.list_namespaces", + "openhuman.memory_list_namespaces", json!({ "extra": true }), ) .await .expect_err("unknown param should fail"); - assert!(err.contains("unknown field `extra`") || err.contains("extra")); + assert!(err.contains("extra")); } #[tokio::test] async fn invoke_memory_query_namespace_missing_namespace_fails() { let err = invoke_method( default_state(), - "memory.query_namespace", + "openhuman.memory_query_namespace", json!({ "query": "who owns atlas" }), ) .await .expect_err("missing namespace should fail"); - assert!(err.contains("missing field `namespace`") || err.contains("namespace")); + assert!(err.contains("namespace")); } #[tokio::test] async fn invoke_memory_recall_memories_rejects_unknown_param() { let err = invoke_method( default_state(), - "memory.recall_memories", + "openhuman.memory_recall_memories", json!({ "namespace": "team", "extra": true }), ) .await .expect_err("unknown param should fail"); - assert!(err.contains("unknown field `extra`") || err.contains("extra")); + assert!(err.contains("extra")); } #[tokio::test] diff --git a/src/openhuman/memory/mod.rs b/src/openhuman/memory/mod.rs index 9083420a7..440aec89f 100644 --- a/src/openhuman/memory/mod.rs +++ b/src/openhuman/memory/mod.rs @@ -4,6 +4,7 @@ pub mod ingestion; pub mod ops; pub(crate) mod relex; pub mod rpc_models; +pub mod schemas; pub mod store; pub mod traits; @@ -14,6 +15,10 @@ pub use ingestion::{ pub use ops as rpc; pub use ops::*; pub use rpc_models::*; +pub use schemas::{ + all_controller_schemas as all_memory_controller_schemas, + all_registered_controllers as all_memory_registered_controllers, +}; pub use store::{ create_memory, create_memory_for_migration, create_memory_with_storage, create_memory_with_storage_and_routes, effective_memory_backend_name, MemoryClient, diff --git a/src/openhuman/memory/schemas.rs b/src/openhuman/memory/schemas.rs new file mode 100644 index 000000000..d47f0c800 --- /dev/null +++ b/src/openhuman/memory/schemas.rs @@ -0,0 +1,1076 @@ +use serde::de::DeserializeOwned; +use serde_json::{Map, Value}; + +use crate::core::all::{ControllerFuture, RegisteredController}; +use crate::core::{ControllerSchema, FieldSchema, TypeSchema}; +use crate::openhuman::memory::rpc::{ + self, DeleteDocParams, GraphQueryParams, GraphUpsertParams, IngestDocParams, KvGetDeleteParams, + KvSetParams, NamespaceOnlyParams, PutDocParams, QueryNamespaceParams, RecallNamespaceParams, +}; +use crate::openhuman::memory::{ + DeleteDocumentRequest, EmptyRequest, ListDocumentsRequest, ListMemoryFilesRequest, + MemoryInitRequest, QueryNamespaceRequest, ReadMemoryFileRequest, RecallContextRequest, + RecallMemoriesRequest, WriteMemoryFileRequest, +}; +use crate::rpc::RpcOutcome; + +// --------------------------------------------------------------------------- +// Public entry points +// --------------------------------------------------------------------------- + +pub fn all_controller_schemas() -> Vec { + vec![ + schemas("init"), + schemas("list_documents"), + schemas("list_namespaces"), + schemas("delete_document"), + schemas("query_namespace"), + schemas("recall_context"), + schemas("recall_memories"), + schemas("list_files"), + schemas("read_file"), + schemas("write_file"), + schemas("namespace_list"), + schemas("doc_put"), + schemas("doc_ingest"), + schemas("doc_list"), + schemas("doc_delete"), + schemas("context_query"), + schemas("context_recall"), + schemas("kv_set"), + schemas("kv_get"), + schemas("kv_delete"), + schemas("kv_list_namespace"), + schemas("graph_upsert"), + schemas("graph_query"), + ] +} + +pub fn all_registered_controllers() -> Vec { + vec![ + RegisteredController { + schema: schemas("init"), + handler: handle_init, + }, + RegisteredController { + schema: schemas("list_documents"), + handler: handle_list_documents, + }, + RegisteredController { + schema: schemas("list_namespaces"), + handler: handle_list_namespaces, + }, + RegisteredController { + schema: schemas("delete_document"), + handler: handle_delete_document, + }, + RegisteredController { + schema: schemas("query_namespace"), + handler: handle_query_namespace, + }, + RegisteredController { + schema: schemas("recall_context"), + handler: handle_recall_context, + }, + RegisteredController { + schema: schemas("recall_memories"), + handler: handle_recall_memories, + }, + RegisteredController { + schema: schemas("list_files"), + handler: handle_list_files, + }, + RegisteredController { + schema: schemas("read_file"), + handler: handle_read_file, + }, + RegisteredController { + schema: schemas("write_file"), + handler: handle_write_file, + }, + RegisteredController { + schema: schemas("namespace_list"), + handler: handle_namespace_list, + }, + RegisteredController { + schema: schemas("doc_put"), + handler: handle_doc_put, + }, + RegisteredController { + schema: schemas("doc_ingest"), + handler: handle_doc_ingest, + }, + RegisteredController { + schema: schemas("doc_list"), + handler: handle_doc_list, + }, + RegisteredController { + schema: schemas("doc_delete"), + handler: handle_doc_delete, + }, + RegisteredController { + schema: schemas("context_query"), + handler: handle_context_query, + }, + RegisteredController { + schema: schemas("context_recall"), + handler: handle_context_recall, + }, + RegisteredController { + schema: schemas("kv_set"), + handler: handle_kv_set, + }, + RegisteredController { + schema: schemas("kv_get"), + handler: handle_kv_get, + }, + RegisteredController { + schema: schemas("kv_delete"), + handler: handle_kv_delete, + }, + RegisteredController { + schema: schemas("kv_list_namespace"), + handler: handle_kv_list_namespace, + }, + RegisteredController { + schema: schemas("graph_upsert"), + handler: handle_graph_upsert, + }, + RegisteredController { + schema: schemas("graph_query"), + handler: handle_graph_query, + }, + ] +} + +// --------------------------------------------------------------------------- +// Schema definitions +// --------------------------------------------------------------------------- + +pub fn schemas(function: &str) -> ControllerSchema { + match function { + // ----- legacy envelope-style methods ----- + "init" => ControllerSchema { + namespace: "memory", + function: "init", + description: "Initialise the memory subsystem for the current workspace.", + inputs: vec![FieldSchema { + name: "jwt_token", + ty: TypeSchema::String, + comment: "JWT token for authenticating the memory session.", + required: true, + }], + outputs: vec![FieldSchema { + name: "result", + ty: TypeSchema::Json, + comment: "Envelope with initialisation status, workspace and memory paths.", + required: true, + }], + }, + "list_documents" => ControllerSchema { + namespace: "memory", + function: "list_documents", + description: "List documents stored in memory, optionally filtered by namespace.", + inputs: vec![FieldSchema { + name: "namespace", + ty: TypeSchema::Option(Box::new(TypeSchema::String)), + comment: "Namespace to filter documents by.", + required: false, + }], + outputs: vec![FieldSchema { + name: "result", + ty: TypeSchema::Json, + comment: "Envelope with documents array and count.", + required: true, + }], + }, + "list_namespaces" => ControllerSchema { + namespace: "memory", + function: "list_namespaces", + description: "List all namespaces that contain memory documents.", + inputs: vec![], + outputs: vec![FieldSchema { + name: "result", + ty: TypeSchema::Json, + comment: "Envelope with namespaces array and count.", + required: true, + }], + }, + "delete_document" => ControllerSchema { + namespace: "memory", + function: "delete_document", + description: "Delete a specific document from a namespace.", + inputs: vec![ + FieldSchema { + name: "namespace", + ty: TypeSchema::String, + comment: "Namespace containing the document.", + required: true, + }, + FieldSchema { + name: "document_id", + ty: TypeSchema::String, + comment: "Identifier of the document to delete.", + required: true, + }, + ], + outputs: vec![FieldSchema { + name: "result", + ty: TypeSchema::Json, + comment: "Envelope with deletion status.", + required: true, + }], + }, + "query_namespace" => ControllerSchema { + namespace: "memory", + function: "query_namespace", + description: "Semantic query against a namespace with optional reference data.", + inputs: vec![ + FieldSchema { + name: "namespace", + ty: TypeSchema::String, + comment: "Namespace to query.", + required: true, + }, + FieldSchema { + name: "query", + ty: TypeSchema::String, + comment: "Natural-language query string.", + required: true, + }, + FieldSchema { + name: "include_references", + ty: TypeSchema::Option(Box::new(TypeSchema::Bool)), + comment: "Whether to include entity/relation context in the response.", + required: false, + }, + FieldSchema { + name: "document_ids", + ty: TypeSchema::Option(Box::new(TypeSchema::Array(Box::new( + TypeSchema::String, + )))), + comment: "Restrict results to these document IDs.", + required: false, + }, + FieldSchema { + name: "limit", + ty: TypeSchema::Option(Box::new(TypeSchema::U64)), + comment: "Maximum number of results to return.", + required: false, + }, + FieldSchema { + name: "max_chunks", + ty: TypeSchema::Option(Box::new(TypeSchema::U64)), + comment: "Maximum number of chunks to return (alias for limit).", + required: false, + }, + ], + outputs: vec![FieldSchema { + name: "result", + ty: TypeSchema::Json, + comment: "Envelope with retrieval context and LLM context message.", + required: true, + }], + }, + "recall_context" => ControllerSchema { + namespace: "memory", + function: "recall_context", + description: "Recall contextual data from a namespace without a specific query.", + inputs: vec![ + FieldSchema { + name: "namespace", + ty: TypeSchema::String, + comment: "Namespace to recall from.", + required: true, + }, + FieldSchema { + name: "include_references", + ty: TypeSchema::Option(Box::new(TypeSchema::Bool)), + comment: "Whether to include entity/relation context.", + required: false, + }, + FieldSchema { + name: "limit", + ty: TypeSchema::Option(Box::new(TypeSchema::U64)), + comment: "Maximum number of results.", + required: false, + }, + FieldSchema { + name: "max_chunks", + ty: TypeSchema::Option(Box::new(TypeSchema::U64)), + comment: "Maximum number of chunks (alias for limit).", + required: false, + }, + ], + outputs: vec![FieldSchema { + name: "result", + ty: TypeSchema::Json, + comment: "Envelope with retrieval context and LLM context message.", + required: true, + }], + }, + "recall_memories" => ControllerSchema { + namespace: "memory", + function: "recall_memories", + description: "Recall memory items from a namespace with optional retention filtering.", + inputs: vec![ + FieldSchema { + name: "namespace", + ty: TypeSchema::String, + comment: "Namespace to recall memories from.", + required: true, + }, + FieldSchema { + name: "min_retention", + ty: TypeSchema::Option(Box::new(TypeSchema::F64)), + comment: "Minimum retention score (forward-compat, currently ignored).", + required: false, + }, + FieldSchema { + name: "as_of", + ty: TypeSchema::Option(Box::new(TypeSchema::F64)), + comment: "Temporal recall timestamp (forward-compat, currently ignored).", + required: false, + }, + FieldSchema { + name: "limit", + ty: TypeSchema::Option(Box::new(TypeSchema::U64)), + comment: "Maximum number of results.", + required: false, + }, + FieldSchema { + name: "max_chunks", + ty: TypeSchema::Option(Box::new(TypeSchema::U64)), + comment: "Maximum number of chunks (alias for limit).", + required: false, + }, + FieldSchema { + name: "top_k", + ty: TypeSchema::Option(Box::new(TypeSchema::U64)), + comment: "Top-k override (alias for limit).", + required: false, + }, + ], + outputs: vec![FieldSchema { + name: "result", + ty: TypeSchema::Json, + comment: "Envelope with recalled memory items.", + required: true, + }], + }, + + // ----- file-based memory methods ----- + "list_files" => ControllerSchema { + namespace: "memory", + function: "list_files", + description: "List files in a memory directory.", + inputs: vec![FieldSchema { + name: "relative_dir", + ty: TypeSchema::String, + comment: "Relative directory path under the workspace (default: \"memory\").", + required: false, + }], + outputs: vec![FieldSchema { + name: "result", + ty: TypeSchema::Json, + comment: "Envelope with file listing.", + required: true, + }], + }, + "read_file" => ControllerSchema { + namespace: "memory", + function: "read_file", + description: "Read the contents of a memory file.", + inputs: vec![FieldSchema { + name: "relative_path", + ty: TypeSchema::String, + comment: "Relative path to the file under the workspace memory directory.", + required: true, + }], + outputs: vec![FieldSchema { + name: "result", + ty: TypeSchema::Json, + comment: "Envelope with file content.", + required: true, + }], + }, + "write_file" => ControllerSchema { + namespace: "memory", + function: "write_file", + description: "Write content to a memory file.", + inputs: vec![ + FieldSchema { + name: "relative_path", + ty: TypeSchema::String, + comment: "Relative path to the file under the workspace memory directory.", + required: true, + }, + FieldSchema { + name: "content", + ty: TypeSchema::String, + comment: "Content to write to the file.", + required: true, + }, + ], + outputs: vec![FieldSchema { + name: "result", + ty: TypeSchema::Json, + comment: "Envelope with write confirmation and bytes written.", + required: true, + }], + }, + + // ----- unified memory API methods ----- + "namespace_list" => ControllerSchema { + namespace: "memory", + function: "namespace_list", + description: "List all namespaces in the unified memory store.", + inputs: vec![], + outputs: vec![FieldSchema { + name: "namespaces", + ty: TypeSchema::Array(Box::new(TypeSchema::String)), + comment: "Namespace names.", + required: true, + }], + }, + "doc_put" => ControllerSchema { + namespace: "memory", + function: "doc_put", + description: "Upsert a document into a namespace.", + inputs: vec![ + FieldSchema { + name: "namespace", + ty: TypeSchema::String, + comment: "Target namespace.", + required: true, + }, + FieldSchema { + name: "key", + ty: TypeSchema::String, + comment: "Document key for upsert deduplication.", + required: true, + }, + FieldSchema { + name: "title", + ty: TypeSchema::String, + comment: "Human-readable title.", + required: true, + }, + FieldSchema { + name: "content", + ty: TypeSchema::String, + comment: "Document body content.", + required: true, + }, + FieldSchema { + name: "source_type", + ty: TypeSchema::String, + comment: "Source type label (default: \"doc\").", + required: false, + }, + FieldSchema { + name: "priority", + ty: TypeSchema::String, + comment: "Priority level (default: \"medium\").", + required: false, + }, + FieldSchema { + name: "tags", + ty: TypeSchema::Array(Box::new(TypeSchema::String)), + comment: "Tags for categorisation (default: []).", + required: false, + }, + FieldSchema { + name: "metadata", + ty: TypeSchema::Json, + comment: "Arbitrary metadata (default: {}).", + required: false, + }, + FieldSchema { + name: "category", + ty: TypeSchema::String, + comment: "Memory category (default: \"core\").", + required: false, + }, + FieldSchema { + name: "session_id", + ty: TypeSchema::Option(Box::new(TypeSchema::String)), + comment: "Optional session ID for provenance tracking.", + required: false, + }, + FieldSchema { + name: "document_id", + ty: TypeSchema::Option(Box::new(TypeSchema::String)), + comment: "Optional explicit document ID; generated if omitted.", + required: false, + }, + ], + outputs: vec![FieldSchema { + name: "document_id", + ty: TypeSchema::String, + comment: "ID of the upserted document.", + required: true, + }], + }, + "doc_ingest" => ControllerSchema { + namespace: "memory", + function: "doc_ingest", + description: "Ingest a document with entity/relation extraction and chunk embedding.", + inputs: vec![ + FieldSchema { + name: "namespace", + ty: TypeSchema::String, + comment: "Target namespace.", + required: true, + }, + FieldSchema { + name: "key", + ty: TypeSchema::String, + comment: "Document key.", + required: true, + }, + FieldSchema { + name: "title", + ty: TypeSchema::String, + comment: "Human-readable title.", + required: true, + }, + FieldSchema { + name: "content", + ty: TypeSchema::String, + comment: "Document body content.", + required: true, + }, + FieldSchema { + name: "source_type", + ty: TypeSchema::String, + comment: "Source type label (default: \"doc\").", + required: false, + }, + FieldSchema { + name: "priority", + ty: TypeSchema::String, + comment: "Priority level (default: \"medium\").", + required: false, + }, + FieldSchema { + name: "tags", + ty: TypeSchema::Array(Box::new(TypeSchema::String)), + comment: "Tags for categorisation (default: []).", + required: false, + }, + FieldSchema { + name: "metadata", + ty: TypeSchema::Json, + comment: "Arbitrary metadata (default: {}).", + required: false, + }, + FieldSchema { + name: "category", + ty: TypeSchema::String, + comment: "Memory category (default: \"core\").", + required: false, + }, + FieldSchema { + name: "session_id", + ty: TypeSchema::Option(Box::new(TypeSchema::String)), + comment: "Optional session ID.", + required: false, + }, + FieldSchema { + name: "document_id", + ty: TypeSchema::Option(Box::new(TypeSchema::String)), + comment: "Optional explicit document ID.", + required: false, + }, + FieldSchema { + name: "config", + ty: TypeSchema::Option(Box::new(TypeSchema::Json)), + comment: "Optional ingestion configuration overrides.", + required: false, + }, + ], + outputs: vec![FieldSchema { + name: "result", + ty: TypeSchema::Json, + comment: "Ingestion result with entity, relation and chunk counts.", + required: true, + }], + }, + "doc_list" => ControllerSchema { + namespace: "memory", + function: "doc_list", + description: "List documents in the unified memory store, optionally by namespace.", + inputs: vec![FieldSchema { + name: "namespace", + ty: TypeSchema::Option(Box::new(TypeSchema::String)), + comment: "Optional namespace filter.", + required: false, + }], + outputs: vec![FieldSchema { + name: "result", + ty: TypeSchema::Json, + comment: "Document listing.", + required: true, + }], + }, + "doc_delete" => ControllerSchema { + namespace: "memory", + function: "doc_delete", + description: "Delete a document from the unified memory store.", + inputs: vec![ + FieldSchema { + name: "namespace", + ty: TypeSchema::String, + comment: "Namespace containing the document.", + required: true, + }, + FieldSchema { + name: "document_id", + ty: TypeSchema::String, + comment: "Document identifier to delete.", + required: true, + }, + ], + outputs: vec![FieldSchema { + name: "result", + ty: TypeSchema::Json, + comment: "Deletion result.", + required: true, + }], + }, + "context_query" => ControllerSchema { + namespace: "memory", + function: "context_query", + description: "Query a namespace for contextual information.", + inputs: vec![ + FieldSchema { + name: "namespace", + ty: TypeSchema::String, + comment: "Namespace to query.", + required: true, + }, + FieldSchema { + name: "query", + ty: TypeSchema::String, + comment: "Natural-language query string.", + required: true, + }, + FieldSchema { + name: "limit", + ty: TypeSchema::Option(Box::new(TypeSchema::U64)), + comment: "Maximum number of results.", + required: false, + }, + ], + outputs: vec![FieldSchema { + name: "result", + ty: TypeSchema::String, + comment: "Contextual query result string.", + required: true, + }], + }, + "context_recall" => ControllerSchema { + namespace: "memory", + function: "context_recall", + description: "Recall context from a namespace.", + inputs: vec![ + FieldSchema { + name: "namespace", + ty: TypeSchema::String, + comment: "Namespace to recall from.", + required: true, + }, + FieldSchema { + name: "limit", + ty: TypeSchema::Option(Box::new(TypeSchema::U64)), + comment: "Maximum number of results.", + required: false, + }, + ], + outputs: vec![FieldSchema { + name: "result", + ty: TypeSchema::Json, + comment: "Recalled context (may be null if empty).", + required: true, + }], + }, + + // ----- key-value methods ----- + "kv_set" => ControllerSchema { + namespace: "memory", + function: "kv_set", + description: "Set a key-value pair in the memory store.", + inputs: vec![ + FieldSchema { + name: "namespace", + ty: TypeSchema::Option(Box::new(TypeSchema::String)), + comment: "Optional namespace scope.", + required: false, + }, + FieldSchema { + name: "key", + ty: TypeSchema::String, + comment: "Key to set.", + required: true, + }, + FieldSchema { + name: "value", + ty: TypeSchema::Json, + comment: "JSON value to store.", + required: true, + }, + ], + outputs: vec![FieldSchema { + name: "result", + ty: TypeSchema::Bool, + comment: "True when the value was stored.", + required: true, + }], + }, + "kv_get" => ControllerSchema { + namespace: "memory", + function: "kv_get", + description: "Get a value by key from the memory store.", + inputs: vec![ + FieldSchema { + name: "namespace", + ty: TypeSchema::Option(Box::new(TypeSchema::String)), + comment: "Optional namespace scope.", + required: false, + }, + FieldSchema { + name: "key", + ty: TypeSchema::String, + comment: "Key to retrieve.", + required: true, + }, + ], + outputs: vec![FieldSchema { + name: "result", + ty: TypeSchema::Json, + comment: "Stored value or null if not found.", + required: true, + }], + }, + "kv_delete" => ControllerSchema { + namespace: "memory", + function: "kv_delete", + description: "Delete a key-value pair from the memory store.", + inputs: vec![ + FieldSchema { + name: "namespace", + ty: TypeSchema::Option(Box::new(TypeSchema::String)), + comment: "Optional namespace scope.", + required: false, + }, + FieldSchema { + name: "key", + ty: TypeSchema::String, + comment: "Key to delete.", + required: true, + }, + ], + outputs: vec![FieldSchema { + name: "result", + ty: TypeSchema::Bool, + comment: "True when the key was deleted.", + required: true, + }], + }, + "kv_list_namespace" => ControllerSchema { + namespace: "memory", + function: "kv_list_namespace", + description: "List all key-value entries in a namespace.", + inputs: vec![FieldSchema { + name: "namespace", + ty: TypeSchema::String, + comment: "Namespace to list.", + required: true, + }], + outputs: vec![FieldSchema { + name: "result", + ty: TypeSchema::Json, + comment: "Array of key-value entries.", + required: true, + }], + }, + + // ----- graph methods ----- + "graph_upsert" => ControllerSchema { + namespace: "memory", + function: "graph_upsert", + description: "Upsert a relation triple in the knowledge graph.", + inputs: vec![ + FieldSchema { + name: "namespace", + ty: TypeSchema::Option(Box::new(TypeSchema::String)), + comment: "Optional namespace scope.", + required: false, + }, + FieldSchema { + name: "subject", + ty: TypeSchema::String, + comment: "Subject entity of the relation.", + required: true, + }, + FieldSchema { + name: "predicate", + ty: TypeSchema::String, + comment: "Relation predicate.", + required: true, + }, + FieldSchema { + name: "object", + ty: TypeSchema::String, + comment: "Object entity of the relation.", + required: true, + }, + FieldSchema { + name: "attrs", + ty: TypeSchema::Json, + comment: "Extra attributes on the relation (default: {}).", + required: false, + }, + ], + outputs: vec![FieldSchema { + name: "result", + ty: TypeSchema::Bool, + comment: "True when the relation was upserted.", + required: true, + }], + }, + "graph_query" => ControllerSchema { + namespace: "memory", + function: "graph_query", + description: "Query relations from the knowledge graph.", + inputs: vec![ + FieldSchema { + name: "namespace", + ty: TypeSchema::Option(Box::new(TypeSchema::String)), + comment: "Optional namespace scope.", + required: false, + }, + FieldSchema { + name: "subject", + ty: TypeSchema::Option(Box::new(TypeSchema::String)), + comment: "Filter by subject entity.", + required: false, + }, + FieldSchema { + name: "predicate", + ty: TypeSchema::Option(Box::new(TypeSchema::String)), + comment: "Filter by relation predicate.", + required: false, + }, + ], + outputs: vec![FieldSchema { + name: "result", + ty: TypeSchema::Json, + comment: "Array of matching relation records.", + required: true, + }], + }, + + // ----- fallback ----- + _other => ControllerSchema { + namespace: "memory", + function: "unknown", + description: "Unknown memory controller function.", + inputs: vec![FieldSchema { + name: "function", + ty: TypeSchema::String, + comment: "Unknown function requested for schema lookup.", + required: true, + }], + outputs: vec![FieldSchema { + name: "error", + ty: TypeSchema::String, + comment: "Lookup error details.", + required: true, + }], + }, + } +} + +// --------------------------------------------------------------------------- +// Handlers +// --------------------------------------------------------------------------- + +fn handle_init(params: Map) -> ControllerFuture { + Box::pin(async move { + let payload = parse_params::(params)?; + to_json(rpc::memory_init(payload).await?) + }) +} + +fn handle_list_documents(params: Map) -> ControllerFuture { + Box::pin(async move { + let payload = parse_params::(params)?; + to_json(rpc::memory_list_documents(payload).await?) + }) +} + +fn handle_list_namespaces(_params: Map) -> ControllerFuture { + Box::pin(async move { to_json(rpc::memory_list_namespaces(EmptyRequest {}).await?) }) +} + +fn handle_delete_document(params: Map) -> ControllerFuture { + Box::pin(async move { + let payload = parse_params::(params)?; + to_json(rpc::memory_delete_document(payload).await?) + }) +} + +fn handle_query_namespace(params: Map) -> ControllerFuture { + Box::pin(async move { + let payload = parse_params::(params)?; + to_json(rpc::memory_query_namespace(payload).await?) + }) +} + +fn handle_recall_context(params: Map) -> ControllerFuture { + Box::pin(async move { + let payload = parse_params::(params)?; + to_json(rpc::memory_recall_context(payload).await?) + }) +} + +fn handle_recall_memories(params: Map) -> ControllerFuture { + Box::pin(async move { + let payload = parse_params::(params)?; + to_json(rpc::memory_recall_memories(payload).await?) + }) +} + +fn handle_list_files(params: Map) -> ControllerFuture { + Box::pin(async move { + let relative_dir = params + .get("relative_dir") + .and_then(|v| v.as_str()) + .unwrap_or("memory") + .to_string(); + let payload = ListMemoryFilesRequest { relative_dir }; + to_json(rpc::ai_list_memory_files(payload).await?) + }) +} + +fn handle_read_file(params: Map) -> ControllerFuture { + Box::pin(async move { + let payload = parse_params::(params)?; + to_json(rpc::ai_read_memory_file(payload).await?) + }) +} + +fn handle_write_file(params: Map) -> ControllerFuture { + Box::pin(async move { + let payload = parse_params::(params)?; + to_json(rpc::ai_write_memory_file(payload).await?) + }) +} + +fn handle_namespace_list(_params: Map) -> ControllerFuture { + Box::pin(async move { to_json(rpc::namespace_list().await?) }) +} + +fn handle_doc_put(params: Map) -> ControllerFuture { + Box::pin(async move { + let payload = parse_params::(params)?; + to_json(rpc::doc_put(payload).await?) + }) +} + +fn handle_doc_ingest(params: Map) -> ControllerFuture { + Box::pin(async move { + let payload = parse_params::(params)?; + to_json(rpc::doc_ingest(payload).await?) + }) +} + +fn handle_doc_list(params: Map) -> ControllerFuture { + Box::pin(async move { + let namespace = + params + .get("namespace") + .and_then(|v| v.as_str()) + .map(|ns| NamespaceOnlyParams { + namespace: ns.to_string(), + }); + to_json(rpc::doc_list(namespace).await?) + }) +} + +fn handle_doc_delete(params: Map) -> ControllerFuture { + Box::pin(async move { + let payload = parse_params::(params)?; + to_json(rpc::doc_delete(payload).await?) + }) +} + +fn handle_context_query(params: Map) -> ControllerFuture { + Box::pin(async move { + let payload = parse_params::(params)?; + to_json(rpc::context_query(payload).await?) + }) +} + +fn handle_context_recall(params: Map) -> ControllerFuture { + Box::pin(async move { + let payload = parse_params::(params)?; + to_json(rpc::context_recall(payload).await?) + }) +} + +fn handle_kv_set(params: Map) -> ControllerFuture { + Box::pin(async move { + let payload = parse_params::(params)?; + to_json(rpc::kv_set(payload).await?) + }) +} + +fn handle_kv_get(params: Map) -> ControllerFuture { + Box::pin(async move { + let payload = parse_params::(params)?; + to_json(rpc::kv_get(payload).await?) + }) +} + +fn handle_kv_delete(params: Map) -> ControllerFuture { + Box::pin(async move { + let payload = parse_params::(params)?; + to_json(rpc::kv_delete(payload).await?) + }) +} + +fn handle_kv_list_namespace(params: Map) -> ControllerFuture { + Box::pin(async move { + let payload = parse_params::(params)?; + to_json(rpc::kv_list_namespace(payload).await?) + }) +} + +fn handle_graph_upsert(params: Map) -> ControllerFuture { + Box::pin(async move { + let payload = parse_params::(params)?; + to_json(rpc::graph_upsert(payload).await?) + }) +} + +fn handle_graph_query(params: Map) -> ControllerFuture { + Box::pin(async move { + let payload = parse_params::(params)?; + to_json(rpc::graph_query(payload).await?) + }) +} + +// --------------------------------------------------------------------------- +// Helpers +// --------------------------------------------------------------------------- + +fn parse_params(params: Map) -> Result { + serde_json::from_value(Value::Object(params)).map_err(|e| format!("invalid params: {e}")) +} + +fn to_json(outcome: RpcOutcome) -> Result { + outcome.into_cli_compatible_json() +} diff --git a/src/rpc/dispatch.rs b/src/rpc/dispatch.rs index dd3a3216c..91a431e10 100644 --- a/src/rpc/dispatch.rs +++ b/src/rpc/dispatch.rs @@ -1,222 +1,16 @@ -use serde::de::DeserializeOwned; -use serde::{Deserialize, Serialize}; +use serde::Serialize; use crate::rpc::RpcOutcome; -fn parse_params(params: serde_json::Value) -> Result { - serde_json::from_value(params).map_err(|e| format!("invalid params: {e}")) -} - fn rpc_json(outcome: RpcOutcome) -> Result { outcome.into_cli_compatible_json() } -#[derive(Debug, Deserialize)] -struct MemoryDocListParams { - namespace: Option, -} - pub async fn try_dispatch( method: &str, - params: serde_json::Value, + _params: serde_json::Value, ) -> Option> { match method { - "memory.init" => Some( - async move { - let payload: crate::openhuman::memory::MemoryInitRequest = parse_params(params)?; - rpc_json(crate::openhuman::memory::rpc::memory_init(payload).await?) - } - .await, - ), - - "memory.list_documents" => Some( - async move { - let payload: crate::openhuman::memory::ListDocumentsRequest = parse_params(params)?; - rpc_json(crate::openhuman::memory::rpc::memory_list_documents(payload).await?) - } - .await, - ), - - "memory.list_namespaces" => Some( - async move { - let payload: crate::openhuman::memory::EmptyRequest = parse_params(params)?; - rpc_json(crate::openhuman::memory::rpc::memory_list_namespaces(payload).await?) - } - .await, - ), - - "memory.delete_document" => Some( - async move { - let payload: crate::openhuman::memory::DeleteDocumentRequest = - parse_params(params)?; - rpc_json(crate::openhuman::memory::rpc::memory_delete_document(payload).await?) - } - .await, - ), - - "memory.query_namespace" => Some( - async move { - let payload: crate::openhuman::memory::QueryNamespaceRequest = - parse_params(params)?; - rpc_json(crate::openhuman::memory::rpc::memory_query_namespace(payload).await?) - } - .await, - ), - - "memory.recall_context" | "memory.recall_namespace" => Some( - async move { - let payload: crate::openhuman::memory::RecallContextRequest = parse_params(params)?; - rpc_json(crate::openhuman::memory::rpc::memory_recall_context(payload).await?) - } - .await, - ), - - "memory.recall_memories" => Some( - async move { - let payload: crate::openhuman::memory::RecallMemoriesRequest = - parse_params(params)?; - rpc_json(crate::openhuman::memory::rpc::memory_recall_memories(payload).await?) - } - .await, - ), - - "ai.list_memory_files" => Some( - async move { - let payload: crate::openhuman::memory::ListMemoryFilesRequest = - parse_params(params)?; - rpc_json(crate::openhuman::memory::rpc::ai_list_memory_files(payload).await?) - } - .await, - ), - - "ai.read_memory_file" => Some( - async move { - let payload: crate::openhuman::memory::ReadMemoryFileRequest = - parse_params(params)?; - rpc_json(crate::openhuman::memory::rpc::ai_read_memory_file(payload).await?) - } - .await, - ), - - "ai.write_memory_file" => Some( - async move { - let payload: crate::openhuman::memory::WriteMemoryFileRequest = - parse_params(params)?; - rpc_json(crate::openhuman::memory::rpc::ai_write_memory_file(payload).await?) - } - .await, - ), - - "memory.namespace.list" => Some( - async move { rpc_json(crate::openhuman::memory::rpc::namespace_list().await?) }.await, - ), - - "memory.doc.put" => Some( - async move { - let payload: crate::openhuman::memory::rpc::PutDocParams = parse_params(params)?; - rpc_json(crate::openhuman::memory::rpc::doc_put(payload).await?) - } - .await, - ), - - "memory.doc.ingest" => Some( - async move { - let payload: crate::openhuman::memory::rpc::IngestDocParams = parse_params(params)?; - rpc_json(crate::openhuman::memory::rpc::doc_ingest(payload).await?) - } - .await, - ), - - "memory.doc.list" => Some( - async move { - let payload: MemoryDocListParams = parse_params(params)?; - let namespace_params = payload.namespace.map(|namespace| { - crate::openhuman::memory::rpc::NamespaceOnlyParams { namespace } - }); - rpc_json(crate::openhuman::memory::rpc::doc_list(namespace_params).await?) - } - .await, - ), - - "memory.doc.delete" => Some( - async move { - let payload: crate::openhuman::memory::rpc::DeleteDocParams = parse_params(params)?; - rpc_json(crate::openhuman::memory::rpc::doc_delete(payload).await?) - } - .await, - ), - - "memory.context.query" => Some( - async move { - let payload: crate::openhuman::memory::rpc::QueryNamespaceParams = - parse_params(params)?; - rpc_json(crate::openhuman::memory::rpc::context_query(payload).await?) - } - .await, - ), - - "memory.context.recall" => Some( - async move { - let payload: crate::openhuman::memory::rpc::RecallNamespaceParams = - parse_params(params)?; - rpc_json(crate::openhuman::memory::rpc::context_recall(payload).await?) - } - .await, - ), - - "memory.kv.set" => Some( - async move { - let payload: crate::openhuman::memory::rpc::KvSetParams = parse_params(params)?; - rpc_json(crate::openhuman::memory::rpc::kv_set(payload).await?) - } - .await, - ), - - "memory.kv.get" => Some( - async move { - let payload: crate::openhuman::memory::rpc::KvGetDeleteParams = - parse_params(params)?; - rpc_json(crate::openhuman::memory::rpc::kv_get(payload).await?) - } - .await, - ), - - "memory.kv.delete" => Some( - async move { - let payload: crate::openhuman::memory::rpc::KvGetDeleteParams = - parse_params(params)?; - rpc_json(crate::openhuman::memory::rpc::kv_delete(payload).await?) - } - .await, - ), - - "memory.kv.list_namespace" => Some( - async move { - let payload: crate::openhuman::memory::rpc::NamespaceOnlyParams = - parse_params(params)?; - rpc_json(crate::openhuman::memory::rpc::kv_list_namespace(payload).await?) - } - .await, - ), - - "memory.graph.upsert" => Some( - async move { - let payload: crate::openhuman::memory::rpc::GraphUpsertParams = - parse_params(params)?; - rpc_json(crate::openhuman::memory::rpc::graph_upsert(payload).await?) - } - .await, - ), - - "memory.graph.query" => Some( - async move { - let payload: crate::openhuman::memory::rpc::GraphQueryParams = - parse_params(params)?; - rpc_json(crate::openhuman::memory::rpc::graph_query(payload).await?) - } - .await, - ), - "openhuman.security_policy_info" => Some(rpc_json( crate::openhuman::security::rpc::security_policy_info(), )), @@ -231,57 +25,10 @@ mod tests { use super::try_dispatch; - /// Verify that the dispatcher recognises `memory.doc.ingest` (returns `Some`). - /// The inner handler will fail because no memory client is initialised, - /// but the route being present (not `None`) is what we need to assert. - #[tokio::test] - async fn dispatch_routes_memory_doc_ingest() { - let params = json!({ - "namespace": "test", - "key": "k1", - "title": "Title", - "content": "body" - }); - let result = try_dispatch("memory.doc.ingest", params).await; - assert!( - result.is_some(), - "memory.doc.ingest should be routed by dispatch" - ); - } - - /// Verify that `memory.graph.query` is routed. - #[tokio::test] - async fn dispatch_routes_memory_graph_query() { - let params = json!({}); - let result = try_dispatch("memory.graph.query", params).await; - assert!( - result.is_some(), - "memory.graph.query should be routed by dispatch" - ); - } - /// Unknown methods must return `None` so callers can fall through. #[tokio::test] async fn dispatch_returns_none_for_unknown_method() { let result = try_dispatch("nonexistent.method", json!({})).await; assert!(result.is_none(), "unknown methods should return None"); } - - /// Verify that params deserialization errors surface as `Some(Err(...))`. - #[tokio::test] - async fn dispatch_memory_doc_ingest_rejects_invalid_params() { - // Missing required fields → should be Some(Err) - let result = try_dispatch("memory.doc.ingest", json!({})).await; - assert!(result.is_some()); - let inner = result.unwrap(); - assert!( - inner.is_err(), - "missing required fields should produce a deserialization error" - ); - let err_msg = inner.unwrap_err(); - assert!( - err_msg.contains("invalid params"), - "error should mention invalid params, got: {err_msg}" - ); - } }