event_processor.py 6.0 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185
  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
  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.initialize import app, logger
  15. __author__ = 'James Iter'
  16. __date__ = '2017/4/15'
  17. __contact__ = 'james.iter.cn@gmail.com'
  18. __copyright__ = '(c) 2017 by James Iter.'
  19. class EventProcessor(object):
  20. message = None
  21. log = Log()
  22. guest = Guest()
  23. disk = Disk()
  24. config = Config()
  25. @classmethod
  26. def log_processor(cls):
  27. cls.log.set(type=cls.message['type'], timestamp=cls.message['timestamp'], host=cls.message['host'],
  28. message=cls.message['message'])
  29. cls.log.create()
  30. @classmethod
  31. def guest_event_processor(cls):
  32. cls.guest.uuid = cls.message['message']['uuid']
  33. cls.guest.get_by('uuid')
  34. cls.guest.status = cls.message['type']
  35. cls.guest.on_host = cls.message['host']
  36. cls.guest.update()
  37. @classmethod
  38. def host_event_processor(cls):
  39. key = cls.message['message']['node_id']
  40. value = {
  41. 'hostname': cls.message['host'],
  42. 'timestamp': ji.Common.ts()
  43. }
  44. db.r.hset(app.config['hosts_info'], key=key, value=json.dumps(value, ensure_ascii=False))
  45. db.r.expire(app.config['hosts_info'], 10)
  46. @classmethod
  47. def response_processor(cls):
  48. action = cls.message['message']['action']
  49. uuid = cls.message['message']['uuid']
  50. state = cls.message['type']
  51. data = cls.message['message']['data']
  52. if action == 'create_guest':
  53. if state == ResponseState.success.value:
  54. cls.disk.uuid = uuid
  55. cls.disk.get_by('uuid')
  56. cls.disk.guest_uuid = uuid
  57. cls.disk.state = DiskState.mounted.value
  58. # disk_info['virtual-size'] 的单位为Byte,需要除以 1024 的 3 次方,换算成单位为 GB 的值
  59. cls.disk.size = data['disk_info']['virtual-size'] / (1024 ** 3)
  60. cls.disk.update()
  61. else:
  62. cls.guest.uuid = uuid
  63. cls.guest.get_by('uuid')
  64. cls.guest.status = GuestState.dirty.value
  65. cls.guest.update()
  66. elif action == 'migrate':
  67. pass
  68. elif action == 'delete_guest':
  69. if state == ResponseState.success.value:
  70. cls.config.id = 1
  71. cls.config.get()
  72. cls.guest.uuid = uuid
  73. cls.guest.get_by('uuid')
  74. if IP(cls.config.start_ip).int() <= IP(cls.guest.ip).int() <= IP(cls.config.end_ip).int():
  75. if db.r.srem(app.config['ip_used_set'], cls.guest.ip):
  76. db.r.sadd(app.config['ip_available_set'], cls.guest.ip)
  77. if (cls.guest.vnc_port - cls.config.start_vnc_port) <= \
  78. (IP(cls.config.end_ip).int() - IP(cls.config.start_ip).int()):
  79. if db.r.srem(app.config['vnc_port_used_set'], cls.guest.vnc_port):
  80. db.r.sadd(app.config['vnc_port_available_set'], cls.guest.vnc_port)
  81. cls.guest.delete()
  82. cls.disk.uuid = uuid
  83. cls.disk.get_by('uuid')
  84. cls.disk.delete()
  85. elif action == 'create_disk':
  86. cls.disk.uuid = uuid
  87. cls.disk.get_by('uuid')
  88. if state == ResponseState.success.value:
  89. cls.disk.state = DiskState.idle.value
  90. else:
  91. cls.disk.state = DiskState.dirty.value
  92. cls.disk.update()
  93. elif action == 'resize_disk':
  94. if state == ResponseState.success.value:
  95. cls.disk.uuid = uuid
  96. cls.disk.get_by('uuid')
  97. cls.disk.size = cls.message['message']['passback_parameters']['size']
  98. cls.disk.update()
  99. elif action == 'attach_disk':
  100. cls.disk.uuid = cls.message['message']['passback_parameters']['disk_uuid']
  101. cls.disk.get_by('uuid')
  102. if state == ResponseState.success.value:
  103. cls.disk.guest_uuid = uuid
  104. cls.disk.sequence = cls.message['message']['passback_parameters']['sequence']
  105. cls.disk.state = DiskState.mounted.value
  106. cls.disk.update()
  107. elif action == 'detach_disk':
  108. cls.disk.uuid = cls.message['message']['passback_parameters']['disk_uuid']
  109. cls.disk.get_by('uuid')
  110. if state == ResponseState.success.value:
  111. cls.disk.guest_uuid = ''
  112. cls.disk.sequence = -1
  113. cls.disk.state = DiskState.idle.value
  114. cls.disk.update()
  115. elif action == 'delete_disk':
  116. cls.disk.uuid = uuid
  117. cls.disk.get_by('uuid')
  118. cls.disk.delete()
  119. else:
  120. pass
  121. @classmethod
  122. def launch(cls):
  123. while True:
  124. if Utils.exit_flag:
  125. Utils.thread_counter -= 1
  126. print 'Thread EventProcessor say bye-bye'
  127. return
  128. try:
  129. report = db.r.lpop(app.config['upstream_queue'])
  130. if report is None:
  131. time.sleep(1)
  132. continue
  133. cls.message = json.loads(report)
  134. if cls.message['kind'] == EmitKind.log.value:
  135. cls.log_processor()
  136. elif cls.message['kind'] == EmitKind.guest_event.value:
  137. cls.guest_event_processor()
  138. elif cls.message['kind'] == EmitKind.host_event.value:
  139. cls.host_event_processor()
  140. elif cls.message['kind'] == EmitKind.response.value:
  141. cls.response_processor()
  142. else:
  143. pass
  144. except Exception as e:
  145. logger.error(e.message)