disk.py 17 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613
  1. #!/usr/bin/env python
  2. # -*- coding: utf-8 -*-
  3. from math import ceil
  4. from flask import Blueprint, request, url_for
  5. import json
  6. import requests
  7. from uuid import uuid4
  8. import jimit as ji
  9. from models import Guest, DiskState, Host
  10. from models.initialize import dev_table
  11. from models import Config
  12. from models import Disk
  13. from models import Rules
  14. from models import Utils
  15. from models.status import StorageMode
  16. from base import Base
  17. __author__ = 'James Iter'
  18. __date__ = '2017/4/24'
  19. __contact__ = 'james.iter.cn@gmail.com'
  20. __copyright__ = '(c) 2017 by James Iter.'
  21. blueprint = Blueprint(
  22. 'api_disk',
  23. __name__,
  24. url_prefix='/api/disk'
  25. )
  26. blueprints = Blueprint(
  27. 'api_disks',
  28. __name__,
  29. url_prefix='/api/disks'
  30. )
  31. disk_base = Base(the_class=Disk, the_blueprint=blueprint, the_blueprints=blueprints)
  32. @Utils.dumps2response
  33. def r_create():
  34. args_rules = [
  35. Rules.DISK_SIZE.value,
  36. Rules.REMARK.value,
  37. Rules.QUANTITY.value
  38. ]
  39. config = Config()
  40. config.id = 1
  41. config.get()
  42. # 非共享模式,必须指定 node_id
  43. if config.storage_mode not in [StorageMode.shared_mount.value, StorageMode.ceph.value,
  44. StorageMode.glusterfs.value]:
  45. args_rules.append(
  46. Rules.NODE_ID.value
  47. )
  48. try:
  49. ji.Check.previewing(args_rules, request.json)
  50. size = request.json['size']
  51. quantity = request.json['quantity']
  52. ret = dict()
  53. ret['state'] = ji.Common.exchange_state(20000)
  54. # 如果是共享模式,则让负载最轻的计算节点去创建磁盘
  55. if config.storage_mode in [StorageMode.shared_mount.value, StorageMode.ceph.value,
  56. StorageMode.glusterfs.value]:
  57. available_hosts = Host.get_available_hosts()
  58. if available_hosts.__len__() == 0:
  59. ret['state'] = ji.Common.exchange_state(50351)
  60. return ret
  61. # 在可用计算节点中平均分配任务
  62. chosen_host = available_hosts[quantity % available_hosts.__len__()]
  63. request.json['node_id'] = chosen_host['node_id']
  64. node_id = request.json['node_id']
  65. if size < 1:
  66. ret['state'] = ji.Common.exchange_state(41255)
  67. return ret
  68. while quantity:
  69. quantity -= 1
  70. disk = Disk()
  71. disk.guest_uuid = ''
  72. disk.size = size
  73. disk.uuid = uuid4().__str__()
  74. disk.remark = request.json.get('remark', '')
  75. disk.node_id = int(node_id)
  76. disk.sequence = -1
  77. disk.format = 'qcow2'
  78. disk.path = config.storage_path + '/' + disk.uuid + '.' + disk.format
  79. disk.quota(config=config)
  80. message = {
  81. '_object': 'disk',
  82. 'action': 'create',
  83. 'uuid': disk.uuid,
  84. 'storage_mode': config.storage_mode,
  85. 'dfs_volume': config.dfs_volume,
  86. 'node_id': disk.node_id,
  87. 'image_path': disk.path,
  88. 'size': disk.size
  89. }
  90. Utils.emit_instruction(message=json.dumps(message, ensure_ascii=False))
  91. disk.create()
  92. return ret
  93. except ji.PreviewingError, e:
  94. return json.loads(e.message)
  95. @Utils.dumps2response
  96. def r_resize(uuid, size):
  97. args_rules = [
  98. Rules.UUID.value,
  99. Rules.DISK_SIZE_STR.value
  100. ]
  101. try:
  102. ji.Check.previewing(args_rules, {'uuid': uuid, 'size': size})
  103. disk = Disk()
  104. disk.uuid = uuid
  105. disk.get_by('uuid')
  106. ret = dict()
  107. ret['state'] = ji.Common.exchange_state(20000)
  108. if disk.size >= int(size):
  109. ret['state'] = ji.Common.exchange_state(41257)
  110. return ret
  111. config = Config()
  112. config.id = 1
  113. config.get()
  114. disk.size = int(size)
  115. disk.quota(config=config)
  116. # 将在事件返回层(models/event_processor.py:224 附近),更新数据库中 disk 对象
  117. message = {
  118. '_object': 'disk',
  119. 'action': 'resize',
  120. 'uuid': disk.uuid,
  121. 'guest_uuid': disk.guest_uuid,
  122. 'storage_mode': config.storage_mode,
  123. 'size': disk.size,
  124. 'dfs_volume': config.dfs_volume,
  125. 'node_id': disk.node_id,
  126. 'image_path': disk.path,
  127. 'disks': [disk.__dict__],
  128. 'passback_parameters': {'size': disk.size}
  129. }
  130. if config.storage_mode in [StorageMode.shared_mount.value, StorageMode.ceph.value,
  131. StorageMode.glusterfs.value]:
  132. message['node_id'] = Host.get_lightest_host()['node_id']
  133. if disk.guest_uuid.__len__() == 36:
  134. message['device_node'] = dev_table[disk.sequence]
  135. Utils.emit_instruction(message=json.dumps(message, ensure_ascii=False))
  136. return ret
  137. except ji.PreviewingError, e:
  138. return json.loads(e.message)
  139. @Utils.dumps2response
  140. def r_delete(uuids):
  141. args_rules = [
  142. Rules.UUIDS.value
  143. ]
  144. try:
  145. ji.Check.previewing(args_rules, {'uuids': uuids})
  146. ret = dict()
  147. ret['state'] = ji.Common.exchange_state(20000)
  148. disk = Disk()
  149. # 检测所指定的 UUDIs 磁盘都存在
  150. for uuid in uuids.split(','):
  151. disk.uuid = uuid
  152. disk.get_by('uuid')
  153. # 判断磁盘是否与虚拟机处于离状态
  154. if disk.state not in [DiskState.idle.value, DiskState.dirty.value]:
  155. ret['state'] = ji.Common.exchange_state(41256)
  156. return ret
  157. config = Config()
  158. config.id = 1
  159. config.get()
  160. # 执行删除操作
  161. for uuid in uuids.split(','):
  162. disk.uuid = uuid
  163. disk.get_by('uuid')
  164. message = {
  165. '_object': 'disk',
  166. 'action': 'delete',
  167. 'uuid': disk.uuid,
  168. 'storage_mode': config.storage_mode,
  169. 'dfs_volume': config.dfs_volume,
  170. 'node_id': disk.node_id,
  171. 'image_path': disk.path
  172. }
  173. if config.storage_mode in [StorageMode.shared_mount.value, StorageMode.ceph.value,
  174. StorageMode.glusterfs.value]:
  175. message['node_id'] = Host.get_lightest_host()['node_id']
  176. Utils.emit_instruction(message=json.dumps(message, ensure_ascii=False))
  177. return ret
  178. except ji.PreviewingError, e:
  179. return json.loads(e.message)
  180. def add_device(func):
  181. from functools import wraps
  182. @wraps(func)
  183. def _add_device(*args, **kwargs):
  184. ret = func(*args, **kwargs)
  185. if ret['data'].__len__() > 0:
  186. if isinstance(ret['data'], list):
  187. for i, item in enumerate(ret['data']):
  188. ret['data'][i][u'device'] = u'/dev/' + dev_table[item['sequence']]
  189. if item['sequence'] < 0:
  190. ret['data'][i][u'device'] = None
  191. elif isinstance(ret['data'], dict):
  192. ret['data'][u'device'] = u'/dev/' + dev_table[ret['data']['sequence']]
  193. if ret['data']['sequence'] < 0:
  194. ret['data'][u'device'] = None
  195. else:
  196. raise json.dumps(ret)
  197. return ret
  198. return _add_device
  199. @Utils.dumps2response
  200. @add_device
  201. def r_get(uuids):
  202. return disk_base.get(ids=uuids, ids_rule=Rules.UUIDS.value, by_field='uuid')
  203. @Utils.dumps2response
  204. @add_device
  205. def r_get_by_filter():
  206. return disk_base.get_by_filter()
  207. @Utils.dumps2response
  208. @add_device
  209. def r_content_search():
  210. return disk_base.content_search()
  211. @Utils.dumps2response
  212. def r_update(uuids):
  213. ret = dict()
  214. ret['state'] = ji.Common.exchange_state(20000)
  215. ret['data'] = list()
  216. args_rules = [
  217. Rules.UUIDS.value
  218. ]
  219. if 'remark' in request.json:
  220. args_rules.append(
  221. Rules.REMARK.value
  222. )
  223. if 'iops' in request.json:
  224. args_rules.append(
  225. Rules.IOPS.value
  226. )
  227. if 'iops_rd' in request.json:
  228. args_rules.append(
  229. Rules.IOPS_RD.value
  230. )
  231. if 'iops_wr' in request.json:
  232. args_rules.append(
  233. Rules.IOPS_WR.value
  234. )
  235. if 'iops_max' in request.json:
  236. args_rules.append(
  237. Rules.IOPS_MAX.value
  238. )
  239. if 'iops_max_length' in request.json:
  240. args_rules.append(
  241. Rules.IOPS_MAX_LENGTH.value
  242. )
  243. if 'bps' in request.json:
  244. args_rules.append(
  245. Rules.BPS.value
  246. )
  247. if 'bps_rd' in request.json:
  248. args_rules.append(
  249. Rules.BPS_RD.value
  250. )
  251. if 'bps_wr' in request.json:
  252. args_rules.append(
  253. Rules.BPS_WR.value
  254. )
  255. if 'bps_max' in request.json:
  256. args_rules.append(
  257. Rules.BPS_MAX.value
  258. )
  259. if 'bps_max_length' in request.json:
  260. args_rules.append(
  261. Rules.BPS_MAX_LENGTH.value
  262. )
  263. if args_rules.__len__() < 2:
  264. return ret
  265. request.json['uuids'] = uuids
  266. need_update_quota = False
  267. need_update_quota_parameters = ['iops', 'iops_rd', 'iops_wr', 'iops_max', 'iops_max_length',
  268. 'bps', 'bps_rd', 'bps_wr', 'bps_max', 'bps_max_length']
  269. if filter(lambda p: p in request.json, need_update_quota_parameters).__len__() > 0:
  270. need_update_quota = True
  271. try:
  272. ji.Check.previewing(args_rules, request.json)
  273. disk = Disk()
  274. # 检测所指定的 UUDIs 磁盘都存在
  275. for uuid in uuids.split(','):
  276. disk.uuid = uuid
  277. disk.get_by('uuid')
  278. for uuid in uuids.split(','):
  279. disk.uuid = uuid
  280. disk.get_by('uuid')
  281. disk.remark = request.json.get('remark', disk.remark)
  282. disk.iops = request.json.get('iops', disk.iops)
  283. disk.iops_rd = request.json.get('iops_rd', disk.iops_rd)
  284. disk.iops_wr = request.json.get('iops_wr', disk.iops_wr)
  285. disk.iops_max = request.json.get('iops_max', disk.iops_max)
  286. disk.iops_max_length = request.json.get('iops_max_length', disk.iops_max_length)
  287. disk.bps = request.json.get('bps', disk.bps)
  288. disk.bps_rd = request.json.get('bps_rd', disk.bps_rd)
  289. disk.bps_wr = request.json.get('bps_wr', disk.bps_wr)
  290. disk.bps_max = request.json.get('bps_max', disk.bps_max)
  291. disk.bps_max_length = request.json.get('bps_max_length', disk.bps_max_length)
  292. disk.update()
  293. disk.get()
  294. if disk.sequence >= 0 and need_update_quota:
  295. message = {
  296. '_object': 'disk',
  297. 'action': 'quota',
  298. 'uuid': disk.uuid,
  299. 'guest_uuid': disk.guest_uuid,
  300. 'node_id': disk.node_id,
  301. 'disks': [disk.__dict__]
  302. }
  303. Utils.emit_instruction(message=json.dumps(message))
  304. ret['data'].append(disk.__dict__)
  305. return ret
  306. except ji.PreviewingError, e:
  307. return json.loads(e.message)
  308. @Utils.dumps2response
  309. def r_distribute_count():
  310. from models import Disk
  311. rows, count = Disk.get_all()
  312. ret = dict()
  313. ret['state'] = ji.Common.exchange_state(20000)
  314. ret['data'] = {
  315. 'kind': {'system': 0, 'data_mounted': 0, 'data_idle': 0},
  316. 'total_size': 0,
  317. 'disks': rows.__len__()
  318. }
  319. for disk in rows:
  320. if disk['sequence'] == 0:
  321. ret['data']['kind']['system'] += 1
  322. elif disk['sequence'] < 0:
  323. ret['data']['kind']['data_idle'] += 1
  324. else:
  325. ret['data']['kind']['data_mounted'] += 1
  326. ret['data']['total_size'] += disk['size']
  327. return ret
  328. @Utils.dumps2response
  329. def r_show():
  330. args = list()
  331. page = request.args.get('page', 1)
  332. if page == '':
  333. page = 1
  334. page = int(page)
  335. page_size = int(request.args.get('page_size', 10))
  336. keyword = request.args.get('keyword', None)
  337. show_area = request.args.get('show_area', 'unmount')
  338. guest_uuid = request.args.get('guest_uuid', None)
  339. sequence = request.args.get('sequence', None)
  340. order_by = request.args.get('order_by', None)
  341. order = request.args.get('order', None)
  342. filters = list()
  343. if page is not None:
  344. args.append('page=' + page.__str__())
  345. if page_size is not None:
  346. args.append('page_size=' + page_size.__str__())
  347. if keyword is not None:
  348. args.append('keyword=' + keyword.__str__())
  349. if guest_uuid is not None:
  350. filters.append('guest_uuid:in:' + guest_uuid.__str__())
  351. show_area = 'all'
  352. if sequence is not None:
  353. filters.append('sequence:in:' + sequence.__str__())
  354. show_area = 'all'
  355. if show_area in ['unmount', 'data_disk', 'all']:
  356. if show_area == 'unmount':
  357. filters.append('sequence:eq:-1')
  358. elif show_area == 'data_disk':
  359. filters.append('sequence:gt:0')
  360. else:
  361. pass
  362. else:
  363. # 与前端页面相照应,首次打开时,默认只显示未挂载的磁盘
  364. filters.append('sequence:eq:-1')
  365. if order_by is not None:
  366. args.append('order_by=' + order_by)
  367. if order is not None:
  368. args.append('order=' + order)
  369. if filters.__len__() > 0:
  370. args.append('filter=' + ','.join(filters))
  371. hosts_url = url_for('api_hosts.r_get_by_filter', _external=True)
  372. disks_url = url_for('api_disks.r_get_by_filter', _external=True)
  373. if keyword is not None:
  374. disks_url = url_for('api_disks.r_content_search', _external=True)
  375. # 关键字检索,不支持显示域过滤
  376. show_area = 'all'
  377. hosts_ret = requests.get(url=hosts_url, cookies=request.cookies)
  378. hosts_ret = json.loads(hosts_ret.content)
  379. hosts_mapping_by_node_id = dict()
  380. for host in hosts_ret['data']:
  381. hosts_mapping_by_node_id[int(host['node_id'])] = host
  382. if args.__len__() > 0:
  383. disks_url = disks_url + '?' + '&'.join(args)
  384. disks_ret = requests.get(url=disks_url, cookies=request.cookies)
  385. disks_ret = json.loads(disks_ret.content)
  386. guests_uuid = list()
  387. disks_uuid = list()
  388. for disk in disks_ret['data']:
  389. disks_uuid.append(disk['uuid'])
  390. if disk['guest_uuid'].__len__() == 36:
  391. guests_uuid.append(disk['guest_uuid'])
  392. if guests_uuid.__len__() > 0:
  393. guests, _ = Guest.get_by_filter(filter_str='uuid:in:' + ','.join(guests_uuid))
  394. guests_uuid_mapping = dict()
  395. for guest in guests:
  396. guests_uuid_mapping[guest['uuid']] = guest
  397. for i, disk in enumerate(disks_ret['data']):
  398. if disk['guest_uuid'].__len__() == 36:
  399. disks_ret['data'][i]['guest'] = guests_uuid_mapping[disk['guest_uuid']]
  400. if disks_uuid.__len__() > 0:
  401. snapshots_id_mapping_by_disks_uuid_url = url_for('api_snapshots.r_get_snapshots_by_disks_uuid',
  402. disks_uuid=','.join(disks_uuid), _external=True)
  403. snapshots_id_mapping_by_disks_uuid_ret = requests.get(url=snapshots_id_mapping_by_disks_uuid_url,
  404. cookies=request.cookies)
  405. snapshots_id_mapping_by_disks_uuid_ret = json.loads(snapshots_id_mapping_by_disks_uuid_ret.content)
  406. snapshots_id_mapping_by_disk_uuid = dict()
  407. for snapshot_id_mapping_by_disk_uuid in snapshots_id_mapping_by_disks_uuid_ret['data']:
  408. disk_uuid = snapshot_id_mapping_by_disk_uuid['disk_uuid']
  409. snapshot_id = snapshot_id_mapping_by_disk_uuid['snapshot_id']
  410. if disk_uuid not in snapshots_id_mapping_by_disk_uuid:
  411. snapshots_id_mapping_by_disk_uuid[disk_uuid] = list()
  412. snapshots_id_mapping_by_disk_uuid[disk_uuid].append(snapshot_id)
  413. for i, disk in enumerate(disks_ret['data']):
  414. if disk['uuid'] in snapshots_id_mapping_by_disk_uuid:
  415. disks_ret['data'][i]['snapshot'] = snapshots_id_mapping_by_disk_uuid[disk['uuid']]
  416. config = Config()
  417. config.id = 1
  418. config.get()
  419. show_on_host = False
  420. if config.storage_mode == StorageMode.local.value:
  421. show_on_host = True
  422. last_page = int(ceil(disks_ret['paging']['total'] / float(page_size)))
  423. page_length = 5
  424. pages = list()
  425. if page < int(ceil(page_length / 2.0)):
  426. for i in range(1, page_length + 1):
  427. pages.append(i)
  428. if i == last_page or last_page == 0:
  429. break
  430. elif last_page - page < page_length / 2:
  431. for i in range(last_page - page_length + 1, last_page + 1):
  432. if i < 1:
  433. continue
  434. pages.append(i)
  435. else:
  436. for i in range(page - page_length / 2, page + int(ceil(page_length / 2.0))):
  437. pages.append(i)
  438. if i == last_page or last_page == 0:
  439. break
  440. ret = dict()
  441. ret['state'] = ji.Common.exchange_state(20000)
  442. ret['data'] = {
  443. 'disks': disks_ret['data'],
  444. 'hosts_mapping_by_node_id': hosts_mapping_by_node_id,
  445. 'order_by': order_by,
  446. 'order': order,
  447. 'show_area': show_area,
  448. 'config': config.__dict__,
  449. 'show_on_host': show_on_host,
  450. 'paging': disks_ret['paging'],
  451. 'page': page,
  452. 'page_size': page_size,
  453. 'keyword': keyword,
  454. 'pages': pages,
  455. 'last_page': last_page
  456. }
  457. return ret