Bryan Duxbury | a48b7d6 | 2011-03-09 18:05:58 +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 | # |
| 19 | |
| 20 | |
| 21 | import logging |
Konrad Grochowski | 3a724e3 | 2014-08-12 11:48:29 -0400 | [diff] [blame] | 22 | |
Nobuaki Sukegawa | 10308cb | 2016-02-03 01:57:03 +0900 | [diff] [blame] | 23 | from multiprocessing import Process, Value, Condition |
Bryan Duxbury | a48b7d6 | 2011-03-09 18:05:58 +0000 | [diff] [blame] | 24 | |
Nobuaki Sukegawa | 760511f | 2015-11-06 21:24:16 +0900 | [diff] [blame] | 25 | from .TServer import TServer |
Bryan Duxbury | a48b7d6 | 2011-03-09 18:05:58 +0000 | [diff] [blame] | 26 | from thrift.transport.TTransport import TTransportException |
| 27 | |
Nobuaki Sukegawa | 10308cb | 2016-02-03 01:57:03 +0900 | [diff] [blame] | 28 | logger = logging.getLogger(__name__) |
| 29 | |
Bryan Duxbury | a48b7d6 | 2011-03-09 18:05:58 +0000 | [diff] [blame] | 30 | |
Bryan Duxbury | 6972041 | 2012-01-03 17:32:30 +0000 | [diff] [blame] | 31 | class TProcessPoolServer(TServer): |
| 32 | """Server with a fixed size pool of worker subprocesses to service requests |
| 33 | |
Bryan Duxbury | a48b7d6 | 2011-03-09 18:05:58 +0000 | [diff] [blame] | 34 | Note that if you need shared state between the handlers - it's up to you! |
| 35 | Written by Dvir Volk, doat.com |
| 36 | """ |
Bryan Duxbury | 6972041 | 2012-01-03 17:32:30 +0000 | [diff] [blame] | 37 | def __init__(self, *args): |
Bryan Duxbury | a48b7d6 | 2011-03-09 18:05:58 +0000 | [diff] [blame] | 38 | TServer.__init__(self, *args) |
| 39 | self.numWorkers = 10 |
| 40 | self.workers = [] |
| 41 | self.isRunning = Value('b', False) |
| 42 | self.stopCondition = Condition() |
| 43 | self.postForkCallback = None |
| 44 | |
Yiyang Zhou | da1e19b | 2021-08-13 16:30:58 +0800 | [diff] [blame] | 45 | def __getstate__(self): |
Jens Geyer | 443a03c | 2021-11-15 19:30:13 +0100 | [diff] [blame] | 46 | state = self.__dict__.copy() |
Yiyang Zhou | da1e19b | 2021-08-13 16:30:58 +0800 | [diff] [blame] | 47 | state['workers'] = None |
| 48 | return state |
| 49 | |
Bryan Duxbury | a48b7d6 | 2011-03-09 18:05:58 +0000 | [diff] [blame] | 50 | def setPostForkCallback(self, callback): |
| 51 | if not callable(callback): |
| 52 | raise TypeError("This is not a callback!") |
| 53 | self.postForkCallback = callback |
| 54 | |
| 55 | def setNumWorkers(self, num): |
| 56 | """Set the number of worker threads that should be created""" |
| 57 | self.numWorkers = num |
| 58 | |
| 59 | def workerProcess(self): |
Bryan Duxbury | 6972041 | 2012-01-03 17:32:30 +0000 | [diff] [blame] | 60 | """Loop getting clients from the shared queue and process them""" |
Bryan Duxbury | a48b7d6 | 2011-03-09 18:05:58 +0000 | [diff] [blame] | 61 | if self.postForkCallback: |
| 62 | self.postForkCallback() |
| 63 | |
Bryan Duxbury | 6972041 | 2012-01-03 17:32:30 +0000 | [diff] [blame] | 64 | while self.isRunning.value: |
Bryan Duxbury | a48b7d6 | 2011-03-09 18:05:58 +0000 | [diff] [blame] | 65 | try: |
| 66 | client = self.serverTransport.accept() |
Roger Meier | ab2793a | 2014-04-21 21:20:00 +0200 | [diff] [blame] | 67 | if not client: |
Nobuaki Sukegawa | 10308cb | 2016-02-03 01:57:03 +0900 | [diff] [blame] | 68 | continue |
Bryan Duxbury | a48b7d6 | 2011-03-09 18:05:58 +0000 | [diff] [blame] | 69 | self.serveClient(client) |
| 70 | except (KeyboardInterrupt, SystemExit): |
| 71 | return 0 |
jfarrell | d565e2f | 2015-03-18 21:02:47 -0400 | [diff] [blame] | 72 | except Exception as x: |
Konrad Grochowski | 3a724e3 | 2014-08-12 11:48:29 -0400 | [diff] [blame] | 73 | logger.exception(x) |
Bryan Duxbury | a48b7d6 | 2011-03-09 18:05:58 +0000 | [diff] [blame] | 74 | |
| 75 | def serveClient(self, client): |
| 76 | """Process input/output from a client for as long as possible""" |
| 77 | itrans = self.inputTransportFactory.getTransport(client) |
| 78 | otrans = self.outputTransportFactory.getTransport(client) |
| 79 | iprot = self.inputProtocolFactory.getProtocol(itrans) |
| 80 | oprot = self.outputProtocolFactory.getProtocol(otrans) |
| 81 | |
| 82 | try: |
| 83 | while True: |
| 84 | self.processor.process(iprot, oprot) |
Nobuaki Sukegawa | 10308cb | 2016-02-03 01:57:03 +0900 | [diff] [blame] | 85 | except TTransportException: |
Bryan Duxbury | a48b7d6 | 2011-03-09 18:05:58 +0000 | [diff] [blame] | 86 | pass |
jfarrell | d565e2f | 2015-03-18 21:02:47 -0400 | [diff] [blame] | 87 | except Exception as x: |
Konrad Grochowski | 3a724e3 | 2014-08-12 11:48:29 -0400 | [diff] [blame] | 88 | logger.exception(x) |
Bryan Duxbury | a48b7d6 | 2011-03-09 18:05:58 +0000 | [diff] [blame] | 89 | |
| 90 | itrans.close() |
| 91 | otrans.close() |
| 92 | |
Bryan Duxbury | a48b7d6 | 2011-03-09 18:05:58 +0000 | [diff] [blame] | 93 | def serve(self): |
Bryan Duxbury | 6972041 | 2012-01-03 17:32:30 +0000 | [diff] [blame] | 94 | """Start workers and put into queue""" |
| 95 | # this is a shared state that can tell the workers to exit when False |
Bryan Duxbury | a48b7d6 | 2011-03-09 18:05:58 +0000 | [diff] [blame] | 96 | self.isRunning.value = True |
| 97 | |
Bryan Duxbury | 6972041 | 2012-01-03 17:32:30 +0000 | [diff] [blame] | 98 | # first bind and listen to the port |
Bryan Duxbury | a48b7d6 | 2011-03-09 18:05:58 +0000 | [diff] [blame] | 99 | self.serverTransport.listen() |
| 100 | |
Bryan Duxbury | 6972041 | 2012-01-03 17:32:30 +0000 | [diff] [blame] | 101 | # fork the children |
Bryan Duxbury | a48b7d6 | 2011-03-09 18:05:58 +0000 | [diff] [blame] | 102 | for i in range(self.numWorkers): |
| 103 | try: |
| 104 | w = Process(target=self.workerProcess) |
| 105 | w.daemon = True |
| 106 | w.start() |
| 107 | self.workers.append(w) |
jfarrell | d565e2f | 2015-03-18 21:02:47 -0400 | [diff] [blame] | 108 | except Exception as x: |
Konrad Grochowski | 3a724e3 | 2014-08-12 11:48:29 -0400 | [diff] [blame] | 109 | logger.exception(x) |
Bryan Duxbury | a48b7d6 | 2011-03-09 18:05:58 +0000 | [diff] [blame] | 110 | |
Bryan Duxbury | 6972041 | 2012-01-03 17:32:30 +0000 | [diff] [blame] | 111 | # wait until the condition is set by stop() |
Bryan Duxbury | a48b7d6 | 2011-03-09 18:05:58 +0000 | [diff] [blame] | 112 | while True: |
Bryan Duxbury | a48b7d6 | 2011-03-09 18:05:58 +0000 | [diff] [blame] | 113 | self.stopCondition.acquire() |
| 114 | try: |
| 115 | self.stopCondition.wait() |
| 116 | break |
| 117 | except (SystemExit, KeyboardInterrupt): |
Bryan Duxbury | 6972041 | 2012-01-03 17:32:30 +0000 | [diff] [blame] | 118 | break |
jfarrell | d565e2f | 2015-03-18 21:02:47 -0400 | [diff] [blame] | 119 | except Exception as x: |
Konrad Grochowski | 3a724e3 | 2014-08-12 11:48:29 -0400 | [diff] [blame] | 120 | logger.exception(x) |
Bryan Duxbury | a48b7d6 | 2011-03-09 18:05:58 +0000 | [diff] [blame] | 121 | |
| 122 | self.isRunning.value = False |
| 123 | |
| 124 | def stop(self): |
| 125 | self.isRunning.value = False |
| 126 | self.stopCondition.acquire() |
| 127 | self.stopCondition.notify() |
| 128 | self.stopCondition.release() |