Source code for asyncutils.iterclasses

 1import asyncutils as A
 2from asyncutils.constants import _NO_DEFAULT
 3from asyncutils._internal import helpers as H, patch as P
 4from asyncutils._internal.submodules import iterclasses_all as __all__
 5from _collections import defaultdict, deque
 6from sys import maxsize as I
[docs] 7@H.subscriptable 8class AChain: 9 __slots__ = '__its', 10 @classmethod 11 async def _flatten_it_of_its(cls, I, /): 12 async for i in A.iter_to_agen(I): 13 if isinstance(i, cls): 14 async for _ in cls._flatten_it_of_its(i.__its): yield _ 15 else: yield i
[docs] 16 @classmethod 17 def from_iterable(cls, it_of_its): (s := super().__new__(cls)).__its = cls._flatten_it_of_its(it_of_its); return s
18 def __new__(cls, *its): return cls.from_iterable(its)
[docs] 19 async def __aiter__(self): 20 async for i in self.__its: # ty: ignore[unresolved-attribute] 21 async for _ in i: yield _
[docs] 22@H.subscriptable 23class APeekable(H.LoopMixinBase): 24 __slots__ = '__ca', '__it' 25 def __init__(self, it=()): self.__it, self.__ca = A.iter_to_agen(it), deque(); super().__init__()
[docs] 26 def __aiter__(self): return self
[docs] 27 async def can_peek(self): 28 try: await self.peek(); return True 29 except StopAsyncIteration: return False
[docs] 30 async def peek(self, default=_NO_DEFAULT): 31 if not (c := self.__ca): 32 try: c.append(await anext(self.__it)) 33 except StopAsyncIteration: 34 if default is _NO_DEFAULT: raise 35 return default 36 return c[0]
[docs] 37 def prepend(self, /, *i): self.__ca.extendleft(reversed(i))
[docs] 38 async def __anext__(self): 39 if (c := self.__ca): return c.popleft() 40 return await anext(self.__it)
[docs] 41 async def __getitem__(self, i, /, _=~I): 42 f = (C := self.__ca).append 43 if isinstance(i, slice): 44 if (c := 1 if (s := i.step) is None else int(s)) > 0: a, b = 0 if (s := i.start) is None else int(s), I if (s := i.stop) is None else int(s) 45 elif c < 0: a, b = -1 if (s := i.start) is None else int(s), _ if (s := i.stop) is None else int(s) 46 else: raise ValueError('asyncutils.iterclasses.APeekable: slice step cannot be zero') 47 if a < 0 or b < 0: 48 async for s in A.iter_to_agen(self.__it): f(s) 49 elif (d := min(max(a, b)+1, I)-len(C)) >= 0: 50 async for s in A.take(self.__it, d): f(s) 51 return tuple(C)[a:b:c] 52 async for s in A.iter_to_agen(self.__it) if (i := i.__index__()) < 0 else A.empty_agen() if i < (l := len(C)) else A.take(self.__it, i-l+1): f(s) 53 return C[i]
54 P.patch_method_signatures((__getitem__, 'idx, /'))
[docs] 55@H.subscriptable 56class ABucket: 57 __slots__ = '__ca', '__it', '__key', '__vd' 58 def __init__(self, it, key, validator=None): super().__init__(); self.__it, self.__key, self.__ca, self.__vd = A.iter_to_agen(it), key, defaultdict(deque), validator or (lambda _: True)
[docs] 59 async def contains(self, k, /): 60 if not self.__vd(k): return False 61 try: i = await anext(self[k]) 62 except StopAsyncIteration: return False 63 self.__ca[k].append(i); return True
[docs] 64 async def __aiter__(self): 65 K, V, C = self.__key, self.__vd, self.__ca 66 async for i in self.__it: 67 if V(k := K(i)): C[k].append(i) 68 for k in C: yield k
[docs] 69 async def __getitem__(self, k, /): 70 if not (V := self.__vd)(k): return 71 p, I, K = (a := (C := self.__ca)[k]).popleft, self.__it, self.__key 72 while True: 73 if a: yield p(); continue 74 while True: 75 try: i = await anext(I) 76 except StopAsyncIteration: 77 if not a: del C[k] 78 return 79 if (c := K(i)) == k: yield i; break 80 elif V(c): C[c].append(i)
81del P