api.py 40 KB


  1. #!/usr/bin/env python
  2. # -- coding: utf-8 --
  3. # Licensed to Cloudera, Inc. under one
  4. # or more contributor license agreements. See the NOTICE file
  5. # distributed with this work for additional information
  6. # regarding copyright ownership. Cloudera, Inc. licenses this file
  7. # to you under the Apache License, Version 2.0 (the
  8. # "License"); you may not use this file except in compliance
  9. # with the License. You may obtain a copy of the License at
  10. #
  11. # http://www.apache.org/licenses/LICENSE-2.0
  12. #
  13. # Unless required by applicable law or agreed to in writing, software
  14. # distributed under the License is distributed on an "AS IS" BASIS,
  15. # WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
  16. # See the License for the specific language governing permissions and
  17. # limitations under the License.
  18. import json
  19. import logging
  20. import re
  21. import urllib
  22. from itertools import groupby
  23. from django.utils.translation import ugettext as _
  24. from dashboard.facet_builder import _compute_range_facet
  25. from dashboard.models import Collection2
  26. from desktop.lib.exceptions_renderable import PopupException
  27. from desktop.conf import SERVER_USER
  28. from desktop.lib.i18n import force_unicode
  29. from desktop.lib.rest.http_client import HttpClient, RestException
  30. from desktop.lib.rest import resource
  31. from libsolr.conf import SSL_CERT_CA_VERIFY
  32. LOG = logging.getLogger(__name__)
  33. try:
  34. from search.conf import EMPTY_QUERY, SECURITY_ENABLED, SOLR_URL
  35. except ImportError, e:
  36. LOG.warn('Solr Search is not enabled')
  37. def utf_quoter(what):
  38. return urllib.quote(unicode(what).encode('utf-8'), safe='~@#$&()*!+=;,.?/\'')
  39. class SolrApi(object):
  40. """
  41. http://wiki.apache.org/solr/CoreAdmin#CoreAdminHandler
  42. """
  43. def __init__(self, solr_url=None, user=None, security_enabled=False, ssl_cert_ca_verify=SSL_CERT_CA_VERIFY.get()):
  44. if solr_url is None and hasattr(SOLR_URL, 'get'):
  45. solr_url = SOLR_URL.get()
  46. if solr_url:
  47. self._url = solr_url
  48. self._user = user
  49. self._client = HttpClient(self._url, logger=LOG)
  50. self.security_enabled = security_enabled or SECURITY_ENABLED.get()
  51. if self.security_enabled:
  52. self._client.set_kerberos_auth()
  53. self._client.set_verify(ssl_cert_ca_verify)
  54. self._root = resource.Resource(self._client)
  55. # The Kerberos handshake requires two requests in order to authenticate,
  56. # but if our first request is a PUT/POST, it might flat-out reject the
  57. # first request if the body is too large. So, connect here in order to get
  58. # a cookie so future PUT/POSTs will be pre-authenticated.
  59. if self.security_enabled:
  60. self._root.invoke('HEAD', '/')
  61. def query(self, collection, query):
  62. solr_query = {}
  63. solr_query['collection'] = collection['name']
  64. if query.get('download'):
  65. solr_query['rows'] = 1000
  66. solr_query['start'] = 0
  67. else:
  68. solr_query['rows'] = int(collection['template']['rows'] or 10)
  69. solr_query['start'] = int(query['start'])
  70. solr_query['rows'] = min(solr_query['rows'], 1000)
  71. solr_query['start'] = min(solr_query['start'], 10000)
  72. params = self._get_params() + (
  73. ('q', self._get_q(query)),
  74. ('wt', 'json'),
  75. ('rows', solr_query['rows']),
  76. ('start', solr_query['start']),
  77. )
  78. if any(collection['facets']):
  79. params += (
  80. ('facet', 'true'),
  81. ('facet.mincount', 0),
  82. ('facet.limit', 10),
  83. )
  84. json_facets = {}
  85. timeFilter = self._get_range_borders(collection, query)
  86. for facet in collection['facets']:
  87. if facet['type'] == 'query':
  88. params += (('facet.query', '%s' % facet['field']),)
  89. elif facet['type'] == 'range' or facet['type'] == 'range-up':
  90. keys = {
  91. 'id': '%(id)s' % facet,
  92. 'field': facet['field'],
  93. 'key': '%(field)s-%(id)s' % facet,
  94. 'start': facet['properties']['start'],
  95. 'end': facet['properties']['end'],
  96. 'gap': facet['properties']['gap'],
  97. 'mincount': int(facet['properties']['mincount'])
  98. }
  99. if timeFilter and timeFilter['time_field'] == facet['field'] and (facet['id'] not in timeFilter['time_filter_overrides'] or facet['widgetType'] != 'histogram-widget'):
  100. keys.update(self._get_time_filter_query(timeFilter, facet))
  101. params += (
  102. ('facet.range', '{!key=%(key)s ex=%(id)s f.%(field)s.facet.range.start=%(start)s f.%(field)s.facet.range.end=%(end)s f.%(field)s.facet.range.gap=%(gap)s f.%(field)s.facet.mincount=%(mincount)s}%(field)s' % keys),
  103. )
  104. elif facet['type'] == 'field':
  105. keys = {
  106. 'id': '%(id)s' % facet,
  107. 'field': facet['field'],
  108. 'key': '%(field)s-%(id)s' % facet,
  109. 'limit': int(facet['properties'].get('limit', 10)) + (1 if facet['widgetType'] == 'facet-widget' else 0),
  110. 'mincount': int(facet['properties']['mincount'])
  111. }
  112. params += (
  113. ('facet.field', '{!key=%(key)s ex=%(id)s f.%(field)s.facet.limit=%(limit)s f.%(field)s.facet.mincount=%(mincount)s}%(field)s' % keys),
  114. )
  115. elif facet['type'] == 'nested':
  116. _f = {}
  117. if facet['properties']['facets']:
  118. self._n_facet_dimension(facet, _f, facet['properties']['facets'], 1, timeFilter)
  119. if facet['properties'].get('domain'):
  120. if facet['properties']['domain'].get('blockParent') or facet['properties']['domain'].get('blockChildren'):
  121. _f['domain'] = {}
  122. if facet['properties']['domain'].get('blockParent'):
  123. _f['domain']['blockParent'] = ' OR '.join(facet['properties']['domain']['blockParent'])
  124. if facet['properties']['domain'].get('blockChildren'):
  125. _f['domain']['blockChildren'] = ' OR '.join(facet['properties']['domain']['blockChildren'])
  126. if _f:
  127. sort = {'count': facet['properties']['facets'][0]['sort']}
  128. for i, agg in enumerate(self._get_dimension_aggregates(facet['properties']['facets'][1:])):
  129. if agg['sort'] != 'default':
  130. agg_function = self._get_aggregate_function(agg)
  131. sort = {'agg_%02d_%02d:%s' % (1, i, agg_function): agg['sort']}
  132. if sort.get('count') == 'default':
  133. sort['count'] = 'desc'
  134. dim_key = [key for key in _f['facet'].keys() if 'dim' in key][0]
  135. _f['facet'][dim_key].update({
  136. 'excludeTags': facet['id'],
  137. 'offset': 0,
  138. 'numBuckets': True,
  139. 'allBuckets': True,
  140. 'sort': sort
  141. #'prefix': '' # Forbidden on numeric fields
  142. })
  143. json_facets[facet['id']] = _f['facet'][dim_key]
  144. elif facet['type'] == 'function':
  145. if facet['properties']['facets']:
  146. json_facets[facet['id']] = self._get_aggregate_function(facet['properties']['facets'][0])
  147. if facet['properties']['compare']['is_enabled']:
  148. # TODO: global compare override
  149. unit = re.split('\d+', facet['properties']['compare']['gap'])[1]
  150. json_facets[facet['id']] = {
  151. 'type': 'range',
  152. 'field': collection['timeFilter'].get('field'),
  153. 'start': 'NOW/%s-%s-%s' % (unit, facet['properties']['compare']['gap'], facet['properties']['compare']['gap']),
  154. 'end': 'NOW/%s' % unit,
  155. 'gap': '+%(gap)s' % facet['properties']['compare'],
  156. 'facet': {facet['id']: json_facets[facet['id']]}
  157. }
  158. if facet['properties']['filter']['is_enabled']:
  159. json_facets[facet['id']] = {
  160. 'type': 'query',
  161. 'q': facet['properties']['filter']['query'] or EMPTY_QUERY.get(),
  162. 'facet': {facet['id']: json_facets[facet['id']]}
  163. }
  164. json_facets['processEmpty'] = True
  165. elif facet['type'] == 'pivot':
  166. if facet['properties']['facets'] or facet['widgetType'] == 'map-widget':
  167. fields = facet['field']
  168. fields_limits = []
  169. for f in facet['properties']['facets']:
  170. fields_limits.append('f.%s.facet.limit=%s' % (f['field'], f['limit']))
  171. fields_limits.append('f.%s.facet.mincount=%s' % (f['field'], f['mincount']))
  172. fields += ',' + f['field']
  173. keys = {
  174. 'id': '%(id)s' % facet,
  175. 'key': '%(field)s-%(id)s' % facet,
  176. 'field': facet['field'],
  177. 'fields': fields,
  178. 'limit': int(facet['properties'].get('limit', 10)),
  179. 'mincount': int(facet['properties']['mincount']),
  180. 'fields_limits': ' '.join(fields_limits)
  181. }
  182. params += (
  183. ('facet.pivot', '{!key=%(key)s ex=%(id)s f.%(field)s.facet.limit=%(limit)s f.%(field)s.facet.mincount=%(mincount)s %(fields_limits)s}%(fields)s' % keys),
  184. )
  185. if json_facets:
  186. params += (
  187. ('json.facet', json.dumps(json_facets)),
  188. )
  189. params += self._get_fq(collection, query)
  190. fl = urllib.unquote(utf_quoter(','.join(Collection2.get_field_list(collection))))
  191. nested_fields = self._get_nested_fields(collection)
  192. if nested_fields:
  193. fl += urllib.unquote(utf_quoter(',[child parentFilter="%s"]' % ' OR '.join(nested_fields)))
  194. params += (('fl', fl),)
  195. params += (
  196. ('hl', 'true'),
  197. ('hl.fl', '*'),
  198. ('hl.snippets', 5),
  199. ('hl.fragsize', 1000),
  200. )
  201. if collection['template']['fieldsSelected']:
  202. fields = []
  203. for field in collection['template']['fieldsSelected']:
  204. attribute_field = filter(lambda attribute: field == attribute['name'], collection['template']['fieldsAttributes'])
  205. if attribute_field:
  206. if attribute_field[0]['sort']['direction']:
  207. fields.append('%s %s' % (field, attribute_field[0]['sort']['direction']))
  208. if fields:
  209. params += (
  210. ('sort', ','.join(fields)),
  211. )
  212. response = self._root.get('%(collection)s/select' % solr_query, params)
  213. return self._get_json(response)
  214. def _n_facet_dimension(self, widget, _f, facets, dim, timeFilter):
  215. facet = facets[0]
  216. f_name = 'dim_%02d:%s' % (dim, facet['field'])
  217. if facet['aggregate']['function'] == 'count':
  218. if 'facet' not in _f:
  219. _f['facet'] = {f_name: {}}
  220. else:
  221. _f['facet'][f_name] = {}
  222. _f = _f['facet']
  223. sort = {'count': facet['sort']}
  224. for i, agg in enumerate(self._get_dimension_aggregates(facets)):
  225. if agg['sort'] != 'default':
  226. agg_function = self._get_aggregate_function(agg)
  227. sort = {'agg_%02d_%02d:%s' % (dim, i, agg_function): agg['sort']}
  228. if sort.get('count') == 'default':
  229. sort['count'] = 'desc'
  230. _f[f_name] = {
  231. 'type': 'terms',
  232. 'field': '%(field)s' % facet,
  233. 'limit': int(facet.get('limit', 10)),
  234. 'mincount': int(facet['mincount']),
  235. 'numBuckets': True,
  236. 'allBuckets': True,
  237. 'sort': sort,
  238. 'missing': facet.get('missing', False)
  239. #'prefix': '' # Forbidden on numeric fields
  240. }
  241. if 'start' in facet and not facet.get('type') == 'field':
  242. _f[f_name].update({
  243. 'type': 'range',
  244. 'start': facet['start'],
  245. 'end': facet['end'],
  246. 'gap': facet['gap']
  247. })
  248. # Only on dim 1 currently
  249. if timeFilter and timeFilter['time_field'] == facet['field'] and (widget['id'] not in timeFilter['time_filter_overrides']): # or facet['widgetType'] != 'bucket-widget'):
  250. facet['widgetType'] = widget['widgetType']
  251. _f[f_name].update(self._get_time_filter_query(timeFilter, facet))
  252. if widget['widgetType'] == 'tree2-widget' and facets[-1]['aggregate']['function'] != 'count':
  253. _f['subcount'] = self._get_aggregate_function(facets[-1])
  254. if len(facets) > 1: # Get n+1 dimension
  255. if facets[1]['aggregate']['function'] == 'count':
  256. self._n_facet_dimension(widget, _f[f_name], facets[1:], dim + 1, timeFilter)
  257. else:
  258. self._n_facet_dimension(widget, _f[f_name], facets[1:], dim, timeFilter)
  259. else:
  260. agg_function = self._get_aggregate_function(facet)
  261. _f['facet'] = {
  262. 'agg_%02d_00:%s' % (dim, agg_function): agg_function
  263. }
  264. for i, _f_agg in enumerate(facets[1:], 1):
  265. if _f_agg['aggregate']['function'] != 'count':
  266. agg_function = self._get_aggregate_function(_f_agg)
  267. _f['facet']['agg_%02d_%02d:%s' % (dim, i, agg_function)] = agg_function
  268. else:
  269. self._n_facet_dimension(widget, _f, facets[i:], dim + 1, timeFilter) # Get n+1 dimension
  270. break
  271. def select(self, collection, query=None, rows=100, start=0):
  272. if query is None:
  273. query = EMPTY_QUERY.get()
  274. params = self._get_params() + (
  275. ('q', query),
  276. ('wt', 'json'),
  277. ('rows', rows),
  278. ('start', start),
  279. )
  280. response = self._root.get('%s/select' % collection, params)
  281. return self._get_json(response)
  282. def suggest(self, collection, query):
  283. try:
  284. params = self._get_params() + (
  285. ('suggest', 'true'),
  286. ('suggest.build', 'true'),
  287. ('suggest.q', query['q']),
  288. ('wt', 'json'),
  289. )
  290. if query.get('dictionary'):
  291. params += (
  292. ('suggest.dictionary', query['dictionary']),
  293. )
  294. response = self._root.get('%s/suggest' % collection, params)
  295. return self._get_json(response)
  296. except RestException, e:
  297. raise PopupException(e, title=_('Error while accessing Solr'))
  298. def collections(self): # To drop, used in indexer v1
  299. try:
  300. params = self._get_params() + (
  301. ('detail', 'true'),
  302. ('path', '/clusterstate.json'),
  303. )
  304. response = self._root.get('zookeeper', params=params)
  305. return json.loads(response['znode'].get('data', '{}'))
  306. except RestException, e:
  307. raise PopupException(e, title=_('Error while accessing Solr'))
  308. def collections2(self):
  309. try:
  310. params = self._get_params() + (
  311. ('action', 'LIST'),
  312. ('wt', 'json'),
  313. )
  314. return self._root.get('admin/collections', params=params)['collections']
  315. except RestException, e:
  316. raise PopupException(e, title=_('Error while accessing Solr'))
  317. def config(self, name):
  318. try:
  319. params = self._get_params() + (
  320. ('wt', 'json'),
  321. )
  322. response = self._root.get('%s/config' % name, params=params)
  323. return self._get_json(response)['config']
  324. except RestException, e:
  325. raise PopupException(e, title=_('Error while accessing Solr'))
  326. def configs(self):
  327. try:
  328. params = self._get_params() + (
  329. ('action', 'LIST'),
  330. ('wt', 'json'),
  331. )
  332. return self._root.get('admin/configs', params=params)['configSets']
  333. except RestException, e:
  334. raise PopupException(e, title=_('Error while accessing Solr'))
  335. def create_config(self, name, base_config, immutable=False):
  336. try:
  337. params = self._get_params() + (
  338. ('action', 'CREATE'),
  339. ('name', name),
  340. ('baseConfigSet', base_config),
  341. ('configSetProp.immutable', immutable),
  342. ('wt', 'json'),
  343. )
  344. return self._root.post('admin/configs', params=params, contenttype='application/json')
  345. except RestException, e:
  346. raise PopupException(e, title=_('Error while accessing Solr'))
  347. def delete_config(self, name):
  348. response = {'status': -1, 'message': ''}
  349. try:
  350. params = self._get_params() + (
  351. ('action', 'DELETE'),
  352. ('name', name),
  353. ('wt', 'json')
  354. )
  355. data = self._root.get('admin/configs', params=params)
  356. if data['responseHeader']['status'] == 0:
  357. response['status'] = 0
  358. else:
  359. response['message'] = "Could not remove config: %s" % data
  360. except RestException, e:
  361. raise PopupException(e, title=_('Error while accessing Solr'))
  362. return response
  363. def list_aliases(self):
  364. try:
  365. params = self._get_params() + (
  366. ('action', 'LISTALIASES'),
  367. ('wt', 'json'),
  368. )
  369. return self._root.get('admin/collections', params=params)['aliases']
  370. except RestException, e:
  371. raise PopupException(e, title=_('Error while accessing Solr'))
  372. def collection_or_core(self, hue_collection):
  373. if hue_collection.is_core_only:
  374. return self.core(hue_collection.name)
  375. else:
  376. return self.collection(hue_collection.name)
  377. def collection(self, name):
  378. try:
  379. collections = self.collections()
  380. return collections[name]
  381. except Exception, e:
  382. raise PopupException(e, title=_('Error while accessing Solr'))
  383. def create_collection2(self, name, config_name=None, shards=1, replication=1, **kwargs):
  384. try:
  385. params = self._get_params() + (
  386. ('action', 'CREATE'),
  387. ('name', name),
  388. ('numShards', shards),
  389. ('replicationFactor', replication),
  390. ('wt', 'json')
  391. )
  392. if config_name:
  393. params += (
  394. ('collection.configName', config_name),
  395. )
  396. if kwargs:
  397. params += tuple(((key, val) for key, val in kwargs.iteritems()))
  398. response = self._root.post('admin/collections', params=params, contenttype='application/json')
  399. response_data = self._get_json(response)
  400. if response_data.get('failure'):
  401. raise PopupException(_('Collection could not be created: %(failure)s') % response_data)
  402. else:
  403. return response_data
  404. except RestException, e:
  405. raise PopupException(e, title=_('Error while accessing Solr'))
  406. def update_config(self, name, properties):
  407. try:
  408. params = self._get_params() + (
  409. ('wt', 'json'),
  410. )
  411. response = self._root.post('%(collection)s/config' % {'collection': name}, params=params, data=json.dumps(properties), contenttype='application/json')
  412. return self._get_json(response)
  413. except RestException, e:
  414. raise PopupException(e, title=_('Error while accessing Solr'))
  415. def add_fields(self, name, fields):
  416. try:
  417. params = self._get_params() + (
  418. ('wt', 'json'),
  419. )
  420. data = {'add-field': fields}
  421. response = self._root.post('%(collection)s/schema' % {'collection': name}, params=params, data=json.dumps(data), contenttype='application/json')
  422. return self._get_json(response)
  423. except RestException, e:
  424. raise PopupException(e, title=_('Error while accessing Solr'))
  425. def create_core(self, name, instance_dir, shards=1, replication=1):
  426. try:
  427. params = self._get_params() + (
  428. ('action', 'CREATE'),
  429. ('name', name),
  430. ('instanceDir', instance_dir),
  431. ('wt', 'json'),
  432. )
  433. response = self._root.post('admin/cores', params=params, contenttype='application/json')
  434. if response.get('responseHeader', {}).get('status', -1) == 0:
  435. return True
  436. else:
  437. LOG.error("Could not create core. Check response:\n%s" % json.dumps(response, indent=2))
  438. return False
  439. except RestException, e:
  440. if 'already exists' in e.message:
  441. LOG.warn("Could not create collection.", exc_info=True)
  442. return False
  443. else:
  444. raise PopupException(e, title=_('Error while accessing Solr'))
  445. def create_alias(self, name, collections):
  446. try:
  447. params = self._get_params() + (
  448. ('action', 'CREATEALIAS'),
  449. ('name', name),
  450. ('collections', ','.join(collections)),
  451. ('wt', 'json'),
  452. )
  453. response = self._root.post('admin/collections', params=params, contenttype='application/json')
  454. if response.get('responseHeader', {}).get('status', -1) != 0:
  455. raise PopupException(_("Could not create or edit alias: %s") % response)
  456. else:
  457. return response
  458. except RestException, e:
  459. raise PopupException(e, title=_('Error while accessing Solr'))
  460. def delete_alias(self, name):
  461. try:
  462. params = self._get_params() + (
  463. ('action', 'DELETEALIAS'),
  464. ('name', name),
  465. ('wt', 'json'),
  466. )
  467. response = self._root.post('admin/collections', params=params, contenttype='application/json')
  468. if response.get('responseHeader', {}).get('status', -1) != 0:
  469. msg = _("Could not delete alias. Check response:\n%s") % json.dumps(response, indent=2)
  470. LOG.error(msg)
  471. raise PopupException(msg)
  472. except RestException, e:
  473. raise PopupException(e, title=_('Error while accessing Solr'))
  474. def delete_collection(self, name):
  475. response = {'status': -1, 'message': ''}
  476. try:
  477. params = self._get_params() + (
  478. ('action', 'DELETE'),
  479. ('name', name),
  480. ('wt', 'json')
  481. )
  482. data = self._root.post('admin/collections', params=params, contenttype='application/json')
  483. if data['responseHeader']['status'] == 0:
  484. response['status'] = 0
  485. else:
  486. response['message'] = "Could not remove collection: %s" % data
  487. except RestException, e:
  488. raise PopupException(e, title=_('Error while accessing Solr'))
  489. return response
  490. def remove_core(self, name):
  491. try:
  492. params = self._get_params() + (
  493. ('action', 'UNLOAD'),
  494. ('name', name),
  495. ('deleteIndex', 'true'),
  496. ('wt', 'json')
  497. )
  498. response = self._root.post('admin/cores', params=params, contenttype='application/json')
  499. if 'success' in response:
  500. return True
  501. else:
  502. LOG.error("Could not remove core. Check response:\n%s" % json.dumps(response, indent=2))
  503. return False
  504. except RestException, e:
  505. raise PopupException(e, title=_('Error while accessing Solr'))
  506. def cores(self):
  507. try:
  508. params = self._get_params() + (
  509. ('wt', 'json'),
  510. )
  511. return self._root.get('admin/cores', params=params)['status']
  512. except RestException, e:
  513. raise PopupException(e, title=_('Error while accessing Solr'))
  514. def core(self, core):
  515. try:
  516. params = self._get_params() + (
  517. ('wt', 'json'),
  518. ('core', core),
  519. )
  520. return self._root.get('admin/cores', params=params)
  521. except RestException, e:
  522. raise PopupException(e, title=_('Error while accessing Solr'))
  523. def get_schema(self, collection):
  524. try:
  525. params = self._get_params() + (
  526. ('wt', 'json'),
  527. )
  528. response = self._root.get('%(core)s/schema' % {'core': collection}, params=params)
  529. return self._get_json(response)['schema']
  530. except RestException, e:
  531. raise PopupException(e, title=_('Error while accessing Solr'))
  532. # Deprecated
  533. def schema(self, core):
  534. try:
  535. params = self._get_params() + (
  536. ('wt', 'json'),
  537. ('file', 'schema.xml'),
  538. )
  539. return self._root.get('%(core)s/admin/file' % {'core': core}, params=params)
  540. except RestException, e:
  541. raise PopupException(e, title=_('Error while accessing Solr'))
  542. def fields(self, core, dynamic=False):
  543. try:
  544. params = self._get_params() + (
  545. ('wt', 'json'),
  546. ('fl', '*'),
  547. )
  548. if not dynamic:
  549. params += (('show', 'schema'),)
  550. response = self._root.get('%(core)s/admin/luke' % {'core': core}, params=params)
  551. return self._get_json(response)
  552. except RestException, e:
  553. raise PopupException(e, title=_('Error while accessing Solr'))
  554. def luke(self, core):
  555. try:
  556. params = self._get_params() + (
  557. ('wt', 'json'),
  558. )
  559. response = self._root.get('%(core)s/admin/luke' % {'core': core}, params=params)
  560. return self._get_json(response)
  561. except RestException, e:
  562. raise PopupException(e, title=_('Error while accessing Solr'))
  563. def schema_fields(self, core):
  564. try:
  565. params = self._get_params() + (
  566. ('wt', 'json'),
  567. )
  568. response = self._root.get('%(core)s/schema/fields' % {'core': core}, params=params)
  569. return self._get_json(response)
  570. except RestException, e:
  571. raise PopupException(e, title=_('Error while accessing Solr'))
  572. def stats(self, core, fields, query=None, facet=''):
  573. try:
  574. params = self._get_params() + (
  575. ('q', self._get_q(query) if query is not None else EMPTY_QUERY.get()),
  576. ('wt', 'json'),
  577. ('rows', 0),
  578. ('stats', 'true'),
  579. )
  580. if query is not None:
  581. params += self._get_fq(None, query)
  582. if facet:
  583. params += (('stats.facet', facet),)
  584. params += tuple([('stats.field', field) for field in fields])
  585. response = self._root.get('%(core)s/select' % {'core': core}, params=params)
  586. return self._get_json(response)
  587. except RestException, e:
  588. raise PopupException(e, title=_('Error while accessing Solr'))
  589. def terms(self, core, field, properties=None):
  590. try:
  591. params = self._get_params() + (
  592. ('wt', 'json'),
  593. ('rows', 0),
  594. ('terms.fl', field),
  595. )
  596. if properties:
  597. for key, val in properties.iteritems():
  598. params += ((key, val),)
  599. response = self._root.get('%(core)s/terms' % {'core': core}, params=params)
  600. return self._get_json(response)
  601. except RestException, e:
  602. raise PopupException(e, title=_('Error while accessing Solr'))
  603. def info_system(self):
  604. try:
  605. params = self._get_params() + (
  606. ('wt', 'json'),
  607. )
  608. response = self._root.get('admin/info/system', params=params)
  609. return self._get_json(response)
  610. except RestException, e:
  611. raise PopupException(e, title=_('Error while accessing Solr'))
  612. def sql(self, collection, statement):
  613. try:
  614. if 'limit' not in statement.lower(): # rows is not supported
  615. statement = statement + ' LIMIT 100'
  616. params = self._get_params() + (
  617. ('wt', 'json'),
  618. ('rows', 0),
  619. ('stmt', statement),
  620. ('rows', 100),
  621. ('start', 0),
  622. )
  623. response = self._root.get('%(collection)s/sql' % {'collection': collection}, params=params)
  624. return self._get_json(response)
  625. except RestException, e:
  626. raise PopupException(e, title=_('Error while accessing Solr'))
  627. def get(self, core, doc_id):
  628. collection_name = core['name']
  629. try:
  630. params = self._get_params() + (
  631. ('id', doc_id),
  632. ('wt', 'json'),
  633. )
  634. response = self._root.get('%(core)s/get' % {'core': collection_name}, params=params)
  635. return self._get_json(response)
  636. except RestException, e:
  637. raise PopupException(e, title=_('Error while accessing Solr'))
  638. def export(self, name, query, fl, sort, rows=100):
  639. try:
  640. params = self._get_params() + (
  641. ('q', query),
  642. ('fl', fl),
  643. ('sort', sort),
  644. ('rows', rows),
  645. ('wt', 'json'),
  646. )
  647. response = self._root.get('%(name)s/export' % {'name': name}, params=params)
  648. return self._get_json(response)
  649. except RestException, e:
  650. raise PopupException(e, title=_('Error while accessing Solr'))
  651. def update(self, collection_or_core_name, data, content_type='csv', version=None, **kwargs):
  652. if content_type == 'csv':
  653. content_type = 'application/csv'
  654. elif content_type == 'json':
  655. content_type = 'application/json'
  656. else:
  657. LOG.error("Trying to update collection %s with content type %s. Allowed content types: csv/json" % (collection_or_core_name, content_type))
  658. params = self._get_params() + (
  659. ('wt', 'json'),
  660. ('overwrite', 'true'),
  661. ('commit', 'true'),
  662. )
  663. if version is not None:
  664. params += (
  665. ('_version_', version),
  666. ('versions', 'true')
  667. )
  668. if kwargs:
  669. params += tuple(((key, val) for key, val in kwargs.iteritems()))
  670. response = self._root.post('%s/update' % collection_or_core_name, contenttype=content_type, params=params, data=data)
  671. return self._get_json(response)
  672. # Deprecated
  673. def aliases(self):
  674. try:
  675. params = self._get_params() + ( # Waiting for SOLR-4968
  676. ('detail', 'true'),
  677. ('path', '/aliases.json'),
  678. )
  679. response = self._root.get('zookeeper', params=params)
  680. return json.loads(response['znode'].get('data', '{}')).get('collection', {})
  681. except RestException, e:
  682. raise PopupException(e, title=_('Error while accessing Solr'))
  683. # Deprecated
  684. def create_collection(self, name, shards=1, replication=1):
  685. try:
  686. params = self._get_params() + (
  687. ('action', 'CREATE'),
  688. ('name', name),
  689. ('numShards', shards),
  690. ('replicationFactor', replication),
  691. ('collection.configName', name),
  692. ('wt', 'json')
  693. )
  694. response = self._root.post('admin/collections', params=params, contenttype='application/json')
  695. if 'success' in response:
  696. return True
  697. else:
  698. LOG.error("Could not create collection. Check response:\n%s" % json.dumps(response, indent=2))
  699. return False
  700. except RestException, e:
  701. raise PopupException(e, title=_('Error while accessing Solr'))
  702. # Deprecated
  703. def remove_collection(self, name):
  704. try:
  705. params = self._get_params() + (
  706. ('action', 'DELETE'),
  707. ('name', name),
  708. ('wt', 'json')
  709. )
  710. response = self._root.post('admin/collections', params=params, contenttype='application/json')
  711. if 'success' in response:
  712. return True
  713. else:
  714. LOG.error("Could not remove collection. Check response:\n%s" % json.dumps(response, indent=2))
  715. return False
  716. except RestException, e:
  717. raise PopupException(e, title=_('Error while accessing Solr'))
  718. def _get_params(self):
  719. if self.security_enabled:
  720. return (('doAs', self._user ),)
  721. return (('user.name', SERVER_USER.get()), ('doAs', self._user),)
  722. def _get_q(self, query):
  723. q_template = '(%s)' if len(query['qs']) >= 2 else '%s'
  724. return 'OR'.join([q_template % (q['q'] or EMPTY_QUERY.get()) for q in query['qs']]).encode('utf-8')
  725. @classmethod
  726. def _get_aggregate_function(cls, facet):
  727. f = facet['aggregate']
  728. if f['formula']:
  729. return f['formula']
  730. elif f['function'] == 'field':
  731. return f['value']
  732. else:
  733. fields = [facet['field']]
  734. if f['function'] == 'median':
  735. f['function'] = 'percentile'
  736. fields.append('50')
  737. elif f['function'] == 'percentile':
  738. fields.append(str(f['percentile']))
  739. f['function'] = 'percentile'
  740. return '%s(%s)' % (f['function'], ','.join(fields))
  741. def _get_range_borders(self, collection, query):
  742. props = {}
  743. time_field = collection['timeFilter'].get('field')
  744. if time_field and (collection['timeFilter']['value'] != 'all' or collection['timeFilter']['type'] == 'fixed'):
  745. # fqs overrides main time filter
  746. fq_time_ids = [fq['id'] for fq in query['fqs'] if fq['field'] == time_field]
  747. props['time_filter_overrides'] = fq_time_ids
  748. props['time_field'] = time_field
  749. if collection['timeFilter']['type'] == 'rolling':
  750. props['field'] = collection['timeFilter']['field']
  751. props['from'] = 'NOW-%s' % collection['timeFilter']['value']
  752. props['to'] = 'NOW'
  753. props['gap'] = GAPS.get(collection['timeFilter']['value'])
  754. elif collection['timeFilter']['type'] == 'fixed':
  755. props['field'] = collection['timeFilter']['field']
  756. props['from'] = collection['timeFilter'].get('from', 'NOW-7DAYS')
  757. props['to'] = collection['timeFilter'].get('to', 'NOW')
  758. props['fixed'] = True
  759. return props
  760. def _get_time_filter_query(self, timeFilter, facet):
  761. if 'fixed' in timeFilter:
  762. props = {}
  763. stat_facet = {'min': timeFilter['from'], 'max': timeFilter['to']}
  764. _compute_range_facet(facet['widgetType'], stat_facet, props, stat_facet['min'], stat_facet['max'])
  765. gap = props['gap']
  766. unit = re.split('\d+', gap)[1]
  767. return {
  768. 'start': '%(from)s/%(unit)s' % {'from': timeFilter['from'], 'unit': unit},
  769. 'end': '%(to)s/%(unit)s' % {'to': timeFilter['to'], 'unit': unit},
  770. 'gap': '%(gap)s' % props, # add a 'auto'
  771. }
  772. else:
  773. gap = timeFilter['gap'][facet['widgetType']]
  774. return {
  775. 'start': '%(from)s/%(unit)s' % {'from': timeFilter['from'], 'unit': gap['unit']},
  776. 'end': '%(to)s/%(unit)s' % {'to': timeFilter['to'], 'unit': gap['unit']},
  777. 'gap': '%(coeff)s%(unit)s/%(unit)s' % gap, # add a 'auto'
  778. }
  779. def _get_fq(self, collection, query):
  780. params = ()
  781. timeFilter = {}
  782. if collection:
  783. timeFilter = self._get_range_borders(collection, query)
  784. if timeFilter and not timeFilter.get('time_filter_overrides'):
  785. params += (('fq', urllib.unquote(utf_quoter('%(field)s:[%(from)s TO %(to)s]' % timeFilter))),)
  786. # Merge facets queries on same fields
  787. grouped_fqs = groupby(query['fqs'], lambda x: (x['type'], x['field']))
  788. merged_fqs = []
  789. for key, group in grouped_fqs:
  790. field_fq = next(group)
  791. for fq in group:
  792. for f in fq['filter']:
  793. field_fq['filter'].append(f)
  794. merged_fqs.append(field_fq)
  795. for fq in merged_fqs:
  796. if fq['type'] == 'field':
  797. fields = fq['field'] if type(fq['field']) == list else [fq['field']] # 2D facets support
  798. for field in fields:
  799. f = []
  800. for _filter in fq['filter']:
  801. values = _filter['value'] if type(_filter['value']) == list else [_filter['value']] # 2D facets support
  802. if fields.index(field) < len(values): # Lowest common field denominator
  803. value = values[fields.index(field)]
  804. if value:
  805. exclude = '-' if _filter['exclude'] else ''
  806. if value is not None and ' ' in force_unicode(value):
  807. value = force_unicode(value).replace('"', '\\"')
  808. f.append('%s%s:"%s"' % (exclude, field, value))
  809. else:
  810. f.append('%s{!field f=%s}%s' % (exclude, field, value))
  811. else: # Handle empty value selection that are returned using solr facet.missing
  812. value = "*"
  813. exclude = '-'
  814. f.append('%s%s:%s' % (exclude, field, value))
  815. _params ='{!tag=%(id)s}' % fq + ' '.join(f)
  816. params += (('fq', urllib.unquote(utf_quoter(_params))),)
  817. elif fq['type'] == 'range':
  818. params += (('fq', '{!tag=%(id)s}' % fq + ' '.join([urllib.unquote(
  819. utf_quoter('%s%s:[%s TO %s}' % ('-' if field['exclude'] else '', fq['field'], f['from'], f['to']))) for field, f in zip(fq['filter'], fq['properties'])])),)
  820. elif fq['type'] == 'range-up':
  821. params += (('fq', '{!tag=%(id)s}' % fq + ' '.join([urllib.unquote(
  822. utf_quoter('%s%s:[%s TO %s}' % ('-' if field['exclude'] else '', fq['field'], f['from'] if fq['is_up'] else '*', '*' if fq['is_up'] else f['from'])))
  823. for field, f in zip(fq['filter'], fq['properties'])])),)
  824. elif fq['type'] == 'map':
  825. _keys = fq.copy()
  826. _keys.update(fq['properties'])
  827. params += (('fq', '{!tag=%(id)s}' % fq + urllib.unquote(
  828. utf_quoter('%(lat)s:[%(lat_sw)s TO %(lat_ne)s} AND %(lon)s:[%(lon_sw)s TO %(lon_ne)s}' % _keys))),)
  829. nested_fields = self._get_nested_fields(collection)
  830. if nested_fields:
  831. params += (('fq', urllib.unquote(utf_quoter(' OR '.join(nested_fields)))),)
  832. return params
  833. def _get_dimension_aggregates(self, facets):
  834. aggregates = []
  835. for agg in facets:
  836. if agg['aggregate']['function'] != 'count':
  837. aggregates.append(agg)
  838. else:
  839. return aggregates
  840. return aggregates
  841. def _get_nested_fields(self, collection):
  842. if collection and collection.get('nested') and collection['nested']['enabled']:
  843. return [field['filter'] for field in self._flatten_schema(collection['nested']['schema']) if field['selected']]
  844. else:
  845. return []
  846. def _flatten_schema(self, level):
  847. fields = []
  848. for field in level:
  849. fields.append(field)
  850. if field['values']:
  851. fields.extend(self._flatten_schema(field['values']))
  852. return fields
  853. @classmethod
  854. def _get_json(cls, response):
  855. if type(response) != dict:
  856. # Got 'plain/text' mimetype instead of 'application/json'
  857. try:
  858. response = json.loads(response)
  859. except ValueError, e:
  860. # Got some null bytes in the response
  861. LOG.error('%s: %s' % (unicode(e), repr(response)))
  862. response = json.loads(response.replace('\x00', ''))
  863. return response
  864. def uniquekey(self, collection):
  865. try:
  866. params = self._get_params() + (
  867. ('wt', 'json'),
  868. )
  869. response = self._root.get('%s/schema/uniquekey' % collection, params=params)
  870. return self._get_json(response)['uniqueKey']
  871. except RestException, e:
  872. raise PopupException(e, title=_('Error while accessing Solr'))
  873. GAPS = {
  874. '5MINUTES': {
  875. 'histogram-widget': {'coeff': '+3', 'unit': 'SECONDS'}, # ~100 slots
  876. 'timeline-widget': {'coeff': '+3', 'unit': 'SECONDS'}, # ~100 slots
  877. 'bucket-widget': {'coeff': '+3', 'unit': 'SECONDS'}, # ~100 slots
  878. 'bar-widget': {'coeff': '+3', 'unit': 'SECONDS'}, # ~100 slots
  879. 'facet-widget': {'coeff': '+1', 'unit': 'MINUTES'}, # ~10 slots
  880. },
  881. '30MINUTES': {
  882. 'histogram-widget': {'coeff': '+20', 'unit': 'SECONDS'},
  883. 'timeline-widget': {'coeff': '+20', 'unit': 'SECONDS'},
  884. 'bucket-widget': {'coeff': '+20', 'unit': 'SECONDS'},
  885. 'bar-widget': {'coeff': '+20', 'unit': 'SECONDS'},
  886. 'facet-widget': {'coeff': '+5', 'unit': 'MINUTES'},
  887. },
  888. '1HOURS': {
  889. 'histogram-widget': {'coeff': '+30', 'unit': 'SECONDS'},
  890. 'timeline-widget': {'coeff': '+30', 'unit': 'SECONDS'},
  891. 'bucket-widget': {'coeff': '+30', 'unit': 'SECONDS'},
  892. 'bar-widget': {'coeff': '+30', 'unit': 'SECONDS'},
  893. 'facet-widget': {'coeff': '+10', 'unit': 'MINUTES'},
  894. },
  895. '12HOURS': {
  896. 'histogram-widget': {'coeff': '+7', 'unit': 'MINUTES'},
  897. 'timeline-widget': {'coeff': '+7', 'unit': 'MINUTES'},
  898. 'bucket-widget': {'coeff': '+7', 'unit': 'MINUTES'},
  899. 'bar-widget': {'coeff': '+7', 'unit': 'MINUTES'},
  900. 'facet-widget': {'coeff': '+1', 'unit': 'HOURS'},
  901. },
  902. '1DAYS': {
  903. 'histogram-widget': {'coeff': '+15', 'unit': 'MINUTES'},
  904. 'timeline-widget': {'coeff': '+15', 'unit': 'MINUTES'},
  905. 'bucket-widget': {'coeff': '+15', 'unit': 'MINUTES'},
  906. 'bar-widget': {'coeff': '+15', 'unit': 'MINUTES'},
  907. 'facet-widget': {'coeff': '+3', 'unit': 'HOURS'},
  908. },
  909. '2DAYS': {
  910. 'histogram-widget': {'coeff': '+30', 'unit': 'MINUTES'},
  911. 'timeline-widget': {'coeff': '+30', 'unit': 'MINUTES'},
  912. 'bucket-widget': {'coeff': '+30', 'unit': 'MINUTES'},
  913. 'bar-widget': {'coeff': '+30', 'unit': 'MINUTES'},
  914. 'facet-widget': {'coeff': '+6', 'unit': 'HOURS'},
  915. },
  916. '7DAYS': {
  917. 'histogram-widget': {'coeff': '+3', 'unit': 'HOURS'},
  918. 'timeline-widget': {'coeff': '+3', 'unit': 'HOURS'},
  919. 'bucket-widget': {'coeff': '+3', 'unit': 'HOURS'},
  920. 'bar-widget': {'coeff': '+3', 'unit': 'HOURS'},
  921. 'facet-widget': {'coeff': '+1', 'unit': 'DAYS'},
  922. },
  923. '1MONTHS': {
  924. 'histogram-widget': {'coeff': '+12', 'unit': 'HOURS'},
  925. 'timeline-widget': {'coeff': '+12', 'unit': 'HOURS'},
  926. 'bucket-widget': {'coeff': '+12', 'unit': 'HOURS'},
  927. 'bar-widget': {'coeff': '+12', 'unit': 'HOURS'},
  928. 'facet-widget': {'coeff': '+5', 'unit': 'DAYS'},
  929. },
  930. '3MONTHS': {
  931. 'histogram-widget': {'coeff': '+1', 'unit': 'DAYS'},
  932. 'timeline-widget': {'coeff': '+1', 'unit': 'DAYS'},
  933. 'bucket-widget': {'coeff': '+1', 'unit': 'DAYS'},
  934. 'bar-widget': {'coeff': '+1', 'unit': 'DAYS'},
  935. 'facet-widget': {'coeff': '+30', 'unit': 'DAYS'},
  936. },
  937. '1YEARS': {
  938. 'histogram-widget': {'coeff': '+3', 'unit': 'DAYS'},
  939. 'timeline-widget': {'coeff': '+3', 'unit': 'DAYS'},
  940. 'bucket-widget': {'coeff': '+3', 'unit': 'DAYS'},
  941. 'bar-widget': {'coeff': '+3', 'unit': 'DAYS'},
  942. 'facet-widget': {'coeff': '+12', 'unit': 'MONTHS'},
  943. },
  944. '2YEARS': {
  945. 'histogram-widget': {'coeff': '+7', 'unit': 'DAYS'},
  946. 'timeline-widget': {'coeff': '+7', 'unit': 'DAYS'},
  947. 'bucket-widget': {'coeff': '+7', 'unit': 'DAYS'},
  948. 'bar-widget': {'coeff': '+7', 'unit': 'DAYS'},
  949. 'facet-widget': {'coeff': '+3', 'unit': 'MONTHS'},
  950. },
  951. '10YEARS': {
  952. 'histogram-widget': {'coeff': '+1', 'unit': 'MONTHS'},
  953. 'timeline-widget': {'coeff': '+1', 'unit': 'MONTHS'},
  954. 'bucket-widget': {'coeff': '+1', 'unit': 'MONTHS'},
  955. 'bar-widget': {'coeff': '+1', 'unit': 'MONTHS'},
  956. 'facet-widget': {'coeff': '+1', 'unit': 'YEARS'},
  957. }
  958. }