Files
opensearch-pyd/elasticsearch/connection/base.py
T

186 lines
5.9 KiB
Python
Raw Normal View History

2013-08-25 16:39:56 +02:00
import logging
2019-05-10 09:16:33 -06:00
from platform import python_version
try:
import simplejson as json
except ImportError:
import json
2013-08-25 16:39:56 +02:00
from ..exceptions import TransportError, HTTP_EXCEPTIONS
from .. import __versionstr__
2013-08-25 16:39:56 +02:00
2019-05-10 09:16:33 -06:00
logger = logging.getLogger("elasticsearch")
# create the elasticsearch.trace logger, but only set propagate to False if the
# logger hasn't already been configured
2019-05-10 09:16:33 -06:00
_tracer_already_configured = "elasticsearch.trace" in logging.Logger.manager.loggerDict
tracer = logging.getLogger("elasticsearch.trace")
if not _tracer_already_configured:
tracer.propagate = False
2013-08-25 16:39:56 +02:00
class Connection(object):
"""
Class responsible for maintaining a connection to an Elasticsearch node. It
holds persistent connection pool to it and it's main interface
(`perform_request`) is thread-safe.
Also responsible for logging.
"""
2019-05-10 09:16:33 -06:00
def __init__(
self,
host="localhost",
port=9200,
use_ssl=False,
url_prefix="",
timeout=10,
**kwargs
):
2013-08-25 16:39:56 +02:00
"""
:arg host: hostname of the node (default: localhost)
2014-12-30 17:21:22 +01:00
:arg port: port to use (integer, default: 9200)
2013-08-25 16:39:56 +02:00
:arg url_prefix: optional url prefix for elasticsearch
2014-12-30 17:21:22 +01:00
:arg timeout: default timeout in seconds (float, default: 10)
2013-08-25 16:39:56 +02:00
"""
2019-05-10 09:16:33 -06:00
scheme = kwargs.get("scheme", "http")
if use_ssl or scheme == "https":
scheme = "https"
2017-05-15 21:56:49 +02:00
use_ssl = True
self.use_ssl = use_ssl
2019-05-10 09:16:33 -06:00
self.host = "%s://%s:%s" % (scheme, host, port)
2013-08-25 16:39:56 +02:00
if url_prefix:
2019-05-10 09:16:33 -06:00
url_prefix = "/" + url_prefix.strip("/")
2013-08-25 16:39:56 +02:00
self.url_prefix = url_prefix
self.timeout = timeout
2013-08-25 16:39:56 +02:00
def __repr__(self):
2019-05-10 09:16:33 -06:00
return "<%s: %s>" % (self.__class__.__name__, self.host)
2013-08-25 16:39:56 +02:00
def __eq__(self, other):
if not isinstance(other, Connection):
raise TypeError(
"Unsupported equality check for %s and %s" % (self, other)
)
return True
def __hash__(self):
return id(self)
def _pretty_json(self, data):
# pretty JSON in tracer curl logs
try:
2019-05-10 09:16:33 -06:00
return json.dumps(
json.loads(data), sort_keys=True, indent=2, separators=(",", ": ")
).replace("'", r"\u0027")
except (ValueError, TypeError):
# non-json data or a bulk request
return data
def _log_trace(self, method, path, body, status_code, response, duration):
if not tracer.isEnabledFor(logging.INFO) or not tracer.handlers:
return
# include pretty in trace curls
2019-05-10 09:16:33 -06:00
path = path.replace("?", "?pretty&", 1) if "?" in path else path + "?pretty"
if self.url_prefix:
2019-05-10 09:16:33 -06:00
path = path.replace(self.url_prefix, "", 1)
tracer.info(
"curl %s-X%s 'http://localhost:9200%s' -d '%s'",
"-H 'Content-Type: application/json' " if body else "",
method,
path,
self._pretty_json(body) if body else "",
)
if tracer.isEnabledFor(logging.DEBUG):
2019-05-10 09:16:33 -06:00
tracer.debug(
"#[%s] (%.3fs)\n#%s",
status_code,
duration,
self._pretty_json(response).replace("\n", "\n#") if response else "",
)
2019-05-10 09:16:33 -06:00
def log_request_success(
self, method, full_url, path, body, status_code, response, duration
):
2013-08-25 16:39:56 +02:00
""" Log a successful API call. """
# TODO: optionally pass in params instead of full_url and do urlencode only when needed
# body has already been serialized to utf-8, deserialize it for logging
# TODO: find a better way to avoid (de)encoding the body back and forth
if body:
2018-12-10 17:32:54 -06:00
try:
2019-05-10 09:16:33 -06:00
body = body.decode("utf-8", "ignore")
2018-12-10 17:32:54 -06:00
except AttributeError:
pass
2013-08-25 16:39:56 +02:00
logger.info(
2019-05-10 09:16:33 -06:00
"%s %s [status:%s request:%.3fs]", method, full_url, status_code, duration
2013-08-25 16:39:56 +02:00
)
2019-05-10 09:16:33 -06:00
logger.debug("> %s", body)
logger.debug("< %s", response)
2013-08-25 16:39:56 +02:00
self._log_trace(method, path, body, status_code, response, duration)
2013-08-25 16:39:56 +02:00
2019-05-10 09:16:33 -06:00
def log_request_fail(
self,
method,
full_url,
path,
body,
duration,
status_code=None,
response=None,
exception=None,
):
2013-08-25 16:39:56 +02:00
""" Log an unsuccessful API call. """
2016-04-27 13:21:47 +02:00
# do not log 404s on HEAD requests
2019-05-10 09:16:33 -06:00
if method == "HEAD" and status_code == 404:
2016-04-27 13:21:47 +02:00
return
2013-08-25 16:39:56 +02:00
logger.warning(
2019-05-10 09:16:33 -06:00
"%s %s [status:%s request:%.3fs]",
method,
full_url,
status_code or "N/A",
duration,
exc_info=exception is not None,
2013-08-25 16:39:56 +02:00
)
# body has already been serialized to utf-8, deserialize it for logging
# TODO: find a better way to avoid (de)encoding the body back and forth
if body:
2018-12-10 17:32:54 -06:00
try:
2019-05-10 09:16:33 -06:00
body = body.decode("utf-8", "ignore")
2018-12-10 17:32:54 -06:00
except AttributeError:
pass
2019-05-10 09:16:33 -06:00
logger.debug("> %s", body)
2013-08-25 16:39:56 +02:00
self._log_trace(method, path, body, status_code, response, duration)
2016-01-26 18:45:52 -08:00
if response is not None:
2019-05-10 09:16:33 -06:00
logger.debug("< %s", response)
2016-01-26 18:45:52 -08:00
2013-08-25 16:39:56 +02:00
def _raise_error(self, status_code, raw_data):
""" Locate appropriate exception and raise it. """
error_message = raw_data
additional_info = None
try:
2016-06-04 14:52:53 -07:00
if raw_data:
additional_info = json.loads(raw_data)
2019-05-10 09:16:33 -06:00
error_message = additional_info.get("error", error_message)
if isinstance(error_message, dict) and "type" in error_message:
error_message = error_message["type"]
2016-02-24 16:34:50 -05:00
except (ValueError, TypeError) as err:
2019-05-10 09:16:33 -06:00
logger.warning("Undecodable raw error response from server: %s", err)
2013-08-25 16:39:56 +02:00
2019-05-10 09:16:33 -06:00
raise HTTP_EXCEPTIONS.get(status_code, TransportError)(
status_code, error_message, additional_info
)
def _get_default_user_agent(self):
return "elasticsearch-py/%s (Python %s)" % (__versionstr__, python_version())