mirror of
https://github.com/agent-substrate/substrate.git
synced 2026-10-02 03:24:42 +08:00
atelet: pin the multi-actor stats fold and keep CPU baselines across pending sweeps (#1885)
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.
This commit is contained in:
+39
-47
@@ -48,24 +48,19 @@ import (
|
||||
const workerPoolLabel = "ate.dev/worker-pool"
|
||||
|
||||
// minActorStatsPollInterval is the floor a configured poll interval is clamped
|
||||
// to. It is the worst-case duration of one ateom's sweep on the micro-VM
|
||||
// runtime -- maxActorContainers containers at statsCallTimeout each, constants
|
||||
// that live with that runtime -- so a shorter interval could start a new poll
|
||||
// into a guest agent still serving the previous one.
|
||||
// to: one actor's worst-case micro-VM read, maxActorContainers (25) containers
|
||||
// at statsCallTimeout (2s) each. Sweeps never overlap, so it bounds
|
||||
// guest-agent load, not correctness.
|
||||
const minActorStatsPollInterval = 50 * time.Second
|
||||
|
||||
// statsRPCTimeout bounds one ateom's GetActiveWorkloadStats call. It has to
|
||||
// cover the ateom's own worst-case sweep (see minActorStatsPollInterval);
|
||||
// anything still unanswered past that is a stuck socket, not a slow guest.
|
||||
// statsRPCTimeout bounds one ateom's GetActiveWorkloadStats call. The micro-VM
|
||||
// ateom sets its statsSweepBudget (45s) below this timeout and reports guests
|
||||
// it did not reach as pending, so the call does not grow with actor count.
|
||||
const statsRPCTimeout = 55 * time.Second
|
||||
|
||||
// statsSweepConcurrency bounds how many ateoms one sweep probes at once. The
|
||||
// interval floor protects a single guest from overlapping polls; probing
|
||||
// DISTINCT ateoms concurrently puts one probe on each guest, so the only
|
||||
// stacking the cap prevents is on atelet itself -- without it, a node of
|
||||
// stuck-but-accepting sockets would hold one hung call per ateom for the full
|
||||
// statsRPCTimeout. With it, such a node degrades the sweep to
|
||||
// ceil(n/statsSweepConcurrency) timeouts instead of n.
|
||||
// statsSweepConcurrency bounds how many ateoms one sweep probes at once, one
|
||||
// call per ateom whatever it hosts. It caps how many stuck sockets can hold a
|
||||
// hung call open on atelet at the same time.
|
||||
const statsSweepConcurrency = 8
|
||||
|
||||
// workerPoolListTimeout bounds the per-sweep pod list that fetches worker
|
||||
@@ -74,11 +69,10 @@ const statsSweepConcurrency = 8
|
||||
const workerPoolListTimeout = 10 * time.Second
|
||||
|
||||
// clampActorStatsPollInterval enforces the floor on a nonzero configured
|
||||
// interval, warning rather than obeying: an interval below the worst-case
|
||||
// sweep would pile overlapping polls onto the same guest agent.
|
||||
// interval, warning rather than obeying.
|
||||
func clampActorStatsPollInterval(ctx context.Context, configured time.Duration) time.Duration {
|
||||
if configured > 0 && configured < minActorStatsPollInterval {
|
||||
slog.WarnContext(ctx, "actor-stats-poll-interval below the worst-case sweep; clamping",
|
||||
slog.WarnContext(ctx, "actor-stats-poll-interval below the floor; clamping",
|
||||
slog.Duration("configured", configured), slog.Duration("clamped_to", minActorStatsPollInterval))
|
||||
return minActorStatsPollInterval
|
||||
}
|
||||
@@ -142,13 +136,9 @@ type statsPoller struct {
|
||||
// exist, the same bound lastCPU keeps.
|
||||
cachedPools map[string]workerPoolRef
|
||||
|
||||
// lastCPU is the previous sweep's cpu_usage_usec per actor uid, the
|
||||
// baseline the next sweep's deltas are computed against. Only the sweep
|
||||
// loop touches it (under collect's mutex), and entries for actors a sweep
|
||||
// did not see are dropped at its end -- an actor that leaves the node
|
||||
// stops occupying memory here, and one that comes BACK later simply
|
||||
// re-baselines. Empty after an atelet restart, so the first sweep
|
||||
// contributes zero deltas: an undercount, never an overcount.
|
||||
// lastCPU is the last measured cpu_usage_usec per actor uid, the baseline
|
||||
// for the next delta; read-only during a sweep and replaced whole at the
|
||||
// end. Pending actors keep their entry, unseen actors are dropped.
|
||||
lastCPU map[string]uint64
|
||||
}
|
||||
|
||||
@@ -242,8 +232,9 @@ func (p *statsPoller) tick(ctx context.Context) {
|
||||
// worker pods (nothing garbage-collects them eagerly), ateoms that have made
|
||||
// their directory but not yet listened, and workers torn down mid-sweep. The
|
||||
// no-sample answers are equally routine: an empty samples list is an idle
|
||||
// worker, and a pending entry (source UNSPECIFIED) is a boot or restore in
|
||||
// progress -- both are skips by the RPC's own contract.
|
||||
// worker, and a pending entry (source UNSPECIFIED) is a workload the ateom
|
||||
// could not measure (boot, restore, teardown, or an unreached guest), which
|
||||
// adds nothing but keeps its CPU baseline.
|
||||
func (p *statsPoller) collect(ctx context.Context) map[templateKey]*templateAggregate {
|
||||
entries, err := os.ReadDir(p.ateomsDir)
|
||||
if err != nil {
|
||||
@@ -291,17 +282,20 @@ func (p *statsPoller) collect(ctx context.Context) map[templateKey]*templateAggr
|
||||
return nil
|
||||
}
|
||||
|
||||
// One entry per workload the ateom is executing; empty when it is
|
||||
// available. Today that is at most one entry -- multi-actor
|
||||
// workers will grow it, and nothing here assumes otherwise.
|
||||
// One entry per workload the ateom is hosting; empty when it is
|
||||
// available.
|
||||
for _, sample := range resp.GetSamples() {
|
||||
if sample.GetSource() == ateompb.StatsSource_STATS_SOURCE_UNSPECIFIED {
|
||||
// A pending workload: attributed but not measured (boot,
|
||||
// restore, teardown, or a transition underneath the
|
||||
// ateom's read). It contributes nothing anywhere -- not
|
||||
// to the aggregates (sampled_actors keeps meaning
|
||||
// "measured"), not to the CPU baselines, and no event --
|
||||
// exactly as the old no-sample answers behaved.
|
||||
// Pending: hosted but not measured, so it adds nothing
|
||||
// and its CPU baseline carries forward. A measured value
|
||||
// from another worker this sweep (restore in flight) wins.
|
||||
if last, ok := p.lastCPU[sample.GetActorUid()]; ok {
|
||||
mu.Lock()
|
||||
if _, measured := seenCPU[sample.GetActorUid()]; !measured {
|
||||
seenCPU[sample.GetActorUid()] = last
|
||||
}
|
||||
mu.Unlock()
|
||||
}
|
||||
continue
|
||||
}
|
||||
p.eventEmitter.emit(ctx, eventKindPeriodic, sample, pools[podUID])
|
||||
@@ -323,22 +317,20 @@ func (p *statsPoller) collect(ctx context.Context) map[templateKey]*templateAggr
|
||||
agg.memoryCurrentBytes = addSat(agg.memoryCurrentBytes, sample.GetMemoryCurrentBytes())
|
||||
agg.memoryWorkingSetBytes = addSat(agg.memoryWorkingSetBytes, sample.GetMemoryWorkingSetBytes())
|
||||
|
||||
// The counter increase this sample represents. A decrease means
|
||||
// the epoch reset underneath us (the cgroup source restarts at
|
||||
// zero on restore), so the new value IS the usage since the
|
||||
// reset. A sample with NO baseline charges nothing and only
|
||||
// records one: atelet cannot tell a new actor from its own
|
||||
// restart, and charging the whole epoch-so-far would re-count
|
||||
// hours of usage the previous atelet already counted, as one
|
||||
// artificial spike. The bounded price is that every actor's
|
||||
// boot-to-first-poll usage goes uncounted -- the events channel
|
||||
// carries per-actor precision.
|
||||
// The counter increase since the baseline. On a decrease, a
|
||||
// cgroup counter restarted at zero, so the new value is the
|
||||
// usage since then. A guest-agent decrease is ambiguous (a
|
||||
// resume at another guest's value, a restart at zero, or a
|
||||
// container exit), so it charges nothing. A sample with no
|
||||
// baseline charges nothing: atelet cannot tell a new actor
|
||||
// from its own restart.
|
||||
cpu := sample.GetCpuUsageUsec()
|
||||
seenCPU[sample.GetActorUid()] = cpu
|
||||
if last, ok := p.lastCPU[sample.GetActorUid()]; ok {
|
||||
if last <= cpu {
|
||||
switch {
|
||||
case last <= cpu:
|
||||
agg.cpuDeltaUsec = addSat(agg.cpuDeltaUsec, cpu-last)
|
||||
} else {
|
||||
case sample.GetSource() == ateompb.StatsSource_STATS_SOURCE_CGROUP:
|
||||
agg.cpuDeltaUsec = addSat(agg.cpuDeltaUsec, cpu)
|
||||
}
|
||||
}
|
||||
|
||||
@@ -23,6 +23,7 @@ import (
|
||||
"math"
|
||||
"os"
|
||||
"path/filepath"
|
||||
"sort"
|
||||
"sync"
|
||||
"testing"
|
||||
"time"
|
||||
@@ -52,9 +53,19 @@ type fakeStatsAteom struct {
|
||||
// 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()
|
||||
@@ -76,6 +87,31 @@ func executingResponse(templateNS, templateName string, class ateompb.SandboxCla
|
||||
}
|
||||
}
|
||||
|
||||
// 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{}
|
||||
@@ -95,15 +131,21 @@ func pendingResponse(actorUID string) *ateompb.GetActiveWorkloadStatsResponse {
|
||||
}
|
||||
|
||||
// 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
|
||||
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
|
||||
}
|
||||
|
||||
@@ -314,8 +356,8 @@ func cpuResponse(actorUID string, cpuUsec uint64) *ateompb.GetActiveWorkloadStat
|
||||
// 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 decrease is an epoch reset whose new value is the
|
||||
// usage since the reset, and an actor that disappears stops contributing and
|
||||
// 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"}
|
||||
@@ -639,3 +681,271 @@ func TestStatsPollerPoolCachePartialListFallsBack(t *testing.T) {
|
||||
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)
|
||||
}
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user