event_processor.py 12 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290
  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, CPUMemory, Traffic, DiskIO
  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. from models.status import CollectionPerformanceDataKind
  17. __author__ = 'James Iter'
  18. __date__ = '2017/4/15'
  19. __contact__ = 'james.iter.cn@gmail.com'
  20. __copyright__ = '(c) 2017 by James Iter.'
  21. class EventProcessor(object):
  22. message = None
  23. log = Log()
  24. guest = Guest()
  25. guest_migrate_info = GuestMigrateInfo()
  26. disk = Disk()
  27. config = Config()
  28. cpu_memory = CPUMemory()
  29. traffic = Traffic()
  30. disk_io = DiskIO()
  31. @classmethod
  32. def log_processor(cls):
  33. cls.log.set(type=cls.message['type'], timestamp=cls.message['timestamp'], host=cls.message['host'],
  34. message=cls.message['message'])
  35. cls.log.create()
  36. @classmethod
  37. def guest_event_processor(cls):
  38. cls.guest.uuid = cls.message['message']['uuid']
  39. cls.guest.get_by('uuid')
  40. cls.guest.on_host = cls.message['host']
  41. last_status = cls.guest.status
  42. cls.guest.status = cls.message['type']
  43. if cls.message['type'] == GuestState.update.value:
  44. # 更新事件不改变 Guest 的状态
  45. cls.guest.status = last_status
  46. cls.guest.xml = cls.message['message']['xml']
  47. elif cls.guest.status == GuestState.migrating.value:
  48. try:
  49. cls.guest_migrate_info.uuid = cls.guest.uuid
  50. cls.guest_migrate_info.get_by('uuid')
  51. cls.guest_migrate_info.type = cls.message['message']['migrating_info']['type']
  52. cls.guest_migrate_info.time_elapsed = cls.message['message']['migrating_info']['time_elapsed']
  53. cls.guest_migrate_info.time_remaining = cls.message['message']['migrating_info']['time_remaining']
  54. cls.guest_migrate_info.data_total = cls.message['message']['migrating_info']['data_total']
  55. cls.guest_migrate_info.data_processed = cls.message['message']['migrating_info']['data_processed']
  56. cls.guest_migrate_info.data_remaining = cls.message['message']['migrating_info']['data_remaining']
  57. cls.guest_migrate_info.mem_total = cls.message['message']['migrating_info']['mem_total']
  58. cls.guest_migrate_info.mem_processed = cls.message['message']['migrating_info']['mem_processed']
  59. cls.guest_migrate_info.mem_remaining = cls.message['message']['migrating_info']['mem_remaining']
  60. cls.guest_migrate_info.file_total = cls.message['message']['migrating_info']['file_total']
  61. cls.guest_migrate_info.file_processed = cls.message['message']['migrating_info']['file_processed']
  62. cls.guest_migrate_info.file_remaining = cls.message['message']['migrating_info']['file_remaining']
  63. cls.guest_migrate_info.update()
  64. except ji.PreviewingError as e:
  65. ret = json.loads(e.message)
  66. if ret['state']['code'] == '404':
  67. cls.guest_migrate_info.type = cls.message['message']['migrating_info']['type']
  68. cls.guest_migrate_info.time_elapsed = cls.message['message']['migrating_info']['time_elapsed']
  69. cls.guest_migrate_info.time_remaining = cls.message['message']['migrating_info']['time_remaining']
  70. cls.guest_migrate_info.data_total = cls.message['message']['migrating_info']['data_total']
  71. cls.guest_migrate_info.data_processed = cls.message['message']['migrating_info']['data_processed']
  72. cls.guest_migrate_info.data_remaining = cls.message['message']['migrating_info']['data_remaining']
  73. cls.guest_migrate_info.mem_total = cls.message['message']['migrating_info']['mem_total']
  74. cls.guest_migrate_info.mem_processed = cls.message['message']['migrating_info']['mem_processed']
  75. cls.guest_migrate_info.mem_remaining = cls.message['message']['migrating_info']['mem_remaining']
  76. cls.guest_migrate_info.file_total = cls.message['message']['migrating_info']['file_total']
  77. cls.guest_migrate_info.file_processed = cls.message['message']['migrating_info']['file_processed']
  78. cls.guest_migrate_info.file_remaining = cls.message['message']['migrating_info']['file_remaining']
  79. cls.guest_migrate_info.create()
  80. cls.guest.update()
  81. @classmethod
  82. def host_event_processor(cls):
  83. key = cls.message['message']['node_id']
  84. value = {
  85. 'hostname': cls.message['host'],
  86. 'timestamp': ji.Common.ts()
  87. }
  88. db.r.hset(app.config['hosts_info'], key=key, value=json.dumps(value, ensure_ascii=False))
  89. db.r.expire(app.config['hosts_info'], 10)
  90. @classmethod
  91. def response_processor(cls):
  92. action = cls.message['message']['action']
  93. uuid = cls.message['message']['uuid']
  94. state = cls.message['type']
  95. data = cls.message['message']['data']
  96. if action == 'create_guest':
  97. if state == ResponseState.success.value:
  98. # 系统盘的 UUID 与其 Guest 的 UUID 相同
  99. cls.disk.uuid = uuid
  100. cls.disk.get_by('uuid')
  101. cls.disk.guest_uuid = uuid
  102. cls.disk.state = DiskState.mounted.value
  103. # disk_info['virtual-size'] 的单位为Byte,需要除以 1024 的 3 次方,换算成单位为 GB 的值
  104. cls.disk.size = data['disk_info']['virtual-size'] / (1024 ** 3)
  105. cls.disk.update()
  106. else:
  107. cls.guest.uuid = uuid
  108. cls.guest.get_by('uuid')
  109. cls.guest.status = GuestState.dirty.value
  110. cls.guest.update()
  111. elif action == 'migrate':
  112. pass
  113. elif action == 'delete_guest':
  114. if state == ResponseState.success.value:
  115. cls.config.id = 1
  116. cls.config.get()
  117. cls.guest.uuid = uuid
  118. cls.guest.get_by('uuid')
  119. if IP(cls.config.start_ip).int() <= IP(cls.guest.ip).int() <= IP(cls.config.end_ip).int():
  120. if db.r.srem(app.config['ip_used_set'], cls.guest.ip):
  121. db.r.sadd(app.config['ip_available_set'], cls.guest.ip)
  122. if (cls.guest.vnc_port - cls.config.start_vnc_port) <= \
  123. (IP(cls.config.end_ip).int() - IP(cls.config.start_ip).int()):
  124. if db.r.srem(app.config['vnc_port_used_set'], cls.guest.vnc_port):
  125. db.r.sadd(app.config['vnc_port_available_set'], cls.guest.vnc_port)
  126. cls.guest.delete()
  127. cls.disk.uuid = uuid
  128. cls.disk.get_by('uuid')
  129. cls.disk.delete()
  130. elif action == 'create_disk':
  131. cls.disk.uuid = uuid
  132. cls.disk.get_by('uuid')
  133. if state == ResponseState.success.value:
  134. cls.disk.state = DiskState.idle.value
  135. else:
  136. cls.disk.state = DiskState.dirty.value
  137. cls.disk.update()
  138. elif action == 'resize_disk':
  139. if state == ResponseState.success.value:
  140. cls.disk.uuid = uuid
  141. cls.disk.get_by('uuid')
  142. cls.disk.size = cls.message['message']['passback_parameters']['size']
  143. cls.disk.update()
  144. elif action == 'attach_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 = uuid
  149. cls.disk.sequence = cls.message['message']['passback_parameters']['sequence']
  150. cls.disk.state = DiskState.mounted.value
  151. cls.disk.update()
  152. elif action == 'detach_disk':
  153. cls.disk.uuid = cls.message['message']['passback_parameters']['disk_uuid']
  154. cls.disk.get_by('uuid')
  155. if state == ResponseState.success.value:
  156. cls.disk.guest_uuid = ''
  157. cls.disk.sequence = -1
  158. cls.disk.state = DiskState.idle.value
  159. cls.disk.update()
  160. elif action == 'delete_disk':
  161. cls.disk.uuid = uuid
  162. cls.disk.get_by('uuid')
  163. cls.disk.delete()
  164. elif action == 'boot':
  165. boot_jobs_id = cls.message['message']['passback_parameters']['boot_jobs_id']
  166. if state == ResponseState.success.value:
  167. cls.guest.uuid = uuid
  168. cls.guest.delete_boot_jobs(boot_jobs_id=boot_jobs_id)
  169. else:
  170. pass
  171. @classmethod
  172. def collection_performance_processor(cls):
  173. data_kind = cls.message['type']
  174. timestamp = ji.Common.ts()
  175. timestamp -= (timestamp % 60)
  176. data = cls.message['message']['data']
  177. if data_kind == CollectionPerformanceDataKind.cpu_memory.value:
  178. for item in data:
  179. cls.cpu_memory.guest_uuid = item['guest_uuid']
  180. cls.cpu_memory.cpu_load = item['cpu_load']
  181. cls.cpu_memory.memory_available = item['memory_available']
  182. cls.cpu_memory.memory_unused = item['memory_unused']
  183. cls.cpu_memory.timestamp = timestamp
  184. cls.cpu_memory.create()
  185. if data_kind == CollectionPerformanceDataKind.traffic.value:
  186. for item in data:
  187. cls.traffic.guest_uuid = item['guest_uuid']
  188. cls.traffic.name = item['name']
  189. cls.traffic.rx_bytes = item['rx_bytes']
  190. cls.traffic.rx_packets = item['rx_packets']
  191. cls.traffic.rx_errs = item['rx_errs']
  192. cls.traffic.rx_drop = item['rx_drop']
  193. cls.traffic.tx_bytes = item['tx_bytes']
  194. cls.traffic.tx_packets = item['tx_packets']
  195. cls.traffic.tx_errs = item['tx_errs']
  196. cls.traffic.tx_drop = item['tx_drop']
  197. cls.traffic.timestamp = timestamp
  198. cls.traffic.create()
  199. if data_kind == CollectionPerformanceDataKind.disk_io.value:
  200. for item in data:
  201. cls.disk_io.disk_uuid = item['disk_uuid']
  202. cls.disk_io.rd_req = item['rd_req']
  203. cls.disk_io.rd_bytes = item['rd_bytes']
  204. cls.disk_io.wr_req = item['wr_req']
  205. cls.disk_io.wr_bytes = item['wr_bytes']
  206. cls.disk_io.timestamp = timestamp
  207. cls.disk_io.create()
  208. else:
  209. pass
  210. @classmethod
  211. def launch(cls):
  212. while True:
  213. if Utils.exit_flag:
  214. print 'Thread EventProcessor say bye-bye'
  215. return
  216. try:
  217. report = db.r.lpop(app.config['upstream_queue'])
  218. if report is None:
  219. time.sleep(1)
  220. continue
  221. cls.message = json.loads(report)
  222. if cls.message['kind'] == EmitKind.log.value:
  223. cls.log_processor()
  224. elif cls.message['kind'] == EmitKind.guest_event.value:
  225. cls.guest_event_processor()
  226. elif cls.message['kind'] == EmitKind.host_event.value:
  227. cls.host_event_processor()
  228. elif cls.message['kind'] == EmitKind.response.value:
  229. cls.response_processor()
  230. elif cls.message['kind'] == EmitKind.collection_performance.value:
  231. cls.collection_performance_processor()
  232. else:
  233. pass
  234. except Exception as e:
  235. logger.error(e.message)