Просмотр исходного кода

[oozie] Create deployment directory when creating a new workflow

Fix coordinatoor submission popup asking for name of a dataset
Variable names need to be at least 2 char for being used in a coordinator
Romain Rigaux 13 лет назад
Родитель
Сommit
3aba683

+ 15 - 8
apps/oozie/src/oozie/models.py

@@ -50,7 +50,7 @@ LOG = logging.getLogger(__name__)
 
 PATH_MAX = 512
 name_validator = RegexValidator(regex='[a-zA-Z_][\-_a-zA-Z0-9]{1,39}',
-                                message=_('Please enter a valid value: 40 alphanum chars starting by an alpha'))
+                                message=_('Please enter a valid value: combination of 2 and 40 letters and digits starting by a letter'))
 
 
 """
@@ -153,7 +153,8 @@ class Job(models.Model):
       return _('personal')
 
   def find_all_parameters(self):
-    params = dict([(param, '') for param in self.find_parameters()])
+    params = self.find_parameters()
+
     for param in self.get_parameters():
       params[param['name'].strip()] = param['value']
 
@@ -199,7 +200,7 @@ class WorkflowManager(models.Manager):
 
     # Recheck if deployement dir exists
     oozie_setup.create_data_dir(fs)
-    Submission(workflow.owner, workflow, fs, {})._create_deployment_dir()
+    Submission(workflow.owner, workflow, fs, {})._create_dir(workflow.deployment_dir)
 
     return workflow
 
@@ -471,7 +472,7 @@ class Workflow(Job):
       if hasattr(node, 'find_parameters'):
         params.update(node.find_parameters())
 
-    return list(params)
+    return dict([(param, '') for param in list(params)])
 
   @property
   def actions(self):
@@ -1101,14 +1102,20 @@ class Coordinator(Job):
     return '%(number)d %(unit)s' % {'unit': self.frequency_unit, 'number': self.frequency_number}
 
   def find_parameters(self):
-    params = set()
+    params = self.workflow.find_parameters()
 
     for dataset in self.dataset_set.all():
-      params.update(set(find_parameters(dataset, ['uri'])))
+      for param in find_parameters(dataset, ['uri']):
+        if param not in set(['MINUTE', 'DAY', 'MONTH', 'YEAR']):
+          params[param] = ''
+
+    for ds in self.datainput_set.all():
+      params[ds.name] = '%s [dataset]' % ds.dataset
 
-    params.update(set(self.workflow.find_parameters()))
+    for ds in self.dataoutput_set.all():
+      params[ds.name] = '%s [dataset]' % ds.dataset
 
-    return list(params - set(['MINUTE', 'DAY', 'MONTH', 'YEAR']))
+    return params
 
 
 def utc_datetime_format(utc_time):

+ 28 - 4
apps/oozie/src/oozie/tests.py

@@ -880,15 +880,29 @@ class TestEditor:
     coord = create_coordinator(self.wf)
     create_dataset(coord)
 
-    response = self.c.post(reverse('oozie:create_coordinator_data', args=[coord.id, 'input']),
-                           {u'input-name': [u'input_dir'], u'input-dataset': [u'1']})
-    data = json.loads(response.content)
-    assert_equal(0, data['status'], data['data'])
+    create_coordinator_data(coord)
+
 
   def test_setup_app(self):
     self.c.post(reverse('oozie:setup_app'))
 
 
+  def test_get_workflow_parameters(self):
+    assert_equal([{'name': u'output', 'value': ''}, {'name': u'SLEEP', 'value': ''}, {'name': u'market', 'value': u'US'}],
+                 self.wf.find_all_parameters())
+
+
+  def test_get_coordinator_parameters(self):
+    coord = create_coordinator(self.wf)
+
+    create_dataset(coord)
+    create_coordinator_data(coord)
+
+    assert_equal([{'name': u'output', 'value': ''}, {'name': u'SLEEP', 'value': ''}, {'name': u'market', 'value': u'US,France'},
+                  {'name': u'input_dir', 'value': 'MyDataset [dataset]'}],
+                 coord.find_all_parameters())
+
+
 # Utils
 WORKFLOW_DICT = {u'deployment_dir': [u''], u'name': [u'wf-name-1'], u'description': [u''],
                  u'schema_version': [u'uri:oozie:workflow:0.2'],
@@ -921,6 +935,7 @@ def create_workflow():
 
   wf = Workflow.objects.get()
   assert_not_equal('', wf.deployment_dir)
+  # TODO test for existence on HDFS
 
   action1 = add_action(wf.id, wf.start.id, 'action-name-1')
   action2 = add_action(wf.id, action1.id, 'action-name-2')
@@ -969,6 +984,15 @@ def create_dataset(coord):
   assert_equal(0, data['status'], data['data'])
 
 
+def create_coordinator_data(coord):
+  c = make_logged_in_client()
+
+  response = c.post(reverse('oozie:create_coordinator_data', args=[coord.id, 'input']),
+                         {u'input-name': [u'input_dir'], u'input-dataset': [u'1']})
+  data = json.loads(response.content)
+  assert_equal(0, data['status'], data['data'])
+
+
 def move(c, wf, direction, action):
   try:
     LOG.info(wf.get_hierarchy())

+ 2 - 1
apps/oozie/src/oozie/views/editor.py

@@ -609,7 +609,8 @@ def edit_coordinator(request, coordinator):
     new_data_input_formset = NewDataInputFormSet(request.POST, request.FILES, instance=coordinator, prefix='input')
     new_data_output_formset = NewDataOutputFormSet(request.POST, request.FILES, instance=coordinator, prefix='output')
 
-    if coordinator_form.is_valid() and dataset_formset.is_valid() and data_input_formset.is_valid() and data_output_formset.is_valid() and new_data_input_formset.is_valid() and new_data_output_formset.is_valid():
+    if coordinator_form.is_valid() and dataset_formset.is_valid() and data_input_formset.is_valid() and data_output_formset.is_valid() \
+        and new_data_input_formset.is_valid() and new_data_output_formset.is_valid():
       coordinator = coordinator_form.save()
       dataset_formset.save()
       data_input_formset.save()

+ 10 - 4
desktop/libs/liboozie/src/liboozie/submittion.py

@@ -103,24 +103,30 @@ class Submission(object):
   def _create_deployment_dir(self):
     """
     Return the job deployment directory in HDFS, creating it if necessary.
+    The actual deployment dir should be 0711 owned by the user
     """
     path = Hdfs.join(REMOTE_DEPLOYMENT_DIR.get(), '_%s_-oozie-%s-%s' % (self.user.username, self.job.id, time.time()))
+    self._create_dir(path)
+    return path
 
+  def _create_dir(self, path, perms=0711):
+    """
+    Return the directory in HDFS, creating it if necessary.
+    """
     try:
       statbuf = self.fs.stats(path)
       if not statbuf.isDir:
-        msg = _("Workflow deployment path is not a directory: %s") % (path,)
+        msg = _("Path is not a directory: %s") % (path,)
         LOG.error(msg)
         raise Exception(msg)
       return path
     except IOError, ex:
       if ex.errno != errno.ENOENT:
-        msg = _("Error accessing workflow directory '%s': %s") % (path, ex)
+        msg = _("Error accessing directory '%s': %s") % (path, ex)
         LOG.exception(msg)
         raise IOError(ex.errno, msg)
-    # The actual deployment dir should be 0711 owned by the user
     if not self.fs.exists(path):
-      self._do_as(self.user.username , self.fs.mkdir, path, 0711)
+      self._do_as(self.user.username , self.fs.mkdir, path, perms)
     return path
 
   def _copy_files(self, deployment_dir, oozie_xml):