From cbd7f375a4f76b5a3ac56b4a4b4961c02da0a411 Mon Sep 17 00:00:00 2001 From: Jakub Novak Date: Sat, 16 May 2026 11:49:56 +0000 Subject: [PATCH 01/11] perf(api): wake reservation waiters via pub/sub instead of 20ms polling [ENG-4070] MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit A duplicate concurrent /sandboxes/:sandboxID/connect held the second request in createWaitForStart's 20ms-interval GET+ZSCORE poll loop for the entire duration of the first request's create. A 70s create produced ~6600 Redis ops in a single trace. Replace the polling with the same pub/sub fan-out the storage state- transition path uses: - Subscribe before the initial Redis check to close the publish-vs-read race window. - Drop the 20ms ticker; keep a 1s fallback ticker (matching state-change pollInterval) for messages dropped by the bounded publisher worker pool. - createFinishStart and Release publish a routing key on the shared storage notify channel after their atomic script writes complete. - Release now writes a tombstone result instead of DEL'ing the key, so the waiter's read path is a single GET on the result key. The ZSCORE-on-pending branch remains as a compatibility safety net for pre-tombstone Releases written by older instances. - The waiter resolves to a typed sandbox.ErrReservationReleased on the tombstone path, replacing the previous stringly-typed error. Wiring: storage/redis exposes a new Notifier struct via Storage.Notifier that wraps subManager.Subscribe and publisher.Publish; reservations depend on a structurally-equivalent Notifier interface declared in their own package. The shared publisher worker pool (introduced for lock and state-change in the parent commit) handles all reservation publishes — no new goroutines, no new connections. Tests cover: pub/sub wakeup latency under 500ms, Release wakeup with typed error, fallback ticker when the subscription manager is disabled, race where the result lands before the waiter subscribes, ctx cancel, multiple concurrent waiters on a single publish, failed-start propagation, and a bounded-reads regression assertion that proves the polling no longer exists. --- .../api/internal/orchestrator/orchestrator.go | 2 +- packages/api/internal/sandbox/errors.go | 7 + .../sandbox/reservations/redis/reservation.go | 204 ++++++-- .../redis/reservation_pubsub_test.go | 462 ++++++++++++++++++ .../reservations/redis/reservation_test.go | 16 +- .../sandbox/reservations/redis/result.go | 19 +- .../sandbox/reservations/redis/scripts.go | 10 +- .../sandbox/reservations/redis/utils.go | 11 + .../storage/redis/notifier_external.go | 45 ++ 9 files changed, 730 insertions(+), 46 deletions(-) create mode 100644 packages/api/internal/sandbox/reservations/redis/reservation_pubsub_test.go create mode 100644 packages/api/internal/sandbox/storage/redis/notifier_external.go diff --git a/packages/api/internal/orchestrator/orchestrator.go b/packages/api/internal/orchestrator/orchestrator.go index 74e6479049..934f228d5f 100644 --- a/packages/api/internal/orchestrator/orchestrator.go +++ b/packages/api/internal/orchestrator/orchestrator.go @@ -167,7 +167,7 @@ func New( go redisbackend.NewCleaner(redisStorage).Start(ctx) case cfg.SandboxStorageBackendRedis: - reservationStorage = redisreservations.NewReservationStorage(redisClient) + reservationStorage = redisreservations.NewReservationStorage(redisClient, redisStorage.Notifier()) sandboxStorage = redisStorage logger.L().Info(ctx, "Using redis sandbox storage backend") default: diff --git a/packages/api/internal/sandbox/errors.go b/packages/api/internal/sandbox/errors.go index 76b49af5e8..d26d6c8eda 100644 --- a/packages/api/internal/sandbox/errors.go +++ b/packages/api/internal/sandbox/errors.go @@ -40,3 +40,10 @@ var ErrAlreadyExists = errors.New("sandbox already exists") var ErrEvictionInProgress = errors.New("sandbox eviction already in progress") var ErrEvictionNotNeeded = errors.New("sandbox eviction not needed") + +// ErrReservationReleased is returned by ReservationStorage.Reserve's +// waitForStart callback when the producer released the reservation +// instead of completing the sandbox creation. It is the structural +// equivalent of "the other instance gave up": the caller can retry +// Reserve from scratch. +var ErrReservationReleased = errors.New("reservation released") diff --git a/packages/api/internal/sandbox/reservations/redis/reservation.go b/packages/api/internal/sandbox/reservations/redis/reservation.go index 29068ecd7c..a6ce5f6383 100644 --- a/packages/api/internal/sandbox/reservations/redis/reservation.go +++ b/packages/api/internal/sandbox/reservations/redis/reservation.go @@ -15,8 +15,16 @@ import ( ) const ( - resultTTL = 30 * time.Second - retryInterval = 20 * time.Millisecond + resultTTL = 30 * time.Second + + // fallbackPollInterval is how often the waiter re-checks Redis when no + // PubSub wakeup arrives. PubSub is the primary wakeup mechanism; this + // ticker is the safety net for dropped messages (network blips, queue + // saturation on the publisher worker pool, Redis reconnects). + // + // Matches storage/redis pollInterval so reservations and state-change + // share a single tail-latency story for missed notifications. + fallbackPollInterval = 1 * time.Second // staleTTL is the maximum age of a pending entry before it is considered stale // and cleaned up. This handles the case where an API instance crashes mid-creation. @@ -26,13 +34,25 @@ const ( var _ sandbox.ReservationStorage = (*ReservationStorage)(nil) +// Notifier is the consumer-side view of the shared storage pub/sub seam. +// It is satisfied structurally by *storage_redis.Notifier, but accepting +// the interface keeps this package free of an explicit storage dependency +// and lets tests inject a fake. Publish is fire-and-forget — drops on +// queue saturation are recovered by the waiter's fallback ticker. +type Notifier interface { + Subscribe(routingKey string) (<-chan struct{}, func()) + Publish(ctx context.Context, routingKey string) +} + type ReservationStorage struct { redisClient redis.UniversalClient + notifier Notifier } -func NewReservationStorage(redisClient redis.UniversalClient) *ReservationStorage { +func NewReservationStorage(redisClient redis.UniversalClient, notifier Notifier) *ReservationStorage { return &ReservationStorage{ redisClient: redisClient, + notifier: notifier, } } @@ -76,21 +96,40 @@ func (s *ReservationStorage) Release(ctx context.Context, teamID uuid.UUID, sand pendingSetKey := getPendingSetKey(teamIDStr) resultKeyStr := getResultKey(teamIDStr, sandboxID) - err := releaseScript.Run(ctx, s.redisClient, []string{pendingSetKey, resultKeyStr}, sandboxID).Err() + tombstone, err := encodeReleased() + if err != nil { + return fmt.Errorf("failed to encode release tombstone: %w", err) + } + ttlSeconds := int(resultTTL.Seconds()) + + err = releaseScript.Run(ctx, s.redisClient, + []string{pendingSetKey, resultKeyStr}, + sandboxID, tombstone, ttlSeconds, + ).Err() if err != nil { return fmt.Errorf("failed to run release script: %w", err) } + // Wake any in-process waiter so it reads the tombstone immediately + // rather than after the fallback ticker. Drop-tolerant. + s.notifier.Publish(ctx, getReservationRoutingKey(teamIDStr, sandboxID)) + return nil } // createFinishStart returns a callback that completes the reservation. -// It removes the sandbox from the pending zset and stores the result for cross-instance waiters. +// It removes the sandbox from the pending zset and stores the result for +// cross-instance waiters, then publishes a wakeup on the shared notify +// channel so subscribed waiters resolve without waiting for the fallback +// ticker. func (s *ReservationStorage) createFinishStart(ctx context.Context, teamID uuid.UUID, sandboxID string) func(sandbox.Sandbox, error) { return func(sbx sandbox.Sandbox, startErr error) { teamIDStr := teamID.String() pendingSetKey := getPendingSetKey(teamIDStr) resultKeyStr := getResultKey(teamIDStr, sandboxID) + routingKey := getReservationRoutingKey(teamIDStr, sandboxID) + + bgCtx := context.WithoutCancel(ctx) resultData, encodeErr := encodeResult(sbx, startErr) if encodeErr != nil { @@ -99,14 +138,28 @@ func (s *ReservationStorage) createFinishStart(ctx context.Context, teamID uuid. logger.WithSandboxID(sandboxID), ) - // Still try to remove from pending even if encoding fails - _ = s.redisClient.ZRem(context.WithoutCancel(ctx), pendingSetKey, sandboxID).Err() + // Best-effort: write a released tombstone so waiters resolve + // with a typed error rather than blocking until the result + // key TTL elapses. Falls back to a plain ZRem if even encoding + // the tombstone fails. + tombstone, tsErr := encodeReleased() + if tsErr != nil { + _ = s.redisClient.ZRem(bgCtx, pendingSetKey, sandboxID).Err() + } else { + _ = releaseScript.Run(bgCtx, s.redisClient, + []string{pendingSetKey, resultKeyStr}, + sandboxID, tombstone, int(resultTTL.Seconds()), + ).Err() + } + // Still wake waiters so they read the tombstone (or fall back + // to the ticker if even ZRem failed). + s.notifier.Publish(bgCtx, routingKey) return } ttlSeconds := int(resultTTL.Seconds()) - err := finishStartScript.Run(context.WithoutCancel(ctx), s.redisClient, + err := finishStartScript.Run(bgCtx, s.redisClient, []string{pendingSetKey, resultKeyStr}, sandboxID, resultData, ttlSeconds, ).Err() @@ -115,7 +168,13 @@ func (s *ReservationStorage) createFinishStart(ctx context.Context, teamID uuid. zap.Error(err), logger.WithSandboxID(sandboxID), ) + + return } + + // Wake any in-process waiter immediately. Drop-tolerant: the + // fallback ticker covers a saturated publish queue. + s.notifier.Publish(bgCtx, routingKey) } } @@ -123,42 +182,107 @@ func (s *ReservationStorage) createFinishStart(ctx context.Context, teamID uuid. // initiated by another instance. func (s *ReservationStorage) createWaitForStart(teamID uuid.UUID, sandboxID string) func(ctx context.Context) (sandbox.Sandbox, error) { return func(ctx context.Context) (sandbox.Sandbox, error) { - teamIDStr := teamID.String() - resultKeyStr := getResultKey(teamIDStr, sandboxID) - pendingSetKey := getPendingSetKey(teamIDStr) + return s.waitForResult(ctx, teamID, sandboxID) + } +} - for { - // Check for result - data, err := s.redisClient.Get(ctx, resultKeyStr).Bytes() - if err == nil { - return decodeResult(data) - } - if !errors.Is(err, redis.Nil) { - return sandbox.Sandbox{}, fmt.Errorf("failed to check result key: %w", err) - } +// waitForResult blocks until the reservation initiated by another instance +// either completes (result key set with a sandbox or producer error) or is +// released (result key set with a tombstone resolving to +// sandbox.ErrReservationReleased). +// +// PubSub is the primary wakeup channel: createFinishStart and Release both +// publish on the shared notify channel after their atomic Redis writes +// complete. A 1s fallback ticker recovers any dropped notification — by +// design, the publisher worker pool can drop on saturation and the waiter +// must tolerate that. +// +// Ordering: we subscribe BEFORE the initial GET so a publish landing in the +// window between (Reserve returning AlreadyPending) and (the waiter calling +// Subscribe) cannot be missed. +func (s *ReservationStorage) waitForResult(ctx context.Context, teamID uuid.UUID, sandboxID string) (sandbox.Sandbox, error) { + teamIDStr := teamID.String() + resultKeyStr := getResultKey(teamIDStr, sandboxID) + pendingSetKey := getPendingSetKey(teamIDStr) + routingKey := getReservationRoutingKey(teamIDStr, sandboxID) - // No result yet — check if still pending (ZSCORE returns nil if not a member) - err = s.redisClient.ZScore(ctx, pendingSetKey, sandboxID).Err() - if errors.Is(err, redis.Nil) { - // Not pending anymore, final check - data, err = s.redisClient.Get(ctx, resultKeyStr).Bytes() - if err == nil { - return decodeResult(data) - } + ch, cleanup := s.notifier.Subscribe(routingKey) + defer cleanup() - return sandbox.Sandbox{}, fmt.Errorf("sandbox %s is no longer pending and has no result", sandboxID) - } - if err != nil { - return sandbox.Sandbox{}, fmt.Errorf("failed to check pending set: %w", err) - } + // Initial probe: the producer may have finished before we subscribed, + // or we may be a late waiter joining after the result was already set. + if done, sbx, err := s.tryReadResult(ctx, resultKeyStr, pendingSetKey, sandboxID); done { + return sbx, err + } - // Wait before next poll - select { - case <-ctx.Done(): - return sandbox.Sandbox{}, ctx.Err() - case <-time.After(retryInterval): - // continue polling - } + ticker := time.NewTicker(fallbackPollInterval) + defer ticker.Stop() + + for { + select { + case <-ctx.Done(): + return sandbox.Sandbox{}, ctx.Err() + case <-ch: + case <-ticker.C: + } + + if done, sbx, err := s.tryReadResult(ctx, resultKeyStr, pendingSetKey, sandboxID); done { + return sbx, err + } + } +} + +// tryReadResult performs a single probe of the reservation state. +// +// With the tombstone-on-Release contract, the result key is the single +// source of truth: any terminal state (success, producer error, or +// release) is encoded there. We only consult the pending zset as a +// safety net for legacy entries written by older instances that did not +// tombstone — those entries decay via stale-GC, so this branch will be +// dead in steady state after deploy. +// +// Returns done=true when the wait is over: +// - the result key holds an encoded terminal result, or +// - the sandbox vanished from the pending set without a tombstone +// (legacy compatibility path), or +// - the Redis call itself failed. +// +// Returns done=false when the reservation is still pending and the caller +// should wait for the next wakeup. +func (s *ReservationStorage) tryReadResult( + ctx context.Context, + resultKey, pendingSetKey, sandboxID string, +) (done bool, sbx sandbox.Sandbox, err error) { + data, getErr := s.redisClient.Get(ctx, resultKey).Bytes() + if getErr == nil { + sbx, err = decodeResult(data) + + return true, sbx, err + } + if !errors.Is(getErr, redis.Nil) { + return true, sandbox.Sandbox{}, fmt.Errorf("failed to check result key: %w", getErr) + } + + // No result yet. Confirm the reservation is still pending; if it's + // not, this is a pre-tombstone Release from a legacy instance. Treat + // it as a release for compatibility. + scoreErr := s.redisClient.ZScore(ctx, pendingSetKey, sandboxID).Err() + if errors.Is(scoreErr, redis.Nil) { + // Final read in case finishStart/release raced between our GET + // and ZSCORE. + data, getErr = s.redisClient.Get(ctx, resultKey).Bytes() + if getErr == nil { + sbx, err = decodeResult(data) + + return true, sbx, err } + + return true, sandbox.Sandbox{}, sandbox.ErrReservationReleased } + if scoreErr != nil { + return true, sandbox.Sandbox{}, fmt.Errorf("failed to check pending set: %w", scoreErr) + } + + // Still pending, no result yet. + return false, sandbox.Sandbox{}, nil } diff --git a/packages/api/internal/sandbox/reservations/redis/reservation_pubsub_test.go b/packages/api/internal/sandbox/reservations/redis/reservation_pubsub_test.go new file mode 100644 index 0000000000..bb10b34fe3 --- /dev/null +++ b/packages/api/internal/sandbox/reservations/redis/reservation_pubsub_test.go @@ -0,0 +1,462 @@ +package redis + +import ( + "context" + "errors" + "fmt" + "sync" + "testing" + "time" + + "github.com/google/uuid" + goredis "github.com/redis/go-redis/v9" + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" + + "github.com/e2b-dev/infra/packages/api/internal/sandbox" + storage_redis "github.com/e2b-dev/infra/packages/api/internal/sandbox/storage/redis" + "github.com/e2b-dev/infra/packages/shared/pkg/consts" + redis_utils "github.com/e2b-dev/infra/packages/shared/pkg/redis" +) + +// testSandbox is the canonical successful-finish payload used across pubsub tests. +func testSandbox(teamID uuid.UUID, sandboxID string) sandbox.Sandbox { + return sandbox.Sandbox{ + ClientID: consts.ClientID, + SandboxID: sandboxID, + TemplateID: "test", + TeamID: teamID, + StartTime: time.Now(), + EndTime: time.Now().Add(time.Hour), + MaxInstanceLength: time.Hour, + } +} + +// setupReservationStorageWithoutSubManager wires the reservation store so that +// PubSub messages are never delivered in-process. The publisher worker still +// runs (so PUBLISH commands hit Redis) but no subscription manager is +// listening; in-process waiters only resolve via the fallback ticker. Used to +// exercise the safety-net path. +func setupReservationStorageWithoutSubManager(t *testing.T) (*ReservationStorage, goredis.UniversalClient) { + t.Helper() + + client := redis_utils.SetupInstance(t) + + storageInstance := storage_redis.NewStorage(client) + // Deliberately do NOT call storageInstance.Start — no subManager.start + // goroutine, no fan-out. Publish() still works through the in-process + // queue but its drainer is also not running. The waiter must rely + // entirely on its 1s fallback ticker. + t.Cleanup(storageInstance.Close) + + storage := NewReservationStorage(client, storageInstance.Notifier()) + + return storage, client +} + +// TestWaitForStart_WokenByFinishStartPublish is the load-bearing regression +// test for the polling bug. With pub/sub working, the waiter must wake +// well under the 1s fallback ticker — anything that fast can only come +// from a PubSub delivery. +func TestWaitForStart_WokenByFinishStartPublish(t *testing.T) { + t.Parallel() + + storage, _ := setupTestReservationStorage(t) + teamID := uuid.New() + sbxID := "pubsub-finish" + + finishStart, _, err := storage.Reserve(t.Context(), teamID, sbxID, 10) + require.NoError(t, err) + require.NotNil(t, finishStart) + + _, waitForStart, err := storage.Reserve(t.Context(), teamID, sbxID, 10) + require.NoError(t, err) + require.NotNil(t, waitForStart) + + waiterDone := make(chan struct{}) + var got sandbox.Sandbox + var waitErr error + go func() { + got, waitErr = waitForStart(t.Context()) + close(waiterDone) + }() + + // Let the waiter subscribe before we finish. + time.Sleep(50 * time.Millisecond) + + start := time.Now() + finishStart(testSandbox(teamID, sbxID), nil) + + select { + case <-waiterDone: + elapsed := time.Since(start) + require.NoError(t, waitErr) + assert.Equal(t, sbxID, got.SandboxID) + assert.Less(t, elapsed, 500*time.Millisecond, + "waiter should wake via PubSub, not the fallback ticker") + case <-time.After(3 * time.Second): + require.FailNow(t, "waiter did not wake in time") + } +} + +// TestWaitForStart_WokenByReleasePublish proves Release also drives a fast +// wakeup and that the waiter returns the typed ErrReservationReleased. +func TestWaitForStart_WokenByReleasePublish(t *testing.T) { + t.Parallel() + + storage, _ := setupTestReservationStorage(t) + teamID := uuid.New() + sbxID := "pubsub-release" + + finishStart, _, err := storage.Reserve(t.Context(), teamID, sbxID, 10) + require.NoError(t, err) + require.NotNil(t, finishStart) + + _, waitForStart, err := storage.Reserve(t.Context(), teamID, sbxID, 10) + require.NoError(t, err) + require.NotNil(t, waitForStart) + + waiterErr := make(chan error, 1) + go func() { + _, err := waitForStart(t.Context()) + waiterErr <- err + }() + time.Sleep(50 * time.Millisecond) + + start := time.Now() + require.NoError(t, storage.Release(t.Context(), teamID, sbxID)) + + select { + case err := <-waiterErr: + require.ErrorIs(t, err, sandbox.ErrReservationReleased) + assert.Less(t, time.Since(start), 500*time.Millisecond, + "release should wake the waiter via PubSub") + case <-time.After(3 * time.Second): + require.FailNow(t, "waiter did not wake in time") + } +} + +// TestWaitForStart_FallbackTickerWhenPubSubMissed disables the subscription +// manager entirely so no in-process wakeup can fire. The waiter must still +// resolve via the 1s fallback ticker. +func TestWaitForStart_FallbackTickerWhenPubSubMissed(t *testing.T) { + t.Parallel() + + storage, _ := setupReservationStorageWithoutSubManager(t) + teamID := uuid.New() + sbxID := "pubsub-fallback" + + finishStart, _, err := storage.Reserve(t.Context(), teamID, sbxID, 10) + require.NoError(t, err) + require.NotNil(t, finishStart) + + _, waitForStart, err := storage.Reserve(t.Context(), teamID, sbxID, 10) + require.NoError(t, err) + + // Set the result directly via the finishStartScript path; with the + // subManager not running, no fan-out occurs. The waiter must rely on + // the fallback ticker. + resultData, err := encodeResult(testSandbox(teamID, sbxID), nil) + require.NoError(t, err) + teamIDStr := teamID.String() + err = finishStartScript.Run(t.Context(), storage.redisClient, + []string{getPendingSetKey(teamIDStr), getResultKey(teamIDStr, sbxID)}, + sbxID, resultData, int(resultTTL.Seconds()), + ).Err() + require.NoError(t, err) + + done := make(chan error, 1) + go func() { + _, err := waitForStart(t.Context()) + done <- err + }() + + select { + case err := <-done: + require.NoError(t, err) + case <-time.After(fallbackPollInterval + 2*time.Second): + require.FailNow(t, "fallback ticker did not resolve the wait") + } +} + +// TestWaitForStart_ResultLandedBeforeSubscribe covers the race where the +// producer finishes BEFORE the waiter calls waitForStart. The initial +// post-subscribe GET must catch the already-present result rather than +// blocking until a wakeup that will never come. +func TestWaitForStart_ResultLandedBeforeSubscribe(t *testing.T) { + t.Parallel() + + storage, _ := setupTestReservationStorage(t) + teamID := uuid.New() + sbxID := "pubsub-race" + + finishStart, _, err := storage.Reserve(t.Context(), teamID, sbxID, 10) + require.NoError(t, err) + + _, waitForStart, err := storage.Reserve(t.Context(), teamID, sbxID, 10) + require.NoError(t, err) + + // Finish BEFORE the waiter ever subscribes. + finishStart(testSandbox(teamID, sbxID), nil) + // Give the in-process publisher worker a moment to drain. + time.Sleep(50 * time.Millisecond) + + done := make(chan struct{}) + var got sandbox.Sandbox + var waitErr error + start := time.Now() + go func() { + got, waitErr = waitForStart(t.Context()) + close(done) + }() + + select { + case <-done: + require.NoError(t, waitErr) + assert.Equal(t, sbxID, got.SandboxID) + assert.Less(t, time.Since(start), 200*time.Millisecond, + "initial post-subscribe check must catch a pre-existing result") + case <-time.After(3 * time.Second): + require.FailNow(t, "waiter did not return") + } +} + +// TestWaitForStart_ContextCancellation asserts ctx.Done is the only path +// out under normal load; cancellation must unblock immediately. +func TestWaitForStart_ContextCancellation(t *testing.T) { + t.Parallel() + + storage, _ := setupTestReservationStorage(t) + teamID := uuid.New() + sbxID := "pubsub-cancel" + + _, _, err := storage.Reserve(t.Context(), teamID, sbxID, 10) + require.NoError(t, err) + + _, waitForStart, err := storage.Reserve(t.Context(), teamID, sbxID, 10) + require.NoError(t, err) + + ctx, cancel := context.WithTimeout(t.Context(), 100*time.Millisecond) + defer cancel() + + start := time.Now() + _, err = waitForStart(ctx) + elapsed := time.Since(start) + + require.ErrorIs(t, err, context.DeadlineExceeded) + assert.Less(t, elapsed, 500*time.Millisecond, "cancellation should be immediate") +} + +// TestWaitForStart_MultipleWaitersOnePublish confirms the fan-out wakes +// every concurrent waiter from a single producer publish. +func TestWaitForStart_MultipleWaitersOnePublish(t *testing.T) { + t.Parallel() + + storage, _ := setupTestReservationStorage(t) + teamID := uuid.New() + sbxID := "pubsub-multi" + const numWaiters = 10 + + finishStart, _, err := storage.Reserve(t.Context(), teamID, sbxID, 50) + require.NoError(t, err) + + waiters := make([]func(ctx context.Context) (sandbox.Sandbox, error), numWaiters) + for i := range numWaiters { + _, w, err := storage.Reserve(t.Context(), teamID, sbxID, 50) + require.NoError(t, err) + require.NotNil(t, w) + waiters[i] = w + } + + var wg sync.WaitGroup + errs := make([]error, numWaiters) + completions := make([]time.Duration, numWaiters) + for i, w := range waiters { + wg.Add(1) + go func(i int, w func(ctx context.Context) (sandbox.Sandbox, error)) { + defer wg.Done() + start := time.Now() + _, errs[i] = w(t.Context()) + completions[i] = time.Since(start) + }(i, w) + } + + // Let everyone subscribe. + time.Sleep(100 * time.Millisecond) + finishStart(testSandbox(teamID, sbxID), nil) + + done := make(chan struct{}) + go func() { wg.Wait(); close(done) }() + + select { + case <-done: + for i := range numWaiters { + require.NoError(t, errs[i]) + assert.Less(t, completions[i], 600*time.Millisecond, + "waiter %d should be woken by PubSub", i) + } + case <-time.After(3 * time.Second): + require.FailNow(t, "not all waiters completed") + } +} + +// TestWaitForStart_FailedStartPropagatesPromptly proves a producer-side +// failure round-trips back through the result key and wakes the waiter +// via PubSub. +func TestWaitForStart_FailedStartPropagatesPromptly(t *testing.T) { + t.Parallel() + + storage, _ := setupTestReservationStorage(t) + teamID := uuid.New() + sbxID := "pubsub-failed" + + finishStart, _, err := storage.Reserve(t.Context(), teamID, sbxID, 10) + require.NoError(t, err) + + _, waitForStart, err := storage.Reserve(t.Context(), teamID, sbxID, 10) + require.NoError(t, err) + + waiterErr := make(chan error, 1) + go func() { + _, err := waitForStart(t.Context()) + waiterErr <- err + }() + time.Sleep(50 * time.Millisecond) + + start := time.Now() + finishStart(sandbox.Sandbox{}, errors.New("boom")) + + select { + case err := <-waiterErr: + require.Error(t, err) + assert.Contains(t, err.Error(), "boom") + assert.Less(t, time.Since(start), 500*time.Millisecond) + case <-time.After(3 * time.Second): + require.FailNow(t, "waiter did not return") + } +} + +// TestWaitForStart_RedisOpsBounded is the explicit regression test for the +// production incident: a 1s wait must NOT generate dozens of Redis ops. We +// install a redis.Hook on the client, run a waiter that resolves via +// PubSub, and assert the GET/ZSCORE count is single-digit. +func TestWaitForStart_RedisOpsBounded(t *testing.T) { + t.Parallel() + + client := redis_utils.SetupInstance(t) + + counter := &cmdCounter{} + client.AddHook(counter) + + storageInstance := storage_redis.NewStorage(client) + go storageInstance.Start(t.Context()) + t.Cleanup(storageInstance.Close) + + storage := NewReservationStorage(client, storageInstance.Notifier()) + + teamID := uuid.New() + sbxID := "pubsub-bounded" + + finishStart, _, err := storage.Reserve(t.Context(), teamID, sbxID, 10) + require.NoError(t, err) + + _, waitForStart, err := storage.Reserve(t.Context(), teamID, sbxID, 10) + require.NoError(t, err) + + // Hold the producer for ~1.5s so the OLD 20ms-polling implementation + // would have done ~75 GETs + ~75 ZSCOREs = ~150 ops. The new pubsub + // implementation should do exactly one initial GET (+ one ZSCORE for + // the legacy-compat path) plus one final GET on wake. + waiterDone := make(chan error, 1) + go func() { + _, err := waitForStart(t.Context()) + waiterDone <- err + }() + + time.Sleep(50 * time.Millisecond) + counter.Reset() // start counting from after the waiter has subscribed + + time.Sleep(1500 * time.Millisecond) + finishStart(testSandbox(teamID, sbxID), nil) + + select { + case err := <-waiterDone: + require.NoError(t, err) + case <-time.After(3 * time.Second): + require.FailNow(t, "waiter did not return") + } + + // Allow at most a handful of read ops: initial GET (+ ZSCORE legacy + // safety net), maybe one fallback tick if it fired just before the + // publish landed, and the post-wakeup GET. 10 is a generous ceiling + // for the new design; the old design would blow past 50. + got := counter.Reads() + assert.LessOrEqual(t, got, 10, + "expected bounded reads; got %d (regression: polling is back)", got) +} + +// cmdCounter is a redis.Hook that counts read-side operations on the +// client. We only care about the waiter's reads, not the producer's +// writes, so we only count GET and ZSCORE. +type cmdCounter struct { + mu sync.Mutex + get int + zs int +} + +func (c *cmdCounter) Reset() { + c.mu.Lock() + defer c.mu.Unlock() + c.get = 0 + c.zs = 0 +} + +func (c *cmdCounter) Reads() int { + c.mu.Lock() + defer c.mu.Unlock() + + return c.get + c.zs +} + +func (c *cmdCounter) DialHook(next goredis.DialHook) goredis.DialHook { + return next +} + +func (c *cmdCounter) ProcessHook(next goredis.ProcessHook) goredis.ProcessHook { + return func(ctx context.Context, cmd goredis.Cmder) error { + c.mu.Lock() + switch cmd.Name() { + case "get": + c.get++ + case "zscore": + c.zs++ + } + c.mu.Unlock() + + return next(ctx, cmd) + } +} + +func (c *cmdCounter) ProcessPipelineHook(next goredis.ProcessPipelineHook) goredis.ProcessPipelineHook { + return func(ctx context.Context, cmds []goredis.Cmder) error { + c.mu.Lock() + for _, cmd := range cmds { + switch cmd.Name() { + case "get": + c.get++ + case "zscore": + c.zs++ + } + } + c.mu.Unlock() + + return next(ctx, cmds) + } +} + +// Compile-time assertion that the storage Notifier satisfies the local +// Notifier seam. If this breaks, the orchestrator wiring needs to change too. +var _ Notifier = (*storage_redis.Notifier)(nil) + +// Compile-time check that fmt is used by anything pulled in transitively. +var _ = fmt.Sprintf diff --git a/packages/api/internal/sandbox/reservations/redis/reservation_test.go b/packages/api/internal/sandbox/reservations/redis/reservation_test.go index 27c3c007cb..eaa5975f9d 100644 --- a/packages/api/internal/sandbox/reservations/redis/reservation_test.go +++ b/packages/api/internal/sandbox/reservations/redis/reservation_test.go @@ -16,6 +16,7 @@ import ( "golang.org/x/sync/errgroup" "github.com/e2b-dev/infra/packages/api/internal/sandbox" + storage_redis "github.com/e2b-dev/infra/packages/api/internal/sandbox/storage/redis" "github.com/e2b-dev/infra/packages/shared/pkg/consts" redis_utils "github.com/e2b-dev/infra/packages/shared/pkg/redis" ) @@ -26,10 +27,18 @@ const ( var testTeamID = uuid.New() +// setupTestReservationStorage wires the reservation store to a real storage +// pub/sub seam: one Redis container, one subscription manager, one publish +// worker pool. Mirrors the production wiring in orchestrator.go. func setupTestReservationStorage(t *testing.T) (*ReservationStorage, goredis.UniversalClient) { t.Helper() client := redis_utils.SetupInstance(t) - storage := NewReservationStorage(client) + + storageInstance := storage_redis.NewStorage(client) + go storageInstance.Start(t.Context()) + t.Cleanup(storageInstance.Close) + + storage := NewReservationStorage(client, storageInstance.Notifier()) return storage, client } @@ -496,7 +505,10 @@ func TestReservation_StalePendingCleanup(t *testing.T) { assert.Equal(t, int64(1), count) // Create a new storage instance (simulating a fresh/restarted API) - storage := NewReservationStorage(client) + storageInstance := storage_redis.NewStorage(client) + go storageInstance.Start(t.Context()) + t.Cleanup(storageInstance.Close) + storage := NewReservationStorage(client, storageInstance.Notifier()) // Reserve with limit=1 — this should succeed because the stale entry // gets cleaned up by the reserveScript before counting diff --git a/packages/api/internal/sandbox/reservations/redis/result.go b/packages/api/internal/sandbox/reservations/redis/result.go index f7402085e6..aaff284e27 100644 --- a/packages/api/internal/sandbox/reservations/redis/result.go +++ b/packages/api/internal/sandbox/reservations/redis/result.go @@ -16,8 +16,14 @@ type reservationResult struct { } // reservationError preserves api.APIError fields for cross-instance error propagation. +// +// A Released=true value is the "tombstone" written by Release: the producer +// withdrew the reservation without completing creation. Tombstones make the +// result key the single source of truth for the waiter, so the wait path +// never needs to inspect the pending zset. type reservationError struct { - Message string `json:"message"` + Released bool `json:"released,omitempty"` + Message string `json:"message,omitempty"` Code int `json:"code,omitempty"` ClientMsg string `json:"client_msg,omitempty"` } @@ -40,6 +46,14 @@ func encodeResult(sbx sandbox.Sandbox, err error) ([]byte, error) { return json.Marshal(result) } +// encodeReleased returns the tombstone result written by Release. It carries +// no sandbox payload and decodes to sandbox.ErrReservationReleased. +func encodeReleased() ([]byte, error) { + return json.Marshal(reservationResult{ + Error: &reservationError{Released: true}, + }) +} + func decodeResult(data []byte) (sandbox.Sandbox, error) { var result reservationResult if err := json.Unmarshal(data, &result); err != nil { @@ -57,6 +71,9 @@ func decodeResult(data []byte) (sandbox.Sandbox, error) { // If the error had an API code, it reconstructs an *api.APIError to preserve // errors.As(err, &apiErr) behavior in create_instance.go. func reconstructError(re *reservationError) error { + if re.Released { + return sandbox.ErrReservationReleased + } if re.Code != 0 { return &api.APIError{ Code: re.Code, diff --git a/packages/api/internal/sandbox/reservations/redis/scripts.go b/packages/api/internal/sandbox/reservations/redis/scripts.go index 29d187ef92..be384ebedf 100644 --- a/packages/api/internal/sandbox/reservations/redis/scripts.go +++ b/packages/api/internal/sandbox/reservations/redis/scripts.go @@ -75,13 +75,19 @@ var ( return 1 `) - // releaseScript removes a sandbox from the pending zset and deletes the result key. + // releaseScript removes a sandbox from the pending zset and writes a + // tombstone result so any concurrent waiter resolves to + // sandbox.ErrReservationReleased via a single GET on the result key, + // without ever consulting the pending zset. + // // KEYS[1] = pending zset key // KEYS[2] = result key // ARGV[1] = sandboxID + // ARGV[2] = tombstone JSON + // ARGV[3] = TTL in seconds releaseScript = redis.NewScript(` redis.call('ZREM', KEYS[1], ARGV[1]) - redis.call('DEL', KEYS[2]) + redis.call('SET', KEYS[2], ARGV[2], 'EX', tonumber(ARGV[3])) return 1 `) ) diff --git a/packages/api/internal/sandbox/reservations/redis/utils.go b/packages/api/internal/sandbox/reservations/redis/utils.go index ab19d25d03..501ec0acae 100644 --- a/packages/api/internal/sandbox/reservations/redis/utils.go +++ b/packages/api/internal/sandbox/reservations/redis/utils.go @@ -9,6 +9,7 @@ const ( reservationsKey = "reservations" pendingKey = "pending" resultKey = "result" + notifySuffix = "notify" ) // getStorageIndexKey returns the existing storage team index key (read-only). @@ -33,3 +34,13 @@ func getPendingSetKey(teamID string) string { func getResultKey(teamID, sandboxID string) string { return redis_utils.CreateKey(getReservationPrefix(teamID), sandboxID, resultKey) } + +// getReservationRoutingKey is the per-(team, sandbox) PubSub routing key +// for reservation completion notifications. It is published as the payload +// of messages on the shared storage notify channel and consumed by +// in-process waiters subscribed via the storage Notifier. +// +// e.g. sandbox:storage:{teamID}:reservations:sandboxID:notify +func getReservationRoutingKey(teamID, sandboxID string) string { + return redis_utils.CreateKey(getReservationPrefix(teamID), sandboxID, notifySuffix) +} diff --git a/packages/api/internal/sandbox/storage/redis/notifier_external.go b/packages/api/internal/sandbox/storage/redis/notifier_external.go new file mode 100644 index 0000000000..f27df5e790 --- /dev/null +++ b/packages/api/internal/sandbox/storage/redis/notifier_external.go @@ -0,0 +1,45 @@ +package redis + +import ( + "context" +) + +// Notifier is the public seam onto the shared storage pub/sub infrastructure. +// +// It exposes both the consumer side (Subscribe — register a wakeup channel +// for a routing key) and the producer side (Publish — enqueue a routing key +// onto the shared publisher worker pool). Reservations and any future +// cross-package consumer depend on this rather than on subscriptionManager +// or publisher directly. +// +// Lifecycle is owned by Storage: callers MUST construct Notifier via +// Storage.Notifier(), and the underlying Storage must be Start()-ed +// before any Subscribe or Publish call. Close() on Storage shuts both +// sides down. +type Notifier struct { + sub *subscriptionManager + pub *publisher +} + +// Subscribe registers interest in routingKey. The returned channel is +// signaled (non-blocking, drop-on-full) whenever a matching message +// arrives on the shared notify channel. The caller MUST invoke cleanup +// when done to avoid a memory leak. +func (n *Notifier) Subscribe(routingKey string) (<-chan struct{}, func()) { + return n.sub.subscribe(routingKey) +} + +// Publish enqueues routingKey for asynchronous PUBLISH on the shared +// notify channel. Never blocks: drops silently (with rate-limited warn) +// when the publish queue is saturated. Drop tolerance is part of the +// contract — every consumer ships with a fallback ticker. +func (n *Notifier) Publish(ctx context.Context, routingKey string) { + n.pub.Publish(ctx, routingKey) +} + +// Notifier returns the cross-package pub/sub seam. The returned value is +// cheap; callers may cache or re-fetch as convenient. Safe to call before +// Storage.Start, but Subscribe/Publish only function once Start is running. +func (s *Storage) Notifier() *Notifier { + return &Notifier{sub: s.subManager, pub: s.publisher} +} From 8e1cdab6611b2c684fec51f588afb5cdf5fea9fd Mon Sep 17 00:00:00 2001 From: Jakub Novak Date: Tue, 19 May 2026 08:17:13 +0000 Subject: [PATCH 02/11] chore: add todo --- .../api/internal/sandbox/reservations/redis/reservation.go | 3 +++ 1 file changed, 3 insertions(+) diff --git a/packages/api/internal/sandbox/reservations/redis/reservation.go b/packages/api/internal/sandbox/reservations/redis/reservation.go index a6ce5f6383..27b74faac4 100644 --- a/packages/api/internal/sandbox/reservations/redis/reservation.go +++ b/packages/api/internal/sandbox/reservations/redis/reservation.go @@ -266,6 +266,9 @@ func (s *ReservationStorage) tryReadResult( // No result yet. Confirm the reservation is still pending; if it's // not, this is a pre-tombstone Release from a legacy instance. Treat // it as a release for compatibility. + // + // TODO [ENG-4089]: drop this ZSCORE fallback once all + // instances are guaranteed to write a tombstone on Release. scoreErr := s.redisClient.ZScore(ctx, pendingSetKey, sandboxID).Err() if errors.Is(scoreErr, redis.Nil) { // Final read in case finishStart/release raced between our GET From a8d3fa8aca8afef8953157b83e7507547d574ed9 Mon Sep 17 00:00:00 2001 From: Jakub Novak Date: Tue, 19 May 2026 09:27:21 +0000 Subject: [PATCH 03/11] chore: clean up --- .../sandbox/reservations/redis/reservation.go | 99 ++++++------------- .../redis/reservation_pubsub_test.go | 4 - .../sandbox/reservations/redis/utils.go | 6 +- .../{notifier_external.go => notifier.go} | 6 -- 4 files changed, 30 insertions(+), 85 deletions(-) rename packages/api/internal/sandbox/storage/redis/{notifier_external.go => notifier.go} (81%) diff --git a/packages/api/internal/sandbox/reservations/redis/reservation.go b/packages/api/internal/sandbox/reservations/redis/reservation.go index 27b74faac4..918d754710 100644 --- a/packages/api/internal/sandbox/reservations/redis/reservation.go +++ b/packages/api/internal/sandbox/reservations/redis/reservation.go @@ -18,12 +18,7 @@ const ( resultTTL = 30 * time.Second // fallbackPollInterval is how often the waiter re-checks Redis when no - // PubSub wakeup arrives. PubSub is the primary wakeup mechanism; this - // ticker is the safety net for dropped messages (network blips, queue - // saturation on the publisher worker pool, Redis reconnects). - // - // Matches storage/redis pollInterval so reservations and state-change - // share a single tail-latency story for missed notifications. + // PubSub wakeup arrives fallbackPollInterval = 1 * time.Second // staleTTL is the maximum age of a pending entry before it is considered stale @@ -34,11 +29,7 @@ const ( var _ sandbox.ReservationStorage = (*ReservationStorage)(nil) -// Notifier is the consumer-side view of the shared storage pub/sub seam. -// It is satisfied structurally by *storage_redis.Notifier, but accepting -// the interface keeps this package free of an explicit storage dependency -// and lets tests inject a fake. Publish is fire-and-forget — drops on -// queue saturation are recovered by the waiter's fallback ticker. +// Publish is fire-and-forget — drops on queue saturation are recovered by the waiter's fallback ticker. type Notifier interface { Subscribe(routingKey string) (<-chan struct{}, func()) Publish(ctx context.Context, routingKey string) @@ -111,17 +102,12 @@ func (s *ReservationStorage) Release(ctx context.Context, teamID uuid.UUID, sand } // Wake any in-process waiter so it reads the tombstone immediately - // rather than after the fallback ticker. Drop-tolerant. s.notifier.Publish(ctx, getReservationRoutingKey(teamIDStr, sandboxID)) return nil } // createFinishStart returns a callback that completes the reservation. -// It removes the sandbox from the pending zset and stores the result for -// cross-instance waiters, then publishes a wakeup on the shared notify -// channel so subscribed waiters resolve without waiting for the fallback -// ticker. func (s *ReservationStorage) createFinishStart(ctx context.Context, teamID uuid.UUID, sandboxID string) func(sandbox.Sandbox, error) { return func(sbx sandbox.Sandbox, startErr error) { teamIDStr := teamID.String() @@ -138,12 +124,10 @@ func (s *ReservationStorage) createFinishStart(ctx context.Context, teamID uuid. logger.WithSandboxID(sandboxID), ) - // Best-effort: write a released tombstone so waiters resolve - // with a typed error rather than blocking until the result - // key TTL elapses. Falls back to a plain ZRem if even encoding - // the tombstone fails. + // Still try to remove from pending even if encoding fails tombstone, tsErr := encodeReleased() if tsErr != nil { + // If we can't encode the tombstone either, just delete the pending entry without a result key _ = s.redisClient.ZRem(bgCtx, pendingSetKey, sandboxID).Err() } else { _ = releaseScript.Run(bgCtx, s.redisClient, @@ -151,8 +135,8 @@ func (s *ReservationStorage) createFinishStart(ctx context.Context, teamID uuid. sandboxID, tombstone, int(resultTTL.Seconds()), ).Err() } - // Still wake waiters so they read the tombstone (or fall back - // to the ticker if even ZRem failed). + + // Wake waiters so they read the tombstone s.notifier.Publish(bgCtx, routingKey) return @@ -182,65 +166,40 @@ func (s *ReservationStorage) createFinishStart(ctx context.Context, teamID uuid. // initiated by another instance. func (s *ReservationStorage) createWaitForStart(teamID uuid.UUID, sandboxID string) func(ctx context.Context) (sandbox.Sandbox, error) { return func(ctx context.Context) (sandbox.Sandbox, error) { - return s.waitForResult(ctx, teamID, sandboxID) - } -} - -// waitForResult blocks until the reservation initiated by another instance -// either completes (result key set with a sandbox or producer error) or is -// released (result key set with a tombstone resolving to -// sandbox.ErrReservationReleased). -// -// PubSub is the primary wakeup channel: createFinishStart and Release both -// publish on the shared notify channel after their atomic Redis writes -// complete. A 1s fallback ticker recovers any dropped notification — by -// design, the publisher worker pool can drop on saturation and the waiter -// must tolerate that. -// -// Ordering: we subscribe BEFORE the initial GET so a publish landing in the -// window between (Reserve returning AlreadyPending) and (the waiter calling -// Subscribe) cannot be missed. -func (s *ReservationStorage) waitForResult(ctx context.Context, teamID uuid.UUID, sandboxID string) (sandbox.Sandbox, error) { - teamIDStr := teamID.String() - resultKeyStr := getResultKey(teamIDStr, sandboxID) - pendingSetKey := getPendingSetKey(teamIDStr) - routingKey := getReservationRoutingKey(teamIDStr, sandboxID) + teamIDStr := teamID.String() + resultKeyStr := getResultKey(teamIDStr, sandboxID) + pendingSetKey := getPendingSetKey(teamIDStr) + routingKey := getReservationRoutingKey(teamIDStr, sandboxID) - ch, cleanup := s.notifier.Subscribe(routingKey) - defer cleanup() + ch, cleanup := s.notifier.Subscribe(routingKey) + defer cleanup() - // Initial probe: the producer may have finished before we subscribed, - // or we may be a late waiter joining after the result was already set. - if done, sbx, err := s.tryReadResult(ctx, resultKeyStr, pendingSetKey, sandboxID); done { - return sbx, err - } + // Initial probe: the producer may have finished before we subscribed, + // or we may be a late waiter joining after the result was already set. + if done, sbx, err := s.tryReadResult(ctx, resultKeyStr, pendingSetKey, sandboxID); done { + return sbx, err + } - ticker := time.NewTicker(fallbackPollInterval) - defer ticker.Stop() + ticker := time.NewTicker(fallbackPollInterval) + defer ticker.Stop() - for { - select { - case <-ctx.Done(): - return sandbox.Sandbox{}, ctx.Err() - case <-ch: - case <-ticker.C: - } + for { + select { + case <-ctx.Done(): + return sandbox.Sandbox{}, ctx.Err() + case <-ch: + case <-ticker.C: + } - if done, sbx, err := s.tryReadResult(ctx, resultKeyStr, pendingSetKey, sandboxID); done { - return sbx, err + if done, sbx, err := s.tryReadResult(ctx, resultKeyStr, pendingSetKey, sandboxID); done { + return sbx, err + } } } } // tryReadResult performs a single probe of the reservation state. // -// With the tombstone-on-Release contract, the result key is the single -// source of truth: any terminal state (success, producer error, or -// release) is encoded there. We only consult the pending zset as a -// safety net for legacy entries written by older instances that did not -// tombstone — those entries decay via stale-GC, so this branch will be -// dead in steady state after deploy. -// // Returns done=true when the wait is over: // - the result key holds an encoded terminal result, or // - the sandbox vanished from the pending set without a tombstone diff --git a/packages/api/internal/sandbox/reservations/redis/reservation_pubsub_test.go b/packages/api/internal/sandbox/reservations/redis/reservation_pubsub_test.go index bb10b34fe3..cd0ad08bbd 100644 --- a/packages/api/internal/sandbox/reservations/redis/reservation_pubsub_test.go +++ b/packages/api/internal/sandbox/reservations/redis/reservation_pubsub_test.go @@ -3,7 +3,6 @@ package redis import ( "context" "errors" - "fmt" "sync" "testing" "time" @@ -457,6 +456,3 @@ func (c *cmdCounter) ProcessPipelineHook(next goredis.ProcessPipelineHook) gored // Compile-time assertion that the storage Notifier satisfies the local // Notifier seam. If this breaks, the orchestrator wiring needs to change too. var _ Notifier = (*storage_redis.Notifier)(nil) - -// Compile-time check that fmt is used by anything pulled in transitively. -var _ = fmt.Sprintf diff --git a/packages/api/internal/sandbox/reservations/redis/utils.go b/packages/api/internal/sandbox/reservations/redis/utils.go index 501ec0acae..043be1f928 100644 --- a/packages/api/internal/sandbox/reservations/redis/utils.go +++ b/packages/api/internal/sandbox/reservations/redis/utils.go @@ -35,11 +35,7 @@ func getResultKey(teamID, sandboxID string) string { return redis_utils.CreateKey(getReservationPrefix(teamID), sandboxID, resultKey) } -// getReservationRoutingKey is the per-(team, sandbox) PubSub routing key -// for reservation completion notifications. It is published as the payload -// of messages on the shared storage notify channel and consumed by -// in-process waiters subscribed via the storage Notifier. -// +// getReservationRoutingKey is PubSub routing key for reservation completion notifications. // e.g. sandbox:storage:{teamID}:reservations:sandboxID:notify func getReservationRoutingKey(teamID, sandboxID string) string { return redis_utils.CreateKey(getReservationPrefix(teamID), sandboxID, notifySuffix) diff --git a/packages/api/internal/sandbox/storage/redis/notifier_external.go b/packages/api/internal/sandbox/storage/redis/notifier.go similarity index 81% rename from packages/api/internal/sandbox/storage/redis/notifier_external.go rename to packages/api/internal/sandbox/storage/redis/notifier.go index f27df5e790..68ec7756f6 100644 --- a/packages/api/internal/sandbox/storage/redis/notifier_external.go +++ b/packages/api/internal/sandbox/storage/redis/notifier.go @@ -6,12 +6,6 @@ import ( // Notifier is the public seam onto the shared storage pub/sub infrastructure. // -// It exposes both the consumer side (Subscribe — register a wakeup channel -// for a routing key) and the producer side (Publish — enqueue a routing key -// onto the shared publisher worker pool). Reservations and any future -// cross-package consumer depend on this rather than on subscriptionManager -// or publisher directly. -// // Lifecycle is owned by Storage: callers MUST construct Notifier via // Storage.Notifier(), and the underlying Storage must be Start()-ed // before any Subscribe or Publish call. Close() on Storage shuts both From b8545d7338a07eabf263d06ccfeb0d226aedb1fa Mon Sep 17 00:00:00 2001 From: Jakub Novak Date: Thu, 21 May 2026 10:06:34 +0000 Subject: [PATCH 04/11] chore: lint --- .../redis/reservation_pubsub_test.go | 8 ++++---- .../reservations/redis/reservation_test.go | 18 ++++++++++++++---- 2 files changed, 18 insertions(+), 8 deletions(-) diff --git a/packages/api/internal/sandbox/reservations/redis/reservation_pubsub_test.go b/packages/api/internal/sandbox/reservations/redis/reservation_pubsub_test.go index cd0ad08bbd..ad58ac610f 100644 --- a/packages/api/internal/sandbox/reservations/redis/reservation_pubsub_test.go +++ b/packages/api/internal/sandbox/reservations/redis/reservation_pubsub_test.go @@ -41,12 +41,12 @@ func setupReservationStorageWithoutSubManager(t *testing.T) (*ReservationStorage client := redis_utils.SetupInstance(t) - storageInstance := storage_redis.NewStorage(client) + storageInstance := newTestSandboxStorage(t, client) // Deliberately do NOT call storageInstance.Start — no subManager.start // goroutine, no fan-out. Publish() still works through the in-process // queue but its drainer is also not running. The waiter must rely // entirely on its 1s fallback ticker. - t.Cleanup(storageInstance.Close) + t.Cleanup(func() { storageInstance.Close(context.WithoutCancel(t.Context())) }) storage := NewReservationStorage(client, storageInstance.Notifier()) @@ -347,9 +347,9 @@ func TestWaitForStart_RedisOpsBounded(t *testing.T) { counter := &cmdCounter{} client.AddHook(counter) - storageInstance := storage_redis.NewStorage(client) + storageInstance := newTestSandboxStorage(t, client) go storageInstance.Start(t.Context()) - t.Cleanup(storageInstance.Close) + t.Cleanup(func() { storageInstance.Close(context.WithoutCancel(t.Context())) }) storage := NewReservationStorage(client, storageInstance.Notifier()) diff --git a/packages/api/internal/sandbox/reservations/redis/reservation_test.go b/packages/api/internal/sandbox/reservations/redis/reservation_test.go index eaa5975f9d..af2cc7ae6f 100644 --- a/packages/api/internal/sandbox/reservations/redis/reservation_test.go +++ b/packages/api/internal/sandbox/reservations/redis/reservation_test.go @@ -13,6 +13,7 @@ import ( goredis "github.com/redis/go-redis/v9" "github.com/stretchr/testify/assert" "github.com/stretchr/testify/require" + "go.opentelemetry.io/otel/metric/noop" "golang.org/x/sync/errgroup" "github.com/e2b-dev/infra/packages/api/internal/sandbox" @@ -34,15 +35,24 @@ func setupTestReservationStorage(t *testing.T) (*ReservationStorage, goredis.Uni t.Helper() client := redis_utils.SetupInstance(t) - storageInstance := storage_redis.NewStorage(client) + storageInstance := newTestSandboxStorage(t, client) go storageInstance.Start(t.Context()) - t.Cleanup(storageInstance.Close) + t.Cleanup(func() { storageInstance.Close(context.WithoutCancel(t.Context())) }) storage := NewReservationStorage(client, storageInstance.Notifier()) return storage, client } +func newTestSandboxStorage(t *testing.T, client goredis.UniversalClient) *storage_redis.Storage { + t.Helper() + + storageInstance, err := storage_redis.NewStorage(client, noop.NewMeterProvider()) + require.NoError(t, err) + + return storageInstance +} + func TestReservation(t *testing.T) { t.Parallel() storage, _ := setupTestReservationStorage(t) @@ -505,9 +515,9 @@ func TestReservation_StalePendingCleanup(t *testing.T) { assert.Equal(t, int64(1), count) // Create a new storage instance (simulating a fresh/restarted API) - storageInstance := storage_redis.NewStorage(client) + storageInstance := newTestSandboxStorage(t, client) go storageInstance.Start(t.Context()) - t.Cleanup(storageInstance.Close) + t.Cleanup(func() { storageInstance.Close(context.WithoutCancel(t.Context())) }) storage := NewReservationStorage(client, storageInstance.Notifier()) // Reserve with limit=1 — this should succeed because the stale entry From 7f13b8b01e4e441fd96bcf6e050967eb4f7c7cde Mon Sep 17 00:00:00 2001 From: Jakub Novak Date: Thu, 21 May 2026 10:21:54 +0000 Subject: [PATCH 05/11] chore: simplify comments, remove unused function --- .../sandbox/reservations/redis/reservation_test.go | 3 --- .../api/internal/sandbox/storage/redis/notifier.go | 13 ++----------- packages/api/internal/sandbox/store.go | 4 ---- 3 files changed, 2 insertions(+), 18 deletions(-) diff --git a/packages/api/internal/sandbox/reservations/redis/reservation_test.go b/packages/api/internal/sandbox/reservations/redis/reservation_test.go index af2cc7ae6f..bb252f054b 100644 --- a/packages/api/internal/sandbox/reservations/redis/reservation_test.go +++ b/packages/api/internal/sandbox/reservations/redis/reservation_test.go @@ -28,9 +28,6 @@ const ( var testTeamID = uuid.New() -// setupTestReservationStorage wires the reservation store to a real storage -// pub/sub seam: one Redis container, one subscription manager, one publish -// worker pool. Mirrors the production wiring in orchestrator.go. func setupTestReservationStorage(t *testing.T) (*ReservationStorage, goredis.UniversalClient) { t.Helper() client := redis_utils.SetupInstance(t) diff --git a/packages/api/internal/sandbox/storage/redis/notifier.go b/packages/api/internal/sandbox/storage/redis/notifier.go index 68ec7756f6..c731d3c2d8 100644 --- a/packages/api/internal/sandbox/storage/redis/notifier.go +++ b/packages/api/internal/sandbox/storage/redis/notifier.go @@ -5,11 +5,6 @@ import ( ) // Notifier is the public seam onto the shared storage pub/sub infrastructure. -// -// Lifecycle is owned by Storage: callers MUST construct Notifier via -// Storage.Notifier(), and the underlying Storage must be Start()-ed -// before any Subscribe or Publish call. Close() on Storage shuts both -// sides down. type Notifier struct { sub *subscriptionManager pub *publisher @@ -24,16 +19,12 @@ func (n *Notifier) Subscribe(routingKey string) (<-chan struct{}, func()) { } // Publish enqueues routingKey for asynchronous PUBLISH on the shared -// notify channel. Never blocks: drops silently (with rate-limited warn) -// when the publish queue is saturated. Drop tolerance is part of the -// contract — every consumer ships with a fallback ticker. +// notify channel. Every consumer should use a fallback ticker. func (n *Notifier) Publish(ctx context.Context, routingKey string) { n.pub.Publish(ctx, routingKey) } -// Notifier returns the cross-package pub/sub seam. The returned value is -// cheap; callers may cache or re-fetch as convenient. Safe to call before -// Storage.Start, but Subscribe/Publish only function once Start is running. +// Notifier returns the cross-package pub/sub seam func (s *Storage) Notifier() *Notifier { return &Notifier{sub: s.subManager, pub: s.publisher} } diff --git a/packages/api/internal/sandbox/store.go b/packages/api/internal/sandbox/store.go index cb988090ec..d5a5aa633c 100644 --- a/packages/api/internal/sandbox/store.go +++ b/packages/api/internal/sandbox/store.go @@ -220,7 +220,3 @@ func (s *Store) Reserve(ctx context.Context, teamID uuid.UUID, sandboxID string, return finishStart, waitForStart, nil } - -func (s *Store) Release(ctx context.Context, teamID uuid.UUID, sandboxID string) error { - return s.reservations.Release(ctx, teamID, sandboxID) -} From 24d2caa560ad1885d4d72f9b6db9b876cff7bb5f Mon Sep 17 00:00:00 2001 From: Jakub Novak Date: Thu, 21 May 2026 10:22:10 +0000 Subject: [PATCH 06/11] chore: check Redis error --- .../api/internal/sandbox/reservations/redis/reservation.go | 3 +++ 1 file changed, 3 insertions(+) diff --git a/packages/api/internal/sandbox/reservations/redis/reservation.go b/packages/api/internal/sandbox/reservations/redis/reservation.go index 918d754710..83f1a729f5 100644 --- a/packages/api/internal/sandbox/reservations/redis/reservation.go +++ b/packages/api/internal/sandbox/reservations/redis/reservation.go @@ -238,6 +238,9 @@ func (s *ReservationStorage) tryReadResult( return true, sbx, err } + if !errors.Is(getErr, redis.Nil) { + return true, sandbox.Sandbox{}, fmt.Errorf("failed to check result key: %w", getErr) + } return true, sandbox.Sandbox{}, sandbox.ErrReservationReleased } From 8eceefbe23fef4b5a2dc5bf30cb6727b90a258e8 Mon Sep 17 00:00:00 2001 From: Jakub Novak Date: Thu, 21 May 2026 10:22:24 +0000 Subject: [PATCH 07/11] chore: add readme for reservation store --- .../sandbox/reservations/redis/README.md | 26 +++++++++++++++++++ 1 file changed, 26 insertions(+) create mode 100644 packages/api/internal/sandbox/reservations/redis/README.md diff --git a/packages/api/internal/sandbox/reservations/redis/README.md b/packages/api/internal/sandbox/reservations/redis/README.md new file mode 100644 index 0000000000..2193aa36cb --- /dev/null +++ b/packages/api/internal/sandbox/reservations/redis/README.md @@ -0,0 +1,26 @@ +# Redis Reservation Storage + +This package coordinates sandbox creation reservations across API instances. + +## Keys + +- Storage index: `sandbox:storage:{teamID}:index` +- Pending zset: `sandbox:storage:{teamID}:reservations:pending` +- Result key: `sandbox:storage:{teamID}:reservations:{sandboxID}:result` +- PubSub routing key: `sandbox:storage:{teamID}:reservations:{sandboxID}:notify` + +## Flow + +`Reserve` runs a Lua script that atomically removes stale pending entries, checks whether the sandbox already exists or has already pending start, enforces the team limit using `SCARD(storage index) + ZCARD(pending zset)`, deletes any stale result key, and adds the sandbox ID to the pending zset. + +When creation completes, it removes the sandbox from the pending zset, writes a TTL result key containing either the sandbox or the creation error, and publishes the routing key. + +A waiter subscribes to the routing key, probes the result key immediately, then waits for PubSub notifications or the 1 second fallback ticker. PubSub is best-effort; the fallback ticker is required for correctness. + +`Release` is called when the sandbox is removed from storage (`Store.Remove`). It removes the sandbox from the pending zset, writes a TTL tombstone result key, and publishes the routing key. The tombstone decodes to `sandbox.ErrReservationReleased`. + +## Migration Note + +The waiter still has a legacy compatibility path: if the result key is missing and the sandbox is no longer in the pending zset, it treats the reservation as released. This supports old instances that removed pending entries without writing tombstones. + +After all instances write tombstones on release, this waiter-side `ZSCORE` fallback can be removed. The pending zset itself must remain unless reservation and limit accounting are redesigned. From 0b051be96f841517ffb1992a1862e41f3d2f9e81 Mon Sep 17 00:00:00 2001 From: Jakub Novak Date: Thu, 21 May 2026 12:46:11 +0000 Subject: [PATCH 08/11] chore: improve comment --- .../sandbox/reservations/redis/reservation.go | 15 ++++++++++----- 1 file changed, 10 insertions(+), 5 deletions(-) diff --git a/packages/api/internal/sandbox/reservations/redis/reservation.go b/packages/api/internal/sandbox/reservations/redis/reservation.go index 83f1a729f5..f65f30f9d7 100644 --- a/packages/api/internal/sandbox/reservations/redis/reservation.go +++ b/packages/api/internal/sandbox/reservations/redis/reservation.go @@ -87,6 +87,9 @@ func (s *ReservationStorage) Release(ctx context.Context, teamID uuid.UUID, sand pendingSetKey := getPendingSetKey(teamIDStr) resultKeyStr := getResultKey(teamIDStr, sandboxID) + // Add a tombstone to prevent a race condition when a waiter doesn't check the pending set in time and we remove the pending entry + // but before we set the result key. The waiter treats a missing pending entry as a release, + // so without the tombstone it might miss the release and end up waiting indefinitely for a result that will never come. tombstone, err := encodeReleased() if err != nil { return fmt.Errorf("failed to encode release tombstone: %w", err) @@ -222,16 +225,18 @@ func (s *ReservationStorage) tryReadResult( return true, sandbox.Sandbox{}, fmt.Errorf("failed to check result key: %w", getErr) } - // No result yet. Confirm the reservation is still pending; if it's - // not, this is a pre-tombstone Release from a legacy instance. Treat - // it as a release for compatibility. + // No result yet. New Release calls write a tombstone result key, so a + // missing result normally means the reservation is still pending. This + // ZSCORE check is only for compatibility with legacy instances that + // released reservations by removing the pending entry without writing a + // tombstone; in that case, a missing pending entry means released. // // TODO [ENG-4089]: drop this ZSCORE fallback once all // instances are guaranteed to write a tombstone on Release. scoreErr := s.redisClient.ZScore(ctx, pendingSetKey, sandboxID).Err() if errors.Is(scoreErr, redis.Nil) { - // Final read in case finishStart/release raced between our GET - // and ZSCORE. + // Re-read the result in case finishStart or a new Release wrote it + // between the initial GET and the legacy pending-set check. data, getErr = s.redisClient.Get(ctx, resultKey).Bytes() if getErr == nil { sbx, err = decodeResult(data) From 379f9a6f0e3a79f77e796064c40048c615b8becf Mon Sep 17 00:00:00 2001 From: Jakub Novak Date: Thu, 21 May 2026 12:46:52 +0000 Subject: [PATCH 09/11] chore: clean up readme --- packages/api/internal/sandbox/reservations/redis/README.md | 6 ------ 1 file changed, 6 deletions(-) diff --git a/packages/api/internal/sandbox/reservations/redis/README.md b/packages/api/internal/sandbox/reservations/redis/README.md index 2193aa36cb..04f0f1282e 100644 --- a/packages/api/internal/sandbox/reservations/redis/README.md +++ b/packages/api/internal/sandbox/reservations/redis/README.md @@ -18,9 +18,3 @@ When creation completes, it removes the sandbox from the pending zset, writes a A waiter subscribes to the routing key, probes the result key immediately, then waits for PubSub notifications or the 1 second fallback ticker. PubSub is best-effort; the fallback ticker is required for correctness. `Release` is called when the sandbox is removed from storage (`Store.Remove`). It removes the sandbox from the pending zset, writes a TTL tombstone result key, and publishes the routing key. The tombstone decodes to `sandbox.ErrReservationReleased`. - -## Migration Note - -The waiter still has a legacy compatibility path: if the result key is missing and the sandbox is no longer in the pending zset, it treats the reservation as released. This supports old instances that removed pending entries without writing tombstones. - -After all instances write tombstones on release, this waiter-side `ZSCORE` fallback can be removed. The pending zset itself must remain unless reservation and limit accounting are redesigned. From abe93622b4ffac3e52507809f8c3cc022e44de22 Mon Sep 17 00:00:00 2001 From: Jakub Novak Date: Thu, 21 May 2026 12:50:52 +0000 Subject: [PATCH 10/11] chore: clean up comments --- .../redis/reservation_pubsub_test.go | 26 +------------------ 1 file changed, 1 insertion(+), 25 deletions(-) diff --git a/packages/api/internal/sandbox/reservations/redis/reservation_pubsub_test.go b/packages/api/internal/sandbox/reservations/redis/reservation_pubsub_test.go index ad58ac610f..4fd39617bc 100644 --- a/packages/api/internal/sandbox/reservations/redis/reservation_pubsub_test.go +++ b/packages/api/internal/sandbox/reservations/redis/reservation_pubsub_test.go @@ -32,20 +32,13 @@ func testSandbox(teamID uuid.UUID, sandboxID string) sandbox.Sandbox { } // setupReservationStorageWithoutSubManager wires the reservation store so that -// PubSub messages are never delivered in-process. The publisher worker still -// runs (so PUBLISH commands hit Redis) but no subscription manager is -// listening; in-process waiters only resolve via the fallback ticker. Used to -// exercise the safety-net path. +// PubSub messages are never delivered in-process. Used to exercise the safety-net path. func setupReservationStorageWithoutSubManager(t *testing.T) (*ReservationStorage, goredis.UniversalClient) { t.Helper() client := redis_utils.SetupInstance(t) storageInstance := newTestSandboxStorage(t, client) - // Deliberately do NOT call storageInstance.Start — no subManager.start - // goroutine, no fan-out. Publish() still works through the in-process - // queue but its drainer is also not running. The waiter must rely - // entirely on its 1s fallback ticker. t.Cleanup(func() { storageInstance.Close(context.WithoutCancel(t.Context())) }) storage := NewReservationStorage(client, storageInstance.Notifier()) @@ -53,10 +46,6 @@ func setupReservationStorageWithoutSubManager(t *testing.T) (*ReservationStorage return storage, client } -// TestWaitForStart_WokenByFinishStartPublish is the load-bearing regression -// test for the polling bug. With pub/sub working, the waiter must wake -// well under the 1s fallback ticker — anything that fast can only come -// from a PubSub delivery. func TestWaitForStart_WokenByFinishStartPublish(t *testing.T) { t.Parallel() @@ -98,8 +87,6 @@ func TestWaitForStart_WokenByFinishStartPublish(t *testing.T) { } } -// TestWaitForStart_WokenByReleasePublish proves Release also drives a fast -// wakeup and that the waiter returns the typed ErrReservationReleased. func TestWaitForStart_WokenByReleasePublish(t *testing.T) { t.Parallel() @@ -135,9 +122,6 @@ func TestWaitForStart_WokenByReleasePublish(t *testing.T) { } } -// TestWaitForStart_FallbackTickerWhenPubSubMissed disables the subscription -// manager entirely so no in-process wakeup can fire. The waiter must still -// resolve via the 1s fallback ticker. func TestWaitForStart_FallbackTickerWhenPubSubMissed(t *testing.T) { t.Parallel() @@ -335,10 +319,6 @@ func TestWaitForStart_FailedStartPropagatesPromptly(t *testing.T) { } } -// TestWaitForStart_RedisOpsBounded is the explicit regression test for the -// production incident: a 1s wait must NOT generate dozens of Redis ops. We -// install a redis.Hook on the client, run a waiter that resolves via -// PubSub, and assert the GET/ZSCORE count is single-digit. func TestWaitForStart_RedisOpsBounded(t *testing.T) { t.Parallel() @@ -362,10 +342,6 @@ func TestWaitForStart_RedisOpsBounded(t *testing.T) { _, waitForStart, err := storage.Reserve(t.Context(), teamID, sbxID, 10) require.NoError(t, err) - // Hold the producer for ~1.5s so the OLD 20ms-polling implementation - // would have done ~75 GETs + ~75 ZSCOREs = ~150 ops. The new pubsub - // implementation should do exactly one initial GET (+ one ZSCORE for - // the legacy-compat path) plus one final GET on wake. waiterDone := make(chan error, 1) go func() { _, err := waitForStart(t.Context()) From 5ce87e55052a8215d64aea091b63ca1d5133cb81 Mon Sep 17 00:00:00 2001 From: Jakub Novak Date: Fri, 22 May 2026 07:41:50 +0000 Subject: [PATCH 11/11] chore: keep tombstone out of pubsub PR --- packages/api/internal/sandbox/errors.go | 7 --- .../sandbox/reservations/redis/README.md | 2 +- .../sandbox/reservations/redis/reservation.go | 44 ++++--------------- .../redis/reservation_pubsub_test.go | 3 +- .../sandbox/reservations/redis/result.go | 19 +------- .../sandbox/reservations/redis/scripts.go | 10 +---- 6 files changed, 15 insertions(+), 70 deletions(-) diff --git a/packages/api/internal/sandbox/errors.go b/packages/api/internal/sandbox/errors.go index d26d6c8eda..76b49af5e8 100644 --- a/packages/api/internal/sandbox/errors.go +++ b/packages/api/internal/sandbox/errors.go @@ -40,10 +40,3 @@ var ErrAlreadyExists = errors.New("sandbox already exists") var ErrEvictionInProgress = errors.New("sandbox eviction already in progress") var ErrEvictionNotNeeded = errors.New("sandbox eviction not needed") - -// ErrReservationReleased is returned by ReservationStorage.Reserve's -// waitForStart callback when the producer released the reservation -// instead of completing the sandbox creation. It is the structural -// equivalent of "the other instance gave up": the caller can retry -// Reserve from scratch. -var ErrReservationReleased = errors.New("reservation released") diff --git a/packages/api/internal/sandbox/reservations/redis/README.md b/packages/api/internal/sandbox/reservations/redis/README.md index 04f0f1282e..67c48b0f72 100644 --- a/packages/api/internal/sandbox/reservations/redis/README.md +++ b/packages/api/internal/sandbox/reservations/redis/README.md @@ -17,4 +17,4 @@ When creation completes, it removes the sandbox from the pending zset, writes a A waiter subscribes to the routing key, probes the result key immediately, then waits for PubSub notifications or the 1 second fallback ticker. PubSub is best-effort; the fallback ticker is required for correctness. -`Release` is called when the sandbox is removed from storage (`Store.Remove`). It removes the sandbox from the pending zset, writes a TTL tombstone result key, and publishes the routing key. The tombstone decodes to `sandbox.ErrReservationReleased`. +`Release` is called when the sandbox is removed from storage (`Store.Remove`). It removes the sandbox from the pending zset, deletes the result key, and publishes the routing key. diff --git a/packages/api/internal/sandbox/reservations/redis/reservation.go b/packages/api/internal/sandbox/reservations/redis/reservation.go index f65f30f9d7..922057cb85 100644 --- a/packages/api/internal/sandbox/reservations/redis/reservation.go +++ b/packages/api/internal/sandbox/reservations/redis/reservation.go @@ -87,24 +87,15 @@ func (s *ReservationStorage) Release(ctx context.Context, teamID uuid.UUID, sand pendingSetKey := getPendingSetKey(teamIDStr) resultKeyStr := getResultKey(teamIDStr, sandboxID) - // Add a tombstone to prevent a race condition when a waiter doesn't check the pending set in time and we remove the pending entry - // but before we set the result key. The waiter treats a missing pending entry as a release, - // so without the tombstone it might miss the release and end up waiting indefinitely for a result that will never come. - tombstone, err := encodeReleased() - if err != nil { - return fmt.Errorf("failed to encode release tombstone: %w", err) - } - ttlSeconds := int(resultTTL.Seconds()) - - err = releaseScript.Run(ctx, s.redisClient, + err := releaseScript.Run(ctx, s.redisClient, []string{pendingSetKey, resultKeyStr}, - sandboxID, tombstone, ttlSeconds, + sandboxID, ).Err() if err != nil { return fmt.Errorf("failed to run release script: %w", err) } - // Wake any in-process waiter so it reads the tombstone immediately + // Wake any in-process waiter so it checks the pending set immediately. s.notifier.Publish(ctx, getReservationRoutingKey(teamIDStr, sandboxID)) return nil @@ -127,19 +118,10 @@ func (s *ReservationStorage) createFinishStart(ctx context.Context, teamID uuid. logger.WithSandboxID(sandboxID), ) - // Still try to remove from pending even if encoding fails - tombstone, tsErr := encodeReleased() - if tsErr != nil { - // If we can't encode the tombstone either, just delete the pending entry without a result key - _ = s.redisClient.ZRem(bgCtx, pendingSetKey, sandboxID).Err() - } else { - _ = releaseScript.Run(bgCtx, s.redisClient, - []string{pendingSetKey, resultKeyStr}, - sandboxID, tombstone, int(resultTTL.Seconds()), - ).Err() - } + // Still try to remove from pending even if encoding fails. + _ = s.redisClient.ZRem(bgCtx, pendingSetKey, sandboxID).Err() - // Wake waiters so they read the tombstone + // Wake waiters so they can observe that the reservation is gone. s.notifier.Publish(bgCtx, routingKey) return @@ -205,8 +187,7 @@ func (s *ReservationStorage) createWaitForStart(teamID uuid.UUID, sandboxID stri // // Returns done=true when the wait is over: // - the result key holds an encoded terminal result, or -// - the sandbox vanished from the pending set without a tombstone -// (legacy compatibility path), or +// - the sandbox vanished from the pending set without a result, or // - the Redis call itself failed. // // Returns done=false when the reservation is still pending and the caller @@ -225,14 +206,7 @@ func (s *ReservationStorage) tryReadResult( return true, sandbox.Sandbox{}, fmt.Errorf("failed to check result key: %w", getErr) } - // No result yet. New Release calls write a tombstone result key, so a - // missing result normally means the reservation is still pending. This - // ZSCORE check is only for compatibility with legacy instances that - // released reservations by removing the pending entry without writing a - // tombstone; in that case, a missing pending entry means released. - // - // TODO [ENG-4089]: drop this ZSCORE fallback once all - // instances are guaranteed to write a tombstone on Release. + // No result yet, so check whether another instance is still creating the sandbox. scoreErr := s.redisClient.ZScore(ctx, pendingSetKey, sandboxID).Err() if errors.Is(scoreErr, redis.Nil) { // Re-read the result in case finishStart or a new Release wrote it @@ -247,7 +221,7 @@ func (s *ReservationStorage) tryReadResult( return true, sandbox.Sandbox{}, fmt.Errorf("failed to check result key: %w", getErr) } - return true, sandbox.Sandbox{}, sandbox.ErrReservationReleased + return true, sandbox.Sandbox{}, fmt.Errorf("sandbox %s is no longer pending and has no result", sandboxID) } if scoreErr != nil { return true, sandbox.Sandbox{}, fmt.Errorf("failed to check pending set: %w", scoreErr) diff --git a/packages/api/internal/sandbox/reservations/redis/reservation_pubsub_test.go b/packages/api/internal/sandbox/reservations/redis/reservation_pubsub_test.go index 4fd39617bc..d93322b1b3 100644 --- a/packages/api/internal/sandbox/reservations/redis/reservation_pubsub_test.go +++ b/packages/api/internal/sandbox/reservations/redis/reservation_pubsub_test.go @@ -114,7 +114,8 @@ func TestWaitForStart_WokenByReleasePublish(t *testing.T) { select { case err := <-waiterErr: - require.ErrorIs(t, err, sandbox.ErrReservationReleased) + require.Error(t, err) + assert.Contains(t, err.Error(), "no longer pending") assert.Less(t, time.Since(start), 500*time.Millisecond, "release should wake the waiter via PubSub") case <-time.After(3 * time.Second): diff --git a/packages/api/internal/sandbox/reservations/redis/result.go b/packages/api/internal/sandbox/reservations/redis/result.go index aaff284e27..f7402085e6 100644 --- a/packages/api/internal/sandbox/reservations/redis/result.go +++ b/packages/api/internal/sandbox/reservations/redis/result.go @@ -16,14 +16,8 @@ type reservationResult struct { } // reservationError preserves api.APIError fields for cross-instance error propagation. -// -// A Released=true value is the "tombstone" written by Release: the producer -// withdrew the reservation without completing creation. Tombstones make the -// result key the single source of truth for the waiter, so the wait path -// never needs to inspect the pending zset. type reservationError struct { - Released bool `json:"released,omitempty"` - Message string `json:"message,omitempty"` + Message string `json:"message"` Code int `json:"code,omitempty"` ClientMsg string `json:"client_msg,omitempty"` } @@ -46,14 +40,6 @@ func encodeResult(sbx sandbox.Sandbox, err error) ([]byte, error) { return json.Marshal(result) } -// encodeReleased returns the tombstone result written by Release. It carries -// no sandbox payload and decodes to sandbox.ErrReservationReleased. -func encodeReleased() ([]byte, error) { - return json.Marshal(reservationResult{ - Error: &reservationError{Released: true}, - }) -} - func decodeResult(data []byte) (sandbox.Sandbox, error) { var result reservationResult if err := json.Unmarshal(data, &result); err != nil { @@ -71,9 +57,6 @@ func decodeResult(data []byte) (sandbox.Sandbox, error) { // If the error had an API code, it reconstructs an *api.APIError to preserve // errors.As(err, &apiErr) behavior in create_instance.go. func reconstructError(re *reservationError) error { - if re.Released { - return sandbox.ErrReservationReleased - } if re.Code != 0 { return &api.APIError{ Code: re.Code, diff --git a/packages/api/internal/sandbox/reservations/redis/scripts.go b/packages/api/internal/sandbox/reservations/redis/scripts.go index be384ebedf..29d187ef92 100644 --- a/packages/api/internal/sandbox/reservations/redis/scripts.go +++ b/packages/api/internal/sandbox/reservations/redis/scripts.go @@ -75,19 +75,13 @@ var ( return 1 `) - // releaseScript removes a sandbox from the pending zset and writes a - // tombstone result so any concurrent waiter resolves to - // sandbox.ErrReservationReleased via a single GET on the result key, - // without ever consulting the pending zset. - // + // releaseScript removes a sandbox from the pending zset and deletes the result key. // KEYS[1] = pending zset key // KEYS[2] = result key // ARGV[1] = sandboxID - // ARGV[2] = tombstone JSON - // ARGV[3] = TTL in seconds releaseScript = redis.NewScript(` redis.call('ZREM', KEYS[1], ARGV[1]) - redis.call('SET', KEYS[2], ARGV[2], 'EX', tonumber(ARGV[3])) + redis.call('DEL', KEYS[2]) return 1 `) )