| 1 | //! The container registry: OCI Distribution 1.1 on `g1t.sh/v2/`. |
| 2 | //! |
| 3 | //! A client is sent to `/v2/token` first (the `WWW-Authenticate` |
| 4 | //! challenge), and comes back with a bearer token naming what it may do; |
| 5 | //! Basic credentials and g1t tokens are taken on every endpoint too, for |
| 6 | //! clients that send them straight away. Pulls of public images need |
| 7 | //! neither. |
| 8 | |
| 9 | use futures_util::StreamExt; |
| 10 | use g1t_contracts::audit::AuditActor; |
| 11 | use g1t_contracts::events::PackageEvent; |
| 12 | use g1t_contracts::packages::Ecosystem; |
| 13 | use g1t_contracts::time::rfc3339; |
| 14 | use g1t_contracts::{Viewer, new_id}; |
| 15 | use g1t_kit::now_ms; |
| 16 | use serde_json::{Value, json}; |
| 17 | use worker::{Context, Headers, Method, Request, Response, ResponseBody, Result, Url}; |
| 18 | |
| 19 | use crate::access::{self, Action}; |
| 20 | use crate::limits; |
| 21 | use crate::db::{NewFile, NewVersion, PackageRow, UploadRow}; |
| 22 | use crate::digest::{Digest, Sha256}; |
| 23 | use crate::manifest::{self, Kind}; |
| 24 | use crate::names::{self, ImageName, Reference, Route}; |
| 25 | use crate::range; |
| 26 | use crate::store::BlobStore; |
| 27 | use crate::token::{self, Claims, Grant}; |
| 28 | use crate::upload::{self, Finished, Progress, Writer}; |
| 29 | use crate::{Caller, Packages, TargetOf}; |
| 30 | |
| 31 | const CONTAINER: &str = "container"; |
| 32 | /// Blobs larger than this are downloaded from the store directly when it |
| 33 | /// can sign a URL. |
| 34 | const REDIRECT_BYTES: u64 = 4 * 1024 * 1024; |
| 35 | /// How long a signed download URL works. |
| 36 | const REDIRECT_SECONDS: u32 = 10 * 60; |
| 37 | /// The most tags one page lists. |
| 38 | const MAX_TAGS_PAGE: u32 = 1000; |
| 39 | const DOCS: &str = "https://docs.g1t.sh/guides/containers/"; |
| 40 | |
| 41 | /// An answer in the registry's error format: `{"errors": [...]}`. |
| 42 | pub fn error(status: u16, code: &str, message: impl Into<String>) -> Result<Response> { |
| 43 | let body = json!({ "errors": [{ "code": code, "message": message.into(), "detail": null }] }); |
| 44 | Ok(Response::from_json(&body)?.with_status(status)) |
| 45 | } |
| 46 | |
| 47 | fn respond(status: u16, headers: &[(&str, String)], body: ResponseBody) -> Result<Response> { |
| 48 | let set = Headers::new(); |
| 49 | for (name, value) in headers { |
| 50 | set.set(name, value)?; |
| 51 | } |
| 52 | Ok(Response::from_body(body)?.with_status(status).with_headers(set)) |
| 53 | } |
| 54 | |
| 55 | fn empty(status: u16, headers: &[(&str, String)]) -> Result<Response> { |
| 56 | respond(status, headers, ResponseBody::Empty) |
| 57 | } |
| 58 | |
| 59 | /// Who the request comes from, as its headers say. |
| 60 | enum Credentials { |
| 61 | None, |
| 62 | /// One of this registry's tokens. |
| 63 | Token(Claims), |
| 64 | /// A g1t token or password, resolved. |
| 65 | Viewer(Viewer), |
| 66 | /// Credentials that are wrong, or a token that expired. |
| 67 | Bad, |
| 68 | } |
| 69 | |
| 70 | impl Credentials { |
| 71 | fn anonymous(&self) -> bool { |
| 72 | matches!(self, Credentials::None) || matches!(self, Credentials::Token(claims) if claims.actor.is_none()) |
| 73 | } |
| 74 | } |
| 75 | |
| 76 | /// `scheme://host`, as the client reached the registry. |
| 77 | fn origin(url: &Url) -> String { |
| 78 | let host = url.host_str().unwrap_or("g1t.sh"); |
| 79 | match url.port() { |
| 80 | Some(port) => format!("{}://{host}:{port}", url.scheme()), |
| 81 | None => format!("{}://{host}", url.scheme()), |
| 82 | } |
| 83 | } |
| 84 | |
| 85 | fn service_name(url: &Url) -> String { |
| 86 | match url.port() { |
| 87 | Some(port) => format!("{}:{port}", url.host_str().unwrap_or("g1t.sh")), |
| 88 | None => url.host_str().unwrap_or("g1t.sh").to_owned(), |
| 89 | } |
| 90 | } |
| 91 | |
| 92 | /// The challenge that sends a client to get a token. |
| 93 | fn challenge(url: &Url, scope: Option<String>) -> String { |
| 94 | let mut text = format!("Bearer realm=\"{}/v2/token\",service=\"{}\"", origin(url), service_name(url)); |
| 95 | if let Some(scope) = scope { |
| 96 | text.push_str(&format!(",scope=\"{scope}\"")); |
| 97 | } |
| 98 | text |
| 99 | } |
| 100 | |
| 101 | fn unauthorized(url: &Url, scope: Option<String>, message: &str) -> Result<Response> { |
| 102 | let mut response = error(401, "UNAUTHORIZED", message)?; |
| 103 | response.headers_mut().set("www-authenticate", &challenge(url, scope))?; |
| 104 | Ok(response) |
| 105 | } |
| 106 | |
| 107 | fn query(url: &Url, key: &str) -> Option<String> { |
| 108 | url.query_pairs().find(|(k, _)| k == key).map(|(_, v)| v.into_owned()) |
| 109 | } |
| 110 | |
| 111 | /// The audit actor a version's `published_by` names. |
| 112 | fn published_by(caller: &Caller) -> Option<String> { |
| 113 | caller.actor.as_ref().map(|actor| actor.on_behalf_of.clone().unwrap_or_else(|| actor.actor.clone())) |
| 114 | } |
| 115 | |
| 116 | fn too_large(limit: u64) -> Result<Response> { |
| 117 | let mb = limit / 1_000_000; |
| 118 | error( |
| 119 | 413, |
| 120 | "SIZE_INVALID", |
| 121 | format!( |
| 122 | "A request to this registry may hold at most {mb} MB, and this one holds more. `docker push` sends each layer in one request, so a layer has to be under {mb} MB. See {DOCS}#the-{mb}-mb-limit" |
| 123 | ), |
| 124 | ) |
| 125 | } |
| 126 | |
| 127 | impl Packages { |
| 128 | /// Answers a registry request. |
| 129 | pub async fn registry(&self, request: Request, ctx: &Context) -> Result<Response> { |
| 130 | let url = request.url()?; |
| 131 | let Some(route) = names::route(url.path()) else { |
| 132 | return error(404, "NAME_UNKNOWN", "There is nothing at this address."); |
| 133 | }; |
| 134 | let mut response = match self.route(request, &url, route, ctx).await { |
| 135 | Ok(response) => response, |
| 136 | Err(problem) => { |
| 137 | worker::console_error!("packages: {} {}: {problem}", url.path(), problem); |
| 138 | error(500, "UNKNOWN", "Something went wrong on our side. Try again in a moment.")? |
| 139 | } |
| 140 | }; |
| 141 | response.headers_mut().set("docker-distribution-api-version", "registry/2.0")?; |
| 142 | Ok(response) |
| 143 | } |
| 144 | |
| 145 | async fn route(&self, mut request: Request, url: &Url, route: Route, ctx: &Context) -> Result<Response> { |
| 146 | let method = request.method(); |
| 147 | let credentials = self.credentials(&request).await?; |
| 148 | if limits::counts(&method, &route) |
| 149 | && let Some(refused) = self.limited(&request, &credentials).await? |
| 150 | { |
| 151 | return Ok(refused); |
| 152 | } |
| 153 | if let Route::Base = route { |
| 154 | return match credentials { |
| 155 | Credentials::Token(_) | Credentials::Viewer(Some(_)) => Ok(Response::from_json(&json!({}))?), |
| 156 | _ => unauthorized(url, None, "Sign in with `docker login`, a g1t token as the password."), |
| 157 | }; |
| 158 | } |
| 159 | if let Route::Token = route { |
| 160 | // Docker's OAuth form: `grant_type=password` with the username |
| 161 | // and password in the body, and the scopes beside them. |
| 162 | if method == Method::Post { |
| 163 | let form = request.text().await.unwrap_or_default(); |
| 164 | let fields: Vec<(String, String)> = Url::parse(&format!("http://form/?{form}"))? |
| 165 | .query_pairs() |
| 166 | .map(|(k, v)| (k.into_owned(), v.into_owned())) |
| 167 | .collect(); |
| 168 | let field = |key: &str| fields.iter().find(|(k, _)| k == key).map(|(_, v)| v.clone()); |
| 169 | let credentials = match (field("grant_type").as_deref(), field("username"), field("password")) { |
| 170 | (Some("password"), Some(username), Some(password)) => match self.viewer_for(&username, &password).await? { |
| 171 | Some(user) => Credentials::Viewer(Some(user)), |
| 172 | None => Credentials::Bad, |
| 173 | }, |
| 174 | (Some("password"), ..) | (Some("refresh_token"), ..) => Credentials::Bad, |
| 175 | _ => credentials, |
| 176 | }; |
| 177 | let scopes: Vec<String> = fields.iter().filter(|(k, _)| k == "scope").map(|(_, v)| v.clone()).collect(); |
| 178 | return self.issue(url, credentials, &scopes).await; |
| 179 | } |
| 180 | let scopes: Vec<String> = url.query_pairs().filter(|(k, _)| k == "scope").map(|(_, v)| v.into_owned()).collect(); |
| 181 | return self.issue(url, credentials, &scopes).await; |
| 182 | } |
| 183 | let (name, action) = match &route { |
| 184 | Route::Manifest { name, .. } => (name, match method { |
| 185 | Method::Get | Method::Head => Action::Pull, |
| 186 | Method::Put => Action::Push, |
| 187 | Method::Delete => Action::Delete, |
| 188 | _ => return error(405, "UNSUPPORTED", "Not a method manifests take."), |
| 189 | }), |
| 190 | Route::Blob { name, .. } => (name, match method { |
| 191 | Method::Get | Method::Head => Action::Pull, |
| 192 | Method::Delete => Action::Delete, |
| 193 | _ => return error(405, "UNSUPPORTED", "Not a method blobs take."), |
| 194 | }), |
| 195 | Route::Uploads { name } | Route::Upload { name, .. } => (name, Action::Push), |
| 196 | Route::Tags { name } | Route::Referrers { name, .. } => (name, Action::Pull), |
| 197 | Route::Base | Route::Token => unreachable!("answered above"), |
| 198 | }; |
| 199 | let name = match names::parse_name(name) { |
| 200 | Ok(name) => name, |
| 201 | Err(message) => return error(400, "NAME_INVALID", message), |
| 202 | }; |
| 203 | let found = self.db.package(&name.workspace, CONTAINER, &name.name).await?; |
| 204 | // A deleted workspace's images are gone for everyone until it is |
| 205 | // restored: pulls find nothing, and nothing new is pushed to it. |
| 206 | let hidden = match &found { |
| 207 | Some(package) => package.hidden(), |
| 208 | None => action == Action::Push && self.db.workspace_hidden(&name.workspace).await?, |
| 209 | }; |
| 210 | if hidden { |
| 211 | return if action == Action::Pull { |
| 212 | error(404, "NAME_UNKNOWN", format!("There is no image {}.", name.full())) |
| 213 | } else { |
| 214 | error(403, "DENIED", format!("The workspace {} is deleted; nothing can be pushed to it or deleted from it.", name.workspace)) |
| 215 | }; |
| 216 | } |
| 217 | let caller = match self.authorize(url, &credentials, &name, found.as_ref(), action).await? { |
| 218 | Ok(caller) => caller, |
| 219 | Err(refused) => return Ok(refused), |
| 220 | }; |
| 221 | match route { |
| 222 | Route::Manifest { reference, .. } => match method { |
| 223 | Method::Put => self.put_manifest(&mut request, &name, found, &reference, &caller).await, |
| 224 | Method::Delete => self.delete_manifest(&name, found, &reference, &caller).await, |
| 225 | _ => self.get_manifest(&name, found, &reference, method == Method::Head, ctx).await, |
| 226 | }, |
| 227 | Route::Blob { digest, .. } => { |
| 228 | let Some(digest) = Digest::parse(&digest) else { |
| 229 | return error(400, "DIGEST_INVALID", format!("{digest} is not a sha256 digest.")); |
| 230 | }; |
| 231 | let Some(package) = found else { |
| 232 | return error(404, "BLOB_UNKNOWN", "No such blob in this image."); |
| 233 | }; |
| 234 | match method { |
| 235 | Method::Delete => self.delete_blob(&package, &digest).await, |
| 236 | _ => self.get_blob(&request, &package, &digest, method == Method::Head).await, |
| 237 | } |
| 238 | } |
| 239 | Route::Uploads { .. } => { |
| 240 | if method != Method::Post { |
| 241 | return error(405, "UNSUPPORTED", "Start an upload with POST."); |
| 242 | } |
| 243 | self.start_upload(request, url, &name, found, &credentials, &caller).await |
| 244 | } |
| 245 | Route::Upload { id, .. } => { |
| 246 | let Some(row) = self.db.upload(&id).await?.filter(|row| row.package == name.full()) else { |
| 247 | return error(404, "BLOB_UPLOAD_UNKNOWN", "No such upload; it may have finished or expired. Start again."); |
| 248 | }; |
| 249 | match method { |
| 250 | Method::Patch => self.patch_upload(request, &name, row).await, |
| 251 | Method::Put => match found { |
| 252 | Some(package) => self.finish_upload(request, url, &name, &package, row).await, |
| 253 | None => { |
| 254 | self.cancel_upload(row).await?; |
| 255 | error(404, "NAME_UNKNOWN", format!("{} was deleted while this upload was open.", name.full())) |
| 256 | } |
| 257 | }, |
| 258 | Method::Get => upload_status(&name, &row), |
| 259 | Method::Delete => self.cancel_upload(row).await, |
| 260 | _ => error(405, "UNSUPPORTED", "Not a method uploads take."), |
| 261 | } |
| 262 | } |
| 263 | Route::Tags { .. } => self.tags(url, &name, found).await, |
| 264 | Route::Referrers { digest, .. } => self.referrers(url, &name, found, &digest).await, |
| 265 | Route::Base | Route::Token => unreachable!("answered above"), |
| 266 | } |
| 267 | } |
| 268 | |
| 269 | /// The 429 for a client past its limit, if it is. A limit that is not |
| 270 | /// configured (self-hosted) or cannot be asked lets the request through. |
| 271 | async fn limited(&self, request: &Request, credentials: &Credentials) -> Result<Option<Response>> { |
| 272 | let subject = match credentials { |
| 273 | Credentials::Token(claims) => claims.actor.as_ref().map(|actor| actor.actor_id.clone()), |
| 274 | Credentials::Viewer(Some(user)) => Some(user.id.clone()), |
| 275 | _ => None, |
| 276 | }; |
| 277 | let address = request.headers().get("cf-connecting-ip")?; |
| 278 | let (limit, key) = limits::key(subject.as_deref(), address.as_deref()); |
| 279 | let Ok(limiter) = self.env.rate_limiter(limit.binding()) else { |
| 280 | return Ok(None); |
| 281 | }; |
| 282 | match limiter.limit(key).await { |
| 283 | Ok(outcome) if !outcome.success => { |
| 284 | let mut response = error( |
| 285 | 429, |
| 286 | "TOOMANYREQUESTS", |
| 287 | match limit { |
| 288 | limits::Limit::Anonymous => "Too many requests from this address. Wait a minute, or sign in with `docker login g1t.sh` for a higher limit.", |
| 289 | limits::Limit::Signed => "Too many requests. Wait a minute and try again.", |
| 290 | }, |
| 291 | )?; |
| 292 | response.headers_mut().set("retry-after", &limits::RETRY_AFTER_SECONDS.to_string())?; |
| 293 | Ok(Some(response)) |
| 294 | } |
| 295 | Ok(_) => Ok(None), |
| 296 | Err(problem) => { |
| 297 | worker::console_error!("packages: the rate limit could not be asked: {problem}"); |
| 298 | Ok(None) |
| 299 | } |
| 300 | } |
| 301 | } |
| 302 | |
| 303 | async fn credentials(&self, request: &Request) -> Result<Credentials> { |
| 304 | let Some(header) = request.headers().get("authorization")? else { |
| 305 | return Ok(Credentials::None); |
| 306 | }; |
| 307 | if let Some(bearer) = token::bearer(&header) { |
| 308 | if token::is_registry_token(bearer) { |
| 309 | return Ok(match token::verify(bearer, &self.secret, now_ms() / 1000) { |
| 310 | Some(claims) => Credentials::Token(claims), |
| 311 | None => Credentials::Bad, |
| 312 | }); |
| 313 | } |
| 314 | return Ok(match self.viewer_for("token", bearer).await? { |
| 315 | Some(user) => Credentials::Viewer(Some(user)), |
| 316 | None => Credentials::Bad, |
| 317 | }); |
| 318 | } |
| 319 | if let Some((username, secret)) = token::basic(&header) { |
| 320 | return Ok(match self.viewer_for(&username, &secret).await? { |
| 321 | Some(user) => Credentials::Viewer(Some(user)), |
| 322 | None => Credentials::Bad, |
| 323 | }); |
| 324 | } |
| 325 | Ok(Credentials::Bad) |
| 326 | } |
| 327 | |
| 328 | /// Whether the request may do `action` to the image, and as whom. |
| 329 | async fn authorize( |
| 330 | &self, |
| 331 | url: &Url, |
| 332 | credentials: &Credentials, |
| 333 | name: &ImageName, |
| 334 | found: Option<&PackageRow>, |
| 335 | action: Action, |
| 336 | ) -> Result<std::result::Result<Caller, Response>> { |
| 337 | let scope = Some(format!("repository:{}:{}", name.full(), match action { |
| 338 | Action::Pull => "pull", |
| 339 | Action::Push => "pull,push", |
| 340 | Action::Delete => "delete", |
| 341 | })); |
| 342 | let viewer = match credentials { |
| 343 | Credentials::Bad => { |
| 344 | return Ok(Err(unauthorized(url, scope, "The token or password is not right, or has expired. Sign in again with `docker login`.")?)); |
| 345 | } |
| 346 | Credentials::Token(claims) => { |
| 347 | if claims.allows(&name.full(), action) { |
| 348 | return Ok(Ok(Caller { actor: claims.actor.clone() })); |
| 349 | } |
| 350 | // A token for something else still pulls a public image. |
| 351 | if action == Action::Pull && found.is_some_and(|package| package.public()) { |
| 352 | return Ok(Ok(Caller { actor: claims.actor.clone() })); |
| 353 | } |
| 354 | if credentials.anonymous() { |
| 355 | return Ok(Err(unauthorized(url, scope, "Sign in with `docker login` to do that.")?)); |
| 356 | } |
| 357 | return Ok(Err(error(403, "DENIED", format!("This token may not {} {}.", action.as_str(), name.full()))?)); |
| 358 | } |
| 359 | Credentials::Viewer(viewer) => viewer.clone(), |
| 360 | Credentials::None => None, |
| 361 | }; |
| 362 | let target = self.target(name, found).await?; |
| 363 | let decision = access::decide(viewer.as_ref(), &target.view(), action); |
| 364 | if decision.allowed { |
| 365 | return Ok(Ok(Caller { actor: viewer.as_ref().map(AuditActor::of) })); |
| 366 | } |
| 367 | let reason = decision.reason.unwrap_or_else(|| "Not allowed.".to_owned()); |
| 368 | if viewer.is_none() { |
| 369 | return Ok(Err(unauthorized(url, scope, &reason)?)); |
| 370 | } |
| 371 | Ok(Err(error(403, "DENIED", reason)?)) |
| 372 | } |
| 373 | |
| 374 | /// `GET /v2/token`: a token for each scope asked for, cut down to what |
| 375 | /// the credentials may do. |
| 376 | async fn issue(&self, url: &Url, credentials: Credentials, asked: &[String]) -> Result<Response> { |
| 377 | let viewer = match credentials { |
| 378 | Credentials::Viewer(viewer) => viewer, |
| 379 | Credentials::None => None, |
| 380 | // A registry token is not a way to get another. |
| 381 | Credentials::Token(_) | Credentials::Bad => { |
| 382 | return unauthorized(url, None, "The username or token is not right. Use a g1t token as the password."); |
| 383 | } |
| 384 | }; |
| 385 | let mut access: Vec<Grant> = Vec::new(); |
| 386 | let scopes: Vec<String> = asked.iter().flat_map(|value| value.split(' ').map(str::to_owned)).collect(); |
| 387 | for scope in scopes.iter().take(20) { |
| 388 | let Some((full, wanted)) = token::parse_scope(scope) else { continue }; |
| 389 | let Ok(name) = names::parse_name(&full) else { continue }; |
| 390 | let found = self.db.package(&name.workspace, CONTAINER, &name.name).await?; |
| 391 | // Nothing is granted on a deleted workspace's images. |
| 392 | if found.as_ref().is_some_and(|package| package.hidden()) { |
| 393 | continue; |
| 394 | } |
| 395 | let target = self.target(&name, found.as_ref()).await?; |
| 396 | let actions: Vec<Action> = wanted |
| 397 | .into_iter() |
| 398 | .filter(|action| access::decide(viewer.as_ref(), &target.view(), *action).allowed) |
| 399 | .collect(); |
| 400 | if let Some(grant) = access.iter_mut().find(|grant| grant.name == full) { |
| 401 | for action in actions { |
| 402 | if !grant.actions.contains(&action) { |
| 403 | grant.actions.push(action); |
| 404 | } |
| 405 | } |
| 406 | } else { |
| 407 | access.push(Grant { name: full, actions }); |
| 408 | } |
| 409 | } |
| 410 | let now = now_ms(); |
| 411 | let claims = Claims { |
| 412 | actor: viewer.as_ref().map(AuditActor::of), |
| 413 | access, |
| 414 | iat: now / 1000, |
| 415 | exp: now / 1000 + token::TTL_SECONDS, |
| 416 | }; |
| 417 | let signed = token::sign(&claims, &self.secret); |
| 418 | Response::from_json(&json!({ |
| 419 | "token": signed, |
| 420 | "access_token": signed, |
| 421 | "expires_in": token::TTL_SECONDS, |
| 422 | "issued_at": rfc3339(now), |
| 423 | })) |
| 424 | } |
| 425 | |
| 426 | /// The package for an image, made on its first push and linked to the |
| 427 | /// repository of the image's name when there is one. |
| 428 | async fn package_for_push(&self, name: &ImageName, found: Option<PackageRow>, caller: &Caller) -> Result<PackageRow> { |
| 429 | if let Some(found) = found { |
| 430 | return Ok(found); |
| 431 | } |
| 432 | let repo = self.repo_by_name(&name.workspace, name.repo_name()).await?; |
| 433 | self.db |
| 434 | .create_package( |
| 435 | &new_id("pkg", now_ms()), |
| 436 | &name.workspace, |
| 437 | CONTAINER, |
| 438 | &name.name, |
| 439 | repo.as_ref().map(|repo| (repo.id.as_str(), repo.name.as_str(), repo.is_private)), |
| 440 | caller.actor.as_ref().map_or("", |actor| actor.actor_id.as_str()), |
| 441 | now_ms(), |
| 442 | ) |
| 443 | .await |
| 444 | } |
| 445 | |
| 446 | async fn get_manifest( |
| 447 | &self, |
| 448 | name: &ImageName, |
| 449 | found: Option<PackageRow>, |
| 450 | reference: &str, |
| 451 | head: bool, |
| 452 | ctx: &Context, |
| 453 | ) -> Result<Response> { |
| 454 | let Some(package) = found else { |
| 455 | return error(404, "NAME_UNKNOWN", format!("There is no image {}.", name.full())); |
| 456 | }; |
| 457 | let version = match Reference::parse(reference) { |
| 458 | Some(Reference::Tag(tag)) => self.db.version_by_tag(&package.id, &tag).await?, |
| 459 | Some(Reference::Digest(digest)) => self.db.version_by_digest(&package.id, digest.as_str()).await?, |
| 460 | None => return error(400, "MANIFEST_INVALID", format!("{reference} is not a tag or a digest.")), |
| 461 | }; |
| 462 | let Some(version) = version else { |
| 463 | return error(404, "MANIFEST_UNKNOWN", format!("{}:{reference} is not there.", name.full())); |
| 464 | }; |
| 465 | let digest = Digest::parse(&version.digest).ok_or_else(|| worker::Error::RustError("a stored digest is malformed".into()))?; |
| 466 | let Some(blob) = self.db.blob(&digest).await? else { |
| 467 | return error(404, "MANIFEST_UNKNOWN", format!("{}:{reference} is not there.", name.full())); |
| 468 | }; |
| 469 | let headers = [ |
| 470 | ("content-type", version.media_type().unwrap_or_else(|| manifest::OCI_MANIFEST.to_owned())), |
| 471 | ("docker-content-digest", version.digest.clone()), |
| 472 | ("etag", format!("\"{}\"", version.digest)), |
| 473 | ("content-length", blob.size.to_string()), |
| 474 | ]; |
| 475 | if head { |
| 476 | return empty(200, &headers); |
| 477 | } |
| 478 | let Some(bytes) = self.store.read(&blob.object_key).await? else { |
| 479 | return error(404, "MANIFEST_UNKNOWN", format!("{}:{reference} is not there.", name.full())); |
| 480 | }; |
| 481 | self.count_download(&package.id, ctx); |
| 482 | respond(200, &headers, ResponseBody::Body(bytes)) |
| 483 | } |
| 484 | |
| 485 | async fn put_manifest( |
| 486 | &self, |
| 487 | request: &mut Request, |
| 488 | name: &ImageName, |
| 489 | found: Option<PackageRow>, |
| 490 | reference: &str, |
| 491 | caller: &Caller, |
| 492 | ) -> Result<Response> { |
| 493 | let reference = match Reference::parse(reference) { |
| 494 | Some(reference) => reference, |
| 495 | None => return error(400, "TAG_INVALID", format!("{reference} is not a valid tag.")), |
| 496 | }; |
| 497 | let bytes = request.bytes().await?; |
| 498 | if bytes.len() > manifest::MAX_MANIFEST_BYTES { |
| 499 | return error(413, "SIZE_INVALID", "A manifest is at most 4 MiB."); |
| 500 | } |
| 501 | let digest = Digest::of(&bytes); |
| 502 | if let Reference::Digest(given) = &reference |
| 503 | && given != &digest |
| 504 | { |
| 505 | return error(400, "DIGEST_INVALID", format!("The manifest's digest is {digest}, not {given}.")); |
| 506 | } |
| 507 | let content_type = request.headers().get("content-type")?; |
| 508 | let parsed = match manifest::parse(&bytes, content_type.as_deref()) { |
| 509 | Ok(parsed) => parsed, |
| 510 | Err(manifest::Refused::Invalid(message)) => return error(400, "MANIFEST_INVALID", message), |
| 511 | Err(manifest::Refused::Unsupported(message)) => return error(415, "UNSUPPORTED", message), |
| 512 | }; |
| 513 | let package = self.package_for_push(name, found, caller).await?; |
| 514 | let mut files = vec![NewFile { |
| 515 | name: "manifest".to_owned(), |
| 516 | digest: digest.to_string(), |
| 517 | size: bytes.len() as u64, |
| 518 | media_type: Some(parsed.media_type.clone()), |
| 519 | }]; |
| 520 | match parsed.kind { |
| 521 | Kind::Image => { |
| 522 | let mut layer = 0; |
| 523 | for blob in &parsed.blobs { |
| 524 | let Some(stored) = self.db.package_blob(&package.id, &blob.digest).await? else { |
| 525 | return error(400, "MANIFEST_BLOB_UNKNOWN", format!("Blob {} is not in {}: push it first.", blob.digest, name.full())); |
| 526 | }; |
| 527 | let file_name = if blob.role == "config" { |
| 528 | "config".to_owned() |
| 529 | } else { |
| 530 | layer += 1; |
| 531 | format!("layer:{layer}") |
| 532 | }; |
| 533 | files.push(NewFile { name: file_name, digest: blob.digest.to_string(), size: stored.size, media_type: blob.media_type.clone() }); |
| 534 | } |
| 535 | } |
| 536 | Kind::Index => { |
| 537 | for child in &parsed.manifests { |
| 538 | if self.db.version_by_digest(&package.id, child.digest.as_str()).await?.is_none() { |
| 539 | return error(400, "MANIFEST_UNKNOWN", format!("Manifest {} is not in {}: push it first.", child.digest, name.full())); |
| 540 | } |
| 541 | } |
| 542 | } |
| 543 | } |
| 544 | let pushed: Vec<(String, u64)> = files.iter().map(|file| (file.digest.clone(), file.size)).collect(); |
| 545 | if let Some(refused) = self.storage_refusal(&package, &pushed).await? { |
| 546 | return error(403, "DENIED", refused); |
| 547 | } |
| 548 | if self.db.blob(&digest).await?.is_none() { |
| 549 | self.store.put(&digest.object_key(), bytes.clone()).await?; |
| 550 | } |
| 551 | let now = now_ms(); |
| 552 | self.db |
| 553 | .keep_blob(&package.id, &digest, bytes.len() as u64, Some(&parsed.media_type), &digest.object_key(), now) |
| 554 | .await?; |
| 555 | let size = files.iter().map(|file| file.size).sum(); |
| 556 | let metadata = json!({ |
| 557 | "media_type": parsed.media_type, |
| 558 | "artifact_type": parsed.artifact_type, |
| 559 | "annotations": parsed.annotations, |
| 560 | "platforms": parsed.platforms, |
| 561 | }); |
| 562 | let tag = match &reference { |
| 563 | Reference::Tag(tag) => Some(tag.as_str()), |
| 564 | Reference::Digest(_) => None, |
| 565 | }; |
| 566 | let (version, changed) = self |
| 567 | .db |
| 568 | .publish( |
| 569 | NewVersion { |
| 570 | id: new_id("ver", now), |
| 571 | package_id: package.id.clone(), |
| 572 | version: digest.to_string(), |
| 573 | digest: digest.to_string(), |
| 574 | size, |
| 575 | metadata: metadata.to_string(), |
| 576 | subject: parsed.subject.as_ref().map(Digest::to_string), |
| 577 | published_by: published_by(caller), |
| 578 | files, |
| 579 | }, |
| 580 | tag, |
| 581 | now, |
| 582 | ) |
| 583 | .await?; |
| 584 | self.db.measure(&package.workspace).await?; |
| 585 | if changed { |
| 586 | let tags: Vec<String> = self |
| 587 | .db |
| 588 | .tags(&package.id) |
| 589 | .await? |
| 590 | .into_iter() |
| 591 | .filter(|row| row.version_id == version.id) |
| 592 | .map(|row| row.tag) |
| 593 | .collect(); |
| 594 | let event = PackageEvent { |
| 595 | version: Some(version.version.clone()), |
| 596 | digest: Some(version.digest.clone()), |
| 597 | size: Some(version.size), |
| 598 | tags: Some(tags), |
| 599 | ..self.event_of(&package) |
| 600 | }; |
| 601 | self.announce("package.published", &package, event, caller).await; |
| 602 | self.audit(caller, "package.publish", &package, Some(&format!("{}@{}", name.full(), version.digest)), None) |
| 603 | .await; |
| 604 | } |
| 605 | let mut headers = vec![ |
| 606 | ("location", format!("/v2/{}/manifests/{digest}", name.full())), |
| 607 | ("docker-content-digest", digest.to_string()), |
| 608 | ]; |
| 609 | if let Some(subject) = &parsed.subject { |
| 610 | headers.push(("oci-subject", subject.to_string())); |
| 611 | } |
| 612 | empty(201, &headers) |
| 613 | } |
| 614 | |
| 615 | async fn delete_manifest(&self, name: &ImageName, found: Option<PackageRow>, reference: &str, caller: &Caller) -> Result<Response> { |
| 616 | let Some(package) = found else { |
| 617 | return error(404, "NAME_UNKNOWN", format!("There is no image {}.", name.full())); |
| 618 | }; |
| 619 | match Reference::parse(reference) { |
| 620 | Some(Reference::Tag(tag)) => { |
| 621 | if self.db.version_by_tag(&package.id, &tag).await?.is_none() { |
| 622 | return error(404, "MANIFEST_UNKNOWN", format!("{}:{tag} is not there.", name.full())); |
| 623 | } |
| 624 | self.db.delete_tag(&package.id, &tag).await?; |
| 625 | self.audit(caller, "package.untag", &package, Some(&format!("{}:{tag}", name.full())), None).await; |
| 626 | empty(202, &[]) |
| 627 | } |
| 628 | Some(Reference::Digest(digest)) => { |
| 629 | let Some(version) = self.db.version_by_digest(&package.id, digest.as_str()).await? else { |
| 630 | return error(404, "MANIFEST_UNKNOWN", format!("{}@{digest} is not there.", name.full())); |
| 631 | }; |
| 632 | self.remove_version(&package, &version, caller).await?; |
| 633 | empty(202, &[]) |
| 634 | } |
| 635 | None => error(400, "MANIFEST_INVALID", format!("{reference} is not a tag or a digest.")), |
| 636 | } |
| 637 | } |
| 638 | |
| 639 | async fn get_blob(&self, request: &Request, package: &PackageRow, digest: &Digest, head: bool) -> Result<Response> { |
| 640 | let Some(blob) = self.db.package_blob(&package.id, digest).await? else { |
| 641 | return error(404, "BLOB_UNKNOWN", format!("Blob {digest} is not in this image.")); |
| 642 | }; |
| 643 | let mut headers = vec![ |
| 644 | ("docker-content-digest", digest.to_string()), |
| 645 | ("content-type", "application/octet-stream".to_owned()), |
| 646 | ("accept-ranges", "bytes".to_owned()), |
| 647 | ("etag", format!("\"{digest}\"")), |
| 648 | ("cache-control", "max-age=31536000".to_owned()), |
| 649 | ]; |
| 650 | if head { |
| 651 | headers.push(("content-length", blob.size.to_string())); |
| 652 | return empty(200, &headers); |
| 653 | } |
| 654 | let wanted = match range::parse_range(request.headers().get("range")?.as_deref(), blob.size) { |
| 655 | Ok(wanted) => wanted, |
| 656 | Err(()) => { |
| 657 | headers.push(("content-range", format!("bytes */{}", blob.size))); |
| 658 | return respond(416, &headers, ResponseBody::Empty); |
| 659 | } |
| 660 | }; |
| 661 | if wanted.is_none() |
| 662 | && blob.size > REDIRECT_BYTES |
| 663 | && let Some(url) = self.store.presign_get(&blob.object_key, REDIRECT_SECONDS, now_ms()) |
| 664 | { |
| 665 | headers.push(("location", url)); |
| 666 | return empty(307, &headers); |
| 667 | } |
| 668 | let Some(got) = self.store.get(&blob.object_key, wanted).await? else { |
| 669 | return error(404, "BLOB_UNKNOWN", format!("Blob {digest} is not in this image.")); |
| 670 | }; |
| 671 | match wanted { |
| 672 | Some(wanted) => { |
| 673 | headers.push(("content-range", wanted.content_range(got.size))); |
| 674 | headers.push(("content-length", wanted.length.to_string())); |
| 675 | respond(206, &headers, got.body) |
| 676 | } |
| 677 | None => { |
| 678 | headers.push(("content-length", got.size.to_string())); |
| 679 | respond(200, &headers, got.body) |
| 680 | } |
| 681 | } |
| 682 | } |
| 683 | |
| 684 | async fn delete_blob(&self, package: &PackageRow, digest: &Digest) -> Result<Response> { |
| 685 | if self.db.package_blob(&package.id, digest).await?.is_none() { |
| 686 | return error(404, "BLOB_UNKNOWN", format!("Blob {digest} is not in this image.")); |
| 687 | } |
| 688 | if self.db.blob_in_use(&package.id, digest).await? { |
| 689 | return error(405, "DENIED", format!("Blob {digest} is used by a manifest of this image; delete the manifest instead.")); |
| 690 | } |
| 691 | self.db.unlink_blob(&package.id, digest).await?; |
| 692 | empty(202, &[]) |
| 693 | } |
| 694 | |
| 695 | async fn start_upload( |
| 696 | &self, |
| 697 | request: Request, |
| 698 | url: &Url, |
| 699 | name: &ImageName, |
| 700 | found: Option<PackageRow>, |
| 701 | credentials: &Credentials, |
| 702 | caller: &Caller, |
| 703 | ) -> Result<Response> { |
| 704 | let package = self.package_for_push(name, found, caller).await?; |
| 705 | // A blob another image of the workspace has: linked, not copied. |
| 706 | if let (Some(mount), Some(from)) = (query(url, "mount"), query(url, "from")) |
| 707 | && let (Some(digest), Ok(from)) = (Digest::parse(&mount), names::parse_name(&from)) |
| 708 | && from.workspace == name.workspace |
| 709 | && let Some(source) = self.db.package(&from.workspace, CONTAINER, &from.name).await? |
| 710 | && self.may_pull(credentials, &from, &source).await? |
| 711 | && let Some(blob) = self.db.package_blob(&source.id, &digest).await? |
| 712 | { |
| 713 | if let Some(refused) = self.storage_refusal(&package, &[(digest.to_string(), blob.size)]).await? { |
| 714 | return error(403, "DENIED", refused); |
| 715 | } |
| 716 | self.db.link_blob(&package.id, &digest, now_ms()).await?; |
| 717 | return empty(201, &[ |
| 718 | ("location", format!("/v2/{}/blobs/{digest}", name.full())), |
| 719 | ("docker-content-digest", digest.to_string()), |
| 720 | ]); |
| 721 | } |
| 722 | let id = new_id("upl", now_ms()); |
| 723 | let row = UploadRow { |
| 724 | id: id.clone(), |
| 725 | workspace: name.workspace.clone(), |
| 726 | package_id: package.id.clone(), |
| 727 | package: name.full(), |
| 728 | multipart_id: None, |
| 729 | parts: "[]".to_owned(), |
| 730 | offset: 0, |
| 731 | tail: 0, |
| 732 | hash_state: Sha256::new().save(), |
| 733 | }; |
| 734 | // The whole blob in this one request. |
| 735 | if let Some(digest) = query(url, "digest") { |
| 736 | let Some(digest) = Digest::parse(&digest) else { |
| 737 | return error(400, "DIGEST_INVALID", format!("{digest} is not a sha256 digest.")); |
| 738 | }; |
| 739 | let writer = Writer::resume(&self.store, Progress::new(&id)).await?; |
| 740 | let writer = match self.take_body(request, writer).await? { |
| 741 | Ok(writer) => writer, |
| 742 | Err(refused) => return Ok(refused), |
| 743 | }; |
| 744 | return self.complete(writer, name, &package, &digest).await; |
| 745 | } |
| 746 | self.db.create_upload(&row, now_ms()).await?; |
| 747 | upload_status_with(name, &id, 0, 202) |
| 748 | } |
| 749 | |
| 750 | /// Whether the credentials may pull another image, for a mount. |
| 751 | async fn may_pull(&self, credentials: &Credentials, name: &ImageName, package: &PackageRow) -> Result<bool> { |
| 752 | Ok(match credentials { |
| 753 | Credentials::Token(claims) => claims.allows(&name.full(), Action::Pull) || package.public(), |
| 754 | Credentials::Viewer(viewer) => { |
| 755 | let target = TargetOf::package(package); |
| 756 | access::decide(viewer.as_ref(), &target.view(), Action::Pull).allowed |
| 757 | } |
| 758 | Credentials::None | Credentials::Bad => package.public(), |
| 759 | }) |
| 760 | } |
| 761 | |
| 762 | /// Writes the request's body into the upload, refusing one over the |
| 763 | /// limit. On a refusal or a failure the whole upload is let go, the |
| 764 | /// parts this request sent and the tail an earlier one kept included, |
| 765 | /// so nothing is left in the store; the caller forgets its row. |
| 766 | async fn take_body<'a, S: BlobStore>( |
| 767 | &self, |
| 768 | mut request: Request, |
| 769 | mut writer: Writer<'a, S>, |
| 770 | ) -> Result<std::result::Result<Writer<'a, S>, Response>> { |
| 771 | let declared = request.headers().get("content-length")?.and_then(|n| n.parse::<u64>().ok()); |
| 772 | if declared.is_some_and(|n| n > self.max_request) { |
| 773 | writer.abandon().await?; |
| 774 | return Ok(Err(too_large(self.max_request)?)); |
| 775 | } |
| 776 | let mut received = 0u64; |
| 777 | let mut stream = match request.stream() { |
| 778 | Ok(stream) => stream, |
| 779 | // No body at all. |
| 780 | Err(_) => return Ok(Ok(writer)), |
| 781 | }; |
| 782 | while let Some(chunk) = stream.next().await { |
| 783 | let written = match chunk { |
| 784 | Ok(chunk) => { |
| 785 | received += chunk.len() as u64; |
| 786 | if received > self.max_request { |
| 787 | writer.abandon().await?; |
| 788 | return Ok(Err(too_large(self.max_request)?)); |
| 789 | } |
| 790 | writer.write(&chunk).await |
| 791 | } |
| 792 | Err(error) => Err(error), |
| 793 | }; |
| 794 | if let Err(error) = written { |
| 795 | let _ = writer.abandon().await; |
| 796 | return Err(error); |
| 797 | } |
| 798 | } |
| 799 | Ok(Ok(writer)) |
| 800 | } |
| 801 | |
| 802 | async fn patch_upload(&self, request: Request, name: &ImageName, row: UploadRow) -> Result<Response> { |
| 803 | let Some(progress) = row.progress() else { |
| 804 | return error(404, "BLOB_UPLOAD_INVALID", "This upload cannot be continued. Start again."); |
| 805 | }; |
| 806 | if let Some(header) = request.headers().get("content-range")? { |
| 807 | match range::parse_content_range(&header) { |
| 808 | Some(chunk) if chunk.start == progress.offset => {} |
| 809 | _ => { |
| 810 | return empty(416, &[ |
| 811 | ("location", format!("/v2/{}/blobs/uploads/{}", name.full(), row.id)), |
| 812 | ("range", range::upload_range(progress.offset)), |
| 813 | ("docker-upload-uuid", row.id.clone()), |
| 814 | ]); |
| 815 | } |
| 816 | } |
| 817 | } |
| 818 | let writer = Writer::resume(&self.store, progress).await?; |
| 819 | // A refused or failed chunk ends the upload: the store holds |
| 820 | // nothing of it any more (take_body, pause), so neither does its row. |
| 821 | let progress = match self.take_body(request, writer).await { |
| 822 | Ok(Ok(writer)) => writer.pause().await, |
| 823 | Ok(Err(refused)) => { |
| 824 | self.db.delete_upload(&row.id).await?; |
| 825 | return Ok(refused); |
| 826 | } |
| 827 | Err(error) => Err(error), |
| 828 | }; |
| 829 | let progress = match progress { |
| 830 | Ok(progress) => progress, |
| 831 | Err(error) => { |
| 832 | self.db.delete_upload(&row.id).await?; |
| 833 | return Err(error); |
| 834 | } |
| 835 | }; |
| 836 | self.db.save_progress(&progress, now_ms()).await?; |
| 837 | upload_status_with(name, &row.id, progress.offset, 202) |
| 838 | } |
| 839 | |
| 840 | async fn finish_upload(&self, request: Request, url: &Url, name: &ImageName, package: &PackageRow, row: UploadRow) -> Result<Response> { |
| 841 | let Some(digest) = query(url, "digest").and_then(|d| Digest::parse(&d)) else { |
| 842 | return error(400, "DIGEST_INVALID", "Finish an upload with ?digest=sha256:<hex>."); |
| 843 | }; |
| 844 | let Some(progress) = row.progress() else { |
| 845 | return error(404, "BLOB_UPLOAD_INVALID", "This upload cannot be continued. Start again."); |
| 846 | }; |
| 847 | let writer = Writer::resume(&self.store, progress).await?; |
| 848 | // However it ends, the upload is over and its row goes. |
| 849 | let answer = match self.take_body(request, writer).await { |
| 850 | Ok(Ok(writer)) => self.complete(writer, name, package, &digest).await, |
| 851 | Ok(Err(refused)) => Ok(refused), |
| 852 | Err(error) => Err(error), |
| 853 | }; |
| 854 | self.db.delete_upload(&row.id).await?; |
| 855 | answer |
| 856 | } |
| 857 | |
| 858 | /// Ends an upload whose bytes are all in: checks the digest and the |
| 859 | /// workspace's free storage, and keeps the blob unless it was already |
| 860 | /// kept. Anything refused is let go. |
| 861 | async fn complete<S: BlobStore>(&self, writer: Writer<'_, S>, name: &ImageName, package: &PackageRow, digest: &Digest) -> Result<Response> { |
| 862 | let package_id = package.id.as_str(); |
| 863 | let refused = self.storage_refusal(package, &[(digest.to_string(), writer.received())]).await; |
| 864 | match refused { |
| 865 | Ok(None) => {} |
| 866 | Ok(Some(refused)) => { |
| 867 | writer.abandon().await?; |
| 868 | return error(403, "DENIED", refused); |
| 869 | } |
| 870 | Err(error) => { |
| 871 | let _ = writer.abandon().await; |
| 872 | return Err(error); |
| 873 | } |
| 874 | } |
| 875 | // Recorded and still in the store: these bytes are not needed. A |
| 876 | // blob whose object went missing is stored again. |
| 877 | let stored = match self.db.blob(digest).await { |
| 878 | Ok(Some(blob)) => self.store.head(&blob.object_key).await.map(|size| size.is_some()), |
| 879 | Ok(None) => Ok(false), |
| 880 | Err(error) => Err(error), |
| 881 | }; |
| 882 | let stored = match stored { |
| 883 | Ok(stored) => stored, |
| 884 | Err(error) => { |
| 885 | let _ = writer.abandon().await; |
| 886 | return Err(error); |
| 887 | } |
| 888 | }; |
| 889 | let now = now_ms(); |
| 890 | match writer.finish(digest, stored).await? { |
| 891 | Finished::Mismatch { actual } => { |
| 892 | return error(400, "DIGEST_INVALID", format!("The upload's digest is {actual}, not {digest}.")); |
| 893 | } |
| 894 | Finished::Duplicate { .. } => self.db.link_blob(package_id, digest, now).await?, |
| 895 | Finished::Stored { key, size } => self.db.keep_blob(package_id, digest, size, None, &key, now).await?, |
| 896 | } |
| 897 | empty(201, &[ |
| 898 | ("location", format!("/v2/{}/blobs/{digest}", name.full())), |
| 899 | ("docker-content-digest", digest.to_string()), |
| 900 | ]) |
| 901 | } |
| 902 | |
| 903 | async fn cancel_upload(&self, row: UploadRow) -> Result<Response> { |
| 904 | if let Some(progress) = row.progress() { |
| 905 | upload::abort(&self.store, &progress).await?; |
| 906 | } |
| 907 | self.db.delete_upload(&row.id).await?; |
| 908 | empty(204, &[]) |
| 909 | } |
| 910 | |
| 911 | async fn tags(&self, url: &Url, name: &ImageName, found: Option<PackageRow>) -> Result<Response> { |
| 912 | let Some(package) = found else { |
| 913 | return error(404, "NAME_UNKNOWN", format!("There is no image {}.", name.full())); |
| 914 | }; |
| 915 | let n = query(url, "n").and_then(|n| n.parse::<u32>().ok()).unwrap_or(MAX_TAGS_PAGE).min(MAX_TAGS_PAGE); |
| 916 | let last = query(url, "last"); |
| 917 | let tags = if n == 0 { Vec::new() } else { self.db.tag_names(&package.id, last.as_deref(), n).await? }; |
| 918 | let mut response = Response::from_json(&json!({ "name": name.full(), "tags": tags }))?; |
| 919 | if tags.len() as u32 == n |
| 920 | && let Some(final_tag) = tags.last() |
| 921 | { |
| 922 | response |
| 923 | .headers_mut() |
| 924 | .set("link", &format!("</v2/{}/tags/list?n={n}&last={final_tag}>; rel=\"next\"", name.full()))?; |
| 925 | } |
| 926 | Ok(response) |
| 927 | } |
| 928 | |
| 929 | async fn referrers(&self, url: &Url, name: &ImageName, found: Option<PackageRow>, digest: &str) -> Result<Response> { |
| 930 | let Some(digest) = Digest::parse(digest) else { |
| 931 | return error(400, "DIGEST_INVALID", format!("{digest} is not a sha256 digest.")); |
| 932 | }; |
| 933 | let Some(package) = found else { |
| 934 | return error(404, "NAME_UNKNOWN", format!("There is no image {}.", name.full())); |
| 935 | }; |
| 936 | let wanted = query(url, "artifactType"); |
| 937 | let mut manifests: Vec<Value> = Vec::new(); |
| 938 | for version in self.db.referrers(&package.id, digest.as_str()).await? { |
| 939 | let meta = version.meta(); |
| 940 | let artifact_type = meta["artifact_type"].as_str().map(str::to_owned); |
| 941 | if wanted.is_some() && artifact_type != wanted { |
| 942 | continue; |
| 943 | } |
| 944 | let size = self.db.blob(&Digest::parse(&version.digest).unwrap_or(digest.clone())).await?.map_or(0, |b| b.size); |
| 945 | let mut entry = json!({ |
| 946 | "mediaType": meta["media_type"].as_str().unwrap_or(manifest::OCI_MANIFEST), |
| 947 | "digest": version.digest, |
| 948 | "size": size, |
| 949 | }); |
| 950 | if let Some(artifact_type) = artifact_type { |
| 951 | entry["artifactType"] = json!(artifact_type); |
| 952 | } |
| 953 | if meta["annotations"].is_object() { |
| 954 | entry["annotations"] = meta["annotations"].clone(); |
| 955 | } |
| 956 | manifests.push(entry); |
| 957 | } |
| 958 | let body = json!({ "schemaVersion": 2, "mediaType": manifest::OCI_INDEX, "manifests": manifests }); |
| 959 | let mut response = Response::from_json(&body)?; |
| 960 | response.headers_mut().set("content-type", manifest::OCI_INDEX)?; |
| 961 | if wanted.is_some() { |
| 962 | response.headers_mut().set("oci-filters-applied", "artifactType")?; |
| 963 | } |
| 964 | Ok(response) |
| 965 | } |
| 966 | |
| 967 | /// Deletes a version of a package, with its tags, and says so. |
| 968 | pub(crate) async fn remove_version(&self, package: &PackageRow, version: &crate::db::VersionRow, caller: &Caller) -> Result<()> { |
| 969 | self.db.delete_version(&version.id).await?; |
| 970 | self.db.measure(&package.workspace).await?; |
| 971 | let event = PackageEvent { |
| 972 | version: Some(version.version.clone()), |
| 973 | digest: Some(version.digest.clone()), |
| 974 | ..self.event_of(package) |
| 975 | }; |
| 976 | self.announce("package.version_deleted", package, event, caller).await; |
| 977 | self.audit(caller, "package.delete_version", package, Some(&format!("{}/{}@{}", package.workspace, package.name, version.digest)), None) |
| 978 | .await; |
| 979 | Ok(()) |
| 980 | } |
| 981 | |
| 982 | pub(crate) fn event_of(&self, package: &PackageRow) -> PackageEvent { |
| 983 | PackageEvent { |
| 984 | package_id: package.id.clone(), |
| 985 | workspace: package.workspace.clone(), |
| 986 | ecosystem: Ecosystem::parse(&package.ecosystem).unwrap_or(Ecosystem::Container).as_str().to_owned(), |
| 987 | name: package.name.clone(), |
| 988 | repo_id: package.repo_id.clone(), |
| 989 | ..PackageEvent::default() |
| 990 | } |
| 991 | } |
| 992 | } |
| 993 | |
| 994 | fn upload_status(name: &ImageName, row: &UploadRow) -> Result<Response> { |
| 995 | upload_status_with(name, &row.id, row.offset, 204) |
| 996 | } |
| 997 | |
| 998 | fn upload_status_with(name: &ImageName, id: &str, offset: u64, status: u16) -> Result<Response> { |
| 999 | empty(status, &[ |
| 1000 | ("location", format!("/v2/{}/blobs/uploads/{id}", name.full())), |
| 1001 | ("range", range::upload_range(offset)), |
| 1002 | ("docker-upload-uuid", id.to_owned()), |
| 1003 | ("content-length", "0".to_owned()), |
| 1004 | ]) |
| 1005 | } |
| 1006 | |
| 1007 | #[cfg(test)] |
| 1008 | mod tests { |
| 1009 | use super::*; |
| 1010 | |
| 1011 | #[test] |
| 1012 | fn the_challenge_names_the_token_endpoint_on_the_same_host() { |
| 1013 | let url = Url::parse("https://g1t.sh/v2/").unwrap(); |
| 1014 | assert_eq!(challenge(&url, None), "Bearer realm=\"https://g1t.sh/v2/token\",service=\"g1t.sh\""); |
| 1015 | let local = Url::parse("http://localhost:8790/v2/acme/web/manifests/latest").unwrap(); |
| 1016 | assert_eq!( |
| 1017 | challenge(&local, Some("repository:acme/web:pull".into())), |
| 1018 | "Bearer realm=\"http://localhost:8790/v2/token\",service=\"localhost:8790\",scope=\"repository:acme/web:pull\"" |
| 1019 | ); |
| 1020 | } |
| 1021 | } |