Przeglądaj źródła

[oozie] Set implicitly the user used in the API

Dashboard doAs username equalled to the default user previously
Romain Rigaux 12 lat temu
rodzic
commit
17d1fd1

+ 1 - 1
apps/oozie/src/oozie/conf.py

@@ -58,7 +58,7 @@ OOZIE_JOBS_COUNT = Config(
 def config_validator(user):
   res = []
 
-  status = get_oozie_status()
+  status = get_oozie_status(user)
 
   if 'NORMAL' not in status:
     res.append((NICE_NAME, _("The app won't work without a running Oozie server")))

+ 1 - 1
apps/oozie/src/oozie/templates/editor/action_utils.mako

@@ -110,7 +110,7 @@
     </p>
   % endif
   % if node_type == 'java':
-    % if get_oozie().security_enabled:
+    % if get_oozie(user).security_enabled:
     <br style="clear: both" />
       <p class="alert alert-warn span5">
         ${ _('The delegation token needs to be propagated from the launcher job to the MR job') }.

+ 28 - 27
apps/oozie/src/oozie/views/dashboard.py

@@ -67,7 +67,7 @@ def manage_oozie_jobs(request, job_id, action):
   response = {'status': -1, 'data': ''}
 
   try:
-    response['data'] = get_oozie().job_control(job_id, action)
+    response['data'] = get_oozie(request.user).job_control(job_id, action)
     response['status'] = 0
     if 'notification' in request.POST:
       request.info(_(request.POST.get('notification')))
@@ -95,19 +95,19 @@ def list_oozie_workflows(request):
   if not has_dashboard_jobs_access(request.user):
     kwargs['user'] = request.user.username
 
-  workflows = get_oozie().get_workflows(**kwargs)
+  workflows = get_oozie(request.user).get_workflows(**kwargs)
 
   if request.GET.get('format') == 'json':
     json_jobs = workflows.jobs
     if request.GET.get('type') == 'running':
-      json_jobs = split_oozie_jobs(workflows.jobs)['running_jobs']
+      json_jobs = split_oozie_jobs(request.user, workflows.jobs)['running_jobs']
     if request.GET.get('type') == 'completed':
-      json_jobs = split_oozie_jobs(workflows.jobs)['completed_jobs']
+      json_jobs = split_oozie_jobs(request.user, workflows.jobs)['completed_jobs']
     return HttpResponse(encode_json_for_js(massaged_oozie_jobs_for_json(json_jobs, request.user)), mimetype="application/json")
 
   return render('dashboard/list_oozie_workflows.mako', request, {
     'user': request.user,
-    'jobs': split_oozie_jobs(workflows.jobs),
+    'jobs': split_oozie_jobs(request.user, workflows.jobs),
     'has_job_edition_permission':  has_job_edition_permission,
   })
 
@@ -118,18 +118,18 @@ def list_oozie_coordinators(request):
   if not has_dashboard_jobs_access(request.user):
     kwargs['user'] = request.user.username
 
-  coordinators = get_oozie().get_coordinators(**kwargs)
+  coordinators = get_oozie(request.user).get_coordinators(**kwargs)
 
   if request.GET.get('format') == 'json':
     json_jobs = coordinators.jobs
     if request.GET.get('type') == 'running':
-      json_jobs = split_oozie_jobs(coordinators.jobs)['running_jobs']
+      json_jobs = split_oozie_jobs(request.user, coordinators.jobs)['running_jobs']
     if request.GET.get('type') == 'completed':
-      json_jobs = split_oozie_jobs(coordinators.jobs)['completed_jobs']
+      json_jobs = split_oozie_jobs(request.user, coordinators.jobs)['completed_jobs']
     return HttpResponse(json.dumps(massaged_oozie_jobs_for_json(json_jobs, request.user)).replace('\\\\', '\\'), mimetype="application/json")
 
   return render('dashboard/list_oozie_coordinators.mako', request, {
-    'jobs': split_oozie_jobs(coordinators.jobs),
+    'jobs': split_oozie_jobs(request.user, coordinators.jobs),
     'has_job_edition_permission': has_job_edition_permission,
   })
 
@@ -140,18 +140,18 @@ def list_oozie_bundles(request):
   if not has_dashboard_jobs_access(request.user):
     kwargs['user'] = request.user.username
 
-  bundles = get_oozie().get_bundles(**kwargs)
+  bundles = get_oozie(request.user).get_bundles(**kwargs)
 
   if request.GET.get('format') == 'json':
     json_jobs = bundles.jobs
     if request.GET.get('type') == 'running':
-      json_jobs = split_oozie_jobs(bundles.jobs)['running_jobs']
+      json_jobs = split_oozie_jobs(request.user, bundles.jobs)['running_jobs']
     if request.GET.get('type') == 'completed':
-      json_jobs = split_oozie_jobs(bundles.jobs)['completed_jobs']
+      json_jobs = split_oozie_jobs(request.user, bundles.jobs)['completed_jobs']
     return HttpResponse(json.dumps(massaged_oozie_jobs_for_json(json_jobs, request.user)).replace('\\\\', '\\'), mimetype="application/json")
 
   return render('dashboard/list_oozie_bundles.mako', request, {
-    'jobs': split_oozie_jobs(bundles.jobs),
+    'jobs': split_oozie_jobs(request.user, bundles.jobs),
     'has_job_edition_permission': has_job_edition_permission,
   })
 
@@ -283,7 +283,7 @@ def list_oozie_bundle(request, job_id):
 @show_oozie_error
 def list_oozie_workflow_action(request, action, coordinator_job_id=None, bundle_job_id=None):
   try:
-    action = get_oozie().get_action(action)
+    action = get_oozie(request.user).get_action(action)
     workflow = check_job_access_permission(request, action.id.split('@')[0])
   except RestException, ex:
     raise PopupException(_("Error accessing Oozie action %s.") % (action,),
@@ -310,10 +310,11 @@ def list_oozie_workflow_action(request, action, coordinator_job_id=None, bundle_
 
 @show_oozie_error
 def list_oozie_info(request):
+  api = get_oozie(request.user)
 
-  instrumentation = get_oozie().get_instrumentation()
-  configuration = get_oozie().get_configuration()
-  oozie_status = get_oozie().get_oozie_status()
+  instrumentation = api.get_instrumentation()
+  configuration = api.get_configuration()
+  oozie_status = api.get_oozie_status()
 
   return render('dashboard/list_oozie_info.mako', request, {
     'instrumentation': instrumentation,
@@ -598,11 +599,11 @@ def massaged_oozie_jobs_for_json(oozie_jobs, user):
   for job in oozie_jobs:
     if job.is_running():
       if job.type == 'Workflow':
-        job = get_oozie().get_job(job.id)
+        job = get_oozie(user).get_job(job.id)
       elif job.type == 'Coordinator':
-        job = get_oozie().get_coordinator(job.id)
+        job = get_oozie(user).get_coordinator(job.id)
       else:
-        job = get_oozie().get_bundle(job.id)
+        job = get_oozie(user).get_bundle(job.id)
 
     massaged_job = {
       'id': job.id,
@@ -632,7 +633,7 @@ def massaged_oozie_jobs_for_json(oozie_jobs, user):
   return jobs
 
 
-def split_oozie_jobs(oozie_jobs):
+def split_oozie_jobs(user, oozie_jobs):
   jobs = {}
   jobs_running = []
   jobs_completed = []
@@ -641,11 +642,11 @@ def split_oozie_jobs(oozie_jobs):
     if job.appName != 'pig-app-hue-script':
       if job.is_running():
         if job.type == 'Workflow':
-          job = get_oozie().get_job(job.id)
+          job = get_oozie(user).get_job(job.id)
         elif job.type == 'Coordinator':
-          job = get_oozie().get_coordinator(job.id)
+          job = get_oozie(user).get_coordinator(job.id)
         else:
-          job = get_oozie().get_bundle(job.id)
+          job = get_oozie(user).get_bundle(job.id)
         jobs_running.append(job)
       else:
         jobs_completed.append(job)
@@ -667,11 +668,11 @@ def check_job_access_permission(request, job_id):
   """
   if job_id is not None:
     if job_id.endswith('W'):
-      get_job = get_oozie().get_job
+      get_job = get_oozie(request.user).get_job
     elif job_id.endswith('C'):
-      get_job = get_oozie().get_coordinator
+      get_job = get_oozie(request.user).get_coordinator
     else:
-      get_job = get_oozie().get_bundle
+      get_job = get_oozie(request.user).get_bundle
 
     try:
       oozie_job = get_job(job_id)

+ 3 - 3
apps/pig/src/pig/api.py

@@ -132,14 +132,14 @@ class OozieApi:
     return pig_params
 
   def stop(self, job_id):
-    return get_oozie().job_control(job_id, 'kill')
+    return get_oozie(self.user).job_control(job_id, 'kill')
 
   def get_jobs(self):
     kwargs = {'cnt': OozieApi.MAX_DASHBOARD_JOBS,}
     kwargs['user'] = self.user.username
     kwargs['name'] = OozieApi.WORKFLOW_NAME
 
-    return get_oozie().get_workflows(**kwargs).jobs
+    return get_oozie(self.user).get_workflows(**kwargs).jobs
 
   def get_log(self, request, oozie_workflow):
     logs = {}
@@ -208,7 +208,7 @@ class OozieApi:
 
     for job in oozie_jobs:
       if job.is_running():
-        job = get_oozie().get_job(job.id)
+        job = get_oozie(self.user).get_job(job.id)
         get_copy = request.GET.copy() # Hacky, would need to refactor JobBrowser get logs
         get_copy['format'] = 'python'
         request.GET = get_copy

+ 3 - 3
apps/spark/src/spark/api.py

@@ -135,14 +135,14 @@ cd %(spark_home)s
     return spark_params
 
   def stop(self, job_id):
-    return get_oozie().job_control(job_id, 'kill')
+    return get_oozie(self.user).job_control(job_id, 'kill')
 
   def get_jobs(self):
     kwargs = {'cnt': OozieSparkApi.MAX_DASHBOARD_JOBS,}
     kwargs['user'] = self.user.username
     kwargs['name'] = OozieSparkApi.WORKFLOW_NAME
 
-    return get_oozie().get_workflows(**kwargs).jobs
+    return get_oozie(self.user).get_workflows(**kwargs).jobs
 
   def get_log(self, request, oozie_workflow):
     logs = {}
@@ -211,7 +211,7 @@ cd %(spark_home)s
 
     for job in oozie_jobs:
       if job.is_running():
-        job = get_oozie().get_job(job.id)
+        job = get_oozie(self.user).get_job(job.id)
         get_copy = request.GET.copy() # Hacky, would need to refactor JobBrowser get logs
         get_copy['format'] = 'python'
         request.GET = get_copy

+ 4 - 3
desktop/libs/liboozie/src/liboozie/conf.py

@@ -39,14 +39,15 @@ REMOTE_DEPLOYMENT_DIR = Config(
   default="/user/hue/oozie/deployments",
   help=_t("Location on HDFS where the workflows/coordinators are deployed when submitted by a non-owner."))
 
-def get_oozie_status():
+
+def get_oozie_status(user):
   from liboozie.oozie_api import get_oozie
 
   status = 'down'
 
   try:
     if not 'test' in sys.argv: # Avoid tests hanging
-      status = str(get_oozie().get_oozie_status())
+      status = str(get_oozie(user).get_oozie_status())
   except:
     pass
 
@@ -62,7 +63,7 @@ def config_validator(user):
 
   res = []
 
-  status = get_oozie_status()
+  status = get_oozie_status(user)
   if 'NORMAL' not in status:
     res.append((status, _('The Oozie server is not available')))
 

+ 7 - 16
desktop/libs/liboozie/src/liboozie/oozie_api.py

@@ -40,8 +40,7 @@ _api_cache = None
 _api_cache_lock = threading.Lock()
 
 
-def get_oozie():
-  """Return a cached OozieApi"""
+def get_oozie(user):
   global _api_cache
   if _api_cache is None:
     _api_cache_lock.acquire()
@@ -51,6 +50,7 @@ def get_oozie():
         _api_cache = OozieApi(OOZIE_URL.get(), secure)
     finally:
       _api_cache_lock.release()
+  _api_cache.setuser(user)
   return _api_cache
 
 
@@ -62,9 +62,7 @@ class OozieApi(object):
       self._client.set_kerberos_auth()
     self._root = Resource(self._client)
     self._security_enabled = security_enabled
-
-    # To store user info
-    self._thread_local = threading.local()
+    self.user = None # username actually
 
   def __str__(self):
     return "OozieApi at %s" % (self._url,)
@@ -77,18 +75,11 @@ class OozieApi(object):
   def security_enabled(self):
     return self._security_enabled
 
-  @property
-  def user(self):
-    try:
-      return self._thread_local.user
-    except AttributeError:
-      return DEFAULT_USER
-
   def setuser(self, user):
-    """Return the previous user"""
-    prev = self.user
-    self._thread_local.user = user
-    return prev
+    if hasattr(user, 'username'):
+      self.user = user.username
+    else:
+      self.user = user
 
   def _get_params(self):
     if self.security_enabled:

+ 7 - 4
desktop/libs/liboozie/src/liboozie/oozie_api_test.py

@@ -190,7 +190,7 @@ class OozieServerProvider(object):
         status = None
         try:
           LOG.info('Check Oozie status...')
-          status = get_oozie().get_oozie_status()
+          status = get_oozie(cluster.superuser).get_oozie_status()
           if status['systemMode'] == 'NORMAL':
             started = True
             break
@@ -214,12 +214,15 @@ class OozieServerProvider(object):
 
     _oozie_lock.release()
 
-    return get_oozie(), callback
+    cluster = pseudo_hdfs4.shared_cluster()
+    return get_oozie(cluster.superuser), callback
 
 
 class TestMiniOozie(OozieServerProvider):
 
   def test_oozie_status(self):
-    assert_equal(get_oozie().get_oozie_status()['systemMode'], 'NORMAL')
+    user = getpass.getuser()
 
-    assert_true(self.cluster.fs.exists('/user/%(user)s/share/lib' % {'user': getpass.getuser()}))
+    assert_equal(get_oozie(user).get_oozie_status()['systemMode'], 'NORMAL')
+
+    assert_true(self.cluster.fs.exists('/user/%(user)s/share/lib' % {'user': user}))

+ 26 - 41
desktop/libs/liboozie/src/liboozie/submittion.py

@@ -45,6 +45,7 @@ class Submission(object):
     self.fs = fs
     self.jt = jt
     self.oozie_id = oozie_id
+    self.api = get_oozie(self.user)
 
     if properties is not None:
       self.properties = properties
@@ -72,40 +73,32 @@ class Submission(object):
 
     deployment_dir = self.deploy()
 
-    try:
-      prev = get_oozie().setuser(self.user.username)
-      self._update_properties(jobtracker, deployment_dir)
-      self.oozie_id = get_oozie().submit_job(self.properties)
-      LOG.info("Submitted: %s" % (self,))
-
-      if self.job.get_type() == 'workflow':
-        get_oozie().job_control(self.oozie_id, 'start')
-        LOG.info("Started: %s" % (self,))
-    finally:
-      get_oozie().setuser(prev)
+    self._update_properties(jobtracker, deployment_dir)
+    self.oozie_id = self.api.submit_job(self.properties)
+    LOG.info("Submitted: %s" % (self,))
+
+    if self.job.get_type() == 'workflow':
+      self.api.job_control(self.oozie_id, 'start')
+      LOG.info("Started: %s" % (self,))
 
     return self.oozie_id
 
   def rerun(self, deployment_dir, fail_nodes=None, skip_nodes=None):
     jobtracker = cluster.get_cluster_addr_for_job_submission()
 
-    try:
-      prev = get_oozie().setuser(self.user.username)
-      self._update_properties(jobtracker, deployment_dir)
-      self.properties.update({'oozie.wf.application.path': deployment_dir})
+    self._update_properties(jobtracker, deployment_dir)
+    self.properties.update({'oozie.wf.application.path': deployment_dir})
 
-      if fail_nodes:
-        self.properties.update({'oozie.wf.rerun.failnodes': fail_nodes})
-      elif not skip_nodes:
-        self.properties.update({'oozie.wf.rerun.failnodes': 'false'}) # Case empty 'skip_nodes' list
-      else:
-        self.properties.update({'oozie.wf.rerun.skip.nodes': skip_nodes})
+    if fail_nodes:
+      self.properties.update({'oozie.wf.rerun.failnodes': fail_nodes})
+    elif not skip_nodes:
+      self.properties.update({'oozie.wf.rerun.failnodes': 'false'}) # Case empty 'skip_nodes' list
+    else:
+      self.properties.update({'oozie.wf.rerun.skip.nodes': skip_nodes})
 
-      get_oozie().rerun(self.oozie_id, properties=self.properties)
+    self.api.rerun(self.oozie_id, properties=self.properties)
 
-      LOG.info("Rerun: %s" % (self,))
-    finally:
-      get_oozie().setuser(prev)
+    LOG.info("Rerun: %s" % (self,))
 
     return self.oozie_id
 
@@ -113,15 +106,11 @@ class Submission(object):
   def rerun_coord(self, deployment_dir, params):
     jobtracker = cluster.get_cluster_addr_for_job_submission()
 
-    try:
-      prev = get_oozie().setuser(self.user.username)
-      self._update_properties(jobtracker, deployment_dir)
-      self.properties.update({'oozie.coord.application.path': deployment_dir})
+    self._update_properties(jobtracker, deployment_dir)
+    self.properties.update({'oozie.coord.application.path': deployment_dir})
 
-      get_oozie().job_control(self.oozie_id, action='coord-rerun', properties=self.properties, parameters=params)
-      LOG.info("Rerun: %s" % (self,))
-    finally:
-      get_oozie().setuser(prev)
+    self.api.job_control(self.oozie_id, action='coord-rerun', properties=self.properties, parameters=params)
+    LOG.info("Rerun: %s" % (self,))
 
     return self.oozie_id
 
@@ -129,14 +118,10 @@ class Submission(object):
   def rerun_bundle(self, deployment_dir, params):
     jobtracker = cluster.get_cluster_addr_for_job_submission()
 
-    try:
-      prev = get_oozie().setuser(self.user.username)
-      self._update_properties(jobtracker, deployment_dir)
-      self.properties.update({'oozie.bundle.application.path': deployment_dir})
-      get_oozie().job_control(self.oozie_id, action='bundle-rerun', properties=self.properties, parameters=params)
-      LOG.info("Rerun: %s" % (self,))
-    finally:
-      get_oozie().setuser(prev)
+    self._update_properties(jobtracker, deployment_dir)
+    self.properties.update({'oozie.bundle.application.path': deployment_dir})
+    self.api.job_control(self.oozie_id, action='bundle-rerun', properties=self.properties, parameters=params)
+    LOG.info("Rerun: %s" % (self,))
 
     return self.oozie_id