34#include <condition_variable>
43using namespace std::string_literals;
55 void HandleResponse(XrdCl::XRootDStatus *status,
56 XrdCl::AnyObject *response)
override
58 if (status !=
nullptr)
60 const std::unique_ptr<XrdCl::XRootDStatus> owned_status{status};
63 if (response !=
nullptr)
65 const std::unique_ptr<XrdCl::AnyObject> owned_response{response};
69 const std::lock_guard<std::mutex> lock{mutex};
73 is_ready_changed.notify_all();
77 bool wait(std::chrono::seconds timeout = std::chrono::seconds(0))
79 std::unique_lock<std::mutex> lock{mutex};
83 if (timeout.count() > 0)
84 ready = is_ready_changed.wait_for(lock, timeout, [
this]{
return is_ready; });
86 is_ready_changed.wait(lock, [
this]{
return is_ready; });
95 std::condition_variable is_ready_changed{};
104 explicit ExtraHeaderCallout(HeaderList extra) : extra(std::move(extra)) {}
106 std::shared_ptr<HeaderList> GetHeaders(
const std::string &,
108 const HeaderList &headers)
override
110 auto result = std::make_shared<HeaderList>(headers);
111 result->insert(result->end(), extra.begin(), extra.end());
126 m_logger->Debug(
kLogXrdClHttp,
"Constructing filesystem object with base URL %s", url.c_str());
132 m_url.SetParams(map);
138Filesystem::QueueOperation(std::unique_ptr<CurlOperation> operation,
139 const char *description)
143 m_queue->Produce(std::move(operation));
145 catch(
const std::exception &ex)
148 "Failed to add %s to queue: %s", description, ex.what());
152 return XrdCl::XRootDStatus();
163 auto full_url = GetCurrentURL(path);
165 m_logger->Debug(
kLogXrdClHttp,
"Filesystem::DirList path %s", path.c_str());
166 auto listdirOp = std::make_unique<XrdClHttp::CurlListdirOp>(
168 m_url.GetHostName() +
":" + std::to_string(m_url.GetPort()),
169 SendResponseInfo(), ts, m_logger,
170 GetConnCallout(), m_header_callout.load(std::memory_order_acquire));
171 return QueueOperation(std::move(listdirOp),
"directory list operation");
175Filesystem::GetConnCallout()
const {
176 std::string pointer_str;
177 if (!
GetProperty(
"XrdClConnectionCallout", pointer_str) && pointer_str.empty()) {
182 pointer = std::stoll(pointer_str,
nullptr, 16);
194 std::string &value)
const
196 std::shared_lock lock(m_properties_mutex);
198 const auto p = m_properties.find(name);
199 if (p == std::end(m_properties)) {
220 auto locateInfo = std::make_unique<XrdCl::LocationInfo>();
223 auto obj = std::make_unique<XrdCl::AnyObject>();
224 obj->Set(locateInfo.release());
238 auto full_url = GetCurrentURL(path);
239 m_logger->Debug(
kLogXrdClHttp,
"Filesystem::MkDir path %s", full_url.c_str());
241 auto mkdirOp = std::make_unique<CurlMkcolOp>(
242 handler, full_url, ts, m_logger, SendResponseInfo(), GetConnCallout(),
243 m_header_callout.load(std::memory_order_acquire));
244 return QueueOperation(std::move(mkdirOp),
"filesystem mkdir operation");
248 const std::vector<std::string> &fileList,
259 EINVAL,
"missing prepare file list");
263 auto tapeOp = std::make_unique<CurlTapePrepareOp>(
264 handler, m_url.GetURL(), fileList, flags, ts, m_logger,
266 m_header_callout.load(std::memory_order_acquire));
267 return QueueOperation(std::move(tapeOp),
"Tape prepare operation");
276 std::unique_ptr<CurlOperation> operation;
277 const char *description =
nullptr;
283 operation = std::make_unique<CurlTapeQueryOp>(
284 handler, m_url.GetURL(), queryCode, arg, ts, m_logger,
286 m_header_callout.load(std::memory_order_acquire));
287 description =
"Tape query operation";
291 const auto url = GetCurrentURL(arg.
ToString());
293 "XrdClHttp::Filesystem::Query checksum path %s", url.c_str());
299 const auto iter = urlObj.
GetParams().find(
"cks.type");
306 "Unknown checksum type %s", iter->second.c_str());
309 "unknown checksum type '" + iter->second +
"'");
312 operation = std::make_unique<CurlChecksumOp>(
313 handler, url, preferred, ts, m_logger, SendResponseInfo(),
315 m_header_callout.load(std::memory_order_acquire));
316 description =
"checksum operation";
321 const std::string path = arg.
ToString();
323 "XrdClHttp::Filesystem::Query xattr full_url %s, path %s",
324 m_url.GetURL().c_str(), path.c_str());
325 operation = std::make_unique<CurlQueryOp>(
326 handler, path, ts, m_logger, SendResponseInfo(),
327 GetConnCallout(), queryCode,
328 m_header_callout.load(std::memory_order_acquire));
329 description =
"xattr query operation";
337 return QueueOperation(std::move(operation), description);
347 auto full_url = GetCurrentURL(path);
348 m_logger->Debug(
kLogXrdClHttp,
"Filesystem::Rm path %s", full_url.c_str());
350 auto deleteOp = std::make_unique<CurlDeleteOp>(
351 handler, full_url, ts, m_logger, SendResponseInfo(),
352 GetConnCallout(), m_header_callout.load(std::memory_order_acquire));
353 return QueueOperation(std::move(deleteOp),
"filesystem delete operation");
361 return Rm(path, handler, timeout);
366 const std::string &value)
368 if (name ==
"XrdClHttpHeaderCallout") {
371 pointer = std::stoll(value,
nullptr, 16);
381 std::unique_lock lock(m_properties_mutex);
382 m_properties[name] = value;
393 auto full_url = GetCurrentURL(path);
394 m_logger->Debug(
kLogXrdClHttp,
"Filesystem::Stat path %s", full_url.c_str());
396 auto statOp = std::make_unique<CurlStatOp>(
397 handler, full_url, ts, m_logger, SendResponseInfo(),
398 GetConnCallout(), m_header_callout.load(std::memory_order_acquire));
399 return QueueOperation(std::move(statOp),
"filesystem stat operation");
407 return protocol ==
"http" || protocol ==
"https" || protocol ==
"dav" || protocol ==
"davs";
431 const std::string &url,
439 const auto callout = std::make_shared<ExtraHeaderCallout>(hdrs);
440 const auto rh = std::make_shared<SyncResponseHandler>();
441 const auto op = std::make_shared<StatOp>(rh.get(), url, timespec{timeout, 0}, log,
true,
nullptr, callout.get());
445 queue.
Produce(std::shared_ptr<StatOp>(op.get(), [callout, rh, op](
auto){}));
447 if (!rh->wait(std::chrono::seconds(timeout)))
449 else if (!op->IsDone() || op->HasFailed())
452 return static_cast<std::size_t
>(op->GetStatInfo().first);
454 catch (
const std::exception &e) {
463 const std::string &dest,
469 log->
Debug(
kLogXrdClHttp,
"XrdClHttp::ThirdPartyCopy src %s dst %s", source.c_str(), dest.c_str());
473 log->
Error(
kLogXrdClHttp,
"Third party copy can only be done between http(s) protocols");
475 "Third party copy can only be done between http(s) protocols");
479 int number_of_streams = 1;
480 time_t init_timeout = 0;
481 time_t tpc_timeout = 0;
482 std::string token_file;
483 bool force_overwrite =
false;
488 properties->
Get(
"initTimeout", init_timeout);
489 properties->
Get(
"tpcTimeout", tpc_timeout);
490 token_file = properties->
Get<std::string>(
"thirdPartyTokenFile");
491 properties->
Get(
"force", force_overwrite);
495 if (tpc_timeout == 0)
496 tpc_timeout = timeout;
499 env->GetInt(
"SubStreamsPerChannel", number_of_streams);
505 if (!token_file.empty() && !
ParseTokenFile(token_file, src_hdrs, dst_hdrs))
507 "Failed to parse the token file '" + token_file +
"'");
512 std::any_of(hdrs.begin(), hdrs.end(),
513 [] (
const auto &header) { return header.first ==
"Authorization"s; });
516 if (is_unsecure(src_hdrs, source) || is_unsecure(dst_hdrs, dest))
518 log->
Error(
kLogXrdClHttp,
"Refusing to send an Authorization header over an unencrypted http URL");
520 "Refusing to send an Authorization header over an unencrypted http URL");
527 stat_size(*m_queue, dest, dst_hdrs, init_timeout, log,
"destination"))
531 "The destination exists; use --force to replace it");
534 headers.emplace_back(
"X-Number-Of-Streams"s, std::to_string(number_of_streams));
538 dst_hdrs.emplace_back(
"Overwrite"s, force_overwrite ?
"T"s :
"F"s);
540 std::size_t size = 0;
542 if (progress_handler)
544 size =
stat_size(*m_queue, source, src_hdrs, init_timeout, log,
"source").value_or(0);
548 const auto rh = std::make_shared<SyncResponseHandler>();
549 const auto op = std::make_shared<CopyOp>(rh.get(), source, src_hdrs, dest, dst_hdrs, headers,
550 mode, timespec{tpc_timeout, 0}, log,
nullptr);
551 op->SetProgressHandler(progress_handler);
555 m_queue->Produce(std::shared_ptr<CopyOp>(op.get(), [rh, op](
auto){}));
557 catch (
const std::exception &e) {
562 if (!rh->wait(std::chrono::seconds(tpc_timeout)))
565 if (op->IsDone() && !op->IsSentSuccessfully())
568 if (progress_handler && size > 0)
574bool Filesystem::SendResponseInfo()
const {
579std::string Filesystem::GetCurrentURL(
const std::string &path)
const {
582 auto prefix = m_url.
GetURL();
583 std::string_view prefix_view = prefix;
584 while (!prefix_view.empty() && prefix_view[prefix_view.size() - 1] ==
'/')
585 prefix_view = prefix_view.substr(0, prefix_view.size() - 1);
588 std::string_view path_view = path;
589 while (!path_view.empty() && path_view[0] ==
'/')
590 path_view = path_view.substr(1);
591 auto retval = std::string(prefix_view) +
"/" + std::string(path_view);
595 std::shared_lock lock(m_properties_mutex);
596 auto iter = m_properties.find(
"XrdClHttpQueryParam");
597 if (iter != m_properties.end() && !iter->second.empty()) {
598 retval += ((retval.find(
'?') == std::string::npos) ?
'?' :
':') + iter->second;
static bool is_http_url(const std::string &url)
#define ResponseInfoProperty
std::vector< std::pair< std::string, std::string > > Headers
static struct timespec GetHeaderTimeoutWithDefault(time_t oper_timeout)
virtual XrdCl::XRootDStatus MkDir(const std::string &path, XrdCl::MkDirFlags::Flags flags, XrdCl::Access::Mode mode, XrdCl::ResponseHandler *handler, time_t timeout) override
virtual bool SetProperty(const std::string &name, const std::string &value) override
virtual XrdCl::XRootDStatus Prepare(const std::vector< std::string > &fileList, XrdCl::PrepareFlags::Flags flags, uint8_t priority, XrdCl::ResponseHandler *handler, time_t timeout) override
XrdCl::XRootDStatus DirList(const std::string &path, XrdCl::DirListFlags::Flags flags, XrdCl::ResponseHandler *handler, time_t timeout) override
virtual XrdCl::XRootDStatus Rm(const std::string &path, XrdCl::ResponseHandler *handler, time_t timeout) override
virtual XrdCl::XRootDStatus Stat(const std::string &path, XrdCl::ResponseHandler *handler, time_t timeout) override
virtual XrdCl::XRootDStatus ThirdPartyCopy(const std::string &source, const std::string &dest, const XrdCl::PropertyList *properties, XrdCl::ProgressHandler *progress_handler, time_t timeout=0) override
virtual XrdCl::XRootDStatus Locate(const std::string &path, XrdCl::OpenFlags::Flags flags, XrdCl::ResponseHandler *handler, time_t timeout) override
virtual ~Filesystem() noexcept
virtual XrdCl::XRootDStatus RmDir(const std::string &path, XrdCl::ResponseHandler *handler, time_t timeout) override
virtual bool GetProperty(const std::string &name, std::string &value) const override
virtual XrdCl::XRootDStatus Query(XrdCl::QueryCode::Code queryCode, const XrdCl::Buffer &arg, XrdCl::ResponseHandler *handler, time_t timeout) override
Filesystem(const std::string &, std::shared_ptr< HandlerQueue > queue, XrdCl::Log *log)
void Produce(std::shared_ptr< CurlOperation > handler)
Binary blob representation.
std::string ToString() const
Convert the buffer to a string.
static Log * GetLog()
Get default log.
static Env * GetEnv()
Get default client environment.
@ Read
read access is allowed
@ ServerOnline
server node where the file is online
void Error(uint64_t topic, const char *format,...)
Report an error.
void Warning(uint64_t topic, const char *format,...)
Report a warning.
void Debug(uint64_t topic, const char *format,...)
Print a debug message.
Handle the progress of an asynchronous operation.
virtual void HandleProgress(std::size_t processed, std::size_t total=0)=0
A key-value pair map storing both keys and values as strings.
bool Get(const std::string &name, Item &item) const
Handle an async response.
virtual void HandleResponse(XRootDStatus *status, AnyObject *response)
std::map< std::string, std::string > ParamsMap
bool FromString(const std::string &url)
Parse a string and fill the URL fields.
std::string GetURL() const
Get the URL.
const ParamsMap & GetParams() const
Get the URL params.
const std::string & GetProtocol() const
Get the protocol.
ConnectionCallout *(*)(const std::string &, const ResponseInfo &) CreateConnCalloutType
ChecksumType GetTypeFromString(const std::string &str)
const uint64_t kLogXrdClHttp
static std::optional< std::size_t > stat_size(HandlerQueue &queue, const std::string &url, const CurlCopyOp::Headers &hdrs, time_t timeout, XrdCl::Log *log, const char *what)
bool ParseTokenFile(const std::string &token_file, CurlCopyOp::Headers &src_hdrs, CurlCopyOp::Headers &dst_hdrs)
const uint16_t errErrorResponse
const uint16_t errOperationExpired
const uint16_t errNotImplemented
Operation is not implemented.
const uint16_t stError
An error occurred that could potentially be retried.
const uint16_t errInternal
Internal error.
const uint16_t errOSError
const uint16_t errPipelineFailed
Pipeline failed and operation couldn't be executed.
const uint16_t errInvalidArgs
const uint16_t errAuthFailed
Flags
Open flags, may be or'd when appropriate.
Code
XRootD query request codes.
@ XAttr
Query file extended attributes.
@ Opaque
Implementation dependent.
@ Checksum
Query file checksum.
@ Prepare
Query prepare status.