diff --git a/pkg/cli/probe.go b/pkg/cli/probe.go index b032841..a7ef80d 100644 --- a/pkg/cli/probe.go +++ b/pkg/cli/probe.go @@ -1,6 +1,8 @@ package cli import ( + "strings" + "github.com/mattfenwick/cyclonus/pkg/connectivity" "github.com/mattfenwick/cyclonus/pkg/connectivity/probe" "github.com/mattfenwick/cyclonus/pkg/generator" @@ -11,7 +13,6 @@ import ( v1 "k8s.io/api/core/v1" networkingv1 "k8s.io/api/networking/v1" "k8s.io/apimachinery/pkg/util/intstr" - "strings" ) type ProbeArgs struct { diff --git a/pkg/connectivity/comparisontable.go b/pkg/connectivity/comparisontable.go index 2c9b337..7675b70 100644 --- a/pkg/connectivity/comparisontable.go +++ b/pkg/connectivity/comparisontable.go @@ -44,6 +44,12 @@ func NewComparisonTable(items []string) *ComparisonTable { } func NewComparisonTableFrom(kubeProbe *probe.Table, simulatedProbe *probe.Table) *ComparisonTable { + if kubeProbe == nil { + panic(errors.Errorf("kubeprobe is nil")) + } + if simulatedProbe == nil { + panic(errors.Errorf("sim probe is nil")) + } if len(kubeProbe.Wrapped.Froms) != len(simulatedProbe.Wrapped.Froms) || len(kubeProbe.Wrapped.Tos) != len(simulatedProbe.Wrapped.Tos) { panic(errors.Errorf("cannot compare tables of different dimensions")) } diff --git a/pkg/connectivity/interpreter.go b/pkg/connectivity/interpreter.go index 9d532fe..9f4fb68 100644 --- a/pkg/connectivity/interpreter.go +++ b/pkg/connectivity/interpreter.go @@ -2,6 +2,8 @@ package connectivity import ( "fmt" + "time" + "github.com/mattfenwick/cyclonus/pkg/connectivity/probe" "github.com/mattfenwick/cyclonus/pkg/generator" "github.com/mattfenwick/cyclonus/pkg/kube" @@ -9,7 +11,6 @@ import ( "github.com/pkg/errors" "github.com/sirupsen/logrus" networkingv1 "k8s.io/api/networking/v1" - "time" ) const ( @@ -116,6 +117,10 @@ func (t *Interpreter) ExecuteTestCase(testCase *generator.TestCase) *Result { err = testCaseState.SetPodLabels(ns, pod, labels) } else if action.DeletePod != nil { err = testCaseState.DeletePod(action.DeletePod.Namespace, action.DeletePod.Pod) + } else if action.CreateService != nil { + err = testCaseState.CreateService(action.CreateService.Service) + } else if action.DeleteService != nil { + err = testCaseState.DeleteService(action.DeleteService.Service) } else { err = errors.Errorf("invalid Action at step %d, action %d", stepIndex, actionIndex) } diff --git a/pkg/connectivity/probe/jobbuilder.go b/pkg/connectivity/probe/jobbuilder.go index 1dc516c..e110ba1 100644 --- a/pkg/connectivity/probe/jobbuilder.go +++ b/pkg/connectivity/probe/jobbuilder.go @@ -3,6 +3,7 @@ package probe import ( "github.com/mattfenwick/cyclonus/pkg/generator" "github.com/pkg/errors" + "github.com/sirupsen/logrus" v1 "k8s.io/api/core/v1" "k8s.io/apimachinery/pkg/util/intstr" ) @@ -12,8 +13,11 @@ type JobBuilder struct { } func (j *JobBuilder) GetJobsForProbeConfig(resources *Resources, config *generator.ProbeConfig) *Jobs { + logrus.Debugf("getting jobs for probe config %v", config) if config.AllAvailable { return j.GetJobsAllAvailableServers(resources, config.Mode) + } else if config.Mode == generator.ProbeModeNodeIP { + return j.GetJobsForNodeIP(resources, config) } else if config.PortProtocol != nil { return j.GetJobsForNamedPortProtocol(resources, config.PortProtocol.Port, config.PortProtocol.Protocol, config.Mode) } else { @@ -21,10 +25,46 @@ func (j *JobBuilder) GetJobsForProbeConfig(resources *Resources, config *generat } } +func (j *JobBuilder) GetJobsForNodeIP(resources *Resources, config *generator.ProbeConfig) *Jobs { + jobs := &Jobs{} + logrus.Debugf("getting jobs for node ip %v", config) + + for _, podFrom := range resources.Pods { + for _, node := range resources.Nodes { + job := &Job{ + FromKey: podFrom.PodString().String(), + FromNamespace: podFrom.Namespace, + FromNamespaceLabels: resources.Namespaces[podFrom.Namespace], + FromPod: podFrom.Name, + FromPodLabels: podFrom.Labels, + FromContainer: podFrom.Containers[0].Name, + FromIP: podFrom.IP, + ToKey: node.Name, + ToHost: node.IP, + ToNamespace: "node", + ToNamespaceLabels: map[string]string{}, + ToPodLabels: map[string]string{}, + ToIP: node.IP, + ResolvedPort: config.PortProtocol.Port.IntValue(), + ResolvedPortName: "Custom", + Protocol: config.PortProtocol.Protocol, + TimeoutSeconds: j.TimeoutSeconds, + } + jobs.Valid = append(jobs.Valid, job) + + } + } + + return jobs +} + func (j *JobBuilder) GetJobsForNamedPortProtocol(resources *Resources, port intstr.IntOrString, protocol v1.Protocol, mode generator.ProbeMode) *Jobs { jobs := &Jobs{} + logrus.Debugf("named port getting jobs for resources %+v", resources) for _, podFrom := range resources.Pods { + logrus.Debugf("named port getting jobs for podfrom %+v", podFrom) for _, podTo := range resources.Pods { + logrus.Debugf("named port getting jobs for podTo %+v", podTo) job := &Job{ FromKey: podFrom.PodString().String(), FromNamespace: podFrom.Namespace, @@ -76,8 +116,11 @@ func (j *JobBuilder) GetJobsForNamedPortProtocol(resources *Resources, port ints func (j *JobBuilder) GetJobsAllAvailableServers(resources *Resources, mode generator.ProbeMode) *Jobs { var jobs []*Job + logrus.Debugf("all available getting jobs for resources %+v", resources) for _, podFrom := range resources.Pods { + logrus.Debugf("all available getting jobs for podfrom %+v", podFrom) for _, podTo := range resources.Pods { + logrus.Debugf("all available getting jobs for podTo %+v", podTo) for _, contTo := range podTo.Containers { jobs = append(jobs, &Job{ FromKey: podFrom.PodString().String(), diff --git a/pkg/connectivity/probe/jobrunner.go b/pkg/connectivity/probe/jobrunner.go index efa11b6..148484f 100644 --- a/pkg/connectivity/probe/jobrunner.go +++ b/pkg/connectivity/probe/jobrunner.go @@ -1,13 +1,14 @@ package probe import ( + "strings" + "github.com/mattfenwick/cyclonus/pkg/generator" "github.com/mattfenwick/cyclonus/pkg/kube" "github.com/mattfenwick/cyclonus/pkg/matcher" "github.com/mattfenwick/cyclonus/pkg/utils" "github.com/mattfenwick/cyclonus/pkg/worker" "github.com/sirupsen/logrus" - "strings" ) type Runner struct { @@ -28,12 +29,24 @@ func NewKubeBatchRunner(kubernetes kube.IKubernetes, workers int, jobBuilder *Jo } func (p *Runner) RunProbeForConfig(probeConfig *generator.ProbeConfig, resources *Resources) *Table { - return NewTableFromJobResults(resources, p.runProbe(p.JobBuilder.GetJobsForProbeConfig(resources, probeConfig))) + jobs := p.JobBuilder.GetJobsForProbeConfig(resources, probeConfig) + logrus.Debugf("got jobs %+v", jobs) + jobresults := p.runProbe(jobs) + if probeConfig.Mode == generator.ProbeModeNodeIP { + return NewNodeTableFromJobResults(resources, jobresults) + } else { + return NewPodTableFromJobResults(resources, jobresults) + } } func (p *Runner) runProbe(jobs *Jobs) []*JobResult { + logrus.Debugf("running probe for job %+v", jobs) resultSlice := p.JobRunner.RunJobs(jobs.Valid) + for _, res := range resultSlice { + logrus.Debugf("resultslice combined: %+v, ingress: %+v, egress %+v", res.Combined, res.Ingress, res.Egress) + } + invalidPP := ConnectivityInvalidPortProtocol unknown := ConnectivityUnknown for _, j := range jobs.BadPortProtocol { @@ -100,6 +113,7 @@ type KubeJobRunner struct { } func (k *KubeJobRunner) RunJobs(jobs []*Job) []*JobResult { + logrus.Debugf("run job single with %+v", jobs) size := len(jobs) jobsChan := make(chan *Job, size) resultsChan := make(chan *JobResult, size) @@ -124,6 +138,7 @@ func (k *KubeJobRunner) RunJobs(jobs []*Job) []*JobResult { // it only writes pass/fail status to a channel and has no failure side effects, this is by design since we do not want to fail inside a goroutine. func (k *KubeJobRunner) worker(jobs <-chan *Job, results chan<- *JobResult) { for job := range jobs { + logrus.Debugf("probing connectivity for job %+v", job) connectivity, _ := probeConnectivity(k.Kubernetes, job) results <- &JobResult{ Job: job, @@ -135,13 +150,13 @@ func (k *KubeJobRunner) worker(jobs <-chan *Job, results chan<- *JobResult) { func probeConnectivity(k8s kube.IKubernetes, job *Job) (Connectivity, string) { commandDebugString := strings.Join(job.KubeExecCommand(), " ") stdout, stderr, commandErr, err := k8s.ExecuteRemoteCommand(job.FromNamespace, job.FromPod, job.FromContainer, job.ClientCommand()) - logrus.Debugf("stdout, stderr from %s: \n%s\n%s", commandDebugString, stdout, stderr) + logrus.Debugf("stdout, stderr from [%s]: \n%s\n%s", commandDebugString, stdout, stderr) if err != nil { - logrus.Errorf("unable to set up command %s: %+v", commandDebugString, err) + logrus.Errorf("unable to set up command [%s]: %+v", commandDebugString, err) return ConnectivityCheckFailed, commandDebugString } if commandErr != nil { - logrus.Debugf("unable to run command %s: %+v", commandDebugString, commandErr) + logrus.Debugf("unable to run command [%s]: %+v", commandDebugString, commandErr) return ConnectivityBlocked, commandDebugString } return ConnectivityAllowed, commandDebugString diff --git a/pkg/connectivity/probe/node.go b/pkg/connectivity/probe/node.go new file mode 100644 index 0000000..65278a3 --- /dev/null +++ b/pkg/connectivity/probe/node.go @@ -0,0 +1,15 @@ +package probe + +type Node struct { + Name string + Labels map[string]string + IP string +} + +func NewNode(name string, labels map[string]string, ip string) *Node { + return &Node{ + Name: name, + Labels: make(map[string]string), + IP: ip, + } +} diff --git a/pkg/connectivity/probe/pod.go b/pkg/connectivity/probe/pod.go index 32b3482..1371d13 100644 --- a/pkg/connectivity/probe/pod.go +++ b/pkg/connectivity/probe/pod.go @@ -2,13 +2,14 @@ package probe import ( "fmt" + "strings" + "github.com/mattfenwick/collections/pkg/slice" "github.com/mattfenwick/cyclonus/pkg/generator" "github.com/mattfenwick/cyclonus/pkg/kube" "github.com/pkg/errors" v1 "k8s.io/api/core/v1" metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" - "strings" ) const ( @@ -48,6 +49,7 @@ type Pod struct { Labels map[string]string ServiceIP string IP string + NodeIP string Containers []*Container } @@ -59,6 +61,8 @@ func (p *Pod) Host(probeMode generator.ProbeMode) string { return p.IP case generator.ProbeModeServiceIP: return p.ServiceIP + case generator.ProbeModeNodeIP: + return p.NodeIP default: panic(errors.Errorf("invalid mode %s", probeMode)) } @@ -89,6 +93,10 @@ func (p *Pod) ServiceName() string { return fmt.Sprintf("s-%s-%s", p.Namespace, p.Name) } +func (p *Pod) ServiceNameLoadBalancer() string { + return fmt.Sprintf("s-%s-%s-lb", p.Namespace, p.Name) +} + func (p *Pod) KubePod() *v1.Pod { zero := int64(0) return &v1.Pod{ @@ -117,6 +125,22 @@ func (p *Pod) KubeService() *v1.Service { } } +func (p *Pod) KubeServiceLoadBalancer() *v1.Service { + ports := slice.Map(func(cont *Container) v1.ServicePort { return cont.KubeServicePort() }, p.Containers) + tcpPorts := slice.Filter(func(port v1.ServicePort) bool { return port.Protocol == v1.ProtocolTCP }, ports) + return &v1.Service{ + ObjectMeta: metav1.ObjectMeta{ + Name: p.ServiceNameLoadBalancer(), + Namespace: p.Namespace, + }, + Spec: v1.ServiceSpec{ + Ports: tcpPorts, + Selector: p.Labels, + Type: v1.ServiceTypeLoadBalancer, + }, + } +} + func (p *Pod) KubeContainers() []v1.Container { return slice.Map(func(cont *Container) v1.Container { return cont.KubeContainer() }, p.Containers) } @@ -188,6 +212,15 @@ func (c *Container) KubeServicePort() v1.ServicePort { } } +// when using load balancer types, cannot contain more than 1 protocol +func (c *Container) KubeServicePortTCP() v1.ServicePort { + return v1.ServicePort{ + Name: fmt.Sprintf("service-port-%s-%d", strings.ToLower(string(v1.ProtocolTCP)), c.Port), + Protocol: v1.ProtocolTCP, + Port: int32(c.Port), + } +} + func (c *Container) Image() string { if c.BatchJobs { return cyclonusWorkerImage diff --git a/pkg/connectivity/probe/resources.go b/pkg/connectivity/probe/resources.go index 5ee86df..0d17d27 100644 --- a/pkg/connectivity/probe/resources.go +++ b/pkg/connectivity/probe/resources.go @@ -1,6 +1,8 @@ package probe import ( + "time" + "github.com/mattfenwick/collections/pkg/slice" "github.com/mattfenwick/cyclonus/pkg/kube" "github.com/pkg/errors" @@ -8,12 +10,13 @@ import ( "golang.org/x/exp/maps" v1 "k8s.io/api/core/v1" metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" - "time" ) type Resources struct { Namespaces map[string]map[string]string Pods []*Pod + Nodes []*Node + Services map[string]*v1.Service //ExternalIPs []string } @@ -22,6 +25,8 @@ func NewDefaultResources(kubernetes kube.IKubernetes, namespaces []string, podNa r := &Resources{ Namespaces: map[string]map[string]string{}, + Services: make(map[string]*v1.Service), + //ExternalIPs: externalIPs, } @@ -139,6 +144,7 @@ func (r *Resources) CreateNamespace(ns string, labels map[string]string) (*Resou return &Resources{ Namespaces: newNamespaces, Pods: r.Pods, + Nodes: r.Nodes, }, nil } @@ -155,6 +161,7 @@ func (r *Resources) UpdateNamespaceLabels(ns string, labels map[string]string) ( return &Resources{ Namespaces: newNamespaces, Pods: r.Pods, + Nodes: r.Nodes, }, nil } @@ -180,9 +187,57 @@ func (r *Resources) DeleteNamespace(ns string) (*Resources, error) { return &Resources{ Namespaces: newNamespaces, Pods: pods, + Nodes: r.Nodes, + }, nil +} + +// CreateServce returns a new object with a new namespace. It should not affect the original Resources object. +func (r *Resources) CreateService(svc *v1.Service) (*Resources, error) { + if _, ok := r.Services[svc.Name]; ok { + return nil, errors.Errorf("service %s already found", svc.Name) + } + newServices := map[string]*v1.Service{} + for oldServiceName, oldService := range r.Services { + newServices[oldServiceName] = oldService // Note: service type is pointer, duplicate resource type needs to be deep copied + } + newServices[svc.Name] = svc + return &Resources{ + Services: newServices, + Pods: r.Pods, + Nodes: r.Nodes, + Namespaces: r.Namespaces, + }, nil +} + +// DeleteNamespace returns a new object without the namespace. It should not affect the original Resources object. +func (r *Resources) DeleteService(svc *v1.Service) (*Resources, error) { + if _, ok := r.Services[svc.Name]; !ok { + return nil, errors.Errorf("service %s/%s not found in test state", svc.Namespace, svc.Name) + } + newServices := map[string]*v1.Service{} + for oldServiceName, oldService := range r.Services { + if oldServiceName != svc.Name { + newServices[oldServiceName] = oldService + } + } + + return &Resources{ + Services: newServices, }, nil } +func (r *Resources) addNodes(nodes *v1.NodeList) { + for _, node := range nodes.Items { + nodeips := node.Status.Addresses + if len(nodeips) > 0 { + logrus.Debugf("loading node name %s and ip %+v", node.Name, nodeips[0].Address) + r.Nodes = append(r.Nodes, NewNode(node.Name, node.Labels, nodeips[0].Address)) + } else { + logrus.Errorf("node %s has no ip's", node.Name) + } + } +} + // CreatePod returns a new object with a new pod. It should not affect the original Resources object. func (r *Resources) CreatePod(ns string, podName string, labels map[string]string) (*Resources, error) { // TODO this needs to be improved @@ -192,6 +247,7 @@ func (r *Resources) CreatePod(ns string, podName string, labels map[string]strin } return &Resources{ Namespaces: r.Namespaces, + Nodes: r.Nodes, Pods: append(append([]*Pod{}, r.Pods...), NewPod(ns, podName, labels, "TODO", r.Pods[0].Containers)), //ExternalIPs: r.ExternalIPs, }, nil @@ -214,6 +270,7 @@ func (r *Resources) SetPodLabels(ns string, podName string, labels map[string]st } return &Resources{ Namespaces: r.Namespaces, + Nodes: r.Nodes, Pods: pods, //ExternalIPs: r.ExternalIPs, }, nil @@ -236,6 +293,7 @@ func (r *Resources) DeletePod(ns string, podName string) (*Resources, error) { return &Resources{ Namespaces: r.Namespaces, Pods: newPods, + Nodes: r.Nodes, //ExternalIPs: r.ExternalIPs, }, nil } @@ -246,6 +304,12 @@ func (r *Resources) SortedPodNames() []string { r.Pods)) } +func (r *Resources) SortedNodeNames() []string { + return slice.Sort(slice.Map( + func(n *Node) string { return n.Name }, + r.Nodes)) +} + func (r *Resources) NamespacesSlice() []string { return maps.Keys(r.Namespaces) } @@ -269,15 +333,27 @@ func (r *Resources) CreateResourcesInKube(kubernetes kube.IKubernetes) error { } } kubeService := pod.KubeService() + kubeServiceLoadBalancer := pod.KubeServiceLoadBalancer() _, err = kubernetes.GetService(kubeService.Namespace, kubeService.Name) if err != nil { _, err = kubernetes.CreateService(kubeService) if err != nil { return err } + _, err = kubernetes.CreateService(kubeServiceLoadBalancer) + if err != nil { + return err + } } } - return nil + + nodes, err := kubernetes.GetNodes() + if err != nil { + return err + } + r.addNodes(nodes) + + return err } func KubeNamespace(ns string, labels map[string]string) *v1.Namespace { diff --git a/pkg/connectivity/probe/resources_test.go b/pkg/connectivity/probe/resources_test.go index b9ca433..f13befb 100644 --- a/pkg/connectivity/probe/resources_test.go +++ b/pkg/connectivity/probe/resources_test.go @@ -3,6 +3,7 @@ package probe import ( . "github.com/onsi/ginkgo/v2" . "github.com/onsi/gomega" + v1 "k8s.io/api/core/v1" ) func RunResourcesTests() { @@ -49,5 +50,17 @@ func RunResourcesTests() { Expect(r.Pods[0].Labels).To(Equal(labels)) Expect(r2.Pods[0].Labels).To(Equal(map[string]string{})) }) + + It("Should create a service nondestructively", func() { + r := &Resources{ + Services: make(map[string]*v1.Service), + } + svc := v1.Service{} + r2, err := r.CreateService(&svc) + Expect(err).To(Succeed()) + + Expect(r.Services).To(HaveLen(0)) + Expect(r2.Services).To(HaveLen(1)) + }) }) } diff --git a/pkg/connectivity/probe/table.go b/pkg/connectivity/probe/table.go index 1ff0a77..9a98d07 100644 --- a/pkg/connectivity/probe/table.go +++ b/pkg/connectivity/probe/table.go @@ -1,11 +1,13 @@ package probe import ( + "strings" + "github.com/mattfenwick/collections/pkg/slice" "github.com/mattfenwick/cyclonus/pkg/utils" "github.com/pkg/errors" + "github.com/sirupsen/logrus" "golang.org/x/exp/maps" - "strings" ) type Item struct { @@ -36,7 +38,7 @@ func NewTable(items []string) *Table { })} } -func NewTableFromJobResults(resources *Resources, jobResults []*JobResult) *Table { +func NewPodTableFromJobResults(resources *Resources, jobResults []*JobResult) *Table { table := NewTable(resources.SortedPodNames()) for _, result := range jobResults { fr := result.Job.FromKey @@ -48,6 +50,21 @@ func NewTableFromJobResults(resources *Resources, jobResults []*JobResult) *Tabl return table } +// Note: The UX here needs some work to accurately portray pod->node matrix +func NewNodeTableFromJobResults(resources *Resources, jobResults []*JobResult) *Table { + res := append(resources.SortedNodeNames(), resources.SortedPodNames()...) + logrus.Debugf("merged table %+v", res) + table := NewTable(res) + for _, result := range jobResults { + fr := result.Job.FromKey + to := result.Job.ToKey + pp := table.Get(fr, to) + // this really shouldn't happen, so let's not recover from it + utils.DoOrDie(pp.AddJobResult(result)) + } + return table +} + //func (t *Table) Set(from string, to string, value *Item) { // t.Wrapped.Set(from, to, value) //} diff --git a/pkg/connectivity/stepresult.go b/pkg/connectivity/stepresult.go index df0f2e9..74c7c3c 100644 --- a/pkg/connectivity/stepresult.go +++ b/pkg/connectivity/stepresult.go @@ -3,6 +3,7 @@ package connectivity import ( "github.com/mattfenwick/cyclonus/pkg/connectivity/probe" "github.com/mattfenwick/cyclonus/pkg/matcher" + "github.com/sirupsen/logrus" networkingv1 "k8s.io/api/networking/v1" ) @@ -29,6 +30,7 @@ func (s *StepResult) AddKubeProbe(kubeProbe *probe.Table) { func (s *StepResult) Comparison(i int) *ComparisonTable { if s.comparisons[i] == nil { + logrus.Debugf("comparing i [%d] in kubeprobes %+v", i, s.KubeProbes) s.comparisons[i] = NewComparisonTableFrom(s.KubeProbes[i], s.SimulatedProbe) } return s.comparisons[i] diff --git a/pkg/connectivity/summary.go b/pkg/connectivity/summary.go index 2ec65d9..37ed85e 100644 --- a/pkg/connectivity/summary.go +++ b/pkg/connectivity/summary.go @@ -2,6 +2,7 @@ package connectivity import ( "fmt" + v1 "k8s.io/api/core/v1" ) diff --git a/pkg/connectivity/testcasestate.go b/pkg/connectivity/testcasestate.go index 5868041..304ba0d 100644 --- a/pkg/connectivity/testcasestate.go +++ b/pkg/connectivity/testcasestate.go @@ -1,13 +1,16 @@ package connectivity import ( + "fmt" + "log" + "time" + "github.com/mattfenwick/cyclonus/pkg/connectivity/probe" "github.com/mattfenwick/cyclonus/pkg/kube" "github.com/pkg/errors" "github.com/sirupsen/logrus" v1 "k8s.io/api/core/v1" networkingv1 "k8s.io/api/networking/v1" - "time" ) type TestCaseState struct { @@ -49,6 +52,26 @@ func (t *TestCaseState) UpdatePolicy(policy *networkingv1.NetworkPolicy) error { return err } +func (t *TestCaseState) CreateService(svc *v1.Service) error { + newResources, err := t.Resources.CreateService(svc) + if err != nil { + return err + } + t.Resources = newResources + fmt.Printf("creating service %+v", svc.Name) + _, err = t.Kubernetes.CreateService(svc) + return err +} + +func (t *TestCaseState) DeleteService(svc *v1.Service) error { + newResources, err := t.Resources.DeleteService(svc) + if err != nil { + return err + } + t.Resources = newResources + return t.Kubernetes.DeleteService(svc.Name, svc.Namespace) +} + func (t *TestCaseState) CreateNamespace(ns string, labels map[string]string) error { newResources, err := t.Resources.CreateNamespace(ns, labels) if err != nil { @@ -96,6 +119,12 @@ func (t *TestCaseState) CreatePod(ns string, pod string, labels map[string]strin if err != nil { return err } + + _, err = t.Kubernetes.CreateService(newPod.KubeServiceLoadBalancer()) + if err != nil { + return err + } + // wait for ready, get ip for i := 0; i < 12; i++ { kubePod, err := t.Kubernetes.GetPod(ns, pod) @@ -104,6 +133,8 @@ func (t *TestCaseState) CreatePod(ns string, pod string, labels map[string]strin } if kubePod.Status.Phase == "Running" && kubePod.Status.PodIP != "" { newPod.IP = kubePod.Status.PodIP + newPod.NodeIP = "10.240.0.4" + log.Printf("setting node ip to %v", newPod.NodeIP) return nil } time.Sleep(5 * time.Second) diff --git a/pkg/generator/action.go b/pkg/generator/action.go index 5b32580..55eacf8 100644 --- a/pkg/generator/action.go +++ b/pkg/generator/action.go @@ -1,6 +1,9 @@ package generator -import networkingv1 "k8s.io/api/networking/v1" +import ( + v1 "k8s.io/api/core/v1" + networkingv1 "k8s.io/api/networking/v1" +) // Action models a sum type (discriminated union): exactly one field must be non-null. type Action struct { @@ -17,6 +20,10 @@ type Action struct { CreatePod *CreatePodAction SetPodLabels *SetPodLabelsAction DeletePod *DeletePodAction + + CreateService *CreateServiceAction + UpdateService *UpdateServiceAction + DeleteService *DeleteServiceAction } type CreatePolicyAction struct { @@ -117,3 +124,33 @@ func DeletePod(namespace string, pod string) *Action { Pod: pod, }} } + +type CreateServiceAction struct { + Service *v1.Service +} + +func CreateService(svc *v1.Service) *Action { + return &Action{CreateService: &CreateServiceAction{ + Service: svc, + }} +} + +type UpdateServiceAction struct { + Service *v1.Service +} + +func UpdateService(svc *v1.Service) *Action { + return &Action{UpdateService: &UpdateServiceAction{ + Service: svc, + }} +} + +type DeleteServiceAction struct { + Service *v1.Service +} + +func DeleteService(svc *v1.Service) *Action { + return &Action{DeleteService: &DeleteServiceAction{ + Service: svc, + }} +} diff --git a/pkg/generator/feature.go b/pkg/generator/feature.go index 13df07e..5ec81ac 100644 --- a/pkg/generator/feature.go +++ b/pkg/generator/feature.go @@ -20,6 +20,9 @@ const ( ActionFeatureCreatePod = "action: create pod" ActionFeatureSetPodLabels = "action: set pod labels" ActionFeatureDeletePod = "action: delete pod" + + ActionFeatureCreateService = "action: create service" + ActionFeatureDeleteService = "action: delete service" ) const ( diff --git a/pkg/generator/loadbalancertestcases.go b/pkg/generator/loadbalancertestcases.go new file mode 100644 index 0000000..0bad04c --- /dev/null +++ b/pkg/generator/loadbalancertestcases.go @@ -0,0 +1,66 @@ +package generator + +import ( + v1 "k8s.io/api/core/v1" + networkingv1 "k8s.io/api/networking/v1" + metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" + "k8s.io/apimachinery/pkg/util/intstr" +) + +func (t *TestCaseGenerator) LoadBalancerTestCase() []*TestCase { + svc1 := &v1.Service{ + ObjectMeta: metav1.ObjectMeta{ + Name: "test-service-name", + Namespace: "x", + }, + Spec: v1.ServiceSpec{ + Type: v1.ServiceTypeNodePort, + Ports: []v1.ServicePort{ + { + Protocol: v1.ProtocolTCP, + Port: 81, + NodePort: 32086, + }, + }, + Selector: map[string]string{"pod": "a"}, + }, + } + probe := &ProbeConfig{ + AllAvailable: false, + PortProtocol: &PortProtocol{ + Protocol: v1.ProtocolTCP, + Port: intstr.FromInt(32086), + }, + Mode: ProbeModeNodeIP, + } + denyAll := &networkingv1.NetworkPolicy{ + ObjectMeta: metav1.ObjectMeta{ + Name: "deny-all-policy", + Namespace: "x", + }, + Spec: networkingv1.NetworkPolicySpec{ + PodSelector: metav1.LabelSelector{}, + PolicyTypes: []networkingv1.PolicyType{ + networkingv1.PolicyTypeIngress, + }, + }, + } + + // these tests cases are currently nonblocking the suite, + // but there is still work needed to measure the success/failure of the + // probe in the node+pod table + return []*TestCase{ + NewTestCase("should allow access to nodeport with no netpols applied", + NewStringSet(TagLoadBalancer), + NewTestStep(probe, + CreateService(svc1), + ), + ), + NewTestCase("should deny access to nodeport with netpols applied", + NewStringSet(TagLoadBalancer), + NewTestStep(probe, + CreatePolicy(denyAll), + ), + ), + } +} diff --git a/pkg/generator/tags.go b/pkg/generator/tags.go index 2f725ad..8a6988c 100644 --- a/pkg/generator/tags.go +++ b/pkg/generator/tags.go @@ -1,11 +1,12 @@ package generator import ( + "sort" + "strings" + "github.com/mattfenwick/collections/pkg/slice" "github.com/pkg/errors" "golang.org/x/exp/maps" - "sort" - "strings" ) const ( @@ -84,6 +85,7 @@ const ( TagConflict = "conflict" TagExample = "example" TagUpstreamE2E = "upstream-e2e" + TagLoadBalancer = "loadbalancer" ) var AllTags = map[string][]string{ @@ -143,6 +145,7 @@ var AllTags = map[string][]string{ TagConflict, TagExample, TagUpstreamE2E, + TagLoadBalancer, }, } diff --git a/pkg/generator/testcase.go b/pkg/generator/testcase.go index 3c9ba03..b091283 100644 --- a/pkg/generator/testcase.go +++ b/pkg/generator/testcase.go @@ -1,12 +1,14 @@ package generator import ( + "fmt" + "sort" + "strings" + "github.com/pkg/errors" v1 "k8s.io/api/core/v1" networkingv1 "k8s.io/api/networking/v1" "k8s.io/apimachinery/pkg/util/intstr" - "sort" - "strings" ) type TestCase struct { @@ -64,8 +66,12 @@ func (t *TestCase) collectActionsAndPolicies() (map[string]bool, []*networkingv1 features[ActionFeatureSetPodLabels] = true } else if action.DeletePod != nil { features[ActionFeatureDeletePod] = true + } else if action.CreateService != nil { + features[ActionFeatureCreateService] = true + } else if action.DeleteService != nil { + features[ActionFeatureDeleteService] = true } else { - panic("invalid Action") + panic(fmt.Sprintf("invalid Action: %v", action)) } } } @@ -114,6 +120,7 @@ const ( ProbeModeServiceName = "service-name" ProbeModeServiceIP = "service-ip" ProbeModePodIP = "pod-ip" + ProbeModeNodeIP = "node-ip" ) var AllProbeModes = []string{ @@ -135,8 +142,7 @@ func ParseProbeMode(mode string) (ProbeMode, error) { } // ProbeConfig: exactly one field must be non-null (or, in AllAvailable's case, non-false). This -// -// models a discriminated union (sum type). +// models a discriminated union (sum type). type ProbeConfig struct { AllAvailable bool PortProtocol *PortProtocol diff --git a/pkg/generator/testcasegenerator.go b/pkg/generator/testcasegenerator.go index d575797..97de01e 100644 --- a/pkg/generator/testcasegenerator.go +++ b/pkg/generator/testcasegenerator.go @@ -72,7 +72,8 @@ func (t *TestCaseGenerator) GenerateAllTestCases() []*TestCase { t.ActionTestCases(), t.ConflictTestCases(), t.NamespaceTestCases(), - t.UpstreamE2ETestCases()) + t.UpstreamE2ETestCases(), + t.LoadBalancerTestCase()) } func (t *TestCaseGenerator) GenerateTestCases() []*TestCase { diff --git a/pkg/generator/testcasegenerator_tests.go b/pkg/generator/testcasegenerator_tests.go index f2f0fb2..a99515f 100644 --- a/pkg/generator/testcasegenerator_tests.go +++ b/pkg/generator/testcasegenerator_tests.go @@ -20,7 +20,7 @@ func RunTestCaseGeneratorTests() { Expect(len(gen.ConflictTestCases())).To(Equal(16)) Expect(len(gen.NamespaceTestCases())).To(Equal(2)) - Expect(len(gen.GenerateTestCases())).To(Equal(230)) + Expect(len(gen.GenerateTestCases())).To(Equal(232)) }) }) } diff --git a/pkg/kube/ikubernetes.go b/pkg/kube/ikubernetes.go index 497277a..cefc91e 100644 --- a/pkg/kube/ikubernetes.go +++ b/pkg/kube/ikubernetes.go @@ -2,12 +2,13 @@ package kube import ( "fmt" + "math/rand" + "github.com/mattfenwick/cyclonus/pkg/utils" "github.com/pkg/errors" v1 "k8s.io/api/core/v1" networkingv1 "k8s.io/api/networking/v1" metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" - "math/rand" ) type IKubernetes interface { @@ -34,6 +35,8 @@ type IKubernetes interface { SetPodLabels(namespace string, pod string, labels map[string]string) (*v1.Pod, error) GetPodsInNamespace(namespace string) ([]v1.Pod, error) + GetNodes() (*v1.NodeList, error) + ExecuteRemoteCommand(namespace string, pod string, container string, command []string) (string, string, error, error) } @@ -104,6 +107,10 @@ func NewMockKubernetes(passRate float64) *MockKubernetes { } } +func (m *MockKubernetes) GetNodes() (*v1.NodeList, error) { + return nil, errors.New("not implemented") +} + func (m *MockKubernetes) getNamespaceObject(namespace string) (*MockNamespace, error) { if ns, ok := m.Namespaces[namespace]; ok { return ns, nil diff --git a/pkg/kube/kubernetes.go b/pkg/kube/kubernetes.go index ba6d066..f6601fd 100644 --- a/pkg/kube/kubernetes.go +++ b/pkg/kube/kubernetes.go @@ -3,6 +3,7 @@ package kube import ( "bytes" "context" + "github.com/pkg/errors" log "github.com/sirupsen/logrus" v1 "k8s.io/api/core/v1" @@ -39,6 +40,11 @@ func NewKubernetesForContext(context string) (*Kubernetes, error) { }, nil } +func (k *Kubernetes) GetNodes() (*v1.NodeList, error) { + nodes, err := k.ClientSet.CoreV1().Nodes().List(context.TODO(), metav1.ListOptions{}) + return nodes, errors.Wrapf(err, "unable to get nodes") +} + func (k *Kubernetes) GetNamespace(namespace string) (*v1.Namespace, error) { ns, err := k.ClientSet.CoreV1().Namespaces().Get(context.TODO(), namespace, metav1.GetOptions{}) return ns, errors.Wrapf(err, "unable to get namespace %s", namespace) diff --git a/pkg/worker/worker.go b/pkg/worker/worker.go index be40179..afb0308 100644 --- a/pkg/worker/worker.go +++ b/pkg/worker/worker.go @@ -2,9 +2,10 @@ package worker import ( "encoding/json" + "os/exec" + "github.com/pkg/errors" v1 "k8s.io/api/core/v1" - "os/exec" ) var (