Skip to content

Commit 880e5e7

Browse files
adababysclaude
andcommitted
fix(orchestrator): early-bind gRPC port + serve /health during startup reclaim
On a host with many leaked ns-* in /run/netns (e.g. after draining a node that had run many sandboxes), startup reclaim tears down namespaces serially and runs before the listener binds, so GRPC_PORT stays unbound for 20-30 min. Nomad health checks get connection-refused and the node reads as failed. Bind the cmux listener and start the HTTP /health server before the long sandbox-runtime init. serviceInfo now starts Unhealthy, so /health returns 503 (Nomad sees "starting", not "failed") until the gRPC server is wired up, at which point the status flips to Healthy. The gRPC listener is matched at early-bind but only Served after all RegisterService calls complete; early gRPC connections buffer in the cmux matcher. The late /upload handler is wired through an atomic pointer so the mux stays immutable once serving. Refs #3615 Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
1 parent 16bd4e3 commit 880e5e7

2 files changed

Lines changed: 90 additions & 60 deletions

File tree

‎packages/orchestrator/pkg/factories/run.go‎

Lines changed: 84 additions & 59 deletions
Original file line numberDiff line numberDiff line change
@@ -15,6 +15,7 @@ import (
1515
"slices"
1616
"strconv"
1717
"strings"
18+
"sync/atomic"
1819
"syscall"
1920
"time"
2021

@@ -539,6 +540,76 @@ func run(config cfg.Config, opts Options) (success bool) {
539540

540541
var closers []closer
541542

543+
// Early-bind the service port and start the HTTP /health endpoint BEFORE the
544+
// potentially long sandbox-runtime init below (startup reclaim can take tens
545+
// of minutes on a host with many leaked ns-* in /run/netns). serviceInfo
546+
// starts Unhealthy, so /health returns 503 until the gRPC server is wired up
547+
// and we flip it to Healthy — the port is open (Nomad sees "starting", not a
548+
// connection-refused "failed") but no new work is routed here yet.
549+
//
550+
// cmux matchers must be created before Serve() (Match mutates state Serve
551+
// reads). The gRPC listener is matched here but only served later, after all
552+
// RegisterService calls complete — gRPC forbids RegisterService after Serve,
553+
// so early gRPC connections buffer in the matcher until then.
554+
cmuxServer, err := NewCMUXServer(ctx, config.GRPCPort, tel.MeterProvider)
555+
if err != nil {
556+
logger.L().Fatal(ctx, "failed to create cmux server", zap.Error(err))
557+
}
558+
httpListener := cmuxServer.Match(cmux.HTTP1Fast())
559+
grpcListener := cmuxServer.Match(cmux.Any()) // the rest are GRPC requests
560+
561+
startService("cmux server", func() error {
562+
logger.L().Info(ctx, "Starting network server", zap.Uint16("port", config.GRPCPort))
563+
err := cmuxServer.Serve()
564+
if err != nil && strings.Contains(err.Error(), "use of closed network connection") {
565+
return nil
566+
}
567+
568+
return err
569+
})
570+
closers = append(closers, closer{"cmux server", func(context.Context) error {
571+
logger.L().Info(ctx, "Shutting down cmux server")
572+
cmuxServer.Close()
573+
574+
return nil
575+
}})
576+
577+
healthcheck, err := e2bhealthcheck.NewHealthcheck(serviceInfo)
578+
if err != nil {
579+
logger.L().Fatal(ctx, "failed to create healthcheck", zap.Error(err))
580+
}
581+
582+
// The /upload handler (local build storage) is created later during template
583+
// manager setup. Register a stable indirection now so the mux is immutable
584+
// once Serve starts; it 503s until the real handler is stored.
585+
var uploadHandlerPtr atomic.Pointer[localupload.Handler]
586+
httpMux := http.NewServeMux()
587+
httpMux.Handle("/health", healthcheck.CreateHandler())
588+
httpMux.HandleFunc("/upload", func(w http.ResponseWriter, r *http.Request) {
589+
if h := uploadHandlerPtr.Load(); h != nil {
590+
h.ServeHTTP(w, r)
591+
592+
return
593+
}
594+
595+
http.Error(w, "upload handler not ready", http.StatusServiceUnavailable)
596+
})
597+
598+
httpServer := NewHTTPServer()
599+
httpServer.Handler = httpMux
600+
startService("http server", func() error {
601+
err := httpServer.Serve(httpListener)
602+
switch {
603+
case errors.Is(err, cmux.ErrServerClosed):
604+
return nil
605+
case errors.Is(err, http.ErrServerClosed):
606+
return nil
607+
default:
608+
return err
609+
}
610+
})
611+
closers = append(closers, closer{"http server", httpServer.Shutdown})
612+
542613
// The sandbox map is shared between the server and the proxy
543614
// to propagate information about sandbox routing.
544615
sandboxes := sandbox.NewSandboxesMap()
@@ -932,14 +1003,17 @@ func run(config cfg.Config, opts Options) (success bool) {
9321003

9331004
// template manager
9341005
var tmpl *tmplserver.ServerStore
935-
var localUploadHandler *localupload.Handler
9361006
if services.RunsTemplateManager() {
9371007
buildPersistence, uploadHandler, err := setupBuildStorage(ctx, limiter, config)
9381008
if err != nil {
9391009
logger.L().Fatal(ctx, "failed to setup build storage", zap.Error(err))
9401010
}
9411011

942-
localUploadHandler = uploadHandler
1012+
// Publish the real /upload handler to the indirection registered on the
1013+
// early HTTP mux (nil until now → served 503).
1014+
if uploadHandler != nil {
1015+
uploadHandlerPtr.Store(uploadHandler)
1016+
}
9431017

9441018
tmpl, err = tmplserver.New(
9451019
ctx,
@@ -970,33 +1044,6 @@ func run(config cfg.Config, opts Options) (success bool) {
9701044
grpcHealth := health.NewServer()
9711045
grpc_health_v1.RegisterHealthServer(grpcServer, grpcHealth)
9721046

973-
// cmux server, allows us to reuse the same TCP port between grpc and HTTP requests
974-
cmuxServer, err := NewCMUXServer(ctx, config.GRPCPort, tel.MeterProvider)
975-
if err != nil {
976-
logger.L().Fatal(ctx, "failed to create cmux server", zap.Error(err))
977-
}
978-
979-
// Create all matchers BEFORE starting Serve() to avoid data race.
980-
// cmux.Match() modifies internal state that Serve() reads from.
981-
httpListener := cmuxServer.Match(cmux.HTTP1Fast())
982-
grpcListener := cmuxServer.Match(cmux.Any()) // the rest are GRPC requests
983-
984-
startService("cmux server", func() error {
985-
logger.L().Info(ctx, "Starting network server", zap.Uint16("port", config.GRPCPort))
986-
err := cmuxServer.Serve()
987-
if err != nil && strings.Contains(err.Error(), "use of closed network connection") {
988-
return nil
989-
}
990-
991-
return err
992-
})
993-
closers = append(closers, closer{"cmux server", func(context.Context) error {
994-
logger.L().Info(ctx, "Shutting down cmux server")
995-
cmuxServer.Close()
996-
997-
return nil
998-
}})
999-
10001047
pprofServer := telemetry.NewPprofServer()
10011048
// We handle the pprof in a separate goroutine to prevent any interaction with the main server.
10021049
go func() {
@@ -1008,36 +1055,9 @@ func run(config cfg.Config, opts Options) (success bool) {
10081055
}()
10091056
closers = append(closers, closer{"pprof server", pprofServer.Shutdown})
10101057

1011-
// http server
1012-
healthcheck, err := e2bhealthcheck.NewHealthcheck(serviceInfo)
1013-
if err != nil {
1014-
logger.L().Fatal(ctx, "failed to create healthcheck", zap.Error(err))
1015-
}
1016-
1017-
httpMux := http.NewServeMux()
1018-
httpMux.Handle("/health", healthcheck.CreateHandler())
1019-
1020-
if localUploadHandler != nil {
1021-
httpMux.Handle("/upload", localUploadHandler)
1022-
}
1023-
1024-
httpServer := NewHTTPServer()
1025-
httpServer.Handler = httpMux
1026-
1027-
startService("http server", func() error {
1028-
err := httpServer.Serve(httpListener)
1029-
switch {
1030-
case errors.Is(err, cmux.ErrServerClosed):
1031-
return nil
1032-
case errors.Is(err, http.ErrServerClosed):
1033-
return nil
1034-
default:
1035-
return err
1036-
}
1037-
})
1038-
closers = append(closers, closer{"http server", httpServer.Shutdown})
1039-
1040-
// grpc server
1058+
// grpc server. All RegisterService calls above are complete, so it is now
1059+
// safe to Serve the grpc listener that cmux has been matching since early
1060+
// bind (any connections received meanwhile are buffered in the matcher).
10411061
startService("grpc server", func() error {
10421062
return grpcServer.Serve(grpcListener)
10431063
})
@@ -1048,6 +1068,11 @@ func run(config cfg.Config, opts Options) (success bool) {
10481068
return nil
10491069
}})
10501070

1071+
// Sandbox runtime is fully wired up and the gRPC server is serving. Flip to
1072+
// Healthy so /health returns 200 and the edge starts routing new work here.
1073+
// Until this point the port was open but reported Unhealthy (503).
1074+
serviceInfo.SetStatus(ctx, orchestratorinfo.ServiceInfoStatus_Healthy)
1075+
10511076
// Wait for the shutdown signal or if some service fails
10521077
select {
10531078
case <-sig.Done():

‎packages/orchestrator/pkg/service/info.go‎

Lines changed: 6 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -91,7 +91,12 @@ func NewInfoContainer(clientId string, version string, commit string, instanceID
9191
ClientId: clientId,
9292
ServiceId: instanceID,
9393

94-
status: ServiceStatus{Status: orchestratorinfo.ServiceInfoStatus_Healthy, ChangedAt: startup},
94+
// Start Unhealthy: the /health endpoint and gRPC port may bind before the
95+
// sandbox runtime finishes initializing (startup reclaim, network pool
96+
// populate). Reporting Unhealthy until the gRPC server is actually serving
97+
// keeps the edge/service discovery from routing new work to a node that is
98+
// still coming up; run.go flips this to Healthy once wiring completes.
99+
status: ServiceStatus{Status: orchestratorinfo.ServiceInfoStatus_Unhealthy, ChangedAt: startup},
95100

96101
Startup: startup,
97102
Roles: serviceRoles,

0 commit comments

Comments
 (0)