event_processor.py 9.7 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240
  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. last_status = cls.guest.status
  38. cls.guest.status = cls.message['type']
  39. if cls.message['type'] == GuestState.update.value:
  40. # 更新事件不改变 Guest 的状态
  41. cls.guest.status = last_status
  42. cls.guest.xml = cls.message['message']['xml']
  43. elif cls.guest.status == GuestState.migrating.value:
  44. try:
  45. cls.guest_migrate_info.uuid = cls.guest.uuid
  46. cls.guest_migrate_info.get_by('uuid')
  47. cls.guest_migrate_info.type = cls.message['message']['migrating_info']['type']
  48. cls.guest_migrate_info.time_elapsed = cls.message['message']['migrating_info']['time_elapsed']
  49. cls.guest_migrate_info.time_remaining = cls.message['message']['migrating_info']['time_remaining']
  50. cls.guest_migrate_info.data_total = cls.message['message']['migrating_info']['data_total']
  51. cls.guest_migrate_info.data_processed = cls.message['message']['migrating_info']['data_processed']
  52. cls.guest_migrate_info.data_remaining = cls.message['message']['migrating_info']['data_remaining']
  53. cls.guest_migrate_info.mem_total = cls.message['message']['migrating_info']['mem_total']
  54. cls.guest_migrate_info.mem_processed = cls.message['message']['migrating_info']['mem_processed']
  55. cls.guest_migrate_info.mem_remaining = cls.message['message']['migrating_info']['mem_remaining']
  56. cls.guest_migrate_info.file_total = cls.message['message']['migrating_info']['file_total']
  57. cls.guest_migrate_info.file_processed = cls.message['message']['migrating_info']['file_processed']
  58. cls.guest_migrate_info.file_remaining = cls.message['message']['migrating_info']['file_remaining']
  59. cls.guest_migrate_info.update()
  60. except ji.PreviewingError as e:
  61. ret = json.loads(e.message)
  62. if ret['state']['code'] == '404':
  63. cls.guest_migrate_info.type = cls.message['message']['migrating_info']['type']
  64. cls.guest_migrate_info.time_elapsed = cls.message['message']['migrating_info']['time_elapsed']
  65. cls.guest_migrate_info.time_remaining = cls.message['message']['migrating_info']['time_remaining']
  66. cls.guest_migrate_info.data_total = cls.message['message']['migrating_info']['data_total']
  67. cls.guest_migrate_info.data_processed = cls.message['message']['migrating_info']['data_processed']
  68. cls.guest_migrate_info.data_remaining = cls.message['message']['migrating_info']['data_remaining']
  69. cls.guest_migrate_info.mem_total = cls.message['message']['migrating_info']['mem_total']
  70. cls.guest_migrate_info.mem_processed = cls.message['message']['migrating_info']['mem_processed']
  71. cls.guest_migrate_info.mem_remaining = cls.message['message']['migrating_info']['mem_remaining']
  72. cls.guest_migrate_info.file_total = cls.message['message']['migrating_info']['file_total']
  73. cls.guest_migrate_info.file_processed = cls.message['message']['migrating_info']['file_processed']
  74. cls.guest_migrate_info.file_remaining = cls.message['message']['migrating_info']['file_remaining']
  75. cls.guest_migrate_info.create()
  76. cls.guest.update()
  77. @classmethod
  78. def host_event_processor(cls):
  79. key = cls.message['message']['node_id']
  80. value = {
  81. 'hostname': cls.message['host'],
  82. 'timestamp': ji.Common.ts()
  83. }
  84. db.r.hset(app.config['hosts_info'], key=key, value=json.dumps(value, ensure_ascii=False))
  85. db.r.expire(app.config['hosts_info'], 10)
  86. @classmethod
  87. def response_processor(cls):
  88. action = cls.message['message']['action']
  89. uuid = cls.message['message']['uuid']
  90. state = cls.message['type']
  91. data = cls.message['message']['data']
  92. if action == 'create_guest':
  93. if state == ResponseState.success.value:
  94. # 系统盘的 UUID 与其 Guest 的 UUID 相同
  95. cls.disk.uuid = uuid
  96. cls.disk.get_by('uuid')
  97. cls.disk.guest_uuid = uuid
  98. cls.disk.state = DiskState.mounted.value
  99. # disk_info['virtual-size'] 的单位为Byte,需要除以 1024 的 3 次方,换算成单位为 GB 的值
  100. cls.disk.size = data['disk_info']['virtual-size'] / (1024 ** 3)
  101. cls.disk.update()
  102. else:
  103. cls.guest.uuid = uuid
  104. cls.guest.get_by('uuid')
  105. cls.guest.status = GuestState.dirty.value
  106. cls.guest.update()
  107. elif action == 'migrate':
  108. pass
  109. elif action == 'delete_guest':
  110. if state == ResponseState.success.value:
  111. cls.config.id = 1
  112. cls.config.get()
  113. cls.guest.uuid = uuid
  114. cls.guest.get_by('uuid')
  115. if IP(cls.config.start_ip).int() <= IP(cls.guest.ip).int() <= IP(cls.config.end_ip).int():
  116. if db.r.srem(app.config['ip_used_set'], cls.guest.ip):
  117. db.r.sadd(app.config['ip_available_set'], cls.guest.ip)
  118. if (cls.guest.vnc_port - cls.config.start_vnc_port) <= \
  119. (IP(cls.config.end_ip).int() - IP(cls.config.start_ip).int()):
  120. if db.r.srem(app.config['vnc_port_used_set'], cls.guest.vnc_port):
  121. db.r.sadd(app.config['vnc_port_available_set'], cls.guest.vnc_port)
  122. cls.guest.delete()
  123. cls.disk.uuid = uuid
  124. cls.disk.get_by('uuid')
  125. cls.disk.delete()
  126. elif action == 'create_disk':
  127. cls.disk.uuid = uuid
  128. cls.disk.get_by('uuid')
  129. if state == ResponseState.success.value:
  130. cls.disk.state = DiskState.idle.value
  131. else:
  132. cls.disk.state = DiskState.dirty.value
  133. cls.disk.update()
  134. elif action == 'resize_disk':
  135. if state == ResponseState.success.value:
  136. cls.disk.uuid = uuid
  137. cls.disk.get_by('uuid')
  138. cls.disk.size = cls.message['message']['passback_parameters']['size']
  139. cls.disk.update()
  140. elif action == 'attach_disk':
  141. cls.disk.uuid = cls.message['message']['passback_parameters']['disk_uuid']
  142. cls.disk.get_by('uuid')
  143. if state == ResponseState.success.value:
  144. cls.disk.guest_uuid = uuid
  145. cls.disk.sequence = cls.message['message']['passback_parameters']['sequence']
  146. cls.disk.state = DiskState.mounted.value
  147. cls.disk.update()
  148. elif action == 'detach_disk':
  149. cls.disk.uuid = cls.message['message']['passback_parameters']['disk_uuid']
  150. cls.disk.get_by('uuid')
  151. if state == ResponseState.success.value:
  152. cls.disk.guest_uuid = ''
  153. cls.disk.sequence = -1
  154. cls.disk.state = DiskState.idle.value
  155. cls.disk.update()
  156. elif action == 'delete_disk':
  157. cls.disk.uuid = uuid
  158. cls.disk.get_by('uuid')
  159. cls.disk.delete()
  160. elif action == 'boot':
  161. boot_jobs_id = cls.message['message']['passback_parameters']['boot_jobs_id']
  162. if state == ResponseState.success.value:
  163. cls.guest.uuid = uuid
  164. cls.guest.delete_boot_jobs(boot_jobs_id=boot_jobs_id)
  165. else:
  166. pass
  167. @classmethod
  168. def launch(cls):
  169. while True:
  170. if Utils.exit_flag:
  171. Utils.thread_counter -= 1
  172. print 'Thread EventProcessor say bye-bye'
  173. return
  174. try:
  175. report = db.r.lpop(app.config['upstream_queue'])
  176. if report is None:
  177. time.sleep(1)
  178. continue
  179. cls.message = json.loads(report)
  180. if cls.message['kind'] == EmitKind.log.value:
  181. cls.log_processor()
  182. elif cls.message['kind'] == EmitKind.guest_event.value:
  183. cls.guest_event_processor()
  184. elif cls.message['kind'] == EmitKind.host_event.value:
  185. cls.host_event_processor()
  186. elif cls.message['kind'] == EmitKind.response.value:
  187. cls.response_processor()
  188. else:
  189. pass
  190. except Exception as e:
  191. logger.error(e.message)