From a74ef7dcc1bc3db18b486a2c468c4ecdfced560a Mon Sep 17 00:00:00 2001 From: CMGS Date: Fri, 25 Sep 2026 15:14:24 +0800 Subject: [PATCH 1/3] fix: a claim and a fork child are indexed by claim id, so e2b lookups ask their node first Claim recorded only the name key, and Fork recorded nothing, while the e2b surface and the envd proxy resolve by claim id. So the first by-id call after every e2b create or fork swept the fleet's inventories and, before the node published, asked every node. Claim now also records the claim id, and Fork records each child by name and by claim id, as scaling-design.md's claim-time index already states. --- pkg/scale/sandboxstore_claim_test.go | 3 ++- pkg/scale/sandboxstore_impl.go | 1 + pkg/scale/sandboxstore_lifecycle.go | 2 ++ pkg/scale/sandboxstore_live_test.go | 38 ++++++++++++++++++++++++++++ 4 files changed, 43 insertions(+), 1 deletion(-) 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{ From 0ac84d9228df2b0ae6abb797a69a755c0c517405 Mon Sep 17 00:00:00 2001 From: CMGS Date: Fri, 25 Sep 2026 15:27:00 +0800 Subject: [PATCH 2/3] review: a skipped node's inventory error logs at warn, the level of a degraded read --- pkg/e2bcompat/lifecycle.go | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/pkg/e2bcompat/lifecycle.go b/pkg/e2bcompat/lifecycle.go index 3925060..84dc8ab 100644 --- a/pkg/e2bcompat/lifecycle.go +++ b/pkg/e2bcompat/lifecycle.go @@ -260,7 +260,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 @@ -356,7 +356,7 @@ func (s *Server) inventories(r *http.Request) ([]*scale.NodeInventory, error) { 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) From 1d6c607a48964f597b7f3f42022abc3b014c80df Mon Sep 17 00:00:00 2001 From: CMGS Date: Fri, 25 Sep 2026 15:28:15 +0800 Subject: [PATCH 3/3] review: godoc and inline comments fit the one-line budget --- cmd/sandbox-apiserver/main.go | 15 ++---------- examples/lifecycle/example.go | 21 +++-------------- pkg/e2bcompat/lifecycle.go | 35 ++++------------------------ pkg/e2bcompat/server.go | 43 ++++++++--------------------------- pkg/envdproxy/route.go | 5 +--- 5 files changed, 21 insertions(+), 98 deletions(-) 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 84dc8ab..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) { @@ -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,8 +334,6 @@ 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").Warnf(r.Context(), "e2b: node inventory unavailable node=%s err=%v", node, err) continue } @@ -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))