use serde_json::{json, Map, Value};
use crate::error::{AreevError, Result};
use crate::http::HttpClient;
pub struct Grains<'a> {
http: &'a HttpClient,
memory_id: String,
}
impl<'a> Grains<'a> {
pub(crate) fn new(http: &'a HttpClient, memory_id: String) -> Self {
Self { http, memory_id }
}
pub async fn add(
&self,
grain_type: &str,
fields: Value,
namespace: Option<&str>,
) -> Result<String> {
let mut body = json!({ "grain_type": grain_type, "fields": fields });
if let Some(ns) = namespace {
body.as_object_mut()
.unwrap()
.insert("namespace".to_string(), Value::String(ns.to_string()));
}
let path = format!("/memories/{}/add", self.memory_id);
let resp = self.http._post(&path, Some(&body)).await?;
extract_string(&resp, &["blob_hash", "hash"]).ok_or_else(missing_hash_error)
}
pub async fn batch_add(&self, grains: Vec<Value>) -> Result<Value> {
let body = json!({ "grains": grains });
let path = format!("/memories/{}/batch-add", self.memory_id);
self.http._post(&path, Some(&body)).await
}
pub fn recall(&self, query: &str) -> RecallBuilder<'_, 'a> {
RecallBuilder::new(self, query)
}
pub async fn recall_chain(&self, question: &str, limit: u32) -> Result<Value> {
let body = json!({ "query": question, "limit": limit });
let path = format!("/memories/{}/recall-chain", self.memory_id);
self.http._post(&path, Some(&body)).await
}
pub async fn get(&self, hash: &str) -> Result<Value> {
let path = format!("/memories/{}/grains/{}", self.memory_id, hash);
self.http._get(&path, None).await
}
pub async fn forget(&self, hash: &str) -> Result<()> {
let body = json!({ "blob_hash": hash });
let path = format!("/memories/{}/forget", self.memory_id);
self.http._post(&path, Some(&body)).await.map(|_| ())
}
pub async fn forget_many(&self, hashes: &[&str]) -> Result<()> {
let body = json!({ "hashes": hashes });
let path = format!("/memories/{}/forget", self.memory_id);
self.http._post(&path, Some(&body)).await.map(|_| ())
}
pub async fn forget_user(&self, user_id: &str, confirm: bool) -> Result<Value> {
if !confirm {
return Err(confirm_error("grains.forget_user"));
}
let body = json!({ "user_id": user_id });
let path = format!("/memories/{}/forget", self.memory_id);
self.http._post_sensitive(&path, Some(&body), false).await
}
pub async fn get_raw(&self, hash: &str) -> Result<Vec<u8>> {
let path = format!("/memories/{}/grains/{}/raw", self.memory_id, hash);
self.http._get_bytes(&path, None).await
}
pub async fn list(&self, filters: Option<&Value>) -> Result<Value> {
let path = format!("/memories/{}/grains", self.memory_id);
self.http._get(&path, filters).await
}
pub async fn cal(&self, query: &str) -> Result<Value> {
let path = format!("/memories/{}/cal", self.memory_id);
self.http
._post_raw(&path, query.to_string(), "text/cal")
.await
}
pub async fn cal_ast(&self, ast: &Value) -> Result<Value> {
let path = format!("/memories/{}/cal", self.memory_id);
let body = serde_json::to_string(ast).map_err(AreevError::from)?;
self.http
._post_raw(&path, body, "application/json+cal")
.await
}
pub async fn predicate_registry(&self) -> Result<Value> {
let path = format!("/memories/{}/predicate-registry", self.memory_id);
self.http._get(&path, None).await
}
pub async fn predicate_hint(
&self,
relation: &str,
hint: &str,
confidence: f64,
user_id: &str,
) -> Result<Value> {
let body = json!({
"relation": relation,
"hint": hint,
"confidence": confidence,
"user_id": user_id,
});
let path = format!("/memories/{}/predicate-hint", self.memory_id);
self.http._post(&path, Some(&body)).await
}
pub async fn supersede(&self, old_hash: &str, fields: Value) -> Result<String> {
let body = json!({
"blob_hash": old_hash,
"fields": fields,
});
let path = format!("/memories/{}/supersede", self.memory_id);
let resp = self.http._post(&path, Some(&body)).await?;
extract_string(&resp, &["new_hash", "blob_hash", "hash"]).ok_or_else(missing_hash_error)
}
pub fn accumulate<'b>(
&'b self,
grain_type: &str,
subject: &str,
relation: &str,
) -> AccumulateBuilder<'b, 'a> {
AccumulateBuilder::new(self, grain_type, subject, relation)
}
pub async fn detect_pii(&self, text: &str) -> Result<Value> {
let body = json!({ "text": text });
let path = format!("/memories/{}/detect-pii", self.memory_id);
self.http._post_sensitive(&path, Some(&body), true).await
}
}
pub struct RecallBuilder<'b, 'a> {
grains: &'b Grains<'a>,
body: Map<String, Value>,
}
impl<'b, 'a> RecallBuilder<'b, 'a> {
fn new(grains: &'b Grains<'a>, query: &str) -> Self {
let mut body = Map::new();
body.insert("query".to_string(), Value::String(query.to_string()));
body.insert("limit".to_string(), Value::from(10u32));
Self { grains, body }
}
pub fn limit(mut self, limit: u32) -> Self {
self.body.insert("limit".to_string(), Value::from(limit));
self
}
pub fn namespace(mut self, namespace: &str) -> Self {
self.body.insert(
"namespace".to_string(),
Value::String(namespace.to_string()),
);
self
}
pub fn tags(mut self, tags: &[&str]) -> Self {
let arr: Vec<Value> = tags.iter().map(|t| Value::String(t.to_string())).collect();
self.body.insert("tags".to_string(), Value::Array(arr));
self
}
pub fn grain_type(mut self, grain_type: &str) -> Self {
self.body.insert(
"grain_type".to_string(),
Value::String(grain_type.to_string()),
);
self
}
pub fn user_id(mut self, user_id: &str) -> Self {
self.body
.insert("user_id".to_string(), Value::String(user_id.to_string()));
self
}
pub fn extra(mut self, key: &str, value: Value) -> Self {
self.body.insert(key.to_string(), value);
self
}
pub async fn send(self) -> Result<Value> {
let path = format!("/memories/{}/recall", self.grains.memory_id);
self.grains
.http
._post(&path, Some(&Value::Object(self.body)))
.await
}
}
pub struct AccumulateBuilder<'b, 'a> {
grains: &'b Grains<'a>,
body: Map<String, Value>,
}
impl<'b, 'a> AccumulateBuilder<'b, 'a> {
fn new(grains: &'b Grains<'a>, grain_type: &str, subject: &str, relation: &str) -> Self {
let mut body = Map::new();
body.insert(
"grain_type".to_string(),
Value::String(grain_type.to_string()),
);
body.insert(
"target".to_string(),
json!({
"kind": "tip_resolved",
"subject": subject,
"relation": relation,
}),
);
body.insert("add".to_string(), Value::Object(Map::new()));
body.insert("reason".to_string(), Value::String(String::new()));
Self { grains, body }
}
pub fn deltas(mut self, deltas: &[(&str, f64)]) -> Self {
let mut map = Map::new();
for (k, v) in deltas {
map.insert((*k).to_string(), Value::from(*v));
}
self.body.insert("add".to_string(), Value::Object(map));
self
}
pub fn reason(mut self, reason: &str) -> Self {
self.body
.insert("reason".to_string(), Value::String(reason.to_string()));
self
}
pub async fn send(self) -> Result<Value> {
let path = format!("/memories/{}/accumulate", self.grains.memory_id);
self.grains
.http
._post(&path, Some(&Value::Object(self.body)))
.await
}
}
fn extract_string(v: &Value, keys: &[&str]) -> Option<String> {
let map = v.as_object()?;
for k in keys {
if let Some(s) = map.get(*k).and_then(|x| x.as_str()) {
if !s.is_empty() {
return Some(s.to_string());
}
}
}
None
}
fn missing_hash_error() -> AreevError {
AreevError::Server {
code: Some("SDK-E012".to_string()),
http_status: 502,
message: "unexpected server response: no grain hash".to_string(),
body: Value::Null,
request_id: None,
}
}
pub(crate) fn confirm_error(op: &str) -> AreevError {
AreevError::Validation {
code: Some("SDK-E002".to_string()),
http_status: 0,
message: format!(
"{op} requires confirm=true — this is irreversible (crypto-erasure / drops grains)"
),
body: serde_json::Value::Null,
request_id: None,
}
}