1414
1515#include " google/cloud/storage/internal/hedged_object_read_source.h"
1616#include " google/cloud/internal/make_status.h"
17+ #include < opentelemetry/metrics/provider.h>
1718#include < atomic>
1819#include < cstring>
1920#include < future>
@@ -24,12 +25,45 @@ namespace cloud {
2425namespace storage {
2526GOOGLE_CLOUD_CPP_INLINE_NAMESPACE_BEGIN
2627namespace internal {
28+
29+ HedgedReadMetrics::HedgedReadMetrics ()
30+ : HedgedReadMetrics(opentelemetry::metrics::Provider::GetMeterProvider()) {}
31+
32+ HedgedReadMetrics::HedgedReadMetrics (
33+ opentelemetry::nostd::shared_ptr<
34+ opentelemetry::metrics::MeterProvider> const & provider) {
35+ if (!provider) return ;
36+ auto meter = provider->GetMeter (" gl-cpp" , version_string ());
37+ if (!meter) return ;
38+
39+ hedges_dispatched_ = meter->CreateUInt64Counter (
40+ " storage.read_hedging.hedges_dispatched" ,
41+ " Total number of speculative hedge read attempts dispatched" , " {hedge}" );
42+ hedge_won_ = meter->CreateUInt64Counter (
43+ " storage.read_hedging.hedge_won" ,
44+ " Total number of hedged read operations won by a secondary hedge attempt" ,
45+ " {request}" );
46+ }
47+
48+ void HedgedReadMetrics::IncrementHedgesDispatched () {
49+ IncrementHedgesDispatched (1 );
50+ }
51+ void HedgedReadMetrics::IncrementHedgesDispatched (std::uint64_t count) {
52+ if (hedges_dispatched_) hedges_dispatched_->Add (count);
53+ }
54+
55+ void HedgedReadMetrics::IncrementHedgeWon () { IncrementHedgeWon (1 ); }
56+ void HedgedReadMetrics::IncrementHedgeWon (std::uint64_t count) {
57+ if (hedge_won_) hedge_won_->Add (count);
58+ }
59+
2760namespace {
2861
2962struct RaceResult {
3063 StatusOr<ReadSourceResult> result;
3164 std::unique_ptr<ObjectReadSource> source;
3265 std::unique_ptr<char []> buffer;
66+ bool is_primary;
3367};
3468
3569struct RaceState {
@@ -44,7 +78,8 @@ struct RaceState {
4478void RunAttempt (std::shared_ptr<RaceState> const & state,
4579 HedgedObjectReadSource::ChildFactory const & factory,
4680 std::size_t n, bool resolve_on_open_error,
47- std::shared_ptr<HedgingThreadPool> release_slot) {
81+ std::shared_ptr<HedgingThreadPool> release_slot,
82+ bool is_primary) {
4883 // Releases the acquired hedge concurrency slot upon function exit across
4984 // all code paths (early return on open/allocation error, race winner, or
5085 // race loser). For primary attempts, release_slot is nullptr.
@@ -55,13 +90,13 @@ void RunAttempt(std::shared_ptr<RaceState> const& state,
5590 }
5691 } guard{std::move (release_slot)};
5792
58- auto source = factory ();
93+ StatusOr<std::unique_ptr<ObjectReadSource>> source = factory ();
5994 if (!source) {
6095 if (!resolve_on_open_error) return ;
6196 bool expected = false ;
6297 if (state->resolved .compare_exchange_strong (expected, true )) {
6398 state->promise .set_value (
64- RaceResult{std::move (source).status (), nullptr , {}});
99+ RaceResult{std::move (source).status (), nullptr , {}, is_primary });
65100 }
66101 return ;
67102 }
@@ -74,15 +109,16 @@ void RunAttempt(std::shared_ptr<RaceState> const& state,
74109 google::cloud::internal::ResourceExhaustedError (
75110 " Out of memory allocating hedge buffer" , GCP_ERROR_INFO ()),
76111 nullptr ,
77- {}});
112+ {},
113+ is_primary});
78114 }
79115 return ;
80116 }
81- auto result = (*source)->Read (buffer.get (), n);
117+ StatusOr<ReadSourceResult> result = (*source)->Read (buffer.get (), n);
82118 bool expected = false ;
83119 if (state->resolved .compare_exchange_strong (expected, true )) {
84- state->promise .set_value (
85- RaceResult{ std::move (result), *std::move (source), std::move (buffer)});
120+ state->promise .set_value (RaceResult{
121+ std::move (result), *std::move (source), std::move (buffer), is_primary });
86122 } else {
87123 (*source)->Close ();
88124 }
@@ -94,12 +130,23 @@ HedgedObjectReadSource::HedgedObjectReadSource(
94130 std::shared_ptr<ThreadPool> read_pool,
95131 std::shared_ptr<HedgingThreadPool> hedge_pool, ChildFactory child_factory,
96132 std::chrono::milliseconds delay, int max_hedges, std::size_t max_buffer)
133+ : HedgedObjectReadSource(std::move(read_pool), std::move(hedge_pool),
134+ std::move (child_factory), delay, max_hedges,
135+ max_buffer,
136+ std::make_shared<HedgedReadMetrics>()) {}
137+
138+ HedgedObjectReadSource::HedgedObjectReadSource (
139+ std::shared_ptr<ThreadPool> read_pool,
140+ std::shared_ptr<HedgingThreadPool> hedge_pool, ChildFactory child_factory,
141+ std::chrono::milliseconds delay, int max_hedges, std::size_t max_buffer,
142+ std::shared_ptr<HedgedReadMetrics> metrics)
97143 : read_pool_(std::move(read_pool)),
98144 hedge_pool_(std::move(hedge_pool)),
99145 child_factory_(std::move(child_factory)),
100146 delay_(delay),
101147 max_hedges_(max_hedges),
102- max_buffer_(max_buffer) {}
148+ max_buffer_(max_buffer),
149+ metrics_(std::move(metrics)) {}
103150
104151bool HedgedObjectReadSource::IsOpen () const {
105152 if (active_child_) return active_child_->IsOpen ();
@@ -130,24 +177,26 @@ StatusOr<ReadSourceResult> HedgedObjectReadSource::Read(char* buf,
130177 // the tail latency it avoids, so open the stream without hedging and read
131178 // straight into the caller's buffer.
132179 if (n > max_buffer_) {
133- auto child = child_factory_ ();
180+ StatusOr<std::unique_ptr<ObjectReadSource>> child = child_factory_ ();
134181 if (!child) return std::move (child).status ();
135182 active_child_ = *std::move (child);
136183 return active_child_->Read (buf, n);
137184 }
138185
139186 auto state = std::make_shared<RaceState>();
140- auto future = state->promise .get_future ();
187+ std::future<RaceResult> future = state->promise .get_future ();
141188
142189 auto primary = [state, factory = child_factory_, n] {
143- RunAttempt (state, factory, n, /* resolve_on_open_error=*/ true , nullptr );
190+ RunAttempt (state, factory, n, /* resolve_on_open_error=*/ true , nullptr ,
191+ /* is_primary=*/ true );
144192 };
145193 // The primary attempt is scheduled on the dedicated read pool.
146194 // If the pool is shutting down run the attempt inline, the read must
147195 // complete either way.
148196 if (!read_pool_->Enqueue (primary)) primary ();
149197
150- for (int hedges_dispatched = 0 ; hedges_dispatched < max_hedges_;) {
198+ int hedges_dispatched = 0 ;
199+ for (; hedges_dispatched < max_hedges_;) {
151200 if (future.wait_for (delay_) != std::future_status::timeout) break ;
152201 if (!hedge_pool_->TryAcquireHedgeToken ()) {
153202 // When delay_ is 0ms (or token acquisition fails), back off briefly on
@@ -162,16 +211,22 @@ StatusOr<ReadSourceResult> HedgedObjectReadSource::Read(char* buf,
162211 continue ;
163212 }
164213 auto hedge = [state, factory = child_factory_, n, pool = hedge_pool_] {
165- RunAttempt (state, factory, n, /* resolve_on_open_error=*/ false , pool);
214+ RunAttempt (state, factory, n, /* resolve_on_open_error=*/ false , pool,
215+ /* is_primary=*/ false );
166216 };
167217 if (!hedge_pool_->Enqueue (hedge)) {
168218 hedge_pool_->ReleaseHedgeSlot ();
169219 break ;
170220 }
171221 ++hedges_dispatched;
222+ if (metrics_) metrics_->IncrementHedgesDispatched ();
223+ }
224+
225+ RaceResult race = future.get ();
226+ if (metrics_ && !race.is_primary ) {
227+ metrics_->IncrementHedgeWon ();
172228 }
173229
174- auto race = future.get ();
175230 active_child_ = std::move (race.source );
176231 if (race.result .ok () && race.result ->bytes_received > 0 ) {
177232 std::memcpy (buf, race.buffer .get (), race.result ->bytes_received );
0 commit comments