@@ -109,6 +109,7 @@ class Reader(Stoppable):
109109 on_msg_coro : OnMsgCoro = attrs .field (validator = Validators .not_none ())
110110 on_close_coro : OnCloseCoro = attrs .field (validator = Validators .not_none ())
111111 _buffer : bytearray = attrs .field (init = False , factory = bytearray )
112+ _read_pos : int = attrs .field (init = False , default = 0 )
112113 _drain_buffer : bytearray | None = attrs .field (init = False , default = None )
113114 _task : asyncio .Task = attrs .field (init = False , default = None )
114115 _stopped : bool = attrs .field (init = False , default = False )
@@ -130,6 +131,10 @@ async def buffer_until_drained(self, discard_buffer: bool = False):
130131 finally :
131132 self ._drain_mode = None
132133 if not discard_buffer :
134+ # Compact before extending with drain buffer
135+ if self ._read_pos > 0 :
136+ del self ._buffer [:self ._read_pos ]
137+ self ._read_pos = 0
133138 self ._buffer .extend (self ._drain_buffer )
134139 self ._drain_buffer = None
135140
@@ -153,11 +158,11 @@ def is_stopped(self):
153158
154159 async def _process (self ):
155160 while not self ._stopped :
156- len_before = len_after = len (self ._buffer )
157- if len_before > 0 :
161+ available_before = available_after = len (self ._buffer ) - self . _read_pos
162+ if available_before > 0 :
158163 await self ._process_1 ()
159- len_after = len (self ._buffer )
160- if self ._drain_mode and not self ._drain_mode .is_set () and len_after == len_before :
164+ available_after = len (self ._buffer ) - self . _read_pos
165+ if self ._drain_mode and not self ._drain_mode .is_set () and available_after == available_before :
161166 self ._drain_mode .set ()
162167
163168 await asyncio .sleep (0.0001 )
0 commit comments