Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
1 change: 1 addition & 0 deletions pulse-metrics/src/admin/server.rs
Original file line number Diff line number Diff line change
Expand Up @@ -344,6 +344,7 @@ impl MetaStatsEmitter {
prom_remote_write.auth.into_option(),
PROM_REMOTE_WRITE_HEADERS,
prom_remote_write.request_headers,
None,
)
.await?,
}) as Box<dyn MetaStatsSender>,
Expand Down
9 changes: 8 additions & 1 deletion pulse-metrics/src/clients/http.rs
Original file line number Diff line number Diff line change
Expand Up @@ -168,10 +168,17 @@ impl HyperHttpRemoteWriteClient {
auth_config: Option<HttpRemoteWriteAuthConfig>,
core_request_headers: &[(&str, &str)],
config_request_headers: Vec<RequestHeader>,
pool_idle_timeout: Option<Duration>,
) -> Result<Self> {
Ok(Self {
// TODO(mattklein123): Make connect timeout configurable.
inner: Client::builder(TokioExecutor::new()).build(make_tls_connector(250.milliseconds())),
inner: Client::builder(TokioExecutor::new())
.pool_idle_timeout(
pool_idle_timeout
.unwrap_or_else(|| 90.seconds())
.unsigned_abs(),
)
.build(make_tls_connector(250.milliseconds())),
endpoint,
timeout,
auth: Self::create_auth(auth_config).await?,
Expand Down
6 changes: 5 additions & 1 deletion pulse-metrics/src/pipeline/outflow/http/remote_write.rs
Original file line number Diff line number Diff line change
Expand Up @@ -24,7 +24,7 @@ use backoff::backoff::Backoff;
use bd_log::warn_every;
use bd_server_stats::stats::{AutoGauge, Scope};
use bd_shutdown::{ComponentShutdown, ComponentStatus};
use bd_time::TimeDurationExt;
use bd_time::{ProtoDurationExt, TimeDurationExt};
use bytes::Bytes;
use http::HeaderMap;
use prometheus::{Histogram, IntCounter, IntGauge};
Expand Down Expand Up @@ -124,6 +124,7 @@ impl HttpRemoteWriteOutflow {
core_request_headers: &[(&str, &str)],
config_request_headers: Vec<RequestHeader>,
context: OutflowFactoryContext,
pool_idle_timeout: MessageField<Duration>,
) -> anyhow::Result<Arc<Self>> {
let request_timeout = request_timeout.unwrap_duration_or(DEFAULT_REQUEST_TIMEOUT);
let client = Arc::new(
Expand All @@ -133,6 +134,9 @@ impl HttpRemoteWriteOutflow {
auth_config,
core_request_headers,
config_request_headers,
pool_idle_timeout
.map(|d| d.to_time_duration())
.into_option(),
)
.await?,
);
Expand Down
1 change: 1 addition & 0 deletions pulse-metrics/src/pipeline/outflow/otlp/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -88,6 +88,7 @@ pub async fn make_otlp_outflow(
&[],
config.request_headers,
context,
None.into(),
)
.await
}
Expand Down
1 change: 1 addition & 0 deletions pulse-metrics/src/pipeline/outflow/prom/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -96,6 +96,7 @@ pub async fn make_prom_outflow(
PROM_REMOTE_WRITE_HEADERS,
config.request_headers,
context,
config.pool_idle_timeout,
)
.await
}
Expand Down
1 change: 1 addition & 0 deletions pulse-promcli/src/main.rs
Original file line number Diff line number Diff line change
Expand Up @@ -107,6 +107,7 @@ async fn main() -> anyhow::Result<()> {
None,
PROM_REMOTE_WRITE_HEADERS,
vec![],
None,
)
.await?;
client
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -64,4 +64,7 @@ message PromRemoteWriteClientConfig {

// Lyft specific configuration. To be generalized later.
LyftSpecificConfig lyft_specific_config = 11;

// Idle time for outbound connections in the connection pool. If not set, defaults to 90s.
google.protobuf.Duration pool_idle_timeout = 12;
}
Original file line number Diff line number Diff line change
Expand Up @@ -57,6 +57,8 @@ pub struct PromRemoteWriteClientConfig {
pub retry_policy: ::protobuf::MessageField<super::retry::RetryPolicy>,
// @@protoc_insertion_point(field:pulse.config.outflow.v1.PromRemoteWriteClientConfig.lyft_specific_config)
pub lyft_specific_config: ::protobuf::MessageField<prom_remote_write_client_config::LyftSpecificConfig>,
// @@protoc_insertion_point(field:pulse.config.outflow.v1.PromRemoteWriteClientConfig.pool_idle_timeout)
pub pool_idle_timeout: ::protobuf::MessageField<::protobuf::well_known_types::duration::Duration>,
// special fields
// @@protoc_insertion_point(special_field:pulse.config.outflow.v1.PromRemoteWriteClientConfig.special_fields)
pub special_fields: ::protobuf::SpecialFields,
Expand All @@ -74,7 +76,7 @@ impl PromRemoteWriteClientConfig {
}

fn generated_message_descriptor_data() -> ::protobuf::reflect::GeneratedMessageDescriptorData {
let mut fields = ::std::vec::Vec::with_capacity(11);
let mut fields = ::std::vec::Vec::with_capacity(12);
let mut oneofs = ::std::vec::Vec::with_capacity(0);
fields.push(::protobuf::reflect::rt::v2::make_simpler_field_accessor::<_, _>(
"send_to",
Expand Down Expand Up @@ -131,6 +133,11 @@ impl PromRemoteWriteClientConfig {
|m: &PromRemoteWriteClientConfig| { &m.lyft_specific_config },
|m: &mut PromRemoteWriteClientConfig| { &mut m.lyft_specific_config },
));
fields.push(::protobuf::reflect::rt::v2::make_message_field_accessor::<_, ::protobuf::well_known_types::duration::Duration>(
"pool_idle_timeout",
|m: &PromRemoteWriteClientConfig| { &m.pool_idle_timeout },
|m: &mut PromRemoteWriteClientConfig| { &mut m.pool_idle_timeout },
));
::protobuf::reflect::GeneratedMessageDescriptorData::new_2::<PromRemoteWriteClientConfig>(
"PromRemoteWriteClientConfig",
fields,
Expand Down Expand Up @@ -182,6 +189,9 @@ impl ::protobuf::Message for PromRemoteWriteClientConfig {
90 => {
::protobuf::rt::read_singular_message_into_field(is, &mut self.lyft_specific_config)?;
},
98 => {
::protobuf::rt::read_singular_message_into_field(is, &mut self.pool_idle_timeout)?;
},
tag => {
::protobuf::rt::read_unknown_or_skip_group(tag, is, self.special_fields.mut_unknown_fields())?;
},
Expand Down Expand Up @@ -233,6 +243,10 @@ impl ::protobuf::Message for PromRemoteWriteClientConfig {
let len = v.compute_size();
my_size += 1 + ::protobuf::rt::compute_raw_varint64_size(len) + len;
}
if let Some(v) = self.pool_idle_timeout.as_ref() {
let len = v.compute_size();
my_size += 1 + ::protobuf::rt::compute_raw_varint64_size(len) + len;
}
my_size += ::protobuf::rt::unknown_fields_size(self.special_fields.unknown_fields());
self.special_fields.cached_size().set(my_size as u32);
my_size
Expand Down Expand Up @@ -272,6 +286,9 @@ impl ::protobuf::Message for PromRemoteWriteClientConfig {
if let Some(v) = self.lyft_specific_config.as_ref() {
::protobuf::rt::write_message_field_with_cached_size(11, v, os)?;
}
if let Some(v) = self.pool_idle_timeout.as_ref() {
::protobuf::rt::write_message_field_with_cached_size(12, v, os)?;
}
os.write_unknown_fields(self.special_fields.unknown_fields())?;
::std::result::Result::Ok(())
}
Expand Down Expand Up @@ -300,6 +317,7 @@ impl ::protobuf::Message for PromRemoteWriteClientConfig {
self.convert_metric_name = ::std::option::Option::None;
self.retry_policy.clear();
self.lyft_specific_config.clear();
self.pool_idle_timeout.clear();
self.special_fields.clear();
}

Expand All @@ -316,6 +334,7 @@ impl ::protobuf::Message for PromRemoteWriteClientConfig {
convert_metric_name: ::std::option::Option::None,
retry_policy: ::protobuf::MessageField::none(),
lyft_specific_config: ::protobuf::MessageField::none(),
pool_idle_timeout: ::protobuf::MessageField::none(),
special_fields: ::protobuf::SpecialFields::new(),
};
&instance
Expand Down Expand Up @@ -505,30 +524,31 @@ static file_descriptor_proto_data: &'static [u8] = b"\
utflow.v1\x1a#pulse/config/common/v1/common.proto\x1a\"pulse/config/comm\
on/v1/retry.proto\x1a*pulse/config/outflow/v1/queue_policy.proto\x1a,pul\
se/config/outflow/v1/outflow_common.proto\x1a\x1egoogle/protobuf/duratio\
n.proto\x1a\x17validate/validate.proto\"\xe2\x08\n\x1bPromRemoteWriteCli\
entConfig\x12\x20\n\x07send_to\x18\x01\x20\x01(\tR\x06sendToB\x07\xfaB\
\x04r\x02\x10\x01\x12L\n\x0frequest_timeout\x18\x02\x20\x01(\x0b2\x19.go\
ogle.protobuf.DurationR\x0erequestTimeoutB\x08\xfaB\x05\xaa\x01\x02*\0\
\x12'\n\rmax_in_flight\x18\x03\x20\x01(\x04H\0R\x0bmaxInFlight\x88\x01\
\x01\x12G\n\x0cqueue_policy\x18\x04\x20\x01(\x0b2$.pulse.config.outflow.\
v1.QueuePolicyR\x0bqueuePolicy\x12/\n\x11batch_max_samples\x18\x05\x20\
\x01(\x04H\x01R\x0fbatchMaxSamples\x88\x01\x01\x12#\n\rmetadata_only\x18\
\x06\x20\x01(\x08R\x0cmetadataOnly\x12F\n\x04auth\x18\x07\x20\x01(\x0b22\
.pulse.config.outflow.v1.HttpRemoteWriteAuthConfigR\x04auth\x12O\n\x0fre\
quest_headers\x18\x08\x20\x03(\x0b2&.pulse.config.outflow.v1.RequestHead\
erR\x0erequestHeaders\x123\n\x13convert_metric_name\x18\t\x20\x01(\x08H\
\x02R\x11convertMetricName\x88\x01\x01\x12F\n\x0cretry_policy\x18\n\x20\
\x01(\x0b2#.pulse.config.common.v1.RetryPolicyR\x0bretryPolicy\x12y\n\
\x14lyft_specific_config\x18\x0b\x20\x01(\x0b2G.pulse.config.outflow.v1.\
PromRemoteWriteClientConfig.LyftSpecificConfigR\x12lyftSpecificConfig\
\x1a\xb9\x02\n\x12LyftSpecificConfig\x12=\n\x16general_storage_policy\
\x18\x01\x20\x01(\tR\x14generalStoragePolicyB\x07\xfaB\x04r\x02\x10\x01\
\x12J\n\x1finstance_metrics_storage_policy\x18\x02\x20\x01(\tH\0R\x1cins\
tanceMetricsStoragePolicy\x88\x01\x01\x12N\n!cloudwatch_metrics_storage_\
policy\x18\x03\x20\x01(\tH\x01R\x1ecloudwatchMetricsStoragePolicy\x88\
\x01\x01B\"\n\x20_instance_metrics_storage_policyB$\n\"_cloudwatch_metri\
cs_storage_policyB\x10\n\x0e_max_in_flightB\x14\n\x12_batch_max_samplesB\
\x16\n\x14_convert_metric_nameb\x06proto3\
n.proto\x1a\x17validate/validate.proto\"\xa9\t\n\x1bPromRemoteWriteClien\
tConfig\x12\x20\n\x07send_to\x18\x01\x20\x01(\tR\x06sendToB\x07\xfaB\x04\
r\x02\x10\x01\x12L\n\x0frequest_timeout\x18\x02\x20\x01(\x0b2\x19.google\
.protobuf.DurationR\x0erequestTimeoutB\x08\xfaB\x05\xaa\x01\x02*\0\x12'\
\n\rmax_in_flight\x18\x03\x20\x01(\x04H\0R\x0bmaxInFlight\x88\x01\x01\
\x12G\n\x0cqueue_policy\x18\x04\x20\x01(\x0b2$.pulse.config.outflow.v1.Q\
ueuePolicyR\x0bqueuePolicy\x12/\n\x11batch_max_samples\x18\x05\x20\x01(\
\x04H\x01R\x0fbatchMaxSamples\x88\x01\x01\x12#\n\rmetadata_only\x18\x06\
\x20\x01(\x08R\x0cmetadataOnly\x12F\n\x04auth\x18\x07\x20\x01(\x0b22.pul\
se.config.outflow.v1.HttpRemoteWriteAuthConfigR\x04auth\x12O\n\x0freques\
t_headers\x18\x08\x20\x03(\x0b2&.pulse.config.outflow.v1.RequestHeaderR\
\x0erequestHeaders\x123\n\x13convert_metric_name\x18\t\x20\x01(\x08H\x02\
R\x11convertMetricName\x88\x01\x01\x12F\n\x0cretry_policy\x18\n\x20\x01(\
\x0b2#.pulse.config.common.v1.RetryPolicyR\x0bretryPolicy\x12y\n\x14lyft\
_specific_config\x18\x0b\x20\x01(\x0b2G.pulse.config.outflow.v1.PromRemo\
teWriteClientConfig.LyftSpecificConfigR\x12lyftSpecificConfig\x12E\n\x11\
pool_idle_timeout\x18\x0c\x20\x01(\x0b2\x19.google.protobuf.DurationR\
\x0fpoolIdleTimeout\x1a\xb9\x02\n\x12LyftSpecificConfig\x12=\n\x16genera\
l_storage_policy\x18\x01\x20\x01(\tR\x14generalStoragePolicyB\x07\xfaB\
\x04r\x02\x10\x01\x12J\n\x1finstance_metrics_storage_policy\x18\x02\x20\
\x01(\tH\0R\x1cinstanceMetricsStoragePolicy\x88\x01\x01\x12N\n!cloudwatc\
h_metrics_storage_policy\x18\x03\x20\x01(\tH\x01R\x1ecloudwatchMetricsSt\
oragePolicy\x88\x01\x01B\"\n\x20_instance_metrics_storage_policyB$\n\"_c\
loudwatch_metrics_storage_policyB\x10\n\x0e_max_in_flightB\x14\n\x12_bat\
ch_max_samplesB\x16\n\x14_convert_metric_nameb\x06proto3\
";

/// `FileDescriptorProto` object which was a source for this generated file
Expand Down
2 changes: 2 additions & 0 deletions pulse-proxy/src/test/integration/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -579,6 +579,7 @@ impl PromClient {
None,
PROM_REMOTE_WRITE_HEADERS,
vec![],
None,
)
.await
.unwrap(),
Expand Down Expand Up @@ -620,6 +621,7 @@ impl OtlpClient {
None,
&[],
vec![],
None,
)
.await
.unwrap(),
Expand Down
Loading