aboutsummaryrefslogtreecommitdiffstats
path: root/library/cpp/grpc/server/grpc_request.cpp
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