From 65d1605c09cc413f8bcf9f0ebc4fd3e43c2d5790 Mon Sep 17 00:00:00 2001 From: barry Date: Sat, 1 Aug 2026 15:55:30 +0800 Subject: [PATCH 1/3] refactor(server): extract data mover crate --- Cargo.lock | 41 +++++++++-- Cargo.toml | 3 + crates/server/curvine-data-mover/Cargo.toml | 42 ++++++++++++ .../curvine-data-mover/src/common/mod.rs | 20 ++++++ .../src/common/ufs_client.rs | 2 +- .../src/common/ufs_factory.rs | 4 +- .../src/common/ufs_manager.rs | 2 +- .../src}/job_worker_client.rs | 0 crates/server/curvine-data-mover/src/lib.rs | 23 +++++++ .../curvine-data-mover/src}/rpc_context.rs | 0 .../src/transfer/backend.rs | 0 .../src/transfer/cluster_cache.rs | 2 +- .../src/transfer/handler.rs | 0 .../src/transfer/job_snapshot.rs | 0 .../src/transfer/memory_store.rs | 0 .../src/transfer/metrics.rs | 0 .../curvine-data-mover/src/transfer/mod.rs | 67 ++++++++++++++++++ .../src/transfer/mysql_store.rs | 0 .../src/transfer/planner.rs | 2 +- .../src/transfer/router_handler.rs | 0 .../src/transfer/scheduler.rs | 0 .../src/transfer/service.rs | 0 .../src/transfer/sqlite_store.rs | 0 .../curvine-data-mover}/src/transfer/store.rs | 0 .../src/transfer/tests/planner_test.rs | 0 .../src/transfer/transfer_server.rs | 2 +- .../curvine-data-mover/src/worker/mod.rs | 15 ++++ .../src/worker/task/load_task_runner.rs | 10 +-- .../curvine-data-mover/src/worker/task/mod.rs | 28 ++++++++ .../src/worker/task/task_context.rs | 0 .../src/worker/task/task_manager.rs | 4 +- .../src/worker/task/task_store.rs | 0 curvine-master/Cargo.toml | 3 +- curvine-master/src/common/mod.rs | 3 +- curvine-master/src/master/job/job_manager.rs | 2 +- curvine-master/src/master/job/job_runner.rs | 2 +- curvine-master/src/master/job/mod.rs | 3 +- .../src/master/journal/ufs_loader.rs | 2 +- curvine-master/src/master/mod.rs | 3 +- curvine-server/Cargo.toml | 6 +- curvine-server/src/common/mod.rs | 7 +- curvine-server/src/test/mini_cluster.rs | 2 +- curvine-server/src/transfer/mod.rs | 68 +------------------ .../tests/load_task_runner_fault_test.rs | 2 +- curvine-web/Cargo.toml | 2 +- curvine-web/src/router/load_handler.rs | 2 +- curvine-worker/Cargo.toml | 4 +- curvine-worker/src/lib.rs | 4 +- .../src/worker/block/block_actor.rs | 2 +- .../src/worker/block/master_client.rs | 2 +- .../replication/worker_replication_manager.rs | 4 +- curvine-worker/src/worker/task/mod.rs | 15 +--- 52 files changed, 275 insertions(+), 130 deletions(-) create mode 100644 crates/server/curvine-data-mover/Cargo.toml create mode 100644 crates/server/curvine-data-mover/src/common/mod.rs rename {curvine-server => crates/server/curvine-data-mover}/src/common/ufs_client.rs (97%) rename {curvine-master => crates/server/curvine-data-mover}/src/common/ufs_factory.rs (95%) rename {curvine-server => crates/server/curvine-data-mover}/src/common/ufs_manager.rs (99%) rename {curvine-master/src/master/job => crates/server/curvine-data-mover/src}/job_worker_client.rs (100%) create mode 100644 crates/server/curvine-data-mover/src/lib.rs rename {curvine-master/src/master => crates/server/curvine-data-mover/src}/rpc_context.rs (100%) rename {curvine-server => crates/server/curvine-data-mover}/src/transfer/backend.rs (100%) rename {curvine-server => crates/server/curvine-data-mover}/src/transfer/cluster_cache.rs (99%) rename {curvine-server => crates/server/curvine-data-mover}/src/transfer/handler.rs (100%) rename {curvine-server => crates/server/curvine-data-mover}/src/transfer/job_snapshot.rs (100%) rename {curvine-server => crates/server/curvine-data-mover}/src/transfer/memory_store.rs (100%) rename {curvine-server => crates/server/curvine-data-mover}/src/transfer/metrics.rs (100%) create mode 100644 crates/server/curvine-data-mover/src/transfer/mod.rs rename {curvine-server => crates/server/curvine-data-mover}/src/transfer/mysql_store.rs (100%) rename {curvine-server => crates/server/curvine-data-mover}/src/transfer/planner.rs (99%) rename {curvine-server => crates/server/curvine-data-mover}/src/transfer/router_handler.rs (100%) rename {curvine-server => crates/server/curvine-data-mover}/src/transfer/scheduler.rs (100%) rename {curvine-server => crates/server/curvine-data-mover}/src/transfer/service.rs (100%) rename {curvine-server => crates/server/curvine-data-mover}/src/transfer/sqlite_store.rs (100%) rename {curvine-server => crates/server/curvine-data-mover}/src/transfer/store.rs (100%) rename {curvine-server => crates/server/curvine-data-mover}/src/transfer/tests/planner_test.rs (100%) rename {curvine-server => crates/server/curvine-data-mover}/src/transfer/transfer_server.rs (99%) create mode 100644 crates/server/curvine-data-mover/src/worker/mod.rs rename {curvine-worker => crates/server/curvine-data-mover}/src/worker/task/load_task_runner.rs (99%) create mode 100644 crates/server/curvine-data-mover/src/worker/task/mod.rs rename {curvine-worker => crates/server/curvine-data-mover}/src/worker/task/task_context.rs (100%) rename {curvine-worker => crates/server/curvine-data-mover}/src/worker/task/task_manager.rs (99%) rename {curvine-worker => crates/server/curvine-data-mover}/src/worker/task/task_store.rs (100%) diff --git a/Cargo.lock b/Cargo.lock index 7651f7170..1407efbd4 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -2168,6 +2168,34 @@ dependencies = [ "toml", ] +[[package]] +name = "curvine-data-mover" +version = "0.2.0" +dependencies = [ + "axum", + "bytes", + "crossbeam", + "curvine-client-core", + "curvine-common", + "curvine-job-client", + "curvine-ufs-api", + "curvine-unified-fs", + "curvine-web", + "dashmap 5.5.3", + "futures", + "log", + "mysql", + "once_cell", + "orpc", + "parking_lot", + "prost 0.11.9", + "rusqlite", + "serde", + "serde_json", + "tokio", + "uuid", +] + [[package]] name = "curvine-error" version = "0.2.0" @@ -2361,11 +2389,12 @@ version = "0.2.0" dependencies = [ "axum", "bytes", - "curvine-client", "curvine-common", + "curvine-data-mover", "curvine-error", "curvine-fault", "curvine-rocksdb", + "curvine-unified-fs", "curvine-web", "dashmap 5.5.3", "futures", @@ -2477,10 +2506,12 @@ dependencies = [ "chrono", "clap", "crossbeam", - "curvine-client", + "curvine-client-core", "curvine-common", + "curvine-data-mover", "curvine-error", "curvine-fault", + "curvine-job-client", "curvine-master", "curvine-rocksdb", "curvine-ufs-api", @@ -2667,8 +2698,8 @@ name = "curvine-web" version = "0.2.0" dependencies = [ "axum", - "curvine-client", "curvine-common", + "curvine-job-client", "log", "orpc", "serde", @@ -2684,10 +2715,10 @@ version = "0.2.0" dependencies = [ "axum", "bytes", - "curvine-client", + "curvine-client-core", "curvine-common", + "curvine-data-mover", "curvine-fault", - "curvine-master", "curvine-storage-local", "curvine-storage-spdk", "curvine-web", diff --git a/Cargo.toml b/Cargo.toml index 2d0cae238..4b0d612c2 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -32,6 +32,7 @@ members = [ "crates/adapters/curvine-ufs-opendal", "crates/adapters/curvine-ufs-oss-hdfs", "crates/adapters/curvine-hdfs-jni", + "crates/server/curvine-data-mover", "curvine-master", "curvine-worker", "curvine-common", @@ -74,6 +75,7 @@ default-members = [ "crates/infra/curvine-io", "crates/metadata/curvine-raft", "crates/adapters/curvine-storage-local", + "crates/server/curvine-data-mover", "curvine-master", "curvine-worker", "curvine-common", @@ -112,6 +114,7 @@ curvine-storage-spdk = { path = "crates/adapters/curvine-storage-spdk" } curvine-ufs-opendal = { path = "crates/adapters/curvine-ufs-opendal" } curvine-ufs-oss-hdfs = { path = "crates/adapters/curvine-ufs-oss-hdfs" } curvine-hdfs-jni = { path = "crates/adapters/curvine-hdfs-jni" } +curvine-data-mover = { path = "crates/server/curvine-data-mover" } curvine-fault = { path = "crates/infra/curvine-fault" } curvine-error = { path = "crates/common/curvine-error" } curvine-proto = { path = "crates/common/curvine-proto" } diff --git a/crates/server/curvine-data-mover/Cargo.toml b/crates/server/curvine-data-mover/Cargo.toml new file mode 100644 index 000000000..b9ac214b0 --- /dev/null +++ b/crates/server/curvine-data-mover/Cargo.toml @@ -0,0 +1,42 @@ +[package] +name = "curvine-data-mover" +version.workspace = true +edition.workspace = true +license.workspace = true + +[features] +default = [] +opendal = ["curvine-unified-fs/opendal"] +opendal-s3 = ["opendal", "curvine-unified-fs/opendal-s3"] +opendal-oss = ["opendal", "curvine-unified-fs/opendal-oss"] +opendal-gcs = ["opendal", "curvine-unified-fs/opendal-gcs"] +opendal-azblob = ["opendal", "curvine-unified-fs/opendal-azblob"] +opendal-cos = ["opendal", "curvine-unified-fs/opendal-cos"] +opendal-hdfs = ["opendal", "curvine-unified-fs/opendal-hdfs"] +opendal-webhdfs = ["opendal", "curvine-unified-fs/opendal-webhdfs"] +opendal-hdfs-native = ["opendal", "curvine-unified-fs/opendal-hdfs-native"] +oss-hdfs = ["curvine-unified-fs/oss-hdfs"] + +[dependencies] +orpc = { workspace = true } +curvine-common = { workspace = true } +curvine-client-core = { workspace = true } +curvine-job-client = { workspace = true } +curvine-unified-fs = { workspace = true } +curvine-ufs-api = { workspace = true } +curvine-web = { workspace = true } +prost = { workspace = true } +bytes = { workspace = true } +serde = { workspace = true } +serde_json = { workspace = true } +log = { workspace = true } +tokio = { workspace = true } +futures = { workspace = true } +dashmap = { workspace = true } +crossbeam = { workspace = true } +parking_lot = { workspace = true } +once_cell = { workspace = true } +axum = { workspace = true } +rusqlite = { workspace = true } +mysql = { workspace = true } +uuid = { workspace = true } diff --git a/crates/server/curvine-data-mover/src/common/mod.rs b/crates/server/curvine-data-mover/src/common/mod.rs new file mode 100644 index 000000000..6f11607c3 --- /dev/null +++ b/crates/server/curvine-data-mover/src/common/mod.rs @@ -0,0 +1,20 @@ +// Copyright 2025 OPPO. +// +// 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. + +pub mod ufs_client; +pub mod ufs_manager; + +mod ufs_factory; +pub use self::ufs_client::UfsClient; +pub use self::ufs_factory::UfsFactory; diff --git a/curvine-server/src/common/ufs_client.rs b/crates/server/curvine-data-mover/src/common/ufs_client.rs similarity index 97% rename from curvine-server/src/common/ufs_client.rs rename to crates/server/curvine-data-mover/src/common/ufs_client.rs index b61e1953a..164c11f75 100644 --- a/curvine-server/src/common/ufs_client.rs +++ b/crates/server/curvine-data-mover/src/common/ufs_client.rs @@ -14,12 +14,12 @@ #![allow(unused)] -use curvine_client::unified::{UfsFileSystem, UnifiedReader, UnifiedWriter}; use curvine_common::error::FsError; use curvine_common::fs::{CurvineURI, FileSystem, Path}; use curvine_common::state::FileStatus; use curvine_common::FsResult; use curvine_ufs_api::fs::ufs_context::UFSContext; +use curvine_unified_fs::{UfsFileSystem, UnifiedReader, UnifiedWriter}; use std::sync::Arc; #[derive(Clone)] diff --git a/curvine-master/src/common/ufs_factory.rs b/crates/server/curvine-data-mover/src/common/ufs_factory.rs similarity index 95% rename from curvine-master/src/common/ufs_factory.rs rename to crates/server/curvine-data-mover/src/common/ufs_factory.rs index 5e10780ca..48c8a773b 100644 --- a/curvine-master/src/common/ufs_factory.rs +++ b/crates/server/curvine-data-mover/src/common/ufs_factory.rs @@ -12,11 +12,11 @@ // See the License for the specific language governing permissions and // limitations under the License. -use crate::master::JobWorkerClient; -use curvine_client::unified::{MountValue, UfsFileSystem}; +use crate::JobWorkerClient; use curvine_common::conf::ClientConf; use curvine_common::state::{MountInfo, WorkerAddress}; use curvine_common::FsResult; +use curvine_unified_fs::{MountValue, UfsFileSystem}; use orpc::client::ClientFactory; use orpc::io::net::InetAddr; use orpc::runtime::Runtime; diff --git a/curvine-server/src/common/ufs_manager.rs b/crates/server/curvine-data-mover/src/common/ufs_manager.rs similarity index 99% rename from curvine-server/src/common/ufs_manager.rs rename to crates/server/curvine-data-mover/src/common/ufs_manager.rs index e65861066..e1f03fba9 100644 --- a/curvine-server/src/common/ufs_manager.rs +++ b/crates/server/curvine-data-mover/src/common/ufs_manager.rs @@ -13,7 +13,7 @@ // limitations under the License. use crate::common::UfsClient; -use curvine_client::file::FsClient; +use curvine_client_core::file::FsClient; use curvine_common::conf::{UfsConf, UfsConfBuilder}; use curvine_common::fs::{CurvineURI, Path}; use curvine_common::state::MountInfo; diff --git a/curvine-master/src/master/job/job_worker_client.rs b/crates/server/curvine-data-mover/src/job_worker_client.rs similarity index 100% rename from curvine-master/src/master/job/job_worker_client.rs rename to crates/server/curvine-data-mover/src/job_worker_client.rs diff --git a/crates/server/curvine-data-mover/src/lib.rs b/crates/server/curvine-data-mover/src/lib.rs new file mode 100644 index 000000000..5ae768949 --- /dev/null +++ b/crates/server/curvine-data-mover/src/lib.rs @@ -0,0 +1,23 @@ +// Copyright 2025 OPPO. +// +// 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. + +pub mod common; +pub mod transfer; +pub mod worker; + +mod job_worker_client; +pub use self::job_worker_client::JobWorkerClient; + +mod rpc_context; +pub use self::rpc_context::RpcContext; diff --git a/curvine-master/src/master/rpc_context.rs b/crates/server/curvine-data-mover/src/rpc_context.rs similarity index 100% rename from curvine-master/src/master/rpc_context.rs rename to crates/server/curvine-data-mover/src/rpc_context.rs diff --git a/curvine-server/src/transfer/backend.rs b/crates/server/curvine-data-mover/src/transfer/backend.rs similarity index 100% rename from curvine-server/src/transfer/backend.rs rename to crates/server/curvine-data-mover/src/transfer/backend.rs diff --git a/curvine-server/src/transfer/cluster_cache.rs b/crates/server/curvine-data-mover/src/transfer/cluster_cache.rs similarity index 99% rename from curvine-server/src/transfer/cluster_cache.rs rename to crates/server/curvine-data-mover/src/transfer/cluster_cache.rs index 760407e38..a77819baf 100644 --- a/curvine-server/src/transfer/cluster_cache.rs +++ b/crates/server/curvine-data-mover/src/transfer/cluster_cache.rs @@ -12,7 +12,7 @@ // See the License for the specific language governing permissions and // limitations under the License. -use curvine_client::file::CurvineFileSystem; +use curvine_client_core::file::CurvineFileSystem; use curvine_common::error::FsError; use curvine_common::fs::Path; use curvine_common::state::{MountInfo, TransferKind, WorkerInfo}; diff --git a/curvine-server/src/transfer/handler.rs b/crates/server/curvine-data-mover/src/transfer/handler.rs similarity index 100% rename from curvine-server/src/transfer/handler.rs rename to crates/server/curvine-data-mover/src/transfer/handler.rs diff --git a/curvine-server/src/transfer/job_snapshot.rs b/crates/server/curvine-data-mover/src/transfer/job_snapshot.rs similarity index 100% rename from curvine-server/src/transfer/job_snapshot.rs rename to crates/server/curvine-data-mover/src/transfer/job_snapshot.rs diff --git a/curvine-server/src/transfer/memory_store.rs b/crates/server/curvine-data-mover/src/transfer/memory_store.rs similarity index 100% rename from curvine-server/src/transfer/memory_store.rs rename to crates/server/curvine-data-mover/src/transfer/memory_store.rs diff --git a/curvine-server/src/transfer/metrics.rs b/crates/server/curvine-data-mover/src/transfer/metrics.rs similarity index 100% rename from curvine-server/src/transfer/metrics.rs rename to crates/server/curvine-data-mover/src/transfer/metrics.rs diff --git a/crates/server/curvine-data-mover/src/transfer/mod.rs b/crates/server/curvine-data-mover/src/transfer/mod.rs new file mode 100644 index 000000000..ab08e8a7e --- /dev/null +++ b/crates/server/curvine-data-mover/src/transfer/mod.rs @@ -0,0 +1,67 @@ +mod store; +pub use self::store::*; + +mod memory_store; +pub use self::memory_store::MemoryTransferStore; + +mod sqlite_store; +pub use self::sqlite_store::SqliteTransferStore; + +mod mysql_store; +pub use self::mysql_store::MysqlTransferStore; + +mod metrics; +pub use self::metrics::TransferMetrics; + +mod backend; +pub(crate) use self::backend::is_store_unavailable_error; +pub use self::backend::TransferStoreBackend; +pub(crate) use curvine_common::transfer::transfer_failure_message; + +mod cluster_cache; +pub use self::cluster_cache::ClusterMetadataCache; + +mod job_snapshot; +pub use self::job_snapshot::job_mount_snapshot; + +mod planner; +pub use self::planner::{PlannedTransfer, TransferPlanner}; + +#[cfg(test)] +#[path = "tests/planner_test.rs"] +mod planner_test; + +mod service; +pub use self::service::{progress_to_proto, task_summary_to_proto, TransferService}; + +mod scheduler; +pub use self::scheduler::TransferScheduler; + +mod handler; +pub use self::handler::TransferHandler; + +mod router_handler; +pub use self::router_handler::TransferRouterHandler; + +mod transfer_server; +pub use self::transfer_server::TransferServer; + +pub(crate) fn apply_task_report_progress( + summary: &mut curvine_common::state::TransferProgress, + previous: &curvine_common::state::TransferProgress, + current: &curvine_common::state::TransferProgress, + now_ms: i64, +) { + summary.loaded_size = summary + .loaded_size + .saturating_sub(previous.loaded_size) + .saturating_add(current.loaded_size) + .max(0); + summary.total_size = summary + .total_size + .saturating_sub(previous.total_size) + .saturating_add(current.total_size) + .max(0); + summary.update_time = now_ms; + summary.message = current.message.clone(); +} diff --git a/curvine-server/src/transfer/mysql_store.rs b/crates/server/curvine-data-mover/src/transfer/mysql_store.rs similarity index 100% rename from curvine-server/src/transfer/mysql_store.rs rename to crates/server/curvine-data-mover/src/transfer/mysql_store.rs diff --git a/curvine-server/src/transfer/planner.rs b/crates/server/curvine-data-mover/src/transfer/planner.rs similarity index 99% rename from curvine-server/src/transfer/planner.rs rename to crates/server/curvine-data-mover/src/transfer/planner.rs index 7f4cf63dc..0e71d824e 100644 --- a/curvine-server/src/transfer/planner.rs +++ b/crates/server/curvine-data-mover/src/transfer/planner.rs @@ -14,7 +14,7 @@ use crate::common::UfsFactory; use crate::transfer::{job_mount_snapshot, ClusterMetadataCache, TransferMetrics}; -use curvine_client::file::CurvineFileSystem; +use curvine_client_core::file::CurvineFileSystem; use curvine_common::conf::ClientConf; use curvine_common::error::FsError; use curvine_common::fs::{FileSystem, Path}; diff --git a/curvine-server/src/transfer/router_handler.rs b/crates/server/curvine-data-mover/src/transfer/router_handler.rs similarity index 100% rename from curvine-server/src/transfer/router_handler.rs rename to crates/server/curvine-data-mover/src/transfer/router_handler.rs diff --git a/curvine-server/src/transfer/scheduler.rs b/crates/server/curvine-data-mover/src/transfer/scheduler.rs similarity index 100% rename from curvine-server/src/transfer/scheduler.rs rename to crates/server/curvine-data-mover/src/transfer/scheduler.rs diff --git a/curvine-server/src/transfer/service.rs b/crates/server/curvine-data-mover/src/transfer/service.rs similarity index 100% rename from curvine-server/src/transfer/service.rs rename to crates/server/curvine-data-mover/src/transfer/service.rs diff --git a/curvine-server/src/transfer/sqlite_store.rs b/crates/server/curvine-data-mover/src/transfer/sqlite_store.rs similarity index 100% rename from curvine-server/src/transfer/sqlite_store.rs rename to crates/server/curvine-data-mover/src/transfer/sqlite_store.rs diff --git a/curvine-server/src/transfer/store.rs b/crates/server/curvine-data-mover/src/transfer/store.rs similarity index 100% rename from curvine-server/src/transfer/store.rs rename to crates/server/curvine-data-mover/src/transfer/store.rs diff --git a/curvine-server/src/transfer/tests/planner_test.rs b/crates/server/curvine-data-mover/src/transfer/tests/planner_test.rs similarity index 100% rename from curvine-server/src/transfer/tests/planner_test.rs rename to crates/server/curvine-data-mover/src/transfer/tests/planner_test.rs diff --git a/curvine-server/src/transfer/transfer_server.rs b/crates/server/curvine-data-mover/src/transfer/transfer_server.rs similarity index 99% rename from curvine-server/src/transfer/transfer_server.rs rename to crates/server/curvine-data-mover/src/transfer/transfer_server.rs index f9bb3650d..e12a59e6c 100644 --- a/curvine-server/src/transfer/transfer_server.rs +++ b/crates/server/curvine-data-mover/src/transfer/transfer_server.rs @@ -16,7 +16,7 @@ use std::sync::Arc; use std::time::Duration; use crate::common::UfsFactory; -use curvine_client::file::CurvineFileSystem; +use curvine_client_core::file::CurvineFileSystem; use curvine_common::conf::{ClusterConf, TransferStoreType}; use curvine_common::error::FsError; use curvine_web::server::{WebHandlerService, WebServer}; diff --git a/crates/server/curvine-data-mover/src/worker/mod.rs b/crates/server/curvine-data-mover/src/worker/mod.rs new file mode 100644 index 000000000..a63661e69 --- /dev/null +++ b/crates/server/curvine-data-mover/src/worker/mod.rs @@ -0,0 +1,15 @@ +// Copyright 2025 OPPO. +// +// 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. + +pub mod task; diff --git a/curvine-worker/src/worker/task/load_task_runner.rs b/crates/server/curvine-data-mover/src/worker/task/load_task_runner.rs similarity index 99% rename from curvine-worker/src/worker/task/load_task_runner.rs rename to crates/server/curvine-data-mover/src/worker/task/load_task_runner.rs index ff454c2d7..d9f6e2cb6 100644 --- a/curvine-worker/src/worker/task/load_task_runner.rs +++ b/crates/server/curvine-data-mover/src/worker/task/load_task_runner.rs @@ -15,10 +15,7 @@ use crate::common::UfsFactory; use crate::transfer::transfer_failure_message; use crate::worker::task::TaskContext; -use curvine_client::file::{CurvineFileSystem, FsReader}; -use curvine_client::rpc::JobMasterClient; -use curvine_client::rpc::TransferClient; -use curvine_client::unified::{UfsFileSystem, UnifiedReader, UnifiedWriter}; +use curvine_client_core::file::{CurvineFileSystem, FsReader}; use curvine_common::error::FsError; use curvine_common::fs::{FileSystem, Path, Reader, Writer}; use curvine_common::state::{ @@ -26,6 +23,9 @@ use curvine_common::state::{ SetAttrOptsBuilder, TRANSFER_TEMP_PATH_MARKER, }; use curvine_common::FsResult; +use curvine_job_client::JobMasterClient; +use curvine_job_client::TransferClient; +use curvine_unified_fs::{UfsFileSystem, UnifiedReader, UnifiedWriter}; use log::{debug, error, info, warn}; use orpc::common::{LocalTime, TimeSpent}; use orpc::err_box; @@ -805,8 +805,8 @@ fn xattr_equals(status: &FileStatus, key: &str, expected: &[u8]) -> bool { #[cfg(test)] mod tests { use super::rename_ufs_output; - use curvine_client::unified::UfsFileSystem; use curvine_common::fs::Path; + use curvine_unified_fs::UfsFileSystem; use orpc::runtime::{AsyncRuntime, RpcRuntime}; use std::collections::HashMap; use std::fs; diff --git a/crates/server/curvine-data-mover/src/worker/task/mod.rs b/crates/server/curvine-data-mover/src/worker/task/mod.rs new file mode 100644 index 000000000..0b0d023a2 --- /dev/null +++ b/crates/server/curvine-data-mover/src/worker/task/mod.rs @@ -0,0 +1,28 @@ +// Copyright 2025 OPPO. +// +// 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. + +//! Worker side data loading module +//! Responsible for loading data from external storage systems (such as S3) to local storage + +mod load_task_runner; +pub use self::load_task_runner::LoadTaskRunner; + +mod task_manager; +pub use self::task_manager::TaskManager; + +mod task_store; +pub use self::task_store::TaskStore; + +mod task_context; +pub use self::task_context::TaskContext; diff --git a/curvine-worker/src/worker/task/task_context.rs b/crates/server/curvine-data-mover/src/worker/task/task_context.rs similarity index 100% rename from curvine-worker/src/worker/task/task_context.rs rename to crates/server/curvine-data-mover/src/worker/task/task_context.rs diff --git a/curvine-worker/src/worker/task/task_manager.rs b/crates/server/curvine-data-mover/src/worker/task/task_manager.rs similarity index 99% rename from curvine-worker/src/worker/task/task_manager.rs rename to crates/server/curvine-data-mover/src/worker/task/task_manager.rs index f6f575be7..7298acf9a 100644 --- a/curvine-worker/src/worker/task/task_manager.rs +++ b/crates/server/curvine-data-mover/src/worker/task/task_manager.rs @@ -15,11 +15,11 @@ use crate::common::UfsFactory; use crate::worker::task::load_task_runner::LoadTaskRunner; use crate::worker::task::{TaskContext, TaskStore}; -use curvine_client::file::{CurvineFileSystem, FsContext}; -use curvine_client::rpc::TransferClient; +use curvine_client_core::file::{CurvineFileSystem, FsContext}; use curvine_common::conf::ClusterConf; use curvine_common::state::{JobTaskProgress, JobTaskState, LoadTaskInfo, TransferTaskReportInfo}; use curvine_common::FsResult; +use curvine_job_client::TransferClient; use dashmap::mapref::entry::Entry; use log::{debug, info, warn}; use orpc::runtime::{RpcRuntime, Runtime}; diff --git a/curvine-worker/src/worker/task/task_store.rs b/crates/server/curvine-data-mover/src/worker/task/task_store.rs similarity index 100% rename from curvine-worker/src/worker/task/task_store.rs rename to crates/server/curvine-data-mover/src/worker/task/task_store.rs diff --git a/curvine-master/Cargo.toml b/curvine-master/Cargo.toml index 7dbc20eb9..a91d894d1 100644 --- a/curvine-master/Cargo.toml +++ b/curvine-master/Cargo.toml @@ -13,7 +13,8 @@ orpc = { workspace = true } curvine-common = { workspace = true } curvine-rocksdb = { workspace = true } curvine-error = { workspace = true, features = ["axum-response"] } -curvine-client = { workspace = true, features = ["job-client", "unified", "opendal-s3"] } +curvine-data-mover = { workspace = true, features = ["opendal-s3"] } +curvine-unified-fs = { workspace = true, features = ["opendal-s3"] } curvine-web = { workspace = true } curvine-fault = { workspace = true } prost = { workspace = true } diff --git a/curvine-master/src/common/mod.rs b/curvine-master/src/common/mod.rs index 5ea697aa6..02013813b 100644 --- a/curvine-master/src/common/mod.rs +++ b/curvine-master/src/common/mod.rs @@ -12,5 +12,4 @@ // See the License for the specific language governing permissions and // limitations under the License. -mod ufs_factory; -pub use self::ufs_factory::UfsFactory; +pub use curvine_data_mover::common::UfsFactory; diff --git a/curvine-master/src/master/job/job_manager.rs b/curvine-master/src/master/job/job_manager.rs index 719a4e6c2..02948f491 100644 --- a/curvine-master/src/master/job/job_manager.rs +++ b/curvine-master/src/master/job/job_manager.rs @@ -16,7 +16,6 @@ use crate::common::UfsFactory; use crate::master::fs::MasterFilesystem; use crate::master::{JobStore, LoadJobRunner, MountManager}; use core::time::Duration; -use curvine_client::unified::MountValue; use curvine_common::conf::ClusterConf; use curvine_common::error::FsError; use curvine_common::executor::ScheduledExecutor; @@ -25,6 +24,7 @@ use curvine_common::state::{ JobStatus, JobTaskProgress, JobTaskState, LoadJobCommand, LoadJobResult, }; use curvine_common::FsResult; +use curvine_unified_fs::MountValue; use log::{debug, info, warn}; use orpc::common::LocalTime; use orpc::runtime::{LoopTask, RpcRuntime, Runtime}; diff --git a/curvine-master/src/master/job/job_runner.rs b/curvine-master/src/master/job/job_runner.rs index 83c90cf00..630df45d5 100644 --- a/curvine-master/src/master/job/job_runner.rs +++ b/curvine-master/src/master/job/job_runner.rs @@ -16,7 +16,6 @@ use crate::common::UfsFactory; use crate::master::fs::policy::ChooseContext; use crate::master::fs::MasterFilesystem; use crate::master::{JobContext, JobStore, TaskDetail}; -use curvine_client::unified::MountValue; use curvine_common::conf::ClientConf; use curvine_common::error::FsError; use curvine_common::fs::{FileSystem, Path}; @@ -26,6 +25,7 @@ use curvine_common::state::{ }; use curvine_common::utils::CommonUtils; use curvine_common::FsResult; +use curvine_unified_fs::MountValue; use dashmap::mapref::entry::Entry; use futures::future; use log::{debug, error, info, warn}; diff --git a/curvine-master/src/master/job/mod.rs b/curvine-master/src/master/job/mod.rs index a0023bc9b..eee89c62d 100644 --- a/curvine-master/src/master/job/mod.rs +++ b/curvine-master/src/master/job/mod.rs @@ -18,8 +18,7 @@ pub use crate::master::job::job_manager::JobManager; mod job_handler; pub use job_handler::JobHandler; -mod job_worker_client; -pub use self::job_worker_client::JobWorkerClient; +pub use curvine_data_mover::JobWorkerClient; mod job_store; pub use job_store::JobStore; diff --git a/curvine-master/src/master/journal/ufs_loader.rs b/curvine-master/src/master/journal/ufs_loader.rs index 18a7d7b02..b65741b75 100644 --- a/curvine-master/src/master/journal/ufs_loader.rs +++ b/curvine-master/src/master/journal/ufs_loader.rs @@ -16,12 +16,12 @@ use crate::master::journal::{ CompleteFileEntry, DeleteEntry, JournalEntry, MkdirEntry, RenameEntry, }; use crate::master::JobManager; -use curvine_client::unified::MountValue; use curvine_common::conf::JournalConf; use curvine_common::error::FsError; use curvine_common::fs::{FileSystem, Path}; use curvine_common::state::{JobTaskState, LoadJobCommand}; use curvine_common::FsResult; +use curvine_unified_fs::MountValue; use log::{info, warn}; use orpc::common::DurationUnit; use orpc::{err_box, CommonResult}; diff --git a/curvine-master/src/master/mod.rs b/curvine-master/src/master/mod.rs index b50616bfd..a0f32d0cf 100644 --- a/curvine-master/src/master/mod.rs +++ b/curvine-master/src/master/mod.rs @@ -45,8 +45,7 @@ pub use self::master_metrics::*; mod router_handler; pub use self::router_handler::*; -mod rpc_context; -pub use rpc_context::RpcContext; +pub use curvine_data_mover::RpcContext; pub mod mount; diff --git a/curvine-server/Cargo.toml b/curvine-server/Cargo.toml index fdf2bf09c..67fb6ed6b 100644 --- a/curvine-server/Cargo.toml +++ b/curvine-server/Cargo.toml @@ -8,7 +8,7 @@ license.workspace = true [features] default = [] -jni = ["curvine-client/opendal-hdfs"] +jni = ["curvine-data-mover/opendal-hdfs"] spdk = [ "curvine-worker/spdk", ] @@ -29,7 +29,9 @@ curvine-master = { workspace = true } curvine-worker = { workspace = true } curvine-error = { workspace = true, features = ["axum-response"] } curvine-fault = { workspace = true } -curvine-client = { workspace = true, features = ["job-client", "unified", "opendal-s3"] } +curvine-client-core = { workspace = true } +curvine-data-mover = { workspace = true, features = ["opendal-s3"] } +curvine-job-client = { workspace = true } curvine-ufs-api = { workspace = true } curvine-web = { workspace = true } prost = { workspace = true } diff --git a/curvine-server/src/common/mod.rs b/curvine-server/src/common/mod.rs index 722b74cea..3bfeb703d 100644 --- a/curvine-server/src/common/mod.rs +++ b/curvine-server/src/common/mod.rs @@ -12,9 +12,4 @@ // See the License for the specific language governing permissions and // limitations under the License. -pub mod ufs_manager; - -pub mod ufs_client; -pub use self::ufs_client::UfsClient; - -pub use curvine_master::UfsFactory; +pub use curvine_data_mover::common::*; diff --git a/curvine-server/src/test/mini_cluster.rs b/curvine-server/src/test/mini_cluster.rs index 1c71781bb..ba3c92538 100644 --- a/curvine-server/src/test/mini_cluster.rs +++ b/curvine-server/src/test/mini_cluster.rs @@ -16,7 +16,7 @@ use crate::master::fs::MasterFilesystem; use crate::master::replication::master_replication_manager::MasterReplicationManager; use crate::master::Master; use crate::worker::Worker; -use curvine_client::file::CurvineFileSystem; +use curvine_client_core::file::CurvineFileSystem; use curvine_common::conf::ClusterConf; use curvine_common::raft::{NodeId, RaftPeer}; use curvine_common::FsResult; diff --git a/curvine-server/src/transfer/mod.rs b/curvine-server/src/transfer/mod.rs index ab08e8a7e..f71489f5a 100644 --- a/curvine-server/src/transfer/mod.rs +++ b/curvine-server/src/transfer/mod.rs @@ -1,67 +1 @@ -mod store; -pub use self::store::*; - -mod memory_store; -pub use self::memory_store::MemoryTransferStore; - -mod sqlite_store; -pub use self::sqlite_store::SqliteTransferStore; - -mod mysql_store; -pub use self::mysql_store::MysqlTransferStore; - -mod metrics; -pub use self::metrics::TransferMetrics; - -mod backend; -pub(crate) use self::backend::is_store_unavailable_error; -pub use self::backend::TransferStoreBackend; -pub(crate) use curvine_common::transfer::transfer_failure_message; - -mod cluster_cache; -pub use self::cluster_cache::ClusterMetadataCache; - -mod job_snapshot; -pub use self::job_snapshot::job_mount_snapshot; - -mod planner; -pub use self::planner::{PlannedTransfer, TransferPlanner}; - -#[cfg(test)] -#[path = "tests/planner_test.rs"] -mod planner_test; - -mod service; -pub use self::service::{progress_to_proto, task_summary_to_proto, TransferService}; - -mod scheduler; -pub use self::scheduler::TransferScheduler; - -mod handler; -pub use self::handler::TransferHandler; - -mod router_handler; -pub use self::router_handler::TransferRouterHandler; - -mod transfer_server; -pub use self::transfer_server::TransferServer; - -pub(crate) fn apply_task_report_progress( - summary: &mut curvine_common::state::TransferProgress, - previous: &curvine_common::state::TransferProgress, - current: &curvine_common::state::TransferProgress, - now_ms: i64, -) { - summary.loaded_size = summary - .loaded_size - .saturating_sub(previous.loaded_size) - .saturating_add(current.loaded_size) - .max(0); - summary.total_size = summary - .total_size - .saturating_sub(previous.total_size) - .saturating_add(current.total_size) - .max(0); - summary.update_time = now_ms; - summary.message = current.message.clone(); -} +pub use curvine_data_mover::transfer::*; diff --git a/curvine-server/tests/load_task_runner_fault_test.rs b/curvine-server/tests/load_task_runner_fault_test.rs index f13c9dea1..5d5b30aa7 100644 --- a/curvine-server/tests/load_task_runner_fault_test.rs +++ b/curvine-server/tests/load_task_runner_fault_test.rs @@ -14,11 +14,11 @@ #![cfg(feature = "fault-injection")] -use curvine_client::rpc::JobMasterClient; use curvine_common::conf::ClusterConf; use curvine_common::fs::Path; use curvine_common::state::{JobTaskState, LoadJobCommand, MountOptions}; use curvine_fault::{FaultHttpController, FaultRuleBuilder, FaultRuntime, FaultTestSession}; +use curvine_job_client::JobMasterClient; use curvine_server::test::MiniCluster; use orpc::common::Utils; use orpc::runtime::RpcRuntime; diff --git a/curvine-web/Cargo.toml b/curvine-web/Cargo.toml index 2c72206c6..07e80b374 100644 --- a/curvine-web/Cargo.toml +++ b/curvine-web/Cargo.toml @@ -7,7 +7,7 @@ license.workspace = true # See more keys and their definitions at https://doc.rust-lang.org/cargo/reference/manifest.html [dependencies] -curvine-client = { workspace = true, features = ["job-client"] } +curvine-job-client = { workspace = true } axum = { workspace = true } tokio = { workspace = true } curvine-common = { workspace = true } diff --git a/curvine-web/src/router/load_handler.rs b/curvine-web/src/router/load_handler.rs index 5bbf7331c..20988ff9c 100644 --- a/curvine-web/src/router/load_handler.rs +++ b/curvine-web/src/router/load_handler.rs @@ -19,7 +19,7 @@ use axum::http::StatusCode; use axum::response::{IntoResponse, Response}; use axum::routing::{get, post}; use axum::Router; -use curvine_client::rpc::JobMasterClient; +use curvine_job_client::JobMasterClient; use log::{debug, info}; use serde::{Deserialize, Serialize}; use std::collections::HashMap; diff --git a/curvine-worker/Cargo.toml b/curvine-worker/Cargo.toml index 19e41f196..18607830f 100644 --- a/curvine-worker/Cargo.toml +++ b/curvine-worker/Cargo.toml @@ -23,8 +23,8 @@ fault-injection = [ [dependencies] orpc = { workspace = true } curvine-common = { workspace = true } -curvine-master = { workspace = true } -curvine-client = { workspace = true, features = ["job-client", "unified", "opendal-s3"] } +curvine-client-core = { workspace = true } +curvine-data-mover = { workspace = true, features = ["opendal-s3"] } curvine-web = { workspace = true } curvine-fault = { workspace = true } curvine-storage-local = { workspace = true } diff --git a/curvine-worker/src/lib.rs b/curvine-worker/src/lib.rs index ac4808e42..04e28b635 100644 --- a/curvine-worker/src/lib.rs +++ b/curvine-worker/src/lib.rs @@ -16,11 +16,11 @@ pub mod worker; pub use worker::*; pub mod common { - pub use curvine_master::UfsFactory; + pub use curvine_data_mover::common::UfsFactory; } pub mod master { - pub use curvine_master::master::RpcContext; + pub use curvine_data_mover::RpcContext; } pub mod transfer { diff --git a/curvine-worker/src/worker/block/block_actor.rs b/curvine-worker/src/worker/block/block_actor.rs index e986df359..a33ac855d 100644 --- a/curvine-worker/src/worker/block/block_actor.rs +++ b/curvine-worker/src/worker/block/block_actor.rs @@ -13,7 +13,7 @@ // limitations under the License. use crate::worker::block::{BlockStore, HeartbeatTask, MasterClient}; -use curvine_client::file::FsContext; +use curvine_client_core::file::FsContext; use curvine_common::conf::ClusterConf; use curvine_common::executor::ScheduledExecutor; use curvine_common::state::{BlockReportInfo, HeartbeatStatus, WorkerAddress}; diff --git a/curvine-worker/src/worker/block/master_client.rs b/curvine-worker/src/worker/block/master_client.rs index cd0ddde1a..32a2f08e8 100644 --- a/curvine-worker/src/worker/block/master_client.rs +++ b/curvine-worker/src/worker/block/master_client.rs @@ -13,7 +13,7 @@ // limitations under the License. use crate::worker::block::{BlockMeta, BlockState}; -use curvine_client::file::{FsClient, FsContext}; +use curvine_client_core::file::{FsClient, FsContext}; use curvine_common::fs::RpcCode; use curvine_common::proto::*; use curvine_common::state::{ diff --git a/curvine-worker/src/worker/replication/worker_replication_manager.rs b/curvine-worker/src/worker/replication/worker_replication_manager.rs index c80f607e6..1a92990a8 100644 --- a/curvine-worker/src/worker/replication/worker_replication_manager.rs +++ b/curvine-worker/src/worker/replication/worker_replication_manager.rs @@ -14,8 +14,8 @@ use crate::worker::block::{BlockState, BlockStore, MasterClient}; use crate::worker::replication::replication_job::ReplicationJob; -use curvine_client::block::BlockWriterRemote; -use curvine_client::file::FsContext; +use curvine_client_core::block::BlockWriterRemote; +use curvine_client_core::file::FsContext; use curvine_common::conf::ClusterConf; use curvine_common::fs::RpcCode; use curvine_common::proto::{ReportBlockReplicationRequest, ReportBlockReplicationResponse}; diff --git a/curvine-worker/src/worker/task/mod.rs b/curvine-worker/src/worker/task/mod.rs index 0b0d023a2..1c8e21ed9 100644 --- a/curvine-worker/src/worker/task/mod.rs +++ b/curvine-worker/src/worker/task/mod.rs @@ -12,17 +12,4 @@ // See the License for the specific language governing permissions and // limitations under the License. -//! Worker side data loading module -//! Responsible for loading data from external storage systems (such as S3) to local storage - -mod load_task_runner; -pub use self::load_task_runner::LoadTaskRunner; - -mod task_manager; -pub use self::task_manager::TaskManager; - -mod task_store; -pub use self::task_store::TaskStore; - -mod task_context; -pub use self::task_context::TaskContext; +pub use curvine_data_mover::worker::task::*; From 3736e28ea6f35fa71943db9c8e3869424e798692 Mon Sep 17 00:00:00 2001 From: barry Date: Sat, 1 Aug 2026 16:18:40 +0800 Subject: [PATCH 2/3] refactor(server): rename data transfer crate --- Cargo.lock | 8 ++++---- Cargo.toml | 6 +++--- .../Cargo.toml | 2 +- .../src/common/mod.rs | 0 .../src/common/ufs_client.rs | 0 .../src/common/ufs_factory.rs | 0 .../src/common/ufs_manager.rs | 0 .../src/job_worker_client.rs | 0 .../src/lib.rs | 0 .../src/rpc_context.rs | 0 .../src/transfer/backend.rs | 0 .../src/transfer/cluster_cache.rs | 0 .../src/transfer/handler.rs | 0 .../src/transfer/job_snapshot.rs | 0 .../src/transfer/memory_store.rs | 0 .../src/transfer/metrics.rs | 0 .../src/transfer/mod.rs | 0 .../src/transfer/mysql_store.rs | 0 .../src/transfer/planner.rs | 0 .../src/transfer/router_handler.rs | 0 .../src/transfer/scheduler.rs | 0 .../src/transfer/service.rs | 0 .../src/transfer/sqlite_store.rs | 0 .../src/transfer/store.rs | 0 .../src/transfer/tests/planner_test.rs | 0 .../src/transfer/transfer_server.rs | 0 .../src/worker/mod.rs | 0 .../src/worker/task/load_task_runner.rs | 0 .../src/worker/task/mod.rs | 0 .../src/worker/task/task_context.rs | 0 .../src/worker/task/task_manager.rs | 0 .../src/worker/task/task_store.rs | 0 curvine-master/Cargo.toml | 2 +- curvine-master/src/common/mod.rs | 2 +- curvine-master/src/master/job/mod.rs | 2 +- curvine-master/src/master/mod.rs | 2 +- curvine-server/Cargo.toml | 4 ++-- curvine-server/src/common/mod.rs | 2 +- curvine-server/src/transfer/mod.rs | 2 +- curvine-worker/Cargo.toml | 2 +- curvine-worker/src/lib.rs | 4 ++-- curvine-worker/src/worker/task/mod.rs | 2 +- 42 files changed, 20 insertions(+), 20 deletions(-) rename crates/server/{curvine-data-mover => curvine-data-transfer}/Cargo.toml (97%) rename crates/server/{curvine-data-mover => curvine-data-transfer}/src/common/mod.rs (100%) rename crates/server/{curvine-data-mover => curvine-data-transfer}/src/common/ufs_client.rs (100%) rename crates/server/{curvine-data-mover => curvine-data-transfer}/src/common/ufs_factory.rs (100%) rename crates/server/{curvine-data-mover => curvine-data-transfer}/src/common/ufs_manager.rs (100%) rename crates/server/{curvine-data-mover => curvine-data-transfer}/src/job_worker_client.rs (100%) rename crates/server/{curvine-data-mover => curvine-data-transfer}/src/lib.rs (100%) rename crates/server/{curvine-data-mover => curvine-data-transfer}/src/rpc_context.rs (100%) rename crates/server/{curvine-data-mover => curvine-data-transfer}/src/transfer/backend.rs (100%) rename crates/server/{curvine-data-mover => curvine-data-transfer}/src/transfer/cluster_cache.rs (100%) rename crates/server/{curvine-data-mover => curvine-data-transfer}/src/transfer/handler.rs (100%) rename crates/server/{curvine-data-mover => curvine-data-transfer}/src/transfer/job_snapshot.rs (100%) rename crates/server/{curvine-data-mover => curvine-data-transfer}/src/transfer/memory_store.rs (100%) rename crates/server/{curvine-data-mover => curvine-data-transfer}/src/transfer/metrics.rs (100%) rename crates/server/{curvine-data-mover => curvine-data-transfer}/src/transfer/mod.rs (100%) rename crates/server/{curvine-data-mover => curvine-data-transfer}/src/transfer/mysql_store.rs (100%) rename crates/server/{curvine-data-mover => curvine-data-transfer}/src/transfer/planner.rs (100%) rename crates/server/{curvine-data-mover => curvine-data-transfer}/src/transfer/router_handler.rs (100%) rename crates/server/{curvine-data-mover => curvine-data-transfer}/src/transfer/scheduler.rs (100%) rename crates/server/{curvine-data-mover => curvine-data-transfer}/src/transfer/service.rs (100%) rename crates/server/{curvine-data-mover => curvine-data-transfer}/src/transfer/sqlite_store.rs (100%) rename crates/server/{curvine-data-mover => curvine-data-transfer}/src/transfer/store.rs (100%) rename crates/server/{curvine-data-mover => curvine-data-transfer}/src/transfer/tests/planner_test.rs (100%) rename crates/server/{curvine-data-mover => curvine-data-transfer}/src/transfer/transfer_server.rs (100%) rename crates/server/{curvine-data-mover => curvine-data-transfer}/src/worker/mod.rs (100%) rename crates/server/{curvine-data-mover => curvine-data-transfer}/src/worker/task/load_task_runner.rs (100%) rename crates/server/{curvine-data-mover => curvine-data-transfer}/src/worker/task/mod.rs (100%) rename crates/server/{curvine-data-mover => curvine-data-transfer}/src/worker/task/task_context.rs (100%) rename crates/server/{curvine-data-mover => curvine-data-transfer}/src/worker/task/task_manager.rs (100%) rename crates/server/{curvine-data-mover => curvine-data-transfer}/src/worker/task/task_store.rs (100%) diff --git a/Cargo.lock b/Cargo.lock index 1407efbd4..d31620154 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -2169,7 +2169,7 @@ dependencies = [ ] [[package]] -name = "curvine-data-mover" +name = "curvine-data-transfer" version = "0.2.0" dependencies = [ "axum", @@ -2390,7 +2390,7 @@ dependencies = [ "axum", "bytes", "curvine-common", - "curvine-data-mover", + "curvine-data-transfer", "curvine-error", "curvine-fault", "curvine-rocksdb", @@ -2508,7 +2508,7 @@ dependencies = [ "crossbeam", "curvine-client-core", "curvine-common", - "curvine-data-mover", + "curvine-data-transfer", "curvine-error", "curvine-fault", "curvine-job-client", @@ -2717,7 +2717,7 @@ dependencies = [ "bytes", "curvine-client-core", "curvine-common", - "curvine-data-mover", + "curvine-data-transfer", "curvine-fault", "curvine-storage-local", "curvine-storage-spdk", diff --git a/Cargo.toml b/Cargo.toml index 4b0d612c2..66b404e82 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -32,7 +32,7 @@ members = [ "crates/adapters/curvine-ufs-opendal", "crates/adapters/curvine-ufs-oss-hdfs", "crates/adapters/curvine-hdfs-jni", - "crates/server/curvine-data-mover", + "crates/server/curvine-data-transfer", "curvine-master", "curvine-worker", "curvine-common", @@ -75,7 +75,7 @@ default-members = [ "crates/infra/curvine-io", "crates/metadata/curvine-raft", "crates/adapters/curvine-storage-local", - "crates/server/curvine-data-mover", + "crates/server/curvine-data-transfer", "curvine-master", "curvine-worker", "curvine-common", @@ -114,7 +114,7 @@ curvine-storage-spdk = { path = "crates/adapters/curvine-storage-spdk" } curvine-ufs-opendal = { path = "crates/adapters/curvine-ufs-opendal" } curvine-ufs-oss-hdfs = { path = "crates/adapters/curvine-ufs-oss-hdfs" } curvine-hdfs-jni = { path = "crates/adapters/curvine-hdfs-jni" } -curvine-data-mover = { path = "crates/server/curvine-data-mover" } +curvine-data-transfer = { path = "crates/server/curvine-data-transfer" } curvine-fault = { path = "crates/infra/curvine-fault" } curvine-error = { path = "crates/common/curvine-error" } curvine-proto = { path = "crates/common/curvine-proto" } diff --git a/crates/server/curvine-data-mover/Cargo.toml b/crates/server/curvine-data-transfer/Cargo.toml similarity index 97% rename from crates/server/curvine-data-mover/Cargo.toml rename to crates/server/curvine-data-transfer/Cargo.toml index b9ac214b0..f9d21531b 100644 --- a/crates/server/curvine-data-mover/Cargo.toml +++ b/crates/server/curvine-data-transfer/Cargo.toml @@ -1,5 +1,5 @@ [package] -name = "curvine-data-mover" +name = "curvine-data-transfer" version.workspace = true edition.workspace = true license.workspace = true diff --git a/crates/server/curvine-data-mover/src/common/mod.rs b/crates/server/curvine-data-transfer/src/common/mod.rs similarity index 100% rename from crates/server/curvine-data-mover/src/common/mod.rs rename to crates/server/curvine-data-transfer/src/common/mod.rs diff --git a/crates/server/curvine-data-mover/src/common/ufs_client.rs b/crates/server/curvine-data-transfer/src/common/ufs_client.rs similarity index 100% rename from crates/server/curvine-data-mover/src/common/ufs_client.rs rename to crates/server/curvine-data-transfer/src/common/ufs_client.rs diff --git a/crates/server/curvine-data-mover/src/common/ufs_factory.rs b/crates/server/curvine-data-transfer/src/common/ufs_factory.rs similarity index 100% rename from crates/server/curvine-data-mover/src/common/ufs_factory.rs rename to crates/server/curvine-data-transfer/src/common/ufs_factory.rs diff --git a/crates/server/curvine-data-mover/src/common/ufs_manager.rs b/crates/server/curvine-data-transfer/src/common/ufs_manager.rs similarity index 100% rename from crates/server/curvine-data-mover/src/common/ufs_manager.rs rename to crates/server/curvine-data-transfer/src/common/ufs_manager.rs diff --git a/crates/server/curvine-data-mover/src/job_worker_client.rs b/crates/server/curvine-data-transfer/src/job_worker_client.rs similarity index 100% rename from crates/server/curvine-data-mover/src/job_worker_client.rs rename to crates/server/curvine-data-transfer/src/job_worker_client.rs diff --git a/crates/server/curvine-data-mover/src/lib.rs b/crates/server/curvine-data-transfer/src/lib.rs similarity index 100% rename from crates/server/curvine-data-mover/src/lib.rs rename to crates/server/curvine-data-transfer/src/lib.rs diff --git a/crates/server/curvine-data-mover/src/rpc_context.rs b/crates/server/curvine-data-transfer/src/rpc_context.rs similarity index 100% rename from crates/server/curvine-data-mover/src/rpc_context.rs rename to crates/server/curvine-data-transfer/src/rpc_context.rs diff --git a/crates/server/curvine-data-mover/src/transfer/backend.rs b/crates/server/curvine-data-transfer/src/transfer/backend.rs similarity index 100% rename from crates/server/curvine-data-mover/src/transfer/backend.rs rename to crates/server/curvine-data-transfer/src/transfer/backend.rs diff --git a/crates/server/curvine-data-mover/src/transfer/cluster_cache.rs b/crates/server/curvine-data-transfer/src/transfer/cluster_cache.rs similarity index 100% rename from crates/server/curvine-data-mover/src/transfer/cluster_cache.rs rename to crates/server/curvine-data-transfer/src/transfer/cluster_cache.rs diff --git a/crates/server/curvine-data-mover/src/transfer/handler.rs b/crates/server/curvine-data-transfer/src/transfer/handler.rs similarity index 100% rename from crates/server/curvine-data-mover/src/transfer/handler.rs rename to crates/server/curvine-data-transfer/src/transfer/handler.rs diff --git a/crates/server/curvine-data-mover/src/transfer/job_snapshot.rs b/crates/server/curvine-data-transfer/src/transfer/job_snapshot.rs similarity index 100% rename from crates/server/curvine-data-mover/src/transfer/job_snapshot.rs rename to crates/server/curvine-data-transfer/src/transfer/job_snapshot.rs diff --git a/crates/server/curvine-data-mover/src/transfer/memory_store.rs b/crates/server/curvine-data-transfer/src/transfer/memory_store.rs similarity index 100% rename from crates/server/curvine-data-mover/src/transfer/memory_store.rs rename to crates/server/curvine-data-transfer/src/transfer/memory_store.rs diff --git a/crates/server/curvine-data-mover/src/transfer/metrics.rs b/crates/server/curvine-data-transfer/src/transfer/metrics.rs similarity index 100% rename from crates/server/curvine-data-mover/src/transfer/metrics.rs rename to crates/server/curvine-data-transfer/src/transfer/metrics.rs diff --git a/crates/server/curvine-data-mover/src/transfer/mod.rs b/crates/server/curvine-data-transfer/src/transfer/mod.rs similarity index 100% rename from crates/server/curvine-data-mover/src/transfer/mod.rs rename to crates/server/curvine-data-transfer/src/transfer/mod.rs diff --git a/crates/server/curvine-data-mover/src/transfer/mysql_store.rs b/crates/server/curvine-data-transfer/src/transfer/mysql_store.rs similarity index 100% rename from crates/server/curvine-data-mover/src/transfer/mysql_store.rs rename to crates/server/curvine-data-transfer/src/transfer/mysql_store.rs diff --git a/crates/server/curvine-data-mover/src/transfer/planner.rs b/crates/server/curvine-data-transfer/src/transfer/planner.rs similarity index 100% rename from crates/server/curvine-data-mover/src/transfer/planner.rs rename to crates/server/curvine-data-transfer/src/transfer/planner.rs diff --git a/crates/server/curvine-data-mover/src/transfer/router_handler.rs b/crates/server/curvine-data-transfer/src/transfer/router_handler.rs similarity index 100% rename from crates/server/curvine-data-mover/src/transfer/router_handler.rs rename to crates/server/curvine-data-transfer/src/transfer/router_handler.rs diff --git a/crates/server/curvine-data-mover/src/transfer/scheduler.rs b/crates/server/curvine-data-transfer/src/transfer/scheduler.rs similarity index 100% rename from crates/server/curvine-data-mover/src/transfer/scheduler.rs rename to crates/server/curvine-data-transfer/src/transfer/scheduler.rs diff --git a/crates/server/curvine-data-mover/src/transfer/service.rs b/crates/server/curvine-data-transfer/src/transfer/service.rs similarity index 100% rename from crates/server/curvine-data-mover/src/transfer/service.rs rename to crates/server/curvine-data-transfer/src/transfer/service.rs diff --git a/crates/server/curvine-data-mover/src/transfer/sqlite_store.rs b/crates/server/curvine-data-transfer/src/transfer/sqlite_store.rs similarity index 100% rename from crates/server/curvine-data-mover/src/transfer/sqlite_store.rs rename to crates/server/curvine-data-transfer/src/transfer/sqlite_store.rs diff --git a/crates/server/curvine-data-mover/src/transfer/store.rs b/crates/server/curvine-data-transfer/src/transfer/store.rs similarity index 100% rename from crates/server/curvine-data-mover/src/transfer/store.rs rename to crates/server/curvine-data-transfer/src/transfer/store.rs diff --git a/crates/server/curvine-data-mover/src/transfer/tests/planner_test.rs b/crates/server/curvine-data-transfer/src/transfer/tests/planner_test.rs similarity index 100% rename from crates/server/curvine-data-mover/src/transfer/tests/planner_test.rs rename to crates/server/curvine-data-transfer/src/transfer/tests/planner_test.rs diff --git a/crates/server/curvine-data-mover/src/transfer/transfer_server.rs b/crates/server/curvine-data-transfer/src/transfer/transfer_server.rs similarity index 100% rename from crates/server/curvine-data-mover/src/transfer/transfer_server.rs rename to crates/server/curvine-data-transfer/src/transfer/transfer_server.rs diff --git a/crates/server/curvine-data-mover/src/worker/mod.rs b/crates/server/curvine-data-transfer/src/worker/mod.rs similarity index 100% rename from crates/server/curvine-data-mover/src/worker/mod.rs rename to crates/server/curvine-data-transfer/src/worker/mod.rs diff --git a/crates/server/curvine-data-mover/src/worker/task/load_task_runner.rs b/crates/server/curvine-data-transfer/src/worker/task/load_task_runner.rs similarity index 100% rename from crates/server/curvine-data-mover/src/worker/task/load_task_runner.rs rename to crates/server/curvine-data-transfer/src/worker/task/load_task_runner.rs diff --git a/crates/server/curvine-data-mover/src/worker/task/mod.rs b/crates/server/curvine-data-transfer/src/worker/task/mod.rs similarity index 100% rename from crates/server/curvine-data-mover/src/worker/task/mod.rs rename to crates/server/curvine-data-transfer/src/worker/task/mod.rs diff --git a/crates/server/curvine-data-mover/src/worker/task/task_context.rs b/crates/server/curvine-data-transfer/src/worker/task/task_context.rs similarity index 100% rename from crates/server/curvine-data-mover/src/worker/task/task_context.rs rename to crates/server/curvine-data-transfer/src/worker/task/task_context.rs diff --git a/crates/server/curvine-data-mover/src/worker/task/task_manager.rs b/crates/server/curvine-data-transfer/src/worker/task/task_manager.rs similarity index 100% rename from crates/server/curvine-data-mover/src/worker/task/task_manager.rs rename to crates/server/curvine-data-transfer/src/worker/task/task_manager.rs diff --git a/crates/server/curvine-data-mover/src/worker/task/task_store.rs b/crates/server/curvine-data-transfer/src/worker/task/task_store.rs similarity index 100% rename from crates/server/curvine-data-mover/src/worker/task/task_store.rs rename to crates/server/curvine-data-transfer/src/worker/task/task_store.rs diff --git a/curvine-master/Cargo.toml b/curvine-master/Cargo.toml index a91d894d1..c2c98dc46 100644 --- a/curvine-master/Cargo.toml +++ b/curvine-master/Cargo.toml @@ -13,7 +13,7 @@ orpc = { workspace = true } curvine-common = { workspace = true } curvine-rocksdb = { workspace = true } curvine-error = { workspace = true, features = ["axum-response"] } -curvine-data-mover = { workspace = true, features = ["opendal-s3"] } +curvine-data-transfer = { workspace = true, features = ["opendal-s3"] } curvine-unified-fs = { workspace = true, features = ["opendal-s3"] } curvine-web = { workspace = true } curvine-fault = { workspace = true } diff --git a/curvine-master/src/common/mod.rs b/curvine-master/src/common/mod.rs index 02013813b..43f09cf4d 100644 --- a/curvine-master/src/common/mod.rs +++ b/curvine-master/src/common/mod.rs @@ -12,4 +12,4 @@ // See the License for the specific language governing permissions and // limitations under the License. -pub use curvine_data_mover::common::UfsFactory; +pub use curvine_data_transfer::common::UfsFactory; diff --git a/curvine-master/src/master/job/mod.rs b/curvine-master/src/master/job/mod.rs index eee89c62d..c1b22f0b3 100644 --- a/curvine-master/src/master/job/mod.rs +++ b/curvine-master/src/master/job/mod.rs @@ -18,7 +18,7 @@ pub use crate::master::job::job_manager::JobManager; mod job_handler; pub use job_handler::JobHandler; -pub use curvine_data_mover::JobWorkerClient; +pub use curvine_data_transfer::JobWorkerClient; mod job_store; pub use job_store::JobStore; diff --git a/curvine-master/src/master/mod.rs b/curvine-master/src/master/mod.rs index a0f32d0cf..84e21d970 100644 --- a/curvine-master/src/master/mod.rs +++ b/curvine-master/src/master/mod.rs @@ -45,7 +45,7 @@ pub use self::master_metrics::*; mod router_handler; pub use self::router_handler::*; -pub use curvine_data_mover::RpcContext; +pub use curvine_data_transfer::RpcContext; pub mod mount; diff --git a/curvine-server/Cargo.toml b/curvine-server/Cargo.toml index 67fb6ed6b..1397215ec 100644 --- a/curvine-server/Cargo.toml +++ b/curvine-server/Cargo.toml @@ -8,7 +8,7 @@ license.workspace = true [features] default = [] -jni = ["curvine-data-mover/opendal-hdfs"] +jni = ["curvine-data-transfer/opendal-hdfs"] spdk = [ "curvine-worker/spdk", ] @@ -30,7 +30,7 @@ curvine-worker = { workspace = true } curvine-error = { workspace = true, features = ["axum-response"] } curvine-fault = { workspace = true } curvine-client-core = { workspace = true } -curvine-data-mover = { workspace = true, features = ["opendal-s3"] } +curvine-data-transfer = { workspace = true, features = ["opendal-s3"] } curvine-job-client = { workspace = true } curvine-ufs-api = { workspace = true } curvine-web = { workspace = true } diff --git a/curvine-server/src/common/mod.rs b/curvine-server/src/common/mod.rs index 3bfeb703d..20b9c5790 100644 --- a/curvine-server/src/common/mod.rs +++ b/curvine-server/src/common/mod.rs @@ -12,4 +12,4 @@ // See the License for the specific language governing permissions and // limitations under the License. -pub use curvine_data_mover::common::*; +pub use curvine_data_transfer::common::*; diff --git a/curvine-server/src/transfer/mod.rs b/curvine-server/src/transfer/mod.rs index f71489f5a..fe6c000c8 100644 --- a/curvine-server/src/transfer/mod.rs +++ b/curvine-server/src/transfer/mod.rs @@ -1 +1 @@ -pub use curvine_data_mover::transfer::*; +pub use curvine_data_transfer::transfer::*; diff --git a/curvine-worker/Cargo.toml b/curvine-worker/Cargo.toml index 18607830f..e7e7b19c6 100644 --- a/curvine-worker/Cargo.toml +++ b/curvine-worker/Cargo.toml @@ -24,7 +24,7 @@ fault-injection = [ orpc = { workspace = true } curvine-common = { workspace = true } curvine-client-core = { workspace = true } -curvine-data-mover = { workspace = true, features = ["opendal-s3"] } +curvine-data-transfer = { workspace = true, features = ["opendal-s3"] } curvine-web = { workspace = true } curvine-fault = { workspace = true } curvine-storage-local = { workspace = true } diff --git a/curvine-worker/src/lib.rs b/curvine-worker/src/lib.rs index 04e28b635..4a53adcdf 100644 --- a/curvine-worker/src/lib.rs +++ b/curvine-worker/src/lib.rs @@ -16,11 +16,11 @@ pub mod worker; pub use worker::*; pub mod common { - pub use curvine_data_mover::common::UfsFactory; + pub use curvine_data_transfer::common::UfsFactory; } pub mod master { - pub use curvine_data_mover::RpcContext; + pub use curvine_data_transfer::RpcContext; } pub mod transfer { diff --git a/curvine-worker/src/worker/task/mod.rs b/curvine-worker/src/worker/task/mod.rs index 1c8e21ed9..9e4ed1e56 100644 --- a/curvine-worker/src/worker/task/mod.rs +++ b/curvine-worker/src/worker/task/mod.rs @@ -12,4 +12,4 @@ // See the License for the specific language governing permissions and // limitations under the License. -pub use curvine_data_mover::worker::task::*; +pub use curvine_data_transfer::worker::task::*; From eca53b2f2196c57b65ec69951879fece74c4faf9 Mon Sep 17 00:00:00 2001 From: barry Date: Sat, 1 Aug 2026 16:49:29 +0800 Subject: [PATCH 3/3] refactor(server): gate data transfer server deps --- .../server/curvine-data-transfer/Cargo.toml | 19 +++-- .../server/curvine-data-transfer/src/lib.rs | 1 + .../src/transfer/backend.rs | 70 ++++++++++++++++++- .../curvine-data-transfer/src/transfer/mod.rs | 8 +++ .../src/worker/task/load_task_runner.rs | 2 +- curvine-server/Cargo.toml | 2 +- 6 files changed, 93 insertions(+), 9 deletions(-) diff --git a/crates/server/curvine-data-transfer/Cargo.toml b/crates/server/curvine-data-transfer/Cargo.toml index f9d21531b..231f79d70 100644 --- a/crates/server/curvine-data-transfer/Cargo.toml +++ b/crates/server/curvine-data-transfer/Cargo.toml @@ -6,6 +6,15 @@ license.workspace = true [features] default = [] +transfer = ["dep:uuid"] +transfer-web = ["transfer", "dep:axum", "dep:curvine-web"] +transfer-store-sqlite = ["transfer", "dep:rusqlite"] +transfer-store-mysql = ["transfer", "dep:mysql"] +transfer-server = [ + "transfer-web", + "transfer-store-sqlite", + "transfer-store-mysql", +] opendal = ["curvine-unified-fs/opendal"] opendal-s3 = ["opendal", "curvine-unified-fs/opendal-s3"] opendal-oss = ["opendal", "curvine-unified-fs/opendal-oss"] @@ -24,7 +33,7 @@ curvine-client-core = { workspace = true } curvine-job-client = { workspace = true } curvine-unified-fs = { workspace = true } curvine-ufs-api = { workspace = true } -curvine-web = { workspace = true } +curvine-web = { workspace = true, optional = true } prost = { workspace = true } bytes = { workspace = true } serde = { workspace = true } @@ -36,7 +45,7 @@ dashmap = { workspace = true } crossbeam = { workspace = true } parking_lot = { workspace = true } once_cell = { workspace = true } -axum = { workspace = true } -rusqlite = { workspace = true } -mysql = { workspace = true } -uuid = { workspace = true } +axum = { workspace = true, optional = true } +rusqlite = { workspace = true, optional = true } +mysql = { workspace = true, optional = true } +uuid = { workspace = true, optional = true } diff --git a/crates/server/curvine-data-transfer/src/lib.rs b/crates/server/curvine-data-transfer/src/lib.rs index 5ae768949..9d3bb17db 100644 --- a/crates/server/curvine-data-transfer/src/lib.rs +++ b/crates/server/curvine-data-transfer/src/lib.rs @@ -13,6 +13,7 @@ // limitations under the License. pub mod common; +#[cfg(feature = "transfer")] pub mod transfer; pub mod worker; diff --git a/crates/server/curvine-data-transfer/src/transfer/backend.rs b/crates/server/curvine-data-transfer/src/transfer/backend.rs index 24c665fda..91340ad3a 100644 --- a/crates/server/curvine-data-transfer/src/transfer/backend.rs +++ b/crates/server/curvine-data-transfer/src/transfer/backend.rs @@ -20,14 +20,20 @@ use curvine_common::state::{ use curvine_common::FsResult; use std::time::Instant; +#[cfg(feature = "transfer-store-mysql")] +use crate::transfer::MysqlTransferStore; +#[cfg(feature = "transfer-store-sqlite")] +use crate::transfer::SqliteTransferStore; use crate::transfer::{ - MemoryTransferStore, MysqlTransferStore, SqliteTransferStore, TransferMetrics, - TransferPlannedTasks, TransferRequeueUpdate, TransferStore, TransferTaskStateUpdate, + MemoryTransferStore, TransferMetrics, TransferPlannedTasks, TransferRequeueUpdate, + TransferStore, TransferTaskStateUpdate, }; pub enum TransferStoreBackend { Memory(MemoryTransferStore), + #[cfg(feature = "transfer-store-sqlite")] Sqlite(SqliteTransferStore), + #[cfg(feature = "transfer-store-mysql")] Mysql(MysqlTransferStore), } @@ -35,7 +41,9 @@ impl TransferStore for TransferStoreBackend { fn check_available(&self) -> FsResult<()> { self.record_store_operation("check_available", || match self { Self::Memory(store) => store.check_available(), + #[cfg(feature = "transfer-store-sqlite")] Self::Sqlite(store) => store.check_available(), + #[cfg(feature = "transfer-store-mysql")] Self::Mysql(store) => store.check_available(), }) } @@ -43,7 +51,9 @@ impl TransferStore for TransferStoreBackend { fn create_or_get_by_request_id(&self, job: TransferJobRecord) -> FsResult { self.record_store_operation("create_or_get_by_request_id", || match self { Self::Memory(store) => store.create_or_get_by_request_id(job), + #[cfg(feature = "transfer-store-sqlite")] Self::Sqlite(store) => store.create_or_get_by_request_id(job), + #[cfg(feature = "transfer-store-mysql")] Self::Mysql(store) => store.create_or_get_by_request_id(job), }) } @@ -54,7 +64,9 @@ impl TransferStore for TransferStoreBackend { ) -> FsResult { self.record_store_operation("create_or_get_by_request_id_checked", || match self { Self::Memory(store) => store.create_or_get_by_request_id_checked(job), + #[cfg(feature = "transfer-store-sqlite")] Self::Sqlite(store) => store.create_or_get_by_request_id_checked(job), + #[cfg(feature = "transfer-store-mysql")] Self::Mysql(store) => store.create_or_get_by_request_id_checked(job), }) } @@ -62,7 +74,9 @@ impl TransferStore for TransferStoreBackend { fn get_transfer(&self, job_id: &str) -> FsResult> { self.record_store_operation("get_transfer", || match self { Self::Memory(store) => store.get_transfer(job_id), + #[cfg(feature = "transfer-store-sqlite")] Self::Sqlite(store) => store.get_transfer(job_id), + #[cfg(feature = "transfer-store-mysql")] Self::Mysql(store) => store.get_transfer(job_id), }) } @@ -74,7 +88,9 @@ impl TransferStore for TransferStoreBackend { ) -> FsResult> { self.record_store_operation("get_transfer_by_request", || match self { Self::Memory(store) => store.get_transfer_by_request(submitter, client_request_id), + #[cfg(feature = "transfer-store-sqlite")] Self::Sqlite(store) => store.get_transfer_by_request(submitter, client_request_id), + #[cfg(feature = "transfer-store-mysql")] Self::Mysql(store) => store.get_transfer_by_request(submitter, client_request_id), }) } @@ -82,7 +98,9 @@ impl TransferStore for TransferStoreBackend { fn list_active_transfers(&self) -> FsResult> { self.record_store_operation("list_active_transfers", || match self { Self::Memory(store) => store.list_active_transfers(), + #[cfg(feature = "transfer-store-sqlite")] Self::Sqlite(store) => store.list_active_transfers(), + #[cfg(feature = "transfer-store-mysql")] Self::Mysql(store) => store.list_active_transfers(), }) } @@ -97,9 +115,11 @@ impl TransferStore for TransferStoreBackend { Self::Memory(store) => { store.find_conflicting_active_transfer(target_path, submitter, client_request_id) } + #[cfg(feature = "transfer-store-sqlite")] Self::Sqlite(store) => { store.find_conflicting_active_transfer(target_path, submitter, client_request_id) } + #[cfg(feature = "transfer-store-mysql")] Self::Mysql(store) => { store.find_conflicting_active_transfer(target_path, submitter, client_request_id) } @@ -109,7 +129,9 @@ impl TransferStore for TransferStoreBackend { fn count_active_transfers(&self) -> FsResult { self.record_store_operation("count_active_transfers", || match self { Self::Memory(store) => store.count_active_transfers(), + #[cfg(feature = "transfer-store-sqlite")] Self::Sqlite(store) => store.count_active_transfers(), + #[cfg(feature = "transfer-store-mysql")] Self::Mysql(store) => store.count_active_transfers(), }) } @@ -117,7 +139,9 @@ impl TransferStore for TransferStoreBackend { fn count_executing_transfers(&self) -> FsResult { self.record_store_operation("count_executing_transfers", || match self { Self::Memory(store) => store.count_executing_transfers(), + #[cfg(feature = "transfer-store-sqlite")] Self::Sqlite(store) => store.count_executing_transfers(), + #[cfg(feature = "transfer-store-mysql")] Self::Mysql(store) => store.count_executing_transfers(), }) } @@ -125,7 +149,9 @@ impl TransferStore for TransferStoreBackend { fn list_transfers(&self, filter: TransferListFilter) -> FsResult> { self.record_store_operation("list_transfers", || match self { Self::Memory(store) => store.list_transfers(filter), + #[cfg(feature = "transfer-store-sqlite")] Self::Sqlite(store) => store.list_transfers(filter), + #[cfg(feature = "transfer-store-mysql")] Self::Mysql(store) => store.list_transfers(filter), }) } @@ -137,7 +163,9 @@ impl TransferStore for TransferStoreBackend { ) -> FsResult> { self.record_store_operation("list_tenant_summaries", || match self { Self::Memory(store) => store.list_tenant_summaries(limit, offset), + #[cfg(feature = "transfer-store-sqlite")] Self::Sqlite(store) => store.list_tenant_summaries(limit, offset), + #[cfg(feature = "transfer-store-mysql")] Self::Mysql(store) => store.list_tenant_summaries(limit, offset), }) } @@ -145,7 +173,9 @@ impl TransferStore for TransferStoreBackend { fn purge_terminal_transfers(&self, older_than_ms: i64, limit: usize) -> FsResult { self.record_store_operation("purge_terminal_transfers", || match self { Self::Memory(store) => store.purge_terminal_transfers(older_than_ms, limit), + #[cfg(feature = "transfer-store-sqlite")] Self::Sqlite(store) => store.purge_terminal_transfers(older_than_ms, limit), + #[cfg(feature = "transfer-store-mysql")] Self::Mysql(store) => store.purge_terminal_transfers(older_than_ms, limit), }) } @@ -153,7 +183,9 @@ impl TransferStore for TransferStoreBackend { fn list_transfer_tasks(&self, job_id: &str, run_id: u64) -> FsResult> { self.record_store_operation("list_transfer_tasks", || match self { Self::Memory(store) => store.list_transfer_tasks(job_id, run_id), + #[cfg(feature = "transfer-store-sqlite")] Self::Sqlite(store) => store.list_transfer_tasks(job_id, run_id), + #[cfg(feature = "transfer-store-mysql")] Self::Mysql(store) => store.list_transfer_tasks(job_id, run_id), }) } @@ -161,7 +193,9 @@ impl TransferStore for TransferStoreBackend { fn request_cancel(&self, job_id: &str, run_id: u64, now_ms: i64) -> FsResult { self.record_store_operation("request_cancel", || match self { Self::Memory(store) => store.request_cancel(job_id, run_id, now_ms), + #[cfg(feature = "transfer-store-sqlite")] Self::Sqlite(store) => store.request_cancel(job_id, run_id, now_ms), + #[cfg(feature = "transfer-store-mysql")] Self::Mysql(store) => store.request_cancel(job_id, run_id, now_ms), }) } @@ -177,9 +211,11 @@ impl TransferStore for TransferStoreBackend { Self::Memory(store) => { store.acquire_runnable_transfer(owner, lease_ms, now_ms, max_executing_transfers) } + #[cfg(feature = "transfer-store-sqlite")] Self::Sqlite(store) => { store.acquire_runnable_transfer(owner, lease_ms, now_ms, max_executing_transfers) } + #[cfg(feature = "transfer-store-mysql")] Self::Mysql(store) => { store.acquire_runnable_transfer(owner, lease_ms, now_ms, max_executing_transfers) } @@ -199,9 +235,11 @@ impl TransferStore for TransferStoreBackend { Self::Memory(store) => { store.renew_lease(job_id, run_id, owner, lease_epoch, lease_ms, now_ms) } + #[cfg(feature = "transfer-store-sqlite")] Self::Sqlite(store) => { store.renew_lease(job_id, run_id, owner, lease_epoch, lease_ms, now_ms) } + #[cfg(feature = "transfer-store-mysql")] Self::Mysql(store) => { store.renew_lease(job_id, run_id, owner, lease_epoch, lease_ms, now_ms) } @@ -211,7 +249,9 @@ impl TransferStore for TransferStoreBackend { fn update_transfer_state(&self, update: TransferStateUpdate) -> FsResult { self.record_store_operation("update_transfer_state", || match self { Self::Memory(store) => store.update_transfer_state(update), + #[cfg(feature = "transfer-store-sqlite")] Self::Sqlite(store) => store.update_transfer_state(update), + #[cfg(feature = "transfer-store-mysql")] Self::Mysql(store) => store.update_transfer_state(update), }) } @@ -219,7 +259,9 @@ impl TransferStore for TransferStoreBackend { fn requeue_transfer(&self, update: TransferRequeueUpdate) -> FsResult { self.record_store_operation("requeue_transfer", || match self { Self::Memory(store) => store.requeue_transfer(update), + #[cfg(feature = "transfer-store-sqlite")] Self::Sqlite(store) => store.requeue_transfer(update), + #[cfg(feature = "transfer-store-mysql")] Self::Mysql(store) => store.requeue_transfer(update), }) } @@ -242,6 +284,7 @@ impl TransferStore for TransferStoreBackend { cv_metadata_epoch, now_ms, ), + #[cfg(feature = "transfer-store-sqlite")] Self::Sqlite(store) => store.set_transfer_cv_metadata_epoch( job_id, run_id, @@ -250,6 +293,7 @@ impl TransferStore for TransferStoreBackend { cv_metadata_epoch, now_ms, ), + #[cfg(feature = "transfer-store-mysql")] Self::Mysql(store) => store.set_transfer_cv_metadata_epoch( job_id, run_id, @@ -264,7 +308,9 @@ impl TransferStore for TransferStoreBackend { fn insert_tasks(&self, tasks: Vec) -> FsResult<()> { self.record_store_operation("insert_tasks", || match self { Self::Memory(store) => store.insert_tasks(tasks), + #[cfg(feature = "transfer-store-sqlite")] Self::Sqlite(store) => store.insert_tasks(tasks), + #[cfg(feature = "transfer-store-mysql")] Self::Mysql(store) => store.insert_tasks(tasks), }) } @@ -272,7 +318,9 @@ impl TransferStore for TransferStoreBackend { fn persist_planned_tasks(&self, update: TransferPlannedTasks) -> FsResult { self.record_store_operation("persist_planned_tasks", || match self { Self::Memory(store) => store.persist_planned_tasks(update), + #[cfg(feature = "transfer-store-sqlite")] Self::Sqlite(store) => store.persist_planned_tasks(update), + #[cfg(feature = "transfer-store-mysql")] Self::Mysql(store) => store.persist_planned_tasks(update), }) } @@ -280,7 +328,9 @@ impl TransferStore for TransferStoreBackend { fn update_task_state(&self, update: TransferTaskStateUpdate) -> FsResult { self.record_store_operation("update_task_state", || match self { Self::Memory(store) => store.update_task_state(update), + #[cfg(feature = "transfer-store-sqlite")] Self::Sqlite(store) => store.update_task_state(update), + #[cfg(feature = "transfer-store-mysql")] Self::Mysql(store) => store.update_task_state(update), }) } @@ -293,7 +343,9 @@ impl TransferStore for TransferStoreBackend { ) -> FsResult> { self.record_store_operation("claim_pending_tasks", || match self { Self::Memory(store) => store.claim_pending_tasks(job_id, run_id, limit), + #[cfg(feature = "transfer-store-sqlite")] Self::Sqlite(store) => store.claim_pending_tasks(job_id, run_id, limit), + #[cfg(feature = "transfer-store-mysql")] Self::Mysql(store) => store.claim_pending_tasks(job_id, run_id, limit), }) } @@ -311,9 +363,11 @@ impl TransferStore for TransferStoreBackend { Self::Memory(store) => { store.mark_stale_attempts(job_id, run_id, owner, lease_epoch, now_ms, limit) } + #[cfg(feature = "transfer-store-sqlite")] Self::Sqlite(store) => { store.mark_stale_attempts(job_id, run_id, owner, lease_epoch, now_ms, limit) } + #[cfg(feature = "transfer-store-mysql")] Self::Mysql(store) => { store.mark_stale_attempts(job_id, run_id, owner, lease_epoch, now_ms, limit) } @@ -331,9 +385,11 @@ impl TransferStore for TransferStoreBackend { Self::Memory(store) => { store.list_stale_running_tasks(job_id, run_id, stale_before_ms, limit) } + #[cfg(feature = "transfer-store-sqlite")] Self::Sqlite(store) => { store.list_stale_running_tasks(job_id, run_id, stale_before_ms, limit) } + #[cfg(feature = "transfer-store-mysql")] Self::Mysql(store) => { store.list_stale_running_tasks(job_id, run_id, stale_before_ms, limit) } @@ -343,7 +399,9 @@ impl TransferStore for TransferStoreBackend { fn has_failed_tasks(&self, job_id: &str, run_id: u64) -> FsResult { self.record_store_operation("has_failed_tasks", || match self { Self::Memory(store) => store.has_failed_tasks(job_id, run_id), + #[cfg(feature = "transfer-store-sqlite")] Self::Sqlite(store) => store.has_failed_tasks(job_id, run_id), + #[cfg(feature = "transfer-store-mysql")] Self::Mysql(store) => store.has_failed_tasks(job_id, run_id), }) } @@ -351,7 +409,9 @@ impl TransferStore for TransferStoreBackend { fn has_recoverable_tasks(&self, job_id: &str, run_id: u64) -> FsResult { self.record_store_operation("has_recoverable_tasks", || match self { Self::Memory(store) => store.has_recoverable_tasks(job_id, run_id), + #[cfg(feature = "transfer-store-sqlite")] Self::Sqlite(store) => store.has_recoverable_tasks(job_id, run_id), + #[cfg(feature = "transfer-store-mysql")] Self::Mysql(store) => store.has_recoverable_tasks(job_id, run_id), }) } @@ -359,7 +419,9 @@ impl TransferStore for TransferStoreBackend { fn start_task_attempt(&self, start: TaskAttemptStart) -> FsResult { self.record_store_operation("start_task_attempt", || match self { Self::Memory(store) => store.start_task_attempt(start), + #[cfg(feature = "transfer-store-sqlite")] Self::Sqlite(store) => store.start_task_attempt(start), + #[cfg(feature = "transfer-store-mysql")] Self::Mysql(store) => store.start_task_attempt(start), }) } @@ -367,7 +429,9 @@ impl TransferStore for TransferStoreBackend { fn update_task_report(&self, report: TransferTaskReport) -> FsResult { self.record_store_operation("update_task_report", || match self { Self::Memory(store) => store.update_task_report(report), + #[cfg(feature = "transfer-store-sqlite")] Self::Sqlite(store) => store.update_task_report(report), + #[cfg(feature = "transfer-store-mysql")] Self::Mysql(store) => store.update_task_report(report), }) } @@ -377,7 +441,9 @@ impl TransferStoreBackend { pub fn backend_label(&self) -> &'static str { match self { Self::Memory(_) => "memory", + #[cfg(feature = "transfer-store-sqlite")] Self::Sqlite(_) => "sqlite", + #[cfg(feature = "transfer-store-mysql")] Self::Mysql(_) => "mysql", } } diff --git a/crates/server/curvine-data-transfer/src/transfer/mod.rs b/crates/server/curvine-data-transfer/src/transfer/mod.rs index ab08e8a7e..c0eeb0221 100644 --- a/crates/server/curvine-data-transfer/src/transfer/mod.rs +++ b/crates/server/curvine-data-transfer/src/transfer/mod.rs @@ -4,10 +4,14 @@ pub use self::store::*; mod memory_store; pub use self::memory_store::MemoryTransferStore; +#[cfg(feature = "transfer-store-sqlite")] mod sqlite_store; +#[cfg(feature = "transfer-store-sqlite")] pub use self::sqlite_store::SqliteTransferStore; +#[cfg(feature = "transfer-store-mysql")] mod mysql_store; +#[cfg(feature = "transfer-store-mysql")] pub use self::mysql_store::MysqlTransferStore; mod metrics; @@ -40,10 +44,14 @@ pub use self::scheduler::TransferScheduler; mod handler; pub use self::handler::TransferHandler; +#[cfg(feature = "transfer-web")] mod router_handler; +#[cfg(feature = "transfer-web")] pub use self::router_handler::TransferRouterHandler; +#[cfg(feature = "transfer-server")] mod transfer_server; +#[cfg(feature = "transfer-server")] pub use self::transfer_server::TransferServer; pub(crate) fn apply_task_report_progress( diff --git a/crates/server/curvine-data-transfer/src/worker/task/load_task_runner.rs b/crates/server/curvine-data-transfer/src/worker/task/load_task_runner.rs index d9f6e2cb6..8009fd2af 100644 --- a/crates/server/curvine-data-transfer/src/worker/task/load_task_runner.rs +++ b/crates/server/curvine-data-transfer/src/worker/task/load_task_runner.rs @@ -13,7 +13,6 @@ // limitations under the License. use crate::common::UfsFactory; -use crate::transfer::transfer_failure_message; use crate::worker::task::TaskContext; use curvine_client_core::file::{CurvineFileSystem, FsReader}; use curvine_common::error::FsError; @@ -22,6 +21,7 @@ use curvine_common::state::{ CreateFileOptsBuilder, FileBlocks, FileStatus, JobTaskProgress, JobTaskState, SetAttrOptsBuilder, TRANSFER_TEMP_PATH_MARKER, }; +use curvine_common::transfer::transfer_failure_message; use curvine_common::FsResult; use curvine_job_client::JobMasterClient; use curvine_job_client::TransferClient; diff --git a/curvine-server/Cargo.toml b/curvine-server/Cargo.toml index 1397215ec..9607c6004 100644 --- a/curvine-server/Cargo.toml +++ b/curvine-server/Cargo.toml @@ -30,7 +30,7 @@ curvine-worker = { workspace = true } curvine-error = { workspace = true, features = ["axum-response"] } curvine-fault = { workspace = true } curvine-client-core = { workspace = true } -curvine-data-transfer = { workspace = true, features = ["opendal-s3"] } +curvine-data-transfer = { workspace = true, features = ["opendal-s3", "transfer-server"] } curvine-job-client = { workspace = true } curvine-ufs-api = { workspace = true } curvine-web = { workspace = true }