g1t/services/packages/src/oci.rs

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