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(rename_all = "snake_case")]
pub enum UploadMode {
#[default]
ServiceProxied,
DirectPut,
DirectMultipart,
}
impl UploadMode {
pub fn as_str(self) -> &'static str {
match self {
Self::ServiceProxied => "service_proxied",
Self::DirectPut => "direct_put",
Self::DirectMultipart => "direct_multipart",
}
}
}
#[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 {
#[serde(default, skip_serializing_if = "Option::is_none")]
#[cfg_attr(feature = "openapi", schema(nullable = false))]
size_bytes: Option<u64>,
},
#[cfg_attr(feature = "openapi", schema(title = "BeginUploadDirectMultipart"))]
DirectMultipart {
#[serde(default, skip_serializing_if = "Option::is_none")]
#[cfg_attr(feature = "openapi", schema(nullable = false))]
part_size_bytes: Option<u64>,
},
}
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))]
#[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,
checksum_algorithm: ChecksumAlgorithm,
access: ObjectTransferAccess,
},
#[cfg_attr(
feature = "openapi",
schema(title = "BeginUploadResponseDirectMultipart")
)]
DirectMultipart {
namespace_id: NamespaceId,
upload_id: UploadId,
part_size_bytes: u64,
checksum_algorithm: ChecksumAlgorithm,
},
}
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, PartialEq, Eq, Serialize, Deserialize)]
#[cfg_attr(feature = "openapi", derive(utoipa::ToSchema))]
#[serde(tag = "mode", rename_all = "snake_case", deny_unknown_fields)]
pub enum CompleteUploadRequest {
#[cfg_attr(feature = "openapi", schema(title = "CompleteUploadServiceProxied"))]
ServiceProxied {},
#[cfg_attr(feature = "openapi", schema(title = "CompleteUploadDirectPut"))]
DirectPut {
content: UploadContentClaim,
},
#[cfg_attr(feature = "openapi", schema(title = "CompleteUploadDirectMultipart"))]
DirectMultipart {
content: UploadContentClaim,
parts: Vec<CompletedUploadPart>,
},
}
impl CompleteUploadRequest {
pub const 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(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")]
#[cfg_attr(feature = "openapi", schema(nullable = false))]
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 UploadSession {
pub namespace_id: NamespaceId,
pub upload_id: UploadId,
pub mode: UploadMode,
#[serde(flatten)]
pub status: UploadSessionStatus,
}
impl UploadSession {
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, CompleteUploadRequest, ContentToken,
ObjectTransferAccess, UploadContentClaim, UploadMode, UploadSession, UploadSessionStatus,
};
use crate::{Checksum, ChecksumAlgorithm, ContentId, ContentRef, NamespaceId, UploadId};
use std::collections::BTreeMap;
#[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","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","content":{"size_bytes":5,"checksum":{"algorithm":"sha256","value":"2cf24dba5fb0a30e26e83b2ac5b9e29e1b161e5c1fa7425e73043362938b9824"}}}"#,
r#"{"mode":"direct_put","part_size_bytes":8388608}"#,
r#"{"mode":"direct_multipart","size_bytes":5}"#,
] {
assert!(
serde_json::from_str::<BeginUploadRequest>(body).is_err(),
"decoded a begin request that mixes modes: {body}"
);
}
}
#[test]
fn a_multipart_begin_names_its_part_size_beside_the_mode() {
assert_eq!(
serde_json::from_str::<BeginUploadRequest>(
r#"{"mode":"direct_multipart","part_size_bytes":8388608}"#
)
.expect("decode multipart begin request"),
BeginUploadRequest::DirectMultipart {
part_size_bytes: Some(8 * 1024 * 1024),
}
);
assert_eq!(
serde_json::from_str::<BeginUploadRequest>(r#"{"mode":"direct_multipart"}"#)
.expect("decode multipart begin without a part size"),
BeginUploadRequest::DirectMultipart {
part_size_bytes: None,
}
);
assert_eq!(
serde_json::to_value(BeginUploadRequest::DirectMultipart {
part_size_bytes: None,
})
.expect("serialize multipart begin request"),
serde_json::json!({ "mode": "direct_multipart" })
);
}
#[test]
fn completion_requests_are_tagged_and_mode_specific() {
assert_eq!(
serde_json::from_str::<CompleteUploadRequest>(r#"{"mode":"service_proxied"}"#)
.expect("decode proxied completion"),
CompleteUploadRequest::ServiceProxied {}
);
let direct_put = CompleteUploadRequest::DirectPut {
content: UploadContentClaim {
size_bytes: 5,
checksum: Checksum::crc32c(b"hello"),
},
};
assert_eq!(
serde_json::to_value(&direct_put).expect("encode direct-put completion"),
serde_json::json!({
"mode": "direct_put",
"content": {
"size_bytes": 5,
"checksum": Checksum::crc32c(b"hello"),
},
})
);
for body in [
r#"{}"#,
r#"{"mode":"service_proxied","content":{"size_bytes":5,"checksum":{"algorithm":"crc64nvme","value":"0123456789abcdef"}},"parts":[]}"#,
r#"{"mode":"direct_put"}"#,
r#"{"mode":"direct_multipart"}"#,
] {
assert!(
serde_json::from_str::<CompleteUploadRequest>(body).is_err(),
"decoded an invalid completion request: {body}"
);
}
let missing_parts = r#"{"mode":"direct_multipart","content":{"size_bytes":5,"checksum":{"algorithm":"crc64nvme","value":"0123456789abcdef"}}}"#;
let error = serde_json::from_str::<CompleteUploadRequest>(missing_parts)
.expect_err("multipart parts are required");
assert!(
error.to_string().contains("parts"),
"the rejection should name the missing field: {error}"
);
let multipart = CompleteUploadRequest::DirectMultipart {
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_eq!(
serde_json::from_str::<serde_json::Value>(&encoded).expect("decode multipart JSON"),
serde_json::json!({
"mode": "direct_multipart",
"content": {
"size_bytes": 5,
"checksum": Checksum::crc64nvme(b"hello"),
},
"parts": [],
})
);
}
#[test]
fn a_begin_response_carries_only_its_transports_fields() {
let namespace_id = NamespaceId::parse("demo").expect("namespace id");
let upload_id =
UploadId::parse("upl_00000000000000000000000000000001").expect("valid upload id");
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(),
checksum_algorithm: ChecksumAlgorithm::Crc64nvme,
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",
"checksum_algorithm": "crc64nvme",
"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,
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",
"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","size_bytes":5}"#)
.expect("decode direct-put begin request");
assert_eq!(
request,
BeginUploadRequest::DirectPut {
size_bytes: Some(5),
}
);
assert_eq!(
serde_json::from_str::<BeginUploadRequest>(r#"{"mode":"direct_put"}"#)
.expect("decode direct-put begin without a size"),
BeginUploadRequest::DirectPut { size_bytes: None }
);
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_session_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(UploadSession {
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(UploadSession {
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(UploadSession {
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(UploadSession {
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());
}
}