Skip to content
Merged
128 changes: 128 additions & 0 deletions integration-tests/metrics_test.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,128 @@
package tests

import (
"encoding/json"
"net/http"
"net/http/httptest"
"sync"
"testing"
"time"

"github.com/stretchr/testify/assert"

"github.com/upsun/cli/pkg/mockapi"
)

func TestMetricsLatest(t *testing.T) {
authServer := mockapi.NewAuthServer(t)
defer authServer.Close()

apiHandler := mockapi.NewHandler(t)
apiServer := httptest.NewServer(apiHandler)
defer apiServer.Close()

projectID := mockapi.ProjectID()
apiHandler.SetProjects([]*mockapi.Project{{
ID: projectID,
Links: mockapi.MakeHALLinks("self=/projects/"+projectID,
"environments=/projects/"+projectID+"/environments"),
DefaultBranch: "main",
}})

envPath := "/projects/" + projectID + "/environments/main"
main := makeEnv(projectID, "main", "production", "active", nil)
main.Links["#observability-pipeline"] = mockapi.HALLink{HREF: apiServer.URL + envPath + "/observability"}
main.SetCurrentDeployment(&mockapi.Deployment{
WebApps: map[string]mockapi.App{"app": {Name: "app", Type: "golang:1.23", Size: "AUTO"}},
Services: map[string]mockapi.App{"db": {Name: "db", Type: "mariadb:11.4", Size: "AUTO"}},
Workers: map[string]mockapi.Worker{},
Routes: map[string]any{},
Links: mockapi.MakeHALLinks("self=" + envPath + "/deployment/current"),
})
apiHandler.SetEnvironments([]*mockapi.Environment{main})

cpu := func(used, limit float64) map[string]any {
return map[string]any{"cpu_used": map[string]any{"avg": used}, "cpu_limit": map[string]any{"max": limit}}
}
// Timestamps are relative to the request time, as recent points are treated differently.
var (
mu sync.Mutex
now time.Time
data func() []map[string]any
)
ts := func(minutesAgo int) time.Time { return now.Add(-time.Duration(minutesAgo) * time.Minute) }
point := func(minutesAgo int, services map[string]any) map[string]any {
p := map[string]any{"timestamp": ts(minutesAgo).Unix()}
if services != nil {
p["services"] = services
}
return p
}
row := func(minutesAgo int, rest string) string {
mu.Lock()
defer mu.Unlock()
return ts(minutesAgo).Format("2006-01-02T15:04:05+00:00") + "\t" + rest
}

// Modeled on the API: recent points lack services that have not reported yet,
// and the in-progress point has no "services" key at all.
data = func() []map[string]any {
return []map[string]any{
point(3, map[string]any{"app": cpu(0.1, 1), "db": cpu(0.1, 1), "router": cpu(0.01, 0.1)}),
point(2, map[string]any{"app": cpu(0.2, 1), "db": cpu(0.3, 1), "router": cpu(0.02, 0.1)}),
point(1, map[string]any{"db": cpu(0.4, 1)}),
point(0, nil),
}
}
setData := func(d func() []map[string]any) {
mu.Lock()
defer mu.Unlock()
data = d
}
apiHandler.Get(envPath+"/observability/resources/overview", func(w http.ResponseWriter, _ *http.Request) {
mu.Lock()
defer mu.Unlock()
now = time.Now().UTC()
_ = json.NewEncoder(w).Encode(map[string]any{
"_grain": 60,
"_from": ts(10).Unix(),
"_to": now.Unix(),
"data": data(),
})
})

f := newCommandFactory(t, apiServer.URL, authServer.URL)
latest := func() string {
return f.Run("metrics:cpu", "-p", projectID, "-e", "main", "--latest", "--format", "tsv", "--no-header")
}

out := latest()
assertTrimmed(t, row(2, "app\t0.2\t1\t20.0%")+"\n"+
row(2, "db\t0.3\t1\t30.0%")+"\n"+
row(2, "router\t0.02\t0.1\t20.0%"), out)

out = f.Run("metrics:cpu", "-p", projectID, "-e", "main", "--format", "tsv")
assert.Contains(t, out, row(1, "db\t0.4\t1\t40.0%"))

// A service that stopped reporting before the recent points is ignored.
setData(func() []map[string]any {
return []map[string]any{
point(4, map[string]any{"app": cpu(0.1, 1), "db": cpu(0.1, 1)}),
point(3, map[string]any{"app": cpu(0.2, 1)}),
point(2, map[string]any{"app": cpu(0.3, 1)}),
point(1, map[string]any{"app": cpu(0.4, 1)}),
}
})
out = latest()
assertTrimmed(t, row(1, "app\t0.4\t1\t40.0%"), out)

// Older points are not skipped.
setData(func() []map[string]any {
return []map[string]any{
point(4, map[string]any{"app": cpu(0.1, 1), "db": cpu(0.1, 1)}),
point(3, map[string]any{"app": cpu(0.2, 1)}),
}
})
out = latest()
assertTrimmed(t, row(3, "app\t0.2\t1\t20.0%"), out)
}
25 changes: 20 additions & 5 deletions legacy/src/Command/Metrics/MetricsCommandBase.php
Original file line number Diff line number Diff line change
Expand Up @@ -42,6 +42,9 @@ abstract class MetricsCommandBase extends CommandBase
public const MIN_RANGE = 300; // 5 minutes
public const DEFAULT_RANGE = 600;

// Data points that started within this many seconds of now may still be missing services.
private const LATEST_SETTLE_TIME = 120;

/**
* @var bool whether services have been identified that use high memory
*/
Expand Down Expand Up @@ -97,7 +100,7 @@ protected function addMetricsOptions(): self
. "\n" . \sprintf('Minimum <comment>%s</comment>.', $duration->humanize(self::MIN_INTERVAL)),
);
$this->addOption('to', null, InputOption::VALUE_REQUIRED, 'The end time. Defaults to now.');
$this->addOption('latest', '1', InputOption::VALUE_NONE, 'Show only the latest single data point');
$this->addOption('latest', '1', InputOption::VALUE_NONE, 'Show only the latest single data point' . "\n" . 'Points that started in the last 2 minutes are skipped if they have fewer services than an older point.');
$this->addOption('service', 's', InputOption::VALUE_REQUIRED | InputOption::VALUE_IS_ARRAY, 'Filter by service or application name' . "\n" . Wildcard::HELP);
$this->addOption('type', null, InputOption::VALUE_REQUIRED | InputOption::VALUE_IS_ARRAY, 'Filter by service type (if --service is not provided). The version is not required.' . "\n" . Wildcard::HELP);

Expand Down Expand Up @@ -184,14 +187,26 @@ protected function processQuery(InputInterface $input, array $metricTypes, array
throw new \RuntimeException('No data points were found in the metrics response.');
}

// Filter to only the latest timestamp if --latest is given.
// Filter to the latest complete data point if --latest is given.
// Services' metrics can take a minute or two to arrive, so a point
// that started recently is skipped if an older one has more services.
if ($input->getOption('latest')) {
$settledBefore = time() - self::LATEST_SETTLE_TIME;
$latest = null;
foreach (array_reverse($items['data']) as $item) {
if (isset($item['services'])) {
$items['data'] = [$item];
if (empty($item['services'])) {
continue;
}
if ($latest === null || \count($item['services']) > \count($latest['services'])) {
$latest = $item;
}
if ((int) $item['timestamp'] <= $settledBefore) {
Comment thread
pjcdawkins marked this conversation as resolved.
break;
}
}
if ($latest !== null) {
$items['data'] = [$latest];
}
}

// It's possible that there is nothing to display, e.g. if the router
Expand Down Expand Up @@ -298,7 +313,7 @@ protected function validateTimeInput(InputInterface $input): false|TimeSpec
$interval = (int) (new Duration())->toSeconds($intervalString);

if (empty($interval)) {
$this->stdErr->writeln('Invalid --range: <error>' . $intervalString . '</error>');
$this->stdErr->writeln('Invalid --interval: <error>' . $intervalString . '</error>');

return false;
}
Expand Down
Loading