LCOV - code coverage report
Current view: top level - root/contrail/src/contrail-common/base - queue_task.h (source / functions) Hit Total Coverage
Test: OpenSDN C/C++ coverage (all TARGET_SET jobs) Lines: 262 300 87.3 %
Date: 2026-10-05 02:12:29 Functions: 688 2301 29.9 %
Legend: Lines: hit not hit

          Line data    Source code
       1             : /*
       2             :  * Copyright (c) 2013 Juniper Networks, Inc. All rights reserved.
       3             :  */
       4             : 
       5             : // queue_task.h
       6             : //
       7             : // Task based queue processor implementing thread safe enqueue and dequeue
       8             : // using concurrent queues. If queue is empty, enqueue creates a dequeue task
       9             : // that drains the queue. The dequeue task runs a maximum of kMaxIterations
      10             : // before yielding.
      11             : //
      12             : #ifndef __QUEUE_TASK_H__
      13             : #define __QUEUE_TASK_H__
      14             : 
      15             : #include <iostream>
      16             : #include <sstream>
      17             : #include <algorithm>
      18             : #include <vector>
      19             : #include <set>
      20             : #include <mutex>
      21             : #include <atomic>
      22             : 
      23             : #include <tbb/concurrent_queue.h>
      24             : #include <tbb/spin_rw_mutex.h>
      25             : 
      26             : #include <base/task.h>
      27             : #include <base/time_util.h>
      28             : #include <base/watermark.h>
      29             : 
      30             : template <typename QueueEntryT, typename QueueT>
      31             : class QueueTaskRunner : public Task {
      32             : public:
      33      395179 :     QueueTaskRunner(QueueT *queue)
      34      395179 :         : Task(queue->GetTaskId(), queue->GetTaskInstance()), queue_(queue) {
      35      395178 :     }
      36             : 
      37     1100550 :     bool Run() {
      38             :         // Check if this run needs to be deferred
      39     1100550 :         if (!queue_->OnEntry()) {
      40         741 :             return false;
      41             :         }
      42     1099840 :         return RunQueue();
      43             :         // No more client callbacks after updating
      44             :         // queue running_ and current_runner_ in RunQueue to
      45             :         // avoid client callbacks running concurrently
      46             :     }
      47             : 
      48           0 :     virtual std::string Description() const {
      49           0 :         return queue_->Description();
      50             :     }
      51             : 
      52             : private:
      53     1099809 :     bool RunQueue() {
      54             :         // Check if we need to abort
      55     1099809 :         if (queue_->RunnerAbort()) {
      56          19 :             return queue_->RunnerDone();
      57             :         }
      58             : 
      59     1099845 :         uint64_t start = 0;
      60     1099845 :         if (queue_->measure_busy_time_)
      61           0 :             start = ClockMonotonicUsec();
      62             : 
      63     1099834 :         QueueEntryT entry = QueueEntryT();
      64     1099828 :         size_t count = 0;
      65     3874281 :         while (queue_->Dequeue(&entry)) {
      66             :             // Process the entry
      67     3479443 :             if (!queue_->GetCallback()(entry)) {
      68      668364 :                 break;
      69             :             }
      70     2811123 :             if (++count == queue_->max_iterations_) {
      71       36670 :                 if (start)
      72           0 :                     queue_->add_busy_time(ClockMonotonicUsec() - start);
      73       36670 :                 return queue_->RunnerDone();
      74             :             }
      75             :         }
      76             : 
      77     1063118 :         if (start)
      78           0 :             queue_->add_busy_time(ClockMonotonicUsec() - start);
      79             : 
      80             :         // Running is done if queue_ is empty
      81             :         // While notification is being run, its possible that more entries
      82             :         // are added into queue_
      83     1063118 :         return queue_->RunnerDone();
      84      194625 :     }
      85             : 
      86             :     QueueT *queue_;
      87             : };
      88             : 
      89             : template <typename QueueEntryT>
      90             : struct WorkQueueDelete {
      91             :     template <typename QueueT>
      92       22306 :     void operator()(QueueT &, bool) {}
      93             : };
      94             : 
      95             : template <typename QueueEntryT>
      96             : struct WorkQueueDelete<QueueEntryT *> {
      97             :     template <typename QueueT>
      98       43125 :     void operator()(QueueT &q, bool delete_entry) {
      99             :         QueueEntryT *entry;
     100       43125 :         while (q.try_pop(entry)) {
     101           0 :             if (delete_entry) {
     102           0 :                 delete entry;
     103             :             }
     104             :         }
     105       43125 :     }
     106             : };
     107             : 
     108             : template <typename QueueEntryT>
     109             : class WorkQueue {
     110             : public:
     111             :     static const int kMaxSize = 1024;
     112             :     static const int kMaxIterations = 32;
     113             :     typedef tbb::concurrent_queue<QueueEntryT> Queue;
     114             :     typedef boost::function<bool (QueueEntryT)> Callback;
     115             :     typedef boost::function<bool (void)> StartRunnerFunc;
     116             :     typedef boost::function<void (bool)> TaskExitCallback;
     117             :     typedef boost::function<bool ()> TaskEntryCallback;
     118             : 
     119       64128 :     WorkQueue(int taskId, int taskInstance, Callback callback,
     120             :               size_t size = kMaxSize,
     121             :               size_t max_iterations = kMaxIterations) :
     122       64128 :         running_(false),
     123       64128 :         taskId_(taskId),
     124       64128 :         taskInstance_(taskInstance),
     125       64128 :         name_(""),
     126       64128 :         callback_(callback),
     127       64128 :         on_entry_cb_(0),
     128       64128 :         on_exit_cb_(0),
     129       64128 :         start_runner_(0),
     130       64128 :         current_runner_(NULL),
     131       64128 :         on_entry_defer_count_(0),
     132       64128 :         deleted_(false),
     133       64128 :         enqueues_(0),
     134       64128 :         dequeues_(0),
     135       64128 :         drops_(0),
     136       64128 :         max_iterations_(max_iterations),
     137       64128 :         size_(size),
     138       64128 :         bounded_(false),
     139       64128 :         shutdown_scheduled_(false),
     140       64128 :         delete_entries_on_shutdown_(true),
     141       64128 :         task_starts_(0),
     142       64128 :         max_queue_len_(0),
     143       64128 :         busy_time_(0),
     144      128256 :         measure_busy_time_(false) {
     145       64128 :         count_ = 0;
     146       64128 :         disabled_ = false;
     147       64128 :     }
     148             : 
     149             :     // Concurrency - should be called from a task whose policy
     150             :     // assures that the dequeue task - QueueTaskRunner is not running
     151             :     // concurrently
     152       63087 :     void Shutdown(bool delete_entries = true) {
     153       63087 :         std::scoped_lock lock(mutex_);
     154       63087 :         ShutdownLocked(delete_entries);
     155       63087 :     }
     156             : 
     157             :     // Concurrency - can be called from any context
     158             :     // Schedule shutdown of the WorkQueue, shutdown may happen asynchronously
     159             :     // or in the caller's context also
     160           3 :     void ScheduleShutdown(bool delete_entries = true) {
     161           3 :         std::scoped_lock lock(mutex_);
     162           3 :         if (shutdown_scheduled_) {
     163           0 :             return;
     164             :         }
     165           3 :         shutdown_scheduled_ = true;
     166           3 :         delete_entries_on_shutdown_ = delete_entries;
     167             : 
     168             :         // Cancel QueueTaskRunner
     169           3 :         if (running_) {
     170           2 :             assert(current_runner_);
     171           2 :             TaskScheduler *scheduler = TaskScheduler::GetInstance();
     172             :             TaskScheduler::CancelReturnCode cancel_code =
     173           2 :                 scheduler->Cancel(current_runner_);
     174           2 :             if (cancel_code == TaskScheduler::CANCELLED) {
     175           1 :                 running_ = false;
     176           1 :                 current_runner_ = NULL;
     177           1 :                 ShutdownLocked(delete_entries);
     178             :             } else {
     179           1 :                 assert(cancel_code == TaskScheduler::QUEUED);
     180             :             }
     181             :         } else {
     182           1 :             ShutdownLocked(delete_entries);
     183             :         }
     184           3 :     }
     185             : 
     186       64112 :     ~WorkQueue() {
     187       64112 :         std::scoped_lock lock(mutex_);
     188             :         // Shutdown() needs to be called before deleting
     189             :         //assert(!running_ && deleted_);
     190       64112 :     }
     191             : 
     192         158 :     void SetStartRunnerFunc(StartRunnerFunc start_runner_fn) {
     193         158 :         start_runner_ = start_runner_fn;
     194         158 :     }
     195             : 
     196           1 :     void SetSize(size_t size) {
     197           1 :         size_ = size;
     198           1 :     }
     199             : 
     200           1 :     void SetBounded(bool bounded) {
     201           1 :         bounded_ = bounded;
     202           1 :     }
     203             : 
     204             :     bool GetBounded() const {
     205             :         return bounded_;
     206             :     }
     207             : 
     208           1 :     void SetHighWaterMark(const WaterMarkInfos &high_water) {
     209           1 :         std::scoped_lock lock(water_mutex_);
     210           1 :         watermarks_.SetHighWaterMark(high_water);
     211           1 :     }
     212             : 
     213         291 :     void SetHighWaterMark(const WaterMarkInfo& hwm_info) {
     214         291 :         std::scoped_lock lock(water_mutex_);
     215         291 :         watermarks_.SetHighWaterMark(hwm_info);
     216         291 :     }
     217             : 
     218       65547 :     void ResetHighWaterMark() {
     219       65547 :         std::scoped_lock lock(water_mutex_);
     220       65547 :         watermarks_.ResetHighWaterMark();
     221       65547 :     }
     222             : 
     223           1 :     WaterMarkInfos GetHighWaterMark() const {
     224           1 :         std::scoped_lock lock(water_mutex_);
     225           2 :         return watermarks_.GetHighWaterMark();
     226           1 :     }
     227             : 
     228           4 :     void SetLowWaterMark(const WaterMarkInfos &low_water) {
     229           4 :         std::scoped_lock lock(water_mutex_);
     230           4 :         watermarks_.SetLowWaterMark(low_water);
     231           4 :      }
     232             : 
     233         276 :     void SetLowWaterMark(const WaterMarkInfo& lwm_info) {
     234         276 :         std::scoped_lock lock(water_mutex_);
     235         276 :         watermarks_.SetLowWaterMark(lwm_info);
     236         276 :      }
     237             : 
     238       65546 :     void ResetLowWaterMark() {
     239       65546 :         std::scoped_lock lock(water_mutex_);
     240       65546 :         watermarks_.ResetLowWaterMark();
     241       65546 :     }
     242             : 
     243           1 :     WaterMarkInfos GetLowWaterMark() const {
     244           1 :         std::scoped_lock lock(water_mutex_);
     245           2 :         return watermarks_.GetLowWaterMark();
     246           1 :     }
     247             : 
     248     3444740 :     bool Enqueue(QueueEntryT entry) {
     249     3444740 :         if (bounded_) {
     250           5 :             if (AreWaterMarksSet()) {
     251           0 :                 return EnqueueBoundedLocked(entry);
     252             :             } else {
     253           5 :                 return EnqueueBounded(entry);
     254             :             }
     255             :         } else {
     256     3444735 :             if (AreWaterMarksSet()) {
     257      100583 :                 return EnqueueInternalLocked(entry);
     258             :             } else {
     259     3344125 :                 return EnqueueInternal(entry);
     260             :             }
     261             :         }
     262             :     }
     263             : 
     264             :     // Returns true if pop is successful.
     265     3874229 :     bool Dequeue(QueueEntryT *entry) {
     266     3874229 :         if (AreWaterMarksSet()) {
     267      100927 :             return DequeueInternalLocked(entry);
     268             :         } else {
     269     3773273 :             return DequeueInternal(entry);
     270             :         }
     271             :     }
     272             : 
     273      395176 :     int GetTaskId() const {
     274      395176 :         return taskId_;
     275             :     }
     276             : 
     277      395178 :     int GetTaskInstance() const {
     278      395178 :         return taskInstance_;
     279             :     }
     280             : 
     281     3496989 :     void MayBeStartRunner() {
     282     3496989 :         std::scoped_lock lock(mutex_);
     283     3497016 :         if (running_ || queue_.empty() || deleted_ || RunnerAbortLocked()) {
     284     3101766 :             return;
     285             :         }
     286      395181 :         task_starts_++;
     287      395181 :         running_ = true;
     288      395181 :         assert(current_runner_ == NULL);
     289      395175 :         current_runner_ =
     290      395181 :             new QueueTaskRunner<QueueEntryT, WorkQueue<QueueEntryT> >(this);
     291      395175 :         TaskScheduler *scheduler = TaskScheduler::GetInstance();
     292      395170 :         scheduler->Enqueue(current_runner_);
     293     3497037 :     }
     294             : 
     295     3479436 :     Callback GetCallback() const {
     296     3479436 :         return callback_;
     297             :     }
     298             : 
     299          40 :     void SetEntryCallback(TaskEntryCallback on_entry) {
     300          40 :         on_entry_cb_ = on_entry;
     301          40 :     }
     302             : 
     303        1877 :     void SetExitCallback(TaskExitCallback on_exit) {
     304        1877 :         on_exit_cb_ = on_exit;
     305        1877 :     }
     306             : 
     307       17935 :     void set_name(const std::string &name) {
     308       17935 :         name_ = name;
     309       17935 :     }
     310           0 :     std::string Description() const {
     311           0 :         if (name_.empty() == false)
     312           0 :             return name_;
     313             : 
     314           0 :         std::ostringstream str;
     315           0 :         str << "Function " << callback_;
     316           0 :         return str.str();
     317           0 :     }
     318             : 
     319       34226 :     void set_disable(bool disabled) {
     320       34226 :         if (disabled_ != disabled) {
     321       34224 :             disabled_ = disabled;
     322       34224 :             if (!disabled_) {
     323       17112 :                 MayBeStartRunner();
     324             :             }
     325             :         }
     326       34226 :     }
     327             : 
     328             :     bool IsDisabled() const {
     329             :         return disabled_;
     330             :     }
     331             : 
     332             :     size_t on_entry_defer_count() const {
     333             :         return on_entry_defer_count_;
     334             :     }
     335             : 
     336     1100579 :     bool OnEntry() {
     337     1100579 :         bool run = (on_entry_cb_.empty() || on_entry_cb_());
     338             : 
     339             :         // Track number of times this queue run is deferred
     340     1100583 :         if (!run) {
     341         741 :             on_entry_defer_count_++;
     342             :         }
     343     1100583 :         return run;
     344             :     }
     345             : 
     346     1099850 :     void OnExit(bool done) {
     347     1099850 :         if (!on_exit_cb_.empty()) {
     348         866 :             on_exit_cb_(done);
     349             :         }
     350     1099848 :     }
     351             : 
     352      121372 :     bool IsQueueEmpty() const {
     353      121372 :         return queue_.empty();
     354             :     }
     355             : 
     356         414 :     size_t Length() const {
     357         414 :         return count_;
     358             :     }
     359             : 
     360          20 :     size_t NumEnqueues() const {
     361          20 :         return enqueues_;
     362             :     }
     363             : 
     364          21 :     size_t NumDequeues() const {
     365          21 :         return dequeues_;
     366             :     }
     367             : 
     368             :     size_t NumDrops() const {
     369             :         return drops_;
     370             :     }
     371             : 
     372           0 :     bool deleted() const {
     373           0 :         return deleted_;
     374             :     }
     375             : 
     376           0 :     uint32_t task_starts() const { return task_starts_; }
     377          11 :     size_t max_queue_len() const { return max_queue_len_; }
     378             :     bool measure_busy_time() const { return measure_busy_time_; }
     379           0 :     void set_measure_busy_time(bool val) const { measure_busy_time_ = val; }
     380           0 :     uint64_t busy_time() const { return busy_time_; }
     381           0 :     void add_busy_time(uint64_t t) { busy_time_ += t; }
     382           0 :     void ClearStats() const {
     383           0 :         max_queue_len_ = 0;
     384           0 :         enqueues_ = 0;
     385           0 :         dequeues_ = 0;
     386           0 :         busy_time_ = 0;
     387           0 :         task_starts_ = 0;
     388           0 :     }
     389             : private:
     390             :     // Returns true if pop is successful.
     391     3874168 :     bool DequeueInternal(QueueEntryT *entry) {
     392     3874168 :         bool success = queue_.try_pop(*entry);
     393     3874250 :         if (success) {
     394     3479472 :             dequeues_++;
     395     3479472 :             size_t ncount(AtomicDecrementQueueCount(entry));
     396     3479490 :             ProcessLowWaterMarks(ncount);
     397             :         }
     398     3874243 :         return success;
     399             :     }
     400             : 
     401      100927 :     bool DequeueInternalLocked(QueueEntryT *entry) {
     402      100927 :         std::scoped_lock lock(water_mutex_);
     403      201854 :         return DequeueInternal(entry);
     404      100927 :     }
     405             : 
     406     7353946 :     bool AreWaterMarksSet() const {
     407     7353946 :         return watermarks_.AreWaterMarksSet();
     408             :     }
     409             : 
     410       65546 :     void ShutdownLocked(bool delete_entries) {
     411             :         // Cancel QueueTaskRunner from the scheduler
     412       65546 :         assert(!deleted_);
     413       65546 :         if (running_) {
     414         315 :             running_ = false;
     415         315 :             assert(current_runner_);
     416         315 :             TaskScheduler *scheduler = TaskScheduler::GetInstance();
     417             :             TaskScheduler::CancelReturnCode cancel_code =
     418         315 :                 scheduler->Cancel(current_runner_);
     419         315 :             assert(cancel_code == TaskScheduler::CANCELLED);
     420         315 :             current_runner_ = NULL;
     421             :         }
     422       65546 :         ResetHighWaterMark();
     423       65546 :         ResetLowWaterMark();
     424             :         WorkQueueDelete<QueueEntryT> deleter;
     425       65546 :         deleter(queue_, delete_entries);
     426       65546 :         queue_.clear();
     427       65546 :         count_ = 0;
     428       65546 :         deleted_ = true;
     429       65546 :     }
     430             : 
     431     3444193 :     size_t AtomicIncrementQueueCount(QueueEntryT *entry) {
     432     6888386 :         return count_.fetch_add(1) + 1;
     433             :     }
     434             : 
     435     3478924 :     size_t AtomicDecrementQueueCount(QueueEntryT *entry) {
     436     6957848 :         return count_.fetch_sub(1) - 1;
     437             :     }
     438             : 
     439     3444790 :     void ProcessHighWaterMarks(size_t count) {
     440     3444790 :         watermarks_.ProcessHighWaterMarks(count);
     441     3444756 :     }
     442             : 
     443     3479474 :     void ProcessLowWaterMarks(size_t count) {
     444     3479474 :         watermarks_.ProcessLowWaterMarks(count);
     445     3479448 :     }
     446             : 
     447     3444732 :     bool EnqueueInternal(QueueEntryT entry) {
     448     3444732 :         enqueues_++;
     449     3444732 :         size_t ncount(AtomicIncrementQueueCount(&entry));
     450     3444795 :         if (ncount > max_queue_len_)
     451      643162 :             max_queue_len_ = ncount;
     452     3444795 :         ProcessHighWaterMarks(ncount);
     453     3444755 :         queue_.push(entry);
     454     3444755 :         MayBeStartRunner();
     455     3444795 :         return ncount < size_;
     456             :     }
     457             : 
     458      100583 :     bool EnqueueInternalLocked(QueueEntryT entry) {
     459      100583 :         std::scoped_lock lock(water_mutex_);
     460      201166 :         return EnqueueInternal(entry);
     461      100583 :     }
     462             : 
     463           5 :     bool EnqueueBounded(QueueEntryT entry) {
     464           5 :         size_t ncount(AtomicIncrementQueueCount(&entry));
     465           5 :         if (ncount > max_queue_len_)
     466           5 :             max_queue_len_ = ncount;
     467           5 :         if (ncount < size_) {
     468           0 :             enqueues_++;
     469           0 :             ProcessHighWaterMarks(ncount);
     470           0 :             queue_.push(entry);
     471           0 :             MayBeStartRunner();
     472           0 :             return true;
     473             :         }
     474           5 :         AtomicDecrementQueueCount(&entry);
     475           5 :         drops_++;
     476           5 :         max_queue_len_ = count_;
     477           5 :         return false;
     478             :     }
     479             : 
     480           0 :     bool EnqueueBoundedLocked(QueueEntryT entry) {
     481           0 :         std::scoped_lock lock(water_mutex_);
     482           0 :         return EnqueueBounded(entry);
     483           0 :     }
     484             : 
     485     2224023 :     bool RunnerAbortLocked() {
     486     4424000 :         return (disabled_ || shutdown_scheduled_ ||
     487     4423983 :                 (!start_runner_.empty() && !start_runner_()));
     488             :     }
     489             : 
     490     1099828 :     bool RunnerAbort() {
     491     1099828 :         std::scoped_lock lock(mutex_);
     492     2199727 :         return RunnerAbortLocked();
     493     1099854 :     }
     494             : 
     495     1099868 :     bool RunnerDone() {
     496     1099868 :         std::scoped_lock lock(mutex_);
     497     1099870 :         bool done = false;
     498     1099870 :         if (queue_.empty() || RunnerAbortLocked()) {
     499      394926 :             done = true;
     500      394926 :             OnExit(done);
     501      394925 :             current_runner_ = NULL;
     502      394925 :             running_ = false;
     503      394925 :             if (shutdown_scheduled_) {
     504           1 :                 ShutdownLocked(delete_entries_on_shutdown_);
     505             :             }
     506             :         } else {
     507      704923 :             OnExit(done);
     508      704923 :             running_ = true;
     509             :         }
     510     1099867 :         return done;
     511     1099848 :     }
     512             : 
     513             :     Queue queue_;
     514             :     std::atomic<size_t> count_;
     515             :     std::mutex mutex_;
     516             :     bool running_;
     517             :     int taskId_;
     518             :     int taskInstance_;
     519             :     std::string name_;
     520             :     Callback callback_;
     521             :     TaskEntryCallback on_entry_cb_;
     522             :     TaskExitCallback on_exit_cb_;
     523             :     StartRunnerFunc start_runner_;
     524             :     QueueTaskRunner<QueueEntryT, WorkQueue<QueueEntryT> > *current_runner_;
     525             :     size_t on_entry_defer_count_;
     526             :     std::atomic<bool> disabled_;
     527             :     bool deleted_;
     528             :     mutable size_t enqueues_;
     529             :     mutable size_t dequeues_;
     530             :     size_t drops_;
     531             :     size_t max_iterations_;
     532             :     size_t size_;
     533             :     bool bounded_;
     534             :     bool shutdown_scheduled_;
     535             :     bool delete_entries_on_shutdown_;
     536             :     WaterMarkTuple watermarks_;
     537             :     mutable std::mutex water_mutex_;
     538             :     mutable uint32_t task_starts_;
     539             :     mutable size_t max_queue_len_;
     540             :     mutable uint64_t busy_time_;
     541             :     mutable bool measure_busy_time_;
     542             : 
     543             :     friend class QueueTaskTest;
     544             :     friend class QueueTaskShutdownTest;
     545             :     friend class QueueTaskWaterMarkTest;
     546             :     friend class QueueTaskRunner<QueueEntryT, WorkQueue<QueueEntryT> >;
     547             : 
     548             :     DISALLOW_COPY_AND_ASSIGN(WorkQueue);
     549             : };
     550             : 
     551             : #endif /* __QUEUE_TASK_H__ */

Generated by: LCOV version 1.14