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__ */