Skip to content
1,056 linesCodeBlameRaw
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
9use futures_util::StreamExt;
10use g1t_contracts::audit::AuditActor;
11use g1t_contracts::events::PackageEvent;
12use g1t_contracts::packages::Ecosystem;
13use g1t_contracts::time::rfc3339;
14use g1t_contracts::{Viewer, new_id};
15use g1t_kit::now_ms;
16use serde_json::{Value, json};
17use worker::{Context, Headers, Method, Request, Response, ResponseBody, Result, Url};
18
19use crate::access::Action;
20use crate::limits;
21use crate::db::{NewFile, NewVersion, PackageRow, UploadRow};
22use crate::digest::{Digest, Sha256};
23use crate::manifest::{self, Kind};
24use crate::names::{self, ImageName, Reference, Route};
25use crate::range;
26use crate::store::BlobStore;
27use crate::token::{self, Claims, Grant};
28use crate::upload::{self, Finished, Progress, Writer};
29use crate::{Caller, Packages, TargetOf};
30
31const CONTAINER: &str = "container";
32/// Blobs larger than this are downloaded from the store directly when it
33/// can sign a URL.
34const REDIRECT_BYTES: u64 = 4 * 1024 * 1024;
35/// How long a signed download URL works.
36const REDIRECT_SECONDS: u32 = 10 * 60;
37/// The most tags one page lists.
38const MAX_TAGS_PAGE: u32 = 1000;
39const DOCS: &str = "https://docs.g1t.sh/guides/containers/";
40
41/// An answer in the registry's error format: `{"errors": [...]}`.
42pub 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
47fn 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
55fn 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.
60pub(crate) enum Credentials {
61 None,
62 /// One of this registry's tokens.
63 Token(Claims),
64 /// A g1t token or password, resolved.
65 Viewer(Viewer),
66 /// Credentials that are wrong, or a token that expired.
67 Bad,
68}
69
70impl 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.
77pub(crate) fn origin(url: &Url) -> String {
78 let host = url.host_str().unwrap_or("g1t.sh");
79 match url.port() {
80 Some(port) => format!("{}://{host}:{port}", url.scheme()),
81 None => format!("{}://{host}", url.scheme()),
82 }
83}
84
85fn 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.
93fn 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
101fn 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
107fn 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.
112pub(crate) fn published_by(caller: &Caller) -> Option<String> {
113 caller.actor.as_ref().map(|actor| actor.on_behalf_of.clone().unwrap_or_else(|| actor.actor.clone()))
114}
115
116fn 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
127impl Packages {
128 /// Answers a registry request.
129 pub async fn registry(&self, request: Request, ctx: &Context) -> Result<Response> {
130 let url = request.url()?;
131 let Some(route) = names::route(url.path()) else {
132 return error(404, "NAME_UNKNOWN", "There is nothing at this address.");
133 };
134 let mut response = match self.route(request, &url, route, ctx).await {
135 Ok(response) => response,
136 Err(problem) => {
137 worker::console_error!("packages: {} {}: {problem}", url.path(), problem);
138 error(500, "UNKNOWN", "Something went wrong on our side. Try again in a moment.")?
139 }
140 };
141 response.headers_mut().set("docker-distribution-api-version", "registry/2.0")?;
142 Ok(response)
143 }
144
145 async fn route(&self, mut request: Request, url: &Url, route: Route, ctx: &Context) -> Result<Response> {
146 let method = request.method();
147 let credentials = self.credentials(&request).await?;
148 if limits::counts(&method, &route)
149 && let Some(refused) = self.limited(&request, &credentials, "`docker login g1t.sh`").await?
150 {
151 return Ok(refused);
152 }
153 if let Route::Base = route {
154 return match credentials {
155 Credentials::Token(_) | Credentials::Viewer(Some(_)) => Ok(Response::from_json(&json!({}))?),
156 _ => unauthorized(url, None, "Sign in with `docker login`, a g1t token as the password."),
157 };
158 }
159 if let Route::Token = route {
160 // Docker's OAuth form: `grant_type=password` with the username
161 // and password in the body, and the scopes beside them.
162 if method == Method::Post {
163 let form = request.text().await.unwrap_or_default();
164 let fields: Vec<(String, String)> = Url::parse(&format!("http://form/?{form}"))?
165 .query_pairs()
166 .map(|(k, v)| (k.into_owned(), v.into_owned()))
167 .collect();
168 let field = |key: &str| fields.iter().find(|(k, _)| k == key).map(|(_, v)| v.clone());
169 let credentials = match (field("grant_type").as_deref(), field("username"), field("password")) {
170 (Some("password"), Some(username), Some(password)) => match self.viewer_for(&username, &password).await? {
171 Some(user) => Credentials::Viewer(Some(user)),
172 None => Credentials::Bad,
173 },
174 (Some("password"), ..) | (Some("refresh_token"), ..) => Credentials::Bad,
175 _ => credentials,
176 };
177 let scopes: Vec<String> = fields.iter().filter(|(k, _)| k == "scope").map(|(_, v)| v.clone()).collect();
178 return self.issue(url, credentials, &scopes).await;
179 }
180 let scopes: Vec<String> = url.query_pairs().filter(|(k, _)| k == "scope").map(|(_, v)| v.into_owned()).collect();
181 return self.issue(url, credentials, &scopes).await;
182 }
183 let (name, action) = match &route {
184 Route::Manifest { name, .. } => (name, match method {
185 Method::Get | Method::Head => Action::Pull,
186 Method::Put => Action::Push,
187 Method::Delete => Action::Delete,
188 _ => return error(405, "UNSUPPORTED", "Not a method manifests take."),
189 }),
190 Route::Blob { name, .. } => (name, match method {
191 Method::Get | Method::Head => Action::Pull,
192 Method::Delete => Action::Delete,
193 _ => return error(405, "UNSUPPORTED", "Not a method blobs take."),
194 }),
195 Route::Uploads { name } | Route::Upload { name, .. } => (name, Action::Push),
196 Route::Tags { name } | Route::Referrers { name, .. } => (name, Action::Pull),
197 Route::Base | Route::Token => unreachable!("answered above"),
198 };
199 let name = match names::parse_name(name) {
200 Ok(name) => name,
201 Err(message) => return error(400, "NAME_INVALID", message),
202 };
203 let found = self.db.package(&name.workspace, CONTAINER, &name.name).await?;
204 // A deleted 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).
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 {
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)),
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.
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>> {
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 {
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(),
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
305 pub(crate) async fn credentials(&self, request: &Request) -> Result<Credentials> {
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",
342 Action::Delete | Action::Admin | Action::Settings => "delete",
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) => {
349 let caller = || Caller { actor: claims.actor.clone(), viewer: claims.viewer.as_deref().cloned() };
350 if claims.allows(&name.full(), action) {
351 return Ok(Ok(caller()));
352 }
353 // A token for something else still pulls a public image.
354 if action == Action::Pull && found.is_some_and(|package| package.public()) {
355 return Ok(Ok(caller()));
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 };
365 let mut target = self.target(name, found).await?;
366 let decision = self.decide(viewer.as_ref(), &mut target, action).await?;
367 if decision.allowed {
368 return Ok(Ok(Caller::of(viewer.as_ref())));
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 }
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 }
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();
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();
424 let claims = Claims {
425 actor: viewer.as_ref().map(AuditActor::of),
426 viewer: viewer.as_ref().filter(|_| !pushes.is_empty()).map(|user| Box::new(token::slim(user, &pushes))),
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?;
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))))
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 };
515 self.count_version_download(&package.id, &version.id, ctx);
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?;
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 };
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) => {
799 let mut target = TargetOf::package(package);
800 self.decide(viewer.as_ref(), &mut target, Action::Pull).await?.allowed
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
1011 /// Deletes a version of a package, and says so: hidden with its tags,
1012 /// restorable for 30 days (see settings.rs).
1013 pub(crate) async fn remove_version(&self, package: &PackageRow, version: &crate::db::VersionRow, caller: &Caller) -> Result<()> {
1014 self.remove_version_from(package, version, caller, None).await
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
1029fn upload_status(name: &ImageName, row: &UploadRow) -> Result<Response> {
1030 upload_status_with(name, &row.id, row.offset, 204)
1031}
1032
1033fn 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)]
1043mod 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}