event_processor.py 22 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515
  1. #!/usr/bin/env python
  2. # -*- coding: utf-8 -*-
  3. import traceback
  4. import json
  5. import time
  6. import jimit as ji
  7. from jimvc.models import logger, app_config
  8. from jimvc.models import Database as db, Config, GuestCPUMemory, GuestTraffic, GuestDiskIO, SSHKeyGuestMapping
  9. from jimvc.models import Guest
  10. from jimvc.models import Disk
  11. from jimvc.models import Snapshot, SnapshotDiskMapping, OSTemplateImage
  12. from jimvc.models import Log
  13. from jimvc.models import Utils
  14. from jimvc.models import EmitKind
  15. from jimvc.models import ResponseState, GuestState, DiskState
  16. from jimvc.models import GuestMigrateInfo
  17. from jimvc.models import GuestCollectionPerformanceDataKind, HostCollectionPerformanceDataKind
  18. from jimvc.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. snapshot = Snapshot()
  30. snapshot_disk_mapping = SnapshotDiskMapping()
  31. os_template_image = OSTemplateImage()
  32. config = Config()
  33. config.id = 1
  34. guest_cpu_memory = GuestCPUMemory()
  35. guest_traffic = GuestTraffic()
  36. guest_disk_io = GuestDiskIO()
  37. host_cpu_memory = HostCPUMemory()
  38. host_traffic = HostTraffic()
  39. host_disk_usage_io = HostDiskUsageIO()
  40. @classmethod
  41. def log_processor(cls):
  42. cls.log.set(type=cls.message['type'], timestamp=cls.message['timestamp'], host=cls.message['host'],
  43. message=cls.message['message'],
  44. full_message='' if cls.message['message'].__len__() < 255 else cls.message['message'])
  45. cls.log.create()
  46. pass
  47. @classmethod
  48. def guest_event_processor(cls):
  49. cls.guest.uuid = cls.message['message']['uuid']
  50. cls.guest.get_by('uuid')
  51. cls.guest.node_id = cls.message['node_id']
  52. last_status = cls.guest.status
  53. cls.guest.status = cls.message['type']
  54. if cls.message['type'] == GuestState.update.value:
  55. # 更新事件不改变 Guest 的状态
  56. cls.guest.status = last_status
  57. cls.guest.xml = cls.message['message']['xml']
  58. elif cls.guest.status == GuestState.migrating.value:
  59. try:
  60. cls.guest_migrate_info.uuid = cls.guest.uuid
  61. cls.guest_migrate_info.get_by('uuid')
  62. cls.guest_migrate_info.type = cls.message['message']['migrating_info']['type']
  63. cls.guest_migrate_info.time_elapsed = cls.message['message']['migrating_info']['time_elapsed']
  64. cls.guest_migrate_info.time_remaining = cls.message['message']['migrating_info']['time_remaining']
  65. cls.guest_migrate_info.data_total = cls.message['message']['migrating_info']['data_total']
  66. cls.guest_migrate_info.data_processed = cls.message['message']['migrating_info']['data_processed']
  67. cls.guest_migrate_info.data_remaining = cls.message['message']['migrating_info']['data_remaining']
  68. cls.guest_migrate_info.mem_total = cls.message['message']['migrating_info']['mem_total']
  69. cls.guest_migrate_info.mem_processed = cls.message['message']['migrating_info']['mem_processed']
  70. cls.guest_migrate_info.mem_remaining = cls.message['message']['migrating_info']['mem_remaining']
  71. cls.guest_migrate_info.file_total = cls.message['message']['migrating_info']['file_total']
  72. cls.guest_migrate_info.file_processed = cls.message['message']['migrating_info']['file_processed']
  73. cls.guest_migrate_info.file_remaining = cls.message['message']['migrating_info']['file_remaining']
  74. cls.guest_migrate_info.update()
  75. except ji.PreviewingError as e:
  76. ret = json.loads(e.message)
  77. if ret['state']['code'] == '404':
  78. cls.guest_migrate_info.type = cls.message['message']['migrating_info']['type']
  79. cls.guest_migrate_info.time_elapsed = cls.message['message']['migrating_info']['time_elapsed']
  80. cls.guest_migrate_info.time_remaining = cls.message['message']['migrating_info']['time_remaining']
  81. cls.guest_migrate_info.data_total = cls.message['message']['migrating_info']['data_total']
  82. cls.guest_migrate_info.data_processed = cls.message['message']['migrating_info']['data_processed']
  83. cls.guest_migrate_info.data_remaining = cls.message['message']['migrating_info']['data_remaining']
  84. cls.guest_migrate_info.mem_total = cls.message['message']['migrating_info']['mem_total']
  85. cls.guest_migrate_info.mem_processed = cls.message['message']['migrating_info']['mem_processed']
  86. cls.guest_migrate_info.mem_remaining = cls.message['message']['migrating_info']['mem_remaining']
  87. cls.guest_migrate_info.file_total = cls.message['message']['migrating_info']['file_total']
  88. cls.guest_migrate_info.file_processed = cls.message['message']['migrating_info']['file_processed']
  89. cls.guest_migrate_info.file_remaining = cls.message['message']['migrating_info']['file_remaining']
  90. cls.guest_migrate_info.create()
  91. elif cls.guest.status == GuestState.creating.value:
  92. if cls.message['message']['progress'] <= cls.guest.progress:
  93. return
  94. cls.guest.progress = cls.message['message']['progress']
  95. elif cls.guest.status == GuestState.snapshot_converting.value:
  96. cls.os_template_image.id = cls.message['message']['os_template_image_id']
  97. cls.os_template_image.get()
  98. if cls.message['message']['progress'] <= cls.os_template_image.progress:
  99. return
  100. cls.os_template_image.progress = cls.message['message']['progress']
  101. cls.os_template_image.update()
  102. return
  103. cls.guest.update()
  104. # 限定特殊情况下更新磁盘所属 Guest,避免迁移、创建时频繁被无意义的更新
  105. if cls.guest.status in [GuestState.running.value, GuestState.shutoff.value]:
  106. cls.disk.update_by_filter({'node_id': cls.guest.node_id}, filter_str='guest_uuid:eq:' + cls.guest.uuid)
  107. @classmethod
  108. def host_event_processor(cls):
  109. key = cls.message['message']['node_id']
  110. value = {
  111. 'hostname': cls.message['host'],
  112. 'cpu': cls.message['message']['cpu'],
  113. 'cpuinfo': cls.message['message'].get('cpuinfo'),
  114. 'system_load': cls.message['message']['system_load'],
  115. 'memory': cls.message['message']['memory'],
  116. 'memory_available': cls.message['message']['memory_available'],
  117. 'dmidecode': cls.message['message']['dmidecode'],
  118. 'interfaces': cls.message['message']['interfaces'],
  119. 'disks': cls.message['message']['disks'],
  120. 'boot_time': cls.message['message']['boot_time'],
  121. 'nonrandom': False,
  122. 'threads_status': cls.message['message']['threads_status'],
  123. 'version': cls.message['message']['version'],
  124. 'timestamp': ji.Common.ts()
  125. }
  126. db.r.hset(app_config['hosts_info'], key=key, value=json.dumps(value, ensure_ascii=False))
  127. @classmethod
  128. def response_processor(cls):
  129. _object = cls.message['message']['_object']
  130. action = cls.message['message']['action']
  131. uuid = cls.message['message']['uuid']
  132. state = cls.message['type']
  133. data = cls.message['message']['data']
  134. node_id = cls.message['node_id']
  135. if _object == 'guest':
  136. if action == 'create':
  137. if state == ResponseState.success.value:
  138. # 系统盘的 UUID 与其 Guest 的 UUID 相同
  139. cls.disk.uuid = uuid
  140. cls.disk.get_by('uuid')
  141. cls.disk.guest_uuid = uuid
  142. cls.disk.state = DiskState.mounted.value
  143. # disk_info['virtual-size'] 的单位为Byte,需要除以 1024 的 3 次方,换算成单位为 GB 的值
  144. cls.disk.size = data['disk_info']['virtual-size'] / (1024 ** 3)
  145. cls.disk.update()
  146. cls.guest.uuid = uuid
  147. cls.guest.get_by('uuid')
  148. cls.guest.progress = 100
  149. cls.guest.update()
  150. else:
  151. cls.guest.uuid = uuid
  152. cls.guest.get_by('uuid')
  153. cls.guest.status = GuestState.dirty.value
  154. cls.guest.update()
  155. elif action == 'migrate':
  156. pass
  157. elif action == 'sync_data':
  158. if state == ResponseState.success.value:
  159. cls.guest.uuid = uuid
  160. cls.guest.get_by('uuid')
  161. cls.guest.autostart = bool(data['autostart'])
  162. cls.guest.update()
  163. elif action == 'delete':
  164. if state == ResponseState.success.value:
  165. cls.config.get()
  166. cls.guest.uuid = uuid
  167. cls.guest.get_by('uuid')
  168. cls.guest.delete()
  169. # TODO: 加入是否删除使用的数据磁盘开关,如果为True,则顺便删除使用的磁盘。否则解除该磁盘被使用的状态。
  170. cls.disk.uuid = uuid
  171. cls.disk.get_by('uuid')
  172. cls.disk.delete()
  173. cls.disk.update_by_filter({'guest_uuid': '', 'sequence': -1, 'state': DiskState.idle.value},
  174. filter_str='guest_uuid:eq:' + cls.guest.uuid)
  175. SSHKeyGuestMapping.delete_by_filter(filter_str=':'.join(['guest_uuid', 'eq', cls.guest.uuid]))
  176. elif action == 'reset_password':
  177. if state == ResponseState.success.value:
  178. cls.guest.uuid = uuid
  179. cls.guest.get_by('uuid')
  180. cls.guest.password = cls.message['message']['passback_parameters']['password']
  181. cls.guest.update()
  182. elif action == 'autostart':
  183. if state == ResponseState.success.value:
  184. cls.guest.uuid = uuid
  185. cls.guest.get_by('uuid')
  186. cls.guest.autostart = cls.message['message']['passback_parameters']['autostart']
  187. cls.guest.update()
  188. elif action == 'attach_disk':
  189. cls.disk.uuid = cls.message['message']['passback_parameters']['disk_uuid']
  190. cls.disk.get_by('uuid')
  191. if state == ResponseState.success.value:
  192. cls.disk.guest_uuid = uuid
  193. cls.disk.sequence = cls.message['message']['passback_parameters']['sequence']
  194. cls.disk.state = DiskState.mounted.value
  195. cls.disk.update()
  196. elif action == 'detach_disk':
  197. cls.disk.uuid = cls.message['message']['passback_parameters']['disk_uuid']
  198. cls.disk.get_by('uuid')
  199. if state == ResponseState.success.value:
  200. cls.disk.guest_uuid = ''
  201. cls.disk.sequence = -1
  202. cls.disk.state = DiskState.idle.value
  203. cls.disk.update()
  204. elif action == 'boot':
  205. if state == ResponseState.success.value:
  206. pass
  207. elif action == 'allocate_bandwidth':
  208. if state == ResponseState.success.value:
  209. cls.guest.uuid = uuid
  210. cls.guest.get_by('uuid')
  211. cls.guest.bandwidth = cls.message['message']['passback_parameters']['bandwidth']
  212. cls.guest.update()
  213. elif action == 'adjust_ability':
  214. if state == ResponseState.success.value:
  215. cls.guest.uuid = uuid
  216. cls.guest.get_by('uuid')
  217. cls.guest.cpu = cls.message['message']['passback_parameters']['cpu']
  218. cls.guest.memory = cls.message['message']['passback_parameters']['memory']
  219. cls.guest.update()
  220. elif action == 'change_vlan':
  221. if state == ResponseState.success.value:
  222. cls.guest.uuid = uuid
  223. cls.guest.get_by('uuid')
  224. cls.guest.vlan_id = cls.message['message']['passback_parameters']['vlan_id']
  225. cls.guest.network = 'vlan' + cls.guest.vlan_id.__str__()
  226. cls.guest.update()
  227. elif _object == 'disk':
  228. if action == 'create':
  229. cls.disk.uuid = uuid
  230. cls.disk.get_by('uuid')
  231. cls.disk.node_id = node_id
  232. if state == ResponseState.success.value:
  233. cls.disk.state = DiskState.idle.value
  234. else:
  235. cls.disk.state = DiskState.dirty.value
  236. cls.disk.update()
  237. elif action == 'resize':
  238. if state == ResponseState.success.value:
  239. cls.config.get()
  240. cls.disk.uuid = uuid
  241. cls.disk.get_by('uuid')
  242. cls.disk.size = cls.message['message']['passback_parameters']['size']
  243. cls.disk.quota(config=cls.config)
  244. cls.disk.update()
  245. elif action == 'delete':
  246. cls.disk.uuid = uuid
  247. cls.disk.get_by('uuid')
  248. cls.disk.delete()
  249. elif _object == 'snapshot':
  250. if action == 'create':
  251. cls.snapshot.id = cls.message['message']['passback_parameters']['id']
  252. cls.snapshot.get()
  253. if state == ResponseState.success.value:
  254. cls.snapshot.snapshot_id = data['snapshot_id']
  255. cls.snapshot.parent_id = data['parent_id']
  256. cls.snapshot.xml = data['xml']
  257. cls.snapshot.progress = 100
  258. cls.snapshot.update()
  259. disks, _ = Disk.get_by_filter(filter_str='guest_uuid:eq:' + cls.snapshot.guest_uuid)
  260. for disk in disks:
  261. cls.snapshot_disk_mapping.snapshot_id = cls.snapshot.snapshot_id
  262. cls.snapshot_disk_mapping.disk_uuid = disk['uuid']
  263. cls.snapshot_disk_mapping.create()
  264. else:
  265. cls.snapshot.progress = 255
  266. cls.snapshot.update()
  267. if action == 'delete':
  268. if state == ResponseState.success.value:
  269. cls.snapshot.id = cls.message['message']['passback_parameters']['id']
  270. cls.snapshot.get()
  271. # 更新子快照的 parent_id 为,当前快照的 parent_id。因为当前快照已被删除。
  272. Snapshot.update_by_filter({'parent_id': cls.snapshot.parent_id},
  273. filter_str='parent_id:eq:' + cls.snapshot.snapshot_id)
  274. SnapshotDiskMapping.delete_by_filter(
  275. filter_str=':'.join(['snapshot_id', 'eq', cls.snapshot.snapshot_id]))
  276. cls.snapshot.delete()
  277. else:
  278. pass
  279. if action == 'revert':
  280. # 不论恢复成功与否,都使快照恢复至正常状态。
  281. cls.snapshot.id = cls.message['message']['passback_parameters']['id']
  282. cls.snapshot.get()
  283. cls.snapshot.progress = 100
  284. cls.snapshot.update()
  285. if action == 'convert':
  286. cls.snapshot.snapshot_id = cls.message['message']['passback_parameters']['id']
  287. cls.snapshot.get_by('snapshot_id')
  288. cls.snapshot.progress = 100
  289. cls.snapshot.update()
  290. cls.os_template_image.id = cls.message['message']['passback_parameters']['os_template_image_id']
  291. cls.os_template_image.get()
  292. if state == ResponseState.success.value:
  293. cls.os_template_image.progress = 100
  294. else:
  295. cls.os_template_image.progress = 255
  296. cls.os_template_image.update()
  297. elif _object == 'os_template_image':
  298. if action == 'delete':
  299. cls.os_template_image.id = cls.message['message']['passback_parameters']['id']
  300. cls.os_template_image.get()
  301. if state == ResponseState.success.value:
  302. cls.os_template_image.delete()
  303. else:
  304. pass
  305. else:
  306. pass
  307. @classmethod
  308. def guest_collection_performance_processor(cls):
  309. data_kind = cls.message['type']
  310. timestamp = ji.Common.ts()
  311. timestamp -= (timestamp % 60)
  312. data = cls.message['message']['data']
  313. if data_kind == GuestCollectionPerformanceDataKind.cpu_memory.value:
  314. for item in data:
  315. cls.guest_cpu_memory.guest_uuid = item['guest_uuid']
  316. cls.guest_cpu_memory.cpu_load = item['cpu_load']
  317. cls.guest_cpu_memory.memory_available = item['memory_available']
  318. cls.guest_cpu_memory.memory_rate = item['memory_rate']
  319. cls.guest_cpu_memory.timestamp = timestamp
  320. cls.guest_cpu_memory.create()
  321. if data_kind == GuestCollectionPerformanceDataKind.traffic.value:
  322. for item in data:
  323. cls.guest_traffic.guest_uuid = item['guest_uuid']
  324. cls.guest_traffic.name = item['name']
  325. cls.guest_traffic.rx_bytes = item['rx_bytes']
  326. cls.guest_traffic.rx_packets = item['rx_packets']
  327. cls.guest_traffic.rx_errs = item['rx_errs']
  328. cls.guest_traffic.rx_drop = item['rx_drop']
  329. cls.guest_traffic.tx_bytes = item['tx_bytes']
  330. cls.guest_traffic.tx_packets = item['tx_packets']
  331. cls.guest_traffic.tx_errs = item['tx_errs']
  332. cls.guest_traffic.tx_drop = item['tx_drop']
  333. cls.guest_traffic.timestamp = timestamp
  334. cls.guest_traffic.create()
  335. if data_kind == GuestCollectionPerformanceDataKind.disk_io.value:
  336. for item in data:
  337. cls.guest_disk_io.disk_uuid = item['disk_uuid']
  338. cls.guest_disk_io.rd_req = item['rd_req']
  339. cls.guest_disk_io.rd_bytes = item['rd_bytes']
  340. cls.guest_disk_io.wr_req = item['wr_req']
  341. cls.guest_disk_io.wr_bytes = item['wr_bytes']
  342. cls.guest_disk_io.timestamp = timestamp
  343. cls.guest_disk_io.create()
  344. else:
  345. pass
  346. @classmethod
  347. def host_collection_performance_processor(cls):
  348. data_kind = cls.message['type']
  349. timestamp = ji.Common.ts()
  350. timestamp -= (timestamp % 60)
  351. data = cls.message['message']['data']
  352. if data_kind == HostCollectionPerformanceDataKind.cpu_memory.value:
  353. cls.host_cpu_memory.node_id = data['node_id']
  354. cls.host_cpu_memory.cpu_load = data['cpu_load']
  355. cls.host_cpu_memory.memory_available = data['memory_available']
  356. cls.host_cpu_memory.timestamp = timestamp
  357. cls.host_cpu_memory.create()
  358. if data_kind == HostCollectionPerformanceDataKind.traffic.value:
  359. for item in data:
  360. cls.host_traffic.node_id = item['node_id']
  361. cls.host_traffic.name = item['name']
  362. cls.host_traffic.rx_bytes = item['rx_bytes']
  363. cls.host_traffic.rx_packets = item['rx_packets']
  364. cls.host_traffic.rx_errs = item['rx_errs']
  365. cls.host_traffic.rx_drop = item['rx_drop']
  366. cls.host_traffic.tx_bytes = item['tx_bytes']
  367. cls.host_traffic.tx_packets = item['tx_packets']
  368. cls.host_traffic.tx_errs = item['tx_errs']
  369. cls.host_traffic.tx_drop = item['tx_drop']
  370. cls.host_traffic.timestamp = timestamp
  371. cls.host_traffic.create()
  372. if data_kind == HostCollectionPerformanceDataKind.disk_usage_io.value:
  373. for item in data:
  374. cls.host_disk_usage_io.node_id = item['node_id']
  375. cls.host_disk_usage_io.mountpoint = item['mountpoint']
  376. cls.host_disk_usage_io.used = item['used']
  377. cls.host_disk_usage_io.rd_req = item['rd_req']
  378. cls.host_disk_usage_io.rd_bytes = item['rd_bytes']
  379. cls.host_disk_usage_io.wr_req = item['wr_req']
  380. cls.host_disk_usage_io.wr_bytes = item['wr_bytes']
  381. cls.host_disk_usage_io.timestamp = timestamp
  382. cls.host_disk_usage_io.create()
  383. else:
  384. pass
  385. @classmethod
  386. def launch(cls):
  387. logger.info(msg='Thread EventProcessor is launched.')
  388. while True:
  389. if Utils.exit_flag:
  390. msg = 'Thread EventProcessor say bye-bye'
  391. print msg
  392. logger.info(msg=msg)
  393. return
  394. try:
  395. report = db.r.lpop(app_config['upstream_queue'])
  396. if report is None:
  397. time.sleep(1)
  398. continue
  399. cls.message = json.loads(report)
  400. if cls.message['kind'] == EmitKind.log.value:
  401. cls.log_processor()
  402. elif cls.message['kind'] == EmitKind.guest_event.value:
  403. cls.guest_event_processor()
  404. elif cls.message['kind'] == EmitKind.host_event.value:
  405. cls.host_event_processor()
  406. elif cls.message['kind'] == EmitKind.response.value:
  407. cls.response_processor()
  408. elif cls.message['kind'] == EmitKind.guest_collection_performance.value:
  409. cls.guest_collection_performance_processor()
  410. elif cls.message['kind'] == EmitKind.host_collection_performance.value:
  411. cls.host_collection_performance_processor()
  412. else:
  413. pass
  414. except AttributeError as e:
  415. logger.error(traceback.format_exc())
  416. time.sleep(1)
  417. if db.r is None:
  418. db.init_conn_redis()
  419. except Exception as e:
  420. logger.error(traceback.format_exc())
  421. time.sleep(1)