Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
23 commits
Select commit Hold shift + click to select a range
ff7fc5f
review: inline the single log.WithFunc call in meshinventory.tick
CMGS Sep 28, 2026
addc8a0
cut(e2bcompat): drop the unread force field of the build start body
CMGS Sep 28, 2026
96a957e
fix(e2bcompat): a full e2b/ key in a create or a build's fromImage is…
CMGS Sep 28, 2026
db24afc
fix(e2bbuild): an ENV value is the expansion's stdout alone
CMGS Sep 28, 2026
db084fe
fix(e2bcompat): a publish deletes only the holders older than its own…
CMGS Sep 28, 2026
2376abf
fix(sandboxd): promote, checkpoint, fork, hibernate and wake ride a c…
CMGS Sep 28, 2026
5547e03
fix(sandboxd): the port relay speaks the node's scheme and can presen…
CMGS Sep 28, 2026
ee77d3c
fix(e2b): build commands ride the claim's own relay, so a long build …
CMGS Sep 28, 2026
40002c8
fix(envdproxy): a signed file URL relays with the node token, so a st…
CMGS Sep 28, 2026
8356c6d
review: named func types at the build seams, exported var first, mode…
CMGS Sep 28, 2026
4b0fd35
simplify(scale): one fleet read per claim, direct PoolKey conversions…
CMGS Sep 28, 2026
9902e12
simplify(envdproxy): the probe returns the owner, the resolver builds…
CMGS Sep 28, 2026
0051eb0
simplify(e2b): one relay env map, the build lease from the executor, …
CMGS Sep 28, 2026
0c60192
simplify(cmd): one sandboxd token flag pair for the three binaries
CMGS Sep 28, 2026
b76e04f
docs(e2b): built templates expect a checkpoint_store per node
CMGS Sep 28, 2026
590c5fc
fix(envdproxy): a signed file URL goes without the token only on envd…
CMGS Sep 28, 2026
288fa28
fix(scale): a stamped name is checked against sandboxd's grammar befo…
CMGS Sep 28, 2026
1c2726a
fix(e2b): a publish deletes older holders by the digest it observed a…
CMGS Sep 28, 2026
6cf75ea
fix(sandboxd): the relay target keeps an IPv6 origin's brackets once …
CMGS Sep 28, 2026
fc5d212
tests(sandboxd): the client sends the observed digest and surfaces a …
CMGS Sep 28, 2026
a164694
fix(e2b): a tag write names the generation it was meant for
CMGS Sep 28, 2026
2e9cc0a
tests(e2b): the publish regression tests 1c2726aa named
CMGS Sep 28, 2026
94d7696
fix(e2bbuild): a failed publish keeps its node-free message
CMGS Sep 28, 2026
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
13 changes: 6 additions & 7 deletions cmd/sandbox-apiserver/main.go
Original file line number Diff line number Diff line change
Expand Up @@ -44,6 +44,8 @@ const (
watchDrainGrace = 2 * time.Second
)

type managerBuilder func() (manager.Runnable, error)

// options has no etcd option because this server stores nothing.
type options struct {
SecureServing *genericoptions.SecureServingOptionsWithLoopback
Expand Down Expand Up @@ -85,10 +87,7 @@ func (o *options) addFlags(fs *pflag.FlagSet) {
o.Authentication.AddFlags(fs)
o.Authorization.AddFlags(fs)
o.Features.AddFlags(fs)
fs.StringVar(&o.SandboxdToken, "sandboxd-token", o.SandboxdToken,
"Uniform fleet-wide sandboxd api_token presented on node-local claim/release. Prefer --sandboxd-token-file for a Secret mount.")
fs.StringVar(&o.SandboxdTokenFile, "sandboxd-token-file", o.SandboxdTokenFile,
"Path to a file (Secret mount) holding the sandboxd api_token; overrides --sandboxd-token when set.")
sandboxd.AddTokenFlags(fs, &o.SandboxdToken, &o.SandboxdTokenFile)
fs.BoolVar(&o.WarmPoolDriver, "enable-warm-pool-driver", o.WarmPoolDriver,
"Run the in-process SandboxWarmPool → sandboxd pool reconcile loop (control-plane warm-capacity surface; pool-level, never per-sandbox).")
fs.DurationVar(&o.WarmPoolInterval, "warm-pool-sync-interval", o.WarmPoolInterval,
Expand Down Expand Up @@ -140,15 +139,15 @@ func run() error {
if err != nil {
return fmt.Errorf("load kube config: %w", err)
}
reader, err := kubeinventory.NewCache(ctx, restCfg)
informers, err := kubeinventory.NewCache(ctx, restCfg)
if err != nil {
return err
}
token, err := sandboxd.TokenFrom(o.SandboxdToken, o.SandboxdTokenFile)
if err != nil {
return err
}
invSource, err := kubeinventory.New(ctx, reader, o.Inventory)
invSource, err := kubeinventory.New(ctx, informers, o.Inventory)
if err != nil {
return err
}
Expand Down Expand Up @@ -221,7 +220,7 @@ func startWarmPoolDriver(ctx context.Context, restCfg *restclient.Config, token
return nil
}

func runRestarting(ctx context.Context, r manager.Runnable, build func() (manager.Runnable, error), delay time.Duration) {
func runRestarting(ctx context.Context, r manager.Runnable, build managerBuilder, delay time.Duration) {
logger := log.WithFunc("main.runRestarting")
for {
if r != nil {
Expand Down
9 changes: 3 additions & 6 deletions cmd/sandbox-e2b/main.go
Original file line number Diff line number Diff line change
Expand Up @@ -35,10 +35,7 @@ type options struct {
func (o *options) addFlags(fs *pflag.FlagSet) {
fs.StringSliceVar(&o.Seeds, "sandboxd-seeds", o.Seeds,
"Comma-separated sandboxd addresses dialed at start, each naming one node; each node reports its own advertise_addr as its key, and the rest of the mesh is found through their gossip.")
fs.StringVar(&o.SandboxdToken, "sandboxd-token", o.SandboxdToken,
"Fleet root sandboxd api_token (GET /v1/info needs root). Prefer --sandboxd-token-file for a Secret mount.")
fs.StringVar(&o.SandboxdTokenFile, "sandboxd-token-file", o.SandboxdTokenFile,
"Path to a file (Secret mount) holding the sandboxd api_token; overrides --sandboxd-token when set.")
sandboxd.AddTokenFlags(fs, &o.SandboxdToken, &o.SandboxdTokenFile)
fs.DurationVar(&o.PollInterval, "inventory-poll", o.PollInterval,
"How often every node's info and sandbox list are read; List and Watch lag a change by up to one tick.")
o.E2B.AddFlags(fs)
Expand Down Expand Up @@ -94,11 +91,11 @@ func run() error {
<-ctx.Done()
return nil
}
resolver, err := envdproxy.NewResolver(scale.NewScatterGatherStore(src), store, src, "", opts.EnvdSecret)
resolver, err := envdproxy.NewResolver(store, src, "", opts.EnvdSecret)
if err != nil {
return err
}
proxy, err := envdproxy.NewServer(resolver, envdproxy.Options{Domain: o.E2B.Domain, GuestHTTP2: o.Proxy.GuestHTTP2})
proxy, err := envdproxy.NewServer(resolver, envdproxy.Options{Domain: o.E2B.Domain, GuestHTTP2: o.Proxy.GuestHTTP2, NodeToken: token})
if err != nil {
return err
}
Expand Down
14 changes: 6 additions & 8 deletions cmd/sandbox-envd-proxy/main.go
Original file line number Diff line number Diff line change
Expand Up @@ -45,10 +45,7 @@ func (o *options) addFlags(fs *pflag.FlagSet) {
o.Inventory.AddFlags(fs)
fs.StringVar(&o.Domain, "domain", o.Domain,
"Base domain sandbox hosts are derived from, as {port}-{sandboxID}.{domain}. Must match the apiserver's --e2b-domain.")
fs.StringVar(&o.SandboxdToken, "sandboxd-token", o.SandboxdToken,
"Fleet root sandboxd api_token, which reads a sandbox's claim token to verify its envd access token. Prefer --sandboxd-token-file for a Secret mount.")
fs.StringVar(&o.SandboxdTokenFile, "sandboxd-token-file", o.SandboxdTokenFile,
"Path to a file (Secret mount) holding the sandboxd api_token; overrides --sandboxd-token when set.")
sandboxd.AddTokenFlags(fs, &o.SandboxdToken, &o.SandboxdTokenFile)
e2bcompat.AddEnvdSecretFlag(fs, &o.EnvdSecretFile)
fs.StringVar(&o.Namespace, "namespace", o.Namespace,
"Namespace inventory lookups are filtered to; empty matches every namespace. Not an access boundary: a caller holding a sandbox's token reaches it in any namespace.")
Expand All @@ -68,7 +65,7 @@ func main() {
o.addFlags(fs)
_ = fs.Parse(os.Args[1:])
if err := run(ctx, o); err != nil {
log.WithFunc("main").Fatalf(ctx, err, "sandbox-envd-proxy exited")
log.WithFunc("main.main").Fatalf(ctx, err, "sandbox-envd-proxy exited")
}
}

Expand All @@ -78,11 +75,11 @@ func run(ctx context.Context, o *options) error {
if err != nil {
return fmt.Errorf("load kube config: %w", err)
}
reader, err := kubeinventory.NewCache(ctx, restCfg)
informers, err := kubeinventory.NewCache(ctx, restCfg)
if err != nil {
return err
}
inv, err := kubeinventory.New(ctx, reader, o.Inventory)
inv, err := kubeinventory.New(ctx, informers, o.Inventory)
if err != nil {
return err
}
Expand All @@ -95,13 +92,14 @@ func run(ctx context.Context, o *options) error {
return err
}
routed := scale.NewScatterGatherStore(inv, scale.WithClaimRouting(token, scale.NewSandboxdClientFactory()))
resolver, err := envdproxy.NewResolver(scale.NewScatterGatherStore(inv), routed, inv, o.Namespace, secret)
resolver, err := envdproxy.NewResolver(routed, inv, o.Namespace, secret)
if err != nil {
return err
}
srv, err := envdproxy.NewServer(resolver, envdproxy.Options{
Domain: o.Domain,
GuestHTTP2: o.Proxy.GuestHTTP2,
NodeToken: token,
})
if err != nil {
return err
Expand Down
7 changes: 7 additions & 0 deletions docs/e2b-compat.md
Original file line number Diff line number Diff line change
Expand Up @@ -325,6 +325,13 @@ needs its own stickiness (or run one replica). A finished build is kept for an h
loses a build in flight; the SDK's next poll throws, and a rebuild converges
because a promote replaces.

Built templates expect a `checkpoint_store` per node. A node reads the store's
template records once at start and advertises them as its own, so with one S3
store shared by every node each node advertises every node's templates as of its
own start, and a publish that deletes an older holder removes a record its
siblings still advertise; a create of one fails over to the next advertiser, and
the template list lags until those nodes restart.

## Limits worth knowing

- **Reaching `envd` (the in-sandbox data plane).** The SDK derives the sandbox
Expand Down
13 changes: 8 additions & 5 deletions pkg/e2bbuild/copy.go
Original file line number Diff line number Diff line change
Expand Up @@ -2,6 +2,7 @@ package e2bbuild

import (
"path"
"slices"
"strings"
)

Expand Down Expand Up @@ -47,7 +48,7 @@ func copyScript(state Command, s Step) string {
scratch := "/tmp/" + s.FilesHash
unpack := scratch + "/unpack"
vars := [][2]string{
{"archive", scratch + ".tar"},
{"archive", archivePath(s.FilesHash)},
{"scratch", scratch},
{"unpack", unpack},
{"src", s.Args[0]},
Expand All @@ -65,13 +66,15 @@ func copyScript(state Command, s Step) string {
return b.String() + copyMove
}

func archivePath(hash string) string {
return "/tmp/" + hash + ".tar"
}

// globBase is src up to its first segment with a glob; the SDK expands globs into the archive itself.
func globBase(src string) string {
segs := strings.Split(strings.TrimSuffix(src, "/"), "/")
for i, seg := range segs {
if strings.ContainsAny(seg, "*?[{") {
return strings.Join(segs[:i], "/")
}
if i := slices.IndexFunc(segs, func(seg string) bool { return strings.ContainsAny(seg, "*?[{") }); i >= 0 {
segs = segs[:i]
}
return strings.Join(segs, "/")
}
77 changes: 42 additions & 35 deletions pkg/e2bbuild/executor.go
Original file line number Diff line number Diff line change
Expand Up @@ -14,6 +14,7 @@ import (
"maps"
"path"
"regexp"
"slices"
"strconv"
"strings"
"sync"
Expand Down Expand Up @@ -55,11 +56,21 @@ var (
// ErrUnknownBuild is a build this process never registered, or dropped an hour after it finished or was left unstarted.
ErrUnknownBuild = errors.New("e2bbuild: build not found")

envEscaper = strings.NewReplacer(`\`, `\\`, `"`, `\"`, "`", "\\`", "$(", `\$(`)
// FilesHash is the shape of the digest the SDK names a COPY upload by.
FilesHash = regexp.MustCompile(`^[0-9a-f]{64}$`)

envEscaper = strings.NewReplacer(`\`, `\\`, `"`, `\"`, "`", "\\`", "$(", `\$(`)
)

// LineFunc receives one line of a command's output.
type LineFunc func(line string)

// ArchiveFunc opens the upload a COPY step's files hash names.
type ArchiveFunc func(ctx context.Context, hash string) (io.ReadCloser, error)

// PublishFunc runs after the promote with the template's node, key and content digest.
type PublishFunc func(ctx context.Context, node string, key scale.PoolKey, digest string) error

// Request is what a build asked for before it starts.
type Request struct {
Size string
Expand Down Expand Up @@ -99,19 +110,17 @@ func (s Step) problem(canCopy bool) string {
return ""
}

// Spec is how a registered build runs; RelayEnvs seed the steps when the claim relays through the node's proxy, StartCmd runs in the background and ReadyCmd until it exits 0, both before the promote; Publish runs after the promote and before the build reads ready.
// Spec is how a registered build runs; StartCmd runs in the background and ReadyCmd until it exits 0, both before the promote; Publish runs after the promote and before the build reads ready.
type Spec struct {
Namespace string
ClaimName string
Pool scale.PoolKey
Template string
TTLSeconds int
Steps []Step
StartCmd string
ReadyCmd string
RelayEnvs map[string]string
Archive func(ctx context.Context, hash string) (io.ReadCloser, error)
Publish func(ctx context.Context, node string, key scale.PoolKey, digest string) error
Namespace string
ClaimName string
Pool scale.PoolKey
Template string
Steps []Step
StartCmd string
ReadyCmd string
Archive ArchiveFunc
Publish PublishFunc
}

// Command is one shell line run as User, in Workdir, with Envs; an empty Workdir is the user's home.
Expand All @@ -124,8 +133,8 @@ type Command struct {

// Guest runs a build's commands inside its claimed sandbox.
type Guest interface {
// Run runs cmd to its end, passes each output line to out, and returns its exit code.
Run(ctx context.Context, a scale.Assignment, cmd Command, out func(line string)) (int, error)
// Run runs cmd to its end, passes each line of its output to stdout or stderr, and returns its exit code.
Run(ctx context.Context, a scale.Assignment, cmd Command, stdout, stderr LineFunc) (int, error)
// Start starts cmd and returns while it runs.
Start(ctx context.Context, a scale.Assignment, cmd Command) error
// Init makes the user, workdir and envs of defaults what every later process in the sandbox gets.
Expand Down Expand Up @@ -227,7 +236,7 @@ func (e *Executor) Status(id string) (Info, bool) {
return Info{}, false
}
info := r.info
info.Logs = append([]LogEntry(nil), r.info.Logs...)
info.Logs = slices.Clip(r.info.Logs)
return info, true
}

Expand All @@ -236,7 +245,7 @@ func (e *Executor) run(ctx context.Context, id string, spec Spec) {
ctx, cancel := context.WithTimeout(ctx, e.timeout)
defer cancel()
e.logf(id, PhaseClaim, "claiming a sandbox from %s (%s, %s)", spec.Pool.Template, spec.Pool.Net, spec.Pool.Size)
a, err := e.store.Claim(ctx, spec.Namespace, spec.ClaimName, spec.Pool, scale.ClaimOptions{TTLSeconds: spec.TTLSeconds})
a, err := e.store.Claim(ctx, spec.Namespace, spec.ClaimName, spec.Pool, scale.ClaimOptions{TTLSeconds: int(e.timeout / time.Second)})
if err != nil {
msg := fmt.Sprintf("could not claim a sandbox of %s (%s, %s)", spec.Pool.Template, spec.Pool.Net, spec.Pool.Size)
if scale.IsNoWarmCapacity(err) {
Expand All @@ -261,7 +270,7 @@ func (e *Executor) run(ctx context.Context, id string, spec Spec) {
func (e *Executor) prepare(ctx context.Context, id string, a scale.Assignment, spec Spec) (string, string, error) {
state := Command{User: defaultUser, Envs: map[string]string{}}
if a.NetRoute == sandboxd.NetRouteRelay {
maps.Copy(state.Envs, spec.RelayEnvs)
maps.Copy(state.Envs, sandboxd.RelayEnv)
}
for i, step := range spec.Steps {
phase := strconv.Itoa(i + 1)
Expand Down Expand Up @@ -298,12 +307,12 @@ func (e *Executor) step(ctx context.Context, id, phase string, a scale.Assignmen
if len(s.Args) > 1 {
cmd.User = s.Args[1]
}
msg, err := e.sh(ctx, a, cmd, logged)
msg, err := e.sh(ctx, a, cmd, logged, logged)
return state, msg, err
case stepEnv:
envs := maps.Clone(state.Envs)
for i := 0; i+1 < len(s.Args); i += 2 {
v, msg, err := e.expand(ctx, a, root, s.Args[i+1])
v, msg, err := e.expand(ctx, a, root, s.Args[i+1], logged)
if err != nil {
return state, msg, err
}
Expand All @@ -316,22 +325,22 @@ func (e *Executor) step(ctx context.Context, id, phase string, a scale.Assignmen
dir = path.Join(cmp.Or(state.Workdir, "/"), dir)
}
script := fmt.Sprintf(`t=%[1]s; [ -d "$t" ] && exit 0; n=$t; while [ ! -d "$(dirname "$n")" ]; do n=$(dirname "$n"); done; mkdir -p "$t" && chown -R %[2]s: "$n"`, shellQuote(dir), shellQuote(state.User))
if msg, err := e.sh(ctx, a, withLine(root, script), logged); err != nil {
if msg, err := e.sh(ctx, a, withLine(root, script), logged, logged); err != nil {
return state, msg, err
}
state.Workdir = dir
case stepUser:
name := shellQuote(s.Args[0])
script := fmt.Sprintf("id -u %[1]s >/dev/null 2>&1 || useradd --create-home --shell /bin/bash %[1]s", name)
if msg, err := e.sh(ctx, a, withLine(root, script), logged); err != nil {
if msg, err := e.sh(ctx, a, withLine(root, script), logged, logged); err != nil {
return state, msg, err
}
state.User = s.Args[0]
case stepCopy:
if msg, err := e.copyIn(ctx, a, spec, s.FilesHash); err != nil {
return state, msg, err
}
if _, err := e.sh(ctx, a, withLine(root, copyScript(state, s)), logged); err != nil {
if _, err := e.sh(ctx, a, withLine(root, copyScript(state, s)), logged, logged); err != nil {
return state, fmt.Sprintf("could not copy %s to %s", s.Args[0], s.Args[1]), err
}
}
Expand All @@ -345,15 +354,15 @@ func (e *Executor) copyIn(ctx context.Context, a scale.Assignment, spec Spec, ha
return "could not read the uploaded files", err
}
defer func() { _ = archive.Close() }()
if err := e.guest.Write(ctx, a, "/tmp/"+hash+".tar", archive); err != nil {
if err := e.guest.Write(ctx, a, archivePath(hash), archive); err != nil {
return "could not copy the uploaded files into the build sandbox", err
}
return "", nil
}

// sh returns the failure the caller sees next to the error; a non-zero exit is a failure.
func (e *Executor) sh(ctx context.Context, a scale.Assignment, cmd Command, out func(string)) (string, error) {
code, err := e.guest.Run(ctx, a, cmd, out)
func (e *Executor) sh(ctx context.Context, a scale.Assignment, cmd Command, stdout, stderr LineFunc) (string, error) {
code, err := e.guest.Run(ctx, a, cmd, stdout, stderr)
if err != nil {
return "could not run a command in the build sandbox", err
}
Expand All @@ -364,18 +373,18 @@ func (e *Executor) sh(ctx context.Context, a scale.Assignment, cmd Command, out
return "", nil
}

// expand evaluates an ENV value in the guest's shell, so $VAR references resolve and command substitution does not run.
func (e *Executor) expand(ctx context.Context, a scale.Assignment, root Command, value string) (string, string, error) {
// expand evaluates an ENV value in the guest's shell, so $VAR references resolve and command substitution does not run; the shell's own stderr goes to the log, never into the value.
func (e *Executor) expand(ctx context.Context, a scale.Assignment, root Command, value string, stderr LineFunc) (string, string, error) {
var lines []string
if _, err := e.sh(ctx, a, withLine(root, `printf "%s" "`+envEscaper.Replace(value)+`"`), func(line string) { lines = append(lines, line) }); err != nil {
if _, err := e.sh(ctx, a, withLine(root, `printf "%s" "`+envEscaper.Replace(value)+`"`), func(line string) { lines = append(lines, line) }, stderr); err != nil {
return "", fmt.Sprintf("could not evaluate the value %q", value), err
}
return strings.Join(lines, "\n"), "", nil
}

func (e *Executor) awaitReady(ctx context.Context, a scale.Assignment, cmd Command) error {
for {
code, err := e.guest.Run(ctx, a, cmd, func(string) {})
code, err := e.guest.Run(ctx, a, cmd, func(string) {}, func(string) {})
if err == nil && code == 0 {
return nil
}
Expand Down Expand Up @@ -440,13 +449,11 @@ func (e *Executor) log(id, step, level, message string) {
}

func (e *Executor) sweep(now time.Time) {
for id, r := range e.builds {
maps.DeleteFunc(e.builds, func(_ string, r *record) bool {
done := !r.finished.IsZero() && now.Sub(r.finished) > recordTTL
abandoned := r.info.Status == StatusWaiting && now.Sub(r.registered) > recordTTL
if done || abandoned {
delete(e.builds, id)
}
}
return done || abandoned
})
}

// Invalid names the first step a build cannot run, empty when every step can; a COPY needs canCopy.
Expand Down
Loading
Loading