tests.py 85 KB


  1. #!/usr/bin/env python
  2. # Licensed to Cloudera, Inc. under one
  3. # or more contributor license agreements. See the NOTICE file
  4. # distributed with this work for additional information
  5. # regarding copyright ownership. Cloudera, Inc. licenses this file
  6. # to you under the Apache License, Version 2.0 (the
  7. # "License"); you may not use this file except in compliance
  8. # with the License. You may obtain a copy of the License at
  9. #
  10. # http://www.apache.org/licenses/LICENSE-2.0
  11. #
  12. # Unless required by applicable law or agreed to in writing, software
  13. # distributed under the License is distributed on an "AS IS" BASIS,
  14. # WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
  15. # See the License for the specific language governing permissions and
  16. # limitations under the License.
  17. try:
  18. import json
  19. except ImportError:
  20. import simplejson as json
  21. import logging
  22. import re
  23. import os
  24. from nose.plugins.skip import SkipTest
  25. from nose.tools import assert_true, assert_false, assert_equal, assert_not_equal
  26. from django.contrib.auth.models import User
  27. from django.core.urlresolvers import reverse
  28. from desktop.lib.django_test_util import make_logged_in_client
  29. from desktop.lib.test_utils import grant_access, add_permission
  30. from jobsub.management.commands import jobsub_setup
  31. from jobsub.models import OozieDesign
  32. from liboozie import oozie_api
  33. from liboozie.conf import OOZIE_URL
  34. from liboozie.oozie_api_test import OozieServerProvider
  35. from liboozie.types import WorkflowList, Workflow as OozieWorkflow, Coordinator as OozieCoordinator,\
  36. CoordinatorList, WorkflowAction
  37. from oozie.models import Workflow, Node, Kill, Link, Job, Coordinator, History,\
  38. find_parameters, NODE_TYPES
  39. from oozie.conf import SHARE_JOBS
  40. from oozie.utils import workflow_to_dict, model_to_dict
  41. from oozie.import_workflow import import_workflow
  42. LOG = logging.getLogger(__name__)
  43. _INITIALIZED = False
  44. class MockOozieApi:
  45. 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'},
  46. {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'},
  47. {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'},
  48. {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'}]
  49. WORKFLOW_IDS = [wf['id'] for wf in JSON_WORKFLOW_LIST]
  50. WORKFLOW_DICT = dict([(wf['id'], wf) for wf in JSON_WORKFLOW_LIST])
  51. 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'},
  52. {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'},
  53. {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'},
  54. {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'}]
  55. COORDINATOR_IDS = [coord['coordJobId'] for coord in JSON_COORDINATOR_LIST]
  56. COORDINATOR_DICT = dict([(coord['coordJobId'], coord) for coord in JSON_COORDINATOR_LIST])
  57. 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.2">\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'}
  58. def __init__(self, *args, **kwargs):
  59. pass
  60. def setuser(self, user):
  61. pass
  62. def submit_job(self, properties):
  63. return 'ONE-OOZIE-ID-W'
  64. def get_workflows(self, **kwargs):
  65. workflows = MockOozieApi.JSON_WORKFLOW_LIST
  66. if 'user' in kwargs:
  67. workflows = filter(lambda wf: wf['user'] == kwargs['user'], workflows)
  68. return WorkflowList(self, {'offset': 0, 'total': 4, 'workflows': workflows})
  69. def get_coordinators(self, **kwargs):
  70. coordinatorjobs = MockOozieApi.JSON_COORDINATOR_LIST
  71. if 'user' in kwargs:
  72. coordinatorjobs = filter(lambda coord: coord['user'] == kwargs['user'], coordinatorjobs)
  73. return CoordinatorList(self, {'offset': 0, 'total': 5, 'coordinatorjobs': coordinatorjobs})
  74. def get_job(self, job_id):
  75. if job_id in MockOozieApi.WORKFLOW_DICT:
  76. return OozieWorkflow(self, MockOozieApi.WORKFLOW_DICT[job_id])
  77. else:
  78. return OozieWorkflow(self, {'id': job_id, 'actions': []})
  79. def get_coordinator(self, job_id):
  80. if job_id in MockOozieApi.COORDINATOR_DICT:
  81. return OozieCoordinator(self, MockOozieApi.COORDINATOR_DICT[job_id])
  82. else:
  83. return OozieCoordinator(self, {'id': job_id, 'actions': []})
  84. def get_action(self, action_id):
  85. return WorkflowAction(MockOozieApi.WORKFLOW_ACTION)
  86. def job_control(self, job_id, action):
  87. return 'Done'
  88. def get_job_definition(self, jobid):
  89. return '<xml></xml>'
  90. def get_job_log(self, jobid):
  91. return '<xml></xml>'
  92. class OozieMockBase(object):
  93. def setUp(self):
  94. # Beware: Monkey patch Oozie/LibOozie with Mock API
  95. if not hasattr(oozie_api, 'OriginalOozieApi'):
  96. oozie_api.OriginalOozieApi = oozie_api.OozieApi
  97. if not hasattr(Workflow.objects, 'original_check_workspace'):
  98. Workflow.objects.original_check_workspace = Workflow.objects.check_workspace
  99. Workflow.objects.check_workspace = lambda a, b: None
  100. oozie_api.OozieApi = MockOozieApi
  101. oozie_api._api_cache = None
  102. Coordinator.objects.all().delete()
  103. self.c = make_logged_in_client(is_superuser=False)
  104. grant_access("test", "test", "oozie")
  105. self.user = User.objects.get(username='test')
  106. self.wf = create_workflow(self.c)
  107. def tearDown(self):
  108. oozie_api.OozieApi = oozie_api.OriginalOozieApi
  109. Workflow.objects.check_workspace = Workflow.objects.original_check_workspace
  110. oozie_api._api_cache = None
  111. def setup_simple_workflow(self):
  112. """ Creates a linear workflow """
  113. Link.objects.filter(parent__workflow=self.wf).delete()
  114. Link(parent=self.wf.start, child=self.wf.end, name="related").save()
  115. action1 = add_node(self.wf, 'action-name-1', 'mapreduce', [self.wf.start], {
  116. 'description': '',
  117. 'files': '[]',
  118. 'jar_path': '/user/hue/oozie/examples/lib/hadoop-examples.jar',
  119. 'job_properties': '[{"name":"sleep","value":"${SLEEP}"}]',
  120. 'prepares': '[{"value":"${output}","type":"delete"},{"value":"/test","type":"mkdir"}]',
  121. 'archives': '[]',
  122. })
  123. action2 = add_node(self.wf, 'action-name-2', 'mapreduce', [action1], {
  124. 'description': '',
  125. 'files': '[]',
  126. 'jar_path': '/user/hue/oozie/examples/lib/hadoop-examples.jar',
  127. 'job_properties': '[{"name":"sleep","value":"${SLEEP}"}]',
  128. 'prepares': '[{"value":"${output}","type":"delete"},{"value":"/test","type":"mkdir"}]',
  129. 'archives': '[]',
  130. })
  131. action3 = add_node(self.wf, 'action-name-3', 'mapreduce', [action2], {
  132. 'description': '',
  133. 'files': '[]',
  134. 'jar_path': '/user/hue/oozie/examples/lib/hadoop-examples.jar',
  135. 'job_properties': '[{"name":"sleep","value":"${SLEEP}"}]',
  136. 'prepares': '[{"value":"${output}","type":"delete"},{"value":"/test","type":"mkdir"}]',
  137. 'archives': '[]',
  138. })
  139. Link(parent=action3, child=self.wf.end, name="ok").save()
  140. def setup_forking_workflow(self):
  141. """ Creates a workflow with a fork """
  142. Link.objects.filter(parent__workflow=self.wf).delete()
  143. Link(parent=self.wf.start, child=self.wf.end, name="related").save()
  144. fork1 = add_node(self.wf, 'fork-name-1', 'fork', [self.wf.start])
  145. action1 = add_node(self.wf, 'action-name-1', 'mapreduce', [fork1])
  146. action2 = add_node(self.wf, 'action-name-2', 'mapreduce', [fork1])
  147. join1 = add_node(self.wf, 'join-name-1', 'join', [action1, action2])
  148. Link(parent=fork1, child=join1, name="related").save()
  149. action3 = add_node(self.wf, 'action-name-3', 'mapreduce', [join1])
  150. Link(parent=action3, child=self.wf.end, name="ok").save()
  151. class OozieBase(OozieServerProvider):
  152. requires_hadoop = True
  153. def setUp(self):
  154. OozieServerProvider.setup_class()
  155. self.c = make_logged_in_client(is_superuser=False)
  156. self.user = User.objects.get(username="test")
  157. grant_access("test", "test", "oozie")
  158. self.cluster = OozieServerProvider.cluster
  159. self.install_examples()
  160. # Ensure access to MR folder
  161. self.cluster.fs.do_as_superuser(self.cluster.fs.chmod, '/tmp', 0777, recursive=True)
  162. def install_examples(self):
  163. global _INITIALIZED
  164. if _INITIALIZED:
  165. return
  166. self.c.post(reverse('oozie:setup_app'))
  167. self.cluster.fs.do_as_user('test', self.cluster.fs.create_home_dir, '/user/test')
  168. self.cluster.fs.do_as_superuser(self.cluster.fs.chmod, '/user/test', 0777, True)
  169. hue = User.objects.create_user('hue', 'hue' + '@localhost', 'hue')
  170. Workflow.objects.update(owner=hue)
  171. _INITIALIZED = True
  172. def setup_simple_workflow(self):
  173. """ Creates a linear workflow """
  174. Link.objects.filter(parent__workflow=self.wf).delete()
  175. Link(parent=self.wf.start, child=self.wf.end, name="related").save()
  176. action1 = add_node(self.wf, 'action-name-1', 'mapreduce', [self.wf.start], {
  177. 'description': '',
  178. 'files': '[]',
  179. 'jar_path': '/user/hue/oozie/examples/lib/hadoop-examples.jar',
  180. 'job_properties': '[{"name":"sleep","value":"${SLEEP}"}]',
  181. 'prepares': '[{"value":"${output}","type":"delete"},{"value":"/test","type":"mkdir"}]',
  182. 'archives': '[]',
  183. })
  184. action2 = add_node(self.wf, 'action-name-2', 'mapreduce', [action1], {
  185. 'description': '',
  186. 'files': '[]',
  187. 'jar_path': '/user/hue/oozie/examples/lib/hadoop-examples.jar',
  188. 'job_properties': '[{"name":"sleep","value":"${SLEEP}"}]',
  189. 'prepares': '[{"value":"${output}","type":"delete"},{"value":"/test","type":"mkdir"}]',
  190. 'archives': '[]',
  191. })
  192. action3 = add_node(self.wf, 'action-name-3', 'mapreduce', [action2], {
  193. 'description': '',
  194. 'files': '[]',
  195. 'jar_path': '/user/hue/oozie/examples/lib/hadoop-examples.jar',
  196. 'job_properties': '[{"name":"sleep","value":"${SLEEP}"}]',
  197. 'prepares': '[{"value":"${output}","type":"delete"},{"value":"/test","type":"mkdir"}]',
  198. 'archives': '[]',
  199. })
  200. Link(parent=action3, child=self.wf.end, name="ok").save()
  201. def setup_forking_workflow(self):
  202. """ Creates a workflow with a fork """
  203. Link.objects.filter(parent__workflow=self.wf).delete()
  204. Link(parent=self.wf.start, child=self.wf.end, name="related").save()
  205. fork1 = add_node(self.wf, 'fork-name-1', 'fork', [self.wf.start])
  206. action1 = add_node(self.wf, 'action-name-1', 'mapreduce', [fork1])
  207. action2 = add_node(self.wf, 'action-name-2', 'mapreduce', [fork1])
  208. join1 = add_node(self.wf, 'join-name-1', 'join', [action1, action2])
  209. Link(parent=fork1, child=join1, name="related").save()
  210. action3 = add_node(self.wf, 'action-name-3', 'mapreduce', [join1])
  211. Link(parent=action3, child=self.wf.end, name="ok").save()
  212. class TestAPI(OozieMockBase):
  213. def setUp(self):
  214. OozieMockBase.setUp(self)
  215. # When updating wf, update wf_json as well!
  216. self.wf = Workflow.objects.get(name='wf-name-1')
  217. def test_workflow_save(self):
  218. self.setup_simple_workflow()
  219. workflow_dict = workflow_to_dict(self.wf)
  220. workflow_dict = remove_related_fields( workflow_dict )
  221. workflow_json = json.dumps(workflow_dict)
  222. response = self.c.post(reverse('oozie:workflow_save', kwargs={'workflow': self.wf.pk}), data={'workflow': workflow_json})
  223. test_response_json = response.content
  224. test_response_json_object = json.loads(test_response_json)
  225. assert_equal(0, test_response_json_object['status'])
  226. # Change property and save
  227. workflow_dict = workflow_to_dict(self.wf)
  228. workflow_dict = remove_related_fields( workflow_dict )
  229. workflow_dict['description'] = 'test'
  230. workflow_json = json.dumps(workflow_dict)
  231. response = self.c.post(reverse('oozie:workflow_save', kwargs={'workflow': self.wf.pk}), data={'workflow': workflow_json})
  232. test_response_json = response.content
  233. test_response_json_object = json.loads(test_response_json)
  234. assert_equal(0, test_response_json_object['status'])
  235. wf = Workflow.objects.get(id=self.wf.id)
  236. assert_equal('test', wf.description)
  237. assert_equal(self.wf.name, wf.name)
  238. # Change node and save
  239. workflow_dict = workflow_to_dict(self.wf)
  240. workflow_dict = remove_related_fields( workflow_dict )
  241. workflow_dict['nodes'][2]['name'] = 'new-name'
  242. node_id = workflow_dict['nodes'][2]['id']
  243. workflow_json = json.dumps(workflow_dict)
  244. response = self.c.post(reverse('oozie:workflow_save', kwargs={'workflow': self.wf.pk}), data={'workflow': workflow_json})
  245. test_response_json = response.content
  246. test_response_json_object = json.loads(test_response_json)
  247. assert_equal(0, test_response_json_object['status'])
  248. node = Node.objects.get(id=node_id)
  249. assert_equal('new-name', node.name)
  250. def test_workflow_save_fail(self):
  251. self.setup_simple_workflow()
  252. # Bad workflow name
  253. workflow_dict = workflow_to_dict(self.wf)
  254. del workflow_dict['name']
  255. workflow_dict = remove_related_fields( workflow_dict )
  256. workflow_json = json.dumps(workflow_dict)
  257. response = self.c.post(reverse('oozie:workflow_save', kwargs={'workflow': self.wf.pk}), data={'workflow': workflow_json}, HTTP_X_REQUESTED_WITH='XMLHttpRequest')
  258. assert_equal(400, response.status_code)
  259. # Bad node name
  260. workflow_dict = workflow_to_dict(self.wf)
  261. del workflow_dict['nodes'][2]['name']
  262. workflow_dict = remove_related_fields( workflow_dict )
  263. workflow_json = json.dumps(workflow_dict)
  264. response = self.c.post(reverse('oozie:workflow_save', kwargs={'workflow': self.wf.pk}), data={'workflow': workflow_json}, HTTP_X_REQUESTED_WITH='XMLHttpRequest')
  265. assert_equal(400, response.status_code)
  266. def test_workflow(self):
  267. response = self.c.get(reverse('oozie:workflow', kwargs={'workflow': self.wf.pk}))
  268. test_response_json = response.content
  269. test_response_json_object = json.loads(test_response_json)
  270. assert_equal(0, test_response_json_object['status'])
  271. def test_workflow_validate_action(self):
  272. 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"}
  273. response = self.c.post(reverse('oozie:workflow_validate_action', kwargs={'workflow': self.wf.pk, 'node_type': 'hive'}), data={'node': json.dumps(data)}, HTTP_X_REQUESTED_WITH='XMLHttpRequest')
  274. test_response_json = response.content
  275. test_response_json_object = json.loads(test_response_json)
  276. assert_equal(0, test_response_json_object['status'])
  277. assert_true('name' in test_response_json_object['data'], test_response_json_object['data'])
  278. assert_equal(0, len(test_response_json_object['data']['name']), test_response_json_object['data'])
  279. assert_true('description' in test_response_json_object['data'], test_response_json_object['data'])
  280. assert_equal(0, len(test_response_json_object['data']['description']), test_response_json_object['data'])
  281. assert_true('script_path' in test_response_json_object['data'], test_response_json_object['data'])
  282. assert_equal(0, len(test_response_json_object['data']['script_path']), test_response_json_object['data'])
  283. assert_true('job_xml' in test_response_json_object['data'], test_response_json_object['data'])
  284. assert_equal(0, len(test_response_json_object['data']['job_xml']), test_response_json_object['data'])
  285. assert_true('job_properties' in test_response_json_object['data'], test_response_json_object['data'])
  286. assert_equal(0, len(test_response_json_object['data']['job_properties']), test_response_json_object['data'])
  287. assert_true('files' in test_response_json_object['data'], test_response_json_object['data'])
  288. assert_equal(0, len(test_response_json_object['data']['files']), test_response_json_object['data'])
  289. assert_true('params' in test_response_json_object['data'], test_response_json_object['data'])
  290. assert_equal(0, len(test_response_json_object['data']['params']), test_response_json_object['data'])
  291. assert_true('prepares' in test_response_json_object['data'], test_response_json_object['data'])
  292. assert_equal(0, len(test_response_json_object['data']['prepares']), test_response_json_object['data'])
  293. assert_true('archives' in test_response_json_object['data'], test_response_json_object['data'])
  294. assert_equal(0, len(test_response_json_object['data']['archives']), test_response_json_object['data'])
  295. def test_workflow_validate_action_fail(self):
  296. # Empty files field
  297. 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"}
  298. response = self.c.post(reverse('oozie:workflow_validate_action', kwargs={'workflow': self.wf.pk, 'node_type': 'hive'}), data={'node': json.dumps(data)}, HTTP_X_REQUESTED_WITH='XMLHttpRequest')
  299. test_response_json = response.content
  300. test_response_json_object = json.loads(test_response_json)
  301. assert_equal(-1, test_response_json_object['status'])
  302. assert_true('name' in test_response_json_object['data'], test_response_json_object['data'])
  303. assert_equal(0, len(test_response_json_object['data']['name']), test_response_json_object['data'])
  304. assert_true('description' in test_response_json_object['data'], test_response_json_object['data'])
  305. assert_equal(0, len(test_response_json_object['data']['description']), test_response_json_object['data'])
  306. assert_true('script_path' in test_response_json_object['data'], test_response_json_object['data'])
  307. assert_equal(0, len(test_response_json_object['data']['script_path']), test_response_json_object['data'])
  308. assert_true('job_xml' in test_response_json_object['data'], test_response_json_object['data'])
  309. assert_equal(0, len(test_response_json_object['data']['job_xml']), test_response_json_object['data'])
  310. assert_true('job_properties' in test_response_json_object['data'], test_response_json_object['data'])
  311. assert_equal(0, len(test_response_json_object['data']['job_properties']), test_response_json_object['data'])
  312. assert_true('files' in test_response_json_object['data'], test_response_json_object['data'])
  313. assert_equal(1, len(test_response_json_object['data']['files']), test_response_json_object['data'])
  314. assert_true('params' in test_response_json_object['data'], test_response_json_object['data'])
  315. assert_equal(0, len(test_response_json_object['data']['params']), test_response_json_object['data'])
  316. assert_true('prepares' in test_response_json_object['data'], test_response_json_object['data'])
  317. assert_equal(0, len(test_response_json_object['data']['prepares']), test_response_json_object['data'])
  318. assert_true('archives' in test_response_json_object['data'], test_response_json_object['data'])
  319. assert_equal(0, len(test_response_json_object['data']['archives']), test_response_json_object['data'])
  320. # Empty script path
  321. 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"}
  322. response = self.c.post(reverse('oozie:workflow_validate_action', kwargs={'workflow': self.wf.pk, 'node_type': 'hive'}), data={'node': json.dumps(data)}, HTTP_X_REQUESTED_WITH='XMLHttpRequest')
  323. test_response_json = response.content
  324. test_response_json_object = json.loads(test_response_json)
  325. assert_equal(-1, test_response_json_object['status'])
  326. assert_true('name' in test_response_json_object['data'], test_response_json_object['data'])
  327. assert_equal(0, len(test_response_json_object['data']['name']), test_response_json_object['data'])
  328. assert_true('description' in test_response_json_object['data'], test_response_json_object['data'])
  329. assert_equal(0, len(test_response_json_object['data']['description']), test_response_json_object['data'])
  330. assert_true('script_path' in test_response_json_object['data'], test_response_json_object['data'])
  331. assert_equal(1, len(test_response_json_object['data']['script_path']), test_response_json_object['data'])
  332. assert_true('job_xml' in test_response_json_object['data'], test_response_json_object['data'])
  333. assert_equal(0, len(test_response_json_object['data']['job_xml']), test_response_json_object['data'])
  334. assert_true('job_properties' in test_response_json_object['data'], test_response_json_object['data'])
  335. assert_equal(0, len(test_response_json_object['data']['job_properties']), test_response_json_object['data'])
  336. assert_true('files' in test_response_json_object['data'], test_response_json_object['data'])
  337. assert_equal(0, len(test_response_json_object['data']['files']), test_response_json_object['data'])
  338. assert_true('params' in test_response_json_object['data'], test_response_json_object['data'])
  339. assert_equal(0, len(test_response_json_object['data']['params']), test_response_json_object['data'])
  340. assert_true('prepares' in test_response_json_object['data'], test_response_json_object['data'])
  341. assert_equal(0, len(test_response_json_object['data']['prepares']), test_response_json_object['data'])
  342. assert_true('archives' in test_response_json_object['data'], test_response_json_object['data'])
  343. assert_equal(0, len(test_response_json_object['data']['archives']), test_response_json_object['data'])
  344. class TestApiPermissionsWithOozie(OozieBase):
  345. def setUp(self):
  346. OozieBase.setUp(self)
  347. # When updating wf, update wf_json as well!
  348. self.wf = Workflow.objects.get(name='MapReduce').clone(self.cluster.fs, self.user)
  349. def test_workflow_save(self):
  350. # Share
  351. self.wf.is_shared = True
  352. self.wf.save()
  353. Workflow.objects.check_workspace(self.wf, self.cluster.fs)
  354. # Login as someone else
  355. client_not_me = make_logged_in_client(username='not_me', is_superuser=False, groupname='test')
  356. grant_access("not_me", "test", "oozie")
  357. workflow_dict = workflow_to_dict(self.wf)
  358. workflow_dict = remove_related_fields(workflow_dict)
  359. workflow_json = json.dumps(workflow_dict)
  360. response = client_not_me.post(reverse('oozie:workflow_save', kwargs={'workflow': self.wf.pk}), data={'workflow': workflow_json}, HTTP_X_REQUESTED_WITH='XMLHttpRequest')
  361. assert_equal(401, response.status_code, response.status_code)
  362. response = self.c.post(reverse('oozie:workflow_save', kwargs={'workflow': self.wf.pk}), data={'workflow': workflow_json}, HTTP_X_REQUESTED_WITH='XMLHttpRequest')
  363. test_response_json = response.content
  364. test_response_json_object = json.loads(test_response_json)
  365. assert_equal(200, response.status_code, response)
  366. assert_equal(0, test_response_json_object['status'])
  367. def test_workflow_save_fail(self):
  368. # Unshare
  369. self.wf.is_shared = False
  370. self.wf.save()
  371. Workflow.objects.check_workspace(self.wf, self.cluster.fs)
  372. # Login as someone else
  373. client_not_me = make_logged_in_client(username='not_me', is_superuser=False, groupname='test')
  374. grant_access("not_me", "test", "oozie")
  375. workflow_dict = workflow_to_dict(self.wf)
  376. del workflow_dict['name']
  377. workflow_dict = remove_related_fields( workflow_dict )
  378. workflow_json = json.dumps(workflow_dict)
  379. response = client_not_me.post(reverse('oozie:workflow_save', kwargs={'workflow': self.wf.pk}), data={'workflow': workflow_json}, HTTP_X_REQUESTED_WITH='XMLHttpRequest')
  380. assert_equal(401, response.status_code, response)
  381. def test_workflow(self):
  382. # Share
  383. self.wf.is_shared = True
  384. self.wf.save()
  385. Workflow.objects.check_workspace(self.wf, self.cluster.fs)
  386. # Login as someone else
  387. client_not_me = make_logged_in_client(username='not_me', is_superuser=False, groupname='test')
  388. grant_access("not_me", "test", "oozie")
  389. response = client_not_me.get(reverse('oozie:workflow', kwargs={'workflow': self.wf.pk}), HTTP_X_REQUESTED_WITH='XMLHttpRequest')
  390. test_response_json = response.content
  391. test_response_json_object = json.loads(test_response_json)
  392. assert_equal(200, response.status_code, response)
  393. assert_equal(0, test_response_json_object['status'])
  394. def test_workflow_fail(self):
  395. # Unshare
  396. self.wf.is_shared = False
  397. self.wf.save()
  398. Workflow.objects.check_workspace(self.wf, self.cluster.fs)
  399. # Login as someone else
  400. client_not_me = make_logged_in_client(username='not_me', is_superuser=False, groupname='test')
  401. grant_access("not_me", "test", "oozie")
  402. response = client_not_me.get(reverse('oozie:workflow', kwargs={'workflow': self.wf.pk}), HTTP_X_REQUESTED_WITH='XMLHttpRequest')
  403. assert_equal(401, response.status_code)
  404. class TestEditor(OozieMockBase):
  405. def setUp(self):
  406. super(TestEditor, self).setUp()
  407. self.setup_simple_workflow()
  408. def test_find_parameters(self):
  409. jobs = [Job(name="$a"),
  410. Job(name="foo ${b} $$"),
  411. Job(name="${foo}", description="xxx ${foo}")]
  412. result = [find_parameters(job, ['name', 'description']) for job in jobs]
  413. assert_equal(set(["b", "foo"]), reduce(lambda x, y: x | set(y), result, set()))
  414. def test_find_all_parameters(self):
  415. assert_equal([{'name': u'output', 'value': u''}, {'name': u'SLEEP', 'value': ''}, {'name': u'market', 'value': u'US'}],
  416. self.wf.find_all_parameters())
  417. def test_workflow_has_cycle(self):
  418. action1 = Node.objects.get(name='action-name-1')
  419. action3 = Node.objects.get(name='action-name-3')
  420. assert_false(self.wf.has_cycle())
  421. ok = action3.get_link('ok')
  422. ok.child = action1
  423. ok.save()
  424. assert_true(self.wf.has_cycle())
  425. def test_workflow_gen_xml(self):
  426. assert_equal(
  427. '<workflow-app name="wf-name-1" xmlns="uri:oozie:workflow:0.2">\n'
  428. ' <global>\n'
  429. ' <job-xml>jobconf.xml</job-xml>\n'
  430. ' <configuration>\n'
  431. ' <property>\n'
  432. ' <name>sleep-all</name>\n'
  433. ' <value>${SLEEP}</value>\n'
  434. ' </property>\n'
  435. ' </configuration>\n'
  436. ' </global>\n'
  437. ' <start to="action-name-1"/>\n'
  438. ' <action name="action-name-1">\n'
  439. ' <map-reduce>\n'
  440. ' <job-tracker>${jobTracker}</job-tracker>\n'
  441. ' <name-node>${nameNode}</name-node>\n'
  442. ' <prepare>\n'
  443. ' <delete path="${nameNode}${output}"/>\n'
  444. ' <mkdir path="${nameNode}/test"/>\n'
  445. ' </prepare>\n'
  446. ' <configuration>\n'
  447. ' <property>\n'
  448. ' <name>sleep</name>\n'
  449. ' <value>${SLEEP}</value>\n'
  450. ' </property>\n'
  451. ' </configuration>\n'
  452. ' </map-reduce>\n'
  453. ' <ok to="action-name-2"/>\n'
  454. ' <error to="kill"/>\n'
  455. ' </action>\n'
  456. ' <action name="action-name-2">\n'
  457. ' <map-reduce>\n'
  458. ' <job-tracker>${jobTracker}</job-tracker>\n'
  459. ' <name-node>${nameNode}</name-node>\n'
  460. ' <prepare>\n'
  461. ' <delete path="${nameNode}${output}"/>\n'
  462. ' <mkdir path="${nameNode}/test"/>\n'
  463. ' </prepare>\n'
  464. ' <configuration>\n'
  465. ' <property>\n'
  466. ' <name>sleep</name>\n'
  467. ' <value>${SLEEP}</value>\n'
  468. ' </property>\n'
  469. ' </configuration>\n'
  470. ' </map-reduce>\n'
  471. ' <ok to="action-name-3"/>\n'
  472. ' <error to="kill"/>\n'
  473. ' </action>\n'
  474. ' <action name="action-name-3">\n'
  475. ' <map-reduce>\n'
  476. ' <job-tracker>${jobTracker}</job-tracker>\n'
  477. ' <name-node>${nameNode}</name-node>\n'
  478. ' <prepare>\n'
  479. ' <delete path="${nameNode}${output}"/>\n'
  480. ' <mkdir path="${nameNode}/test"/>\n'
  481. ' </prepare>\n'
  482. ' <configuration>\n'
  483. ' <property>\n'
  484. ' <name>sleep</name>\n'
  485. ' <value>${SLEEP}</value>\n'
  486. ' </property>\n'
  487. ' </configuration>\n'
  488. ' </map-reduce>\n'
  489. ' <ok to="end"/>\n'
  490. ' <error to="kill"/>\n'
  491. ' </action>\n'
  492. ' <kill name="kill">\n'
  493. ' <message>Action failed, error message[${wf:errorMessage(wf:lastErrorNode())}]</message>\n'
  494. ' </kill>\n'
  495. ' <end name="end"/>\n'
  496. '</workflow-app>'.split(), self.wf.to_xml().split())
  497. def test_workflow_shell_gen_xml(self):
  498. self.wf.node_set.filter(name='action-name-1').delete()
  499. action1 = add_node(self.wf, 'action-name-1', 'shell', [self.wf.start], {
  500. u'job_xml': 'my-job.xml',
  501. u'files': '["hello.py"]',
  502. u'name': 'Shell',
  503. u'job_properties': '[]',
  504. u'capture_output': 'on',
  505. u'command': 'hello.py',
  506. u'archives': '[]',
  507. u'prepares': '[]',
  508. u'params': '[{"value":"World!","type":"argument"}]',
  509. u'description': 'Execute a Python script printing its arguments'
  510. })
  511. Link(parent=action1, child=self.wf.end, name="ok").save()
  512. xml = self.wf.to_xml()
  513. assert_true("""
  514. <shell xmlns="uri:oozie:shell-action:0.1">
  515. <job-tracker>${jobTracker}</job-tracker>
  516. <name-node>${nameNode}</name-node>
  517. <job-xml>my-job.xml</job-xml>
  518. <exec>hello.py</exec>
  519. <argument>World!</argument>
  520. <file>hello.py#hello.py</file>
  521. <capture-output/>
  522. </shell>""" in xml, xml)
  523. action1.capture_output = False
  524. action1.save()
  525. xml = self.wf.to_xml()
  526. assert_true("""
  527. <shell xmlns="uri:oozie:shell-action:0.1">
  528. <job-tracker>${jobTracker}</job-tracker>
  529. <name-node>${nameNode}</name-node>
  530. <job-xml>my-job.xml</job-xml>
  531. <exec>hello.py</exec>
  532. <argument>World!</argument>
  533. <file>hello.py#hello.py</file>
  534. </shell>""" in xml, xml)
  535. def test_workflow_fs_gen_xml(self):
  536. self.wf.node_set.filter(name='action-name-1').delete()
  537. action1 = add_node(self.wf, 'action-name-1', 'fs', [self.wf.start], {
  538. u'name': 'MyFs',
  539. u'description': 'Execute a Fs action that manage files',
  540. u'deletes': '[{"name":"/to/delete"},{"name":"to/delete2"}]',
  541. u'mkdirs': '[{"name":"/to/mkdir"},{"name":"${mkdir2}"}]',
  542. u'moves': '[{"source":"/to/move/source","destination":"/to/move/destination"},{"source":"/to/move/source2","destination":"/to/move/destination2"}]',
  543. u'chmods': '[{"path":"/to/chmod","recursive":true,"permissions":"-rwxrw-rw-"},{"path":"/to/chmod2","recursive":false,"permissions":"755"}]',
  544. u'touchzs': '[{"name":"/to/touchz"},{"name":"/to/touchz2"}]'
  545. })
  546. Link(parent=action1, child=self.wf.end, name="ok").save()
  547. xml = self.wf.to_xml()
  548. assert_true("""
  549. <action name="MyFs">
  550. <fs>
  551. <delete path='${nameNode}/to/delete'/>
  552. <delete path='${nameNode}/user/${wf:user()}/to/delete2'/>
  553. <mkdir path='${nameNode}/to/mkdir'/>
  554. <mkdir path='${nameNode}${mkdir2}'/>
  555. <move source='${nameNode}/to/move/source' target='${nameNode}/to/move/destination'/>
  556. <move source='${nameNode}/to/move/source2' target='${nameNode}/to/move/destination2'/>
  557. <chmod path='${nameNode}/to/chmod' permissions='-rwxrw-rw-' dir-files='true'/>
  558. <chmod path='${nameNode}/to/chmod2' permissions='755' dir-files='false'/>
  559. <touchz path='${nameNode}/to/touchz'/>
  560. <touchz path='${nameNode}/to/touchz2'/>
  561. </fs>
  562. <ok to="end"/>
  563. <error to="kill"/>
  564. </action>""" in xml, xml)
  565. def test_workflow_email_gen_xml(self):
  566. self.wf.node_set.filter(name='action-name-1').delete()
  567. action1 = add_node(self.wf, 'action-name-1', 'email', [self.wf.start], {
  568. u'name': 'MyEmail',
  569. u'description': 'Execute an Email action',
  570. u'to': 'hue@hue.org,django@python.org',
  571. u'cc': '',
  572. u'subject': 'My subject',
  573. u'body': 'My body'
  574. })
  575. Link(parent=action1, child=self.wf.end, name="ok").save()
  576. xml = self.wf.to_xml()
  577. assert_true("""
  578. <action name="MyEmail">
  579. <email xmlns="uri:oozie:email-action:0.1">
  580. <to>hue@hue.org,django@python.org</to>
  581. <subject>My subject</subject>
  582. <body>My body</body>
  583. </email>
  584. <ok to="end"/>
  585. <error to="kill"/>
  586. </action>""" in xml, xml)
  587. action1.cc = 'lambda@python.org'
  588. action1.save()
  589. xml = self.wf.to_xml()
  590. assert_true("""
  591. <action name="MyEmail">
  592. <email xmlns="uri:oozie:email-action:0.1">
  593. <to>hue@hue.org,django@python.org</to>
  594. <cc>lambda@python.org</cc>
  595. <subject>My subject</subject>
  596. <body>My body</body>
  597. </email>
  598. <ok to="end"/>
  599. <error to="kill"/>
  600. </action>""" in xml, xml)
  601. def test_workflow_subworkflow_gen_xml(self):
  602. self.wf.node_set.filter(name='action-name-1').delete()
  603. wf_dict = WORKFLOW_DICT.copy()
  604. wf_dict['name'] = [u'wf-name-2']
  605. wf2 = create_workflow(self.c, wf_dict)
  606. action1 = add_node(self.wf, 'action-name-1', 'subworkflow', [self.wf.start], {
  607. u'name': 'MySubworkflow',
  608. u'description': 'Execute a subworkflow action',
  609. u'sub_workflow': wf2,
  610. u'propagate_configuration': True,
  611. u'job_properties': '[{"value":"World!","name":"argument"}]'
  612. })
  613. Link(parent=action1, child=self.wf.end, name="ok").save()
  614. xml = self.wf.to_xml()
  615. assert_true(re.search(
  616. '<sub-workflow>\W+'
  617. '<app-path>\${nameNode}/user/hue/oozie/workspaces/_test_-oozie-(.+?)</app-path>\W+'
  618. '<propagate-configuration/>\W+'
  619. '<configuration>\W+'
  620. '<property>\W+'
  621. '<name>argument</name>\W+'
  622. '<value>World!</value>\W+'
  623. '</property>\W+'
  624. '</configuration>\W+'
  625. '</sub-workflow>', xml, re.MULTILINE), xml)
  626. wf2.delete()
  627. def test_workflow_flatten_list(self):
  628. assert_equal('[<Start: start>, <Mapreduce: action-name-1>, <Mapreduce: action-name-2>, <Mapreduce: action-name-3>, '
  629. '<Kill: kill>, <End: end>]',
  630. str(self.wf.node_list))
  631. # 1 2
  632. # 3
  633. self.setup_forking_workflow()
  634. assert_equal('[<Start: start>, <Fork: fork-name-1>, <Mapreduce: action-name-1>, <Mapreduce: action-name-2>, '
  635. '<Join: join-name-1>, <Mapreduce: action-name-3>, <Kill: kill>, <End: end>]',
  636. str(self.wf.node_list))
  637. def test_create_coordinator(self):
  638. create_coordinator(self.wf, self.c)
  639. def test_clone_coordinator(self):
  640. coord = create_coordinator(self.wf, self.c)
  641. coordinator_count = Coordinator.objects.count()
  642. response = self.c.post(reverse('oozie:clone_coordinator', args=[coord.id]), {}, follow=True)
  643. coord2 = Coordinator.objects.latest('id')
  644. assert_not_equal(coord.id, coord2.id)
  645. assert_equal(coordinator_count + 1, Coordinator.objects.count(), response)
  646. assert_equal(coord.dataset_set.count(), coord2.dataset_set.count())
  647. assert_equal(coord.datainput_set.count(), coord2.datainput_set.count())
  648. assert_equal(coord.dataoutput_set.count(), coord2.dataoutput_set.count())
  649. ds_ids = set(coord.dataset_set.values_list('id', flat=True))
  650. for node in coord2.dataset_set.all():
  651. assert_false(node.id in ds_ids)
  652. data_input_ids = set(coord.datainput_set.values_list('id', flat=True))
  653. for node in coord2.datainput_set.all():
  654. assert_false(node.id in data_input_ids)
  655. data_output_ids = set(coord.dataoutput_set.values_list('id', flat=True))
  656. for node in coord2.dataoutput_set.all():
  657. assert_false(node.id in data_output_ids)
  658. assert_not_equal(coord.deployment_dir, coord2.deployment_dir)
  659. assert_not_equal('', coord2.deployment_dir)
  660. def test_coordinator_workflow_access_permissions(self):
  661. self.wf.is_shared = True
  662. self.wf.save()
  663. # Login as someone else not superuser
  664. client_another_me = make_logged_in_client(username='another_me', is_superuser=False, groupname='test')
  665. grant_access("another_me", "test", "oozie")
  666. coord = create_coordinator(self.wf, client_another_me)
  667. response = client_another_me.get(reverse('oozie:edit_coordinator', args=[coord.id]))
  668. assert_true('Editor' in response.content, response.content)
  669. assert_true('value="Save"' in response.content, response.content)
  670. # Check can schedule a non personal/shared workflow
  671. workflow_select = '%s</option>' % self.wf
  672. response = client_another_me.get(reverse('oozie:edit_coordinator', args=[coord.id]))
  673. assert_true(workflow_select in response.content, response.content)
  674. self.wf.is_shared = False
  675. self.wf.save()
  676. response = client_another_me.get(reverse('oozie:edit_coordinator', args=[coord.id]))
  677. assert_false(workflow_select in response.content, response.content)
  678. self.wf.is_shared = True
  679. self.wf.save()
  680. # Edit
  681. finish = SHARE_JOBS.set_for_testing(True)
  682. try:
  683. response = client_another_me.post(reverse('oozie:edit_coordinator', args=[coord.id]))
  684. assert_true(workflow_select in response.content, response.content)
  685. assert_true('value="Save"' in response.content, response.content)
  686. finally:
  687. finish()
  688. finish = SHARE_JOBS.set_for_testing(False)
  689. try:
  690. response = client_another_me.post(reverse('oozie:edit_coordinator', args=[coord.id]))
  691. assert_true('This field is required' in response.content, response.content)
  692. assert_false(workflow_select in response.content, response.content)
  693. assert_true('value="Save"' in response.content, response.content)
  694. finally:
  695. finish()
  696. def test_coordinator_gen_xml(self):
  697. coord = create_coordinator(self.wf, self.c)
  698. assert_equal(
  699. '<coordinator-app name="MyCoord"\n'
  700. ' frequency="${coord:days(1)}"\n'
  701. ' start="2012-07-01T00:00Z" end="2012-07-04T00:00Z" timezone="America/Los_Angeles"\n'
  702. ' xmlns="uri:oozie:coordinator:0.1">\n'
  703. ' <controls>\n'
  704. ' <timeout>100</timeout>\n'
  705. ' <concurrency>3</concurrency>\n'
  706. ' <execution>FIFO</execution>\n'
  707. ' <throttle>10</throttle>\n'
  708. ' </controls>\n'
  709. ' <action>\n'
  710. ' <workflow>\n'
  711. ' <app-path>${wf_application_path}</app-path>\n'
  712. ' </workflow>\n'
  713. ' </action>\n'
  714. '</coordinator-app>\n'.split(), coord.to_xml().split())
  715. def test_coordinator_with_data_input_gen_xml(self):
  716. coord = create_coordinator(self.wf, self.c)
  717. create_dataset(coord, self.c)
  718. create_coordinator_data(coord, self.c)
  719. assert_equal(
  720. ['<coordinator-app', 'name="MyCoord"', 'frequency="${coord:days(1)}"', 'start="2012-07-01T00:00Z"', 'end="2012-07-04T00:00Z"',
  721. 'timezone="America/Los_Angeles"',
  722. 'xmlns="uri:oozie:coordinator:0.1">',
  723. '<controls>',
  724. '<timeout>100</timeout>',
  725. '<concurrency>3</concurrency>',
  726. '<execution>FIFO</execution>',
  727. '<throttle>10</throttle>',
  728. '</controls>',
  729. '<datasets>',
  730. '<dataset', 'name="MyDataset"', 'frequency="${coord:days(1)}"', 'initial-instance="2012-07-01T00:00Z"', 'timezone="America/Los_Angeles">',
  731. '<uri-template>/data/${YEAR}${MONTH}${DAY}</uri-template>',
  732. '<done-flag></done-flag>',
  733. '</dataset>',
  734. '</datasets>',
  735. '<input-events>',
  736. '<data-in', 'name="input_dir"', 'dataset="MyDataset">',
  737. '<instance>${coord:current(0)}</instance>',
  738. '</data-in>',
  739. '</input-events>',
  740. '<action>',
  741. '<workflow>',
  742. '<app-path>${wf_application_path}</app-path>',
  743. '<configuration>',
  744. '<property>',
  745. '<name>input_dir</name>',
  746. "<value>${coord:dataIn('input_dir')}</value>",
  747. '</property>',
  748. '</configuration>',
  749. '</workflow>',
  750. '</action>',
  751. '</coordinator-app>'], coord.to_xml().split())
  752. def test_create_coordinator_dataset(self):
  753. coord = create_coordinator(self.wf, self.c)
  754. create_dataset(coord, self.c)
  755. def test_create_coordinator_input_data(self):
  756. coord = create_coordinator(self.wf, self.c)
  757. create_dataset(coord, self.c)
  758. create_coordinator_data(coord, self.c)
  759. def test_setup_app(self):
  760. self.c.post(reverse('oozie:setup_app'))
  761. def test_workflow_prepare(self):
  762. action1 = Node.objects.get(name='action-name-1').get_full_node()
  763. action1.prepares = json.dumps([
  764. {"type": "delete","value": "${output}"},
  765. {"type": "delete","value": "out"},
  766. {"type": "delete","value": "/user/test/out"},
  767. {"type": "delete","value": "hdfs://localhost:8020/user/test/out"}])
  768. action1.save()
  769. xml = self.wf.to_xml()
  770. assert_true('<delete path="${nameNode}${output}"/>' in xml, xml)
  771. assert_true('<delete path="${nameNode}/user/${wf:user()}/out"/>' in xml, xml)
  772. assert_true('<delete path="${nameNode}/user/test/out"/>' in xml, xml)
  773. assert_true('<delete path="hdfs://localhost:8020/user/test/out"/>' in xml, xml)
  774. def test_get_workflow_parameters(self):
  775. assert_equal([{'name': u'output', 'value': ''}, {'name': u'SLEEP', 'value': ''}, {'name': u'market', 'value': u'US'}],
  776. self.wf.find_all_parameters())
  777. def test_get_coordinator_parameters(self):
  778. coord = create_coordinator(self.wf, self.c)
  779. create_dataset(coord, self.c)
  780. create_coordinator_data(coord, self.c)
  781. assert_equal([{'name': u'output', 'value': ''}, {'name': u'SLEEP', 'value': ''}, {'name': u'market', 'value': u'US,France'}],
  782. coord.find_all_parameters())
  783. def test_workflow_data_binds(self):
  784. response = self.c.get(reverse('oozie:edit_workflow', args=[self.wf.id]))
  785. assert_equal(1, response.content.count('checked: is_shared'), response.content)
  786. assert_true('checked: capture_output' in response.content, response.content)
  787. def test_import_workflow_basic(self):
  788. workflow = Workflow.objects.new_workflow(self.user)
  789. workflow.save()
  790. f = open('apps/oozie/src/oozie/test_data/0.4/test-basic.xml')
  791. import_workflow(workflow, f.read(), schema_version=0.4)
  792. f.close()
  793. workflow.save()
  794. assert_equal(2, len(Node.objects.filter(workflow=workflow)))
  795. assert_equal(1, len(Link.objects.filter(parent__workflow=workflow)))
  796. assert_equal('done', Node.objects.get(workflow=workflow, node_type='end').name)
  797. workflow.delete()
  798. def test_import_workflow_decision(self):
  799. workflow = Workflow.objects.new_workflow(self.user)
  800. workflow.save()
  801. f = open('apps/oozie/src/oozie/test_data/0.4/test-decision.xml')
  802. import_workflow(workflow, f.read(), schema_version=0.4)
  803. f.close()
  804. workflow.save()
  805. assert_equal(11, len(Node.objects.filter(workflow=workflow)))
  806. assert_equal(19, len(Link.objects.filter(parent__workflow=workflow)))
  807. assert_equal(1, len(Link.objects.filter(parent__workflow=workflow, parent__node_type='decision', comment='${1 gt 2}', name='start')))
  808. assert_equal(1, len(Link.objects.filter(parent__workflow=workflow, parent__node_type='decision', comment='', name='start')))
  809. assert_equal(1, len(Link.objects.filter(parent__workflow=workflow, parent__node_type='decision', name='default')))
  810. assert_equal(1, len(Link.objects.filter(parent__workflow=workflow, parent__node_type='decision', child__node_type='end', name='related')))
  811. workflow.delete()
  812. def test_import_workflow_distcp(self):
  813. workflow = Workflow.objects.new_workflow(self.user)
  814. workflow.save()
  815. f = open('apps/oozie/src/oozie/test_data/0.4/test-distcp.0.1.xml')
  816. import_workflow(workflow, f.read(), schema_version=0.4)
  817. f.close()
  818. workflow.save()
  819. assert_equal(4, len(Node.objects.filter(workflow=workflow)))
  820. assert_equal(3, len(Link.objects.filter(parent__workflow=workflow)))
  821. 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)
  822. workflow.delete()
  823. def test_import_workflow_forks(self):
  824. workflow = Workflow.objects.new_workflow(self.user)
  825. workflow.save()
  826. f = open('apps/oozie/src/oozie/test_data/0.4/test-forks.xml')
  827. import_workflow(workflow, f.read(), schema_version=0.4)
  828. f.close()
  829. workflow.save()
  830. assert_equal(12, len(Node.objects.filter(workflow=workflow)))
  831. assert_equal(19, len(Link.objects.filter(parent__workflow=workflow)))
  832. assert_equal(6, len(Link.objects.filter(parent__workflow=workflow, parent__node_type='fork')))
  833. assert_equal(4, len(Link.objects.filter(parent__workflow=workflow, parent__node_type='fork', name='start')))
  834. assert_equal(2, len(Link.objects.filter(parent__workflow=workflow, parent__node_type='fork', child__node_type='join', name='related')))
  835. workflow.delete()
  836. def test_import_workflow_mapreduce(self):
  837. workflow = Workflow.objects.new_workflow(self.user)
  838. workflow.save()
  839. f = open('apps/oozie/src/oozie/test_data/0.4/test-mapreduce.xml')
  840. import_workflow(workflow, f.read(), schema_version=0.4)
  841. f.close()
  842. workflow.save()
  843. assert_equal(4, len(Node.objects.filter(workflow=workflow)))
  844. assert_equal(3, len(Link.objects.filter(parent__workflow=workflow)))
  845. 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)
  846. workflow.delete()
  847. def test_import_workflow_pig(self):
  848. workflow = Workflow.objects.new_workflow(self.user)
  849. workflow.save()
  850. f = open('apps/oozie/src/oozie/test_data/0.4/test-pig.xml')
  851. import_workflow(workflow, f.read(), schema_version=0.4)
  852. f.close()
  853. workflow.save()
  854. node = Node.objects.get(workflow=workflow, node_type='pig').get_full_node()
  855. assert_equal(4, len(Node.objects.filter(workflow=workflow)))
  856. assert_equal(3, len(Link.objects.filter(parent__workflow=workflow)))
  857. assert_equal('aggregate.pig', node.script_path)
  858. 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)
  859. workflow.delete()
  860. def test_import_workflow_sqoop(self):
  861. workflow = Workflow.objects.new_workflow(self.user)
  862. workflow.save()
  863. f = open('apps/oozie/src/oozie/test_data/0.4/test-sqoop.0.2.xml')
  864. import_workflow(workflow, f.read(), schema_version=0.4)
  865. f.close()
  866. workflow.save()
  867. assert_equal(4, len(Node.objects.filter(workflow=workflow)))
  868. assert_equal(3, len(Link.objects.filter(parent__workflow=workflow)))
  869. node = Node.objects.get(workflow=workflow, node_type='sqoop').get_full_node()
  870. assert_equal('["db.hsqldb.properties#db.hsqldb.properties","db.hsqldb.script#db.hsqldb.script"]', node.files)
  871. assert_equal('import --connect jdbc:hsqldb:file:db.hsqldb --table TT --target-dir ${output} -m 1', node.script_path)
  872. workflow.delete()
  873. def test_import_workflow_java(self):
  874. workflow = Workflow.objects.new_workflow(self.user)
  875. workflow.save()
  876. f = open('apps/oozie/src/oozie/test_data/0.4/test-java.xml')
  877. import_workflow(workflow, f.read(), schema_version=0.4)
  878. f.close()
  879. workflow.save()
  880. assert_equal(5, len(Node.objects.filter(workflow=workflow)))
  881. assert_equal(5, len(Link.objects.filter(parent__workflow=workflow)))
  882. nodes = [Node.objects.filter(workflow=workflow, node_type='java')[0].get_full_node(),
  883. Node.objects.filter(workflow=workflow, node_type='java')[1].get_full_node()]
  884. assert_equal('org.apache.hadoop.examples.terasort.TeraGen', nodes[0].main_class)
  885. assert_equal('["${records}","${output_dir}/teragen"]', nodes[0].args)
  886. assert_equal('org.apache.hadoop.examples.terasort.TeraSort', nodes[1].main_class)
  887. assert_equal('["${output_dir}/teragen","${output_dir}/terasort"]', nodes[1].args)
  888. workflow.delete()
  889. class TestPermissions(OozieBase):
  890. def setUp(self):
  891. self.c = make_logged_in_client()
  892. self.wf = create_workflow(self.c)
  893. self.setup_simple_workflow()
  894. def test_workflow_permissions(self):
  895. response = self.c.get(reverse('oozie:edit_workflow', args=[self.wf.id]))
  896. assert_true('Editor' in response.content, response.content)
  897. assert_true('Save' in response.content, response.content)
  898. assert_false(self.wf.is_shared)
  899. # Login as someone else
  900. client_not_me = make_logged_in_client(username='not_me', is_superuser=False, groupname='test')
  901. grant_access("not_me", "test", "oozie")
  902. # List
  903. finish = SHARE_JOBS.set_for_testing(True)
  904. try:
  905. response = client_not_me.get(reverse('oozie:list_workflows'))
  906. assert_false('wf-name-1' in response.content, response.content)
  907. finally:
  908. finish()
  909. finish = SHARE_JOBS.set_for_testing(False)
  910. try:
  911. response = client_not_me.get(reverse('oozie:list_workflows'))
  912. assert_false('wf-name-1' in response.content, response.content)
  913. finally:
  914. finish()
  915. # View
  916. finish = SHARE_JOBS.set_for_testing(True)
  917. try:
  918. response = client_not_me.get(reverse('oozie:edit_workflow', args=[self.wf.id]))
  919. assert_true('Permission denied' in response.content, response.content)
  920. finally:
  921. finish()
  922. finish = SHARE_JOBS.set_for_testing(False)
  923. try:
  924. response = client_not_me.get(reverse('oozie:edit_workflow', args=[self.wf.id]))
  925. assert_true('Permission denied' in response.content, response.content)
  926. finally:
  927. finish()
  928. # Share it !
  929. self.wf = Workflow.objects.get(name='wf-name-1')
  930. self.wf.is_shared = True
  931. self.wf.save()
  932. Workflow.objects.check_workspace(self.wf, self.cluster.fs)
  933. # List
  934. finish = SHARE_JOBS.set_for_testing(True)
  935. try:
  936. response = client_not_me.get(reverse('oozie:list_workflows'))
  937. assert_equal(200, response.status_code)
  938. assert_true('wf-name-1' in response.content, response.content)
  939. finally:
  940. finish()
  941. # View
  942. finish = SHARE_JOBS.set_for_testing(True)
  943. try:
  944. response = client_not_me.get(reverse('oozie:edit_workflow', args=[self.wf.id]))
  945. assert_false('Permission denied' in response.content, response.content)
  946. assert_true('Save' in response.content, response.content)
  947. finally:
  948. finish()
  949. finish = SHARE_JOBS.set_for_testing(False)
  950. try:
  951. response = client_not_me.get(reverse('oozie:edit_workflow', args=[self.wf.id]))
  952. assert_true('Permission denied' in response.content, response.content)
  953. finally:
  954. finish()
  955. # Submit
  956. finish = SHARE_JOBS.set_for_testing(False)
  957. try:
  958. response = client_not_me.post(reverse('oozie:submit_workflow', args=[self.wf.id]))
  959. assert_true('Permission denied' in response.content, response.content)
  960. finally:
  961. finish()
  962. finish = SHARE_JOBS.set_for_testing(True)
  963. try:
  964. try:
  965. response = client_not_me.post(reverse('oozie:submit_workflow', args=[self.wf.id]))
  966. assert_false('Permission denied' in response.content, response.content)
  967. except IOError:
  968. pass
  969. finally:
  970. finish()
  971. # Delete
  972. finish = SHARE_JOBS.set_for_testing(False)
  973. try:
  974. response = client_not_me.post(reverse('oozie:delete_workflow', args=[self.wf.id]))
  975. assert_true('Permission denied' in response.content, response.content)
  976. finally:
  977. finish()
  978. response = self.c.post(reverse('oozie:delete_workflow', args=[self.wf.id]), follow=True)
  979. assert_equal(200, response.status_code)
  980. def test_coordinator_permissions(self):
  981. coord = create_coordinator(self.wf, self.c)
  982. response = self.c.get(reverse('oozie:edit_coordinator', args=[coord.id]))
  983. assert_true('Editor' in response.content, response.content)
  984. assert_true('value="Save"' in response.content, response.content)
  985. # Login as someone else
  986. client_not_me = make_logged_in_client(username='not_me', is_superuser=False, groupname='test')
  987. grant_access("not_me", "test", "oozie")
  988. # List
  989. finish = SHARE_JOBS.set_for_testing(True)
  990. try:
  991. response = client_not_me.get(reverse('oozie:list_coordinators'))
  992. assert_false('MyCoord' in response.content, response.content)
  993. finally:
  994. finish()
  995. finish = SHARE_JOBS.set_for_testing(False)
  996. try:
  997. response = client_not_me.get(reverse('oozie:list_coordinators'))
  998. assert_false('MyCoord' in response.content, response.content)
  999. finally:
  1000. finish()
  1001. # View
  1002. finish = SHARE_JOBS.set_for_testing(True)
  1003. try:
  1004. response = client_not_me.get(reverse('oozie:edit_coordinator', args=[coord.id]))
  1005. assert_true('Permission denied' in response.content, response.content)
  1006. finally:
  1007. finish()
  1008. finish = SHARE_JOBS.set_for_testing(False)
  1009. try:
  1010. response = client_not_me.get(reverse('oozie:edit_coordinator', args=[coord.id]))
  1011. assert_false('MyCoord' in response.content, response.content)
  1012. finally:
  1013. finish()
  1014. # Share it !
  1015. wf = Workflow.objects.get(id=coord.workflow.id)
  1016. wf.is_shared = True
  1017. wf.save()
  1018. Workflow.objects.check_workspace(wf, self.cluster.fs)
  1019. post = COORDINATOR_DICT.copy()
  1020. post.update({
  1021. u'datainput_set-TOTAL_FORMS': [u'0'], u'datainput_set-INITIAL_FORMS': [u'0'], u'dataset_set-INITIAL_FORMS': [u'0'],
  1022. u'dataoutput_set-INITIAL_FORMS': [u'0'], u'datainput_set-MAX_NUM_FORMS': [u'0'], u'output-MAX_NUM_FORMS': [u''],
  1023. u'output-INITIAL_FORMS': [u'0'], u'dataoutput_set-TOTAL_FORMS': [u'0'], u'input-TOTAL_FORMS': [u'0'],
  1024. u'dataset_set-MAX_NUM_FORMS': [u'0'], u'dataoutput_set-MAX_NUM_FORMS': [u'0'], u'input-MAX_NUM_FORMS': [u''],
  1025. u'dataset_set-TOTAL_FORMS': [u'0'], u'input-INITIAL_FORMS': [u'0'], u'output-TOTAL_FORMS': [u'0']})
  1026. post['is_shared'] = [u'on']
  1027. post['workflow'] = coord.workflow.id
  1028. self.c.post(reverse('oozie:edit_coordinator', args=[coord.id]), post)
  1029. coord = Coordinator.objects.get(id=coord.id)
  1030. assert_true(coord.is_shared)
  1031. # List
  1032. finish = SHARE_JOBS.set_for_testing(True)
  1033. try:
  1034. response = client_not_me.get(reverse('oozie:list_coordinators'))
  1035. assert_equal(200, response.status_code)
  1036. assert_true('MyCoord' in response.content, response.content)
  1037. finally:
  1038. finish()
  1039. # View
  1040. finish = SHARE_JOBS.set_for_testing(True)
  1041. try:
  1042. response = client_not_me.get(reverse('oozie:edit_coordinator', args=[coord.id]))
  1043. assert_false('Permission denied' in response.content, response.content)
  1044. assert_false('value="Save"' in response.content, response.content)
  1045. finally:
  1046. finish()
  1047. finish = SHARE_JOBS.set_for_testing(False)
  1048. try:
  1049. response = client_not_me.get(reverse('oozie:edit_coordinator', args=[coord.id]))
  1050. assert_true('Permission denied' in response.content, response.content)
  1051. finally:
  1052. finish()
  1053. # Edit
  1054. finish = SHARE_JOBS.set_for_testing(True)
  1055. try:
  1056. response = client_not_me.post(reverse('oozie:edit_coordinator', args=[coord.id]))
  1057. assert_false('MyCoord' in response.content, response.content)
  1058. assert_true('Not allowed' in response.content, response.content)
  1059. finally:
  1060. finish()
  1061. # Submit
  1062. finish = SHARE_JOBS.set_for_testing(False)
  1063. try:
  1064. response = client_not_me.post(reverse('oozie:submit_coordinator', args=[coord.id]))
  1065. assert_true('Permission denied' in response.content, response.content)
  1066. finally:
  1067. finish()
  1068. finish = SHARE_JOBS.set_for_testing(True)
  1069. try:
  1070. try:
  1071. response = client_not_me.post(reverse('oozie:submit_coordinator', args=[coord.id]))
  1072. assert_false('Permission denied' in response.content, response.content)
  1073. except IOError:
  1074. pass
  1075. finally:
  1076. finish()
  1077. # Resubmit
  1078. finish = SHARE_JOBS.set_for_testing(False)
  1079. try:
  1080. history, created = History.objects.get_or_create(job=coord, oozie_job_id=MockOozieApi.COORDINATOR_IDS[0],
  1081. defaults={'submitter': User.objects.get(username='test'), 'properties': '[]'})
  1082. oozie_job_id = history.oozie_job_id
  1083. response = client_not_me.post(reverse('oozie:resubmit_coordinator', args=[oozie_job_id]))
  1084. assert_true('Permission denied' in response.content, response.content)
  1085. finally:
  1086. finish()
  1087. finish = SHARE_JOBS.set_for_testing(True)
  1088. try:
  1089. try:
  1090. history, created = History.objects.get_or_create(job=coord, oozie_job_id=MockOozieApi.COORDINATOR_IDS[0],
  1091. defaults={'submitter': User.objects.get(username='test'), 'properties': '[]'})
  1092. oozie_job_id = history.oozie_job_id
  1093. response = client_not_me.post(reverse('oozie:resubmit_coordinator', args=[oozie_job_id]))
  1094. assert_false('Permission denied' in response.content, response.content)
  1095. except IOError:
  1096. pass
  1097. finally:
  1098. finish()
  1099. # Delete
  1100. finish = SHARE_JOBS.set_for_testing(False)
  1101. try:
  1102. response = client_not_me.post(reverse('oozie:delete_coordinator', args=[coord.id]))
  1103. assert_true('Permission denied' in response.content, response.content)
  1104. finally:
  1105. finish()
  1106. response = self.c.post(reverse('oozie:delete_coordinator', args=[coord.id]), follow=True)
  1107. assert_equal(200, response.status_code)
  1108. class TestEditorWithOozie(OozieBase):
  1109. def setUp(self):
  1110. self.c = make_logged_in_client()
  1111. self.wf = create_workflow(self.c)
  1112. self.setup_simple_workflow()
  1113. def tearDown(self):
  1114. self.wf.delete()
  1115. def test_create_workflow(self):
  1116. dir_stat = self.cluster.fs.stats(self.wf.deployment_dir)
  1117. assert_equal('test', dir_stat.user)
  1118. assert_equal('hue', dir_stat.group)
  1119. assert_equal('40711', '%o' % dir_stat.mode)
  1120. def test_clone_workflow(self):
  1121. workflow_count = Workflow.objects.count()
  1122. response = self.c.post(reverse('oozie:clone_workflow', args=[self.wf.id]), {}, follow=True)
  1123. assert_equal(workflow_count + 1, Workflow.objects.count(), response)
  1124. wf2 = Workflow.objects.latest('id')
  1125. assert_not_equal(self.wf.id, wf2.id)
  1126. assert_equal(self.wf.node_set.count(), wf2.node_set.count())
  1127. node_ids = set(self.wf.node_set.values_list('id', flat=True))
  1128. for node in wf2.node_set.all():
  1129. assert_false(node.id in node_ids)
  1130. assert_not_equal(self.wf.deployment_dir, wf2.deployment_dir)
  1131. assert_not_equal('', wf2.deployment_dir)
  1132. def test_import_action(self):
  1133. raise SkipTest
  1134. # Setup jobsub examples
  1135. if not jobsub_setup.Command().has_been_setup():
  1136. jobsub_setup.Command().handle()
  1137. # There should be 3 from examples
  1138. jobsub_design = OozieDesign.objects.all()[0]
  1139. node_size = len(Node.objects.all())
  1140. kwargs = dict(workflow=self.wf.id, parent_action_id=self.wf.end.get_parents()[0].id)
  1141. response = self.c.post(reverse('oozie:import_action', kwargs=kwargs), {'action_id': jobsub_design.id})
  1142. assert_equal(302, response.status_code)
  1143. assert_equal(node_size + 1, len(Node.objects.all()))
  1144. # There should now be an imported action at the end of Node list
  1145. # Need to test properties to make sure we got it right
  1146. # Must also make sure that jobsub field values are translated
  1147. translation_regex = re.compile('(?<!\$)\$(\w+)')
  1148. node = Node.objects.all()[len(Node.objects.all())-1].get_full_node()
  1149. for field in node.PARAM_FIELDS:
  1150. assert_equal(translation_regex.sub(r'${\1}', getattr(jobsub_design.get_root_action(), field)), getattr(node, field))
  1151. def test_import_workflow(self):
  1152. workflow_count = Workflow.objects.count()
  1153. # Create
  1154. filename = os.path.abspath(os.path.dirname(__file__) + "/test_data/0.4/test-mapreduce.xml")
  1155. fh = open(filename)
  1156. response = self.c.post(reverse('oozie:import_workflow'), {
  1157. 'job_xml': [''],
  1158. 'name': ['test_workflow'],
  1159. 'parameters': ['[{"name":"oozie.use.system.libpath","value":"true"}]'],
  1160. 'deployment_dir': [''],
  1161. 'job_properties': ['[]'],
  1162. 'schema_version': ['0.4'],
  1163. 'definition_file': [fh],
  1164. 'description': ['']
  1165. }, follow=True)
  1166. fh.close()
  1167. assert_equal(workflow_count + 1, Workflow.objects.count(), response)
  1168. class TestOozieSubmissions(OozieBase):
  1169. def test_submit_mapreduce_action(self):
  1170. wf = Workflow.objects.get(name='MapReduce')
  1171. post_data = {u'form-MAX_NUM_FORMS': [u''], u'form-INITIAL_FORMS': [u'1'],
  1172. u'form-0-name': [u'REDUCER_SLEEP_TIME'], u'form-0-value': [u'1'], u'form-TOTAL_FORMS': [u'1']}
  1173. response = self.c.post(reverse('oozie:submit_workflow', args=[wf.id]), data=post_data, follow=True)
  1174. job = OozieServerProvider.wait_until_completion(response.context['oozie_workflow'].id)
  1175. assert_equal('SUCCEEDED', job.status)
  1176. # Rerun with default options
  1177. post_data.update({u'rerun_form_choice': [u'skip_nodes']})
  1178. response = self.c.post(reverse('oozie:rerun_oozie_job', kwargs={'job_id': job.id, 'app_path': job.appPath}), data=post_data, follow=True)
  1179. job = OozieServerProvider.wait_until_completion(response.context['oozie_workflow'].id)
  1180. assert_equal('SUCCEEDED', job.status)
  1181. # Rerun with skip OK actions skipped
  1182. post_data.update({u'rerun_form_choice': [u'skip_nodes'], u'skip_nodes': [u'Sleep']})
  1183. response = self.c.post(reverse('oozie:rerun_oozie_job', kwargs={'job_id': job.id, 'app_path': job.appPath}), data=post_data, follow=True)
  1184. job = OozieServerProvider.wait_until_completion(response.context['oozie_workflow'].id)
  1185. assert_equal('SUCCEEDED', job.status)
  1186. # Rerun with failed nodes too
  1187. post_data.update({u'rerun_form_choice': [u'failed_nodes']})
  1188. response = self.c.post(reverse('oozie:rerun_oozie_job', kwargs={'job_id': job.id, 'app_path': job.appPath}), data=post_data, follow=True)
  1189. job = OozieServerProvider.wait_until_completion(response.context['oozie_workflow'].id)
  1190. def test_submit_java_action(self):
  1191. wf = Workflow.objects.get(name='Sequential Java')
  1192. response = self.c.post(reverse('oozie:submit_workflow', args=[wf.id]),
  1193. data={u'form-MAX_NUM_FORMS': [u''],
  1194. u'form-0-name': [u'records'], u'form-0-value': [u'10'],
  1195. u'form-1-name': [u' output_dir '], u'form-1-value': [u'${nameNode}/user/test/out/terasort'],
  1196. u'form-INITIAL_FORMS': [u'2'], u'form-TOTAL_FORMS': [u'2']},
  1197. follow=True)
  1198. job = OozieServerProvider.wait_until_completion(response.context['oozie_workflow'].id)
  1199. assert_equal('SUCCEEDED', job.status)
  1200. def test_submit_distcp_action(self):
  1201. wf = Workflow.objects.get(name='DistCp')
  1202. response = self.c.post(reverse('oozie:submit_workflow', args=[wf.id]),
  1203. data= {u'form-MAX_NUM_FORMS': [u''], u'form-TOTAL_FORMS': [u'3'], u'form-INITIAL_FORMS': [u'3'],
  1204. u'form-0-name': [u'oozie.use.system.libpath'], u'form-0-value': [u'true'],
  1205. u'form-1-name': [u'OUTPUT'], u'form-1-value': [u'${nameNode}/user/test/out/distcp'],
  1206. u'form-2-name': [u'MAP_NUMBER'], u'form-2-value': [u'5'],
  1207. },
  1208. follow=True)
  1209. job = OozieServerProvider.wait_until_completion(response.context['oozie_workflow'].id)
  1210. assert_equal('SUCCEEDED', job.status)
  1211. class TestDashboardNoMocking:
  1212. def test_oozie_not_running_message(self):
  1213. c = make_logged_in_client(is_superuser=False)
  1214. grant_access("test", "test", "oozie")
  1215. finish = OOZIE_URL.set_for_testing('http://not_localhost:11000/bad')
  1216. try:
  1217. response = c.get(reverse('oozie:list_oozie_workflows'))
  1218. assert_true('The Oozie server is not running' in response.content, response.content)
  1219. finally:
  1220. finish()
  1221. class TestDashboard(OozieMockBase):
  1222. def test_manage_workflow_dashboard(self):
  1223. # Kill button in response
  1224. response = self.c.get(reverse('oozie:list_oozie_workflow', args=[MockOozieApi.WORKFLOW_IDS[0]]), {}, follow=True)
  1225. assert_true(('%s/kill' % MockOozieApi.WORKFLOW_IDS[0]) in response.content, response.content)
  1226. assert_false('Rerun' in response.content, response.content)
  1227. # Rerun button in response
  1228. response = self.c.get(reverse('oozie:list_oozie_workflow', args=[MockOozieApi.WORKFLOW_IDS[1]]), {}, follow=True)
  1229. assert_false(('%s/kill' % MockOozieApi.WORKFLOW_IDS[1]) in response.content, response.content)
  1230. assert_true('Rerun' in response.content, response.content)
  1231. def test_manage_coordinator_dashboard(self):
  1232. # Kill button in response
  1233. response = self.c.get(reverse('oozie:list_oozie_coordinator', args=[MockOozieApi.COORDINATOR_IDS[0]]), {}, follow=True)
  1234. assert_true(('%s/kill' % MockOozieApi.COORDINATOR_IDS[0]) in response.content, response.content)
  1235. assert_false('Rerun' in response.content, response.content)
  1236. # Rerun button in response
  1237. response = self.c.get(reverse('oozie:list_oozie_coordinator', args=[MockOozieApi.COORDINATOR_IDS[1]]), {}, follow=True)
  1238. assert_false(('%s/kill' % MockOozieApi.COORDINATOR_IDS[1]) in response.content, response.content)
  1239. assert_true('Resubmit' in response.content, response.content)
  1240. def test_list_workflows(self):
  1241. response = self.c.get(reverse('oozie:list_oozie_workflows'))
  1242. for wf_id in MockOozieApi.WORKFLOW_IDS:
  1243. assert_true(wf_id in response.content, response.content)
  1244. def test_list_coordinators(self):
  1245. response = self.c.get(reverse('oozie:list_oozie_coordinators'))
  1246. for coord_id in MockOozieApi.COORDINATOR_IDS:
  1247. assert_true(coord_id in response.content, response.content)
  1248. def test_list_workflow(self):
  1249. response = self.c.get(reverse('oozie:list_oozie_workflow', args=[MockOozieApi.WORKFLOW_IDS[0]]))
  1250. assert_true('Workflow WordCount1' in response.content, response.content)
  1251. assert_true('Workflow' in response.content, response.content)
  1252. response = self.c.get(reverse('oozie:list_oozie_workflow', args=[MockOozieApi.WORKFLOW_IDS[0], MockOozieApi.COORDINATOR_IDS[0]]))
  1253. assert_true('Workflow WordCount1' in response.content, response.content)
  1254. assert_true('Workflow' in response.content, response.content)
  1255. assert_true('DailyWordCount1' in response.content, response.content)
  1256. assert_true('Coordinator' in response.content, response.content)
  1257. def test_list_coordinator(self):
  1258. response = self.c.get(reverse('oozie:list_oozie_coordinator', args=[MockOozieApi.COORDINATOR_IDS[0]]))
  1259. assert_true('Coordinator DailyWordCount1' in response.content, response.content)
  1260. assert_true('Workflow' in response.content, response.content)
  1261. def test_manage_oozie_jobs(self):
  1262. try:
  1263. self.c.get(reverse('oozie:manage_oozie_jobs', args=[MockOozieApi.COORDINATOR_IDS[0], 'kill']))
  1264. assert False
  1265. except:
  1266. pass
  1267. response = self.c.post(reverse('oozie:manage_oozie_jobs', args=[MockOozieApi.COORDINATOR_IDS[0], 'kill']))
  1268. data = json.loads(response.content)
  1269. assert_equal(0, data['status'])
  1270. def test_workflows_permissions(self):
  1271. response = self.c.get(reverse('oozie:list_oozie_workflows'))
  1272. assert_true('WordCount1' in response.content, response.content)
  1273. # Rerun
  1274. response = self.c.get(reverse('oozie:rerun_oozie_job', kwargs={'job_id': MockOozieApi.WORKFLOW_IDS[0],
  1275. 'app_path': MockOozieApi.JSON_WORKFLOW_LIST[0]['appPath']}))
  1276. assert_false('Permission denied.' in response.content, response.content)
  1277. # Login as someone else
  1278. client_not_me = make_logged_in_client(username='not_me', is_superuser=False, groupname='test', recreate=True)
  1279. grant_access("not_me", "not_me", "oozie")
  1280. response = client_not_me.get(reverse('oozie:list_oozie_workflows'))
  1281. assert_false('WordCount1' in response.content, response.content)
  1282. # Rerun
  1283. response = client_not_me.get(reverse('oozie:rerun_oozie_job', kwargs={'job_id': MockOozieApi.WORKFLOW_IDS[0],
  1284. 'app_path': MockOozieApi.JSON_WORKFLOW_LIST[0]['appPath']}))
  1285. assert_true('Permission denied.' in response.content, response.content)
  1286. # Add read only access
  1287. add_permission("not_me", "dashboard_jobs_access", "dashboard_jobs_access", "oozie")
  1288. response = client_not_me.get(reverse('oozie:list_oozie_workflows'))
  1289. assert_true('WordCount1' in response.content, response.content)
  1290. # Rerun
  1291. response = client_not_me.get(reverse('oozie:rerun_oozie_job', kwargs={'job_id': MockOozieApi.WORKFLOW_IDS[0],
  1292. 'app_path': MockOozieApi.JSON_WORKFLOW_LIST[0]['appPath']}))
  1293. assert_false('Permission denied.' in response.content, response.content)
  1294. def test_workflow_permissions(self):
  1295. response = self.c.get(reverse('oozie:list_oozie_workflow', args=[MockOozieApi.WORKFLOW_IDS[0]]))
  1296. assert_true('WordCount1' in response.content, response.content)
  1297. assert_false('Permission denied' in response.content, response.content)
  1298. response = self.c.get(reverse('oozie:list_oozie_workflow_action', args=['XXX']))
  1299. assert_false('Permission denied' in response.content, response.content)
  1300. # Login as someone else
  1301. client_not_me = make_logged_in_client(username='not_me', is_superuser=False, groupname='test', recreate=True)
  1302. grant_access("not_me", "not_me", "oozie")
  1303. response = client_not_me.get(reverse('oozie:list_oozie_workflow', args=[MockOozieApi.WORKFLOW_IDS[0]]))
  1304. assert_true('Permission denied' in response.content, response.content)
  1305. response = client_not_me.get(reverse('oozie:list_oozie_workflow_action', args=['XXX']))
  1306. assert_true('Permission denied' in response.content, response.content)
  1307. # Add read only access
  1308. add_permission("not_me", "dashboard_jobs_access", "dashboard_jobs_access", "oozie")
  1309. response = client_not_me.get(reverse('oozie:list_oozie_workflow', args=[MockOozieApi.WORKFLOW_IDS[0]]))
  1310. assert_false('Permission denied' in response.content, response.content)
  1311. def test_coordinators_permissions(self):
  1312. response = self.c.get(reverse('oozie:list_oozie_coordinators'))
  1313. assert_true('DailyWordCount1' in response.content, response.content)
  1314. # Login as someone else
  1315. client_not_me = make_logged_in_client(username='not_me', is_superuser=False, groupname='test', recreate=True)
  1316. grant_access("not_me", "not_me", "oozie")
  1317. response = client_not_me.get(reverse('oozie:list_oozie_coordinators'))
  1318. assert_false('DailyWordCount1' in response.content, response.content)
  1319. # Add read only access
  1320. add_permission("not_me", "dashboard_jobs_access", "dashboard_jobs_access", "oozie")
  1321. response = client_not_me.get(reverse('oozie:list_oozie_coordinators'))
  1322. assert_true('DailyWordCount1' in response.content, response.content)
  1323. def test_coordinator_permissions(self):
  1324. response = self.c.get(reverse('oozie:list_oozie_coordinator', args=[MockOozieApi.COORDINATOR_IDS[0]]))
  1325. assert_true('DailyWordCount1' in response.content, response.content)
  1326. assert_false('Permission denied' in response.content, response.content)
  1327. # Login as someone else
  1328. client_not_me = make_logged_in_client(username='not_me', is_superuser=False, groupname='test', recreate=True)
  1329. grant_access("not_me", "not_me", "oozie")
  1330. response = client_not_me.get(reverse('oozie:list_oozie_coordinator', args=[MockOozieApi.COORDINATOR_IDS[0]]))
  1331. assert_true('Permission denied' in response.content, response.content)
  1332. # Add read only access
  1333. add_permission("not_me", "dashboard_jobs_access", "dashboard_jobs_access", "oozie")
  1334. response = client_not_me.get(reverse('oozie:list_oozie_coordinator', args=[MockOozieApi.COORDINATOR_IDS[0]]))
  1335. assert_false('Permission denied' in response.content, response.content)
  1336. class TestUtils(OozieMockBase):
  1337. def setUp(self):
  1338. OozieMockBase.setUp(self)
  1339. # When updating wf, update wf_json as well!
  1340. self.wf = Workflow.objects.get(name='wf-name-1')
  1341. def test_workflow_to_dict(self):
  1342. workflow_dict = workflow_to_dict(self.wf)
  1343. # Test properties
  1344. assert_true('job_xml' in workflow_dict, workflow_dict)
  1345. assert_true('is_shared' in workflow_dict, workflow_dict)
  1346. assert_true('end' in workflow_dict, workflow_dict)
  1347. assert_true('description' in workflow_dict, workflow_dict)
  1348. assert_true('parameters' in workflow_dict, workflow_dict)
  1349. assert_true('is_single' in workflow_dict, workflow_dict)
  1350. assert_true('deployment_dir' in workflow_dict, workflow_dict)
  1351. assert_true('schema_version' in workflow_dict, workflow_dict)
  1352. assert_true('job_properties' in workflow_dict, workflow_dict)
  1353. assert_true('start' in workflow_dict, workflow_dict)
  1354. assert_true('nodes' in workflow_dict, workflow_dict)
  1355. assert_true('id' in workflow_dict, workflow_dict)
  1356. assert_true('name' in workflow_dict, workflow_dict)
  1357. # Check links
  1358. for node in workflow_dict['nodes']:
  1359. assert_true('child_links' in node, node)
  1360. for link in node['child_links']:
  1361. assert_true('name' in link, link)
  1362. assert_true('comment' in link, link)
  1363. assert_true('parent' in link, link)
  1364. assert_true('child' in link, link)
  1365. def test_model_to_dict(self):
  1366. node_dict = model_to_dict(self.wf.node_set.filter(node_type='start')[0])
  1367. # Test properties
  1368. assert_true('id' in node_dict)
  1369. assert_true('name' in node_dict)
  1370. assert_true('description' in node_dict)
  1371. assert_true('node_type' in node_dict)
  1372. assert_true('workflow' in node_dict)
  1373. # Utils
  1374. WORKFLOW_DICT = {u'deployment_dir': [u''], u'name': [u'wf-name-1'], u'description': [u''],
  1375. u'schema_version': [u'uri:oozie:workflow:0.2'],
  1376. u'parameters': [u'[{"name":"market","value":"US"}]'],
  1377. u'job_xml': [u'jobconf.xml'],
  1378. u'job_properties': [u'[{"name":"sleep-all","value":"${SLEEP}"}]']
  1379. }
  1380. COORDINATOR_DICT = {u'name': [u'MyCoord'], u'description': [u'Description of my coodinator'],
  1381. u'workflow': [u'1'],
  1382. u'frequency_number': [u'1'], u'frequency_unit': [u'days'],
  1383. u'start_0': [u'07/01/2012'], u'start_1': [u'12:00 AM'],
  1384. u'end_0': [u'07/04/2012'], u'end_1': [u'12:00 AM'],
  1385. u'timezone': [u'America/Los_Angeles'],
  1386. u'parameters': [u'[{"name":"market","value":"US,France"}]'],
  1387. u'timeout': [u'100'],
  1388. u'concurrency': [u'3'],
  1389. u'execution': [u'FIFO'],
  1390. u'throttle': [u'10'],
  1391. u'schema_version': [u'uri:oozie:coordinator:0.1']
  1392. }
  1393. def remove_related_fields(workflow_dict):
  1394. """
  1395. workflow_dict is a workflow that has been converted into a dictionary via workflow_to_dict
  1396. """
  1397. del workflow_dict['owner']
  1398. del workflow_dict['job_ptr']
  1399. for node in workflow_dict['nodes']:
  1400. del node['node_ptr']
  1401. return workflow_dict
  1402. def add_node(workflow, name, node_type, parents, attrs={}):
  1403. """
  1404. create a node of type node_type and associate the listed parents.
  1405. """
  1406. NodeClass = NODE_TYPES[node_type]
  1407. node = NodeClass(workflow=workflow, node_type=node_type, name=name)
  1408. for attr in attrs:
  1409. setattr(node, attr, attrs[attr])
  1410. node.save()
  1411. # Add parent
  1412. # If skipped, remember to preserve order: regular links first, then error link
  1413. if parents:
  1414. for parent in parents:
  1415. name = 'ok'
  1416. if parent.node_type == 'start' or parent.node_type == 'join':
  1417. name = 'to'
  1418. elif parent.node_type == 'fork' or parent.node_type == 'decision':
  1419. name = 'start'
  1420. link = Link(parent=parent, child=node, name=name)
  1421. link.save()
  1422. # Create error link
  1423. if node_type != 'fork' and node_type != 'decision' and node_type != 'join':
  1424. link = Link(parent=node, child=Kill.objects.get(name='kill', workflow=workflow), name="error")
  1425. link.save()
  1426. return node
  1427. def create_workflow(client, workflow_dict=WORKFLOW_DICT):
  1428. name = str(workflow_dict['name'][0])
  1429. Node.objects.filter(workflow__name=name).delete()
  1430. Workflow.objects.filter(name=name).delete()
  1431. workflow_count = Workflow.objects.count()
  1432. response = client.get(reverse('oozie:create_workflow'))
  1433. assert_equal(workflow_count, Workflow.objects.count(), response)
  1434. response = client.post(reverse('oozie:create_workflow'), workflow_dict, follow=True)
  1435. assert_equal(200, response.status_code)
  1436. assert_equal(workflow_count + 1, Workflow.objects.count(), response)
  1437. wf = Workflow.objects.get(name=name)
  1438. assert_not_equal('', wf.deployment_dir)
  1439. return wf
  1440. def create_coordinator(workflow, client):
  1441. coord_count = Coordinator.objects.count()
  1442. response = client.get(reverse('oozie:create_coordinator'))
  1443. assert_equal(coord_count, Coordinator.objects.count(), response)
  1444. post = COORDINATOR_DICT.copy()
  1445. post['workflow'] = workflow.id
  1446. response = client.post(reverse('oozie:create_coordinator'), post)
  1447. assert_equal(coord_count + 1, Coordinator.objects.count(), response)
  1448. return Coordinator.objects.get(name='MyCoord')
  1449. def create_dataset(coord, client):
  1450. response = client.post(reverse('oozie:create_coordinator_dataset', args=[coord.id]), {
  1451. u'create-name': [u'MyDataset'], u'create-frequency_number': [u'1'], u'create-frequency_unit': [u'days'],
  1452. u'create-uri': [u'/data/${YEAR}${MONTH}${DAY}'],
  1453. u'create-start_0': [u'07/01/2012'], u'create-start_1': [u'12:00 AM'],
  1454. u'create-timezone': [u'America/Los_Angeles'], u'create-done_flag': [u''],
  1455. u'create-description': [u'']})
  1456. data = json.loads(response.content)
  1457. assert_equal(0, data['status'], data['data'])
  1458. def create_coordinator_data(coord, client):
  1459. response = client.post(reverse('oozie:create_coordinator_data', args=[coord.id, 'input']),
  1460. {u'input-name': [u'input_dir'], u'input-dataset': [u'1']})
  1461. data = json.loads(response.content)
  1462. assert_equal(0, data['status'], data['data'])