Skip to content

Commit 0e21c59

Browse files
committed
improve
1 parent 302b1c4 commit 0e21c59

11 files changed

Lines changed: 96 additions & 152 deletions

File tree

objectstore-server/src/endpoints/objects.rs

Lines changed: 17 additions & 14 deletions
Original file line numberDiff line numberDiff line change
@@ -121,23 +121,26 @@ async fn object_get(
121121
let stream = state.meter_stream(stream, &context);
122122
let metadata_headers = metadata.to_headers("").map_err(ServiceError::from)?;
123123

124-
let is_partial = !content_range.is_full();
125-
let status = if is_partial {
126-
StatusCode::PARTIAL_CONTENT
127-
} else {
128-
StatusCode::OK
124+
let mut response = match content_range {
125+
Some(ref content_range) => {
126+
let mut resp = (
127+
StatusCode::PARTIAL_CONTENT,
128+
metadata_headers,
129+
Body::from_stream(stream),
130+
)
131+
.into_response();
132+
let headers = resp.headers_mut();
133+
headers.insert(
134+
http::header::CONTENT_LENGTH,
135+
content_range.len_to_header_value(),
136+
);
137+
headers.insert(http::header::CONTENT_RANGE, content_range.to_header_value());
138+
resp
139+
}
140+
None => (StatusCode::OK, metadata_headers, Body::from_stream(stream)).into_response(),
129141
};
130-
let mut response = (status, metadata_headers, Body::from_stream(stream)).into_response();
131142

132143
insert_accept_ranges(&mut response);
133-
if is_partial {
134-
let headers = response.headers_mut();
135-
headers.insert(
136-
http::header::CONTENT_LENGTH,
137-
content_range.len_to_header_value(),
138-
);
139-
headers.insert(http::header::CONTENT_RANGE, content_range.to_header_value());
140-
}
141144

142145
Ok(response)
143146
}

objectstore-server/tests/range_requests.rs

Lines changed: 0 additions & 23 deletions
Original file line numberDiff line numberDiff line change
@@ -233,26 +233,3 @@ async fn range_on_nonexistent_object_returns_404() -> Result<()> {
233233
assert_eq!(resp.status(), reqwest::StatusCode::NOT_FOUND);
234234
Ok(())
235235
}
236-
237-
#[tokio::test]
238-
async fn full_range_returns_200() -> Result<()> {
239-
let (server, key) = setup().await;
240-
let client = reqwest::Client::new();
241-
242-
// Request the full object as a range — should still get 200 since it's the full content.
243-
let resp = client
244-
.get(server.url(&format!("/v1/objects/test/org=1/{key}")))
245-
.header("range", "bytes=0-21")
246-
.send()
247-
.await?;
248-
249-
assert_eq!(resp.status(), reqwest::StatusCode::OK);
250-
assert!(
251-
resp.headers().get("content-range").is_none(),
252-
"full-object range should be 200 without Content-Range"
253-
);
254-
255-
let body = resp.text().await?;
256-
assert_eq!(body, "Hello, Range Requests!");
257-
Ok(())
258-
}

objectstore-service/src/backend/bigtable.rs

Lines changed: 2 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -38,7 +38,7 @@ use bigtable_rs::google::bigtable::v2::{self, mutation};
3838
use bytes::Bytes;
3939
use futures_util::TryStreamExt;
4040
use objectstore_types::metadata::{ExpirationPolicy, Metadata};
41-
use objectstore_types::range::{ByteRange, ContentRange};
41+
use objectstore_types::range::ByteRange;
4242
use serde::{Deserialize, Serialize};
4343
use tonic::Code;
4444

@@ -1015,8 +1015,7 @@ impl HighVolumeBackend for BigTableBackend {
10151015
RowData::Object { metadata, payload } => {
10161016
let mut metadata = metadata;
10171017
metadata.size = Some(payload.len());
1018-
let content_range = ContentRange::full(payload.len() as u64);
1019-
TieredGet::Object(metadata, content_range, crate::stream::single(payload))
1018+
TieredGet::Object(metadata, None, crate::stream::single(payload))
10201019
}
10211020
})
10221021
}

objectstore-service/src/backend/common.rs

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -24,7 +24,7 @@ pub const USER_AGENT: &str = concat!("sentry-objectstore/", env!("CARGO_PKG_VERS
2424
/// Backend response for put operations.
2525
pub type PutResponse = ();
2626
/// Backend response for get operations.
27-
pub type GetResponse = Option<(Metadata, ContentRange, PayloadStream)>;
27+
pub type GetResponse = Option<(Metadata, Option<ContentRange>, PayloadStream)>;
2828
/// Backend response for metadata-only get operations.
2929
pub type MetadataResponse = Option<Metadata>;
3030
/// Backend response for delete operations.
@@ -212,7 +212,7 @@ pub struct Tombstone {
212212
/// Typed response from [`HighVolumeBackend::get_tiered_object`].
213213
pub enum TieredGet {
214214
/// A real object was found.
215-
Object(Metadata, ContentRange, PayloadStream),
215+
Object(Metadata, Option<ContentRange>, PayloadStream),
216216
/// A redirect tombstone was found; the real object lives in the long-term backend.
217217
Tombstone(Tombstone),
218218
/// No entry exists at this key.

objectstore-service/src/backend/gcs.rs

Lines changed: 12 additions & 10 deletions
Original file line numberDiff line numberDiff line change
@@ -708,17 +708,19 @@ impl Backend for GcsBackend {
708708
.await?;
709709

710710
let content_range = if payload_response.status() == StatusCode::PARTIAL_CONTENT {
711-
payload_response
712-
.headers()
713-
.get(header::CONTENT_RANGE)
714-
.and_then(|v| v.to_str().ok())
715-
.and_then(|s| s.parse::<ContentRange>().ok())
716-
.ok_or_else(|| Error::Generic {
717-
context: "GCS: 206 response missing valid Content-Range header".to_owned(),
718-
cause: None,
719-
})?
711+
Some(
712+
payload_response
713+
.headers()
714+
.get(header::CONTENT_RANGE)
715+
.and_then(|v| v.to_str().ok())
716+
.and_then(|s| s.parse::<ContentRange>().ok())
717+
.ok_or_else(|| Error::Generic {
718+
context: "GCS: 206 response missing valid Content-Range header".to_owned(),
719+
cause: None,
720+
})?,
721+
)
720722
} else {
721-
ContentRange::full(metadata.size.unwrap_or(0) as u64)
723+
None
722724
};
723725

724726
let stream = payload_response

objectstore-service/src/backend/in_memory.rs

Lines changed: 21 additions & 21 deletions
Original file line numberDiff line numberDiff line change
@@ -9,7 +9,7 @@ use std::collections::{BTreeMap, HashMap};
99
use std::sync::{Arc, Mutex};
1010
use std::time::SystemTime;
1111

12-
use objectstore_types::range::{ByteRange, ContentRange};
12+
use objectstore_types::range::ByteRange;
1313

1414
use bytes::{Bytes, BytesMut};
1515
use futures_util::TryStreamExt;
@@ -132,16 +132,16 @@ impl super::common::Backend for InMemoryBackend {
132132
Some(StoreEntry::Object(mut metadata, bytes)) => {
133133
let total = bytes.len() as u64;
134134
metadata.size = Some(bytes.len());
135-
let content_range = match range {
136-
Some(range) => range
137-
.resolve(total)
138-
.ok_or(Error::RangeNotSatisfiable { total })?,
139-
None => ContentRange::full(total),
140-
};
141-
let payload = if content_range.is_full() {
142-
bytes
143-
} else {
144-
bytes.slice(content_range.start as usize..=content_range.end as usize)
135+
let (content_range, payload) = match range {
136+
Some(range) => {
137+
let content_range = range
138+
.resolve(total)
139+
.ok_or(Error::RangeNotSatisfiable { total })?;
140+
let sliced =
141+
bytes.slice(content_range.start as usize..=content_range.end as usize);
142+
(Some(content_range), sliced)
143+
}
144+
None => (None, bytes),
145145
};
146146
Ok(Some((
147147
metadata,
@@ -189,16 +189,16 @@ impl HighVolumeBackend for InMemoryBackend {
189189
Some(StoreEntry::Object(mut metadata, bytes)) => {
190190
let total = bytes.len() as u64;
191191
metadata.size = Some(bytes.len());
192-
let content_range = match range {
193-
Some(range) => range
194-
.resolve(total)
195-
.ok_or(Error::RangeNotSatisfiable { total })?,
196-
None => ContentRange::full(total),
197-
};
198-
let payload = if content_range.is_full() {
199-
bytes
200-
} else {
201-
bytes.slice(content_range.start as usize..=content_range.end as usize)
192+
let (content_range, payload) = match range {
193+
Some(range) => {
194+
let content_range = range
195+
.resolve(total)
196+
.ok_or(Error::RangeNotSatisfiable { total })?;
197+
let sliced =
198+
bytes.slice(content_range.start as usize..=content_range.end as usize);
199+
(Some(content_range), sliced)
200+
}
201+
None => (None, bytes),
202202
};
203203
TieredGet::Object(metadata, content_range, crate::stream::single(payload))
204204
}

objectstore-service/src/backend/local_fs.rs

Lines changed: 17 additions & 22 deletions
Original file line numberDiff line numberDiff line change
@@ -8,7 +8,7 @@ use std::time::SystemTime;
88

99
use futures_util::StreamExt;
1010
use objectstore_types::metadata::Metadata;
11-
use objectstore_types::range::{ByteRange, ContentRange};
11+
use objectstore_types::range::ByteRange;
1212
use tokio::fs::OpenOptions;
1313
use tokio::io::{AsyncBufReadExt, AsyncReadExt, AsyncSeekExt, AsyncWriteExt, BufReader, BufWriter};
1414
use tokio_util::io::{ReaderStream, StreamReader};
@@ -147,28 +147,20 @@ impl Backend for LocalFsBackend {
147147
.ok_or_else(|| Error::generic("local-fs file corrupted: shorter than header"))?;
148148
metadata.size = Some(payload_size as usize);
149149

150-
let content_range = match range {
151-
Some(byte_range) => match byte_range.resolve(payload_size) {
152-
Some(content_range) => content_range,
153-
None => {
154-
return Err(Error::RangeNotSatisfiable {
150+
let (content_range, stream) = match range {
151+
Some(byte_range) => {
152+
let content_range = byte_range
153+
.resolve(payload_size)
154+
.ok_or(Error::RangeNotSatisfiable {
155155
total: payload_size,
156-
});
157-
}
158-
},
159-
None => ContentRange::full(payload_size),
160-
};
161-
162-
let stream = if content_range.is_full() {
163-
let stream = ReaderStream::new(reader);
164-
stream.boxed()
165-
} else {
166-
reader
167-
.seek(std::io::SeekFrom::Current(content_range.start as i64))
168-
.await?;
169-
let limited = reader.take(content_range.len());
170-
let stream = ReaderStream::new(limited);
171-
stream.boxed()
156+
})?;
157+
reader
158+
.seek(std::io::SeekFrom::Current(content_range.start as i64))
159+
.await?;
160+
let limited = reader.take(content_range.len());
161+
(Some(content_range), ReaderStream::new(limited).boxed())
162+
}
163+
None => (None, ReaderStream::new(reader).boxed()),
172164
};
173165
Ok(Some((metadata, content_range, stream)))
174166
}
@@ -777,6 +769,7 @@ mod tests {
777769
let data: BytesMut = body.try_collect().await.unwrap();
778770

779771
assert_eq!(data.as_ref(), b"range");
772+
let content_range = content_range.unwrap();
780773
assert_eq!(content_range.start, 7);
781774
assert_eq!(content_range.end, 11);
782775
assert_eq!(content_range.total, payload.len() as u64);
@@ -803,6 +796,7 @@ mod tests {
803796
let data: BytesMut = body.try_collect().await.unwrap();
804797

805798
assert_eq!(data.as_ref(), b"range requests!");
799+
let content_range = content_range.unwrap();
806800
assert_eq!(content_range.start, 7);
807801
assert_eq!(content_range.end, 21);
808802
assert_eq!(content_range.total, payload.len() as u64);
@@ -829,6 +823,7 @@ mod tests {
829823
let data: BytesMut = body.try_collect().await.unwrap();
830824

831825
assert_eq!(data.as_ref(), b"requests!");
826+
let content_range = content_range.unwrap();
832827
assert_eq!(content_range.start, 13);
833828
assert_eq!(content_range.end, 21);
834829
assert_eq!(content_range.total, payload.len() as u64);

objectstore-service/src/backend/s3_compatible.rs

Lines changed: 5 additions & 11 deletions
Original file line numberDiff line numberDiff line change
@@ -165,7 +165,7 @@ where
165165
method: Method,
166166
id: &ObjectId,
167167
range: Option<ByteRange>,
168-
) -> Result<Option<(Metadata, ContentRange, reqwest::Response)>> {
168+
) -> Result<Option<(Metadata, Option<ContentRange>, reqwest::Response)>> {
169169
let object_url = self.object_url(id);
170170

171171
let mut builder = self.request(method, &object_url).await?;
@@ -219,18 +219,12 @@ where
219219
cause: None,
220220
})?;
221221
metadata.size = Some(range.total as usize);
222-
range
222+
Some(range)
223223
} else {
224-
match response.content_length() {
225-
Some(len) => {
226-
metadata.size = Some(len as usize);
227-
ContentRange::full(len)
228-
}
229-
None => {
230-
objectstore_log::warn!("S3: 200 response missing Content-Length header");
231-
ContentRange::full(metadata.size.unwrap_or(0) as u64)
232-
}
224+
if let Some(len) = response.content_length() {
225+
metadata.size = Some(len as usize);
233226
}
227+
None
234228
};
235229

236230
// TODO: Schedule into background persistently so this doesn't get lost on restarts

objectstore-service/src/backend/tiered.rs

Lines changed: 12 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -443,13 +443,18 @@ impl Backend for TieredStorage {
443443
backend_type = backend_type,
444444
);
445445

446-
if let Some((_, ref content_range, _)) = result {
447-
objectstore_metrics::record!(
448-
"get.size" = content_range.len(),
449-
usecase = id.usecase().to_owned(),
450-
backend_choice = backend_choice.as_str(),
451-
backend_type = backend_type,
452-
);
446+
if let Some((ref metadata, ref content_range, _)) = result {
447+
let size = content_range
448+
.map(|cr| cr.len() as usize)
449+
.or(metadata.size);
450+
if let Some(size) = size {
451+
objectstore_metrics::record!(
452+
"get.size" = size,
453+
usecase = id.usecase().to_owned(),
454+
backend_choice = backend_choice.as_str(),
455+
backend_type = backend_type,
456+
);
457+
}
453458
}
454459

455460
Ok(result)

objectstore-service/src/service.rs

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -23,7 +23,7 @@ use crate::stream::{ClientStream, PayloadStream};
2323
use crate::streaming::StreamExecutor;
2424

2525
/// Service response for [`StorageService::get_object`].
26-
pub type GetResponse = Option<(Metadata, ContentRange, PayloadStream)>;
26+
pub type GetResponse = Option<(Metadata, Option<ContentRange>, PayloadStream)>;
2727
/// Service response for [`StorageService::get_metadata`].
2828
pub type MetadataResponse = Option<Metadata>;
2929
/// Service response for [`StorageService::insert_object`].

0 commit comments

Comments
 (0)