mirror of
https://github.com/agent-substrate/substrate.git
synced 2026-10-02 03:24:42 +08:00
Closes the remaining items in #1640: the atelet side of multi-actor telemetry. Since #1836 an ateom returns one `GetActiveWorkloadStats` entry per hosted actor. atelet's sweep already folded every entry, but nothing pinned the behavior at occupancy above one, and reviewing that fold surfaced three CPU-accounting defects: two in how baselines survive a pending sweep, one in the decrease rule. ## What changed **Tests pin the multi-actor fold.** One worker, several actors: same-template entries sum under one key, a different template gets its own, a pending sibling contributes nothing to aggregates or events, CPU baselines are per actor (independent advance, independent epoch reset), one probe serves the whole response. Each test was verified to fail against the corresponding mutation. **Three CPU-accounting defects found and fixed along the way:** 1. *Undercount.* A pending entry (boot, restore, or a micro-VM guest the ateom did not reach in its 45s budget) never wrote the actor's baseline, and `lastCPU` is rebuilt each sweep — so measured → pending → measured lost the whole interval. On the guest-agent source the counter survives a restore, so that was real CPU time dropped, and on a full micro-VM worker it would happen routinely. Pending now carries the baseline forward; only actors absent from every response are dropped. 2. *Overcount opened by fix 1.* Carrying the baseline unconditionally meant that during a restore spanning one sweep (source still reports the actor measured at C1, destination reports it pending), goroutine order could let the stale C0 overwrite C1, and the next sweep charged the interval twice. The pending branch now writes only if no worker has measured the actor this sweep. 3. *Overcount on a guest-agent decrease.* A CPU decrease was charged as the whole new value. That is right for a cgroup counter, which restarts at zero on restore. It is not always right for a guest-agent counter: a FULL or DATA_ON_GOLDEN restore resumes a guest at that guest's counter (after a golden restore, main billed the golden guest's CPU to the actor), a cold boot restarts it at zero, and a container exit lowers the sum. A sample cannot tell these apart, so a guest-agent decrease now charges nothing and re-baselines. A cgroup decrease still charges the new value, with or without a pending sweep between. **Trade-off, new versus main:** a guest-agent decrease loses the usage from the event to the first sample. Main counted it correctly only after a cold boot (a DATA restore, or a Run after a crash); in every other case (golden restore, FULL restore from an older snapshot, container exit) main overcounted it. Remaining limit, documented in the proto: an epoch that starts *above* the last reported value cannot be detected from consecutive samples. Main has this for guest-agent restores between sweeps; carrying the baseline extends it to restores that span a pending sweep. The overcount is bounded by the golden guest's counter minus the actor's last value, so it needs an actor that used less CPU than the golden guest had at snapshot. The test for (2) forces both fold orders deterministically: the sweep closes a probe's connection right after folding its response, so the fixture's closer releases the gated probe from `Close` — no sleeps, no polling. **Timeout rationale reconciled.** The poll floor (50s) is one actor's worst-case micro-VM read: 25 containers at 2s each. It never prevented overlap, because sweeps are sequential; it bounds guest-agent load, and the comment now says so. The RPC timeout (55s) was documented as covering that single-actor worst case. With several actors per worker, the micro-VM ateom caps its own read at `statsSweepBudget` (45s, chosen to stay under atelet's 55s) and reports unreached guests as pending, so the timeout stays a constant; its comment now points at that budget. ## Not in scope, recorded - Two workers both *measuring* one actor in one sweep double-charges. Pre-existing and unchanged here. It is rare: it needs the same probe skew as defect 2, with the destination read after its restore finished. - `ateom.proto` still says the discovery list holds "at most one entry until multi-actor workers land". Stale; fixed in a separate change. ## Verification `gofmt`, `go vet`, and race-enabled tests pass. Mutation checks: dropping the carry, the measured-wins guard, or either side of the per-source decrease rule each fails its test.
952 lines
39 KiB
Go
952 lines
39 KiB
Go
// Copyright 2026 Google LLC
|
|
//
|
|
// Licensed under the Apache License, Version 2.0 (the "License");
|
|
// you may not use this file except in compliance with the License.
|
|
// You may obtain a copy of the License at
|
|
//
|
|
// http://www.apache.org/licenses/LICENSE-2.0
|
|
//
|
|
// Unless required by applicable law or agreed to in writing, software
|
|
// distributed under the License is distributed on an "AS IS" BASIS,
|
|
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
|
|
// See the License for the specific language governing permissions and
|
|
// limitations under the License.
|
|
|
|
package main
|
|
|
|
import (
|
|
"bytes"
|
|
"context"
|
|
"encoding/json"
|
|
"errors"
|
|
"io"
|
|
"math"
|
|
"os"
|
|
"path/filepath"
|
|
"sort"
|
|
"sync"
|
|
"testing"
|
|
"time"
|
|
|
|
"github.com/google/go-cmp/cmp"
|
|
sdkmetric "go.opentelemetry.io/otel/sdk/metric"
|
|
"go.opentelemetry.io/otel/sdk/metric/metricdata"
|
|
"google.golang.org/grpc"
|
|
|
|
corev1 "k8s.io/api/core/v1"
|
|
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
|
|
k8sfake "k8s.io/client-go/kubernetes/fake"
|
|
|
|
"github.com/agent-substrate/substrate/internal/proto/ateompb"
|
|
)
|
|
|
|
// fakeStatsAteom answers GetActiveWorkloadStats with a canned response or
|
|
// error, standing in for one ateom socket.
|
|
type fakeStatsAteom struct {
|
|
resp *ateompb.GetActiveWorkloadStatsResponse
|
|
err error
|
|
|
|
// mu guards the recordings below: the sweep probes ateoms concurrently.
|
|
mu sync.Mutex
|
|
// calls counts probes, so tests can tell "skipped" from "never found".
|
|
calls int
|
|
// sawDeadline records whether the probe's context carried one, pinning the
|
|
// per-call timeout.
|
|
sawDeadline bool
|
|
// gate, when set, is received from before answering, so a test can order
|
|
// this probe after another worker's response has been folded.
|
|
gate <-chan struct{}
|
|
}
|
|
|
|
func (f *fakeStatsAteom) GetActiveWorkloadStats(ctx context.Context, req *ateompb.GetActiveWorkloadStatsRequest, opts ...grpc.CallOption) (*ateompb.GetActiveWorkloadStatsResponse, error) {
|
|
if f.gate != nil {
|
|
select {
|
|
case <-f.gate:
|
|
case <-ctx.Done():
|
|
return nil, ctx.Err() // a misconfigured gate fails the probe, not the run
|
|
}
|
|
}
|
|
f.mu.Lock()
|
|
f.calls++
|
|
_, f.sawDeadline = ctx.Deadline()
|
|
f.mu.Unlock()
|
|
return f.resp, f.err
|
|
}
|
|
|
|
// executingResponse builds the sample an executing ateom would echo.
|
|
func executingResponse(templateNS, templateName string, class ateompb.SandboxClass, source ateompb.StatsSource, current, workingSet uint64) *ateompb.GetActiveWorkloadStatsResponse {
|
|
return &ateompb.GetActiveWorkloadStatsResponse{
|
|
Samples: []*ateompb.WorkloadStatsSample{{
|
|
ActorTemplateAtespace: templateNS,
|
|
ActorTemplateName: templateName,
|
|
SandboxClass: class,
|
|
Source: source,
|
|
MemoryCurrentBytes: current,
|
|
MemoryWorkingSetBytes: workingSet,
|
|
}},
|
|
}
|
|
}
|
|
|
|
// measuredSample is one measured entry of a multi-actor response.
|
|
func measuredSample(actorUID, templateNS, templateName string, current, workingSet, cpuUsec uint64) *ateompb.WorkloadStatsSample {
|
|
return &ateompb.WorkloadStatsSample{
|
|
ActorUid: actorUID,
|
|
ActorTemplateAtespace: templateNS,
|
|
ActorTemplateName: templateName,
|
|
SandboxClass: ateompb.SandboxClass_SANDBOX_CLASS_GVISOR,
|
|
Source: ateompb.StatsSource_STATS_SOURCE_CGROUP,
|
|
MemoryCurrentBytes: current,
|
|
MemoryWorkingSetBytes: workingSet,
|
|
CpuUsageUsec: cpuUsec,
|
|
}
|
|
}
|
|
|
|
// pendingSample is one pending entry: attribution present, source
|
|
// UNSPECIFIED, measurements absent.
|
|
func pendingSample(actorUID, templateNS, templateName string) *ateompb.WorkloadStatsSample {
|
|
return &ateompb.WorkloadStatsSample{
|
|
ActorUid: actorUID,
|
|
ActorTemplateAtespace: templateNS,
|
|
ActorTemplateName: templateName,
|
|
SandboxClass: ateompb.SandboxClass_SANDBOX_CLASS_GVISOR,
|
|
}
|
|
}
|
|
|
|
// availableResponse is an idle ateom's answer: the empty list.
|
|
func availableResponse() *ateompb.GetActiveWorkloadStatsResponse {
|
|
return &ateompb.GetActiveWorkloadStatsResponse{}
|
|
}
|
|
|
|
// pendingResponse is a workload with no numbers yet, per the response
|
|
// contract: attribution present, source UNSPECIFIED, measurements absent.
|
|
func pendingResponse(actorUID string) *ateompb.GetActiveWorkloadStatsResponse {
|
|
return &ateompb.GetActiveWorkloadStatsResponse{
|
|
Samples: []*ateompb.WorkloadStatsSample{{
|
|
ActorUid: actorUID,
|
|
ActorTemplateAtespace: "ns-a",
|
|
ActorTemplateName: "tmpl-a",
|
|
SandboxClass: ateompb.SandboxClass_SANDBOX_CLASS_GVISOR,
|
|
}},
|
|
}
|
|
}
|
|
|
|
// closeRecorder counts Close calls, standing in for a probe's connection.
|
|
// collect closes it after folding the response, so onClose is the point at
|
|
// which that worker's entries are in the aggregates.
|
|
type closeRecorder struct {
|
|
mu sync.Mutex
|
|
closes int
|
|
onClose func()
|
|
}
|
|
|
|
func (c *closeRecorder) Close() error {
|
|
c.mu.Lock()
|
|
c.closes++
|
|
c.mu.Unlock()
|
|
if c.onClose != nil {
|
|
c.onClose()
|
|
}
|
|
return nil
|
|
}
|
|
|
|
// newPollerFixture builds a poller over a fixture ateoms directory with one
|
|
// subdirectory (and one fake) per entry in fakes. Dialing a UID without a fake
|
|
// fails, which is the shape of a stale directory whose socket is gone. Every
|
|
// successful dial hands out a recorded closer; assertClosed checks the
|
|
// connections-live-exactly-one-probe contract.
|
|
func newPollerFixture(t *testing.T, fakes map[string]*fakeStatsAteom) (*statsPoller, map[string]*closeRecorder) {
|
|
t.Helper()
|
|
dir := t.TempDir()
|
|
closers := make(map[string]*closeRecorder)
|
|
for uid := range fakes {
|
|
if err := os.Mkdir(filepath.Join(dir, uid), 0o700); err != nil {
|
|
t.Fatalf("creating fixture ateom dir %q: %v", uid, err)
|
|
}
|
|
closers[uid] = &closeRecorder{}
|
|
}
|
|
return &statsPoller{
|
|
ateomsDir: dir,
|
|
dial: func(_ context.Context, podUID string) (activeStatsClient, io.Closer, error) {
|
|
f, ok := fakes[podUID]
|
|
if !ok || f == nil {
|
|
return nil, nil, errors.New("no such socket")
|
|
}
|
|
return f, closers[podUID], nil
|
|
},
|
|
}, closers
|
|
}
|
|
|
|
// assertClosed checks that every successfully dialed probe closed its
|
|
// connection exactly once per sweep -- the RPC failing must not leak it.
|
|
func assertClosed(t *testing.T, fakes map[string]*fakeStatsAteom, closers map[string]*closeRecorder, sweeps int) {
|
|
t.Helper()
|
|
for uid, f := range fakes {
|
|
if f == nil {
|
|
continue // dial fails: no connection to close
|
|
}
|
|
if got := closers[uid].closes; got != sweeps {
|
|
t.Errorf("ateom %s connection closed %d times over %d sweeps, want %d", uid, got, sweeps, sweeps)
|
|
}
|
|
}
|
|
}
|
|
|
|
func TestStatsPollerCollectAggregates(t *testing.T) {
|
|
// Two actors of the same template on this node, one of another, one idle
|
|
// worker, one mid-boot: the same-template pair sums, the others contribute
|
|
// nothing.
|
|
fakes := map[string]*fakeStatsAteom{
|
|
"uid-1": {resp: executingResponse("ns-a", "tmpl-a", ateompb.SandboxClass_SANDBOX_CLASS_GVISOR, ateompb.StatsSource_STATS_SOURCE_CGROUP, 1000, 700)},
|
|
"uid-2": {resp: executingResponse("ns-a", "tmpl-a", ateompb.SandboxClass_SANDBOX_CLASS_GVISOR, ateompb.StatsSource_STATS_SOURCE_CGROUP, 500, 300)},
|
|
"uid-3": {resp: executingResponse("ns-b", "tmpl-b", ateompb.SandboxClass_SANDBOX_CLASS_MICROVM, ateompb.StatsSource_STATS_SOURCE_GUEST_AGENT, 42, 40)},
|
|
"uid-4": {resp: availableResponse()},
|
|
"uid-5": {resp: pendingResponse("uid-5-actor")},
|
|
}
|
|
p, closers := newPollerFixture(t, fakes)
|
|
|
|
got := p.collect(context.Background())
|
|
|
|
want := map[templateKey]*templateAggregate{
|
|
{templateNamespace: "ns-a", templateName: "tmpl-a", sandboxClass: "gvisor", source: "cgroup"}: {
|
|
sampledActors: 2, memoryCurrentBytes: 1500, memoryWorkingSetBytes: 1000,
|
|
},
|
|
{templateNamespace: "ns-b", templateName: "tmpl-b", sandboxClass: "microvm", source: "guest-agent"}: {
|
|
sampledActors: 1, memoryCurrentBytes: 42, memoryWorkingSetBytes: 40,
|
|
},
|
|
}
|
|
if diff := cmp.Diff(want, got, cmp.AllowUnexported(templateAggregate{}, templateKey{}, workerPoolRef{})); diff != "" {
|
|
t.Errorf("collect() mismatch (-want +got):\n%s", diff)
|
|
}
|
|
|
|
for uid, f := range fakes {
|
|
if f.calls != 1 {
|
|
t.Errorf("ateom %s probed %d times, want 1", uid, f.calls)
|
|
}
|
|
if !f.sawDeadline {
|
|
t.Errorf("ateom %s probed without a deadline; every probe must carry the per-call timeout", uid)
|
|
}
|
|
}
|
|
assertClosed(t, fakes, closers, 1)
|
|
}
|
|
|
|
// TestStatsPollerCollectSkipsFailures pins the scan's one tolerance rule: a
|
|
// dial or call failure means "not a target this tick", never a failed sweep.
|
|
// The healthy ateom's sample must still be aggregated.
|
|
func TestStatsPollerCollectSkipsFailures(t *testing.T) {
|
|
fakes := map[string]*fakeStatsAteom{
|
|
"uid-healthy": {resp: executingResponse("ns-a", "tmpl-a", ateompb.SandboxClass_SANDBOX_CLASS_GVISOR, ateompb.StatsSource_STATS_SOURCE_CGROUP, 100, 80)},
|
|
"uid-stale": nil, // directory with no reachable socket: dial fails
|
|
"uid-broken": {err: errors.New("rpc error: connection refused")},
|
|
}
|
|
p, closers := newPollerFixture(t, fakes)
|
|
|
|
got := p.collect(context.Background())
|
|
|
|
want := map[templateKey]*templateAggregate{
|
|
{templateNamespace: "ns-a", templateName: "tmpl-a", sandboxClass: "gvisor", source: "cgroup"}: {
|
|
sampledActors: 1, memoryCurrentBytes: 100, memoryWorkingSetBytes: 80,
|
|
},
|
|
}
|
|
if diff := cmp.Diff(want, got, cmp.AllowUnexported(templateAggregate{}, templateKey{}, workerPoolRef{})); diff != "" {
|
|
t.Errorf("collect() mismatch (-want +got):\n%s", diff)
|
|
}
|
|
// The broken ateom's RPC failed, but its connection was dialed -- it must
|
|
// be closed all the same.
|
|
assertClosed(t, fakes, closers, 1)
|
|
}
|
|
|
|
// TestStatsPollerCollectNoAteomsDir: a node whose first workload has not
|
|
// arrived has no ateoms directory, which is empty coverage, not an error.
|
|
func TestStatsPollerCollectNoAteomsDir(t *testing.T) {
|
|
p := &statsPoller{ateomsDir: filepath.Join(t.TempDir(), "does-not-exist")}
|
|
if got := p.collect(context.Background()); len(got) != 0 {
|
|
t.Errorf("collect() with no ateoms dir = %v, want empty", got)
|
|
}
|
|
}
|
|
|
|
func TestClampActorStatsPollInterval(t *testing.T) {
|
|
tests := []struct {
|
|
name string
|
|
in time.Duration
|
|
want time.Duration
|
|
}{
|
|
{name: "zero stays disabled", in: 0, want: 0},
|
|
{name: "below floor clamps", in: time.Second, want: minActorStatsPollInterval},
|
|
{name: "at floor passes", in: minActorStatsPollInterval, want: minActorStatsPollInterval},
|
|
{name: "above floor passes", in: 5 * time.Minute, want: 5 * time.Minute},
|
|
}
|
|
for _, tc := range tests {
|
|
t.Run(tc.name, func(t *testing.T) {
|
|
if got := clampActorStatsPollInterval(context.Background(), tc.in); got != tc.want {
|
|
t.Errorf("clampActorStatsPollInterval(%v) = %v, want %v", tc.in, got, tc.want)
|
|
}
|
|
})
|
|
}
|
|
}
|
|
|
|
// TestStatsInstrumentsObserveLatestSnapshotOnly pins the reason the gauges are
|
|
// observable rather than synchronous: each collection reports exactly the
|
|
// groups the latest sweep found. A synchronous gauge would re-export its last
|
|
// recorded value on every collection until process exit, so a template whose
|
|
// actors left the node would keep reporting their memory forever.
|
|
func TestStatsInstrumentsObserveLatestSnapshotOnly(t *testing.T) {
|
|
reader := sdkmetric.NewManualReader()
|
|
mp := sdkmetric.NewMeterProvider(sdkmetric.WithReader(reader))
|
|
defer mp.Shutdown(context.Background())
|
|
|
|
inst, err := newStatsInstruments(mp.Meter("test"))
|
|
if err != nil {
|
|
t.Fatalf("newStatsInstruments() error = %v", err)
|
|
}
|
|
|
|
key := templateKey{templateNamespace: "ns-a", templateName: "tmpl-a", sandboxClass: "gvisor", source: "cgroup"}
|
|
inst.publish(map[templateKey]*templateAggregate{
|
|
key: {sampledActors: 2, memoryCurrentBytes: 1500, memoryWorkingSetBytes: 1000},
|
|
})
|
|
|
|
if got := gaugePointCount(t, reader, workingSetMetric); got != 1 {
|
|
t.Fatalf("after publish: %s has %d datapoints, want 1", workingSetMetric, got)
|
|
}
|
|
|
|
// The template's actors leave the node: an empty sweep must make the
|
|
// series disappear, not freeze at its last value.
|
|
inst.publish(map[templateKey]*templateAggregate{})
|
|
if got := gaugePointCount(t, reader, workingSetMetric); got != 0 {
|
|
t.Errorf("after empty sweep: %s has %d datapoints, want 0", workingSetMetric, got)
|
|
}
|
|
}
|
|
|
|
// gaugePointCount collects once and returns how many datapoints name has.
|
|
func gaugePointCount(t *testing.T, reader *sdkmetric.ManualReader, name string) int {
|
|
t.Helper()
|
|
var rm metricdata.ResourceMetrics
|
|
if err := reader.Collect(context.Background(), &rm); err != nil {
|
|
t.Fatalf("Collect() error = %v", err)
|
|
}
|
|
for _, sm := range rm.ScopeMetrics {
|
|
for _, m := range sm.Metrics {
|
|
if m.Name != name {
|
|
continue
|
|
}
|
|
g, ok := m.Data.(metricdata.Gauge[int64])
|
|
if !ok {
|
|
t.Fatalf("metric %s has data type %T, want Gauge[int64]", name, m.Data)
|
|
}
|
|
return len(g.DataPoints)
|
|
}
|
|
}
|
|
return 0
|
|
}
|
|
|
|
// cpuResponse is executingResponse with only the CPU counter set, for the
|
|
// delta tests.
|
|
func cpuResponse(actorUID string, cpuUsec uint64) *ateompb.GetActiveWorkloadStatsResponse {
|
|
return &ateompb.GetActiveWorkloadStatsResponse{
|
|
Samples: []*ateompb.WorkloadStatsSample{{
|
|
ActorUid: actorUID,
|
|
ActorTemplateAtespace: "ns-a",
|
|
ActorTemplateName: "tmpl-a",
|
|
SandboxClass: ateompb.SandboxClass_SANDBOX_CLASS_GVISOR,
|
|
Source: ateompb.StatsSource_STATS_SOURCE_CGROUP,
|
|
CpuUsageUsec: cpuUsec,
|
|
}},
|
|
}
|
|
}
|
|
|
|
// TestStatsPollerCPUDeltas pins the increase computation across sweeps: the
|
|
// first sight of an actor establishes a baseline and charges nothing (atelet
|
|
// cannot tell a new actor from its own restart, and re-charging an epoch the
|
|
// previous atelet counted would spike the counter), a later sweep charges
|
|
// only the increase, a cgroup decrease is an epoch reset whose new value is
|
|
// the usage since the reset, and an actor that disappears stops contributing and
|
|
// is dropped from the baselines.
|
|
func TestStatsPollerCPUDeltas(t *testing.T) {
|
|
key := templateKey{templateNamespace: "ns-a", templateName: "tmpl-a", sandboxClass: "gvisor", source: "cgroup"}
|
|
fake := &fakeStatsAteom{resp: cpuResponse("uid-a", 1000)}
|
|
p, _ := newPollerFixture(t, map[string]*fakeStatsAteom{"uid-1": fake})
|
|
|
|
if got := p.collect(context.Background())[key].cpuDeltaUsec; got != 0 {
|
|
t.Errorf("first sweep delta = %d, want 0 (baseline only on first sight)", got)
|
|
}
|
|
|
|
fake.resp = cpuResponse("uid-a", 1600)
|
|
if got := p.collect(context.Background())[key].cpuDeltaUsec; got != 600 {
|
|
t.Errorf("second sweep delta = %d, want 600 (the increase)", got)
|
|
}
|
|
|
|
// Epoch reset: the counter went backwards, so the new value is the usage
|
|
// since the reset.
|
|
fake.resp = cpuResponse("uid-a", 250)
|
|
if got := p.collect(context.Background())[key].cpuDeltaUsec; got != 250 {
|
|
t.Errorf("post-reset sweep delta = %d, want 250", got)
|
|
}
|
|
|
|
// The actor leaves: nothing to contribute, and its baseline must be
|
|
// dropped so a later return re-baselines instead of comparing against a
|
|
// dead value.
|
|
fake.resp = availableResponse()
|
|
if got := p.collect(context.Background()); len(got) != 0 {
|
|
t.Errorf("empty sweep aggregates = %v, want none", got)
|
|
}
|
|
if len(p.lastCPU) != 0 {
|
|
t.Errorf("baselines after empty sweep = %v, want pruned empty", p.lastCPU)
|
|
}
|
|
}
|
|
|
|
// TestStatsPollerWorkerPoolLabels pins the pool enrichment: a resolved pod
|
|
// groups under its pool, an unresolved one groups without pool labels rather
|
|
// than vanishing, and the two never merge.
|
|
func TestStatsPollerWorkerPoolLabels(t *testing.T) {
|
|
fakes := map[string]*fakeStatsAteom{
|
|
"uid-pooled": {resp: executingResponse("ns-a", "tmpl-a", ateompb.SandboxClass_SANDBOX_CLASS_GVISOR, ateompb.StatsSource_STATS_SOURCE_CGROUP, 100, 80)},
|
|
"uid-unpooled": {resp: executingResponse("ns-a", "tmpl-a", ateompb.SandboxClass_SANDBOX_CLASS_GVISOR, ateompb.StatsSource_STATS_SOURCE_CGROUP, 10, 8)},
|
|
}
|
|
p, _ := newPollerFixture(t, fakes)
|
|
p.fetchWorkerPools = func(context.Context) map[string]workerPoolRef {
|
|
return map[string]workerPoolRef{"uid-pooled": {namespace: "pool-ns", name: "pool-a"}}
|
|
}
|
|
|
|
got := p.collect(context.Background())
|
|
|
|
base := templateKey{templateNamespace: "ns-a", templateName: "tmpl-a", sandboxClass: "gvisor", source: "cgroup"}
|
|
pooled := base
|
|
pooled.workerPool = workerPoolRef{namespace: "pool-ns", name: "pool-a"}
|
|
want := map[templateKey]*templateAggregate{
|
|
pooled: {sampledActors: 1, memoryCurrentBytes: 100, memoryWorkingSetBytes: 80},
|
|
base: {sampledActors: 1, memoryCurrentBytes: 10, memoryWorkingSetBytes: 8},
|
|
}
|
|
if diff := cmp.Diff(want, got, cmp.AllowUnexported(templateAggregate{}, templateKey{}, workerPoolRef{})); diff != "" {
|
|
t.Errorf("collect() mismatch (-want +got):\n%s", diff)
|
|
}
|
|
}
|
|
|
|
// TestStatsPollerPeriodicEvents pins the events channel: one event per
|
|
// executing sample per sweep, none for idle or mid-boot ateoms, identity
|
|
// taken from the echo, pool labels from the sweep's own resolution.
|
|
func TestStatsPollerPeriodicEvents(t *testing.T) {
|
|
fakes := map[string]*fakeStatsAteom{
|
|
"uid-1": {resp: executingResponse("ns-a", "tmpl-a", ateompb.SandboxClass_SANDBOX_CLASS_GVISOR, ateompb.StatsSource_STATS_SOURCE_CGROUP, 1000, 700)},
|
|
"uid-2": {resp: availableResponse()},
|
|
}
|
|
p, _ := newPollerFixture(t, fakes)
|
|
var buf syncBuffer
|
|
p.eventEmitter = newBufferEmitter(&buf, false)
|
|
p.fetchWorkerPools = func(context.Context) map[string]workerPoolRef {
|
|
return map[string]workerPoolRef{"uid-1": {namespace: "pool-ns", name: "pool-a"}}
|
|
}
|
|
|
|
p.collect(context.Background())
|
|
|
|
lines := bytes.Count(bytes.TrimSpace(buf.Bytes()), []byte("\n")) + 1
|
|
if buf.Len() == 0 {
|
|
t.Fatal("no periodic event emitted for the executing ateom")
|
|
}
|
|
if lines != 1 {
|
|
t.Fatalf("emitted %d events, want 1 (idle ateoms emit nothing): %q", lines, buf.String())
|
|
}
|
|
var rec map[string]any
|
|
if err := json.Unmarshal(bytes.TrimSpace(buf.Bytes()), &rec); err != nil {
|
|
t.Fatalf("unmarshal: %v", err)
|
|
}
|
|
if got := rec["kind"]; got != "periodic" {
|
|
t.Errorf("kind = %v, want periodic", got)
|
|
}
|
|
labels, _ := rec["labels"].(map[string]any)
|
|
if got := labels["ate.workerpool.name"]; got != "pool-a" {
|
|
t.Errorf("labels[ate.workerpool.name] = %v, want pool-a", got)
|
|
}
|
|
}
|
|
|
|
func TestAddSat(t *testing.T) {
|
|
tests := []struct {
|
|
name string
|
|
agg int64
|
|
v uint64
|
|
want int64
|
|
}{
|
|
{name: "normal add", agg: 100, v: 50, want: 150},
|
|
{name: "zero add", agg: 100, v: 0, want: 100},
|
|
// A wire value above MaxInt64 -- a corrupt or hostile guest reading --
|
|
// must pin at the ceiling, not wrap the aggregate negative.
|
|
{name: "value above MaxInt64 saturates", agg: 0, v: math.MaxUint64, want: math.MaxInt64},
|
|
// The addition itself can also overflow once inputs are clamped.
|
|
{name: "sum overflow saturates", agg: math.MaxInt64 - 10, v: 100, want: math.MaxInt64},
|
|
{name: "exactly at ceiling", agg: math.MaxInt64 - 5, v: 5, want: math.MaxInt64},
|
|
}
|
|
for _, tc := range tests {
|
|
t.Run(tc.name, func(t *testing.T) {
|
|
if got := addSat(tc.agg, tc.v); got != tc.want {
|
|
t.Errorf("addSat(%d, %d) = %d, want %d", tc.agg, tc.v, tc.want, got)
|
|
}
|
|
})
|
|
}
|
|
}
|
|
|
|
// TestStatsPollerCollectSaturatesCorruptSamples pins the end-to-end behavior:
|
|
// one guest reporting absurd counters must not flip a template's aggregates
|
|
// negative -- a negative gauge misreads as "no memory", and a negative CPU
|
|
// delta is a spec-violating counter Add. Everything pins at MaxInt64 instead.
|
|
func TestStatsPollerCollectSaturatesCorruptSamples(t *testing.T) {
|
|
key := templateKey{templateNamespace: "ns-a", templateName: "tmpl-a", sandboxClass: "gvisor", source: "cgroup"}
|
|
fakes := map[string]*fakeStatsAteom{
|
|
"uid-1": {resp: executingResponse("ns-a", "tmpl-a", ateompb.SandboxClass_SANDBOX_CLASS_GVISOR, ateompb.StatsSource_STATS_SOURCE_CGROUP, math.MaxUint64, math.MaxUint64)},
|
|
"uid-2": {resp: executingResponse("ns-a", "tmpl-a", ateompb.SandboxClass_SANDBOX_CLASS_GVISOR, ateompb.StatsSource_STATS_SOURCE_CGROUP, 1000, 700)},
|
|
}
|
|
p, _ := newPollerFixture(t, fakes)
|
|
|
|
got := p.collect(context.Background())[key]
|
|
if got == nil {
|
|
t.Fatal("collect() returned no aggregate for the template")
|
|
}
|
|
if got.memoryCurrentBytes != math.MaxInt64 || got.memoryWorkingSetBytes != math.MaxInt64 {
|
|
t.Errorf("memory aggregates = %d/%d, want both pinned at MaxInt64",
|
|
got.memoryCurrentBytes, got.memoryWorkingSetBytes)
|
|
}
|
|
if got.memoryCurrentBytes < 0 || got.memoryWorkingSetBytes < 0 || got.cpuDeltaUsec < 0 {
|
|
t.Errorf("aggregate went negative: %+v", got)
|
|
}
|
|
}
|
|
|
|
// TestStatsPollerCPUDeltaSaturatesCorruptCounter: a baseline followed by an
|
|
// absurd counter value is a huge "increase"; it must clamp, not go negative.
|
|
func TestStatsPollerCPUDeltaSaturatesCorruptCounter(t *testing.T) {
|
|
key := templateKey{templateNamespace: "ns-a", templateName: "tmpl-a", sandboxClass: "gvisor", source: "cgroup"}
|
|
fake := &fakeStatsAteom{resp: cpuResponse("uid-a", 1000)}
|
|
p, _ := newPollerFixture(t, map[string]*fakeStatsAteom{"uid-1": fake})
|
|
|
|
if got := p.collect(context.Background())[key].cpuDeltaUsec; got != 0 {
|
|
t.Fatalf("first sweep delta = %d, want 0", got)
|
|
}
|
|
|
|
fake.resp = cpuResponse("uid-a", math.MaxUint64)
|
|
got := p.collect(context.Background())[key].cpuDeltaUsec
|
|
if got != math.MaxInt64 {
|
|
t.Errorf("corrupt-counter sweep delta = %d, want pinned at MaxInt64", got)
|
|
}
|
|
}
|
|
|
|
// TestNewWorkerPoolFetcher pins the fetcher's ingestion rules: a labeled worker
|
|
// maps by pod UID, an empty label value names no pool and never enters the
|
|
// map (the presence-only selector matches it anyway), and unlabeled pods are
|
|
// not workers at all. The fake clientset honors label selectors but not the
|
|
// spec.nodeName field selector, so node scoping is not assertable here.
|
|
func TestNewWorkerPoolFetcher(t *testing.T) {
|
|
client := k8sfake.NewSimpleClientset(
|
|
&corev1.Pod{ObjectMeta: metav1.ObjectMeta{
|
|
Name: "worker-a", Namespace: "pool-ns", UID: "uid-a",
|
|
Labels: map[string]string{workerPoolLabel: "pool-a"},
|
|
}},
|
|
&corev1.Pod{ObjectMeta: metav1.ObjectMeta{
|
|
Name: "worker-empty", Namespace: "pool-ns", UID: "uid-empty",
|
|
Labels: map[string]string{workerPoolLabel: ""},
|
|
}},
|
|
&corev1.Pod{ObjectMeta: metav1.ObjectMeta{
|
|
Name: "bystander", Namespace: "other-ns", UID: "uid-bystander",
|
|
}},
|
|
)
|
|
|
|
got := newWorkerPoolFetcher(client, "node-1")(context.Background())
|
|
|
|
want := map[string]workerPoolRef{"uid-a": {namespace: "pool-ns", name: "pool-a"}}
|
|
if diff := cmp.Diff(want, got, cmp.AllowUnexported(workerPoolRef{})); diff != "" {
|
|
t.Errorf("newWorkerPoolFetcher mismatch (-want +got):\n%s", diff)
|
|
}
|
|
}
|
|
|
|
// TestStatsPollerPoolCacheSurvivesListFlap pins the fix for the label-set
|
|
// split: a failed list must answer from the cache and keep that tick's
|
|
// samples -- including the CPU increase computed during the flap, the value
|
|
// that feeds the monotonic counter and can never be re-attributed -- on the
|
|
// pooled label set.
|
|
func TestStatsPollerPoolCacheSurvivesListFlap(t *testing.T) {
|
|
resp := executingResponse("ns-a", "tmpl-a", ateompb.SandboxClass_SANDBOX_CLASS_GVISOR, ateompb.StatsSource_STATS_SOURCE_CGROUP, 100, 80)
|
|
resp.GetSamples()[0].ActorUid = "uid-a"
|
|
resp.GetSamples()[0].CpuUsageUsec = 1000
|
|
fake := &fakeStatsAteom{resp: resp}
|
|
p, _ := newPollerFixture(t, map[string]*fakeStatsAteom{"uid-1": fake})
|
|
listOK := true
|
|
p.fetchWorkerPools = func(context.Context) map[string]workerPoolRef {
|
|
if !listOK {
|
|
return nil // the apiserver list failed this sweep
|
|
}
|
|
return map[string]workerPoolRef{"uid-1": {namespace: "pool-ns", name: "pool-a"}}
|
|
}
|
|
|
|
pooled := templateKey{templateNamespace: "ns-a", templateName: "tmpl-a", sandboxClass: "gvisor", source: "cgroup",
|
|
workerPool: workerPoolRef{namespace: "pool-ns", name: "pool-a"}}
|
|
|
|
// Sweep 1 resolves, seeds the pool cache, and baselines the CPU counter.
|
|
if got := p.collect(context.Background()); got[pooled] == nil {
|
|
t.Fatalf("sweep 1: no pooled aggregate; got %v", got)
|
|
}
|
|
|
|
// Sweep 2: the list fails AND the actor consumed CPU. Both the sample and
|
|
// its delta must still group under the pool.
|
|
listOK = false
|
|
resp2 := executingResponse("ns-a", "tmpl-a", ateompb.SandboxClass_SANDBOX_CLASS_GVISOR, ateompb.StatsSource_STATS_SOURCE_CGROUP, 100, 80)
|
|
resp2.GetSamples()[0].ActorUid = "uid-a"
|
|
resp2.GetSamples()[0].CpuUsageUsec = 1600
|
|
fake.resp = resp2
|
|
got := p.collect(context.Background())
|
|
if got[pooled] == nil {
|
|
t.Fatalf("sweep 2 (list flap): samples left the pooled label set; got %v", got)
|
|
}
|
|
if len(got) != 1 {
|
|
t.Errorf("sweep 2 (list flap): %d label sets, want 1 (no pool-less split)", len(got))
|
|
}
|
|
if got[pooled].cpuDeltaUsec != 600 {
|
|
t.Errorf("sweep 2 (list flap): pooled cpu delta = %d, want 600 -- the counter increment must land on the pooled series", got[pooled].cpuDeltaUsec)
|
|
}
|
|
}
|
|
|
|
// TestStatsPollerPoolCachePrunes: the cache is rebuilt against the pods whose
|
|
// ateom directories exist, so a departed pod's entry does not linger.
|
|
func TestStatsPollerPoolCachePrunes(t *testing.T) {
|
|
fakes := map[string]*fakeStatsAteom{
|
|
"uid-1": {resp: availableResponse()},
|
|
}
|
|
p, _ := newPollerFixture(t, fakes)
|
|
p.fetchWorkerPools = func(context.Context) map[string]workerPoolRef {
|
|
return map[string]workerPoolRef{
|
|
"uid-1": {namespace: "pool-ns", name: "pool-a"},
|
|
"uid-gone": {namespace: "pool-ns", name: "pool-a"}, // no ateom dir
|
|
}
|
|
}
|
|
|
|
p.collect(context.Background())
|
|
|
|
if _, ok := p.cachedPools["uid-1"]; !ok {
|
|
t.Errorf("cachedPools lost the live pod's entry: %v", p.cachedPools)
|
|
}
|
|
if _, ok := p.cachedPools["uid-gone"]; ok {
|
|
t.Errorf("cachedPools kept an entry with no ateom directory: %v", p.cachedPools)
|
|
}
|
|
}
|
|
|
|
// TestStatsPollerPoolCacheMissDuringOutage: a pod first seen while the list
|
|
// is failing has no cache entry to fall back to -- it groups without pool
|
|
// labels (the residual, documented case) and heals on the next good list.
|
|
func TestStatsPollerPoolCacheMissDuringOutage(t *testing.T) {
|
|
fakes := map[string]*fakeStatsAteom{
|
|
"uid-new": {resp: executingResponse("ns-a", "tmpl-a", ateompb.SandboxClass_SANDBOX_CLASS_GVISOR, ateompb.StatsSource_STATS_SOURCE_CGROUP, 10, 8)},
|
|
}
|
|
p, _ := newPollerFixture(t, fakes)
|
|
p.fetchWorkerPools = func(context.Context) map[string]workerPoolRef { return nil }
|
|
|
|
got := p.collect(context.Background())
|
|
|
|
bare := templateKey{templateNamespace: "ns-a", templateName: "tmpl-a", sandboxClass: "gvisor", source: "cgroup"}
|
|
if got[bare] == nil || len(got) != 1 {
|
|
t.Errorf("collect() during outage = %v, want the one pool-less group", got)
|
|
}
|
|
}
|
|
|
|
// TestStatsPollerPoolCachePartialListFallsBack pins the per-pod half of the
|
|
// promised fallback ("a failed OR PARTIAL list"): a fetch that succeeds but
|
|
// omits a cached pod must not strand that pod's samples. Distinct from the
|
|
// all-nil flap test above -- a whole-map fallback would pass that test and
|
|
// fail this one.
|
|
func TestStatsPollerPoolCachePartialListFallsBack(t *testing.T) {
|
|
resp := executingResponse("ns-a", "tmpl-a", ateompb.SandboxClass_SANDBOX_CLASS_GVISOR, ateompb.StatsSource_STATS_SOURCE_CGROUP, 100, 80)
|
|
resp.GetSamples()[0].ActorUid = "uid-a"
|
|
resp.GetSamples()[0].CpuUsageUsec = 1000
|
|
fake := &fakeStatsAteom{resp: resp}
|
|
p, _ := newPollerFixture(t, map[string]*fakeStatsAteom{"uid-1": fake})
|
|
full := true
|
|
p.fetchWorkerPools = func(context.Context) map[string]workerPoolRef {
|
|
if !full {
|
|
// A successful list that no longer contains uid-1 (races between
|
|
// the pod store and the dir scan look exactly like this).
|
|
return map[string]workerPoolRef{"uid-other": {namespace: "pool-ns", name: "pool-b"}}
|
|
}
|
|
return map[string]workerPoolRef{"uid-1": {namespace: "pool-ns", name: "pool-a"}}
|
|
}
|
|
|
|
pooled := templateKey{templateNamespace: "ns-a", templateName: "tmpl-a", sandboxClass: "gvisor", source: "cgroup",
|
|
workerPool: workerPoolRef{namespace: "pool-ns", name: "pool-a"}}
|
|
|
|
if got := p.collect(context.Background()); got[pooled] == nil {
|
|
t.Fatalf("sweep 1: no pooled aggregate; got %v", got)
|
|
}
|
|
|
|
full = false
|
|
resp2 := executingResponse("ns-a", "tmpl-a", ateompb.SandboxClass_SANDBOX_CLASS_GVISOR, ateompb.StatsSource_STATS_SOURCE_CGROUP, 100, 80)
|
|
resp2.GetSamples()[0].ActorUid = "uid-a"
|
|
resp2.GetSamples()[0].CpuUsageUsec = 1250
|
|
fake.resp = resp2
|
|
got := p.collect(context.Background())
|
|
if got[pooled] == nil || len(got) != 1 {
|
|
t.Errorf("sweep 2 (partial list): samples left the pooled label set; got %v", got)
|
|
}
|
|
if got[pooled] != nil && got[pooled].cpuDeltaUsec != 250 {
|
|
t.Errorf("sweep 2 (partial list): pooled cpu delta = %d, want 250", got[pooled].cpuDeltaUsec)
|
|
}
|
|
}
|
|
|
|
// TestStatsPollerFoldsMultiActorWorker pins the fold over a multi-actor
|
|
// worker: same-template entries sum, another template gets its own key, and
|
|
// a pending entry contributes nothing.
|
|
func TestStatsPollerFoldsMultiActorWorker(t *testing.T) {
|
|
fakes := map[string]*fakeStatsAteom{
|
|
"uid-w1": {resp: &ateompb.GetActiveWorkloadStatsResponse{Samples: []*ateompb.WorkloadStatsSample{
|
|
measuredSample("actor-1", "ns-a", "tmpl-a", 1000, 700, 0),
|
|
measuredSample("actor-2", "ns-a", "tmpl-a", 500, 300, 0),
|
|
measuredSample("actor-3", "ns-b", "tmpl-b", 42, 40, 0),
|
|
pendingSample("actor-4", "ns-a", "tmpl-a"),
|
|
}}},
|
|
}
|
|
p, closers := newPollerFixture(t, fakes)
|
|
|
|
got := p.collect(context.Background())
|
|
|
|
want := map[templateKey]*templateAggregate{
|
|
{templateNamespace: "ns-a", templateName: "tmpl-a", sandboxClass: "gvisor", source: "cgroup"}: {
|
|
sampledActors: 2, memoryCurrentBytes: 1500, memoryWorkingSetBytes: 1000,
|
|
},
|
|
{templateNamespace: "ns-b", templateName: "tmpl-b", sandboxClass: "gvisor", source: "cgroup"}: {
|
|
sampledActors: 1, memoryCurrentBytes: 42, memoryWorkingSetBytes: 40,
|
|
},
|
|
}
|
|
if diff := cmp.Diff(want, got, cmp.AllowUnexported(templateAggregate{}, templateKey{}, workerPoolRef{})); diff != "" {
|
|
t.Errorf("collect() mismatch (-want +got):\n%s", diff)
|
|
}
|
|
// One probe serves all four entries.
|
|
if fakes["uid-w1"].calls != 1 {
|
|
t.Errorf("worker probed %d times, want 1", fakes["uid-w1"].calls)
|
|
}
|
|
assertClosed(t, fakes, closers, 1)
|
|
}
|
|
|
|
// TestStatsPollerMultiActorCPUBaselines pins that CPU baselines are per
|
|
// actor, not per worker: two actors on one worker advance independently, and
|
|
// a pending sibling neither gains a baseline nor disturbs the others'.
|
|
func TestStatsPollerMultiActorCPUBaselines(t *testing.T) {
|
|
key := templateKey{templateNamespace: "ns-a", templateName: "tmpl-a", sandboxClass: "gvisor", source: "cgroup"}
|
|
fake := &fakeStatsAteom{resp: &ateompb.GetActiveWorkloadStatsResponse{Samples: []*ateompb.WorkloadStatsSample{
|
|
measuredSample("actor-1", "ns-a", "tmpl-a", 1, 1, 1000),
|
|
measuredSample("actor-2", "ns-a", "tmpl-a", 1, 1, 5000),
|
|
pendingSample("actor-3", "ns-a", "tmpl-a"),
|
|
}}}
|
|
p, _ := newPollerFixture(t, map[string]*fakeStatsAteom{"uid-w1": fake})
|
|
|
|
if got := p.collect(context.Background())[key].cpuDeltaUsec; got != 0 {
|
|
t.Fatalf("first sweep delta = %d, want 0 (baselines only)", got)
|
|
}
|
|
if _, ok := p.lastCPU["actor-3"]; ok {
|
|
t.Error("pending actor was given a CPU baseline")
|
|
}
|
|
|
|
// actor-1 advances by 600, actor-2 by 250; the sum lands on the shared key.
|
|
fake.resp = &ateompb.GetActiveWorkloadStatsResponse{Samples: []*ateompb.WorkloadStatsSample{
|
|
measuredSample("actor-1", "ns-a", "tmpl-a", 1, 1, 1600),
|
|
measuredSample("actor-2", "ns-a", "tmpl-a", 1, 1, 5250),
|
|
pendingSample("actor-3", "ns-a", "tmpl-a"),
|
|
}}
|
|
if got := p.collect(context.Background())[key].cpuDeltaUsec; got != 850 {
|
|
t.Errorf("second sweep delta = %d, want 850 (600 + 250)", got)
|
|
}
|
|
|
|
// actor-2 resets its epoch; actor-1 keeps advancing. The reset charges
|
|
// the new value, not a negative delta.
|
|
fake.resp = &ateompb.GetActiveWorkloadStatsResponse{Samples: []*ateompb.WorkloadStatsSample{
|
|
measuredSample("actor-1", "ns-a", "tmpl-a", 1, 1, 1700),
|
|
measuredSample("actor-2", "ns-a", "tmpl-a", 1, 1, 40),
|
|
}}
|
|
if got := p.collect(context.Background())[key].cpuDeltaUsec; got != 140 {
|
|
t.Errorf("reset sweep delta = %d, want 140 (100 + 40)", got)
|
|
}
|
|
if len(p.lastCPU) != 2 {
|
|
t.Errorf("baselines = %v, want exactly the two measured actors", p.lastCPU)
|
|
}
|
|
}
|
|
|
|
// TestStatsPollerMultiActorEvents pins one usage event per measured entry
|
|
// and none for a pending one.
|
|
func TestStatsPollerMultiActorEvents(t *testing.T) {
|
|
fakes := map[string]*fakeStatsAteom{
|
|
"uid-w1": {resp: &ateompb.GetActiveWorkloadStatsResponse{Samples: []*ateompb.WorkloadStatsSample{
|
|
measuredSample("actor-1", "ns-a", "tmpl-a", 1000, 700, 10),
|
|
measuredSample("actor-2", "ns-a", "tmpl-a", 500, 300, 20),
|
|
pendingSample("actor-3", "ns-a", "tmpl-a"),
|
|
}}},
|
|
}
|
|
p, _ := newPollerFixture(t, fakes)
|
|
var buf syncBuffer
|
|
p.eventEmitter = newBufferEmitter(&buf, false)
|
|
|
|
p.collect(context.Background())
|
|
|
|
var uids []string
|
|
for _, line := range bytes.Split(bytes.TrimSpace(buf.Bytes()), []byte("\n")) {
|
|
var rec map[string]any
|
|
if err := json.Unmarshal(line, &rec); err != nil {
|
|
t.Fatalf("event line is not JSON: %v: %s", err, line)
|
|
}
|
|
labels, _ := rec["labels"].(map[string]any)
|
|
uid, _ := labels["ate.actor.uid"].(string)
|
|
uids = append(uids, uid)
|
|
}
|
|
sort.Strings(uids)
|
|
if want := []string{"actor-1", "actor-2"}; !cmp.Equal(want, uids) {
|
|
t.Errorf("event actor uids = %v, want %v (one per measured entry, none for pending)", uids, want)
|
|
}
|
|
}
|
|
|
|
// TestStatsPollerPendingKeepsCPUBaseline pins that a sweep in which an actor
|
|
// reports pending carries its baseline forward: the interval is charged when
|
|
// the actor is measured again, not lost.
|
|
func TestStatsPollerPendingKeepsCPUBaseline(t *testing.T) {
|
|
key := templateKey{templateNamespace: "ns-a", templateName: "tmpl-a", sandboxClass: "gvisor", source: "cgroup"}
|
|
fake := &fakeStatsAteom{resp: &ateompb.GetActiveWorkloadStatsResponse{Samples: []*ateompb.WorkloadStatsSample{
|
|
measuredSample("actor-1", "ns-a", "tmpl-a", 1, 1, 1000),
|
|
}}}
|
|
p, _ := newPollerFixture(t, map[string]*fakeStatsAteom{"uid-w1": fake})
|
|
|
|
p.collect(context.Background()) // baseline 1000
|
|
|
|
fake.resp = &ateompb.GetActiveWorkloadStatsResponse{Samples: []*ateompb.WorkloadStatsSample{
|
|
pendingSample("actor-1", "ns-a", "tmpl-a"),
|
|
}}
|
|
if got := p.collect(context.Background()); len(got) != 0 {
|
|
t.Fatalf("pending sweep produced aggregates %v, want none", got)
|
|
}
|
|
if got, ok := p.lastCPU["actor-1"]; !ok || got != 1000 {
|
|
t.Fatalf("baseline after pending sweep = %d (present=%v), want 1000 kept", got, ok)
|
|
}
|
|
|
|
fake.resp = &ateompb.GetActiveWorkloadStatsResponse{Samples: []*ateompb.WorkloadStatsSample{
|
|
measuredSample("actor-1", "ns-a", "tmpl-a", 1, 1, 1600),
|
|
}}
|
|
if got := p.collect(context.Background())[key].cpuDeltaUsec; got != 600 {
|
|
t.Errorf("delta after pending sweep = %d, want 600 (measured against the kept baseline)", got)
|
|
}
|
|
}
|
|
|
|
// TestStatsPollerCPUDecreaseBySource pins the decrease rule, with and without
|
|
// a pending sweep between the samples: a cgroup decrease charges the new value,
|
|
// since that counter restarts at zero, and a guest-agent decrease charges
|
|
// nothing, since that counter can resume at another guest's value.
|
|
func TestStatsPollerCPUDecreaseBySource(t *testing.T) {
|
|
guestAgent := func(s *ateompb.WorkloadStatsSample) *ateompb.WorkloadStatsSample {
|
|
s.SandboxClass = ateompb.SandboxClass_SANDBOX_CLASS_MICROVM
|
|
s.Source = ateompb.StatsSource_STATS_SOURCE_GUEST_AGENT
|
|
return s
|
|
}
|
|
for _, tc := range []struct {
|
|
name string
|
|
sample func(cpu uint64) *ateompb.WorkloadStatsSample
|
|
key templateKey
|
|
pending bool
|
|
wantDelta int64
|
|
}{
|
|
{
|
|
name: "cgroup",
|
|
sample: func(cpu uint64) *ateompb.WorkloadStatsSample {
|
|
return measuredSample("actor-1", "ns-a", "tmpl-a", 1, 1, cpu)
|
|
},
|
|
key: templateKey{templateNamespace: "ns-a", templateName: "tmpl-a", sandboxClass: "gvisor", source: "cgroup"},
|
|
wantDelta: 400,
|
|
},
|
|
{
|
|
name: "cgroup across pending",
|
|
sample: func(cpu uint64) *ateompb.WorkloadStatsSample {
|
|
return measuredSample("actor-1", "ns-a", "tmpl-a", 1, 1, cpu)
|
|
},
|
|
key: templateKey{templateNamespace: "ns-a", templateName: "tmpl-a", sandboxClass: "gvisor", source: "cgroup"},
|
|
pending: true,
|
|
wantDelta: 400,
|
|
},
|
|
{
|
|
name: "guest-agent",
|
|
sample: func(cpu uint64) *ateompb.WorkloadStatsSample {
|
|
return guestAgent(measuredSample("actor-1", "ns-a", "tmpl-a", 1, 1, cpu))
|
|
},
|
|
key: templateKey{templateNamespace: "ns-a", templateName: "tmpl-a", sandboxClass: "microvm", source: "guest-agent"},
|
|
wantDelta: 0,
|
|
},
|
|
{
|
|
name: "guest-agent across pending",
|
|
sample: func(cpu uint64) *ateompb.WorkloadStatsSample {
|
|
return guestAgent(measuredSample("actor-1", "ns-a", "tmpl-a", 1, 1, cpu))
|
|
},
|
|
key: templateKey{templateNamespace: "ns-a", templateName: "tmpl-a", sandboxClass: "microvm", source: "guest-agent"},
|
|
pending: true,
|
|
wantDelta: 0,
|
|
},
|
|
} {
|
|
t.Run(tc.name, func(t *testing.T) {
|
|
fake := &fakeStatsAteom{}
|
|
p, _ := newPollerFixture(t, map[string]*fakeStatsAteom{"uid-w1": fake})
|
|
sweep := func(sample *ateompb.WorkloadStatsSample) map[templateKey]*templateAggregate {
|
|
fake.resp = &ateompb.GetActiveWorkloadStatsResponse{Samples: []*ateompb.WorkloadStatsSample{sample}}
|
|
return p.collect(context.Background())
|
|
}
|
|
|
|
sweep(tc.sample(1000)) // baseline 1000
|
|
if tc.pending {
|
|
sweep(pendingSample("actor-1", "ns-a", "tmpl-a"))
|
|
}
|
|
if got := sweep(tc.sample(400))[tc.key].cpuDeltaUsec; got != tc.wantDelta {
|
|
t.Errorf("delta for a decrease 1000 -> 400 = %d, want %d", got, tc.wantDelta)
|
|
}
|
|
// Either way 400 is the new baseline.
|
|
if got := sweep(tc.sample(500))[tc.key].cpuDeltaUsec; got != 100 {
|
|
t.Errorf("delta after the decrease = %d, want 100", got)
|
|
}
|
|
})
|
|
}
|
|
}
|
|
|
|
// TestStatsPollerRestoreInFlightMeasuredWins pins that when one worker
|
|
// measures an actor and another reports it pending in the same sweep, the
|
|
// measured value is the baseline in either fold order.
|
|
func TestStatsPollerRestoreInFlightMeasuredWins(t *testing.T) {
|
|
key := templateKey{templateNamespace: "ns-a", templateName: "tmpl-a", sandboxClass: "gvisor", source: "cgroup"}
|
|
for _, order := range []string{"measured first", "pending first", "unordered"} {
|
|
t.Run(order, func(t *testing.T) {
|
|
src := &fakeStatsAteom{resp: &ateompb.GetActiveWorkloadStatsResponse{Samples: []*ateompb.WorkloadStatsSample{
|
|
measuredSample("actor-1", "ns-a", "tmpl-a", 1, 1, 1000),
|
|
}}}
|
|
dst := &fakeStatsAteom{resp: availableResponse()}
|
|
p, closers := newPollerFixture(t, map[string]*fakeStatsAteom{"uid-src": src, "uid-dst": dst})
|
|
|
|
p.collect(context.Background()) // baseline C0 = 1000
|
|
|
|
// Sweep N: source measured at C1 = 1300; destination pending.
|
|
src.resp = &ateompb.GetActiveWorkloadStatsResponse{Samples: []*ateompb.WorkloadStatsSample{
|
|
measuredSample("actor-1", "ns-a", "tmpl-a", 1, 1, 1300),
|
|
}}
|
|
dst.resp = &ateompb.GetActiveWorkloadStatsResponse{Samples: []*ateompb.WorkloadStatsSample{
|
|
pendingSample("actor-1", "ns-a", "tmpl-a"),
|
|
}}
|
|
// The gated probe answers only after the other worker's response has
|
|
// been folded: collect closes a probe's connection after the fold.
|
|
gate := make(chan struct{})
|
|
release := sync.OnceFunc(func() { close(gate) })
|
|
switch order {
|
|
case "measured first":
|
|
dst.gate = gate
|
|
closers["uid-src"].onClose = release
|
|
case "pending first":
|
|
src.gate = gate
|
|
closers["uid-dst"].onClose = release
|
|
}
|
|
if got := p.collect(context.Background())[key].cpuDeltaUsec; got != 300 {
|
|
t.Fatalf("sweep N delta = %d, want 300", got)
|
|
}
|
|
if got := p.lastCPU["actor-1"]; got != 1300 {
|
|
t.Fatalf("baseline after sweep N = %d, want the measured 1300, not the carried 1000", got)
|
|
}
|
|
|
|
// Sweep N+1: only the destination hosts it now, measured at C2 = 1500.
|
|
closers["uid-src"].onClose, closers["uid-dst"].onClose = nil, nil
|
|
src.resp, src.gate = availableResponse(), nil
|
|
dst.resp, dst.gate = &ateompb.GetActiveWorkloadStatsResponse{Samples: []*ateompb.WorkloadStatsSample{
|
|
measuredSample("actor-1", "ns-a", "tmpl-a", 1, 1, 1500),
|
|
}}, nil
|
|
if got := p.collect(context.Background())[key].cpuDeltaUsec; got != 200 {
|
|
t.Errorf("sweep N+1 delta = %d, want 200 (C2-C1); 500 would double-charge C1-C0", got)
|
|
}
|
|
})
|
|
}
|
|
}
|