| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641 |
- #!/usr/bin/env python
- # Licensed to Cloudera, Inc. under one
- # or more contributor license agreements. See the NOTICE file
- # distributed with this work for additional information
- # regarding copyright ownership. Cloudera, Inc. licenses this file
- # to you under the Apache License, Version 2.0 (the
- # "License"); you may not use this file except in compliance
- # with the License. You may obtain a copy of the License at
- #
- # http://www.apache.org/licenses/LICENSE-2.0
- #
- # Unless required by applicable law or agreed to in writing, software
- # distributed under the License is distributed on an "AS IS" BASIS,
- # WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
- # See the License for the specific language governing permissions and
- # limitations under the License.
- try:
- import json
- except ImportError:
- import simplejson as json
- import logging
- from django.forms.models import inlineformset_factory, modelformset_factory
- from django.core.urlresolvers import reverse
- from django.db.models import Q
- from django.http import HttpResponse
- from django.shortcuts import redirect
- from django.utils.functional import wraps
- from django.utils.translation import ugettext as _
- from desktop.lib.django_util import render, PopupException, extract_field_data
- from desktop.lib.rest.http_client import RestException
- from desktop.log.access import access_warn
- from hadoop.fs.exceptions import WebHdfsException
- from liboozie.submittion import Submission
- from oozie.models import Workflow, Node, Link, History, Coordinator,\
- Dataset, DataInput, DataOutput, Job, _STD_PROPERTIES_JSON
- from oozie.forms import NodeForm, WorkflowForm, CoordinatorForm, DatasetForm,\
- DataInputForm, DataInputSetForm, DataOutputForm, DataOutputSetForm, LinkForm,\
- DefaultLinkForm, design_form_by_type
- from oozie.conf import SHARE_JOBS
- LOG = logging.getLogger(__name__)
- def can_access_job(request, job_id):
- """
- Logic for testing if a user can access a certain Workflow / Coordinator.
- """
- if job_id is None:
- return
- try:
- job = Job.objects.select_related().get(pk=job_id).get_full_node()
- if not SHARE_JOBS.get() and not request.user.is_superuser \
- and job.owner != request.user.username:
- # TODO is shared perms
- message = _("Permission denied. %(username)s don't have the permissions to access job %(id)s") % \
- {'username': request.user.username, 'id': job.id}
- access_warn(request, message)
- raise PopupException(message)
- else:
- return job
- except Job.DoesNotExist:
- raise PopupException(_('job %(id)s not found') % {'id': job_id})
- def can_modify_job(request, job):
- """Only owners or admins can modify a job."""
- return request.user.is_superuser or job.owner.id == request.user.id
- def check_job_modification(request, job):
- if not can_modify_job(request, job):
- raise PopupException(_('Not allowed to modified this job'))
- def check_job_modification_permission(view_func):
- """
- Decorator ensuring that the user has the permissions to modify a workflow or coordinator.
- Need to appear below @check_job_access_permission
- """
- def decorate(request, *args, **kwargs):
- if 'workflow' in kwargs:
- job_type = 'workflow'
- else:
- job_type = 'coordinator'
- job = kwargs.get(job_type)
- if job is not None:
- check_job_modification(request, job)
- return view_func(request, *args, **kwargs)
- return wraps(view_func)(decorate)
- def check_job_access_permission(view_func):
- """
- Decorator ensuring that the user has access to the workflow or coordinator.
- Arg: 'workflow' or 'coordinator' id.
- Return: the workflow of coordinator or raise an exception
- Notice: its gets an id in input and returns the full object in output (not an id).
- """
- def decorate(request, *args, **kwargs):
- if 'workflow' in kwargs:
- job_type = 'workflow'
- else:
- job_type = 'coordinator'
- job = kwargs.get(job_type)
- if job is not None:
- job = can_access_job(request, job)
- kwargs[job_type] = job
- return view_func(request, *args, **kwargs)
- return wraps(view_func)(decorate)
- def check_action_access_permission(view_func):
- """
- Decorator ensuring that the user has access to the workflow action.
- Arg: 'workflow action' id.
- Return: the workflow action or raise an exception
- Notice: its gets an id in input and returns the full object in output (not an id).
- """
- def decorate(request, *args, **kwargs):
- action_id = kwargs.get('action')
- action = Node.objects.get(id=action_id).get_full_node()
- can_access_job(request, action.workflow.id)
- kwargs['action'] = action
- return view_func(request, *args, **kwargs)
- return wraps(view_func)(decorate)
- def check_action_modification_permission(view_func):
- """
- Decorator ensuring that the user has the permissions to modify a workflow action.
- Need to appear below @check_action_access_permission
- """
- def decorate(request, *args, **kwargs):
- action = kwargs.get('action')
- check_job_modification(request, action.workflow)
- return view_func(request, *args, **kwargs)
- return wraps(view_func)(decorate)
- def list_workflows(request, job_type='workflow'):
- show_install_examples = True
- if job_type == 'coordinators':
- data = Coordinator.objects
- template = "editor/list_coordinators.mako"
- else:
- data = Workflow.objects
- template = "editor/list_workflows.mako"
- if not SHARE_JOBS.get() and not request.user.is_superuser:
- data = data.filter(owner=request.user)
- else:
- data = data.filter(Q(is_shared=True) | Q(owner=request.user))
- data = data.order_by('-last_modified')
- return render(template, request, {
- 'workflows': list(data),
- 'currentuser': request.user,
- 'show_install_examples': show_install_examples,
- })
- def create_workflow(request):
- workflow = Workflow.objects.new_workflow(request.user)
- if request.method == 'POST':
- workflow_form = WorkflowForm(request.POST, instance=workflow)
- if workflow_form.is_valid():
- wf = workflow_form.save()
- Workflow.objects.initialize(wf, request.fs)
- return redirect(reverse('oozie:edit_workflow', kwargs={'workflow': workflow.id}))
- else:
- workflow_form = WorkflowForm(instance=workflow)
- return render('editor/create_workflow.mako', request, {
- 'workflow_form': workflow_form,
- 'workflow': workflow,
- })
- @check_job_access_permission
- def edit_workflow(request, workflow):
- WorkflowFormSet = inlineformset_factory(Workflow, Node, form=NodeForm, max_num=0, can_order=False, can_delete=False)
- history = History.objects.filter(submitter=request.user, job=workflow)
- if request.method == 'POST' and can_modify_job(request, workflow):
- try:
- workflow_form = WorkflowForm(request.POST, instance=workflow)
- actions_formset = WorkflowFormSet(request.POST, request.FILES, instance=workflow)
- if 'clone_action' in request.POST: return clone_action(request, action=request.POST['clone_action'])
- if 'delete_action' in request.POST: return delete_action(request, action=request.POST['delete_action'])
- if 'move_up_action' in request.POST: return move_up_action(request, action=request.POST['move_up_action'])
- if 'move_down_action' in request.POST: return move_down_action(request, action=request.POST['move_down_action'])
- if workflow_form.is_valid() and actions_formset.is_valid():
- workflow_form.save()
- actions_formset.save()
- return redirect(reverse('oozie:list_workflows'))
- except Exception, e:
- request.error(_('Sorry, this operation is not supported: %(error)s') % {'error': e})
- else:
- workflow_form = WorkflowForm(instance=workflow)
- actions_formset = WorkflowFormSet(instance=workflow)
- return render('editor/edit_workflow.mako', request, {
- 'workflow_form': workflow_form,
- 'workflow': workflow,
- 'actions_formset': actions_formset,
- 'graph': workflow.gen_graph(actions_formset.forms),
- 'history': history,
- })
- @check_job_access_permission
- @check_job_modification_permission
- def delete_workflow(request, workflow):
- if request.method != 'POST':
- raise PopupException(_('A POST request is required.'))
- try:
- workflow.coordinator_set.update(workflow=None) # In Django 1.3 could do ON DELETE set NULL
- workflow.save()
- workflow.delete()
- Submission(workflow, request.fs, {}).remove_deployment_dir()
- except Workflow.DoesNotExist:
- LOG.error("Trying to delete non-existent workflow (id %s)" % (workflow,))
- raise PopupException(_('Workflow not found'))
- # TODO notification
- return redirect(reverse('oozie:list_workflows'))
- @check_job_access_permission
- def clone_workflow(request, workflow):
- if request.method != 'POST':
- raise PopupException(_('A POST request is required.'))
- clone = workflow.clone(request.user)
- response = {'url': reverse('oozie:edit_workflow', kwargs={'workflow': clone.id})}
- return HttpResponse(json.dumps(response), mimetype="application/json")
- @check_job_access_permission
- def submit_workflow(request, workflow):
- if request.method != 'POST':
- raise PopupException(_('A POST request is required.'))
- try:
- mapping = dict(request.POST.iteritems())
- submission = Submission(workflow, request.fs, mapping)
- job_id = submission.run()
- except RestException, ex:
- raise PopupException(_("Error submitting workflow %s") % (workflow,),
- detail=ex._headers.get('oozie-error-message', ex))
- History.objects.create_from_submission(submission)
- request.info(_('Workflow submitted'))
- return redirect(reverse('oozie:list_oozie_workflow', kwargs={'job_id': job_id}))
- @check_job_access_permission
- def get_workflow_parameters(request, workflow):
- """
- Return the parameters found in the workflow as a JSON dictionary of {param_key : label}.
- This expects an Ajax call.
- """
- params = workflow.find_parameters()
- params_with_labels = dict((p, p.upper()) for p in params)
- return render('dont_care_for_ajax', request, { 'params': params_with_labels })
- @check_job_access_permission
- def new_action(request, workflow, node_type, parent_action_id):
- ActionForm = design_form_by_type(node_type)
- if request.method == 'POST':
- action_form = ActionForm(request.POST)
- if action_form.is_valid():
- action = action_form.save(commit=False)
- action.node_type = node_type
- action.workflow = workflow
- action.save()
- workflow.add_action(action, parent_action_id)
- return redirect(reverse('oozie:edit_workflow', kwargs={'workflow': workflow.id}))
- else:
- action_form = ActionForm()
- return render('editor/edit_workflow_action.mako', request, {
- 'workflow': workflow,
- 'job_properties': extract_field_data(action_form['job_properties']),
- 'files': extract_field_data(action_form['files']),
- 'archives': extract_field_data(action_form['archives']),
- 'params': 'params' in action_form.fields and extract_field_data(action_form['params']) or '[]',
- 'action_form': action_form,
- 'node_type': node_type,
- 'properties_hint': _STD_PROPERTIES_JSON,
- 'form_url': reverse('oozie:new_action', kwargs={'workflow': workflow.id,
- 'node_type': node_type,
- 'parent_action_id': parent_action_id}),
- })
- @check_action_access_permission
- def edit_action(request, action):
- ActionForm = design_form_by_type(action.node_type)
- if request.method == 'POST' and can_modify_job(request, action.workflow):
- action_form = ActionForm(request.POST, instance=action)
- if action_form.is_valid():
- action = action_form.save()
- return redirect(reverse('oozie:edit_workflow', kwargs={'workflow': action.workflow.id}))
- else:
- action_form = ActionForm(instance=action)
- return render('editor/edit_workflow_action.mako', request, {
- 'workflow': action.workflow,
- 'job_properties': extract_field_data(action_form['job_properties']),
- 'files': extract_field_data(action_form['files']),
- 'archives': extract_field_data(action_form['archives']),
- 'params': 'params' in action_form.fields and extract_field_data(action_form['params']) or '[]',
- 'action_form': action_form,
- 'node_type': action.node_type,
- 'properties_hint': _STD_PROPERTIES_JSON,
- 'form_url': reverse('oozie:edit_action', kwargs={'action': action.id}),
- })
- @check_action_access_permission
- @check_action_modification_permission
- def edit_workflow_fork(request, action):
- fork = action
- LinkFormSet = modelformset_factory(Link, form=LinkForm, max_num=0)
- if request.method == 'POST':
- link_formset = LinkFormSet(request.POST)
- default_link_form = DefaultLinkForm(request.POST, action=fork)
- if link_formset.is_valid():
- is_decision = fork.has_decisions()
- link_formset.save()
- if not is_decision and fork.has_decisions():
- default_link = default_link_form.save(commit=False)
- default_link.parent = fork
- default_link.name = 'default'
- default_link.comment = 'default'
- default_link.save()
- fork.convert_to_decision()
- fork.update_description()
- return redirect(reverse('oozie:edit_workflow', kwargs={'workflow': fork.workflow.id}))
- else:
- link_formset = LinkFormSet(queryset=fork.get_children_links().exclude(name__in=['related', 'default']))
- default_link = Link(parent=fork, name='default', comment='default')
- default_link_form = DefaultLinkForm(action=fork, instance=default_link)
- return render('editor/edit_workflow_fork.mako', request, {
- 'workflow': fork.workflow,
- 'fork': fork,
- 'link_formset': link_formset,
- 'default_link_form': default_link_form,
- })
- @check_action_access_permission
- @check_action_modification_permission
- def delete_action(request, action):
- if request.method == 'POST':
- action.workflow.delete_action(action)
- return redirect(reverse('oozie:edit_workflow', kwargs={'workflow': action.workflow.id}))
- else:
- raise PopupException(_('A POST request is required.'))
- @check_action_access_permission
- def clone_action(request, action):
- # Really weird: action is like a clone object with the old id here
- action_id = action.id
- workflow = action.workflow
- clone = action.clone()
- workflow.add_action(clone, action_id)
- return redirect(reverse('oozie:edit_workflow', kwargs={'workflow': workflow.id}))
- @check_action_access_permission
- @check_action_modification_permission
- def move_up_action(request, action):
- if request.method == 'POST':
- action.workflow.move_action_up(action)
- return redirect(reverse('oozie:edit_workflow', kwargs={'workflow': action.workflow.id}))
- else:
- raise PopupException(_('A POST request is required.'))
- @check_action_access_permission
- @check_action_modification_permission
- def move_down_action(request, action):
- if request.method == 'POST':
- action.workflow.move_action_down(action)
- return redirect(reverse('oozie:edit_workflow', kwargs={'workflow': action.workflow.id}))
- else:
- raise PopupException(_('A POST request is required.'))
- @check_job_access_permission
- def create_coordinator(request, workflow=None):
- if workflow is not None:
- coordinator = Coordinator(owner=request.user, workflow=workflow)
- else:
- coordinator = Coordinator(owner=request.user)
- if request.method == 'POST':
- coordinator_form = CoordinatorForm(request.POST, instance=coordinator)
- if coordinator_form.is_valid():
- coordinator = coordinator_form.save()
- return redirect(reverse('oozie:edit_coordinator', kwargs={'coordinator': coordinator.id}))
- else:
- coordinator_form = CoordinatorForm(instance=coordinator)
- return render('editor/create_coordinator.mako', request, {
- 'coordinator': coordinator,
- 'coordinator_form': coordinator_form,
- })
- @check_job_access_permission
- @check_job_modification_permission
- def edit_coordinator(request, coordinator):
- history = History.objects.filter(submitter=request.user, job=coordinator)
- DatasetFormSet = inlineformset_factory(Coordinator, Dataset, form=DatasetForm, max_num=0, can_order=False, can_delete=True)
- DataInputFormSet = inlineformset_factory(Coordinator, DataInput, form=DataInputSetForm, max_num=0, can_order=False, can_delete=True)
- DataOutputFormSet = inlineformset_factory(Coordinator, DataOutput, form=DataOutputSetForm, max_num=0, can_order=False, can_delete=True)
- dataset = Dataset(coordinator=coordinator)
- dataset_form = DatasetForm(instance=dataset)
- data_input = DataInput(coordinator=coordinator)
- data_input_form = DataInputForm(instance=data_input, coordinator=coordinator)
- data_output = DataOutput(coordinator=coordinator)
- data_output_form = DataOutputForm(instance=data_output, coordinator=coordinator)
- if request.method == 'POST':
- coordinator_form = CoordinatorForm(request.POST, instance=coordinator)
- dataset_formset = DatasetFormSet(request.POST, request.FILES, instance=coordinator)
- data_input_formset = DataInputFormSet(request.POST, request.FILES, instance=coordinator)
- data_output_formset = DataOutputFormSet(request.POST, request.FILES, instance=coordinator)
- if coordinator_form.is_valid() and dataset_formset.is_valid() and data_input_formset.is_valid() and data_output_formset.is_valid():
- coordinator = coordinator_form.save()
- dataset_formset.save()
- data_input_formset.save()
- data_output_formset.save()
- return redirect(reverse('oozie:list_coordinator'))
- else:
- coordinator_form = CoordinatorForm(instance=coordinator)
- dataset_formset = DatasetFormSet(instance=coordinator)
- data_input_formset = DataInputFormSet(instance=coordinator)
- data_output_formset = DataOutputFormSet(instance=coordinator)
- return render('editor/edit_coordinator.mako', request, {
- 'coordinator': coordinator,
- 'coordinator_form': coordinator_form,
- 'dataset_formset': dataset_formset,
- 'data_input_formset': data_input_formset,
- 'data_output_formset': data_output_formset,
- 'dataset_form': dataset_form,
- 'data_input_form': data_input_form,
- 'data_output_form': data_output_form,
- 'history': history,
- })
- @check_job_access_permission
- @check_job_modification_permission
- def create_coordinator_dataset(request, coordinator):
- """Returns {'status' 0/1, data:html or url}"""
- dataset = Dataset(coordinator=coordinator)
- response = {'status': -1, 'data': 'None'}
- if request.method == 'POST':
- dataset_form = DatasetForm(request.POST, instance=dataset)
- if dataset_form.is_valid():
- dataset_form.save()
- response['status'] = 0
- response['data'] = reverse('oozie:edit_coordinator', kwargs={'coordinator': coordinator.id})
- request.info(_('Dataset created'));
- else:
- dataset_form = DatasetForm(request.POST, instance=dataset)
- else:
- ## Bad
- response['data'] = _('A POST request is required.')
- if response['status'] != 0:
- response['data'] = render('editor/create_coordinator_dataset.mako', request, {
- 'coordinator': coordinator,
- 'dataset_form': dataset_form,
- }, force_template=True).content
- return HttpResponse(json.dumps(response), mimetype="application/json")
- @check_job_access_permission
- @check_job_modification_permission
- def create_coordinator_data(request, coordinator, data_type):
- """Returns {'status' 0/1, data:html or url}"""
- if data_type == 'input':
- data_instance = DataInput(coordinator=coordinator)
- DataForm = DataInputForm
- else:
- data_instance = DataOutput(coordinator=coordinator)
- DataForm = DataOutputForm
- response = {'status': -1, 'data': 'None'}
- if request.method == 'POST':
- data_form = DataForm(request.POST, instance=data_instance, coordinator=coordinator)
- if data_form.is_valid():
- data_form.save()
- response['status'] = 0
- response['data'] = reverse('oozie:edit_coordinator', kwargs={'coordinator': coordinator.id})
- request.info(_('Coordinator data created'));
- else:
- data_form = DataForm(request.POST, instance=data_instance, coordinator=coordinator)
- else:
- ## Bad
- response['data'] = _('A POST request is required.')
- if response['status'] != 0:
- response['data'] = render('editor/create_coordinator_data.mako', request, {
- 'coordinator': coordinator,
- 'form': data_form, },
- force_template=True).content
- return HttpResponse(json.dumps(response), mimetype="application/json")
- @check_job_access_permission
- def submit_coordinator(request, coordinator):
- if request.method != 'POST':
- raise PopupException(_('A POST request is required.'))
- try:
- if not coordinator.workflow.is_deployed(request.fs):
- submission = Submission(coordinator.workflow, request.fs, request.POST)
- wf_dir = submission.deploy()
- coordinator.workflow.deployment_dir = wf_dir
- coordinator.workflow.save()
- coordinator.deployment_dir = coordinator.workflow.deployment_dir
- properties = {'wf_application_path': coordinator.workflow.deployment_dir}
- properties.update(dict(request.POST.iteritems()))
- submission = Submission(coordinator, request.fs, properties=properties)
- job_id = submission.run()
- except RestException, ex:
- raise PopupException(_("Error submitting coordinator %s") % (coordinator,),
- detail=ex._headers['oozie-error-message'])
- History.objects.create_from_submission(submission)
- request.info(_('Coordinator submitted'))
- return redirect(reverse('oozie:list_oozie_coordinator', kwargs={'job_id': job_id}))
- def list_history(request):
- """
- List the job submission history.
- Normal users can only look at their own submissions.
- """
- history = History.objects
- if not request.user.is_superuser:
- history = history.filter(submitter=request.user)
- history = history.order_by('-submission_date')
- return render('editor/list_history.mako', request, {
- 'history': history,
- })
- def list_history_record(request, record_id):
- """
- List a job submission history.
- Normal users can only look at their own jobs.
- """
- history = History.objects
- if not request.user.is_superuser:
- history.filter(submitter=request.user)
- history = history.get(id=record_id)
- return render('editor/list_history_record.mako', request, {
- 'record': history,
- })
- def setup(request):
- """Installs oozie examples."""
- if request.method != 'POST':
- raise PopupException(_('A POST request is required.'))
- try:
- # Warning: below will modify fs.user
- #oozie_setup.Command().handle_noargs()
- pass
- except WebHdfsException, e:
- raise PopupException(_('The examples could not be installed.'), detail=e)
- return redirect(reverse('oozie:list_workflows'))
|