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: 268 300 89.3 %
Date: 2026-09-28 02:13:17 Functions: 1090 2301 47.4 %
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      217169 :     QueueTaskRunner(QueueT *queue)
      34      217169 :         : Task(queue->GetTaskId(), queue->GetTaskInstance()), queue_(queue) {
      35      217133 :     }
      36             : 
      37      744655 :     bool Run() {
      38             :         // Check if this run needs to be deferred
      39      744655 :         if (!queue_->OnEntry()) {
      40         136 :             return false;
      41             :         }
      42      744549 :         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      744533 :     bool RunQueue() {
      54             :         // Check if we need to abort
      55      744533 :         if (queue_->RunnerAbort()) {
      56          19 :             return queue_->RunnerDone();
      57             :         }
      58             : 
      59      744561 :         uint64_t start = 0;
      60      744561 :         if (queue_->measure_busy_time_)
      61           0 :             start = ClockMonotonicUsec();
      62             : 
      63      744548 :         QueueEntryT entry = QueueEntryT();
      64      744517 :         size_t count = 0;
      65     2841816 :         while (queue_->Dequeue(&entry)) {
      66             :             // Process the entry
      67     2625011 :             if (!queue_->GetCallback()(entry)) {
      68      493042 :                 break;
      69             :             }
      70     2132077 :             if (++count == queue_->max_iterations_) {
      71       34778 :                 if (start)
      72           0 :                     queue_->add_busy_time(ClockMonotonicUsec() - start);
      73       34778 :                 return queue_->RunnerDone();
      74             :             }
      75             :         }
      76             : 
      77      709743 :         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      709743 :         return queue_->RunnerDone();
      84       29974 :     }
      85             : 
      86             :     QueueT *queue_;
      87             : };
      88             : 
      89             : template <typename QueueEntryT>
      90             : struct WorkQueueDelete {
      91             :     template <typename QueueT>
      92       21099 :     void operator()(QueueT &, bool) {}
      93             : };
      94             : 
      95             : template <typename QueueEntryT>
      96             : struct WorkQueueDelete<QueueEntryT *> {
      97             :     template <typename QueueT>
      98       41839 :     void operator()(QueueT &q, bool delete_entry) {
      99             :         QueueEntryT *entry;
     100       41839 :         while (q.try_pop(entry)) {
     101           0 :             if (delete_entry) {
     102           0 :                 delete entry;
     103             :             }
     104             :         }
     105       41839 :     }
     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       61219 :     WorkQueue(int taskId, int taskInstance, Callback callback,
     120             :               size_t size = kMaxSize,
     121             :               size_t max_iterations = kMaxIterations) :
     122       61219 :         running_(false),
     123       61219 :         taskId_(taskId),
     124       61219 :         taskInstance_(taskInstance),
     125       61219 :         name_(""),
     126       61219 :         callback_(callback),
     127       61219 :         on_entry_cb_(0),
     128       61219 :         on_exit_cb_(0),
     129       61219 :         start_runner_(0),
     130       61219 :         current_runner_(NULL),
     131       61219 :         on_entry_defer_count_(0),
     132       61219 :         deleted_(false),
     133       61219 :         enqueues_(0),
     134       61219 :         dequeues_(0),
     135       61219 :         drops_(0),
     136       61219 :         max_iterations_(max_iterations),
     137       61219 :         size_(size),
     138       61219 :         bounded_(false),
     139       61219 :         shutdown_scheduled_(false),
     140       61219 :         delete_entries_on_shutdown_(true),
     141       61219 :         task_starts_(0),
     142       61219 :         max_queue_len_(0),
     143       61219 :         busy_time_(0),
     144      122438 :         measure_busy_time_(false) {
     145       61219 :         count_ = 0;
     146       61219 :         disabled_ = false;
     147       61219 :     }
     148             : 
     149             :     // Concurrency - should be called from a task whose policy
     150             :     // assures that the dequeue task - QueueTaskRunner is not running
     151             :     // concurrently
     152       60592 :     void Shutdown(bool delete_entries = true) {
     153       60592 :         std::scoped_lock lock(mutex_);
     154       60592 :         ShutdownLocked(delete_entries);
     155       60592 :     }
     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       61197 :     ~WorkQueue() {
     187       61197 :         std::scoped_lock lock(mutex_);
     188             :         // Shutdown() needs to be called before deleting
     189             :         //assert(!running_ && deleted_);
     190       61197 :     }
     191             : 
     192         162 :     void SetStartRunnerFunc(StartRunnerFunc start_runner_fn) {
     193         162 :         start_runner_ = start_runner_fn;
     194         162 :     }
     195             : 
     196           1 :     void SetSize(size_t size) {
     197           1 :         size_ = size;
     198           1 :     }
     199             : 
     200           4 :     void SetBounded(bool bounded) {
     201           4 :         bounded_ = bounded;
     202           4 :     }
     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         279 :     void SetHighWaterMark(const WaterMarkInfo& hwm_info) {
     214         279 :         std::scoped_lock lock(water_mutex_);
     215         279 :         watermarks_.SetHighWaterMark(hwm_info);
     216         279 :     }
     217             : 
     218       63052 :     void ResetHighWaterMark() {
     219       63052 :         std::scoped_lock lock(water_mutex_);
     220       63052 :         watermarks_.ResetHighWaterMark();
     221       63052 :     }
     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         264 :     void SetLowWaterMark(const WaterMarkInfo& lwm_info) {
     234         264 :         std::scoped_lock lock(water_mutex_);
     235         264 :         watermarks_.SetLowWaterMark(lwm_info);
     236         264 :      }
     237             : 
     238       63051 :     void ResetLowWaterMark() {
     239       63051 :         std::scoped_lock lock(water_mutex_);
     240       63051 :         watermarks_.ResetLowWaterMark();
     241       63051 :     }
     242             : 
     243           1 :     WaterMarkInfos GetLowWaterMark() const {
     244           1 :         std::scoped_lock lock(water_mutex_);
     245           2 :         return watermarks_.GetLowWaterMark();
     246           1 :     }
     247             : 
     248     2590736 :     bool Enqueue(QueueEntryT entry) {
     249     2590736 :         if (bounded_) {
     250          77 :             if (AreWaterMarksSet()) {
     251           0 :                 return EnqueueBoundedLocked(entry);
     252             :             } else {
     253          77 :                 return EnqueueBounded(entry);
     254             :             }
     255             :         } else {
     256     2590659 :             if (AreWaterMarksSet()) {
     257      100572 :                 return EnqueueInternalLocked(entry);
     258             :             } else {
     259     2490008 :                 return EnqueueInternal(entry);
     260             :             }
     261             :         }
     262             :     }
     263             : 
     264             :     // Returns true if pop is successful.
     265     2841753 :     bool Dequeue(QueueEntryT *entry) {
     266     2841753 :         if (AreWaterMarksSet()) {
     267      100855 :             return DequeueInternalLocked(entry);
     268             :         } else {
     269     2740828 :             return DequeueInternal(entry);
     270             :         }
     271             :     }
     272             : 
     273      217158 :     int GetTaskId() const {
     274      217158 :         return taskId_;
     275             :     }
     276             : 
     277      217162 :     int GetTaskInstance() const {
     278      217162 :         return taskInstance_;
     279             :     }
     280             : 
     281     2642527 :     void MayBeStartRunner() {
     282     2642527 :         std::scoped_lock lock(mutex_);
     283     2642589 :         if (running_ || queue_.empty() || deleted_ || RunnerAbortLocked()) {
     284     2425316 :             return;
     285             :         }
     286      217126 :         task_starts_++;
     287      217126 :         running_ = true;
     288      217126 :         assert(current_runner_ == NULL);
     289      217128 :         current_runner_ =
     290      217126 :             new QueueTaskRunner<QueueEntryT, WorkQueue<QueueEntryT> >(this);
     291      217128 :         TaskScheduler *scheduler = TaskScheduler::GetInstance();
     292      217123 :         scheduler->Enqueue(current_runner_);
     293     2642673 :     }
     294             : 
     295     2624998 :     Callback GetCallback() const {
     296     2624998 :         return callback_;
     297             :     }
     298             : 
     299           3 :     void SetEntryCallback(TaskEntryCallback on_entry) {
     300           3 :         on_entry_cb_ = on_entry;
     301           3 :     }
     302             : 
     303        1538 :     void SetExitCallback(TaskExitCallback on_exit) {
     304        1538 :         on_exit_cb_ = on_exit;
     305        1538 :     }
     306             : 
     307       17469 :     void set_name(const std::string &name) {
     308       17469 :         name_ = name;
     309       17469 :     }
     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       34227 :     void set_disable(bool disabled) {
     320       34227 :         if (disabled_ != disabled) {
     321       34226 :             disabled_ = disabled;
     322       34226 :             if (!disabled_) {
     323       17113 :                 MayBeStartRunner();
     324             :             }
     325             :         }
     326       34227 :     }
     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      744664 :     bool OnEntry() {
     337      744664 :         bool run = (on_entry_cb_.empty() || on_entry_cb_());
     338             : 
     339             :         // Track number of times this queue run is deferred
     340      744687 :         if (!run) {
     341         136 :             on_entry_defer_count_++;
     342             :         }
     343      744687 :         return run;
     344             :     }
     345             : 
     346      744539 :     void OnExit(bool done) {
     347      744539 :         if (!on_exit_cb_.empty()) {
     348        1151 :             on_exit_cb_(done);
     349             :         }
     350      744537 :     }
     351             : 
     352       87579 :     bool IsQueueEmpty() const {
     353       87579 :         return queue_.empty();
     354             :     }
     355             : 
     356         216 :     size_t Length() const {
     357         216 :         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          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     2841645 :     bool DequeueInternal(QueueEntryT *entry) {
     392     2841645 :         bool success = queue_.try_pop(*entry);
     393     2841806 :         if (success) {
     394     2625082 :             dequeues_++;
     395     2625082 :             size_t ncount(AtomicDecrementQueueCount(entry));
     396     2625177 :             ProcessLowWaterMarks(ncount);
     397             :         }
     398     2841761 :         return success;
     399             :     }
     400             : 
     401      100855 :     bool DequeueInternalLocked(QueueEntryT *entry) {
     402      100855 :         std::scoped_lock lock(water_mutex_);
     403      201710 :         return DequeueInternal(entry);
     404      100855 :     }
     405             : 
     406     5466945 :     bool AreWaterMarksSet() const {
     407     5466945 :         return watermarks_.AreWaterMarksSet();
     408             :     }
     409             : 
     410       63051 :     void ShutdownLocked(bool delete_entries) {
     411             :         // Cancel QueueTaskRunner from the scheduler
     412       63051 :         assert(!deleted_);
     413       63051 :         if (running_) {
     414         275 :             running_ = false;
     415         275 :             assert(current_runner_);
     416         275 :             TaskScheduler *scheduler = TaskScheduler::GetInstance();
     417             :             TaskScheduler::CancelReturnCode cancel_code =
     418         275 :                 scheduler->Cancel(current_runner_);
     419         275 :             assert(cancel_code == TaskScheduler::CANCELLED);
     420         275 :             current_runner_ = NULL;
     421             :         }
     422       63051 :         ResetHighWaterMark();
     423       63051 :         ResetLowWaterMark();
     424             :         WorkQueueDelete<QueueEntryT> deleter;
     425       63051 :         deleter(queue_, delete_entries);
     426       63051 :         queue_.clear();
     427       63051 :         count_ = 0;
     428       63051 :         deleted_ = true;
     429       63051 :     }
     430             : 
     431     2590202 :     size_t AtomicIncrementQueueCount(QueueEntryT *entry) {
     432     5180404 :         return count_.fetch_add(1) + 1;
     433             :     }
     434             : 
     435     2624538 :     size_t AtomicDecrementQueueCount(QueueEntryT *entry) {
     436     5249076 :         return count_.fetch_sub(1) - 1;
     437             :     }
     438             : 
     439     2590888 :     void ProcessHighWaterMarks(size_t count) {
     440     2590888 :         watermarks_.ProcessHighWaterMarks(count);
     441     2590754 :     }
     442             : 
     443     2625146 :     void ProcessLowWaterMarks(size_t count) {
     444     2625146 :         watermarks_.ProcessLowWaterMarks(count);
     445     2625040 :     }
     446             : 
     447     2590655 :     bool EnqueueInternal(QueueEntryT entry) {
     448     2590655 :         enqueues_++;
     449     2590655 :         size_t ncount(AtomicIncrementQueueCount(&entry));
     450     2590833 :         if (ncount > max_queue_len_)
     451      634121 :             max_queue_len_ = ncount;
     452     2590833 :         ProcessHighWaterMarks(ncount);
     453     2590673 :         queue_.push(entry);
     454     2590718 :         MayBeStartRunner();
     455     2590843 :         return ncount < size_;
     456             :     }
     457             : 
     458      100572 :     bool EnqueueInternalLocked(QueueEntryT entry) {
     459      100572 :         std::scoped_lock lock(water_mutex_);
     460      201144 :         return EnqueueInternal(entry);
     461      100572 :     }
     462             : 
     463          77 :     bool EnqueueBounded(QueueEntryT entry) {
     464          77 :         size_t ncount(AtomicIncrementQueueCount(&entry));
     465          77 :         if (ncount > max_queue_len_)
     466          20 :             max_queue_len_ = ncount;
     467          77 :         if (ncount < size_) {
     468          72 :             enqueues_++;
     469          72 :             ProcessHighWaterMarks(ncount);
     470          72 :             queue_.push(entry);
     471          72 :             MayBeStartRunner();
     472          72 :             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     1513323 :     bool RunnerAbortLocked() {
     486     3002527 :         return (disabled_ || shutdown_scheduled_ ||
     487     3002452 :                 (!start_runner_.empty() && !start_runner_()));
     488             :     }
     489             : 
     490      744521 :     bool RunnerAbort() {
     491      744521 :         std::scoped_lock lock(mutex_);
     492     1489174 :         return RunnerAbortLocked();
     493      744535 :     }
     494             : 
     495      744577 :     bool RunnerDone() {
     496      744577 :         std::scoped_lock lock(mutex_);
     497      744601 :         bool done = false;
     498      744601 :         if (queue_.empty() || RunnerAbortLocked()) {
     499      217004 :             done = true;
     500      217004 :             OnExit(done);
     501      217001 :             current_runner_ = NULL;
     502      217001 :             running_ = false;
     503      217001 :             if (shutdown_scheduled_) {
     504           1 :                 ShutdownLocked(delete_entries_on_shutdown_);
     505             :             }
     506             :         } else {
     507      527541 :             OnExit(done);
     508      527541 :             running_ = true;
     509             :         }
     510      744595 :         return done;
     511      744542 :     }
     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