2015-09-02 21:31:44 +00:00
|
|
|
WHENCE_ABSOLUTE = 0
|
|
|
|
WHENCE_RELATIVE = 1
|
|
|
|
WHENCE_RELATIVE_END = 2
|
|
|
|
|
|
|
|
READ_UNTIL_END = -1
|
|
|
|
|
|
|
|
|
|
|
|
class BaseStreamFilelike(object):
|
|
|
|
def __init__(self, fileobj):
|
|
|
|
self._fileobj = fileobj
|
|
|
|
self._cursor_position = 0
|
|
|
|
|
2015-09-28 19:27:56 +00:00
|
|
|
def close(self):
|
|
|
|
self._fileobj.close()
|
|
|
|
|
2015-09-02 21:31:44 +00:00
|
|
|
def read(self, size=READ_UNTIL_END):
|
|
|
|
buf = self._fileobj.read(size)
|
|
|
|
self._cursor_position += len(buf)
|
|
|
|
return buf
|
|
|
|
|
|
|
|
def tell(self):
|
|
|
|
return self._cursor_position
|
|
|
|
|
|
|
|
def seek(self, index, whence=WHENCE_ABSOLUTE):
|
|
|
|
num_bytes_to_ff = 0
|
|
|
|
if whence == WHENCE_ABSOLUTE:
|
|
|
|
if index < self._cursor_position:
|
|
|
|
raise IOError('Cannot seek backwards')
|
|
|
|
num_bytes_to_ff = index - self._cursor_position
|
|
|
|
|
|
|
|
elif whence == WHENCE_RELATIVE:
|
|
|
|
if index < 0:
|
|
|
|
raise IOError('Cannnot seek backwards')
|
|
|
|
num_bytes_to_ff = index
|
|
|
|
|
|
|
|
elif whence == WHENCE_RELATIVE_END:
|
|
|
|
raise IOError('Stream does not have a known end point')
|
|
|
|
|
2015-09-30 18:19:25 +00:00
|
|
|
bytes_forward = num_bytes_to_ff
|
2015-09-02 21:31:44 +00:00
|
|
|
while num_bytes_to_ff > 0:
|
|
|
|
buf = self._fileobj.read(num_bytes_to_ff)
|
|
|
|
if not buf:
|
|
|
|
raise IOError('Seek past end of file')
|
|
|
|
num_bytes_to_ff -= len(buf)
|
|
|
|
|
2015-09-30 18:19:25 +00:00
|
|
|
self._cursor_position += bytes_forward
|
|
|
|
return bytes_forward
|
|
|
|
|
2015-09-02 21:31:44 +00:00
|
|
|
|
|
|
|
class SocketReader(BaseStreamFilelike):
|
|
|
|
def __init__(self, fileobj):
|
|
|
|
super(SocketReader, self).__init__(fileobj)
|
2015-08-12 20:39:32 +00:00
|
|
|
self.handlers = []
|
|
|
|
|
|
|
|
def add_handler(self, handler):
|
|
|
|
self.handlers.append(handler)
|
|
|
|
|
2015-09-02 21:31:44 +00:00
|
|
|
def read(self, size=READ_UNTIL_END):
|
|
|
|
buf = super(SocketReader, self).read(size)
|
2015-08-12 20:39:32 +00:00
|
|
|
for handler in self.handlers:
|
|
|
|
handler(buf)
|
|
|
|
return buf
|
|
|
|
|
|
|
|
|
2015-08-25 19:53:13 +00:00
|
|
|
def wrap_with_handler(in_fp, handler):
|
2015-08-12 20:39:32 +00:00
|
|
|
wrapper = SocketReader(in_fp)
|
2015-08-25 19:53:13 +00:00
|
|
|
wrapper.add_handler(handler)
|
2015-08-12 20:39:32 +00:00
|
|
|
return wrapper
|
2015-09-02 21:31:44 +00:00
|
|
|
|
|
|
|
|
2015-09-30 18:19:25 +00:00
|
|
|
class FilelikeStreamConcat(object):
|
|
|
|
""" A file-like object which concats all the file-like objects in the specified generator into
|
|
|
|
a single stream.
|
|
|
|
"""
|
2015-09-02 21:31:44 +00:00
|
|
|
def __init__(self, file_generator):
|
|
|
|
self._file_generator = file_generator
|
|
|
|
self._current_file = file_generator.next()
|
2015-09-30 18:19:25 +00:00
|
|
|
self._current_position = 0
|
|
|
|
self._closed = False
|
2015-09-02 21:31:44 +00:00
|
|
|
|
2015-09-30 18:19:25 +00:00
|
|
|
def tell(self):
|
|
|
|
return self._current_position
|
2015-09-02 21:31:44 +00:00
|
|
|
|
2015-09-30 18:19:25 +00:00
|
|
|
def close(self):
|
|
|
|
self._closed = True
|
2015-09-25 18:57:14 +00:00
|
|
|
|
|
|
|
def read(self, size=READ_UNTIL_END):
|
2015-09-30 18:19:25 +00:00
|
|
|
buf = ''
|
|
|
|
current_size = size
|
|
|
|
|
|
|
|
while size == READ_UNTIL_END or len(buf) < size:
|
|
|
|
current_buf = self._current_file.read(current_size)
|
|
|
|
if current_buf:
|
|
|
|
buf += current_buf
|
|
|
|
self._current_position += len(current_buf)
|
|
|
|
if size != READ_UNTIL_END:
|
|
|
|
current_size -= len(current_buf)
|
|
|
|
|
|
|
|
else:
|
|
|
|
# That file was out of data, prime a new one
|
|
|
|
self._current_file.close()
|
|
|
|
try:
|
|
|
|
self._current_file = self._file_generator.next()
|
|
|
|
except StopIteration:
|
|
|
|
return buf
|
2015-09-25 18:57:14 +00:00
|
|
|
|
2015-09-30 18:19:25 +00:00
|
|
|
return buf
|
2015-09-25 18:57:14 +00:00
|
|
|
|
|
|
|
|
2015-09-02 21:31:44 +00:00
|
|
|
class StreamSlice(BaseStreamFilelike):
|
2015-09-30 18:19:25 +00:00
|
|
|
""" A file-like object which returns a file-like object that represents a slice of the data in
|
|
|
|
the specified file obj. All methods will act as if the slice is its own file.
|
|
|
|
"""
|
|
|
|
|
2015-09-02 21:31:44 +00:00
|
|
|
def __init__(self, fileobj, start_offset=0, end_offset_exclusive=READ_UNTIL_END):
|
|
|
|
super(StreamSlice, self).__init__(fileobj)
|
|
|
|
self._end_offset_exclusive = end_offset_exclusive
|
2015-09-30 18:19:25 +00:00
|
|
|
self._start_offset = start_offset
|
2015-09-02 21:31:44 +00:00
|
|
|
|
|
|
|
if start_offset > 0:
|
|
|
|
self.seek(start_offset)
|
|
|
|
|
|
|
|
def read(self, size=READ_UNTIL_END):
|
|
|
|
if self._end_offset_exclusive == READ_UNTIL_END:
|
|
|
|
# We weren't asked to limit the end of the stream
|
|
|
|
return super(StreamSlice, self).read(size)
|
|
|
|
|
|
|
|
# Compute the max bytes to read until the end or until we reach the user requested max
|
2015-09-30 18:19:25 +00:00
|
|
|
max_bytes_to_read = self._end_offset_exclusive - super(StreamSlice, self).tell()
|
2015-09-02 21:31:44 +00:00
|
|
|
if size != READ_UNTIL_END:
|
|
|
|
max_bytes_to_read = min(max_bytes_to_read, size)
|
|
|
|
|
|
|
|
return super(StreamSlice, self).read(max_bytes_to_read)
|
2015-09-30 18:19:25 +00:00
|
|
|
|
|
|
|
def _file_min(self, first, second):
|
|
|
|
if first == READ_UNTIL_END:
|
|
|
|
return second
|
|
|
|
|
|
|
|
if second == READ_UNTIL_END:
|
|
|
|
return first
|
|
|
|
|
|
|
|
return min(first, second)
|
|
|
|
|
|
|
|
def tell(self):
|
|
|
|
return super(StreamSlice, self).tell() - self._start_offset
|
|
|
|
|
|
|
|
def seek(self, index, whence=WHENCE_ABSOLUTE):
|
|
|
|
index = self._file_min(self._end_offset_exclusive, index)
|
|
|
|
super(StreamSlice, self).seek(index, whence)
|
|
|
|
|
|
|
|
|
|
|
|
class LimitingStream(StreamSlice):
|
|
|
|
""" A file-like object which mimics the specified file stream being limited to the given number
|
|
|
|
of bytes. All calls after that limit (if specified) will act as if the file has no additional
|
|
|
|
data.
|
|
|
|
"""
|
|
|
|
def __init__(self, fileobj, read_limit=READ_UNTIL_END):
|
|
|
|
super(LimitingStream, self).__init__(fileobj, 0, read_limit)
|