event_processor.py 16 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366
  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, HostCollectionPerformanceDataKind
  17. from models import HostCPUMemory, HostTraffic, HostDiskUsageIO
  18. __author__ = 'James Iter'
  19. __date__ = '2017/4/15'
  20. __contact__ = 'james.iter.cn@gmail.com'
  21. __copyright__ = '(c) 2017 by James Iter.'
  22. class EventProcessor(object):
  23. message = None
  24. log = Log()
  25. guest = Guest()
  26. guest_migrate_info = GuestMigrateInfo()
  27. disk = Disk()
  28. config = Config()
  29. cpu_memory = CPUMemory()
  30. traffic = Traffic()
  31. disk_io = DiskIO()
  32. host_cpu_memory = HostCPUMemory()
  33. host_traffic = HostTraffic()
  34. host_disk_usage_io = HostDiskUsageIO()
  35. @classmethod
  36. def log_processor(cls):
  37. cls.log.set(type=cls.message['type'], timestamp=cls.message['timestamp'], host=cls.message['host'],
  38. message=cls.message['message'],
  39. full_message='' if cls.message['message'].__len__() < 255 else cls.message['message'])
  40. cls.log.create()
  41. @classmethod
  42. def guest_event_processor(cls):
  43. cls.guest.uuid = cls.message['message']['uuid']
  44. cls.guest.get_by('uuid')
  45. cls.guest.on_host = cls.message['host']
  46. last_status = cls.guest.status
  47. cls.guest.status = cls.message['type']
  48. if cls.message['type'] == GuestState.update.value:
  49. # 更新事件不改变 Guest 的状态
  50. cls.guest.status = last_status
  51. cls.guest.xml = cls.message['message']['xml']
  52. elif cls.guest.status == GuestState.migrating.value:
  53. try:
  54. cls.guest_migrate_info.uuid = cls.guest.uuid
  55. cls.guest_migrate_info.get_by('uuid')
  56. cls.guest_migrate_info.type = cls.message['message']['migrating_info']['type']
  57. cls.guest_migrate_info.time_elapsed = cls.message['message']['migrating_info']['time_elapsed']
  58. cls.guest_migrate_info.time_remaining = cls.message['message']['migrating_info']['time_remaining']
  59. cls.guest_migrate_info.data_total = cls.message['message']['migrating_info']['data_total']
  60. cls.guest_migrate_info.data_processed = cls.message['message']['migrating_info']['data_processed']
  61. cls.guest_migrate_info.data_remaining = cls.message['message']['migrating_info']['data_remaining']
  62. cls.guest_migrate_info.mem_total = cls.message['message']['migrating_info']['mem_total']
  63. cls.guest_migrate_info.mem_processed = cls.message['message']['migrating_info']['mem_processed']
  64. cls.guest_migrate_info.mem_remaining = cls.message['message']['migrating_info']['mem_remaining']
  65. cls.guest_migrate_info.file_total = cls.message['message']['migrating_info']['file_total']
  66. cls.guest_migrate_info.file_processed = cls.message['message']['migrating_info']['file_processed']
  67. cls.guest_migrate_info.file_remaining = cls.message['message']['migrating_info']['file_remaining']
  68. cls.guest_migrate_info.update()
  69. except ji.PreviewingError as e:
  70. ret = json.loads(e.message)
  71. if ret['state']['code'] == '404':
  72. cls.guest_migrate_info.type = cls.message['message']['migrating_info']['type']
  73. cls.guest_migrate_info.time_elapsed = cls.message['message']['migrating_info']['time_elapsed']
  74. cls.guest_migrate_info.time_remaining = cls.message['message']['migrating_info']['time_remaining']
  75. cls.guest_migrate_info.data_total = cls.message['message']['migrating_info']['data_total']
  76. cls.guest_migrate_info.data_processed = cls.message['message']['migrating_info']['data_processed']
  77. cls.guest_migrate_info.data_remaining = cls.message['message']['migrating_info']['data_remaining']
  78. cls.guest_migrate_info.mem_total = cls.message['message']['migrating_info']['mem_total']
  79. cls.guest_migrate_info.mem_processed = cls.message['message']['migrating_info']['mem_processed']
  80. cls.guest_migrate_info.mem_remaining = cls.message['message']['migrating_info']['mem_remaining']
  81. cls.guest_migrate_info.file_total = cls.message['message']['migrating_info']['file_total']
  82. cls.guest_migrate_info.file_processed = cls.message['message']['migrating_info']['file_processed']
  83. cls.guest_migrate_info.file_remaining = cls.message['message']['migrating_info']['file_remaining']
  84. cls.guest_migrate_info.create()
  85. elif cls.guest.status == GuestState.creating.value:
  86. cls.guest.progress = cls.message['message']['progress']
  87. cls.guest.update()
  88. # 限定特殊情况下更新磁盘所属 Guest,避免迁移、创建时频繁被无意义的更新
  89. if cls.guest.status in [GuestState.running.value, GuestState.shutoff.value]:
  90. cls.disk.update_by_filter({'on_host': cls.guest.on_host}, filter_str='guest_uuid:eq:' + cls.guest.uuid)
  91. @classmethod
  92. def host_event_processor(cls):
  93. key = cls.message['message']['node_id']
  94. value = {
  95. 'hostname': cls.message['host'],
  96. 'cpu': cls.message['message']['cpu'],
  97. 'system_load': cls.message['message']['system_load'],
  98. 'memory': cls.message['message']['memory'],
  99. 'memory_available': cls.message['message']['memory_available'],
  100. 'interfaces': cls.message['message']['interfaces'],
  101. 'disks': cls.message['message']['disks'],
  102. 'boot_time': cls.message['message']['boot_time'],
  103. 'timestamp': ji.Common.ts()
  104. }
  105. db.r.hset(app.config['hosts_info'], key=key, value=json.dumps(value, ensure_ascii=False))
  106. @classmethod
  107. def response_processor(cls):
  108. _object = cls.message['message']['_object']
  109. action = cls.message['message']['action']
  110. uuid = cls.message['message']['uuid']
  111. state = cls.message['type']
  112. data = cls.message['message']['data']
  113. hostname = cls.message['host']
  114. if _object == 'guest':
  115. if action == 'create':
  116. if state == ResponseState.success.value:
  117. # 系统盘的 UUID 与其 Guest 的 UUID 相同
  118. cls.disk.uuid = uuid
  119. cls.disk.get_by('uuid')
  120. cls.disk.guest_uuid = uuid
  121. cls.disk.state = DiskState.mounted.value
  122. # disk_info['virtual-size'] 的单位为Byte,需要除以 1024 的 3 次方,换算成单位为 GB 的值
  123. cls.disk.size = data['disk_info']['virtual-size'] / (1024 ** 3)
  124. cls.disk.update()
  125. else:
  126. cls.guest.uuid = uuid
  127. cls.guest.get_by('uuid')
  128. cls.guest.status = GuestState.dirty.value
  129. cls.guest.update()
  130. elif action == 'migrate':
  131. pass
  132. elif action == 'delete':
  133. if state == ResponseState.success.value:
  134. cls.config.id = 1
  135. cls.config.get()
  136. cls.guest.uuid = uuid
  137. cls.guest.get_by('uuid')
  138. if IP(cls.config.start_ip).int() <= IP(cls.guest.ip).int() <= IP(cls.config.end_ip).int():
  139. if db.r.srem(app.config['ip_used_set'], cls.guest.ip):
  140. db.r.sadd(app.config['ip_available_set'], cls.guest.ip)
  141. if (cls.guest.vnc_port - cls.config.start_vnc_port) <= \
  142. (IP(cls.config.end_ip).int() - IP(cls.config.start_ip).int()):
  143. if db.r.srem(app.config['vnc_port_used_set'], cls.guest.vnc_port):
  144. db.r.sadd(app.config['vnc_port_available_set'], cls.guest.vnc_port)
  145. cls.guest.delete()
  146. # TODO: 加入是否删除使用的数据磁盘开关,如果为True,则顺便删除使用的磁盘。否则解除该磁盘被使用的状态。
  147. cls.disk.uuid = uuid
  148. cls.disk.get_by('uuid')
  149. cls.disk.delete()
  150. cls.disk.update_by_filter({'guest_uuid': '', 'sequence': -1, 'state': DiskState.idle.value},
  151. filter_str='guest_uuid:eq:' + cls.guest.uuid)
  152. elif action == 'attach_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 = uuid
  157. cls.disk.sequence = cls.message['message']['passback_parameters']['sequence']
  158. cls.disk.state = DiskState.mounted.value
  159. cls.disk.update()
  160. elif action == 'detach_disk':
  161. cls.disk.uuid = cls.message['message']['passback_parameters']['disk_uuid']
  162. cls.disk.get_by('uuid')
  163. if state == ResponseState.success.value:
  164. cls.disk.guest_uuid = ''
  165. cls.disk.sequence = -1
  166. cls.disk.state = DiskState.idle.value
  167. cls.disk.update()
  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. elif _object == 'disk':
  174. if action == 'create':
  175. cls.disk.uuid = uuid
  176. cls.disk.get_by('uuid')
  177. cls.disk.on_host = hostname
  178. if state == ResponseState.success.value:
  179. cls.disk.state = DiskState.idle.value
  180. else:
  181. cls.disk.state = DiskState.dirty.value
  182. cls.disk.update()
  183. elif action == 'resize':
  184. if state == ResponseState.success.value:
  185. cls.disk.uuid = uuid
  186. cls.disk.get_by('uuid')
  187. cls.disk.size = cls.message['message']['passback_parameters']['size']
  188. cls.disk.update()
  189. elif action == 'delete':
  190. cls.disk.uuid = uuid
  191. cls.disk.get_by('uuid')
  192. cls.disk.delete()
  193. else:
  194. pass
  195. @classmethod
  196. def collection_performance_processor(cls):
  197. data_kind = cls.message['type']
  198. timestamp = ji.Common.ts()
  199. timestamp -= (timestamp % 60)
  200. data = cls.message['message']['data']
  201. if data_kind == CollectionPerformanceDataKind.cpu_memory.value:
  202. for item in data:
  203. cls.cpu_memory.guest_uuid = item['guest_uuid']
  204. cls.cpu_memory.cpu_load = item['cpu_load']
  205. cls.cpu_memory.memory_available = item['memory_available']
  206. cls.cpu_memory.memory_unused = item['memory_unused']
  207. cls.cpu_memory.timestamp = timestamp
  208. cls.cpu_memory.create()
  209. if data_kind == CollectionPerformanceDataKind.traffic.value:
  210. for item in data:
  211. cls.traffic.guest_uuid = item['guest_uuid']
  212. cls.traffic.name = item['name']
  213. cls.traffic.rx_bytes = item['rx_bytes']
  214. cls.traffic.rx_packets = item['rx_packets']
  215. cls.traffic.rx_errs = item['rx_errs']
  216. cls.traffic.rx_drop = item['rx_drop']
  217. cls.traffic.tx_bytes = item['tx_bytes']
  218. cls.traffic.tx_packets = item['tx_packets']
  219. cls.traffic.tx_errs = item['tx_errs']
  220. cls.traffic.tx_drop = item['tx_drop']
  221. cls.traffic.timestamp = timestamp
  222. cls.traffic.create()
  223. if data_kind == CollectionPerformanceDataKind.disk_io.value:
  224. for item in data:
  225. cls.disk_io.disk_uuid = item['disk_uuid']
  226. cls.disk_io.rd_req = item['rd_req']
  227. cls.disk_io.rd_bytes = item['rd_bytes']
  228. cls.disk_io.wr_req = item['wr_req']
  229. cls.disk_io.wr_bytes = item['wr_bytes']
  230. cls.disk_io.timestamp = timestamp
  231. cls.disk_io.create()
  232. else:
  233. pass
  234. @classmethod
  235. def host_collection_performance_processor(cls):
  236. data_kind = cls.message['type']
  237. timestamp = ji.Common.ts()
  238. timestamp -= (timestamp % 60)
  239. data = cls.message['message']['data']
  240. if data_kind == HostCollectionPerformanceDataKind.cpu_memory.value:
  241. cls.host_cpu_memory.node_id = data['node_id']
  242. cls.host_cpu_memory.cpu_load = data['cpu_load']
  243. cls.host_cpu_memory.memory_available = data['memory_available']
  244. cls.host_cpu_memory.timestamp = timestamp
  245. cls.host_cpu_memory.create()
  246. if data_kind == HostCollectionPerformanceDataKind.traffic.value:
  247. for item in data:
  248. cls.host_traffic.node_id = item['node_id']
  249. cls.host_traffic.name = item['name']
  250. cls.host_traffic.rx_bytes = item['rx_bytes']
  251. cls.host_traffic.rx_packets = item['rx_packets']
  252. cls.host_traffic.rx_errs = item['rx_errs']
  253. cls.host_traffic.rx_drop = item['rx_drop']
  254. cls.host_traffic.tx_bytes = item['tx_bytes']
  255. cls.host_traffic.tx_packets = item['tx_packets']
  256. cls.host_traffic.tx_errs = item['tx_errs']
  257. cls.host_traffic.tx_drop = item['tx_drop']
  258. cls.host_traffic.timestamp = timestamp
  259. cls.host_traffic.create()
  260. if data_kind == HostCollectionPerformanceDataKind.disk_usage_io.value:
  261. for item in data:
  262. cls.host_disk_usage_io.node_id = item['node_id']
  263. cls.host_disk_usage_io.mountpoint = item['mountpoint']
  264. cls.host_disk_usage_io.used = item['used']
  265. cls.host_disk_usage_io.rd_req = item['rd_req']
  266. cls.host_disk_usage_io.rd_bytes = item['rd_bytes']
  267. cls.host_disk_usage_io.wr_req = item['wr_req']
  268. cls.host_disk_usage_io.wr_bytes = item['wr_bytes']
  269. cls.host_disk_usage_io.timestamp = timestamp
  270. cls.host_disk_usage_io.create()
  271. else:
  272. pass
  273. @classmethod
  274. def launch(cls):
  275. while True:
  276. if Utils.exit_flag:
  277. msg = 'Thread EventProcessor say bye-bye'
  278. print msg
  279. logger.info(msg=msg)
  280. return
  281. try:
  282. report = db.r.lpop(app.config['upstream_queue'])
  283. if report is None:
  284. time.sleep(1)
  285. continue
  286. cls.message = json.loads(report)
  287. if cls.message['kind'] == EmitKind.log.value:
  288. cls.log_processor()
  289. elif cls.message['kind'] == EmitKind.guest_event.value:
  290. cls.guest_event_processor()
  291. elif cls.message['kind'] == EmitKind.host_event.value:
  292. cls.host_event_processor()
  293. elif cls.message['kind'] == EmitKind.response.value:
  294. cls.response_processor()
  295. elif cls.message['kind'] == EmitKind.collection_performance.value:
  296. cls.collection_performance_processor()
  297. elif cls.message['kind'] == EmitKind.host_collection_performance.value:
  298. cls.host_collection_performance_processor()
  299. else:
  300. pass
  301. except Exception as e:
  302. logger.error(e.message)