pull/370/head
zijiren233 3 months ago
parent 9e0ab44000
commit 85105fd2d3
No known key found for this signature in database
GPG Key ID: 534E082AAA9B39DC

123
Cargo.lock generated

@ -586,7 +586,7 @@ dependencies = [
"sha1 0.10.7",
"sync_wrapper",
"tokio",
"tokio-tungstenite",
"tokio-tungstenite 0.29.0",
"tower",
"tower-layer",
"tower-service",
@ -1313,12 +1313,13 @@ source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "217698eaf96b4a3f0bc4f3662aaa55bdf913cd54d7204591faa790070c6d0853"
[[package]]
name = "crc32c"
version = "0.6.8"
name = "crc-fast"
version = "1.10.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "3a47af21622d091a8f0fb295b88bc886ac74efcc613efc19f5d0b21de5c89e47"
checksum = "e75b2483e97a5a7da73ac68a05b629f9c53cff58d8ed1c77866079e18b00dba5"
dependencies = [
"rustc_version",
"digest 0.10.7",
"spin 0.10.0",
]
[[package]]
@ -2291,7 +2292,7 @@ checksum = "5e139bc46ca777eb5efaf62df0ab8cc5fd400866427e56c68b22e414e53bd3be"
dependencies = [
"futures-core",
"futures-sink",
"spin",
"spin 0.9.8",
]
[[package]]
@ -3546,7 +3547,7 @@ version = "1.5.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "bbd2bcb4c963f2ddae06a2efc7e9f3591312473c50c6685e1f298068316e66fe"
dependencies = [
"spin",
"spin 0.9.8",
]
[[package]]
@ -4173,9 +4174,9 @@ dependencies = [
[[package]]
name = "opendal"
version = "0.57.0"
version = "0.58.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "96c9c85ce253ff87225e7669979d877a20c98a06604ec9d6dd5f4473e08f1ae1"
checksum = "77d02c6564e376d3670aaf66ad886cd34f83c0aca407b6364777c624a63e0e2d"
dependencies = [
"opendal-core",
"opendal-service-s3",
@ -4183,24 +4184,22 @@ dependencies = [
[[package]]
name = "opendal-core"
version = "0.57.0"
version = "0.58.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "c4f8607c90e2c963a91467f50fb49fbc7fb3d573f88cea219ca59ccd3740b309"
checksum = "f8564bd76b75d2aea59178cb5b6e9f770be1da73b1bca062720ec151cc56215a"
dependencies = [
"anyhow",
"base64 0.22.1",
"bytes",
"futures",
"http",
"http-body",
"jiff",
"log",
"md-5",
"mea",
"percent-encoding",
"quick-xml 0.39.4",
"quick-xml",
"reqsign-core",
"reqwest",
"serde",
"serde_json",
"tokio",
@ -4209,20 +4208,34 @@ dependencies = [
"web-time",
]
[[package]]
name = "opendal-http-transport-reqwest"
version = "0.58.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "5b7fd001b204df76be5d2b3f7a81c48b5cf25b28e2cfb3b8a819bb89cdc0ea3a"
dependencies = [
"bytes",
"futures",
"http",
"http-body",
"opendal-core",
"reqwest",
]
[[package]]
name = "opendal-service-s3"
version = "0.57.0"
version = "0.58.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "313d46c9f5ae70bca26b7c3e3fbb9b639292625f28af73aa016f47e788af9deb"
checksum = "750d9cc8588c19b27c4c2c5d06d7eb6f2bff62549c48652f1f9a7603cd872377"
dependencies = [
"base64 0.22.1",
"bytes",
"crc32c",
"crc-fast",
"http",
"log",
"md-5",
"opendal-core",
"quick-xml 0.39.4",
"quick-xml",
"reqsign-aws-v4",
"reqsign-core",
"reqsign-file-read-tokio",
@ -4911,9 +4924,9 @@ dependencies = [
[[package]]
name = "prost-protovalidate"
version = "0.5.0"
version = "0.6.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "2fbb2b6b23d0677360b58c543d0c7bd2aa806beaf9f5ccf0620994bb12ff76e8"
checksum = "bd223ccf88bbaf6d15f5b93eefa4c955da02f0f1da58896e146ef8d837539405"
dependencies = [
"cel",
"chrono",
@ -4929,9 +4942,9 @@ dependencies = [
[[package]]
name = "prost-protovalidate-types"
version = "0.5.0"
version = "0.6.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "4461a14164494eb78f4319a5650576bf0dd7a27781fb55ebb4e031e3e34c8744"
checksum = "4406a3b6ef7227ab6d5d264b9ae93fef5264a624954c41db12b0a91b7d098d8f"
dependencies = [
"prost",
"prost-build",
@ -4943,9 +4956,9 @@ dependencies = [
[[package]]
name = "prost-reflect"
version = "0.16.4"
version = "0.16.5"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "590aa145fee8f7a26b5a6055365e7c5e89a5c1caae9869de76ec0ee73181a2f9"
checksum = "01b80ea363c31af2de2b92e3c07ed1156628f7838c4afb4df75ee78a37fedbd1"
dependencies = [
"base64 0.22.1",
"prost",
@ -4957,9 +4970,9 @@ dependencies = [
[[package]]
name = "prost-reflect-build"
version = "0.16.0"
version = "0.16.1"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "8214ae2c30bbac390db0134d08300e770ef89b6d4e5abf855e8d300eded87e28"
checksum = "95a9e8261adf6617d5dc2a5a9e75cce5ab9d546a007f6f870f809a1ad25386b6"
dependencies = [
"prost-build",
"prost-reflect",
@ -5155,16 +5168,6 @@ version = "2.0.1"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "a993555f31e5a609f617c12db6250dedcac1b0a85076912c436e6fc9b2c8e6a3"
[[package]]
name = "quick-xml"
version = "0.39.4"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "cdcc8dd4e2f670d309a5f0e83fe36dfdc05af317008fea29144da1a2ac858e5e"
dependencies = [
"memchr",
"serde",
]
[[package]]
name = "quick-xml"
version = "0.41.0"
@ -5539,7 +5542,7 @@ dependencies = [
"http",
"log",
"percent-encoding",
"quick-xml 0.41.0",
"quick-xml",
"reqsign-core",
"rust-ini",
"serde",
@ -6422,6 +6425,12 @@ dependencies = [
"lock_api",
]
[[package]]
name = "spin"
version = "0.10.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "d5fe4ccb98d9c292d56fec89a5e07da7fc4cf0dc11e156b41793132775d3e591"
[[package]]
name = "spinning_top"
version = "0.3.0"
@ -6949,7 +6958,7 @@ dependencies = [
"tokio",
"tokio-rustls",
"tokio-stream",
"tokio-tungstenite",
"tokio-tungstenite 0.30.0",
"tokio-util",
"toml",
"tonic",
@ -7034,7 +7043,7 @@ dependencies = [
"thiserror 2.0.18",
"tokio",
"tokio-stream",
"tokio-tungstenite",
"tokio-tungstenite 0.30.0",
"tokio-util",
"tonic",
"tonic-health",
@ -7148,6 +7157,7 @@ dependencies = [
"oauth2",
"opaque-ke",
"opendal",
"opendal-http-transport-reqwest",
"parking_lot",
"paste",
"percent-encoding",
@ -7327,7 +7337,7 @@ dependencies = [
"tokio",
"tokio-rustls",
"tokio-stream",
"tokio-tungstenite",
"tokio-tungstenite 0.30.0",
"tonic",
"tonic-health",
"tonic-prost",
@ -7453,6 +7463,7 @@ dependencies = [
"indexmap 2.14.0",
"lru",
"opendal",
"opendal-http-transport-reqwest",
"parking_lot",
"rand 0.10.2",
"serde",
@ -7783,6 +7794,18 @@ name = "tokio-tungstenite"
version = "0.29.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "8f72a05e828585856dacd553fba484c242c46e391fb0e58917c942ee9202915c"
dependencies = [
"futures-util",
"log",
"tokio",
"tungstenite 0.29.0",
]
[[package]]
name = "tokio-tungstenite"
version = "0.30.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "17a073bfed563fa236697a068031408a93cd9522e08abf9933ead3e73411bd71"
dependencies = [
"futures-util",
"log",
@ -7791,7 +7814,7 @@ dependencies = [
"rustls-pki-types",
"tokio",
"tokio-rustls",
"tungstenite",
"tungstenite 0.30.0",
"webpki-roots 0.26.11",
]
@ -8152,9 +8175,25 @@ dependencies = [
"httparse",
"log",
"rand 0.9.5",
"sha1 0.10.7",
"thiserror 2.0.18",
]
[[package]]
name = "tungstenite"
version = "0.30.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "e48ac77174b19c110a50ab2128b24215ac9cb40e0e12e093fb602d175c569d22"
dependencies = [
"bytes",
"data-encoding",
"http",
"httparse",
"log",
"rand 0.10.2",
"rustls",
"rustls-pki-types",
"sha1 0.10.7",
"sha1 0.11.0",
"thiserror 2.0.18",
]

@ -98,9 +98,9 @@ tonic-prost-build = "=0.14.6"
protoc-bin-vendored = "=3.2.0"
prost = "=0.14.4"
prost-types = "=0.14.4"
prost-reflect = "=0.16.4"
prost-reflect-build = "=0.16.0"
prost-protovalidate = "=0.5.0"
prost-reflect = "=0.16.5"
prost-reflect-build = "=0.16.1"
prost-protovalidate = "=0.6.0"
pbjson = "=0.9.0"
pbjson-build = "=0.9.0"
pbjson-types = "=0.9.0"
@ -200,10 +200,11 @@ iana-time-zone = "=0.1.65"
synctv-xiu = { path = "synctv-xiu", default-features = false }
# Object Storage
opendal = { version = "=0.57.0", default-features = false, features = [
opendal = { version = "=0.58.0", default-features = false, features = [
"executors-tokio",
"services-s3",
] }
opendal-http-transport-reqwest = { version = "=0.58.0", default-features = false, features = ["rustls-no-provider"] }
# Logging & Tracing
tracing = "=0.1.44"
@ -308,7 +309,7 @@ testcontainers-modules = { version = "=0.15.0", default-features = false, featur
"redis",
"rustfs",
] }
tokio-tungstenite = { version = "=0.29.0", default-features = false, features = ["connect"] }
tokio-tungstenite = { version = "=0.30.0", default-features = false, features = ["connect"] }
# Benchmarking
criterion = { version = "=0.8.2", features = ["html_reports", "async_tokio"] }

@ -66,7 +66,7 @@ impl GoogleApiError {
crate::impls::ApiError::RangeNotSatisfiable { total_size } => {
details.set_resource_info(
"byte_range",
total_size.max(&0).to_string(),
total_size.to_string(),
"",
"Requested byte range is not satisfiable",
);

@ -265,7 +265,7 @@ impl From<crate::impls::ApiError> for AppError {
extra_headers: Vec::new(),
};
if let crate::impls::ApiError::RangeNotSatisfiable { total_size } = err {
if let Ok(value) = HeaderValue::from_str(&format!("bytes */{}", total_size.max(0))) {
if let Ok(value) = HeaderValue::from_str(&format!("bytes */{total_size}")) {
app_err = app_err.with_header(header::CONTENT_RANGE, value);
}
}

@ -153,9 +153,9 @@ pub(crate) fn optional_file_range(
("", "") => Err(AppError::bad_request("Invalid Range header")),
("", suffix) => {
let length = suffix
.parse::<i64>()
.parse::<u64>()
.map_err(|_| AppError::bad_request("Invalid Range header"))?;
if length <= 0 {
if length == 0 {
return Err(AppError::bad_request("Invalid Range header"));
}
Ok(Some(synctv_core::models::FileRangeRequest::Suffix {
@ -164,21 +164,18 @@ pub(crate) fn optional_file_range(
}
(start, "") => {
let start = start
.parse::<i64>()
.parse::<u64>()
.map_err(|_| AppError::bad_request("Invalid Range header"))?;
if start < 0 {
return Err(AppError::bad_request("Invalid Range header"));
}
Ok(Some(synctv_core::models::FileRangeRequest::From { start }))
}
(start, end) => {
let start = start
.parse::<i64>()
.parse::<u64>()
.map_err(|_| AppError::bad_request("Invalid Range header"))?;
let end_inclusive = end
.parse::<i64>()
.parse::<u64>()
.map_err(|_| AppError::bad_request("Invalid Range header"))?;
if start < 0 || end_inclusive < start {
if end_inclusive < start {
return Err(AppError::bad_request("Invalid Range header"));
}
Ok(Some(synctv_core::models::FileRangeRequest::Exact(

@ -325,9 +325,7 @@ fn map_proxy_execution_error(err: anyhow::Error) -> AppError {
}
Some(synctv_proxy::ProxyErrorKind::RangeNotSatisfiable) => {
if let Some(total_size) = synctv_proxy::proxy_range_not_satisfiable_total_size(&err) {
AppError::from(ApiError::RangeNotSatisfiable {
total_size: i64::try_from(total_size).unwrap_or(i64::MAX),
})
AppError::from(ApiError::RangeNotSatisfiable { total_size })
} else {
AppError::new(StatusCode::RANGE_NOT_SATISFIABLE, err.to_string())
}

@ -467,7 +467,7 @@ pub enum ApiError {
violations: Vec<ApiFieldViolation>,
},
RangeNotSatisfiable {
total_size: i64,
total_size: u64,
},
BadGateway(String),
RequestTimeout(String),

@ -1023,9 +1023,7 @@ fn map_proxy_execution_error(err: &anyhow::Error) -> ApiError {
ApiError::Authorization("Proxy target is not allowed by SSRF policy".to_string())
}
Some(synctv_proxy::ProxyErrorKind::RangeNotSatisfiable) => ApiError::RangeNotSatisfiable {
total_size: synctv_proxy::proxy_range_not_satisfiable_total_size(err)
.and_then(|size| i64::try_from(size).ok())
.unwrap_or(0),
total_size: synctv_proxy::proxy_range_not_satisfiable_total_size(err).unwrap_or(0),
},
Some(synctv_proxy::ProxyErrorKind::InvalidRequest) => {
ApiError::InvalidInput(err.to_string())

@ -59,6 +59,7 @@ indexmap.workspace = true
# Utilities
anyhow.workspace = true
opendal.workspace = true
opendal-http-transport-reqwest.workspace = true
thiserror.workspace = true
bytes.workspace = true
parking_lot.workspace = true

@ -42,7 +42,7 @@ pub enum Error {
InvalidInput(String),
#[error("Range not satisfiable: total size {total_size}")]
RangeNotSatisfiable { total_size: i64 },
RangeNotSatisfiable { total_size: u64 },
#[error("Rate limited: {0}")]
RateLimited(String),

@ -727,13 +727,13 @@ pub enum StoreFileUploadResult {
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
pub struct FileByteRange {
pub start: i64,
pub end_inclusive: i64,
pub start: u64,
pub end_inclusive: u64,
}
impl FileByteRange {
#[must_use]
pub const fn size_bytes(self) -> i64 {
pub const fn size_bytes(self) -> u64 {
self.end_inclusive - self.start + 1
}
}
@ -741,8 +741,8 @@ impl FileByteRange {
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
pub enum FileRangeRequest {
Exact(FileByteRange),
From { start: i64 },
Suffix { length: i64 },
From { start: u64 },
Suffix { length: u64 },
}
#[derive(Debug, Clone, Serialize, Deserialize)]

@ -3,13 +3,13 @@ use sqlx::{PgPool, Postgres, Transaction};
use crate::{
models::{
FileBlob, FileBlobCompression, FileBlobPart, FileCleanupJob, FileCleanupMetadata,
FileMetadata, FileObject, FileObjectGroup, FileObjectVariant, FileReferenceMetadata,
FileReferenceTarget, FileUploadSessionKind, FileUploadSessionMetadata,
FileUploadSessionPart, FileUploadSessionRecord, FileVariantMetadata, StoredFileReference,
UserId, FILE_CLEANUP_ORIGIN_MAX_CHARS, FILE_OBJECT_KEY_MAX_CHARS,
FILE_REFERENCE_ID_MAX_CHARS, FILE_REFERENCE_KIND_MAX_CHARS, FILE_SHA256_HEX_CHARS,
FILE_STORAGE_BACKEND_MAX_CHARS,
FileBlob, FileBlobCompression, FileBlobPart, FileByteRange, FileCleanupJob,
FileCleanupMetadata, FileMetadata, FileObject, FileObjectGroup, FileObjectVariant,
FileReferenceMetadata, FileReferenceTarget, FileUploadSessionKind,
FileUploadSessionMetadata, FileUploadSessionPart, FileUploadSessionRecord,
FileVariantMetadata, StoredFileReference, UserId, FILE_CLEANUP_ORIGIN_MAX_CHARS,
FILE_OBJECT_KEY_MAX_CHARS, FILE_REFERENCE_ID_MAX_CHARS, FILE_REFERENCE_KIND_MAX_CHARS,
FILE_SHA256_HEX_CHARS, FILE_STORAGE_BACKEND_MAX_CHARS,
},
Error, Result,
};
@ -1389,12 +1389,15 @@ impl FileStorageRepository {
&self,
storage_backend: &str,
object_key: &str,
start: i64,
end_inclusive: i64,
range: FileByteRange,
) -> Result<Vec<FileBlobPart>> {
if start < 0 || end_inclusive < start {
if range.end_inclusive < range.start {
return Err(Error::InvalidInput("file range is invalid".to_string()));
}
let start = i64::try_from(range.start)
.map_err(|_| Error::InvalidInput("file range exceeds database limits".to_string()))?;
let end_inclusive = i64::try_from(range.end_inclusive)
.map_err(|_| Error::InvalidInput("file range exceeds database limits".to_string()))?;
let rows = sqlx::query_as!(
FileBlobPart,
r#"

@ -1350,8 +1350,7 @@ async fn s3_file_storage_rejects_tampered_upload_session_image() {
let operator = ok(
Operator::new(opendal::services::Memory::default()),
"memory operator should build",
)
.finish();
);
let service = ok(
S3CompatibleFileStorageService::new_with_repository(
test_s3_file_storage_config(),
@ -1422,8 +1421,7 @@ async fn s3_file_storage_creates_resumable_upload_session() {
let operator = ok(
Operator::new(opendal::services::Memory::default()),
"memory operator should build",
)
.finish();
);
let service = ok(
S3CompatibleFileStorageService::new_with_repository(
test_s3_file_storage_config(),
@ -1513,8 +1511,7 @@ async fn image_upload_sessions_resume_pending_session_for_reused_client_ids() {
let operator = ok(
Operator::new(opendal::services::Memory::default()),
"memory operator should build",
)
.finish();
);
let service = ok(
S3CompatibleFileStorageService::new_with_repository(
test_s3_file_storage_config(),
@ -1589,8 +1586,7 @@ async fn s3_file_storage_reuses_registered_object_with_ownership_proof() {
let operator = ok(
Operator::new(opendal::services::Memory::default()),
"memory operator should build",
)
.finish();
);
ok(
operator.write(&object_key, b"data".to_vec()).await,
"object should be written",

@ -333,34 +333,30 @@ pub(super) fn upload_manifest_is_single_object(
pub(super) fn resolve_file_range(
request: Option<FileRangeRequest>,
total_size_bytes: i64,
total_size_bytes: u64,
) -> Result<Option<FileByteRange>> {
let Some(request) = request else {
return Ok(None);
};
if total_size_bytes <= 0 {
return Err(Error::RangeNotSatisfiable {
total_size: total_size_bytes.max(0),
});
if total_size_bytes == 0 {
return Err(Error::RangeNotSatisfiable { total_size: 0 });
}
let range = match request {
FileRangeRequest::Exact(range) => {
if range.start < 0 || range.end_inclusive < range.start {
return Err(Error::InvalidInput("file range is invalid".to_string()));
}
range
}
FileRangeRequest::From { start } => {
if start < 0 {
if range.end_inclusive < range.start {
return Err(Error::InvalidInput("file range is invalid".to_string()));
}
FileByteRange {
start,
end_inclusive: total_size_bytes - 1,
start: range.start,
end_inclusive: range.end_inclusive.min(total_size_bytes - 1),
}
}
FileRangeRequest::From { start } => FileByteRange {
start,
end_inclusive: total_size_bytes - 1,
},
FileRangeRequest::Suffix { length } => {
if length <= 0 {
if length == 0 {
return Err(Error::InvalidInput("file range is invalid".to_string()));
}
let size = length.min(total_size_bytes);

@ -127,26 +127,42 @@ impl DatabaseFileStorageService {
else {
return Err(Error::NotFound("File object not found".to_string()));
};
let range = super::resolve_file_range(request.range, object.size_bytes)?;
let total_size = u64::try_from(object.size_bytes)
.map_err(|_| Error::Internal("file object size is invalid".to_string()))?;
let range = super::resolve_file_range(request.range, total_size)?;
if total_size == 0 {
return Ok(FileObjectDownload {
metadata: FileObjectMetadata {
storage_backend: self.storage_backend.clone(),
object_key,
mime_type: object.mime_type,
size_bytes: 0,
total_size_bytes: 0,
content_manifest_sha256: object.content_manifest_sha256,
compression: FileBlobCompression::None,
range: None,
metadata: object.metadata,
created_at: object.created_at,
},
stream: futures::stream::empty().boxed(),
});
}
let read_range = range.unwrap_or(FileByteRange {
start: 0,
end_inclusive: object.size_bytes - 1,
end_inclusive: total_size - 1,
});
let parts = self
.repository
.list_blob_parts_overlapping_range(
&self.storage_backend,
&object_key,
read_range.start,
read_range.end_inclusive,
)
.list_blob_parts_overlapping_range(&self.storage_backend, &object_key, read_range)
.await?;
ensure_database_parts_cover_range(&parts, read_range)?;
let metadata = FileObjectMetadata {
storage_backend: self.storage_backend.clone(),
object_key,
mime_type: object.mime_type,
size_bytes: read_range.size_bytes(),
size_bytes: i64::try_from(read_range.size_bytes()).map_err(|_| {
Error::Internal("file range size exceeds database limits".to_string())
})?,
total_size_bytes: object.size_bytes,
content_manifest_sha256: object.content_manifest_sha256,
compression: FileBlobCompression::None,
@ -329,7 +345,7 @@ struct DecompressedBlobPart {
data: Vec<u8>,
}
fn range_end_exclusive(range: FileByteRange) -> Result<i64> {
fn range_end_exclusive(range: FileByteRange) -> Result<u64> {
range
.end_inclusive
.checked_add(1)
@ -345,14 +361,17 @@ fn ensure_database_parts_cover_range(parts: &[FileBlobPart], range: FileByteRang
let end_exclusive = range_end_exclusive(range)?;
let mut expected = range.start;
for part in parts {
let part_end = part
.offset_bytes
.checked_add(part.size_bytes)
let part_start = u64::try_from(part.offset_bytes)
.map_err(|_| Error::Internal("file blob part offset is invalid".to_string()))?;
let part_size = u64::try_from(part.size_bytes)
.map_err(|_| Error::Internal("file blob part size is invalid".to_string()))?;
let part_end = part_start
.checked_add(part_size)
.ok_or_else(|| Error::Internal("file blob part offset overflow".to_string()))?;
if part_end <= expected {
continue;
}
if part.offset_bytes > expected {
if part_start > expected {
return Err(Error::Internal(
"file object is missing one or more blob parts".to_string(),
));
@ -369,10 +388,12 @@ fn ensure_database_parts_cover_range(parts: &[FileBlobPart], range: FileByteRang
async fn database_part_chunk(part: FileBlobPart, range: FileByteRange) -> Result<bytes::Bytes> {
let part_data = decompress_blob_part(part).await?;
let part_start_absolute = part_data.offset_bytes;
let part_end_absolute = part_data
.offset_bytes
.checked_add(part_data.size_bytes)
let part_start_absolute = u64::try_from(part_data.offset_bytes)
.map_err(|_| Error::Internal("file blob part offset is invalid".to_string()))?;
let part_size = u64::try_from(part_data.size_bytes)
.map_err(|_| Error::Internal("file blob part size is invalid".to_string()))?;
let part_end_absolute = part_start_absolute
.checked_add(part_size)
.ok_or_else(|| Error::Internal("file blob part offset overflow".to_string()))?;
let read_start_absolute = range.start.max(part_start_absolute);
let read_end_absolute = range_end_exclusive(range)?.min(part_end_absolute);
@ -1310,12 +1331,24 @@ impl FileStorageService for DatabaseFileStorageService {
let ranges = ownership_proof.ranges.clone();
let mut chunks = Vec::with_capacity(ranges.len());
for range in &ranges {
let start = u64::try_from(range.offset).map_err(|_| {
Error::Internal("file ownership proof range is invalid".to_string())
})?;
let length = u64::try_from(range.length).map_err(|_| {
Error::Internal("file ownership proof range is invalid".to_string())
})?;
let end_inclusive = start
.checked_add(length)
.and_then(|end| end.checked_sub(1))
.ok_or_else(|| {
Error::Internal("file ownership proof range is invalid".to_string())
})?;
let data = self
.load_range_data(
&object_key,
Some(FileRangeRequest::Exact(FileByteRange {
start: range.offset,
end_inclusive: range.offset + i64::from(range.length) - 1,
start,
end_inclusive,
})),
)
.await?;
@ -1482,16 +1515,15 @@ impl FileStorageService for DatabaseFileStorageService {
.and_then(|end| end.checked_sub(1))
.ok_or_else(|| Error::Internal("file reader range overflow".to_string()))?;
let range = FileByteRange {
start: offset,
end_inclusive,
start: u64::try_from(offset).map_err(|_| {
Error::Internal("file reader offset is invalid".to_string())
})?,
end_inclusive: u64::try_from(end_inclusive).map_err(|_| {
Error::Internal("file reader range is invalid".to_string())
})?,
};
let parts = repository
.list_blob_parts_overlapping_range(
&storage_backend,
&object_key,
range.start,
range.end_inclusive,
)
.list_blob_parts_overlapping_range(&storage_backend, &object_key, range)
.await?;
ensure_database_parts_cover_range(&parts, range)?;
let mut out = bytes::BytesMut::with_capacity(length);

@ -122,10 +122,13 @@ impl S3CompatibleFileStorageService {
} else {
self.stat_object(&object_key).await
};
let total_size_bytes = object.as_ref().map(|object| object.size_bytes).or_else(|| {
stat.as_ref()
.and_then(|stat| i64::try_from(stat.content_length()).ok())
});
let total_size_bytes = match object.as_ref() {
Some(object) => Some(
u64::try_from(object.size_bytes)
.map_err(|_| Error::Internal("file object size is invalid".to_string()))?,
),
None => stat.as_ref().map(opendal::Metadata::content_length),
};
let Some(total_size_bytes) = total_size_bytes else {
if request.range.is_some() {
return Err(Error::InvalidInput(
@ -172,18 +175,22 @@ impl S3CompatibleFileStorageService {
stream: futures::stream::once(async move { Ok(data.to_bytes()) }).boxed(),
});
};
let range = super::super::resolve_file_range(request.range, total_size_bytes)?;
let read_range = range.unwrap_or(FileByteRange {
start: 0,
end_inclusive: total_size_bytes - 1,
});
let start = u64::try_from(read_range.start)
.map_err(|_| Error::InvalidInput("file range is invalid".to_string()))?;
let end = read_range
.end_inclusive
.checked_add(1)
.and_then(|end| u64::try_from(end).ok())
.ok_or_else(|| Error::InvalidInput("file range is invalid".to_string()))?;
let requested_range = request.range;
let range = super::super::resolve_file_range(requested_range, total_size_bytes)?;
let storage_range = match requested_range {
None => opendal::BytesRange::default(),
Some(crate::models::FileRangeRequest::Exact(_)) => {
let resolved = range
.ok_or_else(|| Error::Internal("resolved file range is missing".to_string()))?;
opendal::BytesRange::new(resolved.start, Some(resolved.size_bytes()))
}
Some(crate::models::FileRangeRequest::From { start }) => {
opendal::BytesRange::new(start, None)
}
Some(crate::models::FileRangeRequest::Suffix { length }) => {
opendal::BytesRange::suffix(length)
}
};
let mime_type = object
.as_ref()
.map(|object| object.mime_type.clone())
@ -208,7 +215,7 @@ impl S3CompatibleFileStorageService {
.await
.map_err(|error| Error::NotFound(format!("File object not found: {error}")))?;
let bytes_stream = reader
.into_bytes_stream(start..end)
.into_bytes_stream(storage_range)
.await
.map_err(|error| Error::NotFound(format!("File object not found: {error}")))?;
let stream = bytes_stream
@ -223,8 +230,15 @@ impl S3CompatibleFileStorageService {
storage_backend: self.config.storage_backend.clone(),
object_key,
mime_type,
size_bytes: read_range.size_bytes(),
total_size_bytes,
size_bytes: i64::try_from(
range.map_or(total_size_bytes, FileByteRange::size_bytes),
)
.map_err(|_| {
Error::Internal("file range size exceeds storage limits".to_string())
})?,
total_size_bytes: i64::try_from(total_size_bytes).map_err(|_| {
Error::Internal("file object size exceeds storage limits".to_string())
})?,
content_manifest_sha256,
compression: FileBlobCompression::None,
range,

@ -4,6 +4,10 @@ use super::S3FileStorageConfig;
use crate::{Error, Result};
pub(super) fn s3_operator_from_config(config: &S3FileStorageConfig) -> Result<Operator> {
crate::install_process_crypto_provider();
opendal::HttpTransporter::install_default(
opendal_http_transport_reqwest::ReqwestTransport::default(),
);
let mut builder = S3::default()
.endpoint(config.endpoint.trim())
.access_key_id(config.access_key_id.trim())
@ -18,7 +22,6 @@ pub(super) fn s3_operator_from_config(config: &S3FileStorageConfig) -> Result<Op
}
Operator::new(builder)
.map(opendal::OperatorBuilder::finish)
.map_err(|error| Error::Internal(format!("failed to initialize S3 file storage: {error}")))
}

@ -44,6 +44,68 @@ fn empty_reference_metadata() -> FileReferenceMetadata {
FileReferenceMetadata::File(FileMetadata::default())
}
#[test]
fn file_ranges_resolve_to_consistent_closed_intervals() {
assert_eq!(
resolve_file_range(Some(FileRangeRequest::From { start: 900 }), 1_000)
.expect("from range should resolve"),
Some(FileByteRange {
start: 900,
end_inclusive: 999,
})
);
assert_eq!(
resolve_file_range(Some(FileRangeRequest::Suffix { length: 100 }), 1_000)
.expect("suffix range should resolve"),
Some(FileByteRange {
start: 900,
end_inclusive: 999,
})
);
assert_eq!(
resolve_file_range(Some(FileRangeRequest::Suffix { length: 100 }), 80)
.expect("large suffix should resolve to the full object"),
Some(FileByteRange {
start: 0,
end_inclusive: 79,
})
);
assert_eq!(
resolve_file_range(
Some(FileRangeRequest::Exact(FileByteRange {
start: 60,
end_inclusive: 100,
})),
80,
)
.expect("exact range should clamp to the object end"),
Some(FileByteRange {
start: 60,
end_inclusive: 79,
})
);
assert_eq!(
resolve_file_range(None, 0).expect("full empty object should resolve"),
None
);
}
#[test]
fn file_ranges_reject_invalid_and_unsatisfiable_requests() {
assert!(matches!(
resolve_file_range(Some(FileRangeRequest::Suffix { length: 0 }), 100),
Err(Error::InvalidInput(_))
));
assert!(matches!(
resolve_file_range(Some(FileRangeRequest::From { start: 100 }), 100),
Err(Error::RangeNotSatisfiable { total_size: 100 })
));
assert!(matches!(
resolve_file_range(Some(FileRangeRequest::Suffix { length: 1 }), 0),
Err(Error::RangeNotSatisfiable { total_size: 0 })
));
}
fn single_manifest_part(
size_bytes: i64,
checksum_sha256: impl Into<String>,
@ -2690,6 +2752,24 @@ async fn database_storage_resumable_upload_completes_after_all_parts() {
assert_eq!(ranged.range, Some(requested));
assert_eq!(ranged.total_size_bytes, payload_size(&payload));
assert_eq!(ranged.data, payload[3..9]);
let suffix = ok(
storage
.get_object(GetFileObject {
encoded_object_key: read_encoded_object_key.clone(),
read_token: read_token.clone(),
range: Some(FileRangeRequest::Suffix { length: 4 }),
})
.await,
"database suffix range should read",
);
assert_eq!(suffix.data, payload[payload.len() - 4..]);
assert_eq!(
suffix.range,
Some(FileByteRange {
start: u64::try_from(payload.len() - 4).expect("test offset should fit"),
end_inclusive: u64::try_from(payload.len() - 1).expect("test offset should fit"),
})
);
}
#[tokio::test]
@ -3162,7 +3242,10 @@ async fn database_storage_range_reads_from_permanent_blob_parts() {
);
assert_eq!(loaded.range, Some(requested));
assert_eq!(loaded.total_size_bytes, payload_size(&payload));
assert_eq!(loaded.size_bytes, requested.size_bytes());
assert_eq!(
loaded.size_bytes,
i64::try_from(requested.size_bytes()).expect("test range should fit i64")
);
assert_eq!(loaded.data, payload[1024..4096]);
}
@ -3171,8 +3254,7 @@ async fn s3_storage_multipart_session_returns_native_part_urls() {
let (_postgres, pool) = synctv_core_testing::create_test_pool().await;
let repository = Arc::new(FileStorageRepository::new(pool));
let operator = ok(
opendal::Operator::new(opendal::services::Memory::default())
.map(opendal::OperatorBuilder::finish),
opendal::Operator::new(opendal::services::Memory::default()),
"memory operator should build",
);
let storage = ok(
@ -3252,8 +3334,7 @@ async fn s3_storage_single_object_session_uses_backend_proxy_upload() {
let (_postgres, pool) = synctv_core_testing::create_test_pool().await;
let repository = Arc::new(FileStorageRepository::new(pool));
let operator = ok(
opendal::Operator::new(opendal::services::Memory::default())
.map(opendal::OperatorBuilder::finish),
opendal::Operator::new(opendal::services::Memory::default()),
"memory operator should build",
);
let storage = ok(
@ -3371,8 +3452,7 @@ async fn s3_storage_streams_range_from_backend_proxy_path() {
let (_postgres, pool) = synctv_core_testing::create_test_pool().await;
let repository = Arc::new(FileStorageRepository::new(pool));
let operator = ok(
opendal::Operator::new(opendal::services::Memory::default())
.map(opendal::OperatorBuilder::finish),
opendal::Operator::new(opendal::services::Memory::default()),
"memory operator should build",
);
let storage = ok(
@ -3427,15 +3507,18 @@ async fn s3_storage_streams_range_from_backend_proxy_path() {
let download = ok(
storage
.get_object_stream(GetFileObject {
encoded_object_key,
read_token,
encoded_object_key: encoded_object_key.clone(),
read_token: read_token.clone(),
range: Some(FileRangeRequest::Exact(requested)),
})
.await,
"S3 object range should stream",
);
assert_eq!(download.metadata.range, Some(requested));
assert_eq!(download.metadata.size_bytes, requested.size_bytes());
assert_eq!(
download.metadata.size_bytes,
i64::try_from(requested.size_bytes()).expect("test range should fit i64")
);
assert_eq!(download.metadata.total_size_bytes, payload_size(&payload));
assert_eq!(download.metadata.mime_type, "application/pdf");
let chunks = ok(
@ -3447,6 +3530,24 @@ async fn s3_storage_streams_range_from_backend_proxy_path() {
.flat_map(std::iter::IntoIterator::into_iter)
.collect::<Vec<_>>();
assert_eq!(collected, payload[4..12]);
let suffix = ok(
storage
.get_object(GetFileObject {
encoded_object_key,
read_token,
range: Some(FileRangeRequest::Suffix { length: 100 }),
})
.await,
"S3 suffix range should read",
);
assert_eq!(suffix.data, payload);
assert_eq!(
suffix.range,
Some(FileByteRange {
start: 0,
end_inclusive: 15,
})
);
}
#[tokio::test]
@ -3454,8 +3555,7 @@ async fn s3_storage_multipart_completion_uses_part_manifest_digest() {
let (_postgres, pool) = synctv_core_testing::create_test_pool().await;
let repository = Arc::new(FileStorageRepository::new(pool));
let operator = ok(
opendal::Operator::new(opendal::services::Memory::default())
.map(opendal::OperatorBuilder::finish),
opendal::Operator::new(opendal::services::Memory::default()),
"memory operator should build",
);
let storage = ok(
@ -3583,8 +3683,7 @@ async fn s3_storage_store_upload_accepts_server_mediated_parts() {
let (_postgres, pool) = synctv_core_testing::create_test_pool().await;
let repository = Arc::new(FileStorageRepository::new(pool));
let operator = ok(
opendal::Operator::new(opendal::services::Memory::default())
.map(opendal::OperatorBuilder::finish),
opendal::Operator::new(opendal::services::Memory::default()),
"memory operator should build",
);
let storage = ok(
@ -3704,8 +3803,7 @@ async fn s3_storage_rejects_part_outside_declared_manifest() {
let (_postgres, pool) = synctv_core_testing::create_test_pool().await;
let repository = Arc::new(FileStorageRepository::new(pool));
let operator = ok(
opendal::Operator::new(opendal::services::Memory::default())
.map(opendal::OperatorBuilder::finish),
opendal::Operator::new(opendal::services::Memory::default()),
"memory operator should build",
);
let storage = ok(
@ -3783,8 +3881,7 @@ async fn s3_storage_server_mediated_upload_is_bound_to_session_key() {
let (_postgres, pool) = synctv_core_testing::create_test_pool().await;
let repository = Arc::new(FileStorageRepository::new(pool));
let operator = ok(
opendal::Operator::new(opendal::services::Memory::default())
.map(opendal::OperatorBuilder::finish),
opendal::Operator::new(opendal::services::Memory::default()),
"memory operator should build",
);
let storage = ok(
@ -3915,8 +4012,7 @@ async fn s3_storage_multipart_completion_rejects_manifest_mismatch() {
let (_postgres, pool) = synctv_core_testing::create_test_pool().await;
let repository = Arc::new(FileStorageRepository::new(pool));
let operator = ok(
opendal::Operator::new(opendal::services::Memory::default())
.map(opendal::OperatorBuilder::finish),
opendal::Operator::new(opendal::services::Memory::default()),
"memory operator should build",
);
let storage = ok(
@ -4018,8 +4114,7 @@ async fn s3_storage_multipart_completion_rejects_manifest_mismatch() {
#[tokio::test]
async fn s3_public_constructor_requires_repository_for_upload_session() {
let operator = ok(
opendal::Operator::new(opendal::services::Memory::default())
.map(opendal::OperatorBuilder::finish),
opendal::Operator::new(opendal::services::Memory::default()),
"memory operator should build",
);
let storage = ok(
@ -4067,8 +4162,7 @@ async fn s3_multipart_completion_uses_all_recorded_parts() {
let (_postgres, pool) = synctv_core_testing::create_test_pool().await;
let repository = Arc::new(FileStorageRepository::new(pool.clone()));
let operator = ok(
opendal::Operator::new(opendal::services::Memory::default())
.map(opendal::OperatorBuilder::finish),
opendal::Operator::new(opendal::services::Memory::default()),
"memory operator should build",
);
let storage = ok(

@ -99,9 +99,7 @@ async fn cloudreve_v4_container_matches_login_user_and_list_contracts() -> anyho
.await?;
ensure_cloudreve_success(&create_folder, "folder creation")?;
let listing = client
.list(&token.access_token, "", 1, None, 20)
.await?;
let listing = client.list(&token.access_token, "", 1, None, 20).await?;
let folder = listing
.files
.iter()

@ -2904,15 +2904,15 @@ message FileUploadRange {
}
message FileByteRange {
int64 start = 1;
int64 end_inclusive = 2;
uint64 start = 1;
uint64 end_inclusive = 2;
}
message FileRangeRequest {
oneof range {
FileByteRange exact = 1;
int64 from_start = 2;
int64 suffix_length = 3;
uint64 from_start = 2;
uint64 suffix_length = 3;
}
}

@ -7,7 +7,7 @@ license = "MIT"
[features]
default = []
oss = ["dep:opendal"]
oss = ["dep:opendal", "dep:opendal-http-transport-reqwest"]
tls-aws-lc = ["synctv-common/tls-aws-lc"]
tls-ring = ["synctv-common/tls-ring"]
tls-webpki-roots = ["synctv-common/tls-webpki-roots"]
@ -69,6 +69,7 @@ backon.workspace = true
# Object storage (optional, behind `oss` feature)
opendal = { workspace = true, optional = true }
opendal-http-transport-reqwest = { workspace = true, optional = true }
[dev-dependencies]
tempfile.workspace = true

@ -70,7 +70,10 @@ mod inner {
}
// Build operator
let operator = Operator::new(builder)?.finish();
opendal::HttpTransporter::install_default(
opendal_http_transport_reqwest::ReqwestTransport::default(),
);
let operator = Operator::new(builder)?;
Ok(Self {
config,

Loading…
Cancel
Save