blob: 6f177973bf24c20ee9b48e2510bfd3749c52a684 (
plain)
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
|
//
// Copyright Aliaksei Levin (levlam@telegram.org), Arseny Smirnov (arseny30@gmail.com) 2014-2023
//
// Distributed under the Boost Software License, Version 1.0. (See accompanying
// file LICENSE_1_0.txt or copy at http://www.boost.org/LICENSE_1_0.txt)
//
#include "td/actor/MultiTimeout.h"
#include "td/utils/logging.h"
namespace td {
bool MultiTimeout::has_timeout(int64 key) const {
return items_.count(Item(key)) > 0;
}
void MultiTimeout::set_timeout_at(int64 key, double timeout) {
LOG(DEBUG) << "Set " << get_name() << " for " << key << " in " << timeout - Time::now();
auto item = items_.emplace(key);
auto heap_node = static_cast<HeapNode *>(const_cast<Item *>(&*item.first));
if (heap_node->in_heap()) {
CHECK(!item.second);
bool need_update_timeout = heap_node->is_top();
timeout_queue_.fix(timeout, heap_node);
if (need_update_timeout || heap_node->is_top()) {
update_timeout("set_timeout");
}
} else {
CHECK(item.second);
timeout_queue_.insert(timeout, heap_node);
if (heap_node->is_top()) {
update_timeout("set_timeout 2");
}
}
}
void MultiTimeout::add_timeout_at(int64 key, double timeout) {
LOG(DEBUG) << "Add " << get_name() << " for " << key << " in " << timeout - Time::now();
auto item = items_.emplace(key);
auto heap_node = static_cast<HeapNode *>(const_cast<Item *>(&*item.first));
if (heap_node->in_heap()) {
CHECK(!item.second);
} else {
CHECK(item.second);
timeout_queue_.insert(timeout, heap_node);
if (heap_node->is_top()) {
update_timeout("add_timeout");
}
}
}
void MultiTimeout::cancel_timeout(int64 key) {
LOG(DEBUG) << "Cancel " << get_name() << " for " << key;
auto item = items_.find(Item(key));
if (item != items_.end()) {
auto heap_node = static_cast<HeapNode *>(const_cast<Item *>(&*item));
CHECK(heap_node->in_heap());
bool need_update_timeout = heap_node->is_top();
timeout_queue_.erase(heap_node);
items_.erase(item);
if (need_update_timeout) {
update_timeout("cancel_timeout");
}
}
}
void MultiTimeout::update_timeout(const char *source) {
if (items_.empty()) {
LOG(DEBUG) << "Cancel timeout of " << get_name();
LOG_CHECK(timeout_queue_.empty()) << get_name() << ' ' << source;
LOG_CHECK(Actor::has_timeout()) << get_name() << ' ' << source;
Actor::cancel_timeout();
} else {
LOG(DEBUG) << "Set timeout of " << get_name() << " in " << timeout_queue_.top_key() - Time::now_cached();
Actor::set_timeout_at(timeout_queue_.top_key());
}
}
vector<int64> MultiTimeout::get_expired_keys(double now) {
vector<int64> expired_keys;
while (!timeout_queue_.empty() && timeout_queue_.top_key() < now) {
int64 key = static_cast<Item *>(timeout_queue_.pop())->key;
items_.erase(Item(key));
expired_keys.push_back(key);
}
return expired_keys;
}
void MultiTimeout::timeout_expired() {
vector<int64> expired_keys = get_expired_keys(Time::now_cached());
if (!items_.empty()) {
update_timeout("timeout_expired");
}
for (auto key : expired_keys) {
callback_(data_, key);
}
}
void MultiTimeout::run_all() {
vector<int64> expired_keys = get_expired_keys(Time::now_cached() + 1e10);
if (!expired_keys.empty()) {
update_timeout("run_all");
}
for (auto key : expired_keys) {
callback_(data_, key);
}
}
} // namespace td
|