tests.py 161 KB


  1. #!/usr/bin/env python
  2. ## -*- coding: utf-8 -*-
  3. # Licensed to Cloudera, Inc. under one
  4. # or more contributor license agreements. See the NOTICE file
  5. # distributed with this work for additional information
  6. # regarding copyright ownership. Cloudera, Inc. licenses this file
  7. # to you under the Apache License, Version 2.0 (the
  8. # "License"); you may not use this file except in compliance
  9. # with the License. You may obtain a copy of the License at
  10. #
  11. # http://www.apache.org/licenses/LICENSE-2.0
  12. #
  13. # Unless required by applicable law or agreed to in writing, software
  14. # distributed under the License is distributed on an "AS IS" BASIS,
  15. # WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
  16. # See the License for the specific language governing permissions and
  17. # limitations under the License.
  18. import json
  19. import logging
  20. import re
  21. import os
  22. import StringIO
  23. import shutil
  24. import tempfile
  25. import zipfile
  26. from datetime import datetime
  27. from itertools import chain
  28. from nose.plugins.skip import SkipTest
  29. from nose.tools import raises, assert_true, assert_false, assert_equal, assert_not_equal
  30. from django.contrib.auth.models import User
  31. from django.core.urlresolvers import reverse
  32. from desktop.lib.django_test_util import make_logged_in_client
  33. from desktop.lib.test_utils import grant_access, add_permission, add_to_group, reformat_json, reformat_xml
  34. from desktop.models import Document
  35. from jobsub.models import OozieDesign, OozieMapreduceAction
  36. from liboozie import oozie_api
  37. from liboozie.conf import OOZIE_URL
  38. from liboozie.oozie_api_test import OozieServerProvider
  39. from liboozie.types import WorkflowList, Workflow as OozieWorkflow, Coordinator as OozieCoordinator,\
  40. Bundle as OozieBundle, CoordinatorList, WorkflowAction, BundleList
  41. from oozie.models import Workflow, Node, Kill, Link, Job, Coordinator, History,\
  42. find_parameters, NODE_TYPES, Bundle
  43. from oozie.utils import workflow_to_dict, model_to_dict, smart_path
  44. from oozie.importlib.workflows import import_workflow
  45. from oozie.importlib.jobdesigner import convert_jobsub_design
  46. LOG = logging.getLogger(__name__)
  47. _INITIALIZED = False
  48. class MockOozieApi:
  49. JSON_WORKFLOW_LIST = [{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=0000012-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[0000012-120725142744176-oozie-oozi-W] status[SUCCEEDED]', u'endTime': u'Mon, 30 Jul 2012 22:37:00 GMT', u'id': u'0000012-120725142744176-oozie-oozi-W', u'group': None, u'user': u'test'},
  50. {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'},
  51. {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'},
  52. {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'}]
  53. WORKFLOW_IDS = [wf['id'] for wf in JSON_WORKFLOW_LIST]
  54. WORKFLOW_DICT = dict([(wf['id'], wf) for wf in JSON_WORKFLOW_LIST])
  55. JSON_COORDINATOR_LIST = [{u'startTime': u'Sun, 01 Jul 2012 00:00:00 GMT', u'actions': [], u'frequency': 1, u'concurrency': 1, u'pauseTime': None, u'group': None, u'toString': u'Coornidator application id[0000041-120717205528122-oozie-oozi-C] status[RUNNING]', u'consoleUrl': None, u'mat_throttling': 0, u'status': u'RUNNING', u'conf': None, u'user': u'test', u'timeOut': 120, u'coordJobPath': u'hdfs://localhost:8020/user/test/demo2', u'timeUnit': u'DAY', u'coordJobId': u'0000041-120717205528122-oozie-oozi-C', u'coordJobName': u'DailyWordCount1', u'nextMaterializedTime': u'Wed, 04 Jul 2012 00:00:00 GMT', u'coordExternalId': None, u'acl': None, u'lastAction': u'Wed, 04 Jul 2012 00:00:00 GMT', u'executionPolicy': u'FIFO', u'timeZone': u'America/Los_Angeles', u'endTime': u'Wed, 04 Jul 2012 00:00:00 GMT'},
  56. {u'startTime': u'Sun, 01 Jul 2012 00:00:00 GMT', u'actions': [], u'frequency': 1, u'concurrency': 1, u'pauseTime': None, u'group': None, u'toString': u'Coornidator application id[0000011-120706144403213-oozie-oozi-C] status[DONEWITHERROR]', u'consoleUrl': None, u'mat_throttling': 0, u'status': u'DONEWITHERROR', u'conf': None, u'user': u'test', u'timeOut': 120, u'coordJobPath': u'hdfs://localhost:8020/user/hue/jobsub/_romain_-design-2', u'timeUnit': u'DAY', u'coordJobId': u'0000011-120706144403213-oozie-oozi-C', u'coordJobName': u'DailyWordCount2', u'nextMaterializedTime': u'Thu, 05 Jul 2012 00:00:00 GMT', u'coordExternalId': None, u'acl': None, u'lastAction': u'Thu, 05 Jul 2012 00:00:00 GMT', u'executionPolicy': u'FIFO', u'timeZone': u'America/Los_Angeles', u'endTime': u'Wed, 04 Jul 2012 18:54:00 GMT'},
  57. {u'startTime': u'Sun, 01 Jul 2012 00:00:00 GMT', u'actions': [], u'frequency': 1, u'concurrency': 1, u'pauseTime': None, u'group': None, u'toString': u'Coornidator application id[0000010-120706144403213-oozie-oozi-C] status[DONEWITHERROR]', u'consoleUrl': None, u'mat_throttling': 0, u'status': u'DONEWITHERROR', u'conf': None, u'user': u'test', u'timeOut': 120, u'coordJobPath': u'hdfs://localhost:8020/user/hue/jobsub/_romain_-design-2', u'timeUnit': u'DAY', u'coordJobId': u'0000010-120706144403213-oozie-oozi-C', u'coordJobName': u'DailyWordCount3', u'nextMaterializedTime': u'Thu, 05 Jul 2012 00:00:00 GMT', u'coordExternalId': None, u'acl': None, u'lastAction': u'Thu, 05 Jul 2012 00:00:00 GMT', u'executionPolicy': u'FIFO', u'timeZone': u'America/Los_Angeles', u'endTime': u'Wed, 04 Jul 2012 18:54:00 GMT'},
  58. {u'startTime': u'Sun, 01 Jul 2012 00:00:00 GMT', u'actions': [], u'frequency': 1, u'concurrency': 1, u'pauseTime': None, u'group': None, u'toString': u'Coornidator application id[0000009-120706144403213-oozie-oozi-C] status[DONEWITHERROR]', u'consoleUrl': None, u'mat_throttling': 0, u'status': u'DONEWITHERROR', u'conf': None, u'user': u'test', u'timeOut': 120, u'coordJobPath': u'hdfs://localhost:8020/user/hue/jobsub/_romain_-design-2', u'timeUnit': u'DAY', u'coordJobId': u'0000009-120706144403213-oozie-oozi-C', u'coordJobName': u'DailyWordCount4', u'nextMaterializedTime': u'Thu, 05 Jul 2012 00:00:00 GMT', u'coordExternalId': None, u'acl': None, u'lastAction': u'Thu, 05 Jul 2012 00:00:00 GMT', u'executionPolicy': u'FIFO', u'timeZone': u'America/Los_Angeles', u'endTime': u'Wed, 04 Jul 2012 18:54:00 GMT'},
  59. {u'startTime': u'Sun, 01 Jul 2012 00:00:00 GMT', u'actions': [], u'frequency': 1, u'concurrency': 1, u'pauseTime': None, u'group': None, u'toString': u'Coornidator application id[00000012-120706144403213-oozie-oozi-C] status[DONEWITHERROR]', u'consoleUrl': None, u'mat_throttling': 0, u'status': u'DONEWITHERROR', u'conf': None, u'user': u'test', u'timeOut': 120, u'coordJobPath': u'hdfs://localhost:8020/user/hue/jobsub/_romain_-design-2', u'timeUnit': u'DAY', u'coordJobId': u'0000009-120706144403213-oozie-oozi-C', u'coordJobName': u'DåilyWordCount5', u'nextMaterializedTime': u'Thu, 05 Jul 2012 00:00:00 GMT', u'coordExternalId': None, u'acl': None, u'lastAction': u'Thu, 05 Jul 2012 00:00:00 GMT', u'executionPolicy': u'FIFO', u'timeZone': u'America/Los_Angeles', u'endTime': u'Wed, 04 Jul 2012 18:54:00 GMT'}]
  60. COORDINATOR_IDS = [coord['coordJobId'] for coord in JSON_COORDINATOR_LIST]
  61. COORDINATOR_DICT = dict([(coord['coordJobId'], coord) for coord in JSON_COORDINATOR_LIST])
  62. JSON_BUNDLE_LIST = [
  63. {u'status': u'SUCCEEDED', u'toString': u'Bundle id[0000021-130210132208494-oozie-oozi-B] status[SUCCEEDED]', u'group': None, u'conf': u'<configuration>\r\n <property>\r\n <name>oozie.bundle.application.path</name>\r\n <value>hdfs://localhost:8020/user/hue/oozie/workspaces/_romain_-oozie-22-1360636939.69</value>\r\n </property>\r\n <property>\r\n <name>user.name</name>\r\n <value>romain</value>\r\n </property>\r\n <property>\r\n <name>oozie.use.system.libpath</name>\r\n <value>true</value>\r\n </property>\r\n <property>\r\n <name>nameNode</name>\r\n <value>hdfs://localhost:8020</value>\r\n </property>\r\n <property>\r\n <name>wf_application_path</name>\r\n <value>hdfs://localhost:8020/user/hue/oozie/deployments/_romain_-oozie-5-1360649203.07</value>\r\n </property>\r\n <property>\r\n <name>jobTracker</name>\r\n <value>localhost:8021</value>\r\n </property>\r\n <property>\r\n <name>hue-id-b</name>\r\n <value>22</value>\r\n </property>\r\n</configuration>', u'bundleJobName': u'MyBundle1', u'startTime': None, u'bundleCoordJobs': [], u'kickoffTime': u'Mon, 11 Feb 2013 08:33:00 PST', u'acl': None, u'bundleJobPath': u'hdfs://localhost:8020/user/hue/oozie/workspaces/_romain_-oozie-22-1360636939.69', u'createdTime': u'Mon, 11 Feb 2013 22:06:44 PST', u'timeOut': 0, u'consoleUrl': None, u'bundleExternalId': None, u'timeUnit': u'NONE', u'pauseTime': None, u'bundleJobId': u'0000021-130210132208494-oozie-oozi-B', u'endTime': None, u'user': u'test'},
  64. {u'status': u'KILLED', u'toString': u'Bundle id[0000020-130210132208494-oozie-oozi-B] status[KILLED]', u'group': None, u'conf': u'<configuration>\r\n <property>\r\n <name>oozie.bundle.application.path</name>\r\n <value>hdfs://localhost:8020/user/hue/oozie/workspaces/_romain_-oozie-22-1360636939.69</value>\r\n </property>\r\n <property>\r\n <name>user.name</name>\r\n <value>romain</value>\r\n </property>\r\n <property>\r\n <name>oozie.use.system.libpath</name>\r\n <value>true</value>\r\n </property>\r\n <property>\r\n <name>nameNode</name>\r\n <value>hdfs://localhost:8020</value>\r\n </property>\r\n <property>\r\n <name>wf_application_path</name>\r\n <value>hdfs://localhost:8020/user/hue/oozie/deployments/_romain_-oozie-5-1360648977.2</value>\r\n </property>\r\n <property>\r\n <name>jobTracker</name>\r\n <value>localhost:8021</value>\r\n </property>\r\n <property>\r\n <name>hue-id-b</name>\r\n <value>22</value>\r\n </property>\r\n</configuration>', u'bundleJobName': u'MyBundle2', u'startTime': None, u'bundleCoordJobs': [], u'kickoffTime': u'Mon, 11 Feb 2013 08:33:00 PST', u'acl': None, u'bundleJobPath': u'hdfs://localhost:8020/user/hue/oozie/workspaces/_romain_-oozie-22-1360636939.69', u'createdTime': u'Mon, 11 Feb 2013 22:02:58 PST', u'timeOut': 0, u'consoleUrl': None, u'bundleExternalId': None, u'timeUnit': u'NONE', u'pauseTime': None, u'bundleJobId': u'0000020-130210132208494-oozie-oozi-B', u'endTime': None, u'user': u'test'},
  65. {u'status': u'KILLED', u'toString': u'Bundle id[0000019-130210132208494-oozie-oozi-B] status[KILLED]', u'group': None, u'conf': u'<configuration>\r\n <property>\r\n <name>oozie.bundle.application.path</name>\r\n <value>hdfs://localhost:8020/user/hue/oozie/workspaces/_romain_-oozie-22-1360636939.69</value>\r\n </property>\r\n <property>\r\n <name>user.name</name>\r\n <value>romain</value>\r\n </property>\r\n <property>\r\n <name>oozie.use.system.libpath</name>\r\n <value>true</value>\r\n </property>\r\n <property>\r\n <name>nameNode</name>\r\n <value>hdfs://localhost:8020</value>\r\n </property>\r\n <property>\r\n <name>wf_application_path</name>\r\n <value>hdfs://localhost:8020/user/hue/oozie/deployments/_romain_-oozie-5-1360639372.41</value>\r\n </property>\r\n <property>\r\n <name>jobTracker</name>\r\n <value>localhost:8021</value>\r\n </property>\r\n <property>\r\n <name>hue-id-b</name>\r\n <value>22</value>\r\n </property>\r\n</configuration>', u'bundleJobName': u'MyBundle3', u'startTime': None, u'bundleCoordJobs': [], u'kickoffTime': u'Mon, 11 Feb 2013 08:33:00 PST', u'acl': None, u'bundleJobPath': u'hdfs://localhost:8020/user/hue/oozie/workspaces/_romain_-oozie-22-1360636939.69', u'createdTime': u'Mon, 11 Feb 2013 19:22:53 PST', u'timeOut': 0, u'consoleUrl': None, u'bundleExternalId': None, u'timeUnit': u'NONE', u'pauseTime': None, u'bundleJobId': u'0000019-130210132208494-oozie-oozi-B', u'endTime': None, u'user': u'test'}
  66. ]
  67. BUNDLE_IDS = [bundle['bundleJobId'] for bundle in JSON_BUNDLE_LIST]
  68. BUNDLE_DICT = dict([(bundle['bundleJobId'], bundle) for bundle in JSON_BUNDLE_LIST])
  69. WORKFLOW_ACTION = {u'status': u'OK', u'retries': 0, u'transition': u'end', u'stats': None, u'startTime': u'Fri, 10 Aug 2012 05:24:21 GMT', u'toString': u'Action name[WordCount] status[OK]', u'cred': u'null', u'errorMessage': None, u'errorCode': None, u'consoleUrl': u'http://localhost:50030/jobdetails.jsp?jobid=job_201208072118_0044', u'externalId': u'job_201208072118_0044', u'externalStatus': u'SUCCEEDED', u'conf': u'<map-reduce xmlns="uri:oozie:workflow:0.4">\r\n <job-tracker>localhost:8021</job-tracker>\r\n <name-node>hdfs://localhost:8020</name-node>\r\n <configuration>\r\n <property>\r\n <name>mapred.mapper.regex</name>\r\n <value>dream</value>\r\n </property>\r\n <property>\r\n <name>mapred.input.dir</name>\r\n <value>/user/test/words/20120702</value>\r\n </property>\r\n <property>\r\n <name>mapred.output.dir</name>\r\n <value>/user/test/out/rrwords/20120702</value>\r\n </property>\r\n <property>\r\n <name>mapred.mapper.class</name>\r\n <value>org.apache.hadoop.mapred.lib.RegexMapper</value>\r\n </property>\r\n <property>\r\n <name>mapred.combiner.class</name>\r\n <value>org.apache.hadoop.mapred.lib.LongSumReducer</value>\r\n </property>\r\n <property>\r\n <name>mapred.reducer.class</name>\r\n <value>org.apache.hadoop.mapred.lib.LongSumReducer</value>\r\n </property>\r\n <property>\r\n <name>mapred.output.key.class</name>\r\n <value>org.apache.hadoop.io.Text</value>\r\n </property>\r\n <property>\r\n <name>mapred.output.value.class</name>\r\n <value>org.apache.hadoop.io.LongWritable</value>\r\n </property>\r\n </configuration>\r\n</map-reduce>', u'type': u'map-reduce', u'trackerUri': u'localhost:8021', u'externalChildIDs': None, u'endTime': u'Fri, 10 Aug 2012 05:24:38 GMT', u'data': None, u'id': u'0000012-120725142744176-oozie-oozi-W@WordCount', u'name': u'WordCount', u'externalChildIDs': u'job_201302280955_0018,job_201302280955_0019,job_201302280955_0020'}
  70. JSON_ACTION_BUNDLE_LIST = [
  71. {u'status': u'SUCCEEDED', u'toString': u'WorkflowAction name[0000000-130117135211239-oozie-oozi-C@1] status[SUCCEEDED]', u'runConf': None, u'errorMessage': None, u'missingDependencies': u'', u'coordJobId': u'0000000-130117135211239-oozie-oozi-C', u'errorCode': None, u'actionNumber': 1, u'consoleUrl': None, u'nominalTime': u'Mon, 31 Dec 2012 16:00:00 PST', u'externalStatus': u'', u'createdConf': None, u'createdTime': u'Fri, 25 Jan 2013 10:53:39 PST', u'externalId': u'0000035-130124125317829-oozie-oozi-W', u'lastModifiedTime': u'Fri, 25 Jan 2013 10:53:51 PST', u'type': None, u'id': u'0000000-130117135211239-oozie-oozi-C@1', u'trackerUri': None},
  72. {u'status': u'SUCCEEDED', u'toString': u'WorkflowAction name[0000000-130117135211239-oozie-oozi-C@2] status[SUCCEEDED]', u'runConf': None, u'errorMessage': None, u'missingDependencies': u'', u'coordJobId': u'0000000-130117135211239-oozie-oozi-C', u'errorCode': None, u'actionNumber': 2, u'consoleUrl': None, u'nominalTime': u'Tue, 01 Jan 2013 16:00:00 PST', u'externalStatus': u'', u'createdConf': None, u'createdTime': u'Fri, 25 Jan 2013 10:56:27 PST', u'externalId': u'0000038-130124125317829-oozie-oozi-W', u'lastModifiedTime': u'Fri, 25 Jan 2013 10:56:41 PST', u'type': None, u'id': u'0000000-130117135211239-oozie-oozi-C@2', u'trackerUri': None},
  73. {u'status': u'SUCCEEDED', u'toString': u'WorkflowAction name[0000000-130117135211239-oozie-oozi-C@3] status[SUCCEEDED]', u'runConf': None, u'errorMessage': None, u'missingDependencies': u'', u'coordJobId': u'0000000-130117135211239-oozie-oozi-C', u'errorCode': None, u'actionNumber': 3, u'consoleUrl': None, u'nominalTime': u'Wed, 02 Jan 2013 16:00:00 PST', u'externalStatus': u'', u'createdConf': None, u'createdTime': u'Fri, 25 Jan 2013 08:59:38 PST', u'externalId': u'0000026-130124125317829-oozie-oozi-W', u'lastModifiedTime': u'Fri, 25 Jan 2013 09:00:05 PST', u'type': None, u'id': u'0000000-130117135211239-oozie-oozi-C@3', u'trackerUri': None},
  74. {u'status': u'SUCCEEDED', u'toString': u'WorkflowAction name[0000000-130117135211239-oozie-oozi-C@4] status[SUCCEEDED]', u'runConf': None, u'errorMessage': None, u'missingDependencies': u'', u'coordJobId': u'0000000-130117135211239-oozie-oozi-C', u'errorCode': None, u'actionNumber': 4, u'consoleUrl': None, u'nominalTime': u'Thu, 03 Jan 2013 16:00:00 PST', u'externalStatus': u'', u'createdConf': None, u'createdTime': u'Fri, 25 Jan 2013 10:53:39 PST', u'externalId': u'0000037-130124125317829-oozie-oozi-W', u'lastModifiedTime': u'Fri, 25 Jan 2013 10:54:17 PST', u'type': None, u'id': u'0000000-130117135211239-oozie-oozi-C@4', u'trackerUri': None}
  75. ]
  76. BUNDLE_ACTION = {u'startTime': u'Mon, 31 Dec 2012 16:00:00 PST', u'actions': [], u'frequency': 1, u'concurrency': 1, u'pauseTime': None, u'group': None, u'toString': u'Coordinator application id[0000022-130210132208494-oozie-oozi-C] status[SUCCEEDED]', u'consoleUrl': None, u'mat_throttling': 12, u'status': u'SUCCEEDED', u'conf': u'<configuration>\r\n <property>\r\n <name>oozie.coord.application.path</name>\r\n <value>hdfs://localhost:8020/user/hue/oozie/deployments/_romain_-oozie-6-1360649203.56</value>\r\n </property>\r\n <property>\r\n <name>oozie.bundle.application.path</name>\r\n <value>hdfs://localhost:8020/user/hue/oozie/workspaces/_romain_-oozie-22-1360636939.69</value>\r\n </property>\r\n <property>\r\n <name>market</name>\r\n <value>France</value>\r\n </property>\r\n <property>\r\n <name>user.name</name>\r\n <value>romain</value>\r\n </property>\r\n <property>\r\n <name>oozie.use.system.libpath</name>\r\n <value>true</value>\r\n </property>\r\n <property>\r\n <name>oozie.bundle.id</name>\r\n <value>0000021-130210132208494-oozie-oozi-B</value>\r\n </property>\r\n <property>\r\n <name>nameNode</name>\r\n <value>hdfs://localhost:8020</value>\r\n </property>\r\n <property>\r\n <name>wf_application_path</name>\r\n <value>hdfs://localhost:8020/user/hue/oozie/deployments/_romain_-oozie-5-1360649203.07</value>\r\n </property>\r\n <property>\r\n <name>jobTracker</name>\r\n <value>localhost:8021</value>\r\n </property>\r\n <property>\r\n <name>hue-id-b</name>\r\n <value>22</value>\r\n </property>\r\n</configuration>', u'user': u'romain', u'timeOut': 120, u'coordJobPath': u'hdfs://localhost:8020/user/hue/oozie/deployments/_romain_-oozie-6-1360649203.56', u'timeUnit': u'DAY', u'coordJobId': u'0000022-130210132208494-oozie-oozi-C', u'coordJobName': u'DailySleep', u'nextMaterializedTime': u'Fri, 04 Jan 2013 16:00:00 PST', u'coordExternalId': None, u'acl': None, u'lastAction': u'Fri, 04 Jan 2013 16:00:00 PST', u'executionPolicy': u'FIFO', u'timeZone': u'America/Los_Angeles', u'endTime': u'Fri, 04 Jan 2013 16:00:00 PST'}
  77. WORKFLOWS_SLAS = [
  78. {u'actualDuration': 68406, u'appType': u'WORKFLOW_JOB', u'appName': u'Forks', u'actualStart': u'Fri, 06 Dec 2013 14:01:53 PST', u'jobStatus': u'SUCCEEDED', u'id': u'0000002-131206135002457-oozie-oozi-W', u'expectedDuration': 1800000, u'nominalTime': u'Mon, 17 Jun 2013 17:01:00 PDT', u'slaStatus': u'MISS', u'lastModified': u'Fri, 06 Dec 2013 14:03:05 PST', u'actualEnd': u'Fri, 06 Dec 2013 14:03:01 PST', u'expectedEnd': u'Mon, 17 Jun 2013 17:31:00 PDT', u'expectedStart': u'Mon, 17 Jun 2013 17:11:00 PDT', u'user': u'romain'}
  79. ]
  80. def __init__(self, *args, **kwargs):
  81. self.api_version = 'v2'
  82. def setuser(self, user):
  83. pass
  84. @property
  85. def security_enabled(self):
  86. return False
  87. def submit_job(self, properties):
  88. return 'ONE-OOZIE-ID-W'
  89. def get_workflows(self, **kwargs):
  90. workflows = MockOozieApi.JSON_WORKFLOW_LIST
  91. if 'user' in kwargs:
  92. workflows = filter(lambda wf: wf['user'] == kwargs['user'], workflows)
  93. return WorkflowList(self, {'offset': 0, 'total': 4, 'workflows': workflows})
  94. def get_coordinators(self, **kwargs):
  95. coordinatorjobs = MockOozieApi.JSON_COORDINATOR_LIST
  96. if 'user' in kwargs:
  97. coordinatorjobs = filter(lambda coord: coord['user'] == kwargs['user'], coordinatorjobs)
  98. return CoordinatorList(self, {'offset': 0, 'total': 5, 'coordinatorjobs': coordinatorjobs})
  99. def get_bundles(self, **kwargs):
  100. bundlejobs = MockOozieApi.JSON_BUNDLE_LIST
  101. if 'user' in kwargs:
  102. bundlejobs = filter(lambda coord: coord['user'] == kwargs['user'], bundlejobs)
  103. return BundleList(self, {'offset': 0, 'total': 4, 'bundlejobs': bundlejobs})
  104. def get_job(self, job_id):
  105. if job_id in MockOozieApi.WORKFLOW_DICT:
  106. return OozieWorkflow(self, MockOozieApi.WORKFLOW_DICT[job_id])
  107. else:
  108. return OozieWorkflow(self, {'id': job_id, 'actions': []})
  109. def get_coordinator(self, job_id):
  110. if job_id in MockOozieApi.COORDINATOR_DICT:
  111. return OozieCoordinator(self, MockOozieApi.COORDINATOR_DICT[job_id])
  112. else:
  113. return OozieCoordinator(self, {'id': job_id, 'actions': []})
  114. def get_bundle(self, job_id):
  115. return OozieBundle(self, MockOozieApi.JSON_BUNDLE_LIST[0])
  116. def get_action(self, action_id):
  117. return WorkflowAction(MockOozieApi.WORKFLOW_ACTION)
  118. def job_control(self, job_id, action):
  119. return 'Done'
  120. def get_job_definition(self, jobid):
  121. if jobid == MockOozieApi.WORKFLOW_IDS[0]:
  122. return """<workflow-app name="MapReduce" xmlns="uri:oozie:workflow:0.4">
  123. <start to="Sleep"/>
  124. <action name="Sleep">
  125. <map-reduce>
  126. <job-tracker>${jobTracker}</job-tracker>
  127. <name-node>${nameNode}</name-node>
  128. <configuration>
  129. <property>
  130. <name>mapred.reduce.tasks</name>
  131. <value>1</value>
  132. </property>
  133. <property>
  134. <name>sleep.job.reduce.sleep.time</name>
  135. <value>${REDUCER_SLEEP_TIME}</value>
  136. </property>
  137. </configuration>
  138. </map-reduce>
  139. <ok to="end"/>
  140. <error to="kill"/>
  141. </action>
  142. <kill name="kill">
  143. <message>Action failed, error message[${wf:errorMessage(wf:lastErrorNode())}]</message>
  144. </kill>
  145. <end name="end"/>
  146. </workflow-app>"""
  147. else:
  148. return """<workflow-app name="MapReduce" xmlns="uri:oozie:workflow:0.4">BAD</workflow-app>"""
  149. def get_job_log(self, jobid):
  150. return '2013-01-08 16:28:06,487 INFO ActionStartXCommand:539 - USER[romain] GROUP[-] TOKEN[] APP[MapReduce] JOB[0000002-130108101138395-oozie-oozi-W] ACTION[0000002-130108101138395-oozie-oozi-W@:start:] Start action [0000002-130108101138395-oozie-oozi-W@:start:] with user-retry state : userRetryCount [0], userRetryMax [0], userRetryInterval [10]'
  151. def get_oozie_slas(self, **kwargs):
  152. return MockOozieApi.WORKFLOWS_SLAS
  153. class OozieMockBase(object):
  154. def setUp(self):
  155. # Beware: Monkey patch Oozie/LibOozie with Mock API
  156. if not hasattr(oozie_api, 'OriginalOozieApi'):
  157. oozie_api.OriginalOozieApi = oozie_api.OozieApi
  158. if not hasattr(Workflow.objects, 'original_check_workspace'):
  159. Workflow.objects.original_check_workspace = Workflow.objects.check_workspace
  160. Workflow.objects.check_workspace = lambda a, b: None
  161. oozie_api.OozieApi = MockOozieApi
  162. oozie_api._api_cache = None
  163. self.c = make_logged_in_client(is_superuser=False)
  164. grant_access("test", "test", "oozie")
  165. add_to_group("test")
  166. self.user = User.objects.get(username='test')
  167. self.wf = create_workflow(self.c, self.user)
  168. def tearDown(self):
  169. oozie_api.OozieApi = oozie_api.OriginalOozieApi
  170. Workflow.objects.check_workspace = Workflow.objects.original_check_workspace
  171. oozie_api._api_cache = None
  172. History.objects.all().delete()
  173. for coordinator in Coordinator.objects.all():
  174. coordinator.delete(skip_trash=True)
  175. for bundle in Bundle.objects.all():
  176. bundle.delete(skip_trash=True)
  177. def setup_simple_workflow(self):
  178. """ Creates a linear workflow """
  179. Link.objects.filter(parent__workflow=self.wf).delete()
  180. Link(parent=self.wf.start, child=self.wf.end, name="related").save()
  181. action1 = add_node(self.wf, 'action-name-1', 'mapreduce', [self.wf.start], {
  182. 'description': '',
  183. 'files': '[]',
  184. 'jar_path': '/user/hue/oozie/examples/lib/hadoop-examples.jar',
  185. 'job_properties': '[{"name":"sleep","value":"${SLEEP}"}]',
  186. 'prepares': '[{"value":"${output}","type":"delete"},{"value":"/test","type":"mkdir"}]',
  187. 'archives': '[]',
  188. })
  189. action2 = add_node(self.wf, 'action-name-2', 'mapreduce', [action1], {
  190. 'description': '',
  191. 'files': '[]',
  192. 'jar_path': '/user/hue/oozie/examples/lib/hadoop-examples.jar',
  193. 'job_properties': '[{"name":"sleep","value":"${SLEEP}"}]',
  194. 'prepares': '[{"value":"${output}","type":"delete"},{"value":"/test","type":"mkdir"}]',
  195. 'archives': '[]',
  196. })
  197. action3 = add_node(self.wf, 'action-name-3', 'mapreduce', [action2], {
  198. 'description': '',
  199. 'files': '[]',
  200. 'jar_path': '/user/hue/oozie/examples/lib/hadoop-examples.jar',
  201. 'job_properties': '[{"name":"sleep","value":"${SLEEP}"}]',
  202. 'prepares': '[{"value":"${output}","type":"delete"},{"value":"/test","type":"mkdir"}]',
  203. 'archives': '[]',
  204. })
  205. Link(parent=action3, child=self.wf.end, name="ok").save()
  206. def setup_forking_workflow(self):
  207. Link.objects.filter(parent__workflow=self.wf).delete()
  208. Link(parent=self.wf.start, child=self.wf.end, name="related").save()
  209. fork1 = add_node(self.wf, 'fork-name-1', 'fork', [self.wf.start])
  210. action1 = add_node(self.wf, 'action-name-1', 'mapreduce', [fork1])
  211. action2 = add_node(self.wf, 'action-name-2', 'mapreduce', [fork1])
  212. join1 = add_node(self.wf, 'join-name-1', 'join', [action1, action2])
  213. Link(parent=fork1, child=join1, name="related").save()
  214. action3 = add_node(self.wf, 'action-name-3', 'mapreduce', [join1])
  215. Link(parent=action3, child=self.wf.end, name="ok").save()
  216. def create_noop_workflow(self, name='noop-test'):
  217. Node.objects.filter(workflow__name=name).delete()
  218. Workflow.objects.filter(name=name).delete()
  219. if Document.objects.get_docs(self.user, Workflow).filter(name=name).exists():
  220. for doc in Document.objects.get_docs(self.user, Workflow).filter(name=name):
  221. if doc.content_object:
  222. self.c.post(reverse('oozie:delete_workflow') + '?skip_trash=true', {'job_selection': [doc.content_object.id]}, follow=True)
  223. else:
  224. doc.delete()
  225. wf = Workflow.objects.new_workflow(self.user)
  226. wf.name = name
  227. wf.save()
  228. wf.start.workflow = wf
  229. wf.end.workflow = wf
  230. wf.start.save()
  231. wf.end.save()
  232. Document.objects.link(wf, owner=wf.owner, name=wf.name, description=wf.description)
  233. Kill.objects.create(name='kill', workflow=wf, node_type=Kill.node_type)
  234. Link.objects.create(parent=wf.start, child=wf.end, name='related')
  235. Link.objects.create(parent=wf.start, child=wf.end, name="to")
  236. return wf
  237. class OozieBase(OozieServerProvider):
  238. requires_hadoop = True
  239. def setUp(self):
  240. OozieServerProvider.setup_class()
  241. self.c = make_logged_in_client(is_superuser=False)
  242. self.user = User.objects.get(username="test")
  243. grant_access("test", "test", "oozie")
  244. add_to_group("test")
  245. self.cluster = OozieServerProvider.cluster
  246. self.install_examples()
  247. def install_examples(self):
  248. global _INITIALIZED
  249. if _INITIALIZED:
  250. return
  251. self.c.post(reverse('oozie:install_examples'))
  252. self.cluster.fs.do_as_user('test', self.cluster.fs.create_home_dir, '/user/test')
  253. _INITIALIZED = True
  254. def setup_simple_workflow(self):
  255. """ Creates a linear workflow """
  256. Link.objects.filter(parent__workflow=self.wf).delete()
  257. Link(parent=self.wf.start, child=self.wf.end, name="related").save()
  258. action1 = add_node(self.wf, 'action-name-1', 'mapreduce', [self.wf.start], {
  259. 'description': '',
  260. 'files': '[]',
  261. 'jar_path': '/user/hue/oozie/examples/lib/hadoop-examples.jar',
  262. 'job_properties': '[{"name":"sleep","value":"${SLEEP}"}]',
  263. 'prepares': '[{"value":"${output}","type":"delete"},{"value":"/test","type":"mkdir"}]',
  264. 'archives': '[]',
  265. })
  266. action2 = add_node(self.wf, 'action-name-2', 'mapreduce', [action1], {
  267. 'description': '',
  268. 'files': '[]',
  269. 'jar_path': '/user/hue/oozie/examples/lib/hadoop-examples.jar',
  270. 'job_properties': '[{"name":"sleep","value":"${SLEEP}"}]',
  271. 'prepares': '[{"value":"${output}","type":"delete"},{"value":"/test","type":"mkdir"}]',
  272. 'archives': '[]',
  273. })
  274. action3 = add_node(self.wf, 'action-name-3', 'mapreduce', [action2], {
  275. 'description': '',
  276. 'files': '[]',
  277. 'jar_path': '/user/hue/oozie/examples/lib/hadoop-examples.jar',
  278. 'job_properties': '[{"name":"sleep","value":"${SLEEP}"}]',
  279. 'prepares': '[{"value":"${output}","type":"delete"},{"value":"/test","type":"mkdir"}]',
  280. 'archives': '[]',
  281. })
  282. Link(parent=action3, child=self.wf.end, name="ok").save()
  283. def setup_forking_workflow(self):
  284. """ Creates a workflow with a fork """
  285. Link.objects.filter(parent__workflow=self.wf).delete()
  286. Link(parent=self.wf.start, child=self.wf.end, name="related").save()
  287. fork1 = add_node(self.wf, 'fork-name-1', 'fork', [self.wf.start])
  288. action1 = add_node(self.wf, 'action-name-1', 'mapreduce', [fork1])
  289. action2 = add_node(self.wf, 'action-name-2', 'mapreduce', [fork1])
  290. join1 = add_node(self.wf, 'join-name-1', 'join', [action1, action2])
  291. Link(parent=fork1, child=join1, name="related").save()
  292. action3 = add_node(self.wf, 'action-name-3', 'mapreduce', [join1])
  293. Link(parent=action3, child=self.wf.end, name="ok").save()
  294. class TestAPI(OozieMockBase):
  295. def setUp(self):
  296. OozieMockBase.setUp(self)
  297. self.wf = Workflow.objects.get(name='wf-name-1', managed=True)
  298. def test_workflow_save(self):
  299. self.setup_simple_workflow()
  300. workflow_dict = workflow_to_dict(self.wf)
  301. workflow_json = json.dumps(workflow_dict)
  302. response = self.c.post(reverse('oozie:workflow_save', kwargs={'workflow': self.wf.pk}), data={'workflow': workflow_json})
  303. test_response_json = response.content
  304. test_response_json_object = json.loads(test_response_json)
  305. assert_equal(0, test_response_json_object['status'])
  306. # Change property and save
  307. workflow_dict = workflow_to_dict(self.wf)
  308. workflow_dict['description'] = 'test'
  309. workflow_json = json.dumps(workflow_dict)
  310. response = self.c.post(reverse('oozie:workflow_save', kwargs={'workflow': self.wf.pk}), data={'workflow': workflow_json})
  311. test_response_json = response.content
  312. test_response_json_object = json.loads(test_response_json)
  313. assert_equal(0, test_response_json_object['status'])
  314. wf = Workflow.objects.get(id=self.wf.id)
  315. assert_equal('test', wf.description)
  316. assert_equal(self.wf.name, wf.name)
  317. # Change node and save
  318. workflow_dict = workflow_to_dict(self.wf)
  319. workflow_dict['nodes'][2]['name'] = 'new-name'
  320. node_id = workflow_dict['nodes'][2]['id']
  321. workflow_json = json.dumps(workflow_dict)
  322. response = self.c.post(reverse('oozie:workflow_save', kwargs={'workflow': self.wf.pk}), data={'workflow': workflow_json})
  323. test_response_json = response.content
  324. test_response_json_object = json.loads(test_response_json)
  325. assert_equal(0, test_response_json_object['status'])
  326. node = Node.objects.get(id=node_id)
  327. assert_equal('new-name', node.name)
  328. def test_workflow_save_fail(self):
  329. self.setup_simple_workflow()
  330. # Bad workflow name
  331. workflow_dict = workflow_to_dict(self.wf)
  332. del workflow_dict['name']
  333. workflow_json = json.dumps(workflow_dict)
  334. response = self.c.post(reverse('oozie:workflow_save', kwargs={'workflow': self.wf.pk}), data={'workflow': workflow_json}, HTTP_X_REQUESTED_WITH='XMLHttpRequest')
  335. assert_equal(400, response.status_code)
  336. # Bad node name
  337. workflow_dict = workflow_to_dict(self.wf)
  338. del workflow_dict['nodes'][2]['name']
  339. workflow_json = json.dumps(workflow_dict)
  340. response = self.c.post(reverse('oozie:workflow_save', kwargs={'workflow': self.wf.pk}), data={'workflow': workflow_json}, HTTP_X_REQUESTED_WITH='XMLHttpRequest')
  341. assert_equal(400, response.status_code)
  342. # Bad control node name should still go through
  343. workflow_dict = workflow_to_dict(self.wf)
  344. del workflow_dict['nodes'][0]['name']
  345. workflow_json = json.dumps(workflow_dict)
  346. response = self.c.post(reverse('oozie:workflow_save', kwargs={'workflow': self.wf.pk}), data={'workflow': workflow_json}, HTTP_X_REQUESTED_WITH='XMLHttpRequest')
  347. assert_equal(200, response.status_code)
  348. def test_workflow_add_subworkflow_node(self):
  349. wf = self.create_noop_workflow()
  350. try:
  351. subworkflow_name = "subworkflow-1"
  352. subworkflow_id = "subworkflow:1"
  353. subworkflow_json = """{
  354. "description": "",
  355. "workflow": %(workflow)d,
  356. "child_links": [
  357. {
  358. "comment": "",
  359. "name": "ok",
  360. "parent": "%(id)s",
  361. "child": %(end)d
  362. },
  363. {
  364. "comment": "",
  365. "name": "error",
  366. "parent": "%(id)s",
  367. "child": %(kill)d
  368. }
  369. ],
  370. "node_type": "subworkflow",
  371. "sub_workflow": %(subworkflow)d,
  372. "job_properties": "[]",
  373. "name": "%(name)s",
  374. "id": "%(id)s",
  375. "propagate_configuration": true
  376. }"""
  377. subworkflow_json = subworkflow_json % {
  378. 'workflow': wf.id,
  379. 'subworkflow': self.wf.id,
  380. 'end': wf.end.id,
  381. 'kill': Kill.objects.get(workflow=wf).id,
  382. 'name': subworkflow_name,
  383. 'id': subworkflow_id
  384. }
  385. workflow_dict = workflow_to_dict(wf)
  386. workflow_dict['nodes'].append(json.loads(subworkflow_json))
  387. workflow_dict['nodes'][0]['child_links'][1]['child'] = subworkflow_id
  388. del workflow_dict['nodes'][0]['child_links'][1]['id']
  389. workflow_json = json.dumps(workflow_dict)
  390. response = self.c.post(reverse('oozie:workflow_save', kwargs={'workflow': wf.pk}), data={'workflow': workflow_json})
  391. test_response_json = response.content
  392. test_response_json_object = json.loads(test_response_json)
  393. assert_equal(0, test_response_json_object['status'], workflow_json)
  394. finally:
  395. wf.delete(skip_trash=True)
  396. def test_workflow_add_mapreduce_node(self):
  397. wf = self.create_noop_workflow()
  398. try:
  399. node_name = "mr-1"
  400. node_id = "mapreduce:1"
  401. node_json = """{
  402. "description": "",
  403. "workflow": %(workflow)d,
  404. "child_links": [
  405. {
  406. "comment": "",
  407. "name": "ok",
  408. "parent": "%(id)s",
  409. "child": %(end)d
  410. },
  411. {
  412. "comment": "",
  413. "name": "error",
  414. "parent": "%(id)s",
  415. "child": %(kill)d
  416. }
  417. ],
  418. "node_type": "mapreduce",
  419. "jar_path": "test",
  420. "job_properties": "[]",
  421. "files": "[]",
  422. "archives": "[]",
  423. "prepares": "[]",
  424. "name": "%(name)s",
  425. "id": "%(id)s"
  426. }"""
  427. node_json = node_json % {
  428. 'workflow': wf.id,
  429. 'end': wf.end.id,
  430. 'kill': Kill.objects.get(workflow=wf).id,
  431. 'name': node_name,
  432. 'id': node_id
  433. }
  434. workflow_dict = workflow_to_dict(wf)
  435. workflow_dict['nodes'].append(json.loads(node_json))
  436. workflow_dict['nodes'][0]['child_links'][1]['child'] = node_id
  437. del workflow_dict['nodes'][0]['child_links'][1]['id']
  438. workflow_json = json.dumps(workflow_dict)
  439. response = self.c.post(reverse('oozie:workflow_save', kwargs={'workflow': wf.pk}), data={'workflow': workflow_json})
  440. test_response_json = response.content
  441. test_response_json_object = json.loads(test_response_json)
  442. assert_equal(0, test_response_json_object['status'], workflow_json)
  443. finally:
  444. wf.delete(skip_trash=True)
  445. def test_workflow(self):
  446. response = self.c.get(reverse('oozie:workflow', kwargs={'workflow': self.wf.pk}))
  447. test_response_json = response.content
  448. test_response_json_object = json.loads(test_response_json)
  449. assert_equal(0, test_response_json_object['status'])
  450. def test_workflow_validate_node(self):
  451. data = {"files":"[\"hive-site.xml\"]","job_xml":"hive-site.xml","description":"Show databases","workflow":17,"child_links":[{"comment":"","name":"ok","id":106,"parent":76,"child":74},{"comment":"","name":"error","id":107,"parent":76,"child":73}],"job_properties":"[{\"name\":\"oozie.hive.defaults\",\"value\":\"hive-site.xml\"}]","node_type":"hive","params":"[{\"value\":\"INPUT=/user/hue/oozie/workspaces/data\",\"type\":\"param\"}]","archives":"[]","node_ptr":76,"prepares":"[]","script_path":"hive.sql","id":76,"name":"Hive"}
  452. response = self.c.post(reverse('oozie:workflow_validate_node', kwargs={'workflow': self.wf.pk, 'node_type': 'hive'}), data={'node': json.dumps(data)}, HTTP_X_REQUESTED_WITH='XMLHttpRequest')
  453. test_response_json = response.content
  454. test_response_json_object = json.loads(test_response_json)
  455. assert_equal(0, test_response_json_object['status'])
  456. assert_true('name' in test_response_json_object['data'], test_response_json_object['data'])
  457. assert_equal(0, len(test_response_json_object['data']['name']), test_response_json_object['data'])
  458. assert_true('description' in test_response_json_object['data'], test_response_json_object['data'])
  459. assert_equal(0, len(test_response_json_object['data']['description']), test_response_json_object['data'])
  460. assert_true('script_path' in test_response_json_object['data'], test_response_json_object['data'])
  461. assert_equal(0, len(test_response_json_object['data']['script_path']), test_response_json_object['data'])
  462. assert_true('job_xml' in test_response_json_object['data'], test_response_json_object['data'])
  463. assert_equal(0, len(test_response_json_object['data']['job_xml']), test_response_json_object['data'])
  464. assert_true('job_properties' in test_response_json_object['data'], test_response_json_object['data'])
  465. assert_equal(0, len(test_response_json_object['data']['job_properties']), test_response_json_object['data'])
  466. assert_true('files' in test_response_json_object['data'], test_response_json_object['data'])
  467. assert_equal(0, len(test_response_json_object['data']['files']), test_response_json_object['data'])
  468. assert_true('params' in test_response_json_object['data'], test_response_json_object['data'])
  469. assert_equal(0, len(test_response_json_object['data']['params']), test_response_json_object['data'])
  470. assert_true('prepares' in test_response_json_object['data'], test_response_json_object['data'])
  471. assert_equal(0, len(test_response_json_object['data']['prepares']), test_response_json_object['data'])
  472. assert_true('archives' in test_response_json_object['data'], test_response_json_object['data'])
  473. assert_equal(0, len(test_response_json_object['data']['archives']), test_response_json_object['data'])
  474. def test_workflow_validate_node_fail(self):
  475. # Empty files field
  476. data = {"job_xml":"hive-site.xml","description":"Show databases","workflow":17,"child_links":[{"comment":"","name":"ok","id":106,"parent":76,"child":74},{"comment":"","name":"error","id":107,"parent":76,"child":73}],"job_properties":"[{\"name\":\"oozie.hive.defaults\",\"value\":\"hive-site.xml\"}]","node_type":"hive","params":"[{\"value\":\"INPUT=/user/hue/oozie/workspaces/data\",\"type\":\"param\"}]","archives":"[]","node_ptr":76,"prepares":"[]","script_path":"hive.sql","id":76,"name":"Hive"}
  477. response = self.c.post(reverse('oozie:workflow_validate_node', kwargs={'workflow': self.wf.pk, 'node_type': 'hive'}), data={'node': json.dumps(data)}, HTTP_X_REQUESTED_WITH='XMLHttpRequest')
  478. test_response_json = response.content
  479. test_response_json_object = json.loads(test_response_json)
  480. assert_equal(-1, test_response_json_object['status'])
  481. assert_true('name' in test_response_json_object['data'], test_response_json_object['data'])
  482. assert_equal(0, len(test_response_json_object['data']['name']), test_response_json_object['data'])
  483. assert_true('description' in test_response_json_object['data'], test_response_json_object['data'])
  484. assert_equal(0, len(test_response_json_object['data']['description']), test_response_json_object['data'])
  485. assert_true('script_path' in test_response_json_object['data'], test_response_json_object['data'])
  486. assert_equal(0, len(test_response_json_object['data']['script_path']), test_response_json_object['data'])
  487. assert_true('job_xml' in test_response_json_object['data'], test_response_json_object['data'])
  488. assert_equal(0, len(test_response_json_object['data']['job_xml']), test_response_json_object['data'])
  489. assert_true('job_properties' in test_response_json_object['data'], test_response_json_object['data'])
  490. assert_equal(0, len(test_response_json_object['data']['job_properties']), test_response_json_object['data'])
  491. assert_true('files' in test_response_json_object['data'], test_response_json_object['data'])
  492. assert_equal(1, len(test_response_json_object['data']['files']), test_response_json_object['data'])
  493. assert_true('params' in test_response_json_object['data'], test_response_json_object['data'])
  494. assert_equal(0, len(test_response_json_object['data']['params']), test_response_json_object['data'])
  495. assert_true('prepares' in test_response_json_object['data'], test_response_json_object['data'])
  496. assert_equal(0, len(test_response_json_object['data']['prepares']), test_response_json_object['data'])
  497. assert_true('archives' in test_response_json_object['data'], test_response_json_object['data'])
  498. assert_equal(0, len(test_response_json_object['data']['archives']), test_response_json_object['data'])
  499. # Empty script path
  500. data = {"files":"[\"hive-site.xml\"]","job_xml":"hive-site.xml","description":"Show databases","workflow":17,"child_links":[{"comment":"","name":"ok","id":106,"parent":76,"child":74},{"comment":"","name":"error","id":107,"parent":76,"child":73}],"job_properties":"[{\"name\":\"oozie.hive.defaults\",\"value\":\"hive-site.xml\"}]","node_type":"hive","params":"[{\"value\":\"INPUT=/user/hue/oozie/workspaces/data\",\"type\":\"param\"}]","archives":"[]","node_ptr":76,"prepares":"[]","script_path":"","id":76,"name":"Hive"}
  501. response = self.c.post(reverse('oozie:workflow_validate_node', kwargs={'workflow': self.wf.pk, 'node_type': 'hive'}), data={'node': json.dumps(data)}, HTTP_X_REQUESTED_WITH='XMLHttpRequest')
  502. test_response_json = response.content
  503. test_response_json_object = json.loads(test_response_json)
  504. assert_equal(-1, test_response_json_object['status'])
  505. assert_true('name' in test_response_json_object['data'], test_response_json_object['data'])
  506. assert_equal(0, len(test_response_json_object['data']['name']), test_response_json_object['data'])
  507. assert_true('description' in test_response_json_object['data'], test_response_json_object['data'])
  508. assert_equal(0, len(test_response_json_object['data']['description']), test_response_json_object['data'])
  509. assert_true('script_path' in test_response_json_object['data'], test_response_json_object['data'])
  510. assert_equal(1, len(test_response_json_object['data']['script_path']), test_response_json_object['data'])
  511. assert_true('job_xml' in test_response_json_object['data'], test_response_json_object['data'])
  512. assert_equal(0, len(test_response_json_object['data']['job_xml']), test_response_json_object['data'])
  513. assert_true('job_properties' in test_response_json_object['data'], test_response_json_object['data'])
  514. assert_equal(0, len(test_response_json_object['data']['job_properties']), test_response_json_object['data'])
  515. assert_true('files' in test_response_json_object['data'], test_response_json_object['data'])
  516. assert_equal(0, len(test_response_json_object['data']['files']), test_response_json_object['data'])
  517. assert_true('params' in test_response_json_object['data'], test_response_json_object['data'])
  518. assert_equal(0, len(test_response_json_object['data']['params']), test_response_json_object['data'])
  519. assert_true('prepares' in test_response_json_object['data'], test_response_json_object['data'])
  520. assert_equal(0, len(test_response_json_object['data']['prepares']), test_response_json_object['data'])
  521. assert_true('archives' in test_response_json_object['data'], test_response_json_object['data'])
  522. assert_equal(0, len(test_response_json_object['data']['archives']), test_response_json_object['data'])
  523. def test_workflows(self):
  524. response = self.c.get(reverse('oozie:workflows') + "?managed=true", HTTP_X_REQUESTED_WITH='XMLHttpRequest')
  525. response_json_dict = json.loads(response.content)
  526. assert_equal(0, response_json_dict['status'])
  527. assert_equal(1, len(response_json_dict['data']['workflows']))
  528. response = self.c.get(reverse('oozie:workflows') + "?managed=false", HTTP_X_REQUESTED_WITH='XMLHttpRequest')
  529. response_json_dict = json.loads(response.content)
  530. assert_equal(0, response_json_dict['status'])
  531. assert_equal(0, len(response_json_dict['data']['workflows']))
  532. def test_workflow_actions(self):
  533. response = self.c.get(reverse('oozie:workflow_actions', kwargs={'workflow': self.wf.pk}), HTTP_X_REQUESTED_WITH='XMLHttpRequest')
  534. response_json_dict = json.loads(response.content)
  535. assert_equal(0, response_json_dict['status'])
  536. assert_equal(0, len(response_json_dict['data']['actions']))
  537. self.setup_simple_workflow()
  538. response = self.c.get(reverse('oozie:workflow_actions', kwargs={'workflow': self.wf.pk}), HTTP_X_REQUESTED_WITH='XMLHttpRequest')
  539. response_json_dict = json.loads(response.content)
  540. assert_equal(0, response_json_dict['status'])
  541. assert_equal(3, len(response_json_dict['data']['actions']))
  542. def test_autocomplete(self):
  543. response = self.c.get(reverse('oozie:autocomplete_properties'))
  544. test_response_json = response.content
  545. assert_true('mapred.input.dir' in test_response_json)
  546. class TestApiPermissionsWithOozie(OozieBase):
  547. def setUp(self):
  548. OozieBase.setUp(self)
  549. # When updating wf, update wf_json as well!
  550. self.wf = Workflow.objects.get(name='MapReduce', managed=True).clone(self.cluster.fs, self.user)
  551. def test_workflow_save(self):
  552. # Share
  553. self.wf.is_shared = True
  554. self.wf.save()
  555. Workflow.objects.check_workspace(self.wf, self.cluster.fs)
  556. # Login as someone else
  557. client_not_me = make_logged_in_client(username='not_me', is_superuser=False, groupname='test')
  558. grant_access("not_me", "test", "oozie")
  559. workflow_dict = workflow_to_dict(self.wf)
  560. workflow_json = json.dumps(workflow_dict)
  561. response = client_not_me.post(reverse('oozie:workflow_save', kwargs={'workflow': self.wf.pk}), data={'workflow': workflow_json}, HTTP_X_REQUESTED_WITH='XMLHttpRequest')
  562. assert_equal(401, response.status_code, response.status_code)
  563. response = self.c.post(reverse('oozie:workflow_save', kwargs={'workflow': self.wf.pk}), data={'workflow': workflow_json}, HTTP_X_REQUESTED_WITH='XMLHttpRequest')
  564. test_response_json = response.content
  565. test_response_json_object = json.loads(test_response_json)
  566. assert_equal(200, response.status_code, response)
  567. assert_equal(0, test_response_json_object['status'])
  568. def test_workflow_save_fail(self):
  569. # Unshare
  570. self.wf.is_shared = False
  571. self.wf.save()
  572. Workflow.objects.check_workspace(self.wf, self.cluster.fs)
  573. # Login as someone else
  574. client_not_me = make_logged_in_client(username='not_me', is_superuser=False, groupname='test')
  575. grant_access("not_me", "test", "oozie")
  576. workflow_dict = workflow_to_dict(self.wf)
  577. del workflow_dict['name']
  578. workflow_json = json.dumps(workflow_dict)
  579. response = client_not_me.post(reverse('oozie:workflow_save', kwargs={'workflow': self.wf.pk}), data={'workflow': workflow_json}, HTTP_X_REQUESTED_WITH='XMLHttpRequest')
  580. assert_equal(401, response.status_code, response)
  581. def test_workflow(self):
  582. # Share
  583. self.wf.is_shared = True
  584. self.wf.doc.get().share_to_default()
  585. self.wf.save()
  586. Workflow.objects.check_workspace(self.wf, self.cluster.fs)
  587. # Login as someone else
  588. client_not_me = make_logged_in_client(username='not_me', is_superuser=False, groupname='default', recreate=True)
  589. grant_access("not_me", "test", "oozie")
  590. response = client_not_me.get(reverse('oozie:workflow', kwargs={'workflow': self.wf.pk}), HTTP_X_REQUESTED_WITH='XMLHttpRequest')
  591. test_response_json = response.content
  592. test_response_json_object = json.loads(test_response_json)
  593. assert_equal(200, response.status_code, response)
  594. assert_equal(0, test_response_json_object['status'])
  595. def test_workflow_fail(self):
  596. # Unshare
  597. self.wf.is_shared = False
  598. self.wf.save()
  599. Workflow.objects.check_workspace(self.wf, self.cluster.fs)
  600. # Login as someone else
  601. client_not_me = make_logged_in_client(username='not_me', is_superuser=False, groupname='test')
  602. grant_access("not_me", "test", "oozie")
  603. response = client_not_me.get(reverse('oozie:workflow', kwargs={'workflow': self.wf.pk}), HTTP_X_REQUESTED_WITH='XMLHttpRequest')
  604. assert_equal(401, response.status_code)
  605. class TestEditor(OozieMockBase):
  606. def setUp(self):
  607. super(TestEditor, self).setUp()
  608. self.setup_simple_workflow()
  609. def test_workflow_name(self):
  610. try:
  611. workflow_dict = WORKFLOW_DICT.copy()
  612. workflow_count = Document.objects.available_docs(Workflow, self.user).count()
  613. workflow_dict['name'][0] = 'bad workflow name'
  614. response = self.c.post(reverse('oozie:create_workflow'), workflow_dict, follow=True)
  615. assert_equal(200, response.status_code)
  616. assert_equal(workflow_count, Document.objects.available_docs(Workflow, self.user).count(), response)
  617. workflow_dict['name'][0] = 'good-workflow-name'
  618. response = self.c.post(reverse('oozie:create_workflow'), workflow_dict, follow=True)
  619. assert_equal(200, response.status_code)
  620. assert_equal(workflow_count + 1, Document.objects.available_docs(Workflow, self.user).count(), response)
  621. finally:
  622. name = 'bad workflow name'
  623. if Workflow.objects.filter(name=name).exists():
  624. Node.objects.filter(workflow__name=name).delete()
  625. Workflow.objects.filter(name=name).delete()
  626. name = 'good-workflow-name'
  627. if Workflow.objects.filter(name=name).exists():
  628. Node.objects.filter(workflow__name=name).delete()
  629. Workflow.objects.filter(name=name).delete()
  630. def test_find_parameters(self):
  631. jobs = [Job(name="$a"),
  632. Job(name="foo ${b} $$"),
  633. Job(name="${foo}", description="xxx ${foo}")]
  634. result = [find_parameters(job, ['name', 'description']) for job in jobs]
  635. assert_equal(set(["b", "foo"]), reduce(lambda x, y: x | set(y), result, set()))
  636. def test_find_all_parameters(self):
  637. assert_equal([{'name': u'output', 'value': u''}, {'name': u'SLEEP', 'value': ''}, {'name': u'market', 'value': u'US'}],
  638. self.wf.find_all_parameters())
  639. def test_workflow_has_cycle(self):
  640. action1 = Node.objects.get(workflow=self.wf, name='action-name-1')
  641. action3 = Node.objects.get(workflow=self.wf, name='action-name-3')
  642. assert_false(self.wf.has_cycle())
  643. ok = action3.get_link('ok')
  644. ok.child = action1
  645. ok.save()
  646. assert_true(self.wf.has_cycle())
  647. def test_workflow_gen_xml(self):
  648. assert_equal(
  649. '<workflow-app name="wf-name-1" xmlns="uri:oozie:workflow:0.4">\n'
  650. ' <global>\n'
  651. ' <job-xml>jobconf.xml</job-xml>\n'
  652. ' <configuration>\n'
  653. ' <property>\n'
  654. ' <name>sleep-all</name>\n'
  655. ' <value>${SLEEP}</value>\n'
  656. ' </property>\n'
  657. ' </configuration>\n'
  658. ' </global>\n'
  659. ' <start to="action-name-1"/>\n'
  660. ' <action name="action-name-1">\n'
  661. ' <map-reduce>\n'
  662. ' <job-tracker>${jobTracker}</job-tracker>\n'
  663. ' <name-node>${nameNode}</name-node>\n'
  664. ' <prepare>\n'
  665. ' <delete path="${nameNode}${output}"/>\n'
  666. ' <mkdir path="${nameNode}/test"/>\n'
  667. ' </prepare>\n'
  668. ' <configuration>\n'
  669. ' <property>\n'
  670. ' <name>sleep</name>\n'
  671. ' <value>${SLEEP}</value>\n'
  672. ' </property>\n'
  673. ' </configuration>\n'
  674. ' </map-reduce>\n'
  675. ' <ok to="action-name-2"/>\n'
  676. ' <error to="kill"/>\n'
  677. ' </action>\n'
  678. ' <action name="action-name-2">\n'
  679. ' <map-reduce>\n'
  680. ' <job-tracker>${jobTracker}</job-tracker>\n'
  681. ' <name-node>${nameNode}</name-node>\n'
  682. ' <prepare>\n'
  683. ' <delete path="${nameNode}${output}"/>\n'
  684. ' <mkdir path="${nameNode}/test"/>\n'
  685. ' </prepare>\n'
  686. ' <configuration>\n'
  687. ' <property>\n'
  688. ' <name>sleep</name>\n'
  689. ' <value>${SLEEP}</value>\n'
  690. ' </property>\n'
  691. ' </configuration>\n'
  692. ' </map-reduce>\n'
  693. ' <ok to="action-name-3"/>\n'
  694. ' <error to="kill"/>\n'
  695. ' </action>\n'
  696. ' <action name="action-name-3">\n'
  697. ' <map-reduce>\n'
  698. ' <job-tracker>${jobTracker}</job-tracker>\n'
  699. ' <name-node>${nameNode}</name-node>\n'
  700. ' <prepare>\n'
  701. ' <delete path="${nameNode}${output}"/>\n'
  702. ' <mkdir path="${nameNode}/test"/>\n'
  703. ' </prepare>\n'
  704. ' <configuration>\n'
  705. ' <property>\n'
  706. ' <name>sleep</name>\n'
  707. ' <value>${SLEEP}</value>\n'
  708. ' </property>\n'
  709. ' </configuration>\n'
  710. ' </map-reduce>\n'
  711. ' <ok to="end"/>\n'
  712. ' <error to="kill"/>\n'
  713. ' </action>\n'
  714. ' <kill name="kill">\n'
  715. ' <message>Action failed, error message[${wf:errorMessage(wf:lastErrorNode())}]</message>\n'
  716. ' </kill>\n'
  717. ' <end name="end"/>\n'
  718. '</workflow-app>'.split(), self.wf.to_xml({'output': '/path'}).split())
  719. def test_workflow_java_gen_xml(self):
  720. self.wf.node_set.filter(name='action-name-1').delete()
  721. action1 = add_node(self.wf, 'action-name-1', 'java', [self.wf.start], {
  722. u'name': 'MyTeragen',
  723. "description":"Generate N number of records",
  724. "main_class":"org.apache.hadoop.examples.terasort.TeraGen",
  725. "args":"1000 ${output_dir}/teragen",
  726. "files":'["my_file","my_file2"]',
  727. "job_xml":"",
  728. "java_opts":"-Dexample-property=natty",
  729. "jar_path":"/user/hue/oozie/workspaces/lib/hadoop-examples.jar",
  730. "prepares":'[{"value":"/test","type":"mkdir"}]',
  731. "archives":'[{"dummy":"","name":"my_archive"},{"dummy":"","name":"my_archive2"}]',
  732. "capture_output": "on",
  733. })
  734. Link(parent=action1, child=self.wf.end, name="ok").save()
  735. xml = self.wf.to_xml({'output_dir': '/path'})
  736. assert_true("""
  737. <action name="MyTeragen">
  738. <java>
  739. <job-tracker>${jobTracker}</job-tracker>
  740. <name-node>${nameNode}</name-node>
  741. <prepare>
  742. <mkdir path="${nameNode}/test"/>
  743. </prepare>
  744. <main-class>org.apache.hadoop.examples.terasort.TeraGen</main-class>
  745. <java-opts>-Dexample-property=natty</java-opts>
  746. <arg>1000</arg>
  747. <arg>${output_dir}/teragen</arg>
  748. <file>my_file#my_file</file>
  749. <file>my_file2#my_file2</file>
  750. <archive>my_archive#my_archive</archive>
  751. <archive>my_archive2#my_archive2</archive>
  752. <capture-output/>
  753. </java>
  754. <ok to="end"/>
  755. <error to="kill"/>
  756. </action>""" in xml, xml)
  757. def test_workflow_streaming_gen_xml(self):
  758. self.wf.node_set.filter(name='action-name-1').delete()
  759. action1 = add_node(self.wf, 'action-name-1', 'streaming', [self.wf.start], {
  760. u'name': 'MyStreaming',
  761. "description": "Generate N number of records",
  762. "main_class": "org.apache.hadoop.examples.terasort.TeraGen",
  763. "mapper": "MyMapper",
  764. "reducer": "MyReducer",
  765. "files": '["my_file"]',
  766. "archives":'[{"dummy":"","name":"my_archive"}]',
  767. })
  768. Link(parent=action1, child=self.wf.end, name="ok").save()
  769. xml = self.wf.to_xml()
  770. assert_true("""
  771. <action name="MyStreaming">
  772. <map-reduce>
  773. <job-tracker>${jobTracker}</job-tracker>
  774. <name-node>${nameNode}</name-node>
  775. <streaming>
  776. <mapper>MyMapper</mapper>
  777. <reducer>MyReducer</reducer>
  778. </streaming>
  779. <file>my_file#my_file</file>
  780. <archive>my_archive#my_archive</archive>
  781. </map-reduce>
  782. <ok to="end"/>
  783. <error to="kill"/>
  784. </action>""" in xml, xml)
  785. def test_workflow_shell_gen_xml(self):
  786. self.wf.node_set.filter(name='action-name-1').delete()
  787. action1 = add_node(self.wf, 'action-name-1', 'shell', [self.wf.start], {
  788. u'job_xml': 'my-job.xml',
  789. u'files': '["hello.py"]',
  790. u'name': 'Shell',
  791. u'job_properties': '[]',
  792. u'capture_output': 'on',
  793. u'command': 'hello.py',
  794. u'archives': '[]',
  795. u'prepares': '[]',
  796. u'params': '[{"value":"World!","type":"argument"}]',
  797. u'description': 'Execute a Python script printing its arguments'
  798. })
  799. Link(parent=action1, child=self.wf.end, name="ok").save()
  800. xml = self.wf.to_xml()
  801. assert_true("""
  802. <shell xmlns="uri:oozie:shell-action:0.1">
  803. <job-tracker>${jobTracker}</job-tracker>
  804. <name-node>${nameNode}</name-node>
  805. <job-xml>my-job.xml</job-xml>
  806. <exec>hello.py</exec>
  807. <argument>World!</argument>
  808. <file>hello.py#hello.py</file>
  809. <capture-output/>
  810. </shell>""" in xml, xml)
  811. action1.capture_output = False
  812. action1.save()
  813. xml = self.wf.to_xml()
  814. assert_true("""
  815. <shell xmlns="uri:oozie:shell-action:0.1">
  816. <job-tracker>${jobTracker}</job-tracker>
  817. <name-node>${nameNode}</name-node>
  818. <job-xml>my-job.xml</job-xml>
  819. <exec>hello.py</exec>
  820. <argument>World!</argument>
  821. <file>hello.py#hello.py</file>
  822. </shell>""" in xml, xml)
  823. def test_workflow_fs_gen_xml(self):
  824. self.wf.node_set.filter(name='action-name-1').delete()
  825. action1 = add_node(self.wf, 'action-name-1', 'fs', [self.wf.start], {
  826. u'name': 'MyFs',
  827. u'description': 'Execute a Fs action that manage files',
  828. u'deletes': '[{"name":"/to/delete"},{"name":"to/delete2"}]',
  829. u'mkdirs': '[{"name":"/to/mkdir"},{"name":"${mkdir2}"}]',
  830. u'moves': '[{"source":"/to/move/source","destination":"/to/move/destination"},{"source":"/to/move/source2","destination":"/to/move/destination2"}]',
  831. u'chmods': '[{"path":"/to/chmod","recursive":true,"permissions":"-rwxrw-rw-"},{"path":"/to/chmod2","recursive":false,"permissions":"755"}]',
  832. u'touchzs': '[{"name":"/to/touchz"},{"name":"/to/touchz2"}]'
  833. })
  834. Link(parent=action1, child=self.wf.end, name="ok").save()
  835. xml = self.wf.to_xml({'mkdir2': '/path'})
  836. assert_true("""
  837. <action name="MyFs">
  838. <fs>
  839. <delete path='${nameNode}/to/delete'/>
  840. <delete path='${nameNode}/user/${wf:user()}/to/delete2'/>
  841. <mkdir path='${nameNode}/to/mkdir'/>
  842. <mkdir path='${nameNode}${mkdir2}'/>
  843. <move source='${nameNode}/to/move/source' target='${nameNode}/to/move/destination'/>
  844. <move source='${nameNode}/to/move/source2' target='${nameNode}/to/move/destination2'/>
  845. <chmod path='${nameNode}/to/chmod' permissions='-rwxrw-rw-' dir-files='true'/>
  846. <chmod path='${nameNode}/to/chmod2' permissions='755' dir-files='false'/>
  847. <touchz path='${nameNode}/to/touchz'/>
  848. <touchz path='${nameNode}/to/touchz2'/>
  849. </fs>
  850. <ok to="end"/>
  851. <error to="kill"/>
  852. </action>""" in xml, xml)
  853. def test_workflow_email_gen_xml(self):
  854. self.wf.node_set.filter(name='action-name-1').delete()
  855. action1 = add_node(self.wf, 'action-name-1', 'email', [self.wf.start], {
  856. u'name': 'MyEmail',
  857. u'description': 'Execute an Email action',
  858. u'to': 'hue@hue.org,django@python.org',
  859. u'cc': '',
  860. u'subject': 'My subject',
  861. u'body': 'My body'
  862. })
  863. Link(parent=action1, child=self.wf.end, name="ok").save()
  864. xml = self.wf.to_xml()
  865. assert_true("""
  866. <action name="MyEmail">
  867. <email xmlns="uri:oozie:email-action:0.1">
  868. <to>hue@hue.org,django@python.org</to>
  869. <subject>My subject</subject>
  870. <body>My body</body>
  871. </email>
  872. <ok to="end"/>
  873. <error to="kill"/>
  874. </action>""" in xml, xml)
  875. action1.cc = 'lambda@python.org'
  876. action1.save()
  877. xml = self.wf.to_xml()
  878. assert_true("""
  879. <action name="MyEmail">
  880. <email xmlns="uri:oozie:email-action:0.1">
  881. <to>hue@hue.org,django@python.org</to>
  882. <cc>lambda@python.org</cc>
  883. <subject>My subject</subject>
  884. <body>My body</body>
  885. </email>
  886. <ok to="end"/>
  887. <error to="kill"/>
  888. </action>""" in xml, xml)
  889. def test_workflow_subworkflow_gen_xml(self):
  890. self.wf.node_set.filter(name='action-name-1').delete()
  891. wf_dict = WORKFLOW_DICT.copy()
  892. wf_dict['name'] = [u'wf-name-2']
  893. wf2 = create_workflow(self.c, self.user, wf_dict)
  894. action1 = add_node(self.wf, 'action-name-1', 'subworkflow', [self.wf.start], {
  895. u'name': 'MySubworkflow',
  896. u'description': 'Execute a subworkflow action',
  897. u'sub_workflow': wf2,
  898. u'propagate_configuration': True,
  899. u'job_properties': '[{"value":"World!","name":"argument"}]'
  900. })
  901. Link(parent=action1, child=self.wf.end, name="ok").save()
  902. xml = self.wf.to_xml()
  903. assert_true(re.search(
  904. '<sub-workflow>\W+'
  905. '<app-path>\${nameNode}/user/hue/oozie/workspaces/_test_-oozie-(.+?)</app-path>\W+'
  906. '<propagate-configuration/>\W+'
  907. '<configuration>\W+'
  908. '<property>\W+'
  909. '<name>argument</name>\W+'
  910. '<value>World!</value>\W+'
  911. '</property>\W+'
  912. '</configuration>\W+'
  913. '</sub-workflow>', xml, re.MULTILINE), xml)
  914. wf2.delete(skip_trash=True)
  915. def test_workflow_flatten_list(self):
  916. assert_equal('[<Start: start>, <Mapreduce: action-name-1>, <Mapreduce: action-name-2>, <Mapreduce: action-name-3>, '
  917. '<Kill: kill>, <End: end>]',
  918. str(self.wf.node_list))
  919. # 1 2
  920. # 3
  921. self.setup_forking_workflow()
  922. assert_equal('[<Start: start>, <Fork: fork-name-1>, <Mapreduce: action-name-1>, <Mapreduce: action-name-2>, '
  923. '<Join: join-name-1>, <Mapreduce: action-name-3>, <Kill: kill>, <End: end>]',
  924. str(self.wf.node_list))
  925. def test_workflow_generic_gen_xml(self):
  926. self.wf.node_set.filter(name='action-name-1').delete()
  927. action1 = add_node(self.wf, 'action-name-1', 'generic', [self.wf.start], {
  928. u'name': 'Generic',
  929. u'description': 'Execute a Generic email action',
  930. u'xml': """
  931. <email xmlns="uri:oozie:email-action:0.1">
  932. <to>hue@hue.org,django@python.org</to>
  933. <subject>My subject</subject>
  934. <body>My body</body>
  935. </email>""",
  936. })
  937. Link(parent=action1, child=self.wf.end, name="ok").save()
  938. xml = self.wf.to_xml()
  939. assert_true("""
  940. <action name="Generic">
  941. <email xmlns="uri:oozie:email-action:0.1">
  942. <to>hue@hue.org,django@python.org</to>
  943. <subject>My subject</subject>
  944. <body>My body</body>
  945. </email>
  946. <ok to="end"/>
  947. <error to="kill"/>
  948. </action>""" in xml, xml)
  949. def test_workflow_hive_gen_xml(self):
  950. self.wf.node_set.filter(name='action-name-1').delete()
  951. action1 = add_node(self.wf, 'action-name-1', 'hive', [self.wf.start], {
  952. u'job_xml': 'my-job.xml',
  953. u'files': '["hello.py"]',
  954. u'name': 'MyHive',
  955. u'job_properties': '[]',
  956. u'script_path': 'hello.sql',
  957. u'archives': '[]',
  958. u'prepares': '[]',
  959. u'params': '[{"value":"World!","type":"argument"}]',
  960. u'description': ''
  961. })
  962. Link(parent=action1, child=self.wf.end, name="ok").save()
  963. xml = self.wf.to_xml()
  964. assert_true("""
  965. <workflow-app name="wf-name-1" xmlns="uri:oozie:workflow:0.4">
  966. <global>
  967. <job-xml>jobconf.xml</job-xml>
  968. <configuration>
  969. <property>
  970. <name>sleep-all</name>
  971. <value>${SLEEP}</value>
  972. </property>
  973. </configuration>
  974. </global>
  975. <start to="MyHive"/>
  976. <action name="MyHive">
  977. <hive xmlns="uri:oozie:hive-action:0.2">
  978. <job-tracker>${jobTracker}</job-tracker>
  979. <name-node>${nameNode}</name-node>
  980. <job-xml>my-job.xml</job-xml>
  981. <script>hello.sql</script>
  982. <argument>World!</argument>
  983. <file>hello.py#hello.py</file>
  984. </hive>
  985. <ok to="end"/>
  986. <error to="kill"/>
  987. </action>
  988. <kill name="kill">
  989. <message>Action failed, error message[${wf:errorMessage(wf:lastErrorNode())}]</message>
  990. </kill>
  991. <end name="end"/>
  992. </workflow-app>""" in xml, xml)
  993. import beeswax
  994. from beeswax.tests import hive_site_xml
  995. tmpdir = tempfile.mkdtemp()
  996. saved = None
  997. try:
  998. # We just replace the Beeswax conf variable
  999. class Getter(object):
  1000. def get(self):
  1001. return tmpdir
  1002. xml = hive_site_xml(is_local=False, use_sasl=True, kerberos_principal='hive/_HOST@test.com')
  1003. file(os.path.join(tmpdir, 'hive-site.xml'), 'w').write(xml)
  1004. beeswax.hive_site.reset()
  1005. saved = beeswax.conf.HIVE_CONF_DIR
  1006. beeswax.conf.HIVE_CONF_DIR = Getter()
  1007. xml = self.wf.to_xml(mapping={
  1008. 'is_kerberized_hive': True,
  1009. 'credential_type': 'hcat',
  1010. 'thrift_server': 'thrift://darkside-1234:9999',
  1011. 'hive_principal': 'hive/darkside-1234@test.com'
  1012. })
  1013. assert_true("""
  1014. <workflow-app name="wf-name-1" xmlns="uri:oozie:workflow:0.4">
  1015. <global>
  1016. <job-xml>jobconf.xml</job-xml>
  1017. <configuration>
  1018. <property>
  1019. <name>sleep-all</name>
  1020. <value>${SLEEP}</value>
  1021. </property>
  1022. </configuration>
  1023. </global>
  1024. <credentials>
  1025. <credential name='hive_credentials' type='hcat'>
  1026. <property>
  1027. <name>hcat.metastore.uri</name>
  1028. <value>thrift://darkside-1234:9999</value>
  1029. </property>
  1030. <property>
  1031. <name>hcat.metastore.principal</name>
  1032. <value>hive/darkside-1234@test.com</value>
  1033. </property>
  1034. </credential>
  1035. </credentials>
  1036. <start to="MyHive"/>
  1037. <action name="MyHive" cred='hive_credentials'>
  1038. <hive xmlns="uri:oozie:hive-action:0.2">
  1039. <job-tracker>${jobTracker}</job-tracker>
  1040. <name-node>${nameNode}</name-node>
  1041. <job-xml>my-job.xml</job-xml>
  1042. <script>hello.sql</script>
  1043. <argument>World!</argument>
  1044. <file>hello.py#hello.py</file>
  1045. </hive>
  1046. <ok to="end"/>
  1047. <error to="kill"/>
  1048. </action>
  1049. <kill name="kill">
  1050. <message>Action failed, error message[${wf:errorMessage(wf:lastErrorNode())}]</message>
  1051. </kill>
  1052. <end name="end"/>
  1053. </workflow-app>""" in xml, xml)
  1054. finally:
  1055. beeswax.hive_site.reset()
  1056. if saved is not None:
  1057. beeswax.conf.HIVE_CONF_DIR = saved
  1058. shutil.rmtree(tmpdir)
  1059. self.wf.node_set.filter(name='action-name-1').delete()
  1060. def test_workflow_gen_workflow_sla(self):
  1061. xml = self.wf.to_xml({'output': '/path'})
  1062. assert_false('<sla' in xml, xml)
  1063. assert_false('xmlns="uri:oozie:workflow:0.5"' in xml, xml)
  1064. assert_false('xmlns:sla="uri:oozie:sla:0.2"' in xml, xml)
  1065. sla = self.wf.sla
  1066. sla[0]['value'] = True
  1067. sla[1]['value'] = 'now' # nominal-time
  1068. sla[3]['value'] = '${ 10 * MINUTES}' # should-end
  1069. self.wf.set_sla(sla)
  1070. self.wf.save()
  1071. xml = self.wf.to_xml({'output': '/path'})
  1072. assert_true('xmlns="uri:oozie:workflow:0.5"' in xml, xml)
  1073. assert_true('xmlns:sla="uri:oozie:sla:0.2"' in xml, xml)
  1074. assert_true("""<end name="end"/>
  1075. <sla:info>
  1076. <sla:nominal-time>now</sla:nominal-time>
  1077. <sla:should-end>${ 10 * MINUTES}</sla:should-end>
  1078. </sla:info>
  1079. </workflow-app>""" in xml, xml)
  1080. def test_workflow_gen_action_sla(self):
  1081. xml = self.wf.to_xml({'output': '/path'})
  1082. assert_false('<sla' in xml, xml)
  1083. assert_false('xmlns="uri:oozie:workflow:0.5"' in xml, xml)
  1084. assert_false('xmlns:sla="uri:oozie:sla:0.2"' in xml, xml)
  1085. self.wf.node_set.filter(name='action-name-1').delete()
  1086. action1 = add_node(self.wf, 'action-name-1', 'hive', [self.wf.start], {
  1087. u'job_xml': 'my-job.xml',
  1088. u'files': '["hello.py"]',
  1089. u'name': 'MyHive',
  1090. u'job_properties': '[]',
  1091. u'script_path': 'hello.sql',
  1092. u'archives': '[]',
  1093. u'prepares': '[]',
  1094. u'params': '[{"value":"World!","type":"argument"}]',
  1095. u'description': ''
  1096. })
  1097. Link(parent=action1, child=self.wf.end, name="ok").save()
  1098. xml = self.wf.to_xml()
  1099. sla = action1.sla
  1100. sla[0]['value'] = True
  1101. sla[1]['value'] = 'now' # nominal-time
  1102. sla[3]['value'] = '${ 10 * MINUTES}' # should-end
  1103. action1.set_sla(sla)
  1104. action1.save()
  1105. xml = self.wf.to_xml({'output': '/path'})
  1106. assert_true('xmlns="uri:oozie:workflow:0.5"' in xml, xml)
  1107. assert_true('xmlns:sla="uri:oozie:sla:0.2"' in xml, xml)
  1108. assert_true("""<error to="kill"/>
  1109. <sla:info>
  1110. <sla:nominal-time>now</sla:nominal-time>
  1111. <sla:should-end>${ 10 * MINUTES}</sla:should-end>
  1112. </sla:info>
  1113. </action>""" in xml, xml)
  1114. def test_create_coordinator(self):
  1115. create_coordinator(self.wf, self.c, self.user)
  1116. def test_clone_coordinator(self):
  1117. coord = create_coordinator(self.wf, self.c, self.user)
  1118. coordinator_count = Document.objects.available_docs(Coordinator, self.user).count()
  1119. response = self.c.post(reverse('oozie:clone_coordinator', args=[coord.id]), {}, follow=True)
  1120. coord2 = Coordinator.objects.latest('id')
  1121. assert_not_equal(coord.id, coord2.id)
  1122. assert_equal(coordinator_count + 1, Document.objects.available_docs(Coordinator, self.user).count(), response)
  1123. assert_equal(coord.dataset_set.count(), coord2.dataset_set.count())
  1124. assert_equal(coord.datainput_set.count(), coord2.datainput_set.count())
  1125. assert_equal(coord.dataoutput_set.count(), coord2.dataoutput_set.count())
  1126. ds_ids = set(coord.dataset_set.values_list('id', flat=True))
  1127. for node in coord2.dataset_set.all():
  1128. assert_false(node.id in ds_ids)
  1129. data_input_ids = set(coord.datainput_set.values_list('id', flat=True))
  1130. for node in coord2.datainput_set.all():
  1131. assert_false(node.id in data_input_ids)
  1132. data_output_ids = set(coord.dataoutput_set.values_list('id', flat=True))
  1133. for node in coord2.dataoutput_set.all():
  1134. assert_false(node.id in data_output_ids)
  1135. assert_not_equal(coord.deployment_dir, coord2.deployment_dir)
  1136. assert_not_equal('', coord2.deployment_dir)
  1137. # Bulk delete
  1138. response = self.c.post(reverse('oozie:delete_coordinator'), {'job_selection': [coord.id, coord2.id]}, follow=True)
  1139. assert_equal(coordinator_count - 1, Document.objects.available_docs(Coordinator, self.user).count(), response)
  1140. def test_coordinator_workflow_access_permissions(self):
  1141. raise SkipTest
  1142. self.wf.is_shared = True
  1143. self.wf.save()
  1144. # Login as someone else not superuser
  1145. client_another_me = make_logged_in_client(username='another_me', is_superuser=False, groupname='test')
  1146. grant_access("another_me", "test", "oozie")
  1147. coord = create_coordinator(self.wf, client_another_me, self.user)
  1148. response = client_another_me.get(reverse('oozie:edit_coordinator', args=[coord.id]))
  1149. assert_true('Editor' in response.content, response.content)
  1150. assert_true('Save coordinator' in response.content, response.content)
  1151. # Check can schedule a non personal/shared workflow
  1152. workflow_select = '%s</option>' % self.wf
  1153. response = client_another_me.get(reverse('oozie:edit_coordinator', args=[coord.id]))
  1154. assert_true(workflow_select in response.content, response.content)
  1155. self.wf.is_shared = False
  1156. self.wf.save()
  1157. response = client_another_me.get(reverse('oozie:edit_coordinator', args=[coord.id]))
  1158. assert_false(workflow_select in response.content, response.content)
  1159. self.wf.is_shared = True
  1160. self.wf.save()
  1161. # Edit
  1162. finish = SHARE_JOBS.set_for_testing(True)
  1163. try:
  1164. response = client_another_me.post(reverse('oozie:edit_coordinator', args=[coord.id]))
  1165. assert_true(workflow_select in response.content, response.content)
  1166. assert_true('Save coordinator' in response.content, response.content)
  1167. finally:
  1168. finish()
  1169. finish = SHARE_JOBS.set_for_testing(False)
  1170. try:
  1171. response = client_another_me.post(reverse('oozie:edit_coordinator', args=[coord.id]))
  1172. assert_true('This field is required' in response.content, response.content)
  1173. assert_false(workflow_select in response.content, response.content)
  1174. assert_true('Save coordinator' in response.content, response.content)
  1175. finally:
  1176. finish()
  1177. def test_coordinator_gen_xml(self):
  1178. coord = create_coordinator(self.wf, self.c, self.user)
  1179. assert_true(
  1180. """<coordinator-app name="MyCoord"
  1181. frequency="${coord:days(1)}"
  1182. start="2012-07-01T00:00Z" end="2012-07-04T00:00Z" timezone="America/Los_Angeles"
  1183. xmlns="uri:oozie:coordinator:0.2">
  1184. <controls>
  1185. <timeout>100</timeout>
  1186. <concurrency>3</concurrency>
  1187. <execution>FIFO</execution>
  1188. <throttle>10</throttle>
  1189. </controls>
  1190. <action>
  1191. <workflow>
  1192. <app-path>${wf_application_path}</app-path>
  1193. <configuration>
  1194. <property>
  1195. <name>username</name>
  1196. <value>${coord:user()}</value>
  1197. </property>
  1198. <property>
  1199. <name>SLEEP</name>
  1200. <value>1000</value>
  1201. </property>
  1202. <property>
  1203. <name>market</name>
  1204. <value>US</value>
  1205. </property>
  1206. </configuration>
  1207. </workflow>
  1208. </action>
  1209. </coordinator-app>""" in coord.to_xml(), coord.to_xml())
  1210. def test_coordinator_gen_sla(self):
  1211. coord = create_coordinator(self.wf, self.c, self.user)
  1212. xml = coord.to_xml()
  1213. assert_false('<sla' in xml, xml)
  1214. assert_false('xmlns="uri:oozie:coordinator:0.4"' in xml, xml)
  1215. assert_false('xmlns:sla="uri:oozie:sla:0.2"' in xml, xml)
  1216. sla = coord.sla
  1217. sla[0]['value'] = True
  1218. sla[1]['value'] = 'now' # nominal-time
  1219. sla[3]['value'] = '${ 10 * MINUTES}' # should-end
  1220. coord.set_sla(sla)
  1221. coord.save()
  1222. xml = coord.to_xml()
  1223. assert_true('xmlns="uri:oozie:coordinator:0.4"' in xml, xml)
  1224. assert_true('xmlns:sla="uri:oozie:sla:0.2"' in xml, xml)
  1225. assert_true("""</workflow>
  1226. <sla:info>
  1227. <sla:nominal-time>now</sla:nominal-time>
  1228. <sla:should-end>${ 10 * MINUTES}</sla:should-end>
  1229. </sla:info>
  1230. </action>""" in xml, xml)
  1231. def test_coordinator_with_data_input_gen_xml(self):
  1232. coord = create_coordinator(self.wf, self.c, self.user)
  1233. create_dataset(coord, self.c)
  1234. create_coordinator_data(coord, self.c)
  1235. self.c.post(reverse('oozie:create_coordinator_dataset', args=[coord.id]), {
  1236. u'create-name': [u'MyDataset2'], u'create-frequency_number': [u'1'], u'create-frequency_unit': [u'days'],
  1237. u'create-uri': [u's3n://a-server/data/out/${YEAR}${MONTH}${DAY}'],
  1238. u'create-instance_choice': [u'single'],
  1239. u'instance_start': [u'-1'],
  1240. u'create-advanced_start_instance': [u'0'],
  1241. u'create-advanced_end_instance': [u'0'],
  1242. u'create-start_0': [u'07/01/2012'], u'create-start_1': [u'12:00 AM'],
  1243. u'create-timezone': [u'America/Los_Angeles'], u'create-done_flag': [u''],
  1244. u'create-description': [u'']})
  1245. self.c.post(reverse('oozie:create_coordinator_data', args=[coord.id, 'output']),
  1246. {u'output-name': [u'output_dir'], u'output-dataset': [u'2']})
  1247. assert_true(
  1248. """<coordinator-app name="MyCoord"
  1249. frequency="${coord:days(1)}"
  1250. start="2012-07-01T00:00Z" end="2012-07-04T00:00Z" timezone="America/Los_Angeles"
  1251. xmlns="uri:oozie:coordinator:0.2">
  1252. <controls>
  1253. <timeout>100</timeout>
  1254. <concurrency>3</concurrency>
  1255. <execution>FIFO</execution>
  1256. <throttle>10</throttle>
  1257. </controls>
  1258. <datasets>
  1259. <dataset name="MyDataset" frequency="${coord:days(1)}"
  1260. initial-instance="2012-07-01T00:00Z" timezone="America/Los_Angeles">
  1261. <uri-template>${nameNode}/data/${YEAR}${MONTH}${DAY}</uri-template>
  1262. <done-flag></done-flag>
  1263. </dataset>
  1264. <dataset name="MyDataset2" frequency="${coord:days(1)}"
  1265. initial-instance="2012-07-01T00:00Z" timezone="America/Los_Angeles">
  1266. <uri-template>s3n://a-server/data/out/${YEAR}${MONTH}${DAY}</uri-template>
  1267. <done-flag></done-flag>
  1268. </dataset>
  1269. </datasets>
  1270. <input-events>
  1271. <data-in name="input_dir" dataset="MyDataset">
  1272. <start-instance>
  1273. ${coord:current(-1)}
  1274. </start-instance>
  1275. <end-instance>
  1276. ${coord:current(1)}
  1277. </end-instance>
  1278. </data-in>
  1279. </input-events>
  1280. <output-events>
  1281. <data-out name="output_dir" dataset="MyDataset2">
  1282. <instance>${coord:current(0)}</instance>
  1283. </data-out>
  1284. </output-events>
  1285. <action>
  1286. <workflow>
  1287. <app-path>${wf_application_path}</app-path>
  1288. <configuration>
  1289. <property>
  1290. <name>input_dir</name>
  1291. <value>${coord:dataIn('input_dir')}</value>
  1292. </property>
  1293. <property>
  1294. <name>output_dir</name>
  1295. <value>${coord:dataOut('output_dir')}</value>
  1296. </property>
  1297. <property>
  1298. <name>username</name>
  1299. <value>${coord:user()}</value>
  1300. </property>
  1301. <property>
  1302. <name>SLEEP</name>
  1303. <value>1000</value>
  1304. </property>
  1305. <property>
  1306. <name>market</name>
  1307. <value>US</value>
  1308. </property>
  1309. </configuration>
  1310. </workflow>
  1311. </action>
  1312. </coordinator-app>""" in coord.to_xml(), coord.to_xml())
  1313. def test_create_coordinator_dataset(self):
  1314. coord = create_coordinator(self.wf, self.c, self.user)
  1315. create_dataset(coord, self.c)
  1316. def test_edit_coordinator_dataset(self):
  1317. coord = create_coordinator(self.wf, self.c, self.user)
  1318. create_dataset(coord, self.c)
  1319. response = self.c.post(reverse('oozie:edit_coordinator_dataset', args=[1]), {
  1320. u'edit-name': [u'MyDataset'], u'edit-frequency_number': [u'1'], u'edit-frequency_unit': [u'days'],
  1321. u'edit-uri': [u'/data/${YEAR}${MONTH}${DAY}'],
  1322. u'edit-start_0': [u'07/01/2012'], u'edit-start_1': [u'12:00 AM'],
  1323. u'edit-instance_choice': [u'range'],
  1324. u'edit-advanced_start_instance': [u'-1'],
  1325. u'edit-advanced_end_instance': [u'${coord:current(1)}'],
  1326. u'edit-start_0': [u'07/01/2012'], u'edit-start_1': [u'12:00 AM'],
  1327. u'edit-timezone': [u'America/Los_Angeles'], u'edit-done_flag': [u''],
  1328. u'edit-description': [u'']}, follow=True)
  1329. data = json.loads(response.content)
  1330. assert_equal(0, data['status'], data['status'])
  1331. def test_create_coordinator_input_data(self):
  1332. coord = create_coordinator(self.wf, self.c, self.user)
  1333. create_dataset(coord, self.c)
  1334. create_coordinator_data(coord, self.c)
  1335. def test_install_examples(self):
  1336. self.c.post(reverse('oozie:install_examples'))
  1337. def test_workflow_prepare(self):
  1338. action1 = Node.objects.get(workflow=self.wf, name='action-name-1').get_full_node()
  1339. action1.prepares = json.dumps([
  1340. {"type": "delete","value": "${output}"},
  1341. {"type": "delete","value": "out"},
  1342. {"type": "delete","value": "/user/test/out"},
  1343. {"type": "delete","value": "hdfs://localhost:8020/user/test/out"}])
  1344. action1.save()
  1345. xml = self.wf.to_xml({'output': '/path'})
  1346. assert_true('<delete path="${nameNode}${output}"/>' in xml, xml)
  1347. assert_true('<delete path="${nameNode}/user/${wf:user()}/out"/>' in xml, xml)
  1348. assert_true('<delete path="${nameNode}/user/test/out"/>' in xml, xml)
  1349. assert_true('<delete path="hdfs://localhost:8020/user/test/out"/>' in xml, xml)
  1350. def test_get_workflow_parameters(self):
  1351. assert_equal([{'name': u'output', 'value': ''}, {'name': u'SLEEP', 'value': ''}, {'name': u'market', 'value': u'US'}],
  1352. self.wf.find_all_parameters())
  1353. def test_get_coordinator_parameters(self):
  1354. coord = create_coordinator(self.wf, self.c, self.user)
  1355. create_dataset(coord, self.c)
  1356. create_coordinator_data(coord, self.c)
  1357. assert_equal([{'name': u'output', 'value': ''}, {'name': u'market', 'value': u'US'}],
  1358. coord.find_all_parameters())
  1359. def test_workflow_data_binds(self):
  1360. response = self.c.get(reverse('oozie:edit_workflow', args=[self.wf.id]))
  1361. assert_equal(1, response.content.count('checked: is_shared'), response.content)
  1362. assert_true('checked: capture_output' in response.content, response.content)
  1363. def test_xss_escape_js(self):
  1364. hacked = '[{"name":"oozie.use.system.libpath","value":"true"}, {"name": "123\\"><script>alert(1)</script>", "value": "\'hacked\'"}]'
  1365. escaped = '[{"name": "oozie.use.system.libpath", "value": "true"}, {"name": "123\\"\\u003e\\u003cscript\\u003ealert(1)\\u003c/script\\u003e", "value": "\'hacked\'"}]'
  1366. self.wf.job_properties = hacked
  1367. self.wf.parameters = hacked
  1368. assert_equal(escaped, self.wf._escapejs_parameters_list(hacked))
  1369. assert_equal(escaped, self.wf.job_properties_escapejs)
  1370. assert_equal(escaped, self.wf.parameters_escapejs)
  1371. def test_xss_html_escaping(self):
  1372. data = WORKFLOW_DICT.copy()
  1373. data['description'] = [u'"><script>alert(1);</script>']
  1374. self.wf = create_workflow(self.c, self.user, workflow_dict=data)
  1375. resp = self.c.get('/oozie/list_workflows/')
  1376. assert_false('"><script>alert(1);</script>' in resp.content, resp.content)
  1377. assert_true('&quot;&gt;&lt;script&gt;alert(1);&lt;/script&gt;' in resp.content, resp.content)
  1378. def test_submit_workflow(self):
  1379. # Check param popup
  1380. response = self.c.get(reverse('oozie:submit_workflow', args=[self.wf.id]))
  1381. assert_equal([{'name': u'output', 'value': ''},
  1382. {'name': u'SLEEP', 'value': ''},
  1383. {'name': u'market', 'value': u'US'}
  1384. ],
  1385. response.context['params_form'].initial)
  1386. def test_submit_coordinator(self):
  1387. coord = create_coordinator(self.wf, self.c, self.user)
  1388. # Check param popup, SLEEP is set by coordinator so not shown in the popup
  1389. response = self.c.get(reverse('oozie:submit_coordinator', args=[coord.id]))
  1390. assert_equal([{'name': u'output', 'value': ''},
  1391. {'name': u'market', 'value': u'US'}
  1392. ],
  1393. response.context['params_form'].initial)
  1394. def test_trash_workflow(self):
  1395. previous_trashed = Document.objects.trashed_docs(Workflow, self.user).count()
  1396. previous_available = Document.objects.available_docs(Workflow, self.user).count()
  1397. response = self.c.post(reverse('oozie:delete_workflow'), {'job_selection': [self.wf.id]}, follow=True)
  1398. assert_equal(200, response.status_code, response)
  1399. assert_equal(previous_trashed + 1, Document.objects.trashed_docs(Workflow, self.user).count())
  1400. assert_equal(previous_available - 1, Document.objects.available_docs(Workflow, self.user).count())
  1401. def test_workflow_export(self):
  1402. response = self.c.get(reverse('oozie:export_workflow', args=[self.wf.id]))
  1403. zfile = zipfile.ZipFile(StringIO.StringIO(response.content))
  1404. assert_true('workflow.xml' in zfile.namelist(), 'workflow.xml not in response')
  1405. assert_true('workflow-metadata.json' in zfile.namelist(), 'workflow-metadata.json not in response')
  1406. assert_equal(2, len(zfile.namelist()))
  1407. workflow_xml = reformat_xml("""<workflow-app name="wf-name-1" xmlns="uri:oozie:workflow:0.4">
  1408. <global>
  1409. <job-xml>jobconf.xml</job-xml>
  1410. <configuration>
  1411. <property>
  1412. <name>sleep-all</name>
  1413. <value>${SLEEP}</value>
  1414. </property>
  1415. </configuration>
  1416. </global>
  1417. <start to="action-name-1"/>
  1418. <action name="action-name-1">
  1419. <map-reduce>
  1420. <job-tracker>${jobTracker}</job-tracker>
  1421. <name-node>${nameNode}</name-node>
  1422. <prepare>
  1423. <delete path="${nameNode}${output}"/>
  1424. <mkdir path="${nameNode}/test"/>
  1425. </prepare>
  1426. <configuration>
  1427. <property>
  1428. <name>sleep</name>
  1429. <value>${SLEEP}</value>
  1430. </property>
  1431. </configuration>
  1432. </map-reduce>
  1433. <ok to="action-name-2"/>
  1434. <error to="kill"/>
  1435. </action>
  1436. <action name="action-name-2">
  1437. <map-reduce>
  1438. <job-tracker>${jobTracker}</job-tracker>
  1439. <name-node>${nameNode}</name-node>
  1440. <prepare>
  1441. <delete path="${nameNode}${output}"/>
  1442. <mkdir path="${nameNode}/test"/>
  1443. </prepare>
  1444. <configuration>
  1445. <property>
  1446. <name>sleep</name>
  1447. <value>${SLEEP}</value>
  1448. </property>
  1449. </configuration>
  1450. </map-reduce>
  1451. <ok to="action-name-3"/>
  1452. <error to="kill"/>
  1453. </action>
  1454. <action name="action-name-3">
  1455. <map-reduce>
  1456. <job-tracker>${jobTracker}</job-tracker>
  1457. <name-node>${nameNode}</name-node>
  1458. <prepare>
  1459. <delete path="${nameNode}${output}"/>
  1460. <mkdir path="${nameNode}/test"/>
  1461. </prepare>
  1462. <configuration>
  1463. <property>
  1464. <name>sleep</name>
  1465. <value>${SLEEP}</value>
  1466. </property>
  1467. </configuration>
  1468. </map-reduce>
  1469. <ok to="end"/>
  1470. <error to="kill"/>
  1471. </action>
  1472. <kill name="kill">
  1473. <message>Action failed, error message[${wf:errorMessage(wf:lastErrorNode())}]</message>
  1474. </kill>
  1475. <end name="end"/>
  1476. </workflow-app>""")
  1477. workflow_metadata_json = reformat_json("""{
  1478. "attributes": {
  1479. "deployment_dir": "/user/hue/oozie/workspaces/_test_-oozie-13-1383539302.62",
  1480. "description": ""
  1481. },
  1482. "nodes": {
  1483. "action-name-1": {
  1484. "attributes": {
  1485. "jar_path": "/user/hue/oozie/examples/lib/hadoop-examples.jar"
  1486. }
  1487. },
  1488. "action-name-2": {
  1489. "attributes": {
  1490. "jar_path": "/user/hue/oozie/examples/lib/hadoop-examples.jar"
  1491. }
  1492. },
  1493. "action-name-3": {
  1494. "attributes": {
  1495. "jar_path": "/user/hue/oozie/examples/lib/hadoop-examples.jar"
  1496. }
  1497. }
  1498. },
  1499. "version": "0.0.1"
  1500. }""")
  1501. result_workflow_metadata_json = reformat_json(zfile.read('workflow-metadata.json'))
  1502. workflow_metadata_json = synchronize_workflow_attributes(workflow_metadata_json, result_workflow_metadata_json)
  1503. assert_equal(workflow_xml, reformat_xml(zfile.read('workflow.xml')))
  1504. assert_equal(workflow_metadata_json, result_workflow_metadata_json)
  1505. class TestEditorBundle(OozieMockBase):
  1506. def setUp(self):
  1507. super(TestEditorBundle, self).setUp()
  1508. self.setup_simple_workflow()
  1509. def test_create_bundle(self):
  1510. create_bundle(self.c, self.user)
  1511. def test_clone_bundle(self):
  1512. bundle = create_bundle(self.c, self.user)
  1513. bundle_count = Document.objects.available_docs(Bundle, self.user).count()
  1514. response = self.c.post(reverse('oozie:clone_bundle', args=[bundle.id]), {}, follow=True)
  1515. bundle2 = Bundle.objects.latest('id')
  1516. assert_not_equal(bundle.id, bundle2.id)
  1517. assert_equal(bundle_count + 1, Document.objects.available_docs(Bundle, self.user).count(), response)
  1518. coord_ids = set(bundle.coordinators.values_list('id', flat=True))
  1519. coord2_ids = set(bundle2.coordinators.values_list('id', flat=True))
  1520. if coord_ids or coord2_ids:
  1521. assert_not_equal(coord_ids, coord2_ids)
  1522. assert_not_equal(bundle.deployment_dir, bundle2.deployment_dir)
  1523. assert_not_equal('', bundle2.deployment_dir)
  1524. # Bulk delete
  1525. response = self.c.post(reverse('oozie:delete_bundle'), {'job_selection': [bundle.id, bundle2.id]}, follow=True)
  1526. assert_equal(bundle_count - 1, Document.objects.available_docs(Bundle, self.user).count(), response)
  1527. def test_delete_bundle(self):
  1528. bundle = create_bundle(self.c, self.user)
  1529. bundle_count = Document.objects.available_docs(Bundle, self.user).count()
  1530. response = self.c.post(reverse('oozie:delete_bundle'), {'job_selection': [bundle.id]}, follow=True)
  1531. assert_equal(bundle_count - 1, Document.objects.available_docs(Bundle, self.user).count(), response)
  1532. def test_bundle_gen_xml(self):
  1533. bundle = create_bundle(self.c, self.user)
  1534. assert_true(
  1535. """<bundle-app name="MyBundle"
  1536. xmlns:xsi='http://www.w3.org/2001/XMLSchema-instance'
  1537. xmlns="uri:oozie:coordinator:0.2">
  1538. <parameters>
  1539. <property>
  1540. <name>market</name>
  1541. <value>US,France</value>
  1542. </property>
  1543. </parameters>
  1544. <controls>
  1545. <kick-off-time>2012-07-01T00:00Z</kick-off-time>
  1546. </controls>
  1547. </bundle-app>
  1548. """ in bundle.to_xml(), bundle.to_xml())
  1549. def test_create_bundled_coordinator(self):
  1550. bundle = create_bundle(self.c, self.user)
  1551. coord = create_coordinator(self.wf, self.c, self.user)
  1552. post = {
  1553. u'name': [u'test2'], u'kick_off_time_0': [u'02/12/2013'], u'kick_off_time_1': [u'05:05 PM'],
  1554. u'create-bundled-coordinator-parameters': [u'[{"name":"market","value":"US"}]'],
  1555. u'schema_version': [u'uri:oozie:bundle:0.2', u'uri:oozie:bundle:0.2'], u'coordinators-MAX_NUM_FORMS': [u'0'],
  1556. u'coordinators-INITIAL_FORMS': [u'0'],
  1557. u'parameters': [u'[{"name":"oozie.use.system.libpath","value":"true"}]'], u'coordinators-TOTAL_FORMS': [u'0'],
  1558. u'description': [u'ss']
  1559. }
  1560. response = self.c.get(reverse('oozie:create_bundled_coordinator', args=[bundle.id]))
  1561. assert_true('Add coordinator' in response.content, response.content)
  1562. response = self.c.post(reverse('oozie:create_bundled_coordinator', args=[bundle.id]), post, follow=True)
  1563. assert_true('This field is required' in response.content, response.content)
  1564. post['create-bundled-coordinator-coordinator'] = ['%s' % coord.id]
  1565. response = self.c.post(reverse('oozie:create_bundled_coordinator', args=[bundle.id]), post, follow=True)
  1566. assert_true('Coordinators' in response.content, response.content)
  1567. xml = bundle.to_xml({
  1568. 'wf_%s_dir' % self.wf.id: '/deployment_path_wf',
  1569. 'coord_%s_dir' % coord.id: '/deployment_path_coord'
  1570. })
  1571. assert_true(
  1572. """<bundle-app name="MyBundle"
  1573. xmlns:xsi='http://www.w3.org/2001/XMLSchema-instance'
  1574. xmlns="uri:oozie:coordinator:0.2">
  1575. <parameters>
  1576. <property>
  1577. <name>market</name>
  1578. <value>US,France</value>
  1579. </property>
  1580. </parameters>
  1581. <controls>
  1582. <kick-off-time>2012-07-01T00:00Z</kick-off-time>
  1583. </controls>
  1584. <coordinator name='MyCoord-1' >
  1585. <app-path>${nameNode}/deployment_path_coord</app-path>
  1586. <configuration>
  1587. <property>
  1588. <name>wf_application_path</name>
  1589. <value>/deployment_path_wf</value>
  1590. </property>
  1591. <property>
  1592. <name>market</name>
  1593. <value>US</value>
  1594. </property>
  1595. </configuration>
  1596. </coordinator>
  1597. </bundle-app>""" in xml, xml)
  1598. class TestImportWorkflow04(OozieMockBase):
  1599. def setUp(self):
  1600. raise SkipTest
  1601. super(TestImportWorkflow04, self).setUp()
  1602. self.setup_simple_workflow()
  1603. @raises(RuntimeError)
  1604. def test_import_workflow_namespace_error(self):
  1605. """
  1606. Validates import for most basic workflow with an error.
  1607. """
  1608. workflow = Workflow.objects.new_workflow(self.user)
  1609. workflow.save()
  1610. f = open('apps/oozie/src/oozie/test_data/workflows/0.4/test-basic-namespace-missing.xml')
  1611. contents = f.read()
  1612. f.close()
  1613. # Should throw PopupException
  1614. import_workflow(workflow, contents)
  1615. def test_import_workflow_basic(self):
  1616. """
  1617. Validates import for most basic workflow: start and end.
  1618. """
  1619. workflow = Workflow.objects.new_workflow(self.user)
  1620. workflow.save()
  1621. f = open('apps/oozie/src/oozie/test_data/workflows/0.4/test-basic.xml')
  1622. import_workflow(workflow, f.read())
  1623. f.close()
  1624. workflow.save()
  1625. assert_equal(2, len(Node.objects.filter(workflow=workflow)))
  1626. assert_equal(2, len(Link.objects.filter(parent__workflow=workflow)))
  1627. assert_equal('done', Node.objects.get(workflow=workflow, node_type='end').name)
  1628. assert_equal('uri:oozie:workflow:0.4', workflow.schema_version)
  1629. workflow.delete(skip_trash=True)
  1630. def test_import_workflow_basic_global_config(self):
  1631. """
  1632. Validates import for basic workflow: start, end, and global configuration.
  1633. """
  1634. workflow = Workflow.objects.new_workflow(self.user)
  1635. workflow.save()
  1636. f = open('apps/oozie/src/oozie/test_data/workflows/0.4/test-basic-global-config.xml')
  1637. import_workflow(workflow, f.read())
  1638. f.close()
  1639. workflow.save()
  1640. assert_equal(2, len(Node.objects.filter(workflow=workflow)))
  1641. assert_equal(2, len(Link.objects.filter(parent__workflow=workflow)))
  1642. assert_equal('done', Node.objects.get(workflow=workflow, node_type='end').name)
  1643. assert_equal('uri:oozie:workflow:0.4', workflow.schema_version)
  1644. workflow.delete(skip_trash=True)
  1645. def test_import_workflow_decision(self):
  1646. """
  1647. Validates import for decision node: link comments (conditions), default link, decision end.
  1648. """
  1649. workflow = Workflow.objects.new_workflow(self.user)
  1650. workflow.save()
  1651. f = open('apps/oozie/src/oozie/test_data/workflows/0.4/test-decision.xml')
  1652. import_workflow(workflow, f.read())
  1653. f.close()
  1654. workflow.save()
  1655. assert_equal(12, len(Node.objects.filter(workflow=workflow)))
  1656. assert_equal(21, len(Link.objects.filter(parent__workflow=workflow)))
  1657. assert_equal(1, len(Link.objects.filter(parent__workflow=workflow, parent__node_type='decision', comment='${1 gt 2}', name='start')))
  1658. assert_equal(1, len(Link.objects.filter(parent__workflow=workflow, parent__node_type='decision', comment='', name='start')))
  1659. assert_equal(1, len(Link.objects.filter(parent__workflow=workflow, parent__node_type='decision', name='default')))
  1660. assert_equal(1, len(Link.objects.filter(parent__workflow=workflow, parent__node_type='decision', child__node_type='decisionend', name='related')))
  1661. workflow.delete(skip_trash=True)
  1662. def test_import_workflow_decision_complex(self):
  1663. workflow = Workflow.objects.new_workflow(self.user)
  1664. workflow.save()
  1665. f = open('apps/oozie/src/oozie/test_data/workflows/0.4/test-decision-complex.xml')
  1666. import_workflow(workflow, f.read())
  1667. f.close()
  1668. workflow.save()
  1669. assert_equal(14, len(Node.objects.filter(workflow=workflow)))
  1670. assert_equal(27, len(Link.objects.filter(parent__workflow=workflow)))
  1671. assert_equal(3, len(Link.objects.filter(parent__workflow=workflow, parent__node_type='decision', comment='${ 1 gt 2 }', name='start')))
  1672. assert_equal(0, len(Link.objects.filter(parent__workflow=workflow, parent__node_type='decision', comment='', name='start')))
  1673. assert_equal(3, len(Link.objects.filter(parent__workflow=workflow, parent__node_type='decision', name='default')))
  1674. assert_equal(3, len(Link.objects.filter(parent__workflow=workflow, parent__node_type='decision', child__node_type='decisionend', name='related')))
  1675. workflow.delete(skip_trash=True)
  1676. def test_import_workflow_distcp(self):
  1677. """
  1678. Validates import for distcp node: params.
  1679. """
  1680. workflow = Workflow.objects.new_workflow(self.user)
  1681. workflow.save()
  1682. f = open('apps/oozie/src/oozie/test_data/workflows/0.4/test-distcp.0.1.xml')
  1683. import_workflow(workflow, f.read())
  1684. f.close()
  1685. workflow.save()
  1686. assert_equal(4, len(Node.objects.filter(workflow=workflow)))
  1687. assert_equal(4, len(Link.objects.filter(parent__workflow=workflow)))
  1688. assert_equal('[{"type":"arg","value":"-overwrite"},{"type":"arg","value":"-m"},{"type":"arg","value":"${MAP_NUMBER}"},{"type":"arg","value":"/user/hue/oozie/workspaces/data"},{"type":"arg","value":"${OUTPUT}"}]', Node.objects.get(workflow=workflow, node_type='distcp').get_full_node().params)
  1689. workflow.delete(skip_trash=True)
  1690. def test_import_workflow_forks(self):
  1691. workflow = Workflow.objects.new_workflow(self.user)
  1692. workflow.save()
  1693. f = open('apps/oozie/src/oozie/test_data/workflows/0.4/test-forks.xml')
  1694. import_workflow(workflow, f.read())
  1695. f.close()
  1696. workflow.save()
  1697. assert_equal(12, len(Node.objects.filter(workflow=workflow)))
  1698. assert_equal(20, len(Link.objects.filter(parent__workflow=workflow)))
  1699. assert_equal(6, len(Link.objects.filter(parent__workflow=workflow, parent__node_type='fork')))
  1700. assert_equal(4, len(Link.objects.filter(parent__workflow=workflow, parent__node_type='fork', name='start')))
  1701. assert_equal(2, len(Link.objects.filter(parent__workflow=workflow, parent__node_type='fork', child__node_type='join', name='related')))
  1702. workflow.delete(skip_trash=True)
  1703. def test_import_workflow_mapreduce(self):
  1704. """
  1705. Validates import for mapreduce node: job_properties.
  1706. """
  1707. workflow = Workflow.objects.new_workflow(self.user)
  1708. workflow.save()
  1709. f = open('apps/oozie/src/oozie/test_data/workflows/0.4/test-mapreduce.xml')
  1710. import_workflow(workflow, f.read())
  1711. f.close()
  1712. workflow.save()
  1713. assert_equal(4, len(Node.objects.filter(workflow=workflow)))
  1714. assert_equal(4, len(Link.objects.filter(parent__workflow=workflow)))
  1715. assert_equal('[{"name":"mapred.reduce.tasks","value":"1"},{"name":"mapred.mapper.class","value":"org.apache.hadoop.examples.SleepJob"},{"name":"mapred.reducer.class","value":"org.apache.hadoop.examples.SleepJob"},{"name":"mapred.mapoutput.key.class","value":"org.apache.hadoop.io.IntWritable"},{"name":"mapred.mapoutput.value.class","value":"org.apache.hadoop.io.NullWritable"},{"name":"mapred.output.format.class","value":"org.apache.hadoop.mapred.lib.NullOutputFormat"},{"name":"mapred.input.format.class","value":"org.apache.hadoop.examples.SleepJob$SleepInputFormat"},{"name":"mapred.partitioner.class","value":"org.apache.hadoop.examples.SleepJob"},{"name":"mapred.speculative.execution","value":"false"},{"name":"sleep.job.map.sleep.time","value":"0"},{"name":"sleep.job.reduce.sleep.time","value":"1"}]', Node.objects.get(workflow=workflow, node_type='mapreduce').get_full_node().job_properties)
  1716. workflow.delete(skip_trash=True)
  1717. def test_import_workflow_pig(self):
  1718. """
  1719. Validates import for pig node: params.
  1720. """
  1721. workflow = Workflow.objects.new_workflow(self.user)
  1722. workflow.save()
  1723. f = open('apps/oozie/src/oozie/test_data/workflows/0.4/test-pig.xml')
  1724. import_workflow(workflow, f.read())
  1725. f.close()
  1726. workflow.save()
  1727. node = Node.objects.get(workflow=workflow, node_type='pig').get_full_node()
  1728. assert_equal(4, len(Node.objects.filter(workflow=workflow)))
  1729. assert_equal(4, len(Link.objects.filter(parent__workflow=workflow)))
  1730. assert_equal('aggregate.pig', node.script_path)
  1731. assert_equal('[{"type":"argument","value":"-param"},{"type":"argument","value":"INPUT=/user/hue/oozie/workspaces/data"},{"type":"argument","value":"-param"},{"type":"argument","value":"OUTPUT=${output}"}]', node.params)
  1732. workflow.delete(skip_trash=True)
  1733. def test_import_workflow_sqoop(self):
  1734. """
  1735. Validates import for sqoop node: script_path, files.
  1736. """
  1737. workflow = Workflow.objects.new_workflow(self.user)
  1738. workflow.save()
  1739. f = open('apps/oozie/src/oozie/test_data/workflows/0.4/test-sqoop.0.2.xml')
  1740. import_workflow(workflow, f.read())
  1741. f.close()
  1742. workflow.save()
  1743. assert_equal(4, len(Node.objects.filter(workflow=workflow)))
  1744. assert_equal(4, len(Link.objects.filter(parent__workflow=workflow)))
  1745. node = Node.objects.get(workflow=workflow, node_type='sqoop').get_full_node()
  1746. assert_equal('["db.hsqldb.properties#db.hsqldb.properties","db.hsqldb.script#db.hsqldb.script"]', node.files)
  1747. assert_equal('import --connect jdbc:hsqldb:file:db.hsqldb --table TT --target-dir ${output} -m 1', node.script_path)
  1748. assert_equal('[{"type":"arg","value":"My invalid arg"},{"type":"arg","value":"My invalid arg 2"}]', node.params)
  1749. workflow.delete(skip_trash=True)
  1750. def test_import_workflow_java(self):
  1751. """
  1752. Validates import for java node: main_class, args.
  1753. """
  1754. workflow = Workflow.objects.new_workflow(self.user)
  1755. workflow.save()
  1756. f = open('apps/oozie/src/oozie/test_data/workflows/0.4/test-java.xml')
  1757. import_workflow(workflow, f.read())
  1758. f.close()
  1759. workflow.save()
  1760. assert_equal(5, len(Node.objects.filter(workflow=workflow)))
  1761. assert_equal(6, len(Link.objects.filter(parent__workflow=workflow)))
  1762. nodes = [Node.objects.filter(workflow=workflow, node_type='java')[0].get_full_node(),
  1763. Node.objects.filter(workflow=workflow, node_type='java')[1].get_full_node()]
  1764. assert_equal('org.apache.hadoop.examples.terasort.TeraGen', nodes[0].main_class)
  1765. assert_equal('${records} ${output_dir}/teragen', nodes[0].args)
  1766. assert_equal('org.apache.hadoop.examples.terasort.TeraSort', nodes[1].main_class)
  1767. assert_equal('${output_dir}/teragen ${output_dir}/terasort', nodes[1].args)
  1768. assert_true(nodes[0].capture_output)
  1769. assert_false(nodes[1].capture_output)
  1770. workflow.delete(skip_trash=True)
  1771. def test_import_workflow_shell(self):
  1772. workflow = Workflow.objects.new_workflow(self.user)
  1773. workflow.save()
  1774. f = open('apps/oozie/src/oozie/test_data/workflows/0.4/test-shell.xml')
  1775. import_workflow(workflow, f.read())
  1776. f.close()
  1777. workflow.save()
  1778. assert_equal(5, len(Node.objects.filter(workflow=workflow)))
  1779. assert_equal(6, len(Link.objects.filter(parent__workflow=workflow)))
  1780. nodes = [Node.objects.filter(workflow=workflow, node_type='shell')[0].get_full_node(),
  1781. Node.objects.filter(workflow=workflow, node_type='shell')[1].get_full_node()]
  1782. assert_equal('shell-1', nodes[0].name)
  1783. assert_equal('shell-2', nodes[1].name)
  1784. assert_equal('my-job.xml', nodes[0].job_xml)
  1785. assert_equal('hello.py', nodes[0].command)
  1786. assert_equal('[{"type":"argument","value":"World!"}]', nodes[0].params)
  1787. assert_true(nodes[0].capture_output)
  1788. assert_false(nodes[1].capture_output)
  1789. workflow.delete(skip_trash=True)
  1790. def test_import_workflow_fs(self):
  1791. """
  1792. Validates import for fs node: chmods, deletes, mkdirs, moves, touchzs.
  1793. """
  1794. workflow = Workflow.objects.new_workflow(self.user)
  1795. workflow.save()
  1796. f = open('apps/oozie/src/oozie/test_data/workflows/0.4/test-fs.xml')
  1797. import_workflow(workflow, f.read())
  1798. f.close()
  1799. workflow.save()
  1800. assert_equal(4, len(Node.objects.filter(workflow=workflow)))
  1801. assert_equal(4, len(Link.objects.filter(parent__workflow=workflow)))
  1802. node = Node.objects.get(workflow=workflow, node_type='fs').get_full_node()
  1803. assert_equal('[{"path":"${nameNode}${output}/testfs/renamed","permissions":"700","recursive":"false"}]', node.chmods)
  1804. assert_equal('[{"name":"${nameNode}${output}/testfs"}]', node.deletes)
  1805. assert_equal('["${nameNode}${output}/testfs","${nameNode}${output}/testfs/source"]', node.mkdirs)
  1806. assert_equal('[{"source":"${nameNode}${output}/testfs/source","destination":"${nameNode}${output}/testfs/renamed"}]', node.moves)
  1807. assert_equal('["${nameNode}${output}/testfs/new_file"]', node.touchzs)
  1808. workflow.delete(skip_trash=True)
  1809. def test_import_workflow_email(self):
  1810. """
  1811. Validates import for email node: to, css, subject, body.
  1812. """
  1813. workflow = Workflow.objects.new_workflow(self.user)
  1814. workflow.save()
  1815. f = open('apps/oozie/src/oozie/test_data/workflows/0.4/test-email.0.1.xml')
  1816. import_workflow(workflow, f.read())
  1817. f.close()
  1818. workflow.save()
  1819. assert_equal(4, len(Node.objects.filter(workflow=workflow)))
  1820. assert_equal(4, len(Link.objects.filter(parent__workflow=workflow)))
  1821. node = Node.objects.get(workflow=workflow, node_type='email').get_full_node()
  1822. assert_equal('example@example.org', node.to)
  1823. assert_equal('', node.cc)
  1824. assert_equal('I love', node.subject)
  1825. assert_equal('Hue', node.body)
  1826. workflow.delete(skip_trash=True)
  1827. def test_import_workflow_generic(self):
  1828. """
  1829. Validates import for generic node: xml.
  1830. """
  1831. workflow = Workflow.objects.new_workflow(self.user)
  1832. workflow.save()
  1833. f = open('apps/oozie/src/oozie/test_data/workflows/0.4/test-generic.xml')
  1834. import_workflow(workflow, f.read())
  1835. f.close()
  1836. workflow.save()
  1837. assert_equal(4, len(Node.objects.filter(workflow=workflow)))
  1838. assert_equal(4, len(Link.objects.filter(parent__workflow=workflow)))
  1839. node = Node.objects.get(workflow=workflow, node_type='generic').get_full_node()
  1840. assert_equal("<bleh test=\"test\">\n <test>test</test>\n </bleh>", node.xml)
  1841. workflow.delete(skip_trash=True)
  1842. def test_import_workflow_multi_kill_node(self):
  1843. """
  1844. Validates import for multiple kill nodes: xml.
  1845. Kill nodes should be skipped and a single kill node should be created.
  1846. """
  1847. workflow = Workflow.objects.new_workflow(self.user)
  1848. workflow.save()
  1849. f = open('apps/oozie/src/oozie/test_data/workflows/0.4/test-java-multiple-kill.xml')
  1850. import_workflow(workflow, f.read())
  1851. f.close()
  1852. workflow.save()
  1853. assert_equal('kill', Kill.objects.get(workflow=workflow).name)
  1854. assert_equal(5, len(Node.objects.filter(workflow=workflow)))
  1855. assert_equal(6, len(Link.objects.filter(parent__workflow=workflow)))
  1856. nodes = [Node.objects.filter(workflow=workflow, node_type='java')[0].get_full_node(),
  1857. Node.objects.filter(workflow=workflow, node_type='java')[1].get_full_node()]
  1858. assert_equal('org.apache.hadoop.examples.terasort.TeraGen', nodes[0].main_class)
  1859. assert_equal('${records} ${output_dir}/teragen', nodes[0].args)
  1860. assert_equal('org.apache.hadoop.examples.terasort.TeraSort', nodes[1].main_class)
  1861. assert_equal('${output_dir}/teragen ${output_dir}/terasort', nodes[1].args)
  1862. assert_true(nodes[0].capture_output)
  1863. assert_false(nodes[1].capture_output)
  1864. workflow.delete(skip_trash=True)
  1865. def test_import_workflow_different_error_link(self):
  1866. """
  1867. Validates import with error link to end: main_class, args.
  1868. If an error link cannot be resolved, default to 'kill' node.
  1869. """
  1870. workflow = Workflow.objects.new_workflow(self.user)
  1871. workflow.save()
  1872. f = open('apps/oozie/src/oozie/test_data/workflows/0.4/test-java-different-error-links.xml')
  1873. import_workflow(workflow, f.read())
  1874. f.close()
  1875. workflow.save()
  1876. assert_equal(5, len(Node.objects.filter(workflow=workflow)))
  1877. assert_equal(6, len(Link.objects.filter(parent__workflow=workflow)))
  1878. nodes = [Node.objects.filter(workflow=workflow, node_type='java')[0].get_full_node(),
  1879. Node.objects.filter(workflow=workflow, node_type='java')[1].get_full_node()]
  1880. assert_equal('org.apache.hadoop.examples.terasort.TeraGen', nodes[0].main_class)
  1881. assert_equal('${records} ${output_dir}/teragen', nodes[0].args)
  1882. assert_equal('org.apache.hadoop.examples.terasort.TeraSort', nodes[1].main_class)
  1883. assert_equal('${output_dir}/teragen ${output_dir}/terasort', nodes[1].args)
  1884. assert_true(nodes[0].capture_output)
  1885. assert_false(nodes[1].capture_output)
  1886. assert_equal(1, len(Link.objects.filter(parent__workflow=workflow).filter(parent__name='TeraGenWorkflow').filter(name='error').filter(child__node_type='java')))
  1887. assert_equal(1, len(Link.objects.filter(parent__workflow=workflow).filter(parent__name='TeraSort').filter(name='error').filter(child__node_type='kill')))
  1888. workflow.delete(skip_trash=True)
  1889. class TestImportCoordinator02(OozieMockBase):
  1890. def setUp(self):
  1891. super(TestImportCoordinator02, self).setUp()
  1892. self.setup_simple_workflow()
  1893. def test_import_coordinator_simple(self):
  1894. coordinator_count = Document.objects.available_docs(Coordinator, self.user).count()
  1895. # Create
  1896. filename = os.path.abspath(os.path.dirname(__file__) + "/test_data/coordinators/0.2/test-basic.xml")
  1897. fh = open(filename)
  1898. response = self.c.post(reverse('oozie:import_coordinator'), {
  1899. 'name': ['test_coordinator'],
  1900. 'workflow': Workflow.objects.get(name='wf-name-1').pk,
  1901. 'definition_file': [fh],
  1902. 'description': ['test description']
  1903. }, follow=True)
  1904. fh.close()
  1905. assert_equal(coordinator_count + 1, Document.objects.available_docs(Coordinator, self.user).count(), response)
  1906. coordinator = Coordinator.objects.get(name='test_coordinator')
  1907. assert_equal('[{"name":"oozie.use.system.libpath","value":"true"}]', coordinator.parameters)
  1908. assert_equal('uri:oozie:coordinator:0.2', coordinator.schema_version)
  1909. assert_equal('test description', coordinator.description)
  1910. assert_equal(datetime.strptime('2013-06-03T00:00Z', '%Y-%m-%dT%H:%MZ'), coordinator.start)
  1911. assert_equal(datetime.strptime('2013-06-05T00:00Z', '%Y-%m-%dT%H:%MZ'), coordinator.end)
  1912. assert_equal('America/Los_Angeles', coordinator.timezone)
  1913. assert_equal('days', coordinator.frequency_unit)
  1914. assert_equal(1, coordinator.frequency_number)
  1915. assert_equal(None, coordinator.timeout)
  1916. assert_equal(None, coordinator.concurrency)
  1917. assert_equal(None, coordinator.execution)
  1918. assert_equal(None, coordinator.throttle)
  1919. assert_not_equal(None, coordinator.deployment_dir)
  1920. class TestPermissions(OozieBase):
  1921. def setUp(self):
  1922. super(TestPermissions, self).setUp()
  1923. self.wf = create_workflow(self.c, self.user)
  1924. self.setup_simple_workflow()
  1925. def tearDown(self):
  1926. try:
  1927. self.wf.delete(skip_trash=True)
  1928. except:
  1929. pass
  1930. def test_workflow_permissions(self):
  1931. raise SkipTest
  1932. response = self.c.get(reverse('oozie:edit_workflow', args=[self.wf.id]))
  1933. assert_true('Editor' in response.content, response.content)
  1934. assert_true('Save' in response.content, response.content)
  1935. assert_false(self.wf.is_shared)
  1936. # Login as someone else
  1937. client_not_me = make_logged_in_client(username='not_me', is_superuser=False, groupname='test')
  1938. grant_access("not_me", "test", "oozie")
  1939. # List
  1940. finish = SHARE_JOBS.set_for_testing(True)
  1941. try:
  1942. response = client_not_me.get(reverse('oozie:list_workflows'))
  1943. assert_false('wf-name-1' in response.content, response.content)
  1944. finally:
  1945. finish()
  1946. finish = SHARE_JOBS.set_for_testing(False)
  1947. try:
  1948. response = client_not_me.get(reverse('oozie:list_workflows'))
  1949. assert_false('wf-name-1' in response.content, response.content)
  1950. finally:
  1951. finish()
  1952. # View
  1953. finish = SHARE_JOBS.set_for_testing(True)
  1954. try:
  1955. response = client_not_me.get(reverse('oozie:edit_workflow', args=[self.wf.id]))
  1956. assert_true('Permission denied' in response.content, response.content)
  1957. finally:
  1958. finish()
  1959. finish = SHARE_JOBS.set_for_testing(False)
  1960. try:
  1961. response = client_not_me.get(reverse('oozie:edit_workflow', args=[self.wf.id]))
  1962. assert_true('Permission denied' in response.content, response.content)
  1963. finally:
  1964. finish()
  1965. # Share it !
  1966. self.wf = Workflow.objects.get(name='wf-name-1', managed=True)
  1967. self.wf.is_shared = True
  1968. self.wf.save()
  1969. Workflow.objects.check_workspace(self.wf, self.cluster.fs)
  1970. # List
  1971. finish = SHARE_JOBS.set_for_testing(True)
  1972. try:
  1973. response = client_not_me.get(reverse('oozie:list_workflows'))
  1974. assert_equal(200, response.status_code)
  1975. assert_true('wf-name-1' in response.content, response.content)
  1976. finally:
  1977. finish()
  1978. # View
  1979. finish = SHARE_JOBS.set_for_testing(True)
  1980. try:
  1981. response = client_not_me.get(reverse('oozie:edit_workflow', args=[self.wf.id]))
  1982. assert_false('Permission denied' in response.content, response.content)
  1983. finally:
  1984. finish()
  1985. finish = SHARE_JOBS.set_for_testing(False)
  1986. try:
  1987. response = client_not_me.get(reverse('oozie:edit_workflow', args=[self.wf.id]))
  1988. assert_true('Permission denied' in response.content, response.content)
  1989. finally:
  1990. finish()
  1991. # Submit
  1992. finish = SHARE_JOBS.set_for_testing(False)
  1993. try:
  1994. response = client_not_me.post(reverse('oozie:submit_workflow', args=[self.wf.id]))
  1995. assert_true('Permission denied' in response.content, response.content)
  1996. finally:
  1997. finish()
  1998. finish = SHARE_JOBS.set_for_testing(True)
  1999. try:
  2000. try:
  2001. response = client_not_me.post(reverse('oozie:submit_workflow', args=[self.wf.id]))
  2002. assert_false('Permission denied' in response.content, response.content)
  2003. except IOError:
  2004. pass
  2005. finally:
  2006. finish()
  2007. # Move to trash
  2008. finish = SHARE_JOBS.set_for_testing(False)
  2009. try:
  2010. response = client_not_me.post(reverse('oozie:delete_workflow'), {'job_selection': [self.wf.id]})
  2011. assert_true('Permission denied' in response.content, response.content)
  2012. finally:
  2013. finish()
  2014. response = self.c.post(reverse('oozie:delete_workflow'), {'job_selection': [self.wf.id]}, follow=True)
  2015. assert_equal(200, response.status_code)
  2016. # Trash
  2017. finish = SHARE_JOBS.set_for_testing(False)
  2018. try:
  2019. response = client_not_me.get(reverse('oozie:list_trashed_workflows'))
  2020. assert_false(self.wf.name in response.content, response.content)
  2021. finally:
  2022. finish()
  2023. response = self.c.get(reverse('oozie:list_trashed_workflows'))
  2024. assert_true(self.wf.name in response.content, response.content)
  2025. # Restore
  2026. finish = SHARE_JOBS.set_for_testing(False)
  2027. try:
  2028. response = client_not_me.post(reverse('oozie:restore_workflow'), {'job_selection': [self.wf.id]})
  2029. assert_true('Permission denied' in response.content, response.content)
  2030. finally:
  2031. finish()
  2032. response = self.c.post(reverse('oozie:restore_workflow'), {'job_selection': [self.wf.id]}, follow=True)
  2033. assert_equal(200, response.status_code)
  2034. def test_coordinator_permissions(self):
  2035. raise SkipTest
  2036. coord = create_coordinator(self.wf, self.c, self.user)
  2037. response = self.c.get(reverse('oozie:edit_coordinator', args=[coord.id]))
  2038. assert_true('Editor' in response.content, response.content)
  2039. assert_true('Save coordinator' in response.content, response.content)
  2040. # Login as someone else
  2041. client_not_me = make_logged_in_client(username='not_me', is_superuser=False, groupname='test')
  2042. grant_access("not_me", "test", "oozie")
  2043. # List
  2044. finish = SHARE_JOBS.set_for_testing(True)
  2045. try:
  2046. response = client_not_me.get(reverse('oozie:list_coordinators'))
  2047. assert_false('MyCoord' in response.content, response.content)
  2048. finally:
  2049. finish()
  2050. finish = SHARE_JOBS.set_for_testing(False)
  2051. try:
  2052. response = client_not_me.get(reverse('oozie:list_coordinators'))
  2053. assert_false('MyCoord' in response.content, response.content)
  2054. finally:
  2055. finish()
  2056. # View
  2057. finish = SHARE_JOBS.set_for_testing(True)
  2058. try:
  2059. response = client_not_me.get(reverse('oozie:edit_coordinator', args=[coord.id]))
  2060. assert_true('Permission denied' in response.content, response.content)
  2061. finally:
  2062. finish()
  2063. finish = SHARE_JOBS.set_for_testing(False)
  2064. try:
  2065. response = client_not_me.get(reverse('oozie:edit_coordinator', args=[coord.id]))
  2066. assert_false('MyCoord' in response.content, response.content)
  2067. finally:
  2068. finish()
  2069. # Share it !
  2070. wf = Workflow.objects.get(id=coord.workflow.id, managed=True)
  2071. wf.is_shared = True
  2072. wf.save()
  2073. Workflow.objects.check_workspace(wf, self.cluster.fs)
  2074. post = COORDINATOR_DICT.copy()
  2075. post.update({
  2076. u'datainput_set-TOTAL_FORMS': [u'0'], u'datainput_set-INITIAL_FORMS': [u'0'], u'dataset_set-INITIAL_FORMS': [u'0'],
  2077. u'dataoutput_set-INITIAL_FORMS': [u'0'], u'datainput_set-MAX_NUM_FORMS': [u'0'], u'output-MAX_NUM_FORMS': [u''],
  2078. u'output-INITIAL_FORMS': [u'0'], u'dataoutput_set-TOTAL_FORMS': [u'0'], u'input-TOTAL_FORMS': [u'0'],
  2079. u'dataset_set-MAX_NUM_FORMS': [u'0'], u'dataoutput_set-MAX_NUM_FORMS': [u'0'], u'input-MAX_NUM_FORMS': [u''],
  2080. u'dataset_set-TOTAL_FORMS': [u'0'], u'input-INITIAL_FORMS': [u'0'], u'output-TOTAL_FORMS': [u'0']})
  2081. post['is_shared'] = [u'on']
  2082. post['workflow'] = coord.workflow.id
  2083. self.c.post(reverse('oozie:edit_coordinator', args=[coord.id]), post)
  2084. coord = Coordinator.objects.get(id=coord.id)
  2085. assert_true(coord.is_shared)
  2086. # List
  2087. finish = SHARE_JOBS.set_for_testing(True)
  2088. try:
  2089. response = client_not_me.get(reverse('oozie:list_coordinators'))
  2090. assert_equal(200, response.status_code)
  2091. assert_true('MyCoord' in response.content, response.content)
  2092. finally:
  2093. finish()
  2094. # View
  2095. finish = SHARE_JOBS.set_for_testing(True)
  2096. try:
  2097. response = client_not_me.get(reverse('oozie:edit_coordinator', args=[coord.id]))
  2098. assert_false('Permission denied' in response.content, response.content)
  2099. assert_false('Save coordinator' in response.content, response.content)
  2100. finally:
  2101. finish()
  2102. finish = SHARE_JOBS.set_for_testing(False)
  2103. try:
  2104. response = client_not_me.get(reverse('oozie:edit_coordinator', args=[coord.id]))
  2105. assert_true('Permission denied' in response.content, response.content)
  2106. finally:
  2107. finish()
  2108. # Edit
  2109. finish = SHARE_JOBS.set_for_testing(True)
  2110. try:
  2111. response = client_not_me.post(reverse('oozie:edit_coordinator', args=[coord.id]))
  2112. assert_false('MyCoord' in response.content, response.content)
  2113. assert_true('Not allowed' in response.content, response.content)
  2114. finally:
  2115. finish()
  2116. # Submit
  2117. finish = SHARE_JOBS.set_for_testing(False)
  2118. try:
  2119. response = client_not_me.post(reverse('oozie:submit_coordinator', args=[coord.id]))
  2120. assert_true('Permission denied' in response.content, response.content)
  2121. finally:
  2122. finish()
  2123. finish = SHARE_JOBS.set_for_testing(True)
  2124. try:
  2125. try:
  2126. response = client_not_me.post(reverse('oozie:submit_coordinator', args=[coord.id]))
  2127. assert_false('Permission denied' in response.content, response.content)
  2128. except IOError:
  2129. pass
  2130. finally:
  2131. finish()
  2132. # Move to trash
  2133. finish = SHARE_JOBS.set_for_testing(False)
  2134. try:
  2135. response = client_not_me.post(reverse('oozie:delete_coordinator'), {'job_selection': [coord.id]})
  2136. assert_true('Permission denied' in response.content, response.content)
  2137. finally:
  2138. finish()
  2139. response = self.c.post(reverse('oozie:delete_coordinator'), {'job_selection': [coord.id]}, follow=True)
  2140. assert_equal(200, response.status_code)
  2141. # List trash
  2142. finish = SHARE_JOBS.set_for_testing(True)
  2143. try:
  2144. response = client_not_me.get(reverse('oozie:list_trashed_coordinators'))
  2145. assert_true(coord.name in response.content, response.content)
  2146. finally:
  2147. finish()
  2148. finish = SHARE_JOBS.set_for_testing(False)
  2149. response = client_not_me.get(reverse('oozie:list_trashed_coordinators'))
  2150. assert_false(coord.name in response.content, response.content)
  2151. # Restore
  2152. finish = SHARE_JOBS.set_for_testing(False)
  2153. try:
  2154. response = client_not_me.post(reverse('oozie:restore_coordinator'), {'job_selection': [coord.id]})
  2155. assert_true('Permission denied' in response.content, response.content)
  2156. finally:
  2157. finish()
  2158. response = self.c.post(reverse('oozie:restore_coordinator'), {'job_selection': [coord.id]}, follow=True)
  2159. assert_equal(200, response.status_code)
  2160. def test_bundle_permissions(self):
  2161. raise SkipTest
  2162. bundle = create_bundle(self.c, self.user)
  2163. response = self.c.get(reverse('oozie:edit_bundle', args=[bundle.id]))
  2164. assert_true('Editor' in response.content, response.content)
  2165. assert_true('MyBundle' in response.content, response.content)
  2166. assert_true('Save' in response.content, response.content)
  2167. assert_false(bundle.is_shared)
  2168. # Login as someone else
  2169. client_not_me = make_logged_in_client(username='not_me', is_superuser=False, groupname='test')
  2170. grant_access("not_me", "test", "oozie")
  2171. # List
  2172. finish = SHARE_JOBS.set_for_testing(True)
  2173. try:
  2174. response = client_not_me.get(reverse('oozie:list_bundles'))
  2175. assert_false('MyBundle' in response.content, response.content)
  2176. finally:
  2177. finish()
  2178. finish = SHARE_JOBS.set_for_testing(False)
  2179. try:
  2180. response = client_not_me.get(reverse('oozie:list_bundles'))
  2181. assert_false('MyBundle' in response.content, response.content)
  2182. finally:
  2183. finish()
  2184. # View
  2185. finish = SHARE_JOBS.set_for_testing(True)
  2186. try:
  2187. response = client_not_me.get(reverse('oozie:edit_bundle', args=[bundle.id]))
  2188. assert_true('Permission denied' in response.content, response.content)
  2189. finally:
  2190. finish()
  2191. finish = SHARE_JOBS.set_for_testing(False)
  2192. try:
  2193. response = client_not_me.get(reverse('oozie:edit_bundle', args=[bundle.id]))
  2194. assert_true('Permission denied' in response.content, response.content)
  2195. finally:
  2196. finish()
  2197. # Share it !
  2198. bundle.is_shared = True
  2199. bundle.save()
  2200. # List
  2201. finish = SHARE_JOBS.set_for_testing(True)
  2202. try:
  2203. response = client_not_me.get(reverse('oozie:list_bundles'))
  2204. assert_equal(200, response.status_code)
  2205. assert_true('MyBundle' in response.content, response.content)
  2206. finally:
  2207. finish()
  2208. # View
  2209. finish = SHARE_JOBS.set_for_testing(True)
  2210. try:
  2211. response = client_not_me.get(reverse('oozie:edit_bundle', args=[bundle.id]))
  2212. assert_false('Permission denied' in response.content, response.content)
  2213. finally:
  2214. finish()
  2215. finish = SHARE_JOBS.set_for_testing(False)
  2216. try:
  2217. response = client_not_me.get(reverse('oozie:edit_bundle', args=[bundle.id]))
  2218. assert_true('Permission denied' in response.content, response.content)
  2219. finally:
  2220. finish()
  2221. # Submit
  2222. finish = SHARE_JOBS.set_for_testing(False)
  2223. try:
  2224. response = client_not_me.post(reverse('oozie:submit_bundle', args=[bundle.id]),{
  2225. u'form-MAX_NUM_FORMS': [u''], u'form-INITIAL_FORMS': [u'0'], u'form-TOTAL_FORMS': [u'0']
  2226. })
  2227. assert_true('Permission denied' in response.content, response.content)
  2228. finally:
  2229. finish()
  2230. finish = SHARE_JOBS.set_for_testing(True)
  2231. try:
  2232. try:
  2233. response = client_not_me.post(reverse('oozie:submit_bundle', args=[bundle.id]), {
  2234. u'form-MAX_NUM_FORMS': [u''], u'form-INITIAL_FORMS': [u'0'], u'form-TOTAL_FORMS': [u'0']
  2235. })
  2236. assert_false('Permission denied' in response.content, response.content)
  2237. except IOError:
  2238. pass
  2239. finally:
  2240. finish()
  2241. # Delete
  2242. finish = SHARE_JOBS.set_for_testing(False)
  2243. try:
  2244. response = client_not_me.post(reverse('oozie:delete_bundle'), {'job_selection': [bundle.id]})
  2245. assert_true('Permission denied' in response.content, response.content)
  2246. finally:
  2247. finish()
  2248. response = self.c.post(reverse('oozie:delete_bundle'), {'job_selection': [bundle.id]}, follow=True)
  2249. assert_equal(200, response.status_code)
  2250. # List trash
  2251. finish = SHARE_JOBS.set_for_testing(True)
  2252. try:
  2253. response = client_not_me.get(reverse('oozie:list_trashed_bundles'))
  2254. assert_true(bundle.name in response.content, response.content)
  2255. finally:
  2256. finish()
  2257. finish = SHARE_JOBS.set_for_testing(False)
  2258. response = client_not_me.get(reverse('oozie:list_trashed_bundles'))
  2259. assert_false(bundle.name in response.content, response.content)
  2260. # Restore
  2261. finish = SHARE_JOBS.set_for_testing(False)
  2262. try:
  2263. response = client_not_me.post(reverse('oozie:restore_bundle'), {'job_selection': [bundle.id]})
  2264. assert_true('Permission denied' in response.content, response.content)
  2265. finally:
  2266. finish()
  2267. response = self.c.post(reverse('oozie:restore_bundle'), {'job_selection': [bundle.id]}, follow=True)
  2268. assert_equal(200, response.status_code)
  2269. class TestEditorWithOozie(OozieBase):
  2270. def setUp(self):
  2271. OozieBase.setUp(self)
  2272. self.c = make_logged_in_client()
  2273. self.wf = create_workflow(self.c, self.user)
  2274. self.setup_simple_workflow()
  2275. def tearDown(self):
  2276. try:
  2277. self.wf.delete(skip_trash=True)
  2278. except:
  2279. pass
  2280. def test_create_workflow(self):
  2281. dir_stat = self.cluster.fs.stats(self.wf.deployment_dir)
  2282. assert_equal('test', dir_stat.user)
  2283. assert_equal('hue', dir_stat.group)
  2284. assert_equal('40711', '%o' % dir_stat.mode)
  2285. def test_clone_workflow(self):
  2286. workflow_count = Document.objects.available_docs(Workflow, self.user).count()
  2287. response = self.c.post(reverse('oozie:clone_workflow', args=[self.wf.id]), {}, follow=True)
  2288. assert_equal(workflow_count + 1, Document.objects.available_docs(Workflow, self.user).count(), response)
  2289. wf2 = Workflow.objects.latest('id')
  2290. assert_not_equal(self.wf.id, wf2.id)
  2291. assert_equal(self.wf.node_set.count(), wf2.node_set.count())
  2292. node_ids = set(self.wf.node_set.values_list('id', flat=True))
  2293. for node in wf2.node_set.all():
  2294. assert_false(node.id in node_ids)
  2295. assert_not_equal(self.wf.deployment_dir, wf2.deployment_dir)
  2296. assert_not_equal('', wf2.deployment_dir)
  2297. # Bulk delete
  2298. response = self.c.post(reverse('oozie:delete_workflow'), {'job_selection': [self.wf.id, wf2.id]}, follow=True)
  2299. assert_equal(workflow_count - 1, Document.objects.available_docs(Workflow, self.user).count(), response)
  2300. def test_import_workflow(self):
  2301. workflow_count = Document.objects.available_docs(Workflow, self.user).count()
  2302. # Create
  2303. filename = os.path.abspath(os.path.dirname(__file__) + "/test_data/workflows/0.4/test-mapreduce.xml")
  2304. fh = open(filename)
  2305. response = self.c.post(reverse('oozie:import_workflow'), {
  2306. 'job_xml': [''],
  2307. 'name': ['test_workflow'],
  2308. 'parameters': ['[{"name":"oozie.use.system.libpath","value":"true"}]'],
  2309. 'deployment_dir': [''],
  2310. 'job_properties': ['[]'],
  2311. 'schema_version': ['0.4'],
  2312. 'definition_file': [fh],
  2313. 'description': ['']
  2314. }, follow=True)
  2315. fh.close()
  2316. assert_equal(workflow_count + 1, Document.objects.available_docs(Workflow, self.user).count(), response)
  2317. def test_delete_workflow(self):
  2318. previous_trashed = Document.objects.trashed_docs(Workflow, self.user).count()
  2319. previous_available = Document.objects.available_docs(Workflow, self.user).count()
  2320. response = self.c.post(reverse('oozie:delete_workflow') + "?skip_trash=true", {'job_selection': [self.wf.id]}, follow=True)
  2321. assert_equal(200, response.status_code, response)
  2322. assert_equal(previous_trashed, Document.objects.trashed_docs(Workflow, self.user).count())
  2323. assert_equal(previous_available - 1, Document.objects.available_docs(Workflow, self.user).count())
  2324. class TestImportWorkflow04WithOozie(OozieBase):
  2325. def setUp(self):
  2326. OozieBase.setUp(self)
  2327. self.c = make_logged_in_client()
  2328. self.wf = create_workflow(self.c, self.user)
  2329. self.setup_simple_workflow()
  2330. # in order to reference examples in Subworkflow, must be owned by current user.
  2331. Workflow.objects.update(owner=self.user)
  2332. def tearDown(self):
  2333. self.wf.delete(skip_trash=True)
  2334. def test_import_workflow_subworkflow(self):
  2335. """
  2336. Validates import for subworkflow node: propagate_configuration.
  2337. """
  2338. workflow = Workflow.objects.new_workflow(self.user)
  2339. workflow.save()
  2340. f = open('apps/oozie/src/oozie/test_data/workflows/0.4/test-subworkflow.xml')
  2341. import_workflow(workflow, f.read(), None, self.cluster.fs)
  2342. f.close()
  2343. workflow.save()
  2344. assert_equal(4, len(Node.objects.filter(workflow=workflow)))
  2345. assert_equal(4, len(Link.objects.filter(parent__workflow=workflow)))
  2346. node = Node.objects.get(workflow=workflow, node_type='subworkflow').get_full_node()
  2347. assert_equal(True, node.propagate_configuration)
  2348. workflow.delete(skip_trash=True)
  2349. class TestOozieSubmissions(OozieBase):
  2350. def test_submit_mapreduce_action(self):
  2351. wf = Document.objects.get_docs(self.user, Workflow).get(name='MapReduce', owner__username='sample', extra='').content_object
  2352. wf.owner = User.objects.get(username='sample')
  2353. wf.save()
  2354. post_data = {u'form-MAX_NUM_FORMS': [u''], u'form-INITIAL_FORMS': [u'1'],
  2355. u'form-0-name': [u'REDUCER_SLEEP_TIME'], u'form-0-value': [u'1'],
  2356. u'form-TOTAL_FORMS': [u'1']}
  2357. assert_equal('sample', wf.owner.username)
  2358. response = self.c.post(reverse('oozie:submit_workflow', args=[wf.id]), data=post_data, follow=True)
  2359. job = OozieServerProvider.wait_until_completion(response.context['oozie_workflow'].id)
  2360. assert_equal('SUCCEEDED', job.status)
  2361. assert_equal(100, job.get_progress())
  2362. # Rerun with default options
  2363. post_data.update({u'rerun_form_choice': [u'skip_nodes']})
  2364. response = self.c.post(reverse('oozie:rerun_oozie_job', kwargs={'job_id': job.id, 'app_path': job.appPath}), data=post_data, follow=True)
  2365. job = OozieServerProvider.wait_until_completion(response.context['oozie_workflow'].id)
  2366. assert_equal('SUCCEEDED', job.status)
  2367. assert_equal(100, job.get_progress())
  2368. # Rerun with skip OK actions skipped
  2369. post_data.update({u'rerun_form_choice': [u'skip_nodes'], u'skip_nodes': [u'Sleep']})
  2370. response = self.c.post(reverse('oozie:rerun_oozie_job', kwargs={'job_id': job.id, 'app_path': job.appPath}), data=post_data, follow=True)
  2371. job = OozieServerProvider.wait_until_completion(response.context['oozie_workflow'].id)
  2372. assert_equal('SUCCEEDED', job.status)
  2373. assert_equal(100, job.get_progress())
  2374. # Rerun with failed nodes too
  2375. post_data.update({u'rerun_form_choice': [u'failed_nodes']})
  2376. response = self.c.post(reverse('oozie:rerun_oozie_job', kwargs={'job_id': job.id, 'app_path': job.appPath}), data=post_data, follow=True)
  2377. job = OozieServerProvider.wait_until_completion(response.context['oozie_workflow'].id)
  2378. def test_submit_java_action(self):
  2379. wf = Document.objects.get_docs(self.user, Workflow).get(name='Sequential Java', owner__username='sample', extra='').content_object
  2380. wf.owner = User.objects.get(username='sample')
  2381. wf.save()
  2382. response = self.c.post(reverse('oozie:submit_workflow', args=[wf.id]),
  2383. data={u'form-MAX_NUM_FORMS': [u''],
  2384. u'form-0-name': [u'records'], u'form-0-value': [u'10'],
  2385. u'form-1-name': [u' output_dir '], u'form-1-value': [u'${nameNode}/user/test/out/terasort'],
  2386. u'form-INITIAL_FORMS': [u'2'], u'form-TOTAL_FORMS': [u'2']},
  2387. follow=True)
  2388. job = OozieServerProvider.wait_until_completion(response.context['oozie_workflow'].id)
  2389. assert_equal('SUCCEEDED', job.status)
  2390. def test_submit_distcp_action(self):
  2391. wf = Document.objects.get_docs(self.user, Workflow).get(name='DistCp', owner__username='sample', extra='').content_object
  2392. wf.owner = User.objects.get(username='sample')
  2393. wf.save()
  2394. response = self.c.post(reverse('oozie:submit_workflow', args=[wf.id]),
  2395. data={
  2396. u'form-MAX_NUM_FORMS': [u''], u'form-TOTAL_FORMS': [u'3'], u'form-INITIAL_FORMS': [u'3'],
  2397. u'form-0-name': [u'oozie.use.system.libpath'], u'form-0-value': [u'true'],
  2398. u'form-1-name': [u'OUTPUT'], u'form-1-value': [u'${nameNode}/user/test/out/distcp'],
  2399. u'form-2-name': [u'MAP_NUMBER'], u'form-2-value': [u'5'],
  2400. },
  2401. follow=True)
  2402. job = OozieServerProvider.wait_until_completion(response.context['oozie_workflow'].id)
  2403. assert_equal('SUCCEEDED', job.status)
  2404. def test_oozie_page(self):
  2405. response = self.c.get(reverse('oozie:list_oozie_info'))
  2406. assert_true('version' in response.content, response.content)
  2407. assert_true('NORMAL' in response.content, response.content)
  2408. assert_true('variables' in response.content, response.content)
  2409. assert_true('timers' in response.content, response.content)
  2410. assert_true('counters' in response.content, response.content)
  2411. assert_true('ownMinTime' in response.content, response.content)
  2412. assert_true('oozie.base.url' in response.content, response.content)
  2413. class TestDashboardWithOozie(OozieBase):
  2414. def setUp(self):
  2415. super(TestDashboardWithOozie, self).setUp()
  2416. self.c = make_logged_in_client()
  2417. self.wf = create_workflow(self.c, self.user)
  2418. self.setup_simple_workflow()
  2419. def tearDown(self):
  2420. try:
  2421. self.wf.delete(skip_trash=True)
  2422. except:
  2423. pass
  2424. def test_submit_external_workflow(self):
  2425. # Check popup and reading workflow.xml and job.properties
  2426. oozie_xml = self.wf.to_xml({'output': '/path'})
  2427. deployment_dir = self.cluster.fs.mktemp(prefix='test_submit_external_workflow')
  2428. application_path = deployment_dir + '/workflow.xml'
  2429. self.cluster.fs.create(application_path, data=oozie_xml)
  2430. response = self.c.get(reverse('oozie:submit_external_job', kwargs={'application_path': application_path}))
  2431. assert_equal([{'name': 'SLEEP', 'value': ''}, {'name': 'output', 'value': ''}],
  2432. response.context['params_form'].initial)
  2433. oozie_properties = """
  2434. #
  2435. # Licensed to the Hue
  2436. #
  2437. nameNode=hdfs://localhost:8020
  2438. jobTracker=localhost:8021
  2439. my_prop_not_filtered=10
  2440. """
  2441. self.cluster.fs.create(deployment_dir + '/job.properties', data=oozie_properties)
  2442. response = self.c.get(reverse('oozie:submit_external_job', kwargs={'application_path': application_path}))
  2443. assert_equal([{'name': 'SLEEP', 'value': ''}, {'name': 'my_prop_not_filtered', 'value': '10'}, {'name': 'output', 'value': ''}],
  2444. response.context['params_form'].initial)
  2445. # Submit, just check if submittion worked
  2446. response = self.c.post(reverse('oozie:submit_external_job', kwargs={'application_path': application_path}), {
  2447. u'form-MAX_NUM_FORMS': [u''],
  2448. u'form-TOTAL_FORMS': [u'3'],
  2449. u'form-INITIAL_FORMS': [u'3'],
  2450. u'form-0-name': [u'SLEEP'],
  2451. u'form-0-value': [u'ilovesleep'],
  2452. u'form-1-name': [u'my_prop_not_filtered'],
  2453. u'form-1-value': [u'10'],
  2454. u'form-2-name': [u'output'],
  2455. u'form-2-value': [u'/path/output'],
  2456. }, follow=True)
  2457. assert_true(response.context['oozie_workflow'], response.content)
  2458. def test_oozie_not_running_message(self):
  2459. raise SkipTest # Not reseting the oozie url for some reason
  2460. finish = OOZIE_URL.set_for_testing('http://not_localhost:11000/bad')
  2461. try:
  2462. response = self.c.get(reverse('oozie:list_oozie_workflows'))
  2463. assert_true('The Oozie server is not running' in response.content, response.content)
  2464. finally:
  2465. finish()
  2466. class TestDashboard(OozieMockBase):
  2467. def test_manage_workflow_dashboard(self):
  2468. # Display of buttons happens in js now
  2469. response = self.c.get(reverse('oozie:list_oozie_workflow', args=[MockOozieApi.WORKFLOW_IDS[0]]), {}, follow=True)
  2470. assert_true(('%s/kill' % MockOozieApi.WORKFLOW_IDS[0]) in response.content, response.content)
  2471. assert_true(('rerun_oozie_job/%s' % MockOozieApi.WORKFLOW_IDS[0]) in response.content, response.content)
  2472. assert_true(('%s/suspend' % MockOozieApi.WORKFLOW_IDS[0]) in response.content, response.content)
  2473. assert_true(('%s/resume' % MockOozieApi.WORKFLOW_IDS[0]) in response.content, response.content)
  2474. response = self.c.get(reverse('oozie:list_oozie_workflow', args=[MockOozieApi.WORKFLOW_IDS[1]]), {}, follow=True)
  2475. assert_true(('%s/kill' % MockOozieApi.WORKFLOW_IDS[1]) in response.content, response.content)
  2476. assert_true(('rerun_oozie_job/%s' % MockOozieApi.WORKFLOW_IDS[1]) in response.content, response.content)
  2477. def test_manage_coordinator_dashboard(self):
  2478. # Display of buttons happens in js now
  2479. response = self.c.get(reverse('oozie:list_oozie_coordinator', args=[MockOozieApi.COORDINATOR_IDS[0]]), {}, follow=True)
  2480. assert_true(('%s/kill' % MockOozieApi.COORDINATOR_IDS[0]) in response.content, response.content)
  2481. assert_true(('rerun_oozie_coord/%s' % MockOozieApi.COORDINATOR_IDS[0]) in response.content, response.content)
  2482. assert_true(('%s/suspend' % MockOozieApi.COORDINATOR_IDS[0]) in response.content, response.content)
  2483. assert_true(('%s/resume' % MockOozieApi.COORDINATOR_IDS[0]) in response.content, response.content)
  2484. def test_manage_bundles_dashboard(self):
  2485. # Display of buttons happens in js now
  2486. response = self.c.get(reverse('oozie:list_oozie_bundle', args=[MockOozieApi.BUNDLE_IDS[0]]), {}, follow=True)
  2487. assert_true(('%s/kill' % MockOozieApi.BUNDLE_IDS[0]) in response.content, response.content)
  2488. assert_true(('rerun_oozie_bundle/%s' % MockOozieApi.BUNDLE_IDS[0]) in response.content, response.content)
  2489. assert_true(('%s/suspend' % MockOozieApi.BUNDLE_IDS[0]) in response.content, response.content)
  2490. assert_true(('%s/resume' % MockOozieApi.BUNDLE_IDS[0]) in response.content, response.content)
  2491. def test_rerun_coordinator(self):
  2492. response = self.c.get(reverse('oozie:rerun_oozie_coord', args=[MockOozieApi.WORKFLOW_IDS[0], '/path']))
  2493. assert_true('Select actions to rerun' in response.content, response.content)
  2494. def test_rerun_coordinator_permissions(self):
  2495. post_data = {
  2496. u'form-MAX_NUM_FORMS': [u''],
  2497. u'nocleanup': [u'on'],
  2498. u'refresh': [u'on'],
  2499. u'form-TOTAL_FORMS': [u'0'],
  2500. u'actions': [u'1'],
  2501. u'form-INITIAL_FORMS': [u'0']
  2502. }
  2503. response = self.c.post(reverse('oozie:rerun_oozie_coord', args=[MockOozieApi.COORDINATOR_IDS[0], '/path']), post_data)
  2504. assert_false('Permission denied' in response.content, response.content)
  2505. # Login as someone else
  2506. client_not_me = make_logged_in_client(username='not_me', is_superuser=False, groupname='test')
  2507. grant_access("not_me", "test", "oozie")
  2508. response = client_not_me.post(reverse('oozie:rerun_oozie_coord', args=[MockOozieApi.COORDINATOR_IDS[0], '/path']), post_data)
  2509. assert_true('Permission denied' in response.content, response.content)
  2510. def test_rerun_bundle(self):
  2511. response = self.c.get(reverse('oozie:rerun_oozie_coord', args=[MockOozieApi.WORKFLOW_IDS[0], '/path']))
  2512. assert_true('Select actions to rerun' in response.content, response.content)
  2513. def test_rerun_bundle_permissions(self):
  2514. post_data = {
  2515. u'end_1': [u'01:55 PM'],
  2516. u'end_0': [u'03/02/2013'],
  2517. u'refresh': [u'on'],
  2518. u'nocleanup': [u'on'],
  2519. u'start_0': [u'02/27/2013'],
  2520. u'start_1': [u'01:55 PM'],
  2521. u'coordinators': [u'DailySleep'],
  2522. u'form-MAX_NUM_FORMS': [u''],
  2523. u'form-TOTAL_FORMS': [u'3'],
  2524. u'form-INITIAL_FORMS': [u'3'],
  2525. u'form-0-name': [u'oozie.use.system.libpath'],
  2526. u'form-0-value': [u'true'],
  2527. u'form-1-name': [u'hue-id-b'],
  2528. u'form-1-value': [u'22'],
  2529. u'form-2-name': [u'oozie.bundle.application.path'],
  2530. u'form-2-value': [u'hdfs://localhost:8020/path'],
  2531. }
  2532. response = self.c.post(reverse('oozie:rerun_oozie_bundle', args=[MockOozieApi.BUNDLE_IDS[0], '/path']), post_data)
  2533. assert_false('Permission denied' in response.content, response.content)
  2534. # Login as someone else
  2535. client_not_me = make_logged_in_client(username='not_me', is_superuser=False, groupname='test')
  2536. grant_access("not_me", "test", "oozie")
  2537. response = client_not_me.post(reverse('oozie:rerun_oozie_bundle', args=[MockOozieApi.BUNDLE_IDS[0], '/path']), post_data)
  2538. assert_true('Permission denied' in response.content, response.content)
  2539. def test_list_workflows(self):
  2540. response = self.c.get(reverse('oozie:list_oozie_workflows') + "?format=json")
  2541. for wf_id in MockOozieApi.WORKFLOW_IDS:
  2542. assert_true(wf_id in response.content, response.content)
  2543. def test_list_coordinators(self):
  2544. response = self.c.get(reverse('oozie:list_oozie_coordinators') + "?format=json")
  2545. for coord_id in MockOozieApi.COORDINATOR_IDS:
  2546. assert_true(coord_id in response.content, response.content)
  2547. def test_list_bundles(self):
  2548. response = self.c.get(reverse('oozie:list_oozie_bundles') + "?format=json")
  2549. for coord_id in MockOozieApi.BUNDLE_IDS:
  2550. assert_true(coord_id in response.content, response.content)
  2551. def test_list_workflow(self):
  2552. response = self.c.get(reverse('oozie:list_oozie_workflow', args=[MockOozieApi.WORKFLOW_IDS[0]]))
  2553. assert_true('Workflow WordCount1' in response.content, response.content)
  2554. assert_true('Workflow' in response.content, response.content)
  2555. response = self.c.get(reverse('oozie:list_oozie_workflow', args=[MockOozieApi.WORKFLOW_IDS[0], MockOozieApi.COORDINATOR_IDS[0]]))
  2556. assert_true('Workflow WordCount1' in response.content, response.content)
  2557. assert_true('Workflow' in response.content, response.content)
  2558. assert_true('DailyWordCount1' in response.content, response.content)
  2559. assert_true('Coordinator' in response.content, response.content)
  2560. def test_list_workflow_action(self):
  2561. response = self.c.get(reverse('oozie:list_oozie_workflow_action', args=['XXX']))
  2562. assert_true('Action WordCount' in response.content, response.content)
  2563. assert_true('job_201302280955_0018' in response.content, response.content)
  2564. assert_true('job_201302280955_0019' in response.content, response.content)
  2565. assert_true('job_201302280955_0020' in response.content, response.content)
  2566. response = self.c.get(reverse('oozie:list_oozie_workflow_action', args=['XXX', MockOozieApi.COORDINATOR_IDS[0], MockOozieApi.BUNDLE_IDS[0]]))
  2567. assert_true('Bundle' in response.content, response.content)
  2568. assert_true('MyBundle1' in response.content, response.content)
  2569. assert_true('Coordinator' in response.content, response.content)
  2570. assert_true('DailyWordCount1' in response.content, response.content)
  2571. assert_true('Workflow' in response.content, response.content)
  2572. assert_true('WordCount1' in response.content, response.content)
  2573. def test_list_coordinator(self):
  2574. response = self.c.get(reverse('oozie:list_oozie_coordinator', args=[MockOozieApi.COORDINATOR_IDS[4]]))
  2575. assert_true(u'Coordinator DåilyWordCount5' in response.content.decode('utf-8', 'replace'), response.content.decode('utf-8', 'replace'))
  2576. assert_true('Workflow' in response.content, response.content)
  2577. def test_list_bundle(self):
  2578. response = self.c.get(reverse('oozie:list_oozie_bundle', args=[MockOozieApi.BUNDLE_IDS[0]]))
  2579. assert_true('Bundle MyBundle1' in response.content, response.content)
  2580. assert_true('Coordinators' in response.content, response.content)
  2581. def test_manage_oozie_jobs(self):
  2582. try:
  2583. self.c.get(reverse('oozie:manage_oozie_jobs', args=[MockOozieApi.COORDINATOR_IDS[0], 'kill']))
  2584. assert False
  2585. except:
  2586. pass
  2587. response = self.c.post(reverse('oozie:manage_oozie_jobs', args=[MockOozieApi.COORDINATOR_IDS[0], 'kill']))
  2588. data = json.loads(response.content)
  2589. assert_equal(0, data['status'])
  2590. def test_workflows_permissions(self):
  2591. response = self.c.get(reverse('oozie:list_oozie_workflows')+"?format=json")
  2592. assert_true('WordCount1' in response.content, response.content)
  2593. # Rerun
  2594. response = self.c.get(reverse('oozie:rerun_oozie_job', kwargs={'job_id': MockOozieApi.WORKFLOW_IDS[0],
  2595. 'app_path': MockOozieApi.JSON_WORKFLOW_LIST[0]['appPath']}))
  2596. assert_false('Permission denied.' in response.content, response.content)
  2597. # Login as someone else
  2598. client_not_me = make_logged_in_client(username='not_me', is_superuser=False, groupname='test', recreate=True)
  2599. grant_access("not_me", "not_me", "oozie")
  2600. response = client_not_me.get(reverse('oozie:list_oozie_workflows')+"?format=json")
  2601. assert_false('WordCount1' in response.content, response.content)
  2602. # Rerun
  2603. response = client_not_me.get(reverse('oozie:rerun_oozie_job', kwargs={'job_id': MockOozieApi.WORKFLOW_IDS[0],
  2604. 'app_path': MockOozieApi.JSON_WORKFLOW_LIST[0]['appPath']}))
  2605. assert_true('Permission denied.' in response.content, response.content)
  2606. # Add read only access
  2607. add_permission("not_me", "dashboard_jobs_access", "dashboard_jobs_access", "oozie")
  2608. response = client_not_me.get(reverse('oozie:list_oozie_workflows')+"?format=json")
  2609. assert_true('WordCount1' in response.content, response.content)
  2610. # Rerun
  2611. response = client_not_me.get(reverse('oozie:rerun_oozie_job', kwargs={'job_id': MockOozieApi.WORKFLOW_IDS[0],
  2612. 'app_path': MockOozieApi.JSON_WORKFLOW_LIST[0]['appPath']}))
  2613. assert_true('Permission denied.' in response.content, response.content)
  2614. def test_workflow_permissions(self):
  2615. response = self.c.get(reverse('oozie:list_oozie_workflow', args=[MockOozieApi.WORKFLOW_IDS[0]]))
  2616. assert_true('WordCount1' in response.content, response.content)
  2617. assert_false('Permission denied' in response.content, response.content)
  2618. response = self.c.get(reverse('oozie:list_oozie_workflow_action', args=['XXX']))
  2619. assert_false('Permission denied' in response.content, response.content)
  2620. # Login as someone else
  2621. client_not_me = make_logged_in_client(username='not_me', is_superuser=False, groupname='test', recreate=True)
  2622. grant_access("not_me", "not_me", "oozie")
  2623. response = client_not_me.get(reverse('oozie:list_oozie_workflow', args=[MockOozieApi.WORKFLOW_IDS[0]]))
  2624. assert_true('Permission denied' in response.content, response.content)
  2625. response = client_not_me.get(reverse('oozie:list_oozie_workflow_action', args=['XXX']))
  2626. assert_true('Permission denied' in response.content, response.content)
  2627. # Add read only access
  2628. add_permission("not_me", "dashboard_jobs_access", "dashboard_jobs_access", "oozie")
  2629. response = client_not_me.get(reverse('oozie:list_oozie_workflow', args=[MockOozieApi.WORKFLOW_IDS[0]]))
  2630. assert_false('Permission denied' in response.content, response.content)
  2631. def test_coordinators_permissions(self):
  2632. response = self.c.get(reverse('oozie:list_oozie_coordinators')+"?format=json")
  2633. assert_true('DailyWordCount1' in response.content, response.content)
  2634. # Login as someone else
  2635. client_not_me = make_logged_in_client(username='not_me', is_superuser=False, groupname='test', recreate=True)
  2636. grant_access("not_me", "not_me", "oozie")
  2637. response = client_not_me.get(reverse('oozie:list_oozie_coordinators')+"?format=json")
  2638. assert_false('DailyWordCount1' in response.content, response.content)
  2639. # Add read only access
  2640. add_permission("not_me", "dashboard_jobs_access", "dashboard_jobs_access", "oozie")
  2641. response = client_not_me.get(reverse('oozie:list_oozie_coordinators')+"?format=json")
  2642. assert_true('DailyWordCount1' in response.content, response.content)
  2643. def test_coordinator_permissions(self):
  2644. response = self.c.get(reverse('oozie:list_oozie_coordinator', args=[MockOozieApi.COORDINATOR_IDS[0]]))
  2645. assert_true('DailyWordCount1' in response.content, response.content)
  2646. assert_false('Permission denied' in response.content, response.content)
  2647. # Login as someone else
  2648. client_not_me = make_logged_in_client(username='not_me', is_superuser=False, groupname='test', recreate=True)
  2649. grant_access("not_me", "not_me", "oozie")
  2650. response = client_not_me.get(reverse('oozie:list_oozie_coordinator', args=[MockOozieApi.COORDINATOR_IDS[0]]))
  2651. assert_true('Permission denied' in response.content, response.content)
  2652. # Add read only access
  2653. add_permission("not_me", "dashboard_jobs_access", "dashboard_jobs_access", "oozie")
  2654. response = client_not_me.get(reverse('oozie:list_oozie_coordinator', args=[MockOozieApi.COORDINATOR_IDS[0]]))
  2655. assert_false('Permission denied' in response.content, response.content)
  2656. def test_bundles_permissions(self):
  2657. response = self.c.get(reverse('oozie:list_oozie_bundles') + "?format=json")
  2658. assert_true('MyBundle1' in response.content, response.content)
  2659. # Login as someone else
  2660. client_not_me = make_logged_in_client(username='not_me', is_superuser=False, groupname='test', recreate=True)
  2661. grant_access("not_me", "not_me", "oozie")
  2662. response = client_not_me.get(reverse('oozie:list_oozie_bundles')+"?format=json")
  2663. assert_false('MyBundle1' in response.content, response.content)
  2664. # Add read only access
  2665. add_permission("not_me", "dashboard_jobs_access", "dashboard_jobs_access", "oozie")
  2666. response = client_not_me.get(reverse('oozie:list_oozie_bundles')+"?format=json")
  2667. assert_true('MyBundle1' in response.content, response.content)
  2668. def test_bundle_permissions(self):
  2669. response = self.c.get(reverse('oozie:list_oozie_bundle', args=[MockOozieApi.BUNDLE_IDS[0]]))
  2670. assert_true('MyBundle1' in response.content, response.content)
  2671. assert_false('Permission denied' in response.content, response.content)
  2672. # Login as someone else
  2673. client_not_me = make_logged_in_client(username='not_me', is_superuser=False, groupname='test', recreate=True)
  2674. grant_access("not_me", "not_me", "oozie")
  2675. response = client_not_me.get(reverse('oozie:list_oozie_bundle', args=[MockOozieApi.BUNDLE_IDS[0]]))
  2676. assert_true('Permission denied' in response.content, response.content)
  2677. # Add read only access
  2678. add_permission("not_me", "dashboard_jobs_access", "dashboard_jobs_access", "oozie")
  2679. response = client_not_me.get(reverse('oozie:list_oozie_bundle', args=[MockOozieApi.BUNDLE_IDS[0]]))
  2680. assert_false('Permission denied' in response.content, response.content)
  2681. def test_good_workflow_status_graph(self):
  2682. workflow_count = Document.objects.available_docs(Workflow, self.user).count()
  2683. response = self.c.get(reverse('oozie:list_oozie_workflow', args=[MockOozieApi.WORKFLOW_IDS[0]]), {})
  2684. assert_true(response.context['workflow_graph'])
  2685. assert_equal(Document.objects.available_docs(Workflow, self.user).count(), workflow_count)
  2686. def test_bad_workflow_status_graph(self):
  2687. workflow_count = Document.objects.available_docs(Workflow, self.user).count()
  2688. response = self.c.get(reverse('oozie:list_oozie_workflow', args=[MockOozieApi.WORKFLOW_IDS[1]]), {})
  2689. assert_true(response.context['workflow_graph'] is None)
  2690. assert_equal(Document.objects.available_docs(Workflow, self.user).count(), workflow_count)
  2691. def test_list_oozie_sla(self):
  2692. response = self.c.get(reverse('oozie:list_oozie_sla'))
  2693. assert_true('Oozie Dashboard' in response.content, response.content)
  2694. response = self.c.get(reverse('oozie:list_oozie_sla') + "?format=json")
  2695. for sla in MockOozieApi.WORKFLOWS_SLAS:
  2696. assert_equal({"oozie_slas": []}, json.loads(response.content), response.content)
  2697. response = self.c.post(reverse('oozie:list_oozie_sla') + "?format=json", {'job_name': 'kochang'})
  2698. for sla in MockOozieApi.WORKFLOWS_SLAS:
  2699. assert_true('MISS' in response.content, response.content)
  2700. class GeneralTestsWithOozie(OozieBase):
  2701. def setUp(self):
  2702. OozieBase.setUp(self)
  2703. def test_import_jobsub_actions(self):
  2704. design = OozieDesign(owner=self.user, name="test")
  2705. action = OozieMapreduceAction(jar_path='/tmp/test.jar')
  2706. action.action_type = OozieMapreduceAction.ACTION_TYPE
  2707. action.save()
  2708. design.root_action = action
  2709. design.save()
  2710. try:
  2711. # There should be 3 from examples
  2712. action = convert_jobsub_design(design)
  2713. assert_equal(design.name, action.name)
  2714. assert_equal(design.description, action.description)
  2715. assert_equal(OozieMapreduceAction.ACTION_TYPE, action.node_type)
  2716. finally:
  2717. OozieDesign.objects.all().delete()
  2718. OozieMapreduceAction.objects.all().delete()
  2719. class TestUtils(OozieMockBase):
  2720. def setUp(self):
  2721. OozieMockBase.setUp(self)
  2722. # When updating wf, update wf_json as well!
  2723. self.wf = Document.objects.get_docs(self.user, Workflow).get(name='wf-name-1').content_object
  2724. def test_workflow_to_dict(self):
  2725. workflow_dict = workflow_to_dict(self.wf)
  2726. # Test properties
  2727. assert_true('job_xml' in workflow_dict, workflow_dict)
  2728. assert_true('is_shared' in workflow_dict, workflow_dict)
  2729. assert_true('end' in workflow_dict, workflow_dict)
  2730. assert_true('description' in workflow_dict, workflow_dict)
  2731. assert_true('parameters' in workflow_dict, workflow_dict)
  2732. assert_true('is_single' in workflow_dict, workflow_dict)
  2733. assert_true('deployment_dir' in workflow_dict, workflow_dict)
  2734. assert_true('schema_version' in workflow_dict, workflow_dict)
  2735. assert_true('job_properties' in workflow_dict, workflow_dict)
  2736. assert_true('start' in workflow_dict, workflow_dict)
  2737. assert_true('nodes' in workflow_dict, workflow_dict)
  2738. assert_true('id' in workflow_dict, workflow_dict)
  2739. assert_true('name' in workflow_dict, workflow_dict)
  2740. # Check links
  2741. for node in workflow_dict['nodes']:
  2742. assert_true('child_links' in node, node)
  2743. for link in node['child_links']:
  2744. assert_true('name' in link, link)
  2745. assert_true('comment' in link, link)
  2746. assert_true('parent' in link, link)
  2747. assert_true('child' in link, link)
  2748. def test_model_to_dict(self):
  2749. node_dict = model_to_dict(self.wf.node_set.filter(node_type='start')[0])
  2750. # Test properties
  2751. assert_true('id' in node_dict)
  2752. assert_true('name' in node_dict)
  2753. assert_true('description' in node_dict)
  2754. assert_true('node_type' in node_dict)
  2755. assert_true('workflow' in node_dict)
  2756. def test_smart_path(self):
  2757. assert_equal('${nameNode}/user/${wf:user()}/out', smart_path('out', {'output': '/path/out'}))
  2758. assert_equal('${nameNode}/path', smart_path('/path', {'output': '/path/out'}))
  2759. assert_equal('${nameNode}/path', smart_path('/path', {}))
  2760. assert_equal('${nameNode}${output}', smart_path('${output}', {'output': '/path/out'}))
  2761. assert_equal('hdfs://nn${output}', smart_path('hdfs://nn${output}', {'output': '/path/out'}))
  2762. assert_equal('${output}', smart_path('${output}', {}))
  2763. assert_equal('${output}', smart_path('${output}', {'output': 'hdfs://nn/path/out'}))
  2764. assert_equal('${output}', smart_path('${output}', {'output': '${path}'}))
  2765. assert_equal('${output_dir}', smart_path('${output_dir}', {'output': '/path/out', 'output_dir': 'hdfs://nn/path/out'}))
  2766. # Utils
  2767. WORKFLOW_DICT = {
  2768. u'deployment_dir': [u''], u'name': [u'wf-name-1'], u'description': [u''],
  2769. u'schema_version': [u'uri:oozie:workflow:0.4'],
  2770. u'parameters': [u'[{"name":"market","value":"US"}]'],
  2771. u'job_xml': [u'jobconf.xml'],
  2772. u'job_properties': [u'[{"name":"sleep-all","value":"${SLEEP}"}]']
  2773. }
  2774. COORDINATOR_DICT = {
  2775. u'name': [u'MyCoord'], u'description': [u'Description of my coordinator'],
  2776. u'workflow': [u'1'],
  2777. u'frequency_number': [u'1'], u'frequency_unit': [u'days'],
  2778. u'start_0': [u'07/01/2012'], u'start_1': [u'12:00 AM'],
  2779. u'end_0': [u'07/04/2012'], u'end_1': [u'12:00 AM'],
  2780. u'timezone': [u'America/Los_Angeles'],
  2781. u'parameters': [u'[{"name":"market","value":"US"}]'],
  2782. u'job_properties': [u'[{"name":"username","value":"${coord:user()}"},{"name":"SLEEP","value":"1000"}]'],
  2783. u'timeout': [u'100'],
  2784. u'concurrency': [u'3'],
  2785. u'execution': [u'FIFO'],
  2786. u'throttle': [u'10'],
  2787. u'schema_version': [u'uri:oozie:coordinator:0.2']
  2788. }
  2789. BUNDLE_DICT = {
  2790. u'name': [u'MyBundle'], u'description': [u'Description of my bundle'],
  2791. u'parameters': [u'[{"name":"market","value":"US,France"}]'],
  2792. u'kick_off_time_0': [u'07/01/2012'], u'kick_off_time_1': [u'12:00 AM'],
  2793. u'schema_version': [u'uri:oozie:coordinator:0.2']
  2794. }
  2795. def add_node(workflow, name, node_type, parents, attrs={}):
  2796. """
  2797. create a node of type node_type and associate the listed parents.
  2798. """
  2799. NodeClass = NODE_TYPES[node_type]
  2800. node = NodeClass(workflow=workflow, node_type=node_type, name=name)
  2801. for attr in attrs:
  2802. setattr(node, attr, attrs[attr])
  2803. node.save()
  2804. # Add parent
  2805. # If skipped, remember to preserve order: regular links first, then error link
  2806. if parents:
  2807. for parent in parents:
  2808. name = 'ok'
  2809. if parent.node_type == 'start' or parent.node_type == 'join':
  2810. name = 'to'
  2811. elif parent.node_type == 'fork' or parent.node_type == 'decision':
  2812. name = 'start'
  2813. link = Link(parent=parent, child=node, name=name)
  2814. link.save()
  2815. # Create error link
  2816. if node_type != 'fork' and node_type != 'decision' and node_type != 'join':
  2817. link = Link(parent=node, child=Kill.objects.get(name='kill', workflow=workflow), name="error")
  2818. link.save()
  2819. return node
  2820. def create_workflow(client, user, workflow_dict=WORKFLOW_DICT):
  2821. name = str(workflow_dict['name'][0])
  2822. # If not infinite looping
  2823. Node.objects.filter(workflow__name=name).delete()
  2824. # Leaking here for some reason
  2825. for doc in list(chain(Document.objects.get_docs(user, Workflow).filter(name=name, extra=''),
  2826. Document.objects.filter(name='mapreduce1', owner__username='jobsub_test').all(),
  2827. Document.objects.filter(name='sleep_job-copy', owner__username='jobsub_test').all())):
  2828. if doc.content_object:
  2829. client.post(reverse('oozie:delete_workflow') + '?skip_trash=true', {'job_selection': [doc.content_object.id]}, follow=True)
  2830. else:
  2831. doc.delete()
  2832. workflow_count = Document.objects.available_docs(Workflow, user).count()
  2833. response = client.get(reverse('oozie:create_workflow'))
  2834. assert_equal(workflow_count, Document.objects.available_docs(Workflow, user).count(), response)
  2835. response = client.post(reverse('oozie:create_workflow'), workflow_dict, follow=True)
  2836. assert_equal(200, response.status_code)
  2837. assert_equal(workflow_count + 1, Document.objects.available_docs(Workflow, user).count())
  2838. wf = Document.objects.get_docs(user, Workflow).get(name=name, extra='').content_object
  2839. assert_not_equal('', wf.deployment_dir)
  2840. assert_true(wf.managed)
  2841. return wf
  2842. def create_coordinator(workflow, client, user):
  2843. name = str(COORDINATOR_DICT['name'][0])
  2844. if Document.objects.get_docs(user, Coordinator).filter(name=name).exists():
  2845. for doc in Document.objects.get_docs(user, Coordinator).filter(name=name):
  2846. if doc.content_object:
  2847. client.post(reverse('oozie:delete_coordinator') + '?skip_trash=true', {'job_selection': [doc.content_object.id]}, follow=True)
  2848. else:
  2849. doc.delete()
  2850. coord_count = Document.objects.available_docs(Coordinator, user).count()
  2851. response = client.get(reverse('oozie:create_coordinator'))
  2852. assert_equal(coord_count, Document.objects.available_docs(Coordinator, user).count(), response)
  2853. post = COORDINATOR_DICT.copy()
  2854. post['workflow'] = workflow.id
  2855. response = client.post(reverse('oozie:create_coordinator'), post)
  2856. assert_equal(coord_count + 1, Document.objects.available_docs(Coordinator, user).count(), response)
  2857. return Document.objects.available_docs(Coordinator, user).get(name=name).content_object
  2858. def create_bundle(client, user):
  2859. name = str(BUNDLE_DICT['name'][0])
  2860. if Document.objects.get_docs(user, Bundle).filter(name=name).exists():
  2861. for doc in Document.objects.get_docs(user, Bundle).filter(name=name):
  2862. if doc.content_object:
  2863. client.post(reverse('oozie:delete_bundle') + '?skip_trash=true', {'job_selection': [doc.content_object.id]}, follow=True)
  2864. else:
  2865. doc.delete()
  2866. bundle_count = Document.objects.available_docs(Bundle, user).count()
  2867. response = client.get(reverse('oozie:create_bundle'))
  2868. assert_equal(bundle_count, Document.objects.available_docs(Bundle, user).count(), response)
  2869. post = BUNDLE_DICT.copy()
  2870. response = client.post(reverse('oozie:create_bundle'), post)
  2871. assert_equal(bundle_count + 1, Document.objects.available_docs(Bundle, user).count(), response)
  2872. return Document.objects.available_docs(Bundle, user).get(name=name).content_object
  2873. def create_dataset(coord, client):
  2874. response = client.post(reverse('oozie:create_coordinator_dataset', args=[coord.id]), {
  2875. u'create-name': [u'MyDataset'], u'create-frequency_number': [u'1'], u'create-frequency_unit': [u'days'],
  2876. u'create-uri': [u'/data/${YEAR}${MONTH}${DAY}'],
  2877. u'create-instance_choice': [u'range'],
  2878. u'create-advanced_start_instance': [u'-1'],
  2879. u'create-advanced_end_instance': [u'${coord:current(1)}'],
  2880. u'create-start_0': [u'07/01/2012'], u'create-start_1': [u'12:00 AM'],
  2881. u'create-timezone': [u'America/Los_Angeles'], u'create-done_flag': [u''],
  2882. u'create-description': [u'']})
  2883. data = json.loads(response.content)
  2884. assert_equal(0, data['status'], data['data'])
  2885. def create_coordinator_data(coord, client):
  2886. response = client.post(reverse('oozie:create_coordinator_data', args=[coord.id, 'input']),
  2887. {u'input-name': [u'input_dir'], u'input-dataset': [u'1']})
  2888. data = json.loads(response.content)
  2889. assert_equal(0, data['status'], data['data'])
  2890. def synchronize_workflow_attributes(workflow_json, correct_workflow_json):
  2891. if isinstance(workflow_json, basestring):
  2892. workflow_dict = json.loads(workflow_json)
  2893. else:
  2894. workflow_dict = workflow_json
  2895. if isinstance(correct_workflow_json, basestring):
  2896. correct_workflow_dict = json.loads(correct_workflow_json)
  2897. else:
  2898. correct_workflow_dict = correct_workflow_json
  2899. if 'attributes' in workflow_dict and 'attributes' in correct_workflow_dict:
  2900. workflow_dict['attributes']['deployment_dir'] = correct_workflow_dict['attributes']['deployment_dir']
  2901. return reformat_json(workflow_dict)