g1t/services/packages/src/oci.rs

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