areev 0.1.1

Rust SDK for the Areev knowledge database — gRPC and HTTP transports
Documentation
use crate::error::{AreevError, Result};
use crate::types::*;

/// HTTP transport client for Areev.
pub struct HttpClient {
    client: reqwest::Client,
    base_url: String,
    memory_id: String,
    api_key: Option<String>,
}

impl HttpClient {
    /// Create a new HTTP client.
    ///
    /// # Arguments
    /// - `base_url` — Areev server URL (e.g., `http://localhost:4009`)
    /// - `memory_id` — Memory database ID
    /// - `api_key` — Optional API key for authentication
    pub fn new(base_url: &str, memory_id: &str, api_key: Option<&str>) -> Self {
        Self {
            client: reqwest::Client::new(),
            base_url: base_url.trim_end_matches('/').to_string(),
            memory_id: memory_id.to_string(),
            api_key: api_key.map(|s| s.to_string()),
        }
    }

    fn url(&self, path: &str) -> String {
        format!(
            "{}/api/memories/{}/{}",
            self.base_url, self.memory_id, path
        )
    }

    fn request(&self, method: reqwest::Method, path: &str) -> reqwest::RequestBuilder {
        let mut req = self.client.request(method, self.url(path));
        if let Some(ref key) = self.api_key {
            req = req.header("X-API-Key", key);
        }
        req
    }

    /// Add a grain to memory.
    pub async fn add(&self, req: &AddRequest) -> Result<AddResponse> {
        let resp = self
            .request(reqwest::Method::POST, "add")
            .json(req)
            .send()
            .await?;
        self.handle_response(resp).await
    }

    /// Query (recall) grains from memory.
    pub async fn recall(&self, req: &RecallRequest) -> Result<RecallResponse> {
        let resp = self
            .request(reqwest::Method::GET, "recall")
            .query(req)
            .send()
            .await?;
        self.handle_response(resp).await
    }

    /// Get a grain by its hash.
    pub async fn get(&self, hash: &str) -> Result<GetResponse> {
        let resp = self
            .request(reqwest::Method::GET, &format!("grains/{hash}"))
            .send()
            .await?;
        self.handle_response(resp).await
    }

    /// Forget (delete) a grain.
    pub async fn forget(&self, hash: &str) -> Result<()> {
        let resp = self
            .request(reqwest::Method::DELETE, &format!("grains/{hash}"))
            .send()
            .await?;
        if resp.status().is_success() {
            Ok(())
        } else {
            let status = resp.status().as_u16();
            let message = resp.text().await.unwrap_or_default();
            Err(AreevError::Api { status, message })
        }
    }

    /// Remember natural language text as memory.
    ///
    /// Ingests text, creates a source Observation grain, and optionally
    /// extracts structured beliefs (sync or async via Axtion).
    pub async fn remember(&self, req: &RememberRequest) -> Result<RememberResponse> {
        let resp = self
            .request(reqwest::Method::POST, "remember")
            .json(req)
            .send()
            .await?;
        self.handle_response(resp).await
    }

    /// Supersede (update) a grain.
    pub async fn supersede(&self, req: &SupersedeRequest) -> Result<SupersedeResponse> {
        let resp = self
            .request(reqwest::Method::POST, "supersede")
            .json(req)
            .send()
            .await?;
        self.handle_response(resp).await
    }

    // -- Harness chat (HPL) ---------------------------------------------------

    /// Run one harness turn (Flow-A). Returns a [`HarnessChatResponse`]
    /// whose `status` is either `completed` (terminal) or `requires_action`
    /// (paused on `client://` tools — caller must execute them and call
    /// [`HttpClient::harness_chat_resume`]).
    pub async fn harness_chat(
        &self,
        slug: &str,
        req: &HarnessChatRequest,
    ) -> Result<HarnessChatResponse> {
        let resp = self
            .request(
                reqwest::Method::POST,
                &format!("harnesses/{slug}/chat"),
            )
            .json(req)
            .send()
            .await?;
        self.handle_response(resp).await
    }

    /// Resume a paused harness turn with tool outputs. All
    /// `pending_tool_calls[*].tool_call_id` values from the prior
    /// `requires_action` response must appear exactly once in `req.tool_outputs`.
    pub async fn harness_chat_resume(
        &self,
        slug: &str,
        req: &ChatResumeRequest,
    ) -> Result<HarnessChatResponse> {
        let resp = self
            .request(
                reqwest::Method::POST,
                &format!("harnesses/{slug}/chat/resume"),
            )
            .json(req)
            .send()
            .await?;
        self.handle_response(resp).await
    }

    /// Cancel a paused harness chat session. Idempotent — returns Ok even
    /// when the session does not exist for this caller (HRN-E015 is
    /// intentionally opaque).
    pub async fn cancel_harness_chat_session(
        &self,
        slug: &str,
        session_id: &str,
    ) -> Result<()> {
        let resp = self
            .request(
                reqwest::Method::DELETE,
                &format!("harnesses/{slug}/chat/sessions/{session_id}"),
            )
            .send()
            .await?;
        if resp.status().is_success() || resp.status().as_u16() == 404 {
            Ok(())
        } else {
            let status = resp.status().as_u16();
            let message = resp.text().await.unwrap_or_default();
            Err(AreevError::Api { status, message })
        }
    }

    /// Get health status.
    pub async fn health(&self) -> Result<HealthResponse> {
        let resp = self
            .client
            .get(format!("{}/health", self.base_url))
            .send()
            .await?;
        self.handle_response(resp).await
    }

    /// Get database statistics.
    pub async fn stats(&self) -> Result<StatsResponse> {
        let resp = self
            .request(reqwest::Method::GET, "stats")
            .send()
            .await?;
        self.handle_response(resp).await
    }

    /// Flush write buffer.
    pub async fn flush(&self) -> Result<()> {
        let resp = self
            .request(reqwest::Method::POST, "flush")
            .send()
            .await?;
        if resp.status().is_success() {
            Ok(())
        } else {
            let status = resp.status().as_u16();
            let message = resp.text().await.unwrap_or_default();
            Err(AreevError::Api { status, message })
        }
    }

    async fn handle_response<T: serde::de::DeserializeOwned>(
        &self,
        resp: reqwest::Response,
    ) -> Result<T> {
        if resp.status().is_success() {
            Ok(resp.json().await?)
        } else {
            let status = resp.status().as_u16();
            let message = resp.text().await.unwrap_or_default();
            Err(AreevError::Api { status, message })
        }
    }
}