use crate::{ContentRef, NamespaceId, UploadId};
use serde::{Deserialize, Serialize};
use std::collections::BTreeMap;
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
#[cfg_attr(feature = "openapi", derive(utoipa::ToSchema))]
#[serde(deny_unknown_fields)]
pub struct DirectPutContentClaim {
pub size_bytes: u64,
pub sha256: String,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
#[cfg_attr(feature = "openapi", derive(utoipa::ToSchema))]
#[serde(deny_unknown_fields)]
pub struct DirectMultipartContentClaim {
pub size_bytes: u64,
pub crc64nvme: String,
}
#[derive(Debug, Clone, Copy, Default, PartialEq, Eq, Serialize, Deserialize)]
#[cfg_attr(feature = "openapi", derive(utoipa::ToSchema))]
#[serde(deny_unknown_fields)]
pub struct DirectMultipartUploadOptions {
#[serde(default, skip_serializing_if = "Option::is_none")]
pub part_size_bytes: Option<u64>,
}
#[derive(Debug, Clone, Copy, Default, PartialEq, Eq, Serialize, Deserialize)]
#[cfg_attr(feature = "openapi", derive(utoipa::ToSchema))]
#[serde(rename_all = "snake_case")]
pub enum UploadMode {
#[default]
ServiceProxied,
DirectPut,
DirectMultipart,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
#[cfg_attr(feature = "openapi", derive(utoipa::ToSchema))]
#[serde(tag = "mode", rename_all = "snake_case", deny_unknown_fields)]
pub enum BeginUploadRequest {
#[cfg_attr(feature = "openapi", schema(title = "BeginUploadServiceProxied"))]
ServiceProxied {},
#[cfg_attr(feature = "openapi", schema(title = "BeginUploadDirectPut"))]
DirectPut {
content: DirectPutContentClaim,
},
#[cfg_attr(feature = "openapi", schema(title = "BeginUploadDirectMultipart"))]
DirectMultipart {
#[serde(default, skip_serializing_if = "Option::is_none")]
multipart: Option<DirectMultipartUploadOptions>,
},
}
impl BeginUploadRequest {
pub fn mode(&self) -> UploadMode {
match self {
Self::ServiceProxied {} => UploadMode::ServiceProxied,
Self::DirectPut { .. } => UploadMode::DirectPut,
Self::DirectMultipart { .. } => UploadMode::DirectMultipart,
}
}
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
#[cfg_attr(feature = "openapi", derive(utoipa::ToSchema))]
#[serde(tag = "kind", rename_all = "snake_case")]
pub enum ObjectTransferAccess {
#[cfg_attr(
feature = "openapi",
schema(title = "ObjectTransferAccessPresignedUrl")
)]
PresignedUrl {
method: String,
url: String,
#[serde(default, skip_serializing_if = "BTreeMap::is_empty")]
headers: BTreeMap<String, String>,
expires_at_ms: u64,
},
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
#[cfg_attr(feature = "openapi", derive(utoipa::ToSchema))]
pub struct DirectPutUpload {
pub content_ref: ContentRef,
pub access: ObjectTransferAccess,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[cfg_attr(feature = "openapi", derive(utoipa::ToSchema))]
pub struct DirectMultipartUpload {
pub part_size_bytes: u64,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
#[cfg_attr(feature = "openapi", derive(utoipa::ToSchema))]
#[serde(deny_unknown_fields)]
pub struct UploadPartChecksumClaim {
pub part_number: u32,
pub crc64nvme: String,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
#[cfg_attr(feature = "openapi", derive(utoipa::ToSchema))]
#[serde(deny_unknown_fields)]
pub struct SignUploadPartsRequest {
pub parts: Vec<UploadPartChecksumClaim>,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
#[cfg_attr(feature = "openapi", derive(utoipa::ToSchema))]
pub struct SignedUploadPart {
pub part_number: u32,
pub access: ObjectTransferAccess,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
#[cfg_attr(feature = "openapi", derive(utoipa::ToSchema))]
pub struct SignUploadPartsResponse {
pub namespace_id: NamespaceId,
pub upload_id: UploadId,
pub parts: Vec<SignedUploadPart>,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
#[cfg_attr(feature = "openapi", derive(utoipa::ToSchema))]
#[serde(deny_unknown_fields)]
pub struct CompletedUploadPart {
pub part_number: u32,
pub etag: String,
pub crc64nvme: String,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
#[cfg_attr(feature = "openapi", derive(utoipa::ToSchema))]
pub struct ValidatedContentToken {
pub content_ref: ContentRef,
pub token: String,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
#[cfg_attr(feature = "openapi", derive(utoipa::ToSchema))]
pub struct BeginUploadResponse {
pub namespace_id: NamespaceId,
pub upload_id: UploadId,
pub mode: UploadMode,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub direct_put: Option<DirectPutUpload>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub direct_multipart: Option<DirectMultipartUpload>,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
#[cfg_attr(feature = "openapi", derive(utoipa::ToSchema))]
pub struct UploadContentResponse {
pub namespace_id: NamespaceId,
pub upload_id: UploadId,
pub content_ref: ContentRef,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
#[cfg_attr(feature = "openapi", derive(utoipa::ToSchema))]
#[serde(tag = "completion", rename_all = "snake_case", deny_unknown_fields)]
pub enum CompleteUploadRequest {
#[cfg_attr(feature = "openapi", schema(title = "CompleteUploadContentRef"))]
ContentRef {
content_ref: ContentRef,
},
#[cfg_attr(feature = "openapi", schema(title = "CompleteUploadMultipart"))]
Multipart {
multipart: DirectMultipartContentClaim,
parts: Vec<CompletedUploadPart>,
},
}
impl CompleteUploadRequest {
pub fn for_content_ref(content_ref: ContentRef) -> Self {
Self::ContentRef { content_ref }
}
pub fn for_multipart(
claim: DirectMultipartContentClaim,
parts: Vec<CompletedUploadPart>,
) -> Self {
Self::Multipart {
multipart: claim,
parts,
}
}
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
#[cfg_attr(feature = "openapi", derive(utoipa::ToSchema))]
pub struct CompleteUploadResponse {
pub namespace_id: NamespaceId,
pub upload_id: UploadId,
pub content_ref: ContentRef,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub validated_content_token: Option<String>,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
#[cfg_attr(feature = "openapi", derive(utoipa::ToSchema))]
#[serde(tag = "state", rename_all = "snake_case")]
pub enum UploadSessionStatus {
#[cfg_attr(feature = "openapi", schema(title = "UploadSessionStatusOpen"))]
Open {
expires_at_ms: u64,
},
#[cfg_attr(feature = "openapi", schema(title = "UploadSessionStatusCompleted"))]
Completed {
completed_at_ms: u64,
content_ref: ContentRef,
#[serde(default, skip_serializing_if = "Option::is_none")]
validated_content_token: Option<String>,
},
#[cfg_attr(feature = "openapi", schema(title = "UploadSessionStatusAborted"))]
Aborted {
aborted_at_ms: u64,
},
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
#[cfg_attr(feature = "openapi", derive(utoipa::ToSchema))]
pub struct UploadStatusResponse {
pub namespace_id: NamespaceId,
pub upload_id: UploadId,
pub status: UploadSessionStatus,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
#[cfg_attr(feature = "openapi", derive(utoipa::ToSchema))]
pub struct AbortUploadResponse {
pub namespace_id: NamespaceId,
pub upload_id: UploadId,
pub aborted_at_ms: u64,
}
#[cfg(test)]
mod tests {
use super::{
BeginUploadRequest, BeginUploadResponse, CompleteUploadRequest, DirectPutContentClaim,
DirectPutUpload, ObjectTransferAccess, UploadMode, UploadSessionStatus,
};
use crate::{ContentId, ContentRef, NamespaceId, UploadId};
use std::collections::BTreeMap;
#[test]
fn direct_put_upload_mode_serializes_as_expected() {
assert_eq!(
serde_json::to_string(&UploadMode::DirectPut).expect("serialize mode"),
r#""direct_put""#
);
}
#[test]
fn a_begin_request_without_a_mode_does_not_decode() {
assert!(serde_json::from_str::<BeginUploadRequest>("{}").is_err());
assert_eq!(
serde_json::from_str::<BeginUploadRequest>(r#"{"mode":"service_proxied"}"#)
.expect("decode proxied begin request"),
BeginUploadRequest::ServiceProxied {}
);
}
#[test]
fn a_begin_request_carrying_another_modes_fields_does_not_decode() {
for body in [
r#"{"mode":"service_proxied","multipart":{"part_size_bytes":8388608}}"#,
r#"{"mode":"service_proxied","content":{"size_bytes":5,"sha256":"2cf24dba5fb0a30e26e83b2ac5b9e29e1b161e5c1fa7425e73043362938b9824"}}"#,
r#"{"mode":"direct_multipart","content":{"size_bytes":5,"sha256":"2cf24dba5fb0a30e26e83b2ac5b9e29e1b161e5c1fa7425e73043362938b9824"}}"#,
r#"{"mode":"direct_put"}"#,
] {
assert!(
serde_json::from_str::<BeginUploadRequest>(body).is_err(),
"decoded a begin request that mixes modes: {body}"
);
}
}
#[test]
fn a_completion_mixing_its_two_shapes_does_not_decode() {
for body in [
r#"{"completion":"multipart","multipart":{"size_bytes":5,"crc64nvme":"0123456789abcdef"},"parts":[],"content_ref":{"kind":"blob_v1","content_id":"con_0123456789abcdef0123456789abcdef","size_bytes":5,"storage_checksum":{"algorithm":"sha256","value":"2cf24dba5fb0a30e26e83b2ac5b9e29e1b161e5c1fa7425e73043362938b9824"}}}"#,
r#"{"completion":"multipart","multipart":{"size_bytes":5,"crc64nvme":"0123456789abcdef"}}"#,
r#"{"completion":"multipart","parts":[]}"#,
r#"{"completion":"content_ref"}"#,
] {
assert!(
serde_json::from_str::<CompleteUploadRequest>(body).is_err(),
"decoded a completion that mixes shapes: {body}"
);
}
}
#[test]
fn direct_put_response_exposes_only_presigned_access() {
let response = BeginUploadResponse {
namespace_id: NamespaceId::parse("demo").expect("namespace id"),
upload_id: UploadId::parse("upl_00000000000000000000000000000001")
.expect("valid upload id"),
mode: UploadMode::DirectPut,
direct_put: Some(DirectPutUpload {
content_ref: ContentRef::blob_v1(ContentId::generate(), b"hello"),
access: ObjectTransferAccess::PresignedUrl {
method: "PUT".to_owned(),
url: "https://bucket.example/object?X-Amz-Signature=abc".to_owned(),
headers: BTreeMap::from([
("if-none-match".to_owned(), "*".to_owned()),
(
"x-provider-checksum".to_owned(),
"LPJNul+wow4m6DsqxbninhsWHlwfp0JecwQzYpOLmCQ=".to_owned(),
),
]),
expires_at_ms: 1,
},
}),
direct_multipart: None,
};
let json = serde_json::to_string(&response).expect("serialize response");
assert!(json.contains(r#""kind":"presigned_url""#));
assert!(!json.contains("object_key"));
}
#[test]
fn a_direct_put_claim_names_only_size_and_digest() {
let request: BeginUploadRequest = serde_json::from_str(
r#"{"mode":"direct_put","content":{"size_bytes":5,"sha256":"2cf24dba5fb0a30e26e83b2ac5b9e29e1b161e5c1fa7425e73043362938b9824"}}"#,
)
.expect("decode direct-put begin request");
assert_eq!(
request,
BeginUploadRequest::DirectPut {
content: DirectPutContentClaim {
size_bytes: 5,
sha256: "2cf24dba5fb0a30e26e83b2ac5b9e29e1b161e5c1fa7425e73043362938b9824"
.to_owned(),
},
}
);
assert!(
serde_json::from_str::<DirectPutContentClaim>(
r#"{"size_bytes":5,"sha256":"2cf24dba5fb0a30e26e83b2ac5b9e29e1b161e5c1fa7425e73043362938b9824","content_id":"con_0123456789abcdef0123456789abcdef"}"#
)
.is_err(),
"a client must not be able to name the content object"
);
}
#[test]
fn upload_status_names_its_state_on_the_wire() {
let open = serde_json::to_value(UploadSessionStatus::Open {
expires_at_ms: 1_000,
})
.expect("serialize open status");
assert_eq!(open["state"], "open");
let aborted = serde_json::to_value(UploadSessionStatus::Aborted {
aborted_at_ms: 2_000,
})
.expect("serialize aborted status");
assert_eq!(aborted["state"], "aborted");
let completed = serde_json::to_value(UploadSessionStatus::Completed {
completed_at_ms: 3_000,
content_ref: ContentRef::blob_v1(ContentId::generate(), b"hello"),
validated_content_token: None,
})
.expect("serialize completed status");
assert_eq!(completed["state"], "completed");
assert!(
completed.get("validated_content_token").is_none(),
"a session past its receipt window reports no token at all"
);
}
}