event_processor.py 12 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294
  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. 'cpu': cls.message['message']['cpu'],
  87. 'memory': cls.message['message']['memory'],
  88. 'interfaces': cls.message['message']['interfaces'],
  89. 'disks': cls.message['message']['disks'],
  90. 'timestamp': ji.Common.ts()
  91. }
  92. db.r.hset(app.config['hosts_info'], key=key, value=json.dumps(value, ensure_ascii=False))
  93. db.r.expire(app.config['hosts_info'], 10)
  94. @classmethod
  95. def response_processor(cls):
  96. action = cls.message['message']['action']
  97. uuid = cls.message['message']['uuid']
  98. state = cls.message['type']
  99. data = cls.message['message']['data']
  100. if action == 'create_guest':
  101. if state == ResponseState.success.value:
  102. # 系统盘的 UUID 与其 Guest 的 UUID 相同
  103. cls.disk.uuid = uuid
  104. cls.disk.get_by('uuid')
  105. cls.disk.guest_uuid = uuid
  106. cls.disk.state = DiskState.mounted.value
  107. # disk_info['virtual-size'] 的单位为Byte,需要除以 1024 的 3 次方,换算成单位为 GB 的值
  108. cls.disk.size = data['disk_info']['virtual-size'] / (1024 ** 3)
  109. cls.disk.update()
  110. else:
  111. cls.guest.uuid = uuid
  112. cls.guest.get_by('uuid')
  113. cls.guest.status = GuestState.dirty.value
  114. cls.guest.update()
  115. elif action == 'migrate':
  116. pass
  117. elif action == 'delete_guest':
  118. if state == ResponseState.success.value:
  119. cls.config.id = 1
  120. cls.config.get()
  121. cls.guest.uuid = uuid
  122. cls.guest.get_by('uuid')
  123. if IP(cls.config.start_ip).int() <= IP(cls.guest.ip).int() <= IP(cls.config.end_ip).int():
  124. if db.r.srem(app.config['ip_used_set'], cls.guest.ip):
  125. db.r.sadd(app.config['ip_available_set'], cls.guest.ip)
  126. if (cls.guest.vnc_port - cls.config.start_vnc_port) <= \
  127. (IP(cls.config.end_ip).int() - IP(cls.config.start_ip).int()):
  128. if db.r.srem(app.config['vnc_port_used_set'], cls.guest.vnc_port):
  129. db.r.sadd(app.config['vnc_port_available_set'], cls.guest.vnc_port)
  130. cls.guest.delete()
  131. cls.disk.uuid = uuid
  132. cls.disk.get_by('uuid')
  133. cls.disk.delete()
  134. elif action == 'create_disk':
  135. cls.disk.uuid = uuid
  136. cls.disk.get_by('uuid')
  137. if state == ResponseState.success.value:
  138. cls.disk.state = DiskState.idle.value
  139. else:
  140. cls.disk.state = DiskState.dirty.value
  141. cls.disk.update()
  142. elif action == 'resize_disk':
  143. if state == ResponseState.success.value:
  144. cls.disk.uuid = uuid
  145. cls.disk.get_by('uuid')
  146. cls.disk.size = cls.message['message']['passback_parameters']['size']
  147. cls.disk.update()
  148. elif action == 'attach_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 = uuid
  153. cls.disk.sequence = cls.message['message']['passback_parameters']['sequence']
  154. cls.disk.state = DiskState.mounted.value
  155. cls.disk.update()
  156. elif action == 'detach_disk':
  157. cls.disk.uuid = cls.message['message']['passback_parameters']['disk_uuid']
  158. cls.disk.get_by('uuid')
  159. if state == ResponseState.success.value:
  160. cls.disk.guest_uuid = ''
  161. cls.disk.sequence = -1
  162. cls.disk.state = DiskState.idle.value
  163. cls.disk.update()
  164. elif action == 'delete_disk':
  165. cls.disk.uuid = uuid
  166. cls.disk.get_by('uuid')
  167. cls.disk.delete()
  168. elif action == 'boot':
  169. boot_jobs_id = cls.message['message']['passback_parameters']['boot_jobs_id']
  170. if state == ResponseState.success.value:
  171. cls.guest.uuid = uuid
  172. cls.guest.delete_boot_jobs(boot_jobs_id=boot_jobs_id)
  173. else:
  174. pass
  175. @classmethod
  176. def collection_performance_processor(cls):
  177. data_kind = cls.message['type']
  178. timestamp = ji.Common.ts()
  179. timestamp -= (timestamp % 60)
  180. data = cls.message['message']['data']
  181. if data_kind == CollectionPerformanceDataKind.cpu_memory.value:
  182. for item in data:
  183. cls.cpu_memory.guest_uuid = item['guest_uuid']
  184. cls.cpu_memory.cpu_load = item['cpu_load']
  185. cls.cpu_memory.memory_available = item['memory_available']
  186. cls.cpu_memory.memory_unused = item['memory_unused']
  187. cls.cpu_memory.timestamp = timestamp
  188. cls.cpu_memory.create()
  189. if data_kind == CollectionPerformanceDataKind.traffic.value:
  190. for item in data:
  191. cls.traffic.guest_uuid = item['guest_uuid']
  192. cls.traffic.name = item['name']
  193. cls.traffic.rx_bytes = item['rx_bytes']
  194. cls.traffic.rx_packets = item['rx_packets']
  195. cls.traffic.rx_errs = item['rx_errs']
  196. cls.traffic.rx_drop = item['rx_drop']
  197. cls.traffic.tx_bytes = item['tx_bytes']
  198. cls.traffic.tx_packets = item['tx_packets']
  199. cls.traffic.tx_errs = item['tx_errs']
  200. cls.traffic.tx_drop = item['tx_drop']
  201. cls.traffic.timestamp = timestamp
  202. cls.traffic.create()
  203. if data_kind == CollectionPerformanceDataKind.disk_io.value:
  204. for item in data:
  205. cls.disk_io.disk_uuid = item['disk_uuid']
  206. cls.disk_io.rd_req = item['rd_req']
  207. cls.disk_io.rd_bytes = item['rd_bytes']
  208. cls.disk_io.wr_req = item['wr_req']
  209. cls.disk_io.wr_bytes = item['wr_bytes']
  210. cls.disk_io.timestamp = timestamp
  211. cls.disk_io.create()
  212. else:
  213. pass
  214. @classmethod
  215. def launch(cls):
  216. while True:
  217. if Utils.exit_flag:
  218. print 'Thread EventProcessor say bye-bye'
  219. return
  220. try:
  221. report = db.r.lpop(app.config['upstream_queue'])
  222. if report is None:
  223. time.sleep(1)
  224. continue
  225. cls.message = json.loads(report)
  226. if cls.message['kind'] == EmitKind.log.value:
  227. cls.log_processor()
  228. elif cls.message['kind'] == EmitKind.guest_event.value:
  229. cls.guest_event_processor()
  230. elif cls.message['kind'] == EmitKind.host_event.value:
  231. cls.host_event_processor()
  232. elif cls.message['kind'] == EmitKind.response.value:
  233. cls.response_processor()
  234. elif cls.message['kind'] == EmitKind.collection_performance.value:
  235. cls.collection_performance_processor()
  236. else:
  237. pass
  238. except Exception as e:
  239. logger.error(e.message)