aboutsummaryrefslogtreecommitdiffstats
path: root/library/cpp/grpc/server/grpc_request.cpp
blob: 33264fe6f29738ceabf57f5cc3c2096b3adee865 (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_VERIFY(!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_VERIFY(!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