diff --git a/cmd/sandbox-apiserver/main.go b/cmd/sandbox-apiserver/main.go index f1302e7..019b46b 100644 --- a/cmd/sandbox-apiserver/main.go +++ b/cmd/sandbox-apiserver/main.go @@ -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 @@ -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 { diff --git a/examples/lifecycle/example.go b/examples/lifecycle/example.go index c441613..82add92 100644 --- a/examples/lifecycle/example.go +++ b/examples/lifecycle/example.go @@ -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 @@ -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 { @@ -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 { @@ -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). diff --git a/pkg/e2bcompat/lifecycle.go b/pkg/e2bcompat/lifecycle.go index 3925060..be87d72 100644 --- a/pkg/e2bcompat/lifecycle.go +++ b/pkg/e2bcompat/lifecycle.go @@ -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") @@ -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) { @@ -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) { @@ -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 @@ -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) @@ -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 { @@ -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) @@ -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) diff --git a/pkg/e2bcompat/server.go b/pkg/e2bcompat/server.go index 297fd17..bd4285e 100644 --- a/pkg/e2bcompat/server.go +++ b/pkg/e2bcompat/server.go @@ -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" @@ -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 } @@ -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] } @@ -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: diff --git a/pkg/envdproxy/route.go b/pkg/envdproxy/route.go index d56f9be..d5c1329 100644 --- a/pkg/envdproxy/route.go +++ b/pkg/envdproxy/route.go @@ -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)) diff --git a/pkg/scale/sandboxstore_claim_test.go b/pkg/scale/sandboxstore_claim_test.go index 5ae04e4..fb44502 100644 --- a/pkg/scale/sandboxstore_claim_test.go +++ b/pkg/scale/sandboxstore_claim_test.go @@ -311,6 +311,7 @@ type recordingFactory struct { releaseCalls int claimResult sandboxd.ClaimResult + forkResult sandboxd.ForkResult claimErr error releaseErr error verbErr error @@ -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) { diff --git a/pkg/scale/sandboxstore_impl.go b/pkg/scale/sandboxstore_impl.go index 1195c75..89ae4a0 100644 --- a/pkg/scale/sandboxstore_impl.go +++ b/pkg/scale/sandboxstore_impl.go @@ -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) { diff --git a/pkg/scale/sandboxstore_lifecycle.go b/pkg/scale/sandboxstore_lifecycle.go index d83dcce..317cd7a 100644 --- a/pkg/scale/sandboxstore_lifecycle.go +++ b/pkg/scale/sandboxstore_lifecycle.go @@ -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, diff --git a/pkg/scale/sandboxstore_live_test.go b/pkg/scale/sandboxstore_live_test.go index f17100d..e971485 100644 --- a/pkg/scale/sandboxstore_live_test.go +++ b/pkg/scale/sandboxstore_live_test.go @@ -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{