Skip to content

Commit 8a73f02

Browse files
dobrace2b-bot[bot]
authored andcommitted
refactor(api): retire the api's second discovery mechanism
GitOrigin-RevId: 136062996f2a78dd6cf552ddd817ef59c4e723e8
1 parent 8a1c824 commit 8a73f02

22 files changed

Lines changed: 590 additions & 572 deletions

‎packages/api/go.mod‎

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -56,8 +56,6 @@ require (
5656
golang.org/x/sync v0.22.0
5757
google.golang.org/grpc v1.83.0
5858
google.golang.org/protobuf v1.36.11
59-
k8s.io/api v0.35.3
60-
k8s.io/apimachinery v0.35.3
6159
k8s.io/client-go v0.35.3
6260
)
6361

@@ -417,6 +415,8 @@ require (
417415
gopkg.in/inf.v0 v0.9.1 // indirect
418416
gopkg.in/yaml.v2 v2.4.0 // indirect
419417
gopkg.in/yaml.v3 v3.0.1 // indirect
418+
k8s.io/api v0.35.3 // indirect
419+
k8s.io/apimachinery v0.35.3 // indirect
420420
k8s.io/klog/v2 v2.140.0 // indirect
421421
k8s.io/kube-openapi v0.0.0-20250910181357-589584f1c912 // indirect
422422
k8s.io/utils v0.0.0-20251002143259-bc988d571ff4 // indirect

‎packages/api/internal/cfg/model.go‎

Lines changed: 7 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -46,11 +46,15 @@ type Config struct {
4646
LokiURL string `env:"LOKI_URL,required"`
4747
LokiUser string `env:"LOKI_USER"`
4848

49-
// ServiceDiscoveryProvider selects how the API discovers orchestrator and template-manager instances.
49+
// ServiceDiscoveryProvider selects how the API discovers orchestrator and
50+
// template-manager instances. Left unset it is nomad, except in a local
51+
// environment, which has no Nomad agent — see serviceDiscoveryProvider. It
52+
// carries no envDefault precisely so that "unset" stays distinguishable
53+
// from "set to nomad", which an operator running a local Nomad agent needs.
5054
// Allowed values:
5155
// "nomad" (default) - query the local Nomad agent's HTTP API.
5256
// "kubernetes" - list pods via the in-cluster K8s API.
53-
ServiceDiscoveryProvider string `env:"SERVICE_DISCOVERY_PROVIDER" envDefault:"nomad"`
57+
ServiceDiscoveryProvider string `env:"SERVICE_DISCOVERY_PROVIDER"`
5458

5559
NomadAddress string `env:"NOMAD_ADDRESS" envDefault:"http://localhost:4646"`
5660
NomadToken string `env:"NOMAD_TOKEN"`
@@ -271,7 +275,7 @@ func Parse() (Config, error) {
271275
config.AuthDBConnectionString = config.PostgresConnectionString
272276
}
273277

274-
if !slices.Contains([]string{ServiceDiscoveryProviderNomad, ServiceDiscoveryProviderKubernetes, ServiceDiscoveryProviderLocal}, config.ServiceDiscoveryProvider) {
278+
if !slices.Contains([]string{"", ServiceDiscoveryProviderNomad, ServiceDiscoveryProviderKubernetes, ServiceDiscoveryProviderLocal}, config.ServiceDiscoveryProvider) {
275279
return config, newFailureError(
276280
FailureConditionInvalidServiceDiscoveryProvider,
277281
fmt.Sprintf("invalid service discovery provider: %s", config.ServiceDiscoveryProvider),

‎packages/api/internal/clusters/cluster.go‎

Lines changed: 67 additions & 18 deletions
Original file line numberDiff line numberDiff line change
@@ -5,9 +5,7 @@ import (
55
"errors"
66
"fmt"
77
"math/rand/v2"
8-
"net"
98
"net/http"
10-
"strconv"
119
"sync"
1210
"time"
1311

@@ -16,7 +14,6 @@ import (
1614
"go.uber.org/zap"
1715

1816
"github.com/e2b-dev/infra/packages/api/internal/cfg"
19-
"github.com/e2b-dev/infra/packages/api/internal/clusters/discovery"
2017
clickhouse "github.com/e2b-dev/infra/packages/clickhouse/pkg"
2118
"github.com/e2b-dev/infra/packages/shared/pkg/consts"
2219
"github.com/e2b-dev/infra/packages/shared/pkg/featureflags"
@@ -25,6 +22,7 @@ import (
2522
"github.com/e2b-dev/infra/packages/shared/pkg/logger"
2623
"github.com/e2b-dev/infra/packages/shared/pkg/logs/loki"
2724
"github.com/e2b-dev/infra/packages/shared/pkg/machineinfo"
25+
"github.com/e2b-dev/infra/packages/shared/pkg/servicediscovery"
2826
"github.com/e2b-dev/infra/packages/shared/pkg/smap"
2927
"github.com/e2b-dev/infra/packages/shared/pkg/synchronization"
3028
"github.com/e2b-dev/infra/packages/shared/pkg/telemetry"
@@ -43,7 +41,7 @@ type Cluster struct {
4341
AuthOrgID string
4442

4543
instances *smap.Map[*Instance]
46-
synchronization *synchronization.Synchronize[discovery.Item, *Instance]
44+
synchronization *synchronization.Synchronize[servicediscovery.Instance, *Instance]
4745
resources ClusterResource
4846
}
4947

@@ -57,7 +55,7 @@ func NewCluster(
5755
domain *string,
5856
authOrgID string,
5957
sandboxes *smap.Map[*Instance],
60-
synchronization *synchronization.Synchronize[discovery.Item, *Instance],
58+
synchronization *synchronization.Synchronize[servicediscovery.Instance, *Instance],
6159
resources ClusterResource,
6260
) *Cluster {
6361
return &Cluster{
@@ -73,7 +71,7 @@ func NewCluster(
7371
func newLocalCluster(
7472
ctx context.Context,
7573
tel *telemetry.Client,
76-
storeDiscovery discovery.Discovery,
74+
storeDiscovery servicediscovery.Discoverer,
7775
clickhouse clickhouse.Clickhouse,
7876
queryLogsProvider *loki.LokiQueryProvider,
7977
sandboxLogsReader ClickhouseLogsReader,
@@ -83,9 +81,9 @@ func newLocalCluster(
8381
clusterID := consts.LocalClusterID
8482

8583
instances := smap.New[*Instance]()
86-
instanceCreation := func(ctx context.Context, item discovery.Item) (*Instance, error) {
84+
instanceCreation := func(ctx context.Context, item servicediscovery.Instance) (*Instance, error) {
8785
// For local cluster we are doing direct connection to instance IP and API port and without additional cluster auth.
88-
return newInstance(ctx, tel, nil, clusterID, item, net.JoinHostPort(item.LocalIPAddress, strconv.FormatUint(uint64(item.LocalInstanceApiPort), 10)), false)
86+
return newInstance(ctx, tel, nil, clusterID, item, item.Address(), false)
8987
}
9088

9189
store := instancesSyncStore{clusterID: clusterID, instances: instances, discovery: storeDiscovery, instanceCreation: instanceCreation}
@@ -105,6 +103,14 @@ func newLocalCluster(
105103
return c
106104
}
107105

106+
// remoteInstanceAuthorization routes a call through the remote proxy to one
107+
// instance. The discovery ID is that instance's service id: the edge API
108+
// reports it per run, which is what the proxy keys on — the machine would send
109+
// the call to whatever is on it now.
110+
func remoteInstanceAuthorization(secret string, tls bool, sd servicediscovery.Instance) *instanceAuthorization {
111+
return &instanceAuthorization{secret: secret, tls: tls, serviceInstanceID: sd.WorkloadID}
112+
}
113+
108114
func newRemoteCluster(
109115
ctx context.Context,
110116
tel *telemetry.Client,
@@ -142,14 +148,13 @@ func newRemoteCluster(
142148
}
143149

144150
instances := smap.New[*Instance]()
145-
instanceCreation := func(ctx context.Context, item discovery.Item) (*Instance, error) {
146-
// For remote cluster we are doing connection to endpoint that works as gRPC proxy and handles auth and routing for us.
147-
auth := &instanceAuthorization{secret: secret, tls: endpointTLS, serviceInstanceID: item.InstanceID}
148-
149-
return newInstance(ctx, tel, auth, clusterID, item, endpoint, endpointTLS)
151+
instanceCreation := func(ctx context.Context, item servicediscovery.Instance) (*Instance, error) {
152+
// A remote cluster is reached through an endpoint acting as a gRPC proxy,
153+
// which handles auth and routing for us.
154+
return newInstance(ctx, tel, remoteInstanceAuthorization(secret, endpointTLS, item), clusterID, item, endpoint, endpointTLS)
150155
}
151156

152-
storeDiscovery := discovery.NewRemoteServiceDiscovery(clusterID, httpClient)
157+
storeDiscovery := servicediscovery.NewRemote(httpClient)
153158
store := instancesSyncStore{clusterID: clusterID, instances: instances, instanceCreation: instanceCreation, discovery: storeDiscovery}
154159

155160
c := NewCluster(
@@ -191,17 +196,61 @@ func (c *Cluster) Close(ctx context.Context) error {
191196
return nil
192197
}
193198

199+
// Scanned rather than looked up: the pool is keyed by the running process, and
200+
// a build persists the machine it was placed on (env_builds.cluster_node_id),
201+
// so this resolves the machine to whichever process is on it now.
194202
func (c *Cluster) GetTemplateBuilderByNodeID(nodeID string) (*Instance, error) {
195-
instance, found := c.instances.Get(nodeID)
203+
instance, found := builderOnNode(c.instances, nodeID)
196204
if !found {
197205
return nil, ErrTemplateBuilderNotFound
198206
}
199207

200-
if info := instance.GetInfo(); info.Status == infogrpc.ServiceInfoStatus_Unhealthy || !info.IsBuilder {
201-
return nil, ErrTemplateBuilderNotFound
208+
return instance, nil
209+
}
210+
211+
// instanceOnNode finds an instance running on a machine, preferring one that is
212+
// not unhealthy. Like builderOnNode this cannot be a map lookup, but it applies
213+
// no further filter: callers that only need to reach the machine should not
214+
// start rejecting instances the old lookup accepted.
215+
func instanceOnNode(instances *smap.Map[*Instance], nodeID string) (*Instance, bool) {
216+
var unhealthy *Instance
217+
218+
for _, instance := range instances.Items() {
219+
if instance.NodeID != nodeID {
220+
continue
221+
}
222+
223+
if instance.GetInfo().Status == infogrpc.ServiceInfoStatus_Unhealthy {
224+
unhealthy = instance
225+
226+
continue
227+
}
228+
229+
return instance, true
202230
}
203231

204-
return instance, nil
232+
return unhealthy, unhealthy != nil
233+
}
234+
235+
// builderOnNode finds the usable template builder running on a machine. The
236+
// pool is keyed by process, so this cannot be a map lookup, and every entry on
237+
// the machine has to be considered rather than the first: a restart leaves the
238+
// dead run and the live one sharing a machine for a sync round, and stopping
239+
// early would answer with whichever the map happened to yield.
240+
func builderOnNode(instances *smap.Map[*Instance], nodeID string) (*Instance, bool) {
241+
for _, instance := range instances.Items() {
242+
if instance.NodeID != nodeID {
243+
continue
244+
}
245+
246+
if info := instance.GetInfo(); info.Status == infogrpc.ServiceInfoStatus_Unhealthy || !info.IsBuilder {
247+
continue
248+
}
249+
250+
return instance, true
251+
}
252+
253+
return nil, false
205254
}
206255

207256
func (c *Cluster) GetByServiceInstanceID(serviceInstanceID string) (*Instance, bool) {
Lines changed: 134 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,134 @@
1+
package clusters
2+
3+
import (
4+
"testing"
5+
6+
"github.com/google/uuid"
7+
"github.com/stretchr/testify/assert"
8+
"github.com/stretchr/testify/require"
9+
10+
infogrpc "github.com/e2b-dev/infra/packages/shared/pkg/grpc/orchestrator-info"
11+
"github.com/e2b-dev/infra/packages/shared/pkg/smap"
12+
)
13+
14+
// clusterWithInstances builds only the fields the builder lookup reads.
15+
func clusterWithInstances(t *testing.T, instances ...*Instance) *Cluster {
16+
t.Helper()
17+
18+
pool := smap.New[*Instance]()
19+
for _, i := range instances {
20+
pool.Insert(i.workloadID, i)
21+
}
22+
23+
return &Cluster{ID: uuid.New(), instances: pool}
24+
}
25+
26+
func builderOn(nodeID, processID string) *Instance {
27+
return &Instance{
28+
workloadID: processID,
29+
NodeID: nodeID,
30+
isBuilder: true,
31+
status: infogrpc.ServiceInfoStatus_Healthy,
32+
}
33+
}
34+
35+
// A build persists the machine it was placed on (env_builds.cluster_node_id)
36+
// and re-resolves it later. The pool is keyed by process now, so this lookup
37+
// has to find the machine some other way or every in-flight build breaks.
38+
func TestGetTemplateBuilderByNodeID_ResolvesAPersistedMachine(t *testing.T) {
39+
t.Parallel()
40+
41+
wanted := builderOn("node-b", "alloc-2")
42+
c := clusterWithInstances(t, builderOn("node-a", "alloc-1"), wanted, builderOn("node-c", "alloc-3"))
43+
44+
got, err := c.GetTemplateBuilderByNodeID("node-b")
45+
require.NoError(t, err)
46+
assert.Same(t, wanted, got)
47+
}
48+
49+
// After a restart the machine is the same and the process is not, which is the
50+
// case the re-keying exists for: the persisted node id must resolve to whatever
51+
// is running there now.
52+
func TestGetTemplateBuilderByNodeID_ResolvesToTheCurrentProcessAfterARestart(t *testing.T) {
53+
t.Parallel()
54+
55+
replacement := builderOn("node-a", "alloc-2")
56+
c := clusterWithInstances(t, replacement)
57+
58+
got, err := c.GetTemplateBuilderByNodeID("node-a")
59+
require.NoError(t, err)
60+
assert.Equal(t, "alloc-2", got.workloadID)
61+
}
62+
63+
func TestGetTemplateBuilderByNodeID_RejectsAMachineItDoesNotHave(t *testing.T) {
64+
t.Parallel()
65+
66+
c := clusterWithInstances(t, builderOn("node-a", "alloc-1"))
67+
68+
_, err := c.GetTemplateBuilderByNodeID("node-z")
69+
require.ErrorIs(t, err, ErrTemplateBuilderNotFound)
70+
}
71+
72+
func TestGetTemplateBuilderByNodeID_RejectsANonBuilderAndAnUnhealthyOne(t *testing.T) {
73+
t.Parallel()
74+
75+
notABuilder := builderOn("node-a", "alloc-1")
76+
notABuilder.isBuilder = false
77+
78+
unhealthy := builderOn("node-b", "alloc-2")
79+
unhealthy.status = infogrpc.ServiceInfoStatus_Unhealthy
80+
81+
c := clusterWithInstances(t, notABuilder, unhealthy)
82+
83+
_, err := c.GetTemplateBuilderByNodeID("node-a")
84+
require.ErrorIs(t, err, ErrTemplateBuilderNotFound)
85+
86+
_, err = c.GetTemplateBuilderByNodeID("node-b")
87+
require.ErrorIs(t, err, ErrTemplateBuilderNotFound)
88+
}
89+
90+
// A restart leaves the dead run and the live one on the same machine for a sync
91+
// round — the state the pool re-keying deliberately allows. Answering with
92+
// whichever entry the map yielded first sent builds to the dead one, or
93+
// reported the machine as having no builder at all.
94+
//
95+
// Repeated because the pool is a map: one pass only exercises whichever order
96+
// Go happened to randomise into, and the wrong answer is the order-dependent
97+
// one.
98+
func TestGetTemplateBuilderByNodeID_SkipsADeadRunSharingTheMachine(t *testing.T) {
99+
t.Parallel()
100+
101+
dead := builderOn("node-a", "alloc-1")
102+
dead.status = infogrpc.ServiceInfoStatus_Unhealthy
103+
notABuilder := builderOn("node-a", "alloc-3")
104+
notABuilder.isBuilder = false
105+
live := builderOn("node-a", "alloc-2")
106+
107+
c := clusterWithInstances(t, dead, notABuilder, live)
108+
109+
for range 200 {
110+
got, err := c.GetTemplateBuilderByNodeID("node-a")
111+
require.NoError(t, err, "the live builder must be found whatever order the pool yields")
112+
require.Equal(t, "alloc-2", got.workloadID)
113+
}
114+
}
115+
116+
// Build-log streaming resolves the builder from the machine the build was
117+
// placed on, which is what the database persists. The pool is keyed by process
118+
// now, so a map lookup on the machine finds nothing on Nomad and Kubernetes,
119+
// where an allocation ID or pod name is never a node name — and the failure is
120+
// silent, falling back to persistent logs instead of erroring.
121+
func TestBuildLogs_ResolveTheBuilderByMachineNotByPoolKey(t *testing.T) {
122+
t.Parallel()
123+
124+
builder := builderOn("node-a", "alloc-1")
125+
pool := smap.New[*Instance]()
126+
pool.Insert(builder.workloadID, builder)
127+
128+
_, foundByPoolKey := pool.Get("node-a")
129+
require.False(t, foundByPoolKey, "the machine is not the pool key; this is the lookup that broke")
130+
131+
got, found := instanceOnNode(pool, "node-a")
132+
require.True(t, found, "the builder on the machine must still be reachable")
133+
assert.Equal(t, "alloc-1", got.workloadID)
134+
}

‎packages/api/internal/clusters/clusters_sync.go‎

Lines changed: 3 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -9,14 +9,14 @@ import (
99
"go.uber.org/zap"
1010

1111
"github.com/e2b-dev/infra/packages/api/internal/cfg"
12-
"github.com/e2b-dev/infra/packages/api/internal/clusters/discovery"
1312
clickhouse "github.com/e2b-dev/infra/packages/clickhouse/pkg"
1413
"github.com/e2b-dev/infra/packages/db/client"
1514
"github.com/e2b-dev/infra/packages/db/queries"
1615
"github.com/e2b-dev/infra/packages/shared/pkg/consts"
1716
"github.com/e2b-dev/infra/packages/shared/pkg/featureflags"
1817
"github.com/e2b-dev/infra/packages/shared/pkg/logger"
1918
"github.com/e2b-dev/infra/packages/shared/pkg/logs/loki"
19+
"github.com/e2b-dev/infra/packages/shared/pkg/servicediscovery"
2020
"github.com/e2b-dev/infra/packages/shared/pkg/smap"
2121
"github.com/e2b-dev/infra/packages/shared/pkg/synchronization"
2222
"github.com/e2b-dev/infra/packages/shared/pkg/telemetry"
@@ -47,7 +47,7 @@ func NewPool(
4747
ctx context.Context,
4848
tel *telemetry.Client,
4949
db *client.Client,
50-
localDiscovery discovery.Discovery,
50+
localDiscovery servicediscovery.Discoverer,
5151
queryMetricsProvider clickhouse.Clickhouse,
5252
queryLogsProvider *loki.LokiQueryProvider,
5353
sandboxLogsReader ClickhouseLogsReader,
@@ -116,7 +116,7 @@ type clustersSyncStore struct {
116116
tel *telemetry.Client
117117
clusters *smap.Map[*Cluster]
118118
local *queries.Cluster
119-
localDiscovery discovery.Discovery
119+
localDiscovery servicediscovery.Discoverer
120120
queryMetricsProvider clickhouse.Clickhouse
121121
queryLogsProvider *loki.LokiQueryProvider
122122
sandboxLogsReader ClickhouseLogsReader

0 commit comments

Comments
 (0)