fix(rpc): bind local mutations to the listener instance

This commit is contained in:
overtrue
2026-09-06 10:58:11 +08:00
parent 7f631ec378
commit 11c9fd64ce
15 changed files with 1408 additions and 157 deletions
+35 -2
View File
@@ -26,13 +26,14 @@
//! server is not ready rather than that another server's global context applies.
use super::global::{AppContext, get_global_app_context};
use crate::app::storage_api::context::ECStore;
use crate::app::storage_api::context::{BootstrapLocalTarget, ECStore, InstanceContext};
use std::sync::{Arc, OnceLock};
/// Late-bound, per-server handle to the application context.
#[derive(Default)]
pub struct ServerContextSlot {
app_context: OnceLock<Arc<AppContext>>,
bootstrap_target: Option<BootstrapLocalTarget>,
heal_topology_fingerprint: Arc<tokio::sync::OnceCell<String>>,
}
@@ -50,15 +51,47 @@ impl ServerContextSlot {
pub fn new() -> Arc<Self> {
Arc::new(Self {
app_context: OnceLock::new(),
bootstrap_target: None,
heal_topology_fingerprint: Arc::new(tokio::sync::OnceCell::new()),
})
}
/// Bind the listener to its foundation before it can accept requests.
pub fn with_instance_context(ctx: Arc<InstanceContext>) -> Arc<Self> {
Arc::new(Self {
bootstrap_target: Some(BootstrapLocalTarget::new(ctx)),
..Self::default()
})
}
/// Install this server's application context (once). Returns `false` if
/// the slot was already installed; the first installation wins, matching
/// the process-global singleton's `get_or_init` semantics.
pub fn install(&self, context: Arc<AppContext>) -> bool {
self.app_context.set(context).is_ok()
self.try_install(context).is_ok()
}
/// Claim the slot before any process-global application publication.
/// Repeated installation, even of the same Arc, is an explicit conflict.
pub fn try_install(&self, context: Arc<AppContext>) -> std::io::Result<()> {
if self
.bootstrap_target
.as_ref()
.is_some_and(|target| !target.is_for_store(&context.object_store()))
{
return Err(std::io::Error::new(
std::io::ErrorKind::InvalidInput,
"application context does not belong to this server foundation",
));
}
self.app_context.set(context).map_err(|_| {
std::io::Error::new(std::io::ErrorKind::AlreadyExists, "server application context is already installed")
})
}
/// Immutable, restricted startup capability; never resolves an ambient store.
pub fn bootstrap_target(&self) -> Option<BootstrapLocalTarget> {
self.bootstrap_target.clone()
}
/// This server's installed application context, if startup has completed.
+2 -2
View File
@@ -37,8 +37,8 @@ impl AppContext {
// also publishes to the process default (first server wins) so legacy
// free-function readers keep resolving the first server's context.
let context = Arc::new(AppContext::with_default_interfaces(store, iam, kms_interface));
publish_global_app_context(context.clone());
let _ = server_ctx.install(context);
server_ctx.try_install(context.clone())?;
publish_global_app_context(context);
Ok(())
}
}
+1 -1
View File
@@ -1259,7 +1259,7 @@ pub(crate) mod context {
pub(crate) use super::EndpointServerPools;
pub(crate) use super::bucket;
pub(crate) use super::runtime;
pub(crate) use crate::storage::storage_api::{ECStore, EndpointServerPools};
pub(crate) use crate::storage::storage_api::{BootstrapLocalTarget, ECStore, EndpointServerPools, InstanceContext};
#[cfg(test)]
pub(crate) use crate::storage::storage_api::{Endpoint, Endpoints, PoolEndpoints};
}
+3 -1
View File
@@ -36,7 +36,9 @@ use crate::server::{
};
use crate::storage_api::server::http as storage;
use crate::storage_api::server::http::rpc::InternodeRpcService;
#[cfg(test)]
use crate::storage_api::server::http::tonic_service::make_server;
use crate::storage_api::server::http::tonic_service::make_server_for_slot;
use crate::storage_api::server::http::{
ServerContextSlot, TONIC_RPC_PREFIX, normalize_tonic_rpc_audience, tonic_boot_epoch_challenge,
tonic_boot_epoch_response_headers, verify_tonic_rpc_signature_with_bootstrap,
@@ -1834,7 +1836,7 @@ fn process_connection(
// each service in the auth interceptor.
let rpc_max_message_size = rustfs_protos::internode_rpc_max_message_size();
let node_service = InterceptedService::new(
NodeServiceServer::new(make_server())
NodeServiceServer::new(make_server_for_slot(Arc::clone(&server_ctx)))
.max_decoding_message_size(rpc_max_message_size)
.max_encoding_message_size(rpc_max_message_size),
check_auth,
+1 -3
View File
@@ -124,9 +124,6 @@ pub(crate) async fn run_embedded_startup(args: EmbeddedStartupArgs) -> Result<Em
} else {
bootstrap_instance_ctx()
};
// This server's request-path context slot (backlog#1052 S2).
let server_ctx = ServerContextSlot::new();
let EmbeddedStartupConfig {
config,
identity,
@@ -151,6 +148,7 @@ pub(crate) async fn run_embedded_startup(args: EmbeddedStartupArgs) -> Result<Em
.await
.map_err(init_error)?;
let server_ctx = ServerContextSlot::with_instance_context(instance_ctx.clone());
let http_server = start_embedded_http_server(&config, listen_context.readiness.clone(), server_ctx.clone()).await?;
let shutdown_handle = http_server.shutdown_handle;
let bound_addr = http_server.bound_addr;
+1 -4
View File
@@ -141,10 +141,6 @@ async fn run(config: Config) -> Result<()> {
// the storage path explicitly (Phase 5 follow-up, backlog#1052); a future
// multi-instance server constructs its own context here instead.
let instance_ctx = bootstrap_instance_ctx();
// This server's request-path context slot (backlog#1052 S2): handed to the
// HTTP service now, installed once IAM bootstrap completes.
let server_ctx = ServerContextSlot::new();
let StartupListenContext {
readiness,
server_addr,
@@ -152,6 +148,7 @@ async fn run(config: Config) -> Result<()> {
} = init_startup_listen_context(&config, &instance_ctx).await?;
let endpoint_pools = init_startup_storage_foundation(&server_address, &config.volumes, &instance_ctx).await?;
let server_ctx = ServerContextSlot::with_instance_context(instance_ctx.clone());
let StartupHttpServers {
state_manager,
s3_shutdown_tx,
+446 -3
View File
@@ -28,9 +28,9 @@ use crate::storage::storage_api::rpc_consumer::node_service::{
SCANNER_PUBLICATION_LEASE_TTL_MS, SERVICE_SIGNAL_REFRESH_CONFIG, SERVICE_SIGNAL_RELOAD_DYNAMIC, StorageDiskRpcExt as _,
StorageResult, all_local_disk_path, find_local_disk_by_ref, reload_transition_tier_config,
};
use crate::storage::storage_api::runtime_sources_consumer::{EndpointServerPools, runtime_sources};
use crate::storage::storage_api::runtime_sources_consumer::{EndpointServerPools, ServerContextSlot, runtime_sources};
use crate::storage::storage_api::{
sign_tonic_rpc_response_proof, verify_tonic_canonical_body_digest, verify_tonic_mutation_body_digest,
BootstrapLocalTarget, sign_tonic_rpc_response_proof, verify_tonic_canonical_body_digest, verify_tonic_mutation_body_digest,
verify_tonic_mutation_body_digest_reject_unsigned,
};
use bytes::Bytes;
@@ -482,6 +482,13 @@ mod metrics;
pub struct NodeService {
local_peer: LocalPeerS3Client,
context: Option<Arc<runtime_sources::AppContext>>,
server_ctx: Option<Arc<ServerContextSlot>>,
}
enum LocalMutationTarget {
Ready(Arc<ECStore>),
Bootstrap(BootstrapLocalTarget),
Unbound,
}
impl std::fmt::Debug for NodeService {
@@ -507,7 +514,19 @@ pub fn make_server() -> NodeService {
pub fn make_server_for_context(context: Option<Arc<runtime_sources::AppContext>>) -> NodeService {
let local_peer = LocalPeerS3Client::new(None, None);
NodeService { local_peer, context }
NodeService {
local_peer,
context,
server_ctx: None,
}
}
pub(crate) fn make_server_for_slot(server_ctx: Arc<ServerContextSlot>) -> NodeService {
// Unrelated RPCs retain their existing context policy. Target mutations
// resolve exclusively through this listener slot on each request.
let mut service = make_server();
service.server_ctx = Some(server_ctx);
service
}
#[derive(Clone, Debug, Default)]
@@ -1074,6 +1093,24 @@ impl heal_control_service_server::HealControlService for HealControlRpcService {
}
impl NodeService {
fn local_mutation_target(&self) -> LocalMutationTarget {
if let Some(slot) = &self.server_ctx {
// Capture exactly once per request, not at connection acceptance.
// A captured Bootstrap request cannot upgrade across a later await.
if let Some(store) = slot.installed_object_store() {
LocalMutationTarget::Ready(store)
} else if let Some(target) = slot.bootstrap_target() {
LocalMutationTarget::Bootstrap(target)
} else {
LocalMutationTarget::Unbound
}
} else if let Some(context) = &self.context {
LocalMutationTarget::Ready(context.object_store())
} else {
LocalMutationTarget::Unbound
}
}
fn resolve_object_store(&self) -> Option<Arc<ECStore>> {
let context = self.context.clone().or_else(runtime_sources::current_app_context);
runtime_sources::current_object_store_handle_for_context(context.as_deref())
@@ -2680,6 +2717,7 @@ mod tests {
validate_admin_heal_control_start,
};
use crate::storage::rpc::node_service::heal::heal_topology_fingerprint;
use crate::storage::storage_api::ecstore_disk::DiskAPI as _;
use crate::storage::storage_api::rpc_consumer::node_service::{DiskError, HealBucketInfo};
use crate::storage::storage_api::set_tonic_canonical_body_digest;
use crate::storage::storage_api::{
@@ -4530,6 +4568,411 @@ mod tests {
assert!(rename_response.error.is_some());
}
struct TargetRpcFixture {
_root: tempfile::TempDir,
env: rustfs_test_utils::TestECStoreEnv,
instance: Arc<crate::storage::storage_api::InstanceContext>,
context: Arc<crate::runtime_sources::AppContext>,
iam: Arc<rustfs_iam::sys::IamSys<ObjectStore>>,
}
async fn target_rpc_fixture() -> TargetRpcFixture {
super::timeout(Duration::from_secs(90), async {
let root = tempfile::tempdir().expect("target RPC root");
let env = rustfs_test_utils::TestECStoreEnv::builder()
.base_dir(root.path())
.init_bucket_metadata(false)
.build()
.await;
ObjectStore::new(env.ecstore.clone())
.save_iam_config(serde_json::json!({"version": 1}), format!("{}/format.json", *IAM_CONFIG_PREFIX))
.await
.expect("seed real IAM format");
let iam = rustfs_iam::build_iam_sys(env.ecstore.clone())
.await
.expect("build fixture IAM");
let context = Arc::new(crate::runtime_sources::AppContext::with_default_interfaces(
env.ecstore.clone(),
iam.clone(),
Arc::new(KmsServiceManager::new()),
));
let instance = crate::storage::storage_api::bootstrap_instance_ctx();
assert!(
super::BootstrapLocalTarget::new(instance.clone()).is_for_store(&env.ecstore),
"the standard builder must use this exact instance context"
);
super::timeout(Duration::from_secs(10), async {
while env.ecstore.scanner_data_usage_publication_blocked().await {
tokio::task::yield_now().await;
}
})
.await
.expect("startup namespace commits drain before test");
TargetRpcFixture {
_root: root,
env,
instance,
context,
iam,
}
})
.await
.expect("bounded real fixture initialization")
}
async fn stage_target_rpc(fixture: &TargetRpcFixture) -> (super::DiskStore, rustfs_filemeta::FileInfo, Vec<u8>) {
use crate::storage::storage_api::ecstore_disk::{DiskAPI, ReadOptions};
let disk = fixture
.instance
.local_disk_map()
.read()
.await
.values()
.find_map(Clone::clone)
.expect("local target");
let mut fi = rustfs_filemeta::FileInfo::new("destination", 1, 0);
fi.erasure.index = 1;
fi.version_id = Some(Uuid::new_v4());
fi.mod_time = Some(OffsetDateTime::now_utc());
fi.size = 17;
fi.parts = vec![rustfs_filemeta::ObjectPartInfo {
number: 1,
size: 17,
actual_size: 17,
..Default::default()
}];
fi.data = Some(Bytes::from_static(b"target-rpc-inline"));
fi.set_inline_data();
disk.make_volume("target-rpc").await.expect("target volume");
disk.write_metadata("target-rpc", "target-rpc", "staged", fi.clone())
.await
.expect("stage real inline body");
let read = disk
.read_version(
"target-rpc",
"target-rpc",
"staged",
&fi.version_id.expect("version").to_string(),
&ReadOptions {
read_data: true,
..Default::default()
},
)
.await
.expect("read staged body before mutation");
assert_eq!(read.data, fi.data);
let before = tokio::fs::read(disk.path().join("target-rpc/staged/xl.meta"))
.await
.expect("staged bytes");
(disk, fi, before)
}
fn target_rename_request(disk: &super::DiskStore, fi: &rustfs_filemeta::FileInfo) -> Request<RenameDataRequest> {
let mut request = Request::new(RenameDataRequest {
disk: disk.endpoint().to_string(),
src_volume: "target-rpc".to_string(),
src_path: "staged".to_string(),
dst_volume: "target-rpc".to_string(),
dst_path: "destination".to_string(),
file_info: serde_json::to_string(fi).expect("real FileInfo JSON"),
..Default::default()
});
let body = rustfs_protos::canonical_rename_data_request_body(request.get_ref()).expect("canonical target body");
set_tonic_canonical_body_digest(&mut request, &body).expect("body digest");
// Direct-handler precondition only; this does not stand in for wire authentication.
mark_v2_authenticated(&mut request);
request
}
#[tokio::test]
async fn target_slot_rejects_mismatched_and_repeated_install_before_global_publication() {
let fixture = target_rpc_fixture().await;
assert!(
crate::runtime_sources::current_app_context().is_none(),
"requires a separate nextest process"
);
let wrong = super::ServerContextSlot::with_instance_context(crate::storage::storage_api::new_instance_ctx());
let error = crate::runtime_sources::AppContext::ensure_startup_after_iam(
fixture.env.ecstore.clone(),
Arc::new(KmsServiceManager::new()),
&wrong,
fixture.iam.clone(),
)
.expect_err("mismatched startup must fail");
assert_eq!(error.kind(), std::io::ErrorKind::InvalidInput);
assert!(wrong.installed_app_context().is_none());
assert!(
crate::runtime_sources::current_app_context().is_none(),
"failed install must not publish globally"
);
assert!(!wrong.install(fixture.context.clone()), "bool adapter cannot bypass identity checks");
let slot = super::ServerContextSlot::with_instance_context(fixture.instance.clone());
crate::runtime_sources::AppContext::ensure_startup_after_iam(
fixture.env.ecstore.clone(),
Arc::new(KmsServiceManager::new()),
&slot,
fixture.iam.clone(),
)
.expect("matching startup installation");
let installed = slot.installed_app_context().expect("installed A");
assert!(Arc::ptr_eq(
&crate::runtime_sources::current_app_context().expect("published A"),
&installed
));
assert_eq!(
slot.try_install(installed.clone())
.expect_err("same Arc is still a duplicate")
.kind(),
std::io::ErrorKind::AlreadyExists
);
assert!(!slot.install(installed.clone()));
assert!(Arc::ptr_eq(&slot.installed_app_context().expect("first winner retained"), &installed));
}
#[tokio::test]
async fn target_slot_captures_bootstrap_once_and_next_request_observes_ready() {
let fixture = target_rpc_fixture().await;
let (disk, fi, before) = stage_target_rpc(&fixture).await;
let slot = super::ServerContextSlot::with_instance_context(fixture.instance.clone());
let service = super::make_server_for_slot(slot.clone());
let captured = service.local_mutation_target();
slot.try_install(fixture.context.clone())
.expect("install after the request captures bootstrap");
let super::LocalMutationTarget::Bootstrap(target) = captured else {
panic!("pre-install request must capture bootstrap");
};
assert!(
target
.rename_local_data(
&disk.endpoint().to_string(),
("target-rpc", "staged"),
&fi,
("target-rpc", "destination"),
None
)
.await
.is_err(),
"captured request cannot acquire Ready privileges"
);
assert_eq!(
tokio::fs::read(disk.path().join("target-rpc/staged/xl.meta"))
.await
.expect("original source"),
before
);
assert!(!disk.path().join("target-rpc/destination").exists());
assert!(
matches!(service.local_mutation_target(), super::LocalMutationTarget::Ready(_)),
"the same service must read the installed slot for its next request"
);
let result = service
.rename_data(target_rename_request(&disk, &fi))
.await
.expect("ready handler")
.into_inner();
assert!(result.success, "{:?}", result.error);
}
#[tokio::test]
async fn target_unbound_slot_never_mutates_a_published_global_store() {
let fixture = target_rpc_fixture().await;
let (disk, fi, before) = stage_target_rpc(&fixture).await;
let published = crate::runtime_sources::publish_test_app_context(fixture.context.clone());
assert!(Arc::ptr_eq(&published, &fixture.context));
let service = super::make_server_for_slot(super::ServerContextSlot::new());
let result = service
.rename_data(target_rename_request(&disk, &fi))
.await
.expect("handler reply")
.into_inner();
assert!(!result.success);
assert!(result.error.is_some());
assert_eq!(
tokio::fs::read(disk.path().join("target-rpc/staged/xl.meta"))
.await
.expect("source remains"),
before
);
assert!(!disk.path().join("target-rpc/destination").exists());
assert!(!fixture.env.ecstore.scanner_data_usage_publication_blocked().await);
}
#[tokio::test]
async fn target_undo_rejects_force_delete_marker_before_mutation() {
let fixture = target_rpc_fixture().await;
let (disk, fi, before) = stage_target_rpc(&fixture).await;
let service = make_server_for_context(Some(fixture.context.clone()));
let opts = crate::storage::storage_api::ecstore_disk::DeleteOptions {
undo_write: true,
..Default::default()
};
let mut request = Request::new(DeleteVersionRequest {
disk: disk.endpoint().to_string(),
volume: "target-rpc".to_string(),
path: "staged".to_string(),
file_info: serde_json::to_string(&fi).expect("FileInfo"),
opts: serde_json::to_string(&opts).expect("opts"),
force_del_marker: true,
..Default::default()
});
let body = rustfs_protos::canonical_delete_version_request_body(request.get_ref()).expect("canonical undo body");
set_tonic_canonical_body_digest(&mut request, &body).expect("body digest");
mark_v2_authenticated(&mut request);
let result = service.delete_version(request).await.expect("handler reply").into_inner();
assert!(!result.success);
assert!(result.error.is_some());
assert_eq!(
tokio::fs::read(disk.path().join("target-rpc/staged/xl.meta"))
.await
.expect("source remains"),
before
);
assert!(!fixture.env.ecstore.scanner_data_usage_publication_blocked().await);
}
#[cfg(not(windows))]
#[tokio::test]
async fn target_handler_cancellation_retains_namespace_through_physical_rename() {
use crate::storage::storage_api::{
LocalPublicationPause, LocalPublicationStage,
ecstore_disk::{DiskAPI, ReadOptions},
};
let fixture = target_rpc_fixture().await;
let (disk, fi, _) = stage_target_rpc(&fixture).await;
let slot = super::ServerContextSlot::with_instance_context(fixture.instance.clone());
slot.try_install(fixture.context.clone()).expect("ready target");
let service = super::make_server_for_slot(slot);
let mut pause =
LocalPublicationPause::install(&disk, "target-rpc", "destination/xl.meta", LocalPublicationStage::PreparedRename)
.expect("install scoped physical pause");
let mut handler = Box::pin(service.rename_data(target_rename_request(&disk, &fi)));
super::timeout(Duration::from_secs(10), async {
tokio::select! {
result = &mut handler => panic!("handler completed before physical entry: {result:?}"),
entered = pause.entered() => entered.expect("physical executor entered"),
}
})
.await
.expect("bounded physical entry");
drop(handler);
assert!(
fixture.env.ecstore.scanner_data_usage_publication_blocked().await,
"dropping the actual target handler must not release its physical owner"
);
drop(pause);
super::timeout(Duration::from_secs(10), async {
while fixture.env.ecstore.scanner_data_usage_publication_blocked().await {
tokio::task::yield_now().await;
}
})
.await
.expect("physical owner must drain");
let read = disk
.read_version(
"target-rpc",
"target-rpc",
"destination",
&fi.version_id.expect("version").to_string(),
&ReadOptions {
read_data: true,
..Default::default()
},
)
.await
.expect("read real late commit");
assert_eq!(read.data, fi.data);
}
#[cfg(not(windows))]
#[tokio::test]
async fn target_undo_handler_cancellation_retains_owner_until_backup_restoration() {
use crate::storage::storage_api::{
LocalPublicationPause, LocalPublicationStage,
ecstore_disk::{DeleteOptions, DiskAPI, ReadOptions},
};
let fixture = target_rpc_fixture().await;
let (disk, fi, _) = stage_target_rpc(&fixture).await;
let mut old = fi.clone();
old.data = Some(Bytes::from_static(b"previous-rpc-body"));
assert_eq!(old.data.as_ref().expect("old body").len(), 17);
disk.write_metadata("target-rpc", "target-rpc", "destination", old.clone())
.await
.expect("old actual version");
let old_bytes = tokio::fs::read(disk.path().join("target-rpc/destination/xl.meta"))
.await
.expect("old metadata bytes");
let committed = fixture
.env
.ecstore
.rename_local_data(
&disk.endpoint().to_string(),
("target-rpc", "staged"),
&fi,
("target-rpc", "destination"),
None,
)
.await
.expect("real overwrite creates rollback backup");
let opts = DeleteOptions {
undo_write: true,
old_data_dir: Some(committed.rollback_data_dir.expect("real rollback backup")),
..Default::default()
};
let service = make_server_for_context(Some(fixture.context.clone()));
let mut request = Request::new(DeleteVersionRequest {
disk: disk.endpoint().to_string(),
volume: "target-rpc".to_string(),
path: "destination".to_string(),
file_info: serde_json::to_string(&fi).expect("FileInfo"),
opts: serde_json::to_string(&opts).expect("undo options"),
..Default::default()
});
let body = rustfs_protos::canonical_delete_version_request_body(request.get_ref()).expect("canonical undo body");
set_tonic_canonical_body_digest(&mut request, &body).expect("body digest");
mark_v2_authenticated(&mut request);
let mut pause = LocalPublicationPause::install(&disk, "target-rpc", "destination/xl.meta", LocalPublicationStage::Rename)
.expect("pause actual backup restoration");
let mut handler = Box::pin(service.delete_version(request));
super::timeout(Duration::from_secs(10), async {
tokio::select! {
result = &mut handler => panic!("undo completed before physical entry: {result:?}"),
entered = pause.entered() => entered.expect("physical restore entered"),
}
})
.await
.expect("bounded physical restore entry");
drop(handler);
assert!(fixture.env.ecstore.scanner_data_usage_publication_blocked().await);
drop(pause);
super::timeout(Duration::from_secs(10), async {
while fixture.env.ecstore.scanner_data_usage_publication_blocked().await {
tokio::task::yield_now().await;
}
})
.await
.expect("restore owner drains");
assert_eq!(
tokio::fs::read(disk.path().join("target-rpc/destination/xl.meta"))
.await
.expect("restored bytes"),
old_bytes
);
let read = disk
.read_version(
"target-rpc",
"target-rpc",
"destination",
&fi.version_id.expect("version").to_string(),
&ReadOptions {
read_data: true,
..Default::default()
},
)
.await
.expect("restored readable version");
assert_eq!(read.data, old.data);
}
#[tokio::test]
async fn rename_data_same_uuid_uses_captured_instance_instead_of_global_disk() {
use crate::storage::storage_api::{
+131 -126
View File
@@ -12,7 +12,7 @@
// See the License for the specific language governing permissions and
// limitations under the License.
use super::NodeService;
use super::{LocalMutationTarget, NodeService};
use crate::storage::storage_api::rpc_consumer::node_service::{
BatchReadVersionReq, BatchReadVersionResp, DeleteOptions, DiskError, DiskInfoOptions, FileInfoVersions, ReadMultipleReq,
ReadMultipleResp, ReadOptions, StorageDiskRpcExt as _, UpdateMetadataOpts, validate_batch_read_version_item_count,
@@ -39,6 +39,46 @@ use tonic::{Request, Response, Status};
use tracing::debug;
use uuid::Uuid;
impl LocalMutationTarget {
async fn rename_local_data(
&self,
disk_ref: &str,
source: (&str, &str),
fi: &FileInfo,
destination: (&str, &str),
scanner_token: Option<Uuid>,
) -> Result<RenameDataResp, DiskError> {
match self {
Self::Ready(store) => {
store
.rename_local_data(disk_ref, source, fi, destination, scanner_token)
.await
}
Self::Bootstrap(target) => {
target
.rename_local_data(disk_ref, source, fi, destination, scanner_token)
.await
}
Self::Unbound => Err(DiskError::other("target disk instance is unavailable")),
}
}
async fn undo_local_write(
&self,
disk_ref: &str,
volume: &str,
path: &str,
fi: FileInfo,
opts: DeleteOptions,
) -> Result<(), DiskError> {
match self {
Self::Ready(store) => store.undo_local_write(disk_ref, volume, path, fi, opts).await,
Self::Bootstrap(target) => target.undo_local_write(disk_ref, volume, path, fi, opts).await,
Self::Unbound => Err(DiskError::other("target disk instance is unavailable")),
}
}
}
/// Initial capacity hint (bytes) for typical small msgpack requests and responses.
const MSGPACK_ENCODE_CAPACITY_HINT: usize = 512;
const FILE_INFO_MSGPACK_ENCODE_CAPACITY_HINT: usize = 1024;
@@ -670,55 +710,59 @@ impl NodeService {
"delete_version",
)?;
let request = request.into_inner();
if let Some(disk) = self.find_disk(&request.disk).await {
let file_info = match decode_msgpack_or_json::<FileInfo>(&request.file_info_bin, &request.file_info, "FileInfo") {
Ok(file_info) => file_info,
Err(err) => {
return Ok(Response::new(DeleteVersionResponse {
success: false,
raw_file_info: "".to_string(),
error: Some(DiskError::other(format!("decode FileInfo failed: {err}")).into()),
}));
}
};
let opts = match decode_msgpack_or_json::<DeleteOptions>(&request.opts_bin, &request.opts, "DeleteOptions") {
Ok(opts) => opts,
Err(err) => {
return Ok(Response::new(DeleteVersionResponse {
success: false,
raw_file_info: "".to_string(),
error: Some(DiskError::other(format!("decode DeleteOptions failed: {err}")).into()),
}));
}
};
match disk
.delete_version(&request.volume, &request.path, file_info, request.force_del_marker, opts)
let file_info = match decode_msgpack_or_json::<FileInfo>(&request.file_info_bin, &request.file_info, "FileInfo") {
Ok(file_info) => file_info,
Err(err) => {
return Ok(Response::new(DeleteVersionResponse {
success: false,
raw_file_info: "".to_string(),
error: Some(DiskError::other(format!("decode FileInfo failed: {err}")).into()),
}));
}
};
let opts = match decode_msgpack_or_json::<DeleteOptions>(&request.opts_bin, &request.opts, "DeleteOptions") {
Ok(opts) => opts,
Err(err) => {
return Ok(Response::new(DeleteVersionResponse {
success: false,
raw_file_info: "".to_string(),
error: Some(DiskError::other(format!("decode DeleteOptions failed: {err}")).into()),
}));
}
};
let result = if opts.undo_write {
if request.force_del_marker {
Err(DiskError::other("undo_write cannot force a delete marker"))
} else {
let target = self.local_mutation_target();
target
.undo_local_write(&request.disk, &request.volume, &request.path, file_info, opts)
.await
}
} else if let Some(disk) = self.find_disk(&request.disk).await {
disk.delete_version(&request.volume, &request.path, file_info, request.force_del_marker, opts)
.await
{
Ok(raw_file_info) => match serde_json::to_string(&raw_file_info) {
Ok(raw_file_info) => Ok(Response::new(DeleteVersionResponse {
success: true,
raw_file_info,
error: None,
})),
Err(err) => Ok(Response::new(DeleteVersionResponse {
success: false,
raw_file_info: "".to_string(),
error: Some(DiskError::other(format!("encode data failed: {err}")).into()),
})),
},
} else {
Err(DiskError::other("cannot find disk"))
};
match result {
Ok(raw_file_info) => match serde_json::to_string(&raw_file_info) {
Ok(raw_file_info) => Ok(Response::new(DeleteVersionResponse {
success: true,
raw_file_info,
error: None,
})),
Err(err) => Ok(Response::new(DeleteVersionResponse {
success: false,
raw_file_info: "".to_string(),
error: Some(err.into()),
error: Some(DiskError::other(format!("encode data failed: {err}")).into()),
})),
}
} else {
Ok(Response::new(DeleteVersionResponse {
},
Err(err) => Ok(Response::new(DeleteVersionResponse {
success: false,
raw_file_info: "".to_string(),
error: Some(DiskError::other("cannot find disk".to_string()).into()),
}))
error: Some(err.into()),
})),
}
}
@@ -1206,98 +1250,59 @@ impl NodeService {
"rename_data",
)?;
let request = request.into_inner();
if let Some(disk) = self.find_disk(&request.disk).await {
let decoded_file_info = match decode_rename_data_request_file_info(&request.file_info_bin, &request.file_info) {
Ok(file_info) => file_info,
Err(err) => {
return Ok(Response::new(RenameDataResponse {
success: false,
rename_data_resp: String::new(),
rename_data_resp_bin: Vec::new().into(),
error: Some(DiskError::other(format!("decode FileInfo failed: {err}")).into()),
}));
}
};
let scanner_publication_lease_token = if request.scanner_publication_lease_token.is_empty() {
None
} else {
let token = Uuid::from_slice(&request.scanner_publication_lease_token)
.map_err(|_| Status::invalid_argument("scanner publication lease token must be a UUID"))?;
if token.is_nil() {
return Err(Status::invalid_argument("scanner publication lease token must not be nil"));
}
Some(token)
};
// The target owns this read guard. It must span the complete
// disk rename, not merely the preflight, so a movement transition
// cannot restart after validation and before rename linearization.
let scanner_publication_lease_guard: Option<Arc<dyn Send + Sync>> =
if let Some(token) = scanner_publication_lease_token {
let Some(store) = self.resolve_object_store() else {
return Ok(Response::new(RenameDataResponse {
success: false,
rename_data_resp: String::new(),
rename_data_resp_bin: Vec::new().into(),
error: Some(DiskError::other("scanner publication lease owner is unavailable").into()),
}));
};
match store.acquire_scanner_publication_lease_guard(token).await {
Ok(guard) => Some(Arc::new(guard)),
Err(err) => {
return Ok(Response::new(RenameDataResponse {
success: false,
rename_data_resp: String::new(),
rename_data_resp_bin: Vec::new().into(),
error: Some(DiskError::other(err.to_string()).into()),
}));
}
}
} else {
None
};
let request_decoded_from_msgpack = decoded_file_info.from_msgpack;
match disk
.rename_data_borrowed_with_fence_and_guard(
&request.src_volume,
&request.src_path,
&decoded_file_info.value,
&request.dst_volume,
&request.dst_path,
scanner_publication_lease_token,
scanner_publication_lease_guard,
)
.await
{
Ok(rename_data_resp) => {
match encode_rename_data_response_payloads(&rename_data_resp, request_decoded_from_msgpack) {
Ok((rename_data_resp, rename_data_resp_bin)) => Ok(Response::new(RenameDataResponse {
success: true,
rename_data_resp,
rename_data_resp_bin: rename_data_resp_bin.into(),
error: None,
})),
Err(err) => Ok(Response::new(RenameDataResponse {
success: false,
rename_data_resp: String::new(),
rename_data_resp_bin: Vec::new().into(),
error: Some(err.into()),
})),
}
}
let target = self.local_mutation_target();
let decoded_file_info = match decode_rename_data_request_file_info(&request.file_info_bin, &request.file_info) {
Ok(file_info) => file_info,
Err(err) => {
return Ok(Response::new(RenameDataResponse {
success: false,
rename_data_resp: String::new(),
rename_data_resp_bin: Vec::new().into(),
error: Some(DiskError::other(format!("decode FileInfo failed: {err}")).into()),
}));
}
};
let scanner_publication_lease_token = if request.scanner_publication_lease_token.is_empty() {
None
} else {
let token = Uuid::from_slice(&request.scanner_publication_lease_token)
.map_err(|_| Status::invalid_argument("scanner publication lease token must be a UUID"))?;
if token.is_nil() {
return Err(Status::invalid_argument("scanner publication lease token must not be nil"));
}
Some(token)
};
let request_decoded_from_msgpack = decoded_file_info.from_msgpack;
match target
.rename_local_data(
&request.disk,
(&request.src_volume, &request.src_path),
&decoded_file_info.value,
(&request.dst_volume, &request.dst_path),
scanner_publication_lease_token,
)
.await
{
Ok(rename_data_resp) => match encode_rename_data_response_payloads(&rename_data_resp, request_decoded_from_msgpack) {
Ok((rename_data_resp, rename_data_resp_bin)) => Ok(Response::new(RenameDataResponse {
success: true,
rename_data_resp,
rename_data_resp_bin: rename_data_resp_bin.into(),
error: None,
})),
Err(err) => Ok(Response::new(RenameDataResponse {
success: false,
rename_data_resp: String::new(),
rename_data_resp_bin: Vec::new().into(),
error: Some(err.into()),
})),
}
} else {
Ok(Response::new(RenameDataResponse {
},
Err(err) => Ok(Response::new(RenameDataResponse {
success: false,
rename_data_resp: String::new(),
rename_data_resp_bin: Vec::new().into(),
error: Some(DiskError::other("cannot find disk".to_string()).into()),
}))
error: Some(err.into()),
})),
}
}
+8 -3
View File
@@ -377,10 +377,12 @@ pub(crate) mod timeout_wrapper_consumer {
}
pub(crate) mod tonic_service_consumer {
#[cfg(test)]
pub(crate) use super::super::tonic_service::make_server;
#[cfg(test)]
pub(crate) use super::super::tonic_service::{heal_topology_fingerprint, make_heal_control_server_for_source};
pub(crate) use super::super::tonic_service::{
make_heal_control_server_with_cache, make_scanner_control_server, make_server, make_tier_mutation_control_server,
make_heal_control_server_with_cache, make_scanner_control_server, make_server_for_slot, make_tier_mutation_control_server,
};
}
@@ -600,8 +602,8 @@ pub(crate) mod ecstore_storage {
#[cfg(test)]
pub(crate) use rustfs_ecstore::api::storage::init_local_disks;
pub(crate) use rustfs_ecstore::api::storage::{
ECStore, SCANNER_PUBLICATION_LEASE_TTL_MS, ScannerDataMovementPauseStatus, all_local_disk, all_local_disk_path,
find_local_disk_by_ref, init_local_disks_with_instance_ctx, init_lock_clients,
BootstrapLocalTarget, ECStore, SCANNER_PUBLICATION_LEASE_TTL_MS, ScannerDataMovementPauseStatus, all_local_disk,
all_local_disk_path, find_local_disk_by_ref, init_local_disks_with_instance_ctx, init_lock_clients,
prewarm_local_disk_id_map_with_instance_ctx,
};
}
@@ -677,6 +679,9 @@ type EcstoreReplicationStats = ecstore_bucket::replication::ReplicationStats;
pub(crate) type DynReplicationPool = StorageReplicationPoolHandle;
pub(crate) type DynReader = ecstore_rio::DynReader;
pub(crate) type ECStore = ecstore_storage::ECStore;
pub(crate) type BootstrapLocalTarget = ecstore_storage::BootstrapLocalTarget;
#[cfg(all(test, not(windows)))]
pub(crate) use rustfs_ecstore::api::disk::{LocalPublicationPause, LocalPublicationStage};
pub(crate) type Endpoint = ecstore_disk::endpoint::Endpoint;
#[cfg(test)]
pub(crate) type Endpoints = ecstore_layout::Endpoints;
+1 -1
View File
@@ -13,8 +13,8 @@
// limitations under the License.
pub(crate) use crate::storage::rpc::node_service::make_heal_control_server_with_cache;
pub(crate) use crate::storage::rpc::node_service::make_scanner_control_server;
#[cfg(test)]
pub(crate) use crate::storage::rpc::node_service::{heal::heal_topology_fingerprint, make_heal_control_server_for_source};
pub(crate) use crate::storage::rpc::node_service::{make_scanner_control_server, make_server_for_slot};
pub use crate::storage::rpc::{make_heal_control_server, make_server, make_tier_mutation_control_server};
pub type NodeService = crate::storage::rpc::NodeService;
+4 -1
View File
@@ -171,12 +171,15 @@ pub(crate) mod server {
}
pub(crate) mod tonic_service {
#[cfg(test)]
pub(crate) use crate::storage::storage_api::tonic_service_consumer::make_server;
#[cfg(test)]
pub(crate) use crate::storage::storage_api::tonic_service_consumer::{
heal_topology_fingerprint, make_heal_control_server_for_source,
};
pub(crate) use crate::storage::storage_api::tonic_service_consumer::{
make_heal_control_server_with_cache, make_scanner_control_server, make_server, make_tier_mutation_control_server,
make_heal_control_server_with_cache, make_scanner_control_server, make_server_for_slot,
make_tier_mutation_control_server,
};
}
}