(&self, message: String, agent_type: &str)
| 208 | } |
| 209 | |
| 210 | async fn send_message(&self, message: String, agent_type: &str) -> Result<String> { |
| 211 | let session_id = self.ensure_session(agent_type).await?; |
| 212 | tracing::info!("Sending message to session {}: {}", session_id, message); |
| 213 | |
| 214 | // Generate a turn_id |
| 215 | let turn_id = uuid::Uuid::new_v4().to_string(); |
| 216 | |
| 217 | // Store current turn_id for cancellation |
| 218 | { |
| 219 | let mut turn_guard = self.current_turn_id.lock().await; |
| 220 | *turn_guard = Some(turn_id.clone()); |
| 221 | } |
| 222 | |
| 223 | // Start the dialog turn — this is async, events will arrive via EventQueue |
| 224 | let start_result = self |
| 225 | .coordinator |
| 226 | .start_dialog_turn( |
| 227 | session_id.clone(), |
| 228 | message.clone(), |
| 229 | None, |
| 230 | Some(turn_id.clone()), |
| 231 | agent_type.to_string(), |
| 232 | Some(self.workspace_path_string()), |
| 233 | None, |
| 234 | None, |
| 235 | DialogSubmissionPolicy::for_source(DialogTriggerSource::Cli), |
| 236 | None, |
| 237 | ) |
| 238 | .await; |
| 239 | |
| 240 | if let Err(err) = start_result { |
| 241 | if Self::is_session_not_found_error(&err.to_string()) { |
| 242 | tracing::warn!( |
| 243 | "Session missing when starting turn, attempting recovery and retry: session_id={}, error={}", |
| 244 | session_id, |
| 245 | err |
| 246 | ); |
| 247 | self.ensure_backend_session_alive(&session_id, agent_type) |
| 248 | .await?; |
| 249 | self.coordinator |
| 250 | .start_dialog_turn( |
| 251 | session_id, |
| 252 | message, |
| 253 | None, |
| 254 | Some(turn_id.clone()), |
| 255 | agent_type.to_string(), |
| 256 | Some(self.workspace_path_string()), |
| 257 | None, |
| 258 | None, |
| 259 | DialogSubmissionPolicy::for_source(DialogTriggerSource::Cli), |
| 260 | None, |
| 261 | ) |
| 262 | .await?; |
| 263 | } else { |
| 264 | return Err(err.into()); |
| 265 | } |
| 266 | } |
| 267 |
no test coverage detected