userver: userver/server/request/response_base.hpp Source File
Loading...
Searching...
No Matches
response_base.hpp
Go to the documentation of this file.
1#pragma once
2
3/// @file userver/server/request/response_base.hpp
4/// @brief @copybrief server::request::ResponseBase
5
6#include <atomic>
7#include <chrono>
8#include <functional>
9#include <limits>
10#include <optional>
11#include <string>
12
13#include <userver/concurrent/queue.hpp>
14#include <userver/concurrent/striped_counter.hpp>
15#include <userver/engine/single_consumer_event.hpp>
16#include <userver/utils/fast_pimpl.hpp>
17
18USERVER_NAMESPACE_BEGIN
19
20/// @cond
21// TODO: server internals. remove from a public interface
22namespace server::http::impl {
23
24struct Http2StreamEvent {
25 std::int32_t stream_id{-1};
26 std::string body_part{};
27 bool is_end{false};
28};
29
30// The order is fifo in the context of a single producer. So we are tolerant to
31// reordering between producers
32using Http2StreamEventQueue = concurrent::NonFifoMpscQueue<Http2StreamEvent>;
33
34class Http2StreamEventProducer final {
35public:
36 Http2StreamEventProducer(Http2StreamEventQueue& queue, engine::SingleConsumerEvent& event);
37
38 void PushEvent(Http2StreamEvent event, engine::Deadline deadline = {});
39
40 void CloseStream(std::int32_t id);
41
42private:
43 Http2StreamEventQueue::Producer producer_;
44 engine::SingleConsumerEvent& event_;
45};
46
47} // namespace server::http::impl
48/// @endcond
49
50namespace engine::io {
51class RwBase;
52} // namespace engine::io
53
54namespace server::request {
55
56class ResponseDataAccounter final {
57public:
58 void StartRequest(std::chrono::steady_clock::time_point create_time);
59
60 void StopRequest(std::size_t size, std::chrono::steady_clock::time_point create_time);
61
62 void ReaccountRequest(
63 std::size_t old_size,
64 std::chrono::steady_clock::time_point old_create_time,
65 std::size_t new_size,
66 std::chrono::steady_clock::time_point new_create_time
67 );
68
69 std::size_t GetPendingResponsesSizeInBytes() const { return pending_responses_size_in_bytes_; }
70
71 std::size_t GetPendingResponsesCount() const { return pending_responses_count_.NonNegativeRead(); }
72
73 std::size_t GetMaxPendingResponsesSizeInBytes() const { return max_pending_responses_size_in_bytes_; }
74
75 void SetMaxPendingResponsesSizeInBytes(size_t size) { max_pending_responses_size_in_bytes_ = size; }
76
77 std::chrono::milliseconds GetAvgRequestTime() const;
78
79private:
80 std::atomic<std::size_t> pending_responses_size_in_bytes_{0};
81 std::atomic<std::size_t> max_pending_responses_size_in_bytes_{std::numeric_limits<std::size_t>::max()};
82 concurrent::StripedCounter pending_responses_count_{};
83 concurrent::StripedCounter time_sum_{};
84};
85
86// TODO: merge with HttpResponse
87
88/// @brief Base class for all the server responses.
90public:
91 explicit ResponseBase(ResponseDataAccounter& data_accounter);
92 ResponseBase(const ResponseBase&) = delete;
93 ResponseBase(ResponseBase&&) = delete;
94 virtual ~ResponseBase() noexcept;
95
96 void SetData(std::string data);
97 const std::string& GetData() const { return data_; }
98 std::string&& ExtractData() { return std::move(data_); }
99
100 virtual bool IsBodyStreamed() const = 0;
101 virtual bool WaitForHeadersEnd() = 0;
102 virtual void SetHeadersEnd() = 0;
103
104 /// @cond
105 // TODO: server internals. remove from a public interface
106 void SetReady();
107 void SetReady(std::chrono::steady_clock::time_point now);
108 virtual void SetSendFailed(std::chrono::steady_clock::time_point failure_time);
109 bool IsLimitReached() const;
110
111 bool IsReady() const { return ready_time_ != kUnset; }
112 bool IsSent() const noexcept { return sent_time_ != kUnset; }
113 size_t BytesSent() const { return bytes_sent_; }
114 std::chrono::steady_clock::time_point ReadyTime() const { return ready_time_; }
115 std::chrono::steady_clock::time_point SentTime() const { return sent_time_; }
116 virtual void SendResponse(engine::io::RwBase& socket) = 0;
117
118 virtual void SetStatusServiceUnavailable() = 0;
119 virtual void SetStatusOk() = 0;
120 virtual void SetStatusNotFound() = 0;
121
122 // HTTP/2.0 only
123 void SetStreamId(std::int32_t stream_id);
124 std::optional<std::int32_t> GetStreamId() const { return stream_id_; }
125 void SetStreamProdicer(http::impl::Http2StreamEventProducer&& producer);
126 http::impl::Http2StreamEventProducer GetStreamProducer();
127 /// @endcond
128
129protected:
130 ResponseBase(ResponseDataAccounter& data_account, std::chrono::steady_clock::time_point now);
131
132 void SetSent(std::size_t bytes_sent, std::chrono::steady_clock::time_point sent_time);
133
134private:
135 static constexpr auto kUnset = std::chrono::steady_clock::time_point::min();
136
137 ResponseDataAccounter& accounter_;
138 std::string data_;
139 std::chrono::steady_clock::time_point create_time_;
140 std::chrono::steady_clock::time_point ready_time_{kUnset};
141 std::chrono::steady_clock::time_point sent_time_{kUnset};
142 std::size_t accounted_size_ = 0;
143 size_t bytes_sent_ = 0;
144 std::optional<std::int32_t> stream_id_;
145 std::optional<http::impl::Http2StreamEventProducer> producer_{};
146};
147
148} // namespace server::request
149
150USERVER_NAMESPACE_END