Skip to content

Commit 36b0bc5

Browse files
committed
prom remote write: add configurable pool idle timeout
Signed-off-by: Matt Klein <mklein@bitdrift.io>
1 parent f77a606 commit 36b0bc5

9 files changed

Lines changed: 63 additions & 27 deletions

File tree

pulse-metrics/src/admin/server.rs

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -344,6 +344,7 @@ impl MetaStatsEmitter {
344344
prom_remote_write.auth.into_option(),
345345
PROM_REMOTE_WRITE_HEADERS,
346346
prom_remote_write.request_headers,
347+
None,
347348
)
348349
.await?,
349350
}) as Box<dyn MetaStatsSender>,

pulse-metrics/src/clients/http.rs

Lines changed: 4 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -168,10 +168,13 @@ impl HyperHttpRemoteWriteClient {
168168
auth_config: Option<HttpRemoteWriteAuthConfig>,
169169
core_request_headers: &[(&str, &str)],
170170
config_request_headers: Vec<RequestHeader>,
171+
pool_idle_timeout: Option<Duration>,
171172
) -> Result<Self> {
172173
Ok(Self {
173174
// TODO(mattklein123): Make connect timeout configurable.
174-
inner: Client::builder(TokioExecutor::new()).build(make_tls_connector(250.milliseconds())),
175+
inner: Client::builder(TokioExecutor::new())
176+
.pool_idle_timeout(pool_idle_timeout.unwrap_or_else(|| 90.seconds()).unsigned_abs())
177+
.build(make_tls_connector(250.milliseconds())),
175178
endpoint,
176179
timeout,
177180
auth: Self::create_auth(auth_config).await?,

pulse-metrics/src/pipeline/outflow/http/remote_write.rs

Lines changed: 5 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -24,7 +24,7 @@ use backoff::backoff::Backoff;
2424
use bd_log::warn_every;
2525
use bd_server_stats::stats::{AutoGauge, Scope};
2626
use bd_shutdown::{ComponentShutdown, ComponentStatus};
27-
use bd_time::TimeDurationExt;
27+
use bd_time::{ProtoDurationExt, TimeDurationExt};
2828
use bytes::Bytes;
2929
use http::HeaderMap;
3030
use prometheus::{Histogram, IntCounter, IntGauge};
@@ -124,6 +124,7 @@ impl HttpRemoteWriteOutflow {
124124
core_request_headers: &[(&str, &str)],
125125
config_request_headers: Vec<RequestHeader>,
126126
context: OutflowFactoryContext,
127+
pool_idle_timeout: MessageField<Duration>,
127128
) -> anyhow::Result<Arc<Self>> {
128129
let request_timeout = request_timeout.unwrap_duration_or(DEFAULT_REQUEST_TIMEOUT);
129130
let client = Arc::new(
@@ -133,6 +134,9 @@ impl HttpRemoteWriteOutflow {
133134
auth_config,
134135
core_request_headers,
135136
config_request_headers,
137+
pool_idle_timeout
138+
.map(|d| d.to_time_duration())
139+
.into_option(),
136140
)
137141
.await?,
138142
);

pulse-metrics/src/pipeline/outflow/otlp/mod.rs

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -88,6 +88,7 @@ pub async fn make_otlp_outflow(
8888
&[],
8989
config.request_headers,
9090
context,
91+
None.into(),
9192
)
9293
.await
9394
}

pulse-metrics/src/pipeline/outflow/prom/mod.rs

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -96,6 +96,7 @@ pub async fn make_prom_outflow(
9696
PROM_REMOTE_WRITE_HEADERS,
9797
config.request_headers,
9898
context,
99+
config.pool_idle_timeout,
99100
)
100101
.await
101102
}

pulse-promcli/src/main.rs

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -107,6 +107,7 @@ async fn main() -> anyhow::Result<()> {
107107
None,
108108
PROM_REMOTE_WRITE_HEADERS,
109109
vec![],
110+
None,
110111
)
111112
.await?;
112113
client

pulse-protobuf/proto/pulse/config/outflow/v1/prom_remote_write.proto

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -64,4 +64,7 @@ message PromRemoteWriteClientConfig {
6464

6565
// Lyft specific configuration. To be generalized later.
6666
LyftSpecificConfig lyft_specific_config = 11;
67+
68+
// Idle time for outbound connections in the connection pool. If not set, defaults to 90s.
69+
google.protobuf.Duration pool_idle_timeout = 12;
6770
}

pulse-protobuf/src/protos/pulse/config/outflow/v1/prom_remote_write.rs

Lines changed: 45 additions & 25 deletions
Original file line numberDiff line numberDiff line change
@@ -57,6 +57,8 @@ pub struct PromRemoteWriteClientConfig {
5757
pub retry_policy: ::protobuf::MessageField<super::retry::RetryPolicy>,
5858
// @@protoc_insertion_point(field:pulse.config.outflow.v1.PromRemoteWriteClientConfig.lyft_specific_config)
5959
pub lyft_specific_config: ::protobuf::MessageField<prom_remote_write_client_config::LyftSpecificConfig>,
60+
// @@protoc_insertion_point(field:pulse.config.outflow.v1.PromRemoteWriteClientConfig.pool_idle_timeout)
61+
pub pool_idle_timeout: ::protobuf::MessageField<::protobuf::well_known_types::duration::Duration>,
6062
// special fields
6163
// @@protoc_insertion_point(special_field:pulse.config.outflow.v1.PromRemoteWriteClientConfig.special_fields)
6264
pub special_fields: ::protobuf::SpecialFields,
@@ -74,7 +76,7 @@ impl PromRemoteWriteClientConfig {
7476
}
7577

7678
fn generated_message_descriptor_data() -> ::protobuf::reflect::GeneratedMessageDescriptorData {
77-
let mut fields = ::std::vec::Vec::with_capacity(11);
79+
let mut fields = ::std::vec::Vec::with_capacity(12);
7880
let mut oneofs = ::std::vec::Vec::with_capacity(0);
7981
fields.push(::protobuf::reflect::rt::v2::make_simpler_field_accessor::<_, _>(
8082
"send_to",
@@ -131,6 +133,11 @@ impl PromRemoteWriteClientConfig {
131133
|m: &PromRemoteWriteClientConfig| { &m.lyft_specific_config },
132134
|m: &mut PromRemoteWriteClientConfig| { &mut m.lyft_specific_config },
133135
));
136+
fields.push(::protobuf::reflect::rt::v2::make_message_field_accessor::<_, ::protobuf::well_known_types::duration::Duration>(
137+
"pool_idle_timeout",
138+
|m: &PromRemoteWriteClientConfig| { &m.pool_idle_timeout },
139+
|m: &mut PromRemoteWriteClientConfig| { &mut m.pool_idle_timeout },
140+
));
134141
::protobuf::reflect::GeneratedMessageDescriptorData::new_2::<PromRemoteWriteClientConfig>(
135142
"PromRemoteWriteClientConfig",
136143
fields,
@@ -182,6 +189,9 @@ impl ::protobuf::Message for PromRemoteWriteClientConfig {
182189
90 => {
183190
::protobuf::rt::read_singular_message_into_field(is, &mut self.lyft_specific_config)?;
184191
},
192+
98 => {
193+
::protobuf::rt::read_singular_message_into_field(is, &mut self.pool_idle_timeout)?;
194+
},
185195
tag => {
186196
::protobuf::rt::read_unknown_or_skip_group(tag, is, self.special_fields.mut_unknown_fields())?;
187197
},
@@ -233,6 +243,10 @@ impl ::protobuf::Message for PromRemoteWriteClientConfig {
233243
let len = v.compute_size();
234244
my_size += 1 + ::protobuf::rt::compute_raw_varint64_size(len) + len;
235245
}
246+
if let Some(v) = self.pool_idle_timeout.as_ref() {
247+
let len = v.compute_size();
248+
my_size += 1 + ::protobuf::rt::compute_raw_varint64_size(len) + len;
249+
}
236250
my_size += ::protobuf::rt::unknown_fields_size(self.special_fields.unknown_fields());
237251
self.special_fields.cached_size().set(my_size as u32);
238252
my_size
@@ -272,6 +286,9 @@ impl ::protobuf::Message for PromRemoteWriteClientConfig {
272286
if let Some(v) = self.lyft_specific_config.as_ref() {
273287
::protobuf::rt::write_message_field_with_cached_size(11, v, os)?;
274288
}
289+
if let Some(v) = self.pool_idle_timeout.as_ref() {
290+
::protobuf::rt::write_message_field_with_cached_size(12, v, os)?;
291+
}
275292
os.write_unknown_fields(self.special_fields.unknown_fields())?;
276293
::std::result::Result::Ok(())
277294
}
@@ -300,6 +317,7 @@ impl ::protobuf::Message for PromRemoteWriteClientConfig {
300317
self.convert_metric_name = ::std::option::Option::None;
301318
self.retry_policy.clear();
302319
self.lyft_specific_config.clear();
320+
self.pool_idle_timeout.clear();
303321
self.special_fields.clear();
304322
}
305323

@@ -316,6 +334,7 @@ impl ::protobuf::Message for PromRemoteWriteClientConfig {
316334
convert_metric_name: ::std::option::Option::None,
317335
retry_policy: ::protobuf::MessageField::none(),
318336
lyft_specific_config: ::protobuf::MessageField::none(),
337+
pool_idle_timeout: ::protobuf::MessageField::none(),
319338
special_fields: ::protobuf::SpecialFields::new(),
320339
};
321340
&instance
@@ -505,30 +524,31 @@ static file_descriptor_proto_data: &'static [u8] = b"\
505524
utflow.v1\x1a#pulse/config/common/v1/common.proto\x1a\"pulse/config/comm\
506525
on/v1/retry.proto\x1a*pulse/config/outflow/v1/queue_policy.proto\x1a,pul\
507526
se/config/outflow/v1/outflow_common.proto\x1a\x1egoogle/protobuf/duratio\
508-
n.proto\x1a\x17validate/validate.proto\"\xe2\x08\n\x1bPromRemoteWriteCli\
509-
entConfig\x12\x20\n\x07send_to\x18\x01\x20\x01(\tR\x06sendToB\x07\xfaB\
510-
\x04r\x02\x10\x01\x12L\n\x0frequest_timeout\x18\x02\x20\x01(\x0b2\x19.go\
511-
ogle.protobuf.DurationR\x0erequestTimeoutB\x08\xfaB\x05\xaa\x01\x02*\0\
512-
\x12'\n\rmax_in_flight\x18\x03\x20\x01(\x04H\0R\x0bmaxInFlight\x88\x01\
513-
\x01\x12G\n\x0cqueue_policy\x18\x04\x20\x01(\x0b2$.pulse.config.outflow.\
514-
v1.QueuePolicyR\x0bqueuePolicy\x12/\n\x11batch_max_samples\x18\x05\x20\
515-
\x01(\x04H\x01R\x0fbatchMaxSamples\x88\x01\x01\x12#\n\rmetadata_only\x18\
516-
\x06\x20\x01(\x08R\x0cmetadataOnly\x12F\n\x04auth\x18\x07\x20\x01(\x0b22\
517-
.pulse.config.outflow.v1.HttpRemoteWriteAuthConfigR\x04auth\x12O\n\x0fre\
518-
quest_headers\x18\x08\x20\x03(\x0b2&.pulse.config.outflow.v1.RequestHead\
519-
erR\x0erequestHeaders\x123\n\x13convert_metric_name\x18\t\x20\x01(\x08H\
520-
\x02R\x11convertMetricName\x88\x01\x01\x12F\n\x0cretry_policy\x18\n\x20\
521-
\x01(\x0b2#.pulse.config.common.v1.RetryPolicyR\x0bretryPolicy\x12y\n\
522-
\x14lyft_specific_config\x18\x0b\x20\x01(\x0b2G.pulse.config.outflow.v1.\
523-
PromRemoteWriteClientConfig.LyftSpecificConfigR\x12lyftSpecificConfig\
524-
\x1a\xb9\x02\n\x12LyftSpecificConfig\x12=\n\x16general_storage_policy\
525-
\x18\x01\x20\x01(\tR\x14generalStoragePolicyB\x07\xfaB\x04r\x02\x10\x01\
526-
\x12J\n\x1finstance_metrics_storage_policy\x18\x02\x20\x01(\tH\0R\x1cins\
527-
tanceMetricsStoragePolicy\x88\x01\x01\x12N\n!cloudwatch_metrics_storage_\
528-
policy\x18\x03\x20\x01(\tH\x01R\x1ecloudwatchMetricsStoragePolicy\x88\
529-
\x01\x01B\"\n\x20_instance_metrics_storage_policyB$\n\"_cloudwatch_metri\
530-
cs_storage_policyB\x10\n\x0e_max_in_flightB\x14\n\x12_batch_max_samplesB\
531-
\x16\n\x14_convert_metric_nameb\x06proto3\
527+
n.proto\x1a\x17validate/validate.proto\"\xa9\t\n\x1bPromRemoteWriteClien\
528+
tConfig\x12\x20\n\x07send_to\x18\x01\x20\x01(\tR\x06sendToB\x07\xfaB\x04\
529+
r\x02\x10\x01\x12L\n\x0frequest_timeout\x18\x02\x20\x01(\x0b2\x19.google\
530+
.protobuf.DurationR\x0erequestTimeoutB\x08\xfaB\x05\xaa\x01\x02*\0\x12'\
531+
\n\rmax_in_flight\x18\x03\x20\x01(\x04H\0R\x0bmaxInFlight\x88\x01\x01\
532+
\x12G\n\x0cqueue_policy\x18\x04\x20\x01(\x0b2$.pulse.config.outflow.v1.Q\
533+
ueuePolicyR\x0bqueuePolicy\x12/\n\x11batch_max_samples\x18\x05\x20\x01(\
534+
\x04H\x01R\x0fbatchMaxSamples\x88\x01\x01\x12#\n\rmetadata_only\x18\x06\
535+
\x20\x01(\x08R\x0cmetadataOnly\x12F\n\x04auth\x18\x07\x20\x01(\x0b22.pul\
536+
se.config.outflow.v1.HttpRemoteWriteAuthConfigR\x04auth\x12O\n\x0freques\
537+
t_headers\x18\x08\x20\x03(\x0b2&.pulse.config.outflow.v1.RequestHeaderR\
538+
\x0erequestHeaders\x123\n\x13convert_metric_name\x18\t\x20\x01(\x08H\x02\
539+
R\x11convertMetricName\x88\x01\x01\x12F\n\x0cretry_policy\x18\n\x20\x01(\
540+
\x0b2#.pulse.config.common.v1.RetryPolicyR\x0bretryPolicy\x12y\n\x14lyft\
541+
_specific_config\x18\x0b\x20\x01(\x0b2G.pulse.config.outflow.v1.PromRemo\
542+
teWriteClientConfig.LyftSpecificConfigR\x12lyftSpecificConfig\x12E\n\x11\
543+
pool_idle_timeout\x18\x0c\x20\x01(\x0b2\x19.google.protobuf.DurationR\
544+
\x0fpoolIdleTimeout\x1a\xb9\x02\n\x12LyftSpecificConfig\x12=\n\x16genera\
545+
l_storage_policy\x18\x01\x20\x01(\tR\x14generalStoragePolicyB\x07\xfaB\
546+
\x04r\x02\x10\x01\x12J\n\x1finstance_metrics_storage_policy\x18\x02\x20\
547+
\x01(\tH\0R\x1cinstanceMetricsStoragePolicy\x88\x01\x01\x12N\n!cloudwatc\
548+
h_metrics_storage_policy\x18\x03\x20\x01(\tH\x01R\x1ecloudwatchMetricsSt\
549+
oragePolicy\x88\x01\x01B\"\n\x20_instance_metrics_storage_policyB$\n\"_c\
550+
loudwatch_metrics_storage_policyB\x10\n\x0e_max_in_flightB\x14\n\x12_bat\
551+
ch_max_samplesB\x16\n\x14_convert_metric_nameb\x06proto3\
532552
";
533553

534554
/// `FileDescriptorProto` object which was a source for this generated file

pulse-proxy/src/test/integration/mod.rs

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -579,6 +579,7 @@ impl PromClient {
579579
None,
580580
PROM_REMOTE_WRITE_HEADERS,
581581
vec![],
582+
None,
582583
)
583584
.await
584585
.unwrap(),
@@ -620,6 +621,7 @@ impl OtlpClient {
620621
None,
621622
&[],
622623
vec![],
624+
None,
623625
)
624626
.await
625627
.unwrap(),

0 commit comments

Comments
 (0)