Skip to main content

p_memory/
protocol.rs

1//! Versioned JSON boundary for thin foreign-language bindings. Domain behavior
2//! lives in the typed Rust stores; this module only decodes and dispatches.
3use crate::{KnowledgeBase, Result, Error, types::*, memory::DecayPolicy};
4use serde::{de::DeserializeOwned, Deserialize, Serialize};
5use serde_json::{json, Value};
6
7pub const PROTOCOL_VERSION:u32=1;
8fn decode<T:DeserializeOwned>(value:Value)->Result<T>{Ok(serde_json::from_value(value)?)}
9fn encode<T:Serialize>(value:T)->Result<Value>{Ok(serde_json::to_value(value)?)}
10fn field<T:DeserializeOwned>(value:&Value,name:&str)->Result<T>{decode(value.get(name).cloned().ok_or_else(||Error::Validation(format!("missing {name}")))?)}
11fn optional<T:DeserializeOwned+Default>(value:&Value,name:&str)->Result<T>{value.get(name).cloned().map(decode).transpose().map(|v|v.unwrap_or_default())}
12#[derive(Deserialize)]
13struct IdRequest { id:i64, #[serde(default)] filter:ReadFilter }
14#[derive(Deserialize)]
15struct GraphId { id:i64, kind:RecordKind, #[serde(default)] filter:ReadFilter }
16
17pub fn dispatch(kb:&KnowledgeBase,operation:&str,args:Value)->Result<Value>{
18    match operation {
19        "health"=>encode(kb.health()?),
20        "rebuild_indexes"=>encode(kb.rebuild_indexes()?),
21        "rebuild_progress"=>encode(kb.rebuild_progress()?),
22        "update_index"=>encode(kb.update_index()?),
23        "backup"=>{kb.backup(field::<String>(&args,"path")?)?;Ok(Value::Null)},
24        "close"=>{kb.close()?;Ok(Value::Null)},
25        "search"=>encode(kb.search(&decode(args)?)?),
26        "preset"=>encode(kb.search_preset(&decode(args)?)?),
27        "memories.upsert"=>encode(kb.memories().upsert(decode(args)?)?),
28        "memories.upsert_many"=>encode(kb.memories().upsert_many(&decode::<Vec<_>>(args)?)?),
29        "memories.upsert_by_judgment"=>encode(kb.memories().upsert_by_judgment(decode(args)?)?),
30        "memories.get"=>{let r:IdRequest=decode(args)?;encode(kb.memories().get(r.id,&r.filter)?)},
31        "memories.list"=>encode(kb.memories().list(&decode(args)?)?),
32        "memories.delete"=>{let r:IdRequest=decode(args)?;encode(kb.memories().delete(r.id,&r.filter)?)},
33        "memories.feedback"=>encode(kb.memories().feedback(&decode(args)?)?),
34        "memories.decay"=>encode(kb.memories().decay(&optional::<ReadFilter>(&args,"filter")?,&optional::<DecayPolicy>(&args,"policy")?,optional(&args,"now_us")?)?),
35        "graph.apply_batch"=>encode(kb.graph().apply_batch(&decode(args)?)?),
36        "graph.get"=>{let r:GraphId=decode(args)?;encode(kb.graph().get(r.kind,r.id,&r.filter)?)},
37        "graph.list"=>encode(kb.graph().list(field(&args,"kind")?,&optional::<PageRequest>(&args,"page")?)?),
38        "graph.delete"=>{let r:GraphId=decode(args)?;encode(kb.graph().delete(r.kind,r.id,&r.filter)?)},
39        "graph.resolve"=>encode(kb.graph().resolve(&field::<String>(&args,"name")?,&optional::<ReadFilter>(&args,"filter")?,args.get("limit").cloned().map(decode).transpose()?.unwrap_or(10))?),
40        "graph.set_predicate_equivalents"=>encode(kb.graph().set_predicate_equivalents(&field::<String>(&args,"namespace")?,&field::<Vec<Vec<String>>>(&args,"groups")?)?),
41        "graph.predicate_equivalents"=>encode(kb.graph().predicate_equivalents(&field::<String>(&args,"namespace")?)?),
42        "graph.expand_query"=>encode(kb.graph().expand_query(&field::<String>(&args,"namespace")?,&field::<String>(&args,"text")?)?),
43        "graph.neighbors"=>encode(kb.graph().neighbors(field::<i64>(&args,"id")?,&optional::<ReadFilter>(&args,"filter")?,args.get("limit").cloned().map(decode).transpose()?.unwrap_or(50))?),
44        "graph.events_for_entity"=>encode(kb.graph().events_for_entity(field::<i64>(&args,"id")?,&optional::<ReadFilter>(&args,"filter")?,args.get("limit").cloned().map(decode).transpose()?.unwrap_or(50))?),
45        "graph.ego"=>{let id:i64=field(&args,"id")?;let depth:usize=args.get("depth").cloned().map(decode).transpose()?.unwrap_or(1);let limit:usize=args.get("limit").cloned().map(decode).transpose()?.unwrap_or(50);encode(kb.graph().ego(id,depth,&optional::<ReadFilter>(&args,"filter")?,limit)?)},
46        "graph.path"=>encode(kb.graph().path(field::<i64>(&args,"from")?,field::<i64>(&args,"to")?,&optional::<ReadFilter>(&args,"filter")?)?),
47        "graph.strongly_connected"=>encode(kb.graph().strongly_connected(&optional::<ReadFilter>(&args,"filter")?)?),
48        "graph.component_count"=>encode(kb.graph().build_graph(&optional::<ReadFilter>(&args,"filter")?)?.component_count()),
49        "notes.upsert_file"=>encode(kb.notes().upsert_file(decode(args)?)?),
50        "notes.set_root"=>{kb.notes().set_root(&field::<String>(&args,"namespace")?,&field::<String>(&args,"root")?)?;Ok(Value::Null)},
51        "notes.root"=>encode(kb.notes().root(&field::<String>(&args,"namespace")?)?),
52        "notes.get"=>{let r:IdRequest=decode(args)?;encode(kb.notes().get(r.id,&r.filter)?)},
53        "notes.list"=>encode(kb.notes().list(&decode(args)?)?),
54        "notes.delete"=>{let r:IdRequest=decode(args)?;encode(kb.notes().delete(r.id,&r.filter)?)},
55        "notes.chunks"=>{let r:IdRequest=decode(args)?;encode(kb.notes().chunks(r.id,&r.filter)?)},
56        "notes.get_chunk"=>{let r:IdRequest=decode(args)?;encode(kb.notes().get_chunk(r.id,&r.filter)?)},
57        "embeddings.register_space"=>encode(kb.embeddings().register_space(decode(args)?)?),
58        "embeddings.spaces"=>encode(kb.embeddings().spaces()?),
59        "embeddings.embedder_space"=>encode(kb.embeddings().embedder_space(&field::<String>(&args,"space_id")?)?),
60        "embeddings.sync"=>encode(kb.embeddings().sync(&field::<String>(&args,"space_id")?,field(&args,"batch")?)?),
61        "embeddings.vector_ready"=>encode(kb.embeddings().vector_ready(&field::<String>(&args,"namespace")?,&field::<String>(&args,"space_id")?,&field::<String>(&args,"target")?)?),
62        "embeddings.unregister_embedder"=>encode(kb.embeddings().unregister_embedder(&field::<String>(&args,"space_id")?)?),
63        "embeddings.namespace_vectorization"=>encode(kb.embeddings().namespace_vectorization(&field::<String>(&args,"namespace")?)?),
64        "embeddings.set_namespace_vectorization"=>encode(kb.embeddings().set_namespace_vectorization(&field::<String>(&args,"namespace")?,field(&args,"enabled")?)?),
65        "embeddings.vectorization"=>encode(kb.embeddings().vectorization(&field::<String>(&args,"namespace")?,&field::<String>(&args,"target")?)?),
66        "embeddings.set_vectorization"=>encode(kb.embeddings().set_vectorization(&field::<String>(&args,"namespace")?,&field::<String>(&args,"target")?,field(&args,"enabled")?)?),
67        "embeddings.delete_space"=>encode(kb.embeddings().delete_space(&field::<String>(&args,"id")?)?),
68        _=>Err(Error::Validation(format!("unknown operation: {operation}"))),
69    }
70}
71
72pub fn envelope(result:Result<Value>)->Value {
73    match result {Ok(value)=>json!({"ok":true,"result":value}),Err(error)=>json!({"ok":false,"error":{"code":error.code(),"message":error.to_string()}})}
74}