use crate::activity_log::{ActionType, ActivityEntry};
use crate::audit::AuditEntry;
use crate::auth::{enforce_namespace_scope, NamespaceAuthority};
use crate::circuit_breaker::CircuitBreakerRegistry;
use crate::config::basic_auth_header;
use crate::registry::docker_auth::DockerAuth;
use crate::registry::{circuit_open_response, method_not_allowed, ProxyError};
use crate::registry::{
content_length, sha256_of_file, stream_body_to_file, StreamOutcome, TempFileGuard,
};
use crate::secrets::expose_opt;
use crate::storage::Storage;
use crate::validation::{
ends_with_ci, validate_digest, validate_docker_name, validate_docker_reference,
};
use crate::AppState;
use axum::{
body::{Body, Bytes},
extract::{Path, State},
http::{header, HeaderName, Method, StatusCode, Uri},
response::{IntoResponse, Response},
routing::get,
Extension, Json, Router,
};
use futures::StreamExt;
use parking_lot::RwLock;
use serde::{Deserialize, Serialize};
use serde_json::{json, Value};
use std::collections::HashMap;
use std::sync::Arc;
use std::time::{Duration, SystemTime, UNIX_EPOCH};
use tokio_util::io::ReaderStream;
fn blob_key(namespace: Option<&str>, name: &str, digest: &str) -> String {
match namespace {
Some(ns) => format!("docker/{}/{}/blobs/{}", ns, name, digest),
None => format!("docker/{}/blobs/{}", name, digest),
}
}
fn manifest_key(namespace: Option<&str>, name: &str, reference: &str) -> String {
match namespace {
Some(ns) => format!("docker/{}/{}/manifests/{}.json", ns, name, reference),
None => format!("docker/{}/manifests/{}.json", name, reference),
}
}
fn manifest_meta_key(namespace: Option<&str>, name: &str, reference: &str) -> String {
match namespace {
Some(ns) => format!("docker/{}/{}/manifests/{}.meta.json", ns, name, reference),
None => format!("docker/{}/manifests/{}.meta.json", name, reference),
}
}
fn manifest_prefix(namespace: Option<&str>, name: &str) -> String {
match namespace {
Some(ns) => format!("docker/{}/{}/manifests/", ns, name),
None => format!("docker/{}/manifests/", name),
}
}
pub(crate) struct Canonical {
pub name: String,
pub namespace: Option<String>,
matched_upstream_idx: Option<usize>,
pub denied: bool,
}
impl Canonical {
pub fn upstreams_to_try<'a>(
&self,
upstreams: &'a [crate::config::DockerUpstream],
) -> Vec<&'a crate::config::DockerUpstream> {
match self.matched_upstream_idx {
Some(idx) => vec![&upstreams[idx]],
None => upstreams.iter().collect(),
}
}
pub fn denied_response(&self) -> Option<axum::response::Response> {
if self.denied {
tracing::warn!(
name = %self.name,
"Docker request denied: image name did not match any configured upstream prefix"
);
Some(
(
axum::http::StatusCode::FORBIDDEN,
axum::Json(serde_json::json!({
"errors": [{
"code": "DENIED",
"message": "image name does not match any configured upstream prefix",
"detail": "default_action is set to deny"
}]
})),
)
.into_response(),
)
} else {
None
}
}
}
pub(crate) fn canonicalize(
raw_name: &str,
docker_config: &crate::config::DockerConfig,
) -> Canonical {
let upstreams = &docker_config.upstreams;
let deny_mode = docker_config.default_action == crate::config::DefaultAction::Deny;
if let Some((first_segment, rest)) = raw_name.split_once('/') {
for (idx, upstream) in upstreams.iter().enumerate() {
if let Some(ref prefix) = upstream.prefix {
if first_segment == prefix {
tracing::debug!(
prefix = %prefix,
upstream = %upstream.url,
stripped_name = %rest,
routing = "prefix",
"Docker path-based upstream routing"
);
return Canonical {
name: rest.to_string(),
namespace: Some(upstream.resolved_namespace()),
matched_upstream_idx: Some(idx),
denied: false,
};
}
}
}
if first_segment.contains('.') && !rest.is_empty() {
for (idx, upstream) in upstreams.iter().enumerate() {
let ns = upstream.resolved_namespace();
if first_segment == ns {
return Canonical {
name: rest.to_string(),
namespace: Some(ns),
matched_upstream_idx: Some(idx),
denied: false,
};
}
}
let ns = upstreams.first().map(|u| u.resolved_namespace());
return Canonical {
name: rest.to_string(),
namespace: ns,
matched_upstream_idx: None,
denied: deny_mode,
};
}
}
let ns = upstreams.first().map(|u| u.resolved_namespace());
Canonical {
name: raw_name.to_string(),
namespace: ns,
matched_upstream_idx: None,
denied: deny_mode,
}
}
async fn storage_get_with_fallback(
storage: &Storage,
ns_key: &str,
legacy_key: &str,
) -> Result<Bytes, crate::storage::StorageError> {
match storage.get(ns_key).await {
Err(crate::storage::StorageError::NotFound) if ns_key != legacy_key => {
storage.get(legacy_key).await
}
other => other,
}
}
async fn storage_get_reader_with_fallback(
storage: &Storage,
ns_key: &str,
legacy_key: &str,
) -> Result<
(
u64,
std::pin::Pin<Box<dyn tokio::io::AsyncRead + Send + Unpin>>,
),
crate::storage::StorageError,
> {
match storage.get_reader(ns_key).await {
Err(crate::storage::StorageError::NotFound) if ns_key != legacy_key => {
storage.get_reader(legacy_key).await
}
other => other,
}
}
struct VerifyingReader<R> {
inner: R,
hasher: sha2::Sha256,
expected_hex: String,
finished: bool,
}
impl<R> VerifyingReader<R> {
fn new(inner: R, digest: &str) -> Self {
let expected_hex = digest
.strip_prefix("sha256:")
.unwrap_or(digest)
.to_ascii_lowercase();
Self {
inner,
hasher: sha2::Sha256::default(),
expected_hex,
finished: false,
}
}
}
impl<R: tokio::io::AsyncRead + Unpin> tokio::io::AsyncRead for VerifyingReader<R> {
fn poll_read(
self: std::pin::Pin<&mut Self>,
cx: &mut std::task::Context<'_>,
buf: &mut tokio::io::ReadBuf<'_>,
) -> std::task::Poll<std::io::Result<()>> {
use sha2::Digest as _;
use std::task::Poll;
let this = self.get_mut();
if this.finished {
return Poll::Ready(Ok(()));
}
let before = buf.filled().len();
match std::pin::Pin::new(&mut this.inner).poll_read(cx, buf) {
Poll::Ready(Ok(())) => {
let filled = buf.filled();
if filled.len() > before {
this.hasher.update(&filled[before..]);
Poll::Ready(Ok(()))
} else {
this.finished = true;
let got = hex::encode(this.hasher.clone().finalize());
if got == this.expected_hex {
Poll::Ready(Ok(()))
} else {
tracing::error!(
expected = %this.expected_hex,
got = %got,
"blob integrity verification failed while streaming — aborting response"
);
Poll::Ready(Err(std::io::Error::new(
std::io::ErrorKind::InvalidData,
"blob integrity verification failed",
)))
}
}
}
other => other,
}
}
}
async fn storage_stat_with_fallback(
storage: &Storage,
ns_key: &str,
legacy_key: &str,
) -> Option<crate::storage::FileMeta> {
if let Some(meta) = storage.stat(ns_key).await {
return Some(meta);
}
if ns_key != legacy_key {
return storage.stat(legacy_key).await;
}
None
}
#[derive(Debug, Clone, Serialize, Deserialize, Default)]
pub struct ImageMetadata {
pub push_timestamp: u64,
pub last_pulled: u64,
pub downloads: u64,
pub size_bytes: u64,
pub os: String,
pub arch: String,
#[serde(skip_serializing_if = "Option::is_none")]
pub variant: Option<String>,
pub layers: Vec<LayerInfo>,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct LayerInfo {
pub digest: String,
pub size: u64,
}
pub struct UploadSession {
temp_path: std::path::PathBuf,
size: u64,
name: String,
created_at: std::time::Instant,
}
const DEFAULT_MAX_UPLOAD_SESSIONS: usize = 100;
const DEFAULT_MAX_SESSION_SIZE_MB: usize = 2048;
const SESSION_TTL: Duration = Duration::from_secs(30 * 60);
fn max_upload_sessions() -> usize {
std::env::var("NORA_MAX_UPLOAD_SESSIONS")
.ok()
.and_then(|v| v.parse().ok())
.unwrap_or(DEFAULT_MAX_UPLOAD_SESSIONS)
}
fn too_many_uploads() -> Response {
let body = json!({
"errors": [{
"code": "TOOMANYREQUESTS",
"message": "too many concurrent uploads",
"detail": { "limit": max_upload_sessions() }
}]
});
let retry_after = rand::Rng::gen_range(&mut rand::thread_rng(), 3..=10).to_string();
(
StatusCode::TOO_MANY_REQUESTS,
[
(header::RETRY_AFTER, retry_after.as_str()),
(header::CONTENT_TYPE, "application/json"),
],
body.to_string(),
)
.into_response()
}
fn max_session_size() -> usize {
let mb = std::env::var("NORA_MAX_UPLOAD_SESSION_SIZE_MB")
.ok()
.and_then(|v| v.parse::<usize>().ok())
.unwrap_or(DEFAULT_MAX_SESSION_SIZE_MB);
mb.saturating_mul(1024 * 1024)
}
const MAX_MANIFEST_BYTES: usize = 4 * 1024 * 1024;
fn effective_upload_cap(body_limit_mb: usize) -> u64 {
let body_limit = (body_limit_mb as u64).saturating_mul(1024 * 1024);
std::cmp::min(max_session_size() as u64, body_limit)
}
static IN_FLIGHT_UPLOADS: std::sync::atomic::AtomicUsize = std::sync::atomic::AtomicUsize::new(0);
struct InFlightGuard;
impl InFlightGuard {
fn try_acquire() -> Option<Self> {
use std::sync::atomic::Ordering;
let prev = IN_FLIGHT_UPLOADS.fetch_add(1, Ordering::AcqRel);
if prev >= max_upload_sessions() {
IN_FLIGHT_UPLOADS.fetch_sub(1, Ordering::AcqRel);
return None;
}
Some(InFlightGuard)
}
}
impl Drop for InFlightGuard {
fn drop(&mut self) {
IN_FLIGHT_UPLOADS.fetch_sub(1, std::sync::atomic::Ordering::AcqRel);
}
}
pub fn in_flight_uploads() -> usize {
IN_FLIGHT_UPLOADS.load(std::sync::atomic::Ordering::Acquire)
}
struct ProxyDownloadGuard;
impl Drop for ProxyDownloadGuard {
fn drop(&mut self) {
crate::metrics::PROXY_ACTIVE_DOWNLOADS.dec();
}
}
fn validate_upload_uuid(uuid: &str) -> Result<(), &'static str> {
if uuid.is_empty() || uuid.len() > 36 {
return Err("invalid upload UUID length");
}
if uuid
.bytes()
.any(|b| !matches!(b, b'0'..=b'9' | b'a'..=b'f' | b'-'))
{
return Err("invalid upload UUID characters");
}
Ok(())
}
fn upload_temp_dir(data_dir: &str) -> std::path::PathBuf {
let dir = std::path::PathBuf::from(data_dir).join("tmp/docker-uploads");
if let Err(e) = std::fs::create_dir_all(&dir) {
tracing::error!(path = %dir.display(), error = %e, "failed to create upload temp directory");
}
dir
}
fn proxy_temp_dir(data_dir: &str) -> std::path::PathBuf {
let dir = std::path::PathBuf::from(data_dir).join("tmp/docker-proxy");
if let Err(e) = std::fs::create_dir_all(&dir) {
tracing::error!(path = %dir.display(), error = %e, "failed to create proxy temp directory");
}
dir
}
pub fn cleanup_expired_sessions(sessions: &RwLock<HashMap<String, UploadSession>>) {
let mut guard = sessions.write();
let before = guard.len();
guard.retain(|_, s| {
if s.created_at.elapsed() < SESSION_TTL {
return true;
}
let temp_active = std::fs::metadata(&s.temp_path)
.and_then(|m| m.modified())
.ok()
.and_then(|t| t.elapsed().ok())
.is_some_and(|age| age < SESSION_TTL);
if temp_active {
return true;
}
let _ = std::fs::remove_file(&s.temp_path);
false
});
let removed = before - guard.len();
if removed > 0 {
tracing::info!(
removed = removed,
remaining = guard.len(),
"Cleaned up expired upload sessions"
);
}
}
pub fn cleanup_upload_temp_dir(data_dir: &str) {
let dir = std::path::PathBuf::from(data_dir).join("tmp/docker-uploads");
let entries = match std::fs::read_dir(&dir) {
Ok(e) => e,
Err(e) if e.kind() == std::io::ErrorKind::NotFound => return,
Err(e) => {
tracing::warn!(path = %dir.display(), error = %e, "Failed to read upload temp directory for cleanup");
return;
}
};
let mut removed = 0u64;
for entry in entries.flatten() {
let is_stale = entry
.metadata()
.ok()
.and_then(|m| m.modified().ok())
.and_then(|t| t.elapsed().ok())
.is_some_and(|age| age >= SESSION_TTL);
if is_stale && std::fs::remove_file(entry.path()).is_ok() {
removed += 1;
}
}
if removed > 0 {
tracing::info!(removed, dir = %dir.display(), "Cleaned up stale Docker upload temp files");
}
}
const PROXY_TEMP_MAX_AGE: std::time::Duration = std::time::Duration::from_secs(4 * 60 * 60);
pub fn cleanup_proxy_temp_dir(data_dir: &str) {
let dir = std::path::PathBuf::from(data_dir).join("tmp/docker-proxy");
let entries = match std::fs::read_dir(&dir) {
Ok(e) => e,
Err(e) if e.kind() == std::io::ErrorKind::NotFound => return,
Err(e) => {
tracing::warn!(path = %dir.display(), error = %e, "Failed to read proxy temp directory for cleanup");
return;
}
};
let mut removed = 0u64;
for entry in entries.flatten() {
let is_stale = entry
.metadata()
.ok()
.and_then(|m| m.modified().ok())
.and_then(|t| t.elapsed().ok())
.is_some_and(|age| age >= PROXY_TEMP_MAX_AGE);
if is_stale && std::fs::remove_file(entry.path()).is_ok() {
removed += 1;
}
}
if removed > 0 {
tracing::info!(removed, dir = %dir.display(), "Cleaned up stale proxy temp files");
}
}
fn resolve_quarantine(state: &AppState) -> (crate::digest_quarantine::QuarantineMode, i64) {
use crate::digest_quarantine::QuarantineMode;
let mode = state
.config
.curation
.docker
.quarantine
.as_ref()
.or(state.config.curation.quarantine.as_ref())
.cloned()
.unwrap_or(QuarantineMode::Off);
if matches!(mode, QuarantineMode::Off) {
return (QuarantineMode::Off, 0);
}
let ttl_str = state
.config
.curation
.docker
.quarantine_ttl
.as_deref()
.or(state.config.curation.quarantine_ttl.as_deref())
.unwrap_or("14d");
let secs = crate::curation::parse_duration(ttl_str).unwrap_or(14 * 86400);
(mode, secs)
}
fn quarantine_forbidden(
digest: &str,
status: &crate::digest_quarantine::QuarantineStatus,
quarantine_secs: i64,
) -> Response {
let remaining = match status {
crate::digest_quarantine::QuarantineStatus::New => quarantine_secs,
crate::digest_quarantine::QuarantineStatus::Pending { remaining_secs } => *remaining_secs,
crate::digest_quarantine::QuarantineStatus::Mature => 0,
};
let quarantine_until = chrono::Utc::now().timestamp() + remaining;
let body = json!({
"errors": [{
"code": "DENIED",
"message": "held by quarantine policy on registry 'docker': digest must age past the quarantine window",
"detail": {
"registry": "docker",
"digest": digest,
"policy": {
"control": "quarantine",
"quarantine_ttl_secs": quarantine_secs,
},
"quarantine_until": quarantine_until,
}
}]
});
(
StatusCode::FORBIDDEN,
[
(
HeaderName::from_static("x-nora-quarantine"),
status.header_value(),
),
(header::CONTENT_TYPE, "application/json"),
],
body.to_string(),
)
.into_response()
}
fn manifest_unknown_response(name: &str, reference: &str) -> Response {
(
StatusCode::NOT_FOUND,
Json(json!({
"errors": [{
"code": "MANIFEST_UNKNOWN",
"message": "manifest unknown",
"detail": { "name": name, "reference": reference }
}]
})),
)
.into_response()
}
fn blob_unknown_response(digest: &str) -> Response {
(
StatusCode::NOT_FOUND,
Json(json!({
"errors": [{
"code": "BLOB_UNKNOWN",
"message": "blob unknown to registry",
"detail": { "digest": digest }
}]
})),
)
.into_response()
}
#[must_use = "the returned response blocks a quarantined artifact; dropping it serves it"]
fn quarantine_cache_serve_gate(state: &AppState, digest: &str) -> Option<Response> {
let (q_mode, q_secs) = resolve_quarantine(state);
if matches!(q_mode, crate::digest_quarantine::QuarantineMode::Off) {
return None;
}
let status = state.digest_store.check("docker", digest, q_secs);
if let crate::digest_quarantine::QuarantineStatus::Pending { .. } = status {
let outcome = if matches!(q_mode, crate::digest_quarantine::QuarantineMode::Enforce) {
"blocked"
} else {
"observed"
};
crate::metrics::QUARANTINE_HOLDS_TOTAL
.with_label_values(&["docker", outcome])
.inc();
tracing::warn!(
digest = %digest,
status = %status.header_value(),
mode = ?q_mode,
quarantine_ttl_secs = q_secs,
"Quarantine: held cached artifact (proxy cooldown)"
);
if matches!(q_mode, crate::digest_quarantine::QuarantineMode::Enforce) {
return Some(quarantine_forbidden(digest, &status, q_secs));
}
}
None
}
#[must_use = "the returned response blocks a quarantined artifact; dropping it serves it"]
fn quarantine_proxy_fetch_gate(state: &AppState, digest: &str, upstream: &str) -> Option<Response> {
let (q_mode, q_secs) = resolve_quarantine(state);
if matches!(q_mode, crate::digest_quarantine::QuarantineMode::Off) {
return None;
}
state.digest_store.record("docker", digest, upstream, None);
let status = state.digest_store.check("docker", digest, q_secs);
if !matches!(status, crate::digest_quarantine::QuarantineStatus::Mature) {
let outcome = if matches!(q_mode, crate::digest_quarantine::QuarantineMode::Enforce) {
"blocked"
} else {
"observed"
};
crate::metrics::QUARANTINE_HOLDS_TOTAL
.with_label_values(&["docker", outcome])
.inc();
tracing::warn!(
digest = %digest,
upstream = %upstream,
status = %status.header_value(),
mode = ?q_mode,
quarantine_ttl_secs = q_secs,
"Quarantine: proxy-fetched blob held (new to this mirror)"
);
if matches!(q_mode, crate::digest_quarantine::QuarantineMode::Enforce) {
return Some(quarantine_forbidden(digest, &status, q_secs));
}
}
None
}
pub fn routes() -> Router<AppState> {
Router::new()
.route(
"/v2/",
get(check).fallback(|| async { method_not_allowed("GET") }),
)
.route("/v2/_catalog", get(catalog))
.route("/v2/{*rest}", axum::routing::any(docker_v2_dispatch))
}
async fn docker_v2_dispatch(
state: State<AppState>,
method: Method,
Path(wildcard): Path<String>,
Extension(authority): Extension<NamespaceAuthority>,
uri: Uri,
headers: axum::http::HeaderMap,
body: Body,
) -> Response {
let rest = wildcard.trim_start_matches('/');
if rest.is_empty() {
return StatusCode::NOT_FOUND.into_response();
}
let is_write = matches!(
method,
Method::POST | Method::PUT | Method::PATCH | Method::DELETE
);
if let Some((name, after)) = rest.rsplit_once("/blobs/uploads/") {
if name.is_empty() {
return StatusCode::NOT_FOUND.into_response();
}
if validate_docker_name(name).is_err() {
return (StatusCode::BAD_REQUEST, "Invalid image name").into_response();
}
if is_write && enforce_namespace_scope(&authority, name).is_err() {
return StatusCode::FORBIDDEN.into_response();
}
if is_write {
if let Some(len) = content_length(&headers) {
if len > effective_upload_cap(state.0.config.server.body_limit_mb) {
return (StatusCode::PAYLOAD_TOO_LARGE, "Upload exceeds size limit")
.into_response();
}
}
}
return if after.is_empty() {
match method {
Method::POST => {
let params = parse_query_string(uri.query());
if params.contains_key("digest") {
let upload_id = uuid::Uuid::new_v4().to_string();
upload_blob(
state,
Path((name.to_string(), upload_id)),
axum::extract::Query(params),
body,
)
.await
} else if let (Some(digest), Some(from)) =
(params.get("mount"), params.get("from"))
{
mount_blob(state, name, from, digest).await
} else {
start_upload(state, Path(name.to_string())).await
}
}
_ => method_not_allowed("POST"),
}
} else {
match method {
Method::PATCH => {
patch_blob(state, Path((name.to_string(), after.to_string())), body).await
}
Method::PUT => {
let params = parse_query_string(uri.query());
upload_blob(
state,
Path((name.to_string(), after.to_string())),
axum::extract::Query(params),
body,
)
.await
}
Method::DELETE => {
cancel_upload(state, Path((name.to_string(), after.to_string()))).await
}
_ => method_not_allowed("PATCH, PUT, DELETE"),
}
};
}
if let Some((name, digest)) = rest.rsplit_once("/blobs/") {
if name.is_empty() || digest.is_empty() {
return StatusCode::NOT_FOUND.into_response();
}
if validate_docker_name(name).is_err() {
return (StatusCode::BAD_REQUEST, "Invalid image name").into_response();
}
if is_write && enforce_namespace_scope(&authority, name).is_err() {
return StatusCode::FORBIDDEN.into_response();
}
return match method {
Method::HEAD => check_blob(state, Path((name.to_string(), digest.to_string()))).await,
Method::GET => {
download_blob(state, headers, Path((name.to_string(), digest.to_string()))).await
}
Method::DELETE => {
delete_blob(state, Path((name.to_string(), digest.to_string()))).await
}
_ => method_not_allowed("GET, HEAD, DELETE"),
};
}
if let Some((name, reference)) = rest.rsplit_once("/manifests/") {
if name.is_empty() || reference.is_empty() {
return StatusCode::NOT_FOUND.into_response();
}
if validate_docker_name(name).is_err() {
return (StatusCode::BAD_REQUEST, "Invalid image name").into_response();
}
if is_write && enforce_namespace_scope(&authority, name).is_err() {
return StatusCode::FORBIDDEN.into_response();
}
return match method {
Method::GET | Method::HEAD => {
let resp = get_manifest(
state,
headers,
Path((name.to_string(), reference.to_string())),
)
.await;
if method == Method::HEAD {
let (parts, _) = resp.into_parts();
Response::from_parts(parts, axum::body::Body::empty())
} else {
resp
}
}
Method::PUT => match axum::body::to_bytes(body, MAX_MANIFEST_BYTES).await {
Ok(b) => {
put_manifest(state, Path((name.to_string(), reference.to_string())), b).await
}
Err(_) => (StatusCode::PAYLOAD_TOO_LARGE, "Manifest too large").into_response(),
},
Method::DELETE => {
delete_manifest(state, Path((name.to_string(), reference.to_string()))).await
}
_ => method_not_allowed("GET, HEAD, PUT, DELETE"),
};
}
if let Some(name) = rest.strip_suffix("/tags/list") {
if name.is_empty() {
return StatusCode::NOT_FOUND.into_response();
}
if validate_docker_name(name).is_err() {
return (StatusCode::BAD_REQUEST, "Invalid image name").into_response();
}
return match method {
Method::GET => list_tags(state, Path(name.to_string())).await,
_ => method_not_allowed("GET"),
};
}
StatusCode::NOT_FOUND.into_response()
}
fn parse_query_string(query: Option<&str>) -> HashMap<String, String> {
use percent_encoding::percent_decode_str;
query
.unwrap_or_default()
.split('&')
.filter(|s| !s.is_empty())
.filter_map(|pair| {
let (k, v) = pair.split_once('=')?;
Some((
percent_decode_str(k).decode_utf8_lossy().into_owned(),
percent_decode_str(v).decode_utf8_lossy().into_owned(),
))
})
.collect()
}
async fn check() -> impl IntoResponse {
(
StatusCode::OK,
[(
HeaderName::from_static("docker-distribution-api-version"),
"registry/2.0",
)],
Json(json!({})),
)
}
pub(crate) fn strip_docker_namespace(name: &str) -> &str {
if let Some((first, rest)) = name.split_once('/') {
if first.contains('.') && !rest.is_empty() {
return rest;
}
}
name
}
async fn catalog(State(state): State<AppState>) -> Response {
let keys = match state.storage.list("docker/").await {
Ok(k) => k,
Err(e) => {
tracing::error!(error = ?e, "docker: failed to list storage for catalog");
return StatusCode::SERVICE_UNAVAILABLE.into_response();
}
};
let mut repos: Vec<String> = keys
.iter()
.filter_map(|k| {
let rest = k.strip_prefix("docker/")?;
let name = if let Some(idx) = rest.find("/manifests/") {
&rest[..idx]
} else if let Some(idx) = rest.find("/blobs/") {
&rest[..idx]
} else {
return None;
};
if name.is_empty() {
return None;
}
Some(canonicalize(name, &state.config.docker).name)
})
.collect();
repos.sort();
repos.dedup();
Json(json!({ "repositories": repos })).into_response()
}
async fn check_blob(
State(state): State<AppState>,
Path((name, digest)): Path<(String, String)>,
) -> Response {
let c = canonicalize(&name, &state.config.docker);
if let Some(r) = c.denied_response() {
return r;
}
let name = c.name;
if let Err(e) = validate_docker_name(&name) {
return (StatusCode::BAD_REQUEST, e.to_string()).into_response();
}
if let Err(e) = validate_digest(&digest) {
return (StatusCode::BAD_REQUEST, e.to_string()).into_response();
}
let key = blob_key(c.namespace.as_deref(), &name, &digest);
let legacy_key = blob_key(None, &name, &digest);
match storage_stat_with_fallback(&state.storage, &key, &legacy_key).await {
Some(meta) => {
if let Some(resp) = quarantine_cache_serve_gate(&state, &digest) {
return resp;
}
(
StatusCode::OK,
[(header::CONTENT_LENGTH, meta.size.to_string())],
)
.into_response()
}
None => StatusCode::NOT_FOUND.into_response(),
}
}
async fn download_blob(
State(state): State<AppState>,
headers: axum::http::HeaderMap,
Path((name, digest)): Path<(String, String)>,
) -> Response {
let c = canonicalize(&name, &state.config.docker);
if let Some(r) = c.denied_response() {
return r;
}
let upstreams_to_try = c.upstreams_to_try(&state.config.docker.upstreams);
let ns = c.namespace;
let name = c.name;
if let Err(e) = validate_docker_name(&name) {
return (StatusCode::BAD_REQUEST, e.to_string()).into_response();
}
if let Err(e) = validate_digest(&digest) {
return (StatusCode::BAD_REQUEST, e.to_string()).into_response();
}
let internal = crate::curation::is_internal_namespace(
&state.curation().curation_engine,
crate::curation::RegistryType::Docker,
&name,
);
if !internal {
if let Some(response) = crate::curation::check_download(
&state.curation().curation_engine,
state.bypass_token().as_deref(),
&headers,
crate::curation::RegistryType::Docker,
&name,
Some(&digest),
None,
) {
return response;
}
}
let key = blob_key(ns.as_deref(), &name, &digest);
let legacy_key = blob_key(None, &name, &digest);
if let Ok((size, reader)) =
storage_get_reader_with_fallback(&state.storage, &key, &legacy_key).await
{
if let Some(response) = crate::curation::verify_integrity_by_hash(
&state.curation().curation_engine,
crate::curation::RegistryType::Docker,
&name,
Some(&digest),
&digest,
) {
return response;
}
if let Some(resp) = quarantine_cache_serve_gate(&state, &digest) {
return resp;
}
if let Some(response) = crate::registry::range::range_response(
&state.storage,
&[&key, &legacy_key],
&headers,
size,
"application/octet-stream",
&[
(
header::CACHE_CONTROL,
"public, max-age=31536000, immutable".to_string(),
),
(
HeaderName::from_static("docker-content-digest"),
digest.clone(),
),
],
)
.await
{
if response.status() == StatusCode::PARTIAL_CONTENT {
state.metrics.record_download("docker");
state.metrics.record_cache_hit("docker");
}
return response;
}
state.metrics.record_download("docker");
state.metrics.record_cache_hit("docker");
state.activity.push(ActivityEntry::new(
ActionType::Pull,
format!("{}@{}", name, &digest[..19.min(digest.len())]),
crate::registry_type::RegistryType::Docker,
"LOCAL",
));
let stream = ReaderStream::new(VerifyingReader::new(reader, &digest));
return Response::builder()
.status(StatusCode::OK)
.header(header::CONTENT_TYPE, "application/octet-stream")
.header(header::CONTENT_LENGTH, size)
.header(header::ACCEPT_RANGES, "bytes")
.header(header::CACHE_CONTROL, "public, max-age=31536000, immutable")
.header("docker-content-digest", &digest)
.body(Body::from_stream(stream))
.unwrap_or_else(|_| StatusCode::INTERNAL_SERVER_ERROR.into_response());
}
if internal {
crate::metrics::record_namespace_isolation_refused("docker");
return blob_unknown_response(&digest);
}
let temp_dir = proxy_temp_dir(&state.config.storage.path);
let names_to_try: Vec<String> = if name.contains('/') {
vec![name.clone()]
} else {
vec![name.clone(), format!("library/{}", name)]
};
for try_name in &names_to_try {
for upstream in &upstreams_to_try {
match fetch_blob_from_upstream(
&state.http_client,
&upstream.url,
try_name,
&digest,
&state.docker_auth,
state.config.docker.proxy_timeout,
state.config.docker.read_timeout,
expose_opt(&upstream.auth),
&state.circuit_breaker,
&temp_dir,
)
.await
{
Ok(mut fetched) => {
let expected_hash = digest.strip_prefix("sha256:").unwrap_or(&digest);
if fetched.sha256 != expected_hash {
tracing::warn!(
name = %try_name,
digest = %digest,
expected = %expected_hash,
actual = %fetched.sha256,
"Docker blob SHA-256 mismatch from upstream — rejecting"
);
return StatusCode::BAD_GATEWAY.into_response();
}
let hash_with_prefix = format!("sha256:{}", fetched.sha256);
if let Some(response) = crate::curation::verify_integrity_by_hash(
&state.curation().curation_engine,
crate::curation::RegistryType::Docker,
try_name,
Some(&digest),
&hash_with_prefix,
) {
return response;
}
state.metrics.record_download("docker");
state.metrics.record_cache_miss("docker");
state.activity.push(ActivityEntry::new(
ActionType::ProxyFetch,
format!("{}@{}", try_name, &digest[..19.min(digest.len())]),
crate::registry_type::RegistryType::Docker,
"PROXY",
));
let file_size = tokio::fs::metadata(&fetched.path)
.await
.map(|m| m.len())
.unwrap_or(0);
let sha_for_pin = fetched.sha256.clone();
match state
.storage
.put_from_path(&key, &fetched.path, Some(&sha_for_pin))
.await
{
Ok(()) => {
fetched._guard.disarm();
state.repo_index.invalidate("docker");
}
Err(e) => {
tracing::error!(
error = %e,
key = %key,
"Failed to store proxied blob — serving from upstream anyway"
);
}
}
if let Some(resp) = quarantine_proxy_fetch_gate(&state, &digest, &upstream.url)
{
return resp;
}
if fetched._guard.path.is_none() {
match state.storage.get_reader(&key).await {
Ok((size, reader)) => {
let stream =
ReaderStream::new(VerifyingReader::new(reader, &digest));
return Response::builder()
.status(StatusCode::OK)
.header(header::CONTENT_TYPE, "application/octet-stream")
.header(header::CONTENT_LENGTH, size)
.header(
header::CACHE_CONTROL,
"public, max-age=31536000, immutable",
)
.header("docker-content-digest", &digest)
.body(Body::from_stream(stream))
.unwrap_or_else(|_| {
StatusCode::INTERNAL_SERVER_ERROR.into_response()
});
}
Err(e) => {
tracing::error!(error = %e, key = %key, "Failed to read just-stored blob");
return StatusCode::INTERNAL_SERVER_ERROR.into_response();
}
}
} else {
match tokio::fs::File::open(&fetched.path).await {
Ok(file) => {
let stream = ReaderStream::new(VerifyingReader::new(file, &digest));
return Response::builder()
.status(StatusCode::OK)
.header(header::CONTENT_TYPE, "application/octet-stream")
.header(header::CONTENT_LENGTH, file_size)
.header("docker-content-digest", &digest)
.body(Body::from_stream(stream))
.unwrap_or_else(|_| {
StatusCode::INTERNAL_SERVER_ERROR.into_response()
});
}
Err(e) => {
tracing::error!(error = %e, "Failed to open temp blob file for streaming");
return StatusCode::INTERNAL_SERVER_ERROR.into_response();
}
}
}
}
Err(ProxyError::CircuitOpen(reg)) => return circuit_open_response(®),
Err(e) => {
tracing::debug!(error = ?e, upstream = %upstream.url, name = %try_name, "Docker blob proxy fetch failed, trying next");
continue;
}
}
}
}
if !state.config.docker.upstreams.is_empty() {
tracing::warn!(registry = "docker", name = %name, digest = %digest, "Proxy failed, returning 404");
}
StatusCode::NOT_FOUND.into_response()
}
fn blob_created(name: &str, digest: &str) -> Response {
(
StatusCode::CREATED,
[
(header::LOCATION, format!("/v2/{}/blobs/{}", name, digest)),
(
HeaderName::from_static("docker-content-digest"),
digest.to_string(),
),
],
)
.into_response()
}
async fn mount_blob(state: State<AppState>, raw_name: &str, from: &str, digest: &str) -> Response {
let c = canonicalize(raw_name, &state.config.docker);
if let Some(r) = c.denied_response() {
return r;
}
let name = c.name;
if let Err(e) = validate_docker_name(&name) {
return (StatusCode::BAD_REQUEST, e.to_string()).into_response();
}
let Some(hex) = digest.strip_prefix("sha256:") else {
return start_upload(state, Path(raw_name.to_string())).await;
};
if validate_digest(digest).is_err() {
return start_upload(state, Path(raw_name.to_string())).await;
}
let c_from = canonicalize(from, &state.config.docker);
if c_from.denied || validate_docker_name(&c_from.name).is_err() {
return start_upload(state, Path(raw_name.to_string())).await;
}
let dst_key = format!("docker/{}/blobs/{}", name, digest);
if state.storage.stat(&dst_key).await.is_some() {
return blob_created(&name, digest);
}
if quarantine_cache_serve_gate(&state, digest).is_some() {
return start_upload(state, Path(raw_name.to_string())).await;
}
let ns_key = blob_key(c_from.namespace.as_deref(), &c_from.name, digest);
let legacy_key = blob_key(None, &c_from.name, digest);
let src_key = if state.storage.stat(&ns_key).await.is_some() {
ns_key
} else if ns_key != legacy_key && state.storage.stat(&legacy_key).await.is_some() {
legacy_key
} else {
return start_upload(state, Path(raw_name.to_string())).await;
};
if let Err(e) = state.storage.copy(&src_key, &dst_key, Some(hex)).await {
tracing::warn!(error = %e, src = %src_key, dst = %dst_key, "Blob mount failed, falling back to upload");
return start_upload(state, Path(raw_name.to_string())).await;
}
state.audit.log(AuditEntry::new(
"mount",
"api",
&format!("{}@{}", name, digest),
"docker",
"blob",
));
state.activity.push(ActivityEntry::new(
ActionType::Push,
format!("{}@{}", name, &digest[..19.min(digest.len())]),
crate::registry_type::RegistryType::Docker,
"LOCAL",
));
state.repo_index.invalidate("docker");
blob_created(&name, digest)
}
async fn start_upload(State(state): State<AppState>, Path(name): Path<String>) -> Response {
let c = canonicalize(&name, &state.config.docker);
if let Some(r) = c.denied_response() {
return r;
}
let name = c.name;
if let Err(e) = validate_docker_name(&name) {
return (StatusCode::BAD_REQUEST, e.to_string()).into_response();
}
let uuid = uuid::Uuid::new_v4().to_string();
let temp_dir = upload_temp_dir(&state.config.storage.path);
let temp_path = temp_dir.join(&uuid);
if let Err(e) = tokio::fs::File::create(&temp_path).await {
tracing::error!(path = %temp_path.display(), error = %e, "Failed to create upload temp file");
return StatusCode::INTERNAL_SERVER_ERROR.into_response();
}
let rejected_temp = {
let mut sessions = state.upload_sessions.write();
let max_sessions = max_upload_sessions();
if sessions.len() >= max_sessions {
tracing::warn!(
max = max_sessions,
current = sessions.len(),
"Upload session limit reached — rejecting new upload"
);
Some(temp_path)
} else {
sessions.insert(
uuid.clone(),
UploadSession {
temp_path,
size: 0,
name: name.clone(),
created_at: std::time::Instant::now(),
},
);
None
}
};
if let Some(p) = rejected_temp {
let _ = tokio::fs::remove_file(&p).await;
return too_many_uploads();
}
let location = format!("/v2/{}/blobs/uploads/{}", name, uuid);
(
StatusCode::ACCEPTED,
[
(header::LOCATION, location),
(HeaderName::from_static("docker-upload-uuid"), uuid),
],
)
.into_response()
}
async fn patch_blob(
State(state): State<AppState>,
Path((name, uuid)): Path<(String, String)>,
body: Body,
) -> Response {
let c = canonicalize(&name, &state.config.docker);
if let Some(r) = c.denied_response() {
return r;
}
let name = c.name;
if let Err(e) = validate_docker_name(&name) {
return (StatusCode::BAD_REQUEST, e.to_string()).into_response();
}
if let Err(e) = validate_upload_uuid(&uuid) {
return (StatusCode::BAD_REQUEST, e).into_response();
}
let _inflight = match InFlightGuard::try_acquire() {
Some(g) => g,
None => return too_many_uploads(),
};
let (temp_path, current_size) = {
let mut sessions = state.upload_sessions.write();
let session = match sessions.get_mut(&uuid) {
Some(s) => s,
None => {
return (StatusCode::NOT_FOUND, "Upload session not found or expired")
.into_response();
}
};
if session.name != name {
tracing::warn!(
session_name = %session.name,
request_name = %name,
"SECURITY: upload session name mismatch — possible session fixation"
);
return (
StatusCode::BAD_REQUEST,
"Session does not belong to this repository",
)
.into_response();
}
if session.created_at.elapsed() >= SESSION_TTL {
let _ = std::fs::remove_file(&session.temp_path);
sessions.remove(&uuid);
return (StatusCode::NOT_FOUND, "Upload session expired").into_response();
}
(session.temp_path.clone(), session.size)
};
let budget =
effective_upload_cap(state.config.server.body_limit_mb).saturating_sub(current_size);
let mut file = match tokio::fs::OpenOptions::new()
.append(true)
.open(&temp_path)
.await
{
Ok(f) => f,
Err(e) if e.kind() == std::io::ErrorKind::NotFound => {
tracing::warn!(uuid = %uuid, "Upload temp file deleted by cleanup during PATCH — session race");
state.upload_sessions.write().remove(&uuid);
return (StatusCode::NOT_FOUND, "Upload session expired").into_response();
}
Err(e) => {
tracing::error!(error = %e, "Failed to open upload temp file");
state.upload_sessions.write().remove(&uuid);
return StatusCode::INTERNAL_SERVER_ERROR.into_response();
}
};
let mut guard = TempFileGuard::new(temp_path.clone());
let written = match stream_body_to_file(body, &mut file, budget).await {
StreamOutcome::Ok(n) => {
guard.disarm();
n
}
StreamOutcome::TooLarge => {
state.upload_sessions.write().remove(&uuid);
return (
StatusCode::PAYLOAD_TOO_LARGE,
"Upload session exceeds size limit",
)
.into_response();
}
StreamOutcome::ClientGone => {
state.upload_sessions.write().remove(&uuid);
return StatusCode::BAD_REQUEST.into_response();
}
StreamOutcome::Io(e) => {
tracing::error!(error = %e, "Failed to write to upload temp file");
state.upload_sessions.write().remove(&uuid);
return StatusCode::INTERNAL_SERVER_ERROR.into_response();
}
};
let new_size = current_size + written;
{
let mut sessions = state.upload_sessions.write();
match sessions.get_mut(&uuid) {
Some(session) => session.size = new_size,
None => {
tracing::warn!(uuid = %uuid, "Upload session disappeared between Phase 1 and Phase 3");
return (StatusCode::NOT_FOUND, "Upload session not found or expired")
.into_response();
}
}
}
let total_size = new_size;
let location = format!("/v2/{}/blobs/uploads/{}", name, uuid);
let range = if total_size > 0 {
format!("0-{}", total_size - 1)
} else {
"0-0".to_string()
};
(
StatusCode::ACCEPTED,
[
(header::LOCATION, location),
(header::RANGE, range),
(HeaderName::from_static("docker-upload-uuid"), uuid),
],
)
.into_response()
}
async fn cancel_upload(
State(state): State<AppState>,
Path((name, uuid)): Path<(String, String)>,
) -> Response {
let c = canonicalize(&name, &state.config.docker);
if let Some(r) = c.denied_response() {
return r;
}
let name = c.name;
if let Err(e) = validate_docker_name(&name) {
return (StatusCode::BAD_REQUEST, e.to_string()).into_response();
}
if let Err(e) = validate_upload_uuid(&uuid) {
return (StatusCode::BAD_REQUEST, e).into_response();
}
let temp_path = {
let mut sessions = state.upload_sessions.write();
match sessions.get(&uuid).map(|s| s.name == name) {
None => None,
Some(false) => {
tracing::warn!(
request_name = %name,
uuid = %uuid,
"SECURITY: upload cancel name mismatch — possible session fixation"
);
return (
StatusCode::BAD_REQUEST,
"Session does not belong to this repository",
)
.into_response();
}
Some(true) => sessions.remove(&uuid).map(|s| s.temp_path),
}
};
match temp_path {
Some(p) => {
let _ = tokio::fs::remove_file(&p).await;
StatusCode::NO_CONTENT.into_response()
}
None => (StatusCode::NOT_FOUND, "Upload session not found or expired").into_response(),
}
}
async fn upload_blob(
State(state): State<AppState>,
Path((name, uuid)): Path<(String, String)>,
axum::extract::Query(params): axum::extract::Query<std::collections::HashMap<String, String>>,
body: Body,
) -> Response {
let c = canonicalize(&name, &state.config.docker);
if let Some(r) = c.denied_response() {
return r;
}
let name = c.name;
if let Err(e) = validate_docker_name(&name) {
return (StatusCode::BAD_REQUEST, e.to_string()).into_response();
}
if let Err(e) = validate_upload_uuid(&uuid) {
return (StatusCode::BAD_REQUEST, e).into_response();
}
let _inflight = match InFlightGuard::try_acquire() {
Some(g) => g,
None => return too_many_uploads(),
};
let digest = match params.get("digest") {
Some(d) => d,
None => return (StatusCode::BAD_REQUEST, "Missing digest parameter").into_response(),
};
if let Err(e) = validate_digest(digest) {
return (StatusCode::BAD_REQUEST, e.to_string()).into_response();
}
let session_opt = {
let mut sessions = state.upload_sessions.write();
sessions.remove(&uuid)
};
if !digest.starts_with("sha256:") {
return (
StatusCode::BAD_REQUEST,
"Only sha256 digests are supported for blob uploads",
)
.into_response();
}
let cap = effective_upload_cap(state.config.server.body_limit_mb);
let (temp_path, mut guard) = if let Some(session) = session_opt {
if session.name != name {
tracing::warn!(
session_name = %session.name,
request_name = %name,
"SECURITY: upload finalization name mismatch"
);
let _ = tokio::fs::remove_file(&session.temp_path).await;
return (
StatusCode::BAD_REQUEST,
"Session does not belong to this repository",
)
.into_response();
}
let mut file = match tokio::fs::OpenOptions::new()
.create(true)
.append(true)
.open(&session.temp_path)
.await
{
Ok(f) => f,
Err(e) => {
tracing::error!(error = %e, "Failed to open temp file for PUT body");
let _ = tokio::fs::remove_file(&session.temp_path).await;
return StatusCode::INTERNAL_SERVER_ERROR.into_response();
}
};
let guard = TempFileGuard::new(session.temp_path.clone());
let budget = cap.saturating_sub(session.size);
match stream_body_to_file(body, &mut file, budget).await {
StreamOutcome::Ok(_) => {}
StreamOutcome::TooLarge => {
return (
StatusCode::PAYLOAD_TOO_LARGE,
"Upload session exceeds size limit",
)
.into_response()
}
StreamOutcome::ClientGone => return StatusCode::BAD_REQUEST.into_response(),
StreamOutcome::Io(e) => {
tracing::error!(error = %e, "Failed to append PUT body to temp file");
return StatusCode::INTERNAL_SERVER_ERROR.into_response();
}
}
(session.temp_path, guard)
} else {
let temp_dir = upload_temp_dir(&state.config.storage.path);
let temp_path = temp_dir.join(format!("mono-{}", uuid));
let mut file = match tokio::fs::File::create(&temp_path).await {
Ok(f) => f,
Err(e) => {
tracing::error!(error = %e, "Failed to create monolithic upload temp file");
return StatusCode::INTERNAL_SERVER_ERROR.into_response();
}
};
let guard = TempFileGuard::new(temp_path.clone());
match stream_body_to_file(body, &mut file, cap).await {
StreamOutcome::Ok(_) => {}
StreamOutcome::TooLarge => {
return (StatusCode::PAYLOAD_TOO_LARGE, "Upload exceeds size limit").into_response()
}
StreamOutcome::ClientGone => return StatusCode::BAD_REQUEST.into_response(),
StreamOutcome::Io(e) => {
tracing::error!(error = %e, "Failed to write monolithic upload temp file");
return StatusCode::INTERNAL_SERVER_ERROR.into_response();
}
}
(temp_path, guard)
};
{
let computed = match sha256_of_file(&temp_path).await {
Ok(h) => format!("sha256:{h}"),
Err(e) => {
tracing::error!(error = %e, "Failed to hash temp file for digest verification");
let _ = tokio::fs::remove_file(&temp_path).await;
return StatusCode::INTERNAL_SERVER_ERROR.into_response();
}
};
if computed != *digest {
tracing::warn!(
expected = %digest,
computed = %computed,
name = %name,
"SECURITY: blob digest mismatch — rejecting upload"
);
let _ = tokio::fs::remove_file(&temp_path).await;
return (
StatusCode::BAD_REQUEST,
Json(json!({
"errors": [{
"code": "DIGEST_INVALID",
"message": "provided digest did not match uploaded content",
"detail": { "expected": digest, "computed": computed }
}]
})),
)
.into_response();
}
}
let key = format!("docker/{}/blobs/{}", name, digest);
match state.storage.put_from_path(&key, &temp_path, None).await {
Ok(()) => {
guard.disarm(); state.metrics.record_upload("docker");
state.audit.log(AuditEntry::new(
"push",
"api",
&format!("{}@{}", name, digest),
"docker",
"blob",
));
state.activity.push(ActivityEntry::new(
ActionType::Push,
format!("{}@{}", name, &digest[..19.min(digest.len())]),
crate::registry_type::RegistryType::Docker,
"LOCAL",
));
state.repo_index.invalidate("docker");
blob_created(&name, digest)
}
Err(e) => {
tracing::error!(error = %e, key = %key, name = %name, "Failed to store blob");
StatusCode::INTERNAL_SERVER_ERROR.into_response()
}
}
}
async fn try_fetch_and_cache(
state: &AppState,
upstreams: &[&crate::config::DockerUpstream],
upstream_name: &str,
name: &str,
reference: &str,
cache_key: &str,
) -> Option<Response> {
for upstream in upstreams {
tracing::debug!(upstream_url = %upstream.url, upstream_name = %upstream_name, "Trying upstream");
match fetch_manifest_from_upstream(
&state.http_client,
&upstream.url,
upstream_name,
reference,
&state.docker_auth,
state.config.docker.proxy_timeout,
expose_opt(&upstream.auth),
&state.circuit_breaker,
)
.await
{
Ok((data, content_type)) => {
state.metrics.record_download("docker");
state.metrics.record_cache_miss("docker");
state.activity.push(ActivityEntry::new(
ActionType::ProxyFetch,
format!("{}:{}", name, reference),
crate::registry_type::RegistryType::Docker,
"PROXY",
));
use sha2::Digest;
let digest = format!("sha256:{}", hex::encode(sha2::Sha256::digest(&data)));
let (q_mode, q_secs) = resolve_quarantine(state);
if !matches!(q_mode, crate::digest_quarantine::QuarantineMode::Off) {
state
.digest_store
.record("docker", &digest, &upstream.url, None);
let q_status = state.digest_store.check("docker", &digest, q_secs);
match &q_status {
crate::digest_quarantine::QuarantineStatus::Mature => {}
_ => {
let outcome = if matches!(
q_mode,
crate::digest_quarantine::QuarantineMode::Enforce
) {
"blocked"
} else {
"observed"
};
crate::metrics::QUARANTINE_HOLDS_TOTAL
.with_label_values(&["docker", outcome])
.inc();
tracing::warn!(
digest = %digest,
upstream = %upstream.url,
status = %q_status.header_value(),
mode = ?q_mode,
quarantine_ttl_secs = q_secs,
"Quarantine: proxy-fetched manifest"
);
}
}
if matches!(q_mode, crate::digest_quarantine::QuarantineMode::Enforce)
&& !matches!(q_status, crate::digest_quarantine::QuarantineStatus::Mature)
{
let storage = state.storage.clone();
let key_clone = cache_key.to_string();
let repo_index = Arc::clone(&state.repo_index);
tokio::spawn(async move {
if let Err(e) = storage.put(&key_clone, &data).await {
tracing::warn!(key = %key_clone, error = %e, "cache write failed (quarantine pre-cache)");
crate::metrics::CACHE_WRITE_ERRORS
.with_label_values(&["docker", "manifest"])
.inc();
}
repo_index.invalidate("docker");
});
return Some(quarantine_forbidden(&digest, &q_status, q_secs));
}
}
let upstream_ns = Some(upstream.resolved_namespace());
let storage = state.storage.clone();
let key_clone = cache_key.to_string();
let data_clone = data.clone();
let name_clone = name.to_string();
let reference_clone = reference.to_string();
let digest_clone = digest.clone();
let repo_index = Arc::clone(&state.repo_index);
tokio::spawn(async move {
if let Err(e) = storage.put(&key_clone, &data_clone).await {
tracing::warn!(key = %key_clone, error = %e, "cache write failed (manifest by tag)");
crate::metrics::CACHE_WRITE_ERRORS
.with_label_values(&["docker", "manifest"])
.inc();
}
let digest_key =
manifest_key(upstream_ns.as_deref(), &name_clone, &digest_clone);
if let Err(e) = storage.put(&digest_key, &data_clone).await {
tracing::warn!(key = %digest_key, error = %e, "cache write failed (manifest by digest)");
crate::metrics::CACHE_WRITE_ERRORS
.with_label_values(&["docker", "manifest"])
.inc();
}
let metadata = extract_metadata(&data_clone, &storage, &name_clone).await;
if let Ok(meta_json) = serde_json::to_vec(&metadata) {
let meta_key = manifest_meta_key(
upstream_ns.as_deref(),
&name_clone,
&reference_clone,
);
if let Err(e) = storage.put(&meta_key, &meta_json).await {
tracing::warn!(key = %meta_key, error = %e, "cache write failed (metadata by tag)");
crate::metrics::CACHE_WRITE_ERRORS
.with_label_values(&["docker", "metadata"])
.inc();
}
let digest_meta_key =
manifest_meta_key(upstream_ns.as_deref(), &name_clone, &digest_clone);
if let Err(e) = storage.put(&digest_meta_key, &meta_json).await {
tracing::warn!(key = %digest_meta_key, error = %e, "cache write failed (metadata by digest)");
crate::metrics::CACHE_WRITE_ERRORS
.with_label_values(&["docker", "metadata"])
.inc();
}
}
repo_index.invalidate("docker");
});
return Some(manifest_response(data, content_type, digest));
}
Err(ProxyError::CircuitOpen(reg)) => return Some(circuit_open_response(®)),
Err(e) => {
tracing::debug!(error = ?e, upstream = %upstream.url, name = %upstream_name, reference = %reference, "Docker manifest proxy fetch failed, trying next");
continue;
}
}
}
None
}
fn manifest_cache_fresh(
is_digest: bool,
has_upstream: bool,
metadata_ttl: i64,
modified: Option<u64>,
) -> bool {
if is_digest || !has_upstream {
return true;
}
metadata_ttl > 0
&& modified
.map(|m| crate::cache_ttl::is_within_ttl(m, metadata_ttl))
.unwrap_or(false)
}
async fn get_manifest(
State(state): State<AppState>,
headers: axum::http::HeaderMap,
Path((name, reference)): Path<(String, String)>,
) -> Response {
let c = canonicalize(&name, &state.config.docker);
if let Some(r) = c.denied_response() {
return r;
}
let upstreams_to_try = c.upstreams_to_try(&state.config.docker.upstreams);
let ns = c.namespace;
let name = c.name;
if let Err(e) = validate_docker_name(&name) {
return (StatusCode::BAD_REQUEST, e.to_string()).into_response();
}
if let Err(e) = validate_docker_reference(&reference) {
return (StatusCode::BAD_REQUEST, e.to_string()).into_response();
}
let publish_date = extract_docker_publish_date(
&state.storage,
&name,
&reference,
state.config.docker.upstreams.is_empty(),
ns.as_deref(),
)
.await;
let internal = crate::curation::is_internal_namespace(
&state.curation().curation_engine,
crate::curation::RegistryType::Docker,
&name,
);
if !internal {
if let Some(response) = crate::curation::check_download(
&state.curation().curation_engine,
state.bypass_token().as_deref(),
&headers,
crate::curation::RegistryType::Docker,
&name,
Some(&reference),
publish_date,
) {
return response;
}
}
let key = manifest_key(ns.as_deref(), &name, &reference);
let legacy_key = manifest_key(None, &name, &reference);
use crate::storage::StorageError;
let (cached, hosted) = match state.storage.get(&key).await {
Ok(data) => (Some(data), key == legacy_key),
Err(StorageError::NotFound) if key != legacy_key => {
match state.storage.get(&legacy_key).await {
Ok(data) => (Some(data), true),
Err(StorageError::NotFound) => (None, true),
Err(e) => {
tracing::error!(error = %e, key = %legacy_key, "manifest read failed");
return StatusCode::INTERNAL_SERVER_ERROR.into_response();
}
}
}
Err(StorageError::NotFound) => (None, true),
Err(e) => {
tracing::error!(error = %e, key = %key, "manifest read failed");
return StatusCode::INTERNAL_SERVER_ERROR.into_response();
}
};
let is_digest = reference.starts_with("sha256:");
let revalidate = !hosted && !upstreams_to_try.is_empty();
let cache_fresh = if cached.is_some() {
let modified = if revalidate && !is_digest {
storage_stat_with_fallback(&state.storage, &key, &legacy_key)
.await
.map(|m| m.modified)
} else {
None
};
manifest_cache_fresh(
is_digest,
revalidate,
state.config.docker.metadata_ttl,
modified,
)
} else {
false
};
if let Some(ref data) = cached {
if cache_fresh {
return serve_cached_manifest(&state, data, &name, &reference, ns.as_deref());
}
}
if internal {
if let Some(ref data) = cached {
return serve_cached_manifest(&state, data, &name, &reference, ns.as_deref());
}
crate::metrics::record_namespace_isolation_refused("docker");
return manifest_unknown_response(&name, &reference);
}
tracing::debug!(
upstreams_count = upstreams_to_try.len(),
"Trying upstream proxies"
);
if let Some(response) =
try_fetch_and_cache(&state, &upstreams_to_try, &name, &name, &reference, &key).await
{
return response;
}
if !name.contains('/') {
let library_name = format!("library/{}", name);
if let Some(response) = try_fetch_and_cache(
&state,
&upstreams_to_try,
&library_name,
&name,
&reference,
&key,
)
.await
{
return response;
}
}
if let Some(ref data) = cached {
if state.config.docker.serve_stale {
tracing::warn!(
registry = "docker",
name = %name,
reference = %reference,
"Upstream failed, serving stale cached manifest"
);
let mut response =
serve_cached_manifest(&state, data, &name, &reference, ns.as_deref());
response.headers_mut().insert(
axum::http::header::HeaderName::from_static("x-nora-stale"),
axum::http::header::HeaderValue::from_static("true"),
);
response.headers_mut().insert(
axum::http::header::CACHE_CONTROL,
axum::http::header::HeaderValue::from_static("public, max-age=0, must-revalidate"),
);
return response;
}
}
if !state.config.docker.upstreams.is_empty() {
tracing::warn!(registry = "docker", name = %name, reference = %reference, "Proxy failed, returning 404");
}
StatusCode::NOT_FOUND.into_response()
}
fn serve_cached_manifest(
state: &AppState,
data: &[u8],
name: &str,
reference: &str,
ns: Option<&str>,
) -> Response {
if let Some(response) = crate::curation::verify_integrity(
&state.curation().curation_engine,
crate::curation::RegistryType::Docker,
name,
Some(reference),
data,
) {
return response;
}
state.metrics.record_download("docker");
state.metrics.record_cache_hit("docker");
state.activity.push(ActivityEntry::new(
ActionType::Pull,
format!("{}:{}", name, reference),
crate::registry_type::RegistryType::Docker,
"LOCAL",
));
use sha2::Digest;
let digest = format!("sha256:{}", hex::encode(sha2::Sha256::digest(data)));
if let Some(resp) = quarantine_cache_serve_gate(state, &digest) {
return resp;
}
let content_type = detect_manifest_media_type(data);
let meta_key = manifest_meta_key(ns, name, reference);
let state_clone = state.clone();
let storage_clone = state.storage.clone();
tokio::spawn(update_metadata_on_pull(
state_clone,
storage_clone,
meta_key,
));
manifest_response(Bytes::copy_from_slice(data), content_type, digest)
}
async fn missing_manifest_ref(body: &[u8], storage: &Storage, name: &str) -> Option<String> {
let json = serde_json::from_slice::<serde_json::Value>(body).ok()?;
if let Some(manifests) = json.get("manifests").and_then(|v| v.as_array()) {
for m in manifests {
if let Some(d) = m.get("digest").and_then(|v| v.as_str()) {
if storage
.stat(&format!("docker/{}/manifests/{}.json", name, d))
.await
.is_none()
{
return Some(d.to_string());
}
}
}
return None;
}
let mut refs: Vec<&str> = Vec::new();
if let Some(d) = json
.get("config")
.and_then(|c| c.get("digest"))
.and_then(|v| v.as_str())
{
refs.push(d);
}
if let Some(layers) = json.get("layers").and_then(|v| v.as_array()) {
for l in layers {
if let Some(d) = l.get("digest").and_then(|v| v.as_str()) {
refs.push(d);
}
}
}
for d in refs {
if storage
.stat(&format!("docker/{}/blobs/{}", name, d))
.await
.is_none()
{
return Some(d.to_string());
}
}
None
}
async fn put_manifest(
State(state): State<AppState>,
Path((name, reference)): Path<(String, String)>,
body: Bytes,
) -> Response {
let c = canonicalize(&name, &state.config.docker);
if let Some(r) = c.denied_response() {
return r;
}
let name = c.name;
let ns = c.namespace;
if let Err(e) = validate_docker_name(&name) {
return (StatusCode::BAD_REQUEST, e.to_string()).into_response();
}
if let Err(e) = validate_docker_reference(&reference) {
return (StatusCode::BAD_REQUEST, e.to_string()).into_response();
}
use sha2::Digest;
let digest = format!("sha256:{}", hex::encode(sha2::Sha256::digest(&body)));
if let Some(missing) = missing_manifest_ref(&body, &state.storage, &name).await {
tracing::warn!(
name = %name,
reference = %reference,
missing = %missing,
"rejecting manifest push: references an absent blob/sub-manifest"
);
return (
StatusCode::BAD_REQUEST,
Json(json!({
"errors": [{
"code": "MANIFEST_BLOB_UNKNOWN",
"message": "manifest references an unknown blob",
"detail": { "digest": missing }
}]
})),
)
.into_response();
}
let key = format!("docker/{}/manifests/{}.json", name, reference);
let manifest_lock = state.publish_lock(&manifest_key(ns.as_deref(), &name, &reference));
let _manifest_guard = manifest_lock.lock().await;
if state.storage.put(&key, &body).await.is_err() {
return StatusCode::INTERNAL_SERVER_ERROR.into_response();
}
let digest_key = format!("docker/{}/manifests/{}.json", name, digest);
if state.storage.put(&digest_key, &body).await.is_err() {
return StatusCode::INTERNAL_SERVER_ERROR.into_response();
}
let metadata = extract_metadata(&body, &state.storage, &name).await;
let meta_key = format!("docker/{}/manifests/{}.meta.json", name, reference);
if let Ok(meta_json) = serde_json::to_vec(&metadata) {
if let Err(e) = state.storage.put(&meta_key, &meta_json).await {
tracing::warn!(key = %meta_key, error = %e, "cache write failed (push metadata by tag)");
crate::metrics::CACHE_WRITE_ERRORS
.with_label_values(&["docker", "metadata"])
.inc();
}
let digest_meta_key = format!("docker/{}/manifests/{}.meta.json", name, digest);
if let Err(e) = state.storage.put(&digest_meta_key, &meta_json).await {
tracing::warn!(key = %digest_meta_key, error = %e, "cache write failed (push metadata by digest)");
crate::metrics::CACHE_WRITE_ERRORS
.with_label_values(&["docker", "metadata"])
.inc();
}
}
state.metrics.record_upload("docker");
state.activity.push(ActivityEntry::new(
ActionType::Push,
format!("{}:{}", name, reference),
crate::registry_type::RegistryType::Docker,
"LOCAL",
));
state.audit.log(AuditEntry::new(
"push",
"api",
&format!("{}:{}", name, reference),
"docker",
"manifest",
));
state.repo_index.invalidate("docker");
let location = format!("/v2/{}/manifests/{}", name, reference);
(
StatusCode::CREATED,
[
(header::LOCATION, location),
(HeaderName::from_static("docker-content-digest"), digest),
],
)
.into_response()
}
async fn list_tags(State(state): State<AppState>, Path(name): Path<String>) -> Response {
let c = canonicalize(&name, &state.config.docker);
if let Some(r) = c.denied_response() {
return r;
}
let ns = c.namespace;
let name = c.name;
if let Err(e) = validate_docker_name(&name) {
return (StatusCode::BAD_REQUEST, e.to_string()).into_response();
}
let prefix = manifest_prefix(ns.as_deref(), &name);
let legacy_prefix = manifest_prefix(None, &name);
let mut keys = match state.storage.list(&prefix).await {
Ok(k) => k,
Err(e) => {
tracing::error!(error = ?e, "docker: failed to list manifests for tags");
return StatusCode::SERVICE_UNAVAILABLE.into_response();
}
};
if prefix != legacy_prefix {
match state.storage.list(&legacy_prefix).await {
Ok(legacy_keys) => {
keys.extend(legacy_keys);
keys.sort();
keys.dedup();
}
Err(e) => {
tracing::warn!(error = ?e, "docker: failed to list legacy manifests, continuing with namespaced only");
}
}
}
let tags: Vec<String> = keys
.iter()
.filter_map(|k| {
k.strip_prefix(&prefix)
.or_else(|| k.strip_prefix(&legacy_prefix))
.and_then(|t| t.strip_suffix(".json"))
.map(String::from)
})
.filter(|t| !ends_with_ci(t, ".meta") && !t.contains(".meta."))
.collect();
(StatusCode::OK, Json(json!({"name": name, "tags": tags}))).into_response()
}
async fn delete_tags_for_digest(state: &AppState, ns: Option<&str>, name: &str, digest: &str) {
use sha2::Digest as _;
let prefix = manifest_prefix(ns, name);
let legacy_prefix = manifest_prefix(None, name);
let mut keys = state.storage.list(&prefix).await.unwrap_or_default();
if prefix != legacy_prefix {
if let Ok(legacy) = state.storage.list(&legacy_prefix).await {
keys.extend(legacy);
}
}
let mut tags: Vec<String> = keys
.iter()
.filter_map(|k| {
k.strip_prefix(&prefix)
.or_else(|| k.strip_prefix(&legacy_prefix))
.and_then(|t| t.strip_suffix(".json"))
.map(String::from)
})
.filter(|t| !t.starts_with("sha256:") && !ends_with_ci(t, ".meta") && !t.contains(".meta."))
.collect();
tags.sort();
tags.dedup();
for tag in tags {
let key = manifest_key(ns, name, &tag);
let legacy_key = manifest_key(None, name, &tag);
let lock = state.publish_lock(&key);
let _guard = lock.lock().await;
let data = match storage_get_with_fallback(&state.storage, &key, &legacy_key).await {
Ok(d) => d,
Err(_) => continue,
};
let resolved = format!("sha256:{}", hex::encode(sha2::Sha256::digest(&data)));
if resolved != digest {
continue;
}
let _ = state.storage.delete(&key).await;
let _ = state.storage.delete(&legacy_key).await;
let _ = state
.storage
.delete(&manifest_meta_key(ns, name, &tag))
.await;
let _ = state
.storage
.delete(&manifest_meta_key(None, name, &tag))
.await;
tracing::info!(
name = %name, tag = %tag, digest = %digest,
"Docker tag removed because its manifest was deleted by digest (#658)"
);
}
}
async fn delete_manifest(
State(state): State<AppState>,
Path((name, reference)): Path<(String, String)>,
) -> Response {
let c = canonicalize(&name, &state.config.docker);
if let Some(r) = c.denied_response() {
return r;
}
let ns = c.namespace;
let name = c.name;
if let Err(e) = validate_docker_name(&name) {
return (StatusCode::BAD_REQUEST, e.to_string()).into_response();
}
if let Err(e) = validate_docker_reference(&reference) {
return (StatusCode::BAD_REQUEST, e.to_string()).into_response();
}
let key = manifest_key(ns.as_deref(), &name, &reference);
let legacy_key = manifest_key(None, &name, &reference);
let lock = state.publish_lock(&key);
let _guard = lock.lock().await;
let is_tag = !reference.starts_with("sha256:");
if is_tag {
if let Ok(data) = storage_get_with_fallback(&state.storage, &key, &legacy_key).await {
use sha2::Digest;
let digest = format!("sha256:{}", hex::encode(sha2::Sha256::digest(&data)));
let _ = state
.storage
.delete(&manifest_key(ns.as_deref(), &name, &digest))
.await;
let _ = state
.storage
.delete(&manifest_key(None, &name, &digest))
.await;
let _ = state
.storage
.delete(&manifest_meta_key(ns.as_deref(), &name, &digest))
.await;
let _ = state
.storage
.delete(&manifest_meta_key(None, &name, &digest))
.await;
}
} else {
delete_tags_for_digest(&state, ns.as_deref(), &name, &reference).await;
}
match state.storage.delete(&key).await {
Ok(()) => {
let meta_key = manifest_meta_key(ns.as_deref(), &name, &reference);
let _ = state
.storage
.delete(&manifest_meta_key(None, &name, &reference))
.await;
let _ = state.storage.delete(&meta_key).await;
state.audit.log(AuditEntry::new(
"delete",
"api",
&format!("{}:{}", name, reference),
"docker",
"manifest",
));
state.repo_index.invalidate("docker");
tracing::info!(name = %name, reference = %reference, "Docker manifest deleted");
StatusCode::ACCEPTED.into_response()
}
Err(crate::storage::StorageError::NotFound) if key != legacy_key => {
match state.storage.delete(&legacy_key).await {
Ok(()) => {
let _ = state
.storage
.delete(&manifest_meta_key(None, &name, &reference))
.await;
state.audit.log(AuditEntry::new(
"delete",
"api",
&format!("{}:{}", name, reference),
"docker",
"manifest",
));
state.repo_index.invalidate("docker");
tracing::info!(name = %name, reference = %reference, "Docker manifest deleted (legacy key)");
StatusCode::ACCEPTED.into_response()
}
_ => manifest_unknown_response(&name, &reference),
}
}
Err(crate::storage::StorageError::NotFound) => manifest_unknown_response(&name, &reference),
Err(e) => {
tracing::error!(error = %e, key = %key, name = %name, reference = %reference, "Failed to delete manifest");
StatusCode::INTERNAL_SERVER_ERROR.into_response()
}
}
}
async fn delete_blob(
State(state): State<AppState>,
Path((name, digest)): Path<(String, String)>,
) -> Response {
let c = canonicalize(&name, &state.config.docker);
if let Some(r) = c.denied_response() {
return r;
}
let ns = c.namespace;
let name = c.name;
if let Err(e) = validate_docker_name(&name) {
return (StatusCode::BAD_REQUEST, e.to_string()).into_response();
}
if let Err(e) = validate_digest(&digest) {
return (StatusCode::BAD_REQUEST, e.to_string()).into_response();
}
let key = blob_key(ns.as_deref(), &name, &digest);
let legacy_key = blob_key(None, &name, &digest);
if key != legacy_key {
let _ = state.storage.delete(&legacy_key).await;
}
match state.storage.delete(&key).await {
Ok(()) => {
state.audit.log(AuditEntry::new(
"delete",
"api",
&format!("{}@{}", name, &digest[..19.min(digest.len())]),
"docker",
"blob",
));
state.repo_index.invalidate("docker");
tracing::info!(name = %name, digest = %digest, "Docker blob deleted");
StatusCode::ACCEPTED.into_response()
}
Err(crate::storage::StorageError::NotFound) => blob_unknown_response(&digest),
Err(e) => {
tracing::error!(error = %e, key = %key, name = %name, digest = %digest, "Failed to delete blob");
StatusCode::INTERNAL_SERVER_ERROR.into_response()
}
}
}
pub struct FetchedBlob {
pub path: std::path::PathBuf,
pub sha256: String,
#[allow(dead_code)]
pub content_length: Option<u64>,
pub _guard: TempFileGuard,
}
#[allow(clippy::too_many_arguments)]
pub async fn fetch_blob_from_upstream(
client: &reqwest::Client,
upstream_url: &str,
name: &str,
digest: &str,
docker_auth: &DockerAuth,
timeout: u64,
read_timeout: u64,
basic_auth: Option<&str>,
cb: &CircuitBreakerRegistry,
temp_dir: &std::path::Path,
) -> Result<FetchedBlob, ProxyError> {
use crate::metrics::{PROXY_ACTIVE_DOWNLOADS, PROXY_DOWNLOAD_BYTES};
PROXY_ACTIVE_DOWNLOADS.inc();
let _download_gauge_guard = ProxyDownloadGuard;
tracing::info!(
blob.name = %name,
blob.digest = %digest,
upstream = %upstream_url,
"Proxy blob download started"
);
let cb_key = format!("docker:{}", upstream_url.trim_end_matches('/'));
let probe = cb.check(&cb_key)?;
let url = format!(
"{}/v2/{}/blobs/{}",
upstream_url.trim_end_matches('/'),
name,
digest
);
let mut request = client.get(&url).timeout(Duration::from_secs(timeout));
if let Some(credentials) = basic_auth {
request = request.header("Authorization", basic_auth_header(credentials));
}
let response = request.send().await.map_err(|e| {
cb.record_failure(&cb_key, probe);
ProxyError::Network(e.to_string())
})?;
let response = if response.status() == reqwest::StatusCode::UNAUTHORIZED {
let www_auth = response
.headers()
.get("www-authenticate")
.and_then(|v| v.to_str().ok())
.map(String::from);
if let Some(token) = docker_auth
.get_token(upstream_url, name, www_auth.as_deref(), basic_auth)
.await
{
client
.get(&url)
.header("Authorization", format!("Bearer {}", token))
.send()
.await
.map_err(|e| {
cb.record_failure(&cb_key, probe);
ProxyError::Network(e.to_string())
})?
} else {
return Err(ProxyError::Network("token fetch failed".into()));
}
} else {
response
};
if !response.status().is_success() {
let status = response.status().as_u16();
if (400..500).contains(&status) {
cb.record_alive(&cb_key, probe);
} else {
cb.record_failure(&cb_key, probe);
}
return Err(ProxyError::Upstream(status));
}
let upstream_content_length = response.content_length();
let temp_path = temp_dir.join(format!("proxy-{}", uuid::Uuid::new_v4()));
let guard = TempFileGuard::new(temp_path.clone());
let mut file = tokio::fs::File::create(&temp_path)
.await
.map_err(|e| ProxyError::Network(format!("temp file create: {}", e)))?;
use sha2::Digest;
let mut stream = response.bytes_stream();
let mut hasher = sha2::Sha256::new();
let chunk_timeout = Duration::from_secs(read_timeout);
let mut bytes_written: u64 = 0;
loop {
match tokio::time::timeout(chunk_timeout, stream.next()).await {
Ok(Some(Ok(chunk))) => {
hasher.update(&chunk);
use tokio::io::AsyncWriteExt;
file.write_all(&chunk)
.await
.map_err(|e| ProxyError::Network(format!("temp file write: {}", e)))?;
bytes_written += chunk.len() as u64;
}
Ok(Some(Err(e))) => {
cb.record_failure(&cb_key, probe);
return Err(ProxyError::Network(format!("chunk read error: {}", e)));
}
Ok(None) => break, Err(_) => {
cb.record_failure(&cb_key, probe);
return Err(ProxyError::Network(format!(
"read timeout ({}s per chunk)",
read_timeout
)));
}
}
}
use tokio::io::AsyncWriteExt;
file.flush()
.await
.map_err(|e| ProxyError::Network(format!("temp file flush: {}", e)))?;
drop(file);
let sha256 = hex::encode(sha2::Digest::finalize(hasher));
cb.record_success(&cb_key, probe);
PROXY_DOWNLOAD_BYTES.inc_by(bytes_written);
tracing::info!(
bytes = bytes_written,
content_length = ?upstream_content_length,
"Proxy blob download complete"
);
Ok(FetchedBlob {
path: temp_path,
sha256,
content_length: upstream_content_length,
_guard: guard,
})
}
#[allow(clippy::too_many_arguments)]
pub async fn fetch_manifest_from_upstream(
client: &reqwest::Client,
upstream_url: &str,
name: &str,
reference: &str,
docker_auth: &DockerAuth,
timeout: u64,
basic_auth: Option<&str>,
cb: &CircuitBreakerRegistry,
) -> Result<(Vec<u8>, String), ProxyError> {
let cb_key = format!("docker:{}", upstream_url.trim_end_matches('/'));
let probe = cb.check(&cb_key)?;
let url = format!(
"{}/v2/{}/manifests/{}",
upstream_url.trim_end_matches('/'),
name,
reference
);
tracing::debug!(url = %url, "Fetching manifest from upstream");
let accept_header = "application/vnd.docker.distribution.manifest.v2+json, \
application/vnd.docker.distribution.manifest.list.v2+json, \
application/vnd.oci.image.manifest.v1+json, \
application/vnd.oci.image.index.v1+json";
let mut request = client
.get(&url)
.timeout(Duration::from_secs(timeout))
.header("Accept", accept_header);
if let Some(credentials) = basic_auth {
request = request.header("Authorization", basic_auth_header(credentials));
}
let response = request.send().await.map_err(|e| {
tracing::error!(error = %e, url = %url, "Failed to send request to upstream");
cb.record_failure(&cb_key, probe);
ProxyError::Network(e.to_string())
})?;
tracing::debug!(status = %response.status(), "Initial upstream response");
let response = if response.status() == reqwest::StatusCode::UNAUTHORIZED {
let www_auth = response
.headers()
.get("www-authenticate")
.and_then(|v| v.to_str().ok())
.map(String::from);
tracing::debug!(www_auth = ?www_auth, "Got 401, fetching token");
if let Some(token) = docker_auth
.get_token(upstream_url, name, www_auth.as_deref(), basic_auth)
.await
{
tracing::debug!("Token acquired, retrying with auth");
client
.get(&url)
.header("Accept", accept_header)
.header("Authorization", format!("Bearer {}", token))
.send()
.await
.map_err(|e| {
tracing::error!(error = %e, "Failed to send authenticated request");
cb.record_failure(&cb_key, probe);
ProxyError::Network(e.to_string())
})?
} else {
tracing::error!("Failed to acquire token");
return Err(ProxyError::Network("token fetch failed".into()));
}
} else {
response
};
tracing::debug!(status = %response.status(), "Final upstream response");
if !response.status().is_success() {
let status = response.status().as_u16();
tracing::warn!(status = %response.status(), "Upstream returned non-success status");
if (400..500).contains(&status) {
cb.record_alive(&cb_key, probe);
} else {
cb.record_failure(&cb_key, probe);
}
return Err(ProxyError::Upstream(status));
}
let content_type = response
.headers()
.get("content-type")
.and_then(|v| v.to_str().ok())
.unwrap_or("application/vnd.docker.distribution.manifest.v2+json")
.to_string();
let bytes = response.bytes().await.map_err(|e| {
cb.record_failure(&cb_key, probe);
ProxyError::Network(e.to_string())
})?;
cb.record_success(&cb_key, probe);
Ok((bytes.to_vec(), content_type))
}
fn manifest_response(data: impl Into<Bytes>, content_type: String, digest: String) -> Response {
let body: Bytes = data.into();
let content_length = body.len().to_string();
(
StatusCode::OK,
[
(header::CONTENT_TYPE, content_type),
(HeaderName::from_static("docker-content-digest"), digest),
(header::CONTENT_LENGTH, content_length),
],
body,
)
.into_response()
}
fn detect_manifest_media_type(data: &[u8]) -> String {
if let Ok(json) = serde_json::from_slice::<Value>(data) {
if let Some(media_type) = json.get("mediaType").and_then(|v| v.as_str()) {
return media_type.to_string();
}
if let Some(schema_version) = json.get("schemaVersion").and_then(|v| v.as_u64()) {
if schema_version == 1 {
return "application/vnd.docker.distribution.manifest.v1+json".to_string();
}
if let Some(config) = json.get("config") {
if let Some(config_mt) = config.get("mediaType").and_then(|v| v.as_str()) {
if config_mt.starts_with("application/vnd.docker.") {
return "application/vnd.docker.distribution.manifest.v2+json".to_string();
}
return "application/vnd.oci.image.manifest.v1+json".to_string();
}
return "application/vnd.docker.distribution.manifest.v2+json".to_string();
}
if json.get("manifests").is_some() {
return "application/vnd.oci.image.index.v1+json".to_string();
}
}
}
"application/vnd.docker.distribution.manifest.v2+json".to_string()
}
async fn extract_docker_publish_date(
storage: &Storage,
name: &str,
reference: &str,
upstreams_empty: bool,
ns: Option<&str>,
) -> Option<i64> {
let meta = manifest_meta_key(ns, name, reference);
let legacy_meta = manifest_meta_key(None, name, reference);
if let Ok(data) = storage_get_with_fallback(storage, &meta, &legacy_meta).await {
if let Ok(meta) = serde_json::from_slice::<ImageMetadata>(&data) {
if meta.push_timestamp > 0 {
return Some(meta.push_timestamp as i64);
}
}
}
if upstreams_empty {
let key = manifest_key(ns, name, reference);
let legacy = manifest_key(None, name, reference);
if let Some(date) = crate::curation::extract_mtime_as_publish_date(storage, &key).await {
return Some(date);
}
if key != legacy {
return crate::curation::extract_mtime_as_publish_date(storage, &legacy).await;
}
}
None
}
async fn extract_metadata(manifest: &[u8], storage: &Storage, name: &str) -> ImageMetadata {
let now = SystemTime::now()
.duration_since(UNIX_EPOCH)
.unwrap_or_default()
.as_secs();
let mut metadata = ImageMetadata {
push_timestamp: now,
last_pulled: 0,
downloads: 0,
..Default::default()
};
let Ok(json) = serde_json::from_slice::<Value>(manifest) else {
return metadata;
};
if json.get("manifests").is_some() {
if let Some(manifests) = json.get("manifests").and_then(|m| m.as_array()) {
let total_size: u64 = manifests
.iter()
.filter_map(|m| m.get("size").and_then(|s| s.as_u64()))
.sum();
metadata.size_bytes = total_size;
if let Some(first) = manifests.first() {
if let Some(platform) = first.get("platform") {
metadata.os = platform
.get("os")
.and_then(|v| v.as_str())
.unwrap_or("multi-arch")
.to_string();
metadata.arch = platform
.get("architecture")
.and_then(|v| v.as_str())
.unwrap_or("multi")
.to_string();
metadata.variant = platform
.get("variant")
.and_then(|v| v.as_str())
.map(String::from);
}
}
}
return metadata;
}
if let Some(layers) = json.get("layers").and_then(|l| l.as_array()) {
let mut total_size: u64 = 0;
for layer in layers {
let digest = layer
.get("digest")
.and_then(|d| d.as_str())
.unwrap_or("")
.to_string();
let size = layer.get("size").and_then(|s| s.as_u64()).unwrap_or(0);
total_size += size;
metadata.layers.push(LayerInfo { digest, size });
}
metadata.size_bytes = total_size;
}
if let Some(config) = json.get("config") {
if let Some(config_digest) = config.get("digest").and_then(|d| d.as_str()) {
let (os, arch, variant) = get_config_info(storage, name, config_digest).await;
metadata.os = os;
metadata.arch = arch;
metadata.variant = variant;
}
}
if metadata.os.is_empty() {
metadata.os = "unknown".to_string();
}
if metadata.arch.is_empty() {
metadata.arch = "unknown".to_string();
}
metadata
}
async fn get_config_info(
storage: &Storage,
name: &str,
config_digest: &str,
) -> (String, String, Option<String>) {
let key = format!("docker/{}/blobs/{}", name, config_digest);
let Ok(data) = storage.get(&key).await else {
return ("unknown".to_string(), "unknown".to_string(), None);
};
let Ok(config) = serde_json::from_slice::<Value>(&data) else {
return ("unknown".to_string(), "unknown".to_string(), None);
};
let os = config
.get("os")
.and_then(|v| v.as_str())
.unwrap_or("unknown")
.to_string();
let arch = config
.get("architecture")
.and_then(|v| v.as_str())
.unwrap_or("unknown")
.to_string();
let variant = config
.get("variant")
.and_then(|v| v.as_str())
.map(String::from);
(os, arch, variant)
}
async fn update_metadata_on_pull(state: AppState, storage: Storage, meta_key: String) {
let lock = state.publish_lock(&meta_key);
let _guard = lock.lock().await;
let now = SystemTime::now()
.duration_since(UNIX_EPOCH)
.unwrap_or_default()
.as_secs();
let mut metadata = if let Ok(data) = storage.get(&meta_key).await {
serde_json::from_slice::<ImageMetadata>(&data).unwrap_or_default()
} else {
ImageMetadata::default()
};
metadata.downloads += 1;
metadata.last_pulled = now;
if let Ok(json) = serde_json::to_vec(&metadata) {
if let Err(e) = storage.put(&meta_key, &json).await {
tracing::warn!(key = %meta_key, error = %e, "cache write failed (pull stats update)");
crate::metrics::CACHE_WRITE_ERRORS
.with_label_values(&["docker", "metadata"])
.inc();
}
}
}
#[cfg(test)]
#[allow(clippy::unwrap_used)]
mod tests {
use super::*;
#[tokio::test]
async fn verifying_reader_passes_match_and_aborts_mismatch() {
use sha2::Digest as _;
use tokio::io::AsyncReadExt;
let data = b"hello-blob-content-1234567890";
let good = format!("sha256:{}", hex::encode(sha2::Sha256::digest(data)));
let mut ok = VerifyingReader::new(&data[..], &good);
let mut buf = Vec::new();
assert!(ok.read_to_end(&mut buf).await.is_ok());
assert_eq!(buf, data);
let bad = format!(
"sha256:{}",
hex::encode(sha2::Sha256::digest(b"different-bytes"))
);
let mut tampered = VerifyingReader::new(&data[..], &bad);
let mut buf2 = Vec::new();
assert!(
tampered.read_to_end(&mut buf2).await.is_err(),
"a digest mismatch must error the stream, not complete cleanly"
);
}
#[test]
fn test_image_metadata_default() {
let meta = ImageMetadata::default();
assert_eq!(meta.push_timestamp, 0);
assert_eq!(meta.last_pulled, 0);
assert_eq!(meta.downloads, 0);
assert_eq!(meta.size_bytes, 0);
assert_eq!(meta.os, "");
assert_eq!(meta.arch, "");
assert!(meta.variant.is_none());
assert!(meta.layers.is_empty());
}
#[test]
fn test_image_metadata_serialization() {
let meta = ImageMetadata {
push_timestamp: 1700000000,
last_pulled: 1700001000,
downloads: 42,
size_bytes: 1024000,
os: "linux".to_string(),
arch: "amd64".to_string(),
variant: None,
layers: vec![LayerInfo {
digest: "sha256:abc123".to_string(),
size: 512000,
}],
};
let json = serde_json::to_string(&meta).unwrap();
assert!(json.contains("\"os\":\"linux\""));
assert!(json.contains("\"arch\":\"amd64\""));
assert!(!json.contains("variant")); }
#[test]
fn test_image_metadata_with_variant() {
let meta = ImageMetadata {
variant: Some("v8".to_string()),
..Default::default()
};
let json = serde_json::to_string(&meta).unwrap();
assert!(json.contains("\"variant\":\"v8\""));
}
#[test]
fn test_image_metadata_deserialization() {
let json = r#"{
"push_timestamp": 1700000000,
"last_pulled": 0,
"downloads": 5,
"size_bytes": 2048,
"os": "linux",
"arch": "arm64",
"variant": "v8",
"layers": [
{"digest": "sha256:aaa", "size": 1024},
{"digest": "sha256:bbb", "size": 1024}
]
}"#;
let meta: ImageMetadata = serde_json::from_str(json).unwrap();
assert_eq!(meta.os, "linux");
assert_eq!(meta.arch, "arm64");
assert_eq!(meta.variant, Some("v8".to_string()));
assert_eq!(meta.layers.len(), 2);
assert_eq!(meta.layers[0].digest, "sha256:aaa");
assert_eq!(meta.layers[1].size, 1024);
}
#[test]
fn test_layer_info_serialization_roundtrip() {
let layer = LayerInfo {
digest: "sha256:deadbeef".to_string(),
size: 999999,
};
let json = serde_json::to_value(&layer).unwrap();
let restored: LayerInfo = serde_json::from_value(json).unwrap();
assert_eq!(layer.digest, restored.digest);
assert_eq!(layer.size, restored.size);
}
#[test]
fn test_cleanup_expired_sessions_empty() {
let sessions: RwLock<HashMap<String, UploadSession>> = RwLock::new(HashMap::new());
cleanup_expired_sessions(&sessions);
assert_eq!(sessions.read().len(), 0);
}
#[test]
fn test_cleanup_expired_sessions_fresh() {
let sessions: RwLock<HashMap<String, UploadSession>> = RwLock::new(HashMap::new());
let temp_dir = tempfile::TempDir::new().unwrap();
let temp_path = temp_dir.path().join("uuid-1");
std::fs::write(&temp_path, b"test data").unwrap();
sessions.write().insert(
"uuid-1".to_string(),
UploadSession {
temp_path,
size: 9,
name: "test/image".to_string(),
created_at: std::time::Instant::now(),
},
);
cleanup_expired_sessions(&sessions);
assert_eq!(sessions.read().len(), 1); }
#[test]
fn test_max_upload_sessions_default() {
let max = max_upload_sessions();
assert!(max > 0);
assert_eq!(max, DEFAULT_MAX_UPLOAD_SESSIONS);
}
#[tokio::test]
async fn test_cancel_upload_frees_session_and_temp() {
use crate::test_helpers::{create_test_context, send};
let ctx = create_test_context();
let resp = send(
&ctx.app,
Method::POST,
"/v2/alpine/blobs/uploads/",
Body::empty(),
)
.await;
assert_eq!(resp.status(), StatusCode::ACCEPTED);
let location = resp.headers()[header::LOCATION]
.to_str()
.unwrap()
.to_string();
let temp_path = {
let sessions = ctx.state.upload_sessions.read();
assert_eq!(sessions.len(), 1);
sessions.values().next().unwrap().temp_path.clone()
};
assert!(temp_path.exists());
let resp = send(&ctx.app, Method::DELETE, &location, Body::empty()).await;
assert_eq!(
resp.status(),
StatusCode::NO_CONTENT,
"cancel must be honoured, not 405 — a refused cancel holds a session slot until SESSION_TTL"
);
assert_eq!(ctx.state.upload_sessions.read().len(), 0);
assert!(!temp_path.exists(), "cancel must reclaim the temp file");
let resp = send(&ctx.app, Method::DELETE, &location, Body::empty()).await;
assert_eq!(resp.status(), StatusCode::NOT_FOUND);
}
#[tokio::test]
async fn test_upload_session_limit_returns_oci_429() {
use crate::test_helpers::{body_bytes, create_test_context, send};
let ctx = create_test_context();
{
let mut sessions = ctx.state.upload_sessions.write();
for i in 0..max_upload_sessions() {
sessions.insert(
format!("fill-{i}"),
UploadSession {
temp_path: std::env::temp_dir().join(format!("fill-{i}")),
size: 0,
name: "alpine".to_string(),
created_at: std::time::Instant::now(),
},
);
}
}
let resp = send(
&ctx.app,
Method::POST,
"/v2/alpine/blobs/uploads/",
Body::empty(),
)
.await;
assert_eq!(resp.status(), StatusCode::TOO_MANY_REQUESTS);
assert_eq!(
resp.headers()[header::CONTENT_TYPE].to_str().unwrap(),
"application/json"
);
assert!(
resp.headers().contains_key(header::RETRY_AFTER),
"429 must carry Retry-After so clients back off"
);
let body: serde_json::Value = serde_json::from_slice(&body_bytes(resp).await).unwrap();
assert_eq!(body["errors"][0]["code"], "TOOMANYREQUESTS");
assert_eq!(
body["errors"][0]["detail"]["limit"],
max_upload_sessions() as u64
);
}
#[test]
fn test_max_session_size_default() {
let max = max_session_size();
assert_eq!(max, DEFAULT_MAX_SESSION_SIZE_MB * 1024 * 1024);
}
#[test]
fn test_validate_upload_uuid_valid() {
assert!(validate_upload_uuid("550e8400-e29b-41d4-a716-446655440000").is_ok());
assert!(validate_upload_uuid("abcdef01-2345-4678-9abc-def012345678").is_ok());
}
#[test]
fn test_validate_upload_uuid_rejects_path_traversal() {
assert!(validate_upload_uuid("../../../etc/passwd").is_err());
assert!(validate_upload_uuid("x/../../etc/cron.d/backdoor").is_err());
assert!(validate_upload_uuid("..").is_err());
}
#[test]
fn test_validate_upload_uuid_rejects_invalid() {
assert!(validate_upload_uuid("").is_err()); assert!(validate_upload_uuid("ABCDEF01-2345-4678-9ABC-DEF012345678").is_err()); assert!(validate_upload_uuid("hello world").is_err()); assert!(validate_upload_uuid("a".repeat(37).as_str()).is_err()); }
#[test]
fn test_cleanup_upload_temp_dir_removes_old_files() {
use std::fs::FileTimes;
use std::time::SystemTime;
let temp_dir = tempfile::TempDir::new().unwrap();
let upload_dir = temp_dir.path().join("tmp/docker-uploads");
std::fs::create_dir_all(&upload_dir).unwrap();
let old_file = upload_dir.join("old-uuid");
std::fs::write(&old_file, b"stale data").unwrap();
let old_time = SystemTime::now() - std::time::Duration::from_secs(7200);
let times = FileTimes::new().set_modified(old_time);
std::fs::File::options()
.write(true)
.open(&old_file)
.unwrap()
.set_times(times)
.unwrap();
let new_file = upload_dir.join("new-uuid");
std::fs::write(&new_file, b"fresh data").unwrap();
cleanup_upload_temp_dir(temp_dir.path().to_str().unwrap());
assert!(!old_file.exists(), "old file should be removed");
assert!(new_file.exists(), "recent file should be preserved");
}
#[test]
fn test_cleanup_upload_temp_dir_nonexistent_dir() {
cleanup_upload_temp_dir("/nonexistent/path/that/does/not/exist");
}
#[test]
fn test_detect_manifest_explicit_media_type() {
let manifest = serde_json::json!({
"mediaType": "application/vnd.docker.distribution.manifest.v2+json",
"schemaVersion": 2
});
let result = detect_manifest_media_type(manifest.to_string().as_bytes());
assert_eq!(
result,
"application/vnd.docker.distribution.manifest.v2+json"
);
}
#[test]
fn test_detect_manifest_oci_media_type() {
let manifest = serde_json::json!({
"mediaType": "application/vnd.oci.image.manifest.v1+json",
"schemaVersion": 2
});
let result = detect_manifest_media_type(manifest.to_string().as_bytes());
assert_eq!(result, "application/vnd.oci.image.manifest.v1+json");
}
#[test]
fn test_detect_manifest_schema_v1() {
let manifest = serde_json::json!({
"schemaVersion": 1,
"name": "test/image"
});
let result = detect_manifest_media_type(manifest.to_string().as_bytes());
assert_eq!(
result,
"application/vnd.docker.distribution.manifest.v1+json"
);
}
#[test]
fn test_detect_manifest_docker_v2_from_config() {
let manifest = serde_json::json!({
"schemaVersion": 2,
"config": {
"mediaType": "application/vnd.docker.container.image.v1+json",
"digest": "sha256:abc"
}
});
let result = detect_manifest_media_type(manifest.to_string().as_bytes());
assert_eq!(
result,
"application/vnd.docker.distribution.manifest.v2+json"
);
}
#[test]
fn test_detect_manifest_oci_from_config() {
let manifest = serde_json::json!({
"schemaVersion": 2,
"config": {
"mediaType": "application/vnd.oci.image.config.v1+json",
"digest": "sha256:abc"
}
});
let result = detect_manifest_media_type(manifest.to_string().as_bytes());
assert_eq!(result, "application/vnd.oci.image.manifest.v1+json");
}
#[test]
fn test_detect_manifest_no_config_media_type() {
let manifest = serde_json::json!({
"schemaVersion": 2,
"config": {
"digest": "sha256:abc"
}
});
let result = detect_manifest_media_type(manifest.to_string().as_bytes());
assert_eq!(
result,
"application/vnd.docker.distribution.manifest.v2+json"
);
}
#[test]
fn test_detect_manifest_index() {
let manifest = serde_json::json!({
"schemaVersion": 2,
"manifests": [
{"digest": "sha256:aaa", "platform": {"os": "linux", "architecture": "amd64"}}
]
});
let result = detect_manifest_media_type(manifest.to_string().as_bytes());
assert_eq!(result, "application/vnd.oci.image.index.v1+json");
}
#[test]
fn test_detect_manifest_invalid_json() {
let result = detect_manifest_media_type(b"not json at all");
assert_eq!(
result,
"application/vnd.docker.distribution.manifest.v2+json"
);
}
#[test]
fn test_detect_manifest_empty() {
let result = detect_manifest_media_type(b"{}");
assert_eq!(
result,
"application/vnd.docker.distribution.manifest.v2+json"
);
}
#[test]
fn test_detect_manifest_helm_chart() {
let manifest = serde_json::json!({
"schemaVersion": 2,
"config": {
"mediaType": "application/vnd.cncf.helm.config.v1+json",
"digest": "sha256:abc"
}
});
let result = detect_manifest_media_type(manifest.to_string().as_bytes());
assert_eq!(result, "application/vnd.oci.image.manifest.v1+json");
}
}
#[cfg(test)]
#[allow(clippy::unwrap_used)]
mod integration_tests {
use crate::circuit_breaker::ProbeToken;
use crate::test_helpers::{
body_bytes, create_test_context, create_test_context_with_config, send,
};
use axum::body::Body;
use axum::http::{header, Method, StatusCode};
use sha2::Digest;
#[tokio::test]
async fn test_docker_namespace_scope_enforced() {
use crate::auth::NamespaceAuthority;
use crate::config::ScopeEnforcement;
use axum::extract::{Path, State};
use axum::http::Uri;
use axum::Extension;
let ctx = create_test_context();
let scoped = NamespaceAuthority::from_oidc_scope(
"ci",
&["myorg/**".to_string()],
ScopeEnforcement::Enforce,
);
let resp = super::docker_v2_dispatch(
State(ctx.state.clone()),
Method::POST,
Path("other/app/blobs/uploads/".to_string()),
Extension(scoped.clone()),
"/v2/other/app/blobs/uploads/".parse::<Uri>().unwrap(),
axum::http::HeaderMap::new(),
Body::empty(),
)
.await;
assert_eq!(resp.status(), StatusCode::FORBIDDEN);
let resp = super::docker_v2_dispatch(
State(ctx.state.clone()),
Method::POST,
Path("myorg/app/blobs/uploads/".to_string()),
Extension(scoped.clone()),
"/v2/myorg/app/blobs/uploads/".parse::<Uri>().unwrap(),
axum::http::HeaderMap::new(),
Body::empty(),
)
.await;
assert_ne!(resp.status(), StatusCode::FORBIDDEN);
let resp = super::docker_v2_dispatch(
State(ctx.state.clone()),
Method::GET,
Path("other/app/manifests/latest".to_string()),
Extension(scoped.clone()),
"/v2/other/app/manifests/latest".parse::<Uri>().unwrap(),
axum::http::HeaderMap::new(),
Body::empty(),
)
.await;
assert_ne!(resp.status(), StatusCode::FORBIDDEN);
}
#[tokio::test]
async fn test_docker_v2_check() {
let ctx = create_test_context();
let resp = send(&ctx.app, Method::GET, "/v2/", Body::empty()).await;
assert_eq!(resp.status(), StatusCode::OK);
}
#[tokio::test]
async fn test_docker_catalog_empty() {
let ctx = create_test_context();
let resp = send(&ctx.app, Method::GET, "/v2/_catalog", Body::empty()).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["repositories"].as_array().unwrap().is_empty());
}
async fn seed_zero_config(state: &crate::AppState, name: &str) {
let _ = state
.storage
.put(
&format!(
"docker/{}/blobs/sha256:0000000000000000000000000000000000000000000000000000000000000000",
name
),
b"x",
)
.await;
}
#[tokio::test]
async fn test_docker_put_get_manifest() {
let ctx = create_test_context();
let manifest = serde_json::json!({
"schemaVersion": 2,
"mediaType": "application/vnd.docker.distribution.manifest.v2+json",
"config": {
"mediaType": "application/vnd.docker.container.image.v1+json",
"size": 0,
"digest": "sha256:0000000000000000000000000000000000000000000000000000000000000000"
},
"layers": []
});
let manifest_bytes = serde_json::to_vec(&manifest).unwrap();
seed_zero_config(&ctx.state, "alpine").await;
let put_resp = send(
&ctx.app,
Method::PUT,
"/v2/alpine/manifests/latest",
Body::from(manifest_bytes.clone()),
)
.await;
assert_eq!(put_resp.status(), StatusCode::CREATED);
let digest_header = put_resp
.headers()
.get("docker-content-digest")
.unwrap()
.to_str()
.unwrap()
.to_string();
assert!(digest_header.starts_with("sha256:"));
let get_resp = send(
&ctx.app,
Method::GET,
"/v2/alpine/manifests/latest",
Body::empty(),
)
.await;
assert_eq!(get_resp.status(), StatusCode::OK);
let get_digest = get_resp
.headers()
.get("docker-content-digest")
.unwrap()
.to_str()
.unwrap()
.to_string();
assert_eq!(get_digest, digest_header);
let body = body_bytes(get_resp).await;
assert_eq!(body.as_ref(), manifest_bytes.as_slice());
}
#[tokio::test]
async fn test_hosted_manifest_skips_upstream_revalidation() {
use crate::config::DockerUpstream;
use crate::test_helpers::create_test_context_with_config;
let ctx = create_test_context_with_config(|cfg| {
cfg.docker.serve_stale = false;
cfg.docker.upstreams = vec![DockerUpstream {
url: "http://127.0.0.1:1".into(),
auth: None,
namespace: None,
prefix: None,
}];
});
let manifest = serde_json::json!({
"schemaVersion": 2,
"mediaType": "application/vnd.docker.distribution.manifest.v2+json",
"config": {
"mediaType": "application/vnd.docker.container.image.v1+json",
"size": 0,
"digest": "sha256:0000000000000000000000000000000000000000000000000000000000000000"
},
"layers": []
});
let manifest_bytes = serde_json::to_vec(&manifest).unwrap();
seed_zero_config(&ctx.state, "hosted-img").await;
let put_resp = send(
&ctx.app,
Method::PUT,
"/v2/hosted-img/manifests/v1",
Body::from(manifest_bytes.clone()),
)
.await;
assert_eq!(put_resp.status(), StatusCode::CREATED);
let resp = send(
&ctx.app,
Method::GET,
"/v2/hosted-img/manifests/v1",
Body::empty(),
)
.await;
assert_eq!(resp.status(), StatusCode::OK);
assert!(
resp.headers().get("x-nora-stale").is_none(),
"hosted manifest must not be served via the stale-while-error path"
);
let body = body_bytes(resp).await;
assert_eq!(body.as_ref(), manifest_bytes.as_slice());
}
#[tokio::test]
async fn test_local_push_writes_bare_manifest_key() {
let ctx = create_test_context();
seed_zero_config(&ctx.state, "prov-img").await;
let manifest = serde_json::json!({
"schemaVersion": 2,
"mediaType": "application/vnd.docker.distribution.manifest.v2+json",
"config": {
"mediaType": "application/vnd.docker.container.image.v1+json",
"size": 0,
"digest": "sha256:0000000000000000000000000000000000000000000000000000000000000000"
},
"layers": []
});
let manifest_bytes = serde_json::to_vec(&manifest).unwrap();
let put_resp = send(
&ctx.app,
Method::PUT,
"/v2/prov-img/manifests/v1",
Body::from(manifest_bytes),
)
.await;
assert_eq!(put_resp.status(), StatusCode::CREATED);
assert!(
ctx.state
.storage
.get("docker/prov-img/manifests/v1.json")
.await
.is_ok(),
"local push must write the bare manifest key (hosted provenance)"
);
}
#[tokio::test]
async fn test_manifest_push_rejects_absent_blob() {
let ctx = create_test_context();
let manifest = serde_json::json!({
"schemaVersion": 2,
"mediaType": "application/vnd.docker.distribution.manifest.v2+json",
"config": {
"mediaType": "application/vnd.docker.container.image.v1+json",
"size": 0,
"digest": "sha256:1111111111111111111111111111111111111111111111111111111111111111"
},
"layers": []
});
let resp = send(
&ctx.app,
Method::PUT,
"/v2/reject/manifests/latest",
Body::from(serde_json::to_vec(&manifest).unwrap()),
)
.await;
assert_eq!(resp.status(), StatusCode::BAD_REQUEST);
let body = body_bytes(resp).await;
let json: serde_json::Value = serde_json::from_slice(&body).unwrap();
assert_eq!(json["errors"][0]["code"], "MANIFEST_BLOB_UNKNOWN");
}
#[tokio::test]
async fn test_docker_list_tags() {
let ctx = create_test_context();
let manifest = serde_json::json!({
"schemaVersion": 2,
"mediaType": "application/vnd.docker.distribution.manifest.v2+json",
"config": {
"mediaType": "application/vnd.docker.container.image.v1+json",
"size": 0,
"digest": "sha256:0000000000000000000000000000000000000000000000000000000000000000"
},
"layers": []
});
seed_zero_config(&ctx.state, "alpine").await;
send(
&ctx.app,
Method::PUT,
"/v2/alpine/manifests/latest",
Body::from(serde_json::to_vec(&manifest).unwrap()),
)
.await;
let list_resp = send(&ctx.app, Method::GET, "/v2/alpine/tags/list", Body::empty()).await;
assert_eq!(list_resp.status(), StatusCode::OK);
let body = body_bytes(list_resp).await;
let json: serde_json::Value = serde_json::from_slice(&body).unwrap();
assert_eq!(json["name"], "alpine");
let tags = json["tags"].as_array().unwrap();
assert!(tags.contains(&serde_json::json!("latest")));
}
#[tokio::test]
async fn test_docker_delete_manifest() {
let ctx = create_test_context();
let manifest = serde_json::json!({
"schemaVersion": 2,
"mediaType": "application/vnd.docker.distribution.manifest.v2+json",
"config": {
"mediaType": "application/vnd.docker.container.image.v1+json",
"size": 0,
"digest": "sha256:0000000000000000000000000000000000000000000000000000000000000000"
},
"layers": []
});
seed_zero_config(&ctx.state, "alpine").await;
let put_resp = send(
&ctx.app,
Method::PUT,
"/v2/alpine/manifests/latest",
Body::from(serde_json::to_vec(&manifest).unwrap()),
)
.await;
let digest = put_resp
.headers()
.get("docker-content-digest")
.unwrap()
.to_str()
.unwrap()
.to_string();
let del = send(
&ctx.app,
Method::DELETE,
&format!("/v2/alpine/manifests/{}", digest),
Body::empty(),
)
.await;
assert_eq!(del.status(), StatusCode::ACCEPTED);
}
#[tokio::test]
async fn test_docker_delete_by_digest_removes_tag() {
let ctx = create_test_context();
let manifest = serde_json::json!({
"schemaVersion": 2,
"mediaType": "application/vnd.docker.distribution.manifest.v2+json",
"config": {
"mediaType": "application/vnd.docker.container.image.v1+json",
"size": 0,
"digest": "sha256:0000000000000000000000000000000000000000000000000000000000000000"
},
"layers": []
});
seed_zero_config(&ctx.state, "alpine").await;
let put_resp = send(
&ctx.app,
Method::PUT,
"/v2/alpine/manifests/v1",
Body::from(serde_json::to_vec(&manifest).unwrap()),
)
.await;
let digest = put_resp
.headers()
.get("docker-content-digest")
.unwrap()
.to_str()
.unwrap()
.to_string();
let del = send(
&ctx.app,
Method::DELETE,
&format!("/v2/alpine/manifests/{}", digest),
Body::empty(),
)
.await;
assert_eq!(del.status(), StatusCode::ACCEPTED);
let tag_get = send(
&ctx.app,
Method::GET,
"/v2/alpine/manifests/v1",
Body::empty(),
)
.await;
assert_eq!(
tag_get.status(),
StatusCode::NOT_FOUND,
"tag must 404 after its manifest is deleted by digest (#658)"
);
let list = send(&ctx.app, Method::GET, "/v2/alpine/tags/list", Body::empty()).await;
let body = body_bytes(list).await;
let json: serde_json::Value = serde_json::from_slice(&body).unwrap();
let tags = json["tags"].as_array().unwrap();
assert!(
!tags.contains(&serde_json::json!("v1")),
"v1 must be gone from tags/list after digest delete (#658)"
);
}
#[tokio::test]
async fn test_docker_monolithic_upload() {
let ctx = create_test_context();
let blob_data = b"test blob data";
let digest = format!("sha256:{}", hex::encode(sha2::Sha256::digest(blob_data)));
let post_resp = send(
&ctx.app,
Method::POST,
"/v2/alpine/blobs/uploads/",
Body::empty(),
)
.await;
assert_eq!(post_resp.status(), StatusCode::ACCEPTED);
let location = post_resp
.headers()
.get("location")
.unwrap()
.to_str()
.unwrap()
.to_string();
let uuid = location.rsplit('/').next().unwrap();
let put_url = format!("/v2/alpine/blobs/uploads/{}?digest={}", uuid, digest);
let put_resp = send(&ctx.app, Method::PUT, &put_url, Body::from(&blob_data[..])).await;
assert_eq!(put_resp.status(), StatusCode::CREATED);
}
#[tokio::test]
async fn test_docker_single_post_monolithic_upload() {
let ctx = create_test_context();
let blob_data = b"single-post monolithic blob";
let digest = format!("sha256:{}", hex::encode(sha2::Sha256::digest(blob_data)));
let resp = send(
&ctx.app,
Method::POST,
&format!("/v2/alpine/blobs/uploads/?digest={}", digest),
Body::from(&blob_data[..]),
)
.await;
assert_eq!(
resp.status(),
StatusCode::CREATED,
"single-POST monolithic upload must return 201"
);
assert_eq!(
resp.headers()
.get("docker-content-digest")
.unwrap()
.to_str()
.unwrap(),
digest
);
let head = send(
&ctx.app,
Method::HEAD,
&format!("/v2/alpine/blobs/{}", digest),
Body::empty(),
)
.await;
assert_eq!(
head.status(),
StatusCode::OK,
"blob must exist after a single-POST upload"
);
}
#[tokio::test]
async fn test_docker_cross_repo_blob_mount() {
let ctx = create_test_context();
let blob_data = b"cross-repo mounted layer";
let digest = format!("sha256:{}", hex::encode(sha2::Sha256::digest(blob_data)));
let resp = send(
&ctx.app,
Method::POST,
&format!("/v2/srcrepo/blobs/uploads/?digest={}", digest),
Body::from(&blob_data[..]),
)
.await;
assert_eq!(resp.status(), StatusCode::CREATED);
let mount_url = format!("/v2/dstrepo/blobs/uploads/?mount={}&from=srcrepo", digest);
let resp = send(&ctx.app, Method::POST, &mount_url, Body::empty()).await;
assert_eq!(
resp.status(),
StatusCode::CREATED,
"a blob the source repo holds must mount without an upload"
);
assert_eq!(
resp.headers()
.get(header::LOCATION)
.unwrap()
.to_str()
.unwrap(),
format!("/v2/dstrepo/blobs/{}", digest)
);
assert_eq!(
resp.headers()
.get("docker-content-digest")
.unwrap()
.to_str()
.unwrap(),
digest
);
let get = send(
&ctx.app,
Method::GET,
&format!("/v2/dstrepo/blobs/{}", digest),
Body::empty(),
)
.await;
assert_eq!(get.status(), StatusCode::OK);
assert_eq!(&body_bytes(get).await[..], &blob_data[..]);
let resp = send(
&ctx.app,
Method::POST,
&format!("/v2/dstrepo/blobs/uploads/?mount={}&from=otherrepo", digest),
Body::empty(),
)
.await;
assert_eq!(resp.status(), StatusCode::CREATED);
}
#[tokio::test]
async fn test_docker_mount_unknown_blob_falls_back_to_upload() {
let ctx = create_test_context();
let digest = format!("sha256:{}", "b".repeat(64));
let resp = send(
&ctx.app,
Method::POST,
&format!("/v2/dstrepo/blobs/uploads/?mount={}&from=srcrepo", digest),
Body::empty(),
)
.await;
assert_eq!(
resp.status(),
StatusCode::ACCEPTED,
"an unmountable blob must degrade to a normal upload session"
);
assert!(resp.headers().contains_key("docker-upload-uuid"));
}
#[tokio::test]
async fn test_docker_mount_invalid_source_repo_falls_back_to_upload() {
let ctx = create_test_context();
let blob_data = b"blob behind an unusable source name";
let digest = format!("sha256:{}", hex::encode(sha2::Sha256::digest(blob_data)));
let resp = send(
&ctx.app,
Method::POST,
&format!("/v2/srcrepo/blobs/uploads/?digest={}", digest),
Body::from(&blob_data[..]),
)
.await;
assert_eq!(resp.status(), StatusCode::CREATED);
let resp = send(
&ctx.app,
Method::POST,
&format!("/v2/dstrepo/blobs/uploads/?mount={}&from=SRCREPO", digest),
Body::empty(),
)
.await;
assert_eq!(
resp.status(),
StatusCode::ACCEPTED,
"an invalid source repository must degrade to an upload session, not 4xx"
);
assert!(resp.headers().contains_key("docker-upload-uuid"));
}
#[tokio::test]
async fn test_docker_single_post_digest_mismatch_rejected() {
let ctx = create_test_context();
let wrong = format!("sha256:{}", "0".repeat(64));
let resp = send(
&ctx.app,
Method::POST,
&format!("/v2/alpine/blobs/uploads/?digest={}", wrong),
Body::from(&b"some other bytes"[..]),
)
.await;
assert_eq!(resp.status(), StatusCode::BAD_REQUEST);
}
#[tokio::test]
async fn test_docker_chunked_upload() {
let ctx = create_test_context();
let blob_data = b"test chunked blob";
let digest = format!("sha256:{}", hex::encode(sha2::Sha256::digest(blob_data)));
let post_resp = send(
&ctx.app,
Method::POST,
"/v2/alpine/blobs/uploads/",
Body::empty(),
)
.await;
assert_eq!(post_resp.status(), StatusCode::ACCEPTED);
let location = post_resp
.headers()
.get("location")
.unwrap()
.to_str()
.unwrap()
.to_string();
let uuid = location.rsplit('/').next().unwrap();
let patch_url = format!("/v2/alpine/blobs/uploads/{}", uuid);
let patch_resp = send(
&ctx.app,
Method::PATCH,
&patch_url,
Body::from(&blob_data[..]),
)
.await;
assert_eq!(patch_resp.status(), StatusCode::ACCEPTED);
let put_url = format!("/v2/alpine/blobs/uploads/{}?digest={}", uuid, digest);
let put_resp = send(&ctx.app, Method::PUT, &put_url, Body::empty()).await;
assert_eq!(put_resp.status(), StatusCode::CREATED);
}
#[tokio::test]
async fn test_monolithic_upload_over_cap_rejected_fast() {
let ctx = create_test_context_with_config(|c| c.server.body_limit_mb = 1); let blob = vec![0u8; 2 * 1024 * 1024]; let digest = format!("sha256:{}", hex::encode(sha2::Sha256::digest(&blob)));
let resp = send(
&ctx.app,
Method::POST,
&format!("/v2/alpine/blobs/uploads/?digest={}", digest),
Body::from(blob),
)
.await;
assert_eq!(resp.status(), StatusCode::PAYLOAD_TOO_LARGE);
}
#[tokio::test]
async fn test_monolithic_upload_over_cap_streamed_rejected() {
let ctx = create_test_context_with_config(|c| c.server.body_limit_mb = 1); let chunks: Vec<Result<axum::body::Bytes, std::io::Error>> = (0..4)
.map(|_| Ok(axum::body::Bytes::from(vec![0u8; 512 * 1024])))
.collect(); let body = Body::from_stream(futures::stream::iter(chunks));
let digest = format!("sha256:{}", "0".repeat(64)); let resp = send(
&ctx.app,
Method::POST,
&format!("/v2/alpine/blobs/uploads/?digest={}", digest),
body,
)
.await;
assert_eq!(
resp.status(),
StatusCode::PAYLOAD_TOO_LARGE,
"streamed body over cap must 413 via the incremental guard, not OOM"
);
}
#[tokio::test]
async fn test_patch_over_cap_rejected() {
let ctx = create_test_context_with_config(|c| c.server.body_limit_mb = 1); let post = send(
&ctx.app,
Method::POST,
"/v2/alpine/blobs/uploads/",
Body::empty(),
)
.await;
assert_eq!(post.status(), StatusCode::ACCEPTED);
let loc = post
.headers()
.get("location")
.unwrap()
.to_str()
.unwrap()
.to_string();
let uuid = loc.rsplit('/').next().unwrap();
let big = vec![0u8; 2 * 1024 * 1024]; let resp = send(
&ctx.app,
Method::PATCH,
&format!("/v2/alpine/blobs/uploads/{}", uuid),
Body::from(big),
)
.await;
assert_eq!(resp.status(), StatusCode::PAYLOAD_TOO_LARGE);
}
#[tokio::test]
async fn test_manifest_over_cap_rejected() {
let ctx = create_test_context();
let huge = vec![b'{'; super::MAX_MANIFEST_BYTES + 1]; let resp = send(
&ctx.app,
Method::PUT,
"/v2/alpine/manifests/latest",
Body::from(huge),
)
.await;
assert_eq!(resp.status(), StatusCode::PAYLOAD_TOO_LARGE);
}
#[tokio::test]
async fn test_streaming_upload_multiframe_ok() {
let ctx = create_test_context();
let blob = vec![7u8; 1024 * 1024]; let digest = format!("sha256:{}", hex::encode(sha2::Sha256::digest(&blob)));
let chunks: Vec<Result<axum::body::Bytes, std::io::Error>> = blob
.chunks(64 * 1024)
.map(|c| Ok(axum::body::Bytes::copy_from_slice(c)))
.collect();
let body = Body::from_stream(futures::stream::iter(chunks));
let resp = send(
&ctx.app,
Method::POST,
&format!("/v2/alpine/blobs/uploads/?digest={}", digest),
body,
)
.await;
assert_eq!(
resp.status(),
StatusCode::CREATED,
"multi-frame streamed blob must store and verify"
);
}
#[tokio::test]
async fn test_docker_check_blob() {
let ctx = create_test_context();
let blob_data = b"test blob for head";
let digest = format!("sha256:{}", hex::encode(sha2::Sha256::digest(blob_data)));
let post_resp = send(
&ctx.app,
Method::POST,
"/v2/alpine/blobs/uploads/",
Body::empty(),
)
.await;
let location = post_resp
.headers()
.get("location")
.unwrap()
.to_str()
.unwrap()
.to_string();
let uuid = location.rsplit('/').next().unwrap();
let put_url = format!("/v2/alpine/blobs/uploads/{}?digest={}", uuid, digest);
send(&ctx.app, Method::PUT, &put_url, Body::from(&blob_data[..])).await;
let head_url = format!("/v2/alpine/blobs/{}", digest);
let head_resp = send(&ctx.app, Method::HEAD, &head_url, Body::empty()).await;
assert_eq!(head_resp.status(), StatusCode::OK);
let cl = head_resp
.headers()
.get(header::CONTENT_LENGTH)
.unwrap()
.to_str()
.unwrap()
.parse::<usize>()
.unwrap();
assert_eq!(cl, blob_data.len());
}
#[tokio::test]
async fn test_docker_download_blob() {
let ctx = create_test_context();
let blob_data = b"test blob for download";
let digest = format!("sha256:{}", hex::encode(sha2::Sha256::digest(blob_data)));
let post_resp = send(
&ctx.app,
Method::POST,
"/v2/alpine/blobs/uploads/",
Body::empty(),
)
.await;
let location = post_resp
.headers()
.get("location")
.unwrap()
.to_str()
.unwrap()
.to_string();
let uuid = location.rsplit('/').next().unwrap();
let put_url = format!("/v2/alpine/blobs/uploads/{}?digest={}", uuid, digest);
send(&ctx.app, Method::PUT, &put_url, Body::from(&blob_data[..])).await;
let get_url = format!("/v2/alpine/blobs/{}", digest);
let get_resp = send(&ctx.app, Method::GET, &get_url, Body::empty()).await;
assert_eq!(get_resp.status(), StatusCode::OK);
let body = body_bytes(get_resp).await;
assert_eq!(body.as_ref(), &blob_data[..]);
}
#[tokio::test]
async fn test_docker_blob_range_request() {
use tower::ServiceExt;
let ctx = create_test_context();
let blob = b"0123456789abcdef";
let digest = format!("sha256:{}", hex::encode(sha2::Sha256::digest(blob)));
let post = send(
&ctx.app,
Method::POST,
"/v2/rng/blobs/uploads/",
Body::empty(),
)
.await;
let loc = post
.headers()
.get("location")
.unwrap()
.to_str()
.unwrap()
.to_string();
let uuid = loc.rsplit('/').next().unwrap();
let put_url = format!("/v2/rng/blobs/uploads/{}?digest={}", uuid, digest);
send(&ctx.app, Method::PUT, &put_url, Body::from(&blob[..])).await;
let req = axum::http::Request::builder()
.method(Method::GET)
.uri(format!("/v2/rng/blobs/{}", digest))
.header(header::RANGE, "bytes=4-7")
.body(Body::empty())
.unwrap();
let resp = ctx.app.clone().oneshot(req).await.unwrap();
assert_eq!(resp.status(), StatusCode::PARTIAL_CONTENT);
assert_eq!(
resp.headers()
.get(header::CONTENT_RANGE)
.unwrap()
.to_str()
.unwrap(),
format!("bytes 4-7/{}", blob.len())
);
let body = body_bytes(resp).await;
assert_eq!(body.as_ref(), &blob[4..=7]);
let req = axum::http::Request::builder()
.method(Method::GET)
.uri(format!("/v2/rng/blobs/{}", digest))
.header(header::RANGE, format!("bytes={}-", blob.len()))
.body(Body::empty())
.unwrap();
let resp = ctx.app.clone().oneshot(req).await.unwrap();
assert_eq!(resp.status(), StatusCode::RANGE_NOT_SATISFIABLE);
assert_eq!(
resp.headers()
.get(header::CONTENT_RANGE)
.unwrap()
.to_str()
.unwrap(),
format!("bytes */{}", blob.len())
);
let req = axum::http::Request::builder()
.method(Method::GET)
.uri(format!("/v2/rng/blobs/{}", digest))
.header(header::RANGE, "kilobytes=1-2")
.body(Body::empty())
.unwrap();
let resp = ctx.app.clone().oneshot(req).await.unwrap();
assert_eq!(resp.status(), StatusCode::OK);
}
#[tokio::test]
async fn test_docker_blob_not_found() {
let ctx = create_test_context();
let fake_digest = "sha256:0000000000000000000000000000000000000000000000000000000000000000";
let head_url = format!("/v2/alpine/blobs/{}", fake_digest);
let resp = send(&ctx.app, Method::HEAD, &head_url, Body::empty()).await;
assert_eq!(resp.status(), StatusCode::NOT_FOUND);
}
#[tokio::test]
async fn test_docker_delete_blob() {
let ctx = create_test_context();
let blob_data = b"test blob for delete";
let digest = format!("sha256:{}", hex::encode(sha2::Sha256::digest(blob_data)));
let post_resp = send(
&ctx.app,
Method::POST,
"/v2/alpine/blobs/uploads/",
Body::empty(),
)
.await;
let location = post_resp
.headers()
.get("location")
.unwrap()
.to_str()
.unwrap()
.to_string();
let uuid = location.rsplit('/').next().unwrap();
let put_url = format!("/v2/alpine/blobs/uploads/{}?digest={}", uuid, digest);
send(&ctx.app, Method::PUT, &put_url, Body::from(&blob_data[..])).await;
let delete_url = format!("/v2/alpine/blobs/{}", digest);
let delete_resp = send(&ctx.app, Method::DELETE, &delete_url, Body::empty()).await;
assert_eq!(delete_resp.status(), StatusCode::ACCEPTED);
}
#[tokio::test]
async fn test_docker_namespaced_routes() {
let ctx = create_test_context();
let manifest = serde_json::json!({
"schemaVersion": 2,
"mediaType": "application/vnd.docker.distribution.manifest.v2+json",
"config": {
"mediaType": "application/vnd.docker.container.image.v1+json",
"size": 0,
"digest": "sha256:0000000000000000000000000000000000000000000000000000000000000000"
},
"layers": []
});
seed_zero_config(&ctx.state, "library/alpine").await;
let put_resp = send(
&ctx.app,
Method::PUT,
"/v2/library/alpine/manifests/latest",
Body::from(serde_json::to_vec(&manifest).unwrap()),
)
.await;
assert_eq!(put_resp.status(), StatusCode::CREATED);
assert!(put_resp
.headers()
.get("docker-content-digest")
.unwrap()
.to_str()
.unwrap()
.starts_with("sha256:"));
}
#[tokio::test]
async fn test_extract_docker_publish_date_from_meta() {
let dir = tempfile::tempdir().unwrap();
let storage = crate::storage::Storage::new_local(dir.path().join("data").to_str().unwrap());
let meta = super::ImageMetadata {
push_timestamp: 1700000000,
..Default::default()
};
storage
.put(
"docker/library/nginx/manifests/latest.meta.json",
serde_json::to_vec(&meta).unwrap().as_slice(),
)
.await
.unwrap();
let result = super::extract_docker_publish_date(
&storage,
"library/nginx",
"latest",
true, None, )
.await;
assert_eq!(result, Some(1700000000));
}
#[tokio::test]
async fn test_extract_docker_publish_date_mtime_fallback() {
let dir = tempfile::tempdir().unwrap();
let storage = crate::storage::Storage::new_local(dir.path().join("data").to_str().unwrap());
storage
.put("docker/library/nginx/manifests/latest.json", b"{}")
.await
.unwrap();
let result = super::extract_docker_publish_date(
&storage,
"library/nginx",
"latest",
true, None, )
.await;
assert!(result.is_some());
assert!(result.unwrap() > 0);
}
#[tokio::test]
async fn test_extract_docker_publish_date_proxy_no_fallback() {
let dir = tempfile::tempdir().unwrap();
let storage = crate::storage::Storage::new_local(dir.path().join("data").to_str().unwrap());
storage
.put("docker/library/nginx/manifests/latest.json", b"{}")
.await
.unwrap();
let result = super::extract_docker_publish_date(
&storage,
"library/nginx",
"latest",
false, None, )
.await;
assert!(result.is_none());
}
#[tokio::test]
async fn test_docker_circuit_breaker_trips() {
use crate::config::DockerUpstream;
use crate::test_helpers::{body_bytes, create_test_context_with_config, send};
let ctx = create_test_context_with_config(|cfg| {
cfg.circuit_breaker.enabled = true;
cfg.circuit_breaker.failure_threshold = 2;
cfg.circuit_breaker.reset_timeout = 3600;
cfg.docker.upstreams = vec![DockerUpstream {
url: "http://127.0.0.1:1".into(),
auth: None,
namespace: None,
prefix: None,
}];
});
ctx.state
.circuit_breaker
.record_failure("docker:http://127.0.0.1:1", ProbeToken::BACKGROUND);
ctx.state
.circuit_breaker
.record_failure("docker:http://127.0.0.1:1", ProbeToken::BACKGROUND);
let response = send(
&ctx.app,
Method::GET,
"/v2/library/nonexistent/manifests/latest",
Body::empty(),
)
.await;
assert_eq!(response.status(), StatusCode::SERVICE_UNAVAILABLE);
assert_eq!(
response
.headers()
.get("retry-after")
.and_then(|v| v.to_str().ok()),
Some("30")
);
let body = body_bytes(response).await;
assert!(String::from_utf8_lossy(&body).contains("temporarily unavailable"));
}
#[tokio::test]
async fn test_oci_v2_api_version_header() {
let ctx = create_test_context();
let resp = send(&ctx.app, Method::GET, "/v2/", Body::empty()).await;
assert_eq!(resp.status(), StatusCode::OK);
let api_ver = resp
.headers()
.get("docker-distribution-api-version")
.expect("OCI spec requires Docker-Distribution-API-Version header")
.to_str()
.unwrap();
assert_eq!(api_ver, "registry/2.0");
}
#[tokio::test]
async fn test_oci_catalog_json_structure() {
let ctx = create_test_context();
let resp = send(&ctx.app, Method::GET, "/v2/_catalog", Body::empty()).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("repositories").is_some(),
"OCI spec requires 'repositories' key in catalog response"
);
assert!(json["repositories"].is_array());
}
#[tokio::test]
async fn test_oci_tags_list_json_structure() {
let ctx = create_test_context();
let manifest = serde_json::json!({
"schemaVersion": 2,
"mediaType": "application/vnd.docker.distribution.manifest.v2+json",
"config": {
"mediaType": "application/vnd.docker.container.image.v1+json",
"size": 0,
"digest": "sha256:0000000000000000000000000000000000000000000000000000000000000000"
},
"layers": []
});
send(
&ctx.app,
Method::PUT,
"/v2/myapp/manifests/v1",
Body::from(serde_json::to_vec(&manifest).unwrap()),
)
.await;
let resp = send(&ctx.app, Method::GET, "/v2/myapp/tags/list", Body::empty()).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("name").is_some(),
"OCI spec requires 'name' in tags/list"
);
assert!(
json.get("tags").is_some(),
"OCI spec requires 'tags' in tags/list"
);
assert_eq!(json["name"], "myapp");
assert!(json["tags"].is_array());
}
#[tokio::test]
async fn test_oci_manifest_digest_matches_body() {
let ctx = create_test_context();
let manifest = serde_json::json!({
"schemaVersion": 2,
"mediaType": "application/vnd.docker.distribution.manifest.v2+json",
"config": {
"mediaType": "application/vnd.docker.container.image.v1+json",
"size": 0,
"digest": "sha256:0000000000000000000000000000000000000000000000000000000000000000"
},
"layers": []
});
let manifest_bytes = serde_json::to_vec(&manifest).unwrap();
seed_zero_config(&ctx.state, "verify").await;
send(
&ctx.app,
Method::PUT,
"/v2/verify/manifests/latest",
Body::from(manifest_bytes),
)
.await;
let resp = send(
&ctx.app,
Method::GET,
"/v2/verify/manifests/latest",
Body::empty(),
)
.await;
assert_eq!(resp.status(), StatusCode::OK);
let digest_header = resp
.headers()
.get("docker-content-digest")
.expect("OCI spec requires Docker-Content-Digest header")
.to_str()
.unwrap()
.to_string();
let body = body_bytes(resp).await;
let computed = format!("sha256:{}", hex::encode(sha2::Sha256::digest(&body)));
assert_eq!(
digest_header, computed,
"Docker-Content-Digest must equal sha256 of response body"
);
}
#[tokio::test]
async fn test_oci_manifest_content_type_matches_media_type() {
let ctx = create_test_context();
let manifest = serde_json::json!({
"schemaVersion": 2,
"mediaType": "application/vnd.docker.distribution.manifest.v2+json",
"config": {
"mediaType": "application/vnd.docker.container.image.v1+json",
"size": 0,
"digest": "sha256:0000000000000000000000000000000000000000000000000000000000000000"
},
"layers": []
});
seed_zero_config(&ctx.state, "ctcheck").await;
send(
&ctx.app,
Method::PUT,
"/v2/ctcheck/manifests/v1",
Body::from(serde_json::to_vec(&manifest).unwrap()),
)
.await;
let resp = send(
&ctx.app,
Method::GET,
"/v2/ctcheck/manifests/v1",
Body::empty(),
)
.await;
assert_eq!(resp.status(), StatusCode::OK);
let ct = resp
.headers()
.get(header::CONTENT_TYPE)
.expect("OCI spec requires Content-Type header on manifest")
.to_str()
.unwrap();
assert_eq!(ct, "application/vnd.docker.distribution.manifest.v2+json");
}
#[tokio::test]
async fn test_oci_manifest_content_length_matches() {
let ctx = create_test_context();
let manifest = serde_json::json!({
"schemaVersion": 2,
"mediaType": "application/vnd.docker.distribution.manifest.v2+json",
"config": {
"mediaType": "application/vnd.docker.container.image.v1+json",
"size": 0,
"digest": "sha256:0000000000000000000000000000000000000000000000000000000000000000"
},
"layers": []
});
seed_zero_config(&ctx.state, "clcheck").await;
send(
&ctx.app,
Method::PUT,
"/v2/clcheck/manifests/v1",
Body::from(serde_json::to_vec(&manifest).unwrap()),
)
.await;
let resp = send(
&ctx.app,
Method::GET,
"/v2/clcheck/manifests/v1",
Body::empty(),
)
.await;
assert_eq!(resp.status(), StatusCode::OK);
let cl: usize = resp
.headers()
.get(header::CONTENT_LENGTH)
.expect("OCI spec requires Content-Length on manifest")
.to_str()
.unwrap()
.parse()
.unwrap();
let body = body_bytes(resp).await;
assert_eq!(cl, body.len(), "Content-Length must match body size");
}
#[tokio::test]
async fn test_docker_circuit_breaker_per_upstream() {
use crate::config::DockerUpstream;
use crate::test_helpers::create_test_context_with_config;
let ctx = create_test_context_with_config(|cfg| {
cfg.circuit_breaker.enabled = true;
cfg.circuit_breaker.failure_threshold = 2;
cfg.circuit_breaker.reset_timeout = 3600;
cfg.docker.upstreams = vec![
DockerUpstream {
url: "http://127.0.0.1:1".into(), auth: None,
namespace: None,
prefix: None,
},
DockerUpstream {
url: "http://127.0.0.1:2".into(), auth: None,
namespace: None,
prefix: None,
},
];
});
ctx.state
.circuit_breaker
.record_failure("docker:http://127.0.0.1:1", ProbeToken::BACKGROUND);
ctx.state
.circuit_breaker
.record_failure("docker:http://127.0.0.1:1", ProbeToken::BACKGROUND);
assert!(ctx
.state
.circuit_breaker
.check("docker:http://127.0.0.1:1")
.is_err());
assert!(ctx
.state
.circuit_breaker
.check("docker:http://127.0.0.1:2")
.is_ok());
}
#[test]
fn test_strip_docker_namespace() {
assert_eq!(
super::strip_docker_namespace("docker.io/library/nginx"),
"library/nginx"
);
assert_eq!(
super::strip_docker_namespace("ghcr.io/requarks/wiki"),
"requarks/wiki"
);
assert_eq!(
super::strip_docker_namespace("registry.example.com/myapp"),
"myapp"
);
assert_eq!(
super::strip_docker_namespace("library/nginx"),
"library/nginx"
);
assert_eq!(super::strip_docker_namespace("alpine"), "alpine");
assert_eq!(super::strip_docker_namespace(""), "");
assert_eq!(super::strip_docker_namespace("docker.io/"), "docker.io/");
assert_eq!(
super::strip_docker_namespace("myorg/myimage"),
"myorg/myimage"
);
}
#[tokio::test]
async fn test_catalog_dedup_across_namespaces() {
use crate::test_helpers::create_test_context;
let ctx = create_test_context();
ctx.state
.storage
.put(
"docker/docker.io/library/nginx/manifests/latest.json",
b"{}",
)
.await
.unwrap();
ctx.state
.storage
.put("docker/ghcr.io/library/nginx/manifests/v1.json", b"{}")
.await
.unwrap();
ctx.state
.storage
.put("docker/library/nginx/manifests/old.json", b"{}")
.await
.unwrap();
let keys = ctx.state.storage.list("docker/").await.unwrap();
let mut repos: Vec<String> = keys
.iter()
.filter_map(|k| {
let rest = k.strip_prefix("docker/")?;
let name = if let Some(idx) = rest.find("/manifests/") {
&rest[..idx]
} else {
return None;
};
if name.is_empty() {
return None;
}
Some(super::strip_docker_namespace(name).to_string())
})
.collect();
repos.sort();
repos.dedup();
assert_eq!(repos, vec!["library/nginx"]);
}
fn docker_config_allow(
upstreams: Vec<crate::config::DockerUpstream>,
) -> crate::config::DockerConfig {
crate::config::DockerConfig {
upstreams,
default_action: crate::config::DefaultAction::Allow,
..Default::default()
}
}
fn docker_config_deny(
upstreams: Vec<crate::config::DockerUpstream>,
) -> crate::config::DockerConfig {
crate::config::DockerConfig {
upstreams,
default_action: crate::config::DefaultAction::Deny,
..Default::default()
}
}
#[test]
fn test_canonicalize_prefix_routing() {
let upstreams = vec![crate::config::DockerUpstream {
url: "https://registry-1.docker.io".to_string(),
auth: None,
namespace: Some("docker.io".to_string()),
prefix: Some("docker-hub".to_string()),
}];
let cfg = docker_config_allow(upstreams.clone());
let c = super::canonicalize("docker-hub/library/nginx", &cfg);
assert_eq!(c.name, "library/nginx");
assert_eq!(c.namespace.as_deref(), Some("docker.io"));
assert_eq!(c.upstreams_to_try(&upstreams).len(), 1);
assert!(!c.denied);
}
#[test]
fn test_canonicalize_hostname_detection() {
let upstreams = vec![crate::config::DockerUpstream {
url: "https://registry-1.docker.io".to_string(),
auth: None,
namespace: Some("docker.io".to_string()),
prefix: None,
}];
let cfg = docker_config_allow(upstreams.clone());
let c = super::canonicalize("docker.io/library/nginx", &cfg);
assert_eq!(c.name, "library/nginx");
assert_eq!(c.namespace.as_deref(), Some("docker.io"));
assert_eq!(c.upstreams_to_try(&upstreams).len(), 1);
assert!(!c.denied);
let c2 = super::canonicalize("ghcr.io/requarks/wiki", &cfg);
assert_eq!(c2.name, "requarks/wiki");
assert_eq!(c2.namespace.as_deref(), Some("docker.io"));
assert_eq!(c2.upstreams_to_try(&upstreams).len(), 1); assert!(!c2.denied); }
#[test]
fn test_canonicalize_fallback() {
let upstreams = vec![crate::config::DockerUpstream {
url: "https://registry-1.docker.io".to_string(),
auth: None,
namespace: Some("docker.io".to_string()),
prefix: None,
}];
let cfg = docker_config_allow(upstreams.clone());
let c = super::canonicalize("library/nginx", &cfg);
assert_eq!(c.name, "library/nginx");
assert_eq!(c.namespace.as_deref(), Some("docker.io"));
assert_eq!(c.upstreams_to_try(&upstreams).len(), 1);
assert!(!c.denied);
let c2 = super::canonicalize("alpine", &cfg);
assert_eq!(c2.name, "alpine");
assert_eq!(c2.namespace.as_deref(), Some("docker.io"));
assert!(!c2.denied);
}
#[test]
fn test_canonicalize_empty_upstreams() {
let upstreams: Vec<crate::config::DockerUpstream> = vec![];
let cfg = docker_config_allow(upstreams.clone());
let c = super::canonicalize("library/nginx", &cfg);
assert_eq!(c.name, "library/nginx");
assert!(c.namespace.is_none());
assert!(c.upstreams_to_try(&upstreams).is_empty());
assert!(!c.denied);
}
#[test]
fn test_manifest_cache_fresh_tag_vs_digest() {
use super::manifest_cache_fresh;
let now = std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.map(|d| d.as_secs())
.unwrap_or(0);
assert!(manifest_cache_fresh(true, true, -1, Some(now)));
assert!(manifest_cache_fresh(true, true, 0, None));
assert!(manifest_cache_fresh(false, false, -1, Some(now)));
assert!(!manifest_cache_fresh(false, true, -1, Some(now)));
assert!(!manifest_cache_fresh(false, true, 0, Some(now)));
assert!(manifest_cache_fresh(false, true, 3600, Some(now)));
assert!(!manifest_cache_fresh(
false,
true,
3600,
Some(now.saturating_sub(7200))
));
assert!(!manifest_cache_fresh(false, true, 3600, None));
}
#[test]
fn test_canonicalize_multi_upstream() {
let upstreams = vec![
crate::config::DockerUpstream {
url: "https://registry-1.docker.io".to_string(),
auth: None,
namespace: Some("docker.io".to_string()),
prefix: Some("docker-hub".to_string()),
},
crate::config::DockerUpstream {
url: "https://ghcr.io".to_string(),
auth: None,
namespace: Some("ghcr.io".to_string()),
prefix: Some("ghcr".to_string()),
},
];
let cfg = docker_config_allow(upstreams.clone());
let c1 = super::canonicalize("ghcr/requarks/wiki", &cfg);
assert_eq!(c1.name, "requarks/wiki");
assert_eq!(c1.namespace.as_deref(), Some("ghcr.io"));
let targets = c1.upstreams_to_try(&upstreams);
assert_eq!(targets.len(), 1);
assert_eq!(targets[0].url, "https://ghcr.io");
assert!(!c1.denied);
let c2 = super::canonicalize("library/nginx", &cfg);
assert_eq!(c2.name, "library/nginx");
assert_eq!(c2.upstreams_to_try(&upstreams).len(), 2);
assert!(!c2.denied);
}
#[test]
fn test_canonicalize_deny_mode_blocks_unmatched() {
let upstreams = vec![
crate::config::DockerUpstream {
url: "https://registry-1.docker.io".to_string(),
auth: None,
namespace: Some("docker.io".to_string()),
prefix: Some("docker-hub".to_string()),
},
crate::config::DockerUpstream {
url: "https://ghcr.io".to_string(),
auth: None,
namespace: Some("ghcr.io".to_string()),
prefix: Some("ghcr".to_string()),
},
];
let cfg = docker_config_deny(upstreams.clone());
let c1 = super::canonicalize("ghcr/requarks/wiki", &cfg);
assert!(!c1.denied);
assert_eq!(c1.name, "requarks/wiki");
let c2 = super::canonicalize("library/nginx", &cfg);
assert!(c2.denied);
assert!(c2.denied_response().is_some());
let c3 = super::canonicalize("alpine", &cfg);
assert!(c3.denied);
let c4 = super::canonicalize("quay.io/prometheus/node-exporter", &cfg);
assert!(c4.denied);
}
#[test]
fn test_canonicalize_deny_mode_allows_known_hostname() {
let upstreams = vec![crate::config::DockerUpstream {
url: "https://registry-1.docker.io".to_string(),
auth: None,
namespace: Some("docker.io".to_string()),
prefix: None,
}];
let cfg = docker_config_deny(upstreams.clone());
let c = super::canonicalize("docker.io/library/nginx", &cfg);
assert!(!c.denied);
assert_eq!(c.name, "library/nginx");
assert_eq!(c.namespace.as_deref(), Some("docker.io"));
}
#[test]
fn test_denied_response_format() {
let upstreams = vec![crate::config::DockerUpstream {
url: "https://registry-1.docker.io".to_string(),
auth: None,
namespace: Some("docker.io".to_string()),
prefix: Some("hub".to_string()),
}];
let cfg = docker_config_deny(upstreams);
let c = super::canonicalize("library/nginx", &cfg);
assert!(c.denied);
let resp = c.denied_response();
assert!(resp.is_some());
let c2 = super::canonicalize("hub/library/nginx", &cfg);
assert!(!c2.denied);
assert!(c2.denied_response().is_none());
}
#[tokio::test]
async fn test_docker_proxy_blob_sha256_mismatch_rejected() {
use crate::config::DockerUpstream;
use crate::test_helpers::create_test_context_with_config;
use wiremock::matchers::{method, path};
use wiremock::{Mock, MockServer, ResponseTemplate};
let upstream = MockServer::start().await;
let requested_digest = format!(
"sha256:{}",
hex::encode(sha2::Sha256::digest(b"the real layer"))
);
let poisoned = b"poisoned bytes that do not hash to the requested digest".to_vec();
Mock::given(method("GET"))
.and(path(format!("/v2/library/test/blobs/{requested_digest}")))
.respond_with(ResponseTemplate::new(200).set_body_bytes(poisoned))
.mount(&upstream)
.await;
let ctx = create_test_context_with_config(|cfg| {
cfg.docker.upstreams = vec![DockerUpstream {
url: upstream.uri(),
auth: None,
namespace: None,
prefix: None,
}];
});
let response = send(
&ctx.app,
Method::GET,
&format!("/v2/library/test/blobs/{requested_digest}"),
Body::empty(),
)
.await;
assert_eq!(response.status(), StatusCode::BAD_GATEWAY);
let c = super::canonicalize("library/test", &ctx.state.config.docker);
let key = super::blob_key(c.namespace.as_deref(), &c.name, &requested_digest);
assert!(
ctx.state.storage.get(&key).await.is_err(),
"poisoned blob must not be cached"
);
let proxy_tmp =
std::path::Path::new(&ctx.state.config.storage.path).join("tmp/docker-proxy");
let leftover = std::fs::read_dir(&proxy_tmp)
.map(|rd| rd.filter_map(|e| e.ok()).count())
.unwrap_or(0);
assert_eq!(leftover, 0, "temp files must be cleaned up after rejection");
}
#[tokio::test]
async fn test_docker_proxy_tag_revalidates_on_upstream_change() {
use crate::config::DockerUpstream;
use crate::test_helpers::create_test_context_with_config;
use wiremock::matchers::{method, path};
use wiremock::{Mock, MockServer, ResponseTemplate};
let upstream = MockServer::start().await;
let ct = "application/vnd.oci.image.manifest.v1+json";
let manifest_old =
br#"{"schemaVersion":2,"mediaType":"application/vnd.oci.image.manifest.v1+json","config":{"digest":"sha256:aaaa"},"layers":[]}"#.to_vec();
let manifest_new =
br#"{"schemaVersion":2,"mediaType":"application/vnd.oci.image.manifest.v1+json","config":{"digest":"sha256:bbbb"},"layers":[{"digest":"sha256:cccc"}]}"#.to_vec();
Mock::given(method("GET"))
.and(path("/v2/library/test/manifests/latest"))
.respond_with(
ResponseTemplate::new(200)
.insert_header("Content-Type", ct)
.set_body_bytes(manifest_new.clone()),
)
.mount(&upstream)
.await;
let ctx = create_test_context_with_config(|cfg| {
cfg.docker.upstreams = vec![DockerUpstream {
url: upstream.uri(),
auth: None,
namespace: None,
prefix: None,
}];
});
let c = super::canonicalize("library/test", &ctx.state.config.docker);
let key = super::manifest_key(c.namespace.as_deref(), &c.name, "latest");
ctx.state
.storage
.put(&key, &manifest_old)
.await
.expect("seed cache");
let resp = send(
&ctx.app,
Method::GET,
"/v2/library/test/manifests/latest",
Body::empty(),
)
.await;
assert_eq!(resp.status(), StatusCode::OK);
assert_eq!(
body_bytes(resp).await.as_ref(),
manifest_new.as_slice(),
"a proxied tag with an outdated cached manifest must revalidate and return the upstream version (#638)"
);
}
#[tokio::test]
async fn test_docker_proxy_blob_sha256_match_served_and_cached() {
use crate::config::DockerUpstream;
use crate::test_helpers::create_test_context_with_config;
use wiremock::matchers::{method, path};
use wiremock::{Mock, MockServer, ResponseTemplate};
let upstream = MockServer::start().await;
let content = b"a genuine docker layer payload".to_vec();
let digest = format!("sha256:{}", hex::encode(sha2::Sha256::digest(&content)));
Mock::given(method("GET"))
.and(path(format!("/v2/library/ok/blobs/{digest}")))
.respond_with(ResponseTemplate::new(200).set_body_bytes(content.clone()))
.mount(&upstream)
.await;
let ctx = create_test_context_with_config(|cfg| {
cfg.docker.upstreams = vec![DockerUpstream {
url: upstream.uri(),
auth: None,
namespace: None,
prefix: None,
}];
});
let response = send(
&ctx.app,
Method::GET,
&format!("/v2/library/ok/blobs/{digest}"),
Body::empty(),
)
.await;
assert_eq!(response.status(), StatusCode::OK);
let body = body_bytes(response).await;
assert_eq!(
body.as_ref(),
content.as_slice(),
"served body must match upstream"
);
let c = super::canonicalize("library/ok", &ctx.state.config.docker);
let key = super::blob_key(c.namespace.as_deref(), &c.name, &digest);
assert!(
ctx.state.storage.get(&key).await.is_ok(),
"verified blob must be cached"
);
}
struct FailingPrefixBackend {
inner: crate::storage::LocalStorage,
fail_prefix: String,
}
#[async_trait::async_trait]
impl crate::storage::StorageBackend for FailingPrefixBackend {
async fn put(&self, key: &str, data: &[u8]) -> crate::storage::Result<()> {
self.inner.put(key, data).await
}
async fn get(&self, key: &str) -> crate::storage::Result<axum::body::Bytes> {
if key.starts_with(&self.fail_prefix) {
return Err(crate::storage::StorageError::Network("injected".into()));
}
self.inner.get(key).await
}
async fn delete(&self, key: &str) -> crate::storage::Result<()> {
self.inner.delete(key).await
}
async fn list(&self, prefix: &str) -> crate::storage::Result<Vec<String>> {
self.inner.list(prefix).await
}
async fn stat(&self, key: &str) -> Option<crate::storage::FileMeta> {
self.inner.stat(key).await
}
async fn health_check(&self) -> bool {
true
}
async fn total_size(&self) -> u64 {
self.inner.total_size().await
}
fn backend_name(&self) -> &'static str {
"failing-prefix-test"
}
async fn copy(&self, src: &str, dst: &str) -> crate::storage::Result<()> {
self.inner.copy(src, dst).await
}
async fn put_from_path(
&self,
key: &str,
src: &std::path::Path,
) -> crate::storage::Result<()> {
self.inner.put_from_path(key, src).await
}
async fn get_reader(
&self,
key: &str,
) -> crate::storage::Result<(
u64,
std::pin::Pin<Box<dyn tokio::io::AsyncRead + Send + Unpin>>,
)> {
if key.starts_with(&self.fail_prefix) {
return Err(crate::storage::StorageError::Network("injected".into()));
}
self.inner.get_reader(key).await
}
}
fn unreachable_upstream() -> crate::config::DockerUpstream {
crate::config::DockerUpstream {
url: "http://127.0.0.1:9".to_string(),
auth: None,
namespace: None,
prefix: None,
}
}
async fn get_manifest_via_dispatch(
state: &crate::AppState,
path: &str,
) -> axum::response::Response {
use axum::extract::{Path, State};
use axum::http::Uri;
use axum::Extension;
super::docker_v2_dispatch(
State(state.clone()),
Method::GET,
Path(path.to_string()),
Extension(crate::auth::NamespaceAuthority::Unrestricted),
format!("/v2/{path}").parse::<Uri>().unwrap(),
axum::http::HeaderMap::new(),
Body::empty(),
)
.await
}
#[tokio::test]
async fn manifest_get_fails_closed_on_transient_storage_error() {
use std::sync::Arc;
let ctx = create_test_context_with_config(|cfg| {
cfg.docker.upstreams = vec![unreachable_upstream()];
});
ctx.state
.storage
.put("docker/app/manifests/latest.json", b"{\"legacy\":true}")
.await
.unwrap();
let mut flaky = ctx.state.clone();
flaky.storage = crate::storage::Storage::from_backend(Arc::new(FailingPrefixBackend {
inner: crate::storage::LocalStorage::new(ctx._tempdir.path().to_str().unwrap()),
fail_prefix: "docker/127.0.0.1/".to_string(),
}));
let resp = get_manifest_via_dispatch(&flaky, "app/manifests/latest").await;
assert_eq!(
resp.status(),
StatusCode::INTERNAL_SERVER_ERROR,
"transient error on the namespaced key must fail closed, not fall back"
);
}
#[tokio::test]
async fn manifest_get_serves_hosted_legacy_on_ns_not_found() {
let ctx = create_test_context_with_config(|cfg| {
cfg.docker.upstreams = vec![unreachable_upstream()];
});
ctx.state
.storage
.put("docker/app/manifests/latest.json", b"{\"legacy\":true}")
.await
.unwrap();
let resp = get_manifest_via_dispatch(&ctx.state, "app/manifests/latest").await;
assert_eq!(resp.status(), StatusCode::OK);
assert!(
resp.headers().get("x-nora-stale").is_none(),
"hosted copy is authoritative — must not be served via the stale path"
);
let body = body_bytes(resp).await;
assert_eq!(body.as_ref(), b"{\"legacy\":true}");
}
}