1use 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 IdsRequest { ids: Vec<i64>, #[serde(default)] filter:ReadFilter }
16#[derive(Deserialize)]
17struct GraphId { id:i64, kind:RecordKind, #[serde(default)] filter:ReadFilter }
18
19pub fn dispatch(kb:&KnowledgeBase,operation:&str,args:Value)->Result<Value>{
20 match operation {
21 "health"=>encode(kb.health()?),
22 "rebuild_indexes"=>encode(kb.rebuild_indexes()?),
23 "rebuild_progress"=>encode(kb.rebuild_progress()?),
24 "update_index"=>encode(kb.update_index()?),
25 "backup"=>{kb.backup(field::<String>(&args,"path")?)?;Ok(Value::Null)},
26 "close"=>{kb.close()?;Ok(Value::Null)},
27 "search"=>encode(kb.search(&decode(args)?)?),
28 "preset"=>encode(kb.search_preset(&decode(args)?)?),
29 "memories.upsert"=>encode(kb.memories().upsert(decode(args)?)?),
30 "memories.upsert_many"=>encode(kb.memories().upsert_many(&decode::<Vec<_>>(args)?)?),
31 "memories.upsert_by_judgment"=>encode(kb.memories().upsert_by_judgment(decode(args)?)?),
32 "memories.get"=>{let r:IdRequest=decode(args)?;encode(kb.memories().get(r.id,&r.filter)?)},
33 "memories.list"=>encode(kb.memories().list(&decode(args)?)?),
34 "memories.delete"=>{let r:IdRequest=decode(args)?;encode(kb.memories().delete(r.id,&r.filter)?)},
35 "memories.delete_by_filter"=>encode(kb.memories().delete_by_filter(&optional::<ReadFilter>(&args,"filter")?)?),
36 "memories.feedback"=>encode(kb.memories().feedback(&decode(args)?)?),
37 "memories.decay"=>encode(kb.memories().decay(&optional::<ReadFilter>(&args,"filter")?,&optional::<DecayPolicy>(&args,"policy")?,optional(&args,"now_us")?)?),
38 "graph.apply_batch"=>encode(kb.graph().apply_batch(&decode(args)?)?),
39 "graph.get"=>{let r:GraphId=decode(args)?;encode(kb.graph().get(r.kind,r.id,&r.filter)?)},
40 "graph.list"=>encode(kb.graph().list(field(&args,"kind")?,&optional::<PageRequest>(&args,"page")?)?),
41 "graph.delete"=>{let r:GraphId=decode(args)?;encode(kb.graph().delete(r.kind,r.id,&r.filter)?)},
42 "graph.delete_by_filter"=>encode(kb.graph().delete_by_filter(&optional::<ReadFilter>(&args,"filter")?)?),
43 "graph.resolve"=>encode(kb.graph().resolve(&field::<String>(&args,"name")?,&optional::<ReadFilter>(&args,"filter")?,args.get("limit").cloned().map(decode).transpose()?.unwrap_or(10))?),
44 "graph.set_predicate_equivalents"=>encode(kb.graph().set_predicate_equivalents(&field::<String>(&args,"namespace")?,&field::<Vec<Vec<String>>>(&args,"groups")?)?),
45 "graph.delete_predicate_equivalents"=>{let namespace=field::<String>(&args,"namespace")?;let predicates=match args.get("predicates"){None|Some(Value::Null)=>None,Some(value)=>Some(decode::<Vec<String>>(value.clone())?)};encode(kb.graph().delete_predicate_equivalents(&namespace,predicates.as_deref())?)},
46 "graph.predicate_equivalents"=>encode(kb.graph().predicate_equivalents(&field::<String>(&args,"namespace")?)?),
47 "graph.set_predicate_rule"=>{let predicate=field::<String>(&args,"predicate")?;let inverse=args.get("inverse").cloned().map(decode::<String>).transpose()?;let symmetric=args.get("symmetric").cloned().map(decode::<bool>).transpose()?.unwrap_or(false);encode(kb.graph().set_predicate_rule(&predicate,inverse.as_deref(),symmetric)?)},
48 "graph.delete_predicate_rule"=>encode(kb.graph().delete_predicate_rule(&field::<String>(&args,"predicate")?)?),
49 "graph.expand_query"=>encode(kb.graph().expand_query(&field::<String>(&args,"namespace")?,&field::<String>(&args,"text")?)?),
50 "graph.neighbors"=>encode(kb.graph().neighbors(field::<i64>(&args,"id")?,&optional::<ReadFilter>(&args,"filter")?,args.get("limit").cloned().map(decode).transpose()?.unwrap_or(50))?),
51 "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))?),
52 "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)?)},
53 "graph.path"=>encode(kb.graph().path(field::<i64>(&args,"from")?,field::<i64>(&args,"to")?,&optional::<ReadFilter>(&args,"filter")?)?),
54 "graph.strongly_connected"=>encode(kb.graph().strongly_connected(&optional::<ReadFilter>(&args,"filter")?)?),
55 "graph.component_count"=>encode(kb.graph().build_graph(&optional::<ReadFilter>(&args,"filter")?)?.component_count()),
56 "notes.upsert_file"=>encode(kb.notes().upsert_file(decode(args)?)?),
57 "notes.set_root"=>{kb.notes().set_root(&field::<String>(&args,"namespace")?,&field::<String>(&args,"root")?)?;Ok(Value::Null)},
58 "notes.unset_root"=>encode(kb.notes().unset_root(&field::<String>(&args,"namespace")?)?),
59 "notes.root"=>encode(kb.notes().root(&field::<String>(&args,"namespace")?)?),
60 "notes.get"=>{let r:IdRequest=decode(args)?;encode(kb.notes().get(r.id,&r.filter)?)},
61 "notes.get_many"=>{let r:IdsRequest=decode(args)?;encode(kb.notes().get_many(&r.ids,&r.filter)?)},
62 "notes.list"=>encode(kb.notes().list(&decode(args)?)?),
63 "notes.delete"=>{let r:IdRequest=decode(args)?;encode(kb.notes().delete(r.id,&r.filter)?)},
64 "notes.delete_by_filter"=>encode(kb.notes().delete_by_filter(&optional::<ReadFilter>(&args,"filter")?)?),
65 "notes.chunks"=>{let r:IdRequest=decode(args)?;encode(kb.notes().chunks(r.id,&r.filter)?)},
66 "notes.get_chunk"=>{let r:IdRequest=decode(args)?;encode(kb.notes().get_chunk(r.id,&r.filter)?)},
67 "embeddings.register_space"=>encode(kb.embeddings().register_space(decode(args)?)?),
68 "embeddings.spaces"=>encode(kb.embeddings().spaces()?),
69 "embeddings.embedder_space"=>encode(kb.embeddings().embedder_space(&field::<String>(&args,"space_id")?)?),
70 "embeddings.sync"=>encode(kb.embeddings().sync(&field::<String>(&args,"space_id")?,field(&args,"batch")?)?),
71 "embeddings.vector_ready"=>encode(kb.embeddings().vector_ready(&field::<String>(&args,"namespace")?,&field::<String>(&args,"space_id")?,&field::<String>(&args,"target")?)?),
72 "embeddings.unregister_embedder"=>encode(kb.embeddings().unregister_embedder(&field::<String>(&args,"space_id")?)?),
73 "embeddings.namespace_vectorization"=>encode(kb.embeddings().namespace_vectorization(&field::<String>(&args,"namespace")?)?),
74 "embeddings.set_namespace_vectorization"=>encode(kb.embeddings().set_namespace_vectorization(&field::<String>(&args,"namespace")?,field(&args,"enabled")?)?),
75 "embeddings.vectorization"=>encode(kb.embeddings().vectorization(&field::<String>(&args,"namespace")?,&field::<String>(&args,"target")?)?),
76 "embeddings.set_vectorization"=>encode(kb.embeddings().set_vectorization(&field::<String>(&args,"namespace")?,&field::<String>(&args,"target")?,field(&args,"enabled")?)?),
77 "embeddings.delete_space"=>encode(kb.embeddings().delete_space(&field::<String>(&args,"id")?)?),
78 _=>Err(Error::Validation(format!("unknown operation: {operation}"))),
79 }
80}
81
82pub fn envelope(result:Result<Value>)->Value {
83 match result {Ok(value)=>json!({"ok":true,"result":value}),Err(error)=>json!({"ok":false,"error":{"code":error.code(),"message":error.to_string()}})}
84}