diff --git a/cmd/deploy/deploy.go b/cmd/deploy/deploy.go index f55ad53fe..3a96c31c9 100644 --- a/cmd/deploy/deploy.go +++ b/cmd/deploy/deploy.go @@ -110,6 +110,7 @@ func deployFn(cmd *cobra.Command, args []string) error { return err } d.Progress = viper.GetBool("progress") + d.Timeout = viper.GetDuration("deploy_timeout") d.ReportUsage = viper.GetBool("report_usage") d.ReportUsageProjectID = viper.GetString("report_usage_project_id") d.ReportUsageTopicID = viper.GetString("report_usage_topic_id") diff --git a/cmd/root.go b/cmd/root.go index 7827251a8..f26a58cc9 100644 --- a/cmd/root.go +++ b/cmd/root.go @@ -45,6 +45,7 @@ environment.`, root.PersistentFlags().String("report_usage_project_id", "", "Project to report anonymous usage metrics to") root.PersistentFlags().String("report_usage_topic_id", "", "Topic to report anonymous usage metrics to") root.PersistentFlags().Bool("progress", false, "Display progress of container bringup") + root.PersistentFlags().Duration("deploy_timeout", 0, "Overall timeout for deploying the ingress, CNI and controllers, 0 for the default") root.PersistentPreRunE = func(cmd *cobra.Command, args []string) error { if *cfgFile == "" { return nil diff --git a/deploy/deploy.go b/deploy/deploy.go index 540741583..5f8f88331 100644 --- a/deploy/deploy.go +++ b/deploy/deploy.go @@ -5,6 +5,7 @@ import ( "crypto/rand" "encoding/base64" "encoding/json" + "errors" "fmt" "net" "os" @@ -28,6 +29,7 @@ import ( epb "github.com/openconfig/kne/proto/event" "github.com/pborman/uuid" metallbv1 "go.universe.tf/metallb/api/v1beta1" + "golang.org/x/sync/errgroup" appsv1 "k8s.io/api/apps/v1" corev1 "k8s.io/api/core/v1" metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" @@ -51,7 +53,23 @@ var ( setPIDMaxScript = filepath.Join(homedir.HomeDir(), "kne-internal", "set_pid_max.sh") pullRetryDelay = time.Second poolRetryDelay = 5 * time.Second - healthTimeout = time.Minute + // healthTimeout is how long a single component is given to become + // healthy, measured from the moment that component finished deploying. + // + // This is deliberately generous. Operator readiness is mostly a fixed + // cost rather than noise, and reproduces closely between runs: MetalLB + // has been measured at 69.15s and 69.25s on two separate deployments, + // with Lemming at ~45s and IxiaTG at ~44s. MetalLB therefore never fit + // in a one minute budget; it simply was never held to one, because the + // wait that matters happens inside its Deploy rather than in Healthy. + // defaultDeployTimeout, not this, is the backstop against a deployment + // that is truly stuck. + healthTimeout = 3 * time.Minute + // defaultDeployTimeout bounds the concurrent deployment of the ingress, + // CNI and controllers as a whole. It is a backstop for the cases the + // per-component budgets cannot catch, such as a Deploy that hangs: those + // run kubectl, which has no timeout of its own. + defaultDeployTimeout = 10 * time.Minute // Stubs for testing. execLookPath = exec.LookPath @@ -88,6 +106,15 @@ type Controller interface { Healthy(context.Context) error } +// A component is the part of Ingress, CNI and Controller that a deployment +// drives: bring yourself up, then report when you are ready. Ingress, CNI and +// Controller all satisfy it, which lets them be deployed uniformly and +// concurrently. +type component interface { + Deploy(context.Context) error + Healthy(context.Context) error +} + type Deployment struct { Cluster Cluster `kne:"cluster"` Ingress Ingress `kne:"ingress"` @@ -98,6 +125,12 @@ type Deployment struct { // standard output. Progress bool + // Timeout bounds the concurrent deployment of the ingress, CNI and + // controllers. It does not cover bringing up the cluster itself, which + // happens first and is largely not interruptible. If zero, + // defaultDeployTimeout is used. + Timeout time.Duration + // If ReportUsage is true then anonymous usage metrics will be // published using Cloud PubSub. ReportUsage bool @@ -255,47 +288,123 @@ func (d *Deployment) Deploy(ctx context.Context, kubecfg string) (rerr error) { }() } + // Inject the clients into every component up front. Ingress needs its + // client during Deploy, and the others only happen to get away without + // one because their Deploy is a plain kubectl apply today. d.Ingress.SetKClient(kClient) d.Ingress.SetRCfg(rCfg) d.Ingress.SetDockerNetworkResourceName(d.Cluster.GetDockerNetworkResourceName()) - - log.Infof("Deploying ingress...") - if err := d.Ingress.Deploy(ctx); err != nil { - return fmt.Errorf("failed to deploy ingress: %w", err) - } - tCtx, cancel := context.WithTimeout(ctx, healthTimeout) - defer cancel() - if err := d.Ingress.Healthy(tCtx); err != nil { - return fmt.Errorf("failed to check if ingress is healthy: %w", err) - } - log.Infof("Ingress healthy") - log.Infof("Deploying CNI...") - if err := d.CNI.Deploy(ctx); err != nil { - return fmt.Errorf("failed to deploy CNI: %w", err) - } d.CNI.SetKClient(kClient) - tCtx, cancel = context.WithTimeout(ctx, healthTimeout) - defer cancel() - if err := d.CNI.Healthy(tCtx); err != nil { - return fmt.Errorf("failed to check if CNI is healthy: %w", err) - } - log.Infof("CNI healthy") for _, c := range d.Controllers { - log.Infof("Deploying controller...") - if err := c.Deploy(ctx); err != nil { - return fmt.Errorf("failed to deploy controller: %w", err) - } c.SetKClient(kClient) - tCtx, cancel = context.WithTimeout(ctx, healthTimeout) - defer cancel() - if err := c.Healthy(tCtx); err != nil { - return fmt.Errorf("failed to check if controller is healthy: %w", err) - } } - log.Infof("Controllers deployed and healthy") + + if err := d.deployComponents(ctx); err != nil { + return err + } + log.Infof("Ingress, CNI and controllers deployed and healthy") return nil } +// deployComponents brings up the ingress, CNI and controllers. +// +// They are independent of each other: each applies its own manifests and then +// waits for its own workloads in its own namespace. Deploying them +// concurrently makes a deployment cost the slowest component rather than the +// sum of all of them. +func (d *Deployment) deployComponents(ctx context.Context) error { + dCtx, dCancel := context.WithTimeout(ctx, d.timeout()) + defer dCancel() + g, gCtx := errgroup.WithContext(dCtx) + + start := func(name string, c component) { + // A component that was never configured has nothing to deploy. The + // sequential code often never reached a nil component because an + // earlier one failed first; running them together always reaches it, + // so say so rather than panicking. + if c == nil { + log.Warningf("No %s configured, skipping", name) + return + } + g.Go(func() error { + // Tag any command this component runs, so its output can be + // picked out of the interleaved log. + cCtx := run.WithLabel(gCtx, name) + log.Infof("Deploying %s...", name) + if err := c.Deploy(cCtx); err != nil { + return componentErr(name, dCtx, gCtx, nil, fmt.Errorf("failed to deploy %s: %w", name, err)) + } + log.Infof("%s deployed", name) + + // Each component gets its own health budget, started when that + // component finished deploying. A single shared deadline would + // charge a component for a slow sibling, and would expire for + // everything still pending at once, hiding which component was + // actually stuck. + hCtx, hCancel := context.WithTimeout(gCtx, healthTimeout) + defer hCancel() + if err := c.Healthy(hCtx); err != nil { + return componentErr(name, dCtx, gCtx, hCtx, fmt.Errorf("failed to check if %s is healthy: %w", name, err)) + } + log.Infof("%s healthy", name) + return nil + }) + } + + start("ingress", d.Ingress) + start("CNI", d.CNI) + for _, c := range d.Controllers { + start(controllerName(c), c) + } + return g.Wait() +} + +// timeout returns the budget for the deployment as a whole. +func (d *Deployment) timeout() time.Duration { + if d.Timeout > 0 { + return d.Timeout + } + return defaultDeployTimeout +} + +// controllerName returns a name for c suitable for logs and errors, derived +// from its type: a *CEOSLabSpec becomes "CEOSLab controller". Controllers are +// otherwise anonymous, and with several of them coming up at once "controller" +// alone does not say which one failed. +func controllerName(c Controller) string { + n := fmt.Sprintf("%T", c) + n = n[strings.LastIndex(n, ".")+1:] + if n = strings.TrimSuffix(n, "Spec"); n == "" { + return "controller" + } + return n + " controller" +} + +// componentErr explains why a component stopped early. The components share a +// context, so the first failure cancels the rest; without this every other +// component would report an indistinguishable "context canceled" and bury the +// one real error. +// +// overall is the deployment-wide context, group the errgroup context, and own +// the component's own health context, which is nil while it is still +// deploying. +func componentErr(name string, overall, group, own context.Context, err error) error { + switch { + case errors.Is(overall.Err(), context.DeadlineExceeded): + // The whole deployment ran out of time, so this component's own + // budget is not the interesting fact. + return fmt.Errorf("%s did not finish before the overall deployment timeout expired: %w", name, err) + case own != nil && errors.Is(own.Err(), context.DeadlineExceeded): + return fmt.Errorf("%s was not healthy within %v: %w", name, healthTimeout, err) + case group.Err() != nil: + // Either a sibling failed or the deployment was canceled, e.g. by the + // pod watcher seeing a container fail to start. + return fmt.Errorf("%s abandoned before completing: %w", name, err) + default: + return err + } +} + func validateKubectlVersion() error { output, err := run.OutCommand("kubectl", "version", "--output=yaml") if err != nil { @@ -418,27 +527,35 @@ func (d *Deployment) Healthy(ctx context.Context) error { return fmt.Errorf("failed to check cluster is healthy: %w", err) } log.Infof("Cluster healthy") - tCtx, cancel := context.WithTimeout(ctx, healthTimeout) - defer cancel() - if err := d.Ingress.Healthy(tCtx); err != nil { - return fmt.Errorf("failed to check ingress is healthy: %w", err) - } - log.Infof("Ingress healthy") - tCtx, cancel = context.WithTimeout(ctx, healthTimeout) - defer cancel() - if err := d.CNI.Healthy(tCtx); err != nil { - return fmt.Errorf("failed to check CNI is healthy: %w", err) - } - log.Infof("CNI healthy") + + // As in Deploy, the components are independent, so check them + // concurrently and give each its own budget. + hCtx, hCancel := context.WithTimeout(ctx, d.timeout()) + defer hCancel() + g, gCtx := errgroup.WithContext(hCtx) + + check := func(name string, c component) { + if c == nil { + log.Warningf("No %s configured, skipping health check", name) + return + } + g.Go(func() error { + cCtx, cancel := context.WithTimeout(gCtx, healthTimeout) + defer cancel() + if err := c.Healthy(cCtx); err != nil { + return componentErr(name, hCtx, gCtx, cCtx, fmt.Errorf("failed to check %s is healthy: %w", name, err)) + } + log.Infof("%s healthy", name) + return nil + }) + } + + check("ingress", d.Ingress) + check("CNI", d.CNI) for _, c := range d.Controllers { - tCtx, cancel = context.WithTimeout(ctx, healthTimeout) - defer cancel() - if err := c.Healthy(tCtx); err != nil { - return fmt.Errorf("failed to check controller is healthy: %w", err) - } + check(controllerName(c), c) } - log.Infof("Controllers healthy") - return nil + return g.Wait() } func init() { @@ -987,7 +1104,7 @@ func (m *MetalLBSpec) Deploy(ctx context.Context) error { m.Manifest = filepath.Join(m.ManifestDir, "metallb-native.yaml") } log.Infof("Deploying MetalLB from: %s", m.Manifest) - if err := run.LogCommand("kubectl", "apply", "-f", m.Manifest); err != nil { + if err := run.LogCommandContext(ctx, "kubectl", "apply", "-f", m.Manifest); err != nil { return fmt.Errorf("failed to deploy metallb: %w", err) } if _, err := m.kClient.CoreV1().Secrets("metallb-system").Get(ctx, "memberlist", metav1.GetOptions{}); err != nil { @@ -1120,7 +1237,7 @@ func (m *MeshnetSpec) Deploy(ctx context.Context) error { m.Manifest = filepath.Join(m.ManifestDir, "manifest.yaml") } log.Infof("Deploying Meshnet from: %s", m.Manifest) - if err := run.LogCommand("kubectl", "apply", "-f", m.Manifest); err != nil { + if err := run.LogCommandContext(ctx, "kubectl", "apply", "-f", m.Manifest); err != nil { return fmt.Errorf("failed to deploy meshnet: %w", err) } log.Infof("Meshnet Deployed") @@ -1197,7 +1314,7 @@ func (c *CEOSLabSpec) Deploy(ctx context.Context) error { c.Operator = filepath.Join(c.ManifestDir, "manifest.yaml") } log.Infof("Deploying CEOSLab controller from: %s", c.Operator) - if err := run.LogCommand("kubectl", "apply", "-f", c.Operator); err != nil { + if err := run.LogCommandContext(ctx, "kubectl", "apply", "-f", c.Operator); err != nil { return fmt.Errorf("failed to deploy ceoslab operator: %w", err) } log.Infof("CEOSLab controller deployed") @@ -1246,7 +1363,7 @@ func (l *LemmingSpec) Deploy(ctx context.Context) error { l.Operator = filepath.Join(l.ManifestDir, "manifest.yaml") } log.Infof("Deploying Lemming controller from: %s", l.Operator) - if err := run.LogCommand("kubectl", "apply", "-f", l.Operator); err != nil { + if err := run.LogCommandContext(ctx, "kubectl", "apply", "-f", l.Operator); err != nil { return fmt.Errorf("failed to deploy lemming operator: %w", err) } log.Infof("Lemming controller deployed") @@ -1295,7 +1412,7 @@ func (s *SRLinuxSpec) Deploy(ctx context.Context) error { s.Operator = filepath.Join(s.ManifestDir, "manifest.yaml") } log.Infof("Deploying SRLinux controller from: %s", s.Operator) - if err := run.LogCommand("kubectl", "apply", "-f", s.Operator); err != nil { + if err := run.LogCommandContext(ctx, "kubectl", "apply", "-f", s.Operator); err != nil { return fmt.Errorf("failed to deploy srlinux operator: %w", err) } log.Infof("SRLinux controller deployed") @@ -1346,7 +1463,7 @@ func (i *IxiaTGSpec) Deploy(ctx context.Context) error { i.Operator = filepath.Join(i.ManifestDir, "ixiatg-operator.yaml") } log.Infof("Deploying IxiaTG controller from: %s", i.Operator) - if err := run.LogCommand("kubectl", "apply", "-f", i.Operator); err != nil { + if err := run.LogCommandContext(ctx, "kubectl", "apply", "-f", i.Operator); err != nil { return fmt.Errorf("failed to deploy ixiatg operator: %w", err) } @@ -1380,7 +1497,7 @@ func (i *IxiaTGSpec) Deploy(ctx context.Context) error { i.ConfigMap = f.Name() } log.Infof("Deploying IxiaTG config map from: %s", i.ConfigMap) - if err := run.LogCommand("kubectl", "apply", "-f", i.ConfigMap); err != nil { + if err := run.LogCommandContext(ctx, "kubectl", "apply", "-f", i.ConfigMap); err != nil { return fmt.Errorf("failed to deploy ixiatg config map: %w", err) } log.Infof("IxiaTG controller deployed") @@ -1424,7 +1541,7 @@ func (c *CdnosSpec) Deploy(ctx context.Context) error { c.Operator = f.Name() } log.Infof("Deploying Cdnos controller from: %s", c.Operator) - if err := run.LogCommand("kubectl", "apply", "-f", c.Operator); err != nil { + if err := run.LogCommandContext(ctx, "kubectl", "apply", "-f", c.Operator); err != nil { return fmt.Errorf("failed to deploy cdnos operator: %w", err) } log.Infof("Cdnos controller deployed") diff --git a/deploy/deploy_test.go b/deploy/deploy_test.go index 0fa3ac073..7f5d09663 100644 --- a/deploy/deploy_test.go +++ b/deploy/deploy_test.go @@ -4,6 +4,8 @@ import ( "context" "fmt" "net/netip" + "strings" + "sync" "testing" "time" @@ -16,6 +18,7 @@ import ( "github.com/openconfig/kne/deploy/mocks" kexec "github.com/openconfig/kne/exec" fexec "github.com/openconfig/kne/exec/fake" + "github.com/openconfig/kne/exec/run" epb "github.com/openconfig/kne/proto/event" "github.com/pkg/errors" metallbv1 "go.universe.tf/metallb/api/v1beta1" @@ -2171,6 +2174,73 @@ func (m *mockCluster) GetName() string { return m.name } func (m *mockCluster) GetDockerNetworkResourceName() string { return "network" } func (m *mockCluster) Apply([]byte) error { return nil } +// mockComponent records when it ran and can stall, so tests can tell +// concurrent execution from sequential execution. +type mockComponent struct { + deployErr error + healthyErr error + + // deployDelay and healthyDelay stall the respective call. A stalled call + // returns early if its context is canceled, reporting the context error, + // which is how the timeout cases are exercised. + deployDelay time.Duration + healthyDelay time.Duration + + // onDeploy, if set, is called at the start of Deploy with the context the + // component was given. + onDeploy func(context.Context) + + mu sync.Mutex + deployStart time.Time + deployEnd time.Time +} + +func (m *mockComponent) Deploy(ctx context.Context) error { + if m.onDeploy != nil { + m.onDeploy(ctx) + } + m.mu.Lock() + m.deployStart = time.Now() + m.mu.Unlock() + if err := stall(ctx, m.deployDelay); err != nil { + return err + } + m.mu.Lock() + m.deployEnd = time.Now() + m.mu.Unlock() + return m.deployErr +} + +func (m *mockComponent) Healthy(ctx context.Context) error { + if err := stall(ctx, m.healthyDelay); err != nil { + return err + } + return m.healthyErr +} + +func (m *mockComponent) SetKClient(kubernetes.Interface) {} + +func (m *mockComponent) window() (time.Time, time.Time) { + m.mu.Lock() + defer m.mu.Unlock() + return m.deployStart, m.deployEnd +} + +// stall waits for d, or until ctx ends. +func stall(ctx context.Context, d time.Duration) error { + if d == 0 { + return nil + } + t := time.NewTimer(d) + defer t.Stop() + select { + case <-t.C: + return nil + case <-ctx.Done(): + return ctx.Err() + } +} + type mockIngress struct { deployErr error healthyErr error @@ -2182,6 +2252,14 @@ func (m *mockIngress) Healthy(context.Context) error { return m.healthyErr func (m *mockIngress) SetRCfg(*rest.Config) {} func (m *mockIngress) SetDockerNetworkResourceName(string) {} +// mockIngressComponent adapts mockComponent to the Ingress interface. +type mockIngressComponent struct { + *mockComponent +} + +func (m *mockIngressComponent) SetRCfg(*rest.Config) {} +func (m *mockIngressComponent) SetDockerNetworkResourceName(string) {} + type mockCNI struct { deployErr error healthyErr error @@ -2433,3 +2511,226 @@ func TestExtractVersionFromImage(t *testing.T) { }) } } + +func TestControllerName(t *testing.T) { + for _, tt := range []struct { + c Controller + want string + }{ + {c: &CEOSLabSpec{}, want: "CEOSLab controller"}, + {c: &SRLinuxSpec{}, want: "SRLinux controller"}, + {c: &IxiaTGSpec{}, want: "IxiaTG controller"}, + {c: &LemmingSpec{}, want: "Lemming controller"}, + {c: &CdnosSpec{}, want: "Cdnos controller"}, + {c: &mockController{}, want: "mockController controller"}, + } { + if got := controllerName(tt.c); got != tt.want { + t.Errorf("controllerName(%T) = %q, want %q", tt.c, got, tt.want) + } + } +} + +// TestDeployComponentsConcurrent checks that the components really do overlap, +// rather than each waiting for the previous one to report healthy. +func TestDeployComponentsConcurrent(t *testing.T) { + const delay = 200 * time.Millisecond + + ingress := &mockComponent{deployDelay: delay, healthyDelay: delay} + cni := &mockComponent{deployDelay: delay, healthyDelay: delay} + ctrl1 := &mockComponent{deployDelay: delay, healthyDelay: delay} + ctrl2 := &mockComponent{deployDelay: delay, healthyDelay: delay} + + d := &Deployment{ + Ingress: &mockIngressComponent{ingress}, + CNI: cni, + Controllers: []Controller{ctrl1, ctrl2}, + } + + start := time.Now() + if err := d.deployComponents(context.Background()); err != nil { + t.Fatalf("deployComponents() error = %v, want nil", err) + } + elapsed := time.Since(start) + + // Serialized this would be 4 components x 2 delays = 8*delay. Concurrent + // it is ~2*delay. Allow generous slack so the test is not timing flaky; + // it only needs to distinguish 2 from 8. + if max := 5 * delay; elapsed > max { + t.Errorf("deployComponents() took %v, want under %v: components do not appear to overlap", elapsed, max) + } + + // Every component should have been deploying at the same moment. + var latestStart, earliestEnd time.Time + for _, m := range []*mockComponent{ingress, cni, ctrl1, ctrl2} { + s, e := m.window() + if s.IsZero() || e.IsZero() { + t.Fatalf("component never deployed: start %v end %v", s, e) + } + if s.After(latestStart) { + latestStart = s + } + if earliestEnd.IsZero() || e.Before(earliestEnd) { + earliestEnd = e + } + } + if !latestStart.Before(earliestEnd) { + t.Errorf("deploys did not overlap: last started at %v, first finished at %v", latestStart, earliestEnd) + } +} + +func TestDeployComponents(t *testing.T) { + deployErr := errors.New("deploy failed") + healthyErr := errors.New("not healthy") + + for _, tt := range []struct { + name string + d *Deployment + // wantErr is a substring the error must contain, "" means no error. + wantErr string + }{ + { + name: "all healthy", + d: &Deployment{ + Ingress: &mockIngressComponent{&mockComponent{}}, + CNI: &mockComponent{}, + Controllers: []Controller{&mockComponent{}}, + }, + }, + { + name: "no controllers", + d: &Deployment{ + Ingress: &mockIngressComponent{&mockComponent{}}, + CNI: &mockComponent{}, + }, + }, + { + name: "ingress deploy fails", + d: &Deployment{ + Ingress: &mockIngressComponent{&mockComponent{deployErr: deployErr}}, + CNI: &mockComponent{}, + }, + wantErr: "failed to deploy ingress", + }, + { + name: "cni unhealthy", + d: &Deployment{ + Ingress: &mockIngressComponent{&mockComponent{}}, + CNI: &mockComponent{healthyErr: healthyErr}, + }, + wantErr: "failed to check if CNI is healthy", + }, + { + // The failing component must name itself, even though the others + // are cancelled at the same time. + name: "controller failure is attributed", + d: &Deployment{ + Ingress: &mockIngressComponent{&mockComponent{}}, + CNI: &mockComponent{}, + Controllers: []Controller{ + &mockComponent{healthyDelay: time.Hour}, + &mockComponent{deployErr: deployErr}, + }, + }, + wantErr: "failed to deploy mockComponent controller", + }, + { + // A nil component is skipped rather than panicking. + name: "nil cni", + d: &Deployment{ + Ingress: &mockIngressComponent{&mockComponent{}}, + }, + }, + } { + t.Run(tt.name, func(t *testing.T) { + err := tt.d.deployComponents(context.Background()) + switch { + case tt.wantErr == "" && err != nil: + t.Errorf("deployComponents() error = %v, want nil", err) + case tt.wantErr != "" && err == nil: + t.Errorf("deployComponents() error = nil, want error containing %q", tt.wantErr) + case tt.wantErr != "" && !strings.Contains(err.Error(), tt.wantErr): + t.Errorf("deployComponents() error = %v, want error containing %q", err, tt.wantErr) + } + }) + } +} + +// TestDeployComponentsTimeouts covers the two ways a deployment can run out of +// time, and checks the errors say which one happened. +func TestDeployComponentsTimeouts(t *testing.T) { + origHealthTimeout := healthTimeout + defer func() { healthTimeout = origHealthTimeout }() + + t.Run("component exceeds its own health budget", func(t *testing.T) { + healthTimeout = 100 * time.Millisecond + d := &Deployment{ + Ingress: &mockIngressComponent{&mockComponent{}}, + CNI: &mockComponent{healthyDelay: time.Hour}, + // Overall budget is far larger, so only the CNI's own budget + // can be what expired. + Timeout: time.Minute, + } + err := d.deployComponents(context.Background()) + if err == nil { + t.Fatalf("deployComponents() error = nil, want a timeout") + } + if want := "CNI was not healthy within"; !strings.Contains(err.Error(), want) { + t.Errorf("deployComponents() error = %v, want it to contain %q", err, want) + } + }) + + t.Run("overall timeout expires", func(t *testing.T) { + // A Deploy that hangs is not covered by the per-component health + // budget, so only the overall timeout can stop it. + healthTimeout = time.Hour + d := &Deployment{ + Ingress: &mockIngressComponent{&mockComponent{deployDelay: time.Hour}}, + CNI: &mockComponent{}, + Timeout: 100 * time.Millisecond, + } + start := time.Now() + err := d.deployComponents(context.Background()) + if err == nil { + t.Fatalf("deployComponents() error = nil, want a timeout") + } + if want := "overall deployment timeout"; !strings.Contains(err.Error(), want) { + t.Errorf("deployComponents() error = %v, want it to contain %q", err, want) + } + if elapsed := time.Since(start); elapsed > 10*time.Second { + t.Errorf("deployComponents() took %v, want it to give up promptly", elapsed) + } + }) +} + +// TestDeployComponentsLabelsCommands checks that each component is given a +// context labelled with its own name, so that the kubectl output of several +// components deploying at once can be told apart. +func TestDeployComponentsLabelsCommands(t *testing.T) { + var mu sync.Mutex + got := map[string]bool{} + record := func(ctx context.Context) { + mu.Lock() + defer mu.Unlock() + got[run.Label(ctx)] = true + } + + d := &Deployment{ + Ingress: &mockIngressComponent{&mockComponent{onDeploy: record}}, + CNI: &mockComponent{onDeploy: record}, + Controllers: []Controller{ + &mockComponent{onDeploy: record}, + }, + } + if err := d.deployComponents(context.Background()); err != nil { + t.Fatalf("deployComponents() error = %v, want nil", err) + } + + want := map[string]bool{ + "ingress": true, + "CNI": true, + "mockComponent controller": true, + } + if diff := cmp.Diff(want, got); diff != "" { + t.Errorf("command labels mismatch (-want +got):\n%s", diff) + } +} diff --git a/exec/fake/fake.go b/exec/fake/fake.go index 8574db084..ecbbe669e 100644 --- a/exec/fake/fake.go +++ b/exec/fake/fake.go @@ -19,6 +19,13 @@ // ... test code ... // // } +// +// A Command is safe for concurrent use. Each call to Command returns an +// independent exec.Cmd, and the responses are consumed under a lock, so code +// under test may run commands from multiple goroutines. Note that responses +// are still matched in order by default: when the commands may be issued +// concurrently, the order they arrive in is not deterministic, so the +// corresponding responses must be marked OutOfOrder. package fake import ( @@ -26,6 +33,7 @@ import ( "fmt" "io" "strings" + "sync" "github.com/openconfig/kne/exec" ) @@ -70,20 +78,31 @@ func (r Response) String() string { return buf.String() } -// A Command is an implementation of exec.Cmd that is used to return -// predefined results when exec.Cmd.Run is called. +// A Command hands out exec.Cmd implementations that return predefined results +// when exec.Cmd.Run is called. The zero value is not useful; use Commands. +// +// A Command is safe for concurrent use by multiple goroutines. type Command struct { - Name string // if set it is included in errors - cmd string - args []string + Name string // if set it is included in errors + + mu sync.Mutex responses []Response unexpected []Response - stdout io.Writer - stderr io.Writer - stdin io.Reader cnt int } +// An invocation is a single command produced by Command.Command. Each +// invocation holds its own stdio so that concurrent commands do not interfere +// with each other; the shared response bookkeeping lives on the parent. +type invocation struct { + parent *Command + cmd string + args []string + stdout io.Writer + stderr io.Writer + stdin io.Reader +} + // Commands returns a Command that is primed with the provided responses. // It's Command method can be used to override exec.Command. func Commands(resp []Response) *Command { @@ -92,21 +111,23 @@ func Commands(resp []Response) *Command { } } -// Command resets the command associated with c. +// Command returns a new exec.Cmd that draws its result from c when run. func (c *Command) Command(cmd string, args ...string) exec.Cmd { - c.cmd = cmd - c.args = args - return c + return &invocation{ + parent: c, + cmd: cmd, + args: args, + } } // Stdout sets standard out to w. -func (c *Command) SetStdout(w io.Writer) { c.stdout = w } +func (i *invocation) SetStdout(w io.Writer) { i.stdout = w } // Stderr sets standard err to w. -func (c *Command) SetStderr(w io.Writer) { c.stderr = w } +func (i *invocation) SetStderr(w io.Writer) { i.stderr = w } // Stdin sets standard in to r. -func (c *Command) SetStdin(r io.Reader) { c.stdin = r } +func (i *invocation) SetStdin(r io.Reader) { i.stdin = r } // LogCommand is called with the string representation of the command that is // running. The test program can optionally set this to their own function. @@ -118,40 +139,51 @@ var LogCommand = func(string) {} // calls LogCommand with a string representation of a Response that matches this // command. // -// Run returns nil if no matching response is found. Use c.Done to detect +// Run returns nil if no matching response is found. Use Command.Done to detect // these errors. -func (c *Command) Run() error { +func (i *invocation) Run() error { + return i.parent.run(i) +} + +// run consumes the response matching i. The lock is held for the whole call, +// which both keeps the response bookkeeping consistent and serializes the +// LogCommand callback, so a test hook that records commands does not need its +// own synchronization. +func (c *Command) run(i *invocation) error { + c.mu.Lock() + defer c.mu.Unlock() + c.cnt++ call := Response{ - Cmd: c.cmd, - Args: c.args, + Cmd: i.cmd, + Args: i.args, } defer func() { LogCommand(call.String()) }() if len(c.responses) == 0 { - c.unexpected = append(c.unexpected, Response{Cmd: c.cmd, Args: c.args}) + c.unexpected = append(c.unexpected, Response{Cmd: i.cmd, Args: i.args}) return nil } // Always check to see if we match the next expected response. // If we don't then look to see if there is an OutOfOrder response that we match. r := c.responses[0] - if c.matches(r) { + if i.matches(r) { c.responses = c.responses[1:] } else { matched := false - var i int - for i, r = range c.responses { - if r.OutOfOrder && c.matches(r) { + var n int + for n, r = range c.responses { + if r.OutOfOrder && i.matches(r) { matched = true - c.responses = append(c.responses[:i], c.responses[i+1:]...) + c.responses = append(c.responses[:n], c.responses[n+1:]...) break } } if !matched { - c.unexpected = append(c.unexpected, Response{Cmd: c.cmd, Args: c.args}) + c.unexpected = append(c.unexpected, Response{Cmd: i.cmd, Args: i.args}) return nil } } @@ -164,11 +196,13 @@ func (c *Command) Run() error { call.OutOfOrder = r.OutOfOrder call.Optional = r.Optional - if c.stdout != nil && r.Stdout != "" { - fmt.Fprint(c.stdout, r.Stdout) + // The writers belong to the caller, and a fake has nothing useful to do + // about a failure to write to them, so the results are discarded. + if i.stdout != nil && r.Stdout != "" { + _, _ = fmt.Fprint(i.stdout, r.Stdout) } - if c.stderr != nil && r.Stderr != "" { - fmt.Fprint(c.stderr, r.Stderr) + if i.stderr != nil && r.Stderr != "" { + _, _ = fmt.Fprint(i.stderr, r.Stderr) } switch e := r.Err.(type) { case string: @@ -213,8 +247,12 @@ func (e *DoneError) Error() string { // Done returns an error if there were any unexpected commands called on c or if // there are any non-optional responses left. // -// Done should be called once the test has finished calling c.Command. +// Done should be called once the test has finished calling c.Command and all +// commands it handed out have finished running. func (c *Command) Done() error { + c.mu.Lock() + defer c.mu.Unlock() + left := c.left() if len(left) == 0 && len(c.unexpected) == 0 { return nil @@ -230,7 +268,7 @@ func (c *Command) Done() error { } } -// left returns any non-optional unused responses +// left returns any non-optional unused responses. c.mu must be held. func (c *Command) left() []Response { var resp []Response for _, r := range c.responses { @@ -241,15 +279,15 @@ func (c *Command) left() []Response { return resp } -// matches returns true if the current command in c matches r. -func (c *Command) matches(r Response) bool { - if c.cmd != r.Cmd && r.Cmd != "" { +// matches returns true if the command in i matches r. +func (i *invocation) matches(r Response) bool { + if i.cmd != r.Cmd && r.Cmd != "" { return false } - return compareArgs(c.args, r.Args) + return compareArgs(i.args, r.Args) } -// compareArgs compares the two list of arguments to determin if they are the +// compareArgs compares the two list of arguments to determine if they are the // same or not. The values in wantArgs can have a ".*" as the suffix or prefix // to indicate a prefix or suffix match should be used instead of equality. func compareArgs(gotArgs, wantArgs []string) bool { diff --git a/exec/fake/fake_test.go b/exec/fake/fake_test.go index 351518788..1d8b92cd9 100644 --- a/exec/fake/fake_test.go +++ b/exec/fake/fake_test.go @@ -2,7 +2,9 @@ package fake import ( "errors" + "fmt" "strings" + "sync" "testing" "github.com/openconfig/gnmi/errdiff" @@ -284,11 +286,11 @@ func TestFailed(t *testing.T) { } } for _, u := range cmds.unexpected { - var ue Response + var wantResp Response if len(tt.unexpected) > 0 { - ue = tt.unexpected[0] + wantResp = tt.unexpected[0] } - t.Logf("Compare %v and %v", u, ue) + t.Logf("Compare %v and %v", u, wantResp) if len(tt.unexpected) > 0 && tt.unexpected[0].String() == u.String() { tt.unexpected = tt.unexpected[1:] continue @@ -415,3 +417,62 @@ func TestComparArg(t *testing.T) { } } } + +// TestConcurrent verifies that a single Command can serve commands run from +// multiple goroutines. Each goroutine runs a distinct command whose response +// is marked OutOfOrder, since the order the commands arrive in is not +// deterministic. Run with -race to check for data races. +func TestConcurrent(t *testing.T) { + const n = 16 + + resp := make([]Response, n) + for i := range resp { + resp[i] = Response{ + Cmd: fmt.Sprintf("cmd%d", i), + Args: []string{fmt.Sprintf("arg%d", i)}, + Stdout: fmt.Sprintf("out%d", i), + OutOfOrder: true, + } + } + cmds := Commands(resp) + + // LogCommand is called while the lock is held, so an unsynchronized + // callback like this one must not race. + var logged []string + oLog := LogCommand + defer func() { LogCommand = oLog }() + LogCommand = func(s string) { logged = append(logged, s) } + + var wg sync.WaitGroup + errs := make([]error, n) + outs := make([]strings.Builder, n) + start := make(chan struct{}) + for i := 0; i < n; i++ { + wg.Add(1) + go func() { + defer wg.Done() + c := cmds.Command(fmt.Sprintf("cmd%d", i), fmt.Sprintf("arg%d", i)) + c.SetStdout(&outs[i]) + <-start // maximize the overlap between goroutines + errs[i] = c.Run() + }() + } + close(start) + wg.Wait() + + for i := 0; i < n; i++ { + if errs[i] != nil { + t.Errorf("#%d: unexpected error: %v", i, errs[i]) + } + // Each goroutine must see only its own response's output. + if got, want := outs[i].String(), fmt.Sprintf("out%d", i); got != want { + t.Errorf("#%d: got stdout %q, want %q", i, got, want) + } + } + if len(logged) != n { + t.Errorf("got %d logged commands, want %d", len(logged), n) + } + if err := cmds.Done(); err != nil { + t.Errorf("%v", err) + } +} diff --git a/exec/run/run.go b/exec/run/run.go index 0449aede5..24a28c9d4 100644 --- a/exec/run/run.go +++ b/exec/run/run.go @@ -2,6 +2,7 @@ package run import ( "bytes" + "context" "io" kexec "github.com/openconfig/kne/exec" @@ -14,19 +15,54 @@ var ( logWarning = log.Warning ) +// labelKey is the context key under which a command label is stored. +type labelKey struct{} + +// WithLabel returns a copy of ctx that tags the output of commands run with it +// as coming from label. Callers that run commands concurrently should set a +// label; without one the interleaved output of several commands is +// indistinguishable, since every line is attributed only to the binary that +// produced it. +// +// The label affects logging only. It does not tie the command's lifetime to +// ctx: commands are not canceled when ctx ends. +func WithLabel(ctx context.Context, label string) context.Context { + return context.WithValue(ctx, labelKey{}, label) +} + +// Label returns the label attached to ctx by WithLabel, or "" if there is +// none. Components can use it to tag their own log messages the same way +// their command output is tagged. +func Label(ctx context.Context) string { + if ctx == nil { + return "" + } + s, _ := ctx.Value(labelKey{}).(string) + return s +} + +// logPrefix returns the prefix to put in front of the output of cmd. +func logPrefix(ctx context.Context, cmd string) string { + if l := Label(ctx); l != "" { + return "(" + cmd + "/" + l + "): " + } + return "(" + cmd + "): " +} + // runCommand is a wrapper utility function that creates and runs a command with // various inputs and settings. -func runCommand(writeLogs bool, in []byte, cmd string, args ...string) ([]byte, error) { +func runCommand(ctx context.Context, writeLogs bool, in []byte, cmd string, args ...string) ([]byte, error) { c := kexec.Command(cmd, args...) var out bytes.Buffer c.SetStdout(&out) c.SetStderr(&out) if writeLogs { + prefix := logPrefix(ctx, cmd) outLog := logshim.New(func(v ...interface{}) { - logInfo(append([]interface{}{"(" + cmd + "): "}, v...)...) + logInfo(append([]interface{}{prefix}, v...)...) }) errLog := logshim.New(func(v ...interface{}) { - logWarning(append([]interface{}{"(" + cmd + "): "}, v...)...) + logWarning(append([]interface{}{prefix}, v...)...) }) defer func() { outLog.Close() @@ -45,7 +81,14 @@ func runCommand(writeLogs bool, in []byte, cmd string, args ...string) ([]byte, // LogCommand runs the specified command but records standard output // with log.Info and standard error with log.Warning. func LogCommand(cmd string, args ...string) error { - _, err := runCommand(true, nil, cmd, args...) + _, err := runCommand(context.Background(), true, nil, cmd, args...) + return err +} + +// LogCommandContext is LogCommand with any label attached to ctx applied to +// the logged output. +func LogCommandContext(ctx context.Context, cmd string, args ...string) error { + _, err := runCommand(ctx, true, nil, cmd, args...) return err } @@ -53,7 +96,7 @@ func LogCommand(cmd string, args ...string) error { // with log.Info and standard error with log.Warning. in is sent to // the standard input of the command. func LogCommandWithInput(in []byte, cmd string, args ...string) error { - _, err := runCommand(true, in, cmd, args...) + _, err := runCommand(context.Background(), true, in, cmd, args...) return err } @@ -61,11 +104,11 @@ func LogCommandWithInput(in []byte, cmd string, args ...string) error { // with log.Info and standard error with log.Warning. Standard output // and standard error are also returned. func OutLogCommand(cmd string, args ...string) ([]byte, error) { - return runCommand(true, nil, cmd, args...) + return runCommand(context.Background(), true, nil, cmd, args...) } // OutCommand runs the specified command and returns any standard output // as well as any errors. func OutCommand(cmd string, args ...string) ([]byte, error) { - return runCommand(false, nil, cmd, args...) + return runCommand(context.Background(), false, nil, cmd, args...) } diff --git a/exec/run/run_test.go b/exec/run/run_test.go index b57a6bc92..e9632c59c 100644 --- a/exec/run/run_test.go +++ b/exec/run/run_test.go @@ -2,6 +2,7 @@ package run import ( "bytes" + "context" "fmt" "testing" @@ -14,6 +15,7 @@ func TestRunCommand(t *testing.T) { tests := []struct { desc string writeLogs bool + label string in string cmd string args []string @@ -69,6 +71,29 @@ func TestRunCommand(t *testing.T) { {Cmd: "cat", Stdout: "hello"}, }, want: "hello", + }, { + desc: "log with label", + writeLogs: true, + label: "SRLinux controller", + cmd: "kubectl", + args: []string{"apply"}, + resp: []fexec.Response{ + {Cmd: "kubectl", Args: []string{"apply"}, Stdout: "applied"}, + }, + want: "applied", + wantInfos: "(kubectl/SRLinux controller): applied", + }, { + desc: "log with label and stderr", + writeLogs: true, + label: "ingress", + cmd: "kubectl", + args: []string{"apply"}, + resp: []fexec.Response{ + {Cmd: "kubectl", Args: []string{"apply"}, Stdout: "out", Stderr: "err"}, + }, + want: "outerr", + wantInfos: "(kubectl/ingress): out", + wantWarnings: "(kubectl/ingress): err", }, { desc: "failed command", cmd: "false", @@ -104,7 +129,11 @@ func TestRunCommand(t *testing.T) { fmt.Fprint(&warnings, args...) } - got, err := runCommand(tt.writeLogs, []byte(tt.in), tt.cmd, tt.args...) + ctx := context.Background() + if tt.label != "" { + ctx = WithLabel(ctx, tt.label) + } + got, err := runCommand(ctx, tt.writeLogs, []byte(tt.in), tt.cmd, tt.args...) if s := errdiff.Substring(err, tt.wantErr); s != "" { t.Fatalf("unexpected error: %s", s) }