48 m_default_put_handler(new PutDefaultHandler(*this))
51 virtual ~File() noexcept;
53 virtual
XrdCl::XRootDStatus
Open(const std::
string &url,
56 XrdCl::ResponseHandler *handler,
57 time_t timeout) override;
60 time_t timeout) override;
62 virtual
XrdCl::XRootDStatus
Stat(
bool force,
63 XrdCl::ResponseHandler *handler,
64 time_t timeout) override;
67 XrdCl::ResponseHandler *handler,
68 time_t timeout) override;
70 virtual
XrdCl::XRootDStatus
Read(uint64_t offset,
73 XrdCl::ResponseHandler *handler,
74 time_t timeout) override;
79 XrdCl::ResponseHandler *handler,
80 time_t timeout) override;
84 XrdCl::ResponseHandler *handler,
85 time_t timeout ) override;
87 virtual
XrdCl::XRootDStatus
Write(uint64_t offset,
90 XrdCl::ResponseHandler *handler,
91 time_t timeout) override;
93 virtual
XrdCl::XRootDStatus
Write(uint64_t offset,
94 XrdCl::Buffer &&buffer,
95 XrdCl::ResponseHandler *handler,
96 time_t timeout) override;
98 virtual
bool IsOpen() const override;
101 const std::
string &value ) override;
104 std::
string &value ) const override;
110 static void SetMinimumHeaderTimeout(
struct timespec &ts) {m_min_client_timeout.tv_sec = ts.tv_sec; m_min_client_timeout.tv_nsec = ts.tv_nsec;}
116 static void SetDefaultHeaderTimeout(
struct timespec &ts) {m_default_header_timeout.tv_sec = ts.tv_sec; m_default_header_timeout.tv_nsec = ts.tv_nsec;}
122 void SetHeaderTimeout(
const struct timespec &ts) {m_header_timeout.tv_sec = ts.tv_sec; m_header_timeout.tv_nsec = ts.tv_nsec;}
149 std::tuple<XrdCl::XRootDStatus, bool> ReadPrefetch(uint64_t offset, uint64_t size,
void *buffer,
XrdCl::ResponseHandler *handler, time_t timeout,
bool isPgRead);
154 bool SendResponseInfo()
const;
165 const std::string GetCurrentURL()
const;
170 void CalculateCurrentURL(
const std::string &value)
const;
172 bool m_is_opened{
false};
173 std::atomic<bool> m_full_download{
false};
179 std::string m_last_url;
180 mutable std::string m_url_current;
181 std::shared_ptr<XrdClHttp::HandlerQueue> m_queue;
182 XrdCl::Log *m_logger{
nullptr};
183 std::unordered_map<std::string, std::string> m_properties;
186 mutable std::shared_mutex m_properties_mutex;
189 struct timespec m_timeout{0, 0};
192 static struct timespec m_min_client_timeout;
195 static struct timespec m_default_header_timeout;
198 struct timespec m_header_timeout;
201 static struct timespec m_fed_timeout;
208 std::shared_ptr<XrdClHttp::CurlPutOp> m_put_op;
211 class PutResponseHandler :
public XrdCl::ResponseHandler {
213 PutResponseHandler(XrdCl::ResponseHandler *handler);
215 virtual void HandleResponse(XrdCl::XRootDStatus *status_raw, XrdCl::AnyObject *response_raw)
override;
224 XrdCl::XRootDStatus QueueWrite(std::variant<std::pair<const void *, size_t>, XrdCl::Buffer> buffer,
225 XrdCl::ResponseHandler *handler,
struct timespec timeout);
227 void SetOp(std::shared_ptr<XrdClHttp::CurlPutOp> op) {m_op = op;}
229 void WaitForCompletion();
233 bool m_initial{
true};
234 std::shared_ptr<XrdClHttp::CurlPutOp> m_op;
235 XrdCl::ResponseHandler *m_active_handler;
236 std::condition_variable m_cv;
238 std::deque<std::tuple<std::variant<std::pair<const void *, size_t>, XrdCl::Buffer>, XrdCl::ResponseHandler*,
struct timespec>> m_pending_writes;
248 std::atomic<PutResponseHandler *>m_put_handler{
nullptr};
255 class PutDefaultHandler :
public XrdCl::ResponseHandler {
257 PutDefaultHandler(
File &file) : m_logger(file.m_logger) {}
259 virtual void HandleResponse(XrdCl::XRootDStatus *status, XrdCl::AnyObject *response);
262 XrdCl::Log *m_logger{
nullptr};
266 std::shared_ptr<PutDefaultHandler> m_default_put_handler;
274 std::shared_ptr<XrdClHttp::CurlReadOp> m_prefetch_op;
278 std::atomic<off_t> m_prefetch_offset{0};
285 class PrefetchResponseHandler :
public XrdCl::ResponseHandler {
307 PrefetchResponseHandler(
File &parent,
308 off_t offset,
size_t size, std::atomic<off_t> *prefetch_offset,
char *buffer, XrdCl::ResponseHandler *handler,
309 std::unique_lock<std::mutex> *lock, time_t timeout);
311 virtual void HandleResponse(XrdCl::XRootDStatus *status, XrdCl::AnyObject *response);
315 void ResubmitOperation();
322 XrdCl::ResponseHandler *m_handler;
325 PrefetchResponseHandler *m_next{
nullptr};
328 char *m_buffer{
nullptr};
338 std::atomic<off_t> *m_prefetch_offset{
nullptr};
346 PrefetchResponseHandler *m_last_prefetch_handler{
nullptr};
349 off_t m_prefetch_size{-1};
352 std::atomic<off_t> m_put_offset{0};
360 class PrefetchDefaultHandler :
public XrdCl::ResponseHandler {
362 PrefetchDefaultHandler(
File &file) : m_logger(file.m_logger), m_url(file.m_url) {}
364 virtual void HandleResponse(XrdCl::XRootDStatus *status, XrdCl::AnyObject *response);
367 void DisablePrefetch() {
368 auto enabled = m_prefetch_enabled.load(std::memory_order_relaxed);
370 std::unique_lock lock(m_prefetch_mutex);
371 m_prefetch_enabled.store(
false, std::memory_order_relaxed);
376 bool IsPrefetching()
const {
377 auto enabled = m_prefetch_enabled.load(std::memory_order_relaxed);
379 std::unique_lock lock(m_prefetch_mutex);
380 return m_prefetch_enabled.load(std::memory_order_relaxed);
385 XrdCl::Log *m_logger{
nullptr};
390 mutable std::mutex m_prefetch_mutex;
398 mutable std::atomic<bool> m_prefetch_enabled{
true};
405 std::shared_ptr<PrefetchDefaultHandler> m_default_prefetch_handler;
408 std::atomic<XrdClHttp::HeaderCallout *> m_header_callout{
nullptr};
411 class HeaderCallout :
public XrdClHttp::HeaderCallout {
413 HeaderCallout(
File &fs) : m_parent(fs)
416 virtual ~HeaderCallout() noexcept = default;
418 virtual std::shared_ptr<
HeaderList> GetHeaders(const std::
string &verb,
419 const std::
string &url,
426 HeaderCallout m_default_header_callout{*
this};
428 static std::atomic<uint64_t> m_prefetch_count;
429 static std::atomic<uint64_t> m_prefetch_expired_count;
430 static std::atomic<uint64_t> m_prefetch_failed_count;
431 static std::atomic<uint64_t> m_prefetch_reads_hit;
432 static std::atomic<uint64_t> m_prefetch_reads_miss;
433 static std::atomic<uint64_t> m_prefetch_bytes_used;