Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
16 changes: 16 additions & 0 deletions clients/python/src/objectstore_client/metadata.py
Original file line number Diff line number Diff line change
Expand Up @@ -14,6 +14,7 @@
HEADER_TIME_EXPIRES = "x-sn-time-expires"
HEADER_ORIGIN = "x-sn-origin"
HEADER_FILENAME = "x-sn-filename"
HEADER_SIZE = "x-sn-size"
HEADER_META_PREFIX = "x-snme-"


Expand Down Expand Up @@ -70,6 +71,17 @@ class Metadata:
prompting browsers and download tools to save the file under this name.
"""

size: int | None
"""
The size of the complete stored object in bytes.

This is always the size of the whole object, even when only a range of it was
requested. For a compressed object it is the compressed size, matching the bytes on
the wire.

This field is computed by the server, it cannot be set by clients.
"""

custom: dict[str, str]

@classmethod
Expand All @@ -81,6 +93,7 @@ def from_headers(cls, headers: Mapping[str, str]) -> Metadata:
time_expires = None
origin = None
filename = None
size = None
custom_metadata = {}

for k, v in headers.items():
Expand All @@ -98,6 +111,8 @@ def from_headers(cls, headers: Mapping[str, str]) -> Metadata:
origin = v
elif k == HEADER_FILENAME:
filename = v
elif k == HEADER_SIZE:
size = int(v)
elif k.startswith(HEADER_META_PREFIX):
custom_metadata[k[len(HEADER_META_PREFIX) :]] = v

Expand All @@ -109,6 +124,7 @@ def from_headers(cls, headers: Mapping[str, str]) -> Metadata:
time_expires=time_expires,
origin=origin,
filename=filename,
size=size,
custom=custom_metadata,
)

Expand Down
7 changes: 5 additions & 2 deletions clients/python/tests/test_e2e.py
Original file line number Diff line number Diff line change
Expand Up @@ -185,12 +185,15 @@ def test_head(server_url: str) -> None:

session = client.session(test_usecase, org=42, project=1337)

object_key = session.put(b"test data", origin="203.0.113.42")
payload = b"test data"
object_key = session.put(payload, origin="203.0.113.42", compression="none")

metadata = session.head(object_key)
assert metadata is not None
assert metadata.time_created is not None
assert metadata.origin == "203.0.113.42"
# `size` reports the stored bytes, which match the payload only without compression.
assert metadata.size == len(payload)

session.delete(object_key)

Expand Down Expand Up @@ -741,7 +744,7 @@ def test_presigned_head_succeeds(server_url: str) -> None:
)
status, _ = _fetch(url, method="HEAD")

assert status == 204
assert status == 200


def test_presigned_case_insensitive_method(server_url: str) -> None:
Expand Down
2 changes: 1 addition & 1 deletion objectstore-server/src/endpoints/batch.rs
Original file line number Diff line number Diff line change
Expand Up @@ -191,7 +191,7 @@ async fn convert_to_part(
idx,
&key,
"head",
StatusCode::NO_CONTENT,
StatusCode::OK,
None,
Bytes::new(),
Some(headers),
Comment thread
jan-auer marked this conversation as resolved.
Expand Down
43 changes: 30 additions & 13 deletions objectstore-server/src/endpoints/objects.rs
Original file line number Diff line number Diff line change
Expand Up @@ -88,25 +88,27 @@ async fn object_get(
};

let stream = state.meter_stream(stream, &context);
let metadata_headers = metadata.to_headers("").map_err(ServiceError::from)?;
let mut metadata_headers = metadata.to_headers("").map_err(ServiceError::from)?;

let mut response = match content_range {
Some(ref content_range) => {
let mut resp = (
metadata_headers.insert(
http::header::CONTENT_LENGTH,
content_range.len_to_header_value(),
);
metadata_headers.insert(http::header::CONTENT_RANGE, content_range.to_header_value());

(
StatusCode::PARTIAL_CONTENT,
metadata_headers,
Body::from_stream(stream),
)
.into_response();
let headers = resp.headers_mut();
headers.insert(
http::header::CONTENT_LENGTH,
content_range.len_to_header_value(),
);
headers.insert(http::header::CONTENT_RANGE, content_range.to_header_value());
resp
.into_response()
}
None => {
insert_content_length(&mut metadata_headers, &metadata);
(StatusCode::OK, metadata_headers, Body::from_stream(stream)).into_response()
}
None => (StatusCode::OK, metadata_headers, Body::from_stream(stream)).into_response(),
};

insert_content_disposition(&mut response, &metadata);
Expand All @@ -120,14 +122,29 @@ async fn object_head(service: AuthAwareService, Xt(id): Xt<ObjectId>) -> ApiResu
return Ok(StatusCode::NOT_FOUND.into_response());
};

let headers = metadata.to_headers("").map_err(ServiceError::from)?;
let mut headers = metadata.to_headers("").map_err(ServiceError::from)?;
insert_content_length(&mut headers, &metadata);

let mut response = (StatusCode::NO_CONTENT, headers).into_response();
let mut response = (StatusCode::OK, headers).into_response();
Comment thread
cursor[bot] marked this conversation as resolved.
Comment thread
cursor[bot] marked this conversation as resolved.
insert_content_disposition(&mut response, &metadata);
insert_accept_ranges(&mut response);
Ok(response)
}

/// Inserts a `Content-Length` header covering the complete object.
///
/// Only valid for responses whose body is the whole object, and for `HEAD` responses, which
/// describe what a `GET` would have returned. Ranged responses announce the length of the range
/// instead and must not use this.
fn insert_content_length(headers: &mut HeaderMap, metadata: &Metadata) {
if let Some(size) = metadata.size {
headers.insert(
http::header::CONTENT_LENGTH,
http::HeaderValue::from(size as u64),
);
}
}

fn insert_content_disposition(response: &mut Response, metadata: &Metadata) {
if let Some(val) = metadata
.filename
Expand Down
2 changes: 1 addition & 1 deletion objectstore-server/tests/objects.rs
Original file line number Diff line number Diff line change
Expand Up @@ -43,7 +43,7 @@ async fn filename_produces_content_disposition() -> Result<()> {
.head(server.url("/v1/objects/test/org=1/cd-key"))
.send()
.await?;
assert_eq!(resp.status(), reqwest::StatusCode::NO_CONTENT);
assert_eq!(resp.status(), reqwest::StatusCode::OK);
assert_eq!(resp.headers().get("x-sn-filename").unwrap(), "report.pdf");
assert_eq!(
resp.headers().get("content-disposition").unwrap(),
Expand Down
2 changes: 1 addition & 1 deletion objectstore-server/tests/presigned.rs
Original file line number Diff line number Diff line change
Expand Up @@ -144,7 +144,7 @@ async fn presigned_head_succeeds() -> Result<()> {
.send()
.await?;

assert_eq!(resp.status(), reqwest::StatusCode::NO_CONTENT);
assert_eq!(resp.status(), reqwest::StatusCode::OK);
Ok(())
}

Expand Down
68 changes: 67 additions & 1 deletion objectstore-server/tests/range_requests.rs
Original file line number Diff line number Diff line change
Expand Up @@ -66,14 +66,80 @@ async fn head_returns_accept_ranges() -> Result<()> {
.send()
.await?;

assert_eq!(resp.status(), reqwest::StatusCode::NO_CONTENT);
assert_eq!(resp.status(), reqwest::StatusCode::OK);
assert_eq!(
resp.headers().get("accept-ranges").unwrap().to_str()?,
"bytes"
);
Ok(())
}

/// `HEAD` announces the object size, so clients can learn it without a download.
#[tokio::test]
async fn head_reports_size() -> Result<()> {
let (server, key) = setup().await;
let client = reqwest::Client::new();

let resp = client
.head(server.url(&format!("/v1/objects/test/org=1/{key}")))
.send()
.await?;

assert_eq!(resp.status(), reqwest::StatusCode::OK);
assert_eq!(
resp.headers().get("content-length").unwrap().to_str()?,
"22"
);
assert_eq!(resp.headers().get("x-sn-size").unwrap().to_str()?, "22");
Ok(())
}

/// A full `GET` announces the same length as `HEAD`, rather than falling back to chunked.
#[tokio::test]
async fn full_get_reports_size() -> Result<()> {
let (server, key) = setup().await;
let client = reqwest::Client::new();

let resp = client
.get(server.url(&format!("/v1/objects/test/org=1/{key}")))
.send()
.await?;

assert_eq!(resp.status(), reqwest::StatusCode::OK);
assert_eq!(
resp.headers().get("content-length").unwrap().to_str()?,
"22"
);
assert_eq!(resp.headers().get("x-sn-size").unwrap().to_str()?, "22");

let body = resp.text().await?;
assert_eq!(body, "Hello, Range Requests!");
Ok(())
}

/// On a ranged response `Content-Length` describes the range, while `x-sn-size` stays the size of
/// the complete object.
#[tokio::test]
async fn range_content_length_covers_only_the_range() -> Result<()> {
let (server, key) = setup().await;
let client = reqwest::Client::new();

let resp = client
.get(server.url(&format!("/v1/objects/test/org=1/{key}")))
.header("range", "bytes=0-4")
.send()
.await?;

assert_eq!(resp.status(), reqwest::StatusCode::PARTIAL_CONTENT);
assert_eq!(resp.headers().get("content-length").unwrap().to_str()?, "5");
assert_eq!(
resp.headers().get("content-range").unwrap().to_str()?,
"bytes 0-4/22"
);
assert_eq!(resp.headers().get("x-sn-size").unwrap().to_str()?, "22");
Ok(())
}

#[tokio::test]
async fn range_prefix_returns_206() -> Result<()> {
let (server, key) = setup().await;
Expand Down
59 changes: 55 additions & 4 deletions objectstore-service/src/backend/s3_compatible.rs
Original file line number Diff line number Diff line change
Expand Up @@ -4,9 +4,9 @@ use std::time::SystemTime;
use std::{fmt, io};

use futures_util::{StreamExt, TryStreamExt};
use objectstore_types::metadata::Metadata;
use objectstore_types::metadata::{HEADER_SIZE, Metadata};
use objectstore_types::range::{ByteRange, ContentRange};
use reqwest::header::HeaderMap;
use reqwest::header::{HeaderMap, HeaderName};
use reqwest::{Body, IntoUrl, Method, RequestBuilder, Response, StatusCode};

use super::extensions::{ResponseExt, SendTraced};
Expand Down Expand Up @@ -125,6 +125,13 @@ fn metadata_to_gcs_headers(
prefix: &str,
) -> Result<HeaderMap, objectstore_types::metadata::Error> {
let mut headers = metadata.to_headers(prefix)?;

// The size is derived from the native `Content-Length` on every read, so it must not be
// persisted: metadata updates rewrite *all* stored metadata, and a stored `x-sn-size` key
// is rejected by the GCS JSON backend when it deserializes the object.
let size = HeaderName::try_from(format!("{prefix}{HEADER_SIZE}"))?;
headers.remove(&size);
Comment thread
lcian marked this conversation as resolved.

// GCS custom-time for lifecycle expiration
if let Some(expires_at) = metadata.time_expires {
let expires_at = humantime::format_rfc3339_seconds(expires_at);
Expand Down Expand Up @@ -217,8 +224,20 @@ where
metadata.size = Some(range.total as usize);
Some(range)
} else {
if let Some(len) = response.content_length() {
metadata.size = Some(len as usize);
// NB: Read the header rather than `Response::content_length`, which reports the
// length of the decoded body and is therefore always zero for a HEAD response.
let size = headers
.get(reqwest::header::CONTENT_LENGTH)
.and_then(|value| value.to_str().ok())
.map(|value| value.parse::<usize>())
.transpose()
.map_err(|cause| Error::Generic {
context: "S3: failed to parse Content-Length from object response".to_string(),
Comment thread
jan-auer marked this conversation as resolved.
cause: Some(Box::new(cause)),
})?;

if let Some(size) = size {
metadata.size = Some(size);
} else {
objectstore_log::warn!("S3: 200 response missing Content-Length header");
}
Expand Down Expand Up @@ -403,6 +422,19 @@ mod tests {
})
}

#[test]
fn metadata_to_gcs_headers_omits_size() {
let metadata = Metadata {
size: Some(4096),
..Default::default()
};

let headers = metadata_to_gcs_headers(&metadata, GCS_CUSTOM_PREFIX).unwrap();

// Persisting the size would store a key that the GCS JSON backend rejects on read.
assert!(headers.get("x-goog-meta-x-sn-size").is_none());
}

#[test]
fn metadata_to_gcs_headers_uses_time_expires() {
let expires = SystemTime::now() + Duration::from_hours(1);
Expand All @@ -429,6 +461,25 @@ mod tests {
Ok(())
}

#[tokio::test]
#[ignore = "MinIO does not support streaming bodies (requires Content-Length)"]
async fn test_get_metadata_reports_size() -> Result<()> {
let backend = create_test_backend();
let id = make_id();
let payload = "hello, world";

backend
.put_object(&id, &Metadata::default(), stream::single(payload))
.await?;

// The size must come from the `Content-Length` header, not from the (empty) body of
// the HEAD response.
let metadata = backend.get_metadata(&id).await?.expect("object exists");
assert_eq!(metadata.size, Some(payload.len()));

Ok(())
}

#[tokio::test]
#[ignore = "MinIO does not support streaming bodies (requires Content-Length)"]
async fn test_ttl_immediate() -> Result<()> {
Expand Down
Loading
Loading