Pick any line to see why it is the way it is: the commit, the pull request and issue it came from, and what the agent was thinking.
| Packages, with a container registry on g1t.sh; workspaces deleted whole and kept 30 days; Members for every member | 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. | |
| npm on g1t.sh: publish and install @<workspace>/<name> with the npm CLI and a g1t token | 60 | pub(crate) enum Credentials { |
| Packages, with a container registry on g1t.sh; workspaces deleted whole and kept 30 days; Members for every member | 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. | |
| npm on g1t.sh: publish and install @<workspace>/<name> with the npm CLI and a g1t token | 77 | pub(crate) fn origin(url: &Url) -> String { |
| Packages, with a container registry on g1t.sh; workspaces deleted whole and kept 30 days; Members for every member | 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. | |
| npm on g1t.sh: publish and install @<workspace>/<name> with the npm CLI and a g1t token | 112 | pub(crate) fn published_by(caller: &Caller) -> Option<String> { |
| Packages, with a container registry on g1t.sh; workspaces deleted whole and kept 30 days; Members for every member | 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) | |
| npm on g1t.sh: publish and install @<workspace>/<name> with the npm CLI and a g1t token | 149 | && let Some(refused) = self.limited(&request, &credentials, "`docker login g1t.sh`").await? |
| Packages, with a container registry on g1t.sh; workspaces deleted whole and kept 30 days; Members for every member | 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. | |
| npm on g1t.sh: publish and install @<workspace>/<name> with the npm CLI and a g1t token | 271 | /// `sign_in` is how a client of this registry signs in, for the message. |
| 272 | pub(crate) async fn limited(&self, request: &Request, credentials: &Credentials, sign_in: &str) -> Result<Option<Response>> { | |
| Packages, with a container registry on g1t.sh; workspaces deleted whole and kept 30 days; Members for every member | 273 | let subject = match credentials { |
| 274 | Credentials::Token(claims) => claims.actor.as_ref().map(|actor| actor.actor_id.clone()), | |
| 275 | Credentials::Viewer(Some(user)) => Some(user.id.clone()), | |
| 276 | _ => None, | |
| 277 | }; | |
| 278 | let address = request.headers().get("cf-connecting-ip")?; | |
| 279 | let (limit, key) = limits::key(subject.as_deref(), address.as_deref()); | |
| 280 | let Ok(limiter) = self.env.rate_limiter(limit.binding()) else { | |
| 281 | return Ok(None); | |
| 282 | }; | |
| 283 | match limiter.limit(key).await { | |
| 284 | Ok(outcome) if !outcome.success => { | |
| 285 | let mut response = error( | |
| 286 | 429, | |
| 287 | "TOOMANYREQUESTS", | |
| 288 | match limit { | |
| npm on g1t.sh: publish and install @<workspace>/<name> with the npm CLI and a g1t token | 289 | limits::Limit::Anonymous => format!("Too many requests from this address. Wait a minute, or sign in with {sign_in} for a higher limit."), |
| 290 | limits::Limit::Signed => "Too many requests. Wait a minute and try again.".to_owned(), | |
| Packages, with a container registry on g1t.sh; workspaces deleted whole and kept 30 days; Members for every member | 291 | }, |
| 292 | )?; | |
| 293 | response.headers_mut().set("retry-after", &limits::RETRY_AFTER_SECONDS.to_string())?; | |
| 294 | Ok(Some(response)) | |
| 295 | } | |
| 296 | Ok(_) => Ok(None), | |
| 297 | Err(problem) => { | |
| 298 | worker::console_error!("packages: the rate limit could not be asked: {problem}"); | |
| 299 | Ok(None) | |
| 300 | } | |
| 301 | } | |
| 302 | } | |
| 303 | ||
| npm on g1t.sh: publish and install @<workspace>/<name> with the npm CLI and a g1t token | 304 | pub(crate) async fn credentials(&self, request: &Request) -> Result<Credentials> { |
| Packages, with a container registry on g1t.sh; workspaces deleted whole and kept 30 days; Members for every member | 305 | let Some(header) = request.headers().get("authorization")? else { |
| 306 | return Ok(Credentials::None); | |
| 307 | }; | |
| 308 | if let Some(bearer) = token::bearer(&header) { | |
| 309 | if token::is_registry_token(bearer) { | |
| 310 | return Ok(match token::verify(bearer, &self.secret, now_ms() / 1000) { | |
| 311 | Some(claims) => Credentials::Token(claims), | |
| 312 | None => Credentials::Bad, | |
| 313 | }); | |
| 314 | } | |
| 315 | return Ok(match self.viewer_for("token", bearer).await? { | |
| 316 | Some(user) => Credentials::Viewer(Some(user)), | |
| 317 | None => Credentials::Bad, | |
| 318 | }); | |
| 319 | } | |
| 320 | if let Some((username, secret)) = token::basic(&header) { | |
| 321 | return Ok(match self.viewer_for(&username, &secret).await? { | |
| 322 | Some(user) => Credentials::Viewer(Some(user)), | |
| 323 | None => Credentials::Bad, | |
| 324 | }); | |
| 325 | } | |
| 326 | Ok(Credentials::Bad) | |
| 327 | } | |
| 328 | ||
| 329 | /// Whether the request may do `action` to the image, and as whom. | |
| 330 | async fn authorize( | |
| 331 | &self, | |
| 332 | url: &Url, | |
| 333 | credentials: &Credentials, | |
| 334 | name: &ImageName, | |
| 335 | found: Option<&PackageRow>, | |
| 336 | action: Action, | |
| 337 | ) -> Result<std::result::Result<Caller, Response>> { | |
| 338 | let scope = Some(format!("repository:{}:{}", name.full(), match action { | |
| 339 | Action::Pull => "pull", | |
| 340 | Action::Push => "pull,push", | |
| 341 | Action::Delete => "delete", | |
| 342 | })); | |
| 343 | let viewer = match credentials { | |
| 344 | Credentials::Bad => { | |
| 345 | return Ok(Err(unauthorized(url, scope, "The token or password is not right, or has expired. Sign in again with `docker login`.")?)); | |
| 346 | } | |
| 347 | Credentials::Token(claims) => { | |
| 348 | if claims.allows(&name.full(), action) { | |
| 349 | return Ok(Ok(Caller { actor: claims.actor.clone() })); | |
| 350 | } | |
| 351 | // A token for something else still pulls a public image. | |
| 352 | if action == Action::Pull && found.is_some_and(|package| package.public()) { | |
| 353 | return Ok(Ok(Caller { actor: claims.actor.clone() })); | |
| 354 | } | |
| 355 | if credentials.anonymous() { | |
| 356 | return Ok(Err(unauthorized(url, scope, "Sign in with `docker login` to do that.")?)); | |
| 357 | } | |
| 358 | return Ok(Err(error(403, "DENIED", format!("This token may not {} {}.", action.as_str(), name.full()))?)); | |
| 359 | } | |
| 360 | Credentials::Viewer(viewer) => viewer.clone(), | |
| 361 | Credentials::None => None, | |
| 362 | }; | |
| 363 | let target = self.target(name, found).await?; | |
| 364 | let decision = access::decide(viewer.as_ref(), &target.view(), action); | |
| 365 | if decision.allowed { | |
| 366 | return Ok(Ok(Caller { actor: viewer.as_ref().map(AuditActor::of) })); | |
| 367 | } | |
| 368 | let reason = decision.reason.unwrap_or_else(|| "Not allowed.".to_owned()); | |
| 369 | if viewer.is_none() { | |
| 370 | return Ok(Err(unauthorized(url, scope, &reason)?)); | |
| 371 | } | |
| 372 | Ok(Err(error(403, "DENIED", reason)?)) | |
| 373 | } | |
| 374 | ||
| 375 | /// `GET /v2/token`: a token for each scope asked for, cut down to what | |
| 376 | /// the credentials may do. | |
| 377 | async fn issue(&self, url: &Url, credentials: Credentials, asked: &[String]) -> Result<Response> { | |
| 378 | let viewer = match credentials { | |
| 379 | Credentials::Viewer(viewer) => viewer, | |
| 380 | Credentials::None => None, | |
| 381 | // A registry token is not a way to get another. | |
| 382 | Credentials::Token(_) | Credentials::Bad => { | |
| 383 | return unauthorized(url, None, "The username or token is not right. Use a g1t token as the password."); | |
| 384 | } | |
| 385 | }; | |
| 386 | let mut access: Vec<Grant> = Vec::new(); | |
| 387 | let scopes: Vec<String> = asked.iter().flat_map(|value| value.split(' ').map(str::to_owned)).collect(); | |
| 388 | for scope in scopes.iter().take(20) { | |
| 389 | let Some((full, wanted)) = token::parse_scope(scope) else { continue }; | |
| 390 | let Ok(name) = names::parse_name(&full) else { continue }; | |
| 391 | let found = self.db.package(&name.workspace, CONTAINER, &name.name).await?; | |
| 392 | // Nothing is granted on a deleted workspace's images. | |
| 393 | if found.as_ref().is_some_and(|package| package.hidden()) { | |
| 394 | continue; | |
| 395 | } | |
| 396 | let target = self.target(&name, found.as_ref()).await?; | |
| 397 | let actions: Vec<Action> = wanted | |
| 398 | .into_iter() | |
| 399 | .filter(|action| access::decide(viewer.as_ref(), &target.view(), *action).allowed) | |
| 400 | .collect(); | |
| 401 | if let Some(grant) = access.iter_mut().find(|grant| grant.name == full) { | |
| 402 | for action in actions { | |
| 403 | if !grant.actions.contains(&action) { | |
| 404 | grant.actions.push(action); | |
| 405 | } | |
| 406 | } | |
| 407 | } else { | |
| 408 | access.push(Grant { name: full, actions }); | |
| 409 | } | |
| 410 | } | |
| 411 | let now = now_ms(); | |
| 412 | let claims = Claims { | |
| 413 | actor: viewer.as_ref().map(AuditActor::of), | |
| 414 | access, | |
| 415 | iat: now / 1000, | |
| 416 | exp: now / 1000 + token::TTL_SECONDS, | |
| 417 | }; | |
| 418 | let signed = token::sign(&claims, &self.secret); | |
| 419 | Response::from_json(&json!({ | |
| 420 | "token": signed, | |
| 421 | "access_token": signed, | |
| 422 | "expires_in": token::TTL_SECONDS, | |
| 423 | "issued_at": rfc3339(now), | |
| 424 | })) | |
| 425 | } | |
| 426 | ||
| 427 | /// The package for an image, made on its first push and linked to the | |
| 428 | /// repository of the image's name when there is one. | |
| 429 | async fn package_for_push(&self, name: &ImageName, found: Option<PackageRow>, caller: &Caller) -> Result<PackageRow> { | |
| 430 | if let Some(found) = found { | |
| 431 | return Ok(found); | |
| 432 | } | |
| 433 | let repo = self.repo_by_name(&name.workspace, name.repo_name()).await?; | |
| 434 | self.db | |
| 435 | .create_package( | |
| 436 | &new_id("pkg", now_ms()), | |
| 437 | &name.workspace, | |
| 438 | CONTAINER, | |
| 439 | &name.name, | |
| 440 | repo.as_ref().map(|repo| (repo.id.as_str(), repo.name.as_str(), repo.is_private)), | |
| 441 | caller.actor.as_ref().map_or("", |actor| actor.actor_id.as_str()), | |
| 442 | now_ms(), | |
| 443 | ) | |
| 444 | .await | |
| 445 | } | |
| 446 | ||
| 447 | async fn get_manifest( | |
| 448 | &self, | |
| 449 | name: &ImageName, | |
| 450 | found: Option<PackageRow>, | |
| 451 | reference: &str, | |
| 452 | head: bool, | |
| 453 | ctx: &Context, | |
| 454 | ) -> Result<Response> { | |
| 455 | let Some(package) = found else { | |
| 456 | return error(404, "NAME_UNKNOWN", format!("There is no image {}.", name.full())); | |
| 457 | }; | |
| 458 | let version = match Reference::parse(reference) { | |
| 459 | Some(Reference::Tag(tag)) => self.db.version_by_tag(&package.id, &tag).await?, | |
| 460 | Some(Reference::Digest(digest)) => self.db.version_by_digest(&package.id, digest.as_str()).await?, | |
| 461 | None => return error(400, "MANIFEST_INVALID", format!("{reference} is not a tag or a digest.")), | |
| 462 | }; | |
| 463 | let Some(version) = version else { | |
| 464 | return error(404, "MANIFEST_UNKNOWN", format!("{}:{reference} is not there.", name.full())); | |
| 465 | }; | |
| 466 | let digest = Digest::parse(&version.digest).ok_or_else(|| worker::Error::RustError("a stored digest is malformed".into()))?; | |
| 467 | let Some(blob) = self.db.blob(&digest).await? else { | |
| 468 | return error(404, "MANIFEST_UNKNOWN", format!("{}:{reference} is not there.", name.full())); | |
| 469 | }; | |
| 470 | let headers = [ | |
| 471 | ("content-type", version.media_type().unwrap_or_else(|| manifest::OCI_MANIFEST.to_owned())), | |
| 472 | ("docker-content-digest", version.digest.clone()), | |
| 473 | ("etag", format!("\"{}\"", version.digest)), | |
| 474 | ("content-length", blob.size.to_string()), | |
| 475 | ]; | |
| 476 | if head { | |
| 477 | return empty(200, &headers); | |
| 478 | } | |
| 479 | let Some(bytes) = self.store.read(&blob.object_key).await? else { | |
| 480 | return error(404, "MANIFEST_UNKNOWN", format!("{}:{reference} is not there.", name.full())); | |
| 481 | }; | |
| 482 | self.count_download(&package.id, ctx); | |
| 483 | respond(200, &headers, ResponseBody::Body(bytes)) | |
| 484 | } | |
| 485 | ||
| 486 | async fn put_manifest( | |
| 487 | &self, | |
| 488 | request: &mut Request, | |
| 489 | name: &ImageName, | |
| 490 | found: Option<PackageRow>, | |
| 491 | reference: &str, | |
| 492 | caller: &Caller, | |
| 493 | ) -> Result<Response> { | |
| 494 | let reference = match Reference::parse(reference) { | |
| 495 | Some(reference) => reference, | |
| 496 | None => return error(400, "TAG_INVALID", format!("{reference} is not a valid tag.")), | |
| 497 | }; | |
| 498 | let bytes = request.bytes().await?; | |
| 499 | if bytes.len() > manifest::MAX_MANIFEST_BYTES { | |
| 500 | return error(413, "SIZE_INVALID", "A manifest is at most 4 MiB."); | |
| 501 | } | |
| 502 | let digest = Digest::of(&bytes); | |
| 503 | if let Reference::Digest(given) = &reference | |
| 504 | && given != &digest | |
| 505 | { | |
| 506 | return error(400, "DIGEST_INVALID", format!("The manifest's digest is {digest}, not {given}.")); | |
| 507 | } | |
| 508 | let content_type = request.headers().get("content-type")?; | |
| 509 | let parsed = match manifest::parse(&bytes, content_type.as_deref()) { | |
| 510 | Ok(parsed) => parsed, | |
| 511 | Err(manifest::Refused::Invalid(message)) => return error(400, "MANIFEST_INVALID", message), | |
| 512 | Err(manifest::Refused::Unsupported(message)) => return error(415, "UNSUPPORTED", message), | |
| 513 | }; | |
| 514 | let package = self.package_for_push(name, found, caller).await?; | |
| 515 | let mut files = vec![NewFile { | |
| 516 | name: "manifest".to_owned(), | |
| 517 | digest: digest.to_string(), | |
| 518 | size: bytes.len() as u64, | |
| 519 | media_type: Some(parsed.media_type.clone()), | |
| 520 | }]; | |
| 521 | match parsed.kind { | |
| 522 | Kind::Image => { | |
| 523 | let mut layer = 0; | |
| 524 | for blob in &parsed.blobs { | |
| 525 | let Some(stored) = self.db.package_blob(&package.id, &blob.digest).await? else { | |
| 526 | return error(400, "MANIFEST_BLOB_UNKNOWN", format!("Blob {} is not in {}: push it first.", blob.digest, name.full())); | |
| 527 | }; | |
| 528 | let file_name = if blob.role == "config" { | |
| 529 | "config".to_owned() | |
| 530 | } else { | |
| 531 | layer += 1; | |
| 532 | format!("layer:{layer}") | |
| 533 | }; | |
| 534 | files.push(NewFile { name: file_name, digest: blob.digest.to_string(), size: stored.size, media_type: blob.media_type.clone() }); | |
| 535 | } | |
| 536 | } | |
| 537 | Kind::Index => { | |
| 538 | for child in &parsed.manifests { | |
| 539 | if self.db.version_by_digest(&package.id, child.digest.as_str()).await?.is_none() { | |
| 540 | return error(400, "MANIFEST_UNKNOWN", format!("Manifest {} is not in {}: push it first.", child.digest, name.full())); | |
| 541 | } | |
| 542 | } | |
| 543 | } | |
| 544 | } | |
| 545 | let pushed: Vec<(String, u64)> = files.iter().map(|file| (file.digest.clone(), file.size)).collect(); | |
| 546 | if let Some(refused) = self.storage_refusal(&package, &pushed).await? { | |
| 547 | return error(403, "DENIED", refused); | |
| 548 | } | |
| 549 | if self.db.blob(&digest).await?.is_none() { | |
| 550 | self.store.put(&digest.object_key(), bytes.clone()).await?; | |
| 551 | } | |
| 552 | let now = now_ms(); | |
| 553 | self.db | |
| 554 | .keep_blob(&package.id, &digest, bytes.len() as u64, Some(&parsed.media_type), &digest.object_key(), now) | |
| 555 | .await?; | |
| 556 | let size = files.iter().map(|file| file.size).sum(); | |
| 557 | let metadata = json!({ | |
| 558 | "media_type": parsed.media_type, | |
| 559 | "artifact_type": parsed.artifact_type, | |
| 560 | "annotations": parsed.annotations, | |
| 561 | "platforms": parsed.platforms, | |
| 562 | }); | |
| 563 | let tag = match &reference { | |
| 564 | Reference::Tag(tag) => Some(tag.as_str()), | |
| 565 | Reference::Digest(_) => None, | |
| 566 | }; | |
| 567 | let (version, changed) = self | |
| 568 | .db | |
| 569 | .publish( | |
| 570 | NewVersion { | |
| 571 | id: new_id("ver", now), | |
| 572 | package_id: package.id.clone(), | |
| 573 | version: digest.to_string(), | |
| 574 | digest: digest.to_string(), | |
| 575 | size, | |
| 576 | metadata: metadata.to_string(), | |
| 577 | subject: parsed.subject.as_ref().map(Digest::to_string), | |
| 578 | published_by: published_by(caller), | |
| 579 | files, | |
| 580 | }, | |
| 581 | tag, | |
| 582 | now, | |
| 583 | ) | |
| 584 | .await?; | |
| 585 | self.db.measure(&package.workspace).await?; | |
| 586 | if changed { | |
| 587 | let tags: Vec<String> = self | |
| 588 | .db | |
| 589 | .tags(&package.id) | |
| 590 | .await? | |
| 591 | .into_iter() | |
| 592 | .filter(|row| row.version_id == version.id) | |
| 593 | .map(|row| row.tag) | |
| 594 | .collect(); | |
| 595 | let event = PackageEvent { | |
| 596 | version: Some(version.version.clone()), | |
| 597 | digest: Some(version.digest.clone()), | |
| 598 | size: Some(version.size), | |
| 599 | tags: Some(tags), | |
| 600 | ..self.event_of(&package) | |
| 601 | }; | |
| 602 | self.announce("package.published", &package, event, caller).await; | |
| 603 | self.audit(caller, "package.publish", &package, Some(&format!("{}@{}", name.full(), version.digest)), None) | |
| 604 | .await; | |
| 605 | } | |
| 606 | let mut headers = vec![ | |
| 607 | ("location", format!("/v2/{}/manifests/{digest}", name.full())), | |
| 608 | ("docker-content-digest", digest.to_string()), | |
| 609 | ]; | |
| 610 | if let Some(subject) = &parsed.subject { | |
| 611 | headers.push(("oci-subject", subject.to_string())); | |
| 612 | } | |
| 613 | empty(201, &headers) | |
| 614 | } | |
| 615 | ||
| 616 | async fn delete_manifest(&self, name: &ImageName, found: Option<PackageRow>, reference: &str, caller: &Caller) -> Result<Response> { | |
| 617 | let Some(package) = found else { | |
| 618 | return error(404, "NAME_UNKNOWN", format!("There is no image {}.", name.full())); | |
| 619 | }; | |
| 620 | match Reference::parse(reference) { | |
| 621 | Some(Reference::Tag(tag)) => { | |
| 622 | if self.db.version_by_tag(&package.id, &tag).await?.is_none() { | |
| 623 | return error(404, "MANIFEST_UNKNOWN", format!("{}:{tag} is not there.", name.full())); | |
| 624 | } | |
| 625 | self.db.delete_tag(&package.id, &tag).await?; | |
| 626 | self.audit(caller, "package.untag", &package, Some(&format!("{}:{tag}", name.full())), None).await; | |
| 627 | empty(202, &[]) | |
| 628 | } | |
| 629 | Some(Reference::Digest(digest)) => { | |
| 630 | let Some(version) = self.db.version_by_digest(&package.id, digest.as_str()).await? else { | |
| 631 | return error(404, "MANIFEST_UNKNOWN", format!("{}@{digest} is not there.", name.full())); | |
| 632 | }; | |
| 633 | self.remove_version(&package, &version, caller).await?; | |
| 634 | empty(202, &[]) | |
| 635 | } | |
| 636 | None => error(400, "MANIFEST_INVALID", format!("{reference} is not a tag or a digest.")), | |
| 637 | } | |
| 638 | } | |
| 639 | ||
| 640 | async fn get_blob(&self, request: &Request, package: &PackageRow, digest: &Digest, head: bool) -> Result<Response> { | |
| 641 | let Some(blob) = self.db.package_blob(&package.id, digest).await? else { | |
| 642 | return error(404, "BLOB_UNKNOWN", format!("Blob {digest} is not in this image.")); | |
| 643 | }; | |
| 644 | let mut headers = vec![ | |
| 645 | ("docker-content-digest", digest.to_string()), | |
| 646 | ("content-type", "application/octet-stream".to_owned()), | |
| 647 | ("accept-ranges", "bytes".to_owned()), | |
| 648 | ("etag", format!("\"{digest}\"")), | |
| 649 | ("cache-control", "max-age=31536000".to_owned()), | |
| 650 | ]; | |
| 651 | if head { | |
| 652 | headers.push(("content-length", blob.size.to_string())); | |
| 653 | return empty(200, &headers); | |
| 654 | } | |
| 655 | let wanted = match range::parse_range(request.headers().get("range")?.as_deref(), blob.size) { | |
| 656 | Ok(wanted) => wanted, | |
| 657 | Err(()) => { | |
| 658 | headers.push(("content-range", format!("bytes */{}", blob.size))); | |
| 659 | return respond(416, &headers, ResponseBody::Empty); | |
| 660 | } | |
| 661 | }; | |
| 662 | if wanted.is_none() | |
| 663 | && blob.size > REDIRECT_BYTES | |
| 664 | && let Some(url) = self.store.presign_get(&blob.object_key, REDIRECT_SECONDS, now_ms()) | |
| 665 | { | |
| 666 | headers.push(("location", url)); | |
| 667 | return empty(307, &headers); | |
| 668 | } | |
| 669 | let Some(got) = self.store.get(&blob.object_key, wanted).await? else { | |
| 670 | return error(404, "BLOB_UNKNOWN", format!("Blob {digest} is not in this image.")); | |
| 671 | }; | |
| 672 | match wanted { | |
| 673 | Some(wanted) => { | |
| 674 | headers.push(("content-range", wanted.content_range(got.size))); | |
| 675 | headers.push(("content-length", wanted.length.to_string())); | |
| 676 | respond(206, &headers, got.body) | |
| 677 | } | |
| 678 | None => { | |
| 679 | headers.push(("content-length", got.size.to_string())); | |
| 680 | respond(200, &headers, got.body) | |
| 681 | } | |
| 682 | } | |
| 683 | } | |
| 684 | ||
| 685 | async fn delete_blob(&self, package: &PackageRow, digest: &Digest) -> Result<Response> { | |
| 686 | if self.db.package_blob(&package.id, digest).await?.is_none() { | |
| 687 | return error(404, "BLOB_UNKNOWN", format!("Blob {digest} is not in this image.")); | |
| 688 | } | |
| 689 | if self.db.blob_in_use(&package.id, digest).await? { | |
| 690 | return error(405, "DENIED", format!("Blob {digest} is used by a manifest of this image; delete the manifest instead.")); | |
| 691 | } | |
| 692 | self.db.unlink_blob(&package.id, digest).await?; | |
| 693 | empty(202, &[]) | |
| 694 | } | |
| 695 | ||
| 696 | async fn start_upload( | |
| 697 | &self, | |
| 698 | request: Request, | |
| 699 | url: &Url, | |
| 700 | name: &ImageName, | |
| 701 | found: Option<PackageRow>, | |
| 702 | credentials: &Credentials, | |
| 703 | caller: &Caller, | |
| 704 | ) -> Result<Response> { | |
| 705 | let package = self.package_for_push(name, found, caller).await?; | |
| 706 | // A blob another image of the workspace has: linked, not copied. | |
| 707 | if let (Some(mount), Some(from)) = (query(url, "mount"), query(url, "from")) | |
| 708 | && let (Some(digest), Ok(from)) = (Digest::parse(&mount), names::parse_name(&from)) | |
| 709 | && from.workspace == name.workspace | |
| 710 | && let Some(source) = self.db.package(&from.workspace, CONTAINER, &from.name).await? | |
| 711 | && self.may_pull(credentials, &from, &source).await? | |
| 712 | && let Some(blob) = self.db.package_blob(&source.id, &digest).await? | |
| 713 | { | |
| 714 | if let Some(refused) = self.storage_refusal(&package, &[(digest.to_string(), blob.size)]).await? { | |
| 715 | return error(403, "DENIED", refused); | |
| 716 | } | |
| 717 | self.db.link_blob(&package.id, &digest, now_ms()).await?; | |
| 718 | return empty(201, &[ | |
| 719 | ("location", format!("/v2/{}/blobs/{digest}", name.full())), | |
| 720 | ("docker-content-digest", digest.to_string()), | |
| 721 | ]); | |
| 722 | } | |
| 723 | let id = new_id("upl", now_ms()); | |
| 724 | let row = UploadRow { | |
| 725 | id: id.clone(), | |
| 726 | workspace: name.workspace.clone(), | |
| 727 | package_id: package.id.clone(), | |
| 728 | package: name.full(), | |
| 729 | multipart_id: None, | |
| 730 | parts: "[]".to_owned(), | |
| 731 | offset: 0, | |
| 732 | tail: 0, | |
| 733 | hash_state: Sha256::new().save(), | |
| 734 | }; | |
| 735 | // The whole blob in this one request. | |
| 736 | if let Some(digest) = query(url, "digest") { | |
| 737 | let Some(digest) = Digest::parse(&digest) else { | |
| 738 | return error(400, "DIGEST_INVALID", format!("{digest} is not a sha256 digest.")); | |
| 739 | }; | |
| 740 | let writer = Writer::resume(&self.store, Progress::new(&id)).await?; | |
| 741 | let writer = match self.take_body(request, writer).await? { | |
| 742 | Ok(writer) => writer, | |
| 743 | Err(refused) => return Ok(refused), | |
| 744 | }; | |
| 745 | return self.complete(writer, name, &package, &digest).await; | |
| 746 | } | |
| 747 | self.db.create_upload(&row, now_ms()).await?; | |
| 748 | upload_status_with(name, &id, 0, 202) | |
| 749 | } | |
| 750 | ||
| 751 | /// Whether the credentials may pull another image, for a mount. | |
| 752 | async fn may_pull(&self, credentials: &Credentials, name: &ImageName, package: &PackageRow) -> Result<bool> { | |
| 753 | Ok(match credentials { | |
| 754 | Credentials::Token(claims) => claims.allows(&name.full(), Action::Pull) || package.public(), | |
| 755 | Credentials::Viewer(viewer) => { | |
| 756 | let target = TargetOf::package(package); | |
| 757 | access::decide(viewer.as_ref(), &target.view(), Action::Pull).allowed | |
| 758 | } | |
| 759 | Credentials::None | Credentials::Bad => package.public(), | |
| 760 | }) | |
| 761 | } | |
| 762 | ||
| 763 | /// Writes the request's body into the upload, refusing one over the | |
| 764 | /// limit. On a refusal or a failure the whole upload is let go, the | |
| 765 | /// parts this request sent and the tail an earlier one kept included, | |
| 766 | /// so nothing is left in the store; the caller forgets its row. | |
| 767 | async fn take_body<'a, S: BlobStore>( | |
| 768 | &self, | |
| 769 | mut request: Request, | |
| 770 | mut writer: Writer<'a, S>, | |
| 771 | ) -> Result<std::result::Result<Writer<'a, S>, Response>> { | |
| 772 | let declared = request.headers().get("content-length")?.and_then(|n| n.parse::<u64>().ok()); | |
| 773 | if declared.is_some_and(|n| n > self.max_request) { | |
| 774 | writer.abandon().await?; | |
| 775 | return Ok(Err(too_large(self.max_request)?)); | |
| 776 | } | |
| 777 | let mut received = 0u64; | |
| 778 | let mut stream = match request.stream() { | |
| 779 | Ok(stream) => stream, | |
| 780 | // No body at all. | |
| 781 | Err(_) => return Ok(Ok(writer)), | |
| 782 | }; | |
| 783 | while let Some(chunk) = stream.next().await { | |
| 784 | let written = match chunk { | |
| 785 | Ok(chunk) => { | |
| 786 | received += chunk.len() as u64; | |
| 787 | if received > self.max_request { | |
| 788 | writer.abandon().await?; | |
| 789 | return Ok(Err(too_large(self.max_request)?)); | |
| 790 | } | |
| 791 | writer.write(&chunk).await | |
| 792 | } | |
| 793 | Err(error) => Err(error), | |
| 794 | }; | |
| 795 | if let Err(error) = written { | |
| 796 | let _ = writer.abandon().await; | |
| 797 | return Err(error); | |
| 798 | } | |
| 799 | } | |
| 800 | Ok(Ok(writer)) | |
| 801 | } | |
| 802 | ||
| 803 | async fn patch_upload(&self, request: Request, name: &ImageName, row: UploadRow) -> Result<Response> { | |
| 804 | let Some(progress) = row.progress() else { | |
| 805 | return error(404, "BLOB_UPLOAD_INVALID", "This upload cannot be continued. Start again."); | |
| 806 | }; | |
| 807 | if let Some(header) = request.headers().get("content-range")? { | |
| 808 | match range::parse_content_range(&header) { | |
| 809 | Some(chunk) if chunk.start == progress.offset => {} | |
| 810 | _ => { | |
| 811 | return empty(416, &[ | |
| 812 | ("location", format!("/v2/{}/blobs/uploads/{}", name.full(), row.id)), | |
| 813 | ("range", range::upload_range(progress.offset)), | |
| 814 | ("docker-upload-uuid", row.id.clone()), | |
| 815 | ]); | |
| 816 | } | |
| 817 | } | |
| 818 | } | |
| 819 | let writer = Writer::resume(&self.store, progress).await?; | |
| 820 | // A refused or failed chunk ends the upload: the store holds | |
| 821 | // nothing of it any more (take_body, pause), so neither does its row. | |
| 822 | let progress = match self.take_body(request, writer).await { | |
| 823 | Ok(Ok(writer)) => writer.pause().await, | |
| 824 | Ok(Err(refused)) => { | |
| 825 | self.db.delete_upload(&row.id).await?; | |
| 826 | return Ok(refused); | |
| 827 | } | |
| 828 | Err(error) => Err(error), | |
| 829 | }; | |
| 830 | let progress = match progress { | |
| 831 | Ok(progress) => progress, | |
| 832 | Err(error) => { | |
| 833 | self.db.delete_upload(&row.id).await?; | |
| 834 | return Err(error); | |
| 835 | } | |
| 836 | }; | |
| 837 | self.db.save_progress(&progress, now_ms()).await?; | |
| 838 | upload_status_with(name, &row.id, progress.offset, 202) | |
| 839 | } | |
| 840 | ||
| 841 | async fn finish_upload(&self, request: Request, url: &Url, name: &ImageName, package: &PackageRow, row: UploadRow) -> Result<Response> { | |
| 842 | let Some(digest) = query(url, "digest").and_then(|d| Digest::parse(&d)) else { | |
| 843 | return error(400, "DIGEST_INVALID", "Finish an upload with ?digest=sha256:<hex>."); | |
| 844 | }; | |
| 845 | let Some(progress) = row.progress() else { | |
| 846 | return error(404, "BLOB_UPLOAD_INVALID", "This upload cannot be continued. Start again."); | |
| 847 | }; | |
| 848 | let writer = Writer::resume(&self.store, progress).await?; | |
| 849 | // However it ends, the upload is over and its row goes. | |
| 850 | let answer = match self.take_body(request, writer).await { | |
| 851 | Ok(Ok(writer)) => self.complete(writer, name, package, &digest).await, | |
| 852 | Ok(Err(refused)) => Ok(refused), | |
| 853 | Err(error) => Err(error), | |
| 854 | }; | |
| 855 | self.db.delete_upload(&row.id).await?; | |
| 856 | answer | |
| 857 | } | |
| 858 | ||
| 859 | /// Ends an upload whose bytes are all in: checks the digest and the | |
| 860 | /// workspace's free storage, and keeps the blob unless it was already | |
| 861 | /// kept. Anything refused is let go. | |
| 862 | async fn complete<S: BlobStore>(&self, writer: Writer<'_, S>, name: &ImageName, package: &PackageRow, digest: &Digest) -> Result<Response> { | |
| 863 | let package_id = package.id.as_str(); | |
| 864 | let refused = self.storage_refusal(package, &[(digest.to_string(), writer.received())]).await; | |
| 865 | match refused { | |
| 866 | Ok(None) => {} | |
| 867 | Ok(Some(refused)) => { | |
| 868 | writer.abandon().await?; | |
| 869 | return error(403, "DENIED", refused); | |
| 870 | } | |
| 871 | Err(error) => { | |
| 872 | let _ = writer.abandon().await; | |
| 873 | return Err(error); | |
| 874 | } | |
| 875 | } | |
| 876 | // Recorded and still in the store: these bytes are not needed. A | |
| 877 | // blob whose object went missing is stored again. | |
| 878 | let stored = match self.db.blob(digest).await { | |
| 879 | Ok(Some(blob)) => self.store.head(&blob.object_key).await.map(|size| size.is_some()), | |
| 880 | Ok(None) => Ok(false), | |
| 881 | Err(error) => Err(error), | |
| 882 | }; | |
| 883 | let stored = match stored { | |
| 884 | Ok(stored) => stored, | |
| 885 | Err(error) => { | |
| 886 | let _ = writer.abandon().await; | |
| 887 | return Err(error); | |
| 888 | } | |
| 889 | }; | |
| 890 | let now = now_ms(); | |
| 891 | match writer.finish(digest, stored).await? { | |
| 892 | Finished::Mismatch { actual } => { | |
| 893 | return error(400, "DIGEST_INVALID", format!("The upload's digest is {actual}, not {digest}.")); | |
| 894 | } | |
| 895 | Finished::Duplicate { .. } => self.db.link_blob(package_id, digest, now).await?, | |
| 896 | Finished::Stored { key, size } => self.db.keep_blob(package_id, digest, size, None, &key, now).await?, | |
| 897 | } | |
| 898 | empty(201, &[ | |
| 899 | ("location", format!("/v2/{}/blobs/{digest}", name.full())), | |
| 900 | ("docker-content-digest", digest.to_string()), | |
| 901 | ]) | |
| 902 | } | |
| 903 | ||
| 904 | async fn cancel_upload(&self, row: UploadRow) -> Result<Response> { | |
| 905 | if let Some(progress) = row.progress() { | |
| 906 | upload::abort(&self.store, &progress).await?; | |
| 907 | } | |
| 908 | self.db.delete_upload(&row.id).await?; | |
| 909 | empty(204, &[]) | |
| 910 | } | |
| 911 | ||
| 912 | async fn tags(&self, url: &Url, name: &ImageName, found: Option<PackageRow>) -> Result<Response> { | |
| 913 | let Some(package) = found else { | |
| 914 | return error(404, "NAME_UNKNOWN", format!("There is no image {}.", name.full())); | |
| 915 | }; | |
| 916 | let n = query(url, "n").and_then(|n| n.parse::<u32>().ok()).unwrap_or(MAX_TAGS_PAGE).min(MAX_TAGS_PAGE); | |
| 917 | let last = query(url, "last"); | |
| 918 | let tags = if n == 0 { Vec::new() } else { self.db.tag_names(&package.id, last.as_deref(), n).await? }; | |
| 919 | let mut response = Response::from_json(&json!({ "name": name.full(), "tags": tags }))?; | |
| 920 | if tags.len() as u32 == n | |
| 921 | && let Some(final_tag) = tags.last() | |
| 922 | { | |
| 923 | response | |
| 924 | .headers_mut() | |
| 925 | .set("link", &format!("</v2/{}/tags/list?n={n}&last={final_tag}>; rel=\"next\"", name.full()))?; | |
| 926 | } | |
| 927 | Ok(response) | |
| 928 | } | |
| 929 | ||
| 930 | async fn referrers(&self, url: &Url, name: &ImageName, found: Option<PackageRow>, digest: &str) -> Result<Response> { | |
| 931 | let Some(digest) = Digest::parse(digest) else { | |
| 932 | return error(400, "DIGEST_INVALID", format!("{digest} is not a sha256 digest.")); | |
| 933 | }; | |
| 934 | let Some(package) = found else { | |
| 935 | return error(404, "NAME_UNKNOWN", format!("There is no image {}.", name.full())); | |
| 936 | }; | |
| 937 | let wanted = query(url, "artifactType"); | |
| 938 | let mut manifests: Vec<Value> = Vec::new(); | |
| 939 | for version in self.db.referrers(&package.id, digest.as_str()).await? { | |
| 940 | let meta = version.meta(); | |
| 941 | let artifact_type = meta["artifact_type"].as_str().map(str::to_owned); | |
| 942 | if wanted.is_some() && artifact_type != wanted { | |
| 943 | continue; | |
| 944 | } | |
| 945 | let size = self.db.blob(&Digest::parse(&version.digest).unwrap_or(digest.clone())).await?.map_or(0, |b| b.size); | |
| 946 | let mut entry = json!({ | |
| 947 | "mediaType": meta["media_type"].as_str().unwrap_or(manifest::OCI_MANIFEST), | |
| 948 | "digest": version.digest, | |
| 949 | "size": size, | |
| 950 | }); | |
| 951 | if let Some(artifact_type) = artifact_type { | |
| 952 | entry["artifactType"] = json!(artifact_type); | |
| 953 | } | |
| 954 | if meta["annotations"].is_object() { | |
| 955 | entry["annotations"] = meta["annotations"].clone(); | |
| 956 | } | |
| 957 | manifests.push(entry); | |
| 958 | } | |
| 959 | let body = json!({ "schemaVersion": 2, "mediaType": manifest::OCI_INDEX, "manifests": manifests }); | |
| 960 | let mut response = Response::from_json(&body)?; | |
| 961 | response.headers_mut().set("content-type", manifest::OCI_INDEX)?; | |
| 962 | if wanted.is_some() { | |
| 963 | response.headers_mut().set("oci-filters-applied", "artifactType")?; | |
| 964 | } | |
| 965 | Ok(response) | |
| 966 | } | |
| 967 | ||
| 968 | /// Deletes a version of a package, with its tags, and says so. | |
| 969 | pub(crate) async fn remove_version(&self, package: &PackageRow, version: &crate::db::VersionRow, caller: &Caller) -> Result<()> { | |
| 970 | self.db.delete_version(&version.id).await?; | |
| 971 | self.db.measure(&package.workspace).await?; | |
| 972 | let event = PackageEvent { | |
| 973 | version: Some(version.version.clone()), | |
| 974 | digest: Some(version.digest.clone()), | |
| 975 | ..self.event_of(package) | |
| 976 | }; | |
| 977 | self.announce("package.version_deleted", package, event, caller).await; | |
| Cargo on g1t.sh: a sparse registry per workspace, cargo publish/add/yank/search with a g1t token; npm rows show their version | 978 | // npm and Cargo name a version by its number; an image by its digest. |
| npm on g1t.sh: publish and install @<workspace>/<name> with the npm CLI and a g1t token | 979 | let path = if package.ecosystem == "npm" { |
| 980 | format!("@{}/{}@{}", package.workspace, package.name, version.version) | |
| Cargo on g1t.sh: a sparse registry per workspace, cargo publish/add/yank/search with a g1t token; npm rows show their version | 981 | } else if package.ecosystem == "cargo" { |
| 982 | format!("{}/{}@{}", package.workspace, package.name, version.version) | |
| npm on g1t.sh: publish and install @<workspace>/<name> with the npm CLI and a g1t token | 983 | } else { |
| 984 | format!("{}/{}@{}", package.workspace, package.name, version.digest) | |
| 985 | }; | |
| 986 | self.audit(caller, "package.delete_version", package, Some(&path), None).await; | |
| Packages, with a container registry on g1t.sh; workspaces deleted whole and kept 30 days; Members for every member | 987 | Ok(()) |
| 988 | } | |
| 989 | ||
| 990 | pub(crate) fn event_of(&self, package: &PackageRow) -> PackageEvent { | |
| 991 | PackageEvent { | |
| 992 | package_id: package.id.clone(), | |
| 993 | workspace: package.workspace.clone(), | |
| 994 | ecosystem: Ecosystem::parse(&package.ecosystem).unwrap_or(Ecosystem::Container).as_str().to_owned(), | |
| 995 | name: package.name.clone(), | |
| 996 | repo_id: package.repo_id.clone(), | |
| 997 | ..PackageEvent::default() | |
| 998 | } | |
| 999 | } | |
| 1000 | } | |
| 1001 | ||
| 1002 | fn upload_status(name: &ImageName, row: &UploadRow) -> Result<Response> { | |
| 1003 | upload_status_with(name, &row.id, row.offset, 204) | |
| 1004 | } | |
| 1005 | ||
| 1006 | fn upload_status_with(name: &ImageName, id: &str, offset: u64, status: u16) -> Result<Response> { | |
| 1007 | empty(status, &[ | |
| 1008 | ("location", format!("/v2/{}/blobs/uploads/{id}", name.full())), | |
| 1009 | ("range", range::upload_range(offset)), | |
| 1010 | ("docker-upload-uuid", id.to_owned()), | |
| 1011 | ("content-length", "0".to_owned()), | |
| 1012 | ]) | |
| 1013 | } | |
| 1014 | ||
| 1015 | #[cfg(test)] | |
| 1016 | mod tests { | |
| 1017 | use super::*; | |
| 1018 | ||
| 1019 | #[test] | |
| 1020 | fn the_challenge_names_the_token_endpoint_on_the_same_host() { | |
| 1021 | let url = Url::parse("https://g1t.sh/v2/").unwrap(); | |
| 1022 | assert_eq!(challenge(&url, None), "Bearer realm=\"https://g1t.sh/v2/token\",service=\"g1t.sh\""); | |
| 1023 | let local = Url::parse("http://localhost:8790/v2/acme/web/manifests/latest").unwrap(); | |
| 1024 | assert_eq!( | |
| 1025 | challenge(&local, Some("repository:acme/web:pull".into())), | |
| 1026 | "Bearer realm=\"http://localhost:8790/v2/token\",service=\"localhost:8790\",scope=\"repository:acme/web:pull\"" | |
| 1027 | ); | |
| 1028 | } | |
| 1029 | } |