use std::time::Duration;
use reqwest::{header, Client as ReqwestClient, Method, Url};
use crate::config::Config;
use crate::errors::{from_http, CleanLibraryError, TransportError};
use crate::types::{AuditResponse, PolicyPreviewRequest, PolicyPreviewResponse, Verdict};
const DEFAULT_TIMEOUT_SECS: u64 = 30;
#[derive(Debug, Clone)]
pub struct Client {
http: ReqwestClient,
base_url: Url,
api_key: Option<String>,
api_version: String,
}
impl Client {
pub fn from_config(config: &Config) -> Result<Self, CleanLibraryError> {
Self::build(&config.endpoint.url, config.auth.api_key.clone(), &config.endpoint.api_version)
}
pub fn new(endpoint: &str, api_key: Option<String>) -> Result<Self, CleanLibraryError> {
Self::build(endpoint, api_key, "v1")
}
fn build(
endpoint: &str,
api_key: Option<String>,
api_version: &str,
) -> Result<Self, CleanLibraryError> {
let base_url = Url::parse(endpoint)
.map_err(|e| TransportError::InvalidUrl(format!("{}: {}", endpoint, e)))?;
let is_localhost = matches!(
base_url.host_str(),
Some("localhost") | Some("127.0.0.1") | Some("::1")
);
if base_url.scheme() != "https" && !is_localhost {
return Err(TransportError::TlsRequired(endpoint.to_string()).into());
}
let http = ReqwestClient::builder()
.timeout(Duration::from_secs(DEFAULT_TIMEOUT_SECS))
.user_agent(concat!("cleanlib-cli/", env!("CARGO_PKG_VERSION")))
.build()
.map_err(TransportError::Network)?;
Ok(Self {
http,
base_url,
api_key,
api_version: api_version.to_string(),
})
}
pub async fn fetch_verdict(
&self,
ecosystem: &str,
package: &str,
version: &str,
) -> Result<Verdict, CleanLibraryError> {
let path = format!(
"{}/customer/verdicts/{}/{}/{}",
self.api_version,
urlencode(ecosystem),
urlencode(package),
urlencode(version),
);
let url = self
.base_url
.join(&path)
.map_err(|e| TransportError::InvalidUrl(format!("{}: {}", path, e)))?;
let response = self.send(Method::GET, url).await?;
let status = response.status();
let headers = response.headers().clone();
let body = response.text().await.map_err(TransportError::Network)?;
if !status.is_success() {
return Err(from_http(status.as_u16(), &headers, &body));
}
serde_json::from_str(&body)
.map_err(|e| CleanLibraryError::Parse(format!("verdict response: {}", e)))
}
pub async fn policy_preview(
&self,
req: &PolicyPreviewRequest,
) -> Result<PolicyPreviewResponse, CleanLibraryError> {
let path = format!("{}/customer/policy/preview", self.api_version);
let url = self
.base_url
.join(&path)
.map_err(|e| TransportError::InvalidUrl(format!("{}: {}", path, e)))?;
let body = serde_json::to_vec(req)
.map_err(|e| CleanLibraryError::Parse(format!("policy_preview body: {}", e)))?;
let response = self.send_with_body(Method::POST, url, body, "application/json").await?;
let status = response.status();
let headers = response.headers().clone();
let body = response.text().await.map_err(TransportError::Network)?;
if !status.is_success() {
return Err(from_http(status.as_u16(), &headers, &body));
}
serde_json::from_str(&body)
.map_err(|e| CleanLibraryError::Parse(format!("policy_preview response: {}", e)))
}
pub async fn fetch_artifact(
&self,
ecosystem: &str,
package: &str,
version: &str,
) -> Result<Vec<u8>, CleanLibraryError> {
let url = build_fetch_url(&self.base_url, ecosystem, package, version)?;
let response = self.send(Method::GET, url).await?;
let status = response.status();
let headers = response.headers().clone();
if !status.is_success() {
let body = response.text().await.map_err(TransportError::Network)?;
return Err(from_http(status.as_u16(), &headers, &body));
}
emit_decision_headers(&headers);
let bytes = response.bytes().await.map_err(TransportError::Network)?;
Ok(bytes.to_vec())
}
pub async fn fetch_artifact_stream<W>(
&self,
ecosystem: &str,
package: &str,
version: &str,
writer: &mut W,
) -> Result<u64, CleanLibraryError>
where
W: tokio::io::AsyncWrite + Unpin,
{
use futures_util::StreamExt;
use tokio::io::AsyncWriteExt;
let url = build_fetch_url(&self.base_url, ecosystem, package, version)?;
let response = self.send(Method::GET, url).await?;
let status = response.status();
let headers = response.headers().clone();
if !status.is_success() {
let body = response.text().await.map_err(TransportError::Network)?;
return Err(from_http(status.as_u16(), &headers, &body));
}
emit_decision_headers(&headers);
let mut total: u64 = 0;
let mut stream = response.bytes_stream();
while let Some(chunk) = stream.next().await {
let bytes = chunk.map_err(TransportError::Network)?;
writer
.write_all(&bytes)
.await
.map_err(|e| CleanLibraryError::Parse(format!("write artifact chunk: {}", e)))?;
total += bytes.len() as u64;
}
writer
.flush()
.await
.map_err(|e| CleanLibraryError::Parse(format!("flush artifact stream: {}", e)))?;
Ok(total)
}
pub async fn audit(
&self,
since: Option<&str>,
decision: Option<&str>,
ecosystem: Option<&str>,
) -> Result<AuditResponse, CleanLibraryError> {
let path = format!("{}/customer/audit", self.api_version);
let mut url = self
.base_url
.join(&path)
.map_err(|e| TransportError::InvalidUrl(format!("{}: {}", path, e)))?;
{
let mut q = url.query_pairs_mut();
if let Some(s) = since {
q.append_pair("since", s);
}
if let Some(d) = decision {
q.append_pair("decision", d);
}
if let Some(e) = ecosystem {
q.append_pair("ecosystem", e);
}
}
let response = self.send(Method::GET, url).await?;
let status = response.status();
let headers = response.headers().clone();
let body = response.text().await.map_err(TransportError::Network)?;
if !status.is_success() {
return Err(from_http(status.as_u16(), &headers, &body));
}
serde_json::from_str(&body)
.map_err(|e| CleanLibraryError::Parse(format!("audit response: {}", e)))
}
async fn send_with_body(
&self,
method: Method,
url: Url,
body: Vec<u8>,
content_type: &str,
) -> Result<reqwest::Response, CleanLibraryError> {
let mut req = self.http.request(method, url).body(body);
req = req.header(header::CONTENT_TYPE, content_type);
if let Some(key) = &self.api_key {
req = req.header(header::AUTHORIZATION, format!("Bearer {}", key));
}
req.send().await.map_err(|e| {
if e.is_timeout() {
CleanLibraryError::Transport(TransportError::Timeout)
} else {
CleanLibraryError::Transport(TransportError::Network(e))
}
})
}
pub async fn send(
&self,
method: Method,
url: Url,
) -> Result<reqwest::Response, CleanLibraryError> {
let mut req = self.http.request(method, url);
if let Some(key) = &self.api_key {
req = req.header(header::AUTHORIZATION, format!("Bearer {}", key));
}
req.send().await.map_err(|e| {
if e.is_timeout() {
CleanLibraryError::Transport(TransportError::Timeout)
} else {
CleanLibraryError::Transport(TransportError::Network(e))
}
})
}
pub fn base_url(&self) -> &Url {
&self.base_url
}
}
fn emit_decision_headers(headers: &reqwest::header::HeaderMap) {
if let Some(decision) = headers
.get("X-CleanLibrary-Decision")
.and_then(|v| v.to_str().ok())
{
eprintln!("# decision: {}", decision);
}
if let Some(reason) = headers
.get("X-CleanLibrary-Reason")
.and_then(|v| v.to_str().ok())
{
eprintln!("# reason: {}", reason);
}
}
fn build_fetch_url(
base: &Url,
ecosystem: &str,
package: &str,
version: &str,
) -> Result<Url, TransportError> {
let path = match ecosystem {
"npm" => {
if let Some(stripped) = package.strip_prefix('@') {
let (scope, name) = stripped.split_once('/').ok_or_else(|| {
TransportError::InvalidUrl(format!(
"npm scoped package missing /name: {}",
package
))
})?;
format!("npm/@{}/{}/-/{}-{}.tgz", scope, name, name, version)
} else {
format!("npm/{}/-/{}-{}.tgz", package, package, version)
}
}
"go" => format!("go/{}/@v/{}.zip", package, version),
"pypi" => {
format!("pypi/{}/{}-{}.tar.gz", package, package, version)
}
other => {
return Err(TransportError::InvalidUrl(format!(
"ecosystem not yet supported by `cleanlib fetch`: {} (Phase 1 Tier A: npm, pypi, go)",
other
)))
}
};
base.join(&path)
.map_err(|e| TransportError::InvalidUrl(format!("{}: {}", path, e)))
}
fn urlencode(s: &str) -> String {
let mut out = String::with_capacity(s.len());
for c in s.chars() {
match c {
'A'..='Z' | 'a'..='z' | '0'..='9' | '-' | '_' | '.' | '~' => out.push(c),
_ => {
let mut buf = [0u8; 4];
let encoded = c.encode_utf8(&mut buf);
for b in encoded.bytes() {
out.push_str(&format!("%{:02X}", b));
}
}
}
}
out
}
#[cfg(test)]
mod tests {
use super::*;
use crate::config::{Config, EndpointConfig};
fn cfg(url: &str) -> Config {
let mut c = Config::default();
c.endpoint = EndpointConfig {
url: url.to_string(),
api_version: "v1".to_string(),
};
c
}
#[test]
fn refuses_remote_plaintext() {
let err = Client::from_config(&cfg("http://cleanapp.clnstrt.dev")).unwrap_err();
assert!(matches!(
err,
CleanLibraryError::Transport(TransportError::TlsRequired(_))
));
}
#[test]
fn allows_localhost_plaintext_for_testing() {
let client = Client::new("http://localhost:8080", None).unwrap();
assert_eq!(client.base_url().host_str(), Some("localhost"));
}
#[test]
fn allows_127_loopback_plaintext() {
let client = Client::new("http://127.0.0.1:8080", None).unwrap();
assert_eq!(client.base_url().host_str(), Some("127.0.0.1"));
}
#[test]
fn accepts_https_endpoint() {
let client = Client::from_config(&cfg("https://cleanapp.clnstrt.dev")).unwrap();
assert_eq!(client.base_url().as_str(), "https://cleanapp.clnstrt.dev/");
}
#[test]
fn rejects_invalid_url() {
let err = Client::from_config(&cfg("not a url")).unwrap_err();
assert!(matches!(
err,
CleanLibraryError::Transport(TransportError::InvalidUrl(_))
));
}
#[test]
fn urlencode_npm_scoped_pkg() {
assert_eq!(urlencode("@my-org/foo"), "%40my-org%2Ffoo");
}
#[test]
fn urlencode_passes_simple() {
assert_eq!(urlencode("lodash"), "lodash");
assert_eq!(urlencode("4.17.21"), "4.17.21");
assert_eq!(urlencode("github.com/sirupsen/logrus"), "github.com%2Fsirupsen%2Flogrus");
}
#[test]
fn urlencode_handles_unicode() {
assert_eq!(urlencode("é"), "%C3%A9");
}
fn base() -> Url {
Url::parse("https://cleanapp.clnstrt.dev").unwrap()
}
#[test]
fn fetch_url_npm_bare() {
let url = build_fetch_url(&base(), "npm", "lodash", "4.17.21").unwrap();
assert_eq!(
url.as_str(),
"https://cleanapp.clnstrt.dev/npm/lodash/-/lodash-4.17.21.tgz"
);
}
#[test]
fn fetch_url_npm_scoped() {
let url = build_fetch_url(&base(), "npm", "@my-org/foo", "1.0.0").unwrap();
assert_eq!(
url.as_str(),
"https://cleanapp.clnstrt.dev/npm/@my-org/foo/-/foo-1.0.0.tgz"
);
}
#[test]
fn fetch_url_npm_scoped_malformed() {
let err = build_fetch_url(&base(), "npm", "@my-org", "1.0.0").unwrap_err();
assert!(matches!(err, TransportError::InvalidUrl(_)));
}
#[test]
fn fetch_url_go() {
let url = build_fetch_url(&base(), "go", "github.com/sirupsen/logrus", "v1.9.0").unwrap();
assert_eq!(
url.as_str(),
"https://cleanapp.clnstrt.dev/go/github.com/sirupsen/logrus/@v/v1.9.0.zip"
);
}
#[test]
fn fetch_url_pypi() {
let url = build_fetch_url(&base(), "pypi", "requests", "2.31.0").unwrap();
assert_eq!(
url.as_str(),
"https://cleanapp.clnstrt.dev/pypi/requests/requests-2.31.0.tar.gz"
);
}
#[test]
fn fetch_url_unknown_ecosystem_errors() {
let err = build_fetch_url(&base(), "maven", "junit:junit", "4.13.2").unwrap_err();
match err {
TransportError::InvalidUrl(msg) => {
assert!(msg.contains("Phase 1 Tier A"));
}
other => panic!("expected InvalidUrl, got {:?}", other),
}
}
}