feat(object): serve GET misses from the on-demand migration source (#7084)

feat(rustfs): serve GET misses from the migration source

Wire the on-demand migration read-through into the GET path, after the
local read and the replication proxy have both missed (rustfs/backlog#2156).

A source HEAD supplies size, validators and metadata. Conditional headers
are evaluated locally against it and never forwarded, so a source 304/412
cannot be mistaken for a source failure. An object within inline_max_bytes
is teed: the primary streams to the client while the secondary commits the
local copy in a background task, so a client disconnect still stores the
whole object and a failed write-back never touches the client stream.
Range reads and larger objects stream straight through and queue a
background pull per policy. Concurrent misses of one key share the
singleflight slot: the leader tees, followers re-read local after it
commits or degrade to passthrough after first_byte_ms.

Version reads, partNumber reads, anti-loop marked requests and a respected
local delete marker keep their original 404. Source answers carry
x-rustfs-on-demand-migration: source; local hits are untouched, and the
local hit path gains no await or lock.
This commit is contained in:
Zhengchao An
2026-09-03 07:21:51 +08:00
committed by GitHub
parent 37344a84da
commit d04611ba38
6 changed files with 1601 additions and 21 deletions
@@ -22,7 +22,7 @@
//! not exercised by the harness self-test.
use crate::common::{RustFSTestEnvironment, signed_request};
use crate::fake_s3_target::{FAKE_ACCESS_KEY, FAKE_SECRET_KEY, FakeS3Target, FakeS3TargetOptions, SeedMetadata};
use crate::fake_s3_target::{FAKE_ACCESS_KEY, FAKE_SECRET_KEY, FakeS3Target, FakeS3TargetOptions, Operation, SeedMetadata};
use aws_config::retry::RetryConfig;
use aws_sdk_s3::Client;
use aws_sdk_s3::config::{Credentials, Region};
@@ -30,12 +30,16 @@ use aws_smithy_http_client::Builder as SmithyHttpClientBuilder;
use bytes::Bytes;
use serde::Serialize;
use std::fmt;
use std::time::{Duration, Instant};
pub type BoxError = Box<dyn std::error::Error + Send + Sync>;
/// Module switch the server reads at startup (`false` before GA). The harness
/// turns it on so scenario tests exercise the feature without repeating it.
pub const ODM_MODULE_SWITCH_ENV: &str = "RUSTFS_ON_DEMAND_MIGRATION_ENABLED";
/// The source client shares the replication egress guard, which rejects the
/// fake source's loopback endpoint unless this switch is set.
pub const ALLOW_LOOPBACK_SOURCE_ENV: &str = "RUSTFS_REPLICATION_ALLOW_LOOPBACK_TARGET";
/// Admin route prefix; the bucket name is appended as one path segment.
pub const ODM_ADMIN_ROUTE: &str = "/rustfs/admin/v3/on-demand-migration";
/// Region the fake source is addressed with (it accepts any SigV4 region).
@@ -235,6 +239,21 @@ impl AdminResponse {
}
}
/// Raw S3 response (status, headers, body) for assertions on headers the
/// SDK does not surface, such as `x-rustfs-on-demand-migration`.
#[derive(Debug, Clone)]
pub struct RawResponse {
pub status: u16,
pub headers: http::HeaderMap,
pub body: Bytes,
}
impl RawResponse {
pub fn header(&self, name: &str) -> Option<&str> {
self.headers.get(name).and_then(|value| value.to_str().ok())
}
}
/// One object to seed into the source.
#[derive(Clone)]
pub struct SeedObject {
@@ -277,7 +296,7 @@ impl OdmTestEnv {
let source = FakeS3Target::start_with_options(options).await?;
let mut rustfs = RustFSTestEnvironment::new().await?;
rustfs
.start_rustfs_server_with_env(vec![], &[(ODM_MODULE_SWITCH_ENV, "true")])
.start_rustfs_server_with_env(vec![], &[(ODM_MODULE_SWITCH_ENV, "true"), (ALLOW_LOOPBACK_SOURCE_ENV, "true")])
.await?;
let client = rustfs.create_s3_client();
Ok(Self { rustfs, source, client })
@@ -412,6 +431,52 @@ impl OdmTestEnv {
assert_eq!(body.as_ref(), expected, "{bucket}/{key} local content mismatch");
}
/// Polls the listing until `key` is present locally or `timeout` elapses
/// (background pulls land after the response that triggered them).
pub async fn wait_local_listed(&self, bucket: &str, key: &str, timeout: Duration) -> Result<bool, BoxError> {
let deadline = Instant::now() + timeout;
loop {
if self.local_key_listed(bucket, key).await? {
return Ok(true);
}
if Instant::now() >= deadline {
return Ok(false);
}
tokio::time::sleep(Duration::from_millis(100)).await;
}
}
/// Raw signed `GET /{bucket}/{key}` against the RustFS under test.
pub async fn raw_get(&self, bucket: &str, key: &str) -> Result<RawResponse, BoxError> {
let url = format!("{}/{bucket}/{key}", self.rustfs.url);
let response =
signed_request(http::Method::GET, &url, &self.rustfs.access_key, &self.rustfs.secret_key, None, None).await?;
Ok(RawResponse {
status: response.status().as_u16(),
headers: response.headers().clone(),
body: response.bytes().await?,
})
}
/// Waits until the runtime consults the source for `bucket`: a config
/// install is applied asynchronously after the admin call returns. The
/// probe is a HEAD on a key that exists nowhere, so nothing is pulled and
/// only that key enters the negative cache.
pub async fn wait_until_source_consulted(&self, bucket: &str) -> Result<(), BoxError> {
const PROBE_KEY: &str = "_odm-readiness-probe";
let deadline = Instant::now() + Duration::from_secs(30);
loop {
let _ = self.client.head_object().bucket(bucket).key(PROBE_KEY).send().await;
if self.source.count_requests(Operation::HeadObject, PROBE_KEY) > 0 {
return Ok(());
}
if Instant::now() >= deadline {
return Err(format!("on-demand migration runtime for {bucket} did not consult the source in time").into());
}
tokio::time::sleep(Duration::from_millis(50)).await;
}
}
/// Panics if `key` is listed locally.
pub async fn assert_local_absent(&self, bucket: &str, key: &str) {
assert!(
@@ -0,0 +1,265 @@
// Copyright 2024 RustFS Team
//
// Licensed under the Apache License, Version 2.0 (the "License");
// you may not use this file except in compliance with the License.
// You may obtain a copy of the License at
//
// http://www.apache.org/licenses/LICENSE-2.0
//
// Unless required by applicable law or agreed to in writing, software
// distributed under the License is distributed on an "AS IS" BASIS,
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
// See the License for the specific language governing permissions and
// limitations under the License.
//! Basic GET read-through scenarios (rustfs/backlog#2156): inline pull and
//! local persistence, large-object passthrough with background backfill,
//! Range passthrough, source 404, `versionId` reads, and a disabled bucket.
//! Every source-side expectation is asserted on the fake source's journal.
use super::common::{BoxError, OdmSourceSpec, OdmTestEnv, SeedObject};
use crate::fake_s3_target::{BucketMode, Operation};
use aws_sdk_s3::error::ProvideErrorMetadata;
use aws_sdk_s3::types::{BucketVersioningStatus, VersioningConfiguration};
use bytes::Bytes;
use std::time::Duration;
type TestResult = Result<(), BoxError>;
const SOURCE_BUCKET: &str = "odm-get-source";
const ODM_RESPONSE_HEADER: &str = "x-rustfs-on-demand-migration";
/// Background pulls run after the response; generous for a loaded CI host.
const BACKFILL_WAIT: Duration = Duration::from_secs(60);
/// Position-dependent payload so a misaligned or truncated copy is caught.
fn payload(len: usize) -> Bytes {
(0..len).map(|index| (index % 251) as u8).collect::<Vec<u8>>().into()
}
/// RustFS with `local_bucket` migrating from `SOURCE_BUCKET` on the fake
/// source (unversioned, like a plain migration source); `adjust` tweaks the
/// policy before it is installed. Returns once the runtime consults the
/// source.
async fn configured_env(local_bucket: &str, adjust: impl FnOnce(&mut OdmSourceSpec)) -> Result<OdmTestEnv, BoxError> {
let env = OdmTestEnv::start().await?;
env.source.create_bucket_with_mode(SOURCE_BUCKET, BucketMode::Unversioned);
env.rustfs.create_test_bucket(local_bucket).await?;
let mut spec = env.fake_source_spec(SOURCE_BUCKET);
adjust(&mut spec);
let response = env.configure_source(local_bucket, &spec).await?;
assert_eq!(response.status, 200, "configure on-demand migration: {}", response.body);
env.wait_until_source_consulted(local_bucket).await?;
Ok(env)
}
fn source_get_ranges(env: &OdmTestEnv, key: &str) -> Vec<Option<String>> {
env.source
.requests()
.into_iter()
.filter(|record| record.operation == Operation::GetObject && record.key.as_deref() == Some(key))
.map(|record| record.range)
.collect()
}
#[tokio::test]
async fn get_miss_pulls_inline_and_serves_locally_afterwards() -> TestResult {
let bucket = "odm-get-inline";
let env = configured_env(bucket, |_| {}).await?;
let key = "inline/report.bin";
let body = payload(200 * 1024);
let etag = env
.seed_source(SOURCE_BUCKET, &[SeedObject::new(key, body.clone())])
.remove(0);
let quoted_etag = format!("\"{etag}\"");
let first = env.raw_get(bucket, key).await?;
assert_eq!(first.status, 200, "{}", String::from_utf8_lossy(&first.body));
assert_eq!(first.header(ODM_RESPONSE_HEADER), Some("source"), "a source answer is marked");
assert_eq!(first.header("etag"), Some(quoted_etag.as_str()), "inline answers carry the source ETag");
assert_eq!(first.header("content-length"), Some(body.len().to_string().as_str()));
assert_eq!(first.header("accept-ranges"), Some("bytes"));
assert_eq!(first.body, body, "the client receives the source bytes");
assert_eq!(env.source.count_requests(Operation::GetObject, key), 1, "exactly one source GET");
assert_eq!(env.source.count_requests(Operation::HeadObject, key), 1);
assert!(
env.wait_local_listed(bucket, key, BACKFILL_WAIT).await?,
"the inline pull must store the object locally"
);
let second = env.raw_get(bucket, key).await?;
assert_eq!(second.status, 200);
assert_eq!(second.header(ODM_RESPONSE_HEADER), None, "a local hit carries no source marker");
assert_eq!(second.body, body, "the local copy is the source bytes");
assert_eq!(second.header("etag"), Some(quoted_etag.as_str()), "preserve_etag keeps the source ETag");
assert_eq!(
env.source.count_requests(Operation::GetObject, key),
1,
"the second GET is served locally"
);
assert_eq!(env.source.count_requests(Operation::HeadObject, key), 1);
Ok(())
}
#[tokio::test]
async fn get_large_object_streams_through_and_backfills_in_background() -> TestResult {
let bucket = "odm-get-large";
let env = configured_env(bucket, |spec| spec.policy.inline_max_bytes = 4096).await?;
let key = "large/archive.bin";
let body = payload(512 * 1024);
env.seed_source(SOURCE_BUCKET, &[SeedObject::new(key, body.clone())]);
let response = env.raw_get(bucket, key).await?;
assert_eq!(response.status, 200, "{}", String::from_utf8_lossy(&response.body));
assert_eq!(response.header(ODM_RESPONSE_HEADER), Some("source"));
assert_eq!(response.header("content-length"), Some(body.len().to_string().as_str()));
assert_eq!(response.body, body, "the passthrough streams the whole object");
assert!(
env.wait_local_listed(bucket, key, BACKFILL_WAIT).await?,
"the background pull must store the object locally"
);
env.assert_local_present(bucket, key, &body).await;
assert_eq!(
source_get_ranges(&env, key),
vec![None, None],
"one passthrough GET plus one background pull, both unranged"
);
Ok(())
}
#[tokio::test]
async fn get_range_streams_206_and_backfills_the_whole_object() -> TestResult {
let bucket = "odm-get-range";
let env = configured_env(bucket, |_| {}).await?;
let key = "range/video.bin";
let body = payload(100_000);
let etag = env
.seed_source(SOURCE_BUCKET, &[SeedObject::new(key, body.clone())])
.remove(0);
let response = env
.client
.get_object()
.bucket(bucket)
.key(key)
.range("bytes=10-19")
.send()
.await?;
assert_eq!(response.content_range(), Some("bytes 10-19/100000"), "the source's 206 is passed through");
assert_eq!(response.content_length(), Some(10));
assert_eq!(response.e_tag(), Some(format!("\"{etag}\"").as_str()));
assert_eq!(response.body.collect().await?.into_bytes(), body.slice(10..20));
assert_eq!(
source_get_ranges(&env, key),
vec![Some("bytes=10-19".to_string())],
"the Range is forwarded"
);
assert!(
env.wait_local_listed(bucket, key, BACKFILL_WAIT).await?,
"serve_and_backfill must pull the whole object"
);
env.assert_local_present(bucket, key, &body).await;
assert_eq!(
source_get_ranges(&env, key),
vec![Some("bytes=10-19".to_string()), None],
"the background pull fetches the whole object"
);
Ok(())
}
#[tokio::test]
async fn get_source_not_found_is_404_and_negative_cached() -> TestResult {
let bucket = "odm-get-missing";
let env = configured_env(bucket, |_| {}).await?;
let key = "missing/nowhere.bin";
for attempt in 1..=2 {
let err = env
.client
.get_object()
.bucket(bucket)
.key(key)
.send()
.await
.expect_err("a key missing on both sides is 404");
assert_eq!(err.code(), Some("NoSuchKey"), "attempt {attempt}: {err:?}");
}
assert_eq!(env.source.count_requests(Operation::GetObject, key), 0, "a source miss never pulls");
assert_eq!(
env.source.count_requests(Operation::HeadObject, key),
1,
"the second miss stops at the negative cache"
);
env.assert_local_absent(bucket, key).await;
Ok(())
}
#[tokio::test]
async fn get_with_version_id_does_not_consult_the_source() -> TestResult {
let bucket = "odm-get-versioned";
let env = configured_env(bucket, |_| {}).await?;
env.client
.put_bucket_versioning()
.bucket(bucket)
.versioning_configuration(
VersioningConfiguration::builder()
.status(BucketVersioningStatus::Enabled)
.build(),
)
.send()
.await?;
let key = "versioned/doc.bin";
let body = payload(1024);
env.seed_source(SOURCE_BUCKET, &[SeedObject::new(key, body.clone())]);
let err = env
.client
.get_object()
.bucket(bucket)
.key(key)
.version_id("11111111-2222-4333-8444-555555555555")
.send()
.await
.expect_err("a version read cannot be answered by the source");
assert!(
matches!(err.code(), Some("NoSuchVersion") | Some("NoSuchKey")),
"unexpected error: {err:?}"
);
assert_eq!(env.source.count_requests(Operation::HeadObject, key), 0);
assert_eq!(env.source.count_requests(Operation::GetObject, key), 0);
env.assert_local_absent(bucket, key).await;
// The same key without versionId is still migrated: the gate is per request.
let response = env.raw_get(bucket, key).await?;
assert_eq!(response.status, 200, "{}", String::from_utf8_lossy(&response.body));
assert_eq!(response.header(ODM_RESPONSE_HEADER), Some("source"));
assert_eq!(response.body, body);
assert_eq!(env.source.count_requests(Operation::GetObject, key), 1);
Ok(())
}
#[tokio::test]
async fn get_after_disable_does_not_consult_the_source() -> TestResult {
let bucket = "odm-get-disabled";
let env = configured_env(bucket, |_| {}).await?;
let key = "disabled/doc.bin";
env.seed_source(SOURCE_BUCKET, &[SeedObject::new(key, payload(1024))]);
let response = env.disable(bucket).await?;
assert_eq!(response.status, 204, "{}", response.body);
let err = env
.client
.get_object()
.bucket(bucket)
.key(key)
.send()
.await
.expect_err("a disabled bucket answers locally");
assert_eq!(err.code(), Some("NoSuchKey"), "{err:?}");
assert_eq!(env.source.count_requests(Operation::HeadObject, key), 0);
assert_eq!(env.source.count_requests(Operation::GetObject, key), 0);
env.assert_local_absent(bucket, key).await;
Ok(())
}
@@ -16,9 +16,10 @@
//!
//! `common` is the shared environment: one RustFS under test, one programmable
//! fake S3 source, admin-API wrappers, seeding and local-state assertions.
//! `harness_self_test` proves the harness itself; ODM behavior scenarios are
//! separate modules wired by later tasks.
//! `harness_self_test` proves the harness itself; `get_basic_test` covers the
//! GET read-through (rustfs/backlog#2156).
pub mod common;
mod get_basic_test;
mod harness_self_test;
File diff suppressed because it is too large Load Diff
+86 -1
View File
@@ -15,7 +15,9 @@
//! Cross-cutting helpers shared by the object use-case modules.
use super::*;
use crate::app::storage_api::object_usecase::bucket::on_demand_migration::{OdmStateError, PolicyConfig, SourceErrorPolicy};
use crate::app::storage_api::object_usecase::bucket::on_demand_migration::{
OdmStateError, PolicyConfig, SourceErrorPolicy, SourceHead,
};
pub(super) const RUSTFS_EXPECTED_CURRENT_VERSION_ID: &str = "x-rustfs-expected-current-version-id";
@@ -870,6 +872,21 @@ pub(crate) fn mark_on_demand_migration_response(headers: &mut HeaderMap) {
headers.insert(ON_DEMAND_MIGRATION_HEADER, ON_DEMAND_MIGRATION_SOURCE);
}
/// Evaluates the request's conditional headers (`If-Match`, `If-None-Match`,
/// `If-Modified-Since`, `If-Unmodified-Since`) against the source's view of
/// the object, with the same S3 semantics the local path applies to
/// `ObjectInfo` (rustfs/backlog#2156). Conditional headers are never
/// forwarded to the source: a 304/412 answered by the source would be
/// indistinguishable from a source failure.
pub(crate) fn odm_check_source_preconditions(headers: &HeaderMap, head: &SourceHead) -> S3Result<()> {
let info = ObjectInfo {
etag: head.etag.clone(),
mod_time: head.last_modified.map(OffsetDateTime::from),
..Default::default()
};
check_preconditions(headers, &info)
}
#[cfg(test)]
mod tests {
use super::*;
@@ -1859,4 +1876,72 @@ mod on_demand_migration_tests {
mark_on_demand_migration_response(&mut headers);
assert_eq!(headers.get("x-rustfs-on-demand-migration").and_then(|v| v.to_str().ok()), Some("source"));
}
fn conditional_source_head() -> SourceHead {
SourceHead {
etag: Some("0123456789abcdef0123456789abcdef".to_string()),
size: 7,
// Thu, 01 Jan 2026 00:00:00 GMT
last_modified: Some(std::time::SystemTime::UNIX_EPOCH + Duration::from_secs(1_767_225_600)),
..Default::default()
}
}
fn headers_with(name: http::header::HeaderName, value: &'static str) -> HeaderMap {
let mut headers = HeaderMap::new();
headers.insert(name, HeaderValue::from_static(value));
headers
}
#[test]
fn odm_source_preconditions_follow_local_semantics() {
let head = conditional_source_head();
assert!(odm_check_source_preconditions(&HeaderMap::new(), &head).is_ok());
// If-None-Match on the source ETag: 304 carrying the source ETag and Last-Modified.
let err = odm_check_source_preconditions(
&headers_with(http::header::IF_NONE_MATCH, "\"0123456789abcdef0123456789abcdef\""),
&head,
)
.expect_err("matching If-None-Match is 304");
assert_eq!(err.code(), &S3ErrorCode::NotModified);
assert_eq!(err.status_code(), Some(StatusCode::NOT_MODIFIED));
let echoed = err.headers().expect("304 echoes validators");
assert_eq!(
echoed.get("etag").and_then(|v| v.to_str().ok()),
Some("\"0123456789abcdef0123456789abcdef\"")
);
assert_eq!(
echoed.get("last-modified").and_then(|v| v.to_str().ok()),
Some("Thu, 01 Jan 2026 00:00:00 GMT")
);
assert!(odm_check_source_preconditions(&headers_with(http::header::IF_NONE_MATCH, "\"other\""), &head).is_ok());
// If-Match on another ETag: 412.
let err = odm_check_source_preconditions(&headers_with(http::header::IF_MATCH, "\"other\""), &head)
.expect_err("mismatching If-Match is 412");
assert_eq!(err.code(), &S3ErrorCode::PreconditionFailed);
assert!(
odm_check_source_preconditions(&headers_with(http::header::IF_MATCH, "\"0123456789abcdef0123456789abcdef\""), &head)
.is_ok()
);
// Date validators compare against the source Last-Modified.
let err = odm_check_source_preconditions(
&headers_with(http::header::IF_MODIFIED_SINCE, "Fri, 02 Jan 2026 00:00:00 GMT"),
&head,
)
.expect_err("not modified since a later date is 304");
assert_eq!(err.code(), &S3ErrorCode::NotModified);
let err = odm_check_source_preconditions(
&headers_with(http::header::IF_UNMODIFIED_SINCE, "Wed, 31 Dec 2025 00:00:00 GMT"),
&head,
)
.expect_err("modified since an earlier date is 412");
assert_eq!(err.code(), &S3ErrorCode::PreconditionFailed);
// A source without validators cannot fail a precondition.
let bare = SourceHead::default();
assert!(odm_check_source_preconditions(&headers_with(http::header::IF_MATCH, "\"other\""), &bare).is_ok());
}
}
+2 -2
View File
@@ -626,7 +626,7 @@ pub(crate) mod bucket {
pub(crate) mod on_demand_migration {
pub(crate) use crate::storage::storage_api::ecstore_bucket::on_demand_migration::source_client::{
SourceClient, SourceError, SourceHead,
SourceClient, SourceError, SourceGet, SourceHead,
};
#[cfg(test)]
pub(crate) use crate::storage::storage_api::ecstore_bucket::on_demand_migration::{
@@ -635,7 +635,7 @@ pub(crate) mod bucket {
};
pub(crate) use crate::storage::storage_api::ecstore_bucket::on_demand_migration::{
BucketOdmState, HeadPolicy, OdmLookup, OdmOp, OdmOutcome, OdmStateError, OnDemandMigrationSys, PolicyConfig,
SourceErrorPolicy,
PullError, PullLeader, PullOutcome, PullReason, PullSlot, RangeGetPolicy, SourceErrorPolicy, commit_inline,
};
}