models.py 18 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526
  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 base64
  18. import datetime
  19. import json
  20. import logging
  21. from django.db import models
  22. from django.contrib.auth.models import User
  23. from django.contrib.contenttypes.fields import GenericRelation
  24. from django.urls import reverse
  25. from django.utils.translation import ugettext as _, ugettext_lazy as _t
  26. from enum import Enum
  27. from TCLIService.ttypes import TSessionHandle, THandleIdentifier, TOperationState, TOperationHandle, TOperationType
  28. from desktop.redaction import global_redaction_engine
  29. from desktop.lib.exceptions_renderable import PopupException
  30. from desktop.models import Document
  31. from librdbms.server import dbms as librdbms_dbms
  32. from beeswax.design import HQLdesign
  33. LOG = logging.getLogger(__name__)
  34. QUERY_SUBMISSION_TIMEOUT = datetime.timedelta(0, 60 * 60) # 1 hour
  35. # Constants for DB fields, hue ini
  36. BEESWAX = 'beeswax'
  37. HIVE_SERVER2 = 'hiveserver2'
  38. QUERY_TYPES = (HQL, IMPALA, RDBMS, SPARK) = range(4)
  39. class QueryHistory(models.Model):
  40. """
  41. Holds metadata about all queries that have been executed.
  42. """
  43. class STATE(Enum):
  44. submitted = 0
  45. running = 1
  46. available = 2
  47. failed = 3
  48. expired = 4
  49. SERVER_TYPE = ((BEESWAX, 'Beeswax'), (HIVE_SERVER2, 'Hive Server 2'),
  50. (librdbms_dbms.MYSQL, 'MySQL'), (librdbms_dbms.POSTGRESQL, 'PostgreSQL'),
  51. (librdbms_dbms.SQLITE, 'sqlite'), (librdbms_dbms.ORACLE, 'oracle'))
  52. owner = models.ForeignKey(User, db_index=True)
  53. query = models.TextField()
  54. last_state = models.IntegerField(db_index=True)
  55. has_results = models.BooleanField(default=False) # If true, this query will eventually return tabular results.
  56. submission_date = models.DateTimeField(auto_now_add=True)
  57. # In case of multi statements in a query, these are the id of the currently running statement
  58. server_id = models.CharField(max_length=1024, null=True) # Aka secret, only query in the "submitted" state is allowed to have no server_id
  59. server_guid = models.CharField(max_length=1024, null=True, default=None)
  60. statement_number = models.SmallIntegerField(default=0) # The index of the currently running statement
  61. operation_type = models.SmallIntegerField(null=True)
  62. modified_row_count = models.FloatField(null=True)
  63. log_context = models.CharField(max_length=1024, null=True)
  64. server_host = models.CharField(max_length=128, help_text=_('Host of the query server.'), default='')
  65. server_port = models.PositiveIntegerField(help_text=_('Port of the query server.'), default=10000)
  66. server_name = models.CharField(max_length=128, help_text=_('Name of the query server.'), default='')
  67. server_type = models.CharField(max_length=128, help_text=_('Type of the query server.'), default=BEESWAX, choices=SERVER_TYPE)
  68. query_type = models.SmallIntegerField(help_text=_('Type of the query.'), default=HQL, choices=((HQL, 'HQL'), (IMPALA, 'IMPALA')))
  69. design = models.ForeignKey('SavedQuery', to_field='id', null=True) # Some queries (like read/create table) don't have a design
  70. notify = models.BooleanField(default=False) # Notify on completion
  71. is_redacted = models.BooleanField(default=False)
  72. extra = models.TextField(default='{}') # Json fields for extra properties
  73. is_cleared = models.BooleanField(default=False)
  74. class Meta:
  75. ordering = ['-submission_date']
  76. @staticmethod
  77. def build(*args, **kwargs):
  78. return HiveServerQueryHistory(*args, **kwargs)
  79. def get_full_object(self):
  80. return HiveServerQueryHistory.objects.get(id=self.id)
  81. @staticmethod
  82. def get(id):
  83. return HiveServerQueryHistory.objects.get(id=id)
  84. @staticmethod
  85. def get_type_name(query_type):
  86. if query_type == IMPALA:
  87. return 'impala'
  88. elif query_type == RDBMS:
  89. return 'rdbms'
  90. elif query_type == SPARK:
  91. return 'spark'
  92. else:
  93. return 'beeswax'
  94. def get_query_server_config(self):
  95. from beeswax.server.dbms import get_query_server_config
  96. query_server = get_query_server_config(QueryHistory.get_type_name(self.query_type))
  97. query_server.update({
  98. 'server_name': self.server_name,
  99. # 'server_host': self.server_host, # Always use the live server configuration as the session is currently tied to the connection
  100. # 'server_port': int(self.server_port),
  101. 'server_type': self.server_type,
  102. })
  103. return query_server
  104. def get_current_statement(self):
  105. if self.design is not None:
  106. design = self.design.get_design()
  107. return design.get_query_statement(self.statement_number)
  108. else:
  109. return self.query
  110. def refresh_design(self, hql_query):
  111. # Refresh only HQL query part
  112. query = self.design.get_design()
  113. query.hql_query = hql_query
  114. self.design.data = query.dumps()
  115. self.query = hql_query
  116. def is_finished(self):
  117. is_statement_finished = not self.is_running()
  118. if self.design is not None:
  119. design = self.design.get_design()
  120. return is_statement_finished and self.statement_number + 1 == design.statement_count # Last statement
  121. else:
  122. return is_statement_finished
  123. def is_running(self):
  124. return self.last_state in (QueryHistory.STATE.running.value, QueryHistory.STATE.submitted.value)
  125. def is_success(self):
  126. return self.last_state in (QueryHistory.STATE.available.value,)
  127. def is_failure(self):
  128. return self.last_state in (QueryHistory.STATE.expired.value, QueryHistory.STATE.failed.value)
  129. def is_expired(self):
  130. return self.last_state in (QueryHistory.STATE.expired.value,)
  131. def set_to_running(self):
  132. self.last_state = QueryHistory.STATE.running.value
  133. def set_to_failed(self):
  134. self.last_state = QueryHistory.STATE.failed.value
  135. def set_to_available(self):
  136. self.last_state = QueryHistory.STATE.available.value
  137. def set_to_expired(self):
  138. self.last_state = QueryHistory.STATE.expired.value
  139. def save(self, *args, **kwargs):
  140. """
  141. Override `save` to optionally mask out the query from being saved to the
  142. database. This is because if the beeswax database contains sensitive
  143. information like personally identifiable information, that information
  144. could be leaked into the Hue database and logfiles.
  145. """
  146. if global_redaction_engine.is_enabled():
  147. redacted_query = global_redaction_engine.redact(self.query)
  148. if self.query != redacted_query:
  149. self.query = redacted_query
  150. self.is_redacted = True
  151. super(QueryHistory, self).save(*args, **kwargs)
  152. def update_extra(self, key, val):
  153. extra = json.loads(self.extra)
  154. extra[key] = val
  155. self.extra = json.dumps(extra)
  156. def get_extra(self, key):
  157. return json.loads(self.extra).get(key)
  158. def make_query_context(type, info):
  159. """
  160. ``type`` is one of "table" and "design", and ``info`` is the table name or design id.
  161. Returns a value suitable for GET param.
  162. """
  163. if type == 'table':
  164. return "%s:%s" % (type, info)
  165. elif type == 'design':
  166. # Use int() to validate that info is a number
  167. return "%s:%s" % (type, int(info))
  168. LOG.error("Invalid query context type: %s" % (type,))
  169. return '' # Empty string is safer than None
  170. class HiveServerQueryHistory(QueryHistory):
  171. # Map from (thrift) server state
  172. STATE_MAP = {
  173. TOperationState.INITIALIZED_STATE : QueryHistory.STATE.submitted,
  174. TOperationState.RUNNING_STATE : QueryHistory.STATE.running,
  175. TOperationState.FINISHED_STATE : QueryHistory.STATE.available,
  176. TOperationState.CANCELED_STATE : QueryHistory.STATE.failed,
  177. TOperationState.CLOSED_STATE : QueryHistory.STATE.expired,
  178. TOperationState.ERROR_STATE : QueryHistory.STATE.failed,
  179. TOperationState.UKNOWN_STATE : QueryHistory.STATE.failed,
  180. TOperationState.PENDING_STATE : QueryHistory.STATE.submitted,
  181. }
  182. node_type = HIVE_SERVER2
  183. class Meta:
  184. proxy = True
  185. def get_handle(self):
  186. secret, guid = HiveServerQueryHandle.get_decoded(self.server_id, self.server_guid)
  187. return HiveServerQueryHandle(secret=secret,
  188. guid=guid,
  189. has_result_set=self.has_results,
  190. operation_type=self.operation_type,
  191. modified_row_count=self.modified_row_count)
  192. def save_state(self, new_state):
  193. self.last_state = new_state.value
  194. self.save()
  195. @classmethod
  196. def is_canceled(self, res):
  197. return res.operationState in (TOperationState.CANCELED_STATE, TOperationState.CLOSED_STATE)
  198. class SavedQuery(models.Model):
  199. """
  200. Stores the query that people have save or submitted.
  201. Note that this used to be called QueryDesign. Any references to 'design'
  202. probably mean a SavedQuery.
  203. """
  204. DEFAULT_NEW_DESIGN_NAME = _('My saved query')
  205. AUTO_DESIGN_SUFFIX = _(' (new)')
  206. TYPES = QUERY_TYPES
  207. TYPES_MAPPING = {'beeswax': HQL, 'hql': HQL, 'impala': IMPALA, 'rdbms': RDBMS, 'spark': SPARK}
  208. type = models.IntegerField(null=False)
  209. owner = models.ForeignKey(User, db_index=True)
  210. # Data is a json of dictionary. See the beeswax.design module.
  211. data = models.TextField(max_length=65536)
  212. name = models.CharField(max_length=80)
  213. desc = models.TextField(max_length=1024)
  214. mtime = models.DateTimeField(auto_now=True)
  215. # An auto design is a place-holder for things users submit but not saved.
  216. # We still want to store it as a design to allow users to save them later.
  217. is_auto = models.BooleanField(default=False, db_index=True)
  218. is_trashed = models.BooleanField(default=False, db_index=True, verbose_name=_t('Is trashed'),
  219. help_text=_t('If this query is trashed.'))
  220. is_redacted = models.BooleanField(default=False)
  221. doc = GenericRelation(Document, related_query_name='hql_doc')
  222. class Meta:
  223. ordering = ['-mtime']
  224. def get_design(self):
  225. try:
  226. return HQLdesign.loads(self.data)
  227. except ValueError:
  228. # data is empty
  229. pass
  230. def clone(self, new_owner=None):
  231. if new_owner is None:
  232. new_owner = self.owner
  233. design = SavedQuery(type=self.type, owner=new_owner)
  234. design.data = self.data
  235. design.name = self.name
  236. design.desc = self.desc
  237. design.is_auto = self.is_auto
  238. return design
  239. @classmethod
  240. def create_empty(cls, app_name, owner, data):
  241. query_type = SavedQuery.TYPES_MAPPING[app_name]
  242. design = SavedQuery(owner=owner, type=query_type)
  243. design.name = SavedQuery.DEFAULT_NEW_DESIGN_NAME
  244. design.desc = ''
  245. if global_redaction_engine.is_enabled():
  246. design.data = global_redaction_engine.redact(data)
  247. else:
  248. design.data = data
  249. design.is_auto = True
  250. design.save()
  251. Document.objects.link(design, owner=design.owner, extra=design.type, name=design.name, description=design.desc)
  252. design.doc.get().add_to_history()
  253. return design
  254. @staticmethod
  255. def get(id, owner=None, type=None):
  256. """
  257. get(id, owner=None, type=None) -> SavedQuery object
  258. Checks that the owner and type match (when given).
  259. May raise PopupException (type/owner mismatch).
  260. May raise SavedQuery.DoesNotExist.
  261. """
  262. try:
  263. design = SavedQuery.objects.get(id=id)
  264. except SavedQuery.DoesNotExist, err:
  265. msg = _('Cannot retrieve query id %(id)s.') % {'id': id}
  266. raise err
  267. if owner is not None and design.owner != owner:
  268. msg = _('Query id %(id)s does not belong to user %(user)s.') % {'id': id, 'user': owner}
  269. LOG.error(msg)
  270. raise PopupException(msg)
  271. if type is not None and design.type != type:
  272. msg = _('Type mismatch for design id %(id)s (owner %(owner)s) - Expected %(expected_type)s, got %(real_type)s.') % \
  273. {'id': id, 'owner': owner, 'expected_type': design.type, 'real_type': type}
  274. LOG.error(msg)
  275. raise PopupException(msg)
  276. return design
  277. def __str__(self):
  278. return '%s %s' % (self.name, self.owner)
  279. def get_query_context(self):
  280. try:
  281. return make_query_context('design', self.id)
  282. except:
  283. LOG.exception('failed to make query context')
  284. return ""
  285. def get_absolute_url(self):
  286. return reverse(QueryHistory.get_type_name(self.type) + ':execute_design', kwargs={'design_id': self.id})
  287. def save(self, *args, **kwargs):
  288. """
  289. Override `save` to optionally mask out the query from being saved to the
  290. database. This is because if the beeswax database contains sensitive
  291. information like personally identifiable information, that information
  292. could be leaked into the Hue database and logfiles.
  293. """
  294. if global_redaction_engine.is_enabled():
  295. data = json.loads(self.data)
  296. try:
  297. query = data['query']['query']
  298. except KeyError:
  299. pass
  300. else:
  301. redacted_query = global_redaction_engine.redact(query)
  302. if query != redacted_query:
  303. data['query']['query'] = redacted_query
  304. self.is_redacted = True
  305. self.data = json.dumps(data)
  306. super(SavedQuery, self).save(*args, **kwargs)
  307. class SessionManager(models.Manager):
  308. def get_session(self, user, application='beeswax', filter_open=True):
  309. try:
  310. q = self.filter(owner=user, application=application).exclude(guid='').exclude(secret='')
  311. if filter_open:
  312. q = q.filter(status_code=0)
  313. return q.latest("last_used")
  314. except Session.DoesNotExist, e:
  315. return None
  316. def get_n_sessions(self, user, n, application='beeswax', filter_open=True):
  317. q = self.filter(owner=user, application=application).exclude(guid='').exclude(secret='')
  318. if filter_open:
  319. q = q.filter(status_code=0)
  320. return q.order_by("-last_used")[0:n]
  321. class Session(models.Model):
  322. """
  323. A sessions is bound to a user and an application (e.g. Bob with the Impala application).
  324. """
  325. owner = models.ForeignKey(User, db_index=True)
  326. status_code = models.PositiveSmallIntegerField() # ttypes.TStatusCode
  327. secret = models.TextField(max_length='100')
  328. guid = models.TextField(max_length='100')
  329. server_protocol_version = models.SmallIntegerField(default=0)
  330. last_used = models.DateTimeField(auto_now=True, db_index=True, verbose_name=_t('Last used'))
  331. application = models.CharField(max_length=128, help_text=_t('Application we communicate with.'), default='beeswax')
  332. properties = models.TextField(default='{}')
  333. objects = SessionManager()
  334. def get_handle(self):
  335. secret, guid = HiveServerQueryHandle.get_decoded(secret=self.secret, guid=self.guid)
  336. handle_id = THandleIdentifier(secret=secret, guid=guid)
  337. return TSessionHandle(sessionId=handle_id)
  338. def get_properties(self):
  339. return json.loads(self.properties) if self.properties else {}
  340. def get_formatted_properties(self):
  341. return [dict({'key': key, 'value': value}) for key, value in self.get_properties().items()]
  342. def __str__(self):
  343. return '%s %s' % (self.owner, self.last_used)
  344. class QueryHandle(object):
  345. def __init__(self, secret=None, guid=None, operation_type=None, has_result_set=None, modified_row_count=None, log_context=None, session_guid=None):
  346. self.secret = secret
  347. self.guid = guid
  348. self.operation_type = operation_type
  349. self.has_result_set = has_result_set
  350. self.modified_row_count = modified_row_count
  351. self.log_context = log_context
  352. def is_valid(self):
  353. return sum([bool(obj) for obj in [self.get()]]) > 0
  354. def __str__(self):
  355. return '%s %s' % (self.secret, self.guid)
  356. class HiveServerQueryHandle(QueryHandle):
  357. """
  358. QueryHandle for Hive Server 2.
  359. Store THandleIdentifier base64 encoded in order to be unicode compatible with Django.
  360. Also store session handle if provided.
  361. """
  362. def __init__(self, **kwargs):
  363. super(HiveServerQueryHandle, self).__init__(**kwargs)
  364. self.secret, self.guid = self.get_encoded()
  365. self.session_guid = kwargs.get('session_guid')
  366. def get(self):
  367. return self.secret, self.guid
  368. def get_rpc_handle(self):
  369. secret, guid = self.get_decoded(self.secret, self.guid)
  370. operation = getattr(TOperationType, TOperationType._NAMES_TO_VALUES.get(self.operation_type, 'EXECUTE_STATEMENT'))
  371. return TOperationHandle(operationId=THandleIdentifier(guid=guid, secret=secret),
  372. operationType=operation,
  373. hasResultSet=self.has_result_set,
  374. modifiedRowCount=self.modified_row_count)
  375. @classmethod
  376. def get_decoded(cls, secret, guid):
  377. return base64.decodestring(secret), base64.decodestring(guid)
  378. def get_encoded(self):
  379. return base64.encodestring(self.secret), base64.encodestring(self.guid)
  380. # Deprecated. Could be removed.
  381. class BeeswaxQueryHandle(QueryHandle):
  382. """
  383. QueryHandle for Beeswax.
  384. """
  385. def __init__(self, secret, has_result_set, log_context):
  386. super(BeeswaxQueryHandle, self).__init__(secret=secret,
  387. has_result_set=has_result_set,
  388. log_context=log_context)
  389. def get(self):
  390. return self.secret, None
  391. def get_rpc_handle(self):
  392. return BeeswaxdQueryHandle(id=self.secret, log_context=self.log_context)
  393. # TODO remove
  394. def get_encoded(self):
  395. return self.get(), None
  396. class MetaInstall(models.Model):
  397. """
  398. Metadata about the installation. Should have at most one row.
  399. """
  400. installed_example = models.BooleanField(default=False)
  401. @staticmethod
  402. def get():
  403. """
  404. MetaInstall.get() -> MetaInstall object
  405. It helps dealing with that this table has just one row.
  406. """
  407. try:
  408. return MetaInstall.objects.get(id=1)
  409. except MetaInstall.DoesNotExist:
  410. return MetaInstall(id=1)