diff --git a/CHANGELOG.md b/CHANGELOG.md index 994b55951..8455b91a1 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -1,5 +1,9 @@ # Changelog +## 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). + ## v0.53.0 / 2026-07-24 * [Fixed] Add DD_DOGSTATSD_URL support for Unix and UDP URLs. See [#968](https://github.com/DataDog/datadogpy/pull/968). diff --git a/datadog/dogstatsd/base.py b/datadog/dogstatsd/base.py index 5922e1183..e622c9a1e 100644 --- a/datadog/dogstatsd/base.py +++ b/datadog/dogstatsd/base.py @@ -23,16 +23,12 @@ if sys.version_info[:2] >= (3, 5): from typing import TYPE_CHECKING # noqa: F401 -try: - import queue -except ImportError: - # pypy has the same module, but capitalized. - import Queue as queue # type: ignore[no-redef] - # pylint: disable=unused-import if sys.version_info[:2] >= (3, 5): - from typing import Any, Optional, List, Text, Tuple, Type, Union, Iterable, Callable, overload # noqa: F401 + from typing import ( # noqa: F401 + Any, Callable, Iterable, List, Optional, Text, Tuple, Type, Union, overload, + ) try: from typing import SupportsIndex @@ -49,7 +45,16 @@ ) from datadog.dogstatsd.route import get_default_route from datadog.dogstatsd.container import Cgroup -from datadog.util.compat import text, urlparse +from datadog.dogstatsd.sender_queue import ( + SenderQueue, + PendingPayload, + Stop, + payload_text, +) + +if sys.version_info[:2] >= (3, 5): + from datadog.dogstatsd.sender_queue import QueuedItem # noqa: F401 +from datadog.util.compat import monotonic, text, urlparse from datadog.util.format import normalize_tags, validate_cardinality from datadog.version import __version__ @@ -183,6 +188,14 @@ def reverse(self): 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. +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 # reconnect-and-resend attempt instead of dropping the packet outright. @@ -224,8 +237,6 @@ def reverse(self): ] ) + "\n" -Stop = object() - SUPPORTS_FORKING = hasattr(os, "register_at_fork") and not os.environ.get("DD_DOGSTATSD_DISABLE_FORK_SUPPORT", None) TRACK_INSTANCES = not os.environ.get("DD_DOGSTATSD_DISABLE_INSTANCE_TRACKING", None) @@ -306,6 +317,7 @@ def __init__( disable_background_sender=True, # type: bool sender_queue_size=0, # type: int 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] ): # type: (...) -> None @@ -490,8 +502,8 @@ def __init__( Default: True. :type disable_background_sender: boolean - :param sender_queue_size: Set the maximum number of packets to queue for the sender. Optional - How may packets to queue before blocking or dropping the packet if the packet queue is already full. + :param sender_queue_size: Set the maximum number of packets to queue for the sender. Optional. + How many packets to queue before blocking or dropping the packet if the packet queue is already full. Default: 0 (unlimited). :type sender_queue_size: integer @@ -502,6 +514,13 @@ def __init__( Default: 0 (no wait) :type sender_queue_timeout: float + :param sender_queue_expiry_seconds: How long, in seconds, a non-replay-safe payload + may sit in the 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 until they can be sent. + Default: PENDING_PAYLOAD_EXPIRY_SECONDS (10.0). + :type sender_queue_expiry_seconds: float + :param track_instance: Keep track of this instance and automatically handle cleanup when os.fork() is called, if supported. Default: True. @@ -608,8 +627,6 @@ def __init__( self._telemetry = not disable_telemetry self._last_flush_time = time.time() - self._current_buffer_total_size = 0 - self._buffer = [] # type: List[Text] self._buffer_lock = RLock() self._reset_buffer() @@ -634,12 +651,14 @@ def __init__( else: log.debug("Statsd buffering and aggregation is disabled") - self._queue = None # type: Optional[queue.Queue[Union[str, object]]] + self._queue = None # type: Optional[SenderQueue] self._sender_thread = None # type: Optional[threading.Thread] self._sender_enabled = False if not disable_background_sender: - self.enable_background_sender(sender_queue_size, sender_queue_timeout) + self.enable_background_sender( + sender_queue_size, sender_queue_timeout, sender_queue_expiry_seconds + ) if TRACK_INSTANCES and track_instance: _instances.add(self) @@ -703,8 +722,13 @@ def telemetry_socket(self, t_socket): log.info("Unexpected telemetry socket provided with no support for getsockopt") self._telemetry_socket_kind = None - def enable_background_sender(self, sender_queue_size=0, sender_queue_timeout=0): - # type: (int, Optional[float]) -> None + def enable_background_sender( + self, + sender_queue_size=0, + sender_queue_timeout=0, + sender_queue_expiry_seconds=PENDING_PAYLOAD_EXPIRY_SECONDS, + ): + # type: (int, Optional[float], float) -> None """ Use a background thread to communicate with the dogstatsd server. When enabled, a background thread will be used to send metric payloads to the Agent. @@ -724,17 +748,18 @@ def enable_background_sender(self, sender_queue_size=0, sender_queue_timeout=0): If set to None, wait forever. If set to zero drop the packet immediately if the queue is full. Default: 0 (no wait). :type sender_queue_timeout: float, optional + :param sender_queue_expiry_seconds: How long, in seconds, a non-replay-safe payload may sit in the + sender queue before it's considered stale and dropped instead of sent. Replay-safe payloads + (those carrying their own explicit timestamp) are exempt and are kept until they can be sent. + Default: PENDING_PAYLOAD_EXPIRY_SECONDS (10.0). + :type sender_queue_expiry_seconds: float, optional """ with self._config_lock: self._sender_enabled = True self._sender_queue_size = sender_queue_size - if sender_queue_timeout is None: - self._queue_blocking = True - self._queue_timeout = None - else: - self._queue_blocking = sender_queue_timeout > 0 - self._queue_timeout = max(0, sender_queue_timeout) + self._sender_queue_timeout = sender_queue_timeout + self._sender_queue_expiry_seconds = sender_queue_expiry_seconds self._start_sender_thread() @@ -1167,23 +1192,60 @@ def close_buffer(self): def _reset_buffer(self): # type: () -> None with self._buffer_lock: - self._current_buffer_total_size = 0 - self._buffer = [] + # Buffered lines are kept in two separate batches, one per expiry + # policy: replay-safe metrics carry their own timestamp, the rest + # are stamped on receipt. + self._buffer = [] # type: List[Text] # non-replay-safe + self._buffer_rs = [] # type: List[Text] # replay-safe + # Running packet size per buffer, each including the newline that + # will join its lines, so both stay under _max_payload_size + # independently. + self._buffer_size = 0 # type: int + self._buffer_rs_size = 0 # type: int def flush(self): # type: () -> None self.flush_buffered_metrics() + def _flush_one_buffer(self, replay_safe): + # type: (bool) -> None + """Flush just the batch holding lines with the given expiry policy. + + Caller must hold self._buffer_lock (an RLock, so a re-entrant flush + from _send_to_buffer() is fine). Only the named batch is touched: the + other one keeps accumulating, which is the whole point of splitting + them. + """ + if replay_safe: + lines = self._buffer_rs + if not lines: + return + self._buffer_rs = [] + self._buffer_rs_size = 0 + else: + lines = self._buffer + if not lines: + return + self._buffer = [] + self._buffer_size = 0 + self._send_to_server("\n".join(lines), replay_safe) + def flush_buffered_metrics(self): # type: () -> None """ Flush the metrics buffer by sending the data to the server. + + Emits up to two packets, one per expiry policy. Lines keep their + relative order within each packet, but ordering *between* the two is + not preserved: replay-safe payloads carry their own explicit + timestamp, and non-replay-safe ones are timestamped on receipt, so + neither one's meaning depends on where the other lands. """ with self._buffer_lock: - # Only send packets if there are packets to send - if self._buffer: - self._send_to_server("\n".join(self._buffer)) - self._reset_buffer() + # Non-replay-safe first: it's the only batch subject to staleness + # expiry, so give it the earliest queue position. + self._flush_one_buffer(False) + self._flush_one_buffer(True) def flush_aggregated_metrics(self): # type: () -> None @@ -1569,8 +1631,13 @@ def _report(self, metric, metric_type, value, tags, sample_rate, timestamp=0, sa metric, metric_type, value, tags, sample_rate, timestamp, cardinality ) + # A metric carrying its own explicit timestamp is replay-safe: sending + # it late (e.g. after sitting in the background sender queue) doesn't + # change what it means. + replay_safe = timestamp > 0 + # Send it - self._send(payload) + self._send(payload, replay_safe) def _reset_telemetry(self): # type: () -> None @@ -1580,21 +1647,35 @@ def _reset_telemetry(self): self.bytes_sent = 0 self.bytes_dropped_queue = 0 self.bytes_dropped_writer = 0 + self.bytes_dropped_expired = 0 self.packets_sent = 0 self.packets_dropped_queue = 0 self.packets_dropped_writer = 0 + self.packets_dropped_expired = 0 self._last_flush_time = time.time() # Aliases for backwards compatibility. @property def packets_dropped(self): # type: () -> int - return self.packets_dropped_queue + self.packets_dropped_writer + return self.packets_dropped_queue + self.packets_dropped_writer + self.packets_dropped_expired @property def bytes_dropped(self): # type: () -> int - return self.bytes_dropped_queue + self.bytes_dropped_writer + return self.bytes_dropped_queue + self.bytes_dropped_writer + self.bytes_dropped_expired + + def _account_dropped_queue_full(self, item): + # type: (QueuedItem) -> None + """A payload was evicted from the sender queue to make room for a new one.""" + self.packets_dropped_queue += 1 + self.bytes_dropped_queue += len(payload_text(item).encode(self.encoding)) + + def _account_dropped_expired(self, item): + # type: (QueuedItem) -> None + """A payload sat in the sender queue longer than the configured expiry (sender_queue_expiry_seconds).""" + self.packets_dropped_expired += 1 + self.bytes_dropped_expired += len(payload_text(item).encode(self.encoding)) def _flush_telemetry(self): # type: () -> str @@ -1603,6 +1684,9 @@ def _flush_telemetry(self): tags.extend(self.constant_tags) telemetry_tags = ",".join(tags) + bytes_dropped_queue = self.bytes_dropped_queue + self.bytes_dropped_expired + packets_dropped_queue = self.packets_dropped_queue + self.packets_dropped_expired + return TELEMETRY_FORMATTING_STR % ( self.metrics_count, telemetry_tags, @@ -1612,17 +1696,17 @@ def _flush_telemetry(self): telemetry_tags, self.bytes_sent, telemetry_tags, - self.bytes_dropped_queue + self.bytes_dropped_writer, + bytes_dropped_queue + self.bytes_dropped_writer, telemetry_tags, - self.bytes_dropped_queue, + bytes_dropped_queue, telemetry_tags, self.bytes_dropped_writer, telemetry_tags, self.packets_sent, telemetry_tags, - self.packets_dropped_queue + self.packets_dropped_writer, + packets_dropped_queue + self.packets_dropped_writer, telemetry_tags, - self.packets_dropped_queue, + packets_dropped_queue, telemetry_tags, self.packets_dropped_writer, telemetry_tags, @@ -1633,26 +1717,38 @@ def _is_telemetry_flush_time(self): return self._telemetry and \ self._last_flush_time + self._telemetry_flush_interval < time.time() - def _send_to_server(self, packet): - # type: (str) -> None + 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. with self._buffer_lock: packet_with_newline = packet + '\n' if self._queue is not None: - try: - self._queue.put(packet_with_newline, self._queue_blocking, self._queue_timeout) - except queue.Full: - self.packets_dropped_queue += 1 - self.bytes_dropped_queue += len(packet_with_newline.encode(self.encoding)) + if replay_safe: + # Never expires, so it needs no enqueued_at and no + # wrapper at all: queue the bare string and let the + # queue infer replay-safety from the type. Saves the + # PendingPayload object (~56 bytes) per entry and + # keeps these out of the cyclic GC's traversal set. + self._queue.put(packet_with_newline) + else: + self._queue.put(PendingPayload(packet_with_newline, monotonic())) return self._xmit_packet_with_telemetry(packet + '\n') - def _xmit_packet_with_telemetry(self, packet): - # type: (str) -> None - self._xmit_packet(packet, False) + def _xmit_packet_with_telemetry(self, packet, queue_mode=False): + # type: (str, bool) -> Optional[bool] + """Send one packet, optionally piggy-backing a telemetry flush. + + :param queue_mode: True when called from the background sender + thread on behalf of a queued PendingPayload. In that mode, a + connection failure is reported back as None (rather than being + accounted for and dropped) so the caller can requeue the payload + and retry once reconnected, instead of losing it. + """ + sent = self._xmit_packet(packet, False, queue_mode=queue_mode) if self._is_telemetry_flush_time(): telemetry = self._flush_telemetry() @@ -1666,6 +1762,8 @@ def _xmit_packet_with_telemetry(self, packet): self.bytes_dropped_writer += len(telemetry) self.packets_dropped_writer += 1 + return sent + def _installed_socket(self, is_telemetry): # type: (bool) -> Optional[_Socket] """ @@ -1680,8 +1778,16 @@ def _installed_socket(self, is_telemetry): return self.telemetry_socket return self.socket - def _xmit_packet(self, packet, is_telemetry): - # type: (str, bool) -> bool + 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. + + 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. + """ if is_telemetry and self._dedicated_telemetry_destination(): uses_uds = self.telemetry_socket_path is not None @@ -1693,6 +1799,7 @@ def _xmit_packet(self, packet, is_telemetry): retry_deadline = time.time() + self.socket_connect_timeout 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 ( @@ -1704,6 +1811,7 @@ def _xmit_packet(self, packet, is_telemetry): "Gave up reconnecting after socket_connect_timeout (%ss), dropping the packet", self.socket_connect_timeout, ) + sent = None break sent = self._xmit_packet_attempt( @@ -1730,6 +1838,12 @@ def _xmit_packet(self, packet, is_telemetry): time.sleep(min(backoff, remaining)) backoff = min(backoff * 2, UDS_CONNECT_RETRY_MAX_BACKOFF) + if sent is None and queue_mode: + # Connection trouble, and the caller is the background sender + # queue: let it requeue the payload and retry once reconnected, + # instead of dropping it here. + return None + if not is_telemetry and self._telemetry: self.bytes_dropped_writer += len(packet) self.packets_dropped_writer += 1 @@ -1836,22 +1950,43 @@ def _xmit_packet_attempt(self, packet, is_telemetry, retry_eligible, retry_deadl return False - def _send_to_buffer(self, packet): - # type: (str) -> None + def _send_to_buffer(self, packet, replay_safe=False): + # type: (str, bool) -> None + """Append one serialized line to the batch matching its expiry policy. + + Deliberately written out per branch rather than indexing a dict by + replay_safe, and with the size check inlined rather than delegated to + _should_flush(): both cost real CPU at once-per-metric frequency. See + _reset_buffer() for the measurements. The bool() coercion the dict form + needed is gone too -- the branch treats any truthy value correctly. + """ with self._buffer_lock: - if self._should_flush(len(packet)): - self.flush_buffered_metrics() + # Length including the newline that will join this line to the + # next, so the running total anticipates the final packet size. + length = len(packet) + 1 + + if replay_safe: + if self._buffer_rs_size + length > self._max_payload_size: + self._flush_one_buffer(True) + self._buffer_rs.append(packet) + self._buffer_rs_size += length + else: + if self._buffer_size + length > self._max_payload_size: + self._flush_one_buffer(False) + self._buffer.append(packet) + self._buffer_size += length - self._buffer.append(packet) - # Update the current buffer length, including line break to anticipate - # the final packet size - self._current_buffer_total_size += len(packet) + 1 + def _should_flush(self, length_to_be_added, replay_safe=False): + # type: (int, bool) -> bool + """Whether adding a line of this length would overflow its batch. - def _should_flush(self, length_to_be_added): - # type: (int) -> bool - if self._current_buffer_total_size + length_to_be_added + 1 > self._max_payload_size: - return True - return False + NOT used by _send_to_buffer(), which inlines the same comparison to + keep a function call off the per-metric path. Retained because it is + part of the pre-existing surface and is convenient in tests; keep the + two in step if either changes. + """ + current = self._buffer_rs_size if replay_safe else self._buffer_size + return current + length_to_be_added + 1 > self._max_payload_size @staticmethod def _escape_event_content(string): @@ -1938,7 +2073,9 @@ def event( if self._telemetry: self.events_count += 1 - self._send(string) + # An event carrying its own explicit date_happened is replay-safe: + # sending it late doesn't change what it means. + self._send(string, replay_safe=bool(date_happened)) def service_check( self, @@ -1984,7 +2121,9 @@ def service_check( if self._telemetry: self.service_checks_count += 1 - self._send(string) + # A service check carrying its own explicit timestamp is replay-safe: + # sending it late doesn't change what it means. + self._send(string, replay_safe=bool(timestamp)) @staticmethod def _normalize_and_join_tags(tags): @@ -2062,7 +2201,13 @@ def _start_sender_thread(self): if self._queue is not None: return - self._queue = queue.Queue(self._sender_queue_size) + self._queue = SenderQueue( + self._sender_queue_size, + self._sender_queue_expiry_seconds, + self._account_dropped_queue_full, + self._account_dropped_expired, + put_timeout=self._sender_queue_timeout, + ) log.debug("Starting background sender thread") self._sender_thread = threading.Thread( @@ -2086,19 +2231,37 @@ def _stop_sender_thread(self): self._sender_thread.join() self._sender_thread = None - def _sender_main_loop(self, queue): - # type: (queue.Queue[Union[str, object]]) -> None + def _sender_main_loop(self, pending_queue): + # type: (SenderQueue) -> None + backoff = UDS_CONNECT_RETRY_INITIAL_BACKOFF while True: - item = queue.get() + item = pending_queue.get() if item is Stop: - queue.task_done() + pending_queue.task_done(item) return - # next line has type ignore because the type checker cannot - # know that 'if item is Stop' is the only case where item is - # of object type. - self._xmit_packet_with_telemetry(item) # type: ignore[arg-type] # noqa: F821 - queue.task_done() + # payload_text() also narrows the type: 'if item is Stop' above is + # the only case where item is the bare object sentinel, which the + # type checker can't know on its own. + sent = self._xmit_packet_with_telemetry( + payload_text(item), queue_mode=True # type: ignore[arg-type] + ) + + if sent is None: + # Connection trouble: 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) + 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 def wait_for_pending(self): # type: () -> None diff --git a/datadog/dogstatsd/context.py b/datadog/dogstatsd/context.py index ec27b30dd..3e2df8d86 100644 --- a/datadog/dogstatsd/context.py +++ b/datadog/dogstatsd/context.py @@ -6,14 +6,9 @@ import sys -try: - from time import monotonic # type: ignore[attr-defined] -except ImportError: - from time import time as monotonic - # datadog from datadog.dogstatsd.context_async import _get_wrapped_co -from datadog.util.compat import iscoroutinefunction +from datadog.util.compat import iscoroutinefunction, monotonic if sys.version_info[:2] >= (3, 5): diff --git a/datadog/dogstatsd/context_async.py b/datadog/dogstatsd/context_async.py index 7cb0f826c..af976f3d3 100644 --- a/datadog/dogstatsd/context_async.py +++ b/datadog/dogstatsd/context_async.py @@ -23,10 +23,8 @@ # https://github.com/python/mypy/issues/6897 ASYNC_SOURCE = r''' from functools import wraps -try: - from time import monotonic -except ImportError: - from time import time as monotonic + +from datadog.util.compat import monotonic def _get_wrapped_co(self, func): diff --git a/datadog/dogstatsd/sender_queue.py b/datadog/dogstatsd/sender_queue.py new file mode 100644 index 000000000..f6f699e76 --- /dev/null +++ b/datadog/dogstatsd/sender_queue.py @@ -0,0 +1,361 @@ +import collections +import logging +import sys +import threading + +from datadog.util.compat import monotonic + +log = logging.getLogger("datadog.dogstatsd") + +if sys.version_info[:2] >= (3, 5): + from typing import Callable, Dict, Optional, Union # noqa: F401 + + +# Sentinel telling the background sender thread to shut down. +Stop = object() + +# What the queue can hold. A payload is either a bare string (replay-safe, no +# expiry state needed) or a PendingPayload (subject to expiry); Stop is the +# only other thing that ever goes in, and is matched by identity. +if sys.version_info[:2] >= (3, 5): + QueuedItem = Union[str, "PendingPayload"] # noqa: F401 + QueuedItemOrStop = Union[str, "PendingPayload", object] # noqa: F401 + + +class PendingPayload(object): + """A packet queued for the background sender that can go stale. + + Only payloads subject to expiry are wrapped in this. A replay-safe + payload -- one carrying its own explicit timestamp, so that delivering it + late doesn't change what it means -- is queued as the bare packet string + instead, because it needs none of the state here. The queue therefore + reads replay-safety off the entry's *type* rather than a stored flag (see + SenderQueue._expired() and is_replay_safe()), which keeps ~56 bytes per + replay-safe entry out of the queue and keeps those entries out of the + cyclic GC's traversal set entirely, since str holds no references. + + :ivar payload: The already-serialized packet text (including its + trailing newline), ready to be written to the socket. + :ivar enqueued_at: A monotonic timestamp recorded when the payload + became eligible for sending (i.e. when it was put on the queue). + Used to decide whether it has been sitting in the queue for too + long to still be worth sending. + """ + + __slots__ = ("payload", "enqueued_at") + + def __init__(self, payload, enqueued_at): + # type: (str, float) -> None + self.payload = payload + self.enqueued_at = enqueued_at + + +def is_replay_safe(item): + # type: (Union[str, PendingPayload]) -> bool + """True when this queue entry is exempt from expiry. + + Replay-safe entries are queued as bare strings; everything subject to + expiry is wrapped in PendingPayload. Centralised here so the type test + isn't repeated at every site that cares. + """ + return not isinstance(item, PendingPayload) + + +def payload_text(item): + # type: (Union[str, PendingPayload]) -> str + """The serialized packet text of a queue entry, whichever form it took.""" + if isinstance(item, PendingPayload): + return item.payload + return item + + +class SenderQueue(object): + """Bounded hand-off queue between application threads and the background sender thread. + + put() never rejects a payload outright. When the queue is already at its + maximum size, what happens depends on put_timeout: + - 0 (the default): no waiting at all -- the oldest entry is dropped + immediately to make room, along with any additional expired entries + left at the front. + - None: put() blocks the calling thread indefinitely, waiting for the + sender thread to drain a slot. It will wait forever if nothing ever + does -- this is an explicit opt-in to unbounded backpressure on the + calling thread. + - a positive number: put() blocks the calling thread for up to that + many seconds waiting for a slot; if the wait times out without one + opening up, it falls back to the same drop-oldest eviction as the + 0 case. + + get() drops expired entries lazily too, from the front, before returning + the next payload actually worth handing to the sender. + + A payload that fails to send (e.g. because the connection is down) can be + handed back with requeue_front() so it's retried first. That still + respects both the expiry check and the size limit though: the queue + must never grow past maxsize, and a payload that's gone stale while it + was being (re)tried is dropped rather than requeued. requeue_front() + never blocks on put_timeout. + """ + + def __init__(self, maxsize, expiry_seconds, on_drop_queue_full, on_drop_expired, put_timeout=0): + # type: (int, float, Callable[[QueuedItem], None], Callable[[QueuedItem], None], Optional[float]) -> None # noqa: E501 + self._maxsize = maxsize + self._expiry_seconds = expiry_seconds + self._on_drop_queue_full = on_drop_queue_full + self._on_drop_expired = on_drop_expired + self._put_timeout = put_timeout + self._deque = collections.deque() # type: collections.deque + self._lock = threading.Lock() + self._not_empty = threading.Condition(self._lock) + self._not_full = threading.Condition(self._lock) + self._all_tasks_done = threading.Condition(self._lock) + + # Keep track of the tasks that are being processed. A task pulled from the queue may + # be returned if the connection fails, so we don't consider the queue empty until + # all tasks have been dropped or sent. + self._unfinished_tasks = 0 + + # 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 + # normally holds at most one entry at a time. Tracking it lets + # requeue_front()/task_done() verify their precondition -- that the + # item they're handed really is the one get() currently has out of + # the queue -- so the misuse that would otherwise silently corrupt + # _unfinished_tasks (a double finish, a double requeue, or a requeue + # of an already-finished item) fails loudly instead of drifting the + # counter. The item itself is held in the value to keep it alive + # (and its id stable) for as long as the entry exists. + self._in_flight = {} # type: Dict[int, QueuedItemOrStop] + + def _expired(self, item, now): + # type: (QueuedItem, float) -> bool + if not isinstance(item, PendingPayload): + # A bare string is a replay-safe payload: it carries its own + # timestamp, so it never goes stale (see PendingPayload). + return False + return (now - item.enqueued_at) > self._expiry_seconds + + def _make_room_locked(self): + # type: () -> None + """Drop the oldest entry, plus any further expired entries at the front. + + Called with self._lock already held, and only when the queue is at + capacity. Never touches the Stop sentinel: by the time it's queued, + nothing else is ever put on the queue again, so it can only ever be + the newest entry, never the one being evicted here. + """ + if not self._deque or self._deque[0] is Stop: + return + + now = monotonic() + oldest = self._deque.popleft() + # The oldest entry is always dropped to make room. + if self._expired(oldest, now): + self._on_drop_expired(oldest) + else: + self._on_drop_queue_full(oldest) + self._finish_task_locked() + + # Keep clearing out additional stale entries left at the front. + # If any additional entries were cleared out notify not_full as the + # queue will now have available space for additional entries. + reclaimed = 0 + while self._deque and self._deque[0] is not Stop and self._expired(self._deque[0], now): + self._on_drop_expired(self._deque.popleft()) + self._finish_task_locked() + reclaimed += 1 + + if reclaimed: + self._not_full.notify(reclaimed) + + def put(self, item): + # type: (QueuedItemOrStop) -> None + """Queue a payload (or the Stop sentinel). + + If the queue is full: waits for room according to put_timeout -- + forever if it's None, up to put_timeout seconds if it's a positive + 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. + """ + 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: + # Wait forever: an explicit opt-in to unbounded + # backpressure on the calling thread. + while len(self._deque) >= self._maxsize: + self._not_full.wait() + elif self._put_timeout > 0: + deadline = monotonic() + self._put_timeout + while len(self._deque) >= self._maxsize: + remaining = deadline - monotonic() + if remaining <= 0: + break + self._not_full.wait(remaining) + # else: put_timeout is 0 (or negative) -- no wait at all, + # straight to eviction below. + + if len(self._deque) >= self._maxsize: + self._make_room_locked() + + self._deque.append(item) + self._unfinished_tasks += 1 + self._not_empty.notify() + + def requeue_front(self, item): + # type: (QueuedItem) -> None + """Put an in-flight payload back at the front after a failed send attempt. + + The payload was already accounted for by the put() that originally + queued it (its task isn't done yet), so a successful requeue here + doesn't touch _unfinished_tasks. But it's still subject to the same + rules as any other entry: an item that's expired while it was being + (re)tried is dropped instead of requeued, and the queue is never + allowed to grow past maxsize -- if it's already full, the requeue is + dropped too rather than evicting something else to make room for it. + Either way, a drop here finishes the task that put() started. + """ + with self._not_empty: + # The item handed back must be exactly the one get() currently has + # out of the queue. SenderQueue is single-consumer; a double + # requeue or a requeue of an already-finished item would otherwise + # corrupt _unfinished_tasks. If the item isn't in flight, log and + # bail out without touching the deque or the counter -- it's + # already been accounted for elsewhere, so this is a no-op rather + # than a crash. (Releasing it from in flight here is correct in + # every branch below: it's either requeued back onto the deque -- + # where a future get() will pick it up again -- or dropped for + # good.) + if not self._release_in_flight_locked(item, "requeue_front"): + return + + if self._expired(item, monotonic()): + self._on_drop_expired(item) + self._finish_task_locked() + return + + if self._maxsize > 0 and len(self._deque) >= self._maxsize: + self._on_drop_queue_full(item) + self._finish_task_locked() + return + + self._deque.appendleft(item) + self._not_empty.notify() + + def get(self): + # type: () -> QueuedItemOrStop + """Block for the next payload, silently dropping expired entries along the way.""" + while True: + with self._not_empty: + while not self._deque: + self._not_empty.wait() + item = self._deque.popleft() + # A slot just opened up: wake one thread blocked in put()'s + # wait-for-room loop, if any (harmless no-op otherwise). + self._not_full.notify() + # Record this item as in flight, owned by the current thread, + # until requeue_front() or task_done() releases it (see + # _in_flight in __init__). + self._take_in_flight_locked(item) + + if item is Stop: + return item + + # Guard the monotonic() call on the type test rather than letting + # _expired() do it: the argument is evaluated BEFORE the call, so + # `self._expired(item, monotonic())` read the clock on every get() + # including for bare strings, which are replay-safe and can never + # expire, so the value was computed and immediately discarded. + if isinstance(item, PendingPayload) and self._expired(item, monotonic()): + self._on_drop_expired(item) + self.task_done(item) + continue + + return item + + def _take_in_flight_locked(self, item): + # type: (QueuedItemOrStop) -> None + # Caller already holds self._lock (shared by _not_empty / _all_tasks_done). + # Records `item` as the one currently handed out by get(). A duplicate + # here means a previous get() was never finished (or the same object + # was queued twice); we log it and overwrite so the new handout is the + # one tracked, rather than crashing the sender thread. + key = id(item) + if key in self._in_flight: + log.error( + "dogstatsd sender queue: get() handed out an item already tracked as " + "in flight; a previous get() was never finished with task_done() / " + "requeue_front(), or the same object was queued more than once. " + "Counter bookkeeping may drift." + ) + self._in_flight[key] = item + + def _release_in_flight_locked(self, item, action): + # type: (QueuedItemOrStop, str) -> bool + # Caller already holds self._lock (shared by _not_empty / _all_tasks_done). + # Verifies `item` is currently in flight, then drops it from the + # in-flight map. `action` names the caller ("requeue_front"/ + # "task_done") for the log message. Returns False (after logging) when + # the item is not in flight, so the caller can skip the counter/deque + # mutation that would otherwise drift _unfinished_tasks -- without + # crashing the sender thread. + key = id(item) + if key not in self._in_flight: + log.error( + "dogstatsd sender queue: %s() was called on an item that is not " + "currently in flight (it was never returned by get(), or was " + "already finished). Ignoring it to keep the task counter consistent.", + action, + ) + return False + del self._in_flight[key] + return True + + def _finish_task_locked(self): + # type: () -> None + # Caller already holds self._lock (shared by _not_empty / _all_tasks_done). + unfinished = self._unfinished_tasks - 1 + if unfinished < 0: + # More finishes than puts: a real bookkeeping bug. Log it and + # clamp at zero rather than raising, so the sender thread stays + # alive. Notify in case a join() is waiting, so it doesn't hang. + log.error( + "dogstatsd sender queue: task accounting went negative " + "(_unfinished_tasks below zero); clamping. This indicates a " + "double finish or a finish without a matching put()." + ) + unfinished = 0 + self._unfinished_tasks = unfinished + if unfinished == 0: + self._all_tasks_done.notify_all() + + def task_done(self, item): + # type: (QueuedItemOrStop) -> None + with self._all_tasks_done: + # The item being finished must be the one get() currently has out + # of the queue. This is the counterpart to get()'s + # _take_in_flight_locked(); a second task_done() (double finish) + # would otherwise let _unfinished_tasks drift. If it's not in + # flight, log and bail out without decrementing -- the task was + # already finished elsewhere -- rather than crashing the sender. + if not self._release_in_flight_locked(item, "task_done"): + return + self._finish_task_locked() + + def join(self): + # type: () -> None + with self._all_tasks_done: + while self._unfinished_tasks: + self._all_tasks_done.wait() + + def qsize(self): + # type: () -> int + with self._lock: + return len(self._deque) + + def empty(self): + # type: () -> bool + with self._lock: + return not self._deque diff --git a/datadog/threadstats/base.py b/datadog/threadstats/base.py index 9dbbdb9be..69f474383 100644 --- a/datadog/threadstats/base.py +++ b/datadog/threadstats/base.py @@ -16,17 +16,13 @@ from functools import wraps from time import time -try: - from time import monotonic # type: ignore[attr-defined] -except ImportError: - from time import time as monotonic - # datadog from datadog.api.exceptions import ApiNotInitialized from datadog.threadstats.constants import MetricType from datadog.threadstats.events import EventsAggregator from datadog.threadstats.metrics import MetricsAggregator, Counter, Gauge, Histogram, Timing, Distribution, Set from datadog.threadstats.reporters import HttpReporter +from datadog.util.compat import monotonic # Loggers log = logging.getLogger("datadog.threadstats") diff --git a/datadog/util/compat.py b/datadog/util/compat.py index febb804b9..4b863e8e7 100644 --- a/datadog/util/compat.py +++ b/datadog/util/compat.py @@ -109,6 +109,16 @@ def emit(self, record): pass +# Python >= 3.3 +if sys.version_info >= (3, 3): + from time import monotonic +# Python 2.x: there is no monotonic clock, so fall back to the wall clock. +# Callers that compare two readings (elapsed time, queue entry age) are +# therefore sensitive to the clock being stepped backwards on Python 2. +else: + from time import time as monotonic + + def _is_py_version_higher_than(major, minor=0): # type: (int, int) -> bool """ diff --git a/tests/integration/dogstatsd/test_statsd_sender.py b/tests/integration/dogstatsd/test_statsd_sender.py index 1cd5a50ed..5841cd772 100644 --- a/tests/integration/dogstatsd/test_statsd_sender.py +++ b/tests/integration/dogstatsd/test_statsd_sender.py @@ -103,7 +103,8 @@ def test_fork_hooks(disable_background_sender, disable_buffering): assert statsd._flush_thread is None assert statsd._sender_thread is None assert statsd._queue is None or statsd._queue.empty() - assert len(statsd._buffer) == 0 + # Buffered lines are split by expiry policy, so check every batch. + assert not statsd._buffer and not statsd._buffer_rs statsd.post_fork_parent() diff --git a/tests/unit/dogstatsd/test_statsd.py b/tests/unit/dogstatsd/test_statsd.py index 3045f3286..d6179841e 100644 --- a/tests/unit/dogstatsd/test_statsd.py +++ b/tests/unit/dogstatsd/test_statsd.py @@ -9,7 +9,8 @@ """ # Standard libraries from collections import deque -from contextlib import closing +from contextlib import closing, contextmanager +import logging import struct from threading import Thread import errno @@ -30,7 +31,9 @@ # 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, 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, UDS_CONNECT_RETRY_INITIAL_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 from datadog.util.compat import is_higher_py35, is_p3k from tests.util.contextmanagers import preserve_environment_variable, EnvVars @@ -123,20 +126,26 @@ def __init__(self): super(OverflownSocket, self).__init__(errno.EAGAIN) -def telemetry_metrics(metrics=1, events=0, service_checks=0, bytes_sent=0, bytes_dropped_writer=0, packets_sent=1, packets_dropped_writer=0, transport="udp", tags="", bytes_dropped_queue=0, packets_dropped_queue=0): +def telemetry_metrics(metrics=1, events=0, service_checks=0, bytes_sent=0, bytes_dropped_writer=0, packets_sent=1, packets_dropped_writer=0, transport="udp", tags="", bytes_dropped_queue=0, packets_dropped_queue=0, bytes_dropped_expired=0, packets_dropped_expired=0): tags = "," + tags if tags else "" + # Expired drops have no dedicated wire metric: they're folded into the + # *_dropped_queue lines (and totals) reported to the Agent. See + # DogStatsd._flush_telemetry(). + reported_bytes_dropped_queue = bytes_dropped_queue + bytes_dropped_expired + reported_packets_dropped_queue = packets_dropped_queue + packets_dropped_expired + return "\n".join([ "datadog.dogstatsd.client.metrics:{}|c|#client:py,client_version:{},client_transport:{}{}".format(metrics, version, transport, tags), "datadog.dogstatsd.client.events:{}|c|#client:py,client_version:{},client_transport:{}{}".format(events, version, transport, tags), "datadog.dogstatsd.client.service_checks:{}|c|#client:py,client_version:{},client_transport:{}{}".format(service_checks, version, transport, tags), "datadog.dogstatsd.client.bytes_sent:{}|c|#client:py,client_version:{},client_transport:{}{}".format(bytes_sent, version, transport, tags), - "datadog.dogstatsd.client.bytes_dropped:{}|c|#client:py,client_version:{},client_transport:{}{}".format(bytes_dropped_queue + bytes_dropped_writer, version, transport, tags), - "datadog.dogstatsd.client.bytes_dropped_queue:{}|c|#client:py,client_version:{},client_transport:{}{}".format(bytes_dropped_queue, version, transport, tags), + "datadog.dogstatsd.client.bytes_dropped:{}|c|#client:py,client_version:{},client_transport:{}{}".format(reported_bytes_dropped_queue + bytes_dropped_writer, version, transport, tags), + "datadog.dogstatsd.client.bytes_dropped_queue:{}|c|#client:py,client_version:{},client_transport:{}{}".format(reported_bytes_dropped_queue, version, transport, tags), "datadog.dogstatsd.client.bytes_dropped_writer:{}|c|#client:py,client_version:{},client_transport:{}{}".format(bytes_dropped_writer, version, transport, tags), "datadog.dogstatsd.client.packets_sent:{}|c|#client:py,client_version:{},client_transport:{}{}".format(packets_sent, version, transport, tags), - "datadog.dogstatsd.client.packets_dropped:{}|c|#client:py,client_version:{},client_transport:{}{}".format(packets_dropped_queue + packets_dropped_writer, version, transport, tags), - "datadog.dogstatsd.client.packets_dropped_queue:{}|c|#client:py,client_version:{},client_transport:{}{}".format(packets_dropped_queue, version, transport, tags), + "datadog.dogstatsd.client.packets_dropped:{}|c|#client:py,client_version:{},client_transport:{}{}".format(reported_packets_dropped_queue + packets_dropped_writer, version, transport, tags), + "datadog.dogstatsd.client.packets_dropped_queue:{}|c|#client:py,client_version:{},client_transport:{}{}".format(reported_packets_dropped_queue, version, transport, tags), "datadog.dogstatsd.client.packets_dropped_writer:{}|c|#client:py,client_version:{},client_transport:{}{}".format(packets_dropped_writer, version, transport, tags), ]) + "\n" @@ -167,6 +176,28 @@ def tearDown(self): """ self._procfs_mock.stop() + @contextmanager + def _capture_error_logs(self): + # assertLogs() is Python 3.4+ only, but this suite still runs on + # Python 2.7 / pypy2.7. Capture ERROR records on the dogstatsd logger + # with a plain handler instead, so the guard tests work everywhere. + captured = [] + + class _CaptureHandler(logging.Handler): + def emit(self, record): + captured.append(record) + + dogstatsd_logger = logging.getLogger("datadog.dogstatsd") + handler = _CaptureHandler(level=logging.ERROR) + dogstatsd_logger.addHandler(handler) + prev_level = dogstatsd_logger.level + dogstatsd_logger.setLevel(min(prev_level, logging.ERROR)) + try: + yield captured + finally: + dogstatsd_logger.removeHandler(handler) + dogstatsd_logger.setLevel(prev_level) + def assert_equal_telemetry(self, expected_payload, actual_payload, telemetry=None, **kwargs): if telemetry is None: telemetry = telemetry_metrics(bytes_sent=len(expected_payload), **kwargs) @@ -1742,6 +1773,91 @@ def test_manual_buffer_ops_deprecation(self, mock_warn): self.statsd.close_buffer() self.assertEqual(mock_warn.call_count, 2) + def test_mixed_batch_splits_by_replay_safety(self): + # A queued packet expires as a single unit, so every line in it has to + # share one expiry policy. Batching timestamped lines together with + # plain ones would make the whole packet non-replay-safe and strip the + # timestamped lines of the staleness exemption they're supposed to + # have, dropping them with the batch after ~10s of backlog. The buffer + # must split by policy instead. + sent = [] + self.statsd._send_to_server = lambda packet, replay_safe=False: sent.append((packet, replay_safe)) + + self.statsd.open_buffer() + self.statsd.gauge_with_timestamp("ts.one", 1, timestamp=1700000000) + self.statsd.gauge("plain.one", 2) + self.statsd.gauge_with_timestamp("ts.two", 3, timestamp=1700000001) + self.statsd.gauge("plain.two", 4) + self.statsd.close_buffer() + + # One packet per expiry policy, not one per metric: interleaving must + # not defeat batching. + self.assertEqual(len(sent), 2, "expected exactly one packet per expiry policy, got: {!r}".format(sent)) + + by_policy = dict((replay_safe, packet) for packet, replay_safe in sent) + self.assertEqual(sorted(by_policy.keys()), [False, True]) + + # Assert on structure rather than exact packet text: constant/origin + # tags vary by environment, but which lines land in which packet, and + # in what order, does not. + def names(packet): + return [line.split(":")[0] for line in packet.split("\n")] + + self.assertEqual(names(by_policy[True]), ["ts.one", "ts.two"]) + self.assertEqual(names(by_policy[False]), ["plain.one", "plain.two"]) + + # The invariant that actually matters: no packet mixes the two, and no + # timestamped line ever rides in an expiring packet. + for packet, replay_safe in sent: + lines = packet.split("\n") + timestamped = [line for line in lines if "|T" in line] + if replay_safe: + self.assertEqual(timestamped, lines, "replay-safe packet must be entirely timestamped lines") + else: + self.assertEqual(timestamped, [], "timestamped line leaked into an expiring packet") + + def test_mixed_batch_respects_max_payload_size_per_buffer(self): + # Each buffer has to stay under _max_payload_size on its own, and one + # buffer overflowing must not drag the other one out with it. + sent = [] + self.statsd._send_to_server = lambda packet, replay_safe=False: sent.append((packet, replay_safe)) + + # Measure a real serialised line and size the cap from it. Hard-coding + # a byte count would make the test depend on how long constant/origin + # tags happen to make each line in this environment: too small and a + # single line breaches the cap, too large and nothing ever overflows. + self.statsd.open_buffer() + self.statsd.gauge("plain.filler.0", 0) + line_size = self.statsd._buffer_size + self.statsd.close_buffer() + del sent[:] + + # Room for two lines, so every third one forces a flush. + self.statsd._max_payload_size = line_size * 2 + 1 + + self.statsd.open_buffer() + # One small replay-safe line that should still be buffered while the + # plain buffer churns through several flushes. + self.statsd.gauge_with_timestamp("ts.keep", 1, timestamp=1700000000) + for i in range(12): + self.statsd.gauge("plain.filler.{}".format(i), i) + flushes_before_close = len(sent) + self.statsd.close_buffer() + + self.assertGreater(flushes_before_close, 0, "the plain buffer should have overflowed at least once") + self.assertTrue( + all(not replay_safe for _, replay_safe in sent[:flushes_before_close]), + "overflow of the plain buffer must not flush the replay-safe buffer", + ) + for packet, _ in sent: + self.assertLessEqual(len(packet) + 1, self.statsd._max_payload_size) + + # The replay-safe line survived to the final flush, intact and alone. + final_packet, final_replay_safe = sent[-1] + self.assertTrue(final_replay_safe) + self.assertEqual([line.split(":")[0] for line in final_packet.split("\n")], ["ts.keep"]) + self.assertIn("|T1700000000", final_packet) + def test_batching_sequential(self): self.statsd.open_buffer() self.statsd.gauge('discarded.data', 123) @@ -1954,6 +2070,51 @@ def test_telemetry(self): self.assertEqual(0, self.statsd.bytes_dropped_queue) self.assertEqual(0, self.statsd.packets_dropped_queue) + def test_telemetry_folds_expired_drops_into_dropped_queue(self): + # There's no dedicated wire metric for expired drops: they're + # reported to the Agent as part of *_dropped_queue (and the combined + # *_dropped total), alongside capacity-based queue drops, since both + # never reach a socket write attempt. The distinction is still + # available in-process via bytes_dropped_expired/packets_dropped_expired. + # Avoid any real container-id auto-detected from the host/sandbox + # cgroup leaking into the expected payload below -- this test is + # about the telemetry counters, not the container-id field. + self.statsd._container_id = None + + self.statsd.bytes_dropped_queue = 8 + self.statsd.packets_dropped_queue = 9 + self.statsd.bytes_dropped_expired = 10 + self.statsd.packets_dropped_expired = 11 + self.statsd.bytes_dropped_writer = 5 + self.statsd.packets_dropped_writer = 7 + + self.statsd.open_buffer() + self.statsd.gauge('page.views', 123) + self.statsd.close_buffer() + + payload = 'page.views:123|g\n' + telemetry = telemetry_metrics( + metrics=1, + bytes_sent=len(payload), + packets_sent=1, + bytes_dropped_queue=8, + packets_dropped_queue=9, + bytes_dropped_expired=10, + packets_dropped_expired=11, + bytes_dropped_writer=5, + packets_dropped_writer=7, + ) + + self.assert_equal_telemetry(payload, self.recv(2), telemetry=telemetry) + + # The in-process counters stay separate even after the flush resets + # them -- confirming the fold happens only in the wire output, not + # by merging the underlying attributes. + self.assertEqual(0, self.statsd.bytes_dropped_queue) + self.assertEqual(0, self.statsd.packets_dropped_queue) + self.assertEqual(0, self.statsd.bytes_dropped_expired) + self.assertEqual(0, self.statsd.packets_dropped_expired) + def test_telemetry_flush_interval(self): dogstatsd = DogStatsd(disable_buffering=False) fake_socket = FakeSocket() @@ -2611,6 +2772,28 @@ def test_sender_mode(self): statsd = DogStatsd(disable_background_sender=False) self.assertIsNotNone(statsd._queue) + def test_sender_queue_expiry_seconds_defaults_to_constant(self): + statsd = DogStatsd(disable_background_sender=False) + self.assertEqual(statsd._sender_queue_expiry_seconds, PENDING_PAYLOAD_EXPIRY_SECONDS) + self.assertEqual(statsd._queue._expiry_seconds, PENDING_PAYLOAD_EXPIRY_SECONDS) + statsd.stop() + + def test_sender_queue_expiry_seconds_is_configurable(self): + statsd = DogStatsd( + disable_background_sender=False, + sender_queue_expiry_seconds=0.5, + ) + self.assertEqual(statsd._sender_queue_expiry_seconds, 0.5) + self.assertEqual(statsd._queue._expiry_seconds, 0.5) + statsd.stop() + + def test_enable_background_sender_accepts_expiry_seconds(self): + statsd = DogStatsd(disable_background_sender=True) + statsd.enable_background_sender(sender_queue_expiry_seconds=1.5) + self.assertEqual(statsd._sender_queue_expiry_seconds, 1.5) + self.assertEqual(statsd._queue._expiry_seconds, 1.5) + statsd.stop() + def test_sender_calls_task_done(self): statsd = DogStatsd(disable_background_sender=False) statsd.socket = OverflownSocket() @@ -2619,28 +2802,712 @@ def test_sender_calls_task_done(self): def test_sender_queue_no_timeout(self): statsd = DogStatsd(disable_background_sender=False, sender_queue_timeout=None) + statsd.stop() - def test_bytes_dropped_queue_counts_actual_bytes(self): - # Use a queue of size 1 and a non-blocking timeout so packets are dropped - # when the queue is full, then verify bytes_dropped_queue reflects the real - # byte length of the dropped packet (including the appended newline). + 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) + # block waiting for room, not just an internal SenderQueue detail. statsd = DogStatsd( disable_background_sender=False, sender_queue_size=1, - sender_queue_timeout=0, + sender_queue_timeout=5.0, ) + # No socket assigned: the sender thread can never drain anything by + # actually sending, so the only way room opens up is via get() + # pulling an item off (which happens immediately, since nothing can + # succeed in sending it -- it gets hard-dropped as a writer failure + # and the sender loop moves on to the next get()). statsd.socket = FakeSocket() - # Build a packet whose serialised form we know, then compute its length. - metric_name = "test.metric" + statsd.increment("first") + + t0 = time.time() + statsd.increment("second") + elapsed = time.time() - t0 + + self.assertLess(elapsed, 5.0, "should not have waited out the full 5s timeout") + statsd.wait_for_pending() + statsd.stop() + + def test_bytes_dropped_queue_counts_actual_bytes(self): + # No sender thread: a live one could drain the first payload before the + # third is queued, so nothing would be evicted and the counters below + # would describe a schedule that never happened. Size 2 rather than 1 + # so the eviction order is observable -- with a single slot the evicted + # entry is both the oldest and the newest. + statsd = DogStatsd(disable_background_sender=True) + statsd._queue = SenderQueue( + 2, + PENDING_PAYLOAD_EXPIRY_SECONDS, + statsd._account_dropped_queue_full, + statsd._account_dropped_expired, + ) - # Send two packets: the first fills the queue, the second is dropped. - statsd._send_to_server(metric_name) - statsd._send_to_server(metric_name) + first, second, third = "test.metric.first", "test.metric.second", "test.metric.third" + statsd._send_to_server(first) + statsd._send_to_server(second) + statsd._send_to_server(third) # evicts the oldest (first) to make room - expected_bytes = len((metric_name + '\n').encode("utf-8")) - self.assertEqual(statsd.bytes_dropped_queue, expected_bytes) + # bytes_dropped_queue is the real byte length, including the newline + # _send_to_server() appends. + self.assertEqual(statsd.bytes_dropped_queue, len((first + "\n").encode("utf-8"))) self.assertEqual(statsd.packets_dropped_queue, 1) + self.assertEqual(statsd.bytes_dropped_expired, 0) + self.assertEqual(statsd.packets_dropped_expired, 0) + + # Dropping the oldest leaves the two newest queued, in order. + survivors = [statsd._queue.get().payload, statsd._queue.get().payload] + self.assertEqual(survivors, [second + "\n", third + "\n"]) + + statsd.stop() + + def test_sender_queue_put_timeout_default_evicts_immediately(self): + # Default put_timeout (0, whether omitted or explicit): no waiting + # at all, same as before this feature existed. Deliberately omits + # put_timeout here to prove the *default* -- not just 0 -- means + # "don't wait", since None means something very different (wait + # forever) and must not be the implicit default for anyone who + # constructs a SenderQueue without thinking about put_timeout at all. + dropped_queue_full = [] + pending_queue = SenderQueue( + maxsize=1, + expiry_seconds=100.0, + on_drop_queue_full=dropped_queue_full.append, + on_drop_expired=lambda item: self.fail("unexpected expiry drop"), + ) + + pending_queue.put(PendingPayload("first\n", sender_queue_clock())) + + t0 = time.time() + pending_queue.put(PendingPayload("second\n", sender_queue_clock())) + elapsed = time.time() - t0 + + self.assertLess(elapsed, 0.05, "put() should not have waited at all") + self.assertEqual([p.payload for p in dropped_queue_full], ["first\n"]) + self.assertEqual(pending_queue.get().payload, "second\n") + + def test_sender_queue_put_timeout_zero_evicts_immediately(self): + # Same as the default, but with put_timeout=0 passed explicitly. + dropped_queue_full = [] + pending_queue = SenderQueue( + maxsize=1, + expiry_seconds=100.0, + on_drop_queue_full=dropped_queue_full.append, + on_drop_expired=lambda item: self.fail("unexpected expiry drop"), + put_timeout=0, + ) + + pending_queue.put(PendingPayload("first\n", sender_queue_clock())) + + t0 = time.time() + pending_queue.put(PendingPayload("second\n", sender_queue_clock())) + elapsed = time.time() - t0 + + self.assertLess(elapsed, 0.05, "put() should not have waited at all") + self.assertEqual([p.payload for p in dropped_queue_full], ["first\n"]) + self.assertEqual(pending_queue.get().payload, "second\n") + + def test_sender_queue_put_timeout_none_waits_forever_and_never_evicts(self): + # put_timeout=None is an explicit opt-in to unbounded blocking: put() + # must keep waiting indefinitely -- not fall back to eviction after + # some internal default -- until room actually opens up. + dropped_queue_full = [] + pending_queue = SenderQueue( + maxsize=1, + expiry_seconds=100.0, + on_drop_queue_full=lambda item: dropped_queue_full.append(item), + on_drop_expired=lambda item: self.fail("unexpected expiry drop"), + put_timeout=None, + ) + pending_queue.put(PendingPayload("first\n", sender_queue_clock())) + + result = {} + + def blocked_put(): + t0 = time.time() + pending_queue.put(PendingPayload("second\n", sender_queue_clock())) + result["elapsed"] = time.time() - t0 + + t = threading.Thread(target=blocked_put) + t.start() + try: + # Nothing is draining the queue: with a real timeout this would + # have already fired and evicted "first" well before 1s. With + # None it must still be waiting. + time.sleep(1.0) + self.assertTrue(t.is_alive(), "put(timeout=None) must keep waiting, never fall back to eviction on its own") + self.assertEqual(dropped_queue_full, []) + + # Now free up room: the blocked put() should wake up and + # succeed without ever having dropped anything. + first = pending_queue.get() + self.assertEqual(first.payload, "first\n") + pending_queue.task_done(first) + finally: + t.join(timeout=5.0) + + self.assertFalse(t.is_alive()) + self.assertEqual(dropped_queue_full, [], "put_timeout=None must never fall back to eviction") + self.assertEqual(pending_queue.get().payload, "second\n") + + def test_sender_queue_put_timeout_wakes_up_when_room_opens(self): + # A slot freed by get() (well within put_timeout) should wake a + # blocked put() immediately rather than making it wait out the full + # timeout, and nothing should be dropped. + dropped_queue_full = [] + pending_queue = SenderQueue( + maxsize=1, + expiry_seconds=100.0, + on_drop_queue_full=dropped_queue_full.append, + on_drop_expired=lambda item: self.fail("unexpected expiry drop"), + put_timeout=5.0, + ) + pending_queue.put(PendingPayload("first\n", sender_queue_clock())) + + result = {} + + def blocked_put(): + t0 = time.time() + pending_queue.put(PendingPayload("second\n", sender_queue_clock())) + result["elapsed"] = time.time() - t0 + + t = threading.Thread(target=blocked_put) + t.start() + time.sleep(0.2) + self.assertTrue(t.is_alive(), "put() should still be waiting for room") + + # Drain the one slot: the blocked put() should wake up promptly. + first = pending_queue.get() + self.assertEqual(first.payload, "first\n") + pending_queue.task_done(first) + + t.join(timeout=5.0) + self.assertFalse(t.is_alive()) + self.assertLess(result["elapsed"], 5.0, "should have woken up well before the 5s timeout") + self.assertEqual(dropped_queue_full, [], "nothing should have been dropped: room opened up in time") + self.assertEqual(pending_queue.get().payload, "second\n") + + def test_sender_queue_put_timeout_falls_back_to_eviction(self): + # If room never opens up within put_timeout, put() falls back to + # the same drop-oldest eviction as the immediate (no-wait) case. + dropped_queue_full = [] + pending_queue = SenderQueue( + maxsize=1, + expiry_seconds=100.0, + on_drop_queue_full=dropped_queue_full.append, + on_drop_expired=lambda item: self.fail("unexpected expiry drop"), + put_timeout=0.2, + ) + pending_queue.put(PendingPayload("first\n", sender_queue_clock())) + + t0 = time.time() + pending_queue.put(PendingPayload("second\n", sender_queue_clock())) + elapsed = time.time() - t0 + + self.assertGreaterEqual(elapsed, 0.2) + self.assertEqual([p.payload for p in dropped_queue_full], ["first\n"]) + self.assertEqual(pending_queue.get().payload, "second\n") + + def test_sender_queue_bulk_expired_reclaim_wakes_blocked_producers(self): + # When the eviction path's cleanup loop reclaims *more* than the one + # slot its caller needs, the surplus is real free capacity. Producers + # already parked in put()'s wait-for-room loop have to be told about + # it, otherwise they sleep out their full put_timeout while the queue + # sits half empty. + put_timeout = 1.0 + maxsize = 4 + dropped_expired = [] + pending_queue = SenderQueue( + maxsize=maxsize, + expiry_seconds=100.0, + on_drop_queue_full=lambda item: self.fail("entries are stale: expect expiry drops, not full drops"), + on_drop_expired=dropped_expired.append, + put_timeout=put_timeout, + ) + # Fill to capacity with entries that are already stale, so the + # cleanup loop has something to reclaim beyond the mandatory one. + stale_clock = sender_queue_clock() - 1000.0 + for i in range(maxsize): + pending_queue.put(PendingPayload("stale-{}\n".format(i), stale_clock)) + + result = {} + + def evictor(): + # Queue is full and nothing drains it, so this waits out + # put_timeout and then falls back to eviction, whose cleanup loop + # reclaims all remaining stale entries in one go. + pending_queue.put(PendingPayload("evictor\n", sender_queue_clock())) + + def late_waiter(): + t0 = time.time() + pending_queue.put(PendingPayload("late\n", sender_queue_clock())) + result["elapsed"] = time.time() - t0 + + t_evictor = threading.Thread(target=evictor) + t_evictor.start() + # Start the second producer halfway through the first one's timeout so + # its own deadline is strictly later: it must be woken by the bulk + # reclaim, not by its own timeout firing. + time.sleep(put_timeout / 2.0) + t_late = threading.Thread(target=late_waiter) + t_late.start() + + t_evictor.join(timeout=5.0) + t_late.join(timeout=5.0) + self.assertFalse(t_evictor.is_alive()) + self.assertFalse(t_late.is_alive()) + + # All four stale entries went out through the cleanup path. + self.assertEqual( + [p.payload for p in dropped_expired], + ["stale-0\n", "stale-1\n", "stale-2\n", "stale-3\n"], + ) + # Both live payloads made it, and the queue is well under maxsize. + self.assertEqual(pending_queue.qsize(), 2) + + # The heart of it: the late producer had roughly put_timeout/2 left on + # its own clock when capacity opened up. Waking on the reclaim means + # ~put_timeout/2 elapsed; sleeping through it means the full + # put_timeout. Assert it beat its own deadline by a clear margin. + self.assertLess( + result["elapsed"], + put_timeout * 0.9, + "blocked producer slept through its put_timeout despite the bulk reclaim freeing capacity", + ) + + def test_sender_queue_requeue_front_never_blocks_on_put_timeout(self): + # requeue_front() runs on the background sender thread; it must + # never wait on put_timeout, or one stuck retry would stall every + # other queued payload behind it. + dropped_queue_full = [] + pending_queue = SenderQueue( + maxsize=1, + expiry_seconds=100.0, + on_drop_queue_full=dropped_queue_full.append, + on_drop_expired=lambda item: self.fail("unexpected expiry drop"), + put_timeout=5.0, + ) + in_flight = PendingPayload("in-flight\n", sender_queue_clock()) + pending_queue.put(in_flight) + got = pending_queue.get() + pending_queue.put(PendingPayload("new\n", sender_queue_clock())) # fills the one slot again + + t0 = time.time() + pending_queue.requeue_front(got) + elapsed = time.time() - t0 + + self.assertLess(elapsed, 0.05, "requeue_front() must not block on put_timeout") + self.assertEqual([p.payload for p in dropped_queue_full], ["in-flight\n"]) + self.assertEqual(pending_queue.get().payload, "new\n") + + def test_sender_queue_drops_oldest_and_stale_entries_on_overflow(self): + dropped_queue_full = [] + dropped_expired = [] + + pending_queue = SenderQueue( + maxsize=2, + expiry_seconds=20.0, + on_drop_queue_full=dropped_queue_full.append, + on_drop_expired=dropped_expired.append, + ) + + now = sender_queue_clock() + fresh = PendingPayload("fresh\n", now) + stale = PendingPayload("stale\n", now - 100) + newest = PendingPayload("newest\n", now) + + # Fill the queue: [fresh, stale] (stale is already expired, but that + # doesn't matter until something tries to make room or pull it off). + pending_queue.put(fresh) + pending_queue.put(stale) + self.assertEqual(pending_queue.qsize(), 2) + + # Queue is full: the oldest entry (fresh) is evicted to make room, and + # since the next entry at the front (stale) is also expired, it gets + # opportunistically cleared out too. + pending_queue.put(newest) + + self.assertEqual([p.payload for p in dropped_queue_full], ["fresh\n"]) + self.assertEqual([p.payload for p in dropped_expired], ["stale\n"]) + self.assertEqual(pending_queue.qsize(), 1) + self.assertEqual(pending_queue.get().payload, "newest\n") + + def test_sender_queue_overflow_attributes_stale_oldest_entry_to_expiry(self): + dropped_queue_full = [] + dropped_expired = [] + + pending_queue = SenderQueue( + maxsize=1, + expiry_seconds=20.0, + on_drop_queue_full=dropped_queue_full.append, + on_drop_expired=dropped_expired.append, + ) + + stale = PendingPayload("stale\n", sender_queue_clock() - 100) + pending_queue.put(stale) + + # The oldest (and only) entry being evicted is itself already + # expired: that's a staleness drop, not a queue-full drop. + pending_queue.put(PendingPayload("newest\n", sender_queue_clock())) + + self.assertEqual(dropped_queue_full, []) + self.assertEqual([p.payload for p in dropped_expired], ["stale\n"]) + + def test_sender_queue_get_drops_expired_entries(self): + dropped_expired = [] + + pending_queue = SenderQueue( + maxsize=0, + expiry_seconds=20.0, + on_drop_queue_full=lambda item: self.fail("unexpected queue-full drop"), + on_drop_expired=dropped_expired.append, + ) + + now = sender_queue_clock() + pending_queue.put(PendingPayload("stale-1\n", now - 100)) + pending_queue.put(PendingPayload("stale-2\n", now - 100)) + pending_queue.put(PendingPayload("fresh\n", now)) + + # get() lazily drains every stale entry at the front before handing + # back the next payload actually worth sending. + item = pending_queue.get() + self.assertEqual(item.payload, "fresh\n") + self.assertEqual([p.payload for p in dropped_expired], ["stale-1\n", "stale-2\n"]) + + def test_sender_queue_replay_safe_payload_never_expires(self): + pending_queue = SenderQueue( + maxsize=0, + expiry_seconds=20.0, + on_drop_queue_full=lambda item: self.fail("unexpected queue-full drop"), + on_drop_expired=lambda item: self.fail("replay-safe payload should not expire"), + ) + + # A replay-safe entry is the bare string, so there is no enqueued_at to + # age against at all -- it can never be dropped for staleness however + # long it sits there. + old_but_replay_safe = "timestamped\n" + pending_queue.put(old_but_replay_safe) + + self.assertTrue(is_replay_safe(old_but_replay_safe)) + self.assertIs(pending_queue.get(), old_but_replay_safe) + + def test_sender_queue_requeue_front_when_room_available(self): + pending_queue = SenderQueue( + maxsize=2, + expiry_seconds=20.0, + on_drop_queue_full=lambda item: self.fail("unexpected queue-full drop"), + on_drop_expired=lambda item: self.fail("unexpected expiry drop"), + ) + + in_flight = PendingPayload("in-flight\n", sender_queue_clock()) + pending_queue.put(in_flight) + + # Simulate the sender thread picking it up and failing to send it. + got = pending_queue.get() + self.assertIs(got, in_flight) + pending_queue.requeue_front(got) + + # There was room for it: it's retried first, ahead of anything newer. + pending_queue.put(PendingPayload("new\n", sender_queue_clock())) + self.assertEqual(pending_queue.get().payload, "in-flight\n") + self.assertEqual(pending_queue.get().payload, "new\n") + + def test_sender_queue_requeue_front_drops_when_queue_is_full(self): + dropped_queue_full = [] + + pending_queue = SenderQueue( + maxsize=1, + expiry_seconds=20.0, + on_drop_queue_full=dropped_queue_full.append, + on_drop_expired=lambda item: self.fail("unexpected expiry drop"), + ) + + in_flight = PendingPayload("in-flight\n", sender_queue_clock()) + pending_queue.put(in_flight) + + # Simulate the sender thread picking it up, failing to send it, and a + # fresh payload filling the now-empty slot in the meantime. + got = pending_queue.get() + self.assertIs(got, in_flight) + pending_queue.put(PendingPayload("new\n", sender_queue_clock())) + + # The queue is already at maxsize: the requeue is dropped rather than + # growing the queue past its limit or evicting the newer entry. + pending_queue.requeue_front(got) + + self.assertEqual([p.payload for p in dropped_queue_full], ["in-flight\n"]) + self.assertEqual(pending_queue.qsize(), 1) + self.assertEqual(pending_queue.get().payload, "new\n") + + def test_sender_queue_requeue_front_drops_when_expired(self): + dropped_expired = [] + + pending_queue = SenderQueue( + maxsize=0, + expiry_seconds=0.01, + on_drop_queue_full=lambda item: self.fail("unexpected queue-full drop"), + on_drop_expired=dropped_expired.append, + ) + + # Simulate the sender thread picking up a payload and failing to + # send it, with enough time passing in between that it's now stale. + # Unbounded queue (so it's never "full") isolates the expiry check. + in_flight = PendingPayload("stale\n", sender_queue_clock()) + pending_queue.put(in_flight) + got = pending_queue.get() + time.sleep(0.02) + + pending_queue.requeue_front(got) + + self.assertEqual([p.payload for p in dropped_expired], ["stale\n"]) + self.assertEqual(pending_queue.qsize(), 0) + + def test_sender_queue_requeue_front_ignores_item_not_in_flight(self): + # requeue_front() must be called on the exact item get() returned, and + # only while it's still in flight. Requeuing something that was never + # get() (or was already finished) would otherwise corrupt + # _unfinished_tasks. The guard logs the misuse and ignores the call + # (no deque change, no counter change) so the sender thread stays + # alive rather than crashing on a bookkeeping bug. + pending_queue = SenderQueue( + maxsize=0, + expiry_seconds=20.0, + on_drop_queue_full=lambda item: self.fail("unexpected queue-full drop"), + on_drop_expired=lambda item: self.fail("unexpected expiry drop"), + ) + + never_got = PendingPayload("never-get\n", sender_queue_clock()) + with self._capture_error_logs() as captured: + pending_queue.requeue_front(never_got) + self.assertTrue(captured, "the not-in-flight requeue should have logged an error") + self.assertEqual(pending_queue.qsize(), 0, "never-got item was not added to the deque") + self.assertEqual(pending_queue._unfinished_tasks, 0, "counter untouched") + + # And requeuing an item that was already finished (get then task_done) + # is equally ignored: it's no longer in flight, the counter stays put. + finished = PendingPayload("finished\n", sender_queue_clock()) + pending_queue.put(finished) + got = pending_queue.get() + self.assertIs(got, finished) + pending_queue.task_done(got) + self.assertEqual(pending_queue._unfinished_tasks, 0, "finished -> counter back to zero") + with self._capture_error_logs() as captured: + pending_queue.requeue_front(got) + self.assertTrue(captured, "the already-finished requeue should have logged an error") + self.assertEqual(pending_queue._unfinished_tasks, 0, "counter did not drift on the ignored requeue") + self.assertEqual(pending_queue.qsize(), 0, "already-finished item was not re-added") + + def test_sender_queue_requeue_front_ignores_double_requeue(self): + # A second requeue_front() of the same item (e.g. two threads both + # handed the same in-flight reference) would put it in the deque twice + # while only one task was ever counted. The in-flight guard logs the + # second call and ignores it, so the deque and counter stay consistent. + pending_queue = SenderQueue( + maxsize=0, + expiry_seconds=20.0, + on_drop_queue_full=lambda item: self.fail("unexpected queue-full drop"), + on_drop_expired=lambda item: self.fail("unexpected expiry drop"), + ) + + item = PendingPayload("once\n", sender_queue_clock()) + pending_queue.put(item) + got = pending_queue.get() + self.assertIs(got, item) + pending_queue.requeue_front(got) # first requeue: fine, back in the deque + self.assertEqual(pending_queue.qsize(), 1) + with self._capture_error_logs() as captured: + pending_queue.requeue_front(got) # second: no longer in flight + self.assertTrue(captured, "the double requeue should have logged an error") + self.assertEqual(pending_queue.qsize(), 1, "item was NOT added to the deque a second time") + self.assertEqual(pending_queue._unfinished_tasks, 1, "counter unchanged") + + def test_sender_queue_task_done_ignores_double_finish(self): + # A second task_done() on the same item (double finish) would drive + # _unfinished_tasks negative. The in-flight guard logs the second call + # and skips the decrement, so the counter stays at zero instead of + # going to -1 -- and the sender thread stays alive. + pending_queue = SenderQueue( + maxsize=0, + expiry_seconds=20.0, + on_drop_queue_full=lambda item: self.fail("unexpected queue-full drop"), + on_drop_expired=lambda item: self.fail("unexpected expiry drop"), + ) + + item = PendingPayload("once\n", sender_queue_clock()) + pending_queue.put(item) + got = pending_queue.get() + pending_queue.task_done(got) + self.assertEqual(pending_queue._unfinished_tasks, 0, "first finish -> counter at zero") + with self._capture_error_logs() as captured: + pending_queue.task_done(got) + self.assertTrue(captured, "the double finish should have logged an error") + self.assertEqual(pending_queue._unfinished_tasks, 0, "second finish did NOT drive the counter negative") + + def test_sender_queue_many_threads_get_and_requeue_never_logs_error(self): + # 100 threads, split into getters and requeuers. Each getter runs + # get() and hands the item to a requeuer via a thread-safe handoff; + # each requeuer takes that item and calls requeue_front() on it. So + # the thread that get() the item is NOT the thread that requeues it -- + # the item crosses thread boundaries. The identity guard is + # thread-agnostic by design, so this must stay clean: no ERROR log, + # no counter drift, no crash. This is the future-proofing proof: a + # legitimate cross-thread get/requeue workload stays clean. + try: + import queue as queue_mod + except ImportError: # Python 2 + import Queue as queue_mod + + pending_queue = SenderQueue( + maxsize=0, # unbounded: requeue never drops for capacity + expiry_seconds=100.0, + on_drop_queue_full=lambda item: self.fail("unexpected queue-full drop"), + on_drop_expired=lambda item: self.fail("unexpected expiry drop"), + ) + + n_items = 200 + n_getters = 50 + n_requeuers = 50 + iterations = 100 + + for i in range(n_items): + pending_queue.put(PendingPayload("item-{}\n".format(i), sender_queue_clock())) + self.assertEqual(pending_queue._unfinished_tasks, n_items) + + handoff = queue_mod.Queue() + SENTINEL = object() + + with self._capture_error_logs() as captured: + def getter(): + for _ in range(iterations): + item = pending_queue.get() + handoff.put(item) + + def requeuer(): + while True: + item = handoff.get() + if item is SENTINEL: + return + pending_queue.requeue_front(item) + + getters = [threading.Thread(target=getter) for _ in range(n_getters)] + requeuers = [threading.Thread(target=requeuer) for _ in range(n_requeuers)] + + # Start requeuers first so they're draining the handoff before + # getters begin filling it; otherwise the handoff could grow + # unbounded and the SenderQueue could drain to empty (getters + # would then block in get() until requeuers put items back). + for t in requeuers: + t.start() + for t in getters: + t.start() + + for t in getters: + t.join() + # All getters done: every real item is either in the handoff or + # already requeued. Sentinels go behind them, so no real item is + # stranded. + for _ in range(n_requeuers): + handoff.put(SENTINEL) + for t in requeuers: + t.join() + + # No ERROR log ever fired: the in-flight guard never tripped even + # though every item crossed from the getter thread to a different + # requeuer thread. + self.assertEqual( + captured, [], + "no error should be logged when get() and requeue_front() run on " + "different threads; got: {!r}".format([r.getMessage() for r in captured]), + ) + + # Every get() was matched by a requeue_front() (no task_done, no + # drops on the unbounded queue), so the queue is fully populated and + # the task counter is unchanged -- no drift. + self.assertEqual(pending_queue.qsize(), n_items, "all items back in the queue") + self.assertEqual( + pending_queue._unfinished_tasks, n_items, + "_unfinished_tasks never drifted under cross-thread get/requeue contention", + ) + + def test_replay_safety_is_carried_by_the_queued_entry_type(self): + # Replay-safety is not a stored flag: an entry subject to expiry is a + # PendingPayload (carrying the enqueued_at it will be judged against), + # and a replay-safe one is the bare packet string. That keeps the + # wrapper -- and the GC traversal it implies -- off replay-safe entries + # entirely. + statsd = DogStatsd(disable_background_sender=False) + statsd.socket = FakeSocket() + + captured = [] + original_put = statsd._queue.put + + def capture_put(item): + if item is not Stop: + captured.append(item) + return original_put(item) + + statsd._queue.put = capture_put + + statsd.increment("no.timestamp") + statsd.gauge_with_timestamp("with.timestamp", 1, int(time.time())) + statsd.wait_for_pending() + + self.assertEqual(len(captured), 2) + + expiring, replay_safe = captured + self.assertIsInstance(expiring, PendingPayload) + self.assertFalse(is_replay_safe(expiring)) + self.assertIsInstance( + expiring.enqueued_at, float, + "an expiring payload needs a real timestamp to be judged against", + ) + + # Deliberately not asserting a concrete string type here: on Python 2 + # the serialized packet is unicode, not str. What matters is that the + # entry is the bare payload rather than a wrapper, which is exactly + # what is_replay_safe()/payload_text() key off. + self.assertNotIsInstance(replay_safe, PendingPayload) + self.assertTrue(is_replay_safe(replay_safe)) + self.assertIs(payload_text(replay_safe), replay_safe) + + # Either form still yields its packet text the same way. + self.assertTrue(payload_text(expiring).startswith("no.timestamp")) + self.assertTrue(payload_text(replay_safe).startswith("with.timestamp")) + + 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. + working_socket = FakeSocket() + attempts = {"count": 0} + + def flaky_get_uds_socket(cls, socket_path, timeout, connect_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)): + statsd = DogStatsd( + socket_path="/tmp/dogstatsd-test-requeue.sock", + disable_telemetry=True, + disable_background_sender=False, + ) + statsd.socket_connect_timeout = 0.05 + + statsd.gauge("eventually.sent", 1) + statsd.wait_for_pending() + + # The payload survived every failed reconnect attempt and was sent + # once a working socket was finally available -- it was never + # dropped as a writer failure or expired out of the queue. + self.assertGreaterEqual(attempts["count"], 4) + self.assertEqual(statsd.packets_dropped_writer, 0) + self.assertEqual(statsd.packets_dropped_expired, 0) + self.assertTrue(working_socket.payloads[0].decode("utf-8").startswith("eventually.sent:1|g")) statsd.stop()