This is an automated email from the ASF dual-hosted git repository.
Xuanwo pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/opendal.git
The following commit(s) were added to refs/heads/main by this push:
new c6c3b7b4b feat(core): return metadata from read (#7624)
c6c3b7b4b is described below
commit c6c3b7b4b5c681c7b2a931908939732ff1c2e05a
Author: Xuanwo <[email protected]>
AuthorDate: Wed May 27 21:56:24 2026 +0800
feat(core): return metadata from read (#7624)
* feat(core): return metadata from read
* refactor(core): set read metadata once
* docs(core): clarify read metadata semantics
* refactor(core): adjust read metadata by range
* refactor(core): use real metadata for chunked reads
* refactor(core): inline buffer stream state types
* fix(core): reuse chunked read for metadata
* refactor(core): move read metadata into cache
* refactor(core): ignore duplicate read metadata cache
* refactor(core): keep single read metadata cache
* refactor(core): remove content range from metadata
* refactor(core): cache stream metadata once
* refactor(core): centralize read metadata cache
* refactor(core): reuse chunked read for metadata
* refactor(core): wait chunked read metadata
* refactor(core): remove buffered stream metadata result
* fix(core): handle in-flight stream metadata
* refactor(core): model stream metadata opening
* refactor(core): simplify chunked read input
* test(core): cover stream read metadata
* test(core): inline stream metadata cases
* fix: Return object size in read metadata
* fix: Preserve object size for range read metadata
* feat(core): make read metadata optional
* fix(core): keep read metadata changes msrv compatible
* fix(object_store): fallback on range read metadata miss
* fix(object_store): document read metadata fallback
---
.../src/docs/rfcs/5871_read_returns_metadata.md | 158 +++-----
core/core/src/layers/correctness_check.rs | 5 +-
core/core/src/raw/http_util/header.rs | 4 +-
core/core/src/raw/rps.rs | 59 +--
core/core/src/services/memory/backend.rs | 15 +-
core/core/src/types/context/read.rs | 27 +-
core/core/src/types/metadata.rs | 30 +-
core/core/src/types/read/buffer_stream.rs | 136 ++++++-
core/core/src/types/read/futures_bytes_stream.rs | 9 +
core/core/src/types/read/reader.rs | 68 ++++
core/layers/concurrent-limit/src/lib.rs | 5 +-
core/layers/foyer/src/full.rs | 12 +-
core/layers/retry/src/lib.rs | 2 +-
core/layers/timeout/src/lib.rs | 5 +-
core/services/aliyun-drive/src/backend.rs | 7 +-
core/services/alluxio/src/backend.rs | 5 +-
core/services/azblob/src/backend.rs | 5 +-
core/services/azdls/src/backend.rs | 5 +-
core/services/azfile/src/backend.rs | 5 +-
core/services/b2/src/backend.rs | 7 +-
core/services/cacache/src/backend.rs | 4 +-
core/services/cloudflare-kv/src/backend.rs | 4 +-
core/services/compfs/src/backend.rs | 10 +-
core/services/cos/src/backend.rs | 7 +-
core/services/d1/src/backend.rs | 4 +-
core/services/dashmap/src/backend.rs | 4 +-
core/services/dropbox/src/backend.rs | 7 +-
core/services/etcd/src/backend.rs | 7 +-
core/services/foundationdb/src/backend.rs | 4 +-
core/services/foyer/src/backend.rs | 4 +-
core/services/fs/src/backend.rs | 11 +-
core/services/ftp/src/backend.rs | 3 +-
core/services/gcs/src/backend.rs | 7 +-
core/services/gdrive/src/backend.rs | 12 +-
core/services/ghac/src/backend.rs | 17 +-
core/services/github/src/backend.rs | 7 +-
core/services/goosefs/src/backend.rs | 10 +-
core/services/goosefs/src/reader.rs | 19 +-
core/services/gridfs/src/backend.rs | 4 +-
core/services/hdfs-native/src/backend.rs | 6 +-
core/services/hdfs/src/backend.rs | 2 +-
core/services/hf/src/backend.rs | 4 +-
core/services/hf/src/reader.rs | 23 +-
core/services/http/src/backend.rs | 7 +-
core/services/ipfs/src/backend.rs | 7 +-
core/services/koofr/src/backend.rs | 7 +-
core/services/lakefs/src/backend.rs | 7 +-
core/services/memcached/src/backend.rs | 4 +-
core/services/mini_moka/src/backend.rs | 7 +-
core/services/moka/src/backend.rs | 4 +-
core/services/mongodb/src/backend.rs | 4 +-
core/services/mysql/src/backend.rs | 4 +-
core/services/obs/src/backend.rs | 7 +-
core/services/onedrive/src/backend.rs | 7 +-
core/services/opfs/src/backend.rs | 3 +-
core/services/oss/src/backend.rs | 7 +-
core/services/pcloud/src/backend.rs | 7 +-
core/services/persy/src/backend.rs | 4 +-
core/services/postgresql/src/backend.rs | 4 +-
core/services/redb/src/backend.rs | 4 +-
core/services/redis/src/backend.rs | 13 +-
core/services/rocksdb/src/backend.rs | 4 +-
core/services/s3/src/backend.rs | 7 +-
core/services/s3/src/error.rs | 9 +
core/services/seafile/src/backend.rs | 7 +-
core/services/sled/src/backend.rs | 4 +-
core/services/sqlite/src/backend.rs | 12 +-
core/services/sqlite/src/core.rs | 9 +-
core/services/surrealdb/src/backend.rs | 4 +-
core/services/swift/src/backend.rs | 5 +-
core/services/tikv/src/backend.rs | 4 +-
core/services/tos/src/backend.rs | 7 +-
core/services/upyun/src/backend.rs | 7 +-
core/services/vercel-artifacts/src/backend.rs | 7 +-
core/services/vercel-blob/src/backend.rs | 7 +-
core/services/webdav/src/backend.rs | 7 +-
core/services/yandex-disk/src/backend.rs | 5 +-
core/tests/behavior/async_read.rs | 174 +++++++++
integrations/object_store/src/service/reader.rs | 21 +-
integrations/object_store/src/store.rs | 403 +++++++++++++++------
80 files changed, 1077 insertions(+), 482 deletions(-)
diff --git a/core/core/src/docs/rfcs/5871_read_returns_metadata.md
b/core/core/src/docs/rfcs/5871_read_returns_metadata.md
index f854a76e4..fbef79ff7 100644
--- a/core/core/src/docs/rfcs/5871_read_returns_metadata.md
+++ b/core/core/src/docs/rfcs/5871_read_returns_metadata.md
@@ -5,9 +5,9 @@
# Summary
-Expose the metadata returned by read operations without requiring users or
adapters to issue a separate `stat` request.
+Expose metadata returned natively by read operations without requiring users
or adapters to issue a separate `stat` request.
-Read metadata is bound to a concrete read stream. Calling `metadata().await`
on the stream opens the underlying read request, stores the returned metadata
and reader, and lets subsequent stream consumption reuse the same reader.
+Read metadata is bound to a concrete read stream. Calling `metadata().await`
on the stream opens the underlying read request, stores the returned reader,
and caches metadata if the service returned it. Subsequent stream consumption
reuses the same reader.
`Reader` also maintains a best-effort cache of complete object metadata that
has been observed from successful read opens.
@@ -15,7 +15,9 @@ Read metadata is bound to a concrete read stream. Calling
`metadata().await` on
Today, users who need metadata while reading must issue an additional `stat()`
request. This is inefficient for services that already return useful metadata
from their read APIs, such as S3 `GetObject`, GCS object download, and Azure
Blob download. It can also observe a different object version if the object
changes between read and stat.
-This is especially important for adapters such as `object_store_opendal`.
`object_store::GetResult` must contain `ObjectMeta`, returned range,
attributes, and a streaming payload at the time `get_opts()` returns. A
stat-first implementation can satisfy that contract, but it adds an extra
request before every ranged read. OpenDAL should expose read response metadata
through its public API so these adapters can construct their result without
issuing a redundant `stat`.
+However, not every service can return metadata as part of read open. Forcing
those services to call `stat`, `size`, or `metadata` internally would move the
extra request into OpenDAL's core read path and make normal reads more
expensive.
+
+This is especially important for adapters such as `object_store_opendal`.
`object_store::GetResult` must contain `ObjectMeta`, returned range,
attributes, and a streaming payload at the time `get_opts()` returns. A
stat-first implementation can satisfy that contract, but it adds an extra
request before every ranged read. OpenDAL should expose read-open object
metadata when it is available, so these adapters can avoid redundant `stat` on
metadata-capable services while preserving the exis [...]
# Guide-level explanation
@@ -44,7 +46,9 @@ if let Some(etag) = meta.etag() {
let data = stream.try_collect::<Vec<_>>().await?;
```
-`metadata().await` opens the read request if it has not been opened yet. The
stream caches both the metadata and the underlying reader, so consuming the
stream afterwards does not open a second read request.
+`metadata().await` opens the read request if it has not been opened yet. If
the underlying service returns metadata while opening the read, the stream
caches both the metadata and the underlying reader, so consuming the stream
afterwards does not open a second read request.
+
+If the service does not return metadata while opening the read,
`metadata().await` returns `ErrorKind::Unsupported`. The opened reader is still
reused by later stream consumption.
If users do not call `metadata().await`, streams keep their existing lazy
behavior: the underlying read request is opened when the stream is first polled.
@@ -59,7 +63,7 @@ if let Some(meta) = reader.metadata() {
}
```
-`Reader::metadata()` does not perform I/O. It returns `None` until a read
opened by this reader has observed enough information to build complete object
metadata.
+`Reader::metadata()` does not perform I/O. It returns `None` until a read
opened by this reader has observed enough information to build complete object
metadata. It also returns `None` for services that do not return metadata while
opening read operations.
# Reference-level explanation
@@ -79,17 +83,13 @@ Semantics:
- It never performs I/O.
- It initially returns `None`.
-- It returns object metadata, not read response metadata.
+- It returns object metadata.
- Its `content_length` is always the full object size.
-- Its `content_range` is always `None`.
-- It is updated only when a read response can be normalized into complete
object metadata.
-- If several reads are opened by the same reader, the cache stores the latest
complete object metadata observed by those reads.
-
-For a full read response, the object size is `read_metadata.content_length()`.
-
-For a ranged read response, the object size is
`read_metadata.content_range().and_then(|v| v.size())`. If the ranged response
does not provide the full object size, `Reader::metadata()` must not be updated
from that response.
+- It is updated from metadata returned by successful read opens.
+- It remains `None` if the service does not return metadata while opening read
operations.
+- If several reads are opened by the same reader, the cache stores the first
object metadata observed by those reads.
-This API is an ergonomic cache for users who read through `Reader::read(...)`
and then want the complete object metadata observed along the way. It is not
the API adapters should use when metadata must be tied to a concrete streaming
payload.
+This API is an ergonomic cache for users who read through `Reader::read(...)`
and then want the complete object metadata observed along the way.
## Stream API
@@ -107,88 +107,55 @@ impl FuturesBytesStream {
Semantics:
-- The returned metadata describes this concrete stream and its requested range.
+- The returned metadata describes the object being read.
+- Its `content_length` is always the full object size, even for ranged reads.
- The first `metadata().await` call prepares the stream by opening the
underlying read request.
-- The stream stores the returned `raw::RpRead::metadata` and the returned
`oio::Reader`.
+- The stream stores the returned `oio::Reader` and caches
`raw::RpRead::metadata` if present.
- Later `metadata().await` calls return the cached metadata.
- Later `poll_next` calls consume the cached reader and must not open the same
read request again.
- `metadata().await` does not read the response body. It only opens the read
and observes response metadata.
-
-The stream metadata describes the read response. For ranged reads, its
`content_length` is the payload length and its `content_range` carries the
range and full object size when available.
+- If `raw::RpRead::metadata` is absent, `metadata().await` returns
`ErrorKind::Unsupported`.
## Changes to `raw::RpRead`
-`raw::RpRead` becomes the metadata carrier for a successful read open:
+`raw::RpRead` becomes an optional metadata carrier for a successful read open:
```rust,ignore
pub struct RpRead {
- metadata: Metadata,
+ metadata: Option<Metadata>,
}
impl RpRead {
pub fn new(metadata: Metadata) -> Self;
- pub fn metadata(&self) -> &Metadata;
+ pub fn metadata(&self) -> Option<&Metadata>;
- pub fn into_metadata(self) -> Metadata;
+ pub fn into_metadata(self) -> Option<Metadata>;
}
```
-`RpRead::size()` and `RpRead::range()` should be removed. Their meanings
overlap with `Metadata::content_length()` and `Metadata::content_range()`, and
keeping two parallel surfaces would make the range-read semantics ambiguous.
+Services should set `RpRead::metadata` only when the read operation natively
returns enough information to build complete object metadata. They must not
issue an extra request only to fill `RpRead::metadata`.
-Callers that need the payload size should use:
+`RpRead::size()` and `RpRead::range()` should be removed. Metadata should
expose object metadata only; returned ranges are derived from the read request
and the object size.
-```rust,ignore
-rp.metadata().content_length()
-```
-
-Callers that need the read range or full object size for a ranged read should
use:
+Callers that need the object size should use:
```rust,ignore
-let range = rp.metadata().content_range();
-let object_size = range.and_then(|v| v.size());
+rp.metadata().map(|m| m.content_length())
```
## `content_length` Semantics
-`Metadata::content_length` should describe the body length of the entity
represented by this metadata.
-
-For stat and list metadata, the represented entity is the object:
+`Metadata::content_length` always describes the full object size.
```text
content_length = full object size
-content_range = None
-```
-
-For read metadata, the represented entity is the read response:
-
-```text
-content_length = read response payload length
-content_range = read response range within the full object
```
-For example, reading `1024..2048` from a 10 MiB object should produce:
+For example, reading `1024..2048` from a 10 MiB object should still produce:
```text
-content_length = 1024
-content_range = bytes 1024-2047/10485760
-```
-
-The full object size for a ranged read is available from
`content_range.size()`, not from `content_length`.
-
-When updating `Reader::metadata()` from read metadata, OpenDAL must normalize
the read response metadata into object metadata:
-
-```text
-full read metadata:
- reader.content_length = read.content_length
- reader.content_range = None
-
-range read metadata with total object size:
- reader.content_length = read.content_range.size
- reader.content_range = None
-
-range read metadata without total object size:
- reader metadata is not updated
+content_length = 10485760
```
HTTP response mapping:
@@ -196,41 +163,32 @@ HTTP response mapping:
```text
HTTP 200 + Content-Length: N
-> content_length = N
- -> content_range = None
HTTP 206 + Content-Range: bytes start-end/total
- -> content_length = end - start + 1
- -> content_range = start..end/total
+ -> content_length = total
HTTP 206 + Content-Range: bytes start-end/*
- -> content_length = end - start + 1
- -> content_range = start..end/*
+ -> content_length cannot be derived from the response alone
HTTP 206 + Content-Length: M, no Content-Range
- -> content_length = M
- -> content_range = None
+ -> content_length cannot be derived from the response alone
```
-Adapters that need full object size, such as `object_store_opendal`, must read
it from `content_range.size()` for ranged reads and fall back to
`content_length()` for full reads.
+Backends must not expose HTTP `Content-Range` through public `Metadata`. They
may parse it internally to populate the full object size.
## Implementation Details
-Backends must return metadata for every successful read open.
-
For services that return metadata in read responses:
- Capture metadata from the read response.
-- Populate `content_length` according to the read response payload length.
-- Populate `content_range` when the response provides a valid content range.
+- Populate `content_length` with the full object size. HTTP services may
derive it from `Content-Length` for full responses or from `Content-Range` for
ranged responses.
- Populate fields such as `content_type`, `etag`, `version`, `last_modified`,
and user metadata when available from the response.
For services whose read primitive does not naturally return metadata:
-- The backend must still satisfy the `RpRead { metadata }` contract.
-- File-like services may use file metadata, seek, or an equivalent
backend-local primitive during read open.
-- Remote file-like services such as SFTP, HDFS, and FTP may need an additional
status or size request to satisfy the contract.
-
-This cost is part of making read metadata a mandatory read-open contract
instead of an optional hint.
+- Return `RpRead::default()`.
+- Do not call `stat`, `size`, `metadata`, `seek`, or equivalent extra APIs
only to populate `RpRead::metadata`.
+- Keep any existing service-specific requests that are required for read
correctness itself, such as resolving an open-ended range that the underlying
API cannot express.
## Stream State
@@ -240,19 +198,19 @@ This cost is part of making read metadata a mandatory
read-open contract instead
metadata().await
-> prepare/open if needed
-> accessor.read(path, args_with_range).await
- -> cache RpRead.metadata
- -> update Reader metadata cache if complete object metadata can be derived
+ -> cache RpRead.metadata if present
+ -> update Reader metadata cache if metadata is present
-> cache returned oio::Reader
- -> return cached metadata
+ -> return cached metadata, or Unsupported if metadata is absent
poll_next()
- -> prepare/open if needed
- -> consume cached oio::Reader
+ -> open if needed
+ -> consume cached or newly opened oio::Reader
```
For non-chunked streaming reads, one stream maps to one underlying
`accessor.read`.
-For chunked or concurrent reads, a single physical read response may only
describe one chunk. If metadata is exposed for a logical chunked stream, the
stream layer must synthesize logical stream metadata from the requested range
and the observed object size. Paths that need exact metadata before returning a
streaming payload can choose the non-chunked stream path.
+For chunked or concurrent reads, a single physical read response may only
return one chunk body, but its metadata still describes the full object.
## `object_store_opendal` Integration
@@ -270,21 +228,16 @@ let reader = self.inner.reader_with(&raw_location)
let mut stream = reader.into_bytes_stream(read_range.clone()).await?;
let meta = stream.metadata().await?;
-let object_size = meta
- .content_range()
- .and_then(|v| v.size())
- .unwrap_or_else(|| meta.content_length());
-
let result = GetResult {
payload: GetResultPayload::Stream(Box::pin(stream.map_err(...))),
meta: ObjectMeta {
location: location.clone(),
last_modified:
meta.last_modified().and_then(timestamp_to_datetime).unwrap_or_default(),
- size: object_size,
+ size: meta.content_length(),
e_tag: meta.etag().map(ToString::to_string),
version: meta.version().map(ToString::to_string),
},
- range: build_returned_range(&meta, read_range),
+ range: build_returned_range(read_range, meta.content_length()),
attributes: build_attributes(&meta),
};
```
@@ -293,24 +246,25 @@ Cases that still need stat-first behavior:
- `head = true`, because the caller explicitly asks for no body.
- suffix ranges, because OpenDAL's public range input does not currently
express suffix ranges.
+- read streams whose services do not return metadata while opening the read
operation.
- responses that cannot provide the full object size required by
`ObjectMeta::size`.
- compatibility paths that intentionally preserve existing out-of-range
behavior.
# Drawbacks
-- `RpRead` becomes a stronger backend contract: successful read open must
return metadata.
-- Some backends will need extra work during read open to obtain size or status.
+- `metadata().await` is best-effort across services instead of universally
available.
+- Adapters that require metadata still need a fallback path for services
without read-open metadata.
- Stream implementations become more complex because they need a reusable
prepare/open state.
-- `Metadata::content_length` becomes operation-scoped: it describes object
size for stat/list metadata and response payload length for read metadata. This
must be documented clearly.
-- There are two public metadata views: stream metadata describes a read
response, while reader metadata describes the complete object metadata cache.
# Rationale and alternatives
-This design keeps existing `Reader::read()` and `Operator::read()` return
types unchanged. Users can opt into exact read response metadata by working
with read streams, or inspect the complete object metadata already observed by
`Reader`.
+This design keeps existing `Reader::read()` and `Operator::read()` return
types unchanged. Users can opt into read-open metadata by working with read
streams, or inspect the complete object metadata already observed by `Reader`.
+
+Using `Reader::metadata()` as the only metadata API is not enough: it is not
tied to a concrete read open and cannot satisfy APIs that must return metadata
before handing out a streaming payload. Another alternative is returning a new
`ReadResult` carrier from read APIs, but that would be a broader public API
change and would duplicate the existing stream abstraction.
-Using `Reader::metadata()` as the only metadata API is not enough: it is not
tied to a concrete range and cannot satisfy APIs that must return metadata
before handing out a streaming payload. Another alternative is returning a new
`ReadResult` carrier from read APIs, but that would be a broader public API
change and would duplicate the existing stream abstraction.
+Making `RpRead::metadata` mandatory was rejected because it forces services
without native read metadata to add extra status calls to the read path. That
violates OpenDAL's zero-extra-cost read design.
-Removing `RpRead::size()` and `RpRead::range()` avoids a second metadata
surface. Payload length is represented by `metadata.content_length()`, and
range information is represented by `metadata.content_range()`.
+Removing `RpRead::size()` and `RpRead::range()` avoids a second metadata
surface. Object size is represented by `metadata.content_length()`, and
returned range information is derived from the read request.
# Prior art
@@ -323,9 +277,9 @@ Similar patterns exist in other storage SDKs:
# Unresolved questions
- Whether OpenDAL should add a public suffix-range input in the future so
suffix reads can also use the no-stat path.
-- How much logical metadata synthesis should be supported for chunked and
concurrent streams.
+- Whether OpenDAL should expose a capability bit for services that can return
read-open metadata.
# Future possibilities
-- `object_store_opendal` can use stream metadata to avoid stat-first bounded
range reads.
-- `ReadContext::parse_into_range` can use read metadata or user-provided
content length hints to avoid additional stat calls for open-ended ranges.
+- `object_store_opendal` can use stream metadata to avoid stat-first bounded
range reads when services support read-open metadata.
+- `ReadContext::parse_into_range` can use user-provided content length hints
to avoid additional stat calls for open-ended ranges.
diff --git a/core/core/src/layers/correctness_check.rs
b/core/core/src/layers/correctness_check.rs
index 020acf407..e2af2ddda 100644
--- a/core/core/src/layers/correctness_check.rs
+++ b/core/core/src/layers/correctness_check.rs
@@ -325,7 +325,10 @@ mod tests {
}
async fn read(&self, _: &str, _: OpRead) -> Result<(RpRead,
Self::Reader)> {
- Ok((RpRead::new(), Box::new(bytes::Bytes::new())))
+ Ok((
+
RpRead::new(Metadata::new(EntryMode::FILE).with_content_length(0)),
+ Box::new(bytes::Bytes::new()),
+ ))
}
async fn write(&self, _: &str, _: OpWrite) -> Result<(RpWrite,
Self::Writer)> {
diff --git a/core/core/src/raw/http_util/header.rs
b/core/core/src/raw/http_util/header.rs
index 32553ab06..ce2e29b8b 100644
--- a/core/core/src/raw/http_util/header.rs
+++ b/core/core/src/raw/http_util/header.rs
@@ -175,8 +175,8 @@ pub fn parse_into_metadata(path: &str, headers: &HeaderMap)
-> Result<Metadata>
m.set_content_encoding(v);
}
- if let Some(v) = parse_content_range(headers)? {
- m.set_content_range(v);
+ if let Some(v) = parse_content_range(headers)?.and_then(|v| v.size()) {
+ m.set_content_length(v);
}
if let Some(v) = parse_etag(headers)? {
diff --git a/core/core/src/raw/rps.rs b/core/core/src/raw/rps.rs
index 560c88899..3462113df 100644
--- a/core/core/src/raw/rps.rs
+++ b/core/core/src/raw/rps.rs
@@ -17,7 +17,6 @@
use http::Request;
-use crate::raw::*;
use crate::*;
/// Reply for `create_dir` operation
@@ -98,58 +97,32 @@ impl<T: Default> From<PresignedRequest> for Request<T> {
}
/// Reply for `read` operation.
+///
+/// `RpRead` can carry metadata observed while opening this read operation.
+/// Services should set it only when the metadata is returned natively by the
+/// read operation. In particular, `metadata.content_length()` is the full
+/// object size, even if this read only returns a range of the object.
#[derive(Debug, Clone, Default)]
pub struct RpRead {
- /// Size is the size of the reader returned by this read operation.
- ///
- /// - `Some(size)` means the reader has at most size bytes.
- /// - `None` means the reader has unknown size.
- ///
- /// It's ok to leave size as empty, but it's recommended to set size if
possible. We will use
- /// this size as hint to do some optimization like avoid an extra stat or
read.
- size: Option<u64>,
- /// Range is the range of the reader returned by this read operation.
- ///
- /// - `Some(range)` means the reader's content range inside the whole file.
- /// - `None` means the reader's content range is unknown.
- ///
- /// It's ok to leave range as empty, but it's recommended to set range if
possible. We will use
- /// this range as hint to do some optimization like avoid an extra stat or
read.
- range: Option<BytesContentRange>,
+ metadata: Option<Metadata>,
}
impl RpRead {
/// Create a new reply for `read`.
- pub fn new() -> Self {
- RpRead::default()
- }
-
- /// Got the size of the reader returned by this read operation.
- ///
- /// - `Some(size)` means the reader has at most size bytes.
- /// - `None` means the reader has unknown size.
- pub fn size(&self) -> Option<u64> {
- self.size
- }
-
- /// Set the size of the reader returned by this read operation.
- pub fn with_size(mut self, size: Option<u64>) -> Self {
- self.size = size;
- self
+ pub fn new(metadata: Metadata) -> Self {
+ Self {
+ metadata: Some(metadata),
+ }
}
- /// Got the range of the reader returned by this read operation.
- ///
- /// - `Some(range)` means the reader has content range inside the whole
file.
- /// - `None` means the reader has unknown size.
- pub fn range(&self) -> Option<BytesContentRange> {
- self.range
+ /// Get metadata returned by this read operation.
+ pub fn metadata(&self) -> Option<&Metadata> {
+ self.metadata.as_ref()
}
- /// Set the range of the reader returned by this read operation.
- pub fn with_range(mut self, range: Option<BytesContentRange>) -> Self {
- self.range = range;
- self
+ /// Consume RpRead to get the inner metadata.
+ pub fn into_metadata(self) -> Option<Metadata> {
+ self.metadata
}
}
diff --git a/core/core/src/services/memory/backend.rs
b/core/core/src/services/memory/backend.rs
index 9eda17fc9..a122b45b5 100644
--- a/core/core/src/services/memory/backend.rs
+++ b/core/core/src/services/memory/backend.rs
@@ -139,10 +139,17 @@ impl Access for MemoryBackend {
}
};
- Ok((
- RpRead::new(),
- value.content.slice(args.range().to_range_as_usize()),
- ))
+ let total_size = value.content.len() as u64;
+ let range = args.range();
+ let start = range.offset().min(total_size) as usize;
+ let end = match range.size() {
+ Some(size) => range.offset().saturating_add(size).min(total_size),
+ None => total_size,
+ } as usize;
+ let content = value.content.slice(start..end);
+ let metadata =
Metadata::new(EntryMode::FILE).with_content_length(total_size);
+
+ Ok((RpRead::new(metadata), content))
}
async fn write(&self, path: &str, args: OpWrite) -> Result<(RpWrite,
Self::Writer)> {
diff --git a/core/core/src/types/context/read.rs
b/core/core/src/types/context/read.rs
index b16abddb1..04fbc0f73 100644
--- a/core/core/src/types/context/read.rs
+++ b/core/core/src/types/context/read.rs
@@ -19,6 +19,7 @@ use std::ops::Bound;
use std::ops::Range;
use std::ops::RangeBounds;
use std::sync::Arc;
+use std::sync::OnceLock;
use crate::raw::*;
use crate::*;
@@ -33,6 +34,8 @@ pub struct ReadContext {
args: OpRead,
/// Options for the reader.
options: OpReader,
+ /// Complete object metadata observed from successful read opens.
+ metadata: OnceLock<Metadata>,
}
impl ReadContext {
@@ -44,6 +47,7 @@ impl ReadContext {
path,
args,
options,
+ metadata: OnceLock::new(),
}
}
@@ -71,6 +75,17 @@ impl ReadContext {
&self.options
}
+ /// Get complete object metadata observed by this reader.
+ #[inline]
+ pub fn metadata(&self) -> Option<&Metadata> {
+ self.metadata.get()
+ }
+
+ /// Set cached object metadata observed from successful read opens once.
+ pub(crate) fn set_metadata(&self, metadata: Metadata) {
+ let _ = self.metadata.set(metadata);
+ }
+
/// Parse the range bounds into a range.
pub(crate) async fn parse_into_range(
&self,
@@ -179,9 +194,19 @@ impl ReadGenerator {
};
let args = self.ctx.args.clone().with_range(range);
- let (_, r) = self.ctx.acc.read(&self.ctx.path, args).await?;
+ let (rp, r) = self.ctx.acc.read(&self.ctx.path, args).await?;
+ if let Some(metadata) = rp.into_metadata() {
+ if self.ctx.metadata().is_none() {
+ self.ctx.set_metadata(metadata);
+ }
+ }
Ok(Some(r))
}
+
+ /// Get metadata observed by generated readers.
+ pub(crate) fn metadata(&self) -> Option<&Metadata> {
+ self.ctx.metadata()
+ }
}
#[cfg(test)]
diff --git a/core/core/src/types/metadata.rs b/core/core/src/types/metadata.rs
index 9125c3b47..6d3e99037 100644
--- a/core/core/src/types/metadata.rs
+++ b/core/core/src/types/metadata.rs
@@ -56,7 +56,6 @@ pub struct Metadata {
content_disposition: Option<String>,
content_length: Option<u64>,
content_md5: Option<String>,
- content_range: Option<BytesContentRange>,
content_type: Option<String>,
content_encoding: Option<String>,
etag: Option<String>,
@@ -89,9 +88,6 @@ impl fmt::Debug for Metadata {
if let Some(content_md5) = &self.content_md5 {
ds.field("content_md5", content_md5);
}
- if let Some(content_range) = self.content_range {
- ds.field("content_range", &content_range);
- }
if let Some(content_type) = &self.content_type {
ds.field("content_type", content_type);
}
@@ -129,7 +125,6 @@ impl Metadata {
content_md5: None,
content_type: None,
content_encoding: None,
- content_range: None,
last_modified: None,
etag: None,
content_disposition: None,
@@ -265,6 +260,10 @@ impl Metadata {
///
/// Refer to [MDN
Content-Length](https://developer.mozilla.org/en-US/docs/Web/HTTP/Headers/Content-Length)
for more information.
///
+ /// For file metadata returned by stat, list, or read operations, this
value
+ /// represents the full object size, even if the read operation only
returns
+ /// a range of the object.
+ ///
/// # Returns
///
/// Content length of this entry. It will be `0` if the content length is
not set by the storage services.
@@ -347,27 +346,6 @@ impl Metadata {
self
}
- /// Content Range of this entry.
- ///
- /// Content Range is defined by [RFC
9110](https://httpwg.org/specs/rfc9110.html#field.content-range).
- ///
- /// Refer to [MDN
Content-Range](https://developer.mozilla.org/en-US/docs/Web/HTTP/Headers/Content-Range)
for more information.
- pub fn content_range(&self) -> Option<BytesContentRange> {
- self.content_range
- }
-
- /// Set Content Range of this entry.
- pub fn set_content_range(&mut self, v: BytesContentRange) -> &mut Self {
- self.content_range = Some(v);
- self
- }
-
- /// Set Content Range of this entry.
- pub fn with_content_range(mut self, v: BytesContentRange) -> Self {
- self.content_range = Some(v);
- self
- }
-
/// Last modified of this entry.
///
/// `Last-Modified` is defined by [RFC
7232](https://httpwg.org/specs/rfc7232.html#header.last-modified)
diff --git a/core/core/src/types/read/buffer_stream.rs
b/core/core/src/types/read/buffer_stream.rs
index 971c33574..5c644848e 100644
--- a/core/core/src/types/read/buffer_stream.rs
+++ b/core/core/src/types/read/buffer_stream.rs
@@ -47,6 +47,27 @@ impl StreamingReader {
reader: None,
}
}
+
+ async fn prepare_metadata(&mut self) -> Result<()> {
+ if self.generator.metadata().is_some() {
+ return Ok(());
+ }
+
+ if self.reader.is_none() {
+ self.reader = self.generator.next_reader().await?;
+ }
+
+ Ok(())
+ }
+
+ async fn metadata(&mut self) -> Result<Metadata> {
+ self.prepare_metadata().await?;
+
+ self.generator
+ .metadata()
+ .cloned()
+ .ok_or_else(|| Error::new(ErrorKind::Unsupported, "read metadata
is not available"))
+ }
}
impl oio::Read for StreamingReader {
@@ -55,6 +76,7 @@ impl oio::Read for StreamingReader {
if self.reader.is_none() {
self.reader = self.generator.next_reader().await?;
}
+
let Some(r) = self.reader.as_mut() else {
return Ok(Buffer::new());
};
@@ -71,10 +93,10 @@ impl oio::Read for StreamingReader {
}
}
-/// Input for a chunked read task.
struct ChunkedReadInput {
ctx: Arc<ReadContext>,
range: BytesRange,
+ reader: Option<oio::Reader>,
}
/// ChunkedReader will read the file in chunks.
@@ -84,6 +106,7 @@ pub struct ChunkedReader {
ctx: Arc<ReadContext>,
offset: u64,
remaining: Option<u64>,
+ opened: Option<ChunkedReadInput>,
tasks: ConcurrentTasks<ChunkedReadInput, Buffer>,
done: bool,
}
@@ -99,11 +122,20 @@ impl ChunkedReader {
ctx.accessor().info().executor(),
ctx.options().concurrent(),
ctx.options().prefetch(),
- |input: ChunkedReadInput| {
+ |mut input: ChunkedReadInput| {
Box::pin(async move {
- let args =
input.ctx.args().clone().with_range(input.range);
let result = async {
- let (_, mut r) =
input.ctx.accessor().read(input.ctx.path(), args).await?;
+ if let Some(mut reader) = input.reader.take() {
+ return reader.read_all().await;
+ }
+
+ let args =
input.ctx.args().clone().with_range(input.range);
+ let (rp, mut r) =
input.ctx.accessor().read(input.ctx.path(), args).await?;
+ if let Some(metadata) = rp.into_metadata() {
+ if input.ctx.metadata().is_none() {
+ input.ctx.set_metadata(metadata);
+ }
+ }
r.read_all().await
}
.await;
@@ -115,11 +147,48 @@ impl ChunkedReader {
ctx,
offset: range.offset(),
remaining: range.size(),
+ opened: None,
tasks,
done: false,
}
}
+ async fn prepare_metadata(&mut self) -> Result<()> {
+ if self.ctx.metadata().is_some() {
+ return Ok(());
+ }
+
+ if self.opened.is_none() {
+ if let Some(range) = self.next_range() {
+ let args = self.ctx.args().clone().with_range(range);
+ let (rp, reader) = self.ctx.accessor().read(self.ctx.path(),
args).await?;
+ if let Some(metadata) = rp.into_metadata() {
+ if self.ctx.metadata().is_none() {
+ self.ctx.set_metadata(metadata);
+ }
+ }
+ self.opened = Some(ChunkedReadInput {
+ ctx: self.ctx.clone(),
+ range,
+ reader: Some(reader),
+ });
+ } else {
+ self.done = true;
+ }
+ }
+
+ Ok(())
+ }
+
+ async fn metadata(&mut self) -> Result<Metadata> {
+ self.prepare_metadata().await?;
+
+ self.ctx
+ .metadata()
+ .cloned()
+ .ok_or_else(|| Error::new(ErrorKind::Unsupported, "read metadata
is not available"))
+ }
+
/// Generate the next range to read, advancing internal state.
fn next_range(&mut self) -> Option<BytesRange> {
if self.remaining == Some(0) {
@@ -151,11 +220,14 @@ impl ChunkedReader {
impl oio::Read for ChunkedReader {
async fn read(&mut self) -> Result<Buffer> {
while self.tasks.has_remaining() && !self.done {
- if let Some(range) = self.next_range() {
+ if let Some(input) = self.opened.take() {
+ self.tasks.execute(input).await?;
+ } else if let Some(range) = self.next_range() {
self.tasks
.execute(ChunkedReadInput {
ctx: self.ctx.clone(),
range,
+ reader: None,
})
.await?;
} else {
@@ -166,7 +238,12 @@ impl oio::Read for ChunkedReader {
break;
}
}
- Ok(self.tasks.next().await.transpose()?.unwrap_or_default())
+
+ let Some(buffer) = self.tasks.next().await.transpose()? else {
+ return Ok(Buffer::new());
+ };
+
+ Ok(buffer)
}
}
@@ -174,6 +251,7 @@ impl oio::Read for ChunkedReader {
///
/// `BufferStream` implements `Stream` trait.
pub struct BufferStream {
+ ctx: Arc<ReadContext>,
/// # Notes to maintainers
///
/// The underlying reader is either a StreamingReader or a ChunkedReader.
@@ -184,6 +262,7 @@ pub struct BufferStream {
state: State,
}
+#[allow(clippy::large_enum_variant)]
enum State {
Idle(Option<TwoWays<StreamingReader, ChunkedReader>>),
Reading(BoxedStaticFuture<(TwoWays<StreamingReader, ChunkedReader>,
Result<Buffer>)>),
@@ -198,12 +277,19 @@ impl BufferStream {
);
let reader = if ctx.options().chunk().is_some() {
- TwoWays::Two(ChunkedReader::new(ctx, BytesRange::new(offset,
size)))
+ TwoWays::Two(ChunkedReader::new(
+ ctx.clone(),
+ BytesRange::new(offset, size),
+ ))
} else {
- TwoWays::One(StreamingReader::new(ctx, BytesRange::new(offset,
size)))
+ TwoWays::One(StreamingReader::new(
+ ctx.clone(),
+ BytesRange::new(offset, size),
+ ))
};
Self {
+ ctx,
state: State::Idle(Some(reader)),
}
}
@@ -218,15 +304,45 @@ impl BufferStream {
) -> Result<Self> {
let reader = if ctx.options().chunk().is_some() {
let range = ctx.parse_into_range(range).await?;
- TwoWays::Two(ChunkedReader::new(ctx, range.into()))
+ TwoWays::Two(ChunkedReader::new(ctx.clone(), range.into()))
} else {
- TwoWays::One(StreamingReader::new(ctx, range.into()))
+ TwoWays::One(StreamingReader::new(ctx.clone(), range.into()))
};
Ok(Self {
+ ctx,
state: State::Idle(Some(reader)),
})
}
+
+ /// Get metadata for this stream.
+ ///
+ /// Calling this method opens the underlying read request if needed.
+ /// Returns [`ErrorKind::Unsupported`] if the underlying service doesn't
+ /// return metadata while opening the read operation.
+ pub async fn metadata(&mut self) -> Result<Metadata> {
+ if let Some(metadata) = self.ctx.metadata() {
+ return Ok(metadata.clone());
+ }
+
+ match std::mem::replace(&mut self.state, State::Idle(None)) {
+ State::Idle(reader) => {
+ let mut reader = reader.expect("reader must exist while idle");
+ let prepared = match &mut reader {
+ TwoWays::One(v) => v.metadata().await,
+ TwoWays::Two(v) => v.metadata().await,
+ };
+ self.state = State::Idle(Some(reader));
+ prepared
+ }
+ State::Reading(fut) => {
+ self.state = State::Reading(fut);
+ self.ctx.metadata().cloned().ok_or_else(|| {
+ Error::new(ErrorKind::Unsupported, "read metadata is not
available")
+ })
+ }
+ }
+ }
}
impl Stream for BufferStream {
diff --git a/core/core/src/types/read/futures_bytes_stream.rs
b/core/core/src/types/read/futures_bytes_stream.rs
index 1a4e7ac11..8839abd8b 100644
--- a/core/core/src/types/read/futures_bytes_stream.rs
+++ b/core/core/src/types/read/futures_bytes_stream.rs
@@ -54,6 +54,15 @@ impl FuturesBytesStream {
buf: Buffer::new(),
})
}
+
+ /// Get metadata for this stream.
+ ///
+ /// Calling this method opens the underlying read request if needed.
+ /// Returns [`ErrorKind::Unsupported`] if the underlying service doesn't
+ /// return metadata while opening the read operation.
+ pub async fn metadata(&mut self) -> Result<Metadata> {
+ self.stream.metadata().await
+ }
}
impl Stream for FuturesBytesStream {
diff --git a/core/core/src/types/read/reader.rs
b/core/core/src/types/read/reader.rs
index f54e604a1..8dd14f840 100644
--- a/core/core/src/types/read/reader.rs
+++ b/core/core/src/types/read/reader.rs
@@ -105,6 +105,15 @@ impl Reader {
Reader { ctx: Arc::new(ctx) }
}
+ /// Get complete object metadata observed by this reader.
+ ///
+ /// This method doesn't perform I/O. It returns `None` if no read has
+ /// observed complete object metadata yet, or if the underlying service
+ /// doesn't return metadata while opening read operations.
+ pub fn metadata(&self) -> Option<&Metadata> {
+ self.ctx.metadata()
+ }
+
/// Read give range from reader into [`Buffer`].
///
/// This operation is zero-copy, which means it keeps the [`bytes::Bytes`]
returned by underlying
@@ -457,6 +466,65 @@ mod tests {
Ok(())
}
+ #[tokio::test]
+ async fn test_reader_metadata_after_read() -> Result<()> {
+ let op = Operator::via_iter(services::MEMORY_SCHEME, [])?;
+ op.write("test", Buffer::from("HelloWorld")).await?;
+
+ let reader = op.reader("test").await?;
+ assert_eq!(reader.metadata(), None);
+
+ let buf = reader.read(4..8).await?;
+ assert_eq!(&buf.to_vec(), b"oWor");
+
+ let meta = reader.metadata().expect("metadata must be observed");
+ assert_eq!(meta.content_length(), 10);
+
+ Ok(())
+ }
+
+ #[tokio::test]
+ async fn test_stream_metadata_updates_reader_metadata() -> Result<()> {
+ let op = Operator::via_iter(services::MEMORY_SCHEME, [])?;
+ op.write("test", Buffer::from("HelloWorld")).await?;
+
+ let reader = op.reader("test").await?;
+ let mut stream = reader.clone().into_stream(4..8).await?;
+
+ let stream_meta = stream.metadata().await?;
+ assert_eq!(stream_meta.content_length(), 10);
+
+ let bufs: Vec<_> = stream.try_collect().await?;
+ let buf: Buffer = bufs.into_iter().flatten().collect();
+ assert_eq!(&buf.to_vec(), b"oWor");
+
+ let reader_meta = reader.metadata().expect("reader metadata must be
observed");
+ assert_eq!(reader_meta.content_length(), 10);
+
+ Ok(())
+ }
+
+ #[tokio::test]
+ async fn test_chunked_stream_metadata_updates_reader_metadata() ->
Result<()> {
+ let op = Operator::via_iter(services::MEMORY_SCHEME, [])?;
+ op.write("test", Buffer::from("HelloWorld")).await?;
+
+ let reader = op.reader_with("test").chunk(2).await?;
+ let mut stream = reader.clone().into_stream(4..8).await?;
+
+ let stream_meta = stream.metadata().await?;
+ assert_eq!(stream_meta.content_length(), 10);
+
+ let bufs: Vec<_> = stream.try_collect().await?;
+ let buf: Buffer = bufs.into_iter().flatten().collect();
+ assert_eq!(&buf.to_vec(), b"oWor");
+
+ let reader_meta = reader.metadata().expect("reader metadata must be
observed");
+ assert_eq!(reader_meta.content_length(), 10);
+
+ Ok(())
+ }
+
fn gen_random_bytes() -> Vec<u8> {
let mut rng = rand::rng();
// Generate size between 1B..16MB.
diff --git a/core/layers/concurrent-limit/src/lib.rs
b/core/layers/concurrent-limit/src/lib.rs
index 00707a0cd..bda3f5f50 100644
--- a/core/layers/concurrent-limit/src/lib.rs
+++ b/core/layers/concurrent-limit/src/lib.rs
@@ -613,7 +613,10 @@ mod tests {
};
let req =
http::Request::get("http://fake").body(data).unwrap();
let resp = self.info.http_client().fetch(req).await?;
- Ok((RpRead::default(), resp.into_body()))
+ Ok((
+
RpRead::new(Metadata::new(EntryMode::FILE).with_content_length(0)),
+ resp.into_body(),
+ ))
}
async fn stat(&self, _: &str, _: OpStat) -> Result<RpStat> {
diff --git a/core/layers/foyer/src/full.rs b/core/layers/foyer/src/full.rs
index 8922f1a3f..12378609d 100644
--- a/core/layers/foyer/src/full.rs
+++ b/core/layers/foyer/src/full.rs
@@ -18,10 +18,11 @@
use std::sync::Arc;
use opendal_core::Buffer;
+use opendal_core::EntryMode;
use opendal_core::Error;
+use opendal_core::Metadata;
use opendal_core::Result;
use opendal_core::raw::Access;
-use opendal_core::raw::BytesContentRange;
use opendal_core::raw::BytesRange;
use opendal_core::raw::OpRead;
use opendal_core::raw::OpStat;
@@ -124,13 +125,10 @@ impl<A: Access> FullReader<A> {
match result {
Ok(entry) => {
let end = range_end.unwrap_or(entry.len() as u64);
- let range = BytesContentRange::default()
- .with_range(range_start, end - 1)
- .with_size(entry.len() as _);
let buffer = entry.slice(range_start as usize..end as usize);
- let rp = RpRead::new()
- .with_size(Some(buffer.len() as _))
- .with_range(Some(range));
+ let rp = RpRead::new(
+
Metadata::new(EntryMode::FILE).with_content_length(entry.len() as _),
+ );
Ok((rp, buffer))
}
Err(e) => match e.downcast_ref::<FetchError>() {
diff --git a/core/layers/retry/src/lib.rs b/core/layers/retry/src/lib.rs
index eb4a4fc28..1b2537862 100644
--- a/core/layers/retry/src/lib.rs
+++ b/core/layers/retry/src/lib.rs
@@ -919,7 +919,7 @@ mod tests {
async fn read(&self, _: &str, args: OpRead) -> Result<(RpRead,
Self::Reader)> {
Ok((
- RpRead::new(),
+
RpRead::new(Metadata::new(EntryMode::FILE).with_content_length(0)),
MockReader {
buf: Bytes::from("Hello, World!").into(),
range: args.range(),
diff --git a/core/layers/timeout/src/lib.rs b/core/layers/timeout/src/lib.rs
index b942e8245..774548c74 100644
--- a/core/layers/timeout/src/lib.rs
+++ b/core/layers/timeout/src/lib.rs
@@ -420,7 +420,10 @@ mod tests {
/// This function will build a reader that always return pending.
async fn read(&self, _: &str, _: OpRead) -> Result<(RpRead,
Self::Reader)> {
- Ok((RpRead::new(), Box::new(MockReader)))
+ Ok((
+
RpRead::new(Metadata::new(EntryMode::FILE).with_content_length(0)),
+ Box::new(MockReader),
+ ))
}
/// This function will never return.
diff --git a/core/services/aliyun-drive/src/backend.rs
b/core/services/aliyun-drive/src/backend.rs
index 854a01afa..2d7aa746c 100644
--- a/core/services/aliyun-drive/src/backend.rs
+++ b/core/services/aliyun-drive/src/backend.rs
@@ -327,9 +327,10 @@ impl Access for AliyunDriveBackend {
let status = resp.status();
match status {
- StatusCode::OK | StatusCode::PARTIAL_CONTENT => {
- Ok((RpRead::default(), resp.into_body()))
- }
+ StatusCode::OK | StatusCode::PARTIAL_CONTENT => Ok((
+ RpRead::new(parse_into_metadata(path, resp.headers())?),
+ resp.into_body(),
+ )),
_ => {
let (part, mut body) = resp.into_parts();
let buf = body.to_buffer().await?;
diff --git a/core/services/alluxio/src/backend.rs
b/core/services/alluxio/src/backend.rs
index e5e16a799..b68a8611c 100644
--- a/core/services/alluxio/src/backend.rs
+++ b/core/services/alluxio/src/backend.rs
@@ -165,7 +165,10 @@ impl Access for AlluxioBackend {
let buf = body.to_buffer().await?;
return Err(parse_error(Response::from_parts(part, buf)));
}
- Ok((RpRead::new(), resp.into_body()))
+ Ok((
+ RpRead::new(parse_into_metadata(path, resp.headers())?),
+ resp.into_body(),
+ ))
}
async fn write(&self, path: &str, args: OpWrite) -> Result<(RpWrite,
Self::Writer)> {
diff --git a/core/services/azblob/src/backend.rs
b/core/services/azblob/src/backend.rs
index db1e765f6..1b76f5dfc 100644
--- a/core/services/azblob/src/backend.rs
+++ b/core/services/azblob/src/backend.rs
@@ -504,7 +504,10 @@ impl Access for AzblobBackend {
let status = resp.status();
match status {
- StatusCode::OK | StatusCode::PARTIAL_CONTENT => Ok((RpRead::new(),
resp.into_body())),
+ StatusCode::OK | StatusCode::PARTIAL_CONTENT => Ok((
+ RpRead::new(parse_into_metadata(path, resp.headers())?),
+ resp.into_body(),
+ )),
_ => {
let (part, mut body) = resp.into_parts();
let buf = body.to_buffer().await?;
diff --git a/core/services/azdls/src/backend.rs
b/core/services/azdls/src/backend.rs
index a835b350c..f29fe6be8 100644
--- a/core/services/azdls/src/backend.rs
+++ b/core/services/azdls/src/backend.rs
@@ -407,7 +407,10 @@ impl Access for AzdlsBackend {
let status = resp.status();
match status {
- StatusCode::OK | StatusCode::PARTIAL_CONTENT => Ok((RpRead::new(),
resp.into_body())),
+ StatusCode::OK | StatusCode::PARTIAL_CONTENT => Ok((
+ RpRead::new(parse_into_metadata(path, resp.headers())?),
+ resp.into_body(),
+ )),
_ => {
let (part, mut body) = resp.into_parts();
let buf = body.to_buffer().await?;
diff --git a/core/services/azfile/src/backend.rs
b/core/services/azfile/src/backend.rs
index 06412c1d3..be7d3b43e 100644
--- a/core/services/azfile/src/backend.rs
+++ b/core/services/azfile/src/backend.rs
@@ -337,7 +337,10 @@ impl Access for AzfileBackend {
let status = resp.status();
match status {
- StatusCode::OK | StatusCode::PARTIAL_CONTENT => Ok((RpRead::new(),
resp.into_body())),
+ StatusCode::OK | StatusCode::PARTIAL_CONTENT => Ok((
+ RpRead::new(parse_into_metadata(path, resp.headers())?),
+ resp.into_body(),
+ )),
_ => {
let (part, mut body) = resp.into_parts();
let buf = body.to_buffer().await?;
diff --git a/core/services/b2/src/backend.rs b/core/services/b2/src/backend.rs
index 76fb9eac6..3ea5f203a 100644
--- a/core/services/b2/src/backend.rs
+++ b/core/services/b2/src/backend.rs
@@ -257,9 +257,10 @@ impl Access for B2Backend {
let status = resp.status();
match status {
- StatusCode::OK | StatusCode::PARTIAL_CONTENT => {
- Ok((RpRead::default(), resp.into_body()))
- }
+ StatusCode::OK | StatusCode::PARTIAL_CONTENT => Ok((
+ RpRead::new(parse_into_metadata(path, resp.headers())?),
+ resp.into_body(),
+ )),
_ => {
let (part, mut body) = resp.into_parts();
let buf = body.to_buffer().await?;
diff --git a/core/services/cacache/src/backend.rs
b/core/services/cacache/src/backend.rs
index 2b3b5d880..12003fe10 100644
--- a/core/services/cacache/src/backend.rs
+++ b/core/services/cacache/src/backend.rs
@@ -117,6 +117,7 @@ impl Access for CacacheBackend {
match data {
Some(bytes) => {
let range = args.range();
+ let content_length = bytes.len() as u64;
let buffer = if range.is_full() {
Buffer::from(bytes)
} else {
@@ -127,7 +128,8 @@ impl Access for CacacheBackend {
};
Buffer::from(bytes.slice(start..end.min(bytes.len())))
};
- Ok((RpRead::new(), buffer))
+ let metadata =
Metadata::new(EntryMode::FILE).with_content_length(content_length);
+ Ok((RpRead::new(metadata), buffer))
}
None => Err(Error::new(ErrorKind::NotFound, "entry not found")),
}
diff --git a/core/services/cloudflare-kv/src/backend.rs
b/core/services/cloudflare-kv/src/backend.rs
index c7af780e3..eae7e2cec 100644
--- a/core/services/cloudflare-kv/src/backend.rs
+++ b/core/services/cloudflare-kv/src/backend.rs
@@ -445,6 +445,7 @@ impl Access for CloudflareKvBackend {
}
let range = args.range();
+ let total_size = resp_body.len() as u64;
let buffer = if range.is_full() {
resp_body
} else {
@@ -455,7 +456,8 @@ impl Access for CloudflareKvBackend {
};
resp_body.slice(start..end.min(resp_body.len()))
};
- Ok((RpRead::new(), buffer))
+ let metadata =
Metadata::new(EntryMode::FILE).with_content_length(total_size);
+ Ok((RpRead::new(metadata), buffer))
}
async fn write(&self, path: &str, _: OpWrite) -> Result<(RpWrite,
Self::Writer)> {
diff --git a/core/services/compfs/src/backend.rs
b/core/services/compfs/src/backend.rs
index 8b23445a5..3a3d8ffe1 100644
--- a/core/services/compfs/src/backend.rs
+++ b/core/services/compfs/src/backend.rs
@@ -226,11 +226,17 @@ impl Access for CompfsBackend {
let file = self
.core
- .exec(|| async move {
compio::fs::OpenOptions::new().read(true).open(&path).await })
+ .exec(|| async move {
+ let file = compio::fs::OpenOptions::new()
+ .read(true)
+ .open(&path)
+ .await?;
+ Ok(file)
+ })
.await?;
let r = CompfsReader::new(self.core.clone(), file, op.range());
- Ok((RpRead::new(), r))
+ Ok((RpRead::default(), r))
}
async fn write(&self, path: &str, args: OpWrite) -> Result<(RpWrite,
Self::Writer)> {
diff --git a/core/services/cos/src/backend.rs b/core/services/cos/src/backend.rs
index 8c863500b..bc1e9b2ba 100644
--- a/core/services/cos/src/backend.rs
+++ b/core/services/cos/src/backend.rs
@@ -356,9 +356,10 @@ impl Access for CosBackend {
let status = resp.status();
match status {
- StatusCode::OK | StatusCode::PARTIAL_CONTENT => {
- Ok((RpRead::default(), resp.into_body()))
- }
+ StatusCode::OK | StatusCode::PARTIAL_CONTENT => Ok((
+ RpRead::new(parse_into_metadata(path, resp.headers())?),
+ resp.into_body(),
+ )),
_ => {
let (part, mut body) = resp.into_parts();
let buf = body.to_buffer().await?;
diff --git a/core/services/d1/src/backend.rs b/core/services/d1/src/backend.rs
index 1df7101e7..cfdbddeb1 100644
--- a/core/services/d1/src/backend.rs
+++ b/core/services/d1/src/backend.rs
@@ -257,7 +257,9 @@ impl Access for D1Backend {
return Err(Error::new(ErrorKind::NotFound, "kv not found in
d1"));
}
};
- Ok((RpRead::new(), bs.slice(args.range().to_range_as_usize())))
+ let content = bs.slice(args.range().to_range_as_usize());
+ let metadata =
Metadata::new(EntryMode::FILE).with_content_length(bs.len() as u64);
+ Ok((RpRead::new(metadata), content))
}
async fn write(&self, path: &str, _: OpWrite) -> Result<(RpWrite,
Self::Writer)> {
diff --git a/core/services/dashmap/src/backend.rs
b/core/services/dashmap/src/backend.rs
index b7c33b5da..c309c92f2 100644
--- a/core/services/dashmap/src/backend.rs
+++ b/core/services/dashmap/src/backend.rs
@@ -147,6 +147,7 @@ impl Access for DashmapBackend {
match self.core.get(&p)? {
Some(value) => {
+ let total_size = value.content.len() as u64;
let buffer = if args.range().is_full() {
value.content
} else {
@@ -158,7 +159,8 @@ impl Access for DashmapBackend {
};
value.content.slice(start..end.min(value.content.len()))
};
- Ok((RpRead::new(), buffer))
+ let metadata =
Metadata::new(EntryMode::FILE).with_content_length(total_size);
+ Ok((RpRead::new(metadata), buffer))
}
None => Err(Error::new(ErrorKind::NotFound, "key not found in
dashmap")),
}
diff --git a/core/services/dropbox/src/backend.rs
b/core/services/dropbox/src/backend.rs
index 1a05deed7..44c3e3310 100644
--- a/core/services/dropbox/src/backend.rs
+++ b/core/services/dropbox/src/backend.rs
@@ -110,9 +110,10 @@ impl Access for DropboxBackend {
let status = resp.status();
match status {
- StatusCode::OK | StatusCode::PARTIAL_CONTENT => {
- Ok((RpRead::default(), resp.into_body()))
- }
+ StatusCode::OK | StatusCode::PARTIAL_CONTENT => Ok((
+ RpRead::new(parse_into_metadata(path, resp.headers())?),
+ resp.into_body(),
+ )),
_ => {
let (part, mut body) = resp.into_parts();
let buf = body.to_buffer().await?;
diff --git a/core/services/etcd/src/backend.rs
b/core/services/etcd/src/backend.rs
index cd1fea9b1..93e0a8c9d 100644
--- a/core/services/etcd/src/backend.rs
+++ b/core/services/etcd/src/backend.rs
@@ -268,10 +268,12 @@ impl Access for EtcdBackend {
match self.core.get(&abs_path).await? {
Some(buffer) => {
let range = op.range();
+ let total_size = buffer.len() as u64;
// If range is full, return the buffer directly
if range.is_full() {
- return Ok((RpRead::new(), buffer));
+ let metadata =
Metadata::new(EntryMode::FILE).with_content_length(total_size);
+ return Ok((RpRead::new(metadata), buffer));
}
// Handle range requests
@@ -286,8 +288,9 @@ impl Access for EtcdBackend {
let size = range.size().map(|s| s as usize);
let end = size.map_or(buffer.len(), |s| (offset +
s).min(buffer.len()));
let sliced_buffer = buffer.slice(offset..end);
+ let metadata =
Metadata::new(EntryMode::FILE).with_content_length(total_size);
- Ok((RpRead::new(), sliced_buffer))
+ Ok((RpRead::new(metadata), sliced_buffer))
}
None => Err(Error::new(ErrorKind::NotFound, "path not found")),
}
diff --git a/core/services/foundationdb/src/backend.rs
b/core/services/foundationdb/src/backend.rs
index 865f92985..33404587d 100644
--- a/core/services/foundationdb/src/backend.rs
+++ b/core/services/foundationdb/src/backend.rs
@@ -161,7 +161,9 @@ impl Access for FoundationdbBackend {
));
}
};
- Ok((RpRead::new(), bs.slice(args.range().to_range_as_usize())))
+ let content = bs.slice(args.range().to_range_as_usize());
+ let metadata =
Metadata::new(EntryMode::FILE).with_content_length(bs.len() as u64);
+ Ok((RpRead::new(metadata), content))
}
async fn write(&self, path: &str, _: OpWrite) -> Result<(RpWrite,
Self::Writer)> {
diff --git a/core/services/foyer/src/backend.rs
b/core/services/foyer/src/backend.rs
index 303a2054f..90fd9d9a1 100644
--- a/core/services/foyer/src/backend.rs
+++ b/core/services/foyer/src/backend.rs
@@ -272,6 +272,7 @@ impl Access for FoyerBackend {
Some(bs) => bs,
None => return Err(Error::new(ErrorKind::NotFound, "key not found
in foyer")),
};
+ let content_length = buffer.len() as u64;
let buffer = if args.range().is_full() {
buffer
@@ -285,7 +286,8 @@ impl Access for FoyerBackend {
buffer.slice(start..end.min(buffer.len()))
};
- Ok((RpRead::new(), buffer))
+ let metadata =
Metadata::new(EntryMode::FILE).with_content_length(content_length);
+ Ok((RpRead::new(metadata), buffer))
}
async fn write(&self, path: &str, _: OpWrite) -> Result<(RpWrite,
Self::Writer)> {
diff --git a/core/services/fs/src/backend.rs b/core/services/fs/src/backend.rs
index fd01b1d57..8384e0134 100644
--- a/core/services/fs/src/backend.rs
+++ b/core/services/fs/src/backend.rs
@@ -204,15 +204,6 @@ impl Access for FsBackend {
Ok(RpStat::new(m))
}
- /// # Notes
- ///
- /// There are three ways to get the total file length:
- ///
- /// - call std::fs::metadata directly and then open. (400ns)
- /// - open file first, and then use `f.metadata()` (300ns)
- /// - open file first, and then use `seek`. (100ns)
- ///
- /// Benchmark could be found
[here](https://gist.github.com/Xuanwo/48f9cfbc3022ea5f865388bb62e1a70f)
async fn read(&self, path: &str, args: OpRead) -> Result<(RpRead,
Self::Reader)> {
let f = self.core.fs_read(path, &args).await?;
let r = FsReader::new(
@@ -220,7 +211,7 @@ impl Access for FsBackend {
f,
args.range().size().unwrap_or(u64::MAX) as _,
);
- Ok((RpRead::new(), r))
+ Ok((RpRead::default(), r))
}
async fn write(&self, path: &str, op: OpWrite) -> Result<(RpWrite,
Self::Writer)> {
diff --git a/core/services/ftp/src/backend.rs b/core/services/ftp/src/backend.rs
index 56458fffd..47798b58b 100644
--- a/core/services/ftp/src/backend.rs
+++ b/core/services/ftp/src/backend.rs
@@ -246,9 +246,8 @@ impl Access for FtpBackend {
async fn read(&self, path: &str, args: OpRead) -> Result<(RpRead,
Self::Reader)> {
let ftp_stream = self.core.ftp_connect(Operation::Read).await?;
-
let reader = FtpReader::new(ftp_stream, path.to_string(), args).await?;
- Ok((RpRead::new(), reader))
+ Ok((RpRead::default(), reader))
}
async fn write(&self, path: &str, op: OpWrite) -> Result<(RpWrite,
Self::Writer)> {
diff --git a/core/services/gcs/src/backend.rs b/core/services/gcs/src/backend.rs
index 54ab8ff2c..60502218a 100644
--- a/core/services/gcs/src/backend.rs
+++ b/core/services/gcs/src/backend.rs
@@ -458,9 +458,10 @@ impl Access for GcsBackend {
let status = resp.status();
match status {
- StatusCode::OK | StatusCode::PARTIAL_CONTENT => {
- Ok((RpRead::default(), resp.into_body()))
- }
+ StatusCode::OK | StatusCode::PARTIAL_CONTENT => Ok((
+ RpRead::new(parse_into_metadata(path, resp.headers())?),
+ resp.into_body(),
+ )),
_ => {
let (part, mut body) = resp.into_parts();
let buf = body.to_buffer().await?;
diff --git a/core/services/gdrive/src/backend.rs
b/core/services/gdrive/src/backend.rs
index b928a8c01..909087c59 100644
--- a/core/services/gdrive/src/backend.rs
+++ b/core/services/gdrive/src/backend.rs
@@ -148,15 +148,19 @@ impl Access for GdriveBackend {
let status = resp.status();
match status {
- StatusCode::OK | StatusCode::PARTIAL_CONTENT => Ok((RpRead::new(),
resp.into_body())),
+ StatusCode::OK | StatusCode::PARTIAL_CONTENT => Ok((
+ RpRead::new(parse_into_metadata(path, resp.headers())?),
+ resp.into_body(),
+ )),
StatusCode::NOT_FOUND => {
self.core.refresh_path(&abs_path).await;
let resp = self.core.gdrive_get(path, args.range()).await?;
let status = resp.status();
match status {
- StatusCode::OK | StatusCode::PARTIAL_CONTENT => {
- Ok((RpRead::new(), resp.into_body()))
- }
+ StatusCode::OK | StatusCode::PARTIAL_CONTENT => Ok((
+ RpRead::new(parse_into_metadata(path,
resp.headers())?),
+ resp.into_body(),
+ )),
_ => {
let (part, mut body) = resp.into_parts();
let buf = body.to_buffer().await?;
diff --git a/core/services/ghac/src/backend.rs
b/core/services/ghac/src/backend.rs
index fed800c4e..f41ea70e4 100644
--- a/core/services/ghac/src/backend.rs
+++ b/core/services/ghac/src/backend.rs
@@ -231,15 +231,7 @@ impl Access for GhacBackend {
let status = resp.status();
match status {
StatusCode::OK | StatusCode::PARTIAL_CONTENT |
StatusCode::RANGE_NOT_SATISFIABLE => {
- let mut meta = parse_into_metadata(path, resp.headers())?;
- // Correct content length via returning content range.
- meta.set_content_length(
- meta.content_range()
- .expect("content range must be valid")
- .size()
- .expect("content range must contains size"),
- );
-
+ let meta = parse_into_metadata(path, resp.headers())?;
Ok(RpStat::new(meta))
}
_ => Err(parse_error(resp)),
@@ -251,9 +243,10 @@ impl Access for GhacBackend {
let status = resp.status();
match status {
- StatusCode::OK | StatusCode::PARTIAL_CONTENT => {
- Ok((RpRead::default(), resp.into_body()))
- }
+ StatusCode::OK | StatusCode::PARTIAL_CONTENT => Ok((
+ RpRead::new(parse_into_metadata(path, resp.headers())?),
+ resp.into_body(),
+ )),
_ => {
let (part, mut body) = resp.into_parts();
let buf = body.to_buffer().await?;
diff --git a/core/services/github/src/backend.rs
b/core/services/github/src/backend.rs
index dff01aaa3..3b3920831 100644
--- a/core/services/github/src/backend.rs
+++ b/core/services/github/src/backend.rs
@@ -218,9 +218,10 @@ impl Access for GithubBackend {
let status = resp.status();
match status {
- StatusCode::OK | StatusCode::PARTIAL_CONTENT => {
- Ok((RpRead::default(), resp.into_body()))
- }
+ StatusCode::OK | StatusCode::PARTIAL_CONTENT => Ok((
+ RpRead::new(parse_into_metadata(path, resp.headers())?),
+ resp.into_body(),
+ )),
_ => {
let (part, mut body) = resp.into_parts();
let buf = body.to_buffer().await?;
diff --git a/core/services/goosefs/src/backend.rs
b/core/services/goosefs/src/backend.rs
index fa5323405..86ac1fc14 100644
--- a/core/services/goosefs/src/backend.rs
+++ b/core/services/goosefs/src/backend.rs
@@ -324,8 +324,14 @@ impl Access for GoosefsBackend {
}
async fn read(&self, path: &str, args: OpRead) -> Result<(RpRead,
Self::Reader)> {
- let reader = GoosefsReader::new(self.core.clone(), path.to_string(),
args);
- Ok((RpRead::new(), reader))
+ let content_length = if args.range().offset() != 0 &&
args.range().size().is_none() {
+ let file_info = self.core.get_status(path).await?;
+ Some(self.core.file_info_to_metadata(&file_info).content_length())
+ } else {
+ None
+ };
+ let reader = GoosefsReader::new(self.core.clone(), path.to_string(),
args, content_length);
+ Ok((RpRead::default(), reader))
}
async fn write(&self, path: &str, args: OpWrite) -> Result<(RpWrite,
Self::Writer)> {
diff --git a/core/services/goosefs/src/reader.rs
b/core/services/goosefs/src/reader.rs
index d3c1d75b5..22332833e 100644
--- a/core/services/goosefs/src/reader.rs
+++ b/core/services/goosefs/src/reader.rs
@@ -72,6 +72,7 @@ pub struct GoosefsReader {
core: Arc<GoosefsCore>,
path: String,
args: OpRead,
+ content_length: Option<u64>,
/// Lazily-opened SDK reader. `None` until the first `read()` call.
inner: Option<SdkReader>,
/// Terminal flag: once the underlying stream has returned `None`,
@@ -82,11 +83,17 @@ pub struct GoosefsReader {
}
impl GoosefsReader {
- pub fn new(core: Arc<GoosefsCore>, path: String, args: OpRead) -> Self {
+ pub fn new(
+ core: Arc<GoosefsCore>,
+ path: String,
+ args: OpRead,
+ content_length: Option<u64>,
+ ) -> Self {
GoosefsReader {
core,
path,
args,
+ content_length,
inner: None,
done: false,
}
@@ -110,9 +117,13 @@ impl GoosefsReader {
(0, None) => self.core.open_reader(&self.path).await,
(off, Some(len)) => self.core.open_range_reader(&self.path, off,
len).await,
(off, None) => {
- let info = self.core.get_status(&self.path).await?;
- let file_len = info.length.unwrap_or(0) as u64;
- let len = file_len.saturating_sub(off);
+ let content_length = self.content_length.ok_or_else(|| {
+ Error::new(
+ ErrorKind::Unexpected,
+ "content length must be known for offset reads",
+ )
+ })?;
+ let len = content_length.saturating_sub(off);
if len == 0 {
// Empty tail — short-circuit with a zero-length
// ranged open so the very next `read_next_block`
diff --git a/core/services/gridfs/src/backend.rs
b/core/services/gridfs/src/backend.rs
index cc31f0e0c..c827297a4 100644
--- a/core/services/gridfs/src/backend.rs
+++ b/core/services/gridfs/src/backend.rs
@@ -221,7 +221,9 @@ impl Access for GridfsBackend {
return Err(Error::new(ErrorKind::NotFound, "kv not found in
gridfs"));
}
};
- Ok((RpRead::new(), bs.slice(args.range().to_range_as_usize())))
+ let content = bs.slice(args.range().to_range_as_usize());
+ let metadata =
Metadata::new(EntryMode::FILE).with_content_length(bs.len() as u64);
+ Ok((RpRead::new(metadata), content))
}
async fn write(&self, path: &str, _: OpWrite) -> Result<(RpWrite,
Self::Writer)> {
diff --git a/core/services/hdfs-native/src/backend.rs
b/core/services/hdfs-native/src/backend.rs
index 1d32cd1d0..7306761c8 100644
--- a/core/services/hdfs-native/src/backend.rs
+++ b/core/services/hdfs-native/src/backend.rs
@@ -191,10 +191,14 @@ impl Access for HdfsNativeBackend {
async fn read(&self, path: &str, args: OpRead) -> Result<(RpRead,
Self::Reader)> {
let (f, offset, size) = self.core.hdfs_read(path, &args).await?;
+ let content_length = f.file_length() as u64;
let r = HdfsNativeReader::new(f, offset as _, size as _);
- Ok((RpRead::new(), r))
+ Ok((
+
RpRead::new(Metadata::new(EntryMode::FILE).with_content_length(content_length)),
+ r,
+ ))
}
async fn write(&self, path: &str, args: OpWrite) -> Result<(RpWrite,
Self::Writer)> {
diff --git a/core/services/hdfs/src/backend.rs
b/core/services/hdfs/src/backend.rs
index a5d111b00..84750a8fa 100644
--- a/core/services/hdfs/src/backend.rs
+++ b/core/services/hdfs/src/backend.rs
@@ -222,7 +222,7 @@ impl Access for HdfsBackend {
let f = self.core.hdfs_read(path, &args).await?;
Ok((
- RpRead::new(),
+ RpRead::default(),
HdfsReader::new(f, args.range().size().unwrap_or(u64::MAX) as _),
))
}
diff --git a/core/services/hf/src/backend.rs b/core/services/hf/src/backend.rs
index 6297b7164..2b7f860ba 100644
--- a/core/services/hf/src/backend.rs
+++ b/core/services/hf/src/backend.rs
@@ -268,8 +268,8 @@ impl Access for HfBackend {
}
async fn read(&self, path: &str, args: OpRead) -> Result<(RpRead,
Self::Reader)> {
- let reader = HfReader::try_new(&self.core, path, args.range()).await?;
- Ok((RpRead::default(), reader))
+ let (metadata, reader) = HfReader::try_new(&self.core, path,
args.range()).await?;
+ Ok((RpRead::new(metadata), reader))
}
async fn list(&self, path: &str, args: OpList) -> Result<(RpList,
Self::Lister)> {
diff --git a/core/services/hf/src/reader.rs b/core/services/hf/src/reader.rs
index ec74f26ee..0073d90e2 100644
--- a/core/services/hf/src/reader.rs
+++ b/core/services/hf/src/reader.rs
@@ -36,7 +36,7 @@ impl HfReader {
/// Buckets always use XET. For other repo types, a HEAD request
/// probes for the `X-Xet-Hash` header. Files stored on XET are
/// downloaded via the CAS protocol; all others fall back to HTTP GET.
- pub async fn try_new(core: &HfCore, path: &str, range: BytesRange) ->
Result<Self> {
+ pub async fn try_new(core: &HfCore, path: &str, range: BytesRange) ->
Result<(Metadata, Self)> {
if let Some(xet_file) = core.maybe_xet_file(path).await? {
return Self::try_new_xet(core, &xet_file, range).await;
}
@@ -51,7 +51,11 @@ impl HfReader {
Self::try_new_http(core, path, range).await
}
- pub async fn try_new_http(core: &HfCore, path: &str, range: BytesRange) ->
Result<Self> {
+ pub async fn try_new_http(
+ core: &HfCore,
+ path: &str,
+ range: BytesRange,
+ ) -> Result<(Metadata, Self)> {
let client = core.info.http_client();
let uri = core.uri(path);
let url = uri.resolve_url(&core.endpoint);
@@ -68,7 +72,10 @@ impl HfReader {
let status = resp.status();
match status {
- StatusCode::OK | StatusCode::PARTIAL_CONTENT =>
Ok(Self::Http(resp.into_body())),
+ StatusCode::OK | StatusCode::PARTIAL_CONTENT => Ok((
+ parse_into_metadata(path, resp.headers())?,
+ Self::Http(resp.into_body()),
+ )),
_ => {
let (part, mut body) = resp.into_parts();
let buf = body.to_buffer().await?;
@@ -81,7 +88,7 @@ impl HfReader {
core: &HfCore,
file_info: &XetFileInfo,
range: BytesRange,
- ) -> Result<Self> {
+ ) -> Result<(Metadata, Self)> {
let group = core.xet_download_group().await?;
let xet_range = if range.is_full() {
@@ -103,7 +110,13 @@ impl HfReader {
.set_source(err)
})?;
stream.start();
- Ok(Self::Xet(stream))
+
+ let total_size = file_info.file_size.unwrap_or_default();
+ let metadata = Metadata::new(EntryMode::FILE)
+ .with_content_length(total_size)
+ .with_etag(file_info.hash().to_string());
+
+ Ok((metadata, Self::Xet(stream)))
}
}
diff --git a/core/services/http/src/backend.rs
b/core/services/http/src/backend.rs
index 96cfe74bc..ae9a2091a 100644
--- a/core/services/http/src/backend.rs
+++ b/core/services/http/src/backend.rs
@@ -207,9 +207,10 @@ impl Access for HttpBackend {
let status = resp.status();
match status {
- StatusCode::OK | StatusCode::PARTIAL_CONTENT => {
- Ok((RpRead::default(), resp.into_body()))
- }
+ StatusCode::OK | StatusCode::PARTIAL_CONTENT => Ok((
+ RpRead::new(parse_into_metadata(path, resp.headers())?),
+ resp.into_body(),
+ )),
_ => {
let (part, mut body) = resp.into_parts();
let buf = body.to_buffer().await?;
diff --git a/core/services/ipfs/src/backend.rs
b/core/services/ipfs/src/backend.rs
index fb60ce3b5..4c65b8af6 100644
--- a/core/services/ipfs/src/backend.rs
+++ b/core/services/ipfs/src/backend.rs
@@ -164,9 +164,10 @@ impl Access for IpfsBackend {
let status = resp.status();
match status {
- StatusCode::OK | StatusCode::PARTIAL_CONTENT => {
- Ok((RpRead::default(), resp.into_body()))
- }
+ StatusCode::OK | StatusCode::PARTIAL_CONTENT => Ok((
+ RpRead::new(parse_into_metadata(path, resp.headers())?),
+ resp.into_body(),
+ )),
_ => {
let (part, mut body) = resp.into_parts();
let buf = body.to_buffer().await?;
diff --git a/core/services/koofr/src/backend.rs
b/core/services/koofr/src/backend.rs
index 075f2d539..ece316bbc 100644
--- a/core/services/koofr/src/backend.rs
+++ b/core/services/koofr/src/backend.rs
@@ -244,9 +244,10 @@ impl Access for KoofrBackend {
let status = resp.status();
match status {
- StatusCode::OK | StatusCode::PARTIAL_CONTENT => {
- Ok((RpRead::default(), resp.into_body()))
- }
+ StatusCode::OK | StatusCode::PARTIAL_CONTENT => Ok((
+ RpRead::new(parse_into_metadata(path, resp.headers())?),
+ resp.into_body(),
+ )),
_ => {
let (part, mut body) = resp.into_parts();
let buf = body.to_buffer().await?;
diff --git a/core/services/lakefs/src/backend.rs
b/core/services/lakefs/src/backend.rs
index 1b62ea4a6..176c398f5 100644
--- a/core/services/lakefs/src/backend.rs
+++ b/core/services/lakefs/src/backend.rs
@@ -234,9 +234,10 @@ impl Access for LakefsBackend {
let status = resp.status();
match status {
- StatusCode::OK | StatusCode::PARTIAL_CONTENT => {
- Ok((RpRead::default(), resp.into_body()))
- }
+ StatusCode::OK | StatusCode::PARTIAL_CONTENT => Ok((
+ RpRead::new(parse_into_metadata(path, resp.headers())?),
+ resp.into_body(),
+ )),
_ => {
let (part, mut body) = resp.into_parts();
let buf = body.to_buffer().await?;
diff --git a/core/services/memcached/src/backend.rs
b/core/services/memcached/src/backend.rs
index 9062267a8..004861352 100644
--- a/core/services/memcached/src/backend.rs
+++ b/core/services/memcached/src/backend.rs
@@ -248,7 +248,9 @@ impl Access for MemcachedBackend {
Some(bs) => bs,
None => return Err(Error::new(ErrorKind::NotFound, "kv not found
in memcached")),
};
- Ok((RpRead::new(), bs.slice(args.range().to_range_as_usize())))
+ let content = bs.slice(args.range().to_range_as_usize());
+ let metadata =
Metadata::new(EntryMode::FILE).with_content_length(bs.len() as u64);
+ Ok((RpRead::new(metadata), content))
}
async fn write(&self, path: &str, _: OpWrite) -> Result<(RpWrite,
Self::Writer)> {
diff --git a/core/services/mini_moka/src/backend.rs
b/core/services/mini_moka/src/backend.rs
index 1458f5439..139cfb0bc 100644
--- a/core/services/mini_moka/src/backend.rs
+++ b/core/services/mini_moka/src/backend.rs
@@ -195,10 +195,12 @@ impl Access for MiniMokaBackend {
match self.core.get(&p) {
Some(value) => {
let range = op.range();
+ let total_size = value.content.len() as u64;
// If range is full, return the content buffer directly
if range.is_full() {
- return Ok((RpRead::new(), value.content));
+ let metadata =
Metadata::new(EntryMode::FILE).with_content_length(total_size);
+ return Ok((RpRead::new(metadata), value.content));
}
let offset = range.offset() as usize;
@@ -214,8 +216,9 @@ impl Access for MiniMokaBackend {
(offset + s).min(value.content.len())
});
let sliced_content = value.content.slice(offset..end);
+ let metadata =
Metadata::new(EntryMode::FILE).with_content_length(total_size);
- Ok((RpRead::new(), sliced_content))
+ Ok((RpRead::new(metadata), sliced_content))
}
None => Err(Error::new(ErrorKind::NotFound, "path not found")),
}
diff --git a/core/services/moka/src/backend.rs
b/core/services/moka/src/backend.rs
index d54075a65..4ca8e44f2 100644
--- a/core/services/moka/src/backend.rs
+++ b/core/services/moka/src/backend.rs
@@ -278,6 +278,7 @@ impl Access for MokaBackend {
match self.core.get(&p).await? {
Some(value) => {
+ let total_size = value.content.len() as u64;
let buffer = if args.range().is_full() {
value.content
} else {
@@ -289,7 +290,8 @@ impl Access for MokaBackend {
};
value.content.slice(start..end.min(value.content.len()))
};
- Ok((RpRead::new(), buffer))
+ let metadata =
Metadata::new(EntryMode::FILE).with_content_length(total_size);
+ Ok((RpRead::new(metadata), buffer))
}
None => Err(Error::new(ErrorKind::NotFound, "key not found in
moka")),
}
diff --git a/core/services/mongodb/src/backend.rs
b/core/services/mongodb/src/backend.rs
index be3a3f93f..75d6037b8 100644
--- a/core/services/mongodb/src/backend.rs
+++ b/core/services/mongodb/src/backend.rs
@@ -239,7 +239,9 @@ impl Access for MongodbBackend {
return Err(Error::new(ErrorKind::NotFound, "kv not found in
mongodb"));
}
};
- Ok((RpRead::new(), bs.slice(args.range().to_range_as_usize())))
+ let content = bs.slice(args.range().to_range_as_usize());
+ let metadata =
Metadata::new(EntryMode::FILE).with_content_length(bs.len() as u64);
+ Ok((RpRead::new(metadata), content))
}
async fn write(&self, path: &str, _: OpWrite) -> Result<(RpWrite,
Self::Writer)> {
diff --git a/core/services/mysql/src/backend.rs
b/core/services/mysql/src/backend.rs
index 5ce0b4a6a..0746fa044 100644
--- a/core/services/mysql/src/backend.rs
+++ b/core/services/mysql/src/backend.rs
@@ -222,7 +222,9 @@ impl Access for MysqlBackend {
Some(bs) => bs,
None => return Err(Error::new(ErrorKind::NotFound, "kv not found
in mysql")),
};
- Ok((RpRead::new(), bs.slice(args.range().to_range_as_usize())))
+ let content = bs.slice(args.range().to_range_as_usize());
+ let metadata =
Metadata::new(EntryMode::FILE).with_content_length(bs.len() as u64);
+ Ok((RpRead::new(metadata), content))
}
async fn write(&self, path: &str, _: OpWrite) -> Result<(RpWrite,
Self::Writer)> {
diff --git a/core/services/obs/src/backend.rs b/core/services/obs/src/backend.rs
index 7deafba0f..3859e0bcd 100644
--- a/core/services/obs/src/backend.rs
+++ b/core/services/obs/src/backend.rs
@@ -327,9 +327,10 @@ impl Access for ObsBackend {
let status = resp.status();
match status {
- StatusCode::OK | StatusCode::PARTIAL_CONTENT => {
- Ok((RpRead::default(), resp.into_body()))
- }
+ StatusCode::OK | StatusCode::PARTIAL_CONTENT => Ok((
+ RpRead::new(parse_into_metadata(path, resp.headers())?),
+ resp.into_body(),
+ )),
_ => {
let (part, mut body) = resp.into_parts();
let buf = body.to_buffer().await?;
diff --git a/core/services/onedrive/src/backend.rs
b/core/services/onedrive/src/backend.rs
index 31b0a2b97..8c7d8d9e7 100644
--- a/core/services/onedrive/src/backend.rs
+++ b/core/services/onedrive/src/backend.rs
@@ -67,9 +67,10 @@ impl Access for OnedriveBackend {
async fn read(&self, path: &str, args: OpRead) -> Result<(RpRead,
Self::Reader)> {
let response = self.core.onedrive_get_content(path, &args).await?;
match response.status() {
- StatusCode::OK | StatusCode::PARTIAL_CONTENT => {
- Ok((RpRead::default(), response.into_body()))
- }
+ StatusCode::OK | StatusCode::PARTIAL_CONTENT => Ok((
+ RpRead::new(parse_into_metadata(path, response.headers())?),
+ response.into_body(),
+ )),
_ => {
let (part, mut body) = response.into_parts();
let buf = body.to_buffer().await?;
diff --git a/core/services/opfs/src/backend.rs
b/core/services/opfs/src/backend.rs
index 63a80bbda..97dd2faf8 100644
--- a/core/services/opfs/src/backend.rs
+++ b/core/services/opfs/src/backend.rs
@@ -108,8 +108,7 @@ impl Access for OpfsBackend {
async fn read(&self, path: &str, args: OpRead) -> Result<(RpRead,
Self::Reader)> {
let p = build_abs_path(&self.core.root, path);
let handle = get_file_handle(&p, false).await?;
-
- Ok((RpRead::new(), OpfsReader::new(handle, args.range())))
+ Ok((RpRead::default(), OpfsReader::new(handle, args.range())))
}
async fn list(&self, path: &str, _args: OpList) -> Result<(RpList,
Self::Lister)> {
diff --git a/core/services/oss/src/backend.rs b/core/services/oss/src/backend.rs
index 544dd41fb..8e6e1bbce 100644
--- a/core/services/oss/src/backend.rs
+++ b/core/services/oss/src/backend.rs
@@ -676,9 +676,10 @@ impl Access for OssBackend {
let status = resp.status();
match status {
- StatusCode::OK | StatusCode::PARTIAL_CONTENT => {
- Ok((RpRead::default(), resp.into_body()))
- }
+ StatusCode::OK | StatusCode::PARTIAL_CONTENT => Ok((
+ RpRead::new(parse_into_metadata(path, resp.headers())?),
+ resp.into_body(),
+ )),
_ => {
let (part, mut body) = resp.into_parts();
let buf = body.to_buffer().await?;
diff --git a/core/services/pcloud/src/backend.rs
b/core/services/pcloud/src/backend.rs
index 7e6d1898d..acd525e72 100644
--- a/core/services/pcloud/src/backend.rs
+++ b/core/services/pcloud/src/backend.rs
@@ -231,9 +231,10 @@ impl Access for PcloudBackend {
let status = resp.status();
match status {
- StatusCode::OK | StatusCode::PARTIAL_CONTENT => {
- Ok((RpRead::default(), resp.into_body()))
- }
+ StatusCode::OK | StatusCode::PARTIAL_CONTENT => Ok((
+ RpRead::new(parse_into_metadata(path, resp.headers())?),
+ resp.into_body(),
+ )),
_ => {
let (part, mut body) = resp.into_parts();
let buf = body.to_buffer().await?;
diff --git a/core/services/persy/src/backend.rs
b/core/services/persy/src/backend.rs
index 519c29dc9..756a790e1 100644
--- a/core/services/persy/src/backend.rs
+++ b/core/services/persy/src/backend.rs
@@ -184,7 +184,9 @@ impl Access for PersyBackend {
return Err(Error::new(ErrorKind::NotFound, "kv not found in
persy"));
}
};
- Ok((RpRead::new(), bs.slice(args.range().to_range_as_usize())))
+ let content = bs.slice(args.range().to_range_as_usize());
+ let metadata =
Metadata::new(EntryMode::FILE).with_content_length(bs.len() as u64);
+ Ok((RpRead::new(metadata), content))
}
async fn write(&self, path: &str, _: OpWrite) -> Result<(RpWrite,
Self::Writer)> {
diff --git a/core/services/postgresql/src/backend.rs
b/core/services/postgresql/src/backend.rs
index ea1b0c5f5..98c7e876a 100644
--- a/core/services/postgresql/src/backend.rs
+++ b/core/services/postgresql/src/backend.rs
@@ -222,7 +222,9 @@ impl Access for PostgresqlBackend {
));
}
};
- Ok((RpRead::new(), bs.slice(args.range().to_range_as_usize())))
+ let content = bs.slice(args.range().to_range_as_usize());
+ let metadata =
Metadata::new(EntryMode::FILE).with_content_length(bs.len() as u64);
+ Ok((RpRead::new(metadata), content))
}
async fn write(&self, path: &str, _: OpWrite) -> Result<(RpWrite,
Self::Writer)> {
diff --git a/core/services/redb/src/backend.rs
b/core/services/redb/src/backend.rs
index 11d9f1c2c..61a13a01b 100644
--- a/core/services/redb/src/backend.rs
+++ b/core/services/redb/src/backend.rs
@@ -202,7 +202,9 @@ impl Access for RedbBackend {
return Err(Error::new(ErrorKind::NotFound, "kv not found in
redb"));
}
};
- Ok((RpRead::new(), bs.slice(args.range().to_range_as_usize())))
+ let content = bs.slice(args.range().to_range_as_usize());
+ let metadata =
Metadata::new(EntryMode::FILE).with_content_length(bs.len() as u64);
+ Ok((RpRead::new(metadata), content))
}
async fn write(&self, path: &str, _: OpWrite) -> Result<(RpWrite,
Self::Writer)> {
diff --git a/core/services/redis/src/backend.rs
b/core/services/redis/src/backend.rs
index a06012348..11d401606 100644
--- a/core/services/redis/src/backend.rs
+++ b/core/services/redis/src/backend.rs
@@ -338,10 +338,14 @@ impl Access for RedisBackend {
let p = build_abs_path(&self.root, path);
let range = args.range();
- let buffer = if range.is_full() {
+ let (buffer, metadata) = if range.is_full() {
// Full read - use GET
match self.core.get(&p).await? {
- Some(bs) => bs,
+ Some(bs) => {
+ let metadata =
+
Metadata::new(EntryMode::FILE).with_content_length(bs.len() as u64);
+ (bs, Some(metadata))
+ }
None => return Err(Error::new(ErrorKind::NotFound, "key not
found in redis")),
}
} else {
@@ -353,12 +357,13 @@ impl Access for RedisBackend {
};
match self.core.get_range(&p, start, end).await? {
- Some(bs) => bs,
+ Some(bs) => (bs, None),
None => return Err(Error::new(ErrorKind::NotFound, "key not
found in redis")),
}
};
- Ok((RpRead::new(), buffer))
+ let rp = metadata.map_or_else(RpRead::default, RpRead::new);
+ Ok((rp, buffer))
}
async fn write(&self, path: &str, _: OpWrite) -> Result<(RpWrite,
Self::Writer)> {
diff --git a/core/services/rocksdb/src/backend.rs
b/core/services/rocksdb/src/backend.rs
index 8fb30be06..f0a48126c 100644
--- a/core/services/rocksdb/src/backend.rs
+++ b/core/services/rocksdb/src/backend.rs
@@ -152,7 +152,9 @@ impl Access for RocksdbBackend {
return Err(Error::new(ErrorKind::NotFound, "kv not found in
rocksdb"));
}
};
- Ok((RpRead::new(), bs.slice(args.range().to_range_as_usize())))
+ let content = bs.slice(args.range().to_range_as_usize());
+ let metadata =
Metadata::new(EntryMode::FILE).with_content_length(bs.len() as u64);
+ Ok((RpRead::new(metadata), content))
}
async fn write(&self, path: &str, _: OpWrite) -> Result<(RpWrite,
Self::Writer)> {
diff --git a/core/services/s3/src/backend.rs b/core/services/s3/src/backend.rs
index 86431e533..a7db1b514 100644
--- a/core/services/s3/src/backend.rs
+++ b/core/services/s3/src/backend.rs
@@ -1082,9 +1082,10 @@ impl Access for S3Backend {
let status = resp.status();
match status {
- StatusCode::OK | StatusCode::PARTIAL_CONTENT => {
- Ok((RpRead::default(), resp.into_body()))
- }
+ StatusCode::OK | StatusCode::PARTIAL_CONTENT => Ok((
+ RpRead::new(parse_into_metadata(path, resp.headers())?),
+ resp.into_body(),
+ )),
_ => {
let (part, mut body) = resp.into_parts();
let buf = body.to_buffer().await?;
diff --git a/core/services/s3/src/error.rs b/core/services/s3/src/error.rs
index da675c6de..137851a88 100644
--- a/core/services/s3/src/error.rs
+++ b/core/services/s3/src/error.rs
@@ -131,6 +131,7 @@ pub fn parse_s3_error_code(code: &str) ->
Option<(ErrorKind, bool)> {
| "ExceedAccountRateLimit"
| "ExceedBucketQPSLimit"
| "ExceedBucketRateLimit" => Some((ErrorKind::RateLimited, true)),
+ "InvalidRange" => Some((ErrorKind::RangeNotSatisfied, false)),
_ => None,
}
}
@@ -180,4 +181,12 @@ mod tests {
let out: S3Error = de::from_reader(bs.reader()).expect("must success");
assert_eq!(out, S3Error::default());
}
+
+ #[test]
+ fn test_parse_s3_error_code_invalid_range() {
+ assert_eq!(
+ parse_s3_error_code("InvalidRange"),
+ Some((ErrorKind::RangeNotSatisfied, false))
+ );
+ }
}
diff --git a/core/services/seafile/src/backend.rs
b/core/services/seafile/src/backend.rs
index 636d78595..f634755c4 100644
--- a/core/services/seafile/src/backend.rs
+++ b/core/services/seafile/src/backend.rs
@@ -237,9 +237,10 @@ impl Access for SeafileBackend {
let status = resp.status();
match status {
- StatusCode::OK | StatusCode::PARTIAL_CONTENT => {
- Ok((RpRead::default(), resp.into_body()))
- }
+ StatusCode::OK | StatusCode::PARTIAL_CONTENT => Ok((
+ RpRead::new(parse_into_metadata(path, resp.headers())?),
+ resp.into_body(),
+ )),
_ => {
let (part, mut body) = resp.into_parts();
let buf = body.to_buffer().await?;
diff --git a/core/services/sled/src/backend.rs
b/core/services/sled/src/backend.rs
index 3832363e5..b59165cb3 100644
--- a/core/services/sled/src/backend.rs
+++ b/core/services/sled/src/backend.rs
@@ -176,7 +176,9 @@ impl Access for SledBackend {
return Err(Error::new(ErrorKind::NotFound, "kv not found in
sled"));
}
};
- Ok((RpRead::new(), bs.slice(args.range().to_range_as_usize())))
+ let content = bs.slice(args.range().to_range_as_usize());
+ let metadata =
Metadata::new(EntryMode::FILE).with_content_length(bs.len() as u64);
+ Ok((RpRead::new(metadata), content))
}
async fn write(&self, path: &str, _: OpWrite) -> Result<(RpWrite,
Self::Writer)> {
diff --git a/core/services/sqlite/src/backend.rs
b/core/services/sqlite/src/backend.rs
index 9bdbc73df..c53a07d5c 100644
--- a/core/services/sqlite/src/backend.rs
+++ b/core/services/sqlite/src/backend.rs
@@ -260,10 +260,13 @@ impl Access for SqliteBackend {
let p = build_abs_path(&self.root, path);
let range = args.range();
- let buffer = if range.is_full() {
+ let (buffer, content_length) = if range.is_full() {
// Full read - use GET
match self.core.get(&p).await? {
- Some(bs) => bs,
+ Some(bs) => {
+ let content_length = bs.len() as u64;
+ (bs, content_length)
+ }
None => return Err(Error::new(ErrorKind::NotFound, "key not
found in sqlite")),
}
} else {
@@ -275,12 +278,13 @@ impl Access for SqliteBackend {
};
match self.core.get_range(&p, start, limit).await? {
- Some(bs) => bs,
+ Some((bs, content_length)) => (bs, content_length),
None => return Err(Error::new(ErrorKind::NotFound, "key not
found in sqlite")),
}
};
- Ok((RpRead::new(), buffer))
+ let metadata =
Metadata::new(EntryMode::FILE).with_content_length(content_length);
+ Ok((RpRead::new(metadata), buffer))
}
async fn write(&self, path: &str, _: OpWrite) -> Result<(RpWrite,
Self::Writer)> {
diff --git a/core/services/sqlite/src/core.rs b/core/services/sqlite/src/core.rs
index 56499db6c..95766b1d4 100644
--- a/core/services/sqlite/src/core.rs
+++ b/core/services/sqlite/src/core.rs
@@ -65,13 +65,14 @@ impl SqliteCore {
path: &str,
start: isize,
limit: isize,
- ) -> Result<Option<Buffer>> {
+ ) -> Result<Option<(Buffer, u64)>> {
let pool = self.get_client().await?;
- let value: Option<Vec<u8>> = sqlx::query_scalar(&format!(
- "SELECT SUBSTR(`{}`, {}, {}) FROM `{}` WHERE `{}` = $1 LIMIT 1",
+ let value: Option<(Vec<u8>, i64)> = sqlx::query_as(&format!(
+ "SELECT SUBSTR(`{}`, {}, {}), LENGTH(`{}`) FROM `{}` WHERE `{}` =
$1 LIMIT 1",
self.value_field,
start + 1,
limit,
+ self.value_field,
self.table,
self.key_field
))
@@ -80,7 +81,7 @@ impl SqliteCore {
.await
.map_err(parse_sqlite_error)?;
- Ok(value.map(Buffer::from))
+ Ok(value.map(|(bs, size)| (Buffer::from(bs), size as u64)))
}
pub async fn set(&self, path: &str, value: Buffer) -> Result<()> {
diff --git a/core/services/surrealdb/src/backend.rs
b/core/services/surrealdb/src/backend.rs
index 286f737c8..7753357fe 100644
--- a/core/services/surrealdb/src/backend.rs
+++ b/core/services/surrealdb/src/backend.rs
@@ -270,7 +270,9 @@ impl Access for SurrealdbBackend {
return Err(Error::new(ErrorKind::NotFound, "kv not found in
surrealdb"));
}
};
- Ok((RpRead::new(), bs.slice(args.range().to_range_as_usize())))
+ let content = bs.slice(args.range().to_range_as_usize());
+ let metadata =
Metadata::new(EntryMode::FILE).with_content_length(bs.len() as u64);
+ Ok((RpRead::new(metadata), content))
}
async fn write(&self, path: &str, _: OpWrite) -> Result<(RpWrite,
Self::Writer)> {
diff --git a/core/services/swift/src/backend.rs
b/core/services/swift/src/backend.rs
index 8aad1f7e4..7db62df0b 100644
--- a/core/services/swift/src/backend.rs
+++ b/core/services/swift/src/backend.rs
@@ -271,7 +271,10 @@ impl Access for SwiftBackend {
let status = resp.status();
match status {
- StatusCode::OK | StatusCode::PARTIAL_CONTENT => Ok((RpRead::new(),
resp.into_body())),
+ StatusCode::OK | StatusCode::PARTIAL_CONTENT => Ok((
+ RpRead::new(parse_into_metadata(path, resp.headers())?),
+ resp.into_body(),
+ )),
_ => {
let (part, mut body) = resp.into_parts();
let buf = body.to_buffer().await?;
diff --git a/core/services/tikv/src/backend.rs
b/core/services/tikv/src/backend.rs
index 4268c91cf..3a29e38d3 100644
--- a/core/services/tikv/src/backend.rs
+++ b/core/services/tikv/src/backend.rs
@@ -176,7 +176,9 @@ impl Access for TikvBackend {
Some(bs) => bs,
None => return Err(Error::new(ErrorKind::NotFound, "kv not found
in tikv")),
};
- Ok((RpRead::new(), bs.slice(args.range().to_range_as_usize())))
+ let content = bs.slice(args.range().to_range_as_usize());
+ let metadata =
Metadata::new(EntryMode::FILE).with_content_length(bs.len() as u64);
+ Ok((RpRead::new(metadata), content))
}
async fn write(&self, path: &str, _: OpWrite) -> Result<(RpWrite,
Self::Writer)> {
diff --git a/core/services/tos/src/backend.rs b/core/services/tos/src/backend.rs
index 872482af5..054bb6fb8 100644
--- a/core/services/tos/src/backend.rs
+++ b/core/services/tos/src/backend.rs
@@ -322,9 +322,10 @@ impl Access for TosBackend {
let status = resp.status();
match status {
- StatusCode::OK | StatusCode::PARTIAL_CONTENT => {
- Ok((RpRead::default(), resp.into_body()))
- }
+ StatusCode::OK | StatusCode::PARTIAL_CONTENT => Ok((
+ RpRead::new(parse_into_metadata(path, resp.headers())?),
+ resp.into_body(),
+ )),
_ => {
let (part, mut body) = resp.into_parts();
let buf = body.to_buffer().await?;
diff --git a/core/services/upyun/src/backend.rs
b/core/services/upyun/src/backend.rs
index bf17987f6..34cc0641f 100644
--- a/core/services/upyun/src/backend.rs
+++ b/core/services/upyun/src/backend.rs
@@ -227,9 +227,10 @@ impl Access for UpyunBackend {
let status = resp.status();
match status {
- StatusCode::OK | StatusCode::PARTIAL_CONTENT => {
- Ok((RpRead::default(), resp.into_body()))
- }
+ StatusCode::OK | StatusCode::PARTIAL_CONTENT => Ok((
+ RpRead::new(parse_into_metadata(path, resp.headers())?),
+ resp.into_body(),
+ )),
_ => {
let (part, mut body) = resp.into_parts();
let buf = body.to_buffer().await?;
diff --git a/core/services/vercel-artifacts/src/backend.rs
b/core/services/vercel-artifacts/src/backend.rs
index a2ef7db2d..9b8c9c24d 100644
--- a/core/services/vercel-artifacts/src/backend.rs
+++ b/core/services/vercel-artifacts/src/backend.rs
@@ -67,9 +67,10 @@ impl Access for VercelArtifactsBackend {
let status = response.status();
match status {
- StatusCode::OK | StatusCode::PARTIAL_CONTENT => {
- Ok((RpRead::new(), response.into_body()))
- }
+ StatusCode::OK | StatusCode::PARTIAL_CONTENT => Ok((
+ RpRead::new(parse_into_metadata(path, response.headers())?),
+ response.into_body(),
+ )),
_ => {
let (part, mut body) = response.into_parts();
let buf = body.to_buffer().await?;
diff --git a/core/services/vercel-blob/src/backend.rs
b/core/services/vercel-blob/src/backend.rs
index d9f792fe6..a5a441fd2 100644
--- a/core/services/vercel-blob/src/backend.rs
+++ b/core/services/vercel-blob/src/backend.rs
@@ -172,9 +172,10 @@ impl Access for VercelBlobBackend {
let status = resp.status();
match status {
- StatusCode::OK | StatusCode::PARTIAL_CONTENT => {
- Ok((RpRead::default(), resp.into_body()))
- }
+ StatusCode::OK | StatusCode::PARTIAL_CONTENT => Ok((
+ RpRead::new(parse_into_metadata(path, resp.headers())?),
+ resp.into_body(),
+ )),
_ => {
let (part, mut body) = resp.into_parts();
let buf = body.to_buffer().await?;
diff --git a/core/services/webdav/src/backend.rs
b/core/services/webdav/src/backend.rs
index 1a1969687..2eb9b41fa 100644
--- a/core/services/webdav/src/backend.rs
+++ b/core/services/webdav/src/backend.rs
@@ -276,9 +276,10 @@ impl Access for WebdavBackend {
let status = resp.status();
match status {
- StatusCode::OK | StatusCode::PARTIAL_CONTENT => {
- Ok((RpRead::default(), resp.into_body()))
- }
+ StatusCode::OK | StatusCode::PARTIAL_CONTENT => Ok((
+ RpRead::new(parse_into_metadata(path, resp.headers())?),
+ resp.into_body(),
+ )),
_ => {
let (part, mut body) = resp.into_parts();
let buf = body.to_buffer().await?;
diff --git a/core/services/yandex-disk/src/backend.rs
b/core/services/yandex-disk/src/backend.rs
index 19fd116f1..ca67770fe 100644
--- a/core/services/yandex-disk/src/backend.rs
+++ b/core/services/yandex-disk/src/backend.rs
@@ -190,7 +190,10 @@ impl Access for YandexDiskBackend {
let status = resp.status();
match status {
- StatusCode::OK | StatusCode::PARTIAL_CONTENT => Ok((RpRead::new(),
resp.into_body())),
+ StatusCode::OK | StatusCode::PARTIAL_CONTENT => Ok((
+ RpRead::new(parse_into_metadata(path, resp.headers())?),
+ resp.into_body(),
+ )),
_ => {
let (part, mut body) = resp.into_parts();
let buf = body.to_buffer().await?;
diff --git a/core/tests/behavior/async_read.rs
b/core/tests/behavior/async_read.rs
index 268795301..89fedf164 100644
--- a/core/tests/behavior/async_read.rs
+++ b/core/tests/behavior/async_read.rs
@@ -35,6 +35,10 @@ pub fn tests(op: &Operator, tests: &mut Vec<Trial>) {
test_read_full,
test_read_range,
test_reader,
+ test_buffer_stream_metadata,
+ test_buffer_stream_metadata_with_concurrent,
+ test_futures_bytes_stream_metadata,
+ test_futures_bytes_stream_metadata_with_concurrent,
test_reader_with_if_match,
test_reader_with_if_none_match,
test_reader_with_if_modified_since,
@@ -150,6 +154,176 @@ pub async fn test_reader(op: Operator) ->
anyhow::Result<()> {
Ok(())
}
+/// BufferStream should return complete object metadata before reading if
supported.
+pub async fn test_buffer_stream_metadata(op: Operator) -> anyhow::Result<()> {
+ let path = TEST_FIXTURE.new_file_path();
+ let content = gen_fixed_bytes(1024);
+ let start = 128;
+ let end = 640;
+
+ op.write(&path, content.clone())
+ .await
+ .expect("write must succeed");
+
+ let mut stream = op
+ .reader(&path)
+ .await?
+ .into_stream(start as u64..end as u64)
+ .await?;
+
+ match stream.metadata().await {
+ Ok(meta) => assert_eq!(
+ meta.content_length(),
+ content.len() as u64,
+ "metadata content length"
+ ),
+ Err(err) if err.kind() == ErrorKind::Unsupported => {}
+ Err(err) => return Err(err.into()),
+ }
+
+ let bs: Vec<_> = stream.try_collect().await?;
+ let bs: Buffer = bs.into_iter().flatten().collect();
+ assert_eq!(bs.len(), end - start, "read size");
+ assert_eq!(
+ sha256_digest(bs.to_bytes()),
+ sha256_digest(&content[start..end]),
+ "read content"
+ );
+
+ Ok(())
+}
+
+/// BufferStream should return complete object metadata with concurrent reads
if supported.
+pub async fn test_buffer_stream_metadata_with_concurrent(op: Operator) ->
anyhow::Result<()> {
+ let path = TEST_FIXTURE.new_file_path();
+ let content = gen_fixed_bytes(1024);
+ let start = 128;
+ let end = 640;
+
+ op.write(&path, content.clone())
+ .await
+ .expect("write must succeed");
+
+ let mut stream = op
+ .reader_with(&path)
+ .chunk(128)
+ .concurrent(4)
+ .await?
+ .into_stream(start as u64..end as u64)
+ .await?;
+
+ match stream.metadata().await {
+ Ok(meta) => assert_eq!(
+ meta.content_length(),
+ content.len() as u64,
+ "metadata content length"
+ ),
+ Err(err) if err.kind() == ErrorKind::Unsupported => {}
+ Err(err) => return Err(err.into()),
+ }
+
+ let bs: Vec<_> = stream.try_collect().await?;
+ let bs: Buffer = bs.into_iter().flatten().collect();
+ assert_eq!(bs.len(), end - start, "read size");
+ assert_eq!(
+ sha256_digest(bs.to_bytes()),
+ sha256_digest(&content[start..end]),
+ "read content"
+ );
+
+ Ok(())
+}
+
+/// FuturesBytesStream should return complete object metadata before reading
if supported.
+pub async fn test_futures_bytes_stream_metadata(op: Operator) ->
anyhow::Result<()> {
+ let path = TEST_FIXTURE.new_file_path();
+ let content = gen_fixed_bytes(1024);
+ let start = 128;
+ let end = 640;
+
+ op.write(&path, content.clone())
+ .await
+ .expect("write must succeed");
+
+ let mut stream = op
+ .reader(&path)
+ .await?
+ .into_bytes_stream(start as u64..end as u64)
+ .await?;
+
+ match stream.metadata().await {
+ Ok(meta) => assert_eq!(
+ meta.content_length(),
+ content.len() as u64,
+ "metadata content length"
+ ),
+ Err(err) if err.kind() == ErrorKind::Unsupported => {}
+ Err(err) => return Err(err.into()),
+ }
+
+ let bs = stream
+ .try_fold(Vec::new(), |mut acc, chunk| {
+ acc.extend_from_slice(&chunk);
+ async { Ok(acc) }
+ })
+ .await?;
+ assert_eq!(bs.len(), end - start, "read size");
+ assert_eq!(
+ sha256_digest(&bs),
+ sha256_digest(&content[start..end]),
+ "read content"
+ );
+
+ Ok(())
+}
+
+/// FuturesBytesStream should return complete object metadata with concurrent
reads if supported.
+pub async fn test_futures_bytes_stream_metadata_with_concurrent(
+ op: Operator,
+) -> anyhow::Result<()> {
+ let path = TEST_FIXTURE.new_file_path();
+ let content = gen_fixed_bytes(1024);
+ let start = 128;
+ let end = 640;
+
+ op.write(&path, content.clone())
+ .await
+ .expect("write must succeed");
+
+ let mut stream = op
+ .reader_with(&path)
+ .chunk(128)
+ .concurrent(4)
+ .await?
+ .into_bytes_stream(start as u64..end as u64)
+ .await?;
+
+ match stream.metadata().await {
+ Ok(meta) => assert_eq!(
+ meta.content_length(),
+ content.len() as u64,
+ "metadata content length"
+ ),
+ Err(err) if err.kind() == ErrorKind::Unsupported => {}
+ Err(err) => return Err(err.into()),
+ }
+
+ let bs = stream
+ .try_fold(Vec::new(), |mut acc, chunk| {
+ acc.extend_from_slice(&chunk);
+ async { Ok(acc) }
+ })
+ .await?;
+ assert_eq!(bs.len(), end - start, "read size");
+ assert_eq!(
+ sha256_digest(&bs),
+ sha256_digest(&content[start..end]),
+ "read content"
+ );
+
+ Ok(())
+}
+
/// Read not exist file should return NotFound
pub async fn test_read_not_exist(op: Operator) -> anyhow::Result<()> {
let path = uuid::Uuid::new_v4().to_string();
diff --git a/integrations/object_store/src/service/reader.rs
b/integrations/object_store/src/service/reader.rs
index 30f0114a5..852c8910d 100644
--- a/integrations/object_store/src/service/reader.rs
+++ b/integrations/object_store/src/service/reader.rs
@@ -26,6 +26,7 @@ use object_store::path::Path as ObjectStorePath;
use opendal::raw::*;
use opendal::*;
+use super::core::format_metadata;
use super::core::parse_op_read;
use super::error::parse_error;
@@ -33,7 +34,6 @@ use super::error::parse_error;
pub struct ObjectStoreReader {
bytes_stream: BoxStream<'static, object_store::Result<Bytes>>,
meta: object_store::ObjectMeta,
- args: OpRead,
}
impl ObjectStoreReader {
@@ -47,26 +47,11 @@ impl ObjectStoreReader {
let result = store.get_opts(&path, opts).await.map_err(parse_error)?;
let meta = result.meta.clone();
let bytes_stream = result.into_stream();
- Ok(Self {
- bytes_stream,
- meta,
- args,
- })
+ Ok(Self { bytes_stream, meta })
}
pub(crate) fn rp(&self) -> RpRead {
- let mut rp = RpRead::new().with_size(Some(self.meta.size));
- if !self.args.range().is_full() {
- let range = self.args.range();
- let size = match range.size() {
- Some(size) => size,
- None => self.meta.size,
- };
- rp = rp.with_range(Some(
- BytesContentRange::default().with_range(range.offset(),
range.offset() + size - 1),
- ));
- }
- rp
+ RpRead::new(format_metadata(&self.meta))
}
}
diff --git a/integrations/object_store/src/store.rs
b/integrations/object_store/src/store.rs
index 17c1a4129..e336b8cf0 100644
--- a/integrations/object_store/src/store.rs
+++ b/integrations/object_store/src/store.rs
@@ -21,8 +21,8 @@ use std::io;
use std::ops::Range;
use std::sync::Arc;
+use crate::datetime_to_timestamp;
use crate::utils::*;
-use crate::{datetime_to_timestamp, timestamp_to_datetime};
use async_trait::async_trait;
use bytes::Bytes;
use futures::FutureExt;
@@ -31,6 +31,7 @@ use futures::TryStreamExt;
use futures::stream::BoxStream;
use mea::mutex::Mutex;
use mea::oneshot;
+use object_store::Attributes;
use object_store::CopyMode as ObjectStoreCopyMode;
use object_store::CopyOptions as ObjectStoreCopyOptions;
use object_store::ListResult;
@@ -48,12 +49,85 @@ use object_store::{GetResult, PutMode};
use opendal::Buffer;
use opendal::Writer;
use opendal::options::CopyOptions;
+use opendal::options::ReaderOptions;
+use opendal::options::StatOptions;
use opendal::raw::percent_decode_path;
use opendal::{Operator, OperatorInfo};
use std::collections::HashMap;
const DEFAULT_CONCURRENT: usize = 8;
+fn format_object_attributes(meta: &opendal::Metadata) -> Attributes {
+ let mut attributes = Attributes::new();
+ if let Some(user_meta) = meta.user_metadata() {
+ for (key, value) in user_meta {
+ attributes.insert(
+ object_store::Attribute::Metadata(key.clone().into()),
+ value.clone().into(),
+ );
+ }
+ }
+
+ attributes
+}
+
+fn format_reader_options(options: &GetOptions, content_length_hint:
Option<u64>) -> ReaderOptions {
+ ReaderOptions {
+ version: options.version.clone(),
+ if_match: options.if_match.clone(),
+ if_none_match: options.if_none_match.clone(),
+ if_modified_since:
options.if_modified_since.and_then(datetime_to_timestamp),
+ if_unmodified_since:
options.if_unmodified_since.and_then(datetime_to_timestamp),
+ content_length_hint,
+ ..Default::default()
+ }
+}
+
+fn format_stat_options(options: &GetOptions) -> StatOptions {
+ StatOptions {
+ version: options.version.clone(),
+ if_match: options.if_match.clone(),
+ if_none_match: options.if_none_match.clone(),
+ if_modified_since:
options.if_modified_since.and_then(datetime_to_timestamp),
+ if_unmodified_since:
options.if_unmodified_since.and_then(datetime_to_timestamp),
+ ..Default::default()
+ }
+}
+
+fn format_read_range(range: Option<&GetRange>, size: u64) -> Range<u64> {
+ match range {
+ Some(GetRange::Bounded(r)) => {
+ if r.start >= r.end || r.start >= size {
+ 0..0
+ } else {
+ r.start..r.end.min(size)
+ }
+ }
+ Some(GetRange::Offset(offset)) => {
+ if *offset < size {
+ *offset..size
+ } else {
+ 0..0
+ }
+ }
+ Some(GetRange::Suffix(suffix)) if *suffix < size => (size -
*suffix)..size,
+ _ => 0..size,
+ }
+}
+
+fn format_without_stat_error(err: opendal::Error, path: &str) ->
object_store::Error {
+ match err.kind() {
+ // Ask get_opts to fall back to the stat path when read-open can't
provide
+ // enough metadata to build a valid GetResult, such as ranges beyond
EOF.
+ opendal::ErrorKind::Unsupported |
opendal::ErrorKind::RangeNotSatisfied => {
+ object_store::Error::NotSupported {
+ source: Box::new(err),
+ }
+ }
+ _ => format_object_store_error(err, path),
+ }
+}
+
/// OpendalStore implements ObjectStore trait by using opendal.
///
/// This allows users to use opendal as an object store without extra cost.
@@ -121,6 +195,147 @@ impl OpendalStore {
self.info.as_ref()
}
+ async fn get_opts_without_stat(
+ &self,
+ location: &Path,
+ raw_location: &str,
+ options: &GetOptions,
+ ) -> object_store::Result<GetResult> {
+ let reader = self
+ .inner
+ .reader_options(raw_location, format_reader_options(options, None))
+ .into_send()
+ .await
+ .map_err(|err| format_without_stat_error(err, location.as_ref()))?;
+
+ let mut stream = match options.range.as_ref() {
+ Some(GetRange::Bounded(range)) => {
+ reader
+ .into_bytes_stream(range.start..range.end)
+ .into_send()
+ .await
+ }
+ Some(GetRange::Offset(offset)) =>
reader.into_bytes_stream(*offset..).into_send().await,
+ Some(GetRange::Suffix(_)) => unreachable!("suffix range needs
object metadata"),
+ None => reader.into_bytes_stream(..).into_send().await,
+ }
+ .map_err(|err| format_without_stat_error(err, location.as_ref()))?;
+
+ let metadata = stream
+ .metadata()
+ .into_send()
+ .await
+ .map_err(|err| format_without_stat_error(err, location.as_ref()))?;
+ let attributes = format_object_attributes(&metadata);
+ let meta = format_object_meta(location.as_ref(), &metadata);
+ let read_range = format_read_range(options.range.as_ref(), meta.size);
+
+ if read_range.start >= read_range.end {
+ return Ok(GetResult {
+ payload:
GetResultPayload::Stream(Box::pin(futures::stream::empty())),
+ range: read_range,
+ meta,
+ attributes,
+ });
+ }
+
+ if matches!(
+ options.range.as_ref(),
+ Some(GetRange::Bounded(range)) if range.end > meta.size
+ ) {
+ let reader = self
+ .inner
+ .reader_options(
+ raw_location,
+ format_reader_options(options, Some(meta.size)),
+ )
+ .into_send()
+ .await
+ .map_err(|err| format_object_store_error(err,
location.as_ref()))?;
+ stream = reader
+ .into_bytes_stream(read_range.start..read_range.end)
+ .into_send()
+ .await
+ .map_err(|err| format_object_store_error(err,
location.as_ref()))?;
+ }
+
+ let stream = stream
+ .into_send()
+ .map_err(|err: io::Error| object_store::Error::Generic {
+ store: "IoError",
+ source: Box::new(err),
+ });
+
+ Ok(GetResult {
+ payload: GetResultPayload::Stream(Box::pin(stream)),
+ range: read_range,
+ meta,
+ attributes,
+ })
+ }
+
+ async fn get_opts_with_stat(
+ &self,
+ location: &Path,
+ raw_location: &str,
+ options: &GetOptions,
+ ) -> object_store::Result<GetResult> {
+ let metadata = self
+ .inner
+ .stat_options(raw_location, format_stat_options(options))
+ .into_send()
+ .await
+ .map_err(|err| format_object_store_error(err, location.as_ref()))?;
+ let attributes = format_object_attributes(&metadata);
+ let meta = format_object_meta(location.as_ref(), &metadata);
+
+ if options.head {
+ return Ok(GetResult {
+ payload:
GetResultPayload::Stream(Box::pin(futures::stream::empty())),
+ range: 0..0,
+ meta,
+ attributes,
+ });
+ }
+
+ let read_range = format_read_range(options.range.as_ref(), meta.size);
+ if read_range.start >= read_range.end {
+ return Ok(GetResult {
+ payload:
GetResultPayload::Stream(Box::pin(futures::stream::empty())),
+ range: read_range,
+ meta,
+ attributes,
+ });
+ }
+
+ let reader = self
+ .inner
+ .reader_options(
+ raw_location,
+ format_reader_options(options, Some(meta.size)),
+ )
+ .into_send()
+ .await
+ .map_err(|err| format_object_store_error(err, location.as_ref()))?;
+ let stream = reader
+ .into_bytes_stream(read_range.start..read_range.end)
+ .into_send()
+ .await
+ .map_err(|err| format_object_store_error(err, location.as_ref()))?
+ .into_send()
+ .map_err(|err: io::Error| object_store::Error::Generic {
+ store: "IoError",
+ source: Box::new(err),
+ });
+
+ Ok(GetResult {
+ payload: GetResultPayload::Stream(Box::pin(stream)),
+ range: read_range,
+ meta,
+ attributes,
+ })
+ }
+
/// Copy a file from one location to another
async fn copy_request(
&self,
@@ -296,126 +511,26 @@ impl ObjectStore for OpendalStore {
options: GetOptions,
) -> object_store::Result<GetResult> {
let raw_location = percent_decode_path(location.as_ref());
- let meta = {
- let mut s = self.inner.stat_with(&raw_location);
- if let Some(version) = &options.version {
- s = s.version(version.as_str())
- }
- if let Some(if_match) = &options.if_match {
- s = s.if_match(if_match.as_str());
- }
- if let Some(if_none_match) = &options.if_none_match {
- s = s.if_none_match(if_none_match.as_str());
- }
- if let Some(if_modified_since) =
- options.if_modified_since.and_then(datetime_to_timestamp)
- {
- s = s.if_modified_since(if_modified_since);
- }
- if let Some(if_unmodified_since) =
- options.if_unmodified_since.and_then(datetime_to_timestamp)
- {
- s = s.if_unmodified_since(if_unmodified_since);
- }
- s.into_send()
- .await
- .map_err(|err| format_object_store_error(err,
location.as_ref()))?
- };
-
- // Convert user defined metadata from OpenDAL to object_store
attributes
- let mut attributes = object_store::Attributes::new();
- if let Some(user_meta) = meta.user_metadata() {
- for (key, value) in user_meta {
- attributes.insert(
- object_store::Attribute::Metadata(key.clone().into()),
- value.clone().into(),
- );
- }
- }
-
- let meta = ObjectMeta {
- location: location.clone(),
- last_modified: meta
- .last_modified()
- .and_then(timestamp_to_datetime)
- .unwrap_or_default(),
- size: meta.content_length(),
- e_tag: meta.etag().map(|x| x.to_string()),
- version: meta.version().map(|x| x.to_string()),
- };
if options.head {
- return Ok(GetResult {
- payload:
GetResultPayload::Stream(Box::pin(futures::stream::empty())),
- range: 0..0,
- meta,
- attributes,
- });
+ return self
+ .get_opts_with_stat(location, &raw_location, &options)
+ .await;
}
- let reader = {
- let mut r = self.inner.reader_with(raw_location.as_ref());
- if let Some(version) = options.version {
- r = r.version(version.as_str());
- }
- if let Some(if_match) = options.if_match {
- r = r.if_match(if_match.as_str());
- }
- if let Some(if_none_match) = options.if_none_match {
- r = r.if_none_match(if_none_match.as_str());
- }
- if let Some(if_modified_since) =
- options.if_modified_since.and_then(datetime_to_timestamp)
- {
- r = r.if_modified_since(if_modified_since);
- }
- if let Some(if_unmodified_since) =
- options.if_unmodified_since.and_then(datetime_to_timestamp)
- {
- r = r.if_unmodified_since(if_unmodified_since);
- }
- r.into_send()
+ if !matches!(options.range.as_ref(), Some(GetRange::Suffix(_))) {
+ match self
+ .get_opts_without_stat(location, &raw_location, &options)
.await
- .map_err(|err| format_object_store_error(err,
location.as_ref()))?
- };
-
- let read_range = match options.range {
- Some(GetRange::Bounded(r)) => {
- if r.start >= r.end || r.start >= meta.size {
- 0..0
- } else {
- let end = r.end.min(meta.size);
- r.start..end
- }
- }
- Some(GetRange::Offset(r)) => {
- if r < meta.size {
- r..meta.size
- } else {
- 0..0
- }
+ {
+ Ok(result) => return Ok(result),
+ Err(object_store::Error::NotSupported { .. }) => {}
+ Err(err) => return Err(err),
}
- Some(GetRange::Suffix(r)) if r < meta.size => (meta.size -
r)..meta.size,
- _ => 0..meta.size,
- };
+ }
- let stream = reader
- .into_bytes_stream(read_range.start..read_range.end)
- .into_send()
+ self.get_opts_with_stat(location, &raw_location, &options)
.await
- .map_err(|err| format_object_store_error(err, location.as_ref()))?
- .into_send()
- .map_err(|err: io::Error| object_store::Error::Generic {
- store: "IoError",
- source: Box::new(err),
- });
-
- Ok(GetResult {
- payload: GetResultPayload::Stream(Box::pin(stream)),
- range: read_range.start..read_range.end,
- meta,
- attributes,
- })
}
async fn get_ranges(
@@ -1017,17 +1132,75 @@ mod tests {
// Reset counter
stat_count.store(0, Ordering::SeqCst);
- // Test 2: get_opts SHOULD call stat() to get metadata
+ // Test 2: get_opts should NOT call stat() for bounded range reads
let opts = object_store::GetOptions {
range: Some(object_store::GetRange::Bounded(0..5)),
..Default::default()
};
let ret = store.get_opts(&location, opts).await.unwrap();
+ assert_eq!(ret.meta.size, value.len() as u64);
+ assert_eq!(ret.range, 0..5);
let data = ret.bytes().await.unwrap();
assert_eq!(Bytes::from_static(b"Hello"), data);
+ assert_eq!(
+ stat_count.load(Ordering::SeqCst),
+ 0,
+ "get_opts should not call stat() for bounded range reads"
+ );
+
+ // Reset counter
+ stat_count.store(0, Ordering::SeqCst);
+
+ // Test 3: get_opts should NOT call stat() for offset range reads
+ let opts = object_store::GetOptions {
+ range: Some(object_store::GetRange::Offset(7)),
+ ..Default::default()
+ };
+ let ret = store.get_opts(&location, opts).await.unwrap();
+ assert_eq!(ret.meta.size, value.len() as u64);
+ assert_eq!(ret.range, 7..value.len() as u64);
+ let data = ret.bytes().await.unwrap();
+ assert_eq!(Bytes::from_static(b"world!"), data);
+ assert_eq!(
+ stat_count.load(Ordering::SeqCst),
+ 0,
+ "get_opts should not call stat() for offset range reads"
+ );
+
+ // Reset counter
+ stat_count.store(0, Ordering::SeqCst);
+
+ // Test 4: get_opts should NOT call stat() for full reads
+ let ret = store
+ .get_opts(&location, object_store::GetOptions::default())
+ .await
+ .unwrap();
+ assert_eq!(ret.meta.size, value.len() as u64);
+ assert_eq!(ret.range, 0..value.len() as u64);
+ let data = ret.bytes().await.unwrap();
+ assert_eq!(value, data);
+ assert_eq!(
+ stat_count.load(Ordering::SeqCst),
+ 0,
+ "get_opts should not call stat() for full reads"
+ );
+
+ // Reset counter
+ stat_count.store(0, Ordering::SeqCst);
+
+ // Test 5: get_opts should still call stat() for suffix range reads
+ let opts = object_store::GetOptions {
+ range: Some(object_store::GetRange::Suffix(6)),
+ ..Default::default()
+ };
+ let ret = store.get_opts(&location, opts).await.unwrap();
+ assert_eq!(ret.meta.size, value.len() as u64);
+ assert_eq!(ret.range, 7..value.len() as u64);
+ let data = ret.bytes().await.unwrap();
+ assert_eq!(Bytes::from_static(b"world!"), data);
assert!(
stat_count.load(Ordering::SeqCst) > 0,
- "get_opts should call stat() to get metadata"
+ "get_opts should call stat() for suffix range reads"
);
// Cleanup