1414
1515#include " google/cloud/storage/internal/hedged_object_read_source.h"
1616#include " google/cloud/internal/make_status.h"
17+ #ifdef GOOGLE_CLOUD_CPP_HAVE_OPENTELEMETRY
18+ #include < opentelemetry/metrics/meter.h>
19+ #include < opentelemetry/metrics/provider.h>
20+ #endif // GOOGLE_CLOUD_CPP_HAVE_OPENTELEMETRY
1721#include < atomic>
1822#include < cstring>
1923#include < future>
@@ -24,12 +28,44 @@ namespace cloud {
2428namespace storage {
2529GOOGLE_CLOUD_CPP_INLINE_NAMESPACE_BEGIN
2630namespace internal {
31+
32+ #ifdef GOOGLE_CLOUD_CPP_HAVE_OPENTELEMETRY
33+ HedgedReadMetrics::HedgedReadMetrics ()
34+ : HedgedReadMetrics(opentelemetry::metrics::Provider::GetMeterProvider()) {}
35+
36+ HedgedReadMetrics::HedgedReadMetrics (
37+ opentelemetry::nostd::shared_ptr<
38+ opentelemetry::metrics::MeterProvider> const & provider) {
39+ if (!provider) return ;
40+ opentelemetry::nostd::shared_ptr<opentelemetry::metrics::Meter> meter =
41+ provider->GetMeter (" gl-cpp" , version_string ());
42+ if (!meter) return ;
43+
44+ hedges_dispatched_ = meter->CreateUInt64Counter (
45+ " storage.read_hedging.hedges_dispatched" ,
46+ " Total number of speculative hedge read attempts dispatched" , " {hedge}" );
47+ hedge_won_ = meter->CreateUInt64Counter (
48+ " storage.read_hedging.hedge_won" ,
49+ " Total number of hedged read operations won by a secondary hedge attempt" ,
50+ " {request}" );
51+ }
52+
53+ void HedgedReadMetrics::IncrementHedgesDispatched () {
54+ if (hedges_dispatched_) hedges_dispatched_->Add (1 );
55+ }
56+
57+ void HedgedReadMetrics::IncrementHedgeWon () {
58+ if (hedge_won_) hedge_won_->Add (1 );
59+ }
60+ #endif // GOOGLE_CLOUD_CPP_HAVE_OPENTELEMETRY
61+
2762namespace {
2863
2964struct RaceResult {
3065 StatusOr<ReadSourceResult> result;
3166 std::unique_ptr<ObjectReadSource> source;
3267 std::unique_ptr<char []> buffer;
68+ bool is_primary;
3369};
3470
3571struct RaceState {
@@ -44,7 +80,8 @@ struct RaceState {
4480void RunAttempt (std::shared_ptr<RaceState> const & state,
4581 HedgedObjectReadSource::ChildFactory const & factory,
4682 std::size_t n, bool resolve_on_open_error,
47- std::shared_ptr<HedgingThreadPool> release_slot) {
83+ std::shared_ptr<HedgingThreadPool> release_slot,
84+ bool is_primary) {
4885 // Releases the acquired hedge concurrency slot upon function exit across
4986 // all code paths (early return on open/allocation error, race winner, or
5087 // race loser). For primary attempts, release_slot is nullptr.
@@ -55,13 +92,13 @@ void RunAttempt(std::shared_ptr<RaceState> const& state,
5592 }
5693 } guard{std::move (release_slot)};
5794
58- auto source = factory ();
95+ StatusOr<std::unique_ptr<ObjectReadSource>> source = factory ();
5996 if (!source) {
6097 if (!resolve_on_open_error) return ;
6198 bool expected = false ;
6299 if (state->resolved .compare_exchange_strong (expected, true )) {
63100 state->promise .set_value (
64- RaceResult{std::move (source).status (), nullptr , {}});
101+ RaceResult{std::move (source).status (), nullptr , {}, is_primary });
65102 }
66103 return ;
67104 }
@@ -74,32 +111,84 @@ void RunAttempt(std::shared_ptr<RaceState> const& state,
74111 google::cloud::internal::ResourceExhaustedError (
75112 " Out of memory allocating hedge buffer" , GCP_ERROR_INFO ()),
76113 nullptr ,
77- {}});
114+ {},
115+ is_primary});
78116 }
79117 return ;
80118 }
81- auto result = (*source)->Read (buffer.get (), n);
119+ StatusOr<ReadSourceResult> result = (*source)->Read (buffer.get (), n);
82120 bool expected = false ;
83121 if (state->resolved .compare_exchange_strong (expected, true )) {
84- state->promise .set_value (
85- RaceResult{ std::move (result), * std::move (source), std::move (buffer)});
122+ state->promise .set_value (RaceResult{ std::move (result), * std::move (source),
123+ std::move (buffer), is_primary });
86124 } else {
87125 (*source)->Close ();
88126 }
89127}
90128
129+ void DispatchHedges (std::shared_ptr<RaceState> const & state,
130+ std::future<RaceResult>& future, std::size_t n,
131+ int max_hedges, std::chrono::milliseconds delay,
132+ std::shared_ptr<HedgingThreadPool> const & hedge_pool,
133+ HedgedObjectReadSource::ChildFactory const & child_factory,
134+ HedgedReadMetrics& metrics) {
135+ for (int hedges_dispatched = 0 ; hedges_dispatched < max_hedges;) {
136+ if (future.wait_for (delay) != std::future_status::timeout) break ;
137+ if (!hedge_pool->TryAcquireHedgeToken ()) {
138+ // When delay is 0ms (or token acquisition fails), back off briefly on
139+ // the future instead of busy-spinning if tokens or concurrency slots are
140+ // temporarily exhausted.
141+ if (delay == std::chrono::milliseconds::zero () &&
142+ future.wait_for (std::chrono::milliseconds (10 )) !=
143+ std::future_status::timeout) {
144+ break ;
145+ }
146+ continue ;
147+ }
148+ auto hedge = [state, factory = child_factory, n, pool = hedge_pool] {
149+ RunAttempt (state, factory, n, /* resolve_on_open_error=*/ false , pool,
150+ /* is_primary=*/ false );
151+ };
152+ if (!hedge_pool->Enqueue (hedge)) {
153+ hedge_pool->ReleaseHedgeSlot ();
154+ break ;
155+ }
156+ ++hedges_dispatched;
157+ metrics.IncrementHedgesDispatched ();
158+ }
159+ }
160+
91161} // namespace
92162
93163HedgedObjectReadSource::HedgedObjectReadSource (
94164 std::shared_ptr<ThreadPool> read_pool,
95165 std::shared_ptr<HedgingThreadPool> hedge_pool, ChildFactory child_factory,
96166 std::chrono::milliseconds delay, int max_hedges, std::size_t max_buffer)
167+ : HedgedObjectReadSource(std::move(read_pool), std::move(hedge_pool),
168+ std::move (child_factory), delay, max_hedges,
169+ max_buffer,
170+ std::make_shared<HedgedReadMetrics>()) {}
171+
172+ HedgedObjectReadSource::HedgedObjectReadSource (
173+ std::shared_ptr<ThreadPool> read_pool,
174+ std::shared_ptr<HedgingThreadPool> hedge_pool, ChildFactory child_factory,
175+ std::chrono::milliseconds delay, int max_hedges, std::size_t max_buffer,
176+ std::shared_ptr<HedgedReadMetrics> metrics)
97177 : read_pool_(std::move(read_pool)),
98178 hedge_pool_(std::move(hedge_pool)),
99179 child_factory_(std::move(child_factory)),
100180 delay_(delay),
101181 max_hedges_(max_hedges),
102- max_buffer_(max_buffer) {}
182+ max_buffer_(max_buffer),
183+ #ifdef GOOGLE_CLOUD_CPP_HAVE_OPENTELEMETRY
184+ metrics_ (metrics ? std::move(metrics)
185+ : std::make_shared<HedgedReadMetrics>(nullptr )) {
186+ }
187+ #else
188+ metrics_ (metrics ? std::move(metrics)
189+ : std::make_shared<HedgedReadMetrics>()) {
190+ }
191+ #endif
103192
104193bool HedgedObjectReadSource::IsOpen () const {
105194 if (active_child_) return active_child_->IsOpen ();
@@ -130,48 +219,30 @@ StatusOr<ReadSourceResult> HedgedObjectReadSource::Read(char* buf,
130219 // the tail latency it avoids, so open the stream without hedging and read
131220 // straight into the caller's buffer.
132221 if (n > max_buffer_) {
133- auto child = child_factory_ ();
222+ StatusOr<std::unique_ptr<ObjectReadSource>> child = child_factory_ ();
134223 if (!child) return std::move (child).status ();
135224 active_child_ = *std::move (child);
136225 return active_child_->Read (buf, n);
137226 }
138227
139228 auto state = std::make_shared<RaceState>();
140- auto future = state->promise .get_future ();
229+ std::future<RaceResult> future = state->promise .get_future ();
141230
142231 auto primary = [state, factory = child_factory_, n] {
143- RunAttempt (state, factory, n, /* resolve_on_open_error=*/ true , nullptr );
232+ RunAttempt (state, factory, n, /* resolve_on_open_error=*/ true , nullptr ,
233+ /* is_primary=*/ true );
144234 };
145235 // The primary attempt is scheduled on the dedicated read pool.
146236 // If the pool is shutting down run the attempt inline, the read must
147237 // complete either way.
148238 if (!read_pool_->Enqueue (primary)) primary ();
149239
150- for (int hedges_dispatched = 0 ; hedges_dispatched < max_hedges_;) {
151- if (future.wait_for (delay_) != std::future_status::timeout) break ;
152- if (!hedge_pool_->TryAcquireHedgeToken ()) {
153- // When delay_ is 0ms (or token acquisition fails), back off briefly on
154- // the future instead of busy-spinning if tokens or concurrency slots are
155- // temporarily exhausted.
156- if (delay_ == std::chrono::milliseconds::zero ()) {
157- if (future.wait_for (std::chrono::milliseconds (10 )) !=
158- std::future_status::timeout) {
159- break ;
160- }
161- }
162- continue ;
163- }
164- auto hedge = [state, factory = child_factory_, n, pool = hedge_pool_] {
165- RunAttempt (state, factory, n, /* resolve_on_open_error=*/ false , pool);
166- };
167- if (!hedge_pool_->Enqueue (hedge)) {
168- hedge_pool_->ReleaseHedgeSlot ();
169- break ;
170- }
171- ++hedges_dispatched;
172- }
240+ DispatchHedges (state, future, n, max_hedges_, delay_, hedge_pool_,
241+ child_factory_, *metrics_);
242+
243+ RaceResult race = future.get ();
244+ if (!race.is_primary ) metrics_->IncrementHedgeWon ();
173245
174- auto race = future.get ();
175246 active_child_ = std::move (race.source );
176247 if (race.result .ok () && race.result ->bytes_received > 0 ) {
177248 std::memcpy (buf, race.buffer .get (), race.result ->bytes_received );
0 commit comments