Skip to content
18 changes: 18 additions & 0 deletions changelog/fragments/1787924366-sync-enrollment-write.yaml
Original file line number Diff line number Diff line change
@@ -0,0 +1,18 @@
kind: enhancement

summary: Add synchronous enrollment write strategy to prevent ghost agent documents

description: |
Under high enrollment load, a network interruption after fleet-server committed an
agent record but before the response reached the agent could cause the agent to retry
on a different fleet-server instance. If the record was not yet visible to search, the
retry would create a second, duplicate agent document ("ghost agent"). Ghost agents
accumulate over time and do not resolve on their own.

A new feature flag, inputs[].server.feature_flags._sync_enrollment_write (default:
false), selects the enrollment write strategy. When set to true, fleet-server writes
the agent document synchronously with refresh=wait_for, ensuring it is immediately
searchable before the enrollment response is sent. Any retry on any instance finds the
existing record rather than creating a new one.
Comment thread
ycombinator marked this conversation as resolved.
Outdated

component: fleet-server
31 changes: 29 additions & 2 deletions internal/pkg/api/handleEnroll.go
Original file line number Diff line number Diff line change
Expand Up @@ -5,6 +5,7 @@
package api

import (
"bytes"
"context"
"crypto/hmac"
"crypto/pbkdf2"
Expand Down Expand Up @@ -33,6 +34,7 @@ import (
"github.com/elastic/fleet-server/v7/internal/pkg/model"
"github.com/elastic/fleet-server/v7/internal/pkg/rollback"
"github.com/elastic/fleet-server/v7/internal/pkg/sqn"
"github.com/elastic/go-elasticsearch/v8/esapi"

"github.com/gofrs/uuid/v5"
"github.com/hashicorp/go-version"
Expand Down Expand Up @@ -423,7 +425,7 @@ func (et *EnrollerT) _enroll(
ReplaceToken: replaceHash,
}

err = createFleetAgent(ctx, et.bulker, agentID, agent)
err = createFleetAgent(ctx, et.bulker, agentID, agent, et.cfg.Features.SyncEnrollmentWrite)
if err != nil {
return nil, err
}
Expand Down Expand Up @@ -651,7 +653,7 @@ func updateFleetAgent(ctx context.Context, bulker bulk.Bulk, id string, doc bulk
return bulker.Update(ctx, dl.FleetAgents, id, body, bulk.WithRefresh(), bulk.WithRetryOnConflict(3))
}

func createFleetAgent(ctx context.Context, bulker bulk.Bulk, id string, agent model.Agent) error {
func createFleetAgent(ctx context.Context, bulker bulk.Bulk, id string, agent model.Agent, syncWrite bool) error {
Comment thread
ycombinator marked this conversation as resolved.
span, ctx := apm.StartSpan(ctx, "createAgent", "create")
defer span.End()

Expand All @@ -660,6 +662,31 @@ func createFleetAgent(ctx context.Context, bulker bulk.Bulk, id string, agent mo
return err
}

if syncWrite {
// Sync path: write directly with refresh=wait_for so the document is immediately
// visible to any retry pod's FindAgent search, preventing ghost agents.
req := esapi.IndexRequest{
Index: dl.FleetAgents,
DocumentID: id,
Body: bytes.NewReader(data),
OpType: "create",
Refresh: "wait_for",
}
res, err := req.Do(ctx, bulker.Client())
if err != nil {
return err
}
defer res.Body.Close()
if res.StatusCode == http.StatusConflict {
zerolog.Ctx(ctx).Debug().Str("agent_id", id).Msg("agent document already exists on enrollment create, treating as success")
return nil
}
if res.IsError() {
return fmt.Errorf("createFleetAgent: %s", res.String())
}
return nil
Comment thread
ycombinator marked this conversation as resolved.
}

_, err = bulker.Create(ctx, dl.FleetAgents, id, data, bulk.WithRefresh())
// A 409 on op_type:create is a definitive signal the document already exists; treat as success.
if errors.Is(err, es.ErrElasticVersionConflict) {
Expand Down
44 changes: 43 additions & 1 deletion internal/pkg/api/handleEnroll_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -11,13 +11,17 @@ import (
"encoding/json"
"errors"
"fmt"
"io"
"net/http"
"reflect"
"strings"
"testing"

"github.com/elastic/go-elasticsearch/v8"
"github.com/rs/zerolog"
"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/mock"
"github.com/stretchr/testify/require"

"github.com/elastic/fleet-server/v7/internal/pkg/apikey"
"github.com/elastic/fleet-server/v7/internal/pkg/bulk"
Expand Down Expand Up @@ -600,10 +604,48 @@ func TestCreateFleetAgentVersionConflictSucceeds(t *testing.T) {
bulker.On("Create", mock.Anything, mock.Anything, mock.Anything, mock.Anything, mock.Anything).
Return("", es.ErrElasticVersionConflict)

err := createFleetAgent(t.Context(), bulker, "test-agent-id", model.Agent{})
err := createFleetAgent(t.Context(), bulker, "test-agent-id", model.Agent{}, false)
assert.NoError(t, err)
}

func TestCreateFleetAgentSyncWrite409Succeeds(t *testing.T) {
mt := &MockTransport{}
mt.RoundTripFn = func(req *http.Request) (*http.Response, error) {
return &http.Response{
StatusCode: http.StatusConflict,
Body: io.NopCloser(strings.NewReader(`{}`)),
Header: http.Header{"X-Elastic-Product": []string{"Elasticsearch"}},
}, nil
}
Comment thread
ycombinator marked this conversation as resolved.
cli, err := elasticsearch.NewClient(elasticsearch.Config{Transport: mt})
require.NoError(t, err)

bulker := ftesting.NewMockBulk()
bulker.On("Client").Return(cli)

err = createFleetAgent(t.Context(), bulker, "test-agent-id", model.Agent{}, true)
require.NoError(t, err)
}

func TestCreateFleetAgentSyncWriteErrorSurfaces(t *testing.T) {
mt := &MockTransport{}
mt.RoundTripFn = func(req *http.Request) (*http.Response, error) {
return &http.Response{
StatusCode: http.StatusInternalServerError,
Body: io.NopCloser(strings.NewReader(`{}`)),
Header: http.Header{"X-Elastic-Product": []string{"Elasticsearch"}},
}, nil
}
cli, err := elasticsearch.NewClient(elasticsearch.Config{Transport: mt})
require.NoError(t, err)

bulker := ftesting.NewMockBulk()
bulker.On("Client").Return(cli)

err = createFleetAgent(t.Context(), bulker, "test-agent-id", model.Agent{}, true)
require.Error(t, err)
}

func TestValidateEnrollRequest(t *testing.T) {
t.Run("invalid json", func(t *testing.T) {
req, err := validateRequest(context.Background(), strings.NewReader("not a json"))
Expand Down
6 changes: 6 additions & 0 deletions internal/pkg/config/input.go
Original file line number Diff line number Diff line change
Expand Up @@ -123,6 +123,12 @@ type (
// GracefulForceUnenroll configures the three-step escalation applied to agents that
// check in with an invalid or disabled API key.
GracefulForceUnenroll GracefulForceUnenrollConfig `config:"graceful_force_unenroll"`

// SyncEnrollmentWrite (config key: feature_flags._sync_enrollment_write) selects the enrollment write strategy.
// When false (default), agent documents are written via the async bulk queue — existing behaviour.
Comment thread
ycombinator marked this conversation as resolved.
Outdated
// When true, agent documents are written synchronously with refresh=wait_for, making the document
// immediately searchable before the enrollment response is sent and eliminating ghost agents.
SyncEnrollmentWrite bool `config:"_sync_enrollment_write"`
}

// GracefulForceUnenrollConfig controls the graceful-force-unenroll feature.
Expand Down
Loading