From 9786a3ac8033e14219af446d2cb2e5e4fee3fdda Mon Sep 17 00:00:00 2001 From: Armin Ronacher Date: Fri, 31 Aug 2018 22:17:04 +0200 Subject: [PATCH 1/8] ref: Use atexit only in init --- sentry_sdk/api.py | 6 +++++- sentry_sdk/transport.py | 2 +- 2 files changed, 6 insertions(+), 2 deletions(-) diff --git a/sentry_sdk/api.py b/sentry_sdk/api.py index a4c43f55b1..fab3487f3a 100644 --- a/sentry_sdk/api.py +++ b/sentry_sdk/api.py @@ -1,3 +1,5 @@ +import atexit + from .hub import Hub from .utils import EventHint from .client import Client, get_options @@ -27,7 +29,9 @@ def _init_on_hub(hub, args, kwargs): def init(*args, **kwargs): """Initializes the SDK and optionally integrations.""" - return _init_on_hub(Hub.main, args, kwargs) + guard = _init_on_hub(Hub.main, args, kwargs) + atexit.register(guard._client.close) + return guard def _init_on_current(*args, **kwargs): diff --git a/sentry_sdk/transport.py b/sentry_sdk/transport.py index 9c86aa3428..cc40b44267 100644 --- a/sentry_sdk/transport.py +++ b/sentry_sdk/transport.py @@ -93,7 +93,7 @@ def thread(): try: disabled_until = send_event(transport._pool, item, auth) except Exception: - print("Could not send sentry event", file=sys.stderr) + # XXX: use the logger print(traceback.format_exc(), file=sys.stderr) finally: queue.task_done() From ebb5a922ec5486af9f1281b461e10c96f12ecc51 Mon Sep 17 00:00:00 2001 From: Armin Ronacher Date: Fri, 31 Aug 2018 22:36:57 +0200 Subject: [PATCH 2/8] feat: Avoid using a demon thread --- sentry_sdk/transport.py | 31 +++++++++++++++++++++++-------- 1 file changed, 23 insertions(+), 8 deletions(-) diff --git a/sentry_sdk/transport.py b/sentry_sdk/transport.py index cc40b44267..c254d04ece 100644 --- a/sentry_sdk/transport.py +++ b/sentry_sdk/transport.py @@ -1,6 +1,8 @@ from __future__ import print_function +import atexit import json +import time import io import urllib3 import logging @@ -37,10 +39,17 @@ def _make_pool(dsn, http_proxy, https_proxy): return urllib3.PoolManager(**opts) +_PYTHON_SHUTTING_DOWN = False _SHUTDOWN = object() _retry = urllib3.util.Retry() +@atexit.register +def _learn_about_shutting_down(): + global _PYTHON_SHUTTING_DOWN + _PYTHON_SHUTTING_DOWN = True + + def send_event(pool, event, auth): body = io.BytesIO() with gzip.GzipFile(fileobj=body, mode="w") as f: @@ -76,17 +85,23 @@ def thread(): disabled_until = None # copy to local var in case transport._queue is set to None - queue = transport._queue + q = transport._queue while 1: - item = queue.get() + try: + item = q.get(timeout=0.1) + except queue.Empty: + if _SHUTDOWN: + break + continue + if item is _SHUTDOWN: - queue.task_done() + q.task_done() break if disabled_until is not None: if datetime.utcnow() < disabled_until: - queue.task_done() + q.task_done() continue disabled_until = None @@ -96,10 +111,9 @@ def thread(): # XXX: use the logger print(traceback.format_exc(), file=sys.stderr) finally: - queue.task_done() + q.task_done() t = threading.Thread(target=thread) - t.setDaemon(True) t.start() @@ -133,9 +147,10 @@ def close(self): def drain_events(self, timeout): q = self._queue if q is not None: + started = time.time() with q.all_tasks_done: - while q.unfinished_tasks: - q.all_tasks_done.wait(timeout) + while q.unfinished_tasks and (time.time() - started) < timeout: + q.all_tasks_done.wait(0.1) def __del__(self): self.close() From 6d40fa72273b40acf1d6dff4c2c7dbef01a03d03 Mon Sep 17 00:00:00 2001 From: Armin Ronacher Date: Fri, 31 Aug 2018 22:39:17 +0200 Subject: [PATCH 3/8] feat: Use logger for transport --- sentry_sdk/transport.py | 7 +++---- 1 file changed, 3 insertions(+), 4 deletions(-) diff --git a/sentry_sdk/transport.py b/sentry_sdk/transport.py index c254d04ece..2dd5b09e99 100644 --- a/sentry_sdk/transport.py +++ b/sentry_sdk/transport.py @@ -9,13 +9,13 @@ import threading import certifi import sys -import traceback import gzip from datetime import datetime, timedelta from ._compat import queue from .consts import VERSION +from .utils import get_logger try: from urllib.request import getproxies @@ -23,7 +23,7 @@ from urllib import getproxies -logger = logging.getLogger(__name__) +logger = get_logger(__name__) def _make_pool(dsn, http_proxy, https_proxy): @@ -108,8 +108,7 @@ def thread(): try: disabled_until = send_event(transport._pool, item, auth) except Exception: - # XXX: use the logger - print(traceback.format_exc(), file=sys.stderr) + logger.error('Could not send event', exc_info=sys.exc_info()) finally: q.task_done() From 2881ffce6b11f663b41c79ed0e369368c054fb02 Mon Sep 17 00:00:00 2001 From: Armin Ronacher Date: Sat, 1 Sep 2018 01:53:03 +0200 Subject: [PATCH 4/8] feat: Port old background worker code over --- sentry_sdk/_compat.py | 24 +++ sentry_sdk/api.py | 6 +- sentry_sdk/client.py | 20 +- sentry_sdk/consts.py | 9 + sentry_sdk/hub.py | 3 +- sentry_sdk/transport.py | 172 +++++++----------- sentry_sdk/worker.py | 99 ++++++++++ tests/conftest.py | 6 +- .../excepthook/test_excepthook.py | 5 +- tests/test_client.py | 6 +- 10 files changed, 210 insertions(+), 140 deletions(-) create mode 100644 sentry_sdk/worker.py diff --git a/sentry_sdk/_compat.py b/sentry_sdk/_compat.py index b796adf608..7e9986a08d 100644 --- a/sentry_sdk/_compat.py +++ b/sentry_sdk/_compat.py @@ -53,3 +53,27 @@ def __new__(cls, name, this_bases, d): return meta(name, bases, d) return type.__new__(metaclass, "temporary_class", (), {}) + + +def check_thread_support(): + try: + from uwsgi import opt + except ImportError: + return + + # When `threads` is passed in as a uwsgi option, + # `enable-threads` is implied on. + if "threads" in opt: + return + + if str(opt.get("enable-threads", "0")).lower() in ("false", "off", "no", "0"): + from warnings import warn + + warn( + Warning( + "We detected the use of uwsgi with disabled threads. " + "This will cause issues with the transport you are " + "trying to use. Please enable threading for uwsgi. " + '(Enable the "enable-threads" flag).' + ) + ) diff --git a/sentry_sdk/api.py b/sentry_sdk/api.py index fab3487f3a..a4c43f55b1 100644 --- a/sentry_sdk/api.py +++ b/sentry_sdk/api.py @@ -1,5 +1,3 @@ -import atexit - from .hub import Hub from .utils import EventHint from .client import Client, get_options @@ -29,9 +27,7 @@ def _init_on_hub(hub, args, kwargs): def init(*args, **kwargs): """Initializes the SDK and optionally integrations.""" - guard = _init_on_hub(Hub.main, args, kwargs) - atexit.register(guard._client.close) - return guard + return _init_on_hub(Hub.main, args, kwargs) def _init_on_current(*args, **kwargs): diff --git a/sentry_sdk/client.py b/sentry_sdk/client.py index 51ba02667a..93e960e258 100644 --- a/sentry_sdk/client.py +++ b/sentry_sdk/client.py @@ -42,7 +42,7 @@ def get_options(*args, **kwargs): class Client(object): def __init__(self, *args, **kwargs): self.options = options = get_options(*args, **kwargs) - self._transport = make_transport(options) + self.transport = make_transport(options) request_bodies = ("always", "never", "small", "medium") if options["request_bodies"] not in request_bodies: @@ -52,9 +52,6 @@ def __init__(self, *args, **kwargs): ) ) - # XXX: we should probably only do this for the init()ed client - atexit.register(self.close) - @property def dsn(self): """Returns the configured dsn.""" @@ -129,7 +126,7 @@ def _should_capture(self, event, hint=None, scope=None): def capture_event(self, event, hint=None, scope=None): """Captures an event.""" - if self._transport is None: + if self.transport is None: return rv = event.get("event_id") if rv is None: @@ -137,16 +134,5 @@ def capture_event(self, event, hint=None, scope=None): if self._should_capture(event, hint, scope): event = self._prepare_event(event, hint, scope) if event is not None: - self._transport.capture_event(event) + self.transport.capture_event(event) return rv - - def drain_events(self, timeout=None): - if timeout is None: - timeout = self.options["shutdown_timeout"] - if self._transport is not None: - self._transport.drain_events(timeout) - - def close(self): - self.drain_events() - if self._transport is not None: - self._transport.close() diff --git a/sentry_sdk/consts.py b/sentry_sdk/consts.py index 8e4bf4981f..38b6d815a2 100644 --- a/sentry_sdk/consts.py +++ b/sentry_sdk/consts.py @@ -1,5 +1,13 @@ +import os import socket + +def default_shutdown_callback(pending, timeout): + print("Sentry is attempting to send %i pending error messages" % pending) + print("Waiting up to %s seconds" % timeout) + print("Press Ctrl-%s to quit" % (os.name == "nt" and "Break" or "C")) + + VERSION = "0.1" DEFAULT_SERVER_NAME = socket.gethostname() if hasattr(socket, "gethostname") else None DEFAULT_OPTIONS = { @@ -10,6 +18,7 @@ "environment": None, "server_name": DEFAULT_SERVER_NAME, "shutdown_timeout": 2.0, + "shutdown_callback": default_shutdown_callback, "integrations": [], "in_app_include": [], "in_app_exclude": [], diff --git a/sentry_sdk/hub.py b/sentry_sdk/hub.py index 8eba749f07..37711e4c30 100644 --- a/sentry_sdk/hub.py +++ b/sentry_sdk/hub.py @@ -9,8 +9,7 @@ _local = ContextVar("sentry_current_hub") - -logger = get_logger(__name__) +logger = get_logger("sentry_sdk.errors") @contextmanager diff --git a/sentry_sdk/transport.py b/sentry_sdk/transport.py index 7a5360fef5..5ea601051e 100644 --- a/sentry_sdk/transport.py +++ b/sentry_sdk/transport.py @@ -13,9 +13,10 @@ from datetime import datetime, timedelta -from ._compat import queue from .consts import VERSION -from .utils import get_logger, Dsn +from .utils import Dsn +from .worker import BackgroundWorker +from .hub import _internal_exceptions try: from urllib.request import getproxies @@ -23,9 +24,6 @@ from urllib import getproxies -logger = get_logger(__name__) - - def _make_pool(parsed_dsn, http_proxy, https_proxy): proxy = https_proxy if parsed_dsn == "https" else http_proxy if not proxy: @@ -39,81 +37,15 @@ def _make_pool(parsed_dsn, http_proxy, https_proxy): return urllib3.PoolManager(**opts) -_PYTHON_SHUTTING_DOWN = False -_SHUTDOWN = object() -_retry = urllib3.util.Retry() - - @atexit.register -def _learn_about_shutting_down(): - global _PYTHON_SHUTTING_DOWN - _PYTHON_SHUTTING_DOWN = True - - -def send_event(pool, event, auth): - body = io.BytesIO() - with gzip.GzipFile(fileobj=body, mode="w") as f: - f.write(json.dumps(event).encode("utf-8")) - - response = pool.request( - "POST", - str(auth.store_api_url), - body=body.getvalue(), - headers={ - "X-Sentry-Auth": str(auth.to_header()), - "Content-Type": "application/json", - "Content-Encoding": "gzip", - }, - ) - - try: - if response.status == 429: - return datetime.utcnow() + timedelta( - seconds=_retry.get_retry_after(response) - ) - - if response.status >= 300 or response.status < 200: - raise ValueError("Unexpected status code: %s" % response.status) - finally: - response.close() - - -def spawn_thread(transport): - auth = transport.parsed_dsn.to_auth("sentry-python/%s" % VERSION) +def _shutdown(): + from .hub import Hub - def thread(): - disabled_until = None - - # copy to local var in case transport._queue is set to None - q = transport._queue - - while 1: - try: - item = q.get(timeout=0.1) - except queue.Empty: - if _SHUTDOWN: - break - continue - - if item is _SHUTDOWN: - q.task_done() - break - - if disabled_until is not None: - if datetime.utcnow() < disabled_until: - q.task_done() - continue - disabled_until = None - - try: - disabled_until = send_event(transport._pool, item, auth) - except Exception: - logger.error('Could not send event', exc_info=sys.exc_info()) - finally: - q.task_done() - - t = threading.Thread(target=thread) - t.start() + main_client = Hub.main.client + if main_client is not None: + transport = main_client.transport + if transport is not None: + transport.wait_and_close() class Transport(object): @@ -127,10 +59,10 @@ def __init__(self, options=None): def capture_event(self, event): raise NotImplementedError() - def close(self): - pass + def wait_and_close(self): + self.close() - def drain_events(self, timeout): + def close(self): pass def __del__(self): @@ -140,45 +72,69 @@ def __del__(self): class HttpTransport(Transport): def __init__(self, options): Transport.__init__(self, options) - self._queue = None + self._worker = BackgroundWorker( + shutdown_timeout=options["shutdown_timeout"], + shutdown_callback=options["shutdown_callback"], + ) + self._auth = self.parsed_dsn.to_auth("sentry-python/%s" % VERSION) self._pool = _make_pool( self.parsed_dsn, http_proxy=options["http_proxy"], https_proxy=options["https_proxy"], ) - self.start() + self._disabled_until = None + self._retry = urllib3.util.Retry() + + def _send_event(self, event): + if self._disabled_until is not None: + if datetime.utcnow() < self._disabled_until: + return + self._disabled_until = None + + with _internal_exceptions(): + body = io.BytesIO() + with gzip.GzipFile(fileobj=body, mode="w") as f: + f.write(json.dumps(event).encode("utf-8")) + + response = self._pool.request( + "POST", + str(self._auth.store_api_url), + body=body.getvalue(), + headers={ + "X-Sentry-Auth": str(self._auth.to_header()), + "Content-Type": "application/json", + "Content-Encoding": "gzip", + }, + ) + + try: + if response.status == 429: + self._disabled_until = datetime.utcnow() + timedelta( + seconds=_retry.get_retry_after(response) + ) + return + + elif response.status >= 300 or response.status < 200: + raise ValueError("Unexpected status code: %s" % response.status) + finally: + response.close() - def start(self): - if self._queue is None: - self._queue = queue.Queue(30) - spawn_thread(self) + self._disabled_until = None def capture_event(self, event): - if self._queue is None: - raise RuntimeError("Transport shut down") - try: - self._queue.put_nowait(event) - except queue.Full: - pass + self._worker.submit(lambda: self._send_event(event)) + + def wait_and_close(self): + self._worker.shutdown() def close(self): - if self._queue is not None: - try: - self._queue.put_nowait(_SHUTDOWN) - except queue.Full: - pass - self._queue = None - - def drain_events(self, timeout): - q = self._queue - if q is not None: - started = time.time() - with q.all_tasks_done: - while q.unfinished_tasks and (time.time() - started) < timeout: - q.all_tasks_done.wait(0.1) + self._worker.stop() def __del__(self): - self.close() + try: + self.close() + except Exception: + pass class _FunctionTransport(Transport): diff --git a/sentry_sdk/worker.py b/sentry_sdk/worker.py new file mode 100644 index 0000000000..cc4846cf35 --- /dev/null +++ b/sentry_sdk/worker.py @@ -0,0 +1,99 @@ +import logging +import threading +import os + +from time import sleep, time +from ._compat import queue, check_thread_support +from .utils import get_logger + + +_TERMINATOR = object() +logger = get_logger("sentry_sdk.errors") + + +class BackgroundWorker(object): + def __init__( + self, shutdown_timeout=10, initial_timeout=0.2, shutdown_callback=None + ): + check_thread_support() + self._queue = queue.Queue(-1) + self._lock = threading.Lock() + self._thread = None + self._thread_for_pid = None + self.initial_timeout = initial_timeout + self.shutdown_timeout = shutdown_timeout + self.shutdown_callback = shutdown_callback + + @property + def is_alive(self): + if self._thread_for_pid != os.getpid(): + return False + return self._thread and self._thread.is_alive() + + def _ensure_thread(self): + if not self.is_alive: + self.start() + + def _timed_queue_join(self, timeout): + deadline = time() + timeout + queue = self._queue + queue.all_tasks_done.acquire() + try: + while queue.unfinished_tasks: + delay = deadline - time() + if delay <= 0: + return False + queue.all_tasks_done.wait(timeout=delay) + return True + finally: + queue.all_tasks_done.release() + + def start(self): + with self._lock: + if not self.is_alive: + self._thread = threading.Thread( + target=self._target, name="raven-sentry.BackgroundWorker" + ) + self._thread.setDaemon(True) + self._thread.start() + self._thread_for_pid = os.getpid() + + def stop(self, timeout=None): + with self._lock: + if self._thread: + self._queue.put_nowait(_TERMINATOR) + if timeout is not None: + self._thread.join(timeout=timeout) + self._thread = None + self._thread_for_pid = None + + def shutdown(self): + with self._lock: + if not self.is_alive: + return + self._queue.put_nowait(_TERMINATOR) + timeout = self.shutdown_timeout + initial_timeout = min(self.initial_timeout, timeout) + if not self._timed_queue_join(initial_timeout): + if self.shutdown_callback is not None: + self.shutdown_callback(self._queue.qsize(), timeout) + self._timed_queue_join(timeout - initial_timeout) + self._thread = None + + def submit(self, callback): + self._ensure_thread() + self._queue.put_nowait(callback) + + def _target(self): + while True: + callback = self._queue.get() + try: + if callback is _TERMINATOR: + break + try: + callback() + except Exception: + logger.error("Failed processing job", exc_info=True) + finally: + self._queue.task_done() + sleep(0) diff --git a/tests/conftest.py b/tests/conftest.py index ca12b8c78c..b632dd1de4 100644 --- a/tests/conftest.py +++ b/tests/conftest.py @@ -33,7 +33,7 @@ def capture_internal_exception(exc_info): def monkeypatch_test_transport(monkeypatch, assert_semaphore_acceptance): def inner(client): monkeypatch.setattr( - client, "_transport", TestTransport(assert_semaphore_acceptance) + client, "transport", TestTransport(assert_semaphore_acceptance) ) return inner @@ -106,13 +106,13 @@ def capture_events(monkeypatch): def inner(): events = [] test_client = sentry_sdk.Hub.current.client - old_capture_event = test_client._transport.capture_event + old_capture_event = test_client.transport.capture_event def append(event): events.append(event) return old_capture_event(event) - monkeypatch.setattr(test_client._transport, "capture_event", append) + monkeypatch.setattr(test_client.transport, "capture_event", append) return events return inner diff --git a/tests/integrations/excepthook/test_excepthook.py b/tests/integrations/excepthook/test_excepthook.py index bf201418cb..3a376e434b 100644 --- a/tests/integrations/excepthook/test_excepthook.py +++ b/tests/integrations/excepthook/test_excepthook.py @@ -12,11 +12,11 @@ def test_excepthook(tmpdir): """ from sentry_sdk import init, transport - def send_event(pool, event, auth): + def send_event(self, event): print("capture event was called") print(event) - transport.send_event = send_event + transport.HttpTransport._send_event = send_event init("http://foobar@localhost/123") @@ -31,6 +31,7 @@ def send_event(pool, event, auth): subprocess.check_output([sys.executable, str(app)], stderr=subprocess.STDOUT) output = excinfo.value.output + print(output) assert b"ZeroDivisionError" in output assert b"LOL" in output diff --git a/tests/test_client.py b/tests/test_client.py index 5866b26705..22da5f3163 100644 --- a/tests/test_client.py +++ b/tests/test_client.py @@ -75,11 +75,11 @@ def test_atexit(tmpdir, monkeypatch, num_messages): import time from sentry_sdk import init, transport, capture_message - def send_event(pool, event, auth): + def send_event(self, event): time.sleep(0.1) print(event["message"]) - transport.send_event = send_event + transport.HttpTransport._send_event = send_event init("http://foobar@localhost/123", shutdown_timeout={num_messages}) for _ in range({num_messages}): @@ -148,7 +148,7 @@ def test_transport_works(httpserver, request, capsys): add_breadcrumb(level="info", message="i like bread", timestamp=datetime.now()) capture_message("löl") - client.drain_events() + client.transport.wait_and_close() out, err = capsys.readouterr() assert not err and not out From 9458049b0086ae9829ce0389f0bf6aa019887311 Mon Sep 17 00:00:00 2001 From: Armin Ronacher Date: Sat, 1 Sep 2018 02:01:12 +0200 Subject: [PATCH 5/8] ref: Added close method to client which shuts down orderly. --- sentry_sdk/client.py | 13 +++++++++++++ sentry_sdk/transport.py | 25 +++++++++++-------------- sentry_sdk/worker.py | 4 +--- tests/conftest.py | 9 +-------- tests/test_client.py | 2 +- 5 files changed, 27 insertions(+), 26 deletions(-) diff --git a/sentry_sdk/client.py b/sentry_sdk/client.py index 93e960e258..639aacb39b 100644 --- a/sentry_sdk/client.py +++ b/sentry_sdk/client.py @@ -136,3 +136,16 @@ def capture_event(self, event, hint=None, scope=None): if event is not None: self.transport.capture_event(event) return rv + + def close(self): + """Closes the client which shuts down the transport in an + orderly manner. + """ + if self.transport is not None: + self.transport.shutdown() + + def __enter__(self): + return self + + def __exit__(self, exc_type, exc_value, tb): + self.close() diff --git a/sentry_sdk/transport.py b/sentry_sdk/transport.py index 5ea601051e..9caeee202b 100644 --- a/sentry_sdk/transport.py +++ b/sentry_sdk/transport.py @@ -45,7 +45,7 @@ def _shutdown(): if main_client is not None: transport = main_client.transport if transport is not None: - transport.wait_and_close() + transport.shutdown() class Transport(object): @@ -59,14 +59,17 @@ def __init__(self, options=None): def capture_event(self, event): raise NotImplementedError() - def wait_and_close(self): - self.close() + def shutdown(self): + self.kill() - def close(self): + def kill(self): pass def __del__(self): - self.close() + try: + self.kill() + except Exception: + pass class HttpTransport(Transport): @@ -124,17 +127,11 @@ def _send_event(self, event): def capture_event(self, event): self._worker.submit(lambda: self._send_event(event)) - def wait_and_close(self): + def shutdown(self): self._worker.shutdown() - def close(self): - self._worker.stop() - - def __del__(self): - try: - self.close() - except Exception: - pass + def kill(self): + self._worker.kill() class _FunctionTransport(Transport): diff --git a/sentry_sdk/worker.py b/sentry_sdk/worker.py index cc4846cf35..bc7c531a8c 100644 --- a/sentry_sdk/worker.py +++ b/sentry_sdk/worker.py @@ -58,12 +58,10 @@ def start(self): self._thread.start() self._thread_for_pid = os.getpid() - def stop(self, timeout=None): + def kill(self): with self._lock: if self._thread: self._queue.put_nowait(_TERMINATOR) - if timeout is not None: - self._thread.join(timeout=timeout) self._thread = None self._thread_for_pid = None diff --git a/tests/conftest.py b/tests/conftest.py index b632dd1de4..c3cf483770 100644 --- a/tests/conftest.py +++ b/tests/conftest.py @@ -89,17 +89,10 @@ def inner(*a, **kw): class TestTransport(Transport): def __init__(self, capture_event_callback): + Transport.__init__(self) self.capture_event = capture_event_callback self._queue = None - def start(self): - pass - - def close(self): - pass - - dsn = "LOL" - @pytest.fixture def capture_events(monkeypatch): diff --git a/tests/test_client.py b/tests/test_client.py index 22da5f3163..010b982368 100644 --- a/tests/test_client.py +++ b/tests/test_client.py @@ -148,7 +148,7 @@ def test_transport_works(httpserver, request, capsys): add_breadcrumb(level="info", message="i like bread", timestamp=datetime.now()) capture_message("löl") - client.transport.wait_and_close() + client.close() out, err = capsys.readouterr() assert not err and not out From 4fc4fadd9f48fb57773b87d92eb5ae99a50ad05f Mon Sep 17 00:00:00 2001 From: Armin Ronacher Date: Sat, 1 Sep 2018 14:51:23 +0200 Subject: [PATCH 6/8] feat: Add more debug logs --- sentry_sdk/client.py | 9 ++++++++- sentry_sdk/hub.py | 4 ++++ sentry_sdk/scope.py | 23 +++++++++++++++++------ sentry_sdk/transport.py | 20 ++++++++++++++++---- sentry_sdk/utils.py | 4 +++- sentry_sdk/worker.py | 7 ++++++- 6 files changed, 54 insertions(+), 13 deletions(-) diff --git a/sentry_sdk/client.py b/sentry_sdk/client.py index 639aacb39b..320c7e539e 100644 --- a/sentry_sdk/client.py +++ b/sentry_sdk/client.py @@ -11,11 +11,15 @@ convert_types, handle_in_app, get_type_name, + get_logger, ) from .transport import make_transport from .consts import DEFAULT_OPTIONS, SDK_INFO +logger = get_logger("sentry_sdk.errors") + + def get_options(*args, **kwargs): if args and (isinstance(args[0], string_types) or args[0] is None): dsn = args[0] @@ -81,7 +85,10 @@ def _prepare_event(self, event, hint, scope): before_send = self.options["before_send"] if before_send is not None: - event = before_send(event) + new_event = before_send(event) + if new_event is None: + logger.info("before send dropped event (%s)", event) + event = new_event # Postprocess the event in the very end so that annotated types do # generally not surface in before_send diff --git a/sentry_sdk/hub.py b/sentry_sdk/hub.py index 37711e4c30..caeb4bd733 100644 --- a/sentry_sdk/hub.py +++ b/sentry_sdk/hub.py @@ -159,6 +159,7 @@ def add_breadcrumb(self, *args, **kwargs): """Adds a breadcrumb.""" client, scope = self._stack[-1] if client is None: + logger.info("Dropped breadcrumb because no client bound") return if not kwargs and len(args) == 1 and callable(args[0]): @@ -173,11 +174,14 @@ def add_breadcrumb(self, *args, **kwargs): if crumb.get("type") is None: crumb["type"] = "default" + original_crumb = crumb if client.options["before_breadcrumb"] is not None: crumb = client.options["before_breadcrumb"](crumb) if crumb is not None: scope._breadcrumbs.append(crumb) + else: + logger.info("before breadcrumb dropped breadcrumb (%s)", original_crumb) while len(scope._breadcrumbs) >= client.options["max_breadcrumbs"]: scope._breadcrumbs.popleft() diff --git a/sentry_sdk/scope.py b/sentry_sdk/scope.py index f8c20be519..45a0719fd8 100644 --- a/sentry_sdk/scope.py +++ b/sentry_sdk/scope.py @@ -1,3 +1,9 @@ +from .utils import get_logger + + +logger = get_logger("sentry_sdk.errors") + + class Scope(object): __slots__ = ["_data", "_breadcrumbs", "_event_processors", "_error_processors"] @@ -73,6 +79,9 @@ def func(event, exc_info): self._error_processors.append(func) def apply_to_event(self, event, hint=None): + def _drop(event, cause, ty): + logger.info("%s (%s) dropped event (%s)", ty, cause, event) + event.setdefault("breadcrumbs", []).extend(self._breadcrumbs) if event.get("user") is None and "user" in self._data: event["user"] = self._data["user"] @@ -99,14 +108,16 @@ def apply_to_event(self, event, hint=None): if hint is not None and hint.exc_info is not None: exc_info = hint.exc_info for processor in self._error_processors: - event = processor(event, exc_info) - if event is None: - return + new_event = processor(event, exc_info) + if new_event is None: + return _drop(event, processor, "error processor") + event = new_event for processor in self._event_processors: - event = processor(event) - if event is None: - return None + new_event = processor(event) + if new_event is None: + return _drop(event, processor, "event processor") + event = new_event return event diff --git a/sentry_sdk/transport.py b/sentry_sdk/transport.py index 9caeee202b..3facd7d87d 100644 --- a/sentry_sdk/transport.py +++ b/sentry_sdk/transport.py @@ -14,7 +14,7 @@ from datetime import datetime, timedelta from .consts import VERSION -from .utils import Dsn +from .utils import Dsn, get_logger from .worker import BackgroundWorker from .hub import _internal_exceptions @@ -24,6 +24,9 @@ from urllib import getproxies +logger = get_logger("sentry_sdk.errors") + + def _make_pool(parsed_dsn, http_proxy, https_proxy): proxy = https_proxy if parsed_dsn == "https" else http_proxy if not proxy: @@ -43,9 +46,7 @@ def _shutdown(): main_client = Hub.main.client if main_client is not None: - transport = main_client.transport - if transport is not None: - transport.shutdown() + main_client.close() class Transport(object): @@ -99,6 +100,15 @@ def _send_event(self, event): with gzip.GzipFile(fileobj=body, mode="w") as f: f.write(json.dumps(event).encode("utf-8")) + logger.debug( + "Sending %s event [%s] to %s project:%s" + % ( + event.get("level") or "error", + event["event_id"], + self.parsed_dsn.host, + self.parsed_dsn.project_id, + ) + ) response = self._pool.request( "POST", str(self._auth.store_api_url), @@ -128,9 +138,11 @@ def capture_event(self, event): self._worker.submit(lambda: self._send_event(event)) def shutdown(self): + logger.debug("Shutting down HTTP transport orderly") self._worker.shutdown() def kill(self): + logger.debug("Killing HTTP transport") self._worker.kill() diff --git a/sentry_sdk/utils.py b/sentry_sdk/utils.py index 87eefd1328..3214a54e86 100644 --- a/sentry_sdk/utils.py +++ b/sentry_sdk/utils.py @@ -527,7 +527,9 @@ def strip_string(value, assume_length=None, max_length=512): def get_logger(name): rv = logging.getLogger(name) if not rv.handlers: - rv.addHandler(logging.StreamHandler(sys.stderr)) + handler = logging.StreamHandler(sys.stderr) + handler.setFormatter(logging.Formatter(" [sentry] %(levelname)s: %(message)s")) + rv.addHandler(handler) rv.setLevel(logging.DEBUG) return rv diff --git a/sentry_sdk/worker.py b/sentry_sdk/worker.py index bc7c531a8c..a374b3c381 100644 --- a/sentry_sdk/worker.py +++ b/sentry_sdk/worker.py @@ -59,6 +59,7 @@ def start(self): self._thread_for_pid = os.getpid() def kill(self): + logger.debug("Transport got kill request") with self._lock: if self._thread: self._queue.put_nowait(_TERMINATOR) @@ -66,6 +67,7 @@ def kill(self): self._thread_for_pid = None def shutdown(self): + logger.debug("Transport got shutdown request") with self._lock: if not self.is_alive: return @@ -73,10 +75,13 @@ def shutdown(self): timeout = self.shutdown_timeout initial_timeout = min(self.initial_timeout, timeout) if not self._timed_queue_join(initial_timeout): + pending = self._queue.qsize() + logger.debug("%d event(s) pending on shutdown", pending) if self.shutdown_callback is not None: - self.shutdown_callback(self._queue.qsize(), timeout) + self.shutdown_callback(pending, timeout) self._timed_queue_join(timeout - initial_timeout) self._thread = None + logger.debug("Transport shut down") def submit(self, callback): self._ensure_thread() From 0b65532ddc000f5927d897798f18b5e1eddcb8ef Mon Sep 17 00:00:00 2001 From: Armin Ronacher Date: Sat, 1 Sep 2018 17:09:24 +0200 Subject: [PATCH 7/8] ref: Moved atexit handling into an integration --- sentry_sdk/client.py | 6 ++++-- sentry_sdk/consts.py | 7 ------- sentry_sdk/integrations/__init__.py | 2 ++ sentry_sdk/integrations/atexit.py | 30 +++++++++++++++++++++++++++++ sentry_sdk/transport.py | 20 ++++--------------- sentry_sdk/utils.py | 2 +- sentry_sdk/worker.py | 22 ++++++++------------- 7 files changed, 49 insertions(+), 40 deletions(-) create mode 100644 sentry_sdk/integrations/atexit.py diff --git a/sentry_sdk/client.py b/sentry_sdk/client.py index 11dbbd78eb..3f3c28d61f 100644 --- a/sentry_sdk/client.py +++ b/sentry_sdk/client.py @@ -140,12 +140,14 @@ def capture_event(self, event, hint=None, scope=None): self.transport.capture_event(event) return rv - def close(self): + def close(self, timeout=None, shutdown_callback=None): """Closes the client which shuts down the transport in an orderly manner. """ if self.transport is not None: - self.transport.shutdown() + if timeout is None: + timeout = self.options["shutdown_timeout"] + self.transport.shutdown(timeout=timeout, callback=shutdown_callback) def __enter__(self): return self diff --git a/sentry_sdk/consts.py b/sentry_sdk/consts.py index 38b6d815a2..d0ffb7d93b 100644 --- a/sentry_sdk/consts.py +++ b/sentry_sdk/consts.py @@ -2,12 +2,6 @@ import socket -def default_shutdown_callback(pending, timeout): - print("Sentry is attempting to send %i pending error messages" % pending) - print("Waiting up to %s seconds" % timeout) - print("Press Ctrl-%s to quit" % (os.name == "nt" and "Break" or "C")) - - VERSION = "0.1" DEFAULT_SERVER_NAME = socket.gethostname() if hasattr(socket, "gethostname") else None DEFAULT_OPTIONS = { @@ -18,7 +12,6 @@ def default_shutdown_callback(pending, timeout): "environment": None, "server_name": DEFAULT_SERVER_NAME, "shutdown_timeout": 2.0, - "shutdown_callback": default_shutdown_callback, "integrations": [], "in_app_include": [], "in_app_exclude": [], diff --git a/sentry_sdk/integrations/__init__.py b/sentry_sdk/integrations/__init__.py index 2d6526daa4..b895688468 100644 --- a/sentry_sdk/integrations/__init__.py +++ b/sentry_sdk/integrations/__init__.py @@ -11,10 +11,12 @@ def _get_default_integrations(): from .logging import LoggingIntegration from .excepthook import ExcepthookIntegration from .dedupe import DedupeIntegration + from .atexit import AtexitIntegration yield LoggingIntegration yield ExcepthookIntegration yield DedupeIntegration + yield AtexitIntegration def setup_integrations(options): diff --git a/sentry_sdk/integrations/atexit.py b/sentry_sdk/integrations/atexit.py new file mode 100644 index 0000000000..17ba68d25c --- /dev/null +++ b/sentry_sdk/integrations/atexit.py @@ -0,0 +1,30 @@ +import os +import atexit + +from sentry_sdk.hub import Hub +from sentry_sdk.utils import logger +from . import Integration + + +def default_shutdown_callback(pending, timeout): + print("Sentry is attempting to send %i pending error messages" % pending) + print("Waiting up to %s seconds" % timeout) + print("Press Ctrl-%s to quit" % (os.name == "nt" and "Break" or "C")) + + +class AtexitIntegration(Integration): + identifier = "atexit" + + def __init__(self, callback=None): + if callback is None: + callback = default_shutdown_callback + self.callback = callback + + def install(self): + @atexit.register + def _shutdown(): + main_client = Hub.main.client + logger.debug("atexit: got shutdown signal") + if main_client is not None: + logger.debug("atexit: shutting down client") + main_client.close(shutdown_callback=self.callback) diff --git a/sentry_sdk/transport.py b/sentry_sdk/transport.py index 4f37663cff..1ab745c4e1 100644 --- a/sentry_sdk/transport.py +++ b/sentry_sdk/transport.py @@ -33,15 +33,6 @@ def _make_pool(parsed_dsn, http_proxy, https_proxy): return urllib3.PoolManager(**opts) -@atexit.register -def _shutdown(): - from .hub import Hub - - main_client = Hub.main.client - if main_client is not None: - main_client.close() - - class Transport(object): def __init__(self, options=None): self.options = options @@ -53,7 +44,7 @@ def __init__(self, options=None): def capture_event(self, event): raise NotImplementedError() - def shutdown(self): + def shutdown(self, timeout, callback=None): self.kill() def kill(self): @@ -69,10 +60,7 @@ def __del__(self): class HttpTransport(Transport): def __init__(self, options): Transport.__init__(self, options) - self._worker = BackgroundWorker( - shutdown_timeout=options["shutdown_timeout"], - shutdown_callback=options["shutdown_callback"], - ) + self._worker = BackgroundWorker() self._auth = self.parsed_dsn.to_auth("sentry-python/%s" % VERSION) self._pool = _make_pool( self.parsed_dsn, @@ -130,9 +118,9 @@ def _send_event(self, event): def capture_event(self, event): self._worker.submit(lambda: self._send_event(event)) - def shutdown(self): + def shutdown(self, timeout, callback=None): logger.debug("Shutting down HTTP transport orderly") - self._worker.shutdown() + self._worker.shutdown(timeout, callback) def kill(self): logger.debug("Killing HTTP transport") diff --git a/sentry_sdk/utils.py b/sentry_sdk/utils.py index c71ca365cf..fb5e66bc1d 100644 --- a/sentry_sdk/utils.py +++ b/sentry_sdk/utils.py @@ -523,7 +523,7 @@ def strip_string(value, assume_length=None, max_length=512): return value[:max_length] -logger = logging.getLogger("sentry.errors") +logger = logging.getLogger("sentry_sdk.errors") if not logger.handlers: _handler = logging.StreamHandler(sys.stderr) _handler.setFormatter(logging.Formatter(" [sentry] %(levelname)s: %(message)s")) diff --git a/sentry_sdk/worker.py b/sentry_sdk/worker.py index 87b8a019bd..7dd6083cbf 100644 --- a/sentry_sdk/worker.py +++ b/sentry_sdk/worker.py @@ -10,17 +10,12 @@ class BackgroundWorker(object): - def __init__( - self, shutdown_timeout=10, initial_timeout=0.2, shutdown_callback=None - ): + def __init__(self): check_thread_support() self._queue = queue.Queue(-1) self._lock = threading.Lock() self._thread = None self._thread_for_pid = None - self.initial_timeout = initial_timeout - self.shutdown_timeout = shutdown_timeout - self.shutdown_callback = shutdown_callback @property def is_alive(self): @@ -57,29 +52,28 @@ def start(self): self._thread_for_pid = os.getpid() def kill(self): - logger.debug("Transport got kill request") + logger.debug("background worker got kill request") with self._lock: if self._thread: self._queue.put_nowait(_TERMINATOR) self._thread = None self._thread_for_pid = None - def shutdown(self): - logger.debug("Transport got shutdown request") + def shutdown(self, timeout, callback=None): + logger.debug("background worker got shutdown request") with self._lock: if not self.is_alive: return self._queue.put_nowait(_TERMINATOR) - timeout = self.shutdown_timeout - initial_timeout = min(self.initial_timeout, timeout) + initial_timeout = min(0.1, timeout) if not self._timed_queue_join(initial_timeout): pending = self._queue.qsize() logger.debug("%d event(s) pending on shutdown", pending) - if self.shutdown_callback is not None: - self.shutdown_callback(pending, timeout) + if callback is not None: + callback(pending, timeout) self._timed_queue_join(timeout - initial_timeout) self._thread = None - logger.debug("Transport shut down") + logger.debug("background worker shut down") def submit(self, callback): self._ensure_thread() From 6b86a3e6da2ffd36c0a3fb87c72a61cd25d94584 Mon Sep 17 00:00:00 2001 From: Armin Ronacher Date: Sat, 1 Sep 2018 17:37:37 +0200 Subject: [PATCH 8/8] fix: Lint failures --- sentry_sdk/consts.py | 1 - sentry_sdk/integrations/atexit.py | 2 ++ sentry_sdk/transport.py | 1 - 3 files changed, 2 insertions(+), 2 deletions(-) diff --git a/sentry_sdk/consts.py b/sentry_sdk/consts.py index d0ffb7d93b..f175d17266 100644 --- a/sentry_sdk/consts.py +++ b/sentry_sdk/consts.py @@ -1,4 +1,3 @@ -import os import socket diff --git a/sentry_sdk/integrations/atexit.py b/sentry_sdk/integrations/atexit.py index 17ba68d25c..964c7d58f2 100644 --- a/sentry_sdk/integrations/atexit.py +++ b/sentry_sdk/integrations/atexit.py @@ -1,3 +1,5 @@ +from __future__ import absolute_import + import os import atexit diff --git a/sentry_sdk/transport.py b/sentry_sdk/transport.py index 1ab745c4e1..90f399aeb1 100644 --- a/sentry_sdk/transport.py +++ b/sentry_sdk/transport.py @@ -1,6 +1,5 @@ from __future__ import print_function -import atexit import json import io import urllib3