| 682 | } |
| 683 | |
| 684 | void AsyncConnection::DelayedDelivery::discard() { |
| 685 | stop_dispatch = true; |
| 686 | center->submit_to(center->get_id(), |
| 687 | [this]() mutable { |
| 688 | std::lock_guard<std::mutex> l(delay_lock); |
| 689 | while (!delay_queue.empty()) { |
| 690 | Message *m = delay_queue.front(); |
| 691 | dispatch_queue->dispatch_throttle_release( |
| 692 | m->get_dispatch_throttle_size()); |
| 693 | m->put(); |
| 694 | delay_queue.pop_front(); |
| 695 | } |
| 696 | for (auto i : register_time_events) |
| 697 | center->delete_time_event(i); |
| 698 | register_time_events.clear(); |
| 699 | stop_dispatch = false; |
| 700 | }, |
| 701 | true); |
| 702 | } |
| 703 | |
| 704 | void AsyncConnection::DelayedDelivery::flush() { |
| 705 | stop_dispatch = true; |
no test coverage detected