From b7b4a7e7129bea36212d3b5496fa31cba7c41c55 Mon Sep 17 00:00:00 2001 From: sitao Date: Wed, 18 Dec 2024 23:14:05 -0500 Subject: [PATCH 01/18] Add TaskContext as first argument of task --- src/spider/CMakeLists.txt | 7 ++++- src/spider/worker/FunctionManager.hpp | 37 +++++++++++++++++++++------ src/spider/worker/task_executor.cpp | 4 ++- tests/CMakeLists.txt | 14 ++++++++-- tests/worker/test-FunctionManager.cpp | 19 +++++++++----- tests/worker/worker-test.cpp | 5 ++-- 6 files changed, 66 insertions(+), 20 deletions(-) diff --git a/src/spider/CMakeLists.txt b/src/spider/CMakeLists.txt index 4504bd913..9c48ec045 100644 --- a/src/spider/CMakeLists.txt +++ b/src/spider/CMakeLists.txt @@ -83,7 +83,12 @@ set(SPIDER_TASK_EXECUTOR_SOURCES add_executable(spider_task_executor) target_sources(spider_task_executor PRIVATE ${SPIDER_TASK_EXECUTOR_SOURCES}) -target_link_libraries(spider_task_executor PRIVATE spider_core) +target_link_libraries( + spider_task_executor + PRIVATE + spider_core + spider_client_lib +) target_link_libraries( spider_task_executor PRIVATE diff --git a/src/spider/worker/FunctionManager.hpp b/src/spider/worker/FunctionManager.hpp index 6d012e485..97c318aae 100644 --- a/src/spider/worker/FunctionManager.hpp +++ b/src/spider/worker/FunctionManager.hpp @@ -17,6 +17,8 @@ #include #include +#include "../client/task.hpp" +#include "../client/TaskContext.hpp" #include "../io/MsgPack.hpp" // IWYU pragma: keep #include "../io/Serializer.hpp" #include "TaskExecutorMessage.hpp" @@ -36,7 +38,7 @@ using ArgsBuffer = msgpack::sbuffer; using ResultBuffer = msgpack::sbuffer; -using Function = std::function; +using Function = std::function; using FunctionMap = absl::flat_hash_map; @@ -255,12 +257,25 @@ inline auto create_args_request(std::vector const& args_buffer template class FunctionInvoker { public: - static auto apply(F const& function, ArgsBuffer const& args_buffer) -> ResultBuffer { + static auto + apply(F const& function, TaskContext context, ArgsBuffer const& args_buffer) -> ResultBuffer { // NOLINTBEGIN(cppcoreguidelines-pro-type-union-access,cppcoreguidelines-pro-bounds-pointer-arithmetic) using ArgsTuple = signature::args_t; using ReturnType = signature::ret_t; - ArgsTuple args_tuple{}; + static_assert(TaskIo, "Return type must be TaskIo"); + static_assert( + std::is_same_v>, + "First argument must be TaskContext" + ); + for_n - 1>([&](auto i) { + static_assert( + TaskIo>, + "Other arguments must be TaskIo" + ); + }); + + ArgsTuple args_tuple; try { msgpack::object_handle const handle = msgpack::unpack(args_buffer.data(), args_buffer.size()); @@ -273,7 +288,7 @@ class FunctionInvoker { ); } - if (std::tuple_size_v != object.via.array.size) { + if (std::tuple_size_v - 1 != object.via.array.size) { return create_error_response( FunctionInvokeError::WrongNumberOfArguments, fmt::format( @@ -284,10 +299,11 @@ class FunctionInvoker { ); } - for_n>([&](auto i) { + std::get<0>(args_tuple) = context; + for_n - 1>([&](auto i) { msgpack::object arg = object.via.array.ptr[i.cValue]; - std::get(args_tuple) - = arg.as>(); + std::get(args_tuple) + = arg.as>(); }); } catch (msgpack::type_error& e) { return create_error_response( @@ -334,7 +350,12 @@ class FunctionManager { return m_map .emplace( name, - std::bind(&FunctionInvoker::apply, std::move(f), std::placeholders::_1) + std::bind( + &FunctionInvoker::apply, + std::move(f), + std::placeholders::_1, + std::placeholders::_2 + ) ) .second; } diff --git a/src/spider/worker/task_executor.cpp b/src/spider/worker/task_executor.cpp index 7901808f4..f60cc2d64 100644 --- a/src/spider/worker/task_executor.cpp +++ b/src/spider/worker/task_executor.cpp @@ -16,6 +16,7 @@ #include // IWYU pragma: keep #include +#include "../client/TaskContext.hpp" #include "../io/BoostAsio.hpp" // IWYU pragma: keep #include "../io/MsgPack.hpp" // IWYU pragma: keep #include "DllLoader.hpp" @@ -126,7 +127,8 @@ auto main(int const argc, char** argv) -> int { ); return cResultSendErr; } - msgpack::sbuffer const result_buffer = (*function)(args_buffer); + spider::TaskContext task_context{}; + msgpack::sbuffer const result_buffer = (*function)(task_context, args_buffer); spdlog::debug("Function executed"); // Write result buffer to stdout diff --git a/tests/CMakeLists.txt b/tests/CMakeLists.txt index 4734ff0d7..7efe04160 100644 --- a/tests/CMakeLists.txt +++ b/tests/CMakeLists.txt @@ -31,7 +31,12 @@ endforeach() target_sources(unitTest PRIVATE ${SPIDER_TEST_SCHEDULER_SOURCES}) target_link_libraries(unitTest PRIVATE Catch2::Catch2WithMain) -target_link_libraries(unitTest PRIVATE spider_core) +target_link_libraries( + unitTest + PRIVATE + spider_core + spider_client_lib +) target_link_libraries( unitTest PRIVATE @@ -46,7 +51,12 @@ add_dependencies(unitTest worker_test) add_library(worker_test SHARED) target_sources(worker_test PRIVATE worker/worker-test.cpp) -target_link_libraries(worker_test PRIVATE spider_core) +target_link_libraries( + worker_test + PRIVATE + spider_core + spider_client_lib +) add_custom_target(integrationTest ALL) add_custom_command( diff --git a/tests/worker/test-FunctionManager.cpp b/tests/worker/test-FunctionManager.cpp index 34e041674..4341a1c81 100644 --- a/tests/worker/test-FunctionManager.cpp +++ b/tests/worker/test-FunctionManager.cpp @@ -5,15 +5,17 @@ #include +#include "../../src/spider/client/TaskContext.hpp" #include "../../src/spider/io/MsgPack.hpp" // IWYU pragma: keep #include "../../src/spider/worker/FunctionManager.hpp" namespace { -auto int_test(int const x, int const y) -> int { +auto int_test(spider::TaskContext /*context*/, int const x, int const y) -> int { return x + y; } -auto tuple_ret_test(std::string const& str, int const x) -> std::tuple { +auto tuple_ret_test(spider::TaskContext /*context*/, std::string const& str, int const x) + -> std::tuple { return std::make_tuple(str, x); } @@ -28,18 +30,21 @@ TEST_CASE("Register and run function with POD inputs", "[core]") { // Get the registered function should succeed spider::core::Function const* function = manager.get_function("int_test"); + REQUIRE(nullptr != function); + + spider::TaskContext context{}; // Run function with two ints should succeed spider::core::ArgsBuffer const args_buffers = spider::core::create_args_buffers(2, 3); constexpr int cExpected = 2 + 3; - msgpack::sbuffer const result = (*function)(args_buffers); + msgpack::sbuffer const result = (*function)(context, args_buffers); msgpack::sbuffer buffer{}; msgpack::pack(buffer, cExpected); REQUIRE(cExpected == spider::core::response_get_result(result).value_or(0)); // Run function with wrong number of inputs should fail spider::core::ArgsBuffer wrong_args_buffers = spider::core::create_args_buffers(1); - msgpack::sbuffer wrong_result = (*function)(wrong_args_buffers); + msgpack::sbuffer wrong_result = (*function)(context, wrong_args_buffers); std::optional> wrong_result_option = spider::core::response_get_error(wrong_result); REQUIRE(wrong_result_option.has_value()); @@ -50,7 +55,7 @@ TEST_CASE("Register and run function with POD inputs", "[core]") { // Run function with wrong type of inputs should fail wrong_args_buffers = spider::core::create_args_buffers(0, "test"); - wrong_result = (*function)(wrong_args_buffers); + wrong_result = (*function)(context, wrong_args_buffers); wrong_result_option = spider::core::response_get_error(wrong_result); REQUIRE(wrong_result_option.has_value()); if (wrong_result_option.has_value()) { @@ -60,12 +65,14 @@ TEST_CASE("Register and run function with POD inputs", "[core]") { } TEST_CASE("Register and run function with tuple return", "[core]") { + spider::TaskContext context{}; + spider::core::FunctionManager const& manager = spider::core::FunctionManager::get_instance(); spider::core::Function const* function = manager.get_function("tuple_ret_test"); spider::core::ArgsBuffer const args_buffers = spider::core::create_args_buffers("test", 3); - msgpack::sbuffer const result = (*function)(args_buffers); + msgpack::sbuffer const result = (*function)(context, args_buffers); REQUIRE(std::make_tuple("test", 3) == spider::core::response_get_result(result).value_or( std::make_tuple("", 0) diff --git a/tests/worker/worker-test.cpp b/tests/worker/worker-test.cpp index 04d97cf1a..30dd2219e 100644 --- a/tests/worker/worker-test.cpp +++ b/tests/worker/worker-test.cpp @@ -1,15 +1,16 @@ #include #include +#include "../../src/spider/client/TaskContext.hpp" #include "../../src/spider/worker/FunctionManager.hpp" namespace { -auto sum_test(int const x, int const y) -> int { +auto sum_test(spider::TaskContext /*context*/, int const x, int const y) -> int { std::cerr << x << " + " << y << " = " << x + y << "\n"; return x + y; } -auto error_test(int const /*x*/) -> int { +auto error_test(spider::TaskContext /*context*/, int const /*x*/) -> int { throw std::runtime_error("Simulated error in worker"); } } // namespace From aa8115c1f5c789ca8d06bcd59e5c7414542f06ef Mon Sep 17 00:00:00 2001 From: sitao Date: Thu, 19 Dec 2024 10:20:14 -0500 Subject: [PATCH 02/18] Add data storage set locality and basic client data implementation --- src/spider/CMakeLists.txt | 1 + src/spider/client/Data.cpp | 49 +++++++++++++++++++++++++++++ src/spider/client/Data.hpp | 13 ++++++-- src/spider/core/Data.hpp | 2 +- src/spider/storage/DataStorage.hpp | 1 + src/spider/storage/MysqlStorage.cpp | 30 ++++++++++++++++++ src/spider/storage/MysqlStorage.hpp | 1 + 7 files changed, 94 insertions(+), 3 deletions(-) create mode 100644 src/spider/client/Data.cpp diff --git a/src/spider/CMakeLists.txt b/src/spider/CMakeLists.txt index 9c48ec045..ad8f86313 100644 --- a/src/spider/CMakeLists.txt +++ b/src/spider/CMakeLists.txt @@ -1,5 +1,6 @@ # set variable as CACHE INTERNAL to access it from other scope set(SPIDER_CORE_SOURCES + client/Data.cpp storage/MysqlStorage.cpp worker/FunctionManager.cpp io/msgpack_message.cpp diff --git a/src/spider/client/Data.cpp b/src/spider/client/Data.cpp new file mode 100644 index 000000000..925733206 --- /dev/null +++ b/src/spider/client/Data.cpp @@ -0,0 +1,49 @@ +#include "Data.hpp" + +#include +#include + +#include "../core/Data.hpp" +#include "../io/MsgPack.hpp" // IWYU pragma: keep +#include "../io/Serializer.hpp" +#include "../storage/DataStorage.hpp" + +namespace spider { + +template +auto Data::get() -> T { + std::string const& value = m_impl->get_value(); + return msgpack::unpack(value.data(), value.size()).get().as(); +} + +template +auto Data::set_locality(std::vector const& nodes, bool hard) -> void { + m_impl->set_locality(nodes); + m_impl->set_hard_locality(hard); + // TODO: update data storage +} + +template +auto Data::Builder::set_locality(std::vector const& nodes, bool hard) -> Builder& { + m_nodes = nodes; + m_hard_locality = hard; + return *this; +} + +template +auto Data::Builder::set_cleanup_func(std::function const& f) -> Builder& { + m_cleanup_func = f; + return *this; +} + +template +auto Data::Builder::build(T const& t) -> Data { + auto impl = std::make_shared(t); + impl->set_locality(m_nodes); + impl->set_hard_locality(m_hard_locality); + impl->set_cleanup_func(m_cleanup_func); + // TODO: update data storage + return Data(impl); +} + +} // namespace spider diff --git a/src/spider/client/Data.hpp b/src/spider/client/Data.hpp index 8ee4181e6..2f694bd58 100644 --- a/src/spider/client/Data.hpp +++ b/src/spider/client/Data.hpp @@ -9,7 +9,11 @@ #include "../io/Serializer.hpp" namespace spider { -class DataImpl; + +namespace core { +class Data; +class DataStorage; +} // namespace core /** * A representation of data stored on external storage. This class allows the user to define: @@ -72,10 +76,15 @@ class Data { * @throw spider::ConnectionException */ auto build(T const& t) -> Data; + + private: + std::vector m_nodes; + bool m_hard_locality = false; + std::function m_cleanup_func; }; private: - std::unique_ptr m_impl; + std::unique_ptr m_impl; }; } // namespace spider diff --git a/src/spider/core/Data.hpp b/src/spider/core/Data.hpp index 2ca23f8ae..360012285 100644 --- a/src/spider/core/Data.hpp +++ b/src/spider/core/Data.hpp @@ -21,7 +21,7 @@ class Data { [[nodiscard]] auto get_id() const -> boost::uuids::uuid { return m_id; } - [[nodiscard]] auto get_value() const -> std::string { return m_value; } + [[nodiscard]] auto get_value() const -> std::string const& { return m_value; } [[nodiscard]] auto get_locality() const -> std::vector const& { return m_locality; diff --git a/src/spider/storage/DataStorage.hpp b/src/spider/storage/DataStorage.hpp index 29be26e54..a3f86be9f 100644 --- a/src/spider/storage/DataStorage.hpp +++ b/src/spider/storage/DataStorage.hpp @@ -26,6 +26,7 @@ class DataStorage { virtual auto add_driver_data(boost::uuids::uuid driver_id, Data const& data) -> StorageErr = 0; virtual auto add_task_data(boost::uuids::uuid task_id, Data const& data) -> StorageErr = 0; virtual auto get_data(boost::uuids::uuid id, Data* data) -> StorageErr = 0; + virtual auto set_data_locality(Data const& data) -> StorageErr = 0; virtual auto remove_data(boost::uuids::uuid id) -> StorageErr = 0; virtual auto add_task_reference(boost::uuids::uuid id, boost::uuids::uuid task_id) -> StorageErr = 0; diff --git a/src/spider/storage/MysqlStorage.cpp b/src/spider/storage/MysqlStorage.cpp index a9efbe498..5b56fd732 100644 --- a/src/spider/storage/MysqlStorage.cpp +++ b/src/spider/storage/MysqlStorage.cpp @@ -1475,6 +1475,36 @@ auto MySqlDataStorage::get_data(boost::uuids::uuid id, Data* data) -> StorageErr return StorageErr{}; } +auto MySqlDataStorage::set_data_locality(Data const& data) -> StorageErr { + try { + std::unique_ptr const delete_statement( + m_conn->prepareStatement("DELETE FROM `data_locality` WHERE `id` = ?") + ); + sql::bytes id_bytes = uuid_get_bytes(data.get_id()); + delete_statement->setBytes(1, &id_bytes); + delete_statement->executeUpdate(); + std::unique_ptr const insert_statement(m_conn->prepareStatement( + "INSERT INTO `data_locality` (`id`, `address`) VALUES(?, ?)" + )); + for (std::string const& addr : data.get_locality()) { + insert_statement->setBytes(1, &id_bytes); + insert_statement->setString(2, addr); + insert_statement->executeUpdate(); + } + std::unique_ptr const hard_locality_statement( + m_conn->prepareStatement("UPDATE `data` SET `hard_locality` = ? WHERE `id` = ?") + ); + hard_locality_statement->setBoolean(1, data.is_hard_locality()); + hard_locality_statement->setBytes(2, &id_bytes); + hard_locality_statement->executeUpdate(); + } catch (sql::SQLException& e) { + m_conn->rollback(); + return StorageErr{StorageErrType::OtherErr, e.what()}; + } + m_conn->commit(); + return StorageErr{}; +} + auto MySqlDataStorage::remove_data(boost::uuids::uuid id) -> StorageErr { try { std::unique_ptr statement( diff --git a/src/spider/storage/MysqlStorage.hpp b/src/spider/storage/MysqlStorage.hpp index 79244a70e..c729d3888 100644 --- a/src/spider/storage/MysqlStorage.hpp +++ b/src/spider/storage/MysqlStorage.hpp @@ -88,6 +88,7 @@ class MySqlDataStorage : public DataStorage { auto add_driver_data(boost::uuids::uuid driver_id, Data const& data) -> StorageErr override; auto add_task_data(boost::uuids::uuid task_id, Data const& data) -> StorageErr override; auto get_data(boost::uuids::uuid id, Data* data) -> StorageErr override; + auto set_data_locality(Data const& data) -> StorageErr override; auto remove_data(boost::uuids::uuid id) -> StorageErr override; auto add_task_reference(boost::uuids::uuid id, boost::uuids::uuid task_id) -> StorageErr override; From 0330fde50b19286ef568a3e91324791d172b7274 Mon Sep 17 00:00:00 2001 From: sitaowang1998 Date: Thu, 19 Dec 2024 12:43:04 -0500 Subject: [PATCH 03/18] Add serializer for data --- src/spider/CMakeLists.txt | 1 + src/spider/client/Data.hpp | 2 ++ src/spider/io/DataSerializer.hpp | 19 +++++++++++++++++++ 3 files changed, 22 insertions(+) create mode 100644 src/spider/io/DataSerializer.hpp diff --git a/src/spider/CMakeLists.txt b/src/spider/CMakeLists.txt index ad8f86313..3c7c2e8ae 100644 --- a/src/spider/CMakeLists.txt +++ b/src/spider/CMakeLists.txt @@ -4,6 +4,7 @@ set(SPIDER_CORE_SOURCES storage/MysqlStorage.cpp worker/FunctionManager.cpp io/msgpack_message.cpp + io/DataSerializer.hpp CACHE INTERNAL "spider core source files" ) diff --git a/src/spider/client/Data.hpp b/src/spider/client/Data.hpp index 2f694bd58..2345b9ff3 100644 --- a/src/spider/client/Data.hpp +++ b/src/spider/client/Data.hpp @@ -83,6 +83,8 @@ class Data { std::function m_cleanup_func; }; + auto get_impl() const -> std::unique_ptr const& { return m_impl; } + private: std::unique_ptr m_impl; }; diff --git a/src/spider/io/DataSerializer.hpp b/src/spider/io/DataSerializer.hpp new file mode 100644 index 000000000..84dec5be6 --- /dev/null +++ b/src/spider/io/DataSerializer.hpp @@ -0,0 +1,19 @@ +#ifndef SPIDER_CLIENT_DATASERIALIZER_HPP +#define SPIDER_CLIENT_DATASERIALIZER_HPP + +#include "../client/Data.hpp" +#include "../core/Data.hpp" // IWYU pragma: keep +#include "MsgPack.hpp" // IWYU pragma: keep + +template +struct msgpack::adaptor::pack> { + template + auto operator()(msgpack::packer& packer, spider::Data const& data) const + -> msgpack::packer& { + packer.pack_map(1); + packer.pack(data.get_impl()->get_id()); + return packer; + } +}; + +#endif From b9d9db3f0b8f5c074bda7e7dc6c378b68b6e0fb2 Mon Sep 17 00:00:00 2001 From: sitaowang1998 Date: Thu, 19 Dec 2024 12:45:48 -0500 Subject: [PATCH 04/18] Fix clang tidy --- src/spider/client/Data.cpp | 1 - src/spider/worker/task_executor.cpp | 2 +- tests/worker/test-FunctionManager.cpp | 4 ++-- 3 files changed, 3 insertions(+), 4 deletions(-) diff --git a/src/spider/client/Data.cpp b/src/spider/client/Data.cpp index 925733206..081d7d06f 100644 --- a/src/spider/client/Data.cpp +++ b/src/spider/client/Data.cpp @@ -6,7 +6,6 @@ #include "../core/Data.hpp" #include "../io/MsgPack.hpp" // IWYU pragma: keep #include "../io/Serializer.hpp" -#include "../storage/DataStorage.hpp" namespace spider { diff --git a/src/spider/worker/task_executor.cpp b/src/spider/worker/task_executor.cpp index f60cc2d64..69247317f 100644 --- a/src/spider/worker/task_executor.cpp +++ b/src/spider/worker/task_executor.cpp @@ -127,7 +127,7 @@ auto main(int const argc, char** argv) -> int { ); return cResultSendErr; } - spider::TaskContext task_context{}; + spider::TaskContext const task_context{}; msgpack::sbuffer const result_buffer = (*function)(task_context, args_buffer); spdlog::debug("Function executed"); diff --git a/tests/worker/test-FunctionManager.cpp b/tests/worker/test-FunctionManager.cpp index 4341a1c81..88ad4abdd 100644 --- a/tests/worker/test-FunctionManager.cpp +++ b/tests/worker/test-FunctionManager.cpp @@ -32,7 +32,7 @@ TEST_CASE("Register and run function with POD inputs", "[core]") { spider::core::Function const* function = manager.get_function("int_test"); REQUIRE(nullptr != function); - spider::TaskContext context{}; + spider::TaskContext const context{}; // Run function with two ints should succeed spider::core::ArgsBuffer const args_buffers = spider::core::create_args_buffers(2, 3); @@ -65,7 +65,7 @@ TEST_CASE("Register and run function with POD inputs", "[core]") { } TEST_CASE("Register and run function with tuple return", "[core]") { - spider::TaskContext context{}; + spider::TaskContext const context{}; spider::core::FunctionManager const& manager = spider::core::FunctionManager::get_instance(); From 84d5d7fb9823b6f99a60948054cf1ca697698d58 Mon Sep 17 00:00:00 2001 From: sitao Date: Thu, 19 Dec 2024 13:57:32 -0500 Subject: [PATCH 05/18] Fix clang tidy --- src/spider/client/Data.cpp | 1 + src/spider/client/Data.hpp | 2 +- 2 files changed, 2 insertions(+), 1 deletion(-) diff --git a/src/spider/client/Data.cpp b/src/spider/client/Data.cpp index 081d7d06f..4ece324ea 100644 --- a/src/spider/client/Data.cpp +++ b/src/spider/client/Data.cpp @@ -1,5 +1,6 @@ #include "Data.hpp" +#include #include #include diff --git a/src/spider/client/Data.hpp b/src/spider/client/Data.hpp index 2345b9ff3..1d181005d 100644 --- a/src/spider/client/Data.hpp +++ b/src/spider/client/Data.hpp @@ -83,7 +83,7 @@ class Data { std::function m_cleanup_func; }; - auto get_impl() const -> std::unique_ptr const& { return m_impl; } + [[nodiscard]] auto get_impl() const -> std::unique_ptr const& { return m_impl; } private: std::unique_ptr m_impl; From 988f2c7a632ffdbbb1820b205a9e4ad5f129e720 Mon Sep 17 00:00:00 2001 From: sitao Date: Thu, 19 Dec 2024 20:18:22 -0500 Subject: [PATCH 06/18] Add data implementation --- src/spider/client/Data.cpp | 30 +++++++++++++++++++++++------- src/spider/client/Data.hpp | 33 ++++++++++++++++++++++++++++++++- 2 files changed, 55 insertions(+), 8 deletions(-) diff --git a/src/spider/client/Data.cpp b/src/spider/client/Data.cpp index 4ece324ea..2e6f94e10 100644 --- a/src/spider/client/Data.cpp +++ b/src/spider/client/Data.cpp @@ -7,6 +7,8 @@ #include "../core/Data.hpp" #include "../io/MsgPack.hpp" // IWYU pragma: keep #include "../io/Serializer.hpp" +#include "../storage/DataStorage.hpp" +#include "Exception.hpp" namespace spider { @@ -20,7 +22,7 @@ template auto Data::set_locality(std::vector const& nodes, bool hard) -> void { m_impl->set_locality(nodes); m_impl->set_hard_locality(hard); - // TODO: update data storage + m_data_store->set_data_locality(*m_impl); } template @@ -38,12 +40,26 @@ auto Data::Builder::set_cleanup_func(std::function const& f) template auto Data::Builder::build(T const& t) -> Data { - auto impl = std::make_shared(t); - impl->set_locality(m_nodes); - impl->set_hard_locality(m_hard_locality); - impl->set_cleanup_func(m_cleanup_func); - // TODO: update data storage - return Data(impl); + auto data = std::make_unique(t); + data->set_locality(m_nodes); + data->set_hard_locality(m_hard_locality); + data->set_cleanup_func(m_cleanup_func); + core::StorageErr err; + switch (m_data_source) { + case DataSource::Driver: + err = m_data_store->add_driver_data(m_source_id, *data); + if (!err.success()) { + throw ConnectionException(err.description); + } + break; + case DataSource::TaskContext: + err = m_data_store->add_task_data(m_source_id, *data); + if (!err.success()) { + throw ConnectionException(err.description); + } + break; + } + return Data{data, m_data_store}; } } // namespace spider diff --git a/src/spider/client/Data.hpp b/src/spider/client/Data.hpp index 1d181005d..7496702b6 100644 --- a/src/spider/client/Data.hpp +++ b/src/spider/client/Data.hpp @@ -6,7 +6,10 @@ #include #include +#include + #include "../io/Serializer.hpp" +#include "../storage/DataStorage.hpp" namespace spider { @@ -44,6 +47,7 @@ class Data { * @param nodes * @param hard Whether the data is only accessible from the given nodes (i.e., the locality is a * hard requirement). + * @throw spider::ConnectionException */ void set_locality(std::vector const& nodes, bool hard); @@ -78,15 +82,42 @@ class Data { auto build(T const& t) -> Data; private: + enum class DataSource { + Driver, + TaskContext + }; + explicit Builder( + std::shared_ptr data_store, + boost::uuids::uuid const source_id, + DataSource const data_source + ) + : m_data_store{std::move(data_store)}, + m_source_id{source_id}, + m_data_source{data_source} {}; + std::vector m_nodes; bool m_hard_locality = false; std::function m_cleanup_func; + + std::shared_ptr m_data_store; + boost::uuids::uuid m_source_id; + DataSource m_data_source; + + friend class Driver; + friend class TaskContext; }; +private: + Data(std::unique_ptr impl, std::shared_ptr data_store) + : m_impl{std::move(impl)}, + m_data_store{std::move(data_store)} {} + [[nodiscard]] auto get_impl() const -> std::unique_ptr const& { return m_impl; } -private: std::unique_ptr m_impl; + std::shared_ptr m_data_store; + + friend class msgpack::adaptor::pack; }; } // namespace spider From 7a424a3166c034cdee84ebbbaac714003827c8ec Mon Sep 17 00:00:00 2001 From: sitao Date: Thu, 19 Dec 2024 21:07:12 -0500 Subject: [PATCH 07/18] Add create data builder interface in driver and task context --- src/spider/client/Driver.hpp | 6 ++++++ src/spider/client/TaskContext.hpp | 7 +++++++ 2 files changed, 13 insertions(+) diff --git a/src/spider/client/Driver.hpp b/src/spider/client/Driver.hpp index 96b1c3b36..35f2e2182 100644 --- a/src/spider/client/Driver.hpp +++ b/src/spider/client/Driver.hpp @@ -52,6 +52,12 @@ class Driver { */ Driver(std::string const& storage_url, boost::uuids::uuid id); + /** + * @return Data builder. + */ + template + auto get_data_builder() -> Data::Builder; + /** * Inserts the given key-value pair into the key-value store, overwriting any existing value. * diff --git a/src/spider/client/TaskContext.hpp b/src/spider/client/TaskContext.hpp index db8363670..cbba39e07 100644 --- a/src/spider/client/TaskContext.hpp +++ b/src/spider/client/TaskContext.hpp @@ -7,6 +7,7 @@ #include +#include "Data.hpp" #include "Job.hpp" #include "task.hpp" #include "TaskGraph.hpp" @@ -32,6 +33,12 @@ class TaskContext { */ [[nodiscard]] auto get_id() const -> boost::uuids::uuid; + /** + * @return Data builder. + */ + template + auto get_data_builder() -> Data::Builder; + /** * Inserts the given key-value pair into the key-value store, overwriting any existing value. * From c6ee8fe44ed24e8b37718a217c41e96a54ecadaf Mon Sep 17 00:00:00 2001 From: sitao Date: Thu, 19 Dec 2024 22:07:15 -0500 Subject: [PATCH 08/18] Remove direct Data serializer becuase data serialization is ambiguous --- src/spider/CMakeLists.txt | 3 +-- src/spider/io/DataSerializer.hpp | 19 ------------------- 2 files changed, 1 insertion(+), 21 deletions(-) delete mode 100644 src/spider/io/DataSerializer.hpp diff --git a/src/spider/CMakeLists.txt b/src/spider/CMakeLists.txt index 3c7c2e8ae..dd56a21a1 100644 --- a/src/spider/CMakeLists.txt +++ b/src/spider/CMakeLists.txt @@ -4,8 +4,7 @@ set(SPIDER_CORE_SOURCES storage/MysqlStorage.cpp worker/FunctionManager.cpp io/msgpack_message.cpp - io/DataSerializer.hpp - CACHE INTERNAL + CACHE INTERNAL "spider core source files" ) diff --git a/src/spider/io/DataSerializer.hpp b/src/spider/io/DataSerializer.hpp deleted file mode 100644 index 84dec5be6..000000000 --- a/src/spider/io/DataSerializer.hpp +++ /dev/null @@ -1,19 +0,0 @@ -#ifndef SPIDER_CLIENT_DATASERIALIZER_HPP -#define SPIDER_CLIENT_DATASERIALIZER_HPP - -#include "../client/Data.hpp" -#include "../core/Data.hpp" // IWYU pragma: keep -#include "MsgPack.hpp" // IWYU pragma: keep - -template -struct msgpack::adaptor::pack> { - template - auto operator()(msgpack::packer& packer, spider::Data const& data) const - -> msgpack::packer& { - packer.pack_map(1); - packer.pack(data.get_impl()->get_id()); - return packer; - } -}; - -#endif From a791b1ed9f6ca094be8281049d46221b8dbddeb6 Mon Sep 17 00:00:00 2001 From: sitao Date: Fri, 20 Dec 2024 00:30:48 -0500 Subject: [PATCH 09/18] Add data serializer --- src/spider/CMakeLists.txt | 1 + src/spider/client/Data.hpp | 3 ++- src/spider/io/DataSerializer.hpp | 26 ++++++++++++++++++++++++++ 3 files changed, 29 insertions(+), 1 deletion(-) create mode 100644 src/spider/io/DataSerializer.hpp diff --git a/src/spider/CMakeLists.txt b/src/spider/CMakeLists.txt index dd56a21a1..7f518611c 100644 --- a/src/spider/CMakeLists.txt +++ b/src/spider/CMakeLists.txt @@ -20,6 +20,7 @@ set(SPIDER_CORE_HEADERS io/MsgPack.hpp io/msgpack_message.hpp io/Serializer.hpp + io/DataSerializer.hpp utils/TimedCache.hpp storage/MetadataStorage.hpp storage/DataStorage.hpp diff --git a/src/spider/client/Data.hpp b/src/spider/client/Data.hpp index 7496702b6..45085ab57 100644 --- a/src/spider/client/Data.hpp +++ b/src/spider/client/Data.hpp @@ -16,6 +16,7 @@ namespace spider { namespace core { class Data; class DataStorage; +class DataSerializer; } // namespace core /** @@ -117,7 +118,7 @@ class Data { std::unique_ptr m_impl; std::shared_ptr m_data_store; - friend class msgpack::adaptor::pack; + friend class core::DataSerializer; }; } // namespace spider diff --git a/src/spider/io/DataSerializer.hpp b/src/spider/io/DataSerializer.hpp new file mode 100644 index 000000000..886427d05 --- /dev/null +++ b/src/spider/io/DataSerializer.hpp @@ -0,0 +1,26 @@ +#ifndef SPIDER_CORE_DATASEIALIZER_HPP +#define SPIDER_CORE_DATASEIALIZER_HPP + +#include + +#include "../client/Data.hpp" +#include "MsgPack.hpp" // IWYU pragma: keep +#include "Serializer.hpp" // IWYU pragma: keep + +namespace spider::core { +class DataSerializer { +public: + template + static auto serialize_id(msgpack::packer& packer, spider::Data const& data) -> void { + packer.pack(data.get_impl()->get_id()); + } + + template + static auto data_get_id(spider::Data const& data) -> boost::uuids::uuid { + return data.get_impl()->get_id(); + } +}; + +} // namespace spider::core + +#endif From c34e0ebe6425051ef6191eef015b18a6516704f6 Mon Sep 17 00:00:00 2001 From: sitao Date: Fri, 20 Dec 2024 01:09:44 -0500 Subject: [PATCH 10/18] Reformat cmake files --- src/spider/CMakeLists.txt | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/src/spider/CMakeLists.txt b/src/spider/CMakeLists.txt index 7f518611c..54f98b498 100644 --- a/src/spider/CMakeLists.txt +++ b/src/spider/CMakeLists.txt @@ -4,7 +4,7 @@ set(SPIDER_CORE_SOURCES storage/MysqlStorage.cpp worker/FunctionManager.cpp io/msgpack_message.cpp - CACHE INTERNAL + CACHE INTERNAL "spider core source files" ) From cb662e78811b7f9768568703ce93d22644a0c3b9 Mon Sep 17 00:00:00 2001 From: sitao Date: Fri, 20 Dec 2024 14:35:20 -0500 Subject: [PATCH 11/18] Add task id into args buffer message --- src/spider/worker/FunctionManager.hpp | 12 ++++++--- src/spider/worker/TaskExecutor.hpp | 7 +++-- src/spider/worker/task_executor.cpp | 38 ++++++++++++++++++++++++--- 3 files changed, 47 insertions(+), 10 deletions(-) diff --git a/src/spider/worker/FunctionManager.hpp b/src/spider/worker/FunctionManager.hpp index 97c318aae..bdbb559e7 100644 --- a/src/spider/worker/FunctionManager.hpp +++ b/src/spider/worker/FunctionManager.hpp @@ -229,23 +229,27 @@ auto create_args_buffers(Args&&... args) -> ArgsBuffer { } template -auto create_args_request(Args&&... args) -> msgpack::sbuffer { +auto create_args_request(boost::uuids::uuid task_id, Args&&... args) -> msgpack::sbuffer { msgpack::sbuffer buffer; msgpack::packer packer{buffer}; packer.pack_array(2); packer.pack(worker::TaskExecutorRequestType::Arguments); - packer.pack_array(sizeof...(args)); + packer.pack_array(sizeof...(args) + 1); + packer.pack(task_id); ([&] { packer.pack(args); }(), ...); return buffer; } -inline auto create_args_request(std::vector const& args_buffers +inline auto create_args_request( + boost::uuids::uuid task_id, + std::vector const& args_buffers ) -> msgpack::sbuffer { msgpack::sbuffer buffer; msgpack::packer packer{buffer}; packer.pack_array(2); packer.pack(worker::TaskExecutorRequestType::Arguments); - packer.pack_array(args_buffers.size()); + packer.pack_array(args_buffers.size() + 1); + packer.pack(task_id); for (msgpack::sbuffer const& args_buffer : args_buffers) { buffer.write(args_buffer.data(), args_buffer.size()); } diff --git a/src/spider/worker/TaskExecutor.hpp b/src/spider/worker/TaskExecutor.hpp index 25fc93b14..4f0ed7141 100644 --- a/src/spider/worker/TaskExecutor.hpp +++ b/src/spider/worker/TaskExecutor.hpp @@ -16,6 +16,7 @@ #include #include #include +#include #include "../io/BoostAsio.hpp" // IWYU pragma: keep #include "../io/MsgPack.hpp" // IWYU pragma: keep @@ -37,6 +38,7 @@ class TaskExecutor { template TaskExecutor( boost::asio::io_context& context, + boost::uuids::uuid task_id, std::string const& func_name, std::vector const& libs, absl::flat_hash_map< @@ -69,12 +71,13 @@ class TaskExecutor { // Send args msgpack::sbuffer const args_request - = core::create_args_request(std::forward(args)...); + = core::create_args_request(task_id, std::forward(args)...); send_message(m_write_pipe, args_request); } TaskExecutor( boost::asio::io_context& context, + boost::uuids::uuid task_id, std::string const& func_name, std::vector const& libs, absl::flat_hash_map< @@ -106,7 +109,7 @@ class TaskExecutor { boost::asio::co_spawn(context, process_output_handler(), boost::asio::detached); // Send args - msgpack::sbuffer const args_request = core::create_args_request(args_buffers); + msgpack::sbuffer const args_request = core::create_args_request(task_id, args_buffers); send_message(m_write_pipe, args_request); } diff --git a/src/spider/worker/task_executor.cpp b/src/spider/worker/task_executor.cpp index 69247317f..e9db29c40 100644 --- a/src/spider/worker/task_executor.cpp +++ b/src/spider/worker/task_executor.cpp @@ -19,6 +19,9 @@ #include "../client/TaskContext.hpp" #include "../io/BoostAsio.hpp" // IWYU pragma: keep #include "../io/MsgPack.hpp" // IWYU pragma: keep +#include "../storage/DataStorage.hpp" +#include "../storage/MetadataStorage.hpp" +#include "../storage/MysqlStorage.hpp" #include "DllLoader.hpp" #include "FunctionManager.hpp" #include "message_pipe.hpp" @@ -35,6 +38,11 @@ auto parse_arg(int const argc, char** const& argv) -> boost::program_options::va boost::program_options::value>(), "dynamic libraries that include the spider tasks" ); + desc.add_options()( + "storage_url", + boost::program_options::value(), + "storage server url" + ); boost::program_options::variables_map variables; boost::program_options::store( @@ -49,10 +57,11 @@ auto parse_arg(int const argc, char** const& argv) -> boost::program_options::va } // namespace constexpr int cCmdArgParseErr = 1; -constexpr int cDllErr = 2; -constexpr int cFuncArgParseErr = 3; -constexpr int cResultSendErr = 4; -constexpr int cOtherErr = 5; +constexpr int cStorageErr = 2; +constexpr int cDllErr = 3; +constexpr int cFuncArgParseErr = 4; +constexpr int cResultSendErr = 5; +constexpr int cOtherErr = 6; auto main(int const argc, char** argv) -> int { // Set up spdlog to write to stderr @@ -66,11 +75,16 @@ auto main(int const argc, char** argv) -> int { boost::program_options::variables_map const args = parse_arg(argc, argv); std::string func_name; + std::string storage_url; try { if (!args.contains("func")) { return cCmdArgParseErr; } func_name = args["func"].as(); + if (!args.contains("storage_url")) { + return cCmdArgParseErr; + } + storage_url = args["storage_url"].as(); if (!args.contains("libs")) { return cCmdArgParseErr; } @@ -90,6 +104,22 @@ auto main(int const argc, char** argv) -> int { spdlog::debug("Function to run: {}", func_name); try { + // Set up storage + std::shared_ptr metadata_store + = std::make_shared(); + std::shared_ptr data_store + = std::make_shared(); + spider::core::StorageErr err = metadata_store->connect(storage_url); + if (!err.success()) { + spdlog::error("Failed to connect to storage {}", storage_url); + return cStorageErr; + } + err = data_store->connect(storage_url); + if (!err.success()) { + spdlog::error("Failed to connect to storage {}", storage_url); + return cStorageErr; + } + // Set up asio boost::asio::io_context context; boost::asio::posix::stream_descriptor in(context, dup(STDIN_FILENO)); From ef1104bb99e59d5ab24fa3337adc2ec8aa1f7a22 Mon Sep 17 00:00:00 2001 From: sitao Date: Fri, 20 Dec 2024 21:17:02 -0500 Subject: [PATCH 12/18] Add data support in task --- src/spider/CMakeLists.txt | 3 +- src/spider/client/Data.hpp | 6 ++- src/spider/client/TaskContext.hpp | 25 +++++++++++ src/spider/core/Data.hpp | 2 - src/spider/core/DataImpl.hpp | 27 ++++++++++++ src/spider/core/TaskContextImpl.hpp | 32 ++++++++++++++ src/spider/io/DataSerializer.hpp | 26 ----------- src/spider/worker/FunctionManager.hpp | 63 ++++++++++++++++++--------- src/spider/worker/TaskExecutor.hpp | 15 ++++--- src/spider/worker/task_executor.cpp | 3 +- src/spider/worker/worker.cpp | 12 ++++- tests/worker/test-FunctionManager.cpp | 35 ++++++++++++--- tests/worker/test-TaskExecutor.cpp | 38 ++++++++++++---- 13 files changed, 212 insertions(+), 75 deletions(-) create mode 100644 src/spider/core/DataImpl.hpp create mode 100644 src/spider/core/TaskContextImpl.hpp delete mode 100644 src/spider/io/DataSerializer.hpp diff --git a/src/spider/CMakeLists.txt b/src/spider/CMakeLists.txt index 54f98b498..a444ca029 100644 --- a/src/spider/CMakeLists.txt +++ b/src/spider/CMakeLists.txt @@ -11,16 +11,17 @@ set(SPIDER_CORE_SOURCES set(SPIDER_CORE_HEADERS core/Error.hpp core/Data.hpp + core/DataImpl.hpp core/Driver.hpp core/KeyValueData.hpp core/Task.hpp + core/TaskContextImpl.hpp core/TaskGraph.hpp core/JobMetadata.hpp io/BoostAsio.hpp io/MsgPack.hpp io/msgpack_message.hpp io/Serializer.hpp - io/DataSerializer.hpp utils/TimedCache.hpp storage/MetadataStorage.hpp storage/DataStorage.hpp diff --git a/src/spider/client/Data.hpp b/src/spider/client/Data.hpp index 45085ab57..df1318387 100644 --- a/src/spider/client/Data.hpp +++ b/src/spider/client/Data.hpp @@ -16,7 +16,7 @@ namespace spider { namespace core { class Data; class DataStorage; -class DataSerializer; +class DataImpl; } // namespace core /** @@ -108,6 +108,8 @@ class Data { friend class TaskContext; }; + Data() = default; + private: Data(std::unique_ptr impl, std::shared_ptr data_store) : m_impl{std::move(impl)}, @@ -118,7 +120,7 @@ class Data { std::unique_ptr m_impl; std::shared_ptr m_data_store; - friend class core::DataSerializer; + friend class core::DataImpl; }; } // namespace spider diff --git a/src/spider/client/TaskContext.hpp b/src/spider/client/TaskContext.hpp index cbba39e07..115cd7128 100644 --- a/src/spider/client/TaskContext.hpp +++ b/src/spider/client/TaskContext.hpp @@ -13,6 +13,12 @@ #include "TaskGraph.hpp" namespace spider { +namespace core { +class DataStorage; +class MetadataStorage; +class TaskContextImpl; +} // namespace core + /** * TaskContext provides a task with all Spider functionalities, e.g. getting task instance id, * accessing data storage, creating and waiting for new jobs, etc. @@ -118,6 +124,25 @@ class TaskContext { * @throw spider::ConnectionException */ auto get_jobs() -> std::vector; + + TaskContext() = default; + +private: + TaskContext( + std::shared_ptr data_store, + std::shared_ptr metadata_store + ) + : m_data_store{std::move(data_store)}, + m_metadata_store{std::move(metadata_store)} {} + + auto get_data_store() -> std::shared_ptr { return m_data_store; } + + auto get_metadata_store() -> std::shared_ptr { return m_metadata_store; } + + std::shared_ptr m_data_store; + std::shared_ptr m_metadata_store; + + friend class core::TaskContextImpl; }; } // namespace spider diff --git a/src/spider/core/Data.hpp b/src/spider/core/Data.hpp index 360012285..d278794b0 100644 --- a/src/spider/core/Data.hpp +++ b/src/spider/core/Data.hpp @@ -17,8 +17,6 @@ class Data { Data(boost::uuids::uuid const id, std::string value) : m_id(id), m_value(std::move(value)) {} - static auto is_data() -> bool { return true; } - [[nodiscard]] auto get_id() const -> boost::uuids::uuid { return m_id; } [[nodiscard]] auto get_value() const -> std::string const& { return m_value; } diff --git a/src/spider/core/DataImpl.hpp b/src/spider/core/DataImpl.hpp new file mode 100644 index 000000000..1a7c339cd --- /dev/null +++ b/src/spider/core/DataImpl.hpp @@ -0,0 +1,27 @@ +#ifndef SPIDER_CORE_DATAIMPL_HPP +#define SPIDER_CORE_DATAIMPL_HPP + +#include + +#include "../client/Data.hpp" +#include "../core/Data.hpp" + +namespace spider::core { + +class DataImpl { +public: + template + static auto create_data(std::unique_ptr data, std::shared_ptr data_store) + -> spider::Data { + return spider::Data{std::move(data), data_store}; + } + + template + static auto get_impl(spider::Data const& data) -> std::shared_ptr { + return data.get_impl(); + } +}; + +} // namespace spider::core + +#endif diff --git a/src/spider/core/TaskContextImpl.hpp b/src/spider/core/TaskContextImpl.hpp new file mode 100644 index 000000000..91d72fb77 --- /dev/null +++ b/src/spider/core/TaskContextImpl.hpp @@ -0,0 +1,32 @@ +#ifndef SPIDER_CORE_TASKCONTEXTIMPL_HPP +#define SPIDER_CORE_TASKCONTEXTIMPL_HPP + +#include + +#include "../client/TaskContext.hpp" +#include "../storage/DataStorage.hpp" +#include "../storage/MetadataStorage.hpp" + +namespace spider::core { +class TaskContextImpl { +public: + static auto create_task_context( + std::shared_ptr const& data_storage, + std::shared_ptr const& metadata_storage + ) -> TaskContext { + return TaskContext{data_storage, metadata_storage}; + } + + static auto get_data_store(TaskContext const& task_context) -> std::shared_ptr { + return task_context.m_data_store; + } + + static auto get_metadata_store(TaskContext const& task_context + ) -> std::shared_ptr { + return task_context.m_metadata_store; + } +}; + +} // namespace spider::core + +#endif diff --git a/src/spider/io/DataSerializer.hpp b/src/spider/io/DataSerializer.hpp deleted file mode 100644 index 886427d05..000000000 --- a/src/spider/io/DataSerializer.hpp +++ /dev/null @@ -1,26 +0,0 @@ -#ifndef SPIDER_CORE_DATASEIALIZER_HPP -#define SPIDER_CORE_DATASEIALIZER_HPP - -#include - -#include "../client/Data.hpp" -#include "MsgPack.hpp" // IWYU pragma: keep -#include "Serializer.hpp" // IWYU pragma: keep - -namespace spider::core { -class DataSerializer { -public: - template - static auto serialize_id(msgpack::packer& packer, spider::Data const& data) -> void { - packer.pack(data.get_impl()->get_id()); - } - - template - static auto data_get_id(spider::Data const& data) -> boost::uuids::uuid { - return data.get_impl()->get_id(); - } -}; - -} // namespace spider::core - -#endif diff --git a/src/spider/worker/FunctionManager.hpp b/src/spider/worker/FunctionManager.hpp index bdbb559e7..d751577e2 100644 --- a/src/spider/worker/FunctionManager.hpp +++ b/src/spider/worker/FunctionManager.hpp @@ -6,6 +6,7 @@ #include #include #include +#include #include #include #include @@ -19,8 +20,12 @@ #include "../client/task.hpp" #include "../client/TaskContext.hpp" +#include "../core/DataImpl.hpp" +#include "../core/TaskContextImpl.hpp" #include "../io/MsgPack.hpp" // IWYU pragma: keep #include "../io/Serializer.hpp" +#include "../storage/DataStorage.hpp" +#include "../storage/MetadataStorage.hpp" #include "TaskExecutorMessage.hpp" // NOLINTBEGIN(cppcoreguidelines-macro-usage) @@ -42,6 +47,17 @@ using Function = std::function; +template +struct TemplateParameter; + +template