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
15 changes: 2 additions & 13 deletions cmd/sandbox-apiserver/main.go
Original file line number Diff line number Diff line change
Expand Up @@ -43,10 +43,7 @@ const (
e2bShutdownTimeout = 10 * time.Second
)

// options are the standard aggregated-apiserver options: secure serving plus
// delegated authentication/authorization (token/SAR review against the host
// kube-apiserver). There is deliberately no etcd option — this server stores
// nothing. The sandboxd token wires the node-local claim/release write path.
// options has no etcd option because this server stores nothing.
type options struct {
SecureServing *genericoptions.SecureServingOptionsWithLoopback
Authentication *genericoptions.DelegatingAuthenticationOptions
Expand Down Expand Up @@ -237,15 +234,7 @@ func run() error {
return err
}

// startWarmPoolDriver runs the SandboxWarmPool driver as a controller inside a
// controller-runtime manager: it WATCHES SandboxWarmPool and NodeInventory, so a
// `kubectl apply/patch/delete` reconciles in milliseconds instead of waiting for
// a poll tick (the only latency that ever mattered — the node side fills a pool
// in under a second). Leader election makes exactly one of the apiserver replicas
// drive the pools. The manager's own metrics/health servers are disabled; the
// aggregated apiserver owns the serving port. inv is the process-wide cache-fed
// inventory source; the manager's own client would read NodeInventory
// unstructured and so bypass its cache on every node read.
// startWarmPoolDriver takes the cache-fed inv because the manager client reads NodeInventory unstructured, uncached.
func startWarmPoolDriver(ctx context.Context, fail context.CancelCauseFunc, restCfg *restclient.Config, token string, interval time.Duration, inv scale.InventorySource) error {
scheme := runtime.NewScheme()
if err := extv1beta1.AddToScheme(scheme); err != nil {
Expand Down
21 changes: 3 additions & 18 deletions examples/lifecycle/example.go
Original file line number Diff line number Diff line change
Expand Up @@ -63,9 +63,7 @@ type options struct {
keep bool
}

// e2bClient is a minimal client for the e2b REST contract. The real e2b SDKs
// (JS and Python) speak exactly this and need no changes — point E2B_API_URL at
// the server. This exists because there is no official Go SDK.
// e2bClient hand-rolls the e2b REST contract because e2b has no official Go SDK.
type e2bClient struct {
base string
key string
Expand Down Expand Up @@ -214,9 +212,7 @@ func newScheme() (*runtime.Scheme, error) {
return scheme, nil
}

// newClient builds a controller-runtime client that knows this operator's
// types. Any Kubernetes client works — client-go, the dynamic client, or
// kubectl; nothing here is specific to controller-runtime.
// newClient uses controller-runtime, but any Kubernetes client works.
func newClient(kubeconfig string, scheme *runtime.Scheme) (client.Client, error) {
cfg, err := loadConfig(kubeconfig)
if err != nil {
Expand All @@ -225,9 +221,6 @@ func newClient(kubeconfig string, scheme *runtime.Scheme) (client.Client, error)
return client.New(cfg, client.Options{Scheme: scheme})
}

// discoverTemplate reads the fleet's advertised warm pools and returns a
// template that actually has capacity, so the walk-through does not depend on
// a hard-coded image.
func discoverTemplate(ctx context.Context, c client.Client) (string, error) {
var inventories cocoonv1beta1.NodeInventoryList
if err := c.List(ctx, &inventories); err != nil {
Expand Down Expand Up @@ -325,15 +318,7 @@ func deleteCheckpoints(ctx context.Context, e *e2bClient, ids ...string) error {
return nil
}

// post invokes an action subresource. These are POST-only verbs (the
// pods/eviction shape), which is why they are not fields on SandboxSpec: the
// standard agent-sandbox schema stays untouched, so an unmodified upstream
// client keeps working against this server.
//
// A raw REST client is used rather than controller-runtime's SubResource
// helper because these actions have DIFFERENT request and response types
// (SandboxForkOptions in, SandboxForkResult out); the helper decodes the reply
// back into the object it was given, which cannot express that.
// post uses a raw REST client because the SubResource helper decodes the reply into the request object.
func post(ctx context.Context, rc rest.Interface, ns, name, sub string, body, out runtime.Object) error {
req := rc.Post().
Namespace(ns).
Expand Down
39 changes: 7 additions & 32 deletions pkg/e2bcompat/lifecycle.go
Original file line number Diff line number Diff line change
Expand Up @@ -22,18 +22,12 @@ import (
// cannot serialize a handler into the minutes.
const maxNodeConcurrency = 16

// pauseSandbox hibernates the sandbox: its memory is written out and the VM
// stops, so the cost is proportional to guest RAM. e2b's contract is specific
// about the already-paused case — the SDK reads 409 as "already paused" and
// returns false rather than raising — so that state is reported, not retried.
// pauseSandbox answers 409 when already paused because the e2b SDK reads 409 as "already paused".
func (s *Server) pauseSandbox(w http.ResponseWriter, r *http.Request) {
var req SandboxPauseRequest
if !decodeOptionalBody(w, r, &req) {
return
}
// memory=false asks for a filesystem-only snapshot whose resume cold-boots.
// The node's hibernate always captures memory, so honoring it would mean
// silently giving back a different sandbox than asked for.
if req.Memory != nil && !*req.Memory {
writeError(w, http.StatusBadRequest,
"filesystem-only pause (memory=false) is not supported; this backend always snapshots memory")
Expand Down Expand Up @@ -61,10 +55,7 @@ func (s *Server) pauseSandbox(w http.ResponseWriter, r *http.Request) {
w.WriteHeader(http.StatusNoContent)
}

// connectSandbox is the SDK's resume: it returns the sandbox's connection
// details, restoring it first when paused. 200 means it was already running,
// 201 that it was paused and got resumed — the SDK accepts either, and the
// distinction is what tells an operator whether a restore actually happened.
// connectSandbox is the e2b SDK's resume; the SDK accepts both 200 and 201.
func (s *Server) connectSandbox(w http.ResponseWriter, r *http.Request) {
var req ConnectSandbox
if !decodeOptionalBody(w, r, &req) {
Expand Down Expand Up @@ -114,12 +105,7 @@ func (s *Server) connectSandbox(w http.ResponseWriter, r *http.Request) {
})
}

// forkSandbox branches the sandbox into count children. The parent is
// checkpointed in place and keeps running; every child is a fresh sandbox with
// its own id and lease. Per e2b's contract a partial failure is still a 201
// carrying per-child detail — a non-201 means nothing was attempted — so the
// node's all-or-nothing fork is reported as a whole-request failure only when
// it rejects the request outright.
// forkSandbox sets no per-child error because the node fork is all-or-nothing.
func (s *Server) forkSandbox(w http.ResponseWriter, r *http.Request) {
var req SandboxForkRequest
if !decodeOptionalBody(w, r, &req) {
Expand Down Expand Up @@ -260,7 +246,7 @@ func (s *Server) snapshotsOf(r *http.Request) (snaps []scale.Snapshot, complete
g.Go(func() error {
snaps, err := s.store.Snapshots(r.Context(), node)
if err != nil {
log.WithFunc("e2bcompat.snapshotsOf").Errorf(r.Context(), err, "e2b snapshots: node failed node=%s", node)
log.WithFunc("e2bcompat.snapshotsOf").Warnf(r.Context(), "e2b snapshots: node failed node=%s err=%v", node, err)
return nil
}
answered[i] = true
Expand All @@ -277,9 +263,6 @@ func (s *Server) snapshotsOf(r *http.Request) (snaps []scale.Snapshot, complete
return slices.Concat(perNode...), !slices.Contains(answered, false), nil
}

// sandboxMetrics reports one sandbox's resource usage. e2b's schema requires
// every field, so all are emitted; the ones this backend cannot measure are
// reported as zero rather than invented (see SandboxStats).
func (s *Server) sandboxMetrics(w http.ResponseWriter, r *http.Request) {
id := r.PathValue("sandboxID")
sb, err := s.lookup(r, id)
Expand All @@ -305,10 +288,7 @@ func (s *Server) sandboxMetrics(w http.ResponseWriter, r *http.Request) {
}})
}

// listTemplates reports the pools this fleet can serve claims from. e2b's
// templates are build artifacts with their own lifecycle; the equivalent here
// is the set of warm-pool keys nodes advertise, which is what a caller can
// actually pass as templateID on create.
// listTemplates reports the advertised warm-pool keys, the values create accepts as templateID.
func (s *Server) listTemplates(w http.ResponseWriter, r *http.Request) {
nodes, err := s.inventories(r)
if err != nil {
Expand Down Expand Up @@ -354,9 +334,7 @@ func (s *Server) inventories(r *http.Request) ([]*scale.NodeInventory, error) {
for _, node := range nodes {
inv, err := s.opts.Inventory.NodeInventory(r.Context(), node)
if err != nil {
// A partitioned node is skipped, not fatal — the same rule the
// aggregated read path applies.
log.WithFunc("e2bcompat.inventories").Errorf(r.Context(), err, "e2b: node inventory unavailable node=%s", node)
log.WithFunc("e2bcompat.inventories").Warnf(r.Context(), "e2b: node inventory unavailable node=%s err=%v", node, err)
continue
}
out = append(out, inv)
Expand All @@ -372,10 +350,7 @@ func (s *Server) nodesWithSandboxes(r *http.Request) ([]string, error) {
return s.opts.Inventory.ListNodes(r.Context())
}

// isPaused asks the owning node: the listed phase comes from NodeInventory,
// which lags a pause by up to its publish cadence, and a stale Running would
// turn e2b's 409 "already paused" into a second 204. An unreachable node falls
// back to the cached label.
// isPaused asks the owning node because the NodeInventory phase lags a pause by up to one publish.
func (s *Server) isPaused(ctx context.Context, sb *sandboxv1beta1.Sandbox) (bool, error) {
if node, id := sb.Status.NodeName, claimIDOf(sb); node != "" && id != "" {
rec, err := s.store.Read(ctx, node, id)
Expand Down
43 changes: 10 additions & 33 deletions pkg/e2bcompat/server.go
Original file line number Diff line number Diff line change
Expand Up @@ -40,13 +40,9 @@ import (
)

const (
// DefaultEnvdVersion is reported to the SDK when no version is configured.
// The SDK version-compares this before choosing the envd auth style, so it
// must be a real semver at or above the modern-auth cutoff (0.4.0).
// DefaultEnvdVersion is the modern-auth floor the SDK version-compares when none is configured.
DefaultEnvdVersion = "0.4.0"
// DefaultTimeoutSeconds matches the node's own default lease. The e2b SDK
// defaults to 15s, which reaps a sandbox before a first exec on a cold
// client, so an omitted timeout takes the node's default instead.
// DefaultTimeoutSeconds is the node's default lease; the SDK's own 15s reaps a cold client's sandbox.
DefaultTimeoutSeconds = 300
// apiKeyHeader is the header the e2b SDKs authenticate with.
apiKeyHeader = "X-API-KEY"
Expand All @@ -59,34 +55,19 @@ var (

// Options configures the compat server.
type Options struct {
// Namespace is where a key that names no namespace claims, and where
// anonymous claims land.
// Namespace is where anonymous claims and claims by a key that names no namespace land.
Namespace string
// Domain is echoed as the sandbox `domain`, from which the SDK derives the
// envd host as "{port}-{sandboxID}.{domain}". It is required: a sandbox
// handed out without one has no address its client can reach.
//
// The sandbox ids published here are DNS-label safe (see sandboxid.go), so
// that host form is valid; wildcard DNS and a proxy must still route the host
// or the E2b-Sandbox-Id / E2b-Sandbox-Port headers the SDK sends.
// Domain is the required base domain the SDK derives the envd host from, as "{port}-{sandboxID}.{domain}".
Domain string
// EnvdVersion overrides DefaultEnvdVersion. It must name the envd actually
// installed in the pool's image: the SDK version-compares it and kills the
// sandbox when it cannot parse one.
// EnvdVersion overrides DefaultEnvdVersion and must name the envd installed in the pool's image.
EnvdVersion string
// DefaultTimeoutSeconds overrides DefaultTimeoutSeconds for a create that
// names no timeout, and is the lease a refresh grants.
// DefaultTimeoutSeconds overrides DefaultTimeoutSeconds for a create without a timeout and for a refresh.
DefaultTimeoutSeconds int
// APIKeys, when non-empty, is the set of accepted X-API-KEY values, each
// "key" or "key namespace" (Namespace when none is given); a key sees
// nothing outside its namespace. Empty is refused unless AllowAnonymous.
// APIKeys holds the accepted X-API-KEY values, each "key" or "key namespace"; empty requires AllowAnonymous.
APIKeys []string
// AllowAnonymous permits serving with no API key (local development).
AllowAnonymous bool
// Inventory enumerates the fleet's nodes and their advertised pools. It is
// required by the surfaces that are fleet-wide rather than sandbox-scoped
// (template listing, snapshot listing); without it those report an error
// instead of an empty list, so a missing dependency cannot read as "none".
// Inventory enumerates the fleet's nodes; without it template and snapshot listing fail rather than report none.
Inventory scale.InventorySource
}

Expand Down Expand Up @@ -449,9 +430,7 @@ func (f listFilter) keeps(d SandboxDetail) bool {
return true
}

// templateOf reports the pool template a sandbox was claimed from: the label the
// store stamps, which is the only place it survives (a synthesized Sandbox holds
// no pod spec).
// templateOf reads the store's label because a synthesized Sandbox holds no pod spec.
func templateOf(sb *sandboxv1beta1.Sandbox) string {
return sb.Labels[scale.TemplateLabel]
}
Expand All @@ -464,9 +443,7 @@ func netFor(allowInternet *bool) string {
return scale.NetDefault
}

// unsupportedCreateOption names the first requested option this backend cannot
// honor. Honoring it silently would hand back a different sandbox than asked
// for, which is how an SDK ends up trusting a guarantee that does not hold.
// unsupportedCreateOption names the first requested option this backend cannot honor.
func unsupportedCreateOption(req NewSandbox) (string, bool) {
switch {
case req.Secure != nil && !*req.Secure:
Expand Down
5 changes: 1 addition & 4 deletions pkg/envdproxy/route.go
Original file line number Diff line number Diff line change
Expand Up @@ -29,10 +29,7 @@ type route struct {
port uint16
}

// routeOf reads the destination from the two forms the SDK sends: the derived
// host "{port}-{sandboxID}.{domain}", and the header pair used when every
// sandbox shares one host. Headers win, because a client that sends them has
// already decided the host does not carry the address.
// routeOf prefers the header pair because a client that sends it does not address the sandbox by host.
func routeOf(r *http.Request, domain string) (route, bool) {
if id := strings.TrimSpace(r.Header.Get(sandboxIDHeader)); id != "" {
port, ok := parsePort(r.Header.Get(sandboxPortHeader))
Expand Down
3 changes: 2 additions & 1 deletion pkg/scale/sandboxstore_claim_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -311,6 +311,7 @@ type recordingFactory struct {
releaseCalls int

claimResult sandboxd.ClaimResult
forkResult sandboxd.ForkResult
claimErr error
releaseErr error
verbErr error
Expand Down Expand Up @@ -362,7 +363,7 @@ func (c *recordingClient) Renew(context.Context, string, sandboxd.RenewSpec) (ti

func (c *recordingClient) Fork(_ context.Context, _ string, spec sandboxd.ForkSpec) (sandboxd.ForkResult, error) {
c.f.forkSpec = spec
return sandboxd.ForkResult{}, c.f.verbErr
return c.f.forkResult, c.f.verbErr
}

func (c *recordingClient) Checkpoint(context.Context, string, sandboxd.CheckpointSpec) (sandboxd.Checkpoint, error) {
Expand Down
1 change: 1 addition & 0 deletions pkg/scale/sandboxstore_impl.go
Original file line number Diff line number Diff line change
Expand Up @@ -237,6 +237,7 @@ func (s *scatterGatherStore) Claim(ctx context.Context, namespace, name string,
})
if claimErr == nil {
s.index.remember(nameKey(namespace, name), best.node)
s.index.remember(claimKey(namespace, res.ID), best.node)
return Assignment{SandboxName: res.ID, Node: best.node, Address: res.OwnerAddr, Token: res.Token, Deadline: res.Deadline}, nil
}
if !claimUndelivered(claimErr) {
Expand Down
2 changes: 2 additions & 0 deletions pkg/scale/sandboxstore_lifecycle.go
Original file line number Diff line number Diff line change
Expand Up @@ -65,6 +65,8 @@ func (s *scatterGatherStore) Fork(ctx context.Context, namespace, node, id strin
}
out := make([]Assignment, 0, len(res.Children))
for _, c := range res.Children {
s.index.remember(nameKey(namespace, c.ID), node)
s.index.remember(claimKey(namespace, c.ID), node)
out = append(out, Assignment{
SandboxName: c.ID,
Node: node,
Expand Down
38 changes: 38 additions & 0 deletions pkg/scale/sandboxstore_live_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -68,6 +68,44 @@ func TestGetAsksTheClaimingNodeFirst(t *testing.T) {
assert.Zero(t, src.lists.Load(), "a hit on the claiming node must not sweep the fleet's inventories")
}

func TestGetByClaimIDAsksTheClaimingNodeFirst(t *testing.T) {
f := &recordingFactory{claimResult: sandboxd.ClaimResult{ID: "sb_1", Token: "tok"}}
store, src := unpublishedStore(f)
src.Put(poolInv("n1", "n1:7777", PoolCapacity{Template: "img", Warm: 1, Target: 1}))

_, err := store.Claim(t.Context(), "ns", "s1", PoolKey{Template: "img"}, 0)
require.NoError(t, err)
f.rows = map[string][]sandboxd.SandboxSummary{
"n1:7777": {{ID: "sb_1", ClaimRef: "ns/s1"}},
"n2:7777": {{ID: "sb_2", ClaimRef: "ns/s2"}},
}
src.lists.Store(0)

got, err := store.GetByClaimID(t.Context(), "ns", "sb_1", func(id string) bool { return id == "sb_1" })
require.NoError(t, err)
assert.Equal(t, "n1", got.Status.NodeName)
assert.Equal(t, []string{"n1:7777"}, f.rowReads, "only the node that served the claim is asked")
assert.Zero(t, src.lists.Load(), "a hit on the claiming node must not sweep the fleet's inventories")
}

func TestAForkChildIsLookedUpOnItsNode(t *testing.T) {
f := &recordingFactory{forkResult: sandboxd.ForkResult{Children: []sandboxd.ClaimResult{{ID: "sb_c1"}}}}
store, src := unpublishedStore(f)

_, err := store.Fork(t.Context(), "ns", "n1", "sb_parent", 1, 0)
require.NoError(t, err)
f.rows = map[string][]sandboxd.SandboxSummary{"n1:7777": {{ID: "sb_c1", ClaimRef: "ns/sb_c1"}}}
src.lists.Store(0)

byID, err := store.GetByClaimID(t.Context(), "ns", "sb_c1", func(id string) bool { return id == "sb_c1" })
require.NoError(t, err)
byName, err := store.Get(t.Context(), "ns", "sb_c1")
require.NoError(t, err)
assert.Equal(t, "n1", byID.Status.NodeName)
assert.Equal(t, "n1", byName.Status.NodeName)
assert.Zero(t, src.lists.Load(), "a fork child must not sweep the fleet's inventories")
}

func TestGetFindsWhatAnotherReplicaClaimed(t *testing.T) {
claimed := time.Date(2026, 9, 23, 2, 47, 45, 0, time.UTC)
f := &recordingFactory{rows: map[string][]sandboxd.SandboxSummary{
Expand Down
Loading