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