use crate::error::{error_for_status, extract_message, reason_phrase, Error, TransportError};
use crate::options::ConvertRequest;
use crate::result::ConversionResult;
use crate::transport::{HttpRequest, HttpResponse, Sleeper, ThreadSleeper, Transport};
use std::sync::Arc;
use std::time::Duration;
pub const DEFAULT_BASE_URL: &str = "https://api.labelzoom.com";
pub const API_KEY_ENV_VAR: &str = "LABELZOOM_API_KEY";
pub const VERSION: &str = env!("CARGO_PKG_VERSION");
const REQUEST_ID_HEADER: &str = "x-lz-request-id";
#[derive(Clone)]
pub struct LabelZoomClient {
base_url: String,
credential: Option<String>,
max_retries: u32,
user_agent: String,
transport: Arc<dyn Transport>,
sleeper: Arc<dyn Sleeper>,
jitter: bool,
}
impl std::fmt::Debug for LabelZoomClient {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("LabelZoomClient")
.field("base_url", &self.base_url)
.field("authenticated", &self.credential.is_some())
.field("max_retries", &self.max_retries)
.field("user_agent", &self.user_agent)
.finish_non_exhaustive()
}
}
pub struct ClientBuilder {
api_key: Option<Option<String>>,
base_url: String,
max_retries: u32,
timeout: Option<Duration>,
user_agent_suffix: Option<String>,
transport: Option<Arc<dyn Transport>>,
sleeper: Arc<dyn Sleeper>,
jitter: bool,
env: EnvLookup,
}
type EnvLookup = Arc<dyn Fn(&str) -> Option<String> + Send + Sync>;
impl Default for ClientBuilder {
fn default() -> Self {
Self {
api_key: None,
base_url: DEFAULT_BASE_URL.to_owned(),
max_retries: 2,
timeout: None,
user_agent_suffix: None,
transport: None,
sleeper: Arc::new(ThreadSleeper),
jitter: true,
env: Arc::new(|key| std::env::var(key).ok()),
}
}
}
impl ClientBuilder {
#[must_use]
pub fn new() -> Self {
Self::default()
}
#[must_use]
pub fn api_key(mut self, api_key: impl Into<String>) -> Self {
let api_key = api_key.into();
self.api_key = Some(if api_key.is_empty() {
None
} else {
Some(api_key)
});
self
}
#[must_use]
pub fn anonymous(mut self) -> Self {
self.api_key = Some(None);
self
}
#[must_use]
pub fn base_url(mut self, base_url: impl Into<String>) -> Self {
self.base_url = base_url.into();
self
}
#[must_use]
pub fn max_retries(mut self, max_retries: u32) -> Self {
self.max_retries = max_retries;
self
}
#[must_use]
pub fn timeout(mut self, timeout: Duration) -> Self {
self.timeout = Some(timeout);
self
}
#[must_use]
pub fn user_agent_suffix(mut self, suffix: impl Into<String>) -> Self {
self.user_agent_suffix = Some(suffix.into());
self
}
#[must_use]
pub fn transport(mut self, transport: Arc<dyn Transport>) -> Self {
self.transport = Some(transport);
self
}
#[must_use]
pub fn sleeper(mut self, sleeper: Arc<dyn Sleeper>) -> Self {
self.sleeper = sleeper;
self
}
#[must_use]
pub fn jitter(mut self, jitter: bool) -> Self {
self.jitter = jitter;
self
}
#[must_use]
pub fn env_lookup(
mut self,
lookup: impl Fn(&str) -> Option<String> + Send + Sync + 'static,
) -> Self {
self.env = Arc::new(lookup);
self
}
pub fn build(self) -> Result<LabelZoomClient, Error> {
let transport = match self.transport {
Some(transport) => transport,
#[cfg(feature = "ureq-transport")]
None => match self.timeout {
Some(timeout) => Arc::new(crate::transport::UreqTransport::with_timeout(timeout)),
None => Arc::new(crate::transport::UreqTransport::new()),
},
#[cfg(not(feature = "ureq-transport"))]
None => {
return Err(Error::validation(
"transport",
"This build has no bundled HTTP backend (the `ureq-transport` feature \
is off), so a Transport must be supplied with ClientBuilder::transport.",
))
}
};
let credential = match self.api_key {
None => (self.env)(API_KEY_ENV_VAR).filter(|value| !value.is_empty()),
Some(api_key) => api_key,
};
let mut user_agent = format!(
"labelzoom-rust-sdk/{VERSION} (rust; {} {})",
std::env::consts::OS,
std::env::consts::ARCH
);
if let Some(suffix) = self.user_agent_suffix.as_deref().map(str::trim) {
if !suffix.is_empty() {
user_agent.push(' ');
user_agent.push_str(suffix);
}
}
Ok(LabelZoomClient {
base_url: self.base_url.trim_end_matches('/').to_owned(),
credential,
max_retries: self.max_retries,
user_agent,
transport,
sleeper: self.sleeper,
jitter: self.jitter,
})
}
}
impl LabelZoomClient {
#[must_use]
pub fn new() -> Self {
Self::builder()
.build()
.expect("the default build always has a transport")
}
pub fn builder() -> ClientBuilder {
ClientBuilder::new()
}
pub fn is_authenticated(&self) -> bool {
self.credential.is_some()
}
pub fn request_url(&self, request: &ConvertRequest) -> Result<String, Error> {
let mut url = format!(
"{}/api/v2/convert/{}/to/{}",
self.base_url,
request.source.wire_token(),
request.target.wire_token()
);
if let Some(params) = request.options.serialize()? {
url.push_str("?params=");
url.push_str(&percent_encode(¶ms));
}
Ok(url)
}
pub fn convert(&self, request: &ConvertRequest) -> Result<ConversionResult, Error> {
if request.body.is_empty() {
return Err(Error::validation(
"body",
"Source body cannot be empty; the API rejects zero-length requests.",
));
}
let url = self.request_url(request)?;
let content_type = if request.base64_text {
"text/plain"
} else {
request.source.media_type()
};
let mut headers = vec![
("content-type".to_owned(), content_type.to_owned()),
("accept".to_owned(), "*/*".to_owned()),
("user-agent".to_owned(), self.user_agent.clone()),
];
if let Some(credential) = &self.credential {
headers.push(("authorization".to_owned(), format!("Bearer {credential}")));
}
let attempts = self.max_retries + 1;
for attempt in 1..=attempts {
let outgoing = HttpRequest {
method: "POST",
url: url.clone(),
headers: headers.clone(),
body: request.body.clone(),
};
let response = match self.transport.execute(outgoing) {
Ok(response) => response,
Err(error) => {
if attempt >= attempts {
return Err(Error::Transport(error));
}
self.delay(attempt, None);
continue;
}
};
if (200..300).contains(&response.status) {
return Ok(read_result(response));
}
let retry_after = retry_after_seconds(&response);
if attempt >= attempts || !is_retryable(response.status) {
return Err(Error::Api(read_error(&response, retry_after)));
}
self.delay(attempt, retry_after);
}
Err(Error::Transport(TransportError::new(
"the retry loop exited without a result",
)))
}
fn delay(&self, attempt: u32, retry_after: Option<f64>) {
let backoff = 2f64.powi(i32::try_from(attempt).unwrap_or(i32::MAX).saturating_sub(1));
let mut wait = if self.jitter {
backoff * pseudo_random()
} else {
backoff
};
if let Some(seconds) = retry_after {
if seconds > wait {
wait = seconds;
}
}
self.sleeper.sleep(wait);
}
}
impl Default for LabelZoomClient {
fn default() -> Self {
Self::new()
}
}
fn is_retryable(status: u16) -> bool {
status == 429 || status >= 500
}
fn read_result(response: HttpResponse) -> ConversionResult {
ConversionResult {
content_type: response.header("content-type").map(str::to_owned),
request_id: response.header(REQUEST_ID_HEADER).map(str::to_owned),
status: response.status,
bytes: response.body,
}
}
fn read_error(response: &HttpResponse, retry_after: Option<f64>) -> crate::error::ApiError {
let request_id = response.header(REQUEST_ID_HEADER).map(str::to_owned);
let raw_body = String::from_utf8_lossy(&response.body).into_owned();
let message = extract_message(&raw_body, reason_phrase(response.status));
error_for_status(response.status, message, request_id, raw_body, retry_after)
}
fn retry_after_seconds(response: &HttpResponse) -> Option<f64> {
response.header("retry-after")?.trim().parse::<f64>().ok()
}
fn percent_encode(value: &str) -> String {
let mut encoded = String::with_capacity(value.len());
for byte in value.bytes() {
match byte {
b'A'..=b'Z' | b'a'..=b'z' | b'0'..=b'9' | b'-' | b'_' | b'.' | b'~' => {
encoded.push(byte as char);
}
other => {
use std::fmt::Write as _;
let _ = write!(encoded, "%{other:02X}");
}
}
}
encoded
}
fn pseudo_random() -> f64 {
use std::time::{SystemTime, UNIX_EPOCH};
let nanos = SystemTime::now()
.duration_since(UNIX_EPOCH)
.map_or(0, |elapsed| elapsed.subsec_nanos());
let mixed = u64::from(nanos).wrapping_mul(6_364_136_223_846_793_005) >> 11;
f64::from(u32::try_from(mixed % 1_000_000).unwrap_or(0)) / 1_000_000.0
}