editor.py 31 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.core.urlresolvers import reverse
  23. from django.db.models import Q
  24. from django.forms.formsets import formset_factory
  25. from django.forms.models import inlineformset_factory, modelformset_factory
  26. from django.http import HttpResponse
  27. from django.shortcuts import redirect
  28. from django.utils.functional import curry, wraps
  29. from django.utils.translation import ugettext as _
  30. from desktop.lib.django_util import render, extract_field_data
  31. from desktop.lib.exceptions import PopupException
  32. from desktop.lib.rest.http_client import RestException
  33. from hadoop.fs.exceptions import WebHdfsException
  34. from jobsub.models import OozieDesign
  35. from liboozie.submittion import Submission
  36. from oozie.conf import SHARE_JOBS
  37. from oozie.import_jobsub import convert_jobsub_design
  38. from oozie.management.commands import oozie_setup
  39. from oozie.models import Job, Workflow, Node, Link, History, Coordinator,\
  40. Mapreduce, Java, Streaming, Dataset, DataInput, DataOutput,\
  41. _STD_PROPERTIES_JSON
  42. from oozie.forms import NodeForm, WorkflowForm, CoordinatorForm, DatasetForm,\
  43. DataInputForm, DataInputSetForm, DataOutputForm, DataOutputSetForm, LinkForm,\
  44. DefaultLinkForm, design_form_by_type, ImportJobsubDesignForm, ParameterForm
  45. LOG = logging.getLogger(__name__)
  46. def check_job_access_permission(view_func):
  47. """
  48. Decorator ensuring that the user has access to the workflow or coordinator.
  49. Arg: 'workflow' or 'coordinator' id.
  50. Return: the workflow of coordinator or raise an exception
  51. Notice: its gets an id in input and returns the full object in output (not an id).
  52. """
  53. def decorate(request, *args, **kwargs):
  54. if 'workflow' in kwargs:
  55. job_type = 'workflow'
  56. else:
  57. job_type = 'coordinator'
  58. job = kwargs.get(job_type)
  59. if job is not None:
  60. job = Job.objects.is_accessible_or_exception(request, job)
  61. kwargs[job_type] = job
  62. return view_func(request, *args, **kwargs)
  63. return wraps(view_func)(decorate)
  64. def check_job_edition_permission(authorize_get=False):
  65. """
  66. Decorator ensuring that the user has the permissions to modify a workflow or coordinator.
  67. Need to appear below @check_job_access_permission
  68. """
  69. def inner(view_func):
  70. def decorate(request, *args, **kwargs):
  71. if 'workflow' in kwargs:
  72. job_type = 'workflow'
  73. else:
  74. job_type = 'coordinator'
  75. job = kwargs.get(job_type)
  76. if job is not None and not (authorize_get and request.method == 'GET'):
  77. Job.objects.can_edit_or_exception(request, job)
  78. return view_func(request, *args, **kwargs)
  79. return wraps(view_func)(decorate)
  80. return inner
  81. def check_action_access_permission(view_func):
  82. """
  83. Decorator ensuring that the user has access to the workflow action.
  84. Arg: 'workflow action' id.
  85. Return: the workflow action 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. action_id = kwargs.get('action')
  90. action = Node.objects.get(id=action_id).get_full_node()
  91. Job.objects.is_accessible_or_exception(request, action.workflow.id)
  92. kwargs['action'] = action
  93. return view_func(request, *args, **kwargs)
  94. return wraps(view_func)(decorate)
  95. def check_action_edition_permission(view_func):
  96. """
  97. Decorator ensuring that the user has the permissions to modify a workflow action.
  98. Need to appear below @check_action_access_permission
  99. """
  100. def decorate(request, *args, **kwargs):
  101. action = kwargs.get('action')
  102. Job.objects.can_edit_or_exception(request, action.workflow)
  103. return view_func(request, *args, **kwargs)
  104. return wraps(view_func)(decorate)
  105. def check_dataset_access_permission(view_func):
  106. """
  107. Decorator ensuring that the user has access to dataset.
  108. Arg: 'dataset'.
  109. Return: the dataset or raise an exception
  110. Notice: its gets an id in input and returns the full object in output (not an id).
  111. """
  112. def decorate(request, *args, **kwargs):
  113. dataset = kwargs.get('dataset')
  114. if dataset is not None:
  115. dataset = Dataset.objects.is_accessible_or_exception(request, dataset)
  116. kwargs['dataset'] = dataset
  117. return view_func(request, *args, **kwargs)
  118. return wraps(view_func)(decorate)
  119. def check_dataset_edition_permission(authorize_get=False):
  120. """
  121. Decorator ensuring that the user has the permissions to modify a dataset.
  122. A dataset can be edited if the coordinator that owns the dataset can be edited.
  123. Need to appear below @check_dataset_access_permission
  124. """
  125. def inner(view_func):
  126. def decorate(request, *args, **kwargs):
  127. dataset = kwargs.get('dataset')
  128. if dataset is not None and not (authorize_get and request.method == 'GET'):
  129. Job.objects.can_edit_or_exception(request, dataset.coordinator)
  130. return view_func(request, *args, **kwargs)
  131. return wraps(view_func)(decorate)
  132. return inner
  133. def list_workflows(request):
  134. show_setup_app = True
  135. data = Workflow.objects
  136. if not SHARE_JOBS.get() and not request.user.is_superuser:
  137. data = data.filter(owner=request.user)
  138. else:
  139. data = data.filter(Q(is_shared=True) | Q(owner=request.user))
  140. data = data.order_by('-last_modified')
  141. return render('editor/list_workflows.mako', request, {
  142. 'jobs': list(data),
  143. 'currentuser': request.user,
  144. 'show_setup_app': show_setup_app,
  145. })
  146. def list_coordinators(request, workflow_id=None):
  147. show_setup_app = True
  148. data = Coordinator.objects
  149. if workflow_id is not None:
  150. data = data.filter(workflow__id=workflow_id)
  151. if not SHARE_JOBS.get() and not request.user.is_superuser:
  152. data = data.filter(owner=request.user)
  153. else:
  154. data = data.filter(Q(is_shared=True) | Q(owner=request.user))
  155. data = data.order_by('-last_modified')
  156. return render('editor/list_coordinators.mako', request, {
  157. 'jobs': list(data),
  158. 'currentuser': request.user,
  159. 'show_setup_app': show_setup_app,
  160. })
  161. def create_workflow(request):
  162. workflow = Workflow.objects.new_workflow(request.user)
  163. if request.method == 'POST':
  164. workflow_form = WorkflowForm(request.POST, instance=workflow)
  165. if workflow_form.is_valid():
  166. wf = workflow_form.save()
  167. Workflow.objects.initialize(wf, request.fs)
  168. return redirect(reverse('oozie:edit_workflow', kwargs={'workflow': workflow.id}))
  169. else:
  170. request.error(_('Errors on the form: %s') % workflow_form.errors)
  171. else:
  172. workflow_form = WorkflowForm(instance=workflow)
  173. return render('editor/create_workflow.mako', request, {
  174. 'workflow_form': workflow_form,
  175. 'workflow': workflow,
  176. })
  177. @check_job_access_permission
  178. def edit_workflow(request, workflow):
  179. WorkflowFormSet = inlineformset_factory(Workflow, Node, form=NodeForm, max_num=0, can_order=False, can_delete=False)
  180. history = History.objects.filter(submitter=request.user, job=workflow).order_by('-submission_date')
  181. if request.method == 'POST' and Job.objects.can_edit_or_exception(request, workflow):
  182. try:
  183. workflow_form = WorkflowForm(request.POST, instance=workflow)
  184. actions_formset = WorkflowFormSet(request.POST, request.FILES, instance=workflow)
  185. if 'clone_action' in request.POST: return clone_action(request, action=request.POST['clone_action'])
  186. if 'delete_action' in request.POST: return delete_action(request, action=request.POST['delete_action'])
  187. if 'move_up_action' in request.POST: return move_up_action(request, action=request.POST['move_up_action'])
  188. if 'move_down_action' in request.POST: return move_down_action(request, action=request.POST['move_down_action'])
  189. if workflow_form.is_valid() and actions_formset.is_valid():
  190. workflow_form.save()
  191. actions_formset.save()
  192. if workflow.has_cycle():
  193. raise PopupException(_('Sorry, this operation is not creating a cycle which would break the workflow.'))
  194. request.info(_("Workflow saved!"))
  195. return redirect(reverse('oozie:edit_workflow', kwargs={'workflow': workflow.id}))
  196. except Exception, e:
  197. request.error(_('Sorry, this operation is not supported: %(error)s') % {'error': e})
  198. workflow_form = WorkflowForm(instance=workflow)
  199. actions_formset = WorkflowFormSet(instance=workflow)
  200. graph_options = {}
  201. user_can_edit_job = workflow.is_editable(request.user)
  202. if not user_can_edit_job:
  203. graph_options = {'template': 'editor/gen/workflow-graph-readonly.xml.mako'}
  204. graph = workflow.gen_graph(actions_formset.forms, **graph_options)
  205. return render('editor/edit_workflow.mako', request, {
  206. 'workflow_form': workflow_form,
  207. 'workflow': workflow,
  208. 'actions_formset': actions_formset,
  209. 'graph': graph,
  210. 'history': history,
  211. 'user_can_edit_job': user_can_edit_job,
  212. 'parameters': extract_field_data(workflow_form['parameters']),
  213. 'job_properties': extract_field_data(workflow_form['job_properties'])
  214. })
  215. @check_job_access_permission
  216. @check_job_edition_permission()
  217. def delete_workflow(request, workflow):
  218. if request.method != 'POST':
  219. raise PopupException(_('A POST request is required.'))
  220. Workflow.objects.destroy(workflow, request.fs)
  221. request.info(_('Workflow deleted!'))
  222. return redirect(reverse('oozie:list_workflows'))
  223. @check_job_access_permission
  224. def clone_workflow(request, workflow):
  225. if request.method != 'POST':
  226. raise PopupException(_('A POST request is required.'))
  227. clone = workflow.clone(request.fs, request.user)
  228. response = {'url': reverse('oozie:edit_workflow', kwargs={'workflow': clone.id})}
  229. return HttpResponse(json.dumps(response), mimetype="application/json")
  230. @check_job_access_permission
  231. def submit_workflow(request, workflow):
  232. ParametersFormSet = formset_factory(ParameterForm, extra=0)
  233. if request.method == 'POST':
  234. params_form = ParametersFormSet(request.POST)
  235. if params_form.is_valid():
  236. mapping = dict([(param['name'], param['value']) for param in params_form.cleaned_data])
  237. job_id = _submit_workflow(request, workflow, mapping)
  238. request.info(_('Workflow submitted'))
  239. return redirect(reverse('oozie:list_oozie_workflow', kwargs={'job_id': job_id}))
  240. else:
  241. request.error(_('Invalid submission form: %s' % params_form.errors))
  242. else:
  243. parameters = workflow.find_all_parameters()
  244. params_form = ParametersFormSet(initial=parameters)
  245. popup = render('editor/submit_job_popup.mako', request, {
  246. 'params_form': params_form,
  247. 'action': reverse('oozie:submit_workflow', kwargs={'workflow': workflow.id})
  248. }, force_template=True).content
  249. return HttpResponse(json.dumps(popup), mimetype="application/json")
  250. def _submit_workflow(request, workflow, mapping):
  251. try:
  252. submission = Submission(request.user, workflow, request.fs, mapping)
  253. job_id = submission.run()
  254. History.objects.create_from_submission(submission)
  255. return job_id
  256. except RestException, ex:
  257. raise PopupException(_("Error submitting workflow %s") % (workflow,),
  258. detail=ex._headers.get('oozie-error-message', ex))
  259. request.info(_('Workflow submitted'))
  260. return redirect(reverse('oozie:list_oozie_workflow', kwargs={'job_id': job_id}))
  261. def resubmit_workflow(request, oozie_wf_id):
  262. if request.method != 'POST':
  263. raise PopupException(_('A POST request is required.'))
  264. history = History.objects.get(oozie_job_id=oozie_wf_id)
  265. Job.objects.is_accessible_or_exception(request, history.job.id)
  266. workflow = history.get_workflow().get_full_node()
  267. properties = history.properties_dict
  268. job_id = _submit_workflow(request, workflow, properties)
  269. request.info(_('Workflow re-submitted'))
  270. return redirect(reverse('oozie:list_oozie_workflow', kwargs={'job_id': job_id}))
  271. @check_job_access_permission
  272. def schedule_workflow(request, workflow):
  273. if Coordinator.objects.filter(workflow=workflow).exists():
  274. request.info(_('You already have some coordinators for this workflow. Please submit one or create a new one.'))
  275. return list_coordinators(request, workflow_id=workflow.id)
  276. else:
  277. return create_coordinator(request, workflow=workflow.id)
  278. @check_job_access_permission
  279. def new_action(request, workflow, node_type, parent_action_id):
  280. ActionForm = design_form_by_type(node_type)
  281. if request.method == 'POST':
  282. action_form = ActionForm(request.POST)
  283. if action_form.is_valid():
  284. action = action_form.save(commit=False)
  285. action.node_type = node_type
  286. action.workflow = workflow
  287. action.save()
  288. workflow.add_action(action, parent_action_id)
  289. return redirect(reverse('oozie:edit_workflow', kwargs={'workflow': workflow.id}))
  290. else:
  291. action_form = ActionForm()
  292. return render('editor/edit_workflow_action.mako', request, {
  293. 'workflow': workflow,
  294. 'job_properties': 'job_properties' in action_form.fields and extract_field_data(action_form['job_properties']) or '[]',
  295. 'files': 'files' in action_form.fields and extract_field_data(action_form['files']) or '[]',
  296. 'archives': 'archives' in action_form.fields and extract_field_data(action_form['archives']) or '[]',
  297. 'params': 'params' in action_form.fields and extract_field_data(action_form['params']) or '[]',
  298. 'prepares': 'prepares' in action_form.fields and extract_field_data(action_form['prepares']) or '[]',
  299. 'action_form': action_form,
  300. 'node_type': node_type,
  301. 'properties_hint': _STD_PROPERTIES_JSON,
  302. 'form_url': reverse('oozie:new_action', kwargs={'workflow': workflow.id,
  303. 'node_type': node_type,
  304. 'parent_action_id': parent_action_id}),
  305. 'can_edit_action': True,
  306. })
  307. @check_action_access_permission
  308. def edit_action(request, action):
  309. ActionForm = design_form_by_type(action.node_type)
  310. if request.method == 'POST' and Job.objects.can_edit_or_exception(request, action.workflow):
  311. action_form = ActionForm(request.POST, instance=action)
  312. if action_form.is_valid():
  313. action = action_form.save()
  314. return redirect(reverse('oozie:edit_workflow', kwargs={'workflow': action.workflow.id}))
  315. else:
  316. action_form = ActionForm(instance=action)
  317. return render('editor/edit_workflow_action.mako', request, {
  318. 'workflow': action.workflow,
  319. 'job_properties': 'job_properties' in action_form.fields and extract_field_data(action_form['job_properties']) or '[]',
  320. 'files': 'files' in action_form.fields and extract_field_data(action_form['files']) or '[]',
  321. 'archives': 'archives' in action_form.fields and extract_field_data(action_form['archives']) or '[]',
  322. 'params': 'params' in action_form.fields and extract_field_data(action_form['params']) or '[]',
  323. 'prepares': 'prepares' in action_form.fields and extract_field_data(action_form['prepares']) or '[]',
  324. 'action_form': action_form,
  325. 'node_type': action.node_type,
  326. 'properties_hint': _STD_PROPERTIES_JSON,
  327. 'form_url': reverse('oozie:edit_action', kwargs={'action': action.id}),
  328. 'can_edit_action': action.workflow.is_editable(request.user)
  329. })
  330. @check_job_access_permission
  331. def import_action(request, workflow, parent_action_id):
  332. available_actions = OozieDesign.objects.all()
  333. if request.method == 'POST':
  334. form = ImportJobsubDesignForm(data=request.POST, choices=[(action.id, action.name) for action in available_actions])
  335. if form.is_valid():
  336. try:
  337. design = OozieDesign.objects.get(id=form.cleaned_data['action_id'])
  338. action = convert_jobsub_design(design)
  339. action.workflow = workflow
  340. action.save()
  341. workflow.add_action(action, parent_action_id)
  342. except OozieDesign.DoesNotExist:
  343. request.error(_('Jobsub design doesn\'t exist.'))
  344. except (Mapreduce.DoesNotExist, Streaming.DoesNotExist, Java.DoesNotExist):
  345. request.error(_('Could not convert jobsub design'))
  346. except:
  347. request.error(_('Could not convert jobsub design or add action to workflow'))
  348. return redirect(reverse('oozie:edit_workflow', kwargs={'workflow': action.workflow.id}))
  349. return render('editor/import_workflow_action.mako', request, {
  350. 'workflow': workflow,
  351. 'available_actions': available_actions,
  352. 'form_url': reverse('oozie:import_action', kwargs={'workflow': workflow.id, 'parent_action_id': parent_action_id}),
  353. })
  354. @check_action_access_permission
  355. @check_action_edition_permission
  356. def edit_workflow_fork(request, action):
  357. fork = action
  358. LinkFormSet = modelformset_factory(Link, form=LinkForm, max_num=0)
  359. if request.method == 'POST':
  360. link_formset = LinkFormSet(request.POST)
  361. default_link_form = DefaultLinkForm(request.POST, action=fork)
  362. if link_formset.is_valid():
  363. is_decision = fork.has_decisions()
  364. link_formset.save()
  365. if not is_decision and fork.has_decisions():
  366. default_link = default_link_form.save(commit=False)
  367. default_link.parent = fork
  368. default_link.name = 'default'
  369. default_link.comment = 'default'
  370. default_link.save()
  371. fork.convert_to_decision()
  372. fork.update_description()
  373. return redirect(reverse('oozie:edit_workflow', kwargs={'workflow': fork.workflow.id}))
  374. else:
  375. if filter(lambda link: link.child.id != action.workflow.end.id,
  376. [link for link in fork.get_child_join().get_children_links()]):
  377. raise PopupException(_('Sorry, this Fork has some other actions below its Join and cannot be converted. '
  378. 'Please delete the nodes below the Join.'))
  379. link_formset = LinkFormSet(queryset=fork.get_children_links())
  380. default_link = Link(parent=fork, name='default', comment='default')
  381. default_link_form = DefaultLinkForm(action=fork, instance=default_link)
  382. return render('editor/edit_workflow_fork.mako', request, {
  383. 'workflow': fork.workflow,
  384. 'fork': fork,
  385. 'link_formset': link_formset,
  386. 'default_link_form': default_link_form,
  387. })
  388. @check_action_access_permission
  389. @check_action_edition_permission
  390. def delete_action(request, action):
  391. if request.method == 'POST':
  392. action.workflow.delete_action(action)
  393. return redirect(reverse('oozie:edit_workflow', kwargs={'workflow': action.workflow.id}))
  394. else:
  395. raise PopupException(_('A POST request is required.'))
  396. @check_action_access_permission
  397. @check_action_edition_permission
  398. def clone_action(request, action):
  399. if request.method == 'POST':
  400. # Really weird: action is like a clone object with the old id here
  401. action_id = action.id
  402. workflow = action.workflow
  403. clone = action.clone()
  404. workflow.add_action(clone, action_id)
  405. return redirect(reverse('oozie:edit_workflow', kwargs={'workflow': workflow.id}))
  406. else:
  407. raise PopupException(_('A POST request is required.'))
  408. @check_action_access_permission
  409. @check_action_edition_permission
  410. def move_up_action(request, action):
  411. if request.method == 'POST':
  412. action.workflow.move_action_up(action)
  413. return redirect(reverse('oozie:edit_workflow', kwargs={'workflow': action.workflow.id}))
  414. else:
  415. raise PopupException(_('A POST request is required.'))
  416. @check_action_access_permission
  417. @check_action_edition_permission
  418. def move_down_action(request, action):
  419. if request.method == 'POST':
  420. action.workflow.move_action_down(action)
  421. return redirect(reverse('oozie:edit_workflow', kwargs={'workflow': action.workflow.id}))
  422. else:
  423. raise PopupException(_('A POST request is required.'))
  424. @check_job_access_permission
  425. def create_coordinator(request, workflow=None):
  426. if workflow is not None:
  427. coordinator = Coordinator(owner=request.user, schema_version="uri:oozie:coordinator:0.1", workflow=workflow)
  428. else:
  429. coordinator = Coordinator(owner=request.user, schema_version="uri:oozie:coordinator:0.1")
  430. if request.method == 'POST':
  431. coordinator_form = CoordinatorForm(request.POST, instance=coordinator, user=request.user)
  432. if coordinator_form.is_valid():
  433. coordinator = coordinator_form.save()
  434. return redirect(reverse('oozie:edit_coordinator', kwargs={'coordinator': coordinator.id}))
  435. else:
  436. request.error(_('Errors on the form: %s') % coordinator_form.errors)
  437. else:
  438. coordinator_form = CoordinatorForm(instance=coordinator, user=request.user)
  439. return render('editor/create_coordinator.mako', request, {
  440. 'coordinator': coordinator,
  441. 'coordinator_form': coordinator_form,
  442. })
  443. @check_job_access_permission
  444. @check_job_edition_permission()
  445. def delete_coordinator(request, coordinator):
  446. if request.method != 'POST':
  447. raise PopupException(_('A POST request is required.'))
  448. coordinator.delete()
  449. Submission(request.user, coordinator, request.fs, {}).remove_deployment_dir()
  450. request.info(_('Coordinator deleted!'))
  451. return redirect(reverse('oozie:list_coordinators'))
  452. @check_job_access_permission
  453. @check_job_edition_permission(True)
  454. def edit_coordinator(request, coordinator):
  455. history = History.objects.filter(submitter=request.user, job=coordinator).order_by('-submission_date')
  456. DatasetFormSet = inlineformset_factory(Coordinator, Dataset, form=DatasetForm, max_num=0, can_order=False, can_delete=True)
  457. DataInputFormSet = inlineformset_factory(Coordinator, DataInput, form=DataInputSetForm, max_num=0, can_order=False, can_delete=True)
  458. DataOutputFormSet = inlineformset_factory(Coordinator, DataOutput, form=DataOutputSetForm, max_num=0, can_order=False, can_delete=True)
  459. dataset = Dataset(coordinator=coordinator)
  460. dataset_form = DatasetForm(instance=dataset, prefix='create')
  461. NewDataInputFormSet = inlineformset_factory(Coordinator, DataInput, form=DataInputForm, extra=0, can_order=False, can_delete=False)
  462. NewDataInputFormSet.form = staticmethod(curry(DataInputForm, coordinator=coordinator))
  463. NewDataOutputFormSet = inlineformset_factory(Coordinator, DataOutput, form=DataOutputForm, extra=0, can_order=False, can_delete=False)
  464. NewDataOutputFormSet.form = staticmethod(curry(DataOutputForm, coordinator=coordinator))
  465. if request.method == 'POST':
  466. coordinator_form = CoordinatorForm(request.POST, instance=coordinator, user=request.user)
  467. dataset_formset = DatasetFormSet(request.POST, request.FILES, instance=coordinator)
  468. data_input_formset = DataInputFormSet(request.POST, request.FILES, instance=coordinator)
  469. data_output_formset = DataOutputFormSet(request.POST, request.FILES, instance=coordinator)
  470. new_data_input_formset = NewDataInputFormSet(request.POST, request.FILES, instance=coordinator, prefix='input')
  471. new_data_output_formset = NewDataOutputFormSet(request.POST, request.FILES, instance=coordinator, prefix='output')
  472. if coordinator_form.is_valid() and dataset_formset.is_valid() and data_input_formset.is_valid() and data_output_formset.is_valid() \
  473. and new_data_input_formset.is_valid() and new_data_output_formset.is_valid():
  474. coordinator = coordinator_form.save()
  475. dataset_formset.save()
  476. data_input_formset.save()
  477. data_output_formset.save()
  478. new_data_input_formset.save()
  479. new_data_output_formset.save()
  480. request.info(_('Coordinator saved!'))
  481. return redirect(reverse('oozie:edit_coordinator', kwargs={'coordinator': coordinator.id}))
  482. else:
  483. coordinator_form = CoordinatorForm(instance=coordinator, user=request.user)
  484. dataset_formset = DatasetFormSet(instance=coordinator)
  485. data_input_formset = DataInputFormSet(instance=coordinator)
  486. data_output_formset = DataOutputFormSet(instance=coordinator)
  487. new_data_input_formset = NewDataInputFormSet(queryset=DataInput.objects.none(), instance=coordinator, prefix='input')
  488. new_data_output_formset = NewDataOutputFormSet(queryset=DataOutput.objects.none(), instance=coordinator, prefix='output')
  489. return render('editor/edit_coordinator.mako', request, {
  490. 'coordinator': coordinator,
  491. 'coordinator_form': coordinator_form,
  492. 'dataset_formset': dataset_formset,
  493. 'data_input_formset': data_input_formset,
  494. 'data_output_formset': data_output_formset,
  495. 'dataset_form': dataset_form,
  496. 'new_data_input_formset': new_data_input_formset,
  497. 'new_data_output_formset': new_data_output_formset,
  498. 'history': history,
  499. 'parameters': extract_field_data(coordinator_form['parameters'])
  500. })
  501. @check_job_access_permission
  502. @check_job_edition_permission()
  503. def create_coordinator_dataset(request, coordinator):
  504. """Returns {'status' 0/1, data:html or url}"""
  505. dataset = Dataset(coordinator=coordinator)
  506. response = {'status': -1, 'data': 'None'}
  507. if request.method == 'POST':
  508. dataset_form = DatasetForm(request.POST, instance=dataset, prefix='create')
  509. if dataset_form.is_valid():
  510. dataset_form.save()
  511. response['status'] = 0
  512. response['data'] = reverse('oozie:edit_coordinator', kwargs={'coordinator': coordinator.id})
  513. request.info(_('Dataset created'));
  514. else:
  515. dataset_form = DatasetForm(request.POST, instance=dataset, prefix='create')
  516. else:
  517. ## Bad
  518. response['data'] = _('A POST request is required.')
  519. if response['status'] != 0:
  520. response['data'] = render('editor/create_coordinator_dataset.mako', request, {
  521. 'coordinator': coordinator,
  522. 'dataset_form': dataset_form,
  523. }, force_template=True).content
  524. return HttpResponse(json.dumps(response), mimetype="application/json")
  525. @check_dataset_access_permission
  526. @check_dataset_edition_permission()
  527. def edit_coordinator_dataset(request, dataset):
  528. """Returns HTML for modal to edit datasets"""
  529. if request.method == 'POST':
  530. dataset_form = DatasetForm(request.POST, instance=dataset)
  531. if dataset_form.is_valid():
  532. dataset_form.save()
  533. request.info(_('Dataset modified'));
  534. return redirect(reverse('oozie:edit_coordinator', kwargs={'coordinator': dataset.coordinator.id}))
  535. else:
  536. dataset_form = DatasetForm(request.POST, instance=dataset)
  537. else:
  538. dataset_form = DatasetForm(instance=dataset)
  539. return render('editor/edit_coordinator_dataset.mako', request, {
  540. 'coordinator': dataset.coordinator,
  541. 'dataset_form': dataset_form,
  542. 'path': request.path,
  543. }, force_template=True)
  544. @check_job_access_permission
  545. @check_job_edition_permission()
  546. def create_coordinator_data(request, coordinator, data_type):
  547. """Returns {'status' 0/1, data:html or url}"""
  548. if data_type == 'input':
  549. data_instance = DataInput(coordinator=coordinator)
  550. DataForm = DataInputForm
  551. else:
  552. data_instance = DataOutput(coordinator=coordinator)
  553. DataForm = DataOutputForm
  554. response = {'status': -1, 'data': 'None'}
  555. if request.method == 'POST':
  556. data_form = DataForm(request.POST, instance=data_instance, coordinator=coordinator, prefix=data_type)
  557. if data_form.is_valid():
  558. data_form.save()
  559. response['status'] = 0
  560. response['data'] = reverse('oozie:edit_coordinator', kwargs={'coordinator': coordinator.id})
  561. request.info(_('Coordinator data created'));
  562. else:
  563. data_form = DataForm(request.POST, instance=data_instance, coordinator=coordinator)
  564. response['data'] = data_form.errors
  565. else:
  566. response['data'] = _('A POST request is required.')
  567. return HttpResponse(json.dumps(response), mimetype="application/json")
  568. @check_job_access_permission
  569. def clone_coordinator(request, coordinator):
  570. if request.method != 'POST':
  571. raise PopupException(_('A POST request is required.'))
  572. clone = coordinator.clone(request.user)
  573. response = {'url': reverse('oozie:edit_coordinator', kwargs={'coordinator': clone.id})}
  574. return HttpResponse(json.dumps(response), mimetype="application/json")
  575. @check_job_access_permission
  576. def submit_coordinator(request, coordinator):
  577. ParametersFormSet = formset_factory(ParameterForm, extra=0)
  578. if request.method == 'POST':
  579. params_form = ParametersFormSet(request.POST)
  580. if params_form.is_valid():
  581. mapping = dict([(param['name'], param['value']) for param in params_form.cleaned_data])
  582. job_id = _submit_coordinator(request, coordinator, mapping)
  583. request.info(_('Coordinator submitted'))
  584. return redirect(reverse('oozie:list_oozie_coordinator', kwargs={'job_id': job_id}))
  585. else:
  586. request.error(_('Invalid submission form: %s' % params_form.errors))
  587. else:
  588. parameters = coordinator.find_all_parameters()
  589. params_form = ParametersFormSet(initial=parameters)
  590. popup = render('editor/submit_job_popup.mako', request, {
  591. 'params_form': params_form,
  592. 'action': reverse('oozie:submit_coordinator', kwargs={'coordinator': coordinator.id})
  593. }, force_template=True).content
  594. return HttpResponse(json.dumps(popup), mimetype="application/json")
  595. def _submit_coordinator(request, coordinator, mapping):
  596. try:
  597. submission = Submission(request.user, coordinator.workflow, request.fs, mapping)
  598. wf_dir = submission.deploy()
  599. properties = {'wf_application_path': request.fs.get_hdfs_path(wf_dir)}
  600. properties.update(mapping)
  601. submission = Submission(request.user, coordinator, request.fs, properties=properties)
  602. job_id = submission.run()
  603. History.objects.create_from_submission(submission)
  604. return job_id
  605. except RestException, ex:
  606. raise PopupException(_("Error submitting coordinator %s") % (coordinator,),
  607. detail=ex._headers.get('oozie-error-message', ex))
  608. def resubmit_coordinator(request, oozie_coord_id):
  609. if request.method != 'POST':
  610. raise PopupException(_('A POST request is required.'))
  611. history = History.objects.get(oozie_job_id=oozie_coord_id)
  612. Job.objects.is_accessible_or_exception(request, history.job.id)
  613. coordinator = history.get_coordinator().get_full_node()
  614. properties = history.properties_dict
  615. job_id = _submit_coordinator(request, coordinator, properties)
  616. request.info(_('Coordinator re-submitted'))
  617. return redirect(reverse('oozie:list_oozie_coordinator', kwargs={'job_id': job_id}))
  618. def list_history(request):
  619. """
  620. List the job submission history.
  621. Normal users can only look at their own submissions.
  622. """
  623. history = History.objects
  624. if not request.user.is_superuser:
  625. history = history.filter(submitter=request.user)
  626. history = history.order_by('-submission_date')
  627. return render('editor/list_history.mako', request, {
  628. 'history': history,
  629. })
  630. def list_history_record(request, record_id):
  631. """
  632. List a job submission history.
  633. Normal users can only look at their own jobs.
  634. """
  635. history = History.objects
  636. if not request.user.is_superuser:
  637. history.filter(submitter=request.user)
  638. history = history.get(id=record_id)
  639. return render('editor/list_history_record.mako', request, {
  640. 'record': history,
  641. })
  642. def setup_app(request):
  643. if request.method != 'POST':
  644. raise PopupException(_('A POST request is required.'))
  645. try:
  646. oozie_setup.Command().handle_noargs()
  647. request.info(_('Workspaces and examples installed!'))
  648. except WebHdfsException, e:
  649. raise PopupException(_('The app setup could complete.'), detail=e)
  650. return redirect(reverse('oozie:list_workflows'))