|
|
@@ -14,7 +14,6 @@
|
|
|
# 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 desktop.models import Document
|
|
|
"""
|
|
|
Import an external workflow by providing an XML definition.
|
|
|
The workflow definition is imported via the method 'import_workflow'.
|
|
|
@@ -32,23 +31,20 @@ Action extensions are also versioned.
|
|
|
Every action extension will have its own version via /xslt/<workflow version>/extensions/<name of extensions>.<version>.xslt
|
|
|
"""
|
|
|
|
|
|
-try:
|
|
|
- import json
|
|
|
-except ImportError:
|
|
|
- import simplejson as json
|
|
|
-
|
|
|
import logging
|
|
|
from lxml import etree
|
|
|
+import os
|
|
|
|
|
|
from django.core import serializers
|
|
|
from django.utils.encoding import smart_str
|
|
|
from django.utils.translation import ugettext as _
|
|
|
|
|
|
-from conf import DEFINITION_XSLT_DIR
|
|
|
-from models import Workflow, Node, Link, Start, End,\
|
|
|
- Decision, DecisionEnd, Fork, Join,\
|
|
|
- Kill
|
|
|
-from utils import xml_tag
|
|
|
+from desktop.models import Document
|
|
|
+
|
|
|
+from oozie.conf import DEFINITION_XSLT_DIR
|
|
|
+from oozie.models import Workflow, Node, Link, Start, End,\
|
|
|
+ Decision, DecisionEnd, Fork, Join,\
|
|
|
+ Kill
|
|
|
|
|
|
LOG = logging.getLogger(__name__)
|
|
|
|
|
|
@@ -97,7 +93,7 @@ def _save_links(workflow, root):
|
|
|
if child_el.tag.endswith('global'):
|
|
|
continue
|
|
|
|
|
|
- tag = xml_tag(child_el)
|
|
|
+ tag = etree.QName(child_el).localname
|
|
|
name = child_el.attrib.get('name', tag)
|
|
|
LOG.debug("Getting node with data - XML TAG: %(tag)s\tLINK NAME: %(node_name)s\tWORKFLOW NAME: %(workflow_name)s" % {
|
|
|
'tag': smart_str(tag),
|
|
|
@@ -202,7 +198,7 @@ def _decision_relationships(workflow, parent, child_el):
|
|
|
except Node.DoesNotExist, e:
|
|
|
raise RuntimeError(_("Node %s has not been defined.") % to)
|
|
|
|
|
|
- if xml_tag(case) == 'default':
|
|
|
+ if etree.QName(case).localname == 'default':
|
|
|
name = 'default'
|
|
|
obj = Link.objects.create(name=name, parent=parent, child=child)
|
|
|
|
|
|
@@ -226,7 +222,7 @@ def _node_relationships(workflow, parent, child_el):
|
|
|
continue
|
|
|
|
|
|
# Links
|
|
|
- name = xml_tag(el)
|
|
|
+ name = etree.QName(el).localname
|
|
|
if name in LINKS:
|
|
|
if name == 'path':
|
|
|
if 'start' not in el.attrib:
|
|
|
@@ -483,9 +479,9 @@ def _preprocess_nodes(workflow, transformed_root, workflow_definition_root, node
|
|
|
for action_el in workflow_definition_root:
|
|
|
if 'name' in action_el.attrib and action_el.attrib['name'] == full_node.name:
|
|
|
for subworkflow_el in action_el:
|
|
|
- if xml_tag(subworkflow_el) == 'sub-workflow':
|
|
|
+ if etree.QName(subworkflow_el).localname == 'sub-workflow':
|
|
|
for property_el in subworkflow_el:
|
|
|
- if xml_tag(property_el) == 'app-path':
|
|
|
+ if etree.QName(property_el).localname == 'app-path':
|
|
|
app_path = property_el.text
|
|
|
|
|
|
if app_path is None:
|
|
|
@@ -516,7 +512,7 @@ def _resolve_subworkflow_from_deployment_dir(fs, workflow, app_path):
|
|
|
f = fs.open('%s/workflow.xml' % app_path)
|
|
|
root = etree.parse(f)
|
|
|
f.close()
|
|
|
- return Workflow.objects.get(name=root.attrib['name'])
|
|
|
+ return Workflow.objects.get(name=root.attrib['name'], owner=workflow.owner, managed=True)
|
|
|
except IOError:
|
|
|
pass
|
|
|
except (KeyError, AttributeError), e:
|
|
|
@@ -548,41 +544,76 @@ def _save_nodes(workflow, nodes):
|
|
|
node.save()
|
|
|
|
|
|
|
|
|
-def import_workflow(workflow, workflow_definition, fs=None):
|
|
|
- xslt_definition_fh = open("%(xslt_dir)s/workflow.xslt" % {
|
|
|
- 'xslt_dir': DEFINITION_XSLT_DIR.get()
|
|
|
- })
|
|
|
-
|
|
|
- # Parse Workflow Definition
|
|
|
- workflow_definition_root = etree.fromstring(workflow_definition)
|
|
|
+def _process_metadata(workflow, metadata):
|
|
|
+ # Job attributes
|
|
|
+ attributes = metadata.setdefault('attributes', {})
|
|
|
+ workflow.description = attributes.setdefault('description', workflow.description)
|
|
|
+ workflow.deployment_dir = attributes.setdefault('deployment_dir', workflow.deployment_dir)
|
|
|
|
|
|
- if workflow_definition_root is None:
|
|
|
- raise RuntimeError(_("Could not find any nodes in Workflow definition. Maybe it's malformed?"))
|
|
|
+ # Workflow node attributes
|
|
|
+ nodes = metadata.setdefault('nodes', {})
|
|
|
+ for node_name in nodes:
|
|
|
+ try:
|
|
|
+ node = Node.objects.get(name=node_name, workflow=workflow).get_full_node()
|
|
|
+ node_attributes = nodes[node_name].setdefault('attributes', {})
|
|
|
+ for node_attribute in node_attributes:
|
|
|
+ setattr(node, node_attribute, node_attributes[node_attribute])
|
|
|
+ node.save()
|
|
|
+ except Node.DoesNotExist:
|
|
|
+ # @TODO(abe): Log or throw error?
|
|
|
+ raise
|
|
|
+ except AttributeError:
|
|
|
+ # @TODO(abe): Log or throw error?
|
|
|
+ # Here there was an attribute reference in the metadata
|
|
|
+ # for this node that isn't a member of the node.
|
|
|
+ raise
|
|
|
|
|
|
- ns = workflow_definition_root.tag[:-12] # Remove workflow-app from tag in order to get proper namespace prefix
|
|
|
- schema_version = ns and ns[1:-1] or None
|
|
|
|
|
|
- # Ensure namespace exists
|
|
|
- if schema_version not in OOZIE_NAMESPACES:
|
|
|
- raise RuntimeError(_("Tag with namespace %(namespace)s is not valid. Please use one of the following namespaces: %(namespaces)s") % {
|
|
|
- 'namespace': workflow_definition_root.tag,
|
|
|
- 'namespaces': ', '.join(OOZIE_NAMESPACES)
|
|
|
+def import_workflow_root(workflow, workflow_definition_root, metadata=None, fs=None):
|
|
|
+ try:
|
|
|
+ xslt_definition_fh = open("%(xslt_dir)s/workflow.xslt" % {
|
|
|
+ 'xslt_dir': os.path.join(DEFINITION_XSLT_DIR.get(), 'workflows')
|
|
|
})
|
|
|
|
|
|
- # 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)
|
|
|
-
|
|
|
- # Resolve workflow dependencies and node types and link dependencies
|
|
|
- nodes = _prepare_nodes(workflow, transformed_root)
|
|
|
- _preprocess_nodes(workflow, transformed_root, workflow_definition_root, nodes, fs)
|
|
|
- _save_nodes(workflow, nodes)
|
|
|
- _save_links(workflow, workflow_definition_root)
|
|
|
+ tag = etree.QName(workflow_definition_root.tag)
|
|
|
+ schema_version = tag.namespace
|
|
|
+
|
|
|
+ # Ensure namespace exists
|
|
|
+ if schema_version not in OOZIE_NAMESPACES:
|
|
|
+ raise RuntimeError(_("Tag with namespace %(namespace)s is not valid. Please use one of the following namespaces: %(namespaces)s") % {
|
|
|
+ 'namespace': workflow_definition_root.tag,
|
|
|
+ 'namespaces': ', '.join(OOZIE_NAMESPACES)
|
|
|
+ })
|
|
|
+
|
|
|
+ # 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)
|
|
|
+
|
|
|
+ # Resolve workflow dependencies and node types and link dependencies
|
|
|
+ nodes = _prepare_nodes(workflow, transformed_root)
|
|
|
+ _preprocess_nodes(workflow, transformed_root, workflow_definition_root, nodes, fs)
|
|
|
+ _save_nodes(workflow, nodes)
|
|
|
+ _save_links(workflow, workflow_definition_root)
|
|
|
+ if metadata:
|
|
|
+ _process_metadata(workflow, metadata)
|
|
|
+
|
|
|
+ # Update workflow attributes
|
|
|
+ workflow.schema_version = schema_version
|
|
|
+ workflow.name = workflow_definition_root.get('name')
|
|
|
+ workflow.save()
|
|
|
+ except:
|
|
|
+ workflow.delete()
|
|
|
+ raise
|
|
|
+
|
|
|
+
|
|
|
+def import_workflow(workflow, workflow_definition, metadata=None, fs=None):
|
|
|
+ # Parse Workflow Definition
|
|
|
+ workflow_definition_root = etree.fromstring(workflow_definition)
|
|
|
+ if workflow_definition_root is None:
|
|
|
+ raise RuntimeError(_("Could not find any nodes in Workflow definition. Maybe it's malformed?"))
|
|
|
|
|
|
- # Update schema_version
|
|
|
- workflow.schema_version = schema_version
|
|
|
- workflow.save()
|
|
|
+ return import_workflow_root(workflow, workflow_definition_root, metadata, fs)
|