| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230 |
- #!/usr/bin/env python
- # -*- coding: utf-8 -*-
- import json
- import time
- from IPy import IP
- import jimit as ji
- from models import Database as db, Config
- from models import Guest
- from models import Disk
- from models import Log
- from models import Utils
- from models import EmitKind
- from models import ResponseState, GuestState, DiskState
- from models.guest import GuestMigrateInfo
- from models.initialize import app, logger
- __author__ = 'James Iter'
- __date__ = '2017/4/15'
- __contact__ = 'james.iter.cn@gmail.com'
- __copyright__ = '(c) 2017 by James Iter.'
- class EventProcessor(object):
- message = None
- log = Log()
- guest = Guest()
- guest_migrate_info = GuestMigrateInfo()
- disk = Disk()
- config = Config()
- @classmethod
- def log_processor(cls):
- cls.log.set(type=cls.message['type'], timestamp=cls.message['timestamp'], host=cls.message['host'],
- message=cls.message['message'])
- cls.log.create()
- @classmethod
- def guest_event_processor(cls):
- cls.guest.uuid = cls.message['message']['uuid']
- cls.guest.get_by('uuid')
- if cls.message['type'] == GuestState.update.value:
- cls.message['type'] = cls.guest.status
- cls.guest.xml = cls.message['message']['xml']
- elif cls.guest.status == GuestState.migrating.value:
- try:
- cls.guest_migrate_info.uuid = cls.guest.uuid
- cls.guest_migrate_info.get_by('uuid')
- cls.guest_migrate_info.type = cls.message['message']['migrating_info']['type']
- cls.guest_migrate_info.time_elapsed = cls.message['message']['migrating_info']['time_elapsed']
- cls.guest_migrate_info.time_remaining = cls.message['message']['migrating_info']['time_remaining']
- cls.guest_migrate_info.data_total = cls.message['message']['migrating_info']['data_total']
- cls.guest_migrate_info.data_processed = cls.message['message']['migrating_info']['data_processed']
- cls.guest_migrate_info.data_remaining = cls.message['message']['migrating_info']['data_remaining']
- cls.guest_migrate_info.mem_total = cls.message['message']['migrating_info']['mem_total']
- cls.guest_migrate_info.mem_processed = cls.message['message']['migrating_info']['mem_processed']
- cls.guest_migrate_info.mem_remaining = cls.message['message']['migrating_info']['mem_remaining']
- cls.guest_migrate_info.file_total = cls.message['message']['migrating_info']['file_total']
- cls.guest_migrate_info.file_processed = cls.message['message']['migrating_info']['file_processed']
- cls.guest_migrate_info.file_remaining = cls.message['message']['migrating_info']['file_remaining']
- cls.guest_migrate_info.update()
- except ji.PreviewingError as e:
- ret = json.loads(e.message)
- if ret['state']['code'] == '404':
- cls.guest_migrate_info.type = cls.message['message']['migrating_info']['type']
- cls.guest_migrate_info.time_elapsed = cls.message['message']['migrating_info']['time_elapsed']
- cls.guest_migrate_info.time_remaining = cls.message['message']['migrating_info']['time_remaining']
- cls.guest_migrate_info.data_total = cls.message['message']['migrating_info']['data_total']
- cls.guest_migrate_info.data_processed = cls.message['message']['migrating_info']['data_processed']
- cls.guest_migrate_info.data_remaining = cls.message['message']['migrating_info']['data_remaining']
- cls.guest_migrate_info.mem_total = cls.message['message']['migrating_info']['mem_total']
- cls.guest_migrate_info.mem_processed = cls.message['message']['migrating_info']['mem_processed']
- cls.guest_migrate_info.mem_remaining = cls.message['message']['migrating_info']['mem_remaining']
- cls.guest_migrate_info.file_total = cls.message['message']['migrating_info']['file_total']
- cls.guest_migrate_info.file_processed = cls.message['message']['migrating_info']['file_processed']
- cls.guest_migrate_info.file_remaining = cls.message['message']['migrating_info']['file_remaining']
- cls.guest_migrate_info.create()
- cls.guest.on_host = cls.message['host']
- cls.guest.status = cls.message['type']
- cls.guest.update()
- @classmethod
- def host_event_processor(cls):
- key = cls.message['message']['node_id']
- value = {
- 'hostname': cls.message['host'],
- 'timestamp': ji.Common.ts()
- }
- db.r.hset(app.config['hosts_info'], key=key, value=json.dumps(value, ensure_ascii=False))
- db.r.expire(app.config['hosts_info'], 10)
- @classmethod
- def response_processor(cls):
- action = cls.message['message']['action']
- uuid = cls.message['message']['uuid']
- state = cls.message['type']
- data = cls.message['message']['data']
- if action == 'create_guest':
- if state == ResponseState.success.value:
- cls.disk.uuid = uuid
- cls.disk.get_by('uuid')
- cls.disk.guest_uuid = uuid
- cls.disk.state = DiskState.mounted.value
- # disk_info['virtual-size'] 的单位为Byte,需要除以 1024 的 3 次方,换算成单位为 GB 的值
- cls.disk.size = data['disk_info']['virtual-size'] / (1024 ** 3)
- cls.disk.update()
- else:
- cls.guest.uuid = uuid
- cls.guest.get_by('uuid')
- cls.guest.status = GuestState.dirty.value
- cls.guest.update()
- elif action == 'migrate':
- pass
- elif action == 'delete_guest':
- if state == ResponseState.success.value:
- cls.config.id = 1
- cls.config.get()
- cls.guest.uuid = uuid
- cls.guest.get_by('uuid')
- if IP(cls.config.start_ip).int() <= IP(cls.guest.ip).int() <= IP(cls.config.end_ip).int():
- if db.r.srem(app.config['ip_used_set'], cls.guest.ip):
- db.r.sadd(app.config['ip_available_set'], cls.guest.ip)
- if (cls.guest.vnc_port - cls.config.start_vnc_port) <= \
- (IP(cls.config.end_ip).int() - IP(cls.config.start_ip).int()):
- if db.r.srem(app.config['vnc_port_used_set'], cls.guest.vnc_port):
- db.r.sadd(app.config['vnc_port_available_set'], cls.guest.vnc_port)
- cls.guest.delete()
- cls.disk.uuid = uuid
- cls.disk.get_by('uuid')
- cls.disk.delete()
- elif action == 'create_disk':
- cls.disk.uuid = uuid
- cls.disk.get_by('uuid')
- if state == ResponseState.success.value:
- cls.disk.state = DiskState.idle.value
- else:
- cls.disk.state = DiskState.dirty.value
- cls.disk.update()
- elif action == 'resize_disk':
- if state == ResponseState.success.value:
- cls.disk.uuid = uuid
- cls.disk.get_by('uuid')
- cls.disk.size = cls.message['message']['passback_parameters']['size']
- cls.disk.update()
- elif action == 'attach_disk':
- cls.disk.uuid = cls.message['message']['passback_parameters']['disk_uuid']
- cls.disk.get_by('uuid')
- if state == ResponseState.success.value:
- cls.disk.guest_uuid = uuid
- cls.disk.sequence = cls.message['message']['passback_parameters']['sequence']
- cls.disk.state = DiskState.mounted.value
- cls.disk.update()
- elif action == 'detach_disk':
- cls.disk.uuid = cls.message['message']['passback_parameters']['disk_uuid']
- cls.disk.get_by('uuid')
- if state == ResponseState.success.value:
- cls.disk.guest_uuid = ''
- cls.disk.sequence = -1
- cls.disk.state = DiskState.idle.value
- cls.disk.update()
- elif action == 'delete_disk':
- cls.disk.uuid = uuid
- cls.disk.get_by('uuid')
- cls.disk.delete()
- else:
- pass
- @classmethod
- def launch(cls):
- while True:
- if Utils.exit_flag:
- Utils.thread_counter -= 1
- print 'Thread EventProcessor say bye-bye'
- return
- try:
- report = db.r.lpop(app.config['upstream_queue'])
- if report is None:
- time.sleep(1)
- continue
- cls.message = json.loads(report)
- if cls.message['kind'] == EmitKind.log.value:
- cls.log_processor()
- elif cls.message['kind'] == EmitKind.guest_event.value:
- cls.guest_event_processor()
- elif cls.message['kind'] == EmitKind.host_event.value:
- cls.host_event_processor()
- elif cls.message['kind'] == EmitKind.response.value:
- cls.response_processor()
- else:
- pass
- except Exception as e:
- logger.error(e.message)
|