Skip to content

Commit f2fc482

Browse files
bchaliosclaude
authored andcommitted
feat(orchestrator): harvest a resume prefetch mapping after an in-place checkpoint
Co-authored-by: Claude Fable 5.1 <noreply@anthropic.com> GitOrigin-RevId: f24da54561b4235441a25b86fc981cc59a58900d
1 parent 203ddfe commit f2fc482

9 files changed

Lines changed: 439 additions & 64 deletions

File tree

Lines changed: 61 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,61 @@
1+
package sandbox
2+
3+
import (
4+
"sync"
5+
"testing"
6+
7+
"github.com/stretchr/testify/require"
8+
9+
"github.com/e2b-dev/infra/packages/shared/pkg/grpc/orchestrator"
10+
)
11+
12+
// A clone carries the network config as it was at clone time and its own lock;
13+
// updates to either side afterwards do not cross over.
14+
func TestConfigClone_IsolatesNetwork(t *testing.T) {
15+
t.Parallel()
16+
17+
src := NewConfig(Config{Network: &orchestrator.SandboxNetworkConfig{
18+
Egress: &orchestrator.SandboxNetworkEgressConfig{AllowedDomains: []string{"a.example"}},
19+
}})
20+
21+
clone := src.Clone()
22+
src.SetNetworkEgress(&orchestrator.SandboxNetworkEgressConfig{AllowedDomains: []string{"b.example"}})
23+
clone.SetNetworkEgress(&orchestrator.SandboxNetworkEgressConfig{AllowedDomains: []string{"c.example"}})
24+
25+
require.Equal(t, []string{"b.example"}, src.GetNetworkEgress().GetAllowedDomains())
26+
require.Equal(t, []string{"c.example"}, clone.GetNetworkEgress().GetAllowedDomains())
27+
require.NotSame(t, src.Network, clone.Network)
28+
}
29+
30+
// Cloning while another goroutine rewrites the egress must be race-free; the
31+
// race detector is the assertion here.
32+
func TestConfigClone_ConcurrentWithEgressUpdate(t *testing.T) {
33+
t.Parallel()
34+
35+
src := NewConfig(Config{})
36+
var wg sync.WaitGroup
37+
wg.Add(2)
38+
go func() {
39+
defer wg.Done()
40+
for i := range 500 {
41+
src.SetNetworkEgress(&orchestrator.SandboxNetworkEgressConfig{AllowedDomains: []string{string(rune('a' + i%26))}})
42+
}
43+
}()
44+
go func() {
45+
defer wg.Done()
46+
for range 500 {
47+
c := src.Clone()
48+
_ = c.GetNetworkEgress().GetAllowedDomains()
49+
}
50+
}()
51+
wg.Wait()
52+
}
53+
54+
// A nil network on the source still yields a usable, non-nil network on the
55+
// clone, matching NewConfig's normalisation.
56+
func TestConfigClone_NilNetworkNormalised(t *testing.T) {
57+
t.Parallel()
58+
59+
src := NewConfig(Config{})
60+
require.NotNil(t, src.Clone().Network)
61+
}

‎packages/orchestrator/pkg/sandbox/sandbox.go‎

Lines changed: 19 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -22,6 +22,7 @@ import (
2222
"go.opentelemetry.io/otel/metric"
2323
"go.opentelemetry.io/otel/trace"
2424
"go.uber.org/zap"
25+
"google.golang.org/protobuf/proto"
2526

2627
"github.com/e2b-dev/infra/packages/clickhouse/pkg/hoststats"
2728
"github.com/e2b-dev/infra/packages/orchestrator/pkg/cfg"
@@ -188,6 +189,24 @@ func (c *Config) GetNetworkIngress() *orchestrator.SandboxNetworkIngressConfig {
188189
return c.Network.GetIngress()
189190
}
190191

192+
// Clone returns a copy with its own lock and its own network config, snapshotted
193+
// under this config's lock. A throwaway resumed from the copy therefore starts
194+
// from a consistent egress and ingress even while a concurrent Update rewrites
195+
// the original's, and nothing done to the copy reaches the live sandbox.
196+
func (c *Config) Clone() *Config {
197+
c.mu.RLock()
198+
defer c.mu.RUnlock()
199+
200+
clone := *c
201+
clone.Network, _ = proto.Clone(c.Network).(*orchestrator.SandboxNetworkConfig)
202+
if clone.Network == nil {
203+
clone.Network = &orchestrator.SandboxNetworkConfig{}
204+
}
205+
clone.mu = &sync.RWMutex{}
206+
207+
return &clone
208+
}
209+
191210
type VolumeMountConfig struct {
192211
ID uuid.UUID
193212
Name string
Lines changed: 162 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,162 @@
1+
package server
2+
3+
import (
4+
"context"
5+
"errors"
6+
"sync/atomic"
7+
"testing"
8+
"testing/synctest"
9+
"time"
10+
11+
"github.com/launchdarkly/go-sdk-common/v3/ldvalue"
12+
"github.com/launchdarkly/go-server-sdk/v7/testhelpers/ldtestdata"
13+
"github.com/stretchr/testify/require"
14+
15+
"github.com/e2b-dev/infra/packages/orchestrator/pkg/sandbox/build"
16+
"github.com/e2b-dev/infra/packages/orchestrator/pkg/service"
17+
"github.com/e2b-dev/infra/packages/shared/pkg/featureflags"
18+
"github.com/e2b-dev/infra/packages/shared/pkg/grpc/orchestrator"
19+
"github.com/e2b-dev/infra/packages/shared/pkg/utils"
20+
)
21+
22+
// harvestFlagClient enables the harvest with a one-second budget, consume off.
23+
func harvestFlagClient(t *testing.T) *featureflags.Client {
24+
t.Helper()
25+
26+
td := ldtestdata.DataSource()
27+
td.Update(td.Flag(featureflags.PauseResumePrefetchHarvestFlag.Key()).VariationForAll(true))
28+
td.Update(td.Flag(featureflags.PauseResumePrefetchHarvestTimeoutMsFlag.Key()).ValueForAll(ldvalue.Int(1000)))
29+
ff, err := featureflags.NewClientWithDatasource(td)
30+
require.NoError(t, err)
31+
t.Cleanup(func() { _ = ff.Close(context.WithoutCancel(t.Context())) })
32+
33+
return ff
34+
}
35+
36+
// A filesystem-only checkpoint has no memfile to resume, so no harvest is
37+
// scheduled even with the flag on.
38+
func TestHarvestCheckpointPrefetchAsync_SkipsFilesystemOnly(t *testing.T) {
39+
t.Parallel()
40+
41+
s := &Server{info: &service.ServiceInfo{}, featureFlags: harvestFlagClient(t)}
42+
res := &snapshotResult{rootfsDiff: &build.NoDiff{}}
43+
s.harvestCheckpointPrefetchAsync(t.Context(), testHarvestSandbox(), res, &orchestrator.SandboxCheckpointRequest{BuildId: "build-1", FilesystemOnly: true})
44+
require.Zero(t, s.info.OutstandingWork())
45+
}
46+
47+
// A memory checkpoint taken in place schedules the harvest as tracked work,
48+
// which outlives the request and ends when the snapshot's seal settles.
49+
func TestHarvestCheckpointPrefetchAsync_SchedulesAndTracksWork(t *testing.T) {
50+
t.Parallel()
51+
52+
ff := harvestFlagClient(t)
53+
synctest.Test(t, func(t *testing.T) {
54+
s := &Server{info: &service.ServiceInfo{}, featureFlags: ff}
55+
seal := utils.NewSetOnce[build.Diff]()
56+
res := &snapshotResult{rootfsDiff: build.NewDeferredDiff("rootfs", 4096, seal)}
57+
ctx, cancel := context.WithCancel(t.Context())
58+
defer cancel()
59+
60+
s.harvestCheckpointPrefetchAsync(ctx, testHarvestSandbox(), res, &orchestrator.SandboxCheckpointRequest{BuildId: "build-1"})
61+
require.Equal(t, int64(1), s.info.OutstandingWork())
62+
cancel()
63+
synctest.Wait()
64+
require.Equal(t, int64(1), s.info.OutstandingWork(), "request cancellation must not release harvest work")
65+
require.NoError(t, seal.SetError(build.ErrDeferredSealFailed))
66+
synctest.Wait()
67+
require.Zero(t, s.info.OutstandingWork())
68+
})
69+
}
70+
71+
// An in-place checkpoint through the CoW window hands the harvest a memfile
72+
// that is still being swept: the harvest must not resume until that seal
73+
// settles, and a failed seal skips it without touching the rootfs.
74+
func TestHarvestResumePrefetchAsync_WaitsForDeferredMemorySeal(t *testing.T) {
75+
t.Parallel()
76+
77+
for _, tc := range []struct {
78+
name string
79+
memorySealErr error
80+
}{
81+
{name: "seal fails", memorySealErr: errors.New("window cancelled")},
82+
{name: "seal settles"},
83+
} {
84+
t.Run(tc.name, func(t *testing.T) {
85+
t.Parallel()
86+
87+
ff := harvestFlagClient(t)
88+
synctest.Test(t, func(t *testing.T) {
89+
s := &Server{info: &service.ServiceInfo{}, featureFlags: ff}
90+
memorySealed := make(chan error, 1)
91+
rootfsSeal := utils.NewSetOnce[build.Diff]()
92+
var rootfsWaited atomic.Bool
93+
res := &snapshotResult{
94+
rootfsDiff: build.NewDeferredDiff("rootfs", 4096, rootfsSeal),
95+
memoryExportDeferred: true,
96+
waitMemorySealed: func(ctx context.Context) error {
97+
select {
98+
case err := <-memorySealed:
99+
return err
100+
case <-ctx.Done():
101+
return ctx.Err()
102+
}
103+
},
104+
}
105+
// The rootfs promise is only consulted once memory has sealed; a
106+
// failed rootfs seal then ends the harvest without a resume.
107+
go func() {
108+
<-time.After(10 * time.Millisecond)
109+
rootfsWaited.Store(true)
110+
_ = rootfsSeal.SetError(build.ErrDeferredSealFailed)
111+
}()
112+
113+
s.harvestResumePrefetchAsync(t.Context(), testHarvestSandbox(), res, "build-1", nil, harvestSourceCheckpoint)
114+
require.Equal(t, int64(1), s.info.OutstandingWork())
115+
<-time.After(500 * time.Millisecond)
116+
synctest.Wait()
117+
require.Equal(t, int64(1), s.info.OutstandingWork(), "the harvest must hold until the memory seal settles")
118+
119+
memorySealed <- tc.memorySealErr
120+
synctest.Wait()
121+
require.Zero(t, s.info.OutstandingWork())
122+
require.True(t, rootfsWaited.Load())
123+
})
124+
})
125+
}
126+
}
127+
128+
// waitSnapshotSealed orders the waits memfile first, names which seal failed,
129+
// and is a no-op when nothing was deferred.
130+
func TestWaitSnapshotSealed(t *testing.T) {
131+
t.Parallel()
132+
133+
require.NoError(t, waitSnapshotSealed(t.Context(), &snapshotResult{rootfsDiff: &build.NoDiff{}}))
134+
135+
memErr := errors.New("sweep aborted")
136+
err := waitSnapshotSealed(t.Context(), &snapshotResult{
137+
rootfsDiff: &build.NoDiff{},
138+
memoryExportDeferred: true,
139+
waitMemorySealed: func(context.Context) error { return memErr },
140+
})
141+
require.ErrorIs(t, err, memErr)
142+
require.ErrorContains(t, err, "memory seal")
143+
144+
rootfsSeal := utils.NewSetOnce[build.Diff]()
145+
require.NoError(t, rootfsSeal.SetError(build.ErrDeferredSealFailed))
146+
memoryWaited := false
147+
err = waitSnapshotSealed(t.Context(), &snapshotResult{
148+
rootfsDiff: build.NewDeferredDiff("rootfs", 4096, rootfsSeal),
149+
memoryExportDeferred: true,
150+
waitMemorySealed: func(context.Context) error {
151+
memoryWaited = true
152+
153+
return nil
154+
},
155+
})
156+
require.ErrorIs(t, err, build.ErrDeferredSealFailed)
157+
require.ErrorContains(t, err, "rootfs seal")
158+
require.True(t, memoryWaited, "memory seal is awaited before the rootfs seal")
159+
160+
// A deferred flag with no waiter (older snapshot shape) is treated as sealed.
161+
require.NoError(t, waitSnapshotSealed(t.Context(), &snapshotResult{rootfsDiff: &build.NoDiff{}, memoryExportDeferred: true}))
162+
}

0 commit comments

Comments
 (0)