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
48 changes: 47 additions & 1 deletion packages/api/internal/clusters/discovery/static.go
Original file line number Diff line number Diff line change
@@ -1,6 +1,13 @@
package discovery

import "context"
import (
"context"
"fmt"
"net"
"strconv"

"github.com/e2b-dev/infra/packages/shared/pkg/consts"
)

// StaticServiceDiscovery returns a fixed (possibly empty) list of items on
// every Query. It is used for local development against the darwin dummy
Expand All @@ -15,6 +22,45 @@ func NewStaticDiscovery(items []Item) Discovery {
return &StaticServiceDiscovery{items: items}
}

// NewStaticFromAddress returns a Discovery holding a single instance reachable
// at addr, which may be "host:port" or just "host" (defaulting to
// consts.OrchestratorAPIPort). It mirrors the orchestrator-side
// orchestrator/discovery.NewLocal.
//
// A local orchestrator can serve the template-builder role as well
// (ORCHESTRATOR_SERVICES=orchestrator,template-manager). Whether it actually
// does is decided by the roles it reports over the Info RPC during instance
// sync, so pointing template-builder discovery at an orchestrator that does
// not run template-manager (e.g. the darwin dummy) is harmless: the instance
// registers with IsBuilder=false and is skipped when picking a builder.
func NewStaticFromAddress(addr string) (Discovery, error) {
host, portStr, err := net.SplitHostPort(addr)
if err != nil {
// Allow plain "host" without a port.
host = addr
portStr = strconv.FormatUint(uint64(consts.OrchestratorAPIPort), 10)
}
if host == "" {
return nil, fmt.Errorf("static discovery: empty host in %q", addr)
}

port, err := strconv.ParseUint(portStr, 10, 16)
if err != nil {
return nil, fmt.Errorf("static discovery: invalid port %q: %w", portStr, err)
}

return NewStaticDiscovery([]Item{
{
UniqueIdentifier: "local",
NodeID: "local",
// Populated during instance sync.
InstanceID: "unknown",
LocalIPAddress: host,
LocalInstanceApiPort: uint16(port),
},
}), nil
}

func (sd *StaticServiceDiscovery) Query(_ context.Context) ([]Item, error) {
if sd.items == nil {
return []Item{}, nil
Expand Down
88 changes: 88 additions & 0 deletions packages/api/internal/clusters/discovery/static_test.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,88 @@
package discovery

import (
"context"
"testing"

"github.com/stretchr/testify/require"

"github.com/e2b-dev/infra/packages/shared/pkg/consts"
)

func TestNewStaticFromAddress(t *testing.T) {
t.Parallel()

tests := []struct {
name string
addr string
wantHost string
wantPort uint16
wantErr bool
}{
{
name: "host with port",
addr: "127.0.0.1:5108",
wantHost: "127.0.0.1",
wantPort: 5108,
},
{
name: "dns host with port",
addr: "orchestrator.internal:5008",
wantHost: "orchestrator.internal",
wantPort: 5008,
},
{
name: "ipv6 host with port",
addr: "[::1]:5008",
wantHost: "::1",
wantPort: 5008,
},
{
name: "host without port defaults to the orchestrator API port",
addr: "127.0.0.1",
wantHost: "127.0.0.1",
wantPort: consts.OrchestratorAPIPort,
},
{
name: "empty address",
addr: "",
wantErr: true,
},
{
name: "empty host",
addr: ":5008",
wantErr: true,
},
{
name: "non-numeric port",
addr: "127.0.0.1:grpc",
wantErr: true,
},
{
name: "port out of range",
addr: "127.0.0.1:70000",
wantErr: true,
},
}

for _, tt := range tests {
t.Run(tt.name, func(t *testing.T) {
t.Parallel()

sd, err := NewStaticFromAddress(tt.addr)
if tt.wantErr {
require.Error(t, err)

return
}

require.NoError(t, err)

items, err := sd.Query(context.Background())
require.NoError(t, err)
require.Len(t, items, 1)
require.Equal(t, tt.wantHost, items[0].LocalIPAddress)
require.Equal(t, tt.wantPort, items[0].LocalInstanceApiPort)
})
}
}
13 changes: 10 additions & 3 deletions packages/api/internal/handlers/store.go
Original file line number Diff line number Diff line change
Expand Up @@ -179,9 +179,16 @@ func NewAPIStore(ctx context.Context, tel *telemetry.Client, redisClient redis.U
logger.L().Fatal(ctx, "Initializing local orchestrator discovery", zap.Error(localErr))
}
nodeDiscovery = localND
// No template builders in local dev — return an empty list so the
// clusters pool initializes cleanly.
templateBuilderDiscovery = clustersdiscovery.NewStaticDiscovery(nil)
// The local orchestrator doubles as the template builder when it is
// started with ORCHESTRATOR_SERVICES=orchestrator,template-manager, so
// point builder discovery at the same address. Instances that do not
// report the TemplateBuilder role (the darwin dummy orchestrator) are
// registered with IsBuilder=false and never selected for builds, which
// keeps the dummy setup behaving as before.
templateBuilderDiscovery, err = clustersdiscovery.NewStaticFromAddress(config.LocalOrchestratorAddress)
if err != nil {
logger.L().Fatal(ctx, "Initializing local template builder discovery", zap.Error(err))
}
default: // ServiceDiscoveryProviderNomad
nomadClient, nomadErr := nomadapi.NewClient(&nomadapi.Config{
Address: config.NomadAddress,
Expand Down
23 changes: 23 additions & 0 deletions packages/api/internal/orchestrator/client.go
Original file line number Diff line number Diff line change
Expand Up @@ -43,10 +43,33 @@ func (o *Orchestrator) connectToNode(ctx context.Context, discovered nodemanager
return err
}

// registersClusterOrchestrators reports whether instances discovered through
// the clusters registry of the given cluster may be registered as orchestrator
// nodes.
//
// Unless the node discovery loop is disabled (see
// localClusterOwnsOrchestrators), local-cluster orchestrators are owned by the
// node discovery path (connectToNode), which identifies nodes by the ID they
// report over the Info RPC. The local clusters registry only exists to find
// template builders and identifies instances by their discovery item ID, so an
// instance serving both roles — a single process started with
// ORCHESTRATOR_SERVICES=orchestrator,template-manager, as in local dev — would
// otherwise register twice under two different node IDs and have its capacity
// and sandboxes counted twice.
//
// Remote clusters are always registered from their own registry.
func (o *Orchestrator) registersClusterOrchestrators(clusterID uuid.UUID) bool {
return clusterID != consts.LocalClusterID || o.localClusterOwnsOrchestrators
}

func (o *Orchestrator) connectToClusterNode(ctx context.Context, cluster *clusters.Cluster, i *clusters.Instance) {
ctx, span := tracer.Start(ctx, "connect-to-cluster-node")
defer span.End()

if !o.registersClusterOrchestrators(cluster.ID) {
return
}
Comment thread
cursor[bot] marked this conversation as resolved.

// connectGroup is keyed by scopedNodeID so that concurrent callers targeting
// the same cluster instance share a single dial attempt.
scopedKey := o.scopedNodeID(cluster.ID, i.NodeID)
Expand Down
69 changes: 69 additions & 0 deletions packages/api/internal/orchestrator/client_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -21,6 +21,7 @@ import (
"google.golang.org/protobuf/types/known/emptypb"

"github.com/e2b-dev/infra/packages/api/internal/api"
"github.com/e2b-dev/infra/packages/api/internal/clusters"
"github.com/e2b-dev/infra/packages/api/internal/orchestrator/discovery"
"github.com/e2b-dev/infra/packages/api/internal/orchestrator/nodemanager"
"github.com/e2b-dev/infra/packages/shared/pkg/consts"
Expand Down Expand Up @@ -336,3 +337,71 @@ func TestRegisterNode_NoDuplicates(t *testing.T) {
wg.Wait()
assert.Equal(t, 5, o.nodes.Count())
}

// TestRegistersClusterOrchestrators covers which discovery path owns
// orchestrator nodes. Local-cluster instances are owned by the node discovery
// path (connectToNode) and must not be registered a second time from the
// clusters registry — except when the node discovery loop is disabled
// (ENVIRONMENT=local), where the clusters registry is the only source of
// orchestrator nodes and skipping it would leave the API with zero nodes.
func TestRegistersClusterOrchestrators(t *testing.T) {
t.Parallel()

tests := []struct {
name string
clusterID uuid.UUID
localClusterOwnsOrchestrators bool
want bool
}{
{
name: "local cluster with node discovery running",
clusterID: consts.LocalClusterID,
localClusterOwnsOrchestrators: false,
want: false,
},
{
name: "local cluster without node discovery",
clusterID: consts.LocalClusterID,
localClusterOwnsOrchestrators: true,
want: true,
},
{
name: "remote cluster with node discovery running",
clusterID: uuid.New(),
localClusterOwnsOrchestrators: false,
want: true,
},
{
name: "remote cluster without node discovery",
clusterID: uuid.New(),
localClusterOwnsOrchestrators: true,
want: true,
},
}

for _, tt := range tests {
t.Run(tt.name, func(t *testing.T) {
t.Parallel()

o := newTestOrchestrator(t, nil)
o.localClusterOwnsOrchestrators = tt.localClusterOwnsOrchestrators

assert.Equal(t, tt.want, o.registersClusterOrchestrators(tt.clusterID))
})
}
}

// TestConnectToClusterNode_SkipsLocalCluster verifies that the ownership check
// short-circuits connectToClusterNode before it touches the instance. The nil
// instance is intentional: it makes a regression to unconditional registration
// fail loudly instead of silently duplicating the node.
func TestConnectToClusterNode_SkipsLocalCluster(t *testing.T) {
t.Parallel()

o := newTestOrchestrator(t, nil)
o.localClusterOwnsOrchestrators = false

o.connectToClusterNode(t.Context(), &clusters.Cluster{ID: consts.LocalClusterID}, nil)

assert.Zero(t, o.nodes.Count())
}
26 changes: 23 additions & 3 deletions packages/api/internal/orchestrator/orchestrator.go
Original file line number Diff line number Diff line change
Expand Up @@ -67,6 +67,21 @@ type Orchestrator struct {
snapshotUpsertSem *utils.AdjustableSemaphore
redisStorage *redisbackend.Storage

// localClusterOwnsOrchestrators makes connectToClusterNode register
// local-cluster instances that report the Orchestrator role as nodes.
//
// It is only set when the node discovery loop is disabled (see
// skipNomadSync in New), which is the case in the local environment: there
// the local clusters registry is the single source of orchestrator nodes.
//
// Otherwise local-cluster orchestrators are owned by the node discovery
// path (connectToNode), which identifies nodes by the ID they report over
// the Info RPC, while the clusters registry identifies instances by their
// discovery item ID. An instance serving both the orchestrator and the
// template-builder role would then register twice under two different node
// IDs and have its capacity and sandboxes counted twice.
localClusterOwnsOrchestrators bool

// connectGroup deduplicates concurrent dial+register attempts for the same
// physical node. It is keyed by NomadNodeShortID (Nomad-managed nodes) or
// scopedNodeID(clusterID, instanceNodeID) (cluster nodes) and is held inside
Expand Down Expand Up @@ -133,6 +148,10 @@ func New(
}
go redisStorage.Start(ctx)

// For local development and testing, we skip the Nomad sync
// Local cluster is used for single-node setups instead
skipNomadSync := env.IsLocal()

o := Orchestrator{
httpClient: httpClient,
analytics: analyticsInstance,
Expand All @@ -152,6 +171,10 @@ func New(
createdCounter: createdCounter,

snapshotUpsertSem: snapshotUpsertSem,

// Without the node discovery loop, the local clusters registry is the
// only source of orchestrator nodes.
localClusterOwnsOrchestrators: skipNomadSync,
}

o.sandboxStore = sandbox.NewStore(
Expand Down Expand Up @@ -180,9 +203,6 @@ func New(

o.teamMetricsObserver = teamMetricsObserver

// For local development and testing, we skip the Nomad sync
// Local cluster is used for single-node setups instead
skipNomadSync := env.IsLocal()
go o.keepInSync(ctx, o.sandboxStore, skipNomadSync)

if err := o.setupMetrics(tel.MeterProvider); err != nil {
Expand Down
Loading