From da9620444e32caee27d9a5660493560df7ccf5fe Mon Sep 17 00:00:00 2001 From: Jakub Novak Date: Sun, 17 May 2026 17:05:57 +0000 Subject: [PATCH 01/12] feat(metrics): distinguish joined from regular requests (ENG-4072) Add a request.joined attribute to http.server.duration and the active trace span so requests that intentionally block on another in-flight operation (concurrent CreateSandbox joiners, sandbox state-transition waiters) can be filtered out of normal latency/error dashboards. The marker rides on a request-scoped omitHolder attached to context.Context by the metrics middleware, so the orchestrator and storage layers can flag themselves without depending on *gin.Context. --- .../internal/orchestrator/create_instance.go | 4 + .../sandbox/storage/redis/state_change.go | 4 + .../pkg/middleware/otel/metrics/middleware.go | 10 +- .../pkg/middleware/otel/metrics/omit.go | 51 +++++ .../pkg/middleware/otel/metrics/omit_test.go | 204 ++++++++++++++++++ 5 files changed, 272 insertions(+), 1 deletion(-) create mode 100644 packages/shared/pkg/middleware/otel/metrics/omit.go create mode 100644 packages/shared/pkg/middleware/otel/metrics/omit_test.go diff --git a/packages/api/internal/orchestrator/create_instance.go b/packages/api/internal/orchestrator/create_instance.go index 46855aa912..22656d01cf 100644 --- a/packages/api/internal/orchestrator/create_instance.go +++ b/packages/api/internal/orchestrator/create_instance.go @@ -26,6 +26,7 @@ import ( "github.com/e2b-dev/infra/packages/shared/pkg/featureflags" "github.com/e2b-dev/infra/packages/shared/pkg/grpc/orchestrator" "github.com/e2b-dev/infra/packages/shared/pkg/logger" + "github.com/e2b-dev/infra/packages/shared/pkg/middleware/otel/metrics" sandbox_network "github.com/e2b-dev/infra/packages/shared/pkg/sandbox-network" "github.com/e2b-dev/infra/packages/shared/pkg/telemetry" ut "github.com/e2b-dev/infra/packages/shared/pkg/utils" @@ -159,6 +160,9 @@ func (o *Orchestrator) CreateSandbox( } if waitForStart != nil { + // Mark as a joined request for telemetry purposes + metrics.MarkJoined(ctx) + logger.L().Info(ctx, "sandbox is already being started, waiting for it to be ready", logger.WithSandboxID(sandboxID)) sbx, err = waitForStart(ctx) diff --git a/packages/api/internal/sandbox/storage/redis/state_change.go b/packages/api/internal/sandbox/storage/redis/state_change.go index 483642d6f8..4cb220252e 100644 --- a/packages/api/internal/sandbox/storage/redis/state_change.go +++ b/packages/api/internal/sandbox/storage/redis/state_change.go @@ -14,6 +14,7 @@ import ( "github.com/e2b-dev/infra/packages/api/internal/sandbox" "github.com/e2b-dev/infra/packages/shared/pkg/logger" + "github.com/e2b-dev/infra/packages/shared/pkg/middleware/otel/metrics" redis_utils "github.com/e2b-dev/infra/packages/shared/pkg/redis" ) @@ -249,6 +250,9 @@ func (s *Storage) waitForTransition( sandboxID, transitionID string, ) error { + // Mark as a joined request for telemetry purposes + metrics.MarkJoined(ctx) + routingKey := getTransitionRoutingKey(teamID.String(), sandboxID, transitionID) transitionKey := getTransitionKey(teamID.String(), sandboxID) resultKey := getTransitionResultKey(teamID.String(), sandboxID, transitionID) diff --git a/packages/shared/pkg/middleware/otel/metrics/middleware.go b/packages/shared/pkg/middleware/otel/metrics/middleware.go index 8d02e86363..e9134d928e 100644 --- a/packages/shared/pkg/middleware/otel/metrics/middleware.go +++ b/packages/shared/pkg/middleware/otel/metrics/middleware.go @@ -57,6 +57,12 @@ func Middleware(meterProvider metric.MeterProvider, service string, options ...O return func(ginCtx *gin.Context) { ctx := ginCtx.Request.Context() + // Install the request-scoped omitHolder so any descendant code path + // (orchestrator, storage layer, etc.) can mark the request via the + // package helpers without needing access to *gin.Context. + ctx, holder := WithOmitHolder(ctx) + ginCtx.Request = ginCtx.Request.WithContext(ctx) + route := ginCtx.FullPath() if len(route) == 0 { route = "nonconfigured" @@ -90,7 +96,9 @@ func Middleware(meterProvider metric.MeterProvider, service string, options ...O // Append attributes from ginCtx resAttributes = append(resAttributes, attributesFromGinContext(ginCtx, MetricPrefix)...) - // Use processing start time if set, otherwise fall back to the middleware start time. + // Distinguish between regular and joined requests + resAttributes = append(resAttributes, holder.joinedAttribute()) + effectiveStart := start if processingStart, ok := getProcessingStartTime(ginCtx); ok { effectiveStart = processingStart diff --git a/packages/shared/pkg/middleware/otel/metrics/omit.go b/packages/shared/pkg/middleware/otel/metrics/omit.go new file mode 100644 index 0000000000..01c97dafe5 --- /dev/null +++ b/packages/shared/pkg/middleware/otel/metrics/omit.go @@ -0,0 +1,51 @@ +package metrics + +import ( + "context" + "sync/atomic" + + "go.opentelemetry.io/otel/attribute" + "go.opentelemetry.io/otel/trace" +) + +const requestJoinedAttrKey = "request.joined" + +type omitHolder struct { + joined atomic.Bool +} + +type omitHolderKey struct{} + +// WithOmitHolder installs a fresh omitHolder on ctx and returns the augmented +// context plus a handle to the holder. The metrics middleware calls this +// once per request; non-HTTP callers normally never call it, which leaves +// the package helpers as safe no-ops. +func WithOmitHolder(ctx context.Context) (context.Context, *omitHolder) { + h := &omitHolder{} + return context.WithValue(ctx, omitHolderKey{}, h), h +} + +// MarkJoined marks the current request as having joined an in-flight +// concurrent operation (e.g. waiting for another request to finish a sandbox +// state transition, or joining a concurrent CreateSandbox). The flag is +// emitted as a histogram attribute on http.server.duration and as a span +// attribute on the active trace span (first-write-wins). +// +// Safe to call from any goroutine descended from the request context. +// No-op if ctx has no omitHolder (e.g. non-HTTP callers, tests). +func MarkJoined(ctx context.Context) { + h, ok := ctx.Value(omitHolderKey{}).(*omitHolder) + if !ok { + return + } + + if h.joined.CompareAndSwap(false, true) { + trace.SpanFromContext(ctx).SetAttributes( + attribute.String(requestJoinedAttrKey, "true"), + ) + } +} + +func (h *omitHolder) joinedAttribute() attribute.KeyValue { + return attribute.Bool(requestJoinedAttrKey, h.joined.Load()) +} diff --git a/packages/shared/pkg/middleware/otel/metrics/omit_test.go b/packages/shared/pkg/middleware/otel/metrics/omit_test.go new file mode 100644 index 0000000000..1007868ab7 --- /dev/null +++ b/packages/shared/pkg/middleware/otel/metrics/omit_test.go @@ -0,0 +1,204 @@ +package metrics + +import ( + "context" + "net/http" + "net/http/httptest" + "sync" + "testing" + "time" + + "github.com/gin-gonic/gin" + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" + "go.opentelemetry.io/otel/attribute" +) + +// fakeRecorder captures every ObserveHTTPRequestDuration call so tests can +// inspect the attributes the middleware would have emitted. +type fakeRecorder struct { + mu sync.Mutex + calls []fakeRecorderCall +} + +type fakeRecorderCall struct { + duration time.Duration + attrs []attribute.KeyValue +} + +func (f *fakeRecorder) ObserveHTTPRequestDuration(_ context.Context, duration time.Duration, attrs []attribute.KeyValue) { + f.mu.Lock() + defer f.mu.Unlock() + // Copy attrs because the middleware reuses its slice buffer between + // requests when the test engine handles more than one request. + cp := make([]attribute.KeyValue, len(attrs)) + copy(cp, attrs) + f.calls = append(f.calls, fakeRecorderCall{duration: duration, attrs: cp}) +} + +func (f *fakeRecorder) callCount() int { + f.mu.Lock() + defer f.mu.Unlock() + + return len(f.calls) +} + +func (f *fakeRecorder) attrs(i int) []attribute.KeyValue { + f.mu.Lock() + defer f.mu.Unlock() + + return f.calls[i].attrs +} + +func attrValue(attrs []attribute.KeyValue, key string) (attribute.Value, bool) { + for _, a := range attrs { + if string(a.Key) == key { + return a.Value, true + } + } + + return attribute.Value{}, false +} + +func newTestEngine(t *testing.T, handler gin.HandlerFunc) (*gin.Engine, *fakeRecorder) { + t.Helper() + gin.SetMode(gin.TestMode) + + rec := &fakeRecorder{} + r := gin.New() + r.Use(Middleware(nil, "test", WithRecorder(rec))) + r.POST("/sandboxes/:id/resume", handler) + + return r, rec +} + +func doRequest(t *testing.T, r *gin.Engine) { + t.Helper() + req := httptest.NewRequest(http.MethodPost, "/sandboxes/abc/resume", nil) + w := httptest.NewRecorder() + r.ServeHTTP(w, req) + require.Equal(t, http.StatusOK, w.Code) +} + +// MarkJoined must be safe even when the context carries no omitHolder +// (e.g. non-HTTP callers, tests). +func TestMarkJoined_NoHolder_Noop(t *testing.T) { + MarkJoined(context.Background()) +} + +// Untagged requests must carry request.joined=false. +func TestMiddleware_NormalRequest_JoinedAttrIsFalse(t *testing.T) { + r, rec := newTestEngine(t, func(c *gin.Context) { c.Status(http.StatusOK) }) + + doRequest(t, r) + + require.Equal(t, 1, rec.callCount()) + v, ok := attrValue(rec.attrs(0), "request.joined") + require.True(t, ok, "request.joined must be present on every observation") + assert.False(t, v.AsBool(), "untagged request must carry request.joined=false") +} + +// MarkJoined from the handler's ctx must flip the attribute to true. +func TestMiddleware_MarkJoinedFromHandler_AppearsAsTrue(t *testing.T) { + r, rec := newTestEngine(t, func(c *gin.Context) { + MarkJoined(c.Request.Context()) + c.Status(http.StatusOK) + }) + + doRequest(t, r) + + require.Equal(t, 1, rec.callCount()) + v, ok := attrValue(rec.attrs(0), "request.joined") + require.True(t, ok, "request.joined attribute must be present on the histogram") + assert.True(t, v.AsBool()) +} + +// MarkJoined called from a goroutine descended from the request context +// must still flow through to the histogram. This is the key capability of +// the context-attached holder design vs. a *gin.Context-based marker. +func TestMiddleware_MarkJoinedFromDescendantGoroutine(t *testing.T) { + done := make(chan struct{}) + r, rec := newTestEngine(t, func(c *gin.Context) { + ctx := c.Request.Context() + go func() { + MarkJoined(ctx) + close(done) + }() + <-done + c.Status(http.StatusOK) + }) + + doRequest(t, r) + + require.Equal(t, 1, rec.callCount()) + v, ok := attrValue(rec.attrs(0), "request.joined") + require.True(t, ok) + assert.True(t, v.AsBool()) +} + +// MarkJoined is idempotent: repeated calls within the same request do not +// produce duplicate histogram attributes. +func TestMiddleware_MarkJoinedIdempotent(t *testing.T) { + r, rec := newTestEngine(t, func(c *gin.Context) { + MarkJoined(c.Request.Context()) + MarkJoined(c.Request.Context()) + MarkJoined(c.Request.Context()) + c.Status(http.StatusOK) + }) + + doRequest(t, r) + + require.Equal(t, 1, rec.callCount()) + attrs := rec.attrs(0) + + count := 0 + for _, a := range attrs { + if string(a.Key) == "request.joined" { + count++ + } + } + assert.Equal(t, 1, count, "request.joined must appear exactly once even after repeated MarkJoined calls") +} + +// Tagging must not suppress recording — we only add a label. +func TestMiddleware_Tagging_DoesNotSuppressRecording(t *testing.T) { + r, rec := newTestEngine(t, func(c *gin.Context) { + MarkJoined(c.Request.Context()) + c.Status(http.StatusOK) + }) + + doRequest(t, r) + + assert.Equal(t, 1, rec.callCount(), "histogram must still be recorded; tagging only adds attributes") +} + +// Two distinct requests must not share the holder: tagging one must not +// taint the other. +func TestMiddleware_HolderIsRequestScoped(t *testing.T) { + gin.SetMode(gin.TestMode) + rec := &fakeRecorder{} + r := gin.New() + r.Use(Middleware(nil, "test", WithRecorder(rec))) + r.POST("/joiner", func(c *gin.Context) { + MarkJoined(c.Request.Context()) + c.Status(http.StatusOK) + }) + r.POST("/normal", func(c *gin.Context) { c.Status(http.StatusOK) }) + + for _, path := range []string{"/joiner", "/normal"} { + req := httptest.NewRequest(http.MethodPost, path, nil) + w := httptest.NewRecorder() + r.ServeHTTP(w, req) + require.Equal(t, http.StatusOK, w.Code) + } + + require.Equal(t, 2, rec.callCount()) + + v1, ok := attrValue(rec.attrs(0), "request.joined") + require.True(t, ok) + assert.True(t, v1.AsBool(), "first (tagged) request must carry request.joined=true") + + v2, ok := attrValue(rec.attrs(1), "request.joined") + require.True(t, ok) + assert.False(t, v2.AsBool(), "second (untagged) request must carry request.joined=false") +} From c80ba1c973a92a45a42aae06b8704c001e1613d1 Mon Sep 17 00:00:00 2001 From: Jakub Novak Date: Sun, 17 May 2026 17:39:44 +0000 Subject: [PATCH 02/12] refactor(metrics): rename omit* to joined* and pin attribute to server span - Rename omit.go -> joined.go, omitHolder -> joinedHolder, withOmitHolder -> withJoinedHolder, omitHolderKey -> joinedHolderKey. The 'joined' naming matches the public MarkJoined API and the request.joined attribute. - Capture the HTTP server span in joinedHolder at middleware entry so MarkJoined writes the request.joined attribute onto the top-level request span instead of whatever child span (e.g. "create-sandbox") happens to be active when the helper is called. - Add a test that proves the attribute lands on the server span, not the child span, when MarkJoined is invoked from inside a child span. --- .../pkg/middleware/otel/metrics/joined.go | 63 +++++++++++ .../metrics/{omit_test.go => joined_test.go} | 104 ++++++++++++++++-- .../pkg/middleware/otel/metrics/middleware.go | 8 +- .../pkg/middleware/otel/metrics/omit.go | 51 --------- 4 files changed, 160 insertions(+), 66 deletions(-) create mode 100644 packages/shared/pkg/middleware/otel/metrics/joined.go rename packages/shared/pkg/middleware/otel/metrics/{omit_test.go => joined_test.go} (62%) delete mode 100644 packages/shared/pkg/middleware/otel/metrics/omit.go diff --git a/packages/shared/pkg/middleware/otel/metrics/joined.go b/packages/shared/pkg/middleware/otel/metrics/joined.go new file mode 100644 index 0000000000..a79d86b7f3 --- /dev/null +++ b/packages/shared/pkg/middleware/otel/metrics/joined.go @@ -0,0 +1,63 @@ +package metrics + +import ( + "context" + "sync/atomic" + + "go.opentelemetry.io/otel/attribute" + "go.opentelemetry.io/otel/trace" +) + +const requestJoinedAttrKey = "request.joined" + +type joinedHolder struct { + joined atomic.Bool + // serverSpan is captured at middleware entry so MarkJoined can pin the + // request.joined attribute onto the top-level HTTP server span instead + // of whatever child span (e.g. "create-sandbox") happens to be active + // when the helper is called. Tracing middleware is registered before + // the metrics middleware in every service that mounts both, so the ctx + // passed to withJoinedHolder carries the server span as the active span. + // If no real span is on the ctx, this is a no-op span and SetAttributes + // silently does nothing. + serverSpan trace.Span +} + +type joinedHolderKey struct{} + +// withJoinedHolder installs a fresh joinedHolder on ctx and returns the +// augmented context plus a handle to the holder. The metrics middleware +// calls this once per request; non-HTTP callers normally never call it, +// which leaves the package helpers as safe no-ops. +func withJoinedHolder(ctx context.Context) (context.Context, *joinedHolder) { + h := &joinedHolder{ + serverSpan: trace.SpanFromContext(ctx), + } + + return context.WithValue(ctx, joinedHolderKey{}, h), h +} + +// MarkJoined marks the current request as having joined an in-flight +// concurrent operation (e.g. waiting for another request to finish a sandbox +// state transition, or joining a concurrent CreateSandbox). The flag is +// emitted as a histogram attribute on http.server.duration and as an +// attribute on the top-level HTTP server span (first-write-wins). +// +// Safe to call from any goroutine descended from the request context. +// No-op if ctx has no joinedHolder (e.g. non-HTTP callers, tests). +func MarkJoined(ctx context.Context) { + h, ok := ctx.Value(joinedHolderKey{}).(*joinedHolder) + if !ok { + return + } + + if h.joined.CompareAndSwap(false, true) { + h.serverSpan.SetAttributes( + attribute.String(requestJoinedAttrKey, "true"), + ) + } +} + +func (h *joinedHolder) joinedAttribute() attribute.KeyValue { + return attribute.Bool(requestJoinedAttrKey, h.joined.Load()) +} diff --git a/packages/shared/pkg/middleware/otel/metrics/omit_test.go b/packages/shared/pkg/middleware/otel/metrics/joined_test.go similarity index 62% rename from packages/shared/pkg/middleware/otel/metrics/omit_test.go rename to packages/shared/pkg/middleware/otel/metrics/joined_test.go index 1007868ab7..04f6114bc9 100644 --- a/packages/shared/pkg/middleware/otel/metrics/omit_test.go +++ b/packages/shared/pkg/middleware/otel/metrics/joined_test.go @@ -12,6 +12,8 @@ import ( "github.com/stretchr/testify/assert" "github.com/stretchr/testify/require" "go.opentelemetry.io/otel/attribute" + sdktrace "go.opentelemetry.io/otel/sdk/trace" + "go.opentelemetry.io/otel/sdk/trace/tracetest" ) // fakeRecorder captures every ObserveHTTPRequestDuration call so tests can @@ -50,9 +52,11 @@ func (f *fakeRecorder) attrs(i int) []attribute.KeyValue { return f.calls[i].attrs } -func attrValue(attrs []attribute.KeyValue, key string) (attribute.Value, bool) { +// findRequestJoined returns the value of the `request.joined` attribute from +// the given attribute slice, if present. +func findRequestJoined(attrs []attribute.KeyValue) (attribute.Value, bool) { for _, a := range attrs { - if string(a.Key) == key { + if string(a.Key) == requestJoinedAttrKey { return a.Value, true } } @@ -62,7 +66,6 @@ func attrValue(attrs []attribute.KeyValue, key string) (attribute.Value, bool) { func newTestEngine(t *testing.T, handler gin.HandlerFunc) (*gin.Engine, *fakeRecorder) { t.Helper() - gin.SetMode(gin.TestMode) rec := &fakeRecorder{} r := gin.New() @@ -80,26 +83,31 @@ func doRequest(t *testing.T, r *gin.Engine) { require.Equal(t, http.StatusOK, w.Code) } -// MarkJoined must be safe even when the context carries no omitHolder +// MarkJoined must be safe even when the context carries no joinedHolder // (e.g. non-HTTP callers, tests). func TestMarkJoined_NoHolder_Noop(t *testing.T) { + t.Parallel() MarkJoined(context.Background()) } // Untagged requests must carry request.joined=false. func TestMiddleware_NormalRequest_JoinedAttrIsFalse(t *testing.T) { + t.Parallel() + r, rec := newTestEngine(t, func(c *gin.Context) { c.Status(http.StatusOK) }) doRequest(t, r) require.Equal(t, 1, rec.callCount()) - v, ok := attrValue(rec.attrs(0), "request.joined") + v, ok := findRequestJoined(rec.attrs(0)) require.True(t, ok, "request.joined must be present on every observation") assert.False(t, v.AsBool(), "untagged request must carry request.joined=false") } // MarkJoined from the handler's ctx must flip the attribute to true. func TestMiddleware_MarkJoinedFromHandler_AppearsAsTrue(t *testing.T) { + t.Parallel() + r, rec := newTestEngine(t, func(c *gin.Context) { MarkJoined(c.Request.Context()) c.Status(http.StatusOK) @@ -108,7 +116,7 @@ func TestMiddleware_MarkJoinedFromHandler_AppearsAsTrue(t *testing.T) { doRequest(t, r) require.Equal(t, 1, rec.callCount()) - v, ok := attrValue(rec.attrs(0), "request.joined") + v, ok := findRequestJoined(rec.attrs(0)) require.True(t, ok, "request.joined attribute must be present on the histogram") assert.True(t, v.AsBool()) } @@ -117,6 +125,8 @@ func TestMiddleware_MarkJoinedFromHandler_AppearsAsTrue(t *testing.T) { // must still flow through to the histogram. This is the key capability of // the context-attached holder design vs. a *gin.Context-based marker. func TestMiddleware_MarkJoinedFromDescendantGoroutine(t *testing.T) { + t.Parallel() + done := make(chan struct{}) r, rec := newTestEngine(t, func(c *gin.Context) { ctx := c.Request.Context() @@ -131,7 +141,7 @@ func TestMiddleware_MarkJoinedFromDescendantGoroutine(t *testing.T) { doRequest(t, r) require.Equal(t, 1, rec.callCount()) - v, ok := attrValue(rec.attrs(0), "request.joined") + v, ok := findRequestJoined(rec.attrs(0)) require.True(t, ok) assert.True(t, v.AsBool()) } @@ -139,6 +149,8 @@ func TestMiddleware_MarkJoinedFromDescendantGoroutine(t *testing.T) { // MarkJoined is idempotent: repeated calls within the same request do not // produce duplicate histogram attributes. func TestMiddleware_MarkJoinedIdempotent(t *testing.T) { + t.Parallel() + r, rec := newTestEngine(t, func(c *gin.Context) { MarkJoined(c.Request.Context()) MarkJoined(c.Request.Context()) @@ -153,7 +165,7 @@ func TestMiddleware_MarkJoinedIdempotent(t *testing.T) { count := 0 for _, a := range attrs { - if string(a.Key) == "request.joined" { + if string(a.Key) == requestJoinedAttrKey { count++ } } @@ -162,6 +174,8 @@ func TestMiddleware_MarkJoinedIdempotent(t *testing.T) { // Tagging must not suppress recording — we only add a label. func TestMiddleware_Tagging_DoesNotSuppressRecording(t *testing.T) { + t.Parallel() + r, rec := newTestEngine(t, func(c *gin.Context) { MarkJoined(c.Request.Context()) c.Status(http.StatusOK) @@ -172,10 +186,78 @@ func TestMiddleware_Tagging_DoesNotSuppressRecording(t *testing.T) { assert.Equal(t, 1, rec.callCount(), "histogram must still be recorded; tagging only adds attributes") } +// MarkJoined must pin the request.joined attribute onto the top-level HTTP +// server span (captured at middleware entry), not onto whatever child span +// happens to be active when the helper is called. This is the guarantee +// callers rely on for Tempo filtering by root-span attribute. +func TestMiddleware_MarkJoined_TagsServerSpanNotChildSpan(t *testing.T) { + t.Parallel() + + sr := tracetest.NewSpanRecorder() + tp := sdktrace.NewTracerProvider(sdktrace.WithSpanProcessor(sr)) + tracer := tp.Tracer("test") + + rec := &fakeRecorder{} + r := gin.New() + // Stand-in for the tracing middleware: open the server span before the + // metrics middleware installs its joinedHolder. This mirrors the real + // service wiring (tracing registered before metrics). + r.Use(func(c *gin.Context) { + ctx, span := tracer.Start(c.Request.Context(), "HTTP POST /sandboxes/:id/resume") + defer span.End() + c.Request = c.Request.WithContext(ctx) + c.Next() + }) + r.Use(Middleware(nil, "test", WithRecorder(rec))) + r.POST("/sandboxes/:id/resume", func(c *gin.Context) { + // Open a child span (mirrors orchestrator.CreateSandbox's + // "create-sandbox" child span) and call MarkJoined from inside it. + ctx, child := tracer.Start(c.Request.Context(), "create-sandbox") + MarkJoined(ctx) + child.End() + + c.Status(http.StatusOK) + }) + + req := httptest.NewRequest(http.MethodPost, "/sandboxes/abc/resume", nil) + w := httptest.NewRecorder() + r.ServeHTTP(w, req) + require.Equal(t, http.StatusOK, w.Code) + + spans := sr.Ended() + require.Len(t, spans, 2) + + // Identify server vs child by span name. + var serverSpan, childSpan sdktrace.ReadOnlySpan + for _, s := range spans { + if s.Name() == "create-sandbox" { + childSpan = s + } else { + serverSpan = s + } + } + require.NotNil(t, serverSpan, "server span must be recorded") + require.NotNil(t, childSpan, "child span must be recorded") + + serverAttr, hasServer := findRequestJoined(serverSpan.Attributes()) + require.True(t, hasServer, "request.joined must be on the server span") + assert.Equal(t, "true", serverAttr.AsString()) + + _, hasChild := findRequestJoined(childSpan.Attributes()) + assert.False(t, hasChild, "request.joined must NOT be on the child span") + + // And the histogram must still carry request.joined=true. + require.Equal(t, 1, rec.callCount()) + histAttr, ok := findRequestJoined(rec.attrs(0)) + require.True(t, ok) + assert.True(t, histAttr.AsBool()) +} + // Two distinct requests must not share the holder: tagging one must not // taint the other. func TestMiddleware_HolderIsRequestScoped(t *testing.T) { - gin.SetMode(gin.TestMode) + t.Parallel() + rec := &fakeRecorder{} r := gin.New() r.Use(Middleware(nil, "test", WithRecorder(rec))) @@ -194,11 +276,11 @@ func TestMiddleware_HolderIsRequestScoped(t *testing.T) { require.Equal(t, 2, rec.callCount()) - v1, ok := attrValue(rec.attrs(0), "request.joined") + v1, ok := findRequestJoined(rec.attrs(0)) require.True(t, ok) assert.True(t, v1.AsBool(), "first (tagged) request must carry request.joined=true") - v2, ok := attrValue(rec.attrs(1), "request.joined") + v2, ok := findRequestJoined(rec.attrs(1)) require.True(t, ok) assert.False(t, v2.AsBool(), "second (untagged) request must carry request.joined=false") } diff --git a/packages/shared/pkg/middleware/otel/metrics/middleware.go b/packages/shared/pkg/middleware/otel/metrics/middleware.go index e9134d928e..6fec11b828 100644 --- a/packages/shared/pkg/middleware/otel/metrics/middleware.go +++ b/packages/shared/pkg/middleware/otel/metrics/middleware.go @@ -57,10 +57,10 @@ func Middleware(meterProvider metric.MeterProvider, service string, options ...O return func(ginCtx *gin.Context) { ctx := ginCtx.Request.Context() - // Install the request-scoped omitHolder so any descendant code path - // (orchestrator, storage layer, etc.) can mark the request via the - // package helpers without needing access to *gin.Context. - ctx, holder := WithOmitHolder(ctx) + // Install the request-scoped joinedHolder so any descendant code + // path (orchestrator, storage layer, etc.) can mark the request via + // the package helpers without needing access to *gin.Context. + ctx, holder := withJoinedHolder(ctx) ginCtx.Request = ginCtx.Request.WithContext(ctx) route := ginCtx.FullPath() diff --git a/packages/shared/pkg/middleware/otel/metrics/omit.go b/packages/shared/pkg/middleware/otel/metrics/omit.go deleted file mode 100644 index 01c97dafe5..0000000000 --- a/packages/shared/pkg/middleware/otel/metrics/omit.go +++ /dev/null @@ -1,51 +0,0 @@ -package metrics - -import ( - "context" - "sync/atomic" - - "go.opentelemetry.io/otel/attribute" - "go.opentelemetry.io/otel/trace" -) - -const requestJoinedAttrKey = "request.joined" - -type omitHolder struct { - joined atomic.Bool -} - -type omitHolderKey struct{} - -// WithOmitHolder installs a fresh omitHolder on ctx and returns the augmented -// context plus a handle to the holder. The metrics middleware calls this -// once per request; non-HTTP callers normally never call it, which leaves -// the package helpers as safe no-ops. -func WithOmitHolder(ctx context.Context) (context.Context, *omitHolder) { - h := &omitHolder{} - return context.WithValue(ctx, omitHolderKey{}, h), h -} - -// MarkJoined marks the current request as having joined an in-flight -// concurrent operation (e.g. waiting for another request to finish a sandbox -// state transition, or joining a concurrent CreateSandbox). The flag is -// emitted as a histogram attribute on http.server.duration and as a span -// attribute on the active trace span (first-write-wins). -// -// Safe to call from any goroutine descended from the request context. -// No-op if ctx has no omitHolder (e.g. non-HTTP callers, tests). -func MarkJoined(ctx context.Context) { - h, ok := ctx.Value(omitHolderKey{}).(*omitHolder) - if !ok { - return - } - - if h.joined.CompareAndSwap(false, true) { - trace.SpanFromContext(ctx).SetAttributes( - attribute.String(requestJoinedAttrKey, "true"), - ) - } -} - -func (h *omitHolder) joinedAttribute() attribute.KeyValue { - return attribute.Bool(requestJoinedAttrKey, h.joined.Load()) -} From 353173aee7b036705c6d0a01ae08f88fe39f5fcc Mon Sep 17 00:00:00 2001 From: "github-actions[bot]" Date: Sun, 17 May 2026 17:41:17 +0000 Subject: [PATCH 03/12] chore: auto-commit generated changes --- packages/shared/pkg/middleware/otel/metrics/joined_test.go | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/packages/shared/pkg/middleware/otel/metrics/joined_test.go b/packages/shared/pkg/middleware/otel/metrics/joined_test.go index 04f6114bc9..3629a49d88 100644 --- a/packages/shared/pkg/middleware/otel/metrics/joined_test.go +++ b/packages/shared/pkg/middleware/otel/metrics/joined_test.go @@ -195,7 +195,7 @@ func TestMiddleware_MarkJoined_TagsServerSpanNotChildSpan(t *testing.T) { sr := tracetest.NewSpanRecorder() tp := sdktrace.NewTracerProvider(sdktrace.WithSpanProcessor(sr)) - tracer := tp.Tracer("test") + tracer := tp.Tracer("github.com/e2b-dev/infra/packages/shared/pkg/middleware/otel/metrics") rec := &fakeRecorder{} r := gin.New() From 04e368bd18d5874c8516f05fb7dfb07cec4e9e5a Mon Sep 17 00:00:00 2001 From: Jakub Novak Date: Sun, 17 May 2026 20:00:11 +0000 Subject: [PATCH 04/12] refactor(joined): extract to shared otel package, install in tracing too Move the request.joined marker out of the metrics middleware into a standalone shared package (packages/shared/pkg/middleware/otel/joined). Both the tracing and the metrics middleware now install the request-scoped holder via joined.WithHolder, so the marker works regardless of which middleware is mounted on a given route. The server span captured by the holder is now guaranteed to be the top-level HTTP server span: the tracing middleware installs the holder right after tracer.Start, and the metrics middleware's later WithHolder call is a no-op when the holder is already present. - New package joined with WithHolder, Mark, and Attribute. - Tracing middleware installs the holder right after starting the server span. - Metrics middleware reuses the existing holder and emits joined.Attribute(ctx) on every http.server.duration observation. - Callers (orchestrator joiner branch, state-transition waiter) switched to joined.Mark(ctx). --- .../internal/orchestrator/create_instance.go | 4 +- .../sandbox/storage/redis/state_change.go | 4 +- .../pkg/middleware/otel/joined/joined.go | 86 ++++++ .../pkg/middleware/otel/joined/joined_test.go | 137 +++++++++ .../pkg/middleware/otel/metrics/joined.go | 63 ---- .../middleware/otel/metrics/joined_test.go | 286 ------------------ .../pkg/middleware/otel/metrics/middleware.go | 12 +- .../pkg/middleware/otel/tracing/middleware.go | 8 + 8 files changed, 242 insertions(+), 358 deletions(-) create mode 100644 packages/shared/pkg/middleware/otel/joined/joined.go create mode 100644 packages/shared/pkg/middleware/otel/joined/joined_test.go delete mode 100644 packages/shared/pkg/middleware/otel/metrics/joined.go delete mode 100644 packages/shared/pkg/middleware/otel/metrics/joined_test.go diff --git a/packages/api/internal/orchestrator/create_instance.go b/packages/api/internal/orchestrator/create_instance.go index 22656d01cf..b1b9fb0445 100644 --- a/packages/api/internal/orchestrator/create_instance.go +++ b/packages/api/internal/orchestrator/create_instance.go @@ -26,7 +26,7 @@ import ( "github.com/e2b-dev/infra/packages/shared/pkg/featureflags" "github.com/e2b-dev/infra/packages/shared/pkg/grpc/orchestrator" "github.com/e2b-dev/infra/packages/shared/pkg/logger" - "github.com/e2b-dev/infra/packages/shared/pkg/middleware/otel/metrics" + "github.com/e2b-dev/infra/packages/shared/pkg/middleware/otel/joined" sandbox_network "github.com/e2b-dev/infra/packages/shared/pkg/sandbox-network" "github.com/e2b-dev/infra/packages/shared/pkg/telemetry" ut "github.com/e2b-dev/infra/packages/shared/pkg/utils" @@ -161,7 +161,7 @@ func (o *Orchestrator) CreateSandbox( if waitForStart != nil { // Mark as a joined request for telemetry purposes - metrics.MarkJoined(ctx) + joined.Mark(ctx) logger.L().Info(ctx, "sandbox is already being started, waiting for it to be ready", logger.WithSandboxID(sandboxID)) diff --git a/packages/api/internal/sandbox/storage/redis/state_change.go b/packages/api/internal/sandbox/storage/redis/state_change.go index 4cb220252e..8ab729ab8d 100644 --- a/packages/api/internal/sandbox/storage/redis/state_change.go +++ b/packages/api/internal/sandbox/storage/redis/state_change.go @@ -14,7 +14,7 @@ import ( "github.com/e2b-dev/infra/packages/api/internal/sandbox" "github.com/e2b-dev/infra/packages/shared/pkg/logger" - "github.com/e2b-dev/infra/packages/shared/pkg/middleware/otel/metrics" + "github.com/e2b-dev/infra/packages/shared/pkg/middleware/otel/joined" redis_utils "github.com/e2b-dev/infra/packages/shared/pkg/redis" ) @@ -251,7 +251,7 @@ func (s *Storage) waitForTransition( transitionID string, ) error { // Mark as a joined request for telemetry purposes - metrics.MarkJoined(ctx) + joined.Mark(ctx) routingKey := getTransitionRoutingKey(teamID.String(), sandboxID, transitionID) transitionKey := getTransitionKey(teamID.String(), sandboxID) diff --git a/packages/shared/pkg/middleware/otel/joined/joined.go b/packages/shared/pkg/middleware/otel/joined/joined.go new file mode 100644 index 0000000000..0dc3712b43 --- /dev/null +++ b/packages/shared/pkg/middleware/otel/joined/joined.go @@ -0,0 +1,86 @@ +// Package joined provides a request-scoped, concurrency-safe marker for +// "this HTTP request joined an in-flight concurrent operation rather than +// doing fresh work" (e.g. waiting for another request to finish a sandbox +// state transition, or piggy-backing on a concurrent CreateSandbox). +// +// The marker is installed once per request by either the tracing or the +// metrics middleware (both call WithHolder, which is idempotent). Any code +// path descended from the request context can flip the marker via Mark, +// without needing access to *gin.Context. +// +// - The marker writes request.joined="true" to the top-level HTTP server +// span captured at install time (so the attribute always lands on the +// root span, regardless of which child span is active when Mark fires). +// - The marker is exposed as a boolean histogram attribute via Attribute, +// so the metrics middleware can emit request.joined=true/false on every +// observation and dashboards can filter joiner vs. normal traffic. +package joined + +import ( + "context" + "sync/atomic" + + "go.opentelemetry.io/otel/attribute" + "go.opentelemetry.io/otel/trace" +) + +// AttributeKey is the dotted-lowercase key used for both the histogram +// attribute and the span attribute. +const AttributeKey = "request.joined" + +type holder struct { + joined atomic.Bool + // serverSpan is captured at holder install time. It is the span that is + // active on the ctx when WithHolder is called. In services that mount + // the tracing middleware before the metrics middleware (the convention + // in this repo), this is the HTTP server span. + serverSpan trace.Span +} + +type holderKey struct{} + +// WithHolder installs a fresh holder on ctx if one is not already present. +// Idempotent: when both the tracing and the metrics middleware are mounted, +// whichever runs first installs the holder; the other reuses it. +// +// Call WithHolder *after* the server span has been started so the holder +// captures the correct span for Mark to write attributes onto. +func WithHolder(ctx context.Context) context.Context { + if _, ok := ctx.Value(holderKey{}).(*holder); ok { + return ctx + } + + return context.WithValue(ctx, holderKey{}, &holder{ + serverSpan: trace.SpanFromContext(ctx), + }) +} + +// Mark marks the current request as a joiner. First-write-wins: subsequent +// calls in the same request are no-ops. Safe to call from any goroutine +// descended from the request ctx. No-op if no holder is on ctx (e.g. +// non-HTTP callers, tests). +func Mark(ctx context.Context) { + h, ok := ctx.Value(holderKey{}).(*holder) + if !ok { + return + } + + if h.joined.CompareAndSwap(false, true) { + h.serverSpan.SetAttributes( + attribute.String(AttributeKey, "true"), + ) + } +} + +// Attribute returns a boolean histogram attribute reflecting whether Mark +// has been called on this request's holder. Always returns a valid +// attribute so callers can append it unconditionally; if no holder is on +// ctx the value is false (no joiner status). +func Attribute(ctx context.Context) attribute.KeyValue { + h, ok := ctx.Value(holderKey{}).(*holder) + if !ok { + return attribute.Bool(AttributeKey, false) + } + + return attribute.Bool(AttributeKey, h.joined.Load()) +} diff --git a/packages/shared/pkg/middleware/otel/joined/joined_test.go b/packages/shared/pkg/middleware/otel/joined/joined_test.go new file mode 100644 index 0000000000..84f0471d34 --- /dev/null +++ b/packages/shared/pkg/middleware/otel/joined/joined_test.go @@ -0,0 +1,137 @@ +package joined_test + +import ( + "context" + "testing" + + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" + "go.opentelemetry.io/otel/attribute" + sdktrace "go.opentelemetry.io/otel/sdk/trace" + "go.opentelemetry.io/otel/sdk/trace/tracetest" + + "github.com/e2b-dev/infra/packages/shared/pkg/middleware/otel/joined" +) + +// Mark must be safe even when the context carries no holder. +func TestMark_NoHolder_Noop(t *testing.T) { + t.Parallel() + joined.Mark(context.Background()) +} + +// Attribute must return request.joined=false when no holder is on ctx. +func TestAttribute_NoHolder_ReturnsFalse(t *testing.T) { + t.Parallel() + + a := joined.Attribute(context.Background()) + assert.Equal(t, joined.AttributeKey, string(a.Key)) + assert.False(t, a.Value.AsBool()) +} + +// Attribute must return request.joined=false on a freshly installed holder +// before Mark has been called. +func TestAttribute_FreshHolder_ReturnsFalse(t *testing.T) { + t.Parallel() + + ctx := joined.WithHolder(context.Background()) + + a := joined.Attribute(ctx) + assert.False(t, a.Value.AsBool()) +} + +// Mark must flip Attribute to true on the same ctx. +func TestMark_FlipsAttributeToTrue(t *testing.T) { + t.Parallel() + + ctx := joined.WithHolder(context.Background()) + joined.Mark(ctx) + + a := joined.Attribute(ctx) + assert.True(t, a.Value.AsBool()) +} + +// WithHolder must be idempotent: calling it twice returns a ctx that shares +// the same underlying holder (Mark on the first ctx is visible from the +// second). +func TestWithHolder_Idempotent(t *testing.T) { + t.Parallel() + + ctx1 := joined.WithHolder(context.Background()) + ctx2 := joined.WithHolder(ctx1) + + joined.Mark(ctx1) + + a := joined.Attribute(ctx2) + assert.True(t, a.Value.AsBool(), "second WithHolder must reuse the first holder") +} + +// Mark must pin request.joined="true" onto the server span captured at +// WithHolder install time, not onto whichever child span is active when +// Mark fires. +func TestMark_TagsServerSpanNotChildSpan(t *testing.T) { + t.Parallel() + + sr := tracetest.NewSpanRecorder() + tp := sdktrace.NewTracerProvider(sdktrace.WithSpanProcessor(sr)) + tracer := tp.Tracer("test") + + // Open the server span before installing the holder (mirrors tracing + // middleware ordering: tracer.Start -> WithHolder). + ctx, serverSpan := tracer.Start(context.Background(), "HTTP POST /resume") + ctx = joined.WithHolder(ctx) + + // Open a child span (mirrors orchestrator.CreateSandbox's + // "create-sandbox" child) and call Mark from inside it. + childCtx, childSpan := tracer.Start(ctx, "create-sandbox") + joined.Mark(childCtx) + childSpan.End() + serverSpan.End() + + spans := sr.Ended() + require.Len(t, spans, 2) + + var server, child sdktrace.ReadOnlySpan + for _, s := range spans { + if s.Name() == "create-sandbox" { + child = s + } else { + server = s + } + } + require.NotNil(t, server) + require.NotNil(t, child) + + serverAttr, hasServer := findAttr(server.Attributes(), joined.AttributeKey) + require.True(t, hasServer, "request.joined must be on the server span") + assert.Equal(t, "true", serverAttr.AsString()) + + _, hasChild := findAttr(child.Attributes(), joined.AttributeKey) + assert.False(t, hasChild, "request.joined must NOT be on the child span") +} + +// Mark must be safe when called from a goroutine descended from the +// request context. +func TestMark_DescendantGoroutine(t *testing.T) { + t.Parallel() + + ctx := joined.WithHolder(context.Background()) + done := make(chan struct{}) + go func() { + joined.Mark(ctx) + close(done) + }() + <-done + + a := joined.Attribute(ctx) + assert.True(t, a.Value.AsBool()) +} + +func findAttr(attrs []attribute.KeyValue, key string) (attribute.Value, bool) { + for _, a := range attrs { + if string(a.Key) == key { + return a.Value, true + } + } + + return attribute.Value{}, false +} diff --git a/packages/shared/pkg/middleware/otel/metrics/joined.go b/packages/shared/pkg/middleware/otel/metrics/joined.go deleted file mode 100644 index a79d86b7f3..0000000000 --- a/packages/shared/pkg/middleware/otel/metrics/joined.go +++ /dev/null @@ -1,63 +0,0 @@ -package metrics - -import ( - "context" - "sync/atomic" - - "go.opentelemetry.io/otel/attribute" - "go.opentelemetry.io/otel/trace" -) - -const requestJoinedAttrKey = "request.joined" - -type joinedHolder struct { - joined atomic.Bool - // serverSpan is captured at middleware entry so MarkJoined can pin the - // request.joined attribute onto the top-level HTTP server span instead - // of whatever child span (e.g. "create-sandbox") happens to be active - // when the helper is called. Tracing middleware is registered before - // the metrics middleware in every service that mounts both, so the ctx - // passed to withJoinedHolder carries the server span as the active span. - // If no real span is on the ctx, this is a no-op span and SetAttributes - // silently does nothing. - serverSpan trace.Span -} - -type joinedHolderKey struct{} - -// withJoinedHolder installs a fresh joinedHolder on ctx and returns the -// augmented context plus a handle to the holder. The metrics middleware -// calls this once per request; non-HTTP callers normally never call it, -// which leaves the package helpers as safe no-ops. -func withJoinedHolder(ctx context.Context) (context.Context, *joinedHolder) { - h := &joinedHolder{ - serverSpan: trace.SpanFromContext(ctx), - } - - return context.WithValue(ctx, joinedHolderKey{}, h), h -} - -// MarkJoined marks the current request as having joined an in-flight -// concurrent operation (e.g. waiting for another request to finish a sandbox -// state transition, or joining a concurrent CreateSandbox). The flag is -// emitted as a histogram attribute on http.server.duration and as an -// attribute on the top-level HTTP server span (first-write-wins). -// -// Safe to call from any goroutine descended from the request context. -// No-op if ctx has no joinedHolder (e.g. non-HTTP callers, tests). -func MarkJoined(ctx context.Context) { - h, ok := ctx.Value(joinedHolderKey{}).(*joinedHolder) - if !ok { - return - } - - if h.joined.CompareAndSwap(false, true) { - h.serverSpan.SetAttributes( - attribute.String(requestJoinedAttrKey, "true"), - ) - } -} - -func (h *joinedHolder) joinedAttribute() attribute.KeyValue { - return attribute.Bool(requestJoinedAttrKey, h.joined.Load()) -} diff --git a/packages/shared/pkg/middleware/otel/metrics/joined_test.go b/packages/shared/pkg/middleware/otel/metrics/joined_test.go deleted file mode 100644 index 3629a49d88..0000000000 --- a/packages/shared/pkg/middleware/otel/metrics/joined_test.go +++ /dev/null @@ -1,286 +0,0 @@ -package metrics - -import ( - "context" - "net/http" - "net/http/httptest" - "sync" - "testing" - "time" - - "github.com/gin-gonic/gin" - "github.com/stretchr/testify/assert" - "github.com/stretchr/testify/require" - "go.opentelemetry.io/otel/attribute" - sdktrace "go.opentelemetry.io/otel/sdk/trace" - "go.opentelemetry.io/otel/sdk/trace/tracetest" -) - -// fakeRecorder captures every ObserveHTTPRequestDuration call so tests can -// inspect the attributes the middleware would have emitted. -type fakeRecorder struct { - mu sync.Mutex - calls []fakeRecorderCall -} - -type fakeRecorderCall struct { - duration time.Duration - attrs []attribute.KeyValue -} - -func (f *fakeRecorder) ObserveHTTPRequestDuration(_ context.Context, duration time.Duration, attrs []attribute.KeyValue) { - f.mu.Lock() - defer f.mu.Unlock() - // Copy attrs because the middleware reuses its slice buffer between - // requests when the test engine handles more than one request. - cp := make([]attribute.KeyValue, len(attrs)) - copy(cp, attrs) - f.calls = append(f.calls, fakeRecorderCall{duration: duration, attrs: cp}) -} - -func (f *fakeRecorder) callCount() int { - f.mu.Lock() - defer f.mu.Unlock() - - return len(f.calls) -} - -func (f *fakeRecorder) attrs(i int) []attribute.KeyValue { - f.mu.Lock() - defer f.mu.Unlock() - - return f.calls[i].attrs -} - -// findRequestJoined returns the value of the `request.joined` attribute from -// the given attribute slice, if present. -func findRequestJoined(attrs []attribute.KeyValue) (attribute.Value, bool) { - for _, a := range attrs { - if string(a.Key) == requestJoinedAttrKey { - return a.Value, true - } - } - - return attribute.Value{}, false -} - -func newTestEngine(t *testing.T, handler gin.HandlerFunc) (*gin.Engine, *fakeRecorder) { - t.Helper() - - rec := &fakeRecorder{} - r := gin.New() - r.Use(Middleware(nil, "test", WithRecorder(rec))) - r.POST("/sandboxes/:id/resume", handler) - - return r, rec -} - -func doRequest(t *testing.T, r *gin.Engine) { - t.Helper() - req := httptest.NewRequest(http.MethodPost, "/sandboxes/abc/resume", nil) - w := httptest.NewRecorder() - r.ServeHTTP(w, req) - require.Equal(t, http.StatusOK, w.Code) -} - -// MarkJoined must be safe even when the context carries no joinedHolder -// (e.g. non-HTTP callers, tests). -func TestMarkJoined_NoHolder_Noop(t *testing.T) { - t.Parallel() - MarkJoined(context.Background()) -} - -// Untagged requests must carry request.joined=false. -func TestMiddleware_NormalRequest_JoinedAttrIsFalse(t *testing.T) { - t.Parallel() - - r, rec := newTestEngine(t, func(c *gin.Context) { c.Status(http.StatusOK) }) - - doRequest(t, r) - - require.Equal(t, 1, rec.callCount()) - v, ok := findRequestJoined(rec.attrs(0)) - require.True(t, ok, "request.joined must be present on every observation") - assert.False(t, v.AsBool(), "untagged request must carry request.joined=false") -} - -// MarkJoined from the handler's ctx must flip the attribute to true. -func TestMiddleware_MarkJoinedFromHandler_AppearsAsTrue(t *testing.T) { - t.Parallel() - - r, rec := newTestEngine(t, func(c *gin.Context) { - MarkJoined(c.Request.Context()) - c.Status(http.StatusOK) - }) - - doRequest(t, r) - - require.Equal(t, 1, rec.callCount()) - v, ok := findRequestJoined(rec.attrs(0)) - require.True(t, ok, "request.joined attribute must be present on the histogram") - assert.True(t, v.AsBool()) -} - -// MarkJoined called from a goroutine descended from the request context -// must still flow through to the histogram. This is the key capability of -// the context-attached holder design vs. a *gin.Context-based marker. -func TestMiddleware_MarkJoinedFromDescendantGoroutine(t *testing.T) { - t.Parallel() - - done := make(chan struct{}) - r, rec := newTestEngine(t, func(c *gin.Context) { - ctx := c.Request.Context() - go func() { - MarkJoined(ctx) - close(done) - }() - <-done - c.Status(http.StatusOK) - }) - - doRequest(t, r) - - require.Equal(t, 1, rec.callCount()) - v, ok := findRequestJoined(rec.attrs(0)) - require.True(t, ok) - assert.True(t, v.AsBool()) -} - -// MarkJoined is idempotent: repeated calls within the same request do not -// produce duplicate histogram attributes. -func TestMiddleware_MarkJoinedIdempotent(t *testing.T) { - t.Parallel() - - r, rec := newTestEngine(t, func(c *gin.Context) { - MarkJoined(c.Request.Context()) - MarkJoined(c.Request.Context()) - MarkJoined(c.Request.Context()) - c.Status(http.StatusOK) - }) - - doRequest(t, r) - - require.Equal(t, 1, rec.callCount()) - attrs := rec.attrs(0) - - count := 0 - for _, a := range attrs { - if string(a.Key) == requestJoinedAttrKey { - count++ - } - } - assert.Equal(t, 1, count, "request.joined must appear exactly once even after repeated MarkJoined calls") -} - -// Tagging must not suppress recording — we only add a label. -func TestMiddleware_Tagging_DoesNotSuppressRecording(t *testing.T) { - t.Parallel() - - r, rec := newTestEngine(t, func(c *gin.Context) { - MarkJoined(c.Request.Context()) - c.Status(http.StatusOK) - }) - - doRequest(t, r) - - assert.Equal(t, 1, rec.callCount(), "histogram must still be recorded; tagging only adds attributes") -} - -// MarkJoined must pin the request.joined attribute onto the top-level HTTP -// server span (captured at middleware entry), not onto whatever child span -// happens to be active when the helper is called. This is the guarantee -// callers rely on for Tempo filtering by root-span attribute. -func TestMiddleware_MarkJoined_TagsServerSpanNotChildSpan(t *testing.T) { - t.Parallel() - - sr := tracetest.NewSpanRecorder() - tp := sdktrace.NewTracerProvider(sdktrace.WithSpanProcessor(sr)) - tracer := tp.Tracer("github.com/e2b-dev/infra/packages/shared/pkg/middleware/otel/metrics") - - rec := &fakeRecorder{} - r := gin.New() - // Stand-in for the tracing middleware: open the server span before the - // metrics middleware installs its joinedHolder. This mirrors the real - // service wiring (tracing registered before metrics). - r.Use(func(c *gin.Context) { - ctx, span := tracer.Start(c.Request.Context(), "HTTP POST /sandboxes/:id/resume") - defer span.End() - c.Request = c.Request.WithContext(ctx) - c.Next() - }) - r.Use(Middleware(nil, "test", WithRecorder(rec))) - r.POST("/sandboxes/:id/resume", func(c *gin.Context) { - // Open a child span (mirrors orchestrator.CreateSandbox's - // "create-sandbox" child span) and call MarkJoined from inside it. - ctx, child := tracer.Start(c.Request.Context(), "create-sandbox") - MarkJoined(ctx) - child.End() - - c.Status(http.StatusOK) - }) - - req := httptest.NewRequest(http.MethodPost, "/sandboxes/abc/resume", nil) - w := httptest.NewRecorder() - r.ServeHTTP(w, req) - require.Equal(t, http.StatusOK, w.Code) - - spans := sr.Ended() - require.Len(t, spans, 2) - - // Identify server vs child by span name. - var serverSpan, childSpan sdktrace.ReadOnlySpan - for _, s := range spans { - if s.Name() == "create-sandbox" { - childSpan = s - } else { - serverSpan = s - } - } - require.NotNil(t, serverSpan, "server span must be recorded") - require.NotNil(t, childSpan, "child span must be recorded") - - serverAttr, hasServer := findRequestJoined(serverSpan.Attributes()) - require.True(t, hasServer, "request.joined must be on the server span") - assert.Equal(t, "true", serverAttr.AsString()) - - _, hasChild := findRequestJoined(childSpan.Attributes()) - assert.False(t, hasChild, "request.joined must NOT be on the child span") - - // And the histogram must still carry request.joined=true. - require.Equal(t, 1, rec.callCount()) - histAttr, ok := findRequestJoined(rec.attrs(0)) - require.True(t, ok) - assert.True(t, histAttr.AsBool()) -} - -// Two distinct requests must not share the holder: tagging one must not -// taint the other. -func TestMiddleware_HolderIsRequestScoped(t *testing.T) { - t.Parallel() - - rec := &fakeRecorder{} - r := gin.New() - r.Use(Middleware(nil, "test", WithRecorder(rec))) - r.POST("/joiner", func(c *gin.Context) { - MarkJoined(c.Request.Context()) - c.Status(http.StatusOK) - }) - r.POST("/normal", func(c *gin.Context) { c.Status(http.StatusOK) }) - - for _, path := range []string{"/joiner", "/normal"} { - req := httptest.NewRequest(http.MethodPost, path, nil) - w := httptest.NewRecorder() - r.ServeHTTP(w, req) - require.Equal(t, http.StatusOK, w.Code) - } - - require.Equal(t, 2, rec.callCount()) - - v1, ok := findRequestJoined(rec.attrs(0)) - require.True(t, ok) - assert.True(t, v1.AsBool(), "first (tagged) request must carry request.joined=true") - - v2, ok := findRequestJoined(rec.attrs(1)) - require.True(t, ok) - assert.False(t, v2.AsBool(), "second (untagged) request must carry request.joined=false") -} diff --git a/packages/shared/pkg/middleware/otel/metrics/middleware.go b/packages/shared/pkg/middleware/otel/metrics/middleware.go index 6fec11b828..62bc4a6948 100644 --- a/packages/shared/pkg/middleware/otel/metrics/middleware.go +++ b/packages/shared/pkg/middleware/otel/metrics/middleware.go @@ -13,6 +13,7 @@ import ( semconv "go.opentelemetry.io/otel/semconv/v1.7.0" sharedmiddleware "github.com/e2b-dev/infra/packages/shared/pkg/middleware" + "github.com/e2b-dev/infra/packages/shared/pkg/middleware/otel/joined" ) const MetricPrefix = "metric." @@ -57,10 +58,11 @@ func Middleware(meterProvider metric.MeterProvider, service string, options ...O return func(ginCtx *gin.Context) { ctx := ginCtx.Request.Context() - // Install the request-scoped joinedHolder so any descendant code - // path (orchestrator, storage layer, etc.) can mark the request via - // the package helpers without needing access to *gin.Context. - ctx, holder := withJoinedHolder(ctx) + // Install the request-scoped joined holder so descendant code paths + // (orchestrator, storage layer, etc.) can call joined.Mark via + // context.Context. Idempotent: if the tracing middleware already + // installed the holder, this reuses it. + ctx = joined.WithHolder(ctx) ginCtx.Request = ginCtx.Request.WithContext(ctx) route := ginCtx.FullPath() @@ -97,7 +99,7 @@ func Middleware(meterProvider metric.MeterProvider, service string, options ...O resAttributes = append(resAttributes, attributesFromGinContext(ginCtx, MetricPrefix)...) // Distinguish between regular and joined requests - resAttributes = append(resAttributes, holder.joinedAttribute()) + resAttributes = append(resAttributes, joined.Attribute(ctx)) effectiveStart := start if processingStart, ok := getProcessingStartTime(ginCtx); ok { diff --git a/packages/shared/pkg/middleware/otel/tracing/middleware.go b/packages/shared/pkg/middleware/otel/tracing/middleware.go index 0bbe201e4a..6dbae5130a 100644 --- a/packages/shared/pkg/middleware/otel/tracing/middleware.go +++ b/packages/shared/pkg/middleware/otel/tracing/middleware.go @@ -31,6 +31,7 @@ import ( "github.com/e2b-dev/infra/packages/shared/pkg/logger" sharedmiddleware "github.com/e2b-dev/infra/packages/shared/pkg/middleware" + "github.com/e2b-dev/infra/packages/shared/pkg/middleware/otel/joined" "github.com/e2b-dev/infra/packages/shared/pkg/telemetry" ) @@ -110,6 +111,13 @@ func Middleware(tracerProvider oteltrace.TracerProvider, service string) gin.Han ctx, span := tracer.Start(ctx, spanName, opts...) defer span.End() + // Install the request-scoped joined holder so descendant code paths + // (orchestrator, storage layer, etc.) can call joined.Mark and have + // the request.joined attribute pinned to this (top-level) server + // span. Idempotent: a later joined.WithHolder call from the metrics + // middleware will reuse this holder. + ctx = joined.WithHolder(ctx) + // pass the span through the request context c.Request = c.Request.WithContext(ctx) From a431c307feda448759eacc14523f5849d12e37b4 Mon Sep 17 00:00:00 2001 From: "github-actions[bot]" Date: Sun, 17 May 2026 20:07:59 +0000 Subject: [PATCH 05/12] chore: auto-commit generated changes --- packages/shared/pkg/middleware/otel/joined/joined_test.go | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/packages/shared/pkg/middleware/otel/joined/joined_test.go b/packages/shared/pkg/middleware/otel/joined/joined_test.go index 84f0471d34..20fc22947f 100644 --- a/packages/shared/pkg/middleware/otel/joined/joined_test.go +++ b/packages/shared/pkg/middleware/otel/joined/joined_test.go @@ -73,7 +73,7 @@ func TestMark_TagsServerSpanNotChildSpan(t *testing.T) { sr := tracetest.NewSpanRecorder() tp := sdktrace.NewTracerProvider(sdktrace.WithSpanProcessor(sr)) - tracer := tp.Tracer("test") + tracer := tp.Tracer("github.com/e2b-dev/infra/packages/shared/pkg/middleware/otel/joined") // Open the server span before installing the holder (mirrors tracing // middleware ordering: tracer.Start -> WithHolder). From aa030491ebcdc557fd21cd86281cce5336acd991 Mon Sep 17 00:00:00 2001 From: Jakub Novak Date: Mon, 18 May 2026 07:01:08 +0000 Subject: [PATCH 06/12] feat(joined): set request.joined on every server span via tracing middleware - Change joined.Mark to use attribute.Bool(true) instead of String("true") so the span attribute and the histogram attribute have a consistent type and value across Tempo and Prometheus queries. - Tracing middleware now reads joined.Attribute(ctx) after the request completes and pins it on the server span (true or false). Untagged requests now carry request.joined=false on their span, matching the histogram label semantics and removing the need to special-case the 'attribute absent' state in trace queries. --- packages/shared/pkg/middleware/otel/joined/joined.go | 2 +- packages/shared/pkg/middleware/otel/joined/joined_test.go | 2 +- packages/shared/pkg/middleware/otel/tracing/middleware.go | 7 +++++++ 3 files changed, 9 insertions(+), 2 deletions(-) diff --git a/packages/shared/pkg/middleware/otel/joined/joined.go b/packages/shared/pkg/middleware/otel/joined/joined.go index 0dc3712b43..7b209df775 100644 --- a/packages/shared/pkg/middleware/otel/joined/joined.go +++ b/packages/shared/pkg/middleware/otel/joined/joined.go @@ -67,7 +67,7 @@ func Mark(ctx context.Context) { if h.joined.CompareAndSwap(false, true) { h.serverSpan.SetAttributes( - attribute.String(AttributeKey, "true"), + attribute.Bool(AttributeKey, true), ) } } diff --git a/packages/shared/pkg/middleware/otel/joined/joined_test.go b/packages/shared/pkg/middleware/otel/joined/joined_test.go index 20fc22947f..d7735ea222 100644 --- a/packages/shared/pkg/middleware/otel/joined/joined_test.go +++ b/packages/shared/pkg/middleware/otel/joined/joined_test.go @@ -103,7 +103,7 @@ func TestMark_TagsServerSpanNotChildSpan(t *testing.T) { serverAttr, hasServer := findAttr(server.Attributes(), joined.AttributeKey) require.True(t, hasServer, "request.joined must be on the server span") - assert.Equal(t, "true", serverAttr.AsString()) + assert.True(t, serverAttr.AsBool()) _, hasChild := findAttr(child.Attributes(), joined.AttributeKey) assert.False(t, hasChild, "request.joined must NOT be on the child span") diff --git a/packages/shared/pkg/middleware/otel/tracing/middleware.go b/packages/shared/pkg/middleware/otel/tracing/middleware.go index 6dbae5130a..95476240fa 100644 --- a/packages/shared/pkg/middleware/otel/tracing/middleware.go +++ b/packages/shared/pkg/middleware/otel/tracing/middleware.go @@ -138,6 +138,13 @@ func Middleware(tracerProvider oteltrace.TracerProvider, service string) gin.Han spanStatus, spanMessage := semconv.SpanStatusFromHTTPStatusCode(status) span.SetStatus(spanStatus, spanMessage) + // Ensure every server span carries request.joined (true or false) + // so trace queries can filter joiner vs. normal traffic without + // special-casing the "attribute absent" case. Idempotent: if + // joined.Mark already set it to true earlier, this re-asserts the + // same value; otherwise it pins it to false. + span.SetAttributes(joined.Attribute(ctx)) + if len(c.Errors) > 0 { span.SetAttributes(attribute.String("gin.errors", strings.TrimSpace(c.Errors.String()))) } From 43dd9b394d9981bccf0bc77a0d7ecf03ccbb83bb Mon Sep 17 00:00:00 2001 From: Jakub Novak Date: Mon, 18 May 2026 08:43:40 +0000 Subject: [PATCH 07/12] chore(joined): trim verbose comments and drop covered test The server-span-vs-child-span behavior is verified end-to-end against the live local cluster (Tempo confirms request.joined lands on the root span) so the dedicated unit test was redundant. --- .../pkg/middleware/otel/joined/joined.go | 32 +--------- .../pkg/middleware/otel/joined/joined_test.go | 58 ------------------- .../pkg/middleware/otel/metrics/middleware.go | 1 + .../pkg/middleware/otel/tracing/middleware.go | 12 +--- 4 files changed, 5 insertions(+), 98 deletions(-) diff --git a/packages/shared/pkg/middleware/otel/joined/joined.go b/packages/shared/pkg/middleware/otel/joined/joined.go index 7b209df775..dd7667ab33 100644 --- a/packages/shared/pkg/middleware/otel/joined/joined.go +++ b/packages/shared/pkg/middleware/otel/joined/joined.go @@ -1,19 +1,3 @@ -// Package joined provides a request-scoped, concurrency-safe marker for -// "this HTTP request joined an in-flight concurrent operation rather than -// doing fresh work" (e.g. waiting for another request to finish a sandbox -// state transition, or piggy-backing on a concurrent CreateSandbox). -// -// The marker is installed once per request by either the tracing or the -// metrics middleware (both call WithHolder, which is idempotent). Any code -// path descended from the request context can flip the marker via Mark, -// without needing access to *gin.Context. -// -// - The marker writes request.joined="true" to the top-level HTTP server -// span captured at install time (so the attribute always lands on the -// root span, regardless of which child span is active when Mark fires). -// - The marker is exposed as a boolean histogram attribute via Attribute, -// so the metrics middleware can emit request.joined=true/false on every -// observation and dashboards can filter joiner vs. normal traffic. package joined import ( @@ -30,18 +14,13 @@ const AttributeKey = "request.joined" type holder struct { joined atomic.Bool - // serverSpan is captured at holder install time. It is the span that is - // active on the ctx when WithHolder is called. In services that mount - // the tracing middleware before the metrics middleware (the convention - // in this repo), this is the HTTP server span. + serverSpan trace.Span } type holderKey struct{} // WithHolder installs a fresh holder on ctx if one is not already present. -// Idempotent: when both the tracing and the metrics middleware are mounted, -// whichever runs first installs the holder; the other reuses it. // // Call WithHolder *after* the server span has been started so the holder // captures the correct span for Mark to write attributes onto. @@ -55,10 +34,7 @@ func WithHolder(ctx context.Context) context.Context { }) } -// Mark marks the current request as a joiner. First-write-wins: subsequent -// calls in the same request are no-ops. Safe to call from any goroutine -// descended from the request ctx. No-op if no holder is on ctx (e.g. -// non-HTTP callers, tests). +// Mark marks the current request as a joiner. First-write-wins func Mark(ctx context.Context) { h, ok := ctx.Value(holderKey{}).(*holder) if !ok { @@ -72,10 +48,6 @@ func Mark(ctx context.Context) { } } -// Attribute returns a boolean histogram attribute reflecting whether Mark -// has been called on this request's holder. Always returns a valid -// attribute so callers can append it unconditionally; if no holder is on -// ctx the value is false (no joiner status). func Attribute(ctx context.Context) attribute.KeyValue { h, ok := ctx.Value(holderKey{}).(*holder) if !ok { diff --git a/packages/shared/pkg/middleware/otel/joined/joined_test.go b/packages/shared/pkg/middleware/otel/joined/joined_test.go index d7735ea222..b695f87b5e 100644 --- a/packages/shared/pkg/middleware/otel/joined/joined_test.go +++ b/packages/shared/pkg/middleware/otel/joined/joined_test.go @@ -5,10 +5,6 @@ import ( "testing" "github.com/stretchr/testify/assert" - "github.com/stretchr/testify/require" - "go.opentelemetry.io/otel/attribute" - sdktrace "go.opentelemetry.io/otel/sdk/trace" - "go.opentelemetry.io/otel/sdk/trace/tracetest" "github.com/e2b-dev/infra/packages/shared/pkg/middleware/otel/joined" ) @@ -65,50 +61,6 @@ func TestWithHolder_Idempotent(t *testing.T) { assert.True(t, a.Value.AsBool(), "second WithHolder must reuse the first holder") } -// Mark must pin request.joined="true" onto the server span captured at -// WithHolder install time, not onto whichever child span is active when -// Mark fires. -func TestMark_TagsServerSpanNotChildSpan(t *testing.T) { - t.Parallel() - - sr := tracetest.NewSpanRecorder() - tp := sdktrace.NewTracerProvider(sdktrace.WithSpanProcessor(sr)) - tracer := tp.Tracer("github.com/e2b-dev/infra/packages/shared/pkg/middleware/otel/joined") - - // Open the server span before installing the holder (mirrors tracing - // middleware ordering: tracer.Start -> WithHolder). - ctx, serverSpan := tracer.Start(context.Background(), "HTTP POST /resume") - ctx = joined.WithHolder(ctx) - - // Open a child span (mirrors orchestrator.CreateSandbox's - // "create-sandbox" child) and call Mark from inside it. - childCtx, childSpan := tracer.Start(ctx, "create-sandbox") - joined.Mark(childCtx) - childSpan.End() - serverSpan.End() - - spans := sr.Ended() - require.Len(t, spans, 2) - - var server, child sdktrace.ReadOnlySpan - for _, s := range spans { - if s.Name() == "create-sandbox" { - child = s - } else { - server = s - } - } - require.NotNil(t, server) - require.NotNil(t, child) - - serverAttr, hasServer := findAttr(server.Attributes(), joined.AttributeKey) - require.True(t, hasServer, "request.joined must be on the server span") - assert.True(t, serverAttr.AsBool()) - - _, hasChild := findAttr(child.Attributes(), joined.AttributeKey) - assert.False(t, hasChild, "request.joined must NOT be on the child span") -} - // Mark must be safe when called from a goroutine descended from the // request context. func TestMark_DescendantGoroutine(t *testing.T) { @@ -125,13 +77,3 @@ func TestMark_DescendantGoroutine(t *testing.T) { a := joined.Attribute(ctx) assert.True(t, a.Value.AsBool()) } - -func findAttr(attrs []attribute.KeyValue, key string) (attribute.Value, bool) { - for _, a := range attrs { - if string(a.Key) == key { - return a.Value, true - } - } - - return attribute.Value{}, false -} diff --git a/packages/shared/pkg/middleware/otel/metrics/middleware.go b/packages/shared/pkg/middleware/otel/metrics/middleware.go index 62bc4a6948..f4e7da0f23 100644 --- a/packages/shared/pkg/middleware/otel/metrics/middleware.go +++ b/packages/shared/pkg/middleware/otel/metrics/middleware.go @@ -101,6 +101,7 @@ func Middleware(meterProvider metric.MeterProvider, service string, options ...O // Distinguish between regular and joined requests resAttributes = append(resAttributes, joined.Attribute(ctx)) + // Use processing start time if set, otherwise fall back to the middleware start time. effectiveStart := start if processingStart, ok := getProcessingStartTime(ginCtx); ok { effectiveStart = processingStart diff --git a/packages/shared/pkg/middleware/otel/tracing/middleware.go b/packages/shared/pkg/middleware/otel/tracing/middleware.go index 95476240fa..da1b6f932a 100644 --- a/packages/shared/pkg/middleware/otel/tracing/middleware.go +++ b/packages/shared/pkg/middleware/otel/tracing/middleware.go @@ -111,11 +111,7 @@ func Middleware(tracerProvider oteltrace.TracerProvider, service string) gin.Han ctx, span := tracer.Start(ctx, spanName, opts...) defer span.End() - // Install the request-scoped joined holder so descendant code paths - // (orchestrator, storage layer, etc.) can call joined.Mark and have - // the request.joined attribute pinned to this (top-level) server - // span. Idempotent: a later joined.WithHolder call from the metrics - // middleware will reuse this holder. + // Install the request-scoped joined holder ctx = joined.WithHolder(ctx) // pass the span through the request context @@ -138,11 +134,7 @@ func Middleware(tracerProvider oteltrace.TracerProvider, service string) gin.Han spanStatus, spanMessage := semconv.SpanStatusFromHTTPStatusCode(status) span.SetStatus(spanStatus, spanMessage) - // Ensure every server span carries request.joined (true or false) - // so trace queries can filter joiner vs. normal traffic without - // special-casing the "attribute absent" case. Idempotent: if - // joined.Mark already set it to true earlier, this re-asserts the - // same value; otherwise it pins it to false. + // Marks the joined requests for telemetry purposes. span.SetAttributes(joined.Attribute(ctx)) if len(c.Errors) > 0 { From a4bc153cead0a3fff200bb79db6e28445faaf0e0 Mon Sep 17 00:00:00 2001 From: Jakub Novak Date: Mon, 18 May 2026 08:54:29 +0000 Subject: [PATCH 08/12] fix(joined): tag WaitForStateChange entry in both storage backends Resume's case StatePausing -> WaitForStateChange branch is intentionally racing another in-flight operation regardless of whether the transition is still pending by the time we look it up. Tag the request at the public entry of WaitForStateChange so: - Redis backend: the fast-path that finds the transition already completed (redis.Nil) is still tagged. Previously only waitForTransition tagged, missing this race window. - Memory backend (local dev / tests): was never tagging at all. Now is. Mark is idempotent so the redis path that also tags inside waitForTransition remains correct. --- packages/api/internal/sandbox/storage/memory/operations.go | 5 +++++ packages/api/internal/sandbox/storage/redis/state_change.go | 6 ++++++ 2 files changed, 11 insertions(+) diff --git a/packages/api/internal/sandbox/storage/memory/operations.go b/packages/api/internal/sandbox/storage/memory/operations.go index 74550992bf..0e1378b621 100644 --- a/packages/api/internal/sandbox/storage/memory/operations.go +++ b/packages/api/internal/sandbox/storage/memory/operations.go @@ -12,6 +12,7 @@ import ( "github.com/e2b-dev/infra/packages/api/internal/sandbox" "github.com/e2b-dev/infra/packages/shared/pkg/logger" + "github.com/e2b-dev/infra/packages/shared/pkg/middleware/otel/joined" "github.com/e2b-dev/infra/packages/shared/pkg/utils" ) @@ -251,6 +252,10 @@ func startRemoving(ctx context.Context, sbx *memorySandbox, opts sandbox.RemoveO } func (s *Storage) WaitForStateChange(ctx context.Context, _ uuid.UUID, sandboxID string) error { + // Mark as a joined request for telemetry purposes: the caller explicitly + // took the "wait for state change" branch. + joined.Mark(ctx) + sbx, err := s.get(sandboxID) if err != nil { return fmt.Errorf("failed to get sandbox: %w", err) diff --git a/packages/api/internal/sandbox/storage/redis/state_change.go b/packages/api/internal/sandbox/storage/redis/state_change.go index 8ab729ab8d..f0dfc250fc 100644 --- a/packages/api/internal/sandbox/storage/redis/state_change.go +++ b/packages/api/internal/sandbox/storage/redis/state_change.go @@ -229,6 +229,12 @@ func (s *Storage) restoreToRunning(ctx context.Context, teamID uuid.UUID, sandbo // WaitForStateChange waits for a sandbox state transition to complete. func (s *Storage) WaitForStateChange(ctx context.Context, teamID uuid.UUID, sandboxID string) error { + // Mark as a joined request for telemetry purposes: the caller explicitly + // took the "wait for state change" branch, so it is intentionally racing + // another in-flight operation regardless of whether the transition is + // still pending or already completed by the time we look it up. + joined.Mark(ctx) + transitionKey := getTransitionKey(teamID.String(), sandboxID) transactionID, err := s.redisClient.Get(ctx, transitionKey).Result() if errors.Is(err, redis.Nil) { From af2bb4b44436e363cc40af243d9544c5b6640d96 Mon Sep 17 00:00:00 2001 From: Jakub Novak Date: Mon, 18 May 2026 08:57:02 +0000 Subject: [PATCH 09/12] fix(joined): only mark when actually blocking on another transition Revert the previous "tag at WaitForStateChange entry" change. The request.joined attribute should reflect outcome, not intent: if the main action did not actually join another in-flight operation, the request is not a joiner. - Redis backend: keep the existing Mark inside waitForTransition. The redis.Nil fast-return path stays untagged because no actual joining happened. - Memory backend: move Mark from the public WaitForStateChange entry to inside the if-transition-not-nil branch, so it fires only when the request actually blocks on transition.WaitWithContext. --- .../api/internal/sandbox/storage/memory/operations.go | 8 ++++---- .../api/internal/sandbox/storage/redis/state_change.go | 6 ------ 2 files changed, 4 insertions(+), 10 deletions(-) diff --git a/packages/api/internal/sandbox/storage/memory/operations.go b/packages/api/internal/sandbox/storage/memory/operations.go index 0e1378b621..bd6636a13e 100644 --- a/packages/api/internal/sandbox/storage/memory/operations.go +++ b/packages/api/internal/sandbox/storage/memory/operations.go @@ -252,10 +252,6 @@ func startRemoving(ctx context.Context, sbx *memorySandbox, opts sandbox.RemoveO } func (s *Storage) WaitForStateChange(ctx context.Context, _ uuid.UUID, sandboxID string) error { - // Mark as a joined request for telemetry purposes: the caller explicitly - // took the "wait for state change" branch. - joined.Mark(ctx) - sbx, err := s.get(sandboxID) if err != nil { return fmt.Errorf("failed to get sandbox: %w", err) @@ -272,5 +268,9 @@ func waitForStateChange(ctx context.Context, sbx *memorySandbox) error { return nil } + // Mark as a joined request: we are about to actually block on another + // in-flight transition. + joined.Mark(ctx) + return transition.WaitWithContext(ctx) } diff --git a/packages/api/internal/sandbox/storage/redis/state_change.go b/packages/api/internal/sandbox/storage/redis/state_change.go index f0dfc250fc..8ab729ab8d 100644 --- a/packages/api/internal/sandbox/storage/redis/state_change.go +++ b/packages/api/internal/sandbox/storage/redis/state_change.go @@ -229,12 +229,6 @@ func (s *Storage) restoreToRunning(ctx context.Context, teamID uuid.UUID, sandbo // WaitForStateChange waits for a sandbox state transition to complete. func (s *Storage) WaitForStateChange(ctx context.Context, teamID uuid.UUID, sandboxID string) error { - // Mark as a joined request for telemetry purposes: the caller explicitly - // took the "wait for state change" branch, so it is intentionally racing - // another in-flight operation regardless of whether the transition is - // still pending or already completed by the time we look it up. - joined.Mark(ctx) - transitionKey := getTransitionKey(teamID.String(), sandboxID) transactionID, err := s.redisClient.Get(ctx, transitionKey).Result() if errors.Is(err, redis.Nil) { From 117b140a4d715039969d3aa94a68406cacebcff9 Mon Sep 17 00:00:00 2001 From: Jakub Novak Date: Mon, 18 May 2026 09:07:25 +0000 Subject: [PATCH 10/12] fix(joined): only mark the CreateSandbox reservation joiner A state-change wait (resume seeing StatePausing, connect retrying through a transition, StartRemoving observing an existing transition) is not a joiner: the request still does its own work after the wait. Only the CreateSandbox reservation joiner branch (waitForStart != nil) gets another request's result for free, so it is the only true joiner site. - Remove joined.Mark from waitForTransition (redis backend). - Remove joined.Mark from waitForStateChange helper (memory backend). - Drop now-unused joined imports in both files. The reservation joiner in packages/api/internal/orchestrator/create_instance.go is the sole remaining production call site. --- packages/api/internal/sandbox/storage/memory/operations.go | 5 ----- packages/api/internal/sandbox/storage/redis/state_change.go | 4 ---- 2 files changed, 9 deletions(-) diff --git a/packages/api/internal/sandbox/storage/memory/operations.go b/packages/api/internal/sandbox/storage/memory/operations.go index bd6636a13e..74550992bf 100644 --- a/packages/api/internal/sandbox/storage/memory/operations.go +++ b/packages/api/internal/sandbox/storage/memory/operations.go @@ -12,7 +12,6 @@ import ( "github.com/e2b-dev/infra/packages/api/internal/sandbox" "github.com/e2b-dev/infra/packages/shared/pkg/logger" - "github.com/e2b-dev/infra/packages/shared/pkg/middleware/otel/joined" "github.com/e2b-dev/infra/packages/shared/pkg/utils" ) @@ -268,9 +267,5 @@ func waitForStateChange(ctx context.Context, sbx *memorySandbox) error { return nil } - // Mark as a joined request: we are about to actually block on another - // in-flight transition. - joined.Mark(ctx) - return transition.WaitWithContext(ctx) } diff --git a/packages/api/internal/sandbox/storage/redis/state_change.go b/packages/api/internal/sandbox/storage/redis/state_change.go index 8ab729ab8d..483642d6f8 100644 --- a/packages/api/internal/sandbox/storage/redis/state_change.go +++ b/packages/api/internal/sandbox/storage/redis/state_change.go @@ -14,7 +14,6 @@ import ( "github.com/e2b-dev/infra/packages/api/internal/sandbox" "github.com/e2b-dev/infra/packages/shared/pkg/logger" - "github.com/e2b-dev/infra/packages/shared/pkg/middleware/otel/joined" redis_utils "github.com/e2b-dev/infra/packages/shared/pkg/redis" ) @@ -250,9 +249,6 @@ func (s *Storage) waitForTransition( sandboxID, transitionID string, ) error { - // Mark as a joined request for telemetry purposes - joined.Mark(ctx) - routingKey := getTransitionRoutingKey(teamID.String(), sandboxID, transitionID) transitionKey := getTransitionKey(teamID.String(), sandboxID) resultKey := getTransitionResultKey(teamID.String(), sandboxID, transitionID) From bf1d13d29060ab72c5325b3d0947357389338b4d Mon Sep 17 00:00:00 2001 From: Jakub Novak Date: Mon, 18 May 2026 09:09:55 +0000 Subject: [PATCH 11/12] feat(joined): mark concurrent same-state transitions as joiners A second concurrent pause/kill/snapshot for the same sandbox that targets the same state as an already in-flight transition waits for it and inherits the result via alreadyDone=true, doing no work itself. That is a true joiner, semantically identical to the CreateSandbox reservation joiner. - Redis handleExistingTransition: tag inside the sbx.State == newState branch only. The different-state branch retries with its own work, so remains untagged. - Memory startRemoving: tag inside the currentState == newState case after the wait. The allowed-transition case recurses to do its own work and remains untagged. --- packages/api/internal/sandbox/storage/memory/operations.go | 5 +++++ .../api/internal/sandbox/storage/redis/state_change.go | 7 ++++++- 2 files changed, 11 insertions(+), 1 deletion(-) diff --git a/packages/api/internal/sandbox/storage/memory/operations.go b/packages/api/internal/sandbox/storage/memory/operations.go index 74550992bf..4be36c5428 100644 --- a/packages/api/internal/sandbox/storage/memory/operations.go +++ b/packages/api/internal/sandbox/storage/memory/operations.go @@ -12,6 +12,7 @@ import ( "github.com/e2b-dev/infra/packages/api/internal/sandbox" "github.com/e2b-dev/infra/packages/shared/pkg/logger" + "github.com/e2b-dev/infra/packages/shared/pkg/middleware/otel/joined" "github.com/e2b-dev/infra/packages/shared/pkg/utils" ) @@ -192,6 +193,10 @@ func startRemoving(ctx context.Context, sbx *memorySandbox, opts sandbox.RemoveO // If the transition is to the same state just wait switch { case currentState == newState: + // The caller inherits the in-flight transition's result + // without doing the work itself: this is a joiner. + joined.Mark(ctx) + return true, func(context.Context, error) {}, nil case sandbox.AllowedTransitions[currentState][newState]: return startRemoving(ctx, sbx, sandbox.RemoveOpts{Action: opts.Action}) diff --git a/packages/api/internal/sandbox/storage/redis/state_change.go b/packages/api/internal/sandbox/storage/redis/state_change.go index 483642d6f8..d7ac1c7dcc 100644 --- a/packages/api/internal/sandbox/storage/redis/state_change.go +++ b/packages/api/internal/sandbox/storage/redis/state_change.go @@ -14,6 +14,7 @@ import ( "github.com/e2b-dev/infra/packages/api/internal/sandbox" "github.com/e2b-dev/infra/packages/shared/pkg/logger" + "github.com/e2b-dev/infra/packages/shared/pkg/middleware/otel/joined" redis_utils "github.com/e2b-dev/infra/packages/shared/pkg/redis" ) @@ -317,7 +318,11 @@ func (s *Storage) handleExistingTransition( transactionID string, ) (sandbox.Sandbox, bool, func(context.Context, error), error) { if sbx.State == newState { - // Same target state - wait for completion and return alreadyDone=true + // Same target state - wait for completion and return alreadyDone=true. + // The caller inherits the in-flight transition's result without + // doing the work itself: this is a joiner. + joined.Mark(ctx) + logger.L().Debug(ctx, "State transition already in progress to the same state, waiting", logger.WithSandboxID(sbx.SandboxID), zap.String("state", string(newState))) From 59c7630223657a7fdff2d3979a1b26c5bd0d9b6c Mon Sep 17 00:00:00 2001 From: Jakub Novak Date: Mon, 18 May 2026 09:47:14 +0000 Subject: [PATCH 12/12] fix(joined): mark memory backend joiners before waiting MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Per Codex review on PR #2699 (discussion_r3257683480): in the memory backend, joined.Mark fired after transition.WaitWithContext, so any joiner whose inherited transition failed bailed out of WaitWithContext and never got tagged. The failed-joiner traffic was counted as regular errors in telemetry — exactly the noise this change is meant to filter. Move joined.Mark in front of the wait when currentState == newState, so the request stays tagged regardless of whether the leader's transition succeeded. Mirrors the redis backend's handleExistingTransition which already marks before waitForTransition. --- .../internal/sandbox/storage/memory/operations.go | 12 ++++++++---- 1 file changed, 8 insertions(+), 4 deletions(-) diff --git a/packages/api/internal/sandbox/storage/memory/operations.go b/packages/api/internal/sandbox/storage/memory/operations.go index 4be36c5428..3a4eb2cadd 100644 --- a/packages/api/internal/sandbox/storage/memory/operations.go +++ b/packages/api/internal/sandbox/storage/memory/operations.go @@ -184,6 +184,14 @@ func startRemoving(ctx context.Context, sbx *memorySandbox, opts sandbox.RemoveO return false, nil, &sandbox.InvalidStateTransitionError{CurrentState: currentState, TargetState: newState} } + if currentState == newState { + // The caller will inherit the in-flight transition's result + // without doing the work itself: this is a joiner. Mark before + // waiting so the request stays tagged even if the inherited + // transition fails. + joined.Mark(ctx) + } + logger.L().Debug(ctx, "State transition already in progress to the same state, waiting", logger.WithSandboxID(sbx.SandboxID()), zap.String("state", string(newState))) err = transition.WaitWithContext(ctx) if err != nil { @@ -193,10 +201,6 @@ func startRemoving(ctx context.Context, sbx *memorySandbox, opts sandbox.RemoveO // If the transition is to the same state just wait switch { case currentState == newState: - // The caller inherits the in-flight transition's result - // without doing the work itself: this is a joiner. - joined.Mark(ctx) - return true, func(context.Context, error) {}, nil case sandbox.AllowedTransitions[currentState][newState]: return startRemoving(ctx, sbx, sandbox.RemoveOpts{Action: opts.Action})