scheduler.py 2.9 KB

1234567891011121314151617181920212223242526272829303132333435363738394041424344454647484950515253545556575859606162636465666768697071727374757677787980818283848586878889909192
  1. import time
  2. import multiprocessing
  3. from proxypool.processors.server import app
  4. from proxypool.processors.getter import Getter
  5. from proxypool.processors.tester import Tester
  6. from proxypool.setting import CYCLE_GETTER, CYCLE_TESTER, API_HOST, API_THREADED, API_PORT, ENABLE_SERVER, \
  7. ENABLE_GETTER, ENABLE_TESTER, IS_WINDOWS
  8. from loguru import logger
  9. if IS_WINDOWS:
  10. multiprocessing.freeze_support()
  11. tester_process, getter_process, server_process = None, None, None
  12. class Scheduler():
  13. """
  14. scheduler
  15. """
  16. def run_tester(self, cycle=CYCLE_TESTER):
  17. """
  18. run tester
  19. """
  20. tester = Tester()
  21. loop = 0
  22. while True:
  23. logger.debug(f'tester loop {loop} start...')
  24. tester.run()
  25. loop += 1
  26. time.sleep(cycle)
  27. def run_getter(self, cycle=CYCLE_GETTER):
  28. """
  29. run getter
  30. """
  31. getter = Getter()
  32. loop = 0
  33. while True:
  34. logger.debug(f'getter loop {loop} start...')
  35. getter.run()
  36. loop += 1
  37. time.sleep(cycle)
  38. def run_server(self):
  39. """
  40. run server for api
  41. """
  42. app.run(host=API_HOST, port=API_PORT, threaded=API_THREADED)
  43. def run(self):
  44. global tester_process, getter_process, server_process
  45. try:
  46. logger.info('starting proxypool...')
  47. if ENABLE_TESTER:
  48. tester_process = multiprocessing.Process(target=self.run_tester)
  49. logger.info(f'starting tester, pid {tester_process.pid}...')
  50. tester_process.start()
  51. if ENABLE_GETTER:
  52. getter_process = multiprocessing.Process(target=self.run_getter)
  53. logger.info(f'starting getter, pid{getter_process.pid}...')
  54. getter_process.start()
  55. if ENABLE_SERVER:
  56. server_process = multiprocessing.Process(target=self.run_server)
  57. logger.info(f'starting server, pid{server_process.pid}...')
  58. server_process.start()
  59. tester_process.join()
  60. getter_process.join()
  61. server_process.join()
  62. except KeyboardInterrupt:
  63. logger.info('received keyboard interrupt signal')
  64. tester_process.terminate()
  65. getter_process.terminate()
  66. server_process.terminate()
  67. finally:
  68. # must call join method before calling is_alive
  69. tester_process.join()
  70. getter_process.join()
  71. server_process.join()
  72. logger.info(f'tester is {"alive" if tester_process.is_alive() else "dead"}')
  73. logger.info(f'getter is {"alive" if getter_process.is_alive() else "dead"}')
  74. logger.info(f'server is {"alive" if server_process.is_alive() else "dead"}')
  75. logger.info('proxy terminated')
  76. if __name__ == '__main__':
  77. scheduler = Scheduler()
  78. scheduler.run()