event_processor.py 15 KB

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