photon-backend-fluvio 0.1.2

Fluvio storage adapter for Photon
Documentation
//! Fluvio client connection helpers.

use std::sync::Arc;

use fluvio::{Fluvio, FluvioClusterConfig};
use photon_backend::{map_broker_connect_err, PhotonError, Result};

use crate::config::FluvioConfig;

/// Shared Fluvio client handle.
pub type SharedClient = Arc<Fluvio>;

/// Connect to a Fluvio cluster via SC endpoint.
///
/// # Errors
///
/// Returns an error when the security policy rejects plaintext endpoints or connection fails.
pub async fn connect_fluvio(config: &FluvioConfig) -> Result<SharedClient> {
    validate_endpoint(&config.endpoint)?;
    config.transport_security.check_endpoint(&config.endpoint)?;
    let cluster_config = FluvioClusterConfig::new(config.endpoint.clone());
    let client = Fluvio::connect_with_config(&cluster_config)
        .await
        .map_err(|e| map_broker_connect_err("fluvio connect", &config.endpoint, e))?;
    Ok(Arc::new(client))
}

/// Validate endpoint is non-empty.
///
/// # Errors
///
/// Returns an error when endpoint string is empty after parsing.
pub fn validate_endpoint(endpoint: &str) -> Result<()> {
    if endpoint.trim().is_empty() {
        return Err(PhotonError::Internal(
            "PHOTON_FLUVIO_ENDPOINT empty after parsing".into(),
        ));
    }
    Ok(())
}