event_processor.py 17 KB

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