(
formatter: OpenAIFormatter,
http_request: Request,
submissions: Sequence[
Tuple[
CompletionStream,
Callable[[], None],
]
],
)
| 15788 | stream.close() |
| 15789 | |
| 15790 | async def collect_completion_results( |
| 15791 | formatter: OpenAIFormatter, |
| 15792 | http_request: Request, |
| 15793 | submissions: Sequence[ |
| 15794 | Tuple[ |
| 15795 | CompletionStream, |
| 15796 | Callable[[], None], |
| 15797 | ] |
| 15798 | ], |
| 15799 | ) -> List[OpenAICompletion] | Response: |
| 15800 | streams = [stream for stream, _ in submissions] |
| 15801 | cancel_all = [cancel for _, cancel in submissions] |
| 15802 | |
| 15803 | def cancel_all_requests() -> None: |
| 15804 | for cancel in cancel_all: |
| 15805 | cancel() |
| 15806 | |
| 15807 | disconnect_task = asyncio.create_task( |
| 15808 | watch_http_disconnect( |
| 15809 | http_request, |
| 15810 | cancel_all_requests, |
| 15811 | ) |
| 15812 | ) |
| 15813 | try: |
| 15814 | return await asyncio.gather( |
| 15815 | *( |
| 15816 | asyncio.to_thread(formatter.collect_completion, stream) |
| 15817 | for stream in streams |
| 15818 | ) |
| 15819 | ) |
| 15820 | except asyncio.CancelledError: |
| 15821 | cancel_all_requests() |
| 15822 | raise |
| 15823 | except BaseException as exc: |
| 15824 | return await disconnected_cancelled_response_or_raise(http_request, exc) |
| 15825 | finally: |
| 15826 | disconnect_task.cancel() |
| 15827 | for stream in streams: |
| 15828 | stream.close() |
| 15829 | |
| 15830 | async def stream_sse_chunks( |
| 15831 | formatter: OpenAIFormatter, |
no test coverage detected
searching dependent graphs…