/src/curl_fuzzer/proto_fuzzer/mock_server.cc
Line | Count | Source |
1 | | /* |
2 | | * Copyright (C) Max Dymond, <cmeister2@gmail.com>, et al. |
3 | | * |
4 | | * SPDX-License-Identifier: curl |
5 | | */ |
6 | | |
7 | | /// @file |
8 | | /// @brief Implementation of MockConnection and MockServer. |
9 | | |
10 | | #include "proto_fuzzer/mock_server.h" |
11 | | |
12 | | #include <fcntl.h> |
13 | | #include <string.h> |
14 | | #include <sys/select.h> |
15 | | #include <sys/socket.h> |
16 | | #include <sys/types.h> |
17 | | #include <unistd.h> |
18 | | |
19 | | #include <algorithm> |
20 | | #include <cstddef> |
21 | | #include <cstdint> |
22 | | #include <limits> |
23 | | #include <string> |
24 | | #include <vector> |
25 | | |
26 | | #include "proto_fuzzer/multi_socket_driver.h" |
27 | | #include "proto_fuzzer/scenario_limits.h" |
28 | | #include "proto_fuzzer/ws_frame.h" |
29 | | |
30 | | namespace proto_fuzzer { |
31 | | |
32 | | namespace { |
33 | | |
34 | | // fd_set can only represent file descriptors < FD_SETSIZE. Reject any pair that couldn't participate in select() |
35 | | // without memory corruption. |
36 | 291k | bool FdFitsInFdSet(int fd) { return fd >= 0 && fd < FD_SETSIZE; } |
37 | | |
38 | | } // namespace |
39 | | |
40 | | /// @class proto_fuzzer::MockConnection |
41 | | /// @brief Owns one half of a socketpair used to feed canned responses to libcurl. The destructor closes the server-side |
42 | | /// fd; the client-side fd is handed to libcurl via CURLOPT_OPENSOCKETFUNCTION and becomes curl's to close. |
43 | | |
44 | | /// Construct a non-blocking AF_UNIX/SOCK_STREAM socketpair. Both fds are validated to fit inside FD_SETSIZE; on any |
45 | | /// failure ok() returns false and the instance is unusable. |
46 | 145k | MockConnection::MockConnection() : server_fd_(-1), client_fd_(-1), drain_limit_(0) { |
47 | 145k | int fds[2]; |
48 | | |
49 | 145k | if (socketpair(AF_UNIX, SOCK_STREAM, 0, fds) != 0) { |
50 | 0 | return; |
51 | 0 | } |
52 | | |
53 | | // The fds must be small enough to fit in an fd_set for select(). If not, close them and fail the constructor. |
54 | 145k | if (!FdFitsInFdSet(fds[0]) || !FdFitsInFdSet(fds[1])) { |
55 | 0 | close(fds[0]); |
56 | 0 | close(fds[1]); |
57 | 0 | return; |
58 | 0 | } |
59 | | |
60 | | // Set the server-side fd non-blocking so we can write to it without risk of hanging the fuzzer. |
61 | 145k | int flags = fcntl(fds[0], F_GETFL, 0); |
62 | 145k | if (flags < 0 || fcntl(fds[0], F_SETFL, flags | O_NONBLOCK) < 0) { |
63 | 0 | close(fds[0]); |
64 | 0 | close(fds[1]); |
65 | 0 | return; |
66 | 0 | } |
67 | | |
68 | | // Success: store the file descriptors. |
69 | 145k | server_fd_ = fds[0]; |
70 | 145k | client_fd_ = fds[1]; |
71 | 145k | } |
72 | | |
73 | | /// Close the server-side fd (and the client-side fd if it was never handed off via take_client_fd()). |
74 | 145k | MockConnection::~MockConnection() { |
75 | 145k | if (server_fd_ >= 0) { |
76 | 145k | close(server_fd_); |
77 | 145k | } |
78 | 145k | if (client_fd_ >= 0) { |
79 | 0 | close(client_fd_); |
80 | 0 | } |
81 | 145k | } |
82 | | |
83 | | /// @return true if the underlying socketpair was set up successfully. |
84 | 371k | bool MockConnection::ok() const { return server_fd_ >= 0; } |
85 | | |
86 | | /// @return the server-side fd (still owned by this MockConnection). |
87 | 50.1k | int MockConnection::server_fd() const { return server_fd_; } |
88 | | |
89 | | /// Query the client endpoint while this object still owns it. A zero result is |
90 | | /// deliberately ambiguous between an invalid fd and a platform query failure: |
91 | | /// callers need only distinguish a verified capacity from every unsafe case. |
92 | 3.85k | std::size_t MockConnection::client_send_buffer_size() const { |
93 | 3.85k | if (client_fd_ < 0) { |
94 | 0 | return 0; |
95 | 0 | } |
96 | 3.85k | int size = 0; |
97 | 3.85k | socklen_t length = sizeof(size); |
98 | 3.85k | if (getsockopt(client_fd_, SOL_SOCKET, SO_SNDBUF, &size, &length) != 0 || size <= 0) { |
99 | 0 | return 0; |
100 | 0 | } |
101 | 3.85k | return static_cast<std::size_t>(size); |
102 | 3.85k | } |
103 | | |
104 | | /// Establish and verify the capacity required by a synchronous protocol |
105 | | /// driver. Linux may transform socket-buffer requests, so the post-set query |
106 | | /// is the contract rather than assuming setsockopt accepted the exact value. |
107 | 3.85k | bool MockConnection::EnsureClientSendBufferSize(std::size_t minimum) { |
108 | 3.85k | if (client_send_buffer_size() >= minimum) { |
109 | 3.85k | return true; |
110 | 3.85k | } |
111 | 0 | if (client_fd_ < 0 || minimum > static_cast<std::size_t>(std::numeric_limits<int>::max())) { |
112 | 0 | return false; |
113 | 0 | } |
114 | 0 | const int requested = static_cast<int>(minimum); |
115 | 0 | if (setsockopt(client_fd_, SOL_SOCKET, SO_SNDBUF, &requested, sizeof(requested)) != 0) { |
116 | 0 | return false; |
117 | 0 | } |
118 | 0 | return client_send_buffer_size() >= minimum; |
119 | 0 | } |
120 | | |
121 | | /// Hand the client-side fd to libcurl. After this call the caller owns the fd and the MockConnection will not close it |
122 | | /// on destruction. |
123 | | /// @return the client-side socket fd as a curl_socket_t. |
124 | 145k | curl_socket_t MockConnection::take_client_fd() { |
125 | 145k | int fd = client_fd_; |
126 | 145k | client_fd_ = -1; |
127 | 145k | return static_cast<curl_socket_t>(fd); |
128 | 145k | } |
129 | | |
130 | | /// Write 'size' bytes from 'data' to the server fd, looping until the whole |
131 | | /// buffer is sent or a short/failed write occurs. MSG_NOSIGNAL is a harness |
132 | | /// invariant rather than a curl behavior choice: a response race may leave the |
133 | | /// peer closed, and that must look like an ordinary failed mock write instead |
134 | | /// of terminating the fuzz process with SIGPIPE. |
135 | | /// @param data Buffer to send. |
136 | | /// @param size Number of bytes in 'data'. |
137 | | /// @return false on short or failed write (treat the connection as lost). |
138 | 288k | bool MockConnection::WriteAll(const unsigned char* data, std::size_t size) { |
139 | 288k | if (server_fd_ < 0) { |
140 | 0 | return false; |
141 | 0 | } |
142 | 288k | std::size_t written = 0; |
143 | 574k | while (written < size) { |
144 | 288k | ssize_t n = ::send(server_fd_, data + written, size - written, MSG_NOSIGNAL); |
145 | 288k | if (n <= 0) { |
146 | 1.69k | return false; |
147 | 1.69k | } |
148 | 286k | written += static_cast<std::size_t>(n); |
149 | 286k | } |
150 | 286k | return true; |
151 | 288k | } |
152 | | |
153 | | /// Drain bytes curl has written. When a backpressure drain limit has been |
154 | | /// applied (see ApplyBackpressure), stops after drain_limit_ bytes so the |
155 | | /// kernel recv buffer stays near-full and curl keeps seeing short writes. |
156 | | /// Otherwise drains until read() returns 0/EAGAIN, matching legacy behaviour. |
157 | | /// @return number of bytes consumed during this call. |
158 | 876k | std::size_t MockConnection::DrainIncoming() { |
159 | 876k | if (server_fd_ < 0) { |
160 | 0 | return 0; |
161 | 0 | } |
162 | 876k | unsigned char scratch[4096]; |
163 | 876k | std::size_t drained = 0; |
164 | 1.07M | while (drain_limit_ == 0 || drained < drain_limit_) { |
165 | 1.00M | std::size_t want = sizeof(scratch); |
166 | 1.00M | if (drain_limit_ != 0) { |
167 | 136k | const std::size_t remaining = drain_limit_ - drained; |
168 | 136k | if (remaining < want) { |
169 | 122k | want = remaining; |
170 | 122k | } |
171 | 136k | } |
172 | 1.00M | ssize_t n = ::read(server_fd_, scratch, want); |
173 | 1.00M | if (n <= 0) { |
174 | 806k | break; |
175 | 806k | } |
176 | 197k | drained += static_cast<std::size_t>(n); |
177 | 197k | } |
178 | 876k | return drained; |
179 | 876k | } |
180 | | |
181 | | /// Tighten both halves of the socketpair buffer and/or cap DrainIncoming's |
182 | | /// per-call byte budget. SO_RCVBUF on the server fd caps how much curl can |
183 | | /// push into the pipe; SO_SNDBUF on the client fd (which curl will soon own |
184 | | /// but hasn't yet, so we can still tune it) caps how much curl's send() can |
185 | | /// buffer before short-writing. Linux socketpairs effectively use max(SNDBUF, |
186 | | /// RCVBUF*2) as pipe capacity, so we need both to see short writes reliably. |
187 | | /// See header docs. |
188 | 141k | void MockConnection::ApplyBackpressure(int recv_buf_bytes, std::size_t drain_limit) { |
189 | 141k | if (recv_buf_bytes > 0) { |
190 | 14.8k | if (server_fd_ >= 0) { |
191 | 14.8k | (void)setsockopt(server_fd_, SOL_SOCKET, SO_RCVBUF, &recv_buf_bytes, sizeof(recv_buf_bytes)); |
192 | 14.8k | } |
193 | 14.8k | if (client_fd_ >= 0) { |
194 | 14.8k | (void)setsockopt(client_fd_, SOL_SOCKET, SO_SNDBUF, &recv_buf_bytes, sizeof(recv_buf_bytes)); |
195 | 14.8k | } |
196 | 14.8k | } |
197 | 141k | drain_limit_ = drain_limit; |
198 | 141k | } |
199 | | |
200 | | /// Non-blocking read: append whatever bytes are currently available on the |
201 | | /// server fd to 'out'. Used by the WS handshake path to collect curl's HTTP |
202 | | /// Upgrade request without losing any bytes. |
203 | | /// @param out Destination buffer; unchanged if no bytes are pending. |
204 | 127k | void MockConnection::ReadAvailable(std::string* out) { |
205 | 127k | if (server_fd_ < 0 || out == nullptr) { |
206 | 0 | return; |
207 | 0 | } |
208 | 127k | unsigned char scratch[4096]; |
209 | 175k | while (true) { |
210 | 175k | ssize_t n = ::read(server_fd_, scratch, sizeof(scratch)); |
211 | 175k | if (n <= 0) { |
212 | 127k | break; |
213 | 127k | } |
214 | 47.5k | out->append(reinterpret_cast<const char*>(scratch), static_cast<std::size_t>(n)); |
215 | 47.5k | } |
216 | 127k | } |
217 | | |
218 | | /// Signal end-of-response to libcurl by half-closing the write side. |
219 | 91.4k | void MockConnection::ShutdownWrite() { |
220 | 91.4k | if (server_fd_ < 0) { |
221 | 0 | return; |
222 | 0 | } |
223 | 91.4k | ::shutdown(server_fd_, SHUT_WR); |
224 | 91.4k | } |
225 | | |
226 | | /// @class proto_fuzzer::MockServer |
227 | | /// @brief Orchestrates a bounded sequence of mock HTTP exchanges: installs |
228 | | /// socket callbacks on an easy handle, assigns a response script to each new |
229 | | /// socket, and feeds queued chunks as libcurl reads them. |
230 | | |
231 | | /// Construct an idle MockServer with no scripted responses or open peers. |
232 | | /// DriveScenario() configures it from a Scenario before curl can open a socket. |
233 | | MockServer::MockServer() |
234 | 96.7k | : script_count_(0), |
235 | 96.7k | next_script_(0), |
236 | 96.7k | active_script_(nullptr), |
237 | 96.7k | preload_all_chunks_(false), |
238 | 96.7k | keep_connections_open_(false) {} |
239 | | |
240 | | /// Default destructor; current and previous MockConnections clean up their |
241 | | /// server-side descriptors only after curl has been removed from the multi. |
242 | 96.7k | MockServer::~MockServer() = default; |
243 | | |
244 | | /// Build the complete borrowed script table before entering curl. The primary |
245 | | /// Connection always occupies slot zero for backwards compatibility; only |
246 | | /// three subsequent pointers are retained so protobuf mutations cannot |
247 | | /// allocate socketpairs or prolong redirects in proportion to repeated-field |
248 | | /// size. The Scenario passed by ScenarioRunner outlives this synchronous drive, |
249 | | /// which makes borrowing safe while avoiding response-byte copies. |
250 | | /// @param scenario Source of the primary and bounded follow-on scripts. |
251 | 96.7k | void MockServer::SetScripts(const curl::fuzzer::proto::Scenario& scenario) { |
252 | 96.7k | ResetConnections(); |
253 | 96.7k | script_count_ = 0; |
254 | 96.7k | next_script_ = 0; |
255 | | |
256 | 129k | const auto append_script = [this](const curl::fuzzer::proto::Connection& connection) { |
257 | 129k | ConnectionScript& script = scripts_[script_count_++]; |
258 | 129k | script.connection = &connection; |
259 | 129k | script.raw_chunk_count = std::min<std::size_t>(scenario_limits::kMaxResponseChunks, connection.on_readable_size()); |
260 | 129k | const std::size_t frame_budget = scenario_limits::kMaxResponseChunks - script.raw_chunk_count; |
261 | 129k | script.frame_chunk_count = std::min<std::size_t>(frame_budget, connection.server_frames_size()); |
262 | 129k | script.next_chunk = 0; |
263 | 129k | }; |
264 | | |
265 | 96.7k | append_script(scenario.connection()); |
266 | 96.7k | const std::size_t subsequent_count = std::min<std::size_t>( |
267 | 96.7k | scenario_limits::kMaxConnections - 1, static_cast<std::size_t>(scenario.subsequent_connections_size())); |
268 | 129k | for (std::size_t i = 0; i < subsequent_count; ++i) { |
269 | 32.5k | append_script(scenario.subsequent_connections(static_cast<int>(i))); |
270 | 32.5k | } |
271 | 96.7k | } |
272 | | |
273 | | /// Configure whether a completed response leaves its socket reusable. |
274 | 502 | void MockServer::SetKeepConnectionsOpen(bool keep_open) { keep_connections_open_ = keep_open; } |
275 | | |
276 | | /// @return true if at least one on_readable chunk has not yet been sent. |
277 | 822k | bool MockServer::has_more_chunks() const { |
278 | 822k | return active_script_ != nullptr && active_script_->next_chunk < active_script_->chunk_count(); |
279 | 822k | } |
280 | | |
281 | | /// Keep plaintext transport construction behind a virtual boundary so the |
282 | | /// HTTPS lane can add TLS without copying the HTTP script state machine. |
283 | 85.0k | std::unique_ptr<MockConnection> MockServer::CreateConnection() { return std::make_unique<MockConnection>(); } |
284 | | |
285 | | /// Release all socketpairs while the dynamic transport type is still alive. |
286 | | /// This is separate from SetScripts because a derived destructor may need to |
287 | | /// enforce a stricter order than C++'s derived-member-before-base teardown. |
288 | 120k | void MockServer::ResetConnections() { |
289 | 120k | active_script_ = nullptr; |
290 | 120k | connection_.reset(); |
291 | 120k | previous_connections_.clear(); |
292 | 120k | } |
293 | | |
294 | | /// Plain HTTP has no connection-filter result that must be observed live. |
295 | | /// @param easy Active easy handle, unused by the plaintext transport. |
296 | 303k | void MockServer::ObserveActiveTransfer(CURL* /*easy*/) {} |
297 | | |
298 | | /// Called by the OPENSOCKETFUNCTION trampoline in the base class. Creates the |
299 | | /// MockConnection, writes initial_response into it, and returns the |
300 | | /// client-side fd to hand to libcurl. |
301 | | /// @param purpose Socket role requested by curl; ordinary stream mocks do not |
302 | | /// need to distinguish outbound HTTP roles. |
303 | | /// @param address Intended destination retained by curl; the socketpair is |
304 | | /// already connected and therefore leaves it untouched. |
305 | | /// @return the client-side fd to hand to libcurl, or CURL_SOCKET_BAD on |
306 | | /// failure. |
307 | 120k | curl_socket_t MockServer::HandleOpenSocket(curlsocktype purpose, struct curl_sockaddr* address) { |
308 | 120k | (void)purpose; |
309 | 120k | (void)address; |
310 | 120k | if (next_script_ >= script_count_) { |
311 | | // Refusing a fifth socket keeps redirect loops bounded even if curl's own |
312 | | // redirect limit is mutated upward or an authentication scheme retries. |
313 | 6.81k | return CURL_SOCKET_BAD; |
314 | 6.81k | } |
315 | | |
316 | 113k | if (connection_) { |
317 | | // libcurl owns the corresponding client fd, so retain the server half |
318 | | // until the easy handle leaves the multi instead of closing it at the |
319 | | // moment a redirect or retry opens its replacement. |
320 | 23.8k | previous_connections_.push_back(std::move(connection_)); |
321 | 23.8k | } |
322 | | |
323 | 113k | active_script_ = &scripts_[next_script_++]; |
324 | 113k | const curl::fuzzer::proto::Connection& script_connection = *active_script_->connection; |
325 | 113k | connection_ = CreateConnection(); |
326 | 113k | if (!connection_ || !connection_->ok()) { |
327 | 0 | connection_.reset(); |
328 | 0 | active_script_ = nullptr; |
329 | 0 | return CURL_SOCKET_BAD; |
330 | 0 | } |
331 | | |
332 | | // Target policies clamp/clear these values before execution. Saturating the |
333 | | // compatibility lane's raw uint32 avoids implementation-defined narrowing |
334 | | // while preserving its ability to request any representable socket size. |
335 | 113k | const std::uint32_t int_max = static_cast<std::uint32_t>(std::numeric_limits<int>::max()); |
336 | 113k | const auto& backpressure = script_connection.backpressure(); |
337 | 113k | const int recv_buf_bytes = static_cast<int>(std::min(backpressure.recv_buf_bytes(), int_max)); |
338 | 113k | connection_->ApplyBackpressure(recv_buf_bytes, static_cast<std::size_t>(backpressure.drain_limit())); |
339 | | |
340 | 113k | const std::string& initial_response = script_connection.initial_response(); |
341 | 113k | if (!initial_response.empty()) { |
342 | 63.3k | if (!connection_->WriteAll(reinterpret_cast<const unsigned char*>(initial_response.data()), |
343 | 63.3k | initial_response.size())) { |
344 | 0 | connection_.reset(); |
345 | 0 | active_script_ = nullptr; |
346 | 0 | return CURL_SOCKET_BAD; |
347 | 0 | } |
348 | 63.3k | } |
349 | 113k | if (preload_all_chunks_) { |
350 | 7.15k | while (has_more_chunks()) { |
351 | 5.78k | (void)DeliverNextChunk(); |
352 | 5.78k | } |
353 | 1.36k | } |
354 | 113k | if (active_script_->chunk_count() == 0 && !keep_connections_open_) { |
355 | 51.3k | connection_->ShutdownWrite(); |
356 | 51.3k | } |
357 | 113k | return connection_->take_client_fd(); |
358 | 113k | } |
359 | | |
360 | | /// Preload all bounded response bytes from inside OPENSOCKETFUNCTION, where |
361 | | /// the mock still owns both socketpair ends. The total serialized fuzz input |
362 | | /// is capped by libFuzzer's max_len. If a non-blocking preload fills the local |
363 | | /// socket, curl observes only the successfully queued prefix and the timeouts |
364 | | /// below bound the incomplete response. Unlike every other drive loop, this one |
365 | | /// has no iteration budget of its own, so it reasserts those timeouts after |
366 | | /// scenario setopts and clears the one option that can disable them. |
367 | | /// @param easy Configured easy handle whose open-socket callback targets this |
368 | | /// mock. |
369 | | /// @param scenario Bounded response script to preload before performing. |
370 | 945 | void MockServer::DriveEasyScenario(CURL* easy, const curl::fuzzer::proto::Scenario& scenario) { |
371 | 945 | SetScripts(scenario); |
372 | 945 | preload_all_chunks_ = true; |
373 | 945 | (void)curl_easy_setopt(easy, CURLOPT_TIMEOUT_MS, 50L); |
374 | 945 | (void)curl_easy_setopt(easy, CURLOPT_CONNECTTIMEOUT_MS, 50L); |
375 | | // CONNECT_ONLY makes curl's timeleft check report "no limit" for the whole |
376 | | // post-connect phase, so the timeouts above stop being enforced; value 2 |
377 | | // additionally skips the connect-only shortcut and runs a complete transfer. |
378 | | // Clearing it is not a coverage loss: the multi-driven WebSocket lanes still |
379 | | // exercise CONNECT_ONLY under their own iteration budget. Zero rather than |
380 | | // one because CONNECT_ONLY=1 also disables the timeout and stays bounded only |
381 | | // through a curl-internal shortcut this harness should not depend on. |
382 | 945 | (void)curl_easy_setopt(easy, CURLOPT_CONNECT_ONLY, 0L); |
383 | 945 | (void)curl_easy_perform(easy); |
384 | 945 | preload_all_chunks_ = false; |
385 | 945 | } |
386 | | |
387 | | /// Push the next queued chunk. Called by the drive loop when curl is ready |
388 | | /// for more data. |
389 | | /// @return true when a chunk was consumed from the script. |
390 | 226k | bool MockServer::DeliverNextChunk() { |
391 | 226k | if (!connection_ || !has_more_chunks()) { |
392 | 0 | return false; |
393 | 0 | } |
394 | 226k | connection_->DrainIncoming(); |
395 | 226k | const std::size_t chunk_index = active_script_->next_chunk++; |
396 | 226k | const auto& script_connection = *active_script_->connection; |
397 | 226k | if (chunk_index < active_script_->raw_chunk_count) { |
398 | 178k | const std::string& chunk = script_connection.on_readable(static_cast<int>(chunk_index)); |
399 | 178k | if (!chunk.empty()) { |
400 | 170k | connection_->WriteAll(reinterpret_cast<const unsigned char*>(chunk.data()), chunk.size()); |
401 | 170k | } |
402 | 178k | } else { |
403 | | // HTTP scenarios historically accept structured WebSocket frames as raw |
404 | | // response bytes after every on_readable chunk. Serialize only the frame |
405 | | // curl is about to receive: follow-on scripts and capped suffix frames may |
406 | | // never be consumed, so eagerly materialising all of them wastes mutations. |
407 | 47.6k | const std::size_t frame_index = chunk_index - active_script_->raw_chunk_count; |
408 | 47.6k | const std::string chunk = SerializeWebSocketFrame(script_connection.server_frames(static_cast<int>(frame_index))); |
409 | 47.6k | if (!chunk.empty()) { |
410 | 47.6k | connection_->WriteAll(reinterpret_cast<const unsigned char*>(chunk.data()), chunk.size()); |
411 | 47.6k | } |
412 | 47.6k | } |
413 | 226k | if (!has_more_chunks() && !keep_connections_open_) { |
414 | 49.3k | connection_->ShutdownWrite(); |
415 | 49.3k | } |
416 | 226k | return true; |
417 | 226k | } |
418 | | |
419 | | /// Drain every live mock peer. Previous connections are normally quiescent, |
420 | | /// but curl can finish sending a request body or close one after it has begun |
421 | | /// resolving/opening the redirect target. Servicing both sides makes that |
422 | | /// overlap deterministic without conflating their response scripts. |
423 | 363k | std::size_t MockServer::DrainIncomingConnections() { |
424 | 363k | std::size_t drained = 0; |
425 | 363k | for (const auto& previous : previous_connections_) { |
426 | 130k | drained += previous->DrainIncoming(); |
427 | 130k | } |
428 | 363k | if (connection_) { |
429 | 363k | drained += connection_->DrainIncoming(); |
430 | 363k | } |
431 | 363k | return drained; |
432 | 363k | } |
433 | | |
434 | | /// Keep peer servicing identical across the perform and socket-action APIs. |
435 | | /// Draining first creates request-side space before a response can provoke |
436 | | /// another write, and releasing only one chunk preserves the scenario's |
437 | | /// mutation-controlled parser boundaries. |
438 | 363k | bool MockServer::ServiceConnections() { |
439 | 363k | bool made_progress = DrainIncomingConnections() != 0; |
440 | 363k | if (has_more_chunks()) { |
441 | 220k | made_progress = DeliverNextChunk() || made_progress; |
442 | 220k | } |
443 | 363k | return made_progress; |
444 | 363k | } |
445 | | |
446 | 502 | std::size_t MockServer::opened_connection_count() const { return next_script_; } |
447 | | |
448 | | /// Run a zero-wait application event loop around curl_multi_socket_action. |
449 | | /// Positive timers are deliberately not slept: the API lane clears timing |
450 | | /// controls and exists to cover event-driven dispatch, while the dedicated |
451 | | /// timing lane remains responsible for clock-dependent behavior. |
452 | 873 | void MockServer::RunSocketActionLoop(CURLM* multi, CURL* easy) { |
453 | 873 | MultiSocketDriver* driver = multi_socket_driver(); |
454 | 873 | if (driver == nullptr) { |
455 | 0 | return; |
456 | 0 | } |
457 | | |
458 | 873 | int still_running = 1; |
459 | 873 | CURLMcode rc = driver->Start(&still_running); |
460 | 873 | ObserveActiveTransfer(easy); |
461 | 873 | int idle_iterations = 0; |
462 | 873 | int drive_iterations = 0; |
463 | 6.09k | while (rc == CURLM_OK && still_running && idle_iterations < kMaxIdleIterations && |
464 | 5.22k | drive_iterations++ < kMaxDriveIterations) { |
465 | 5.22k | bool made_progress = ServiceConnections(); |
466 | 5.22k | const MultiSocketDriver::DriveResult drive_result = driver->DriveReady(&still_running); |
467 | 5.22k | rc = drive_result.code; |
468 | 5.22k | ObserveActiveTransfer(easy); |
469 | 5.22k | made_progress = drive_result.made_progress || made_progress; |
470 | 5.22k | if (made_progress) { |
471 | 5.15k | idle_iterations = 0; |
472 | 5.15k | } else { |
473 | 72 | ++idle_iterations; |
474 | 72 | } |
475 | 5.22k | } |
476 | 873 | } |
477 | | |
478 | | /// Seed the mock from the scenario, then drive the perform loop until curl is |
479 | | /// done or a deterministic operation/idle budget is hit. Ordinary scenarios |
480 | | /// never wait; explicit backpressure scenarios may use short select() waits. |
481 | | /// @param multi caller-owned multi; 'easy' is already added. |
482 | | /// @param easy the curl easy handle attached to this mock. |
483 | | /// @param scenario source of the initial_response and on_readable chunks. |
484 | 95.3k | void MockServer::RunLoop(CURLM* multi, CURL* easy, const curl::fuzzer::proto::Scenario& scenario) { |
485 | 95.3k | SetScripts(scenario); |
486 | | |
487 | 95.3k | if (multi_socket_driver() != nullptr) { |
488 | 873 | RunSocketActionLoop(multi, easy); |
489 | 873 | return; |
490 | 873 | } |
491 | | |
492 | 94.4k | int still_running = 1; |
493 | 94.4k | int idle_iterations = 0; |
494 | 94.4k | int drive_iterations = 0; |
495 | 94.4k | bool pollset_probed = false; |
496 | | // Each newly opened socket can carry its own pressure settings. Inspect only |
497 | | // the borrowed, runtime-bounded scripts: ignored repeated protobuf entries must |
498 | | // not opt a fast transfer into timed waits. |
499 | 94.4k | const bool timed_drive = |
500 | 118k | std::any_of(scripts_.begin(), scripts_.begin() + script_count_, [](const ConnectionScript& script) { |
501 | 118k | const auto& backpressure = script.connection->backpressure(); |
502 | 118k | return backpressure.recv_buf_bytes() != 0 || backpressure.drain_limit() != 0; |
503 | 118k | }); |
504 | 94.4k | const int idle_limit = timed_drive ? kMaxTimedIdleIterations : kMaxIdleIterations; |
505 | 94.4k | CURLMcode rc = CURLM_OK; |
506 | | |
507 | 442k | while (still_running && idle_iterations < idle_limit && drive_iterations++ < kMaxDriveIterations) { |
508 | 438k | bool made_progress = false; |
509 | 438k | const int running_before = still_running; |
510 | 438k | rc = curl_multi_perform(multi, &still_running); |
511 | 438k | if (rc != CURLM_OK) { |
512 | 0 | break; |
513 | 0 | } |
514 | 438k | ObserveActiveTransfer(easy); |
515 | 438k | made_progress = still_running != running_before; |
516 | 438k | if (timed_drive && still_running && !pollset_probed) { |
517 | | // Probe only after curl has built the socket/filter chain. A zero-timeout |
518 | | // poll adds no wall-clock wait; restricting it to the timing lane keeps |
519 | | // the fixed HTTP target free of a per-input polling syscall. |
520 | 5.55k | ProbeMultiPollset(multi); |
521 | 5.55k | pollset_probed = true; |
522 | 5.55k | } |
523 | 438k | if (!still_running) { |
524 | 89.7k | break; |
525 | 89.7k | } |
526 | | |
527 | | // Always drain whatever curl has written. Under backpressure the kernel |
528 | | // recv buffer would otherwise stay full — curl short-writes, the mock |
529 | | // never consumes, and the transfer wedges until the drive budget. With |
530 | | // drain_limit set this still honours the per-tick byte budget. |
531 | 348k | made_progress = ServiceConnections() || made_progress; |
532 | | |
533 | 348k | if (made_progress) { |
534 | 286k | idle_iterations = 0; |
535 | 286k | continue; |
536 | 286k | } |
537 | | |
538 | | // Ordinary socketpair scenarios never sleep: repeated no-progress |
539 | | // performs are enough to settle curl's local state machine. Only an |
540 | | // explicit BackpressureConfig opts into short waits so timeout-related |
541 | | // behaviour remains fuzzable without taxing every corpus entry. |
542 | 61.5k | if (timed_drive) { |
543 | 24.0k | (void)WaitOnMultiFdset(multi, &rc); |
544 | 24.0k | if (rc != CURLM_OK) { |
545 | 0 | break; |
546 | 0 | } |
547 | 24.0k | } |
548 | 61.5k | ++idle_iterations; |
549 | 61.5k | } |
550 | 94.4k | } |
551 | | |
552 | | } // namespace proto_fuzzer |