api.py 20 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608
  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 json
  18. import logging
  19. from django.core.urlresolvers import reverse
  20. from django.forms import ValidationError
  21. from django.http import HttpResponseBadRequest, HttpResponseRedirect
  22. from django.utils.translation import ugettext as _
  23. from django.views.decorators.http import require_GET, require_POST
  24. from desktop.lib.django_util import JsonResponse
  25. from desktop.models import Document2, Document
  26. from notebook.connectors.base import get_api, Notebook, QueryExpired, SessionExpired
  27. from notebook.decorators import api_error_handler, check_document_access_permission, check_document_modify_permission
  28. from notebook.github import GithubClient
  29. from notebook.models import escape_rows
  30. from notebook.views import upgrade_session_properties
  31. LOG = logging.getLogger(__name__)
  32. DEFAULT_HISTORY_NAME = ''
  33. @require_POST
  34. @api_error_handler
  35. def create_notebook(request):
  36. response = {'status': -1}
  37. editor_type = request.POST.get('type', 'notebook')
  38. directory_uuid = request.POST.get('directory_uuid')
  39. editor = Notebook()
  40. data = editor.get_data()
  41. if editor_type != 'notebook':
  42. data['name'] = ''
  43. data['type'] = 'query-%s' % editor_type # TODO: Add handling for non-SQL types
  44. data['directoryUuid'] = directory_uuid
  45. editor.data = json.dumps(data)
  46. response['notebook'] = editor.get_data()
  47. response['status'] = 0
  48. return JsonResponse(response)
  49. @require_POST
  50. @check_document_access_permission()
  51. @api_error_handler
  52. def create_session(request):
  53. response = {'status': -1}
  54. notebook = json.loads(request.POST.get('notebook', '{}'))
  55. session = json.loads(request.POST.get('session', '{}'))
  56. properties = session.get('properties', [])
  57. response['session'] = get_api(request, session).create_session(lang=session['type'], properties=properties)
  58. response['status'] = 0
  59. return JsonResponse(response)
  60. @require_POST
  61. @check_document_access_permission()
  62. @api_error_handler
  63. def close_session(request):
  64. response = {'status': -1}
  65. session = json.loads(request.POST.get('session', '{}'))
  66. response['session'] = get_api(request, {'type': session['type']}).close_session(session=session)
  67. response['status'] = 0
  68. return JsonResponse(response)
  69. @require_POST
  70. @check_document_access_permission()
  71. @api_error_handler
  72. def execute(request):
  73. response = {'status': -1}
  74. result = None
  75. notebook = json.loads(request.POST.get('notebook', '{}'))
  76. snippet = json.loads(request.POST.get('snippet', '{}'))
  77. try:
  78. response['handle'] = get_api(request, snippet).execute(notebook, snippet)
  79. # Retrieve and remove the result from the handle
  80. if response['handle'].get('sync'):
  81. result = response['handle'].pop('result')
  82. finally:
  83. if notebook['type'].startswith('query-'):
  84. _snippet = [s for s in notebook['snippets'] if s['id'] == snippet['id']][0]
  85. if 'handle' in response: # No failure
  86. _snippet['result']['handle'] = response['handle']
  87. _snippet['result']['statements_count'] = response['handle'].get('statements_count', 1)
  88. _snippet['result']['statement_id'] = response['handle'].get('statement_id', 0)
  89. _snippet['result']['handle']['statement'] = response['handle'].get('statement', snippet['statement']) # For non HS2, as non multi query yet
  90. else:
  91. _snippet['status'] = 'failed'
  92. history = _historify(notebook, request.user)
  93. response['history_id'] = history.id
  94. response['history_uuid'] = history.uuid
  95. if notebook['isSaved']: # Keep track of history of saved queries
  96. response['history_parent_uuid'] = history.dependencies.filter(type__startswith='query-').latest('last_modified').uuid
  97. # Inject and HTML escape results
  98. if result is not None:
  99. response['result'] = result
  100. response['result']['data'] = escape_rows(result['data'])
  101. response['status'] = 0
  102. return JsonResponse(response)
  103. @require_POST
  104. @check_document_access_permission()
  105. @api_error_handler
  106. def check_status(request):
  107. response = {'status': -1}
  108. notebook = json.loads(request.POST.get('notebook', '{}'))
  109. snippet = json.loads(request.POST.get('snippet', '{}'))
  110. if not snippet:
  111. nb_doc = Document2.objects.get_by_uuid(user=request.user, uuid=notebook['id'])
  112. notebook = Notebook(document=nb_doc).get_data()
  113. snippet = notebook['snippets'][0]
  114. try:
  115. response['query_status'] = get_api(request, snippet).check_status(notebook, snippet)
  116. response['status'] = 0
  117. except SessionExpired:
  118. response['status'] = 'expired'
  119. raise
  120. except QueryExpired:
  121. response['status'] = 'expired'
  122. raise
  123. finally:
  124. if response['status'] == 0 and snippet['status'] != response['query_status']:
  125. status = response['query_status']['status']
  126. elif response['status'] == 'expired':
  127. status = 'expired'
  128. else:
  129. status = 'failed'
  130. if notebook['type'].startswith('query'):
  131. nb_doc = Document2.objects.get(id=notebook['id'])
  132. nb_doc.can_write_or_exception(request.user)
  133. nb = Notebook(document=nb_doc).get_data()
  134. nb['snippets'][0]['status'] = status
  135. nb_doc.update_data(nb)
  136. nb_doc.save()
  137. return JsonResponse(response)
  138. @require_POST
  139. @check_document_access_permission()
  140. @api_error_handler
  141. def fetch_result_data(request):
  142. response = {'status': -1}
  143. notebook = json.loads(request.POST.get('notebook', '{}'))
  144. snippet = json.loads(request.POST.get('snippet', '{}'))
  145. rows = json.loads(request.POST.get('rows', 100))
  146. start_over = json.loads(request.POST.get('startOver', False))
  147. response['result'] = get_api(request, snippet).fetch_result(notebook, snippet, rows, start_over)
  148. # Materialize and HTML escape results
  149. if response['result'].get('data') and response['result'].get('type') == 'table':
  150. response['result']['data'] = escape_rows(response['result']['data'])
  151. response['status'] = 0
  152. return JsonResponse(response)
  153. @require_POST
  154. @check_document_access_permission()
  155. @api_error_handler
  156. def fetch_result_metadata(request):
  157. response = {'status': -1}
  158. notebook = json.loads(request.POST.get('notebook', '{}'))
  159. snippet = json.loads(request.POST.get('snippet', '{}'))
  160. response['result'] = get_api(request, snippet).fetch_result_metadata(notebook, snippet)
  161. response['status'] = 0
  162. return JsonResponse(response)
  163. @require_POST
  164. @check_document_access_permission()
  165. @api_error_handler
  166. def cancel_statement(request):
  167. response = {'status': -1}
  168. notebook = json.loads(request.POST.get('notebook', '{}'))
  169. snippet = json.loads(request.POST.get('snippet', '{}'))
  170. response['result'] = get_api(request, snippet).cancel(notebook, snippet)
  171. response['status'] = 0
  172. return JsonResponse(response)
  173. @require_POST
  174. @check_document_access_permission()
  175. @api_error_handler
  176. def get_logs(request):
  177. response = {'status': -1}
  178. notebook = json.loads(request.POST.get('notebook', '{}'))
  179. snippet = json.loads(request.POST.get('snippet', '{}'))
  180. startFrom = request.POST.get('from')
  181. startFrom = int(startFrom) if startFrom else None
  182. size = request.POST.get('size')
  183. size = int(size) if size else None
  184. db = get_api(request, snippet)
  185. logs = db.get_log(notebook, snippet, startFrom=startFrom, size=size)
  186. jobs = json.loads(request.POST.get('jobs', '[]'))
  187. # Get any new jobs from current logs snippet
  188. new_jobs = db.get_jobs(notebook, snippet, logs)
  189. # Append new jobs to known jobs and get the unique set
  190. if new_jobs:
  191. all_jobs = jobs + new_jobs
  192. jobs = dict((job['name'], job) for job in all_jobs).values()
  193. # Retrieve full log for job progress parsing
  194. full_log = request.POST.get('full_log', logs)
  195. response['logs'] = logs
  196. response['progress'] = db.progress(snippet, full_log) if snippet['status'] != 'available' and snippet['status'] != 'success' else 100
  197. response['jobs'] = jobs
  198. response['status'] = 0
  199. return JsonResponse(response)
  200. @require_POST
  201. @check_document_modify_permission()
  202. def save_notebook(request):
  203. response = {'status': -1}
  204. notebook = json.loads(request.POST.get('notebook', '{}'))
  205. notebook_type = notebook.get('type', 'notebook')
  206. if notebook.get('parentSavedQueryUuid'): # We save into the original saved query, not into the query history
  207. notebook_doc = Document2.objects.get_by_uuid(user=request.user, uuid=notebook['parentSavedQueryUuid'])
  208. elif notebook.get('id'):
  209. notebook_doc = Document2.objects.get(id=notebook['id'])
  210. else:
  211. notebook_doc = Document2.objects.create(name=notebook['name'], uuid=notebook['uuid'], type=notebook_type, owner=request.user)
  212. Document.objects.link(notebook_doc, owner=notebook_doc.owner, name=notebook_doc.name, description=notebook_doc.description, extra=notebook_type)
  213. if notebook.get('directoryUuid'):
  214. notebook_doc.parent_directory = Document2.objects.get_by_uuid(user=request.user, uuid=notebook.get('directoryUuid'), perm_type='write')
  215. else:
  216. notebook_doc.parent_directory = Document2.objects.get_home_directory(request.user)
  217. notebook['isSaved'] = True
  218. notebook['isHistory'] = False
  219. notebook['id'] = notebook_doc.id
  220. notebook_doc1 = notebook_doc.doc.get()
  221. notebook_doc.update_data(notebook)
  222. notebook_doc.name = notebook_doc1.name = notebook['name']
  223. notebook_doc.description = notebook_doc1.description = notebook['description']
  224. notebook_doc.save()
  225. notebook_doc1.save()
  226. response['status'] = 0
  227. response['id'] = notebook_doc.id
  228. response['message'] = request.POST.get('editorMode') == 'true' and _('Query saved successfully') or _('Notebook saved successfully')
  229. return JsonResponse(response)
  230. def _historify(notebook, user):
  231. query_type = notebook['type']
  232. name = notebook['name'] if (notebook['name'] and notebook['name'].strip() != '') else DEFAULT_HISTORY_NAME
  233. history_doc = Document2.objects.create(
  234. name=name,
  235. type=query_type,
  236. owner=user,
  237. is_history=True
  238. )
  239. # Link history of saved query
  240. if notebook['isSaved']:
  241. parent_doc = Document2.objects.get(uuid=notebook.get('parentSavedQueryUuid') or notebook['uuid']) # From previous history query or initial saved query
  242. notebook['parentSavedQueryUuid'] = parent_doc.uuid
  243. history_doc.dependencies.add(parent_doc)
  244. Document.objects.link(
  245. history_doc,
  246. name=history_doc.name,
  247. owner=history_doc.owner,
  248. description=history_doc.description,
  249. extra=query_type
  250. )
  251. notebook['uuid'] = history_doc.uuid
  252. history_doc.update_data(notebook)
  253. history_doc.save()
  254. return history_doc
  255. @require_GET
  256. @api_error_handler
  257. @check_document_access_permission()
  258. def get_history(request):
  259. response = {'status': -1}
  260. doc_type = request.GET.get('doc_type')
  261. limit = min(request.GET.get('len', 50), 100)
  262. response['status'] = 0
  263. history = []
  264. for doc in Document2.objects.get_history(doc_type='query-%s' % doc_type, user=request.user).order_by('-last_modified')[:limit]:
  265. notebook = Notebook(document=doc).get_data()
  266. if 'snippets' in notebook:
  267. try:
  268. statement = notebook['snippets'][0]['result']['handle']['statement']
  269. if type(statement) == dict: # Old format
  270. statement = notebook['snippets'][0]['statement_raw']
  271. except KeyError: # Old format
  272. statement = notebook['snippets'][0]['statement_raw']
  273. history.append({
  274. 'name': doc.name,
  275. 'id': doc.id,
  276. 'uuid': doc.uuid,
  277. 'type': doc.type,
  278. 'data': {
  279. 'statement': statement[:1001] if statement else '',
  280. 'lastExecuted': notebook['snippets'][0]['lastExecuted'],
  281. 'status': notebook['snippets'][0]['status'],
  282. 'parentSavedQueryUuid': notebook.get('parentSavedQueryUuid', '')
  283. } if notebook['snippets'] else {},
  284. 'absoluteUrl': doc.get_absolute_url(),
  285. })
  286. else:
  287. LOG.error('Incomplete History Notebook: %s' % notebook)
  288. response['history'] = sorted(history, key=lambda row: row['data']['lastExecuted'], reverse=True)
  289. response['message'] = _('History fetched')
  290. return JsonResponse(response)
  291. @require_POST
  292. @api_error_handler
  293. @check_document_modify_permission()
  294. def clear_history(request):
  295. response = {'status': -1}
  296. notebook = json.loads(request.POST.get('notebook'), '{}')
  297. doc_type = request.POST.get('doc_type')
  298. history = Document2.objects.get_history(doc_type='query-%s' % doc_type, user=request.user)
  299. response['updated'] = history.delete()
  300. response['message'] = _('History cleared !')
  301. response['status'] = 0
  302. return JsonResponse(response)
  303. @require_GET
  304. @check_document_access_permission()
  305. def open_notebook(request):
  306. response = {'status': -1}
  307. notebook_id = request.GET.get('notebook')
  308. notebook = Notebook(document=Document2.objects.get(id=notebook_id))
  309. notebook = upgrade_session_properties(request, notebook)
  310. response['status'] = 0
  311. response['notebook'] = notebook.get_json()
  312. response['message'] = _('Notebook loaded successfully')
  313. @require_POST
  314. @check_document_access_permission()
  315. def close_notebook(request):
  316. response = {'status': -1, 'result': []}
  317. notebook = json.loads(request.POST.get('notebook', '{}'))
  318. for session in [_s for _s in notebook['sessions'] if _s['type'] in ('scala', 'spark', 'pyspark', 'sparkr')]:
  319. try:
  320. response['result'].append(get_api(request, session).close_session(session))
  321. except QueryExpired:
  322. pass
  323. except Exception, e:
  324. LOG.exception('Error closing session %s' % str(e))
  325. for snippet in [_s for _s in notebook['snippets'] if _s['type'] in ('hive', 'impala')]:
  326. try:
  327. response['result'] = get_api(request, snippet).close_statement(snippet)
  328. except QueryExpired:
  329. pass
  330. except Exception, e:
  331. LOG.exception('Error closing statement %s' % str(e))
  332. response['status'] = 0
  333. response['message'] = _('Notebook closed successfully')
  334. return JsonResponse(response)
  335. @require_POST
  336. @check_document_access_permission()
  337. def close_statement(request):
  338. response = {'status': -1}
  339. # Passed by check_document_access_permission but unused by APIs
  340. notebook = json.loads(request.POST.get('notebook', '{}'))
  341. snippet = json.loads(request.POST.get('snippet', '{}'))
  342. try:
  343. response['result'] = get_api(request, snippet).close_statement(snippet)
  344. except QueryExpired:
  345. pass
  346. response['status'] = 0
  347. response['message'] = _('Statement closed !')
  348. return JsonResponse(response)
  349. @require_POST
  350. @check_document_access_permission()
  351. @api_error_handler
  352. def autocomplete(request, server=None, database=None, table=None, column=None, nested=None):
  353. response = {'status': -1}
  354. # Passed by check_document_access_permission but unused by APIs
  355. notebook = json.loads(request.POST.get('notebook', '{}'))
  356. snippet = json.loads(request.POST.get('snippet', '{}'))
  357. try:
  358. autocomplete_data = get_api(request, snippet).autocomplete(snippet, database, table, column, nested)
  359. response.update(autocomplete_data)
  360. except QueryExpired:
  361. pass
  362. response['status'] = 0
  363. return JsonResponse(response)
  364. @require_POST
  365. @check_document_access_permission()
  366. @api_error_handler
  367. def get_sample_data(request, server=None, database=None, table=None, column=None):
  368. response = {'status': -1}
  369. # Passed by check_document_access_permission but unused by APIs
  370. notebook = json.loads(request.POST.get('notebook', '{}'))
  371. snippet = json.loads(request.POST.get('snippet', '{}'))
  372. sample_data = get_api(request, snippet).get_sample_data(snippet, database, table, column)
  373. response.update(sample_data)
  374. response['status'] = 0
  375. return JsonResponse(response)
  376. @require_POST
  377. @check_document_access_permission()
  378. @api_error_handler
  379. def explain(request):
  380. response = {'status': -1}
  381. notebook = json.loads(request.POST.get('notebook', '{}'))
  382. snippet = json.loads(request.POST.get('snippet', '{}'))
  383. response = get_api(request, snippet).explain(notebook, snippet)
  384. return JsonResponse(response)
  385. @require_GET
  386. @api_error_handler
  387. def github_fetch(request):
  388. response = {'status': -1}
  389. api = GithubClient(access_token=request.session.get('github_access_token'))
  390. response['url'] = url = request.GET.get('url')
  391. if url:
  392. owner, repo, branch, filepath = api.parse_github_url(url)
  393. content = api.get_file_contents(owner, repo, filepath, branch)
  394. try:
  395. response['content'] = json.loads(content)
  396. except ValueError:
  397. # Content is not JSON-encoded so return plain-text
  398. response['content'] = content
  399. response['status'] = 0
  400. else:
  401. return HttpResponseBadRequest(_('url param is required'))
  402. return JsonResponse(response)
  403. @api_error_handler
  404. def github_authorize(request):
  405. access_token = request.session.get('github_access_token')
  406. if access_token and GithubClient.is_authenticated(access_token):
  407. response = {
  408. 'status': 0,
  409. 'message': _('User is already authenticated to GitHub.')
  410. }
  411. return JsonResponse(response)
  412. else:
  413. auth_url = GithubClient.get_authorization_url()
  414. request.session['github_callback_redirect'] = request.GET.get('currentURL')
  415. request.session['github_callback_fetch'] = request.GET.get('fetchURL')
  416. response = {
  417. 'status': -1,
  418. 'auth_url':auth_url
  419. }
  420. if (request.is_ajax()):
  421. return JsonResponse(response)
  422. return HttpResponseRedirect(auth_url)
  423. @api_error_handler
  424. def github_callback(request):
  425. redirect_base = request.session['github_callback_redirect'] + "&github_status="
  426. if 'code' in request.GET:
  427. session_code = request.GET.get('code')
  428. request.session['github_access_token'] = GithubClient.get_access_token(session_code)
  429. return HttpResponseRedirect(redirect_base + "0&github_fetch=" + request.session['github_callback_fetch'])
  430. else:
  431. return HttpResponseRedirect(redirect_base + "-1&github_fetch=" + request.session['github_callback_fetch'])
  432. @require_POST
  433. @check_document_access_permission()
  434. @api_error_handler
  435. def export_result(request):
  436. response = {'status': -1, 'message': _('Exporting result failed.')}
  437. # Passed by check_document_access_permission but unused by APIs
  438. notebook = json.loads(request.POST.get('notebook', '{}'))
  439. snippet = json.loads(request.POST.get('snippet', '{}'))
  440. data_format = json.loads(request.POST.get('format', 'hdfs-file'))
  441. destination = json.loads(request.POST.get('destination', ''))
  442. overwrite = json.loads(request.POST.get('overwrite', False))
  443. api = get_api(request, snippet)
  444. if data_format == 'hdfs-file':
  445. if overwrite and request.fs.exists(destination):
  446. if request.fs.isfile(destination):
  447. request.fs.do_as_user(request.user.username, request.fs.rmtree, destination)
  448. else:
  449. raise ValidationError(_("The target path is a directory"))
  450. response['watch_url'] = api.export_data_as_hdfs_file(snippet, destination, overwrite)
  451. response['status'] = 0
  452. elif data_format == 'hive-table':
  453. notebook_id = notebook['id'] or request.GET.get('editor', request.GET.get('notebook'))
  454. response['watch_url'] = reverse('notebook:execute_and_watch') + '?action=save_as_table&notebook=' + str(notebook_id) + '&snippet=0&destination=' + destination
  455. response['status'] = 0
  456. elif data_format == 'hdfs-directory':
  457. notebook_id = notebook['id'] or request.GET.get('editor', request.GET.get('notebook'))
  458. response['watch_url'] = reverse('notebook:execute_and_watch') + '?action=insert_as_query&notebook=' + str(notebook_id) + '&snippet=0&destination=' + destination
  459. response['status'] = 0
  460. return JsonResponse(response)