scheduler.cpp raw

   1  // Copyright (c) 2015-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 <scheduler.h>
   6  
   7  #include <sync.h>
   8  #include <util/time.h>
   9  
  10  #include <cassert>
  11  #include <functional>
  12  #include <utility>
  13  
  14  CScheduler::CScheduler() = default;
  15  
  16  CScheduler::~CScheduler()
  17  {
  18      assert(nThreadsServicingQueue == 0);
  19      if (stopWhenEmpty) assert(taskQueue.empty());
  20  }
  21  
  22  
  23  void CScheduler::serviceQueue()
  24  {
  25      WAIT_LOCK(newTaskMutex, lock);
  26      ++nThreadsServicingQueue;
  27  
  28      // newTaskMutex is locked throughout this loop EXCEPT
  29      // when the thread is waiting or when the user's function
  30      // is called.
  31      while (!shouldStop()) {
  32          try {
  33              while (!shouldStop() && taskQueue.empty()) {
  34                  // Wait until there is something to do.
  35                  newTaskScheduled.wait(lock);
  36              }
  37  
  38              // Wait until either there is a new task, or until
  39              // the time of the first item on the queue:
  40  
  41              while (!shouldStop() && !taskQueue.empty()) {
  42                  std::chrono::steady_clock::time_point timeToWaitFor = taskQueue.begin()->first;
  43                  if (newTaskScheduled.wait_until(lock, timeToWaitFor) == std::cv_status::timeout) {
  44                      break; // Exit loop after timeout, it means we reached the time of the event
  45                  }
  46              }
  47  
  48              // If there are multiple threads, the queue can empty while we're waiting (another
  49              // thread may service the task we were waiting on).
  50              if (shouldStop() || taskQueue.empty())
  51                  continue;
  52  
  53              Function f = taskQueue.begin()->second;
  54              taskQueue.erase(taskQueue.begin());
  55  
  56              {
  57                  // Unlock before calling f, so it can reschedule itself or another task
  58                  // without deadlocking:
  59                  REVERSE_LOCK(lock, newTaskMutex);
  60                  f();
  61              }
  62          } catch (...) {
  63              --nThreadsServicingQueue;
  64              throw;
  65          }
  66      }
  67      --nThreadsServicingQueue;
  68      newTaskScheduled.notify_one();
  69  }
  70  
  71  void CScheduler::schedule(CScheduler::Function f, std::chrono::steady_clock::time_point t)
  72  {
  73      {
  74          LOCK(newTaskMutex);
  75          taskQueue.insert(std::make_pair(t, f));
  76      }
  77      newTaskScheduled.notify_one();
  78  }
  79  
  80  void CScheduler::MockForward(std::chrono::seconds delta_seconds)
  81  {
  82      assert(delta_seconds > 0s && delta_seconds <= 1h);
  83  
  84      {
  85          LOCK(newTaskMutex);
  86  
  87          // use temp_queue to maintain updated schedule
  88          std::multimap<std::chrono::steady_clock::time_point, Function> temp_queue;
  89  
  90          for (const auto& element : taskQueue) {
  91              temp_queue.emplace_hint(temp_queue.cend(), element.first - delta_seconds, element.second);
  92          }
  93  
  94          // point taskQueue to temp_queue
  95          taskQueue = std::move(temp_queue);
  96      }
  97  
  98      // notify that the taskQueue needs to be processed
  99      newTaskScheduled.notify_one();
 100  }
 101  
 102  static void Repeat(CScheduler& s, CScheduler::Function f, std::chrono::milliseconds delta)
 103  {
 104      f();
 105      s.scheduleFromNow([=, &s] { Repeat(s, f, delta); }, delta);
 106  }
 107  
 108  void CScheduler::scheduleEvery(CScheduler::Function f, std::chrono::milliseconds delta)
 109  {
 110      scheduleFromNow([this, f, delta] { Repeat(*this, f, delta); }, delta);
 111  }
 112  
 113  size_t CScheduler::getQueueInfo(std::chrono::steady_clock::time_point& first,
 114                                  std::chrono::steady_clock::time_point& last) const
 115  {
 116      LOCK(newTaskMutex);
 117      size_t result = taskQueue.size();
 118      if (!taskQueue.empty()) {
 119          first = taskQueue.begin()->first;
 120          last = taskQueue.rbegin()->first;
 121      }
 122      return result;
 123  }
 124  
 125  bool CScheduler::AreThreadsServicingQueue() const
 126  {
 127      LOCK(newTaskMutex);
 128      return nThreadsServicingQueue;
 129  }
 130  
 131  
 132  void SerialTaskRunner::MaybeScheduleProcessQueue()
 133  {
 134      {
 135          LOCK(m_callbacks_mutex);
 136          // Try to avoid scheduling too many copies here, but if we
 137          // accidentally have two ProcessQueue's scheduled at once its
 138          // not a big deal.
 139          if (m_are_callbacks_running) return;
 140          if (m_callbacks_pending.empty()) return;
 141      }
 142      m_scheduler.schedule([this] { this->ProcessQueue(); }, std::chrono::steady_clock::now());
 143  }
 144  
 145  void SerialTaskRunner::ProcessQueue()
 146  {
 147      std::function<void()> callback;
 148      {
 149          LOCK(m_callbacks_mutex);
 150          if (m_are_callbacks_running) return;
 151          if (m_callbacks_pending.empty()) return;
 152          m_are_callbacks_running = true;
 153  
 154          callback = std::move(m_callbacks_pending.front());
 155          m_callbacks_pending.pop_front();
 156      }
 157  
 158      // RAII the setting of fCallbacksRunning and calling MaybeScheduleProcessQueue
 159      // to ensure both happen safely even if callback() throws.
 160      struct RAIICallbacksRunning {
 161          SerialTaskRunner* instance;
 162          explicit RAIICallbacksRunning(SerialTaskRunner* _instance) : instance(_instance) {}
 163          ~RAIICallbacksRunning()
 164          {
 165              {
 166                  LOCK(instance->m_callbacks_mutex);
 167                  instance->m_are_callbacks_running = false;
 168              }
 169              instance->MaybeScheduleProcessQueue();
 170          }
 171      } raiicallbacksrunning(this);
 172  
 173      callback();
 174  }
 175  
 176  void SerialTaskRunner::insert(std::function<void()> func)
 177  {
 178      {
 179          LOCK(m_callbacks_mutex);
 180          m_callbacks_pending.emplace_back(std::move(func));
 181      }
 182      MaybeScheduleProcessQueue();
 183  }
 184  
 185  void SerialTaskRunner::flush()
 186  {
 187      assert(!m_scheduler.AreThreadsServicingQueue());
 188      bool should_continue = true;
 189      while (should_continue) {
 190          ProcessQueue();
 191          LOCK(m_callbacks_mutex);
 192          should_continue = !m_callbacks_pending.empty();
 193      }
 194  }
 195  
 196  size_t SerialTaskRunner::size()
 197  {
 198      LOCK(m_callbacks_mutex);
 199      return m_callbacks_pending.size();
 200  }
 201