Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
85 changes: 80 additions & 5 deletions cmd/agent.go
Original file line number Diff line number Diff line change
Expand Up @@ -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"
Expand Down Expand Up @@ -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/<hash of the config path>")
agentCmd.Flags().String("instance-id", "", "Pin this instance's UUID (not persisted); overrides CCF_INSTANCE_ID")

return agentCmd
}

Expand Down Expand Up @@ -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()
Expand Down Expand Up @@ -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
Expand All @@ -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",
Expand Down Expand Up @@ -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())
}
Expand Down
35 changes: 35 additions & 0 deletions cmd/reconciler_test.go
Original file line number Diff line number Diff line change
@@ -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)
}
}
16 changes: 16 additions & 0 deletions docs/configuration.md
Original file line number Diff line number Diff line change
Expand Up @@ -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/<key>/`, relative to the working directory, where
`<key>` 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.
15 changes: 13 additions & 2 deletions docs/running_as_a_service.md
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand All @@ -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
Expand Down Expand Up @@ -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

Expand Down
Loading
Loading