/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 | | #include <curl/curl.h> |
13 | | |
14 | | #include <algorithm> |
15 | | #include <array> |
16 | | #include <cstddef> |
17 | | #include <memory> |
18 | | #include <string> |
19 | | |
20 | | #include "proto_fuzzer/curl_raii.h" |
21 | | #include "proto_fuzzer/mock_server.h" |
22 | | #include "proto_fuzzer/multi_socket_driver.h" |
23 | | #include "proto_fuzzer/option_apply.h" |
24 | | #include "proto_fuzzer/request_data.h" |
25 | | #include "proto_fuzzer/scenario_limits.h" |
26 | | |
27 | | namespace proto_fuzzer { |
28 | | |
29 | | namespace { |
30 | | |
31 | | /// One easy handle and every caller-owned pointer installed on it. Declaration |
32 | | /// order makes reverse destruction detach request data, clean the easy, then |
33 | | /// release CONNECT_TO storage if an early return bypasses explicit teardown. |
34 | | struct TransferState { |
35 | | CurlSlistPtr connect_to; |
36 | | CurlEasyPtr easy; |
37 | | std::unique_ptr<ScenarioRequestData> request_data; |
38 | | bool attached = false; |
39 | | }; |
40 | | |
41 | | constexpr int kMaxDriveIterations = 512; |
42 | | constexpr int kMaxIdleIterations = 8; |
43 | | |
44 | 502 | std::size_t RuntimeTransferCount(const curl::fuzzer::proto::MultiPlan& plan) { |
45 | 502 | return std::max(scenario_limits::kMinMultiTransfers, |
46 | 502 | std::min(static_cast<std::size_t>(plan.transfer_count()), scenario_limits::kMaxMultiTransfers)); |
47 | 502 | } |
48 | | |
49 | | /// Consume every currently queued completion before handle removal can discard |
50 | | /// it. Re-added handles may complete more than once, so count messages rather |
51 | | /// than only distinct easy pointers. |
52 | | bool DrainCompletionMessages(CURLM* multi, std::array<TransferState, scenario_limits::kMaxMultiTransfers>* transfers, |
53 | 10.9k | std::size_t transfer_count, MultiTransferRunStats* stats) { |
54 | 10.9k | bool consumed = false; |
55 | 10.9k | int messages_remaining = 0; |
56 | 10.9k | CURLMsg* message = nullptr; |
57 | 12.3k | while ((message = curl_multi_info_read(multi, &messages_remaining)) != nullptr) { |
58 | 1.32k | consumed = true; |
59 | 1.32k | if (message->msg != CURLMSG_DONE) { |
60 | 0 | continue; |
61 | 0 | } |
62 | 3.00k | for (std::size_t index = 0; index < transfer_count; ++index) { |
63 | 3.00k | TransferState& transfer = (*transfers)[index]; |
64 | 3.00k | if (transfer.easy.get() == message->easy_handle) { |
65 | 1.32k | ++stats->completion_messages; |
66 | 1.32k | break; |
67 | 1.32k | } |
68 | 3.00k | } |
69 | 1.32k | } |
70 | 10.9k | return consumed; |
71 | 10.9k | } |
72 | | |
73 | | /// Apply one schema-safe lifecycle transition. Return true only when libcurl |
74 | | /// accepted a state change; consuming an ineffective action is tracked |
75 | | /// separately so later ordered actions remain reachable. |
76 | | bool ApplyAction(const curl::fuzzer::proto::MultiAction& action, CURLM* multi, |
77 | | std::array<TransferState, scenario_limits::kMaxMultiTransfers>* transfers, |
78 | 1.42k | std::size_t transfer_count) { |
79 | 1.42k | if (transfer_count == 0) { |
80 | 0 | return false; |
81 | 0 | } |
82 | 1.42k | TransferState& transfer = (*transfers)[static_cast<std::size_t>(action.transfer_selector()) % transfer_count]; |
83 | 1.42k | CURL* easy = transfer.easy.get(); |
84 | 1.42k | if (easy == nullptr) { |
85 | 0 | return false; |
86 | 0 | } |
87 | | |
88 | 1.42k | switch (action.kind()) { |
89 | 149 | case curl::fuzzer::proto::MULTI_ACTION_PAUSE_RECV: |
90 | 149 | return curl_easy_pause(easy, CURLPAUSE_RECV) == CURLE_OK; |
91 | 103 | case curl::fuzzer::proto::MULTI_ACTION_PAUSE_SEND: |
92 | 103 | return curl_easy_pause(easy, CURLPAUSE_SEND) == CURLE_OK; |
93 | 225 | case curl::fuzzer::proto::MULTI_ACTION_PAUSE_ALL: |
94 | 225 | return curl_easy_pause(easy, CURLPAUSE_ALL) == CURLE_OK; |
95 | 126 | case curl::fuzzer::proto::MULTI_ACTION_RESUME: |
96 | 126 | return curl_easy_pause(easy, CURLPAUSE_CONT) == CURLE_OK; |
97 | 357 | case curl::fuzzer::proto::MULTI_ACTION_REMOVE: |
98 | 357 | if (transfer.attached && curl_multi_remove_handle(multi, easy) == CURLM_OK) { |
99 | 280 | transfer.attached = false; |
100 | 280 | return true; |
101 | 280 | } |
102 | 77 | return false; |
103 | 245 | case curl::fuzzer::proto::MULTI_ACTION_READD: |
104 | 245 | if (!transfer.attached && curl_multi_add_handle(multi, easy) == CURLM_OK) { |
105 | 134 | transfer.attached = true; |
106 | 134 | return true; |
107 | 134 | } |
108 | 111 | return false; |
109 | 220 | case curl::fuzzer::proto::MULTI_ACTION_NONE: |
110 | 220 | default: |
111 | 220 | return false; |
112 | 1.42k | } |
113 | 1.42k | } |
114 | | |
115 | | /// Cover the public wait/poll/fdset queries from a valid shared-multi state |
116 | | /// without introducing wall-clock delay. |
117 | 502 | void ProbeMultiWaitApis(CURLM* multi) { |
118 | 502 | int numfds = 0; |
119 | 502 | (void)curl_multi_poll(multi, nullptr, 0, 0, &numfds); |
120 | 502 | (void)curl_multi_wait(multi, nullptr, 0, 0, &numfds); |
121 | | |
122 | 502 | fd_set readfds; |
123 | 502 | fd_set writefds; |
124 | 502 | fd_set exceptfds; |
125 | 502 | FD_ZERO(&readfds); |
126 | 502 | FD_ZERO(&writefds); |
127 | 502 | FD_ZERO(&exceptfds); |
128 | 502 | int maxfd = -1; |
129 | 502 | (void)curl_multi_fdset(multi, &readfds, &writefds, &exceptfds, &maxfd); |
130 | 502 | long timeout_ms = -1; |
131 | 502 | (void)curl_multi_timeout(multi, &timeout_ms); |
132 | 502 | } |
133 | | |
134 | | } // namespace |
135 | | |
136 | | /// Construct a stateless runner; all ownership is scoped to Run(). |
137 | 502 | MultiTransferRunner::MultiTransferRunner() = default; |
138 | | |
139 | | /// Default destructor; Run() dismantles each shared-multi lifecycle in place. |
140 | 502 | MultiTransferRunner::~MultiTransferRunner() = default; |
141 | | |
142 | 502 | MultiTransferRunStats MultiTransferRunner::Run(const curl::fuzzer::proto::Scenario& scenario) { |
143 | 502 | MultiTransferRunStats stats; |
144 | 502 | const auto& plan = scenario.multi_plan(); |
145 | 502 | const std::size_t transfer_count = RuntimeTransferCount(plan); |
146 | 502 | stats.configured_handles = transfer_count; |
147 | | |
148 | | // The peer and socket callback state must outlive multi cleanup, which can |
149 | | // synchronously announce socket removals and close cached connections. |
150 | 502 | MockServer mock; |
151 | 502 | mock.SetScripts(scenario); |
152 | 502 | mock.SetKeepConnectionsOpen(plan.keep_connections_open()); |
153 | 502 | MultiSocketDriver socket_driver; |
154 | 502 | CurlMultiPtr multi(curl_multi_init()); |
155 | 502 | if (!multi) { |
156 | 0 | return stats; |
157 | 0 | } |
158 | | |
159 | 502 | const long max_host_connections = |
160 | 502 | static_cast<long>(std::min<std::size_t>(plan.max_host_connections(), transfer_count)); |
161 | 502 | const long max_total_connections = |
162 | 502 | static_cast<long>(std::min<std::size_t>(plan.max_total_connections(), transfer_count)); |
163 | 502 | const long cache_size = |
164 | 502 | static_cast<long>(std::min<std::size_t>(plan.connection_cache_size(), scenario_limits::kMaxMultiTransfers * 2)); |
165 | 502 | (void)curl_multi_setopt(multi.get(), CURLMOPT_MAX_HOST_CONNECTIONS, max_host_connections); |
166 | 502 | (void)curl_multi_setopt(multi.get(), CURLMOPT_MAX_TOTAL_CONNECTIONS, max_total_connections); |
167 | 502 | (void)curl_multi_setopt(multi.get(), CURLMOPT_MAXCONNECTS, cache_size); |
168 | 502 | (void)curl_multi_setopt(multi.get(), CURLMOPT_PIPELINING, |
169 | 502 | plan.multiplex() ? static_cast<long>(CURLPIPE_MULTIPLEX) : 0L); |
170 | | |
171 | 502 | const bool socket_mode = plan.drive_mode() == curl::fuzzer::proto::MULTI_DRIVE_SOCKET; |
172 | 502 | const bool socket_driver_installed = socket_mode && socket_driver.Install(multi.get()); |
173 | 502 | std::array<TransferState, scenario_limits::kMaxMultiTransfers> transfers; |
174 | 502 | const std::string url = "http://" + scenario.host_path(); |
175 | 2.15k | for (std::size_t index = 0; index < transfer_count; ++index) { |
176 | 1.64k | TransferState& transfer = transfers[index]; |
177 | 1.64k | transfer.easy.reset(curl_easy_init()); |
178 | 1.64k | if (!transfer.easy) { |
179 | 0 | continue; |
180 | 0 | } |
181 | 1.64k | transfer.connect_to.reset(ApplyBaselineOptions(transfer.easy.get(), curl::fuzzer::proto::SCHEME_HTTP)); |
182 | 1.64k | (void)curl_easy_setopt(transfer.easy.get(), CURLOPT_URL, url.c_str()); |
183 | 1.64k | mock.Install(transfer.easy.get()); |
184 | 1.64k | (void)ApplyScenarioOptions(transfer.easy.get(), scenario); |
185 | 1.64k | transfer.request_data = std::make_unique<ScenarioRequestData>(transfer.easy.get(), scenario); |
186 | 1.64k | mock.ConfigureRequestData(transfer.request_data.get()); |
187 | 1.64k | if (curl_multi_add_handle(multi.get(), transfer.easy.get()) == CURLM_OK) { |
188 | 1.64k | transfer.attached = true; |
189 | 1.64k | ++stats.added_handles; |
190 | 1.64k | } |
191 | 1.64k | } |
192 | | |
193 | 502 | if (stats.added_handles == 0) { |
194 | 0 | multi.reset(); |
195 | 0 | return stats; |
196 | 0 | } |
197 | | |
198 | 502 | const std::size_t action_count = std::min<std::size_t>(plan.actions_size(), scenario_limits::kMaxMultiActions); |
199 | 502 | std::size_t next_action = 0; |
200 | 502 | bool initial_action_changed = false; |
201 | 502 | if (next_action < action_count) { |
202 | 281 | initial_action_changed = |
203 | 281 | ApplyAction(plan.actions(static_cast<int>(next_action++)), multi.get(), &transfers, transfer_count); |
204 | 281 | ++stats.actions_consumed; |
205 | 281 | } |
206 | | |
207 | 502 | int running_handles = 0; |
208 | 502 | CURLMcode result = socket_driver_installed ? socket_driver.Start(&running_handles) |
209 | 502 | : curl_multi_perform(multi.get(), &running_handles); |
210 | 502 | bool completion_progress = DrainCompletionMessages(multi.get(), &transfers, transfer_count, &stats); |
211 | 502 | ProbeMultiWaitApis(multi.get()); |
212 | 502 | if (plan.wake_multi()) { |
213 | 172 | if (socket_driver_installed) { |
214 | 74 | socket_driver.ProbeControlApis(); |
215 | 98 | } else { |
216 | 98 | (void)curl_multi_wakeup(multi.get()); |
217 | 98 | } |
218 | 172 | } |
219 | | |
220 | 502 | int idle_iterations = initial_action_changed || completion_progress ? 0 : 1; |
221 | 10.3k | for (int iteration = 0; result == CURLM_OK && iteration < kMaxDriveIterations; ++iteration) { |
222 | 10.3k | if (running_handles == 0 && next_action >= action_count) { |
223 | 352 | break; |
224 | 352 | } |
225 | | |
226 | 9.97k | bool made_progress = mock.ServiceConnections(); |
227 | 9.97k | if (next_action < action_count) { |
228 | 1.14k | made_progress = |
229 | 1.14k | ApplyAction(plan.actions(static_cast<int>(next_action++)), multi.get(), &transfers, transfer_count) || |
230 | 544 | made_progress; |
231 | 1.14k | ++stats.actions_consumed; |
232 | | // Consuming an ordered no-op is still harness progress: continue to the |
233 | | // next action instead of letting an idle transport hide its suffix. |
234 | 1.14k | made_progress = true; |
235 | 1.14k | } |
236 | | |
237 | 9.97k | const int running_before = running_handles; |
238 | 9.97k | if (socket_driver_installed) { |
239 | 8.53k | const MultiSocketDriver::DriveResult drive = socket_driver.DriveReady(&running_handles); |
240 | 8.53k | result = drive.code; |
241 | 8.53k | made_progress = drive.made_progress || made_progress; |
242 | 8.53k | } else { |
243 | 1.43k | result = curl_multi_perform(multi.get(), &running_handles); |
244 | 1.43k | made_progress = running_handles != running_before || made_progress; |
245 | 1.43k | } |
246 | 9.97k | made_progress = DrainCompletionMessages(multi.get(), &transfers, transfer_count, &stats) || made_progress; |
247 | | |
248 | 9.97k | if (made_progress) { |
249 | 8.88k | idle_iterations = 0; |
250 | 8.88k | } else if (++idle_iterations >= kMaxIdleIterations) { |
251 | 136 | break; |
252 | 136 | } |
253 | 9.97k | } |
254 | | |
255 | 502 | (void)DrainCompletionMessages(multi.get(), &transfers, transfer_count, &stats); |
256 | 502 | stats.opened_connections = mock.opened_connection_count(); |
257 | 2.15k | for (std::size_t index = 0; index < transfer_count; ++index) { |
258 | 1.64k | TransferState& transfer = transfers[index]; |
259 | 1.64k | if (!transfer.easy) { |
260 | 0 | continue; |
261 | 0 | } |
262 | 1.64k | (void)curl_easy_pause(transfer.easy.get(), CURLPAUSE_CONT); |
263 | 1.64k | if (transfer.attached) { |
264 | 1.50k | (void)curl_multi_remove_handle(multi.get(), transfer.easy.get()); |
265 | 1.50k | transfer.attached = false; |
266 | 1.50k | } |
267 | 1.64k | } |
268 | | |
269 | | // Destroy the multi while both its socket callback storage and mock peer are |
270 | | // live, then detach request-data pointers before cleaning each easy handle. |
271 | 502 | multi.reset(); |
272 | 2.15k | for (std::size_t index = 0; index < transfer_count; ++index) { |
273 | 1.64k | transfers[index].request_data.reset(); |
274 | 1.64k | transfers[index].easy.reset(); |
275 | 1.64k | transfers[index].connect_to.reset(); |
276 | 1.64k | } |
277 | 502 | return stats; |
278 | 502 | } |
279 | | |
280 | | } // namespace proto_fuzzer |