Explorar o código

HUE-2659 [oozie] Workflow submitted from coordinator gets parent graph

krish %!s(int64=10) %!d(string=hai) anos
pai
achega
d62a7165ca

+ 30 - 0
apps/oozie/src/oozie/importlib/workflows.py

@@ -687,6 +687,36 @@ def import_workflow(workflow, workflow_definition, metadata=None, fs=None):
   return import_workflow_root(workflow, workflow_definition_root, metadata, fs)
 
 
+def generate_v2_graph_nodes(workflow_definition):
+  # Parse Workflow Definition
+  workflow_definition_root = etree.fromstring(workflow_definition)
+  if workflow_definition_root is None:
+    raise MalformedWfDefException()
+
+  xslt_definition_fh = open("%(xslt_dir)s/workflow.xslt" % {
+      'xslt_dir': os.path.join(DEFINITION_XSLT2_DIR.get(), 'workflows')
+    })
+
+  tag = etree.QName(workflow_definition_root.tag)
+  schema_version = tag.namespace
+
+  # Ensure namespace exists
+  if schema_version not in OOZIE_NAMESPACES:
+    raise InvalidTagWithNamespaceException(workflow_definition_root.tag)
+
+  # Get XSLT
+  xslt = etree.parse(xslt_definition_fh)
+  xslt_definition_fh.close()
+  transform = etree.XSLT(xslt)
+
+  # Transform XML using XSLT
+  transformed_root = transform(workflow_definition_root)
+  node_list = str(transformed_root).replace('\n', '').replace(' ', '')
+  node_list = json.loads(node_list)
+
+  return node_list
+
+
 class MalformedWfDefException(Exception):
   pass
 

+ 1 - 1
apps/oozie/src/oozie/importlib/xslt2/workflows/0.5/action.xslt

@@ -10,12 +10,12 @@
 <xsl:import href="extensions/sqoop.0.1.xslt"/>
 <xsl:import href="extensions/sqoop.0.2.xslt"/>
 <xsl:import href="extensions/ssh.0.1.xslt"/>
+<xsl:import href="extensions/spark.0.1.xslt"/>
 <xsl:import href="nodes/fs.xslt"/>
 <xsl:import href="nodes/java.xslt"/>
 <xsl:import href="nodes/mapreduce.xslt"/>
 <xsl:import href="nodes/pig.xslt"/>
 <xsl:import href="nodes/streaming.xslt"/>
-<xsl:import href="nodes/generic.xslt"/>
 <xsl:import href="nodes/subworkflow.xslt"/>
 
 <xsl:template match="workflow:action" xmlns:workflow="uri:oozie:workflow:0.5">

+ 8 - 3
apps/oozie/src/oozie/importlib/xslt2/workflows/0.5/extensions/distcp.0.1.xslt

@@ -2,11 +2,16 @@
 
 <xsl:stylesheet version="1.0" xmlns:xsl="http://www.w3.org/1999/XSL/Transform" xmlns:workflow="uri:oozie:workflow:0.5" xmlns:distcp="uri:oozie:distcp-action:0.1" exclude-result-prefixes="workflow distcp">
 
-<xsl:import href="../nodes/fields/job_xml.xslt"/>
-
 <xsl:template match="distcp:distcp">
 
-  ,"distcp": <xsl:call-template name="job_xml"/>
+  ,"distcp": {
+        <xsl:for-each select="arg">
+          "path<xsl:value-of select='position()'/>": "<xsl:value-of select="arg"/>"
+          <xsl:if  test="position() &lt; last()">
+            ,
+          </xsl:if>
+        </xsl:for-each>
+    }
 
 </xsl:template>
 

+ 5 - 1
apps/oozie/src/oozie/importlib/xslt2/workflows/0.5/extensions/email.0.1.xslt

@@ -4,7 +4,11 @@
 
 <xsl:template match="email:email">
 
-  ,"email": {"subject": "<xsl:value-of select="*[local-name()='subject']"/>"}
+  ,"email": {
+    "to": "<xsl:value-of select="*[local-name()='to']"/>",
+    "subject": "<xsl:value-of select="*[local-name()='subject']"/>",
+    "body": "<xsl:value-of select="*[local-name()='body']"/>"
+   }
 
 </xsl:template>
 

+ 3 - 5
apps/oozie/src/oozie/importlib/xslt2/workflows/0.5/extensions/hive.0.1.xslt

@@ -1,12 +1,10 @@
 <?xml version="1.0"?>
 
-<xsl:stylesheet version="1.0" xmlns:xsl="http://www.w3.org/1999/XSL/Transform" xmlns:workflow="uri:oozie:workflow:0.5" xmlns:hive="uri:oozie:hive-action:0.1" exclude-result-prefixes="workflow hive">
+<xsl:stylesheet version="1.0" xmlns:xsl="http://www.w3.org/1999/XSL/Transform" xmlns:workflow="uri:oozie:workflow:0.5" xmlns:hive2="uri:oozie:hive2-action:0.1" exclude-result-prefixes="workflow hive">
 
-<xsl:import href="../nodes/fields/script_path.xslt"/>
+<xsl:template match="hive2:hive2">
 
-<xsl:template match="hive:hive">
-
-  ,"hive": {<xsl:call-template name="script_path"/>}
+  ,"hive2": {"script": "<xsl:value-of select="*[local-name()='script']"/>"}
 
 </xsl:template>
 

+ 1 - 3
apps/oozie/src/oozie/importlib/xslt2/workflows/0.5/extensions/hive.0.2.xslt

@@ -2,11 +2,9 @@
 
 <xsl:stylesheet version="1.0" xmlns:xsl="http://www.w3.org/1999/XSL/Transform" xmlns:workflow="uri:oozie:workflow:0.5" xmlns:hive="uri:oozie:hive-action:0.2" exclude-result-prefixes="workflow hive">
 
-<xsl:import href="../nodes/fields/script_path.xslt"/>
-
 <xsl:template match="hive:hive">
 
-  ,"hive": {<xsl:call-template name="script_path"/>}
+  ,"hive": {"script": "<xsl:value-of select="*[local-name()='script']"/>"}
 
 </xsl:template>
 

+ 17 - 0
apps/oozie/src/oozie/importlib/xslt2/workflows/0.5/extensions/spark.0.1.xslt

@@ -0,0 +1,17 @@
+<?xml version="1.0"?>
+
+<xsl:stylesheet version="1.0" xmlns:xsl="http://www.w3.org/1999/XSL/Transform" xmlns:workflow="uri:oozie:workflow:0.5" xmlns:spark="uri:oozie:spark-action:0.1" exclude-result-prefixes="workflow spark">
+
+<xsl:template match="spark:spark">
+
+  ,"spark": {
+    "name": "<xsl:value-of select="*[local-name()='name']"/>",
+    "master": "<xsl:value-of select="*[local-name()='master']"/>",
+    "mode": "<xsl:value-of select="*[local-name()='mode']"/>",
+    "class": "<xsl:value-of select="*[local-name()='class']"/>",
+    "jar": "<xsl:value-of select="*[local-name()='jar']"/>"
+  }
+
+</xsl:template>
+
+</xsl:stylesheet>

+ 5 - 1
apps/oozie/src/oozie/importlib/xslt2/workflows/0.5/extensions/ssh.0.1.xslt

@@ -6,7 +6,11 @@
 
 <xsl:template match="ssh:ssh">
 
-  ,"ssh": "<xsl:call-template name="command"/>"
+  ,"ssh": {
+    <xsl:call-template name="command"/>,
+    "user":"<xsl:value-of select="*[local-name()='user']"/>",
+    "host":"<xsl:value-of select="*[local-name()='host']"/>"
+    }
 
 </xsl:template>
 

+ 4 - 6
apps/oozie/src/oozie/importlib/xslt2/workflows/0.5/nodes/fields/touchzs.xslt

@@ -4,17 +4,15 @@
 
 <xsl:template name="touchzs">
 
-  "touchzs":
+  "touchzs": {
 
-    <xsl:text>[</xsl:text>
     <xsl:for-each select="*[local-name()='touchz']">
-      <xsl:text><![CDATA[{"name":"]]></xsl:text><xsl:value-of select="@path" /><xsl:text><![CDATA["}]]></xsl:text>
+      "path<xsl:value-of select='position()'/>": "<xsl:value-of select="@path"/>"
       <xsl:if  test="position() &lt; last()">
-        "<xsl:text>,</xsl:text>"
+        ,
       </xsl:if>
     </xsl:for-each>
-    <xsl:text>]</xsl:text>
-
+  }
 </xsl:template>
 
 </xsl:stylesheet>

+ 2 - 1
apps/oozie/src/oozie/importlib/xslt2/workflows/0.5/nodes/streaming.xslt

@@ -9,7 +9,8 @@
 
   ,"streaming": {
         <xsl:call-template name="mapper"/>,
-        <xsl:call-template name="reducer"/> }
+        <xsl:call-template name="reducer"/>
+    }
 
 </xsl:template>
 

+ 1 - 3
apps/oozie/src/oozie/importlib/xslt2/workflows/0.5/nodes/subworkflow.xslt

@@ -2,10 +2,8 @@
 
 <xsl:stylesheet version="1.0" xmlns:xsl="http://www.w3.org/1999/XSL/Transform" xmlns:workflow="uri:oozie:workflow:0.5" exclude-result-prefixes="workflow">
 
-<xsl:import href="fields/job_properties.xslt"/>
-
 <xsl:template match="workflow:sub-workflow" xmlns:workflow="uri:oozie:workflow:0.5">
-  , <xsl:call-template name="job_properties"/>
+  , "sub-workflow": { "app-path":"<xsl:value-of select="*[local-name()='app-path']"/>" }
 </xsl:template>
 
 </xsl:stylesheet>

+ 261 - 1
apps/oozie/src/oozie/models2.py

@@ -36,12 +36,14 @@ from desktop.lib.json_utils import JSONEncoderForHTML
 from desktop.models import Document2
 
 from hadoop.fs.hadoopfs import Hdfs
+from hadoop.fs.exceptions import WebHdfsException
+
 from liboozie.submission2 import Submission
 from liboozie.submission2 import create_directories
 
 from oozie.conf import REMOTE_SAMPLE_DIR
 from oozie.utils import utc_datetime_format, UTC_TIME_FORMAT, convert_to_server_timezone
-from hadoop.fs.exceptions import WebHdfsException
+from oozie.importlib.workflows import generate_v2_graph_nodes, MalformedWfDefException, InvalidTagWithNamespaceException
 
 
 LOG = logging.getLogger(__name__)
@@ -346,6 +348,264 @@ class Workflow(Job):
   def get_application_path_key(cls):
     return 'oozie.wf.application.path'
 
+  @classmethod
+  def gen_workflow_data_from_xml(cls, user, oozie_workflow):
+    node_list = []
+    try:
+      node_list = generate_v2_graph_nodes(oozie_workflow.definition)
+    except MalformedWfDefException, e:
+      LOG.exception("Could not find any nodes in Workflow definition. Maybe it's malformed?")
+    except InvalidTagWithNamespaceException, e:
+      LOG.exception("Tag with namespace %(namespace)s is not valid. Please use one of the following namespaces: %(namespaces)s" % {
+      'namespace': e.namespace,
+      'namespaces': e.namespaces
+    })
+
+    adj_list = _create_graph_adjaceny_list(node_list)
+
+    node_hierarchy = ['start']
+    _get_hierarchy_from_adj_list(adj_list, adj_list['start']['ok_to'], node_hierarchy)
+
+    _add_uuids_to_adj_list(adj_list)
+
+    wf_rows = _create_workflow_layout(node_hierarchy, adj_list)
+    data = {'layout': [{}], 'workflow': {}}
+    if wf_rows:
+      data['layout'][0]['rows'] = wf_rows
+
+    wf_nodes = []
+    _dig_nodes(node_hierarchy, adj_list, user, wf_nodes)
+    data['workflow']['nodes'] = wf_nodes
+    data['workflow']['id'] = "123"
+    data['workflow']['properties'] = json.loads("""{
+      "job_xml": "",
+      "description": "",
+      "wf1_id": null,
+      "sla_enabled": false,
+      "deployment_dir": "/user/hue/oozie/workspaces/hue-oozie-1452553957.19",
+      "schema_version": "uri:oozie:workflow:0.5",
+      "sla": [
+        {
+          "key": "enabled",
+          "value": false
+        },
+        {
+          "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,
+      "parameters": [
+        {
+          "name": "oozie.use.system.libpath",
+          "value": true
+        }
+      ],
+      "properties": []
+    }""")
+
+    return data
+
+def _add_uuids_to_adj_list(adj_list):
+  uuids = {}
+  id = 1
+  for node in adj_list.keys():
+    adj_list[node]['id'] = id
+
+    if adj_list[node]['node_type'] == 'kill':
+      adj_list[node]['uuid'] = '17c9c895-5a16-7443-bb81-f34b30b21548'
+    elif adj_list[node]['node_type'] == 'start':
+      adj_list[node]['uuid'] = '3f107997-04cc-8733-60a9-a4bb62cebffc'
+    elif adj_list[node]['node_type'] == 'end':
+      adj_list[node]['uuid'] = '33430f0f-ebfa-c3ec-f237-3e77efa03d0a'
+    else:
+      adj_list[node]['uuid'] = node[-4:] + str(uuid.uuid4())[4:]
+
+    uuids[id] = adj_list[node]['uuid']
+    id += 1
+  return adj_list
+
+def _dig_nodes(nodes, adj_list, user, wf_nodes):
+  for node in nodes:
+    if type(node) != list:
+      node = adj_list[node]
+      properties = {}
+      if '%s-widget' % node['node_type'] in NODES:
+        properties = dict(NODES['%s-widget' % node['node_type']].get_fields())
+
+      if node['node_type'] == 'pig':
+        properties['script_path'] = node.get('script_path')
+      elif node['node_type'] == 'spark':
+        properties['class'] = node.get('spark').get('class')
+        #TBD: jar_path
+        properties['jar_path'] = node.get('spark').get('jar')
+      elif node['node_type'] == 'hive' or node['node_type'] == 'hive2':
+        properties['script_path'] = node.get('hive2').get('script')
+      elif node['node_type'] == 'java':
+        properties['jar_path'] = node.get('jar_path')
+      elif node['node_type'] == 'sqoop':
+        properties['command'] = node.get('script_path')
+      elif node['node_type'] == 'mapreduce':
+        properties['jar_path'] = node.get('jar_path')
+      elif node['node_type'] == 'shell':
+        properties['shell_command'] = node.get('shell').get('command')
+      elif node['node_type'] == 'ssh':
+        properties['user'] = '%s@%s' % (node.get('ssh').get('user'), node.get('ssh').get('host'))
+        properties['ssh_command'] = node.get('ssh').get('command')
+      elif node['node_type'] == 'fs':
+        #TBD: all
+        properties['deletes'] = [{'value': f['name']} for f in json.loads(node.get('deletes'))]
+        properties['mkdirs'] = [{'value': f['name']} for f in json.loads(node.get('mkdirs'))]
+        properties['moves'] = json.loads(node.get('moves'))
+        properties['touchzs'] = [{'value': f['name']} for f in json.loads(node.get('touchzs'))]
+      elif node['node_type'] == 'email':
+        properties['to'] = node.get('email').get('to')
+        properties['subject'] = node.get('email').get('subject')
+        #TBD: body doesn't show up
+        properties['body'] = node.get('email').get('body')
+      elif node['node_type'] == 'streaming':
+        properties['mapper'] = node.get('streaming').get('mapper')
+        properties['reducer'] = node.get('streaming').get('reducer')
+      elif node['node_type'] == 'distcp':
+        #TBD: both
+        properties['source'] = node.get('distcp').get('source')
+        properties['destination'] = node.get('distcp').get('source')
+      elif node['node_type'] == 'sub-workflow':
+        properties['app-path'] = node.get('sub-workflow').get('app-path')
+        properties['workflow'] = node.get('uuid')
+        properties['job_properties'] = []
+        properties['sla'] = ''
+
+      children = []
+      if node['node_type'] == 'fork':
+        for key in node.keys():
+          if key.startswith('path'):
+            children.append({'to': adj_list[node[key]]['uuid'], 'condition': '${ 1 gt 0 }'})
+      else:
+        if node.get('ok_to'):
+          children.append({'to': adj_list[node['ok_to']]['uuid']})
+        if node.get('error_to'):
+          children.append({'error': adj_list[node['error_to']]['uuid']})
+
+      wf_nodes.append({
+          "id": node['uuid'],
+          "name": '%s-%s' % (node['node_type'].split('-')[0], node['uuid'][:4]),
+          "type": "%s-widget" % node['node_type'],
+          "properties": properties,
+          "children": children
+      })
+    else:
+      _dig_nodes(node, adj_list, user, wf_nodes)
+
+def _create_workflow_layout(nodes, adj_list, size=12):
+  wf_rows = []
+  for node in nodes:
+    if type(node) == list and len(node) == 1:
+      node = node[0]
+    if type(node) != list:
+      wf_rows.append({"widgets":[{"size":size, "name": adj_list[node]['node_type'], "id":  adj_list[node]['uuid'], "widgetType": "%s-widget" % adj_list[node]['node_type'] if adj_list[node]['node_type'] != 'sub-workflow' else 'subworkflow-widget', "properties":{}, "offset":0, "isLoading":False, "klass":"card card-widget span%s" % size, "columns":[]}]})
+    else:
+      if adj_list[node[0]]['node_type'] == 'fork':
+        wf_rows.append({"widgets":[{"size":size, "name": 'Fork', "id":  adj_list[node[0]]['uuid'], "widgetType": "%s-widget" % adj_list[node[0]]['node_type'], "properties":{}, "offset":0, "isLoading":False, "klass":"card card-widget span%s" % size, "columns":[]}]})
+
+        wf_rows.append({
+          "id": str(uuid.uuid4()),
+          "widgets":[
+
+          ],
+          "columns":[
+             {
+                "id": str(uuid.uuid4()),
+                "size": (size / len(node[1])),
+                "rows":
+                   [{
+                      "id": str(uuid.uuid4()),
+                      "widgets": c['widgets'],
+                      "columns":c.get('columns') or []
+                    } for c in col],
+                "klass":"card card-home card-column span%s" % (size / len(node[1]))
+             }
+             for col in [_create_workflow_layout(item, adj_list, size) for item in node[1]]
+          ]
+        })
+
+        wf_rows.append({"widgets":[{"size":size, "name": 'Join', "id":  adj_list[node[2]]['uuid'], "widgetType": "%s-widget" % adj_list[node[2]]['node_type'], "properties":{}, "offset":0, "isLoading":False, "klass":"card card-widget span%s" % size, "columns":[]}]})
+      else:
+        wf_rows.append(_create_workflow_layout(node, adj_list, size))
+  return wf_rows
+
+
+def _get_hierarchy_from_adj_list(adj_list, curr_node, node_hierarchy):
+
+  if adj_list[curr_node]['node_type'] == 'join':
+    return curr_node
+
+  elif adj_list[curr_node]['node_type'] == 'end':
+    node_hierarchy.append(['Kill'])
+    node_hierarchy.append(['End'])
+    return node_hierarchy
+
+  elif adj_list[curr_node]['node_type'] == 'fork':
+    fork_nodes = []
+    fork_nodes.append(curr_node)
+
+    join_node = None
+    children = []
+    for key in adj_list[curr_node].keys():
+      if key.startswith('path'):
+        child = []
+        join_node = _get_hierarchy_from_adj_list(adj_list, adj_list[curr_node][key], child)
+        children.append(child)
+
+    fork_nodes.append(children)
+    fork_nodes.append(join_node)
+
+    node_hierarchy.append(fork_nodes)
+    return _get_hierarchy_from_adj_list(adj_list, adj_list[join_node]['ok_to'], node_hierarchy)
+
+  else:
+    node_hierarchy.append(curr_node)
+    return _get_hierarchy_from_adj_list(adj_list, adj_list[curr_node]['ok_to'], node_hierarchy)
+
+
+def _create_graph_adjaceny_list(nodes):
+  start_node = [node for node in nodes if node.get('node_type') == 'start'][0]
+  adj_list = {'start': start_node}
+
+  for node in nodes:
+    if node and node.get('node_type') != 'start':
+      adj_list[node['name']] = node
+
+  return adj_list
+
 
 class Node():
   def __init__(self, data):

+ 36 - 0
apps/oozie/src/oozie/test_data/xslt2/test-workflow.xml

@@ -0,0 +1,36 @@
+<workflow-app name="Sub" xmlns="uri:oozie:workflow:0.5">
+    <start to="fork-68d4"/>
+    <kill name="Kill">
+        <message>Action failed, error message[${wf:errorMessage(wf:lastErrorNode())}]</message>
+    </kill>
+    <action name="subworkflow-a13f">
+        <sub-workflow>
+            <app-path>${nameNode}/user/hue/oozie/deployments/_admin_-oozie-50001-1427488969.48</app-path>
+              <propagate-configuration/>
+            <configuration>
+                <property>
+                    <name>hue-id-w</name>
+                    <value>50001</value>
+                </property>
+            </configuration>
+        </sub-workflow>
+        <ok to="join-775e"/>
+        <error to="Kill"/>
+    </action>
+    <action name="shell-0f44">
+        <shell xmlns="uri:oozie:shell-action:0.1">
+            <job-tracker>${jobTracker}</job-tracker>
+            <name-node>${nameNode}</name-node>
+            <exec>ls</exec>
+              <capture-output/>
+        </shell>
+        <ok to="join-775e"/>
+        <error to="Kill"/>
+    </action>
+    <fork name="fork-68d4">
+        <path start="subworkflow-a13f" />
+        <path start="shell-0f44" />
+    </fork>
+    <join name="join-775e" to="End"/>
+    <end name="End"/>
+</workflow-app>

+ 169 - 1
apps/oozie/src/oozie/tests2.py

@@ -24,7 +24,8 @@ from django.core.urlresolvers import reverse
 from nose.tools import assert_true, assert_false, assert_equal, assert_not_equal
 
 from oozie.conf import ENABLE_V2
-from oozie.models2 import Workflow, find_dollar_variables, find_dollar_braced_variables, Node
+from oozie.importlib.workflows import generate_v2_graph_nodes
+from oozie.models2 import Workflow, find_dollar_variables, find_dollar_braced_variables, Node, _create_graph_adjaceny_list, _get_hierarchy_from_adj_list, _create_workflow_layout, _add_uuids_to_adj_list
 from oozie.tests import OozieMockBase, save_temp_workflow, MockOozieApi
 
 
@@ -292,3 +293,170 @@ LIMIT $limit"""))
     finally:
       reset()
       wf_doc.delete()
+
+
+class TestExternalWorkflowGraph():
+
+  def setUp(self):
+    self.wf = Workflow()
+
+  def test_graph_generation_from_xml(self):
+    f = open('apps/oozie/src/oozie/test_data/xslt2/test-workflow.xml')
+    self.wf.definition = f.read()
+    self.node_list = [{u'node_type': u'start', u'ok_to': u'fork-68d4', u'name': u''}, {u'node_type': u'kill', u'ok_to': u'', u'name': u'Kill'}, {u'path2': u'shell-0f44', u'node_type': u'fork', u'ok_to': u'', u'name': u'fork-68d4', u'path1': u'subworkflow-a13f'}, {u'node_type': u'join', u'ok_to': u'End', u'name': u'join-775e'}, {u'node_type': u'end', u'ok_to': u'', u'name': u'End'}, {u'node_type': u'sub-workflow', u'ok_to': u'join-775e', u'sub-workflow': {u'app-path': u'${nameNode}/user/hue/oozie/deployments/_admin_-oozie-50001-1427488969.48'}, u'name': u'subworkflow-a13f', u'error_to': u'Kill'}, {u'shell': {u'command': u'ls'}, u'node_type': u'shell', u'ok_to': u'join-775e', u'name': u'shell-0f44', u'error_to': u'Kill'}, {}]
+    assert_equal(self.node_list, generate_v2_graph_nodes(self.wf.definition))
+
+  def test_get_graph_adjacency_list(self):
+    self.node_list = [{u'node_type': u'start', u'ok_to': u'fork-68d4', u'name': u''}, {u'node_type': u'kill', u'ok_to': u'', u'name': u'Kill'}, {u'path2': u'shell-0f44', u'node_type': u'fork', u'ok_to': u'', u'name': u'fork-68d4', u'path1': u'subworkflow-a13f'}, {u'node_type': u'join', u'ok_to': u'End', u'name': u'join-775e'}, {u'node_type': u'end', u'ok_to': u'', u'name': u'End'}, {u'node_type': u'sub-workflow', u'ok_to': u'join-775e', u'sub-workflow': {u'app-path': u'${nameNode}/user/hue/oozie/deployments/_admin_-oozie-50001-1427488969.48'}, u'name': u'subworkflow-a13f', u'error_to': u'Kill'}, {u'shell': {u'command': u'ls'}, u'node_type': u'shell', u'ok_to': u'join-775e', u'name': u'shell-0f44', u'error_to': u'Kill'}, {}]
+    adj_list = _create_graph_adjaceny_list(self.node_list)
+
+    assert_true(len(adj_list) == 7)
+    assert_true('subworkflow-a13f' in adj_list.keys())
+    assert_true(adj_list['shell-0f44']['shell']['command'] == 'ls')
+    assert_equal(adj_list['fork-68d4'], {u'path2': u'shell-0f44', u'node_type': u'fork', u'ok_to': u'', u'name': u'fork-68d4', u'path1': u'subworkflow-a13f'})
+
+  def test_get_hierarchy_from_adj_list(self):
+    self.wf.definition = """<workflow-app name="ls-4thread" xmlns="uri:oozie:workflow:0.5">
+        <start to="fork-fe93"/>
+        <kill name="Kill">
+            <message>Action failed, error message[${wf:errorMessage(wf:lastErrorNode())}]</message>
+        </kill>
+        <action name="shell-5429">
+            <shell xmlns="uri:oozie:shell-action:0.1">
+                <job-tracker>${jobTracker}</job-tracker>
+                <name-node>${nameNode}</name-node>
+                <exec>ls</exec>
+                  <capture-output/>
+            </shell>
+            <ok to="join-7f80"/>
+            <error to="Kill"/>
+        </action>
+        <action name="shell-bd90">
+            <shell xmlns="uri:oozie:shell-action:0.1">
+                <job-tracker>${jobTracker}</job-tracker>
+                <name-node>${nameNode}</name-node>
+                <exec>ls</exec>
+                  <capture-output/>
+            </shell>
+            <ok to="join-7f80"/>
+            <error to="Kill"/>
+        </action>
+        <fork name="fork-fe93">
+            <path start="shell-5429" />
+            <path start="shell-bd90" />
+            <path start="shell-d64c" />
+            <path start="shell-d8cc" />
+        </fork>
+        <join name="join-7f80" to="End"/>
+        <action name="shell-d64c">
+            <shell xmlns="uri:oozie:shell-action:0.1">
+                <job-tracker>${jobTracker}</job-tracker>
+                <name-node>${nameNode}</name-node>
+                <exec>ls</exec>
+                  <capture-output/>
+            </shell>
+            <ok to="join-7f80"/>
+            <error to="Kill"/>
+        </action>
+        <action name="shell-d8cc">
+            <shell xmlns="uri:oozie:shell-action:0.1">
+                <job-tracker>${jobTracker}</job-tracker>
+                <name-node>${nameNode}</name-node>
+                <exec>ls</exec>
+                  <capture-output/>
+            </shell>
+            <ok to="join-7f80"/>
+            <error to="Kill"/>
+        </action>
+        <end name="End"/>
+    </workflow-app>"""
+
+    node_list = generate_v2_graph_nodes(self.wf.definition)
+    adj_list = _create_graph_adjaceny_list(node_list)
+
+    node_hierarchy = ['start']
+    _get_hierarchy_from_adj_list(adj_list, adj_list['start']['ok_to'], node_hierarchy)
+
+    assert_equal(node_hierarchy, ['start', [u'fork-fe93', [[u'shell-bd90'], [u'shell-d64c'], [u'shell-5429'], [u'shell-d8cc']], u'join-7f80'], ['Kill'], ['End']])
+
+  def test_gen_workflow_data_from_xml(self):
+    self.wf.definition = """<workflow-app name="fork-fork-test" xmlns="uri:oozie:workflow:0.5">
+        <start to="fork-949d"/>
+        <kill name="Kill">
+            <message>Action failed, error message[${wf:errorMessage(wf:lastErrorNode())}]</message>
+        </kill>
+        <action name="shell-eadd">
+            <shell xmlns="uri:oozie:shell-action:0.1">
+                <job-tracker>${jobTracker}</job-tracker>
+                <name-node>${nameNode}</name-node>
+                <exec>ls</exec>
+                  <capture-output/>
+            </shell>
+            <ok to="join-1a0f"/>
+            <error to="Kill"/>
+        </action>
+        <action name="shell-f4c1">
+            <shell xmlns="uri:oozie:shell-action:0.1">
+                <job-tracker>${jobTracker}</job-tracker>
+                <name-node>${nameNode}</name-node>
+                <exec>ls</exec>
+                  <capture-output/>
+            </shell>
+            <ok to="join-3bba"/>
+            <error to="Kill"/>
+        </action>
+        <fork name="fork-949d">
+            <path start="fork-e5fa" />
+            <path start="shell-3dd5" />
+        </fork>
+        <join name="join-ca1a" to="End"/>
+        <action name="shell-ef70">
+            <shell xmlns="uri:oozie:shell-action:0.1">
+                <job-tracker>${jobTracker}</job-tracker>
+                <name-node>${nameNode}</name-node>
+                <exec>ls</exec>
+                  <capture-output/>
+            </shell>
+            <ok to="join-1a0f"/>
+            <error to="Kill"/>
+        </action>
+        <fork name="fork-37d7">
+            <path start="shell-eadd" />
+            <path start="shell-ef70" />
+        </fork>
+        <join name="join-1a0f" to="join-ca1a"/>
+        <action name="shell-3dd5">
+            <shell xmlns="uri:oozie:shell-action:0.1">
+                <job-tracker>${jobTracker}</job-tracker>
+                <name-node>${nameNode}</name-node>
+                <exec>ls</exec>
+                  <capture-output/>
+            </shell>
+            <ok to="fork-37d7"/>
+            <error to="Kill"/>
+        </action>
+        <action name="shell-2ba8">
+            <shell xmlns="uri:oozie:shell-action:0.1">
+                <job-tracker>${jobTracker}</job-tracker>
+                <name-node>${nameNode}</name-node>
+                <exec>ls</exec>
+                  <capture-output/>
+            </shell>
+            <ok to="join-3bba"/>
+            <error to="Kill"/>
+        </action>
+        <fork name="fork-e5fa">
+            <path start="shell-f4c1" />
+            <path start="shell-2ba8" />
+        </fork>
+        <join name="join-3bba" to="join-ca1a"/>
+        <end name="End"/>
+    </workflow-app>"""
+
+    workflow_data = Workflow.gen_workflow_data_from_xml('test', self.wf)
+
+    assert_true(len(workflow_data['layout'][0]['rows']) == 6)
+    assert_true(len(workflow_data['workflow']['nodes']) == 14)
+    assert_equal(workflow_data['layout'][0]['rows'][1]['widgets'][0]['widgetType'], 'fork-widget')
+    assert_equal(workflow_data['workflow']['nodes'][0]['name'], 'start-3f10')
+

+ 6 - 5
apps/oozie/src/oozie/views/dashboard.py

@@ -325,7 +325,6 @@ def list_oozie_workflow(request, job_id):
       if hue_workflow: hue_workflow.document.doc.get().can_read_or_exception(request.user)
 
       if hue_workflow:
-        workflow_graph = ''
         full_node_list = hue_workflow.nodes
         workflow_id = hue_workflow.id
         wid = {
@@ -334,11 +333,13 @@ def list_oozie_workflow(request, job_id):
         doc = Document2.objects.get(type='oozie-workflow2', **wid)
         new_workflow = get_workflow()(document=doc)
         workflow_data = new_workflow.get_data()
-        credentials = Credentials()
       else:
-        # For workflows submitted from CLI or deleted in the editor
-        # Until better parsing in https://issues.cloudera.org/browse/HUE-2659
-        workflow_graph, full_node_list = OldWorkflow.gen_status_graph_from_xml(request.user, oozie_workflow)
+        try:
+          workflow_data = Workflow.gen_workflow_data_from_xml(request.user, oozie_workflow)
+        except Exception, e:
+          LOG.exception('Graph data could not be generated from Workflow %s: %s' % (oozie_workflow.id, e))
+      workflow_graph = ''
+      credentials = Credentials()
     except:
       LOG.exception("Error generating full page for running workflow %s" % job_id)
   else: