event_processor.py 2.7 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103
  1. #!/usr/bin/env python
  2. # -*- coding: utf-8 -*-
  3. import json
  4. import time
  5. from models import Database as db
  6. from models import Guest
  7. from models import GuestDisk
  8. from models import Log
  9. from models import Utils
  10. from models import EmitKind
  11. from models import ResponseState, GuestState, DiskState
  12. from models.initialize import app, logger
  13. __author__ = 'James Iter'
  14. __date__ = '2017/4/15'
  15. __contact__ = 'james.iter.cn@gmail.com'
  16. __copyright__ = '(c) 2017 by James Iter.'
  17. class EventProcessor(object):
  18. message = None
  19. log = Log()
  20. guest = Guest()
  21. disk = GuestDisk()
  22. @classmethod
  23. def log_processor(cls):
  24. cls.log.set(type=cls.message['type'], timestamp=cls.message['timestamp'], host=cls.message['host'],
  25. message=cls.message['message'])
  26. cls.log.create()
  27. @classmethod
  28. def event_processor(cls):
  29. cls.guest.uuid = cls.message['message']['uuid']
  30. cls.guest.get_by('uuid')
  31. cls.guest.status = cls.message['type']
  32. cls.guest.on_host = cls.message['host']
  33. cls.guest.update()
  34. @classmethod
  35. def response_processor(cls):
  36. action = cls.message['message']['action']
  37. uuid = cls.message['message']['uuid']
  38. state = cls.message['type']
  39. if action == 'create_vm':
  40. if state != ResponseState.success.value:
  41. cls.guest.uuid = uuid
  42. cls.guest.get_by('uuid')
  43. cls.guest.status = GuestState.dirty.value
  44. cls.guest.update()
  45. elif action == 'create_disk':
  46. cls.disk.uuid = uuid
  47. cls.disk.get_by('uuid')
  48. if state == ResponseState.success.value:
  49. cls.disk.state = DiskState.idle.value
  50. else:
  51. cls.disk.state = DiskState.dirty.value
  52. cls.disk.update()
  53. else:
  54. pass
  55. @classmethod
  56. def launch(cls):
  57. while True:
  58. if Utils.exit_flag:
  59. Utils.thread_counter -= 1
  60. print 'Thread EventProcessor say bye-bye'
  61. return
  62. try:
  63. report = db.r.lpop(app.config['upstream_queue'])
  64. if report is None:
  65. time.sleep(1)
  66. continue
  67. cls.message = json.loads(report)
  68. if cls.message['kind'] == EmitKind.log.value:
  69. cls.log_processor()
  70. elif cls.message['kind'] == EmitKind.event.value:
  71. cls.event_processor()
  72. elif cls.message['kind'] == EmitKind.response.value:
  73. cls.response_processor()
  74. else:
  75. pass
  76. except Exception as e:
  77. logger.error(e.message)