From 691f4696d47c1105e165619191a40ce59dec19d6 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:56:37 -0300 Subject: [PATCH 1/4] feat(agent): pull and apply remote configuration overlays (apply_safe/apply_all) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit With remote_config.mode apply_safe or apply_all the reconciler polls the API for the configuration overlay and applies it on top of the file: - Fetch every poll_interval (±10%) with the cached raw ETag as If-None-Match; a 304 keeps the cached overlay, and the API's error table applies to fetches as to reports (R7, R8). - prepare decodes the overlay strictly (unknown keys and non-string values reject it, R27), runs the agentconfig Classify + WillApply gate against the agent's own base and remote_config before anything touches the network, merges it (RFC 7396) and validates the result. Errors at pointers the overlay touched are strict; file-origin ones keep the R34 tolerance. - A rejected revision is remembered per (ETag, base fingerprint), or per revision + overlay hash without an ETag, with its reason, error and unsafe changes; it is not re-prepared every poll, is re-reported after a restart, and a base change re-classifies it. A failed revision (download, run) is retried with the failed backoff. - Startup tries the fetched overlay, then the cached applied one, then the file alone. A no-op revision is recorded as applied without a restart. - Reports carry the applied and attempted revisions and the unsafe changes; evidence produced under an overlay carries the agent-config-revision prop (R38). ${env:NAME} placeholders in plugin config are resolved in the next PR; until then they reach plugins unchanged, as on main. Co-Authored-By: Claude Opus 5.5 --- cmd/agent.go | 25 +- cmd/config.go | 28 +- cmd/config_test.go | 81 ++++ cmd/reconciler.go | 502 ++++++++++++++++++++++--- cmd/remote_test.go | 700 +++++++++++++++++++++++++++++++++++ cmd/report.go | 27 +- docs/configuration.md | 52 ++- docs/running_as_a_service.md | 2 +- 8 files changed, 1340 insertions(+), 77 deletions(-) diff --git a/cmd/agent.go b/cmd/agent.go index 577674d..664cf73 100644 --- a/cmd/agent.go +++ b/cmd/agent.go @@ -191,6 +191,12 @@ const DefaultProtocolVersion int32 = 1 const RunnerV2ProtocolVersion int32 = 2 const AnnotationProtocolVersionKey = "org.ccf.plugin.protocol.version" +// CCFPropNamespace is the OSCAL prop namespace of CCF props. +const CCFPropNamespace = "https://compliance-framework.github.io/ns" + +// configRevisionPropName stamps evidence with the applied remote configuration revision (R38). +const configRevisionPropName = "agent-config-revision" + // daemonCronStopTimeout bounds the cron stop on SIGINT/SIGTERM before plugins are killed, // also when the signal arrives during a reload drain (R33). var daemonCronStopTimeout = 30 * time.Second @@ -239,7 +245,7 @@ with plugins to ensure continuous compliance.`, agentCmd.Flags().StringP("config", "c", "", "Location of config file") agentCmd.MarkFlagRequired("config") - agentCmd.Flags().String("state-dir", "", "Directory for this instance's state (instance ID); overrides CCF_STATE_DIR. Default: .compliance-framework/state/") + agentCmd.Flags().String("state-dir", "", "Directory for this instance's state (instance ID, remote config cache); overrides CCF_STATE_DIR. Default: .compliance-framework/state/") agentCmd.Flags().String("instance-id", "", "Pin this instance's UUID (not persisted); overrides CCF_INSTANCE_ID") return agentCmd @@ -1450,6 +1456,7 @@ func (ar *AgentRunner) runAllPlugins(ctx context.Context) error { resultsHelper := runner.NewApiHelper(logger, client, labels, pluginName, runner.WithPolicyPaths(policyPaths), runner.WithSources(sourceOf(pluginConfig.Source, source), policySources), + runner.WithEvidenceProps(configRevisionProps(config)...), ) policyBehaviorProto := policyBehaviorToProto(pluginConfig.PolicyBehavior) @@ -1594,6 +1601,7 @@ func (ar *AgentRunner) runPluginWith(ctx context.Context, snap runSnapshot, name resultsHelper := runner.NewApiHelper(pluginLogger, client, labels, name, runner.WithPolicyPaths(policyPaths), runner.WithSources(sourceOf(plugin.Source, pluginExecutable), policySources), + runner.WithEvidenceProps(configRevisionProps(config)...), ) policyBehaviorProto := policyBehaviorToProto(plugin.PolicyBehavior) @@ -1635,6 +1643,20 @@ func (ar *AgentRunner) SendHeartbeat(ctx context.Context, staticAgentUUID uuid.U return nil } +// configRevisionProps returns the evidence prop naming the applied overlay revision, or nil +// when the agent runs the file only (R38). +func configRevisionProps(config *agentConfig) []sdktypes.Property { + meta := config.syncInfo() + if meta.AppliedRevision <= 0 { + return nil + } + return []sdktypes.Property{{ + Ns: CCFPropNamespace, + Name: configRevisionPropName, + Value: strconv.FormatInt(meta.AppliedRevision, 10), + }} +} + // buildHeartbeat builds the heartbeat body. When remote configuration is not off it carries // the applied revision (0 when running the file only, never null) and the effective digest, // which lets the API create the instance row (R11, R45). @@ -1756,6 +1778,7 @@ func (ar *AgentRunner) buildAgentRunEvidence(now time.Time) (*agentEvidenceCreat End: now, Expires: expires, Links: links, + Props: configRevisionProps(config), Status: sdktypes.ObjectiveStatus{ Reason: reason, Remarks: remarks, diff --git a/cmd/config.go b/cmd/config.go index aa237ba..5494c9b 100644 --- a/cmd/config.go +++ b/cmd/config.go @@ -31,7 +31,7 @@ type baseSnapshot struct { warnings []agentconfig.FieldError // skip holds the plugins dropped from the runtime because of a tolerated problem. skip map[string]string - // fingerprint identifies the base for the failed backoff. + // fingerprint identifies the base for the rejected-revision memory. fingerprint string } @@ -240,7 +240,7 @@ func baseFromViper(cmd *cobra.Command, v *viper.Viper, raw []byte) (*baseSnapsho raw: raw, envSourced: envSourcedPointers(v), } - part := partitionByOrigin(declared.Validate()) + part := partitionByOrigin(declared.Validate(), nil) if len(part.fatal) > 0 { return nil, agentconfig.ValidationErrors(part.fatal) } @@ -252,15 +252,18 @@ func baseFromViper(cmd *cobra.Command, v *viper.Viper, raw []byte) (*baseSnapsho // validationPartition is the R34 split of a config's validation errors. type validationPartition struct { + overlay []agentconfig.FieldError // touched by the overlay: strict fatal []agentconfig.FieldError // file-origin, not tolerated: fatal warnings []agentconfig.FieldError // file-origin, tolerated or warn-only: reported skip map[string]string // plugin name -> reason, for tolerated (skip) errors } -// partitionByOrigin splits the validation errors of the file (R34): tolerated rules become -// warnings (and the plugin is skipped), warn-only rules become warnings (nothing is skipped or +// partitionByOrigin splits validation errors by origin (R34). An error at pointer P is +// overlay-origin when some overlay-touched pointer o equals P, is a prefix of P, or has P as a +// prefix (segment-wise). Everything else is file-origin: tolerated rules become warnings +// (and the plugin is skipped), warn-only rules become warnings (nothing is skipped or // changed), the rest is fatal. -func partitionByOrigin(err error) validationPartition { +func partitionByOrigin(err error, overlayTouched []string) validationPartition { var out validationPartition if err == nil { return out @@ -272,6 +275,8 @@ func partitionByOrigin(err error) validationPartition { } for _, e := range errs { switch { + case touchedByOverlay(e.Path, overlayTouched): + out.overlay = append(out.overlay, e) case isToleratedFileRule(e): out.warnings = append(out.warnings, e) if segs := agentconfig.SplitPointer(e.Path); len(segs) >= 2 && segs[0] == "plugins" { @@ -289,6 +294,19 @@ func partitionByOrigin(err error) validationPartition { return out } +// touchedByOverlay compares pointers segment-wise in both directions. +func touchedByOverlay(ptr string, touched []string) bool { + p := agentconfig.SplitPointer(ptr) + for _, o := range touched { + t := agentconfig.SplitPointer(o) + n := min(len(p), len(t)) + if slices.Equal(p[:n], t[:n]) { + return true + } + } + return false +} + // toRuntime converts a merged, env-resolved declared config into the runtime structs. // Disabled plugins and plugins named in skip (R34) are dropped: they get no cron, no download // and no run state, but they stay in the declared form and in reports. diff --git a/cmd/config_test.go b/cmd/config_test.go index d171cb8..7aa9d93 100644 --- a/cmd/config_test.go +++ b/cmd/config_test.go @@ -1,6 +1,7 @@ package cmd import ( + "encoding/json" "errors" "os" "path/filepath" @@ -9,6 +10,7 @@ import ( "github.com/compliance-framework/api/pkg/agentconfig" "github.com/hashicorp/go-hclog" + "google.golang.org/protobuf/proto" ) // writeConfigFile writes content to a temp file with the given extension and returns its path. @@ -65,6 +67,41 @@ func TestLoadBase_WeakDecodingUnchanged(t *testing.T) { } } +// TestWeakDecoding_SurvivesUnrelatedOverlay checks that an overlay touching only the schedule +// leaves the plugin's config and policy_data unchanged on the wire (R51). +func TestWeakDecoding_SurvivesUnrelatedOverlay(t *testing.T) { + base := mustLoadBase(t, "yaml", weakTypedConfig) + fileOnly, err := toRuntime(base.declared, nil) + if err != nil { + t.Fatal(err) + } + merged, err := agentconfig.Merge(base.declared, json.RawMessage(`{"plugins":{"aws":{"schedule":"*/5 * * * *"}}}`)) + if err != nil { + t.Fatal(err) + } + withOverlay, err := toRuntime(merged, nil) + if err != nil { + t.Fatal(err) + } + if !reflect.DeepEqual(fileOnly.Plugins["aws"].Config, withOverlay.Plugins["aws"].Config) { + t.Fatalf("config changed by an unrelated overlay: %#v vs %#v", fileOnly.Plugins["aws"].Config, withOverlay.Plugins["aws"].Config) + } + a, err := mapToStruct(fileOnly.Plugins["aws"].PolicyData) + if err != nil { + t.Fatal(err) + } + b, err := mapToStruct(withOverlay.Plugins["aws"].PolicyData) + if err != nil { + t.Fatal(err) + } + if !proto.Equal(a, b) { + t.Fatalf("policy_data structpb differs: %v vs %v", a, b) + } + if got := *withOverlay.Plugins["aws"].Schedule; got != "*/5 * * * *" { + t.Fatalf("overlay schedule not applied: %q", got) + } +} + func TestEnvSourcedPointers(t *testing.T) { t.Setenv("CCF_PLUGINS_GITHUB_CONFIG_TOKEN", "from-env") base := mustLoadBase(t, "yaml", ` @@ -174,6 +211,16 @@ func TestLoadBase_FileOriginWarnOnly(t *testing.T) { }) } + t.Run("overlay-origin stays strict", func(t *testing.T) { + errs := agentconfig.ValidationErrors{ + {Path: "/verbosity", Code: agentconfig.FieldCodeInvalidValue, Message: "must not be negative"}, + {Path: "/plugins/ssh/labels/team", Code: agentconfig.FieldCodeEnvLocation, Message: "env"}, + } + p := partitionByOrigin(errs, []string{"/verbosity", "/plugins/ssh/labels/team"}) + if len(p.overlay) != 2 || len(p.warnings) != 0 { + t.Fatalf("overlay-introduced values must be strict, got %#v", p) + } + }) } // TestLoadBase_LoadsAsOnMain: YAML that JSON cannot represent loads, and a key the agent does @@ -209,6 +256,40 @@ plugins: } } +func TestPartitionByOrigin(t *testing.T) { + errs := agentconfig.ValidationErrors{ + {Path: "/plugins/ssh/schedule", Code: agentconfig.FieldCodeCron, Message: "bad cron"}, + {Path: "/plugins/github/source", Code: agentconfig.FieldCodeRequired, Message: "source required"}, + } + t.Run("overlay touches another field of the same plugin", func(t *testing.T) { + p := partitionByOrigin(errs, []string{"/plugins/ssh/labels/team"}) + if len(p.overlay) != 0 || len(p.warnings) != 1 || len(p.fatal) != 1 { + t.Fatalf("unexpected partition %#v", p) + } + if _, ok := p.skip["ssh"]; !ok { + t.Fatalf("expected ssh skipped, got %v", p.skip) + } + }) + t.Run("overlay sets the schedule", func(t *testing.T) { + p := partitionByOrigin(errs, []string{"/plugins/ssh/schedule"}) + if len(p.overlay) != 1 || p.overlay[0].Path != "/plugins/ssh/schedule" || len(p.warnings) != 0 { + t.Fatalf("unexpected partition %#v", p) + } + }) + t.Run("overlay adds the plugin", func(t *testing.T) { + p := partitionByOrigin(errs, []string{"/plugins/ssh"}) + if len(p.overlay) != 1 { + t.Fatalf("a prefix pointer must make the error overlay-origin, got %#v", p) + } + }) + t.Run("segment-wise, not string-wise", func(t *testing.T) { + p := partitionByOrigin(errs, []string{"/plugins/ss"}) + if len(p.overlay) != 0 { + t.Fatalf("/plugins/ss must not match /plugins/ssh, got %#v", p) + } + }) +} + func TestToRuntime_DisabledPluginDropped(t *testing.T) { base := mustLoadBase(t, "yaml", ` api: diff --git a/cmd/reconciler.go b/cmd/reconciler.go index 170111b..9d2f9f6 100644 --- a/cmd/reconciler.go +++ b/cmd/reconciler.go @@ -9,6 +9,7 @@ import ( "errors" "fmt" "math/rand" + "slices" "strings" "sync" "sync/atomic" @@ -25,7 +26,7 @@ import ( ) var ( - // remoteRequestTimeout bounds one config report. + // remoteRequestTimeout bounds one config fetch or report; the startup fetch too. remoteRequestTimeout = 30 * time.Second // remoteAuthBackoff is the retry delay after a 404 (API without the feature) or a 401/403 // on a config route (R8, R36). @@ -48,9 +49,10 @@ var ( // once built. type candidate struct { base *baseSnapshot - declared agentconfig.Config // reported and digested - runtime *agentConfig // resolved, enabled-only, skipped plugins removed - digest string // agentconfig.Digest(declared, base.redactOpts()...) (R55) + overlay *agentstate.OverlayRecord // nil = file only + declared agentconfig.Config // merged, ${env:} NOT resolved: reported and digested + runtime *agentConfig // resolved, enabled-only, skipped plugins removed + digest string // agentconfig.Digest(declared, base.redactOpts()...) (R55) // identity changes whenever anything that affects the runtime changes, including the // values the digest masks or omits (api block, secrets). It never leaves the process. identity string @@ -59,12 +61,22 @@ type candidate struct { plugins []agentconfig.PluginReport } +// appliedRevision is the overlay revision the candidate applies, or nil for the file only. +func (c *candidate) appliedRevision() *int64 { + if c == nil || c.overlay == nil { + return nil + } + rev := c.overlay.Revision + return &rev +} + // applyError is why a candidate could not be prepared. Status is agentconfig.StatusRejected // or agentconfig.StatusFailed and Reason is one of agentconfig.Reasons. type applyError struct { Status string Reason string Err error + Unsafe []agentconfig.Change // runtime is the prepared runtime of a download-failed candidate: startup hands it to // onStartupFailure so the startup-failure evidence describes it, as on main. runtime *agentConfig @@ -82,6 +94,10 @@ func (e *applyError) Error() string { func (e *applyError) Unwrap() error { return e.Err } +func rejected(reason string, err error) *applyError { + return &applyError{Status: agentconfig.StatusRejected, Reason: reason, Err: err} +} + func failed(reason string, err error) *applyError { return &applyError{Status: agentconfig.StatusFailed, Reason: reason, Err: err} } @@ -91,6 +107,11 @@ type prefetcher interface { Prefetch(ctx context.Context, cfg *agentConfig) error } +// overlayFetcher is the test seam over sdk.Client.AgentConfig.Get. +type overlayFetcher interface { + Get(ctx context.Context, ifNoneMatch string) (*sdk.AgentConfigResult, error) +} + // configReporter is the test seam over sdk.Client.AgentConfig.Report. type configReporter interface { Report(ctx context.Context, instanceID uuid.UUID, r agentconfig.Report) error @@ -98,6 +119,7 @@ type configReporter interface { // remoteAPI bundles the remote configuration calls. type remoteAPI interface { + overlayFetcher configReporter } @@ -106,6 +128,10 @@ type sdkRemote struct { client *sdk.Client } +func (s sdkRemote) Get(ctx context.Context, ifNoneMatch string) (*sdk.AgentConfigResult, error) { + return s.client.AgentConfig.Get(ctx, ifNoneMatch) +} + func (s sdkRemote) Report(ctx context.Context, instanceID uuid.UUID, r agentconfig.Report) error { return s.client.AgentConfig.Report(ctx, instanceID, r) } @@ -135,9 +161,9 @@ const ( triggerPoll ) -// reconciler is the single writer of the configuration state. Triggers are serialized in one -// goroutine: a trigger builds a complete candidate and only then cancels the running -// configuration; a failure tears nothing down. +// reconciler is the single writer of the configuration state. File and remote triggers are +// serialized in one goroutine: a trigger builds a complete candidate and only then cancels the +// running configuration; a failure tears nothing down. type reconciler struct { cmd *cobra.Command configPath string @@ -169,12 +195,16 @@ type reconciler struct { base *baseSnapshot remote remoteAPI remoteKey string + cache *agentstate.Cache lastOutcome *applyError + attempted *int64 warnedMode bool + fetchBackoffUntil time.Time reportBackoffUntil time.Time loggedOnce map[string]bool + failedKey string // overlayKey of the target in failed backoff failedBase string failedRetryAt time.Time failedInterval time.Duration @@ -210,9 +240,10 @@ func isApplyMode(mode string) bool { return mode == agentconfig.ModeApplySafe || mode == agentconfig.ModeApplyAll } -// setBase installs a new base: warnings are logged and the remote client is rebuilt when the -// api block changed. +// setBase installs a new base: warnings are logged, the remote client is rebuilt when the api +// block changed, and the cache is (re)bound to the base's identity. func (rc *reconciler) setBase(base *baseSnapshot) { + old := rc.base rc.base = base rc.logWarnings(base.warnings) if !rc.warnedMode && base.declared.RemoteConfig != nil { @@ -230,7 +261,22 @@ func (rc *reconciler) setBase(base *baseSnapshot) { rc.remote = rc.newRemote(base.declared) } rc.remoteKey = key - rc.reportBackoffUntil = time.Time{} + rc.fetchBackoffUntil, rc.reportBackoffUntil = time.Time{}, time.Time{} + } + + id := cacheIdentity(base.declared) + if rc.cache == nil || rc.cache.Identity != id { + cache, err := rc.store.LoadCache(id) + rc.cache = cache + if errors.Is(err, agentstate.ErrCacheCorrupt) { + rc.logger.Error("Remote config cache is corrupt; continuing without it", "path", rc.store.CachePath(), "error", err) + rc.lastOutcome = failed(agentconfig.ReasonCacheCorrupt, err) + } + } + if old != nil && old.fingerprint != base.fingerprint && rc.cache.Rejected != nil { + // A base change re-classifies a remembered rejection. + rc.cache.Rejected = nil + rc.saveCache() } } @@ -242,6 +288,26 @@ func remoteKey(c agentconfig.Config) string { return string(raw) } +func cacheIdentity(c agentconfig.Config) agentstate.Identity { + id := agentstate.Identity{} + if c.API != nil { + id.APIURL = strings.TrimSpace(c.API.URL) + if c.API.Auth != nil { + id.ClientID = strings.TrimSpace(c.API.Auth.ClientID) + } + } + return id +} + +func (rc *reconciler) saveCache() { + if rc.cache == nil { + return + } + if err := rc.store.SaveCache(rc.cache); err != nil && rc.logOnce("cache-save") { + rc.logger.Warn("Could not persist the remote config cache", "path", rc.store.CachePath(), "error", err) + } +} + // logOnce reports whether key has not been logged yet, and marks it. func (rc *reconciler) logOnce(key string) bool { if rc.loggedOnce[key] { @@ -251,57 +317,248 @@ func (rc *reconciler) logOnce(key string) bool { return true } -// startup loads the file, prepares it (R32) and reports it. Only an unusable local -// configuration is an error (exit 1, as before). +// startup runs the startup ladder (R32): load the file, fetch (apply modes, bounded), then the +// first candidate that prepares wins among base+fetched, base+applied and base only. Only an +// unusable local configuration is an error (exit 1, as before). func (rc *reconciler) startup(ctx context.Context) (*candidate, error) { base, err := loadBase(rc.cmd, rc.configPath) if err != nil { return nil, fmt.Errorf("config file: %w", err) } rc.setBase(base) + rcfg := rc.rcfg() + if isApplyMode(rcfg.Mode) { + rc.fetch(ctx) + } - active, aerr := rc.prepare(ctx, base) - if aerr != nil { - if aerr.runtime != nil && rc.onStartupFailure != nil { - rc.onStartupFailure(ctx, aerr.runtime, aerr.Err) + var active *candidate + var outcome *applyError + for _, target := range rc.ladder(rcfg.Mode) { + cand, aerr := rc.prepare(ctx, base, target) + if target != nil && target == rc.cache.Fetched { + rev := target.Revision + rc.attempted = &rev } - return nil, aerr + if aerr != nil { + if target == nil { + if aerr.runtime != nil && rc.onStartupFailure != nil { + rc.onStartupFailure(ctx, aerr.runtime, aerr.Err) + } + return nil, aerr + } + rc.logger.Warn("Could not apply the remote configuration at startup", "revision", target.Revision, "status", aerr.Status, "reason", aerr.Reason, "error", aerr.Err) + rc.recordFailure(target, aerr) + if outcome == nil { + outcome = aerr + } + continue + } + active = cand + if target != nil { + rc.cache.Applied = target + rc.saveCache() + } + break + } + if rej := rc.rememberedRejection(rcfg.Mode); rej != nil && outcome == nil { + // The ladder skipped the fetched overlay because it was rejected before the restart: + // report that rejection again, or the API sees an applied, never-attempted revision + // (pending) until a new revision is published (a 304 never re-prepares it). + outcome = rej + } + if outcome != nil { + rc.lastOutcome = outcome } rc.maybeReport(ctx, active, rc.lastOutcome) return active, nil } -// startFailedBackoff starts (or doubles, for the same base) the retry delay of a candidate that -// failed to prepare or to run on base baseFingerprint. -func (rc *reconciler) startFailedBackoff(baseFingerprint string) { - if rc.failedBase == baseFingerprint && rc.failedInterval > 0 { +// ladder lists the startup targets in order; nil is the file only. +func (rc *reconciler) ladder(mode string) []*agentstate.OverlayRecord { + if !isApplyMode(mode) { + return []*agentstate.OverlayRecord{nil} + } + var out []*agentstate.OverlayRecord + if f := rc.cache.Fetched; f != nil && !rc.rememberedRejected(f) { + out = append(out, f) + } + if a := rc.cache.Applied; a != nil && (len(out) == 0 || !sameOverlay(a, out[0])) { + out = append(out, a) + } + return append(out, nil) +} + +func sameOverlay(a, b *agentstate.OverlayRecord) bool { + switch { + case a == nil || b == nil: + return a == b + case a.ETag != "" || b.ETag != "": + return a.ETag == b.ETag && a.Revision == b.Revision + default: + return a.Revision == b.Revision && bytes.Equal(a.Overlay, b.Overlay) + } +} + +// overlayKey identifies an overlay for the rejected memory and the failed backoff: the raw +// ETag, or revision + sha256(overlay) when a response carried no ETag (a stripping proxy), so +// one rejection never blocks every later revision. nil (the file only) has its own key. +func overlayKey(rec *agentstate.OverlayRecord) string { + switch { + case rec == nil: + return "file-only" + case rec.ETag != "": + return "etag:" + rec.ETag + default: + return fmt.Sprintf("rev:%d:%s", rec.Revision, overlayDigest(rec.Overlay)) + } +} + +func overlayDigest(raw []byte) string { + sum := sha256.Sum256(raw) + return hex.EncodeToString(sum[:]) +} + +func (rc *reconciler) rememberedRejected(rec *agentstate.OverlayRecord) bool { + r := rc.cache.Rejected + if r == nil || rec == nil || r.BaseFingerprint != rc.base.fingerprint { + return false + } + if r.ETag != "" || rec.ETag != "" { + return r.ETag == rec.ETag + } + return r.Revision == rec.Revision && r.OverlaySHA256 == overlayDigest(rec.Overlay) +} + +// rememberedRejection is the outcome to report while the fetched overlay is remembered as +// rejected for this base (apply modes only): the persisted rejection, with rc.attempted set to +// its revision. nil when the fetched overlay is not remembered-rejected. +func (rc *reconciler) rememberedRejection(mode string) *applyError { + if !isApplyMode(mode) || rc.cache == nil { + return nil + } + f, r := rc.cache.Fetched, rc.cache.Rejected + if f == nil || !rc.rememberedRejected(f) { + return nil + } + rev := f.Revision + rc.attempted = &rev + aerr := &applyError{ + Status: r.Status, + Reason: r.Reason, + Unsafe: slices.Clone(r.Unsafe), + } + if r.Error != "" { + aerr.Err = errors.New(r.Error) + } + return aerr +} + +// recordFailure remembers a rejected revision for (overlay key, base), or starts the failed +// backoff. A file-only candidate that fails to prepare is backed off too. +func (rc *reconciler) recordFailure(target *agentstate.OverlayRecord, aerr *applyError) { + if target == nil { + rc.startFailedBackoff(nil, rc.base.fingerprint) + return + } + if aerr.Status == agentconfig.StatusRejected { + msg := "" + if aerr.Err != nil { + msg = aerr.Err.Error() + } + rc.cache.Rejected = &agentstate.RejectedRecord{ + Revision: target.Revision, + ETag: target.ETag, + OverlaySHA256: overlayDigest(target.Overlay), + BaseFingerprint: rc.base.fingerprint, + Status: aerr.Status, + Reason: aerr.Reason, + Error: msg, + Unsafe: aerr.Unsafe, + } + rc.saveCache() + return + } + rc.startFailedBackoff(target, rc.base.fingerprint) +} + +// startFailedBackoff starts (or doubles, for the same target and base) the retry delay of a +// target (nil = the file only) that failed to prepare or to run on base baseFingerprint. +func (rc *reconciler) startFailedBackoff(target *agentstate.OverlayRecord, baseFingerprint string) { + key := overlayKey(target) + if rc.failedKey == key && rc.failedBase == baseFingerprint && rc.failedInterval > 0 { rc.failedInterval = min(rc.failedInterval*2, failedRetryMax) } else { rc.failedInterval = failedRetryMin } - rc.failedBase = baseFingerprint + rc.failedKey, rc.failedBase = key, baseFingerprint rc.failedRetryAt = rc.now().Add(rc.failedInterval) } -func (rc *reconciler) inFailedBackoff() bool { - return rc.failedInterval > 0 && rc.failedBase == rc.base.fingerprint && rc.now().Before(rc.failedRetryAt) +func (rc *reconciler) inFailedBackoff(target *agentstate.OverlayRecord) bool { + return rc.failedInterval > 0 && rc.failedKey == overlayKey(target) && + rc.failedBase == rc.base.fingerprint && rc.now().Before(rc.failedRetryAt) } -// clearFailedBackoff forgets the failed backoff after a candidate prepared on the current base, -// unless it is the base in backoff: that one may still fail to RUN, and the retry delay must -// keep growing instead of restarting at failedRetryMin (no flapping every poll). -func (rc *reconciler) clearFailedBackoff() { - if rc.failedBase == rc.base.fingerprint { +// clearFailedBackoff forgets the failed backoff after target prepared on the current base, +// unless it is the target in backoff: that one may still fail to RUN, and the retry delay +// must keep growing instead of restarting at failedRetryMin (no flapping every poll). +func (rc *reconciler) clearFailedBackoff(target *agentstate.OverlayRecord) { + if rc.failedKey == overlayKey(target) && rc.failedBase == rc.base.fingerprint { return } - rc.failedBase, rc.failedInterval = "", 0 + rc.failedKey, rc.failedBase, rc.failedInterval = "", "", 0 } -// prepare builds a candidate from a base. It never touches the running configuration. -func (rc *reconciler) prepare(ctx context.Context, base *baseSnapshot) (*candidate, *applyError) { +// prepare builds a candidate from a base and an optional overlay (G3.3). Cheap checks run +// first and nothing touches the network before the Classify gate passes. It never touches the +// running configuration. +func (rc *reconciler) prepare(ctx context.Context, base *baseSnapshot, ov *agentstate.OverlayRecord) (*candidate, *applyError) { rcfg := base.declared.EffectiveRemoteConfig() + if !isApplyMode(rcfg.Mode) { + ov = nil // report/off: the file only + } + declared := base.declared - runtime, err := toRuntime(declared, base.skip) + var touched []string + if ov != nil { + // Strict decode of the overlay: the only strict decode in the agent (R27, R51). + if err := agentconfig.ValidateOverlay(ov.Overlay); err != nil { + return nil, overlayValidationError(err) + } + changes, err := agentconfig.Classify(base.declared, ov.Overlay, rcfg) + if err != nil { + return nil, rejected(agentconfig.ReasonInvalidConfig, err) + } + if ok, why := agentconfig.WillApply(rcfg, changes); !ok { + aerr := rejected(why, fmt.Errorf("revision %d %s", ov.Revision, strings.ReplaceAll(why, "-", " "))) + for _, c := range changes { + if c.Safety != agentconfig.Safe { + aerr.Unsafe = append(aerr.Unsafe, c) + } + } + return nil, aerr + } + merged, err := agentconfig.Merge(base.declared, ov.Overlay) + if err != nil { + return nil, rejected(agentconfig.ReasonInvalidConfig, err) + } + touched, err = overlayTouched(base.declared, merged) + if err != nil { + return nil, failed(agentconfig.ReasonInternal, err) + } + declared = merged + } + + part := partitionByOrigin(declared.Validate(), touched) + if len(part.overlay) > 0 || len(part.fatal) > 0 { + errs := agentconfig.ValidationErrors(append(append([]agentconfig.FieldError{}, part.overlay...), part.fatal...)) + if ov == nil { + return nil, failed(agentconfig.ReasonInvalidConfig, fmt.Errorf("config file: %w", errs)) + } + return nil, rejected(agentconfig.ReasonInvalidConfig, errs) + } + + runtime, err := toRuntime(declared, part.skip) if err != nil { return nil, failed(agentconfig.ReasonInvalidConfig, err) } @@ -321,18 +578,77 @@ func (rc *reconciler) prepare(ctx context.Context, base *baseSnapshot) (*candida // The digest is over the UNRESOLVED form with the same masking as the reported effective // config (R55): it never changes when a secret rotates. digest := agentconfig.Digest(declared, base.redactOpts()...) - runtime.setSync(syncMeta{Digest: digest, Mode: rcfg.Mode}) + meta := syncMeta{Digest: digest, Mode: rcfg.Mode} + if ov != nil { + meta.AppliedRevision = ov.Revision + } + runtime.setSync(meta) return &candidate{ base: base, + overlay: ov, declared: declared, runtime: runtime, digest: digest, identity: candidateIdentity(declared), - warnings: append([]agentconfig.FieldError{}, base.warnings...), + warnings: append([]agentconfig.FieldError{}, part.warnings...), plugins: plugins, }, nil } +// overlayTouched returns the pointers an overlay changed, computed on the unresolved forms so +// env resolution never counts as an overlay change (R34). +func overlayTouched(base, merged agentconfig.Config) ([]string, error) { + a, err := json.Marshal(base) + if err != nil { + return nil, err + } + b, err := json.Marshal(merged) + if err != nil { + return nil, err + } + diff, err := agentconfig.DiffJSON(a, b) + if err != nil { + return nil, err + } + out := make([]string, 0, len(diff)) + for _, d := range diff { + out = append(out, d.Path) + } + return out, nil +} + +// overlayValidationError maps ValidateOverlay's FieldError codes to a report reason (R43). +// With several errors the first in precedence order wins: forbidden, unknown-field, +// invalid-type, invalid-config. All errors are kept in the message. +func overlayValidationError(err error) *applyError { + var errs agentconfig.ValidationErrors + if !errors.As(err, &errs) { + return rejected(agentconfig.ReasonInvalidConfig, err) + } + rank := map[string]int{ + agentconfig.ReasonForbiddenChanges: 0, + agentconfig.ReasonUnknownField: 1, + agentconfig.ReasonInvalidType: 2, + agentconfig.ReasonInvalidConfig: 3, + } + reason := agentconfig.ReasonInvalidConfig + for _, e := range errs { + r := agentconfig.ReasonInvalidConfig + switch e.Code { + case agentconfig.FieldCodeLockedKey, agentconfig.FieldCodeForbiddenEnv: + r = agentconfig.ReasonForbiddenChanges + case agentconfig.FieldCodeUnknownField: + r = agentconfig.ReasonUnknownField + case agentconfig.FieldCodeInvalidType: + r = agentconfig.ReasonInvalidType + } + if rank[r] < rank[reason] { + reason = r + } + } + return rejected(reason, errs) +} + // candidateIdentity hashes the declared config, including the values the digest masks or // omits, so any change that affects the runtime is a different configuration. func candidateIdentity(c agentconfig.Config) string { @@ -360,8 +676,8 @@ func (rc *reconciler) bind(active *candidate, cancel context.CancelFunc) { // adopt records cand, whose identity equals old's, as the running (or pending) configuration // WITHOUT a restart: the runtime old runs is kept and only its sync metadata changes, so the -// heartbeat and report show cand's, and sameAsActive holds for cand's base on the next trigger -// (no re-prepare every poll). +// heartbeat, evidence and report show cand's applied revision, and sameAsActive holds for +// cand's base and overlay on the next trigger (no re-prepare every poll). func (rc *reconciler) adopt(old, cand *candidate) *candidate { adopted := *cand if old.runtime != nil { @@ -475,8 +791,8 @@ func (rc *reconciler) pollDelay() time.Duration { } // loop is the daemon's reconcile goroutine: config file events (debounced) and the poll -// ticker, which retries a candidate whose failed backoff expired and resends the report when -// due. It returns when ctx is done. +// ticker. The poller lives here, not in the heartbeat cron, because it triggers reloads. It +// returns when ctx is done. func (rc *reconciler) loop(ctx context.Context) { var debounce <-chan time.Time poll := time.NewTimer(rc.pollDelay()) @@ -505,13 +821,22 @@ func (rc *reconciler) loop(ctx context.Context) { func (rc *reconciler) onRunFailed(ctx context.Context, c *candidate) { aerr := failed(agentconfig.ReasonInternal, errors.New("the configuration failed to start; running the previous configuration")) if c != nil { - // Back the candidate off: otherwise every poll re-prepares the new base, cancels the - // healthy configuration, fails and falls back again. + // File-only candidates are backed off too: otherwise every poll re-prepares the new + // base, cancels the healthy configuration, fails and falls back again. baseFingerprint := rc.base.fingerprint if c.base != nil { baseFingerprint = c.base.fingerprint } - rc.startFailedBackoff(baseFingerprint) + rc.startFailedBackoff(c.overlay, baseFingerprint) + } + if c != nil && c.overlay != nil { + if sameOverlay(rc.cache.Applied, c.overlay) { + rc.cache.Applied = nil + if prev := rc.current(); prev != nil && prev.overlay != nil { + rc.cache.Applied = prev.overlay + } + rc.saveCache() + } } rc.lastOutcome = aerr rc.maybeReport(ctx, rc.current(), aerr) @@ -531,37 +856,110 @@ func (rc *reconciler) reconcile(ctx context.Context, t trigger) { } rc.setBase(base) } + rcfg := rc.rcfg() + if isApplyMode(rcfg.Mode) && t == triggerPoll { + rc.fetch(ctx) + } + target := rc.targetOverlay(rcfg.Mode) + if isApplyMode(rcfg.Mode) && target != nil && target == rc.cache.Fetched { + rev := target.Revision + rc.attempted = &rev + } + // While the fetched overlay stays rejected (a 304 never re-prepares it), it remains the + // attempted revision and its rejection the outcome, even when the fallback re-prepares. + remembered := rc.rememberedRejection(rcfg.Mode) active := rc.current() - if rc.sameAsActive(active) || rc.inFailedBackoff() { + if rc.sameAsActive(active, target) { + rc.maybeReport(ctx, active, rc.lastOutcome) + return + } + if target != nil && rc.rememberedRejected(target) { + // The only overlay left (the applied one) was rejected for this base, e.g. after a + // conflicting file edit: keep the last-known-good configuration instead of + // re-preparing it on every poll (§5.4, G3.4). + rc.maybeReport(ctx, active, rc.lastOutcome) + return + } + if rc.inFailedBackoff(target) { rc.maybeReport(ctx, active, rc.lastOutcome) return } - cand, aerr := rc.prepare(ctx, rc.base) + cand, aerr := rc.prepare(ctx, rc.base, target) if aerr != nil { rc.logger.Warn("Could not apply the configuration; keeping the running configuration", "status", aerr.Status, "reason", aerr.Reason, "error", aerr.Err) - rc.startFailedBackoff(rc.base.fingerprint) + rc.recordFailure(target, aerr) rc.lastOutcome = aerr rc.maybeReport(ctx, active, aerr) return } - rc.clearFailedBackoff() - rc.lastOutcome = nil + rc.clearFailedBackoff(target) + if isApplyMode(rcfg.Mode) { + rc.cache.Applied = target + rc.saveCache() + } + rc.lastOutcome = remembered if active != nil && active.identity == cand.identity { - rc.logger.Debug("Trigger did not change the effective configuration; recording it without a restart") + rc.logger.Debug("Trigger did not change the effective configuration; recording it without a restart", "revision", revisionForLog(target)) rc.maybeReport(ctx, rc.adopt(active, cand), rc.lastOutcome) return } - rc.logger.Info("Applying the new configuration") + rc.logger.Info("Applying the new configuration", "revision", revisionForLog(target)) rc.swap(cand) rc.maybeReport(ctx, cand, rc.lastOutcome) } -// sameAsActive reports whether the active candidate already runs this base. -func (rc *reconciler) sameAsActive(active *candidate) bool { +func revisionForLog(ov *agentstate.OverlayRecord) any { + if ov == nil { + return "file-only" + } + return ov.Revision +} + +// targetOverlay is the overlay the agent should run: none in report/off; otherwise the newest +// fetched one unless it was rejected for this base, else the applied one. +func (rc *reconciler) targetOverlay(mode string) *agentstate.OverlayRecord { + if !isApplyMode(mode) || rc.cache == nil { + return nil + } + if f := rc.cache.Fetched; f != nil && !rc.rememberedRejected(f) { + return f + } + return rc.cache.Applied +} + +// sameAsActive reports whether the active candidate already runs this base and target. +func (rc *reconciler) sameAsActive(active *candidate, target *agentstate.OverlayRecord) bool { return active != nil && active.base != nil && - bytes.Equal(active.base.raw, rc.base.raw) && active.base.fingerprint == rc.base.fingerprint + bytes.Equal(active.base.raw, rc.base.raw) && active.base.fingerprint == rc.base.fingerprint && + sameOverlay(active.overlay, target) +} + +// fetch polls the API for the overlay (R8), honoring the error backoffs. The cached overlay +// keeps applying on any error. +func (rc *reconciler) fetch(ctx context.Context) { + if rc.remote == nil || rc.now().Before(rc.fetchBackoffUntil) { + return + } + fetchCtx, cancel := context.WithTimeout(ctx, remoteRequestTimeout) + defer cancel() + res, err := rc.remote.Get(fetchCtx, rc.cache.IfNoneMatch()) + if err != nil { + rc.handleRemoteError("fetch", err, &rc.fetchBackoffUntil) + return + } + switch { + case res.NotModified: + case res.Document != nil: + rc.cache.Fetched = &agentstate.OverlayRecord{ + Revision: res.Document.Revision, + ETag: res.ETag, + Overlay: append(json.RawMessage(nil), res.Document.Overlay...), + FetchedAt: rc.now().UTC(), + } + rc.saveCache() + } } // handleRemoteError applies the R8/R36 error table to a config route error. diff --git a/cmd/remote_test.go b/cmd/remote_test.go index 7d90909..943a456 100644 --- a/cmd/remote_test.go +++ b/cmd/remote_test.go @@ -1,11 +1,15 @@ package cmd import ( + "bytes" "context" "encoding/json" + "errors" "fmt" "os" "path/filepath" + "reflect" + "runtime" "strings" "sync" "testing" @@ -15,6 +19,7 @@ import ( "github.com/compliance-framework/api/pkg/agentconfig" "github.com/compliance-framework/api/sdk" "github.com/google/uuid" + "github.com/hashicorp/go-hclog" ) // fakeRemote is a scripted API for the remote configuration routes. @@ -254,6 +259,39 @@ func TestStartupReport_RedactsAndDescribes(t *testing.T) { } } +func TestReport_ResendPolicy(t *testing.T) { + h := newRemoteHarness(t, remoteConfig("apply_safe", "")) // no document: 404 → file only + h.remote.overlay = json.RawMessage(`{}`) + h.remote.etag = `"r0-x"` + mustStartup(t, h.rc) + if h.remote.reportCount() != 1 { + t.Fatalf("expected the startup report, got %d", h.remote.reportCount()) + } + h.poll(t) + if h.remote.reportCount() != 1 { + t.Fatalf("an unchanged report must not be resent, got %d", h.remote.reportCount()) + } + h.clock.Advance(24*time.Hour + time.Second) + h.poll(t) + if h.remote.reportCount() != 2 { + t.Fatalf("expected a resend after 24h, got %d", h.remote.reportCount()) + } + h.remote.reportErr = func(int, agentconfig.Report) error { return errors.New("network down") } + h.remote.publish(1, `{"plugins":{"ssh":{"schedule":"*/5 * * * *"}}}`) + h.poll(t) + if h.remote.reportCount() != 3 { + t.Fatalf("expected a report after a change, got %d", h.remote.reportCount()) + } + h.remote.reportErr = nil + h.poll(t) + if h.remote.reportCount() != 4 { + t.Fatalf("expected a resend after a send error, got %d", h.remote.reportCount()) + } + if r := h.remote.lastReport(t); r.AppliedRevision == nil || *r.AppliedRevision != 1 { + t.Fatalf("expected applied revision 1, got %+v", r.AppliedRevision) + } +} + func TestModes_ReportAndOff(t *testing.T) { t.Run("report mode never fetches", func(t *testing.T) { h := newRemoteHarness(t, remoteConfig("report", "")) @@ -272,6 +310,47 @@ func TestModes_ReportAndOff(t *testing.T) { t.Fatalf("report mode heartbeat must carry revision 0 and a digest: %+v", hb) } }) + t.Run("credentials without a mode default to report", func(t *testing.T) { + // An agent that applied an overlay under an explicit apply mode... + h := newRemoteHarness(t, remoteConfig("apply_safe", "")) + h.remote.publish(1, `{"plugins":{"ssh":{"schedule":"*/5 * * * *"}}}`) + if active := mustStartup(t, h.rc); active.overlay == nil { + t.Fatal("expected the overlay applied under apply_safe") + } + gets := h.remote.getCount() + + // ...restarts with credentials and no mode, with and without a remote_config block. + noBlock := strings.Replace(remoteBaseConfig, "remote_config:\n mode: %MODE%\n trusted_sources: [\"ghcr.io/trusted/*\"]\n overridable_config_flags: [%FLAGS%]\n", "", 1) + if strings.Contains(noBlock, "remote_config") { + t.Fatal("the fixture still has a remote_config block") + } + for _, tc := range []struct{ name, content string }{ + {"no remote_config block", noBlock}, + {"empty mode", remoteConfig("", "")}, + } { + t.Run(tc.name, func(t *testing.T) { + if err := os.WriteFile(h.path, []byte(tc.content), 0o600); err != nil { + t.Fatal(err) + } + h.rc = h.newReconciler() + active := mustStartup(t, h.rc) + h.poll(t) + if active.runtime.remote.Mode != agentconfig.ModeReport { + t.Fatalf("expected mode report, got %q", active.runtime.remote.Mode) + } + if active.overlay != nil { + t.Fatalf("report mode must not apply the cached overlay, got revision %d", active.overlay.Revision) + } + if h.remote.getCount() != gets { + t.Fatalf("report mode must never fetch, got %d new gets", h.remote.getCount()-gets) + } + r := h.remote.lastReport(t) + if r.Mode != agentconfig.ModeReport || r.Status != agentconfig.StatusNotApplicable || r.AppliedRevision != nil { + t.Fatalf("expected a not-applicable report in mode report, got %+v", r) + } + }) + } + }) t.Run("off sends nothing", func(t *testing.T) { h := newRemoteHarness(t, remoteConfig("off", "")) active := mustStartup(t, h.rc) @@ -308,6 +387,54 @@ func TestHeartbeat_FileOnlyApplySafe(t *testing.T) { } func TestRemoteErrors_Backoffs(t *testing.T) { + t.Run("404 backs off 10 minutes", func(t *testing.T) { + h := newRemoteHarness(t, remoteConfig("apply_safe", "")) + mustStartup(t, h.rc) + h.poll(t) + if h.remote.getCount() != 1 { + t.Fatalf("expected no fetch during the 404 backoff, got %d", h.remote.getCount()) + } + h.clock.Advance(remoteAuthBackoff + time.Second) + h.poll(t) + if h.remote.getCount() != 2 { + t.Fatalf("expected a fetch after the backoff, got %d", h.remote.getCount()) + } + }) + for _, code := range []int{401, 403} { + t.Run(fmt.Sprintf("%d backs off 10 minutes", code), func(t *testing.T) { + h := newRemoteHarness(t, remoteConfig("apply_safe", "")) + h.remote.getErr = &sdk.APIStatusError{StatusCode: code} + mustStartup(t, h.rc) + h.poll(t) + if h.remote.getCount() != 1 { + t.Fatalf("expected no fetch during the backoff, got %d", h.remote.getCount()) + } + h.clock.Advance(remoteAuthBackoff + time.Second) + h.poll(t) + if h.remote.getCount() != 2 { + t.Fatalf("expected a retry after the backoff, got %d", h.remote.getCount()) + } + }) + } + t.Run("409 pauses reports for an hour while polling continues", func(t *testing.T) { + h := newRemoteHarness(t, remoteConfig("apply_safe", "")) + h.remote.publish(0, `{}`) + h.remote.reportErr = func(int, agentconfig.Report) error { return &sdk.APIStatusError{StatusCode: 409} } + mustStartup(t, h.rc) + h.remote.publish(1, `{"plugins":{"ssh":{"schedule":"*/5 * * * *"}}}`) + h.poll(t) + if h.remote.reportCount() != 1 { + t.Fatalf("expected reports paused after a 409, got %d", h.remote.reportCount()) + } + if h.remote.getCount() != 2 { + t.Fatalf("polling must continue after a 409, got %d gets", h.remote.getCount()) + } + h.clock.Advance(reportConflictBackoff + time.Second) + h.poll(t) + if h.remote.reportCount() != 2 { + t.Fatalf("expected a report after the 1h pause, got %d", h.remote.reportCount()) + } + }) t.Run("413 resends truncated", func(t *testing.T) { h := newRemoteHarness(t, remoteConfig("apply_safe", "")) h.remote.publish(0, `{}`) @@ -378,3 +505,576 @@ func TestSetAgentVersion(t *testing.T) { } // --- G3: pull and apply --- + +func TestApply_BadOverlayKeepsOldConfig(t *testing.T) { + h := newRemoteHarness(t, remoteConfig("apply_safe", "")) + h.remote.publish(1, `{"plugins":{"ssh":{"schedule":"*/5 * * * *"}}}`) + first := mustStartup(t, h.rc) + h.remote.publish(2, `{"plugins":{"ssh":{"schedule":"bad cron"}}}`) + if got := h.poll(t); got != first { + t.Fatal("a bad overlay must keep the running config") + } + r := h.remote.lastReport(t) + if r.Status != agentconfig.StatusRejected || r.Reason != agentconfig.ReasonInvalidConfig || *r.AttemptedRevision != 2 || *r.AppliedRevision != 1 { + t.Fatalf("unexpected report %+v", r) + } +} + +func TestApply_DownloadFailureKeepsOldConfig(t *testing.T) { + h := newRemoteHarness(t, remoteConfig("apply_safe", "")) + h.remote.publish(1, `{"plugins":{"ssh":{"schedule":"*/5 * * * *"}}}`) + first := mustStartup(t, h.rc) + h.pf.setErr(errors.New("registry down")) + h.remote.publish(2, `{"plugins":{"ssh":{"schedule":"*/7 * * * *"}}}`) + if got := h.poll(t); got != first { + t.Fatal("a download failure must keep the running config") + } + if r := h.remote.lastReport(t); r.Status != agentconfig.StatusFailed || r.Reason != agentconfig.ReasonDownloadFailed { + t.Fatalf("unexpected report %+v", r) + } +} + +func TestApply_FailedBackoff(t *testing.T) { + h := newRemoteHarness(t, remoteConfig("apply_safe", "")) + h.remote.publish(0, `{}`) + mustStartup(t, h.rc) + h.pf.setErr(errors.New("registry down")) + h.remote.publish(1, `{"plugins":{"ssh":{"schedule":"*/7 * * * *"}}}`) + h.poll(t) + calls := h.pf.callCount() + h.clock.Advance(30 * time.Second) + h.poll(t) + if h.pf.callCount() != calls { + t.Fatal("a failed revision must not be retried before 1m") + } + h.clock.Advance(31 * time.Second) + h.poll(t) + if h.pf.callCount() != calls+1 { + t.Fatal("expected a retry after 1m") + } + for i := 0; i < 6; i++ { // 2m, 4m, 8m, 10m, 10m... + h.clock.Advance(failedRetryMax + time.Second) + h.poll(t) + } + if h.rc.failedInterval != failedRetryMax { + t.Fatalf("backoff must cap at %s, got %s", failedRetryMax, h.rc.failedInterval) + } +} + +func TestApply_ClassifyGate(t *testing.T) { + tests := []struct { + name string + mode string + flags string + overlay string + status string + reason string + }{ + {"api.url is forbidden", "apply_safe", "", `{"api":{"url":"http://evil"}}`, "rejected", "forbidden-changes"}, + {"api.url is forbidden in apply_all", "apply_all", "", `{"api":{"url":"http://evil"}}`, "rejected", "forbidden-changes"}, + {"new untrusted source is unsafe", "apply_safe", "", `{"plugins":{"ssh":{"source":"ghcr.io/other/plugin:v1"}}}`, "rejected", "unsafe-changes"}, + {"config change is unsafe by default", "apply_safe", "", `{"plugins":{"ssh":{"config":{"host":"other"}}}}`, "rejected", "unsafe-changes"}, + {"schedule-only change applies", "apply_safe", "", `{"plugins":{"ssh":{"schedule":"*/5 * * * *"}}}`, "applied", ""}, + {"trusted source applies", "apply_safe", "", `{"plugins":{"ssh":{"source":"ghcr.io/trusted/plugin:v2"}}}`, "applied", ""}, + {"unqualified overridable flag", "apply_safe", `"host"`, `{"plugins":{"ssh":{"config":{"host":"other"}}}}`, "applied", ""}, + {"plugin:key overridable flag", "apply_safe", `"ssh:host"`, `{"plugins":{"ssh":{"config":{"host":"other"}}}}`, "applied", ""}, + {"other plugin's flag does not match", "apply_safe", `"github:host"`, `{"plugins":{"ssh":{"config":{"host":"other"}}}}`, "rejected", "unsafe-changes"}, + {"star overridable flag", "apply_safe", `"*"`, `{"plugins":{"ssh":{"config":{"host":"other"}}}}`, "applied", ""}, + {"trusted new plugin with overridable keys", "apply_safe", `"*"`, `{"plugins":{"new":{"source":"ghcr.io/trusted/new:v1","config":{"a":"b"}}}}`, "applied", ""}, + {"new env ref in an overridable key", "apply_safe", `"*"`, `{"plugins":{"ssh":{"config":{"host":"${env:HOST}"}}}}`, "rejected", "unsafe-changes"}, + {"unsafe applies in apply_all", "apply_all", "", `{"plugins":{"ssh":{"config":{"host":"other"}}}}`, "applied", ""}, + {"R27 non-string config value", "apply_all", "", `{"plugins":{"ssh":{"config":{"port":2222}}}}`, "rejected", "invalid-type"}, + {"R27 unknown field", "apply_all", "", `{"evidence_capture":{}}`, "rejected", "unknown-field"}, + {"R28 mixed-case plugin name", "apply_all", "", `{"plugins":{"GitHub":{"source":"ghcr.io/trusted/gh:v1"}}}`, "rejected", "invalid-config"}, + {"inline policy bundles are not supported", "apply_all", "", `{"policy_bundles":{"ssh":{"modules":{"a.rego":"package compliance_framework.a"}}}}`, "rejected", "unknown-field"}, + } + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + h := newRemoteHarness(t, remoteConfig(tt.mode, tt.flags)) + h.remote.publish(1, tt.overlay) + mustStartup(t, h.rc) + r := h.remote.lastReport(t) + if r.Status != tt.status || r.Reason != tt.reason { + t.Fatalf("got %s/%s (%v), want %s/%s", r.Status, r.Reason, derefString(r.Error), tt.status, tt.reason) + } + if tt.reason == "unsafe-changes" && len(r.Unsafe) == 0 { + t.Fatal("expected the unsafe changes to be listed") + } + }) + } +} + +func derefString(s *string) string { + if s == nil { + return "" + } + return *s +} + +func TestApply_404FallsBackToCacheThenFile(t *testing.T) { + h := newRemoteHarness(t, remoteConfig("apply_safe", "")) + h.remote.publish(1, `{"plugins":{"ssh":{"schedule":"*/5 * * * *"}}}`) + mustStartup(t, h.rc) + + h.remote.overlay = nil // the API now answers 404 + restarted := h.newReconciler() + active := mustStartup(t, restarted) + if active.overlay == nil || active.overlay.Revision != 1 { + t.Fatalf("expected the cached overlay after a 404, got %+v", active.overlay) + } + + if err := os.Remove(filepath.Join(h.dir, "state", "remote-config.json")); err != nil { + t.Fatal(err) + } + fileOnly := mustStartup(t, h.newReconciler()) + if fileOnly.overlay != nil { + t.Fatalf("expected the file only without a cache, got %+v", fileOnly.overlay) + } +} + +func TestApply_OpaqueETag(t *testing.T) { + h := newRemoteHarness(t, remoteConfig("apply_safe", "")) + h.remote.publish(3, `{"plugins":{"ssh":{"schedule":"*/5 * * * *"}}}`) + mustStartup(t, h.rc) + firstETag := h.remote.etag + h.poll(t) + if got := h.remote.gets[len(h.remote.gets)-1]; got != firstETag { + t.Fatalf("If-None-Match must be the raw ETag %q, got %q", firstETag, got) + } + // The API is reset and reuses revision 3 with another overlay and a new ETag. + h.remote.publish(3, `{"plugins":{"ssh":{"schedule":"*/9 * * * *"}}}`) + if got := *h.poll(t).runtime.Plugins["ssh"].Schedule; got != "*/9 * * * *" { + t.Fatalf("a new ETag with a reused revision must be re-evaluated, got schedule %q", got) + } + for _, sent := range h.remote.gets { + if sent == "3" || sent == `"3"` { + t.Fatalf("the agent must never build an ETag from a revision, sent %q", sent) + } + } +} + +func TestCache_IdentityCorruptionAndMode(t *testing.T) { + h := newRemoteHarness(t, remoteConfig("apply_safe", "")) + h.remote.publish(1, `{"plugins":{"ssh":{"schedule":"*/5 * * * *"}}}`) + mustStartup(t, h.rc) + cachePath := filepath.Join(h.dir, "state", "remote-config.json") + info, err := os.Stat(cachePath) + if err != nil { + t.Fatal(err) + } + if runtime.GOOS != "windows" && info.Mode().Perm() != 0o600 { + t.Fatalf("cache mode = %v, want 0600", info.Mode().Perm()) + } + + t.Run("identity mismatch discards the cache", func(t *testing.T) { + store := agentstate.Open(filepath.Join(h.dir, "state"), nil) + c, err := store.LoadCache(agentstate.Identity{APIURL: "http://other", ClientID: "x"}) + if err != nil || c.Applied != nil || c.Fetched != nil { + t.Fatalf("expected an empty cache, got %+v, %v", c, err) + } + }) + + t.Run("corruption is reported and the agent continues", func(t *testing.T) { + raw, err := os.ReadFile(cachePath) + if err != nil { + t.Fatal(err) + } + if err := os.WriteFile(cachePath, []byte(strings.Replace(string(raw), `"revision": 1`, `"revision": 9`, 1)), 0o600); err != nil { + t.Fatal(err) + } + h.remote.overlay = nil // 404: nothing to fetch + active := mustStartup(t, h.newReconciler()) + if active.overlay != nil { + t.Fatalf("a corrupt cache must not be applied, got %+v", active.overlay) + } + if r := h.remote.lastReport(t); r.Status != agentconfig.StatusFailed || r.Reason != agentconfig.ReasonCacheCorrupt { + t.Fatalf("expected failed/cache-corrupt, got %s/%s", r.Status, r.Reason) + } + }) +} + +func TestApply_RejectedRevisionMemory(t *testing.T) { + h := newRemoteHarness(t, remoteConfig("apply_safe", "")) + h.remote.publish(1, `{"plugins":{"ssh":{"source":"ghcr.io/other/plugin:v1"}}}`) + mustStartup(t, h.rc) + if r := h.remote.lastReport(t); r.Reason != agentconfig.ReasonUnsafeChanges { + t.Fatalf("expected unsafe-changes, got %s", r.Reason) + } + calls := h.pf.callCount() + h.poll(t) + h.poll(t) + if h.pf.callCount() != calls { + t.Fatalf("a remembered rejection must not be re-prepared (%d prefetches)", h.pf.callCount()-calls) + } + // A base edit that makes the source "already used" re-classifies and applies. + h.writeConfig(t, remoteConfig("apply_safe", "")+` + other: + source: ghcr.io/other/plugin:v1 +`) + h.rc.reconcile(context.Background(), triggerFile) + if next := h.rc.takePending(); next == nil || next.overlay == nil || next.overlay.Revision != 1 { + t.Fatalf("expected revision 1 to apply after the base edit, got %+v", next) + } +} + +// TestApply_RememberedRejectionReportedAfterRestart: a restart that skips a remembered +// rejection (the fetch answers 304) still reports it, and a later re-prepare of the applied +// overlay (a file edit with the same fingerprint) does not clear it. +func TestApply_RememberedRejectionReportedAfterRestart(t *testing.T) { + h := newRemoteHarness(t, remoteConfig("apply_safe", "")) + h.remote.publish(1, `{"plugins":{"ssh":{"schedule":"*/5 * * * *"}}}`) + mustStartup(t, h.rc) + h.remote.publish(2, `{"plugins":{"ssh":{"source":"ghcr.io/other/plugin:v1"}}}`) + h.poll(t) + before := h.remote.lastReport(t) + if before.Status != agentconfig.StatusRejected || len(before.Unsafe) == 0 { + t.Fatalf("expected revision 2 to be rejected with unsafe changes, got %+v", before) + } + + assertRejected := func(t *testing.T, r agentconfig.Report) { + t.Helper() + if r.Status != agentconfig.StatusRejected || r.Reason != agentconfig.ReasonUnsafeChanges { + t.Fatalf("expected rejected/unsafe-changes, got %s/%s", r.Status, r.Reason) + } + if r.AttemptedRevision == nil || *r.AttemptedRevision != 2 { + t.Fatalf("expected attempted revision 2, got %v", r.AttemptedRevision) + } + if r.AppliedRevision == nil || *r.AppliedRevision != 1 { + t.Fatalf("expected applied revision 1, got %v", r.AppliedRevision) + } + if derefString(r.Error) != derefString(before.Error) { + t.Fatalf("expected error %q, got %q", derefString(before.Error), derefString(r.Error)) + } + if !reflect.DeepEqual(r.Unsafe, before.Unsafe) { + t.Fatalf("expected the persisted unsafe changes %+v, got %+v", before.Unsafe, r.Unsafe) + } + } + + gets, reports := h.remote.getCount(), h.remote.reportCount() + h.rc = h.newReconciler() + calls := h.pf.callCount() + if a := mustStartup(t, h.rc); a.overlay == nil || a.overlay.Revision != 1 { + t.Fatalf("expected the applied revision 1, got %+v", a.overlay) + } + if h.remote.getCount() != gets+1 || h.remote.gets[len(h.remote.gets)-1] != h.remote.etag { + t.Fatal("the restart must fetch conditionally (304)") + } + if h.remote.reportCount() != reports+1 { + t.Fatal("the restart must send a report") + } + assertRejected(t, h.remote.lastReport(t)) + if got := h.pf.callCount() - calls; got != 1 { + t.Fatalf("only the applied revision may be prepared, got %d prepares", got) + } + + h.poll(t) + assertRejected(t, h.remote.lastReport(t)) + + // A comment changes the file but not its fingerprint: the applied overlay is re-prepared + // and the remembered rejection stays the outcome. + h.writeConfig(t, remoteConfig("apply_safe", "")+"# edited\n") + calls = h.pf.callCount() + h.rc.reconcile(context.Background(), triggerFile) + if h.pf.callCount() == calls { + t.Fatal("the file edit must re-prepare the applied revision") + } + assertRejected(t, h.remote.lastReport(t)) +} + +func TestStartupLadder(t *testing.T) { + t.Run("fetched wins", func(t *testing.T) { + h := newRemoteHarness(t, remoteConfig("apply_safe", "")) + h.remote.publish(2, `{"plugins":{"ssh":{"schedule":"*/5 * * * *"}}}`) + if a := mustStartup(t, h.rc); a.overlay == nil || a.overlay.Revision != 2 { + t.Fatalf("got %+v", a.overlay) + } + }) + t.Run("rejected fetched falls back to applied", func(t *testing.T) { + h := newRemoteHarness(t, remoteConfig("apply_safe", "")) + h.remote.publish(1, `{"plugins":{"ssh":{"schedule":"*/5 * * * *"}}}`) + mustStartup(t, h.rc) + h.remote.publish(2, `{"plugins":{"ssh":{"source":"ghcr.io/other/plugin:v1"}}}`) + a := mustStartup(t, h.newReconciler()) + if a.overlay == nil || a.overlay.Revision != 1 { + t.Fatalf("expected the applied revision 1, got %+v", a.overlay) + } + r := h.remote.lastReport(t) + if r.Status != agentconfig.StatusRejected || *r.AttemptedRevision != 2 || *r.AppliedRevision != 1 { + t.Fatalf("unexpected report %+v", r) + } + }) + t.Run("failing fetched and applied fall back to the file", func(t *testing.T) { + h := newRemoteHarness(t, remoteConfig("apply_safe", "")) + h.remote.publish(1, `{"plugins":{"ssh":{"schedule":"*/5 * * * *"}}}`) + mustStartup(t, h.rc) + h.remote.publish(2, `{"plugins":{"ssh":{"schedule":"*/7 * * * *"}}}`) + h.pf.setErr(errors.New("registry down")) + restarted := h.newReconciler() + if _, err := restarted.startup(context.Background()); err == nil { + t.Fatal("when even the file cannot be prepared startup must fail") + } + h.pf.setErr(nil) + }) + t.Run("unusable file fails", func(t *testing.T) { + h := newRemoteHarness(t, "api: {}\n") + if _, err := h.rc.startup(context.Background()); err == nil { + t.Fatal("expected an error") + } + }) +} + +func TestOneShot_FetchApplyReportRun(t *testing.T) { + h := newRemoteHarness(t, strings.Replace(remoteConfig("apply_safe", ""), "daemon: true", "daemon: false", 1)) + h.remote.publish(1, `{"plugins":{"ssh":{"schedule":"*/5 * * * *"}}}`) + active, err := h.rc.startup(context.Background()) + if err != nil { + t.Fatal(err) + } + if h.remote.reportCount() != 1 || h.remote.lastReport(t).Daemon { + t.Fatalf("one-shot must report once with daemon=false before running") + } + runs := 0 + err = h.rc.run(active, func(_ context.Context, cfg *agentConfig) error { + runs++ + if *cfg.Plugins["ssh"].Schedule != "*/5 * * * *" { + t.Fatalf("the overlay was not applied") + } + return nil + }) + if err != nil || runs != 1 { + t.Fatalf("one-shot must run once and exit: runs=%d err=%v", runs, err) + } +} + +func TestReconciler_RaceInterleavedFileEventsAndPolls(t *testing.T) { + h := newRemoteHarness(t, remoteConfig("apply_safe", "")) + h.remote.publish(1, `{"plugins":{"ssh":{"schedule":"*/1 * * * *"}}}`) + h.rc.now = time.Now + active, err := h.rc.startup(context.Background()) + if err != nil { + t.Fatal(err) + } + var mu sync.Mutex + running := 0 + maxRunning := 0 + var last *agentConfig + stop := make(chan struct{}) + done := make(chan error, 1) + go func() { + done <- h.rc.run(active, func(ctx context.Context, cfg *agentConfig) error { + mu.Lock() + running++ + maxRunning = max(maxRunning, running) + last = cfg + mu.Unlock() + defer func() { + mu.Lock() + running-- + mu.Unlock() + }() + select { + case <-ctx.Done(): + return nil + case <-stop: + return errStopRun + } + }) + }() + t.Cleanup(func() { + // Stop rc.run (it returns once the run and its fallback both stop) so the goroutine + // does not leak into other tests. + close(stop) + select { + case <-done: + case <-time.After(5 * time.Second): + t.Error("rc.run did not return") + } + }) + + ctx, cancel := context.WithCancel(context.Background()) + defer cancel() + for i := 2; i <= 200; i++ { + h.remote.publish(int64(i), fmt.Sprintf(`{"plugins":{"ssh":{"schedule":"*/%d * * * *"}}}`, i%59+1)) + h.rc.reconcile(ctx, triggerPoll) + if i%10 == 0 { + h.rc.reconcile(ctx, triggerFile) + } + } + want := fmt.Sprintf("*/%d * * * *", 200%59+1) + deadline := time.Now().Add(5 * time.Second) + for time.Now().Before(deadline) { + mu.Lock() + ok := last != nil && *last.Plugins["ssh"].Schedule == want + mu.Unlock() + if ok { + break + } + time.Sleep(5 * time.Millisecond) + } + mu.Lock() + defer mu.Unlock() + if maxRunning != 1 { + t.Fatalf("expected exactly one active run at a time, saw %d", maxRunning) + } + if last == nil || *last.Plugins["ssh"].Schedule != want { + t.Fatalf("the last revision must win") + } +} + +// TestApply_NoOpRevisionAdoptedWithoutRestart: a revision whose effective config equals the +// running one is recorded as applied without a restart, and is not re-prepared every poll. +func TestApply_NoOpRevisionAdoptedWithoutRestart(t *testing.T) { + h := newRemoteHarness(t, remoteConfig("apply_safe", "")) + h.remote.publish(1, `{"plugins":{"ssh":{"schedule":"*/5 * * * *"}}}`) + first := mustStartup(t, h.rc) + + // verbosity 0 equals the base: the effective config does not change. + h.remote.publish(2, `{"plugins":{"ssh":{"schedule":"*/5 * * * *"}},"verbosity":0}`) + calls := h.pf.callCount() + var cur *candidate + for i := 0; i < 3; i++ { + cur = h.poll(t) + } + if got := h.pf.callCount() - calls; got != 1 { + t.Fatalf("a no-op revision must be prepared once, got %d prefetches", got) + } + if cur.runtime != first.runtime { + t.Fatal("a no-op revision must not restart the running configuration") + } + if rev := cur.runtime.syncInfo().AppliedRevision; rev != 2 { + t.Fatalf("the heartbeat/evidence must show applied revision 2, got %d", rev) + } + if r := h.remote.lastReport(t); r.Status != agentconfig.StatusApplied || r.AppliedRevision == nil || *r.AppliedRevision != 2 { + t.Fatalf("the report must show applied revision 2, got %+v", r) + } + + // A comment-only file edit is adopted the same way. + h.writeConfig(t, "# a comment\n"+remoteConfig("apply_safe", "")) + h.rc.reconcile(context.Background(), triggerFile) + calls = h.pf.callCount() + for i := 0; i < 3; i++ { + cur = h.poll(t) + } + if h.pf.callCount() != calls || cur.runtime != first.runtime { + t.Fatalf("a comment-only edit must not be re-prepared every poll (%d prefetches) nor restart", h.pf.callCount()-calls) + } +} + +// TestApply_RejectedAppliedOverlayNotRePrepared: after a file edit makes the applied overlay +// invalid, the remembered rejection keeps last-known-good instead of re-preparing every poll. +func TestApply_RejectedAppliedOverlayNotRePrepared(t *testing.T) { + h := newRemoteHarness(t, remoteConfig("apply_safe", `"host"`)) + var logs bytes.Buffer + h.rc.logger = hclog.New(&hclog.LoggerOptions{Output: &logs, Level: hclog.Warn}) + h.remote.publish(1, `{"plugins":{"ssh":{"config":{"host":"other"}}}}`) + first := mustStartup(t, h.rc) + if first.overlay == nil { + t.Fatal("revision 1 must apply while host is overridable") + } + + h.writeConfig(t, remoteConfig("apply_safe", "")) // host is no longer overridable + h.rc.reconcile(context.Background(), triggerFile) + if r := h.remote.lastReport(t); r.Status != agentconfig.StatusRejected || r.Reason != agentconfig.ReasonUnsafeChanges { + t.Fatalf("expected rejected/unsafe-changes, got %s/%s", r.Status, r.Reason) + } + prepares := strings.Count(logs.String(), "Could not apply the configuration") + for i := 0; i < 3; i++ { + if got := h.poll(t); got != first { + t.Fatal("the last-known-good configuration must keep running") + } + } + if got := strings.Count(logs.String(), "Could not apply the configuration") - prepares; got != 0 { + t.Fatalf("a remembered rejection must not be re-prepared, got %d prepares", got) + } +} + +// TestApply_EmptyETagDoesNotBlockLaterRevisions: without an ETag, revisions are told apart by +// revision + overlay bytes. +func TestApply_EmptyETagDoesNotBlockLaterRevisions(t *testing.T) { + h := newRemoteHarness(t, remoteConfig("apply_safe", "")) + h.remote.publish(1, `{"plugins":{"ssh":{"source":"ghcr.io/other/plugin:v1"}}}`) + h.remote.etag = "" + mustStartup(t, h.rc) + if r := h.remote.lastReport(t); r.Reason != agentconfig.ReasonUnsafeChanges { + t.Fatalf("expected unsafe-changes, got %s", r.Reason) + } + h.remote.publish(2, `{"plugins":{"ssh":{"schedule":"*/5 * * * *"}}}`) + h.remote.etag = "" + if got := h.poll(t); got.overlay == nil || got.overlay.Revision != 2 { + t.Fatalf("revision 2 must apply despite the empty ETag, got %+v", got.overlay) + } +} + +// TestApply_RunFailureNotifiesAfterFallbackIsBound: onRunFailed sees the fallback as current, +// so the applied overlay reverts to it and the report describes it. +func TestApply_RunFailureNotifiesAfterFallbackIsBound(t *testing.T) { + h := newRemoteHarness(t, remoteConfig("apply_safe", "")) + h.remote.publish(1, `{"plugins":{"ssh":{"schedule":"*/5 * * * *"}}}`) + active, err := h.rc.startup(context.Background()) + if err != nil { + t.Fatal(err) + } + stop := make(chan struct{}) + done := make(chan error, 1) + go func() { + done <- h.rc.run(active, func(ctx context.Context, cfg *agentConfig) error { + if *cfg.Plugins["ssh"].Schedule == "*/7 * * * *" { + return errors.New("failed to start") + } + select { + case <-ctx.Done(): + return nil + case <-stop: + return errStopRun + } + }) + }() + h.remote.publish(2, `{"plugins":{"ssh":{"schedule":"*/7 * * * *"}}}`) + h.rc.reconcile(context.Background(), triggerPoll) + + var failedRun *candidate + select { + case failedRun = <-h.rc.runFailed: + case <-time.After(5 * time.Second): + t.Fatal("the run failure was never notified") + } + if cur := h.rc.current(); cur == nil || cur.overlay == nil || cur.overlay.Revision != 1 { + t.Fatalf("the fallback must be bound before the notification, current is %+v", cur) + } + h.rc.onRunFailed(context.Background(), failedRun) + if a := h.rc.cache.Applied; a == nil || a.Revision != 1 { + t.Fatalf("the applied overlay must revert to revision 1, got %+v", a) + } + if r := h.remote.lastReport(t); r.Status != agentconfig.StatusFailed || r.AppliedRevision == nil || *r.AppliedRevision != 1 { + t.Fatalf("the failure report must describe the fallback, got %+v", r) + } + close(stop) + select { + case <-done: + case <-time.After(5 * time.Second): + t.Fatal("run did not return") + } +} + +// TestApply_PrefetchIsBounded: a hanging registry does not stall the reconciler. +func TestApply_PrefetchIsBounded(t *testing.T) { + old := prepareNetworkTimeout + prepareNetworkTimeout = 50 * time.Millisecond + t.Cleanup(func() { prepareNetworkTimeout = old }) + + h := newRemoteHarness(t, remoteConfig("apply_safe", "")) + h.remote.publish(1, `{"plugins":{"ssh":{"schedule":"*/5 * * * *"}}}`) + first := mustStartup(t, h.rc) + h.pf.setBlock(true) + h.remote.publish(2, `{"plugins":{"ssh":{"schedule":"*/7 * * * *"}}}`) + start := time.Now() + if got := h.poll(t); got != first { + t.Fatal("a hung download must keep the running config") + } + if time.Since(start) > 5*time.Second { + t.Fatalf("prepare was not bounded: %s", time.Since(start)) + } + if r := h.remote.lastReport(t); r.Status != agentconfig.StatusFailed || r.Reason != agentconfig.ReasonDownloadFailed { + t.Fatalf("expected failed/download-failed, got %s/%s", r.Status, r.Reason) + } +} diff --git a/cmd/report.go b/cmd/report.go index 968681c..d2dbc4e 100644 --- a/cmd/report.go +++ b/cmd/report.go @@ -108,6 +108,9 @@ func (rc *reconciler) maybeReport(ctx context.Context, active *candidate, outcom switch { case err == nil: rc.report = reportState{fingerprint: fingerprint, sentAt: rc.now()} + if outcome != nil && outcome.Reason == agentconfig.ReasonCacheCorrupt && rc.lastOutcome == outcome { + rc.lastOutcome = nil // reported once + } case errors.As(err, &statusErr) && statusErr.StatusCode == http.StatusConflict: rc.report.sendFailed = true rc.reportBackoffUntil = rc.now().Add(reportConflictBackoff) @@ -131,20 +134,23 @@ func (rc *reconciler) buildReport(active *candidate, outcome *applyError, rcfg a hostname, _ := os.Hostname() opts := active.base.redactOpts() report := agentconfig.Report{ - Hostname: truncateString(hostname, 255), - AgentVersion: truncateString(agentVersion, 64), - Mode: rcfg.Mode, - Daemon: active.runtime.Daemon, - Base: marshalRaw(agentconfig.Redact(active.base.declared, opts...)), - Effective: marshalRaw(agentconfig.Redact(active.declared, opts...)), - EffectiveDigest: active.digest, - Warnings: active.warnings, - RemoteConfig: &rcfg, - Plugins: active.plugins, + Hostname: truncateString(hostname, 255), + AgentVersion: truncateString(agentVersion, 64), + Mode: rcfg.Mode, + Daemon: active.runtime.Daemon, + AppliedRevision: active.appliedRevision(), + AttemptedRevision: rc.attempted, + Base: marshalRaw(agentconfig.Redact(active.base.declared, opts...)), + Effective: marshalRaw(agentconfig.Redact(active.declared, opts...)), + EffectiveDigest: active.digest, + Warnings: active.warnings, + RemoteConfig: &rcfg, + Plugins: active.plugins, } switch { case !isApplyMode(rcfg.Mode): report.Status = agentconfig.StatusNotApplicable + report.AttemptedRevision = nil case outcome != nil: report.Status = outcome.Status report.Reason = outcome.Reason @@ -153,6 +159,7 @@ func (rc *reconciler) buildReport(active *candidate, outcome *applyError, rcfg a msg = outcome.Err.Error() } report.Error = &msg + report.Unsafe = outcome.Unsafe default: report.Status = agentconfig.StatusApplied } diff --git a/docs/configuration.md b/docs/configuration.md index cca9f12..285acb5 100644 --- a/docs/configuration.md +++ b/docs/configuration.md @@ -154,7 +154,8 @@ periodic agent evidence while keeping `emit_on_run_completion` behavior enabled. no expiry. Set `agent_evidence.emit_on_run_completion` to `false` to disable immediate agent evidence on run completion and startup failures while leaving periodic daemon evidence controlled by `interval`. -The `log_level` is one of the following, defaulting to `0` if not specified: +The `log_level` is one of the following, defaulting to `0` if not specified (a remote overlay's `verbosity` wins over +the `-v` flag): - 0: Shows all ERROR, WARN and INFO - 1: Shows all of 0 plus DEBUG logs - 2: Shows all of 1 plus TRACE logs @@ -172,10 +173,12 @@ A disabled plugin gets no schedule, no download and no run state, but it stays i ## Typing of plugin values Values in the config file keep viper's weak typing exactly as before: `collect_ip_allow_list: false` reaches the plugin -as `"0"`, `account_id: 123456789012` as `"123456789012"`, and `port: 22` as `"22"`. +as `"0"`, `account_id: 123456789012` as `"123456789012"`, and `port: 22` as `"22"`. Values set by a remote overlay +must already be strings: a remote `port: 2222` (a number) is rejected with `invalid-type` (R27, R51). Viper lowercases keys and splits them on dots. Plugin names and config keys in the file are therefore lowercase and -cannot contain dots. +cannot contain dots; a remote overlay that uses `GitHub` addresses a different plugin than the file's `github`. Plugin +names an overlay introduces must match `^[a-z0-9][a-z0-9_-]{0,62}$` (R28). ## Tolerated file problems @@ -183,11 +186,12 @@ A plugin `schedule` in the file that does not parse does not stop the agent: tha and the problem is logged and reported as a warning (R34). A few other file values that always loaded are also only warnings, and are kept unchanged: a negative `verbosity` (`-1` logs WARN and above) and a literal `${env:...}` outside `plugins.*.config`. Every other invalid value in the file (for example a missing `api.url`) still fails startup, and on -a live reload the agent keeps running its last good configuration. +a live reload the agent keeps running its last good configuration. Values set by a remote overlay are always validated +strictly. ## Remote configuration -An agent with `api.auth` credentials reports the configuration it runs to the API. The `remote_config` block +An agent with `api.auth` credentials can pick up a configuration overlay stored in the API. The `remote_config` block controls it. It is **set locally only** (file, host environment, CLI flags), never remotely (R30): ```yaml @@ -200,7 +204,7 @@ remote_config: ``` Defaults (R29): `mode` is `report` when `api.auth` is set and `off` otherwise (no credentials always forces -`off`); `poll_interval` is `60s`; `trusted_sources` and `overridable_config_flags` are empty; +`off`), so an agent applies an overlay only when `mode` is set to `apply_safe` or `apply_all`; `poll_interval` is `60s`; `trusted_sources` and `overridable_config_flags` are empty; `allow_local_sources` is `false`. `CCF_REMOTE_CONFIG_MODE` sets the mode even when the file has no `remote_config` block. @@ -208,6 +212,37 @@ Defaults (R29): `mode` is `report` when `api.auth` is set and `off` otherwise (n |---|---| | `off` | No report, no fetch. The heartbeat carries no configuration fields. | | `report` | The agent reports its configuration (status `not-applicable`) but never fetches an overlay. | +| `apply_safe` | The agent fetches the overlay and applies it only when every change is safe (table below). | +| `apply_all` | The agent applies safe and unsafe changes. Forbidden changes are still rejected. | + +A change is classified as follows (the agent is the authority; the API preview uses the same rules): + +| Change | Class | +|---|---| +| `api`, `daemon` or `remote_config` in the overlay | **forbidden** (the whole revision is rejected in every mode) | +| `verbosity`, `agent_evidence.*` | safe | +| a plugin's `schedule`, `labels`, `policy_behavior`, `protocol_version`, `enabled`, `policy_data` | safe | +| removing a plugin or a policy entry | safe | +| a plugin source or policy entry already used by the file | safe | +| a new source matching `trusted_sources` | safe | +| a new OCI source not in `trusted_sources` | unsafe | +| a new local path | forbidden, unless `apply_all` with `allow_local_sources: true` (then unsafe) | +| a `plugins.

.config.` change matching `overridable_config_flags` (`key`, `plugin:key` or `*`) | safe | +| any other plugin config change | unsafe | +| a new `${env:NAME}` reference | unsafe (`CCF_API_AUTH_*`: forbidden) | + +A rejected or failed revision never interrupts the running configuration: the agent prepares the whole new +configuration (validation, downloads) first and swaps only when it is ready. Every outcome is reported to the API with +a reason (`unsafe-changes`, `forbidden-changes`, `invalid-config`, `invalid-type`, `unknown-field`, +`download-failed`, `cache-corrupt`, `internal`). When a new configuration is applied, +in-flight plugin runs get up to 5 minutes to finish (R33). Evidence produced under an overlay carries the prop +`agent-config-revision` (namespace `https://compliance-framework.github.io/ns`). + +The agent caches the last fetched and applied overlay in `/remote-config.json` (mode 0600, bound to `api.url` +and `api.auth.client_id`), so it keeps running the last good overlay when the API is unreachable. At startup it tries, +in order: the freshly fetched overlay, the cached applied overlay, the file alone. Only an unusable file stops the agent. +A fetched overlay already rejected for the same file is skipped, and its rejection (with the unsafe changes) is reported +again, so the instance still shows as rejected after a restart. When a configuration report is too large for the API, the agent drops its `base` document and marks it truncated; the effective document and digest are kept. @@ -222,8 +257,9 @@ nothing is gated on it. ## State directory and instance ID Each agent instance keeps state in `.compliance-framework/state//`, relative to the working directory, where -`` is derived from the absolute path of the config file (R31): the instance ID (`instance-id`). The OCI download -caches in `.compliance-framework/plugins` and `.compliance-framework/policies` are shared. +`` is derived from the absolute path of the config file (R31): the instance ID (`instance-id`) and the remote +configuration cache. The OCI download caches in `.compliance-framework/plugins` and +`.compliance-framework/policies` are shared. | Setting | Flag | Environment | |---|---|---| diff --git a/docs/running_as_a_service.md b/docs/running_as_a_service.md index 665253d..79055e5 100644 --- a/docs/running_as_a_service.md +++ b/docs/running_as_a_service.md @@ -100,7 +100,7 @@ EOF ``` `WorkingDirectory` and `StateDirectory` give the agent a persistent place for its download caches and its -per-instance state (`.compliance-framework/state/...`: the instance ID). Without +per-instance state (`.compliance-framework/state/...`: the instance ID and the remote configuration cache). Without them the agent writes relative to `/`. The state directory must persist across restarts, otherwise every restart registers a new instance. See [State directory and instance ID](configuration.md#state-directory-and-instance-id). From 10835ca36c8e92c24d3e0446117ebd0df020cb94 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:46 -0300 Subject: [PATCH 2/4] refactor(agent): stamp the revision with runner's agent-owned prop constants The runner drops a plugin's agent-config-revision by these constants, so the agent stamps it with the same ones. Co-Authored-By: Claude Opus 5.5 --- cmd/agent.go | 10 ++-------- 1 file changed, 2 insertions(+), 8 deletions(-) diff --git a/cmd/agent.go b/cmd/agent.go index 664cf73..23fc938 100644 --- a/cmd/agent.go +++ b/cmd/agent.go @@ -191,12 +191,6 @@ const DefaultProtocolVersion int32 = 1 const RunnerV2ProtocolVersion int32 = 2 const AnnotationProtocolVersionKey = "org.ccf.plugin.protocol.version" -// CCFPropNamespace is the OSCAL prop namespace of CCF props. -const CCFPropNamespace = "https://compliance-framework.github.io/ns" - -// configRevisionPropName stamps evidence with the applied remote configuration revision (R38). -const configRevisionPropName = "agent-config-revision" - // daemonCronStopTimeout bounds the cron stop on SIGINT/SIGTERM before plugins are killed, // also when the signal arrives during a reload drain (R33). var daemonCronStopTimeout = 30 * time.Second @@ -1651,8 +1645,8 @@ func configRevisionProps(config *agentConfig) []sdktypes.Property { return nil } return []sdktypes.Property{{ - Ns: CCFPropNamespace, - Name: configRevisionPropName, + Ns: runner.PropNamespace, + Name: runner.PropConfigRevision, Value: strconv.FormatInt(meta.AppliedRevision, 10), }} } From 271355f1b060f791852f12e807df1a911673596d Mon Sep 17 00:00:00 2001 From: "ccf-lisa[bot]" <286799724+ccf-lisa[bot]@users.noreply.github.com> Date: Tue, 6 Oct 2026 12:56:26 -0300 Subject: [PATCH 3/4] fix(agent): key no-ETag overlays on canonical bytes A rejection of an overlay served without an ETag is keyed by revision and sha256(overlay), and sameOverlay compares overlay bytes. fetch now stores the body in agentstate.CanonicalOverlay form, the form the cache saves and loads, so a remembered no-ETag rejection is still recognised after a restart, on the cached copy and on a fresh 200 body. Co-Authored-By: Claude Opus 5.5 --- cmd/reconciler.go | 11 ++++++++++- cmd/remote_test.go | 46 ++++++++++++++++++++++++++++++++++++++++++++++ 2 files changed, 56 insertions(+), 1 deletion(-) diff --git a/cmd/reconciler.go b/cmd/reconciler.go index 9d2f9f6..082d886 100644 --- a/cmd/reconciler.go +++ b/cmd/reconciler.go @@ -402,6 +402,8 @@ func sameOverlay(a, b *agentstate.OverlayRecord) bool { // overlayKey identifies an overlay for the rejected memory and the failed backoff: the raw // ETag, or revision + sha256(overlay) when a response carried no ETag (a stripping proxy), so // one rejection never blocks every later revision. nil (the file only) has its own key. +// Overlay bytes are in agentstate.CanonicalOverlay form (fetch and the cache both ensure it), +// so the key is the same for a fresh body and its cached copy. func overlayKey(rec *agentstate.OverlayRecord) string { switch { case rec == nil: @@ -952,10 +954,17 @@ func (rc *reconciler) fetch(ctx context.Context) { switch { case res.NotModified: case res.Document != nil: + // The canonical bytes are the ones a cached copy has after a restart, so the no-ETag + // keys (overlayKey, sameOverlay, the rejected memory) match across one. + overlay, err := agentstate.CanonicalOverlay(res.Document.Overlay) + if err != nil { + // Not JSON: kept as received; ValidateOverlay rejects it. + overlay = append(json.RawMessage(nil), res.Document.Overlay...) + } rc.cache.Fetched = &agentstate.OverlayRecord{ Revision: res.Document.Revision, ETag: res.ETag, - Overlay: append(json.RawMessage(nil), res.Document.Overlay...), + Overlay: overlay, FetchedAt: rc.now().UTC(), } rc.saveCache() diff --git a/cmd/remote_test.go b/cmd/remote_test.go index 943a456..913fa40 100644 --- a/cmd/remote_test.go +++ b/cmd/remote_test.go @@ -1005,6 +1005,52 @@ func TestApply_EmptyETagDoesNotBlockLaterRevisions(t *testing.T) { } } +// TestApply_EmptyETagRejectionSurvivesRestart: without an ETag a rejection is keyed by revision +// + sha256(overlay). The cache re-indents the overlay it writes, so the key must hash canonical +// bytes for the rejection to be remembered after a restart, whether the restart runs on the +// cached copy (the API is down) or on a fresh 200 body. +func TestApply_EmptyETagRejectionSurvivesRestart(t *testing.T) { + h := newRemoteHarness(t, remoteConfig("apply_safe", "")) + h.remote.publish(1, `{"plugins":{"ssh":{"schedule":"*/5 * * * *"}}}`) + h.remote.etag = "" + mustStartup(t, h.rc) + // Whitespace and <, & that the cache would rewrite. + h.remote.publish(2, "{\n \"plugins\": {\n \"ssh\": {\"source\": \"ghcr.io/other/plugin:v1\", \"config\": {\"q\": \"a Date: Tue, 6 Oct 2026 12:56:26 -0300 Subject: [PATCH 4/4] docs(config): classify enabled changes as api v0.21.0 does Disabling a plugin is safe; re-enabling one the file disables is unsafe unless its source is trusted, and its kept policies, local source and ${env:} references are classified as new. Only enabled plugins' sources count as already used. Co-Authored-By: Claude Opus 5.5 --- docs/configuration.md | 6 ++++-- 1 file changed, 4 insertions(+), 2 deletions(-) diff --git a/docs/configuration.md b/docs/configuration.md index 285acb5..88f2753 100644 --- a/docs/configuration.md +++ b/docs/configuration.md @@ -221,9 +221,11 @@ A change is classified as follows (the agent is the authority; the API preview u |---|---| | `api`, `daemon` or `remote_config` in the overlay | **forbidden** (the whole revision is rejected in every mode) | | `verbosity`, `agent_evidence.*` | safe | -| a plugin's `schedule`, `labels`, `policy_behavior`, `protocol_version`, `enabled`, `policy_data` | safe | +| a plugin's `schedule`, `labels`, `policy_behavior`, `protocol_version`, `policy_data` | safe | +| disabling a plugin (`enabled: false`) | safe | +| re-enabling a plugin the file disables | unsafe, unless its source matches `trusted_sources` (then safe); its kept policies, local source and `${env:}` references are classified as if new (so a local one is forbidden unless `apply_all` with `allow_local_sources: true`) | | removing a plugin or a policy entry | safe | -| a plugin source or policy entry already used by the file | safe | +| a plugin source or policy entry already used by an **enabled** plugin in the file (a disabled plugin's sources do not count) | safe | | a new source matching `trusted_sources` | safe | | a new OCI source not in `trusted_sources` | unsafe | | a new local path | forbidden, unless `apply_all` with `allow_local_sources: true` (then unsafe) |