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: 270 300 90.0 %
Date: 2026-09-07 02:14:47 Functions: 1385 2301 60.2 %
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     1846829 :     QueueTaskRunner(QueueT *queue)
      34     1846829 :         : Task(queue->GetTaskId(), queue->GetTaskInstance()), queue_(queue) {
      35     1846798 :     }
      36             : 
      37     2464377 :     bool Run() {
      38             :         // Check if this run needs to be deferred
      39     2464377 :         if (!queue_->OnEntry()) {
      40       20402 :             return false;
      41             :         }
      42     2443965 :         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     2443949 :     bool RunQueue() {
      54             :         // Check if we need to abort
      55     2443949 :         if (queue_->RunnerAbort()) {
      56          19 :             return queue_->RunnerDone();
      57             :         }
      58             : 
      59     2444101 :         uint64_t start = 0;
      60     2444101 :         if (queue_->measure_busy_time_)
      61           0 :             start = ClockMonotonicUsec();
      62             : 
      63     2444097 :         QueueEntryT entry = QueueEntryT();
      64     2444071 :         size_t count = 0;
      65     9673340 :         while (queue_->Dequeue(&entry)) {
      66             :             // Process the entry
      67     7827428 :             if (!queue_->GetCallback()(entry)) {
      68      551724 :                 break;
      69             :             }
      70     7275915 :             if (++count == queue_->max_iterations_) {
      71       46646 :                 if (start)
      72           0 :                     queue_->add_busy_time(ClockMonotonicUsec() - start);
      73       46646 :                 return queue_->RunnerDone();
      74             :             }
      75             :         }
      76             : 
      77     2397356 :         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     2397356 :         return queue_->RunnerDone();
      84     1361120 :     }
      85             : 
      86             :     QueueT *queue_;
      87             : };
      88             : 
      89             : template <typename QueueEntryT>
      90             : struct WorkQueueDelete {
      91             :     template <typename QueueT>
      92       56394 :     void operator()(QueueT &, bool) {}
      93             : };
      94             : 
      95             : template <typename QueueEntryT>
      96             : struct WorkQueueDelete<QueueEntryT *> {
      97             :     template <typename QueueT>
      98      162940 :     void operator()(QueueT &q, bool delete_entry) {
      99             :         QueueEntryT *entry;
     100      162940 :         while (q.try_pop(entry)) {
     101           0 :             if (delete_entry) {
     102           0 :                 delete entry;
     103             :             }
     104             :         }
     105      162940 :     }
     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      240344 :     WorkQueue(int taskId, int taskInstance, Callback callback,
     120             :               size_t size = kMaxSize,
     121             :               size_t max_iterations = kMaxIterations) :
     122      240347 :         running_(false),
     123      240347 :         taskId_(taskId),
     124      240347 :         taskInstance_(taskInstance),
     125      240347 :         name_(""),
     126      240347 :         callback_(callback),
     127      240340 :         on_entry_cb_(0),
     128      240345 :         on_exit_cb_(0),
     129      240346 :         start_runner_(0),
     130      240344 :         current_runner_(NULL),
     131      240344 :         on_entry_defer_count_(0),
     132      240344 :         deleted_(false),
     133      240344 :         enqueues_(0),
     134      240344 :         dequeues_(0),
     135      240344 :         drops_(0),
     136      240344 :         max_iterations_(max_iterations),
     137      240344 :         size_(size),
     138      240344 :         bounded_(false),
     139      240344 :         shutdown_scheduled_(false),
     140      240344 :         delete_entries_on_shutdown_(true),
     141      240341 :         task_starts_(0),
     142      240341 :         max_queue_len_(0),
     143      240341 :         busy_time_(0),
     144      480688 :         measure_busy_time_(false) {
     145      240341 :         count_ = 0;
     146      240347 :         disabled_ = false;
     147      240346 :     }
     148             : 
     149             :     // Concurrency - should be called from a task whose policy
     150             :     // assures that the dequeue task - QueueTaskRunner is not running
     151             :     // concurrently
     152      219556 :     void Shutdown(bool delete_entries = true) {
     153      219556 :         std::scoped_lock lock(mutex_);
     154      219556 :         ShutdownLocked(delete_entries);
     155      219556 :     }
     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           4 :     void ScheduleShutdown(bool delete_entries = true) {
     161           4 :         std::scoped_lock lock(mutex_);
     162           4 :         if (shutdown_scheduled_) {
     163           0 :             return;
     164             :         }
     165           4 :         shutdown_scheduled_ = true;
     166           4 :         delete_entries_on_shutdown_ = delete_entries;
     167             : 
     168             :         // Cancel QueueTaskRunner
     169           4 :         if (running_) {
     170           3 :             assert(current_runner_);
     171           3 :             TaskScheduler *scheduler = TaskScheduler::GetInstance();
     172             :             TaskScheduler::CancelReturnCode cancel_code =
     173           3 :                 scheduler->Cancel(current_runner_);
     174           3 :             if (cancel_code == TaskScheduler::CANCELLED) {
     175           1 :                 running_ = false;
     176           1 :                 current_runner_ = NULL;
     177           1 :                 ShutdownLocked(delete_entries);
     178             :             } else {
     179           2 :                 assert(cancel_code == TaskScheduler::QUEUED);
     180             :             }
     181             :         } else {
     182           1 :             ShutdownLocked(delete_entries);
     183             :         }
     184           4 :     }
     185             : 
     186      240325 :     ~WorkQueue() {
     187      240325 :         std::scoped_lock lock(mutex_);
     188             :         // Shutdown() needs to be called before deleting
     189             :         //assert(!running_ && deleted_);
     190      240325 :     }
     191             : 
     192         236 :     void SetStartRunnerFunc(StartRunnerFunc start_runner_fn) {
     193         236 :         start_runner_ = start_runner_fn;
     194         236 :     }
     195             : 
     196          11 :     void SetSize(size_t size) {
     197          11 :         size_ = size;
     198          11 :     }
     199             : 
     200          14 :     void SetBounded(bool bounded) {
     201          14 :         bounded_ = bounded;
     202          14 :     }
     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         427 :     void SetHighWaterMark(const WaterMarkInfo& hwm_info) {
     214         427 :         std::scoped_lock lock(water_mutex_);
     215         427 :         watermarks_.SetHighWaterMark(hwm_info);
     216         427 :     }
     217             : 
     218      219561 :     void ResetHighWaterMark() {
     219      219561 :         std::scoped_lock lock(water_mutex_);
     220      219561 :         watermarks_.ResetHighWaterMark();
     221      219561 :     }
     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         412 :     void SetLowWaterMark(const WaterMarkInfo& lwm_info) {
     234         412 :         std::scoped_lock lock(water_mutex_);
     235         412 :         watermarks_.SetLowWaterMark(lwm_info);
     236         412 :      }
     237             : 
     238      219560 :     void ResetLowWaterMark() {
     239      219560 :         std::scoped_lock lock(water_mutex_);
     240      219560 :         watermarks_.ResetLowWaterMark();
     241      219560 :     }
     242             : 
     243           1 :     WaterMarkInfos GetLowWaterMark() const {
     244           1 :         std::scoped_lock lock(water_mutex_);
     245           2 :         return watermarks_.GetLowWaterMark();
     246           1 :     }
     247             : 
     248     7821963 :     bool Enqueue(QueueEntryT entry) {
     249     7821963 :         if (bounded_) {
     250         562 :             if (AreWaterMarksSet()) {
     251           0 :                 return EnqueueBoundedLocked(entry);
     252             :             } else {
     253         562 :                 return EnqueueBounded(entry);
     254             :             }
     255             :         } else {
     256     7821401 :             if (AreWaterMarksSet()) {
     257      112892 :                 return EnqueueInternalLocked(entry);
     258             :             } else {
     259     7708420 :                 return EnqueueInternal(entry);
     260             :             }
     261             :         }
     262             :     }
     263             : 
     264             :     // Returns true if pop is successful.
     265     9673150 :     bool Dequeue(QueueEntryT *entry) {
     266     9673150 :         if (AreWaterMarksSet()) {
     267      119150 :             return DequeueInternalLocked(entry);
     268             :         } else {
     269     9553796 :             return DequeueInternal(entry);
     270             :         }
     271             :     }
     272             : 
     273     1846812 :     int GetTaskId() const {
     274     1846812 :         return taskId_;
     275             :     }
     276             : 
     277     1846823 :     int GetTaskInstance() const {
     278     1846823 :         return taskInstance_;
     279             :     }
     280             : 
     281     7845389 :     void MayBeStartRunner() {
     282     7845389 :         std::scoped_lock lock(mutex_);
     283     7845481 :         if (running_ || queue_.empty() || deleted_ || RunnerAbortLocked()) {
     284     5998484 :             return;
     285             :         }
     286     1846793 :         task_starts_++;
     287     1846793 :         running_ = true;
     288     1846793 :         assert(current_runner_ == NULL);
     289     1846798 :         current_runner_ =
     290     1846793 :             new QueueTaskRunner<QueueEntryT, WorkQueue<QueueEntryT> >(this);
     291     1846798 :         TaskScheduler *scheduler = TaskScheduler::GetInstance();
     292     1846782 :         scheduler->Enqueue(current_runner_);
     293     7845586 :     }
     294             : 
     295     7827393 :     Callback GetCallback() const {
     296     7827393 :         return callback_;
     297             :     }
     298             : 
     299        2087 :     void SetEntryCallback(TaskEntryCallback on_entry) {
     300        2087 :         on_entry_cb_ = on_entry;
     301        2087 :     }
     302             : 
     303       36064 :     void SetExitCallback(TaskExitCallback on_exit) {
     304       36064 :         on_exit_cb_ = on_exit;
     305       36064 :     }
     306             : 
     307       17480 :     void set_name(const std::string &name) {
     308       17480 :         name_ = name;
     309       17480 :     }
     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       34562 :     void set_disable(bool disabled) {
     320       34562 :         if (disabled_ != disabled) {
     321       34558 :             disabled_ = disabled;
     322       34558 :             if (!disabled_) {
     323       17279 :                 MayBeStartRunner();
     324             :             }
     325             :         }
     326       34562 :     }
     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     2464375 :     bool OnEntry() {
     337     2464375 :         bool run = (on_entry_cb_.empty() || on_entry_cb_());
     338             : 
     339             :         // Track number of times this queue run is deferred
     340     2464369 :         if (!run) {
     341       20402 :             on_entry_defer_count_++;
     342             :         }
     343     2464369 :         return run;
     344             :     }
     345             : 
     346     2444061 :     void OnExit(bool done) {
     347     2444061 :         if (!on_exit_cb_.empty()) {
     348       10715 :             on_exit_cb_(done);
     349             :         }
     350     2444064 :     }
     351             : 
     352     8451358 :     bool IsQueueEmpty() const {
     353     8451358 :         return queue_.empty();
     354             :     }
     355             : 
     356        7957 :     size_t Length() const {
     357        7957 :         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           3 :     bool deleted() const {
     373           3 :         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          14 :     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     9672888 :     bool DequeueInternal(QueueEntryT *entry) {
     392     9672888 :         bool success = queue_.try_pop(*entry);
     393     9673156 :         if (success) {
     394     7827498 :             dequeues_++;
     395     7827498 :             size_t ncount(AtomicDecrementQueueCount(entry));
     396     7827726 :             ProcessLowWaterMarks(ncount);
     397             :         }
     398     9673161 :         return success;
     399             :     }
     400             : 
     401      119150 :     bool DequeueInternalLocked(QueueEntryT *entry) {
     402      119150 :         std::scoped_lock lock(water_mutex_);
     403      238300 :         return DequeueInternal(entry);
     404      119150 :     }
     405             : 
     406    17500969 :     bool AreWaterMarksSet() const {
     407    17500969 :         return watermarks_.AreWaterMarksSet();
     408             :     }
     409             : 
     410      219560 :     void ShutdownLocked(bool delete_entries) {
     411             :         // Cancel QueueTaskRunner from the scheduler
     412      219560 :         assert(!deleted_);
     413      219560 :         if (running_) {
     414         434 :             running_ = false;
     415         434 :             assert(current_runner_);
     416         434 :             TaskScheduler *scheduler = TaskScheduler::GetInstance();
     417             :             TaskScheduler::CancelReturnCode cancel_code =
     418         434 :                 scheduler->Cancel(current_runner_);
     419         434 :             assert(cancel_code == TaskScheduler::CANCELLED);
     420         434 :             current_runner_ = NULL;
     421             :         }
     422      219560 :         ResetHighWaterMark();
     423      219560 :         ResetLowWaterMark();
     424             :         WorkQueueDelete<QueueEntryT> deleter;
     425      219560 :         deleter(queue_, delete_entries);
     426      219560 :         queue_.clear();
     427      219560 :         count_ = 0;
     428      219560 :         deleted_ = true;
     429      219560 :     }
     430             : 
     431     7796579 :     size_t AtomicIncrementQueueCount(QueueEntryT *entry) {
     432    15593158 :         return count_.fetch_add(1) + 1;
     433             :     }
     434             : 
     435     7802107 :     size_t AtomicDecrementQueueCount(QueueEntryT *entry) {
     436    15604214 :         return count_.fetch_sub(1) - 1;
     437             :     }
     438             : 
     439     7822171 :     void ProcessHighWaterMarks(size_t count) {
     440     7822171 :         watermarks_.ProcessHighWaterMarks(count);
     441     7822024 :     }
     442             : 
     443     7827665 :     void ProcessLowWaterMarks(size_t count) {
     444     7827665 :         watermarks_.ProcessLowWaterMarks(count);
     445     7827460 :     }
     446             : 
     447     7821392 :     bool EnqueueInternal(QueueEntryT entry) {
     448     7821392 :         enqueues_++;
     449     7821392 :         size_t ncount(AtomicIncrementQueueCount(&entry));
     450     7821618 :         if (ncount > max_queue_len_)
     451      943121 :             max_queue_len_ = ncount;
     452     7821618 :         ProcessHighWaterMarks(ncount);
     453     7821453 :         queue_.push(entry);
     454     7821463 :         MayBeStartRunner();
     455     7821629 :         return ncount < size_;
     456             :     }
     457             : 
     458      112892 :     bool EnqueueInternalLocked(QueueEntryT entry) {
     459      112892 :         std::scoped_lock lock(water_mutex_);
     460      225784 :         return EnqueueInternal(entry);
     461      112892 :     }
     462             : 
     463         562 :     bool EnqueueBounded(QueueEntryT entry) {
     464         562 :         size_t ncount(AtomicIncrementQueueCount(&entry));
     465         562 :         if (ncount > max_queue_len_)
     466          30 :             max_queue_len_ = ncount;
     467         562 :         if (ncount < size_) {
     468         557 :             enqueues_++;
     469         557 :             ProcessHighWaterMarks(ncount);
     470         557 :             queue_.push(entry);
     471         557 :             MayBeStartRunner();
     472         557 :             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     4921303 :     bool RunnerAbortLocked() {
     486     9809670 :         return (disabled_ || shutdown_scheduled_ ||
     487     9809558 :                 (!start_runner_.empty() && !start_runner_()));
     488             :     }
     489             : 
     490     2443941 :     bool RunnerAbort() {
     491     2443941 :         std::scoped_lock lock(mutex_);
     492     4888283 :         return RunnerAbortLocked();
     493     2444037 :     }
     494             : 
     495     2444150 :     bool RunnerDone() {
     496     2444150 :         std::scoped_lock lock(mutex_);
     497     2444152 :         bool done = false;
     498     2444152 :         if (queue_.empty() || RunnerAbortLocked()) {
     499     1846515 :             done = true;
     500     1846515 :             OnExit(done);
     501     1846517 :             current_runner_ = NULL;
     502     1846517 :             running_ = false;
     503     1846517 :             if (shutdown_scheduled_) {
     504           2 :                 ShutdownLocked(delete_entries_on_shutdown_);
     505             :             }
     506             :         } else {
     507      597543 :             OnExit(done);
     508      597547 :             running_ = true;
     509             :         }
     510     2444159 :         return done;
     511     2444064 :     }
     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