浏览代码

HUE-4541 [security] fixing Hue job browser - Kerberos mutual authentication error in Hue

Prakash Ranade 9 年之前
父节点
当前提交
bc0dbe7

+ 5 - 1
apps/jobbrowser/src/jobbrowser/api.py

@@ -174,11 +174,15 @@ class YarnApi(JobBrowserApi):
   """
   def __init__(self, user):
     self.user = user
-    self.resource_manager_api = resource_manager_api.get_resource_manager(user.username)
+    self.resource_manager_api_pool = resource_manager_api.get_resource_manager_pool()
+    self.resource_manager_api = self.resource_manager_api_pool.get(user.username)
     self.mapreduce_api = mapreduce_api.get_mapreduce_api(user.username)
     self.history_server_api = history_server_api.get_history_server_api(user.username)
     self.spark_history_server_api = spark_history_server_api.get_history_server_api()  # Spark HS does not support setuser
 
+  def __del__(self):
+    self.resource_manager_api_pool.put(self.resource_manager_api)
+
   def get_job_link(self, job_id):
     return self.get_job(job_id)
 

+ 4 - 1
apps/jobbrowser/src/jobbrowser/tests.py

@@ -390,7 +390,8 @@ class TestMapReduce2NoHadoop:
   def setUp(self):
     # Beware: Monkey patching
     if not hasattr(resource_manager_api, 'old_get_resource_manager_api'):
-      resource_manager_api.old_get_resource_manager = resource_manager_api.get_resource_manager
+      rm_pool = resource_manager_api.get_resource_manager_pool()
+      resource_manager_api.old_get_resource_manager = rm_pool.get("test2")
     if not hasattr(resource_manager_api, 'old_get_mapreduce_api'):
       mapreduce_api.old_get_mapreduce_api = mapreduce_api.get_mapreduce_api
     if not hasattr(history_server_api, 'old_get_history_server_api'):
@@ -417,6 +418,8 @@ class TestMapReduce2NoHadoop:
 
   def tearDown(self):
     resource_manager_api.get_resource_manager = getattr(resource_manager_api, 'old_get_resource_manager')
+    rm_pool = resource_manager_api.get_resource_manager_pool()
+    rm_pool.put(resource_manager_api.get_resource_manager)
     mapreduce_api.get_mapreduce_api = getattr(mapreduce_api, 'old_get_mapreduce_api')
     history_server_api.get_history_server_api = getattr(history_server_api, 'old_get_history_server_api')
 

+ 4 - 2
apps/jobbrowser/src/jobbrowser/views.py

@@ -72,8 +72,10 @@ def check_job_permission(view_func):
       job = get_api(request.user, request.jt).get_job(jobid=jobid)
     except ApplicationNotRunning, e:
       if e.job.get('state', '').lower() == 'accepted' and 'kill' in request.path:
-        rm_api = resource_manager_api.get_resource_manager(request.user)
+        rm_pool = resource_manager_api.get_resource_manager_pool()
+        rm_api = rm_pool.get(request.user.username)
         job = Application(e.job, rm_api)
+        rm_pool.put(rm_api)
       else:
         # reverse() seems broken, using request.path but beware, it discards GET and POST info
         return job_not_assigned(request, jobid, request.path)
@@ -117,7 +119,7 @@ def jobs(request):
 
   if request.POST.get('format') == 'json':
     try:
-      # Limit number of jobs to be 10,000
+      # Limit number of jobs to be 1000
       jobs = get_api(request.user, request.jt).get_jobs(user=request.user, username=user, state=state, text=text, retired=retired, limit=1000)
     except Exception, ex:
       ex_message = str(ex)

+ 6 - 5
desktop/core/src/desktop/lib/rest/http_client.py

@@ -82,6 +82,7 @@ class HttpClient(object):
     self._exc_class = exc_class or RestException
     self._logger = logger or LOG
     self._session = requests.Session()
+    self._cookies = None
 
   def set_kerberos_auth(self):
     """Set up kerberos auth for the client, based on the current ticket."""
@@ -144,10 +145,7 @@ class HttpClient(object):
         self.logger.warn("GET and DELETE methods do not pass any data. Path '%s'" % path)
         data = None
 
-    request_kwargs = {}
-
-    if not allow_redirects:
-      request_kwargs['allow_redirects'] = False
+    request_kwargs = {'allow_redirects': allow_redirects}
     if headers:
       request_kwargs['headers'] = headers
     if data:
@@ -158,12 +156,15 @@ class HttpClient(object):
     if clear_cookies:
       self._session.cookies.clear()
 
+    if self._cookies:
+      request_kwargs['cookies'] = self._cookies
     try:
       resp = getattr(self._session, http_method.lower())(url, **request_kwargs)
-
       if resp.status_code >= 300:
         resp.raise_for_status()
         raise exceptions.HTTPError(response=resp)
+      # Cache request cookie for the next http_client call.
+      self._cookies = resp.cookies
       return resp
     except (exceptions.ConnectionError,
             exceptions.HTTPError,

+ 0 - 1
desktop/libs/hadoop/src/hadoop/cluster.py

@@ -232,7 +232,6 @@ def get_next_ha_yarncluster():
           if cluster_info['clusterInfo']['haState'] == 'ACTIVE':
             MR_NAME_CACHE = name
             LOG.warn('Picking RM HA: %s' % name)
-            resource_manager_api.API_CACHE = None  # Reset cache
             mapreduce_api.API_CACHE = None
             return (config, rm)
           else:

+ 5 - 2
desktop/libs/hadoop/src/hadoop/yarn/mapreduce_api.py

@@ -24,7 +24,7 @@ from desktop.lib.exceptions_renderable import PopupException
 from desktop.lib.rest.http_client import HttpClient
 from desktop.lib.rest.resource import Resource
 from hadoop import cluster
-from hadoop.yarn.resource_manager_api import get_resource_manager
+from hadoop.yarn.resource_manager_api import get_resource_manager_pool
 
 
 LOG = logging.getLogger(__name__)
@@ -143,4 +143,7 @@ class MapreduceApi(object):
 
   def kill(self, job_id):
     app_id = job_id.replace('job', 'application')
-    get_resource_manager(self.username).kill(app_id) # We need to call the RM
+    pool = get_resource_manager_pool()
+    rmobj = pool.get(self.username)
+    rmobj.kill(app_id) # We need to call the RM
+    pool.put(rmobj)

+ 25 - 28
desktop/libs/hadoop/src/hadoop/yarn/resource_manager_api.py

@@ -18,6 +18,7 @@
 import json
 import logging
 import posixpath
+import Queue
 import threading
 
 from django.utils.translation import ugettext as _
@@ -36,31 +37,42 @@ LOG = logging.getLogger(__name__)
 _API_VERSION = 'v1'
 _JSON_CONTENT_TYPE = 'application/json'
 
-API_CACHE = None
+API_CACHE_POOL = None
 API_CACHE_LOCK = threading.Lock()
 
-
-def get_resource_manager(username):
-  global API_CACHE
-  if API_CACHE is None:
+def get_resource_manager_pool():
+  global API_CACHE_POOL
+  if API_CACHE_POOL is None:
     API_CACHE_LOCK.acquire()
     try:
-      if API_CACHE is None:
+      if API_CACHE_POOL is None:
         yarn_cluster = cluster.get_cluster_conf_for_job_submission()
         if yarn_cluster is None:
           raise PopupException(_('No Resource Manager are available.'))
-        API_CACHE = ResourceManagerApi(yarn_cluster.RESOURCE_MANAGER_API_URL.get(), yarn_cluster.SECURITY_ENABLED.get(), yarn_cluster.SSL_CERT_CA_VERIFY.get())
+        API_CACHE_POOL = ResourceManagerApiPool(yarn_cluster.RESOURCE_MANAGER_API_URL.get(), yarn_cluster.SECURITY_ENABLED.get(), yarn_cluster.SSL_CERT_CA_VERIFY.get())
     finally:
       API_CACHE_LOCK.release()
 
-  API_CACHE.setuser(username) # Set the correct user
-
-  return API_CACHE
-
+  return API_CACHE_POOL
 
 class YarnFailoverOccurred(Exception):
   pass
 
+class ResourceManagerApiPool(object):
+  def __init__(self, api_url, security_enabled, ssl_cert):
+    pool_size = 10
+    self.rmobj_pool = Queue.LifoQueue()
+    for i in range(pool_size):
+      rm_instance = ResourceManagerApi(api_url, security_enabled, ssl_cert)
+      self.rmobj_pool.put(rm_instance)
+
+  def get(self, username):
+    rmobj = self.rmobj_pool.get()
+    rmobj.setuser(username)
+    return rmobj
+
+  def put(self, rmobj):
+    self.rmobj_pool.put(rmobj)
 
 class ResourceManagerApi(object):
 
@@ -160,21 +172,6 @@ class ResourceManagerApi(object):
     response = None
     try:
       response = function(*args, **kwargs)
-    except RestException, e:
-      # YARN-2605: Yarn does not use proper HTTP redirects when the standby RM has
-      # failed back to the master RM.
-      if e.code == 307 and e.message.startswith('This is standby RM'):
-        LOG.info('Received YARN failover redirect response, attempting to resolve redirect.')
-        try:
-          kwargs.update({'allow_redirects': True})
-          response = function(*args, **kwargs)
-        except Exception, e:
-          if response:
-            raise YarnFailoverOccurred(response)
-          else:
-            raise PopupException(_('Failed to resolve YARN RM: %s') % e)
-      else:
-        raise PopupException(_('YARN RM returned a failed response: %s') % e)
-
+    except Exception, e:
+      raise PopupException(_('YARN RM returned a failed response: %s') % e)
     return response
-