2021-03-06 14:50:29 +00:00
|
|
|
import asyncio
|
|
|
|
import logging
|
2021-03-21 14:22:55 +00:00
|
|
|
import socket
|
2021-03-06 14:50:29 +00:00
|
|
|
import typing
|
|
|
|
|
|
|
|
|
2021-05-09 19:48:18 +00:00
|
|
|
class ThrottleGuard:
|
|
|
|
"""Prevent code from executing repeatedly.
|
|
|
|
|
|
|
|
:param threshold: guarding threshold
|
|
|
|
|
|
|
|
This abstraction allows code to execute the first "threshold"
|
|
|
|
times and then only once per "threshold" times afterwards. Use
|
|
|
|
it to ensure that log statements are continuously written during
|
|
|
|
persistent error conditions. The goal is to provide regular
|
|
|
|
feedback while limiting the amount of log spam.
|
|
|
|
|
|
|
|
The following snippet will log the first 100 failures and then
|
|
|
|
once every 100 failures thereafter:
|
|
|
|
|
|
|
|
.. code-block:: python
|
|
|
|
|
|
|
|
executions = 0
|
|
|
|
guard = ThrottleGuard(100)
|
|
|
|
for _ in range(1000):
|
|
|
|
if guard.allow_execution():
|
|
|
|
executions += 1
|
|
|
|
logging.info('called %s times instead of %s times',
|
|
|
|
executions, guard.counter)
|
|
|
|
|
|
|
|
"""
|
|
|
|
def __init__(self, threshold: int):
|
|
|
|
self.counter = 0
|
|
|
|
self.threshold = threshold
|
|
|
|
|
|
|
|
def allow_execution(self) -> bool:
|
|
|
|
"""Should this execution be allowed?"""
|
|
|
|
self.counter += 1
|
|
|
|
allow = (self.counter < self.threshold
|
|
|
|
or (self.counter % self.threshold) == 0)
|
|
|
|
return allow
|
|
|
|
|
|
|
|
def reset(self) -> None:
|
|
|
|
"""Reset counter after error has resolved."""
|
|
|
|
self.counter = 0
|
|
|
|
|
|
|
|
|
2021-04-12 12:08:07 +00:00
|
|
|
class AbstractConnector:
|
|
|
|
"""StatsD connector that does not send metrics or connect.
|
|
|
|
|
|
|
|
Use this connector when you want to maintain the application
|
|
|
|
interface without doing any real work.
|
|
|
|
|
|
|
|
"""
|
|
|
|
async def start(self) -> None:
|
|
|
|
pass
|
|
|
|
|
|
|
|
async def stop(self) -> None:
|
|
|
|
pass
|
|
|
|
|
|
|
|
def inject_metric(self, path: str, value: str, type_code: str) -> None:
|
|
|
|
pass
|
|
|
|
|
|
|
|
def incr(self, path: str, value: int = 1) -> None:
|
|
|
|
"""Increment a counter metric.
|
|
|
|
|
|
|
|
:param path: counter to increment
|
|
|
|
:param value: amount to increment the counter by
|
|
|
|
|
|
|
|
"""
|
|
|
|
self.inject_metric(f'counters.{path}', str(value), 'c')
|
|
|
|
|
|
|
|
def decr(self, path: str, value: int = 1) -> None:
|
|
|
|
"""Decrement a counter metric.
|
|
|
|
|
|
|
|
:param path: counter to decrement
|
|
|
|
:param value: amount to decrement the counter by
|
|
|
|
|
|
|
|
This is equivalent to ``self.incr(path, -value)``.
|
|
|
|
|
|
|
|
"""
|
|
|
|
self.inject_metric(f'counters.{path}', str(-value), 'c')
|
|
|
|
|
|
|
|
def gauge(self, path: str, value: int, delta: bool = False) -> None:
|
|
|
|
"""Manipulate a gauge metric.
|
|
|
|
|
|
|
|
:param path: gauge to adjust
|
|
|
|
:param value: value to send
|
|
|
|
:param delta: is this an adjustment of the gauge?
|
|
|
|
|
|
|
|
If the `delta` parameter is ``False`` (or omitted), then
|
|
|
|
`value` is the new value to set the gauge to. Otherwise,
|
|
|
|
`value` is an adjustment for the current gauge.
|
|
|
|
|
|
|
|
"""
|
|
|
|
if delta:
|
|
|
|
payload = f'{value:+d}'
|
|
|
|
else:
|
|
|
|
payload = str(value)
|
|
|
|
self.inject_metric(f'gauges.{path}', payload, 'g')
|
|
|
|
|
|
|
|
def timing(self, path: str, seconds: float) -> None:
|
|
|
|
"""Send a timer metric.
|
|
|
|
|
|
|
|
:param path: timer to append a value to
|
|
|
|
:param seconds: number of **seconds** to record
|
|
|
|
|
|
|
|
"""
|
|
|
|
self.inject_metric(f'timers.{path}', str(seconds * 1000.0), 'ms')
|
|
|
|
|
|
|
|
|
|
|
|
class Connector(AbstractConnector):
|
2021-03-08 12:27:40 +00:00
|
|
|
"""Sends metrics to a statsd server.
|
|
|
|
|
|
|
|
:param host: statsd server to send metrics to
|
2021-03-21 14:22:55 +00:00
|
|
|
:param port: socket port that the server is listening on
|
|
|
|
:keyword ip_protocol: IP protocol to use for the underlying
|
2021-03-21 21:45:23 +00:00
|
|
|
socket -- either ``socket.IPPROTO_TCP`` for TCP or
|
|
|
|
``socket.IPPROTO_UDP`` for UDP sockets.
|
2021-03-30 11:56:20 +00:00
|
|
|
:keyword prefix: optional string to prepend to metric paths
|
2021-03-11 12:31:24 +00:00
|
|
|
:param kwargs: additional keyword parameters are passed
|
|
|
|
to the :class:`.Processor` initializer
|
2021-03-08 12:27:40 +00:00
|
|
|
|
2021-03-21 14:22:55 +00:00
|
|
|
This class maintains a connection to a statsd server and
|
2021-03-08 12:27:40 +00:00
|
|
|
sends metric lines to it asynchronously. You must call the
|
|
|
|
:meth:`start` method when your application is starting. It
|
|
|
|
creates a :class:`~asyncio.Task` that manages the connection
|
|
|
|
to the statsd server. You must also call :meth:`.stop` before
|
|
|
|
terminating to ensure that all metrics are flushed to the
|
|
|
|
statsd server.
|
|
|
|
|
2021-03-30 11:56:20 +00:00
|
|
|
Metrics are optionally prefixed with :attr:`prefix` before the
|
|
|
|
metric type prefix. This *should* be used to prevent metrics
|
|
|
|
from being overwritten when multiple applications share a StatsD
|
|
|
|
instance. Each metric type is also prefixed by one of the
|
|
|
|
following strings based on the metric type:
|
|
|
|
|
|
|
|
+-------------------+---------------+-----------+
|
|
|
|
| Method call | Prefix | Type code |
|
|
|
|
+-------------------+---------------+-----------+
|
|
|
|
| :meth:`.incr` | ``counters.`` | ``c`` |
|
|
|
|
+-------------------+---------------+-----------+
|
|
|
|
| :meth:`.decr` | ``counters.`` | ``c`` |
|
|
|
|
+-------------------+---------------+-----------+
|
|
|
|
| :meth:`.gauge` | ``gauges.`` | ``g`` |
|
|
|
|
+-------------------+---------------+-----------+
|
|
|
|
| :meth:`.timing` | ``timers.`` | ``ms`` |
|
|
|
|
+-------------------+---------------+-----------+
|
|
|
|
|
|
|
|
When the connector is *should_terminate*, metric payloads are
|
|
|
|
sent by calling the :meth:`.inject_metric` method. The payloads
|
|
|
|
are stored in an internal queue that is consumed whenever the
|
2021-03-08 12:27:40 +00:00
|
|
|
connection to the server is active.
|
|
|
|
|
2021-03-30 11:56:20 +00:00
|
|
|
.. attribute:: prefix
|
|
|
|
:type: str
|
|
|
|
|
|
|
|
String to prefix to all metrics *before* the metric type prefix.
|
|
|
|
|
2021-03-08 12:27:40 +00:00
|
|
|
.. attribute:: processor
|
|
|
|
:type: Processor
|
|
|
|
|
|
|
|
The statsd processor that maintains the connection and
|
|
|
|
sends the metric payloads.
|
|
|
|
|
|
|
|
"""
|
2021-03-30 11:22:22 +00:00
|
|
|
logger: logging.Logger
|
2021-03-30 11:56:20 +00:00
|
|
|
prefix: str
|
2021-03-26 10:40:27 +00:00
|
|
|
processor: 'Processor'
|
|
|
|
|
|
|
|
def __init__(self,
|
|
|
|
host: str,
|
|
|
|
port: int = 8125,
|
2021-03-30 11:56:20 +00:00
|
|
|
*,
|
|
|
|
prefix: str = '',
|
2021-03-26 10:40:27 +00:00
|
|
|
**kwargs: typing.Any) -> None:
|
2021-04-12 12:08:07 +00:00
|
|
|
super().__init__()
|
2021-03-30 11:22:22 +00:00
|
|
|
self.logger = logging.getLogger(__package__).getChild('Connector')
|
2021-03-30 11:56:20 +00:00
|
|
|
self.prefix = f'{prefix}.' if prefix else prefix
|
2021-03-11 12:31:24 +00:00
|
|
|
self.processor = Processor(host=host, port=port, **kwargs)
|
2021-05-09 19:48:18 +00:00
|
|
|
self._enqueue_log_guard = ThrottleGuard(100)
|
2021-03-26 10:40:27 +00:00
|
|
|
self._processor_task: typing.Optional[asyncio.Task[None]] = None
|
2021-03-08 12:27:40 +00:00
|
|
|
|
2021-03-26 10:40:27 +00:00
|
|
|
async def start(self) -> None:
|
2021-03-09 20:06:23 +00:00
|
|
|
"""Start the processor in the background.
|
|
|
|
|
|
|
|
This is a *blocking* method and does not return until the
|
|
|
|
processor task is actually running.
|
|
|
|
|
|
|
|
"""
|
2021-03-08 12:27:40 +00:00
|
|
|
self._processor_task = asyncio.create_task(self.processor.run())
|
2021-03-09 20:06:23 +00:00
|
|
|
await self.processor.running.wait()
|
2021-03-08 12:27:40 +00:00
|
|
|
|
2021-03-26 10:40:27 +00:00
|
|
|
async def stop(self) -> None:
|
2021-03-08 12:27:40 +00:00
|
|
|
"""Stop the background processor.
|
|
|
|
|
|
|
|
Items that are currently in the queue will be flushed to
|
|
|
|
the statsd server if possible. This is a *blocking* method
|
|
|
|
and does not return until the background processor has
|
|
|
|
stopped.
|
|
|
|
|
|
|
|
"""
|
|
|
|
await self.processor.stop()
|
|
|
|
|
2021-03-26 10:40:27 +00:00
|
|
|
def inject_metric(self, path: str, value: str, type_code: str) -> None:
|
2021-03-08 12:27:40 +00:00
|
|
|
"""Send a metric to the statsd server.
|
|
|
|
|
|
|
|
:param path: formatted metric name
|
2021-03-26 10:40:27 +00:00
|
|
|
:param value: formatted metric value
|
2021-03-08 12:27:40 +00:00
|
|
|
:param type_code: type of the metric to send
|
|
|
|
|
|
|
|
This method formats the payload and inserts it on the
|
|
|
|
internal queue for future processing.
|
|
|
|
|
|
|
|
"""
|
2021-03-30 11:56:20 +00:00
|
|
|
payload = f'{self.prefix}{path}:{value}|{type_code}'
|
2021-03-30 11:22:22 +00:00
|
|
|
try:
|
|
|
|
self.processor.enqueue(payload.encode('utf-8'))
|
2021-05-09 19:48:18 +00:00
|
|
|
self._enqueue_log_guard.reset()
|
2021-03-30 11:22:22 +00:00
|
|
|
except asyncio.QueueFull:
|
2021-05-09 19:48:18 +00:00
|
|
|
if self._enqueue_log_guard.allow_execution():
|
|
|
|
self.logger.warning('statsd queue is full, discarding metric')
|
2021-03-08 12:27:40 +00:00
|
|
|
|
|
|
|
|
2021-03-26 10:40:27 +00:00
|
|
|
class StatsdProtocol(asyncio.BaseProtocol):
|
2021-03-19 22:19:03 +00:00
|
|
|
"""Common interface for backend protocols/transports.
|
|
|
|
|
|
|
|
UDP and TCP transports have different interfaces (sendto vs write)
|
|
|
|
so this class adapts them to a common protocol that our code
|
|
|
|
can depend on.
|
|
|
|
|
|
|
|
.. attribute:: buffered_data
|
|
|
|
:type: bytes
|
|
|
|
|
|
|
|
Bytes that are buffered due to low-level transport failures.
|
|
|
|
Since protocols & transports are created anew with each connect
|
|
|
|
attempt, the :class:`.Processor` instance ensures that data
|
|
|
|
buffered on a transport is copied over to the new transport
|
|
|
|
when creating a connection.
|
|
|
|
|
|
|
|
.. attribute:: connected
|
|
|
|
:type: asyncio.Event
|
|
|
|
|
|
|
|
Is the protocol currently connected?
|
|
|
|
|
|
|
|
"""
|
2021-03-26 10:40:27 +00:00
|
|
|
buffered_data: bytes
|
|
|
|
ip_protocol: int = socket.IPPROTO_NONE
|
2021-03-19 22:19:03 +00:00
|
|
|
logger: logging.Logger
|
2021-03-26 10:40:27 +00:00
|
|
|
transport: typing.Optional[asyncio.BaseTransport]
|
2021-03-19 22:19:03 +00:00
|
|
|
|
2021-03-26 10:40:27 +00:00
|
|
|
def __init__(self) -> None:
|
2021-03-19 22:19:03 +00:00
|
|
|
self.buffered_data = b''
|
|
|
|
self.connected = asyncio.Event()
|
|
|
|
self.logger = logging.getLogger(__package__).getChild(
|
|
|
|
self.__class__.__name__)
|
|
|
|
self.transport = None
|
|
|
|
|
|
|
|
def send(self, metric: bytes) -> None:
|
|
|
|
"""Send a metric payload over the transport."""
|
|
|
|
raise NotImplementedError()
|
|
|
|
|
|
|
|
async def shutdown(self) -> None:
|
|
|
|
"""Shutdown the transport and wait for it to close."""
|
|
|
|
raise NotImplementedError()
|
|
|
|
|
2021-03-26 10:40:27 +00:00
|
|
|
def connection_made(self, transport: asyncio.BaseTransport) -> None:
|
2021-03-19 22:19:03 +00:00
|
|
|
"""Capture the new transport and set the connected event."""
|
2021-03-21 21:45:23 +00:00
|
|
|
# NB - this will return a 4-part tuple in some cases
|
|
|
|
server, port = transport.get_extra_info('peername')[:2]
|
2021-03-19 22:19:03 +00:00
|
|
|
self.logger.info('connected to statsd %s:%s', server, port)
|
|
|
|
self.transport = transport
|
2021-03-26 10:40:27 +00:00
|
|
|
self.transport.set_protocol(self)
|
2021-03-19 22:19:03 +00:00
|
|
|
self.connected.set()
|
|
|
|
|
2021-03-26 10:40:27 +00:00
|
|
|
def connection_lost(self, exc: typing.Optional[Exception]) -> None:
|
2021-03-19 22:19:03 +00:00
|
|
|
"""Clear the connected event."""
|
|
|
|
self.logger.warning('statsd server connection lost: %s', exc)
|
|
|
|
self.connected.clear()
|
|
|
|
|
|
|
|
|
|
|
|
class TCPProtocol(StatsdProtocol, asyncio.Protocol):
|
|
|
|
"""StatsdProtocol implementation over a TCP/IP connection."""
|
2021-03-26 10:40:27 +00:00
|
|
|
ip_protocol = socket.IPPROTO_TCP
|
2021-03-19 22:19:03 +00:00
|
|
|
transport: asyncio.WriteTransport
|
|
|
|
|
2021-03-26 10:40:27 +00:00
|
|
|
def eof_received(self) -> None:
|
2021-03-19 22:19:03 +00:00
|
|
|
self.logger.warning('received EOF from statsd server')
|
|
|
|
self.connected.clear()
|
|
|
|
|
|
|
|
def send(self, metric: bytes) -> None:
|
|
|
|
"""Send `metric` to the server.
|
|
|
|
|
|
|
|
If sending the metric fails, it will be saved in
|
|
|
|
``self.buffered_data``. The processor will save and
|
|
|
|
restore the buffered data if it needs to create a
|
|
|
|
new protocol object.
|
|
|
|
|
|
|
|
"""
|
|
|
|
if not self.buffered_data and not metric:
|
|
|
|
return
|
|
|
|
|
|
|
|
self.buffered_data = self.buffered_data + metric + b'\n'
|
|
|
|
while (self.transport is not None and self.connected.is_set()
|
|
|
|
and self.buffered_data):
|
|
|
|
line, maybe_nl, rest = self.buffered_data.partition(b'\n')
|
|
|
|
line += maybe_nl
|
|
|
|
self.transport.write(line)
|
|
|
|
if self.transport.is_closing():
|
|
|
|
self.logger.warning('transport closed during write')
|
|
|
|
break
|
|
|
|
self.buffered_data = rest
|
|
|
|
|
|
|
|
async def shutdown(self) -> None:
|
|
|
|
"""Close the transport after flushing any outstanding data."""
|
|
|
|
self.logger.info('shutting down')
|
|
|
|
if self.connected.is_set():
|
|
|
|
self.send(b'') # flush buffered data
|
|
|
|
self.transport.close()
|
|
|
|
while self.connected.is_set():
|
|
|
|
await asyncio.sleep(0.1)
|
|
|
|
|
|
|
|
|
2021-03-21 14:22:55 +00:00
|
|
|
class UDPProtocol(StatsdProtocol, asyncio.DatagramProtocol):
|
|
|
|
"""StatsdProtocol implementation over a UDP/IP connection."""
|
2021-03-26 10:40:27 +00:00
|
|
|
ip_protocol = socket.IPPROTO_UDP
|
2021-03-21 14:22:55 +00:00
|
|
|
transport: asyncio.DatagramTransport
|
|
|
|
|
|
|
|
def send(self, metric: bytes) -> None:
|
|
|
|
if metric:
|
|
|
|
self.transport.sendto(metric)
|
|
|
|
|
|
|
|
async def shutdown(self) -> None:
|
|
|
|
self.logger.info('shutting down')
|
|
|
|
self.transport.close()
|
|
|
|
|
|
|
|
|
2021-03-19 22:19:03 +00:00
|
|
|
class Processor:
|
2021-03-08 12:27:40 +00:00
|
|
|
"""Maintains the statsd connection and sends metric payloads.
|
|
|
|
|
|
|
|
:param host: statsd server to send metrics to
|
|
|
|
:param port: TCP port that the server is listening on
|
2021-03-30 11:22:22 +00:00
|
|
|
:param max_queue_size: only allow this many elements to be
|
|
|
|
stored in the queue before discarding metrics
|
2021-03-11 12:31:24 +00:00
|
|
|
:param reconnect_sleep: number of seconds to sleep after socket
|
|
|
|
error occurs when connecting
|
|
|
|
:param wait_timeout: number os seconds to wait for a message to
|
|
|
|
arrive on the queue
|
2021-03-08 12:27:40 +00:00
|
|
|
|
|
|
|
This class implements :class:`~asyncio.Protocol` for the statsd
|
|
|
|
TCP connection. The :meth:`.run` method is run as a background
|
|
|
|
:class:`~asyncio.Task` that consumes payloads from an internal
|
|
|
|
queue, connects to the TCP server as required, and sends the
|
|
|
|
already formatted payloads.
|
|
|
|
|
|
|
|
.. attribute:: host
|
|
|
|
:type: str
|
|
|
|
|
|
|
|
IP address or DNS name for the statsd server to send metrics to
|
|
|
|
|
|
|
|
.. attribute:: port
|
|
|
|
:type: int
|
|
|
|
|
|
|
|
TCP port number that the statsd server is listening on
|
|
|
|
|
|
|
|
.. attribute:: should_terminate
|
|
|
|
:type: bool
|
|
|
|
|
|
|
|
Flag that controls whether the background task is active or
|
|
|
|
not. This flag is set to :data:`False` when the task is started.
|
|
|
|
Setting it to :data:`True` will cause the task to shutdown in
|
|
|
|
an orderly fashion.
|
|
|
|
|
|
|
|
.. attribute:: queue
|
|
|
|
:type: asyncio.Queue
|
|
|
|
|
|
|
|
Formatted metric payloads to send to the statsd server. Enqueue
|
|
|
|
payloads to send them to the server.
|
|
|
|
|
2021-03-09 20:06:23 +00:00
|
|
|
.. attribute:: running
|
|
|
|
:type: asyncio.Event
|
|
|
|
|
|
|
|
Is the background task currently running? This is the event that
|
|
|
|
:meth:`.run` sets when it starts and it remains set until the task
|
|
|
|
exits.
|
|
|
|
|
2021-03-08 12:27:40 +00:00
|
|
|
.. attribute:: stopped
|
|
|
|
:type: asyncio.Event
|
|
|
|
|
|
|
|
Is the background task currently stopped? This is the event that
|
|
|
|
:meth:`.run` sets when it exits and that :meth:`.stop` blocks on
|
|
|
|
until the task stops.
|
|
|
|
|
|
|
|
"""
|
2021-03-19 22:19:03 +00:00
|
|
|
|
2021-03-26 10:40:27 +00:00
|
|
|
logger: logging.Logger
|
|
|
|
protocol: typing.Optional[StatsdProtocol]
|
2021-03-30 12:09:53 +00:00
|
|
|
queue: asyncio.Queue
|
2021-03-21 14:22:55 +00:00
|
|
|
_create_transport: typing.Callable[[], typing.Coroutine[
|
|
|
|
typing.Any, typing.Any, typing.Tuple[asyncio.BaseTransport,
|
|
|
|
StatsdProtocol]]]
|
2021-03-19 22:19:03 +00:00
|
|
|
|
2021-03-11 12:31:24 +00:00
|
|
|
def __init__(self,
|
|
|
|
*,
|
2021-03-26 10:40:27 +00:00
|
|
|
host: str,
|
2021-03-11 12:31:24 +00:00
|
|
|
port: int = 8125,
|
2021-03-21 14:22:55 +00:00
|
|
|
ip_protocol: int = socket.IPPROTO_TCP,
|
2021-03-30 11:22:22 +00:00
|
|
|
max_queue_size: int = 1000,
|
|
|
|
reconnect_sleep: float = 1.0,
|
2021-03-26 10:40:27 +00:00
|
|
|
wait_timeout: float = 0.1) -> None:
|
2021-03-08 12:27:40 +00:00
|
|
|
super().__init__()
|
2021-03-09 20:06:23 +00:00
|
|
|
if not host:
|
|
|
|
raise RuntimeError('host must be set')
|
2021-03-19 22:19:03 +00:00
|
|
|
try:
|
|
|
|
port = int(port)
|
|
|
|
if not port or port < 1:
|
|
|
|
raise RuntimeError(
|
|
|
|
f'port must be a positive integer: {port!r}')
|
2021-03-21 14:22:55 +00:00
|
|
|
except (TypeError, ValueError):
|
2021-03-19 22:19:03 +00:00
|
|
|
raise RuntimeError(f'port must be a positive integer: {port!r}')
|
2021-03-09 20:06:23 +00:00
|
|
|
|
2021-03-21 14:22:55 +00:00
|
|
|
transport_creators = {
|
|
|
|
socket.IPPROTO_TCP: self._create_tcp_transport,
|
|
|
|
socket.IPPROTO_UDP: self._create_udp_transport,
|
|
|
|
}
|
|
|
|
try:
|
2021-03-26 10:40:27 +00:00
|
|
|
factory = transport_creators[ip_protocol]
|
2021-03-21 14:22:55 +00:00
|
|
|
except KeyError:
|
|
|
|
raise RuntimeError(f'ip_protocol {ip_protocol} is not supported')
|
2021-03-26 10:40:27 +00:00
|
|
|
else:
|
|
|
|
self._create_transport = factory # type: ignore
|
2021-03-21 14:22:55 +00:00
|
|
|
|
2021-03-06 14:50:29 +00:00
|
|
|
self.host = host
|
|
|
|
self.port = port
|
2021-03-21 14:22:55 +00:00
|
|
|
self._ip_protocol = ip_protocol
|
2021-05-09 19:48:18 +00:00
|
|
|
self._connect_log_guard = ThrottleGuard(100)
|
2021-03-11 12:31:24 +00:00
|
|
|
self._reconnect_sleep = reconnect_sleep
|
|
|
|
self._wait_timeout = wait_timeout
|
2021-03-06 14:50:29 +00:00
|
|
|
|
2021-03-09 20:06:23 +00:00
|
|
|
self.running = asyncio.Event()
|
2021-03-08 12:27:40 +00:00
|
|
|
self.stopped = asyncio.Event()
|
|
|
|
self.stopped.set()
|
2021-03-06 14:50:29 +00:00
|
|
|
self.logger = logging.getLogger(__package__).getChild('Processor')
|
2021-03-08 12:27:40 +00:00
|
|
|
self.should_terminate = False
|
2021-03-19 22:19:03 +00:00
|
|
|
self.protocol = None
|
2021-03-30 11:22:22 +00:00
|
|
|
self.queue = asyncio.Queue(maxsize=max_queue_size)
|
2021-03-06 14:50:29 +00:00
|
|
|
|
2021-03-19 22:19:03 +00:00
|
|
|
@property
|
|
|
|
def connected(self) -> bool:
|
|
|
|
"""Is the processor connected?"""
|
2021-03-26 10:40:27 +00:00
|
|
|
return self.protocol is not None and self.protocol.connected.is_set()
|
|
|
|
|
|
|
|
def enqueue(self, metric: bytes) -> None:
|
|
|
|
self.queue.put_nowait(metric)
|
2021-03-19 22:19:03 +00:00
|
|
|
|
2021-03-26 10:40:27 +00:00
|
|
|
async def run(self) -> None:
|
2021-03-08 12:27:40 +00:00
|
|
|
"""Maintains the connection and processes metric payloads."""
|
2021-03-09 20:06:23 +00:00
|
|
|
self.running.set()
|
2021-03-08 12:27:40 +00:00
|
|
|
self.stopped.clear()
|
|
|
|
self.should_terminate = False
|
|
|
|
while not self.should_terminate:
|
2021-03-06 14:50:29 +00:00
|
|
|
try:
|
|
|
|
await self._connect_if_necessary()
|
2021-03-19 22:19:03 +00:00
|
|
|
if self.connected:
|
2021-03-11 12:03:18 +00:00
|
|
|
await self._process_metric()
|
2021-03-06 14:50:29 +00:00
|
|
|
except asyncio.CancelledError:
|
|
|
|
self.logger.info('task cancelled, exiting')
|
2021-03-21 14:22:55 +00:00
|
|
|
self.should_terminate = True
|
|
|
|
except Exception as error:
|
|
|
|
self.logger.exception('unexpected error occurred: %s', error)
|
|
|
|
self.should_terminate = True
|
2021-03-06 14:50:29 +00:00
|
|
|
|
2021-03-08 12:27:40 +00:00
|
|
|
self.should_terminate = True
|
2021-03-07 19:35:42 +00:00
|
|
|
self.logger.info('loop finished with %d metrics in the queue',
|
2021-03-08 12:27:40 +00:00
|
|
|
self.queue.qsize())
|
2021-03-19 22:19:03 +00:00
|
|
|
if self.connected:
|
|
|
|
num_ready = max(self.queue.qsize(), 1)
|
2021-03-07 19:35:42 +00:00
|
|
|
self.logger.info('draining %d metrics', num_ready)
|
|
|
|
for _ in range(num_ready):
|
|
|
|
await self._process_metric()
|
2021-03-06 14:50:29 +00:00
|
|
|
self.logger.debug('closing transport')
|
2021-03-19 22:19:03 +00:00
|
|
|
if self.protocol is not None:
|
|
|
|
await self.protocol.shutdown()
|
2021-03-06 14:50:29 +00:00
|
|
|
|
2021-03-07 19:35:42 +00:00
|
|
|
self.logger.info('processor is exiting')
|
2021-03-09 20:06:23 +00:00
|
|
|
self.running.clear()
|
2021-03-08 12:27:40 +00:00
|
|
|
self.stopped.set()
|
2021-03-06 14:50:29 +00:00
|
|
|
|
2021-03-26 10:40:27 +00:00
|
|
|
async def stop(self) -> None:
|
2021-03-08 12:27:40 +00:00
|
|
|
"""Stop the processor.
|
2021-03-06 14:50:29 +00:00
|
|
|
|
2021-03-08 12:27:40 +00:00
|
|
|
This is an asynchronous but blocking method. It does not
|
|
|
|
return until enqueued metrics are flushed and the processor
|
|
|
|
connection is closed.
|
|
|
|
|
|
|
|
"""
|
|
|
|
self.should_terminate = True
|
|
|
|
await self.stopped.wait()
|
2021-03-07 19:35:42 +00:00
|
|
|
|
2021-03-21 14:22:55 +00:00
|
|
|
async def _create_tcp_transport(
|
|
|
|
self) -> typing.Tuple[asyncio.BaseTransport, StatsdProtocol]:
|
|
|
|
t, p = await asyncio.get_running_loop().create_connection(
|
|
|
|
protocol_factory=TCPProtocol, host=self.host, port=self.port)
|
|
|
|
return t, typing.cast(StatsdProtocol, p)
|
|
|
|
|
|
|
|
async def _create_udp_transport(
|
|
|
|
self) -> typing.Tuple[asyncio.BaseTransport, StatsdProtocol]:
|
|
|
|
t, p = await asyncio.get_running_loop().create_datagram_endpoint(
|
|
|
|
protocol_factory=UDPProtocol,
|
|
|
|
remote_addr=(self.host, self.port),
|
|
|
|
reuse_port=True)
|
|
|
|
return t, typing.cast(StatsdProtocol, p)
|
|
|
|
|
2021-03-26 10:40:27 +00:00
|
|
|
async def _connect_if_necessary(self) -> None:
|
2021-03-19 22:19:03 +00:00
|
|
|
if self.protocol is not None:
|
2021-03-06 14:50:29 +00:00
|
|
|
try:
|
2021-03-19 22:19:03 +00:00
|
|
|
await asyncio.wait_for(self.protocol.connected.wait(),
|
|
|
|
self._wait_timeout)
|
|
|
|
except asyncio.TimeoutError:
|
|
|
|
self.logger.debug('protocol is no longer connected')
|
|
|
|
|
|
|
|
if not self.connected:
|
|
|
|
try:
|
|
|
|
buffered_data = b''
|
|
|
|
if self.protocol is not None:
|
|
|
|
buffered_data = self.protocol.buffered_data
|
2021-05-09 19:48:18 +00:00
|
|
|
|
2021-03-26 10:40:27 +00:00
|
|
|
t, p = await self._create_transport() # type: ignore[misc]
|
|
|
|
transport, self.protocol = t, p
|
2021-03-19 22:19:03 +00:00
|
|
|
self.protocol.buffered_data = buffered_data
|
2021-05-09 19:48:18 +00:00
|
|
|
self.logger.info(
|
|
|
|
'connection established to %s after %s attempts',
|
|
|
|
transport.get_extra_info('peername'),
|
|
|
|
self._connect_log_guard.counter)
|
|
|
|
self._connect_log_guard.reset()
|
2021-03-06 14:50:29 +00:00
|
|
|
except IOError as error:
|
2021-05-09 19:48:18 +00:00
|
|
|
if self._connect_log_guard.allow_execution():
|
|
|
|
self.logger.warning(
|
|
|
|
'connection to %s:%s failed: %s (%s attempts)',
|
|
|
|
self.host, self.port, error,
|
|
|
|
self._connect_log_guard.counter)
|
2021-03-11 12:31:24 +00:00
|
|
|
await asyncio.sleep(self._reconnect_sleep)
|
2021-03-07 19:35:42 +00:00
|
|
|
|
2021-03-26 10:40:27 +00:00
|
|
|
async def _process_metric(self) -> None:
|
2021-03-19 22:19:03 +00:00
|
|
|
try:
|
|
|
|
metric = await asyncio.wait_for(self.queue.get(),
|
|
|
|
self._wait_timeout)
|
|
|
|
self.logger.debug('received %r from queue', metric)
|
|
|
|
self.queue.task_done()
|
|
|
|
except asyncio.TimeoutError:
|
|
|
|
# we still want to invoke the protocol send in case
|
|
|
|
# it has queued metrics to send
|
|
|
|
metric = b''
|
|
|
|
|
2021-03-26 10:40:27 +00:00
|
|
|
assert self.protocol is not None # AFAICT, this cannot happen
|
2021-03-21 14:22:55 +00:00
|
|
|
self.protocol.send(metric)
|