1use std::collections::BTreeMap;
10use std::fmt;
11use std::future::Future;
12use std::net::IpAddr;
13use std::pin::Pin;
14use std::time::Duration;
15
16use reqwest::header::HeaderMap;
17use serde::{Deserialize, Serialize};
18use sha2::{Digest as _, Sha256};
19
20use super::reference::{ImageReference, ImageSelector, OciDigest, OciPlatform, RegistryName};
21
22pub const REGISTRY_DIAGNOSTIC_SCHEMA_VERSION: u32 = 2;
23pub const DEFAULT_MAX_REQUESTS: usize = 48;
24pub const DEFAULT_MAX_REDIRECTS: usize = 3;
25pub const DEFAULT_MAX_BODY_BYTES: usize = 4 * 1024 * 1024;
26pub const DEFAULT_MAX_MANIFEST_BYTES: usize = 2 * 1024 * 1024;
27pub const DEFAULT_BLOB_SAMPLE_BYTES: usize = 16 * 1024;
28pub const DEFAULT_REQUEST_TIMEOUT: Duration = Duration::from_secs(10);
29pub const DEFAULT_TOTAL_TIMEOUT: Duration = Duration::from_secs(30);
30const HARD_MAX_REQUESTS: usize = 64;
31const HARD_MAX_REDIRECTS: usize = 5;
32const HARD_MAX_BODY_BYTES: usize = 8 * 1024 * 1024;
33const HARD_MAX_MANIFEST_BYTES: usize = 4 * 1024 * 1024;
34const HARD_MAX_BLOB_SAMPLE_BYTES: usize = 1024 * 1024;
35const HARD_MAX_REQUEST_TIMEOUT: Duration = Duration::from_secs(60);
36const HARD_MAX_TOTAL_TIMEOUT: Duration = Duration::from_secs(5 * 60);
37const HARD_MAX_ENDPOINT_BYTES: usize = 2 * 1024;
38const HARD_MAX_ENDPOINT_PATH_BYTES: usize = 512;
39const HARD_MAX_MIRRORS: usize = 8;
40const HARD_MAX_CHALLENGE_BYTES: usize = 4 * 1024;
41const HARD_MAX_HEADER_VALUES: usize = 8;
42
43const ACCEPT_MANIFESTS: &str = concat!(
44 "application/vnd.oci.image.index.v1+json, ",
45 "application/vnd.oci.image.manifest.v1+json, ",
46 "application/vnd.docker.distribution.manifest.list.v2+json, ",
47 "application/vnd.docker.distribution.manifest.v2+json"
48);
49const OCI_INDEX: &str = "application/vnd.oci.image.index.v1+json";
50const OCI_MANIFEST: &str = "application/vnd.oci.image.manifest.v1+json";
51const DOCKER_INDEX: &str = "application/vnd.docker.distribution.manifest.list.v2+json";
52const DOCKER_MANIFEST: &str = "application/vnd.docker.distribution.manifest.v2+json";
53
54pub type RegistryTransportFuture<'a> =
55 Pin<Box<dyn Future<Output = Result<RegistryResponse, RegistryTransportError>> + Send + 'a>>;
56
57#[derive(Clone, PartialEq, Eq)]
60pub struct RegistryEndpoint {
61 url: reqwest::Url,
62 registry: RegistryName,
63 report_origin: String,
64}
65
66impl fmt::Debug for RegistryEndpoint {
67 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
68 formatter
69 .debug_struct("RegistryEndpoint")
70 .field("origin", &self.report_origin)
71 .field(
72 "has_path_prefix",
73 &(self.url.path() != "/" && !self.url.path().is_empty()),
74 )
75 .finish()
76 }
77}
78
79impl RegistryEndpoint {
80 pub fn for_registry(registry: RegistryName) -> Result<Self, RegistryProtocolError> {
81 let transport_authority = if registry.is_docker_hub() {
82 "registry-1.docker.io"
83 } else {
84 registry.as_str()
85 };
86 let mut endpoint = Self::parse(&format!("https://{transport_authority}"), false)?;
87 endpoint.registry = registry;
88 Ok(endpoint)
89 }
90
91 pub fn parse_https(value: &str) -> Result<Self, RegistryProtocolError> {
92 Self::parse(value, false)
93 }
94
95 fn parse(value: &str, allow_loopback_http: bool) -> Result<Self, RegistryProtocolError> {
96 if value.len() > HARD_MAX_ENDPOINT_BYTES
97 || value.contains('%')
98 || value.contains('\\')
99 || value
100 .split('/')
101 .any(|component| matches!(component, "." | ".."))
102 {
103 return Err(RegistryProtocolError::InvalidEndpoint);
104 }
105 let mut url =
106 reqwest::Url::parse(value).map_err(|_| RegistryProtocolError::InvalidEndpoint)?;
107 if !url.username().is_empty() || url.password().is_some() {
108 return Err(RegistryProtocolError::InvalidEndpoint);
109 }
110 if url.query().is_some() || url.fragment().is_some() {
111 return Err(RegistryProtocolError::InvalidEndpoint);
112 }
113 if url.scheme() != "https"
114 && !(allow_loopback_http && url.scheme() == "http" && url_is_loopback(&url))
115 {
116 return Err(RegistryProtocolError::InsecureEndpoint);
117 }
118 let host = url
119 .host_str()
120 .ok_or(RegistryProtocolError::InvalidEndpoint)?;
121 let authority = registry_authority(host, url.port())?;
122 let registry =
123 RegistryName::parse(&authority).map_err(|_| RegistryProtocolError::InvalidEndpoint)?;
124 let path = normalize_endpoint_path(url.path())?;
125 url.set_path(&path);
126 let report_origin = redacted_origin(&url);
127 Ok(Self {
128 url,
129 registry,
130 report_origin,
131 })
132 }
133
134 #[cfg(test)]
135 fn parse_loopback_http(value: &str) -> Result<Self, RegistryProtocolError> {
136 Self::parse(value, true)
137 }
138
139 pub fn registry(&self) -> &RegistryName {
140 &self.registry
141 }
142
143 pub fn report_origin(&self) -> &str {
145 &self.report_origin
146 }
147
148 pub fn is_https(&self) -> bool {
151 self.url.scheme() == "https"
152 }
153
154 fn url_for(&self, path: &str) -> Result<reqwest::Url, RegistryProtocolError> {
155 if !path.starts_with('/')
156 || path.contains(['?', '#', '\\'])
157 || path.split('/').any(|component| component == "..")
158 {
159 return Err(RegistryProtocolError::InvalidPath);
160 }
161 let prefix = self.url.path().trim_end_matches('/');
162 let mut url = self.url.clone();
163 url.set_path(&format!("{prefix}{path}"));
164 Ok(url)
165 }
166
167 fn origin_key(&self) -> String {
168 self.url.origin().ascii_serialization()
169 }
170
171 fn same_origin(&self, url: &reqwest::Url) -> bool {
172 self.origin_key() == url.origin().ascii_serialization()
173 }
174
175 fn allows_anonymous_token_realm(&self, url: &reqwest::Url) -> bool {
176 if self.same_origin(url) {
177 return true;
178 }
179 if self.url.port_or_known_default() != Some(443) || url.port_or_known_default() != Some(443)
180 {
181 return false;
182 }
183 matches!(
184 (
185 self.url.host_str().unwrap_or_default(),
186 url.host_str().unwrap_or_default()
187 ),
188 ("registry-1.docker.io", "auth.docker.io") | ("docker.m.daocloud.io", "m.daocloud.io")
189 )
190 }
191}
192
193fn normalize_endpoint_path(path: &str) -> Result<String, RegistryProtocolError> {
194 if path.len() > HARD_MAX_ENDPOINT_PATH_BYTES
195 || path.contains(['\\', '%'])
196 || path
197 .split('/')
198 .any(|component| matches!(component, "." | ".."))
199 {
200 return Err(RegistryProtocolError::InvalidEndpoint);
201 }
202 if path == "/" {
203 Ok(String::new())
204 } else {
205 Ok(path.trim_end_matches('/').to_owned())
206 }
207}
208
209fn registry_authority(host: &str, port: Option<u16>) -> Result<String, RegistryProtocolError> {
210 let host = if host
211 .parse::<IpAddr>()
212 .is_ok_and(|address| matches!(address, IpAddr::V6(_)))
213 {
214 format!("[{host}]")
215 } else {
216 host.to_owned()
217 };
218 Ok(match port {
219 Some(port) => format!("{host}:{port}"),
220 None => host,
221 })
222}
223
224fn url_is_loopback(url: &reqwest::Url) -> bool {
225 url.host_str().is_some_and(|host| {
226 let host = host.trim_start_matches('[').trim_end_matches(']');
227 host.eq_ignore_ascii_case("localhost")
228 || host
229 .parse::<IpAddr>()
230 .is_ok_and(|address| address.is_loopback())
231 })
232}
233
234#[derive(Clone, Copy, Debug, PartialEq, Eq)]
236pub struct RegistryLimits {
237 pub max_requests: usize,
238 pub max_redirects: usize,
239 pub max_body_bytes: usize,
240 pub max_manifest_bytes: usize,
241 pub blob_sample_bytes: usize,
242 pub request_timeout: Duration,
243 pub total_timeout: Duration,
244}
245
246impl Default for RegistryLimits {
247 fn default() -> Self {
248 Self {
249 max_requests: DEFAULT_MAX_REQUESTS,
250 max_redirects: DEFAULT_MAX_REDIRECTS,
251 max_body_bytes: DEFAULT_MAX_BODY_BYTES,
252 max_manifest_bytes: DEFAULT_MAX_MANIFEST_BYTES,
253 blob_sample_bytes: DEFAULT_BLOB_SAMPLE_BYTES,
254 request_timeout: DEFAULT_REQUEST_TIMEOUT,
255 total_timeout: DEFAULT_TOTAL_TIMEOUT,
256 }
257 }
258}
259
260impl RegistryLimits {
261 fn validate(self) -> Result<Self, RegistryProtocolError> {
262 if self.max_requests == 0
263 || self.max_requests > HARD_MAX_REQUESTS
264 || self.max_redirects > self.max_requests
265 || self.max_redirects > HARD_MAX_REDIRECTS
266 || self.max_body_bytes == 0
267 || self.max_body_bytes > HARD_MAX_BODY_BYTES
268 || self.max_manifest_bytes == 0
269 || self.max_manifest_bytes > HARD_MAX_MANIFEST_BYTES
270 || self.max_manifest_bytes > self.max_body_bytes
271 || self.blob_sample_bytes == 0
272 || self.blob_sample_bytes > HARD_MAX_BLOB_SAMPLE_BYTES
273 || self.blob_sample_bytes > self.max_body_bytes
274 || self.request_timeout.is_zero()
275 || self.total_timeout.is_zero()
276 || self.request_timeout > self.total_timeout
277 || self.request_timeout > HARD_MAX_REQUEST_TIMEOUT
278 || self.total_timeout > HARD_MAX_TOTAL_TIMEOUT
279 {
280 return Err(RegistryProtocolError::InvalidLimits);
281 }
282 Ok(self)
283 }
284}
285
286#[derive(Clone, Debug)]
289pub struct RegistryDiagnosticOptions {
290 pub upstream: RegistryEndpoint,
291 pub mirrors: Vec<RegistryEndpoint>,
292 pub image: Option<ImageReference>,
293 pub platform: Option<OciPlatform>,
294 pub limits: RegistryLimits,
295}
296
297impl RegistryDiagnosticOptions {
298 pub fn new(upstream: RegistryEndpoint) -> Self {
299 Self {
300 upstream,
301 mirrors: Vec::new(),
302 image: None,
303 platform: None,
304 limits: RegistryLimits::default(),
305 }
306 }
307
308 pub fn with_mirrors(mut self, mirrors: Vec<RegistryEndpoint>) -> Self {
309 self.mirrors = mirrors;
310 self
311 }
312
313 pub fn with_image(mut self, image: ImageReference) -> Self {
314 self.image = Some(image);
315 self
316 }
317
318 pub fn with_platform(mut self, platform: OciPlatform) -> Self {
319 self.platform = Some(platform);
320 self
321 }
322
323 pub fn with_limits(mut self, limits: RegistryLimits) -> Self {
324 self.limits = limits;
325 self
326 }
327}
328
329#[derive(Clone, Copy, Debug, PartialEq, Eq)]
330pub enum RegistryMethod {
331 Get,
332 Head,
333}
334
335#[derive(Clone)]
339pub struct RegistryRequest {
340 pub method: RegistryMethod,
341 pub url: reqwest::Url,
342 pub accept: Option<&'static str>,
343 pub range: Option<(u64, u64)>,
344 authorization: Option<BearerToken>,
345 pub max_body_bytes: usize,
346 pub timeout: Duration,
347}
348
349impl fmt::Debug for RegistryRequest {
350 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
351 formatter
352 .debug_struct("RegistryRequest")
353 .field("method", &self.method)
354 .field("origin", &redacted_origin(&self.url))
355 .field("accept", &self.accept)
356 .field("range", &self.range)
357 .field(
358 "authorization",
359 &self.authorization.as_ref().map(|_| "[redacted]"),
360 )
361 .field("max_body_bytes", &self.max_body_bytes)
362 .field("timeout", &self.timeout)
363 .finish()
364 }
365}
366
367impl RegistryRequest {
368 pub fn has_authorization(&self) -> bool {
369 self.authorization.is_some()
370 }
371}
372
373#[derive(Clone, Default)]
374pub struct RegistryResponse {
375 pub status: u16,
376 pub headers: BTreeMap<String, Vec<String>>,
377 pub body: Vec<u8>,
378 pub body_truncated: bool,
380}
381
382impl fmt::Debug for RegistryResponse {
383 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
384 formatter
385 .debug_struct("RegistryResponse")
386 .field("status", &self.status)
387 .field("header_names", &self.headers.keys().collect::<Vec<_>>())
388 .field("body_bytes", &self.body.len())
389 .field("body_truncated", &self.body_truncated)
390 .finish()
391 }
392}
393
394impl RegistryResponse {
395 pub fn new(status: u16) -> Self {
396 Self {
397 status,
398 headers: BTreeMap::new(),
399 body: Vec::new(),
400 body_truncated: false,
401 }
402 }
403
404 pub fn header(mut self, name: &str, value: &str) -> Self {
405 self.headers
406 .entry(name.to_ascii_lowercase())
407 .or_default()
408 .push(value.to_owned());
409 self
410 }
411
412 pub fn body(mut self, body: impl Into<Vec<u8>>) -> Self {
413 self.body = body.into();
414 self
415 }
416
417 #[cfg(test)]
418 fn truncated_body(mut self, body: impl Into<Vec<u8>>) -> Self {
419 self.body = body.into();
420 self.body_truncated = true;
421 self
422 }
423
424 fn header_values(&self, name: &str) -> impl Iterator<Item = &str> {
425 self.headers
426 .get(name)
427 .into_iter()
428 .flatten()
429 .map(String::as_str)
430 }
431
432 fn one_header(&self, name: &str) -> Option<&str> {
433 let mut values = self.header_values(name);
434 let first = values.next()?;
435 values.next().is_none().then_some(first)
436 }
437}
438
439#[derive(Clone, Copy, Debug, PartialEq, Eq, thiserror::Error)]
440pub enum RegistryTransportError {
441 #[error("registry transport failed")]
442 Failed,
443 #[error("registry request timed out")]
444 Timeout,
445 #[error("registry response exceeded its body limit")]
446 BodyTooLarge,
447 #[error("registry redirect violated policy")]
448 RedirectRejected,
449 #[error("registry request budget was exhausted")]
450 RequestLimit,
451}
452
453pub trait RegistryTransport: Send + Sync {
456 fn execute(&self, request: RegistryRequest) -> RegistryTransportFuture<'_>;
457}
458
459#[derive(Clone, Debug)]
462pub struct ReqwestRegistryTransport {
463 client: reqwest::Client,
464}
465
466impl ReqwestRegistryTransport {
467 pub fn new() -> Result<Self, RegistryTransportError> {
468 let client = reqwest::Client::builder()
469 .user_agent(concat!("osdk/", env!("CARGO_PKG_VERSION")))
470 .connect_timeout(DEFAULT_REQUEST_TIMEOUT)
471 .pool_idle_timeout(Duration::from_secs(30))
472 .redirect(reqwest::redirect::Policy::none())
473 .no_proxy()
474 .no_gzip()
475 .https_only(true)
476 .build()
477 .map_err(|_| RegistryTransportError::Failed)?;
478 Ok(Self { client })
479 }
480}
481
482impl RegistryTransport for ReqwestRegistryTransport {
483 fn execute(&self, request: RegistryRequest) -> RegistryTransportFuture<'_> {
484 Box::pin(async move {
485 if request.url.scheme() != "https" {
486 return Err(RegistryTransportError::Failed);
487 }
488 let mut builder = match request.method {
489 RegistryMethod::Get => self.client.get(request.url),
490 RegistryMethod::Head => self.client.head(request.url),
491 };
492 builder = builder
493 .timeout(request.timeout)
494 .header(reqwest::header::ACCEPT_ENCODING, "identity");
495 if let Some(accept) = request.accept {
496 builder = builder.header(reqwest::header::ACCEPT, accept);
497 }
498 if let Some((start, end)) = request.range {
499 builder = builder.header(reqwest::header::RANGE, format!("bytes={start}-{end}"));
500 }
501 if let Some(token) = &request.authorization {
502 builder = builder.bearer_auth(token.expose());
503 }
504 let response = builder.send().await.map_err(map_reqwest_error)?;
505 let status = response.status().as_u16();
506 let headers = copy_safe_response_headers(response.headers());
507 if request.method == RegistryMethod::Head {
508 return Ok(RegistryResponse {
509 status,
510 headers,
511 body: Vec::new(),
512 body_truncated: false,
513 });
514 }
515 let (body, body_truncated) =
516 read_bounded_response(response, request.max_body_bytes).await?;
517 if body_truncated && request.range.is_none() {
518 return Err(RegistryTransportError::BodyTooLarge);
519 }
520 Ok(RegistryResponse {
521 status,
522 headers,
523 body,
524 body_truncated,
525 })
526 })
527 }
528}
529
530fn map_reqwest_error(error: reqwest::Error) -> RegistryTransportError {
531 if error.is_timeout() {
532 RegistryTransportError::Timeout
533 } else {
534 RegistryTransportError::Failed
535 }
536}
537
538async fn read_bounded_response(
539 mut response: reqwest::Response,
540 limit: usize,
541) -> Result<(Vec<u8>, bool), RegistryTransportError> {
542 let declared_oversize = response
543 .content_length()
544 .is_some_and(|length| length > limit as u64);
545 let mut body = Vec::new();
546 while let Some(chunk) = response.chunk().await.map_err(map_reqwest_error)? {
547 let remaining = limit.saturating_sub(body.len());
548 if chunk.len() > remaining {
549 body.extend_from_slice(&chunk[..remaining]);
550 return Ok((body, true));
551 }
552 body.extend_from_slice(&chunk);
553 if body.len() == limit && declared_oversize {
554 return Ok((body, true));
555 }
556 }
557 Ok((body, false))
558}
559
560fn copy_safe_response_headers(headers: &HeaderMap) -> BTreeMap<String, Vec<String>> {
561 const SAFE: [&str; 7] = [
562 "www-authenticate",
563 "location",
564 "content-type",
565 "content-length",
566 "content-range",
567 "docker-content-digest",
568 "etag",
569 ];
570 let mut copied = BTreeMap::new();
571 for name in SAFE {
572 let values = headers
573 .get_all(name)
574 .iter()
575 .take(HARD_MAX_HEADER_VALUES)
576 .filter_map(|value| value.to_str().ok())
577 .filter(|value| value.len() <= HARD_MAX_CHALLENGE_BYTES)
578 .map(str::to_owned)
579 .collect::<Vec<_>>();
580 if !values.is_empty() {
581 copied.insert(name.to_owned(), values);
582 }
583 }
584 copied
585}
586
587#[derive(Clone)]
588struct BearerToken {
589 value: String,
590 bound_origin: String,
591}
592
593impl BearerToken {
594 fn expose(&self) -> &str {
595 &self.value
596 }
597
598 fn is_bound_to_url(&self, url: &reqwest::Url) -> bool {
599 self.bound_origin == url.origin().ascii_serialization()
600 }
601}
602
603impl fmt::Debug for BearerToken {
604 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
605 formatter.write_str("BearerToken([redacted])")
606 }
607}
608
609#[derive(Clone, Copy, Debug, PartialEq, Eq, thiserror::Error)]
610pub enum RegistryProtocolError {
611 #[error("invalid registry endpoint")]
612 InvalidEndpoint,
613 #[error("registry endpoints must use HTTPS")]
614 InsecureEndpoint,
615 #[error("invalid registry request path")]
616 InvalidPath,
617 #[error("invalid registry diagnostic limits")]
618 InvalidLimits,
619 #[error("image registry does not match the diagnostic upstream")]
620 RegistryMismatch,
621 #[error("too many registry mirrors")]
622 TooManyMirrors,
623}
624
625#[derive(Clone, Copy, Debug, PartialEq, Eq, Serialize)]
626#[serde(rename_all = "kebab-case")]
627pub enum RegistryDiagnosticStatus {
628 Healthy,
629 Degraded,
630 AuthenticationRequired,
631 AccessDenied,
632 RateLimited,
633 NotFound,
634 Unreachable,
635 TimedOut,
636 ProtocolError,
637 Corrupt,
638 Unsupported,
639 LimitExceeded,
640}
641
642#[derive(Clone, Copy, Debug, PartialEq, Eq, Serialize)]
643#[serde(rename_all = "kebab-case")]
644pub enum ApiCheckStatus {
645 Available,
646 BearerChallenge,
647 AuthenticationRequired,
648 AccessDenied,
649 RateLimited,
650 NotFound,
651 ServerError,
652 RedirectRejected,
653 InvalidChallenge,
654 UnexpectedResponse,
655 Unreachable,
656 TimedOut,
657 BodyTooLarge,
658 RequestLimit,
659 NotTested,
660}
661
662#[derive(Clone, Copy, Debug, PartialEq, Eq, Serialize)]
663#[serde(rename_all = "kebab-case")]
664pub enum ManifestCheckStatus {
665 Verified,
666 AuthenticationRequired,
667 AccessDenied,
668 RateLimited,
669 NotFound,
670 ServerError,
671 RedirectRejected,
672 InvalidMediaType,
673 InvalidManifest,
674 DigestMismatch,
675 SizeMismatch,
676 BodyTooLarge,
677 PlatformNotFound,
678 Unreachable,
679 TimedOut,
680 RequestLimit,
681 NotRequested,
682 NotTested,
683}
684
685#[derive(Clone, Copy, Debug, PartialEq, Eq, Serialize)]
686#[serde(rename_all = "kebab-case")]
687pub enum BlobRangeStatus {
688 Supported,
689 Ignored,
690 Unsatisfiable,
691 Malformed,
692 DigestMismatch,
693 AuthenticationRequired,
694 AccessDenied,
695 RateLimited,
696 NotFound,
697 ServerError,
698 RedirectRejected,
699 BodyTooLarge,
700 Unreachable,
701 TimedOut,
702 RequestLimit,
703 NotAvailable,
704 NotTested,
705}
706
707#[derive(Clone, Copy, Debug, PartialEq, Eq, Serialize)]
708#[serde(rename_all = "kebab-case")]
709pub enum MirrorCheckStatus {
710 Available,
711 Equivalent,
712 Diverged,
713 AuthenticationRequired,
714 AccessDenied,
715 RateLimited,
716 NotFound,
717 ServerError,
718 RedirectRejected,
719 InvalidResponse,
720 Unreachable,
721 TimedOut,
722 BodyTooLarge,
723 RequestLimit,
724 NotTested,
725}
726
727#[derive(Clone, Copy, Debug, PartialEq, Eq, Serialize)]
728#[serde(rename_all = "kebab-case")]
729pub enum ManifestKind {
730 Image,
731 Index,
732}
733
734#[derive(Clone, Debug, PartialEq, Eq, Serialize)]
735pub struct ApiCheck {
736 pub status: ApiCheckStatus,
737 pub http_status: Option<u16>,
738 pub bearer_challenge: bool,
739 pub challenge_service_present: bool,
740 pub challenge_scope_matches: Option<bool>,
741}
742
743impl Default for ApiCheck {
744 fn default() -> Self {
745 Self {
746 status: ApiCheckStatus::NotTested,
747 http_status: None,
748 bearer_challenge: false,
749 challenge_service_present: false,
750 challenge_scope_matches: None,
751 }
752 }
753}
754
755#[derive(Clone, Debug, PartialEq, Eq, Serialize)]
756pub struct ManifestCheck {
757 pub status: ManifestCheckStatus,
758 pub kind: Option<ManifestKind>,
759 pub media_type: Option<String>,
760 pub digest: Option<OciDigest>,
761 pub child_digest: Option<OciDigest>,
762 pub selected_platform: Option<OciPlatform>,
763 pub selected_os_version: Option<String>,
764 pub byte_size: Option<u64>,
765}
766
767impl Default for ManifestCheck {
768 fn default() -> Self {
769 Self {
770 status: ManifestCheckStatus::NotRequested,
771 kind: None,
772 media_type: None,
773 digest: None,
774 child_digest: None,
775 selected_platform: None,
776 selected_os_version: None,
777 byte_size: None,
778 }
779 }
780}
781
782#[derive(Clone, Debug, PartialEq, Eq, Serialize)]
783pub struct BlobRangeCheck {
784 pub status: BlobRangeStatus,
785 pub digest: Option<OciDigest>,
786 pub requested_bytes: u64,
787 pub received_bytes: u64,
788 pub total_bytes: Option<u64>,
789}
790
791impl Default for BlobRangeCheck {
792 fn default() -> Self {
793 Self {
794 status: BlobRangeStatus::NotTested,
795 digest: None,
796 requested_bytes: 0,
797 received_bytes: 0,
798 total_bytes: None,
799 }
800 }
801}
802
803#[derive(Clone, Debug, PartialEq, Eq, Serialize)]
804pub struct MirrorCheck {
805 pub order: usize,
806 pub origin: String,
807 pub status: MirrorCheckStatus,
808 pub digest: Option<OciDigest>,
809 pub blob_range: BlobRangeCheck,
810 pub elapsed_micros: u64,
811 pub recommended_rank: Option<usize>,
812}
813
814#[derive(Clone, Debug, PartialEq, Eq, Serialize)]
815pub struct RegistryDiagnosticReport {
816 pub schema_version: u32,
817 pub status: RegistryDiagnosticStatus,
818 pub upstream: RegistryName,
819 pub upstream_origin: String,
820 pub image: Option<ImageReference>,
821 pub requested_platform: Option<OciPlatform>,
822 pub api: ApiCheck,
823 pub manifest: ManifestCheck,
824 pub blob_range: BlobRangeCheck,
825 pub mirrors: Vec<MirrorCheck>,
826 pub recommended_mirror_order: Vec<usize>,
827 pub request_count: usize,
828}
829
830struct DiagnosticState<'a, T> {
831 transport: &'a T,
832 options: RegistryDiagnosticOptions,
833 request_count: usize,
834 started: tokio::time::Instant,
835}
836
837#[derive(Clone, Copy)]
838struct RequestPolicy {
839 accept: Option<&'static str>,
840 range: Option<(u64, u64)>,
841 max_body_bytes: usize,
842}
843
844pub async fn diagnose_registry<T: RegistryTransport>(
846 transport: &T,
847 options: RegistryDiagnosticOptions,
848) -> Result<RegistryDiagnosticReport, RegistryProtocolError> {
849 let limits = options.limits.validate()?;
850 if options.mirrors.len() > HARD_MAX_MIRRORS {
851 return Err(RegistryProtocolError::TooManyMirrors);
852 }
853 if let Some(image) = &options.image {
854 if image.registry() != options.upstream.registry() {
855 return Err(RegistryProtocolError::RegistryMismatch);
856 }
857 }
858 let mut state = DiagnosticState {
859 transport,
860 options,
861 request_count: 0,
862 started: tokio::time::Instant::now(),
863 };
864 state.options.limits = limits;
865
866 let mut report = RegistryDiagnosticReport {
867 schema_version: REGISTRY_DIAGNOSTIC_SCHEMA_VERSION,
868 status: RegistryDiagnosticStatus::Healthy,
869 upstream: state.options.upstream.registry().clone(),
870 upstream_origin: state.options.upstream.report_origin().to_owned(),
871 image: state.options.image.clone(),
872 requested_platform: state.options.platform.clone(),
873 api: ApiCheck::default(),
874 manifest: ManifestCheck::default(),
875 blob_range: BlobRangeCheck::default(),
876 mirrors: state
877 .options
878 .mirrors
879 .iter()
880 .enumerate()
881 .map(|(order, mirror)| MirrorCheck {
882 order,
883 origin: mirror.report_origin().to_owned(),
884 status: MirrorCheckStatus::NotTested,
885 digest: None,
886 blob_range: BlobRangeCheck::default(),
887 elapsed_micros: 0,
888 recommended_rank: None,
889 })
890 .collect(),
891 recommended_mirror_order: Vec::new(),
892 request_count: 0,
893 };
894
895 let upstream = state.options.upstream.clone();
896 let mut upstream_auth = None;
897 let api_probe = probe_api(&mut state, &upstream).await;
898 report.api = api_probe.check;
899 if let (Some(challenge), Some(image)) = (api_probe.challenge, state.options.image.clone()) {
900 let scope_matches =
901 valid_challenge_scope(challenge.scope.as_deref(), image.repository().as_str());
902 report.api.challenge_scope_matches = Some(scope_matches);
903 if scope_matches {
904 match obtain_anonymous_token(&mut state, &upstream, &challenge, &image).await {
905 Ok(token) => upstream_auth = Some(token),
906 Err(status) => report.api.status = status,
907 }
908 } else {
909 report.api.status = ApiCheckStatus::AuthenticationRequired;
910 }
911 }
912
913 if let Some(image) = state.options.image.clone() {
914 if matches!(
915 report.api.status,
916 ApiCheckStatus::Available | ApiCheckStatus::BearerChallenge
917 ) {
918 let platform = state.options.platform.clone();
919 let result = inspect_image(
920 &mut state,
921 &upstream,
922 &image,
923 platform.as_ref(),
924 upstream_auth.as_ref(),
925 )
926 .await;
927 report.manifest = result.manifest;
928 report.blob_range = result.blob_range;
929
930 if let Some(resolved) = result.resolved_digest {
931 for (index, mirror) in state.options.mirrors.clone().into_iter().enumerate() {
932 report.mirrors[index] = check_mirror(
933 &mut state,
934 index,
935 &mirror,
936 &image,
937 &resolved,
938 result.probe_manifest.as_ref(),
939 result.probe_blob.as_ref(),
940 )
941 .await;
942 }
943 }
944 }
945 } else {
946 report.manifest.status = ManifestCheckStatus::NotRequested;
947 for (index, mirror) in state.options.mirrors.clone().into_iter().enumerate() {
951 let started = tokio::time::Instant::now();
952 let api = probe_api(&mut state, &mirror).await;
953 report.mirrors[index] = mirror_failure(
954 index,
955 &mirror,
956 mirror_status_from_api(api.check.status),
957 elapsed_micros(started),
958 );
959 }
960 }
961
962 report.request_count = state.request_count;
963 rank_mirror_checks(&mut report);
964 report.status = aggregate_status(&report);
965 Ok(report)
966}
967
968impl<'a, T: RegistryTransport> DiagnosticState<'a, T> {
969 async fn execute(
970 &mut self,
971 endpoint: &RegistryEndpoint,
972 method: RegistryMethod,
973 path: &str,
974 authorization: Option<&BearerToken>,
975 policy: RequestPolicy,
976 ) -> Result<RegistryResponse, RegistryTransportError> {
977 let mut url = endpoint
978 .url_for(path)
979 .map_err(|_| RegistryTransportError::Failed)?;
980 let mut redirects = 0;
981 let authorization_allowed = authorization.is_some_and(|token| token.is_bound_to_url(&url));
982 loop {
983 if self.request_count >= self.options.limits.max_requests {
984 return Err(RegistryTransportError::RequestLimit);
985 }
986 let elapsed = self.started.elapsed();
987 let remaining = self
988 .options
989 .limits
990 .total_timeout
991 .checked_sub(elapsed)
992 .ok_or(RegistryTransportError::Timeout)?;
993 let timeout = remaining.min(self.options.limits.request_timeout);
994 let request_authorization = authorization.filter(|_| authorization_allowed);
995 let request = RegistryRequest {
996 method,
997 url: url.clone(),
998 accept: policy.accept,
999 range: policy.range,
1000 authorization: request_authorization.cloned(),
1001 max_body_bytes: policy.max_body_bytes,
1002 timeout,
1003 };
1004 self.request_count += 1;
1005 let response = tokio::time::timeout(timeout, self.transport.execute(request))
1006 .await
1007 .map_err(|_| RegistryTransportError::Timeout)??;
1008 if response.body.len() > policy.max_body_bytes
1009 || (response.body_truncated && policy.range.is_none())
1010 {
1011 return Err(RegistryTransportError::BodyTooLarge);
1012 }
1013 if !(300..400).contains(&response.status) {
1014 return Ok(response);
1015 }
1016 if redirects >= self.options.limits.max_redirects {
1017 return Err(RegistryTransportError::RedirectRejected);
1018 }
1019 let location = response
1020 .one_header("location")
1021 .ok_or(RegistryTransportError::RedirectRejected)?;
1022 let next = url
1023 .join(location)
1024 .map_err(|_| RegistryTransportError::RedirectRejected)?;
1025 if next.scheme() != "https"
1026 || !endpoint.same_origin(&next)
1027 || !next.username().is_empty()
1028 || next.password().is_some()
1029 || next.fragment().is_some()
1030 || next.host_str().is_none()
1031 {
1032 return Err(RegistryTransportError::RedirectRejected);
1033 }
1034 redirects += 1;
1035 url = next;
1036 }
1037 }
1038}
1039
1040struct ApiProbe {
1041 check: ApiCheck,
1042 challenge: Option<BearerChallenge>,
1043}
1044
1045async fn probe_api<T: RegistryTransport>(
1046 state: &mut DiagnosticState<'_, T>,
1047 endpoint: &RegistryEndpoint,
1048) -> ApiProbe {
1049 let response = state
1050 .execute(
1051 endpoint,
1052 RegistryMethod::Get,
1053 "/v2/",
1054 None,
1055 RequestPolicy {
1056 accept: None,
1057 range: None,
1058 max_body_bytes: state.options.limits.max_body_bytes.min(8 * 1024),
1059 },
1060 )
1061 .await;
1062 let response = match response {
1063 Ok(response) => response,
1064 Err(error) => {
1065 return ApiProbe {
1066 check: ApiCheck {
1067 status: api_transport_status(error),
1068 ..ApiCheck::default()
1069 },
1070 challenge: None,
1071 };
1072 }
1073 };
1074 if response.status == 200 {
1075 return ApiProbe {
1076 check: ApiCheck {
1077 status: ApiCheckStatus::Available,
1078 http_status: Some(200),
1079 ..ApiCheck::default()
1080 },
1081 challenge: None,
1082 };
1083 }
1084 if response.status == 401 {
1085 return match parse_bearer_challenge(&response) {
1086 Some(challenge) => ApiProbe {
1087 check: ApiCheck {
1088 status: ApiCheckStatus::BearerChallenge,
1089 http_status: Some(401),
1090 bearer_challenge: true,
1091 challenge_service_present: challenge.service.is_some(),
1092 challenge_scope_matches: None,
1093 },
1094 challenge: Some(challenge),
1095 },
1096 None => ApiProbe {
1097 check: ApiCheck {
1098 status: ApiCheckStatus::InvalidChallenge,
1099 http_status: Some(401),
1100 ..ApiCheck::default()
1101 },
1102 challenge: None,
1103 },
1104 };
1105 }
1106 ApiProbe {
1107 check: ApiCheck {
1108 status: api_status_for_http(response.status),
1109 http_status: Some(response.status),
1110 ..ApiCheck::default()
1111 },
1112 challenge: None,
1113 }
1114}
1115
1116async fn obtain_anonymous_token<T: RegistryTransport>(
1117 state: &mut DiagnosticState<'_, T>,
1118 endpoint: &RegistryEndpoint,
1119 challenge: &BearerChallenge,
1120 image: &ImageReference,
1121) -> Result<BearerToken, ApiCheckStatus> {
1122 if !valid_challenge_scope(challenge.scope.as_deref(), image.repository().as_str()) {
1123 return Err(ApiCheckStatus::AuthenticationRequired);
1124 }
1125 let realm =
1126 reqwest::Url::parse(&challenge.realm).map_err(|_| ApiCheckStatus::InvalidChallenge)?;
1127 if realm.scheme() != "https"
1128 || !endpoint.allows_anonymous_token_realm(&realm)
1129 || !realm.username().is_empty()
1130 || realm.password().is_some()
1131 || realm.fragment().is_some()
1132 {
1133 return Err(ApiCheckStatus::InvalidChallenge);
1134 }
1135 let mut token_url = realm;
1139 {
1140 let mut query = token_url.query_pairs_mut();
1141 if let Some(service) = challenge.service.as_deref() {
1142 query.append_pair("service", service);
1143 }
1144 query.append_pair(
1145 "scope",
1146 &format!("repository:{}:pull", image.repository().as_str()),
1147 );
1148 }
1149 let elapsed = state.started.elapsed();
1150 let remaining = state
1151 .options
1152 .limits
1153 .total_timeout
1154 .checked_sub(elapsed)
1155 .ok_or(ApiCheckStatus::TimedOut)?;
1156 let timeout = remaining.min(state.options.limits.request_timeout);
1157 if state.request_count >= state.options.limits.max_requests {
1158 return Err(ApiCheckStatus::RequestLimit);
1159 }
1160 let body_limit = state.options.limits.max_body_bytes.min(64 * 1024);
1161 state.request_count += 1;
1162 let response = tokio::time::timeout(
1163 timeout,
1164 state.transport.execute(RegistryRequest {
1165 method: RegistryMethod::Get,
1166 url: token_url,
1167 accept: Some("application/json"),
1168 range: None,
1169 authorization: None,
1170 max_body_bytes: body_limit,
1171 timeout,
1172 }),
1173 )
1174 .await
1175 .map_err(|_| ApiCheckStatus::TimedOut)?
1176 .map_err(api_transport_status)?;
1177 if response.status != 200 {
1178 return Err(api_status_for_http(response.status));
1179 }
1180 if response.body.len() > body_limit {
1181 return Err(ApiCheckStatus::BodyTooLarge);
1182 }
1183 #[derive(Deserialize)]
1184 struct TokenResponse {
1185 #[serde(default)]
1186 token: Option<String>,
1187 #[serde(default)]
1188 access_token: Option<String>,
1189 }
1190 let token: TokenResponse =
1191 serde_json::from_slice(&response.body).map_err(|_| ApiCheckStatus::UnexpectedResponse)?;
1192 let token = token
1193 .token
1194 .or(token.access_token)
1195 .ok_or(ApiCheckStatus::UnexpectedResponse)?;
1196 if token.is_empty()
1197 || token.len() > 16 * 1024
1198 || token.bytes().any(|byte| byte.is_ascii_control())
1199 {
1200 return Err(ApiCheckStatus::UnexpectedResponse);
1201 }
1202 Ok(BearerToken {
1203 value: token,
1204 bound_origin: endpoint.origin_key(),
1205 })
1206}
1207
1208#[derive(Clone)]
1209struct BearerChallenge {
1210 realm: String,
1211 service: Option<String>,
1212 scope: Option<String>,
1213}
1214
1215impl fmt::Debug for BearerChallenge {
1216 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
1217 formatter
1218 .debug_struct("BearerChallenge")
1219 .field("realm", &"[redacted]")
1220 .field("service_present", &self.service.is_some())
1221 .field("scope_present", &self.scope.is_some())
1222 .finish()
1223 }
1224}
1225
1226fn parse_bearer_challenge(response: &RegistryResponse) -> Option<BearerChallenge> {
1227 response
1228 .header_values("www-authenticate")
1229 .find_map(parse_one_bearer_challenge)
1230}
1231
1232fn parse_one_bearer_challenge(value: &str) -> Option<BearerChallenge> {
1233 if value.len() > HARD_MAX_CHALLENGE_BYTES {
1234 return None;
1235 }
1236 let (scheme, parameters) = value.trim().split_once(char::is_whitespace)?;
1237 if !scheme.eq_ignore_ascii_case("bearer") {
1238 return None;
1239 }
1240 let parameters = parse_auth_parameters(parameters)?;
1241 let realm = parameters.get("realm")?.to_owned();
1242 let realm_url = reqwest::Url::parse(&realm).ok()?;
1243 if realm.is_empty()
1244 || realm_url.scheme() != "https"
1245 || realm_url.host_str().is_none()
1246 || !realm_url.username().is_empty()
1247 || realm_url.password().is_some()
1248 || realm_url.fragment().is_some()
1249 {
1250 return None;
1251 }
1252 Some(BearerChallenge {
1253 realm,
1254 service: parameters.get("service").cloned(),
1255 scope: parameters.get("scope").cloned(),
1256 })
1257}
1258
1259fn parse_auth_parameters(value: &str) -> Option<BTreeMap<String, String>> {
1260 let mut result = BTreeMap::new();
1261 let bytes = value.as_bytes();
1262 let mut cursor = 0;
1263 while cursor < bytes.len() {
1264 while bytes
1265 .get(cursor)
1266 .is_some_and(|byte| byte.is_ascii_whitespace() || *byte == b',')
1267 {
1268 cursor += 1;
1269 }
1270 let key_start = cursor;
1271 while bytes
1272 .get(cursor)
1273 .is_some_and(|byte| byte.is_ascii_alphanumeric() || *byte == b'-')
1274 {
1275 cursor += 1;
1276 }
1277 if cursor == key_start || bytes.get(cursor) != Some(&b'=') {
1278 return None;
1279 }
1280 let key = value.get(key_start..cursor)?.to_ascii_lowercase();
1281 cursor += 1;
1282 if bytes.get(cursor) != Some(&b'"') {
1283 return None;
1284 }
1285 cursor += 1;
1286 let mut parsed = String::new();
1287 let mut closed = false;
1288 while let Some(byte) = bytes.get(cursor).copied() {
1289 cursor += 1;
1290 match byte {
1291 b'"' => {
1292 closed = true;
1293 break;
1294 }
1295 b'\\' => {
1296 let escaped = bytes.get(cursor).copied()?;
1297 cursor += 1;
1298 if !escaped.is_ascii() || escaped.is_ascii_control() {
1299 return None;
1300 }
1301 parsed.push(char::from(escaped));
1302 }
1303 byte if byte.is_ascii_control() || !byte.is_ascii() => return None,
1304 byte => parsed.push(char::from(byte)),
1305 }
1306 }
1307 if !closed || result.insert(key, parsed).is_some() {
1308 return None;
1309 }
1310 while bytes.get(cursor).is_some_and(u8::is_ascii_whitespace) {
1311 cursor += 1;
1312 }
1313 if cursor < bytes.len() && bytes[cursor] != b',' {
1314 return None;
1315 }
1316 }
1317 Some(result)
1318}
1319
1320fn valid_challenge_scope(scope: Option<&str>, repository: &str) -> bool {
1321 match scope {
1322 None => true,
1323 Some(scope) => scope == format!("repository:{repository}:pull"),
1324 }
1325}
1326
1327struct ImageInspection {
1328 manifest: ManifestCheck,
1329 blob_range: BlobRangeCheck,
1330 resolved_digest: Option<OciDigest>,
1331 probe_manifest: Option<Descriptor>,
1332 probe_blob: Option<Descriptor>,
1333}
1334
1335async fn inspect_image<T: RegistryTransport>(
1336 state: &mut DiagnosticState<'_, T>,
1337 endpoint: &RegistryEndpoint,
1338 image: &ImageReference,
1339 platform: Option<&OciPlatform>,
1340 authorization: Option<&BearerToken>,
1341) -> ImageInspection {
1342 let initial_path = manifest_path(image, image.selector().as_str());
1343 let initial = fetch_manifest(state, endpoint, &initial_path, authorization, None).await;
1344 let mut manifest = match initial {
1345 Ok(manifest) => manifest,
1346 Err(status) => {
1347 return ImageInspection {
1348 manifest: ManifestCheck {
1349 status,
1350 ..ManifestCheck::default()
1351 },
1352 blob_range: BlobRangeCheck::default(),
1353 resolved_digest: None,
1354 probe_manifest: None,
1355 probe_blob: None,
1356 };
1357 }
1358 };
1359 if let ImageSelector::Digest(expected) = image.selector() {
1360 if &manifest.digest != expected {
1361 return ImageInspection {
1362 manifest: ManifestCheck {
1363 status: ManifestCheckStatus::DigestMismatch,
1364 digest: Some(manifest.digest),
1365 byte_size: Some(manifest.body.len() as u64),
1366 ..ManifestCheck::default()
1367 },
1368 blob_range: BlobRangeCheck::default(),
1369 resolved_digest: None,
1370 probe_manifest: None,
1371 probe_blob: None,
1372 };
1373 }
1374 }
1375 let resolved_digest = manifest.digest.clone();
1376 let top_kind = manifest.kind;
1377 let top_media_type = manifest.media_type.clone();
1378 let top_byte_size = manifest.body.len() as u64;
1379 let mut selected_platform = None;
1380 let mut selected_os_version = None;
1381 let mut child_digest = None;
1382 let mut probe_manifest = None;
1383
1384 if manifest.kind == ManifestKind::Index {
1385 let Some(platform) = platform else {
1386 return ImageInspection {
1387 manifest: ManifestCheck {
1388 status: ManifestCheckStatus::PlatformNotFound,
1389 kind: Some(ManifestKind::Index),
1390 media_type: Some(manifest.media_type),
1391 digest: Some(resolved_digest.clone()),
1392 byte_size: Some(manifest.body.len() as u64),
1393 ..ManifestCheck::default()
1394 },
1395 blob_range: BlobRangeCheck::default(),
1396 resolved_digest: Some(resolved_digest),
1397 probe_manifest: None,
1398 probe_blob: None,
1399 };
1400 };
1401 let descriptor = match select_platform_descriptor(&manifest.body, platform) {
1402 Some(descriptor) => descriptor,
1403 None => {
1404 return ImageInspection {
1405 manifest: ManifestCheck {
1406 status: ManifestCheckStatus::PlatformNotFound,
1407 kind: Some(ManifestKind::Index),
1408 media_type: Some(manifest.media_type),
1409 digest: Some(resolved_digest.clone()),
1410 byte_size: Some(manifest.body.len() as u64),
1411 ..ManifestCheck::default()
1412 },
1413 blob_range: BlobRangeCheck::default(),
1414 resolved_digest: Some(resolved_digest),
1415 probe_manifest: None,
1416 probe_blob: None,
1417 };
1418 }
1419 };
1420 let path = manifest_path(image, descriptor.digest.as_str());
1421 let child = fetch_manifest(
1422 state,
1423 endpoint,
1424 &path,
1425 authorization,
1426 Some((&descriptor.digest, descriptor.size)),
1427 )
1428 .await;
1429 manifest = match child {
1430 Ok(child) if child.kind == ManifestKind::Image => child,
1431 Ok(_) => {
1432 return ImageInspection {
1433 manifest: ManifestCheck {
1434 status: ManifestCheckStatus::InvalidMediaType,
1435 kind: Some(ManifestKind::Index),
1436 digest: Some(resolved_digest.clone()),
1437 child_digest: Some(descriptor.digest),
1438 selected_platform: Some(platform.clone()),
1439 selected_os_version: descriptor
1440 .platform
1441 .as_ref()
1442 .and_then(|platform| platform.os_version.clone()),
1443 ..ManifestCheck::default()
1444 },
1445 blob_range: BlobRangeCheck::default(),
1446 resolved_digest: Some(resolved_digest),
1447 probe_manifest: None,
1448 probe_blob: None,
1449 };
1450 }
1451 Err(status) => {
1452 return ImageInspection {
1453 manifest: ManifestCheck {
1454 status,
1455 kind: Some(ManifestKind::Index),
1456 digest: Some(resolved_digest.clone()),
1457 child_digest: Some(descriptor.digest),
1458 selected_platform: Some(platform.clone()),
1459 selected_os_version: descriptor
1460 .platform
1461 .as_ref()
1462 .and_then(|platform| platform.os_version.clone()),
1463 ..ManifestCheck::default()
1464 },
1465 blob_range: BlobRangeCheck::default(),
1466 resolved_digest: Some(resolved_digest),
1467 probe_manifest: None,
1468 probe_blob: None,
1469 };
1470 }
1471 };
1472 child_digest = Some(descriptor.digest.clone());
1473 probe_manifest = Some(descriptor.clone());
1474 selected_platform = Some(platform.clone());
1475 selected_os_version = descriptor.platform.and_then(|platform| platform.os_version);
1476 }
1477
1478 let blob = first_layer_descriptor(&manifest.body)
1479 .expect("a verified image manifest has at least one valid layer");
1480 let blob_range = check_blob_range(state, endpoint, image, authorization, blob.clone()).await;
1481 ImageInspection {
1482 manifest: ManifestCheck {
1483 status: ManifestCheckStatus::Verified,
1484 kind: Some(top_kind),
1485 media_type: Some(top_media_type),
1486 digest: Some(resolved_digest.clone()),
1487 child_digest,
1488 selected_platform,
1489 selected_os_version,
1490 byte_size: Some(top_byte_size),
1491 },
1492 blob_range,
1493 resolved_digest: Some(resolved_digest),
1494 probe_manifest,
1495 probe_blob: Some(blob),
1496 }
1497}
1498
1499struct VerifiedManifest {
1500 kind: ManifestKind,
1501 media_type: String,
1502 digest: OciDigest,
1503 body: Vec<u8>,
1504}
1505
1506async fn fetch_manifest<T: RegistryTransport>(
1507 state: &mut DiagnosticState<'_, T>,
1508 endpoint: &RegistryEndpoint,
1509 path: &str,
1510 authorization: Option<&BearerToken>,
1511 expected: Option<(&OciDigest, u64)>,
1512) -> Result<VerifiedManifest, ManifestCheckStatus> {
1513 let response = state
1514 .execute(
1515 endpoint,
1516 RegistryMethod::Get,
1517 path,
1518 authorization,
1519 RequestPolicy {
1520 accept: Some(ACCEPT_MANIFESTS),
1521 range: None,
1522 max_body_bytes: state.options.limits.max_manifest_bytes,
1523 },
1524 )
1525 .await
1526 .map_err(manifest_transport_status)?;
1527 if response.status != 200 {
1528 return Err(manifest_status_for_http(response.status));
1529 }
1530 if response.body.len() > state.options.limits.max_manifest_bytes {
1531 return Err(ManifestCheckStatus::BodyTooLarge);
1532 }
1533 if let Some(content_length) = response.one_header("content-length") {
1534 if content_length.parse::<u64>().ok() != Some(response.body.len() as u64) {
1535 return Err(ManifestCheckStatus::SizeMismatch);
1536 }
1537 }
1538 let digest = sha256_digest(&response.body).map_err(|_| ManifestCheckStatus::InvalidManifest)?;
1539 if let Some(header) = response.one_header("docker-content-digest") {
1540 if OciDigest::parse(header).ok().as_ref() != Some(&digest) {
1541 return Err(ManifestCheckStatus::DigestMismatch);
1542 }
1543 }
1544 if let Some((expected_digest, expected_size)) = expected {
1545 if expected_digest != &digest {
1546 return Err(ManifestCheckStatus::DigestMismatch);
1547 }
1548 if expected_size != response.body.len() as u64 {
1549 return Err(ManifestCheckStatus::SizeMismatch);
1550 }
1551 }
1552 let media_type = response
1553 .one_header("content-type")
1554 .and_then(parse_media_type)
1555 .ok_or(ManifestCheckStatus::InvalidMediaType)?;
1556 let kind = manifest_kind(&media_type).ok_or(ManifestCheckStatus::InvalidMediaType)?;
1557 let body_value: ManifestEnvelope =
1558 serde_json::from_slice(&response.body).map_err(|_| ManifestCheckStatus::InvalidManifest)?;
1559 if body_value.schema_version != 2 || body_value.media_type.as_deref() != Some(&media_type) {
1560 return Err(ManifestCheckStatus::InvalidManifest);
1561 }
1562 let structurally_valid = match kind {
1563 ManifestKind::Image => validate_image_manifest(&response.body),
1564 ManifestKind::Index => validate_image_index(&response.body),
1565 };
1566 if !structurally_valid {
1567 return Err(ManifestCheckStatus::InvalidManifest);
1568 }
1569 Ok(VerifiedManifest {
1570 kind,
1571 media_type,
1572 digest,
1573 body: response.body,
1574 })
1575}
1576
1577fn manifest_kind(media_type: &str) -> Option<ManifestKind> {
1578 match media_type {
1579 OCI_INDEX | DOCKER_INDEX => Some(ManifestKind::Index),
1580 OCI_MANIFEST | DOCKER_MANIFEST => Some(ManifestKind::Image),
1581 _ => None,
1582 }
1583}
1584
1585fn parse_media_type(value: &str) -> Option<String> {
1586 let media_type = value.split(';').next()?.trim().to_ascii_lowercase();
1587 manifest_kind(&media_type).map(|_| media_type)
1588}
1589
1590#[derive(Deserialize)]
1591struct ManifestEnvelope {
1592 #[serde(rename = "schemaVersion")]
1593 schema_version: u32,
1594 #[serde(rename = "mediaType", default)]
1595 media_type: Option<String>,
1596}
1597
1598#[derive(Deserialize)]
1599struct ImageIndex {
1600 manifests: Vec<Descriptor>,
1601}
1602
1603#[derive(Deserialize)]
1604struct ImageManifest {
1605 config: Descriptor,
1606 layers: Vec<Descriptor>,
1607}
1608
1609#[derive(Clone, Deserialize)]
1610struct Descriptor {
1611 #[serde(rename = "mediaType")]
1612 media_type: String,
1613 digest: OciDigest,
1614 size: u64,
1615 #[serde(default)]
1616 platform: Option<DescriptorPlatform>,
1617}
1618
1619#[derive(Clone, Deserialize)]
1620struct DescriptorPlatform {
1621 os: String,
1622 architecture: String,
1623 #[serde(default)]
1624 variant: Option<String>,
1625 #[serde(default, rename = "os.version")]
1626 os_version: Option<String>,
1627 #[serde(default, rename = "os.features")]
1628 os_features: Vec<String>,
1629}
1630
1631fn select_platform_descriptor(body: &[u8], requested: &OciPlatform) -> Option<Descriptor> {
1632 let index: ImageIndex = serde_json::from_slice(body).ok()?;
1633 let mut candidates = index.manifests.into_iter().filter(|descriptor| {
1634 if !valid_manifest_descriptor(descriptor) {
1635 return false;
1636 }
1637 let Some(platform) = &descriptor.platform else {
1638 return false;
1639 };
1640 if !platform.os_features.is_empty() {
1641 return false;
1642 }
1643 OciPlatform::new(
1644 &platform.os,
1645 &platform.architecture,
1646 platform.variant.as_deref(),
1647 )
1648 .is_ok_and(|candidate| &candidate == requested)
1649 });
1650 let selected = candidates.next()?;
1651 candidates.next().is_none().then_some(selected)
1652}
1653
1654fn first_layer_descriptor(body: &[u8]) -> Option<Descriptor> {
1655 let manifest = serde_json::from_slice::<ImageManifest>(body).ok()?;
1656 if !valid_config_descriptor(&manifest.config)
1657 || manifest.layers.is_empty()
1658 || manifest
1659 .layers
1660 .iter()
1661 .any(|layer| !valid_layer_descriptor(layer))
1662 {
1663 return None;
1664 }
1665 manifest
1666 .layers
1667 .into_iter()
1668 .min_by_key(|descriptor| descriptor.size)
1669}
1670
1671fn validate_image_manifest(body: &[u8]) -> bool {
1672 let Ok(manifest) = serde_json::from_slice::<ImageManifest>(body) else {
1673 return false;
1674 };
1675 valid_config_descriptor(&manifest.config)
1676 && !manifest.layers.is_empty()
1677 && manifest.layers.iter().all(valid_layer_descriptor)
1678}
1679
1680fn validate_image_index(body: &[u8]) -> bool {
1681 serde_json::from_slice::<ImageIndex>(body).is_ok_and(|index| {
1682 !index.manifests.is_empty() && index.manifests.iter().all(valid_manifest_descriptor)
1683 })
1684}
1685
1686fn valid_manifest_descriptor(descriptor: &Descriptor) -> bool {
1687 descriptor.size > 0
1688 && matches!(
1689 descriptor.media_type.as_str(),
1690 OCI_INDEX | OCI_MANIFEST | DOCKER_INDEX | DOCKER_MANIFEST
1691 )
1692}
1693
1694fn valid_config_descriptor(descriptor: &Descriptor) -> bool {
1695 descriptor.size > 0
1696 && matches!(
1697 descriptor.media_type.as_str(),
1698 "application/vnd.oci.image.config.v1+json"
1699 | "application/vnd.docker.container.image.v1+json"
1700 )
1701 && descriptor.platform.is_none()
1702}
1703
1704fn valid_layer_descriptor(descriptor: &Descriptor) -> bool {
1705 descriptor.size > 0
1706 && descriptor.platform.is_none()
1707 && matches!(
1708 descriptor.media_type.as_str(),
1709 "application/vnd.oci.image.layer.v1.tar"
1710 | "application/vnd.oci.image.layer.v1.tar+gzip"
1711 | "application/vnd.oci.image.layer.v1.tar+zstd"
1712 | "application/vnd.oci.image.layer.nondistributable.v1.tar"
1713 | "application/vnd.oci.image.layer.nondistributable.v1.tar+gzip"
1714 | "application/vnd.oci.image.layer.nondistributable.v1.tar+zstd"
1715 | "application/vnd.docker.image.rootfs.diff.tar.gzip"
1716 | "application/vnd.docker.image.rootfs.foreign.diff.tar.gzip"
1717 )
1718}
1719
1720async fn check_blob_range<T: RegistryTransport>(
1721 state: &mut DiagnosticState<'_, T>,
1722 endpoint: &RegistryEndpoint,
1723 image: &ImageReference,
1724 authorization: Option<&BearerToken>,
1725 descriptor: Descriptor,
1726) -> BlobRangeCheck {
1727 let requested = descriptor
1728 .size
1729 .min(state.options.limits.blob_sample_bytes as u64);
1730 let mut check = BlobRangeCheck {
1731 digest: Some(descriptor.digest.clone()),
1732 requested_bytes: requested,
1733 ..BlobRangeCheck::default()
1734 };
1735 if requested == 0 {
1736 check.status = BlobRangeStatus::Malformed;
1737 return check;
1738 }
1739 let path = blob_path(image, descriptor.digest.as_str());
1740 let response = state
1741 .execute(
1742 endpoint,
1743 RegistryMethod::Get,
1744 &path,
1745 authorization,
1746 RequestPolicy {
1747 accept: None,
1748 range: Some((0, requested - 1)),
1749 max_body_bytes: state.options.limits.blob_sample_bytes,
1750 },
1751 )
1752 .await;
1753 let response = match response {
1754 Ok(response) => response,
1755 Err(error) => {
1756 check.status = blob_transport_status(error);
1757 return check;
1758 }
1759 };
1760 check.received_bytes = response.body.len() as u64;
1761 match response.status {
1762 206 => {
1763 let content_range = response
1764 .one_header("content-range")
1765 .and_then(parse_content_range);
1766 if content_range != Some((0, requested - 1, descriptor.size))
1767 || response.body.len() as u64 != requested
1768 {
1769 check.status = BlobRangeStatus::Malformed;
1770 } else {
1771 check.status = BlobRangeStatus::Supported;
1772 check.total_bytes = Some(descriptor.size);
1773 if requested == descriptor.size
1774 && sha256_digest(&response.body).ok().as_ref() != Some(&descriptor.digest)
1775 {
1776 check.status = BlobRangeStatus::DigestMismatch;
1777 }
1778 }
1779 }
1780 200 => {
1781 check.status = if requested == descriptor.size {
1782 if response.body.len() as u64 != descriptor.size {
1783 BlobRangeStatus::Malformed
1784 } else if sha256_digest(&response.body).ok().as_ref() != Some(&descriptor.digest) {
1785 BlobRangeStatus::DigestMismatch
1786 } else {
1787 BlobRangeStatus::Ignored
1788 }
1789 } else if response.body.len() as u64 == requested && response.body_truncated {
1790 BlobRangeStatus::Ignored
1791 } else {
1792 BlobRangeStatus::Malformed
1793 };
1794 }
1795 416 => check.status = BlobRangeStatus::Unsatisfiable,
1796 status => check.status = blob_status_for_http(status),
1797 }
1798 check
1799}
1800
1801fn parse_content_range(value: &str) -> Option<(u64, u64, u64)> {
1802 let range = value.strip_prefix("bytes ")?;
1803 let (bounds, total) = range.split_once('/')?;
1804 let (start, end) = bounds.split_once('-')?;
1805 let start = start.parse().ok()?;
1806 let end = end.parse().ok()?;
1807 let total = total.parse().ok()?;
1808 (start <= end && end < total).then_some((start, end, total))
1809}
1810
1811async fn check_mirror<T: RegistryTransport>(
1812 state: &mut DiagnosticState<'_, T>,
1813 order: usize,
1814 mirror: &RegistryEndpoint,
1815 image: &ImageReference,
1816 expected: &OciDigest,
1817 probe_manifest: Option<&Descriptor>,
1818 probe_blob: Option<&Descriptor>,
1819) -> MirrorCheck {
1820 let started = tokio::time::Instant::now();
1821 let api_probe = probe_api(state, mirror).await;
1822 let authorization = match api_probe.check.status {
1823 ApiCheckStatus::Available => None,
1824 ApiCheckStatus::BearerChallenge => {
1825 let Some(challenge) = api_probe.challenge else {
1826 return mirror_failure(
1827 order,
1828 mirror,
1829 MirrorCheckStatus::InvalidResponse,
1830 elapsed_micros(started),
1831 );
1832 };
1833 if !valid_challenge_scope(challenge.scope.as_deref(), image.repository().as_str()) {
1834 return mirror_failure(
1835 order,
1836 mirror,
1837 MirrorCheckStatus::AuthenticationRequired,
1838 elapsed_micros(started),
1839 );
1840 }
1841 match obtain_anonymous_token(state, mirror, &challenge, image).await {
1842 Ok(token) => Some(token),
1843 Err(status) => {
1844 return mirror_failure(
1845 order,
1846 mirror,
1847 mirror_status_from_api(status),
1848 elapsed_micros(started),
1849 );
1850 }
1851 }
1852 }
1853 ApiCheckStatus::AuthenticationRequired | ApiCheckStatus::InvalidChallenge => {
1854 return mirror_failure(
1855 order,
1856 mirror,
1857 MirrorCheckStatus::AuthenticationRequired,
1858 elapsed_micros(started),
1859 );
1860 }
1861 ApiCheckStatus::AccessDenied => {
1862 return mirror_failure(
1863 order,
1864 mirror,
1865 MirrorCheckStatus::AccessDenied,
1866 elapsed_micros(started),
1867 );
1868 }
1869 ApiCheckStatus::RateLimited => {
1870 return mirror_failure(
1871 order,
1872 mirror,
1873 MirrorCheckStatus::RateLimited,
1874 elapsed_micros(started),
1875 );
1876 }
1877 ApiCheckStatus::NotFound => {
1878 return mirror_failure(
1879 order,
1880 mirror,
1881 MirrorCheckStatus::NotFound,
1882 elapsed_micros(started),
1883 );
1884 }
1885 ApiCheckStatus::ServerError => {
1886 return mirror_failure(
1887 order,
1888 mirror,
1889 MirrorCheckStatus::ServerError,
1890 elapsed_micros(started),
1891 );
1892 }
1893 ApiCheckStatus::RedirectRejected => {
1894 return mirror_failure(
1895 order,
1896 mirror,
1897 MirrorCheckStatus::RedirectRejected,
1898 elapsed_micros(started),
1899 );
1900 }
1901 ApiCheckStatus::TimedOut => {
1902 return mirror_failure(
1903 order,
1904 mirror,
1905 MirrorCheckStatus::TimedOut,
1906 elapsed_micros(started),
1907 );
1908 }
1909 ApiCheckStatus::BodyTooLarge => {
1910 return mirror_failure(
1911 order,
1912 mirror,
1913 MirrorCheckStatus::BodyTooLarge,
1914 elapsed_micros(started),
1915 );
1916 }
1917 ApiCheckStatus::RequestLimit => {
1918 return mirror_failure(
1919 order,
1920 mirror,
1921 MirrorCheckStatus::RequestLimit,
1922 elapsed_micros(started),
1923 );
1924 }
1925 ApiCheckStatus::Unreachable => {
1926 return mirror_failure(
1927 order,
1928 mirror,
1929 MirrorCheckStatus::Unreachable,
1930 elapsed_micros(started),
1931 );
1932 }
1933 ApiCheckStatus::UnexpectedResponse | ApiCheckStatus::NotTested => {
1934 return mirror_failure(
1935 order,
1936 mirror,
1937 MirrorCheckStatus::InvalidResponse,
1938 elapsed_micros(started),
1939 );
1940 }
1941 };
1942 let path = manifest_path(image, expected.as_str());
1943 let result = fetch_manifest(state, mirror, &path, authorization.as_ref(), None).await;
1944 match result {
1945 Ok(manifest) => {
1946 let equivalent = &manifest.digest == expected;
1947 if !equivalent {
1948 return MirrorCheck {
1949 order,
1950 origin: mirror.report_origin().to_owned(),
1951 status: MirrorCheckStatus::Diverged,
1952 digest: Some(manifest.digest),
1953 blob_range: BlobRangeCheck::default(),
1954 elapsed_micros: elapsed_micros(started),
1955 recommended_rank: None,
1956 };
1957 }
1958 if let Some(descriptor) = probe_manifest {
1959 let child_path = manifest_path(image, descriptor.digest.as_str());
1960 match fetch_manifest(
1961 state,
1962 mirror,
1963 &child_path,
1964 authorization.as_ref(),
1965 Some((&descriptor.digest, descriptor.size)),
1966 )
1967 .await
1968 {
1969 Ok(child) if child.kind == ManifestKind::Image => {}
1970 Ok(_) => {
1971 return mirror_failure(
1972 order,
1973 mirror,
1974 MirrorCheckStatus::InvalidResponse,
1975 elapsed_micros(started),
1976 );
1977 }
1978 Err(
1979 ManifestCheckStatus::DigestMismatch | ManifestCheckStatus::SizeMismatch,
1980 ) => {
1981 return mirror_failure(
1982 order,
1983 mirror,
1984 MirrorCheckStatus::Diverged,
1985 elapsed_micros(started),
1986 );
1987 }
1988 Err(status) => {
1989 return mirror_failure(
1990 order,
1991 mirror,
1992 mirror_status_from_manifest(status),
1993 elapsed_micros(started),
1994 );
1995 }
1996 }
1997 }
1998 let blob_range = {
1999 match probe_blob {
2000 Some(blob) => {
2001 check_blob_range(state, mirror, image, authorization.as_ref(), blob.clone())
2002 .await
2003 }
2004 None => BlobRangeCheck::default(),
2005 }
2006 };
2007 MirrorCheck {
2008 order,
2009 origin: mirror.report_origin().to_owned(),
2010 status: MirrorCheckStatus::Equivalent,
2011 digest: Some(manifest.digest),
2012 blob_range,
2013 elapsed_micros: elapsed_micros(started),
2014 recommended_rank: None,
2015 }
2016 }
2017 Err(status) => MirrorCheck {
2018 order,
2019 origin: mirror.report_origin().to_owned(),
2020 status: mirror_status_from_manifest(status),
2021 digest: None,
2022 blob_range: BlobRangeCheck::default(),
2023 elapsed_micros: elapsed_micros(started),
2024 recommended_rank: None,
2025 },
2026 }
2027}
2028
2029fn mirror_failure(
2030 order: usize,
2031 mirror: &RegistryEndpoint,
2032 status: MirrorCheckStatus,
2033 elapsed_micros: u64,
2034) -> MirrorCheck {
2035 MirrorCheck {
2036 order,
2037 origin: mirror.report_origin().to_owned(),
2038 status,
2039 digest: None,
2040 blob_range: BlobRangeCheck::default(),
2041 elapsed_micros,
2042 recommended_rank: None,
2043 }
2044}
2045
2046fn elapsed_micros(started: tokio::time::Instant) -> u64 {
2047 u64::try_from(started.elapsed().as_micros()).unwrap_or(u64::MAX)
2048}
2049
2050fn rank_mirror_checks(report: &mut RegistryDiagnosticReport) {
2051 let require_content = report.image.is_some();
2052 let mut ranked = report
2053 .mirrors
2054 .iter()
2055 .filter(|mirror| {
2056 if require_content {
2057 mirror.status == MirrorCheckStatus::Equivalent
2058 && matches!(
2059 mirror.blob_range.status,
2060 BlobRangeStatus::Supported | BlobRangeStatus::Ignored
2061 )
2062 && mirror.blob_range.received_bytes > 0
2063 } else {
2064 mirror.status == MirrorCheckStatus::Available
2065 }
2066 })
2067 .map(|mirror| (mirror.elapsed_micros, mirror.order))
2068 .collect::<Vec<_>>();
2069 ranked.sort();
2070 report.recommended_mirror_order = ranked.iter().map(|(_, order)| *order).collect();
2071 for (rank, (_, order)) in ranked.into_iter().enumerate() {
2072 if let Some(mirror) = report
2073 .mirrors
2074 .iter_mut()
2075 .find(|mirror| mirror.order == order)
2076 {
2077 mirror.recommended_rank = Some(rank + 1);
2078 }
2079 }
2080}
2081
2082fn manifest_path(image: &ImageReference, selector: &str) -> String {
2083 format!("/v2/{}/manifests/{selector}", image.repository().as_str())
2084}
2085
2086fn blob_path(image: &ImageReference, digest: &str) -> String {
2087 format!("/v2/{}/blobs/{digest}", image.repository().as_str())
2088}
2089
2090fn sha256_digest(body: &[u8]) -> Result<OciDigest, RegistryProtocolError> {
2091 OciDigest::parse(&format!("sha256:{:x}", Sha256::digest(body)))
2092 .map_err(|_| RegistryProtocolError::InvalidPath)
2093}
2094
2095fn redacted_origin(url: &reqwest::Url) -> String {
2096 let host = url.host_str().unwrap_or("invalid");
2097 let authority = registry_authority(host, url.port()).unwrap_or_else(|_| "invalid".to_owned());
2098 format!("{}://{authority}", url.scheme())
2099}
2100
2101fn api_transport_status(error: RegistryTransportError) -> ApiCheckStatus {
2102 match error {
2103 RegistryTransportError::Timeout => ApiCheckStatus::TimedOut,
2104 RegistryTransportError::BodyTooLarge => ApiCheckStatus::BodyTooLarge,
2105 RegistryTransportError::RedirectRejected => ApiCheckStatus::RedirectRejected,
2106 RegistryTransportError::RequestLimit => ApiCheckStatus::RequestLimit,
2107 RegistryTransportError::Failed => ApiCheckStatus::Unreachable,
2108 }
2109}
2110
2111fn manifest_transport_status(error: RegistryTransportError) -> ManifestCheckStatus {
2112 match error {
2113 RegistryTransportError::Timeout => ManifestCheckStatus::TimedOut,
2114 RegistryTransportError::BodyTooLarge => ManifestCheckStatus::BodyTooLarge,
2115 RegistryTransportError::RedirectRejected => ManifestCheckStatus::RedirectRejected,
2116 RegistryTransportError::RequestLimit => ManifestCheckStatus::RequestLimit,
2117 RegistryTransportError::Failed => ManifestCheckStatus::Unreachable,
2118 }
2119}
2120
2121fn blob_transport_status(error: RegistryTransportError) -> BlobRangeStatus {
2122 match error {
2123 RegistryTransportError::Timeout => BlobRangeStatus::TimedOut,
2124 RegistryTransportError::BodyTooLarge => BlobRangeStatus::BodyTooLarge,
2125 RegistryTransportError::RedirectRejected => BlobRangeStatus::RedirectRejected,
2126 RegistryTransportError::RequestLimit => BlobRangeStatus::RequestLimit,
2127 RegistryTransportError::Failed => BlobRangeStatus::Unreachable,
2128 }
2129}
2130
2131fn api_status_for_http(status: u16) -> ApiCheckStatus {
2132 match status {
2133 401 => ApiCheckStatus::AuthenticationRequired,
2134 403 => ApiCheckStatus::AccessDenied,
2135 404 => ApiCheckStatus::NotFound,
2136 429 => ApiCheckStatus::RateLimited,
2137 500..=599 => ApiCheckStatus::ServerError,
2138 _ => ApiCheckStatus::UnexpectedResponse,
2139 }
2140}
2141
2142fn manifest_status_for_http(status: u16) -> ManifestCheckStatus {
2143 match status {
2144 401 => ManifestCheckStatus::AuthenticationRequired,
2145 403 => ManifestCheckStatus::AccessDenied,
2146 404 => ManifestCheckStatus::NotFound,
2147 429 => ManifestCheckStatus::RateLimited,
2148 500..=599 => ManifestCheckStatus::ServerError,
2149 _ => ManifestCheckStatus::InvalidManifest,
2150 }
2151}
2152
2153fn blob_status_for_http(status: u16) -> BlobRangeStatus {
2154 match status {
2155 401 => BlobRangeStatus::AuthenticationRequired,
2156 403 => BlobRangeStatus::AccessDenied,
2157 404 => BlobRangeStatus::NotFound,
2158 429 => BlobRangeStatus::RateLimited,
2159 500..=599 => BlobRangeStatus::ServerError,
2160 _ => BlobRangeStatus::Malformed,
2161 }
2162}
2163
2164fn mirror_status_from_api(status: ApiCheckStatus) -> MirrorCheckStatus {
2165 match status {
2166 ApiCheckStatus::Available | ApiCheckStatus::BearerChallenge => MirrorCheckStatus::Available,
2167 ApiCheckStatus::AuthenticationRequired | ApiCheckStatus::InvalidChallenge => {
2168 MirrorCheckStatus::AuthenticationRequired
2169 }
2170 ApiCheckStatus::AccessDenied => MirrorCheckStatus::AccessDenied,
2171 ApiCheckStatus::RateLimited => MirrorCheckStatus::RateLimited,
2172 ApiCheckStatus::NotFound => MirrorCheckStatus::NotFound,
2173 ApiCheckStatus::ServerError => MirrorCheckStatus::ServerError,
2174 ApiCheckStatus::RedirectRejected => MirrorCheckStatus::RedirectRejected,
2175 ApiCheckStatus::Unreachable => MirrorCheckStatus::Unreachable,
2176 ApiCheckStatus::TimedOut => MirrorCheckStatus::TimedOut,
2177 ApiCheckStatus::BodyTooLarge => MirrorCheckStatus::BodyTooLarge,
2178 ApiCheckStatus::RequestLimit => MirrorCheckStatus::RequestLimit,
2179 ApiCheckStatus::UnexpectedResponse | ApiCheckStatus::NotTested => {
2180 MirrorCheckStatus::InvalidResponse
2181 }
2182 }
2183}
2184
2185fn mirror_status_from_manifest(status: ManifestCheckStatus) -> MirrorCheckStatus {
2186 match status {
2187 ManifestCheckStatus::AuthenticationRequired => MirrorCheckStatus::AuthenticationRequired,
2188 ManifestCheckStatus::AccessDenied => MirrorCheckStatus::AccessDenied,
2189 ManifestCheckStatus::RateLimited => MirrorCheckStatus::RateLimited,
2190 ManifestCheckStatus::NotFound => MirrorCheckStatus::NotFound,
2191 ManifestCheckStatus::ServerError => MirrorCheckStatus::ServerError,
2192 ManifestCheckStatus::RedirectRejected => MirrorCheckStatus::RedirectRejected,
2193 ManifestCheckStatus::Unreachable => MirrorCheckStatus::Unreachable,
2194 ManifestCheckStatus::TimedOut => MirrorCheckStatus::TimedOut,
2195 ManifestCheckStatus::BodyTooLarge => MirrorCheckStatus::BodyTooLarge,
2196 ManifestCheckStatus::RequestLimit => MirrorCheckStatus::RequestLimit,
2197 _ => MirrorCheckStatus::InvalidResponse,
2198 }
2199}
2200
2201fn aggregate_status(report: &RegistryDiagnosticReport) -> RegistryDiagnosticStatus {
2202 match report.api.status {
2203 ApiCheckStatus::AuthenticationRequired | ApiCheckStatus::InvalidChallenge => {
2204 return RegistryDiagnosticStatus::AuthenticationRequired;
2205 }
2206 ApiCheckStatus::AccessDenied => return RegistryDiagnosticStatus::AccessDenied,
2207 ApiCheckStatus::RateLimited => return RegistryDiagnosticStatus::RateLimited,
2208 ApiCheckStatus::NotFound => return RegistryDiagnosticStatus::NotFound,
2209 ApiCheckStatus::Unreachable => return RegistryDiagnosticStatus::Unreachable,
2210 ApiCheckStatus::TimedOut => return RegistryDiagnosticStatus::TimedOut,
2211 ApiCheckStatus::RequestLimit => return RegistryDiagnosticStatus::LimitExceeded,
2212 ApiCheckStatus::BodyTooLarge
2213 | ApiCheckStatus::ServerError
2214 | ApiCheckStatus::RedirectRejected
2215 | ApiCheckStatus::UnexpectedResponse
2216 | ApiCheckStatus::NotTested => return RegistryDiagnosticStatus::ProtocolError,
2217 ApiCheckStatus::Available | ApiCheckStatus::BearerChallenge => {}
2218 }
2219 match report.manifest.status {
2220 ManifestCheckStatus::DigestMismatch | ManifestCheckStatus::SizeMismatch => {
2221 return RegistryDiagnosticStatus::Corrupt;
2222 }
2223 ManifestCheckStatus::AuthenticationRequired => {
2224 return RegistryDiagnosticStatus::AuthenticationRequired;
2225 }
2226 ManifestCheckStatus::AccessDenied => return RegistryDiagnosticStatus::AccessDenied,
2227 ManifestCheckStatus::RateLimited => return RegistryDiagnosticStatus::RateLimited,
2228 ManifestCheckStatus::NotFound => return RegistryDiagnosticStatus::NotFound,
2229 ManifestCheckStatus::Unreachable => return RegistryDiagnosticStatus::Unreachable,
2230 ManifestCheckStatus::TimedOut => return RegistryDiagnosticStatus::TimedOut,
2231 ManifestCheckStatus::RequestLimit => return RegistryDiagnosticStatus::LimitExceeded,
2232 ManifestCheckStatus::InvalidMediaType
2233 | ManifestCheckStatus::InvalidManifest
2234 | ManifestCheckStatus::BodyTooLarge
2235 | ManifestCheckStatus::ServerError
2236 | ManifestCheckStatus::RedirectRejected => {
2237 return RegistryDiagnosticStatus::ProtocolError;
2238 }
2239 ManifestCheckStatus::PlatformNotFound => return RegistryDiagnosticStatus::Unsupported,
2240 ManifestCheckStatus::Verified
2241 | ManifestCheckStatus::NotRequested
2242 | ManifestCheckStatus::NotTested => {}
2243 }
2244 if report
2245 .mirrors
2246 .iter()
2247 .any(|mirror| mirror.status == MirrorCheckStatus::Diverged)
2248 {
2249 return RegistryDiagnosticStatus::Corrupt;
2250 }
2251 if report.blob_range.status == BlobRangeStatus::DigestMismatch {
2252 return RegistryDiagnosticStatus::Corrupt;
2253 }
2254 if report
2255 .mirrors
2256 .iter()
2257 .any(|mirror| mirror.blob_range.status == BlobRangeStatus::DigestMismatch)
2258 {
2259 return RegistryDiagnosticStatus::Corrupt;
2260 }
2261 if !matches!(
2262 report.blob_range.status,
2263 BlobRangeStatus::Supported
2264 | BlobRangeStatus::Ignored
2265 | BlobRangeStatus::NotAvailable
2266 | BlobRangeStatus::NotTested
2267 ) || report.mirrors.iter().any(|mirror| {
2268 !matches!(
2269 mirror.status,
2270 MirrorCheckStatus::Available | MirrorCheckStatus::Equivalent
2271 ) || !matches!(
2272 mirror.blob_range.status,
2273 BlobRangeStatus::Supported
2274 | BlobRangeStatus::Ignored
2275 | BlobRangeStatus::NotAvailable
2276 | BlobRangeStatus::NotTested
2277 )
2278 }) {
2279 RegistryDiagnosticStatus::Degraded
2280 } else {
2281 RegistryDiagnosticStatus::Healthy
2282 }
2283}
2284
2285#[cfg(test)]
2286mod tests {
2287 use std::collections::VecDeque;
2288 use std::sync::{Arc, Mutex};
2289
2290 use super::*;
2291
2292 #[derive(Clone, Debug, PartialEq, Eq)]
2293 struct SeenRequest {
2294 method: RegistryMethod,
2295 origin: String,
2296 path: String,
2297 query: Option<String>,
2298 accept: Option<&'static str>,
2299 range: Option<(u64, u64)>,
2300 authorized: bool,
2301 max_body_bytes: usize,
2302 }
2303
2304 #[derive(Clone)]
2305 struct MockTransport {
2306 steps: Arc<Mutex<VecDeque<MockStep>>>,
2307 seen: Arc<Mutex<Vec<SeenRequest>>>,
2308 }
2309
2310 enum MockStep {
2311 Response(RegistryResponse),
2312 Error(RegistryTransportError),
2313 }
2314
2315 impl MockTransport {
2316 fn new(steps: impl IntoIterator<Item = MockStep>) -> Self {
2317 Self {
2318 steps: Arc::new(Mutex::new(steps.into_iter().collect())),
2319 seen: Arc::new(Mutex::new(Vec::new())),
2320 }
2321 }
2322
2323 fn seen(&self) -> Vec<SeenRequest> {
2324 self.seen.lock().unwrap().clone()
2325 }
2326 }
2327
2328 impl RegistryTransport for MockTransport {
2329 fn execute(&self, request: RegistryRequest) -> RegistryTransportFuture<'_> {
2330 self.seen.lock().unwrap().push(SeenRequest {
2331 method: request.method,
2332 origin: redacted_origin(&request.url),
2333 path: request.url.path().to_owned(),
2334 query: request.url.query().map(str::to_owned),
2335 accept: request.accept,
2336 range: request.range,
2337 authorized: request.has_authorization(),
2338 max_body_bytes: request.max_body_bytes,
2339 });
2340 let result = self
2341 .steps
2342 .lock()
2343 .unwrap()
2344 .pop_front()
2345 .unwrap_or(MockStep::Error(RegistryTransportError::Failed));
2346 Box::pin(async move {
2347 match result {
2348 MockStep::Response(response) => Ok(response),
2349 MockStep::Error(error) => Err(error),
2350 }
2351 })
2352 }
2353 }
2354
2355 fn endpoint(host: &str) -> RegistryEndpoint {
2356 RegistryEndpoint::parse_https(&format!("https://{host}")).unwrap()
2357 }
2358
2359 fn image(value: &str) -> ImageReference {
2360 ImageReference::parse(value).unwrap()
2361 }
2362
2363 fn response(status: u16) -> MockStep {
2364 MockStep::Response(RegistryResponse::new(status))
2365 }
2366
2367 fn manifest_response(body: Vec<u8>) -> MockStep {
2368 let digest = sha256_digest(&body).unwrap();
2369 MockStep::Response(
2370 RegistryResponse::new(200)
2371 .header("content-type", OCI_MANIFEST)
2372 .header("docker-content-digest", digest.as_str())
2373 .body(body),
2374 )
2375 }
2376
2377 fn index_response(body: Vec<u8>) -> MockStep {
2378 let digest = sha256_digest(&body).unwrap();
2379 MockStep::Response(
2380 RegistryResponse::new(200)
2381 .header("content-type", OCI_INDEX)
2382 .header("docker-content-digest", digest.as_str())
2383 .body(body),
2384 )
2385 }
2386
2387 fn image_manifest(layer: &[u8]) -> (Vec<u8>, OciDigest) {
2388 let digest = sha256_digest(layer).unwrap();
2389 let body = serde_json::to_vec(&serde_json::json!({
2390 "schemaVersion": 2,
2391 "mediaType": OCI_MANIFEST,
2392 "config": {
2393 "mediaType": "application/vnd.oci.image.config.v1+json",
2394 "digest": digest,
2395 "size": layer.len()
2396 },
2397 "layers": [{
2398 "mediaType": "application/vnd.oci.image.layer.v1.tar+gzip",
2399 "digest": digest,
2400 "size": layer.len()
2401 }]
2402 }))
2403 .unwrap();
2404 (body, digest)
2405 }
2406
2407 fn options(image: Option<ImageReference>) -> RegistryDiagnosticOptions {
2408 let mut options = RegistryDiagnosticOptions::new(endpoint("registry.example"));
2409 options.image = image;
2410 options
2411 }
2412
2413 #[test]
2414 fn endpoints_are_https_only_and_redacted() {
2415 assert_eq!(
2416 RegistryEndpoint::parse_https("http://registry.example"),
2417 Err(RegistryProtocolError::InsecureEndpoint)
2418 );
2419 assert_eq!(
2420 RegistryEndpoint::parse_https("https://alice:secret@registry.example"),
2421 Err(RegistryProtocolError::InvalidEndpoint)
2422 );
2423 for invalid in [
2424 "https://registry.example?token=secret",
2425 "https://registry.example#fragment",
2426 "https://registry.example/%2e%2e/private",
2427 "https://registry.example/a/../private",
2428 ] {
2429 assert!(RegistryEndpoint::parse_https(invalid).is_err(), "{invalid}");
2430 }
2431 assert!(RegistryEndpoint::parse_loopback_http("http://127.0.0.1:5000").is_ok());
2432 assert!(RegistryEndpoint::parse_loopback_http("http://[::1]:5000").is_ok());
2433 assert!(RegistryEndpoint::parse_loopback_http("http://example.test:5000").is_err());
2434 assert!(endpoint("registry.example").is_https());
2435
2436 let endpoint = RegistryEndpoint::parse_https("https://registry.example/prefix").unwrap();
2437 assert_eq!(endpoint.report_origin(), "https://registry.example");
2438 let debug = format!("{endpoint:?}");
2439 assert!(!debug.contains("/prefix"));
2440 assert!(debug.contains("has_path_prefix"));
2441
2442 let docker =
2443 RegistryEndpoint::for_registry(RegistryName::parse("docker.io").unwrap()).unwrap();
2444 assert_eq!(docker.registry().as_str(), "docker.io");
2445 assert_eq!(docker.report_origin(), "https://registry-1.docker.io");
2446
2447 let oversized = format!(
2448 "https://example.test/{}",
2449 "a".repeat(HARD_MAX_ENDPOINT_BYTES)
2450 );
2451 assert_eq!(
2452 RegistryEndpoint::parse_https(&oversized),
2453 Err(RegistryProtocolError::InvalidEndpoint)
2454 );
2455 }
2456
2457 #[tokio::test]
2458 async fn v2_statuses_and_bearer_challenges_are_typed() {
2459 for (status, expected) in [
2460 (200, ApiCheckStatus::Available),
2461 (403, ApiCheckStatus::AccessDenied),
2462 (404, ApiCheckStatus::NotFound),
2463 (429, ApiCheckStatus::RateLimited),
2464 (503, ApiCheckStatus::ServerError),
2465 ] {
2466 let transport = MockTransport::new([response(status)]);
2467 let report = diagnose_registry(&transport, options(None)).await.unwrap();
2468 assert_eq!(report.api.status, expected);
2469 assert_eq!(report.api.http_status, Some(status));
2470 assert_eq!(report.request_count, 1);
2471 }
2472
2473 let challenge = RegistryResponse::new(401).header(
2474 "www-authenticate",
2475 r#"Bearer realm="https://registry.example/token",service="registry.example""#,
2476 );
2477 let transport = MockTransport::new([MockStep::Response(challenge)]);
2478 let report = diagnose_registry(&transport, options(None)).await.unwrap();
2479 assert_eq!(report.status, RegistryDiagnosticStatus::Healthy);
2480 assert_eq!(report.api.status, ApiCheckStatus::BearerChallenge);
2481 assert!(report.api.bearer_challenge);
2482 assert!(report.api.challenge_service_present);
2483
2484 let transport = MockTransport::new([MockStep::Response(
2485 RegistryResponse::new(401).header("www-authenticate", "Basic realm=\"private\""),
2486 )]);
2487 let report = diagnose_registry(&transport, options(None)).await.unwrap();
2488 assert_eq!(report.api.status, ApiCheckStatus::InvalidChallenge);
2489 }
2490
2491 #[tokio::test]
2492 async fn api_only_diagnostic_probes_each_mirror_in_configured_order() {
2493 let transport = MockTransport::new([response(200), response(200), response(200)]);
2494 let report = diagnose_registry(
2495 &transport,
2496 options(None).with_mirrors(vec![
2497 endpoint("first.example/prefix"),
2498 endpoint("second.example"),
2499 ]),
2500 )
2501 .await
2502 .unwrap();
2503
2504 assert_eq!(report.status, RegistryDiagnosticStatus::Healthy);
2505 assert_eq!(report.request_count, 3);
2506 assert_eq!(
2507 report
2508 .mirrors
2509 .iter()
2510 .map(|mirror| (mirror.order, mirror.status))
2511 .collect::<Vec<_>>(),
2512 [
2513 (0, MirrorCheckStatus::Available),
2514 (1, MirrorCheckStatus::Available),
2515 ]
2516 );
2517 let mut recommended = report.recommended_mirror_order.clone();
2518 recommended.sort_unstable();
2519 assert_eq!(recommended, [0, 1]);
2520 for (rank, order) in report.recommended_mirror_order.iter().enumerate() {
2521 assert_eq!(report.mirrors[*order].recommended_rank, Some(rank + 1));
2522 }
2523 let seen = transport.seen();
2524 assert_eq!(
2525 seen.iter()
2526 .map(|request| (request.origin.as_str(), request.path.as_str()))
2527 .collect::<Vec<_>>(),
2528 [
2529 ("https://registry.example", "/v2/"),
2530 ("https://first.example", "/prefix/v2/"),
2531 ("https://second.example", "/v2/"),
2532 ]
2533 );
2534 }
2535
2536 #[test]
2537 fn anonymous_token_realms_are_limited_to_audited_registry_pairs() {
2538 let docker =
2539 RegistryEndpoint::for_registry(RegistryName::parse("docker.io").unwrap()).unwrap();
2540 assert!(docker.allows_anonymous_token_realm(
2541 &reqwest::Url::parse("https://auth.docker.io/token").unwrap()
2542 ));
2543 let daocloud = endpoint("docker.m.daocloud.io");
2544 assert!(daocloud.allows_anonymous_token_realm(
2545 &reqwest::Url::parse("https://m.daocloud.io/token").unwrap()
2546 ));
2547 assert!(!daocloud.allows_anonymous_token_realm(
2548 &reqwest::Url::parse("https://auth.example/token").unwrap()
2549 ));
2550 assert!(!daocloud.allows_anonymous_token_realm(
2551 &reqwest::Url::parse("https://m.daocloud.io:444/token").unwrap()
2552 ));
2553 }
2554
2555 #[tokio::test]
2556 async fn anonymous_bearer_token_is_scoped_and_never_serialized() {
2557 let layer = b"tiny layer";
2558 let (manifest, _) = image_manifest(layer);
2559 let challenge = RegistryResponse::new(401).header(
2560 "www-authenticate",
2561 r#"Bearer realm="https://registry.example/token?account=anonymous",service="registry.example",scope="repository:team/app:pull""#,
2562 );
2563 let transport = MockTransport::new([
2564 MockStep::Response(challenge),
2565 MockStep::Response(RegistryResponse::new(200).body(br#"{"token":"secret-token"}"#)),
2566 manifest_response(manifest),
2567 MockStep::Response(
2568 RegistryResponse::new(206)
2569 .header("content-range", "bytes 0-9/10")
2570 .body(layer.to_vec()),
2571 ),
2572 ]);
2573 let report = diagnose_registry(
2574 &transport,
2575 options(Some(image("registry.example/team/app:latest"))),
2576 )
2577 .await
2578 .unwrap();
2579 assert_eq!(report.status, RegistryDiagnosticStatus::Healthy);
2580 assert_eq!(report.api.challenge_scope_matches, Some(true));
2581 assert_eq!(report.manifest.status, ManifestCheckStatus::Verified);
2582 assert_eq!(report.blob_range.status, BlobRangeStatus::Supported);
2583
2584 let seen = transport.seen();
2585 assert_eq!(seen.len(), 4);
2586 assert_eq!(seen[0].path, "/v2/");
2587 assert_eq!(seen[1].origin, "https://registry.example");
2588 assert!(!seen[1].authorized);
2589 assert!(seen[1]
2590 .query
2591 .as_deref()
2592 .unwrap()
2593 .contains("scope=repository%3Ateam%2Fapp%3Apull"));
2594 assert!(seen[2].authorized);
2595 assert!(seen[3].authorized);
2596 let json = serde_json::to_string(&report).unwrap();
2597 assert!(!json.contains("secret-token"));
2598 assert!(!json.contains("account=anonymous"));
2599 }
2600
2601 #[tokio::test]
2602 async fn image_manifest_and_full_small_blob_are_verified_exactly() {
2603 let layer = b"0123456789abcdef";
2604 let (manifest, layer_digest) = image_manifest(layer);
2605 let manifest_digest = sha256_digest(&manifest).unwrap();
2606 let reference = image(&format!("registry.example/team/app@{manifest_digest}"));
2607 let transport = MockTransport::new([
2608 response(200),
2609 manifest_response(manifest),
2610 MockStep::Response(
2611 RegistryResponse::new(206)
2612 .header("content-range", "bytes 0-15/16")
2613 .body(layer.to_vec()),
2614 ),
2615 ]);
2616 let report = diagnose_registry(&transport, options(Some(reference)))
2617 .await
2618 .unwrap();
2619 assert_eq!(report.status, RegistryDiagnosticStatus::Healthy);
2620 assert_eq!(report.manifest.digest, Some(manifest_digest));
2621 assert_eq!(report.blob_range.digest, Some(layer_digest));
2622 assert_eq!(report.blob_range.status, BlobRangeStatus::Supported);
2623 let seen = transport.seen();
2624 assert_eq!(seen[1].accept, Some(ACCEPT_MANIFESTS));
2625 assert_eq!(seen[2].range, Some((0, 15)));
2626 assert_eq!(seen[2].max_body_bytes, DEFAULT_BLOB_SAMPLE_BYTES);
2627 }
2628
2629 #[tokio::test]
2630 async fn index_selects_exactly_one_platform_child() {
2631 let layer = b"arm64 layer";
2632 let (child, _) = image_manifest(layer);
2633 let child_digest = sha256_digest(&child).unwrap();
2634 let index = serde_json::to_vec(&serde_json::json!({
2635 "schemaVersion": 2,
2636 "mediaType": OCI_INDEX,
2637 "manifests": [
2638 {
2639 "mediaType": OCI_MANIFEST,
2640 "digest": sha256_digest(b"unused").unwrap(),
2641 "size": 6,
2642 "platform": {"os": "linux", "architecture": "amd64"}
2643 },
2644 {
2645 "mediaType": OCI_MANIFEST,
2646 "digest": child_digest,
2647 "size": child.len(),
2648 "platform": {"os": "linux", "architecture": "arm64", "variant": "v8"}
2649 },
2650 {
2651 "mediaType": OCI_MANIFEST,
2652 "digest": sha256_digest(b"windows").unwrap(),
2653 "size": 7,
2654 "platform": {"os": "windows", "architecture": "amd64", "os.version": "10.0.20348.0"}
2655 }
2656 ]
2657 }))
2658 .unwrap();
2659 let transport = MockTransport::new([
2660 response(200),
2661 index_response(index),
2662 manifest_response(child),
2663 MockStep::Response(
2664 RegistryResponse::new(206)
2665 .header("content-range", "bytes 0-10/11")
2666 .body(layer.to_vec()),
2667 ),
2668 ]);
2669 let opts = options(Some(image("registry.example/team/app:v1")))
2670 .with_platform(OciPlatform::parse("linux/arm64/v8").unwrap());
2671 let report = diagnose_registry(&transport, opts).await.unwrap();
2672 assert_eq!(report.manifest.status, ManifestCheckStatus::Verified);
2673 assert_eq!(report.manifest.kind, Some(ManifestKind::Index));
2674 assert_eq!(report.manifest.child_digest, Some(child_digest.clone()));
2675 assert_eq!(transport.seen().len(), 4);
2676 assert!(transport.seen()[2].path.ends_with(child_digest.as_str()));
2677 }
2678
2679 #[tokio::test]
2680 async fn index_mirror_verifies_the_selected_child_before_sampling_its_layer() {
2681 let layer = b"amd64 layer";
2682 let (child, _) = image_manifest(layer);
2683 let child_digest = sha256_digest(&child).unwrap();
2684 let index = serde_json::to_vec(&serde_json::json!({
2685 "schemaVersion": 2,
2686 "mediaType": OCI_INDEX,
2687 "manifests": [{
2688 "mediaType": OCI_MANIFEST,
2689 "digest": child_digest,
2690 "size": child.len(),
2691 "platform": {"os": "linux", "architecture": "amd64"}
2692 }]
2693 }))
2694 .unwrap();
2695 let index_digest = sha256_digest(&index).unwrap();
2696 let range = || {
2697 MockStep::Response(
2698 RegistryResponse::new(206)
2699 .header("content-range", "bytes 0-10/11")
2700 .body(layer.to_vec()),
2701 )
2702 };
2703 let transport = MockTransport::new([
2704 response(200),
2705 index_response(index.clone()),
2706 manifest_response(child.clone()),
2707 range(),
2708 response(200),
2709 index_response(index),
2710 manifest_response(child),
2711 range(),
2712 ]);
2713 let report = diagnose_registry(
2714 &transport,
2715 options(Some(image("registry.example/team/app:latest")))
2716 .with_platform(OciPlatform::parse("linux/amd64").unwrap())
2717 .with_mirrors(vec![endpoint("mirror.example")]),
2718 )
2719 .await
2720 .unwrap();
2721
2722 assert_eq!(report.mirrors[0].status, MirrorCheckStatus::Equivalent);
2723 assert_eq!(
2724 report.mirrors[0].blob_range.status,
2725 BlobRangeStatus::Supported
2726 );
2727 assert_eq!(report.recommended_mirror_order, [0]);
2728 let seen = transport.seen();
2729 assert_eq!(seen.len(), 8);
2730 assert!(seen[5].path.ends_with(index_digest.as_str()));
2731 assert!(seen[6].path.ends_with(child_digest.as_str()));
2732 assert!(seen[7].path.contains("/blobs/"));
2733 }
2734
2735 #[tokio::test]
2736 async fn duplicate_platform_matches_are_rejected_as_ambiguous() {
2737 let child = b"child";
2738 let digest = sha256_digest(child).unwrap();
2739 let descriptor = serde_json::json!({
2740 "mediaType": OCI_MANIFEST,
2741 "digest": digest,
2742 "size": child.len(),
2743 "platform": {"os": "linux", "architecture": "amd64"}
2744 });
2745 let index = serde_json::to_vec(&serde_json::json!({
2746 "schemaVersion": 2,
2747 "mediaType": OCI_INDEX,
2748 "manifests": [descriptor.clone(), descriptor]
2749 }))
2750 .unwrap();
2751 let transport = MockTransport::new([response(200), index_response(index)]);
2752 let opts = options(Some(image("registry.example/team/app:v1")))
2753 .with_platform(OciPlatform::parse("linux/amd64").unwrap());
2754 let report = diagnose_registry(&transport, opts).await.unwrap();
2755 assert_eq!(report.status, RegistryDiagnosticStatus::Unsupported);
2756 assert_eq!(
2757 report.manifest.status,
2758 ManifestCheckStatus::PlatformNotFound
2759 );
2760 assert_eq!(transport.seen().len(), 2);
2761 }
2762
2763 #[tokio::test]
2764 async fn windows_platform_accepts_and_reports_one_valid_os_version() {
2765 let layer = b"windows-layer";
2766 let (child, _) = image_manifest(layer);
2767 let digest = sha256_digest(&child).unwrap();
2768 let index = serde_json::to_vec(&serde_json::json!({
2769 "schemaVersion": 2,
2770 "mediaType": OCI_INDEX,
2771 "manifests": [{
2772 "mediaType": OCI_MANIFEST,
2773 "digest": digest,
2774 "size": child.len(),
2775 "platform": {
2776 "os": "windows",
2777 "architecture": "amd64",
2778 "os.version": "10.0.20348.0"
2779 }
2780 }]
2781 }))
2782 .unwrap();
2783 let transport = MockTransport::new([
2784 response(200),
2785 index_response(index),
2786 manifest_response(child),
2787 MockStep::Response(
2788 RegistryResponse::new(206)
2789 .header("content-range", "bytes 0-12/13")
2790 .body(layer.to_vec()),
2791 ),
2792 ]);
2793 let opts = options(Some(image("registry.example/team/app:v1")))
2794 .with_platform(OciPlatform::parse("windows/amd64").unwrap());
2795 let report = diagnose_registry(&transport, opts).await.unwrap();
2796 assert_eq!(report.status, RegistryDiagnosticStatus::Healthy);
2797 assert_eq!(report.manifest.status, ManifestCheckStatus::Verified);
2798 assert_eq!(
2799 report.manifest.selected_os_version.as_deref(),
2800 Some("10.0.20348.0")
2801 );
2802 assert_eq!(transport.seen().len(), 4);
2803 }
2804
2805 #[tokio::test]
2806 async fn manifest_corruption_size_and_body_limits_are_typed() {
2807 let body = serde_json::to_vec(&serde_json::json!({
2808 "schemaVersion": 2,
2809 "mediaType": OCI_MANIFEST,
2810 "layers": []
2811 }))
2812 .unwrap();
2813 let wrong = OciDigest::parse(&format!("sha256:{}", "0".repeat(64))).unwrap();
2814 let transport = MockTransport::new([
2815 response(200),
2816 MockStep::Response(
2817 RegistryResponse::new(200)
2818 .header("content-type", OCI_MANIFEST)
2819 .header("docker-content-digest", wrong.as_str())
2820 .body(body),
2821 ),
2822 ]);
2823 let report = diagnose_registry(
2824 &transport,
2825 options(Some(image("registry.example/team/app:v1"))),
2826 )
2827 .await
2828 .unwrap();
2829 assert_eq!(report.status, RegistryDiagnosticStatus::Corrupt);
2830 assert_eq!(report.manifest.status, ManifestCheckStatus::DigestMismatch);
2831
2832 let transport = MockTransport::new([
2833 response(200),
2834 MockStep::Error(RegistryTransportError::BodyTooLarge),
2835 ]);
2836 let report = diagnose_registry(
2837 &transport,
2838 options(Some(image("registry.example/team/app:v1"))),
2839 )
2840 .await
2841 .unwrap();
2842 assert_eq!(report.manifest.status, ManifestCheckStatus::BodyTooLarge);
2843
2844 let invalid_body = serde_json::to_vec(&serde_json::json!({
2845 "schemaVersion": 2,
2846 "mediaType": OCI_MANIFEST,
2847 "config": {
2848 "mediaType": "text/plain",
2849 "digest": sha256_digest(b"config").unwrap(),
2850 "size": 6
2851 },
2852 "layers": []
2853 }))
2854 .unwrap();
2855 let transport = MockTransport::new([response(200), manifest_response(invalid_body)]);
2856 let report = diagnose_registry(
2857 &transport,
2858 options(Some(image("registry.example/team/app:v1"))),
2859 )
2860 .await
2861 .unwrap();
2862 assert_eq!(report.manifest.status, ManifestCheckStatus::InvalidManifest);
2863 assert_ne!(report.manifest.status, ManifestCheckStatus::Verified);
2864 }
2865
2866 #[tokio::test]
2867 async fn blob_range_outcomes_are_typed() {
2868 let layer = b"0123456789abcdef";
2869 for (range_response, expected) in [
2870 (
2871 RegistryResponse::new(200).body(layer.to_vec()),
2872 BlobRangeStatus::Ignored,
2873 ),
2874 (
2875 RegistryResponse::new(200).body(b"xxxxxxxxxxxxxxxx".to_vec()),
2876 BlobRangeStatus::DigestMismatch,
2877 ),
2878 (RegistryResponse::new(416), BlobRangeStatus::Unsatisfiable),
2879 (
2880 RegistryResponse::new(206)
2881 .header("content-range", "bytes 1-15/16")
2882 .body(layer.to_vec()),
2883 BlobRangeStatus::Malformed,
2884 ),
2885 (
2886 RegistryResponse::new(206)
2887 .header("content-range", "bytes 0-15/16")
2888 .body(b"xxxxxxxxxxxxxxxx".to_vec()),
2889 BlobRangeStatus::DigestMismatch,
2890 ),
2891 ] {
2892 let (manifest, _) = image_manifest(layer);
2893 let transport = MockTransport::new([
2894 response(200),
2895 manifest_response(manifest),
2896 MockStep::Response(range_response),
2897 ]);
2898 let report = diagnose_registry(
2899 &transport,
2900 options(Some(image("registry.example/team/app:v1"))),
2901 )
2902 .await
2903 .unwrap();
2904 assert_eq!(report.blob_range.status, expected);
2905 }
2906 }
2907
2908 #[tokio::test]
2909 async fn ignored_range_keeps_a_bounded_sample_instead_of_failing_body_limit() {
2910 let layer = b"0123456789abcdef";
2911 let (manifest, _) = image_manifest(layer);
2912 let transport = MockTransport::new([
2913 response(200),
2914 manifest_response(manifest),
2915 MockStep::Response(
2916 RegistryResponse::new(200)
2917 .header("content-length", "16")
2918 .truncated_body(layer[..4].to_vec()),
2919 ),
2920 ]);
2921 let limits = RegistryLimits {
2922 blob_sample_bytes: 4,
2923 ..RegistryLimits::default()
2924 };
2925 let report = diagnose_registry(
2926 &transport,
2927 options(Some(image("registry.example/team/app:v1"))).with_limits(limits),
2928 )
2929 .await
2930 .unwrap();
2931 assert_eq!(report.blob_range.status, BlobRangeStatus::Ignored);
2932 assert_eq!(report.blob_range.received_bytes, 4);
2933 }
2934
2935 #[tokio::test]
2936 async fn upstream_tag_is_resolved_once_and_mirrors_use_that_digest_in_order() {
2937 let (manifest, _) = image_manifest(b"layer");
2938 let digest = sha256_digest(&manifest).unwrap();
2939 let (mirror_two_body, _) = image_manifest(b"different-layer");
2940 let transport = MockTransport::new([
2941 response(200),
2942 manifest_response(manifest.clone()),
2943 MockStep::Response(
2944 RegistryResponse::new(206)
2945 .header("content-range", "bytes 0-4/5")
2946 .body(b"layer".to_vec()),
2947 ),
2948 response(200),
2949 manifest_response(manifest),
2950 MockStep::Response(
2951 RegistryResponse::new(206)
2952 .header("content-range", "bytes 0-4/5")
2953 .body(b"layer".to_vec()),
2954 ),
2955 response(200),
2956 manifest_response(mirror_two_body),
2957 ]);
2958 let opts = options(Some(image("registry.example/team/app:moving"))).with_mirrors(vec![
2959 endpoint("mirror-one.example"),
2960 endpoint("mirror-two.example"),
2961 ]);
2962 let report = diagnose_registry(&transport, opts).await.unwrap();
2963 assert_eq!(
2964 report
2965 .mirrors
2966 .iter()
2967 .map(|mirror| mirror.status)
2968 .collect::<Vec<_>>(),
2969 [MirrorCheckStatus::Equivalent, MirrorCheckStatus::Diverged]
2970 );
2971 assert_eq!(report.status, RegistryDiagnosticStatus::Corrupt);
2972 assert_eq!(report.recommended_mirror_order, [0]);
2973 assert_eq!(report.mirrors[0].recommended_rank, Some(1));
2974 assert_eq!(report.mirrors[1].recommended_rank, None);
2975 let manifest_paths = transport
2976 .seen()
2977 .into_iter()
2978 .filter(|request| request.path.contains("/manifests/"))
2979 .collect::<Vec<_>>();
2980 assert_eq!(manifest_paths[0].path, "/v2/team/app/manifests/moving");
2981 assert_eq!(manifest_paths[1].origin, "https://mirror-one.example");
2982 assert_eq!(manifest_paths[2].origin, "https://mirror-two.example");
2983 assert!(manifest_paths[1].path.ends_with(digest.as_str()));
2984 assert!(manifest_paths[2].path.ends_with(digest.as_str()));
2985 }
2986
2987 #[tokio::test]
2988 async fn corrupt_mirror_blob_is_excluded_and_marks_the_report_corrupt() {
2989 let layer = b"layer";
2990 let (manifest, _) = image_manifest(layer);
2991 let transport = MockTransport::new([
2992 response(200),
2993 manifest_response(manifest.clone()),
2994 MockStep::Response(
2995 RegistryResponse::new(206)
2996 .header("content-range", "bytes 0-4/5")
2997 .body(layer.to_vec()),
2998 ),
2999 response(200),
3000 manifest_response(manifest),
3001 MockStep::Response(
3002 RegistryResponse::new(206)
3003 .header("content-range", "bytes 0-4/5")
3004 .body(b"xxxxx".to_vec()),
3005 ),
3006 ]);
3007 let report = diagnose_registry(
3008 &transport,
3009 options(Some(image("registry.example/team/app:v1")))
3010 .with_mirrors(vec![endpoint("mirror.example")]),
3011 )
3012 .await
3013 .unwrap();
3014
3015 assert_eq!(report.mirrors[0].status, MirrorCheckStatus::Equivalent);
3016 assert_eq!(
3017 report.mirrors[0].blob_range.status,
3018 BlobRangeStatus::DigestMismatch
3019 );
3020 assert!(report.recommended_mirror_order.is_empty());
3021 assert_eq!(report.status, RegistryDiagnosticStatus::Corrupt);
3022 }
3023
3024 #[test]
3025 fn recommendation_is_latency_sorted_with_one_based_display_ranks() {
3026 let mut report = RegistryDiagnosticReport {
3027 schema_version: REGISTRY_DIAGNOSTIC_SCHEMA_VERSION,
3028 status: RegistryDiagnosticStatus::Healthy,
3029 upstream: RegistryName::parse("registry.example").unwrap(),
3030 upstream_origin: "https://registry.example".into(),
3031 image: None,
3032 requested_platform: None,
3033 api: ApiCheck::default(),
3034 manifest: ManifestCheck::default(),
3035 blob_range: BlobRangeCheck::default(),
3036 mirrors: vec![
3037 MirrorCheck {
3038 order: 0,
3039 origin: "https://slow.example".into(),
3040 status: MirrorCheckStatus::Available,
3041 digest: None,
3042 blob_range: BlobRangeCheck::default(),
3043 elapsed_micros: 200,
3044 recommended_rank: None,
3045 },
3046 MirrorCheck {
3047 order: 1,
3048 origin: "https://fast.example".into(),
3049 status: MirrorCheckStatus::Available,
3050 digest: None,
3051 blob_range: BlobRangeCheck::default(),
3052 elapsed_micros: 100,
3053 recommended_rank: None,
3054 },
3055 ],
3056 recommended_mirror_order: Vec::new(),
3057 request_count: 0,
3058 };
3059
3060 rank_mirror_checks(&mut report);
3061 assert_eq!(report.recommended_mirror_order, [1, 0]);
3062 assert_eq!(report.mirrors[0].recommended_rank, Some(2));
3063 assert_eq!(report.mirrors[1].recommended_rank, Some(1));
3064 }
3065
3066 #[tokio::test]
3067 async fn upstream_token_never_crosses_to_mirror_origin() {
3068 let (manifest, _) = image_manifest(b"layer");
3069 let upstream_challenge = RegistryResponse::new(401).header(
3070 "www-authenticate",
3071 r#"Bearer realm="https://registry.example/token",service="registry.example""#,
3072 );
3073 let mirror_challenge = RegistryResponse::new(401).header(
3074 "www-authenticate",
3075 r#"Bearer realm="https://mirror.example/token",service="mirror.example""#,
3076 );
3077 let transport = MockTransport::new([
3078 MockStep::Response(upstream_challenge),
3079 MockStep::Response(RegistryResponse::new(200).body(br#"{"token":"upstream-secret"}"#)),
3080 manifest_response(manifest.clone()),
3081 MockStep::Response(
3082 RegistryResponse::new(206)
3083 .header("content-range", "bytes 0-4/5")
3084 .body(b"layer".to_vec()),
3085 ),
3086 MockStep::Response(mirror_challenge),
3087 MockStep::Response(RegistryResponse::new(200).body(br#"{"token":"mirror-secret"}"#)),
3088 manifest_response(manifest),
3089 MockStep::Response(
3090 RegistryResponse::new(206)
3091 .header("content-range", "bytes 0-4/5")
3092 .body(b"layer".to_vec()),
3093 ),
3094 ]);
3095 let opts = options(Some(image("registry.example/team/app:v1")))
3096 .with_mirrors(vec![endpoint("mirror.example")]);
3097 let report = diagnose_registry(&transport, opts).await.unwrap();
3098 assert_eq!(report.mirrors[0].status, MirrorCheckStatus::Equivalent);
3099 let seen = transport.seen();
3100 let mirror_probe = seen
3101 .iter()
3102 .find(|request| request.origin == "https://mirror.example" && request.path == "/v2/")
3103 .unwrap();
3104 assert!(!mirror_probe.authorized);
3105 let mirror_manifest = seen
3106 .iter()
3107 .find(|request| {
3108 request.origin == "https://mirror.example" && request.path.contains("/manifests/")
3109 })
3110 .unwrap();
3111 assert!(mirror_manifest.authorized);
3112 let mirror_blob = seen
3113 .iter()
3114 .find(|request| {
3115 request.origin == "https://mirror.example" && request.path.contains("/blobs/")
3116 })
3117 .unwrap();
3118 assert!(mirror_blob.authorized);
3119 }
3120
3121 #[tokio::test]
3122 async fn cross_origin_redirect_and_bearer_realm_fail_closed() {
3123 let transport = MockTransport::new([MockStep::Response(
3124 RegistryResponse::new(307).header("location", "https://cdn.example/v2/"),
3125 )]);
3126 let report = diagnose_registry(&transport, options(None)).await.unwrap();
3127 assert_eq!(report.api.status, ApiCheckStatus::RedirectRejected);
3128 assert_eq!(transport.seen().len(), 1);
3129
3130 let challenge = RegistryResponse::new(401).header(
3131 "www-authenticate",
3132 r#"Bearer realm="https://auth.example/token",service="registry.example""#,
3133 );
3134 let transport = MockTransport::new([MockStep::Response(challenge)]);
3135 let report = diagnose_registry(
3136 &transport,
3137 options(Some(image("registry.example/team/app:v1"))),
3138 )
3139 .await
3140 .unwrap();
3141 assert_eq!(report.api.status, ApiCheckStatus::InvalidChallenge);
3142 assert_eq!(transport.seen().len(), 1);
3143 }
3144
3145 #[tokio::test]
3146 async fn redirects_and_request_budgets_fail_closed() {
3147 let transport = MockTransport::new([
3148 MockStep::Response(RegistryResponse::new(307).header("location", "/v2/ready")),
3149 response(200),
3150 ]);
3151 let report = diagnose_registry(&transport, options(None)).await.unwrap();
3152 assert_eq!(report.api.status, ApiCheckStatus::Available);
3153 assert_eq!(transport.seen().len(), 2);
3154
3155 let transport = MockTransport::new([MockStep::Response(
3156 RegistryResponse::new(307).header("location", "http://other.example/v2/"),
3157 )]);
3158 let report = diagnose_registry(&transport, options(None)).await.unwrap();
3159 assert_eq!(report.api.status, ApiCheckStatus::RedirectRejected);
3160 assert_eq!(transport.seen().len(), 1);
3161
3162 let limits = RegistryLimits {
3163 max_requests: 1,
3164 max_redirects: 1,
3165 ..RegistryLimits::default()
3166 };
3167 let transport = MockTransport::new([response(200)]);
3168 let report = diagnose_registry(
3169 &transport,
3170 options(Some(image("registry.example/team/app:v1"))).with_limits(limits),
3171 )
3172 .await
3173 .unwrap();
3174 assert_eq!(report.status, RegistryDiagnosticStatus::LimitExceeded);
3175 assert_eq!(report.manifest.status, ManifestCheckStatus::RequestLimit);
3176 assert_eq!(report.request_count, 1);
3177 }
3178
3179 #[test]
3180 fn invalid_limits_and_registry_mismatch_are_rejected_without_io() {
3181 let invalid = RegistryLimits {
3182 max_requests: 0,
3183 ..RegistryLimits::default()
3184 };
3185 assert_eq!(
3186 invalid.validate(),
3187 Err(RegistryProtocolError::InvalidLimits)
3188 );
3189
3190 let transport = MockTransport::new([]);
3191 let error = tokio::runtime::Runtime::new()
3192 .unwrap()
3193 .block_on(diagnose_registry(
3194 &transport,
3195 options(Some(image("other.example/team/app:v1"))),
3196 ))
3197 .unwrap_err();
3198 assert_eq!(error, RegistryProtocolError::RegistryMismatch);
3199 assert!(transport.seen().is_empty());
3200
3201 let mirrors = (0..=HARD_MAX_MIRRORS)
3202 .map(|index| endpoint(&format!("mirror-{index}.example")))
3203 .collect();
3204 let transport = MockTransport::new([]);
3205 let error = tokio::runtime::Runtime::new()
3206 .unwrap()
3207 .block_on(diagnose_registry(
3208 &transport,
3209 options(None).with_mirrors(mirrors),
3210 ))
3211 .unwrap_err();
3212 assert_eq!(error, RegistryProtocolError::TooManyMirrors);
3213 assert!(transport.seen().is_empty());
3214 }
3215
3216 #[tokio::test]
3217 async fn skipped_mirrors_are_prepopulated_in_configured_order() {
3218 let transport = MockTransport::new([response(403)]);
3219 let report = diagnose_registry(
3220 &transport,
3221 options(Some(image("registry.example/team/app:v1")))
3222 .with_mirrors(vec![endpoint("first.example"), endpoint("second.example")]),
3223 )
3224 .await
3225 .unwrap();
3226 assert_eq!(
3227 report
3228 .mirrors
3229 .iter()
3230 .map(|mirror| (mirror.order, mirror.origin.as_str(), mirror.status))
3231 .collect::<Vec<_>>(),
3232 [
3233 (0, "https://first.example", MirrorCheckStatus::NotTested),
3234 (1, "https://second.example", MirrorCheckStatus::NotTested),
3235 ]
3236 );
3237 }
3238
3239 #[test]
3240 fn debug_and_json_never_echo_secrets_or_response_bodies() {
3241 let token = BearerToken {
3242 value: "secret-token".to_owned(),
3243 bound_origin: "https://registry.example".to_owned(),
3244 };
3245 assert!(!format!("{token:?}").contains("secret-token"));
3246 let response = RegistryResponse::new(401)
3247 .header("www-authenticate", "Bearer secret-token")
3248 .body(b"secret-body".to_vec());
3249 let debug = format!("{response:?}");
3250 assert!(!debug.contains("secret-token"));
3251 assert!(!debug.contains("secret-body"));
3252 }
3253
3254 #[test]
3255 fn report_json_has_a_stable_typed_shape() {
3256 let report = RegistryDiagnosticReport {
3257 schema_version: REGISTRY_DIAGNOSTIC_SCHEMA_VERSION,
3258 status: RegistryDiagnosticStatus::Healthy,
3259 upstream: RegistryName::parse("registry.example").unwrap(),
3260 upstream_origin: "https://registry.example".to_owned(),
3261 image: None,
3262 requested_platform: None,
3263 api: ApiCheck {
3264 status: ApiCheckStatus::Available,
3265 http_status: Some(200),
3266 ..ApiCheck::default()
3267 },
3268 manifest: ManifestCheck::default(),
3269 blob_range: BlobRangeCheck::default(),
3270 mirrors: Vec::new(),
3271 recommended_mirror_order: Vec::new(),
3272 request_count: 1,
3273 };
3274 assert_eq!(
3275 serde_json::to_string(&report).unwrap(),
3276 r#"{"schema_version":2,"status":"healthy","upstream":"registry.example","upstream_origin":"https://registry.example","image":null,"requested_platform":null,"api":{"status":"available","http_status":200,"bearer_challenge":false,"challenge_service_present":false,"challenge_scope_matches":null},"manifest":{"status":"not-requested","kind":null,"media_type":null,"digest":null,"child_digest":null,"selected_platform":null,"selected_os_version":null,"byte_size":null},"blob_range":{"status":"not-tested","digest":null,"requested_bytes":0,"received_bytes":0,"total_bytes":null},"mirrors":[],"recommended_mirror_order":[],"request_count":1}"#
3277 );
3278 }
3279}