(&mut self)
| 605 | } |
| 606 | |
| 607 | async fn next(&mut self) -> Result<Option<ChatResponseStream>, RecvError> { |
| 608 | if let Some(ev) = self.peek.take() { |
| 609 | return Ok(Some(ev)); |
| 610 | } |
| 611 | |
| 612 | trace!("Attempting to recv next event"); |
| 613 | let start = Instant::now(); |
| 614 | let result = self.response.recv().await; |
| 615 | let duration = Instant::now().duration_since(start); |
| 616 | match result { |
| 617 | Ok(ev) => { |
| 618 | trace!(?ev, "Received new event"); |
| 619 | |
| 620 | if !self.message_start_pushed { |
| 621 | self.buf |
| 622 | .push(StreamResult::Ok(StreamEvent::MessageStart(MessageStartEvent { |
| 623 | role: Role::Assistant, |
| 624 | }))); |
| 625 | self.message_start_pushed = true; |
| 626 | } |
| 627 | |
| 628 | // Track metadata about the chunk. |
| 629 | self.time_to_first_chunk |
| 630 | .get_or_insert_with(|| self.request_start_time.elapsed()); |
| 631 | self.time_between_chunks.push(duration); |
| 632 | self.received_response_size += ev.as_ref().map(|e| e.len()).unwrap_or_default(); |
| 633 | |
| 634 | Ok(ev) |
| 635 | }, |
| 636 | Err(err) => { |
| 637 | error!(?err, "failed to receive the next event"); |
| 638 | if duration.as_secs() >= 59 { |
| 639 | Err(RecvError::Timeout { source: err, duration }) |
| 640 | } else { |
| 641 | Err(RecvError::Other { source: err }) |
| 642 | } |
| 643 | }, |
| 644 | } |
| 645 | } |
| 646 | |
| 647 | fn recv_error_to_stream_error(&self, err: RecvError) -> StreamError { |
| 648 | match err { |
no test coverage detected