/src/pistache/include/pistache/transport.h
Line | Count | Source |
1 | | /* |
2 | | * SPDX-FileCopyrightText: 2016 Mathieu Stefani |
3 | | * |
4 | | * SPDX-License-Identifier: Apache-2.0 |
5 | | */ |
6 | | |
7 | | /* |
8 | | Mathieu Stefani, 26 janvier 2016 |
9 | | |
10 | | Transport TCP layer |
11 | | */ |
12 | | |
13 | | #pragma once |
14 | | |
15 | | #include <pistache/winornix.h> |
16 | | |
17 | | #include PST_SYS_RESOURCE_HDR // for PST_RUSAGE + PST_GETRUSAGE |
18 | | |
19 | | #include <pistache/async.h> |
20 | | #include <pistache/mailbox.h> |
21 | | #include <pistache/pist_quote.h> |
22 | | #include <pistache/pist_timelog.h> |
23 | | #include <pistache/reactor.h> |
24 | | #include <pistache/stream.h> |
25 | | |
26 | | #include <chrono> |
27 | | #include <deque> |
28 | | #include <memory> |
29 | | #include <mutex> |
30 | | #include <optional> |
31 | | #include <unordered_map> |
32 | | |
33 | | namespace Pistache::Tcp |
34 | | { |
35 | | |
36 | | class Peer; |
37 | | class Handler; |
38 | | |
39 | | class Transport : public Aio::Handler |
40 | | { |
41 | | public: |
42 | | explicit Transport(const std::shared_ptr<Tcp::Handler>& handler); |
43 | | |
44 | | Transport(const Transport&) = delete; |
45 | | Transport& operator=(const Transport&) = delete; |
46 | | |
47 | | ~Transport(); |
48 | | |
49 | | void init(const std::shared_ptr<Tcp::Handler>& handler); |
50 | | |
51 | | void registerPoller(Polling::Epoll& poller) override; |
52 | | void unregisterPoller(Polling::Epoll& poller) override; |
53 | | |
54 | | void handleNewPeer(const std::shared_ptr<Peer>& peer); |
55 | | void onReady(const Aio::FdSet& fds) override; |
56 | | |
57 | | template <typename Buf> |
58 | | Async::Promise<PST_SSIZE_T> asyncWrite(Fd fd, const Buf& buffer, |
59 | | int flags = 0 |
60 | | #ifdef _USE_LIBEVENT_LIKE_APPLE |
61 | | , |
62 | | bool msg_more_style = false |
63 | | #endif |
64 | | ) |
65 | 0 | { |
66 | | // Always enqueue reponses for sending. Giving preference to |
67 | | // consumer context means chunked responses could be sent out of |
68 | | // order. |
69 | | // |
70 | | // Note: fd could be PS_FD_EMPTY |
71 | 0 | return Async::Promise<PST_SSIZE_T>( |
72 | 0 | [&, this](Async::Deferred<PST_SSIZE_T> deferred) mutable { |
73 | 0 | BufferHolder holder { buffer }; |
74 | 0 | WriteEntry write(std::move(deferred), std::move(holder), |
75 | 0 | fd, flags |
76 | | #ifdef _USE_LIBEVENT_LIKE_APPLE |
77 | | , |
78 | | msg_more_style |
79 | | #endif |
80 | 0 | ); |
81 | 0 | writesQueue.push(std::move(write)); |
82 | 0 | }); Unexecuted instantiation: Pistache::Tcp::Transport::asyncWrite<Pistache::RawBuffer>(int, Pistache::RawBuffer const&, int)::{lambda(Pistache::Async::Deferred<long>)#1}::operator()(Pistache::Async::Deferred<long>)Unexecuted instantiation: Pistache::Tcp::Transport::asyncWrite<Pistache::FileBuffer>(int, Pistache::FileBuffer const&, int)::{lambda(Pistache::Async::Deferred<long>)#1}::operator()(Pistache::Async::Deferred<long>) |
83 | 0 | } Unexecuted instantiation: Pistache::Async::Promise<long> Pistache::Tcp::Transport::asyncWrite<Pistache::RawBuffer>(int, Pistache::RawBuffer const&, int) Unexecuted instantiation: Pistache::Async::Promise<long> Pistache::Tcp::Transport::asyncWrite<Pistache::FileBuffer>(int, Pistache::FileBuffer const&, int) |
84 | | |
85 | | Async::Promise<PST_RUSAGE> load() |
86 | 0 | { |
87 | 0 | return Async::Promise<PST_RUSAGE>([this](Async::Deferred<PST_RUSAGE> deferred) { |
88 | 0 | PS_TIMEDBG_START_CURLY; |
89 | 0 |
|
90 | 0 | loadRequest_ = std::move(deferred); |
91 | 0 | notifier.notify(); |
92 | 0 | }); |
93 | 0 | } |
94 | | |
95 | | template <typename Duration> |
96 | | void armTimer(Fd fd, Duration timeout, Async::Deferred<uint64_t> deferred) |
97 | | { |
98 | | PS_LOG_DEBUG_ARGS("Fd %" PIST_QUOTE(PS_FD_PRNTFCD), fd); |
99 | | |
100 | | armTimerMs(fd, |
101 | | std::chrono::duration_cast<std::chrono::milliseconds>(timeout), |
102 | | std::move(deferred)); |
103 | | } |
104 | | |
105 | | void disarmTimer(Fd fd); |
106 | | |
107 | | std::shared_ptr<Aio::Handler> clone() const override; |
108 | | |
109 | | void flush(); |
110 | | |
111 | | std::deque<std::shared_ptr<Peer>> getAllPeer(); |
112 | | |
113 | | #ifdef _USE_LIBEVENT |
114 | | std::shared_ptr<EventMethEpollEquiv> getEventMethEpollEquiv() |
115 | | { |
116 | | return (epoll_fd); |
117 | | } |
118 | | #endif |
119 | | |
120 | | void closeFd(Fd fd); |
121 | | |
122 | | // !!!! Make protected like removePeer |
123 | | void removeAllPeers(); // cleans up toWrite and does CLOSE_FD on each |
124 | | |
125 | | private: |
126 | | enum WriteStatus { FirstTry, |
127 | | Retry }; |
128 | | |
129 | | struct BufferHolder |
130 | | { |
131 | | enum Type { Raw, |
132 | | File }; |
133 | | |
134 | | explicit BufferHolder(const RawBuffer& buffer, off_t offset = 0) |
135 | 0 | : _raw(buffer) |
136 | 0 | , size_(buffer.size()) |
137 | 0 | , offset_(offset) |
138 | 0 | , type(Raw) |
139 | 0 | { } |
140 | | |
141 | | explicit BufferHolder(const FileBuffer& buffer, off_t offset = 0) |
142 | 0 | : _fd(buffer.fd()) |
143 | 0 | , size_(buffer.size()) |
144 | 0 | , offset_(offset) |
145 | 0 | , type(File) |
146 | 0 | { } |
147 | | |
148 | 0 | bool isFile() const { return type == File; } |
149 | 0 | bool isRaw() const { return type == Raw; } |
150 | 0 | size_t size() const { return size_; } |
151 | 0 | size_t offset() const { return static_cast<size_t>(offset_); } |
152 | | |
153 | | int fd() const |
154 | 0 | { |
155 | 0 | if (!isFile()) |
156 | 0 | throw std::runtime_error("Tried to retrieve fd of a non-filebuffer"); |
157 | 0 | return _fd; |
158 | 0 | } |
159 | | |
160 | | RawBuffer raw() const |
161 | 0 | { |
162 | 0 | if (!isRaw()) |
163 | 0 | throw std::runtime_error("Tried to retrieve raw data of a non-buffer"); |
164 | 0 | return _raw; |
165 | 0 | } |
166 | | |
167 | | BufferHolder detach(off_t offset = 0) |
168 | 0 | { |
169 | 0 | if (!isRaw()) |
170 | 0 | return BufferHolder(_fd, size_, offset); |
171 | | |
172 | 0 | auto detached = _raw.copy(static_cast<size_t>(offset)); |
173 | 0 | return BufferHolder(detached); |
174 | 0 | } |
175 | | |
176 | | private: |
177 | | BufferHolder(int fd, // regular file desc ("int") even for libevent |
178 | | size_t size, off_t offset = 0) |
179 | 0 | : _fd(fd) |
180 | 0 | , size_(size) |
181 | 0 | , offset_(offset) |
182 | 0 | , type(File) |
183 | 0 | { } |
184 | | |
185 | | RawBuffer _raw; |
186 | | int _fd; // regular old file desc ("int") even in libevent case |
187 | | |
188 | | size_t size_ = 0; |
189 | | off_t offset_ = 0; |
190 | | Type type; |
191 | | }; |
192 | | |
193 | | struct WriteEntry |
194 | | { |
195 | | WriteEntry(Async::Deferred<PST_SSIZE_T> deferred_, BufferHolder buffer_, |
196 | | Fd peerFd_, int flags_ = 0 |
197 | | #ifdef _USE_LIBEVENT_LIKE_APPLE |
198 | | , |
199 | | bool msg_more_style_ = false |
200 | | #endif |
201 | | ) |
202 | 0 | : deferred(std::move(deferred_)) |
203 | 0 | , buffer(std::move(buffer_)) |
204 | 0 | , flags(flags_) |
205 | | #ifdef _USE_LIBEVENT_LIKE_APPLE |
206 | | , msg_more_style(msg_more_style_) |
207 | | #endif |
208 | 0 | , peerFd(peerFd_) |
209 | 0 | { } |
210 | | |
211 | | Async::Deferred<PST_SSIZE_T> deferred; |
212 | | BufferHolder buffer; |
213 | | int flags = 0; |
214 | | #ifdef _USE_LIBEVENT_LIKE_APPLE |
215 | | bool msg_more_style = false; |
216 | | #endif |
217 | | Fd peerFd = PS_FD_EMPTY; |
218 | | }; |
219 | | |
220 | | struct TimerEntry |
221 | | { |
222 | | TimerEntry(Fd fd_, std::chrono::milliseconds value_, |
223 | | Async::Deferred<uint64_t> deferred_) |
224 | 0 | : fd(fd_) |
225 | 0 | , value(value_) |
226 | 0 | , deferred(std::move(deferred_)) |
227 | 0 | , active() |
228 | 0 | { |
229 | 0 | active.store(true, std::memory_order_relaxed); |
230 | 0 | } |
231 | | |
232 | | TimerEntry(TimerEntry&& other) |
233 | 0 | : fd(other.fd) |
234 | 0 | , value(other.value) |
235 | 0 | , deferred(std::move(other.deferred)) |
236 | 0 | , active(other.active.load()) |
237 | 0 | { } |
238 | | |
239 | 0 | void disable() { active.store(false, std::memory_order_relaxed); } |
240 | | |
241 | 0 | bool isActive() const { return active.load(std::memory_order_relaxed); } |
242 | | |
243 | | Fd fd; |
244 | | std::chrono::milliseconds value; |
245 | | Async::Deferred<uint64_t> deferred; |
246 | | std::atomic<bool> active; |
247 | | }; |
248 | | |
249 | | struct PeerEntry |
250 | | { |
251 | | explicit PeerEntry(std::shared_ptr<Peer> peer_) |
252 | 0 | : peer(std::move(peer_)) |
253 | 0 | { } |
254 | | |
255 | | std::shared_ptr<Peer> peer; |
256 | | }; |
257 | | using Lock = std::mutex; |
258 | | using Guard = std::lock_guard<Lock>; |
259 | | |
260 | | #ifdef _USE_LIBEVENT |
261 | | std::shared_ptr<EventMethEpollEquiv> epoll_fd; |
262 | | #endif |
263 | | |
264 | | PollableQueue<WriteEntry> writesQueue; |
265 | | std::unordered_map<Fd, std::deque<WriteEntry>> toWrite; |
266 | | Lock toWriteLock; |
267 | | |
268 | | PollableQueue<TimerEntry> timersQueue; |
269 | | std::unordered_map<FdConst, TimerEntry> timers; |
270 | | |
271 | | PollableQueue<PeerEntry> peersQueue; |
272 | | |
273 | | Async::Deferred<PST_RUSAGE> loadRequest_; |
274 | | NotifyFd notifier; |
275 | | |
276 | | std::shared_ptr<Tcp::Handler> handler_; |
277 | | |
278 | | #ifdef _USE_LIBEVENT_LIKE_APPLE |
279 | | int tcp_prot_num_; // TCP protocol num on this host per getprotobyname |
280 | | #endif |
281 | | |
282 | | protected: |
283 | | void removePeer(const std::shared_ptr<Peer>& peer); |
284 | | |
285 | | // Without the use of peers_mutex_ to protect peers_, http_server_test |
286 | | // multiple_client_with_requests_to_multithreaded_server fails |
287 | | // intermittently (~1 time in 10 - likely highly environment |
288 | | // dependent). The test is doing 3 client requests to the server, one |
289 | | // Peer per request; it fails when two of the requests are using the |
290 | | // same Peer. Which appears to happen when the peers_ unordered_map |
291 | | // gets messed up due to a threading issue. |
292 | | mutable std::mutex peers_mutex_; |
293 | | std::unordered_map<Fd, std::shared_ptr<Peer>> peers_; |
294 | | |
295 | | private: |
296 | | bool isPeerFd(FdConst fd) const; |
297 | | bool isPeerFdNoPeersMutexLock(FdConst fd) const; |
298 | | bool isTimerFd(FdConst fd) const; |
299 | | bool isPeerFd(Polling::Tag tag) const; |
300 | | bool isTimerFd(Polling::Tag tag) const; |
301 | | |
302 | | std::shared_ptr<Peer> getPeer(FdConst fd); |
303 | | std::shared_ptr<Peer> getPeer(Polling::Tag tag); |
304 | | |
305 | | void armTimerMs(Fd fd, std::chrono::milliseconds value, |
306 | | Async::Deferred<uint64_t> deferred); |
307 | | |
308 | | void armTimerMsImpl(TimerEntry entry); |
309 | | |
310 | | // This will attempt to drain the write queue for the fd |
311 | | void asyncWriteImpl(Fd fd); |
312 | | |
313 | | #ifdef _USE_LIBEVENT_LIKE_APPLE |
314 | | void configureMsgMoreStyle(Fd fd, bool msg_more_style); |
315 | | #endif |
316 | | |
317 | | PST_SSIZE_T sendRawBuffer(Fd fd, const char* buffer, size_t len, int flags |
318 | | #ifdef _USE_LIBEVENT_LIKE_APPLE |
319 | | , |
320 | | bool msg_more_style |
321 | | #endif |
322 | | ); |
323 | | PST_SSIZE_T sendFile(Fd fd, int file, off_t offset, size_t len); |
324 | | |
325 | | void handlePeerDisconnection(const std::shared_ptr<Peer>& peer); |
326 | | void handleIncoming(const std::shared_ptr<Peer>& peer); |
327 | | void handleWriteQueue(bool flush = false); |
328 | | void handleTimerQueue(); |
329 | | void handlePeerQueue(); |
330 | | void handleNotify(); |
331 | | void handleTimer(TimerEntry entry); |
332 | | void handlePeer(const std::shared_ptr<Peer>& peer); |
333 | | }; |
334 | | |
335 | | } // namespace Pistache::Tcp |