api.py 30 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726727728729730731732733734735736737738739740741742743744745746747748749750751752753754755756757758759760761762763764765766767768769770771772773774775776777778779780781782783784785786787788789790791792793794795796797798799800801802803804805806807808809810811812813814815816817818819820821822823824825826827828829830831832833834835836837838839840841842843844845846847848849850851852853854855856857858859860861862863864865866867868869870871872873874875876877878879880881882883884885886887888889890891892893894895896897898899900901902903904905906907
  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. import logging
  18. import json
  19. import re
  20. from django.contrib.auth.models import User
  21. from django.core.urlresolvers import reverse
  22. from django.http import Http404
  23. from django.utils.translation import ugettext as _
  24. from django.views.decorators.http import require_POST
  25. from thrift.transport.TTransport import TTransportException
  26. from desktop.context_processors import get_app_name
  27. from desktop.lib.django_util import JsonResponse
  28. from desktop.lib.exceptions import StructuredThriftTransportException
  29. from desktop.lib.exceptions_renderable import PopupException
  30. from desktop.lib.i18n import force_unicode
  31. from desktop.lib.parameterization import substitute_variables
  32. from metastore import parser
  33. from notebook.models import escape_rows
  34. import beeswax.models
  35. from beeswax.data_export import upload
  36. from beeswax.design import HQLdesign
  37. from beeswax.conf import USE_GET_LOG_API
  38. from beeswax.forms import QueryForm
  39. from beeswax.models import Session, QueryHistory
  40. from beeswax.server import dbms
  41. from beeswax.server.dbms import expand_exception, get_query_server_config, QueryServerException, QueryServerTimeoutException
  42. from beeswax.views import authorized_get_design, authorized_get_query_history, make_parameterization_form,\
  43. safe_get_design, save_design, massage_columns_for_json, _get_query_handle_and_state, \
  44. _parse_out_hadoop_jobs
  45. LOG = logging.getLogger(__name__)
  46. def error_handler(view_fn):
  47. def decorator(request, *args, **kwargs):
  48. try:
  49. return view_fn(request, *args, **kwargs)
  50. except Http404, e:
  51. raise e
  52. except Exception, e:
  53. LOG.exception('error in %s' % view_fn)
  54. if not hasattr(e, 'message') or not e.message:
  55. message = str(e)
  56. else:
  57. message = force_unicode(e.message, strings_only=True, errors='replace')
  58. if 'Invalid OperationHandle' in message and 'id' in kwargs:
  59. # Expired state.
  60. query_history = authorized_get_query_history(request, kwargs['id'], must_exist=False)
  61. if query_history:
  62. query_history.set_to_expired()
  63. query_history.save()
  64. response = {
  65. 'status': -1,
  66. 'message': message,
  67. }
  68. if re.search('database is locked|Invalid query handle|not JSON serializable', message, re.IGNORECASE):
  69. response['status'] = 2 # Frontend will not display this type of error
  70. LOG.warn('error_handler silencing the exception: %s' % e)
  71. return JsonResponse(response)
  72. return decorator
  73. @error_handler
  74. def autocomplete(request, database=None, table=None, column=None, nested=None):
  75. app_name = get_app_name(request)
  76. query_server = get_query_server_config(app_name)
  77. do_as = request.user
  78. if (request.user.is_superuser or request.user.has_hue_permission(action="impersonate", app="security")) and 'doas' in request.GET:
  79. do_as = User.objects.get(username=request.GET.get('doas'))
  80. db = dbms.get(do_as, query_server)
  81. response = _autocomplete(db, database, table, column, nested)
  82. return JsonResponse(response)
  83. def _autocomplete(db, database=None, table=None, column=None, nested=None):
  84. response = {}
  85. try:
  86. if database is None:
  87. response['databases'] = db.get_databases()
  88. elif table is None:
  89. tables_meta = db.get_tables_meta(database=database)
  90. response['tables_meta'] = tables_meta
  91. elif column is None:
  92. t = db.get_table(database, table)
  93. response['hdfs_link'] = t.hdfs_link
  94. response['columns'] = [column.name for column in t.cols]
  95. response['extended_columns'] = massage_columns_for_json(t.cols)
  96. response['partition_keys'] = [{'name': part.name, 'type': part.type} for part in t.partition_keys]
  97. else:
  98. col = db.get_column(database, table, column)
  99. if col:
  100. parse_tree = parser.parse_column(col.name, col.type, col.comment)
  101. if nested:
  102. parse_tree = _extract_nested_type(parse_tree, nested)
  103. response = parse_tree
  104. # If column or nested type is scalar/primitive, add sample of values
  105. if parser.is_scalar_type(parse_tree['type']):
  106. table_obj = db.get_table(database, table)
  107. sample = db.get_sample(database, table_obj, column, nested)
  108. if sample:
  109. sample = set([row[0] for row in sample.rows()])
  110. response['sample'] = sorted(list(sample))
  111. else:
  112. raise Exception('Could not find column `%s`.`%s`.`%s`' % (database, table, column))
  113. except (QueryServerTimeoutException, TTransportException), e:
  114. response['code'] = 503
  115. response['error'] = e.message
  116. except Exception, e:
  117. LOG.warn('Autocomplete data fetching error: %s' % e)
  118. response['code'] = 500
  119. response['error'] = e.message
  120. return response
  121. @error_handler
  122. def parameters(request, design_id=None):
  123. response = {'status': -1, 'message': ''}
  124. # Use POST request to not confine query length.
  125. if request.method != 'POST':
  126. response['message'] = _('A POST request is required.')
  127. parameterization_form_cls = make_parameterization_form(request.POST.get('query-query', ''))
  128. if parameterization_form_cls:
  129. parameterization_form = parameterization_form_cls(prefix="parameterization")
  130. response['parameters'] = [{'parameter': field.html_name, 'name': field.name} for field in parameterization_form]
  131. response['status']= 0
  132. else:
  133. response['parameters'] = []
  134. response['status']= 0
  135. return JsonResponse(response)
  136. @error_handler
  137. def execute_directly(request, query, design, query_server, tablename=None, **kwargs):
  138. if design is not None:
  139. design = authorized_get_design(request, design.id)
  140. parameters = kwargs.pop('parameters', None)
  141. db = dbms.get(request.user, query_server)
  142. database = query.query.get('database', 'default')
  143. db.use(database)
  144. history_obj = db.execute_query(query, design)
  145. watch_url = reverse(get_app_name(request) + ':api_watch_query_refresh_json', kwargs={'id': history_obj.id})
  146. if parameters is not None:
  147. history_obj.update_extra('parameters', parameters)
  148. history_obj.save()
  149. response = {
  150. 'status': 0,
  151. 'id': history_obj.id,
  152. 'watch_url': watch_url,
  153. 'statement': history_obj.get_current_statement(),
  154. 'is_redacted': history_obj.is_redacted
  155. }
  156. return JsonResponse(response)
  157. @error_handler
  158. def watch_query_refresh_json(request, id):
  159. query_history = authorized_get_query_history(request, id, must_exist=True)
  160. db = dbms.get(request.user, query_history.get_query_server_config())
  161. if not request.POST.get('next'): # We need this as multi query would fail as current query is closed
  162. handle, state = _get_query_handle_and_state(query_history)
  163. query_history.save_state(state)
  164. # Go to next statement if asked to continue or when a statement with no dataset finished.
  165. try:
  166. if request.POST.get('next') or (not query_history.is_finished() and query_history.is_success() and not query_history.has_results):
  167. close_operation(request, id)
  168. query_history = db.execute_next_statement(query_history, request.POST.get('query-query'))
  169. handle, state = _get_query_handle_and_state(query_history)
  170. except QueryServerException, ex:
  171. raise ex
  172. except Exception, ex:
  173. LOG.exception(ex)
  174. handle, state = _get_query_handle_and_state(query_history)
  175. try:
  176. start_over = request.POST.get('log-start-over') == 'true'
  177. log = db.get_log(handle, start_over=start_over)
  178. except Exception, ex:
  179. log = str(ex)
  180. jobs = _parse_out_hadoop_jobs(log)
  181. job_urls = massage_job_urls_for_json(jobs)
  182. result = {
  183. 'status': -1,
  184. 'log': log,
  185. 'jobs': jobs,
  186. 'jobUrls': job_urls,
  187. 'isSuccess': query_history.is_success(),
  188. 'isFailure': query_history.is_failure(),
  189. 'id': id,
  190. 'statement': query_history.get_current_statement(),
  191. 'watch_url': reverse(get_app_name(request) + ':api_watch_query_refresh_json', kwargs={'id': query_history.id}),
  192. 'oldLogsApi': USE_GET_LOG_API.get()
  193. }
  194. # Run time error
  195. if query_history.is_failure():
  196. res = db.get_operation_status(handle)
  197. if query_history.is_canceled(res):
  198. result['status'] = 0
  199. elif hasattr(res, 'errorMessage') and res.errorMessage:
  200. result['message'] = res.errorMessage
  201. else:
  202. result['message'] = _('Bad status for request %s:\n%s') % (id, res)
  203. else:
  204. result['status'] = 0
  205. return JsonResponse(result)
  206. def massage_job_urls_for_json(jobs):
  207. massaged_jobs = []
  208. for job in jobs:
  209. massaged_jobs.append({
  210. 'name': job,
  211. 'url': reverse('jobbrowser.views.single_job', kwargs={'job': job})
  212. })
  213. return massaged_jobs
  214. @error_handler
  215. def close_operation(request, query_history_id):
  216. response = {
  217. 'status': -1,
  218. 'message': ''
  219. }
  220. if request.method != 'POST':
  221. response['message'] = _('A POST request is required.')
  222. else:
  223. query_history = authorized_get_query_history(request, query_history_id, must_exist=True)
  224. db = dbms.get(query_history.owner, query_history.get_query_server_config())
  225. handle = query_history.get_handle()
  226. db.close_operation(handle)
  227. query_history.set_to_expired()
  228. query_history.save()
  229. response['status'] = 0
  230. return JsonResponse(response)
  231. @error_handler
  232. def explain_directly(request, query_server, query):
  233. explanation = dbms.get(request.user, query_server).explain(query)
  234. response = {
  235. 'status': 0,
  236. 'explanation': explanation.textual,
  237. 'statement': query.get_query_statement(0),
  238. }
  239. return JsonResponse(response)
  240. @error_handler
  241. def execute(request, design_id=None):
  242. response = {'status': -1, 'message': ''}
  243. if request.method != 'POST':
  244. response['message'] = _('A POST request is required.')
  245. app_name = get_app_name(request)
  246. query_server = get_query_server_config(app_name)
  247. query_type = beeswax.models.SavedQuery.TYPES_MAPPING[app_name]
  248. design = safe_get_design(request, query_type, design_id)
  249. try:
  250. query_form = get_query_form(request)
  251. if query_form.is_valid():
  252. query_str = query_form.query.cleaned_data["query"]
  253. explain = request.GET.get('explain', 'false').lower() == 'true'
  254. design = save_design(request, query_form, query_type, design, False)
  255. if query_form.query.cleaned_data['is_parameterized']:
  256. # Parameterized query
  257. parameterization_form_cls = make_parameterization_form(query_str)
  258. if parameterization_form_cls:
  259. parameterization_form = parameterization_form_cls(request.REQUEST, prefix="parameterization")
  260. if parameterization_form.is_valid():
  261. parameters = parameterization_form.cleaned_data
  262. real_query = substitute_variables(query_str, parameters)
  263. query = HQLdesign(query_form, query_type=query_type)
  264. query._data_dict['query']['query'] = real_query
  265. try:
  266. if explain:
  267. return explain_directly(request, query_server, query)
  268. else:
  269. return execute_directly(request, query, design, query_server, parameters=parameters)
  270. except Exception, ex:
  271. db = dbms.get(request.user, query_server)
  272. error_message, log = expand_exception(ex, db)
  273. response['message'] = error_message
  274. return JsonResponse(response)
  275. else:
  276. response['errors'] = parameterization_form.errors
  277. return JsonResponse(response)
  278. # Non-parameterized query
  279. query = HQLdesign(query_form, query_type=query_type)
  280. if request.GET.get('explain', 'false').lower() == 'true':
  281. return explain_directly(request, query_server, query)
  282. else:
  283. return execute_directly(request, query, design, query_server)
  284. else:
  285. response['message'] = _('There was an error with your query.')
  286. response['errors'] = {
  287. 'query': [query_form.query.errors],
  288. 'settings': query_form.settings.errors,
  289. 'file_resources': query_form.file_resources.errors,
  290. 'functions': query_form.functions.errors,
  291. }
  292. except RuntimeError, e:
  293. response['message']= str(e)
  294. return JsonResponse(response)
  295. @error_handler
  296. def save_query_design(request, design_id=None):
  297. response = {'status': -1, 'message': ''}
  298. if request.method != 'POST':
  299. response['message'] = _('A POST request is required.')
  300. app_name = get_app_name(request)
  301. query_type = beeswax.models.SavedQuery.TYPES_MAPPING[app_name]
  302. design = safe_get_design(request, query_type, design_id)
  303. try:
  304. query_form = get_query_form(request)
  305. if query_form.is_valid():
  306. design = save_design(request, query_form, query_type, design, True)
  307. response['design_id'] = design.id
  308. response['status'] = 0
  309. else:
  310. response['errors'] = {
  311. 'query': [query_form.query.errors],
  312. 'settings': query_form.settings.errors,
  313. 'file_resources': query_form.file_resources.errors,
  314. 'functions': query_form.functions.errors,
  315. 'saveform': query_form.saveform.errors,
  316. }
  317. except RuntimeError, e:
  318. response['message'] = str(e)
  319. return JsonResponse(response)
  320. @error_handler
  321. def fetch_saved_design(request, design_id):
  322. response = {'status': 0, 'message': ''}
  323. if request.method != 'GET':
  324. response['message'] = _('A GET request is required.')
  325. app_name = get_app_name(request)
  326. query_type = beeswax.models.SavedQuery.TYPES_MAPPING[app_name]
  327. design = safe_get_design(request, query_type, design_id)
  328. response['design'] = design_to_dict(design)
  329. return JsonResponse(response)
  330. @error_handler
  331. def fetch_query_history(request, query_history_id):
  332. response = {'status': 0, 'message': ''}
  333. if request.method != 'GET':
  334. response['message'] = _('A GET request is required.')
  335. query = authorized_get_query_history(request, query_history_id, must_exist=True)
  336. response['query_history'] = query_history_to_dict(request, query)
  337. return JsonResponse(response)
  338. @error_handler
  339. def cancel_query(request, query_history_id):
  340. response = {'status': -1, 'message': ''}
  341. if request.method != 'POST':
  342. response['message'] = _('A POST request is required.')
  343. else:
  344. try:
  345. query_history = authorized_get_query_history(request, query_history_id, must_exist=True)
  346. db = dbms.get(request.user, query_history.get_query_server_config())
  347. db.cancel_operation(query_history.get_handle())
  348. query_history.set_to_expired()
  349. response['status'] = 0
  350. except Exception, e:
  351. response['message'] = unicode(e)
  352. return JsonResponse(response)
  353. @error_handler
  354. def save_results_hdfs_directory(request, query_history_id):
  355. """
  356. Save the results of a query to an HDFS directory.
  357. Rerun the query.
  358. """
  359. response = {'status': 0, 'message': ''}
  360. query_history = authorized_get_query_history(request, query_history_id, must_exist=True)
  361. server_id, state = _get_query_handle_and_state(query_history)
  362. query_history.save_state(state)
  363. error_msg, log = None, None
  364. if request.method != 'POST':
  365. response['message'] = _('A POST request is required.')
  366. else:
  367. if not query_history.is_success():
  368. response['message'] = _('This query is %(state)s. Results unavailable.') % {'state': state}
  369. response['status'] = -1
  370. return JsonResponse(response)
  371. db = dbms.get(request.user, query_history.get_query_server_config())
  372. form = beeswax.forms.SaveResultsDirectoryForm({
  373. 'target_dir': request.POST.get('path')
  374. }, fs=request.fs)
  375. if form.is_valid():
  376. target_dir = request.POST.get('path')
  377. try:
  378. response['type'] = 'hdfs-dir'
  379. response['id'] = query_history.id
  380. response['query'] = query_history.query
  381. response['path'] = target_dir
  382. response['success_url'] = '/filebrowser/view=%s' % target_dir
  383. query_history = db.insert_query_into_directory(query_history, target_dir)
  384. response['watch_url'] = reverse(get_app_name(request) + ':api_watch_query_refresh_json', kwargs={'id': query_history.id})
  385. except Exception, ex:
  386. error_msg, log = expand_exception(ex, db)
  387. response['message'] = _('The result could not be saved: %s.') % error_msg
  388. response['status'] = -3
  389. else:
  390. response['status'] = 1
  391. response['errors'] = form.errors
  392. return JsonResponse(response)
  393. @error_handler
  394. def save_results_hdfs_file(request, query_history_id):
  395. """
  396. Save the results of a query to an HDFS file.
  397. Do not rerun the query.
  398. """
  399. response = {'status': 0, 'message': ''}
  400. query_history = authorized_get_query_history(request, query_history_id, must_exist=True)
  401. server_id, state = _get_query_handle_and_state(query_history)
  402. query_history.save_state(state)
  403. error_msg, log = None, None
  404. if request.method != 'POST':
  405. response['message'] = _('A POST request is required.')
  406. else:
  407. if not query_history.is_success():
  408. response['message'] = _('This query is %(state)s. Results unavailable.') % {'state': state}
  409. response['status'] = -1
  410. return JsonResponse(response)
  411. db = dbms.get(request.user, query_history.get_query_server_config())
  412. form = beeswax.forms.SaveResultsFileForm({
  413. 'target_file': request.POST.get('path'),
  414. 'overwrite': request.POST.get('overwrite', False),
  415. })
  416. if form.is_valid():
  417. target_file = form.cleaned_data['target_file']
  418. overwrite = form.cleaned_data['overwrite']
  419. try:
  420. handle, state = _get_query_handle_and_state(query_history)
  421. except Exception, ex:
  422. response['message'] = _('Cannot find query handle and state: %s') % str(query_history)
  423. response['status'] = -2
  424. return JsonResponse(response)
  425. try:
  426. if overwrite and request.fs.exists(target_file):
  427. if request.fs.isfile(target_file):
  428. request.fs.do_as_user(request.user.username, request.fs.rmtree, target_file)
  429. else:
  430. raise PopupException(_("The target path is a directory"))
  431. upload(target_file, handle, request.user, db, request.fs)
  432. response['type'] = 'hdfs-file'
  433. response['id'] = query_history.id
  434. response['query'] = query_history.query
  435. response['path'] = target_file
  436. response['success_url'] = '/filebrowser/view=%s' % target_file
  437. response['watch_url'] = reverse(get_app_name(request) + ':api_watch_query_refresh_json', kwargs={'id': query_history.id})
  438. except Exception, ex:
  439. error_msg, log = expand_exception(ex, db)
  440. response['message'] = _('The result could not be saved: %s.') % error_msg
  441. response['status'] = -3
  442. else:
  443. response['status'] = 1
  444. response['errors'] = form.errors
  445. return JsonResponse(response)
  446. @error_handler
  447. def save_results_hive_table(request, query_history_id):
  448. """
  449. Save the results of a query to a hive table.
  450. Rerun the query.
  451. """
  452. response = {'status': 0, 'message': ''}
  453. query_history = authorized_get_query_history(request, query_history_id, must_exist=True)
  454. server_id, state = _get_query_handle_and_state(query_history)
  455. query_history.save_state(state)
  456. error_msg, log = None, None
  457. if request.method != 'POST':
  458. response['message'] = _('A POST request is required.')
  459. else:
  460. if not query_history.is_success():
  461. response['message'] = _('This query is %(state)s. Results unavailable.') % {'state': state}
  462. response['status'] = -1
  463. return JsonResponse(response)
  464. db = dbms.get(request.user, query_history.get_query_server_config())
  465. database = query_history.design.get_design().query.get('database', 'default')
  466. form = beeswax.forms.SaveResultsTableForm({
  467. 'target_table': request.POST.get('table')
  468. }, db=db, database=database)
  469. if form.is_valid():
  470. try:
  471. handle, state = _get_query_handle_and_state(query_history)
  472. result_meta = db.get_results_metadata(handle)
  473. except Exception, ex:
  474. response['message'] = _('Cannot find query handle and state: %s') % str(query_history)
  475. response['status'] = -2
  476. return JsonResponse(response)
  477. try:
  478. query_history = db.create_table_as_a_select(request, query_history, form.target_database, form.cleaned_data['target_table'], result_meta)
  479. response['id'] = query_history.id
  480. response['query'] = query_history.query
  481. response['type'] = 'hive-table'
  482. response['path'] = form.cleaned_data['target_table']
  483. response['success_url'] = reverse('metastore:describe_table', kwargs={'database': form.target_database, 'table': form.cleaned_data['target_table']})
  484. response['watch_url'] = reverse(get_app_name(request) + ':api_watch_query_refresh_json', kwargs={'id': query_history.id})
  485. except Exception, ex:
  486. error_msg, log = expand_exception(ex, db)
  487. response['message'] = _('The result could not be saved: %s.') % error_msg
  488. response['status'] = -3
  489. else:
  490. response['status'] = 1
  491. response['message'] = '\n'.join(form.errors.values()[0])
  492. return JsonResponse(response)
  493. @error_handler
  494. def clear_history(request):
  495. response = {'status': -1, 'message': ''}
  496. if request.method != 'POST':
  497. response['message'] = _('A POST request is required.')
  498. else:
  499. response['count'] = QueryHistory.objects.filter(owner=request.user, is_cleared=False).update(is_cleared=True)
  500. response['status'] = 0
  501. return JsonResponse(response)
  502. @error_handler
  503. def get_sample_data(request, database, table, column=None):
  504. app_name = get_app_name(request)
  505. query_server = get_query_server_config(app_name)
  506. db = dbms.get(request.user, query_server)
  507. response = _get_sample_data(db, database, table, column)
  508. return JsonResponse(response)
  509. def _get_sample_data(db, database, table, column):
  510. table_obj = db.get_table(database, table)
  511. sample_data = db.get_sample(database, table_obj, column)
  512. response = {'status': -1}
  513. if sample_data:
  514. sample = escape_rows(sample_data.rows(), nulls_only=True)
  515. if column:
  516. sample = set([row[0] for row in sample])
  517. sample = [[item] for item in sorted(list(sample))]
  518. response['status'] = 0
  519. response['headers'] = sample_data.cols()
  520. response['full_headers'] = sample_data.full_cols()
  521. response['rows'] = sample
  522. else:
  523. response['message'] = _('Failed to get sample data.')
  524. return response
  525. @error_handler
  526. def get_indexes(request, database, table):
  527. query_server = dbms.get_query_server_config(get_app_name(request))
  528. db = dbms.get(request.user, query_server)
  529. response = {'status': -1}
  530. indexes = db.get_indexes(database, table)
  531. if indexes:
  532. response['status'] = 0
  533. response['headers'] = indexes.cols()
  534. response['rows'] = escape_rows(indexes.rows(), nulls_only=True)
  535. else:
  536. response['message'] = _('Failed to get indexes.')
  537. return JsonResponse(response)
  538. @error_handler
  539. def get_settings(request):
  540. query_server = dbms.get_query_server_config(get_app_name(request))
  541. db = dbms.get(request.user, query_server)
  542. response = {'status': -1}
  543. settings = db.get_configuration()
  544. if settings:
  545. response['status'] = 0
  546. response['settings'] = settings
  547. else:
  548. response['message'] = _('Failed to get settings.')
  549. return JsonResponse(response)
  550. @error_handler
  551. def get_functions(request):
  552. query_server = dbms.get_query_server_config(get_app_name(request))
  553. db = dbms.get(request.user, query_server)
  554. response = {'status': -1}
  555. prefix = request.GET.get('prefix', None)
  556. functions = db.get_functions(prefix)
  557. if functions:
  558. response['status'] = 0
  559. rows = escape_rows(functions.rows(), nulls_only=True)
  560. response['functions'] = [row[0] for row in rows]
  561. else:
  562. response['message'] = _('Failed to get functions.')
  563. return JsonResponse(response)
  564. @error_handler
  565. def analyze_table(request, database, table, columns=None):
  566. app_name = get_app_name(request)
  567. query_server = get_query_server_config(app_name)
  568. db = dbms.get(request.user, query_server)
  569. response = {'status': -1, 'message': '', 'redirect': ''}
  570. if request.method == "POST":
  571. if columns is None:
  572. query_history = db.analyze_table(database, table)
  573. else:
  574. query_history = db.analyze_table_columns(database, table)
  575. response['watch_url'] = reverse('beeswax:api_watch_query_refresh_json', kwargs={'id': query_history.id})
  576. response['status'] = 0
  577. else:
  578. response['message'] = _('A POST request is required.')
  579. return JsonResponse(response)
  580. @error_handler
  581. def get_table_stats(request, database, table, column=None):
  582. app_name = get_app_name(request)
  583. query_server = get_query_server_config(app_name)
  584. db = dbms.get(request.user, query_server)
  585. response = {'status': -1, 'message': '', 'redirect': ''}
  586. if column is not None:
  587. stats = db.get_table_columns_stats(database, table, column)
  588. else:
  589. table = db.get_table(database, table)
  590. stats = table.stats
  591. response['stats'] = stats
  592. response['status'] = 0
  593. return JsonResponse(response)
  594. @error_handler
  595. def get_top_terms(request, database, table, column, prefix=None):
  596. app_name = get_app_name(request)
  597. query_server = get_query_server_config(app_name)
  598. db = dbms.get(request.user, query_server)
  599. response = {'status': -1, 'message': '', 'redirect': ''}
  600. terms = db.get_top_terms(database, table, column, prefix=prefix, limit=int(request.GET.get('limit', 30)))
  601. response['terms'] = terms
  602. response['status'] = 0
  603. return JsonResponse(response)
  604. @error_handler
  605. def get_session(request, session_id=None):
  606. app_name = get_app_name(request)
  607. query_server = get_query_server_config(app_name)
  608. response = {'status': -1, 'message': ''}
  609. if session_id:
  610. session = Session.objects.get(id=session_id, owner=request.user, application=query_server['server_name'])
  611. else: # get the latest session for given user and server type
  612. session = Session.objects.get_session(request.user, query_server['server_name'])
  613. if session is not None:
  614. properties = json.loads(session.properties)
  615. # Redact passwords
  616. for key, value in properties.items():
  617. if 'password' in key.lower():
  618. properties[key] = '*' * len(value)
  619. response['status'] = 0
  620. response['session'] = {'id': session.id, 'application': session.application, 'status': session.status_code}
  621. response['properties'] = properties
  622. else:
  623. response['message'] = _('Could not find session or no open sessions found.')
  624. return JsonResponse(response)
  625. @require_POST
  626. @error_handler
  627. def close_session(request, session_id):
  628. app_name = get_app_name(request)
  629. query_server = get_query_server_config(app_name)
  630. response = {'status': -1, 'message': ''}
  631. try:
  632. filters = {'id': session_id, 'application': query_server['server_name']}
  633. if not request.user.is_superuser:
  634. filters['owner'] = request.user
  635. session = Session.objects.get(**filters)
  636. except Session.DoesNotExist:
  637. response['message'] = _('Session does not exist or you do not have permissions to close the session.')
  638. if session:
  639. session = dbms.get(request.user, query_server).close_session(session)
  640. response['status'] = 0
  641. response['message'] = _('Session successfully closed.')
  642. response['session'] = {'id': session_id, 'application': session.application, 'status': session.status_code}
  643. return JsonResponse(response)
  644. # Proxy API for Metastore App
  645. def describe_table(request, database, table):
  646. try:
  647. from metastore.views import describe_table
  648. return describe_table(request, database, table)
  649. except Exception, e:
  650. LOG.exception('Describe table failed')
  651. raise PopupException(_('Problem accessing table metadata'), detail=e)
  652. def design_to_dict(design):
  653. hql_design = HQLdesign.loads(design.data)
  654. return {
  655. 'id': design.id,
  656. 'query': hql_design.hql_query,
  657. 'name': design.name,
  658. 'desc': design.desc,
  659. 'database': hql_design.query.get('database', None),
  660. 'settings': hql_design.settings,
  661. 'file_resources': hql_design.file_resources,
  662. 'functions': hql_design.functions,
  663. 'is_parameterized': hql_design.query.get('is_parameterized', True),
  664. 'email_notify': hql_design.query.get('email_notify', True),
  665. 'is_redacted': design.is_redacted
  666. }
  667. def query_history_to_dict(request, query_history):
  668. query_history_dict = {
  669. 'id': query_history.id,
  670. 'state': query_history.last_state,
  671. 'query': query_history.query,
  672. 'has_results': query_history.has_results,
  673. 'statement_number': query_history.statement_number,
  674. 'watch_url': reverse(get_app_name(request) + ':api_watch_query_refresh_json', kwargs={'id': query_history.id}),
  675. 'results_url': reverse(get_app_name(request) + ':view_results', kwargs={'id': query_history.id, 'first_row': 0})
  676. }
  677. if query_history.design:
  678. query_history_dict['design'] = design_to_dict(query_history.design)
  679. return query_history_dict
  680. def get_query_form(request):
  681. try:
  682. try:
  683. # Get database choices
  684. query_server = dbms.get_query_server_config(get_app_name(request))
  685. db = dbms.get(request.user, query_server)
  686. databases = [(database, database) for database in db.get_databases()]
  687. except StructuredThriftTransportException, e:
  688. # If Thrift exception was due to failed authentication, raise corresponding message
  689. if 'TSocket read 0 bytes' in str(e) or 'Error validating the login' in str(e):
  690. raise PopupException(_('Failed to authenticate to query server, check authentication configurations.'), detail=e)
  691. else:
  692. raise e
  693. except Exception, e:
  694. raise PopupException(_('Unable to access databases, Query Server or Metastore may be down.'), detail=e)
  695. if not databases:
  696. raise RuntimeError(_("No databases are available. Permissions could be missing."))
  697. query_form = QueryForm()
  698. query_form.bind(request.POST)
  699. query_form.query.fields['database'].choices = databases # Could not do it in the form
  700. return query_form
  701. """
  702. Utils
  703. """
  704. def _extract_nested_type(parse_tree, nested_path):
  705. nested_tokens = nested_path.strip('/').split('/')
  706. subtree = parse_tree
  707. for token in nested_tokens:
  708. if token in subtree:
  709. subtree = subtree[token]
  710. elif 'fields' in subtree:
  711. for field in subtree['fields']:
  712. if field['name'] == token:
  713. subtree = field
  714. break
  715. else:
  716. raise Exception('Invalid nested type path: %s' % nested_path)
  717. return subtree