Skip to content

Commit 50b115c

Browse files
authored
add configurable connect timeout for remote write (#108)
Signed-off-by: Matt Klein <mklein@bitdrift.io>
1 parent 8976ff6 commit 50b115c

12 files changed

Lines changed: 109 additions & 46 deletions

File tree

pulse-metrics/src/admin/server.rs

Lines changed: 4 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -15,7 +15,7 @@ use crate::clients::http::{
1515
PROM_REMOTE_WRITE_HEADERS,
1616
should_retry,
1717
};
18-
use crate::pipeline::config::DEFAULT_REQUEST_TIMEOUT;
18+
use crate::pipeline::config::{DEFAULT_CONNECT_TIMEOUT, DEFAULT_REQUEST_TIMEOUT};
1919
use crate::pipeline::outflow::prom::compress_write_request;
2020
use crate::pipeline::time::{RealTimeProvider, next_flush_interval};
2121
use anyhow::bail;
@@ -341,6 +341,9 @@ impl MetaStatsEmitter {
341341
prom_remote_write
342342
.request_timeout
343343
.unwrap_duration_or(DEFAULT_REQUEST_TIMEOUT),
344+
prom_remote_write
345+
.connect_timeout
346+
.unwrap_duration_or(DEFAULT_CONNECT_TIMEOUT),
344347
prom_remote_write.auth.into_option(),
345348
PROM_REMOTE_WRITE_HEADERS,
346349
prom_remote_write.request_headers,

pulse-metrics/src/clients/http.rs

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -168,20 +168,20 @@ impl HyperHttpRemoteWriteClient {
168168
pub async fn new(
169169
endpoint: String,
170170
timeout: Duration,
171+
connect_timeout: Duration,
171172
auth_config: Option<HttpRemoteWriteAuthConfig>,
172173
core_request_headers: &[(&str, &str)],
173174
config_request_headers: Vec<RequestHeader>,
174175
pool_idle_timeout: Option<Duration>,
175176
) -> Result<Self> {
176177
Ok(Self {
177-
// TODO(mattklein123): Make connect timeout configurable.
178178
inner: Client::builder(TokioExecutor::new())
179179
.pool_idle_timeout(
180180
pool_idle_timeout
181181
.unwrap_or_else(|| 90.seconds())
182182
.unsigned_abs(),
183183
)
184-
.build(make_tls_connector(250.milliseconds())),
184+
.build(make_tls_connector(connect_timeout)),
185185
endpoint,
186186
timeout,
187187
auth: Self::create_auth(auth_config).await?,

pulse-metrics/src/pipeline/config.rs

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -95,6 +95,7 @@ impl fmt::Display for PipelineRouteType {
9595
}
9696

9797
pub const DEFAULT_REQUEST_TIMEOUT: Duration = Duration::seconds(10);
98+
pub const DEFAULT_CONNECT_TIMEOUT: Duration = Duration::milliseconds(250);
9899

99100
#[must_use]
100101
pub const fn default_max_in_flight() -> u64 {

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

Lines changed: 8 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -18,7 +18,11 @@ use crate::clients::http::{
1818
should_retry,
1919
};
2020
use crate::clients::retry::Retry;
21-
use crate::pipeline::config::{DEFAULT_REQUEST_TIMEOUT, default_max_in_flight};
21+
use crate::pipeline::config::{
22+
DEFAULT_CONNECT_TIMEOUT,
23+
DEFAULT_REQUEST_TIMEOUT,
24+
default_max_in_flight,
25+
};
2226
use crate::pipeline::outflow::http::retry_offload::maybe_queue_for_retry;
2327
use crate::pipeline::outflow::{OutflowFactoryContext, OutflowStats, PipelineOutflow};
2428
use crate::pipeline::time::RealTimeProvider;
@@ -142,6 +146,7 @@ pub struct HttpRemoteWriteOutflow {
142146
impl HttpRemoteWriteOutflow {
143147
pub(crate) async fn new(
144148
request_timeout: MessageField<Duration>,
149+
connect_timeout: MessageField<Duration>,
145150
retry_policy: RetryPolicy,
146151
max_in_flight: Option<u64>,
147152
batch_router: Arc<dyn BatchRouter>,
@@ -154,10 +159,12 @@ impl HttpRemoteWriteOutflow {
154159
request_deserializer: RequestDeserializer,
155160
) -> anyhow::Result<Arc<Self>> {
156161
let request_timeout = request_timeout.unwrap_duration_or(DEFAULT_REQUEST_TIMEOUT);
162+
let connect_timeout = connect_timeout.unwrap_duration_or(DEFAULT_CONNECT_TIMEOUT);
157163
let client = Arc::new(
158164
HyperHttpRemoteWriteClient::new(
159165
send_to,
160166
request_timeout,
167+
connect_timeout,
161168
auth_config,
162169
core_request_headers,
163170
config_request_headers,

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

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -82,6 +82,7 @@ pub async fn make_otlp_outflow(
8282
let compression = config.compression.enum_value_or_default();
8383
HttpRemoteWriteOutflow::new(
8484
config.request_timeout,
85+
config.connect_timeout,
8586
config.retry_policy.unwrap_or_default(),
8687
config.max_in_flight,
8788
batch_router,

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

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -103,6 +103,7 @@ pub async fn make_prom_outflow(
103103
);
104104
HttpRemoteWriteOutflow::new(
105105
config.request_timeout,
106+
config.connect_timeout,
106107
config.retry_policy.unwrap_or_default(),
107108
config.max_in_flight,
108109
batch_router,

pulse-promcli/src/main.rs

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -104,6 +104,7 @@ async fn main() -> anyhow::Result<()> {
104104
let client = HyperHttpRemoteWriteClient::new(
105105
options.endpoint,
106106
10.seconds(),
107+
250.milliseconds(),
107108
None,
108109
PROM_REMOTE_WRITE_HEADERS,
109110
vec![],

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

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -27,6 +27,9 @@ message OtlpClientConfig {
2727
// The request timeout. Defaults to 10s.
2828
google.protobuf.Duration request_timeout = 2 [(validate.rules).duration.gt = {}];
2929

30+
// The connection timeout for establishing a new upstream connection. Defaults to 250ms.
31+
google.protobuf.Duration connect_timeout = 11 [(validate.rules).duration.gt = {}];
32+
3033
// The maximum number of inflight requests. Defaults to 16.
3134
optional uint64 max_in_flight = 3;
3235

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
@@ -35,6 +35,9 @@ message PromRemoteWriteClientConfig {
3535
// The request timeout. Defaults to 10s.
3636
google.protobuf.Duration request_timeout = 2 [(validate.rules).duration.gt = {}];
3737

38+
// The connection timeout for establishing a new upstream connection. Defaults to 250ms.
39+
google.protobuf.Duration connect_timeout = 14 [(validate.rules).duration.gt = {}];
40+
3841
// The maximum number of inflight requests. Defaults to 16.
3942
optional uint64 max_in_flight = 3;
4043

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

Lines changed: 37 additions & 16 deletions
Original file line numberDiff line numberDiff line change
@@ -39,6 +39,8 @@ pub struct OtlpClientConfig {
3939
pub send_to: ::protobuf::Chars,
4040
// @@protoc_insertion_point(field:pulse.config.outflow.v1.OtlpClientConfig.request_timeout)
4141
pub request_timeout: ::protobuf::MessageField<::protobuf::well_known_types::duration::Duration>,
42+
// @@protoc_insertion_point(field:pulse.config.outflow.v1.OtlpClientConfig.connect_timeout)
43+
pub connect_timeout: ::protobuf::MessageField<::protobuf::well_known_types::duration::Duration>,
4244
// @@protoc_insertion_point(field:pulse.config.outflow.v1.OtlpClientConfig.max_in_flight)
4345
pub max_in_flight: ::std::option::Option<u64>,
4446
// @@protoc_insertion_point(field:pulse.config.outflow.v1.OtlpClientConfig.queue_policy)
@@ -72,7 +74,7 @@ impl OtlpClientConfig {
7274
}
7375

7476
fn generated_message_descriptor_data() -> ::protobuf::reflect::GeneratedMessageDescriptorData {
75-
let mut fields = ::std::vec::Vec::with_capacity(10);
77+
let mut fields = ::std::vec::Vec::with_capacity(11);
7678
let mut oneofs = ::std::vec::Vec::with_capacity(0);
7779
fields.push(::protobuf::reflect::rt::v2::make_simpler_field_accessor::<_, _>(
7880
"send_to",
@@ -84,6 +86,11 @@ impl OtlpClientConfig {
8486
|m: &OtlpClientConfig| { &m.request_timeout },
8587
|m: &mut OtlpClientConfig| { &mut m.request_timeout },
8688
));
89+
fields.push(::protobuf::reflect::rt::v2::make_message_field_accessor::<_, ::protobuf::well_known_types::duration::Duration>(
90+
"connect_timeout",
91+
|m: &OtlpClientConfig| { &m.connect_timeout },
92+
|m: &mut OtlpClientConfig| { &mut m.connect_timeout },
93+
));
8794
fields.push(::protobuf::reflect::rt::v2::make_option_accessor::<_, _>(
8895
"max_in_flight",
8996
|m: &OtlpClientConfig| { &m.max_in_flight },
@@ -148,6 +155,9 @@ impl ::protobuf::Message for OtlpClientConfig {
148155
18 => {
149156
::protobuf::rt::read_singular_message_into_field(is, &mut self.request_timeout)?;
150157
},
158+
90 => {
159+
::protobuf::rt::read_singular_message_into_field(is, &mut self.connect_timeout)?;
160+
},
151161
24 => {
152162
self.max_in_flight = ::std::option::Option::Some(is.read_uint64()?);
153163
},
@@ -191,6 +201,10 @@ impl ::protobuf::Message for OtlpClientConfig {
191201
let len = v.compute_size();
192202
my_size += 1 + ::protobuf::rt::compute_raw_varint64_size(len) + len;
193203
}
204+
if let Some(v) = self.connect_timeout.as_ref() {
205+
let len = v.compute_size();
206+
my_size += 1 + ::protobuf::rt::compute_raw_varint64_size(len) + len;
207+
}
194208
if let Some(v) = self.max_in_flight {
195209
my_size += ::protobuf::rt::uint64_size(3, v);
196210
}
@@ -231,6 +245,9 @@ impl ::protobuf::Message for OtlpClientConfig {
231245
if let Some(v) = self.request_timeout.as_ref() {
232246
::protobuf::rt::write_message_field_with_cached_size(2, v, os)?;
233247
}
248+
if let Some(v) = self.connect_timeout.as_ref() {
249+
::protobuf::rt::write_message_field_with_cached_size(11, v, os)?;
250+
}
234251
if let Some(v) = self.max_in_flight {
235252
os.write_uint64(3, v)?;
236253
}
@@ -274,6 +291,7 @@ impl ::protobuf::Message for OtlpClientConfig {
274291
fn clear(&mut self) {
275292
self.send_to.clear();
276293
self.request_timeout.clear();
294+
self.connect_timeout.clear();
277295
self.max_in_flight = ::std::option::Option::None;
278296
self.queue_policy.clear();
279297
self.batch_max_samples = ::std::option::Option::None;
@@ -289,6 +307,7 @@ impl ::protobuf::Message for OtlpClientConfig {
289307
static instance: OtlpClientConfig = OtlpClientConfig {
290308
send_to: ::protobuf::Chars::new(),
291309
request_timeout: ::protobuf::MessageField::none(),
310+
connect_timeout: ::protobuf::MessageField::none(),
292311
max_in_flight: ::std::option::Option::None,
293312
queue_policy: ::protobuf::MessageField::none(),
294313
batch_max_samples: ::std::option::Option::None,
@@ -390,23 +409,25 @@ static file_descriptor_proto_data: &'static [u8] = b"\
390409
\x1a\x1egoogle/protobuf/duration.proto\x1a#pulse/config/common/v1/common\
391410
.proto\x1a\"pulse/config/common/v1/retry.proto\x1a,pulse/config/outflow/\
392411
v1/outflow_common.proto\x1a*pulse/config/outflow/v1/queue_policy.proto\
393-
\x1a\x17validate/validate.proto\"\xf3\x05\n\x10OtlpClientConfig\x12\x20\
412+
\x1a\x17validate/validate.proto\"\xc1\x06\n\x10OtlpClientConfig\x12\x20\
394413
\n\x07send_to\x18\x01\x20\x01(\tR\x06sendToB\x07\xfaB\x04r\x02\x10\x01\
395414
\x12L\n\x0frequest_timeout\x18\x02\x20\x01(\x0b2\x19.google.protobuf.Dur\
396-
ationR\x0erequestTimeoutB\x08\xfaB\x05\xaa\x01\x02*\0\x12'\n\rmax_in_fli\
397-
ght\x18\x03\x20\x01(\x04H\0R\x0bmaxInFlight\x88\x01\x01\x12G\n\x0cqueue_\
398-
policy\x18\x04\x20\x01(\x0b2$.pulse.config.outflow.v1.QueuePolicyR\x0bqu\
399-
euePolicy\x12/\n\x11batch_max_samples\x18\x05\x20\x01(\x04H\x01R\x0fbatc\
400-
hMaxSamples\x88\x01\x01\x12F\n\x04auth\x18\x06\x20\x01(\x0b22.pulse.conf\
401-
ig.outflow.v1.HttpRemoteWriteAuthConfigR\x04auth\x12O\n\x0frequest_heade\
402-
rs\x18\x07\x20\x03(\x0b2&.pulse.config.outflow.v1.RequestHeaderR\x0erequ\
403-
estHeaders\x12F\n\x0cretry_policy\x18\x08\x20\x01(\x0b2#.pulse.config.co\
404-
mmon.v1.RetryPolicyR\x0bretryPolicy\x12[\n\x0bcompression\x18\t\x20\x01(\
405-
\x0e29.pulse.config.outflow.v1.OtlpClientConfig.OtlpCompressionR\x0bcomp\
406-
ression\x12=\n\x1bconvert_names_to_prometheus\x18\n\x20\x01(\x08R\x18con\
407-
vertNamesToPrometheus\"'\n\x0fOtlpCompression\x12\n\n\x06SNAPPY\x10\0\
408-
\x12\x08\n\x04NONE\x10\x01B\x10\n\x0e_max_in_flightB\x14\n\x12_batch_max\
409-
_samplesb\x06proto3\
415+
ationR\x0erequestTimeoutB\x08\xfaB\x05\xaa\x01\x02*\0\x12L\n\x0fconnect_\
416+
timeout\x18\x0b\x20\x01(\x0b2\x19.google.protobuf.DurationR\x0econnectTi\
417+
meoutB\x08\xfaB\x05\xaa\x01\x02*\0\x12'\n\rmax_in_flight\x18\x03\x20\x01\
418+
(\x04H\0R\x0bmaxInFlight\x88\x01\x01\x12G\n\x0cqueue_policy\x18\x04\x20\
419+
\x01(\x0b2$.pulse.config.outflow.v1.QueuePolicyR\x0bqueuePolicy\x12/\n\
420+
\x11batch_max_samples\x18\x05\x20\x01(\x04H\x01R\x0fbatchMaxSamples\x88\
421+
\x01\x01\x12F\n\x04auth\x18\x06\x20\x01(\x0b22.pulse.config.outflow.v1.H\
422+
ttpRemoteWriteAuthConfigR\x04auth\x12O\n\x0frequest_headers\x18\x07\x20\
423+
\x03(\x0b2&.pulse.config.outflow.v1.RequestHeaderR\x0erequestHeaders\x12\
424+
F\n\x0cretry_policy\x18\x08\x20\x01(\x0b2#.pulse.config.common.v1.RetryP\
425+
olicyR\x0bretryPolicy\x12[\n\x0bcompression\x18\t\x20\x01(\x0e29.pulse.c\
426+
onfig.outflow.v1.OtlpClientConfig.OtlpCompressionR\x0bcompression\x12=\n\
427+
\x1bconvert_names_to_prometheus\x18\n\x20\x01(\x08R\x18convertNamesToPro\
428+
metheus\"'\n\x0fOtlpCompression\x12\n\n\x06SNAPPY\x10\0\x12\x08\n\x04NON\
429+
E\x10\x01B\x10\n\x0e_max_in_flightB\x14\n\x12_batch_max_samplesb\x06prot\
430+
o3\
410431
";
411432

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

0 commit comments

Comments
 (0)