Browse Source

HUE-2551 [oozie] Easier update of a scheduled workflow

krish 10 years ago
parent
commit
136c060

+ 16 - 0
apps/oozie/src/oozie/templates/dashboard/list_oozie_coordinator.mako

@@ -149,6 +149,13 @@ ${ layout.menubar(section='coordinators', dashboard=True) }
                   </button>
                   </button>
                 </li>
                 </li>
               % endif
               % endif
+              <li class="white">
+                <button title="${ _('Sync Workflow') }" id="sync-wf-btn"
+                     data-sync-url="${ url('oozie:sync_coord_workflow', job_id=oozie_coordinator.id) }"
+                     class="btn btn-small sync-wf-btn" style="margin-bottom: 5px">
+                    ${ _('Sync Workflow') }
+                  </button>
+              </li>
             </ul>
             </ul>
           </div>
           </div>
         </div>
         </div>
@@ -865,6 +872,15 @@ ${ layout.menubar(section='coordinators', dashboard=True) }
       });
       });
     });
     });
 
 
+    $('#sync-wf-btn, .sync-wf-btn').click(function () {
+
+      $.get($(this).data("sync-url"), function (response) {
+        $('#rerun-coord-modal').html(response);
+        $('#rerun-coord-modal').modal('show');
+      });
+
+    });
+
     function refreshActionsPagination() {
     function refreshActionsPagination() {
       actionTableOffset = 1;
       actionTableOffset = 1;
     }
     }

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

@@ -26,7 +26,11 @@
   ${ csrf_token(request) | n,unicode }
   ${ csrf_token(request) | n,unicode }
   <div class="modal-header">
   <div class="modal-header">
     <a href="#" class="close" data-dismiss="modal">&times;</a>
     <a href="#" class="close" data-dismiss="modal">&times;</a>
-    <h3>${ _('Submit %(job)s?') % {'job': name} }</h3>
+    % if header:
+      <h3>${header}</h3>
+    % else:
+      <h3>${ _('Submit %(job)s?') % {'job': name} }</h3>
+    % endif
   </div>
   </div>
   <div class="modal-body">
   <div class="modal-body">
 
 

+ 87 - 1
apps/oozie/src/oozie/tests.py

@@ -37,6 +37,7 @@ from desktop.lib.django_test_util import make_logged_in_client
 from desktop.lib.test_utils import grant_access, add_permission, add_to_group, reformat_json, reformat_xml
 from desktop.lib.test_utils import grant_access, add_permission, add_to_group, reformat_json, reformat_xml
 from desktop.models import Document, Document2
 from desktop.models import Document, Document2
 
 
+from hadoop import cluster as originalCluster
 from hadoop.pseudo_hdfs4 import is_live_cluster
 from hadoop.pseudo_hdfs4 import is_live_cluster
 from jobsub.models import OozieDesign, OozieMapreduceAction
 from jobsub.models import OozieDesign, OozieMapreduceAction
 from liboozie import oozie_api
 from liboozie import oozie_api
@@ -65,7 +66,8 @@ class MockOozieApi:
                         {u'status': u'KILLED', u'run': 0, u'startTime': u'Mon, 30 Jul 2012 22:31:08 GMT', u'appName': u'WordCount2', u'lastModTime': u'Mon, 30 Jul 2012 22:32:20 GMT', u'actions': [], u'acl': None, u'appPath': None, u'externalId': '-', u'consoleUrl': u'http://runreal:11000/oozie?job=0000011-120725142744176-oozie-oozi-W', u'conf': None, u'parentId': None, u'createdTime': u'Mon, 30 Jul 2012 22:31:08 GMT', u'toString': u'Workflow id[0000011-120725142744176-oozie-oozi-W] status[SUCCEEDED]', u'endTime': u'Mon, 30 Jul 2012 22:32:20 GMT', u'id': u'0000011-120725142744176-oozie-oozi-W', u'group': None, u'user': u'test'},
                         {u'status': u'KILLED', u'run': 0, u'startTime': u'Mon, 30 Jul 2012 22:31:08 GMT', u'appName': u'WordCount2', u'lastModTime': u'Mon, 30 Jul 2012 22:32:20 GMT', u'actions': [], u'acl': None, u'appPath': None, u'externalId': '-', u'consoleUrl': u'http://runreal:11000/oozie?job=0000011-120725142744176-oozie-oozi-W', u'conf': None, u'parentId': None, u'createdTime': u'Mon, 30 Jul 2012 22:31:08 GMT', u'toString': u'Workflow id[0000011-120725142744176-oozie-oozi-W] status[SUCCEEDED]', u'endTime': u'Mon, 30 Jul 2012 22:32:20 GMT', u'id': u'0000011-120725142744176-oozie-oozi-W', u'group': None, u'user': u'test'},
                         {u'status': u'SUCCEEDED', u'run': 0, u'startTime': u'Mon, 30 Jul 2012 22:20:48 GMT', u'appName': u'WordCount3', u'lastModTime': u'Mon, 30 Jul 2012 22:22:00 GMT', u'actions': [], u'acl': None, u'appPath': None, u'externalId': '', u'consoleUrl': u'http://runreal:11000/oozie?job=0000009-120725142744176-oozie-oozi-W', u'conf': None, u'parentId': None, u'createdTime': u'Mon, 30 Jul 2012 22:20:48 GMT', u'toString': u'Workflow id[0000009-120725142744176-oozie-oozi-W] status[SUCCEEDED]', u'endTime': u'Mon, 30 Jul 2012 22:22:00 GMT', u'id': u'0000009-120725142744176-oozie-oozi-W', u'group': None, u'user': u'test'},
                         {u'status': u'SUCCEEDED', u'run': 0, u'startTime': u'Mon, 30 Jul 2012 22:20:48 GMT', u'appName': u'WordCount3', u'lastModTime': u'Mon, 30 Jul 2012 22:22:00 GMT', u'actions': [], u'acl': None, u'appPath': None, u'externalId': '', u'consoleUrl': u'http://runreal:11000/oozie?job=0000009-120725142744176-oozie-oozi-W', u'conf': None, u'parentId': None, u'createdTime': u'Mon, 30 Jul 2012 22:20:48 GMT', u'toString': u'Workflow id[0000009-120725142744176-oozie-oozi-W] status[SUCCEEDED]', u'endTime': u'Mon, 30 Jul 2012 22:22:00 GMT', u'id': u'0000009-120725142744176-oozie-oozi-W', u'group': None, u'user': u'test'},
                         {u'status': u'SUCCEEDED', u'run': 0, u'startTime': u'Mon, 30 Jul 2012 22:16:58 GMT', u'appName': u'WordCount4', u'lastModTime': u'Mon, 30 Jul 2012 22:18:10 GMT', u'actions': [], u'acl': None, u'appPath': None, u'externalId': None, u'consoleUrl': u'http://runreal:11000/oozie?job=0000008-120725142744176-oozie-oozi-W', u'conf': None, u'parentId': None, u'createdTime': u'Mon, 30 Jul 2012 22:16:58 GMT', u'toString': u'Workflow id[0000008-120725142744176-oozie-oozi-W] status[SUCCEEDED]', u'endTime': u'Mon, 30 Jul 2012 22:18:10 GMT', u'id': u'0000008-120725142744176-oozie-oozi-W', u'group': None, u'user': u'test'},
                         {u'status': u'SUCCEEDED', u'run': 0, u'startTime': u'Mon, 30 Jul 2012 22:16:58 GMT', u'appName': u'WordCount4', u'lastModTime': u'Mon, 30 Jul 2012 22:18:10 GMT', u'actions': [], u'acl': None, u'appPath': None, u'externalId': None, u'consoleUrl': u'http://runreal:11000/oozie?job=0000008-120725142744176-oozie-oozi-W', u'conf': None, u'parentId': None, u'createdTime': u'Mon, 30 Jul 2012 22:16:58 GMT', u'toString': u'Workflow id[0000008-120725142744176-oozie-oozi-W] status[SUCCEEDED]', u'endTime': u'Mon, 30 Jul 2012 22:18:10 GMT', u'id': u'0000008-120725142744176-oozie-oozi-W', u'group': None, u'user': u'test'},
-                        {u'status': u'SUCCEEDED', u'run': 0, u'startTime': u'Fri, 24 Jul 2015 14:56:08 AEST', u'appName': u'WordCount5', u'lastModTime': u'Fri, 24 Jul 2015 14:57:17 AEST', u'actions': [], u'acl': None, u'appPath': None, u'externalId': None, u'consoleUrl': u'http://runreal:11000/oozie?job=0000008-120725142744176-oozie-oozi-W', u'conf': None, u'parentId': None, u'createdTime': u'Fri, 24 Jul 2015 14:56:08 AEST', u'toString': u'Workflow id[0000007-120725142744176-oozie-oozi-W] status[SUCCEEDED]', u'endTime': u'Fri, 24 Jul 2015 14:57:17 AEST', u'id': u'0000007-120725142744176-oozie-oozi-W', u'group': None, u'user': u'test'}]
+                        {u'status': u'SUCCEEDED', u'run': 0, u'startTime': u'Fri, 24 Jul 2015 14:56:08 AEST', u'appName': u'WordCount5', u'lastModTime': u'Fri, 24 Jul 2015 14:57:17 AEST', u'actions': [], u'acl': None, u'appPath': None, u'externalId': None, u'consoleUrl': u'http://runreal:11000/oozie?job=0000008-120725142744176-oozie-oozi-W', u'conf': None, u'parentId': None, u'createdTime': u'Fri, 24 Jul 2015 14:56:08 AEST', u'toString': u'Workflow id[0000007-120725142744176-oozie-oozi-W] status[SUCCEEDED]', u'endTime': u'Fri, 24 Jul 2015 14:57:17 AEST', u'id': u'0000007-120725142744176-oozie-oozi-W', u'group': None, u'user': u'test'},
+                        {u'status': u'RUNNING', u'run': 0, u'startTime': u'Mon, 30 Jul 2012 22:35:48 GMT', u'appName': u'WordCount1', u'lastModTime': u'Mon, 30 Jul 2012 22:37:00 GMT', u'actions': [], u'acl': None, u'appPath': None, u'externalId': 'job_201208072118_0044', u'consoleUrl': u'http://runreal:11000/oozie?job=0000006-120725142744176-oozie-oozi-W', u'conf': None, u'parentId': None, u'createdTime': u'Mon, 30 Jul 2012 22:35:48 GMT', u'toString': u'Workflow id[0000006-120725142744176-oozie-oozi-W] status[SUCCEEDED]', u'endTime': u'Mon, 30 Jul 2012 22:37:00 GMT', u'id': u'0000006-120725142744176-oozie-oozi-W', u'group': None, u'user': u'test'}]
 
 
   WORKFLOW_IDS = [wf['id'] for wf in JSON_WORKFLOW_LIST]
   WORKFLOW_IDS = [wf['id'] for wf in JSON_WORKFLOW_LIST]
   WORKFLOW_DICT = dict([(wf['id'], wf) for wf in JSON_WORKFLOW_LIST])
   WORKFLOW_DICT = dict([(wf['id'], wf) for wf in JSON_WORKFLOW_LIST])
@@ -232,6 +234,27 @@ class MockOozieApi:
     return {'status': "RUNNING"}
     return {'status': "RUNNING"}
 
 
 
 
+class MockFs():
+  def __init__(self, logical_name=None):
+
+    self.fs_defaultfs = 'hdfs://curacao:8020'
+    self.logical_name = logical_name if logical_name else ''
+
+  def setuser(self, user):
+    pass
+
+  def join(self, path1, path2):
+    return path1 + path2
+
+  def do_as_user(self, username, fn, *args, **kwargs):
+    return ''
+
+  def read(self, length=1024*1024):
+    return 'data'
+
+  def exists(self, path):
+    return True
+
 class OozieMockBase(object):
 class OozieMockBase(object):
 
 
   def setUp(self):
   def setUp(self):
@@ -249,14 +272,20 @@ class OozieMockBase(object):
     add_to_group("test")
     add_to_group("test")
     self.user = User.objects.get(username='test')
     self.user = User.objects.get(username='test')
     self.wf = create_workflow(self.c, self.user)
     self.wf = create_workflow(self.c, self.user)
+    self.original_fs = originalCluster.FS_CACHE["default"]
+    originalCluster.FS_CACHE["default"] = MockFs()
 
 
 
 
   def tearDown(self):
   def tearDown(self):
     oozie_api.OozieApi = oozie_api.OriginalOozieApi
     oozie_api.OozieApi = oozie_api.OriginalOozieApi
+    if originalCluster.FS_CACHE is None:
+      originalCluster.FS_CACHE = {}
+    originalCluster.FS_CACHE["default"] = self.original_fs
     Workflow.objects.check_workspace = Workflow.objects.original_check_workspace
     Workflow.objects.check_workspace = Workflow.objects.original_check_workspace
     oozie_api._api_cache = None
     oozie_api._api_cache = None
 
 
     History.objects.all().delete()
     History.objects.all().delete()
+
     for coordinator in Coordinator.objects.all():
     for coordinator in Coordinator.objects.all():
       coordinator.delete(skip_trash=True)
       coordinator.delete(skip_trash=True)
     for bundle in Bundle.objects.all():
     for bundle in Bundle.objects.all():
@@ -3283,6 +3312,15 @@ class TestDashboard(OozieMockBase):
     response = self.c.get(reverse('oozie:rerun_oozie_coord', args=[MockOozieApi.WORKFLOW_IDS[0], '/path']))
     response = self.c.get(reverse('oozie:rerun_oozie_coord', args=[MockOozieApi.WORKFLOW_IDS[0], '/path']))
     assert_true('Rerun' in response.content, response.content)
     assert_true('Rerun' in response.content, response.content)
 
 
+  def test_sync_coord_workflow(self):
+    wf_doc = save_temp_workflow(MockOozieApi.JSON_WORKFLOW_LIST[5], self.user)
+    reset = ENABLE_V2.set_for_testing(True)
+    try:
+      response = self.c.get(reverse('oozie:sync_coord_workflow', args=[MockOozieApi.WORKFLOW_IDS[5]]))
+      assert_equal([{'name':'Dryrun', 'value': False}, {'name':'ls_arg', 'value': '-l'}], response.context['params_form'].initial)
+    finally:
+      wf_doc.delete()
+      reset()
 
 
   def test_rerun_coordinator_permissions(self):
   def test_rerun_coordinator_permissions(self):
     post_data = {
     post_data = {
@@ -3914,3 +3952,51 @@ def synchronize_workflow_attributes(workflow_json, correct_workflow_json):
     workflow_dict['attributes']['deployment_dir'] = correct_workflow_dict['attributes']['deployment_dir']
     workflow_dict['attributes']['deployment_dir'] = correct_workflow_dict['attributes']['deployment_dir']
 
 
   return reformat_json(workflow_dict)
   return reformat_json(workflow_dict)
+
+def save_temp_workflow(wf, user):
+    data = json.dumps({
+          'layout': [{
+              "size":12, "rows":[
+                  {"widgets":[{"size":12, "name":"Start", "id":"3f107997-04cc-8733-60a9-a4bb62cebffc", "widgetType":"start-widget", "properties":{}, "offset":0, "isLoading":False, "klass":"card card-widget span12"}]},
+                  {"widgets":[{"size":12, "name":"End", "id":"33430f0f-ebfa-c3ec-f237-3e77efa03d0a", "widgetType":"end-widget", "properties":{}, "offset":0, "isLoading":False, "klass":"card card-widget span12"}]},
+                  {"widgets":[{"size":12, "name":"Kill", "id":"17c9c895-5a16-7443-bb81-f34b30b21548", "widgetType":"kill-widget", "properties":{}, "offset":0, "isLoading":False, "klass":"card card-widget span12"}]}
+              ],
+              "drops":[ "temp"],
+              "klass":"card card-home card-column span12"
+          }],
+          'workflow': {
+              "id": None,
+              "uuid": None,
+              "name": "My Workflow",
+              "properties": {
+                  "deployment_dir": "",
+                  "description": "",
+                  "job_xml": "",
+                  "sla_enabled": False,
+                  "schema_version": "uri:oozie:workflow:0.5",
+                  "properties": [],
+                  "parameters": [{'name':'Dryrun', 'value': False}, {'name':'ls_arg', 'value': '-l'}],
+                  "sla":  [
+      {'key': 'enabled', 'value': False}, # Always first element
+      {'key': 'nominal-time', 'value': '${nominal_time}'},
+      {'key': 'should-start', 'value': ''},
+      {'key': 'should-end', 'value': '${30 * MINUTES}'},
+      {'key': 'max-duration', 'value': ''},
+      {'key': 'alert-events', 'value': ''},
+      {'key': 'alert-contact', 'value': ''},
+      {'key': 'notification-msg', 'value': ''},
+      {'key': 'upstream-apps', 'value': ''},
+  ],
+                  "show_arrows": True,
+                  "wf1_id": None
+              },
+              "nodes":[
+                  {"id":"3f107997-04cc-8733-60a9-a4bb62cebffc","name":"Start","type":"start-widget","properties":{},"children":[{'to': '33430f0f-ebfa-c3ec-f237-3e77efa03d0a'}]},
+                  {"id":"33430f0f-ebfa-c3ec-f237-3e77efa03d0a","name":"End","type":"end-widget","properties":{},"children":[]},
+                  {"id":"17c9c895-5a16-7443-bb81-f34b30b21548","name":"Kill","type":"kill-widget","properties":{'message': 'Action failed, error message[${wf:errorMessage(wf:lastErrorNode())}]'},"children":[]}
+              ]
+          }
+      })
+    workflow_doc = Document2.objects.create(name='test', type='oozie-workflow2', owner=user, data=data)
+    wf[u'conf'] = u'<configuration><property><name>hue-id-w</name><value>' + str(workflow_doc.id) + u'</value></property></configuration>'
+    return workflow_doc

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

@@ -133,6 +133,7 @@ urlpatterns += patterns(
   url(r'^rerun_oozie_job/(?P<job_id>[-\w]+)/(?P<app_path>.+?)$', 'rerun_oozie_job', name='rerun_oozie_job'),
   url(r'^rerun_oozie_job/(?P<job_id>[-\w]+)/(?P<app_path>.+?)$', 'rerun_oozie_job', name='rerun_oozie_job'),
   url(r'^rerun_oozie_coord/(?P<job_id>[-\w]+)/(?P<app_path>.+?)$', 'rerun_oozie_coordinator', name='rerun_oozie_coord'),
   url(r'^rerun_oozie_coord/(?P<job_id>[-\w]+)/(?P<app_path>.+?)$', 'rerun_oozie_coordinator', name='rerun_oozie_coord'),
   url(r'^rerun_oozie_bundle/(?P<job_id>[-\w]+)/(?P<app_path>.+?)$', 'rerun_oozie_bundle', name='rerun_oozie_bundle'),
   url(r'^rerun_oozie_bundle/(?P<job_id>[-\w]+)/(?P<app_path>.+?)$', 'rerun_oozie_bundle', name='rerun_oozie_bundle'),
+  url(r'^sync_coord_workflow/(?P<job_id>[-\w]+)$', 'sync_coord_workflow', name='sync_coord_workflow'),
   url(r'^manage_oozie_jobs/(?P<job_id>[-\w]+)/(?P<action>(start|suspend|resume|kill|rerun|change|ignore))$', 'manage_oozie_jobs', name='manage_oozie_jobs'),
   url(r'^manage_oozie_jobs/(?P<job_id>[-\w]+)/(?P<action>(start|suspend|resume|kill|rerun|change|ignore))$', 'manage_oozie_jobs', name='manage_oozie_jobs'),
   url(r'^bulk_manage_oozie_jobs/$', 'bulk_manage_oozie_jobs', name='bulk_manage_oozie_jobs'),
   url(r'^bulk_manage_oozie_jobs/$', 'bulk_manage_oozie_jobs', name='bulk_manage_oozie_jobs'),
 
 

+ 44 - 1
apps/oozie/src/oozie/views/dashboard.py

@@ -41,7 +41,7 @@ from desktop.models import Document, Document2
 
 
 from liboozie.oozie_api import get_oozie
 from liboozie.oozie_api import get_oozie
 from liboozie.credentials import Credentials
 from liboozie.credentials import Credentials
-from liboozie.submittion import Submission
+from liboozie.submission2 import Submission
 from liboozie.types import Workflow as OozieWorkflow, Coordinator as CoordinatorWorkflow, Bundle as BundleWorkflow
 from liboozie.types import Workflow as OozieWorkflow, Coordinator as CoordinatorWorkflow, Bundle as BundleWorkflow
 
 
 from oozie.conf import OOZIE_JOBS_COUNT, ENABLE_CRON_SCHEDULING, ENABLE_V2
 from oozie.conf import OOZIE_JOBS_COUNT, ENABLE_CRON_SCHEDULING, ENABLE_V2
@@ -638,6 +638,49 @@ def massaged_sla_for_json(sla, request):
   return massaged_sla
   return massaged_sla
 
 
 
 
+@show_oozie_error
+def sync_coord_workflow(request, job_id):
+  ParametersFormSet = formset_factory(ParameterForm, extra=0)
+  job = check_job_access_permission(request, job_id)
+  check_job_edition_permission(job, request.user)
+
+  hue_coord = get_history().get_coordinator_from_config(job.conf_dict)
+  hue_wf = (hue_coord and hue_coord.workflow) or get_history().get_workflow_from_config(job.conf_dict)
+
+  if request.method == 'POST':
+    params_form = ParametersFormSet(request.POST)
+    if params_form.is_valid():
+      mapping = dict([(param['name'], param['value']) for param in params_form.cleaned_data])
+
+      submission = Submission(user=request.user, job=hue_wf, fs=request.fs, jt=request.jt, properties=mapping)
+      submission._sync_definition(hue_wf.deployment_dir, mapping)
+
+      request.info(_('Successfully updated Workflow definition'))
+      return redirect(reverse('oozie:list_oozie_coordinator', kwargs={'job_id': job_id}))
+    else:
+      request.error(_('Invalid submission form: %s' % params_form.errors))
+  else:
+    parameters = hue_wf and hue_wf.find_all_parameters() or []
+    params_dict = dict([(param['name'], param['value']) for param in parameters])
+
+    submission = Submission(user=request.user, job=hue_wf, fs=request.fs, jt=request.jt, properties=None)
+    prev_properties = hue_wf and hue_wf.deployment_dir and \
+                      submission.get_external_parameters(request.fs.join(hue_wf.deployment_dir, hue_wf.XML_FILE_NAME)) or {}
+
+    for key, value in params_dict.iteritems():
+      params_dict[key] = prev_properties[key] if key in prev_properties.keys() else params_dict[key]
+
+    initial_params = ParameterForm.get_initial_params(params_dict)
+    params_form = ParametersFormSet(initial=initial_params)
+
+  popup = render('editor2/submit_job_popup.mako', request, {
+             'params_form': params_form,
+             'name': _('Job'),
+             'header': _('Sync Workflow definition?'),
+             'action': reverse('oozie:sync_coord_workflow', kwargs={'job_id': job_id})
+           }, force_template=True).content
+  return JsonResponse(popup, safe=False)
+
 @show_oozie_error
 @show_oozie_error
 def rerun_oozie_job(request, job_id, app_path):
 def rerun_oozie_job(request, job_id, app_path):
   ParametersFormSet = formset_factory(ParameterForm, extra=0)
   ParametersFormSet = formset_factory(ParameterForm, extra=0)

+ 17 - 6
desktop/libs/liboozie/src/liboozie/submission2.py

@@ -277,13 +277,9 @@ class Submission(object):
     Copy XML and the jar_path files from Java or MR actions to the deployment directory.
     Copy XML and the jar_path files from Java or MR actions to the deployment directory.
     This should run as the workflow user.
     This should run as the workflow user.
     """
     """
-    xml_path = self.fs.join(deployment_dir, self.job.XML_FILE_NAME)
-    self.fs.create(xml_path, overwrite=True, permission=0644, data=smart_str(oozie_xml))
-    LOG.debug("Created %s" % (xml_path,))
 
 
-    properties_path = self.fs.join(deployment_dir, 'job.properties')
-    self.fs.create(properties_path, overwrite=True, permission=0644, data=smart_str('\n'.join(['%s=%s' % (key, val) for key, val in oozie_properties.iteritems()])))
-    LOG.debug("Created %s" % (properties_path,))
+    self._create_file(deployment_dir, self.job.XML_FILE_NAME, oozie_xml)
+    self._create_file(deployment_dir, 'job.properties', data='\n'.join(['%s=%s' % (key, val) for key, val in oozie_properties.iteritems()]))
 
 
     # List jar files
     # List jar files
     files = []
     files = []
@@ -337,7 +333,22 @@ class Submission(object):
     from oozie.models2 import Coordinator
     from oozie.models2 import Coordinator
     return Coordinator.PROPERTY_APP_PATH in self.properties
     return Coordinator.PROPERTY_APP_PATH in self.properties
 
 
+  def _create_file(self, deployment_dir, file_name, data, do_as=False):
+   file_path = self.fs.join(deployment_dir, file_name)
+   if do_as:
+     self.fs.do_as_user(self.user, self.fs.create, file_path, overwrite=True, permission=0644, data=smart_str(data))
+   else:
+     self.fs.create(file_path, overwrite=True, permission=0644, data=smart_str(data))
+   LOG.debug("Created/Updated %s" % (file_path,))
 
 
+  def _sync_definition(self, deployment_dir, mapping):
+    """ This is helper function for 'Sync Workflow' functionality in a Coordinator.
+      It copies updated workflow changes into HDFS """
+
+    self._create_file(deployment_dir, self.job.XML_FILE_NAME, self.job.to_xml(mapping=mapping), do_as=True)
+
+    data_properties = smart_str('\n'.join(['%s=%s' % (key, val) for key, val in mapping.iteritems()]))
+    self._create_file(deployment_dir, 'job.properties', data_properties, do_as=True)
 
 
 def create_directories(fs, directory_list=[]):
 def create_directories(fs, directory_list=[]):
   # If needed, create the remote home, deployment and data directories
   # If needed, create the remote home, deployment and data directories

+ 4 - 0
desktop/libs/liboozie/src/liboozie/submittion2_tests.py

@@ -135,6 +135,10 @@ def test_copy_files():
     assert_not_equal(stats_udf5['fileId'], cluster.fs.stats(deployment_dir + '/udf5.jar')['fileId'])
     assert_not_equal(stats_udf5['fileId'], cluster.fs.stats(deployment_dir + '/udf5.jar')['fileId'])
     assert_equal(stats_udf6['fileId'], cluster.fs.stats(deployment_dir + '/udf6.jar')['fileId'])
     assert_equal(stats_udf6['fileId'], cluster.fs.stats(deployment_dir + '/udf6.jar')['fileId'])
 
 
+    # Test _create_file()
+    submission._create_file(deployment_dir, 'test.txt', data='Test data')
+    assert_true(cluster.fs.exists(deployment_dir + '/test.txt'), list_dir_workspace)
+
   finally:
   finally:
     try:
     try:
       cluster.fs.rmtree(prefix)
       cluster.fs.rmtree(prefix)