From 6d69366840e83f008b63f10d5c9dc01a4125f501 Mon Sep 17 00:00:00 2001 From: sitao Date: Sun, 12 Jan 2025 19:20:06 -0500 Subject: [PATCH 1/5] Remove get_address and pass in address as argument. --- src/spider/client/Driver.cpp | 14 +---- src/spider/core/Driver.hpp | 5 +- src/spider/io/BoostAsio.hpp | 31 ---------- src/spider/scheduler/scheduler.cpp | 17 ++++-- src/spider/storage/MetadataStorage.hpp | 1 - src/spider/storage/MysqlStorage.cpp | 70 +++++----------------- src/spider/storage/MysqlStorage.hpp | 1 - src/spider/worker/worker.cpp | 23 ++++--- tests/integration/test_client.py | 4 ++ tests/integration/test_scheduler_worker.py | 4 ++ tests/scheduler/test-SchedulerPolicy.cpp | 4 +- tests/storage/test-DataStorage.cpp | 8 +-- tests/storage/test-MetadataStorage.cpp | 6 +- tests/worker/test-FunctionManager.cpp | 2 +- tests/worker/test-TaskExecutor.cpp | 2 +- 15 files changed, 61 insertions(+), 131 deletions(-) diff --git a/src/spider/client/Driver.cpp b/src/spider/client/Driver.cpp index 02da9f754..85ffb064c 100644 --- a/src/spider/client/Driver.cpp +++ b/src/spider/client/Driver.cpp @@ -34,12 +34,7 @@ Driver::Driver(std::string const& storage_url) { throw ConnectionException(err.description); } - std::optional const optional_addr = core::get_address(); - if (!optional_addr.has_value()) { - throw ConnectionException("Cannot get machine address"); - } - std::string const& addr = optional_addr.value(); - err = m_metadata_storage->add_driver(core::Driver{m_id, addr}); + err = m_metadata_storage->add_driver(core::Driver{m_id}); if (!err.success()) { if (core::StorageErrType::DuplicateKeyErr == err.type) { throw DriverIdInUseException(m_id); @@ -72,12 +67,7 @@ Driver::Driver(std::string const& storage_url, boost::uuids::uuid const id) : m_ throw ConnectionException(err.description); } - std::optional const optional_addr = core::get_address(); - if (!optional_addr.has_value()) { - throw ConnectionException("Cannot get machine address"); - } - std::string const& addr = optional_addr.value(); - err = m_metadata_storage->add_driver(core::Driver{m_id, addr}); + err = m_metadata_storage->add_driver(core::Driver{m_id}); if (!err.success()) { if (core::StorageErrType::DuplicateKeyErr == err.type) { throw DriverIdInUseException(m_id); diff --git a/src/spider/core/Driver.hpp b/src/spider/core/Driver.hpp index f8996ac56..deb8890c1 100644 --- a/src/spider/core/Driver.hpp +++ b/src/spider/core/Driver.hpp @@ -10,15 +10,12 @@ namespace spider::core { class Driver { public: - Driver(boost::uuids::uuid const id, std::string addr) : m_id{id}, m_addr{std::move(addr)} {} + explicit Driver(boost::uuids::uuid const id) : m_id{id} {} [[nodiscard]] auto get_id() const -> boost::uuids::uuid const& { return m_id; } - [[nodiscard]] auto get_addr() const -> std::string const& { return m_addr; } - private: boost::uuids::uuid m_id; - std::string m_addr; }; class Scheduler { diff --git a/src/spider/io/BoostAsio.hpp b/src/spider/io/BoostAsio.hpp index 23731bc80..b983d7901 100644 --- a/src/spider/io/BoostAsio.hpp +++ b/src/spider/io/BoostAsio.hpp @@ -40,35 +40,4 @@ // IWYU pragma: end_exports // clang-format on -#include - -#include - -namespace spider::core { -inline auto get_address() -> std::optional { - try { - boost::asio::io_context io_context; - boost::asio::ip::tcp::resolver resolver(io_context); - auto const endpoints = resolver.resolve(boost::asio::ip::host_name(), ""); - for (auto const& endpoint : endpoints) { - if (endpoint.endpoint().address().is_v4() - && !endpoint.endpoint().address().is_loopback()) - { - return endpoint.endpoint().address().to_string(); - } - } - // If no non-loopback address found, return loopback address - spdlog::warn("No non-loopback address found, using loopback address"); - for (auto const& endpoint : endpoints) { - if (endpoint.endpoint().address().is_v4()) { - return endpoint.endpoint().address().to_string(); - } - } - return std::nullopt; - } catch (boost::system::system_error const& e) { - return std::nullopt; - } -} -} // namespace spider::core - #endif // SPIDER_CORE_BOOSTASIO_HPP diff --git a/src/spider/scheduler/scheduler.cpp b/src/spider/scheduler/scheduler.cpp index 216d64ec3..e72cbc784 100644 --- a/src/spider/scheduler/scheduler.cpp +++ b/src/spider/scheduler/scheduler.cpp @@ -42,6 +42,11 @@ namespace { auto parse_args(int const argc, char** argv) -> boost::program_options::variables_map { boost::program_options::options_description desc; desc.add_options()("help", "spider scheduler"); + desc.add_options()( + "host", + boost::program_options::value(), + "scheduler host address" + ); desc.add_options()( "port", boost::program_options::value(), @@ -136,6 +141,7 @@ auto main(int argc, char** argv) -> int { boost::program_options::variables_map const args = parse_args(argc, argv); unsigned short port = 0; + std::string scheduler_addr; std::string storage_url; try { if (!args.contains("port")) { @@ -143,6 +149,11 @@ auto main(int argc, char** argv) -> int { return cCmdArgParseErr; } port = args["port"].as(); + if (!args.contains("host")) { + spdlog::error("host is required"); + return cCmdArgParseErr; + } + scheduler_addr = args["host"].as(); if (!args.contains("storage_url")) { spdlog::error("storage_url is required"); return cCmdArgParseErr; @@ -185,12 +196,6 @@ auto main(int argc, char** argv) -> int { // Get scheduler id and addr boost::uuids::random_generator gen; boost::uuids::uuid const scheduler_id = gen(); - std::optional const optional_scheduler_addr = spider::core::get_address(); - if (!optional_scheduler_addr.has_value()) { - spdlog::error("Failed to get scheduler address"); - return cSchedulerAddrErr; - } - std::string const& scheduler_addr = optional_scheduler_addr.value(); // Start scheduler server spider::core::StopToken stop_token; diff --git a/src/spider/storage/MetadataStorage.hpp b/src/spider/storage/MetadataStorage.hpp index 4c5dbd54a..af1440b70 100644 --- a/src/spider/storage/MetadataStorage.hpp +++ b/src/spider/storage/MetadataStorage.hpp @@ -28,7 +28,6 @@ class MetadataStorage { virtual auto add_driver(Driver const& driver) -> StorageErr = 0; virtual auto add_scheduler(Scheduler const& scheduler) -> StorageErr = 0; - virtual auto get_driver(boost::uuids::uuid id, std::string* addr) -> StorageErr = 0; virtual auto get_active_scheduler(std::vector* schedulers) -> StorageErr = 0; virtual auto diff --git a/src/spider/storage/MysqlStorage.cpp b/src/spider/storage/MysqlStorage.cpp index 7ddec2a7c..3b5fa76ed 100644 --- a/src/spider/storage/MysqlStorage.cpp +++ b/src/spider/storage/MysqlStorage.cpp @@ -55,13 +55,13 @@ namespace spider::core { namespace { char const* const cCreateDriverTable = R"(CREATE TABLE IF NOT EXISTS `drivers` ( `id` BINARY(16) NOT NULL, - `address` VARCHAR(40) NOT NULL, `heartbeat` TIMESTAMP NOT NULL DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP, PRIMARY KEY (`id`) ))"; char const* const cCreateSchedulerTable = R"(CREATE TABLE IF NOT EXISTS `schedulers` ( `id` BINARY(16) NOT NULL, + `address` VARCHAR(40) NOT NULL, `port` INT UNSIGNED NOT NULL, `state` ENUM('normal', 'recovery', 'gc') NOT NULL, CONSTRAINT `scheduler_driver_id` FOREIGN KEY (`id`) REFERENCES `drivers` (`id`) ON UPDATE NO ACTION ON DELETE CASCADE, @@ -331,11 +331,10 @@ auto get_sql_string(sql::SQLString const& str) -> std::string { auto MySqlMetadataStorage::add_driver(Driver const& driver) -> StorageErr { try { std::unique_ptr statement( - m_conn->prepareStatement("INSERT INTO `drivers` (`id`, `address`) VALUES (?, ?)") + m_conn->prepareStatement("INSERT INTO `drivers` (`id`) VALUES (?)") ); sql::bytes id_bytes = uuid_get_bytes(driver.get_id()); statement->setBytes(1, &id_bytes); - statement->setString(2, driver.get_addr()); statement->executeUpdate(); } catch (sql::SQLException& e) { m_conn->rollback(); @@ -351,17 +350,18 @@ auto MySqlMetadataStorage::add_driver(Driver const& driver) -> StorageErr { auto MySqlMetadataStorage::add_scheduler(Scheduler const& scheduler) -> StorageErr { try { std::unique_ptr driver_statement( - m_conn->prepareStatement("INSERT INTO `drivers` (`id`, `address`) VALUES (?, ?)") + m_conn->prepareStatement("INSERT INTO `drivers` (`id`) VALUES (?)") ); sql::bytes id_bytes = uuid_get_bytes(scheduler.get_id()); driver_statement->setBytes(1, &id_bytes); - driver_statement->setString(2, scheduler.get_addr()); driver_statement->executeUpdate(); std::unique_ptr scheduler_statement(m_conn->prepareStatement( - "INSERT INTO `schedulers` (`id`, `port`, `state`) VALUES (?, ?, 'normal')" + "INSERT INTO `schedulers` (`id`, `address`, `port`, `state`) " + "VALUES (?, ?, ?, 'normal')" )); scheduler_statement->setBytes(1, &id_bytes); - scheduler_statement->setInt(2, scheduler.get_port()); + scheduler_statement->setString(2, scheduler.get_addr()); + scheduler_statement->setInt(3, scheduler.get_port()); scheduler_statement->executeUpdate(); } catch (sql::SQLException& e) { m_conn->rollback(); @@ -374,31 +374,6 @@ auto MySqlMetadataStorage::add_scheduler(Scheduler const& scheduler) -> StorageE return StorageErr{}; } -auto MySqlMetadataStorage::get_driver(boost::uuids::uuid id, std::string* addr) -> StorageErr { - try { - std::unique_ptr statement( - m_conn->prepareStatement("SELECT `address` FROM `drivers` WHERE `id` = ?") - ); - sql::bytes id_bytes = uuid_get_bytes(id); - statement->setBytes(1, &id_bytes); - std::unique_ptr res(statement->executeQuery()); - if (0 == res->rowsCount()) { - m_conn->rollback(); - return StorageErr{ - StorageErrType::KeyNotFoundErr, - fmt::format("no driver with id {}", boost::uuids::to_string(id)) - }; - } - res->next(); - *addr = get_sql_string(res->getString(1)); - } catch (sql::SQLException& e) { - m_conn->rollback(); - return StorageErr{StorageErrType::OtherErr, e.what()}; - } - m_conn->commit(); - return StorageErr{}; -} - auto MySqlMetadataStorage::get_active_scheduler(std::vector* schedulers) -> StorageErr { try { std::unique_ptr statement(m_conn->createStatement()); @@ -1513,35 +1488,22 @@ auto MySqlMetadataStorage::get_scheduler_state(boost::uuids::uuid id, std::strin auto MySqlMetadataStorage::get_scheduler_addr(boost::uuids::uuid id, std::string* addr, int* port) -> StorageErr { try { - std::unique_ptr addr_statement( - m_conn->prepareStatement("SELECT `address` FROM `drivers` WHERE `id` = ?") - ); + std::unique_ptr statement(m_conn->prepareStatement( + "SELECT `address`, `port` FROM `schedulers` WHERE `id` = ?" + )); sql::bytes id_bytes = uuid_get_bytes(id); - addr_statement->setBytes(1, &id_bytes); - std::unique_ptr addr_res{addr_statement->executeQuery()}; - if (addr_res->rowsCount() == 0) { - m_conn->rollback(); - return StorageErr{ - StorageErrType::KeyNotFoundErr, - fmt::format("no driver with id {}", boost::uuids::to_string(id)) - }; - } - std::unique_ptr port_statement( - m_conn->prepareStatement("SELECT `port` FROM `schedulers` WHERE `id` = ?") - ); - port_statement->setBytes(1, &id_bytes); - std::unique_ptr port_res{port_statement->executeQuery()}; - if (port_res->rowsCount() == 0) { + statement->setBytes(1, &id_bytes); + std::unique_ptr res{statement->executeQuery()}; + if (res->rowsCount() == 0) { m_conn->rollback(); return StorageErr{ StorageErrType::KeyNotFoundErr, fmt::format("no scheduler with id {}", boost::uuids::to_string(id)) }; } - addr_res->next(); - *addr = get_sql_string(addr_res->getString(1)); - port_res->next(); - *port = port_res->getInt(1); + res->next(); + *addr = get_sql_string(res->getString(1)); + *port = res->getInt(2); } catch (sql::SQLException& e) { m_conn->rollback(); return StorageErr{StorageErrType::OtherErr, e.what()}; diff --git a/src/spider/storage/MysqlStorage.hpp b/src/spider/storage/MysqlStorage.hpp index f035dcbf6..4625bfb44 100644 --- a/src/spider/storage/MysqlStorage.hpp +++ b/src/spider/storage/MysqlStorage.hpp @@ -34,7 +34,6 @@ class MySqlMetadataStorage : public MetadataStorage { auto initialize() -> StorageErr override; auto add_driver(Driver const& driver) -> StorageErr override; auto add_scheduler(Scheduler const& scheduler) -> StorageErr override; - auto get_driver(boost::uuids::uuid id, std::string* addr) -> StorageErr override; auto get_active_scheduler(std::vector* schedulers) -> StorageErr override; auto add_job(boost::uuids::uuid job_id, boost::uuids::uuid client_id, TaskGraph const& task_graph diff --git a/src/spider/worker/worker.cpp b/src/spider/worker/worker.cpp index 37d2e4210..64fa339ed 100644 --- a/src/spider/worker/worker.cpp +++ b/src/spider/worker/worker.cpp @@ -63,6 +63,7 @@ auto parse_args(int const argc, char** argv) -> boost::program_options::variable boost::program_options::value>(), "dynamic libraries that include the spider tasks" ); + desc.add_options()("host", boost::program_options::value(), "worker host address"); boost::program_options::variables_map variables; boost::program_options::store( @@ -332,12 +333,22 @@ auto main(int argc, char** argv) -> int { std::string storage_url; std::vector libs; + std::string worker_addr; try { - if (!args.contains("storage_url") || !args.contains("libs")) { - spdlog::error("Error: missing required arguments"); + if (!args.contains("storage_url")) { + spdlog::error("Missing storage_url"); return cCmdArgParseErr; } storage_url = args["storage_url"].as(); + if (!args.contains("host")) { + spdlog::error("Missing host"); + return cCmdArgParseErr; + } + worker_addr = args["host"].as(); + if (!args.contains("libs") || args["libs"].empty()) { + spdlog::error("Missing libs"); + return cCmdArgParseErr; + } libs = args["libs"].as>(); } catch (boost::bad_any_cast const& e) { spdlog::error("Error: {}", e.what()); @@ -362,16 +373,10 @@ auto main(int argc, char** argv) -> int { spdlog::error("Cannot connect to data storage: {}", err.description); return cStorageConnectionErr; } - std::optional const optional_worker_addr = spider::core::get_address(); - if (!optional_worker_addr.has_value()) { - spdlog::error("Failed to get worker address"); - return cWorkerAddrErr; - } - std::string const& worker_addr = optional_worker_addr.value(); boost::uuids::random_generator gen; boost::uuids::uuid const worker_id = gen(); - spider::core::Driver driver{worker_id, worker_addr}; + spider::core::Driver driver{worker_id}; err = metadata_store->add_driver(driver); if (!err.success()) { spdlog::error("Cannot add driver to metadata storage: {}", err.description); diff --git a/tests/integration/test_client.py b/tests/integration/test_client.py index 6dad18592..b6ebde74d 100644 --- a/tests/integration/test_client.py +++ b/tests/integration/test_client.py @@ -19,6 +19,8 @@ def start_scheduler_workers( dir_path = dir_path / ".." / ".." / "src" / "spider" scheduler_cmds = [ str(dir_path / "spider_scheduler"), + "--host", + "127.0.0.1", "--port", str(scheduler_port), "--storage_url", @@ -27,6 +29,8 @@ def start_scheduler_workers( scheduler_process = subprocess.Popen(scheduler_cmds) worker_cmds = [ str(dir_path / "spider_worker"), + "--host", + "127.0.0.1", "--storage_url", storage_url, "--libs", diff --git a/tests/integration/test_scheduler_worker.py b/tests/integration/test_scheduler_worker.py index d820bbed8..f6e9c87f4 100644 --- a/tests/integration/test_scheduler_worker.py +++ b/tests/integration/test_scheduler_worker.py @@ -34,6 +34,8 @@ def start_scheduler_worker( dir_path = dir_path / ".." / ".." / "src" / "spider" scheduler_cmds = [ str(dir_path / "spider_scheduler"), + "--host", + "127.0.0.1", "--port", str(scheduler_port), "--storage_url", @@ -42,6 +44,8 @@ def start_scheduler_worker( scheduler_process = subprocess.Popen(scheduler_cmds) worker_cmds = [ str(dir_path / "spider_worker"), + "--host", + "127.0.0.1", "--storage_url", storage_url, "--libs", diff --git a/tests/scheduler/test-SchedulerPolicy.cpp b/tests/scheduler/test-SchedulerPolicy.cpp index 1185869a7..4a2f5a5d2 100644 --- a/tests/scheduler/test-SchedulerPolicy.cpp +++ b/tests/scheduler/test-SchedulerPolicy.cpp @@ -94,7 +94,7 @@ TEMPLATE_LIST_TEST_CASE( spider::core::Data data{"value"}; data.set_hard_locality(true); data.set_locality({"127.0.0.1"}); - REQUIRE(metadata_store->add_driver(spider::core::Driver{client_id, "127.0.0.1"}).success()); + REQUIRE(metadata_store->add_driver(spider::core::Driver{client_id}).success()); REQUIRE(data_store->add_driver_data(client_id, data).success()); task.add_input(spider::core::TaskInput{data.get_id()}); spider::core::TaskGraph graph; @@ -141,7 +141,7 @@ TEMPLATE_LIST_TEST_CASE( spider::core::Data data; data.set_hard_locality(false); data.set_locality({"127.0.0.1"}); - REQUIRE(metadata_store->add_driver(spider::core::Driver{client_id, "127.0.0.1"}).success()); + REQUIRE(metadata_store->add_driver(spider::core::Driver{client_id}).success()); REQUIRE(data_store->add_driver_data(client_id, data).success()); task.add_input(spider::core::TaskInput{data.get_id()}); spider::core::TaskGraph graph; diff --git a/tests/storage/test-DataStorage.cpp b/tests/storage/test-DataStorage.cpp index e19f23fa1..f24d88ef7 100644 --- a/tests/storage/test-DataStorage.cpp +++ b/tests/storage/test-DataStorage.cpp @@ -25,7 +25,7 @@ TEMPLATE_LIST_TEST_CASE("Add, get and remove data", "[storage]", spider::test::S spider::core::Data const data{"value"}; boost::uuids::random_generator gen; boost::uuids::uuid const driver_id = gen(); - REQUIRE(metadata_storage->add_driver(spider::core::Driver{driver_id, "127.0.0.1"}).success()); + REQUIRE(metadata_storage->add_driver(spider::core::Driver{driver_id}).success()); REQUIRE(data_storage->add_driver_data(driver_id, data).success()); // Add data with same id again should fail @@ -57,7 +57,7 @@ TEMPLATE_LIST_TEST_CASE( // Add driver boost::uuids::random_generator gen; boost::uuids::uuid const driver_id = gen(); - REQUIRE(metadata_storage->add_driver(spider::core::Driver{driver_id, "127.0.0.1"}).success()); + REQUIRE(metadata_storage->add_driver(spider::core::Driver{driver_id}).success()); // Add data spider::core::KeyValueData const data{"key", "value", driver_id}; @@ -176,8 +176,8 @@ TEMPLATE_LIST_TEST_CASE( // Add driver boost::uuids::uuid const driver_id = gen(); boost::uuids::uuid const driver_id_2 = gen(); - REQUIRE(metadata_storage->add_driver(spider::core::Driver{driver_id, "127.0.0.1"}).success()); - REQUIRE(metadata_storage->add_driver(spider::core::Driver{driver_id_2, "127.0.0.1"}).success()); + REQUIRE(metadata_storage->add_driver(spider::core::Driver{driver_id}).success()); + REQUIRE(metadata_storage->add_driver(spider::core::Driver{driver_id_2}).success()); // Add driver reference without data should fail REQUIRE(!data_storage->add_driver_reference(gen(), driver_id).success()); diff --git a/tests/storage/test-MetadataStorage.cpp b/tests/storage/test-MetadataStorage.cpp index c96ad59dc..619e8a09a 100644 --- a/tests/storage/test-MetadataStorage.cpp +++ b/tests/storage/test-MetadataStorage.cpp @@ -31,11 +31,7 @@ TEMPLATE_LIST_TEST_CASE("Driver heartbeat", "[storage]", spider::test::MetadataS // Add driver should succeed boost::uuids::random_generator gen; boost::uuids::uuid const driver_id = gen(); - REQUIRE(storage->add_driver(spider::core::Driver{driver_id, "127.0.0.1"}).success()); - - std::string addr; - REQUIRE(storage->get_driver(driver_id, &addr).success()); - REQUIRE("127.0.0.1" == addr); + REQUIRE(storage->add_driver(spider::core::Driver{driver_id}).success()); std::vector ids{}; // Driver should not time out diff --git a/tests/worker/test-FunctionManager.cpp b/tests/worker/test-FunctionManager.cpp index cbcdee2c9..77d8f8245 100644 --- a/tests/worker/test-FunctionManager.cpp +++ b/tests/worker/test-FunctionManager.cpp @@ -153,7 +153,7 @@ TEMPLATE_LIST_TEST_CASE( spider::core::Data const data{std::string{buffer.data(), buffer.size()}}; boost::uuids::random_generator gen; boost::uuids::uuid const driver_id = gen(); - spider::core::Driver const driver{driver_id, "127.0.0.1"}; + spider::core::Driver const driver{driver_id}; REQUIRE(metadata_storage->add_driver(driver).success()); REQUIRE(data_storage->add_driver_data(driver_id, data).success()); diff --git a/tests/worker/test-TaskExecutor.cpp b/tests/worker/test-TaskExecutor.cpp index eccf0a1e3..15a464e75 100644 --- a/tests/worker/test-TaskExecutor.cpp +++ b/tests/worker/test-TaskExecutor.cpp @@ -152,7 +152,7 @@ TEMPLATE_LIST_TEST_CASE( spider::core::Data const data{std::string{buffer.data(), buffer.size()}}; boost::uuids::random_generator gen; boost::uuids::uuid const driver_id = gen(); - spider::core::Driver const driver{driver_id, "127.0.0.1"}; + spider::core::Driver const driver{driver_id}; REQUIRE(metadata_storage->add_driver(driver).success()); REQUIRE(data_storage->add_driver_data(driver_id, data).success()); From f299de354b0d1e9d220813d59dc06296f3ef3738 Mon Sep 17 00:00:00 2001 From: sitao Date: Sun, 12 Jan 2025 19:36:25 -0500 Subject: [PATCH 2/5] Fix integration tests for removing address in driver --- tests/integration/client.py | 5 +---- tests/integration/test_scheduler_worker.py | 22 +++++++++++++++++----- 2 files changed, 18 insertions(+), 9 deletions(-) diff --git a/tests/integration/client.py b/tests/integration/client.py index 27bda5b78..e3aca9f08 100644 --- a/tests/integration/client.py +++ b/tests/integration/client.py @@ -42,7 +42,6 @@ class TaskGraph: @dataclass class Driver: id: uuid.UUID - addr: str @dataclass @@ -180,9 +179,7 @@ def remove_job(conn, job_id: uuid.UUID): def add_driver(conn, driver: Driver): cursor = conn.cursor() - cursor.execute( - "INSERT INTO drivers (id, address) VALUES (%s, %s)", (driver.id.bytes, driver.addr) - ) + cursor.execute("INSERT INTO drivers (id) VALUES (%s)", (driver.id.bytes,)) conn.commit() cursor.close() diff --git a/tests/integration/test_scheduler_worker.py b/tests/integration/test_scheduler_worker.py index f6e9c87f4..976ccc251 100644 --- a/tests/integration/test_scheduler_worker.py +++ b/tests/integration/test_scheduler_worker.py @@ -106,9 +106,21 @@ def success_job(storage): ) submit_job(storage, uuid.uuid4(), graph) - assert get_task_state(storage, parent_1.id) == "ready" - assert get_task_state(storage, parent_2.id) == "ready" - assert get_task_state(storage, child.id) == "pending" + assert ( + get_task_state(storage, parent_1.id) == "ready" + or get_task_state(storage, parent_1.id) == "running" + or get_task_state(storage, parent_1.id) == "success" + ) + assert ( + get_task_state(storage, parent_2.id) == "ready" + or get_task_state(storage, parent_2.id) == "running" + or get_task_state(storage, parent_2.id) == "success" + ) + assert ( + get_task_state(storage, child.id) == "pending" + or get_task_state(storage, child.id) == "running" + or get_task_state(storage, child.id) == "success" + ) print("success job task ids:", parent_1.id, parent_2.id, child.id) yield graph, parent_1, parent_2, child @@ -144,7 +156,7 @@ def data_job(storage): id=uuid.uuid4(), value=msgpack.packb(2), ) - driver = Driver(id=uuid.uuid4(), addr="127.0.0.1") + driver = Driver(id=uuid.uuid4()) add_driver(storage, driver) add_driver_data(storage, driver, data) @@ -175,7 +187,7 @@ def random_fail_job(storage): id=uuid.uuid4(), value=msgpack.packb(2), ) - driver = Driver(id=uuid.uuid4(), addr="127.0.0.1") + driver = Driver(id=uuid.uuid4()) add_driver(storage, driver) add_driver_data(storage, driver, data) From f51a1dfe89ba51873a1161d2547d9ffe5f3c5ac6 Mon Sep 17 00:00:00 2001 From: sitao Date: Sun, 12 Jan 2025 19:53:46 -0500 Subject: [PATCH 3/5] Add host cmd argument in quick start guide --- docs/quick-start.md | 12 +++++++++--- 1 file changed, 9 insertions(+), 3 deletions(-) diff --git a/docs/quick-start.md b/docs/quick-start.md index a927c842d..96cc07cce 100644 --- a/docs/quick-start.md +++ b/docs/quick-start.md @@ -146,6 +146,7 @@ To start the scheduler, run: build/spider/src/spider/spider_scheduler \ --storage_url \ "jdbc:mariadb://localhost:3306/spider-storage?user=spider&password=password" \ + --host "127.0.0.1" \ --port 6000 ``` @@ -153,6 +154,7 @@ NOTE: * If you used a different set of arguments to set up the storage backend, ensure you update the `storage_url` argument in the command. +* In production, change the host to the real IP address of the machine running the scheduler. * If the scheduler fails to bind to port `6000`, change the port in the command and try again. ## Setting up a worker @@ -169,13 +171,17 @@ To start a worker, run: build/spider/src/spider/spider_worker \ --storage_url \ "jdbc:mariadb://localhost:3306/spider-storage?user=spider&password=password" \ - --port 6000 + --host "127.0.0.1" \ + --libs "build/libtasks.so" ``` NOTE: -If you used a different set of arguments to set up the storage backend, ensure you update the -`storage_url` argument in the command. +* If you used a different set of arguments to set up the storage backend, ensure you update the + `storage_url` argument in the command. +* In production, change the host to the real IP address of the machine running the worker. +* You can specify multiple task libraries to load. The task libraries must be built with linkage + to the Spider client library. > [!TIP] > You can start multiple workers to increase the number of concurrent tasks that can be run on the From 59eed9f643ce320c8a7665469ff6b5419100b480 Mon Sep 17 00:00:00 2001 From: sitao Date: Sun, 12 Jan 2025 20:55:23 -0500 Subject: [PATCH 4/5] Fix clang tidy --- src/spider/io/BoostAsio.hpp | 2 -- 1 file changed, 2 deletions(-) diff --git a/src/spider/io/BoostAsio.hpp b/src/spider/io/BoostAsio.hpp index b983d7901..01285c328 100644 --- a/src/spider/io/BoostAsio.hpp +++ b/src/spider/io/BoostAsio.hpp @@ -1,8 +1,6 @@ #ifndef SPIDER_CORE_BOOSTASIO_HPP #define SPIDER_CORE_BOOSTASIO_HPP -#include - // clang-format off // IWYU pragma: begin_exports From 3982db99f0507b8e1ebf9a83548d9808d2d3fa8e Mon Sep 17 00:00:00 2001 From: sitao Date: Sun, 12 Jan 2025 21:51:38 -0500 Subject: [PATCH 5/5] Fix clang tidy --- src/spider/scheduler/scheduler.cpp | 1 - 1 file changed, 1 deletion(-) diff --git a/src/spider/scheduler/scheduler.cpp b/src/spider/scheduler/scheduler.cpp index e72cbc784..8b778112c 100644 --- a/src/spider/scheduler/scheduler.cpp +++ b/src/spider/scheduler/scheduler.cpp @@ -3,7 +3,6 @@ #include #include #include -#include #include #include #include