areev 0.2.0

Rust SDK for the Areev knowledge database — gRPC and HTTP transports
Documentation
//! `Chat` resource — app-style assistant chat over a memory: threads,
//! messages, tool dispatch, and SSE streaming.
//!
//! Mirrors the Python SDK's `client.chat.*` surface.
//!
//! This is the app's general assistant chat (thread-based). For the
//! harness conversational loop (Flow-A, caller-supplied tools) use
//! [`crate::Areev::harness`].

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

use crate::error::Result;
use crate::http::{HttpClient, SseStream};

/// Assistant chat + thread management.
///
/// Access via [`crate::Areev::chat`].
pub struct Chat<'a> {
    http: &'a HttpClient,
    memory_id: String,
}

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

    /// Start a [`ChatSendBuilder`] for one non-streaming assistant turn.
    pub fn send<'b>(&'b self, messages: Vec<Value>) -> ChatSendBuilder<'b, 'a> {
        ChatSendBuilder::new(self, messages)
    }

    /// Open a token-by-token SSE stream for one assistant turn.
    ///
    /// Returns an [`SseStream`] the caller pumps with
    /// [`SseStream::next`]. **Not auto-retried** — a stream is a single
    /// live connection; restart on failure.
    ///
    /// ```no_run
    /// # #[tokio::main]
    /// # async fn main() -> areev::Result<()> {
    /// # use serde_json::json;
    /// let areev = areev::Areev::from_env();
    /// let mut stream = areev
    ///     .chat()
    ///     .stream(vec![json!({"role": "user", "content": "hi"})], None, None)
    ///     .await?;
    /// while let Some(chunk) = stream.next().await? {
    ///     print!("{chunk}");
    /// }
    /// # Ok(())
    /// # }
    /// ```
    pub async fn stream(
        &self,
        messages: Vec<Value>,
        provider: Option<&str>,
        model: Option<&str>,
    ) -> Result<SseStream> {
        let mut body = Map::new();
        body.insert("messages".into(), Value::Array(messages));
        if let Some(p) = provider {
            body.insert("provider".into(), Value::String(p.to_string()));
        }
        if let Some(m) = model {
            body.insert("model".into(), Value::String(m.to_string()));
        }
        let path = format!("/memories/{}/chat/stream", self.memory_id);
        self.http
            ._stream_sse(&path, Some(&Value::Object(body)))
            .await
    }

    /// Dispatch a single app-chat tool call server-side.
    pub async fn dispatch_tool(&self, tool: &str, args: Value) -> Result<Value> {
        let path = format!("/memories/{}/chat/tools/{}", self.memory_id, tool);
        self.http._post(&path, Some(&args)).await
    }

    // ── Threads ──────────────────────────────────────────────────────

    /// Create a chat thread.
    pub async fn create_thread(&self, title: Option<&str>) -> Result<Value> {
        let mut body = Map::new();
        if let Some(t) = title {
            body.insert("title".into(), Value::String(t.to_string()));
        }
        let path = format!("/memories/{}/chat/threads", self.memory_id);
        self.http._post(&path, Some(&Value::Object(body))).await
    }

    /// List chat threads. `filters` is added as query parameters.
    pub async fn list_threads(&self, filters: Option<&Value>) -> Result<Value> {
        let path = format!("/memories/{}/chat/threads", self.memory_id);
        self.http._get(&path, filters).await
    }

    /// Get a chat thread by id.
    pub async fn get_thread(&self, thread_id: &str) -> Result<Value> {
        let path = format!("/memories/{}/chat/threads/{}", self.memory_id, thread_id);
        self.http._get(&path, None).await
    }

    /// Rename a chat thread.
    pub async fn rename_thread(&self, thread_id: &str, title: &str) -> Result<Value> {
        let body = json!({ "title": title });
        let path = format!("/memories/{}/chat/threads/{}", self.memory_id, thread_id);
        self.http._patch(&path, Some(&body)).await
    }

    /// Delete a chat thread and its messages.
    pub async fn delete_thread(&self, thread_id: &str) -> Result<()> {
        let path = format!("/memories/{}/chat/threads/{}", self.memory_id, thread_id);
        self.http._delete(&path).await.map(|_| ())
    }

    /// List the messages in a thread. `filters` is added as query
    /// parameters.
    pub async fn list_thread_messages(
        &self,
        thread_id: &str,
        filters: Option<&Value>,
    ) -> Result<Value> {
        let path = format!(
            "/memories/{}/chat/threads/{}/messages",
            self.memory_id, thread_id
        );
        self.http._get(&path, filters).await
    }
}

/// Builder for [`Chat::send`].
pub struct ChatSendBuilder<'b, 'a> {
    chat: &'b Chat<'a>,
    body: Map<String, Value>,
}

impl<'b, 'a> ChatSendBuilder<'b, 'a> {
    fn new(chat: &'b Chat<'a>, messages: Vec<Value>) -> Self {
        let mut body = Map::new();
        body.insert("messages".into(), Value::Array(messages));
        Self { chat, body }
    }

    /// Attach the turn to an existing thread.
    pub fn thread_id(mut self, thread_id: &str) -> Self {
        self.body
            .insert("thread_id".into(), Value::String(thread_id.to_string()));
        self
    }

    /// Stable conversation id for grouping turns.
    pub fn conversation_id(mut self, conversation_id: &str) -> Self {
        self.body.insert(
            "conversation_id".into(),
            Value::String(conversation_id.to_string()),
        );
        self
    }

    /// Override the LLM provider for this turn.
    pub fn provider(mut self, provider: &str) -> Self {
        self.body
            .insert("provider".into(), Value::String(provider.to_string()));
        self
    }

    /// Override the LLM model for this turn.
    pub fn model(mut self, model: &str) -> Self {
        self.body
            .insert("model".into(), Value::String(model.to_string()));
        self
    }

    /// Enable / disable server-side tool dispatch for this turn.
    pub fn tools_enabled(mut self, tools_enabled: bool) -> Self {
        self.body
            .insert("tools_enabled".into(), Value::Bool(tools_enabled));
        self
    }

    /// Set an arbitrary additional body field — escape hatch.
    pub fn extra(mut self, key: &str, value: Value) -> Self {
        self.body.insert(key.to_string(), value);
        self
    }

    /// Issue the chat request and return the assembled assistant reply.
    ///
    /// `POST /memories/{id}/chat` is **SSE-only** (`text/event-stream`) —
    /// there is no JSON response. This consumes the stream, concatenates
    /// the per-chunk content deltas into the full assistant message, and
    /// returns:
    ///
    /// ```json
    /// { "role": "assistant", "content": "<full text>", "usage": { … },
    ///   "model": "…", "conversation_id": "…" }
    /// ```
    ///
    /// The `usage` / `model` / `conversation_id` fields are populated from
    /// the terminal event when the server emits one. For token-by-token
    /// consumption use [`Chat::stream`] instead.
    pub async fn send(self) -> Result<Value> {
        let path = format!("/memories/{}/chat", self.chat.memory_id);
        let mut stream = self
            .chat
            .http
            ._stream_sse(&path, Some(&Value::Object(self.body)))
            .await?;

        let mut content = String::new();
        let mut tail = Map::new();
        while let Some(payload) = stream.next().await? {
            let Ok(value) = serde_json::from_str::<Value>(&payload) else {
                // Non-JSON SSE payload — treat the bare text as a delta.
                content.push_str(&payload);
                continue;
            };
            if let Some(delta) = extract_delta(&value) {
                content.push_str(&delta);
            } else if let Value::Object(map) = value {
                // Terminal / metadata frame (usage, model, conversation_id).
                tail = map;
            }
        }

        let mut out = Map::new();
        out.insert("role".into(), Value::String("assistant".into()));
        out.insert("content".into(), Value::String(content));
        for key in ["usage", "model", "conversation_id", "id"] {
            if let Some(v) = tail.remove(key) {
                out.insert(key.into(), v);
            }
        }
        Ok(Value::Object(out))
    }
}

/// Pull the assistant content delta out of one SSE chunk, supporting both
/// the cell's native `{"delta": "…"}` shape and the OpenAI-style
/// `chat.completion.chunk` shape (`choices[0].delta.content`).
fn extract_delta(value: &Value) -> Option<String> {
    if let Some(s) = value.get("delta").and_then(|v| v.as_str()) {
        return Some(s.to_string());
    }
    value
        .get("choices")
        .and_then(|c| c.as_array())
        .and_then(|arr| arr.first())
        .and_then(|c| c.get("delta"))
        .and_then(|d| d.get("content"))
        .and_then(|v| v.as_str())
        .map(|s| s.to_string())
}