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

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
1 change: 1 addition & 0 deletions .gitignore
Original file line number Diff line number Diff line change
Expand Up @@ -89,6 +89,7 @@ java/target
java/test-libs
java/*.log
java/include/org_rocksdb_*.h
juicefs-hadoop-1.3.1.jar

.idea/
*.iml
Expand Down
5 changes: 5 additions & 0 deletions build_tools/pack_server.sh
Original file line number Diff line number Diff line change
Expand Up @@ -168,6 +168,11 @@ if [ -n "$HADOOP_HOME" ]; then
fi
# Pack the jars.
mkdir -p ${pack}/hadoop

# output/hadoop/juicefs-hadoop-1.3.1.jar (match thirdparty)
juicefs_jar="${THIRDPARTY_ROOT}/output/hadoop/juicefs-hadoop-1.3.1.jar"
[ -f "$juicefs_jar" ] || { echo "ERROR: $juicefs_jar not found"; exit 1; }
copy_file "$juicefs_jar" ${pack}/hadoop
for f in ${HADOOP_HOME}/share/hadoop/common/lib/*.jar; do
copy_file $f ${pack}/hadoop
done
Expand Down
41 changes: 32 additions & 9 deletions src/block_service/block_service_manager.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -37,12 +37,19 @@ namespace dsn {
namespace dist {
namespace block_service {

const char *BLOCK_SERVICE_JUICEFS = "juicefs_service";

block_service_registry::block_service_registry()
{
CHECK(utils::factory_store<block_filesystem>::register_factory(
"hdfs_service", block_filesystem::create<hdfs_service>, PROVIDER_TYPE_MAIN),
"register hdfs_service failed");

// juice_service use hdfs_service as default provider

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

juice_service -> juicefs_service

CHECK(utils::factory_store<block_filesystem>::register_factory(
BLOCK_SERVICE_JUICEFS, block_filesystem::create<hdfs_service>, PROVIDER_TYPE_MAIN),
"register juice_service failed");

CHECK(utils::factory_store<block_filesystem>::register_factory(
"local_service", block_filesystem::create<local_service>, PROVIDER_TYPE_MAIN),
"register local_service failed");
Expand Down Expand Up @@ -70,27 +77,43 @@ block_filesystem *block_service_manager::get_or_create_block_filesystem(const st
return iter->second.get();
}

const char *provider_type = dsn_config_get_value_string(
(std::string("block_service.") + provider).c_str(), "type", "", "block service type");
block_filesystem *fs = nullptr;
const char *provider_type = nullptr;
bool isJuicefs = is_juicefs_provider(provider);

block_filesystem *fs =
utils::factory_store<block_filesystem>::create(provider_type, PROVIDER_TYPE_MAIN);
if (isJuicefs) {
provider_type = BLOCK_SERVICE_JUICEFS;
} else {
provider_type = dsn_config_get_value_string(
(std::string("block_service.") + provider).c_str(), "type", "", "block service type");
}
fs = utils::factory_store<block_filesystem>::create(provider_type, PROVIDER_TYPE_MAIN);
if (fs == nullptr) {
LOG_ERROR("acquire block filesystem failed, provider = {}, provider_type = {}",
provider,
provider_type);
return nullptr;
}

const char *arguments = dsn_config_get_value_string(
(std::string("block_service.") + provider).c_str(), "args", "", "args for block_service");

std::vector<std::string> args;
utils::split_args(arguments, args);
std::string args_for_log;
if (isJuicefs) {
// juicefs provider example: jfs://pegasus@ak-bigdata
args = {provider, "/"};
args_for_log = provider + " /";
} else {
const char *arguments =
dsn_config_get_value_string((std::string("block_service.") + provider).c_str(),
"args",
"",
"args for block_service");
utils::split_args(arguments, args);
args_for_log = arguments;
}
dsn::error_code err = fs->initialize(args);

const auto provider_desc = fmt::format(
"provider = {}, provider_type = {}, args = {}", provider, provider_type, arguments);
"provider = {}, provider_type = {}, args = {}", provider, provider_type, args_for_log);
if (dsn::ERR_OK == err) {
LOG_INFO("create block filesystem ok for {}", provider_desc);
_fs_map.emplace(provider, std::unique_ptr<block_filesystem>(fs));
Expand Down
32 changes: 24 additions & 8 deletions src/block_service/block_service_manager.h
Original file line number Diff line number Diff line change
Expand Up @@ -17,18 +17,22 @@

#pragma once

#include <stdint.h>
#include <absl/strings/match.h>
#include <cstddef>
#include <cstdint>
#include <map>
#include <memory>
#include <string>

#include <string_view>
#include "utils/error_code.h"
#include "utils/singleton.h"
#include "utils/zlocks.h"

namespace dsn {
namespace dist {
namespace block_service {
namespace dsn::dist::block_service {

// JuiceFS provider URL prefix (e.g., jfs://volume@cluster_name)
constexpr std::string_view JUICEFS_PROVIDER_PREFIX = "jfs://";

class block_filesystem;

// a singleton for rDSN service_engine to register all blocks, this should be called only once
Expand All @@ -51,6 +55,20 @@ class block_service_manager
~block_service_manager();
block_filesystem *get_or_create_block_filesystem(const std::string &provider);

static bool is_juicefs_provider(const std::string &provider)
{
if (!absl::StartsWith(provider, JUICEFS_PROVIDER_PREFIX)) {
return false;
}
std::string remaining = provider.substr(JUICEFS_PROVIDER_PREFIX.size());
size_t at_pos = remaining.find('@');
if (at_pos == std::string::npos || at_pos == 0) {
return false;
}
std::string host = remaining.substr(at_pos + 1);
return !host.empty();
}

// download files from remote file system
// \return ERR_FILE_OPERATION_FAILED: local file system error
// \return ERR_FS_INTERNAL: remote file system error
Expand Down Expand Up @@ -84,6 +102,4 @@ class block_service_manager
friend class block_service_manager_mock;
};

} // namespace block_service
} // namespace dist
} // namespace dsn
} // namespace dsn::dist::block_service
32 changes: 32 additions & 0 deletions src/block_service/test/block_service_manager_test.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -72,6 +72,38 @@ class block_service_manager_test : public pegasus::encrypt_data_test_base

INSTANTIATE_TEST_SUITE_P(, block_service_manager_test, ::testing::Values(false, true));

TEST(is_juicefs_provider_test, valid_provider)
{
EXPECT_TRUE(block_service_manager::is_juicefs_provider("jfs://volume@cluster_name"));
EXPECT_TRUE(block_service_manager::is_juicefs_provider("jfs://pegasus@ak-bigdata"));
EXPECT_TRUE(block_service_manager::is_juicefs_provider("jfs://admin@192.168.1.1"));
}

TEST(is_juicefs_provider_test, no_volume)
{
EXPECT_FALSE(block_service_manager::is_juicefs_provider("jfs://@cluster_name"));
EXPECT_FALSE(block_service_manager::is_juicefs_provider("jfs://@"));
}

TEST(is_juicefs_provider_test, no_cluster_name)
{
EXPECT_FALSE(block_service_manager::is_juicefs_provider("jfs://volume@"));
}

TEST(is_juicefs_provider_test, no_at_symbol)
{
EXPECT_FALSE(block_service_manager::is_juicefs_provider("jfs://volumecluster"));
EXPECT_FALSE(block_service_manager::is_juicefs_provider("jfs://"));
}

TEST(is_juicefs_provider_test, wrong_prefix)
{
EXPECT_FALSE(block_service_manager::is_juicefs_provider("dfs://volume@cluster_name"));
EXPECT_FALSE(block_service_manager::is_juicefs_provider("hdfs://volume@cluster_name"));
EXPECT_FALSE(block_service_manager::is_juicefs_provider(""));
EXPECT_FALSE(block_service_manager::is_juicefs_provider("jfs:/volume@cluster_name"));
}

TEST_P(block_service_manager_test, remote_file_not_exist)
{
utils::filesystem::remove_path(LOCAL_DIR);
Expand Down
25 changes: 25 additions & 0 deletions src/block_service/test/run.sh
Original file line number Diff line number Diff line change
Expand Up @@ -23,6 +23,31 @@ if [ -z "${REPORT_DIR}" ]; then
REPORT_DIR="."
fi

# By default, unit tests use local_service.
# To connect to HDFS/JuiceFS or other storage systems in unit tests, set PACKAGE_DIR:
# export PACKAGE_DIR=/path/to/pegasus-server-x.x.x-glibc2.17-release
if [ -n "${PACKAGE_DIR}" ]; then
package_dir="${PACKAGE_DIR}"
echo "Using package_dir: $package_dir"

# Set the ld library path
ld_library_path=$package_dir/DSN_ROOT/lib:$package_dir/bin:$LD_LIBRARY_PATH
export LD_LIBRARY_PATH=$ld_library_path

export CLASSPATH=$package_dir/hadoop/
for f in $package_dir/hadoop/*.jar; do
export CLASSPATH=$CLASSPATH:$f
done
JAVA_JVM_LIBRARY_DIR=$(dirname $(find "${JAVA_HOME}/" -name libjvm.so | head -1))
export LD_LIBRARY_PATH=${JAVA_JVM_LIBRARY_DIR}:$LD_LIBRARY_PATH

echo CLASSPATH=$CLASSPATH
echo LD_LIBRARY_PATH=$LD_LIBRARY_PATH
echo JAVA_JVM_LIBRARY_DIR=$JAVA_JVM_LIBRARY_DIR
else
echo "PACKAGE_DIR is not set, running tests with local_service only."
fi

./clear.sh
output_xml="${REPORT_DIR}/dsn_block_service_test.xml"
GTEST_OUTPUT="xml:${output_xml}" ./dsn_block_service_test
11 changes: 11 additions & 0 deletions thirdparty/CMakeLists.txt
Original file line number Diff line number Diff line change
Expand Up @@ -56,6 +56,7 @@ message(STATUS "Setting up third-parties...")

file(MAKE_DIRECTORY ${TP_OUTPUT}/include)
file(MAKE_DIRECTORY ${TP_OUTPUT}/lib)
file(MAKE_DIRECTORY ${TP_OUTPUT}/hadoop)

option(ENABLE_ASAN "enable ASan" OFF)
message(STATUS "ENABLE_ASAN = ${ENABLE_ASAN}")
Expand All @@ -78,6 +79,16 @@ ExternalProject_Add(boost
DOWNLOAD_NO_PROGRESS true
)

# JuiceFS Hadoop JAR -> output/hadoop/juicefs-hadoop-1.3.1.jar (match pack_server.sh)
set(JUICEFS_JAR "output/hadoop/juicefs-hadoop-1.3.1.jar")
set(JUICEFS_DL "${TP_DIR}/build/Download/juicefs-hadoop/juicefs-hadoop-1.3.1.jar")
ExternalProject_Add(juicefs-hadoop
DOWNLOAD_COMMAND ${CMAKE_COMMAND} -E make_directory ${TP_DIR}/build/Download/juicefs-hadoop
COMMAND curl -L -o "${JUICEFS_DL}" "https://github.com/juicedata/juicefs/releases/download/v1.3.1/juicefs-hadoop-1.3.1.jar"
CONFIGURE_COMMAND "" BUILD_COMMAND ""
INSTALL_COMMAND ${CMAKE_COMMAND} -E copy "${JUICEFS_DL}" "${TP_DIR}/${JUICEFS_JAR}"
DOWNLOAD_NO_PROGRESS true)

# header-only
file(MAKE_DIRECTORY ${TP_OUTPUT}/include/concurrentqueue)
ExternalProject_Add(concurrentqueue
Expand Down
Loading