Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
16 changes: 4 additions & 12 deletions packages/envd/internal/services/process/start.go
Original file line number Diff line number Diff line change
Expand Up @@ -61,11 +61,9 @@ func (s *Service) handleStart(ctx context.Context, req *connect.Request[rpc.Star

exitChan := make(chan struct{})

startMultiplexer := handler.NewMultiplexedChannel[rpc.ProcessEvent_Start](0)
defer close(startMultiplexer.Source)

start, startCancel := startMultiplexer.Fork()
defer startCancel()
// Buffered so the send below never blocks when the receiver
// goroutine has already exited on a cancelled context.
start := make(chan rpc.ProcessEvent_Start, 1)

data, dataCancel := proc.DataEvent.Fork()
defer dataCancel()
Expand All @@ -81,13 +79,7 @@ func (s *Service) handleStart(ctx context.Context, req *connect.Request[rpc.Star
cancel(ctx.Err())

return
case event, ok := <-start:
if !ok {
cancel(connect.NewError(connect.CodeUnknown, errors.New("start event channel closed before sending start event")))

return
}

case event := <-start:
streamErr := stream.Send(&rpc.StartResponse{
Event: &rpc.ProcessEvent{
Event: &event,
Expand Down
64 changes: 62 additions & 2 deletions packages/envd/internal/services/process/start_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -21,7 +21,7 @@ import (
"github.com/e2b-dev/infra/packages/envd/internal/utils"
)

func newTestService(t *testing.T) (spec.ProcessClient, func()) {
func newTestService(t *testing.T, middleware ...func(http.Handler) http.Handler) (spec.ProcessClient, func()) {
t.Helper()

// handler.New sets SysProcAttr.Credential to switch uid/gid,
Expand All @@ -46,7 +46,12 @@ func newTestService(t *testing.T) (spec.ProcessClient, func()) {
path, handler := spec.NewProcessHandler(svc)
mux.Handle(path, handler)

srv := httptest.NewServer(mux)
var h http.Handler = mux
for _, mw := range middleware {
h = mw(h)
}
Comment thread
arkamar marked this conversation as resolved.

srv := httptest.NewServer(h)
client := spec.NewProcessClient(srv.Client(), srv.URL)

return client, srv.Close
Expand Down Expand Up @@ -84,6 +89,61 @@ func TestStart_ShortCommand(t *testing.T) {
assert.NotNil(t, events[len(events)-1].GetEnd(), "last event should be End")
}

// TestStart_CancelBeforeStartEvent verifies that the handler returns instead
// of blocking forever when the request context is cancelled before the
// bootstrap start event reaches the sender goroutine.
//
// The middleware cancels the request context at handler entry, exercising the
// ordering where the sender goroutine exits before handleStart emits the start
// event. The process itself still starts because it runs on an independent
// context. A stuck handler keeps its connection open and makes the server
// close at the end of the test hang.
func TestStart_CancelBeforeStartEvent(t *testing.T) {
t.Parallel()

cancelAtEntry := func(next http.Handler) http.Handler {
return http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
ctx, cancel := context.WithCancel(r.Context())
cancel()

next.ServeHTTP(w, r.WithContext(ctx))
})
}

client, cleanup := newTestService(t, cancelAtEntry)

// Bound the drain below in case the handler never responds.
ctx, cancel := context.WithTimeout(t.Context(), 5*time.Second)
defer cancel()

stream, err := client.Start(ctx, connect.NewRequest(&rpc.StartRequest{
Process: &rpc.ProcessConfig{
Cmd: "echo",
Args: []string{"hello"},
},
}))
if err == nil {
for stream.Receive() {
}
_ = stream.Close()
}

// Server close waits for outstanding handlers, so a wedged
// handler makes it hang.
closed := make(chan struct{})
go func() {
defer close(closed)

cleanup()
}()

select {
case <-closed:
case <-time.After(20 * time.Second):
t.Fatal("server close timed out: handler goroutine leaked on cancelled start")
}
}

// TestStart_ClientDisconnectMidStream verifies that when a client
// cancels mid-stream, the handler returns without racing. This is
// the scenario that caused the nil-pointer panic in production:
Expand Down
2 changes: 1 addition & 1 deletion packages/envd/pkg/version.go
Original file line number Diff line number Diff line change
@@ -1,3 +1,3 @@
package pkg

const Version = "0.6.9"
const Version = "0.6.10"
Loading