/src/curl_fuzzer/proto_fuzzer/multi_transfer_runner.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 the bounded concurrent easy-handle runner. |
9 | | |
10 | | #include "proto_fuzzer/multi_transfer_runner.h" |
11 | | |
12 | | #define CURL_ALLOW_OLD_MULTI_SOCKET |
13 | | #include <curl/curl.h> |
14 | | #include <curl/multi.h> |
15 | | |
16 | | #include <algorithm> |
17 | | #include <array> |
18 | | #include <cstddef> |
19 | | #include <memory> |
20 | | #include <string> |
21 | | |
22 | | #include "proto_fuzzer/curl_raii.h" |
23 | | #include "proto_fuzzer/mock_server.h" |
24 | | #include "proto_fuzzer/multi_socket_driver.h" |
25 | | #include "proto_fuzzer/option_apply.h" |
26 | | #include "proto_fuzzer/request_data.h" |
27 | | #include "proto_fuzzer/scenario_limits.h" |
28 | | |
29 | | namespace proto_fuzzer { |
30 | | |
31 | | namespace { |
32 | | |
33 | | /// One easy handle and every caller-owned pointer installed on it. Declaration |
34 | | /// order makes reverse destruction detach request data, clean the easy, then |
35 | | /// release CONNECT_TO storage if an early return bypasses explicit teardown. |
36 | | struct TransferState { |
37 | | CurlSlistPtr connect_to; |
38 | | CurlEasyPtr easy; |
39 | | std::unique_ptr<ScenarioRequestData> request_data; |
40 | | bool attached = false; |
41 | | }; |
42 | | |
43 | | constexpr int kMaxDriveIterations = 512; |
44 | | constexpr int kMaxIdleIterations = 8; |
45 | | |
46 | 19.9k | std::size_t RuntimeTransferCount(const curl::fuzzer::proto::MultiPlan& plan) { |
47 | 19.9k | return std::max(scenario_limits::kMinMultiTransfers, |
48 | 19.9k | std::min(static_cast<std::size_t>(plan.transfer_count()), scenario_limits::kMaxMultiTransfers)); |
49 | 19.9k | } |
50 | | |
51 | 84.5k | void RecordNotification(CURLM*, unsigned int notification, CURL*, void* userdata) { |
52 | 84.5k | auto* stats = static_cast<MultiTransferRunStats*>(userdata); |
53 | 84.5k | if (stats == nullptr) { |
54 | 0 | return; |
55 | 0 | } |
56 | 84.5k | if (notification == CURLMNOTIFY_INFO_READ) { |
57 | 26.8k | ++stats->info_read_notifications; |
58 | 57.7k | } else if (notification == CURLMNOTIFY_EASY_DONE) { |
59 | 57.7k | ++stats->easy_done_notifications; |
60 | 57.7k | } |
61 | 84.5k | } |
62 | | |
63 | | /// Consume every currently queued completion before handle removal can discard |
64 | | /// it. Re-added handles may complete more than once, so count messages rather |
65 | | /// than only distinct easy pointers. |
66 | | bool DrainCompletionMessages(CURLM* multi, std::array<TransferState, scenario_limits::kMaxMultiTransfers>* transfers, |
67 | 211k | std::size_t transfer_count, MultiTransferRunStats* stats) { |
68 | 211k | bool consumed = false; |
69 | 211k | int messages_remaining = 0; |
70 | 211k | CURLMsg* message = nullptr; |
71 | 268k | while ((message = curl_multi_info_read(multi, &messages_remaining)) != nullptr) { |
72 | 57.5k | consumed = true; |
73 | 57.5k | if (message->msg != CURLMSG_DONE) { |
74 | 0 | continue; |
75 | 0 | } |
76 | 121k | for (std::size_t index = 0; index < transfer_count; ++index) { |
77 | 121k | TransferState& transfer = (*transfers)[index]; |
78 | 121k | if (transfer.easy.get() == message->easy_handle) { |
79 | 57.5k | ++stats->completion_messages; |
80 | 57.5k | break; |
81 | 57.5k | } |
82 | 121k | } |
83 | 57.5k | } |
84 | 211k | return consumed; |
85 | 211k | } |
86 | | |
87 | | /// Apply one schema-safe lifecycle transition. Return true only when libcurl |
88 | | /// accepted a state change; consuming an ineffective action is tracked |
89 | | /// separately so later ordered actions remain reachable. |
90 | | bool ApplyAction(const curl::fuzzer::proto::MultiAction& action, CURLM* multi, |
91 | | std::array<TransferState, scenario_limits::kMaxMultiTransfers>* transfers, |
92 | 28.5k | std::size_t transfer_count) { |
93 | 28.5k | if (transfer_count == 0) { |
94 | 0 | return false; |
95 | 0 | } |
96 | 28.5k | TransferState& transfer = (*transfers)[static_cast<std::size_t>(action.transfer_selector()) % transfer_count]; |
97 | 28.5k | CURL* easy = transfer.easy.get(); |
98 | 28.5k | if (easy == nullptr) { |
99 | 0 | return false; |
100 | 0 | } |
101 | | |
102 | 28.5k | switch (action.kind()) { |
103 | 1.31k | case curl::fuzzer::proto::MULTI_ACTION_PAUSE_RECV: |
104 | 1.31k | return curl_easy_pause(easy, CURLPAUSE_RECV) == CURLE_OK; |
105 | 949 | case curl::fuzzer::proto::MULTI_ACTION_PAUSE_SEND: |
106 | 949 | return curl_easy_pause(easy, CURLPAUSE_SEND) == CURLE_OK; |
107 | 1.01k | case curl::fuzzer::proto::MULTI_ACTION_PAUSE_ALL: |
108 | 1.01k | return curl_easy_pause(easy, CURLPAUSE_ALL) == CURLE_OK; |
109 | 694 | case curl::fuzzer::proto::MULTI_ACTION_RESUME: |
110 | 694 | return curl_easy_pause(easy, CURLPAUSE_CONT) == CURLE_OK; |
111 | 7.56k | case curl::fuzzer::proto::MULTI_ACTION_REMOVE: |
112 | 7.56k | if (transfer.attached && curl_multi_remove_handle(multi, easy) == CURLM_OK) { |
113 | 6.92k | transfer.attached = false; |
114 | 6.92k | return true; |
115 | 6.92k | } |
116 | 642 | return false; |
117 | 7.93k | case curl::fuzzer::proto::MULTI_ACTION_READD: |
118 | 7.93k | if (!transfer.attached && curl_multi_add_handle(multi, easy) == CURLM_OK) { |
119 | 6.16k | transfer.attached = true; |
120 | 6.16k | return true; |
121 | 6.16k | } |
122 | 1.77k | return false; |
123 | 9.09k | case curl::fuzzer::proto::MULTI_ACTION_NONE: |
124 | 9.09k | default: |
125 | 9.09k | return false; |
126 | 28.5k | } |
127 | 28.5k | } |
128 | | |
129 | | /// Cover the public wait/poll/fdset queries from a valid shared-multi state |
130 | | /// without introducing wall-clock delay. |
131 | 19.9k | void ProbeMultiWaitApis(CURLM* multi) { |
132 | 19.9k | int numfds = 0; |
133 | 19.9k | (void)curl_multi_poll(multi, nullptr, 0, 0, &numfds); |
134 | 19.9k | (void)curl_multi_wait(multi, nullptr, 0, 0, &numfds); |
135 | | |
136 | 19.9k | unsigned int waitfd_count = 0; |
137 | 19.9k | (void)curl_multi_waitfds(multi, nullptr, 0, &waitfd_count); |
138 | 19.9k | std::array<curl_waitfd, scenario_limits::kMaxMultiTransfers * 2> waitfds{}; |
139 | 19.9k | (void)curl_multi_waitfds(multi, waitfds.data(), static_cast<unsigned int>(waitfds.size()), &waitfd_count); |
140 | | |
141 | 19.9k | fd_set readfds; |
142 | 19.9k | fd_set writefds; |
143 | 19.9k | fd_set exceptfds; |
144 | 19.9k | FD_ZERO(&readfds); |
145 | 19.9k | FD_ZERO(&writefds); |
146 | 19.9k | FD_ZERO(&exceptfds); |
147 | 19.9k | int maxfd = -1; |
148 | 19.9k | (void)curl_multi_fdset(multi, &readfds, &writefds, &exceptfds, &maxfd); |
149 | 19.9k | long timeout_ms = -1; |
150 | 19.9k | (void)curl_multi_timeout(multi, &timeout_ms); |
151 | | |
152 | 19.9k | CURL** handles = curl_multi_get_handles(multi); |
153 | 19.9k | curl_free(handles); |
154 | | |
155 | 19.9k | constexpr CURLMinfo_offt kOffsetInfo[] = { |
156 | 19.9k | CURLMINFO_XFERS_CURRENT, CURLMINFO_XFERS_RUNNING, CURLMINFO_XFERS_PENDING, |
157 | 19.9k | CURLMINFO_XFERS_DONE, CURLMINFO_XFERS_ADDED, |
158 | 19.9k | }; |
159 | 99.8k | for (const CURLMinfo_offt info : kOffsetInfo) { |
160 | 99.8k | curl_off_t value = 0; |
161 | 99.8k | (void)curl_multi_get_offt(multi, info, &value); |
162 | 99.8k | } |
163 | 19.9k | } |
164 | | |
165 | | /// Exercise the two deprecated public entrypoints as real symbols rather than |
166 | | /// curl_multi_socket's compatibility macro. Their output is deliberately not |
167 | | /// used to select the primary runner: either call may advance a transfer, and |
168 | | /// the normal bounded loop consumes whatever state remains. |
169 | 19.9k | void ProbeLegacyMultiSocketApis(CURLM* multi) { |
170 | 19.9k | int running_handles = 0; |
171 | 19.9k | #if defined(__GNUC__) |
172 | 19.9k | #pragma GCC diagnostic push |
173 | 19.9k | #pragma GCC diagnostic ignored "-Wdeprecated-declarations" |
174 | 19.9k | #endif |
175 | 19.9k | (void)curl_multi_socket(multi, CURL_SOCKET_TIMEOUT, &running_handles); |
176 | 19.9k | (void)curl_multi_socket_all(multi, &running_handles); |
177 | 19.9k | #if defined(__GNUC__) |
178 | 19.9k | #pragma GCC diagnostic pop |
179 | 19.9k | #endif |
180 | 19.9k | } |
181 | | |
182 | | } // namespace |
183 | | |
184 | 19.9k | MultiTransferRunStats RunMultiTransferScenario(const curl::fuzzer::proto::Scenario& scenario) { |
185 | 19.9k | MultiTransferRunStats stats; |
186 | 19.9k | const auto& plan = scenario.multi_plan(); |
187 | 19.9k | const std::size_t transfer_count = RuntimeTransferCount(plan); |
188 | 19.9k | stats.configured_handles = transfer_count; |
189 | | |
190 | | // The peer and socket callback state must outlive multi cleanup, which can |
191 | | // synchronously announce socket removals and close cached connections. |
192 | 19.9k | MockServer mock; |
193 | 19.9k | mock.SetScripts(scenario); |
194 | 19.9k | mock.SetKeepConnectionsOpen(plan.keep_connections_open()); |
195 | 19.9k | MultiSocketDriver socket_driver; |
196 | 19.9k | CurlMultiPtr multi(curl_multi_init()); |
197 | 19.9k | if (!multi) { |
198 | 0 | return stats; |
199 | 0 | } |
200 | | |
201 | 19.9k | const long max_host_connections = |
202 | 19.9k | static_cast<long>(std::min<std::size_t>(plan.max_host_connections(), transfer_count)); |
203 | 19.9k | const long max_total_connections = |
204 | 19.9k | static_cast<long>(std::min<std::size_t>(plan.max_total_connections(), transfer_count)); |
205 | 19.9k | const long cache_size = |
206 | 19.9k | static_cast<long>(std::min<std::size_t>(plan.connection_cache_size(), scenario_limits::kMaxMultiTransfers * 2)); |
207 | 19.9k | (void)curl_multi_setopt(multi.get(), CURLMOPT_MAX_HOST_CONNECTIONS, max_host_connections); |
208 | 19.9k | (void)curl_multi_setopt(multi.get(), CURLMOPT_MAX_TOTAL_CONNECTIONS, max_total_connections); |
209 | 19.9k | (void)curl_multi_setopt(multi.get(), CURLMOPT_MAXCONNECTS, cache_size); |
210 | 19.9k | (void)curl_multi_setopt(multi.get(), CURLMOPT_PIPELINING, |
211 | 19.9k | plan.multiplex() ? static_cast<long>(CURLPIPE_MULTIPLEX) : 0L); |
212 | 19.9k | (void)curl_multi_setopt(multi.get(), CURLMOPT_NOTIFYFUNCTION, &RecordNotification); |
213 | 19.9k | (void)curl_multi_setopt(multi.get(), CURLMOPT_NOTIFYDATA, &stats); |
214 | 19.9k | (void)curl_multi_notify_enable(multi.get(), CURLMNOTIFY_INFO_READ); |
215 | 19.9k | (void)curl_multi_notify_enable(multi.get(), CURLMNOTIFY_EASY_DONE); |
216 | 19.9k | (void)curl_multi_notify_disable(multi.get(), CURLMNOTIFY_INFO_READ); |
217 | 19.9k | (void)curl_multi_notify_enable(multi.get(), CURLMNOTIFY_INFO_READ); |
218 | 19.9k | (void)curl_multi_notify_disable(multi.get(), CURLMNOTIFY_EASY_DONE); |
219 | 19.9k | (void)curl_multi_notify_enable(multi.get(), CURLMNOTIFY_EASY_DONE); |
220 | | |
221 | 19.9k | const bool socket_mode = plan.drive_mode() == curl::fuzzer::proto::MULTI_DRIVE_SOCKET; |
222 | 19.9k | const bool socket_driver_installed = socket_mode && socket_driver.Install(multi.get()); |
223 | 19.9k | std::array<TransferState, scenario_limits::kMaxMultiTransfers> transfers; |
224 | 19.9k | const std::string url = "http://" + scenario.host_path(); |
225 | 80.9k | for (std::size_t index = 0; index < transfer_count; ++index) { |
226 | 60.9k | TransferState& transfer = transfers[index]; |
227 | 60.9k | transfer.easy.reset(curl_easy_init()); |
228 | 60.9k | if (!transfer.easy) { |
229 | 0 | continue; |
230 | 0 | } |
231 | 60.9k | transfer.connect_to.reset( |
232 | 60.9k | ApplyBaselineOptions(transfer.easy.get(), curl::fuzzer::proto::SCHEME_HTTP, scenario.trace_ids())); |
233 | 60.9k | (void)curl_easy_setopt(transfer.easy.get(), CURLOPT_URL, url.c_str()); |
234 | 60.9k | mock.Install(transfer.easy.get()); |
235 | 60.9k | (void)ApplyScenarioOptions(transfer.easy.get(), scenario); |
236 | 60.9k | transfer.request_data = std::make_unique<ScenarioRequestData>(transfer.easy.get(), scenario); |
237 | 60.9k | mock.ConfigureRequestData(transfer.request_data.get()); |
238 | 60.9k | if (curl_multi_add_handle(multi.get(), transfer.easy.get()) == CURLM_OK) { |
239 | 60.9k | transfer.attached = true; |
240 | 60.9k | ++stats.added_handles; |
241 | 60.9k | } |
242 | 60.9k | } |
243 | | |
244 | 19.9k | if (stats.added_handles == 0) { |
245 | 0 | multi.reset(); |
246 | 0 | return stats; |
247 | 0 | } |
248 | | |
249 | 19.9k | ProbeLegacyMultiSocketApis(multi.get()); |
250 | | |
251 | 19.9k | const std::size_t action_count = std::min<std::size_t>(plan.actions_size(), scenario_limits::kMaxMultiActions); |
252 | 19.9k | std::size_t next_action = 0; |
253 | 19.9k | bool initial_action_changed = false; |
254 | 19.9k | if (next_action < action_count) { |
255 | 3.42k | initial_action_changed = |
256 | 3.42k | ApplyAction(plan.actions(static_cast<int>(next_action++)), multi.get(), &transfers, transfer_count); |
257 | 3.42k | ++stats.actions_consumed; |
258 | 3.42k | } |
259 | | |
260 | 19.9k | int running_handles = 0; |
261 | 19.9k | CURLMcode result = socket_driver_installed ? socket_driver.Start(&running_handles) |
262 | 19.9k | : curl_multi_perform(multi.get(), &running_handles); |
263 | 19.9k | bool completion_progress = DrainCompletionMessages(multi.get(), &transfers, transfer_count, &stats); |
264 | 19.9k | ProbeMultiWaitApis(multi.get()); |
265 | 19.9k | if (plan.wake_multi()) { |
266 | 330 | if (socket_driver_installed) { |
267 | 92 | socket_driver.ProbeControlApis(); |
268 | 238 | } else { |
269 | 238 | (void)curl_multi_wakeup(multi.get()); |
270 | 238 | } |
271 | 330 | } |
272 | | |
273 | 19.9k | int idle_iterations = initial_action_changed || completion_progress ? 0 : 1; |
274 | 189k | for (int iteration = 0; result == CURLM_OK && iteration < kMaxDriveIterations; ++iteration) { |
275 | 188k | if (running_handles == 0 && next_action >= action_count) { |
276 | 17.7k | break; |
277 | 17.7k | } |
278 | | |
279 | 171k | bool made_progress = mock.ServiceConnections(); |
280 | 171k | if (next_action < action_count) { |
281 | 25.1k | made_progress = |
282 | 25.1k | ApplyAction(plan.actions(static_cast<int>(next_action++)), multi.get(), &transfers, transfer_count) || |
283 | 10.5k | made_progress; |
284 | 25.1k | ++stats.actions_consumed; |
285 | | // Consuming an ordered no-op is still harness progress: continue to the |
286 | | // next action instead of letting an idle transport hide its suffix. |
287 | 25.1k | made_progress = true; |
288 | 25.1k | } |
289 | | |
290 | 171k | const int running_before = running_handles; |
291 | 171k | if (socket_driver_installed) { |
292 | 121k | const MultiSocketDriver::DriveResult drive = socket_driver.DriveReady(&running_handles); |
293 | 121k | result = drive.code; |
294 | 121k | made_progress = drive.made_progress || made_progress; |
295 | 121k | } else { |
296 | 49.6k | result = curl_multi_perform(multi.get(), &running_handles); |
297 | 49.6k | made_progress = running_handles != running_before || made_progress; |
298 | 49.6k | } |
299 | 171k | made_progress = DrainCompletionMessages(multi.get(), &transfers, transfer_count, &stats) || made_progress; |
300 | | |
301 | 171k | if (made_progress) { |
302 | 154k | idle_iterations = 0; |
303 | 154k | } else if (++idle_iterations >= kMaxIdleIterations) { |
304 | 2.06k | break; |
305 | 2.06k | } |
306 | 171k | } |
307 | | |
308 | 19.9k | (void)DrainCompletionMessages(multi.get(), &transfers, transfer_count, &stats); |
309 | 19.9k | stats.opened_connections = mock.opened_connection_count(); |
310 | 80.9k | for (std::size_t index = 0; index < transfer_count; ++index) { |
311 | 60.9k | TransferState& transfer = transfers[index]; |
312 | 60.9k | if (!transfer.easy) { |
313 | 0 | continue; |
314 | 0 | } |
315 | 60.9k | (void)curl_easy_pause(transfer.easy.get(), CURLPAUSE_CONT); |
316 | 60.9k | if (transfer.attached) { |
317 | 60.1k | (void)curl_multi_remove_handle(multi.get(), transfer.easy.get()); |
318 | 60.1k | transfer.attached = false; |
319 | 60.1k | } |
320 | 60.9k | } |
321 | | |
322 | | // Destroy the multi while both its socket callback storage and mock peer are |
323 | | // live, then detach request-data pointers before cleaning each easy handle. |
324 | 19.9k | multi.reset(); |
325 | 80.9k | for (std::size_t index = 0; index < transfer_count; ++index) { |
326 | 60.9k | transfers[index].request_data.reset(); |
327 | 60.9k | transfers[index].easy.reset(); |
328 | 60.9k | transfers[index].connect_to.reset(); |
329 | 60.9k | } |
330 | 19.9k | return stats; |
331 | 19.9k | } |
332 | | |
333 | | } // namespace proto_fuzzer |