Эх сурвалжийг харах

HUE-3520 [jb] Use impersonation to access JHS if security is enabled

Jenny Kim 9 жил өмнө
parent
commit
46ad257

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

@@ -174,7 +174,7 @@ class YarnApi(JobBrowserApi):
     self.user = user
     self.resource_manager_api = resource_manager_api.get_resource_manager(user.username)
     self.mapreduce_api = mapreduce_api.get_mapreduce_api(user.username)
-    self.history_server_api = history_server_api.get_history_server_api()
+    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()
 
   def get_job_link(self, job_id):

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

@@ -406,7 +406,7 @@ class TestMapReduce2NoHadoop:
 
     resource_manager_api.get_resource_manager = lambda username: MockResourceManagerApi(username)
     mapreduce_api.get_mapreduce_api = lambda username: MockMapreduceApi(username)
-    history_server_api.get_history_server_api = lambda: HistoryServerApi()
+    history_server_api.get_history_server_api = lambda username: HistoryServerApi(username)
 
     self.finish = [
         YARN_CLUSTERS['default'].SUBMIT_TO.set_for_testing(True),

+ 55 - 19
desktop/libs/hadoop/src/hadoop/yarn/history_server_api.py

@@ -19,13 +19,14 @@ import logging
 import posixpath
 import threading
 
+from desktop.conf import DEFAULT_USER
+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
 
 
 LOG = logging.getLogger(__name__)
-DEFAULT_USER = 'hue'
 
 _API_VERSION = 'v1'
 _JSON_CONTENT_TYPE = 'application/json'
@@ -33,18 +34,26 @@ _JSON_CONTENT_TYPE = 'application/json'
 _api_cache = None
 _api_cache_lock = threading.Lock()
 
+API_CACHE = None
+API_CACHE_LOCK = threading.Lock()
 
-def get_history_server_api():
-  global _api_cache
-  if _api_cache is None:
-    _api_cache_lock.acquire()
+
+def get_history_server_api(username):
+  global API_CACHE
+  if API_CACHE is None:
+    API_CACHE_LOCK.acquire()
     try:
-      if _api_cache is None:
+      if API_CACHE is None:
         yarn_cluster = cluster.get_cluster_conf_for_job_submission()
-        _api_cache = HistoryServerApi(yarn_cluster.HISTORY_SERVER_API_URL.get(), yarn_cluster.SECURITY_ENABLED.get(), yarn_cluster.SSL_CERT_CA_VERIFY.get())
+        if yarn_cluster is None:
+          raise PopupException(_('YARN cluster is not available.'))
+        API_CACHE = HistoryServerApi(yarn_cluster.HISTORY_SERVER_API_URL.get(), yarn_cluster.SECURITY_ENABLED.get(), yarn_cluster.SSL_CERT_CA_VERIFY.get())
     finally:
-      _api_cache_lock.release()
-  return _api_cache
+      API_CACHE_LOCK.release()
+
+  API_CACHE.setuser(username)  # Set the correct user
+
+  return API_CACHE
 
 
 class HistoryServerApi(object):
@@ -54,6 +63,7 @@ class HistoryServerApi(object):
     self._client = HttpClient(self._url, logger=LOG)
     self._root = Resource(self._client)
     self._security_enabled = security_enabled
+    self._thread_local = threading.local()  # To store user info
 
     if self._security_enabled:
       self._client.set_kerberos_auth()
@@ -63,37 +73,63 @@ class HistoryServerApi(object):
   def __str__(self):
     return "HistoryServerApi at %s" % (self._url,)
 
+  def _get_params(self):
+    params = {}
+
+    if self.username != DEFAULT_USER.get():  # We impersonate if needed
+      params['doAs'] = self.username
+      if not self._security_enabled:
+        params['user.name'] = DEFAULT_USER.get()
+
+    return params
+
   @property
   def url(self):
     return self._url
 
+  @property
+  def user(self):
+    return self.username  # Backward compatibility
+
+  @property
+  def username(self):
+    try:
+      return self._thread_local.user
+    except AttributeError:
+      return DEFAULT_USER.get()
+
+  def setuser(self, user):
+    curr = self.user
+    self._thread_local.user = user
+    return curr
+
   def job(self, user, job_id):
-    return self._root.get('mapreduce/jobs/%(job_id)s' % {'job_id': job_id}, headers={'Accept': _JSON_CONTENT_TYPE})
+    return self._root.get('mapreduce/jobs/%(job_id)s' % {'job_id': job_id}, params=self._get_params(), headers={'Accept': _JSON_CONTENT_TYPE})
 
   def counters(self, job_id):
-    return self._root.get('mapreduce/jobs/%(job_id)s/counters' % {'job_id': job_id}, headers={'Accept': _JSON_CONTENT_TYPE})
+    return self._root.get('mapreduce/jobs/%(job_id)s/counters' % {'job_id': job_id}, params=self._get_params(), headers={'Accept': _JSON_CONTENT_TYPE})
 
   def conf(self, job_id):
-    return self._root.get('mapreduce/jobs/%(job_id)s/conf' % {'job_id': job_id}, headers={'Accept': _JSON_CONTENT_TYPE})
+    return self._root.get('mapreduce/jobs/%(job_id)s/conf' % {'job_id': job_id}, params=self._get_params(), headers={'Accept': _JSON_CONTENT_TYPE})
 
   def job_attempts(self, job_id):
-    return self._root.get('mapreduce/jobs/%(job_id)s/jobattempts' % {'job_id': job_id}, headers={'Accept': _JSON_CONTENT_TYPE})
+    return self._root.get('mapreduce/jobs/%(job_id)s/jobattempts' % {'job_id': job_id}, params=self._get_params(), headers={'Accept': _JSON_CONTENT_TYPE})
 
   def tasks(self, job_id):
-    return self._root.get('mapreduce/jobs/%(job_id)s/tasks' % {'job_id': job_id}, headers={'Accept': _JSON_CONTENT_TYPE})
+    return self._root.get('mapreduce/jobs/%(job_id)s/tasks' % {'job_id': job_id}, params=self._get_params(), headers={'Accept': _JSON_CONTENT_TYPE})
 
   def task(self, job_id, task_id):
-    return self._root.get('mapreduce/jobs/%(job_id)s/tasks/%(task_id)s' % {'job_id': job_id, 'task_id': task_id}, headers={'Accept': _JSON_CONTENT_TYPE})
+    return self._root.get('mapreduce/jobs/%(job_id)s/tasks/%(task_id)s' % {'job_id': job_id, 'task_id': task_id}, params=self._get_params(), headers={'Accept': _JSON_CONTENT_TYPE})
 
   def task_attempts(self, job_id, task_id):
-    return self._root.get('mapreduce/jobs/%(job_id)s/tasks/%(task_id)s/attempts' % {'job_id': job_id, 'task_id': task_id}, headers={'Accept': _JSON_CONTENT_TYPE})
+    return self._root.get('mapreduce/jobs/%(job_id)s/tasks/%(task_id)s/attempts' % {'job_id': job_id, 'task_id': task_id}, params=self._get_params(), headers={'Accept': _JSON_CONTENT_TYPE})
 
   def task_counters(self, job_id, task_id):
     job_id = job_id.replace('application', 'job')
-    return self._root.get('mapreduce/jobs/%(job_id)s/tasks/%(task_id)s/counters' % {'job_id': job_id, 'task_id': task_id}, headers={'Accept': _JSON_CONTENT_TYPE})
+    return self._root.get('mapreduce/jobs/%(job_id)s/tasks/%(task_id)s/counters' % {'job_id': job_id, 'task_id': task_id}, params=self._get_params(), headers={'Accept': _JSON_CONTENT_TYPE})
 
   def task_attempt(self, job_id, task_id, attempt_id):
-    return self._root.get('mapreduce/jobs/%(job_id)s/tasks/%(task_id)s/attempts/%(attempt_id)s' % {'job_id': job_id, 'task_id': task_id, 'attempt_id': attempt_id}, headers={'Accept': _JSON_CONTENT_TYPE})
+    return self._root.get('mapreduce/jobs/%(job_id)s/tasks/%(task_id)s/attempts/%(attempt_id)s' % {'job_id': job_id, 'task_id': task_id, 'attempt_id': attempt_id}, params=self._get_params(), headers={'Accept': _JSON_CONTENT_TYPE})
 
   def task_attempt_counters(self, job_id, task_id, attempt_id):
-    return self._root.get('mapreduce/jobs/%(job_id)s/tasks/%(task_id)s/attempts/%(attempt_id)s/counters' % {'job_id': job_id, 'task_id': task_id, 'attempt_id': attempt_id}, headers={'Accept': _JSON_CONTENT_TYPE})
+    return self._root.get('mapreduce/jobs/%(job_id)s/tasks/%(task_id)s/attempts/%(attempt_id)s/counters' % {'job_id': job_id, 'task_id': task_id, 'attempt_id': attempt_id}, params=self._get_params(), headers={'Accept': _JSON_CONTENT_TYPE})

+ 49 - 20
desktop/libs/hadoop/src/hadoop/yarn/mapreduce_api.py

@@ -19,6 +19,8 @@ import logging
 import posixpath
 import threading
 
+from desktop.conf import DEFAULT_USER
+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
@@ -30,22 +32,26 @@ LOG = logging.getLogger(__name__)
 _API_VERSION = 'v1'
 _JSON_CONTENT_TYPE = 'application/json'
 
-_api_cache = None
-_api_cache_lock = threading.Lock()
+API_CACHE = None
+API_CACHE_LOCK = threading.Lock()
 
 
-def get_mapreduce_api(user):
-  global _api_cache
-  if _api_cache is None:
-    _api_cache_lock.acquire()
+def get_mapreduce_api(username):
+  global API_CACHE
+  if API_CACHE is None:
+    API_CACHE_LOCK.acquire()
     try:
-      if _api_cache is None:
+      if API_CACHE is None:
         yarn_cluster = cluster.get_cluster_conf_for_job_submission()
-        if yarn_cluster is not None:
-          _api_cache = MapreduceApi(user, yarn_cluster.PROXY_API_URL.get(), yarn_cluster.SECURITY_ENABLED.get(), yarn_cluster.SSL_CERT_CA_VERIFY.get())
+        if yarn_cluster is None:
+          raise PopupException(_('No Resource Manager are available.'))
+        API_CACHE = MapreduceApi(username, yarn_cluster.PROXY_API_URL.get(), yarn_cluster.SECURITY_ENABLED.get(), yarn_cluster.SSL_CERT_CA_VERIFY.get())
     finally:
-      _api_cache_lock.release()
-  return _api_cache
+      API_CACHE_LOCK.release()
+
+  API_CACHE.setuser(username)  # Set the correct user
+
+  return API_CACHE
 
 
 class MapreduceApi(object):
@@ -56,6 +62,7 @@ class MapreduceApi(object):
     self._client = HttpClient(self._url, logger=LOG)
     self._root = Resource(self._client)
     self._security_enabled = security_enabled
+    self._thread_local = threading.local()  # To store user info
 
     if self._security_enabled:
       self._client.set_kerberos_auth()
@@ -65,17 +72,39 @@ class MapreduceApi(object):
   def __str__(self):
     return "MapreduceApi at %s" % (self._url,)
 
+  def _get_params(self):
+    params = {}
+
+    if self.username != DEFAULT_USER.get():  # We impersonate if needed
+      params['doAs'] = self.username
+      if not self._security_enabled:
+        params['user.name'] = DEFAULT_USER.get()
+
+    return params
+
   @property
   def url(self):
     return self._url
 
+  @property
+  def username(self):
+    try:
+      return self._thread_local.user
+    except AttributeError:
+      return DEFAULT_USER.get()
+
+  def setuser(self, user):
+    curr = self._user
+    self._thread_local.user = user
+    return curr
+
   def job(self, user, job_id):
     app_id = job_id.replace('job', 'application')
-    return self._root.get('%(app_id)s/ws/%(version)s/mapreduce/jobs/%(job_id)s' % {'app_id': app_id, 'job_id': job_id, 'version': _API_VERSION}, headers={'Accept': _JSON_CONTENT_TYPE})
+    return self._root.get('%(app_id)s/ws/%(version)s/mapreduce/jobs/%(job_id)s' % {'app_id': app_id, 'job_id': job_id, 'version': _API_VERSION}, params=self._get_params(), headers={'Accept': _JSON_CONTENT_TYPE})
 
   def counters(self, job_id):
     app_id = job_id.replace('job', 'application')
-    response = self._root.get('%(app_id)s/ws/%(version)s/mapreduce/jobs/%(job_id)s/counters' % {'app_id': app_id, 'job_id': job_id, 'version': _API_VERSION}, headers={'Accept': _JSON_CONTENT_TYPE})
+    response = self._root.get('%(app_id)s/ws/%(version)s/mapreduce/jobs/%(job_id)s/counters' % {'app_id': app_id, 'job_id': job_id, 'version': _API_VERSION}, params=self._get_params(), headers={'Accept': _JSON_CONTENT_TYPE})
     # If it hits the job history server, it will return HTML.
     # Simply return None in this case because there isn't much data there.
     if isinstance(response, basestring):
@@ -85,33 +114,33 @@ class MapreduceApi(object):
 
   def tasks(self, job_id):
     app_id = job_id.replace('job', 'application')
-    return self._root.get('%(app_id)s/ws/%(version)s/mapreduce/jobs/%(job_id)s/tasks' % {'app_id': app_id, 'job_id': job_id, 'version': _API_VERSION}, headers={'Accept': _JSON_CONTENT_TYPE})
+    return self._root.get('%(app_id)s/ws/%(version)s/mapreduce/jobs/%(job_id)s/tasks' % {'app_id': app_id, 'job_id': job_id, 'version': _API_VERSION}, params=self._get_params(), headers={'Accept': _JSON_CONTENT_TYPE})
 
   def job_attempts(self, job_id):
     app_id = job_id.replace('job', 'application')
-    return self._root.get('%(app_id)s/ws/%(version)s/mapreduce/jobs/%(job_id)s/jobattempts' % {'app_id': app_id, 'job_id': job_id, 'version': _API_VERSION}, headers={'Accept': _JSON_CONTENT_TYPE})
+    return self._root.get('%(app_id)s/ws/%(version)s/mapreduce/jobs/%(job_id)s/jobattempts' % {'app_id': app_id, 'job_id': job_id, 'version': _API_VERSION}, params=self._get_params(), headers={'Accept': _JSON_CONTENT_TYPE})
 
   def conf(self, job_id):
     app_id = job_id.replace('job', 'application')
-    return self._root.get('%(app_id)s/ws/%(version)s/mapreduce/jobs/%(job_id)s/conf' % {'app_id': app_id, 'job_id': job_id, 'version': _API_VERSION}, headers={'Accept': _JSON_CONTENT_TYPE})
+    return self._root.get('%(app_id)s/ws/%(version)s/mapreduce/jobs/%(job_id)s/conf' % {'app_id': app_id, 'job_id': job_id, 'version': _API_VERSION}, params=self._get_params(), headers={'Accept': _JSON_CONTENT_TYPE})
 
   def task(self, job_id, task_id):
     app_id = job_id.replace('job', 'application')
-    return self._root.get('%(app_id)s/ws/%(version)s/mapreduce/jobs/%(job_id)s/tasks/%(task_id)s' % {'app_id': app_id, 'job_id': job_id, 'task_id': task_id, 'version': _API_VERSION}, headers={'Accept': _JSON_CONTENT_TYPE})
+    return self._root.get('%(app_id)s/ws/%(version)s/mapreduce/jobs/%(job_id)s/tasks/%(task_id)s' % {'app_id': app_id, 'job_id': job_id, 'task_id': task_id, 'version': _API_VERSION}, params=self._get_params(), headers={'Accept': _JSON_CONTENT_TYPE})
 
   def task_counters(self, job_id, task_id):
     app_id = job_id.replace('job', 'application')
     job_id = job_id.replace('application', 'job')
-    return self._root.get('%(app_id)s/ws/%(version)s/mapreduce/jobs/%(job_id)s/tasks/%(task_id)s/counters' % {'app_id': app_id, 'job_id': job_id, 'task_id': task_id, 'version': _API_VERSION}, headers={'Accept': _JSON_CONTENT_TYPE})
+    return self._root.get('%(app_id)s/ws/%(version)s/mapreduce/jobs/%(job_id)s/tasks/%(task_id)s/counters' % {'app_id': app_id, 'job_id': job_id, 'task_id': task_id, 'version': _API_VERSION}, params=self._get_params(), headers={'Accept': _JSON_CONTENT_TYPE})
 
   def task_attempts(self, job_id, task_id):
     app_id = job_id.replace('job', 'application')
-    return self._root.get('%(app_id)s/ws/%(version)s/mapreduce/jobs/%(job_id)s/tasks/%(task_id)s/attempts' % {'app_id': app_id, 'job_id': job_id, 'task_id': task_id, 'version': _API_VERSION}, headers={'Accept': _JSON_CONTENT_TYPE})
+    return self._root.get('%(app_id)s/ws/%(version)s/mapreduce/jobs/%(job_id)s/tasks/%(task_id)s/attempts' % {'app_id': app_id, 'job_id': job_id, 'task_id': task_id, 'version': _API_VERSION}, params=self._get_params(), headers={'Accept': _JSON_CONTENT_TYPE})
 
   def task_attempt(self, job_id, task_id, attempt_id):
     app_id = job_id.replace('job', 'application')
     job_id = job_id.replace('application', 'job')
-    return self._root.get('%(app_id)s/ws/%(version)s/mapreduce/jobs/%(job_id)s/tasks/%(task_id)s/attempts/%(attempt_id)s' % {'app_id': app_id, 'job_id': job_id, 'task_id': task_id, 'attempt_id': attempt_id, 'version': _API_VERSION}, headers={'Accept': _JSON_CONTENT_TYPE})
+    return self._root.get('%(app_id)s/ws/%(version)s/mapreduce/jobs/%(job_id)s/tasks/%(task_id)s/attempts/%(attempt_id)s' % {'app_id': app_id, 'job_id': job_id, 'task_id': task_id, 'attempt_id': attempt_id, 'version': _API_VERSION}, params=self._get_params(), headers={'Accept': _JSON_CONTENT_TYPE})
 
   def kill(self, job_id):
     app_id = job_id.replace('job', 'application')