Skip to content

Commit 7873b92

Browse files
committed
feat: made aquila grpc timeout configurable
1 parent f1f3271 commit 7873b92

3 files changed

Lines changed: 32 additions & 6 deletions

File tree

crates/base/src/client/mod.rs

Lines changed: 16 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -22,17 +22,20 @@ pub struct DracoRuntimeStatusService {
2222
const MAX_BACKOFF: u64 = 2000 * 60;
2323
const MAX_RETRIES: i8 = 10;
2424

25-
// Will create a channel and retry if its not possible
26-
pub async fn create_channel_with_retry(channel_name: &str, url: String) -> Channel {
25+
pub async fn create_channel_with_retry(
26+
channel_name: &str,
27+
url: String,
28+
connect_timeout: std::time::Duration,
29+
request_timeout: std::time::Duration,
30+
) -> Channel {
2731
let mut backoff = 100;
2832
let mut retries = 0;
2933

3034
loop {
3135
let channel = match Endpoint::from_shared(url.clone()) {
3236
Ok(c) => {
3337
log::debug!("Creating a new endpoint for the: {} Service", channel_name);
34-
c.connect_timeout(std::time::Duration::from_secs(2))
35-
.timeout(std::time::Duration::from_secs(10))
38+
c.connect_timeout(connect_timeout).timeout(request_timeout)
3639
}
3740
Err(err) => {
3841
panic!(
@@ -67,8 +70,15 @@ pub async fn create_channel_with_retry(channel_name: &str, url: String) -> Chann
6770
}
6871
}
6972
impl DracoRuntimeStatusService {
70-
pub async fn from_url(aquila_url: String, aquila_token: String, identifier: String) -> Self {
71-
let channel = create_channel_with_retry("Aquila", aquila_url).await;
73+
pub async fn from_url(
74+
aquila_url: String,
75+
aquila_token: String,
76+
identifier: String,
77+
connect_timeout: std::time::Duration,
78+
request_timeout: std::time::Duration,
79+
) -> Self {
80+
let channel =
81+
create_channel_with_retry("Aquila", aquila_url, connect_timeout, request_timeout).await;
7282
Self::new(channel, identifier, aquila_token)
7383
}
7484

crates/base/src/config.rs

Lines changed: 12 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -67,6 +67,12 @@ pub struct AdapterConfig {
6767
/// Interval for runtime status heartbeat updates while the adapter is running.
6868
/// Set to 0 to disable periodic heartbeat updates.
6969
pub adapter_status_update_interval_seconds: u64,
70+
71+
/// Timeout in seconds for establishing Aquila gRPC channels.
72+
pub aquila_grpc_connect_timeout_secs: u64,
73+
74+
/// Timeout in seconds for Aquila gRPC requests.
75+
pub aquila_grpc_request_timeout_secs: u64,
7076
}
7177

7278
impl AdapterConfig {
@@ -103,6 +109,10 @@ impl AdapterConfig {
103109
"ADAPTER_STATUS_UPDATE_INTERVAL_SECONDS",
104110
30_u64,
105111
);
112+
let aquila_grpc_connect_timeout_secs =
113+
code0_flow::flow_config::env_with_default("AQUILA_GRPC_CONNECT_TIMEOUT_SECS", 2_u64);
114+
let aquila_grpc_request_timeout_secs =
115+
code0_flow::flow_config::env_with_default("AQUILA_GRPC_REQUEST_TIMEOUT_SECS", 10_u64);
106116
Self {
107117
environment,
108118
nats_bucket,
@@ -116,6 +126,8 @@ impl AdapterConfig {
116126
with_health_service,
117127
draco_variant,
118128
adapter_status_update_interval_seconds,
129+
aquila_grpc_connect_timeout_secs,
130+
aquila_grpc_request_timeout_secs,
119131
}
120132
}
121133

crates/base/src/runner.rs

Lines changed: 4 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -67,6 +67,8 @@ impl<C: LoadConfig> ServerRunner<C> {
6767
config.aquila_url.clone(),
6868
config.aquila_token.clone(),
6969
config.draco_variant.clone(),
70+
Duration::from_secs(config.aquila_grpc_connect_timeout_secs),
71+
Duration::from_secs(config.aquila_grpc_request_timeout_secs),
7072
)
7173
.await,
7274
));
@@ -84,6 +86,8 @@ impl<C: LoadConfig> ServerRunner<C> {
8486
config.aquila_url.clone(),
8587
config.definition_path.as_str(),
8688
config.aquila_token.clone(),
89+
Duration::from_secs(config.aquila_grpc_connect_timeout_secs),
90+
Duration::from_secs(config.aquila_grpc_request_timeout_secs),
8791
)
8892
.await
8993
.with_definition_source(service_name)

0 commit comments

Comments
 (0)