XRootD
Loading...
Searching...
No Matches
XrdClHttpWorker.hh
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#ifndef XRDCLHTTPWORKER_HH
22#define XRDCLHTTPWORKER_HH
23
24#include "XrdClHttpOps.hh"
25
26#include <array>
27#include <atomic>
28#include <chrono>
29#include <condition_variable>
30#include <memory>
31#include <mutex>
32#include <tuple>
33#include <unordered_map>
34#include <unordered_set>
35
36typedef void CURL;
37
38namespace XrdCl {
39
40class Env;
41class Log;
42class ResponseHandler;
43class URL;
44
45}
46
47namespace XrdClHttp {
48
49class HandlerQueue;
50class VerbsCache;
51
53public:
54 CurlWorker(std::shared_ptr<HandlerQueue> queue, VerbsCache &cache, XrdCl::Log* logger);
55
56 CurlWorker(const CurlWorker &) = delete;
57
58 void Run();
59 static void RunStatic(CurlWorker *myself);
60
61 // Passes some initial values to the worker so it can start
62 void Start(std::unique_ptr<XrdClHttp::CurlWorker> self, std::thread tid);
63
64 // Returns the configured X509 client certificate and key file name
65 std::tuple<std::string, std::string> ClientX509CertKeyFile() const;
66
67 // Change the period (in seconds) for queue maintenance.
68 //
69 // Defaults to 5 seconds; smaller values are convenient for unit tests.
70 static void SetMaintenancePeriod(unsigned maint) {
71 m_maintenance_period.store(maint, std::memory_order_relaxed);
72 }
73
74 static std::string GetMonitoringJson();
75
76private:
77 // Invoked by the destructor of one of our static members. This triggers when
78 // the plugin is unloaded, triggers the shutdown of each of the worker threads.
79 static void ShutdownAll();
80
81 // Invoked by ShutdownAll, kills off the current object's thread
82 void Shutdown();
83
84 // A list of all known worker threads -- used to shutdown the process
85 static std::vector<std::unique_ptr<XrdClHttp::CurlWorker>> m_workers;
86 // Protects the data in m_workers
87 static std::mutex m_workers_mutex;
88
89 std::chrono::steady_clock::time_point m_last_prefix_log;
90 VerbsCache &m_cache; // Cache mapping server URLs to list of selected HTTP verbs.
91 std::shared_ptr<HandlerQueue> m_queue;
92
93 // Queue for operations that can be unpaused.
94 // Paused operations occur when a PUT is started but cannot be continued
95 // because more data is needed from the caller.
96 std::shared_ptr<HandlerQueue> m_continue_queue;
97
98 std::unordered_map<CURL*, std::pair<std::shared_ptr<CurlOperation>, std::chrono::system_clock::time_point>> m_op_map;
99 XrdCl::Log* m_logger;
100 std::string m_x509_client_cert_file;
101 std::string m_x509_client_key_file;
102
103 const static unsigned m_max_ops{20};
104 static std::atomic<unsigned> m_maintenance_period;
105
106 // File descriptor pair indicating shutdown is requested.
107 int m_shutdown_pipe_r{-1};
108 int m_shutdown_pipe_w{-1};
109 // Mutex for managing the startup of a worker
110 std::mutex m_start_lock;
111 // Condition variable for a worker to indicate to RunStatic that it is ready
112 std::condition_variable m_start_complete_cv;
113 // Flag indicating that Start has been called.
114 bool m_start_complete{false};
115 // The worker's thread object
116 std::thread m_self_tid;
117
118 // Monitoring statistics
119 struct OpStats {
120 std::atomic<uint64_t> m_conncall_timeout{}; // Timeout due to the connection callout mechanism
121 std::atomic<uint64_t> m_client_timeout{};
122 std::atomic<std::chrono::system_clock::duration::rep> m_duration{};
123 std::atomic<uint64_t> m_error{};
124 std::atomic<uint64_t> m_finished{};
125 std::atomic<std::chrono::steady_clock::duration::rep> m_pause_duration{};
126 std::atomic<uint64_t> m_started{};
127 std::atomic<uint64_t> m_server_timeout{};
128 std::atomic<uint64_t> m_bytes{};
129 };
130
131 enum class OpKind {
132 ConncallTimeout,
133 ClientTimeout,
134 Error,
135 Finish,
136 Start,
137 ServerTimeout,
138 Update
139 };
140 void OpRecord(XrdClHttp::CurlOperation &op, OpKind);
141
142 static std::atomic<uint64_t> m_conncall_errors;
143 static std::atomic<uint64_t> m_conncall_req;
144 static std::atomic<uint64_t> m_conncall_success;
145 static std::atomic<uint64_t> m_conncall_timeout;
146 static std::array<std::array<OpStats, 403>, static_cast<size_t>(XrdClHttp::CurlOperation::HttpVerb::Count)> m_ops;
147 std::atomic<std::chrono::system_clock::rep> m_last_completed_cycle;
148 std::atomic<std::chrono::system_clock::rep> m_oldest_op;
149
150 // Vector tracking known worker statistics.
151 static std::vector<std::atomic<std::chrono::system_clock::rep>*> m_workers_last_completed_cycle;
152 static std::vector<std::atomic<std::chrono::system_clock::rep>*> m_workers_oldest_op;
153 size_t m_stats_offset{0};
154 static std::mutex m_worker_stats_mutex;
155
156 // shutdown trigger
157 static struct initcontrol {
158 initcontrol();
159 ~initcontrol();
160 } m_initcontrol;
161};
162
163}
164
165#endif // XRDCLHTTPWORKER_HH
void CURL
std::tuple< std::string, std::string > ClientX509CertKeyFile() const
CurlWorker(std::shared_ptr< HandlerQueue > queue, VerbsCache &cache, XrdCl::Log *logger)
CurlWorker(const CurlWorker &)=delete
static void RunStatic(CurlWorker *myself)
void Start(std::unique_ptr< XrdClHttp::CurlWorker > self, std::thread tid)
static void SetMaintenancePeriod(unsigned maint)
static std::string GetMonitoringJson()
Handle diagnostics.
Definition XrdClLog.hh:101
Handle an async response.
URL representation.
Definition XrdClURL.hh:31