From f67642e4bdc50e3bb3e908da1bef48318e1a1800 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:14:27 -0300 Subject: [PATCH 1/3] feat(api): agent-facing config overlay route and heartbeat instance registration Eleventh layer of the agent remote-configuration stack (split from #465): GET /api/agent/config returns the authenticated agent's current overlay with an opaque ETag (If-None-Match -> 304), agent JWT only and agent:sync; authenticated heartbeats refresh the instance's last-seen time and, with config_digest/config_revision, register the instance. Co-Authored-By: Claude Opus 5.5 --- docs/docs.go | 93 +++++- docs/swagger.json | 93 +++++- docs/swagger.yaml | 70 +++- internal/api/handler/agent_config_body.go | 8 + internal/api/handler/agent_config_sync.go | 103 ++++++ .../agent_config_sync_integration_test.go | 307 ++++++++++++++++++ internal/api/handler/api.go | 17 +- internal/api/handler/heartbeat.go | 55 +++- 8 files changed, 738 insertions(+), 8 deletions(-) create mode 100644 internal/api/handler/agent_config_body.go create mode 100644 internal/api/handler/agent_config_sync.go create mode 100644 internal/api/handler/agent_config_sync_integration_test.go diff --git a/docs/docs.go b/docs/docs.go index 01d4140f..9c77e160 100644 --- a/docs/docs.go +++ b/docs/docs.go @@ -2578,9 +2578,63 @@ const docTemplate = `{ ] } }, + "/agent/config": { + "get": { + "description": "Returns the authenticated agent's current remote-configuration overlay (an RFC 7396 merge patch over its local config file, snake_case) with an opaque ETag. Send the raw ETag back as If-None-Match to get a 304 when nothing changed; never construct one. Revision 0 means no overlay. Agent JWT only; a 404 means the API predates remote configuration.", + "produces": [ + "application/json" + ], + "tags": [ + "Agents" + ], + "summary": "Get this agent's configuration overlay", + "parameters": [ + { + "type": "string", + "description": "ETag of the overlay the agent already has", + "name": "If-None-Match", + "in": "header" + } + ], + "responses": { + "200": { + "description": "OK", + "schema": { + "$ref": "#/definitions/handler.GenericDataResponse-agentconfig_OverlayDocument" + } + }, + "304": { + "description": "Not Modified" + }, + "401": { + "description": "Unauthorized", + "schema": { + "$ref": "#/definitions/api.Error" + } + }, + "403": { + "description": "Forbidden", + "schema": { + "$ref": "#/definitions/api.Error" + } + }, + "500": { + "description": "Internal Server Error", + "schema": { + "$ref": "#/definitions/api.Error" + } + } + }, + "security": [ + { + "OAuth2Password": [] + } + ] + } + }, "/agent/heartbeat": { "post": { - "description": "Creates a new heartbeat record for monitoring.", + "description": "Creates a new heartbeat record for monitoring. An authenticated agent heartbeat also refreshes the instance's last-seen time; with config_digest (new agents, mode != off) it records config_revision/config_digest and registers the instance if needed.", "consumes": [ "application/json" ], @@ -34330,6 +34384,21 @@ const docTemplate = `{ } }, "definitions": { + "agentconfig.OverlayDocument": { + "type": "object", + "properties": { + "created-at": { + "description": "omitted for revision 0", + "type": "string" + }, + "overlay": { + "type": "object" + }, + "revision": { + "type": "integer" + } + } + }, "api.Error": { "type": "object", "properties": { @@ -36340,6 +36409,19 @@ const docTemplate = `{ "meta": {} } }, + "handler.GenericDataResponse-agentconfig_OverlayDocument": { + "type": "object", + "properties": { + "data": { + "description": "Wrapped response data", + "allOf": [ + { + "$ref": "#/definitions/agentconfig.OverlayDocument" + } + ] + } + } + }, "handler.GenericDataResponse-array_handler_catalogLinkSummary": { "type": "object", "properties": { @@ -37850,10 +37932,19 @@ const docTemplate = `{ "uuid" ], "properties": { + "config_digest": { + "description": "ConfigDigest is agentconfig.Digest of the effective config; sent whenever mode != off.", + "type": "string" + }, + "config_revision": { + "description": "ConfigRevision is the APPLIED configuration revision, 0 when running from the file only\n(R45). Absent for old agents and mode off.", + "type": "integer" + }, "created_at": { "type": "string" }, "uuid": { + "description": "UUID is the agent instance id.", "type": "string" } } diff --git a/docs/swagger.json b/docs/swagger.json index 8aeaa21c..d9ccb50e 100644 --- a/docs/swagger.json +++ b/docs/swagger.json @@ -2572,9 +2572,63 @@ ] } }, + "/agent/config": { + "get": { + "description": "Returns the authenticated agent's current remote-configuration overlay (an RFC 7396 merge patch over its local config file, snake_case) with an opaque ETag. Send the raw ETag back as If-None-Match to get a 304 when nothing changed; never construct one. Revision 0 means no overlay. Agent JWT only; a 404 means the API predates remote configuration.", + "produces": [ + "application/json" + ], + "tags": [ + "Agents" + ], + "summary": "Get this agent's configuration overlay", + "parameters": [ + { + "type": "string", + "description": "ETag of the overlay the agent already has", + "name": "If-None-Match", + "in": "header" + } + ], + "responses": { + "200": { + "description": "OK", + "schema": { + "$ref": "#/definitions/handler.GenericDataResponse-agentconfig_OverlayDocument" + } + }, + "304": { + "description": "Not Modified" + }, + "401": { + "description": "Unauthorized", + "schema": { + "$ref": "#/definitions/api.Error" + } + }, + "403": { + "description": "Forbidden", + "schema": { + "$ref": "#/definitions/api.Error" + } + }, + "500": { + "description": "Internal Server Error", + "schema": { + "$ref": "#/definitions/api.Error" + } + } + }, + "security": [ + { + "OAuth2Password": [] + } + ] + } + }, "/agent/heartbeat": { "post": { - "description": "Creates a new heartbeat record for monitoring.", + "description": "Creates a new heartbeat record for monitoring. An authenticated agent heartbeat also refreshes the instance's last-seen time; with config_digest (new agents, mode != off) it records config_revision/config_digest and registers the instance if needed.", "consumes": [ "application/json" ], @@ -34324,6 +34378,21 @@ } }, "definitions": { + "agentconfig.OverlayDocument": { + "type": "object", + "properties": { + "created-at": { + "description": "omitted for revision 0", + "type": "string" + }, + "overlay": { + "type": "object" + }, + "revision": { + "type": "integer" + } + } + }, "api.Error": { "type": "object", "properties": { @@ -36334,6 +36403,19 @@ "meta": {} } }, + "handler.GenericDataResponse-agentconfig_OverlayDocument": { + "type": "object", + "properties": { + "data": { + "description": "Wrapped response data", + "allOf": [ + { + "$ref": "#/definitions/agentconfig.OverlayDocument" + } + ] + } + } + }, "handler.GenericDataResponse-array_handler_catalogLinkSummary": { "type": "object", "properties": { @@ -37844,10 +37926,19 @@ "uuid" ], "properties": { + "config_digest": { + "description": "ConfigDigest is agentconfig.Digest of the effective config; sent whenever mode != off.", + "type": "string" + }, + "config_revision": { + "description": "ConfigRevision is the APPLIED configuration revision, 0 when running from the file only\n(R45). Absent for old agents and mode off.", + "type": "integer" + }, "created_at": { "type": "string" }, "uuid": { + "description": "UUID is the agent instance id.", "type": "string" } } diff --git a/docs/swagger.yaml b/docs/swagger.yaml index e60a83ad..c93f24e6 100644 --- a/docs/swagger.yaml +++ b/docs/swagger.yaml @@ -2,6 +2,16 @@ basePath: /api consumes: - application/json definitions: + agentconfig.OverlayDocument: + properties: + created-at: + description: omitted for revision 0 + type: string + overlay: + type: object + revision: + type: integer + type: object api.Error: properties: errors: @@ -1418,6 +1428,13 @@ definitions: type: array meta: {} type: object + handler.GenericDataResponse-agentconfig_OverlayDocument: + properties: + data: + allOf: + - $ref: '#/definitions/agentconfig.OverlayDocument' + description: Wrapped response data + type: object handler.GenericDataResponse-array_handler_catalogLinkSummary: properties: data: @@ -2237,9 +2254,19 @@ definitions: type: object handler.HeartbeatCreateRequest: properties: + config_digest: + description: ConfigDigest is agentconfig.Digest of the effective config; sent + whenever mode != off. + type: string + config_revision: + description: |- + ConfigRevision is the APPLIED configuration revision, 0 when running from the file only + (R45). Absent for old agents and mode off. + type: integer created_at: type: string uuid: + description: UUID is the agent instance id. type: string required: - created_at @@ -14097,11 +14124,52 @@ paths: summary: Upload an artifact tags: - Artifacts + /agent/config: + get: + description: Returns the authenticated agent's current remote-configuration + overlay (an RFC 7396 merge patch over its local config file, snake_case) with + an opaque ETag. Send the raw ETag back as If-None-Match to get a 304 when + nothing changed; never construct one. Revision 0 means no overlay. Agent JWT + only; a 404 means the API predates remote configuration. + parameters: + - description: ETag of the overlay the agent already has + in: header + name: If-None-Match + type: string + produces: + - application/json + responses: + "200": + description: OK + schema: + $ref: '#/definitions/handler.GenericDataResponse-agentconfig_OverlayDocument' + "304": + description: Not Modified + "401": + description: Unauthorized + schema: + $ref: '#/definitions/api.Error' + "403": + description: Forbidden + schema: + $ref: '#/definitions/api.Error' + "500": + description: Internal Server Error + schema: + $ref: '#/definitions/api.Error' + security: + - OAuth2Password: [] + summary: Get this agent's configuration overlay + tags: + - Agents /agent/heartbeat: post: consumes: - application/json - description: Creates a new heartbeat record for monitoring. + description: Creates a new heartbeat record for monitoring. An authenticated + agent heartbeat also refreshes the instance's last-seen time; with config_digest + (new agents, mode != off) it records config_revision/config_digest and registers + the instance if needed. parameters: - description: Heartbeat payload in: body diff --git a/internal/api/handler/agent_config_body.go b/internal/api/handler/agent_config_body.go new file mode 100644 index 00000000..a7838b12 --- /dev/null +++ b/internal/api/handler/agent_config_body.go @@ -0,0 +1,8 @@ +package handler + +// HTTP header names used by the agent-configuration routes. +const ( + headerETag = "ETag" + headerIfMatch = "If-Match" + headerIfNoneMatch = "If-None-Match" +) diff --git a/internal/api/handler/agent_config_sync.go b/internal/api/handler/agent_config_sync.go new file mode 100644 index 00000000..e530d441 --- /dev/null +++ b/internal/api/handler/agent_config_sync.go @@ -0,0 +1,103 @@ +package handler + +import ( + "encoding/json" + "errors" + "net/http" + "regexp" + + "github.com/compliance-framework/api/internal/api" + "github.com/compliance-framework/api/internal/api/middleware" + "github.com/compliance-framework/api/internal/service/relational/agentcfg" + "github.com/compliance-framework/api/pkg/agentconfig" + "github.com/google/uuid" + "github.com/labstack/echo/v4" + "go.uber.org/zap" +) + +const ( + // headerRemoteConfig marks responses from the remote-configuration routes, so an agent can + // tell them apart from a proxy's. + headerRemoteConfig = "X-CCF-Remote-Config" +) + +var effectiveDigestPattern = regexp.MustCompile(`^sha256:[0-9a-f]{64}$`) + +// AgentConfigSyncHandler serves the agent-facing remote-configuration routes. They accept +// agent JWTs only (never anonymous, even with public agent endpoints on) and act on the +// authenticated agent's own configuration. +type AgentConfigSyncHandler struct { + sugar *zap.SugaredLogger + svc *agentcfg.Service +} + +func NewAgentConfigSyncHandler(sugar *zap.SugaredLogger, svc *agentcfg.Service) *AgentConfigSyncHandler { + return &AgentConfigSyncHandler{sugar: sugar, svc: svc} +} + +// Register mounts GET /config on the /agent group. Pass the strict agent JWT middleware and +// the agent:sync guard. +func (h *AgentConfigSyncHandler) Register(g *echo.Group, middlewares ...echo.MiddlewareFunc) { + g.GET("/config", h.GetConfig, middlewares...) +} + +// agentAuthFrom returns the authenticated agent, or nil. The handler requires it even though +// the route middleware does too: Cedar grants the agent role to anonymous subjects when +// public agent endpoints are on (defence in depth). +func agentAuthFrom(ctx echo.Context) *middleware.AgentAuthContext { + auth, _ := ctx.Get("agent_auth").(*middleware.AgentAuthContext) + if auth == nil || auth.Agent == nil || auth.Agent.ID == nil { + return nil + } + return auth +} + +func setRemoteConfigHeaders(ctx echo.Context) { + ctx.Response().Header().Set(headerRemoteConfig, "1") + ctx.Response().Header().Set(echo.HeaderCacheControl, "no-cache") +} + +// GetConfig godoc +// +// @Summary Get this agent's configuration overlay +// @Description Returns the authenticated agent's current remote-configuration overlay (an RFC 7396 merge patch over its local config file, snake_case) with an opaque ETag. Send the raw ETag back as If-None-Match to get a 304 when nothing changed; never construct one. Revision 0 means no overlay. Agent JWT only; a 404 means the API predates remote configuration. +// @Tags Agents +// @Produce json +// @Param If-None-Match header string false "ETag of the overlay the agent already has" +// @Success 200 {object} handler.GenericDataResponse[agentconfig.OverlayDocument] +// @Success 304 "Not Modified" +// @Failure 401 {object} api.Error +// @Failure 403 {object} api.Error +// @Failure 500 {object} api.Error +// @Security OAuth2Password +// @Router /agent/config [get] +func (h *AgentConfigSyncHandler) GetConfig(ctx echo.Context) error { + auth := agentAuthFrom(ctx) + if auth == nil { + return ctx.JSON(http.StatusUnauthorized, api.NewError(errors.New("agent authentication required"))) + } + agentID := *auth.Agent.ID + + cur, err := h.svc.Current(ctx.Request().Context(), agentID) + if err != nil { + h.sugar.Errorw("Failed to load agent configuration", "agentID", agentID, "error", err) + return ctx.JSON(http.StatusInternalServerError, api.InternalServerError()) + } + doc := agentconfig.OverlayDocument{Revision: 0, Overlay: json.RawMessage(`{}`)} + rowID := uuid.Nil + if cur != nil { + doc.Revision = cur.Revision + doc.Overlay = json.RawMessage(cur.Overlay) + createdAt := cur.CreatedAt.UTC() + doc.CreatedAt = &createdAt + rowID = *cur.ID + } + etag := agentconfig.ETagForRevision(doc.Revision, rowID, agentID) + + setRemoteConfigHeaders(ctx) + ctx.Response().Header().Set(headerETag, etag) + if agentconfig.MatchIfNoneMatch(ctx.Request().Header.Get("If-None-Match"), etag) { + return ctx.NoContent(http.StatusNotModified) + } + return ctx.JSON(http.StatusOK, GenericDataResponse[agentconfig.OverlayDocument]{Data: doc}) +} diff --git a/internal/api/handler/agent_config_sync_integration_test.go b/internal/api/handler/agent_config_sync_integration_test.go new file mode 100644 index 00000000..cd151de2 --- /dev/null +++ b/internal/api/handler/agent_config_sync_integration_test.go @@ -0,0 +1,307 @@ +//go:build integration + +package handler + +import ( + "bytes" + "context" + "encoding/json" + "fmt" + "net/http" + "net/http/httptest" + "strings" + "testing" + "time" + + "github.com/compliance-framework/api/internal/api" + "github.com/compliance-framework/api/internal/service/relational" + "github.com/compliance-framework/api/internal/service/relational/agentcfg" + "github.com/compliance-framework/api/internal/tests" + "github.com/compliance-framework/api/pkg/agentconfig" + "github.com/google/uuid" + "github.com/labstack/echo/v4" + "github.com/stretchr/testify/suite" + "go.uber.org/zap" +) + +const syncTestDigest = "sha256:abcdef0123456789abcdef0123456789abcdef0123456789abcdef0123456789" + +type AgentConfigSyncIntegrationSuite struct { + tests.IntegrationTestSuite + server *api.Server + svc *agentcfg.Service +} + +func TestAgentConfigSyncAPI(t *testing.T) { + suite.Run(t, new(AgentConfigSyncIntegrationSuite)) +} + +type syncAgent struct { + agent *relational.Agent + key *relational.AgentServiceAccountKey + token string +} + +func (s *AgentConfigSyncIntegrationSuite) SetupTest() { + s.Require().NoError(s.Migrator.Refresh()) + // Public agent endpoints stay ENABLED: the sync routes must still require an agent JWT. + s.Config.StrictDisablePublicAgentEndpoints = false + s.server = s.buildServer() + s.svc = agentcfg.NewService(s.DB, agentcfg.Settings{}, nil) +} + +func (s *AgentConfigSyncIntegrationSuite) buildServer() *api.Server { + logger, _ := zap.NewDevelopment() + metrics := api.NewMetricsHandler(context.Background(), logger.Sugar()) + server := api.NewServer(context.Background(), logger.Sugar(), s.Config, metrics) + RegisterHandlers(server, logger.Sugar(), s.DB, s.Config, &APIServices{}) + return server +} + +func (s *AgentConfigSyncIntegrationSuite) newAgent(name string) syncAgent { + agent, err := s.CreateAgent(name) + s.Require().NoError(err) + key, _, err := s.CreateAgentKey(agent, name+"-key") + s.Require().NoError(err) + token, err := s.GetAgentToken(agent, key) + s.Require().NoError(err) + return syncAgent{agent: agent, key: key, token: *token} +} + +func (s *AgentConfigSyncIntegrationSuite) do(server *api.Server, method, path, token string, body []byte, headers map[string]string) *httptest.ResponseRecorder { + rec := httptest.NewRecorder() + req := httptest.NewRequest(method, path, bytes.NewReader(body)) + if body != nil { + req.Header.Set(echo.HeaderContentType, echo.MIMEApplicationJSON) + } + if token != "" { + req.Header.Set(echo.HeaderAuthorization, "Bearer "+token) + } + for k, v := range headers { + req.Header.Set(k, v) + } + server.E().ServeHTTP(rec, req) + return rec +} + +func (s *AgentConfigSyncIntegrationSuite) getConfig(token, ifNoneMatch string) *httptest.ResponseRecorder { + var headers map[string]string + if ifNoneMatch != "" { + headers = map[string]string{"If-None-Match": ifNoneMatch} + } + return s.do(s.server, http.MethodGet, "/api/agent/config", token, nil, headers) +} + +func (s *AgentConfigSyncIntegrationSuite) createRevision(agentID uuid.UUID, expected int64, overlay string) *relational.AgentConfigRevision { + rev, err := s.svc.CreateRevision(context.Background(), agentcfg.CreateRevisionParams{ + AgentID: agentID, + ExpectedRevision: expected, + Overlay: json.RawMessage(overlay), + CreatedBy: "test@example.com", + }) + s.Require().NoError(err) + return rev +} + +func (s *AgentConfigSyncIntegrationSuite) instance(agentID, instanceID uuid.UUID) (relational.AgentInstance, bool) { + var row relational.AgentInstance + res := s.DB.Where("agent_id = ? AND instance_id = ?", agentID, instanceID).Limit(1).Find(&row) + s.Require().NoError(res.Error) + return row, res.RowsAffected == 1 +} + +func (s *AgentConfigSyncIntegrationSuite) assertRemoteConfigHeaders(rec *httptest.ResponseRecorder) { + s.Equal("1", rec.Header().Get("X-CCF-Remote-Config")) + s.Equal("no-cache", rec.Header().Get(echo.HeaderCacheControl)) +} + +// ---- GET /api/agent/config ---- + +func (s *AgentConfigSyncIntegrationSuite) TestGetConfigRevisionZero() { + a := s.newAgent("rev-zero") + wantTag := fmt.Sprintf(`"r0-%s"`, *a.agent.ID) + + rec := s.getConfig(a.token, "") + s.Require().Equal(http.StatusOK, rec.Code, rec.Body.String()) + s.JSONEq(`{"data":{"revision":0,"overlay":{}}}`, rec.Body.String()) + s.Equal(wantTag, rec.Header().Get("ETag")) + s.assertRemoteConfigHeaders(rec) + + rec = s.getConfig(a.token, wantTag) + s.Require().Equal(http.StatusNotModified, rec.Code) + s.Empty(rec.Body.Bytes()) + s.Equal(wantTag, rec.Header().Get("ETag")) + s.assertRemoteConfigHeaders(rec) + + // A tag for some other agent's revision 0 does not match. + rec = s.getConfig(a.token, fmt.Sprintf(`"r0-%s"`, uuid.New())) + s.Equal(http.StatusOK, rec.Code) +} + +func (s *AgentConfigSyncIntegrationSuite) TestGetConfigAfterRevision() { + a := s.newAgent("with-revision") + r0Tag := fmt.Sprintf(`"r0-%s"`, *a.agent.ID) + overlay := `{"plugins":{"x":{"schedule":"*/5 * * * *"}}}` + rev := s.createRevision(*a.agent.ID, 0, overlay) + r1Tag := fmt.Sprintf(`"r1-%s"`, *rev.ID) + + rec := s.getConfig(a.token, "") + s.Require().Equal(http.StatusOK, rec.Code, rec.Body.String()) + s.Equal(r1Tag, rec.Header().Get("ETag")) + s.assertRemoteConfigHeaders(rec) + var got struct { + Data struct { + Revision int64 `json:"revision"` + Overlay json.RawMessage `json:"overlay"` + CreatedAt *time.Time `json:"created-at"` + } `json:"data"` + } + s.Require().NoError(json.Unmarshal(rec.Body.Bytes(), &got)) + s.Equal(int64(1), got.Data.Revision) + s.JSONEq(overlay, string(got.Data.Overlay)) + s.Require().NotNil(got.Data.CreatedAt) + s.WithinDuration(rev.CreatedAt, *got.Data.CreatedAt, time.Second) + + // The previous revision-0 tag no longer matches. + rec = s.getConfig(a.token, r0Tag) + s.Equal(http.StatusOK, rec.Code) + + bare := strings.Trim(r1Tag, `"`) + for _, inm := range []string{ + r1Tag, + "W/" + r1Tag, + bare, + fmt.Sprintf(`%s, "r7-%s", %s`, r0Tag, uuid.New(), r1Tag), + "*", + } { + rec = s.getConfig(a.token, inm) + s.Equal(http.StatusNotModified, rec.Code, "If-None-Match %q", inm) + s.Empty(rec.Body.Bytes(), "If-None-Match %q", inm) + s.Equal(r1Tag, rec.Header().Get("ETag")) + s.assertRemoteConfigHeaders(rec) + } +} + +func (s *AgentConfigSyncIntegrationSuite) TestGetConfigAfterSimulatedReset() { + a := s.newAgent("reset") + old := s.createRevision(*a.agent.ID, 0, `{"a":1}`) + oldTag := fmt.Sprintf(`"r1-%s"`, *old.ID) + s.Require().Equal(http.StatusNotModified, s.getConfig(a.token, oldTag).Code) + + // The gorm model is append-only (BeforeDelete hook), so reset with raw SQL. + s.Require().NoError(s.DB.Exec("DELETE FROM ccf_agent_config_revisions").Error) + recreated := s.createRevision(*a.agent.ID, 0, `{"a":1}`) + s.Require().NotEqual(*old.ID, *recreated.ID) + + rec := s.getConfig(a.token, oldTag) + s.Require().Equal(http.StatusOK, rec.Code, "same revision number, different row: no false 304") + s.Equal(fmt.Sprintf(`"r1-%s"`, *recreated.ID), rec.Header().Get("ETag")) +} + +func (s *AgentConfigSyncIntegrationSuite) TestGetConfigAuth() { + a := s.newAgent("auth") + + s.Equal(http.StatusUnauthorized, s.getConfig("", "").Code, "no token, public agent endpoints on") + + userToken, err := s.GetAuthToken() + s.Require().NoError(err) + s.Equal(http.StatusUnauthorized, s.getConfig(*userToken, "").Code, "user JWT") + + s.Equal(http.StatusUnauthorized, s.getConfig("not-a-jwt", "").Code, "garbage token") + + s.Require().Equal(http.StatusOK, s.getConfig(a.token, "").Code) + + // Revoked key. + revoked := time.Now().UTC().Add(-time.Minute) + s.Require().NoError(s.DB.Model(&relational.AgentServiceAccountKey{}).Where("id = ?", *a.key.ID).Update("revoked_at", revoked).Error) + s.Equal(http.StatusForbidden, s.getConfig(a.token, "").Code, "revoked key") + + // Inactive agent. + b := s.newAgent("inactive") + s.Require().NoError(s.DB.Exec("UPDATE ccf_agents SET is_active = false WHERE id = ?", *b.agent.ID).Error) + s.Equal(http.StatusForbidden, s.getConfig(b.token, "").Code, "inactive agent") +} + +func (s *AgentConfigSyncIntegrationSuite) TestGetConfigCrossAgentIsolation() { + a := s.newAgent("agent-a") + b := s.newAgent("agent-b") + c := s.newAgent("agent-c") + s.createRevision(*a.agent.ID, 0, `{"who":"a"}`) + s.createRevision(*a.agent.ID, 1, `{"who":"a2"}`) + revB := s.createRevision(*b.agent.ID, 0, `{"who":"b"}`) + + rec := s.getConfig(b.token, "") + s.Require().Equal(http.StatusOK, rec.Code) + var got GenericDataResponse[agentconfig.OverlayDocument] + s.Require().NoError(json.Unmarshal(rec.Body.Bytes(), &got)) + s.Equal(int64(1), got.Data.Revision) + s.JSONEq(`{"who":"b"}`, string(got.Data.Overlay)) + s.Equal(fmt.Sprintf(`"r1-%s"`, *revB.ID), rec.Header().Get("ETag")) + + rec = s.getConfig(a.token, "") + s.Require().Equal(http.StatusOK, rec.Code) + s.Require().NoError(json.Unmarshal(rec.Body.Bytes(), &got)) + s.Equal(int64(2), got.Data.Revision) + s.JSONEq(`{"who":"a2"}`, string(got.Data.Overlay)) + + rec = s.getConfig(c.token, "") + s.Require().Equal(http.StatusOK, rec.Code) + s.JSONEq(`{"data":{"revision":0,"overlay":{}}}`, rec.Body.String()) + s.Equal(fmt.Sprintf(`"r0-%s"`, *c.agent.ID), rec.Header().Get("ETag")) +} + +// ---- PUT /api/agent/instances/:instanceId/config-report ---- + +// ---- POST /api/agent/heartbeat ---- + +func (s *AgentConfigSyncIntegrationSuite) heartbeat(token string, instanceID uuid.UUID, rev *int64, digest *string) *httptest.ResponseRecorder { + raw, err := json.Marshal(HeartbeatCreateRequest{ + UUID: instanceID, + CreatedAt: time.Now().UTC(), + ConfigRevision: rev, + ConfigDigest: digest, + }) + s.Require().NoError(err) + return s.do(s.server, http.MethodPost, "/api/agent/heartbeat", token, raw, nil) +} + +func (s *AgentConfigSyncIntegrationSuite) TestHeartbeatWithDigestRegistersInstance() { + a := s.newAgent("hb-digest") + instanceID := uuid.New() + rev := int64(3) + digest := syncTestDigest + + rec := s.heartbeat(a.token, instanceID, &rev, &digest) + s.Require().Equal(http.StatusCreated, rec.Code, rec.Body.String()) + + row, ok := s.instance(*a.agent.ID, instanceID) + s.Require().True(ok) + s.Require().NotNil(row.HeartbeatConfigDigest) + s.Equal(syncTestDigest, *row.HeartbeatConfigDigest) + s.Require().NotNil(row.HeartbeatConfigRevision) + s.Equal(int64(3), *row.HeartbeatConfigRevision) + s.Empty(row.ReportedStatus) + s.Empty(row.Mode) + s.Nil(row.ReportedAt) + s.Require().NotNil(row.CredentialID) + s.Equal(*a.key.ID, *row.CredentialID) +} + +func (s *AgentConfigSyncIntegrationSuite) TestHeartbeatMalformedDigestAndAnonymous() { + a := s.newAgent("hb-malformed") + instanceID := uuid.New() + rev := int64(1) + bad := "sha256:not-hex" + + rec := s.heartbeat(a.token, instanceID, &rev, &bad) + s.Require().Equal(http.StatusCreated, rec.Code, rec.Body.String()) + _, ok := s.instance(*a.agent.ID, instanceID) + s.False(ok, "a malformed digest is treated as absent") + + digest := syncTestDigest + rec = s.heartbeat("", uuid.New(), &rev, &digest) + s.Require().Equal(http.StatusCreated, rec.Code, rec.Body.String()) + var count int64 + s.Require().NoError(s.DB.Model(&relational.AgentInstance{}).Count(&count).Error) + s.Equal(int64(0), count, "anonymous heartbeats never touch the instance registry") +} diff --git a/internal/api/handler/api.go b/internal/api/handler/api.go index a8a87381..f43a16c8 100644 --- a/internal/api/handler/api.go +++ b/internal/api/handler/api.go @@ -11,6 +11,7 @@ import ( "github.com/compliance-framework/api/internal/config" "github.com/compliance-framework/api/internal/service/digest" "github.com/compliance-framework/api/internal/service/notification" + "github.com/compliance-framework/api/internal/service/relational/agentcfg" artifactsvc "github.com/compliance-framework/api/internal/service/relational/artifacts" evidencesvc "github.com/compliance-framework/api/internal/service/relational/evidence" poamsvc "github.com/compliance-framework/api/internal/service/relational/poam" @@ -94,7 +95,11 @@ func RegisterHandlers(server *api.Server, logger *zap.SugaredLogger, db *gorm.DB lineageGroup.Use(middleware.JWTMiddleware(config.JWTPublicKey)) lineageHandler.Register(lineageGroup, pep.For(authz.ResourceLineage)) - heartbeatHandler := NewHeartbeatHandler(logger, db) + // Agent remote configuration (overlay revisions + reporting instances). + agentCfgSvc := agentcfg.NewService(db, agentcfg.SettingsFromConfig(config), logger) + agentGuard := pep.For(authz.ResourceAgent) + + heartbeatHandler := NewHeartbeatHandler(logger, db).WithAgentInstances(agentCfgSvc) heartbeatGuard := pep.For(authz.ResourceHeartbeat) agentIngestMiddleware := middleware.AgentJWTOrPublicMiddleware(db, config.JWTPublicKey, !config.StrictDisablePublicAgentEndpoints) heartbeatHandler.RegisterCreate(server.API().Group("/agent/heartbeat"), agentIngestMiddleware, heartbeatGuard.Do(authz.ActionIngest)) @@ -215,12 +220,20 @@ func RegisterHandlers(server *api.Server, logger *zap.SugaredLogger, db *gorm.DB // Agent routes are guarded per route (R40): list/get need agent:read, writes and keys stay // admin:manage. The builtin PDP still requires the admin check for users on agent:*. - agentGuard := pep.For(authz.ResourceAgent) agentHandler := NewAgentHandler(logger, db) agentsGroup := server.API().Group("/admin/agents") agentsGroup.Use(middleware.JWTMiddleware(config.JWTPublicKey)) agentHandler.Register(agentsGroup, agentGuard.Read(), pep.Authorize(authz.ResourceAdmin, authz.ActionManage)) + // Agent-facing configuration sync: agent JWT only (strict — it ignores + // StrictDisablePublicAgentEndpoints) and agent:sync. + agentConfigSyncHandler := NewAgentConfigSyncHandler(logger, agentCfgSvc) + agentConfigSyncHandler.Register( + server.API().Group("/agent"), + middleware.AgentJWTMiddleware(db, config.JWTPublicKey), + agentGuard.Do(authz.ActionSync), + ) + userHandler := NewUserHandler(logger, db) adminGroup := server.API().Group("/admin/users") diff --git a/internal/api/handler/heartbeat.go b/internal/api/handler/heartbeat.go index 7ddb5a9d..fd49733d 100644 --- a/internal/api/handler/heartbeat.go +++ b/internal/api/handler/heartbeat.go @@ -6,6 +6,7 @@ import ( "github.com/compliance-framework/api/internal/api" "github.com/compliance-framework/api/internal/service" + "github.com/compliance-framework/api/internal/service/relational/agentcfg" "github.com/google/uuid" "github.com/labstack/echo/v4" "go.uber.org/zap" @@ -13,8 +14,9 @@ import ( ) type HeartbeatHandler struct { - db *gorm.DB - sugar *zap.SugaredLogger + db *gorm.DB + sugar *zap.SugaredLogger + instances *agentcfg.Service } func NewHeartbeatHandler(sugar *zap.SugaredLogger, db *gorm.DB) *HeartbeatHandler { @@ -24,6 +26,14 @@ func NewHeartbeatHandler(sugar *zap.SugaredLogger, db *gorm.DB) *HeartbeatHandle } } +// WithAgentInstances makes authenticated heartbeats update the agent-instance registry +// (R11): last-seen always, and a new instance row only when the heartbeat carries a config +// digest. +func (h *HeartbeatHandler) WithAgentInstances(svc *agentcfg.Service) *HeartbeatHandler { + h.instances = svc + return h +} + func (h *HeartbeatHandler) Register(api *echo.Group) { api.POST("", h.Create) api.GET("/over-time", h.OverTime) @@ -38,14 +48,20 @@ func (h *HeartbeatHandler) RegisterOverTime(api *echo.Group, middlewares ...echo } type HeartbeatCreateRequest struct { + // UUID is the agent instance id. UUID uuid.UUID `json:"uuid,omitempty" validate:"required"` CreatedAt time.Time `json:"created_at,omitempty" validate:"required"` + // ConfigRevision is the APPLIED configuration revision, 0 when running from the file only + // (R45). Absent for old agents and mode off. + ConfigRevision *int64 `json:"config_revision,omitempty"` + // ConfigDigest is agentconfig.Digest of the effective config; sent whenever mode != off. + ConfigDigest *string `json:"config_digest,omitempty"` } // Create godoc // // @Summary Create Heartbeat -// @Description Creates a new heartbeat record for monitoring. +// @Description Creates a new heartbeat record for monitoring. An authenticated agent heartbeat also refreshes the instance's last-seen time; with config_digest (new agents, mode != off) it records config_revision/config_digest and registers the instance if needed. // @Tags Heartbeat // @Accept json // @Produce json @@ -73,10 +89,43 @@ func (h *HeartbeatHandler) Create(ctx echo.Context) error { return ctx.JSON(http.StatusInternalServerError, api.NewError(err)) } + h.touchAgentInstance(ctx, heartbeat) + // Return a 201 Created response with no content. return ctx.NoContent(http.StatusCreated) } +// touchAgentInstance records an authenticated heartbeat against the agent-instance registry. +// Failures are logged and never change the 201; anonymous heartbeats are ignored. +func (h *HeartbeatHandler) touchAgentInstance(ctx echo.Context, heartbeat HeartbeatCreateRequest) { + if h.instances == nil { + return + } + auth := agentAuthFrom(ctx) + if auth == nil { + return + } + var credentialID *uuid.UUID + if auth.Key != nil && auth.Key.ID != nil { + id := *auth.Key.ID + credentialID = &id + } + // Only a well-formed digest may register an instance; anything else is treated as absent + // (last-seen update only). + digest := heartbeat.ConfigDigest + if digest != nil && !effectiveDigestPattern.MatchString(*digest) { + digest = nil + } + revision := heartbeat.ConfigRevision + if revision != nil && *revision < 0 { + revision = nil + } + if err := h.instances.TouchFromHeartbeat(ctx.Request().Context(), *auth.Agent.ID, credentialID, heartbeat.UUID, revision, digest); err != nil { + h.sugar.Warnw("Failed to record agent heartbeat against its instance", + "agentID", *auth.Agent.ID, "instanceID", heartbeat.UUID, "error", err) + } +} + // OverTime godoc // // @Summary Get Heartbeat Metrics Over Time From d71ccd71137ea0cb278f2175925b10dc32086dc7 Mon Sep 17 00:00:00 2001 From: "ccf-lisa[bot]" <286799724+ccf-lisa[bot]@users.noreply.github.com> Date: Mon, 5 Oct 2026 15:00:11 -0300 Subject: [PATCH 2/3] docs: heartbeats register instances under heartbeat:ingest Say in the heartbeat route description and next to the agent role's sync grant that removing sync does not stop heartbeats carrying config_digest from registering instances. Co-Authored-By: Claude Opus 5.5 --- docs/docs.go | 2 +- docs/swagger.json | 2 +- docs/swagger.yaml | 8 +++++--- internal/api/handler/heartbeat.go | 2 +- internal/authz/manifest.yaml | 2 ++ 5 files changed, 10 insertions(+), 6 deletions(-) diff --git a/docs/docs.go b/docs/docs.go index 9c77e160..9aafb6bd 100644 --- a/docs/docs.go +++ b/docs/docs.go @@ -2634,7 +2634,7 @@ const docTemplate = `{ }, "/agent/heartbeat": { "post": { - "description": "Creates a new heartbeat record for monitoring. An authenticated agent heartbeat also refreshes the instance's last-seen time; with config_digest (new agents, mode != off) it records config_revision/config_digest and registers the instance if needed.", + "description": "Creates a new heartbeat record for monitoring. An authenticated agent heartbeat also refreshes the instance's last-seen time; with config_digest (new agents, mode != off) it records config_revision/config_digest and registers the instance if needed. That registration is authorized by heartbeat:ingest, not agent:sync: removing sync from an agent's role stops overlay fetches and config reports, not instance registration by heartbeats.", "consumes": [ "application/json" ], diff --git a/docs/swagger.json b/docs/swagger.json index d9ccb50e..112a3afe 100644 --- a/docs/swagger.json +++ b/docs/swagger.json @@ -2628,7 +2628,7 @@ }, "/agent/heartbeat": { "post": { - "description": "Creates a new heartbeat record for monitoring. An authenticated agent heartbeat also refreshes the instance's last-seen time; with config_digest (new agents, mode != off) it records config_revision/config_digest and registers the instance if needed.", + "description": "Creates a new heartbeat record for monitoring. An authenticated agent heartbeat also refreshes the instance's last-seen time; with config_digest (new agents, mode != off) it records config_revision/config_digest and registers the instance if needed. That registration is authorized by heartbeat:ingest, not agent:sync: removing sync from an agent's role stops overlay fetches and config reports, not instance registration by heartbeats.", "consumes": [ "application/json" ], diff --git a/docs/swagger.yaml b/docs/swagger.yaml index c93f24e6..fb3c3c6c 100644 --- a/docs/swagger.yaml +++ b/docs/swagger.yaml @@ -14166,10 +14166,12 @@ paths: post: consumes: - application/json - description: Creates a new heartbeat record for monitoring. An authenticated - agent heartbeat also refreshes the instance's last-seen time; with config_digest + description: 'Creates a new heartbeat record for monitoring. An authenticated + agent heartbeat also refreshes the instance''s last-seen time; with config_digest (new agents, mode != off) it records config_revision/config_digest and registers - the instance if needed. + the instance if needed. That registration is authorized by heartbeat:ingest, + not agent:sync: removing sync from an agent''s role stops overlay fetches + and config reports, not instance registration by heartbeats.' parameters: - description: Heartbeat payload in: body diff --git a/internal/api/handler/heartbeat.go b/internal/api/handler/heartbeat.go index fd49733d..54121c09 100644 --- a/internal/api/handler/heartbeat.go +++ b/internal/api/handler/heartbeat.go @@ -61,7 +61,7 @@ type HeartbeatCreateRequest struct { // Create godoc // // @Summary Create Heartbeat -// @Description Creates a new heartbeat record for monitoring. An authenticated agent heartbeat also refreshes the instance's last-seen time; with config_digest (new agents, mode != off) it records config_revision/config_digest and registers the instance if needed. +// @Description Creates a new heartbeat record for monitoring. An authenticated agent heartbeat also refreshes the instance's last-seen time; with config_digest (new agents, mode != off) it records config_revision/config_digest and registers the instance if needed. That registration is authorized by heartbeat:ingest, not agent:sync: removing sync from an agent's role stops overlay fetches and config reports, not instance registration by heartbeats. // @Tags Heartbeat // @Accept json // @Produce json diff --git a/internal/authz/manifest.yaml b/internal/authz/manifest.yaml index 57a12195..ca09ba2c 100644 --- a/internal/authz/manifest.yaml +++ b/internal/authz/manifest.yaml @@ -365,6 +365,8 @@ roles: agent: evidence: [create] heartbeat: [ingest] + # sync gates the overlay fetch and config reports. A heartbeat carrying config_digest also + # registers the instance, under heartbeat:ingest, so removing sync does not stop that. agent: [register, ingest, sync] risk-template: [update] subject-template: [update] From 64f2362d85f7bc1f536709461f50748e26fb3446 Mon Sep 17 00:00:00 2001 From: "ccf-lisa[bot]" <286799724+ccf-lisa[bot]@users.noreply.github.com> Date: Tue, 6 Oct 2026 08:25:57 -0300 Subject: [PATCH 3/3] fix(api): answer a 304 config poll without loading the overlay GetConfig loaded the whole current revision (SELECT *, overlay included) before comparing If-None-Match, although nearly every poll is a 304. With an If-None-Match header, read only the revision head (id, revision, created_at; same indexed query) through agentcfg.CurrentHead, compute the ETag and return 304 on a match. The overlay is loaded only on a miss, and the response is unchanged. Co-Authored-By: Claude Opus 5.5 --- internal/api/handler/agent_config_sync.go | 25 +++++++++- ...g_sync_poll_regression_integration_test.go | 46 +++++++++++++++++++ .../service/relational/agentcfg/service.go | 24 ++++++++++ .../agentcfg/service_integration_test.go | 16 +++++++ 4 files changed, 109 insertions(+), 2 deletions(-) create mode 100644 internal/api/handler/agent_config_sync_poll_regression_integration_test.go diff --git a/internal/api/handler/agent_config_sync.go b/internal/api/handler/agent_config_sync.go index e530d441..5dbb0a00 100644 --- a/internal/api/handler/agent_config_sync.go +++ b/internal/api/handler/agent_config_sync.go @@ -77,8 +77,29 @@ func (h *AgentConfigSyncHandler) GetConfig(ctx echo.Context) error { return ctx.JSON(http.StatusUnauthorized, api.NewError(errors.New("agent authentication required"))) } agentID := *auth.Agent.ID + reqCtx := ctx.Request().Context() + ifNoneMatch := ctx.Request().Header.Get("If-None-Match") - cur, err := h.svc.Current(ctx.Request().Context(), agentID) + // Nearly every poll ends in 304: check the agent's ETag against the revision head first, + // and load the overlay only when the agent's copy is stale. + if ifNoneMatch != "" { + head, err := h.svc.CurrentHead(reqCtx, agentID) + if err != nil { + h.sugar.Errorw("Failed to load agent configuration", "agentID", agentID, "error", err) + return ctx.JSON(http.StatusInternalServerError, api.InternalServerError()) + } + rev, rowID := int64(0), uuid.Nil + if head != nil { + rev, rowID = head.Revision, head.ID + } + if etag := agentconfig.ETagForRevision(rev, rowID, agentID); agentconfig.MatchIfNoneMatch(ifNoneMatch, etag) { + setRemoteConfigHeaders(ctx) + ctx.Response().Header().Set(headerETag, etag) + return ctx.NoContent(http.StatusNotModified) + } + } + + cur, err := h.svc.Current(reqCtx, agentID) if err != nil { h.sugar.Errorw("Failed to load agent configuration", "agentID", agentID, "error", err) return ctx.JSON(http.StatusInternalServerError, api.InternalServerError()) @@ -96,7 +117,7 @@ func (h *AgentConfigSyncHandler) GetConfig(ctx echo.Context) error { setRemoteConfigHeaders(ctx) ctx.Response().Header().Set(headerETag, etag) - if agentconfig.MatchIfNoneMatch(ctx.Request().Header.Get("If-None-Match"), etag) { + if agentconfig.MatchIfNoneMatch(ifNoneMatch, etag) { return ctx.NoContent(http.StatusNotModified) } return ctx.JSON(http.StatusOK, GenericDataResponse[agentconfig.OverlayDocument]{Data: doc}) diff --git a/internal/api/handler/agent_config_sync_poll_regression_integration_test.go b/internal/api/handler/agent_config_sync_poll_regression_integration_test.go new file mode 100644 index 00000000..754ee087 --- /dev/null +++ b/internal/api/handler/agent_config_sync_poll_regression_integration_test.go @@ -0,0 +1,46 @@ +//go:build integration + +package handler + +import ( + "net/http" + "sync" + + "gorm.io/gorm" +) + +// Regression (review #479, fp 935fac26ca05): a poll that ends in 304 must not read the +// overlay column; the overlay is only loaded when the agent's ETag is stale. +func (s *AgentConfigSyncIntegrationSuite) TestRegressionNotModifiedPollSkipsOverlay() { + a := s.newAgent("poll-agent") + s.createRevision(*a.agent.ID, 0, `{"verbosity":1}`) + rec := s.getConfig(a.token, "") + s.Require().Equal(http.StatusOK, rec.Code, rec.Body.String()) + etag := rec.Header().Get("ETag") + s.Require().NotEmpty(etag) + + var ( + mu sync.Mutex + queries []string + ) + const cb = "regression:record-revision-queries" + s.Require().NoError(s.DB.Callback().Query().After("gorm:query").Register(cb, func(db *gorm.DB) { + if db.Statement != nil && db.Statement.Table == "ccf_agent_config_revisions" { + mu.Lock() + queries = append(queries, db.Statement.SQL.String()) + mu.Unlock() + } + })) + defer func() { s.Require().NoError(s.DB.Callback().Query().Remove(cb)) }() + + rec = s.getConfig(a.token, etag) + s.Require().Equal(http.StatusNotModified, rec.Code, rec.Body.String()) + + mu.Lock() + defer mu.Unlock() + s.Require().NotEmpty(queries, "the poll reads the current revision") + for _, q := range queries { + s.NotContains(q, "SELECT *", q) + s.NotContains(q, "overlay", q) + } +} diff --git a/internal/service/relational/agentcfg/service.go b/internal/service/relational/agentcfg/service.go index 59852e7a..739ee973 100644 --- a/internal/service/relational/agentcfg/service.go +++ b/internal/service/relational/agentcfg/service.go @@ -146,6 +146,30 @@ func (s *Service) Current(ctx context.Context, agentID uuid.UUID) (*relational.A return &rev, nil } +// RevisionHead identifies a revision without loading its overlay. +type RevisionHead struct { + ID uuid.UUID + Revision int64 + CreatedAt time.Time +} + +// CurrentHead returns the id, revision and creation time of an agent's latest revision, or +// nil (revision 0) when none exists. It is Current's indexed query without the overlay, so +// an agent poll can compute its ETag (and answer 304) without reading the overlay. +func (s *Service) CurrentHead(ctx context.Context, agentID uuid.UUID) (*RevisionHead, error) { + var heads []RevisionHead + err := s.db.WithContext(ctx).Model(&relational.AgentConfigRevision{}). + Select("id, revision, created_at"). + Where("agent_id = ?", agentID). + Order("revision DESC"). + Limit(1). + Find(&heads).Error + if err != nil || len(heads) == 0 { + return nil, err + } + return &heads[0], nil +} + // CurrentRevisionNumber returns the latest revision number, 0 when none exists. func (s *Service) CurrentRevisionNumber(ctx context.Context, agentID uuid.UUID) (int64, error) { var cur int64 diff --git a/internal/service/relational/agentcfg/service_integration_test.go b/internal/service/relational/agentcfg/service_integration_test.go index 00182fbd..850b578f 100644 --- a/internal/service/relational/agentcfg/service_integration_test.go +++ b/internal/service/relational/agentcfg/service_integration_test.go @@ -1055,3 +1055,19 @@ func (s *AgentCfgServiceIntegrationSuite) TestPruneInstances() { s.Require().NoError(err) s.Equal(int64(0), deleted, "idempotent") } + +func (s *AgentCfgServiceIntegrationSuite) TestCurrentHead() { + agentID := s.newAgent("current-head") + head, err := s.svc.CurrentHead(s.ctx, agentID) + s.Require().NoError(err) + s.Nil(head, "no revision yet") + + s.createRevision(agentID, 0, `{"verbosity":1}`) + second := s.createRevision(agentID, 1, `{"verbosity":2}`) + head, err = s.svc.CurrentHead(s.ctx, agentID) + s.Require().NoError(err) + s.Require().NotNil(head) + s.Equal(*second.ID, head.ID) + s.Equal(int64(2), head.Revision) + s.True(second.CreatedAt.Equal(head.CreatedAt)) +}