From a8718ab4f96967e5b8fccc81e720edd39088cae5 Mon Sep 17 00:00:00 2001 From: "ccf-lisa[bot]" <286799724+ccf-lisa[bot]@users.noreply.github.com> Date: Mon, 5 Oct 2026 14:08:19 -0300 Subject: [PATCH 1/5] feat(agentcfg): validation bases, instance pruning and agent deletion Ninth layer of the agent remote-configuration stack (split from #465): ValidationBases (the instances a save validates against, R48) and the bounded PreviewBases set, the River job that prunes stale instances (CCF_AGENT_INSTANCE_PRUNE_*), and deleting an agent now deletes its instances and config revisions. Co-Authored-By: Claude Opus 5.5 --- internal/api/handler/agents.go | 9 + .../service/relational/agentcfg/service.go | 208 ++++++++++++++++ .../agentcfg/service_integration_test.go | 224 ++++++++++++++++++ .../service/worker/agent_instance_prune.go | 74 ++++++ .../agent_instance_prune_integration_test.go | 94 ++++++++ .../worker/agent_instance_prune_test.go | 31 +++ internal/service/worker/service.go | 9 + 7 files changed, 649 insertions(+) create mode 100644 internal/service/worker/agent_instance_prune.go create mode 100644 internal/service/worker/agent_instance_prune_integration_test.go create mode 100644 internal/service/worker/agent_instance_prune_test.go diff --git a/internal/api/handler/agents.go b/internal/api/handler/agents.go index ee07b492..9c33899f 100644 --- a/internal/api/handler/agents.go +++ b/internal/api/handler/agents.go @@ -10,6 +10,7 @@ import ( "github.com/compliance-framework/api/internal/api" "github.com/compliance-framework/api/internal/service/relational" + "github.com/compliance-framework/api/internal/service/relational/agentcfg" "github.com/google/uuid" "github.com/labstack/echo/v4" "go.uber.org/zap" @@ -205,6 +206,14 @@ func (h *AgentHandler) DeleteAgent(ctx echo.Context) error { return err } + if err := agentcfg.DeleteInstancesForAgent(tx, *agent.ID); err != nil { + return err + } + + if err := agentcfg.DeleteRevisionsForAgent(tx, *agent.ID); err != nil { + return err + } + if err := tx.Delete(agent).Error; err != nil { return err } diff --git a/internal/service/relational/agentcfg/service.go b/internal/service/relational/agentcfg/service.go index 8c011c5a..08075d04 100644 --- a/internal/service/relational/agentcfg/service.go +++ b/internal/service/relational/agentcfg/service.go @@ -625,6 +625,203 @@ func (s *Service) GetInstance(ctx context.Context, agentID, instanceID uuid.UUID return &out, nil } +// baseColumns are the instance columns validation and preview need (no report payloads +// besides the base), so loading every instance of a large fleet stays cheap. +var baseColumns = []string{ + "id", "agent_id", "instance_id", "hostname", "mode", "daemon", "first_seen_at", "last_seen_at", + "reported_at", "applied_revision", "attempted_revision", "reported_status", + "base_config", "remote_config", +} + +// InstanceBase is a reported base an overlay is validated or previewed against. +type InstanceBase struct { + Instance relational.AgentInstance + Base agentconfig.Config + Remote agentconfig.RemoteConfig // reported remote-config, or the base's block normalized with hasAuth=true + Stale bool + Validated bool // member of ValidationBases (R48) +} + +var applyModes = []string{agentconfig.ModeApplySafe, agentconfig.ModeApplyAll} + +// ValidationBases is exactly the set PUT and revert validate against (R14, R48): +// 1. all fresh instances (seen within InstanceStaleAfter) with a reported base and an +// apply mode; +// 2. else the single most recently reported instance with a base and an apply mode, +// whatever its age; +// 3. else none, and standalone=true (overlay-level checks only). +// +// Report-mode instances and instances without a base are never validated against. A base +// that no longer decodes is skipped with a warning. +func (s *Service) ValidationBases(ctx context.Context, agentID uuid.UUID) ([]InstanceBase, bool, error) { + now := s.now() + var fresh []relational.AgentInstance + err := s.db.WithContext(ctx). + Select(baseColumns). + Where("agent_id = ? AND base_config IS NOT NULL AND mode IN ? AND last_seen_at >= ?", agentID, applyModes, now.Add(-s.settings.InstanceStaleAfter)). + Order("last_seen_at DESC, instance_id"). + Find(&fresh).Error + if err != nil { + return nil, false, err + } + rows := fresh + if len(rows) == 0 { + var latest []relational.AgentInstance + err := s.db.WithContext(ctx). + Select(baseColumns). + Where("agent_id = ? AND base_config IS NOT NULL AND mode IN ? AND reported_at IS NOT NULL", agentID, applyModes). + Order("reported_at DESC, instance_id"). + Limit(1). + Find(&latest).Error + if err != nil { + return nil, false, err + } + rows = latest + } + bases := s.toBases(rows, now, true) + return bases, len(bases) == 0, nil +} + +// Preview bounds (R14): a preview shows at most PreviewMaxInstances instances and decodes +// at most PreviewMaxConfigBytes of reported base+effective config, so one agent credential +// cannot make a single preview cost minutes of CPU by reporting many large instances. +const ( + PreviewMaxInstances = 50 + PreviewMaxConfigBytes = 16 << 20 +) + +// PreviewSet is what a preview works on. +type PreviewSet struct { + // Validation is ValidationBases: the set a save validates against (R48), in full. + Validation []InstanceBase + // Instances are the instances the preview shows, each marked Validated when it is in + // Validation: the validated ones first, then the others, newest first, within + // PreviewMaxInstances and PreviewMaxConfigBytes. + Instances []InstanceBase + // Omitted counts the instances with a reported base the bounds left out. + Omitted int64 +} + +// PreviewBases returns the validation set and the bounded list of instances with a +// reported base (fresh and stale, flagged) a preview shows. Only the selected instances' +// configs are loaded. +func (s *Service) PreviewBases(ctx context.Context, agentID uuid.UUID) (PreviewSet, error) { + validation, _, err := s.ValidationBases(ctx, agentID) + if err != nil { + return PreviewSet{}, err + } + validated := map[uuid.UUID]bool{} + for _, b := range validation { + validated[b.Instance.InstanceID] = true + } + + type candidate struct { + InstanceID uuid.UUID + ConfigSize int64 + } + var candidates []candidate + if err := s.db.WithContext(ctx). + Model(&relational.AgentInstance{}). + Select("instance_id, octet_length(base_config::text) + COALESCE(octet_length(effective_config::text), 0) AS config_size"). + Where("agent_id = ? AND base_config IS NOT NULL", agentID). + Order("last_seen_at DESC, instance_id"). + Scan(&candidates).Error; err != nil { + return PreviewSet{}, err + } + // Validated instances first (their errors block a save), keeping newest-first order. + slices.SortStableFunc(candidates, func(a, b candidate) int { + switch { + case validated[a.InstanceID] == validated[b.InstanceID]: + return 0 + case validated[a.InstanceID]: + return -1 + default: + return 1 + } + }) + var ( + picked []uuid.UUID + budget int64 = PreviewMaxConfigBytes + ) + for _, c := range candidates { + if len(picked) == PreviewMaxInstances || (len(picked) > 0 && c.ConfigSize > budget) { + break + } + picked = append(picked, c.InstanceID) + budget -= c.ConfigSize + } + set := PreviewSet{Validation: validation, Omitted: int64(len(candidates) - len(picked))} + if len(picked) == 0 { + return set, nil + } + + var rows []relational.AgentInstance + if err := s.db.WithContext(ctx). + Select(append(append([]string{}, baseColumns...), "effective_config")). + Where("agent_id = ? AND instance_id IN ?", agentID, picked). + Find(&rows).Error; err != nil { + return PreviewSet{}, err + } + order := make(map[uuid.UUID]int, len(picked)) + for i, id := range picked { + order[id] = i + } + slices.SortFunc(rows, func(a, b relational.AgentInstance) int { + return order[a.InstanceID] - order[b.InstanceID] + }) + set.Instances = s.toBases(rows, s.now(), false) + for i := range set.Instances { + set.Instances[i].Validated = validated[set.Instances[i].Instance.InstanceID] + } + return set, nil +} + +func (s *Service) toBases(rows []relational.AgentInstance, now time.Time, validated bool) []InstanceBase { + out := make([]InstanceBase, 0, len(rows)) + for _, row := range rows { + base, err := agentconfig.DecodeConfig(row.BaseConfig) + if err != nil { + s.logger.Warnw("Skipping agent instance with an undecodable reported base", + "agentID", row.AgentID, "instanceID", row.InstanceID, "error", err) + continue + } + out = append(out, InstanceBase{ + Instance: row, + Base: base, + Remote: reportedRemote(row, base), + Stale: IsStale(row, now, s.settings), + Validated: validated, + }) + } + return out +} + +// reportedRemote returns the instance's reported remote-config block, or the base's block +// normalized with hasAuth=true (the base is redacted, so it has no client secret). +func reportedRemote(row relational.AgentInstance, base agentconfig.Config) agentconfig.RemoteConfig { + var rc agentconfig.RemoteConfig + if len(row.RemoteConfig) > 0 && string(row.RemoteConfig) != "null" { + if err := json.Unmarshal(row.RemoteConfig, &rc); err == nil { + if rc.Mode == "" { + rc.Mode = row.Mode + } + return rc.Normalize(true) + } + } + if base.RemoteConfig != nil { + rc = *base.RemoteConfig + } + if rc.Mode == "" { + rc.Mode = row.Mode + } + return rc.Normalize(true) +} + +// DeleteInstancesForAgent removes every instance of an agent (agent deletion). +func DeleteInstancesForAgent(tx *gorm.DB, agentID uuid.UUID) error { + return tx.Where("agent_id = ?", agentID).Delete(&relational.AgentInstance{}).Error +} + // DeleteRevisionsForAgent removes every configuration revision of an agent (agent deletion). // It is the purge path for an overlay that held a secret, so it bypasses the append-only // BeforeDelete hook, which still blocks every other delete of a revision. @@ -634,6 +831,17 @@ func DeleteRevisionsForAgent(tx *gorm.DB, agentID uuid.UUID) error { Delete(&relational.AgentConfigRevision{}).Error } +// PruneInstances deletes one-shot instances (daemon=false) not seen for +// OneShotInstanceRetention and any instance not seen for InstanceRetention (R37). +func PruneInstances(ctx context.Context, db *gorm.DB, s Settings, now time.Time) (int64, error) { + s = s.WithDefaults() + res := db.WithContext(ctx). + Where("(daemon = false AND last_seen_at < ?) OR last_seen_at < ?", + now.Add(-s.OneShotInstanceRetention), now.Add(-s.InstanceRetention)). + Delete(&relational.AgentInstance{}) + return res.RowsAffected, res.Error +} + // ---- Derived state ---- // IsStale reports whether an instance was last seen before now - InstanceStaleAfter. diff --git a/internal/service/relational/agentcfg/service_integration_test.go b/internal/service/relational/agentcfg/service_integration_test.go index f0cceba1..784aa290 100644 --- a/internal/service/relational/agentcfg/service_integration_test.go +++ b/internal/service/relational/agentcfg/service_integration_test.go @@ -745,6 +745,198 @@ func (s *AgentCfgServiceIntegrationSuite) TestListInstancesPagesAndCountsAll() { s.Equal(agentcfg.InstanceCounts{}, counts) } +// validationFixture seeds an agent with fresh apply-mode instances plus instances that must +// never be validated against. +type validationFixture struct { + agentID uuid.UUID + freshSafe, freshAll, reportMode, stale uuid.UUID + noBase uuid.UUID +} + +func (s *AgentCfgServiceIntegrationSuite) seedValidationFixture() validationFixture { + f := validationFixture{ + agentID: s.newAgent("validation"), freshSafe: uuid.New(), freshAll: uuid.New(), + reportMode: uuid.New(), stale: uuid.New(), noBase: uuid.New(), + } + safe := applyReport(agentconfig.ModeApplySafe, baseConfig) + safe.RemoteConfig = &agentconfig.RemoteConfig{Mode: agentconfig.ModeApplySafe, TrustedSources: []string{"ghcr.io/x/*"}} + s.Require().NoError(s.reportAt(s.svc, s.now, f.agentID, f.freshSafe, safe)) + + // No reported remote-config: falls back to the base block. + withRemote := `{"daemon":true,"verbosity":0,"remote_config":{"mode":"apply_all","poll_interval":"30s"},"plugins":{}}` + s.Require().NoError(s.reportAt(s.svc, s.now.Add(-5*time.Minute), f.agentID, f.freshAll, applyReport(agentconfig.ModeApplyAll, withRemote))) + + s.Require().NoError(s.reportAt(s.svc, s.now.Add(-time.Minute), f.agentID, f.reportMode, applyReport(agentconfig.ModeReport, baseConfig))) + s.Require().NoError(s.reportAt(s.svc, s.now.Add(-time.Minute), f.agentID, f.noBase, applyReport(agentconfig.ModeApplySafe, ""))) + s.Require().NoError(s.reportAt(s.svc, s.now.Add(-time.Hour), f.agentID, f.stale, applyReport(agentconfig.ModeApplySafe, baseConfig))) + return f +} + +func byInstance(bases []agentcfg.InstanceBase) map[uuid.UUID]agentcfg.InstanceBase { + out := map[uuid.UUID]agentcfg.InstanceBase{} + for _, b := range bases { + out[b.Instance.InstanceID] = b + } + return out +} + +func (s *AgentCfgServiceIntegrationSuite) TestValidationBasesFreshInstances() { + f := s.seedValidationFixture() + + bases, standalone, err := s.svc.ValidationBases(s.ctx, f.agentID) + s.Require().NoError(err) + s.False(standalone) + s.Require().Len(bases, 2) + s.Equal(f.freshSafe, bases[0].Instance.InstanceID, "most recently seen first") + s.Equal(f.freshAll, bases[1].Instance.InstanceID) + + m := byInstance(bases) + for _, b := range bases { + s.False(b.Stale) + s.True(b.Validated) + s.True(b.Base.Daemon, "base decoded") + } + s.Require().NotNil(m[f.freshSafe].Base.API) + s.Equal("http://api:8080", m[f.freshSafe].Base.API.URL) + + // Reported remote-config wins. + safeRemote := m[f.freshSafe].Remote + s.Equal(agentconfig.ModeApplySafe, safeRemote.Mode) + s.Equal([]string{"ghcr.io/x/*"}, safeRemote.TrustedSources) + s.Equal("60s", safeRemote.PollInterval) + + // Falls back to the base's remote_config block, normalized with hasAuth=true. + allRemote := m[f.freshAll].Remote + s.Equal(agentconfig.ModeApplyAll, allRemote.Mode) + s.Equal("30s", allRemote.PollInterval) + s.Equal([]string{}, allRemote.TrustedSources) + s.Equal([]string{}, allRemote.OverridableConfigFlags) + s.False(allRemote.AllowLocalSources) +} + +func (s *AgentCfgServiceIntegrationSuite) TestValidationBasesFallbackAndStandalone() { + agentID := s.newAgent("fallback") + seenRecently, reportedRecently := uuid.New(), uuid.New() + // Reported 200h ago but heartbeated 30m ago (stale, but seen more recently). + s.Require().NoError(s.reportAt(s.svc, s.now.Add(-200*time.Hour), agentID, seenRecently, applyReport(agentconfig.ModeApplySafe, baseConfig))) + saved := s.now + s.now = saved.Add(-30 * time.Minute) + s.Require().NoError(s.svc.TouchFromHeartbeat(s.ctx, agentID, nil, seenRecently, nil, nil)) + s.now = saved + // Most recently reported apply-mode instance, 100h ago. + s.Require().NoError(s.reportAt(s.svc, s.now.Add(-100*time.Hour), agentID, reportedRecently, applyReport(agentconfig.ModeApplyAll, baseConfig))) + // Newer, but report mode (and a fresh report-mode one) or without a base. + s.Require().NoError(s.reportAt(s.svc, s.now.Add(-2*time.Hour), agentID, uuid.New(), applyReport(agentconfig.ModeReport, baseConfig))) + s.Require().NoError(s.reportAt(s.svc, s.now, agentID, uuid.New(), applyReport(agentconfig.ModeReport, baseConfig))) + s.Require().NoError(s.reportAt(s.svc, s.now.Add(-time.Hour), agentID, uuid.New(), applyReport(agentconfig.ModeApplySafe, ""))) + + bases, standalone, err := s.svc.ValidationBases(s.ctx, agentID) + s.Require().NoError(err) + s.False(standalone) + s.Require().Len(bases, 1) + s.Equal(reportedRecently, bases[0].Instance.InstanceID) + s.True(bases[0].Stale) + s.True(bases[0].Validated) + s.Equal(agentconfig.ModeApplyAll, bases[0].Remote.Mode, "no remote block anywhere: the instance mode") + + // Only report-mode / base-less instances: standalone. + onlyIneligible := s.newAgent("standalone") + s.Require().NoError(s.reportAt(s.svc, s.now, onlyIneligible, uuid.New(), applyReport(agentconfig.ModeReport, baseConfig))) + s.Require().NoError(s.reportAt(s.svc, s.now, onlyIneligible, uuid.New(), applyReport(agentconfig.ModeApplyAll, ""))) + s.Require().NoError(s.svc.TouchFromHeartbeat(s.ctx, onlyIneligible, nil, uuid.New(), ptr(int64(1)), ptr("d"))) + bases, standalone, err = s.svc.ValidationBases(s.ctx, onlyIneligible) + s.Require().NoError(err) + s.True(standalone) + s.Empty(bases) + + bases, standalone, err = s.svc.ValidationBases(s.ctx, s.newAgent("no-instances")) + s.Require().NoError(err) + s.True(standalone) + s.Empty(bases) +} + +func (s *AgentCfgServiceIntegrationSuite) TestPreviewBases() { + f := s.seedValidationFixture() + + set, err := s.svc.PreviewBases(s.ctx, f.agentID) + s.Require().NoError(err) + s.Zero(set.Omitted) + s.Len(set.Validation, 2) + m := byInstance(set.Instances) + s.Require().Len(m, 4, "every instance with a base; base-less excluded") + s.NotContains(m, f.noBase) + + s.True(m[f.freshSafe].Validated) + s.True(m[f.freshAll].Validated) + s.False(m[f.reportMode].Validated) + s.False(m[f.stale].Validated) + s.False(m[f.freshSafe].Stale) + s.False(m[f.reportMode].Stale) + s.True(m[f.stale].Stale) + s.Equal(agentconfig.ModeReport, m[f.reportMode].Remote.Mode) + + // Fallback: only the most recently reported stale instance is validated. + agentID := s.newAgent("preview-fallback") + older, newer := uuid.New(), uuid.New() + s.Require().NoError(s.reportAt(s.svc, s.now.Add(-3*time.Hour), agentID, older, applyReport(agentconfig.ModeApplySafe, baseConfig))) + s.Require().NoError(s.reportAt(s.svc, s.now.Add(-2*time.Hour), agentID, newer, applyReport(agentconfig.ModeApplySafe, baseConfig))) + set, err = s.svc.PreviewBases(s.ctx, agentID) + s.Require().NoError(err) + bases := set.Instances + s.Require().Len(bases, 2) + s.Equal(newer, bases[0].Instance.InstanceID) + s.True(bases[0].Stale) + s.True(bases[0].Validated) + s.True(bases[1].Stale) + s.False(bases[1].Validated) +} + +func (s *AgentCfgServiceIntegrationSuite) TestPreviewBasesBounded() { + agentID := s.newAgent("preview-bounded") + // The validated instance is older than every report-mode one, which are not validated + // and outnumber the preview bound. + validated := uuid.New() + s.Require().NoError(s.reportAt(s.svc, s.now.Add(-5*time.Minute), agentID, validated, applyReport(agentconfig.ModeApplySafe, baseConfig))) + var others []uuid.UUID + for i := range agentcfg.PreviewMaxInstances + 5 { + id := uuid.New() + others = append(others, id) + s.Require().NoError(s.reportAt(s.svc, s.now.Add(-time.Duration(i+1)*time.Second), agentID, id, applyReport(agentconfig.ModeReport, baseConfig))) + } + + set, err := s.svc.PreviewBases(s.ctx, agentID) + s.Require().NoError(err) + s.Len(set.Validation, 1) + s.Require().Len(set.Instances, agentcfg.PreviewMaxInstances) + s.EqualValues(6, set.Omitted) + s.Equal(validated, set.Instances[0].Instance.InstanceID, "validated instances first") + s.True(set.Instances[0].Validated) + s.Equal(others[0], set.Instances[1].Instance.InstanceID, "then newest first") + s.False(set.Instances[1].Validated) +} + +func (s *AgentCfgServiceIntegrationSuite) TestDeleteInstancesForAgentKeepsRevisions() { + agentA := s.newAgent("delete-a") + agentB := s.newAgent("delete-b") + s.createRevision(agentA, 0, `{"verbosity":1}`) + s.createRevision(agentA, 1, `{"verbosity":2}`) + report := applyReport(agentconfig.ModeApplySafe, baseConfig) + s.Require().NoError(s.svc.UpsertReport(s.ctx, agentA, nil, uuid.New(), report)) + s.Require().NoError(s.svc.UpsertReport(s.ctx, agentA, nil, uuid.New(), report)) + keep := uuid.New() + s.Require().NoError(s.svc.UpsertReport(s.ctx, agentB, nil, keep, report)) + + s.Require().NoError(agentcfg.DeleteInstancesForAgent(s.DB, agentA)) + + s.Equal(int64(0), s.countInstances(agentA)) + s.Equal(int64(1), s.countInstances(agentB)) + _, err := s.svc.GetInstance(s.ctx, agentB, keep) + s.NoError(err) + n, err := s.svc.CurrentRevisionNumber(s.ctx, agentA) + s.Require().NoError(err) + s.Equal(int64(2), n, "revisions are kept") +} + func (s *AgentCfgServiceIntegrationSuite) TestDeleteRevisionsForAgent() { agentA := s.newAgent("purge-a") agentB := s.newAgent("purge-b") @@ -781,3 +973,35 @@ func (s *AgentCfgServiceIntegrationSuite) TestCreateRevisionReturnsTheStoredRow( s.Require().Len(metas, 1) s.Equal(len(rev.Overlay), metas[0].OverlaySize) } + +func (s *AgentCfgServiceIntegrationSuite) TestPruneInstances() { + agentID := s.newAgent("prune") + now := s.now + oneShotOld := s.insertInstance(agentID, ptr(false), now.Add(-25*time.Hour)) + oneShotRecent := s.insertInstance(agentID, ptr(false), now.Add(-23*time.Hour)) + daemon25h := s.insertInstance(agentID, ptr(true), now.Add(-25*time.Hour)) + daemonOld := s.insertInstance(agentID, ptr(true), now.Add(-721*time.Hour)) + null25h := s.insertInstance(agentID, nil, now.Add(-25*time.Hour)) + nullOld := s.insertInstance(agentID, nil, now.Add(-721*time.Hour)) + + deleted, err := agentcfg.PruneInstances(s.ctx, s.DB, agentcfg.Settings{}, now) + s.Require().NoError(err) + s.Equal(int64(3), deleted) + + remaining := map[uuid.UUID]bool{} + list, err := s.svc.ListInstances(s.ctx, agentID) + s.Require().NoError(err) + for _, i := range list { + remaining[i.InstanceID] = true + } + s.False(remaining[oneShotOld], "daemon=false pruned after 24h") + s.True(remaining[oneShotRecent]) + s.True(remaining[daemon25h], "daemon=true kept at 25h") + s.False(remaining[daemonOld], "daemon=true pruned after 720h") + s.True(remaining[null25h], "daemon NULL treated like a daemon") + s.False(remaining[nullOld]) + + deleted, err = agentcfg.PruneInstances(s.ctx, s.DB, agentcfg.Settings{}, now) + s.Require().NoError(err) + s.Equal(int64(0), deleted, "idempotent") +} diff --git a/internal/service/worker/agent_instance_prune.go b/internal/service/worker/agent_instance_prune.go new file mode 100644 index 00000000..b4178f23 --- /dev/null +++ b/internal/service/worker/agent_instance_prune.go @@ -0,0 +1,74 @@ +package worker + +import ( + "context" + "fmt" + "time" + + "github.com/compliance-framework/api/internal/service/relational/agentcfg" + "github.com/riverqueue/river" + "go.uber.org/zap" + "gorm.io/gorm" +) + +// JobTypeAgentInstancePrune prunes agent instances that are no longer reporting (R37). +const JobTypeAgentInstancePrune = "agent_instance_prune" + +// defaultAgentInstancePruneSchedule is hourly, River 6-field (seconds first). +const defaultAgentInstancePruneSchedule = "0 17 * * * *" + +// AgentInstancePruneArgs are the (empty) args of the periodic prune job. +type AgentInstancePruneArgs struct{} + +// Kind implements river.JobArgs. +func (AgentInstancePruneArgs) Kind() string { return JobTypeAgentInstancePrune } + +// AgentInstancePruneWorker deletes one-shot instances (daemon=false) not seen for the +// one-shot retention (24h) and any instance not seen for the retention (720h). +type AgentInstancePruneWorker struct { + db *gorm.DB + settings agentcfg.Settings + logger *zap.SugaredLogger + now func() time.Time +} + +// NewAgentInstancePruneWorker builds the worker. +func NewAgentInstancePruneWorker(db *gorm.DB, settings agentcfg.Settings, logger *zap.SugaredLogger) *AgentInstancePruneWorker { + if logger == nil { + logger = zap.NewNop().Sugar() + } + return &AgentInstancePruneWorker{db: db, settings: settings.WithDefaults(), logger: logger, now: func() time.Time { return time.Now().UTC() }} +} + +// Work implements the River worker function. +func (w *AgentInstancePruneWorker) Work(ctx context.Context, _ *river.Job[AgentInstancePruneArgs]) error { + deleted, err := agentcfg.PruneInstances(ctx, w.db, w.settings, w.now()) + if err != nil { + return fmt.Errorf("agent instance prune: %w", err) + } + if deleted > 0 { + w.logger.Infow("Pruned agent instances", "deleted", deleted) + } + return nil +} + +// NewAgentInstancePrunePeriodicJob schedules the prune job on the "scheduler" queue. +func NewAgentInstancePrunePeriodicJob(schedule string, logger *zap.SugaredLogger) *river.PeriodicJob { + sched := parseCronScheduleWithFallback(schedule, defaultAgentInstancePruneSchedule, "agent instance prune", logger) + return river.NewPeriodicJob( + sched, + func() (river.JobArgs, *river.InsertOpts) { + return &AgentInstancePruneArgs{}, &river.InsertOpts{ + Queue: "scheduler", + MaxAttempts: 3, + // Deduplicated per hour: the job runs at most hourly, so a + // CCF_AGENT_INSTANCE_PRUNE_SCHEDULE more frequent than hourly is not honored. + UniqueOpts: river.UniqueOpts{ + ByArgs: true, + ByPeriod: time.Hour, + }, + } + }, + &river.PeriodicJobOpts{RunOnStart: false}, + ) +} diff --git a/internal/service/worker/agent_instance_prune_integration_test.go b/internal/service/worker/agent_instance_prune_integration_test.go new file mode 100644 index 00000000..289ec83e --- /dev/null +++ b/internal/service/worker/agent_instance_prune_integration_test.go @@ -0,0 +1,94 @@ +//go:build integration + +package worker + +import ( + "context" + "testing" + "time" + + "github.com/compliance-framework/api/internal/service/relational" + "github.com/compliance-framework/api/internal/service/relational/agentcfg" + "github.com/compliance-framework/api/internal/tests" + "github.com/google/uuid" + "github.com/riverqueue/river" + "github.com/stretchr/testify/suite" +) + +type AgentInstancePruneWorkerIntegrationSuite struct { + tests.IntegrationTestSuite +} + +func TestAgentInstancePruneWorkerIntegrationSuite(t *testing.T) { + suite.Run(t, new(AgentInstancePruneWorkerIntegrationSuite)) +} + +func (s *AgentInstancePruneWorkerIntegrationSuite) SetupTest() { + s.Require().NoError(s.Migrator.Refresh()) +} + +func (s *AgentInstancePruneWorkerIntegrationSuite) insert(agentID uuid.UUID, daemon *bool, lastSeen time.Time) uuid.UUID { + id := uuid.New() + s.Require().NoError(s.DB.Create(&relational.AgentInstance{ + AgentID: agentID, InstanceID: id, Daemon: daemon, + FirstSeenAt: lastSeen, LastSeenAt: lastSeen, CreatedAt: lastSeen, UpdatedAt: lastSeen, + }).Error) + return id +} + +func (s *AgentInstancePruneWorkerIntegrationSuite) remaining() map[uuid.UUID]bool { + var rows []relational.AgentInstance + s.Require().NoError(s.DB.Find(&rows).Error) + out := map[uuid.UUID]bool{} + for _, r := range rows { + out[r.InstanceID] = true + } + return out +} + +func (s *AgentInstancePruneWorkerIntegrationSuite) TestWorkPrunesWithDefaultRetention() { + now := time.Date(2026, 3, 1, 12, 0, 0, 0, time.UTC) + agent, err := s.CreateAgent("prune-worker") + s.Require().NoError(err) + f, tr := false, true + oneShotOld := s.insert(*agent.ID, &f, now.Add(-25*time.Hour)) + oneShotNew := s.insert(*agent.ID, &f, now.Add(-time.Hour)) + daemon25h := s.insert(*agent.ID, &tr, now.Add(-25*time.Hour)) + daemonOld := s.insert(*agent.ID, &tr, now.Add(-721*time.Hour)) + unknown25h := s.insert(*agent.ID, nil, now.Add(-25*time.Hour)) + + w := NewAgentInstancePruneWorker(s.DB, agentcfg.Settings{}, nil) + w.now = func() time.Time { return now } + s.Require().NoError(w.Work(context.Background(), &river.Job[AgentInstancePruneArgs]{})) + + left := s.remaining() + s.Len(left, 3) + s.False(left[oneShotOld]) + s.True(left[oneShotNew]) + s.True(left[daemon25h]) + s.False(left[daemonOld]) + s.True(left[unknown25h]) + + // A second run is a no-op. + s.Require().NoError(w.Work(context.Background(), &river.Job[AgentInstancePruneArgs]{})) + s.Len(s.remaining(), 3) +} + +func (s *AgentInstancePruneWorkerIntegrationSuite) TestWorkHonoursConfiguredRetention() { + now := time.Date(2026, 3, 1, 12, 0, 0, 0, time.UTC) + agent, err := s.CreateAgent("prune-worker-custom") + s.Require().NoError(err) + f, tr := false, true + oneShot := s.insert(*agent.ID, &f, now.Add(-2*time.Hour)) + daemon := s.insert(*agent.ID, &tr, now.Add(-49*time.Hour)) + daemonKept := s.insert(*agent.ID, &tr, now.Add(-47*time.Hour)) + + w := NewAgentInstancePruneWorker(s.DB, agentcfg.Settings{OneShotInstanceRetention: time.Hour, InstanceRetention: 48 * time.Hour}, nil) + w.now = func() time.Time { return now } + s.Require().NoError(w.Work(context.Background(), &river.Job[AgentInstancePruneArgs]{})) + + left := s.remaining() + s.False(left[oneShot]) + s.False(left[daemon]) + s.True(left[daemonKept]) +} diff --git a/internal/service/worker/agent_instance_prune_test.go b/internal/service/worker/agent_instance_prune_test.go new file mode 100644 index 00000000..0143afbd --- /dev/null +++ b/internal/service/worker/agent_instance_prune_test.go @@ -0,0 +1,31 @@ +package worker + +import ( + "testing" + + "github.com/compliance-framework/api/internal/config" + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" + "go.uber.org/zap" +) + +func TestPeriodicJobsFromConfig_AgentInstancePrune(t *testing.T) { + logger := zap.NewNop().Sugar() + + assert.Len(t, periodicJobsFromConfig(&config.Config{}, logger), 0, "nil Agents config") + + disabled := config.DefaultAgentsConfig() + disabled.InstancePruneEnabled = false + assert.Len(t, periodicJobsFromConfig(&config.Config{Agents: disabled}, logger), 0) + + assert.Len(t, periodicJobsFromConfig(&config.Config{Agents: config.DefaultAgentsConfig()}, logger), 1) + + invalidSchedule := config.DefaultAgentsConfig() + invalidSchedule.InstancePruneSchedule = "not a cron" + assert.Len(t, periodicJobsFromConfig(&config.Config{Agents: invalidSchedule}, logger), 1, "falls back to the default schedule") +} + +func TestAgentInstancePruneArgsKind(t *testing.T) { + assert.Equal(t, "agent_instance_prune", AgentInstancePruneArgs{}.Kind()) + require.NotNil(t, NewAgentInstancePrunePeriodicJob("", zap.NewNop().Sugar())) +} diff --git a/internal/service/worker/service.go b/internal/service/worker/service.go index cefed17a..f2f97b58 100644 --- a/internal/service/worker/service.go +++ b/internal/service/worker/service.go @@ -16,6 +16,7 @@ import ( "github.com/compliance-framework/api/internal/service/notification" emailprovider "github.com/compliance-framework/api/internal/service/notification/providers/email" slackprovider "github.com/compliance-framework/api/internal/service/notification/providers/slack" + "github.com/compliance-framework/api/internal/service/relational/agentcfg" riskrel "github.com/compliance-framework/api/internal/service/relational/risks" "github.com/compliance-framework/api/internal/service/relational/workflows" slacksvc "github.com/compliance-framework/api/internal/service/slack" @@ -319,6 +320,10 @@ func NewServiceWithDigest( poamOpenDigestSchedulerWorker := NewPoamOpenDigestSchedulerWorker(db, clientProxy, poamCfg.OpenDigestWindow, logger) river.AddWorker(workers, river.WorkFunc(poamOpenDigestSchedulerWorker.Work)) + // Agent instance pruning (R37) + agentInstancePruneWorker := NewAgentInstancePruneWorker(db, agentcfg.SettingsFromConfig(digestCfg), logger) + river.AddWorker(workers, river.WorkFunc(agentInstancePruneWorker.Work)) + aiEnabled := digestCfg != nil && digestCfg.AI != nil && digestCfg.AI.Enabled if aiEnabled { llmClient := llm.NewAnthropicClient(llm.AnthropicConfig{ @@ -714,6 +719,10 @@ func periodicJobsFromConfig(cfg *config.Config, logger *zap.SugaredLogger) []*ri if cfg.Poam != nil && cfg.Poam.OpenDigestEnabled { periodicJobs = append(periodicJobs, NewPoamOpenDigestPeriodicJob(cfg.Poam.OpenDigestSchedule, logger)) } + // Agent instance pruning (R37); LoadAgentsConfig enables it by default. + if cfg.Agents != nil && cfg.Agents.InstancePruneEnabled { + periodicJobs = append(periodicJobs, NewAgentInstancePrunePeriodicJob(cfg.Agents.InstancePruneSchedule, logger)) + } return periodicJobs } From 3d696106f752d8942e57c3bb0b977b01da875b30 Mon Sep 17 00:00:00 2001 From: "ccf-lisa[bot]" <286799724+ccf-lisa[bot]@users.noreply.github.com> Date: Mon, 5 Oct 2026 14:59:27 -0300 Subject: [PATCH 2/5] fix(agentcfg): take the instance mode from the validated report mode reportedRemote preferred the reported remote-config block's mode, which the API does not validate, over the report's validated top-level mode, so preview could classify with a different or unknown mode than the one ValidationBases selected on. Co-Authored-By: Claude Opus 5.5 --- .../relational/agentcfg/remote_test.go | 35 +++++++++++++++++++ .../service/relational/agentcfg/service.go | 22 +++++------- 2 files changed, 43 insertions(+), 14 deletions(-) create mode 100644 internal/service/relational/agentcfg/remote_test.go diff --git a/internal/service/relational/agentcfg/remote_test.go b/internal/service/relational/agentcfg/remote_test.go new file mode 100644 index 00000000..81de6f67 --- /dev/null +++ b/internal/service/relational/agentcfg/remote_test.go @@ -0,0 +1,35 @@ +package agentcfg + +import ( + "testing" + + "github.com/compliance-framework/api/internal/service/relational" + "github.com/compliance-framework/api/pkg/agentconfig" + "github.com/stretchr/testify/assert" + "gorm.io/datatypes" +) + +// The mode always comes from the validated top-level report mode, never from the reported +// remote-config block, which is not validated. +func TestReportedRemoteUsesRowMode(t *testing.T) { + cases := []struct { + name string + remote string + base agentconfig.Config + }{ + {name: "block says apply_all", remote: `{"mode":"apply_all","trusted_sources":["ghcr.io/x/*"]}`}, + {name: "block has an unknown mode", remote: `{"mode":"bogus","trusted_sources":["ghcr.io/x/*"]}`}, + {name: "no block, base says apply_all", base: agentconfig.Config{RemoteConfig: &agentconfig.RemoteConfig{Mode: agentconfig.ModeApplyAll, TrustedSources: []string{"ghcr.io/x/*"}}}}, + } + for _, tc := range cases { + t.Run(tc.name, func(t *testing.T) { + row := relational.AgentInstance{Mode: agentconfig.ModeApplySafe} + if tc.remote != "" { + row.RemoteConfig = datatypes.JSON(tc.remote) + } + rc := reportedRemote(row, tc.base) + assert.Equal(t, agentconfig.ModeApplySafe, rc.Mode) + assert.Equal(t, []string{"ghcr.io/x/*"}, rc.TrustedSources) + }) + } +} diff --git a/internal/service/relational/agentcfg/service.go b/internal/service/relational/agentcfg/service.go index 08075d04..d0aad687 100644 --- a/internal/service/relational/agentcfg/service.go +++ b/internal/service/relational/agentcfg/service.go @@ -796,24 +796,18 @@ func (s *Service) toBases(rows []relational.AgentInstance, now time.Time, valida return out } -// reportedRemote returns the instance's reported remote-config block, or the base's block -// normalized with hasAuth=true (the base is redacted, so it has no client secret). +// reportedRemote returns the instance's reported remote-config block, or the base's block, +// normalized with hasAuth=true (the base is redacted, so it has no client secret). The mode +// always comes from the report's validated top-level mode (row.Mode), which ValidationBases +// also selects on; the block's own mode is not validated, so it is ignored. func reportedRemote(row relational.AgentInstance, base agentconfig.Config) agentconfig.RemoteConfig { var rc agentconfig.RemoteConfig - if len(row.RemoteConfig) > 0 && string(row.RemoteConfig) != "null" { - if err := json.Unmarshal(row.RemoteConfig, &rc); err == nil { - if rc.Mode == "" { - rc.Mode = row.Mode - } - return rc.Normalize(true) - } - } - if base.RemoteConfig != nil { + decoded := len(row.RemoteConfig) > 0 && string(row.RemoteConfig) != "null" && + json.Unmarshal(row.RemoteConfig, &rc) == nil + if !decoded && base.RemoteConfig != nil { rc = *base.RemoteConfig } - if rc.Mode == "" { - rc.Mode = row.Mode - } + rc.Mode = row.Mode return rc.Normalize(true) } From 44af657935ea3c9dfe7e3bf7c44c20119581eed3 Mon Sep 17 00:00:00 2001 From: "ccf-lisa[bot]" <286799724+ccf-lisa[bot]@users.noreply.github.com> Date: Tue, 6 Oct 2026 07:25:59 -0300 Subject: [PATCH 3/5] fix(worker): take the prune schedule fallback from config.DefaultAgentsConfig --- internal/service/worker/agent_instance_prune.go | 6 ++---- 1 file changed, 2 insertions(+), 4 deletions(-) diff --git a/internal/service/worker/agent_instance_prune.go b/internal/service/worker/agent_instance_prune.go index b4178f23..2fd4a6b3 100644 --- a/internal/service/worker/agent_instance_prune.go +++ b/internal/service/worker/agent_instance_prune.go @@ -5,6 +5,7 @@ import ( "fmt" "time" + "github.com/compliance-framework/api/internal/config" "github.com/compliance-framework/api/internal/service/relational/agentcfg" "github.com/riverqueue/river" "go.uber.org/zap" @@ -14,9 +15,6 @@ import ( // JobTypeAgentInstancePrune prunes agent instances that are no longer reporting (R37). const JobTypeAgentInstancePrune = "agent_instance_prune" -// defaultAgentInstancePruneSchedule is hourly, River 6-field (seconds first). -const defaultAgentInstancePruneSchedule = "0 17 * * * *" - // AgentInstancePruneArgs are the (empty) args of the periodic prune job. type AgentInstancePruneArgs struct{} @@ -54,7 +52,7 @@ func (w *AgentInstancePruneWorker) Work(ctx context.Context, _ *river.Job[AgentI // NewAgentInstancePrunePeriodicJob schedules the prune job on the "scheduler" queue. func NewAgentInstancePrunePeriodicJob(schedule string, logger *zap.SugaredLogger) *river.PeriodicJob { - sched := parseCronScheduleWithFallback(schedule, defaultAgentInstancePruneSchedule, "agent instance prune", logger) + sched := parseCronScheduleWithFallback(schedule, config.DefaultAgentsConfig().InstancePruneSchedule, "agent instance prune", logger) return river.NewPeriodicJob( sched, func() (river.JobArgs, *river.InsertOpts) { From 0e7f7b0f329f4da5ec6542b135739b7e99f8b32f Mon Sep 17 00:00:00 2001 From: "ccf-lisa[bot]" <286799724+ccf-lisa[bot]@users.noreply.github.com> Date: Tue, 6 Oct 2026 08:22:49 -0300 Subject: [PATCH 4/5] fix(agentcfg): load and decode each distinct validation base once ValidationBases loaded base_config for every fresh apply-mode instance (up to the instance cap, each up to 4 MiB), so a save's cost followed the instance count. Group the set by base content in SQL (SHA-256 of the jsonb text), load and decode one base per group in the same read-only snapshot, and share it across the group's instances. Each instance is still returned with its own fields and remote-config, plus a BaseKey so callers validate a base once. Co-Authored-By: Claude Opus 5.5 --- .../service/relational/agentcfg/service.go | 129 +++++++++++++++--- .../agentcfg/service_integration_test.go | 50 +++++++ 2 files changed, 159 insertions(+), 20 deletions(-) diff --git a/internal/service/relational/agentcfg/service.go b/internal/service/relational/agentcfg/service.go index d0aad687..59852e7a 100644 --- a/internal/service/relational/agentcfg/service.go +++ b/internal/service/relational/agentcfg/service.go @@ -4,6 +4,7 @@ package agentcfg import ( "context" + "database/sql" "encoding/json" "errors" "fmt" @@ -640,10 +641,30 @@ type InstanceBase struct { Remote agentconfig.RemoteConfig // reported remote-config, or the base's block normalized with hasAuth=true Stale bool Validated bool // member of ValidationBases (R48) + // BaseKey identifies the content of the reported base (set by ValidationBases only): + // instances with the same BaseKey share one decoded Base, so a caller validates it once. + // Base is shared, not copied: treat it as read-only. + BaseKey string } var applyModes = []string{agentconfig.ModeApplySafe, agentconfig.ModeApplyAll} +// baseKeyExpr is the content key of a reported base: the SHA-256 of its jsonb text, which +// Postgres normalizes (key order, whitespace), so instances that report the same base get +// the same key. +const baseKeyExpr = "encode(sha256(convert_to(base_config::text, 'UTF8')), 'hex') AS base_key" + +// validationMemberColumns are baseColumns without the base itself, plus its content key. +var validationMemberColumns = append(slices.DeleteFunc(slices.Clone(baseColumns), func(c string) bool { + return c == "base_config" +}), baseKeyExpr) + +// validationMember is one instance of the validation set, without its base. +type validationMember struct { + relational.AgentInstance + BaseKey string +} + // ValidationBases is exactly the set PUT and revert validate against (R14, R48): // 1. all fresh instances (seen within InstanceStaleAfter) with a reported base and an // apply mode; @@ -653,33 +674,101 @@ var applyModes = []string{agentconfig.ModeApplySafe, agentconfig.ModeApplyAll} // // Report-mode instances and instances without a base are never validated against. A base // that no longer decodes is skipped with a warning. +// +// The cost follows the distinct bases, not the instance count: the set is grouped by base +// content in SQL (BaseKey), and each distinct base is loaded and decoded once and shared by +// every instance of its group (instances of one fleet usually report the same base). Both +// reads run in one read-only repeatable-read transaction, so the groups and the loaded bases +// are one snapshot. func (s *Service) ValidationBases(ctx context.Context, agentID uuid.UUID) ([]InstanceBase, bool, error) { now := s.now() - var fresh []relational.AgentInstance - err := s.db.WithContext(ctx). - Select(baseColumns). - Where("agent_id = ? AND base_config IS NOT NULL AND mode IN ? AND last_seen_at >= ?", agentID, applyModes, now.Add(-s.settings.InstanceStaleAfter)). - Order("last_seen_at DESC, instance_id"). - Find(&fresh).Error + var bases []InstanceBase + err := s.db.WithContext(ctx).Transaction(func(tx *gorm.DB) error { + var members []validationMember + if err := s.findValidationSet(tx, agentID, now, validationMemberColumns, &members); err != nil { + return err + } + if len(members) == 0 { + return nil + } + // One representative instance per distinct base. + keyOf := map[uuid.UUID]string{} + seen := map[string]bool{} + var repIDs []uuid.UUID + for _, m := range members { + if m.ID == nil || seen[m.BaseKey] { + continue + } + seen[m.BaseKey] = true + keyOf[*m.ID] = m.BaseKey + repIDs = append(repIDs, *m.ID) + } + var reps []relational.AgentInstance + if err := tx.Select("id", "base_config").Where("id IN ?", repIDs).Find(&reps).Error; err != nil { + return err + } + decoded := make(map[string]agentconfig.Config, len(reps)) + for _, row := range reps { + key := keyOf[*row.ID] + base, err := agentconfig.DecodeConfig(row.BaseConfig) + if err != nil { + s.logger.Warnw("Skipping agent instances with an undecodable reported base", + "agentID", agentID, "instanceIDs", instancesWithKey(members, key), "error", err) + continue + } + decoded[key] = base + } + for _, m := range members { + base, ok := decoded[m.BaseKey] + if !ok { + continue + } + bases = append(bases, InstanceBase{ + Instance: m.AgentInstance, + Base: base, + Remote: reportedRemote(m.AgentInstance, base), + Stale: IsStale(m.AgentInstance, now, s.settings), + Validated: true, + BaseKey: m.BaseKey, + }) + } + return nil + }, &sql.TxOptions{Isolation: sql.LevelRepeatableRead, ReadOnly: true}) if err != nil { return nil, false, err } - rows := fresh - if len(rows) == 0 { - var latest []relational.AgentInstance - err := s.db.WithContext(ctx). - Select(baseColumns). - Where("agent_id = ? AND base_config IS NOT NULL AND mode IN ? AND reported_at IS NOT NULL", agentID, applyModes). - Order("reported_at DESC, instance_id"). - Limit(1). - Find(&latest).Error - if err != nil { - return nil, false, err + return bases, len(bases) == 0, nil +} + +// findValidationSet finds the given columns of the ValidationBases set into dest (a pointer +// to a slice of rows): the fresh apply-mode instances with a reported base, newest first, +// else the most recently reported one. +func (s *Service) findValidationSet(db *gorm.DB, agentID uuid.UUID, now time.Time, columns []string, dest any) error { + res := db.Model(&relational.AgentInstance{}). + Select(strings.Join(columns, ", ")). + Where("agent_id = ? AND base_config IS NOT NULL AND mode IN ? AND last_seen_at >= ?", agentID, applyModes, now.Add(-s.settings.InstanceStaleAfter)). + Order("last_seen_at DESC, instance_id"). + Find(dest) + if res.Error != nil || res.RowsAffected > 0 { + return res.Error + } + return db.Model(&relational.AgentInstance{}). + Select(strings.Join(columns, ", ")). + Where("agent_id = ? AND base_config IS NOT NULL AND mode IN ? AND reported_at IS NOT NULL", agentID, applyModes). + Order("reported_at DESC, instance_id"). + Limit(1). + Find(dest).Error +} + +// instancesWithKey lists the instance ids of the members that report the base with key. +func instancesWithKey(members []validationMember, key string) []uuid.UUID { + var ids []uuid.UUID + for _, m := range members { + if m.BaseKey == key { + ids = append(ids, m.InstanceID) } - rows = latest } - bases := s.toBases(rows, now, true) - return bases, len(bases) == 0, nil + return ids } // Preview bounds (R14): a preview shows at most PreviewMaxInstances instances and decodes diff --git a/internal/service/relational/agentcfg/service_integration_test.go b/internal/service/relational/agentcfg/service_integration_test.go index 784aa290..dece305d 100644 --- a/internal/service/relational/agentcfg/service_integration_test.go +++ b/internal/service/relational/agentcfg/service_integration_test.go @@ -855,6 +855,56 @@ func (s *AgentCfgServiceIntegrationSuite) TestValidationBasesFallbackAndStandalo s.Empty(bases) } +// Instances that report the same base (whatever its key order or whitespace) share one +// BaseKey and one decoded base; each keeps its own instance fields and remote-config. +func (s *AgentCfgServiceIntegrationSuite) TestValidationBasesGroupsByBaseContent() { + agentID := s.newAgent("grouped-bases") + a, b, c := uuid.New(), uuid.New(), uuid.New() + reordered := `{ "verbosity":0, "plugins":{"p1":{"policies":["ghcr.io/x/pol:v1"],"source":"ghcr.io/x/p1:v1"}}, "api":{"auth":{"client_id":"cid"},"url":"http://api:8080"}, "daemon":true }` + other := `{"daemon":true,"verbosity":1,"plugins":{}}` + for _, in := range []struct { + id uuid.UUID + base string + host string + at time.Duration + trusted []string + instMode string + }{ + {a, baseConfig, "host-a", 0, []string{"ghcr.io/a/*"}, agentconfig.ModeApplySafe}, + {b, reordered, "host-b", -time.Minute, []string{"ghcr.io/b/*"}, agentconfig.ModeApplyAll}, + {c, other, "host-c", -2 * time.Minute, nil, agentconfig.ModeApplySafe}, + } { + r := applyReport(in.instMode, in.base) + r.Hostname = in.host + r.RemoteConfig = &agentconfig.RemoteConfig{Mode: in.instMode, TrustedSources: in.trusted} + s.Require().NoError(s.reportAt(s.svc, s.now.Add(in.at), agentID, in.id, r)) + } + + bases, standalone, err := s.svc.ValidationBases(s.ctx, agentID) + s.Require().NoError(err) + s.False(standalone) + s.Require().Len(bases, 3, "every instance of the validation set is still listed") + s.Equal([]uuid.UUID{a, b, c}, []uuid.UUID{bases[0].Instance.InstanceID, bases[1].Instance.InstanceID, bases[2].Instance.InstanceID}) + + m := byInstance(bases) + s.NotEmpty(m[a].BaseKey) + s.Equal(m[a].BaseKey, m[b].BaseKey, "same base content, same key") + s.NotEqual(m[a].BaseKey, m[c].BaseKey) + s.Equal(m[a].Base, m[b].Base) + s.Equal(int32(1), m[c].Base.Verbosity) + s.Nil(m[a].Instance.BaseConfig, "the base is loaded once per group, not per instance") + + s.Equal("host-a", *m[a].Instance.Hostname) + s.Equal("host-b", *m[b].Instance.Hostname) + s.Equal([]string{"ghcr.io/a/*"}, m[a].Remote.TrustedSources) + s.Equal([]string{"ghcr.io/b/*"}, m[b].Remote.TrustedSources) + s.Equal(agentconfig.ModeApplyAll, m[b].Remote.Mode) + for _, base := range bases { + s.True(base.Validated) + s.False(base.Stale) + } +} + func (s *AgentCfgServiceIntegrationSuite) TestPreviewBases() { f := s.seedValidationFixture() From b4d34992234ac59146c9b9366d91bf0e474664b2 Mon Sep 17 00:00:00 2001 From: "ccf-lisa[bot]" <286799724+ccf-lisa[bot]@users.noreply.github.com> Date: Tue, 6 Oct 2026 10:55:42 -0300 Subject: [PATCH 5/5] fix(agentcfg): list instances by page in the prune test ListInstances is paginated now (#476); TestPruneInstances reads its six instances from the first page. Co-Authored-By: Claude Opus 5.5 --- .../service/relational/agentcfg/service_integration_test.go | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/internal/service/relational/agentcfg/service_integration_test.go b/internal/service/relational/agentcfg/service_integration_test.go index dece305d..00182fbd 100644 --- a/internal/service/relational/agentcfg/service_integration_test.go +++ b/internal/service/relational/agentcfg/service_integration_test.go @@ -1039,7 +1039,7 @@ func (s *AgentCfgServiceIntegrationSuite) TestPruneInstances() { s.Equal(int64(3), deleted) remaining := map[uuid.UUID]bool{} - list, err := s.svc.ListInstances(s.ctx, agentID) + list, _, err := s.svc.ListInstances(s.ctx, agentID, service.PaginationParams{Page: 1, Limit: agentcfg.InstancesPageLimit}) s.Require().NoError(err) for _, i := range list { remaining[i.InstanceID] = true