+ with self.lk:
+ self.current.remove(th)
+ self.tcond.notify_all()
+ finally:
+ req.close()
+
+ def close(self):
+ while True:
+ with self.lk:
+ if len(self.current) > 0:
+ th = next(iter(self.current))
+ else:
+ return
+ th.join()
+
+class resplex(handler):
+ cname = "rplex"
+
+ def __init__(self, *, max=None, **kw):
+ super().__init__(**kw)
+ self.current = set()
+ self.lk = threading.Lock()
+ self.tcond = threading.Condition(self.lk)
+ self.max = max
+ self.cqueue = queue.Queue(5)
+ self.cnpipe = os.pipe()
+ self.rthread = reqthread(name="Response thread", target=self.handle2)
+ self.rthread.start()
+
+ @classmethod
+ def parseargs(cls, *, max=None, **args):
+ ret = super().parseargs(**args)
+ if max:
+ ret["max"] = int(max)
+ return ret
+
+ def ckflush(self, req):
+ raise Exception("resplex handler does not support the write() function")
+
+ def handle(self, req):
+ with self.lk:
+ if self.max is not None:
+ while len(self.current) >= self.max:
+ self.tcond.wait()
+ th = reqthread(target=self.handle1, args=[req])
+ th.start()
+ while th.is_alive() and th not in self.current:
+ self.tcond.wait(1)
+
+ def handle1(self, req):
+ try:
+ th = threading.current_thread()
+ with self.lk:
+ self.current.add(th)
+ self.tcond.notify_all()
+ try:
+ env = req.mkenv()
+ respobj = req.handlewsgi(env, req.startreq)
+ respiter = iter(respobj)
+ if not req.status:
+ log.error("request handler returned without calling start_request")
+ if hasattr(respiter, "close"):
+ respiter.close()
+ return
+ else:
+ self.cqueue.put((req, respiter))
+ os.write(self.cnpipe[1], b" ")
+ req = None
+ finally:
+ with self.lk:
+ self.current.remove(th)
+ self.tcond.notify_all()