event_processor.py 5.9 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184
  1. #!/usr/bin/env python
  2. # -*- coding: utf-8 -*-
  3. import json
  4. import time
  5. from IPy import IP
  6. from models import Database as db, Config
  7. from models import Guest
  8. from models import Disk
  9. from models import Log
  10. from models import Utils
  11. from models import EmitKind
  12. from models import ResponseState, GuestState, DiskState
  13. from models.initialize import app, logger
  14. __author__ = 'James Iter'
  15. __date__ = '2017/4/15'
  16. __contact__ = 'james.iter.cn@gmail.com'
  17. __copyright__ = '(c) 2017 by James Iter.'
  18. class EventProcessor(object):
  19. message = None
  20. log = Log()
  21. guest = Guest()
  22. disk = Disk()
  23. config = Config()
  24. @classmethod
  25. def log_processor(cls):
  26. cls.log.set(type=cls.message['type'], timestamp=cls.message['timestamp'], host=cls.message['host'],
  27. message=cls.message['message'])
  28. cls.log.create()
  29. @classmethod
  30. def guest_event_processor(cls):
  31. cls.guest.uuid = cls.message['message']['uuid']
  32. cls.guest.get_by('uuid')
  33. cls.guest.status = cls.message['type']
  34. cls.guest.on_host = cls.message['host']
  35. cls.guest.update()
  36. @classmethod
  37. def host_event_processor(cls):
  38. key = cls.message['message']['node_id']
  39. value = {
  40. 'hostname': cls.message['host'],
  41. 'timestamp': cls.message['timestamp']
  42. }
  43. db.r.hset(app.config['hosts_info'], key=key, value=json.dumps(value, ensure_ascii=False))
  44. db.r.expire(app.config['hosts_info'], 10)
  45. @classmethod
  46. def response_processor(cls):
  47. action = cls.message['message']['action']
  48. uuid = cls.message['message']['uuid']
  49. state = cls.message['type']
  50. data = cls.message['message']['data']
  51. if action == 'create_guest':
  52. if state == ResponseState.success.value:
  53. cls.disk.uuid = uuid
  54. cls.disk.get_by('uuid')
  55. cls.disk.guest_uuid = uuid
  56. cls.disk.state = DiskState.mounted.value
  57. # disk_info['virtual-size'] 的单位为Byte,需要除以 1024 的 3 次方,换算成单位为 GB 的值
  58. cls.disk.size = data['disk_info']['virtual-size'] / (1024 ** 3)
  59. cls.disk.update()
  60. else:
  61. cls.guest.uuid = uuid
  62. cls.guest.get_by('uuid')
  63. cls.guest.status = GuestState.dirty.value
  64. cls.guest.update()
  65. elif action == 'migrate':
  66. pass
  67. elif action == 'delete_guest':
  68. if state == ResponseState.success.value:
  69. cls.config.id = 1
  70. cls.config.get()
  71. cls.guest.uuid = uuid
  72. cls.guest.get_by('uuid')
  73. if IP(cls.config.start_ip).int() <= IP(cls.guest.ip).int() <= IP(cls.config.end_ip).int():
  74. if db.r.srem(app.config['ip_used_set'], cls.guest.ip):
  75. db.r.sadd(app.config['ip_available_set'], cls.guest.ip)
  76. if (cls.guest.vnc_port - cls.config.start_vnc_port) <= \
  77. (IP(cls.config.end_ip).int() - IP(cls.config.start_ip).int()):
  78. if db.r.srem(app.config['vnc_port_used_set'], cls.guest.vnc_port):
  79. db.r.sadd(app.config['vnc_port_available_set'], cls.guest.vnc_port)
  80. cls.guest.delete()
  81. cls.disk.uuid = uuid
  82. cls.disk.get_by('uuid')
  83. cls.disk.delete()
  84. elif action == 'create_disk':
  85. cls.disk.uuid = uuid
  86. cls.disk.get_by('uuid')
  87. if state == ResponseState.success.value:
  88. cls.disk.state = DiskState.idle.value
  89. else:
  90. cls.disk.state = DiskState.dirty.value
  91. cls.disk.update()
  92. elif action == 'resize_disk':
  93. if state == ResponseState.success.value:
  94. cls.disk.uuid = uuid
  95. cls.disk.get_by('uuid')
  96. cls.disk.size = cls.message['message']['passback_parameters']['size']
  97. cls.disk.update()
  98. elif action == 'attach_disk':
  99. cls.disk.uuid = cls.message['message']['passback_parameters']['disk_uuid']
  100. cls.disk.get_by('uuid')
  101. if state == ResponseState.success.value:
  102. cls.disk.guest_uuid = uuid
  103. cls.disk.sequence = cls.message['message']['passback_parameters']['sequence']
  104. cls.disk.state = DiskState.mounted.value
  105. cls.disk.update()
  106. elif action == 'detach_disk':
  107. cls.disk.uuid = cls.message['message']['passback_parameters']['disk_uuid']
  108. cls.disk.get_by('uuid')
  109. if state == ResponseState.success.value:
  110. cls.disk.guest_uuid = ''
  111. cls.disk.sequence = -1
  112. cls.disk.state = DiskState.idle.value
  113. cls.disk.update()
  114. elif action == 'delete_disk':
  115. cls.disk.uuid = uuid
  116. cls.disk.get_by('uuid')
  117. cls.disk.delete()
  118. else:
  119. pass
  120. @classmethod
  121. def launch(cls):
  122. while True:
  123. if Utils.exit_flag:
  124. Utils.thread_counter -= 1
  125. print 'Thread EventProcessor say bye-bye'
  126. return
  127. try:
  128. report = db.r.lpop(app.config['upstream_queue'])
  129. if report is None:
  130. time.sleep(1)
  131. continue
  132. cls.message = json.loads(report)
  133. if cls.message['kind'] == EmitKind.log.value:
  134. cls.log_processor()
  135. elif cls.message['kind'] == EmitKind.guest_event.value:
  136. cls.guest_event_processor()
  137. elif cls.message['kind'] == EmitKind.host_event.value:
  138. cls.host_event_processor()
  139. elif cls.message['kind'] == EmitKind.response.value:
  140. cls.response_processor()
  141. else:
  142. pass
  143. except Exception as e:
  144. logger.error(e.message)