Browse Source

[sqoop] Add a live integration test to the Sqoop2 Server

Romain Rigaux 10 years ago
parent
commit
b44fc61

+ 1 - 2
apps/sqoop/src/sqoop/api/job.py

@@ -17,10 +17,10 @@
 
 import json
 import logging
-import socket
 
 from django.utils.encoding import smart_str
 from django.utils.translation import ugettext as _
+from django.views.decorators.cache import never_cache
 
 from sqoop import client, conf
 from sqoop.client.exception import SqoopException
@@ -30,7 +30,6 @@ from desktop.lib.exceptions import StructuredException
 from desktop.lib.rest.http_client import RestException
 from exception import handle_rest_exception
 from utils import list_to_dict
-from django.views.decorators.cache import never_cache
 
 
 __all__ = ['get_jobs', 'create_job', 'update_job', 'job', 'jobs', 'job_clone', 'job_delete', 'job_start', 'job_stop', 'job_status']

+ 24 - 20
apps/sqoop/src/sqoop/test_base.py

@@ -16,21 +16,19 @@
 # limitations under the License.
 
 import atexit
-# import getpass
 import logging
 import os
-import shutil
 import socket
 import subprocess
 import threading
 import time
 
 from django.conf import settings
-
-# from nose.tools import assert_equal, assert_true
+from nose.plugins.skip import SkipTest
 
 from desktop.lib.paths import get_run_root
 from hadoop import pseudo_hdfs4
+from hadoop.pseudo_hdfs4 import is_live_cluster
 
 from sqoop.client import SqoopClient
 from sqoop.conf import SERVER_URL
@@ -54,6 +52,10 @@ class SqoopServerProvider(object):
 
   @classmethod
   def setup_class(cls):
+
+    if not is_live_cluster():
+      raise SkipTest()
+
     cls.cluster = pseudo_hdfs4.shared_cluster()
     cls.client, callback = cls.get_shared_server()
     cls.shutdown = [callback]
@@ -123,24 +125,27 @@ class SqoopServerProvider(object):
     service_lock.acquire()
 
     if not SqoopServerProvider.is_running:
-      LOG.info('\nStarting a Mini Sqoop. Requires "tools/jenkins/jenkins.sh" to be previously ran.\n')
-
-      finish = (
-        SERVER_URL.set_for_testing("http://%s:%s/sqoop" % (socket.getfqdn(), SqoopServerProvider.TEST_PORT)),
-      )
-
       # Setup
       cluster = pseudo_hdfs4.shared_cluster()
 
-      p = cls.start(cluster)
+      if is_live_cluster():
+        finish = ()
+      else:
+        LOG.info('\nStarting a Mini Sqoop. Requires "tools/jenkins/jenkins.sh" to be previously ran.\n')
 
-      def kill():
-        with open(os.path.join(cluster._tmpdir, 'sqoop/sqoop.pid'), 'r') as pidfile:
-          pid = pidfile.read()
-          LOG.info("Killing Sqoop server (pid %s)." % pid)
-          os.kill(int(pid), 9)
-          p.wait()
-      atexit.register(kill)
+        finish = (
+          SERVER_URL.set_for_testing("http://%s:%s/sqoop" % (socket.getfqdn(), SqoopServerProvider.TEST_PORT)),
+        )
+
+        p = cls.start(cluster)
+
+        def kill():
+          with open(os.path.join(cluster._tmpdir, 'sqoop/sqoop.pid'), 'r') as pidfile:
+            pid = pidfile.read()
+            LOG.info("Killing Sqoop server (pid %s)." % pid)
+            os.kill(int(pid), 9)
+            p.wait()
+        atexit.register(kill)
 
       start = time.time()
       started = False
@@ -149,7 +154,6 @@ class SqoopServerProvider(object):
       client = SqoopClient(SERVER_URL.get(), username, language)
 
       while not started and time.time() - start < 60.0:
-        status = None
         try:
           LOG.info('Check Sqoop status...')
           version = client.get_version()
@@ -163,6 +167,7 @@ class SqoopServerProvider(object):
           time.sleep(sleep)
           sleep *= 2
           pass
+
       if not started:
         service_lock.release()
         raise Exception("Sqoop server took too long to come up.")
@@ -174,7 +179,6 @@ class SqoopServerProvider(object):
       callback = shutdown
 
       SqoopServerProvider.is_running = True
-
     else:
       client = SqoopClient(SERVER_URL.get(), username, language)
 

+ 77 - 39
apps/sqoop/src/sqoop/tests.py

@@ -15,48 +15,35 @@
 # limitations under the License.
 
 import logging
+import json
+
+from django.contrib.auth.models import User
+from nose.tools import assert_true, assert_equal
+from nose.plugins.skip import SkipTest
+
+from django.core.urlresolvers import reverse
+from desktop.lib.django_test_util import make_logged_in_client
+from desktop.lib.test_utils import add_to_group, grant_access
 
 from sqoop.client.link import Link
 from sqoop.client.job import Job
 from sqoop.test_base import SqoopServerProvider
 
-from nose.tools import assert_true, assert_equal
-from nose.plugins.skip import SkipTest
-
 
 LOG = logging.getLogger(__name__)
 
 
-LINK_CONFIG_VALUES = {
-  'linkConfig.jdbcDriver': 'org.apache.derby.jdbc.EmbeddedDriver',
-  'linkConfig.String': 'jdbc%3Aderby%3A%2Ftmp%2Ftest',
-  'linkConfig.username': 'abe',
-  'linkConfig.password': 'test',
-  'linkConfig.jdbcProperties': None
-}
-
-FROM_JOB_CONFIG_VALUES = {
-  'fromJobConfig.schemaName': None,
-  'fromJobConfig.tableName': 'test',
-  'fromJobConfig.sql': None,
-  'fromJobConfig.columns': 'name',
-  'fromJobConfig.partitionColumn': 'id',
-  'fromJobConfig.boundaryQuery': None,
-  'fromJobConfig.allowNullValueInPartitionColumn': None
-}
+class TestSqoopServerBase(SqoopServerProvider):
 
-TO_JOB_CONFIG_VALUES = {
-  'toJobConfig.outputFormat': 'TEXT_FILE',
-  'toJobConfig.outputDirectory': '/tmp/test.out',
-  'toJobConfig.storageType': 'HDFS'
-}
+  @classmethod
+  def setup_class(cls):
+    SqoopServerProvider.setup_class()
 
-DRIVER_CONFIG_VALUES = {
-  'throttlingConfig.numExtractor': '3',
-  'throttlingConfig.numLoaders': '3'
-}
+    cls.client = make_logged_in_client(username='test', is_superuser=False)
+    cls.user = User.objects.get(username='test')
+    add_to_group('test')
+    grant_access("test", "test", "sqoop")
 
-class TestSqoopServerBase(SqoopServerProvider):
   def create_link(self, name='test1', connector_id=1):
     link = Link(name, connector_id)
     link.linkConfig = self.client.get_connectors()[0].link_config
@@ -109,9 +96,25 @@ class TestSqoopServerBase(SqoopServerProvider):
     for obj in objects:
       self.delete_sqoop_object(obj)
 
+
+class TestWithSqoopServer(TestSqoopServerBase):
+
+  def test_list_jobs(self):
+
+    resp = self.client.get(reverse('sqoop:jobs'))
+    content = json.loads(resp.content)
+
+    assert_true('jobs' in content, content)
+
+
 class TestSqoopClientLinks(TestSqoopServerBase):
+
+  def setUp(self):
+    raise SkipTest() # These tests are outdated
+
   def test_link(self):
-    raise SkipTest
+    link3 = None
+
     try:
       # Create
       link = self.create_link(name='link1')
@@ -125,22 +128,28 @@ class TestSqoopClientLinks(TestSqoopServerBase):
       link3 = self.client.get_link(link2.id)
       assert_true(link3.id)
       assert_equal(link3.name, link3.name)
-    except:
-      LOG.exception('failed to test link')
-      self.client.delete_link(link3)
+    finally:
+      if link3:
+        self.client.delete_link(link3)
 
   def test_get_links(self):
-    raise SkipTest
+    link = None
+
     try:
       link = self.create_link(name='link2')
       links = self.client.get_links()
       assert_true(len(links) > 0)
     finally:
-      self.client.delete_link(link)
+      if link:
+        self.client.delete_link(link)
+
 
 class TestSqoopClientJobs(TestSqoopServerBase):
+
+  def setUp(self):
+    raise SkipTest() # These tests are outdated
+
   def test_job(self):
-    raise SkipTest
     removable = []
     # Create
     from_link = self.create_link(name='link3from')
@@ -166,12 +175,11 @@ class TestSqoopClientJobs(TestSqoopServerBase):
       self.delete_sqoop_objects(removable)
 
   def test_get_jobs(self):
-    raise SkipTest
     removable = []
     from_link = self.create_link(name='link4from')
     to_link = self.create_link(name='link4to')
-    try:
 
+    try:
       removable.append(from_link)
       removable.append(to_link)
 
@@ -183,3 +191,33 @@ class TestSqoopClientJobs(TestSqoopServerBase):
       assert_true(len(jobs) > 0)
     finally:
       self.delete_sqoop_objects(removable)
+
+
+LINK_CONFIG_VALUES = {
+  'linkConfig.jdbcDriver': 'org.apache.derby.jdbc.EmbeddedDriver',
+  'linkConfig.String': 'jdbc%3Aderby%3A%2Ftmp%2Ftest',
+  'linkConfig.username': 'abe',
+  'linkConfig.password': 'test',
+  'linkConfig.jdbcProperties': None
+}
+
+FROM_JOB_CONFIG_VALUES = {
+  'fromJobConfig.schemaName': None,
+  'fromJobConfig.tableName': 'test',
+  'fromJobConfig.sql': None,
+  'fromJobConfig.columns': 'name',
+  'fromJobConfig.partitionColumn': 'id',
+  'fromJobConfig.boundaryQuery': None,
+  'fromJobConfig.allowNullValueInPartitionColumn': None
+}
+
+TO_JOB_CONFIG_VALUES = {
+  'toJobConfig.outputFormat': 'TEXT_FILE',
+  'toJobConfig.outputDirectory': '/tmp/test.out',
+  'toJobConfig.storageType': 'HDFS'
+}
+
+DRIVER_CONFIG_VALUES = {
+  'throttlingConfig.numExtractor': '3',
+  'throttlingConfig.numLoaders': '3'
+}

+ 1 - 0
desktop/libs/libsentry/src/libsentry/tests.py

@@ -29,6 +29,7 @@ from libsentry.client import SentryClient
 
 
 class TestWithSentry:
+  requires_hadoop = True
 
   @classmethod
   def setup_class(cls):

+ 1 - 0
desktop/libs/libzookeeper/src/libzookeeper/tests.py

@@ -29,6 +29,7 @@ from libzookeeper.conf import zkensemble
 
 
 class TestWithZooKeeper:
+  requires_hadoop = True
 
   @classmethod
   def setup_class(cls):