@@ -105,8 +105,9 @@ def close(self):
105105 if self ._closing :
106106 return
107107 self ._closing = True
108- self ._conn_lost += 1
109108 if not self ._buffer and self ._write_fut is None :
109+ # Nothing left to flush: no more data will be sent.
110+ self ._conn_lost += 1
110111 self ._loop .call_soon (self ._call_connection_lost , None )
111112 if self ._read_fut is not None :
112113 self ._read_fut .cancel ()
@@ -386,6 +387,7 @@ def _loop_writing(self, f=None, data=None):
386387 self ._buffer = None
387388 if not data :
388389 if self ._closing :
390+ self ._conn_lost += 1
389391 self ._loop .call_soon (self ._call_connection_lost , None )
390392 if self ._eof_written :
391393 self ._sock .shutdown (socket .SHUT_WR )
@@ -480,6 +482,11 @@ def get_write_buffer_size(self):
480482 def abort (self ):
481483 self ._force_close (None )
482484
485+ def _force_close (self , exc ):
486+ # The base class drops the buffer; the size is tracked separately.
487+ self ._buffer_size = 0
488+ super ()._force_close (exc )
489+
483490 def sendto (self , data , addr = None ):
484491 if not isinstance (data , (bytes , bytearray , memoryview )):
485492 raise TypeError ('data argument must be bytes-like object (%r)' ,
@@ -509,6 +516,8 @@ def sendto(self, data, addr=None):
509516 def _loop_writing (self , fut = None ):
510517 try :
511518 if self ._conn_lost :
519+ # No more data will be sent: either everything buffered has
520+ # already been flushed, or _force_close() dropped it.
512521 return
513522
514523 assert fut is self ._write_fut
@@ -517,9 +526,10 @@ def _loop_writing(self, fut=None):
517526 # We are in a _loop_writing() done callback, get the result
518527 fut .result ()
519528
520- if not self ._buffer or ( self . _conn_lost and self . _address ) :
521- # The connection has been closed
529+ if not self ._buffer :
530+ # Everything buffered has been sent
522531 if self ._closing :
532+ self ._conn_lost += 1
523533 self ._loop .call_soon (self ._call_connection_lost , None )
524534 return
525535
@@ -534,6 +544,27 @@ def _loop_writing(self, fut=None):
534544 addr = addr )
535545 except OSError as exc :
536546 self ._protocol .error_received (exc )
547+ # error_received() is arbitrary protocol code: it may have sent
548+ # (scheduling a write of its own, directly or via call_soon()),
549+ # closed, or aborted the transport.
550+ if self ._buffer or self ._closing :
551+ # Either data is still queued, or a close() is waiting on
552+ # the write loop to drain it and call connection_lost().
553+ # This write failed, so there is no completion callback
554+ # pending to re-enter the loop -- schedule one (gh-156698).
555+ def write_next ():
556+ # error_received() may have scheduled a write of its own,
557+ # directly or with call_soon(); its completion callback
558+ # will drain the rest of the buffer.
559+ if self ._write_fut is None :
560+ self ._loop_writing ()
561+
562+ self ._loop .call_soon (write_next )
563+ else :
564+ # Nothing left to write, so a paused protocol has to be
565+ # resumed here: the next entry into _loop_writing() returns
566+ # early on an empty buffer without doing it.
567+ self ._maybe_resume_protocol ()
537568 except Exception as exc :
538569 self ._fatal_error (exc , 'Fatal write error on datagram transport' )
539570 else :
@@ -543,28 +574,20 @@ def _loop_writing(self, fut=None):
543574 def _loop_reading (self , fut = None ):
544575 data = None
545576 try :
546- if self ._conn_lost :
577+ if self ._closing :
547578 return
548579
549- assert self ._read_fut is fut or (self ._read_fut is None and
550- self ._closing )
580+ assert self ._read_fut is fut
551581
552582 self ._read_fut = None
553583 if fut is not None :
554584 res = fut .result ()
555585
556- if self ._closing :
557- # since close() has been called we ignore any read data
558- data = None
559- return
560-
561586 if self ._address is not None :
562587 data , addr = res , self ._address
563588 else :
564589 data , addr = res
565590
566- if self ._conn_lost :
567- return
568591 if self ._address is not None :
569592 self ._read_fut = self ._loop ._proactor .recv (self ._sock ,
570593 self .max_size )
0 commit comments