event_processor.py 23 KB

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