XRootD
Loading...
Searching...
No Matches
XrdClHttpOpReadV.cc
Go to the documentation of this file.
1/******************************************************************************/
2/* Copyright (C) 2025, Pelican Project, Morgridge Institute for Research */
3/* */
4/* This file is part of the XrdClHttp client plugin for XRootD. */
5/* */
6/* XRootD is free software: you can redistribute it and/or modify it under */
7/* the terms of the GNU Lesser General Public License as published by the */
8/* Free Software Foundation, either version 3 of the License, or (at your */
9/* option) any later version. */
10/* */
11/* XRootD is distributed in the hope that it will be useful, but WITHOUT */
12/* ANY WARRANTY; without even the implied warranty of MERCHANTABILITY or */
13/* FITNESS FOR A PARTICULAR PURPOSE. See the GNU Lesser General Public */
14/* License for more details. */
15/* */
16/* The copyright holder's institutional names and contributor's names may not */
17/* be used to endorse or promote products derived from this software without */
18/* specific prior written permission of the institution or contributor. */
19/******************************************************************************/
20
21#include "XrdClHttpOps.hh"
22
23#include <XrdCl/XrdClLog.hh>
25#include <XrdOuc/XrdOucCRC.hh>
27
28using namespace XrdClHttp;
29
30CurlVectorReadOp::CurlVectorReadOp(XrdCl::ResponseHandler *handler, const std::string &url, struct timespec timeout,
31 const XrdCl::ChunkList &op_list, XrdCl::Log *logger, CreateConnCalloutType callout,
32 HeaderCallout *header_callout) :
33 CurlOperation(handler, url, timeout, logger, callout, header_callout),
34 m_vr(new XrdCl::VectorReadInfo()),
35 m_chunk_list(op_list)
36 {}
37
38bool
40{
41 if (!CurlOperation::Setup(curl, worker)) return false;
42 curl_easy_setopt(m_curl.get(), CURLOPT_WRITEFUNCTION, CurlVectorReadOp::WriteCallback);
43 curl_easy_setopt(m_curl.get(), CURLOPT_WRITEDATA, this);
44
45 std::stringstream ss;
46 auto multiple = false;
47 for (const auto &chunk : m_chunk_list) {
48 if (!chunk.GetLength()) continue;
49 if (multiple) {ss << ",";}
50 ss << chunk.GetOffset() << "-" << chunk.GetOffset() + chunk.GetLength() - 1;
51 multiple = true;
52 }
53 auto byte_range_val = ss.str();
54 if (byte_range_val.size()) {
55 m_headers_list.emplace_back("Range", "bytes=" + byte_range_val);
56 }
57 return true;
58}
59
60void
61CurlVectorReadOp::Fail(uint16_t errCode, uint32_t errNum, const std::string &msg)
62{
63 std::string custom_msg = msg;
64 SetDone(true);
65 if (m_handler == nullptr) {return;}
66 std::string offset = "(unknown)";
67 std::string length = "(unknown)";
68 if (!m_chunk_list.empty()) {
69 offset = std::to_string(m_chunk_list[0].GetOffset());
70 length = std::to_string(m_chunk_list[0].GetLength());
71 }
72 if (!custom_msg.empty()) {
73 m_logger->Debug(kLogXrdClHttp, "curl operation with vector starting offset %s / length %s failed with message: %s", offset.c_str(), length.c_str(), custom_msg.c_str());
74 custom_msg += " (vector read operation starting at offset " + offset + " / length " + length + ")";
75 } else {
76 m_logger->Debug(kLogXrdClHttp, "curl vector operation starting at offset %s / length %s failed with status code %d", offset.c_str(), length.c_str(), errNum);
77 }
78 auto status = new XrdCl::XRootDStatus(XrdCl::stError, errCode, errNum, custom_msg);
79 auto handle = m_handler;
80 m_handler = nullptr;
81 handle->HandleResponse(status, nullptr);
82}
83
84void
86{
87 SetDone(false);
88 if (m_handler == nullptr) {return;}
89
90 // If there's a partial last response, give it to the client.
92 auto &chunk = m_chunk_list[m_response_idx];
93 m_vr->GetChunks().emplace_back(chunk.GetOffset(), m_chunk_buffer_idx, chunk.GetBuffer());
95 }
96
97 auto status = new XrdCl::XRootDStatus();
99 auto obj = new XrdCl::AnyObject();
100 obj->Set(m_vr.release());
101 auto handle = m_handler;
102 m_handler = nullptr;
103 handle->HandleResponse(status, obj);
104}
105
106void
108{
109 if (m_curl == nullptr) return;
110 curl_easy_setopt(m_curl.get(), CURLOPT_WRITEFUNCTION, nullptr);
111 curl_easy_setopt(m_curl.get(), CURLOPT_WRITEDATA, nullptr);
112 curl_easy_setopt(m_curl.get(), CURLOPT_HTTPHEADER, nullptr);
113 curl_easy_setopt(m_curl.get(), CURLOPT_OPENSOCKETFUNCTION, nullptr);
114 curl_easy_setopt(m_curl.get(), CURLOPT_OPENSOCKETDATA, nullptr);
115 curl_easy_setopt(m_curl.get(), CURLOPT_SOCKOPTFUNCTION, nullptr);
116 curl_easy_setopt(m_curl.get(), CURLOPT_SOCKOPTDATA, nullptr);
118}
119
120size_t
121CurlVectorReadOp::WriteCallback(char *buffer, size_t size, size_t nitems, void *this_ptr)
122{
123 return static_cast<CurlVectorReadOp*>(this_ptr)->Write(buffer, size * nitems);
124}
125
126// Given a buffer of data from curl, parse it and write it to the response buffers.
127size_t
128CurlVectorReadOp::Write(char *orig_buffer, size_t orig_length)
129{
130 UpdateBytes(orig_length);
131 //m_logger->Debug(kLogXrdClHttp, "Received a write of size %ld with contents:\n%s", static_cast<long>(orig_length), std::string(orig_buffer, orig_length).c_str());
132
133 // Handle the (hopefully uncommon) cases where the server responds to a vector read op
134 // with a single response. We set the length of the response to the max as we
135 // don't care how many bytes the server actually sends.
136 if (GetStatusCode() == 200) {
137 m_current_op.first = 0;
138 m_current_op.second = std::numeric_limits<off_t>::max();
139 } else if (HTTPStatusIsError(GetStatusCode())) {
140 return orig_length;
141 } else if (!m_headers.IsMultipartByterange()) {
143 m_current_op.second = std::numeric_limits<off_t>::max();
144 }
145
146 auto buffer = orig_buffer;
147 auto length = orig_length;
148
149 while (length) {
150 // If we're in the middle of a response chunk, copy as much data as possible.
151 // m_multipart_boundary is false while part headers are still being assembled
152 // across curl writes. Content-Range may already have filled m_current_op, but
153 // the terminating blank line (and thus the payload) has not arrived yet.
154 if (m_multipart_boundary && m_current_op.first != -1 && m_current_op.second != -1) {
155 //m_logger->Debug(kLogXrdClHttp, "Processing response buffer of (%lld, %lld)", static_cast<long long>(m_current_op.first), static_cast<long long>(m_current_op.second));
156 if (m_skip_bytes) {
157 //m_logger->Debug(kLogXrdClHttp, "Skipping %lld bytes", static_cast<long long>(m_skip_bytes));
158 auto to_skip = (m_skip_bytes < length) ? m_skip_bytes : length;
159 buffer += to_skip;
160 length -= to_skip;
161 m_skip_bytes -= to_skip;
162 continue;
163 } else {
164 auto &chunk = m_chunk_list[m_response_idx];
165 auto remaining = static_cast<off_t>(chunk.GetLength()) - m_chunk_buffer_idx;
166 if (remaining < 0) {
167 return FailCallback(kXR_ServerError, "Invalid chunk framing");
168 }
169 auto to_copy = (static_cast<size_t>(remaining) < length) ? static_cast<size_t>(remaining) : length;
170 //m_logger->Debug(kLogXrdClHttp, "Copying %lld bytes to request buffer %ld at offset %lld", static_cast<long long>(to_copy), m_response_idx, static_cast<long long>(m_chunk_buffer_idx));
171 memcpy(static_cast<char *>(chunk.GetBuffer()) + m_chunk_buffer_idx, buffer, to_copy);
172 m_chunk_buffer_idx += to_copy;
173 buffer += to_copy;
174 length -= to_copy;
175 // Handle cases where the requested or response chunk is complete
176 if (chunk.GetLength() == m_chunk_buffer_idx) {
177 m_vr->GetChunks().emplace_back(chunk.GetOffset(), m_chunk_buffer_idx, chunk.GetBuffer());
181 if (m_current_op.second == chunk.GetLength()) {
182 m_current_op.first = m_current_op.second = -1;
183 m_multipart_boundary = true;
184 } else {
185 // We may need to skip the remaining bytes or, potentially, the server
186 // coalesced two adjacent requests into one larger response.
187 m_current_op.first += chunk.GetLength();
188 m_current_op.second -= chunk.GetLength();
189 CalculateNextBuffer();
190 continue;
191 }
192 } else if (m_current_op.second == m_chunk_buffer_idx) {
193 // There are no more bytes in the response but the requested chunk hasn't finished.
194 // Add what we have to the results and create a new chunk on the request list from the remainder; perhaps
195 // the server will send it in the future.
196 m_chunk_list.emplace_back(chunk.GetOffset() + m_chunk_buffer_idx, chunk.GetLength() - m_chunk_buffer_idx, static_cast<char*>(chunk.GetBuffer()) + m_chunk_buffer_idx);
197 m_vr->GetChunks().emplace_back(chunk.GetOffset(), m_chunk_buffer_idx, chunk.GetBuffer());
200 m_current_op.first = m_current_op.second = -1;
201 m_multipart_boundary = true;
203 }
204 }
205 }
206 if (m_skip_bytes) {
207 continue;
208 }
209
210 // We are at the boundary between chunks; we must parse header lines to understand the
211 // next thing to do.
212
213 // The following lambda function returns a string view to the next complete header line,
214 // potentially partially from the previous buffer from curl. If the second item in
215 // the returned pair is false, then we ran out of buffer from curl before finding a
216 // complete line.
217 auto get_next_line = [&]() {
218 std::string_view chunk_header(buffer, length);
219 // CRLF may be split across curl writes: previous chunk ended with \r
220 // and this one starts with \n. find("\r\n") would otherwise glue this
221 // line to the next header (e.g. "--123456\r" + "\nContent-type: ...\r\n").
222 if (!m_response_headers.empty() && m_response_headers.back() == '\r' &&
223 !chunk_header.empty() && chunk_header.front() == '\n') {
224 m_response_headers.pop_back();
226 buffer += 1;
227 length -= 1;
228 return std::make_pair(std::string_view(m_header_line), true);
229 }
230 auto pos = chunk_header.find("\r\n");
231 if (pos == std::string_view::npos) {
232 m_response_headers += chunk_header;
233 length = 0;
234 return std::make_pair(std::string_view(), false);
235 } else {
236 auto line_part = chunk_header.substr(0, pos);
237 if (!m_response_headers.empty()) {
238 m_response_headers += line_part;
240 } else {
241 m_header_line.assign(line_part);
242 }
243 buffer += pos + 2;
244 length -= pos + 2;
245 return std::make_pair(std::string_view(m_header_line), true);
246 }
247 };
248
249 // Consume the boundary line.
250 bool last_segment = false;
251 if (m_multipart_boundary) {
252 while (true) {
253 auto [line, ok] = get_next_line();
254 if (!ok) {
255 return orig_length;
256 }
257 // Per RFC7233, Appendix A, Implementation note 1, multiple CRLF might precede the
258 // first boundary string in the body. However, the XRootD server appears to have an
259 // extra CRLF in front of every boundary string.
260 if (line.empty()) {continue;}
261 if (line == m_headers.MultipartSeparator()) {
262 break;
263 }
264 if (line == m_headers.MultipartSeparator() + "--") {
265 last_segment = true;
266 break;
267 }
268 std::stringstream ss;
269 ss << "Server has responded with an invalid boundary line: '" << line << "' (expected '" << m_headers.MultipartSeparator() << "')";
270 return FailCallback(kXR_ServerError, ss.str());
271 }
272 }
273 if (last_segment) {
274 length = 0;
275 break;
276 }
277 // Consume the header lines
278 while (true) {
279 auto [line, ok] = get_next_line();
280 if (!ok) {
281 m_multipart_boundary = false;
282 return orig_length;
283 }
284 if (line.empty()) {
285 break;
286 }
287 auto header_name_end = line.find(':');
288 if (header_name_end == std::string_view::npos) {
289 std::stringstream ss; ss << "Invalid header line in response from server: " << line;
290 return FailCallback(kXR_ServerError, ss.str());
291 }
292 auto header_name = line.substr(0, header_name_end);
293 // Cannot use strcasecmp here as a string_view's data is not necessarily nul-terminated.
294 // len("content-type") == 13
295 if (header_name.size() != 13 || strncasecmp(header_name.data(), "content-range", 13)) {
296 continue;
297 }
298 // We are parsing a Content-Range value.
299 // Example: Content-Range: bytes 7000-7999/8000
300 auto value = line.substr(header_name_end + 1);
301
302 // Advance whitespace
303 while (!value.empty() && value[0] == ' ') {
304 value = value.substr(1);
305 }
306
307 if (value.substr(0, 5) != "bytes") {
308 std::stringstream ss; ss << "Invalid Content-Range value (no 'bytes' unit): " << value;
309 return FailCallback(kXR_ServerError, ss.str());
310 }
311
312 value = value.substr(5);
313 while (!value.empty() && value[0] == ' ') {
314 value = value.substr(1);
315 }
316
317 // Example: 500-999/8000
318 size_t count;
319 long long bytes_val;
320 try {
321 // It may seem strange to see the string_view data being passed to std::stoll here
322 // as it's not guaranteed to be null-terminated. However, by this point, we do know
323 // there's a CRLF in the buffer -- that is sufficient to guarantee the stoll search
324 // terminates before it goes out-of-bounds.
325 bytes_val = std::stoll(value.data(), &count);
326 } catch (std::invalid_argument &) {
327 std::stringstream ss; ss << "Invalid Content-Range value (no integer in range start): " << value;
328 return FailCallback(kXR_ServerError, ss.str());
329 } catch (std::out_of_range &) {
330 std::stringstream ss; ss << "Invalid Content-Range value (out of range): " << value;
331 return FailCallback(kXR_ServerError, ss.str());
332 }
333 if (value.size() <= count || value[count] != '-') {
334 std::stringstream ss; ss << "Invalid Content-Range value (no dash in range): " << value;
335 return FailCallback(kXR_ServerError, ss.str());
336 }
337 m_current_op.first = bytes_val;
338 value = value.substr(count + 1);
339 try {
340 bytes_val = std::stoll(value.data(), &count);
341 } catch (std::invalid_argument &) {
342 std::stringstream ss; ss << "Invalid Content-Range value (no integer in range end): " << value;
343 return FailCallback(kXR_ServerError, ss.str());
344 } catch (std::out_of_range &) {
345 std::stringstream ss; ss << "Invalid Content-Range value (out of range in range end): " << value;
346 return FailCallback(kXR_ServerError, ss.str());
347 }
348 if (value.size() <= count || value[count] != '/') {
349 std::stringstream ss; ss << "Invalid Content-Range value (no trailing /): " << value;
350 return FailCallback(kXR_ServerError, ss.str());
351 }
352 auto length = bytes_val + 1 - m_current_op.first;
353 if (length < 0) {
354 std::stringstream ss; ss << "Invalid Content-Range value (negative length): " << line;
355 return FailCallback(kXR_ServerError, ss.str());
356 }
357 if (length > std::numeric_limits<decltype(m_current_op.second)>::max()) {
358 std::stringstream ss; ss << "Invalid Content-Range value (length too long): " << line;
359 return FailCallback(kXR_ServerError, ss.str());
360 }
361 m_current_op.second = length;
362
363 // We now have a valid response range; locate a buffer where we will copy the bytes into.
364 CalculateNextBuffer();
365 }
366 m_multipart_boundary = true;
367
368 // Check to see if the Content-Range was missing.
369 if (!last_segment && (m_current_op.first == -1 || m_current_op.second == -1)) {
370 return FailCallback(kXR_ServerError, "Response segment is missing a Content-Range header");
371 }
372 }
373 return orig_length;
374}
375
376void CurlVectorReadOp::CalculateNextBuffer() {
377 // Strategy is to select the index where we will throw away the fewest bytes.
378 off_t distance = std::numeric_limits<off_t>::max();
379 auto starting_idx = m_response_idx;
380 for (decltype(m_chunk_list)::size_type ctr=0; ctr<m_chunk_list.size(); ctr++) {
381 auto idx = (starting_idx + ctr) % m_chunk_list.size();
382 if (static_cast<uint64_t>(m_current_op.first) == m_chunk_list[idx].GetOffset()) {
383 m_response_idx = idx;
384 distance = 0;
385 break;
386 }
387 off_t bytes_to_skip = static_cast<off_t>(m_chunk_list[idx].GetOffset()) - m_current_op.first;
388 //m_logger->Debug(kLogXrdClHttp, "Using client request at index %lu would require us to skip %lld bytes", idx, static_cast<long long>(bytes_to_skip));
389 if (bytes_to_skip > 0 && bytes_to_skip < distance) {
390 distance = bytes_to_skip;
391 m_response_idx = idx;
392 // Note we don't break; some other request might be a better fit.
393 }
394 }
396 if (distance > 0) {
397 m_skip_bytes = distance;
398 } else {
399 m_skip_bytes = 0;
400 }
401}
@ kXR_ServerError
void CURL
void SetDone(bool has_failed)
int FailCallback(XErrorCode ecode, const std::string &emsg)
std::unique_ptr< CURL, void(*)(CURL *)> m_curl
virtual void ReleaseHandle()
void UpdateBytes(uint64_t bytes)
std::vector< std::pair< std::string, std::string > > m_headers_list
XrdCl::ResponseHandler * m_handler
CurlOperation(XrdCl::ResponseHandler *handler, const std::string &url, struct timespec timeout, XrdCl::Log *log, CreateConnCalloutType, HeaderCallout *header_callout)
virtual bool Setup(CURL *curl, CurlWorker &)
size_t Write(char *buffer, size_t size)
CurlVectorReadOp(XrdCl::ResponseHandler *handler, const std::string &url, struct timespec timeout, const XrdCl::ChunkList &op_list, XrdCl::Log *logger, CreateConnCalloutType callout, HeaderCallout *header_callout)
XrdCl::ChunkList m_chunk_list
std::pair< off_t, off_t > m_current_op
std::unique_ptr< XrdCl::VectorReadInfo > m_vr
bool Setup(CURL *curl, CurlWorker &) override
uint64_t GetOffset() const
bool IsMultipartByterange() const
const std::string & MultipartSeparator() const
Handle diagnostics.
Definition XrdClLog.hh:101
void Debug(uint64_t topic, const char *format,...)
Print a debug message.
Definition XrdClLog.cc:282
Handle an async response.
virtual void HandleResponse(XRootDStatus *status, AnyObject *response)
ChunkList & GetChunks()
Get chunks.
void SetSize(uint32_t size)
Set size.
ConnectionCallout *(*)(const std::string &, const ResponseInfo &) CreateConnCalloutType
bool HTTPStatusIsError(unsigned status)
const uint64_t kLogXrdClHttp
const uint16_t stError
An error occurred that could potentially be retried.
std::vector< ChunkInfo > ChunkList
List of chunks.