mirror of
https://github.com/netbirdio/netbird.git
synced 2026-09-08 07:51:28 +02:00
[client] Make netbird up wait for the daemon to become ready
The CLI up path only tolerated a not-yet-ready daemon via a 10s blocking dial, so "netbird service start" immediately followed by "netbird up" (e.g. a container entrypoint) failed with a generic "daemon not running" error. The container entrypoint worked around this with a shell poll loop (status --check live) before running up. Move the readiness wait into the CLI, mirroring how the GUI already dials: - DialClientGRPCServer now uses grpc.NewClient with a tuned reconnect backoff and waits for the connection to reach READY (retrying on TRANSIENT_FAILURE) up to a 30s deadline, instead of grpc.DialContext + WithBlock with a hard 10s timeout. - up now polls Status via waitForDaemonStatus, retrying while the RPC is Unavailable (socket up but service not yet registered). Add an explicit daemon-ready signal so clients can wait deterministically instead of heuristically: - New optional StatusResponse.daemonReady field (field 5, wire-compatible with older GUIs/daemons which leave it unset). Regenerated with the pinned protoc v33.1 toolchain so no version churn leaks into the diff. - The server sets ready once Start succeeds and the DaemonService is registered (SetReady, called from the service controller). - waitForDaemonStatus waits for daemonReady=true (or an already-Connected status), with a bounded grace window so older daemons that never set the field are not blocked. Simplify the container entrypoint accordingly: drop the readiness poll loop (up now waits) and the now-dead NB_ENTRYPOINT_SERVICE_TIMEOUT env, keeping only the daemon+up process glue and SIGTERM forwarding for clean shutdown.
This commit is contained in:
176
client/cmd/dial_test.go
Normal file
176
client/cmd/dial_test.go
Normal file
@@ -0,0 +1,176 @@
|
||||
package cmd
|
||||
|
||||
import (
|
||||
"context"
|
||||
"net"
|
||||
"path/filepath"
|
||||
"sync/atomic"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
"google.golang.org/grpc"
|
||||
"google.golang.org/grpc/connectivity"
|
||||
|
||||
"github.com/netbirdio/netbird/client/internal"
|
||||
"github.com/netbirdio/netbird/client/proto"
|
||||
)
|
||||
|
||||
// startUnixGRPCServer starts a bare gRPC server listening on a unix socket at path
|
||||
// and returns a stop function. No services are registered; the connectivity-state
|
||||
// wait only cares about the transport becoming READY.
|
||||
func startUnixGRPCServer(t *testing.T, path string) func() {
|
||||
t.Helper()
|
||||
lis, err := net.Listen("unix", path)
|
||||
if err != nil {
|
||||
t.Fatalf("listen unix %s: %v", path, err)
|
||||
}
|
||||
srv := grpc.NewServer()
|
||||
go func() { _ = srv.Serve(lis) }()
|
||||
return srv.Stop
|
||||
}
|
||||
|
||||
func TestDialClientGRPCServer_ConnectsWhenServing(t *testing.T) {
|
||||
sock := filepath.Join(t.TempDir(), "nb.sock")
|
||||
stop := startUnixGRPCServer(t, sock)
|
||||
defer stop()
|
||||
|
||||
conn, err := dialClientGRPCServer(context.Background(), "unix://"+sock, 5*time.Second)
|
||||
if err != nil {
|
||||
t.Fatalf("expected connection, got error: %v", err)
|
||||
}
|
||||
defer conn.Close()
|
||||
|
||||
if state := conn.GetState(); state != connectivity.Ready {
|
||||
t.Fatalf("expected READY, got %s", state)
|
||||
}
|
||||
}
|
||||
|
||||
// TestDialClientGRPCServer_WaitsForLateServer is the core regression test: the
|
||||
// daemon socket appears only after the dial has already started, mirroring
|
||||
// "netbird service start" immediately followed by "netbird up".
|
||||
func TestDialClientGRPCServer_WaitsForLateServer(t *testing.T) {
|
||||
sock := filepath.Join(t.TempDir(), "nb.sock")
|
||||
|
||||
var stop func()
|
||||
timer := time.AfterFunc(1*time.Second, func() {
|
||||
stop = startUnixGRPCServer(t, sock)
|
||||
})
|
||||
defer timer.Stop()
|
||||
defer func() {
|
||||
if stop != nil {
|
||||
stop()
|
||||
}
|
||||
}()
|
||||
|
||||
start := time.Now()
|
||||
conn, err := dialClientGRPCServer(context.Background(), "unix://"+sock, 10*time.Second)
|
||||
if err != nil {
|
||||
t.Fatalf("expected connection after late server start, got error: %v", err)
|
||||
}
|
||||
defer conn.Close()
|
||||
|
||||
if elapsed := time.Since(start); elapsed < 500*time.Millisecond {
|
||||
t.Fatalf("connected too fast (%s); server should not have been up yet", elapsed)
|
||||
}
|
||||
if state := conn.GetState(); state != connectivity.Ready {
|
||||
t.Fatalf("expected READY, got %s", state)
|
||||
}
|
||||
}
|
||||
|
||||
// fakeStatusServer serves the Status RPC with a programmable response so we can
|
||||
// exercise waitForDaemonStatus without spinning up a real engine.
|
||||
type fakeStatusServer struct {
|
||||
proto.UnimplementedDaemonServiceServer
|
||||
resp func() *proto.StatusResponse
|
||||
}
|
||||
|
||||
func (f *fakeStatusServer) Status(context.Context, *proto.StatusRequest) (*proto.StatusResponse, error) {
|
||||
return f.resp(), nil
|
||||
}
|
||||
|
||||
func startFakeStatusServer(t *testing.T, sock string, resp func() *proto.StatusResponse) func() {
|
||||
t.Helper()
|
||||
lis, err := net.Listen("unix", sock)
|
||||
if err != nil {
|
||||
t.Fatalf("listen unix %s: %v", sock, err)
|
||||
}
|
||||
srv := grpc.NewServer()
|
||||
proto.RegisterDaemonServiceServer(srv, &fakeStatusServer{resp: resp})
|
||||
go func() { _ = srv.Serve(lis) }()
|
||||
return srv.Stop
|
||||
}
|
||||
|
||||
func dialFake(t *testing.T, sock string) proto.DaemonServiceClient {
|
||||
t.Helper()
|
||||
conn, err := dialClientGRPCServer(context.Background(), "unix://"+sock, 5*time.Second)
|
||||
if err != nil {
|
||||
t.Fatalf("dial fake daemon: %v", err)
|
||||
}
|
||||
t.Cleanup(func() { conn.Close() })
|
||||
return proto.NewDaemonServiceClient(conn)
|
||||
}
|
||||
|
||||
// New daemon that flips DaemonReady=true after a couple of polls: waitForDaemonStatus
|
||||
// must block until the flag is set, then return.
|
||||
func TestWaitForDaemonStatus_WaitsForDaemonReady(t *testing.T) {
|
||||
sock := filepath.Join(t.TempDir(), "nb.sock")
|
||||
var polls int32
|
||||
stop := startFakeStatusServer(t, sock, func() *proto.StatusResponse {
|
||||
n := atomic.AddInt32(&polls, 1)
|
||||
return &proto.StatusResponse{
|
||||
Status: string(internal.StatusConnecting),
|
||||
DaemonReady: n >= 3, // ready only from the 3rd poll on
|
||||
}
|
||||
})
|
||||
defer stop()
|
||||
|
||||
client := dialFake(t, sock)
|
||||
status, err := waitForDaemonStatus(context.Background(), client)
|
||||
if err != nil {
|
||||
t.Fatalf("unexpected error: %v", err)
|
||||
}
|
||||
if !status.GetDaemonReady() {
|
||||
t.Fatalf("expected DaemonReady=true, got false")
|
||||
}
|
||||
if got := atomic.LoadInt32(&polls); got < 3 {
|
||||
t.Fatalf("expected at least 3 polls before ready, got %d", got)
|
||||
}
|
||||
}
|
||||
|
||||
// Older daemon that never sets DaemonReady but reports a healthy (Connected)
|
||||
// status: waitForDaemonStatus must return promptly via the readiness fallback,
|
||||
// not block for the whole grace window.
|
||||
func TestWaitForDaemonStatus_OlderDaemonHealthyStatus(t *testing.T) {
|
||||
sock := filepath.Join(t.TempDir(), "nb.sock")
|
||||
stop := startFakeStatusServer(t, sock, func() *proto.StatusResponse {
|
||||
return &proto.StatusResponse{Status: string(internal.StatusConnected)} // DaemonReady unset
|
||||
})
|
||||
defer stop()
|
||||
|
||||
client := dialFake(t, sock)
|
||||
start := time.Now()
|
||||
status, err := waitForDaemonStatus(context.Background(), client)
|
||||
if err != nil {
|
||||
t.Fatalf("unexpected error: %v", err)
|
||||
}
|
||||
if status.GetDaemonReady() {
|
||||
t.Fatalf("expected DaemonReady=false from older daemon")
|
||||
}
|
||||
if elapsed := time.Since(start); elapsed > 2*time.Second {
|
||||
t.Fatalf("returned too slowly (%s); healthy status should short-circuit the grace", elapsed)
|
||||
}
|
||||
}
|
||||
|
||||
func TestDialClientGRPCServer_TimesOutWhenAbsent(t *testing.T) {
|
||||
sock := filepath.Join(t.TempDir(), "never.sock")
|
||||
|
||||
start := time.Now()
|
||||
conn, err := dialClientGRPCServer(context.Background(), "unix://"+sock, 1*time.Second)
|
||||
if err == nil {
|
||||
conn.Close()
|
||||
t.Fatal("expected timeout error, got nil")
|
||||
}
|
||||
if elapsed := time.Since(start); elapsed < 900*time.Millisecond {
|
||||
t.Fatalf("returned too early (%s); should have waited ~timeout", elapsed)
|
||||
}
|
||||
}
|
||||
@@ -20,6 +20,8 @@ import (
|
||||
"github.com/spf13/cobra"
|
||||
"github.com/spf13/pflag"
|
||||
"google.golang.org/grpc"
|
||||
gbackoff "google.golang.org/grpc/backoff"
|
||||
"google.golang.org/grpc/connectivity"
|
||||
"google.golang.org/grpc/credentials/insecure"
|
||||
|
||||
daddr "github.com/netbirdio/netbird/client/internal/daemonaddr"
|
||||
@@ -264,17 +266,70 @@ func FlagNameToEnvVar(cmdFlag string, prefix string) string {
|
||||
return prefix + upper
|
||||
}
|
||||
|
||||
// DialClientGRPCServer returns client connection to the daemon server.
|
||||
func DialClientGRPCServer(ctx context.Context, addr string) (*grpc.ClientConn, error) {
|
||||
ctx, cancel := context.WithTimeout(ctx, time.Second*10)
|
||||
defer cancel()
|
||||
// defaultDaemonDialTimeout is how long DialClientGRPCServer waits for the daemon
|
||||
// to become reachable. It is intentionally generous so that invoking the CLI
|
||||
// right after "netbird service start" (e.g. from a container entrypoint) tolerates
|
||||
// the window where the daemon has created its socket but is not yet serving.
|
||||
const defaultDaemonDialTimeout = 30 * time.Second
|
||||
|
||||
return grpc.DialContext(
|
||||
ctx,
|
||||
// DialClientGRPCServer returns a client connection to the daemon server. It waits
|
||||
// for the daemon to become reachable, retrying with backoff until the connection
|
||||
// reports READY or defaultDaemonDialTimeout elapses. This handles the startup race
|
||||
// where the daemon socket exists (or is about to) before the gRPC server is serving.
|
||||
func DialClientGRPCServer(ctx context.Context, addr string) (*grpc.ClientConn, error) {
|
||||
return dialClientGRPCServer(ctx, addr, defaultDaemonDialTimeout)
|
||||
}
|
||||
|
||||
func dialClientGRPCServer(ctx context.Context, addr string, timeout time.Duration) (*grpc.ClientConn, error) {
|
||||
conn, err := grpc.NewClient(
|
||||
strings.TrimPrefix(addr, "tcp://"),
|
||||
grpc.WithTransportCredentials(insecure.NewCredentials()),
|
||||
grpc.WithBlock(),
|
||||
// Cap reconnect backoff at 5s; gRPC's default 120s MaxDelay would leave the
|
||||
// CLI waiting far too long to notice a freshly-started daemon. Mirrors the GUI.
|
||||
grpc.WithConnectParams(grpc.ConnectParams{
|
||||
Backoff: gbackoff.Config{
|
||||
BaseDelay: 1 * time.Second,
|
||||
Multiplier: 1.6,
|
||||
Jitter: 0.2,
|
||||
MaxDelay: 5 * time.Second,
|
||||
},
|
||||
}),
|
||||
)
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("create daemon gRPC client: %w", err)
|
||||
}
|
||||
|
||||
// grpc.NewClient is lazy: it does not connect until the first RPC or until we
|
||||
// nudge it. Trigger connection attempts and wait until the channel reaches READY.
|
||||
if err := waitForConnReady(ctx, conn, timeout); err != nil {
|
||||
_ = conn.Close()
|
||||
return nil, err
|
||||
}
|
||||
return conn, nil
|
||||
}
|
||||
|
||||
// waitForConnReady drives the gRPC channel out of IDLE and blocks until it becomes
|
||||
// READY, or until timeout/ctx expires. TRANSIENT_FAILURE (daemon not yet serving)
|
||||
// is treated as retryable so the caller keeps waiting within the deadline.
|
||||
func waitForConnReady(ctx context.Context, conn *grpc.ClientConn, timeout time.Duration) error {
|
||||
ctx, cancel := context.WithTimeout(ctx, timeout)
|
||||
defer cancel()
|
||||
|
||||
for {
|
||||
state := conn.GetState()
|
||||
switch state {
|
||||
case connectivity.Ready:
|
||||
return nil
|
||||
case connectivity.Idle:
|
||||
// Kick the lazy channel into connecting.
|
||||
conn.Connect()
|
||||
}
|
||||
|
||||
if !conn.WaitForStateChange(ctx, state) {
|
||||
// ctx expired while in `state`.
|
||||
return fmt.Errorf("timed out after %s waiting for daemon to become ready (last state: %s)", timeout, state)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// WithBackOff execute function in backoff cycle.
|
||||
|
||||
@@ -78,6 +78,10 @@ func (p *program) Start(svc service.Service) error {
|
||||
log.Fatalf("failed to start daemon: %v", err)
|
||||
}
|
||||
proto.RegisterDaemonServiceServer(p.serv, serverInstance)
|
||||
// The engine is started and the service is registered: from here on the
|
||||
// daemon serves RPCs backed by a running engine. Report readiness so
|
||||
// clients (e.g. netbird up) can wait deterministically.
|
||||
serverInstance.SetReady()
|
||||
|
||||
p.serverInstanceMu.Lock()
|
||||
p.serverInstance = serverInstance
|
||||
|
||||
@@ -295,9 +295,7 @@ func runInDaemonMode(ctx context.Context, cmd *cobra.Command, pm *profilemanager
|
||||
|
||||
client := proto.NewDaemonServiceClient(conn)
|
||||
|
||||
status, err := client.Status(ctx, &proto.StatusRequest{
|
||||
WaitForReady: func() *bool { b := true; return &b }(),
|
||||
})
|
||||
status, err := waitForDaemonStatus(ctx, client)
|
||||
if err != nil {
|
||||
return fmt.Errorf("unable to get daemon status: %v", err)
|
||||
}
|
||||
@@ -336,6 +334,79 @@ func runInDaemonMode(ctx context.Context, cmd *cobra.Command, pm *profilemanager
|
||||
return nil
|
||||
}
|
||||
|
||||
// daemonStatusPollTimeout bounds how long we poll the Status RPC waiting for the
|
||||
// daemon to answer coherently. The transport is already READY at this point (see
|
||||
// DialClientGRPCServer), so this only covers the brief window where the gRPC server
|
||||
// is serving but the daemon engine is still starting up and the Status RPC races
|
||||
// against server.Start().
|
||||
const daemonStatusPollTimeout = 15 * time.Second
|
||||
|
||||
// daemonReadyGrace bounds how long we keep polling once the daemon answers but
|
||||
// still reports DaemonReady=false. A freshly-started daemon flips it to true
|
||||
// within this window; an older daemon that never sets the field simply falls
|
||||
// through after the grace elapses, preserving backward compatibility.
|
||||
const daemonReadyGrace = 10 * time.Second
|
||||
|
||||
// waitForDaemonStatus fetches the daemon status, waiting for the daemon to become
|
||||
// ready. It handles two startup races:
|
||||
//
|
||||
// 1. The gRPC server is not yet serving: Status fails with Unavailable — retry.
|
||||
// 2. The server serves but the engine is still starting: a DaemonReady-aware
|
||||
// daemon reports DaemonReady=false until Start finishes; poll until it flips
|
||||
// true (bounded by daemonReadyGrace). Older daemons never set DaemonReady, so
|
||||
// we stop waiting on it after the grace and use the status as-is.
|
||||
//
|
||||
// It gives up after daemonStatusPollTimeout.
|
||||
func waitForDaemonStatus(ctx context.Context, client proto.DaemonServiceClient) (*proto.StatusResponse, error) {
|
||||
ctx, cancel := context.WithTimeout(ctx, daemonStatusPollTimeout)
|
||||
defer cancel()
|
||||
|
||||
waitForReady := true
|
||||
req := &proto.StatusRequest{WaitForReady: &waitForReady}
|
||||
|
||||
var lastErr error
|
||||
var firstAnswer time.Time
|
||||
for {
|
||||
status, err := client.Status(ctx, req)
|
||||
if err != nil {
|
||||
lastErr = err
|
||||
// Only retry while the daemon is not yet answering; surface real errors.
|
||||
if s, ok := gstatus.FromError(err); !ok || s.Code() != codes.Unavailable {
|
||||
return nil, err
|
||||
}
|
||||
} else {
|
||||
// Daemon answered. Explicitly ready (DaemonReady-aware daemon), or
|
||||
// already fully connected — either way, done. Connected is the only
|
||||
// status unambiguous enough to short-circuit on: a DaemonReady-aware
|
||||
// daemon sets the flag at startup, so trusting Connected here can only
|
||||
// help an older daemon that never sets the flag, without overriding a
|
||||
// new daemon that legitimately reports DaemonReady=false while starting.
|
||||
if status.GetDaemonReady() || internal.StatusType(status.GetStatus()) == internal.StatusConnected {
|
||||
return status, nil
|
||||
}
|
||||
// Answered but neither ready-flagged nor connected yet: give a
|
||||
// DaemonReady-aware daemon a bounded window to finish starting, then
|
||||
// fall through so an older daemon that never sets the flag isn't
|
||||
// blocked here.
|
||||
if firstAnswer.IsZero() {
|
||||
firstAnswer = time.Now()
|
||||
} else if time.Since(firstAnswer) >= daemonReadyGrace {
|
||||
return status, nil
|
||||
}
|
||||
lastErr = nil
|
||||
}
|
||||
|
||||
select {
|
||||
case <-ctx.Done():
|
||||
if lastErr != nil {
|
||||
return nil, fmt.Errorf("daemon did not become ready: %w", lastErr)
|
||||
}
|
||||
return nil, fmt.Errorf("daemon did not become ready within %s", daemonStatusPollTimeout)
|
||||
case <-time.After(500 * time.Millisecond):
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func doDaemonUp(ctx context.Context, cmd *cobra.Command, client proto.DaemonServiceClient, pm *profilemanager.ProfileManager, activeProf *profilemanager.Profile, customDNSAddressConverted []byte, username string) error {
|
||||
|
||||
providedSetupKey, err := getSetupKey()
|
||||
|
||||
Reference in New Issue
Block a user