disk.py 12 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439
  1. #!/usr/bin/env python
  2. # -*- coding: utf-8 -*-
  3. from flask import Blueprint, request
  4. import json
  5. from uuid import uuid4
  6. import jimit as ji
  7. from models import Guest, DiskState, Host
  8. from models.initialize import dev_table
  9. from models import Config
  10. from models import Disk
  11. from models import Rules
  12. from models import Utils
  13. from models.status import StorageMode
  14. from base import Base
  15. __author__ = 'James Iter'
  16. __date__ = '2017/4/24'
  17. __contact__ = 'james.iter.cn@gmail.com'
  18. __copyright__ = '(c) 2017 by James Iter.'
  19. blueprint = Blueprint(
  20. 'api_disk',
  21. __name__,
  22. url_prefix='/api/disk'
  23. )
  24. blueprints = Blueprint(
  25. 'api_disks',
  26. __name__,
  27. url_prefix='/api/disks'
  28. )
  29. disk_base = Base(the_class=Disk, the_blueprint=blueprint, the_blueprints=blueprints)
  30. @Utils.dumps2response
  31. def r_create():
  32. args_rules = [
  33. Rules.DISK_SIZE.value,
  34. Rules.REMARK.value,
  35. Rules.QUANTITY.value
  36. ]
  37. config = Config()
  38. config.id = 1
  39. config.get()
  40. # 非共享模式,必须指定 node_id
  41. if config.storage_mode not in [StorageMode.shared_mount.value, StorageMode.ceph.value,
  42. StorageMode.glusterfs.value]:
  43. args_rules.append(
  44. Rules.NODE_ID.value
  45. )
  46. try:
  47. ji.Check.previewing(args_rules, request.json)
  48. size = request.json['size']
  49. quantity = request.json['quantity']
  50. ret = dict()
  51. ret['state'] = ji.Common.exchange_state(20000)
  52. # 如果是共享模式,则让负载最轻的计算节点去创建磁盘
  53. if config.storage_mode in [StorageMode.shared_mount.value, StorageMode.ceph.value,
  54. StorageMode.glusterfs.value]:
  55. available_hosts = Host.get_available_hosts()
  56. if available_hosts.__len__() == 0:
  57. ret['state'] = ji.Common.exchange_state(50351)
  58. return ret
  59. # 在可用计算节点中平均分配任务
  60. chosen_host = available_hosts[quantity % available_hosts.__len__()]
  61. request.json['node_id'] = chosen_host['node_id']
  62. node_id = request.json['node_id']
  63. if size < 1:
  64. ret['state'] = ji.Common.exchange_state(41255)
  65. return ret
  66. while quantity:
  67. quantity -= 1
  68. disk = Disk()
  69. disk.guest_uuid = ''
  70. disk.size = size
  71. disk.uuid = uuid4().__str__()
  72. disk.remark = request.json.get('remark', '')
  73. disk.node_id = int(node_id)
  74. disk.sequence = -1
  75. disk.format = 'qcow2'
  76. disk.path = config.storage_path + '/' + disk.uuid + '.' + disk.format
  77. disk.quota(config=config)
  78. message = {
  79. '_object': 'disk',
  80. 'action': 'create',
  81. 'uuid': disk.uuid,
  82. 'storage_mode': config.storage_mode,
  83. 'dfs_volume': config.dfs_volume,
  84. 'node_id': disk.node_id,
  85. 'image_path': disk.path,
  86. 'size': disk.size
  87. }
  88. Utils.emit_instruction(message=json.dumps(message, ensure_ascii=False))
  89. disk.create()
  90. return ret
  91. except ji.PreviewingError, e:
  92. return json.loads(e.message)
  93. @Utils.dumps2response
  94. def r_resize(uuid, size):
  95. args_rules = [
  96. Rules.UUID.value,
  97. Rules.DISK_SIZE_STR.value
  98. ]
  99. try:
  100. ji.Check.previewing(args_rules, {'uuid': uuid, 'size': size})
  101. disk = Disk()
  102. disk.uuid = uuid
  103. disk.get_by('uuid')
  104. ret = dict()
  105. ret['state'] = ji.Common.exchange_state(20000)
  106. if disk.size >= int(size):
  107. ret['state'] = ji.Common.exchange_state(41257)
  108. return ret
  109. config = Config()
  110. config.id = 1
  111. config.get()
  112. disk.size = int(size)
  113. disk.quota(config=config)
  114. # 将在事件返回层(models/event_processor.py:224 附近),更新数据库中 disk 对象
  115. message = {
  116. '_object': 'disk',
  117. 'action': 'resize',
  118. 'uuid': disk.uuid,
  119. 'guest_uuid': disk.guest_uuid,
  120. 'storage_mode': config.storage_mode,
  121. 'size': disk.size,
  122. 'dfs_volume': config.dfs_volume,
  123. 'node_id': disk.node_id,
  124. 'image_path': disk.path,
  125. 'disks': [disk.__dict__],
  126. 'passback_parameters': {'size': disk.size}
  127. }
  128. if config.storage_mode in [StorageMode.shared_mount.value, StorageMode.ceph.value,
  129. StorageMode.glusterfs.value]:
  130. message['node_id'] = Host.get_lightest_host()['node_id']
  131. if disk.guest_uuid.__len__() == 36:
  132. message['device_node'] = dev_table[disk.sequence]
  133. Utils.emit_instruction(message=json.dumps(message, ensure_ascii=False))
  134. return ret
  135. except ji.PreviewingError, e:
  136. return json.loads(e.message)
  137. @Utils.dumps2response
  138. def r_delete(uuids):
  139. args_rules = [
  140. Rules.UUIDS.value
  141. ]
  142. try:
  143. ji.Check.previewing(args_rules, {'uuids': uuids})
  144. ret = dict()
  145. ret['state'] = ji.Common.exchange_state(20000)
  146. disk = Disk()
  147. # 检测所指定的 UUDIs 磁盘都存在
  148. for uuid in uuids.split(','):
  149. disk.uuid = uuid
  150. disk.get_by('uuid')
  151. # 判断磁盘是否与虚拟机处于离状态
  152. if disk.state != DiskState.idle.value:
  153. ret['state'] = ji.Common.exchange_state(41256)
  154. return ret
  155. config = Config()
  156. config.id = 1
  157. config.get()
  158. # 执行删除操作
  159. for uuid in uuids.split(','):
  160. disk.uuid = uuid
  161. disk.get_by('uuid')
  162. message = {
  163. '_object': 'disk',
  164. 'action': 'delete',
  165. 'uuid': disk.uuid,
  166. 'storage_mode': config.storage_mode,
  167. 'dfs_volume': config.dfs_volume,
  168. 'node_id': disk.node_id,
  169. 'image_path': disk.path
  170. }
  171. if config.storage_mode in [StorageMode.shared_mount.value, StorageMode.ceph.value,
  172. StorageMode.glusterfs.value]:
  173. message['node_id'] = Host.get_lightest_host()['node_id']
  174. Utils.emit_instruction(message=json.dumps(message, ensure_ascii=False))
  175. return ret
  176. except ji.PreviewingError, e:
  177. return json.loads(e.message)
  178. def add_device(func):
  179. from functools import wraps
  180. @wraps(func)
  181. def _add_device(*args, **kwargs):
  182. ret = func(*args, **kwargs)
  183. if ret['data'].__len__() > 0:
  184. if isinstance(ret['data'], list):
  185. for i, item in enumerate(ret['data']):
  186. ret['data'][i][u'device'] = u'/dev/' + dev_table[item['sequence']]
  187. if item['sequence'] < 0:
  188. ret['data'][i][u'device'] = None
  189. elif isinstance(ret['data'], dict):
  190. ret['data'][u'device'] = u'/dev/' + dev_table[ret['data']['sequence']]
  191. if ret['data']['sequence'] < 0:
  192. ret['data'][u'device'] = None
  193. else:
  194. raise json.dumps(ret)
  195. return ret
  196. return _add_device
  197. @Utils.dumps2response
  198. @add_device
  199. def r_get(uuids):
  200. return disk_base.get(ids=uuids, ids_rule=Rules.UUIDS.value, by_field='uuid')
  201. @Utils.dumps2response
  202. @add_device
  203. def r_get_by_filter():
  204. return disk_base.get_by_filter()
  205. @Utils.dumps2response
  206. @add_device
  207. def r_content_search():
  208. return disk_base.content_search()
  209. @Utils.dumps2response
  210. def r_update(uuids):
  211. ret = dict()
  212. ret['state'] = ji.Common.exchange_state(20000)
  213. ret['data'] = list()
  214. args_rules = [
  215. Rules.UUIDS.value
  216. ]
  217. if 'remark' in request.json:
  218. args_rules.append(
  219. Rules.REMARK.value
  220. )
  221. if 'iops' in request.json:
  222. args_rules.append(
  223. Rules.IOPS.value
  224. )
  225. if 'iops_rd' in request.json:
  226. args_rules.append(
  227. Rules.IOPS_RD.value
  228. )
  229. if 'iops_wr' in request.json:
  230. args_rules.append(
  231. Rules.IOPS_WR.value
  232. )
  233. if 'iops_max' in request.json:
  234. args_rules.append(
  235. Rules.IOPS_MAX.value
  236. )
  237. if 'iops_max_length' in request.json:
  238. args_rules.append(
  239. Rules.IOPS_MAX_LENGTH.value
  240. )
  241. if 'bps' in request.json:
  242. args_rules.append(
  243. Rules.BPS.value
  244. )
  245. if 'bps_rd' in request.json:
  246. args_rules.append(
  247. Rules.BPS_RD.value
  248. )
  249. if 'bps_wr' in request.json:
  250. args_rules.append(
  251. Rules.BPS_WR.value
  252. )
  253. if 'bps_max' in request.json:
  254. args_rules.append(
  255. Rules.BPS_MAX.value
  256. )
  257. if 'bps_max_length' in request.json:
  258. args_rules.append(
  259. Rules.BPS_MAX_LENGTH.value
  260. )
  261. if args_rules.__len__() < 2:
  262. return ret
  263. request.json['uuids'] = uuids
  264. need_update_quota = False
  265. need_update_quota_parameters = ['iops', 'iops_rd', 'iops_wr', 'iops_max', 'iops_max_length',
  266. 'bps', 'bps_rd', 'bps_wr', 'bps_max', 'bps_max_length']
  267. if filter(lambda p: p in request.json, need_update_quota_parameters).__len__() > 0:
  268. need_update_quota = True
  269. try:
  270. ji.Check.previewing(args_rules, request.json)
  271. disk = Disk()
  272. # 检测所指定的 UUDIs 磁盘都存在
  273. for uuid in uuids.split(','):
  274. disk.uuid = uuid
  275. disk.get_by('uuid')
  276. for uuid in uuids.split(','):
  277. disk.uuid = uuid
  278. disk.get_by('uuid')
  279. disk.remark = request.json.get('remark', disk.remark)
  280. disk.iops = request.json.get('iops', disk.iops)
  281. disk.iops_rd = request.json.get('iops_rd', disk.iops_rd)
  282. disk.iops_wr = request.json.get('iops_wr', disk.iops_wr)
  283. disk.iops_max = request.json.get('iops_max', disk.iops_max)
  284. disk.iops_max_length = request.json.get('iops_max_length', disk.iops_max_length)
  285. disk.bps = request.json.get('bps', disk.bps)
  286. disk.bps_rd = request.json.get('bps_rd', disk.bps_rd)
  287. disk.bps_wr = request.json.get('bps_wr', disk.bps_wr)
  288. disk.bps_max = request.json.get('bps_max', disk.bps_max)
  289. disk.bps_max_length = request.json.get('bps_max_length', disk.bps_max_length)
  290. disk.update()
  291. disk.get()
  292. if disk.sequence >= 0 and need_update_quota:
  293. message = {
  294. '_object': 'disk',
  295. 'action': 'quota',
  296. 'uuid': disk.uuid,
  297. 'guest_uuid': disk.guest_uuid,
  298. 'node_id': disk.node_id,
  299. 'disks': [disk.__dict__]
  300. }
  301. Utils.emit_instruction(message=json.dumps(message))
  302. ret['data'].append(disk.__dict__)
  303. return ret
  304. except ji.PreviewingError, e:
  305. return json.loads(e.message)
  306. @Utils.dumps2response
  307. def r_distribute_count():
  308. from models import Disk
  309. rows, count = Disk.get_all()
  310. ret = dict()
  311. ret['state'] = ji.Common.exchange_state(20000)
  312. ret['data'] = {
  313. 'kind': {'system': 0, 'data_mounted': 0, 'data_idle': 0},
  314. 'total_size': 0,
  315. 'disks': rows.__len__()
  316. }
  317. for disk in rows:
  318. if disk['sequence'] == 0:
  319. ret['data']['kind']['system'] += 1
  320. elif disk['sequence'] < 0:
  321. ret['data']['kind']['data_idle'] += 1
  322. else:
  323. ret['data']['kind']['data_mounted'] += 1
  324. ret['data']['total_size'] += disk['size']
  325. return ret