From 8dd861ce2784e2b5f9aea33799990724bb9c3a4a Mon Sep 17 00:00:00 2001 From: Vizonex Date: Tue, 29 Sep 2026 16:08:08 -0500 Subject: [PATCH 1/5] update futures and improve speed. --- winloop/future.pxd | 36 +++++ winloop/future.pyx | 271 ++++++++++++++++++++++++++++++++++++ winloop/includes/stdlib.pxi | 11 +- winloop/loop.pxd | 11 +- winloop/loop.pyx | 77 +++++----- 5 files changed, 364 insertions(+), 42 deletions(-) create mode 100644 winloop/future.pxd create mode 100644 winloop/future.pyx diff --git a/winloop/future.pxd b/winloop/future.pxd new file mode 100644 index 0000000..e985e05 --- /dev/null +++ b/winloop/future.pxd @@ -0,0 +1,36 @@ + +cdef enum fut_state: + _PENDING = 0 + _CANCELLED = 1 + _FINISHED = 2 + +cdef class Future: + cdef: + object __weakref__ + fut_state state + public object _result + public object _exception + public Loop _loop + public object _source_traceback + public object _cancel_message + public object _cancelled_exc + public object _exception_tb + + list __callbacks + public bint _asyncio_future_blocking + bint __log_traceback + bint _init + + cdef object ensure_alive(self) + cpdef Loop get_loop(self) + cpdef object _make_cancelled_error(self) + cpdef object cancel(self, object msg=*) + cpdef object __schedule_callbacks(self) + cpdef bint cancelled(self) noexcept + cpdef bint done(self) noexcept + cpdef object result(self) + cpdef object exception(self) + cpdef object set_result(self, object result) + cpdef Py_ssize_t remove_done_callback(self, object fn) except -1 + cpdef object set_result(self, object result) + cpdef object set_exception(self, object exception) diff --git a/winloop/future.pyx b/winloop/future.pyx new file mode 100644 index 0000000..55eca47 --- /dev/null +++ b/winloop/future.pyx @@ -0,0 +1,271 @@ +from cpython.contextvars cimport PyContext_CopyCurrent +from cpython.list cimport PyList_GET_SIZE, PyList_New, PyList_Clear + +from types import coroutine + + +# Cythonic version of the asyncio Future Object meant for researching into ways to make +# A asynchronous awaitable container object behave correctly while making less costly calls. +# This has positive effects on libraries like uvloop & winloop. +# uvloop maintainers have my permission to try this - Vizonex. + + +@coroutine +def __future_iter(Future fut): + while fut.state == _PENDING: + if not fut._asyncio_future_blocking: + fut._asyncio_future_blocking = True + yield fut + continue + raise RuntimeError("await wasn't used with future") + result = fut.result() + return result + +cdef class Future: + + def __init__(self, *, Loop loop): + self.state = _PENDING + self._result = None + self._exception = None + self._source_traceback = None + self._cancel_message = None + self._cancelled_exc = None + self.__log_traceback = False + + self.__callbacks = PyList_New(0) + + + self._loop = loop + if self._loop.get_debug(): + self._source_traceback = extract_stack() + self._init = True + + @property + def _state(self): + # compatibility with python futures + if self.state == _PENDING: + return "PENDING" + elif self.state == _CANCELLED: + return "CANCELLED" + return "FINISHED" + + + @property + def loop(self): + return self._loop + + def __repr__(self): + return aio__future_repr(self) + + def __del__(self): + cdef dict context + cdef object exc + if not self.__log_traceback: + # set_exception() was not called, or result() or exception() + # has consumed the exception + return + exc = self._exception + context = { + 'message': + f'{self.__class__.__name__} exception was never retrieved', + 'exception': exc, + 'future': self, + } + if self._source_traceback: + context['source_traceback'] = self._source_traceback + self._loop.call_exception_handler(context) + + __class_getitem__ = classmethod(types_GenericAlias) + + @property + def _log_traceback(self): + self.ensure_alive() + return self.__log_traceback + + @_log_traceback.setter + def _log_traceback(self, val): + if val: + raise ValueError('_log_traceback can only be set to False') + self.__log_traceback = False + + cdef object ensure_alive(self): + if self._loop is None: + raise RuntimeError("Future object is not initialized.") + + cpdef Loop get_loop(self): + """Return the event loop the Future is bound to.""" + self.ensure_alive() + return self._loop + + cpdef object _make_cancelled_error(self): + """Create the CancelledError to raise if the Future is cancelled. + + This should only be called once when handling a cancellation since + it erases the saved context exception value. + """ + if self._cancelled_exc is not None: + exc = self._cancelled_exc + self._cancelled_exc = None + return exc + + if self._cancel_message is None: + exc = aio_CancelledError() + else: + exc = aio_CancelledError(self._cancel_message) + exc.__context__ = self._cancelled_exc + # Remove the reference since we don't need this anymore. + self._cancelled_exc = None + return exc + + cpdef object cancel(self, object msg=None): + """Cancel the future and schedule callbacks. + + If the future is already done or cancelled, return False. Otherwise, + change the future's state to cancelled, schedule the callbacks and + return True. + """ + self.ensure_alive() + self.__log_traceback = False + if self.state != _PENDING: + return False + self.state = _CANCELLED + self._cancel_message = msg + self.__schedule_callbacks() + return True + + @property + def _callbacks(self): + self.ensure_alive() + return self.__callbacks + + + cpdef object __schedule_callbacks(self): + """Internal: Ask the event loop to call all callbacks. + + The callbacks are scheduled to be called as soon as possible. Also + clears the callback list. + """ + + if not PyList_GET_SIZE(self.__callbacks): + return + callbacks = self._callbacks[:] + PyList_Clear(self.__callbacks) + for callback, ctx in callbacks: + self._loop.call_soon(callback, self, context=ctx) + + cpdef bint cancelled(self) noexcept: + """Return True if the future was cancelled.""" + return self.state == _CANCELLED + + # Don't implement running(); see http://bugs.python.org/issue18699 + + cpdef bint done(self) noexcept: + """Return True if the future is done. + + Done means either that a result / exception are available, or that the + future was cancelled. + """ + return self.state != _PENDING + + cpdef object result(self): + """Return the result this future represents. + + If the future has been cancelled, raises CancelledError. If the + future's result isn't yet available, raises InvalidStateError. If + the future is done and has an exception set, this exception is raised. + """ + if self.state == _CANCELLED: + exc = self._make_cancelled_error() + raise exc + if self.state != _FINISHED: + raise aio_InvalidStateError('Result is not ready.') + self.__log_traceback = False + if self._exception is not None: + raise self._exception.with_traceback(self._exception_tb) + return self._result + + cpdef object exception(self): + """Return the exception that was set on this future. + + The exception (or None if no exception was set) is returned only if + the future is done. If the future has been cancelled, raises + CancelledError. If the future isn't done yet, raises + InvalidStateError. + """ + if self.state == _CANCELLED: + exc = self._make_cancelled_error() + raise exc + if self.state != _FINISHED: + raise aio_InvalidStateError('Exception is not set.') + self.__log_traceback = False + return self._exception + + # TODO: Remove context=None limitation + def add_done_callback(self, fn, *, context=None): + """Add a callback to be run when the future becomes done. + + The callback is called with a single argument - the future object. If + the future is already done when this is called, the callback is + scheduled with call_soon. + """ + if self.state != _PENDING: + self._loop._call_soon(fn, self, context=context) + else: + self.__callbacks.append((fn, PyContext_CopyCurrent() if context is None else context)) + + # New method not in PEP 3148. + + cpdef Py_ssize_t remove_done_callback(self, object fn) except -1: + """Remove all instances of a callback from the "call when done" list. + + Returns the number of callbacks removed. + """ + self.ensure_alive() + cdef Py_ssize_t removed_count + cdef list filtered_callbacks = [(f, ctx) + for (f, ctx) in self.__callbacks + if f != fn] + if removed_count := (PyList_GET_SIZE(self.__callbacks) - PyList_GET_SIZE(filtered_callbacks)): + self.__callbacks[:] = filtered_callbacks + return removed_count + + # So-called internal methods (note: no set_running_or_notify_cancel()). + + cpdef object set_result(self, object result): + """Mark the future done and set its result. + + If the future is already done when this method is called, raises + InvalidStateError. + """ + self.ensure_alive() + if self.state != _PENDING: + raise aio_InvalidStateError(f'{self._state}: {self!r}') + self._result = result + self.state = _FINISHED + self.__schedule_callbacks() + + cpdef object set_exception(self, object exception): + """Mark the future done and set an exception. + + If the future is already done when this method is called, raises + InvalidStateError. + """ + self.ensure_alive() + if self.state != _PENDING: + raise aio_InvalidStateError(f'{self._state}: {self!r}') + if isinstance(exception, type): + exception = exception() + if type(exception) is StopIteration: + raise TypeError("StopIteration interacts badly with generators " + "and cannot be raised into a Future") + self._exception = exception + self._exception_tb = (exception).__traceback__ + self.state = _FINISHED + self.__schedule_callbacks() + self.__log_traceback = True + + def __await__(self): + return __future_iter(self).__await__() + + def __iter__(self): + return __future_iter(self) diff --git a/winloop/includes/stdlib.pxi b/winloop/includes/stdlib.pxi index 074c06c..aad1cff 100644 --- a/winloop/includes/stdlib.pxi +++ b/winloop/includes/stdlib.pxi @@ -2,8 +2,9 @@ import asyncio, asyncio.log, asyncio.base_events, \ - asyncio.sslproto, asyncio.coroutines, \ - asyncio.futures, asyncio.transports + asyncio.base_futures, asyncio.sslproto, \ + asyncio.coroutines, asyncio.futures, \ + asyncio.transports import collections.abc import concurrent.futures import errno @@ -22,6 +23,7 @@ import shlex import threading import traceback import time +import types import warnings import weakref @@ -35,6 +37,7 @@ cdef aio_TimeoutError = asyncio.TimeoutError # for concept ideas SEE: https://gist.github.com/Vizonex/5196ae5fc7f2287df6a6dec8b37edc37 cdef aio_Future = asyncio.Future +cdef aio__future_repr = asyncio.base_futures._future_repr cdef aio_Task = asyncio.Task cdef aio_ensure_future = asyncio.ensure_future cdef aio_gather = asyncio.gather @@ -182,6 +185,8 @@ cdef tb_StackSummary = traceback.StackSummary cdef tb_walk_stack = traceback.walk_stack cdef tb_format_list = traceback.format_list +cdef types_GenericAlias = types.GenericAlias + cdef warnings_warn = warnings.warn # TODO: weakref classes here could be rewritten to Cython to prevent bottlenecks when passing @@ -198,4 +203,4 @@ del asyncio, concurrent, collections, errno del functools, inspect, itertools, socket, os, threading del signal, subprocess, ssl del time, traceback, warnings, weakref -del shlex +del shlex, types diff --git a/winloop/loop.pxd b/winloop/loop.pxd index 7e58862..879e621 100644 --- a/winloop/loop.pxd +++ b/winloop/loop.pxd @@ -181,13 +181,13 @@ cdef class Loop: cdef _has_writer(self, fd) cdef _remove_writer(self, fd) - cdef _sock_recv(self, fut, sock, n) - cdef _sock_recv_into(self, fut, sock, buf) - cdef _sock_sendall(self, fut, sock, data) - cdef _sock_accept(self, fut, sock) + cdef _sock_recv(self, _fut, sock, n) + cdef _sock_recv_into(self, _fut, sock, buf) + cdef _sock_sendall(self, _fut, sock, data) + cdef _sock_accept(self, _fut, sock) cdef _sock_connect(self, sock, address) - cdef _sock_connect_cb(self, fut, sock, address) + cdef _sock_connect_cb(self, _fut, sock, address) cdef _sock_set_reuseport(self, int fd) @@ -206,6 +206,7 @@ cdef class Loop: include "cbhandles.pxd" +include "future.pxd" include "handles/handle.pxd" include "handles/async_.pxd" diff --git a/winloop/loop.pyx b/winloop/loop.pyx index 89871b6..cb1a8a1 100644 --- a/winloop/loop.pyx +++ b/winloop/loop.pyx @@ -45,7 +45,7 @@ from . import _noop include "includes/stdlib.pxi" - +include "future.pyx" include "errors.pyx" cdef: @@ -77,7 +77,7 @@ cdef _is_sock_dgram(sock_type): cdef isfuture(obj): if aio_isfuture is None: - return isinstance(obj, aio_Future) + return isinstance(obj, (Future, aio_Future)) else: return aio_isfuture(obj) @@ -712,6 +712,8 @@ cdef class Loop: "Non-thread-safe operation invoked on an event loop other " "than the current one") + # XXX: still a problem with wrap_future preventing + # further progress. Fixing this is a TODO cdef inline _new_future(self): return aio_Future(loop=self) @@ -945,7 +947,8 @@ cdef class Loop: nr.query(addr, flags) return fut - cdef _sock_recv(self, fut, sock, n): + cdef _sock_recv(self, _fut, sock, n): + cdef _SyncSocketReaderFuture fut = <_SyncSocketReaderFuture>_fut if UVLOOP_DEBUG: if fut.cancelled(): # Shouldn't happen with _SyncSocketReaderFuture. @@ -972,7 +975,8 @@ cdef class Loop: fut.set_result(data) self._remove_reader(sock) - cdef _sock_recv_into(self, fut, sock, buf): + cdef _sock_recv_into(self, _fut, sock, buf): + cdef _SyncSocketReaderFuture fut = <_SyncSocketReaderFuture>_fut if UVLOOP_DEBUG: if fut.cancelled(): # Shouldn't happen with _SyncSocketReaderFuture. @@ -999,10 +1003,11 @@ cdef class Loop: fut.set_result(data) self._remove_reader(sock) - cdef _sock_sendall(self, fut, sock, data): + cdef _sock_sendall(self, _fut, sock, data): cdef: Handle handle int n + _SyncSocketWriterFuture fut = <_SyncSocketWriterFuture>_fut if UVLOOP_DEBUG: if fut.cancelled(): @@ -1047,7 +1052,8 @@ cdef class Loop: self._add_writer(sock, handle) - cdef _sock_accept(self, fut, sock): + cdef _sock_accept(self, _fut, sock): + cdef _SyncSocketReaderFuture fut = _fut try: conn, address = sock.accept() conn.setblocking(False) @@ -1064,9 +1070,12 @@ cdef class Loop: fut.set_result((conn, address)) self._remove_reader(sock) + # TODO: (Alter Signature to _SyncSocketWriterFuture) + # it will eliminate the rsloop vs uvloop bottleneck. cdef _sock_connect(self, sock, address): cdef: Handle handle + _SyncSocketWriterFuture fut try: sock.connect(address) @@ -1082,12 +1091,13 @@ cdef class Loop: self._sock_connect_cb, None, self, - fut, sock, address) + fut, sock, address) self._add_writer(sock, handle) return fut - cdef _sock_connect_cb(self, fut, sock, address): + cdef _sock_connect_cb(self, _fut, sock, address): + cdef Future fut = _fut if UVLOOP_DEBUG: if fut.cancelled(): # Shouldn't happen with _SyncSocketWriterFuture. @@ -3345,50 +3355,49 @@ cdef inline void __loop_free_buffer(Loop loop): # from being turned into a cdef extension class as it may possibly enhance # reading and writing functions in general. -class _SyncSocketReaderFuture(aio_Future): - - def __init__(self, sock, loop): - aio_Future.__init__(self, loop=loop) +cdef class _SyncSocketReaderFuture(Future): + cdef: + Loop __loop + object __sock + def __init__(self, sock, Loop loop): + super().__init__(loop=loop) self.__sock = sock self.__loop = loop - def __remove_reader(self): + cpdef object __remove_reader(self): if self.__sock is not None and self.__sock.fileno() != -1: self.__loop.remove_reader(self.__sock) self.__sock = None - if PY39: - def cancel(self, msg=None): - self.__remove_reader() - aio_Future.cancel(self, msg=msg) - else: - def cancel(self): - self.__remove_reader() - aio_Future.cancel(self) + # XXX: There is no point in fighting the PY39 flags. + cpdef object cancel(self, msg=None): + self.__remove_reader() + super().cancel(msg=msg) + -class _SyncSocketWriterFuture(aio_Future): - def __init__(self, sock, loop): - aio_Future.__init__(self, loop=loop) +cdef class _SyncSocketWriterFuture(Future): + cdef: + Loop __loop + object __sock + + def __init__(self, sock, Loop loop): + super().__init__(loop=loop) self.__sock = sock self.__loop = loop - def __remove_writer(self): + cpdef object __remove_writer(self): if self.__sock is not None and self.__sock.fileno() != -1: - self.__loop.remove_writer(self.__sock) + self.__loop._remove_writer(self.__sock) self.__sock = None - if PY39: - def cancel(self, msg=None): - self.__remove_writer() - aio_Future.cancel(self, msg=msg) + cpdef object cancel(self, msg=None): + self.__remove_writer() + super().cancel(msg=msg) + - else: - def cancel(self): - self.__remove_writer() - aio_Future.cancel(self) include "cbhandles.pyx" From 9a7b22838dbee843733cb1aac674f2d39d379022 Mon Sep 17 00:00:00 2001 From: Vizonex Date: Tue, 29 Sep 2026 16:14:58 -0500 Subject: [PATCH 2/5] make cancelling inlined. --- winloop/loop.pyx | 19 +++++++++++++++++-- 1 file changed, 17 insertions(+), 2 deletions(-) diff --git a/winloop/loop.pyx b/winloop/loop.pyx index cb1a8a1..2c5090e 100644 --- a/winloop/loop.pyx +++ b/winloop/loop.pyx @@ -3373,7 +3373,15 @@ cdef class _SyncSocketReaderFuture(Future): # XXX: There is no point in fighting the PY39 flags. cpdef object cancel(self, msg=None): self.__remove_reader() - super().cancel(msg=msg) + + self.ensure_alive() + self.__log_traceback = False + if self.state != _PENDING: + return False + self.state = _CANCELLED + self._cancel_message = msg + self.__schedule_callbacks() + return True @@ -3395,7 +3403,14 @@ cdef class _SyncSocketWriterFuture(Future): cpdef object cancel(self, msg=None): self.__remove_writer() - super().cancel(msg=msg) + self.ensure_alive() + self.__log_traceback = False + if self.state != _PENDING: + return False + self.state = _CANCELLED + self._cancel_message = msg + self.__schedule_callbacks() + return True From 9647891caf46828ce9dc2eca276251cc232fae7f Mon Sep 17 00:00:00 2001 From: Vizonex Date: Tue, 29 Sep 2026 16:18:26 -0500 Subject: [PATCH 3/5] fallback if _future_repr fails --- winloop/includes/stdlib.pxi | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/winloop/includes/stdlib.pxi b/winloop/includes/stdlib.pxi index aad1cff..3635067 100644 --- a/winloop/includes/stdlib.pxi +++ b/winloop/includes/stdlib.pxi @@ -37,7 +37,7 @@ cdef aio_TimeoutError = asyncio.TimeoutError # for concept ideas SEE: https://gist.github.com/Vizonex/5196ae5fc7f2287df6a6dec8b37edc37 cdef aio_Future = asyncio.Future -cdef aio__future_repr = asyncio.base_futures._future_repr +cdef aio__future_repr = getattr(asyncio.base_futures, "_future_repr", lambda fut: None) cdef aio_Task = asyncio.Task cdef aio_ensure_future = asyncio.ensure_future cdef aio_gather = asyncio.gather From d85c76d0b022eb6f2a4a15c82fe5c2904048df73 Mon Sep 17 00:00:00 2001 From: Vizonex Date: Tue, 29 Sep 2026 16:25:10 -0500 Subject: [PATCH 4/5] remove log_traceback in socket futures --- winloop/loop.pyx | 3 --- 1 file changed, 3 deletions(-) diff --git a/winloop/loop.pyx b/winloop/loop.pyx index 2c5090e..b11b583 100644 --- a/winloop/loop.pyx +++ b/winloop/loop.pyx @@ -3373,9 +3373,7 @@ cdef class _SyncSocketReaderFuture(Future): # XXX: There is no point in fighting the PY39 flags. cpdef object cancel(self, msg=None): self.__remove_reader() - self.ensure_alive() - self.__log_traceback = False if self.state != _PENDING: return False self.state = _CANCELLED @@ -3404,7 +3402,6 @@ cdef class _SyncSocketWriterFuture(Future): cpdef object cancel(self, msg=None): self.__remove_writer() self.ensure_alive() - self.__log_traceback = False if self.state != _PENDING: return False self.state = _CANCELLED From 6d6dbd2ba6dcf991aedaf0b67ecec7505022888e Mon Sep 17 00:00:00 2001 From: Vizonex Date: Wed, 30 Sep 2026 10:34:18 -0500 Subject: [PATCH 5/5] fix and harden the New Cython Future objects --- tests/test_base.py | 4 +++- winloop/future.pxd | 2 +- winloop/future.pyx | 10 +++++----- winloop/loop.pxd | 17 +++++++++++++++++ winloop/loop.pyx | 47 +++++++++++++++++++++++----------------------- 5 files changed, 49 insertions(+), 31 deletions(-) diff --git a/tests/test_base.py b/tests/test_base.py index b0ca5d3..79726d8 100644 --- a/tests/test_base.py +++ b/tests/test_base.py @@ -758,7 +758,9 @@ def scheduler(): class TestBaseUV(_TestBase, UVTestCase): def test_loop_create_future(self): fut = self.loop.create_future() - self.assertTrue(isinstance(fut, asyncio.Future)) + # winloop.loop.Future (Coming soon) + # or asyncio.Future + self.assertTrue(asyncio.isfuture(fut)) self.assertIs(fut._loop, self.loop) fut.cancel() diff --git a/winloop/future.pxd b/winloop/future.pxd index e985e05..6d01054 100644 --- a/winloop/future.pxd +++ b/winloop/future.pxd @@ -25,7 +25,7 @@ cdef class Future: cpdef Loop get_loop(self) cpdef object _make_cancelled_error(self) cpdef object cancel(self, object msg=*) - cpdef object __schedule_callbacks(self) + cpdef object _schedule_callbacks(self) cpdef bint cancelled(self) noexcept cpdef bint done(self) noexcept cpdef object result(self) diff --git a/winloop/future.pyx b/winloop/future.pyx index 55eca47..f4d35cf 100644 --- a/winloop/future.pyx +++ b/winloop/future.pyx @@ -130,7 +130,7 @@ cdef class Future: return False self.state = _CANCELLED self._cancel_message = msg - self.__schedule_callbacks() + self._schedule_callbacks() return True @property @@ -139,7 +139,7 @@ cdef class Future: return self.__callbacks - cpdef object __schedule_callbacks(self): + cpdef object _schedule_callbacks(self): """Internal: Ask the event loop to call all callbacks. The callbacks are scheduled to be called as soon as possible. Also @@ -148,7 +148,7 @@ cdef class Future: if not PyList_GET_SIZE(self.__callbacks): return - callbacks = self._callbacks[:] + callbacks = self.__callbacks[:] PyList_Clear(self.__callbacks) for callback, ctx in callbacks: self._loop.call_soon(callback, self, context=ctx) @@ -242,7 +242,7 @@ cdef class Future: raise aio_InvalidStateError(f'{self._state}: {self!r}') self._result = result self.state = _FINISHED - self.__schedule_callbacks() + self._schedule_callbacks() cpdef object set_exception(self, object exception): """Mark the future done and set an exception. @@ -261,7 +261,7 @@ cdef class Future: self._exception = exception self._exception_tb = (exception).__traceback__ self.state = _FINISHED - self.__schedule_callbacks() + self._schedule_callbacks() self.__log_traceback = True def __await__(self): diff --git a/winloop/loop.pxd b/winloop/loop.pxd index 879e621..e70d923 100644 --- a/winloop/loop.pxd +++ b/winloop/loop.pxd @@ -228,3 +228,20 @@ include "sslproto.pxd" include "handles/udp.pxd" include "server.pxd" + + +cdef class _SyncSocketReaderFuture(Future): + cdef: + Loop __loop + object __sock + + cpdef object __remove_reader(self) + cpdef object cancel(self, object msg=*) + +cdef class _SyncSocketWriterFuture(Future): + cdef: + Loop __loop + object __sock + + cpdef object __remove_writer(self) + cpdef object cancel(self, object msg=*) diff --git a/winloop/loop.pyx b/winloop/loop.pyx index b11b583..8c28ab8 100644 --- a/winloop/loop.pyx +++ b/winloop/loop.pyx @@ -3356,9 +3356,7 @@ cdef inline void __loop_free_buffer(Loop loop): # reading and writing functions in general. cdef class _SyncSocketReaderFuture(Future): - cdef: - Loop __loop - object __sock + def __init__(self, sock, Loop loop): super().__init__(loop=loop) self.__sock = sock @@ -3366,28 +3364,27 @@ cdef class _SyncSocketReaderFuture(Future): cpdef object __remove_reader(self): if self.__sock is not None and self.__sock.fileno() != -1: - self.__loop.remove_reader(self.__sock) + self.__loop._remove_reader(self.__sock) self.__sock = None + # XXX: There is no point in fighting the PY39 flags + # when this can simply be skipped over. PY38 is old + # and winloop stopped maintaining it in order to + # encourage those who can update to update. + cpdef object cancel(self, object msg=None): + """Cancel the future and schedule callbacks. - # XXX: There is no point in fighting the PY39 flags. - cpdef object cancel(self, msg=None): + If the future is already done or cancelled, return False. Otherwise, + change the future's state to cancelled, schedule the callbacks and + return True. + """ self.__remove_reader() - self.ensure_alive() - if self.state != _PENDING: - return False - self.state = _CANCELLED - self._cancel_message = msg - self.__schedule_callbacks() - return True + return Future.cancel(self, msg) cdef class _SyncSocketWriterFuture(Future): - cdef: - Loop __loop - object __sock def __init__(self, sock, Loop loop): super().__init__(loop=loop) @@ -3399,15 +3396,17 @@ cdef class _SyncSocketWriterFuture(Future): self.__loop._remove_writer(self.__sock) self.__sock = None - cpdef object cancel(self, msg=None): + cpdef object cancel(self, object msg=None): + """Cancel the future and schedule callbacks. + + If the future is already done or cancelled, return False. Otherwise, + change the future's state to cancelled, schedule the callbacks and + return True. + """ self.__remove_writer() - self.ensure_alive() - if self.state != _PENDING: - return False - self.state = _CANCELLED - self._cancel_message = msg - self.__schedule_callbacks() - return True + # recast subclass to regain access to anything + # private + return Future.cancel(self, msg)