Skip to content

Commit ae464f8

Browse files
committed
Apply configurable gRPC connection timeout
The recovered-leader fix bounded only one connection path, while other controller and worker clients could still wait indefinitely for tonic to establish a transport connection. Add a global `grpc.connect-timeout` setting with a 10-second default and apply it to every endpoint produced by `grpc_channel_builder`. Remove the leader-specific wrapper while retaining its retry policy, and cover both the default and environment override configuration.
1 parent 871e3a3 commit ae464f8

4 files changed

Lines changed: 50 additions & 22 deletions

File tree

crates/arroyo-controller/src/job_controller/leader_manager.rs

Lines changed: 6 additions & 18 deletions
Original file line numberDiff line numberDiff line change
@@ -18,8 +18,6 @@ use tonic::codegen::InterceptedService;
1818
use tonic::transport::Channel;
1919
use tracing::{info, warn};
2020

21-
const LEADER_CONNECT_TIMEOUT: Duration = Duration::from_secs(10);
22-
2321
pub struct LeaderManager {
2422
leader_client: JobStatusGrpcClient<InterceptedService<Channel, InjectWorkerId>>,
2523
pub job_id: JobId,
@@ -37,23 +35,13 @@ impl LeaderManager {
3735
address: String,
3836
) -> anyhow::Result<Self> {
3937
let leader_client = retry!(
40-
match tokio::time::timeout(
41-
LEADER_CONNECT_TIMEOUT,
42-
job_status_client(
43-
"controller",
44-
&config().worker.tls,
45-
worker_id,
46-
address.clone(),
47-
),
38+
job_status_client(
39+
"controller",
40+
&config().worker.tls,
41+
worker_id,
42+
address.clone()
4843
)
49-
.await
50-
{
51-
Ok(result) => result,
52-
Err(_) => Err(anyhow!(
53-
"timed out connecting to worker leader after {:?}",
54-
LEADER_CONNECT_TIMEOUT
55-
)),
56-
},
44+
.await,
5745
5,
5846
Duration::from_millis(100),
5947
Duration::from_secs(2),

crates/arroyo-rpc/default.toml

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -2,6 +2,9 @@ checkpoint-url = "/tmp/arroyo/checkpoints"
22
default-checkpoint-interval = "10s"
33
job-controller = "controller"
44

5+
[grpc]
6+
connect-timeout = "10s"
7+
58
[pipeline]
69
source-batch-size = 512
710
source-batch-linger = "100ms"

crates/arroyo-rpc/src/config.rs

Lines changed: 35 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -193,6 +193,9 @@ fn load_config(paths: &[PathBuf]) -> Figment {
193193
#[derive(Debug, Deserialize, Serialize, Clone)]
194194
#[serde(rename_all = "kebab-case", deny_unknown_fields)]
195195
pub struct Config {
196+
/// gRPC client configuration
197+
pub grpc: GrpcConfig,
198+
196199
/// API service configuration
197200
pub api: ApiConfig,
198201

@@ -275,6 +278,13 @@ pub struct Config {
275278
pub disable_telemetry: bool,
276279
}
277280

281+
#[derive(Debug, Deserialize, Serialize, Clone)]
282+
#[serde(rename_all = "kebab-case", deny_unknown_fields)]
283+
pub struct GrpcConfig {
284+
/// Maximum time to establish a gRPC connection
285+
pub connect_timeout: HumanReadableDuration,
286+
}
287+
278288
#[derive(Debug, Deserialize, Serialize, Clone, Default)]
279289
#[serde(rename_all = "kebab-case", deny_unknown_fields)]
280290
pub enum JobControllerMode {
@@ -1047,6 +1057,7 @@ impl TlsConfig {
10471057
#[cfg(test)]
10481058
mod tests {
10491059
use crate::config::{Config, DatabaseType, Scheduler, SchemaName, SqliteConfig, load_config};
1060+
use std::time::Duration;
10501061
use url::Url;
10511062

10521063
#[test]
@@ -1073,6 +1084,30 @@ mod tests {
10731084
}
10741085
}
10751086

1087+
#[test]
1088+
#[allow(clippy::result_large_err)]
1089+
fn grpc_connect_timeout_defaults_to_ten_seconds() {
1090+
figment::Jail::expect_with(|_| {
1091+
let config: Config = load_config(&[]).extract().unwrap();
1092+
1093+
assert_eq!(*config.grpc.connect_timeout, Duration::from_secs(10));
1094+
Ok(())
1095+
});
1096+
}
1097+
1098+
#[test]
1099+
#[allow(clippy::result_large_err)]
1100+
fn grpc_connect_timeout_can_be_overridden_with_environment() {
1101+
figment::Jail::expect_with(|jail| {
1102+
jail.set_env("ARROYO__GRPC__CONNECT_TIMEOUT", "3s");
1103+
1104+
let config: Config = load_config(&[]).extract().unwrap();
1105+
1106+
assert_eq!(*config.grpc.connect_timeout, Duration::from_secs(3));
1107+
Ok(())
1108+
});
1109+
}
1110+
10761111
#[test]
10771112
#[allow(clippy::result_large_err)]
10781113
fn test_config() {

crates/arroyo-rpc/src/lib.rs

Lines changed: 6 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -1004,7 +1004,7 @@ pub async fn grpc_channel_builder(
10041004
target_tls: &Option<TlsConfig>,
10051005
) -> Result<Endpoint> {
10061006
let config = config();
1007-
if let Some(target_tls) = config.get_tls_config(target_tls) {
1007+
let endpoint = if let Some(target_tls) = config.get_tls_config(target_tls) {
10081008
let mut endpoint = Url::parse(&endpoint)?;
10091009
endpoint
10101010
.set_scheme("https")
@@ -1031,11 +1031,13 @@ pub async fn grpc_channel_builder(
10311031
config_builder = config_builder.identity(Identity::from_pem(our_tls.cert, our_tls.key));
10321032
}
10331033

1034-
Ok(b.tls_config(config_builder).context("configuring TLS")?)
1034+
b.tls_config(config_builder).context("configuring TLS")?
10351035
} else {
10361036
debug!("connecting to grpc endpoint {endpoint}");
1037-
Ok(Channel::from_shared(endpoint.to_string())?)
1038-
}
1037+
Channel::from_shared(endpoint.to_string())?
1038+
};
1039+
1040+
Ok(endpoint.connect_timeout(*config.grpc.connect_timeout))
10391041
}
10401042

10411043
/// Connect to a gRPC service with optional TLS

0 commit comments

Comments
 (0)