aurabase 0.1.1

Official Rust SDK for Aurabase: high-performance open-source Backend-as-a-Service (BaaS)
Documentation
use crate::auth_store::{is_token_expired, AuthStore, InMemoryAuthStore};
use crate::error::AuraError;
use crate::types::{AuraResponse, Meta, User};
use reqwest::Client as HttpClient;
use std::sync::Arc;
use tokio::sync::Mutex;

// Imports des services
use crate::services::ai::AiService;
use crate::services::auth::{AuthAdminService, AuthService};
use crate::services::db::DatabaseService;
use crate::services::functions::FunctionsService;
use crate::services::notifications::NotificationsService;
use crate::services::realtime::RealtimeService;
use crate::services::storage::StorageService;

#[derive(Debug, Clone, Default)]
pub struct AuraClientOptions {
    pub auth_storage: Option<Arc<dyn AuthStore>>,
}

#[derive(Debug)]
pub struct InnerClient {
    pub base_url: String,
    pub api_key: Option<String>,
    pub http_client: HttpClient,
    pub auth_store: Arc<dyn AuthStore>,
    pub refresh_lock: Mutex<()>,
}

/// Client principal pour interagir avec Aurabase.
/// Cette structure utilise un `Arc` en interne pour être clonable à bas coût
/// et thread-safe.
#[derive(Debug, Clone)]
pub struct AuraClient {
    pub(crate) inner: Arc<InnerClient>,
}

pub enum RequestBody {
    Json(serde_json::Value),
    Multipart(reqwest::multipart::Form),
    None,
}

impl AuraClient {
    /// Crée une nouvelle instance de client Aurabase.
    ///
    /// **L'URL EST le projet.** Chaque projet est servi sur son propre sous-domaine
    /// (`https://<slug>.aurabase.cloud`) et la passerelle résout le tenant depuis l'hôte puis
    /// injecte elle-même le `project_id` dans le chemin
    /// (`aura-gateway/src/middleware/project_subdomain.rs::inject_project_id`). Le client
    /// n'écrit donc, et ne transporte, aucun identifiant de projet.
    ///
    /// # Arguments
    /// * `base_url` — URL du projet (`https://<slug>.aurabase.cloud`).
    /// * `api_key`  — clé API du projet (`aura_anon_…`, `aura_auth_…`, `aura_sk_…`). Une chaîne
    ///   vide est acceptée (seules les routes auth publiques répondent alors), d'où
    ///   l'avertissement.
    pub fn new(base_url: &str, api_key: &str, options: Option<AuraClientOptions>) -> Self {
        let base_url = base_url.trim_end_matches('/').to_string();
        let opts = options.unwrap_or_default();

        let auth_store = opts
            .auth_storage
            .unwrap_or_else(|| Arc::new(InMemoryAuthStore::new()));

        let http_client = HttpClient::builder().build().unwrap_or_default();

        let api_key = if api_key.is_empty() {
            println!(
                "[aurabase] Warning: No API key provided. Most operations require an API key."
            );
            None
        } else {
            Some(api_key.to_string())
        };

        Self {
            inner: Arc::new(InnerClient {
                base_url,
                api_key,
                http_client,
                auth_store,
                refresh_lock: Mutex::new(()),
            }),
        }
    }

    pub fn base_url(&self) -> &str {
        &self.inner.base_url
    }

    pub fn api_key(&self) -> Option<&str> {
        self.inner.api_key.as_deref()
    }

    /// Préfixe `"<project_id>:"` des noms de canaux realtime, `""` s'il est inconnu.
    ///
    /// `aura-realtime` exige que tout canal commence par le `project_id` du JETON
    /// (`handlers/websocket.rs`, `expected_prefix`). Le SDK ne porte plus d'identifiant de
    /// projet — le projet est celui de l'URL — mais le jeton de session, lui, porte ce claim :
    /// on le lit là, dans le MÊME jeton que celui envoyé au serveur. Aucune décision de
    /// sécurité ne repose sur ce décodage local (pas de vérification de signature) : au pire
    /// le serveur rejette le canal. Parité `aurabase-js`
    /// (`RealtimeService._projectPrefix` / `jwt.ts::getTokenProjectId`).
    pub fn project_prefix(&self) -> String {
        match self.token_project_id() {
            Some(pid) => format!("{pid}:"),
            None => String::new(),
        }
    }

    /// Claim `project_id` du jeton de session courant, `None` s'il n'y a pas de jeton ou si le
    /// jeton n'en porte pas.
    pub fn token_project_id(&self) -> Option<String> {
        let token = self.inner.auth_store.token()?;
        let payload = crate::auth_store::decode_jwt_payload(&token)?;
        payload
            .get("project_id")
            .and_then(|v| v.as_str())
            .filter(|s| !s.is_empty())
            .map(str::to_string)
    }

    pub fn auth_store(&self) -> &Arc<dyn AuthStore> {
        &self.inner.auth_store
    }

    // -- Services -----------------------------------------------------------

    pub fn auth(&self) -> AuthService {
        AuthService::new(self.clone())
    }

    pub fn auth_admin(&self) -> AuthAdminService {
        AuthAdminService::new(self.clone())
    }

    pub fn db(&self) -> DatabaseService {
        DatabaseService::new(self.clone())
    }

    pub fn storage(&self) -> StorageService {
        StorageService::new(self.clone())
    }

    pub fn realtime(&self) -> RealtimeService {
        RealtimeService::new(self.clone())
    }

    pub fn functions(&self) -> FunctionsService {
        FunctionsService::new(self.clone())
    }

    pub fn ai(&self) -> AiService {
        AiService::new(self.clone())
    }

    pub fn notifications(&self) -> NotificationsService {
        NotificationsService::new(self.clone())
    }

    // Raccourcis Database (parité avec le SDK JavaScript `aura.from(...)`)
    pub fn from<Row>(&self, table: &str) -> crate::services::db::QueryBuilder<Vec<Row>> {
        self.db().from(table)
    }

    // Raccourci realtime
    pub fn channel(&self, name: &str) -> crate::services::realtime::AuraChannel {
        self.realtime().channel(name, None)
    }

    // -- HTTP helper ---------------------------------------------------------

    pub(crate) async fn request<T>(
        &self,
        method: reqwest::Method,
        path: &str,
        body: RequestBody,
    ) -> Result<AuraResponse<T>, AuraError>
    where
        T: serde::de::DeserializeOwned,
    {
        // 1. Tenter de rafraîchir le token si expiré
        let _ = self.maybe_refresh().await;

        // 2. Construire la requête
        let url = format!("{}{}", self.inner.base_url, path);
        let mut builder = self.inner.http_client.request(method, &url);

        // Clé API
        if let Some(ref api_key) = self.inner.api_key {
            builder = builder.header("apikey", api_key);
        }

        // Bearer Token
        if let Some(token) = self.inner.auth_store.token() {
            builder = builder.header("Authorization", format!("Bearer {}", token));
        }

        // Injecter le body
        builder = match body {
            RequestBody::Json(json_val) => builder.json(&json_val),
            RequestBody::Multipart(form) => builder.multipart(form),
            RequestBody::None => builder,
        };

        // 3. Envoyer
        let res = builder
            .send()
            .await
            .map_err(|e| AuraError::network(&e.to_string()))?;

        let status = res.status().as_u16();

        // 204 No Content
        if status == 204 {
            return Ok(AuraResponse {
                data: None,
                error: None,
                meta: None,
            });
        }

        // Gérer les erreurs HTTP
        if !res.status().is_success() {
            let body_val: Option<serde_json::Value> = res.json().await.ok();
            let mut code = format!("http_{}", status);
            let mut message = format!("HTTP {}", status);
            let mut details = None;

            if let Some(ref body) = body_val {
                details = Some(body.clone());
                if let Some(err_obj) = body.get("error") {
                    if let Some(c) = err_obj.get("code").and_then(|v| v.as_str()) {
                        code = c.to_string();
                    }
                    if let Some(m) = err_obj.get("message").and_then(|v| v.as_str()) {
                        message = m.to_string();
                    }
                } else if let Some(m) = body.get("message").and_then(|v| v.as_str()) {
                    message = m.to_string();
                }
            }

            return Ok(AuraResponse {
                data: None,
                error: Some(AuraError::new(status, &code, &message, details)),
                meta: None,
            });
        }

        // Succès
        match res.json::<serde_json::Value>().await {
            Ok(body) => {
                let data_val = body.get("data").cloned().unwrap_or_else(|| body.clone());
                let meta: Option<Meta> = body
                    .get("meta")
                    .and_then(|m| serde_json::from_value(m.clone()).ok());

                match serde_json::from_value::<T>(data_val.clone()) {
                    Ok(data) => Ok(AuraResponse {
                        data: Some(data),
                        error: None,
                        meta,
                    }),
                    Err(e) => {
                        let err_str = e.to_string();
                        if err_str.contains("expected a sequence") && data_val.is_object() {
                            let array_val = serde_json::Value::Array(vec![data_val]);
                            match serde_json::from_value::<T>(array_val) {
                                Ok(data) => Ok(AuraResponse {
                                    data: Some(data),
                                    error: None,
                                    meta,
                                }),
                                Err(_) => Err(AuraError::serialization(&err_str)),
                            }
                        } else {
                            Err(AuraError::serialization(&err_str))
                        }
                    }
                }
            }
            Err(_) => Ok(AuraResponse {
                data: None,
                error: None,
                meta: None,
            }),
        }
    }

    /// Rafraîchit le token s'il a expiré et qu'un refresh token est disponible.
    async fn maybe_refresh(&self) -> Result<(), AuraError> {
        let store = &self.inner.auth_store;
        let token = store.token();
        let refresh_token = store.refresh_token();

        if token.is_none() || refresh_token.is_none() {
            return Ok(());
        }

        let token_str = token.unwrap();
        let refresh_token_str = refresh_token.unwrap();

        // Si le token expire dans moins de 30 secondes
        if !is_token_expired(&token_str, 30) {
            return Ok(());
        }

        // Verrou pour dédupliquer les appels concurrents
        let _guard = self.inner.refresh_lock.lock().await;

        // Re-vérifier l'expiration après l'acquisition du verrou
        if let Some(current_tok) = store.token() {
            if !is_token_expired(&current_tok, 30) {
                return Ok(());
            }
        }

        let url = format!("{}/v1/auth/refresh", self.inner.base_url);
        let payload = serde_json::json!({ "refresh_token": refresh_token_str });

        // K7 — la clé API doit être posée ICI aussi. Ce chemin court-circuite `request()`, le
        // SEUL autre endroit du client qui la pose. `/v1/auth/{pid}/refresh` figure bien dans
        // `PUBLIC_AUTH_SUFFIXES` du plan données, mais cette liste ne rend PUBLIC que le JWT :
        // l'étape 3 de `data_plane_auth_middleware` y exige la clé API et répond sinon 401
        // « Clé API requise ». Sans elle, aucun rafraîchissement proactif n'aboutissait jamais.
        let mut requete = self.inner.http_client.post(&url).json(&payload);
        if let Some(ref api_key) = self.inner.api_key {
            requete = requete.header("apikey", api_key);
        }
        let res = requete.send().await;

        match res {
            Ok(r) if r.status().is_success() => {
                if let Ok(body) = r.json::<serde_json::Value>().await {
                    let data = body.get("data").unwrap_or(&body);
                    if let (Some(acc), Some(refr), Some(usr)) = (
                        data.get("access_token").and_then(|v| v.as_str()),
                        data.get("refresh_token").and_then(|v| v.as_str()),
                        data.get("user"),
                    ) {
                        if let Ok(user) = serde_json::from_value::<User>(usr.clone()) {
                            store.save(acc.to_string(), refr.to_string(), user);
                        }
                    }
                }
            }
            Ok(r) => {
                // K7, seconde moitié — n'effacer la session QUE sur un rejet DÉFINITIF du
                // refresh token lui-même, jamais sur un 4xx quelconque de cette requête.
                //
                // Le `store.clear()` inconditionnel qui se trouvait ici déconnectait
                // silencieusement l'utilisateur sur le 401 « Clé API requise » que la ligne
                // manquante ci-dessus provoquait — et le ferait encore sur toute rotation de
                // clé anon, toute clé expirée, tout 5xx d'infrastructure.
                //
                // Deux signaux, et deux seulement — mêmes que `aurabase-js`
                // (`packages/core/src/fetch.ts`, `doRefresh`), pour que les quatre SDK aient le
                // même verdict :
                //   1. 401 + `error.code == "REFRESH_TOKEN_REJECTED"` — discriminant
                //      CONTRACTUEL d'aura-auth (`AuraError::RefreshTokenRejected`,
                //      `libs/aura-core/src/error.rs`), émis sur les seules branches qui
                //      condamnent le jeton présenté. Le gateway ne construit jamais cette
                //      variante : un rejet de CLÉ API y produit `UNAUTHORIZED`, code différent.
                //   2. 403 — sur ce chemin la clé part toujours en en-tête, donc le seul 403
                //      atteignable vient d'aura-auth (`project_id` du jeton étranger au projet).
                //
                // Tout le reste PRÉSERVE la session : 401 « UNAUTHORIZED » (indiscernable d'un
                // refus de clé côté passerelle), 400, 5xx, corps illisible.
                let statut = r.status().as_u16();
                let corps: Option<serde_json::Value> = r.json().await.ok();
                let code = corps
                    .as_ref()
                    .and_then(|b| b.get("error"))
                    .and_then(|e| e.get("code"))
                    .and_then(|v| v.as_str());
                let rejet_definitif =
                    statut == 403 || (statut == 401 && code == Some("REFRESH_TOKEN_REJECTED"));

                // Garde de validité : si une connexion ou une rotation concurrente a changé la
                // session pendant l'attente, le verdict porte sur un ANCIEN jeton — ne rien
                // détruire.
                if rejet_definitif
                    && store.refresh_token().as_deref() == Some(refresh_token_str.as_str())
                {
                    store.clear();
                }
            }
            Err(_) => {
                // Erreur réseau -> conserver les tokens existants
            }
        }

        Ok(())
    }
}