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 runCargo.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