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