Skip to content

Commit 49ee6d4

Browse files
dkropachevclaude
andcommitted
Fix _close() on Windows ProactorEventLoop
ProactorEventLoop does not support remove_reader/remove_writer (raises NotImplementedError). Wrap these calls so the socket is always closed regardless, and use try/finally to ensure connected_event is always set even if cleanup fails. Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>
1 parent c7d58ec commit 49ee6d4

1 file changed

Lines changed: 28 additions & 18 deletions

File tree

cassandra/io/asyncioreactor.py

Lines changed: 28 additions & 18 deletions
Original file line numberDiff line numberDiff line change
@@ -158,24 +158,34 @@ def close(self):
158158

159159
async def _close(self):
160160
log.debug("Closing connection (%s) to %s" % (id(self), self.endpoint))
161-
if self._write_watcher:
162-
self._write_watcher.cancel()
163-
if self._read_watcher:
164-
self._read_watcher.cancel()
165-
if self._socket:
166-
self._loop.remove_writer(self._socket.fileno())
167-
self._loop.remove_reader(self._socket.fileno())
168-
self._socket.close()
169-
170-
log.debug("Closed socket to %s" % (self.endpoint,))
171-
172-
if not self.is_defunct:
173-
msg = "Connection to %s was closed" % self.endpoint
174-
if self.last_error:
175-
msg += ": %s" % (self.last_error,)
176-
self.error_all_requests(ConnectionShutdown(msg))
177-
# don't leave in-progress operations hanging
178-
self.connected_event.set()
161+
try:
162+
if self._write_watcher:
163+
self._write_watcher.cancel()
164+
if self._read_watcher:
165+
self._read_watcher.cancel()
166+
if self._socket:
167+
# remove_reader/remove_writer are not supported on Windows
168+
# ProactorEventLoop — ignore failures so the socket still
169+
# gets closed.
170+
try:
171+
self._loop.remove_writer(self._socket.fileno())
172+
except (NotImplementedError, OSError):
173+
pass
174+
try:
175+
self._loop.remove_reader(self._socket.fileno())
176+
except (NotImplementedError, OSError):
177+
pass
178+
self._socket.close()
179+
180+
log.debug("Closed socket to %s" % (self.endpoint,))
181+
finally:
182+
if not self.is_defunct:
183+
msg = "Connection to %s was closed" % self.endpoint
184+
if self.last_error:
185+
msg += ": %s" % (self.last_error,)
186+
self.error_all_requests(ConnectionShutdown(msg))
187+
# don't leave in-progress operations hanging
188+
self.connected_event.set()
179189

180190
def push(self, data):
181191
if self.is_closed or self.is_defunct:

0 commit comments

Comments
 (0)