فهرست منبع

HUE-1263 [oozie] Workflow with subworkflow action can't be saved

- Added test case
Abraham Elmahrek 12 سال پیش
والد
کامیت
0f06888
3فایلهای تغییر یافته به همراه75 افزوده شده و 8 حذف شده
  1. 1 1
      apps/oozie/src/oozie/models.py
  2. 63 0
      apps/oozie/src/oozie/tests.py
  3. 11 7
      apps/oozie/src/oozie/views/api.py

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

@@ -408,7 +408,7 @@ class Workflow(Job):
     return reverse('oozie:edit_workflow', kwargs={'workflow': self.id})
     return reverse('oozie:edit_workflow', kwargs={'workflow': self.id})
 
 
   def get_hierarchy(self):
   def get_hierarchy(self):
-    node = self.start
+    node = Start.objects.get(workflow=self) # Uncached version of start.
     return self.get_hierarchy_rec(node=node) + [[Kill.objects.get(workflow=node.workflow)],
     return self.get_hierarchy_rec(node=node) + [[Kill.objects.get(workflow=node.workflow)],
                                            [End.objects.get(workflow=node.workflow)]]
                                            [End.objects.get(workflow=node.workflow)]]
 
 

+ 63 - 0
apps/oozie/src/oozie/tests.py

@@ -248,6 +248,20 @@ class OozieMockBase(object):
     Link(parent=action3, child=self.wf.end, name="ok").save()
     Link(parent=action3, child=self.wf.end, name="ok").save()
 
 
 
 
+  def create_noop_workflow(self, name='noop-test'):
+    wf = Workflow.objects.new_workflow(self.user)
+    wf.name = name
+    wf.save()
+    wf.start.workflow = wf
+    wf.end.workflow = wf
+    wf.start.save()
+    wf.end.save()
+    Kill.objects.create(name='kill', workflow=wf, node_type=Kill.node_type)
+    Link.objects.create(parent=wf.start, child=wf.end, name='related')
+    Link.objects.create(parent=wf.start, child=wf.end, name="to")
+    return wf
+
+
 class OozieBase(OozieServerProvider):
 class OozieBase(OozieServerProvider):
   requires_hadoop = True
   requires_hadoop = True
 
 
@@ -398,6 +412,55 @@ class TestAPI(OozieMockBase):
     response = self.c.post(reverse('oozie:workflow_save', kwargs={'workflow': self.wf.pk}), data={'workflow': workflow_json}, HTTP_X_REQUESTED_WITH='XMLHttpRequest')
     response = self.c.post(reverse('oozie:workflow_save', kwargs={'workflow': self.wf.pk}), data={'workflow': workflow_json}, HTTP_X_REQUESTED_WITH='XMLHttpRequest')
     assert_equal(200, response.status_code)
     assert_equal(200, response.status_code)
 
 
+  def test_workflow_save_subworkflow(self):
+    subworkflow_name = "subworkflow-1"
+    subworkflow_id = "subworkflow:1"
+    subworkflow_json = """{
+      "description": "",
+      "workflow": %(workflow)d,
+      "child_links": [
+        {
+          "comment": "",
+          "name": "ok",
+          "parent": "%(id)s",
+          "child": %(end)d
+        },
+        {
+          "comment": "",
+          "name": "error",
+          "parent": "%(id)s",
+          "child": %(kill)d
+        }
+      ],
+      "node_type": "subworkflow",
+      "sub_workflow": %(subworkflow)d,
+      "job_properties": "[]",
+      "name": "%(name)s",
+      "id": "%(id)s",
+      "propagate_configuration": true
+    }"""
+
+    wf = self.create_noop_workflow()
+    subworkflow_json = subworkflow_json % {
+      'workflow': wf.id,
+      'subworkflow': self.wf.id,
+      'end': wf.end.id,
+      'kill': Kill.objects.get(workflow=wf).id,
+      'name': subworkflow_name,
+      'id': subworkflow_id
+    }
+    workflow_dict = workflow_to_dict(wf)
+    workflow_dict['nodes'].append(json.loads(subworkflow_json))
+    workflow_dict['nodes'][0]['child_links'][1]['child'] = subworkflow_id
+    del workflow_dict['nodes'][0]['child_links'][1]['id']
+    workflow_json = json.dumps(workflow_dict)
+
+    response = self.c.post(reverse('oozie:workflow_save', kwargs={'workflow': wf.pk}), data={'workflow': workflow_json})
+    test_response_json = response.content
+    test_response_json_object = json.loads(test_response_json)
+    assert_equal(0, test_response_json_object['status'], workflow_json)
+
+
   def test_workflow(self):
   def test_workflow(self):
     response = self.c.get(reverse('oozie:workflow', kwargs={'workflow': self.wf.pk}))
     response = self.c.get(reverse('oozie:workflow', kwargs={'workflow': self.wf.pk}))
     test_response_json = response.content
     test_response_json = response.content

+ 11 - 7
apps/oozie/src/oozie/views/api.py

@@ -40,7 +40,7 @@ from oozie.utils import model_to_dict, format_dict_field_values, format_field_va
 LOG = logging.getLogger(__name__)
 LOG = logging.getLogger(__name__)
 
 
 
 
-def get_or_create_node(workflow, node_data):
+def get_or_create_node(workflow, node_data, save=True):
   node = None
   node = None
   id = str(node_data['id'])
   id = str(node_data['id'])
   separator_index = id.find(':')
   separator_index = id.find(':')
@@ -56,7 +56,10 @@ def get_or_create_node(workflow, node_data):
     node = node_model(**kwargs)
     node = node_model(**kwargs)
   else:
   else:
     raise StructuredException(code="INVALID_REQUEST_ERROR", message=_('Could not find node of type'), data=node_data, error_code=500)
     raise StructuredException(code="INVALID_REQUEST_ERROR", message=_('Could not find node of type'), data=node_data, error_code=500)
-  node.save()
+
+  if save:
+    node.save()
+
   return node
   return node
 
 
 
 
@@ -191,16 +194,17 @@ def _update_workflow_nodes_json(workflow, json_nodes, id_map, user):
   nodes = []
   nodes = []
 
 
   for json_node in json_nodes:
   for json_node in json_nodes:
-    node = get_or_create_node(workflow, json_node)
-
-    if node.node_type == 'fork' and json_node['node_type'] == 'decision':
-      node = node.convert_to_decision()
+    node = get_or_create_node(workflow, json_node, save=False)
 
 
     if node.node_type == 'subworkflow':
     if node.node_type == 'subworkflow':
       try:
       try:
         node.sub_workflow = Workflow.objects.get(id=int(json_node['sub_workflow']))
         node.sub_workflow = Workflow.objects.get(id=int(json_node['sub_workflow']))
+        node.save()
       except Workflow.DoesNotExist:
       except Workflow.DoesNotExist:
-        pass
+        raise StructuredException(code="INVALID_REQUEST_ERROR", message=_('Error saving workflow'), data={'errors': 'Chosen subworkflow does not exist.'}, error_code=400)
+    elif node.node_type == 'fork' and json_node['node_type'] == 'decision':
+      node.save() # Need to save in case database throws error when performing delete.
+      node = node.convert_to_decision()
 
 
     id_map[str(json_node['id'])] = node.id
     id_map[str(json_node['id'])] = node.id