use crate::activity_log::{ActionType, ActivityEntry};
use crate::audit::AuditEntry;
use crate::auth::{enforce_namespace_scope, NamespaceAuthority};
use crate::registry::{
circuit_open_response, method_not_allowed, nora_base_url, proxy_fetch, ProxyError,
};
use crate::registry_type::RegistryType;
use crate::secrets::expose_opt;
use crate::storage::Storage;
use crate::validation::validate_storage_key;
use crate::AppState;
use axum::{
body::Bytes,
extract::{Path, State},
http::{header, HeaderMap, HeaderValue, StatusCode},
response::{IntoResponse, Response},
routing::{get, put},
Extension, Router,
};
use sha2::Digest;
use std::time::Duration;
fn resolve_cargo_quarantine_config(
state: &AppState,
) -> (crate::digest_quarantine::QuarantineMode, i64) {
crate::digest_quarantine::resolve_global(
state.config.curation.cargo.quarantine.as_ref().or(state
.config
.curation
.quarantine
.as_ref()),
state
.config
.curation
.cargo
.quarantine_ttl
.as_deref()
.or(state.config.curation.quarantine_ttl.as_deref()),
)
}
pub fn routes() -> Router<AppState> {
Router::new()
.route("/cargo/index/config.json", get(index_config))
.route("/cargo/index/{*path}", get(sparse_index))
.route("/cargo/api/v1/crates/{crate_name}", get(get_metadata))
.route(
"/cargo/api/v1/crates/{crate_name}/{version}/download",
get(download),
)
.route(
"/cargo/api/v1/crates/new",
put(publish).fallback(|| async { method_not_allowed("PUT") }),
)
}
async fn index_config(State(state): State<AppState>) -> Response {
let base = nora_base_url(&state);
let auth_required = state.config.auth.enabled && !state.config.auth.anonymous_read;
let config = serde_json::json!({
"dl": format!("{}/cargo/api/v1/crates", base),
"api": format!("{}/cargo", base),
"auth-required": auth_required
});
(
StatusCode::OK,
[
(
header::CONTENT_TYPE,
HeaderValue::from_static("application/json"),
),
(
header::CACHE_CONTROL,
HeaderValue::from_static("public, max-age=300"),
),
],
serde_json::to_vec(&config).unwrap_or_default(),
)
.into_response()
}
async fn sparse_index(
State(state): State<AppState>,
headers: HeaderMap,
Path(path): Path<String>,
) -> Response {
let crate_name = match path.rsplit('/').next() {
Some(name) if !name.is_empty() => name.to_lowercase(),
_ => return StatusCode::NOT_FOUND.into_response(),
};
if !is_valid_crate_name(&crate_name) {
return StatusCode::BAD_REQUEST.into_response();
}
let expected_prefix = crate_index_prefix(&crate_name);
if path.to_lowercase() != format!("{}/{}", expected_prefix, crate_name) {
return StatusCode::NOT_FOUND.into_response();
}
let index_key = format!("cargo/index/{}/{}", expected_prefix, crate_name);
let cached = state.storage.get(&index_key).await.ok();
if let Some(ref data) = cached {
let modified = state.storage.stat(&index_key).await.map(|m| m.modified);
if crate::cache_ttl::mutable_ref_fresh(
state.config.cargo.proxy.is_some(),
state.config.cargo.metadata_ttl,
modified,
) {
state.metrics.record_download("cargo");
state.metrics.record_cache_hit("cargo");
state.activity.push(ActivityEntry::new(
ActionType::CacheHit,
crate_name.to_string(),
crate::registry_type::RegistryType::Cargo,
"CACHE",
));
return sparse_index_conditional(data.to_vec(), &headers);
}
}
if crate::curation::is_internal_namespace(
&state.curation().curation_engine,
crate::curation::RegistryType::Cargo,
&crate_name,
) {
if let Some(ref data) = cached {
state.metrics.record_cache_hit("cargo");
return sparse_index_conditional(data.to_vec(), &headers);
}
return crate::curation::check_namespace_isolation(
&state.curation().curation_engine,
crate::curation::RegistryType::Cargo,
&crate_name,
)
.unwrap_or_else(|| StatusCode::NOT_FOUND.into_response());
}
let proxy_url = match &state.config.cargo.proxy {
Some(url) => url.clone(),
None => {
if let Some(ref data) = cached {
state.metrics.record_cache_hit("cargo");
return sparse_index_conditional(data.to_vec(), &headers);
}
return StatusCode::NOT_FOUND.into_response();
}
};
let upstream_index_url = if proxy_url.contains("crates.io") {
format!("https://index.crates.io/{}/{}", expected_prefix, crate_name)
} else {
format!(
"{}/index/{}/{}",
proxy_url.trim_end_matches('/'),
expected_prefix,
crate_name
)
};
match proxy_fetch(
&state.http_client,
&upstream_index_url,
Duration::from_secs(state.config.cargo.proxy_timeout),
expose_opt(&state.config.cargo.proxy_auth),
&state.circuit_breaker,
RegistryType::Cargo,
)
.await
{
Ok(data) => {
state.metrics.record_download("cargo");
state.metrics.record_cache_miss("cargo");
state.activity.push(ActivityEntry::new(
ActionType::ProxyFetch,
crate_name.to_string(),
crate::registry_type::RegistryType::Cargo,
"PROXY",
));
state
.audit
.log(AuditEntry::new("proxy_fetch", "api", "", "cargo", ""));
state.spawn_cache("cargo", index_key, Bytes::from(data.clone()));
sparse_index_conditional(data, &headers)
}
Err(e) => {
if let Some(ref data) = cached {
tracing::warn!(
crate_name,
"Cargo upstream failed, serving stale cached index"
);
let mut response = sparse_index_conditional(data.to_vec(), &headers);
response.headers_mut().insert(
axum::http::header::HeaderName::from_static("x-nora-stale"),
axum::http::header::HeaderValue::from_static("true"),
);
return response;
}
match e {
ProxyError::CircuitOpen(reg) => circuit_open_response(®),
ProxyError::NotFound => StatusCode::NOT_FOUND.into_response(),
_ => {
tracing::debug!(
crate_name,
error = ?e,
"Cargo sparse index upstream error"
);
StatusCode::NOT_FOUND.into_response()
}
}
}
}
}
async fn get_metadata(State(state): State<AppState>, Path(crate_name): Path<String>) -> Response {
if validate_storage_key(&crate_name).is_err() {
return StatusCode::BAD_REQUEST.into_response();
}
let crate_name = crate_name.to_lowercase();
let key = format!("cargo/{}/metadata.json", crate_name);
if let Ok(data) = state.storage.get(&key).await {
return (StatusCode::OK, data).into_response();
}
if let Some(response) = crate::curation::check_namespace_isolation(
&state.curation().curation_engine,
crate::curation::RegistryType::Cargo,
&crate_name,
) {
return response;
}
let proxy_url = match &state.config.cargo.proxy {
Some(url) => url.clone(),
None => return StatusCode::NOT_FOUND.into_response(),
};
let url = format!(
"{}/api/v1/crates/{}",
proxy_url.trim_end_matches('/'),
crate_name
);
match proxy_fetch(
&state.http_client,
&url,
Duration::from_secs(state.config.cargo.proxy_timeout),
expose_opt(&state.config.cargo.proxy_auth),
&state.circuit_breaker,
RegistryType::Cargo,
)
.await
{
Ok(data) => {
state.spawn_cache("cargo", key.clone(), Bytes::from(data.clone()));
(StatusCode::OK, data).into_response()
}
Err(ProxyError::CircuitOpen(reg)) => circuit_open_response(®),
Err(e) => {
tracing::debug!(error = ?e, crate_name = %crate_name, "Cargo metadata proxy fetch failed");
tracing::warn!(registry = "cargo", crate_name = %crate_name, "Proxy failed, returning 404");
StatusCode::NOT_FOUND.into_response()
}
}
}
async fn ensure_cargo_metadata_cached(state: &AppState, crate_name: &str) {
let key = format!("cargo/{}/metadata.json", crate_name);
if state.storage.get(&key).await.is_ok() {
return;
}
if crate::curation::is_internal_namespace(
&state.curation().curation_engine,
crate::curation::RegistryType::Cargo,
crate_name,
) {
return;
}
let proxy_url = match &state.config.cargo.proxy {
Some(url) => url.clone(),
None => return,
};
let url = format!(
"{}/api/v1/crates/{}",
proxy_url.trim_end_matches('/'),
crate_name
);
if let Ok(data) = proxy_fetch(
&state.http_client,
&url,
Duration::from_secs(state.config.cargo.proxy_timeout),
expose_opt(&state.config.cargo.proxy_auth),
&state.circuit_breaker,
RegistryType::Cargo,
)
.await
{
let _ = state.storage.put(&key, &data).await;
}
}
async fn download(
State(state): State<AppState>,
headers: axum::http::HeaderMap,
Path((crate_name, version)): Path<(String, String)>,
) -> Response {
if validate_storage_key(&crate_name).is_err() || validate_storage_key(&version).is_err() {
return StatusCode::BAD_REQUEST.into_response();
}
let crate_name = crate_name.to_lowercase();
let publish_date = {
let meta_key = format!("cargo/{}/metadata.json", crate_name);
if state.config.server.trust_upstream_dates {
ensure_cargo_metadata_cached(&state, &crate_name).await;
}
extract_cargo_publish_date(
&state.storage,
&meta_key,
&version,
state.config.server.trust_upstream_dates,
)
.await
};
let internal = crate::curation::is_internal_namespace(
&state.curation().curation_engine,
crate::curation::RegistryType::Cargo,
&crate_name,
);
if !internal {
if let Some(response) = crate::curation::check_download(
&state.curation().curation_engine,
state.bypass_token().as_deref(),
&headers,
crate::curation::RegistryType::Cargo,
&crate_name,
Some(&version),
publish_date,
) {
return response;
}
}
let key = format!(
"cargo/{}/{}/{}-{}.crate",
crate_name, version, crate_name, version
);
if let Ok(outcome) = state.storage.get_verified(&key).await {
use nora_registry::verified::{verified_body, GateOutcome};
let data = match outcome {
GateOutcome::Verified(blob) => verified_body(blob),
GateOutcome::Unpinned(blob) => blob.into_inner(),
};
if let Some(response) = crate::curation::verify_integrity(
&state.curation().curation_engine,
crate::curation::RegistryType::Cargo,
&crate_name,
Some(&version),
&data,
) {
return response;
}
state.metrics.record_download("cargo");
state.metrics.record_cache_hit("cargo");
state.activity.push(ActivityEntry::new(
ActionType::Pull,
format!("{}@{}", crate_name, version),
crate::registry_type::RegistryType::Cargo,
"LOCAL",
));
state
.audit
.log(AuditEntry::new("pull", "api", "", "cargo", ""));
let (q_mode, q_secs) = resolve_cargo_quarantine_config(&state);
if let Some(resp) = crate::digest_quarantine::proxy_gate_dated(
&state.digest_store,
"cargo",
&data,
&q_mode,
q_secs,
"cache",
publish_date,
) {
return resp;
}
if let Some(response) = crate::registry::range::range_response(
&state.storage,
&[&key],
&headers,
data.len() as u64,
"application/x-tar",
&[(
header::CACHE_CONTROL,
"public, max-age=31536000, immutable".to_string(),
)],
)
.await
{
return response;
}
return (
StatusCode::OK,
[
(
header::CONTENT_TYPE,
HeaderValue::from_static("application/x-tar"),
),
(
header::CACHE_CONTROL,
HeaderValue::from_static("public, max-age=31536000, immutable"),
),
(header::ACCEPT_RANGES, HeaderValue::from_static("bytes")),
],
data,
)
.into_response();
}
if internal {
return crate::curation::check_namespace_isolation(
&state.curation().curation_engine,
crate::curation::RegistryType::Cargo,
&crate_name,
)
.unwrap_or_else(|| StatusCode::NOT_FOUND.into_response());
}
let proxy_url = match &state.config.cargo.proxy {
Some(url) => url.clone(),
None => return StatusCode::NOT_FOUND.into_response(),
};
let url = format!(
"{}/api/v1/crates/{}/{}/download",
proxy_url.trim_end_matches('/'),
crate_name,
version
);
match proxy_fetch(
&state.http_client,
&url,
Duration::from_secs(state.config.cargo.proxy_timeout),
expose_opt(&state.config.cargo.proxy_auth),
&state.circuit_breaker,
RegistryType::Cargo,
)
.await
{
Ok(data) => {
state.spawn_cache("cargo", key.clone(), Bytes::from(data.clone()));
state.metrics.record_download("cargo");
state.metrics.record_cache_miss("cargo");
state.activity.push(ActivityEntry::new(
ActionType::Pull,
format!("{}@{}", crate_name, version),
crate::registry_type::RegistryType::Cargo,
"PROXY",
));
state
.audit
.log(AuditEntry::new("proxy_fetch", "api", "", "cargo", ""));
let (q_mode, q_secs) = resolve_cargo_quarantine_config(&state);
if let Some(resp) = crate::digest_quarantine::proxy_gate_dated(
&state.digest_store,
"cargo",
&data,
&q_mode,
q_secs,
&url,
publish_date,
) {
return resp;
}
(
StatusCode::OK,
[
(
header::CONTENT_TYPE,
HeaderValue::from_static("application/x-tar"),
),
(
header::CACHE_CONTROL,
HeaderValue::from_static("public, max-age=31536000, immutable"),
),
],
data,
)
.into_response()
}
Err(ProxyError::CircuitOpen(reg)) => circuit_open_response(®),
Err(e) => {
tracing::debug!(error = ?e, crate_name = %crate_name, version = %version, "Cargo crate proxy fetch failed");
tracing::warn!(registry = "cargo", crate_name = %crate_name, version = %version, "Proxy failed, returning 404");
StatusCode::NOT_FOUND.into_response()
}
}
}
async fn publish(
State(state): State<AppState>,
Extension(authority): Extension<NamespaceAuthority>,
body: Bytes,
) -> Response {
if body.len() < 8 {
return (StatusCode::BAD_REQUEST, "Payload too small").into_response();
}
let metadata_len = u32::from_le_bytes([body[0], body[1], body[2], body[3]]) as usize;
if body.len() < 4 + metadata_len + 4 {
return (StatusCode::BAD_REQUEST, "Truncated metadata").into_response();
}
let metadata_bytes = &body[4..4 + metadata_len];
let metadata: serde_json::Value = match serde_json::from_slice(metadata_bytes) {
Ok(v) => v,
Err(e) => {
return (
StatusCode::BAD_REQUEST,
format!("Invalid metadata JSON: {}", e),
)
.into_response()
}
};
let crate_len_offset = 4 + metadata_len;
let crate_len = u32::from_le_bytes([
body[crate_len_offset],
body[crate_len_offset + 1],
body[crate_len_offset + 2],
body[crate_len_offset + 3],
]) as usize;
let crate_start = crate_len_offset + 4;
if body.len() < crate_start + crate_len {
return (StatusCode::BAD_REQUEST, "Truncated crate tarball").into_response();
}
let crate_data = &body[crate_start..crate_start + crate_len];
let name = match metadata.get("name").and_then(|n| n.as_str()) {
Some(n) => n,
None => return (StatusCode::BAD_REQUEST, "Missing crate name").into_response(),
};
let vers = match metadata.get("vers").and_then(|v| v.as_str()) {
Some(v) => v,
None => return (StatusCode::BAD_REQUEST, "Missing crate version").into_response(),
};
if !is_valid_crate_name(name) {
return (StatusCode::BAD_REQUEST, "Invalid crate name").into_response();
}
if validate_storage_key(vers).is_err() {
return (StatusCode::BAD_REQUEST, "Invalid version").into_response();
}
let name = name.to_lowercase();
let vers = vers.to_string();
if enforce_namespace_scope(&authority, &name).is_err() {
return (StatusCode::FORBIDDEN, "Outside namespace scope").into_response();
}
let crate_key = format!("cargo/{}/{}/{}-{}.crate", name, vers, name, vers);
let prefix = crate_index_prefix(&name);
let index_lock_key = format!("cargo/index/{}/{}", prefix, name);
let lock = state.publish_lock(&index_lock_key);
let _guard = lock.lock().await;
if state.storage.stat(&crate_key).await.is_some() {
let err = serde_json::json!({
"errors": [{"detail": format!("crate version `{}@{}` already exists", name, vers)}]
});
return (
StatusCode::CONFLICT,
[(
header::CONTENT_TYPE,
HeaderValue::from_static("application/json"),
)],
serde_json::to_vec(&err).unwrap_or_default(),
)
.into_response();
}
let cksum = hex::encode(sha2::Sha256::digest(crate_data));
let deps = metadata
.get("deps")
.and_then(|d| d.as_array())
.map(|arr| {
arr.iter()
.map(|dep| {
let mut d = dep.clone();
if let Some(obj) = d.as_object_mut() {
if let Some(vr) = obj.remove("version_req") {
obj.insert("req".to_string(), vr);
}
if let Some(ent) = obj.remove("explicit_name_in_toml") {
if !ent.is_null() {
obj.insert("package".to_string(), ent);
}
}
}
d
})
.collect::<Vec<_>>()
})
.map(serde_json::Value::Array)
.unwrap_or(serde_json::json!([]));
let features = metadata
.get("features")
.cloned()
.unwrap_or_else(|| serde_json::json!({}));
let features2 = metadata.get("features2").cloned();
let links = metadata.get("links").cloned();
let mut index_entry = serde_json::json!({
"name": name,
"vers": vers,
"deps": deps,
"cksum": cksum,
"features": features,
"yanked": false,
});
if let Some(f2) = features2 {
index_entry["features2"] = f2;
}
if let Some(l) = links {
index_entry["links"] = l;
}
let entry_line = serde_json::to_string(&index_entry).unwrap_or_default();
if state.storage.put(&crate_key, crate_data).await.is_err() {
return StatusCode::INTERNAL_SERVER_ERROR.into_response();
}
migrate_cargo_index(&state, &prefix, &name).await;
let entry_key = format!("cargo/index-entries/{}/{}/{}.json", prefix, name, vers);
if state
.storage
.put(&entry_key, entry_line.as_bytes())
.await
.is_err()
{
let _ = state.storage.delete(&crate_key).await;
return StatusCode::INTERNAL_SERVER_ERROR.into_response();
}
if regenerate_cargo_index(&state.storage, &prefix, &name)
.await
.is_err()
{
tracing::error!(crate_name = %name, "cargo publish: failed to regenerate sparse index");
return StatusCode::INTERNAL_SERVER_ERROR.into_response();
}
state.metrics.record_upload("cargo");
let artifact = format!("{}@{}", name, vers);
state
.audit
.log(AuditEntry::new("push", "api", &artifact, "cargo", ""));
state.activity.push(ActivityEntry::new(
ActionType::Push,
artifact,
crate::registry_type::RegistryType::Cargo,
"LOCAL",
));
state.repo_index.invalidate("cargo");
let response = serde_json::json!({
"warnings": {
"invalid_categories": [],
"invalid_badges": [],
"other": []
}
});
(
StatusCode::OK,
[(
header::CONTENT_TYPE,
HeaderValue::from_static("application/json"),
)],
serde_json::to_vec(&response).unwrap_or_default(),
)
.into_response()
}
async fn regenerate_cargo_index(storage: &Storage, prefix: &str, name: &str) -> Result<(), ()> {
let entries_prefix = format!("cargo/index-entries/{}/{}/", prefix, name);
let mut keys = storage.list(&entries_prefix).await.map_err(|_| ())?;
keys.sort(); let mut index: Vec<u8> = Vec::new();
for key in &keys {
let line = storage.get(key).await.map_err(|_| ())?;
index.extend_from_slice(&line);
if !line.ends_with(b"\n") {
index.push(b'\n');
}
}
let index_key = format!("cargo/index/{}/{}", prefix, name);
storage.put(&index_key, &index).await.map_err(|_| ())
}
async fn migrate_cargo_index(state: &AppState, prefix: &str, name: &str) {
let entries_prefix = format!("cargo/index-entries/{}/{}/", prefix, name);
if !state
.storage
.list(&entries_prefix)
.await
.unwrap_or_default()
.is_empty()
{
return; }
let index_key = format!("cargo/index/{}/{}", prefix, name);
let Ok(data) = state.storage.get(&index_key).await else {
return;
};
for line in data.split(|&b| b == b'\n') {
if line.is_empty() {
continue;
}
if let Ok(entry) = serde_json::from_slice::<serde_json::Value>(line) {
if let Some(vers) = entry.get("vers").and_then(|v| v.as_str()) {
let key = format!("cargo/index-entries/{}/{}/{}.json", prefix, name, vers);
let _ = state.storage.put(&key, line).await;
}
}
}
}
fn crate_index_prefix(name: &str) -> String {
let lower = name.to_lowercase();
match lower.len() {
1 => "1".to_string(),
2 => "2".to_string(),
3 => format!("3/{}", &lower[..1]),
_ => format!("{}/{}", &lower[..2], &lower[2..4]),
}
}
fn is_valid_crate_name(name: &str) -> bool {
if name.is_empty() || name.len() > 64 {
return false;
}
let first = name.chars().next().unwrap_or('\0');
if !first.is_ascii_alphanumeric() {
return false;
}
name.chars()
.all(|c| c.is_ascii_alphanumeric() || c == '-' || c == '_')
}
fn sparse_index_conditional(data: Vec<u8>, req_headers: &HeaderMap) -> Response {
let etag_hash = hex::encode(sha2::Sha256::digest(&data));
let etag_val = format!("\"{}\"", etag_hash);
if let Some(inm) = req_headers
.get(header::IF_NONE_MATCH)
.and_then(|v| v.to_str().ok())
{
if inm.trim() == etag_val || inm.trim() == "*" {
return (
StatusCode::NOT_MODIFIED,
[
(header::ETAG, etag_val),
(header::CACHE_CONTROL, "public, max-age=300".to_string()),
],
)
.into_response();
}
}
(
StatusCode::OK,
[
(header::CONTENT_TYPE, "application/json".to_string()),
(header::CACHE_CONTROL, "public, max-age=300".to_string()),
(header::ETAG, etag_val),
],
data,
)
.into_response()
}
async fn extract_cargo_publish_date(
storage: &crate::storage::Storage,
metadata_key: &str,
version: &str,
trust_upstream: bool,
) -> Option<i64> {
if !trust_upstream {
return crate::curation::extract_mtime_as_publish_date(storage, metadata_key).await;
}
let data = storage.get(metadata_key).await.ok()?;
let json: serde_json::Value = serde_json::from_slice(&data).ok()?;
let versions = json.get("versions")?.as_array()?;
for v in versions {
if v.get("num")?.as_str()? == version {
let date_str = v.get("created_at")?.as_str()?;
return crate::curation::parse_iso8601_to_unix(date_str);
}
}
None
}
#[cfg(test)]
#[allow(clippy::unwrap_used)]
mod tests {
use super::*;
#[test]
fn test_prefix_single_char() {
assert_eq!(crate_index_prefix("a"), "1");
assert_eq!(crate_index_prefix("Z"), "1");
}
#[test]
fn test_prefix_two_chars() {
assert_eq!(crate_index_prefix("ab"), "2");
assert_eq!(crate_index_prefix("IO"), "2");
}
#[test]
fn test_prefix_three_chars() {
assert_eq!(crate_index_prefix("abc"), "3/a");
assert_eq!(crate_index_prefix("Foo"), "3/f");
}
#[test]
fn test_prefix_four_plus_chars() {
assert_eq!(crate_index_prefix("serde"), "se/rd");
assert_eq!(crate_index_prefix("tokio"), "to/ki");
assert_eq!(crate_index_prefix("Axum"), "ax/um");
assert_eq!(crate_index_prefix("ab_cd_ef"), "ab/_c");
}
#[test]
fn test_valid_crate_names() {
assert!(is_valid_crate_name("serde"));
assert!(is_valid_crate_name("my-crate"));
assert!(is_valid_crate_name("my_crate"));
assert!(is_valid_crate_name("a"));
assert!(is_valid_crate_name("crate123"));
}
#[test]
fn test_invalid_crate_names() {
assert!(!is_valid_crate_name(""));
assert!(!is_valid_crate_name("-start"));
assert!(!is_valid_crate_name("_start"));
assert!(!is_valid_crate_name("has space"));
assert!(!is_valid_crate_name("has/slash"));
assert!(!is_valid_crate_name("has..dots"));
assert!(!is_valid_crate_name("has\\backslash"));
assert!(!is_valid_crate_name(&"a".repeat(65)));
}
#[test]
fn test_crate_name_max_length() {
assert!(is_valid_crate_name(&"a".repeat(64)));
assert!(!is_valid_crate_name(&"a".repeat(65)));
}
#[tokio::test]
async fn test_regenerate_cargo_index_aborts_on_read_error() {
use crate::storage::{Result as StorageResult, StorageBackend, StorageError};
use axum::body::Bytes;
use std::path::Path;
use std::pin::Pin;
use tokio::io::AsyncRead;
struct FailingGetBackend;
#[async_trait::async_trait]
impl StorageBackend for FailingGetBackend {
async fn list(&self, _prefix: &str) -> StorageResult<Vec<String>> {
Ok(vec![
"cargo/index-entries/fa/il/failcrate/0.1.0.json".to_string(),
"cargo/index-entries/fa/il/failcrate/0.2.0.json".to_string(),
])
}
async fn get(&self, _key: &str) -> StorageResult<Bytes> {
Err(StorageError::Io(std::io::Error::other(
"injected transient read error",
)))
}
async fn put(&self, _key: &str, _data: &[u8]) -> StorageResult<()> {
panic!("regenerate must abort before writing a truncated index");
}
async fn delete(&self, _key: &str) -> StorageResult<()> {
Ok(())
}
async fn stat(&self, _key: &str) -> Option<crate::storage::FileMeta> {
None
}
async fn health_check(&self) -> bool {
true
}
async fn total_size(&self) -> u64 {
0
}
fn backend_name(&self) -> &'static str {
"failing-get-test"
}
async fn put_from_path(&self, _key: &str, _src: &Path) -> StorageResult<()> {
Ok(())
}
async fn get_reader(
&self,
_key: &str,
) -> StorageResult<(u64, Pin<Box<dyn AsyncRead + Send + Unpin>>)> {
Err(StorageError::NotFound)
}
async fn copy(&self, _src: &str, _dst: &str) -> StorageResult<()> {
Err(StorageError::NotFound)
}
}
let storage = Storage::from_backend(std::sync::Arc::new(FailingGetBackend));
let result = regenerate_cargo_index(&storage, "fa/il", "failcrate").await;
assert!(
result.is_err(),
"a transient read error must abort the rebuild, not publish a truncated index"
);
}
}
#[cfg(test)]
#[allow(clippy::unwrap_used)]
mod integration_tests {
use crate::test_helpers::{body_bytes, create_test_context, send};
#[tokio::test]
async fn test_cargo_namespace_scope_enforced() {
use crate::auth::NamespaceAuthority;
use crate::config::ScopeEnforcement;
use axum::body::Bytes;
use axum::extract::State;
use axum::Extension;
let ctx = create_test_context();
let scoped = NamespaceAuthority::from_oidc_scope(
"ci",
&["acme/**".to_string()],
ScopeEnforcement::Enforce,
);
let metadata =
serde_json::json!({"name":"other-crate","vers":"0.1.0","deps":[],"features":{}});
let payload = build_publish_payload(&metadata, b"data");
let resp = super::publish(
State(ctx.state.clone()),
Extension(scoped.clone()),
Bytes::from(payload),
)
.await;
assert_eq!(resp.status(), StatusCode::FORBIDDEN);
let metadata = serde_json::json!({"name":"acme","vers":"0.1.0","deps":[],"features":{}});
let payload = build_publish_payload(&metadata, b"data");
let resp = super::publish(
State(ctx.state.clone()),
Extension(scoped),
Bytes::from(payload),
)
.await;
assert_ne!(resp.status(), StatusCode::FORBIDDEN);
}
use axum::body::Body;
use axum::http::{Method, StatusCode};
#[tokio::test]
async fn test_cargo_index_config() {
let ctx = create_test_context();
let resp = send(&ctx.app, Method::GET, "/cargo/index/config.json", "").await;
assert_eq!(resp.status(), StatusCode::OK);
let body = body_bytes(resp).await;
let json: serde_json::Value = serde_json::from_slice(&body).unwrap();
let dl = json["dl"].as_str().unwrap();
let api = json["api"].as_str().unwrap();
assert!(
dl.ends_with("/cargo/api/v1/crates"),
"dl must end with /cargo/api/v1/crates, got: {}",
dl
);
assert!(
api.ends_with("/cargo"),
"api must end with /cargo so Cargo can build {{api}}/api/v1/crates/... routes, got: {}",
api
);
assert_eq!(json["auth-required"], serde_json::Value::Bool(false));
}
#[tokio::test]
async fn test_cargo_index_config_auth_required() {
use crate::test_helpers::{
create_test_context_with_anonymous_read, create_test_context_with_auth,
send_with_headers,
};
use base64::{engine::general_purpose::STANDARD, Engine as _};
let basic = format!("Basic {}", STANDARD.encode("alice:hunter2"));
let ctx = create_test_context_with_auth(&[("alice", "hunter2")]);
let resp = send_with_headers(
&ctx.app,
Method::GET,
"/cargo/index/config.json",
vec![("authorization", &basic)],
Vec::new(),
)
.await;
assert_eq!(resp.status(), StatusCode::OK);
let body = body_bytes(resp).await;
let json: serde_json::Value = serde_json::from_slice(&body).unwrap();
assert_eq!(json["auth-required"], serde_json::Value::Bool(true));
let ctx = create_test_context_with_anonymous_read(&[("alice", "hunter2")]);
let resp = send(&ctx.app, Method::GET, "/cargo/index/config.json", "").await;
assert_eq!(resp.status(), StatusCode::OK);
let body = body_bytes(resp).await;
let json: serde_json::Value = serde_json::from_slice(&body).unwrap();
assert_eq!(
json["auth-required"],
serde_json::Value::Bool(false),
"anonymous_read serves index reads without creds"
);
}
#[tokio::test]
async fn test_cargo_config_api_routes_to_metadata() {
let ctx = create_test_context();
let meta = r#"{"name":"serde","versions":[]}"#;
ctx.state
.storage
.put("cargo/serde/metadata.json", meta.as_bytes())
.await
.unwrap();
let config_resp = send(&ctx.app, Method::GET, "/cargo/index/config.json", "").await;
let config_body = body_bytes(config_resp).await;
let config: serde_json::Value = serde_json::from_slice(&config_body).unwrap();
let api_url = config["api"].as_str().unwrap();
let metadata_path = format!(
"{}/api/v1/crates/serde",
api_url.trim_start_matches("http://localhost")
);
let resp = send(&ctx.app, Method::GET, &metadata_path, "").await;
assert_eq!(
resp.status(),
StatusCode::OK,
"metadata request via config.json api field must not 404 (#442)"
);
}
#[tokio::test]
async fn test_cargo_config_api_routes_to_publish() {
let ctx = create_test_context();
let config_resp = send(&ctx.app, Method::GET, "/cargo/index/config.json", "").await;
let config_body = body_bytes(config_resp).await;
let config: serde_json::Value = serde_json::from_slice(&config_body).unwrap();
let api_url = config["api"].as_str().unwrap();
let publish_path = format!(
"{}/api/v1/crates/new",
api_url.trim_start_matches("http://localhost")
);
let metadata = serde_json::json!({
"name": "publish-from-config",
"vers": "0.1.0",
"deps": [],
"features": {},
});
let payload = build_publish_payload(&metadata, b"crate-data");
let resp = send(&ctx.app, Method::PUT, &publish_path, Body::from(payload)).await;
assert_eq!(
resp.status(),
StatusCode::OK,
"publish request via config.json api field must not 404"
);
}
#[tokio::test]
async fn test_cargo_sparse_index_from_storage() {
let ctx = create_test_context();
let index_data = br#"{"name":"serde","vers":"1.0.0","deps":[],"cksum":"abc123","features":{},"yanked":false}"#;
ctx.state
.storage
.put("cargo/index/se/rd/serde", index_data)
.await
.unwrap();
let resp = send(&ctx.app, Method::GET, "/cargo/index/se/rd/serde", "").await;
assert_eq!(resp.status(), StatusCode::OK);
let body = body_bytes(resp).await;
assert_eq!(&body[..], index_data);
}
#[tokio::test]
async fn test_cargo_sparse_index_wrong_prefix() {
let ctx = create_test_context();
let resp = send(&ctx.app, Method::GET, "/cargo/index/1/serde", "").await;
assert_eq!(resp.status(), StatusCode::NOT_FOUND);
}
#[tokio::test]
async fn test_cargo_sparse_index_single_char() {
let ctx = create_test_context();
ctx.state
.storage
.put("cargo/index/1/a", b"index-data")
.await
.unwrap();
let resp = send(&ctx.app, Method::GET, "/cargo/index/1/a", "").await;
assert_eq!(resp.status(), StatusCode::OK);
}
#[tokio::test]
async fn test_cargo_sparse_index_two_char() {
let ctx = create_test_context();
ctx.state
.storage
.put("cargo/index/2/ab", b"index-data")
.await
.unwrap();
let resp = send(&ctx.app, Method::GET, "/cargo/index/2/ab", "").await;
assert_eq!(resp.status(), StatusCode::OK);
}
#[tokio::test]
async fn test_cargo_sparse_index_three_char() {
let ctx = create_test_context();
ctx.state
.storage
.put("cargo/index/3/f/foo", b"index-data")
.await
.unwrap();
let resp = send(&ctx.app, Method::GET, "/cargo/index/3/f/foo", "").await;
assert_eq!(resp.status(), StatusCode::OK);
}
#[tokio::test]
async fn test_cargo_sparse_index_returns_etag() {
use crate::test_helpers::send_with_headers;
let ctx = create_test_context();
let index_data = br#"{"name":"serde","vers":"1.0.0","deps":[]}"#;
ctx.state
.storage
.put("cargo/index/se/rd/serde", index_data)
.await
.unwrap();
let resp = send(&ctx.app, Method::GET, "/cargo/index/se/rd/serde", "").await;
assert_eq!(resp.status(), StatusCode::OK);
let etag = resp
.headers()
.get("etag")
.expect("response must include ETag")
.to_str()
.unwrap()
.to_string();
assert!(
etag.starts_with('"') && etag.ends_with('"'),
"ETag must be quoted"
);
assert!(etag.len() > 2, "ETag must contain a hash");
let resp2 = send_with_headers(
&ctx.app,
Method::GET,
"/cargo/index/se/rd/serde",
vec![("if-none-match", &etag)],
"",
)
.await;
assert_eq!(resp2.status(), StatusCode::NOT_MODIFIED);
let etag2 = resp2.headers().get("etag").expect("304 must include ETag");
assert_eq!(etag2.to_str().unwrap(), etag);
}
#[tokio::test]
async fn test_cargo_sparse_index_etag_mismatch_returns_200() {
use crate::test_helpers::send_with_headers;
let ctx = create_test_context();
let index_data = br#"{"name":"serde","vers":"1.0.0","deps":[]}"#;
ctx.state
.storage
.put("cargo/index/se/rd/serde", index_data)
.await
.unwrap();
let resp = send_with_headers(
&ctx.app,
Method::GET,
"/cargo/index/se/rd/serde",
vec![("if-none-match", "\"stale-etag\"")],
"",
)
.await;
assert_eq!(resp.status(), StatusCode::OK);
assert!(resp.headers().get("etag").is_some());
}
#[tokio::test]
async fn test_cargo_sparse_index_etag_changes_with_data() {
let ctx = create_test_context();
ctx.state
.storage
.put("cargo/index/se/rd/serde", br#"{"vers":"1.0.0"}"#)
.await
.unwrap();
let resp1 = send(&ctx.app, Method::GET, "/cargo/index/se/rd/serde", "").await;
let etag1 = resp1
.headers()
.get("etag")
.unwrap()
.to_str()
.unwrap()
.to_string();
ctx.state
.storage
.put("cargo/index/se/rd/serde", br#"{"vers":"2.0.0"}"#)
.await
.unwrap();
let resp2 = send(&ctx.app, Method::GET, "/cargo/index/se/rd/serde", "").await;
let etag2 = resp2
.headers()
.get("etag")
.unwrap()
.to_str()
.unwrap()
.to_string();
assert_ne!(etag1, etag2, "ETag must change when data changes");
}
#[tokio::test]
async fn test_cargo_sparse_index_not_found_no_proxy() {
let ctx = create_test_context();
let resp = send(&ctx.app, Method::GET, "/cargo/index/se/rd/serde", "").await;
assert_eq!(resp.status(), StatusCode::NOT_FOUND);
}
#[tokio::test]
async fn test_cargo_metadata_not_found() {
let ctx = create_test_context();
let resp = send(
&ctx.app,
Method::GET,
"/cargo/api/v1/crates/nonexistent",
"",
)
.await;
assert_eq!(resp.status(), StatusCode::NOT_FOUND);
}
#[tokio::test]
async fn test_cargo_metadata_from_storage() {
let ctx = create_test_context();
let meta = r#"{"name":"test-crate","versions":[]}"#;
ctx.state
.storage
.put("cargo/test-crate/metadata.json", meta.as_bytes())
.await
.unwrap();
let resp = send(&ctx.app, Method::GET, "/cargo/api/v1/crates/test-crate", "").await;
assert_eq!(resp.status(), StatusCode::OK);
let body = body_bytes(resp).await;
assert_eq!(&body[..], meta.as_bytes());
}
#[tokio::test]
async fn test_cargo_download_not_found() {
let ctx = create_test_context();
let resp = send(
&ctx.app,
Method::GET,
"/cargo/api/v1/crates/missing/1.0.0/download",
"",
)
.await;
assert_eq!(resp.status(), StatusCode::NOT_FOUND);
}
#[tokio::test]
async fn test_cargo_download_from_storage() {
let ctx = create_test_context();
ctx.state
.storage
.put("cargo/my-crate/1.2.3/my-crate-1.2.3.crate", b"crate-data")
.await
.unwrap();
let resp = send(
&ctx.app,
Method::GET,
"/cargo/api/v1/crates/my-crate/1.2.3/download",
"",
)
.await;
assert_eq!(resp.status(), StatusCode::OK);
let body = body_bytes(resp).await;
assert_eq!(&body[..], b"crate-data");
}
fn build_publish_payload(metadata: &serde_json::Value, crate_data: &[u8]) -> Vec<u8> {
let meta_bytes = serde_json::to_vec(metadata).unwrap();
let meta_len = (meta_bytes.len() as u32).to_le_bytes();
let crate_len = (crate_data.len() as u32).to_le_bytes();
let mut payload = Vec::new();
payload.extend_from_slice(&meta_len);
payload.extend_from_slice(&meta_bytes);
payload.extend_from_slice(&crate_len);
payload.extend_from_slice(crate_data);
payload
}
#[tokio::test]
async fn test_cargo_publish_basic() {
let ctx = create_test_context();
let metadata = serde_json::json!({
"name": "my-crate",
"vers": "0.1.0",
"deps": [],
"features": {},
});
let crate_data = b"fake-crate-tarball";
let payload = build_publish_payload(&metadata, crate_data);
let resp = send(
&ctx.app,
Method::PUT,
"/cargo/api/v1/crates/new",
Body::from(payload),
)
.await;
assert_eq!(resp.status(), StatusCode::OK);
let stored = ctx
.state
.storage
.get("cargo/my-crate/0.1.0/my-crate-0.1.0.crate")
.await
.unwrap();
assert_eq!(&stored[..], crate_data);
let index = ctx
.state
.storage
.get("cargo/index/my/-c/my-crate")
.await
.unwrap();
let index_str = String::from_utf8_lossy(&index);
assert!(index_str.contains("\"name\":\"my-crate\""));
assert!(index_str.contains("\"vers\":\"0.1.0\""));
assert!(index_str.contains("\"cksum\":"));
}
#[tokio::test]
async fn test_cargo_publish_version_immutability() {
let ctx = create_test_context();
let metadata = serde_json::json!({
"name": "immut-test",
"vers": "1.0.0",
"deps": [],
"features": {},
});
let payload = build_publish_payload(&metadata, b"crate-v1");
let resp = send(
&ctx.app,
Method::PUT,
"/cargo/api/v1/crates/new",
Body::from(payload),
)
.await;
assert_eq!(resp.status(), StatusCode::OK);
let payload2 = build_publish_payload(&metadata, b"crate-v1-again");
let resp2 = send(
&ctx.app,
Method::PUT,
"/cargo/api/v1/crates/new",
Body::from(payload2),
)
.await;
assert_eq!(resp2.status(), StatusCode::CONFLICT);
}
#[tokio::test]
async fn test_cargo_publish_multiple_versions() {
let ctx = create_test_context();
let m1 =
serde_json::json!({"name": "multi-ver", "vers": "0.1.0", "deps": [], "features": {}});
let p1 = build_publish_payload(&m1, b"crate-01");
let r1 = send(
&ctx.app,
Method::PUT,
"/cargo/api/v1/crates/new",
Body::from(p1),
)
.await;
assert_eq!(r1.status(), StatusCode::OK);
let m2 =
serde_json::json!({"name": "multi-ver", "vers": "0.2.0", "deps": [], "features": {}});
let p2 = build_publish_payload(&m2, b"crate-02");
let r2 = send(
&ctx.app,
Method::PUT,
"/cargo/api/v1/crates/new",
Body::from(p2),
)
.await;
assert_eq!(r2.status(), StatusCode::OK);
let index = ctx
.state
.storage
.get("cargo/index/mu/lt/multi-ver")
.await
.unwrap();
let index_str = String::from_utf8_lossy(&index);
let lines: Vec<&str> = index_str.lines().collect();
assert_eq!(lines.len(), 2);
assert!(lines[0].contains("0.1.0"));
assert!(lines[1].contains("0.2.0"));
assert!(ctx
.state
.storage
.get("cargo/index-entries/mu/lt/multi-ver/0.1.0.json")
.await
.is_ok());
assert!(ctx
.state
.storage
.get("cargo/index-entries/mu/lt/multi-ver/0.2.0.json")
.await
.is_ok());
}
#[tokio::test]
async fn test_cargo_publish_migrates_single_file_index() {
let ctx = create_test_context();
let old_line = serde_json::json!({
"name": "legacy-crate", "vers": "1.0.0", "deps": [], "cksum": "a",
"features": {}, "yanked": false
})
.to_string();
ctx.state
.storage
.put(
"cargo/index/le/ga/legacy-crate",
format!("{}\n", old_line).as_bytes(),
)
.await
.unwrap();
let m = serde_json::json!({"name": "legacy-crate", "vers": "2.0.0", "deps": [], "features": {}});
let resp = send(
&ctx.app,
Method::PUT,
"/cargo/api/v1/crates/new",
Body::from(build_publish_payload(&m, b"crate2")),
)
.await;
assert_eq!(resp.status(), StatusCode::OK);
assert!(
ctx.state
.storage
.get("cargo/index-entries/le/ga/legacy-crate/1.0.0.json")
.await
.is_ok(),
"old single-file version not migrated to a per-version key"
);
let index = ctx
.state
.storage
.get("cargo/index/le/ga/legacy-crate")
.await
.unwrap();
let s = String::from_utf8_lossy(&index);
assert!(s.contains("\"vers\":\"1.0.0\""), "migrated v1 lost");
assert!(s.contains("\"vers\":\"2.0.0\""), "new v2 lost");
}
#[tokio::test]
async fn test_cargo_publish_invalid_name() {
let ctx = create_test_context();
let metadata = serde_json::json!({
"name": "../traversal",
"vers": "1.0.0",
"deps": [],
"features": {},
});
let payload = build_publish_payload(&metadata, b"bad");
let resp = send(
&ctx.app,
Method::PUT,
"/cargo/api/v1/crates/new",
Body::from(payload),
)
.await;
assert_eq!(resp.status(), StatusCode::BAD_REQUEST);
}
#[tokio::test]
async fn test_cargo_publish_truncated_payload() {
let ctx = create_test_context();
let resp = send(
&ctx.app,
Method::PUT,
"/cargo/api/v1/crates/new",
Body::from(vec![0u8; 3]),
)
.await;
assert_eq!(resp.status(), StatusCode::BAD_REQUEST);
}
#[tokio::test]
async fn test_cargo_publish_response_has_warnings() {
let ctx = create_test_context();
let metadata = serde_json::json!({
"name": "warn-test",
"vers": "1.0.0",
"deps": [],
"features": {},
});
let payload = build_publish_payload(&metadata, b"crate-data");
let resp = send(
&ctx.app,
Method::PUT,
"/cargo/api/v1/crates/new",
Body::from(payload),
)
.await;
assert_eq!(resp.status(), StatusCode::OK);
let body = body_bytes(resp).await;
let json: serde_json::Value = serde_json::from_slice(&body).unwrap();
assert!(json.get("warnings").is_some());
}
#[tokio::test]
async fn test_cargo_publish_then_download() {
let ctx = create_test_context();
let metadata = serde_json::json!({
"name": "roundtrip",
"vers": "2.0.0",
"deps": [],
"features": {},
});
let crate_data = b"published-crate-content";
let payload = build_publish_payload(&metadata, crate_data);
let publish_resp = send(
&ctx.app,
Method::PUT,
"/cargo/api/v1/crates/new",
Body::from(payload),
)
.await;
assert_eq!(publish_resp.status(), StatusCode::OK);
let dl_resp = send(
&ctx.app,
Method::GET,
"/cargo/api/v1/crates/roundtrip/2.0.0/download",
"",
)
.await;
assert_eq!(dl_resp.status(), StatusCode::OK);
let body = body_bytes(dl_resp).await;
assert_eq!(&body[..], crate_data);
}
#[tokio::test]
async fn test_cargo_publish_then_sparse_index() {
let ctx = create_test_context();
let metadata = serde_json::json!({
"name": "idx-test",
"vers": "1.0.0",
"deps": [{"name": "serde", "req": "^1", "features": [], "optional": false, "default_features": true, "target": null, "kind": "normal"}],
"features": {"default": ["serde"]},
"links": null,
});
let payload = build_publish_payload(&metadata, b"crate");
let publish_resp = send(
&ctx.app,
Method::PUT,
"/cargo/api/v1/crates/new",
Body::from(payload),
)
.await;
assert_eq!(publish_resp.status(), StatusCode::OK);
let idx_resp = send(&ctx.app, Method::GET, "/cargo/index/id/x-/idx-test", "").await;
assert_eq!(idx_resp.status(), StatusCode::OK);
let body = body_bytes(idx_resp).await;
let line: serde_json::Value =
serde_json::from_str(String::from_utf8_lossy(&body).lines().next().unwrap()).unwrap();
assert_eq!(line["name"], "idx-test");
assert_eq!(line["vers"], "1.0.0");
assert!(line["deps"].as_array().unwrap().len() == 1);
assert!(line["cksum"].as_str().unwrap().len() == 64); }
#[tokio::test]
async fn test_cargo_publish_transforms_deps_version_req_to_req() {
let ctx = create_test_context();
let metadata = serde_json::json!({
"name": "dep-test",
"vers": "1.0.0",
"deps": [{
"name": "serde",
"version_req": "^1.0",
"features": ["derive"],
"optional": false,
"default_features": true,
"target": null,
"kind": "normal",
"registry": null,
"explicit_name_in_toml": null
}, {
"name": "my_serde",
"version_req": "^1.0",
"features": [],
"optional": false,
"default_features": true,
"target": null,
"kind": "normal",
"registry": null,
"explicit_name_in_toml": "serde_json"
}],
"features": {},
});
let payload = build_publish_payload(&metadata, b"crate-data");
let resp = send(
&ctx.app,
Method::PUT,
"/cargo/api/v1/crates/new",
Body::from(payload),
)
.await;
assert_eq!(resp.status(), StatusCode::OK);
let index = ctx
.state
.storage
.get("cargo/index/de/p-/dep-test")
.await
.unwrap();
let line: serde_json::Value =
serde_json::from_str(String::from_utf8_lossy(&index).lines().next().unwrap()).unwrap();
let deps = line["deps"].as_array().unwrap();
assert_eq!(deps.len(), 2);
assert!(
deps[0].get("version_req").is_none(),
"version_req should not be in index"
);
assert_eq!(deps[0]["req"], "^1.0", "version_req must be renamed to req");
assert!(deps[0].get("explicit_name_in_toml").is_none());
assert!(
deps[0].get("package").is_none(),
"null explicit_name_in_toml should not create package field"
);
assert!(deps[1].get("explicit_name_in_toml").is_none());
assert_eq!(
deps[1]["package"], "serde_json",
"explicit_name_in_toml must become package"
);
}
#[tokio::test]
async fn test_cargo_publish_conflict_json_format() {
let ctx = create_test_context();
let metadata = serde_json::json!({
"name": "conflict-fmt",
"vers": "1.0.0",
"deps": [],
"features": {},
});
let payload = build_publish_payload(&metadata, b"v1");
let r1 = send(
&ctx.app,
Method::PUT,
"/cargo/api/v1/crates/new",
Body::from(payload),
)
.await;
assert_eq!(r1.status(), StatusCode::OK);
let payload2 = build_publish_payload(&metadata, b"v1-again");
let r2 = send(
&ctx.app,
Method::PUT,
"/cargo/api/v1/crates/new",
Body::from(payload2),
)
.await;
assert_eq!(r2.status(), StatusCode::CONFLICT);
let body = body_bytes(r2).await;
let json: serde_json::Value = serde_json::from_slice(&body).unwrap();
assert!(json["errors"].as_array().unwrap().len() > 0);
assert!(json["errors"][0]["detail"]
.as_str()
.unwrap()
.contains("already exists"));
}
#[tokio::test]
async fn test_cargo_download_blocked_by_curation() {
use crate::test_helpers::{create_test_context_with_config, send_with_headers};
let blocklist_dir = tempfile::TempDir::new().unwrap();
let blocklist_path = blocklist_dir.path().join("blocklist.json");
let blocklist = serde_json::json!({
"version": 1,
"rules": [{
"registry": "cargo",
"name": "evil-crate",
"version": "*",
"reason": "known malware"
}]
});
std::fs::write(&blocklist_path, serde_json::to_string(&blocklist).unwrap()).unwrap();
let bl_path = blocklist_path.to_str().unwrap().to_string();
let ctx = create_test_context_with_config(move |cfg| {
cfg.curation.mode = crate::config::CurationMode::Enforce;
cfg.curation.blocklist_path = Some(bl_path);
});
ctx.state
.storage
.put(
"cargo/evil-crate/1.0.0/evil-crate-1.0.0.crate",
b"evil-data",
)
.await
.unwrap();
let resp = send_with_headers(
&ctx.app,
Method::GET,
"/cargo/api/v1/crates/evil-crate/1.0.0/download",
vec![],
"",
)
.await;
assert_eq!(resp.status(), StatusCode::FORBIDDEN);
assert_eq!(
resp.headers()
.get("x-nora-decision")
.and_then(|v| v.to_str().ok()),
Some("blocked")
);
}
#[tokio::test]
async fn test_cargo_download_allowed_by_curation() {
use crate::test_helpers::{create_test_context_with_config, send_with_headers};
let blocklist_dir = tempfile::TempDir::new().unwrap();
let blocklist_path = blocklist_dir.path().join("blocklist.json");
let blocklist = serde_json::json!({
"version": 1,
"rules": [{
"registry": "cargo",
"name": "evil-crate",
"version": "*",
"reason": "known malware"
}]
});
std::fs::write(&blocklist_path, serde_json::to_string(&blocklist).unwrap()).unwrap();
let bl_path = blocklist_path.to_str().unwrap().to_string();
let ctx = create_test_context_with_config(move |cfg| {
cfg.curation.mode = crate::config::CurationMode::Enforce;
cfg.curation.blocklist_path = Some(bl_path);
});
ctx.state
.storage
.put(
"cargo/safe-crate/2.0.0/safe-crate-2.0.0.crate",
b"safe-data",
)
.await
.unwrap();
let resp = send_with_headers(
&ctx.app,
Method::GET,
"/cargo/api/v1/crates/safe-crate/2.0.0/download",
vec![],
"",
)
.await;
assert_eq!(resp.status(), StatusCode::OK);
let body = body_bytes(resp).await;
assert_eq!(&body[..], b"safe-data");
}
#[tokio::test]
async fn test_cargo_download_range_request() {
use crate::test_helpers::send_with_headers;
let ctx = create_test_context();
let krate = b"0123456789";
ctx.state
.storage
.put("cargo/rng/1.0.0/rng-1.0.0.crate", krate)
.await
.unwrap();
let url = "/cargo/api/v1/crates/rng/1.0.0/download";
let resp =
send_with_headers(&ctx.app, Method::GET, url, vec![("range", "bytes=2-5")], "").await;
assert_eq!(resp.status(), StatusCode::PARTIAL_CONTENT);
assert_eq!(
resp.headers()
.get("content-range")
.unwrap()
.to_str()
.unwrap(),
"bytes 2-5/10"
);
assert_eq!(
resp.headers()
.get("accept-ranges")
.unwrap()
.to_str()
.unwrap(),
"bytes"
);
assert_eq!(&body_bytes(resp).await[..], b"2345");
let resp =
send_with_headers(&ctx.app, Method::GET, url, vec![("range", "bytes=10-")], "").await;
assert_eq!(resp.status(), StatusCode::RANGE_NOT_SATISFIABLE);
assert_eq!(
resp.headers()
.get("content-range")
.unwrap()
.to_str()
.unwrap(),
"bytes */10"
);
let resp = send(&ctx.app, Method::GET, url, "").await;
assert_eq!(resp.status(), StatusCode::OK);
assert_eq!(
resp.headers()
.get("accept-ranges")
.unwrap()
.to_str()
.unwrap(),
"bytes"
);
assert_eq!(&body_bytes(resp).await[..], krate);
}
}