event_processor.py 4.7 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156
  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 Disk
  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 = Disk()
  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. data = cls.message['message']['data']
  40. if action == 'create_guest':
  41. if state == ResponseState.success.value:
  42. cls.disk.uuid = uuid
  43. cls.disk.get_by('uuid')
  44. cls.disk.guest_uuid = uuid
  45. cls.disk.state = DiskState.mounted.value
  46. # disk_info['virtual-size'] 的单位为Byte,需要除以 1024 的 3 次方,换算成单位为 GB 的值
  47. cls.disk.size = data['disk_info']['virtual-size'] / (1024 ** 3)
  48. cls.disk.update()
  49. else:
  50. cls.guest.uuid = uuid
  51. cls.guest.get_by('uuid')
  52. cls.guest.status = GuestState.dirty.value
  53. cls.guest.update()
  54. elif action == 'migrate':
  55. pass
  56. elif action == 'delete_guest':
  57. if state == ResponseState.success.value:
  58. cls.guest.uuid = uuid
  59. cls.guest.get_by('uuid')
  60. cls.guest.delete()
  61. cls.disk.uuid = uuid
  62. cls.disk.get_by('uuid')
  63. cls.disk.delete()
  64. elif action == 'create_disk':
  65. cls.disk.uuid = uuid
  66. cls.disk.get_by('uuid')
  67. if state == ResponseState.success.value:
  68. cls.disk.state = DiskState.idle.value
  69. else:
  70. cls.disk.state = DiskState.dirty.value
  71. cls.disk.update()
  72. elif action == 'resize_disk':
  73. if state == ResponseState.success.value:
  74. cls.disk.uuid = uuid
  75. cls.disk.get_by('uuid')
  76. cls.disk.size = cls.message['message']['passback_parameters']['size']
  77. cls.disk.update()
  78. elif action == 'attach_disk':
  79. cls.disk.uuid = cls.message['message']['passback_parameters']['disk_uuid']
  80. cls.disk.get_by('uuid')
  81. if state == ResponseState.success.value:
  82. cls.disk.guest_uuid = uuid
  83. cls.disk.sequence = cls.message['message']['passback_parameters']['sequence']
  84. cls.disk.state = DiskState.mounted.value
  85. cls.disk.update()
  86. elif action == 'detach_disk':
  87. cls.disk.uuid = cls.message['message']['passback_parameters']['disk_uuid']
  88. cls.disk.get_by('uuid')
  89. if state == ResponseState.success.value:
  90. cls.disk.guest_uuid = ''
  91. cls.disk.sequence = -1
  92. cls.disk.state = DiskState.idle.value
  93. cls.disk.update()
  94. elif action == 'delete_disk':
  95. cls.disk.uuid = uuid
  96. cls.disk.get_by('uuid')
  97. cls.disk.delete()
  98. else:
  99. pass
  100. @classmethod
  101. def launch(cls):
  102. while True:
  103. if Utils.exit_flag:
  104. Utils.thread_counter -= 1
  105. print 'Thread EventProcessor say bye-bye'
  106. return
  107. try:
  108. report = db.r.lpop(app.config['upstream_queue'])
  109. if report is None:
  110. time.sleep(1)
  111. continue
  112. cls.message = json.loads(report)
  113. if cls.message['kind'] == EmitKind.log.value:
  114. cls.log_processor()
  115. elif cls.message['kind'] == EmitKind.event.value:
  116. cls.event_processor()
  117. elif cls.message['kind'] == EmitKind.response.value:
  118. cls.response_processor()
  119. else:
  120. pass
  121. except Exception as e:
  122. logger.error(e.message)