use crate::{Checksum, ChecksumAlgorithm, 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 UploadContentClaim {
pub size_bytes: u64,
pub checksum: Checksum,
}
#[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: UploadContentClaim,
},
#[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,
pub checksum_algorithm: ChecksumAlgorithm,
}
#[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 checksum: Checksum,
}
#[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 checksum: Checksum,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
#[cfg_attr(feature = "openapi", derive(utoipa::ToSchema))]
#[serde(deny_unknown_fields)]
pub struct ContentToken {
pub content_ref: ContentRef,
pub token: String,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
#[cfg_attr(feature = "openapi", derive(utoipa::ToSchema))]
#[serde(tag = "mode", rename_all = "snake_case")]
pub enum BeginUploadResponse {
#[cfg_attr(
feature = "openapi",
schema(title = "BeginUploadResponseServiceProxied")
)]
ServiceProxied {
namespace_id: NamespaceId,
upload_id: UploadId,
},
#[cfg_attr(feature = "openapi", schema(title = "BeginUploadResponseDirectPut"))]
DirectPut {
namespace_id: NamespaceId,
upload_id: UploadId,
direct_put: DirectPutUpload,
},
#[cfg_attr(
feature = "openapi",
schema(title = "BeginUploadResponseDirectMultipart")
)]
DirectMultipart {
namespace_id: NamespaceId,
upload_id: UploadId,
direct_multipart: DirectMultipartUpload,
},
}
impl BeginUploadResponse {
pub fn namespace_id(&self) -> &NamespaceId {
match self {
Self::ServiceProxied { namespace_id, .. }
| Self::DirectPut { namespace_id, .. }
| Self::DirectMultipart { namespace_id, .. } => namespace_id,
}
}
pub fn upload_id(&self) -> &UploadId {
match self {
Self::ServiceProxied { upload_id, .. }
| Self::DirectPut { upload_id, .. }
| Self::DirectMultipart { upload_id, .. } => upload_id,
}
}
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))]
pub struct UploadContentResponse {
pub namespace_id: NamespaceId,
pub upload_id: UploadId,
pub content_ref: ContentRef,
}
#[derive(Debug, Clone, Copy, Default, PartialEq, Eq, Serialize, Deserialize)]
#[cfg_attr(feature = "openapi", derive(utoipa::ToSchema))]
#[serde(deny_unknown_fields)]
pub struct CompleteKnownContentUploadRequest {}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
#[cfg_attr(feature = "openapi", derive(utoipa::ToSchema))]
#[serde(deny_unknown_fields)]
pub struct CompleteMultipartUploadRequest {
pub content: UploadContentClaim,
pub parts: Vec<CompletedUploadPart>,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
#[cfg_attr(feature = "openapi", derive(utoipa::ToSchema))]
#[serde(tag = "status", 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")]
content_token: Option<ContentToken>,
},
#[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 UploadSessionResponse {
pub namespace_id: NamespaceId,
pub upload_id: UploadId,
pub mode: UploadMode,
#[serde(flatten)]
pub status: UploadSessionStatus,
}
impl UploadSessionResponse {
pub const fn content_ref(&self) -> Option<&ContentRef> {
match &self.status {
UploadSessionStatus::Completed { content_ref, .. } => Some(content_ref),
UploadSessionStatus::Open { .. } | UploadSessionStatus::Aborted { .. } => None,
}
}
pub const fn content_token(&self) -> Option<&ContentToken> {
match &self.status {
UploadSessionStatus::Completed { content_token, .. } => content_token.as_ref(),
UploadSessionStatus::Open { .. } | UploadSessionStatus::Aborted { .. } => None,
}
}
}
#[cfg(test)]
mod tests {
use super::{
BeginUploadRequest, BeginUploadResponse, CompleteKnownContentUploadRequest,
CompleteMultipartUploadRequest, ContentToken, DirectMultipartUpload, DirectPutUpload,
ObjectTransferAccess, UploadContentClaim, UploadMode, UploadSessionResponse,
UploadSessionStatus,
};
use crate::{Checksum, ChecksumAlgorithm, 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,"checksum":{"algorithm":"sha256","value":"2cf24dba5fb0a30e26e83b2ac5b9e29e1b161e5c1fa7425e73043362938b9824"}}}"#,
r#"{"mode":"direct_multipart","content":{"size_bytes":5,"checksum":{"algorithm":"sha256","value":"2cf24dba5fb0a30e26e83b2ac5b9e29e1b161e5c1fa7425e73043362938b9824"}}}"#,
r#"{"mode":"direct_put"}"#,
] {
assert!(
serde_json::from_str::<BeginUploadRequest>(body).is_err(),
"decoded a begin request that mixes modes: {body}"
);
}
}
#[test]
fn completion_request_shapes_reject_fields_the_session_did_not_select() {
assert_eq!(
serde_json::from_str::<CompleteKnownContentUploadRequest>("{}")
.expect("decode known-content completion"),
CompleteKnownContentUploadRequest {}
);
assert_eq!(
serde_json::to_string(&CompleteKnownContentUploadRequest {})
.expect("encode known-content completion"),
"{}"
);
for body in [
r#"{"completion":"content_ref"}"#,
r#"{"content_ref":{"kind":"blob_v1","content_id":"con_0123456789abcdef0123456789abcdef","size_bytes":5,"checksum":{"algorithm":"sha256","value":"2cf24dba5fb0a30e26e83b2ac5b9e29e1b161e5c1fa7425e73043362938b9824"}}}"#,
r#"{"content":{"size_bytes":5,"checksum":{"algorithm":"crc64nvme","value":"0123456789abcdef"}},"parts":[]}"#,
] {
assert!(
serde_json::from_str::<CompleteKnownContentUploadRequest>(body).is_err(),
"decoded a known-content completion carrying fields: {body}"
);
}
let missing_parts = r#"{"content":{"size_bytes":5,"checksum":{"algorithm":"crc64nvme","value":"0123456789abcdef"}}}"#;
let error = serde_json::from_str::<CompleteMultipartUploadRequest>(missing_parts)
.expect_err("multipart parts are required");
assert!(error.to_string().contains("missing field `parts`"));
let multipart = CompleteMultipartUploadRequest {
content: UploadContentClaim {
size_bytes: 5,
checksum: Checksum::crc64nvme(b"hello"),
},
parts: Vec::new(),
};
let encoded = serde_json::to_string(&multipart).expect("encode multipart completion");
assert!(!encoded.contains("completion"));
assert!(!encoded.contains("content_ref"));
}
#[test]
fn direct_put_response_exposes_only_presigned_access() {
let response = BeginUploadResponse::DirectPut {
namespace_id: NamespaceId::parse("demo").expect("namespace id"),
upload_id: UploadId::parse("upl_00000000000000000000000000000001")
.expect("valid upload id"),
direct_put: 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,
},
},
};
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_begin_response_carries_only_its_transports_field() {
let namespace_id = NamespaceId::parse("demo").expect("namespace id");
let upload_id =
UploadId::parse("upl_00000000000000000000000000000001").expect("valid upload id");
let sha256 = "2cf24dba5fb0a30e26e83b2ac5b9e29e1b161e5c1fa7425e73043362938b9824";
assert_eq!(
serde_json::to_value(BeginUploadResponse::ServiceProxied {
namespace_id: namespace_id.clone(),
upload_id: upload_id.clone(),
})
.expect("serialize proxied response"),
serde_json::json!({
"mode": "service_proxied",
"namespace_id": "demo",
"upload_id": "upl_00000000000000000000000000000001"
})
);
assert_eq!(
serde_json::to_value(BeginUploadResponse::DirectPut {
namespace_id: namespace_id.clone(),
upload_id: upload_id.clone(),
direct_put: DirectPutUpload {
content_ref: ContentRef::blob_v1(
ContentId::parse("con_0123456789abcdef0123456789abcdef")
.expect("content id"),
b"hello",
),
access: ObjectTransferAccess::PresignedUrl {
method: "PUT".to_owned(),
url: "https://bucket.example/object".to_owned(),
headers: BTreeMap::new(),
expires_at_ms: 1,
},
},
})
.expect("serialize direct-put response"),
serde_json::json!({
"mode": "direct_put",
"namespace_id": "demo",
"upload_id": "upl_00000000000000000000000000000001",
"direct_put": {
"content_ref": {
"kind": "blob_v1",
"content_id": "con_0123456789abcdef0123456789abcdef",
"size_bytes": 5,
"checksum": { "algorithm": "sha256", "value": sha256 }
},
"access": {
"kind": "presigned_url",
"method": "PUT",
"url": "https://bucket.example/object",
"expires_at_ms": 1
}
}
})
);
assert_eq!(
serde_json::to_value(BeginUploadResponse::DirectMultipart {
namespace_id,
upload_id,
direct_multipart: DirectMultipartUpload {
part_size_bytes: 8 * 1024 * 1024,
checksum_algorithm: ChecksumAlgorithm::Crc64nvme,
},
})
.expect("serialize multipart response"),
serde_json::json!({
"mode": "direct_multipart",
"namespace_id": "demo",
"upload_id": "upl_00000000000000000000000000000001",
"direct_multipart": {
"part_size_bytes": 8 * 1024 * 1024,
"checksum_algorithm": "crc64nvme"
}
})
);
}
#[test]
fn a_begin_response_carrying_a_later_servers_field_still_decodes() {
assert_eq!(
serde_json::from_str::<BeginUploadResponse>(
r#"{"mode":"service_proxied","namespace_id":"demo","upload_id":"upl_00000000000000000000000000000001","invented_later":true}"#
)
.expect("decode a proxied response carrying an unknown field"),
BeginUploadResponse::ServiceProxied {
namespace_id: NamespaceId::parse("demo").expect("namespace id"),
upload_id: UploadId::parse("upl_00000000000000000000000000000001")
.expect("valid upload id"),
}
);
}
#[test]
fn an_upload_content_claim_names_only_size_and_checksum() {
let request: BeginUploadRequest = serde_json::from_str(
r#"{"mode":"direct_put","content":{"size_bytes":5,"checksum":{"algorithm":"sha256","value":"2cf24dba5fb0a30e26e83b2ac5b9e29e1b161e5c1fa7425e73043362938b9824"}}}"#,
)
.expect("decode direct-put begin request");
assert_eq!(
request,
BeginUploadRequest::DirectPut {
content: UploadContentClaim {
size_bytes: 5,
checksum: Checksum {
algorithm: ChecksumAlgorithm::Sha256,
value: "2cf24dba5fb0a30e26e83b2ac5b9e29e1b161e5c1fa7425e73043362938b9824"
.to_owned(),
},
},
}
);
assert!(
serde_json::from_str::<UploadContentClaim>(
r#"{"size_bytes":5,"checksum":{"algorithm":"sha256","value":"2cf24dba5fb0a30e26e83b2ac5b9e29e1b161e5c1fa7425e73043362938b9824"},"content_id":"con_0123456789abcdef0123456789abcdef"}"#
)
.is_err(),
"a client must not be able to name the content object"
);
}
#[test]
fn an_upload_content_claim_carries_the_operations_required_algorithm() {
let claim: UploadContentClaim = serde_json::from_str(
r#"{"size_bytes":5,"checksum":{"algorithm":"crc32c","value":"a1b2c3d4"}}"#,
)
.expect("decode a non-sha256 direct-put claim");
assert_eq!(claim.checksum.algorithm, ChecksumAlgorithm::Crc32c);
}
#[test]
fn upload_session_response_is_flat_and_uses_one_status_vocabulary() {
let namespace_id = NamespaceId::parse("demo").expect("namespace id");
let upload_id = UploadId::parse("upl_00000000000000000000000000000001").expect("upload id");
let open = serde_json::to_value(UploadSessionResponse {
namespace_id: namespace_id.clone(),
upload_id: upload_id.clone(),
mode: UploadMode::DirectMultipart,
status: UploadSessionStatus::Open {
expires_at_ms: 1_000,
},
})
.expect("serialize open status");
assert_eq!(
open,
serde_json::json!({
"namespace_id": "demo",
"upload_id": "upl_00000000000000000000000000000001",
"mode": "direct_multipart",
"status": "open",
"expires_at_ms": 1_000,
})
);
let aborted = serde_json::to_value(UploadSessionResponse {
namespace_id: namespace_id.clone(),
upload_id: upload_id.clone(),
mode: UploadMode::ServiceProxied,
status: UploadSessionStatus::Aborted {
aborted_at_ms: 2_000,
},
})
.expect("serialize aborted status");
assert_eq!(aborted["status"], "aborted");
assert_eq!(aborted["mode"], "service_proxied");
assert_eq!(aborted["aborted_at_ms"], 2_000);
assert!(aborted.get("state").is_none());
let completed = serde_json::to_value(UploadSessionResponse {
namespace_id,
upload_id,
mode: UploadMode::DirectPut,
status: UploadSessionStatus::Completed {
completed_at_ms: 3_000,
content_ref: ContentRef::blob_v1(ContentId::generate(), b"hello"),
content_token: None,
},
})
.expect("serialize completed status");
assert_eq!(completed["status"], "completed");
assert_eq!(completed["mode"], "direct_put");
assert!(completed.get("state").is_none());
assert!(completed.get("status").is_some());
assert!(
completed.get("content_token").is_none(),
"a session past its receipt window reports no token at all"
);
}
#[test]
fn completion_status_and_commit_share_the_exact_content_token_shape() {
let namespace_id = NamespaceId::parse("demo").expect("namespace id");
let upload_id = UploadId::parse("upl_00000000000000000000000000000001").expect("upload id");
let content_ref = ContentRef::blob_v1(
ContentId::parse("con_0123456789abcdef0123456789abcdef").expect("content id"),
b"hello",
);
let content_token = ContentToken {
content_ref: content_ref.clone(),
token: "opaque-server-token".to_owned(),
};
let completion = serde_json::to_value(UploadSessionResponse {
namespace_id: namespace_id.clone(),
upload_id,
mode: UploadMode::ServiceProxied,
status: UploadSessionStatus::Completed {
completed_at_ms: 3_000,
content_ref: content_ref.clone(),
content_token: Some(content_token.clone()),
},
})
.expect("serialize completion");
let status = serde_json::to_value(UploadSessionStatus::Completed {
completed_at_ms: 3_000,
content_ref,
content_token: Some(content_token),
})
.expect("serialize completed status");
let completion_token = completion["content_token"].clone();
let status_token = status["content_token"].clone();
assert_eq!(completion_token, status_token);
assert_eq!(
completion_token,
serde_json::json!({
"content_ref": {
"kind": "blob_v1",
"content_id": "con_0123456789abcdef0123456789abcdef",
"size_bytes": 5,
"checksum": {
"algorithm": "sha256",
"value": "2cf24dba5fb0a30e26e83b2ac5b9e29e1b161e5c1fa7425e73043362938b9824"
}
},
"token": "opaque-server-token"
})
);
let request: crate::v0::CommitRequest = serde_json::from_value(serde_json::json!({
"commit_id": "same-token-shape",
"actor": crate::ActorRef::loonfs_system(),
"content_tokens": [completion_token],
"operations": [{
"kind": "create_directory",
"path": "/proof",
"parents": false
}]
}))
.expect("completion token decodes unchanged in a commit request");
assert_eq!(
serde_json::to_value(&request.content_tokens[0]).expect("serialize commit token"),
status_token
);
}
#[test]
fn a_content_token_rejects_unknown_fields() {
let token = serde_json::json!({
"content_ref": {
"kind": "blob_v1",
"content_id": "con_0123456789abcdef0123456789abcdef",
"size_bytes": 5,
"checksum": {
"algorithm": "sha256",
"value": "2cf24dba5fb0a30e26e83b2ac5b9e29e1b161e5c1fa7425e73043362938b9824"
}
},
"token": "opaque-server-token",
"expires_at_ms": 1
});
assert!(serde_json::from_value::<ContentToken>(token).is_err());
}
}