editor.py 22 KB


  1. #!/usr/bin/env python
  2. # Licensed to Cloudera, Inc. under one
  3. # or more contributor license agreements. See the NOTICE file
  4. # distributed with this work for additional information
  5. # regarding copyright ownership. Cloudera, Inc. licenses this file
  6. # to you under the Apache License, Version 2.0 (the
  7. # "License"); you may not use this file except in compliance
  8. # with the License. You may obtain a copy of the License at
  9. #
  10. # http://www.apache.org/licenses/LICENSE-2.0
  11. #
  12. # Unless required by applicable law or agreed to in writing, software
  13. # distributed under the License is distributed on an "AS IS" BASIS,
  14. # WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
  15. # See the License for the specific language governing permissions and
  16. # limitations under the License.
  17. try:
  18. import json
  19. except ImportError:
  20. import simplejson as json
  21. import logging
  22. from django.forms.models import inlineformset_factory, modelformset_factory
  23. from django.core.urlresolvers import reverse
  24. from django.db.models import Q
  25. from django.http import HttpResponse
  26. from django.shortcuts import redirect
  27. from django.utils.functional import wraps
  28. from django.utils.translation import ugettext as _
  29. from desktop.lib.django_util import render, PopupException, extract_field_data
  30. from desktop.lib.rest.http_client import RestException
  31. from desktop.log.access import access_warn
  32. from hadoop.fs.exceptions import WebHdfsException
  33. from liboozie.submittion import Submission
  34. from oozie.models import Workflow, Node, Link, History, Coordinator,\
  35. Dataset, DataInput, DataOutput, Job, _STD_PROPERTIES_JSON
  36. from oozie.forms import NodeForm, WorkflowForm, CoordinatorForm, DatasetForm,\
  37. DataInputForm, DataInputSetForm, DataOutputForm, DataOutputSetForm, LinkForm,\
  38. DefaultLinkForm, design_form_by_type
  39. from oozie.conf import SHARE_JOBS
  40. LOG = logging.getLogger(__name__)
  41. def can_access_job(request, job_id):
  42. """
  43. Logic for testing if a user can access a certain Workflow / Coordinator.
  44. """
  45. if job_id is None:
  46. return
  47. try:
  48. job = Job.objects.select_related().get(pk=job_id).get_full_node()
  49. if not SHARE_JOBS.get() and not request.user.is_superuser \
  50. and job.owner != request.user.username:
  51. # TODO is shared perms
  52. message = _("Permission denied. %(username)s don't have the permissions to access job %(id)s") % \
  53. {'username': request.user.username, 'id': job.id}
  54. access_warn(request, message)
  55. raise PopupException(message)
  56. else:
  57. return job
  58. except Job.DoesNotExist:
  59. raise PopupException(_('job %(id)s not found') % {'id': job_id})
  60. def can_modify_job(request, job):
  61. """Only owners or admins can modify a job."""
  62. return request.user.is_superuser or job.owner.id == request.user.id
  63. def check_job_modification(request, job):
  64. if not can_modify_job(request, job):
  65. raise PopupException(_('Not allowed to modified this job'))
  66. def check_job_modification_permission(view_func):
  67. """
  68. Decorator ensuring that the user has the permissions to modify a workflow or coordinator.
  69. Need to appear below @check_job_access_permission
  70. """
  71. def decorate(request, *args, **kwargs):
  72. if 'workflow' in kwargs:
  73. job_type = 'workflow'
  74. else:
  75. job_type = 'coordinator'
  76. job = kwargs.get(job_type)
  77. if job is not None:
  78. check_job_modification(request, job)
  79. return view_func(request, *args, **kwargs)
  80. return wraps(view_func)(decorate)
  81. def check_job_access_permission(view_func):
  82. """
  83. Decorator ensuring that the user has access to the workflow or coordinator.
  84. Arg: 'workflow' or 'coordinator' id.
  85. Return: the workflow of coordinator or raise an exception
  86. Notice: its gets an id in input and returns the full object in output (not an id).
  87. """
  88. def decorate(request, *args, **kwargs):
  89. if 'workflow' in kwargs:
  90. job_type = 'workflow'
  91. else:
  92. job_type = 'coordinator'
  93. job = kwargs.get(job_type)
  94. if job is not None:
  95. job = can_access_job(request, job)
  96. kwargs[job_type] = job
  97. return view_func(request, *args, **kwargs)
  98. return wraps(view_func)(decorate)
  99. def check_action_access_permission(view_func):
  100. """
  101. Decorator ensuring that the user has access to the workflow action.
  102. Arg: 'workflow action' id.
  103. Return: the workflow action or raise an exception
  104. Notice: its gets an id in input and returns the full object in output (not an id).
  105. """
  106. def decorate(request, *args, **kwargs):
  107. action_id = kwargs.get('action')
  108. action = Node.objects.get(id=action_id).get_full_node()
  109. can_access_job(request, action.workflow.id)
  110. kwargs['action'] = action
  111. return view_func(request, *args, **kwargs)
  112. return wraps(view_func)(decorate)
  113. def check_action_modification_permission(view_func):
  114. """
  115. Decorator ensuring that the user has the permissions to modify a workflow action.
  116. Need to appear below @check_action_access_permission
  117. """
  118. def decorate(request, *args, **kwargs):
  119. action = kwargs.get('action')
  120. check_job_modification(request, action.workflow)
  121. return view_func(request, *args, **kwargs)
  122. return wraps(view_func)(decorate)
  123. def list_workflows(request, job_type='workflow'):
  124. show_install_examples = True
  125. if job_type == 'coordinators':
  126. data = Coordinator.objects
  127. template = "editor/list_coordinators.mako"
  128. else:
  129. data = Workflow.objects
  130. template = "editor/list_workflows.mako"
  131. if not SHARE_JOBS.get() and not request.user.is_superuser:
  132. data = data.filter(owner=request.user)
  133. else:
  134. data = data.filter(Q(is_shared=True) | Q(owner=request.user))
  135. data = data.order_by('-last_modified')
  136. return render(template, request, {
  137. 'workflows': list(data),
  138. 'currentuser': request.user,
  139. 'show_install_examples': show_install_examples,
  140. })
  141. def create_workflow(request):
  142. workflow = Workflow.objects.new_workflow(request.user)
  143. if request.method == 'POST':
  144. workflow_form = WorkflowForm(request.POST, instance=workflow)
  145. if workflow_form.is_valid():
  146. wf = workflow_form.save()
  147. Workflow.objects.initialize(wf, request.fs)
  148. return redirect(reverse('oozie:edit_workflow', kwargs={'workflow': workflow.id}))
  149. else:
  150. workflow_form = WorkflowForm(instance=workflow)
  151. return render('editor/create_workflow.mako', request, {
  152. 'workflow_form': workflow_form,
  153. 'workflow': workflow,
  154. })
  155. @check_job_access_permission
  156. def edit_workflow(request, workflow):
  157. WorkflowFormSet = inlineformset_factory(Workflow, Node, form=NodeForm, max_num=0, can_order=False, can_delete=False)
  158. history = History.objects.filter(submitter=request.user, job=workflow)
  159. if request.method == 'POST' and can_modify_job(request, workflow):
  160. try:
  161. workflow_form = WorkflowForm(request.POST, instance=workflow)
  162. actions_formset = WorkflowFormSet(request.POST, request.FILES, instance=workflow)
  163. if 'clone_action' in request.POST: return clone_action(request, action=request.POST['clone_action'])
  164. if 'delete_action' in request.POST: return delete_action(request, action=request.POST['delete_action'])
  165. if 'move_up_action' in request.POST: return move_up_action(request, action=request.POST['move_up_action'])
  166. if 'move_down_action' in request.POST: return move_down_action(request, action=request.POST['move_down_action'])
  167. if workflow_form.is_valid() and actions_formset.is_valid():
  168. workflow_form.save()
  169. actions_formset.save()
  170. return redirect(reverse('oozie:list_workflows'))
  171. except Exception, e:
  172. request.error(_('Sorry, this operation is not supported: %(error)s') % {'error': e})
  173. else:
  174. workflow_form = WorkflowForm(instance=workflow)
  175. actions_formset = WorkflowFormSet(instance=workflow)
  176. return render('editor/edit_workflow.mako', request, {
  177. 'workflow_form': workflow_form,
  178. 'workflow': workflow,
  179. 'actions_formset': actions_formset,
  180. 'graph': workflow.gen_graph(actions_formset.forms),
  181. 'history': history,
  182. })
  183. @check_job_access_permission
  184. @check_job_modification_permission
  185. def delete_workflow(request, workflow):
  186. if request.method != 'POST':
  187. raise PopupException(_('A POST request is required.'))
  188. try:
  189. workflow.coordinator_set.update(workflow=None) # In Django 1.3 could do ON DELETE set NULL
  190. workflow.save()
  191. workflow.delete()
  192. Submission(workflow, request.fs, {}).remove_deployment_dir()
  193. except Workflow.DoesNotExist:
  194. LOG.error("Trying to delete non-existent workflow (id %s)" % (workflow,))
  195. raise PopupException(_('Workflow not found'))
  196. # TODO notification
  197. return redirect(reverse('oozie:list_workflows'))
  198. @check_job_access_permission
  199. def clone_workflow(request, workflow):
  200. if request.method != 'POST':
  201. raise PopupException(_('A POST request is required.'))
  202. clone = workflow.clone(request.user)
  203. response = {'url': reverse('oozie:edit_workflow', kwargs={'workflow': clone.id})}
  204. return HttpResponse(json.dumps(response), mimetype="application/json")
  205. @check_job_access_permission
  206. def submit_workflow(request, workflow):
  207. if request.method != 'POST':
  208. raise PopupException(_('A POST request is required.'))
  209. try:
  210. mapping = dict(request.POST.iteritems())
  211. submission = Submission(workflow, request.fs, mapping)
  212. job_id = submission.run()
  213. except RestException, ex:
  214. raise PopupException(_("Error submitting workflow %s") % (workflow,),
  215. detail=ex._headers.get('oozie-error-message', ex))
  216. History.objects.create_from_submission(submission)
  217. request.info(_('Workflow submitted'))
  218. return redirect(reverse('oozie:list_oozie_workflow', kwargs={'job_id': job_id}))
  219. @check_job_access_permission
  220. def get_workflow_parameters(request, workflow):
  221. """
  222. Return the parameters found in the workflow as a JSON dictionary of {param_key : label}.
  223. This expects an Ajax call.
  224. """
  225. params = workflow.find_parameters()
  226. params_with_labels = dict((p, p.upper()) for p in params)
  227. return render('dont_care_for_ajax', request, { 'params': params_with_labels })
  228. @check_job_access_permission
  229. def new_action(request, workflow, node_type, parent_action_id):
  230. ActionForm = design_form_by_type(node_type)
  231. if request.method == 'POST':
  232. action_form = ActionForm(request.POST)
  233. if action_form.is_valid():
  234. action = action_form.save(commit=False)
  235. action.node_type = node_type
  236. action.workflow = workflow
  237. action.save()
  238. workflow.add_action(action, parent_action_id)
  239. return redirect(reverse('oozie:edit_workflow', kwargs={'workflow': workflow.id}))
  240. else:
  241. action_form = ActionForm()
  242. return render('editor/edit_workflow_action.mako', request, {
  243. 'workflow': workflow,
  244. 'job_properties': extract_field_data(action_form['job_properties']),
  245. 'files': extract_field_data(action_form['files']),
  246. 'archives': extract_field_data(action_form['archives']),
  247. 'params': 'params' in action_form.fields and extract_field_data(action_form['params']) or '[]',
  248. 'action_form': action_form,
  249. 'node_type': node_type,
  250. 'properties_hint': _STD_PROPERTIES_JSON,
  251. 'form_url': reverse('oozie:new_action', kwargs={'workflow': workflow.id,
  252. 'node_type': node_type,
  253. 'parent_action_id': parent_action_id}),
  254. })
  255. @check_action_access_permission
  256. def edit_action(request, action):
  257. ActionForm = design_form_by_type(action.node_type)
  258. if request.method == 'POST' and can_modify_job(request, action.workflow):
  259. action_form = ActionForm(request.POST, instance=action)
  260. if action_form.is_valid():
  261. action = action_form.save()
  262. return redirect(reverse('oozie:edit_workflow', kwargs={'workflow': action.workflow.id}))
  263. else:
  264. action_form = ActionForm(instance=action)
  265. return render('editor/edit_workflow_action.mako', request, {
  266. 'workflow': action.workflow,
  267. 'job_properties': extract_field_data(action_form['job_properties']),
  268. 'files': extract_field_data(action_form['files']),
  269. 'archives': extract_field_data(action_form['archives']),
  270. 'params': 'params' in action_form.fields and extract_field_data(action_form['params']) or '[]',
  271. 'action_form': action_form,
  272. 'node_type': action.node_type,
  273. 'properties_hint': _STD_PROPERTIES_JSON,
  274. 'form_url': reverse('oozie:edit_action', kwargs={'action': action.id}),
  275. })
  276. @check_action_access_permission
  277. @check_action_modification_permission
  278. def edit_workflow_fork(request, action):
  279. fork = action
  280. LinkFormSet = modelformset_factory(Link, form=LinkForm, max_num=0)
  281. if request.method == 'POST':
  282. link_formset = LinkFormSet(request.POST)
  283. default_link_form = DefaultLinkForm(request.POST, action=fork)
  284. if link_formset.is_valid():
  285. is_decision = fork.has_decisions()
  286. link_formset.save()
  287. if not is_decision and fork.has_decisions():
  288. default_link = default_link_form.save(commit=False)
  289. default_link.parent = fork
  290. default_link.name = 'default'
  291. default_link.comment = 'default'
  292. default_link.save()
  293. fork.convert_to_decision()
  294. fork.update_description()
  295. return redirect(reverse('oozie:edit_workflow', kwargs={'workflow': fork.workflow.id}))
  296. else:
  297. link_formset = LinkFormSet(queryset=fork.get_children_links().exclude(name__in=['related', 'default']))
  298. default_link = Link(parent=fork, name='default', comment='default')
  299. default_link_form = DefaultLinkForm(action=fork, instance=default_link)
  300. return render('editor/edit_workflow_fork.mako', request, {
  301. 'workflow': fork.workflow,
  302. 'fork': fork,
  303. 'link_formset': link_formset,
  304. 'default_link_form': default_link_form,
  305. })
  306. @check_action_access_permission
  307. @check_action_modification_permission
  308. def delete_action(request, action):
  309. if request.method == 'POST':
  310. action.workflow.delete_action(action)
  311. return redirect(reverse('oozie:edit_workflow', kwargs={'workflow': action.workflow.id}))
  312. else:
  313. raise PopupException(_('A POST request is required.'))
  314. @check_action_access_permission
  315. def clone_action(request, action):
  316. # Really weird: action is like a clone object with the old id here
  317. action_id = action.id
  318. workflow = action.workflow
  319. clone = action.clone()
  320. workflow.add_action(clone, action_id)
  321. return redirect(reverse('oozie:edit_workflow', kwargs={'workflow': workflow.id}))
  322. @check_action_access_permission
  323. @check_action_modification_permission
  324. def move_up_action(request, action):
  325. if request.method == 'POST':
  326. action.workflow.move_action_up(action)
  327. return redirect(reverse('oozie:edit_workflow', kwargs={'workflow': action.workflow.id}))
  328. else:
  329. raise PopupException(_('A POST request is required.'))
  330. @check_action_access_permission
  331. @check_action_modification_permission
  332. def move_down_action(request, action):
  333. if request.method == 'POST':
  334. action.workflow.move_action_down(action)
  335. return redirect(reverse('oozie:edit_workflow', kwargs={'workflow': action.workflow.id}))
  336. else:
  337. raise PopupException(_('A POST request is required.'))
  338. @check_job_access_permission
  339. def create_coordinator(request, workflow=None):
  340. if workflow is not None:
  341. coordinator = Coordinator(owner=request.user, workflow=workflow)
  342. else:
  343. coordinator = Coordinator(owner=request.user)
  344. if request.method == 'POST':
  345. coordinator_form = CoordinatorForm(request.POST, instance=coordinator)
  346. if coordinator_form.is_valid():
  347. coordinator = coordinator_form.save()
  348. return redirect(reverse('oozie:edit_coordinator', kwargs={'coordinator': coordinator.id}))
  349. else:
  350. coordinator_form = CoordinatorForm(instance=coordinator)
  351. return render('editor/create_coordinator.mako', request, {
  352. 'coordinator': coordinator,
  353. 'coordinator_form': coordinator_form,
  354. })
  355. @check_job_access_permission
  356. @check_job_modification_permission
  357. def edit_coordinator(request, coordinator):
  358. history = History.objects.filter(submitter=request.user, job=coordinator)
  359. DatasetFormSet = inlineformset_factory(Coordinator, Dataset, form=DatasetForm, max_num=0, can_order=False, can_delete=True)
  360. DataInputFormSet = inlineformset_factory(Coordinator, DataInput, form=DataInputSetForm, max_num=0, can_order=False, can_delete=True)
  361. DataOutputFormSet = inlineformset_factory(Coordinator, DataOutput, form=DataOutputSetForm, max_num=0, can_order=False, can_delete=True)
  362. dataset = Dataset(coordinator=coordinator)
  363. dataset_form = DatasetForm(instance=dataset)
  364. data_input = DataInput(coordinator=coordinator)
  365. data_input_form = DataInputForm(instance=data_input, coordinator=coordinator)
  366. data_output = DataOutput(coordinator=coordinator)
  367. data_output_form = DataOutputForm(instance=data_output, coordinator=coordinator)
  368. if request.method == 'POST':
  369. coordinator_form = CoordinatorForm(request.POST, instance=coordinator)
  370. dataset_formset = DatasetFormSet(request.POST, request.FILES, instance=coordinator)
  371. data_input_formset = DataInputFormSet(request.POST, request.FILES, instance=coordinator)
  372. data_output_formset = DataOutputFormSet(request.POST, request.FILES, instance=coordinator)
  373. if coordinator_form.is_valid() and dataset_formset.is_valid() and data_input_formset.is_valid() and data_output_formset.is_valid():
  374. coordinator = coordinator_form.save()
  375. dataset_formset.save()
  376. data_input_formset.save()
  377. data_output_formset.save()
  378. return redirect(reverse('oozie:list_coordinator'))
  379. else:
  380. coordinator_form = CoordinatorForm(instance=coordinator)
  381. dataset_formset = DatasetFormSet(instance=coordinator)
  382. data_input_formset = DataInputFormSet(instance=coordinator)
  383. data_output_formset = DataOutputFormSet(instance=coordinator)
  384. return render('editor/edit_coordinator.mako', request, {
  385. 'coordinator': coordinator,
  386. 'coordinator_form': coordinator_form,
  387. 'dataset_formset': dataset_formset,
  388. 'data_input_formset': data_input_formset,
  389. 'data_output_formset': data_output_formset,
  390. 'dataset_form': dataset_form,
  391. 'data_input_form': data_input_form,
  392. 'data_output_form': data_output_form,
  393. 'history': history,
  394. })
  395. @check_job_access_permission
  396. @check_job_modification_permission
  397. def create_coordinator_dataset(request, coordinator):
  398. """Returns {'status' 0/1, data:html or url}"""
  399. dataset = Dataset(coordinator=coordinator)
  400. response = {'status': -1, 'data': 'None'}
  401. if request.method == 'POST':
  402. dataset_form = DatasetForm(request.POST, instance=dataset)
  403. if dataset_form.is_valid():
  404. dataset_form.save()
  405. response['status'] = 0
  406. response['data'] = reverse('oozie:edit_coordinator', kwargs={'coordinator': coordinator.id})
  407. request.info(_('Dataset created'));
  408. else:
  409. dataset_form = DatasetForm(request.POST, instance=dataset)
  410. else:
  411. ## Bad
  412. response['data'] = _('A POST request is required.')
  413. if response['status'] != 0:
  414. response['data'] = render('editor/create_coordinator_dataset.mako', request, {
  415. 'coordinator': coordinator,
  416. 'dataset_form': dataset_form,
  417. }, force_template=True).content
  418. return HttpResponse(json.dumps(response), mimetype="application/json")
  419. @check_job_access_permission
  420. @check_job_modification_permission
  421. def create_coordinator_data(request, coordinator, data_type):
  422. """Returns {'status' 0/1, data:html or url}"""
  423. if data_type == 'input':
  424. data_instance = DataInput(coordinator=coordinator)
  425. DataForm = DataInputForm
  426. else:
  427. data_instance = DataOutput(coordinator=coordinator)
  428. DataForm = DataOutputForm
  429. response = {'status': -1, 'data': 'None'}
  430. if request.method == 'POST':
  431. data_form = DataForm(request.POST, instance=data_instance, coordinator=coordinator)
  432. if data_form.is_valid():
  433. data_form.save()
  434. response['status'] = 0
  435. response['data'] = reverse('oozie:edit_coordinator', kwargs={'coordinator': coordinator.id})
  436. request.info(_('Coordinator data created'));
  437. else:
  438. data_form = DataForm(request.POST, instance=data_instance, coordinator=coordinator)
  439. else:
  440. ## Bad
  441. response['data'] = _('A POST request is required.')
  442. if response['status'] != 0:
  443. response['data'] = render('editor/create_coordinator_data.mako', request, {
  444. 'coordinator': coordinator,
  445. 'form': data_form, },
  446. force_template=True).content
  447. return HttpResponse(json.dumps(response), mimetype="application/json")
  448. @check_job_access_permission
  449. def submit_coordinator(request, coordinator):
  450. if request.method != 'POST':
  451. raise PopupException(_('A POST request is required.'))
  452. try:
  453. if not coordinator.workflow.is_deployed(request.fs):
  454. submission = Submission(coordinator.workflow, request.fs, request.POST)
  455. wf_dir = submission.deploy()
  456. coordinator.workflow.deployment_dir = wf_dir
  457. coordinator.workflow.save()
  458. coordinator.deployment_dir = coordinator.workflow.deployment_dir
  459. properties = {'wf_application_path': coordinator.workflow.deployment_dir}
  460. properties.update(dict(request.POST.iteritems()))
  461. submission = Submission(coordinator, request.fs, properties=properties)
  462. job_id = submission.run()
  463. except RestException, ex:
  464. raise PopupException(_("Error submitting coordinator %s") % (coordinator,),
  465. detail=ex._headers['oozie-error-message'])
  466. History.objects.create_from_submission(submission)
  467. request.info(_('Coordinator submitted'))
  468. return redirect(reverse('oozie:list_oozie_coordinator', kwargs={'job_id': job_id}))
  469. def list_history(request):
  470. """
  471. List the job submission history.
  472. Normal users can only look at their own submissions.
  473. """
  474. history = History.objects
  475. if not request.user.is_superuser:
  476. history = history.filter(submitter=request.user)
  477. history = history.order_by('-submission_date')
  478. return render('editor/list_history.mako', request, {
  479. 'history': history,
  480. })
  481. def list_history_record(request, record_id):
  482. """
  483. List a job submission history.
  484. Normal users can only look at their own jobs.
  485. """
  486. history = History.objects
  487. if not request.user.is_superuser:
  488. history.filter(submitter=request.user)
  489. history = history.get(id=record_id)
  490. return render('editor/list_history_record.mako', request, {
  491. 'record': history,
  492. })
  493. def setup(request):
  494. """Installs oozie examples."""
  495. if request.method != 'POST':
  496. raise PopupException(_('A POST request is required.'))
  497. try:
  498. # Warning: below will modify fs.user
  499. #oozie_setup.Command().handle_noargs()
  500. pass
  501. except WebHdfsException, e:
  502. raise PopupException(_('The examples could not be installed.'), detail=e)
  503. return redirect(reverse('oozie:list_workflows'))