event_processor.py 9.3 KB

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