event_processor.py 9.2 KB

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