Skip to content

Commit 2a83ce7

Browse files
committed
Limit connection timeout to recovered leaders
The global gRPC connection timeout affected every Arroyo service and RPC client even though the observed stall occurs while reconnecting a recovered job to its worker leader. Remove the global configuration and endpoint behavior. Restore the 10-second deadline around each leader connection attempt while preserving the existing retry and recovery behavior.
1 parent ae464f8 commit 2a83ce7

4 files changed

Lines changed: 22 additions & 50 deletions

File tree

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

Lines changed: 18 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -18,6 +18,8 @@ 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+
2123
pub struct LeaderManager {
2224
leader_client: JobStatusGrpcClient<InterceptedService<Channel, InjectWorkerId>>,
2325
pub job_id: JobId,
@@ -35,13 +37,23 @@ impl LeaderManager {
3537
address: String,
3638
) -> anyhow::Result<Self> {
3739
let leader_client = retry!(
38-
job_status_client(
39-
"controller",
40-
&config().worker.tls,
41-
worker_id,
42-
address.clone()
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+
),
4348
)
44-
.await,
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+
},
4557
5,
4658
Duration::from_millis(100),
4759
Duration::from_secs(2),

crates/arroyo-rpc/default.toml

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

5-
[grpc]
6-
connect-timeout = "10s"
7-
85
[pipeline]
96
source-batch-size = 512
107
source-batch-linger = "100ms"

crates/arroyo-rpc/src/config.rs

Lines changed: 0 additions & 35 deletions
Original file line numberDiff line numberDiff line change
@@ -193,9 +193,6 @@ 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-
199196
/// API service configuration
200197
pub api: ApiConfig,
201198

@@ -278,13 +275,6 @@ pub struct Config {
278275
pub disable_telemetry: bool,
279276
}
280277

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-
288278
#[derive(Debug, Deserialize, Serialize, Clone, Default)]
289279
#[serde(rename_all = "kebab-case", deny_unknown_fields)]
290280
pub enum JobControllerMode {
@@ -1057,7 +1047,6 @@ impl TlsConfig {
10571047
#[cfg(test)]
10581048
mod tests {
10591049
use crate::config::{Config, DatabaseType, Scheduler, SchemaName, SqliteConfig, load_config};
1060-
use std::time::Duration;
10611050
use url::Url;
10621051

10631052
#[test]
@@ -1084,30 +1073,6 @@ mod tests {
10841073
}
10851074
}
10861075

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-
11111076
#[test]
11121077
#[allow(clippy::result_large_err)]
11131078
fn test_config() {

crates/arroyo-rpc/src/lib.rs

Lines changed: 4 additions & 6 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-
let endpoint = if let Some(target_tls) = config.get_tls_config(target_tls) {
1007+
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,13 +1031,11 @@ pub async fn grpc_channel_builder(
10311031
config_builder = config_builder.identity(Identity::from_pem(our_tls.cert, our_tls.key));
10321032
}
10331033

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

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

0 commit comments

Comments
 (0)