From 6a22325d6587d9bf18ef0be0197a130aac535917 Mon Sep 17 00:00:00 2001 From: nabil1440 <52530910+nabil1440@users.noreply.github.com> Date: Tue, 22 Sep 2026 15:25:25 +0600 Subject: [PATCH 1/3] feat(agent): keep samples and events on disk until the control plane accepts them - Keep two queues in the state directory: samples.json (at most 1440, the oldest goes first) and events.json. Each change is a crash-safe write. - Make each event id (a ULID) when the event goes into the queue, so a resend has the same id. Keep only the last agent.started. - Keep each value in the range of the contract before it goes into a queue: one bad value makes a 400, and a 400 drops the whole request. - Send the events first, then the samples with the status, at most 100 events and 240 samples in one request, oldest first. - Follow the reply rules: drop on 200, 400 and 422 (samples); keep on 401 (retry each 5 minutes), 429 (Retry-After), 5xx and network errors (wait 1, 2, 4, 8, then 10 minutes). Apply report_interval from each reply. - Send agent.started at once at start, not at the next tick. Refs #27 --- agent_test.go | 34 +++++ cmd/agent.go | 2 +- go.mod | 1 + go.sum | 3 + internal/agent/agent.go | 90 ++++++++++-- internal/agent/agent_test.go | 4 +- internal/agent/clean.go | 64 +++++++++ internal/agent/fakes_test.go | 91 ++++++++++++ internal/agent/outbox.go | 93 ++++++++++++ internal/agent/outbox_test.go | 144 +++++++++++++++++++ internal/agent/send.go | 172 +++++++++++++++++++++++ internal/agent/send_test.go | 258 ++++++++++++++++++++++++++++++++++ internal/agent/wire/wire.go | 91 ++++++++++++ 13 files changed, 1035 insertions(+), 12 deletions(-) create mode 100644 internal/agent/clean.go create mode 100644 internal/agent/fakes_test.go create mode 100644 internal/agent/outbox.go create mode 100644 internal/agent/outbox_test.go create mode 100644 internal/agent/send.go create mode 100644 internal/agent/send_test.go create mode 100644 internal/agent/wire/wire.go diff --git a/agent_test.go b/agent_test.go index 34c56d4..3ad68ac 100644 --- a/agent_test.go +++ b/agent_test.go @@ -5,6 +5,9 @@ package main import ( "bytes" + "encoding/json" + "net/http" + "net/http/httptest" "os" "os/exec" "strings" @@ -89,6 +92,37 @@ func startAgent(t *testing.T, environ []string) (*exec.Cmd, *lockedBuffer) { return cmd, stderr } +func TestAgentSendsAgentStartedToTheControlPlane(t *testing.T) { + type request struct { + path, auth string + body map[string][]map[string]any + } + got := make(chan request, 10) + srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + req := request{path: r.URL.Path, auth: r.Header.Get("Authorization")} + _ = json.NewDecoder(r.Body).Decode(&req.body) + got <- req + _, _ = w.Write([]byte(`{"accepted": 1}`)) + })) + defer srv.Close() + + environ := append(agentEnv(t), "FLY_AGENT_URL="+srv.URL) + cmd, stderr := startAgent(t, environ) + defer func() { _ = cmd.Process.Signal(syscall.SIGTERM); _ = cmd.Wait() }() + + select { + case req := <-got: + if req.path != "/agent/v1/events" || req.auth != "Bearer "+testToken { + t.Errorf("request to %s with Authorization %q, want /agent/v1/events with the token", req.path, req.auth) + } + if events := req.body["events"]; len(events) != 1 || events[0]["name"] != "agent.started" { + t.Errorf("events = %v, want agent.started", events) + } + case <-time.After(10 * time.Second): + t.Fatalf("the control plane got no request within 10s. stderr:\n%s", stderr) + } +} + func TestAgentConfigErrors(t *testing.T) { tests := []struct { name string diff --git a/cmd/agent.go b/cmd/agent.go index b7e44d6..b7f4117 100644 --- a/cmd/agent.go +++ b/cmd/agent.go @@ -33,7 +33,7 @@ FLY_AGENT_SERVER_ID and STATE_DIRECTORY from the environment.`, ctx, stop := signal.NotifyContext(cmd.Context(), syscall.SIGTERM, os.Interrupt) defer stop() - return agent.Run(ctx, cfg, slog.New(slog.NewTextHandler(os.Stderr, nil))) + return agent.Run(ctx, cfg, slog.New(slog.NewTextHandler(os.Stderr, nil)), nil) }, } diff --git a/go.mod b/go.mod index b0bc569..8cf103c 100644 --- a/go.mod +++ b/go.mod @@ -6,6 +6,7 @@ toolchain go1.27.1 require ( github.com/fatih/color v1.19.0 + github.com/oklog/ulid/v2 v2.1.2 github.com/spf13/cobra v1.10.2 golang.org/x/mod v0.41.0 gopkg.in/yaml.v2 v2.4.0 diff --git a/go.sum b/go.sum index a4c3660..567b99e 100644 --- a/go.sum +++ b/go.sum @@ -7,6 +7,9 @@ github.com/mattn/go-colorable v0.1.15 h1:+u9SLTRGnXv73cEsnsmoZBom+dMU88B2M0aDcWy github.com/mattn/go-colorable v0.1.15/go.mod h1:6LmQG8QLFO4G5z1gPvYEzlUgJ2wF+stgPZH1UqBm1s8= github.com/mattn/go-isatty v0.0.24 h1:tGZZoVgT/KiqK1c8ocVLeDS8BSWMRd47J3Lbz7vsReI= github.com/mattn/go-isatty v0.0.24/go.mod h1:nMCL3Zebbrt45jsMDgnfIwz6ydEQApk5oEI3HqDio6A= +github.com/oklog/ulid/v2 v2.1.2 h1:IEclFb9JNvzYA6MW2SCxbLzcHTVsfqm3PrqGQJH5zec= +github.com/oklog/ulid/v2 v2.1.2/go.mod h1:rcEKHmBBKfef9DhnvX7y1HZBYxjXb0cP5ExxNsTT1QQ= +github.com/pborman/getopt v0.0.0-20170112200414-7148bc3a4c30/go.mod h1:85jBQOZwpVEaDAr341tbn15RS4fCAsIst0qp7i8ex1o= github.com/russross/blackfriday/v2 v2.1.0/go.mod h1:+Rmxgy9KzJVeS9/2gXHxylqXiyQDYRxCVz55jmeOWTM= github.com/spf13/cobra v1.10.2 h1:DMTTonx5m65Ic0GOoRY2c16WCbHxOOw6xxezuLaBpcU= github.com/spf13/cobra v1.10.2/go.mod h1:7C1pvHqHw5A4vrJfjNwvOdzYu0Gml16OCs2GRiTUUS4= diff --git a/internal/agent/agent.go b/internal/agent/agent.go index 11f53e9..be667ae 100644 --- a/internal/agent/agent.go +++ b/internal/agent/agent.go @@ -8,8 +8,10 @@ import ( "path/filepath" "time" + "github.com/flywp/server-cli/internal/agent/wire" "github.com/flywp/server-cli/internal/statefile" "github.com/flywp/server-cli/internal/version" + "github.com/oklog/ulid/v2" ) // The report interval is the number of samples in one report. The control @@ -19,14 +21,31 @@ const ( maxReportInterval = 10 ) +// ControlPlane is the part of the control plane that the agent uses. +type ControlPlane interface { + PostMetrics(ctx context.Context, req *wire.MetricsRequest) (*wire.MetricsReply, error) + PostEvents(ctx context.Context, req *wire.EventsRequest) (*wire.EventsReply, error) +} + +// Collector measures the server. +type Collector interface { + // Sample measures the minute that ends at now. + Sample(now time.Time) (wire.Sample, error) + // Status describes the server now. + Status(ctx context.Context) wire.Status +} + // state is the part of the agent state that is not a queue. type state struct { ReportInterval int `json:"report_interval"` } type agent struct { - cfg Config - log *slog.Logger + cfg Config + log *slog.Logger + cp ControlPlane + collector Collector + outbox *outbox // interval is the report interval, and pending is the number of samples // since the last report. @@ -35,11 +54,20 @@ type agent struct { // last is the time of the last tick. last time.Time + + // retryAt is the earliest time of the next send, and failures is the + // number of failed sends in a row. + retryAt time.Time + failures int } // Run runs the agent until ctx is done. Only one agent can run with the same -// state directory. -func Run(ctx context.Context, cfg Config, log *slog.Logger) error { +// state directory. A nil collector takes no samples. +func Run(ctx context.Context, cfg Config, log *slog.Logger, collector Collector) error { + return run(ctx, cfg, log, NewClient(cfg, nil), collector) +} + +func run(ctx context.Context, cfg Config, log *slog.Logger, cp ControlPlane, collector Collector) error { unlock, err := lock(cfg.StateDir) if err != nil { return err @@ -49,8 +77,21 @@ func Run(ctx context.Context, cfg Config, log *slog.Logger) error { // A crash during a write can leave a temporary file. statefile.RemoveTemp(cfg.StateDir) - a := &agent{cfg: cfg, log: log, interval: loadInterval(cfg.StateDir, log)} + a := &agent{ + cfg: cfg, + log: log, + cp: cp, + collector: collector, + outbox: loadOutbox(cfg.StateDir, log), + interval: loadInterval(cfg.StateDir, log), + } log.Info("agent started", "version", version.Version, "offset", cfg.Offset(), "report_interval", a.interval) + + // Send the events at once, not at the next tick: after an update or a + // restart, they hold the result of the command. + a.addEvent(wire.EventAgentStarted, "", &wire.EventData{Version: version.Version}) + a.send(ctx, false) + a.loop(ctx) log.Info("agent stopped") @@ -91,16 +132,47 @@ func nextAfter(now, last time.Time, offset time.Duration) time.Time { func (a *agent) tick(ctx context.Context, now time.Time) { a.log.Debug("tick", "at", now) + if a.collector != nil { + s, err := a.collector.Sample(now) + if err != nil { + a.log.Warn("skipping the sample of this minute", "error", err) + } else { + s.RecordedAt = now + a.outbox.addSample(cleanSample(s)) + } + } + a.pending++ if a.pending >= a.interval { a.pending = 0 - a.report(ctx) + a.log.Debug("report") + a.send(ctx, true) } } -// report sends the queued data to the control plane. -func (a *agent) report(_ context.Context) { - a.log.Debug("report") +// addEvent puts an event in the queue. Its ID is made now, so that each resend +// has the same ID. +func (a *agent) addEvent(name, commandID string, data *wire.EventData) { + a.outbox.addEvent(cleanEvent(wire.Event{ + ID: ulid.Make().String(), + Name: name, + CommandID: commandID, + At: time.Now().UTC(), + Data: data, + })) +} + +// setInterval applies the report interval of a reply. +func (a *agent) setInterval(n int) { + if n < minReportInterval || n > maxReportInterval || n == a.interval { + return + } + + a.log.Info("report interval changed", "from", a.interval, "to", n) + a.interval = n + if err := statefile.Write(filepath.Join(a.cfg.StateDir, "state.json"), state{ReportInterval: n}); err != nil { + a.log.Error("saving the report interval", "error", err) + } } // nextTick returns the first time after now that is offset past a full minute. diff --git a/internal/agent/agent_test.go b/internal/agent/agent_test.go index 0c0cdc8..81cb233 100644 --- a/internal/agent/agent_test.go +++ b/internal/agent/agent_test.go @@ -95,7 +95,7 @@ func TestLoopTicksAtTheOffsetAndReportsEachInterval(t *testing.T) { ctx, cancel := context.WithCancel(context.Background()) done := make(chan error) go func() { - done <- Run(ctx, Config{ServerID: 17, StateDir: dir}, slog.New(rec)) + done <- run(ctx, Config{ServerID: 17, StateDir: dir}, slog.New(rec), &fakeCP{}, nil) }() // The bubble starts at 2000-01-01 00:00:00 UTC. Let 4 minutes pass. @@ -134,7 +134,7 @@ func TestRunRefusesASecondAgent(t *testing.T) { } defer unlock() - if err := Run(context.Background(), Config{StateDir: dir}, slog.New(&recorder{})); err == nil { + if err := run(context.Background(), Config{StateDir: dir}, slog.New(&recorder{}), &fakeCP{}, nil); err == nil { t.Fatal("Run() = nil, want an error while a different agent holds the lock") } } diff --git a/internal/agent/clean.go b/internal/agent/clean.go new file mode 100644 index 0000000..cb923c1 --- /dev/null +++ b/internal/agent/clean.go @@ -0,0 +1,64 @@ +package agent + +import ( + "math" + "time" + "unicode/utf8" + + "github.com/flywp/server-cli/internal/agent/wire" +) + +// The value ranges of the contract (section 4, "Value ranges", and section 6). +// A value outside its range makes the control plane refuse the whole request +// with a 400, and the agent then drops up to 240 samples. So the agent keeps +// each value in its range before the value goes into a queue. +const ( + maxLoad = 999999.99 + maxVersionLen = 32 + maxStatusTextLen = 255 + maxArchLen = 16 + maxEventNameLen = 64 + maxErrorLen = 2000 +) + +func cleanSample(s wire.Sample) wire.Sample { + s.RecordedAt = s.RecordedAt.UTC().Truncate(time.Second) + s.CPUPercent = clamp(s.CPUPercent, 0, 100) + s.Load1 = clamp(s.Load1, 0, maxLoad) + return s +} + +func cleanStatus(s wire.Status) wire.Status { + s.OS = truncate(s.OS, maxStatusTextLen) + s.Kernel = truncate(s.Kernel, maxStatusTextLen) + s.Arch = truncate(s.Arch, maxArchLen) + return s +} + +func cleanEvent(e wire.Event) wire.Event { + e.Name = truncate(e.Name, maxEventNameLen) + if e.Data != nil { + d := *e.Data + d.Version = truncate(d.Version, maxVersionLen) + d.Error = truncate(d.Error, maxErrorLen) + e.Data = &d + } + return e +} + +// clamp keeps v between lo and hi. NaN becomes lo. +func clamp(v, lo, hi float64) float64 { + if math.IsNaN(v) { + return lo + } + return min(max(v, lo), hi) +} + +// truncate cuts s to at most n characters. The control plane counts +// characters, not bytes. +func truncate(s string, n int) string { + if utf8.RuneCountInString(s) <= n { + return s + } + return string([]rune(s)[:n]) +} diff --git a/internal/agent/fakes_test.go b/internal/agent/fakes_test.go new file mode 100644 index 0000000..9ea76d1 --- /dev/null +++ b/internal/agent/fakes_test.go @@ -0,0 +1,91 @@ +package agent + +import ( + "context" + "slices" + "sync" + "time" + + "github.com/flywp/server-cli/internal/agent/wire" +) + +// fakeCP is a control plane in memory. It keeps a copy of each request and +// the time of the request, and it answers with the reply funcs. A nil func +// accepts everything. +type fakeCP struct { + mu sync.Mutex + + metrics []wire.MetricsRequest + metricsAt []time.Time + events []wire.EventsRequest + eventsAt []time.Time + + // The reply funcs get the number of the call, from 0. + metricsReply func(call int, req *wire.MetricsRequest) (*wire.MetricsReply, error) + eventsReply func(call int, req *wire.EventsRequest) (*wire.EventsReply, error) +} + +func (f *fakeCP) PostMetrics(_ context.Context, req *wire.MetricsRequest) (*wire.MetricsReply, error) { + f.mu.Lock() + defer f.mu.Unlock() + + // Copy the samples: the agent reuses the memory of its queue. + c := *req + c.Samples = slices.Clone(req.Samples) + f.metrics = append(f.metrics, c) + f.metricsAt = append(f.metricsAt, time.Now()) + + if f.metricsReply == nil { + return &wire.MetricsReply{Accepted: len(req.Samples)}, nil + } + return f.metricsReply(len(f.metrics)-1, &c) +} + +func (f *fakeCP) PostEvents(_ context.Context, req *wire.EventsRequest) (*wire.EventsReply, error) { + f.mu.Lock() + defer f.mu.Unlock() + + c := wire.EventsRequest{Events: slices.Clone(req.Events)} + f.events = append(f.events, c) + f.eventsAt = append(f.eventsAt, time.Now()) + + if f.eventsReply == nil { + return &wire.EventsReply{Accepted: len(req.Events)}, nil + } + return f.eventsReply(len(f.events)-1, &c) +} + +// sampleCounts returns the number of samples in each metrics request. +func (f *fakeCP) sampleCounts() []int { + f.mu.Lock() + defer f.mu.Unlock() + + var n []int + for _, r := range f.metrics { + n = append(n, len(r.Samples)) + } + return n +} + +// fakeCollector returns samples whose CPU value counts the samples: 1, 2, 3... +type fakeCollector struct { + mu sync.Mutex + n int + status wire.Status + err error +} + +func (c *fakeCollector) Sample(time.Time) (wire.Sample, error) { + c.mu.Lock() + defer c.mu.Unlock() + + if c.err != nil { + return wire.Sample{}, c.err + } + c.n++ + return wire.Sample{CPUPercent: float64(c.n), MemoryTotalBytes: 1 << 30}, nil +} + +func (c *fakeCollector) Status(context.Context) wire.Status { + return c.status +} diff --git a/internal/agent/outbox.go b/internal/agent/outbox.go new file mode 100644 index 0000000..8fc9a7f --- /dev/null +++ b/internal/agent/outbox.go @@ -0,0 +1,93 @@ +package agent + +import ( + "errors" + "io/fs" + "log/slog" + "path/filepath" + "slices" + + "github.com/flywp/server-cli/internal/agent/wire" + "github.com/flywp/server-cli/internal/statefile" +) + +const ( + // maxSamples is 24 hours of samples (contract section 4). Above it, the + // oldest sample goes first. + maxSamples = 1440 + // maxEvents limits the event queue, for example when a crash loop adds + // events while the control plane is down. + maxEvents = 1000 +) + +// outbox keeps the samples and the events that the control plane has not +// accepted yet. Each change goes to disk, so a restart or a reboot loses +// nothing. +type outbox struct { + dir string + log *slog.Logger + samples []wire.Sample + events []wire.Event +} + +func loadOutbox(dir string, log *slog.Logger) *outbox { + o := &outbox{dir: dir, log: log} + o.samples = loadQueue[wire.Sample](o, "samples.json") + o.events = loadQueue[wire.Event](o, "events.json") + return o +} + +// loadQueue reads a queue file. A file that cannot be read gives an empty +// queue: to stop would only make systemd start the agent again. +func loadQueue[T any](o *outbox, name string) []T { + var q []T + err := statefile.Read(filepath.Join(o.dir, name), &q) + if err != nil && !errors.Is(err, fs.ErrNotExist) { + o.log.Error("dropping a queue that cannot be read", "file", name, "error", err) + return nil + } + + return q +} + +func (o *outbox) addSample(s wire.Sample) { + o.samples = append(o.samples, s) + if over := len(o.samples) - maxSamples; over > 0 { + o.samples = slices.Delete(o.samples, 0, over) + } + o.save("samples.json", o.samples) +} + +func (o *outbox) addEvent(e wire.Event) { + // The control plane records nothing for agent.started, so only the + // last one is useful. + if e.Name == wire.EventAgentStarted { + o.events = slices.DeleteFunc(o.events, func(q wire.Event) bool { return q.Name == wire.EventAgentStarted }) + } + + o.events = append(o.events, e) + if over := len(o.events) - maxEvents; over > 0 { + o.events = slices.Delete(o.events, 0, over) + } + o.save("events.json", o.events) +} + +// dropSamples removes the n oldest samples. +func (o *outbox) dropSamples(n int) { + o.samples = slices.Delete(o.samples, 0, n) + o.save("samples.json", o.samples) +} + +// dropEvents removes the n oldest events. +func (o *outbox) dropEvents(n int) { + o.events = slices.Delete(o.events, 0, n) + o.save("events.json", o.events) +} + +// save writes a queue to disk. When the write fails, the queue stays in memory +// and goes to disk with the next change. +func (o *outbox) save(name string, q any) { + if err := statefile.Write(filepath.Join(o.dir, name), q); err != nil { + o.log.Error("saving a queue", "file", name, "error", err) + } +} diff --git a/internal/agent/outbox_test.go b/internal/agent/outbox_test.go new file mode 100644 index 0000000..acf54a3 --- /dev/null +++ b/internal/agent/outbox_test.go @@ -0,0 +1,144 @@ +package agent + +import ( + "fmt" + "log/slog" + "math" + "os" + "path/filepath" + "strings" + "testing" + "time" + + "github.com/flywp/server-cli/internal/agent/wire" + "github.com/flywp/server-cli/internal/statefile" +) + +func TestOutboxDropsTheOldestSample(t *testing.T) { + dir := t.TempDir() + full := make([]wire.Sample, maxSamples) + for i := range full { + full[i].CPUPercent = float64(i) + } + if err := statefile.Write(filepath.Join(dir, "samples.json"), full); err != nil { + t.Fatal(err) + } + + o := loadOutbox(dir, slog.New(&recorder{})) + o.addSample(wire.Sample{CPUPercent: maxSamples}) + + if len(o.samples) != maxSamples { + t.Fatalf("queue holds %d samples, want %d", len(o.samples), maxSamples) + } + if o.samples[0].CPUPercent != 1 { + t.Errorf("oldest sample = %v, want sample 0 dropped", o.samples[0].CPUPercent) + } +} + +func TestOutboxSurvivesARestart(t *testing.T) { + dir := t.TempDir() + o := loadOutbox(dir, slog.New(&recorder{})) + o.addSample(wire.Sample{CPUPercent: 1}) + o.addSample(wire.Sample{CPUPercent: 2}) + o.addEvent(wire.Event{ID: "01JBY0000000000000000000AA", Name: wire.EventCommandCompleted, CommandID: "01JBX0000000000000000000AA"}) + o.dropSamples(1) + + again := loadOutbox(dir, slog.New(&recorder{})) + if len(again.samples) != 1 || again.samples[0].CPUPercent != 2 { + t.Errorf("samples after a restart = %+v, want the second sample only", again.samples) + } + if len(again.events) != 1 || again.events[0].ID != "01JBY0000000000000000000AA" { + t.Errorf("events after a restart = %+v, want the event with its id", again.events) + } +} + +func TestOutboxKeepsOnlyTheLastAgentStarted(t *testing.T) { + o := loadOutbox(t.TempDir(), slog.New(&recorder{})) + o.addEvent(wire.Event{ID: "a", Name: wire.EventAgentStarted}) + o.addEvent(wire.Event{ID: "b", Name: wire.EventCommandCompleted}) + o.addEvent(wire.Event{ID: "c", Name: wire.EventAgentStarted}) + + var ids []string + for _, e := range o.events { + ids = append(ids, e.ID) + } + if strings.Join(ids, ",") != "b,c" { + t.Errorf("events = %v, want b,c: an older agent.started goes", ids) + } +} + +func TestOutboxLimitsEvents(t *testing.T) { + dir := t.TempDir() + full := make([]wire.Event, maxEvents) + for i := range full { + full[i] = wire.Event{ID: fmt.Sprint(i), Name: wire.EventCommandFailed} + } + if err := statefile.Write(filepath.Join(dir, "events.json"), full); err != nil { + t.Fatal(err) + } + + o := loadOutbox(dir, slog.New(&recorder{})) + o.addEvent(wire.Event{ID: "new", Name: wire.EventCommandFailed}) + if len(o.events) != maxEvents || o.events[0].ID != "1" || o.events[maxEvents-1].ID != "new" { + t.Errorf("queue holds %d events from %s to %s, want %d with event 0 dropped", len(o.events), o.events[0].ID, o.events[len(o.events)-1].ID, maxEvents) + } +} + +func TestOutboxWithAQueueThatCannotBeRead(t *testing.T) { + dir := t.TempDir() + if err := os.WriteFile(filepath.Join(dir, "samples.json"), []byte("[{"), 0o600); err != nil { + t.Fatal(err) + } + + rec := &recorder{} + o := loadOutbox(dir, slog.New(rec)) + if len(o.samples) != 0 { + t.Errorf("samples = %d, want an empty queue", len(o.samples)) + } + if len(rec.times("dropping a queue that cannot be read")) != 1 { + t.Error("want an error log line for the queue that cannot be read") + } +} + +func TestCleanSample(t *testing.T) { + recorded := time.Date(2026, 9, 22, 10, 0, 17, 123456789, time.FixedZone("x", 3600)) + tests := []struct { + cpu, load float64 + wantCPU, wantLoad float64 + }{ + {12.5, 0.4, 12.5, 0.4}, + {150, 2e6, 100, maxLoad}, + {-1, -3, 0, 0}, + {math.NaN(), math.NaN(), 0, 0}, + } + + for _, tt := range tests { + got := cleanSample(wire.Sample{RecordedAt: recorded, CPUPercent: tt.cpu, Load1: tt.load}) + if got.CPUPercent != tt.wantCPU || got.Load1 != tt.wantLoad { + t.Errorf("cleanSample(cpu %v, load %v) = %v, %v; want %v, %v", tt.cpu, tt.load, got.CPUPercent, got.Load1, tt.wantCPU, tt.wantLoad) + } + if want := time.Date(2026, 9, 22, 9, 0, 17, 0, time.UTC); !got.RecordedAt.Equal(want) || got.RecordedAt.Location() != time.UTC { + t.Errorf("recorded_at = %s, want %s in UTC", got.RecordedAt, want) + } + } +} + +func TestCleanTextFields(t *testing.T) { + long := strings.Repeat("é", 300) + + s := cleanStatus(wire.Status{OS: long, Kernel: long, Arch: strings.Repeat("x", 20)}) + if n := len([]rune(s.OS)); n != maxStatusTextLen { + t.Errorf("os has %d characters, want %d", n, maxStatusTextLen) + } + if n := len([]rune(s.Kernel)); n != maxStatusTextLen { + t.Errorf("kernel has %d characters, want %d", n, maxStatusTextLen) + } + if len(s.Arch) != maxArchLen { + t.Errorf("arch has %d characters, want %d", len(s.Arch), maxArchLen) + } + + e := cleanEvent(wire.Event{Name: strings.Repeat("n", 70), Data: &wire.EventData{Version: strings.Repeat("v", 40), Error: strings.Repeat("e", 2500)}}) + if len(e.Name) != maxEventNameLen || len(e.Data.Version) != maxVersionLen || len(e.Data.Error) != maxErrorLen { + t.Errorf("event lengths = %d, %d, %d; want %d, %d, %d", len(e.Name), len(e.Data.Version), len(e.Data.Error), maxEventNameLen, maxVersionLen, maxErrorLen) + } +} diff --git a/internal/agent/send.go b/internal/agent/send.go new file mode 100644 index 0000000..4ee6317 --- /dev/null +++ b/internal/agent/send.go @@ -0,0 +1,172 @@ +package agent + +import ( + "context" + "errors" + "net/http" + "time" + + "github.com/flywp/server-cli/internal/agent/wire" + "github.com/flywp/server-cli/internal/version" +) + +const ( + // sendBudget limits the sends of one report, so that the next tick + // comes on time. + sendBudget = 45 * time.Second + + // The contract limits each request. + maxSamplesPerRequest = 240 + maxEventsPerRequest = 100 + + // unauthorizedWait is the wait after a 401: the token is unknown or + // revoked, and the agent keeps trying (contract section 4). + unauthorizedWait = 5 * time.Minute + // throttledWait is the wait after a 429 without a Retry-After header. + throttledWait = time.Minute + // The wait after other failures doubles from 1 minute up to 10 minutes. + firstBackoff = time.Minute + maxBackoff = 10 * time.Minute +) + +// outcome tells what to do with the data of a request. +type outcome int + +const ( + // sent: drop the data, and send the next request. + sent outcome = iota + // refused: the control plane will never accept the data. Drop it, and + // send the next request. + refused + // later: keep the data, and send nothing until retryAt. + later +) + +// send sends the events and then the samples. With report false, it sends +// only the events. +func (a *agent) send(ctx context.Context, report bool) { + if wait := time.Until(a.retryAt); wait > 0 { + a.log.Debug("waiting before the next send", "wait", wait) + return + } + + ctx, cancel := context.WithTimeout(ctx, sendBudget) + defer cancel() + + if a.sendEvents(ctx) && report { + a.sendSamples(ctx) + } +} + +// sendEvents sends the queued events, oldest first. It returns false when the +// agent must send nothing more now. +func (a *agent) sendEvents(ctx context.Context) bool { + for len(a.outbox.events) > 0 { + n := min(len(a.outbox.events), maxEventsPerRequest) + _, err := a.cp.PostEvents(ctx, &wire.EventsRequest{Events: a.outbox.events[:n]}) + if a.outcome(err, "events", n) == later { + return false + } + a.outbox.dropEvents(n) + } + + return true +} + +// sendSamples sends the queued samples, oldest first, with the status of the +// server in each request. +func (a *agent) sendSamples(ctx context.Context) { + if len(a.outbox.samples) == 0 { + return + } + + var status *wire.Status + if a.collector != nil { + s := cleanStatus(a.collector.Status(ctx)) + status = &s + } + + for len(a.outbox.samples) > 0 { + n := min(len(a.outbox.samples), maxSamplesPerRequest) + reply, err := a.cp.PostMetrics(ctx, &wire.MetricsRequest{ + AgentVersion: truncate(version.Version, maxVersionLen), + Status: status, + Samples: a.outbox.samples[:n], + }) + + switch a.outcome(err, "samples", n) { + case later: + return + case sent: + if len(reply.Rejected) > 0 { + a.log.Warn("the control plane rejected some samples", "rejected", len(reply.Rejected), "first_reason", reply.Rejected[0].Reason) + } + a.setInterval(reply.ReportInterval) + } + a.outbox.dropSamples(n) + } +} + +// outcome applies the rules of the contract (sections 4 and 6) to the result +// of a request that carried n items of what. +func (a *agent) outcome(err error, what string, n int) outcome { + if err == nil { + a.failures = 0 + return sent + } + + var statusErr *StatusError + if errors.As(err, &statusErr) { + switch code := statusErr.StatusCode; { + case code == http.StatusBadRequest: + a.failures = 0 + a.log.Error("the control plane refused the "+what+" as not valid; dropping them", "count", n) + return refused + case code == http.StatusUnprocessableEntity && what == "samples": + a.failures = 0 + a.log.Warn("the control plane rejected every sample; dropping them", "count", n) + return refused + case code == http.StatusUnauthorized: + a.retryAt = time.Now().Add(unauthorizedWait) + a.log.Error("the control plane does not accept the token; keeping the data", "retry_in", unauthorizedWait) + return later + case code == http.StatusTooManyRequests: + wait := statusErr.RetryAfter + if wait <= 0 { + wait = throttledWait + } + a.retryAt = time.Now().Add(wait) + a.log.Warn("the control plane asks the agent to wait; keeping the data", "retry_in", wait) + return later + } + } + + // A 5xx, another status or a network error: keep the data and wait longer + // after each failure. + a.failures++ + wait := min(firstBackoff< Date: Tue, 22 Sep 2026 16:20:44 +0600 Subject: [PATCH 2/3] fix(agent): send metrics while events fail, and keep every value in range Fixes from the adversarial review of this layer. - Cap each integer at the signed 64-bit limit: the control plane (PHP) fails a larger value with a 500, and the agent would resend the same sample for 24 hours. - Events and metrics wait on their own after a failure. A broken events route no longer stops the metrics. (The poll for commands, in the next layer, still waits until the events are sent.) - A wait counts from the start of the report, so a 5 minute wait ends at the tick 5 minutes later, also when the control plane is slow. - Samples that could not go are tried again at the next tick that their wait allows, not only after the next full report interval. - When the time of a report runs out, the rest goes with the next report. That is not a failure of the control plane and makes no backoff. - Log the start of an error reply, for example the validation errors of a 400. Log each sample or event that a full queue drops. - Keep a queue file that cannot be read as .corrupt. - Remove an event's command_id that is not a ULID, so a bad id cannot make a 400 that drops the other events. Refs #27 --- internal/agent/agent.go | 22 ++++--- internal/agent/clean.go | 24 ++++++++ internal/agent/client.go | 16 ++++- internal/agent/client_test.go | 5 ++ internal/agent/fakes_test.go | 14 ++++- internal/agent/outbox.go | 10 +++- internal/agent/outbox_test.go | 32 +++++++++- internal/agent/send.go | 106 +++++++++++++++++++++++----------- internal/agent/send_test.go | 89 ++++++++++++++++++++++++++-- 9 files changed, 263 insertions(+), 55 deletions(-) diff --git a/internal/agent/agent.go b/internal/agent/agent.go index be667ae..bfbcc14 100644 --- a/internal/agent/agent.go +++ b/internal/agent/agent.go @@ -55,10 +55,11 @@ type agent struct { // last is the time of the last tick. last time.Time - // retryAt is the earliest time of the next send, and failures is the - // number of failed sends in a row. - retryAt time.Time - failures int + // The waits of the events and the metrics requests after a failure, and + // whether the last report sent all samples. + eventsWait backoff + metricsWait backoff + samplesSent bool } // Run runs the agent until ctx is done. Only one agent can run with the same @@ -143,10 +144,17 @@ func (a *agent) tick(ctx context.Context, now time.Time) { } a.pending++ - if a.pending >= a.interval { + if a.pending < a.interval { + return + } + + a.log.Debug("report") + a.send(ctx, true) + + // Samples that could not go are tried again at the next tick, when their + // wait allows it, not only after the next full interval. + if a.samplesSent { a.pending = 0 - a.log.Debug("report") - a.send(ctx, true) } } diff --git a/internal/agent/clean.go b/internal/agent/clean.go index cb923c1..1268168 100644 --- a/internal/agent/clean.go +++ b/internal/agent/clean.go @@ -6,6 +6,7 @@ import ( "unicode/utf8" "github.com/flywp/server-cli/internal/agent/wire" + "github.com/oklog/ulid/v2" ) // The value ranges of the contract (section 4, "Value ranges", and section 6). @@ -25,6 +26,12 @@ func cleanSample(s wire.Sample) wire.Sample { s.RecordedAt = s.RecordedAt.UTC().Truncate(time.Second) s.CPUPercent = clamp(s.CPUPercent, 0, 100) s.Load1 = clamp(s.Load1, 0, maxLoad) + for _, v := range []*uint64{ + &s.MemoryUsedBytes, &s.MemoryTotalBytes, &s.SwapUsedBytes, &s.SwapTotalBytes, + &s.DiskUsedBytes, &s.DiskTotalBytes, &s.NetInBytes, &s.NetOutBytes, + } { + *v = clampInt(*v) + } return s } @@ -32,10 +39,20 @@ func cleanStatus(s wire.Status) wire.Status { s.OS = truncate(s.OS, maxStatusTextLen) s.Kernel = truncate(s.Kernel, maxStatusTextLen) s.Arch = truncate(s.Arch, maxArchLen) + s.UpdatesTotal = clampInt(s.UpdatesTotal) + s.UpdatesSecurity = clampInt(s.UpdatesSecurity) + s.UptimeSeconds = clampInt(s.UptimeSeconds) return s } func cleanEvent(e wire.Event) wire.Event { + // A command id that is not a ULID makes a 400, and a 400 drops all the + // events of the request. Without the id, the event changes nothing. + if e.CommandID != "" { + if _, err := ulid.ParseStrict(e.CommandID); err != nil { + e.CommandID = "" + } + } e.Name = truncate(e.Name, maxEventNameLen) if e.Data != nil { d := *e.Data @@ -46,6 +63,13 @@ func cleanEvent(e wire.Event) wire.Event { return e } +// clampInt keeps v in the range of a signed 64-bit integer: the control plane +// (PHP) cannot hold a larger integer, and fails the request with a 500. The +// agent would then send the same sample again for 24 hours. +func clampInt(v uint64) uint64 { + return min(v, math.MaxInt64) +} + // clamp keeps v between lo and hi. NaN becomes lo. func clamp(v, lo, hi float64) float64 { if math.IsNaN(v) { diff --git a/internal/agent/client.go b/internal/agent/client.go index a7b8d29..83f112d 100644 --- a/internal/agent/client.go +++ b/internal/agent/client.go @@ -9,6 +9,7 @@ import ( "net/http" "net/url" "strconv" + "strings" "time" "github.com/flywp/server-cli/internal/version" @@ -52,11 +53,18 @@ func noRedirects(*http.Request, []*http.Request) error { return http.ErrUseLastResponse } +// maxErrorBody limits the part of an error reply that the agent keeps for +// its log. +const maxErrorBody = 1 << 10 + // StatusError is a reply from the control plane that is not 200 OK. type StatusError struct { StatusCode int // RetryAfter is the wait that the Retry-After header asks for, or 0. RetryAfter time.Duration + // Body is the start of the reply, for example the validation errors of + // a 400. + Body string } func (e *StatusError) Error() string { @@ -93,8 +101,12 @@ func (c *Client) do(ctx context.Context, method, path string, in, out any) error defer func() { _ = resp.Body.Close() }() if resp.StatusCode != http.StatusOK { - _, _ = io.Copy(io.Discard, io.LimitReader(resp.Body, maxReplySize)) - return &StatusError{StatusCode: resp.StatusCode, RetryAfter: retryAfter(resp.Header.Get("Retry-After"), time.Now())} + body, _ := io.ReadAll(io.LimitReader(resp.Body, maxErrorBody)) + return &StatusError{ + StatusCode: resp.StatusCode, + RetryAfter: retryAfter(resp.Header.Get("Retry-After"), time.Now()), + Body: strings.ToValidUTF8(string(body), "?"), + } } if out == nil { diff --git a/internal/agent/client_test.go b/internal/agent/client_test.go index 2f91bcc..a0fe988 100644 --- a/internal/agent/client_test.go +++ b/internal/agent/client_test.go @@ -7,6 +7,7 @@ import ( "net/http" "net/http/httptest" "net/url" + "strings" "testing" "time" @@ -100,6 +101,7 @@ func TestClientStatusError(t *testing.T) { w.Header().Set("Retry-After", tt.retryAfter) } w.WriteHeader(tt.code) + _, _ = w.Write([]byte(`{"message":"The samples.0.cpu_percent field must be between 0 and 100."}`)) }) err := c.do(context.Background(), http.MethodPost, "agent/v1/events", map[string]int{}, nil) @@ -110,6 +112,9 @@ func TestClientStatusError(t *testing.T) { if statusErr.StatusCode != tt.code || statusErr.RetryAfter != tt.want { t.Errorf("StatusError = %+v, want code %d and wait %v", statusErr, tt.code, tt.want) } + if !strings.Contains(statusErr.Body, "cpu_percent") { + t.Errorf("StatusError.Body = %q, want the reply for the log", statusErr.Body) + } }) } } diff --git a/internal/agent/fakes_test.go b/internal/agent/fakes_test.go index 9ea76d1..79841e6 100644 --- a/internal/agent/fakes_test.go +++ b/internal/agent/fakes_test.go @@ -20,12 +20,24 @@ type fakeCP struct { events []wire.EventsRequest eventsAt []time.Time + // latency is the time of each metrics request. The request ends early + // when its context ends, like a real HTTP request. + latency time.Duration + // The reply funcs get the number of the call, from 0. metricsReply func(call int, req *wire.MetricsRequest) (*wire.MetricsReply, error) eventsReply func(call int, req *wire.EventsRequest) (*wire.EventsReply, error) } -func (f *fakeCP) PostMetrics(_ context.Context, req *wire.MetricsRequest) (*wire.MetricsReply, error) { +func (f *fakeCP) PostMetrics(ctx context.Context, req *wire.MetricsRequest) (*wire.MetricsReply, error) { + if f.latency > 0 { + select { + case <-time.After(f.latency): + case <-ctx.Done(): + return nil, ctx.Err() + } + } + f.mu.Lock() defer f.mu.Unlock() diff --git a/internal/agent/outbox.go b/internal/agent/outbox.go index 8fc9a7f..a7e4b27 100644 --- a/internal/agent/outbox.go +++ b/internal/agent/outbox.go @@ -4,6 +4,7 @@ import ( "errors" "io/fs" "log/slog" + "os" "path/filepath" "slices" @@ -41,9 +42,12 @@ func loadOutbox(dir string, log *slog.Logger) *outbox { // queue: to stop would only make systemd start the agent again. func loadQueue[T any](o *outbox, name string) []T { var q []T - err := statefile.Read(filepath.Join(o.dir, name), &q) + path := filepath.Join(o.dir, name) + err := statefile.Read(path, &q) if err != nil && !errors.Is(err, fs.ErrNotExist) { - o.log.Error("dropping a queue that cannot be read", "file", name, "error", err) + // Keep the file for a person to examine: the next save replaces it. + _ = os.Rename(path, path+".corrupt") + o.log.Error("dropping a queue that cannot be read; the file is kept as "+name+".corrupt", "file", name, "error", err) return nil } @@ -54,6 +58,7 @@ func (o *outbox) addSample(s wire.Sample) { o.samples = append(o.samples, s) if over := len(o.samples) - maxSamples; over > 0 { o.samples = slices.Delete(o.samples, 0, over) + o.log.Warn("the sample queue is full (24 hours); dropping the oldest sample", "dropped", over) } o.save("samples.json", o.samples) } @@ -68,6 +73,7 @@ func (o *outbox) addEvent(e wire.Event) { o.events = append(o.events, e) if over := len(o.events) - maxEvents; over > 0 { o.events = slices.Delete(o.events, 0, over) + o.log.Warn("the event queue is full; dropping the oldest event", "dropped", over) } o.save("events.json", o.events) } diff --git a/internal/agent/outbox_test.go b/internal/agent/outbox_test.go index acf54a3..dbba21e 100644 --- a/internal/agent/outbox_test.go +++ b/internal/agent/outbox_test.go @@ -24,8 +24,12 @@ func TestOutboxDropsTheOldestSample(t *testing.T) { t.Fatal(err) } - o := loadOutbox(dir, slog.New(&recorder{})) + rec := &recorder{} + o := loadOutbox(dir, slog.New(rec)) o.addSample(wire.Sample{CPUPercent: maxSamples}) + if len(rec.times("the sample queue is full (24 hours); dropping the oldest sample")) != 1 { + t.Error("want a warning when the full queue drops a sample") + } if len(o.samples) != maxSamples { t.Fatalf("queue holds %d samples, want %d", len(o.samples), maxSamples) @@ -95,9 +99,12 @@ func TestOutboxWithAQueueThatCannotBeRead(t *testing.T) { if len(o.samples) != 0 { t.Errorf("samples = %d, want an empty queue", len(o.samples)) } - if len(rec.times("dropping a queue that cannot be read")) != 1 { + if len(rec.times("dropping a queue that cannot be read; the file is kept as samples.json.corrupt")) != 1 { t.Error("want an error log line for the queue that cannot be read") } + if data, err := os.ReadFile(filepath.Join(dir, "samples.json.corrupt")); err != nil || string(data) != "[{" { + t.Errorf("samples.json.corrupt = %q, %v; want the file kept for a person to examine", data, err) + } } func TestCleanSample(t *testing.T) { @@ -142,3 +149,24 @@ func TestCleanTextFields(t *testing.T) { t.Errorf("event lengths = %d, %d, %d; want %d, %d, %d", len(e.Name), len(e.Data.Version), len(e.Data.Error), maxEventNameLen, maxVersionLen, maxErrorLen) } } + +func TestCleanKeepsIntegersInTheRangeOfPHP(t *testing.T) { + s := cleanSample(wire.Sample{NetInBytes: math.MaxUint64, MemoryTotalBytes: 1 << 63, DiskUsedBytes: 42}) + if s.NetInBytes != math.MaxInt64 || s.MemoryTotalBytes != math.MaxInt64 || s.DiskUsedBytes != 42 { + t.Errorf("cleanSample() = %d, %d, %d; want %d, %d, 42", s.NetInBytes, s.MemoryTotalBytes, s.DiskUsedBytes, uint64(math.MaxInt64), uint64(math.MaxInt64)) + } + + st := cleanStatus(wire.Status{UpdatesTotal: math.MaxUint64, UpdatesSecurity: 1 << 63, UptimeSeconds: math.MaxUint64}) + if st.UpdatesTotal != math.MaxInt64 || st.UpdatesSecurity != math.MaxInt64 || st.UptimeSeconds != math.MaxInt64 { + t.Errorf("cleanStatus() = %+v, want each count at most %d", st, uint64(math.MaxInt64)) + } +} + +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) + } + if e := cleanEvent(wire.Event{CommandID: "01JBX0000000000000000000AA"}); e.CommandID != "01JBX0000000000000000000AA" { + t.Errorf("command_id = %q, want the ULID kept", e.CommandID) + } +} diff --git a/internal/agent/send.go b/internal/agent/send.go index 4ee6317..86c81d8 100644 --- a/internal/agent/send.go +++ b/internal/agent/send.go @@ -38,33 +38,53 @@ const ( // refused: the control plane will never accept the data. Drop it, and // send the next request. refused - // later: keep the data, and send nothing until retryAt. + // later: keep the data, and send no more of it now. later ) -// send sends the events and then the samples. With report false, it sends -// only the events. -func (a *agent) send(ctx context.Context, report bool) { - if wait := time.Until(a.retryAt); wait > 0 { - a.log.Debug("waiting before the next send", "wait", wait) - return - } +// backoff is the wait of one kind of request after a failure. Each kind +// waits on its own: a broken events route must not stop the metrics. +type backoff struct { + // retryAt is the earliest time of the next request, and failures is the + // number of failed requests in a row. + retryAt time.Time + failures int +} + +// waiting reports whether the request must wait at the time of the report. +func (b *backoff) waiting(at time.Time) bool { + return at.Before(b.retryAt) +} +// send sends the events and then, with report, the samples. It returns true +// when no event is left in the queue, so that the commands can come next. +func (a *agent) send(ctx context.Context, report bool) (eventsSent bool) { + // The waits count from the start of the report, not from the end of a + // request: a 5 minute wait then ends at the tick 5 minutes later. + start := time.Now() ctx, cancel := context.WithTimeout(ctx, sendBudget) defer cancel() - if a.sendEvents(ctx) && report { - a.sendSamples(ctx) + eventsSent = a.sendEvents(ctx, start) + if report { + a.samplesSent = a.sendSamples(ctx, start) } + + return eventsSent } -// sendEvents sends the queued events, oldest first. It returns false when the -// agent must send nothing more now. -func (a *agent) sendEvents(ctx context.Context) bool { +// sendEvents sends the queued events, oldest first. It returns true when the +// event queue is empty. +func (a *agent) sendEvents(ctx context.Context, start time.Time) bool { + if a.eventsWait.waiting(start) { + a.log.Debug("waiting before the next events request", "until", a.eventsWait.retryAt) + return len(a.outbox.events) == 0 + } + for len(a.outbox.events) > 0 { n := min(len(a.outbox.events), maxEventsPerRequest) _, err := a.cp.PostEvents(ctx, &wire.EventsRequest{Events: a.outbox.events[:n]}) - if a.outcome(err, "events", n) == later { + if a.outcome(ctx, &a.eventsWait, start, err, "events", n) == later { return false } a.outbox.dropEvents(n) @@ -74,10 +94,14 @@ func (a *agent) sendEvents(ctx context.Context) bool { } // sendSamples sends the queued samples, oldest first, with the status of the -// server in each request. -func (a *agent) sendSamples(ctx context.Context) { +// server in each request. It returns true when the sample queue is empty. +func (a *agent) sendSamples(ctx context.Context, start time.Time) bool { if len(a.outbox.samples) == 0 { - return + return true + } + if a.metricsWait.waiting(start) { + a.log.Debug("waiting before the next metrics request", "until", a.metricsWait.retryAt) + return false } var status *wire.Status @@ -94,9 +118,9 @@ func (a *agent) sendSamples(ctx context.Context) { Samples: a.outbox.samples[:n], }) - switch a.outcome(err, "samples", n) { + switch a.outcome(ctx, &a.metricsWait, start, err, "samples", n) { case later: - return + return false case sent: if len(reply.Rejected) > 0 { a.log.Warn("the control plane rejected some samples", "rejected", len(reply.Rejected), "first_reason", reply.Rejected[0].Reason) @@ -105,48 +129,60 @@ func (a *agent) sendSamples(ctx context.Context) { } a.outbox.dropSamples(n) } + + return true } -// outcome applies the rules of the contract (sections 4 and 6) to the result -// of a request that carried n items of what. -func (a *agent) outcome(err error, what string, n int) outcome { +// outcome applies the rules of the contract (sections 4 to 6) to the result +// of a request that carried n items of what. The waits go into b and count +// from start. +func (a *agent) outcome(ctx context.Context, b *backoff, start time.Time, err error, what string, n int) outcome { if err == nil { - a.failures = 0 + b.failures = 0 return sent } + // The time of this report ran out, or the agent stops. That is not a + // failure of the control plane: the data goes with the next report. + if ctx.Err() != nil { + if errors.Is(ctx.Err(), context.DeadlineExceeded) { + a.log.Info("the time for this report is used; the rest goes with the next report", "request", what) + } + return later + } + var statusErr *StatusError if errors.As(err, &statusErr) { switch code := statusErr.StatusCode; { case code == http.StatusBadRequest: - a.failures = 0 - a.log.Error("the control plane refused the "+what+" as not valid; dropping them", "count", n) + b.failures = 0 + a.log.Error("the control plane refused the request as not valid; dropping its data", "request", what, "count", n, "reply", statusErr.Body) return refused case code == http.StatusUnprocessableEntity && what == "samples": - a.failures = 0 - a.log.Warn("the control plane rejected every sample; dropping them", "count", n) + b.failures = 0 + a.log.Warn("the control plane rejected every sample; dropping them", "count", n, "reply", statusErr.Body) return refused case code == http.StatusUnauthorized: - a.retryAt = time.Now().Add(unauthorizedWait) - a.log.Error("the control plane does not accept the token; keeping the data", "retry_in", unauthorizedWait) + b.retryAt = start.Add(unauthorizedWait) + a.log.Error("the control plane does not accept the token; keeping the data", "request", what, "retry_in", unauthorizedWait) return later case code == http.StatusTooManyRequests: wait := statusErr.RetryAfter if wait <= 0 { wait = throttledWait } - a.retryAt = time.Now().Add(wait) - a.log.Warn("the control plane asks the agent to wait; keeping the data", "retry_in", wait) + b.retryAt = start.Add(wait) + a.log.Warn("the control plane asks the agent to wait; keeping the data", "request", what, "retry_in", wait) return later } } // A 5xx, another status or a network error: keep the data and wait longer // after each failure. - a.failures++ - wait := min(firstBackoff< Date: Thu, 24 Sep 2026 13:55:14 +0600 Subject: [PATCH 3/3] feat(agent): log when the control plane accepts the requests again After the warnings of a failure (a 5xx or a network error, a 401 or a 429), the first request that succeeds logs one Info line with the time since the first failure. The events and the samples log apart. --- internal/agent/agent_test.go | 28 +++++++++++++++++ internal/agent/send.go | 27 ++++++++++++++--- internal/agent/send_test.go | 59 ++++++++++++++++++++++++++++++++++++ 3 files changed, 109 insertions(+), 5 deletions(-) diff --git a/internal/agent/agent_test.go b/internal/agent/agent_test.go index 81cb233..a3faf9e 100644 --- a/internal/agent/agent_test.go +++ b/internal/agent/agent_test.go @@ -30,6 +30,34 @@ func (r *recorder) Handle(_ context.Context, rec slog.Record) error { return nil } +// recordsOf returns the records with message msg. +func (r *recorder) recordsOf(msg string) []slog.Record { + r.mu.Lock() + defer r.mu.Unlock() + + var out []slog.Record + for _, rec := range r.records { + if rec.Message == msg { + out = append(out, rec) + } + } + return out +} + +// attr returns the value of the attribute key of a record, as a string, or +// nil. +func attr(rec slog.Record, key string) any { + var v any + rec.Attrs(func(a slog.Attr) bool { + if a.Key == key { + v = a.Value.String() + return false + } + return true + }) + return v +} + // times returns the times of the records with message msg. func (r *recorder) times(msg string) []time.Time { r.mu.Lock() diff --git a/internal/agent/send.go b/internal/agent/send.go index 86c81d8..063d4c3 100644 --- a/internal/agent/send.go +++ b/internal/agent/send.go @@ -46,9 +46,18 @@ const ( // waits on its own: a broken events route must not stop the metrics. type backoff struct { // retryAt is the earliest time of the next request, and failures is the - // number of failed requests in a row. - retryAt time.Time - failures int + // number of failed requests in a row. failingSince is the start of the + // first report whose request failed and kept its data, or zero. + retryAt time.Time + failures int + failingSince time.Time +} + +// failed records a failure that keeps the data. +func (b *backoff) failed(start time.Time) { + if b.failingSince.IsZero() { + b.failingSince = start + } } // waiting reports whether the request must wait at the time of the report. @@ -139,6 +148,11 @@ func (a *agent) sendSamples(ctx context.Context, start time.Time) bool { func (a *agent) outcome(ctx context.Context, b *backoff, start time.Time, err error, what string, n int) outcome { if err == nil { b.failures = 0 + // After the warnings of a failure, say that the data goes again. + if !b.failingSince.IsZero() { + a.log.Info("the control plane accepts the requests again", "request", what, "failing_for", start.Sub(b.failingSince).Round(time.Second)) + b.failingSince = time.Time{} + } return sent } @@ -155,15 +169,16 @@ func (a *agent) outcome(ctx context.Context, b *backoff, start time.Time, err er if errors.As(err, &statusErr) { switch code := statusErr.StatusCode; { case code == http.StatusBadRequest: - b.failures = 0 + b.failures, b.failingSince = 0, time.Time{} a.log.Error("the control plane refused the request as not valid; dropping its data", "request", what, "count", n, "reply", statusErr.Body) return refused case code == http.StatusUnprocessableEntity && what == "samples": - b.failures = 0 + b.failures, b.failingSince = 0, time.Time{} a.log.Warn("the control plane rejected every sample; dropping them", "count", n, "reply", statusErr.Body) return refused case code == http.StatusUnauthorized: b.retryAt = start.Add(unauthorizedWait) + b.failed(start) a.log.Error("the control plane does not accept the token; keeping the data", "request", what, "retry_in", unauthorizedWait) return later case code == http.StatusTooManyRequests: @@ -172,6 +187,7 @@ func (a *agent) outcome(ctx context.Context, b *backoff, start time.Time, err er wait = throttledWait } b.retryAt = start.Add(wait) + b.failed(start) a.log.Warn("the control plane asks the agent to wait; keeping the data", "request", what, "retry_in", wait) return later } @@ -182,6 +198,7 @@ func (a *agent) outcome(ctx context.Context, b *backoff, start time.Time, err er b.failures++ wait := min(firstBackoff<