# (c) Copyright 2018 by Coinkite Inc. This file is part of Coldcard # and is covered by GPLv3 license found in COPYING. # # See also: # import uerrno import uselect as select from uasyncio.core import * class PollEventLoop(EventLoop): def __init__(self, len=42): EventLoop.__init__(self, len) self.poller = select.poll() self.objmap = {} def add_reader(self, sock, cb, *args): if args: self.poller.register(sock, select.POLLIN) self.objmap[id(sock)] = (cb, args) else: self.poller.register(sock, select.POLLIN) self.objmap[id(sock)] = cb def remove_reader(self, sock): self.poller.unregister(sock) del self.objmap[id(sock)] def add_writer(self, sock, cb, *args): if args: self.poller.register(sock, select.POLLOUT) self.objmap[id(sock)] = (cb, args) else: self.poller.register(sock, select.POLLOUT) self.objmap[id(sock)] = cb def remove_writer(self, sock): try: self.poller.unregister(sock) self.objmap.pop(id(sock), None) except OSError as e: # StreamWriter.awrite() first tries to write to a socket, # and if that succeeds, yield IOWrite may never be called # for that socket, and it will never be added to poller. So, # ignore such error. if e.args[0] != uerrno.ENOENT: raise def wait(self, delay): # We need one-shot behavior (second arg of 1 to .poll()) res = self.poller.ipoll(delay, 1) # Remove "if res" workaround after # https://github.com/micropython/micropython/issues/2716 fixed. if res: for sock, ev in res: cb = self.objmap[id(sock)] if ev & (select.POLLHUP | select.POLLERR): # These events are returned even if not requested, and # are sticky, i.e. will be returned again and again. # If the caller doesn't do proper error handling and # unregister this sock, we'll busy-loop on it, so we # as well can unregister it now "just in case". self.remove_reader(sock) if isinstance(cb, tuple): cb[0](*cb[1]) else: cb.pend_throw(None) self.call_soon(cb) class StreamReader: def __init__(self, polls, ios=None): if ios is None: ios = polls self.polls = polls self.ios = ios def read(self, n=-1): while True: yield IORead(self.polls) res = self.ios.read(n) if res is not None: break # This should not happen for real sockets, but can easily # happen for stream wrappers (ssl, websockets, etc.) if not res: yield IOReadDone(self.polls) return res def readexactly(self, n): buf = b"" while n: yield IORead(self.polls) res = self.ios.read(n) assert res is not None if not res: yield IOReadDone(self.polls) break buf += res n -= len(res) return buf def readline(self): buf = b"" while True: yield IORead(self.polls) res = self.ios.readline() assert res is not None if not res: yield IOReadDone(self.polls) break buf += res if buf[-1] == 0x0a: break return buf def aclose(self): yield IOReadDone(self.polls) self.ios.close() def __repr__(self): return "" % (self.polls, self.ios) class StreamWriter: def __init__(self, s, extra): self.s = s self.extra = extra def awrite(self, buf, off=0, sz=-1): # This method is called awrite (async write) to not proliferate # incompatibility with original asyncio. Unlike original asyncio # whose .write() method is both not a coroutine and guaranteed # to return immediately (which means it has to buffer all the # data), this method is a coroutine. if sz == -1: sz = len(buf) - off while True: res = self.s.write(buf, off, sz) # If we spooled everything, return immediately if res == sz: return if res is None: res = 0 assert res < sz off += res sz -= res yield IOWrite(self.s) #assert s2.fileno() == self.s.fileno() # Write piecewise content from iterable (usually, a generator) def awriteiter(self, iterable): for buf in iterable: yield from self.awrite(buf) def aclose(self): yield IOWriteDone(self.s) self.s.close() def get_extra_info(self, name, default=None): return self.extra.get(name, default) def __repr__(self): return "" % self.s import uasyncio.core uasyncio.core._event_loop_class = PollEventLoop