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
63 changes: 55 additions & 8 deletions internal/agent/agent.go
Original file line number Diff line number Diff line change
Expand Up @@ -31,9 +31,16 @@ type ControlPlane interface {
PollCommands(ctx context.Context) (*wire.CommandsReply, error)
}

// ErrNoSample is the error of a minute that has no sample, for example a
// first minute that is too short to measure. It is not a problem.
var ErrNoSample = errors.New("no sample for this minute")

// Collector measures the server.
type Collector interface {
// Sample measures the minute that ends at now.
// Read takes a reading between two ticks, for the peaks of the minute.
Read(now time.Time)
// Sample measures the minute that ends at now. It returns an error that
// wraps ErrNoSample when the minute has no sample.
Sample(now time.Time) (wire.Sample, error)
// Status describes the server now.
Status(ctx context.Context) wire.Status
Expand Down Expand Up @@ -133,17 +140,36 @@ func run(ctx context.Context, cfg Config, log *slog.Logger, cp ControlPlane, col
return nil
}

// loop calls tick at the offset second of each minute until ctx is done. It
// returns true when a command ends the process.
// loop calls tick at the offset second of each minute until ctx is done.
// Between two ticks, it reads the server each 10 seconds for the peaks of the
// minute. It returns true when a command ends the process.
func (a *agent) loop(ctx context.Context) (exit bool) {
for {
next := nextAfter(time.Now(), a.last, a.cfg.Offset())
timer := time.NewTimer(time.Until(next))
now := time.Now()
next := nextAfter(now, a.last, a.cfg.Offset())
wake := next
if a.collector != nil {
if r := nextReading(now, a.cfg.Offset()); r.Before(next) {
wake = r
}
}

timer := time.NewTimer(time.Until(wake))
select {
case <-ctx.Done():
timer.Stop()
return false
case <-timer.C:
if wake.Before(next) {
// The reading has the time at which it ran, so that a late
// timer or a step of the wall clock does not change the
// length of a window.
a.collector.Read(time.Now())
// A reading that ran late must not skip the tick after it.
if time.Now().Before(next) {
continue
}
}
a.last = next
if a.tick(ctx, next) {
return true
Expand All @@ -152,6 +178,22 @@ func (a *agent) loop(ctx context.Context) (exit bool) {
}
}

// readEvery is the time between two readings of the server (contract v0.4.0).
const readEvery = 10 * time.Second

// nextReading returns the next reading after now: offset past a full minute,
// and each 10 seconds after it. A tick is also a reading time: the loop then
// runs the tick instead. After a step back of the wall clock, a reading can
// come again; the collector ignores a reading that is not newer than its last.
func nextReading(now time.Time, offset time.Duration) time.Time {
t := now.Truncate(readEvery).Add(offset % readEvery)
for !t.After(now) {
t = t.Add(readEvery)
}

return t
}

// maxStepBack is the largest step back of the wall clock after which the loop
// still skips the minute that already ran. A larger step means that the clock
// was wrong before (for example a VM that booted with its clock ahead, then
Expand Down Expand Up @@ -182,10 +224,15 @@ func (a *agent) tick(ctx context.Context, now time.Time) (exit bool) {
a.log.Debug("tick", "at", now)

if a.collector != nil {
s, err := a.collector.Sample(now)
if err != nil {
// The reading has the time at which it ran; the sample has the time
// of its tick, for its minute on the control plane.
s, err := a.collector.Sample(time.Now())
switch {
case errors.Is(err, ErrNoSample):
a.log.Info("no sample for this minute", "reason", err)
case err != nil:
a.log.Warn("skipping the sample of this minute", "error", err)
} else {
default:
s.RecordedAt = now
a.outbox.addSample(cleanSample(s))
}
Expand Down
91 changes: 91 additions & 0 deletions internal/agent/agent_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -186,6 +186,97 @@ func TestLoopTicksAtTheOffsetAndReportsEachInterval(t *testing.T) {
})
}

func TestNextReading(t *testing.T) {
base := time.Date(2026, 9, 22, 10, 0, 0, 0, time.UTC)
tests := []struct {
now time.Time
offset time.Duration
want time.Time
}{
{base, 17 * time.Second, base.Add(7 * time.Second)},
{base.Add(7 * time.Second), 17 * time.Second, base.Add(17 * time.Second)},
{base.Add(18 * time.Second), 17 * time.Second, base.Add(27 * time.Second)},
{base.Add(58 * time.Second), 17 * time.Second, base.Add(67 * time.Second)},
{base.Add(59*time.Second + 999*time.Millisecond), 0, base.Add(time.Minute)},
{base.Add(3 * time.Second), 43 * time.Second, base.Add(3 * time.Second).Add(10 * time.Second)},
}

for _, tt := range tests {
if got := nextReading(tt.now, tt.offset); !got.Equal(tt.want) {
t.Errorf("nextReading(%s, %v) = %s, want %s", tt.now.Format(time.TimeOnly), tt.offset, got.Format(time.TimeOnly), tt.want.Format(time.TimeOnly))
}
}
}

func TestLoopReadsEach10SecondsBetweenTheTicks(t *testing.T) {
synctest.Test(t, func(t *testing.T) {
collector := &fakeCollector{}
ctx, cancel := context.WithCancel(context.Background())
done := make(chan error)
go func() {
done <- run(ctx, Config{ServerID: 17, StateDir: t.TempDir()}, slog.New(slog.DiscardHandler), &fakeCP{}, collector)
}()

time.Sleep(2 * time.Minute)
cancel()
if err := <-done; err != nil {
t.Fatal(err)
}

// The ticks are at :17. The readings are at :07, :27, :37, :47 and :57:
// the loop never reads at a tick, because the tick takes the reading.
start := time.Date(2000, 1, 1, 0, 0, 0, 0, time.UTC)
var want []time.Time
for s := 7 * time.Second; s < 2*time.Minute; s += 10 * time.Second {
if s%time.Minute != 17*time.Second {
want = append(want, start.Add(s))
}
}
got := collector.readTimes()
if len(got) != len(want) {
t.Fatalf("readings at %v, want %v", got, want)
}
for i := range want {
if !got[i].Equal(want[i]) {
t.Errorf("reading %d at %s, want %s", i, got[i].Format(time.TimeOnly), want[i].Format(time.TimeOnly))
}
}
if collector.n != 2 {
t.Errorf("%d samples, want 2 (at 0:17 and 1:17)", collector.n)
}
})
}

func TestASlowReadingDoesNotSkipTheTick(t *testing.T) {
synctest.Test(t, func(t *testing.T) {
// Each reading takes 10 s: the reading at :07 ends at the tick.
collector := &fakeCollector{readTime: 10 * time.Second}
rec := &recorder{}
ctx, cancel := context.WithCancel(context.Background())
done := make(chan error)
go func() {
done <- run(ctx, Config{ServerID: 17, StateDir: t.TempDir()}, slog.New(rec), &fakeCP{}, collector)
}()

time.Sleep(2*time.Minute + 30*time.Second)
cancel()
if err := <-done; err != nil {
t.Fatal(err)
}

start := time.Date(2000, 1, 1, 0, 0, 0, 0, time.UTC)
ticks := rec.times("tick")
if len(ticks) != 3 {
t.Fatalf("ticks at %v, want 3 ticks", ticks)
}
for i, got := range ticks {
if want := start.Add(time.Duration(i)*time.Minute + 17*time.Second); !got.Equal(want) {
t.Errorf("tick %d at %s, want %s", i, got.Format(time.TimeOnly), want.Format(time.TimeOnly))
}
}
})
}

func TestRunRefusesASecondAgent(t *testing.T) {
dir := t.TempDir()
unlock, err := lock(dir)
Expand Down
15 changes: 15 additions & 0 deletions internal/agent/clean.go
Original file line number Diff line number Diff line change
Expand Up @@ -32,6 +32,12 @@ func cleanSample(s wire.Sample) wire.Sample {
} {
*v = clampInt(*v)
}
s.CPUMaxPercent = clampPtr(s.CPUMaxPercent, 0, 100)
for _, v := range []**uint64{
&s.MemoryUsedMaxBytes, &s.SwapUsedMaxBytes, &s.NetInMaxBytesPerSecond, &s.NetOutMaxBytesPerSecond,
} {
*v = clampIntPtr(*v)
}
return s
}

Expand Down Expand Up @@ -87,6 +93,15 @@ func clamp(v, lo, hi float64) float64 {
return min(max(v, lo), hi)
}

// clampPtr is clamp for a value that can be "not known" (nil).
func clampPtr(v *float64, lo, hi float64) *float64 {
if v == nil {
return nil
}
c := clamp(*v, lo, hi)
return &c
}

// truncate cuts s to at most n characters. The control plane counts
// characters, not bytes.
func truncate(s string, n int) string {
Expand Down
18 changes: 18 additions & 0 deletions internal/agent/fakes_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -122,8 +122,26 @@ func (f *fakeCP) sampleCounts() []int {
type fakeCollector struct {
mu sync.Mutex
n int
reads []time.Time
status wire.Status
err error
// readTime is the time that each reading takes.
readTime time.Duration
}

func (c *fakeCollector) Read(now time.Time) {
c.mu.Lock()
c.reads = append(c.reads, now)
d := c.readTime
c.mu.Unlock()
time.Sleep(d)
}

// readTimes returns the times of the readings between the ticks.
func (c *fakeCollector) readTimes() []time.Time {
c.mu.Lock()
defer c.mu.Unlock()
return slices.Clone(c.reads)
}

func (c *fakeCollector) Sample(time.Time) (wire.Sample, error) {
Expand Down
39 changes: 39 additions & 0 deletions internal/agent/outbox_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -174,6 +174,45 @@ func TestCleanKeepsIntegersInTheRangeOfPHP(t *testing.T) {
}
}

func TestCleanPeaks(t *testing.T) {
cpu, big := 120.0, uint64(math.MaxUint64)
s := cleanSample(wire.Sample{CPUMaxPercent: &cpu, MemoryUsedMaxBytes: &big, NetInMaxBytesPerSecond: &big})
if *s.CPUMaxPercent != 100 || *s.MemoryUsedMaxBytes != math.MaxInt64 || *s.NetInMaxBytesPerSecond != math.MaxInt64 {
t.Errorf("cleanSample() = %v, %d, %d; want 100 and the largest PHP integer", *s.CPUMaxPercent, *s.MemoryUsedMaxBytes, *s.NetInMaxBytesPerSecond)
}
if cpu != 120 {
t.Error("cleanSample() changed the value of the caller")
}
if s.SwapUsedMaxBytes != nil || s.NetOutMaxBytesPerSecond != nil {
t.Error("cleanSample() gave a value to a peak that is not known")
}
}

// A sample that an older agent queued has no peaks. After an update, the new
// agent sends them as null: not known.
func TestQueuedSampleOfAnOlderAgentSendsNullPeaks(t *testing.T) {
dir := t.TempDir()
old := `[{"recorded_at":"2026-09-24T10:00:17Z","cpu_percent":12.5,"load_1":0.4,"memory_used_bytes":1,"memory_total_bytes":2,` +
`"swap_used_bytes":0,"swap_total_bytes":0,"disk_used_bytes":1,"disk_total_bytes":2,"net_in_bytes":5,"net_out_bytes":6,"net_counters_reset":false}]`
if err := os.WriteFile(filepath.Join(dir, "samples.json"), []byte(old), 0o600); err != nil {
t.Fatal(err)
}

o := loadOutbox(dir, slog.New(slog.DiscardHandler))
if len(o.samples) != 1 {
t.Fatalf("samples = %d, want the queued sample", len(o.samples))
}
data, err := json.Marshal(o.samples[0])
if err != nil {
t.Fatal(err)
}
for _, field := range []string{"cpu_max_percent", "memory_used_max_bytes", "swap_used_max_bytes", "net_in_max_bytes_per_second", "net_out_max_bytes_per_second"} {
if !strings.Contains(string(data), `"`+field+`":null`) {
t.Errorf("sample JSON = %s, want %s as null", data, field)
}
}
}

func TestCleanEventDropsACommandIDThatIsNotAULID(t *testing.T) {
if e := cleanEvent(wire.Event{CommandID: "not-a-ulid"}); e.CommandID != "" {
t.Errorf("command_id = %q, want it removed", e.CommandID)
Expand Down
19 changes: 19 additions & 0 deletions internal/agent/send_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -319,6 +319,25 @@ func TestARetryDoesNotWaitForTheNextInterval(t *testing.T) {
})
}

func TestAMinuteWithoutASampleIsNotAWarning(t *testing.T) {
synctest.Test(t, func(t *testing.T) {
collector := &fakeCollector{err: fmt.Errorf("%w: the first minute is only 1s since the start", ErrNoSample)}
cp := &fakeCP{}
rec := runFor(t, 2*time.Minute, t.TempDir(), cp, collector)

if n := len(rec.times("skipping the sample of this minute")); n != 0 {
t.Errorf("%d warnings, want none for a minute without a sample", n)
}
if got := rec.recordsOf("no sample for this minute"); len(got) != 2 || got[0].Level != slog.LevelInfo {
t.Errorf("%d info lines, want 2", len(got))
}
// The reports still go: the events are sent.
if len(cp.eventsAt) == 0 {
t.Error("no events request, want the report to run")
}
})
}

func TestTheReportTimeRunningOutIsNotAFailure(t *testing.T) {
synctest.Test(t, func(t *testing.T) {
dir := t.TempDir()
Expand Down
8 changes: 8 additions & 0 deletions internal/agent/wire/wire.go
Original file line number Diff line number Diff line change
Expand Up @@ -45,6 +45,14 @@ type Sample struct {
NetInBytes uint64 `json:"net_in_bytes"`
NetOutBytes uint64 `json:"net_out_bytes"`
NetCountersReset bool `json:"net_counters_reset"`

// The peaks within the minute, from the readings each 10 seconds
// (contract v0.4.0). nil (JSON null) means "not known".
CPUMaxPercent *float64 `json:"cpu_max_percent"`
MemoryUsedMaxBytes *uint64 `json:"memory_used_max_bytes"`
SwapUsedMaxBytes *uint64 `json:"swap_used_max_bytes"`
NetInMaxBytesPerSecond *uint64 `json:"net_in_max_bytes_per_second"`
NetOutMaxBytesPerSecond *uint64 `json:"net_out_max_bytes_per_second"`
}

// MetricsReply is the reply to POST /agent/v1/metrics.
Expand Down
Loading
Loading