diff --git a/src/ray/object_manager/object_manager.cc b/src/ray/object_manager/object_manager.cc index 927c24893..53f7fc378 100644 --- a/src/ray/object_manager/object_manager.cc +++ b/src/ray/object_manager/object_manager.cc @@ -78,6 +78,12 @@ ObjectManager::ObjectManager(asio::io_service &main_service, const ClientID &sel ObjectManager::~ObjectManager() { StopRpcService(); } +void ObjectManager::Stop() { + if (plasma::plasma_store_runner != nullptr) { + plasma::plasma_store_runner->Stop(); + } +} + void ObjectManager::RunRpcService() { rpc_service_.run(); } void ObjectManager::StartRpcService() { diff --git a/src/ray/object_manager/object_manager.h b/src/ray/object_manager/object_manager.h index a1b821f8b..a193639cd 100644 --- a/src/ray/object_manager/object_manager.h +++ b/src/ray/object_manager/object_manager.h @@ -184,6 +184,10 @@ class ObjectManager : public ObjectManagerInterface, ~ObjectManager(); + /// Stop the Plasma Store eventloop. Currently it is only used to handle + /// signals from Raylet. + void Stop(); + /// Subscribe to notifications of objects added to local store. /// Upon subscribing, the callback will be invoked for all objects that /// diff --git a/src/ray/object_manager/plasma/store_runner.cc b/src/ray/object_manager/plasma/store_runner.cc index 492bffb94..e8ca85378 100644 --- a/src/ray/object_manager/plasma/store_runner.cc +++ b/src/ray/object_manager/plasma/store_runner.cc @@ -1,7 +1,6 @@ #include "ray/object_manager/plasma/store_runner.h" #include -#include #include #include #include @@ -17,20 +16,11 @@ using arrow::util::ArrowLogLevel; void SetMallocGranularity(int value); -void HandleSignal(int signal) { - if (signal == SIGTERM) { - ARROW_LOG(INFO) << "SIGTERM Signal received, closing Plasma Server..."; - plasma_store_runner->Stop(); - } -} - PlasmaStoreRunner::PlasmaStoreRunner(std::string socket_name, int64_t system_memory, bool hugepages_enabled, std::string plasma_directory, const std::string external_store_endpoint): hugepages_enabled_(hugepages_enabled), external_store_endpoint_(external_store_endpoint) { ArrowLog::StartArrowLog("plasma_store", ArrowLogLevel::ARROW_INFO); - ArrowLog::InstallFailureSignalHandler(); - // Sanity check. if (socket_name.empty()) { ARROW_LOG(FATAL) << "please specify socket for incoming connections with -s switch"; @@ -109,13 +99,6 @@ void PlasmaStoreRunner::Start() { } ARROW_LOG(DEBUG) << "starting server listening on " << socket_name_; -#ifndef _WIN32 // TODO(mehrdadn): Is there an equivalent of this we need for Windows? - // Ignore SIGPIPE signals. If we don't do this, then when we attempt to write - // to a client that has already died, the store could die. - signal(SIGPIPE, SIG_IGN); -#endif - signal(SIGTERM, HandleSignal); - // Create the event loop. loop_.reset(new EventLoop); store_.reset(new PlasmaStore(loop_.get(), plasma_directory_, hugepages_enabled_, @@ -147,7 +130,6 @@ void PlasmaStoreRunner::Start() { #ifdef _WINSOCKAPI_ WSACleanup(); #endif - ArrowLog::UninstallSignalAction(); ArrowLog::ShutDownArrowLog(); } diff --git a/src/ray/plasma/store_exec.cc b/src/ray/plasma/store_exec.cc index 6c687dd0d..3b689c6e6 100644 --- a/src/ray/plasma/store_exec.cc +++ b/src/ray/plasma/store_exec.cc @@ -1,4 +1,5 @@ #include +#include #include #include @@ -11,6 +12,15 @@ #define __STDC_FORMAT_MACROS #endif +using arrow::util::ArrowLog; + +void HandleSignal(int signal) { + if (signal == SIGTERM) { + ARROW_LOG(INFO) << "SIGTERM Signal received, closing Plasma Server..."; + plasma::plasma_store_runner->Stop(); + } +} + int main(int argc, char *argv[]) { std::string socket_name; // Directory where plasma memory mapped files are stored. @@ -50,11 +60,20 @@ int main(int argc, char *argv[]) { } if (!keep_idle) { + ArrowLog::InstallFailureSignalHandler(); plasma::plasma_store_runner.reset( new plasma::PlasmaStoreRunner(socket_name, system_memory, hugepages_enabled, plasma_directory, external_store_endpoint)); + // Install signal handler before starting the eventloop. +#ifndef _WIN32 // TODO(mehrdadn): Is there an equivalent of this we need for Windows? + // Ignore SIGPIPE signals. If we don't do this, then when we attempt to write + // to a client that has already died, the store could die. + signal(SIGPIPE, SIG_IGN); +#endif + signal(SIGTERM, HandleSignal); plasma::plasma_store_runner->Start(); plasma::plasma_store_runner.reset(); + ArrowLog::UninstallSignalAction(); } else { printf( "The Plasma Store is started with the '-z' flag, " diff --git a/src/ray/raylet/raylet.cc b/src/ray/raylet/raylet.cc index 5fcbdcc34..4e28b9c3f 100644 --- a/src/ray/raylet/raylet.cc +++ b/src/ray/raylet/raylet.cc @@ -90,6 +90,7 @@ void Raylet::Start() { } void Raylet::Stop() { + object_manager_.Stop(); RAY_CHECK_OK(gcs_client_->Nodes().UnregisterSelf()); acceptor_.close(); }