protocol.cpp raw

   1  // Copyright (c) 2021-present The Bitcoin Core developers
   2  // Distributed under the MIT software license, see the accompanying
   3  // file COPYING or http://www.opensource.org/licenses/mit-license.php.
   4  
   5  #include <interfaces/init.h>
   6  #include <ipc/capnp/context.h>
   7  #include <ipc/capnp/init.capnp.h>
   8  #include <ipc/capnp/init.capnp.proxy.h>
   9  #include <ipc/capnp/protocol.h>
  10  #include <ipc/exception.h>
  11  #include <ipc/protocol.h>
  12  #include <kj/async.h>
  13  #include <mp/proxy-io.h>
  14  #include <mp/proxy-types.h>
  15  #include <mp/util.h>
  16  #include <util/log.h>
  17  #include <util/threadnames.h>
  18  
  19  #include <cassert>
  20  #include <cerrno>
  21  #include <future>
  22  #include <memory>
  23  #include <mutex>
  24  #include <optional>
  25  #include <string>
  26  #include <sys/socket.h>
  27  #include <system_error>
  28  #include <thread>
  29  
  30  namespace ipc {
  31  namespace capnp {
  32  namespace {
  33  
  34  mp::Log GetRequestedIPCLogLevel()
  35  {
  36      if (util::log::ShouldTraceLog(BCLog::IPC)) return mp::Log::Trace;
  37      if (util::log::ShouldDebugLog(BCLog::IPC)) return mp::Log::Debug;
  38  
  39      // Info, Warning, and Error are logged unconditionally
  40      return mp::Log::Info;
  41  }
  42  
  43  void IpcLogFn(mp::LogMessage message)
  44  {
  45      switch (message.level) {
  46      case mp::Log::Trace:
  47          LogTrace(BCLog::IPC, "%s", message.message);
  48          return;
  49      case mp::Log::Debug:
  50          LogDebug(BCLog::IPC, "%s", message.message);
  51          return;
  52      case mp::Log::Info:
  53          LogInfo("ipc: %s", message.message);
  54          return;
  55      case mp::Log::Warning:
  56          LogWarning("ipc: %s", message.message);
  57          return;
  58      case mp::Log::Error:
  59          LogError("ipc: %s", message.message);
  60          return;
  61      case mp::Log::Raise:
  62          LogError("ipc: %s", message.message);
  63          throw Exception(message.message);
  64      } // no default case, so the compiler can warn about missing cases
  65  
  66      // Be conservative and assume that if MP ever adds a new log level, it
  67      // should only be shown at our most verbose level.
  68      LogTrace(BCLog::IPC, "%s", message.message);
  69  }
  70  
  71  class CapnpProtocol : public Protocol
  72  {
  73  public:
  74      ~CapnpProtocol() noexcept(true)
  75      {
  76          m_loop_ref.reset();
  77          if (m_loop_thread.joinable()) m_loop_thread.join();
  78          assert(!m_loop);
  79      };
  80      std::unique_ptr<interfaces::Init> connect(int fd, const char* exe_name) override
  81      {
  82          startLoop(exe_name);
  83          return mp::ConnectStream<messages::Init>(*m_loop, fd);
  84      }
  85      void listen(int listen_fd, const char* exe_name, interfaces::Init& init) override
  86      {
  87          startLoop(exe_name);
  88          if (::listen(listen_fd, /*backlog=*/5) != 0) {
  89              throw std::system_error(errno, std::system_category());
  90          }
  91          mp::ListenConnections<messages::Init>(*m_loop, listen_fd, init);
  92      }
  93      void serve(int fd, const char* exe_name, interfaces::Init& init, const std::function<void()>& ready_fn = {}) override
  94      {
  95          assert(!m_loop);
  96          mp::g_thread_context.thread_name = mp::ThreadName(exe_name);
  97          mp::LogOptions opts = {
  98              .log_fn = IpcLogFn,
  99              .log_level = GetRequestedIPCLogLevel()
 100          };
 101          m_loop.emplace(exe_name, std::move(opts), &m_context);
 102          if (ready_fn) ready_fn();
 103          mp::ServeStream<messages::Init>(*m_loop, fd, init);
 104          m_parent_connection = &m_loop->m_incoming_connections.back();
 105          m_loop->loop();
 106          m_loop.reset();
 107      }
 108      void disconnectIncoming() override
 109      {
 110          if (!m_loop) return;
 111          // Delete incoming connections, except the connection to a parent
 112          // process (if there is one), since a parent process should be able to
 113          // monitor and control this process, even during shutdown.
 114          m_loop->sync([&] {
 115              m_loop->m_incoming_connections.remove_if([this](mp::Connection& c) { return &c != m_parent_connection; });
 116          });
 117      }
 118      void addCleanup(std::type_index type, void* iface, std::function<void()> cleanup) override
 119      {
 120          mp::ProxyTypeRegister::types().at(type)(iface).cleanup_fns.emplace_back(std::move(cleanup));
 121      }
 122      Context& context() override { return m_context; }
 123      void startLoop(const char* exe_name)
 124      {
 125          if (m_loop) return;
 126          std::promise<void> promise;
 127          m_loop_thread = std::thread([&] {
 128              util::ThreadRename("capnp-loop");
 129              mp::LogOptions opts = {
 130                  .log_fn = IpcLogFn,
 131                  .log_level = GetRequestedIPCLogLevel()
 132              };
 133              m_loop.emplace(exe_name, std::move(opts), &m_context);
 134              m_loop_ref.emplace(*m_loop);
 135              promise.set_value();
 136              m_loop->loop();
 137              m_loop.reset();
 138          });
 139          promise.get_future().wait();
 140      }
 141      Context m_context;
 142      std::thread m_loop_thread;
 143      //! EventLoop object which manages I/O events for all connections.
 144      std::optional<mp::EventLoop> m_loop;
 145      //! Reference to the same EventLoop. Increments the loop’s refcount on
 146      //! creation, decrements on destruction. The loop thread exits when the
 147      //! refcount reaches 0. Other IPC objects also hold their own EventLoopRef.
 148      std::optional<mp::EventLoopRef> m_loop_ref;
 149      //! Connection to parent, if this is a child process spawned by a parent process.
 150      mp::Connection* m_parent_connection{nullptr};
 151  };
 152  } // namespace
 153  
 154  std::unique_ptr<Protocol> MakeCapnpProtocol() { return std::make_unique<CapnpProtocol>(); }
 155  } // namespace capnp
 156  } // namespace ipc
 157