event_processor.py 5.4 KB

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