Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
57 changes: 36 additions & 21 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -11,7 +11,7 @@ Here are some key functionalities that this project extends on Locust:
- [custom trends](#custom-trends)
- [timing thresholds](#thresholds)
- [streamlined metric reporting/tagging system](#db-reporting)
(only influxDB is supported right now)
(InfluxDB and Datadog are supported)

## Installation
This package can be installed via pip: `pip install locust-grasshopper`
Expand Down Expand Up @@ -127,6 +127,9 @@ you must specify a host.
- `--grafana_host`: If your grafana is a separate URL from the influxdb, you can
specify it here. If you don't, then the grafana URL will be the same as the
influxdb URL when the grasshopper object generates grafana links.
- Datadog metric reporting requires the `DD_API_KEY` and `DD_ENV` environment
variables. `DD_SITE` defaults to `datadoghq.com`; `DD_SERVICE` and `DD_VERSION`
add stable service and release tags when provided.

<p align="right">(<a href="#top">back to top</a>)</p>

Expand Down Expand Up @@ -282,10 +285,9 @@ in the "checks" table. Here is an example of using a check:

```python
from grasshopper.lib.util.utils import check

...
response = self.client.get(
'https://google.com', name='get google'
)
response = self.client.get("https://google.com", name="get google")
check(
"get google responded with a 200",
response.status_code == 200,
Expand Down Expand Up @@ -381,13 +383,24 @@ This data is also reported to the console at the end of each test.
Additional design details about how a database listener works with grasshopper/locust can be
found in the [Database Listener Design Documentation](./docs/database_listener_design_documentation.md).

When you specify a time series database URL param to `launch_test`, such as
`influx_host`, all metrics will be automatically reported to tables within the `locust`
timeseries database via the specified URL. These tables include:
- `locust_checks`: check name, check passed, etc.
- `locust_events`: test started, test stopped, etc.
- `locust_exceptions`: error messages
- `locust_requests`: HTTP requests and custom trends
When you specify a metrics backend configuration param to `launch_test`, the
corresponding listener will be initialized automatically. For example:
- `influx_host` enables InfluxDB reporting
- `DD_API_KEY` and `DD_ENV` enable Datadog reporting

If both backends are configured, Grasshopper reports to both. The Datadog listener
emits:

- `locust_requests.count`, `locust_requests.response_time`, and
`locust_requests.response_length`
- `locust_requests.error`
- `locust_checks.total`, `locust_checks.passed`, and `locust_checks.failed`
- numeric custom point fields as `<measurement>.<field>`

Datadog reporting is disabled unless both `DD_API_KEY` and `DD_ENV` are configured.
Metric submission runs outside the request path so Datadog API latency does not delay
load test requests. `DD_ENV`, `DD_SERVICE`, and `DD_VERSION` are added as Datadog tags
when configured.

To run the influxdb/grafana locally, you can use the docker-compose file in the example directory:
```shell
Expand All @@ -398,7 +411,7 @@ and then you can access the grafana UI at `localhost`. The default username/pass
To then run a test which reports to this influxdb just add the `--influx_host=localhost` handle.


There are a few ways you can pass in extra tags which
There are a few ways you can pass in extra tags which
will be reported to the time series DB:

1. **HTTP Request Tagging**
Expand All @@ -407,24 +420,26 @@ will be reported to the time series DB:
as a dictionary for the `context` param when making a request. For example:

```python
self.client.get('https://google.com', name='get google', context={'foo':'bar'})
self.client.get("https://google.com", name="get google", context={"foo": "bar"})
```
The tags on this metric would then be: `{'name': 'get google', 'foo': 'bar'}` which
would get forwarded to the database if specified.
The InfluxDB tags on this metric would then be:
`{'name': 'get google', 'foo': 'bar'}`. Datadog request metrics include the request
name, request type, environment, and response code.

2. **Check Tagging**
When defining a check, you can pass in extra tags with the `tags` parameter:
```python
from grasshopper.lib.util.utils import check

...
response = self.client.get(
'https://google.com', name='get google', context={'foo1':'bar1'}
"https://google.com", name="get google", context={"foo1": "bar1"}
)
check(
"get google responded with a 200",
response.status_code == 200,
env=self.environment,
tags = {'foo2': 'bar2'}
"get google responded with a 200",
response.status_code == 200,
env=self.environment,
tags={"foo2": "bar2"},
)
```

Expand Down Expand Up @@ -480,4 +495,4 @@ Don't forget to give the project a star! Thanks again!
6. Push to the Branch (`git push origin feature/AmazingFeature`)
7. Open a Pull Request

<p align="right">(<a href="#top">back to top</a>)</p>
<p align="right">(<a href="#top">back to top</a>)</p>
2 changes: 1 addition & 1 deletion pyproject.toml
Original file line number Diff line number Diff line change
Expand Up @@ -50,7 +50,7 @@ dev = [
"tox == 4.13.0",

# Linting and formatting
"ruff >= 0.14.10",
"ruff == 0.15.8",
]

[project.urls]
Expand Down
25 changes: 25 additions & 0 deletions src/grasshopper/lib/grasshopper.py
Original file line number Diff line number Diff line change
Expand Up @@ -95,6 +95,31 @@ def grafana_configuration(self) -> dict[str, Optional[str]]:
configuration["grafana_host"] = host
return configuration

@property
def datadog_configuration(self) -> dict[str, str | dict[str, str]]:
"""Build Datadog configuration from standard environment variables."""
if not (api_key := os.getenv("DD_API_KEY")) or not (
environment := os.getenv("DD_ENV")
):
return {}

default_tags = {
tag_name: value
for tag_name, value in {
"env": environment,
"service": os.getenv("DD_SERVICE"),
"version": os.getenv("DD_VERSION"),
}.items()
if value
}

return {
"api_key": api_key,
"site": os.getenv("DD_SITE", "datadoghq.com"),
"namespace": "grasshopper",
"default_tags": default_tags,
}

@staticmethod
def launch_test(
weighted_user_classes: Union[
Expand Down
225 changes: 225 additions & 0 deletions src/grasshopper/lib/util/datadog_listener.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,225 @@
"""Datadog metrics listener."""

import json
import logging
import re
from datetime import datetime, timezone
from urllib import error, request

import gevent
from locust.env import Environment

logger = logging.getLogger()


class DatadogApiListener:
"""Forward Locust and custom metrics to the Datadog metrics API."""

def __init__(
self,
environment: Environment,
api_key: str,
site: str = "datadoghq.com",
namespace: str = "grasshopper",
default_tags: dict | None = None,
batch_size: int = 200,
close_timeout: float = 5,
):
"""Register the listener and keep Datadog API settings.

Metrics are cached until `batch_size` series are ready, then submitted from a
gevent greenlet so Datadog network I/O does not slow Locust request handling.
`close_timeout` caps the best-effort shutdown flush.
"""
self.environment = environment
self.api_key = api_key
self.site = site
self.namespace = namespace.strip(".")
self.default_tags = default_tags or {}
self.batch_size = batch_size
self.close_timeout = close_timeout
self.series_buffer = []
self._flush_greenlet = None
environment.events.request.add_listener(self.on_request)

def close(self):
"""Best-effort flush without allowing telemetry to block test shutdown."""
try:
self.environment.events.request.remove_listener(self.on_request)
except (AttributeError, ValueError):
pass

if not self.series_buffer:
return

self._schedule_flush()
if self._flush_greenlet is None:
return
self._flush_greenlet.join(timeout=self.close_timeout)
if not self._flush_greenlet.ready():
logger.warning(
"Datadog metrics flush exceeded %.1f seconds; "
"continuing test shutdown with unsent telemetry.",
self.close_timeout,
)
self._flush_greenlet.kill(block=False)

def on_request(
self,
request_type,
name,
response_time,
response_length,
response,
context,
exception,
**_kwargs,
):
"""Convert a Locust request event into Datadog request metrics."""
status_code = getattr(response, "status_code", None)
timestamp = self._unix_timestamp()
tags = self.default_tags | {
"name": name,
"request_type": request_type,
"environment": getattr(self.environment, "host", None),
"code": str(status_code) if status_code is not None else None,
}

self._buffer_metric("locust_requests.count", 1, "count", tags, timestamp)
self._buffer_metric(
"locust_requests.response_time", response_time, "gauge", tags, timestamp
)
if response_length is not None:
self._buffer_metric(
"locust_requests.response_length",
response_length,
"gauge",
tags,
timestamp,
)
if exception is not None:
self._buffer_metric(
"locust_requests.error",
1,
"count",
tags | {"exception_type": type(exception).__name__},
timestamp,
)

def record_check(
self, check_name: str, check_passed: bool, extra_tags: dict, time=None
):
"""Record total and pass/fail count metrics for one Grasshopper check."""
timestamp = self._unix_timestamp(time)
tags = (
self.default_tags
| extra_tags
| {
"check_name": re.sub(
r"_+", "_", re.sub(r"[^a-z0-9]+", "_", check_name.lower())
).strip("_"),
"environment": getattr(self.environment, "host", None),
}
)
self._buffer_metric("locust_checks.total", 1, "count", tags, timestamp)
metric_suffix = "passed" if check_passed else "failed"
self._buffer_metric(
f"locust_checks.{metric_suffix}", 1, "count", tags, timestamp
)

def record_custom_point(
self, measurement: str, fields: dict, time=None, tags: dict | None = None
):
"""Record numeric custom fields as Datadog gauge metrics."""
timestamp = self._unix_timestamp(time)
metric_tags = self.default_tags | (tags or {})
for field_name, value in fields.items():
if isinstance(value, (int, float)) and not isinstance(value, bool):
self._buffer_metric(
f"{measurement}.{field_name}",
value,
"gauge",
metric_tags,
timestamp,
)

def _buffer_metric(
self,
metric_name: str,
value: float,
metric_type: str,
tags: dict,
timestamp: int | None = None,
):
"""Append one Datadog series payload and flush when the batch is full."""
metric_path = (
f"{self.namespace}.{metric_name}" if self.namespace else metric_name
)
self.series_buffer.append(
{
"metric": metric_path,
"type": metric_type,
"points": [[timestamp or self._unix_timestamp(), value]],
"tags": [
f"{key}:{value}" for key, value in tags.items() if value is not None
],
}
)
if len(self.series_buffer) >= self.batch_size:
self._schedule_flush()

def _schedule_flush(self):
"""Run Datadog HTTP I/O outside the Locust request path."""
if self._flush_greenlet is None or self._flush_greenlet.ready():
self._flush_greenlet = gevent.spawn(self.flush)

def flush(self):
"""Flush buffered metrics to the Datadog API in batches."""
while self.series_buffer:
batch = self.series_buffer[: self.batch_size]
del self.series_buffer[: self.batch_size]
self._submit_series(batch)

def _submit_series(self, series_batch: list[dict]):
"""Submit one already-built Datadog series batch."""
payload = json.dumps({"series": series_batch}).encode("utf-8")
api_request = request.Request(
f"https://api.{self.site}/api/v1/series",
data=payload,
headers={
"Content-Type": "application/json",
"DD-API-KEY": self.api_key,
},
method="POST",
)
try:
with request.urlopen(api_request, timeout=15) as response:
response.read()
logger.info(
"Submitted %s Datadog metric series to `%s`.",
len(series_batch),
self.site,
)
except error.HTTPError as exc:
logger.warning(
"Datadog metrics submission failed with HTTP %s for `%s`: %s",
exc.code,
self.site,
exc.read().decode("utf-8", errors="replace"),
)
except (OSError, TimeoutError, error.URLError) as exc:
logger.warning(
"Failed to submit Datadog metrics batch to `%s`: %s",
self.site,
exc,
)

@staticmethod
def _unix_timestamp(metric_time=None) -> int:
"""Return Datadog-compatible Unix seconds for a metric point."""
timestamp_source = metric_time or datetime.now(timezone.utc)
if isinstance(timestamp_source, datetime):
if timestamp_source.tzinfo is None:
timestamp_source = timestamp_source.replace(tzinfo=timezone.utc)
return int(timestamp_source.timestamp())
return int(timestamp_source)
Loading
Loading