areev 0.2.0

Rust SDK for the Areev knowledge database — gRPC and HTTP transports
Documentation
//! `Connections` resource — the bridge between a connected `connector` and a
//! Knowledge Source.
//!
//! A `ConnectionRecord` ties a stored connector credential (in the Axtion
//! vault) to a principal/org, and its `id` is the `connection_id` you pass to
//! [`crate::resources::KnowledgeSources::create`] for OAuth-backed connectors.
//! Access via [`crate::Areev::connections`].
//!
//! These `/connections` endpoints are a **shadow API** — they are not in the
//! published OpenAPI spec, so this resource is hand-written (like
//! [`crate::resources::Connectors`]) rather than generated.
//!
//! `ConnectionRecord` fields: `id`, `connector_name`,
//! `connector_display_name`, `axtion_credential_id`, `auth_type`, `status`,
//! `created_at_ms`, `memory_id`, `account_email`, `created_by`, `org_id`.
//! `ConnectionRecord.id` is the `connection_id` for `knowledge_sources.create`
//! — there is no other way to obtain it.
//!
//! # End-to-end: connector → connection → Knowledge Source
//!
//! **OAuth connector** (`google-drive` / `dropbox` / `notion` / `confluence`)
//! — the connection id comes from the OAuth handshake:
//!
//! ```no_run
//! # #[tokio::main]
//! # async fn main() -> areev::Result<()> {
//! # let areev = areev::Areev::from_env();
//! // 1. mint the provider OAuth URL and redirect the user to it
//! let oauth = areev
//!     .connectors()
//!     .authorize("google-drive", "https://app.example/oauth/cb")
//!     .await?;
//! // 2. user consents in the browser → poll until the handshake completes
//! let _status = areev.connectors().poll_oauth("google-drive", &oauth.state).await?;
//! // 3. look up the resulting ConnectionRecord to get its id
//! let conns = areev.connections().list(Some("google-drive")).await?;
//! let connection_id = conns["connections"][0]["id"].as_str().unwrap().to_string();
//! // 4. wire it into a Knowledge Source
//! areev
//!     .knowledge_sources()
//!     .create("google-drive", "Team Drive")
//!     .connection_id(&connection_id)
//!     .send()
//!     .await?;
//! # Ok(())
//! # }
//! ```
//!
//! **web-search** (no OAuth, no connection) — grant consent, then create the
//! KS directly with a `url`. Note the underscore (`web_search`) and that `url`
//! is **required**; the server synthesizes a sentinel connection, so you do
//! NOT call `connections` at all:
//!
//! ```no_run
//! # #[tokio::main]
//! # async fn main() -> areev::Result<()> {
//! # let areev = areev::Areev::from_env();
//! areev
//!     .consent()
//!     .grant(
//!         "user-1",
//!         "knowledge_sources:web_search",
//!         None,
//!         Some("web_search-v1-2026-05-09"),
//!     )
//!     .await?;
//! areev
//!     .knowledge_sources()
//!     .create("web_search", "Docs site")
//!     .url("https://example.com")
//!     .send()
//!     .await?;
//! # Ok(())
//! # }
//! ```

use serde_json::{Map, Value};

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

/// Connection-record lifecycle. Access via [`crate::Areev::connections`].
///
/// A `ConnectionRecord` binds a stored connector credential to a principal
/// (and optionally a memory, for harness tool credentials). Its `id` is the
/// `connection_id` argument of [`crate::resources::KnowledgeSources::create`].
pub struct Connections<'a> {
    http: &'a HttpClient,
}

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

    // ── Principal-scoped (Knowledge-Source connections) ─────────────────

    /// List the calling principal's connection records.
    ///
    /// Pass `connector` (the registry slug, e.g. `"google-drive"`) to filter
    /// to a single connector. Returns
    /// `{ "connections": [ConnectionRecord], "count": <int> }`. The `id` of a
    /// record is the `connection_id` for
    /// [`crate::resources::KnowledgeSources::create`].
    pub async fn list(&self, connector: Option<&str>) -> Result<Value> {
        let query = connector.map(|c| serde_json::json!({ "connector": c }));
        self.http._get("/connections", query.as_ref()).await
    }

    /// Start a [`CreateConnectionBuilder`] for a new principal-scoped
    /// connection record.
    ///
    /// Only `connector_name` is required. **Idempotent** per
    /// `(connector, principal, memory_id)` — calling twice returns the same
    /// record rather than creating a duplicate. `send()` returns the
    /// unwrapped `ConnectionRecord` (read `["id"]` directly).
    pub fn create<'b>(&'b self, connector_name: &str) -> CreateConnectionBuilder<'b, 'a> {
        CreateConnectionBuilder::new(self, connector_name, None)
    }

    /// Delete a connection record. **Destructive — not retried.**
    ///
    /// With `cascade = true` this also deletes every Knowledge Source that
    /// depends on the connection. With `forget_grains = true` it additionally
    /// crypto-erases the grains those sources produced (GDPR Art. 17). Because
    /// the no-retry guard is on, a transport failure surfaces to the caller
    /// instead of silently re-firing a cascade delete.
    ///
    /// <div class="warning">
    ///
    /// This is irreversible. `cascade` removes dependent Knowledge Sources and
    /// `forget_grains` crypto-erases their grains — there is no undo.
    ///
    /// </div>
    pub async fn delete(
        &self,
        connection_id: &str,
        cascade: bool,
        forget_grains: bool,
    ) -> Result<Value> {
        let path =
            format!("/connections/{connection_id}?cascade={cascade}&forget_grains={forget_grains}");
        self.http._delete_no_retry(&path).await
    }

    // ── Memory-scoped (harness tool credentials) ────────────────────────

    /// List connection records scoped to a memory (harness tool credentials).
    ///
    /// Returns `{ "connections": [ConnectionRecord], "count": <int> }`.
    pub async fn list_for_memory(&self, memory_id: &str) -> Result<Value> {
        let path = format!("/memories/{memory_id}/connections");
        self.http._get(&path, None).await
    }

    /// Start a [`CreateConnectionBuilder`] for a new memory-scoped connection
    /// record (harness tool credentials).
    ///
    /// Only `connector_name` is required. Idempotent per
    /// `(connector, principal, memory_id)`. `send()` returns the
    /// unwrapped `ConnectionRecord` (read `["id"]` directly).
    pub fn create_for_memory<'b>(
        &'b self,
        memory_id: &str,
        connector_name: &str,
    ) -> CreateConnectionBuilder<'b, 'a> {
        CreateConnectionBuilder::new(self, connector_name, Some(memory_id.to_string()))
    }

    /// Delete a memory-scoped connection record. **Destructive — not retried.**
    ///
    /// With `cascade = true` also removes dependent Knowledge Sources.
    ///
    /// <div class="warning">
    ///
    /// This is irreversible — `cascade` removes dependent Knowledge Sources.
    ///
    /// </div>
    pub async fn delete_for_memory(
        &self,
        memory_id: &str,
        connection_id: &str,
        cascade: bool,
    ) -> Result<Value> {
        let path = format!("/memories/{memory_id}/connections/{connection_id}?cascade={cascade}");
        self.http._delete_no_retry(&path).await
    }
}

/// Builder for [`Connections::create`] / [`Connections::create_for_memory`].
///
/// `connector_name` is required (set by the constructor). Chain optional knobs
/// and call [`send`](CreateConnectionBuilder::send). When built via
/// `create_for_memory`, the request targets the memory-scoped endpoint and the
/// `.memory_id()` knob is a no-op (the path already carries it).
pub struct CreateConnectionBuilder<'b, 'a> {
    conns: &'b Connections<'a>,
    /// `Some(id)` → memory-scoped endpoint; `None` → principal-scoped.
    scope_memory_id: Option<String>,
    body: Map<String, Value>,
}

impl<'b, 'a> CreateConnectionBuilder<'b, 'a> {
    fn new(
        conns: &'b Connections<'a>,
        connector_name: &str,
        scope_memory_id: Option<String>,
    ) -> Self {
        let mut body = Map::new();
        body.insert(
            "connector_name".into(),
            Value::String(connector_name.to_string()),
        );
        body.insert("auth_type".into(), Value::String("api_key".to_string()));
        Self {
            conns,
            scope_memory_id,
            body,
        }
    }

    /// Human-readable connector label (defaults server-side to
    /// `connector_name`).
    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
    }

    /// Axtion vault credential id this connection points at.
    pub fn axtion_credential_id(mut self, axtion_credential_id: &str) -> Self {
        self.body.insert(
            "axtion_credential_id".into(),
            Value::String(axtion_credential_id.to_string()),
        );
        self
    }

    /// Authentication scheme — e.g. `"oauth"`, `"api_key"`. Defaults to
    /// `"api_key"`.
    pub fn auth_type(mut self, auth_type: &str) -> Self {
        self.body
            .insert("auth_type".into(), Value::String(auth_type.to_string()));
        self
    }

    /// Memory the connection is granted under (principal-scoped create only).
    ///
    /// Ignored when the builder was created via
    /// [`Connections::create_for_memory`] — the path already carries the id.
    pub fn memory_id(mut self, memory_id: &str) -> Self {
        if self.scope_memory_id.is_none() {
            self.body
                .insert("memory_id".into(), Value::String(memory_id.to_string()));
        }
        self
    }

    /// Issue the create request. Returns the unwrapped `ConnectionRecord`
    /// value — the server wraps it in `{ "connection": … }`, which this
    /// strips so callers read `["id"]` directly (matching the TS SDK).
    pub async fn send(self) -> Result<Value> {
        let path = match &self.scope_memory_id {
            Some(mid) => format!("/memories/{mid}/connections"),
            None => "/connections".to_string(),
        };
        let resp = self
            .conns
            .http
            ._post(&path, Some(&Value::Object(self.body)))
            .await?;
        Ok(match resp {
            Value::Object(mut map) => map.remove("connection").unwrap_or(Value::Object(map)),
            other => other,
        })
    }
}