David Reiss | ea2cba8 | 2009-03-30 21:35:00 +0000 | [diff] [blame] | 1 | # |
| 2 | # Licensed to the Apache Software Foundation (ASF) under one |
| 3 | # or more contributor license agreements. See the NOTICE file |
| 4 | # distributed with this work for additional information |
| 5 | # regarding copyright ownership. The ASF licenses this file |
| 6 | # to you under the Apache License, Version 2.0 (the |
| 7 | # "License"); you may not use this file except in compliance |
| 8 | # with the License. You may obtain a copy of the License at |
| 9 | # |
| 10 | # http://www.apache.org/licenses/LICENSE-2.0 |
| 11 | # |
| 12 | # Unless required by applicable law or agreed to in writing, |
| 13 | # software distributed under the License is distributed on an |
| 14 | # "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY |
| 15 | # KIND, either express or implied. See the License for the |
| 16 | # specific language governing permissions and limitations |
| 17 | # under the License. |
| 18 | # |
David Reiss | 7442127 | 2008-11-07 23:09:31 +0000 | [diff] [blame] | 19 | """Implementation of non-blocking server. |
| 20 | |
Bryan Duxbury | 6972041 | 2012-01-03 17:32:30 +0000 | [diff] [blame] | 21 | The main idea of the server is to receive and send requests |
| 22 | only from the main thread. |
David Reiss | 7442127 | 2008-11-07 23:09:31 +0000 | [diff] [blame] | 23 | |
Bryan Duxbury | 6972041 | 2012-01-03 17:32:30 +0000 | [diff] [blame] | 24 | The thread poool should be sized for concurrent tasks, not |
| 25 | maximum connections |
David Reiss | 7442127 | 2008-11-07 23:09:31 +0000 | [diff] [blame] | 26 | """ |
Konrad Grochowski | 3a724e3 | 2014-08-12 11:48:29 -0400 | [diff] [blame] | 27 | |
David Reiss | 7442127 | 2008-11-07 23:09:31 +0000 | [diff] [blame] | 28 | import logging |
Nobuaki Sukegawa | 10308cb | 2016-02-03 01:57:03 +0900 | [diff] [blame] | 29 | import select |
| 30 | import socket |
| 31 | import struct |
| 32 | import threading |
David Reiss | 7442127 | 2008-11-07 23:09:31 +0000 | [diff] [blame] | 33 | |
Nobuaki Sukegawa | 4626fd8 | 2017-02-12 21:11:36 +0900 | [diff] [blame] | 34 | from collections import deque |
Nobuaki Sukegawa | 760511f | 2015-11-06 21:24:16 +0900 | [diff] [blame] | 35 | from six.moves import queue |
| 36 | |
David Reiss | 7442127 | 2008-11-07 23:09:31 +0000 | [diff] [blame] | 37 | from thrift.transport import TTransport |
| 38 | from thrift.protocol.TBinaryProtocol import TBinaryProtocolFactory |
| 39 | |
| 40 | __all__ = ['TNonblockingServer'] |
| 41 | |
Nobuaki Sukegawa | 10308cb | 2016-02-03 01:57:03 +0900 | [diff] [blame] | 42 | logger = logging.getLogger(__name__) |
| 43 | |
Bryan Duxbury | 6972041 | 2012-01-03 17:32:30 +0000 | [diff] [blame] | 44 | |
David Reiss | 7442127 | 2008-11-07 23:09:31 +0000 | [diff] [blame] | 45 | class Worker(threading.Thread): |
| 46 | """Worker is a small helper to process incoming connection.""" |
Bryan Duxbury | 6972041 | 2012-01-03 17:32:30 +0000 | [diff] [blame] | 47 | |
David Reiss | 7442127 | 2008-11-07 23:09:31 +0000 | [diff] [blame] | 48 | def __init__(self, queue): |
| 49 | threading.Thread.__init__(self) |
| 50 | self.queue = queue |
| 51 | |
| 52 | def run(self): |
| 53 | """Process queries from task queue, stop if processor is None.""" |
| 54 | while True: |
| 55 | try: |
| 56 | processor, iprot, oprot, otrans, callback = self.queue.get() |
| 57 | if processor is None: |
| 58 | break |
| 59 | processor.process(iprot, oprot) |
| 60 | callback(True, otrans.getvalue()) |
| 61 | except Exception: |
Nobuaki Sukegawa | 4626fd8 | 2017-02-12 21:11:36 +0900 | [diff] [blame] | 62 | logger.exception("Exception while processing request", exc_info=True) |
Nobuaki Sukegawa | 760511f | 2015-11-06 21:24:16 +0900 | [diff] [blame] | 63 | callback(False, b'') |
David Reiss | 7442127 | 2008-11-07 23:09:31 +0000 | [diff] [blame] | 64 | |
James E. King, III | 0ad20bd | 2017-09-30 15:44:16 -0700 | [diff] [blame] | 65 | |
David Reiss | 7442127 | 2008-11-07 23:09:31 +0000 | [diff] [blame] | 66 | WAIT_LEN = 0 |
| 67 | WAIT_MESSAGE = 1 |
| 68 | WAIT_PROCESS = 2 |
| 69 | SEND_ANSWER = 3 |
| 70 | CLOSED = 4 |
| 71 | |
Bryan Duxbury | 6972041 | 2012-01-03 17:32:30 +0000 | [diff] [blame] | 72 | |
David Reiss | 7442127 | 2008-11-07 23:09:31 +0000 | [diff] [blame] | 73 | def locked(func): |
Bryan Duxbury | 6972041 | 2012-01-03 17:32:30 +0000 | [diff] [blame] | 74 | """Decorator which locks self.lock.""" |
David Reiss | 7442127 | 2008-11-07 23:09:31 +0000 | [diff] [blame] | 75 | def nested(self, *args, **kwargs): |
| 76 | self.lock.acquire() |
| 77 | try: |
| 78 | return func(self, *args, **kwargs) |
| 79 | finally: |
| 80 | self.lock.release() |
| 81 | return nested |
| 82 | |
Bryan Duxbury | 6972041 | 2012-01-03 17:32:30 +0000 | [diff] [blame] | 83 | |
David Reiss | 7442127 | 2008-11-07 23:09:31 +0000 | [diff] [blame] | 84 | def socket_exception(func): |
Bryan Duxbury | 6972041 | 2012-01-03 17:32:30 +0000 | [diff] [blame] | 85 | """Decorator close object on socket.error.""" |
David Reiss | 7442127 | 2008-11-07 23:09:31 +0000 | [diff] [blame] | 86 | def read(self, *args, **kwargs): |
| 87 | try: |
| 88 | return func(self, *args, **kwargs) |
| 89 | except socket.error: |
Nobuaki Sukegawa | 4626fd8 | 2017-02-12 21:11:36 +0900 | [diff] [blame] | 90 | logger.debug('ignoring socket exception', exc_info=True) |
David Reiss | 7442127 | 2008-11-07 23:09:31 +0000 | [diff] [blame] | 91 | self.close() |
| 92 | return read |
| 93 | |
Bryan Duxbury | 6972041 | 2012-01-03 17:32:30 +0000 | [diff] [blame] | 94 | |
Nobuaki Sukegawa | 4626fd8 | 2017-02-12 21:11:36 +0900 | [diff] [blame] | 95 | class Message(object): |
| 96 | def __init__(self, offset, len_, header): |
| 97 | self.offset = offset |
| 98 | self.len = len_ |
| 99 | self.buffer = None |
| 100 | self.is_header = header |
| 101 | |
| 102 | @property |
| 103 | def end(self): |
| 104 | return self.offset + self.len |
| 105 | |
| 106 | |
Nobuaki Sukegawa | b9c859a | 2015-12-21 01:10:25 +0900 | [diff] [blame] | 107 | class Connection(object): |
David Reiss | 7442127 | 2008-11-07 23:09:31 +0000 | [diff] [blame] | 108 | """Basic class is represented connection. |
Bryan Duxbury | 6972041 | 2012-01-03 17:32:30 +0000 | [diff] [blame] | 109 | |
David Reiss | 7442127 | 2008-11-07 23:09:31 +0000 | [diff] [blame] | 110 | It can be in state: |
| 111 | WAIT_LEN --- connection is reading request len. |
| 112 | WAIT_MESSAGE --- connection is reading request. |
Bryan Duxbury | 6972041 | 2012-01-03 17:32:30 +0000 | [diff] [blame] | 113 | WAIT_PROCESS --- connection has just read whole request and |
| 114 | waits for call ready routine. |
David Reiss | 7442127 | 2008-11-07 23:09:31 +0000 | [diff] [blame] | 115 | SEND_ANSWER --- connection is sending answer string (including length |
Bryan Duxbury | 6972041 | 2012-01-03 17:32:30 +0000 | [diff] [blame] | 116 | of answer). |
David Reiss | 7442127 | 2008-11-07 23:09:31 +0000 | [diff] [blame] | 117 | CLOSED --- socket was closed and connection should be deleted. |
| 118 | """ |
| 119 | def __init__(self, new_socket, wake_up): |
| 120 | self.socket = new_socket |
| 121 | self.socket.setblocking(False) |
| 122 | self.status = WAIT_LEN |
| 123 | self.len = 0 |
Nobuaki Sukegawa | 4626fd8 | 2017-02-12 21:11:36 +0900 | [diff] [blame] | 124 | self.received = deque() |
| 125 | self._reading = Message(0, 4, True) |
| 126 | self._rbuf = b'' |
| 127 | self._wbuf = b'' |
David Reiss | 7442127 | 2008-11-07 23:09:31 +0000 | [diff] [blame] | 128 | self.lock = threading.Lock() |
| 129 | self.wake_up = wake_up |
Nobuaki Sukegawa | 4626fd8 | 2017-02-12 21:11:36 +0900 | [diff] [blame] | 130 | self.remaining = False |
David Reiss | 7442127 | 2008-11-07 23:09:31 +0000 | [diff] [blame] | 131 | |
| 132 | @socket_exception |
| 133 | def read(self): |
| 134 | """Reads data from stream and switch state.""" |
| 135 | assert self.status in (WAIT_LEN, WAIT_MESSAGE) |
Nobuaki Sukegawa | 4626fd8 | 2017-02-12 21:11:36 +0900 | [diff] [blame] | 136 | assert not self.received |
| 137 | buf_size = 8192 |
| 138 | first = True |
| 139 | done = False |
| 140 | while not done: |
| 141 | read = self.socket.recv(buf_size) |
| 142 | rlen = len(read) |
| 143 | done = rlen < buf_size |
| 144 | self._rbuf += read |
| 145 | if first and rlen == 0: |
| 146 | if self.status != WAIT_LEN or self._rbuf: |
| 147 | logger.error('could not read frame from socket') |
| 148 | else: |
| 149 | logger.debug('read zero length. client might have disconnected') |
David Reiss | 7442127 | 2008-11-07 23:09:31 +0000 | [diff] [blame] | 150 | self.close() |
Nobuaki Sukegawa | 4626fd8 | 2017-02-12 21:11:36 +0900 | [diff] [blame] | 151 | while len(self._rbuf) >= self._reading.end: |
| 152 | if self._reading.is_header: |
| 153 | mlen, = struct.unpack('!i', self._rbuf[:4]) |
shangxu | b664cfe | 2020-11-13 18:03:01 +0800 | [diff] [blame] | 154 | if mlen < 0: |
| 155 | logger.error('could not read the head from frame') |
| 156 | self.close() |
| 157 | break |
Nobuaki Sukegawa | 4626fd8 | 2017-02-12 21:11:36 +0900 | [diff] [blame] | 158 | self._reading = Message(self._reading.end, mlen, False) |
| 159 | self.status = WAIT_MESSAGE |
| 160 | else: |
| 161 | self._reading.buffer = self._rbuf |
| 162 | self.received.append(self._reading) |
| 163 | self._rbuf = self._rbuf[self._reading.end:] |
| 164 | self._reading = Message(0, 4, True) |
| 165 | first = False |
| 166 | if self.received: |
David Reiss | 7442127 | 2008-11-07 23:09:31 +0000 | [diff] [blame] | 167 | self.status = WAIT_PROCESS |
Nobuaki Sukegawa | 4626fd8 | 2017-02-12 21:11:36 +0900 | [diff] [blame] | 168 | break |
| 169 | self.remaining = not done |
David Reiss | 7442127 | 2008-11-07 23:09:31 +0000 | [diff] [blame] | 170 | |
| 171 | @socket_exception |
| 172 | def write(self): |
| 173 | """Writes data from socket and switch state.""" |
| 174 | assert self.status == SEND_ANSWER |
Nobuaki Sukegawa | 4626fd8 | 2017-02-12 21:11:36 +0900 | [diff] [blame] | 175 | sent = self.socket.send(self._wbuf) |
| 176 | if sent == len(self._wbuf): |
David Reiss | 7442127 | 2008-11-07 23:09:31 +0000 | [diff] [blame] | 177 | self.status = WAIT_LEN |
Nobuaki Sukegawa | 4626fd8 | 2017-02-12 21:11:36 +0900 | [diff] [blame] | 178 | self._wbuf = b'' |
David Reiss | 7442127 | 2008-11-07 23:09:31 +0000 | [diff] [blame] | 179 | self.len = 0 |
| 180 | else: |
Yubing Dong (Tom) | 00646bb | 2018-01-18 23:55:24 -0800 | [diff] [blame] | 181 | self._wbuf = self._wbuf[sent:] |
David Reiss | 7442127 | 2008-11-07 23:09:31 +0000 | [diff] [blame] | 182 | |
| 183 | @locked |
| 184 | def ready(self, all_ok, message): |
| 185 | """Callback function for switching state and waking up main thread. |
Bryan Duxbury | 6972041 | 2012-01-03 17:32:30 +0000 | [diff] [blame] | 186 | |
David Reiss | 7442127 | 2008-11-07 23:09:31 +0000 | [diff] [blame] | 187 | This function is the only function witch can be called asynchronous. |
Bryan Duxbury | 6972041 | 2012-01-03 17:32:30 +0000 | [diff] [blame] | 188 | |
David Reiss | 7442127 | 2008-11-07 23:09:31 +0000 | [diff] [blame] | 189 | The ready can switch Connection to three states: |
David Reiss | 6ce401d | 2009-03-24 20:01:58 +0000 | [diff] [blame] | 190 | WAIT_LEN if request was oneway. |
David Reiss | 7442127 | 2008-11-07 23:09:31 +0000 | [diff] [blame] | 191 | SEND_ANSWER if request was processed in normal way. |
| 192 | CLOSED if request throws unexpected exception. |
Bryan Duxbury | 6972041 | 2012-01-03 17:32:30 +0000 | [diff] [blame] | 193 | |
David Reiss | 7442127 | 2008-11-07 23:09:31 +0000 | [diff] [blame] | 194 | The one wakes up main thread. |
| 195 | """ |
| 196 | assert self.status == WAIT_PROCESS |
| 197 | if not all_ok: |
| 198 | self.close() |
| 199 | self.wake_up() |
| 200 | return |
Nobuaki Sukegawa | 760511f | 2015-11-06 21:24:16 +0900 | [diff] [blame] | 201 | self.len = 0 |
David Reiss | 7442127 | 2008-11-07 23:09:31 +0000 | [diff] [blame] | 202 | if len(message) == 0: |
David Reiss | c51986f | 2009-03-24 20:01:25 +0000 | [diff] [blame] | 203 | # it was a oneway request, do not write answer |
Nobuaki Sukegawa | 4626fd8 | 2017-02-12 21:11:36 +0900 | [diff] [blame] | 204 | self._wbuf = b'' |
David Reiss | 7442127 | 2008-11-07 23:09:31 +0000 | [diff] [blame] | 205 | self.status = WAIT_LEN |
| 206 | else: |
Nobuaki Sukegawa | 4626fd8 | 2017-02-12 21:11:36 +0900 | [diff] [blame] | 207 | self._wbuf = struct.pack('!i', len(message)) + message |
David Reiss | 7442127 | 2008-11-07 23:09:31 +0000 | [diff] [blame] | 208 | self.status = SEND_ANSWER |
| 209 | self.wake_up() |
| 210 | |
| 211 | @locked |
| 212 | def is_writeable(self): |
Bryan Duxbury | 6972041 | 2012-01-03 17:32:30 +0000 | [diff] [blame] | 213 | """Return True if connection should be added to write list of select""" |
David Reiss | 7442127 | 2008-11-07 23:09:31 +0000 | [diff] [blame] | 214 | return self.status == SEND_ANSWER |
| 215 | |
| 216 | # it's not necessary, but... |
| 217 | @locked |
| 218 | def is_readable(self): |
Bryan Duxbury | 6972041 | 2012-01-03 17:32:30 +0000 | [diff] [blame] | 219 | """Return True if connection should be added to read list of select""" |
David Reiss | 7442127 | 2008-11-07 23:09:31 +0000 | [diff] [blame] | 220 | return self.status in (WAIT_LEN, WAIT_MESSAGE) |
| 221 | |
| 222 | @locked |
| 223 | def is_closed(self): |
Bryan Duxbury | 6972041 | 2012-01-03 17:32:30 +0000 | [diff] [blame] | 224 | """Returns True if connection is closed.""" |
David Reiss | 7442127 | 2008-11-07 23:09:31 +0000 | [diff] [blame] | 225 | return self.status == CLOSED |
| 226 | |
| 227 | def fileno(self): |
Bryan Duxbury | 6972041 | 2012-01-03 17:32:30 +0000 | [diff] [blame] | 228 | """Returns the file descriptor of the associated socket.""" |
David Reiss | 7442127 | 2008-11-07 23:09:31 +0000 | [diff] [blame] | 229 | return self.socket.fileno() |
| 230 | |
| 231 | def close(self): |
Bryan Duxbury | 6972041 | 2012-01-03 17:32:30 +0000 | [diff] [blame] | 232 | """Closes connection""" |
David Reiss | 7442127 | 2008-11-07 23:09:31 +0000 | [diff] [blame] | 233 | self.status = CLOSED |
| 234 | self.socket.close() |
| 235 | |
Bryan Duxbury | 6972041 | 2012-01-03 17:32:30 +0000 | [diff] [blame] | 236 | |
Nobuaki Sukegawa | b9c859a | 2015-12-21 01:10:25 +0900 | [diff] [blame] | 237 | class TNonblockingServer(object): |
David Reiss | 7442127 | 2008-11-07 23:09:31 +0000 | [diff] [blame] | 238 | """Non-blocking server.""" |
Bryan Duxbury | 6972041 | 2012-01-03 17:32:30 +0000 | [diff] [blame] | 239 | |
| 240 | def __init__(self, |
| 241 | processor, |
| 242 | lsocket, |
| 243 | inputProtocolFactory=None, |
| 244 | outputProtocolFactory=None, |
| 245 | threads=10): |
David Reiss | 7442127 | 2008-11-07 23:09:31 +0000 | [diff] [blame] | 246 | self.processor = processor |
| 247 | self.socket = lsocket |
| 248 | self.in_protocol = inputProtocolFactory or TBinaryProtocolFactory() |
| 249 | self.out_protocol = outputProtocolFactory or self.in_protocol |
| 250 | self.threads = int(threads) |
| 251 | self.clients = {} |
Nobuaki Sukegawa | 760511f | 2015-11-06 21:24:16 +0900 | [diff] [blame] | 252 | self.tasks = queue.Queue() |
David Reiss | 7442127 | 2008-11-07 23:09:31 +0000 | [diff] [blame] | 253 | self._read, self._write = socket.socketpair() |
| 254 | self.prepared = False |
Roger Meier | cfff856 | 2012-04-13 14:24:55 +0000 | [diff] [blame] | 255 | self._stop = False |
Yiyang Zhou | 88a45ac | 2021-08-04 21:55:04 +0800 | [diff] [blame^] | 256 | self.poll = select.poll() if hasattr(select, 'poll') else None |
David Reiss | 7442127 | 2008-11-07 23:09:31 +0000 | [diff] [blame] | 257 | |
| 258 | def setNumThreads(self, num): |
| 259 | """Set the number of worker threads that should be created.""" |
| 260 | # implement ThreadPool interface |
Bryan Duxbury | 6972041 | 2012-01-03 17:32:30 +0000 | [diff] [blame] | 261 | assert not self.prepared, "Can't change number of threads after start" |
David Reiss | 7442127 | 2008-11-07 23:09:31 +0000 | [diff] [blame] | 262 | self.threads = num |
| 263 | |
| 264 | def prepare(self): |
| 265 | """Prepares server for serve requests.""" |
Roger Meier | cfff856 | 2012-04-13 14:24:55 +0000 | [diff] [blame] | 266 | if self.prepared: |
| 267 | return |
David Reiss | 7442127 | 2008-11-07 23:09:31 +0000 | [diff] [blame] | 268 | self.socket.listen() |
Nobuaki Sukegawa | 760511f | 2015-11-06 21:24:16 +0900 | [diff] [blame] | 269 | for _ in range(self.threads): |
David Reiss | 7442127 | 2008-11-07 23:09:31 +0000 | [diff] [blame] | 270 | thread = Worker(self.tasks) |
| 271 | thread.setDaemon(True) |
| 272 | thread.start() |
| 273 | self.prepared = True |
| 274 | |
| 275 | def wake_up(self): |
| 276 | """Wake up main thread. |
Bryan Duxbury | 6972041 | 2012-01-03 17:32:30 +0000 | [diff] [blame] | 277 | |
Konrad Grochowski | 3b5dacb | 2014-11-24 10:55:31 +0100 | [diff] [blame] | 278 | The server usually waits in select call in we should terminate one. |
David Reiss | 7442127 | 2008-11-07 23:09:31 +0000 | [diff] [blame] | 279 | The simplest way is using socketpair. |
Bryan Duxbury | 6972041 | 2012-01-03 17:32:30 +0000 | [diff] [blame] | 280 | |
David Reiss | 7442127 | 2008-11-07 23:09:31 +0000 | [diff] [blame] | 281 | Select always wait to read from the first socket of socketpair. |
Bryan Duxbury | 6972041 | 2012-01-03 17:32:30 +0000 | [diff] [blame] | 282 | |
David Reiss | 7442127 | 2008-11-07 23:09:31 +0000 | [diff] [blame] | 283 | In this case, we can just write anything to the second socket from |
Bryan Duxbury | 6972041 | 2012-01-03 17:32:30 +0000 | [diff] [blame] | 284 | socketpair. |
| 285 | """ |
Nobuaki Sukegawa | 760511f | 2015-11-06 21:24:16 +0900 | [diff] [blame] | 286 | self._write.send(b'1') |
David Reiss | 7442127 | 2008-11-07 23:09:31 +0000 | [diff] [blame] | 287 | |
Roger Meier | cfff856 | 2012-04-13 14:24:55 +0000 | [diff] [blame] | 288 | def stop(self): |
| 289 | """Stop the server. |
| 290 | |
| 291 | This method causes the serve() method to return. stop() may be invoked |
| 292 | from within your handler, or from another thread. |
| 293 | |
| 294 | After stop() is called, serve() will return but the server will still |
| 295 | be listening on the socket. serve() may then be called again to resume |
| 296 | processing requests. Alternatively, close() may be called after |
| 297 | serve() returns to close the server socket and shutdown all worker |
| 298 | threads. |
| 299 | """ |
| 300 | self._stop = True |
| 301 | self.wake_up() |
| 302 | |
David Reiss | 7442127 | 2008-11-07 23:09:31 +0000 | [diff] [blame] | 303 | def _select(self): |
| 304 | """Does select on open connections.""" |
| 305 | readable = [self.socket.handle.fileno(), self._read.fileno()] |
| 306 | writable = [] |
Nobuaki Sukegawa | 4626fd8 | 2017-02-12 21:11:36 +0900 | [diff] [blame] | 307 | remaining = [] |
Nobuaki Sukegawa | 760511f | 2015-11-06 21:24:16 +0900 | [diff] [blame] | 308 | for i, connection in list(self.clients.items()): |
David Reiss | 7442127 | 2008-11-07 23:09:31 +0000 | [diff] [blame] | 309 | if connection.is_readable(): |
| 310 | readable.append(connection.fileno()) |
Nobuaki Sukegawa | 4626fd8 | 2017-02-12 21:11:36 +0900 | [diff] [blame] | 311 | if connection.remaining or connection.received: |
| 312 | remaining.append(connection.fileno()) |
David Reiss | 7442127 | 2008-11-07 23:09:31 +0000 | [diff] [blame] | 313 | if connection.is_writeable(): |
| 314 | writable.append(connection.fileno()) |
| 315 | if connection.is_closed(): |
| 316 | del self.clients[i] |
Nobuaki Sukegawa | 4626fd8 | 2017-02-12 21:11:36 +0900 | [diff] [blame] | 317 | if remaining: |
| 318 | return remaining, [], [], False |
| 319 | else: |
| 320 | return select.select(readable, writable, readable) + (True,) |
Bryan Duxbury | 6972041 | 2012-01-03 17:32:30 +0000 | [diff] [blame] | 321 | |
Yiyang Zhou | 88a45ac | 2021-08-04 21:55:04 +0800 | [diff] [blame^] | 322 | def _poll_select(self): |
| 323 | """Does poll on open connections, if available.""" |
| 324 | remaining = [] |
| 325 | |
| 326 | self.poll.register(self.socket.handle.fileno(), select.POLLIN | select.POLLRDNORM) |
| 327 | self.poll.register(self._read.fileno(), select.POLLIN | select.POLLRDNORM) |
| 328 | |
| 329 | for i, connection in list(self.clients.items()): |
| 330 | if connection.is_readable(): |
| 331 | self.poll.register(connection.fileno(), select.POLLIN | select.POLLRDNORM | select.POLLERR | select.POLLHUP | select.POLLNVAL) |
| 332 | if connection.remaining or connection.received: |
| 333 | remaining.append(connection.fileno()) |
| 334 | if connection.is_writeable(): |
| 335 | self.poll.register(connection.fileno(), select.POLLOUT | select.POLLWRNORM) |
| 336 | if connection.is_closed(): |
| 337 | try: |
| 338 | self.poll.unregister(i) |
| 339 | except KeyError: |
| 340 | logger.debug("KeyError in unregistering connections...") |
| 341 | del self.clients[i] |
| 342 | if remaining: |
| 343 | return remaining, [], [], False |
| 344 | |
| 345 | rlist = [] |
| 346 | wlist = [] |
| 347 | xlist = [] |
| 348 | pollres = self.poll.poll() |
| 349 | for fd, event in pollres: |
| 350 | if event & (select.POLLERR | select.POLLHUP | select.POLLNVAL): |
| 351 | xlist.append(fd) |
| 352 | elif event & (select.POLLOUT | select.POLLWRNORM): |
| 353 | wlist.append(fd) |
| 354 | elif event & (select.POLLIN | select.POLLRDNORM): |
| 355 | rlist.append(fd) |
| 356 | else: # should be impossible |
| 357 | logger.debug("reached an impossible state in _poll_select") |
| 358 | xlist.append(fd) |
| 359 | |
| 360 | return rlist, wlist, xlist, True |
| 361 | |
David Reiss | 7442127 | 2008-11-07 23:09:31 +0000 | [diff] [blame] | 362 | def handle(self): |
| 363 | """Handle requests. |
Bryan Duxbury | 6972041 | 2012-01-03 17:32:30 +0000 | [diff] [blame] | 364 | |
| 365 | WARNING! You must call prepare() BEFORE calling handle() |
David Reiss | 7442127 | 2008-11-07 23:09:31 +0000 | [diff] [blame] | 366 | """ |
| 367 | assert self.prepared, "You have to call prepare before handle" |
Yiyang Zhou | 88a45ac | 2021-08-04 21:55:04 +0800 | [diff] [blame^] | 368 | rset, wset, xset, selected = self._select() if not self.poll else self._poll_select() |
David Reiss | 7442127 | 2008-11-07 23:09:31 +0000 | [diff] [blame] | 369 | for readable in rset: |
| 370 | if readable == self._read.fileno(): |
| 371 | # don't care i just need to clean readable flag |
Bryan Duxbury | 6972041 | 2012-01-03 17:32:30 +0000 | [diff] [blame] | 372 | self._read.recv(1024) |
David Reiss | 7442127 | 2008-11-07 23:09:31 +0000 | [diff] [blame] | 373 | elif readable == self.socket.handle.fileno(): |
Nobuaki Sukegawa | 4626fd8 | 2017-02-12 21:11:36 +0900 | [diff] [blame] | 374 | try: |
| 375 | client = self.socket.accept() |
| 376 | if client: |
| 377 | self.clients[client.handle.fileno()] = Connection(client.handle, |
| 378 | self.wake_up) |
| 379 | except socket.error: |
| 380 | logger.debug('error while accepting', exc_info=True) |
David Reiss | 7442127 | 2008-11-07 23:09:31 +0000 | [diff] [blame] | 381 | else: |
| 382 | connection = self.clients[readable] |
Nobuaki Sukegawa | 4626fd8 | 2017-02-12 21:11:36 +0900 | [diff] [blame] | 383 | if selected: |
| 384 | connection.read() |
| 385 | if connection.received: |
| 386 | connection.status = WAIT_PROCESS |
Yiyang Zhou | 88a45ac | 2021-08-04 21:55:04 +0800 | [diff] [blame^] | 387 | if self.poll: |
| 388 | self.poll.unregister(connection.fileno()) |
Nobuaki Sukegawa | 4626fd8 | 2017-02-12 21:11:36 +0900 | [diff] [blame] | 389 | msg = connection.received.popleft() |
| 390 | itransport = TTransport.TMemoryBuffer(msg.buffer, msg.offset) |
David Reiss | 7442127 | 2008-11-07 23:09:31 +0000 | [diff] [blame] | 391 | otransport = TTransport.TMemoryBuffer() |
| 392 | iprot = self.in_protocol.getProtocol(itransport) |
| 393 | oprot = self.out_protocol.getProtocol(otransport) |
Bryan Duxbury | 6972041 | 2012-01-03 17:32:30 +0000 | [diff] [blame] | 394 | self.tasks.put([self.processor, iprot, oprot, |
David Reiss | 7442127 | 2008-11-07 23:09:31 +0000 | [diff] [blame] | 395 | otransport, connection.ready]) |
| 396 | for writeable in wset: |
| 397 | self.clients[writeable].write() |
| 398 | for oob in xset: |
| 399 | self.clients[oob].close() |
David Reiss | 7442127 | 2008-11-07 23:09:31 +0000 | [diff] [blame] | 400 | |
| 401 | def close(self): |
| 402 | """Closes the server.""" |
Nobuaki Sukegawa | 760511f | 2015-11-06 21:24:16 +0900 | [diff] [blame] | 403 | for _ in range(self.threads): |
David Reiss | 7442127 | 2008-11-07 23:09:31 +0000 | [diff] [blame] | 404 | self.tasks.put([None, None, None, None, None]) |
| 405 | self.socket.close() |
| 406 | self.prepared = False |
Bryan Duxbury | 6972041 | 2012-01-03 17:32:30 +0000 | [diff] [blame] | 407 | |
David Reiss | 7442127 | 2008-11-07 23:09:31 +0000 | [diff] [blame] | 408 | def serve(self): |
Roger Meier | cfff856 | 2012-04-13 14:24:55 +0000 | [diff] [blame] | 409 | """Serve requests. |
| 410 | |
| 411 | Serve requests forever, or until stop() is called. |
| 412 | """ |
| 413 | self._stop = False |
David Reiss | 7442127 | 2008-11-07 23:09:31 +0000 | [diff] [blame] | 414 | self.prepare() |
Roger Meier | cfff856 | 2012-04-13 14:24:55 +0000 | [diff] [blame] | 415 | while not self._stop: |
David Reiss | 7442127 | 2008-11-07 23:09:31 +0000 | [diff] [blame] | 416 | self.handle() |