mirror of
https://github.com/tinyhumansai/openhuman.git
synced 2026-07-27 21:08:00 +00:00
refactor: migrate memory service to controller registry pattern (#138)
Move all 23 memory RPC methods from legacy dispatch (src/rpc/dispatch.rs) to the controller registry pattern with typed schemas. - Create src/openhuman/memory/schemas.rs with 23 ControllerSchema definitions, RegisteredController entries, and handler functions - Wire memory controllers into src/core/all.rs registry builders - Remove all memory.* and ai.* branches from dispatch.rs (only security_policy_info remains) - Update frontend to use openhuman.memory_* method names directly in tauriCommands.ts (no legacy aliases needed) - Move ai.list_memory_files/read/write into memory namespace as openhuman.memory_list_files/read_file/write_file - Update jsonrpc.rs and tauriCommandsMemory test method strings Methods are now accessible via both JSON-RPC (openhuman.memory_*) and CLI (openhuman memory <function>).
This commit is contained in:
@@ -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 });
|
||||
});
|
||||
});
|
||||
|
||||
@@ -199,7 +199,7 @@ export async function syncMemoryClientToken(token: string): Promise<void> {
|
||||
}
|
||||
try {
|
||||
console.debug('[memory] syncMemoryClientToken: payload → memory.init');
|
||||
await callCoreRpc<boolean>({ method: 'memory.init', params: { jwt_token: token } });
|
||||
await callCoreRpc<boolean>({ 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<unknown>
|
||||
if (!isTauri()) {
|
||||
throw new Error('Not running in Tauri');
|
||||
}
|
||||
return await callCoreRpc<unknown>({ method: 'memory.list_documents', params: { namespace } });
|
||||
return await callCoreRpc<unknown>({
|
||||
method: 'openhuman.memory_list_documents',
|
||||
params: { namespace },
|
||||
});
|
||||
}
|
||||
|
||||
export async function memoryListNamespaces(): Promise<string[]> {
|
||||
if (!isTauri()) {
|
||||
throw new Error('Not running in Tauri');
|
||||
}
|
||||
return await callCoreRpc<string[]>({ method: 'memory.list_namespaces' });
|
||||
return await callCoreRpc<string[]>({ 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<unknown>({
|
||||
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<string>({
|
||||
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<string | null>({
|
||||
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<GraphRelation[]>({
|
||||
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<unknown>({ method: 'memory.doc.ingest', params });
|
||||
return await callCoreRpc<unknown>({ method: 'openhuman.memory_doc_ingest', params });
|
||||
}
|
||||
|
||||
export async function aiListMemoryFiles(relativeDir = 'memory'): Promise<string[]> {
|
||||
@@ -318,7 +321,7 @@ export async function aiListMemoryFiles(relativeDir = 'memory'): Promise<string[
|
||||
throw new Error('Not running in Tauri');
|
||||
}
|
||||
return await callCoreRpc<string[]>({
|
||||
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<string> {
|
||||
throw new Error('Not running in Tauri');
|
||||
}
|
||||
return await callCoreRpc<string>({
|
||||
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<boolean>({
|
||||
method: 'ai.write_memory_file',
|
||||
method: 'openhuman.memory_write_file',
|
||||
params: { relative_path: relativePath, content },
|
||||
});
|
||||
}
|
||||
|
||||
@@ -62,6 +62,7 @@ fn build_registered_controllers() -> Vec<RegisteredController> {
|
||||
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<ControllerSchema> {
|
||||
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,
|
||||
}
|
||||
}
|
||||
|
||||
+8
-8
@@ -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]
|
||||
|
||||
@@ -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,
|
||||
|
||||
File diff suppressed because it is too large
Load Diff
+2
-255
@@ -1,222 +1,16 @@
|
||||
use serde::de::DeserializeOwned;
|
||||
use serde::{Deserialize, Serialize};
|
||||
use serde::Serialize;
|
||||
|
||||
use crate::rpc::RpcOutcome;
|
||||
|
||||
fn parse_params<T: DeserializeOwned>(params: serde_json::Value) -> Result<T, String> {
|
||||
serde_json::from_value(params).map_err(|e| format!("invalid params: {e}"))
|
||||
}
|
||||
|
||||
fn rpc_json<T: Serialize>(outcome: RpcOutcome<T>) -> Result<serde_json::Value, String> {
|
||||
outcome.into_cli_compatible_json()
|
||||
}
|
||||
|
||||
#[derive(Debug, Deserialize)]
|
||||
struct MemoryDocListParams {
|
||||
namespace: Option<String>,
|
||||
}
|
||||
|
||||
pub async fn try_dispatch(
|
||||
method: &str,
|
||||
params: serde_json::Value,
|
||||
_params: serde_json::Value,
|
||||
) -> Option<Result<serde_json::Value, String>> {
|
||||
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}"
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user