XRootD
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>
26 #include <XrdSys/XrdSysPageSize.hh>
27 
28 using namespace XrdClHttp;
29 
30 CurlVectorReadOp::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 
38 bool
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 
60 void
61 CurlVectorReadOp::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 
84 void
86 {
87  SetDone(false);
88  if (m_handler == nullptr) {return;}
89 
90  // If there's a partial last response, give it to the client.
91  if (m_chunk_buffer_idx) {
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();
98  m_vr->SetSize(m_bytes_consumed);
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 
106 void
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 
120 size_t
121 CurlVectorReadOp::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.
127 size_t
128 CurlVectorReadOp::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()) {
142  m_current_op.first = m_headers.GetOffset();
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  if (m_current_op.first != -1 && m_current_op.second != -1) {
152  //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));
153  if (m_skip_bytes) {
154  //m_logger->Debug(kLogXrdClHttp, "Skipping %lld bytes", static_cast<long long>(m_skip_bytes));
155  auto to_skip = (m_skip_bytes < length) ? m_skip_bytes : length;
156  buffer += to_skip;
157  length -= to_skip;
158  m_skip_bytes -= to_skip;
159  continue;
160  } else {
161  auto &chunk = m_chunk_list[m_response_idx];
162  auto remaining = static_cast<off_t>(chunk.GetLength()) - m_chunk_buffer_idx;
163  if (remaining < 0) {
164  return FailCallback(kXR_ServerError, "Invalid chunk framing");
165  }
166  auto to_copy = (static_cast<size_t>(remaining) < length) ? static_cast<size_t>(remaining) : length;
167  //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));
168  memcpy(static_cast<char *>(chunk.GetBuffer()) + m_chunk_buffer_idx, buffer, to_copy);
169  m_chunk_buffer_idx += to_copy;
170  buffer += to_copy;
171  length -= to_copy;
172  // Handle cases where the requested or response chunk is complete
173  if (chunk.GetLength() == m_chunk_buffer_idx) {
174  m_vr->GetChunks().emplace_back(chunk.GetOffset(), m_chunk_buffer_idx, chunk.GetBuffer());
176  m_chunk_buffer_idx = 0;
177  m_response_idx++;
178  if (m_current_op.second == chunk.GetLength()) {
179  m_current_op.first = m_current_op.second = -1;
180  m_multipart_boundary = true;
181  } else {
182  // We may need to skip the remaining bytes or, potentially, the server
183  // coalesced two adjacent requests into one larger response.
184  m_current_op.first += chunk.GetLength();
185  m_current_op.second -= chunk.GetLength();
186  CalculateNextBuffer();
187  continue;
188  }
189  } else if (m_current_op.second == m_chunk_buffer_idx) {
190  // There are no more bytes in the response but the requested chunk hasn't finished.
191  // Add what we have to the results and create a new chunk on the request list from the remainder; perhaps
192  // the server will send it in the future.
193  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);
194  m_vr->GetChunks().emplace_back(chunk.GetOffset(), m_chunk_buffer_idx, chunk.GetBuffer());
196  m_chunk_buffer_idx = 0;
197  m_current_op.first = m_current_op.second = -1;
198  m_multipart_boundary = true;
199  m_response_idx++;
200  }
201  }
202  }
203  if (m_skip_bytes) {
204  continue;
205  }
206 
207  // We are at the boundary between chunks; we must parse header lines to understand the
208  // next thing to do.
209 
210  // The following lambda function returns a string view to the next complete header line,
211  // potentially partially from the previous buffer from curl. If the second item in
212  // the returned pair is false, then we ran out of buffer from curl before finding a
213  // complete line.
214  auto get_next_line = [&]() {
215  std::string_view chunk_header(buffer, length);
216  auto pos = chunk_header.find("\r\n");
217  if (pos == std::string_view::npos) {
218  m_response_headers += chunk_header;
219  length = 0;
220  return std::make_pair(std::string_view(), false);
221  } else {
222  auto line_part = chunk_header.substr(0, pos);
223  if (!m_response_headers.empty()) {
224  m_response_headers += line_part;
225  m_header_line = std::move(m_response_headers);
226  } else {
227  m_header_line.assign(line_part);
228  }
229  buffer += pos + 2;
230  length -= pos + 2;
231  return std::make_pair(std::string_view(m_header_line), true);
232  }
233  };
234 
235  // Consume the boundary line.
236  bool last_segment = false;
237  if (m_multipart_boundary) {
238  while (true) {
239  auto [line, ok] = get_next_line();
240  if (!ok) {
241  return orig_length;
242  }
243  // Per RFC7233, Appendix A, Implementation note 1, multiple CRLF might precede the
244  // first boundary string in the body. However, the XRootD server appears to have an
245  // extra CRLF in front of every boundary string.
246  if (line.empty()) {continue;}
247  if (line == m_headers.MultipartSeparator()) {
248  break;
249  }
250  if (line == m_headers.MultipartSeparator() + "--") {
251  last_segment = true;
252  break;
253  }
254  std::stringstream ss;
255  ss << "Server has responded with an invalid boundary line: '" << line << "' (expected '" << m_headers.MultipartSeparator() << "')";
256  return FailCallback(kXR_ServerError, ss.str());
257  }
258  }
259  if (last_segment) {
260  length = 0;
261  break;
262  }
263  // Consume the header lines
264  while (true) {
265  auto [line, ok] = get_next_line();
266  if (!ok) {
267  m_multipart_boundary = false;
268  return orig_length;
269  }
270  if (line.empty()) {
271  break;
272  }
273  auto header_name_end = line.find(':');
274  if (header_name_end == std::string_view::npos) {
275  std::stringstream ss; ss << "Invalid header line in response from server: " << line;
276  return FailCallback(kXR_ServerError, ss.str());
277  }
278  auto header_name = line.substr(0, header_name_end);
279  // Cannot use strcasecmp here as a string_view's data is not necessarily nul-terminated.
280  // len("content-type") == 13
281  if (header_name.size() != 13 || strncasecmp(header_name.data(), "content-range", 13)) {
282  continue;
283  }
284  // We are parsing a Content-Range value.
285  // Example: Content-Range: bytes 7000-7999/8000
286  auto value = line.substr(header_name_end + 1);
287 
288  // Advance whitespace
289  while (!value.empty() && value[0] == ' ') {
290  value = value.substr(1);
291  }
292 
293  if (value.substr(0, 5) != "bytes") {
294  std::stringstream ss; ss << "Invalid Content-Range value (no 'bytes' unit): " << value;
295  return FailCallback(kXR_ServerError, ss.str());
296  }
297 
298  value = value.substr(5);
299  while (!value.empty() && value[0] == ' ') {
300  value = value.substr(1);
301  }
302 
303  // Example: 500-999/8000
304  size_t count;
305  long long bytes_val;
306  try {
307  // It may seem strange to see the string_view data being passed to std::stoll here
308  // as it's not guaranteed to be null-terminated. However, by this point, we do know
309  // there's a CRLF in the buffer -- that is sufficient to guarantee the stoll search
310  // terminates before it goes out-of-bounds.
311  bytes_val = std::stoll(value.data(), &count);
312  } catch (std::invalid_argument &) {
313  std::stringstream ss; ss << "Invalid Content-Range value (no integer in range start): " << value;
314  return FailCallback(kXR_ServerError, ss.str());
315  } catch (std::out_of_range &) {
316  std::stringstream ss; ss << "Invalid Content-Range value (out of range): " << value;
317  return FailCallback(kXR_ServerError, ss.str());
318  }
319  if (value.size() <= count || value[count] != '-') {
320  std::stringstream ss; ss << "Invalid Content-Range value (no dash in range): " << value;
321  return FailCallback(kXR_ServerError, ss.str());
322  }
323  m_current_op.first = bytes_val;
324  value = value.substr(count + 1);
325  try {
326  bytes_val = std::stoll(value.data(), &count);
327  } catch (std::invalid_argument &) {
328  std::stringstream ss; ss << "Invalid Content-Range value (no integer in range end): " << value;
329  return FailCallback(kXR_ServerError, ss.str());
330  } catch (std::out_of_range &) {
331  std::stringstream ss; ss << "Invalid Content-Range value (out of range in range end): " << value;
332  return FailCallback(kXR_ServerError, ss.str());
333  }
334  if (value.size() <= count || value[count] != '/') {
335  std::stringstream ss; ss << "Invalid Content-Range value (no trailing /): " << value;
336  return FailCallback(kXR_ServerError, ss.str());
337  }
338  auto length = bytes_val + 1 - m_current_op.first;
339  if (length < 0) {
340  std::stringstream ss; ss << "Invalid Content-Range value (negative length): " << line;
341  return FailCallback(kXR_ServerError, ss.str());
342  }
343  if (length > std::numeric_limits<decltype(m_current_op.second)>::max()) {
344  std::stringstream ss; ss << "Invalid Content-Range value (length too long): " << line;
345  return FailCallback(kXR_ServerError, ss.str());
346  }
347  m_current_op.second = length;
348 
349  // We now have a valid response range; locate a buffer where we will copy the bytes into.
350  CalculateNextBuffer();
351  }
352  m_multipart_boundary = true;
353 
354  // Check to see if the Content-Range was missing.
355  if (!last_segment && (m_current_op.first == -1 || m_current_op.second == -1)) {
356  return FailCallback(kXR_ServerError, "Response segment is missing a Content-Range header");
357  }
358  }
359  return orig_length;
360 }
361 
362 void CurlVectorReadOp::CalculateNextBuffer() {
363  // Strategy is to select the index where we will throw away the fewest bytes.
364  off_t distance = std::numeric_limits<off_t>::max();
365  auto starting_idx = m_response_idx;
366  for (decltype(m_chunk_list)::size_type ctr=0; ctr<m_chunk_list.size(); ctr++) {
367  auto idx = (starting_idx + ctr) % m_chunk_list.size();
368  if (static_cast<uint64_t>(m_current_op.first) == m_chunk_list[idx].GetOffset()) {
369  m_response_idx = idx;
370  distance = 0;
371  break;
372  }
373  off_t bytes_to_skip = static_cast<off_t>(m_chunk_list[idx].GetOffset()) - m_current_op.first;
374  //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));
375  if (bytes_to_skip > 0 && bytes_to_skip < distance) {
376  distance = bytes_to_skip;
377  m_response_idx = idx;
378  // Note we don't break; some other request might be a better fit.
379  }
380  }
381  m_chunk_buffer_idx = 0;
382  if (distance > 0) {
383  m_skip_bytes = distance;
384  } else {
385  m_skip_bytes = 0;
386  }
387 }
@ kXR_ServerError
Definition: XProtocol.hh:1044
void CURL
if(Avsz)
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
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
void Fail(uint16_t errCode, uint32_t errNum, const std::string &msg) 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)
bool HTTPStatusIsError(unsigned status)
const uint16_t stError
An error occurred that could potentially be retried.
Definition: XrdClStatus.hh:32
std::vector< ChunkInfo > ChunkList
List of chunks.
ConnectionCallout *(*)(const std::string &, const ResponseInfo &) CreateConnCalloutType
const uint64_t kLogXrdClHttp