diff --git a/CHANGELOG.md b/CHANGELOG.md index 8455b91a1..15d86428a 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -3,6 +3,7 @@ ## Unreleased * [Changed] DogStatsD background sender queue now evicts the oldest queued payloads when full, and drops payloads without an explicit timestamp after they have been queued for more than 10 seconds. See [#986](https://github.com/DataDog/datadogpy/pull/986). +* [Changed] Replace DogStatsD's `socket_connect_timeout` option with a boolean `socket_connect_retry`: when enabled, a UDS connection failure from the background sender is retried indefinitely with backoff instead of dropping the payload immediately. Also add an optional `timeout` parameter to `stop()` and `wait_for_pending()` to bound how long they wait for the queue to drain. See [#987](https://github.com/DataDog/datadogpy/pull/987). ## v0.53.0 / 2026-07-24 diff --git a/datadog/dogstatsd/base.py b/datadog/dogstatsd/base.py index e622c9a1e..41827c8e9 100644 --- a/datadog/dogstatsd/base.py +++ b/datadog/dogstatsd/base.py @@ -184,17 +184,22 @@ def reverse(self): # Socket options MIN_SEND_BUFFER_SIZE = 32 * 1024 -DEFAULT_SOCKET_CONNECT_TIMEOUT = 0 -UDS_CONNECT_RETRY_INITIAL_BACKOFF = 0.025 -UDS_CONNECT_RETRY_MAX_BACKOFF = 1.0 -UDS_TRANSIENT_CONNECT_ERRORS = set([errno.ENOENT, errno.ECONNREFUSED]) - -# How long (in seconds) a non-replay-safe payload may sit in the background -# sender queue before it's considered stale and dropped instead of sent. -# Payloads that carry their own explicit timestamp (replay-safe) are exempt: -# delivering those late doesn't change what they mean, so they're kept -# around until they can actually be sent. This is the default for -# sender_queue_expiry_seconds; it can be overridden per client. +# Backoff for the background sender's own retry-by-requeuing loop (see +# _sender_main_loop). Covers every retryable failure it sees, connect and +# mid-send alike -- not just connects, hence the transport-neutral name. Not +# used for direct/synchronous sends, which never retry a connection failure +# regardless of socket_connect_retry. +SENDER_RETRY_INITIAL_BACKOFF = 0.025 +SENDER_RETRY_MAX_BACKOFF = 60.0 +# How long an *unbounded* shutdown (pre_fork(), or stop()/wait_for_pending() +# called with timeout=None) still waits for a payload stuck in the retry +# loop above to resolve, before giving up on it. +SENDER_UNBOUNDED_STOP_GRACE_SECONDS = SENDER_RETRY_MAX_BACKOFF +# How often the sender retries a connection once a shutdown has been +# requested but its deadline (the caller's own timeout, or the grace period +# above) hasn't passed yet. +SENDER_STOP_RETRY_INTERVAL = 0.5 +# Default for sender_queue_expiry_seconds; it can be overridden per client. PENDING_PAYLOAD_EXPIRY_SECONDS = 10.0 # Errors seen while sending on an already-connected socket that indicate the # peer went away (e.g. the agent crashed/restarted). These are worth a single @@ -319,7 +324,7 @@ def __init__( sender_queue_timeout=0, # type: Optional[float] sender_queue_expiry_seconds=PENDING_PAYLOAD_EXPIRY_SECONDS, # type: float track_instance=True, # type: bool - socket_connect_timeout=DEFAULT_SOCKET_CONNECT_TIMEOUT, # type: Optional[float] + socket_connect_retry=False, # type: bool ): # type: (...) -> None """ Initialize a DogStatsd object. @@ -485,10 +490,19 @@ def __init__( This option does not affect hostname resolution when using UDP. :type socket_timeout: float - :param socket_connect_timeout: Set the timeout for connecting to a UNIX socket, in seconds. Optional. - Transient connection failures are retried within this timeout. If set to zero or None, do not retry. - Default: 0 (no retries). - :type socket_connect_timeout: float + :param socket_connect_retry: Only affects the background sender (disable_background_sender=False). + If True, a connection failure while sending a queued payload to a UNIX socket is retried + indefinitely, backing off up to once a minute between attempts, instead of dropping the payload + immediately. Ordinary payloads stuck retrying are eventually dropped as stale by the sender + queue's own expiry; replay-safe payloads never expire, so only a shutdown (stop(), + disable_background_sender(), pre_fork()) can end their retrying -- and it waits up to its own + timeout (or a bounded grace period, SENDER_UNBOUNDED_STOP_GRACE_SECONDS, for an unbounded call + such as pre_fork()) before giving up on one, at which point it is counted as a writer drop and + the shutdown call reports failure rather than success. wait_for_pending() never forces this: it + only waits for the queue's own retries and expiry to run their course. Direct/synchronous sends + (the default mode) always fail fast on a connection error and are unaffected by this setting. + Default: False (fail fast, matching the previous socket_connect_timeout=0 default). + :type socket_connect_retry: bool :param telemetry_socket_timeout: Set timeout for the telemetry socket operations. Optional. Effective only if either telemetry_host or telemetry_socket_path are set. @@ -559,7 +573,7 @@ def __init__( # Connection self._max_buffer_len = max_buffer_len self.socket_timeout = socket_timeout - self.socket_connect_timeout = socket_connect_timeout + self.socket_connect_retry = socket_connect_retry if socket_path is not None: self.socket_path = socket_path # type: Optional[text] self.host = None @@ -653,6 +667,19 @@ def __init__( self._queue = None # type: Optional[SenderQueue] self._sender_thread = None # type: Optional[threading.Thread] + # Set to ask a running sender thread to stop. Also what makes its + # retry backoff interruptible -- see _sender_main_loop. + self._sender_stopping = threading.Event() + # The monotonic deadline by which a requested shutdown must give up + # on a payload stuck in the retry loop, set by _stop_sender_thread() + # and read by _sender_main_loop. None means no shutdown has been + # requested (or a fresh sender hasn't seen one yet). + self._sender_stop_deadline = None # type: Optional[float] + # Set by _sender_main_loop when it gives up on a payload because the + # deadline above passed while still retrying a connection failure -- + # i.e. the queue did NOT fully drain even though the thread exited. + # _stop_sender_thread() reports this as failure rather than success. + self._sender_abandoned_payload = False self._sender_enabled = False if not disable_background_sender: @@ -763,15 +790,31 @@ def enable_background_sender( self._start_sender_thread() - def disable_background_sender(self): - # type: () -> None + def disable_background_sender(self, timeout=None): + # type: (Optional[float]) -> bool """Disable background sender mode. This call will block until all previously queued payloads are sent. + + :param timeout: Maximum number of seconds to wait for the sender thread + to drain the queue and exit. None (the default) waits indefinitely + for a sender that is genuinely busy (e.g. blocked inside a slow + send()), but not for a payload stuck retrying a connection + failure with socket_connect_retry enabled -- that case is bounded + by SENDER_UNBOUNDED_STOP_GRACE_SECONDS regardless of this + parameter, so an Agent that never comes back cannot hang this + call forever. + :type timeout: float, optional + :return: True if the sender thread finished AND the queue actually + drained. False if timeout (or the grace period above) elapsed + first -- either because the sender thread is still running, or + because it gave up on a payload still retrying a connection + failure and dropped it (counted in packets_dropped_writer) + instead of delivering it. """ with self._config_lock: self._sender_enabled = False - self._stop_sender_thread() + return self._stop_sender_thread(timeout) def disable_telemetry(self): # type: () -> None @@ -903,6 +946,15 @@ def _dedicated_telemetry_destination(self): # type: () -> bool return bool(self.telemetry_socket_path or self.telemetry_host) + def _uses_dedicated_telemetry(self, is_telemetry): + # type: (bool) -> bool + """True when this packet goes to a separate telemetry destination. + + Decides which socket/socket_path pair applies to a packet, so the + send path and its transport check can't disagree about it. + """ + return is_telemetry and self._dedicated_telemetry_destination() + # Context manager helper def __enter__(self): # type: () -> DogStatsd @@ -984,23 +1036,14 @@ def resolve_host(host, use_default_route): return get_default_route() - def get_socket(self, telemetry=False, connect_timeout=None): - # type: (bool, Optional[float]) -> _Socket + def get_socket(self, telemetry=False): + # type: (bool) -> _Socket """ Return a connected socket. Note: connect the socket before assigning it to the class instance to avoid bad thread race conditions. - - :param connect_timeout: Optional override for the UDS connect-retry - budget passed to _get_uds_socket, in place of - self.socket_connect_timeout. The send-retry loop in _xmit_packet - uses this to pass down how much of its overall deadline is left, - so a retried attempt doesn't get a brand-new full budget. """ - if connect_timeout is None: - connect_timeout = self.socket_connect_timeout - with self._socket_lock: if telemetry and self._dedicated_telemetry_destination(): if not self.telemetry_socket: @@ -1008,7 +1051,6 @@ def get_socket(self, telemetry=False, connect_timeout=None): self.telemetry_socket = self._get_uds_socket( self.telemetry_socket_path, self.telemetry_socket_timeout, - connect_timeout, ) else: self.telemetry_socket = self._get_udp_socket( @@ -1024,7 +1066,6 @@ def get_socket(self, telemetry=False, connect_timeout=None): self.socket = self._get_uds_socket( self.socket_path, self.socket_timeout, - connect_timeout, ) else: self.socket = self._get_udp_socket( @@ -1060,8 +1101,10 @@ def _ensure_min_send_buffer_size(cls, sock, min_size=MIN_SEND_BUFFER_SIZE): log.debug("Socket send buffer increased to %dkb", min_size / 1024) @classmethod - def _get_uds_socket(cls, socket_path, timeout, connect_timeout): - # type: (Text, Optional[float], Optional[float]) -> _Socket + def _get_uds_socket(cls, socket_path, timeout): + # type: (Text, Optional[float]) -> _Socket + """Make one connect attempt per candidate socket kind. + """ valid_socket_kinds = [socket.SOCK_DGRAM, socket.SOCK_STREAM] if socket_path.startswith(UNIX_ADDRESS_DATAGRAM_SCHEME): valid_socket_kinds = [socket.SOCK_DGRAM] @@ -1072,56 +1115,30 @@ def _get_uds_socket(cls, socket_path, timeout, connect_timeout): elif socket_path.startswith(UNIX_ADDRESS_SCHEME): socket_path = socket_path[len(UNIX_ADDRESS_SCHEME):] - last_error = socket.timeout("timed out connecting to UDS socket") # type: Exception - deadline = None - if connect_timeout and connect_timeout > 0: - deadline = time.time() + connect_timeout - - for socket_kind in valid_socket_kinds: + for index, socket_kind in enumerate(valid_socket_kinds): + is_last_kind = index == len(valid_socket_kinds) - 1 # py2 stores socket kinds differently than py3, determine the name independently from version sk_name = {socket.SOCK_STREAM: "stream", socket.SOCK_DGRAM: "datagram"}[socket_kind] - - backoff = UDS_CONNECT_RETRY_INITIAL_BACKOFF - - while deadline is None or time.time() < deadline: - sock = None - try: - connect_attempt_timeout = timeout - if deadline is not None: - connect_attempt_timeout = deadline - time.time() - if connect_attempt_timeout <= 0: - break - - sock = socket.socket(socket.AF_UNIX, socket_kind) - sock.settimeout(connect_attempt_timeout) - cls._ensure_min_send_buffer_size(sock) - sock.connect(socket_path) - sock.settimeout(timeout) - log.debug("Connected to socket %s with kind %s", socket_path, sk_name) - return sock - except Exception as e: - if sock is not None: - sock.close() - log.debug("Failed to connect to %s with kind %s: %s", socket_path, sk_name, e) - if getattr(e, "errno", None) == errno.EPROTOTYPE: - last_error = e - break - if ( - deadline is not None - and getattr(e, "errno", None) in UDS_TRANSIENT_CONNECT_ERRORS - ): - last_error = e - remaining_time = max(0, deadline - time.time()) - sleep_time = min(backoff, remaining_time) - if sleep_time > 0: - time.sleep(sleep_time) - backoff = min(backoff * 2, UDS_CONNECT_RETRY_MAX_BACKOFF) - continue - raise e - if getattr(last_error, "errno", None) == errno.EPROTOTYPE: - continue - raise last_error - raise last_error + sock = None + try: + sock = socket.socket(socket.AF_UNIX, socket_kind) + sock.settimeout(timeout) + cls._ensure_min_send_buffer_size(sock) + sock.connect(socket_path) + log.debug("Connected to socket %s with kind %s", socket_path, sk_name) + return sock + except Exception as e: + if sock is not None: + sock.close() + log.debug("Failed to connect to %s with kind %s: %s", socket_path, sk_name, e) + if getattr(e, "errno", None) == errno.EPROTOTYPE and not is_last_kind: + # Wrong socket kind for this address -- try the other one. + continue + raise + # Unreachable: valid_socket_kinds is never empty, and every path above + # either returns or raises. Present only so this always has a return + # type of _Socket. + raise socket.error("no usable socket kind for {}".format(socket_path)) @classmethod def _get_udp_socket(cls, host, port, timeout): @@ -1677,6 +1694,12 @@ def _account_dropped_expired(self, item): self.packets_dropped_expired += 1 self.bytes_dropped_expired += len(payload_text(item).encode(self.encoding)) + def _account_dropped_writer(self, item): + # type: (QueuedItem) -> None + """A payload could not be written and is not being retried further.""" + self.packets_dropped_writer += 1 + self.bytes_dropped_writer += len(payload_text(item).encode(self.encoding)) + def _flush_telemetry(self): # type: () -> str tags = self._client_tags[:] @@ -1719,12 +1742,21 @@ def _is_telemetry_flush_time(self): def _send_to_server(self, packet, replay_safe=False): # type: (str, bool) -> None - # Skip the lock if the queue is None. There is no race with enable_background_sender. - if self._queue is not None: - # Prevent a race with disable_background_sender. + # Skip the lock if the queue is None or a shutdown has already been + # requested -- the lock-protected recheck below is what actually has + # to be correct; this is purely an optimization to avoid the lock + # once there is clearly nothing to enqueue onto. + if self._queue is not None and not self._sender_stopping.is_set(): + # Prevent a race with disable_background_sender: _sender_stopping + # is set BEFORE _stop_sender_thread() appends Stop, under this + # same lock, so rechecking it here (not just self._queue) means a + # racing put() either lands strictly before Stop (still open) or + # is rejected outright and falls through to a direct send below + # -- never behind Stop, where it would be silently lost once the + # sender reaches Stop and exits. with self._buffer_lock: packet_with_newline = packet + '\n' - if self._queue is not None: + if self._queue is not None and not self._sender_stopping.is_set(): if replay_safe: # Never expires, so it needs no enqueued_at and no # wrapper at all: queue the bare string and let the @@ -1764,81 +1796,36 @@ def _xmit_packet_with_telemetry(self, packet, queue_mode=False): return sent - def _installed_socket(self, is_telemetry): - # type: (bool) -> Optional[_Socket] - """ - The socket a send for this packet would use, if one is already installed. - - Returns None when a fresh connection would have to be established, which - is the only situation where the socket_connect_timeout budget applies. - Callers that mean to gate on "would we have to connect?" must use this - rather than the deadline alone. - """ - if is_telemetry and self._dedicated_telemetry_destination(): - return self.telemetry_socket - return self.socket - def _xmit_packet(self, packet, is_telemetry, queue_mode=False): # type: (str, bool, bool) -> Optional[bool] - """Attempt to send packet, retrying a reconnect within this call as budget allows. + """Attempt to send a packet, once. Returns True if sent. Otherwise returns False for a definitive, non-retryable failure (already accounted for as a dropped packet), - or -- only when queue_mode is True -- None for a connection failure - that the sender queue should retry by requeuing the payload rather - than have accounted for here as a drop. + or -- only when queue_mode is True, the transport is UDS, and + socket_connect_retry is enabled -- None for a connection failure that + the sender queue should retry by requeuing the payload (with its own + backoff, capped at SENDER_RETRY_MAX_BACKOFF) rather than have + accounted for here as a drop. """ - if is_telemetry and self._dedicated_telemetry_destination(): + if self._uses_dedicated_telemetry(is_telemetry): uses_uds = self.telemetry_socket_path is not None else: uses_uds = self.socket_path is not None - retry_deadline = None - if uses_uds and self.socket_connect_timeout and self.socket_connect_timeout > 0: - retry_deadline = time.time() + self.socket_connect_timeout + # Direct/synchronous sends (queue_mode=False) always fail fast on a + # connection error: there is no queue expiry to protect a calling + # thread from retrying indefinitely, so socket_connect_retry does not + # apply to them. Reconnect-and-retry also stays UDS-only, matching + # UDP's different (connectionless) failure semantics. + retry_eligible = queue_mode and uses_uds and self.socket_connect_retry - backoff = UDS_CONNECT_RETRY_INITIAL_BACKOFF - sent = None # type: Optional[bool] - while True: - # Cheap fast-path check before even trying to acquire _socket_lock. - if ( - retry_deadline is not None - and retry_deadline - time.time() <= 0 - and not self._installed_socket(is_telemetry) - ): - log.warning( - "Gave up reconnecting after socket_connect_timeout (%ss), dropping the packet", - self.socket_connect_timeout, - ) - sent = None - break + sent = self._xmit_packet_attempt(packet, is_telemetry, retry_eligible) + if sent: + return True - sent = self._xmit_packet_attempt( - packet, is_telemetry, retry_eligible=retry_deadline is not None, retry_deadline=retry_deadline - ) - if sent: - return True - # `sent` is False for a definitive failure (already logged/dropped - # above), or None for a transient one that's worth reconnecting - # and retrying, bounded by socket_connect_timeout. - # - # A None result implies retry_eligible, which implies a deadline was - # set; the explicit check keeps that invariant locally provable - # (for readers and for the type checker) instead of implicit. - if sent is not None or retry_deadline is None: - break - remaining = retry_deadline - time.time() - if remaining <= 0: - log.warning( - "Gave up reconnecting after socket_connect_timeout (%ss), dropping the packet", - self.socket_connect_timeout, - ) - break - time.sleep(min(backoff, remaining)) - backoff = min(backoff * 2, UDS_CONNECT_RETRY_MAX_BACKOFF) - - if sent is None and queue_mode: + if sent is None: # Connection trouble, and the caller is the background sender # queue: let it requeue the payload and retry once reconnected, # instead of dropping it here. @@ -1849,48 +1836,24 @@ def _xmit_packet(self, packet, is_telemetry, queue_mode=False): self.packets_dropped_writer += 1 return False - def _xmit_packet_attempt(self, packet, is_telemetry, retry_eligible, retry_deadline=None): - # type: (str, bool, bool, Optional[float]) -> Optional[bool] + def _xmit_packet_attempt(self, packet, is_telemetry, retry_eligible): + # type: (str, bool, bool) -> Optional[bool] """ Attempt to send a single packet. Returns True if the packet was sent, False if it should be dropped without retrying, or None if the failure is transient (the peer went - away), `retry_eligible` is set, and it's worth a reconnect-and-retry. - - :param retry_deadline: Optional absolute time.time()-based deadline - for the overall _xmit_packet retry operation. The connect_timeout - handed to get_socket() is computed from this only after - _socket_lock is actually acquired (not before), so time spent - waiting on a contended lock counts against the budget instead of - silently extending it. + away, which includes a failed reconnect), `retry_eligible` is set, + and it's worth a reconnect-and-retry by the caller. """ socket_kind = None with self._socket_lock: try: - # Captured under the lock, before the deadline check: whether a - # socket already exists decides whether that deadline is even - # relevant to this attempt. - existing_socket = self._installed_socket(is_telemetry) - - connect_timeout = self.socket_connect_timeout - if retry_deadline is not None: - connect_timeout = retry_deadline - time.time() - if connect_timeout <= 0 and not existing_socket: - log.warning( - "Gave up reconnecting after socket_connect_timeout (%ss), dropping the packet", - self.socket_connect_timeout, - ) - return False - - if is_telemetry and self._dedicated_telemetry_destination(): - mysocket = existing_socket or self.get_socket( - telemetry=True, connect_timeout=connect_timeout - ) + if self._uses_dedicated_telemetry(is_telemetry): + mysocket = self.get_socket(telemetry=True) socket_kind = self._telemetry_socket_kind else: - # If set, use socket directly - mysocket = existing_socket or self.get_socket(connect_timeout=connect_timeout) + mysocket = self.get_socket() socket_kind = self._socket_kind encoded_packet = packet.encode(self.encoding) @@ -2201,6 +2164,12 @@ def _start_sender_thread(self): if self._queue is not None: return + # A previous _stop_sender_thread() leaves this set; clear it before the + # new sender starts so it doesn't immediately think it's shutting down. + self._sender_stopping.clear() + self._sender_stop_deadline = None + self._sender_abandoned_payload = False + self._queue = SenderQueue( self._sender_queue_size, self._sender_queue_expiry_seconds, @@ -2218,26 +2187,86 @@ def _start_sender_thread(self): self._sender_thread.daemon = True self._sender_thread.start() - def _stop_sender_thread(self): - # type: () -> None - # Lock ensures that nothing gets added to the queue after we disable it. + def _stop_sender_thread(self, timeout=None): + # type: (Optional[float]) -> bool + # Ask the sender to stop, with a deadline it must honor BEFORE + # abandoning a payload stuck in its retry-by-requeuing loop (see + # _sender_main_loop): a bounded caller's own timeout IS that + # deadline, so the sender keeps retrying for close to as long as the + # caller asked instead of giving up the instant a shutdown is + # requested. An unbounded caller (timeout=None, e.g. pre_fork()) gets + # a bounded grace period instead, so it can't hang forever on a + # payload that can genuinely never succeed. + grace = SENDER_UNBOUNDED_STOP_GRACE_SECONDS if timeout is None else timeout + self._sender_stop_deadline = monotonic() + grace + # Setting makes _send_to_server() reject any FURTHER producer call outright. + self._sender_stopping.set() + + # Set temporary var to protect from concurrent access to self._queue. + queue_to_close = self._queue + if queue_to_close is not None: + queue_to_close.close() + + # Lock ensures that nothing gets added to the queue after the check + # above -- see _send_to_server(), which takes this same lock and + # re-checks _sender_stopping before it puts. with self._buffer_lock: - if not self._queue: - return - self._queue.put(Stop) - self._queue = None + if self._queue is not None: + # put() lets the Stop sentinel past the size limit, so this + # never blocks even when the queue is full. + self._queue.put(Stop) + + thread = self._sender_thread + if thread is None: + # Nothing left to stop: no thread ever runs, so clear any residual + # queue. + with self._buffer_lock: + self._queue = None + return True - if self._sender_thread is not None: - self._sender_thread.join() + thread.join(timeout) + if thread.is_alive(): + # Timed out. Leave _queue in place. + return False + + # The thread exited, but that alone doesn't mean it drained: it may + # have hit the deadline above with a payload still stuck retrying and + # given up on it instead (see _sender_main_loop). That payload was + # never delivered, so report failure rather than claiming success. + abandoned = self._sender_abandoned_payload + + # _sender_main_loop clears this state on its way out when the thread + # has drained the queue, so this may already be a no-op; it also covers + # a thread that exited without draining (e.g. never actually started). + with self._buffer_lock: + self._queue = None self._sender_thread = None + return not abandoned + + def _release_sender_state(self, pending_queue): + # type: (SenderQueue) -> None + """Drop the client's pointers to this sender's queue and thread. + + Called by the sender thread on its way out so the client self-heals + and a later enable_background_sender() can start fresh -- even when + this exit wasn't awaited, e.g. a stop(timeout) that timed out and the + caller moved on. The guards check this thread still owns the fields: a + newer sender may have already taken over. + """ + with self._buffer_lock: + if self._queue is pending_queue: + self._queue = None + if self._sender_thread is threading.current_thread(): + self._sender_thread = None def _sender_main_loop(self, pending_queue): # type: (SenderQueue) -> None - backoff = UDS_CONNECT_RETRY_INITIAL_BACKOFF + backoff = SENDER_RETRY_INITIAL_BACKOFF while True: item = pending_queue.get() if item is Stop: pending_queue.task_done(item) + self._release_sender_state(pending_queue) return # payload_text() also narrows the type: 'if item is Stop' above is @@ -2248,25 +2277,65 @@ def _sender_main_loop(self, pending_queue): ) if sent is None: - # Connection trouble: keep the payload for the next attempt + # Connection trouble. A shutdown may already have been + # requested (see _stop_sender_thread) with a deadline this + # payload must be given a real chance against before it's + # given up on -- check that now, before requeuing, while the + # queue still considers this item in flight and can finish it + # outright instead. + deadline = self._sender_stop_deadline + if deadline is not None and monotonic() >= deadline: + # The deadline has passed: give up on this payload for + # good rather than requeuing it for a retry that will + # never be awaited. Account for it as a writer drop -- + # it was never delivered -- instead of letting it vanish + # with no telemetry at all. + self._account_dropped_writer(item) # type: ignore[arg-type] + pending_queue.task_done(item) # type: ignore[arg-type] + self._sender_abandoned_payload = True + self._release_sender_state(pending_queue) + return + + # Still worth retrying: keep the payload for the next attempt # instead of losing it. The queue's own expiry check (on a # future get()) is what eventually gives up on a payload # that's been stuck for too long, unless it's replay-safe. pending_queue.requeue_front(item) # type: ignore[arg-type] - time.sleep(backoff) - backoff = min(backoff * 2, UDS_CONNECT_RETRY_MAX_BACKOFF) + + if deadline is None: + # No shutdown requested (yet): normal interruptible + # backoff. wait() returns early -- before backoff fully + # elapses -- the instant a shutdown IS requested, so the + # very next iteration's deadline check above fires + # promptly instead of waiting out a long backoff. + self._sender_stopping.wait(backoff) + backoff = min(backoff * 2, SENDER_RETRY_MAX_BACKOFF) + else: + # A shutdown was requested and its deadline hasn't + # passed yet: keep retrying -- a still-recovering Agent + # should still get drained -- but pace attempts instead + # of hammering a connection that keeps failing, and + # never sleep past the deadline. + remaining = deadline - monotonic() + time.sleep(min(SENDER_STOP_RETRY_INTERVAL, max(remaining, 0))) continue # Sent, or a definitive failure that _xmit_packet already # accounted for as a dropped packet -- either way, this # payload's story is over. pending_queue.task_done(item) # type: ignore[arg-type] - backoff = UDS_CONNECT_RETRY_INITIAL_BACKOFF + backoff = SENDER_RETRY_INITIAL_BACKOFF - def wait_for_pending(self): - # type: () -> None + def wait_for_pending(self, timeout=None): + # type: (Optional[float]) -> bool """ Flush the buffer and wait for all queued payloads to be written to the server. + + :param timeout: Maximum number of seconds to wait for the queue to + drain. None (the default) waits indefinitely. + :type timeout: float, optional + :return: True if every queued payload has been sent, dropped or + expired, False if timeout elapsed with payloads still outstanding. """ self.flush_buffered_metrics() @@ -2276,8 +2345,10 @@ def wait_for_pending(self): # check and join later. queue = self._queue - if queue is not None: - queue.join() + if queue is None: + return True + + return queue.join(timeout) def pre_fork(self): # type: () -> None @@ -2285,6 +2356,11 @@ def pre_fork(self): Flush any pending payloads and stop all background threads. + A payload stuck retrying a connection failure (socket_connect_retry) + is given up on -- and counted as a writer drop -- after a bounded + grace period (SENDER_UNBOUNDED_STOP_GRACE_SECONDS) rather than + blocking the fork indefinitely; that payload is not delivered. + The client should not be used from this point until state is restored by calling post_fork_parent() or post_fork_child(). @@ -2330,21 +2406,75 @@ def post_fork_child(self): self._start_flush_thread() self._start_sender_thread() - def stop(self): - # type: () -> None + def stop(self, timeout=None): + # type: (Optional[float]) -> bool """Stop the client. Disable buffering, aggregation, background sender and flush any pending payloads to the server. Client remains usable after this method, but sending metrics may block if socket_timeout is enabled. - """ - self.disable_background_sender() + :param timeout: Maximum number of seconds to wait for the background + sender to drain its queue and exit. None (the default) waits + indefinitely for a sender that is genuinely busy (e.g. blocked + inside a slow send()). It does NOT wait indefinitely for a + payload stuck retrying a connection failure with + socket_connect_retry enabled: that case is bounded by + SENDER_UNBOUNDED_STOP_GRACE_SECONDS regardless of this parameter, + so an Agent that never comes back cannot hang stop() forever. + :type timeout: float, optional + :return: True if the background sender drained and stopped, and the + final flush and socket close ran. False if timeout (or the grace + period above) elapsed first, in which case neither the final + flush nor the socket close ran (see below) -- either because the + sender thread is still running, or because it gave up on a + payload still retrying a connection failure and dropped it + (counted in packets_dropped_writer) instead of delivering it. Do + not call stop() again while that sender is still running: it + queues a second internal shutdown signal that is never drained, + which can make wait_for_pending() block forever on the abandoned + queue. Use wait_for_pending() to wait for the sender instead, + then call stop() again once it has actually stopped. + """ + + stopped = self.disable_background_sender(timeout) self._disable_buffering = True self._disable_aggregation = True + + if not stopped: + if self._sender_abandoned_payload: + # The sender thread did exit, but only by giving up on a + # payload still retrying a connection failure once its + # deadline passed (see _sender_main_loop) -- that payload was + # dropped, not delivered, so this is not a clean stop either. + log.warning( + "stop() gave up on a payload stuck retrying a connection failure after " + "%ss; it was dropped instead of delivered (see packets_dropped_writer). " + "Skipping the final flush and socket close", + timeout, + ) + else: + # We gave up waiting, so the sender thread is still running and + # still owns the socket -- it can be parked inside a send() with + # _socket_lock held. Flushing or closing here would block on that + # same lock for as long as the sender stays wedged, which would + # make timeout meaningless: the caller asked for a bounded stop(). + # Pushing more data through that socket could not succeed anyway, + # and closing it from under a thread mid-write is not safe. Leave + # it open; the OS reclaims the fd when the process exits, and the + # sender thread is a daemon so it never holds up interpreter + # shutdown. + log.warning( + "stop() timed out after %ss with the background sender still running; " + "skipping the final flush and socket close", + timeout, + ) + return False + self.flush_aggregated_metrics() self.flush_buffered_metrics() self.close_socket() + return True statsd = DogStatsd() diff --git a/datadog/dogstatsd/sender_queue.py b/datadog/dogstatsd/sender_queue.py index f6f699e76..23c2873d6 100644 --- a/datadog/dogstatsd/sender_queue.py +++ b/datadog/dogstatsd/sender_queue.py @@ -115,6 +115,12 @@ def __init__(self, maxsize, expiry_seconds, on_drop_queue_full, on_drop_expired, # all tasks have been dropped or sent. self._unfinished_tasks = 0 + # Set by close(): tells a put() that is (or will be) waiting for room + # to stop waiting immediately instead of riding out put_timeout, or + # forever if put_timeout is None. See close()'s docstring for why + # this exists. + self._closing = False + # The items currently handed out by get() and not yet finished via # requeue_front() or task_done(), keyed by id(item). SenderQueue is # single-consumer by design (one background sender thread), so this @@ -178,18 +184,23 @@ def put(self, item): number, or not at all if it's 0 (the default) -- then falls back to evicting the oldest entry (see _make_room_locked()) if the queue is still full once the wait is over. Either way, put() never rejects - the payload outright. + the payload outright. A close() call (from any thread) cuts any of + that waiting short immediately, regardless of put_timeout. """ with self._not_empty: if item is not Stop and self._maxsize > 0 and len(self._deque) >= self._maxsize: - if self._put_timeout is None: + if self._closing: + pass # Already closing: don't wait at all, straight to eviction below. + elif self._put_timeout is None: # Wait forever: an explicit opt-in to unbounded - # backpressure on the calling thread. - while len(self._deque) >= self._maxsize: + # backpressure on the calling thread. close() is what + # keeps this from actually being forever once a shutdown + # is underway. + while len(self._deque) >= self._maxsize and not self._closing: self._not_full.wait() elif self._put_timeout > 0: deadline = monotonic() + self._put_timeout - while len(self._deque) >= self._maxsize: + while len(self._deque) >= self._maxsize and not self._closing: remaining = deadline - monotonic() if remaining <= 0: break @@ -204,6 +215,36 @@ def put(self, item): self._unfinished_tasks += 1 self._not_empty.notify() + def close(self): + # type: () -> None + """Wake any put() currently waiting for room, immediately. + + A put() blocked waiting for space in a full queue holds no lock this + method needs -- Condition.wait() releases the underlying lock while + waiting -- so this always runs promptly, even while some other + thread is stuck inside that wait (that stuck thread is exactly what + this is for). Without it, a put() with put_timeout=None waits + forever for room that will never open up once nothing is draining + the queue, and even a bounded put_timeout can outlast whatever + timeout a caller trying to shut things down asked for -- see + DogStatsd._stop_sender_thread(), which calls this before it needs + the *caller's* lock (_buffer_lock) that a stuck put() would + otherwise be holding for the entire wait. + + Sticky: once closed, no future put() on this queue ever waits for + room again, regardless of put_timeout -- it goes straight to + eviction, like put_timeout=0. There is no matching "reopen": a fresh + shutdown starts with a fresh SenderQueue instead. + + This does not stop put()/get() from working, and does not reject or + drop anything by itself -- it only ends a wait early. Refusing new + payloads outright is the caller's job (see DogStatsd._send_to_server(), + which checks _sender_stopping before ever calling put()). + """ + with self._not_full: + self._closing = True + self._not_full.notify_all() + def requeue_front(self, item): # type: (QueuedItem) -> None """Put an in-flight payload back at the front after a failed send attempt. @@ -344,11 +385,31 @@ def task_done(self, item): return self._finish_task_locked() - def join(self): - # type: () -> None + def join(self, timeout=None): + # type: (Optional[float]) -> bool + """Wait until every queued payload has been sent, dropped or expired. + + :param timeout: Maximum number of seconds to wait. None (the default) + waits indefinitely. + :return: True if nothing is outstanding any more, False if timeout + elapsed while payloads were still in flight. + """ with self._all_tasks_done: + if timeout is None: + while self._unfinished_tasks: + self._all_tasks_done.wait() + return True + + # Condition.wait()'s return value can't be used to detect a + # timeout: on Python 2 it is always None. Track the deadline + # ourselves instead, the same way put() does for put_timeout. + deadline = monotonic() + timeout while self._unfinished_tasks: - self._all_tasks_done.wait() + remaining = deadline - monotonic() + if remaining <= 0: + return False + self._all_tasks_done.wait(remaining) + return True def qsize(self): # type: () -> int diff --git a/tests/unit/dogstatsd/test_statsd.py b/tests/unit/dogstatsd/test_statsd.py index d6179841e..73f877455 100644 --- a/tests/unit/dogstatsd/test_statsd.py +++ b/tests/unit/dogstatsd/test_statsd.py @@ -31,7 +31,7 @@ # Datadog libraries from datadog import initialize, statsd from datadog import __version__ as version -from datadog.dogstatsd.base import DEFAULT_BUFFERING_FLUSH_INTERVAL, DEFAULT_HOST, DEFAULT_PORT, DogStatsd, MIN_SEND_BUFFER_SIZE, PENDING_PAYLOAD_EXPIRY_SECONDS, PendingPayload, SenderQueue, Stop, UDP_OPTIMAL_PAYLOAD_LENGTH, UDS_CONNECT_RETRY_INITIAL_BACKOFF, UDS_OPTIMAL_PAYLOAD_LENGTH +from datadog.dogstatsd.base import DEFAULT_BUFFERING_FLUSH_INTERVAL, DEFAULT_HOST, DEFAULT_PORT, DogStatsd, MIN_SEND_BUFFER_SIZE, PENDING_PAYLOAD_EXPIRY_SECONDS, PendingPayload, SenderQueue, Stop, UDP_OPTIMAL_PAYLOAD_LENGTH, SENDER_RETRY_MAX_BACKOFF, UDS_OPTIMAL_PAYLOAD_LENGTH from datadog.dogstatsd.sender_queue import is_replay_safe, payload_text from datadog.util.compat import monotonic as sender_queue_clock from datadog.dogstatsd.context import TimedContextManagerDecorator @@ -951,79 +951,50 @@ def test_socket_error(self): mock.ANY, ) - def _uds_statsd(self, connect_timeout): + def _uds_statsd(self, socket_connect_retry=False): """ A UDS-backed client whose current socket is already broken. The reconnect-and-retry path in _xmit_packet is deliberately scoped to UDS only, so these tests must not use the default UDP client. """ - statsd = DogStatsd(socket_path='/tmp/dogstatsd-test.sock', telemetry_min_flush_interval=0) - statsd.socket_connect_timeout = connect_timeout + statsd = DogStatsd( + socket_path='/tmp/dogstatsd-test.sock', + disable_telemetry=True, + socket_connect_retry=socket_connect_retry, + ) statsd.socket = BrokenSocket(error_number=errno.ECONNREFUSED) - statsd._reset_telemetry() return statsd @patch('datadog.dogstatsd.base.DogStatsd._get_uds_socket') - def test_socket_connection_error_reconnects_and_resends(self, mock_get_uds_socket): - working_socket = FakeSocket() - mock_get_uds_socket.return_value = working_socket - statsd = self._uds_statsd(connect_timeout=5) - - with mock.patch("datadog.dogstatsd.base.log") as mock_log: - statsd.gauge('reconnected', 1) - statsd.flush() - - mock_log.error.assert_not_called() - mock_log.warning.assert_not_called() - - # The packet was not dropped: it was resent once a fresh socket was obtained. - mock_get_uds_socket.assert_called_once() - self.assertEqual(statsd.packets_dropped_writer, 0) - self.assertEqual(working_socket.payloads[0].decode('utf-8'), 'reconnected:1|g\n') - - def test_socket_connection_error_drops_packet_if_reconnect_also_fails(self): - # Small deadline so the retry loop gives up quickly in the test. - statsd = self._uds_statsd(connect_timeout=0.05) - - with mock.patch.object( - DogStatsd, '_get_uds_socket', side_effect=socket.error(errno.ECONNREFUSED, "still refused") - ): + def test_direct_send_never_retries_uds_even_with_socket_connect_retry_enabled(self, mock_get_uds_socket): + # socket_connect_retry only affects the background sender: a direct/ + # synchronous send (background sender disabled, the default) has no + # queue and no expiry to protect the calling thread, so it must always + # fail fast on a connection error regardless of this setting. + mock_get_uds_socket.return_value = FakeSocket() + statsd = self._uds_statsd(socket_connect_retry=True) + + with mock.patch.object(statsd.socket, 'send', wraps=statsd.socket.send) as mock_send: with mock.patch("datadog.dogstatsd.base.log") as mock_log: - statsd.gauge('no error', 1) - statsd.flush() + statsd.gauge('not reconnected', 1) mock_log.error.assert_not_called() + mock_log.warning.assert_called_once_with( + "Error submitting packet: %s, dropping the packet and closing the socket", + mock.ANY, + ) - # Both the metric and the telemetry flush hit the same broken reconnect and get dropped - # once the retry deadline is exhausted. - self.assertEqual(statsd.packets_dropped_writer, 2) - - @patch('datadog.dogstatsd.base.DogStatsd._get_uds_socket') - def test_socket_connection_error_retries_multiple_times_before_success(self, mock_get_uds_socket): - working_socket = FakeSocket() - mock_get_uds_socket.side_effect = [ - socket.error(errno.ECONNREFUSED, "still refused"), - socket.error(errno.ECONNREFUSED, "still refused"), - working_socket, - ] - # Long enough to cover a few backoff sleeps well under a second. - statsd = self._uds_statsd(connect_timeout=5) - - with mock.patch("datadog.dogstatsd.base.log") as mock_log: - statsd.gauge('reconnected after retries', 1) - - mock_log.warning.assert_not_called() + # A single send attempt on the broken socket, then dropped -- no + # reconnect-and-resend, and no connect was ever attempted either. + mock_send.assert_called_once() - # Two failed reconnect attempts, then a third that finally succeeds. - self.assertEqual(mock_get_uds_socket.call_count, 3) - self.assertEqual(statsd.packets_dropped_writer, 0) - self.assertTrue(working_socket.payloads[0].decode('utf-8').startswith('reconnected after retries:1|g')) + mock_get_uds_socket.assert_not_called() def test_socket_connection_error_does_not_retry_for_udp(self): # Reconnect-and-retry is UDS-only: a UDP client drops the packet on a - # transient connection error even when socket_connect_timeout is set. - self.statsd.socket_connect_timeout = 5 + # transient connection error even with socket_connect_retry enabled. + self.statsd.socket_connect_retry = True broken_socket = BrokenSocket(error_number=errno.ECONNREFUSED) self.statsd.socket = broken_socket @@ -1040,67 +1011,30 @@ def test_socket_connection_error_does_not_retry_for_udp(self): # No reconnect-and-resend: a single send attempt, then dropped. mock_send.assert_called_once() - @patch('datadog.dogstatsd.base.DogStatsd._get_uds_socket') - def test_expired_deadline_sends_on_socket_installed_by_another_thread(self, mock_get_uds_socket): - # socket_connect_timeout budgets *connecting*. If another thread already - # installed a healthy socket while this one waited for _socket_lock, the - # spent budget is irrelevant: sending on it does no connecting, so the - # packet must not be dropped at the moment of recovery. - statsd = self._uds_statsd(connect_timeout=5) - working_socket = FakeSocket() - statsd.socket = working_socket - - with mock.patch("datadog.dogstatsd.base.log") as mock_log: - sent = statsd._xmit_packet_attempt( - 'recovered:1|g\n', - is_telemetry=False, - retry_eligible=True, - retry_deadline=time.time() - 1, # budget consumed while waiting for the lock - ) - - self.assertTrue(sent) - mock_log.warning.assert_not_called() - - # The installed socket was used as-is; no connect was attempted. - mock_get_uds_socket.assert_not_called() - self.assertEqual(working_socket.payloads[0].decode('utf-8'), 'recovered:1|g\n') - - @patch('datadog.dogstatsd.base.DogStatsd._get_uds_socket') - def test_expired_deadline_drops_when_a_connect_would_be_needed(self, mock_get_uds_socket): - # The complement of the case above, and the reason the guard exists at - # all: with no socket installed, an expired budget must not reach - # get_socket(), which treats a <= 0 connect_timeout as "unbounded". - statsd = self._uds_statsd(connect_timeout=5) - statsd.socket = None - - with mock.patch("datadog.dogstatsd.base.log") as mock_log: - sent = statsd._xmit_packet_attempt( - 'dropped:1|g\n', - is_telemetry=False, - retry_eligible=True, - retry_deadline=time.time() - 1, - ) - - self.assertFalse(sent) - mock_log.warning.assert_called_once_with( - "Gave up reconnecting after socket_connect_timeout (%ss), dropping the packet", - 5, - ) - - mock_get_uds_socket.assert_not_called() + def test_socket_connect_retry_defaults_to_false(self): + self.assertFalse(self.statsd.socket_connect_retry) def test_concurrent_reconnect_does_not_drop_backlog_after_recovery(self): - # End-to-end ordering: thread A reconnects slowly while holding - # _socket_lock; B queues behind it and has its entire connect budget - # consumed by the wait. Once A installs a healthy socket, B must send on - # it rather than discard its packet. - statsd = self._uds_statsd(connect_timeout=0.2) + # End-to-end ordering: thread A connects slowly while holding + # _socket_lock (no socket installed yet); B queues right behind it on + # the same lock. Once A installs a healthy socket, B must send on it + # rather than also trying to connect -- get_socket() always prefers an + # already-installed socket over connecting again. Direct sends never + # retry now, so unlike the old version of this test, neither packet + # can survive a *failed* connect; this exercises the still-relevant + # part, concurrent first-time connects sharing one outcome. + statsd = DogStatsd( + socket_path='/tmp/dogstatsd-test-concurrent.sock', + disable_telemetry=True, + ) working_socket = FakeSocket() a_is_connecting = threading.Event() release_connect = threading.Event() + connect_calls = [] def slow_connect(*args, **kwargs): + connect_calls.append(1) a_is_connecting.set() release_connect.wait(5) return working_socket @@ -1113,14 +1047,17 @@ def slow_connect(*args, **kwargs): thread_b = Thread(target=lambda: statsd.gauge('from-b', 1)) thread_b.start() - # Let B's whole 0.2s budget elapse while it blocks on _socket_lock, - # then let A's connect succeed and install the socket. - time.sleep(0.5) + # Give B a moment to reach and block on _socket_lock (held by A's + # in-progress connect), then let A's connect succeed and install + # the socket. + time.sleep(0.2) release_connect.set() thread_a.join(5) thread_b.join(5) self.assertFalse(thread_a.is_alive()) + # B reused the socket A installed rather than connecting itself. + self.assertEqual(connect_calls, [1]) self.assertFalse(thread_b.is_alive()) payloads = [payload.decode('utf-8') for payload in working_socket.payloads] @@ -1159,25 +1096,6 @@ def send(self, payload): closer_thread.join(2) self.assertFalse(closer_thread.is_alive()) - def test_socket_connection_error_does_not_retry_without_connect_timeout(self): - # socket_connect_timeout defaults to 0 (unset): no reconnect attempt should be made - # for this packet, it should be dropped immediately instead. - self.assertEqual(self.statsd.socket_connect_timeout, 0) - broken_socket = BrokenSocket(error_number=errno.ECONNREFUSED) - self.statsd.socket = broken_socket - - with mock.patch.object(broken_socket, 'send', wraps=broken_socket.send) as mock_send: - with mock.patch("datadog.dogstatsd.base.log") as mock_log: - self.statsd.gauge('not reconnected', 1) - - mock_log.warning.assert_called_once_with( - "Error submitting packet: %s, dropping the packet and closing the socket", - mock.ANY, - ) - - # Only one send attempt was made on the broken socket: no reconnect-and-resend. - mock_send.assert_called_once() - def test_socket_overflown(self): self.statsd.socket = OverflownSocket() with mock.patch("datadog.dogstatsd.base.log") as mock_log: @@ -1234,82 +1152,46 @@ def test_uds_socket_ensures_min_receive_buffer(self, mock_socket_create): MIN_SEND_BUFFER_SIZE, ) - @patch('datadog.dogstatsd.base.time.sleep') - @patch('socket.socket') - def test_uds_socket_retries_missing_socket_until_timeout(self, mock_socket_create, mock_sleep): - missing_socket_error = socket.error(errno.ENOENT, os.strerror(errno.ENOENT)) - first_socket = Mock() - first_socket.connect.side_effect = missing_socket_error - first_socket.getsockopt.return_value = MIN_SEND_BUFFER_SIZE - second_socket = Mock() - second_socket.connect.return_value = None - second_socket.getsockopt.return_value = MIN_SEND_BUFFER_SIZE - mock_socket_create.side_effect = [first_socket, second_socket] - - datadog = DogStatsd( - socket_path="/fake/uds/socket/path", - socket_timeout=0.1, - socket_connect_timeout=1, - ) - datadog.gauge('some value', 1) - datadog.flush() - - self.assertEqual(mock_socket_create.call_count, 2) - first_socket.close.assert_called_once() - second_socket.close.assert_not_called() - second_socket.settimeout.assert_called_with(0.1) - mock_sleep.assert_any_call(UDS_CONNECT_RETRY_INITIAL_BACKOFF) - - @patch('datadog.dogstatsd.base.time.sleep') - @patch('socket.socket') - def test_uds_socket_retries_refused_socket_until_timeout(self, mock_socket_create, mock_sleep): - refused_socket_error = socket.error(errno.ECONNREFUSED, os.strerror(errno.ECONNREFUSED)) - first_socket = Mock() - first_socket.connect.side_effect = refused_socket_error - first_socket.getsockopt.return_value = MIN_SEND_BUFFER_SIZE - second_socket = Mock() - second_socket.connect.return_value = None - second_socket.getsockopt.return_value = MIN_SEND_BUFFER_SIZE - mock_socket_create.side_effect = [first_socket, second_socket] - - datadog = DogStatsd( - socket_path="unixstream:///fake/uds/socket/path", - socket_timeout=0.1, - socket_connect_timeout=1, - ) - datadog.gauge('some value', 1) - datadog.flush() - - self.assertEqual(mock_socket_create.call_count, 2) - first_socket.close.assert_called_once() - second_socket.close.assert_not_called() - second_socket.settimeout.assert_called_with(0.1) - mock_sleep.assert_any_call(UDS_CONNECT_RETRY_INITIAL_BACKOFF) - - @patch('datadog.dogstatsd.base.time.sleep') - @patch('datadog.dogstatsd.base.time.time', side_effect=[0, 0, 0, 0.5, 0.9, 1.1]) @patch('socket.socket') - def test_uds_socket_never_sets_expired_deadline(self, mock_socket_create, mock_time, mock_sleep): + def test_uds_socket_missing_socket_fails_immediately_no_retry(self, mock_socket_create): + # _get_uds_socket makes exactly one connect attempt now: retrying a + # connect failure, if at all, is the caller's job (socket_connect_retry + # for the background sender; never for a direct send). missing_socket_error = socket.error(errno.ENOENT, os.strerror(errno.ENOENT)) mock_socket = mock_socket_create.return_value mock_socket.connect.side_effect = missing_socket_error mock_socket.getsockopt.return_value = MIN_SEND_BUFFER_SIZE with self.assertRaises(socket.error) as raised: - DogStatsd._get_uds_socket("unixgram:///fake/uds/socket/path", 0.1, 1) + DogStatsd._get_uds_socket("unixgram:///fake/uds/socket/path", 0.1) self.assertEqual(raised.exception.errno, errno.ENOENT) mock_socket_create.assert_called_once_with(socket.AF_UNIX, socket.SOCK_DGRAM) - mock_socket.settimeout.assert_called_once_with(1) + mock_socket.settimeout.assert_called_once_with(0.1) + mock_socket.close.assert_called_once() - @patch('datadog.dogstatsd.base.time.time', side_effect=[0, 0.9, 1.1]) @patch('socket.socket') - def test_uds_socket_raises_timeout_before_first_attempt(self, mock_socket_create, mock_time): - with self.assertRaises(socket.timeout) as raised: - DogStatsd._get_uds_socket("unixgram:///fake/uds/socket/path", 0.1, 1) - - self.assertEqual(str(raised.exception), "timed out connecting to UDS socket") - mock_socket_create.assert_not_called() + def test_uds_socket_falls_back_to_the_other_kind_on_eprototype(self, mock_socket_create): + # EPROTOTYPE means "wrong socket kind for this address", not a + # transient failure: _get_uds_socket falls back to the other + # candidate kind for it (still one connect attempt per kind), rather + # than retrying the same kind or giving up. + wrong_kind_socket = Mock() + wrong_kind_socket.connect.side_effect = socket.error(errno.EPROTOTYPE, os.strerror(errno.EPROTOTYPE)) + wrong_kind_socket.getsockopt.return_value = MIN_SEND_BUFFER_SIZE + right_kind_socket = Mock() + right_kind_socket.connect.return_value = None + right_kind_socket.getsockopt.return_value = MIN_SEND_BUFFER_SIZE + mock_socket_create.side_effect = [wrong_kind_socket, right_kind_socket] + + sock = DogStatsd._get_uds_socket("/fake/uds/socket/path", 0.1) + + self.assertIs(sock, right_kind_socket) + mock_socket_create.assert_has_calls([ + call(socket.AF_UNIX, socket.SOCK_DGRAM), + call(socket.AF_UNIX, socket.SOCK_STREAM), + ]) + wrong_kind_socket.close.assert_called_once() @patch('socket.socket') def test_udp_socket_ensures_min_receive_buffer(self, mock_socket_create): @@ -2804,6 +2686,544 @@ def test_sender_queue_no_timeout(self): statsd = DogStatsd(disable_background_sender=False, sender_queue_timeout=None) statsd.stop() + def _call_bounded(self, func, args=(), limit=5.0): + """Call func in a worker thread, failing if it doesn't return in time. + + Everything exercised below exists to *bound* a wait, so a regression + that reintroduces an unbounded wait should surface as a clear failure + rather than hanging the whole suite until CI kills the job. + """ + result = {} + + def run(): + result["value"] = func(*args) + + t = threading.Thread(target=run) + t.daemon = True + t.start() + t.join(limit) + self.assertFalse( + t.is_alive(), + "{} did not return within {}s: timeout not honoured".format(getattr(func, "__name__", func), limit), + ) + return result["value"] + + def test_queue_join_timeout(self): + # join(timeout) must report whether the queue actually drained, and must + # not rely on Condition.wait()'s return value (always None on Python 2). + pending_queue = SenderQueue( + maxsize=0, + expiry_seconds=100.0, + on_drop_queue_full=lambda item: self.fail("unexpected full drop"), + on_drop_expired=lambda item: self.fail("unexpected expiry drop"), + ) + self.assertIs(pending_queue.join(0), True) + self.assertIs(pending_queue.join(), True) + + pending_queue.put(PendingPayload("first\n", sender_queue_clock())) + + # Nothing is draining it, so a bounded join must give up and say so + # rather than blocking forever or claiming success. + t0 = time.time() + self.assertIs(self._call_bounded(pending_queue.join, (0.2,)), False) + self.assertGreaterEqual(time.time() - t0, 0.2) + + # timeout=0 is a non-blocking poll. + t0 = time.time() + self.assertIs(self._call_bounded(pending_queue.join, (0,)), False) + self.assertLess(time.time() - t0, 0.2) + + # Once the payload is accounted for, join() succeeds. + got = pending_queue.get() + pending_queue.task_done(got) + self.assertIs(pending_queue.join(0), True) + + def test_queue_join_timeout_returns_as_soon_as_the_queue_drains(self): + # A generous timeout must not be waited out: join() returns as soon as + # the last task is done. + pending_queue = SenderQueue( + maxsize=0, + expiry_seconds=100.0, + on_drop_queue_full=lambda item: None, + on_drop_expired=lambda item: None, + ) + pending_queue.put(PendingPayload("first\n", sender_queue_clock())) + + def drain(): + time.sleep(0.2) + got = pending_queue.get() + pending_queue.task_done(got) + + t = threading.Thread(target=drain) + t.start() + try: + t0 = time.time() + self.assertIs(pending_queue.join(10.0), True) + elapsed = time.time() - t0 + finally: + t.join(timeout=5.0) + self.assertLess(elapsed, 5.0, "join() should return on drain, not wait out the whole timeout") + + def test_wait_for_pending_timeout(self): + # A queue with no sender thread draining it: wait_for_pending() must + # give up and report False rather than blocking forever. Done without a + # thread on purpose -- the default transport is UDP, where a send + # succeeds even with nothing listening, so "assign no socket" would not + # reliably keep a payload pending. + statsd = DogStatsd(disable_background_sender=True, disable_telemetry=True) + statsd._queue = SenderQueue( + 0, + PENDING_PAYLOAD_EXPIRY_SECONDS, + lambda item: None, + lambda item: None, + ) + statsd._send_to_server("test.metric:1|c") + + t0 = time.time() + self.assertIs(self._call_bounded(statsd.wait_for_pending, (0.2,)), False) + self.assertGreaterEqual(time.time() - t0, 0.2) + + # timeout=0 is a non-blocking poll. + self.assertIs(self._call_bounded(statsd.wait_for_pending, (0,)), False) + + def test_wait_for_pending_returns_true_with_no_queue(self): + # Nothing queued (background sender disabled) is trivially "drained". + statsd = DogStatsd(disable_background_sender=True, disable_telemetry=True) + self.assertIsNone(statsd._queue) + self.assertIs(statsd.wait_for_pending(), True) + self.assertIs(statsd.wait_for_pending(0), True) + + def test_stop_timeout_reports_failure_and_keeps_the_thread_joinable(self): + # A wedged sender must not make stop() hang forever when a timeout is + # given, and stop() must say it didn't finish. + statsd = DogStatsd(disable_background_sender=False, disable_telemetry=True) + release = threading.Event() + wedged = statsd._sender_thread + + # Wedge the sender inside a send so it can't observe Stop. + def blocking_xmit(packet, queue_mode=False): + release.wait(10.0) + return True + + statsd._xmit_packet_with_telemetry = blocking_xmit + statsd._send_to_server("test.metric:1|c") + time.sleep(0.1) # let the sender pick it up and wedge + + try: + t0 = time.time() + self.assertIs(self._call_bounded(statsd.stop, (0.2,)), False) + self.assertGreaterEqual(time.time() - t0, 0.2) + # The handle is retained so the thread isn't lost. + self.assertIs(statsd._sender_thread, wedged) + self.assertTrue(wedged.is_alive()) + + # Unwedge: a second stop() now succeeds and clears the handle. + release.set() + self.assertIs(statsd.stop(5.0), True) + self.assertIsNone(statsd._sender_thread) + finally: + release.set() + wedged.join(timeout=5.0) + + def test_metric_call_racing_stop_does_not_land_behind_the_stop_sentinel(self): + # The actual regression this guards against: a metric call racing + # disable_background_sender() must never land in the queue BEHIND + # the Stop sentinel -- where it would be silently lost once the + # sender reaches Stop and exits without ever seeing it. Rejecting + # new puts is gated on _sender_stopping (set before Stop is + # appended, both under the same lock _send_to_server() re-checks -- + # see _stop_sender_thread), not on self._queue itself, which stays + # non-None until the sender actually finishes with it. + statsd = DogStatsd(disable_background_sender=False, disable_telemetry=True) + release = threading.Event() + entered_first_send = threading.Event() + sent = [] + + def blocking_xmit(packet, queue_mode=False): + sent.append(packet) + entered_first_send.set() + release.wait(10.0) + return True + + statsd._xmit_packet_with_telemetry = blocking_xmit + statsd._send_to_server("first:1|c") + self.assertTrue(entered_first_send.wait(5.0), "sender never picked up the first payload") + + stop_result = {} + second_result = {} + + def call_stop(): + stop_result["value"] = statsd.stop(5.0) + + def call_second(): + # Deliberately started only once the queue is confirmed already + # closed below -- this IS the race under test. + statsd._send_to_server("second:1|c") + second_result["done"] = True + + stopper = threading.Thread(target=call_stop) + stopper.start() + second_thread = None + try: + # _stop_sender_thread sets _sender_stopping (which is what + # _send_to_server() rejects new puts on) well before the drain + # can possibly finish -- it only needs to set a flag and append + # Stop, not wait for the wedged send. Poll briefly rather than a + # fixed sleep, to stay fast without being flaky on a slower + # machine. + signalled = False + deadline = time.time() + 2.0 + while time.time() < deadline: + if statsd._sender_stopping.is_set(): + signalled = True + break + time.sleep(0.01) + self.assertTrue(signalled, "the shutdown signal must be set promptly, without waiting for the drain") + self.assertIsNotNone(statsd._queue, "the real, still-draining queue must remain reachable") + self.assertTrue(stopper.is_alive(), "stop() must still be waiting on the wedged sender") + + second_thread = threading.Thread(target=call_second) + second_thread.start() + # Still wedged behind the first send (same mocked entry point, + # same release Event) -- confirms this really took the direct- + # send fallback rather than silently returning after a queue put. + second_thread.join(0.3) + self.assertTrue(second_thread.is_alive(), "expected the direct send to be wedged behind the first one too") + finally: + release.set() + stopper.join(timeout=10.0) + if second_thread is not None: + second_thread.join(timeout=10.0) + + self.assertIs(stop_result.get("value"), True, "stop() should have completed cleanly") + self.assertTrue(second_result.get("done")) + self.assertEqual( + sent, ["first:1|c\n", "second:1|c\n"], + "the racing payload must have been delivered (via a direct send), not lost behind Stop", + ) + + def test_stop_timeout_is_bounded_when_a_producer_is_wedged_waiting_for_room(self): + # The actual regression this guards against: a producer thread + # blocked inside SenderQueue.put(), waiting for room in a full + # queue, holds _buffer_lock for as long as that wait lasts -- up to + # sender_queue_timeout, or forever if it's None. Without + # SenderQueue.close() waking it immediately, _stop_sender_thread() + # can't even acquire that lock to append Stop, so stop(timeout) + # ignores its own timeout entirely (bounded only by whatever else + # eventually frees the stuck producer -- nothing, in the worst case). + statsd = DogStatsd( + disable_background_sender=False, + disable_telemetry=True, + sender_queue_size=1, + sender_queue_timeout=None, # wait forever for room if not interrupted + ) + release = threading.Event() + entered_send = threading.Event() + + def blocking_xmit(packet, queue_mode=False): + entered_send.set() + release.wait(30.0) + return True + + statsd._xmit_packet_with_telemetry = blocking_xmit + statsd._send_to_server("first:1|c") + self.assertTrue(entered_send.wait(5.0), "sender never picked up first") + # get() already popped "first" off the queue while the sender is + # wedged inside the mocked send -- refill it to maxsize=1 so the + # NEXT put() has nowhere to go. + statsd._send_to_server("second:1|c") + self.assertEqual(statsd._queue.qsize(), 1) + + producer_started = threading.Event() + + def producer(): + producer_started.set() + statsd._send_to_server("third:1|c") + + producer_thread = threading.Thread(target=producer) + producer_thread.start() + self.assertTrue(producer_started.wait(5.0)) + # Give the producer a moment to actually reach put()'s internal wait + # (queue full, sender_queue_timeout=None) before racing stop(). + time.sleep(0.2) + + try: + t0 = time.time() + self.assertIs(self._call_bounded(statsd.stop, (1.0,), limit=3.0), False) + elapsed = time.time() - t0 + self.assertLess(elapsed, 3.0, "stop(1) must not hang behind a producer wedged in put()") + self.assertGreaterEqual(elapsed, 1.0, "stop(1) should still take roughly its own timeout, not return early") + finally: + release.set() + producer_thread.join(timeout=5.0) + + def test_stop_timeout_is_bounded_while_the_sender_holds_the_socket_lock(self): + # The wedge that matters in practice: the sender is parked inside a + # blocking send() and therefore owns _socket_lock. stop()'s own + # close_socket()/flush calls want that same lock, so without care they + # block for as long as the sender stays stuck and the timeout means + # nothing. stop() must still return within its timeout. + statsd = DogStatsd(disable_background_sender=False, disable_telemetry=True) + release = threading.Event() + entered_send = threading.Event() + wedged = statsd._sender_thread + + class BlockingSocket(object): + def send(self, data): + entered_send.set() + release.wait(30.0) + return len(data) + + def sendall(self, data): + return self.send(data) + + def close(self): + pass + + def setblocking(self, *args): + pass + + def settimeout(self, *args): + pass + + def getsockopt(self, *args): + return MIN_SEND_BUFFER_SIZE + + def setsockopt(self, *args): + pass + + statsd.socket = BlockingSocket() + for i in range(5): + statsd._send_to_server("test.metric.{}:1|c".format(i)) + self.assertTrue(entered_send.wait(5.0), "sender never reached send()") + + try: + t0 = time.time() + self.assertIs(self._call_bounded(statsd.stop, (0.2,)), False) + elapsed = time.time() - t0 + self.assertGreaterEqual(elapsed, 0.2) + self.assertLess(elapsed, 5.0, "stop() blocked well past its timeout") + + # The socket was deliberately left alone: closing it under a thread + # that is mid-send is both unsafe and the thing that would block. + self.assertIsNotNone(statsd.socket) + self.assertTrue(wedged.is_alive()) + self.assertIs(statsd._sender_thread, wedged) + finally: + release.set() + wedged.join(timeout=5.0) + + def test_stop_timeout_then_restart_does_not_orphan_the_sender(self): + # After a timed-out stop() the sender is still running and the queue is + # intentionally retained. Re-enabling the background sender must not + # mint a SECOND thread alongside the first (which is what happened when + # the queue was dropped on timeout): the client would leak a thread and + # end up with two senders sharing one socket/lock. + statsd = DogStatsd(disable_background_sender=False, disable_telemetry=True) + release = threading.Event() + entered_send = threading.Event() + wedged = statsd._sender_thread + + def blocking_send(self, data): + entered_send.set() + release.wait(30.0) + return len(data) + + statsd.socket = type("W", (object,), { + "send": blocking_send, + "sendall": blocking_send, + "close": lambda self: None, + "setblocking": lambda self, *a: None, + "settimeout": lambda self, *a: None, + "getsockopt": lambda self, *a: MIN_SEND_BUFFER_SIZE, + "setsockopt": lambda self, *a: None, + })() + statsd._send_to_server("test.metric:1|c") + self.assertTrue(entered_send.wait(5.0)) + + try: + self.assertIs(self._call_bounded(statsd.stop, (0.2,)), False) + self.assertTrue(wedged.is_alive()) + self.assertIsNotNone(statsd._queue, "queue must be retained so state stays coherent") + + statsd.enable_background_sender() + # This identity check IS the proof there's no orphan: if a second + # thread had been minted, _sender_thread would now point at it + # instead of at wedged. (Deliberately not also scanning + # threading.enumerate() for same-named threads process-wide: this + # file has other, unrelated tests that spin up daemon + # "DogStatsd_sender_thread"s of their own, so a global count is not + # a property this test can own.) + self.assertIs( + statsd._sender_thread, wedged, + "re-enabling after a timed-out stop must not start a second sender thread", + ) + self.assertTrue(wedged.is_alive()) + finally: + release.set() + wedged.join(timeout=5.0) + + def test_stop_timeout_then_unwedge_then_restart_starts_fresh(self): + # Once the timed-out thread finally drains and exits, the client + # self-heals: the next stop() reports success and cleans up, and a + # re-enabled background sender is a brand new thread over a new queue. + statsd = DogStatsd(disable_background_sender=False, disable_telemetry=True) + release = threading.Event() + entered_send = threading.Event() + wedged = statsd._sender_thread + + def blocking_send(self, data): + entered_send.set() + release.wait(30.0) + return len(data) + + statsd.socket = type("W", (object,), { + "send": blocking_send, + "sendall": blocking_send, + "close": lambda self: None, + "setblocking": lambda self, *a: None, + "settimeout": lambda self, *a: None, + "getsockopt": lambda self, *a: MIN_SEND_BUFFER_SIZE, + "setsockopt": lambda self, *a: None, + })() + statsd._send_to_server("test.metric:1|c") + self.assertTrue(entered_send.wait(5.0)) + + try: + self.assertIs(self._call_bounded(statsd.stop, (0.2,)), False) + release.set() # let the wedged send finish; the thread drains and exits + self.assertIs(self._call_bounded(statsd.stop, (5.0,)), True) + self.assertIsNone(statsd._queue) + self.assertIsNone(statsd._sender_thread) + + statsd.enable_background_sender() + self.assertIsNotNone(statsd._sender_thread) + self.assertIsNot(statsd._sender_thread, wedged, "a fresh thread should start") + self.assertTrue(statsd._sender_thread.is_alive()) + self.assertFalse(wedged.is_alive(), "the old thread must have actually exited, not just been forgotten") + finally: + release.set() + # The fresh thread's queue is empty and has never been sent a Stop + # sentinel, so joining statsd._sender_thread directly would just + # time out waiting on a get() that never returns -- leaking the + # thread into later tests. stop() sends Stop and joins correctly + # regardless of which thread/queue is current. + statsd.stop(5.0) + + def test_sender_self_heals_when_the_timed_out_thread_finishes_on_its_own(self): + # The timed-out thread can exit without a second stop() ever being + # called. The client must detect that and allow a fresh sender, rather + # than keeping the dead thread / stale queue around and silently + # refusing to start a new one (metrics would then queue up and drop). + statsd = DogStatsd(disable_background_sender=False, disable_telemetry=True) + release = threading.Event() + entered_send = threading.Event() + wedged = statsd._sender_thread + + def blocking_send(self, data): + entered_send.set() + release.wait(30.0) + return len(data) + + statsd.socket = type("W", (object,), { + "send": blocking_send, + "sendall": blocking_send, + "close": lambda self: None, + "setblocking": lambda self, *a: None, + "settimeout": lambda self, *a: None, + "getsockopt": lambda self, *a: MIN_SEND_BUFFER_SIZE, + "setsockopt": lambda self, *a: None, + })() + statsd._send_to_server("test.metric:1|c") + self.assertTrue(entered_send.wait(5.0)) + + try: + self.assertIs(self._call_bounded(statsd.stop, (0.2,)), False) + self.assertIsNotNone(statsd._queue) + self.assertIs(statsd._sender_thread, wedged) + + # Caller moved on and never called stop() again; the wedged send + # eventually completes and the thread drains and exits on its own. + release.set() + wedged.join(timeout=5.0) + self.assertFalse(wedged.is_alive()) + + # The exit must have healed the client state so a re-enable works. + self.assertIsNone(statsd._queue) + self.assertIsNone(statsd._sender_thread) + statsd.enable_background_sender() + self.assertIsNot( + statsd._sender_thread, wedged, + "a timed-out thread that finished on its own must not block a restart", + ) + self.assertTrue(statsd._sender_thread.is_alive()) + finally: + release.set() + # Same reason as the sibling test above: the fresh thread started + # by enable_background_sender() has an empty queue and was never + # sent Stop, so joining it directly would time out and leak it. + statsd.stop(5.0) + + def test_wait_for_pending_is_honest_after_a_stop_timeout(self): + # After a timed-out stop(), the queue is retained (not shown as empty), + # so wait_for_pending() still reports the truth -- False while the + # sender is draining, True once it has actually drained -- rather than + # claiming "drained" while a live thread is still sending. + statsd = DogStatsd(disable_background_sender=False, disable_telemetry=True) + release = threading.Event() + entered_send = threading.Event() + wedged = statsd._sender_thread + + def blocking_send(self, data): + entered_send.set() + release.wait(30.0) + return len(data) + + statsd.socket = type("W", (object,), { + "send": blocking_send, + "sendall": blocking_send, + "close": lambda self: None, + "setblocking": lambda self, *a: None, + "settimeout": lambda self, *a: None, + "getsockopt": lambda self, *a: MIN_SEND_BUFFER_SIZE, + "setsockopt": lambda self, *a: None, + })() + statsd._send_to_server("test.metric:1|c") + self.assertTrue(entered_send.wait(5.0)) + + try: + self.assertIs(self._call_bounded(statsd.stop, (0.2,)), False) + self.assertIs( + self._call_bounded(statsd.wait_for_pending, (0.2,)), False, + "wait_for_pending() must not claim drained while the sender is still running", + ) + + release.set() # sender finishes the send, hits Stop, drains, exits + self.assertIs(self._call_bounded(statsd.wait_for_pending, (5.0,)), True) + finally: + release.set() + wedged.join(timeout=5.0) + + def test_stop_and_wait_for_pending_default_to_waiting_forever(self): + # The default must stay unbounded: a slow-but-progressing sender is + # waited out completely, with nothing left pending. + statsd = DogStatsd(disable_background_sender=False, disable_telemetry=True) + sent = [] + + def slow_xmit(packet, queue_mode=False): + time.sleep(0.05) + sent.append(packet) + return True + + statsd._xmit_packet_with_telemetry = slow_xmit + for i in range(5): + statsd._send_to_server("test.metric.{}:1|c".format(i)) + + self.assertIs(statsd.wait_for_pending(), True) + self.assertEqual(len(sent), 5, "unbounded wait_for_pending() must drain everything") + self.assertIs(statsd.stop(), True) + self.assertIsNone(statsd._sender_thread) + def test_sender_queue_timeout_blocks_the_calling_thread_through_the_client(self): # End-to-end: sender_queue_timeout configured on the real client # actually makes statsd.increment() (the calling/application thread) @@ -3478,13 +3898,14 @@ def capture_put(item): statsd.stop() def test_connection_failure_requeues_and_resends_once_reconnected(self): - # A UDS client whose socket is broken, with a small connect budget so - # the internal reconnect-and-retry loop inside _xmit_packet gives up - # quickly and hands off to the sender queue's own retry-by-requeuing. + # A UDS client whose socket is broken. socket_connect_retry=True is + # what makes queue-mode sends hand off a connection failure to the + # sender queue's own retry-by-requeuing instead of dropping the + # payload on the first attempt -- direct sends never do this. working_socket = FakeSocket() attempts = {"count": 0} - def flaky_get_uds_socket(cls, socket_path, timeout, connect_timeout): + def flaky_get_uds_socket(cls, socket_path, timeout): attempts["count"] += 1 if attempts["count"] < 4: raise socket.error(errno.ECONNREFUSED, "still refused") @@ -3495,8 +3916,8 @@ def flaky_get_uds_socket(cls, socket_path, timeout, connect_timeout): socket_path="/tmp/dogstatsd-test-requeue.sock", disable_telemetry=True, disable_background_sender=False, + socket_connect_retry=True, ) - statsd.socket_connect_timeout = 0.05 statsd.gauge("eventually.sent", 1) statsd.wait_for_pending() @@ -3511,6 +3932,176 @@ def flaky_get_uds_socket(cls, socket_path, timeout, connect_timeout): statsd.stop() + def test_stop_timeout_retries_until_the_deadline_before_abandoning_a_payload(self): + # The actual regression this guards against: stop(timeout) must keep + # retrying a connection failure for close to the requested timeout -- + # not abandon on the very first backoff check -- so a payload that + # would have succeeded on a later attempt is not silently lost. + working_socket = FakeSocket() + attempts = {"count": 0} + + def flaky_get_uds_socket(cls, socket_path, timeout): + attempts["count"] += 1 + if attempts["count"] < 4: + raise socket.error(errno.ECONNREFUSED, "still refused") + return working_socket + + with mock.patch.object(DogStatsd, "_get_uds_socket", classmethod(flaky_get_uds_socket)), \ + patch("datadog.dogstatsd.base.SENDER_RETRY_INITIAL_BACKOFF", 0.05): + statsd = DogStatsd( + socket_path="/tmp/dogstatsd-test-stop-retries.sock", + disable_telemetry=True, + disable_background_sender=False, + socket_connect_retry=True, + ) + + statsd.gauge("eventually.sent", 1) + time.sleep(0.05) # let the first attempt fail and start retrying + + # A generous bounded timeout: comfortably enough for a few fast + # (patched-down) backoff doublings to reach the 4th, working + # attempt, but nowhere near SENDER_UNBOUNDED_STOP_GRACE_SECONDS. + # Before the fix, stop() abandoned on the very first backoff + # check regardless of this value, so this would have failed: + # returning near-instantly with the payload never sent. + self.assertIs(self._call_bounded(statsd.stop, (3.0,), limit=5.0), True) + + self.assertGreaterEqual(attempts["count"], 4) + self.assertEqual(statsd.packets_dropped_writer, 0, "the payload must not have been abandoned") + self.assertTrue(working_socket.payloads[0].decode("utf-8").startswith("eventually.sent:1|g")) + + def test_queue_mode_drops_immediately_by_default(self): + # socket_connect_retry defaults to False for the background sender + # too: without it, a connection failure on a queued payload is + # accounted as a writer drop on the very first attempt, exactly like + # a direct send, rather than requeued and retried. + statsd = DogStatsd( + socket_path="/tmp/dogstatsd-test-no-retry-queue.sock", + disable_background_sender=False, + ) + self.assertFalse(statsd.socket_connect_retry) + statsd.socket = BrokenSocket(error_number=errno.ECONNREFUSED) + + statsd.gauge("dropped.immediately", 1) + statsd.wait_for_pending() + + self.assertGreaterEqual(statsd.packets_dropped_writer, 1) + self.assertEqual(statsd.packets_dropped_expired, 0) + statsd.stop() + + def test_queue_mode_retry_backoff_caps_at_one_minute(self): + # "Backoff...to a maximum of once per minute": the doubling algorithm + # in _sender_main_loop is unchanged (already exercised end-to-end by + # test_connection_failure_requeues_and_resends_once_reconnected, + # which drives it through several real doublings below the cap); what + # changed is this constant, from the old 1-second cap to 60. + self.assertEqual(SENDER_RETRY_MAX_BACKOFF, 60.0) + + def _unreachable_retrying_client(self): + """Background-sender client whose UDS path will never exist. + + socket_connect_retry=True, so every queued payload fails to connect + and gets requeued at the FRONT of the queue -- ahead of any Stop + sentinel. + """ + sock_path = os.path.join(tempfile.mkdtemp(), "never-created.sock") + return DogStatsd( + socket_path="unix://" + sock_path, + socket_connect_retry=True, + disable_background_sender=False, + disable_telemetry=True, + # Keep these out of the module-level _instances WeakSet: the + # pre_fork test below abandons its instance with _config_lock + # still held, and the global os.register_at_fork hooks iterate + # _instances -- a later real fork() would deadlock on it. + track_instance=False, + ) + + def test_stop_is_bounded_with_a_stuck_replay_safe_payload(self): + # requeue_front() puts a failed payload back ahead of the Stop + # sentinel, and replay-safe payloads are exempt from the queue's + # expiry, so such a payload never resolves while the Agent is down. + # Without a bounded grace period the sender never reaches Stop at all + # and stop() hangs forever. A small patched grace makes the test fast + # and deterministic instead of depending on the real 60s default. + statsd = self._unreachable_retrying_client() + statsd.gauge_with_timestamp("replay.safe", 1, timestamp=int(time.time())) + time.sleep(0.1) # let the sender pick it up and start retrying + + with patch("datadog.dogstatsd.base.SENDER_UNBOUNDED_STOP_GRACE_SECONDS", 0.3): + t0 = time.time() + self.assertIs(self._call_bounded(statsd.stop, ()), False, "the payload was never delivered") + self.assertLess(time.time() - t0, 5.0, "stop() did not bound the retry loop") + self.assertIsNone(statsd._queue) + self.assertEqual(statsd.packets_dropped_writer, 1, "the abandoned payload must be accounted for") + + def test_pre_fork_is_bounded_with_a_stuck_replay_safe_payload(self): + # Same starvation, reached through pre_fork() -- which matters more: + # it is installed as an os.register_at_fork hook, has no timeout + # parameter, and would therefore block any fork() in the process. + statsd = self._unreachable_retrying_client() + statsd.gauge_with_timestamp("replay.safe", 1, timestamp=int(time.time())) + time.sleep(0.1) + + # Run it on a worker so a regression fails this assertion instead of + # hanging the suite. That means _config_lock -- which pre_fork() + # deliberately acquires and leaves held for post_fork_parent() to + # release -- ends up owned by a thread that then exits, so this + # instance is deliberately abandoned rather than restored. It is + # constructed with track_instance=False precisely so nothing else can + # ever try to take that lock again. + with patch("datadog.dogstatsd.base.SENDER_UNBOUNDED_STOP_GRACE_SECONDS", 0.3): + t0 = time.time() + self._call_bounded(statsd.pre_fork, ()) + self.assertLess(time.time() - t0, 5.0, "pre_fork() would have blocked os.fork()") + self.assertIsNone(statsd._sender_thread, "pre_fork() should have stopped the sender") + self.assertEqual(statsd.packets_dropped_writer, 1, "the abandoned payload must be accounted for") + + def test_stop_interrupts_a_long_retry_backoff_instead_of_waiting_it_out(self): + # The backoff cap is a full minute. A plain time.sleep() would make + # stop() wait out however much of it is left; the shutdown signal must + # cut it short well before its own (also patched down) grace period. + statsd = self._unreachable_retrying_client() + with patch("datadog.dogstatsd.base.SENDER_RETRY_INITIAL_BACKOFF", 30.0), \ + patch("datadog.dogstatsd.base.SENDER_UNBOUNDED_STOP_GRACE_SECONDS", 0.3): + statsd.gauge("ordinary", 1) + time.sleep(0.3) # fail once, then settle into the 30s backoff + + t0 = time.time() + self.assertIs(self._call_bounded(statsd.stop, ()), False, "the payload was never delivered") + self.assertLess(time.time() - t0, 5.0, "stop() waited out the backoff sleep") + + def test_sender_can_restart_after_a_stop_cleared_the_stopping_signal(self): + # The shutdown signal is sticky, so it has to be cleared when a new + # sender starts or the replacement would exit immediately. + statsd = self._unreachable_retrying_client() + statsd.gauge("ordinary", 1) + time.sleep(0.1) + with patch("datadog.dogstatsd.base.SENDER_UNBOUNDED_STOP_GRACE_SECONDS", 0.3): + self.assertIs(self._call_bounded(statsd.stop, ()), False, "the payload was never delivered") + + statsd.enable_background_sender() + try: + self.assertIsNotNone(statsd._queue) + fresh = statsd._sender_thread + self.assertIsNotNone(fresh) + + # The stale signal only bites once the sender reaches its retry + # branch, so it has to actually fail a send here: an idle sender + # just blocks in get() and would mask the bug. With the signal + # left set, this first failure looks like a shutdown request and + # the fresh sender exits silently on it. + statsd.gauge("ordinary", 1) + time.sleep(0.3) + self.assertTrue(fresh.is_alive(), "fresh sender exited on a stale stopping signal") + self.assertIs(statsd._sender_thread, fresh) + finally: + # Bounded cleanup: the fresh sender is genuinely retrying against + # a socket path that will never exist, so a bounded stop() now + # legitimately waits close to its own timeout before giving up + # (see the fix above) -- keep it short so this cleanup stays fast. + statsd.stop(1.0) + def test_set_socket_timeout(self): statsd = DogStatsd(disable_background_sender=False) statsd.socket = FakeSocket()