소스 검색

HUE-817 [oozie] Resubmit a workflow from a certain step

Rerun any workflow
Select a list of successful actions to skip or do not rerun the failed action
Prefilled the list of actions to skip
Preselect previously selected actions to skip
Adding permissions tests
Adding real resubmissions tests
Removed previously resubmitting mechanism
Hiding 'oozie.use.system.libpath' in submission or rerun popups
Fix deployments directories permissions to 711
Romain Rigaux 13 년 전
부모
커밋
9545e7c

+ 37 - 1
apps/oozie/src/oozie/forms.py

@@ -29,7 +29,21 @@ LOG = logging.getLogger(__name__)
 
 class ParameterForm(forms.Form):
   name = forms.CharField(max_length=40, widget=forms.widgets.HiddenInput())
-  value = forms.CharField(max_length=40, required=False)
+  value = forms.CharField(max_length=100, required=False)
+
+  NON_PARAMETERS = ('user.name',
+                    'oozie.wf.rerun.failnodes',
+                    'oozie.wf.rerun.skip.nodes',
+                    'oozie.wf.application.path',
+                    'jobTracker',
+                    'nameNode',
+                    'hue-id-w')
+
+  @staticmethod
+  def get_initial_params(conf_dict):
+    params = filter(lambda key: key not in ParameterForm.NON_PARAMETERS, conf_dict.keys())
+
+    return [{'name': name, 'value': conf_dict[name]} for name in params]
 
 
 class WorkflowForm(forms.ModelForm):
@@ -325,6 +339,28 @@ _node_type_TO_FORM_CLS = {
 }
 
 
+class RerunForm(forms.Form):
+  skip_nodes = forms.MultipleChoiceField(required=False)
+
+  def __init__(self, *args, **kwargs):
+    oozie_workflow = kwargs.pop('oozie_workflow')
+
+    # Build list of skip nodes
+    decisions = filter(lambda node: node.type == 'switch', oozie_workflow.get_control_flow_actions())
+    working_actions = oozie_workflow.get_working_actions()
+    skip_nodes = []
+
+    for action in decisions + working_actions:
+      if action.status == 'OK':
+        skip_nodes.append((action.name, action.name))
+    initial_skip_nodes = oozie_workflow.conf_dict.get('oozie.wf.rerun.skip.nodes', '').split()
+
+    super(RerunForm, self).__init__(*args, **kwargs)
+
+    self.fields['skip_nodes'].choices = skip_nodes
+    self.fields['skip_nodes'].initial = initial_skip_nodes
+
+
 def design_form_by_type(node_type):
   return _node_type_TO_FORM_CLS[node_type]
 

+ 25 - 13
apps/oozie/src/oozie/templates/dashboard/list_oozie_workflow.mako

@@ -42,11 +42,11 @@ ${ layout.menubar(section='dashboard') }
       ${ _('Workflow') }
     </div>
     <div class="span3">
-      %if hue_workflow is not None:
+      % if hue_workflow is not None:
         <a title="${ _('Edit workflow') }" href="${ hue_workflow.get_absolute_url() }">${ hue_workflow }</a>
       % else:
         ${ oozie_workflow.appName }
-      %endif
+      % endif
     </div>
   </div>
 
@@ -113,7 +113,6 @@ ${ layout.menubar(section='dashboard') }
       ${ _('Manage') }
     </div>
     <div class="span3">
-      <form action="${ url('oozie:resubmit_workflow', oozie_wf_id=oozie_workflow.id) }" method="post">
       % if oozie_workflow.is_running():
         <a title="${_('Kill %(workflow)s') % dict(workflow=oozie_workflow.id)}"
           id="kill-workflow"
@@ -126,11 +125,14 @@ ${ layout.menubar(section='dashboard') }
             ${_('Kill')}
         </a>
       % else:
-        <button type="submit" class="btn">
-          ${ _('Resubmit') }
-        </button>
+        % if oozie_workflow.id:
+          <a title="${ _('Rerun the same workflow') }" class="btn" id="rerun-btn"
+            data-rerun-url="${ url('oozie:rerun_oozie_job', job_id=oozie_workflow.id, app_path=oozie_workflow.appPath) }">
+            ${ _('Rerun') }
+          </a>
+        % endif
+        <div id="rerun-wf-modal" class="modal hide"></div>
       % endif
-      </form>
     </div>
   </div>
   % endif
@@ -236,11 +238,11 @@ ${ layout.menubar(section='dashboard') }
           <tbody>
             <tr>
               <td>${ _('Group') }</td>
-              <td>${ oozie_workflow.group }</td>
+              <td>${ oozie_workflow.group or '-' }</td>
             </tr>
             <tr>
               <td>${ _('External Id') }</td>
-              <td>${ oozie_workflow.externalId or "-" }</td>
+              <td>${ oozie_workflow.externalId or '-' }</td>
             </tr>
             <tr>
               <td>${ _('Start Time') }</td>
@@ -294,12 +296,12 @@ ${ layout.menubar(section='dashboard') }
 <script type="text/javascript">
   $(document).ready(function() {
     $(".action-link").click(function(){
-      window.location = $(this).attr('data-edit');
+      window.location = $(this).data('edit');
     });
 
     $(".confirmationModal").click(function(){
       var _this = $(this);
-      $("#confirmation .message").text(_this.attr("data-confirmation-message"));
+      $("#confirmation .message").text(_this.data("confirmation-message"));
       $("#confirmation").modal("show");
       $("#confirmation a.btn-primary").click(function() {
         _this.trigger('confirmation');
@@ -308,8 +310,8 @@ ${ layout.menubar(section='dashboard') }
 
     $("#kill-workflow").bind('confirmation', function() {
       var _this = this;
-      $.post($(this).attr("data-url"),
-        { 'notification': $(this).attr("data-message") },
+      $.post($(this).data("url"),
+        { 'notification': $(this).data("message") },
         function(response) {
           if (response['status'] != 0) {
             $.jHueNotify.error("${ _('Error: ') }" + response['data']);
@@ -321,6 +323,16 @@ ${ layout.menubar(section='dashboard') }
       return false;
     });
 
+    $('#rerun-btn').click(function() {
+      var _action = $(this).data("rerun-url");
+
+      $.get(_action,  function(response) {
+          $('#rerun-wf-modal').html(response);
+          $('#rerun-wf-modal').modal('show');
+        }
+      );
+     });
+
     $("a[data-row-selector='true']").jHueRowSelector();
   });
 </script>

+ 95 - 0
apps/oozie/src/oozie/templates/dashboard/rerun_job_popup.mako

@@ -0,0 +1,95 @@
+## 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.
+
+<%!
+  from django.utils.translation import ugettext as _
+%>
+
+<%namespace name="utils" file="../utils.inc.mako" />
+
+
+<form action="${ action }" method="POST">
+
+  <div class="modal-header">
+    <a href="#" class="close" data-dismiss="modal">&times;</a>
+    <h2>${ _('Rerun this job?') }</h2>
+  </div>
+
+  <div class="modal-body">
+    <fieldset>
+      <div id="config-container">
+        <h3>${ _('Select all actions') }</h3>
+          <div class="fieldWrapper">
+            <div class="row-fluid">
+              <div class="span6">
+                <input type="radio" name="rerun_form_choice" value="skip_nodes" id="skip_nodes" checked>
+                  ${ _('Skip successful') }
+                  ${ utils.render_field(rerun_form['skip_nodes'], show_label=False) }
+                </div>
+              <div class="span6">
+                <input type="radio" name="rerun_form_choice" value="fail_nodes" id="fail_nodes">
+                ${ _('Exclude failed') }
+              </div>
+            </div>
+          </div>
+      </div>
+
+      <div id="param-container">
+        ${ params_form.management_form }
+
+        % if params_form.forms:
+          <h3>${ _('Variables') }</h3>
+          % for form in params_form.forms:
+            % for hidden in form.hidden_fields():
+              ${ hidden }
+            % endfor
+            <div class="fieldWrapper">
+              <div class="row-fluid
+                % if form['name'].form.initial.get('name') == 'oozie.use.system.libpath':
+                  hide
+                % endif
+                ">
+                <div class="span6">
+                  ${ form['name'].form.initial.get('name') }
+                </div>
+                <div class="span6">
+                  ${ utils.render_field(form['value'], show_label=False) }
+                </div>
+              </div>
+            </div>
+          % endfor
+        % endif
+      </div>
+    </fieldset>
+  </div>
+
+  <div class="modal-footer">
+    <a href="#" class="btn secondary" data-dismiss="modal">${ _('Cancel') }</a>
+    <input id="submit-btn" type="submit" class="btn btn-primary" value="${ _('Submit') }"/>
+  </div>
+</form>
+
+<script type="text/javascript" charset="utf-8">
+    $(document).ready(function(){
+        $("#id_skip_nodes").jHueSelector({
+            selectAllLabel: "${_('Select all')}",
+            searchPlaceholder: "${_('Search')}",
+            noChoicesFound: "${_('No successful actions found.')}",
+            width:250,
+            height:100
+        });
+    });
+</script>

+ 5 - 1
apps/oozie/src/oozie/templates/editor/submit_job_popup.mako

@@ -37,7 +37,11 @@
             ${ hidden }
           % endfor
           <div class="fieldWrapper">
-            <div class="row-fluid">
+            <div class="row-fluid
+              % if form['name'].form.initial.get('name') == 'oozie.use.system.libpath':
+                hide
+              % endif
+              ">
               <div class="span6">
                 ${ form['name'].form.initial.get('name') }
               </div>

+ 42 - 34
apps/oozie/src/oozie/tests.py

@@ -829,30 +829,6 @@ class TestPermissions(OozieBase):
     finally:
       finish()
 
-    # Resubmit
-    finish = SHARE_JOBS.set_for_testing(False)
-    try:
-      history, created = History.objects.get_or_create(job=self.wf, oozie_job_id=MockOozieApi.WORKFLOW_IDS[0],
-                                                       defaults={'submitter': User.objects.get(username='test'), 'properties': '[]'})
-      job_id = history.oozie_job_id
-      response = client_not_me.post(reverse('oozie:resubmit_workflow', args=[job_id]))
-      assert_true('Permission denied' in response.content, response.content)
-    finally:
-      finish()
-
-    finish = SHARE_JOBS.set_for_testing(True)
-    try:
-      try:
-        history, created = History.objects.get_or_create(job=self.wf, oozie_job_id=MockOozieApi.WORKFLOW_IDS[0],
-                                                         defaults={'submitter': User.objects.get(username='test'), 'properties': '[]'})
-        job_id = history.oozie_job_id
-        response = client_not_me.post(reverse('oozie:resubmit_workflow', args=[job_id]))
-        assert_false('Permission denied' in response.content, response.content)
-      except IOError:
-        pass
-    finally:
-      finish()
-
     # Delete
     finish = SHARE_JOBS.set_for_testing(False)
     try:
@@ -1080,15 +1056,33 @@ class TestOozieSubmissions(OozieBase):
 
   def test_submit_mapreduce_action(self):
     wf = Workflow.objects.get(name='MapReduce')
+    post_data = {u'form-MAX_NUM_FORMS': [u''], u'form-INITIAL_FORMS': [u'1'],
+                 u'form-0-name': [u'REDUCER_SLEEP_TIME'], u'form-0-value': [u'1'], u'form-TOTAL_FORMS': [u'1']}
 
-    response = self.c.post(reverse('oozie:submit_workflow', args=[wf.id]),
-                           data={u'form-MAX_NUM_FORMS': [u''],
-                                u'form-INITIAL_FORMS': [u'1'], u'form-0-name': [u'REDUCER_SLEEP_TIME'],
-                                u'form-0-value': [u'1'], u'form-TOTAL_FORMS': [u'1']},
-                           follow=True)
+    response = self.c.post(reverse('oozie:submit_workflow', args=[wf.id]), data=post_data, follow=True)
     job = OozieServerProvider.wait_until_completion(response.context['oozie_workflow'].id)
     assert_equal('SUCCEEDED', job.status)
 
+    # Rerun with default options
+    post_data.update({u'rerun_form_choice': [u'skip_nodes']})
+
+    response = self.c.post(reverse('oozie:rerun_oozie_job', kwargs={'job_id': job.id, 'app_path': job.appPath}), data=post_data, follow=True)
+    job = OozieServerProvider.wait_until_completion(response.context['oozie_workflow'].id)
+    assert_equal('SUCCEEDED', job.status)
+
+    # Rerun with skip OK actions skipped
+    post_data.update({u'rerun_form_choice': [u'skip_nodes'], u'skip_nodes': [u'Sleep']})
+
+    response = self.c.post(reverse('oozie:rerun_oozie_job', kwargs={'job_id': job.id, 'app_path': job.appPath}), data=post_data, follow=True)
+    job = OozieServerProvider.wait_until_completion(response.context['oozie_workflow'].id)
+    assert_equal('SUCCEEDED', job.status)
+
+    # Rerun with failed nodes too
+    post_data.update({u'rerun_form_choice': [u'failed_nodes']})
+
+    response = self.c.post(reverse('oozie:rerun_oozie_job', kwargs={'job_id': job.id, 'app_path': job.appPath}), data=post_data, follow=True)
+    job = OozieServerProvider.wait_until_completion(response.context['oozie_workflow'].id)
+
 
   def test_submit_java_action(self):
     wf = Workflow.objects.get(name='Sequential Java')
@@ -1136,21 +1130,21 @@ class TestDashboard(OozieMockBase):
     # Kill button in response
     response = self.c.get(reverse('oozie:list_oozie_workflow', args=[MockOozieApi.WORKFLOW_IDS[0]]), {}, follow=True)
     assert_true(('%s/kill' % MockOozieApi.WORKFLOW_IDS[0]) in response.content, response.content)
-    assert_false('Resubmit' in response.content, response.content)
+    assert_false('Rerun' in response.content, response.content)
 
-    # Resubmit button in response
+    # Rerun button in response
     response = self.c.get(reverse('oozie:list_oozie_workflow', args=[MockOozieApi.WORKFLOW_IDS[1]]), {}, follow=True)
     assert_false(('%s/kill' % MockOozieApi.WORKFLOW_IDS[1]) in response.content, response.content)
-    assert_true('Resubmit' in response.content, response.content)
+    assert_true('Rerun' in response.content, response.content)
 
 
   def test_manage_coordinator_dashboard(self):
     # Kill button in response
     response = self.c.get(reverse('oozie:list_oozie_coordinator', args=[MockOozieApi.COORDINATOR_IDS[0]]), {}, follow=True)
     assert_true(('%s/kill' % MockOozieApi.COORDINATOR_IDS[0]) in response.content, response.content)
-    assert_false('Resubmit' in response.content, response.content)
+    assert_false('Rerun' in response.content, response.content)
 
-    # Resubmit button in response
+    # Rerun button in response
     response = self.c.get(reverse('oozie:list_oozie_coordinator', args=[MockOozieApi.COORDINATOR_IDS[1]]), {}, follow=True)
     assert_false(('%s/kill' % MockOozieApi.COORDINATOR_IDS[1]) in response.content, response.content)
     assert_true('Resubmit' in response.content, response.content)
@@ -1202,6 +1196,11 @@ class TestDashboard(OozieMockBase):
     response = self.c.get(reverse('oozie:list_oozie_workflows'))
     assert_true('WordCount1' in response.content, response.content)
 
+    # Rerun
+    response = self.c.get(reverse('oozie:rerun_oozie_job', kwargs={'job_id': MockOozieApi.WORKFLOW_IDS[0],
+                                                                   'app_path': MockOozieApi.JSON_WORKFLOW_LIST[0]['appPath']}))
+    assert_false('Permission denied.' in response.content, response.content)
+
     # Login as someone else
     client_not_me = make_logged_in_client(username='not_me', is_superuser=False, groupname='test', recreate=True)
     grant_access("not_me", "not_me", "oozie")
@@ -1209,12 +1208,21 @@ class TestDashboard(OozieMockBase):
     response = client_not_me.get(reverse('oozie:list_oozie_workflows'))
     assert_false('WordCount1' in response.content, response.content)
 
+    # Rerun
+    response = client_not_me.get(reverse('oozie:rerun_oozie_job', kwargs={'job_id': MockOozieApi.WORKFLOW_IDS[0],
+                                                                          'app_path': MockOozieApi.JSON_WORKFLOW_LIST[0]['appPath']}))
+    assert_true('Permission denied.' in response.content, response.content)
+
     # Add read only access
     add_permission("not_me", "dashboard_jobs_access", "dashboard_jobs_access", "oozie")
 
     response = client_not_me.get(reverse('oozie:list_oozie_workflows'))
     assert_true('WordCount1' in response.content, response.content)
 
+    # Rerun
+    response = client_not_me.get(reverse('oozie:rerun_oozie_job', kwargs={'job_id': MockOozieApi.WORKFLOW_IDS[0],
+                                                                          'app_path': MockOozieApi.JSON_WORKFLOW_LIST[0]['appPath']}))
+    assert_false('Permission denied.' in response.content, response.content)
 
   def test_workflow_permissions(self):
     response = self.c.get(reverse('oozie:list_oozie_workflow', args=[MockOozieApi.WORKFLOW_IDS[0]]))

+ 1 - 1
apps/oozie/src/oozie/urls.py

@@ -37,7 +37,6 @@ urlpatterns += patterns(
   url(r'^clone_workflow/(?P<workflow>\d+)$', 'clone_workflow', name='clone_workflow'),
   url(r'^submit_workflow/(?P<workflow>\d+)$', 'submit_workflow', name='submit_workflow'),
   url(r'^schedule_workflow/(?P<workflow>\d+)$', 'schedule_workflow', name='schedule_workflow'),
-  url(r'^resubmit_workflow/(?P<oozie_wf_id>[-\w]+)$', 'resubmit_workflow', name='resubmit_workflow'),
 
   url(r'^import_action/(?P<workflow>\d+)/(?P<parent_action_id>\d+)$', 'import_action', name='import_action'),
 
@@ -67,5 +66,6 @@ urlpatterns += patterns(
   url(r'^list_oozie_workflow/(?P<job_id>[-\w]+)/(?P<coordinator_job_id>[-\w]+)?$', 'list_oozie_workflow', name='list_oozie_workflow'),
   url(r'^list_oozie_coordinator/(?P<job_id>[-\w]+)$', 'list_oozie_coordinator', name='list_oozie_coordinator'),
   url(r'^list_oozie_workflow_action/(?P<action>[-\w@]+)$', 'list_oozie_workflow_action', name='list_oozie_workflow_action'),
+  url(r'^rerun_oozie_job/(?P<job_id>[-\w]+)/(?P<app_path>.+?)$', 'rerun_oozie_job', name='rerun_oozie_job'),
   url(r'^manage_oozie_jobs/(?P<job_id>[-\w]+)/(?P<action>(start|suspend|resume|kill|rerun))$', 'manage_oozie_jobs', name='manage_oozie_jobs'),
 )

+ 56 - 0
apps/oozie/src/oozie/views/dashboard.py

@@ -21,15 +21,21 @@ except ImportError:
   import simplejson as json
 import logging
 
+from django.forms.formsets import formset_factory
 from django.http import HttpResponse
 from django.utils.functional import wraps
 from django.utils.translation import ugettext as _
+from django.core.urlresolvers import reverse
+from django.shortcuts import redirect
 
 from desktop.lib.django_util import render
 from desktop.lib.exceptions import PopupException
 from desktop.lib.rest.http_client import RestException
 from desktop.log.access import access_warn
 from liboozie.oozie_api import get_oozie
+from liboozie.submittion import Submission
+from oozie.forms import RerunForm, ParameterForm
+
 
 from oozie.conf import OOZIE_JOBS_COUNT
 from oozie.models import History, Job
@@ -182,6 +188,56 @@ def list_oozie_workflow_action(request, action):
   })
 
 
+@show_oozie_error
+def rerun_oozie_job(request, job_id, app_path):
+  ParametersFormSet = formset_factory(ParameterForm, extra=0)
+  oozie_workflow = check_job_access_permission(request, job_id)
+
+  if request.method == 'POST':
+    rerun_form = RerunForm(request.POST, oozie_workflow=oozie_workflow)
+    params_form = ParametersFormSet(request.POST)
+
+    if sum([rerun_form.is_valid(), params_form.is_valid()]) == 2:
+      args = {}
+
+      if request.POST['rerun_form_choice'] == 'fail_nodes':
+        args['fail_nodes'] = 'false'
+      else:
+        args['skip_nodes'] = ','.join(rerun_form.cleaned_data['skip_nodes'])
+      args['deployment_dir'] = app_path
+
+      mapping = dict([(param['name'], param['value']) for param in params_form.cleaned_data])
+
+      _rerun_workflow(request, job_id, args, mapping)
+
+      request.info(_('Workflow re-running!'))
+      return redirect(reverse('oozie:list_oozie_workflow', kwargs={'job_id': job_id}))
+    else:
+      request.error(_('Invalid submission form: %s %s' % (rerun_form.errors, params_form.errors)))
+  else:
+    rerun_form = RerunForm(oozie_workflow=oozie_workflow)
+    initial_params = ParameterForm.get_initial_params(oozie_workflow.conf_dict)
+    params_form = ParametersFormSet(initial=initial_params)
+
+  popup = render('dashboard/rerun_job_popup.mako', request, {
+                   'rerun_form': rerun_form,
+                   'params_form': params_form,
+                   'action': reverse('oozie:rerun_oozie_job', kwargs={'job_id': job_id, 'app_path': app_path}),
+                 }, force_template=True).content
+
+  return HttpResponse(json.dumps(popup), mimetype="application/json")
+
+
+def _rerun_workflow(request, oozie_id, run_args, mapping):
+  try:
+    submission = Submission(user=request.user, fs=request.fs, properties=mapping, oozie_id=oozie_id)
+    job_id = submission.rerun(**run_args)
+    return job_id
+  except RestException, ex:
+    raise PopupException(_("Error rerunning workflow %s") % (oozie_id,),
+                         detail=ex._headers.get('oozie-error-message', ex))
+
+
 def split_oozie_jobs(oozie_jobs):
   jobs = {}
   jobs_running = []

+ 3 - 18
apps/oozie/src/oozie/views/editor.py

@@ -158,6 +158,7 @@ def clone_workflow(request, workflow):
   return HttpResponse(json.dumps(response), mimetype="application/json")
 
 
+
 @check_job_access_permission()
 def submit_workflow(request, workflow):
   ParametersFormSet = formset_factory(ParameterForm, extra=0)
@@ -179,13 +180,12 @@ def submit_workflow(request, workflow):
     params_form = ParametersFormSet(initial=parameters)
 
   popup = render('editor/submit_job_popup.mako', request, {
-                 'params_form': params_form,
-                 'action': reverse('oozie:submit_workflow', kwargs={'workflow': workflow.id})
+                   'params_form': params_form,
+                   'action': reverse('oozie:submit_workflow', kwargs={'workflow': workflow.id})
                  }, force_template=True).content
   return HttpResponse(json.dumps(popup), mimetype="application/json")
 
 
-
 def _submit_workflow(request, workflow, mapping):
   try:
     submission = Submission(request.user, workflow, request.fs, mapping)
@@ -202,21 +202,6 @@ def _submit_workflow(request, workflow, mapping):
   return redirect(reverse('oozie:list_oozie_workflow', kwargs={'job_id': job_id}))
 
 
-def resubmit_workflow(request, oozie_wf_id):
-  if request.method != 'POST':
-    raise PopupException(_('A POST request is required.'))
-
-  history = History.objects.get(oozie_job_id=oozie_wf_id)
-  Job.objects.is_accessible_or_exception(request, history.job.id)
-
-  workflow = history.get_workflow().get_full_node()
-  properties = history.properties_dict
-  job_id = _submit_workflow(request, workflow, properties)
-
-  request.info(_('Workflow re-submitted'))
-  return redirect(reverse('oozie:list_oozie_workflow', kwargs={'job_id': job_id}))
-
-
 @check_job_access_permission()
 def schedule_workflow(request, workflow):
   if Coordinator.objects.filter(workflow=workflow).exists():

+ 19 - 5
desktop/libs/liboozie/src/liboozie/oozie_api.py

@@ -104,6 +104,15 @@ class OozieApi(object):
       return { 'doAs': self.user }
     return { 'user.name': DEFAULT_USER, 'doAs': self.user }
 
+  def _get_oozie_properties(self, properties=None):
+    defaults = {
+      'user.name': self.user,
+    }
+
+    if properties is not None:
+      defaults.update(properties)
+
+    return defaults
 
   VALID_JOB_FILTERS = ('name', 'user', 'group', 'status')
 
@@ -210,7 +219,7 @@ class OozieApi(object):
     """
     submit_workflow(application_path, properties=None) -> jobid
 
-    Submit a job to Oozie. May raise PopupException.
+    Raise RestException on error.
     """
     defaults = {
       'oozie.wf.application.path': application_path,
@@ -223,11 +232,12 @@ class OozieApi(object):
 
     return self.submit_job(properties)
 
+  # Is name actually submit_coord?
   def submit_job(self, properties=None):
     """
     submit_job(properties=None, id=None) -> jobid
 
-    Submit a job to Oozie. May raise PopupException.
+    Raise RestException on error.
     """
     defaults = {
       'user.name': self.user,
@@ -239,11 +249,15 @@ class OozieApi(object):
     properties = defaults
 
     params = self._get_params()
-    resp = self._root.post('jobs', params,
-                  data=config_gen(properties),
-                  contenttype=_XML_CONTENT_TYPE)
+    resp = self._root.post('jobs', params, data=config_gen(properties), contenttype=_XML_CONTENT_TYPE)
     return resp['id']
 
+  def rerun(self, jobid, properties=None):
+    properties = self._get_oozie_properties(properties)
+    params = self._get_params()
+    params['action'] = 'rerun'
+
+    return self._root.put('job/%s' % jobid, params, data=config_gen(properties), contenttype=_XML_CONTENT_TYPE)
 
   def get_build_version(self):
     """

+ 46 - 8
desktop/libs/liboozie/src/liboozie/submittion.py

@@ -33,12 +33,18 @@ LOG = logging.getLogger(__name__)
 
 
 class Submission(object):
-  """Represents one unique Oozie submission"""
-  def __init__(self, user, job, fs, properties=None):
+  """
+  Represents one unique Oozie submission.
+
+  Actions are:
+  - submit
+  - rerun
+  """
+  def __init__(self, user, job=None, fs=None, properties=None, oozie_id=None):
     self.job = job
     self.user = user
     self.fs = fs
-    self.oozie_id = None
+    self.oozie_id = oozie_id
 
     if properties is not None:
       self.properties = properties
@@ -46,7 +52,10 @@ class Submission(object):
       self.properties = {}
 
   def __str__(self):
-    res = "Submission for job '%s' (id %s, owner %s)" % (self.job.name, self.job.id, self.user)
+    if self.oozie_id:
+      res = "Submission for job '%s'" % (self.oozie_id,)
+    else:
+      res = "Submission for job '%s' (id %s, owner %s)" % (self.job.name, self.job.id, self.user)
     if self.oozie_id:
       res += " -- " + self.oozie_id
     return res
@@ -77,6 +86,30 @@ class Submission(object):
 
     return self.oozie_id
 
+  def rerun(self, deployment_dir, fail_nodes=None, skip_nodes=None):
+    jobtracker = cluster.get_cluster_addr_for_job_submission()
+
+    try:
+      prev = get_oozie().setuser(self.user.username)
+      self._update_properties(jobtracker, deployment_dir)
+      self.properties.update({'oozie.wf.application.path': deployment_dir})
+
+      if fail_nodes:
+        self.properties.update({'oozie.wf.rerun.failnodes': fail_nodes})
+      elif not skip_nodes:
+        self.properties.update({'oozie.wf.rerun.failnodes': 'true'}) # Case empty 'skip_nodes' list
+      else:
+        self.properties.update({'oozie.wf.rerun.skip.nodes': skip_nodes})
+
+      get_oozie().rerun(self.oozie_id, properties=self.properties)
+
+      LOG.info("Rerun: %s" % (self,))
+    finally:
+      get_oozie().setuser(prev)
+
+    return self.oozie_id
+
+
   def deploy(self):
     try:
       deployment_dir = self._create_deployment_dir()
@@ -92,13 +125,18 @@ class Submission(object):
 
   def _update_properties(self, jobtracker_addr, deployment_dir):
     properties = {
-        'jobTracker': jobtracker_addr,
-        'nameNode': self.fs.fs_defaultfs,
+      'jobTracker': jobtracker_addr,
+      'nameNode': self.fs.fs_defaultfs,
+    }
+
+    if self.job:
+      properties.update({
         self.job.get_application_path_key(): self.fs.get_hdfs_path(deployment_dir),
         self.job.HUE_ID: self.job.id
-    }
+      })
 
     properties.update(self.properties)
+
     self.properties = properties
 
   def _create_deployment_dir(self):
@@ -108,7 +146,7 @@ class Submission(object):
     """
     if self.user != self.job.owner:
       path = Hdfs.join(REMOTE_DEPLOYMENT_DIR.get(), '_%s_-oozie-%s-%s' % (self.user.username, self.job.id, time.time()))
-      self.fs.copy_remote_dir(self.job.deployment_dir, path, owner=self.user)
+      self.fs.copy_remote_dir(self.job.deployment_dir, path, owner=self.user, dir_mode=0711)
     else:
       path = self.job.deployment_dir
       self._create_dir(path)

+ 4 - 1
desktop/libs/liboozie/src/liboozie/types.py

@@ -48,7 +48,7 @@ class Action(object):
   def _fixup(self): pass
 
   def is_finished(self):
-    return self.status in ('OK', 'SUCCEEDED')
+    return self.status in ('OK', 'SUCCEEDED', 'DONE')
 
   @classmethod
   def create(self, action_class, action_dict):
@@ -57,6 +57,9 @@ class Action(object):
     else:
       return action_class(action_dict)
 
+  def __str__(self):
+    return '%s - %s' % (self.type, self.name)
+
 
 class ControlFlowAction(Action):
   _ATTRS = [