areev 0.2.0

Rust SDK for the Areev knowledge database — gRPC and HTTP transports
Documentation
//! `KnowledgeSources` resource — connect external content (Drive, Dropbox,
//! Confluence, Notion, Web Search) and review-gate proposed memory changes.
//!
//! Mirrors the Python SDK's `client.knowledge_sources.*` surface.
//!
//! Requires: Scale or Custom plan. Calls are issued unconditionally; the
//! server returns `FTR-E001` (surfaced as
//! [`crate::AreevError::FeatureNotAvailable`]) when the org's tier does
//! not include Knowledge Sources — the SDK never gates client-side.

use serde_json::{json, Map, Value};

use crate::error::Result;
use crate::http::HttpClient;

/// Knowledge-source lifecycle + change-proposal review.
///
/// Access via [`crate::Areev::knowledge_sources`].
pub struct KnowledgeSources<'a> {
    http: &'a HttpClient,
    memory_id: String,
}

impl<'a> KnowledgeSources<'a> {
    /// Internal constructor — use [`crate::Areev::knowledge_sources`].
    pub(crate) fn new(http: &'a HttpClient, memory_id: String) -> Self {
        Self { http, memory_id }
    }

    /// Start a [`CreateKsBuilder`] for a new knowledge source.
    ///
    /// `connector_name` is the registry slug (`google-drive`, `dropbox`,
    /// `confluence`, `notion`, `web-search`). Chain optional knobs and
    /// call `send()`. Requires: Scale or Custom plan.
    pub fn create<'b>(&'b self, connector_name: &str, name: &str) -> CreateKsBuilder<'b, 'a> {
        CreateKsBuilder::new(self, connector_name, name)
    }

    /// List the memory's knowledge sources. Requires: Scale or Custom plan.
    pub async fn list(&self) -> Result<Value> {
        let path = format!("/memories/{}/knowledge-sources", self.memory_id);
        self.http._get(&path, None).await
    }

    /// Get one knowledge source by id. Requires: Scale or Custom plan.
    pub async fn get(&self, ks_id: &str) -> Result<Value> {
        let path = format!("/memories/{}/knowledge-sources/{}", self.memory_id, ks_id);
        self.http._get(&path, None).await
    }

    /// Update a knowledge source's `name`, `scope`, and/or `sync_policy`.
    /// Only the supplied fields are sent. Requires: Scale or Custom plan.
    pub async fn update(
        &self,
        ks_id: &str,
        name: Option<&str>,
        scope: Option<Value>,
        sync_policy: Option<Value>,
    ) -> Result<Value> {
        let mut body = Map::new();
        if let Some(n) = name {
            body.insert("name".into(), Value::String(n.to_string()));
        }
        if let Some(s) = scope {
            body.insert("scope".into(), s);
        }
        if let Some(sp) = sync_policy {
            body.insert("sync_policy".into(), sp);
        }
        let path = format!("/memories/{}/knowledge-sources/{}", self.memory_id, ks_id);
        self.http._put(&path, Some(&Value::Object(body))).await
    }

    /// Disconnect a knowledge source. Requires: Scale or Custom plan.
    pub async fn delete(&self, ks_id: &str) -> Result<()> {
        let path = format!("/memories/{}/knowledge-sources/{}", self.memory_id, ks_id);
        self.http._delete(&path).await.map(|_| ())
    }

    /// Trigger a manual sync — produces review-gated change proposals.
    /// Requires: Scale or Custom plan.
    pub async fn sync(&self, ks_id: &str, opts: Option<Value>) -> Result<Value> {
        let body = opts.unwrap_or_else(|| json!({}));
        let path = format!(
            "/memories/{}/knowledge-sources/{}/sync",
            self.memory_id, ks_id
        );
        self.http._post(&path, Some(&body)).await
    }

    /// Browse the source's available items (folders / pages) for scoping.
    /// Requires: Scale or Custom plan.
    pub async fn browse(&self, ks_id: &str, opts: Option<Value>) -> Result<Value> {
        let body = opts.unwrap_or_else(|| json!({}));
        let path = format!(
            "/memories/{}/knowledge-sources/{}/browse",
            self.memory_id, ks_id
        );
        self.http._post(&path, Some(&body)).await
    }

    /// List pending change proposals awaiting review for a source.
    /// Requires: Scale or Custom plan.
    pub async fn list_proposals(&self, ks_id: &str) -> Result<Value> {
        let path = format!(
            "/memories/{}/knowledge-sources/{}/proposals",
            self.memory_id, ks_id
        );
        self.http._get(&path, None).await
    }

    /// Approve a change proposal — applies the proposed grain changes.
    /// Requires: Scale or Custom plan. Emits an audit event.
    pub async fn approve_proposal(
        &self,
        ks_id: &str,
        proposal_id: &str,
        opts: Option<Value>,
    ) -> Result<Value> {
        let body = opts.unwrap_or_else(|| json!({}));
        let path = format!(
            "/memories/{}/knowledge-sources/{}/proposals/{}/approve",
            self.memory_id, ks_id, proposal_id
        );
        self.http._post(&path, Some(&body)).await
    }

    /// Reject a change proposal — discards the proposed changes.
    /// Requires: Scale or Custom plan. Emits an audit event.
    pub async fn reject_proposal(
        &self,
        ks_id: &str,
        proposal_id: &str,
        reason: Option<&str>,
    ) -> Result<Value> {
        let mut body = Map::new();
        if let Some(r) = reason {
            body.insert("reason".into(), Value::String(r.to_string()));
        }
        let path = format!(
            "/memories/{}/knowledge-sources/{}/proposals/{}/reject",
            self.memory_id, ks_id, proposal_id
        );
        self.http._post(&path, Some(&Value::Object(body))).await
    }
}

/// Builder for [`KnowledgeSources::create`].
pub struct CreateKsBuilder<'b, 'a> {
    ks: &'b KnowledgeSources<'a>,
    body: Map<String, Value>,
}

impl<'b, 'a> CreateKsBuilder<'b, 'a> {
    fn new(ks: &'b KnowledgeSources<'a>, connector_name: &str, name: &str) -> Self {
        let mut body = Map::new();
        body.insert(
            "connector_name".into(),
            Value::String(connector_name.to_string()),
        );
        body.insert("name".into(), Value::String(name.to_string()));
        Self { ks, body }
    }

    /// Existing per-principal connector connection id to reuse.
    pub fn connection_id(mut self, connection_id: &str) -> Self {
        self.body.insert(
            "connection_id".into(),
            Value::String(connection_id.to_string()),
        );
        self
    }

    /// Human-readable connector label.
    pub fn connector_display_name(mut self, display_name: &str) -> Self {
        self.body.insert(
            "connector_display_name".into(),
            Value::String(display_name.to_string()),
        );
        self
    }

    /// Source URL (Web Search / Confluence space / etc.).
    pub fn url(mut self, url: &str) -> Self {
        self.body
            .insert("url".into(), Value::String(url.to_string()));
        self
    }

    /// Scoping selection (folders / pages to ingest).
    pub fn scope(mut self, scope: Value) -> Self {
        self.body.insert("scope".into(), scope);
        self
    }

    /// Sync policy, e.g. `{"mode": "auto", "poll_seconds": 60}`.
    pub fn sync_policy(mut self, sync_policy: Value) -> Self {
        self.body.insert("sync_policy".into(), sync_policy);
        self
    }

    /// Ingest only the main content (strip boilerplate / chrome).
    pub fn main_content_only(mut self, main_content_only: bool) -> Self {
        self.body
            .insert("main_content_only".into(), Value::Bool(main_content_only));
        self
    }

    /// Issue the create request.
    pub async fn send(self) -> Result<Value> {
        let path = format!("/memories/{}/knowledge-sources", self.ks.memory_id);
        self.ks
            .http
            ._post(&path, Some(&Value::Object(self.body)))
            .await
    }
}