event_processor.py 4.1 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140
  1. #!/usr/bin/env python
  2. # -*- coding: utf-8 -*-
  3. import json
  4. import time
  5. from models import Database as db
  6. from models import Guest
  7. from models import GuestDisk
  8. from models import Log
  9. from models import Utils
  10. from models import EmitKind
  11. from models import ResponseState, GuestState, DiskState
  12. from models.initialize import app, logger
  13. __author__ = 'James Iter'
  14. __date__ = '2017/4/15'
  15. __contact__ = 'james.iter.cn@gmail.com'
  16. __copyright__ = '(c) 2017 by James Iter.'
  17. class EventProcessor(object):
  18. message = None
  19. log = Log()
  20. guest = Guest()
  21. disk = GuestDisk()
  22. @classmethod
  23. def log_processor(cls):
  24. cls.log.set(type=cls.message['type'], timestamp=cls.message['timestamp'], host=cls.message['host'],
  25. message=cls.message['message'])
  26. cls.log.create()
  27. @classmethod
  28. def event_processor(cls):
  29. cls.guest.uuid = cls.message['message']['uuid']
  30. cls.guest.get_by('uuid')
  31. cls.guest.status = cls.message['type']
  32. cls.guest.on_host = cls.message['host']
  33. cls.guest.update()
  34. @classmethod
  35. def response_processor(cls):
  36. action = cls.message['message']['action']
  37. uuid = cls.message['message']['uuid']
  38. state = cls.message['type']
  39. if action == 'create_vm':
  40. if state != ResponseState.success.value:
  41. cls.guest.uuid = uuid
  42. cls.guest.get_by('uuid')
  43. cls.guest.status = GuestState.dirty.value
  44. cls.guest.update()
  45. elif action == 'migrate':
  46. pass
  47. elif action == 'delete':
  48. if state == ResponseState.success.value:
  49. cls.guest.uuid = uuid
  50. cls.guest.get_by('uuid')
  51. cls.guest.delete()
  52. elif action == 'create_disk':
  53. cls.disk.uuid = uuid
  54. cls.disk.get_by('uuid')
  55. if state == ResponseState.success.value:
  56. cls.disk.state = DiskState.idle.value
  57. else:
  58. cls.disk.state = DiskState.dirty.value
  59. cls.disk.update()
  60. elif action == 'resize_disk':
  61. if state == ResponseState.success.value:
  62. cls.disk.uuid = uuid
  63. cls.disk.get_by('uuid')
  64. cls.disk.size = cls.message['message']['passback_parameters']['size']
  65. cls.disk.update()
  66. elif action == 'attach_disk':
  67. cls.disk.uuid = cls.message['message']['passback_parameters']['disk_uuid']
  68. cls.disk.get_by('uuid')
  69. if state == ResponseState.success.value:
  70. cls.disk.guest_uuid = uuid
  71. cls.disk.sequence = cls.message['message']['passback_parameters']['sequence']
  72. cls.disk.state = DiskState.mounted.value
  73. cls.disk.update()
  74. elif action == 'detach_disk':
  75. cls.disk.uuid = cls.message['message']['passback_parameters']['disk_uuid']
  76. cls.disk.get_by('uuid')
  77. if state == ResponseState.success.value:
  78. cls.disk.guest_uuid = ''
  79. cls.disk.sequence = -1
  80. cls.disk.state = DiskState.idle.value
  81. cls.disk.update()
  82. elif action == 'delete_disk':
  83. pass
  84. else:
  85. pass
  86. @classmethod
  87. def launch(cls):
  88. while True:
  89. if Utils.exit_flag:
  90. Utils.thread_counter -= 1
  91. print 'Thread EventProcessor say bye-bye'
  92. return
  93. try:
  94. report = db.r.lpop(app.config['upstream_queue'])
  95. if report is None:
  96. time.sleep(1)
  97. continue
  98. cls.message = json.loads(report)
  99. if cls.message['kind'] == EmitKind.log.value:
  100. cls.log_processor()
  101. elif cls.message['kind'] == EmitKind.event.value:
  102. cls.event_processor()
  103. elif cls.message['kind'] == EmitKind.response.value:
  104. cls.response_processor()
  105. else:
  106. pass
  107. except Exception as e:
  108. logger.error(e.message)