| Server IP : 216.250.13.251 / Your IP : 216.73.216.164 Web Server : nginx/1.24.0 System : Linux gunay 6.8.0-134-generic #134-Ubuntu SMP PREEMPT_DYNAMIC Fri Jun 26 18:43:11 UTC 2026 x86_64 User : root ( 0) PHP Version : 7.4.33 Disable Function : pcntl_alarm,pcntl_fork,pcntl_waitpid,pcntl_wait,pcntl_wifexited,pcntl_wifstopped,pcntl_wifsignaled,pcntl_wifcontinued,pcntl_wexitstatus,pcntl_wtermsig,pcntl_wstopsig,pcntl_signal,pcntl_signal_get_handler,pcntl_signal_dispatch,pcntl_get_last_error,pcntl_strerror,pcntl_sigprocmask,pcntl_sigwaitinfo,pcntl_sigtimedwait,pcntl_exec,pcntl_getpriority,pcntl_setpriority,pcntl_async_signals,pcntl_unshare, MySQL : OFF | cURL : ON | WGET : ON | Perl : ON | Python : OFF | Sudo : ON | Pkexec : OFF Directory : /lib/python3/dist-packages/twisted/trial/_dist/ |
Upload File : |
"""
Buffer byte streams.
"""
from itertools import count
from typing import Dict, Iterator, List, TypeVar
from attrs import Factory, define
from twisted.protocols.amp import AMP, Command, Integer, String as Bytes
T = TypeVar("T")
class StreamOpen(Command):
"""
Open a new stream.
"""
response = [(b"streamId", Integer())]
class StreamWrite(Command):
"""
Write a chunk of data to a stream.
"""
arguments = [
(b"streamId", Integer()),
(b"data", Bytes()),
]
@define
class StreamReceiver:
"""
Buffering de-multiplexing byte stream receiver.
"""
_counter: Iterator[int] = count()
_streams: Dict[int, List[bytes]] = Factory(dict)
def open(self) -> int:
"""
Open a new stream and return its unique identifier.
"""
newId = next(self._counter)
self._streams[newId] = []
return newId
def write(self, streamId: int, chunk: bytes) -> None:
"""
Write to an open stream using its unique identifier.
@raise KeyError: If there is no such open stream.
"""
self._streams[streamId].append(chunk)
def finish(self, streamId: int) -> List[bytes]:
"""
Indicate an open stream may receive no further data and return all of
its current contents.
@raise KeyError: If there is no such open stream.
"""
return self._streams.pop(streamId)
def chunk(data: bytes, chunkSize: int) -> Iterator[bytes]:
"""
Break a byte string into pieces of no more than ``chunkSize`` length.
@param data: The byte string.
@param chunkSize: The maximum length of the resulting pieces. All pieces
except possibly the last will be this length.
@return: The pieces.
"""
pos = 0
while pos < len(data):
yield data[pos : pos + chunkSize]
pos += chunkSize
async def stream(amp: AMP, chunks: Iterator[bytes]) -> int:
"""
Send the given stream chunks, one by one, over the given connection.
The chunks are sent using L{StreamWrite} over a stream opened using
L{StreamOpen}.
@return: The identifier of the stream over which the chunks were sent.
"""
streamId = (await amp.callRemote(StreamOpen))["streamId"]
assert isinstance(streamId, int)
for oneChunk in chunks:
await amp.callRemote(StreamWrite, streamId=streamId, data=oneChunk)
return streamId