| Yury Selivanov | b0b0e62 | 2014-02-18 22:27:48 -0500 | [diff] [blame] | 1 | """Selector and proactor event loops for Windows.""" | 
| Guido van Rossum | 27b7c7e | 2013-10-17 13:40:50 -0700 | [diff] [blame] | 2 |  | 
| Victor Stinner | f2e1768 | 2014-01-31 16:25:24 +0100 | [diff] [blame] | 3 | import _winapi | 
| Guido van Rossum | 27b7c7e | 2013-10-17 13:40:50 -0700 | [diff] [blame] | 4 | import errno | 
| Victor Stinner | f2e1768 | 2014-01-31 16:25:24 +0100 | [diff] [blame] | 5 | import math | 
| Guido van Rossum | 27b7c7e | 2013-10-17 13:40:50 -0700 | [diff] [blame] | 6 | import socket | 
| Victor Stinner | f2e1768 | 2014-01-31 16:25:24 +0100 | [diff] [blame] | 7 | import struct | 
| Guido van Rossum | 27b7c7e | 2013-10-17 13:40:50 -0700 | [diff] [blame] | 8 | import weakref | 
| Guido van Rossum | 27b7c7e | 2013-10-17 13:40:50 -0700 | [diff] [blame] | 9 |  | 
| Guido van Rossum | 0eaa5ac | 2013-11-04 15:50:46 -0800 | [diff] [blame] | 10 | from . import events | 
| Guido van Rossum | 5969128 | 2013-10-30 14:52:03 -0700 | [diff] [blame] | 11 | from . import base_subprocess | 
| Guido van Rossum | 27b7c7e | 2013-10-17 13:40:50 -0700 | [diff] [blame] | 12 | from . import futures | 
 | 13 | from . import proactor_events | 
 | 14 | from . import selector_events | 
 | 15 | from . import tasks | 
 | 16 | from . import windows_utils | 
| Guido van Rossum | 5969128 | 2013-10-30 14:52:03 -0700 | [diff] [blame] | 17 | from . import _overlapped | 
| Victor Stinner | f951d28 | 2014-06-29 00:46:45 +0200 | [diff] [blame] | 18 | from .coroutines import coroutine | 
 | 19 | from .log import logger | 
| Guido van Rossum | 27b7c7e | 2013-10-17 13:40:50 -0700 | [diff] [blame] | 20 |  | 
 | 21 |  | 
| Guido van Rossum | 0eaa5ac | 2013-11-04 15:50:46 -0800 | [diff] [blame] | 22 | __all__ = ['SelectorEventLoop', 'ProactorEventLoop', 'IocpProactor', | 
 | 23 |            'DefaultEventLoopPolicy', | 
 | 24 |            ] | 
| Guido van Rossum | 27b7c7e | 2013-10-17 13:40:50 -0700 | [diff] [blame] | 25 |  | 
 | 26 |  | 
 | 27 | NULL = 0 | 
 | 28 | INFINITE = 0xffffffff | 
 | 29 | ERROR_CONNECTION_REFUSED = 1225 | 
 | 30 | ERROR_CONNECTION_ABORTED = 1236 | 
 | 31 |  | 
 | 32 |  | 
 | 33 | class _OverlappedFuture(futures.Future): | 
 | 34 |     """Subclass of Future which represents an overlapped operation. | 
 | 35 |  | 
 | 36 |     Cancelling it will immediately cancel the overlapped operation. | 
 | 37 |     """ | 
 | 38 |  | 
 | 39 |     def __init__(self, ov, *, loop=None): | 
 | 40 |         super().__init__(loop=loop) | 
| Victor Stinner | fea6a10 | 2014-07-25 00:54:53 +0200 | [diff] [blame] | 41 |         if self._source_traceback: | 
 | 42 |             del self._source_traceback[-1] | 
| Victor Stinner | 18a28dc | 2014-07-25 13:05:20 +0200 | [diff] [blame] | 43 |         self._ov = ov | 
| Guido van Rossum | 27b7c7e | 2013-10-17 13:40:50 -0700 | [diff] [blame] | 44 |  | 
| Victor Stinner | e912e65 | 2014-07-12 03:11:53 +0200 | [diff] [blame] | 45 |     def __repr__(self): | 
 | 46 |         info = [self._state.lower()] | 
| Victor Stinner | 18a28dc | 2014-07-25 13:05:20 +0200 | [diff] [blame] | 47 |         if self._ov is not None: | 
 | 48 |             state = 'pending' if self._ov.pending else 'completed' | 
 | 49 |             info.append('overlapped=<%s, %#x>' % (state, self._ov.address)) | 
| Victor Stinner | e912e65 | 2014-07-12 03:11:53 +0200 | [diff] [blame] | 50 |         if self._state == futures._FINISHED: | 
 | 51 |             info.append(self._format_result()) | 
 | 52 |         if self._callbacks: | 
 | 53 |             info.append(self._format_callbacks()) | 
 | 54 |         return '<%s %s>' % (self.__class__.__name__, ' '.join(info)) | 
 | 55 |  | 
| Victor Stinner | 18a28dc | 2014-07-25 13:05:20 +0200 | [diff] [blame] | 56 |     def _cancel_overlapped(self): | 
 | 57 |         if self._ov is None: | 
 | 58 |             return | 
 | 59 |         try: | 
 | 60 |             self._ov.cancel() | 
 | 61 |         except OSError as exc: | 
 | 62 |             context = { | 
 | 63 |                 'message': 'Cancelling an overlapped future failed', | 
 | 64 |                 'exception': exc, | 
 | 65 |                 'future': self, | 
 | 66 |             } | 
 | 67 |             if self._source_traceback: | 
 | 68 |                 context['source_traceback'] = self._source_traceback | 
 | 69 |             self._loop.call_exception_handler(context) | 
 | 70 |         self._ov = None | 
 | 71 |  | 
| Guido van Rossum | 27b7c7e | 2013-10-17 13:40:50 -0700 | [diff] [blame] | 72 |     def cancel(self): | 
| Victor Stinner | 18a28dc | 2014-07-25 13:05:20 +0200 | [diff] [blame] | 73 |         self._cancel_overlapped() | 
| Guido van Rossum | 27b7c7e | 2013-10-17 13:40:50 -0700 | [diff] [blame] | 74 |         return super().cancel() | 
 | 75 |  | 
| Victor Stinner | 18a28dc | 2014-07-25 13:05:20 +0200 | [diff] [blame] | 76 |     def set_exception(self, exception): | 
 | 77 |         super().set_exception(exception) | 
 | 78 |         self._cancel_overlapped() | 
 | 79 |  | 
| Victor Stinner | 51e44ea | 2014-07-26 00:58:34 +0200 | [diff] [blame^] | 80 |     def set_result(self, result): | 
 | 81 |         super().set_result(result) | 
 | 82 |         self._ov = None | 
 | 83 |  | 
| Guido van Rossum | 27b7c7e | 2013-10-17 13:40:50 -0700 | [diff] [blame] | 84 |  | 
| Guido van Rossum | 90fb914 | 2013-10-30 14:44:05 -0700 | [diff] [blame] | 85 | class _WaitHandleFuture(futures.Future): | 
 | 86 |     """Subclass of Future which represents a wait handle.""" | 
 | 87 |  | 
| Victor Stinner | 18a28dc | 2014-07-25 13:05:20 +0200 | [diff] [blame] | 88 |     def __init__(self, handle, wait_handle, *, loop=None): | 
| Guido van Rossum | 90fb914 | 2013-10-30 14:44:05 -0700 | [diff] [blame] | 89 |         super().__init__(loop=loop) | 
| Victor Stinner | 18a28dc | 2014-07-25 13:05:20 +0200 | [diff] [blame] | 90 |         self._handle = handle | 
| Guido van Rossum | 90fb914 | 2013-10-30 14:44:05 -0700 | [diff] [blame] | 91 |         self._wait_handle = wait_handle | 
 | 92 |  | 
| Victor Stinner | 18a28dc | 2014-07-25 13:05:20 +0200 | [diff] [blame] | 93 |     def _poll(self): | 
 | 94 |         # non-blocking wait: use a timeout of 0 millisecond | 
 | 95 |         return (_winapi.WaitForSingleObject(self._handle, 0) == | 
 | 96 |                 _winapi.WAIT_OBJECT_0) | 
 | 97 |  | 
 | 98 |     def __repr__(self): | 
 | 99 |         info = [self._state.lower()] | 
 | 100 |         if self._wait_handle: | 
 | 101 |             state = 'pending' if self._poll() else 'completed' | 
 | 102 |             info.append('wait_handle=<%s, %#x>' % (state, self._wait_handle)) | 
 | 103 |         info.append('handle=<%#x>' % self._handle) | 
 | 104 |         if self._state == futures._FINISHED: | 
 | 105 |             info.append(self._format_result()) | 
 | 106 |         if self._callbacks: | 
 | 107 |             info.append(self._format_callbacks()) | 
 | 108 |         return '<%s %s>' % (self.__class__.__name__, ' '.join(info)) | 
 | 109 |  | 
| Victor Stinner | fea6a10 | 2014-07-25 00:54:53 +0200 | [diff] [blame] | 110 |     def _unregister(self): | 
 | 111 |         if self._wait_handle is None: | 
 | 112 |             return | 
| Guido van Rossum | 90fb914 | 2013-10-30 14:44:05 -0700 | [diff] [blame] | 113 |         try: | 
 | 114 |             _overlapped.UnregisterWait(self._wait_handle) | 
 | 115 |         except OSError as e: | 
 | 116 |             if e.winerror != _overlapped.ERROR_IO_PENDING: | 
 | 117 |                 raise | 
| Victor Stinner | fea6a10 | 2014-07-25 00:54:53 +0200 | [diff] [blame] | 118 |             # ERROR_IO_PENDING is not an error, the wait was unregistered | 
 | 119 |         self._wait_handle = None | 
 | 120 |  | 
 | 121 |     def cancel(self): | 
 | 122 |         self._unregister() | 
| Victor Stinner | 18a28dc | 2014-07-25 13:05:20 +0200 | [diff] [blame] | 123 |         return super().cancel() | 
| Guido van Rossum | 90fb914 | 2013-10-30 14:44:05 -0700 | [diff] [blame] | 124 |  | 
 | 125 |  | 
| Guido van Rossum | 27b7c7e | 2013-10-17 13:40:50 -0700 | [diff] [blame] | 126 | class PipeServer(object): | 
 | 127 |     """Class representing a pipe server. | 
 | 128 |  | 
 | 129 |     This is much like a bound, listening socket. | 
 | 130 |     """ | 
 | 131 |     def __init__(self, address): | 
 | 132 |         self._address = address | 
 | 133 |         self._free_instances = weakref.WeakSet() | 
 | 134 |         self._pipe = self._server_pipe_handle(True) | 
 | 135 |  | 
 | 136 |     def _get_unconnected_pipe(self): | 
 | 137 |         # Create new instance and return previous one.  This ensures | 
 | 138 |         # that (until the server is closed) there is always at least | 
 | 139 |         # one pipe handle for address.  Therefore if a client attempt | 
 | 140 |         # to connect it will not fail with FileNotFoundError. | 
 | 141 |         tmp, self._pipe = self._pipe, self._server_pipe_handle(False) | 
 | 142 |         return tmp | 
 | 143 |  | 
 | 144 |     def _server_pipe_handle(self, first): | 
 | 145 |         # Return a wrapper for a new pipe handle. | 
 | 146 |         if self._address is None: | 
 | 147 |             return None | 
 | 148 |         flags = _winapi.PIPE_ACCESS_DUPLEX | _winapi.FILE_FLAG_OVERLAPPED | 
 | 149 |         if first: | 
 | 150 |             flags |= _winapi.FILE_FLAG_FIRST_PIPE_INSTANCE | 
 | 151 |         h = _winapi.CreateNamedPipe( | 
 | 152 |             self._address, flags, | 
 | 153 |             _winapi.PIPE_TYPE_MESSAGE | _winapi.PIPE_READMODE_MESSAGE | | 
 | 154 |             _winapi.PIPE_WAIT, | 
 | 155 |             _winapi.PIPE_UNLIMITED_INSTANCES, | 
 | 156 |             windows_utils.BUFSIZE, windows_utils.BUFSIZE, | 
 | 157 |             _winapi.NMPWAIT_WAIT_FOREVER, _winapi.NULL) | 
 | 158 |         pipe = windows_utils.PipeHandle(h) | 
 | 159 |         self._free_instances.add(pipe) | 
 | 160 |         return pipe | 
 | 161 |  | 
 | 162 |     def close(self): | 
 | 163 |         # Close all instances which have not been connected to by a client. | 
 | 164 |         if self._address is not None: | 
 | 165 |             for pipe in self._free_instances: | 
 | 166 |                 pipe.close() | 
 | 167 |             self._pipe = None | 
 | 168 |             self._address = None | 
 | 169 |             self._free_instances.clear() | 
 | 170 |  | 
 | 171 |     __del__ = close | 
 | 172 |  | 
 | 173 |  | 
| Guido van Rossum | 0eaa5ac | 2013-11-04 15:50:46 -0800 | [diff] [blame] | 174 | class _WindowsSelectorEventLoop(selector_events.BaseSelectorEventLoop): | 
| Guido van Rossum | 27b7c7e | 2013-10-17 13:40:50 -0700 | [diff] [blame] | 175 |     """Windows version of selector event loop.""" | 
 | 176 |  | 
 | 177 |     def _socketpair(self): | 
 | 178 |         return windows_utils.socketpair() | 
 | 179 |  | 
 | 180 |  | 
 | 181 | class ProactorEventLoop(proactor_events.BaseProactorEventLoop): | 
 | 182 |     """Windows version of proactor event loop using IOCP.""" | 
 | 183 |  | 
 | 184 |     def __init__(self, proactor=None): | 
 | 185 |         if proactor is None: | 
 | 186 |             proactor = IocpProactor() | 
 | 187 |         super().__init__(proactor) | 
 | 188 |  | 
 | 189 |     def _socketpair(self): | 
 | 190 |         return windows_utils.socketpair() | 
 | 191 |  | 
| Victor Stinner | f951d28 | 2014-06-29 00:46:45 +0200 | [diff] [blame] | 192 |     @coroutine | 
| Guido van Rossum | 27b7c7e | 2013-10-17 13:40:50 -0700 | [diff] [blame] | 193 |     def create_pipe_connection(self, protocol_factory, address): | 
 | 194 |         f = self._proactor.connect_pipe(address) | 
 | 195 |         pipe = yield from f | 
 | 196 |         protocol = protocol_factory() | 
 | 197 |         trans = self._make_duplex_pipe_transport(pipe, protocol, | 
 | 198 |                                                  extra={'addr': address}) | 
 | 199 |         return trans, protocol | 
 | 200 |  | 
| Victor Stinner | f951d28 | 2014-06-29 00:46:45 +0200 | [diff] [blame] | 201 |     @coroutine | 
| Guido van Rossum | 27b7c7e | 2013-10-17 13:40:50 -0700 | [diff] [blame] | 202 |     def start_serving_pipe(self, protocol_factory, address): | 
 | 203 |         server = PipeServer(address) | 
| Guido van Rossum | a8d630a | 2013-11-01 14:20:55 -0700 | [diff] [blame] | 204 |  | 
| Guido van Rossum | 27b7c7e | 2013-10-17 13:40:50 -0700 | [diff] [blame] | 205 |         def loop(f=None): | 
 | 206 |             pipe = None | 
 | 207 |             try: | 
 | 208 |                 if f: | 
 | 209 |                     pipe = f.result() | 
 | 210 |                     server._free_instances.discard(pipe) | 
 | 211 |                     protocol = protocol_factory() | 
 | 212 |                     self._make_duplex_pipe_transport( | 
 | 213 |                         pipe, protocol, extra={'addr': address}) | 
 | 214 |                 pipe = server._get_unconnected_pipe() | 
 | 215 |                 if pipe is None: | 
 | 216 |                     return | 
 | 217 |                 f = self._proactor.accept_pipe(pipe) | 
| Yury Selivanov | ff827f0 | 2014-02-18 18:02:19 -0500 | [diff] [blame] | 218 |             except OSError as exc: | 
| Guido van Rossum | 27b7c7e | 2013-10-17 13:40:50 -0700 | [diff] [blame] | 219 |                 if pipe and pipe.fileno() != -1: | 
| Yury Selivanov | ff827f0 | 2014-02-18 18:02:19 -0500 | [diff] [blame] | 220 |                     self.call_exception_handler({ | 
 | 221 |                         'message': 'Pipe accept failed', | 
 | 222 |                         'exception': exc, | 
 | 223 |                         'pipe': pipe, | 
 | 224 |                     }) | 
| Guido van Rossum | 27b7c7e | 2013-10-17 13:40:50 -0700 | [diff] [blame] | 225 |                     pipe.close() | 
 | 226 |             except futures.CancelledError: | 
 | 227 |                 if pipe: | 
 | 228 |                     pipe.close() | 
 | 229 |             else: | 
 | 230 |                 f.add_done_callback(loop) | 
| Guido van Rossum | a8d630a | 2013-11-01 14:20:55 -0700 | [diff] [blame] | 231 |  | 
| Guido van Rossum | 27b7c7e | 2013-10-17 13:40:50 -0700 | [diff] [blame] | 232 |         self.call_soon(loop) | 
 | 233 |         return [server] | 
 | 234 |  | 
| Victor Stinner | f951d28 | 2014-06-29 00:46:45 +0200 | [diff] [blame] | 235 |     @coroutine | 
| Guido van Rossum | 5969128 | 2013-10-30 14:52:03 -0700 | [diff] [blame] | 236 |     def _make_subprocess_transport(self, protocol, args, shell, | 
 | 237 |                                    stdin, stdout, stderr, bufsize, | 
 | 238 |                                    extra=None, **kwargs): | 
 | 239 |         transp = _WindowsSubprocessTransport(self, protocol, args, shell, | 
 | 240 |                                              stdin, stdout, stderr, bufsize, | 
| Victor Stinner | 73f10fd | 2014-01-29 14:32:20 -0800 | [diff] [blame] | 241 |                                              extra=extra, **kwargs) | 
| Guido van Rossum | 5969128 | 2013-10-30 14:52:03 -0700 | [diff] [blame] | 242 |         yield from transp._post_init() | 
 | 243 |         return transp | 
 | 244 |  | 
| Guido van Rossum | 27b7c7e | 2013-10-17 13:40:50 -0700 | [diff] [blame] | 245 |  | 
 | 246 | class IocpProactor: | 
 | 247 |     """Proactor implementation using IOCP.""" | 
 | 248 |  | 
 | 249 |     def __init__(self, concurrency=0xffffffff): | 
 | 250 |         self._loop = None | 
 | 251 |         self._results = [] | 
 | 252 |         self._iocp = _overlapped.CreateIoCompletionPort( | 
 | 253 |             _overlapped.INVALID_HANDLE_VALUE, NULL, 0, concurrency) | 
 | 254 |         self._cache = {} | 
 | 255 |         self._registered = weakref.WeakSet() | 
 | 256 |         self._stopped_serving = weakref.WeakSet() | 
 | 257 |  | 
| Victor Stinner | fea6a10 | 2014-07-25 00:54:53 +0200 | [diff] [blame] | 258 |     def __repr__(self): | 
 | 259 |         return ('<%s overlapped#=%s result#=%s>' | 
 | 260 |                 % (self.__class__.__name__, len(self._cache), | 
 | 261 |                    len(self._results))) | 
 | 262 |  | 
| Guido van Rossum | 27b7c7e | 2013-10-17 13:40:50 -0700 | [diff] [blame] | 263 |     def set_loop(self, loop): | 
 | 264 |         self._loop = loop | 
 | 265 |  | 
 | 266 |     def select(self, timeout=None): | 
 | 267 |         if not self._results: | 
 | 268 |             self._poll(timeout) | 
 | 269 |         tmp = self._results | 
 | 270 |         self._results = [] | 
 | 271 |         return tmp | 
 | 272 |  | 
 | 273 |     def recv(self, conn, nbytes, flags=0): | 
 | 274 |         self._register_with_iocp(conn) | 
 | 275 |         ov = _overlapped.Overlapped(NULL) | 
 | 276 |         if isinstance(conn, socket.socket): | 
 | 277 |             ov.WSARecv(conn.fileno(), nbytes, flags) | 
 | 278 |         else: | 
 | 279 |             ov.ReadFile(conn.fileno(), nbytes) | 
| Guido van Rossum | a8d630a | 2013-11-01 14:20:55 -0700 | [diff] [blame] | 280 |  | 
| Victor Stinner | c89c8a7 | 2014-02-26 17:35:30 +0100 | [diff] [blame] | 281 |         def finish_recv(trans, key, ov): | 
| Guido van Rossum | 27b7c7e | 2013-10-17 13:40:50 -0700 | [diff] [blame] | 282 |             try: | 
 | 283 |                 return ov.getresult() | 
 | 284 |             except OSError as exc: | 
 | 285 |                 if exc.winerror == _overlapped.ERROR_NETNAME_DELETED: | 
 | 286 |                     raise ConnectionResetError(*exc.args) | 
 | 287 |                 else: | 
 | 288 |                     raise | 
| Guido van Rossum | a8d630a | 2013-11-01 14:20:55 -0700 | [diff] [blame] | 289 |  | 
| Victor Stinner | c89c8a7 | 2014-02-26 17:35:30 +0100 | [diff] [blame] | 290 |         return self._register(ov, conn, finish_recv) | 
| Guido van Rossum | 27b7c7e | 2013-10-17 13:40:50 -0700 | [diff] [blame] | 291 |  | 
 | 292 |     def send(self, conn, buf, flags=0): | 
 | 293 |         self._register_with_iocp(conn) | 
 | 294 |         ov = _overlapped.Overlapped(NULL) | 
 | 295 |         if isinstance(conn, socket.socket): | 
 | 296 |             ov.WSASend(conn.fileno(), buf, flags) | 
 | 297 |         else: | 
 | 298 |             ov.WriteFile(conn.fileno(), buf) | 
| Guido van Rossum | a8d630a | 2013-11-01 14:20:55 -0700 | [diff] [blame] | 299 |  | 
| Victor Stinner | c89c8a7 | 2014-02-26 17:35:30 +0100 | [diff] [blame] | 300 |         def finish_send(trans, key, ov): | 
| Guido van Rossum | 27b7c7e | 2013-10-17 13:40:50 -0700 | [diff] [blame] | 301 |             try: | 
 | 302 |                 return ov.getresult() | 
 | 303 |             except OSError as exc: | 
 | 304 |                 if exc.winerror == _overlapped.ERROR_NETNAME_DELETED: | 
 | 305 |                     raise ConnectionResetError(*exc.args) | 
 | 306 |                 else: | 
 | 307 |                     raise | 
| Guido van Rossum | a8d630a | 2013-11-01 14:20:55 -0700 | [diff] [blame] | 308 |  | 
| Victor Stinner | c89c8a7 | 2014-02-26 17:35:30 +0100 | [diff] [blame] | 309 |         return self._register(ov, conn, finish_send) | 
| Guido van Rossum | 27b7c7e | 2013-10-17 13:40:50 -0700 | [diff] [blame] | 310 |  | 
 | 311 |     def accept(self, listener): | 
 | 312 |         self._register_with_iocp(listener) | 
 | 313 |         conn = self._get_accept_socket(listener.family) | 
 | 314 |         ov = _overlapped.Overlapped(NULL) | 
 | 315 |         ov.AcceptEx(listener.fileno(), conn.fileno()) | 
| Guido van Rossum | a8d630a | 2013-11-01 14:20:55 -0700 | [diff] [blame] | 316 |  | 
| Guido van Rossum | 27b7c7e | 2013-10-17 13:40:50 -0700 | [diff] [blame] | 317 |         def finish_accept(trans, key, ov): | 
 | 318 |             ov.getresult() | 
 | 319 |             # Use SO_UPDATE_ACCEPT_CONTEXT so getsockname() etc work. | 
 | 320 |             buf = struct.pack('@P', listener.fileno()) | 
 | 321 |             conn.setsockopt(socket.SOL_SOCKET, | 
 | 322 |                             _overlapped.SO_UPDATE_ACCEPT_CONTEXT, buf) | 
 | 323 |             conn.settimeout(listener.gettimeout()) | 
 | 324 |             return conn, conn.getpeername() | 
| Guido van Rossum | a8d630a | 2013-11-01 14:20:55 -0700 | [diff] [blame] | 325 |  | 
| Victor Stinner | f951d28 | 2014-06-29 00:46:45 +0200 | [diff] [blame] | 326 |         @coroutine | 
| Victor Stinner | 7de2646 | 2014-01-11 00:03:21 +0100 | [diff] [blame] | 327 |         def accept_coro(future, conn): | 
 | 328 |             # Coroutine closing the accept socket if the future is cancelled | 
 | 329 |             try: | 
 | 330 |                 yield from future | 
 | 331 |             except futures.CancelledError: | 
 | 332 |                 conn.close() | 
 | 333 |                 raise | 
 | 334 |  | 
 | 335 |         future = self._register(ov, listener, finish_accept) | 
 | 336 |         coro = accept_coro(future, conn) | 
 | 337 |         tasks.async(coro, loop=self._loop) | 
 | 338 |         return future | 
| Guido van Rossum | 27b7c7e | 2013-10-17 13:40:50 -0700 | [diff] [blame] | 339 |  | 
 | 340 |     def connect(self, conn, address): | 
 | 341 |         self._register_with_iocp(conn) | 
 | 342 |         # The socket needs to be locally bound before we call ConnectEx(). | 
 | 343 |         try: | 
 | 344 |             _overlapped.BindLocal(conn.fileno(), conn.family) | 
 | 345 |         except OSError as e: | 
 | 346 |             if e.winerror != errno.WSAEINVAL: | 
 | 347 |                 raise | 
 | 348 |             # Probably already locally bound; check using getsockname(). | 
 | 349 |             if conn.getsockname()[1] == 0: | 
 | 350 |                 raise | 
 | 351 |         ov = _overlapped.Overlapped(NULL) | 
 | 352 |         ov.ConnectEx(conn.fileno(), address) | 
| Guido van Rossum | a8d630a | 2013-11-01 14:20:55 -0700 | [diff] [blame] | 353 |  | 
| Guido van Rossum | 27b7c7e | 2013-10-17 13:40:50 -0700 | [diff] [blame] | 354 |         def finish_connect(trans, key, ov): | 
 | 355 |             ov.getresult() | 
 | 356 |             # Use SO_UPDATE_CONNECT_CONTEXT so getsockname() etc work. | 
 | 357 |             conn.setsockopt(socket.SOL_SOCKET, | 
 | 358 |                             _overlapped.SO_UPDATE_CONNECT_CONTEXT, 0) | 
 | 359 |             return conn | 
| Guido van Rossum | a8d630a | 2013-11-01 14:20:55 -0700 | [diff] [blame] | 360 |  | 
| Guido van Rossum | 27b7c7e | 2013-10-17 13:40:50 -0700 | [diff] [blame] | 361 |         return self._register(ov, conn, finish_connect) | 
 | 362 |  | 
 | 363 |     def accept_pipe(self, pipe): | 
 | 364 |         self._register_with_iocp(pipe) | 
 | 365 |         ov = _overlapped.Overlapped(NULL) | 
 | 366 |         ov.ConnectNamedPipe(pipe.fileno()) | 
| Guido van Rossum | a8d630a | 2013-11-01 14:20:55 -0700 | [diff] [blame] | 367 |  | 
| Victor Stinner | c89c8a7 | 2014-02-26 17:35:30 +0100 | [diff] [blame] | 368 |         def finish_accept_pipe(trans, key, ov): | 
| Guido van Rossum | 27b7c7e | 2013-10-17 13:40:50 -0700 | [diff] [blame] | 369 |             ov.getresult() | 
 | 370 |             return pipe | 
| Guido van Rossum | a8d630a | 2013-11-01 14:20:55 -0700 | [diff] [blame] | 371 |  | 
| Victor Stinner | c89c8a7 | 2014-02-26 17:35:30 +0100 | [diff] [blame] | 372 |         return self._register(ov, pipe, finish_accept_pipe) | 
| Guido van Rossum | 27b7c7e | 2013-10-17 13:40:50 -0700 | [diff] [blame] | 373 |  | 
 | 374 |     def connect_pipe(self, address): | 
 | 375 |         ov = _overlapped.Overlapped(NULL) | 
 | 376 |         ov.WaitNamedPipeAndConnect(address, self._iocp, ov.address) | 
| Guido van Rossum | a8d630a | 2013-11-01 14:20:55 -0700 | [diff] [blame] | 377 |  | 
| Victor Stinner | c89c8a7 | 2014-02-26 17:35:30 +0100 | [diff] [blame] | 378 |         def finish_connect_pipe(err, handle, ov): | 
| Guido van Rossum | 27b7c7e | 2013-10-17 13:40:50 -0700 | [diff] [blame] | 379 |             # err, handle were arguments passed to PostQueuedCompletionStatus() | 
 | 380 |             # in a function run in a thread pool. | 
 | 381 |             if err == _overlapped.ERROR_SEM_TIMEOUT: | 
 | 382 |                 # Connection did not succeed within time limit. | 
 | 383 |                 msg = _overlapped.FormatMessage(err) | 
 | 384 |                 raise ConnectionRefusedError(0, msg, None, err) | 
 | 385 |             elif err != 0: | 
 | 386 |                 msg = _overlapped.FormatMessage(err) | 
 | 387 |                 raise OSError(0, msg, None, err) | 
 | 388 |             else: | 
 | 389 |                 return windows_utils.PipeHandle(handle) | 
| Guido van Rossum | a8d630a | 2013-11-01 14:20:55 -0700 | [diff] [blame] | 390 |  | 
| Victor Stinner | c89c8a7 | 2014-02-26 17:35:30 +0100 | [diff] [blame] | 391 |         return self._register(ov, None, finish_connect_pipe, wait_for_post=True) | 
| Guido van Rossum | 27b7c7e | 2013-10-17 13:40:50 -0700 | [diff] [blame] | 392 |  | 
| Guido van Rossum | 90fb914 | 2013-10-30 14:44:05 -0700 | [diff] [blame] | 393 |     def wait_for_handle(self, handle, timeout=None): | 
 | 394 |         if timeout is None: | 
 | 395 |             ms = _winapi.INFINITE | 
 | 396 |         else: | 
| Victor Stinner | f2e1768 | 2014-01-31 16:25:24 +0100 | [diff] [blame] | 397 |             # RegisterWaitForSingleObject() has a resolution of 1 millisecond, | 
 | 398 |             # round away from zero to wait *at least* timeout seconds. | 
 | 399 |             ms = math.ceil(timeout * 1e3) | 
| Guido van Rossum | 90fb914 | 2013-10-30 14:44:05 -0700 | [diff] [blame] | 400 |  | 
 | 401 |         # We only create ov so we can use ov.address as a key for the cache. | 
 | 402 |         ov = _overlapped.Overlapped(NULL) | 
 | 403 |         wh = _overlapped.RegisterWaitWithQueue( | 
 | 404 |             handle, self._iocp, ov.address, ms) | 
| Victor Stinner | 18a28dc | 2014-07-25 13:05:20 +0200 | [diff] [blame] | 405 |         f = _WaitHandleFuture(handle, wh, loop=self._loop) | 
| Guido van Rossum | 90fb914 | 2013-10-30 14:44:05 -0700 | [diff] [blame] | 406 |  | 
| Victor Stinner | c89c8a7 | 2014-02-26 17:35:30 +0100 | [diff] [blame] | 407 |         def finish_wait_for_handle(trans, key, ov): | 
| Richard Oudkerk | 71196e7 | 2013-11-24 17:50:40 +0000 | [diff] [blame] | 408 |             # Note that this second wait means that we should only use | 
 | 409 |             # this with handles types where a successful wait has no | 
 | 410 |             # effect.  So events or processes are all right, but locks | 
 | 411 |             # or semaphores are not.  Also note if the handle is | 
 | 412 |             # signalled and then quickly reset, then we may return | 
 | 413 |             # False even though we have not timed out. | 
| Victor Stinner | 18a28dc | 2014-07-25 13:05:20 +0200 | [diff] [blame] | 414 |             try: | 
 | 415 |                 return f._poll() | 
 | 416 |             finally: | 
 | 417 |                 f._unregister() | 
| Guido van Rossum | 90fb914 | 2013-10-30 14:44:05 -0700 | [diff] [blame] | 418 |  | 
| Victor Stinner | c89c8a7 | 2014-02-26 17:35:30 +0100 | [diff] [blame] | 419 |         self._cache[ov.address] = (f, ov, None, finish_wait_for_handle) | 
| Guido van Rossum | 90fb914 | 2013-10-30 14:44:05 -0700 | [diff] [blame] | 420 |         return f | 
 | 421 |  | 
| Guido van Rossum | 27b7c7e | 2013-10-17 13:40:50 -0700 | [diff] [blame] | 422 |     def _register_with_iocp(self, obj): | 
 | 423 |         # To get notifications of finished ops on this objects sent to the | 
 | 424 |         # completion port, were must register the handle. | 
 | 425 |         if obj not in self._registered: | 
 | 426 |             self._registered.add(obj) | 
 | 427 |             _overlapped.CreateIoCompletionPort(obj.fileno(), self._iocp, 0, 0) | 
 | 428 |             # XXX We could also use SetFileCompletionNotificationModes() | 
 | 429 |             # to avoid sending notifications to completion port of ops | 
 | 430 |             # that succeed immediately. | 
 | 431 |  | 
 | 432 |     def _register(self, ov, obj, callback, wait_for_post=False): | 
 | 433 |         # Return a future which will be set with the result of the | 
 | 434 |         # operation when it completes.  The future's value is actually | 
 | 435 |         # the value returned by callback(). | 
 | 436 |         f = _OverlappedFuture(ov, loop=self._loop) | 
 | 437 |         if ov.pending or wait_for_post: | 
 | 438 |             # Register the overlapped operation for later.  Note that | 
 | 439 |             # we only store obj to prevent it from being garbage | 
 | 440 |             # collected too early. | 
 | 441 |             self._cache[ov.address] = (f, ov, obj, callback) | 
 | 442 |         else: | 
 | 443 |             # The operation has completed, so no need to postpone the | 
 | 444 |             # work.  We cannot take this short cut if we need the | 
 | 445 |             # NumberOfBytes, CompletionKey values returned by | 
 | 446 |             # PostQueuedCompletionStatus(). | 
 | 447 |             try: | 
 | 448 |                 value = callback(None, None, ov) | 
 | 449 |             except OSError as e: | 
 | 450 |                 f.set_exception(e) | 
 | 451 |             else: | 
 | 452 |                 f.set_result(value) | 
 | 453 |         return f | 
 | 454 |  | 
 | 455 |     def _get_accept_socket(self, family): | 
 | 456 |         s = socket.socket(family) | 
 | 457 |         s.settimeout(0) | 
 | 458 |         return s | 
 | 459 |  | 
 | 460 |     def _poll(self, timeout=None): | 
 | 461 |         if timeout is None: | 
 | 462 |             ms = INFINITE | 
 | 463 |         elif timeout < 0: | 
 | 464 |             raise ValueError("negative timeout") | 
 | 465 |         else: | 
| Victor Stinner | f2e1768 | 2014-01-31 16:25:24 +0100 | [diff] [blame] | 466 |             # GetQueuedCompletionStatus() has a resolution of 1 millisecond, | 
 | 467 |             # round away from zero to wait *at least* timeout seconds. | 
 | 468 |             ms = math.ceil(timeout * 1e3) | 
| Guido van Rossum | 27b7c7e | 2013-10-17 13:40:50 -0700 | [diff] [blame] | 469 |             if ms >= INFINITE: | 
 | 470 |                 raise ValueError("timeout too big") | 
 | 471 |         while True: | 
 | 472 |             status = _overlapped.GetQueuedCompletionStatus(self._iocp, ms) | 
 | 473 |             if status is None: | 
 | 474 |                 return | 
 | 475 |             err, transferred, key, address = status | 
 | 476 |             try: | 
 | 477 |                 f, ov, obj, callback = self._cache.pop(address) | 
 | 478 |             except KeyError: | 
 | 479 |                 # key is either zero, or it is used to return a pipe | 
 | 480 |                 # handle which should be closed to avoid a leak. | 
 | 481 |                 if key not in (0, _overlapped.INVALID_HANDLE_VALUE): | 
 | 482 |                     _winapi.CloseHandle(key) | 
 | 483 |                 ms = 0 | 
 | 484 |                 continue | 
| Victor Stinner | 51e44ea | 2014-07-26 00:58:34 +0200 | [diff] [blame^] | 485 |  | 
 | 486 |             if ov.pending: | 
 | 487 |                 # False alarm: the overlapped operation is not completed. | 
 | 488 |                 # FIXME: why do we get false alarms? | 
 | 489 |                 self._cache[address] = (f, ov, obj, callback) | 
 | 490 |                 continue | 
 | 491 |  | 
| Guido van Rossum | 27b7c7e | 2013-10-17 13:40:50 -0700 | [diff] [blame] | 492 |             if obj in self._stopped_serving: | 
 | 493 |                 f.cancel() | 
 | 494 |             elif not f.cancelled(): | 
 | 495 |                 try: | 
 | 496 |                     value = callback(transferred, key, ov) | 
 | 497 |                 except OSError as e: | 
 | 498 |                     f.set_exception(e) | 
 | 499 |                     self._results.append(f) | 
 | 500 |                 else: | 
 | 501 |                     f.set_result(value) | 
 | 502 |                     self._results.append(f) | 
 | 503 |             ms = 0 | 
 | 504 |  | 
 | 505 |     def _stop_serving(self, obj): | 
 | 506 |         # obj is a socket or pipe handle.  It will be closed in | 
 | 507 |         # BaseProactorEventLoop._stop_serving() which will make any | 
 | 508 |         # pending operations fail quickly. | 
 | 509 |         self._stopped_serving.add(obj) | 
 | 510 |  | 
 | 511 |     def close(self): | 
 | 512 |         # Cancel remaining registered operations. | 
| Victor Stinner | fea6a10 | 2014-07-25 00:54:53 +0200 | [diff] [blame] | 513 |         for address, (fut, ov, obj, callback) in list(self._cache.items()): | 
| Guido van Rossum | 27b7c7e | 2013-10-17 13:40:50 -0700 | [diff] [blame] | 514 |             if obj is None: | 
 | 515 |                 # The operation was started with connect_pipe() which | 
 | 516 |                 # queues a task to Windows' thread pool.  This cannot | 
 | 517 |                 # be cancelled, so just forget it. | 
 | 518 |                 del self._cache[address] | 
 | 519 |             else: | 
 | 520 |                 try: | 
| Victor Stinner | fea6a10 | 2014-07-25 00:54:53 +0200 | [diff] [blame] | 521 |                     fut.cancel() | 
 | 522 |                 except OSError as exc: | 
 | 523 |                     if self._loop is not None: | 
 | 524 |                         context = { | 
 | 525 |                             'message': 'Cancelling a future failed', | 
 | 526 |                             'exception': exc, | 
 | 527 |                             'future': fut, | 
 | 528 |                         } | 
 | 529 |                         if fut._source_traceback: | 
 | 530 |                             context['source_traceback'] = fut._source_traceback | 
 | 531 |                         self._loop.call_exception_handler(context) | 
| Guido van Rossum | 27b7c7e | 2013-10-17 13:40:50 -0700 | [diff] [blame] | 532 |  | 
 | 533 |         while self._cache: | 
 | 534 |             if not self._poll(1): | 
| Guido van Rossum | fc29e0f | 2013-10-17 15:39:45 -0700 | [diff] [blame] | 535 |                 logger.debug('taking long time to close proactor') | 
| Guido van Rossum | 27b7c7e | 2013-10-17 13:40:50 -0700 | [diff] [blame] | 536 |  | 
 | 537 |         self._results = [] | 
 | 538 |         if self._iocp is not None: | 
 | 539 |             _winapi.CloseHandle(self._iocp) | 
 | 540 |             self._iocp = None | 
| Guido van Rossum | 5969128 | 2013-10-30 14:52:03 -0700 | [diff] [blame] | 541 |  | 
| Victor Stinner | fea6a10 | 2014-07-25 00:54:53 +0200 | [diff] [blame] | 542 |     def __del__(self): | 
 | 543 |         self.close() | 
 | 544 |  | 
| Guido van Rossum | 5969128 | 2013-10-30 14:52:03 -0700 | [diff] [blame] | 545 |  | 
 | 546 | class _WindowsSubprocessTransport(base_subprocess.BaseSubprocessTransport): | 
 | 547 |  | 
 | 548 |     def _start(self, args, shell, stdin, stdout, stderr, bufsize, **kwargs): | 
 | 549 |         self._proc = windows_utils.Popen( | 
 | 550 |             args, shell=shell, stdin=stdin, stdout=stdout, stderr=stderr, | 
 | 551 |             bufsize=bufsize, **kwargs) | 
| Guido van Rossum | a8d630a | 2013-11-01 14:20:55 -0700 | [diff] [blame] | 552 |  | 
| Guido van Rossum | 5969128 | 2013-10-30 14:52:03 -0700 | [diff] [blame] | 553 |         def callback(f): | 
 | 554 |             returncode = self._proc.poll() | 
 | 555 |             self._process_exited(returncode) | 
| Guido van Rossum | a8d630a | 2013-11-01 14:20:55 -0700 | [diff] [blame] | 556 |  | 
| Guido van Rossum | 5969128 | 2013-10-30 14:52:03 -0700 | [diff] [blame] | 557 |         f = self._loop._proactor.wait_for_handle(int(self._proc._handle)) | 
 | 558 |         f.add_done_callback(callback) | 
| Guido van Rossum | 0eaa5ac | 2013-11-04 15:50:46 -0800 | [diff] [blame] | 559 |  | 
 | 560 |  | 
 | 561 | SelectorEventLoop = _WindowsSelectorEventLoop | 
 | 562 |  | 
 | 563 |  | 
 | 564 | class _WindowsDefaultEventLoopPolicy(events.BaseDefaultEventLoopPolicy): | 
 | 565 |     _loop_factory = SelectorEventLoop | 
 | 566 |  | 
 | 567 |  | 
 | 568 | DefaultEventLoopPolicy = _WindowsDefaultEventLoopPolicy |