diff --git a/src/domain/agent/api.rs b/src/domain/agent/api.rs index f886c34..cbd93a8 100644 --- a/src/domain/agent/api.rs +++ b/src/domain/agent/api.rs @@ -75,6 +75,8 @@ pub async fn perform_search( tracing::info!(query = %query, "Performing Tavily web search"); let client = reqwest::Client::builder() .timeout(std::time::Duration::from_secs(30)) + .connect_timeout(std::time::Duration::from_secs(10)) + .pool_idle_timeout(std::time::Duration::from_secs(60)) .build()?; let response = client .post("https://api.tavily.com/search") diff --git a/src/domain/agent/mod.rs b/src/domain/agent/mod.rs index dd87dba..8f738cc 100644 --- a/src/domain/agent/mod.rs +++ b/src/domain/agent/mod.rs @@ -77,6 +77,9 @@ impl Agent { let client = reqwest::Client::builder() .timeout(std::time::Duration::from_secs(120)) + .connect_timeout(std::time::Duration::from_secs(10)) + .tcp_keepalive(std::time::Duration::from_secs(30)) + .pool_idle_timeout(std::time::Duration::from_secs(60)) .build() .map_err(|e| AppError::Internal(format!("Failed to build HTTP client: {}", e)))?; @@ -119,6 +122,8 @@ impl Agent { return Err(AppError::Internal("Agent run exceeded max turns".into())); } + tracing::info!("Turn {}", turns); + turns += 1; let current_role = self .messages @@ -130,15 +135,24 @@ impl Agent { turns, current_role )); - let assistant_message = self.execute_turn().await?; + let assistant_message = + match tokio::time::timeout(Duration::from_secs(180), self.execute_turn()).await { + Ok(res) => res?, + Err(_) => { + tracing::error!("Agent execution turn timed out after 180s"); + return Err(AppError::Internal("Agent execution turn timed out".into())); + } + }; if let Some(tool_calls) = &assistant_message.tool_calls { + tracing::info!("Assistant tool calls: {:#?}", tool_calls); if tool_calls.iter().any(|tc| tc.function.name == "answer") { finished = true; } } if self.answer.is_some() { + tracing::info!("Answer: {}", self.answer.as_ref().unwrap()); finished = true; } @@ -152,6 +166,8 @@ impl Agent { tool_call_id: None, }); } + + tracing::info!("Finished turn"); } self.log("\n--- Execution Finished ---"); @@ -163,6 +179,8 @@ impl Agent { let mut sub_turns = 0; loop { + tracing::info!("Sub turn {}", sub_turns); + sub_turns += 1; if sub_turns > max_sub_turns { return Err(AppError::Internal( @@ -178,15 +196,19 @@ impl Agent { .message .clone(); + tracing::info!("Assistant message: {:#?}", assistant_message); + self.messages.push(assistant_message.clone()); if let Some(content) = &assistant_message.content { + tracing::info!("Assistant content: {}", content); if !content.is_empty() { self.log(&format!("\nAssistant: {}", content)); } } if let Some(tool_calls) = &assistant_message.tool_calls { + tracing::info!("Assistant tool calls: {:#?}", tool_calls); let mut is_final_cycle = false; let mut final_answer = None; @@ -220,6 +242,7 @@ impl Agent { } if is_final_cycle { + tracing::info!("Final answer: {}", final_answer.as_ref().unwrap()); return Ok(Message { role: "assistant".to_string(), content: final_answer.or(assistant_message.content), @@ -242,7 +265,11 @@ impl Agent { tools: self.tools.clone(), }; - let mut request_builder = self.client.post(&self.url).json(&request); + let mut request_builder = self + .client + .post(&self.url) + .json(&request) + .timeout(Duration::from_secs(60)); if let Some(key) = &self.zen_api_key { request_builder = request_builder.header("Authorization", format!("Bearer {}", key)); @@ -252,10 +279,12 @@ impl Agent { let response = request_builder.send().await.map_err(|e| { let duration = start.elapsed(); let is_timeout = e.is_timeout(); + let is_connect = e.is_connect(); tracing::error!( - "Network error after {:?} during LLM call (Timeout: {}): {:?}", + "Network error after {:?} during LLM call (Timeout: {}, Connect: {}): {:?}", duration, is_timeout, + is_connect, e ); AppError::Network(e) diff --git a/src/domain/auth.rs b/src/domain/auth.rs index 10a1255..a180185 100644 --- a/src/domain/auth.rs +++ b/src/domain/auth.rs @@ -46,7 +46,10 @@ pub struct JwksVerifier { impl JwksVerifier { pub async fn new(issuer: String, audience: String) -> Result> { - let client = Client::new(); + let client = Client::builder() + .timeout(std::time::Duration::from_secs(30)) + .connect_timeout(std::time::Duration::from_secs(10)) + .build()?; let discovery_url = format!( "{}/.well-known/openid-configuration", issuer.trim_end_matches('/') @@ -123,7 +126,10 @@ impl Authenticator { client_id: String, client_secret: String, ) -> Result> { - let client = Client::new(); + let client = Client::builder() + .timeout(std::time::Duration::from_secs(30)) + .connect_timeout(std::time::Duration::from_secs(10)) + .build()?; let discovery_url = format!( "{}/.well-known/openid-configuration", issuer.trim_end_matches('/') diff --git a/src/domain/calendar/mod.rs b/src/domain/calendar/mod.rs index d6176d1..c6f49ae 100644 --- a/src/domain/calendar/mod.rs +++ b/src/domain/calendar/mod.rs @@ -36,9 +36,15 @@ pub struct CalendarClient { impl CalendarClient { pub fn new(base_url: String, authenticator: Arc) -> Self { + let client = Client::builder() + .timeout(std::time::Duration::from_secs(30)) + .connect_timeout(std::time::Duration::from_secs(10)) + .build() + .unwrap_or_else(|_| Client::new()); + Self { base_url: base_url.trim_end_matches('/').to_string(), - client: Client::new(), + client, authenticator, token_state: RwLock::new(None), } diff --git a/src/domain/tasks.rs b/src/domain/tasks.rs index 2ea8b26..bd518ef 100644 --- a/src/domain/tasks.rs +++ b/src/domain/tasks.rs @@ -61,7 +61,7 @@ pub async fn execute_agent_run( goal: String, ) -> AppResult { let run_id = Uuid::new_v4(); - tracing::info!(%task_id, %run_id, "Starting agent execution run"); + tracing::info!(%task_id, %run_id, "Starting background agent execution run"); let new_run = task_run::ActiveModel { id: Set(run_id), @@ -98,98 +98,87 @@ pub async fn execute_agent_run( goal.clone(), )?; - let (logs, answer, status) = match agent.run(config).await { - Ok((logs, answer)) => { - tracing::info!(%task_id, %run_id, "Agent execution completed successfully"); - (logs, answer, "completed".to_string()) + let db_bg = db.clone(); + let config_bg = config.clone(); + let scheduler_bg = _scheduler.clone(); + let task_id_bg = task_id; + + tokio::spawn(async move { + let (logs, answer, status) = match agent.run(&config_bg).await { + Ok((logs, answer)) => { + tracing::info!(task_id = %task_id_bg, run_id = %run_id, "Agent execution completed successfully"); + (logs, answer, "completed".to_string()) + } + Err(e) => { + tracing::error!(task_id = %task_id_bg, run_id = %run_id, error = %e, "Agent execution failed"); + ( + format!("Execution failed: {}", e), + None, + "failed".to_string(), + ) + } + }; + + let run_update = task_run::ActiveModel { + id: Set(run_id), + logs: Set(logs), + answer: Set(answer), + status: Set(status.clone()), + ..Default::default() + }; + + if let Err(e) = run_update.update(&db_bg).await { + tracing::error!(task_id = %task_id_bg, run_id = %run_id, error = %e, "Failed to update run record"); } - Err(e) => { - tracing::error!(%task_id, %run_id, error = %e, "Agent execution failed"); - ( - format!("Execution failed: {}", e), - None, - "failed".to_string(), - ) - } - }; - let run: task_run::ActiveModel = TaskRun::find_by_id(run_id) - .one(db) - .await - .map_err(crate::error::AppError::Database)? - .ok_or_else(|| crate::error::AppError::NotFound("Run not found after insert".into()))? - .into(); + if let Ok(task_response) = get_task_inner(task_id_bg, &db_bg).await { + let _ = scheduler_bg + .tx + .send(crate::server::notifications::WsEvent::RunFinished( + task_response.clone(), + )); - let mut run = run; - run.logs = Set(logs.clone()); - run.answer = Set(answer.clone()); - run.status = Set(status.clone()); + // Send Push Notifications to subscribers + if let Ok(subscriptions) = task_subscription::Entity::find() + .filter(task_subscription::Column::TaskId.eq(task_id_bg)) + .all(&db_bg) + .await + { + for sub in subscriptions { + if let Ok(push_subs) = push_subscription::Entity::find() + .filter(push_subscription::Column::UserSub.eq(sub.user_sub.clone())) + .all(&db_bg) + .await + { + for push_sub in push_subs { + let sender = scheduler_bg.push_sender.clone(); + let goal = task_response.goal.clone(); + let status_bg = status.clone(); + let sub_data = crate::domain::notifications::push::PushSubscription { + endpoint: push_sub.endpoint, + p256dh: push_sub.p256dh, + auth: push_sub.auth, + }; - run.update(db) - .await - .map_err(crate::error::AppError::Database)?; - - let task_response = get_task_inner(task_id, db).await?; - let _ = _scheduler - .tx - .send(crate::server::notifications::WsEvent::RunFinished( - task_response.clone(), - )); - - // Send Push Notifications to subscribers - let subscriptions = task_subscription::Entity::find() - .filter(task_subscription::Column::TaskId.eq(task_id)) - .all(db) - .await - .map_err(crate::error::AppError::Database)?; - - tracing::info!( - "Found {} task subscriptions for task {}", - subscriptions.len(), - task_id - ); - - for sub in subscriptions { - let push_subs = push_subscription::Entity::find() - .filter(push_subscription::Column::UserSub.eq(sub.user_sub.clone())) - .all(db) - .await - .map_err(crate::error::AppError::Database)?; - - tracing::info!( - "Found {} push subscriptions for user {}", - push_subs.len(), - sub.user_sub - ); - - for push_sub in push_subs { - let sender = _scheduler.push_sender.clone(); - let goal = task_response.goal.clone(); - let status = status.clone(); - let sub_data = crate::domain::notifications::push::PushSubscription { - endpoint: push_sub.endpoint, - p256dh: push_sub.p256dh, - auth: push_sub.auth, - }; - - tokio::spawn(async move { - if let Err(e) = sender - .send_notification( - &sub_data, - &format!("Task Completed: {}", status), - &goal, - Some(task_id), - Some(run_id), - ) - .await - { - tracing::error!("Failed to send notification in background task: {}", e); + tokio::spawn(async move { + let _ = sender + .send_notification( + &sub_data, + &format!("Task Completed: {}", status_bg), + &goal, + Some(task_id_bg), + Some(run_id), + ) + .await; + }); + } + } } - }); + } } - } + }); - Ok(task_response) + get_task_inner(task_id, db).await } pub async fn get_task_inner(id: Uuid, db: &DatabaseConnection) -> AppResult { diff --git a/src/server/chat.rs b/src/server/chat.rs index d6b6bcd..6e24d4b 100644 --- a/src/server/chat.rs +++ b/src/server/chat.rs @@ -52,7 +52,15 @@ pub async fn chat_handler( )?; tracing::info!("Starting interactive agent turn"); - let assistant_message = agent.execute_turn().await?; + let assistant_message = + match tokio::time::timeout(std::time::Duration::from_secs(180), agent.execute_turn()).await + { + Ok(res) => res?, + Err(_) => { + tracing::error!("Interactive chat agent turn timed out after 180s"); + return Err(AppError::Internal("Agent turn timed out".into())); + } + }; Ok(Json(ChatResult { message: assistant_message,