|
| 1 | +import io |
| 2 | + |
| 3 | +from collections import deque |
| 4 | +from micropython import ringbuffer, const |
| 5 | + |
| 6 | +try: |
| 7 | + from typing import Union, Tuple |
| 8 | +except: |
| 9 | + pass |
| 10 | + |
| 11 | +# From micropython/py/stream.h |
| 12 | +_MP_STREAM_ERROR = const(-1) |
| 13 | +_MP_STREAM_FLUSH = const(1) |
| 14 | +_MP_STREAM_SEEK = const(2) |
| 15 | +_MP_STREAM_POLL = const(3) |
| 16 | +_MP_STREAM_CLOSE = const(4) |
| 17 | +_MP_STREAM_POLL_RD = const(0x0001) |
| 18 | + |
| 19 | + |
| 20 | +def streampair(buffer_size: Union[int, Tuple[int, int]]=256): |
| 21 | + """ |
| 22 | + Returns two bi-directional linked stream objects where writes to one can be read from the other and vice/versa. |
| 23 | + This can be used somewhat similarly to a socket.socketpair in python, like a pipe |
| 24 | + of data that can be used to connect stream consumers (eg. asyncio.StreamWriter, mock Uart) |
| 25 | + """ |
| 26 | + try: |
| 27 | + size_a, size_b = buffer_size |
| 28 | + except TypeError: |
| 29 | + size_a = size_b = buffer_size |
| 30 | + |
| 31 | + a = ringbuffer(size_a) |
| 32 | + b = ringbuffer(size_b) |
| 33 | + return StreamPair(a, b), StreamPair(b, a) |
| 34 | + |
| 35 | + |
| 36 | +class StreamPair(io.IOBase): |
| 37 | + |
| 38 | + def __init__(self, own: ringbuffer, other: ringbuffer): |
| 39 | + self.own = own |
| 40 | + self.other = other |
| 41 | + super().__init__() |
| 42 | + |
| 43 | + def read(self, nbytes=-1): |
| 44 | + return self.own.read(nbytes) |
| 45 | + |
| 46 | + def readline(self): |
| 47 | + return self.own.readline() |
| 48 | + |
| 49 | + def readinto(self, buf, limit=-1): |
| 50 | + return self.own.readinto(buf, limit) |
| 51 | + |
| 52 | + def write(self, data): |
| 53 | + return self.other.write(data) |
| 54 | + |
| 55 | + def seek(self, offset, whence): |
| 56 | + return self.own.seek(offset, whence) |
| 57 | + |
| 58 | + def flush(self): |
| 59 | + self.own.flush() |
| 60 | + self.other.flush() |
| 61 | + |
| 62 | + def close(self): |
| 63 | + self.own.close() |
| 64 | + self.other.close() |
| 65 | + |
| 66 | + def any(self): |
| 67 | + return self.own.any() |
| 68 | + |
| 69 | + def ioctl(self, op, arg): |
| 70 | + if op == _MP_STREAM_POLL: |
| 71 | + if self.any(): |
| 72 | + return _MP_STREAM_POLL_RD |
| 73 | + return 0 |
| 74 | + |
| 75 | + elif op ==_MP_STREAM_FLUSH: |
| 76 | + return self.flush() |
| 77 | + elif op ==_MP_STREAM_SEEK: |
| 78 | + return self.seek(arg[0], arg[1]) |
| 79 | + elif op ==_MP_STREAM_CLOSE: |
| 80 | + return self.close() |
| 81 | + |
| 82 | + else: |
| 83 | + return _MP_STREAM_ERROR |
0 commit comments