server.py 3.0 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105
  1. # -*- coding: utf-8 -*-
  2. from __future__ import absolute_import
  3. import logging
  4. import threading
  5. from thriftpy.protocol import TBinaryProtocolFactory
  6. from thriftpy.transport import (
  7. TBufferedTransportFactory,
  8. TTransportException
  9. )
  10. logger = logging.getLogger(__name__)
  11. class TServer(object):
  12. def __init__(self, processor, trans,
  13. itrans_factory=None, iprot_factory=None,
  14. otrans_factory=None, oprot_factory=None):
  15. self.processor = processor
  16. self.trans = trans
  17. self.itrans_factory = itrans_factory or TBufferedTransportFactory()
  18. self.iprot_factory = iprot_factory or TBinaryProtocolFactory()
  19. self.otrans_factory = otrans_factory or self.itrans_factory
  20. self.oprot_factory = oprot_factory or self.iprot_factory
  21. def serve(self):
  22. pass
  23. def close(self):
  24. pass
  25. class TSimpleServer(TServer):
  26. """Simple single-threaded server that just pumps around one transport."""
  27. def __init__(self, *args):
  28. TServer.__init__(self, *args)
  29. self.closed = False
  30. def serve(self):
  31. self.trans.listen()
  32. while True:
  33. client = self.trans.accept()
  34. itrans = self.itrans_factory.get_transport(client)
  35. otrans = self.otrans_factory.get_transport(client)
  36. iprot = self.iprot_factory.get_protocol(itrans)
  37. oprot = self.oprot_factory.get_protocol(otrans)
  38. try:
  39. while not self.closed:
  40. self.processor.process(iprot, oprot)
  41. except TTransportException:
  42. pass
  43. except Exception as x:
  44. logger.exception(x)
  45. itrans.close()
  46. otrans.close()
  47. def close(self):
  48. self.closed = True
  49. class TThreadedServer(TServer):
  50. """Threaded server that spawns a new thread per each connection."""
  51. def __init__(self, *args, **kwargs):
  52. self.daemon = kwargs.pop("daemon", False)
  53. TServer.__init__(self, *args, **kwargs)
  54. self.closed = False
  55. def serve(self):
  56. self.trans.listen()
  57. while not self.closed:
  58. try:
  59. client = self.trans.accept()
  60. t = threading.Thread(target=self.handle, args=(client,))
  61. t.setDaemon(self.daemon)
  62. t.start()
  63. except KeyboardInterrupt:
  64. raise
  65. except Exception as x:
  66. logger.exception(x)
  67. def handle(self, client):
  68. itrans = self.itrans_factory.get_transport(client)
  69. otrans = self.otrans_factory.get_transport(client)
  70. iprot = self.iprot_factory.get_protocol(itrans)
  71. oprot = self.oprot_factory.get_protocol(otrans)
  72. try:
  73. while True:
  74. self.processor.process(iprot, oprot)
  75. except TTransportException:
  76. pass
  77. except Exception as x:
  78. logger.exception(x)
  79. itrans.close()
  80. otrans.close()
  81. def close(self):
  82. self.closed = True