From cf9d85b3cd867373c44af02eb0086caeb3783c61 Mon Sep 17 00:00:00 2001 From: mrhid6 Date: Fri, 7 Aug 2026 08:56:26 +0100 Subject: [PATCH] feat: store workload reports and route log results --- agent/internal/grpc/pb/workloads.pb.go | 3 + go.work.sum | 23 +-- proto/vantage/v1/vantage.proto | 3 + server/internal/grpc/pb/workloads.pb.go | 3 + server/internal/grpc/server.go | 71 ++++++++ server/internal/services/workloadresults.go | 128 ++++++++++++++ server/internal/services/workloads.go | 187 ++++++++++++++++++++ 7 files changed, 403 insertions(+), 15 deletions(-) create mode 100644 server/internal/services/workloadresults.go create mode 100644 server/internal/services/workloads.go diff --git a/agent/internal/grpc/pb/workloads.pb.go b/agent/internal/grpc/pb/workloads.pb.go index 3222e41..a3bc0da 100644 --- a/agent/internal/grpc/pb/workloads.pb.go +++ b/agent/internal/grpc/pb/workloads.pb.go @@ -38,6 +38,9 @@ type ReportWorkloadsRequest struct { SystemdOk bool `json:"systemd_ok"` SystemdError string `json:"systemd_error,omitempty"` Workloads []Workload `json:"workloads,omitempty"` // empty on the offer call + // Full marks the second call. It is not inferred from an empty Workloads + // slice: a host running nothing sends an empty list as its full report. + Full bool `json:"full,omitempty"` } type ReportWorkloadsResponse struct { diff --git a/go.work.sum b/go.work.sum index 5bdf3f4..cf66984 100644 --- a/go.work.sum +++ b/go.work.sum @@ -5,7 +5,6 @@ github.com/Intevation/jsonpath v0.2.1/go.mod h1:WnZ8weMmwAx/fAO3SutjYFU+v7DFreNY github.com/VividCortex/ewma v1.2.0/go.mod h1:nz4BbCtbLyFDeC9SUHbtcT5644juEuWfUAUnGx7j5l4= github.com/alecthomas/template v0.0.0-20190718012654-fb15b899a751/go.mod h1:LOuyumcjzFXgccqObfd/Ljyb9UuFJ6TxHnclSeseNhc= github.com/alecthomas/units v0.0.0-20211218093645-b94a6e3cc137/go.mod h1:OMCwj8VM1Kc9e19TLln2VL61YJF0x1XFtfdL4JdbSyE= -github.com/aquasecurity/bolt-fixtures v0.0.0-20200903104109-d34e7f983986/go.mod h1:NT+jyeCzXk6vXR5MTkdn4z64TgGfE5HMLC8qfj5unl8= github.com/aquasecurity/go-gem-version v0.0.0-20201115065557-8eed6fe000ce/go.mod h1:HXgVzOPvXhVGLJs4ZKO817idqr/xhwsTcj17CLYY74s= github.com/aquasecurity/go-npm-version v0.0.1/go.mod h1:hxbJZtKlO4P8sZ9nztizR6XLoE33O+BkPmuYQ4ACyz0= github.com/aquasecurity/go-pep440-version v0.0.1/go.mod h1:3naPe+Bp6wi3n4l5iBFCZgS0JG8vY6FT0H4NGhFJ+i4= @@ -19,9 +18,10 @@ github.com/cpuguy83/go-md2man/v2 v2.0.5/go.mod h1:tgQtvFlXSQOSOSIRvRPT7W67SCa46t github.com/envoyproxy/go-control-plane v0.12.0/go.mod h1:ZBTaoJ23lqITozF0M6G4/IragXCQKCnYbmlmtHvwRG0= github.com/envoyproxy/protoc-gen-validate v1.0.4/go.mod h1:qys6tmnRsYrQqIhm2bvKZH4Blx/1gTIZ2UKVY1M+Yew= github.com/fatih/color v1.18.0/go.mod h1:4FelSpRwEGDpQ12mAdzqdOukCy4u8WUtOY6lkT/6HfU= +github.com/go-logr/logr v1.4.2/go.mod h1:9T104GzyrTigFIr8wt5mBrctHMim0Nb2HLGrmQ40KvY= +github.com/go-logr/stdr v1.2.2/go.mod h1:mMo/vtBO5dYbehREoey6XUKy/eSumjCCveDpRre4VKE= github.com/go-openapi/jsonpointer v0.21.0/go.mod h1:IUyH9l/+uyhIYQ/PXVA41Rexl+kOkAPDdXEYns6fzUY= github.com/go-openapi/swag v0.23.0/go.mod h1:esZ8ITTYEsH1V2trKHjAN8Ai7xHb8RV+YSZ577vPjgQ= -github.com/goccy/go-yaml v1.19.0/go.mod h1:XBurs7gK8ATbW4ZPGKgcbrY1Br56PdM69F7LkFRi1kA= github.com/gocsaf/csaf/v3 v3.1.1/go.mod h1:EpUCrQg69i+Y66MphmQvVbcj333GFLjXOYHg1zoXVso= github.com/golang/glog v1.2.0/go.mod h1:6AhwSGph0fcJtXVM/PEHPqZlFeoLxhs7/t5UDAwmO+w= github.com/golang/protobuf v1.5.0/go.mod h1:FsONVRAS9T7sI+LIUmWTfcYkHO4aIWwzhcaSAoJOfIk= @@ -29,6 +29,7 @@ github.com/golang/protobuf v1.5.4/go.mod h1:lnTiLA8Wa4RWRcIUkrtSVa5nRhsEGBg48fD6 github.com/google/go-cmp v0.5.5/go.mod h1:v8dTdLbMG2kIc/vJvl+f65V22dbkXbowE6jgT/gNBxE= github.com/hashicorp/errwrap v1.0.0/go.mod h1:YH+1FKiLXxHSkmPseP+kNlulaMuP3n2brvKWEqk/Jc4= github.com/hashicorp/go-multierror v1.1.1/go.mod h1:iw975J/qwKPdAO1clOe2L8331t/9/fmwbPZ6JB6eMoM= +github.com/inconshreveable/mousetrap v1.1.0/go.mod h1:vpF70FUmC8bwa3OWnCshd2FqLfsEA9PFc4w1p2J65bw= github.com/josephburnett/jd/v2 v2.3.0/go.mod h1:0I5+gbo7y8diuajJjm79AF44eqTheSJy1K7DSbIUFAQ= github.com/josharian/intern v1.0.0/go.mod h1:5DoeVV0s6jJacbCEi61lwdGj/aVlrQvzHFFd8Hwg//Y= github.com/klauspost/cpuid/v2 v2.2.9/go.mod h1:rqkxqrZ1EhYM9G+hXH7YdowN5R5RGN6NK4QwQ3WMXF8= @@ -36,28 +37,20 @@ github.com/mailru/easyjson v0.7.7/go.mod h1:xzfreul335JAWq5oZzymOObrkdz5UnU4kGfJ github.com/masahiro331/go-mvn-version v0.0.0-20250131095131-f4974fa13b8a/go.mod h1:jZ3F25l7DbD7l7DcA8aj7eo1EZ84nbzcQHBB4lCSrI8= github.com/mattn/go-colorable v0.1.14/go.mod h1:6LmQG8QLFO4G5z1gPvYEzlUgJ2wF+stgPZH1UqBm1s8= github.com/mattn/go-runewidth v0.0.16/go.mod h1:Jdepj2loyihRzMpdS35Xk/zdY8IAYHsh153qUoGf23w= -github.com/oklog/ulid/v2 v2.1.1 h1:suPZ4ARWLOJLegGFiZZ1dFAkqzhMjL3J1TzI+5wHz8s= -github.com/oklog/ulid/v2 v2.1.1/go.mod h1:rcEKHmBBKfef9DhnvX7y1HZBYxjXb0cP5ExxNsTT1QQ= github.com/package-url/packageurl-go v0.1.3/go.mod h1:nKAWB8E6uk1MHqiS/lQb9pYBGH2+mdJ2PJc2s50dQY0= github.com/pandatix/go-cvss v0.6.2/go.mod h1:jDXYlQBZrc8nvrMUVVvTG8PhmuShOnKrxP53nOFkt8Q= github.com/rivo/uniseg v0.4.7/go.mod h1:FN3SvrM+Zdj16jyLfmOkMNblXMcoc8DfTHruCPUcx88= github.com/russross/blackfriday v1.6.0/go.mod h1:ti0ldHuxg49ri4ksnFxlkCfN+hvslNlmVHqNRXXJNAY= github.com/russross/blackfriday/v2 v2.1.0/go.mod h1:+Rmxgy9KzJVeS9/2gXHxylqXiyQDYRxCVz55jmeOWTM= -github.com/samber/lo v1.50.0 h1:XrG0xOeHs+4FQ8gJR97zDz5uOFMW7OwFWiFVzqopKgY= -github.com/samber/lo v1.50.0/go.mod h1:RjZyNk6WSnUFRKK6EyOhsRJMqft3G+pg7dCWHQCWvsc= -github.com/samber/oops v1.18.1 h1:qjhZbqbdyhWBKntkY8sxrDNKA8b4c5VHlmI1rli7X7M= -github.com/samber/oops v1.18.1/go.mod h1:xYqvimigkKV70HyLXiBZJFpIWi2CGcc6Xx7eV+2HycI= github.com/santhosh-tekuri/jsonschema/v5 v5.3.1/go.mod h1:uToXkOrWAZ6/Oc07xWQrPOhJotwFIyu2bBVN41fcDUY= github.com/shopspring/decimal v1.4.0/go.mod h1:gawqmDU56v4yIKSwfBSFip1HdCCXN8/+DMd9qYNcwME= -github.com/stretchr/objx v0.5.2 h1:xuMeJ0Sdp5ZMRXx/aWO6RZxdr3beISkG5/G/aIRr3pY= +github.com/spf13/cobra v1.8.1/go.mod h1:wHxEcudfqmLYa8iTfL+OuZPbBZkmvliBWKIezN3kD9Y= +github.com/spf13/pflag v1.0.6/go.mod h1:McXfInJRrz4CZXVZOBLb0bTZqETkiAhM9Iw0y3An2Bg= github.com/urfave/cli v1.22.16/go.mod h1:EeJR6BKodywf4zciqrdw6hpCPk68JO9z5LazXZMn5Po= -go.etcd.io/bbolt v1.4.3 h1:dEadXpI6G79deX5prL3QRNP6JB8UxVkqo4UPnHaNXJo= -go.etcd.io/bbolt v1.4.3/go.mod h1:tKQlpPaYCVFctUIgFKFnAlvbmB3tpy1vkTnDWohtc0E= +go.etcd.io/gofail v0.2.0/go.mod h1:nL3ILMGfkXTekKI3clMBNazKnjUZjYLKmBHzsVAnC1o= go.mongodb.org/mongo-driver/v2 v2.5.0/go.mod h1:yOI9kBsufol30iFsl1slpdq1I0eHPzybRWdyYUs8K/0= -go.opentelemetry.io/otel v1.34.0 h1:zRLXxLCgL1WyKsPVrgbSdMN4c0FMkDAskSTQP+0hdUY= -go.opentelemetry.io/otel v1.34.0/go.mod h1:OWFPOQ+h4G8xpyjgqo4SxJYdDQ/qmRH+wivy7zzx9oI= -go.opentelemetry.io/otel/trace v1.34.0 h1:+ouXS2V8Rd4hp4580a8q23bg0azF2nI8cqLYnC8mh/k= -go.opentelemetry.io/otel/trace v1.34.0/go.mod h1:Svm7lSjQD7kG7KJ/MUHPVXSDGz2OX4h0M2jHBhmSfRE= +go.opentelemetry.io/auto/sdk v1.1.0/go.mod h1:3wSPjt5PWp2RhlCcmmOial7AvC4DQqZb7a7wCow3W8A= +go.opentelemetry.io/otel/metric v1.34.0/go.mod h1:CEDrp0fy2D0MvkXE+dPV7cMi8tWZwX3dmaIhwPOaqHE= go.yaml.in/yaml/v4 v4.0.0-rc.3/go.mod h1:aZqd9kCMsGL7AuUv/m/PvWLdg5sjJsZ4oHDEnfPPfY0= golang.org/x/crypto v0.33.0/go.mod h1:bVdXmD7IV/4GdElGPozy6U7lWdRXA4qyRVGJV57uQ5M= golang.org/x/crypto v0.46.0/go.mod h1:Evb/oLKmMraqjZ2iQTwDwvCtJkczlDuTmdJXoZVzqU0= diff --git a/proto/vantage/v1/vantage.proto b/proto/vantage/v1/vantage.proto index b753dc2..b02fe75 100644 --- a/proto/vantage/v1/vantage.proto +++ b/proto/vantage/v1/vantage.proto @@ -341,6 +341,9 @@ message ReportWorkloadsRequest { bool systemd_ok = 6; string systemd_error = 7; repeated Workload workloads = 8; // empty on the offer call + // full marks the second call. It is not inferred from an empty workloads + // list: a host running nothing sends an empty list as its full report. + bool full = 9; } message ReportWorkloadsResponse { diff --git a/server/internal/grpc/pb/workloads.pb.go b/server/internal/grpc/pb/workloads.pb.go index 176030a..dc6e51d 100644 --- a/server/internal/grpc/pb/workloads.pb.go +++ b/server/internal/grpc/pb/workloads.pb.go @@ -40,6 +40,9 @@ type ReportWorkloadsRequest struct { SystemdOk bool `json:"systemd_ok"` SystemdError string `json:"systemd_error,omitempty"` Workloads []Workload `json:"workloads,omitempty"` // empty on the offer call + // Full marks the second call. It is not inferred from an empty Workloads + // slice: a host running nothing sends an empty list as its full report. + Full bool `json:"full,omitempty"` } type ReportWorkloadsResponse struct { diff --git a/server/internal/grpc/server.go b/server/internal/grpc/server.go index dbfe266..af60999 100644 --- a/server/internal/grpc/server.go +++ b/server/internal/grpc/server.go @@ -163,6 +163,70 @@ func (s *vantageServer) ReportPackages(ctx context.Context, req *pb.ReportPackag return &pb.ReportPackagesResponse{NeedFull: false}, nil } +// ReportWorkloads stores what a server is running. +// +// It is not gated by licence: the workload registry reads as core fleet +// management rather than a premium add-on. If that ever changes, the check +// belongs here — gating collection, not display — for the same reason it does +// in ReportPackages. +func (s *vantageServer) ReportWorkloads(ctx context.Context, req *pb.ReportWorkloadsRequest) (*pb.ReportWorkloadsResponse, error) { + srv, err := services.ValidateAgentToken(req.ServerId, req.AgentToken) + if err != nil { + return nil, status.Errorf(codes.Unauthenticated, "invalid agent token") + } + + // The offer call: a hash and no body. Answering NeedFull=false here is what + // keeps an unchanged 60-second report to one small message. + // + // The offer is identified by Full, not by an empty Workloads slice: a host + // genuinely running nothing sends an empty list as its FULL report, and + // inferring the offer from emptiness would leave that host answering + // NeedFull=true forever and never storing anything. + if !req.Full { + known, err := services.HasWorkloadHash(srv.InstanceID, srv.ServerID, req.Hash) + if err != nil { + log.Printf("workload hash lookup for %s: %v", srv.ServerID, err) + return nil, status.Errorf(codes.Internal, "workload hash lookup failed") + } + return &pb.ReportWorkloadsResponse{NeedFull: !known}, nil + } + + wls := make([]models.Workload, len(req.Workloads)) + for i, w := range req.Workloads { + wls[i] = models.Workload{ + Kind: w.Kind, + ID: w.Id, + Name: w.Name, + State: w.State, + Health: w.Health, + Image: w.Image, + Stack: w.Stack, + Ports: w.Ports, + Restarts: int(w.Restarts), + Protected: w.Protected, + } + if w.StartedAt != "" { + if t, err := time.Parse(time.RFC3339, w.StartedAt); err == nil { + wls[i].StartedAt = t + } + } + } + + if err := storeWorkloadReport(srv.InstanceID, srv.ServerID, req, wls); err != nil { + log.Printf("store workloads for %s: %v", srv.ServerID, err) + return nil, status.Errorf(codes.Internal, "failed to store workloads") + } + return &pb.ReportWorkloadsResponse{NeedFull: false}, nil +} + +func storeWorkloadReport(instanceID, serverID string, req *pb.ReportWorkloadsRequest, wls []models.Workload) error { + if wls == nil { + wls = []models.Workload{} + } + return services.StoreWorkloads(instanceID, serverID, req.Hash, wls, + req.DockerOk, req.DockerError, req.SystemdOk, req.SystemdError) +} + func (s *vantageServer) ReportInventory(ctx context.Context, req *pb.InventoryReport) (*pb.InventoryReportResponse, error) { srv, err := services.ValidateAgentToken(req.ServerId, req.AgentToken) if err != nil { @@ -257,6 +321,13 @@ func (s *vantageServer) CommandStream(stream pb.Vantage_CommandStreamServer) err if m.Result != nil { r := m.Result log.Printf("agent %s cmd %s: success=%v %s", srv.ServerID, r.CommandId, r.Success, r.Message) + // Republished so a control action waiting on another pod sees + // it. Publishing with no subscriber is a no-op, so this is safe + // for every command result rather than only the awaited ones. + services.WorkloadResults.DeliverCommand(r) + } + if m.WorkloadLogsResult != nil { + services.WorkloadResults.Deliver(m.WorkloadLogsResult) } if m.StepResult != nil { services.StepResults.Deliver(m.StepResult) diff --git a/server/internal/services/workloadresults.go b/server/internal/services/workloadresults.go new file mode 100644 index 0000000..f46de80 --- /dev/null +++ b/server/internal/services/workloadresults.go @@ -0,0 +1,128 @@ +package services + +import ( + "context" + "encoding/json" + "log" + + "gitea.hostxtra.co.uk/mrhid6/vantage/server/internal/bus" + "gitea.hostxtra.co.uk/mrhid6/vantage/server/internal/grpc/pb" +) + +// Workload results travel back over the bus for the same reason commands travel +// out over it: the pod serving the HTTP request and the pod holding the agent's +// stream are two different processes, and a map in one cannot be read by the +// other. +// +// Await MUST be called before the command is dispatched, or a fast agent +// answers into a channel nobody is listening on yet. See stepresults.go. + +type workloadResultRegistry struct{} + +var WorkloadResults = &workloadResultRegistry{} + +// Await subscribes to a command's result channel for a log snapshot. +func (r *workloadResultRegistry) Await(commandID string) (<-chan *pb.WorkloadLogsResult, func()) { + out := make(chan *pb.WorkloadLogsResult, 1) + + ctx, cancel := context.WithCancel(context.Background()) + raw, unsub, err := bus.Subscribe(ctx, bus.ResultChannel+commandID) + if err != nil { + log.Printf("workload results: subscribe for %s: %v", commandID, err) + cancel() + close(out) + return out, func() {} + } + + go func() { + defer close(out) + select { + case <-ctx.Done(): + return + case b, ok := <-raw: + if !ok { + return + } + var res pb.WorkloadLogsResult + if err := json.Unmarshal(b, &res); err != nil { + log.Printf("workload results: undecodable result for %s: %v", commandID, err) + return + } + out <- &res + } + }() + + return out, func() { + cancel() + unsub() + } +} + +// AwaitCommand subscribes to a command's result channel for a plain +// CommandResult, which is what a control action answers with. +// +// A control action reuses CommandResult rather than growing a message of its +// own: start, stop and restart succeed or fail, and that is exactly what +// CommandResult already says. +func (r *workloadResultRegistry) AwaitCommand(commandID string) (<-chan *pb.CommandResult, func()) { + out := make(chan *pb.CommandResult, 1) + + ctx, cancel := context.WithCancel(context.Background()) + raw, unsub, err := bus.Subscribe(ctx, bus.ResultChannel+commandID) + if err != nil { + log.Printf("workload results: subscribe for %s: %v", commandID, err) + cancel() + close(out) + return out, func() {} + } + + go func() { + defer close(out) + select { + case <-ctx.Done(): + return + case b, ok := <-raw: + if !ok { + return + } + var res pb.CommandResult + if err := json.Unmarshal(b, &res); err != nil { + log.Printf("workload results: undecodable command result for %s: %v", commandID, err) + return + } + out <- &res + } + }() + + return out, func() { + cancel() + unsub() + } +} + +// Deliver publishes a log result received from an agent. Called on the pod +// holding that agent's stream, which is not usually the pod waiting for it. +func (r *workloadResultRegistry) Deliver(res *pb.WorkloadLogsResult) { + if res == nil || res.CommandId == "" { + return + } + ctx, cancel := context.WithTimeout(context.Background(), dispatchAckTimeout) + defer cancel() + if _, err := bus.Publish(ctx, bus.ResultChannel+res.CommandId, res); err != nil { + log.Printf("workload results: publish for %s: %v", res.CommandId, err) + } +} + +// DeliverCommand republishes a CommandResult onto the bus so a waiting pod can +// see it. Publishing with no subscriber is a no-op, so this is safe to call for +// every CommandResult rather than only the ones somebody is waiting on. +func (r *workloadResultRegistry) DeliverCommand(res *pb.CommandResult) { + if res == nil || res.CommandId == "" { + return + } + ctx, cancel := context.WithTimeout(context.Background(), dispatchAckTimeout) + defer cancel() + if _, err := bus.Publish(ctx, bus.ResultChannel+res.CommandId, res); err != nil { + log.Printf("workload results: publish command result for %s: %v", res.CommandId, err) + } +} diff --git a/server/internal/services/workloads.go b/server/internal/services/workloads.go new file mode 100644 index 0000000..c87e359 --- /dev/null +++ b/server/internal/services/workloads.go @@ -0,0 +1,187 @@ +package services + +import ( + "context" + "fmt" + "time" + + "gitea.hostxtra.co.uk/mrhid6/vantage/server/internal/db" + "gitea.hostxtra.co.uk/mrhid6/vantage/server/internal/grpc/pb" + "gitea.hostxtra.co.uk/mrhid6/vantage/server/internal/models" + "github.com/google/uuid" + "go.mongodb.org/mongo-driver/v2/bson" + "go.mongodb.org/mongo-driver/v2/mongo" + "go.mongodb.org/mongo-driver/v2/mongo/options" +) + +// workloadResultTimeout bounds how long an API request waits for an agent to +// answer a control action or a log read. It is well above the agent's own +// 90-second control timeout and 60-second log timeout, so a slow-but-working +// agent reports its real error rather than being cut off by this side. +const workloadResultTimeout = 120 * time.Second + +func HasWorkloadHash(instanceID, serverID, hash string) (bool, error) { + err := db.Col("server_workloads").FindOne(context.Background(), bson.M{ + "instance_id": instanceID, + "server_id": serverID, + "hash": hash, + }, options.FindOne().SetProjection(bson.M{"_id": 1})).Err() + if err == mongo.ErrNoDocuments { + return false, nil + } + return err == nil, err +} + +// StoreWorkloads replaces a server's workload list. +func StoreWorkloads(instanceID, serverID, hash string, wls []models.Workload, + dockerOK bool, dockerErr string, systemdOK bool, systemdErr string) error { + + _, err := db.Col("server_workloads").UpdateOne(context.Background(), + bson.M{"instance_id": instanceID, "server_id": serverID}, + bson.M{"$set": bson.M{ + "hash": hash, + "workloads": wls, + "collected_at": time.Now(), + "docker_ok": dockerOK, + "docker_error": dockerErr, + "systemd_ok": systemdOK, + "systemd_error": systemdErr, + }}, + options.UpdateOne().SetUpsert(true), + ) + return err +} + +func GetWorkloads(instanceID, serverID string) (*models.ServerWorkloads, error) { + var sw models.ServerWorkloads + err := db.Col("server_workloads").FindOne(context.Background(), bson.M{ + "instance_id": instanceID, + "server_id": serverID, + }).Decode(&sw) + if err == mongo.ErrNoDocuments { + return nil, nil + } + if err != nil { + return nil, err + } + return &sw, nil +} + +type WorkloadHit struct { + ServerID string `json:"server_id"` + Workload models.Workload `json:"workload"` +} + +// SearchWorkloads answers "which servers run image X" — the reason the snapshot +// is stored rather than fetched on demand and discarded. +func SearchWorkloads(instanceID, image, stack, state string) ([]WorkloadHit, error) { + ctx := context.Background() + + filter := bson.M{"instance_id": instanceID} + if image != "" { + filter["workloads.image"] = image + } + + cur, err := db.Col("server_workloads").Find(ctx, filter) + if err != nil { + return nil, err + } + defer cur.Close(ctx) + + var docs []models.ServerWorkloads + if err := cur.All(ctx, &docs); err != nil { + return nil, err + } + + hits := []WorkloadHit{} + for _, d := range docs { + for _, w := range d.Workloads { + if image != "" && w.Image != image { + continue + } + if stack != "" && w.Stack != stack { + continue + } + if state != "" && w.State != state { + continue + } + hits = append(hits, WorkloadHit{ServerID: d.ServerID, Workload: w}) + } + } + return hits, nil +} + +// DispatchRefreshWorkloads asks an agent to report immediately. It returns as +// soon as the owning pod acks; the caller refetches the stored document. +// +// The refresh carries nothing back on purpose: the agent answers through the +// normal ReportWorkloads RPC, so server_workloads has exactly one writer. +func DispatchRefreshWorkloads(serverID string) error { + return Dispatcher.dispatch(serverID, &pb.ServerCommand{ + CommandId: uuid.New().String(), + RefreshWorkloads: &pb.RefreshWorkloadsCmd{}, + }) +} + +// DispatchControlWorkload runs a control action and waits for the agent's +// CommandResult. +// +// Await is called BEFORE dispatch. Reversing those two lines introduces a race +// that only shows under load, on a fast agent answering into a channel nobody +// has joined yet. +func DispatchControlWorkload(serverID, kind, id, action string) error { + commandID := uuid.New().String() + + results, done := WorkloadResults.AwaitCommand(commandID) + defer done() + + if err := Dispatcher.dispatch(serverID, &pb.ServerCommand{ + CommandId: commandID, + ControlWorkload: &pb.ControlWorkloadCmd{Kind: kind, Id: id, Action: action}, + }); err != nil { + return err + } + + select { + case res, ok := <-results: + if !ok || res == nil { + return fmt.Errorf("no result from agent for %s %s", action, id) + } + if !res.Success { + return fmt.Errorf("%s", res.Message) + } + return nil + case <-time.After(workloadResultTimeout): + return fmt.Errorf("timed out waiting for the agent to %s %s", action, id) + } +} + +// DispatchWorkloadLogs fetches a bounded log snapshot. +// +// Await is called BEFORE dispatch, for the same reason as above. +func DispatchWorkloadLogs(serverID, kind, id string, tail int) (string, bool, error) { + commandID := uuid.New().String() + + results, done := WorkloadResults.Await(commandID) + defer done() + + if err := Dispatcher.dispatch(serverID, &pb.ServerCommand{ + CommandId: commandID, + WorkloadLogs: &pb.WorkloadLogsCmd{Kind: kind, Id: id, Tail: int32(tail)}, + }); err != nil { + return "", false, err + } + + select { + case res, ok := <-results: + if !ok || res == nil { + return "", false, fmt.Errorf("no log result from agent for %s", id) + } + if res.Error != "" { + return "", false, fmt.Errorf("%s", res.Error) + } + return res.Text, res.Truncated, nil + case <-time.After(workloadResultTimeout): + return "", false, fmt.Errorf("timed out waiting for logs for %s", id) + } +}