blob: 7030992173de1cd530223bdab0db091f643c45c9 (
plain) (
blame)
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
|
#include "grpc_request.h"
namespace NGrpc {
const char* GRPC_USER_AGENT_HEADER = "user-agent";
class TStreamAdaptor: public IStreamAdaptor {
public:
TStreamAdaptor()
: StreamIsReady_(true)
{}
void Enqueue(std::function<void()>&& fn, bool urgent) override {
with_lock(Mtx_) {
if (!UrgentQueue_.empty() || !NormalQueue_.empty()) {
Y_ABORT_UNLESS(!StreamIsReady_);
}
auto& queue = urgent ? UrgentQueue_ : NormalQueue_;
if (StreamIsReady_ && queue.empty()) {
StreamIsReady_ = false;
} else {
queue.push_back(std::move(fn));
return;
}
}
fn();
}
size_t ProcessNext() override {
size_t left = 0;
std::function<void()> fn;
with_lock(Mtx_) {
Y_ABORT_UNLESS(!StreamIsReady_);
auto& queue = UrgentQueue_.empty() ? NormalQueue_ : UrgentQueue_;
if (queue.empty()) {
// Both queues are empty
StreamIsReady_ = true;
} else {
fn = std::move(queue.front());
queue.pop_front();
left = UrgentQueue_.size() + NormalQueue_.size();
}
}
if (fn)
fn();
return left;
}
private:
bool StreamIsReady_;
TList<std::function<void()>> NormalQueue_;
TList<std::function<void()>> UrgentQueue_;
TMutex Mtx_;
};
IStreamAdaptor::TPtr CreateStreamAdaptor() {
return std::make_unique<TStreamAdaptor>();
}
} // namespace NGrpc
|