Sends a SendMessage request to the backend, returning the response stream to consume. You should repeatedly call [Self::recv] to receive [ResponseEvent]'s until a [ResponseEvent::EndStream] value is returned. # Arguments `client` - api client to make the request with `conversation_state` - the [crate::api_client::model::ConversationState] to send `request_metadata_lock` - a mutex that will be u
(
client: &ApiClient,
conversation_state: ConversationState,
request_metadata_lock: Arc<Mutex<Option<RequestMetadata>>>,
message_meta_tags: Option<Vec<MessageMetaTag>>,
| 205 | /// future is aborted in the sigint case). The task will gracefully end with updating the mutex |
| 206 | /// with [RequestMetadata]. |
| 207 | pub async fn send_message( |
| 208 | client: &ApiClient, |
| 209 | conversation_state: ConversationState, |
| 210 | request_metadata_lock: Arc<Mutex<Option<RequestMetadata>>>, |
| 211 | message_meta_tags: Option<Vec<MessageMetaTag>>, |
| 212 | ) -> Result<Self, SendMessageError> { |
| 213 | let message_id = uuid::Uuid::new_v4().to_string(); |
| 214 | info!(?message_id, "Generated new message id"); |
| 215 | let user_prompt_length = conversation_state.user_input_message.content.len(); |
| 216 | let model_id = conversation_state.user_input_message.model_id.clone(); |
| 217 | let message_meta_tags = message_meta_tags.unwrap_or_default(); |
| 218 | |
| 219 | let cancel_token = CancellationToken::new(); |
| 220 | let cancel_token_clone = cancel_token.clone(); |
| 221 | |
| 222 | let start_time = Instant::now(); |
| 223 | let start_time_sys = SystemTime::now(); |
| 224 | debug!(?start_time, "sending send_message request"); |
| 225 | let response = client |
| 226 | .send_message(conversation_state) |
| 227 | .await |
| 228 | .map_err(|err| SendMessageError { |
| 229 | source: err, |
| 230 | request_metadata: RequestMetadata { |
| 231 | message_id: message_id.clone(), |
| 232 | request_start_timestamp_ms: system_time_to_unix_ms(start_time_sys), |
| 233 | stream_end_timestamp_ms: system_time_to_unix_ms(SystemTime::now()), |
| 234 | model_id: model_id.clone(), |
| 235 | user_prompt_length, |
| 236 | message_meta_tags: message_meta_tags.clone(), |
| 237 | // Other fields are irrelevant if we can't get a successful response |
| 238 | ..Default::default() |
| 239 | }, |
| 240 | })?; |
| 241 | let elapsed = start_time.elapsed(); |
| 242 | debug!(?elapsed, "send_message succeeded"); |
| 243 | |
| 244 | let request_id = response.request_id().map(str::to_string); |
| 245 | let (ev_tx, ev_rx) = mpsc::channel(16); |
| 246 | tokio::spawn(async move { |
| 247 | ResponseParser::new( |
| 248 | response, |
| 249 | message_id, |
| 250 | model_id, |
| 251 | user_prompt_length, |
| 252 | message_meta_tags, |
| 253 | ev_tx, |
| 254 | start_time, |
| 255 | start_time_sys, |
| 256 | cancel_token_clone, |
| 257 | request_metadata_lock, |
| 258 | ) |
| 259 | .try_recv() |
| 260 | .await; |
| 261 | }); |
| 262 | |
| 263 | Ok(Self { |
| 264 | request_id, |
nothing calls this directly
no test coverage detected