event_processor.py 16 KB

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