Command Palette

Search for a command to run...

Rust

High-performance Rust SDK with memory safety and zero-cost abstractions

Installation & Setup
Required dependencies and setup for Rust development

Rust & Cargo Setup

# Install Rust
curl --proto '=https' --tlsv1.2 -sSf https://sh.rustup.rs | sh
source $HOME/.cargo/env

cargo new leakzero-client
cd leakzero-client

export LEAKZERO_API_KEY=your_api_key_here

cargo build --release
cargo run

Cargo.toml

[package]
name = "leakzero-client"
version = "0.1.0"
edition = "2021"

[dependencies]
reqwest = { version = "0.11", features = ["json"] }
serde = { version = "1.0", features = ["derive"] }
serde_json = "1.0"
tokio = { version = "1.0", features = ["full"] }
tracing = "0.1"
tracing-subscriber = "0.3"
uuid = { version = "1.0", features = ["v4"] }
futures = "0.3"

[dev-dependencies]
tokio-test = "0.4"

Dockerfile

FROM rust:1.75 as builder

WORKDIR /app

COPY Cargo.toml Cargo.lock ./

COPY src ./src

RUN cargo build --release

FROM debian:bookworm-slim

RUN apt-get update && apt-get install -y \
    ca-certificates \
    && rm -rf /var/lib/apt/lists/*

WORKDIR /app

COPY --from=builder /app/target/release/leakzero-client /usr/local/bin/leakzero-client

CMD ["leakzero-client"]
Rust Version: This example requires Rust 1.70+ for modern async features. Uses tokio for async runtime and reqwest for HTTP.
Quick Start Example
Simple example to get you started with Rust

main.rs

use std::env;
use tokio;

#[tokio::main]
async fn main() -> Result<(), Box<dyn std::error::Error>> {
    let client = leakzero::Client::new(env::var("LEAKZERO_API_KEY")?)?;

    let options = leakzero::SearchOptions::new("email", "[email protected]")
        .with_row_limit(100);

    let result = client.regular_search(&options).await?;

    println!("Found {} records", result.total_results);
    println!("Cost: {:.6}", result.cost);


    for (i, record) in result.data.iter().take(5).enumerate() {
        let breach = record.get("breach_name")
            .and_then(|v| v.as_str())
            .unwrap_or("Unknown");
        let date = record.get("breach_date")
            .and_then(|v| v.as_str())
            .unwrap_or("Unknown");
        println!("{}. Breach: {}, Date: {}", i + 1, breach, date);
    }


    let analysis = client.analyze_breaches(&result);
    println!("\nUnique breaches: {}",
        analysis["unique_breaches"].as_u64().unwrap_or(0));

    Ok(())
}
Production ReadyComplete Rust Implementation
Full-featured client with async support, error handling, and performance optimizations

lib.rs

use reqwest::{Client as HttpClient, Response};
use serde::{Deserialize, Serialize};
use serde_json::{Map, Value};
use std::collections::HashMap;
use std::env;
use std::error::Error;
use std::fmt;
use std::time::{Duration, SystemTime, UNIX_EPOCH};
use tokio::time::sleep;
use tracing::{error, info, warn};
use uuid::Uuid;


#[derive(Debug, Clone, Serialize)]
pub struct SearchOptions {
    pub field: String,
    pub value: String,
    pub option: String,
    #[serde(rename = "rowLimit")]
    pub row_limit: u32,
    #[serde(rename = "autoExport", skip_serializing_if = "Option::is_none")]
    pub auto_export: Option<Value>,
}

impl SearchOptions {
    pub fn new(field: impl Into<String>, value: impl Into<String>) -> Self {
        Self {
            field: field.into(),
            value: value.into(),
            option: "exact".to_string(),
            row_limit: 1000,
            auto_export: None,
        }
    }

    pub fn with_option(mut self, option: impl Into<String>) -> Self {
        self.option = option.into();
        self
    }

    pub fn with_row_limit(mut self, limit: u32) -> Self {
        self.row_limit = limit.min(10000);
        self
    }

    pub fn with_auto_export(mut self, export: Value) -> Self {
        self.auto_export = Some(export);
        self
    }
}


#[derive(Debug, Clone)]
pub struct SearchResult {
    pub data: Vec<Map<String, Value>>,
    pub total_results: u32,
    pub cost: f64,
    pub request_id: String,
    pub processing_time: u32,
    pub analysis: HashMap<String, Value>,
}

impl SearchResult {
    pub fn new() -> Self {
        Self {
            data: Vec::new(),
            total_results: 0,
            cost: 0.0,
            request_id: String::new(),
            processing_time: 0,
            analysis: HashMap::new(),
        }
    }

    pub fn len(&self) -> usize {
        self.data.len()
    }

    pub fn is_empty(&self) -> bool {
        self.data.is_empty()
    }
}


#[derive(Debug, Clone, Deserialize)]
pub struct BalanceResponse {
    pub balance: f64,
    pub currency: String,
    #[serde(rename = "lastUsage")]
    pub last_usage: String,
}


#[derive(Debug, Clone, Deserialize)]
pub struct UsageStats {
    #[serde(rename = "totalRequests")]
    pub total_requests: u32,
    #[serde(rename = "totalCost")]
    pub total_cost: f64,
    #[serde(rename = "requestsToday")]
    pub requests_today: u32,
    #[serde(rename = "costToday")]
    pub cost_today: f64,
    #[serde(rename = "requestsThisWeek")]
    pub requests_this_week: u32,
    #[serde(rename = "costThisWeek")]
    pub cost_this_week: f64,
}


#[derive(Debug)]
pub struct ApiError {
    pub error_code: String,
    pub message: String,
    pub status_code: u16,
}

impl fmt::Display for ApiError {
    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
        write!(
            f,
            "API Error [{}]: {} - {}",
            self.status_code, self.error_code, self.message
        )
    }
}

impl Error for ApiError {}


pub struct Client {
    api_key: String,
    base_url: String,
    http_client: HttpClient,
    max_retries: u32,
}

impl Client {

    pub fn new(api_key: impl Into<String>) -> Result<Self, Box<dyn Error>> {
        Self::with_config(
            api_key,
            "https://api.leakzero.io/api/v1",
            3,
            Duration::from_secs(120),
        )
    }


    pub fn with_config(
        api_key: impl Into<String>,
        base_url: impl Into<String>,
        max_retries: u32,
        timeout: Duration,
    ) -> Result<Self, Box<dyn Error>> {
        let api_key = api_key.into();
        if api_key.trim().is_empty() {
            return Err("API key is required".into());
        }

        let http_client = HttpClient::builder()
            .timeout(timeout)
            .connect_timeout(Duration::from_secs(30))
            .pool_idle_timeout(Duration::from_secs(300))
            .pool_max_idle_per_host(10)
            .user_agent("LeakZero-Rust-SDK/1.0.0")
            .build()?;

        Ok(Self {
            api_key,
            base_url: base_url.into().trim_end_matches('/').to_string(),
            http_client,
            max_retries,
        })
    }


    fn generate_headers(&self) -> reqwest::header::HeaderMap {
        let mut headers = reqwest::header::HeaderMap::new();

        headers.insert("x-api-key", self.api_key.parse().unwrap());
        headers.insert(
            "x-timestamp",
            SystemTime::now()
                .duration_since(UNIX_EPOCH)
                .unwrap()
                .as_secs()
                .to_string()
                .parse()
                .unwrap(),
        );
        headers.insert("x-request-id", self.generate_request_id().parse().unwrap());
        headers.insert("Content-Type", "application/json".parse().unwrap());

        headers
    }


    fn generate_request_id(&self) -> String {
        format!(
            "req_{}_{}",
            SystemTime::now()
                .duration_since(UNIX_EPOCH)
                .unwrap()
                .as_secs(),
            Uuid::new_v4().to_string()[..8].to_string()
        )
    }


    async fn make_request_with_retry(
        &self,
        method: &str,
        endpoint: &str,
        payload: Option<&impl Serialize>,
    ) -> Result<Response, Box<dyn Error>> {
        let url = format!("{}{}", self.base_url, endpoint);
        let headers = self.generate_headers();

        for attempt in 0..=self.max_retries {
            let mut request = match method {
                "GET" => self.http_client.get(&url),
                "POST" => {
                    let mut req = self.http_client.post(&url);
                    if let Some(data) = payload {
                        req = req.json(data);
                    }
                    req
                }
                _ => return Err(format!("Unsupported HTTP method: {}", method).into()),
            };

            request = request.headers(headers.clone());

            match request.send().await {
                Ok(response) => {

                    if response.status().as_u16() == 429 {
                        let retry_after = response
                            .headers()
                            .get("Retry-After")
                            .and_then(|h| h.to_str().ok())
                            .and_then(|s| s.parse::<u64>().ok())
                            .unwrap_or(5)
                            .min(30);

                        warn!("Rate limited. Waiting {}s", retry_after);
                        sleep(Duration::from_secs(retry_after)).await;
                        continue;
                    }


                    if response.status().as_u16() >= 400 {
                        let status_code = response.status().as_u16();
                        let body = response.text().await.unwrap_or_default();

                        let (error_code, error_message) = if let Ok(error_json) = serde_json::from_str::<Value>(&body) {
                            (
                                error_json["error"].as_str().unwrap_or("HTTP_ERROR").to_string(),
                                error_json["message"].as_str().unwrap_or(&format!("HTTP {}", status_code)).to_string(),
                            )
                        } else {
                            ("HTTP_ERROR".to_string(), format!("HTTP {}", status_code))
                        };

                        return Err(Box::new(ApiError {
                            error_code,
                            message: error_message,
                            status_code,
                        }));
                    }

                    info!("API Request successful: {} {}", method, endpoint);
                    return Ok(response);
                }
                Err(e) => {
                    if attempt == self.max_retries {
                        return Err(format!("Request failed after {} attempts: {}", self.max_retries + 1, e).into());
                    }
                    warn!("Request attempt {} failed: {}", attempt + 1, e);
                    sleep(Duration::from_secs(attempt + 1)).await;
                }
            }
        }

        Err("Max retries exceeded".into())
    }


    pub async fn get_balance(&self) -> Result<BalanceResponse, Box<dyn Error>> {
        let response = self
            .make_request_with_retry("GET", "/public/balance", None::<&()>)
            .await?;

        let balance = response.json::<BalanceResponse>().await?;
        Ok(balance)
    }


    pub async fn regular_search(
        &self,
        options: &SearchOptions,
    ) -> Result<SearchResult, Box<dyn Error>> {
        let mut options = options.clone();
        options.row_limit = options.row_limit.min(10000);

        info!("Searching for {}: {}", options.field, options.value);

        let response = self
            .make_request_with_retry("POST", "/public/search/regular", Some(&options))
            .await?;

        let mut result = SearchResult::new();


        if let Some(rows_returned) = response.headers().get("x-rows-returned") {
            result.total_results = rows_returned.to_str()?.parse()?;
        }

        if let Some(charge) = response.headers().get("x-charge") {
            let charge_int: i64 = charge.to_str()?.parse()?;
            result.cost = charge_int as f64 / 1_000_000_000.0;
        }

        if let Some(request_id) = response.headers().get("x-request-id") {
            result.request_id = request_id.to_str()?.to_string();
        }

        if let Some(processing_time) = response.headers().get("x-processing-time") {
            result.processing_time = processing_time.to_str()?.parse()?;
        }


        let data: Value = response.json().await?;
        if let Some(records) = data["data"].as_array() {
            for record in records {
                if let Some(obj) = record.as_object() {
                    result.data.push(obj.clone());
                }
            }
        }

        info!(
            "Found {} records (Cost: {:.6})",
            result.total_results, result.cost
        );
        Ok(result)
    }


    pub async fn advanced_search(
        &self,
        queries: &[SearchOptions],
        row_limit: u32,
    ) -> Result<SearchResult, Box<dyn Error>> {
        if queries.is_empty() {
            return Err("At least one query is required".into());
        }

        let row_limit = row_limit.min(10000);
        let payload = serde_json::json!({
            "queries": queries,
            "rowLimit": row_limit
        });

        info!("Advanced search with {} queries", queries.len());

        let response = self
            .make_request_with_retry("POST", "/public/search/advanced", Some(&payload))
            .await?;

        let mut result = SearchResult::new();


        if let Some(rows_returned) = response.headers().get("x-rows-returned") {
            result.total_results = rows_returned.to_str()?.parse()?;
        }

        if let Some(charge) = response.headers().get("x-charge") {
            let charge_int: i64 = charge.to_str()?.parse()?;
            result.cost = charge_int as f64 / 1_000_000_000.0;
        }

        let data: Value = response.json().await?;
        if let Some(records) = data["data"].as_array() {
            for record in records {
                if let Some(obj) = record.as_object() {
                    result.data.push(obj.clone());
                }
            }
        }

        info!("Advanced search completed (Cost: {:.6})", result.cost);
        Ok(result)
    }


    pub async fn identity_linking_search(
        &self,
        field: &str,
        value: &str,
        context_domain: Option<&str>,
        max_results: u32,
        include_possible: bool,
    ) -> Result<SearchResult, Box<dyn Error>> {
        let max_results = max_results.min(100);
        let payload = serde_json::json!({
            "field": field,
            "value": value,
            "maxDepth": 2,
            "minConfidence": 0.35,
            "maxResults": max_results,
            "includePossible": include_possible,
            "contextDomain": context_domain
        });

        info!("identity linking search for: {}:{}", field, value);

        let response = self
            .make_request_with_retry("POST", "/public/search/identity-linking", Some(&payload))
            .await?;

        let mut result = SearchResult::new();
        let data: Value = response.json().await?;


        let correlation = &data["data"];
        result.analysis.insert("summary".to_string(), correlation["summary"].clone());
        result.analysis.insert("metadata".to_string(), correlation["metadata"].clone());

        if let Some(records) = correlation["results"].as_array() {
            for record in records {
                if let Some(obj) = record.as_object() {
                    result.data.push(obj.clone());
                }
            }
        }

        let linked_count = result
            .analysis
            .get("summary")
            .and_then(|value| value["totalConnections"].as_u64())
            .unwrap_or(0);

        info!("Identity linking found {} linked identities", linked_count);
        Ok(result)
    }


    pub async fn get_usage_stats(&self) -> Result<UsageStats, Box<dyn Error>> {
        let response = self
            .make_request_with_retry("GET", "/public/usage/stats", None::<&()>)
            .await?;

        let stats = response.json::<UsageStats>().await?;
        Ok(stats)
    }


    pub fn analyze_breaches(&self, result: &SearchResult) -> HashMap<String, Value> {
        if result.is_empty() {
            let mut analysis = HashMap::new();
            analysis.insert("error".to_string(), Value::String("No records to analyze".to_string()));
            return analysis;
        }

        let mut breach_distribution: HashMap<String, u32> = HashMap::new();
        let mut date_distribution: HashMap<String, u32> = HashMap::new();
        let mut field_distribution: HashMap<String, u32> = HashMap::new();

        for record in &result.data {

            let breach_name = record
                .get("breach_name")
                .and_then(|v| v.as_str())
                .unwrap_or("Unknown")
                .to_string();
            *breach_distribution.entry(breach_name).or_insert(0) += 1;


            if let Some(breach_date) = record.get("breach_date").and_then(|v| v.as_str()) {
                if breach_date.len() >= 4 {
                    let year = breach_date[..4].to_string();
                    if year.chars().all(|c| c.is_ascii_digit()) {
                        *date_distribution.entry(year).or_insert(0) += 1;
                    }
                }
            }


            for (field_name, field_value) in record {
                if !field_value.is_null()
                    && field_name != "breach_name"
                    && field_name != "breach_date"
                {
                    *field_distribution.entry(field_name.clone()).or_insert(0) += 1;
                }
            }
        }

        let mut sorted_breaches: Vec<_> = breach_distribution.iter().collect();
        sorted_breaches.sort_by(|a, b| b.1.cmp(a.1));

        let top_breaches: Vec<Value> = sorted_breaches
            .iter()
            .take(5)
            .map(|(name, count)| {
                serde_json::json!({
                    "name": name,
                    "count": count
                })
            })
            .collect();

        let years: Vec<u32> = date_distribution
            .keys()
            .filter_map(|k| k.parse().ok())
            .collect();

        let date_range = serde_json::json!({
            "earliest": years.iter().min(),
            "latest": years.iter().max()
        });

        let mut analysis = HashMap::new();
        analysis.insert("total_records".to_string(), Value::Number(result.len().into()));
        analysis.insert("unique_breaches".to_string(), Value::Number(breach_distribution.len().into()));
        analysis.insert("breach_distribution".to_string(), serde_json::to_value(&breach_distribution).unwrap());
        analysis.insert("date_distribution".to_string(), serde_json::to_value(&date_distribution).unwrap());
        analysis.insert("field_distribution".to_string(), serde_json::to_value(&field_distribution).unwrap());
        analysis.insert("date_range".to_string(), date_range);
        analysis.insert("top_breaches".to_string(), Value::Array(top_breaches));

        analysis
    }


    pub async fn bulk_search(
        &self,
        search_terms: &[String],
        field: &str,
        batch_size: usize,
        delay_ms: u64,
        row_limit: u32,
    ) -> Result<HashMap<String, Value>, Box<dyn Error>> {
        let mut results = Vec::new();
        let total_terms = search_terms.len();

        info!("Starting bulk search for {} terms", total_terms);

        for (batch_index, batch) in search_terms.chunks(batch_size).enumerate() {
            let mut batch_results = Vec::new();

            for (term_index, term) in batch.iter().enumerate() {

                if term_index > 0 {
                    sleep(Duration::from_millis(delay_ms)).await;
                }

                let options = SearchOptions::new(field, term).with_row_limit(row_limit);

                match self.regular_search(&options).await {
                    Ok(result) => {
                        batch_results.push(serde_json::json!({
                            "term": term,
                            "result": {
                                "total_results": result.total_results,
                                "cost": result.cost,
                                "data_count": result.data.len()
                            },
                            "success": true
                        }));
                    }
                    Err(e) => {
                        error!("Failed to search for {}: {}", term, e);
                        batch_results.push(serde_json::json!({
                            "term": term,
                            "error": e.to_string(),
                            "success": false
                        }));
                    }
                }
            }

            results.extend(batch_results);


            let processed = std::cmp::min((batch_index + 1) * batch_size, total_terms);
            info!("📊 Processed {}/{} terms", processed, total_terms);


            if processed < total_terms {
                sleep(Duration::from_millis(delay_ms * 2)).await;
            }
        }

        let successful: Vec<_> = results
            .iter()
            .filter(|r| r["success"].as_bool().unwrap_or(false))
            .cloned()
            .collect();

        let failed: Vec<_> = results
            .iter()
            .filter(|r| !r["success"].as_bool().unwrap_or(true))
            .cloned()
            .collect();

        let success_rate = if results.is_empty() {
            0.0
        } else {
            successful.len() as f64 / results.len() as f64 * 100.0
        };

        info!(
            "Bulk search completed: {} successful, {} failed",
            successful.len(),
            failed.len()
        );

        Ok(serde_json::json!({
            "successful": successful,
            "failed": failed,
            "total_processed": results.len(),
            "success_rate": format!("{:.2}%", success_rate)
        })
        .as_object()
        .unwrap()
        .clone())
    }
}


pub async fn demonstrate_basic_usage() -> Result<(), Box<dyn Error>> {
    let api_key = env::var("LEAKZERO_API_KEY")
        .map_err(|_| "LEAKZERO_API_KEY environment variable is required")?;

    let client = Client::new(api_key)?;

    println!("Checking API balance...");
    let balance = client.get_balance().await?;
    println!("Current balance: {:.6}", balance.balance);

    if balance.balance < 0.001 {
        return Err("Insufficient balance. Please top up your account.".into());
    }

    println!("
Performing regular search...");
    let options = SearchOptions::new("email", "[email protected]")
        .with_option("exact")
        .with_row_limit(100);

    let result = client.regular_search(&options).await?;

    println!("Found {} records", result.total_results);
    println!("Cost: {:.6}", result.cost);


    for (i, record) in result.data.iter().take(5).enumerate() {
        let breach = record
            .get("breach_name")
            .and_then(|v| v.as_str())
            .unwrap_or("Unknown");
        let date = record
            .get("breach_date")
            .and_then(|v| v.as_str())
            .unwrap_or("Unknown");
        println!("{}. Breach: {}, Date: {}", i + 1, breach, date);
    }


    if !result.is_empty() {
        let analysis = client.analyze_breaches(&result);
        println!(
            "
Analysis: {} unique breaches",
            analysis["unique_breaches"].as_u64().unwrap_or(0)
        );

        if let Some(top_breaches) = analysis["top_breaches"].as_array() {
            println!("Top breaches:");
            for breach in top_breaches {
                println!(
                    "  - {}: {} records",
                    breach["name"].as_str().unwrap_or("Unknown"),
                    breach["count"].as_u64().unwrap_or(0)
                );
            }
        }
    }

    println!("
Getting usage statistics...");
    let stats = client.get_usage_stats().await?;
    println!("Total requests: {}", stats.total_requests);
    println!("Total cost: {:.6}", stats.total_cost);

    Ok(())
}


pub async fn demonstrate_concurrent_usage() -> Result<(), Box<dyn Error>> {
    let api_key = env::var("LEAKZERO_API_KEY")
        .map_err(|_| "LEAKZERO_API_KEY environment variable is required")?;

    let client = Client::new(api_key)?;


    let search_terms = vec![
        "[email protected]".to_string(),
        "[email protected]".to_string(),
        "[email protected]".to_string(),
    ];

    let futures: Vec<_> = search_terms
        .iter()
        .map(|term| {
            let options = SearchOptions::new("email", term).with_row_limit(50);
            client.regular_search(&options)
        })
        .collect();

    let results = futures::future::join_all(futures).await;

    let successful_count = results.iter().filter(|r| r.is_ok()).count();
    println!("Completed {} concurrent searches", successful_count);

    Ok(())
}

#[tokio::main]
async fn main() -> Result<(), Box<dyn Error>> {

    tracing_subscriber::fmt::init();


    if let Err(e) = demonstrate_basic_usage().await {
        eprintln!("Basic usage error: {}", e);
    }

    println!("
{}", "=".repeat(50));


    if let Err(e) = demonstrate_concurrent_usage().await {
        eprintln!("Concurrent usage error: {}", e);
    }

    Ok(())
}
Memory Safe: Zero-copy deserialization and compile-time safety guarantees
High Performance: Async/await with tokio runtime for maximum throughput
Type Safety: Strong typing with serde for JSON handling
Key Features

Core Features

  • • Zero-cost abstractions with compile-time optimizations
  • • Memory-safe operations with ownership system
  • • Automatic retry with exponential backoff
  • • Rate limit handling with intelligent delays
  • • Comprehensive error handling with Result types

Performance Features

  • • Tokio async runtime for high concurrency
  • • Reqwest with connection pooling
  • • Serde for zero-copy JSON deserialization
  • • Tracing for structured logging
  • • Futures for concurrent request handling
Best Practices:
  • • Always use environment variables for API keys
  • • Implement proper error handling with Result types
  • • Use async/await for concurrent operations
  • • Monitor your API usage and balance regularly
  • • Enable logging with tracing for debugging
Documentation - LeakZero | LeakZero