Skip to content

Commit 0a1f741

Browse files
kalyazinclaude
andcommitted
refactor(envd): share one workload freezer across freeze paths
The live-upgrade handover added a second workload freeze/unfreeze path (process.Service.FreezeWorkload/UnfreezeWorkload) that swept the user/pty cgroups directly, independent of the HTTP API's existing /freeze, /unfreeze and /init-deferred-thaw. Those API paths serialize their sweep through API.freezeLock precisely so a pause freeze, a rollback unfreeze and the resume thaw can't interleave and strand the workload frozen; the upgrade path took no lock at all, so an upgrade freeze racing the resume thaw could leave the sandbox frozen. Extract the sweep and its lock into a single cgroups.WorkloadFreezer that owns the manager plus a serializing semaphore and exposes Freeze/Unfreeze (best-effort, joined errors; Unfreeze detaches from ctx cancellation so a thaw always lands). main constructs one instance and passes it to both the HTTP API and the process service, so every freeze/unfreeze caller now shares the same lock. The canonical user/pty cgroup list moves to cgroups.WorkloadProcessTypes. The handover now holds that shared lock across the WHOLE window (freeze -> serialize -> execve) via a new FreezeHold/FreezeWorkloadHold, not just the freeze sweep, so a concurrent /init or /unfreeze thaw blocks on the lock until the handover finishes and cannot thaw the workload mid-handover. No behavior change to the existing endpoints: PostFreeze still returns 503 on a cancelled request and 500 on a sweep failure, PostUnfreeze/​/init still attempt every cgroup best-effort. Per-cgroup error logs collapse into one joined log line per call. Signed-off-by: Nikita Kalyazin <nikita.kalyazin@e2b.dev> Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
1 parent 1c22abe commit 0a1f741

13 files changed

Lines changed: 242 additions & 106 deletions

File tree

packages/envd/internal/api/compose_test.go

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -33,7 +33,7 @@ func newComposeTestAPI(t *testing.T) (*API, *user.User) {
3333
User: currentUser.Username,
3434
}
3535

36-
return New(&logger, defaults, nil, false, cgroups.NewNoopManager()), currentUser
36+
return New(&logger, defaults, nil, false, cgroups.NewWorkloadFreezer(cgroups.NewNoopManager())), currentUser
3737
}
3838

3939
func writeSourceFile(t *testing.T, dir string, name string, data []byte) string {

packages/envd/internal/api/download_test.go

Lines changed: 13 additions & 13 deletions
Original file line numberDiff line numberDiff line change
@@ -96,7 +96,7 @@ func TestGetFilesContentDisposition(t *testing.T) {
9696
EnvVars: utils.NewEnvVars(),
9797
User: currentUser.Username,
9898
}
99-
api := New(&logger, defaults, nil, false, cgroups.NewNoopManager())
99+
api := New(&logger, defaults, nil, false, cgroups.NewWorkloadFreezer(cgroups.NewNoopManager()))
100100

101101
// Create request and response recorder
102102
req := httptest.NewRequestWithContext(t.Context(), http.MethodGet, "/files?path="+url.QueryEscape(tempFile), nil)
@@ -145,7 +145,7 @@ func TestGetFilesContentDispositionWithNestedPath(t *testing.T) {
145145
EnvVars: utils.NewEnvVars(),
146146
User: currentUser.Username,
147147
}
148-
api := New(&logger, defaults, nil, false, cgroups.NewNoopManager())
148+
api := New(&logger, defaults, nil, false, cgroups.NewWorkloadFreezer(cgroups.NewNoopManager()))
149149

150150
// Create request and response recorder
151151
req := httptest.NewRequestWithContext(t.Context(), http.MethodGet, "/files?path="+url.QueryEscape(tempFile), nil)
@@ -188,7 +188,7 @@ func TestGetFiles_GzipEncoding_ExplicitIdentityOffWithRange(t *testing.T) {
188188
EnvVars: utils.NewEnvVars(),
189189
User: currentUser.Username,
190190
}
191-
api := New(&logger, defaults, nil, false, cgroups.NewNoopManager())
191+
api := New(&logger, defaults, nil, false, cgroups.NewWorkloadFreezer(cgroups.NewNoopManager()))
192192

193193
// Create request and response recorder
194194
req := httptest.NewRequestWithContext(t.Context(), http.MethodGet, "/files?path="+url.QueryEscape(tempFile), nil)
@@ -229,7 +229,7 @@ func TestGetFiles_GzipDownload(t *testing.T) {
229229
EnvVars: utils.NewEnvVars(),
230230
User: currentUser.Username,
231231
}
232-
api := New(&logger, defaults, nil, false, cgroups.NewNoopManager())
232+
api := New(&logger, defaults, nil, false, cgroups.NewWorkloadFreezer(cgroups.NewNoopManager()))
233233

234234
req := httptest.NewRequestWithContext(t.Context(), http.MethodGet, "/files?path="+url.QueryEscape(tempFile), nil)
235235
req.Header.Set("Accept-Encoding", "gzip")
@@ -294,7 +294,7 @@ func TestPostFiles_GzipUpload(t *testing.T) {
294294
EnvVars: utils.NewEnvVars(),
295295
User: currentUser.Username,
296296
}
297-
api := New(&logger, defaults, nil, false, cgroups.NewNoopManager())
297+
api := New(&logger, defaults, nil, false, cgroups.NewWorkloadFreezer(cgroups.NewNoopManager()))
298298

299299
req := httptest.NewRequestWithContext(t.Context(), http.MethodPost, "/files?path="+url.QueryEscape(destPath), &gzBuf)
300300
req.Header.Set("Content-Type", mpWriter.FormDataContentType())
@@ -334,7 +334,7 @@ func TestPostFiles_RawBodyUpload(t *testing.T) {
334334
EnvVars: utils.NewEnvVars(),
335335
User: currentUser.Username,
336336
}
337-
api := New(&logger, defaults, nil, false, cgroups.NewNoopManager())
337+
api := New(&logger, defaults, nil, false, cgroups.NewWorkloadFreezer(cgroups.NewNoopManager()))
338338

339339
req := httptest.NewRequestWithContext(t.Context(), http.MethodPost, "/files?path="+url.QueryEscape(destPath), bytes.NewReader(originalContent))
340340
req.Header.Set("Content-Type", "application/octet-stream")
@@ -372,7 +372,7 @@ func TestPostFiles_RawBodyUploadCreatesDirectories(t *testing.T) {
372372
EnvVars: utils.NewEnvVars(),
373373
User: currentUser.Username,
374374
}
375-
api := New(&logger, defaults, nil, false, cgroups.NewNoopManager())
375+
api := New(&logger, defaults, nil, false, cgroups.NewWorkloadFreezer(cgroups.NewNoopManager()))
376376

377377
req := httptest.NewRequestWithContext(t.Context(), http.MethodPost, "/files?path="+url.QueryEscape(destPath), bytes.NewReader(originalContent))
378378
req.Header.Set("Content-Type", "application/octet-stream")
@@ -405,7 +405,7 @@ func TestPostFiles_RawBodyUploadRequiresPath(t *testing.T) {
405405
EnvVars: utils.NewEnvVars(),
406406
User: currentUser.Username,
407407
}
408-
api := New(&logger, defaults, nil, false, cgroups.NewNoopManager())
408+
api := New(&logger, defaults, nil, false, cgroups.NewWorkloadFreezer(cgroups.NewNoopManager()))
409409

410410
req := httptest.NewRequestWithContext(t.Context(), http.MethodPost, "/files", bytes.NewReader([]byte("some content")))
411411
req.Header.Set("Content-Type", "application/octet-stream")
@@ -440,7 +440,7 @@ func TestPostFiles_RawBodyUploadOverwritesExisting(t *testing.T) {
440440
EnvVars: utils.NewEnvVars(),
441441
User: currentUser.Username,
442442
}
443-
api := New(&logger, defaults, nil, false, cgroups.NewNoopManager())
443+
api := New(&logger, defaults, nil, false, cgroups.NewWorkloadFreezer(cgroups.NewNoopManager()))
444444

445445
req := httptest.NewRequestWithContext(t.Context(), http.MethodPost, "/files?path="+url.QueryEscape(destPath), bytes.NewReader(newContent))
446446
req.Header.Set("Content-Type", "application/octet-stream")
@@ -486,7 +486,7 @@ func TestPostFiles_RawBodyGzipUpload(t *testing.T) {
486486
EnvVars: utils.NewEnvVars(),
487487
User: currentUser.Username,
488488
}
489-
api := New(&logger, defaults, nil, false, cgroups.NewNoopManager())
489+
api := New(&logger, defaults, nil, false, cgroups.NewWorkloadFreezer(cgroups.NewNoopManager()))
490490

491491
req := httptest.NewRequestWithContext(t.Context(), http.MethodPost, "/files?path="+url.QueryEscape(destPath), &gzBuf)
492492
req.Header.Set("Content-Type", "application/octet-stream")
@@ -520,7 +520,7 @@ func TestPostFiles_UnsupportedContentType(t *testing.T) {
520520
EnvVars: utils.NewEnvVars(),
521521
User: currentUser.Username,
522522
}
523-
api := New(&logger, defaults, nil, false, cgroups.NewNoopManager())
523+
api := New(&logger, defaults, nil, false, cgroups.NewWorkloadFreezer(cgroups.NewNoopManager()))
524524

525525
tempDir := t.TempDir()
526526
destPath := filepath.Join(tempDir, "test.txt")
@@ -566,7 +566,7 @@ func TestPostFiles_MultipartStillWorksWithoutContentType(t *testing.T) {
566566
EnvVars: utils.NewEnvVars(),
567567
User: currentUser.Username,
568568
}
569-
api := New(&logger, defaults, nil, false, cgroups.NewNoopManager())
569+
api := New(&logger, defaults, nil, false, cgroups.NewWorkloadFreezer(cgroups.NewNoopManager()))
570570

571571
req := httptest.NewRequestWithContext(t.Context(), http.MethodPost, "/files?path="+url.QueryEscape(destPath), &multipartBuf)
572572
req.Header.Set("Content-Type", mpWriter.FormDataContentType())
@@ -624,7 +624,7 @@ func TestGzipUploadThenGzipDownload(t *testing.T) {
624624
EnvVars: utils.NewEnvVars(),
625625
User: currentUser.Username,
626626
}
627-
api := New(&logger, defaults, nil, false, cgroups.NewNoopManager())
627+
api := New(&logger, defaults, nil, false, cgroups.NewWorkloadFreezer(cgroups.NewNoopManager()))
628628

629629
uploadReq := httptest.NewRequestWithContext(t.Context(), http.MethodPost, "/files?path="+url.QueryEscape(destPath), &gzBuf)
630630
uploadReq.Header.Set("Content-Type", mpWriter.FormDataContentType())

packages/envd/internal/api/init.go

Lines changed: 15 additions & 45 deletions
Original file line numberDiff line numberDiff line change
@@ -22,7 +22,6 @@ import (
2222
"github.com/e2b-dev/infra/packages/envd/internal/host"
2323
"github.com/e2b-dev/infra/packages/envd/internal/logs"
2424
"github.com/e2b-dev/infra/packages/envd/internal/logs/ratelimit"
25-
"github.com/e2b-dev/infra/packages/envd/internal/services/cgroups"
2625
"github.com/e2b-dev/infra/packages/envd/pkg"
2726
"github.com/e2b-dev/infra/packages/shared/pkg/keys"
2827
)
@@ -268,11 +267,6 @@ func (a *API) SetData(ctx context.Context, logger zerolog.Logger, data PostInitJ
268267
}
269268

270269
// userCgroupsToFreeze is the cgroup set frozen pre-pause and thawed on /init.
271-
var userCgroupsToFreeze = []cgroups.ProcessType{
272-
cgroups.ProcessTypeUser,
273-
cgroups.ProcessTypePTY,
274-
}
275-
276270
// PostFreeze freezes user/pty cgroups directly (no Process.Start / shell).
277271
// Orchestrator calls this just before pause; the frozen state persists into the
278272
// snapshot and /init thaws on resume. Best-effort: tries every cgroup even if
@@ -282,22 +276,16 @@ func (a *API) PostFreeze(w http.ResponseWriter, r *http.Request) {
282276

283277
logger := a.logger.With().Str(string(logs.OperationIDKey), logs.AssignOperationID()).Logger()
284278

285-
if err := a.freezeLock.Acquire(r.Context(), 1); err != nil {
286-
w.WriteHeader(http.StatusServiceUnavailable)
287-
288-
return
289-
}
290-
defer a.freezeLock.Release(1)
279+
if err := a.workloadFreezer.Freeze(r.Context()); err != nil {
280+
// A failed lock acquire means the request ctx was cancelled; a failed
281+
// sweep leaves ctx intact.
282+
if r.Context().Err() != nil {
283+
w.WriteHeader(http.StatusServiceUnavailable)
291284

292-
var errs []error
293-
for _, pt := range userCgroupsToFreeze {
294-
if err := a.cgroupManager.Freeze(pt); err != nil {
295-
logger.Error().Err(err).Msgf("freeze %s cgroup", pt)
296-
errs = append(errs, fmt.Errorf("freeze %s cgroup: %w", pt, err))
285+
return
297286
}
298-
}
299-
if len(errs) > 0 {
300-
jsonError(w, http.StatusInternalServerError, errors.Join(errs...))
287+
logger.Error().Err(err).Msg("freeze workload cgroups")
288+
jsonError(w, http.StatusInternalServerError, err)
301289

302290
return
303291
}
@@ -313,25 +301,11 @@ func (a *API) PostFreeze(w http.ResponseWriter, r *http.Request) {
313301
func (a *API) PostUnfreeze(w http.ResponseWriter, r *http.Request) {
314302
defer r.Body.Close()
315303

316-
ctx := r.Context()
317304
logger := a.logger.With().Str(string(logs.OperationIDKey), logs.AssignOperationID()).Logger()
318305

319-
if err := a.freezeLock.Acquire(context.WithoutCancel(ctx), 1); err != nil {
320-
w.WriteHeader(http.StatusServiceUnavailable)
321-
322-
return
323-
}
324-
defer a.freezeLock.Release(1)
325-
326-
var errs []error
327-
for _, pt := range userCgroupsToFreeze {
328-
if err := a.cgroupManager.Unfreeze(pt); err != nil {
329-
logger.Error().Err(err).Msgf("unfreeze %s cgroup", pt)
330-
errs = append(errs, fmt.Errorf("unfreeze %s cgroup: %w", pt, err))
331-
}
332-
}
333-
if len(errs) > 0 {
334-
jsonError(w, http.StatusInternalServerError, errors.Join(errs...))
306+
if err := a.workloadFreezer.Unfreeze(r.Context()); err != nil {
307+
logger.Error().Err(err).Msg("unfreeze workload cgroups")
308+
jsonError(w, http.StatusInternalServerError, err)
335309

336310
return
337311
}
@@ -341,15 +315,11 @@ func (a *API) PostUnfreeze(w http.ResponseWriter, r *http.Request) {
341315
}
342316

343317
// unfreezeUserCgroups unfreezes user/pty cgroups (idempotent if not frozen).
344-
// Wraps the context with WithoutCancel so the unfreeze always completes.
318+
// The freezer detaches the wait from ctx cancellation so the unfreeze always
319+
// completes.
345320
func (a *API) unfreezeUserCgroups(ctx context.Context, logger zerolog.Logger) {
346-
_ = a.freezeLock.Acquire(context.WithoutCancel(ctx), 1)
347-
defer a.freezeLock.Release(1)
348-
349-
for _, pt := range userCgroupsToFreeze {
350-
if err := a.cgroupManager.Unfreeze(pt); err != nil {
351-
logger.Warn().Err(err).Msgf("unfreeze %s cgroup", pt)
352-
}
321+
if err := a.workloadFreezer.Unfreeze(ctx); err != nil {
322+
logger.Warn().Err(err).Msg("unfreeze workload cgroups")
353323
}
354324
}
355325

packages/envd/internal/api/init_test.go

Lines changed: 6 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -148,7 +148,7 @@ func newTestAPI(accessToken *SecureToken, mmdsClient MMDSClient) *API {
148148
defaults := &execcontext.Defaults{
149149
EnvVars: utils.NewEnvVars(),
150150
}
151-
api := New(&logger, defaults, nil, false, cgroups.NewNoopManager())
151+
api := New(&logger, defaults, nil, false, cgroups.NewWorkloadFreezer(cgroups.NewNoopManager()))
152152
if accessToken != nil {
153153
api.accessToken.TakeFrom(accessToken)
154154
}
@@ -637,7 +637,7 @@ func (f *fakeCgroupManager) Close() error { return nil }
637637
func newAPIWithCgroupManager(mgr cgroups.Manager) *API {
638638
logger := zerolog.Nop()
639639

640-
return New(&logger, &execcontext.Defaults{EnvVars: utils.NewEnvVars()}, nil, false, mgr)
640+
return New(&logger, &execcontext.Defaults{EnvVars: utils.NewEnvVars()}, nil, false, cgroups.NewWorkloadFreezer(mgr))
641641
}
642642

643643
func TestPostFreeze(t *testing.T) {
@@ -654,7 +654,7 @@ func TestPostFreeze(t *testing.T) {
654654
api.PostFreeze(rec, req)
655655

656656
require.Equal(t, http.StatusNoContent, rec.Code)
657-
assert.Equal(t, userCgroupsToFreeze, mgr.frozen)
657+
assert.Equal(t, cgroups.WorkloadProcessTypes, mgr.frozen)
658658
})
659659

660660
t.Run("returns 500 on freeze error", func(t *testing.T) {
@@ -686,7 +686,7 @@ func TestPostUnfreeze(t *testing.T) {
686686
api.PostUnfreeze(rec, req)
687687

688688
require.Equal(t, http.StatusNoContent, rec.Code)
689-
assert.Equal(t, userCgroupsToFreeze, mgr.unfrozen)
689+
assert.Equal(t, cgroups.WorkloadProcessTypes, mgr.unfrozen)
690690
})
691691

692692
t.Run("returns 500 but attempts every cgroup on unfreeze error", func(t *testing.T) {
@@ -701,7 +701,7 @@ func TestPostUnfreeze(t *testing.T) {
701701

702702
assert.Equal(t, http.StatusInternalServerError, rec.Code)
703703
assert.Empty(t, mgr.unfrozen)
704-
assert.Equal(t, userCgroupsToFreeze, mgr.unfreezeAttempts)
704+
assert.Equal(t, cgroups.WorkloadProcessTypes, mgr.unfreezeAttempts)
705705
})
706706
}
707707

@@ -732,7 +732,7 @@ func TestPostInit_UnfreezeOnStaleTimestamp(t *testing.T) {
732732
require.Equal(t, http.StatusNoContent, rec.Code)
733733
_, ok := api.defaults.EnvVars.Load("SHOULD_NOT_BE_SET")
734734
assert.False(t, ok, "stale /init should not apply EnvVars")
735-
assert.Equal(t, userCgroupsToFreeze, mgr.unfrozen, "stale /init must still unfreeze")
735+
assert.Equal(t, cgroups.WorkloadProcessTypes, mgr.unfrozen, "stale /init must still unfreeze")
736736
}
737737

738738
// Unauthorized /init must NOT thaw cgroups.

packages/envd/internal/api/store.go

Lines changed: 9 additions & 11 deletions
Original file line numberDiff line numberDiff line change
@@ -43,14 +43,13 @@ type API struct {
4343
initLock *semaphore.Weighted
4444

4545
caCertInstaller *host.CACertInstaller
46-
cgroupManager cgroups.Manager
47-
// freezeLock serializes the per-cgroup sweep across /freeze, /unfreeze
48-
// and the /init deferred unfreeze. PostFreeze acquires with the request
49-
// ctx; unfreeze paths acquire with Background so they always land
50-
// regardless of HTTP-client cancellation.
51-
freezeLock *semaphore.Weighted
52-
isMountingNFS atomic.Bool
53-
mountedPaths sync.Map // map[path]lifecycleID - tracks which lifecycle each path was mounted for
46+
// workloadFreezer freezes/thaws the user+pty cgroups. Shared with the process
47+
// service (the live-upgrade handover) so every freeze/unfreeze caller — this
48+
// API's /freeze, /unfreeze and /init deferred thaw, plus the upgrade — is
49+
// serialized through one lock.
50+
workloadFreezer *cgroups.WorkloadFreezer
51+
isMountingNFS atomic.Bool
52+
mountedPaths sync.Map // map[path]lifecycleID - tracks which lifecycle each path was mounted for
5453

5554
// fsFreezer freezes/thaws the guest rootfs for filesystem-only pauses;
5655
// fsFreezeLock serializes /fsfreeze and /fsthaw.
@@ -101,7 +100,7 @@ func (a *API) SetHandoverResult(procs, procsFailed, retained, retainedFailed, wa
101100
}
102101
}
103102

104-
func New(l *zerolog.Logger, defaults *execcontext.Defaults, mmdsChan chan *host.MMDSOpts, isNotFC bool, cgroupManager cgroups.Manager) *API {
103+
func New(l *zerolog.Logger, defaults *execcontext.Defaults, mmdsChan chan *host.MMDSOpts, isNotFC bool, workloadFreezer *cgroups.WorkloadFreezer) *API {
105104
return &API{
106105
logger: l,
107106
defaults: defaults,
@@ -111,9 +110,8 @@ func New(l *zerolog.Logger, defaults *execcontext.Defaults, mmdsChan chan *host.
111110
lastSetTime: utils.NewAtomicMax(),
112111
accessToken: &SecureToken{},
113112
caCertInstaller: host.NewCACertInstaller(l),
114-
cgroupManager: cgroupManager,
113+
workloadFreezer: workloadFreezer,
115114
initLock: semaphore.NewWeighted(1),
116-
freezeLock: semaphore.NewWeighted(1),
117115
fsFreezer: fsfreeze.New(),
118116
fsFreezeLock: semaphore.NewWeighted(1),
119117
}

0 commit comments

Comments
 (0)