codekingpro/portable-devtools
114k
1"""Module implements a memory stream for communication between two co-routines.2 3This module provides a way to communicate between two co-routines using a memory4channel. The writer and reader can be in the same event loop or in different event5loops. When they're in different event loops, they will also be in different threads.6 7Useful in situations when there's a mix of synchronous and asynchronous used in the8code.9"""10 11import asyncio12from asyncio import AbstractEventLoop, Queue13from collections.abc import AsyncIterator14from typing import Generic, TypeVar15 16T = TypeVar("T")17 18 19class _SendStream(Generic[T]):20 def __init__(21 self, reader_loop: AbstractEventLoop, queue: Queue, done: object22 ) -> None:23 """Create a writer for the queue and done object.24 25 Args:26 reader_loop: The event loop to use for the writer.27 28 This loop will be used to schedule the writes to the queue.29 queue: The queue to write to.30 31 This is an asyncio queue.32 done: Special sentinel object to indicate that the writer is done.33 """34 self._reader_loop = reader_loop35 self._queue = queue36 self._done = done37 38 async def send(self, item: T) -> None:39 """Schedule the item to be written to the queue using the original loop.40 41 This is a coroutine that can be awaited.42 43 Args:44 item: The item to write to the queue.45 """46 return self.send_nowait(item)47 48 def send_nowait(self, item: T) -> None:49 """Schedule the item to be written to the queue using the original loop.50 51 This is a non-blocking call.52 53 Args:54 item: The item to write to the queue.55 56 Raises:57 RuntimeError: If the event loop is already closed when trying to write to58 the queue.59 """60 try:61 self._reader_loop.call_soon_threadsafe(self._queue.put_nowait, item)62 except RuntimeError:63 if not self._reader_loop.is_closed():64 raise # Raise the exception if the loop is not closed65 66 async def aclose(self) -> None:67 """Async schedule the done object write the queue using the original loop."""68 return self.close()69 70 def close(self) -> None:71 """Schedule the done object write the queue using the original loop.72 73 This is a non-blocking call.74 75 Raises:76 RuntimeError: If the event loop is already closed when trying to write to77 the queue.78 """79 try:80 self._reader_loop.call_soon_threadsafe(self._queue.put_nowait, self._done)81 except RuntimeError:82 if not self._reader_loop.is_closed():83 raise # Raise the exception if the loop is not closed84 85 86class _ReceiveStream(Generic[T]):87 def __init__(self, queue: Queue, done: object) -> None:88 """Create a reader for the queue and done object.89 90 This reader should be used in the same loop as the loop that was passed to the91 channel.92 """93 self._queue = queue94 self._done = done95 self._is_closed = False96 97 async def __aiter__(self) -> AsyncIterator[T]:98 while True:99 item = await self._queue.get()100 if item is self._done:101 self._is_closed = True102 break103 yield item104 105 106class _MemoryStream(Generic[T]):107 """Stream data from a writer to a reader even if they are in different threads.108 109 Uses asyncio queues to communicate between two co-routines. This implementation110 should work even if the writer and reader co-routines belong to two different event111 loops (e.g. one running from an event loop in the main thread and the other running112 in an event loop in a background thread).113 114 This implementation is meant to be used with a single writer and a single reader.115 116 This is an internal implementation to LangChain. Do not use it directly.117 """118 119 def __init__(self, loop: AbstractEventLoop) -> None:120 """Create a channel for the given loop.121 122 Args:123 loop: The event loop to use for the channel.124 125 The reader is assumed to be running in the same loop as the one passed126 to this constructor. This will NOT be validated at run time.127 """128 self._loop = loop129 self._queue: asyncio.Queue = asyncio.Queue(maxsize=0)130 self._done = object()131 132 def get_send_stream(self) -> _SendStream[T]:133 """Get a writer for the channel.134 135 Returns:136 The writer for the channel.137 """138 return _SendStream[T](139 reader_loop=self._loop, queue=self._queue, done=self._done140 )141 142 def get_receive_stream(self) -> _ReceiveStream[T]:143 """Get a reader for the channel.144 145 Returns:146 The reader for the channel.147 """148 return _ReceiveStream[T](queue=self._queue, done=self._done)149 