From 36e18736abb189d7acd3e4bc2b87a5b8b335d6f5 Mon Sep 17 00:00:00 2001 From: Cole MacKenzie Date: Tue, 28 Jul 2026 13:21:45 -0700 Subject: [PATCH 1/5] Bound recovered leader connection attempts Recovered jobs could stop immediately after logging that their state machine had started. A stale worker address left tonic's connection future pending forever because it had no connection deadline. Apply a 10-second timeout to each attempt while retaining the existing retry and recovery behavior. Exhausted attempts now let the state machine handle the unavailable leader instead of remaining stuck before its first state executes. --- .../src/job_controller/leader_manager.rs | 24 ++++++++++++++----- 1 file changed, 18 insertions(+), 6 deletions(-) diff --git a/crates/arroyo-controller/src/job_controller/leader_manager.rs b/crates/arroyo-controller/src/job_controller/leader_manager.rs index 4e7868b4..315c5a22 100644 --- a/crates/arroyo-controller/src/job_controller/leader_manager.rs +++ b/crates/arroyo-controller/src/job_controller/leader_manager.rs @@ -18,6 +18,8 @@ use tonic::codegen::InterceptedService; use tonic::transport::Channel; use tracing::{info, warn}; +const LEADER_CONNECT_TIMEOUT: Duration = Duration::from_secs(10); + pub struct LeaderManager { leader_client: JobStatusGrpcClient>, pub job_id: JobId, @@ -35,13 +37,23 @@ impl LeaderManager { address: String, ) -> anyhow::Result { let leader_client = retry!( - job_status_client( - "controller", - &config().worker.tls, - worker_id, - address.clone() + match tokio::time::timeout( + LEADER_CONNECT_TIMEOUT, + job_status_client( + "controller", + &config().worker.tls, + worker_id, + address.clone(), + ), ) - .await, + .await + { + Ok(result) => result, + Err(_) => Err(anyhow!( + "timed out connecting to worker leader after {:?}", + LEADER_CONNECT_TIMEOUT + )), + }, 5, Duration::from_millis(100), Duration::from_secs(2), From 871e3a3a7ca01dedfd03f1d6facbdb10737a466d Mon Sep 17 00:00:00 2001 From: Cole MacKenzie Date: Tue, 28 Jul 2026 13:23:03 -0700 Subject: [PATCH 2/5] Release job map lock before queue sends A full state-machine queue caused controller RPC handlers to wait while holding the global job map mutex. One stalled job could therefore prevent messages from reaching every other job. Clone the selected job's sender under the map lock, then release the lock before waiting for channel capacity. This preserves per-job backpressure without making it controller-wide. --- crates/arroyo-controller/src/lib.rs | 38 +++++++++++++--------- crates/arroyo-controller/src/states/mod.rs | 4 +++ 2 files changed, 26 insertions(+), 16 deletions(-) diff --git a/crates/arroyo-controller/src/lib.rs b/crates/arroyo-controller/src/lib.rs index 11f100d8..09d1995f 100644 --- a/crates/arroyo-controller/src/lib.rs +++ b/crates/arroyo-controller/src/lib.rs @@ -651,22 +651,28 @@ impl ControllerServer { } async fn send_to_job_queue(&self, job_id: &str, msg: JobMessage) -> Result<(), Status> { - let mut jobs = self.job_state.lock().await; - - if let Some(sm) = jobs.get_mut(job_id) { - if let Err(e) = sm.send(msg).await { - Err(Status::failed_precondition(format!( - "Cannot handle message for {job_id}: {e}" - ))) - } else { - Ok(()) - } - } else { - warn!(message = "Received message for unknown job id", job_id); - Err(Status::failed_precondition(format!( - "No job with id {job_id}" - ))) - } + // Keep per-job backpressure from holding the global job map lock. + let tx = { + let jobs = self.job_state.lock().await; + let Some(sm) = jobs.get(job_id) else { + warn!(message = "Received message for unknown job id", job_id); + return Err(Status::failed_precondition(format!( + "No job with id {job_id}" + ))); + }; + + sm.sender().ok_or_else(|| { + Status::failed_precondition(format!( + "Cannot handle message for {job_id}: State machine is inactive" + )) + })? + }; + + tx.send(msg).await.map_err(|_| { + Status::failed_precondition(format!( + "Cannot handle message for {job_id}: State machine is inactive" + )) + }) } fn start_updater(&self, guard: ShutdownGuard) { diff --git a/crates/arroyo-controller/src/states/mod.rs b/crates/arroyo-controller/src/states/mod.rs index 3d5199a4..d4ce12a3 100644 --- a/crates/arroyo-controller/src/states/mod.rs +++ b/crates/arroyo-controller/src/states/mod.rs @@ -1207,6 +1207,10 @@ impl StateMachine { } } + pub(crate) fn sender(&self) -> Option> { + self.tx.clone() + } + pub fn done(&self) -> bool { if let Some(tx) = &self.tx { tx.is_closed() From ae464f89512aab8ce96afb48bddd41eac7d61cdd Mon Sep 17 00:00:00 2001 From: Cole MacKenzie Date: Tue, 28 Jul 2026 13:33:06 -0700 Subject: [PATCH 3/5] 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. --- .../src/job_controller/leader_manager.rs | 24 ++++--------- crates/arroyo-rpc/default.toml | 3 ++ crates/arroyo-rpc/src/config.rs | 35 +++++++++++++++++++ crates/arroyo-rpc/src/lib.rs | 10 +++--- 4 files changed, 50 insertions(+), 22 deletions(-) diff --git a/crates/arroyo-controller/src/job_controller/leader_manager.rs b/crates/arroyo-controller/src/job_controller/leader_manager.rs index 315c5a22..4e7868b4 100644 --- a/crates/arroyo-controller/src/job_controller/leader_manager.rs +++ b/crates/arroyo-controller/src/job_controller/leader_manager.rs @@ -18,8 +18,6 @@ use tonic::codegen::InterceptedService; use tonic::transport::Channel; use tracing::{info, warn}; -const LEADER_CONNECT_TIMEOUT: Duration = Duration::from_secs(10); - pub struct LeaderManager { leader_client: JobStatusGrpcClient>, pub job_id: JobId, @@ -37,23 +35,13 @@ impl LeaderManager { address: String, ) -> anyhow::Result { let leader_client = retry!( - match tokio::time::timeout( - LEADER_CONNECT_TIMEOUT, - job_status_client( - "controller", - &config().worker.tls, - worker_id, - address.clone(), - ), + job_status_client( + "controller", + &config().worker.tls, + worker_id, + address.clone() ) - .await - { - Ok(result) => result, - Err(_) => Err(anyhow!( - "timed out connecting to worker leader after {:?}", - LEADER_CONNECT_TIMEOUT - )), - }, + .await, 5, Duration::from_millis(100), Duration::from_secs(2), diff --git a/crates/arroyo-rpc/default.toml b/crates/arroyo-rpc/default.toml index 9f3307d0..0a33023a 100644 --- a/crates/arroyo-rpc/default.toml +++ b/crates/arroyo-rpc/default.toml @@ -2,6 +2,9 @@ checkpoint-url = "/tmp/arroyo/checkpoints" default-checkpoint-interval = "10s" job-controller = "controller" +[grpc] +connect-timeout = "10s" + [pipeline] source-batch-size = 512 source-batch-linger = "100ms" diff --git a/crates/arroyo-rpc/src/config.rs b/crates/arroyo-rpc/src/config.rs index 7238b677..db7c6713 100644 --- a/crates/arroyo-rpc/src/config.rs +++ b/crates/arroyo-rpc/src/config.rs @@ -193,6 +193,9 @@ fn load_config(paths: &[PathBuf]) -> Figment { #[derive(Debug, Deserialize, Serialize, Clone)] #[serde(rename_all = "kebab-case", deny_unknown_fields)] pub struct Config { + /// gRPC client configuration + pub grpc: GrpcConfig, + /// API service configuration pub api: ApiConfig, @@ -275,6 +278,13 @@ pub struct Config { pub disable_telemetry: bool, } +#[derive(Debug, Deserialize, Serialize, Clone)] +#[serde(rename_all = "kebab-case", deny_unknown_fields)] +pub struct GrpcConfig { + /// Maximum time to establish a gRPC connection + pub connect_timeout: HumanReadableDuration, +} + #[derive(Debug, Deserialize, Serialize, Clone, Default)] #[serde(rename_all = "kebab-case", deny_unknown_fields)] pub enum JobControllerMode { @@ -1047,6 +1057,7 @@ impl TlsConfig { #[cfg(test)] mod tests { use crate::config::{Config, DatabaseType, Scheduler, SchemaName, SqliteConfig, load_config}; + use std::time::Duration; use url::Url; #[test] @@ -1073,6 +1084,30 @@ mod tests { } } + #[test] + #[allow(clippy::result_large_err)] + fn grpc_connect_timeout_defaults_to_ten_seconds() { + figment::Jail::expect_with(|_| { + let config: Config = load_config(&[]).extract().unwrap(); + + assert_eq!(*config.grpc.connect_timeout, Duration::from_secs(10)); + Ok(()) + }); + } + + #[test] + #[allow(clippy::result_large_err)] + fn grpc_connect_timeout_can_be_overridden_with_environment() { + figment::Jail::expect_with(|jail| { + jail.set_env("ARROYO__GRPC__CONNECT_TIMEOUT", "3s"); + + let config: Config = load_config(&[]).extract().unwrap(); + + assert_eq!(*config.grpc.connect_timeout, Duration::from_secs(3)); + Ok(()) + }); + } + #[test] #[allow(clippy::result_large_err)] fn test_config() { diff --git a/crates/arroyo-rpc/src/lib.rs b/crates/arroyo-rpc/src/lib.rs index 7acc7bb8..10dd665d 100644 --- a/crates/arroyo-rpc/src/lib.rs +++ b/crates/arroyo-rpc/src/lib.rs @@ -1004,7 +1004,7 @@ pub async fn grpc_channel_builder( target_tls: &Option, ) -> Result { let config = config(); - if let Some(target_tls) = config.get_tls_config(target_tls) { + let endpoint = if let Some(target_tls) = config.get_tls_config(target_tls) { let mut endpoint = Url::parse(&endpoint)?; endpoint .set_scheme("https") @@ -1031,11 +1031,13 @@ pub async fn grpc_channel_builder( config_builder = config_builder.identity(Identity::from_pem(our_tls.cert, our_tls.key)); } - Ok(b.tls_config(config_builder).context("configuring TLS")?) + b.tls_config(config_builder).context("configuring TLS")? } else { debug!("connecting to grpc endpoint {endpoint}"); - Ok(Channel::from_shared(endpoint.to_string())?) - } + Channel::from_shared(endpoint.to_string())? + }; + + Ok(endpoint.connect_timeout(*config.grpc.connect_timeout)) } /// Connect to a gRPC service with optional TLS From 2a83ce703a711407b85cadb6d8cc8a1067ffeb4f Mon Sep 17 00:00:00 2001 From: Cole MacKenzie Date: Tue, 28 Jul 2026 15:42:18 -0700 Subject: [PATCH 4/5] 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. --- .../src/job_controller/leader_manager.rs | 24 +++++++++---- crates/arroyo-rpc/default.toml | 3 -- crates/arroyo-rpc/src/config.rs | 35 ------------------- crates/arroyo-rpc/src/lib.rs | 10 +++--- 4 files changed, 22 insertions(+), 50 deletions(-) diff --git a/crates/arroyo-controller/src/job_controller/leader_manager.rs b/crates/arroyo-controller/src/job_controller/leader_manager.rs index 4e7868b4..315c5a22 100644 --- a/crates/arroyo-controller/src/job_controller/leader_manager.rs +++ b/crates/arroyo-controller/src/job_controller/leader_manager.rs @@ -18,6 +18,8 @@ use tonic::codegen::InterceptedService; use tonic::transport::Channel; use tracing::{info, warn}; +const LEADER_CONNECT_TIMEOUT: Duration = Duration::from_secs(10); + pub struct LeaderManager { leader_client: JobStatusGrpcClient>, pub job_id: JobId, @@ -35,13 +37,23 @@ impl LeaderManager { address: String, ) -> anyhow::Result { let leader_client = retry!( - job_status_client( - "controller", - &config().worker.tls, - worker_id, - address.clone() + match tokio::time::timeout( + LEADER_CONNECT_TIMEOUT, + job_status_client( + "controller", + &config().worker.tls, + worker_id, + address.clone(), + ), ) - .await, + .await + { + Ok(result) => result, + Err(_) => Err(anyhow!( + "timed out connecting to worker leader after {:?}", + LEADER_CONNECT_TIMEOUT + )), + }, 5, Duration::from_millis(100), Duration::from_secs(2), diff --git a/crates/arroyo-rpc/default.toml b/crates/arroyo-rpc/default.toml index 0a33023a..9f3307d0 100644 --- a/crates/arroyo-rpc/default.toml +++ b/crates/arroyo-rpc/default.toml @@ -2,9 +2,6 @@ checkpoint-url = "/tmp/arroyo/checkpoints" default-checkpoint-interval = "10s" job-controller = "controller" -[grpc] -connect-timeout = "10s" - [pipeline] source-batch-size = 512 source-batch-linger = "100ms" diff --git a/crates/arroyo-rpc/src/config.rs b/crates/arroyo-rpc/src/config.rs index db7c6713..7238b677 100644 --- a/crates/arroyo-rpc/src/config.rs +++ b/crates/arroyo-rpc/src/config.rs @@ -193,9 +193,6 @@ fn load_config(paths: &[PathBuf]) -> Figment { #[derive(Debug, Deserialize, Serialize, Clone)] #[serde(rename_all = "kebab-case", deny_unknown_fields)] pub struct Config { - /// gRPC client configuration - pub grpc: GrpcConfig, - /// API service configuration pub api: ApiConfig, @@ -278,13 +275,6 @@ pub struct Config { pub disable_telemetry: bool, } -#[derive(Debug, Deserialize, Serialize, Clone)] -#[serde(rename_all = "kebab-case", deny_unknown_fields)] -pub struct GrpcConfig { - /// Maximum time to establish a gRPC connection - pub connect_timeout: HumanReadableDuration, -} - #[derive(Debug, Deserialize, Serialize, Clone, Default)] #[serde(rename_all = "kebab-case", deny_unknown_fields)] pub enum JobControllerMode { @@ -1057,7 +1047,6 @@ impl TlsConfig { #[cfg(test)] mod tests { use crate::config::{Config, DatabaseType, Scheduler, SchemaName, SqliteConfig, load_config}; - use std::time::Duration; use url::Url; #[test] @@ -1084,30 +1073,6 @@ mod tests { } } - #[test] - #[allow(clippy::result_large_err)] - fn grpc_connect_timeout_defaults_to_ten_seconds() { - figment::Jail::expect_with(|_| { - let config: Config = load_config(&[]).extract().unwrap(); - - assert_eq!(*config.grpc.connect_timeout, Duration::from_secs(10)); - Ok(()) - }); - } - - #[test] - #[allow(clippy::result_large_err)] - fn grpc_connect_timeout_can_be_overridden_with_environment() { - figment::Jail::expect_with(|jail| { - jail.set_env("ARROYO__GRPC__CONNECT_TIMEOUT", "3s"); - - let config: Config = load_config(&[]).extract().unwrap(); - - assert_eq!(*config.grpc.connect_timeout, Duration::from_secs(3)); - Ok(()) - }); - } - #[test] #[allow(clippy::result_large_err)] fn test_config() { diff --git a/crates/arroyo-rpc/src/lib.rs b/crates/arroyo-rpc/src/lib.rs index 10dd665d..7acc7bb8 100644 --- a/crates/arroyo-rpc/src/lib.rs +++ b/crates/arroyo-rpc/src/lib.rs @@ -1004,7 +1004,7 @@ pub async fn grpc_channel_builder( target_tls: &Option, ) -> Result { let config = config(); - let endpoint = if let Some(target_tls) = config.get_tls_config(target_tls) { + if let Some(target_tls) = config.get_tls_config(target_tls) { let mut endpoint = Url::parse(&endpoint)?; endpoint .set_scheme("https") @@ -1031,13 +1031,11 @@ pub async fn grpc_channel_builder( config_builder = config_builder.identity(Identity::from_pem(our_tls.cert, our_tls.key)); } - b.tls_config(config_builder).context("configuring TLS")? + Ok(b.tls_config(config_builder).context("configuring TLS")?) } else { debug!("connecting to grpc endpoint {endpoint}"); - Channel::from_shared(endpoint.to_string())? - }; - - Ok(endpoint.connect_timeout(*config.grpc.connect_timeout)) + Ok(Channel::from_shared(endpoint.to_string())?) + } } /// Connect to a gRPC service with optional TLS From d8f67df0b4f152a43bb5aa7daddf7832259e6c82 Mon Sep 17 00:00:00 2001 From: Cole MacKenzie Date: Wed, 29 Jul 2026 08:54:13 -0700 Subject: [PATCH 5/5] Make gRPC connection timeout opt-in A mandatory global timeout changed connection behavior for every service, while a leader-only fallback prevented deployments from choosing tonic's unbounded default. Add an optional global `grpc.connect-timeout` setting that defaults to unset and apply it when building endpoints. Remove the leader-specific fallback so all clients consistently follow the configured policy. --- .../src/job_controller/leader_manager.rs | 24 +++--------- crates/arroyo-rpc/src/config.rs | 39 +++++++++++++++++++ crates/arroyo-rpc/src/lib.rs | 12 ++++-- 3 files changed, 54 insertions(+), 21 deletions(-) diff --git a/crates/arroyo-controller/src/job_controller/leader_manager.rs b/crates/arroyo-controller/src/job_controller/leader_manager.rs index 315c5a22..4e7868b4 100644 --- a/crates/arroyo-controller/src/job_controller/leader_manager.rs +++ b/crates/arroyo-controller/src/job_controller/leader_manager.rs @@ -18,8 +18,6 @@ use tonic::codegen::InterceptedService; use tonic::transport::Channel; use tracing::{info, warn}; -const LEADER_CONNECT_TIMEOUT: Duration = Duration::from_secs(10); - pub struct LeaderManager { leader_client: JobStatusGrpcClient>, pub job_id: JobId, @@ -37,23 +35,13 @@ impl LeaderManager { address: String, ) -> anyhow::Result { let leader_client = retry!( - match tokio::time::timeout( - LEADER_CONNECT_TIMEOUT, - job_status_client( - "controller", - &config().worker.tls, - worker_id, - address.clone(), - ), + job_status_client( + "controller", + &config().worker.tls, + worker_id, + address.clone() ) - .await - { - Ok(result) => result, - Err(_) => Err(anyhow!( - "timed out connecting to worker leader after {:?}", - LEADER_CONNECT_TIMEOUT - )), - }, + .await, 5, Duration::from_millis(100), Duration::from_secs(2), diff --git a/crates/arroyo-rpc/src/config.rs b/crates/arroyo-rpc/src/config.rs index 7238b677..35dfe0ac 100644 --- a/crates/arroyo-rpc/src/config.rs +++ b/crates/arroyo-rpc/src/config.rs @@ -193,6 +193,10 @@ fn load_config(paths: &[PathBuf]) -> Figment { #[derive(Debug, Deserialize, Serialize, Clone)] #[serde(rename_all = "kebab-case", deny_unknown_fields)] pub struct Config { + /// gRPC client configuration + #[serde(default)] + pub grpc: GrpcConfig, + /// API service configuration pub api: ApiConfig, @@ -275,6 +279,13 @@ pub struct Config { pub disable_telemetry: bool, } +#[derive(Debug, Default, Deserialize, Serialize, Clone)] +#[serde(rename_all = "kebab-case", deny_unknown_fields)] +pub struct GrpcConfig { + /// Maximum time to establish a gRPC connection + pub connect_timeout: Option, +} + #[derive(Debug, Deserialize, Serialize, Clone, Default)] #[serde(rename_all = "kebab-case", deny_unknown_fields)] pub enum JobControllerMode { @@ -1047,6 +1058,7 @@ impl TlsConfig { #[cfg(test)] mod tests { use crate::config::{Config, DatabaseType, Scheduler, SchemaName, SqliteConfig, load_config}; + use std::time::Duration; use url::Url; #[test] @@ -1073,6 +1085,33 @@ mod tests { } } + #[test] + #[allow(clippy::result_large_err)] + fn grpc_connect_timeout_is_unset_by_default() { + figment::Jail::expect_with(|_| { + let config: Config = load_config(&[]).extract().unwrap(); + + assert!(config.grpc.connect_timeout.is_none()); + Ok(()) + }); + } + + #[test] + #[allow(clippy::result_large_err)] + fn grpc_connect_timeout_can_be_overridden_with_environment() { + figment::Jail::expect_with(|jail| { + jail.set_env("ARROYO__GRPC__CONNECT_TIMEOUT", "3s"); + + let config: Config = load_config(&[]).extract().unwrap(); + + assert_eq!( + **config.grpc.connect_timeout.as_ref().unwrap(), + Duration::from_secs(3) + ); + Ok(()) + }); + } + #[test] #[allow(clippy::result_large_err)] fn test_config() { diff --git a/crates/arroyo-rpc/src/lib.rs b/crates/arroyo-rpc/src/lib.rs index 7acc7bb8..3ecbef43 100644 --- a/crates/arroyo-rpc/src/lib.rs +++ b/crates/arroyo-rpc/src/lib.rs @@ -1004,7 +1004,7 @@ pub async fn grpc_channel_builder( target_tls: &Option, ) -> Result { let config = config(); - if let Some(target_tls) = config.get_tls_config(target_tls) { + let endpoint = if let Some(target_tls) = config.get_tls_config(target_tls) { let mut endpoint = Url::parse(&endpoint)?; endpoint .set_scheme("https") @@ -1031,10 +1031,16 @@ pub async fn grpc_channel_builder( config_builder = config_builder.identity(Identity::from_pem(our_tls.cert, our_tls.key)); } - Ok(b.tls_config(config_builder).context("configuring TLS")?) + b.tls_config(config_builder).context("configuring TLS")? } else { debug!("connecting to grpc endpoint {endpoint}"); - Ok(Channel::from_shared(endpoint.to_string())?) + Channel::from_shared(endpoint.to_string())? + }; + + if let Some(connect_timeout) = config.grpc.connect_timeout.as_deref() { + Ok(endpoint.connect_timeout(*connect_timeout)) + } else { + Ok(endpoint) } }