diff --git a/Cargo.toml b/Cargo.toml index 2cb726b..c7116b0 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -14,3 +14,5 @@ clap = { version = "4.5.17", features = ["derive", "env"] } serde_json = "1.0.128" base64 = "0.22.1" sha256 = "1.6.0" +uuid = { version = "1.0", features = ["v4"] } +bytes = "1.9.0" diff --git a/src/blobs.rs b/src/blobs.rs index 4740666..e4173c0 100644 --- a/src/blobs.rs +++ b/src/blobs.rs @@ -14,15 +14,16 @@ use serde_json::{json, Value}; use std::sync::Arc; use crate::{ - auth, state, + auth, response, state, storage::{self, write_blob}, }; use axum::{ body::Body, extract::{Path, Query, State}, - http::{HeaderMap, Request, StatusCode}, + http::{HeaderMap, StatusCode}, response::{Json, Response}, }; +use bytes::Bytes; // end-2 GET /v2/:name/blobs/:digest pub(crate) async fn get_blob_by_digest( @@ -144,74 +145,133 @@ pub(crate) async fn head_blob_by_digest( #[derive(Deserialize)] pub(crate) struct PostBlobUploadQueryParams { digest: Option, + #[allow(dead_code)] mount: Option, + #[allow(dead_code)] from: Option, } pub(crate) async fn post_blob_upload( + State(state): State>, Path((org, repo)): Path<(String, String)>, -) -> Response { - log::info!("blobs/post_blob_upload: org: {}, repo: {}", org, repo,); + Query(params): Query, + headers: HeaderMap, + body: Bytes, +) -> Response { + log::info!("blobs/post_blob_upload: org: {}, repo: {}", org, repo); - Response::builder() - .status(202) - .header("Location", format!("/v2/{}/{}/blobs/uploads/", org, repo)) - .body("202 Accepted".to_string()) - .expect("Failed to build response") -} + let host = &state.args.host; -pub(crate) async fn put_blob_upload( - Path((org, repo)): Path<(String, String)>, - query: Query, - body: Request, -) -> Response { - log::info!( - "blobs/put_blob_upload: org: {}, repo: {}, digest: {:#?}, mount: {:#?}, from: {:#?}", - org, - repo, - query.digest, - query.mount, - query.from - ); + if auth::get(State(state.clone()), headers.clone()) + .await + .status() + != StatusCode::OK + { + return Response::builder() + .status(StatusCode::UNAUTHORIZED) + .header( + "WWW-Authenticate", + format!("Basic realm=\"{}\", charset=\"UTF-8\"", host), + ) + .body(Body::from("401 Unauthorized")) + .unwrap(); + } + + // If digest is provided, handle monolithic upload (end-4b) + if let Some(digest_string) = params.digest { + let success = write_blob(&org, &repo, &digest_string, Body::from(body)).await; - let success = match &query.digest { - Some(digest) => write_blob(&org, &repo, digest, body.into_body()).await, - None => false, - }; + if !success { + return Response::builder() + .status(StatusCode::BAD_REQUEST) + .body(Body::from("Digest mismatch or write failed")) + .unwrap(); + } + + let clean_digest = digest_string + .strip_prefix("sha256:") + .unwrap_or(&digest_string); - if !success { return Response::builder() - .status(400) - .body("400 Bad Request".to_string()) - .expect("Failed to build response"); + .status(StatusCode::CREATED) + .header( + "Location", + format!( + "http://{}/v2/{}/{}/blobs/sha256:{}", + host, org, repo, clean_digest + ), + ) + .header("Docker-Content-Digest", format!("sha256:{}", clean_digest)) + .body(Body::empty()) + .unwrap(); } + // Create new upload session (end-4a) + let uuid = uuid::Uuid::new_v4().to_string(); + + if let Err(e) = storage::init_upload_session(&org, &repo, &uuid) { + log::error!("Failed to init upload session: {}", e); + return response::internal_error(); + } + + let location = format!("http://{}/v2/{}/{}/blobs/uploads/{}", host, org, repo, uuid); + Response::builder() - .status(201) - .header("Location", format!("/v2/{}/blobs/uploads/{}", org, repo)) - .header( - "Docker-Content-Digest", - query.digest.as_ref().expect("Digest should exist"), - ) - .body("201 Created".to_string()) - .expect("Failed to build response") + .status(StatusCode::ACCEPTED) + .header("Location", location) + .header("Range", "0-0") + .header("Docker-Upload-UUID", uuid) + .body(Body::empty()) + .unwrap() } // end-5 PATCH /v2/:name/blobs/uploads/:reference pub(crate) async fn patch_blob_upload( - State(data): State>, - Path(name): Path, - Path(reference): Path, -) -> Json { - let status = data.server_status.lock().await; + State(state): State>, + Path((org, repo, uuid)): Path<(String, String, String)>, + headers: HeaderMap, + body: Bytes, +) -> Response { log::info!( - "blobs/patch_blob_upload: name: {}, reference: {}", - name, - reference + "blobs/patch_blob_upload: org: {}, repo: {}, uuid: {}", + org, + repo, + uuid ); - Json(json!({ - "not_implemented": format!("name {} reference {} server_status {}", name, reference, status) - })) + + let host = &state.args.host; + + if auth::get(State(state.clone()), headers).await.status() != StatusCode::OK { + return Response::builder() + .status(StatusCode::UNAUTHORIZED) + .header( + "WWW-Authenticate", + format!("Basic realm=\"{}\", charset=\"UTF-8\"", host), + ) + .body(Body::from("401 Unauthorized")) + .unwrap(); + } + + match storage::append_upload_chunk(&org, &repo, &uuid, &body) { + Ok(total_size) => { + let location = format!("http://{}/v2/{}/{}/blobs/uploads/{}", host, org, repo, uuid); + + Response::builder() + .status(StatusCode::ACCEPTED) + .header("Location", location) + .header("Range", format!("0-{}", total_size.saturating_sub(1))) + .header("Docker-Upload-UUID", &uuid) + .body(Body::empty()) + .unwrap() + } + Err(e) => { + log::error!("Failed to append chunk for upload {}: {}", uuid, e); + Response::builder() + .status(StatusCode::NOT_FOUND) + .body(Body::from("Upload session not found")) + .unwrap() + } + } } // end-6 PUT /v2/:name/blobs/uploads/:reference?digest=:digest @@ -219,22 +279,71 @@ pub(crate) async fn patch_blob_upload( pub(crate) struct End6QueryParams { digest: String, } + pub(crate) async fn put_blob_upload_by_reference( - State(data): State>, - Path(name): Path, - Path(reference): Path, - query: Query, -) -> Json { - let status = data.server_status.lock().await; + State(state): State>, + Path((org, repo, uuid)): Path<(String, String, String)>, + Query(params): Query, + headers: HeaderMap, + body: Bytes, +) -> Response { log::info!( - "blobs/put_blob_upload_by_reference: name: {}, reference: {}, digest: {}", - name, - reference, - query.digest + "blobs/put_blob_upload_by_reference: org: {}, repo: {}, uuid: {}, digest: {}", + org, + repo, + uuid, + params.digest ); - Json(json!({ - "not_implemented": format!("name {} reference {} digest {} server_status {}", name, reference, query.digest, status) - })) + + let host = &state.args.host; + + if auth::get(State(state.clone()), headers).await.status() != StatusCode::OK { + return Response::builder() + .status(StatusCode::UNAUTHORIZED) + .header( + "WWW-Authenticate", + format!("Basic realm=\"{}\", charset=\"UTF-8\"", host), + ) + .body(Body::from("401 Unauthorized")) + .unwrap(); + } + + // Append final chunk if body is not empty + if !body.is_empty() { + if let Err(e) = storage::append_upload_chunk(&org, &repo, &uuid, &body) { + log::error!("Failed to append final chunk: {}", e); + return response::internal_error(); + } + } + + // Finalize upload and validate digest + match storage::finalize_upload(&org, &repo, &uuid, ¶ms.digest) { + Ok(actual_digest) => { + let location = format!( + "http://{}/v2/{}/{}/blobs/sha256:{}", + host, org, repo, actual_digest + ); + + Response::builder() + .status(StatusCode::CREATED) + .header("Location", location) + .header("Docker-Content-Digest", format!("sha256:{}", actual_digest)) + .body(Body::empty()) + .unwrap() + } + Err(e) => { + log::error!("Failed to finalize upload: {}", e); + + // Clean up failed upload + let _ = storage::delete_upload_session(&org, &repo, &uuid); + + if e.contains("Digest mismatch") { + response::digest_mismatch() + } else { + response::internal_error() + } + } + } } // end-10 DELETE /v2/:name/blobs/:digest diff --git a/src/main.rs b/src/main.rs index 4ac6044..21b3d75 100644 --- a/src/main.rs +++ b/src/main.rs @@ -50,10 +50,6 @@ async fn main() { "/v2/{org}/{repo}/blobs/uploads/", post(blobs::post_blob_upload), ) // end-4a, end-4b, end-11 - .route( - "/v2/{org}/{repo}/blobs/uploads/", - put(blobs::put_blob_upload), - ) // returned location for actual file upload .route( "/v2/{org}/{repo}/blobs/uploads/{reference}", patch(blobs::patch_blob_upload), diff --git a/src/response.rs b/src/response.rs index 4b73215..8c58842 100644 --- a/src/response.rs +++ b/src/response.rs @@ -1,4 +1,4 @@ -use axum::http::Response; +use axum::{body::Body, http::Response}; pub(crate) fn unauthorized(host: &str) -> Response { Response::builder() @@ -17,3 +17,25 @@ pub(crate) fn ok() -> Response { .body("200 OK".to_string()) .unwrap() } + +#[allow(dead_code)] +pub(crate) fn bad_request(message: &str) -> Response { + Response::builder() + .status(400) + .body(Body::from(message.to_string())) + .unwrap() +} + +pub(crate) fn internal_error() -> Response { + Response::builder() + .status(500) + .body(Body::from("Internal server error")) + .unwrap() +} + +pub(crate) fn digest_mismatch() -> Response { + Response::builder() + .status(400) + .body(Body::from("Digest mismatch")) + .unwrap() +} diff --git a/src/storage.rs b/src/storage.rs index eff4f87..243786e 100644 --- a/src/storage.rs +++ b/src/storage.rs @@ -181,3 +181,97 @@ pub(crate) fn list_tags(org: &str, repo: &str) -> Result, std::io::E tags.sort(); Ok(tags) } + +pub(crate) fn init_upload_session(org: &str, repo: &str, uuid: &str) -> Result<(), std::io::Error> { + let sanitized_org = sanitize_string(org); + let sanitized_repo = sanitize_string(repo); + let sanitized_uuid = sanitize_string(uuid); + + let upload_dir = format!("./tmp/uploads/{}/{}", sanitized_org, sanitized_repo); + std::fs::create_dir_all(&upload_dir)?; + + let upload_path = format!("{}/{}", upload_dir, sanitized_uuid); + std::fs::File::create(upload_path)?; + Ok(()) +} + +pub(crate) fn append_upload_chunk( + org: &str, + repo: &str, + uuid: &str, + chunk_data: &[u8], +) -> Result { + use std::fs::OpenOptions; + + let sanitized_org = sanitize_string(org); + let sanitized_repo = sanitize_string(repo); + let sanitized_uuid = sanitize_string(uuid); + + let upload_path = format!( + "./tmp/uploads/{}/{}/{}", + sanitized_org, sanitized_repo, sanitized_uuid + ); + + let mut file = OpenOptions::new().append(true).open(&upload_path)?; + + file.write_all(chunk_data)?; + + let metadata = std::fs::metadata(&upload_path)?; + Ok(metadata.len()) +} + +pub(crate) fn finalize_upload( + org: &str, + repo: &str, + uuid: &str, + expected_digest: &str, +) -> Result { + let sanitized_org = sanitize_string(org); + let sanitized_repo = sanitize_string(repo); + let sanitized_uuid = sanitize_string(uuid); + + let upload_path = format!( + "./tmp/uploads/{}/{}/{}", + sanitized_org, sanitized_repo, sanitized_uuid + ); + + let upload_data = + std::fs::read(&upload_path).map_err(|e| format!("Failed to read upload: {}", e))?; + + let actual_digest = sha256::digest(&upload_data); + let clean_expected = expected_digest + .strip_prefix("sha256:") + .unwrap_or(expected_digest); + + if actual_digest != clean_expected { + return Err(format!( + "Digest mismatch: expected {}, got {}", + clean_expected, actual_digest + )); + } + + let blob_dir = format!("./tmp/blobs/{}/{}", sanitized_org, sanitized_repo); + std::fs::create_dir_all(&blob_dir).map_err(|e| format!("Failed to create blob dir: {}", e))?; + + let blob_path = format!("{}/{}", blob_dir, actual_digest); + std::fs::rename(&upload_path, &blob_path) + .map_err(|e| format!("Failed to move upload to blob: {}", e))?; + + Ok(actual_digest) +} + +pub(crate) fn delete_upload_session( + org: &str, + repo: &str, + uuid: &str, +) -> Result<(), std::io::Error> { + let sanitized_org = sanitize_string(org); + let sanitized_repo = sanitize_string(repo); + let sanitized_uuid = sanitize_string(uuid); + + let upload_path = format!( + "./tmp/uploads/{}/{}/{}", + sanitized_org, sanitized_repo, sanitized_uuid + ); + std::fs::remove_file(upload_path) +}