Team Ai
Datasetpublic

codekingpro/portable-devtools

sourceHugging Faceupdated 5mo agoView on Hugging Face
1likes14kdownloads
memory_stream.py149 linesDownload Raw Back to tracers
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 
codekingpro/portable-devtools · Team Ai