Git hosting and a container registry in one Rust binary (axum + Astro)
Merge branch 'worktree-agent-aef42efcdeffcc508'
7 files changed, +2303 -8
+369-0backend/src/registry/blobs.rs
| @@ -0,0 +1,369 @@ | ||
| 1 | +//! Layers and configs. Uploads are staged on local disk (DATA_DIR/uploads), | |
| 2 | +//! verified against their digest, then stored once in R2. Downloads never | |
| 3 | +//! touch this server: GET answers with a redirect to a presigned R2 URL. | |
| 4 | + | |
| 5 | +use std::{path::Path, time::Duration}; | |
| 6 | + | |
| 7 | +use axum::{ | |
| 8 | + body::Body, | |
| 9 | + extract::Request, | |
| 10 | + http::{HeaderValue, StatusCode, header}, | |
| 11 | + response::{IntoResponse, Response}, | |
| 12 | +}; | |
| 13 | +use futures::TryStreamExt; | |
| 14 | +use sha2::{Digest, Sha256}; | |
| 15 | +use tokio::io::AsyncWriteExt; | |
| 16 | +use tokio_util::io::StreamReader; | |
| 17 | +use uuid::Uuid; | |
| 18 | + | |
| 19 | +use super::{Caller, RegError, RegResult, authorize, ensure_package, valid_digest}; | |
| 20 | +use crate::{models::Package, state::AppState, storage::keys}; | |
| 21 | + | |
| 22 | +/// How long a presigned layer download stays valid. | |
| 23 | +const DOWNLOAD_TTL: Duration = Duration::from_secs(20 * 60); | |
| 24 | + | |
| 25 | +fn query_param(request: &Request, key: &str) -> Option<String> { | |
| 26 | + let query = request.uri().query()?; | |
| 27 | + url::form_urlencoded::parse(query.as_bytes()).find(|(k, _)| k == key).map(|(_, v)| v.into_owned()) | |
| 28 | +} | |
| 29 | + | |
| 30 | +fn upload_path(state: &AppState, id: &Uuid) -> std::path::PathBuf { | |
| 31 | + state.config.uploads_dir().join(id.to_string()) | |
| 32 | +} | |
| 33 | + | |
| 34 | +fn range_value(size: u64) -> String { | |
| 35 | + // The spec's inclusive byte range; an empty upload is reported as 0-0. | |
| 36 | + format!("0-{}", size.saturating_sub(1)) | |
| 37 | +} | |
| 38 | + | |
| 39 | +fn upload_response(status: StatusCode, name: &str, id: &Uuid, size: u64) -> Response { | |
| 40 | + let mut response = status.into_response(); | |
| 41 | + let headers = response.headers_mut(); | |
| 42 | + if let Ok(location) = HeaderValue::from_str(&format!("/v2/{name}/blobs/uploads/{id}")) { | |
| 43 | + headers.insert(header::LOCATION, location); | |
| 44 | + } | |
| 45 | + headers.insert(header::RANGE, HeaderValue::from_str(&range_value(size)).expect("ascii")); | |
| 46 | + headers.insert("docker-upload-uuid", HeaderValue::from_str(&id.to_string()).expect("ascii")); | |
| 47 | + headers.insert(header::CONTENT_LENGTH, HeaderValue::from_static("0")); | |
| 48 | + response | |
| 49 | +} | |
| 50 | + | |
| 51 | +fn created_blob(name: &str, digest: &str) -> Response { | |
| 52 | + let mut response = StatusCode::CREATED.into_response(); | |
| 53 | + let headers = response.headers_mut(); | |
| 54 | + if let Ok(location) = HeaderValue::from_str(&format!("/v2/{name}/blobs/{digest}")) { | |
| 55 | + headers.insert(header::LOCATION, location); | |
| 56 | + } | |
| 57 | + if let Ok(value) = HeaderValue::from_str(digest) { | |
| 58 | + headers.insert("docker-content-digest", value); | |
| 59 | + } | |
| 60 | + headers.insert(header::CONTENT_LENGTH, HeaderValue::from_static("0")); | |
| 61 | + response | |
| 62 | +} | |
| 63 | + | |
| 64 | +/// Size of a blob, if it is linked to the package. | |
| 65 | +async fn linked_size(state: &AppState, package_id: i64, digest: &str) -> RegResult<Option<i64>> { | |
| 66 | + Ok(sqlx::query_scalar( | |
| 67 | + "select b.size from package_blobs pb join blobs b on b.digest = pb.digest where pb.package_id = $1 and pb.digest = $2", | |
| 68 | + ) | |
| 69 | + .bind(package_id) | |
| 70 | + .bind(digest) | |
| 71 | + .fetch_optional(&state.db) | |
| 72 | + .await?) | |
| 73 | +} | |
| 74 | + | |
| 75 | +/// HEAD and GET /v2/<name>/blobs/<digest>. | |
| 76 | +pub async fn get(state: &AppState, caller: &Caller, name: &str, digest: &str, head: bool) -> RegResult<Response> { | |
| 77 | + if !valid_digest(digest) { | |
| 78 | + return Err(RegError::digest_invalid("only sha256 digests are supported")); | |
| 79 | + } | |
| 80 | + let target = authorize(state, caller, name, "pull").await?; | |
| 81 | + let Some(package) = &target.package else { return Err(RegError::blob_unknown()) }; | |
| 82 | + let Some(size) = linked_size(state, package.id, digest).await? else { return Err(RegError::blob_unknown()) }; | |
| 83 | + | |
| 84 | + if head { | |
| 85 | + return Ok(Response::builder() | |
| 86 | + .status(StatusCode::OK) | |
| 87 | + .header(header::CONTENT_LENGTH, size) | |
| 88 | + .header(header::CONTENT_TYPE, "application/octet-stream") | |
| 89 | + .header("docker-content-digest", digest) | |
| 90 | + .body(Body::empty()) | |
| 91 | + .map_err(RegError::internal)?); | |
| 92 | + } | |
| 93 | + let url = state.storage.presign_get(&keys::blob(digest), DOWNLOAD_TTL, None); | |
| 94 | + tracing::debug!(package = %package.full_name(), digest, size, "blob download redirected to R2"); | |
| 95 | + Ok(Response::builder() | |
| 96 | + .status(StatusCode::TEMPORARY_REDIRECT) | |
| 97 | + .header(header::LOCATION, url.as_str()) | |
| 98 | + .header("docker-content-digest", digest) | |
| 99 | + .header(header::CONTENT_LENGTH, 0) | |
| 100 | + .body(Body::empty()) | |
| 101 | + .map_err(RegError::internal)?) | |
| 102 | +} | |
| 103 | + | |
| 104 | +/// DELETE /v2/<name>/blobs/<digest>: unlinks the blob from this image. The | |
| 105 | +/// bytes are removed by garbage collection once nothing links them. | |
| 106 | +pub async fn delete(state: &AppState, caller: &Caller, name: &str, digest: &str) -> RegResult<Response> { | |
| 107 | + if !valid_digest(digest) { | |
| 108 | + return Err(RegError::digest_invalid("only sha256 digests are supported")); | |
| 109 | + } | |
| 110 | + let target = authorize(state, caller, name, "delete").await?; | |
| 111 | + let package = target.package()?; | |
| 112 | + let removed = sqlx::query("delete from package_blobs where package_id = $1 and digest = $2") | |
| 113 | + .bind(package.id) | |
| 114 | + .bind(digest) | |
| 115 | + .execute(&state.db) | |
| 116 | + .await? | |
| 117 | + .rows_affected(); | |
| 118 | + if removed == 0 { | |
| 119 | + return Err(RegError::blob_unknown()); | |
| 120 | + } | |
| 121 | + tracing::info!(package = %package.full_name(), digest, user = caller.user_name().unwrap_or("-"), "blob unlinked"); | |
| 122 | + Ok(StatusCode::ACCEPTED.into_response()) | |
| 123 | +} | |
| 124 | + | |
| 125 | +/// POST /v2/<name>/blobs/uploads/: starts a session, finishes a monolithic | |
| 126 | +/// upload (?digest=), or mounts a blob from another image (?mount=&from=). | |
| 127 | +pub async fn start_upload(state: &AppState, caller: &Caller, name: &str, request: Request) -> RegResult<Response> { | |
| 128 | + let target = authorize(state, caller, name, "push").await?; | |
| 129 | + | |
| 130 | + if let (Some(mount), Some(from)) = (query_param(&request, "mount"), query_param(&request, "from")) { | |
| 131 | + if valid_digest(&mount) { | |
| 132 | + if let Some(response) = try_mount(state, caller, &target, &mount, &from).await? { | |
| 133 | + return Ok(response); | |
| 134 | + } | |
| 135 | + } | |
| 136 | + } | |
| 137 | + | |
| 138 | + let package = ensure_package(state, &target, caller).await?; | |
| 139 | + let id = Uuid::new_v4(); | |
| 140 | + let path = upload_path(state, &id); | |
| 141 | + tokio::fs::create_dir_all(state.config.uploads_dir()).await?; | |
| 142 | + tokio::fs::File::create(&path).await?; | |
| 143 | + | |
| 144 | + if let Some(digest) = query_param(&request, "digest") { | |
| 145 | + // Monolithic: the whole blob is in this request. | |
| 146 | + if !valid_digest(&digest) { | |
| 147 | + let _ = tokio::fs::remove_file(&path).await; | |
| 148 | + return Err(RegError::digest_invalid("only sha256 digests are supported")); | |
| 149 | + } | |
| 150 | + let limit = state.config.limits.max_image_layer_bytes; | |
| 151 | + let result = async { | |
| 152 | + append_body(&path, request.into_body(), 0, limit).await?; | |
| 153 | + finalize(state, &package, &path, &digest).await | |
| 154 | + } | |
| 155 | + .await; | |
| 156 | + let _ = tokio::fs::remove_file(&path).await; | |
| 157 | + result?; | |
| 158 | + tracing::info!(package = %package.full_name(), digest, "blob uploaded (monolithic)"); | |
| 159 | + return Ok(created_blob(name, &digest)); | |
| 160 | + } | |
| 161 | + | |
| 162 | + sqlx::query("insert into blob_uploads (id, package_id, user_id) values ($1, $2, $3)") | |
| 163 | + .bind(id) | |
| 164 | + .bind(package.id) | |
| 165 | + .bind(caller.user_id()) | |
| 166 | + .execute(&state.db) | |
| 167 | + .await?; | |
| 168 | + tracing::debug!(package = %package.full_name(), upload = %id, "upload session started"); | |
| 169 | + Ok(upload_response(StatusCode::ACCEPTED, name, &id, 0)) | |
| 170 | +} | |
| 171 | + | |
| 172 | +/// Links `digest` from the `from` image when the caller may read it there. | |
| 173 | +/// Anything else falls back to a normal upload, as the spec allows. | |
| 174 | +async fn try_mount(state: &AppState, caller: &Caller, target: &super::Target, digest: &str, from: &str) -> RegResult<Option<Response>> { | |
| 175 | + let source = match authorize(state, caller, from, "pull").await { | |
| 176 | + Ok(source) => source, | |
| 177 | + Err(_) => return Ok(None), | |
| 178 | + }; | |
| 179 | + let Some(source_package) = &source.package else { return Ok(None) }; | |
| 180 | + if linked_size(state, source_package.id, digest).await?.is_none() { | |
| 181 | + return Ok(None); | |
| 182 | + } | |
| 183 | + let package = ensure_package(state, target, caller).await?; | |
| 184 | + sqlx::query("insert into package_blobs (package_id, digest) values ($1, $2) on conflict do nothing") | |
| 185 | + .bind(package.id) | |
| 186 | + .bind(digest) | |
| 187 | + .execute(&state.db) | |
| 188 | + .await?; | |
| 189 | + tracing::info!(from = %source_package.full_name(), to = %package.full_name(), digest, "blob mounted"); | |
| 190 | + Ok(Some(created_blob(&target.name, digest))) | |
| 191 | +} | |
| 192 | + | |
| 193 | +struct Session { | |
| 194 | + id: Uuid, | |
| 195 | + package: Package, | |
| 196 | + size: u64, | |
| 197 | +} | |
| 198 | + | |
| 199 | +async fn session(state: &AppState, caller: &Caller, name: &str, id: &str) -> RegResult<(super::Target, Session)> { | |
| 200 | + let target = authorize(state, caller, name, "push").await?; | |
| 201 | + let package = target.package()?.clone(); | |
| 202 | + let id = Uuid::parse_str(id).map_err(|_| RegError::upload_unknown())?; | |
| 203 | + let size: Option<i64> = sqlx::query_scalar("select size from blob_uploads where id = $1 and package_id = $2") | |
| 204 | + .bind(id) | |
| 205 | + .bind(package.id) | |
| 206 | + .fetch_optional(&state.db) | |
| 207 | + .await?; | |
| 208 | + let size = size.ok_or_else(RegError::upload_unknown)?; | |
| 209 | + Ok((target, Session { id, package, size: size.max(0) as u64 })) | |
| 210 | +} | |
| 211 | + | |
| 212 | +async fn drop_session(state: &AppState, id: &Uuid) { | |
| 213 | + let _ = sqlx::query("delete from blob_uploads where id = $1").bind(id).execute(&state.db).await; | |
| 214 | + let _ = tokio::fs::remove_file(upload_path(state, id)).await; | |
| 215 | +} | |
| 216 | + | |
| 217 | +/// PATCH: appends a chunk. | |
| 218 | +pub async fn patch_upload(state: &AppState, caller: &Caller, name: &str, id: &str, request: Request) -> RegResult<Response> { | |
| 219 | + let (_, session) = session(state, caller, name, id).await?; | |
| 220 | + | |
| 221 | + // Content-Range "start-end" must continue exactly where we are. | |
| 222 | + if let Some(range) = request.headers().get(header::CONTENT_RANGE).and_then(|v| v.to_str().ok()) { | |
| 223 | + let start = range.trim().trim_start_matches("bytes ").split('-').next().and_then(|s| s.parse::<u64>().ok()); | |
| 224 | + if start != Some(session.size) { | |
| 225 | + let mut response = upload_response(StatusCode::RANGE_NOT_SATISFIABLE, name, &session.id, session.size); | |
| 226 | + *response.status_mut() = StatusCode::RANGE_NOT_SATISFIABLE; | |
| 227 | + return Ok(response); | |
| 228 | + } | |
| 229 | + } | |
| 230 | + | |
| 231 | + let path = upload_path(state, &session.id); | |
| 232 | + let limit = state.config.limits.max_image_layer_bytes; | |
| 233 | + let written = match append_body(&path, request.into_body(), session.size, limit).await { | |
| 234 | + Ok(n) => n, | |
| 235 | + Err(error) => { | |
| 236 | + if error.status == StatusCode::PAYLOAD_TOO_LARGE { | |
| 237 | + drop_session(state, &session.id).await; | |
| 238 | + } | |
| 239 | + return Err(error); | |
| 240 | + } | |
| 241 | + }; | |
| 242 | + let size = session.size + written; | |
| 243 | + sqlx::query("update blob_uploads set size = $2, updated_at = now() where id = $1") | |
| 244 | + .bind(session.id) | |
| 245 | + .bind(size as i64) | |
| 246 | + .execute(&state.db) | |
| 247 | + .await?; | |
| 248 | + tracing::debug!(package = %session.package.full_name(), upload = %session.id, chunk = written, size, "upload chunk"); | |
| 249 | + Ok(upload_response(StatusCode::ACCEPTED, name, &session.id, size)) | |
| 250 | +} | |
| 251 | + | |
| 252 | +/// PUT ?digest=: appends any final chunk, verifies and stores the blob. | |
| 253 | +pub async fn finish_upload(state: &AppState, caller: &Caller, name: &str, id: &str, request: Request) -> RegResult<Response> { | |
| 254 | + let (_, session) = session(state, caller, name, id).await?; | |
| 255 | + let digest = query_param(&request, "digest").ok_or_else(|| RegError::digest_invalid("digest parameter is required"))?; | |
| 256 | + if !valid_digest(&digest) { | |
| 257 | + return Err(RegError::digest_invalid("only sha256 digests are supported")); | |
| 258 | + } | |
| 259 | + let path = upload_path(state, &session.id); | |
| 260 | + let limit = state.config.limits.max_image_layer_bytes; | |
| 261 | + let result = async { | |
| 262 | + append_body(&path, request.into_body(), session.size, limit).await?; | |
| 263 | + finalize(state, &session.package, &path, &digest).await | |
| 264 | + } | |
| 265 | + .await; | |
| 266 | + // Success or failure, the session is over (a bad digest cannot be fixed). | |
| 267 | + drop_session(state, &session.id).await; | |
| 268 | + let size = result?; | |
| 269 | + tracing::info!(package = %session.package.full_name(), digest, size, user = caller.user_name().unwrap_or("-"), "blob uploaded"); | |
| 270 | + Ok(created_blob(name, &digest)) | |
| 271 | +} | |
| 272 | + | |
| 273 | +/// GET: where an interrupted upload stands. | |
| 274 | +pub async fn upload_status(state: &AppState, caller: &Caller, name: &str, id: &str) -> RegResult<Response> { | |
| 275 | + let (_, session) = session(state, caller, name, id).await?; | |
| 276 | + Ok(upload_response(StatusCode::NO_CONTENT, name, &session.id, session.size)) | |
| 277 | +} | |
| 278 | + | |
| 279 | +/// DELETE: abandons an upload. | |
| 280 | +pub async fn cancel_upload(state: &AppState, caller: &Caller, name: &str, id: &str) -> RegResult<Response> { | |
| 281 | + let (_, session) = session(state, caller, name, id).await?; | |
| 282 | + drop_session(state, &session.id).await; | |
| 283 | + Ok(StatusCode::NO_CONTENT.into_response()) | |
| 284 | +} | |
| 285 | + | |
| 286 | +/// Streams a request body onto the end of `path`, refusing to grow past | |
| 287 | +/// `limit`. Returns the number of bytes written. | |
| 288 | +async fn append_body(path: &Path, body: Body, already: u64, limit: u64) -> RegResult<u64> { | |
| 289 | + let mut file = tokio::fs::OpenOptions::new().append(true).open(path).await?; | |
| 290 | + let stream = body.into_data_stream().map_err(std::io::Error::other); | |
| 291 | + let mut reader = StreamReader::new(stream); | |
| 292 | + let mut written = 0u64; | |
| 293 | + let mut buffer = vec![0u8; 256 * 1024]; | |
| 294 | + loop { | |
| 295 | + let n = tokio::io::AsyncReadExt::read(&mut reader, &mut buffer).await?; | |
| 296 | + if n == 0 { | |
| 297 | + break; | |
| 298 | + } | |
| 299 | + written += n as u64; | |
| 300 | + if already + written > limit { | |
| 301 | + return Err(RegError::new( | |
| 302 | + StatusCode::PAYLOAD_TOO_LARGE, | |
| 303 | + "SIZE_INVALID", | |
| 304 | + format!("layers may be at most {}", crate::web::ui::bytes(limit)), | |
| 305 | + )); | |
| 306 | + } | |
| 307 | + file.write_all(&buffer[..n]).await?; | |
| 308 | + } | |
| 309 | + file.flush().await?; | |
| 310 | + Ok(written) | |
| 311 | +} | |
| 312 | + | |
| 313 | +async fn sha256_file(path: &Path) -> std::io::Result<(String, u64)> { | |
| 314 | + let path = path.to_path_buf(); | |
| 315 | + tokio::task::spawn_blocking(move || { | |
| 316 | + use std::io::Read; | |
| 317 | + let mut file = std::fs::File::open(path)?; | |
| 318 | + let mut hasher = Sha256::new(); | |
| 319 | + let mut buffer = vec![0u8; 1024 * 1024]; | |
| 320 | + let mut size = 0u64; | |
| 321 | + loop { | |
| 322 | + let n = file.read(&mut buffer)?; | |
| 323 | + if n == 0 { | |
| 324 | + break; | |
| 325 | + } | |
| 326 | + size += n as u64; | |
| 327 | + hasher.update(&buffer[..n]); | |
| 328 | + } | |
| 329 | + Ok((format!("sha256:{}", hex::encode(hasher.finalize())), size)) | |
| 330 | + }) | |
| 331 | + .await | |
| 332 | + .map_err(std::io::Error::other)? | |
| 333 | +} | |
| 334 | + | |
| 335 | +/// Verifies the staged file, stores it in R2 unless the blob already exists, | |
| 336 | +/// and links it to the package. Returns the size. | |
| 337 | +/// | |
| 338 | +/// A per-digest advisory lock serializes this with garbage collection, so a | |
| 339 | +/// blob being re-uploaded can never be deleted from R2 underneath it. | |
| 340 | +async fn finalize(state: &AppState, package: &Package, path: &Path, digest: &str) -> RegResult<u64> { | |
| 341 | + let (actual, size) = sha256_file(path).await?; | |
| 342 | + if actual != digest { | |
| 343 | + tracing::warn!(package = %package.full_name(), expected = digest, actual, "upload digest mismatch"); | |
| 344 | + return Err(RegError::digest_invalid("uploaded content does not match the digest").detail(serde_json::json!({ "expected": digest, "actual": actual }))); | |
| 345 | + } | |
| 346 | + | |
| 347 | + let mut tx = state.db.begin().await?; | |
| 348 | + sqlx::query("select pg_advisory_xact_lock(hashtextextended($1, 0))").bind(digest).execute(&mut *tx).await?; | |
| 349 | + let exists: bool = sqlx::query_scalar("select exists(select 1 from blobs where digest = $1)").bind(digest).fetch_one(&mut *tx).await?; | |
| 350 | + if !exists { | |
| 351 | + let started = std::time::Instant::now(); | |
| 352 | + state.storage.put_file(&keys::blob(digest), path).await?; | |
| 353 | + tracing::info!(digest, size, ms = started.elapsed().as_millis() as u64, "blob stored in R2"); | |
| 354 | + sqlx::query("insert into blobs (digest, size) values ($1, $2) on conflict do nothing") | |
| 355 | + .bind(digest) | |
| 356 | + .bind(size as i64) | |
| 357 | + .execute(&mut *tx) | |
| 358 | + .await?; | |
| 359 | + } else { | |
| 360 | + tracing::debug!(digest, "blob already stored; linking only"); | |
| 361 | + } | |
| 362 | + sqlx::query("insert into package_blobs (package_id, digest) values ($1, $2) on conflict do nothing") | |
| 363 | + .bind(package.id) | |
| 364 | + .bind(digest) | |
| 365 | + .execute(&mut *tx) | |
| 366 | + .await?; | |
| 367 | + tx.commit().await?; | |
| 368 | + Ok(size) | |
| 369 | +} |
+115-0backend/src/registry/gc.rs
| @@ -0,0 +1,115 @@ | ||
| 1 | +//! Registry garbage collection, hourly: | |
| 2 | +//! 1. abandoned uploads (and stray staging files) older than a day | |
| 3 | +//! 2. blob links no manifest of their package references, older than a day | |
| 4 | +//! (the grace period covers pushes still uploading layers) | |
| 5 | +//! 3. blobs nothing links any more: their R2 object, then their row | |
| 6 | +//! | |
| 7 | +//! Step 3 takes the same per-digest advisory lock as upload finalization, so | |
| 8 | +//! a blob being pushed again can never be deleted from R2 underneath it. | |
| 9 | + | |
| 10 | +use std::time::Duration; | |
| 11 | + | |
| 12 | +use crate::{state::AppState, storage::keys}; | |
| 13 | + | |
| 14 | +pub fn spawn_gc(state: AppState) { | |
| 15 | + tokio::spawn(async move { | |
| 16 | + // Let startup settle before the first pass. | |
| 17 | + tokio::time::sleep(Duration::from_secs(60)).await; | |
| 18 | + let mut tick = tokio::time::interval(Duration::from_secs(3600)); | |
| 19 | + loop { | |
| 20 | + tick.tick().await; | |
| 21 | + match gc_once(&state).await { | |
| 22 | + Ok(summary) => tracing::info!(%summary, "registry gc"), | |
| 23 | + Err(error) => tracing::error!(error = ?error, "registry gc failed"), | |
| 24 | + } | |
| 25 | + } | |
| 26 | + }); | |
| 27 | +} | |
| 28 | + | |
| 29 | +/// One pass. Returns a one-line summary for logs and the admin panel. | |
| 30 | +pub async fn gc_once(state: &AppState) -> anyhow::Result<String> { | |
| 31 | + let started = std::time::Instant::now(); | |
| 32 | + | |
| 33 | + // 1. Uploads nobody finished. | |
| 34 | + let stale: Vec<uuid::Uuid> = sqlx::query_scalar("delete from blob_uploads where updated_at < now() - interval '24 hours' returning id") | |
| 35 | + .fetch_all(&state.db) | |
| 36 | + .await?; | |
| 37 | + for id in &stale { | |
| 38 | + let _ = tokio::fs::remove_file(state.config.uploads_dir().join(id.to_string())).await; | |
| 39 | + } | |
| 40 | + let strays = remove_stray_files(state).await?; | |
| 41 | + | |
| 42 | + // 2. Links no manifest in the package uses. | |
| 43 | + let unlinked = sqlx::query( | |
| 44 | + "delete from package_blobs pb | |
| 45 | + where pb.created_at < now() - interval '24 hours' | |
| 46 | + and not exists (select 1 from manifest_references r join manifests m on m.id = r.manifest_id | |
| 47 | + where m.package_id = pb.package_id and r.digest = pb.digest)", | |
| 48 | + ) | |
| 49 | + .execute(&state.db) | |
| 50 | + .await? | |
| 51 | + .rows_affected(); | |
| 52 | + | |
| 53 | + // 3. Blobs with no links at all. | |
| 54 | + let orphans: Vec<(String, i64)> = sqlx::query_as( | |
| 55 | + "select digest, size from blobs b | |
| 56 | + where b.created_at < now() - interval '24 hours' | |
| 57 | + and not exists (select 1 from package_blobs pb where pb.digest = b.digest) | |
| 58 | + limit 5000", | |
| 59 | + ) | |
| 60 | + .fetch_all(&state.db) | |
| 61 | + .await?; | |
| 62 | + let (mut deleted, mut freed) = (0u64, 0i64); | |
| 63 | + for (digest, size) in orphans { | |
| 64 | + let mut tx = state.db.begin().await?; | |
| 65 | + sqlx::query("select pg_advisory_xact_lock(hashtextextended($1, 0))").bind(&digest).execute(&mut *tx).await?; | |
| 66 | + let still_orphan: bool = sqlx::query_scalar( | |
| 67 | + "select exists(select 1 from blobs b where b.digest = $1 and not exists (select 1 from package_blobs pb where pb.digest = b.digest))", | |
| 68 | + ) | |
| 69 | + .bind(&digest) | |
| 70 | + .fetch_one(&mut *tx) | |
| 71 | + .await?; | |
| 72 | + if !still_orphan { | |
| 73 | + continue; | |
| 74 | + } | |
| 75 | + if let Err(error) = state.storage.delete(&keys::blob(&digest)).await { | |
| 76 | + tracing::warn!(%error, digest, "could not delete blob from R2; will retry next pass"); | |
| 77 | + continue; | |
| 78 | + } | |
| 79 | + sqlx::query("delete from blobs where digest = $1").bind(&digest).execute(&mut *tx).await?; | |
| 80 | + tx.commit().await?; | |
| 81 | + deleted += 1; | |
| 82 | + freed += size; | |
| 83 | + } | |
| 84 | + | |
| 85 | + Ok(format!( | |
| 86 | + "{} stale uploads, {} stray files, {} unused links, {} blobs deleted ({} freed) in {} ms", | |
| 87 | + stale.len(), | |
| 88 | + strays, | |
| 89 | + unlinked, | |
| 90 | + deleted, | |
| 91 | + crate::web::ui::bytes(freed.max(0) as u64), | |
| 92 | + started.elapsed().as_millis() | |
| 93 | + )) | |
| 94 | +} | |
| 95 | + | |
| 96 | +/// Staging files with no upload row (monolithic uploads cut off mid-way). | |
| 97 | +async fn remove_stray_files(state: &AppState) -> anyhow::Result<u64> { | |
| 98 | + let dir = state.config.uploads_dir(); | |
| 99 | + let Ok(mut entries) = tokio::fs::read_dir(&dir).await else { return Ok(0) }; | |
| 100 | + let mut removed = 0; | |
| 101 | + let day = Duration::from_secs(24 * 3600); | |
| 102 | + while let Some(entry) = entries.next_entry().await? { | |
| 103 | + let Ok(meta) = entry.metadata().await else { continue }; | |
| 104 | + let old = meta.modified().ok().and_then(|m| m.elapsed().ok()).is_some_and(|age| age > day); | |
| 105 | + if !old { | |
| 106 | + continue; | |
| 107 | + } | |
| 108 | + let Some(id) = entry.file_name().to_str().and_then(|n| uuid::Uuid::parse_str(n).ok()) else { continue }; | |
| 109 | + let active: bool = sqlx::query_scalar("select exists(select 1 from blob_uploads where id = $1)").bind(id).fetch_one(&state.db).await?; | |
| 110 | + if !active && tokio::fs::remove_file(entry.path()).await.is_ok() { | |
| 111 | + removed += 1; | |
| 112 | + } | |
| 113 | + } | |
| 114 | + Ok(removed) | |
| 115 | +} |
+405-0backend/src/registry/manifests.rs
| @@ -0,0 +1,405 @@ | ||
| 1 | +//! Manifests and tags. Manifests are a few KB of JSON, kept in Postgres with | |
| 2 | +//! the digests they reference, which is what garbage collection walks. | |
| 3 | + | |
| 4 | +use axum::{ | |
| 5 | + Json, | |
| 6 | + body::Body, | |
| 7 | + extract::Request, | |
| 8 | + http::{HeaderValue, StatusCode, header}, | |
| 9 | + response::{IntoResponse, Response}, | |
| 10 | +}; | |
| 11 | +use serde_json::{Value, json}; | |
| 12 | +use sha2::{Digest, Sha256}; | |
| 13 | + | |
| 14 | +use super::{Caller, RegError, RegResult, authorize, ensure_package, valid_digest, valid_tag}; | |
| 15 | +use crate::{analytics, models::audit, state::AppState}; | |
| 16 | + | |
| 17 | +pub const DOCKER_MANIFEST: &str = "application/vnd.docker.distribution.manifest.v2+json"; | |
| 18 | +pub const DOCKER_LIST: &str = "application/vnd.docker.distribution.manifest.list.v2+json"; | |
| 19 | +pub const OCI_MANIFEST: &str = "application/vnd.oci.image.manifest.v1+json"; | |
| 20 | +pub const OCI_INDEX: &str = "application/vnd.oci.image.index.v1+json"; | |
| 21 | + | |
| 22 | +const KNOWN_TYPES: [&str; 4] = [DOCKER_MANIFEST, DOCKER_LIST, OCI_MANIFEST, OCI_INDEX]; | |
| 23 | + | |
| 24 | +#[derive(Debug, Clone, PartialEq, Eq)] | |
| 25 | +pub struct Descriptor { | |
| 26 | + pub digest: String, | |
| 27 | + pub size: i64, | |
| 28 | +} | |
| 29 | + | |
| 30 | +#[derive(Debug, Clone, PartialEq, Eq)] | |
| 31 | +pub struct ParsedManifest { | |
| 32 | + pub media_type: String, | |
| 33 | + /// Config and layers (image manifests). | |
| 34 | + pub blobs: Vec<Descriptor>, | |
| 35 | + /// Platform manifests (indexes and Docker manifest lists). | |
| 36 | + pub children: Vec<Descriptor>, | |
| 37 | +} | |
| 38 | + | |
| 39 | +impl ParsedManifest { | |
| 40 | + pub fn is_index(&self) -> bool { | |
| 41 | + self.media_type == OCI_INDEX || self.media_type == DOCKER_LIST | |
| 42 | + } | |
| 43 | +} | |
| 44 | + | |
| 45 | +fn manifest_invalid(message: impl Into<String>) -> RegError { | |
| 46 | + RegError::new(StatusCode::BAD_REQUEST, "MANIFEST_INVALID", message) | |
| 47 | +} | |
| 48 | + | |
| 49 | +fn descriptor(value: &Value) -> Result<Descriptor, RegError> { | |
| 50 | + let digest = value.get("digest").and_then(Value::as_str).ok_or_else(|| manifest_invalid("descriptor without digest"))?; | |
| 51 | + if !valid_digest(digest) { | |
| 52 | + return Err(manifest_invalid(format!("unsupported digest {digest}"))); | |
| 53 | + } | |
| 54 | + let size = value.get("size").and_then(Value::as_i64).filter(|s| *s >= 0).ok_or_else(|| manifest_invalid("descriptor without size"))?; | |
| 55 | + Ok(Descriptor { digest: digest.to_string(), size }) | |
| 56 | +} | |
| 57 | + | |
| 58 | +/// Reads the references out of a Docker v2 / OCI manifest or index. | |
| 59 | +pub fn parse_manifest(body: &[u8], content_type: Option<&str>) -> Result<ParsedManifest, RegError> { | |
| 60 | + let json: Value = serde_json::from_slice(body).map_err(|e| manifest_invalid(format!("not JSON: {e}")))?; | |
| 61 | + if json.get("schemaVersion").and_then(Value::as_i64) != Some(2) { | |
| 62 | + return Err(manifest_invalid("only schema version 2 manifests are supported")); | |
| 63 | + } | |
| 64 | + let declared = json.get("mediaType").and_then(Value::as_str); | |
| 65 | + let header_type = content_type.map(|c| c.split(';').next().unwrap_or("").trim()).filter(|c| KNOWN_TYPES.contains(c)); | |
| 66 | + let has_manifests = json.get("manifests").is_some_and(Value::is_array); | |
| 67 | + let media_type = header_type | |
| 68 | + .or(declared.filter(|d| KNOWN_TYPES.contains(d))) | |
| 69 | + .unwrap_or(if has_manifests { OCI_INDEX } else { OCI_MANIFEST }) | |
| 70 | + .to_string(); | |
| 71 | + | |
| 72 | + let mut parsed = ParsedManifest { media_type, blobs: Vec::new(), children: Vec::new() }; | |
| 73 | + if has_manifests { | |
| 74 | + for child in json["manifests"].as_array().into_iter().flatten() { | |
| 75 | + parsed.children.push(descriptor(child)?); | |
| 76 | + } | |
| 77 | + return Ok(parsed); | |
| 78 | + } | |
| 79 | + let config = json.get("config").ok_or_else(|| manifest_invalid("manifest has no config"))?; | |
| 80 | + parsed.blobs.push(descriptor(config)?); | |
| 81 | + for layer in json.get("layers").and_then(Value::as_array).into_iter().flatten() { | |
| 82 | + // Foreign layers live elsewhere (their "urls"); nothing to verify here. | |
| 83 | + if layer.get("urls").and_then(Value::as_array).is_some_and(|u| !u.is_empty()) { | |
| 84 | + continue; | |
| 85 | + } | |
| 86 | + parsed.blobs.push(descriptor(layer)?); | |
| 87 | + } | |
| 88 | + Ok(parsed) | |
| 89 | +} | |
| 90 | + | |
| 91 | +pub fn digest_of(body: &[u8]) -> String { | |
| 92 | + format!("sha256:{}", hex::encode(Sha256::digest(body))) | |
| 93 | +} | |
| 94 | + | |
| 95 | +#[derive(sqlx::FromRow)] | |
| 96 | +struct StoredManifest { | |
| 97 | + digest: String, | |
| 98 | + media_type: String, | |
| 99 | + content: Vec<u8>, | |
| 100 | +} | |
| 101 | + | |
| 102 | +/// GET/HEAD /v2/<name>/manifests/<tag or digest>. | |
| 103 | +pub async fn get(state: &AppState, caller: &Caller, name: &str, reference: &str, head: bool, _request: Request) -> RegResult<Response> { | |
| 104 | + let target = authorize(state, caller, name, "pull").await?; | |
| 105 | + let package = target.package()?; | |
| 106 | + let by_tag = !reference.starts_with("sha256:"); | |
| 107 | + let manifest: Option<StoredManifest> = if by_tag { | |
| 108 | + sqlx::query_as( | |
| 109 | + "select m.digest, m.media_type, m.content from tags t join manifests m on m.id = t.manifest_id | |
| 110 | + where t.package_id = $1 and t.name = $2", | |
| 111 | + ) | |
| 112 | + .bind(package.id) | |
| 113 | + .bind(reference) | |
| 114 | + .fetch_optional(&state.db) | |
| 115 | + .await? | |
| 116 | + } else { | |
| 117 | + sqlx::query_as("select digest, media_type, content from manifests where package_id = $1 and digest = $2") | |
| 118 | + .bind(package.id) | |
| 119 | + .bind(reference) | |
| 120 | + .fetch_optional(&state.db) | |
| 121 | + .await? | |
| 122 | + }; | |
| 123 | + let manifest = manifest.ok_or_else(RegError::manifest_unknown)?; | |
| 124 | + | |
| 125 | + if !head { | |
| 126 | + // A pull is a GET of a tag or of a top-level digest; the platform | |
| 127 | + // manifests fetched through an index do not count again. | |
| 128 | + let top_level: bool = by_tag | |
| 129 | + || sqlx::query_scalar( | |
| 130 | + "select not exists (select 1 from manifest_references r join manifests m on m.id = r.manifest_id | |
| 131 | + where m.package_id = $1 and r.digest = $2 and r.kind = 'manifest')", | |
| 132 | + ) | |
| 133 | + .bind(package.id) | |
| 134 | + .bind(&manifest.digest) | |
| 135 | + .fetch_one(&state.db) | |
| 136 | + .await?; | |
| 137 | + if top_level { | |
| 138 | + sqlx::query("update packages set pull_count = pull_count + 1 where id = $1").bind(package.id).execute(&state.db).await?; | |
| 139 | + analytics::track( | |
| 140 | + state, | |
| 141 | + "image_pulled", | |
| 142 | + caller.user_name(), | |
| 143 | + &package.url(), | |
| 144 | + json!({ "image": package.full_name(), "reference": reference, "visibility": package.visibility }), | |
| 145 | + ); | |
| 146 | + tracing::info!(package = %package.full_name(), reference, user = caller.user_name().unwrap_or("-"), "image pulled"); | |
| 147 | + } | |
| 148 | + } | |
| 149 | + | |
| 150 | + let length = manifest.content.len(); | |
| 151 | + let body = if head { Body::empty() } else { Body::from(manifest.content) }; | |
| 152 | + Response::builder() | |
| 153 | + .status(StatusCode::OK) | |
| 154 | + .header(header::CONTENT_TYPE, &manifest.media_type) | |
| 155 | + .header(header::CONTENT_LENGTH, length) | |
| 156 | + .header("docker-content-digest", &manifest.digest) | |
| 157 | + .header(header::ETAG, format!("\"{}\"", manifest.digest)) | |
| 158 | + .body(body) | |
| 159 | + .map_err(RegError::internal) | |
| 160 | +} | |
| 161 | + | |
| 162 | +/// PUT /v2/<name>/manifests/<tag or digest>. | |
| 163 | +pub async fn put(state: &AppState, caller: &Caller, name: &str, reference: &str, request: Request) -> RegResult<Response> { | |
| 164 | + let target = authorize(state, caller, name, "push").await?; | |
| 165 | + let by_tag = !reference.starts_with("sha256:"); | |
| 166 | + if by_tag && !valid_tag(reference) { | |
| 167 | + return Err(RegError::new(StatusCode::BAD_REQUEST, "TAG_INVALID", "invalid tag")); | |
| 168 | + } | |
| 169 | + if !by_tag && !valid_digest(reference) { | |
| 170 | + return Err(RegError::digest_invalid("only sha256 digests are supported")); | |
| 171 | + } | |
| 172 | + let content_type = request.headers().get(header::CONTENT_TYPE).and_then(|v| v.to_str().ok()).map(str::to_string); | |
| 173 | + let limit = state.config.limits.max_manifest_bytes; | |
| 174 | + let body = axum::body::to_bytes(request.into_body(), limit) | |
| 175 | + .await | |
| 176 | + .map_err(|_| RegError::new(StatusCode::PAYLOAD_TOO_LARGE, "SIZE_INVALID", format!("manifests may be at most {limit} bytes")))?; | |
| 177 | + let digest = digest_of(&body); | |
| 178 | + if !by_tag && reference != digest { | |
| 179 | + return Err(RegError::digest_invalid("manifest content does not match the digest in the URL")); | |
| 180 | + } | |
| 181 | + let parsed = parse_manifest(&body, content_type.as_deref())?; | |
| 182 | + let package = ensure_package(state, &target, caller).await?; | |
| 183 | + | |
| 184 | + // Everything referenced must already be in this image. | |
| 185 | + if !parsed.blobs.is_empty() { | |
| 186 | + let wanted: Vec<String> = parsed.blobs.iter().map(|d| d.digest.clone()).collect(); | |
| 187 | + let present: Vec<String> = sqlx::query_scalar("select digest from package_blobs where package_id = $1 and digest = any($2)") | |
| 188 | + .bind(package.id) | |
| 189 | + .bind(&wanted) | |
| 190 | + .fetch_all(&state.db) | |
| 191 | + .await?; | |
| 192 | + if let Some(missing) = wanted.iter().find(|d| !present.contains(d)) { | |
| 193 | + return Err(RegError::new(StatusCode::BAD_REQUEST, "MANIFEST_BLOB_UNKNOWN", "blob unknown to registry") | |
| 194 | + .detail(json!({ "digest": missing }))); | |
| 195 | + } | |
| 196 | + } | |
| 197 | + let mut total_size: i64 = parsed.blobs.iter().map(|d| d.size).sum(); | |
| 198 | + if !parsed.children.is_empty() { | |
| 199 | + let wanted: Vec<String> = parsed.children.iter().map(|d| d.digest.clone()).collect(); | |
| 200 | + let present: Vec<(String, i64)> = | |
| 201 | + sqlx::query_as("select digest, total_size from manifests where package_id = $1 and digest = any($2)") | |
| 202 | + .bind(package.id) | |
| 203 | + .bind(&wanted) | |
| 204 | + .fetch_all(&state.db) | |
| 205 | + .await?; | |
| 206 | + if let Some(missing) = wanted.iter().find(|d| !present.iter().any(|(p, _)| p == *d)) { | |
| 207 | + return Err(RegError::new(StatusCode::BAD_REQUEST, "MANIFEST_BLOB_UNKNOWN", "manifest unknown to registry") | |
| 208 | + .detail(json!({ "digest": missing }))); | |
| 209 | + } | |
| 210 | + total_size = present.iter().map(|(_, size)| size).sum(); | |
| 211 | + } | |
| 212 | + | |
| 213 | + let mut tx = state.db.begin().await?; | |
| 214 | + let manifest_id: i64 = sqlx::query_scalar( | |
| 215 | + "insert into manifests (package_id, digest, media_type, content, total_size) values ($1, $2, $3, $4, $5) | |
| 216 | + on conflict (package_id, digest) do update set media_type = excluded.media_type, total_size = excluded.total_size | |
| 217 | + returning id", | |
| 218 | + ) | |
| 219 | + .bind(package.id) | |
| 220 | + .bind(&digest) | |
| 221 | + .bind(&parsed.media_type) | |
| 222 | + .bind(body.as_ref()) | |
| 223 | + .bind(total_size) | |
| 224 | + .fetch_one(&mut *tx) | |
| 225 | + .await?; | |
| 226 | + let (mut refs, mut kinds) = (Vec::new(), Vec::new()); | |
| 227 | + for blob in &parsed.blobs { | |
| 228 | + refs.push(blob.digest.clone()); | |
| 229 | + kinds.push("blob".to_string()); | |
| 230 | + } | |
| 231 | + for child in &parsed.children { | |
| 232 | + refs.push(child.digest.clone()); | |
| 233 | + kinds.push("manifest".to_string()); | |
| 234 | + } | |
| 235 | + sqlx::query( | |
| 236 | + "insert into manifest_references (manifest_id, digest, kind) select $1, * from unnest($2::text[], $3::text[]) on conflict do nothing", | |
| 237 | + ) | |
| 238 | + .bind(manifest_id) | |
| 239 | + .bind(&refs) | |
| 240 | + .bind(&kinds) | |
| 241 | + .execute(&mut *tx) | |
| 242 | + .await?; | |
| 243 | + if by_tag { | |
| 244 | + sqlx::query( | |
| 245 | + "insert into tags (package_id, name, manifest_id) values ($1, $2, $3) | |
| 246 | + on conflict (package_id, name) do update set manifest_id = excluded.manifest_id, updated_at = now()", | |
| 247 | + ) | |
| 248 | + .bind(package.id) | |
| 249 | + .bind(reference) | |
| 250 | + .bind(manifest_id) | |
| 251 | + .execute(&mut *tx) | |
| 252 | + .await?; | |
| 253 | + } | |
| 254 | + sqlx::query("update packages set updated_at = now() where id = $1").bind(package.id).execute(&mut *tx).await?; | |
| 255 | + tx.commit().await?; | |
| 256 | + | |
| 257 | + tracing::info!( | |
| 258 | + package = %package.full_name(), | |
| 259 | + reference, | |
| 260 | + digest, | |
| 261 | + media_type = %parsed.media_type, | |
| 262 | + size = total_size, | |
| 263 | + user = caller.user_name().unwrap_or("-"), | |
| 264 | + "manifest stored" | |
| 265 | + ); | |
| 266 | + if by_tag { | |
| 267 | + analytics::track( | |
| 268 | + state, | |
| 269 | + "image_pushed", | |
| 270 | + caller.user_name(), | |
| 271 | + &package.url(), | |
| 272 | + json!({ "image": package.full_name(), "tag": reference, "index": parsed.is_index(), "size": total_size }), | |
| 273 | + ); | |
| 274 | + } | |
| 275 | + | |
| 276 | + let mut response = StatusCode::CREATED.into_response(); | |
| 277 | + let headers = response.headers_mut(); | |
| 278 | + if let Ok(location) = HeaderValue::from_str(&format!("/v2/{name}/manifests/{digest}")) { | |
| 279 | + headers.insert(header::LOCATION, location); | |
| 280 | + } | |
| 281 | + headers.insert("docker-content-digest", HeaderValue::from_str(&digest).expect("ascii")); | |
| 282 | + headers.insert(header::CONTENT_LENGTH, HeaderValue::from_static("0")); | |
| 283 | + Ok(response) | |
| 284 | +} | |
| 285 | + | |
| 286 | +/// DELETE by digest removes the manifest (and its tags); by tag, the tag. | |
| 287 | +pub async fn delete(state: &AppState, caller: &Caller, name: &str, reference: &str) -> RegResult<Response> { | |
| 288 | + let target = authorize(state, caller, name, "delete").await?; | |
| 289 | + let package = target.package()?; | |
| 290 | + let removed = if reference.starts_with("sha256:") { | |
| 291 | + sqlx::query("delete from manifests where package_id = $1 and digest = $2") | |
| 292 | + .bind(package.id) | |
| 293 | + .bind(reference) | |
| 294 | + .execute(&state.db) | |
| 295 | + .await? | |
| 296 | + .rows_affected() | |
| 297 | + } else { | |
| 298 | + sqlx::query("delete from tags where package_id = $1 and name = $2") | |
| 299 | + .bind(package.id) | |
| 300 | + .bind(reference) | |
| 301 | + .execute(&state.db) | |
| 302 | + .await? | |
| 303 | + .rows_affected() | |
| 304 | + }; | |
| 305 | + if removed == 0 { | |
| 306 | + return Err(RegError::manifest_unknown()); | |
| 307 | + } | |
| 308 | + tracing::info!(package = %package.full_name(), reference, user = caller.user_name().unwrap_or("-"), "manifest deleted via API"); | |
| 309 | + audit(&state.db, caller.user_id(), "package.manifest_delete", &package.full_name(), json!({ "reference": reference }), None).await; | |
| 310 | + Ok(StatusCode::ACCEPTED.into_response()) | |
| 311 | +} | |
| 312 | + | |
| 313 | +/// GET /v2/<name>/tags/list?n=&last= | |
| 314 | +pub async fn tags_list(state: &AppState, caller: &Caller, name: &str, request: Request) -> RegResult<Response> { | |
| 315 | + let target = authorize(state, caller, name, "pull").await?; | |
| 316 | + let package = target.package()?; | |
| 317 | + let query = request.uri().query().unwrap_or(""); | |
| 318 | + let mut n: Option<i64> = None; | |
| 319 | + let mut last: Option<String> = None; | |
| 320 | + for (key, value) in url::form_urlencoded::parse(query.as_bytes()) { | |
| 321 | + match key.as_ref() { | |
| 322 | + "n" => n = value.parse().ok().filter(|v: &i64| *v >= 0), | |
| 323 | + "last" => last = Some(value.into_owned()), | |
| 324 | + _ => {} | |
| 325 | + } | |
| 326 | + } | |
| 327 | + let limit = n.unwrap_or(10_000).min(10_000); | |
| 328 | + let tags: Vec<String> = sqlx::query_scalar( | |
| 329 | + "select name from tags where package_id = $1 and ($2::text is null or name collate \"C\" > $2 collate \"C\") | |
| 330 | + order by name collate \"C\" limit $3", | |
| 331 | + ) | |
| 332 | + .bind(package.id) | |
| 333 | + .bind(&last) | |
| 334 | + .bind(limit) | |
| 335 | + .fetch_all(&state.db) | |
| 336 | + .await?; | |
| 337 | + let mut response = Json(json!({ "name": name, "tags": tags })).into_response(); | |
| 338 | + if n.is_some() && tags.len() as i64 == limit && limit > 0 { | |
| 339 | + if let Some(last_tag) = tags.last() { | |
| 340 | + let link = format!("</v2/{name}/tags/list?n={limit}&last={}>; rel=\"next\"", crate::auth::urlencode(last_tag)); | |
| 341 | + if let Ok(value) = HeaderValue::from_str(&link) { | |
| 342 | + response.headers_mut().insert(header::LINK, value); | |
| 343 | + } | |
| 344 | + } | |
| 345 | + } | |
| 346 | + Ok(response) | |
| 347 | +} | |
| 348 | + | |
| 349 | +#[cfg(test)] | |
| 350 | +mod tests { | |
| 351 | + use super::*; | |
| 352 | + | |
| 353 | + fn d(c: char) -> String { | |
| 354 | + format!("sha256:{}", c.to_string().repeat(64)) | |
| 355 | + } | |
| 356 | + | |
| 357 | + #[test] | |
| 358 | + fn parses_docker_image_manifest() { | |
| 359 | + let body = json!({ | |
| 360 | + "schemaVersion": 2, | |
| 361 | + "mediaType": DOCKER_MANIFEST, | |
| 362 | + "config": { "mediaType": "application/vnd.docker.container.image.v1+json", "size": 1470, "digest": d('a') }, | |
| 363 | + "layers": [ | |
| 364 | + { "mediaType": "application/vnd.docker.image.rootfs.diff.tar.gzip", "size": 3000, "digest": d('b') }, | |
| 365 | + { "mediaType": "application/vnd.docker.image.rootfs.foreign.diff.tar.gzip", "size": 9, "digest": d('c'), "urls": ["https://x"] } | |
| 366 | + ] | |
| 367 | + }) | |
| 368 | + .to_string(); | |
| 369 | + let parsed = parse_manifest(body.as_bytes(), Some(DOCKER_MANIFEST)).unwrap(); | |
| 370 | + assert_eq!(parsed.media_type, DOCKER_MANIFEST); | |
| 371 | + assert_eq!(parsed.blobs, vec![Descriptor { digest: d('a'), size: 1470 }, Descriptor { digest: d('b'), size: 3000 }]); | |
| 372 | + assert!(parsed.children.is_empty()); | |
| 373 | + assert!(!parsed.is_index()); | |
| 374 | + } | |
| 375 | + | |
| 376 | + #[test] | |
| 377 | + fn parses_oci_index_and_infers_type() { | |
| 378 | + let body = json!({ | |
| 379 | + "schemaVersion": 2, | |
| 380 | + "manifests": [ | |
| 381 | + { "mediaType": OCI_MANIFEST, "size": 500, "digest": d('1'), "platform": { "os": "linux", "architecture": "amd64" } }, | |
| 382 | + { "mediaType": OCI_MANIFEST, "size": 501, "digest": d('2') } | |
| 383 | + ] | |
| 384 | + }) | |
| 385 | + .to_string(); | |
| 386 | + let parsed = parse_manifest(body.as_bytes(), Some("application/json")).unwrap(); | |
| 387 | + assert_eq!(parsed.media_type, OCI_INDEX); | |
| 388 | + assert!(parsed.is_index()); | |
| 389 | + assert_eq!(parsed.children.len(), 2); | |
| 390 | + } | |
| 391 | + | |
| 392 | + #[test] | |
| 393 | + fn rejects_bad_manifests() { | |
| 394 | + assert!(parse_manifest(b"not json", None).is_err()); | |
| 395 | + assert!(parse_manifest(br#"{"schemaVersion":1,"fsLayers":[]}"#, None).is_err()); | |
| 396 | + assert!(parse_manifest(br#"{"schemaVersion":2}"#, None).is_err()); | |
| 397 | + let bad_digest = json!({ "schemaVersion": 2, "config": { "size": 1, "digest": "md5:abc" }, "layers": [] }).to_string(); | |
| 398 | + assert!(parse_manifest(bad_digest.as_bytes(), None).is_err()); | |
| 399 | + } | |
| 400 | + | |
| 401 | + #[test] | |
| 402 | + fn digests_bytes() { | |
| 403 | + assert_eq!(digest_of(b""), "sha256:e3b0c44298fc1c149afbf4c8996fb92427ae41e4649b934ca495991b7852b855"); | |
| 404 | + } | |
| 405 | +} |
+527-8backend/src/registry/mod.rs
| @@ -1,17 +1,536 @@ | ||
| 1 | -//! OCI distribution API (Docker registry) at /v2/. Stub. | |
| 1 | +//! OCI distribution API (the Docker registry) at /v2/. | |
| 2 | +//! | |
| 3 | +//! Image names are `<owner>/<path>`: the owner is any account, the path is | |
| 4 | +//! one or more lowercase components, so `alice/api` and `acme/tools/builder` | |
| 5 | +//! both work. Layers and configs are stored once in R2 by digest and linked | |
| 6 | +//! to the packages allowed to see them; every blob request checks that link. | |
| 7 | +//! Manifests are small and live in Postgres. | |
| 8 | +//! | |
| 9 | +//! Auth follows the Docker token flow: requests without credentials get a | |
| 10 | +//! 401 pointing at /v2/token, which trades a personal access token (or | |
| 11 | +//! nothing, for public images) for a short-lived signed bearer token. | |
| 2 | 12 | |
| 3 | -use axum::Router; | |
| 13 | +mod blobs; | |
| 14 | +mod gc; | |
| 15 | +mod manifests; | |
| 16 | +mod pages; | |
| 17 | +mod token; | |
| 4 | 18 | |
| 5 | -use crate::state::AppState; | |
| 19 | +use axum::{ | |
| 20 | + Json, Router, | |
| 21 | + extract::{Request, State}, | |
| 22 | + http::{HeaderMap, HeaderValue, Method, StatusCode, header}, | |
| 23 | + response::{IntoResponse, Response}, | |
| 24 | + routing::{any, get}, | |
| 25 | +}; | |
| 26 | +use serde_json::json; | |
| 27 | + | |
| 28 | +use crate::{ | |
| 29 | + analytics, | |
| 30 | + auth::{self, HeaderAuth, Viewer}, | |
| 31 | + models::{self, Account, Package, audit}, | |
| 32 | + perm::{self, Access}, | |
| 33 | + state::AppState, | |
| 34 | +}; | |
| 35 | + | |
| 36 | +// gc_once is the admin panel's "Clean up now". | |
| 37 | +#[allow(unused_imports)] | |
| 38 | +pub use gc::{gc_once, spawn_gc}; | |
| 6 | 39 | |
| 7 | 40 | pub fn router() -> Router<AppState> { |
| 8 | 41 | Router::new() |
| 42 | + .route("/v2", get(base)) | |
| 43 | + .route("/v2/", get(base)) | |
| 44 | + .route("/v2/token", get(token::issue).post(token::issue)) | |
| 45 | + .route("/v2/{*rest}", any(dispatch)) | |
| 46 | + .merge(pages::router()) | |
| 47 | +} | |
| 48 | + | |
| 49 | +// --------------------------------------------------------------------------- | |
| 50 | +// Errors, in the OCI wire format. | |
| 51 | + | |
| 52 | +#[derive(Debug)] | |
| 53 | +pub struct RegError { | |
| 54 | + status: StatusCode, | |
| 55 | + code: &'static str, | |
| 56 | + message: String, | |
| 57 | + detail: Option<serde_json::Value>, | |
| 58 | + /// For 401s: the scope the client should request a token for. | |
| 59 | + challenge_scope: Option<String>, | |
| 60 | +} | |
| 61 | + | |
| 62 | +pub type RegResult<T> = Result<T, RegError>; | |
| 63 | + | |
| 64 | +impl RegError { | |
| 65 | + pub fn new(status: StatusCode, code: &'static str, message: impl Into<String>) -> Self { | |
| 66 | + Self { status, code, message: message.into(), detail: None, challenge_scope: None } | |
| 67 | + } | |
| 68 | + | |
| 69 | + pub fn detail(mut self, detail: serde_json::Value) -> Self { | |
| 70 | + self.detail = Some(detail); | |
| 71 | + self | |
| 72 | + } | |
| 73 | + | |
| 74 | + pub fn unauthorized(scope: Option<String>) -> Self { | |
| 75 | + Self { | |
| 76 | + status: StatusCode::UNAUTHORIZED, | |
| 77 | + code: "UNAUTHORIZED", | |
| 78 | + message: "authentication required".into(), | |
| 79 | + detail: None, | |
| 80 | + challenge_scope: scope, | |
| 81 | + } | |
| 82 | + } | |
| 83 | + | |
| 84 | + pub fn denied() -> Self { | |
| 85 | + Self::new(StatusCode::FORBIDDEN, "DENIED", "requested access to the resource is denied") | |
| 86 | + } | |
| 87 | + | |
| 88 | + pub fn name_unknown() -> Self { | |
| 89 | + Self::new(StatusCode::NOT_FOUND, "NAME_UNKNOWN", "repository name not known to registry") | |
| 90 | + } | |
| 91 | + | |
| 92 | + pub fn name_invalid(message: impl Into<String>) -> Self { | |
| 93 | + Self::new(StatusCode::BAD_REQUEST, "NAME_INVALID", message) | |
| 94 | + } | |
| 95 | + | |
| 96 | + pub fn blob_unknown() -> Self { | |
| 97 | + Self::new(StatusCode::NOT_FOUND, "BLOB_UNKNOWN", "blob unknown to registry") | |
| 98 | + } | |
| 99 | + | |
| 100 | + pub fn upload_unknown() -> Self { | |
| 101 | + Self::new(StatusCode::NOT_FOUND, "BLOB_UPLOAD_UNKNOWN", "blob upload unknown to registry") | |
| 102 | + } | |
| 103 | + | |
| 104 | + pub fn manifest_unknown() -> Self { | |
| 105 | + Self::new(StatusCode::NOT_FOUND, "MANIFEST_UNKNOWN", "manifest unknown") | |
| 106 | + } | |
| 107 | + | |
| 108 | + pub fn digest_invalid(message: impl Into<String>) -> Self { | |
| 109 | + Self::new(StatusCode::BAD_REQUEST, "DIGEST_INVALID", message) | |
| 110 | + } | |
| 111 | + | |
| 112 | + pub fn unsupported() -> Self { | |
| 113 | + Self::new(StatusCode::METHOD_NOT_ALLOWED, "UNSUPPORTED", "the operation is unsupported") | |
| 114 | + } | |
| 115 | + | |
| 116 | + pub fn internal(error: impl std::fmt::Debug) -> Self { | |
| 117 | + tracing::error!(error = ?error, "registry internal error"); | |
| 118 | + Self::new(StatusCode::INTERNAL_SERVER_ERROR, "UNKNOWN", "internal error") | |
| 119 | + } | |
| 120 | +} | |
| 121 | + | |
| 122 | +impl From<sqlx::Error> for RegError { | |
| 123 | + fn from(error: sqlx::Error) -> Self { | |
| 124 | + Self::internal(error) | |
| 125 | + } | |
| 126 | +} | |
| 127 | + | |
| 128 | +impl From<anyhow::Error> for RegError { | |
| 129 | + fn from(error: anyhow::Error) -> Self { | |
| 130 | + Self::internal(error) | |
| 131 | + } | |
| 132 | +} | |
| 133 | + | |
| 134 | +impl From<std::io::Error> for RegError { | |
| 135 | + fn from(error: std::io::Error) -> Self { | |
| 136 | + Self::internal(error) | |
| 137 | + } | |
| 138 | +} | |
| 139 | + | |
| 140 | +/// The realm the 401 challenge points at. | |
| 141 | +fn realm(state: &AppState) -> String { | |
| 142 | + format!("{}/v2/token", state.config.base_url()) | |
| 143 | +} | |
| 144 | + | |
| 145 | +fn challenge(state: &AppState, scope: Option<&str>) -> HeaderValue { | |
| 146 | + let mut value = format!("Bearer realm=\"{}\",service=\"{}\"", realm(state), token::SERVICE); | |
| 147 | + if let Some(scope) = scope { | |
| 148 | + value.push_str(&format!(",scope=\"{scope}\"")); | |
| 149 | + } | |
| 150 | + HeaderValue::from_str(&value).unwrap_or_else(|_| HeaderValue::from_static("Bearer")) | |
| 151 | +} | |
| 152 | + | |
| 153 | +/// Renders an error with the realm filled in (needs state, so not IntoResponse). | |
| 154 | +fn error_response(state: &AppState, error: RegError) -> Response { | |
| 155 | + let mut body = json!({ "code": error.code, "message": error.message }); | |
| 156 | + if let Some(detail) = error.detail { | |
| 157 | + body["detail"] = detail; | |
| 158 | + } | |
| 159 | + let mut response = (error.status, Json(json!({ "errors": [body] }))).into_response(); | |
| 160 | + if error.status == StatusCode::UNAUTHORIZED { | |
| 161 | + response.headers_mut().insert(header::WWW_AUTHENTICATE, challenge(state, error.challenge_scope.as_deref())); | |
| 162 | + } | |
| 163 | + with_api_version(response) | |
| 164 | +} | |
| 165 | + | |
| 166 | +pub fn with_api_version(mut response: Response) -> Response { | |
| 167 | + response | |
| 168 | + .headers_mut() | |
| 169 | + .insert("docker-distribution-api-version", HeaderValue::from_static("registry/2.0")); | |
| 170 | + response | |
| 171 | +} | |
| 172 | + | |
| 173 | +// --------------------------------------------------------------------------- | |
| 174 | +// Paths | |
| 175 | + | |
| 176 | +#[derive(Debug, PartialEq, Eq)] | |
| 177 | +pub enum Route { | |
| 178 | + TagsList { name: String }, | |
| 179 | + Manifest { name: String, reference: String }, | |
| 180 | + UploadStart { name: String }, | |
| 181 | + Upload { name: String, id: String }, | |
| 182 | + Blob { name: String, digest: String }, | |
| 183 | +} | |
| 184 | + | |
| 185 | +/// Splits `/v2/{rest}` into a route. Names may contain slashes, so each form | |
| 186 | +/// is matched on its last occurrence and the tail must be a single segment. | |
| 187 | +pub fn parse_route(rest: &str) -> Option<Route> { | |
| 188 | + let rest = rest.trim_start_matches('/'); | |
| 189 | + let single = |s: &str| !s.is_empty() && !s.contains('/'); | |
| 190 | + | |
| 191 | + if let Some(name) = rest.strip_suffix("/tags/list") { | |
| 192 | + return Some(Route::TagsList { name: name.to_string() }); | |
| 193 | + } | |
| 194 | + if let Some(name) = rest.strip_suffix("/blobs/uploads/").or_else(|| rest.strip_suffix("/blobs/uploads")) { | |
| 195 | + return Some(Route::UploadStart { name: name.to_string() }); | |
| 196 | + } | |
| 197 | + if let Some(i) = rest.rfind("/blobs/uploads/") { | |
| 198 | + let id = &rest[i + "/blobs/uploads/".len()..]; | |
| 199 | + if single(id) { | |
| 200 | + return Some(Route::Upload { name: rest[..i].to_string(), id: id.to_string() }); | |
| 201 | + } | |
| 202 | + } | |
| 203 | + if let Some(i) = rest.rfind("/manifests/") { | |
| 204 | + let reference = &rest[i + "/manifests/".len()..]; | |
| 205 | + if single(reference) { | |
| 206 | + return Some(Route::Manifest { name: rest[..i].to_string(), reference: reference.to_string() }); | |
| 207 | + } | |
| 208 | + } | |
| 209 | + if let Some(i) = rest.rfind("/blobs/") { | |
| 210 | + let digest = &rest[i + "/blobs/".len()..]; | |
| 211 | + if single(digest) { | |
| 212 | + return Some(Route::Blob { name: rest[..i].to_string(), digest: digest.to_string() }); | |
| 213 | + } | |
| 214 | + } | |
| 215 | + None | |
| 216 | +} | |
| 217 | + | |
| 218 | +/// "owner/path/to/image" -> ("owner", "path/to/image"), validated. | |
| 219 | +pub fn split_name(name: &str) -> RegResult<(String, String)> { | |
| 220 | + let (owner, image) = name.split_once('/').ok_or_else(|| RegError::name_invalid("image names look like <owner>/<image>"))?; | |
| 221 | + // Account names are lowercase [a-z0-9-]; reserved words can never exist | |
| 222 | + // as accounts, so only the character set matters here. | |
| 223 | + if owner.is_empty() || owner.len() > 39 || !owner.bytes().all(|b| b.is_ascii_lowercase() || b.is_ascii_digit() || b == b'-') { | |
| 224 | + return Err(RegError::name_invalid("invalid owner name")); | |
| 225 | + } | |
| 226 | + models::valid_package_name(image).map_err(RegError::name_invalid)?; | |
| 227 | + Ok((owner.to_string(), image.to_string())) | |
| 228 | +} | |
| 229 | + | |
| 230 | +pub fn valid_digest(digest: &str) -> bool { | |
| 231 | + digest | |
| 232 | + .strip_prefix("sha256:") | |
| 233 | + .is_some_and(|hex| hex.len() == 64 && hex.bytes().all(|b| b.is_ascii_digit() || (b'a'..=b'f').contains(&b))) | |
| 234 | +} | |
| 235 | + | |
| 236 | +pub fn valid_tag(tag: &str) -> bool { | |
| 237 | + let bytes = tag.as_bytes(); | |
| 238 | + !bytes.is_empty() | |
| 239 | + && bytes.len() <= 128 | |
| 240 | + && (bytes[0].is_ascii_alphanumeric() || bytes[0] == b'_') | |
| 241 | + && bytes.iter().all(|b| b.is_ascii_alphanumeric() || matches!(b, b'_' | b'.' | b'-')) | |
| 242 | +} | |
| 243 | + | |
| 244 | +// --------------------------------------------------------------------------- | |
| 245 | +// Who is calling, and what they may do. | |
| 246 | + | |
| 247 | +#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)] | |
| 248 | +pub struct Actions { | |
| 249 | + pub pull: bool, | |
| 250 | + pub push: bool, | |
| 251 | + pub delete: bool, | |
| 252 | +} | |
| 253 | + | |
| 254 | +impl Actions { | |
| 255 | + pub fn allows(&self, action: &str) -> bool { | |
| 256 | + match action { | |
| 257 | + "pull" => self.pull, | |
| 258 | + "push" => self.push, | |
| 259 | + "delete" | "*" => self.delete, | |
| 260 | + _ => false, | |
| 261 | + } | |
| 262 | + } | |
| 263 | + | |
| 264 | + pub fn list(&self) -> Vec<String> { | |
| 265 | + let mut out = Vec::new(); | |
| 266 | + if self.pull { | |
| 267 | + out.push("pull".to_string()); | |
| 268 | + } | |
| 269 | + if self.push { | |
| 270 | + out.push("push".to_string()); | |
| 271 | + } | |
| 272 | + if self.delete { | |
| 273 | + out.push("delete".to_string()); | |
| 274 | + } | |
| 275 | + out | |
| 276 | + } | |
| 277 | +} | |
| 278 | + | |
| 279 | +/// The caller of a /v2 request. | |
| 280 | +pub enum Caller { | |
| 281 | + Anonymous, | |
| 282 | + /// A bearer token minted by /v2/token: grants are fixed in it. | |
| 283 | + Token(token::Claims), | |
| 284 | + /// A personal access token sent directly (Basic or Bearer igp_...). | |
| 285 | + Pat(Viewer), | |
| 286 | +} | |
| 287 | + | |
| 288 | +impl Caller { | |
| 289 | + pub fn user_id(&self) -> Option<i64> { | |
| 290 | + match self { | |
| 291 | + Caller::Anonymous => None, | |
| 292 | + Caller::Token(claims) => claims.uid, | |
| 293 | + Caller::Pat(viewer) => Some(viewer.id), | |
| 294 | + } | |
| 295 | + } | |
| 296 | + | |
| 297 | + pub fn user_name(&self) -> Option<&str> { | |
| 298 | + match self { | |
| 299 | + Caller::Anonymous => None, | |
| 300 | + Caller::Token(claims) => claims.user.as_deref(), | |
| 301 | + Caller::Pat(viewer) => Some(&viewer.name), | |
| 302 | + } | |
| 303 | + } | |
| 304 | + | |
| 305 | + fn is_anonymous(&self) -> bool { | |
| 306 | + matches!(self, Caller::Anonymous) || matches!(self, Caller::Token(c) if c.uid.is_none()) | |
| 307 | + } | |
| 308 | +} | |
| 309 | + | |
| 310 | +pub enum CallerAuth { | |
| 311 | + Caller(Caller), | |
| 312 | + Invalid, | |
| 313 | +} | |
| 314 | + | |
| 315 | +/// Reads the Authorization header: registry bearer tokens first, then PATs. | |
| 316 | +pub async fn caller(state: &AppState, headers: &HeaderMap) -> RegResult<CallerAuth> { | |
| 317 | + if let Some(value) = headers.get(header::AUTHORIZATION).and_then(|v| v.to_str().ok()) { | |
| 318 | + if let Some(bearer) = value.strip_prefix("Bearer ").or_else(|| value.strip_prefix("bearer ")) { | |
| 319 | + let bearer = bearer.trim(); | |
| 320 | + if !bearer.starts_with(auth::TOKEN_PREFIX) { | |
| 321 | + return Ok(match token::verify(state, bearer) { | |
| 322 | + Some(claims) => CallerAuth::Caller(Caller::Token(claims)), | |
| 323 | + None => CallerAuth::Invalid, | |
| 324 | + }); | |
| 325 | + } | |
| 326 | + } | |
| 327 | + } | |
| 328 | + Ok(match auth::authenticate_header(&state.db, headers).await? { | |
| 329 | + HeaderAuth::Anonymous => CallerAuth::Caller(Caller::Anonymous), | |
| 330 | + HeaderAuth::Invalid => CallerAuth::Invalid, | |
| 331 | + HeaderAuth::Viewer(viewer) => CallerAuth::Caller(Caller::Pat(viewer)), | |
| 332 | + }) | |
| 333 | +} | |
| 334 | + | |
| 335 | +/// What `viewer` may do with `owner/image`, whether or not it exists yet. | |
| 336 | +/// A missing image can be pushed (created) by anyone who may write under | |
| 337 | +/// the owner; they may also "pull" it so they see a clean 404. | |
| 338 | +pub async fn actions_for(state: &AppState, viewer: Option<&Viewer>, owner: &Account, package: Option<&Package>) -> RegResult<Actions> { | |
| 339 | + let viewer = viewer.filter(|v| v.scopes.packages); | |
| 340 | + match package { | |
| 341 | + Some(package) => { | |
| 342 | + let access = perm::package_access(&state.db, viewer, package).await?; | |
| 343 | + Ok(Actions { pull: access.can_read(), push: access.can_write(), delete: access.is_admin() }) | |
| 344 | + } | |
| 345 | + None => { | |
| 346 | + let access = perm::owner_access(&state.db, viewer, owner.id).await?; | |
| 347 | + Ok(Actions { pull: access.can_write(), push: access.can_write(), delete: access == Access::Admin }) | |
| 348 | + } | |
| 349 | + } | |
| 350 | +} | |
| 351 | + | |
| 352 | +/// A resolved image name plus the caller's rights on it. | |
| 353 | +pub struct Target { | |
| 354 | + pub name: String, | |
| 355 | + pub owner: Account, | |
| 356 | + pub image: String, | |
| 357 | + pub package: Option<Package>, | |
| 358 | + pub actions: Actions, | |
| 359 | +} | |
| 360 | + | |
| 361 | +impl Target { | |
| 362 | + pub fn package(&self) -> RegResult<&Package> { | |
| 363 | + self.package.as_ref().ok_or_else(RegError::name_unknown) | |
| 364 | + } | |
| 365 | +} | |
| 366 | + | |
| 367 | +/// Resolves a name and checks `action` for the caller, answering 401 for | |
| 368 | +/// anonymous callers (so docker fetches a token) and 403 otherwise. | |
| 369 | +pub async fn authorize(state: &AppState, caller: &Caller, name: &str, action: &str) -> RegResult<Target> { | |
| 370 | + let (owner_name, image) = split_name(name)?; | |
| 371 | + let scope = format!("repository:{name}:{}", if action == "pull" { "pull" } else { "pull,push" }); | |
| 372 | + let owner = Account::by_name(&state.db, &owner_name).await?; | |
| 373 | + let Some(owner) = owner else { | |
| 374 | + return Err(if caller.is_anonymous() { RegError::unauthorized(Some(scope)) } else { RegError::name_unknown() }); | |
| 375 | + }; | |
| 376 | + let package = Package::by_path(&state.db, &owner_name, &image).await?; | |
| 377 | + let actions = match caller { | |
| 378 | + Caller::Token(claims) => claims.actions_for(name), | |
| 379 | + Caller::Pat(viewer) => actions_for(state, Some(viewer), &owner, package.as_ref()).await?, | |
| 380 | + Caller::Anonymous => actions_for(state, None, &owner, package.as_ref()).await?, | |
| 381 | + }; | |
| 382 | + if actions.allows(action) { | |
| 383 | + return Ok(Target { name: name.to_string(), owner, image, package, actions }); | |
| 384 | + } | |
| 385 | + if caller.is_anonymous() { | |
| 386 | + return Err(RegError::unauthorized(Some(scope))); | |
| 387 | + } | |
| 388 | + Err(RegError::denied()) | |
| 389 | +} | |
| 390 | + | |
| 391 | +/// Finds or creates the package a push writes to. New images are private. | |
| 392 | +pub async fn ensure_package(state: &AppState, target: &Target, caller: &Caller) -> RegResult<Package> { | |
| 393 | + if let Some(package) = &target.package { | |
| 394 | + return Ok(package.clone()); | |
| 395 | + } | |
| 396 | + let inserted: Option<i64> = sqlx::query_scalar( | |
| 397 | + "insert into packages (owner_id, name) values ($1, $2) on conflict (owner_id, name) do nothing returning id", | |
| 398 | + ) | |
| 399 | + .bind(target.owner.id) | |
| 400 | + .bind(&target.image) | |
| 401 | + .fetch_optional(&state.db) | |
| 402 | + .await?; | |
| 403 | + let package = Package::by_path(&state.db, &target.owner.name, &target.image) | |
| 404 | + .await? | |
| 405 | + .ok_or_else(|| RegError::internal("package vanished after insert"))?; | |
| 406 | + if inserted.is_some() { | |
| 407 | + tracing::info!(package = %package.full_name(), user = caller.user_name().unwrap_or("-"), "package created by push"); | |
| 408 | + audit(&state.db, caller.user_id(), "package.create", &package.full_name(), json!({ "via": "push" }), None).await; | |
| 409 | + analytics::track(state, "package_created", caller.user_name(), &package.url(), json!({ "owner_kind": target.owner.kind })); | |
| 410 | + } | |
| 411 | + Ok(package) | |
| 412 | +} | |
| 413 | + | |
| 414 | +// --------------------------------------------------------------------------- | |
| 415 | +// Handlers | |
| 416 | + | |
| 417 | +/// GET /v2/: 200 for any valid credential, otherwise the token challenge. | |
| 418 | +async fn base(State(state): State<AppState>, headers: HeaderMap) -> Response { | |
| 419 | + let result = async { | |
| 420 | + match caller(&state, &headers).await? { | |
| 421 | + CallerAuth::Caller(Caller::Anonymous) | CallerAuth::Invalid => Err(RegError::unauthorized(None)), | |
| 422 | + CallerAuth::Caller(_) => Ok(()), | |
| 423 | + } | |
| 424 | + } | |
| 425 | + .await; | |
| 426 | + match result { | |
| 427 | + Ok(()) => with_api_version((StatusCode::OK, Json(json!({}))).into_response()), | |
| 428 | + Err(error) => error_response(&state, error), | |
| 429 | + } | |
| 430 | +} | |
| 431 | + | |
| 432 | +async fn dispatch(State(state): State<AppState>, request: Request) -> Response { | |
| 433 | + let method = request.method().clone(); | |
| 434 | + let path = request.uri().path().to_string(); | |
| 435 | + let rest = path.strip_prefix("/v2/").unwrap_or(""); | |
| 436 | + let started = std::time::Instant::now(); | |
| 437 | + | |
| 438 | + let result = async { | |
| 439 | + let caller = match caller(&state, request.headers()).await? { | |
| 440 | + CallerAuth::Caller(caller) => caller, | |
| 441 | + CallerAuth::Invalid => return Err(RegError::unauthorized(None)), | |
| 442 | + }; | |
| 443 | + let route = parse_route(rest).ok_or_else(|| RegError::new(StatusCode::NOT_FOUND, "NOT_FOUND", "unknown registry endpoint"))?; | |
| 444 | + match (&method, route) { | |
| 445 | + (&Method::GET, Route::TagsList { name }) => manifests::tags_list(&state, &caller, &name, request).await, | |
| 446 | + (&Method::GET | &Method::HEAD, Route::Manifest { name, reference }) => { | |
| 447 | + manifests::get(&state, &caller, &name, &reference, method == Method::HEAD, request).await | |
| 448 | + } | |
| 449 | + (&Method::PUT, Route::Manifest { name, reference }) => manifests::put(&state, &caller, &name, &reference, request).await, | |
| 450 | + (&Method::DELETE, Route::Manifest { name, reference }) => manifests::delete(&state, &caller, &name, &reference).await, | |
| 451 | + (&Method::GET | &Method::HEAD, Route::Blob { name, digest }) => { | |
| 452 | + blobs::get(&state, &caller, &name, &digest, method == Method::HEAD).await | |
| 453 | + } | |
| 454 | + (&Method::DELETE, Route::Blob { name, digest }) => blobs::delete(&state, &caller, &name, &digest).await, | |
| 455 | + (&Method::POST, Route::UploadStart { name }) => blobs::start_upload(&state, &caller, &name, request).await, | |
| 456 | + (&Method::PATCH, Route::Upload { name, id }) => blobs::patch_upload(&state, &caller, &name, &id, request).await, | |
| 457 | + (&Method::PUT, Route::Upload { name, id }) => blobs::finish_upload(&state, &caller, &name, &id, request).await, | |
| 458 | + (&Method::GET, Route::Upload { name, id }) => blobs::upload_status(&state, &caller, &name, &id).await, | |
| 459 | + (&Method::DELETE, Route::Upload { name, id }) => blobs::cancel_upload(&state, &caller, &name, &id).await, | |
| 460 | + _ => Err(RegError::unsupported()), | |
| 461 | + } | |
| 462 | + } | |
| 463 | + .await; | |
| 464 | + | |
| 465 | + let response = match result { | |
| 466 | + Ok(response) => with_api_version(response), | |
| 467 | + Err(error) => { | |
| 468 | + if error.status.is_server_error() { | |
| 469 | + tracing::error!(%method, %path, code = error.code, "registry request failed"); | |
| 470 | + } else { | |
| 471 | + tracing::debug!(%method, %path, code = error.code, status = error.status.as_u16(), "registry request refused"); | |
| 472 | + } | |
| 473 | + error_response(&state, error) | |
| 474 | + } | |
| 475 | + }; | |
| 476 | + tracing::debug!(%method, %path, status = response.status().as_u16(), ms = started.elapsed().as_millis() as u64, "registry"); | |
| 477 | + response | |
| 9 | 478 | } |
| 10 | 479 | |
| 11 | -/// Starts periodic garbage collection of orphaned blobs and stale uploads. | |
| 12 | -pub fn spawn_gc(_state: AppState) {} | |
| 480 | +#[cfg(test)] | |
| 481 | +mod tests { | |
| 482 | + use super::*; | |
| 483 | + | |
| 484 | + #[test] | |
| 485 | + fn parses_routes_with_nested_names() { | |
| 486 | + assert_eq!(parse_route("alice/api/tags/list"), Some(Route::TagsList { name: "alice/api".into() })); | |
| 487 | + assert_eq!( | |
| 488 | + parse_route("acme/tools/builder/manifests/v1.2"), | |
| 489 | + Some(Route::Manifest { name: "acme/tools/builder".into(), reference: "v1.2".into() }) | |
| 490 | + ); | |
| 491 | + assert_eq!(parse_route("alice/api/blobs/uploads/"), Some(Route::UploadStart { name: "alice/api".into() })); | |
| 492 | + assert_eq!(parse_route("alice/api/blobs/uploads"), Some(Route::UploadStart { name: "alice/api".into() })); | |
| 493 | + assert_eq!( | |
| 494 | + parse_route("alice/api/blobs/uploads/0b7f"), | |
| 495 | + Some(Route::Upload { name: "alice/api".into(), id: "0b7f".into() }) | |
| 496 | + ); | |
| 497 | + assert_eq!( | |
| 498 | + parse_route("alice/api/blobs/sha256:abc"), | |
| 499 | + Some(Route::Blob { name: "alice/api".into(), digest: "sha256:abc".into() }) | |
| 500 | + ); | |
| 501 | + // A name component called "manifests" or "blobs" still parses. | |
| 502 | + assert_eq!( | |
| 503 | + parse_route("alice/manifests/blobs/sha256:abc"), | |
| 504 | + Some(Route::Blob { name: "alice/manifests".into(), digest: "sha256:abc".into() }) | |
| 505 | + ); | |
| 506 | + assert_eq!( | |
| 507 | + parse_route("alice/blobs/manifests/latest"), | |
| 508 | + Some(Route::Manifest { name: "alice/blobs".into(), reference: "latest".into() }) | |
| 509 | + ); | |
| 510 | + assert_eq!(parse_route("alice/api"), None); | |
| 511 | + assert_eq!(parse_route(""), None); | |
| 512 | + } | |
| 513 | + | |
| 514 | + #[test] | |
| 515 | + fn splits_and_validates_names() { | |
| 516 | + assert_eq!(split_name("alice/api").unwrap(), ("alice".into(), "api".into())); | |
| 517 | + assert_eq!(split_name("acme/tools/builder").unwrap(), ("acme".into(), "tools/builder".into())); | |
| 518 | + assert!(split_name("alice").is_err()); | |
| 519 | + assert!(split_name("alice/API").is_err()); | |
| 520 | + assert!(split_name("Alice/api").is_err()); | |
| 521 | + assert!(split_name("alice/a//b").is_err()); | |
| 522 | + } | |
| 13 | 523 | |
| 14 | -/// One garbage collection pass (admin "Clean up now"). Returns a summary. | |
| 15 | -pub async fn gc_once(_state: &AppState) -> anyhow::Result<String> { | |
| 16 | - anyhow::bail!("registry garbage collection is not implemented yet") | |
| 524 | + #[test] | |
| 525 | + fn digests_and_tags() { | |
| 526 | + assert!(valid_digest(&format!("sha256:{}", "a".repeat(64)))); | |
| 527 | + assert!(!valid_digest(&format!("sha256:{}", "A".repeat(64)))); | |
| 528 | + assert!(!valid_digest("sha256:abc")); | |
| 529 | + assert!(!valid_digest(&format!("sha512:{}", "a".repeat(128)))); | |
| 530 | + assert!(valid_tag("latest")); | |
| 531 | + assert!(valid_tag("v1.2.3-rc_1")); | |
| 532 | + assert!(!valid_tag(".hidden")); | |
| 533 | + assert!(!valid_tag("-x")); | |
| 534 | + assert!(!valid_tag(&"a".repeat(129))); | |
| 535 | + } | |
| 17 | 536 | } |
+577-0backend/src/registry/pages.rs
| @@ -0,0 +1,577 @@ | ||
| 1 | +//! Web pages for container images: | |
| 2 | +//! /{owner}/-/packages images an account owns | |
| 3 | +//! /{owner}/-/packages/{name} one image: pull command, tags, settings | |
| 4 | +//! | |
| 5 | +//! Settings post back to the same URL with an `action` field. | |
| 6 | + | |
| 7 | +use axum::{ | |
| 8 | + Form, Router, | |
| 9 | + extract::{Path, Query}, | |
| 10 | + http::{HeaderMap, StatusCode, header}, | |
| 11 | + response::{IntoResponse, Redirect, Response}, | |
| 12 | + routing::get, | |
| 13 | +}; | |
| 14 | +use chrono::{DateTime, Utc}; | |
| 15 | +use maud::{Markup, html}; | |
| 16 | +use serde::Deserialize; | |
| 17 | +use serde_json::{Value, json}; | |
| 18 | + | |
| 19 | +use super::manifests::{DOCKER_LIST, OCI_INDEX}; | |
| 20 | +use crate::{ | |
| 21 | + analytics, | |
| 22 | + auth::{self, Viewer}, | |
| 23 | + error::{AppError, AppResult}, | |
| 24 | + models::{Account, Package, Repo, audit}, | |
| 25 | + package_select, | |
| 26 | + perm::{self, Access, Area}, | |
| 27 | + state::AppState, | |
| 28 | + visible_packages, | |
| 29 | + web::{ | |
| 30 | + layout::{Ctx, Page}, | |
| 31 | + ui, | |
| 32 | + }, | |
| 33 | +}; | |
| 34 | + | |
| 35 | +pub fn router() -> Router<AppState> { | |
| 36 | + Router::new() | |
| 37 | + .route("/{owner}/-/packages", get(list)) | |
| 38 | + .route("/{owner}/-/packages/", get(list)) | |
| 39 | + .route("/{owner}/-/packages/{*name}", get(detail).post(act)) | |
| 40 | +} | |
| 41 | + | |
| 42 | +// --------------------------------------------------------------------------- | |
| 43 | +// List | |
| 44 | + | |
| 45 | +async fn list(ctx: Ctx, Path(owner): Path<String>) -> AppResult<Response> { | |
| 46 | + let owner = Account::by_name(&ctx.state.db, &owner).await?.ok_or(AppError::NotFound)?; | |
| 47 | + let (viewer_id, admin) = perm::visibility_binds(ctx.viewer.as_ref(), Area::Packages); | |
| 48 | + let packages: Vec<Package> = sqlx::query_as(concat!( | |
| 49 | + package_select!(), | |
| 50 | + " where p.owner_id = $3 and ", | |
| 51 | + visible_packages!(), | |
| 52 | + " order by p.updated_at desc limit 500" | |
| 53 | + )) | |
| 54 | + .bind(viewer_id) | |
| 55 | + .bind(admin) | |
| 56 | + .bind(owner.id) | |
| 57 | + .fetch_all(&ctx.state.db) | |
| 58 | + .await?; | |
| 59 | + let ids: Vec<i64> = packages.iter().map(|p| p.id).collect(); | |
| 60 | + let tag_counts: Vec<(i64, i64)> = sqlx::query_as("select package_id, count(*) from tags where package_id = any($1) group by package_id") | |
| 61 | + .bind(&ids) | |
| 62 | + .fetch_all(&ctx.state.db) | |
| 63 | + .await?; | |
| 64 | + let can_push = match &ctx.viewer { | |
| 65 | + Some(viewer) => perm::owner_access(&ctx.state.db, Some(viewer), owner.id).await?.can_write(), | |
| 66 | + None => false, | |
| 67 | + }; | |
| 68 | + let host = ctx.state.config.registry_host(); | |
| 69 | + | |
| 70 | + let body = html! { | |
| 71 | + div class="mx-auto max-w-[1100px] px-4 py-5" { | |
| 72 | + div class="mb-4 flex items-center gap-2" { | |
| 73 | + (ui::avatar(&owner.name, owner.avatar_key.as_deref(), 28)) | |
| 74 | + a href={ "/" (owner.name) } class="text-[15px] font-semibold text-ink" { (owner.label()) } | |
| 75 | + span class="text-ink-faint" { "/" } | |
| 76 | + span class="text-[15px] font-semibold" { "Images" } | |
| 77 | + span class="ml-1 text-xs text-ink-faint" { (ui::plural(packages.len() as i64, "image", "images")) } | |
| 78 | + } | |
| 79 | + @if packages.is_empty() { | |
| 80 | + (ui::empty_state("No images yet", html! { | |
| 81 | + @if can_push { | |
| 82 | + p { "Push one with " code class="font-mono text-xs" { "docker push " (host) "/" (owner.name) "/<image>:<tag>" } "." } | |
| 83 | + p class="mt-1" { a href="/docs/registry" { "How the registry works" } } | |
| 84 | + } @else { | |
| 85 | + p { (owner.name) " has no images you can see." } | |
| 86 | + } | |
| 87 | + })) | |
| 88 | + } @else { | |
| 89 | + div class="box divide-y divide-edge" { | |
| 90 | + @for package in &packages { | |
| 91 | + @let tags = tag_counts.iter().find(|(id, _)| *id == package.id).map(|(_, n)| *n).unwrap_or(0); | |
| 92 | + div class="flex flex-wrap items-center gap-x-3 gap-y-1 px-3 py-2" { | |
| 93 | + (ui::icon_package()) | |
| 94 | + a href=(package.url()) class="font-semibold" data-track="package_opened" { (package.name) } | |
| 95 | + (ui::visibility_tag(&package.visibility)) | |
| 96 | + @if !package.description.is_empty() { | |
| 97 | + span class="min-w-0 flex-1 truncate text-ink-dim" { (package.description) } | |
| 98 | + } @else { | |
| 99 | + span class="flex-1" {} | |
| 100 | + } | |
| 101 | + span class="text-xs text-ink-faint" { (ui::plural(tags, "tag", "tags")) } | |
| 102 | + span class="text-xs text-ink-faint" { (ui::plural(package.pull_count, "pull", "pulls")) } | |
| 103 | + span class="text-xs text-ink-faint" { "updated " (ui::time(package.updated_at)) } | |
| 104 | + } | |
| 105 | + } | |
| 106 | + } | |
| 107 | + } | |
| 108 | + } | |
| 109 | + }; | |
| 110 | + Ok(ctx.render(Page::new(format!("{} images", owner.name), body).description(format!("Container images published by {}.", owner.name)))) | |
| 111 | +} | |
| 112 | + | |
| 113 | +// --------------------------------------------------------------------------- | |
| 114 | +// Detail | |
| 115 | + | |
| 116 | +#[derive(Deserialize, Default)] | |
| 117 | +struct DetailQuery { | |
| 118 | + tab: Option<String>, | |
| 119 | + saved: Option<String>, | |
| 120 | +} | |
| 121 | + | |
| 122 | +#[derive(sqlx::FromRow)] | |
| 123 | +struct TagRow { | |
| 124 | + name: String, | |
| 125 | + digest: String, | |
| 126 | + media_type: String, | |
| 127 | + total_size: i64, | |
| 128 | + updated_at: DateTime<Utc>, | |
| 129 | + content: Vec<u8>, | |
| 130 | +} | |
| 131 | + | |
| 132 | +struct Loaded { | |
| 133 | + owner: Account, | |
| 134 | + package: Package, | |
| 135 | + access: Access, | |
| 136 | +} | |
| 137 | + | |
| 138 | +async fn load(ctx: &Ctx, owner: &str, name: &str) -> AppResult<Loaded> { | |
| 139 | + let name = name.trim_end_matches('/'); | |
| 140 | + let owner = Account::by_name(&ctx.state.db, owner).await?.ok_or(AppError::NotFound)?; | |
| 141 | + let package = Package::by_path(&ctx.state.db, &owner.name, name).await?.ok_or(AppError::NotFound)?; | |
| 142 | + let access = perm::package_access(&ctx.state.db, ctx.viewer.as_ref(), &package).await?; | |
| 143 | + if !access.can_read() { | |
| 144 | + return Err(AppError::NotFound); | |
| 145 | + } | |
| 146 | + Ok(Loaded { owner, package, access }) | |
| 147 | +} | |
| 148 | + | |
| 149 | +async fn detail(ctx: Ctx, Path((owner, name)): Path<(String, String)>, Query(query): Query<DetailQuery>) -> AppResult<Response> { | |
| 150 | + let loaded = load(&ctx, &owner, &name).await?; | |
| 151 | + let saved = query.saved.as_deref().map(|what| match what { | |
| 152 | + "visibility" => "Visibility updated.", | |
| 153 | + "description" => "Description saved.", | |
| 154 | + "repo" => "Linked repository updated.", | |
| 155 | + "tag" => "Tag deleted.", | |
| 156 | + _ => "Saved.", | |
| 157 | + }); | |
| 158 | + render_detail(&ctx, loaded, query.tab.as_deref(), saved, None).await | |
| 159 | +} | |
| 160 | + | |
| 161 | +async fn render_detail(ctx: &Ctx, loaded: Loaded, tab: Option<&str>, saved: Option<&str>, error: Option<&str>) -> AppResult<Response> { | |
| 162 | + let Loaded { owner, package, access } = loaded; | |
| 163 | + let state = &ctx.state; | |
| 164 | + let is_admin = access.is_admin(); | |
| 165 | + let settings_tab = is_admin && tab == Some("settings"); | |
| 166 | + | |
| 167 | + let tags: Vec<TagRow> = sqlx::query_as( | |
| 168 | + "select t.name, m.digest, m.media_type, m.total_size, t.updated_at, m.content | |
| 169 | + from tags t join manifests m on m.id = t.manifest_id | |
| 170 | + where t.package_id = $1 order by t.updated_at desc, t.name limit 300", | |
| 171 | + ) | |
| 172 | + .bind(package.id) | |
| 173 | + .fetch_all(&state.db) | |
| 174 | + .await?; | |
| 175 | + let untagged: i64 = sqlx::query_scalar( | |
| 176 | + "select count(*) from manifests m where m.package_id = $1 | |
| 177 | + and not exists (select 1 from tags t where t.manifest_id = m.id) | |
| 178 | + and not exists (select 1 from manifest_references r join manifests p on p.id = r.manifest_id | |
| 179 | + where p.package_id = m.package_id and r.digest = m.digest)", | |
| 180 | + ) | |
| 181 | + .bind(package.id) | |
| 182 | + .fetch_one(&state.db) | |
| 183 | + .await?; | |
| 184 | + let linked_repo = match package.repo_id { | |
| 185 | + Some(id) => Repo::by_id(&state.db, id).await?, | |
| 186 | + None => None, | |
| 187 | + }; | |
| 188 | + // Only show the link to people who may see the repo. | |
| 189 | + let linked_repo = match linked_repo { | |
| 190 | + Some(repo) if perm::repo_access(&state.db, ctx.viewer.as_ref(), &repo, Area::Repo).await?.can_read() => Some(repo), | |
| 191 | + _ => None, | |
| 192 | + }; | |
| 193 | + | |
| 194 | + let host = state.config.registry_host(); | |
| 195 | + let image_ref = format!("{host}/{}", package.full_name()); | |
| 196 | + let latest_tag = tags.iter().find(|t| t.name == "latest").or_else(|| tags.first()).map(|t| t.name.clone()); | |
| 197 | + let pull = match &latest_tag { | |
| 198 | + Some(tag) => format!("docker pull {image_ref}:{tag}"), | |
| 199 | + None => format!("docker pull {image_ref}"), | |
| 200 | + }; | |
| 201 | + let base_url = package.url(); | |
| 202 | + | |
| 203 | + let body = html! { | |
| 204 | + div class="mx-auto max-w-[1100px] px-4 py-5" { | |
| 205 | + div class="flex flex-wrap items-center gap-2" { | |
| 206 | + (ui::avatar(&owner.name, owner.avatar_key.as_deref(), 24)) | |
| 207 | + a href={ "/" (owner.name) } class="text-[15px] text-ink" { (owner.name) } | |
| 208 | + span class="text-ink-faint" { "/" } | |
| 209 | + a href={ "/" (owner.name) "/-/packages" } class="text-[15px] text-ink-dim" { "images" } | |
| 210 | + span class="text-ink-faint" { "/" } | |
| 211 | + a href=(base_url) class="text-[15px] font-semibold text-ink" { (package.name) } | |
| 212 | + (ui::visibility_tag(&package.visibility)) | |
| 213 | + } | |
| 214 | + @if !package.description.is_empty() { | |
| 215 | + p class="mt-1.5 text-ink-dim" { (package.description) } | |
| 216 | + } | |
| 217 | + | |
| 218 | + nav class="tabs mt-4" { | |
| 219 | + a class="tab" href=(base_url) aria-current=[(!settings_tab).then_some("page")] { | |
| 220 | + "Tags" span class="text-xs text-ink-faint" { (tags.len()) } | |
| 221 | + } | |
| 222 | + @if is_admin { | |
| 223 | + a class="tab" href={ (base_url) "?tab=settings" } aria-current=[settings_tab.then_some("page")] data-track="package_settings_opened" { "Settings" } | |
| 224 | + } | |
| 225 | + } | |
| 226 | + | |
| 227 | + div class="mt-4 grid gap-4 md:grid-cols-[1fr_260px]" { | |
| 228 | + div class="min-w-0" { | |
| 229 | + (ui::alert_ok(saved)) | |
| 230 | + (ui::alert_error(error)) | |
| 231 | + @if settings_tab { | |
| 232 | + (settings(state, &owner, &package, &tags).await?) | |
| 233 | + } @else { | |
| 234 | + (tags_view(&package, &tags, &pull, &image_ref, is_admin, untagged, ctx.viewer.as_ref())) | |
| 235 | + } | |
| 236 | + } | |
| 237 | + aside class="space-y-3 text-[13px]" { | |
| 238 | + div class="box" { | |
| 239 | + div class="box-head font-semibold" { "About" } | |
| 240 | + dl class="grid grid-cols-[auto_1fr] gap-x-3 gap-y-1 px-3 py-2" { | |
| 241 | + dt class="text-ink-faint" { "Owner" } | |
| 242 | + dd class="flex min-w-0 items-center gap-1.5" { | |
| 243 | + (ui::avatar(&owner.name, owner.avatar_key.as_deref(), 16)) | |
| 244 | + a href={ "/" (owner.name) } class="truncate" { (owner.name) } | |
| 245 | + } | |
| 246 | + dt class="text-ink-faint" { "Visibility" } | |
| 247 | + dd { (if package.is_public() { "Public" } else { "Private" }) } | |
| 248 | + dt class="text-ink-faint" { "Repository" } | |
| 249 | + dd class="min-w-0 truncate" { | |
| 250 | + @if let Some(repo) = &linked_repo { | |
| 251 | + a href=(repo.url()) { (repo.full_name()) } | |
| 252 | + } @else { | |
| 253 | + span class="text-ink-faint" { "Not linked" } | |
| 254 | + } | |
| 255 | + } | |
| 256 | + dt class="text-ink-faint" { "Pulls" } | |
| 257 | + dd { (package.pull_count) } | |
| 258 | + dt class="text-ink-faint" { "Created" } | |
| 259 | + dd { (ui::time(package.created_at)) } | |
| 260 | + dt class="text-ink-faint" { "Updated" } | |
| 261 | + dd { (ui::time(package.updated_at)) } | |
| 262 | + } | |
| 263 | + } | |
| 264 | + p class="px-1 text-xs text-ink-faint" { | |
| 265 | + "Image visibility is independent of the repository's. " | |
| 266 | + a href="/docs/registry" { "Registry docs" } | |
| 267 | + } | |
| 268 | + } | |
| 269 | + } | |
| 270 | + } | |
| 271 | + }; | |
| 272 | + | |
| 273 | + let mut page = Page::new(package.full_name(), body) | |
| 274 | + .description(if package.description.is_empty() { format!("Container image {image_ref}") } else { package.description.clone() }); | |
| 275 | + if !package.is_public() || settings_tab { | |
| 276 | + page = page.noindex(); | |
| 277 | + } | |
| 278 | + Ok(ctx.render(page)) | |
| 279 | +} | |
| 280 | + | |
| 281 | +/// "linux/amd64, linux/arm64" for an index; nothing for a single image. | |
| 282 | +fn platforms(tag: &TagRow) -> Vec<String> { | |
| 283 | + if tag.media_type != OCI_INDEX && tag.media_type != DOCKER_LIST { | |
| 284 | + return Vec::new(); | |
| 285 | + } | |
| 286 | + let Ok(json) = serde_json::from_slice::<Value>(&tag.content) else { return Vec::new() }; | |
| 287 | + json["manifests"] | |
| 288 | + .as_array() | |
| 289 | + .into_iter() | |
| 290 | + .flatten() | |
| 291 | + .filter_map(|m| { | |
| 292 | + let os = m["platform"]["os"].as_str()?; | |
| 293 | + let arch = m["platform"]["architecture"].as_str()?; | |
| 294 | + if os == "unknown" { | |
| 295 | + return None; // attestation manifests | |
| 296 | + } | |
| 297 | + Some(match m["platform"]["variant"].as_str() { | |
| 298 | + Some(variant) => format!("{os}/{arch}/{variant}"), | |
| 299 | + None => format!("{os}/{arch}"), | |
| 300 | + }) | |
| 301 | + }) | |
| 302 | + .collect() | |
| 303 | +} | |
| 304 | + | |
| 305 | +fn short_digest(digest: &str) -> &str { | |
| 306 | + let hex = digest.trim_start_matches("sha256:"); | |
| 307 | + &hex[..hex.len().min(12)] | |
| 308 | +} | |
| 309 | + | |
| 310 | +fn tags_view(package: &Package, tags: &[TagRow], pull: &str, image_ref: &str, is_admin: bool, untagged: i64, viewer: Option<&Viewer>) -> Markup { | |
| 311 | + html! { | |
| 312 | + div class="box mb-4" { | |
| 313 | + div class="box-head font-semibold" { "Pull" } | |
| 314 | + div class="p-3" { | |
| 315 | + (ui::copy_field("pull-command", pull, "package_pull_copied")) | |
| 316 | + @if !package.is_public() { | |
| 317 | + p class="hint" { | |
| 318 | + "Private image: run " code { "docker login " (image_ref.split('/').next().unwrap_or("")) } " with your username and a personal access token first." | |
| 319 | + } | |
| 320 | + } | |
| 321 | + } | |
| 322 | + } | |
| 323 | + @if tags.is_empty() { | |
| 324 | + (ui::empty_state("No tags", html! { | |
| 325 | + p { "Push a tag to this image:" } | |
| 326 | + pre class="mx-auto mt-2 max-w-[560px] overflow-x-auto rounded-[4px] border border-edge bg-surface-sunken p-2 text-left font-mono text-xs" { | |
| 327 | + "docker tag <local-image> " (image_ref) ":latest\n" | |
| 328 | + "docker push " (image_ref) ":latest" | |
| 329 | + } | |
| 330 | + })) | |
| 331 | + } @else { | |
| 332 | + div class="box overflow-x-auto" { | |
| 333 | + table class="w-full text-[13px]" { | |
| 334 | + thead { | |
| 335 | + tr class="border-b border-edge bg-surface-raised text-left text-xs text-ink-dim" { | |
| 336 | + th class="px-3 py-1.5 font-semibold" { "Tag" } | |
| 337 | + th class="px-3 py-1.5 font-semibold" { "Digest" } | |
| 338 | + th class="hidden px-3 py-1.5 font-semibold sm:table-cell" { "Platforms" } | |
| 339 | + th class="px-3 py-1.5 text-right font-semibold" { "Size" } | |
| 340 | + th class="hidden px-3 py-1.5 font-semibold sm:table-cell" { "Updated" } | |
| 341 | + @if is_admin { th class="px-3 py-1.5" {} } | |
| 342 | + } | |
| 343 | + } | |
| 344 | + tbody class="divide-y divide-edge" { | |
| 345 | + @for tag in tags { | |
| 346 | + @let platforms = platforms(tag); | |
| 347 | + tr class="hover:bg-surface-raised" { | |
| 348 | + td class="px-3 py-1.5 font-mono font-semibold" { | |
| 349 | + button type="button" class="cursor-pointer text-left text-ink hover:text-accent" title="Copy pull command" | |
| 350 | + data-copy={ "docker pull " (image_ref) ":" (tag.name) } data-track="package_tag_pull_copied" { (tag.name) } | |
| 351 | + } | |
| 352 | + td class="px-3 py-1.5 font-mono text-xs text-ink-dim" title=(tag.digest) { | |
| 353 | + button type="button" class="cursor-pointer hover:text-ink" data-copy={ (image_ref) "@" (tag.digest) } title="Copy pinned reference" { (short_digest(&tag.digest)) } | |
| 354 | + } | |
| 355 | + td class="hidden px-3 py-1.5 text-xs text-ink-dim sm:table-cell" { | |
| 356 | + @if platforms.is_empty() { "single" } @else { (platforms.join(", ")) } | |
| 357 | + } | |
| 358 | + td class="px-3 py-1.5 text-right text-xs whitespace-nowrap" { (ui::bytes(tag.total_size.max(0) as u64)) } | |
| 359 | + td class="hidden px-3 py-1.5 text-xs whitespace-nowrap text-ink-dim sm:table-cell" { (ui::time(tag.updated_at)) } | |
| 360 | + @if is_admin { | |
| 361 | + td class="px-3 py-1.5 text-right" { | |
| 362 | + form method="post" action=(package.url()) data-track-submit="package_tag_deleted" { | |
| 363 | + input type="hidden" name="action" value="delete_tag"; | |
| 364 | + input type="hidden" name="tag" value=(tag.name); | |
| 365 | + button type="submit" class="btn btn-sm btn-danger" title={ "Delete tag " (tag.name) } { "Delete" } | |
| 366 | + } | |
| 367 | + } | |
| 368 | + } | |
| 369 | + } | |
| 370 | + } | |
| 371 | + } | |
| 372 | + } | |
| 373 | + } | |
| 374 | + @if untagged > 0 { | |
| 375 | + p class="hint px-1" { (ui::plural(untagged, "untagged version", "untagged versions")) " kept (pullable by digest)." } | |
| 376 | + } | |
| 377 | + } | |
| 378 | + @if viewer.is_none() && !package.is_public() { | |
| 379 | + p class="hint" { "Sign in to manage this image." } | |
| 380 | + } | |
| 381 | + } | |
| 382 | +} | |
| 383 | + | |
| 384 | +async fn settings(state: &AppState, owner: &Account, package: &Package, tags: &[TagRow]) -> AppResult<Markup> { | |
| 385 | + let repos: Vec<(i64, String)> = sqlx::query_as("select id, name::text from repos where owner_id = $1 order by name limit 1000") | |
| 386 | + .bind(owner.id) | |
| 387 | + .fetch_all(&state.db) | |
| 388 | + .await?; | |
| 389 | + let url = package.url(); | |
| 390 | + Ok(html! { | |
| 391 | + div class="space-y-4" { | |
| 392 | + form method="post" action=(url) class="box" data-track-submit="package_visibility_submitted" { | |
| 393 | + div class="box-head font-semibold" { "Visibility" } | |
| 394 | + div class="space-y-2 p-3" { | |
| 395 | + input type="hidden" name="action" value="visibility"; | |
| 396 | + label class="flex items-start gap-2" { | |
| 397 | + input type="radio" name="visibility" value="public" checked[package.is_public()] class="mt-0.5"; | |
| 398 | + span { strong { "Public" } br; span class="text-ink-dim" { "Anyone can pull this image, signed in or not." } } | |
| 399 | + } | |
| 400 | + label class="flex items-start gap-2" { | |
| 401 | + input type="radio" name="visibility" value="private" checked[!package.is_public()] class="mt-0.5"; | |
| 402 | + span { strong { "Private" } br; span class="text-ink-dim" { "Only " (owner.name) (if owner.is_org() { " members" } else { "" }) " and collaborators of the linked repository can pull." } } | |
| 403 | + } | |
| 404 | + p class="hint" { "This setting is separate from the visibility of any repository." } | |
| 405 | + button type="submit" class="btn btn-primary" { "Save visibility" } | |
| 406 | + } | |
| 407 | + } | |
| 408 | + | |
| 409 | + form method="post" action=(url) class="box" { | |
| 410 | + div class="box-head font-semibold" { "Description" } | |
| 411 | + div class="space-y-2 p-3" { | |
| 412 | + input type="hidden" name="action" value="description"; | |
| 413 | + input class="input" name="description" maxlength="350" value=(package.description) placeholder="What this image is for"; | |
| 414 | + button type="submit" class="btn" { "Save description" } | |
| 415 | + } | |
| 416 | + } | |
| 417 | + | |
| 418 | + form method="post" action=(url) class="box" data-track-submit="package_repo_link_submitted" { | |
| 419 | + div class="box-head font-semibold" { "Linked repository" } | |
| 420 | + div class="space-y-2 p-3" { | |
| 421 | + input type="hidden" name="action" value="link_repo"; | |
| 422 | + select class="input" name="repo_id" { | |
| 423 | + option value="" selected[package.repo_id.is_none()] { "Not linked" } | |
| 424 | + @for (id, name) in &repos { | |
| 425 | + option value=(id) selected[package.repo_id == Some(*id)] { (owner.name) "/" (name) } | |
| 426 | + } | |
| 427 | + } | |
| 428 | + p class="hint" { "Collaborators of the linked repository get the same access to this image. Its visibility is not inherited." } | |
| 429 | + button type="submit" class="btn" { "Save link" } | |
| 430 | + } | |
| 431 | + } | |
| 432 | + | |
| 433 | + div class="box border-danger/50" { | |
| 434 | + div class="box-head font-semibold text-danger" { "Danger zone" } | |
| 435 | + form method="post" action=(url) class="space-y-2 p-3" data-track-submit="package_delete_submitted" { | |
| 436 | + input type="hidden" name="action" value="delete_package"; | |
| 437 | + p { "Deleting " strong { (package.full_name()) } " removes all " (tags.len()) " tags and every version. Layers are cleaned up from storage afterwards. This cannot be undone." } | |
| 438 | + label class="label" for="confirm" { "Type " code { (package.name) } " to confirm" } | |
| 439 | + input id="confirm" class="input font-mono" name="confirm" autocomplete="off" data-confirm-value=(package.name); | |
| 440 | + button type="submit" class="btn btn-danger" { "Delete this image" } | |
| 441 | + } | |
| 442 | + } | |
| 443 | + } | |
| 444 | + }) | |
| 445 | +} | |
| 446 | + | |
| 447 | +// --------------------------------------------------------------------------- | |
| 448 | +// Actions | |
| 449 | + | |
| 450 | +#[derive(Deserialize)] | |
| 451 | +struct ActionForm { | |
| 452 | + action: String, | |
| 453 | + visibility: Option<String>, | |
| 454 | + description: Option<String>, | |
| 455 | + repo_id: Option<String>, | |
| 456 | + tag: Option<String>, | |
| 457 | + confirm: Option<String>, | |
| 458 | +} | |
| 459 | + | |
| 460 | +/// Rejects cross-site form posts (SameSite cookies already block most). | |
| 461 | +fn same_origin(state: &AppState, headers: &HeaderMap) -> bool { | |
| 462 | + match headers.get(header::ORIGIN).and_then(|v| v.to_str().ok()) { | |
| 463 | + Some(origin) => origin.trim_end_matches('/') == state.config.base_url(), | |
| 464 | + None => true, | |
| 465 | + } | |
| 466 | +} | |
| 467 | + | |
| 468 | +async fn act(ctx: Ctx, Path((owner, name)): Path<(String, String)>, headers: HeaderMap, Form(form): Form<ActionForm>) -> AppResult<Response> { | |
| 469 | + if !same_origin(&ctx.state, &headers) { | |
| 470 | + return Err(AppError::forbidden("Cross-site request refused.")); | |
| 471 | + } | |
| 472 | + let Some(viewer) = ctx.viewer.clone() else { return Err(AppError::Unauthorized) }; | |
| 473 | + let loaded = load(&ctx, &owner, &name).await?; | |
| 474 | + if !loaded.access.is_admin() { | |
| 475 | + return Err(AppError::forbidden("Only image admins can change it.")); | |
| 476 | + } | |
| 477 | + let state = &ctx.state; | |
| 478 | + let package = loaded.package.clone(); | |
| 479 | + let url = package.url(); | |
| 480 | + let ip = auth::client_ip(&headers); | |
| 481 | + | |
| 482 | + let outcome: Result<Response, String> = match form.action.as_str() { | |
| 483 | + "visibility" => { | |
| 484 | + let visibility = form.visibility.as_deref().unwrap_or(""); | |
| 485 | + if !matches!(visibility, "public" | "private") { | |
| 486 | + Err("Choose public or private.".into()) | |
| 487 | + } else { | |
| 488 | + sqlx::query("update packages set visibility = $2, updated_at = now() where id = $1") | |
| 489 | + .bind(package.id) | |
| 490 | + .bind(visibility) | |
| 491 | + .execute(&state.db) | |
| 492 | + .await?; | |
| 493 | + audit(&state.db, Some(viewer.id), "package.visibility", &package.full_name(), json!({ "visibility": visibility }), ip.as_deref()).await; | |
| 494 | + analytics::track(state, "package_visibility_changed", Some(&viewer.name), &url, json!({ "visibility": visibility })); | |
| 495 | + tracing::info!(package = %package.full_name(), visibility, user = %viewer.name, "package visibility changed"); | |
| 496 | + Ok(Redirect::to(&format!("{url}?tab=settings&saved=visibility")).into_response()) | |
| 497 | + } | |
| 498 | + } | |
| 499 | + "description" => { | |
| 500 | + let description: String = form.description.unwrap_or_default().trim().chars().take(350).collect(); | |
| 501 | + sqlx::query("update packages set description = $2, updated_at = now() where id = $1") | |
| 502 | + .bind(package.id) | |
| 503 | + .bind(&description) | |
| 504 | + .execute(&state.db) | |
| 505 | + .await?; | |
| 506 | + analytics::track(state, "package_description_changed", Some(&viewer.name), &url, json!({})); | |
| 507 | + Ok(Redirect::to(&format!("{url}?tab=settings&saved=description")).into_response()) | |
| 508 | + } | |
| 509 | + "link_repo" => { | |
| 510 | + let repo_id = form.repo_id.as_deref().filter(|v| !v.is_empty()).map(|v| v.parse::<i64>()); | |
| 511 | + match repo_id { | |
| 512 | + Some(Err(_)) => Err("Pick a repository.".into()), | |
| 513 | + Some(Ok(id)) => { | |
| 514 | + let repo = Repo::by_id(&state.db, id).await?; | |
| 515 | + match repo { | |
| 516 | + Some(repo) if repo.owner_id == package.owner_id => { | |
| 517 | + sqlx::query("update packages set repo_id = $2, updated_at = now() where id = $1") | |
| 518 | + .bind(package.id) | |
| 519 | + .bind(repo.id) | |
| 520 | + .execute(&state.db) | |
| 521 | + .await?; | |
| 522 | + audit(&state.db, Some(viewer.id), "package.link_repo", &package.full_name(), json!({ "repo": repo.full_name() }), ip.as_deref()).await; | |
| 523 | + analytics::track(state, "package_repo_linked", Some(&viewer.name), &url, json!({})); | |
| 524 | + Ok(Redirect::to(&format!("{url}?tab=settings&saved=repo")).into_response()) | |
| 525 | + } | |
| 526 | + _ => Err("Images can only link to repositories with the same owner.".into()), | |
| 527 | + } | |
| 528 | + } | |
| 529 | + None => { | |
| 530 | + sqlx::query("update packages set repo_id = null, updated_at = now() where id = $1").bind(package.id).execute(&state.db).await?; | |
| 531 | + audit(&state.db, Some(viewer.id), "package.unlink_repo", &package.full_name(), json!({}), ip.as_deref()).await; | |
| 532 | + Ok(Redirect::to(&format!("{url}?tab=settings&saved=repo")).into_response()) | |
| 533 | + } | |
| 534 | + } | |
| 535 | + } | |
| 536 | + "delete_tag" => { | |
| 537 | + let tag = form.tag.unwrap_or_default(); | |
| 538 | + let removed = sqlx::query("delete from tags where package_id = $1 and name = $2") | |
| 539 | + .bind(package.id) | |
| 540 | + .bind(&tag) | |
| 541 | + .execute(&state.db) | |
| 542 | + .await? | |
| 543 | + .rows_affected(); | |
| 544 | + if removed == 0 { | |
| 545 | + Err(format!("Tag {tag} does not exist.")) | |
| 546 | + } else { | |
| 547 | + audit(&state.db, Some(viewer.id), "package.tag_delete", &package.full_name(), json!({ "tag": tag }), ip.as_deref()).await; | |
| 548 | + analytics::track(state, "package_tag_deleted", Some(&viewer.name), &url, json!({})); | |
| 549 | + tracing::info!(package = %package.full_name(), tag, user = %viewer.name, "tag deleted"); | |
| 550 | + Ok(Redirect::to(&format!("{url}?saved=tag")).into_response()) | |
| 551 | + } | |
| 552 | + } | |
| 553 | + "delete_package" => { | |
| 554 | + if form.confirm.as_deref().map(str::trim) != Some(package.name.as_str()) { | |
| 555 | + Err("Type the image name exactly to delete it.".into()) | |
| 556 | + } else { | |
| 557 | + sqlx::query("delete from packages where id = $1").bind(package.id).execute(&state.db).await?; | |
| 558 | + audit(&state.db, Some(viewer.id), "package.delete", &package.full_name(), json!({}), ip.as_deref()).await; | |
| 559 | + analytics::track(state, "package_deleted", Some(&viewer.name), &url, json!({})); | |
| 560 | + tracing::info!(package = %package.full_name(), user = %viewer.name, "package deleted"); | |
| 561 | + Ok(Redirect::to(&format!("/{}/-/packages", loaded.owner.name)).into_response()) | |
| 562 | + } | |
| 563 | + } | |
| 564 | + other => Err(format!("Unknown action {other}.")), | |
| 565 | + }; | |
| 566 | + | |
| 567 | + match outcome { | |
| 568 | + Ok(response) => Ok(response), | |
| 569 | + Err(message) => { | |
| 570 | + let tab = if form.action == "delete_tag" { None } else { Some("settings") }; | |
| 571 | + let loaded = load(&ctx, &owner, &name).await?; | |
| 572 | + let mut response = render_detail(&ctx, loaded, tab, None, Some(&message)).await?; | |
| 573 | + *response.status_mut() = StatusCode::BAD_REQUEST; | |
| 574 | + Ok(response) | |
| 575 | + } | |
| 576 | + } | |
| 577 | +} |
+239-0backend/src/registry/token.rs
| @@ -0,0 +1,239 @@ | ||
| 1 | +//! /v2/token: trades a personal access token (Basic auth, any username) or | |
| 2 | +//! nothing at all for a short-lived bearer token scoped to the images and | |
| 3 | +//! actions the caller may use. Signed with SECRET_KEY via `crate::signed`. | |
| 4 | + | |
| 5 | +use axum::{ | |
| 6 | + Json, | |
| 7 | + body::Bytes, | |
| 8 | + extract::{RawQuery, State}, | |
| 9 | + http::{HeaderMap, Method, StatusCode}, | |
| 10 | + response::{IntoResponse, Response}, | |
| 11 | +}; | |
| 12 | +use serde::{Deserialize, Serialize}; | |
| 13 | +use serde_json::json; | |
| 14 | + | |
| 15 | +use super::{Actions, RegError, actions_for, error_response, split_name, with_api_version}; | |
| 16 | +use crate::{ | |
| 17 | + analytics, | |
| 18 | + auth::{self, Viewer}, | |
| 19 | + models::{Account, Package}, | |
| 20 | + signed, | |
| 21 | + state::AppState, | |
| 22 | +}; | |
| 23 | + | |
| 24 | +pub const SERVICE: &str = "irongit"; | |
| 25 | +const PURPOSE: &str = "registry"; | |
| 26 | +/// Docker refreshes expired tokens before each request, so short is fine. | |
| 27 | +const TTL_SECS: i64 = 300; | |
| 28 | + | |
| 29 | +#[derive(Debug, Clone, Serialize, Deserialize)] | |
| 30 | +pub struct Grant { | |
| 31 | + pub name: String, | |
| 32 | + pub actions: Vec<String>, | |
| 33 | +} | |
| 34 | + | |
| 35 | +#[derive(Debug, Clone, Serialize, Deserialize)] | |
| 36 | +pub struct Claims { | |
| 37 | + pub uid: Option<i64>, | |
| 38 | + pub user: Option<String>, | |
| 39 | + pub access: Vec<Grant>, | |
| 40 | +} | |
| 41 | + | |
| 42 | +impl Claims { | |
| 43 | + pub fn actions_for(&self, name: &str) -> Actions { | |
| 44 | + let mut actions = Actions::default(); | |
| 45 | + for grant in self.access.iter().filter(|g| g.name == name) { | |
| 46 | + for action in &grant.actions { | |
| 47 | + match action.as_str() { | |
| 48 | + "pull" => actions.pull = true, | |
| 49 | + "push" => actions.push = true, | |
| 50 | + "delete" | "*" => actions.delete = true, | |
| 51 | + _ => {} | |
| 52 | + } | |
| 53 | + } | |
| 54 | + } | |
| 55 | + actions | |
| 56 | + } | |
| 57 | +} | |
| 58 | + | |
| 59 | +pub fn verify(state: &AppState, token: &str) -> Option<Claims> { | |
| 60 | + signed::verify(&state.config.secret_key, PURPOSE, token) | |
| 61 | +} | |
| 62 | + | |
| 63 | +/// A requested scope: `repository:<name>:<action>[,<action>]`. | |
| 64 | +#[derive(Debug, PartialEq, Eq)] | |
| 65 | +pub struct ScopeRequest { | |
| 66 | + pub name: String, | |
| 67 | + pub actions: Vec<String>, | |
| 68 | +} | |
| 69 | + | |
| 70 | +/// Parses every `repository:` scope in a (possibly space-separated) list. | |
| 71 | +/// The name sits between the first and last colon; other types are ignored. | |
| 72 | +pub fn parse_scopes<'a>(values: impl Iterator<Item = &'a str>) -> Vec<ScopeRequest> { | |
| 73 | + values | |
| 74 | + .flat_map(|v| v.split(' ')) | |
| 75 | + .filter_map(|scope| { | |
| 76 | + let (kind, rest) = scope.split_once(':')?; | |
| 77 | + if kind != "repository" { | |
| 78 | + return None; | |
| 79 | + } | |
| 80 | + let (name, actions) = rest.rsplit_once(':')?; | |
| 81 | + Some(ScopeRequest { | |
| 82 | + name: name.to_string(), | |
| 83 | + actions: actions.split(',').map(str::trim).filter(|a| !a.is_empty()).map(String::from).collect(), | |
| 84 | + }) | |
| 85 | + }) | |
| 86 | + .collect() | |
| 87 | +} | |
| 88 | + | |
| 89 | +#[derive(Deserialize)] | |
| 90 | +struct PostForm { | |
| 91 | + username: Option<String>, | |
| 92 | + password: Option<String>, | |
| 93 | + scope: Option<String>, | |
| 94 | +} | |
| 95 | + | |
| 96 | +/// GET (Docker's flow) and POST (the OAuth2 password form some clients use). | |
| 97 | +pub async fn issue(State(state): State<AppState>, method: Method, headers: HeaderMap, RawQuery(query): RawQuery, body: Bytes) -> Response { | |
| 98 | + let query = query.unwrap_or_default(); | |
| 99 | + let mut scopes: Vec<String> = url::form_urlencoded::parse(query.as_bytes()) | |
| 100 | + .filter(|(k, _)| k == "scope") | |
| 101 | + .map(|(_, v)| v.into_owned()) | |
| 102 | + .collect(); | |
| 103 | + | |
| 104 | + // Who is asking: Basic/Bearer PAT, or the POST form's password. | |
| 105 | + let mut secret: Option<String> = None; | |
| 106 | + let mut presented = false; | |
| 107 | + if method == Method::POST { | |
| 108 | + if let Ok(form) = serde_urlencoded_form(&body) { | |
| 109 | + if let Some(password) = form.password { | |
| 110 | + presented = true; | |
| 111 | + secret = Some(password); | |
| 112 | + let _ = form.username; | |
| 113 | + } | |
| 114 | + if let Some(scope) = form.scope { | |
| 115 | + scopes.push(scope); | |
| 116 | + } | |
| 117 | + } | |
| 118 | + } | |
| 119 | + if let Some(value) = headers.get(axum::http::header::AUTHORIZATION).and_then(|v| v.to_str().ok()) { | |
| 120 | + presented = true; | |
| 121 | + let (scheme, rest) = value.split_once(' ').unwrap_or((value, "")); | |
| 122 | + secret = match scheme.to_ascii_lowercase().as_str() { | |
| 123 | + "basic" => auth::basic_credentials(rest.trim()).map(|(_, password)| password), | |
| 124 | + "bearer" | "token" => Some(rest.trim().to_string()), | |
| 125 | + _ => None, | |
| 126 | + }; | |
| 127 | + } | |
| 128 | + | |
| 129 | + let viewer: Option<Viewer> = match (&secret, presented) { | |
| 130 | + (Some(secret), _) => match auth::viewer_from_token(&state.db, secret).await { | |
| 131 | + Ok(Some(viewer)) => Some(viewer), | |
| 132 | + Ok(None) => { | |
| 133 | + tracing::info!("registry token refused: invalid credentials"); | |
| 134 | + return error_response( | |
| 135 | + &state, | |
| 136 | + RegError::new(StatusCode::UNAUTHORIZED, "UNAUTHORIZED", "invalid credentials: log in with your username and a personal access token"), | |
| 137 | + ); | |
| 138 | + } | |
| 139 | + Err(error) => return error_response(&state, RegError::internal(error)), | |
| 140 | + }, | |
| 141 | + (None, true) => { | |
| 142 | + return error_response(&state, RegError::new(StatusCode::UNAUTHORIZED, "UNAUTHORIZED", "unsupported credentials")); | |
| 143 | + } | |
| 144 | + (None, false) => None, | |
| 145 | + }; | |
| 146 | + | |
| 147 | + let requests = parse_scopes(scopes.iter().map(String::as_str)); | |
| 148 | + let mut access = Vec::new(); | |
| 149 | + for request in requests { | |
| 150 | + match grant(&state, viewer.as_ref(), &request).await { | |
| 151 | + Ok(actions) if !actions.is_empty() => access.push(Grant { name: request.name, actions }), | |
| 152 | + Ok(_) => {} | |
| 153 | + Err(error) => return error_response(&state, error), | |
| 154 | + } | |
| 155 | + } | |
| 156 | + | |
| 157 | + if access.is_empty() && viewer.is_some() && scopes.is_empty() { | |
| 158 | + analytics::track(&state, "registry_login", viewer.as_ref().map(|v| v.name.as_str()), "/v2/token", json!({})); | |
| 159 | + } | |
| 160 | + tracing::debug!( | |
| 161 | + user = viewer.as_ref().map(|v| v.name.as_str()).unwrap_or("-"), | |
| 162 | + grants = %access.iter().map(|g| format!("{}:{}", g.name, g.actions.join(","))).collect::<Vec<_>>().join(" "), | |
| 163 | + "registry token issued" | |
| 164 | + ); | |
| 165 | + | |
| 166 | + let claims = Claims { uid: viewer.as_ref().map(|v| v.id), user: viewer.as_ref().map(|v| v.name.clone()), access }; | |
| 167 | + let token = signed::sign(&state.config.secret_key, PURPOSE, &claims, TTL_SECS); | |
| 168 | + with_api_version( | |
| 169 | + Json(json!({ | |
| 170 | + "token": token, | |
| 171 | + "access_token": token, | |
| 172 | + "expires_in": TTL_SECS, | |
| 173 | + "issued_at": chrono::Utc::now().to_rfc3339_opts(chrono::SecondsFormat::Secs, true), | |
| 174 | + })) | |
| 175 | + .into_response(), | |
| 176 | + ) | |
| 177 | +} | |
| 178 | + | |
| 179 | +fn serde_urlencoded_form(body: &[u8]) -> Result<PostForm, ()> { | |
| 180 | + let mut form = PostForm { username: None, password: None, scope: None }; | |
| 181 | + for (key, value) in url::form_urlencoded::parse(body) { | |
| 182 | + match key.as_ref() { | |
| 183 | + "username" => form.username = Some(value.into_owned()), | |
| 184 | + "password" => form.password = Some(value.into_owned()), | |
| 185 | + "scope" => form.scope = Some(value.into_owned()), | |
| 186 | + _ => {} | |
| 187 | + } | |
| 188 | + } | |
| 189 | + Ok(form) | |
| 190 | +} | |
| 191 | + | |
| 192 | +/// The requested actions the caller actually holds on one image. | |
| 193 | +async fn grant(state: &AppState, viewer: Option<&Viewer>, request: &ScopeRequest) -> Result<Vec<String>, RegError> { | |
| 194 | + let Ok((owner_name, image)) = split_name(&request.name) else { return Ok(Vec::new()) }; | |
| 195 | + let Some(owner) = Account::by_name(&state.db, &owner_name).await? else { return Ok(Vec::new()) }; | |
| 196 | + let package = Package::by_path(&state.db, &owner_name, &image).await?; | |
| 197 | + let held = actions_for(state, viewer, &owner, package.as_ref()).await?; | |
| 198 | + let mut out: Vec<String> = request | |
| 199 | + .actions | |
| 200 | + .iter() | |
| 201 | + .filter(|a| held.allows(a)) | |
| 202 | + .map(|a| if a == "*" { "delete".to_string() } else { a.clone() }) | |
| 203 | + .collect(); | |
| 204 | + out.sort(); | |
| 205 | + out.dedup(); | |
| 206 | + Ok(out) | |
| 207 | +} | |
| 208 | + | |
| 209 | +#[cfg(test)] | |
| 210 | +mod tests { | |
| 211 | + use super::*; | |
| 212 | + | |
| 213 | + #[test] | |
| 214 | + fn parses_scopes() { | |
| 215 | + let scopes = parse_scopes(["repository:alice/api:pull,push", "repository:acme/tools/x:pull registry:catalog:*"].into_iter()); | |
| 216 | + assert_eq!( | |
| 217 | + scopes, | |
| 218 | + vec![ | |
| 219 | + ScopeRequest { name: "alice/api".into(), actions: vec!["pull".into(), "push".into()] }, | |
| 220 | + ScopeRequest { name: "acme/tools/x".into(), actions: vec!["pull".into()] }, | |
| 221 | + ] | |
| 222 | + ); | |
| 223 | + } | |
| 224 | + | |
| 225 | + #[test] | |
| 226 | + fn claims_resolve_actions_per_name() { | |
| 227 | + let claims = Claims { | |
| 228 | + uid: Some(1), | |
| 229 | + user: Some("alice".into()), | |
| 230 | + access: vec![ | |
| 231 | + Grant { name: "alice/api".into(), actions: vec!["pull".into(), "push".into()] }, | |
| 232 | + Grant { name: "bob/x".into(), actions: vec!["pull".into()] }, | |
| 233 | + ], | |
| 234 | + }; | |
| 235 | + assert_eq!(claims.actions_for("alice/api"), Actions { pull: true, push: true, delete: false }); | |
| 236 | + assert_eq!(claims.actions_for("bob/x"), Actions { pull: true, push: false, delete: false }); | |
| 237 | + assert_eq!(claims.actions_for("carol/y"), Actions::default()); | |
| 238 | + } | |
| 239 | +} |
+71-0frontend/src/pages/docs/registry.astro
| @@ -0,0 +1,71 @@ | ||
| 1 | +--- | |
| 2 | +import Layout from "../../layouts/Layout.astro"; | |
| 3 | +--- | |
| 4 | + | |
| 5 | +<Layout title="Container registry · irongit docs" description="Push and pull Docker and OCI images with irongit. Image visibility is set separately from code."> | |
| 6 | + <article class="markdown max-w-[760px]"> | |
| 7 | + <p class="font-mono text-xs text-ember">Docs / Registry</p> | |
| 8 | + <h1>Container registry</h1> | |
| 9 | + <p> | |
| 10 | + Every account and organization has a Docker registry namespace. Images are named | |
| 11 | + <code><host>/<owner>/<image></code>, and the image part may have slashes, like | |
| 12 | + <code><span class="host">this-site</span>/acme/tools/builder</code>. Layers are stored in object storage and pulls are | |
| 13 | + served from there directly. | |
| 14 | + </p> | |
| 15 | + | |
| 16 | + <h2>Sign in</h2> | |
| 17 | + <p>Docker signs in with your username and a personal access token (passwords are not accepted):</p> | |
| 18 | + <pre><code>docker login <span class="host">this-site</span> -u <username> | |
| 19 | +# Password: igp_... (a token with the packages scope)</code></pre> | |
| 20 | + <p> | |
| 21 | + Create a token under <a href="/settings/tokens">Settings / Tokens</a>, or let the CLI do it: | |
| 22 | + <code>ig login</code> sets up Docker for you. | |
| 23 | + </p> | |
| 24 | + | |
| 25 | + <h2>Push</h2> | |
| 26 | + <pre><code>docker tag my-app <span class="host">this-site</span>/<owner>/my-app:1.0 | |
| 27 | +docker push <span class="host">this-site</span>/<owner>/my-app:1.0</code></pre> | |
| 28 | + <p> | |
| 29 | + The first push creates the image. New images are <strong>private</strong>. You can push under your own name and under | |
| 30 | + any organization you belong to. | |
| 31 | + </p> | |
| 32 | + | |
| 33 | + <h2>Pull</h2> | |
| 34 | + <pre><code>docker pull <span class="host">this-site</span>/<owner>/my-app:1.0</code></pre> | |
| 35 | + <p>Public images can be pulled without signing in. Private images need <code>docker login</code> first.</p> | |
| 36 | + | |
| 37 | + <h2>Visibility</h2> | |
| 38 | + <p> | |
| 39 | + Each image is public or private on its own, independent of any repository. A private repository can publish a public | |
| 40 | + image, and the reverse. Change it on the image's Settings tab, or with | |
| 41 | + <code>ig image visibility <owner>/<image> public</code>. | |
| 42 | + </p> | |
| 43 | + | |
| 44 | + <h2>Linking an image to a repository</h2> | |
| 45 | + <p> | |
| 46 | + On the image's Settings tab you can link it to a repository with the same owner. Collaborators of that repository then | |
| 47 | + get the same access to the image (read collaborators can pull, write collaborators can push). The repository's | |
| 48 | + visibility is not inherited. | |
| 49 | + </p> | |
| 50 | + | |
| 51 | + <h2>Tags, versions and cleanup</h2> | |
| 52 | + <ul> | |
| 53 | + <li>Multi-platform images (manifest lists and OCI indexes) are supported; the tags table lists their platforms.</li> | |
| 54 | + <li>Deleting a tag keeps the version, which stays pullable by digest. Deleting the image removes everything.</li> | |
| 55 | + <li>Layers no image uses any more are removed from storage automatically after a day.</li> | |
| 56 | + </ul> | |
| 57 | + | |
| 58 | + <h2>Limits</h2> | |
| 59 | + <ul> | |
| 60 | + <li>Single layers are capped (10 GB by default). Manifests are capped at 4 MB.</li> | |
| 61 | + <li> | |
| 62 | + If the site sits behind Cloudflare's proxy, the registry hostname must be DNS-only, since the proxy rejects uploads | |
| 63 | + over 100 MB. | |
| 64 | + </li> | |
| 65 | + </ul> | |
| 66 | + </article> | |
| 67 | +</Layout> | |
| 68 | + | |
| 69 | +<script> | |
| 70 | + document.querySelectorAll(".host").forEach((el) => (el.textContent = location.host)); | |
| 71 | +</script> |