| 702 | } |
| 703 | |
| 704 | void AsyncConnection::DelayedDelivery::flush() { |
| 705 | stop_dispatch = true; |
| 706 | center->submit_to( |
| 707 | center->get_id(), [this] () mutable { |
| 708 | std::lock_guard<std::mutex> l(delay_lock); |
| 709 | while (!delay_queue.empty()) { |
| 710 | Message *m = delay_queue.front(); |
| 711 | if (msgr->ms_can_fast_dispatch(m)) { |
| 712 | dispatch_queue->fast_dispatch(m); |
| 713 | } else { |
| 714 | dispatch_queue->enqueue(m, m->get_priority(), conn_id); |
| 715 | } |
| 716 | delay_queue.pop_front(); |
| 717 | } |
| 718 | for (auto i : register_time_events) |
| 719 | center->delete_time_event(i); |
| 720 | register_time_events.clear(); |
| 721 | stop_dispatch = false; |
| 722 | }, true); |
| 723 | } |
| 724 | |
| 725 | void AsyncConnection::send_keepalive() |
| 726 | { |
no test coverage detected