diff --git a/cmd/agent.go b/cmd/agent.go index 4c6aa29..678235f 100644 --- a/cmd/agent.go +++ b/cmd/agent.go @@ -27,6 +27,7 @@ import ( "github.com/robfig/cron/v3" "github.com/compliance-framework/agent/internal" + "github.com/compliance-framework/agent/internal/agentstate" "github.com/compliance-framework/agent/runner" "github.com/compliance-framework/api/sdk" sdktypes "github.com/compliance-framework/api/sdk/types" @@ -191,6 +192,9 @@ 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("instance-id", "", "Pin this instance's UUID (not persisted); overrides CCF_INSTANCE_ID") + return agentCmd } @@ -306,7 +310,21 @@ func agentRunner(cmd *cobra.Command, args []string) error { Level: hclog.Debug, }) - agentRun := NewAgentRunner() + stateDir, stateDirSource, err := stateDirFrom(cmd, configPath) + if err != nil { + return err + } + idOverride, err := instanceIDOverride(cmd) + if err != nil { + return err + } + store := agentstate.Open(stateDir, logger) + id, persisted := store.InstanceID(idOverride) + // R52: the default state dir depends on the absolute config path, so moving the config + // file silently creates a new instance. Say where state lives. + logger.Info("Agent state", "state_dir", stateDir, "state_dir_source", stateDirSource, "instance_id", id.String(), "instance_id_persisted", persisted) + + agentRun := NewAgentRunner(WithInstanceID(id)) ctx, configCancel := context.WithCancel(context.Background()) defer configCancel() @@ -348,6 +366,39 @@ func agentRunner(cmd *cobra.Command, args []string) error { return nil } +// stateDirFrom resolves the state directory: --state-dir, then CCF_STATE_DIR, then the +// default derived from the absolute config path (R31). It also returns where it came from. +func stateDirFrom(cmd *cobra.Command, configPath string) (dir string, source string, err error) { + if flag := cmd.Flags().Lookup("state-dir"); flag != nil && strings.TrimSpace(flag.Value.String()) != "" { + dir, err = filepath.Abs(strings.TrimSpace(flag.Value.String())) + return dir, "flag", err + } + if env := strings.TrimSpace(os.Getenv("CCF_STATE_DIR")); env != "" { + dir, err = filepath.Abs(env) + return dir, "env", err + } + dir, err = agentstate.DefaultDir(configPath) + return dir, "default(config-path)", err +} + +// instanceIDOverride returns --instance-id or CCF_INSTANCE_ID. An override that is not a +// UUID is an error: silently ignoring it would register a different instance. +func instanceIDOverride(cmd *cobra.Command) (string, error) { + value, source := "", "" + if flag := cmd.Flags().Lookup("instance-id"); flag != nil && strings.TrimSpace(flag.Value.String()) != "" { + value, source = strings.TrimSpace(flag.Value.String()), "--instance-id" + } else if env := strings.TrimSpace(os.Getenv("CCF_INSTANCE_ID")); env != "" { + value, source = env, "CCF_INSTANCE_ID" + } + if value == "" { + return "", nil + } + if _, err := uuid.Parse(value); err != nil { + return "", fmt.Errorf("%s must be a UUID: %w", source, err) + } + return value, nil +} + type AgentRunner struct { logger hclog.Logger stateMu sync.RWMutex @@ -363,25 +414,45 @@ type AgentRunner struct { downloadGroup singleflight.Group fetchAnnotations func(ctx context.Context, source string, option ...remote.Option) (map[string]string, error) runPluginFunc func(ctx context.Context, name string, pluginConfig *agentPlugin) error + sendHeartbeatFunc func(ctx context.Context, instanceID uuid.UUID) error pluginRunMu sync.RWMutex pluginRuns map[string]pluginRunRecord firstAgentEvidenceSendStarted bool queryBundles []*rego.Rego + + // instanceID is this agent instance's stable ID (R31); set once at construction. + instanceID uuid.UUID } -func NewAgentRunner() *AgentRunner { - return &AgentRunner{ +// AgentRunnerOption configures an AgentRunner. +type AgentRunnerOption func(*AgentRunner) + +// WithInstanceID sets the instance ID the heartbeat and config reports use. +func WithInstanceID(id uuid.UUID) AgentRunnerOption { + return func(ar *AgentRunner) { ar.instanceID = id } +} + +func NewAgentRunner(opts ...AgentRunnerOption) *AgentRunner { + ar := &AgentRunner{ pluginLocations: map[string]string{}, policyLocations: map[string]string{}, activePluginClients: map[*plugin.Client]struct{}{}, pluginRuns: map[string]pluginRunRecord{}, fetchAnnotations: internal.GetAnnotations, httpClient: http.DefaultClient, + instanceID: uuid.New(), } + for _, opt := range opts { + opt(ar) + } + return ar } +// InstanceID returns the instance ID. +func (ar *AgentRunner) InstanceID() uuid.UUID { return ar.instanceID } + func (ar *AgentRunner) UpdateConfig(config *agentConfig) { logger := hclog.New(&hclog.LoggerOptions{ Name: "agent-runner", @@ -1089,9 +1160,13 @@ func (ar *AgentRunner) setupHeartbeatCron(ctx context.Context) (*cron.Cron, erro c := cron.New(cron.WithParser(cron.NewParser( cron.SecondOptional | cron.Minute | cron.Hour | cron.Dom | cron.Month | cron.Dow | cron.Descriptor, ))) - staticAgentUUID := uuid.New() + staticAgentUUID := ar.instanceID + sendHeartbeat := ar.SendHeartbeat + if ar.sendHeartbeatFunc != nil { + sendHeartbeat = ar.sendHeartbeatFunc + } _, err := c.AddFunc(fmt.Sprintf("%d * * * * *", staggeredSeconds), func() { - err := ar.SendHeartbeat(ctx, staticAgentUUID) + err := sendHeartbeat(ctx, staticAgentUUID) if err != nil { logger.Error("Failed to send heartbeat", "error", err, "uuid", staticAgentUUID.String()) } diff --git a/cmd/reconciler_test.go b/cmd/reconciler_test.go new file mode 100644 index 0000000..5de8019 --- /dev/null +++ b/cmd/reconciler_test.go @@ -0,0 +1,35 @@ +package cmd + +import ( + "context" + "sync" + "testing" + + "github.com/google/uuid" +) + +func TestHeartbeatUsesStableInstanceIDAcrossRuns(t *testing.T) { + id := uuid.New() + ar := NewAgentRunner(WithInstanceID(id)) + ar.UpdateConfig(newTestAgentConfig("http://example.test", nil)) + var seen []uuid.UUID + var mu sync.Mutex + ar.sendHeartbeatFunc = func(_ context.Context, got uuid.UUID) error { + mu.Lock() + seen = append(seen, got) + mu.Unlock() + return nil + } + for i := 0; i < 2; i++ { + c, err := ar.setupHeartbeatCron(context.Background()) + if err != nil { + t.Fatal(err) + } + for _, e := range c.Entries() { + e.Job.Run() + } + } + if len(seen) != 2 || seen[0] != id || seen[1] != id { + t.Fatalf("heartbeats used %v, want %s twice", seen, id) + } +} diff --git a/docs/configuration.md b/docs/configuration.md index 30b4fb7..712b0ea 100644 --- a/docs/configuration.md +++ b/docs/configuration.md @@ -183,3 +183,19 @@ A plugin `schedule` in the file that does not parse does not stop the agent: tha and the problem is logged 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. + +## 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. + +| Setting | Flag | Environment | +|---|---|---| +| State directory | `--state-dir` | `CCF_STATE_DIR` | +| Instance ID (a UUID; not persisted) | `--instance-id` | `CCF_INSTANCE_ID` | + +Because the default key depends on the config file's path, **moving or renaming the config file creates a new +instance** (R52). The agent logs the state directory, where it came from and the instance ID at startup. Containers and +Helm deployments should pin `CCF_STATE_DIR` to a mounted volume; ephemeral one-shot runs (CI, Kubernetes jobs) can +set `CCF_INSTANCE_ID` so repeated runs report as one instance. diff --git a/docs/running_as_a_service.md b/docs/running_as_a_service.md index 1dc7bfa..665253d 100644 --- a/docs/running_as_a_service.md +++ b/docs/running_as_a_service.md @@ -84,7 +84,9 @@ WantedBy=multi-user.target [Service] Type=notify -ExecStart=/usr/local/bin/ccf-agent agent -d +WorkingDirectory=/var/lib/ccf-agent +StateDirectory=ccf-agent +ExecStart=/usr/local/bin/ccf-agent agent -d -c /etc/ccf-agent/config.yaml KillMode=process Delegate=yes LimitNOFILE=1048576 @@ -97,6 +99,12 @@ RestartSec=5s 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 +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). + Now run the following command to reload the systemd configuration: ```bash @@ -135,7 +143,10 @@ TODO ## Running as a server/container -TODO +Mount a volume for the agent's state and pin it with `CCF_STATE_DIR`: the default state directory is derived from the +config file's absolute path, so a container that mounts its config elsewhere would otherwise get a new instance ID +(R52). In Kubernetes jobs and CI one-shot runs, set `CCF_INSTANCE_ID` to a fixed UUID so repeated runs report as one +instance; one-shot instances are pruned by the API after 24h. ## Running as a serverless process in AWS diff --git a/internal/agentstate/store.go b/internal/agentstate/store.go new file mode 100644 index 0000000..05781d5 --- /dev/null +++ b/internal/agentstate/store.go @@ -0,0 +1,180 @@ +// Package agentstate owns the agent's per-instance state directory: the stable instance ID +// (R31). +// +// Layout (the OCI download caches under .compliance-framework/{plugins,policies} are shared +// and unchanged): +// +// .compliance-framework/state// key = hex(sha256(abs config path))[:16]; dir 0700 +// instance-id 0644, UUID + "\n" +// +// The package is a leaf: it never imports cmd. +package agentstate + +import ( + "crypto/sha256" + "encoding/hex" + "fmt" + "os" + "path/filepath" + "strings" + "sync" + + "github.com/google/uuid" + "github.com/hashicorp/go-hclog" +) + +const ( + // StateRoot is the default parent of every per-config state directory, relative to the + // working directory like the download caches. + StateRoot = ".compliance-framework/state" + + instanceIDFile = "instance-id" +) + +// Store is one agent instance's state directory. A Store whose directory is not writable +// keeps working in memory: the agent never fails because it cannot persist state. +type Store struct { + dir string + logger hclog.Logger + writable bool + + mu sync.Mutex + id uuid.UUID + idSet bool + warned bool +} + +// DefaultDir returns the default state directory for a config file: StateRoot/, where key +// is the first 16 hex characters of sha256 of the absolute config path. Moving or renaming the +// config file therefore changes the directory and the instance ID (R52); containers should +// pin CCF_STATE_DIR. +func DefaultDir(configPath string) (string, error) { + abs, err := filepath.Abs(configPath) + if err != nil { + return "", err + } + sum := sha256.Sum256([]byte(abs)) + return filepath.Abs(filepath.Join(StateRoot, hex.EncodeToString(sum[:])[:16])) +} + +// Open prepares dir (MkdirAll 0700) and probes that it is writable. A failure is logged once +// as a WARN and the store continues in memory; it is never fatal. +func Open(dir string, logger hclog.Logger) *Store { + if logger == nil { + logger = hclog.NewNullLogger() + } + s := &Store{dir: dir, logger: logger.Named("state")} + if err := os.MkdirAll(dir, 0o700); err != nil { + s.logger.Warn("State directory is not usable; state will not persist across restarts", "dir", dir, "error", err) + return s + } + probe, err := os.CreateTemp(dir, ".probe-*") + if err != nil { + s.logger.Warn("State directory is not writable; state will not persist across restarts", "dir", dir, "error", err) + return s + } + name := probe.Name() + _ = probe.Close() + _ = os.Remove(name) + s.writable = true + return s +} + +// Dir returns the state directory. +func (s *Store) Dir() string { return s.dir } + +// Writable reports whether the directory accepted a probe write at Open. +func (s *Store) Writable() bool { return s.writable } + +// InstanceID returns this instance's stable ID and whether it is persisted: +// 1. a valid override (flag or CCF_INSTANCE_ID) wins and is not persisted; +// 2. otherwise the instance-id file; +// 3. otherwise a new ID, written atomically (a corrupt file is replaced); +// 4. if the store is not writable, the new ID lives in memory only (one WARN). +// +// The result is memoized: every call on one Store returns the same ID. +func (s *Store) InstanceID(override string) (uuid.UUID, bool) { + s.mu.Lock() + defer s.mu.Unlock() + + if o := strings.TrimSpace(override); o != "" { + if id, err := uuid.Parse(o); err == nil { + return id, false + } + s.logger.Warn("Ignoring invalid instance ID override", "value", o) + } + + path := filepath.Join(s.dir, instanceIDFile) + if s.idSet { + return s.id, s.writable && fileExists(path) + } + + if raw, err := os.ReadFile(path); err == nil { + if id, err := uuid.Parse(strings.TrimSpace(string(raw))); err == nil { + s.id, s.idSet = id, true + return id, true + } + s.logger.Warn("Instance ID file is corrupt; replacing it", "path", path) + } + + id := uuid.New() + s.id, s.idSet = id, true + if !s.writable { + s.warnOnce("Instance ID is kept in memory only; a restart creates a new instance", "dir", s.dir) + return id, false + } + if err := WriteFileAtomic(path, []byte(id.String()+"\n"), 0o644); err != nil { + s.warnOnce("Could not persist the instance ID; a restart creates a new instance", "path", path, "error", err) + return id, false + } + return id, true +} + +func (s *Store) warnOnce(msg string, args ...any) { + if s.warned { + return + } + s.warned = true + s.logger.Warn(msg, args...) +} + +func fileExists(path string) bool { + _, err := os.Stat(path) + return err == nil +} + +// WriteFileAtomic writes data to a temp file in path's directory, fsyncs it and renames it +// over path, so readers see either the old or the new content. +func WriteFileAtomic(path string, data []byte, perm os.FileMode) error { + dir := filepath.Dir(path) + tmp, err := os.CreateTemp(dir, "."+filepath.Base(path)+".tmp-*") + if err != nil { + return err + } + tmpName := tmp.Name() + cleanup := func() { _ = os.Remove(tmpName) } + if _, err := tmp.Write(data); err != nil { + _ = tmp.Close() + cleanup() + return err + } + if err := tmp.Chmod(perm); err != nil { + _ = tmp.Close() + cleanup() + return err + } + if err := tmp.Sync(); err != nil { + _ = tmp.Close() + cleanup() + return err + } + if err := tmp.Close(); err != nil { + cleanup() + return err + } + if err := os.Rename(tmpName, path); err != nil { + cleanup() + return fmt.Errorf("rename %s: %w", path, err) + } + return nil +} diff --git a/internal/agentstate/store_test.go b/internal/agentstate/store_test.go new file mode 100644 index 0000000..fac5427 --- /dev/null +++ b/internal/agentstate/store_test.go @@ -0,0 +1,112 @@ +package agentstate + +import ( + "bytes" + "os" + "path/filepath" + "runtime" + "strings" + "testing" + + "github.com/google/uuid" + "github.com/hashicorp/go-hclog" +) + +func TestInstanceID_PersistedAndReused(t *testing.T) { + dir := filepath.Join(t.TempDir(), "state") + first, persisted := Open(dir, nil).InstanceID("") + if !persisted { + t.Fatal("expected the ID to be persisted") + } + info, err := os.Stat(dir) + if err != nil { + t.Fatal(err) + } + if runtime.GOOS != "windows" && info.Mode().Perm() != 0o700 { + t.Fatalf("state dir mode = %v, want 0700", info.Mode().Perm()) + } + second, persisted := Open(dir, nil).InstanceID("") + if !persisted || second != first { + t.Fatalf("expected the persisted ID to be reused: %s vs %s", first, second) + } +} + +func TestInstanceID_OverrideNotPersisted(t *testing.T) { + dir := t.TempDir() + override := uuid.New() + id, persisted := Open(dir, nil).InstanceID(override.String()) + if id != override || persisted { + t.Fatalf("override: got %s persisted=%v", id, persisted) + } + if _, err := os.Stat(filepath.Join(dir, instanceIDFile)); !os.IsNotExist(err) { + t.Fatalf("override must not be written, stat err = %v", err) + } +} + +func TestInstanceID_ReadOnlyDirKeepsIDInMemory(t *testing.T) { + if runtime.GOOS == "windows" || os.Geteuid() == 0 { + t.Skip("permission bits are not enforced") + } + dir := t.TempDir() + if err := os.Chmod(dir, 0o500); err != nil { + t.Fatal(err) + } + t.Cleanup(func() { _ = os.Chmod(dir, 0o700) }) + + var logs bytes.Buffer + logger := hclog.New(&hclog.LoggerOptions{Output: &logs, Level: hclog.Warn}) + s := Open(dir, logger) + id, persisted := s.InstanceID("") + if id == uuid.Nil || persisted { + t.Fatalf("expected an in-memory ID, got %s persisted=%v", id, persisted) + } + if again, _ := s.InstanceID(""); again != id { + t.Fatalf("in-memory ID must be stable for the process: %s vs %s", id, again) + } + if !strings.Contains(logs.String(), "[WARN]") { + t.Fatalf("expected a WARN, got %q", logs.String()) + } +} + +func TestInstanceID_CorruptFileReplaced(t *testing.T) { + dir := t.TempDir() + path := filepath.Join(dir, instanceIDFile) + if err := os.WriteFile(path, []byte("not-a-uuid\n"), 0o644); err != nil { + t.Fatal(err) + } + id, persisted := Open(dir, nil).InstanceID("") + if !persisted { + t.Fatal("expected the replacement to be persisted") + } + raw, err := os.ReadFile(path) + if err != nil { + t.Fatal(err) + } + if strings.TrimSpace(string(raw)) != id.String() { + t.Fatalf("file holds %q, want %s", raw, id) + } +} + +func TestDefaultDir_DependsOnConfigPath(t *testing.T) { + a, err := DefaultDir("/etc/ccf/a.yaml") + if err != nil { + t.Fatal(err) + } + b, err := DefaultDir("/etc/ccf/b.yaml") + if err != nil { + t.Fatal(err) + } + if a == b { + t.Fatalf("two config paths must get two state dirs, got %s", a) + } + if !strings.Contains(filepath.ToSlash(a), StateRoot+"/") || len(filepath.Base(a)) != 16 { + t.Fatalf("unexpected layout %s", a) + } + + root := t.TempDir() + idA, _ := Open(filepath.Join(root, filepath.Base(a)), nil).InstanceID("") + idB, _ := Open(filepath.Join(root, filepath.Base(b)), nil).InstanceID("") + if idA == idB { + t.Fatal("two config paths must get two instance IDs") + } +}